From e3007f51ffd0760bb9f58ca0cccd2565bd73f5e9 Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Fri, 11 Sep 2026 19:43:36 -0400 Subject: [PATCH] perf(pds)!: give each repository its own write lock A write holds its repository's own lock across its check, write and commit, and a share of the server-wide lock; only the blob collection pass takes that one whole. Two repositories' writes no longer wait on each other, and the per-repository lock is never held across a policy judgment. Co-Authored-By: Claude Opus 5 (1M context) Change-Id: I02d66b5835fc202a4e60b6b4b28868e51b714d44 --- crates/didbot-pds/src/provision.rs | 644 ++++++++++++++++------------- 1 file changed, 347 insertions(+), 297 deletions(-) diff --git a/crates/didbot-pds/src/provision.rs b/crates/didbot-pds/src/provision.rs index 9ef8678a..3b3f20c8 100644 --- a/crates/didbot-pds/src/provision.rs +++ b/crates/didbot-pds/src/provision.rs @@ -2,7 +2,7 @@ use std::collections::{BTreeMap, BTreeSet}; use std::sync::atomic::{AtomicU64, Ordering}; -use std::sync::{Arc, Mutex}; +use std::sync::{Arc, Mutex, RwLock}; use didbot_attest::{ Assurance, AttestError, AttestationBackend, AttestationClaim, NodeCredentialBackend, Provenance, @@ -1932,14 +1932,17 @@ pub struct Provisioner { /// minters could: a revision is compared with the revisions of the same /// repository and with nothing else, and [`Minter`] never goes backwards. revisions: Minter, - /// Held across a repository's check, its write and its commit. - /// - /// One lock for every repository rather than one each, which is the same - /// shape the record store already has under it and therefore costs no - /// concurrency this server had. What it buys is that `swapCommit` means - /// something: the head a write is checked against cannot move between the - /// check and the write, because every path that moves it is here. - writing: Mutex<()>, + /// Held by every write, shared, so two repositories' writes do not wait + /// on each other. + /// + /// What it excludes is the one pass that is about every repository at + /// once: [`Self::collect_blobs`] takes it whole, so no write is in flight + /// while it decides which blobs nothing references. + writing: RwLock<()>, + /// One lock per repository, held across its check, its write and its + /// commit. See [`Self::writing_to`], which is the only thing that takes + /// either of these two locks. + repositories: Mutex>>>, /// Judges externally-authored writes before they reach `writing`. /// /// [`NoPolicyGate`] until a deployment configures otherwise — see @@ -2111,7 +2114,8 @@ where history: Arc::new(MemoryCommitStore::new()), credentials: Arc::new(MemoryAgentTokenStore::new()), revisions: Minter::new(), - writing: Mutex::new(()), + writing: RwLock::new(()), + repositories: Mutex::new(std::collections::HashMap::new()), account_revoke_hook: Mutex::new(None), policy_gate: Arc::new(NoPolicyGate), write_log: Arc::new(NullWriteSink), @@ -2471,7 +2475,7 @@ where /// in the *concurrency* one. #[tracing::instrument(name = "collect_blobs", skip(self))] pub fn collect_blobs(&self) -> Vec { - let _writing = self.writing(); + let _writing = self.collecting(); let collected = self.blobs.collect_unreferenced(); if !collected.is_empty() { tracing::info!( @@ -2868,40 +2872,40 @@ where 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, - )) - .collect(); - let at = match self - .lookup(did.as_str()) - .and_then(|account| self.commit_write(&account, touched)) - { - Ok(at) => at, - Err(error) => { - tracing::error!( - %error, - did = did.as_str(), - "the registration record was not committed; it will be in the next commit" - ); - // 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.writing_to(did.as_str(), || { + publish()?; + let Some(record) = + self.records + .get(did.as_str(), registration::COLLECTION, registration::RKEY) + else { self.forget_repository(did.as_str()); return Ok(None); - } - }; - drop(writing); - Ok(Some((record, at))) + }; + let touched: BTreeSet = std::iter::once(didbot_repo::tree_key( + registration::COLLECTION, + registration::RKEY, + )) + .collect(); + let at = match self + .lookup(did.as_str()) + .and_then(|account| self.commit_write(&account, touched)) + { + Ok(at) => at, + Err(error) => { + tracing::error!( + %error, + did = did.as_str(), + "the registration record was not committed; it will be in the next commit" + ); + // 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); + } + }; + Ok(Some((record, at))) + }) } /// Writes the registration record with `publish`, commits it, and @@ -3385,15 +3389,56 @@ where .len() } - /// Takes the repository write lock, recovering from a poisoned mutex. + /// Runs one repository's check, write and commit, holding what that needs. + /// + /// Two locks, in this order and no other. This repository's own, so that + /// nothing else commits to it in between — that is what makes + /// `swapCommit` a compare-and-swap, and what lets [`Self::commit_write`] + /// edit the repository held at the head it has just read. And a share of + /// the server-wide `writing` lock, which only [`Self::collect_blobs`] + /// ever takes whole. + /// + /// Neither is ever held across a policy judgment. A judgment may cross a + /// process boundary and take as long as it likes, and it runs in this + /// repository's queue line — see [`crate::writequeue`] — before `write` + /// is called at all. That is also why this is not that line: an operator + /// freezing an account rewrites its registration record, and must not + /// wait on a judgment the agent is in the middle of. + /// + /// Recovers from a poisoned lock, the same posture the stores take: these + /// guard an ordering rather than a half-updated value, so a panic + /// elsewhere is no reason to refuse every later write. + fn writing_to(&self, did: &str, write: impl FnOnce() -> R) -> R { + let repository = self.repository_lock(did); + let _turn = repository + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + let _server = self + .writing + .read() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + write() + } + + /// The lock for one repository's writes, made on first use and kept. /// - /// The same posture the stores take: the lock guards an ordering rather - /// than a half-updated value, so a panic elsewhere is no reason to refuse - /// every later write. - fn writing(&self) -> std::sync::MutexGuard<'_, ()> { - self.writing + /// One mutex per repository this process has written to: the same shape, + /// and the same bounded cost, as the write queue's own lines. + fn repository_lock(&self, did: &str) -> Arc> { + self.repositories .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) + .entry(did.to_owned()) + .or_default() + .clone() + } + + /// Takes `writing` against every repository at once, for the one pass + /// that is about all of them. See [`Self::collect_blobs`]. + fn collecting(&self) -> std::sync::RwLockWriteGuard<'_, ()> { + self.writing + .write() + .unwrap_or_else(|poisoned| poisoned.into_inner()) } /// Whatever currently sits at `collection`/`rkey`, if that key names @@ -4283,51 +4328,51 @@ where let old = self.records.get(account.did.as_str(), collection, rkey); let admit = || -> Result<(bool, Option, Option), ProvisionError> { - let writing = self.writing(); - // The same re-read `put_record_as`'s own `admit` makes, and for - // the same reason: a deletion queued behind a slow judgment must - // not land on an account an operator froze while it waited. A - // server-authored deletion is exempt for the reason its write - // is: no caller, so no authority to check. - if external { - self.require_writable(&self.lookup(did)?)?; - } - self.check_commit_for(&account, swap.commit.as_ref())?; - 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) - .unwrap_or_default(); - let previous = previous_cid(old.as_ref()); - let removed = self - .records - .remove(account.did.as_str(), collection, rkey, &swap.record) - .inspect_err(|error| tracing::info!(%error, "record deletion refused"))?; - if removed && !old_refs.is_empty() { - self.blobs - .unmark_referenced(account.did.as_str(), &old_refs); - } - // Only when something went: a deletion of a record that was not - // there changed no record, and a commit over an unchanged - // repository would move the head for a caller that merely made - // sure of something. - let at = if removed { - let touched: BTreeSet = - std::iter::once(didbot_repo::tree_key(collection, rkey)).collect(); - Some(self.commit_write(&account, touched)?) - } else { - None - }; - drop(writing); - Ok((removed, at, previous)) + self.writing_to(account.did.as_str(), || { + // The same re-read `put_record_as`'s own `admit` makes, and for + // the same reason: a deletion queued behind a slow judgment must + // not land on an account an operator froze while it waited. A + // server-authored deletion is exempt for the reason its write + // is: no caller, so no authority to check. + if external { + self.require_writable(&self.lookup(did)?)?; + } + self.check_commit_for(&account, swap.commit.as_ref())?; + 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) + .unwrap_or_default(); + let previous = previous_cid(old.as_ref()); + let removed = self + .records + .remove(account.did.as_str(), collection, rkey, &swap.record) + .inspect_err(|error| tracing::info!(%error, "record deletion refused"))?; + if removed && !old_refs.is_empty() { + self.blobs + .unmark_referenced(account.did.as_str(), &old_refs); + } + // Only when something went: a deletion of a record that was not + // there changed no record, and a commit over an unchanged + // repository would move the head for a caller that merely made + // sure of something. + let at = if removed { + let touched: BTreeSet = + std::iter::once(didbot_repo::tree_key(collection, rkey)).collect(); + Some(self.commit_write(&account, touched)?) + } else { + None + }; + Ok((removed, at, previous)) + }) }; let (removed, at, previous) = if let (true, Some(old)) = (external, old.as_ref()) { @@ -4429,85 +4474,85 @@ where // for a write a gate refuses. That is what keeps a denied write from // ever needing to be un-committed: nothing here runs for it. let admit = || -> Result<(Written, Committed, Option), ProvisionError> { - let writing = self.writing(); - // `docs/write-pipeline.md`'s stage 3 again, and this is the read - // that counts: the check on the way in saw the account as it was - // when this write arrived, and an externally-authored write can - // wait a long time for its repository's turn while a slow policy - // judgment runs ahead of it. An operator freezing the account in - // that window would otherwise be told the account is frozen - // while a write already past the gate went on to commit. Re-read - // under the store lock, so what refuses here is the same state - // the commit would have landed against. A server-authored write - // is exempt for the reason it never reached a gate either: it - // has no caller whose authority there is anything to check. - if external { - 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. - // `None` when the key would be minted fresh, which can never - // collide with an existing record. See `records::blob_refs`. - let old = self.existing_record(account.did.as_str(), collection, rkey); - let old_refs = old - .as_ref() - .map(crate::records::blob_refs) - .unwrap_or_default(); - let new_refs = crate::records::blob_refs(&record); - // Before the write lands, because `mark_referenced` runs after it - // and cannot refuse — see `BlobStore::missing_blobs`. A record - // whose image this account does not hold is a permanent `404` - // nothing later repairs, so it is refused while the client still - // has the bytes. - self.require_blobs_held(account.did.as_str(), &new_refs)?; - // Named under the same lock and for the same reason: after the - // write nothing downstream of the store can tell a create from - // an update, and the firehose has to say which it was — and, for - // an update, name what was there. A key that resolves to a - // minted one is always a create; nothing holds it yet. - let previous = previous_cid(old.as_ref()); - let written = self - .records - .put( - account.did.as_str(), - collection, - rkey, - record.clone(), - &swap.record, - ) - .inspect_err(|error| { - tracing::info!(%error, "record refused"); - })?; - // Resolved only once the write has actually landed: a refused - // write must not move a reference count that its record never - // reached. - if !old_refs.is_empty() { - self.blobs - .unmark_referenced(account.did.as_str(), &old_refs); - } - if !new_refs.is_empty() { - self.blobs.mark_referenced(account.did.as_str(), &new_refs); - } - let touched: BTreeSet = - std::iter::once(didbot_repo::tree_key(collection, &written.rkey)).collect(); - let at = self.commit_write(&account, touched)?; - drop(writing); - Ok((written, at, previous)) + self.writing_to(account.did.as_str(), || { + // `docs/write-pipeline.md`'s stage 3 again, and this is the read + // that counts: the check on the way in saw the account as it was + // when this write arrived, and an externally-authored write can + // wait a long time for its repository's turn while a slow policy + // judgment runs ahead of it. An operator freezing the account in + // that window would otherwise be told the account is frozen + // while a write already past the gate went on to commit. Re-read + // under the store lock, so what refuses here is the same state + // the commit would have landed against. A server-authored write + // is exempt for the reason it never reached a gate either: it + // has no caller whose authority there is anything to check. + if external { + 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. + // `None` when the key would be minted fresh, which can never + // collide with an existing record. See `records::blob_refs`. + let old = self.existing_record(account.did.as_str(), collection, rkey); + let old_refs = old + .as_ref() + .map(crate::records::blob_refs) + .unwrap_or_default(); + let new_refs = crate::records::blob_refs(&record); + // Before the write lands, because `mark_referenced` runs after it + // and cannot refuse — see `BlobStore::missing_blobs`. A record + // whose image this account does not hold is a permanent `404` + // nothing later repairs, so it is refused while the client still + // has the bytes. + self.require_blobs_held(account.did.as_str(), &new_refs)?; + // Named under the same lock and for the same reason: after the + // write nothing downstream of the store can tell a create from + // an update, and the firehose has to say which it was — and, for + // an update, name what was there. A key that resolves to a + // minted one is always a create; nothing holds it yet. + let previous = previous_cid(old.as_ref()); + let written = self + .records + .put( + account.did.as_str(), + collection, + rkey, + record.clone(), + &swap.record, + ) + .inspect_err(|error| { + tracing::info!(%error, "record refused"); + })?; + // Resolved only once the write has actually landed: a refused + // write must not move a reference count that its record never + // reached. + if !old_refs.is_empty() { + self.blobs + .unmark_referenced(account.did.as_str(), &old_refs); + } + if !new_refs.is_empty() { + self.blobs.mark_referenced(account.did.as_str(), &new_refs); + } + let touched: BTreeSet = + std::iter::once(didbot_repo::tree_key(collection, &written.rkey)).collect(); + let at = self.commit_write(&account, touched)?; + Ok((written, at, previous)) + }) }; let (written, at, previous) = match author { @@ -6846,133 +6891,138 @@ where // operation in the batch has been judged, never while a judgment is // in flight, and not at all for a batch a gate refuses. let admit = || -> Result { - let writing = self.writing(); - // `docs/write-pipeline.md`'s stage 3 again, the same re-read - // `put_record_as` and `delete_record` make in their own `admit`, - // and for the same reason: the check on the way in saw the - // account as it was when the batch arrived, and a batch can wait - // a long time for this repository's turn while its operations are - // judged one by one. An operator freezing the account in that - // window must not find the batch committed anyway. - // - // Refusing here refuses the *whole* batch, which is the same - // all-or-nothing rule the record store already applies: a - // partially applied batch would be a commit no caller asked for. - // - // This is deliberately not the queue's freeze. A freeze a policy - // tripped on one of these operations stops judgment of everything - // after it, freezes this repository's line, and tells the gate; - // an operator's freeze arriving from outside is a different event - // that reaches the caller as `NotWritable` and touches neither. - // Conflating them would have this batch report a policy verdict - // no policy reached. - // - // The `account` binding above is deliberately kept rather than - // replaced: everything under this lock reads only `account.did`, - // which a freeze cannot change, and the one field a freeze does - // move — `state` — is read from the fresh account here and - // 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 - // its blob references have to fall with it. - let before: Vec> = writes - .iter() - .map(|op| self.existing_record_for_op(did, op)) - .collect(); - // Every blob the batch names, checked before any of it applies — - // all-or-nothing, the same rule the record store applies. See - // `require_blobs_held`. - let named: Vec = writes - .iter() - .filter_map(|op| match op { - BatchOp::Create { record, .. } | BatchOp::Update { record, .. } => { - Some(crate::records::blob_refs(record)) - } - BatchOp::Delete { .. } => None, - }) - .flatten() - .collect(); - self.require_blobs_held(did, &named)?; - // Named while the old content is still readable. A firehose op's - // `prev` is what lets a consumer run an update or a delete backwards, - // and after `apply_batch` there is nothing left to name. - let previous: Vec> = before - .iter() - .map(|old| previous_cid(old.as_ref())) - .collect(); - let outcomes = self - .records - .apply_batch(did, &writes, stance) - .inspect_err(|error| { - tracing::info!(%error, "batch refused; nothing was written"); - })?; - // Resolved only now that the batch has actually landed whole — a - // refused batch touched no record and must move no reference count. - for ((op, outcome), old) in writes.iter().zip(&outcomes).zip(&before) { - let old_refs = old - .as_ref() - .map(crate::records::blob_refs) - .unwrap_or_default(); - if !old_refs.is_empty() { - self.blobs.unmark_referenced(did, &old_refs); - } - let wrote = matches!( - (op, outcome), - ( - BatchOp::Create { .. } | BatchOp::Update { .. }, - BatchOutcome::Written(_) - ) - ); - if wrote { - let record = match op { - BatchOp::Create { record, .. } | BatchOp::Update { record, .. } => record, - BatchOp::Delete { .. } => unreachable!("guarded by `wrote` above"), - }; - let new_refs = crate::records::blob_refs(record); - if !new_refs.is_empty() { - self.blobs.mark_referenced(did, &new_refs); - } + self.writing_to(did, || { + // `docs/write-pipeline.md`'s stage 3 again, the same re-read + // `put_record_as` and `delete_record` make in their own `admit`, + // and for the same reason: the check on the way in saw the + // account as it was when the batch arrived, and a batch can wait + // a long time for this repository's turn while its operations are + // judged one by one. An operator freezing the account in that + // window must not find the batch committed anyway. + // + // Refusing here refuses the *whole* batch, which is the same + // all-or-nothing rule the record store already applies: a + // partially applied batch would be a commit no caller asked for. + // + // This is deliberately not the queue's freeze. A freeze a policy + // tripped on one of these operations stops judgment of everything + // after it, freezes this repository's line, and tells the gate; + // an operator's freeze arriving from outside is a different event + // that reaches the caller as `NotWritable` and touches neither. + // Conflating them would have this batch report a policy verdict + // no policy reached. + // + // The `account` binding above is deliberately kept rather than + // replaced: everything under this lock reads only `account.did`, + // which a freeze cannot change, and the one field a freeze does + // move — `state` — is read from the fresh account here and + // 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(), + )?; } - } - // Every key this batch touched, whichever operation touched it. A - // `Create` with no requested `rkey` only learns the key it landed on - // from its own outcome, which is why this zips them rather than - // reading `writes` alone. - let touched: BTreeSet = writes - .iter() - .zip(&outcomes) - .map(|(op, outcome)| match (op, outcome) { - ( - BatchOp::Create { collection, .. } | BatchOp::Update { collection, .. }, - BatchOutcome::Written(written), - ) => didbot_repo::tree_key(collection, &written.rkey), - (BatchOp::Delete { collection, rkey }, _) => { - didbot_repo::tree_key(collection, rkey) + // 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 + // its blob references have to fall with it. + let before: Vec> = writes + .iter() + .map(|op| self.existing_record_for_op(did, op)) + .collect(); + // Every blob the batch names, checked before any of it applies — + // all-or-nothing, the same rule the record store applies. See + // `require_blobs_held`. + let named: Vec = writes + .iter() + .filter_map(|op| match op { + BatchOp::Create { record, .. } | BatchOp::Update { record, .. } => { + Some(crate::records::blob_refs(record)) + } + BatchOp::Delete { .. } => None, + }) + .flatten() + .collect(); + self.require_blobs_held(did, &named)?; + // Named while the old content is still readable. A firehose op's + // `prev` is what lets a consumer run an update or a delete backwards, + // and after `apply_batch` there is nothing left to name. + let previous: Vec> = before + .iter() + .map(|old| previous_cid(old.as_ref())) + .collect(); + let outcomes = + self.records + .apply_batch(did, &writes, stance) + .inspect_err(|error| { + tracing::info!(%error, "batch refused; nothing was written"); + })?; + // Resolved only now that the batch has actually landed whole — a + // refused batch touched no record and must move no reference count. + for ((op, outcome), old) in writes.iter().zip(&outcomes).zip(&before) { + let old_refs = old + .as_ref() + .map(crate::records::blob_refs) + .unwrap_or_default(); + if !old_refs.is_empty() { + self.blobs.unmark_referenced(did, &old_refs); } - (BatchOp::Create { .. } | BatchOp::Update { .. }, BatchOutcome::Deleted) => { - unreachable!("a create or update never produces a delete outcome") + let wrote = matches!( + (op, outcome), + ( + BatchOp::Create { .. } | BatchOp::Update { .. }, + BatchOutcome::Written(_) + ) + ); + if wrote { + let record = match op { + BatchOp::Create { record, .. } | BatchOp::Update { record, .. } => { + record + } + BatchOp::Delete { .. } => unreachable!("guarded by `wrote` above"), + }; + let new_refs = crate::records::blob_refs(record); + if !new_refs.is_empty() { + self.blobs.mark_referenced(did, &new_refs); + } } - }) - .collect(); - let at = self.commit_write(&account, touched)?; - drop(writing); - Ok((outcomes, at, previous)) + } + // Every key this batch touched, whichever operation touched it. A + // `Create` with no requested `rkey` only learns the key it landed on + // from its own outcome, which is why this zips them rather than + // reading `writes` alone. + let touched: BTreeSet = writes + .iter() + .zip(&outcomes) + .map(|(op, outcome)| match (op, outcome) { + ( + BatchOp::Create { collection, .. } | BatchOp::Update { collection, .. }, + BatchOutcome::Written(written), + ) => didbot_repo::tree_key(collection, &written.rkey), + (BatchOp::Delete { collection, rkey }, _) => { + didbot_repo::tree_key(collection, rkey) + } + ( + BatchOp::Create { .. } | BatchOp::Update { .. }, + BatchOutcome::Deleted, + ) => { + unreachable!("a create or update never produces a delete outcome") + } + }) + .collect(); + let at = self.commit_write(&account, touched)?; + Ok((outcomes, at, previous)) + }) }; // One subject per operation, judged before the lock is taken. -- 2.51.2