From 23c3da3ee7b1ce870975176f2c52d06b6db310fa Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Thu, 3 Sep 2026 23:09:21 -0400 Subject: [PATCH] test(pds): pin a compaction against the log it replaces MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replays generated logs into every store, compacts, replays the compaction into fresh stores, and requires the same answers — including the blob reference counts a compaction rebuilds from the order it writes entries in. A second test requires compacting twice to equal compacting once, and both assert the corpus reached each state worth compacting. Change-Id: I2195be00d0e6bec11202211be67741f94c0b9c4c --- crates/didbot-pds/src/durable.rs | 683 +++++++++++++++++++++++++++++++ 1 file changed, 683 insertions(+) diff --git a/crates/didbot-pds/src/durable.rs b/crates/didbot-pds/src/durable.rs index 7fc69374..103f0e9f 100644 --- a/crates/didbot-pds/src/durable.rs +++ b/crates/didbot-pds/src/durable.rs @@ -2169,6 +2169,689 @@ impl AgentTokenStore for FileAgentTokenStore { } } +#[cfg(test)] +mod tests { + use super::*; + + use crate::kind::Retention; + use crate::ledger::LedgerEvent; + use crate::names::Reservation; + + /// Names the generated logs draw from, small enough that a sequence of + /// any length collides on every one of them and the interesting orders + /// — claim after release, put after delete, upload after collect — turn + /// up without being written out by hand. + const DIDS: [&str; 3] = [ + "did:web:one.agents.localhost", + "did:web:two.agents.localhost", + "did:web:three.agents.localhost", + ]; + const COLLECTIONS: [&str; 2] = ["com.example.thing", "com.example.other"]; + const RKEYS: [&str; 3] = ["alpha", "bravo", "charlie"]; + const LABELS: [&str; 8] = [ + "quernstone", + "farthingale", + "sillabub", + "wagonette", + "bezoar", + "gallimaufry", + "shivaree", + "widdershins", + ]; + const CIDS: [&str; 3] = [ + "bafkreialbcp7cyhcpbsyoy2ttdbtj7ijxvxuc2gqwmgnrfouyxrolfmdri", + "bafkreibvjvcv745gucpxq53k5nyfk3ehwyoypp2xrjvzt3xn5x5rjcxvpa", + "bafkreicrbqivsvi5kqiw52yfvinuprdd57r5jf7hebg5ah55x33ah7cage", + ]; + + /// How long a released name stays taken, for every registry here. + /// + /// Long enough that no hold generated below expires while the test runs, + /// so a difference between two captures is a difference the compaction + /// made rather than the clock. + fn hold() -> time::Duration { + time::Duration::days(30) + } + + /// A xorshift, so a failing seed reproduces exactly. + /// + /// In the test rather than from a crate because what is wanted is a fixed + /// sequence for a fixed number, which is the one property a generator's + /// version bump is free to change. + struct Seq(u64); + + impl Seq { + fn new(seed: u64) -> Self { + Self(seed | 1) + } + + fn next(&mut self) -> u64 { + let mut x = self.0; + x ^= x << 13; + x ^= x >> 7; + x ^= x << 17; + self.0 = x; + x + } + + fn pick(&mut self, len: usize) -> usize { + (self.next() % len as u64) as usize + } + + fn of<'a, T>(&mut self, from: &'a [T]) -> &'a T { + &from[self.pick(from.len())] + } + } + + /// The stores a replay pours into, with nothing above them. + /// + /// `Durable::open` is deliberately not used: what is under test is the + /// pair of functions this module holds — `apply` and `snapshot` — and + /// driving them directly is what lets a log be handed over entry by + /// entry rather than reproduced through the operations that wrote it. + struct Replayed { + dir: PathBuf, + accounts: MemoryAccountStore, + records: MemoryRecordStore, + blobs: FileBlobStore, + ledger: MemoryLedger, + commits: MemoryCommitStore, + names: NameRegistry, + credentials: MemoryAgentTokenStore, + counter: DurableCounter, + /// Read off the entries the way `Durable::open` reads it, because no + /// store holds it. + floor: u64, + } + + impl Replayed { + fn open(dir: PathBuf) -> Self { + let _ = std::fs::remove_dir_all(&dir); + crate::wal::create_dir(&dir).expect("a scratch directory"); + let (wal, _) = Wal::open(&dir).expect("a log to hang the blob store off"); + let blobs = FileBlobStore::open(&dir, Arc::new(wal), BlobLimits::default()) + .expect("a blob store"); + Self { + dir, + accounts: MemoryAccountStore::new(), + records: MemoryRecordStore::new(), + blobs, + ledger: MemoryLedger::new(), + commits: MemoryCommitStore::new(), + names: NameRegistry::new(hold()), + credentials: MemoryAgentTokenStore::new(), + counter: DurableCounter::new(), + floor: 0, + } + } + + /// Replays a whole log, in order. + fn feed(&mut self, entries: &[Entry]) { + for entry in entries { + if let Entry::StreamReserved { through } = entry { + self.floor = self.floor.max(*through); + } + apply( + &self.accounts, + &self.records, + &self.blobs, + &self.ledger, + &self.commits, + &self.names, + &self.credentials, + &self.counter, + entry.clone(), + ); + } + } + + /// The log this state would be written as. + /// + /// Destructive — `snapshot` drains the record store — so a capture + /// has to be taken first. + fn compact(&self) -> Vec { + snapshot( + &self.accounts, + &self.records, + &self.blobs, + &self.ledger, + &self.commits, + &self.names, + &self.credentials, + &self.counter, + self.floor, + ) + } + + /// Everything the deployment would answer with, as comparable text. + /// + /// Read through each store's own accessors rather than through + /// `snapshot`, so that a compaction which dropped something and a + /// capture which failed to look for it cannot cancel out. + fn capture(&self) -> Captured { + let now = OffsetDateTime::now_utc(); + // Far enough ahead that the grace window is not what decides the + // answer: what is wanted is which blobs nothing references. + let cutoff = now + time::Duration::days(3650); + + let mut accounts: Vec = self + .accounts + .snapshot_with_keys() + .into_iter() + .map(|(account, key)| format!("{account:?} {:?}", key.to_bytes())) + .collect(); + accounts.sort(); + + let mut records = Vec::new(); + for did in DIDS { + for (collection, rkey, value) in self.records.snapshot(did) { + records.push(format!("{did} {collection} {rkey} {value}")); + } + } + records.sort(); + + let mut blobs: Vec = self + .blobs + .snapshot() + .into_iter() + .map(|(did, reference, at)| format!("{did} {reference:?} {at}")) + .collect(); + blobs.sort(); + + let mut unreferenced: Vec = self + .blobs + .inner + .index + .collect_candidates(cutoff) + .into_iter() + .map(|(did, cid, size)| format!("{did} {cid} {size}")) + .collect(); + unreferenced.sort(); + + let mut ledger: Vec = self + .ledger + .all() + .into_iter() + .map(|agent| format!("{agent:?}")) + .collect(); + ledger.sort(); + + let mut heads: Vec = self + .commits + .heads() + .into_iter() + .map(|(did, head)| format!("{did} {head:?}")) + .collect(); + heads.sort(); + + // The registry's own answers, which is what every caller asks + // it, rather than the entries it would write. + let names: Vec = LABELS + .iter() + .map(|label| { + format!( + "{label} {:?} {:?}", + self.names.check(label, now), + self.names.reservation_of(label) + ) + }) + .collect(); + + let mut credentials: Vec = self + .credentials + .snapshot() + .into_iter() + .map(|(did, hash, expires)| format!("{did} {hash} {expires}")) + .collect(); + credentials.sort(); + + Captured { + accounts, + records, + blobs, + unreferenced, + ledger, + heads, + names, + credentials, + counter: self.counter.peek(), + floor: self.floor, + } + } + } + + impl Drop for Replayed { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.dir); + } + } + + /// Every answer a deployment gives, in a form two of them can be compared + /// by. A field per store, so a failure names which one moved. + #[derive(Debug, PartialEq, Eq)] + struct Captured { + accounts: Vec, + records: Vec, + blobs: Vec, + /// The blobs a collection pass would take: the observable form of + /// the reference counts a compaction rebuilds rather than carries. + unreferenced: Vec, + ledger: Vec, + heads: Vec, + names: Vec, + credentials: Vec, + counter: u64, + floor: u64, + } + + /// A record carrying exactly `cids` as blob references. + fn record_over(cids: &[String]) -> Value { + let attachments: Vec = cids + .iter() + .map(|cid| { + BlobRef { + cid: cid.clone(), + mime_type: "image/png".to_owned(), + size: 7, + } + .to_json() + }) + .collect(); + serde_json::json!({ "text": "quernstone", "attachments": attachments }) + } + + /// A log of `len` entries, drawn from every variant that carries state. + /// + /// The generator mirrors the index as it goes, so that the orders it + /// emits are orders this server can produce. Two of them are load-bearing + /// and neither is arbitrary: + /// + /// * A record only ever references blobs already uploaded under the same + /// DID. `FileBlobStore::commit` renames the bytes and appends + /// `BlobUploaded` before the reference can be written, so a `RecordPut` + /// never precedes the upload it names. + /// * A blob is only collected while no live record references it. + /// `collect_unreferenced` takes candidates from + /// `BlobIndex::collect_candidates`, which is exactly the blobs at zero + /// references. + /// + /// Generating without those produces a log whose replayed reference + /// counts a compaction cannot reproduce — correctly, because the counts + /// came from a history that never happened. + /// + /// `StreamReserved` is generated too, so the reservation a compaction + /// must never drop is part of what is compared. + fn generated_log(seed: u64, len: usize) -> Vec { + let mut seq = Seq::new(seed); + // Anchored to the clock rather than to a constant, because two of + // the states being compared are clock-relative: a released name is + // held only until `at + hold`, and a log stamped two years ago + // replays into a registry where every hold has already expired and + // the held state never appears. The *structure* of the log is still + // the seed's alone; only the instants move. + let base = OffsetDateTime::now_utc() - time::Duration::hours(1); + let mut floor = 0u64; + let mut next = 0u64; + let mut ledger_seq: BTreeMap = BTreeMap::new(); + // The index and the records, as the generated log has left them. + let mut held: BTreeMap> = BTreeMap::new(); + let mut written: BTreeMap<(String, String, String), Vec> = BTreeMap::new(); + let mut entries = Vec::with_capacity(len); + + for step in 0..len { + let did = (*seq.of(&DIDS)).to_owned(); + let at = base + time::Duration::seconds(step as i64); + let entry = match seq.pick(16) { + 0 => { + let parsed = AgentDid::parse(&did).expect("a did the alphabet holds"); + Entry::AccountInserted { + account: Box::new(AgentAccount::server(parsed, at)), + key: SigningKey::generate(), + } + } + 1 => Entry::AccountRemoved { did }, + // The variant no build writes any more, whose effect a + // compaction has to carry as a retention. + 2 => Entry::AccountPinned { + did, + pinned: seq.pick(2) == 0, + }, + 3 => Entry::AccountRetained { + did, + retention: *seq.of(&[ + Retention::Forever, + Retention::Deployment, + Retention::Until { seconds: 3600 }, + ]), + }, + 4 => Entry::AccountStateChanged { + did, + state: *seq.of(&[ + AccountState::Active, + AccountState::Frozen, + AccountState::SoftDeleted, + ]), + }, + 5 => { + let collection = (*seq.of(&COLLECTIONS)).to_owned(); + let rkey = (*seq.of(&RKEYS)).to_owned(); + // Only blobs this DID has already uploaded, which is the + // only kind a client could name. + let available = held.get(&did).cloned().unwrap_or_default(); + let mut refs: Vec = Vec::new(); + for _ in 0..seq.pick(3) { + if available.is_empty() { + break; + } + let cid = seq.of(&available).clone(); + if !refs.contains(&cid) { + refs.push(cid); + } + } + let record = record_over(&refs); + written.insert((did.clone(), collection.clone(), rkey.clone()), refs); + Entry::RecordPut { + did, + collection, + rkey, + record, + } + } + 6 => { + let collection = (*seq.of(&COLLECTIONS)).to_owned(); + let rkey = (*seq.of(&RKEYS)).to_owned(); + written.remove(&(did.clone(), collection.clone(), rkey.clone())); + Entry::RecordRemoved { + did, + collection, + rkey, + } + } + 7 => { + let cid = (*seq.of(&CIDS)).to_owned(); + let account = held.entry(did.clone()).or_default(); + if account.contains(&cid) { + // Already uploaded. An entry that changes nothing + // rather than one this server would not write. + Entry::AccountRemoved { did } + } else { + account.push(cid.clone()); + Entry::BlobUploaded { + did, + cid, + mime_type: "image/png".to_owned(), + size: 7, + at, + } + } + } + 8 => { + // Only a blob nothing points at, which is the only kind + // the collector is ever handed. + let referenced: Vec<&String> = written + .iter() + .filter(|((owner, _, _), _)| owner == &did) + .flat_map(|(_, refs)| refs.iter()) + .collect(); + let free: Vec = held + .get(&did) + .into_iter() + .flatten() + .filter(|cid| !referenced.contains(cid)) + .cloned() + .collect(); + if free.is_empty() { + Entry::AccountRemoved { did } + } else { + let cid = seq.of(&free).clone(); + held.entry(did.clone()) + .or_default() + .retain(|held| held != &cid); + Entry::BlobCollected { did, cid } + } + } + 9 => Entry::RepoCommitted { + did, + rev: format!("3l{step}"), + commit: (*seq.of(&CIDS)).to_owned(), + data: Some((*seq.of(&CIDS)).to_owned()), + prev: (seq.pick(2) == 0).then(|| (*seq.of(&CIDS)).to_owned()), + created: vec![(*seq.of(&CIDS)).to_owned()], + }, + 10 => { + held.remove(&did); + written.retain(|(owner, _, _), _| owner != &did); + Entry::RepoRemoved { did } + } + 11 => { + let seq_no = ledger_seq.entry(did.clone()).or_insert(0); + *seq_no += 1; + Entry::LedgerAppended { + did, + entry: crate::ledger::LedgerEntry { + seq: *seq_no, + at, + event: LedgerEvent::Pinned, + }, + } + } + 12 => Entry::NameClaimed { + name: (*seq.of(&LABELS)).to_owned(), + }, + 13 => { + if seq.pick(4) == 0 { + Entry::NameReserved { + name: (*seq.of(&LABELS)).to_owned(), + reservation: Reservation::Operational { + reason: "generated".to_owned(), + }, + } + } else { + Entry::NameReleased { + name: (*seq.of(&LABELS)).to_owned(), + at, + } + } + } + 14 => { + if seq.pick(3) == 0 { + Entry::AgentTokenRevoked { did } + } else { + Entry::AgentTokenIssued { + did, + token_hash: format!("{:064x}", seq.next()), + expires_at: at + time::Duration::days(30), + } + } + } + _ => { + if seq.pick(2) == 0 { + next += 1 + seq.next() % 4; + Entry::CounterAdvanced { next } + } else { + floor += 1 + seq.next() % 8; + Entry::StreamReserved { through: floor } + } + } + }; + entries.push(entry); + } + entries + } + + /// Which of the states worth compacting the generated corpus reached. + /// + /// A property test over a corpus that never produces a held name, or + /// never produces a blob something references, passes for the wrong + /// reason. This is asserted at the end of the run, so a change to the + /// generator that quietly stops reaching one of these fails rather than + /// narrows the test. + #[derive(Debug, Default)] + struct Reached { + account: bool, + record: bool, + referenced_blob: bool, + unreferenced_blob: bool, + ledger_trail: bool, + head: bool, + live_name: bool, + held_name: bool, + reserved_name: bool, + credential: bool, + counter: bool, + reservation: bool, + } + + impl Reached { + fn note(&mut self, deployment: &Replayed, captured: &Captured) { + let now = OffsetDateTime::now_utc(); + self.account |= !captured.accounts.is_empty(); + self.record |= !captured.records.is_empty(); + self.referenced_blob |= captured.blobs.len() > captured.unreferenced.len(); + self.unreferenced_blob |= !captured.unreferenced.is_empty(); + self.ledger_trail |= deployment + .ledger + .all() + .iter() + .any(|agent| agent.entries.len() >= 2); + self.head |= !captured.heads.is_empty(); + for label in LABELS { + match deployment.names.check(label, now) { + Ok(()) => {} + Err(crate::names::NameError::Live(_)) => self.live_name = true, + Err(crate::names::NameError::Held { .. }) => self.held_name = true, + Err(crate::names::NameError::Reserved { .. }) => self.reserved_name = true, + Err(_) => {} + } + } + self.credential |= !captured.credentials.is_empty(); + self.counter |= captured.counter > 0; + self.reservation |= captured.floor > 0; + } + + fn assert_whole(&self) { + let missed: Vec<&str> = [ + ("an account", self.account), + ("a record", self.record), + ("a blob a record references", self.referenced_blob), + ("a blob nothing references", self.unreferenced_blob), + ("a ledger of more than one entry", self.ledger_trail), + ("a repository head", self.head), + ("a live name", self.live_name), + ("a released name still held", self.held_name), + ("a reserved name", self.reserved_name), + ("an agent token", self.credential), + ("an advanced counter", self.counter), + ("a stream reservation", self.reservation), + ] + .into_iter() + .filter_map(|(what, seen)| (!seen).then_some(what)) + .collect(); + assert!( + missed.is_empty(), + "the generated corpus never produced: {}", + missed.join(", ") + ); + } + } + + /// The invariant the whole compaction rests on, over generated logs. + /// + /// A compaction replaces a log with the entries `snapshot` produces. That + /// is only safe if replaying those entries lands in the same place the + /// log they replaced lands in — for every store at once, because they + /// share one file and one order and a compaction that got the order wrong + /// between two of them (blobs after the records that reference them, say) + /// leaves each store individually plausible and the pair wrong. + /// + /// So each seed builds a log out of every variant that carries state, + /// replays it, records what all eight stores answer, compacts, replays + /// the compaction into fresh stores, and requires the same answers. The + /// generated shapes are the ones a hand-written case is least likely to + /// reach: a record put over a record that referenced a different blob, a + /// name claimed after it was released, an account deleted and + /// reprovisioned, a blob collected and uploaded again. + #[test] + fn a_compacted_log_answers_exactly_what_the_log_it_replaced_did() { + let mut reached = Reached::default(); + for seed in 1..=24u64 { + let log = generated_log(seed, 160); + + let mut original = Replayed::open( + std::env::temp_dir() + .join(format!("didbot-compact-{}-{seed}-a", std::process::id())), + ); + original.feed(&log); + let before = original.capture(); + reached.note(&original, &before); + let compacted = original.compact(); + + let mut rebuilt = Replayed::open( + std::env::temp_dir() + .join(format!("didbot-compact-{}-{seed}-b", std::process::id())), + ); + rebuilt.feed(&compacted); + let after = rebuilt.capture(); + + assert_eq!( + before, + after, + "seed {seed}: a compaction of {} entries into {} changed what the deployment \ + answers", + log.len(), + compacted.len() + ); + } + reached.assert_whole(); + } + + /// A compaction is a fixed point: compacting a compacted log changes + /// nothing. + /// + /// Which is the property that makes it safe to run at every startup. A + /// compaction that added an entry each time — or dropped one each time, + /// slowly — would pass a single-round comparison and lose the deployment + /// over a month of restarts. + #[test] + fn compacting_a_compacted_log_produces_the_same_log_again() { + for seed in 1..=12u64 { + let log = generated_log(seed, 160); + + let mut once = Replayed::open( + std::env::temp_dir().join(format!("didbot-fixed-{}-{seed}-a", std::process::id())), + ); + once.feed(&log); + let first = once.compact(); + + let mut twice = Replayed::open( + std::env::temp_dir().join(format!("didbot-fixed-{}-{seed}-b", std::process::id())), + ); + twice.feed(&first); + let second = twice.compact(); + + // Sorted, because `NameRegistry` holds its names in a hash map + // and the order its snapshot emits them in is not stable between + // two registries. Nothing replays differently for it: what this + // asserts is that no entry is gained or lost, and the one + // ordering that does matter — blobs ahead of the records that + // reference them — is what + // `a_compacted_log_answers_exactly_what_the_log_it_replaced_did` + // measures, through the reference counts it rebuilds. + let render = |entries: &[Entry]| -> Vec { + let mut rendered: Vec = + entries.iter().map(|entry| format!("{entry:?}")).collect(); + rendered.sort(); + rendered + }; + assert_eq!( + render(&first), + render(&second), + "seed {seed}: compacting twice is not the same as compacting once" + ); + } + } +} + /// The OAuth grant store, with every family it mints, rotates and revokes in /// the log. /// -- 2.51.2