Skip to main content

isb_server/
history.rs

1//! The history: everything that happened to what isb runs, kept so the
2//! current state can always be traced back.
3//!
4//! Three sources, one timeline:
5//! - **controller**: every [`crate::stack::controller::Event`] the stack
6//!   controller emits (deploys, rollouts, health, restarts, backups, jobs),
7//!   persisted as it is emitted, through a bounded queue that never blocks
8//!   the controller (a full queue is counted and recorded as a
9//!   `history.dropped` marker);
10//! - **incus**: every lifecycle event incus emits, in every project,
11//!   including changes made outside isb (`incus delete`, `incus image alias
12//!   delete`), with its `requestor`;
13//! - **audit**: the audit rows ([`crate::audit`]): tool calls, sign-ins,
14//!   account changes.
15//!
16//! Markers make what was not seen explicit: `serve.started`,
17//! `serve.stopped`, and `incus.gap` (events between two times not observed:
18//! the daemon was down, or the event stream dropped).
19//!
20//! Rows live in `<state>/audit.db`'s `history` table, append only and hash
21//! chained like the audit rows (on their own chain). Contexts are scrubbed of
22//! anything that could be a secret before they are stored.
23
24use 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
39/// The default retention: 365 days.
40pub const DEFAULT_RETENTION: Duration = Duration::from_secs(365 * 86400);
41/// The default bound on rows: past it, the oldest go first.
42pub const DEFAULT_MAX_ROWS: i64 = 5_000_000;
43/// Queued records before new ones are dropped (and counted).
44const QUEUE: usize = 10_000;
45
46/// A row to append.
47#[derive(Debug, Clone, Default)]
48pub struct NewRecord {
49    /// Unix milliseconds when it happened; 0 is now.
50    pub time: i64,
51    /// `controller`, `incus` or `marker`.
52    pub source: String,
53    /// `None`: host level (platform admins).
54    pub org: Option<String>,
55    /// The incus project, when there is one.
56    pub project: Option<String>,
57    /// The event's kind (`deploy.succeeded`, `instance-deleted`,
58    /// `serve.started`, ...).
59    pub kind: String,
60    /// `stack`, `instance`, `image`, `image-alias`, `storage-volume`, ...
61    pub object_type: Option<String>,
62    pub object: Option<String>,
63    /// Every name the row is about (stack, service, instance), for search.
64    pub objects: Vec<String>,
65    /// Who: an incus requestor's username, or `isb` for the controller.
66    pub actor: Option<String>,
67    pub level: Option<String>,
68    pub message: Option<String>,
69    pub details: Value,
70}
71
72/// A stored row.
73#[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/// A history query. Globs are shell-style.
131#[derive(Debug, Clone, Default, Deserialize)]
132#[serde(default, deny_unknown_fields)]
133pub struct HistoryQuery {
134    /// One org's rows; with `platform`, only host-level rows.
135    pub org: Option<String>,
136    pub platform: bool,
137    /// A name the row is about: an instance, image, volume, stack, app,
138    /// service. Substring, or the whole name with `exact`. Markers (gaps,
139    /// restarts) are kept, so a timeline shows when nothing was watching.
140    pub object: Option<String>,
141    pub exact: bool,
142    /// Glob on the kind (`instance-*`, `deploy.*`) or the audit action.
143    pub kind: Option<String>,
144    /// `audit`, `controller`, `incus`, `marker`; comma-separated for more
145    /// than one. All by default.
146    pub source: Option<String>,
147    /// Glob on who: an incus requestor, an audit actor.
148    pub actor: Option<String>,
149    /// Unix milliseconds, inclusive.
150    pub since: Option<i64>,
151    /// Unix milliseconds, exclusive.
152    pub until: Option<i64>,
153    /// Page cursor: the `next` of the previous page.
154    pub before: Option<String>,
155    /// Default 100, at most 1000.
156    pub limit: Option<usize>,
157    /// Oldest first (a timeline) rather than newest first.
158    pub ascending: bool,
159    /// Link incus instance events to the audit row that likely caused them.
160    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/// One row of the merged timeline: a history row or an audit row.
174#[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    /// `info`/`warn`/`error` for events; the outcome for audit rows.
185    pub level: Option<String>,
186    pub message: Option<String>,
187    pub details: Value,
188    /// For an incus event isb likely caused: which audit row, and why it is
189    /// thought so. Best effort, by time and name.
190    #[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    /// Ordering key: time, then source, then id.
246    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/// A page of the merged timeline and the cursor for the next one.
271#[derive(Debug, Clone, Serialize)]
272pub struct Page {
273    pub items: Vec<Item>,
274    pub next: Option<String>,
275}
276
277impl AuditLog {
278    /// Append one history row, chained to the newest.
279    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    /// The newest history row's time, if any.
362    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    /// Drop rows past the retention and past the row bound, oldest first,
371    /// keeping the chain anchored.
372    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    /// History rows matching `q` that `vis` may see, newest first (or
418    /// oldest first with `ascending`), `time <= upto` when given.
419    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    /// Rows appended after id `after`, oldest first (tailing).
517    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    /// The newest history id, 0 when empty.
534    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    /// Walk the history chain, as [`AuditLog::verify`] does the audit one.
543    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    /// The merged timeline: history rows `hvis` may see and, when `avis` is
602    /// given, audit rows it may see.
603    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        // Fetch a little past the page from each source, then merge.
612        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    /// For incus instance events, the audit row (same org, up to two
668    /// minutes before) whose target names the instance's stack or app.
669    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                // An audit row is written when its call returns, so a
682                // long call (a deploy that waits) lands after what it did.
683                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            // The nearest call before the event, else the nearest after it.
708            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
736// ---- sources ----
737
738/// A controller event as a history row. `log`-level lines (a deployment's
739/// build output) are not kept here: each deployment has its own log.
740pub 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    // Volume snapshots and staged restores are about a volume, not a stack.
754    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
774/// Keys whose values are never stored from an incus context.
775fn 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
791/// A context with anything that could be a secret removed, strings cut
792/// short, and at most 64 keys per level.
793pub 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                    // A command line can carry anything: keep the program
801                    // and how many arguments, never the arguments.
802                    ("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
825/// `/1.0/instances/web?project=isb-acme` → (`instance`, `web`, `isb-acme`).
826pub 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
882/// The org an incus project belongs to: `isb-<org>` only (the default
883/// project and `isb-system` hold host-level things).
884pub 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
894/// An incus lifecycle event (one message of `/1.0/events`) as a history
895/// row, or `None` for anything else.
896pub 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    // Images, pools and networks are the host's whatever project they sit in.
913    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    // An instance's name inside the path of a snapshot or volume.
935    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
961/// `2026-10-03T09:00:00.123456789Z` (or with an offset) as unix ms.
962pub 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    // Days from civil (Howard Hinnant).
988    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
998/// Unix ms as `2026-10-03 09:45:24Z`.
999pub 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    // Civil from days (Howard Hinnant).
1003    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
1020/// A marker row: what the daemon itself saw (started, stopped, a gap).
1021pub 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
1033// ---- the writer ----
1034
1035/// Appends records on its own thread, so emitting never blocks. A full
1036/// queue drops the record, counts it, and the count is recorded as a
1037/// `history.dropped` marker once there is room.
1038pub 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    /// Queue a record; never blocks.
1061    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    /// What the controller calls for every event.
1073    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    /// Write what is queued and stop.
1083    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
1117/// How long repeats of a routine action are folded into one row.
1118pub const REPEAT_WINDOW: Duration = Duration::from_secs(3600);
1119
1120/// Routine actions that repeat all day (isb's own probes exec `systemctl
1121/// is-active` in every replica every few seconds): the first of a kind per
1122/// instance, program and requestor in a window is recorded as it happens,
1123/// the rest are counted and recorded as one summary row when the window
1124/// ends. Anything else (created, deleted, started, updated, ...) is always
1125/// recorded one by one.
1126const 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    /// key → (first time, last time, how many after the first, the first row)
1137    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    /// The row to record now, or `None` when it is a repeat being counted.
1167    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    /// Summary rows for windows that ended (all of them with `all`).
1185    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/// Follow incus' lifecycle events, in every project, until `stop`: each
1220/// one recorded, reconnecting with backoff, and every stretch without a
1221/// connection recorded as an `incus.gap`.
1222#[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        // An exec keeps the program and the argument count, never the
1353        // arguments.
1354        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        // Images and other projects are host level.
1368        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        // Another program is its own row.
1422        assert!(r.admit(ev(11_000, "sh")).is_some());
1423        // Not routine: always recorded.
1424        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        // An org's members see that org; the platform sees the rest.
1494        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        // Without audit visibility: no audit rows, nothing inferred.
1508        let p = log.timeline(&q, &acme, None).unwrap();
1509        assert_eq!(p.items.len(), 2);
1510        assert!(p.items[0].inferred.is_none());
1511        // The dev-base question.
1512        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        // Paging: one at a time, newest first, nothing twice.
1524        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        // Edits are detected.
1541        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        // After shutdown, records are dropped (and counted), never block.
1577        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}