1use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
25use std::sync::mpsc::{Receiver, RecvTimeoutError, SyncSender};
26use std::sync::{Arc, Mutex};
27use std::time::Duration;
28
29use rusqlite::{OptionalExtension, params};
30use serde::{Deserialize, Serialize};
31use serde_json::{Map, Value, json};
32
33use crate::audit::{
34 AuditLog, Entry, GENESIS, Query, Verified, Visibility, canonical, clip, db_err, hex, meta,
35 set_meta, sql_glob,
36};
37use crate::error::Result;
38
39pub const DEFAULT_RETENTION: Duration = Duration::from_secs(365 * 86400);
41pub const DEFAULT_MAX_ROWS: i64 = 5_000_000;
43const QUEUE: usize = 10_000;
45
46#[derive(Debug, Clone, Default)]
48pub struct NewRecord {
49 pub time: i64,
51 pub source: String,
53 pub org: Option<String>,
55 pub project: Option<String>,
57 pub kind: String,
60 pub object_type: Option<String>,
62 pub object: Option<String>,
63 pub objects: Vec<String>,
65 pub actor: Option<String>,
67 pub level: Option<String>,
68 pub message: Option<String>,
69 pub details: Value,
70}
71
72#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
74pub struct Record {
75 pub id: i64,
76 pub time: i64,
77 pub source: String,
78 pub org: Option<String>,
79 pub project: Option<String>,
80 pub kind: String,
81 pub object_type: Option<String>,
82 pub object: Option<String>,
83 pub objects: String,
84 pub actor: Option<String>,
85 pub level: Option<String>,
86 pub message: Option<String>,
87 pub details: Value,
88 pub prev_hash: String,
89 pub hash: String,
90}
91
92const COLS: &str = "id, time, source, org, project, kind, object_type, object, objects, actor, \
93 level, message, details, prev_hash, hash";
94
95fn row(r: &rusqlite::Row) -> rusqlite::Result<Record> {
96 let details: String = r.get(12)?;
97 Ok(Record {
98 id: r.get(0)?,
99 time: r.get(1)?,
100 source: r.get(2)?,
101 org: r.get(3)?,
102 project: r.get(4)?,
103 kind: r.get(5)?,
104 object_type: r.get(6)?,
105 object: r.get(7)?,
106 objects: r.get(8)?,
107 actor: r.get(9)?,
108 level: r.get(10)?,
109 message: r.get(11)?,
110 details: serde_json::from_str(&details).unwrap_or(Value::Null),
111 prev_hash: r.get(13)?,
112 hash: r.get(14)?,
113 })
114}
115
116fn row_hash(e: &Record) -> String {
117 let body = json!({
118 "id": e.id, "time": e.time, "source": e.source, "org": e.org,
119 "project": e.project, "kind": e.kind, "object_type": e.object_type,
120 "object": e.object, "objects": e.objects, "actor": e.actor,
121 "level": e.level, "message": e.message, "details": e.details,
122 });
123 let mut ctx = ring::digest::Context::new(&ring::digest::SHA256);
124 ctx.update(e.prev_hash.as_bytes());
125 ctx.update(b"\n");
126 ctx.update(canonical(&body).as_bytes());
127 hex(ctx.finish().as_ref())
128}
129
130#[derive(Debug, Clone, Default, Deserialize)]
132#[serde(default, deny_unknown_fields)]
133pub struct HistoryQuery {
134 pub org: Option<String>,
136 pub platform: bool,
137 pub object: Option<String>,
141 pub exact: bool,
142 pub kind: Option<String>,
144 pub source: Option<String>,
147 pub actor: Option<String>,
149 pub since: Option<i64>,
151 pub until: Option<i64>,
153 pub before: Option<String>,
155 pub limit: Option<usize>,
157 pub ascending: bool,
159 pub correlate: bool,
161}
162
163impl HistoryQuery {
164 pub fn wants(&self, source: &str) -> bool {
165 match &self.source {
166 None => true,
167 Some(s) if s.trim().is_empty() => true,
168 Some(s) => s.split(',').any(|x| x.trim() == source),
169 }
170 }
171}
172
173#[derive(Debug, Clone, PartialEq, Serialize)]
175pub struct Item {
176 pub source: String,
177 pub id: i64,
178 pub time: i64,
179 pub org: Option<String>,
180 pub kind: String,
181 pub object_type: Option<String>,
182 pub object: Option<String>,
183 pub actor: Option<String>,
184 pub level: Option<String>,
186 pub message: Option<String>,
187 pub details: Value,
188 #[serde(skip_serializing_if = "Option::is_none")]
191 pub inferred: Option<Value>,
192}
193
194impl Item {
195 pub fn from_record(r: Record) -> Item {
196 Item {
197 source: r.source,
198 id: r.id,
199 time: r.time,
200 org: r.org,
201 kind: r.kind,
202 object_type: r.object_type,
203 object: r.object,
204 actor: r.actor,
205 level: r.level,
206 message: r.message,
207 details: r.details,
208 inferred: None,
209 }
210 }
211
212 pub fn from_audit(e: Entry) -> Item {
213 let actor = match &e.token_name {
214 Some(t) => format!("{} (token {t})", e.actor),
215 None => e.actor.clone(),
216 };
217 let mut details = match e.details {
218 Value::Object(m) => m,
219 _ => Map::new(),
220 };
221 details.insert("surface".into(), json!(e.surface));
222 details.insert("actor_kind".into(), json!(e.actor_kind));
223 if let Some(ip) = e.ip {
224 details.insert("ip".into(), json!(ip));
225 }
226 if let Some(r) = e.request_id {
227 details.insert("request_id".into(), json!(r));
228 }
229 Item {
230 source: "audit".into(),
231 id: e.id,
232 time: e.time,
233 org: e.org,
234 kind: e.action,
235 object_type: None,
236 object: e.target,
237 actor: Some(actor),
238 level: Some(e.outcome),
239 message: None,
240 details: Value::Object(details),
241 inferred: None,
242 }
243 }
244
245 fn key(&self) -> (i64, u8, i64) {
247 let rank = match self.source.as_str() {
248 "audit" => 3,
249 "controller" => 2,
250 "incus" => 1,
251 _ => 0,
252 };
253 (self.time, rank, self.id)
254 }
255
256 fn cursor(&self) -> String {
257 let (t, s, i) = self.key();
258 format!("{t}.{s}.{i}")
259 }
260}
261
262fn parse_cursor(s: &str) -> Option<(i64, u8, i64)> {
263 let mut p = s.split('.');
264 let t = p.next()?.parse().ok()?;
265 let r = p.next()?.parse().ok()?;
266 let i = p.next()?.parse().ok()?;
267 Some((t, r, i))
268}
269
270#[derive(Debug, Clone, Serialize)]
272pub struct Page {
273 pub items: Vec<Item>,
274 pub next: Option<String>,
275}
276
277impl AuditLog {
278 pub fn history_append(&self, n: NewRecord) -> Result<Record> {
280 let now = (self.clock)();
281 let mut objects: Vec<String> = n
282 .objects
283 .into_iter()
284 .chain(n.object.clone())
285 .filter(|o| !o.is_empty())
286 .map(|o| clip(o, 128))
287 .collect();
288 objects.sort();
289 objects.dedup();
290 let mut e = Record {
291 id: 0,
292 time: if n.time > 0 { n.time } else { now },
293 source: clip(n.source, 16),
294 org: n.org.map(|o| clip(o, 64)),
295 project: n.project.map(|o| clip(o, 64)),
296 kind: clip(n.kind, 128),
297 object_type: n.object_type.map(|o| clip(o, 64)),
298 object: n.object.map(|o| clip(o, 256)),
299 objects: objects.join(" "),
300 actor: n.actor.map(|a| clip(a, 256)),
301 level: n.level.map(|l| clip(l, 16)),
302 message: n.message.map(|m| clip(m, 2000)),
303 details: match n.details {
304 Value::Null => json!({}),
305 v => v,
306 },
307 prev_hash: String::new(),
308 hash: String::new(),
309 };
310 let db = self.db();
311 db.execute_batch("BEGIN IMMEDIATE")
312 .map_err(|err| db_err("history append", err))?;
313 let r = (|| -> rusqlite::Result<()> {
314 let head: i64 = meta(&db, "history_head_id")?
315 .and_then(|v| v.parse().ok())
316 .unwrap_or(0);
317 e.prev_hash = meta(&db, "history_head_hash")?.unwrap_or_else(|| GENESIS.into());
318 e.id = head + 1;
319 e.hash = row_hash(&e);
320 db.execute(
321 &format!(
322 "INSERT INTO history ({COLS}) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, \
323 ?10, ?11, ?12, ?13, ?14, ?15)"
324 ),
325 params![
326 e.id,
327 e.time,
328 e.source,
329 e.org,
330 e.project,
331 e.kind,
332 e.object_type,
333 e.object,
334 e.objects,
335 e.actor,
336 e.level,
337 e.message,
338 e.details.to_string(),
339 e.prev_hash,
340 e.hash
341 ],
342 )?;
343 set_meta(&db, "history_head_id", &e.id.to_string())?;
344 set_meta(&db, "history_head_hash", &e.hash)?;
345 Ok(())
346 })();
347 match r {
348 Ok(()) => db
349 .execute_batch("COMMIT")
350 .map_err(|err| db_err("history append", err))?,
351 Err(err) => {
352 let _ = db.execute_batch("ROLLBACK");
353 return Err(db_err("history append", err));
354 }
355 }
356 drop(db);
357 self.appended_one();
358 Ok(e)
359 }
360
361 pub fn history_last_time(&self) -> Result<Option<i64>> {
363 let db = self.db();
364 db.query_row("SELECT MAX(time) FROM history", [], |r| {
365 r.get::<_, Option<i64>>(0)
366 })
367 .map_err(|e| db_err("history", e))
368 }
369
370 pub(crate) fn prune_history(&self) -> Result<usize> {
373 let cutoff = (self.clock)() - self.history_retention.as_millis() as i64;
374 let max = self.history_max_rows;
375 let db = self.db();
376 db.execute_batch("BEGIN IMMEDIATE")
377 .map_err(|e| db_err("history prune", e))?;
378 let r = (|| -> rusqlite::Result<usize> {
379 let by_age: Option<i64> = db.query_row(
380 "SELECT MAX(id) FROM history WHERE time < ?1",
381 [cutoff],
382 |r| r.get(0),
383 )?;
384 let head: Option<i64> =
385 db.query_row("SELECT MAX(id) FROM history", [], |r| r.get(0))?;
386 let count: i64 = db.query_row("SELECT COUNT(*) FROM history", [], |r| r.get(0))?;
387 let by_size = (count > max).then(|| head.unwrap_or(0) - max);
388 let Some(through) = by_age.into_iter().chain(by_size).max() else {
389 return Ok(0);
390 };
391 let hash: Option<String> = db
392 .query_row("SELECT hash FROM history WHERE id = ?1", [through], |r| {
393 r.get(0)
394 })
395 .optional()?;
396 let Some(hash) = hash else { return Ok(0) };
397 set_meta(&db, "pruning", "1")?;
398 let n = db.execute("DELETE FROM history WHERE id <= ?1", [through])?;
399 set_meta(&db, "pruning", "0")?;
400 set_meta(&db, "history_pruned_through", &through.to_string())?;
401 set_meta(&db, "history_pruned_hash", &hash)?;
402 Ok(n)
403 })();
404 match r {
405 Ok(n) => {
406 db.execute_batch("COMMIT")
407 .map_err(|e| db_err("history prune", e))?;
408 Ok(n)
409 }
410 Err(err) => {
411 let _ = db.execute_batch("ROLLBACK");
412 Err(db_err("history prune", err))
413 }
414 }
415 }
416
417 pub fn history_list(
420 &self,
421 q: &HistoryQuery,
422 vis: &Visibility,
423 upto: Option<i64>,
424 from: Option<i64>,
425 limit: usize,
426 ) -> Result<Vec<Record>> {
427 use rusqlite::types::Value as V;
428 let mut sql = format!("SELECT {COLS} FROM history WHERE 1=1");
429 let mut args: Vec<V> = Vec::new();
430 let mut push = |sql: &mut String, cond: &str, v: V| {
431 args.push(v);
432 sql.push_str(&cond.replace('?', &format!("?{}", args.len())));
433 };
434 if let Visibility::Orgs(orgs) = vis {
435 if orgs.is_empty() {
436 return Ok(Vec::new());
437 }
438 let list: Vec<String> = orgs
439 .iter()
440 .map(|o| format!("'{}'", o.replace('\'', "")))
441 .collect();
442 sql.push_str(&format!(" AND org IN ({})", list.join(",")));
443 }
444 if q.platform {
445 sql.push_str(" AND org IS NULL");
446 } else if let Some(o) = &q.org {
447 push(&mut sql, " AND org = ?", V::Text(o.clone()));
448 }
449 let sources: Vec<&str> = ["controller", "incus", "marker"]
450 .into_iter()
451 .filter(|s| q.wants(s))
452 .collect();
453 if sources.is_empty() {
454 return Ok(Vec::new());
455 }
456 sql.push_str(&format!(
457 " AND source IN ({})",
458 sources
459 .iter()
460 .map(|s| format!("'{s}'"))
461 .collect::<Vec<_>>()
462 .join(",")
463 ));
464 if let Some(o) = q.object.as_deref().map(str::trim).filter(|o| !o.is_empty()) {
465 if q.exact {
466 push(
467 &mut sql,
468 " AND ((' ' || objects || ' ') LIKE ? ESCAPE '\\' OR source = 'marker')",
469 V::Text(format!("% {} %", like_escape(o))),
470 );
471 } else {
472 push(
473 &mut sql,
474 " AND (objects LIKE ? ESCAPE '\\' OR source = 'marker')",
475 V::Text(format!("%{}%", like_escape(o))),
476 );
477 }
478 }
479 if let Some(k) = &q.kind {
480 push(&mut sql, " AND kind GLOB ?", V::Text(sql_glob(k)));
481 }
482 if let Some(a) = &q.actor {
483 push(
484 &mut sql,
485 " AND IFNULL(actor, '') GLOB ?",
486 V::Text(sql_glob(a)),
487 );
488 }
489 if let Some(t) = q.since {
490 push(&mut sql, " AND time >= ?", V::Integer(t));
491 }
492 if let Some(t) = q.until {
493 push(&mut sql, " AND time < ?", V::Integer(t));
494 }
495 if let Some(t) = upto {
496 push(&mut sql, " AND time <= ?", V::Integer(t));
497 }
498 if let Some(t) = from {
499 push(&mut sql, " AND time >= ?", V::Integer(t));
500 }
501 sql.push_str(if q.ascending {
502 " ORDER BY time ASC, id ASC"
503 } else {
504 " ORDER BY time DESC, id DESC"
505 });
506 sql.push_str(&format!(" LIMIT {}", limit.clamp(1, 5000)));
507 let db = self.db();
508 let mut st = db.prepare(&sql).map_err(|e| db_err("history query", e))?;
509 let rows = st
510 .query_map(rusqlite::params_from_iter(args), row)
511 .map_err(|e| db_err("history query", e))?;
512 rows.collect::<rusqlite::Result<_>>()
513 .map_err(|e| db_err("history query", e))
514 }
515
516 pub fn history_after(&self, after: i64, vis: &Visibility, limit: usize) -> Result<Vec<Record>> {
518 let db = self.db();
519 let mut st = db
520 .prepare(&format!(
521 "SELECT {COLS} FROM history WHERE id > ?1 ORDER BY id ASC LIMIT ?2"
522 ))
523 .map_err(|e| db_err("history query", e))?;
524 let rows = st
525 .query_map(params![after, limit as i64], row)
526 .map_err(|e| db_err("history query", e))?;
527 let all: Vec<Record> = rows
528 .collect::<rusqlite::Result<_>>()
529 .map_err(|e| db_err("history query", e))?;
530 Ok(all.into_iter().filter(|r| visible(vis, &r.org)).collect())
531 }
532
533 pub fn history_head(&self) -> Result<i64> {
535 let db = self.db();
536 Ok(meta(&db, "history_head_id")
537 .map_err(|e| db_err("history head", e))?
538 .and_then(|v| v.parse().ok())
539 .unwrap_or(0))
540 }
541
542 pub fn history_verify(&self) -> Result<Verified> {
544 let db = self.db();
545 let m = |k: &str| meta(&db, k).map_err(|e| db_err("history verify", e));
546 let pruned_through: Option<i64> = m("history_pruned_through")?.and_then(|v| v.parse().ok());
547 let mut expect = m("history_pruned_hash")?.unwrap_or_else(|| GENESIS.into());
548 let head_id: Option<i64> = m("history_head_id")?.and_then(|v| v.parse().ok());
549 let head_hash = m("history_head_hash")?;
550 let mut st = db
551 .prepare(&format!("SELECT {COLS} FROM history ORDER BY id ASC"))
552 .map_err(|e| db_err("history verify", e))?;
553 let rows = st
554 .query_map([], row)
555 .map_err(|e| db_err("history verify", e))?;
556 let mut out = Verified {
557 ok: true,
558 rows: 0,
559 head: None,
560 pruned_through,
561 broken: None,
562 };
563 let mut last = pruned_through.unwrap_or(0);
564 for r in rows {
565 let e = r.map_err(|e| db_err("history verify", e))?;
566 out.rows += 1;
567 let why = if e.id != last + 1 {
568 Some(format!("expected row {} next, found {}", last + 1, e.id))
569 } else if e.prev_hash != expect {
570 Some("prev_hash does not match the row before it".to_string())
571 } else if row_hash(&e) != e.hash {
572 Some("the row's contents do not match its hash".to_string())
573 } else {
574 None
575 };
576 if let Some(w) = why {
577 out.ok = false;
578 out.broken = Some((e.id, w));
579 return Ok(out);
580 }
581 expect = e.hash.clone();
582 last = e.id;
583 out.head = Some((e.id, e.hash));
584 }
585 let tail_ok = match (&out.head, head_id, &head_hash) {
586 (Some((id, h)), Some(hid), Some(hh)) => *id == hid && h == hh,
587 (None, Some(hid), _) => Some(hid) == pruned_through,
588 (None, None, None) => true,
589 _ => false,
590 };
591 if !tail_ok {
592 out.ok = false;
593 out.broken = Some((
594 head_id.unwrap_or(0),
595 "the newest rows are missing (the head does not match)".into(),
596 ));
597 }
598 Ok(out)
599 }
600
601 pub fn timeline(
604 &self,
605 q: &HistoryQuery,
606 hvis: &Visibility,
607 avis: Option<&Visibility>,
608 ) -> Result<Page> {
609 let limit = q.limit.unwrap_or(100).clamp(1, 1000);
610 let cursor = q.before.as_deref().and_then(parse_cursor);
611 let fetch = limit + 50;
613 let (upto, from) = match (cursor, q.ascending) {
614 (Some((t, _, _)), false) => (Some(t), None),
615 (Some((t, _, _)), true) => (None, Some(t)),
616 (None, _) => (None, None),
617 };
618 let mut items: Vec<Item> = self
619 .history_list(q, hvis, upto, from, fetch)?
620 .into_iter()
621 .map(Item::from_record)
622 .collect();
623 if let (Some(av), true) = (avis, q.wants("audit")) {
624 let aq = Query {
625 org: q.org.clone(),
626 platform: q.platform,
627 actor: q.actor.clone(),
628 action: q.kind.clone(),
629 object: q.object.clone(),
630 object_exact: q.exact,
631 since: match (q.since, from) {
632 (Some(a), Some(b)) => Some(a.max(b)),
633 (a, b) => a.or(b),
634 },
635 until: match (q.until, upto) {
636 (Some(a), Some(b)) => Some(a.min(b + 1)),
637 (a, Some(b)) => a.or(Some(b + 1)),
638 (a, None) => a,
639 },
640 ascending: q.ascending,
641 limit: Some(fetch.min(1000)),
642 ..Default::default()
643 };
644 items.extend(self.list(&aq, av)?.into_iter().map(Item::from_audit));
645 }
646 if q.ascending {
647 items.sort_by_key(|i| i.key());
648 if let Some(c) = cursor {
649 items.retain(|i| i.key() > c);
650 }
651 } else {
652 items.sort_by_key(|i| std::cmp::Reverse(i.key()));
653 if let Some(c) = cursor {
654 items.retain(|i| i.key() < c);
655 }
656 }
657 items.truncate(limit);
658 let next = (items.len() == limit)
659 .then(|| items.last().map(Item::cursor))
660 .flatten();
661 if q.correlate {
662 self.correlate(&mut items, avis)?;
663 }
664 Ok(Page { items, next })
665 }
666
667 fn correlate(&self, items: &mut [Item], avis: Option<&Visibility>) -> Result<()> {
670 let Some(av) = avis else { return Ok(()) };
671 for it in items.iter_mut() {
672 if it.source != "incus" || it.object_type.as_deref() != Some("instance") {
673 continue;
674 }
675 let Some(inst) = it.object.clone() else {
676 continue;
677 };
678 let q = Query {
679 org: it.org.clone(),
680 platform: it.org.is_none(),
681 since: Some(it.time - 120_000),
684 until: Some(it.time + 120_000),
685 outcome: Some("ok".into()),
686 limit: Some(200),
687 ..Default::default()
688 };
689 let mut related: Vec<Entry> = self
690 .list(&q, av)?
691 .into_iter()
692 .filter(|e| {
693 let names: Vec<&str> = e
694 .target
695 .iter()
696 .map(String::as_str)
697 .chain(
698 ["name", "app", "stack", "project", "service"]
699 .iter()
700 .filter_map(|k| e.details.get(*k).and_then(Value::as_str)),
701 )
702 .filter(|n| n.len() >= 2)
703 .collect();
704 e.action != "audit_list" && names.iter().any(|n| inst.contains(n))
705 })
706 .collect();
707 related.sort_by_key(|e| (e.time > it.time, (it.time - e.time).abs()));
709 if let Some(e) = related.into_iter().next() {
710 it.inferred = Some(json!({
711 "audit_id": e.id,
712 "action": e.action,
713 "actor": e.actor,
714 "seconds_before": (it.time - e.time) as f64 / 1000.0,
715 "why": "inferred by time and name: an audit row in the same org within two minutes (written when its call returned) naming this instance's stack or app",
716 }));
717 }
718 }
719 Ok(())
720 }
721}
722
723fn visible(vis: &Visibility, org: &Option<String>) -> bool {
724 match vis {
725 Visibility::All => true,
726 Visibility::Orgs(v) => org.as_ref().is_some_and(|o| v.contains(o)),
727 }
728}
729
730fn like_escape(s: &str) -> String {
731 s.replace('\\', "\\\\")
732 .replace('%', "\\%")
733 .replace('_', "\\_")
734}
735
736pub fn from_event(e: &crate::stack::controller::Event) -> Option<NewRecord> {
741 if e.level == "log" {
742 return None;
743 }
744 let (org, stack) = match e.stack.split_once('/') {
745 Some((o, s)) => (o.to_string(), s.to_string()),
746 None => (crate::org::DEFAULT_ORG.to_string(), e.stack.clone()),
747 };
748 let mut objects = vec![stack.clone()];
749 if !e.service.is_empty() {
750 objects.push(e.service.clone());
751 }
752 objects.extend(e.instance.clone());
753 let object_type = match e.kind.as_deref() {
755 Some(k) if k.starts_with("volume.") => "volume",
756 _ => "stack",
757 };
758 Some(NewRecord {
759 time: e.at as i64,
760 source: "controller".into(),
761 org: Some(org),
762 project: None,
763 kind: e.kind.clone().unwrap_or_else(|| "event".into()),
764 object_type: Some(object_type.into()),
765 object: Some(stack),
766 objects,
767 actor: Some("isb".into()),
768 level: Some(e.level.clone()),
769 message: Some(e.message.clone()),
770 details: json!({"seq": e.seq, "service": e.service, "instance": e.instance}),
771 })
772}
773
774fn secret_key(k: &str) -> bool {
776 let l = k.to_ascii_lowercase();
777 l.starts_with("environment.")
778 || l == "environment"
779 || l == "env"
780 || l == "user.isb.create-token"
781 || l.starts_with("cloud-init.")
782 || l == "user.user-data"
783 || l == "user.vendor-data"
784 || l.contains("secret")
785 || l.contains("password")
786 || l.contains("token")
787 || l.contains("passphrase")
788 || l.contains("private")
789}
790
791pub fn scrub(v: &Value, depth: usize) -> Value {
794 match v {
795 Value::Object(m) if depth < 4 => Value::Object(
796 m.iter()
797 .filter(|(k, _)| !secret_key(k))
798 .take(64)
799 .map(|(k, v)| match (k.as_str(), v) {
800 ("command" | "argv" | "args", Value::Array(a)) => (
803 k.clone(),
804 json!({
805 "program": a.first().and_then(Value::as_str).map(|p| {
806 clip(p.rsplit('/').next().unwrap_or(p).to_string(), 64)
807 }),
808 "args": a.len().saturating_sub(1),
809 }),
810 ),
811 _ => (k.clone(), scrub(v, depth + 1)),
812 })
813 .collect(),
814 ),
815 Value::Object(_) => json!("…"),
816 Value::Array(a) if depth < 4 => {
817 Value::Array(a.iter().take(32).map(|v| scrub(v, depth + 1)).collect())
818 }
819 Value::Array(_) => json!("…"),
820 Value::String(s) => json!(clip(s.clone(), 256)),
821 other => other.clone(),
822 }
823}
824
825pub fn parse_source(src: &str) -> (Option<String>, Option<String>, Option<String>) {
827 let (path, query) = src.split_once('?').unwrap_or((src, ""));
828 let project = query.split('&').find_map(|kv| {
829 kv.strip_prefix("project=")
830 .map(|p| percent_decode(p).to_string())
831 });
832 let parts: Vec<String> = path
833 .trim_start_matches("/1.0/")
834 .split('/')
835 .map(percent_decode)
836 .collect();
837 let p: Vec<&str> = parts.iter().map(String::as_str).collect();
838 let (ty, name) = match p.as_slice() {
839 ["instances", n, "snapshots", s, ..] => ("instance-snapshot", format!("{n}/{s}")),
840 ["instances", n, ..] => ("instance", n.to_string()),
841 ["images", "aliases", n, ..] => ("image-alias", n.to_string()),
842 ["images", n, ..] => ("image", n.to_string()),
843 ["storage-pools", pool, "volumes", _, v, "snapshots", s, ..] => {
844 ("storage-volume-snapshot", format!("{pool}/{v}/{s}"))
845 }
846 ["storage-pools", _, "volumes", _, v, ..] => ("storage-volume", v.to_string()),
847 ["storage-pools", n, ..] => ("storage-pool", n.to_string()),
848 ["networks", n, ..] => ("network", n.to_string()),
849 ["network-acls", n, ..] => ("network-acl", n.to_string()),
850 ["network-zones", n, ..] => ("network-zone", n.to_string()),
851 ["profiles", n, ..] => ("profile", n.to_string()),
852 ["projects", n, ..] => ("project", n.to_string()),
853 [first, n, ..] => (first.strip_suffix('s').unwrap_or(first), n.to_string()),
854 [first] => (first.strip_suffix('s').unwrap_or(first), String::new()),
855 [] => ("", String::new()),
856 };
857 (
858 (!ty.is_empty()).then(|| ty.to_string()),
859 (!name.is_empty()).then_some(name),
860 project,
861 )
862}
863
864fn percent_decode(s: &str) -> String {
865 let b = s.as_bytes();
866 let mut out = Vec::with_capacity(b.len());
867 let mut i = 0;
868 while i < b.len() {
869 if b[i] == b'%' && i + 2 < b.len() {
870 if let Ok(v) = u8::from_str_radix(&s[i + 1..i + 3], 16) {
871 out.push(v);
872 i += 3;
873 continue;
874 }
875 }
876 out.push(b[i]);
877 i += 1;
878 }
879 String::from_utf8_lossy(&out).into_owned()
880}
881
882pub fn project_org(project: &str) -> Option<String> {
885 if project == crate::registry::PROJECT {
886 return None;
887 }
888 project
889 .strip_prefix("isb-")
890 .and_then(|o| crate::org::OrgId::new(o).ok())
891 .map(|o| o.to_string())
892}
893
894pub fn from_incus(v: &Value, received: i64) -> Option<NewRecord> {
897 if v.get("type").and_then(Value::as_str) != Some("lifecycle") {
898 return None;
899 }
900 let md = v.get("metadata")?;
901 let action = md.get("action").and_then(Value::as_str)?.to_string();
902 let source = md.get("source").and_then(Value::as_str).unwrap_or("");
903 let (object_type, object, src_project) = parse_source(source);
904 let project = v
905 .get("project")
906 .and_then(Value::as_str)
907 .filter(|p| !p.is_empty())
908 .or_else(|| md.get("project").and_then(Value::as_str))
909 .map(String::from)
910 .or(src_project)
911 .unwrap_or_else(|| "default".into());
912 let host_level = matches!(
914 object_type.as_deref(),
915 Some("image" | "image-alias" | "storage-pool" | "network" | "network-zone" | "project")
916 );
917 let org = (!host_level).then(|| project_org(&project)).flatten();
918 let req = md.get("requestor").cloned().unwrap_or(Value::Null);
919 let actor = req
920 .get("username")
921 .and_then(Value::as_str)
922 .filter(|u| !u.is_empty())
923 .map(|u| match req.get("protocol").and_then(Value::as_str) {
924 Some(p) if p != "unix" && !p.is_empty() => format!("{u} ({p})"),
925 _ => u.to_string(),
926 });
927 let time = v
928 .get("timestamp")
929 .and_then(Value::as_str)
930 .and_then(rfc3339_ms)
931 .unwrap_or(received);
932 let context = md.get("context").map(|c| scrub(c, 0)).unwrap_or(json!({}));
933 let mut objects: Vec<String> = object.iter().cloned().collect();
934 if let Some(o) = &object {
936 if let Some((a, _)) = o.split_once('/') {
937 objects.push(a.to_string());
938 }
939 }
940 Some(NewRecord {
941 time,
942 source: "incus".into(),
943 org,
944 project: Some(project),
945 kind: action,
946 object_type,
947 object,
948 objects,
949 actor,
950 level: None,
951 message: None,
952 details: json!({
953 "source": clip(source.to_string(), 512),
954 "requestor": scrub(&req, 0),
955 "context": context,
956 "location": v.get("location"),
957 }),
958 })
959}
960
961pub fn rfc3339_ms(s: &str) -> Option<i64> {
963 let s = s.trim();
964 if s.len() < 20 {
965 return None;
966 }
967 let num = |a: usize, b: usize| s.get(a..b)?.parse::<i64>().ok();
968 let (y, mo, d) = (num(0, 4)?, num(5, 7)?, num(8, 10)?);
969 let (h, mi, se) = (num(11, 13)?, num(14, 16)?, num(17, 19)?);
970 let rest = &s[19..];
971 let (frac, tz) = match rest.strip_prefix('.') {
972 Some(r) => {
973 let end = r.find(|c: char| !c.is_ascii_digit()).unwrap_or(r.len());
974 (&r[..end], &r[end..])
975 }
976 None => ("", rest),
977 };
978 let ms = format!("{frac:0<3}").get(..3)?.parse::<i64>().ok()?;
979 let offset = match tz {
980 "Z" | "z" | "" => 0,
981 t if t.len() == 6 => {
982 let sign = if t.starts_with('-') { -1 } else { 1 };
983 sign * (t.get(1..3)?.parse::<i64>().ok()? * 60 + t.get(4..6)?.parse::<i64>().ok()?)
984 }
985 _ => return None,
986 };
987 let y2 = if mo <= 2 { y - 1 } else { y };
989 let era = y2.div_euclid(400);
990 let yoe = y2 - era * 400;
991 let mp = (mo + 9) % 12;
992 let doy = (153 * mp + 2) / 5 + d - 1;
993 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
994 let days = era * 146_097 + doe - 719_468;
995 Some(((days * 86400 + h * 3600 + mi * 60 + se) - offset * 60) * 1000 + ms)
996}
997
998pub fn fmt_ms(ms: i64) -> String {
1000 let secs = ms.div_euclid(1000);
1001 let (days, rem) = (secs.div_euclid(86400), secs.rem_euclid(86400));
1002 let z = days + 719_468;
1004 let era = z.div_euclid(146_097);
1005 let doe = z - era * 146_097;
1006 let yoe = (doe - doe / 1460 + doe / 36524 - doe / 146_096) / 365;
1007 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
1008 let mp = (5 * doy + 2) / 153;
1009 let d = doy - (153 * mp + 2) / 5 + 1;
1010 let m = if mp < 10 { mp + 3 } else { mp - 9 };
1011 let y = yoe + era * 400 + i64::from(m <= 2);
1012 format!(
1013 "{y:04}-{m:02}-{d:02} {:02}:{:02}:{:02}Z",
1014 rem / 3600,
1015 rem % 3600 / 60,
1016 rem % 60
1017 )
1018}
1019
1020pub fn marker(kind: &str, message: String, details: Value) -> NewRecord {
1022 NewRecord {
1023 source: "marker".into(),
1024 kind: kind.into(),
1025 actor: Some("isb".into()),
1026 level: Some("info".into()),
1027 message: Some(message),
1028 details,
1029 ..Default::default()
1030 }
1031}
1032
1033pub struct Recorder {
1039 tx: Mutex<Option<SyncSender<NewRecord>>>,
1040 dropped: Arc<AtomicU64>,
1041 worker: Mutex<Option<std::thread::JoinHandle<()>>>,
1042}
1043
1044impl Recorder {
1045 pub fn start(log: Arc<AuditLog>) -> Arc<Recorder> {
1046 let (tx, rx) = std::sync::mpsc::sync_channel::<NewRecord>(QUEUE);
1047 let dropped = Arc::new(AtomicU64::new(0));
1048 let d = dropped.clone();
1049 let worker = std::thread::Builder::new()
1050 .name("isb-history".into())
1051 .spawn(move || write_loop(&log, &rx, &d))
1052 .ok();
1053 Arc::new(Recorder {
1054 tx: Mutex::new(Some(tx)),
1055 dropped,
1056 worker: Mutex::new(worker),
1057 })
1058 }
1059
1060 pub fn record(&self, n: NewRecord) {
1062 let tx = self.tx.lock().unwrap_or_else(|e| e.into_inner());
1063 let sent = match tx.as_ref() {
1064 Some(t) => t.try_send(n).is_ok(),
1065 None => false,
1066 };
1067 if !sent {
1068 self.dropped.fetch_add(1, Ordering::Relaxed);
1069 }
1070 }
1071
1072 pub fn controller_sink(self: &Arc<Self>) -> crate::stack::controller::EventSink {
1074 let me = self.clone();
1075 Arc::new(move |e: &crate::stack::controller::Event| {
1076 if let Some(n) = from_event(e) {
1077 me.record(n);
1078 }
1079 })
1080 }
1081
1082 pub fn shutdown(&self) {
1084 self.tx.lock().unwrap_or_else(|e| e.into_inner()).take();
1085 if let Some(h) = self.worker.lock().unwrap_or_else(|e| e.into_inner()).take() {
1086 let _ = h.join();
1087 }
1088 }
1089}
1090
1091fn write_loop(log: &AuditLog, rx: &Receiver<NewRecord>, dropped: &AtomicU64) {
1092 loop {
1093 let r = rx.recv_timeout(Duration::from_secs(1));
1094 let lost = dropped.swap(0, Ordering::Relaxed);
1095 if lost > 0 {
1096 let m = marker(
1097 "history.dropped",
1098 format!("{lost} events were not recorded: the history queue was full"),
1099 json!({"count": lost}),
1100 );
1101 if let Err(e) = log.history_append(m) {
1102 eprintln!("isb serve: history: {e}");
1103 }
1104 }
1105 match r {
1106 Ok(n) => {
1107 if let Err(e) = log.history_append(n) {
1108 eprintln!("isb serve: history: {e}");
1109 }
1110 }
1111 Err(RecvTimeoutError::Timeout) => {}
1112 Err(RecvTimeoutError::Disconnected) => return,
1113 }
1114 }
1115}
1116
1117pub const REPEAT_WINDOW: Duration = Duration::from_secs(3600);
1119
1120const ROUTINE: &[&str] = &[
1127 "instance-exec",
1128 "instance-file-pushed",
1129 "instance-file-retrieved",
1130 "instance-log-retrieved",
1131 "instance-metrics-retrieved",
1132];
1133
1134pub struct Repeats {
1135 window: i64,
1136 seen: std::collections::HashMap<String, (i64, i64, u64, NewRecord)>,
1138}
1139
1140impl Repeats {
1141 pub fn new(window: Duration) -> Repeats {
1142 Repeats {
1143 window: window.as_millis() as i64,
1144 seen: Default::default(),
1145 }
1146 }
1147
1148 fn key(n: &NewRecord) -> Option<String> {
1149 if n.source != "incus" || !ROUTINE.contains(&n.kind.as_str()) {
1150 return None;
1151 }
1152 let program = n.details["context"]["command"]["program"]
1153 .as_str()
1154 .or_else(|| n.details["context"]["path"].as_str())
1155 .unwrap_or("");
1156 Some(format!(
1157 "{}|{}|{}|{}|{}",
1158 n.kind,
1159 n.project.as_deref().unwrap_or(""),
1160 n.object.as_deref().unwrap_or(""),
1161 program,
1162 n.actor.as_deref().unwrap_or("")
1163 ))
1164 }
1165
1166 pub fn admit(&mut self, n: NewRecord) -> Option<NewRecord> {
1168 let Some(k) = Self::key(&n) else {
1169 return Some(n);
1170 };
1171 match self.seen.get_mut(&k) {
1172 Some((first, last, count, _)) if n.time - *first < self.window => {
1173 *last = n.time;
1174 *count += 1;
1175 None
1176 }
1177 _ => {
1178 self.seen.insert(k, (n.time, n.time, 0, n.clone()));
1179 Some(n)
1180 }
1181 }
1182 }
1183
1184 pub fn flush(&mut self, now: i64, all: bool) -> Vec<NewRecord> {
1186 let window = self.window;
1187 let done: Vec<String> = self
1188 .seen
1189 .iter()
1190 .filter(|(_, (first, ..))| all || now - first >= window)
1191 .map(|(k, _)| k.clone())
1192 .collect();
1193 let mut out = Vec::new();
1194 for k in done {
1195 let Some((first, last, count, row)) = self.seen.remove(&k) else {
1196 continue;
1197 };
1198 if count == 0 {
1199 continue;
1200 }
1201 let mut n = row;
1202 n.time = last;
1203 n.message = Some(format!(
1204 "{count} more {} like this between {first} and {last} (folded)",
1205 n.kind
1206 ));
1207 if let Value::Object(m) = &mut n.details {
1208 m.insert(
1209 "repeats".into(),
1210 json!({"count": count, "from": first, "to": last}),
1211 );
1212 }
1213 out.push(n);
1214 }
1215 out
1216 }
1217}
1218
1219#[expect(
1223 clippy::excessive_nesting,
1224 reason = "predates the lint ratchet; split it when next changed"
1225)]
1226pub fn watch_incus(client: crate::client::Client, rec: Arc<Recorder>, stop: Arc<AtomicBool>) {
1227 let mut repeats = Repeats::new(REPEAT_WINDOW);
1228 let mut down_since: Option<i64> = None;
1229 let mut why = String::new();
1230 let mut attempt: u32 = 0;
1231 while !stop.load(Ordering::Relaxed) {
1232 match client.events_websocket("type=lifecycle&all-projects=true") {
1233 Ok(mut ws) => {
1234 if let Some(from) = down_since.take() {
1235 let to = crate::audit::now_ms();
1236 rec.record(marker(
1237 "incus.gap",
1238 format!(
1239 "incus events between {} and {} were not observed: {why}",
1240 fmt_ms(from),
1241 fmt_ms(to)
1242 ),
1243 json!({"from": from, "to": to, "reason": why}),
1244 ));
1245 }
1246 attempt = 0;
1247 loop {
1248 if stop.load(Ordering::Relaxed) {
1249 let _ = ws.close(None);
1250 for n in repeats.flush(crate::audit::now_ms(), true) {
1251 rec.record(n);
1252 }
1253 return;
1254 }
1255 match ws.read() {
1256 Ok(tungstenite::Message::Text(t)) => {
1257 if let Ok(v) = serde_json::from_str::<Value>(&t) {
1258 let now = crate::audit::now_ms();
1259 if let Some(n) = from_incus(&v, now) {
1260 if let Some(n) = repeats.admit(n) {
1261 rec.record(n);
1262 }
1263 }
1264 for n in repeats.flush(now, false) {
1265 rec.record(n);
1266 }
1267 }
1268 }
1269 Ok(_) => {}
1270 Err(tungstenite::Error::Io(e))
1271 if matches!(
1272 e.kind(),
1273 std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
1274 ) =>
1275 {
1276 for n in repeats.flush(crate::audit::now_ms(), false) {
1277 rec.record(n);
1278 }
1279 }
1280 Err(e) => {
1281 down_since = Some(crate::audit::now_ms());
1282 why = format!("the event stream broke: {e}");
1283 eprintln!("isb serve: history: incus {why}");
1284 break;
1285 }
1286 }
1287 }
1288 }
1289 Err(e) => {
1290 if down_since.is_none() {
1291 down_since = Some(crate::audit::now_ms());
1292 why = format!("cannot follow incus events: {e}");
1293 eprintln!("isb serve: history: {why}");
1294 }
1295 }
1296 }
1297 let wait = Duration::from_secs(1u64 << attempt.min(5));
1298 attempt += 1;
1299 let until = std::time::Instant::now() + wait;
1300 while std::time::Instant::now() < until && !stop.load(Ordering::Relaxed) {
1301 std::thread::sleep(Duration::from_millis(200));
1302 }
1303 }
1304}
1305
1306#[cfg(test)]
1307mod tests {
1308 use super::*;
1309
1310 fn incus_event(action: &str, source: &str, project: &str) -> Value {
1311 json!({
1312 "type": "lifecycle",
1313 "timestamp": "2026-10-03T09:00:00.250000000Z",
1314 "project": project,
1315 "location": "none",
1316 "metadata": {
1317 "action": action,
1318 "source": source,
1319 "context": {
1320 "type": "container",
1321 "config": {"environment.DB_PASSWORD": "hunter2", "limits.cpu": "2", "user.isb.create-token": "abc"},
1322 "api_token": "zzz",
1323 },
1324 "requestor": {"username": "stephan", "protocol": "unix", "address": "@"},
1325 },
1326 })
1327 }
1328
1329 #[test]
1330 fn incus_events_parse_and_scrub() {
1331 let n = from_incus(
1332 &incus_event(
1333 "instance-deleted",
1334 "/1.0/instances/web-1?project=isb-acme",
1335 "isb-acme",
1336 ),
1337 5,
1338 )
1339 .unwrap();
1340 assert_eq!(n.kind, "instance-deleted");
1341 assert_eq!(n.org.as_deref(), Some("acme"));
1342 assert_eq!(n.object_type.as_deref(), Some("instance"));
1343 assert_eq!(n.object.as_deref(), Some("web-1"));
1344 assert_eq!(n.actor.as_deref(), Some("stephan"));
1345 assert_eq!(n.time, rfc3339_ms("2026-10-03T09:00:00.25Z").unwrap());
1346 let d = n.details.to_string();
1347 assert!(
1348 !d.contains("hunter2") && !d.contains("abc") && !d.contains("zzz"),
1349 "{d}"
1350 );
1351 assert!(d.contains("limits.cpu"));
1352 let mut ex = incus_event(
1355 "instance-exec",
1356 "/1.0/instances/web-1?project=isb-acme",
1357 "isb-acme",
1358 );
1359 ex["metadata"]["context"] = json!({"command": ["/bin/sh", "-c", "export API_KEY=hunter3"]});
1360 let n = from_incus(&ex, 5).unwrap();
1361 let d = n.details.to_string();
1362 assert!(!d.contains("hunter3"), "{d}");
1363 assert_eq!(
1364 n.details["context"]["command"],
1365 json!({"program": "sh", "args": 2})
1366 );
1367 let n = from_incus(
1369 &incus_event(
1370 "image-alias-deleted",
1371 "/1.0/images/aliases/dev-base",
1372 "default",
1373 ),
1374 5,
1375 )
1376 .unwrap();
1377 assert_eq!((n.org, n.object.as_deref()), (None, Some("dev-base")));
1378 assert_eq!(n.object_type.as_deref(), Some("image-alias"));
1379 let n = from_incus(
1380 &incus_event("instance-created", "/1.0/instances/x", "default"),
1381 5,
1382 )
1383 .unwrap();
1384 assert_eq!(n.org, None);
1385 let n = from_incus(
1386 &incus_event(
1387 "instance-created",
1388 "/1.0/instances/x?project=isb-system",
1389 "isb-system",
1390 ),
1391 5,
1392 )
1393 .unwrap();
1394 assert_eq!(n.org, None);
1395 assert!(from_incus(&json!({"type": "logging"}), 5).is_none());
1396 let (t, o, p) =
1397 parse_source("/1.0/storage-pools/default/volumes/custom/data%20x?project=isb-a");
1398 assert_eq!(
1399 (t.as_deref(), o.as_deref(), p.as_deref()),
1400 (Some("storage-volume"), Some("data x"), Some("isb-a"))
1401 );
1402 }
1403
1404 #[test]
1405 fn routine_repeats_are_folded() {
1406 let ev = |t: i64, program: &str| {
1407 let mut e = incus_event(
1408 "instance-exec",
1409 "/1.0/instances/web-1?project=isb-acme",
1410 "isb-acme",
1411 );
1412 e["metadata"]["context"] = json!({"command": [program, "is-active", "x"]});
1413 let mut n = from_incus(&e, 0).unwrap();
1414 n.time = t;
1415 n
1416 };
1417 let mut r = Repeats::new(Duration::from_secs(60));
1418 assert!(r.admit(ev(0, "systemctl")).is_some());
1419 assert!(r.admit(ev(5_000, "systemctl")).is_none());
1420 assert!(r.admit(ev(10_000, "systemctl")).is_none());
1421 assert!(r.admit(ev(11_000, "sh")).is_some());
1423 let del = from_incus(
1425 &incus_event(
1426 "instance-deleted",
1427 "/1.0/instances/web-1?project=isb-acme",
1428 "isb-acme",
1429 ),
1430 1,
1431 )
1432 .unwrap();
1433 assert!(r.admit(del.clone()).is_some() && r.admit(del).is_some());
1434 assert!(r.flush(30_000, false).is_empty());
1435 let s = r.flush(61_000, false);
1436 assert_eq!(s.len(), 1);
1437 assert_eq!(s[0].details["repeats"]["count"], 2);
1438 assert!(r.admit(ev(70_000, "systemctl")).is_some());
1439 }
1440
1441 #[test]
1442 fn rfc3339() {
1443 assert_eq!(fmt_ms(951_868_800_500), "2000-03-01 00:00:00Z");
1444 assert_eq!(
1445 fmt_ms(rfc3339_ms("2026-10-03T09:45:24Z").unwrap()),
1446 "2026-10-03 09:45:24Z"
1447 );
1448 assert_eq!(rfc3339_ms("1970-01-01T00:00:00Z"), Some(0));
1449 assert_eq!(rfc3339_ms("1970-01-01T00:00:01.5Z"), Some(1500));
1450 assert_eq!(rfc3339_ms("2000-03-01T00:00:00Z"), Some(951_868_800_000));
1451 assert_eq!(
1452 rfc3339_ms("2000-03-01T02:00:00+02:00"),
1453 Some(951_868_800_000)
1454 );
1455 assert_eq!(rfc3339_ms("nope"), None);
1456 }
1457
1458 fn rec(org: Option<&str>, kind: &str, object: &str, time: i64) -> NewRecord {
1459 NewRecord {
1460 time,
1461 source: "incus".into(),
1462 org: org.map(String::from),
1463 kind: kind.into(),
1464 object_type: Some("instance".into()),
1465 object: Some(object.into()),
1466 actor: Some("stephan".into()),
1467 ..Default::default()
1468 }
1469 }
1470
1471 #[test]
1472 fn history_chains_prunes_and_merges_with_audit() {
1473 let log = AuditLog::in_memory().unwrap().with_clock(|| 10_000_000);
1474 log.history_append(rec(Some("acme"), "instance-created", "web-1", 1_000))
1475 .unwrap();
1476 log.history_append(rec(None, "image-alias-deleted", "dev-base", 2_000))
1477 .unwrap();
1478 log.history_append(rec(Some("beta"), "instance-deleted", "api-1", 3_000))
1479 .unwrap();
1480 let a = log
1481 .append(crate::audit::NewEntry {
1482 org: Some("acme".into()),
1483 actor: crate::audit::Actor::local(Some(1000)),
1484 action: "stack_remove".into(),
1485 target: Some("web".into()),
1486 outcome: "ok".into(),
1487 ..Default::default()
1488 })
1489 .unwrap();
1490 log.history_append(rec(Some("acme"), "instance-deleted", "web-1", a.time + 500))
1491 .unwrap();
1492 assert!(log.history_verify().unwrap().ok);
1493 let acme = Visibility::Orgs(vec!["acme".into()]);
1495 let q = HistoryQuery {
1496 correlate: true,
1497 ..Default::default()
1498 };
1499 let p = log.timeline(&q, &acme, Some(&acme)).unwrap();
1500 let kinds: Vec<&str> = p.items.iter().map(|i| i.kind.as_str()).collect();
1501 assert_eq!(
1502 kinds,
1503 ["instance-deleted", "stack_remove", "instance-created"]
1504 );
1505 let inf = p.items[0].inferred.as_ref().unwrap();
1506 assert_eq!(inf["action"], "stack_remove");
1507 let p = log.timeline(&q, &acme, None).unwrap();
1509 assert_eq!(p.items.len(), 2);
1510 assert!(p.items[0].inferred.is_none());
1511 let q = HistoryQuery {
1513 object: Some("dev-base".into()),
1514 exact: true,
1515 ..Default::default()
1516 };
1517 let p = log
1518 .timeline(&q, &Visibility::All, Some(&Visibility::All))
1519 .unwrap();
1520 assert_eq!(p.items.len(), 1);
1521 assert_eq!(p.items[0].actor.as_deref(), Some("stephan"));
1522 assert!(log.timeline(&q, &acme, None).unwrap().items.is_empty());
1523 let mut seen = Vec::new();
1525 let mut q = HistoryQuery {
1526 limit: Some(1),
1527 ..Default::default()
1528 };
1529 loop {
1530 let p = log
1531 .timeline(&q, &Visibility::All, Some(&Visibility::All))
1532 .unwrap();
1533 seen.extend(p.items.iter().map(|i| (i.source.clone(), i.id)));
1534 match p.next {
1535 Some(n) => q.before = Some(n),
1536 None => break,
1537 }
1538 }
1539 assert_eq!(seen.len(), 5, "{seen:?}");
1540 log.db()
1542 .execute_batch(
1543 "DROP TRIGGER history_no_update; UPDATE history SET actor = 'x' WHERE id = 2;",
1544 )
1545 .unwrap();
1546 assert_eq!(log.history_verify().unwrap().broken.unwrap().0, 2);
1547 }
1548
1549 #[test]
1550 fn history_prunes_by_age_and_size() {
1551 let now = Arc::new(std::sync::atomic::AtomicI64::new(400 * 86_400_000));
1552 let n = now.clone();
1553 let log = AuditLog::in_memory()
1554 .unwrap()
1555 .with_clock(move || n.load(Ordering::SeqCst))
1556 .with_history_limits(DEFAULT_RETENTION, 1000);
1557 log.history_append(rec(None, "old", "x", 1)).unwrap();
1558 for i in 0..1100 {
1559 log.history_append(rec(None, "k", &format!("o{i}"), 0))
1560 .unwrap();
1561 }
1562 log.prune().unwrap();
1563 let v = log.history_verify().unwrap();
1564 assert!(v.ok, "{v:?}");
1565 assert_eq!(v.rows, 1000);
1566 }
1567
1568 #[test]
1569 fn recorder_never_blocks_and_counts_drops() {
1570 let log = Arc::new(AuditLog::in_memory().unwrap());
1571 let r = Recorder::start(log.clone());
1572 for i in 0..20 {
1573 r.record(rec(None, "k", &format!("o{i}"), 0));
1574 }
1575 r.shutdown();
1576 r.record(rec(None, "late", "x", 0));
1578 assert_eq!(r.dropped.load(Ordering::Relaxed), 1);
1579 assert_eq!(log.history_head().unwrap(), 20);
1580 }
1581}