use anyhow::{bail, Result};
use rusqlite::{Connection, TransactionBehavior};
use super::Store;
use crate::audit::merkle::{self, Hash, Tree};
pub(super) const SCHEMA: &str = "
CREATE TABLE IF NOT EXISTS audit_log (
seq INTEGER PRIMARY KEY,
leaf BLOB NOT NULL,
leaf_hash BLOB NOT NULL
);
CREATE TRIGGER IF NOT EXISTS audit_log_no_update
BEFORE UPDATE ON audit_log
BEGIN
SELECT RAISE(ABORT, 'audit_log is append-only');
END;
CREATE TRIGGER IF NOT EXISTS audit_log_no_delete
BEFORE DELETE ON audit_log
BEGIN
SELECT RAISE(ABORT, 'audit_log is append-only');
END;
CREATE TRIGGER IF NOT EXISTS audit_log_next_seq
BEFORE INSERT ON audit_log
WHEN NEW.seq IS NOT (SELECT COALESCE(MAX(seq), -1) + 1 FROM audit_log)
BEGIN
SELECT RAISE(ABORT, 'audit_log is append-only: a leaf takes the next seq');
END;
";
pub(super) struct Loaded {
pub(super) tree: Tree,
pub(super) last_at: String,
}
pub(super) fn load(conn: &Connection) -> Result<Loaded> {
let mut tree = Tree::new();
read_from(conn, 0, |hash| tree.append(hash))?;
let last_at = last_at(conn, tree.size())?;
Ok(Loaded { tree, last_at })
}
fn read_from(conn: &Connection, from: u64, each: impl FnMut(Hash)) -> Result<u64> {
read_range(conn, from, i64::MAX as u64, each)
}
fn read_range(conn: &Connection, from: u64, most: u64, mut each: impl FnMut(Hash)) -> Result<u64> {
let mut stmt = conn.prepare(
"SELECT seq, leaf, leaf_hash FROM audit_log WHERE seq >= ?1 ORDER BY seq LIMIT ?2",
)?;
let mut rows = stmt.query((from as i64, most.min(i64::MAX as u64) as i64))?;
let mut want = from;
while let Some(row) = rows.next()? {
let seq: i64 = row.get(0)?;
if seq != want as i64 {
bail!(
"the audit log is damaged: leaf {want} is missing (the next one stored is {seq}); \
restore the database from a backup"
);
}
let leaf = row.get_ref(1)?.as_blob()?;
let hash = merkle::hash_leaf(leaf);
if row.get_ref(2)?.as_blob()? != hash.as_slice() {
bail!(
"the audit log is damaged: leaf {seq} no longer hashes to its stored leaf_hash; \
restore the database from a backup"
);
}
each(hash);
want += 1;
}
Ok(want - from)
}
fn last_at(conn: &Connection, size: u64) -> Result<String> {
if size == 0 {
return Ok(String::new());
}
let leaf: Vec<u8> = conn.query_row(
"SELECT leaf FROM audit_log WHERE seq = ?1",
(size as i64 - 1,),
|r| r.get(0),
)?;
Ok(serde_json::from_slice::<serde_json::Value>(&leaf)
.ok()
.and_then(|v| v.get("at")?.as_str().map(str::to_string))
.unwrap_or_default())
}
fn catch_up(conn: &Connection, held: u64) -> Result<(Vec<Hash>, Option<String>)> {
let stored: i64 =
conn.query_row("SELECT COALESCE(MAX(seq) + 1, 0) FROM audit_log", [], |r| {
r.get(0)
})?;
if stored < held as i64 {
bail!(
"the audit log holds {stored} leaves, fewer than the {held} this server appended: \
the database was replaced under a running server; restart it"
);
}
let mut hashes = Vec::new();
if stored == held as i64 {
return Ok((hashes, None));
}
let read = read_from(conn, held, |hash| hashes.push(hash))?;
Ok((hashes, Some(last_at(conn, held + read)?)))
}
pub enum Outcome<T> {
Commit(T),
Refuse(T),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AuditEntry {
pub seq: u64,
pub leaf: Vec<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
pub enum ConsistencyError {
#[error("first must be at least 1 and at most second")]
BadRange,
#[error("second is past the end of the log")]
SecondBeyondTreeSize,
}
fn next_at(one: &str, other: &str) -> String {
let newest = one.max(other);
let now = crate::now();
if now.as_str() < newest {
newest.to_string()
} else {
now
}
}
impl Store {
pub fn audited<T>(
&self,
write: impl FnOnce(&rusqlite::Transaction, &str) -> Result<Outcome<T>>,
build_leaf: impl FnOnce(u64, &str, &T) -> Vec<u8>,
) -> Result<T> {
self.audited_each(write, |seq, at, value| vec![build_leaf(seq, at, value)])
}
pub fn audited_each<T>(
&self,
write: impl FnOnce(&rusqlite::Transaction, &str) -> Result<Outcome<T>>,
build_leaves: impl FnOnce(u64, &str, &T) -> Vec<Vec<u8>>,
) -> Result<T> {
self.audited_each_as(
anyhow::Error::from,
anyhow::Error::from,
write,
build_leaves,
)
}
pub(crate) fn audited_each_as<T>(
&self,
begin_failed: impl FnOnce(rusqlite::Error) -> anyhow::Error,
commit_failed: impl FnOnce(rusqlite::Error) -> anyhow::Error,
write: impl FnOnce(&rusqlite::Transaction, &str) -> Result<Outcome<T>>,
build_leaves: impl FnOnce(u64, &str, &T) -> Vec<Vec<u8>>,
) -> Result<T> {
let mut state = self.lock();
let (held, held_at) = (state.audit.size(), state.audit_at.clone());
let tx = state
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(begin_failed)?;
let (caught, caught_at) = catch_up(&tx, held)?;
let seq = held + caught.len() as u64;
let at = next_at(caught_at.as_deref().unwrap_or(""), &held_at);
let value = match write(&tx, &at)? {
Outcome::Commit(value) => value,
Outcome::Refuse(value) => return Ok(value), };
let leaves = build_leaves(seq, &at, &value);
let mut hashes = Vec::with_capacity(leaves.len());
for (i, leaf) in leaves.iter().enumerate() {
let leaf_hash = merkle::hash_leaf(leaf);
tx.execute(
"INSERT INTO audit_log (seq, leaf, leaf_hash) VALUES (?1, ?2, ?3)",
(seq as i64 + i as i64, leaf, leaf_hash.as_slice()),
)?;
hashes.push(leaf_hash);
}
tx.commit().map_err(commit_failed)?;
for leaf_hash in caught.into_iter().chain(hashes) {
state.audit.append(leaf_hash);
}
if !leaves.is_empty() {
state.audit_at = at;
} else if let Some(caught_at) = caught_at.filter(|c| *c > state.audit_at) {
state.audit_at = caught_at;
}
Ok(value)
}
pub(crate) fn audit_read_ahead(&self) -> Result<()> {
self.audit_read_ahead_by(4096)
}
fn audit_read_ahead_by(&self, chunk: u64) -> Result<()> {
let mut state = self.lock();
loop {
let held = state.audit.size();
let mut hashes = Vec::new();
let read = read_range(&state.conn, held, chunk, |hash| hashes.push(hash))?;
if read == 0 {
return Ok(());
}
let at = last_at(&state.conn, held + read)?;
for hash in hashes {
state.audit.append(hash);
}
if at > state.audit_at {
state.audit_at = at;
}
}
}
pub fn audit_append(&self, build_leaf: impl FnOnce(u64, &str) -> Vec<u8>) -> Result<u64> {
let assigned = std::cell::Cell::new(0u64);
self.audited(
|_tx, _at| Ok(Outcome::Commit(())),
|seq, at, ()| {
assigned.set(seq);
build_leaf(seq, at)
},
)?;
Ok(assigned.get())
}
pub fn audit_checkpoint(&self) -> (u64, Hash) {
let state = self.lock();
(state.audit.size(), state.audit.root())
}
pub fn audit_entries(&self, start: u64, end: u64, max_bytes: usize) -> Result<Vec<AuditEntry>> {
let state = self.lock();
let mut stmt = state
.conn
.prepare("SELECT seq, leaf FROM audit_log WHERE seq >= ?1 AND seq < ?2 ORDER BY seq")?;
let mut rows = stmt.query((start as i64, end as i64))?;
let (mut out, mut bytes) = (Vec::new(), 0usize);
while let Some(row) = rows.next()? {
let leaf: Vec<u8> = row.get(1)?;
bytes += leaf.len();
if bytes > max_bytes && !out.is_empty() {
break;
}
out.push(AuditEntry {
seq: row.get::<_, i64>(0)? as u64,
leaf,
});
}
Ok(out)
}
pub fn audit_consistency(
&self,
first: u64,
second: u64,
) -> Result<Result<Vec<Hash>, ConsistencyError>> {
let state = self.lock();
if first == 0 || first > second {
return Ok(Err(ConsistencyError::BadRange));
}
if second > state.audit.size() {
return Ok(Err(ConsistencyError::SecondBeyondTreeSize));
}
Ok(Ok(state.audit.consistency(first, second)))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::audit::leaf;
fn store() -> Store {
Store::open_in_memory().unwrap()
}
fn push_leaf(seq: u64) -> Vec<u8> {
leaf::encode(
seq,
"2026-01-01T00:00:00.000Z",
leaf::action::PUSH,
&leaf::Actor::Operator,
leaf::subject_file(&leaf::FileChange {
project_key: "acme/app",
file_path: "a.md",
deleted: false,
stored_sha256: "abc",
base_sha256: None,
merged: false,
merge_job: None,
}),
None,
)
}
fn append(st: &Store) {
st.audited(
|_tx, _| Ok(Outcome::Commit(())),
|seq, _, ()| push_leaf(seq),
)
.unwrap();
}
const INSERT_FILE: &str = "INSERT INTO memory_files \
(project_key, file_path, content, source_env, updated_at) VALUES ('a','b','c','d','e')";
#[test]
fn a_committed_write_appends_exactly_one_leaf() {
let st = store();
st.audited(
|tx, _| {
tx.execute(INSERT_FILE, [])?;
Ok(Outcome::Commit(()))
},
|seq, _, ()| push_leaf(seq),
)
.unwrap();
let (size, _) = st.audit_checkpoint();
assert_eq!(size, 1);
assert_eq!(st.audit_entries(0, 1, usize::MAX).unwrap().len(), 1);
}
#[test]
fn a_failing_write_appends_no_leaf_and_keeps_no_change() {
let st = store();
let result = st.audited(
|tx, _| {
tx.execute(INSERT_FILE, [])?;
Err(anyhow::Error::from(rusqlite::Error::ExecuteReturnedResults))
},
|seq, _, ()| push_leaf(seq),
);
assert!(result.is_err());
assert_eq!(
st.audit_checkpoint().0,
0,
"no leaf from a rolled-back write"
);
assert!(st.get("a", "b").unwrap().is_none(), "no row either");
}
#[test]
fn a_leaf_that_cannot_be_written_undoes_the_change() {
let st = store();
st.with_raw(|c| {
c.execute_batch(
"CREATE TEMP TRIGGER no_leaves BEFORE INSERT ON audit_log
BEGIN SELECT RAISE(ABORT, 'no leaves today'); END;",
)
})
.unwrap();
let result = st.audited(
|tx, _| {
tx.execute(INSERT_FILE, [])?;
Ok(Outcome::Commit(()))
},
|seq, _, ()| push_leaf(seq),
);
assert!(result.is_err(), "the leaf's insert failed");
assert!(
st.get("a", "b").unwrap().is_none(),
"so the row is not there"
);
assert_eq!(st.audit_checkpoint().0, 0);
}
#[test]
fn a_refused_write_appends_no_leaf() {
let st = store();
let refusal: &str = st
.audited(
|_tx, _| Ok(Outcome::Refuse("no such code")),
|seq, _, _| push_leaf(seq),
)
.unwrap();
assert_eq!(refusal, "no such code");
assert_eq!(st.audit_checkpoint().0, 0);
}
#[test]
fn one_write_can_append_several_leaves_in_order() {
let st = store();
append(&st);
st.audited_each(
|_tx, _| Ok(Outcome::Commit(3u64)),
|seq, _, n| (seq..seq + n).map(push_leaf).collect(),
)
.unwrap();
let seqs: Vec<u64> = st
.audit_entries(0, 4, usize::MAX)
.unwrap()
.iter()
.map(|e| e.seq)
.collect();
assert_eq!(seqs, vec![0, 1, 2, 3]);
assert_eq!(
st.audit_entries(1, 2, usize::MAX).unwrap()[0].leaf,
push_leaf(1)
);
}
#[test]
fn at_never_decreases_along_seq() {
let st = store();
let mut ats = Vec::new();
for _ in 0..3 {
st.audit_append(|seq, at| {
ats.push(at.to_string());
push_leaf(seq)
})
.unwrap();
}
assert!(ats.windows(2).all(|w| w[0] <= w[1]), "{ats:?}");
assert_eq!(ats[0].len(), 24, "the API's timestamp format");
st.lock().audit_at = "2999-01-01T00:00:00.000Z".into();
let mut got = String::new();
st.audit_append(|seq, at| {
got = at.to_string();
push_leaf(seq)
})
.unwrap();
assert_eq!(got, "2999-01-01T00:00:00.000Z");
}
#[test]
fn checkpoint_matches_the_merkle_root_of_every_leaf() {
let st = store();
for _ in 0..5 {
append(&st);
}
let (size, root) = st.audit_checkpoint();
assert_eq!(size, 5);
let leaves: Vec<Hash> = (0..5)
.map(push_leaf)
.map(|l| merkle::hash_leaf(&l))
.collect();
assert_eq!(root, merkle::root(&leaves));
}
#[test]
fn entries_pages_by_seq_and_by_bytes() {
let st = store();
for _ in 0..3 {
append(&st);
}
let got = st.audit_entries(1, 3, usize::MAX).unwrap();
assert_eq!(got.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![1, 2]);
assert_eq!(got[0].leaf, push_leaf(1));
let one = push_leaf(0).len();
let seqs = |max| -> Vec<u64> {
st.audit_entries(0, 3, max)
.unwrap()
.iter()
.map(|e| e.seq)
.collect()
};
assert_eq!(seqs(2 * one), vec![0, 1], "two fit exactly");
assert_eq!(seqs(2 * one - 1), vec![0], "the second would overflow");
assert_eq!(seqs(1), vec![0], "never fewer than one");
}
#[test]
fn consistency_matches_merkle_and_rejects_bad_ranges() {
let st = store();
for _ in 0..8 {
append(&st);
}
let proof = st.audit_consistency(3, 8).unwrap().unwrap();
let leaves: Vec<Hash> = (0..8)
.map(push_leaf)
.map(|l| merkle::hash_leaf(&l))
.collect();
assert_eq!(proof, merkle::consistency(3, 8, &leaves));
assert_eq!(
st.audit_consistency(0, 8).unwrap(),
Err(ConsistencyError::BadRange)
);
assert_eq!(
st.audit_consistency(5, 3).unwrap(),
Err(ConsistencyError::BadRange)
);
assert_eq!(
st.audit_consistency(1, 100).unwrap(),
Err(ConsistencyError::SecondBeyondTreeSize)
);
}
#[test]
fn a_proof_does_not_read_the_table() {
let st = store();
for _ in 0..40 {
append(&st);
}
let want = st.audit_consistency(7, 40).unwrap().unwrap();
st.with_raw(|c| c.execute_batch("ALTER TABLE audit_log RENAME TO audit_log_away"))
.unwrap();
assert_eq!(st.audit_consistency(7, 40).unwrap().unwrap(), want);
}
#[test]
fn nothing_but_the_next_leaf_can_be_written() {
let st = store();
append(&st);
append(&st);
let refused = |sql: &str| {
assert!(st.with_raw(|c| c.execute(sql, [])).is_err(), "{sql}");
};
refused("UPDATE audit_log SET seq = 99 WHERE seq = 0");
refused("DELETE FROM audit_log WHERE seq = 0");
refused(
"INSERT OR REPLACE INTO audit_log (seq, leaf, leaf_hash) \
VALUES (0, CAST('forged' AS BLOB), zeroblob(32))",
);
refused(
"INSERT INTO audit_log (seq, leaf, leaf_hash) VALUES (0, x'01', zeroblob(32)) \
ON CONFLICT(seq) DO UPDATE SET leaf = excluded.leaf",
);
refused("INSERT INTO audit_log (seq, leaf, leaf_hash) VALUES (100, x'02', zeroblob(32))");
refused("INSERT INTO audit_log (leaf, leaf_hash) VALUES (x'02', zeroblob(32))");
assert_eq!(
st.audit_entries(0, 2, usize::MAX).unwrap()[0].leaf,
push_leaf(0),
"leaf 0 is what was written"
);
assert_eq!(st.audit_checkpoint().0, 2);
st.with_raw(|c| {
c.execute(
"INSERT INTO audit_log (seq, leaf, leaf_hash) VALUES (2, x'03', zeroblob(32))",
[],
)
})
.unwrap();
}
#[test]
fn reopening_rebuilds_the_same_tree() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("r.db");
let before = {
let st = Store::open(&path).unwrap();
for _ in 0..5 {
append(&st);
}
st.audit_checkpoint()
};
let st = Store::open(&path).unwrap();
assert_eq!(st.audit_checkpoint(), before);
assert_eq!(st.audit_append(|seq, _| push_leaf(seq)).unwrap(), 5);
}
#[test]
fn a_leaf_another_process_appended_is_read_in_before_the_next() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("r.db");
let server = Store::open(&path).unwrap();
append(&server);
append(&server);
let host = Store::open(&path).unwrap();
assert_eq!(host.audit_append(|seq, _| push_leaf(seq)).unwrap(), 2);
assert_eq!(server.audit_checkpoint().0, 2, "not read until it appends");
assert_eq!(server.audit_append(|seq, _| push_leaf(seq)).unwrap(), 3);
assert_eq!(
server.audit_checkpoint(),
Store::open(&path).unwrap().audit_checkpoint()
);
let blank = Store::with_connection(Connection::open(&path).unwrap()).unwrap();
blank.lock().audit = Tree::new();
assert_eq!(blank.audit_append(|seq, _| push_leaf(seq)).unwrap(), 4);
let older = dir.path().join("older.db");
{
let st = Store::open(&older).unwrap();
append(&st);
}
std::fs::copy(&older, &path).unwrap();
let err = server.audit_append(|seq, _| push_leaf(seq)).unwrap_err();
assert!(format!("{err:#}").contains("restart it"), "{err:#}");
}
#[test]
fn reading_ahead_builds_the_tree_catching_up_would() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("r.db");
let server = Store::open(&path).unwrap();
for _ in 0..5 {
append(&server);
}
let host = Store::with_connection(Connection::open(&path).unwrap()).unwrap();
host.lock().audit = Tree::new();
host.lock().audit_at = String::new();
host.audit_read_ahead_by(2).unwrap();
assert_eq!(host.audit_checkpoint(), server.audit_checkpoint());
assert_eq!(host.lock().audit_at, "2026-01-01T00:00:00.000Z");
append(&server);
host.audit_read_ahead_by(2).unwrap();
assert_eq!(host.audit_checkpoint(), server.audit_checkpoint());
assert_eq!(host.audit_append(|seq, _| push_leaf(seq)).unwrap(), 6);
assert_eq!(server.audit_append(|seq, _| push_leaf(seq)).unwrap(), 7);
assert_eq!(host.audit_append(|seq, _| push_leaf(seq)).unwrap(), 8);
let conn = Connection::open(&path).unwrap();
conn.execute_batch(
"DROP TRIGGER audit_log_no_update;
UPDATE audit_log SET leaf = CAST('forged' AS BLOB) WHERE seq = 3",
)
.unwrap();
let blank =
Store::with_connection(Connection::open(dir.path().join("x.db")).unwrap()).unwrap();
blank.lock().conn = Connection::open(&path).unwrap();
let err = blank.audit_read_ahead_by(2).unwrap_err();
assert!(
format!("{err:#}").contains("leaf 3 no longer hashes"),
"{err:#}"
);
}
#[test]
fn opening_refuses_a_damaged_log() {
let damaged = |how: &str| -> String {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("r.db");
{
let st = Store::open(&path).unwrap();
for _ in 0..4 {
append(&st);
}
}
let conn = Connection::open(&path).unwrap();
conn.execute_batch(&format!(
"DROP TRIGGER audit_log_no_update; DROP TRIGGER audit_log_no_delete; {how}"
))
.unwrap();
drop(conn);
match Store::open(&path) {
Ok(_) => panic!("opened a log damaged by {how}"),
Err(e) => format!("{e:#}"),
}
};
let e = damaged("UPDATE audit_log SET leaf = CAST('forged' AS BLOB) WHERE seq = 1");
assert!(e.contains("leaf 1 no longer hashes"), "{e}");
let e = damaged("DELETE FROM audit_log WHERE seq = 2");
assert!(e.contains("leaf 2 is missing"), "{e}");
let e = damaged("UPDATE audit_log SET leaf_hash = zeroblob(32) WHERE seq = 3");
assert!(e.contains("leaf 3 no longer hashes"), "{e}");
}
}