Skip to main content

cairn_mod/audit/
append.rs

1//! Audit-log append helpers (#39, v1.3; chain extended to span
2//! `pds_admin_audit` in #85, v1.7).
3//!
4//! Two public entry points + one shared connection-level helper. All
5//! three route through [`crate::audit::hash::compute_audit_row_hash`]
6//! so the chain has a single canonical hash implementation.
7//!
8//! ```text
9//!   in-process call sites with their own tx        →  append_in_tx (writer-internal,
10//!                                                                   flag_reporter)
11//!   in-process callers without an existing tx      →  WriterHandle::append_audit
12//!                                                     (retention_sweep)
13//!   cross-process CLI callers (no writer task)     →  append_via_pool (publish,
14//!                                                                      unpublish)
15//! ```
16//!
17//! `prev_hash` is the **most recent chain link across both tables** —
18//! `audit_log.row_hash` (latest non-NULL by id) and
19//! `pds_admin_audit.row_hash` (latest by id). Tie-broken on insertion
20//! timestamp (`audit_log.created_at` vs.
21//! `pds_admin_audit.call_completed_at`); falls back to
22//! [`crate::audit::hash::GENESIS_PREV_HASH`] when both tables are empty
23//! or audit_log only contains pre-v1.3 rows. The unified-chain semantics
24//! live in the crate-internal `read_latest_chain_hash` helper below.
25//!
26//! The chain-integrity invariant is enforced at the SQLite write-lock
27//! layer. [`append_via_pool`] explicitly issues `BEGIN IMMEDIATE` so the
28//! latest-row read + INSERT pair is atomic against any concurrent
29//! appender — including a `cairn serve` writer task running in another
30//! process. [`append_in_tx`] inherits the lock from its caller's already-
31//! open transaction.
32
33use sqlx::sqlite::SqliteConnection;
34use sqlx::{Pool, Sqlite, Transaction};
35
36use super::hash::{
37    AuditRowForHashing, GENESIS_PREV_HASH, compute_audit_row_hash, parse_stored_hash,
38};
39use crate::error::{Error, Result};
40
41/// Owned audit-row payload used as input to the append helpers.
42/// Mirrors [`crate::audit::hash::AuditRowForHashing`] but holds owned
43/// data so it can travel across futures + writer-task message boundaries.
44#[derive(Debug, Clone)]
45pub struct AuditRowForAppend {
46    /// Internal wall-clock epoch-ms.
47    pub created_at: i64,
48    /// Audit-action discriminator.
49    pub action: String,
50    /// Actor DID.
51    pub actor_did: String,
52    /// Optional target identifier.
53    pub target: Option<String>,
54    /// Optional CID pin on the target.
55    pub target_cid: Option<String>,
56    /// `"success"` or `"failure"`.
57    pub outcome: String,
58    /// Optional structured-JSON or free-text payload.
59    pub reason: Option<String>,
60}
61
62impl AuditRowForAppend {
63    /// Borrowed view for hashing — no ownership transfer.
64    fn as_hashing(&self) -> AuditRowForHashing<'_> {
65        AuditRowForHashing {
66            created_at: self.created_at,
67            action: &self.action,
68            actor_did: &self.actor_did,
69            target: self.target.as_deref(),
70            target_cid: self.target_cid.as_deref(),
71            outcome: &self.outcome,
72            reason: self.reason.as_deref(),
73        }
74    }
75}
76
77/// Append an audit row inside an existing transaction. Used by callers
78/// that need to combine the audit row with other writes in the same
79/// commit — writer-internal handlers (apply / negate / resolve_report)
80/// and `flag_reporter` (suppression update + audit row atomically).
81///
82/// The transaction itself is the serialization point: SQLite's
83/// `BEGIN IMMEDIATE` (or any write-acquiring tx) holds the database
84/// write lock, so the latest-row-hash read + INSERT pair is atomic
85/// against any concurrent appender.
86///
87/// Returns the inserted `audit_log.id`.
88pub async fn append_in_tx(
89    tx: &mut Transaction<'_, Sqlite>,
90    row: &AuditRowForAppend,
91) -> Result<i64> {
92    perform_append(tx, row).await
93}
94
95/// Append an audit row in its own transaction. Used by cross-process
96/// CLI callers that don't have a writer task in their address space —
97/// `cairn publish-service-record` and `cairn unpublish-service-record`
98/// run as one-shot CLIs and write to the same SQLite file as a running
99/// `cairn serve`, so the SQLite write lock is what serializes them
100/// against the writer's appends.
101///
102/// `BEGIN IMMEDIATE` is the load-bearing primitive: it acquires the
103/// write lock at transaction-start instead of first-write-statement,
104/// closing the read-then-modify race that would otherwise let two
105/// processes both observe the same `prev_hash` and write rows that
106/// fork the chain. sqlx 0.8's `Pool::begin` issues `BEGIN DEFERRED`
107/// with no override hook, so we reach for `Pool::acquire` + raw SQL.
108pub async fn append_via_pool(pool: &Pool<Sqlite>, row: &AuditRowForAppend) -> Result<i64> {
109    let mut conn = pool
110        .acquire()
111        .await
112        .map_err(|e| Error::Signing(format!("audit acquire: {e}")))?;
113    sqlx::query("BEGIN IMMEDIATE")
114        .execute(&mut *conn)
115        .await
116        .map_err(|e| Error::Signing(format!("audit begin: {e}")))?;
117
118    match perform_append(&mut conn, row).await {
119        Ok(id) => {
120            sqlx::query("COMMIT")
121                .execute(&mut *conn)
122                .await
123                .map_err(|e| Error::Signing(format!("audit commit: {e}")))?;
124            Ok(id)
125        }
126        Err(e) => {
127            // Best-effort rollback — if the rollback itself fails the
128            // connection is in a degraded state but the original error
129            // is what the caller cares about.
130            let _ = sqlx::query("ROLLBACK").execute(&mut *conn).await;
131            Err(e)
132        }
133    }
134}
135
136/// Connection-level append: reads `prev_hash`, computes `row_hash`,
137/// INSERTs the row. The caller is responsible for transactional
138/// framing (either inheriting one via [`append_in_tx`] or issuing
139/// `BEGIN IMMEDIATE` via [`append_via_pool`]).
140async fn perform_append(conn: &mut SqliteConnection, row: &AuditRowForAppend) -> Result<i64> {
141    let prev_hash = read_latest_chain_hash(&mut *conn).await?;
142    let row_hash = compute_audit_row_hash(&prev_hash, &row.as_hashing())?;
143
144    let prev_hash_slice: &[u8] = &prev_hash;
145    let row_hash_slice: &[u8] = &row_hash;
146    let id = sqlx::query_scalar!(
147        "INSERT INTO audit_log
148             (created_at, action, actor_did, target, target_cid, outcome, reason,
149              prev_hash, row_hash)
150         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
151         RETURNING id",
152        row.created_at,
153        row.action,
154        row.actor_did,
155        row.target,
156        row.target_cid,
157        row.outcome,
158        row.reason,
159        prev_hash_slice,
160        row_hash_slice,
161    )
162    .fetch_one(&mut *conn)
163    .await
164    .map_err(|e| Error::Signing(format!("audit append: {e}")))?;
165    Ok(id)
166}
167
168/// Read the latest stored `row_hash` across the unified chain.
169///
170/// Walks four tables since #94 / v1.7: `audit_log` (#39),
171/// `pds_admin_audit` (#85 / §F23), `xrpc_known_callers` and
172/// `xrpc_trusted_pdses` (#94 / §F23 inbound). Skips pre-v1.3
173/// `audit_log` rows that have NULL hashes. Returns
174/// [`GENESIS_PREV_HASH`] when all four tables are empty (or
175/// `audit_log` only contains pre-v1.3 NULL rows and the other
176/// three are empty).
177///
178/// # Tie-break
179///
180/// When multiple candidate rows share the latest timestamp, the
181/// chain tip is the one with the highest **table priority**:
182///
183/// | Table | Priority |
184/// |-------|----------|
185/// | `audit_log` | 0 |
186/// | `pds_admin_audit` | 1 |
187/// | `xrpc_known_callers` | 2 |
188/// | `xrpc_trusted_pdses` | 3 |
189///
190/// Higher priority = "newer in the chain" on ties. Mirrored by
191/// `cairn audit verify`'s walker (#88, extended in #94). The
192/// pairing keeps the read and the walker consistent so the
193/// chain has a single canonical ordering.
194///
195/// **Monotonicity caveat:** if a low-priority row inserts AFTER
196/// a high-priority row at the same millisecond timestamp, the
197/// priority-based tie-break treats the high-priority row as the
198/// chain tip (not the actually-most-recently-inserted low-
199/// priority row), and the next chained row would skip the low-
200/// priority one. The `xrpc_*` write paths in #94 enforce
201/// strict-monotonic timestamps to sidestep this; the existing
202/// `audit_log` + `pds_admin_audit` paths rely on the writer task's
203/// natural sub-millisecond gap between writes (ms-resolution
204/// timestamps tied across these two tables are rare in practice).
205///
206/// `pub(crate)` so the `pds_admin_audit` (#85) and `xrpc_*`
207/// (#94) append helpers can share the same chain-tip read.
208pub(crate) async fn read_latest_chain_hash(conn: &mut SqliteConnection) -> Result<[u8; 32]> {
209    let candidates = read_chain_tip_candidates(&mut *conn).await?;
210    Ok(select_chain_tip(&candidates).unwrap_or(GENESIS_PREV_HASH))
211}
212
213/// Borrowed candidate row across the four chain tables. Each
214/// variant carries the row's chain-ordering timestamp and the
215/// stored `row_hash` (or `None` for pre-v1.3 audit_log rows that
216/// haven't been backfilled by `cairn audit-rebuild`).
217#[derive(Debug, Clone)]
218struct ChainTipCandidate {
219    /// Chain-ordering timestamp (epoch-ms).
220    timestamp: i64,
221    /// Stable per-table priority (see `read_latest_chain_hash`).
222    priority: u8,
223    /// Stored row_hash. `None` only for pre-v1.3 audit_log rows.
224    row_hash: Option<Vec<u8>>,
225}
226
227async fn read_chain_tip_candidates(conn: &mut SqliteConnection) -> Result<Vec<ChainTipCandidate>> {
228    let mut out = Vec::with_capacity(4);
229
230    let audit_log = sqlx::query!(
231        "SELECT row_hash, created_at FROM audit_log
232         WHERE row_hash IS NOT NULL
233         ORDER BY id DESC LIMIT 1"
234    )
235    .fetch_optional(&mut *conn)
236    .await
237    .map_err(|e| Error::Signing(format!("audit_log prev_hash read: {e}")))?;
238    if let Some(r) = audit_log {
239        out.push(ChainTipCandidate {
240            timestamp: r.created_at,
241            priority: 0,
242            row_hash: r.row_hash,
243        });
244    }
245
246    let pds_admin = sqlx::query!(
247        "SELECT row_hash, call_completed_at FROM pds_admin_audit
248         ORDER BY id DESC LIMIT 1"
249    )
250    .fetch_optional(&mut *conn)
251    .await
252    .map_err(|e| Error::Signing(format!("pds_admin_audit prev_hash read: {e}")))?;
253    if let Some(r) = pds_admin {
254        out.push(ChainTipCandidate {
255            timestamp: r.call_completed_at,
256            priority: 1,
257            row_hash: Some(r.row_hash),
258        });
259    }
260
261    let known_callers = sqlx::query!(
262        "SELECT row_hash, added_at FROM xrpc_known_callers
263         ORDER BY added_at DESC LIMIT 1"
264    )
265    .fetch_optional(&mut *conn)
266    .await
267    .map_err(|e| Error::Signing(format!("xrpc_known_callers prev_hash read: {e}")))?;
268    if let Some(r) = known_callers {
269        out.push(ChainTipCandidate {
270            timestamp: r.added_at,
271            priority: 2,
272            row_hash: Some(r.row_hash),
273        });
274    }
275
276    let trusted_pdses = sqlx::query!(
277        "SELECT row_hash, added_at FROM xrpc_trusted_pdses
278         ORDER BY added_at DESC LIMIT 1"
279    )
280    .fetch_optional(&mut *conn)
281    .await
282    .map_err(|e| Error::Signing(format!("xrpc_trusted_pdses prev_hash read: {e}")))?;
283    if let Some(r) = trusted_pdses {
284        out.push(ChainTipCandidate {
285            timestamp: r.added_at,
286            priority: 3,
287            row_hash: Some(r.row_hash),
288        });
289    }
290
291    Ok(out)
292}
293
294/// Read the maximum chain-ordering timestamp across all four
295/// chain tables, in epoch-ms. Returns `None` when all four are
296/// empty.
297///
298/// Used by the #94 `xrpc_*` write paths to enforce strict-
299/// monotonic chain timestamps:
300/// `added_at = max(epoch_ms_now(), max_existing + 1)`. This
301/// sidesteps the priority-tie-break monotonicity caveat
302/// documented on [`read_latest_chain_hash`] for chain extensions
303/// that originate outside the writer task (CLI inserts).
304pub(crate) async fn read_latest_chain_timestamp_ms(
305    conn: &mut SqliteConnection,
306) -> Result<Option<i64>> {
307    let candidates = read_chain_tip_candidates(&mut *conn).await?;
308    Ok(candidates.iter().map(|c| c.timestamp).max())
309}
310
311/// Pick the chain tip from a set of candidates, applying the
312/// `(timestamp DESC, priority DESC)` rule. Returns `None` if no
313/// candidate has a non-NULL `row_hash` (caller defaults to
314/// `GENESIS_PREV_HASH`).
315fn select_chain_tip(candidates: &[ChainTipCandidate]) -> Option<[u8; 32]> {
316    let mut best: Option<&ChainTipCandidate> = None;
317    for c in candidates {
318        if c.row_hash.is_none() {
319            continue;
320        }
321        match best {
322            None => best = Some(c),
323            Some(b) => {
324                if c.timestamp > b.timestamp
325                    || (c.timestamp == b.timestamp && c.priority > b.priority)
326                {
327                    best = Some(c);
328                }
329            }
330        }
331    }
332    let bytes = best?.row_hash.as_ref()?;
333    parse_stored_hash(bytes).ok()
334}
335
336#[cfg(test)]
337mod tests {
338    use super::*;
339    use crate::storage;
340    use tempfile::tempdir;
341
342    async fn fresh_pool() -> Pool<Sqlite> {
343        let dir = tempdir().unwrap();
344        let path = dir.path().join("audit-test.db");
345        let pool = storage::open(&path).await.unwrap();
346        // Hold the tempdir for the pool's lifetime via Box::leak —
347        // the pool's connections reference the on-disk file. Test-only
348        // pattern; production paths own their TempDir explicitly.
349        Box::leak(Box::new(dir));
350        pool
351    }
352
353    fn sample_row(action: &str, actor_did: &str) -> AuditRowForAppend {
354        AuditRowForAppend {
355            created_at: 1_776_902_400_000,
356            action: action.into(),
357            actor_did: actor_did.into(),
358            target: Some("at://did:plc:target/col/r".into()),
359            target_cid: None,
360            outcome: "success".into(),
361            reason: None,
362        }
363    }
364
365    #[tokio::test]
366    async fn first_append_uses_genesis_prev_hash() {
367        let pool = fresh_pool().await;
368        let id = append_via_pool(&pool, &sample_row("label_applied", "did:plc:m1"))
369            .await
370            .unwrap();
371
372        let stored = sqlx::query!(
373            "SELECT prev_hash, row_hash FROM audit_log WHERE id = ?1",
374            id
375        )
376        .fetch_one(&pool)
377        .await
378        .unwrap();
379
380        let prev: &[u8] = stored.prev_hash.as_deref().expect("prev_hash present");
381        assert_eq!(prev, GENESIS_PREV_HASH);
382        assert!(stored.row_hash.is_some());
383    }
384
385    #[tokio::test]
386    async fn second_append_chains_to_first_row_hash() {
387        let pool = fresh_pool().await;
388        let id1 = append_via_pool(&pool, &sample_row("label_applied", "did:plc:m1"))
389            .await
390            .unwrap();
391        let id2 = append_via_pool(&pool, &sample_row("label_negated", "did:plc:m1"))
392            .await
393            .unwrap();
394
395        let row1_hash: Vec<u8> = sqlx::query_scalar!(
396            r#"SELECT row_hash AS "row_hash!" FROM audit_log WHERE id = ?1"#,
397            id1
398        )
399        .fetch_one(&pool)
400        .await
401        .unwrap();
402        let row2_prev: Vec<u8> = sqlx::query_scalar!(
403            r#"SELECT prev_hash AS "prev_hash!" FROM audit_log WHERE id = ?1"#,
404            id2
405        )
406        .fetch_one(&pool)
407        .await
408        .unwrap();
409        assert_eq!(
410            row1_hash, row2_prev,
411            "row 2's prev_hash must match row 1's row_hash"
412        );
413    }
414
415    #[tokio::test]
416    async fn pre_v13_null_rows_are_skipped_for_prev_hash_lookup() {
417        let pool = fresh_pool().await;
418        // Manually insert a pre-v1.3 row (NULL hashes), then append
419        // a v1.3 row and confirm the v1.3 row's prev_hash is GENESIS,
420        // not the NULL value of the previous row.
421        sqlx::query!(
422            "INSERT INTO audit_log (created_at, action, actor_did, outcome)
423             VALUES (?1, ?2, ?3, ?4)",
424            1_i64,
425            "label_applied",
426            "did:plc:m1",
427            "success"
428        )
429        .execute(&pool)
430        .await
431        .unwrap();
432
433        let v13_id = append_via_pool(&pool, &sample_row("label_negated", "did:plc:m1"))
434            .await
435            .unwrap();
436
437        let prev: Option<Vec<u8>> =
438            sqlx::query_scalar!("SELECT prev_hash FROM audit_log WHERE id = ?1", v13_id)
439                .fetch_one(&pool)
440                .await
441                .unwrap();
442        let prev_bytes = prev.expect("prev_hash present on v1.3 row");
443        assert_eq!(
444            prev_bytes.as_slice(),
445            GENESIS_PREV_HASH,
446            "v1.3 row must use GENESIS_PREV_HASH when latest row has NULL row_hash"
447        );
448    }
449
450    #[tokio::test]
451    async fn append_in_tx_and_append_via_pool_produce_same_row_hash() {
452        // Pin the no-drift contract from #39's design: both append
453        // paths must compute the same row_hash for the same row
454        // content + prev_hash. If this test ever fails, the hash
455        // implementation has split between the two paths.
456        let pool_a = fresh_pool().await;
457        let pool_b = fresh_pool().await;
458        let row = sample_row("label_applied", "did:plc:m1");
459
460        // Path 1: append_via_pool on its own.
461        let id_a = append_via_pool(&pool_a, &row).await.unwrap();
462        let hash_a: Vec<u8> = sqlx::query_scalar!(
463            r#"SELECT row_hash AS "row_hash!" FROM audit_log WHERE id = ?1"#,
464            id_a
465        )
466        .fetch_one(&pool_a)
467        .await
468        .unwrap();
469
470        // Path 2: append_in_tx via an explicit BEGIN IMMEDIATE on a
471        // raw connection (matches the production tx pattern that
472        // writer-internal callers and flag_reporter exercise).
473        let mut tx = pool_b.begin().await.unwrap();
474        let id_b = append_in_tx(&mut tx, &row).await.unwrap();
475        tx.commit().await.unwrap();
476        let hash_b: Vec<u8> = sqlx::query_scalar!(
477            r#"SELECT row_hash AS "row_hash!" FROM audit_log WHERE id = ?1"#,
478            id_b
479        )
480        .fetch_one(&pool_b)
481        .await
482        .unwrap();
483
484        assert_eq!(
485            hash_a, hash_b,
486            "append_via_pool and append_in_tx must agree on row_hash for identical input"
487        );
488    }
489}