From e6f3d5ffdeaee79eacee9896f0c13a0672d48706 Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Fri, 11 Sep 2026 19:13:25 -0400 Subject: [PATCH] fix(pds)!: judge the write that lands, and never let a delete skip the gate A write's subject is built from a read taken before its turn in the repository's queue, and `admit` re-reads under the store lock; the two are now compared, and a key that moved refuses the write as `InvalidSwap`. That also closes the delete shortcut: a delete whose first read found an empty key never reaches the gate, so it must not remove a record that landed while it was on its way. Change-Id: Ie719d3b53747e77c2818cf86198d812ac6039c52 --- crates/didbot-pds/src/provision.rs | 149 +++++++++-- crates/didbot-pds/src/records.rs | 10 + crates/didbot-pds/tests/judged_record.rs | 326 +++++++++++++++++++++++ crates/didbot-serve/src/error.rs | 3 +- docs/write-pipeline.md | 9 + 5 files changed, 471 insertions(+), 26 deletions(-) create mode 100644 crates/didbot-pds/tests/judged_record.rs diff --git a/crates/didbot-pds/src/provision.rs b/crates/didbot-pds/src/provision.rs index 83957c5b..f8277bdf 100644 --- a/crates/didbot-pds/src/provision.rs +++ b/crates/didbot-pds/src/provision.rs @@ -334,6 +334,26 @@ pub enum ProvisionError { /// The commit it is actually at. found: String, }, + /// The key held something else by the time the write reached the store. + /// + /// A caller's write is judged on a read taken before its turn in its + /// repository's queue, and the write itself runs against a second read + /// taken under the store lock. Committing on a verdict reached about + /// what used to be at the key would apply a policy decision to a record + /// it was never about — and a delete whose first read found nothing + /// skips the gate entirely, so what it removed would never have been + /// judged at all. Both are refused here instead. + /// + /// `InvalidSwap` on the wire, beside the two compare-and-swap + /// refusals: the caller re-reads the key and retries, which is the same + /// answer and the same remedy. + #[error("the record at {collection}/{rkey} changed after the write was judged")] + JudgedRecordMoved { + /// The collection the write named. + collection: String, + /// The record key it named. + rkey: String, + }, /// A stored record could not be put into a repository. /// /// Only reachable for a record that got past the write path's validation @@ -940,6 +960,40 @@ fn previous_cid(old: Option<&serde_json::Value>) -> Option { } } +/// Refuses a write whose judgment was about a different record. +/// +/// `judged` is what the write path read to build the subject a gate judged, +/// before the write took its turn in its repository's queue; `found` is what +/// the same key holds under the store lock, which is the read the commit is +/// built on. They disagree when something landed at the key in between — a +/// server-authored write, or a second caller write whose own judgment ran +/// alongside this one, since neither read is inside the queue. +/// +/// A delete that found nothing to delete does not reach a gate at all, so +/// for that one `judged` is `None` and this is the only thing standing +/// between an unjudged record and its removal. +fn guard_judged( + judged: Option<&serde_json::Value>, + found: Option<&serde_json::Value>, + collection: &str, + rkey: &str, +) -> Result<(), ProvisionError> { + if judged == found { + return Ok(()); + } + tracing::info!( + collection, + rkey, + judged = judged.is_some(), + found = found.is_some(), + "write refused: the key moved after the write was judged" + ); + Err(ProvisionError::JudgedRecordMoved { + collection: collection.to_owned(), + rkey: rkey.to_owned(), + }) +} + impl Committed { /// The CAR a `com.atproto.sync.subscribeRepos#commit` frame carries. /// @@ -4107,7 +4161,15 @@ where self.require_writable(&self.lookup(did)?)?; } self.check_commit_for(&account, swap.commit.as_ref())?; - let old = self.records.get(account.did.as_str(), collection, rkey); + let found = self.records.get(account.did.as_str(), collection, rkey); + // What the gate judged, checked against what is actually here. + // A delete that read an empty key never reached the gate — see + // the branch below — so this is also what keeps one from + // removing a record that landed while it was on its way. + if external { + guard_judged(old.as_ref(), found.as_ref(), collection, rkey)?; + } + let old = found; let old_refs = old .as_ref() .map(crate::records::blob_refs) @@ -4141,8 +4203,11 @@ where // delete of a key that holds nothing is a no-op the record store // already refuses to log; there is no content for a gate to // judge and routing it through the queue would only cost a turn - // for nothing. A server-authored deletion never reaches a gate, - // as a server-authored write never does. + // for nothing. It stays a no-op because `admit` refuses the + // delete outright if the key has stopped being empty — the gate + // this branch skipped never saw what would otherwise go. A + // server-authored deletion never reaches a gate, as a + // server-authored write never does. let changes = policy_diff(Some(old), None); let changes = crate::policy::borrow(&changes); let pds_did = self.zone.service_did(); @@ -4208,6 +4273,20 @@ where // write to the same key. let stored = record.clone(); + // Read once here, outside any lock, purely to build the subject a + // gate judges — `admit` reads the same key again under + // `self.writing()` for the write itself, which is the one read this + // repository's queue guarantees is consistent with the commit it + // accompanies. Read before `admit` is built so `admit` can hold the + // two side by side and refuse a write judged against neither. + // + // Judged against the tombstone when the key holds one: a create over + // a deleted record is the edit it would have been, and is handed to + // the gate as one. + let judged = external + .then(|| self.judged_record(account.did.as_str(), collection, rkey)) + .flatten(); + // Check, write, commit — under one lock, so the head a `swapCommit` // was compared against is still the head the new commit is built on. // Built once and run from two places: directly for a server-authored @@ -4234,6 +4313,19 @@ where self.require_writable(&self.lookup(did)?)?; } self.check_commit_for(&account, swap.commit.as_ref())?; + // What the gate judged, checked against what the key holds now. + // Read through `judged_record` on both sides, so a tombstone the + // gate was handed as `before` compares against the tombstone and + // not against the empty key the write itself sees. + if external { + guard_judged( + judged.as_ref(), + self.judged_record(account.did.as_str(), collection, rkey) + .as_ref(), + collection, + rkey.unwrap_or_default(), + )?; + } // What this write is about to replace, if the key it resolves to // already holds something — read before the write, because after // it the old content is gone and its blob references with it. @@ -4286,20 +4378,7 @@ where // [`Self::put_record_as`]'s module-level notes. WriteAuthor::Server { .. } => admit()?, WriteAuthor::External { client_id } => { - // Read once here, outside any lock, purely to build the - // subject a gate judges — `admit` reads the same key again - // under `self.writing()` for the write itself, which is the - // one read this repository's queue guarantees is consistent - // with the commit it accompanies. The two can disagree only - // if something outside this repository's queue changed the - // same key between the two reads — today, only a - // `WriteAuthor::Server` write to the same key, which is rare - // and already narrow: see the bypass list. - // - // Judged against the tombstone when the key holds one: a - // create over a deleted record is the edit it would have - // been, and is handed to the gate as one. - let old = self.judged_record(account.did.as_str(), collection, rkey); + let old = &judged; let changes = policy_diff(old.as_ref(), Some(&record)); let action = if old.is_some() { WriteAction::Update @@ -6583,6 +6662,19 @@ where // nothing else can move it in between. This is what makes it a // compare-and-swap: the pre-flight this replaced only ever looked. // + // What each operation's key holds, read outside any lock to build + // the subject each is judged on. `admit` reads every one of them + // again under `self.writing()` and refuses the whole batch if any + // has moved: see `guard_judged`. Read before `admit` is built so + // `admit` can hold the two side by side, and read for every + // operation before any of them applies, which is the state each + // judgment was actually made against — the operations inside a batch + // change each other's keys as they land. + let judged: Vec> = writes + .iter() + .map(|op| self.judged_record_for_op(did, op)) + .collect(); + // Built as a closure and handed to the write queue for the same // reason `put_record_as` builds one: it runs only after every // operation in the batch has been judged, never while a judgment is @@ -6616,6 +6708,19 @@ where // nowhere else. self.require_writable(&self.lookup(did)?)?; self.check_commit_for(&account, swap)?; + // Every operation against what its key held when it was judged, + // all of them before any of them applies. One moved key refuses + // the whole batch, which is the all-or-nothing rule the record + // store already applies: a batch is one commit, and half of it + // judged against the wrong record is not a commit to keep. + for (op, judged) in writes.iter().zip(&judged) { + guard_judged( + judged.as_ref(), + self.judged_record_for_op(did, op).as_ref(), + op.collection(), + op.rkey().unwrap_or_default(), + )?; + } // What each operation's key currently holds, read before any of them // apply — the same reason `put_record` reads it first: content that // is about to be replaced or removed is unreadable afterwards, and @@ -6700,17 +6805,11 @@ where // `crate::writequeue::WriteQueue::submit_batch` for why a batch takes // the line once rather than once per operation. // - // The `before` values here are read a second time inside `admit`, + // The `judged` values above are read a second time inside `admit`, // under `self.writing()`, exactly as `put_record_as` re-reads its // own: the read the commit is built on has to happen under the lock // that commit is taken with, and a judgment may not hold that lock. - // Within this repository's line the two can disagree only for a - // `WriteAuthor::Server` write to the same key — the same narrow, - // documented gap the single-record path already has. - let judged: Vec> = writes - .iter() - .map(|op| self.judged_record_for_op(did, op)) - .collect(); + // A key that moved between the two refuses the batch there. let diffs: Vec<_> = writes .iter() .zip(&judged) diff --git a/crates/didbot-pds/src/records.rs b/crates/didbot-pds/src/records.rs index b0d054f8..5c635c7e 100644 --- a/crates/didbot-pds/src/records.rs +++ b/crates/didbot-pds/src/records.rs @@ -759,6 +759,16 @@ impl BatchOp { | Self::Delete { collection, .. } => collection, } } + + /// The key this operation names, if it names one. A [`Self::Create`] + /// that leaves the key to the collection's strategy has none to name + /// until it has been written. + pub fn rkey(&self) -> Option<&str> { + match self { + Self::Create { rkey, .. } => rkey.as_deref(), + Self::Update { rkey, .. } | Self::Delete { rkey, .. } => Some(rkey), + } + } } /// What one operation in a batch produced. diff --git a/crates/didbot-pds/tests/judged_record.rs b/crates/didbot-pds/tests/judged_record.rs new file mode 100644 index 00000000..2ff9223f --- /dev/null +++ b/crates/didbot-pds/tests/judged_record.rs @@ -0,0 +1,326 @@ +//! What a write is judged on, and what it then lands on, are two reads of +//! the same key taken at different moments. +//! +//! `docs/write-pipeline.md` puts the judgment outside the store lock on +//! purpose: a slow evaluator must not hold the single writer. The cost is +//! that the subject a gate judges is built from a read taken before the +//! write's turn in its repository's queue, while the write itself runs +//! against a read taken under the lock. Anything that lands at the key in +//! between leaves the two disagreeing, and a delete that read an empty key +//! never reaches a gate at all. +//! +//! Each test here arms one key so that the *first* read of it answers +//! empty and every later read answers truthfully. That is exactly what a +//! write landing between the two reads looks like from the second one, +//! without a second thread to make it flaky. + +use std::sync::{Arc, Mutex}; + +use didbot_dns::LoopbackDns; +use didbot_identity::Zone; +use didbot_pds::records::{ + BatchOp, BatchOutcome, ListParams, MemoryRecordStore, Precondition, RecordError, RecordStats, + RecordStore, Stance, Written, +}; +use didbot_pds::{ + MemoryAccountStore, ProvisionError, ProvisionRequest, Provisioner, Registry, Swap, +}; +use serde_json::{json, Value}; + +const ZONE_HOST: &str = "agents.localhost"; +const COLLECTION: &str = "app.bsky.actor.profile"; +const RKEY: &str = "self"; + +/// A record store that answers the first read of one armed key as if it +/// were empty, and every read after it truthfully. +/// +/// Nothing else is changed: a write, a removal and a batch all reach the +/// store underneath unaltered, so the repository ends up holding whatever +/// the write path decided to put there. That is the point — the assertion +/// is on what survives, not on what the store was asked. +struct RacingStore { + inner: Arc, + /// The key whose next read answers empty, taken as it is used. + armed: Mutex>, +} + +impl RacingStore { + fn over(inner: Arc) -> Self { + Self { + inner, + armed: Mutex::new(None), + } + } + + /// The next read of this key answers empty. + fn hide_once(&self, did: &str, collection: &str, rkey: &str) { + *self.armed.lock().expect("not poisoned") = + Some((did.to_owned(), collection.to_owned(), rkey.to_owned())); + } + + fn hiding(&self, did: &str, collection: &str, rkey: &str) -> bool { + let mut armed = self.armed.lock().expect("not poisoned"); + let matches = armed + .as_ref() + .is_some_and(|(d, c, r)| d == did && c == collection && r == rkey); + if matches { + *armed = None; + } + matches + } +} + +impl RecordStore for RacingStore { + fn put( + &self, + did: &str, + collection: &str, + rkey: Option<&str>, + record: Value, + expect: &Precondition, + ) -> Result { + self.inner.put(did, collection, rkey, record, expect) + } + + fn list(&self, did: &str, collection: &str, params: &ListParams) -> Vec<(String, Value)> { + self.inner.list(did, collection, params) + } + + fn put_authored( + &self, + did: &str, + collection: &str, + record: Value, + ) -> Result { + self.inner.put_authored(did, collection, record) + } + + fn get(&self, did: &str, collection: &str, rkey: &str) -> Option { + if self.hiding(did, collection, rkey) { + return None; + } + self.inner.get(did, collection, rkey) + } + + fn tombstone(&self, did: &str, collection: &str, rkey: &str) -> Option { + self.inner.tombstone(did, collection, rkey) + } + + fn remove( + &self, + did: &str, + collection: &str, + rkey: &str, + expect: &Precondition, + ) -> Result { + self.inner.remove(did, collection, rkey, expect) + } + + fn collections(&self, did: &str) -> Vec { + self.inner.collections(did) + } + + fn snapshot(&self, did: &str) -> Vec<(String, String, Value)> { + self.inner.snapshot(did) + } + + fn remove_repo(&self, did: &str) { + self.inner.remove_repo(did); + } + + fn stats(&self) -> RecordStats { + self.inner.stats() + } + + fn apply_batch( + &self, + did: &str, + ops: &[BatchOp], + stance: Stance, + ) -> Result, RecordError> { + self.inner.apply_batch(did, ops, stance) + } +} + +struct Harness { + pds: Provisioner, + records: Arc, + did: String, +} + +fn harness(agent_id: &str) -> Harness { + let zone = Zone::delegated("localhost", ZONE_HOST).expect("the test zone is valid"); + let records = Arc::new(RacingStore::over(Arc::new(MemoryRecordStore::new()))); + let pds = Provisioner::new( + "did:web:operator.example", + zone, + "http://localhost:3000".to_string(), + LoopbackDns::new(), + MemoryAccountStore::new(), + ) + .with_record_store(records.clone()); + let did = pds + .provision(ProvisionRequest::new(agent_id, None)) + .expect("provisioning succeeds") + .account + .did + .as_str() + .to_owned(); + Harness { pds, records, did } +} + +fn unconditional() -> Swap { + Swap { + record: Precondition::Unconditional, + commit: None, + } +} + +fn profile(name: &str) -> Value { + json!({ "$type": COLLECTION, "displayName": name }) +} + +/// A delete whose read found an empty key skips the gate, the queue and the +/// write log — the whole point of that shortcut being that there is nothing +/// there to judge. When something *is* there by the time the delete takes +/// the store lock, removing it would remove a record no gate ever saw. +#[test] +fn a_delete_that_read_an_empty_key_removes_nothing_that_landed_after_it() { + let harness = harness("delete-race"); + harness + .pds + .put_record_from( + &harness.did, + COLLECTION, + Some(RKEY), + profile("quernstone"), + &unconditional(), + None, + ) + .expect("the profile lands"); + + harness.records.hide_once(&harness.did, COLLECTION, RKEY); + let refused = harness + .pds + .delete_record_from(&harness.did, COLLECTION, RKEY, &unconditional(), None) + .expect_err("the delete is refused"); + assert!( + matches!(refused, ProvisionError::JudgedRecordMoved { .. }), + "{refused:?}" + ); + assert_eq!( + harness + .pds + .get_record(&harness.did, COLLECTION, RKEY) + .expect("the repository reads"), + Some(profile("quernstone")), + "the record no gate judged is still there" + ); +} + +/// A write judged as a create, because the key read empty, must not land as +/// an overwrite of whatever is there when it takes the lock: the gate was +/// asked about a record with no `before` at all, and a rule guarding a field +/// against change was never consulted. +#[test] +fn a_write_judged_as_a_create_never_lands_as_an_overwrite() { + let harness = harness("put-race"); + harness + .pds + .put_record_from( + &harness.did, + COLLECTION, + Some(RKEY), + profile("quernstone"), + &unconditional(), + None, + ) + .expect("the profile lands"); + + harness.records.hide_once(&harness.did, COLLECTION, RKEY); + let refused = harness + .pds + .put_record_from( + &harness.did, + COLLECTION, + Some(RKEY), + profile("marlpit"), + &unconditional(), + None, + ) + .expect_err("the write is refused"); + assert!( + matches!(refused, ProvisionError::JudgedRecordMoved { .. }), + "{refused:?}" + ); + assert_eq!( + harness + .pds + .get_record(&harness.did, COLLECTION, RKEY) + .expect("the repository reads"), + Some(profile("quernstone")), + "the record the gate never saw is untouched" + ); +} + +/// One operation judged against a key that moved refuses the whole batch. +/// A batch is one commit, and half of it judged against the wrong record is +/// not a commit to keep. +#[test] +fn a_batch_operation_judged_against_a_moved_key_refuses_the_whole_batch() { + let harness = harness("batch-race"); + harness + .pds + .put_record_from( + &harness.did, + COLLECTION, + Some(RKEY), + profile("quernstone"), + &unconditional(), + None, + ) + .expect("the profile lands"); + + harness.records.hide_once(&harness.did, COLLECTION, RKEY); + let refused = harness + .pds + .apply_writes_from( + &harness.did, + vec![ + BatchOp::Update { + collection: COLLECTION.to_owned(), + rkey: RKEY.to_owned(), + record: profile("marlpit"), + }, + BatchOp::Create { + collection: "app.bsky.feed.generator".to_owned(), + rkey: Some("elsewhere".to_owned()), + record: json!({ "$type": "app.bsky.feed.generator" }), + }, + ], + Stance::default(), + None, + None, + ) + .expect_err("the batch is refused"); + assert!( + matches!(refused, ProvisionError::JudgedRecordMoved { .. }), + "{refused:?}" + ); + assert_eq!( + harness + .pds + .get_record(&harness.did, COLLECTION, RKEY) + .expect("the repository reads"), + Some(profile("quernstone")), + "the operation judged against the wrong record changed nothing" + ); + assert_eq!( + harness + .pds + .get_record(&harness.did, "app.bsky.feed.generator", "elsewhere") + .expect("the repository reads"), + None, + "and neither did the operation beside it" + ); +} diff --git a/crates/didbot-serve/src/error.rs b/crates/didbot-serve/src/error.rs index fedf6d6e..32320328 100644 --- a/crates/didbot-serve/src/error.rs +++ b/crates/didbot-serve/src/error.rs @@ -414,7 +414,8 @@ impl From<&ProvisionError> for ApiError { // here. That is the whole reason `UnsupportedSwap` had to // exist and the whole reason it no longer does. ProvisionError::Record(RecordError::SwapFailed { .. }) - | ProvisionError::SwapCommitFailed { .. } => (StatusCode::BAD_REQUEST, "InvalidSwap"), + | ProvisionError::SwapCommitFailed { .. } + | ProvisionError::JudgedRecordMoved { .. } => (StatusCode::BAD_REQUEST, "InvalidSwap"), // 403 rather than 400, and its own name: the request is well // formed and the answer is that this caller may not do it. A client told // `InvalidRecord` would rewrite its record and try again forever, diff --git a/docs/write-pipeline.md b/docs/write-pipeline.md index 573c282e..6e402d66 100644 --- a/docs/write-pipeline.md +++ b/docs/write-pipeline.md @@ -158,6 +158,15 @@ 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 may delete data an evaluation read. +Because stages 6 and 7 run outside the lock, a write's subject is built from a +read of its key taken before stage 5, and the write is applied against a +second read taken here. This stage compares the two and refuses the write as +`InvalidSwap` when they disagree, so a verdict reached about one record never +admits a different one. The same comparison is what stops a delete that found +an empty key — and therefore skipped stages 5 through 7 entirely — from +removing a record that landed while it was on its way. A caller re-reads the +key and retries, as it would after any other compare-and-swap refusal. + ## A batch is one commit, judged operation by operation `com.atproto.repo.applyWrites` carries several writes and produces exactly one -- 2.51.2