- 1
//! Durable operational incident records. - 2
//! - 3
//! The control-plane snapshot is sampled, but incidents must survive a page - 4
//! refresh and a server restart. This module stores an append-only event log - 5
//! under the server's sessions home and folds it into the current incident - 6
//! state. A repeated observation updates an incident at most once per minute; - 7
//! a signal disappearing records a resolution event instead of deleting the - 8
//! object. - 9
- 10
use std::collections::HashMap; - 11
use std::io::Write; - 12
use std::path::{Path, PathBuf}; - 13
- 14
use chrono::{DateTime, Duration, Utc}; - 15
use serde::{Deserialize, Serialize}; - 16
use uuid::Uuid; - 17
- 18
fn log_path(home: &Path) -> PathBuf { - 19
home.join("operations").join("incidents.jsonl") - 20
} - 21
- 22
fn actions_path(home: &Path) -> PathBuf { - 23
home.join("operations").join("actions.jsonl") - 24
} - 25
- 26
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] - 27
pub struct IncidentCandidate { - 28
pub fingerprint: String, - 29
pub severity: String, - 30
pub source: String, - 31
pub title: String, - 32
pub detail: String, - 33
pub workspace: Option<String>, - 34
pub evidence: Vec<String>, - 35
} - 36
- 37
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] - 38
pub struct IncidentRecord { - 39
pub id: String, - 40
pub fingerprint: String, - 41
pub severity: String, - 42
pub status: String, - 43
pub source: String, - 44
pub title: String, - 45
pub detail: String, - 46
pub first_seen: DateTime<Utc>, - 47
pub last_seen: DateTime<Utc>, - 48
pub occurrences: u64, - 49
pub workspace: Option<String>, - 50
pub evidence: Vec<String>, - 51
pub resolution: Option<String>, - 52
} - 53
- 54
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] - 55
pub struct ActionVerification { - 56
pub status: String, - 57
pub before: String, - 58
pub after: String, - 59
pub detail: String, - 60
} - 61
- 62
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] - 63
pub struct ActionReceipt { - 64
pub receipt_id: String, - 65
pub service: String, - 66
pub action: String, - 67
pub requested_at: DateTime<Utc>, - 68
pub completed_at: DateTime<Utc>, - 69
pub succeeded: bool, - 70
pub verification: ActionVerification, - 71
pub persisted: bool, - 72
} - 73
- 74
#[derive(Debug, Clone, Serialize, Deserialize)] - 75
struct IncidentEvent { - 76
ts: DateTime<Utc>, - 77
event: String, - 78
incident: IncidentRecord, - 79
} - 80
- 81
fn read_events(home: &Path) -> Vec<IncidentEvent> { - 82
let Ok(raw) = std::fs::read_to_string(log_path(home)) else { - 83
return Vec::new(); - 84
}; - 85
raw.lines() - 86
.filter(|line| !line.trim().is_empty()) - 87
.filter_map(|line| serde_json::from_str(line).ok()) - 88
.collect() - 89
} - 90
- 91
fn folded(home: &Path) -> Vec<IncidentRecord> { - 92
let mut records = HashMap::<String, IncidentRecord>::new(); - 93
for event in read_events(home) { - 94
records.insert(event.incident.id.clone(), event.incident); - 95
} - 96
let mut values: Vec<_> = records.into_values().collect(); - 97
values.sort_by_key(|record| std::cmp::Reverse(record.last_seen)); - 98
values - 99
} - 100
- 101
fn append(home: &Path, event: &IncidentEvent) { - 102
let path = log_path(home); - 103
let Some(parent) = path.parent() else { return }; - 104
if std::fs::create_dir_all(parent).is_err() { - 105
return; - 106
} - 107
let Ok(mut line) = serde_json::to_vec(event) else { - 108
return; - 109
}; - 110
line.push(b'\n'); - 111
if let Ok(mut file) = std::fs::OpenOptions::new() - 112
.create(true) - 113
.append(true) - 114
.open(path) - 115
{ - 116
let _ = file.write_all(&line); - 117
} - 118
} - 119
- 120
/// Persist an operation receipt. Receipts are append-only so an operator can - 121
/// prove what was requested, what the service manager returned, and what a - 122
/// follow-up state probe observed. - 123
pub fn record_action(home: &Path, receipt: &ActionReceipt) -> Result<(), String> { - 124
let path = actions_path(home); - 125
let Some(parent) = path.parent() else { - 126
return Err("operations path has no parent".to_string()); - 127
}; - 128
std::fs::create_dir_all(parent) - 129
.map_err(|error| format!("create operations directory: {error}"))?; - 130
let mut line = - 131
serde_json::to_vec(receipt).map_err(|error| format!("encode receipt: {error}"))?; - 132
line.push(b'\n'); - 133
let mut file = std::fs::OpenOptions::new() - 134
.create(true) - 135
.append(true) - 136
.open(path) - 137
.map_err(|error| format!("open action ledger: {error}"))?; - 138
file.write_all(&line) - 139
.map_err(|error| format!("append action receipt: {error}")) - 140
} - 141
- 142
/// Return the newest durable action receipts, newest first. - 143
pub fn recent_actions(home: &Path, limit: usize) -> Vec<ActionReceipt> { - 144
let Ok(raw) = std::fs::read_to_string(actions_path(home)) else { - 145
return Vec::new(); - 146
}; - 147
raw.lines() - 148
.rev() - 149
.filter(|line| !line.trim().is_empty()) - 150
.filter_map(|line| serde_json::from_str::<ActionReceipt>(line).ok()) - 151
.take(limit) - 152
.collect() - 153
} - 154
- 155
fn new_record(candidate: IncidentCandidate, now: DateTime<Utc>) -> IncidentRecord { - 156
IncidentRecord { - 157
id: format!("INC-{}", Uuid::now_v7().simple()), - 158
fingerprint: candidate.fingerprint, - 159
severity: candidate.severity, - 160
status: "investigating".to_string(), - 161
source: candidate.source, - 162
title: candidate.title, - 163
detail: candidate.detail, - 164
first_seen: now, - 165
last_seen: now, - 166
occurrences: 1, - 167
workspace: candidate.workspace, - 168
evidence: candidate.evidence, - 169
resolution: None, - 170
} - 171
} - 172
- 173
/// Reconcile the current probe candidates into durable incident records. - 174
/// Returns open records first, followed by the most recently resolved records. - 175
pub fn reconcile(home: &Path, candidates: Vec<IncidentCandidate>) -> Vec<IncidentRecord> { - 176
let now = Utc::now(); - 177
let mut all = folded(home); - 178
let mut by_fingerprint: HashMap<String, usize> = all - 179
.iter() - 180
.enumerate() - 181
.map(|(index, record)| (record.fingerprint.clone(), index)) - 182
.collect(); - 183
let mut seen = std::collections::HashSet::new(); - 184
- 185
for candidate in candidates { - 186
seen.insert(candidate.fingerprint.clone()); - 187
if let Some(index) = by_fingerprint.get(&candidate.fingerprint).copied() { - 188
let record = &mut all[index]; - 189
let was_resolved = record.status == "resolved"; - 190
let due = now.signed_duration_since(record.last_seen) >= Duration::minutes(1); - 191
record.severity = candidate.severity; - 192
record.source = candidate.source; - 193
record.title = candidate.title; - 194
record.detail = candidate.detail; - 195
record.workspace = candidate.workspace; - 196
record.evidence = candidate.evidence; - 197
record.status = "investigating".to_string(); - 198
record.resolution = None; - 199
if was_resolved || due { - 200
record.last_seen = now; - 201
record.occurrences = record.occurrences.saturating_add(1); - 202
append( - 203
home, - 204
&IncidentEvent { - 205
ts: now, - 206
event: if was_resolved { "reopened" } else { "observed" }.to_string(), - 207
incident: record.clone(), - 208
}, - 209
); - 210
} - 211
} else { - 212
let record = new_record(candidate, now); - 213
append( - 214
home, - 215
&IncidentEvent { - 216
ts: now, - 217
event: "opened".to_string(), - 218
incident: record.clone(), - 219
}, - 220
); - 221
by_fingerprint.insert(record.fingerprint.clone(), all.len()); - 222
all.push(record); - 223
} - 224
} - 225
- 226
for record in &mut all { - 227
if record.status != "resolved" && !seen.contains(&record.fingerprint) { - 228
record.status = "resolved".to_string(); - 229
record.resolution = - 230
Some("No longer present in current control-plane probes".to_string()); - 231
record.last_seen = now; - 232
append( - 233
home, - 234
&IncidentEvent { - 235
ts: now, - 236
event: "resolved".to_string(), - 237
incident: record.clone(), - 238
}, - 239
); - 240
} - 241
} - 242
- 243
all.sort_by(|a, b| { - 244
(a.status != "investigating") - 245
.cmp(&(b.status != "investigating")) - 246
.then_with(|| b.last_seen.cmp(&a.last_seen)) - 247
}); - 248
all - 249
} - 250
- 251
/// Read the folded incident history without creating new records. - 252
pub fn list(home: &Path) -> Vec<IncidentRecord> { - 253
folded(home) - 254
} - 255
- 256
#[cfg(test)] - 257
#[allow(clippy::unwrap_used, clippy::expect_used)] - 258
mod tests { - 259
use super::*; - 260
- 261
fn candidate(fingerprint: &str) -> IncidentCandidate { - 262
IncidentCandidate { - 263
fingerprint: fingerprint.to_string(), - 264
severity: "warning".to_string(), - 265
source: "test".to_string(), - 266
title: "Test incident".to_string(), - 267
detail: "evidence".to_string(), - 268
workspace: Some("/tmp/project".to_string()), - 269
evidence: vec!["check:test".to_string()], - 270
} - 271
} - 272
- 273
#[test] - 274
fn opens_and_folds_a_durable_incident() { - 275
let dir = tempfile::tempdir().unwrap(); - 276
let rows = reconcile(dir.path(), vec![candidate("test")]); - 277
assert_eq!(rows.len(), 1); - 278
assert!(rows[0].id.starts_with("INC-")); - 279
assert_eq!(rows[0].status, "investigating"); - 280
assert_eq!(list(dir.path()), rows); - 281
} - 282
- 283
#[test] - 284
fn repeated_observation_groups_instead_of_creating_duplicates() { - 285
let dir = tempfile::tempdir().unwrap(); - 286
let first = reconcile(dir.path(), vec![candidate("test")]); - 287
let second = reconcile(dir.path(), vec![candidate("test")]); - 288
assert_eq!(second.len(), 1); - 289
assert_eq!(second[0].id, first[0].id); - 290
assert_eq!(second[0].occurrences, 1); - 291
} - 292
- 293
#[test] - 294
fn disappearance_records_resolution_and_reappearance_reopens() { - 295
let dir = tempfile::tempdir().unwrap(); - 296
let first = reconcile(dir.path(), vec![candidate("test")]); - 297
let resolved = reconcile(dir.path(), Vec::new()); - 298
assert_eq!(resolved[0].id, first[0].id); - 299
assert_eq!(resolved[0].status, "resolved"); - 300
let reopened = reconcile(dir.path(), vec![candidate("test")]); - 301
assert_eq!(reopened[0].id, first[0].id); - 302
assert_eq!(reopened[0].status, "investigating"); - 303
assert_eq!(reopened[0].occurrences, 2); - 304
} - 305
- 306
#[test] - 307
fn action_receipts_round_trip_newest_first() { - 308
let dir = tempfile::tempdir().unwrap(); - 309
let receipt = ActionReceipt { - 310
receipt_id: "OP-test".to_string(), - 311
service: "gateway".to_string(), - 312
action: "restart".to_string(), - 313
requested_at: Utc::now(), - 314
completed_at: Utc::now(), - 315
succeeded: true, - 316
verification: ActionVerification { - 317
status: "verified".to_string(), - 318
before: "stopped".to_string(), - 319
after: "running".to_string(), - 320
detail: "probe reached running".to_string(), - 321
}, - 322
persisted: true, - 323
}; - 324
record_action(dir.path(), &receipt).unwrap(); - 325
assert_eq!(recent_actions(dir.path(), 10), vec![receipt]); - 326
} - 327
} - 328
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.