diff --git a/crates/didbot-pds/src/durable.rs b/crates/didbot-pds/src/durable.rs index cec36632..b7c06e26 100644 --- a/crates/didbot-pds/src/durable.rs +++ b/crates/didbot-pds/src/durable.rs @@ -751,33 +751,31 @@ pub struct FileCommitStore { wal: Arc, } -impl CommitStore for FileCommitStore { - fn head(&self, did: &str) -> Option { - self.inner.head(did) - } - - fn advance( +impl FileCommitStore { + /// Advances the store under this one with `advance`, then logs the commit + /// with the created set it answered. + /// + /// Applied and then appended, which is the reverse of every other durable + /// store here and is the one place it is right: what goes in the log is + /// the created set, and the created set is the answer the store under + /// this one works out. There is nothing to append until it has. + /// + /// A log failure is loud rather than fatal. The write this commit is over + /// has already been acknowledged, and refusing to remember where it left + /// the repository would be refusing to serve a repository that exists; + /// what it costs is a `since` cursor that is answered with the whole + /// repository after the next restart. + fn logged( &self, did: &str, head: Head, - blocks: std::collections::BTreeSet, + advance: impl FnOnce(Head) -> std::collections::BTreeSet, ) -> std::collections::BTreeSet { - // Applied and then appended, which is the reverse of every other - // durable store here and is the one place it is right: what goes in - // the log is the created set, and the created set is the answer the - // store under this one works out. There is nothing to append until it - // has. - // - // A log failure is loud rather than fatal. The write this commit is - // over has already been acknowledged, and refusing to remember where - // it left the repository would be refusing to serve a repository that - // exists; what it costs is a `since` cursor that is answered with the - // whole repository after the next restart. let rev = head.rev.clone(); let commit = head.commit.to_string(); let data = Some(head.data.to_string()); let prev = head.prev.as_ref().map(ToString::to_string); - let created = self.inner.advance(did, head, blocks); + let created = advance(head); if let Err(error) = self.wal.append( &Entry::RepoCommitted { did: did.to_owned(), @@ -798,6 +796,33 @@ impl CommitStore for FileCommitStore { } created } +} + +impl CommitStore for FileCommitStore { + fn head(&self, did: &str) -> Option { + self.inner.head(did) + } + + fn advance( + &self, + did: &str, + head: Head, + blocks: std::collections::BTreeSet, + ) -> std::collections::BTreeSet { + self.logged(did, head, |head| self.inner.advance(did, head, blocks)) + } + + fn advance_by( + &self, + did: &str, + head: Head, + changed: &didbot_repo::Rewritten, + blocks: &dyn Fn() -> std::collections::BTreeSet, + ) -> std::collections::BTreeSet { + self.logged(did, head, |head| { + self.inner.advance_by(did, head, changed, blocks) + }) + } fn created_since( &self, diff --git a/crates/didbot-pds/src/export.rs b/crates/didbot-pds/src/export.rs index d414692d..0022a402 100644 --- a/crates/didbot-pds/src/export.rs +++ b/crates/didbot-pds/src/export.rs @@ -11,16 +11,15 @@ //! # The repository is derived, not stored //! //! No tree, no commit and no signature is written to the write-ahead log. -//! An export reads the records back out of the store and rebuilds all three, +//! [`export_at`] reads the records back out of the store and builds all three, //! and this is the whole reason the log format did not have to change and //! its replay still works unaltered. //! -//! It is affordable because a tree is a pure function of its leaves: two -//! exports of the same records are byte-identical, so there is nothing to keep -//! in sync and nothing that can drift. What it costs is that a repository is -//! rebuilt on every export rather than patched, which is the trade -//! and the sentence to come back to when the first repository here is large -//! enough for it to matter. +//! A tree is a pure function of its leaves, so two builds of the same records +//! are byte-identical. The provisioner holds the last build of a repository in +//! memory and edits it one key at a time as writes land, which lands on the +//! same tree; a repository it holds nothing for is built here again. See +//! `Provisioner::commit_write`. //! //! # Where the commit sits is not derived //! diff --git a/crates/didbot-pds/src/history.rs b/crates/didbot-pds/src/history.rs index ff0338b9..c31a010b 100644 --- a/crates/didbot-pds/src/history.rs +++ b/crates/didbot-pds/src/history.rs @@ -12,11 +12,14 @@ //! block names each recent commit *created*. Neither is a second copy of the //! repository: the head is three short strings, and a trail entry is a handful //! of content identifiers with no bytes attached. The records are still the -//! only thing this server stores, and an export still rebuilds from them. +//! only thing this server stores, and a repository is still built from them. //! //! Beside those it holds the block names the head's repository *has*, as this //! process last saw it. That is what makes a commit's created set a -//! subtraction rather than the whole tree. It is names and no bytes as well, +//! subtraction rather than the whole tree, for a commit whose caller only has +//! the new repository; one that edited the repository at the head reports +//! what it changed instead, and the set is kept up by that change. It is names +//! and no bytes as well, //! it is one set per repository rather than one per revision, and a restart //! begins without it, because a log carries what each commit created and not //! what the repository holds — see [`MemoryCommitStore::restore`]. @@ -49,6 +52,7 @@ use std::collections::{BTreeMap, BTreeSet, VecDeque}; use std::sync::Mutex; use didbot_data::Cid; +use didbot_repo::Rewritten; /// How many revisions of block names one repository keeps behind its head. /// @@ -134,6 +138,27 @@ pub trait CommitStore: Send + Sync + std::fmt::Debug { /// already had, rather than being left unable to walk its own tree. fn advance(&self, did: &str, head: Head, blocks: BTreeSet) -> BTreeSet; + /// [`Self::advance`], for a caller that knows what the commit changed. + /// + /// `changed` is measured against the repository at the head this store + /// holds, which the caller had in hand when it made the commit — see + /// [`didbot_repo::Draft`]. Its `created` set is the answer. `blocks` + /// answers every block the new repository holds; a store asks for it only + /// when it has not kept the set for its head. + /// + /// The default hands the whole set to [`Self::advance`], which is correct + /// for any store and costs what that costs. + fn advance_by( + &self, + did: &str, + head: Head, + changed: &Rewritten, + blocks: &dyn Fn() -> BTreeSet, + ) -> BTreeSet { + let _ = changed; + self.advance(did, head, blocks()) + } + /// The blocks created after `since`, or `None` if it cannot be placed. /// /// `Some(empty)` and `None` are different answers and the difference @@ -205,7 +230,10 @@ impl MemoryCommitStore { /// The trail entry is taken as given rather than derived, because the log /// is where it was derived. What the repository *holds* is left unknown — /// a log carries what each commit created and not the whole set — so the - /// first commit after a restore is recorded as creating all of it. + /// first commit after a restore made through [`CommitStore::advance`] is + /// recorded as creating all of it. One made through + /// [`CommitStore::advance_by`] is exact, because its caller held the + /// repository at this head. pub fn restore(&self, did: &str, head: Head, created: BTreeSet) { let mut repos = self.repos(); let kept = repos.entry(did.to_owned()).or_insert_with(|| Kept { @@ -271,6 +299,37 @@ impl CommitStore for MemoryCommitStore { created } + fn advance_by( + &self, + did: &str, + head: Head, + changed: &Rewritten, + blocks: &dyn Fn() -> BTreeSet, + ) -> BTreeSet { + let mut repos = self.repos(); + let kept = repos.entry(did.to_owned()).or_insert_with(|| Kept { + head: head.clone(), + blocks: None, + trail: VecDeque::new(), + }); + match &mut kept.blocks { + Some(held) => { + for cid in &changed.removed { + held.remove(cid); + } + held.extend(changed.created.iter().cloned()); + } + None => kept.blocks = Some(blocks()), + } + kept.trail + .push_back((head.rev.clone(), changed.created.clone())); + while kept.trail.len() > TRAIL { + kept.trail.pop_front(); + } + kept.head = head; + changed.created.clone() + } + fn created_since(&self, did: &str, since: &str) -> Option> { let repos = self.repos(); let kept = repos.get(did)?; diff --git a/crates/didbot-pds/src/provision.rs b/crates/didbot-pds/src/provision.rs index eaa9e306..9ef8678a 100644 --- a/crates/didbot-pds/src/provision.rs +++ b/crates/didbot-pds/src/provision.rs @@ -770,6 +770,11 @@ pub struct RegistryStats { /// consumer has — a firehose frame, then a `getRecord` for its proof — across /// however many accounts are speaking at once, and thirty-two agents mid-turn /// is more than any deployment here has run. +/// +/// It is also what a write edits: a repository held at its head takes one +/// write in time that grows with the depth of its tree, and one that has +/// fallen out of here is built whole again first. See +/// [`Provisioner::commit_write`]. const BUILT: usize = 32; /// The repositories held built, and the order they were built in. @@ -801,8 +806,7 @@ struct Built { /// for a reason that is about the write path rather than about this map: /// nothing changes a repository's records without leaving a kept head behind /// — every record write in [`Provisioner`] goes through -/// [`Provisioner::commit_write`], which calls -/// [`CommitStore::advance`] — so the +/// [`Provisioner::commit_write`], which advances the history — so the /// moment a derived entry could be stale is the moment it stops being read, /// because [`Provisioner::repository`] takes the head branch instead. A /// repository whose account is gone is dropped outright; see @@ -824,6 +828,13 @@ impl Built { .map(|(_, repository)| Arc::clone(repository)) } + /// Takes the repository held for `did` out, with where it was built. + fn take(&mut self, did: &str) -> Option<(At, Arc)> { + let held = self.by_did.remove(did)?; + self.order.retain(|order| order != did); + Some(held) + } + /// Files `repository` under `at`, evicting the oldest past [`BUILT`]. fn insert(&mut self, did: &str, at: &At, repository: &Arc) { if self @@ -1008,13 +1019,7 @@ impl Committed { fn blocks(&self) -> Vec { let mut keep: BTreeSet = self.created.clone(); keep.extend(self.repository.covering_proof(&self.touched).into_keys()); - let held: BTreeSet = self - .repository - .block_cids() - .into_iter() - .filter(|cid| !keep.contains(cid)) - .collect(); - self.repository.to_car_without(&held) + self.repository.to_car_with(&keep) } } @@ -1960,14 +1965,13 @@ pub struct Provisioner { /// [`Provisioner::forget_repository`] beside `CommitStore::forget`. /// /// What it buys is that a read costs nothing a write has already paid - /// for: `getRepo`, `getRecord`'s proof and `getBlocks` each rebuilt the - /// whole repository and signed it, on every call, to answer from a tree - /// the write that preceded them had already built. It does not make a - /// *write* cheaper — that is `plan/repo-scale.md`, and needs a tree that - /// is state rather than a derivation. + /// for, and that a write edits the repository held at its head rather + /// than building it again: see [`Provisioner::commit_write`]. Losing an + /// entry costs a rebuild and nothing else. + /// /// Capped at [`BUILT`] repositories, oldest write evicted first. built: Mutex, - /// Whole-repository rebuilds this provisioner has paid for. + /// Whole-repository builds this provisioner has paid for. /// /// An operator's gauge and a test's assertion. A duration is not portable /// between machines and a count is, so the regression test that keeps a @@ -2574,14 +2578,15 @@ where // Read before the rewrite: afterwards the key holds the new record // and there is nothing left to say the update replaced. let previous = self.registration_cid(&did); - match self.publish_service_identity(&did, None) { - Ok(()) => self.announce_identity(&did, previous), - Err(error) => tracing::error!( + if let Err(error) = + self.announce_identity(&did, previous, || self.publish_service_identity(&did, None)) + { + tracing::error!( did = did.as_str(), %error, "could not publish the server's own registration record; \ the account still serves its DID document, but not a `bot.did.registration`" - ), + ); } tracing::info!( @@ -2647,8 +2652,9 @@ where pub fn confirm_ownership(&self, owner: &str) -> Result<(), ProvisionError> { let did = AgentDid::parse(&self.zone.service_did())?; let previous = self.registration_cid(&did); - self.publish_service_identity(&did, Some(owner))?; - self.announce_identity(&did, previous); + self.announce_identity(&did, previous, || { + self.publish_service_identity(&did, Some(owner)) + })?; tracing::info!(did = did.as_str(), owner, "server confirmed as owned"); Ok(()) } @@ -2840,7 +2846,7 @@ where Ok(()) } - /// Commits whatever the registration record currently says, without + /// Writes the registration record with `publish` and commits it, without /// announcing it. /// /// The commit half of [`Self::announce_identity`], split out so @@ -2850,15 +2856,27 @@ where /// building the repository, and building it needs the account's signing /// key. /// - /// `None` when there is no registration record to commit yet, or when - /// the commit itself failed — logged either way by the caller's absence - /// of a result, not here, so a caller that only wants to know whether it - /// succeeded is not forced to read a log line to find out. - fn commit_identity(&self, did: &AgentDid) -> Option<(serde_json::Value, Committed)> { - let record = - self.records - .get(did.as_str(), registration::COLLECTION, registration::RKEY)?; + /// The write and the commit are under one [`Self::writing`], like every + /// other record write, so no other commit can land between them — see + /// [`Self::commit_write`] on why that matters. + /// + /// An error is `publish`'s. `Ok(None)` when there is no registration + /// record to commit, or when the commit itself failed, which is logged + /// here. + fn commit_identity( + &self, + did: &AgentDid, + publish: impl FnOnce() -> Result<(), ProvisionError>, + ) -> Result, ProvisionError> { let writing = self.writing(); + publish()?; + let Some(record) = + self.records + .get(did.as_str(), registration::COLLECTION, registration::RKEY) + else { + self.forget_repository(did.as_str()); + return Ok(None); + }; let touched: BTreeSet = std::iter::once(didbot_repo::tree_key( registration::COLLECTION, registration::RKEY, @@ -2875,14 +2893,19 @@ where did = did.as_str(), "the registration record was not committed; it will be in the next commit" ); - return None; + // The record landed and no commit names it, so nothing held + // may be edited past it: the next commit builds from the + // records instead. + self.forget_repository(did.as_str()); + return Ok(None); } }; drop(writing); - Some((record, at)) + Ok(Some((record, at))) } - /// Commits and announces the registration record on the commit stream. + /// Writes the registration record with `publish`, commits it, and + /// announces it on the commit stream. An error is `publish`'s. /// /// A subscriber that only ever saw records would otherwise have to /// fetch this record to learn an agent exists at all. Not used at @@ -2890,9 +2913,14 @@ where /// paths that rewrite an existing account's registration record: /// [`Self::refresh_identity`] and the pin/freeze/unfreeze operations /// that call it. - fn announce_identity(&self, did: &AgentDid, previous: Option) { - let Some((record, at)) = self.commit_identity(did) else { - return; + fn announce_identity( + &self, + did: &AgentDid, + previous: Option, + publish: impl FnOnce() -> Result<(), ProvisionError>, + ) -> Result<(), ProvisionError> { + let Some((record, at)) = self.commit_identity(did, publish)? else { + return Ok(()); }; self.announce_write( did.as_str(), @@ -2902,6 +2930,7 @@ where &at, previous, ); + Ok(()) } /// Rewrites the registration record, loudly, without being able to fail. @@ -2912,15 +2941,15 @@ where return; }; let previous = self.registration_cid(did); - if let Err(error) = self.publish_identity(did, &account) { + if let Err(error) = + self.announce_identity(did, previous, || self.publish_identity(did, &account)) + { tracing::error!( did = did.as_str(), %error, "could not refresh the registration record; it is behind the ledger and says so" ); - return; } - self.announce_identity(did, previous); } /// The CID the registration record at `did` is named by right now. @@ -2944,9 +2973,8 @@ where /// Announces a write, with the commit it landed in. /// - /// Naming the commit means building it, which means rebuilding the whole - /// repository — so this does nothing at all when no sink is - /// attached, and a deployment that indexes nothing pays nothing. A frame + /// Does nothing at all when no sink is attached, so a deployment that + /// indexes nothing pays nothing for the frame. A frame /// that could not name its commit is not sent at all: the alternative is /// a frame carrying a record and no way to check it, which is the thing /// the frame shape exists to stop, and a consumer that misses one has the survey @@ -3243,10 +3271,13 @@ where )?) } - /// How many whole-repository rebuilds this provisioner has paid for. + /// How many whole-repository builds this provisioner has paid for. /// - /// One per write, and — since a read is answered from the repository the - /// write left behind — none per read. See the `built` field. + /// None for a write to a repository held at its head, and none for a read + /// answered from what a write left behind. One for a write or a read that + /// finds nothing held, which is a repository's first write in this + /// process or one that has fallen out of what is held. See the `built` + /// field. #[must_use] pub fn rebuilds(&self) -> u64 { self.rebuilds.load(Ordering::Relaxed) @@ -3269,6 +3300,13 @@ where /// `getLatestCommit` anyway, and the one thing that must not happen is it /// being served again from here. A derived entry has no head to /// disagree with, so there is nothing to check. + /// + /// Also refuses one built where the history no longer is. A read that + /// built a repository while a write was committing would otherwise file + /// the older one over the newer, and the next write would find nothing + /// held at its head. The check is under the same lock as the filing, and + /// a write files only after advancing the head, so the two cannot + /// interleave. fn remember(&self, did: &str, at: &At, repository: &Arc) { if let At::Head(head) = at { if &repository.root() != head { @@ -3281,10 +3319,23 @@ where return; } } - self.built + let mut built = self + .built .lock() - .unwrap_or_else(|poisoned| poisoned.into_inner()) - .insert(did, at, repository); + .unwrap_or_else(|poisoned| poisoned.into_inner()); + let now = self + .history + .head(did) + .map_or(At::Derived, |head| At::Head(head.commit)); + if at != &now { + tracing::debug!( + did, + ?at, + "not keeping a repository the history has moved past" + ); + return; + } + built.insert(did, at, repository); } /// Whether what this server would serve for `did` is what the records @@ -3292,16 +3343,14 @@ where /// /// Builds the repository from the records, ignoring anything already /// built, and compares it against what [`Registry::export_repo`] answers. - /// A repository here is derived, so the two cannot differ unless - /// something is holding a repository the records have moved past — which - /// is the one way a repository kept between builds could be wrong. + /// A write edits the repository held at its head rather than building it + /// again, so the two differ only if an edit landed somewhere a build from + /// the records would not. /// /// `false` is never expected. It is a check rather than an assertion /// because the caller is who decides what to do about it, and because the /// answer has to be obtainable rather than only true: - /// `plan/repo-scale.md` asks for a rebuild-and-compare that can be run, - /// and this is the version of it that exists while the tree is still a - /// derivation. + /// `plan/repo-scale.md` asks for a rebuild-and-compare that can be run. pub fn verify_repository(&self, did: &str) -> Result { let account = self.lookup(did)?; self.require_repository(&account)?; @@ -3479,19 +3528,43 @@ where /// before the write would name records that might still be refused, and a /// write left uncommitted is one no `swapCommit` and no `since` can see. /// - /// It costs a rebuild of the whole repository, which is the same rebuild - /// an export costs and is now paid per write rather than per read. That is - /// the trade `plan/repo-scale.md` is about, moved rather than made: this - /// server already rebuilt on every write to name the commit a firehose - /// frame carries. + /// `touched` is every tree key the write changed. When the repository is + /// held built at the current head, each of those keys is read back out + /// of the record store and applied to it, and the result is signed: see + /// [`didbot_repo::Draft`], which lands on the repository a full build from + /// the records would. When nothing is held at the head — the first write + /// in this process, or a repository that has fallen out of [`BUILT`] — + /// the whole repository is built from the records instead. + /// + /// Editing the held repository is only right if it differs from the + /// records by the keys in `touched` and nothing else. That holds because + /// every change to a repository's records is made under [`Self::writing`] + /// and committed here, naming its key, before the lock is released; a path + /// that empties a repository forgets its head, which leaves nothing held + /// at it. [`Self::verify_repository`] is the check that it still holds. fn commit_write( &self, account: &AgentAccount, touched: BTreeSet, ) -> Result { let did = account.did.as_str(); - let key = self.signing_key(account)?; let before = self.history.head(did); + // Whatever was held comes out, and is edited only if it is the + // repository at this head. Taken before anything here can fail: the + // records have already moved, so what was held is either edited onto + // them or is gone, never left behind to be read or edited later. + let held = self + .built + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .take(did) + .filter(|(at, _)| { + before + .as_ref() + .is_some_and(|head| at == &At::Head(head.commit.clone())) + }) + .map(|(_, repository)| repository); + let key = self.signing_key(account)?; // Where a reader is already told this repository is, whether that came // from a kept head or was derived from the record keys. A restart // leaves the minter at zero and only the clock would carry it past @@ -3501,22 +3574,31 @@ where rev: self.revisions.next(), prev: before.as_ref().map(|head| head.commit.clone()), }; - let repository = Arc::new(self.build(did, &key, &at)?); - // What the repository holds now. Which of it is new is the history's - // to work out, because it is the only thing here that saw the - // repository before the write. - let created = self.history.advance( - did, - crate::history::Head { - commit: repository.root(), - data: repository.commit().commit().data.clone(), - rev: at.rev.clone(), - prev: at.prev.clone(), - }, - repository.block_cids(), - ); + let (repository, changed) = match held { + Some(held) => { + let (repository, changed) = self.amend(did, held, &touched, &key, &at)?; + (Arc::new(repository), Some(changed)) + } + None => (Arc::new(self.build(did, &key, &at)?), None), + }; + let head = crate::history::Head { + commit: repository.root(), + data: repository.commit().commit().data.clone(), + rev: at.rev.clone(), + prev: at.prev.clone(), + }; + // What the commit created. An edit knows; a build leaves it to the + // history, the only thing here that saw the repository before the + // write, because the records have already moved on. + let created = match &changed { + Some(changed) => self + .history + .advance_by(did, head, changed, &|| repository.block_cids()), + None => self.history.advance(did, head, repository.block_cids()), + }; // The write has just built the repository its readers are about to - // ask for. Filed here, they do not build it again. + // ask for, and the next write will edit. Filed here, neither builds + // it again. self.remember(did, &At::Head(repository.root()), &repository); Ok(Committed { repository, @@ -3527,6 +3609,51 @@ where }) } + /// Applies the keys a write touched to the repository held at its head, + /// and signs the result at `at`. + /// + /// Each key is read from the record store, so what is signed is what the + /// store holds rather than what the caller meant to write. The held + /// repository is edited in place when nothing else holds it, and copied + /// first when a reader still does. + fn amend( + &self, + did: &str, + held: Arc, + touched: &BTreeSet, + key: &SigningKey, + at: &crate::export::Position, + ) -> Result<(didbot_repo::Repository, didbot_repo::Rewritten), ProvisionError> { + let mut draft = Arc::unwrap_or_clone(held).draft(); + for tree_key in touched { + // Every key in `touched` is `didbot_repo::tree_key`'s, and neither + // a collection nor a record key may hold a slash. + let Some((collection, rkey)) = tree_key.split_once('/') else { + continue; + }; + match self.records.get(did, collection, rkey) { + Some(record) => { + draft.insert(collection, rkey, &record)?; + } + None => { + draft.remove(collection, rkey); + } + } + } + let (repository, changed) = draft.commit(&at.rev, at.prev.clone(), key); + tracing::debug!( + did, + rev = at.rev, + keys = touched.len(), + records = repository.records(), + created = changed.created.len(), + removed = changed.removed.len(), + root = %repository.root(), + "signed a repository commit over the repository held at its head" + ); + Ok((repository, changed)) + } + /// Where a repository is, without its blocks. /// /// Answered out of the kept history when there is one, which costs a map @@ -5209,16 +5336,17 @@ where // A second real commit, over a tree that now holds both records — // its `prevData` is the profile commit's root (or, had that failed, // no commit at all yet), genuinely the tree before it. - self.publish_identity(did, account)?; - let deferred_registration = self.commit_identity(did).map(|(record, at)| DeferredWrite { - collection: registration::COLLECTION.to_owned(), - rkey: registration::RKEY.to_owned(), - record, - at, - // A repository this deployment minted; nothing has ever been - // at this key, so this is a create. - previous: None, - }); + let deferred_registration = self + .commit_identity(did, || self.publish_identity(did, account))? + .map(|(record, at)| DeferredWrite { + collection: registration::COLLECTION.to_owned(), + rkey: registration::RKEY.to_owned(), + record, + at, + // A repository this deployment minted; nothing has ever been + // at this key, so this is a create. + previous: None, + }); Ok([deferred_profile, deferred_registration] .into_iter() .flatten() diff --git a/crates/didbot-pds/tests/incremental_commit.rs b/crates/didbot-pds/tests/incremental_commit.rs new file mode 100644 index 00000000..006b4bd5 --- /dev/null +++ b/crates/didbot-pds/tests/incremental_commit.rs @@ -0,0 +1,250 @@ +//! A write edits the repository held at its head, and what it produces is +//! what a build from the records would. +//! +//! Seeded random sequences of creates, overwrites, deletes and `applyWrites` +//! batches, across collections and literal keys, against on-disk stores that +//! are closed and reopened part way through. After every step: +//! +//! - `verify_repository` builds the repository from the records and compares +//! it with what the server serves; +//! - a `getRepo` with `since` at the previous revision carries exactly the +//! blocks the full exports before and after say were created, which is what +//! the edit reported and the history logged; +//! - the firehose frame names the head and carries every created block. +//! +//! The first commit after a reopen is a build rather than an edit, and the +//! history has no block set to subtract from, so it may report more than it +//! created. That is the one step allowed to, and only in that direction. + +use std::collections::BTreeSet; +use std::path::Path; +use std::sync::Arc; + +use didbot_data::Cid; +use didbot_dns::LoopbackDns; +use didbot_identity::Zone; +use didbot_pds::{ + BatchOp, Durable, ProvisionRequest, Provisioner, RecordingRepoSink, Registry, RepoEvent, + Stance, Swap, +}; +use didbot_repo::car; +use rand::rngs::StdRng; +use rand::{Rng, SeedableRng}; + +type Pds = Provisioner>; + +const COLLECTIONS: [&str; 3] = [ + "com.example.thing", + "app.bsky.feed.post", + "app.bsky.feed.like", +]; + +fn open(dir: &Path) -> (Pds, Durable, Arc) { + let durable = Durable::open(dir, time::Duration::days(30)).expect("a data directory opens"); + let frames = Arc::new(RecordingRepoSink::new()); + let pds = Provisioner::new( + "did:web:owner.example", + Zone::delegated("localhost", "agents.localhost").expect("a test zone"), + "http://localhost:3000".to_owned(), + LoopbackDns::new(), + durable.accounts(), + ) + .with_record_store(durable.records()) + .with_blob_store(durable.blobs()) + .with_commit_history(durable.history()) + .with_repo_sink(frames.clone()); + (pds, durable, frames) +} + +fn record(collection: &str, n: u32) -> serde_json::Value { + serde_json::json!({ + "$type": collection, + "text": format!("anfractuous {n}"), + "createdAt": "2026-09-11T00:00:00.000Z", + }) +} + +/// The blocks a CAR file holds, by CID. +fn blocks(bytes: &[u8]) -> BTreeSet { + car::read(bytes) + .expect("a CAR this server wrote") + .blocks + .into_keys() + .collect() +} + +/// One random change to `did`, through the path a client takes. +fn change(pds: &Pds, did: &str, rng: &mut StdRng, held: &mut Vec<(String, String)>) { + let n = rng.random_range(0..8); + let pick = |rng: &mut StdRng, held: &[(String, String)]| { + (!held.is_empty()).then(|| held[rng.random_range(0..held.len())].clone()) + }; + match rng.random_range(0..10) { + // A create, under a minted key. + 0..=3 => { + let collection = COLLECTIONS[rng.random_range(0..COLLECTIONS.len())]; + let written = pds + .put_record( + did, + collection, + None, + record(collection, n), + &Swap::default(), + ) + .expect("a create"); + held.push((collection.to_owned(), written.rkey)); + } + // An overwrite, often with a record some other key already holds. + 4..=5 => { + if let Some((collection, rkey)) = pick(rng, held) { + pds.put_record( + did, + &collection, + Some(&rkey), + record(&collection, n), + &Swap::default(), + ) + .expect("an overwrite"); + } + } + // A delete. + 6..=7 => { + if let Some((collection, rkey)) = pick(rng, held) { + assert!(pds + .delete_record(did, &collection, &rkey, &Swap::default()) + .expect("a delete")); + held.retain(|key| key != &(collection.clone(), rkey.clone())); + } + } + // The profile, at its literal key. + 8 => { + pds.put_record( + did, + "app.bsky.actor.profile", + Some("self"), + serde_json::json!({"$type": "app.bsky.actor.profile", "description": format!("{n}")}), + &Swap::default(), + ) + .expect("a profile"); + } + // A batch: one create, and a delete of something held. + _ => { + let mut writes = vec![BatchOp::Create { + collection: COLLECTIONS[0].to_owned(), + rkey: None, + record: record(COLLECTIONS[0], n), + }]; + let gone = pick(rng, held); + if let Some((collection, rkey)) = &gone { + writes.push(BatchOp::Delete { + collection: collection.clone(), + rkey: rkey.clone(), + }); + } + let outcomes = pds + .apply_writes(did, writes, Stance::Optimistic, None) + .expect("a batch"); + if let Some(didbot_pds::BatchOutcome::Written(written)) = outcomes.first() { + held.push((COLLECTIONS[0].to_owned(), written.rkey.clone())); + } + if let Some(gone) = gone { + held.retain(|key| key != &gone); + } + } + } +} + +#[test] +fn every_write_serves_what_a_build_from_the_records_would() { + let seeds: u64 = std::env::var("DIDBOT_COMMIT_SEEDS") + .ok() + .and_then(|seeds| seeds.parse().ok()) + .unwrap_or(2); + for seed in 0..seeds { + let dir = std::env::temp_dir().join(format!( + "didbot-incremental-commit-{seed}-{}", + std::process::id() + )); + let _ = std::fs::remove_dir_all(&dir); + let (mut pds, mut durable, mut frames) = open(&dir); + let did = pds + .provision(ProvisionRequest::new("sandpiper", None)) + .expect("provisioning") + .account + .did + .as_str() + .to_owned(); + let mut rng = StdRng::seed_from_u64(seed); + let mut held = Vec::new(); + let mut reopened = false; + for step in 0..160 { + if step == 80 { + drop(pds); + drop(durable); + (pds, durable, frames) = open(&dir); + reopened = true; + } + let before_rev = pds.latest_commit(&did).expect("a head").rev; + let before = blocks(&pds.export_repo(&did, None).expect("an export").car); + let frames_before = frames.events().len(); + + change(&pds, &did, &mut rng, &mut held); + + let at = format!("seed {seed} step {step}"); + assert!( + pds.verify_repository(&did).expect("a verification"), + "{at}: what is served is not what the records build to" + ); + let head = pds.latest_commit(&did).expect("a head"); + if head.rev == before_rev { + // An overwrite or delete that found nothing to do. + continue; + } + let commit = Cid::parse(&head.head).expect("a cid"); + let after = blocks(&pds.export_repo(&did, None).expect("an export").car); + let created: BTreeSet = after.difference(&before).cloned().collect(); + assert!(created.contains(&commit), "{at}: the commit is not new"); + + let since = pds.export_repo(&did, Some(&before_rev)).expect("a diff"); + assert!( + since.delta.is_diff(), + "{at}: the previous revision was not placed" + ); + let diffed = blocks(&since.car); + if reopened { + assert!( + diffed.is_superset(&created), + "{at}: a diff after a reopen lost blocks" + ); + reopened = false; + } else { + assert_eq!( + diffed, created, + "{at}: the diff is not what the commit created" + ); + } + + let events = frames.events(); + let frame = events[frames_before..] + .iter() + .rev() + .find_map(|event| match event { + RepoEvent::Commit(commit) => Some(commit), + _ => None, + }) + .unwrap_or_else(|| panic!("{at}: the write sent no frame")); + assert_eq!( + frame.commit, head.head, + "{at}: the frame names another commit" + ); + assert_eq!(frame.since.as_deref(), Some(before_rev.as_str())); + assert!( + blocks(&frame.blocks).is_superset(&created), + "{at}: the frame lacks a block the commit created" + ); + } + drop(pds); + drop(durable); + let _ = std::fs::remove_dir_all(&dir); + } +} diff --git a/crates/didbot-pds/tests/repo_cost.rs b/crates/didbot-pds/tests/repo_cost.rs index 0609a628..837ac3a5 100644 --- a/crates/didbot-pds/tests/repo_cost.rs +++ b/crates/didbot-pds/tests/repo_cost.rs @@ -1,8 +1,8 @@ //! What a repository rebuild costs, and how often one is paid. //! -//! `plan/pds-writes.md` says a write rebuilds the whole repository and that -//! the numbers saying whether that matters have not been taken. These are -//! those numbers, and the check that keeps them from getting worse. +//! A write edits the repository held at its head, and a read answers from +//! what the write left, so a whole-repository build is paid only for a +//! repository nothing holds. These count the builds, and time them. //! //! Two tests, deliberately of two different kinds: //! @@ -100,13 +100,13 @@ fn a_read_pays_for_no_rebuild_the_write_already_paid_for() { "a read rebuilt a repository the write before it had already built" ); - // And a write is still exactly one rebuild, which is the cost - // `plan/repo-scale.md` is about and this does not touch. + // And a write to a repository held at its head is no rebuild at all: it + // edits the one held, which is `plan/repo-scale.md`'s first item. fill(&pds, &did, 3); assert_eq!( - pds.rebuilds() - at, - 3, - "a write costs something other than one rebuild" + pds.rebuilds(), + at, + "a write rebuilt a repository it could have edited" ); // The read after it answers at the new head rather than the old one. diff --git a/crates/didbot-pds/tests/write_scale.rs b/crates/didbot-pds/tests/write_scale.rs new file mode 100644 index 00000000..99897b1b --- /dev/null +++ b/crates/didbot-pds/tests/write_scale.rs @@ -0,0 +1,225 @@ +//! What a write costs on on-disk stores, as the repository and the number of +//! writing repositories grow. +//! +//! Both tests are `#[ignore]`d benchmarks that print numbers rather than +//! assert them, because a duration depends on the machine and its load: +//! +//! ```text +//! cargo test --release -p didbot-pds --test write_scale -- --ignored --nocapture +//! ``` +//! +//! Every write goes through `Registry::put_record_from`, the path a client's +//! `createRecord` takes, against a `Durable` over a real data directory and +//! with a firehose sink attached, so each write also builds the frame a relay +//! would be sent. Repositories are seeded straight into the record store and +//! then written to once untimed, which commits them; the timed writes are the +//! ones after that. +//! +//! - `DIDBOT_BENCH_SIZES`: records per repository for the latency test, +//! comma-separated. Default `100,1000,10000`. +//! - `DIDBOT_BENCH_WRITERS`: concurrent repositories for the throughput test. +//! Default `1,4,16,32,64`. +//! - `DIDBOT_BENCH_RECORDS`: records in each of those repositories. Default +//! `1000`. +//! - `DIDBOT_BENCH_WRITES`: timed writes per repository. Default `200` for +//! latency and `50` per writer for throughput. + +use std::path::PathBuf; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Barrier}; +use std::time::{Duration, Instant}; + +use didbot_dns::LoopbackDns; +use didbot_identity::Zone; +use didbot_pds::records::{Precondition, RecordStore}; +use didbot_pds::{ + Durable, ProvisionRequest, Provisioner, Registry, RepoAccount, RepoCommit, RepoEventSink, + RepoIdentity, Swap, +}; + +const COLLECTION: &str = "app.bsky.feed.post"; + +type Pds = Provisioner>; + +/// A post of about 300 bytes of JSON. +fn post(n: usize) -> serde_json::Value { + serde_json::json!({ + "$type": COLLECTION, + "text": format!( + "{n:08} sesquipedalian placeholder text for a benchmark post, long enough \ + that the record is about the size of an ordinary one, and nothing more \ + than that: no facets, no embeds, no langs, just words and a timestamp." + ), + "createdAt": "2026-09-11T00:00:00.000Z", + }) +} + +/// A firehose sink that throws frames away, having made the write build them. +#[derive(Default)] +struct Frames { + bytes: AtomicUsize, +} + +impl RepoEventSink for Frames { + fn commit(&self, event: &RepoCommit) { + self.bytes.fetch_add(event.blocks.len(), Ordering::Relaxed); + } + + fn identity(&self, _: &RepoIdentity) {} + + fn account(&self, _: &RepoAccount) {} +} + +/// A deployment over a fresh data directory, and the directory. +fn deployment(name: &str) -> (Pds, Durable, PathBuf) { + let dir = + std::env::temp_dir().join(format!("didbot-write-scale-{name}-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + let durable = Durable::open(&dir, time::Duration::days(30)).expect("a data directory opens"); + let pds = Provisioner::new( + "did:web:owner.example", + Zone::delegated("localhost", "agents.localhost").expect("a test zone"), + "http://localhost:3000".to_owned(), + LoopbackDns::new(), + durable.accounts(), + ) + .with_record_store(durable.records()) + .with_blob_store(durable.blobs()) + .with_commit_history(durable.history()) + .with_repo_sink(Arc::new(Frames::default())); + (pds, durable, dir) +} + +/// An account holding `records` posts, committed once. +fn seeded(pds: &Pds, durable: &Durable, name: &str, records: usize) -> String { + let did = pds + .provision(ProvisionRequest::new(name, None)) + .expect("provisioning") + .account + .did + .as_str() + .to_owned(); + let store = durable.records(); + for n in 0..records { + store + .put( + &did, + COLLECTION, + None, + post(n), + &Precondition::Unconditional, + ) + .expect("a seeded record"); + } + write(pds, &did, records); + did +} + +/// One write, as a client makes it. +fn write(pds: &Pds, did: &str, n: usize) -> Duration { + let started = Instant::now(); + pds.put_record_from(did, COLLECTION, None, post(n), &Swap::default(), None) + .expect("a write"); + started.elapsed() +} + +/// The `p`th percentile of some durations, which it sorts. +fn percentile(times: &mut [Duration], p: f64) -> Duration { + times.sort_unstable(); + let at = ((times.len() as f64 - 1.0) * p).round() as usize; + times.get(at).copied().unwrap_or_default() +} + +fn setting(name: &str, default: &[usize]) -> Vec { + std::env::var(name) + .ok() + .map(|value| { + value + .split(',') + .filter_map(|size| size.trim().parse().ok()) + .collect() + }) + .unwrap_or_else(|| default.to_vec()) +} + +/// Latency of one repository's writes, by the size of the repository. +#[test] +#[ignore = "a benchmark: run with --release --ignored --nocapture"] +fn write_latency_by_repository_size() { + let writes = setting("DIDBOT_BENCH_WRITES", &[200])[0]; + println!( + "{:>8} {:>10} {:>10} {:>10} {:>10}", + "records", "p50", "p90", "p99", "max" + ); + for size in setting("DIDBOT_BENCH_SIZES", &[100, 1_000, 10_000]) { + let (pds, durable, dir) = deployment(&format!("latency{size}")); + let did = seeded(&pds, &durable, &format!("bench{size}"), size); + let mut times: Vec = (0..writes) + .map(|n| write(&pds, &did, size + 1 + n)) + .collect(); + println!( + "{size:>8} {:>10?} {:>10?} {:>10?} {:>10?}", + percentile(&mut times, 0.5), + percentile(&mut times, 0.9), + percentile(&mut times, 0.99), + percentile(&mut times, 1.0), + ); + drop(pds); + drop(durable); + let _ = std::fs::remove_dir_all(&dir); + } +} + +/// Throughput of the whole deployment, by how many repositories write at +/// once, each from its own thread. +#[test] +#[ignore = "a benchmark: run with --release --ignored --nocapture"] +fn write_throughput_by_concurrent_repositories() { + let records = setting("DIDBOT_BENCH_RECORDS", &[1_000])[0]; + let writes = setting("DIDBOT_BENCH_WRITES", &[50])[0]; + println!("{records} records per repository, {writes} writes per writer"); + println!( + "{:>8} {:>10} {:>10} {:>10} {:>10}", + "writers", "writes/s", "p50", "p99", "rebuilds" + ); + for writers in setting("DIDBOT_BENCH_WRITERS", &[1, 4, 16, 32, 64]) { + let (pds, durable, dir) = deployment(&format!("throughput{writers}")); + let dids: Vec = (0..writers) + .map(|n| seeded(&pds, &durable, &format!("writer{n}"), records)) + .collect(); + let rebuilt = pds.rebuilds(); + let start = Barrier::new(writers + 1); + let (elapsed, mut times) = std::thread::scope(|scope| { + let handles: Vec<_> = dids + .iter() + .map(|did| { + let (pds, start) = (&pds, &start); + scope.spawn(move || { + start.wait(); + (0..writes) + .map(|n| write(pds, did, records + 1 + n)) + .collect::>() + }) + }) + .collect(); + start.wait(); + let started = Instant::now(); + let times: Vec = handles + .into_iter() + .flat_map(|handle| handle.join().expect("a writer")) + .collect(); + (started.elapsed(), times) + }); + let total = times.len() as f64; + println!( + "{writers:>8} {:>10.0} {:>10?} {:>10?} {:>10}", + total / elapsed.as_secs_f64(), + percentile(&mut times, 0.5), + percentile(&mut times, 0.99), + pds.rebuilds() - rebuilt, + ); + drop(pds); + drop(durable); + let _ = std::fs::remove_dir_all(&dir); + } +} diff --git a/docs/write-pipeline.md b/docs/write-pipeline.md index 6e402d66..d144eadd 100644 --- a/docs/write-pipeline.md +++ b/docs/write-pipeline.md @@ -153,6 +153,12 @@ The **only** stage that holds the store lock. Evaluation happens outside it, so a slow policy never holds the single writer, and a denied write never needs un-committing — which the write-ahead log cannot do cheaply. +The commit applies the keys the write changed to the repository held in +memory at its head, and signs the result. Its cost grows with the depth of +the tree, not the number of records. A repository with nothing held — its +first write since the process started, or one pushed out of the +held set by writes to others — is built from its records first. + The sequence number is assigned at admission, not arrival. The policy version and the `now` that judged the write are recorded with it, because decisions are not reproducible — an evaluator may have a model in the loop, and an agent