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}