From 6cd4098597c52a9d7536cba68917f64f0ea9103d Mon Sep 17 00:00:00 2001 From: Michael Yankelev Date: Mon, 14 Sep 2026 09:33:59 +0200 Subject: [PATCH 1/2] test(engine): cover the drain's refusal of an unanchorable scope root Build a `Drain` harness in the drain's own unit module over the test kit seam fakes, so a pass's reject arms are drivable from where they live. The harness owns every value a pass borrows and seeds the snapshot cache, the gateway block store and the floors from one owner-root fixture. Cover `load_scope_root`: a root whose grant section does not decode, a root that carries none, and a name the cache holds nothing at, each against a control that loads the intact fixture and comes back with the ratchet. --- crates/engine/src/sync/drain.rs | 308 ++++++++++++++++++++++++++++++++ 1 file changed, 308 insertions(+) diff --git a/crates/engine/src/sync/drain.rs b/crates/engine/src/sync/drain.rs index 7fe3bdbda3..2f5f914603 100644 --- a/crates/engine/src/sync/drain.rs +++ b/crates/engine/src/sync/drain.rs @@ -7382,4 +7382,312 @@ mod tests { Err(Halt::UploadAttempt), )); } + + // ----------------------------------------------------------------------- + // The drain harness: one real `Drain` over the test kit's seam fakes, so + // the pass's own reject arms are drivable from where they live. + // ----------------------------------------------------------------------- + + use cipherbox_core::ipns::IpnsRecord; + use cipherbox_core::seal::{encode_envelope, set_grant_section}; + + use crate::content::DAG_ROOT_CODEC; + use crate::seams::{EndpointId, OwnerScopedFloorStore, QueueGenerationStore}; + use crate::testkit::fakes::{ + InMemoryCredentialStore, InMemoryFloorStore, InMemoryRecordStore, InMemorySnapshotCache, + InMemoryStagingStore, ScriptedHttp, VirtualScheduler, + }; + use crate::testkit::{ + OWNER_ROOT_EPOCH, OWNER_ROOT_SCOPE_SEED, OWNER_ROOT_WRITE_SCOPE_SEED, OwnerRootSpec, + SeededEntropy, block_on, gateway, owner_root_fixture, serve, + }; + + /// The login secret the harness's owner identity, enc secret and bin keys + /// derive from. + const HARNESS_SECRET: [u8; 32] = [7u8; 32]; + /// The scope and root node ids [`owner_root_fixture`] seals the vault root + /// under. + const HARNESS_SCOPE: [u8; 16] = [0u8; 16]; + const HARNESS_ROOT: NodeId = NodeId([0u8; 16]); + /// The sequence the cached root record carries. `Strictness::AtFloor` + /// admits the floor exactly, so the harness raises the floor to match. + const HARNESS_SEQUENCE: u64 = 1; + const HARNESS_TTL_NANOS: u64 = 2_000_000_000; + const HARNESS_EOL: &str = "2099-01-01T00:00:00Z"; + + type FakeDrain<'a> = Drain< + 'a, + InMemoryRecordStore, + ScriptedHttp, + InMemoryCredentialStore, + OwnerScopedFloorStore, + InMemorySnapshotCache, + QueueGenerationStore, + VirtualScheduler, + >; + + /// Every value a [`Drain`] borrows, owned in one place so a test can hand + /// out a pass and still drive the seams behind it. + struct DrainHarness { + transport: InMemoryRecordStore, + api: ApiClient, + floors: OwnerScopedFloorStore, + snapshot_cache: InMemorySnapshotCache, + staging: QueueGenerationStore, + scheduler: VirtualScheduler, + http: ScriptedHttp, + gateway: Gateway, + deadlines: DeadlinePolicy, + placement: PlacementDecision, + profile: SyncTimingProfile, + storage_policy: StoragePolicy, + content_profile: ContentProfile, + live_blocks: RefCell, + entropy: RefCell>, + base: BaseSnapshot, + held: RefCell, + blocked: RefCell>, + settings_hold: RefCell>, + pending_reclaim: Cell, + reclaim_stalls: RefCell>, + bookkeeping: RefCell, + orphan_heads: OrphanHeads, + converged_tick: Cell, + cancels: RefCell, + events: mpsc::UnboundedSender, + /// Held open: an events channel whose receiver dropped refuses sends. + _event_stream: mpsc::UnboundedReceiver, + bin_keys: BinIndexKeys, + dead_letters: RefCell, + observed_unlinks: RefCell>, + bin_index_record: RefCell>, + bin_index_hold: RefCell>, + pending_scope_exits: RefCell>, + root_name: IpnsName, + read_scope_seed: Zeroizing<[u8; 32]>, + write_scope_seed: Zeroizing<[u8; 32]>, + scope_roots: Vec, + keyless_roots: Vec, + enc_secret: X25519Secret, + owner_identity: EcdsaVerifier, + } + + impl DrainHarness { + fn drain(&self) -> FakeDrain<'_> { + Drain { + transport: &self.transport, + api: &self.api, + floors: &self.floors, + snapshot_cache: &self.snapshot_cache, + staging: &self.staging, + scheduler: &self.scheduler, + http: &self.http, + gateway: &self.gateway, + deadlines: &self.deadlines, + placement: &self.placement, + profile: &self.profile, + storage_policy: &self.storage_policy, + live_blocks: &self.live_blocks, + content_profile: &self.content_profile, + entropy: &self.entropy, + base: &self.base, + held: &self.held, + blocked: &self.blocked, + settings_hold: &self.settings_hold, + pending_reclaim: &self.pending_reclaim, + reclaim_stalls: &self.reclaim_stalls, + bookkeeping: &self.bookkeeping, + orphan_heads: &self.orphan_heads, + converged_tick: &self.converged_tick, + cancels: &self.cancels, + events: &self.events, + bin_keys: &self.bin_keys, + bin_retention_days: None, + dead_letters: &self.dead_letters, + bin_index_record: &self.bin_index_record, + bin_index_hold: &self.bin_index_hold, + established_bin_index: RefCell::new(None), + observed_unlinks: &self.observed_unlinks, + pending_scope_exits: &self.pending_scope_exits, + } + } + + fn scope(&self) -> DrainScope<'_> { + DrainScope { + source: ScopeEnd { + root: HARNESS_ROOT, + root_name: &self.root_name, + read_scope_seed: &self.read_scope_seed, + write_scope_seed: &self.write_scope_seed, + ascent_node_seed: None, + }, + destination: None, + scope_roots: &self.scope_roots, + keyless_roots: &self.keyless_roots, + charges_the_identity: false, + enc_secret: &self.enc_secret, + owner_identity: &self.owner_identity, + } + } + } + + /// The vault root this harness anchors on, as the fixture authors it. + fn harness_root_envelope() -> Envelope { + owner_root_fixture(OwnerRootSpec { + owner_identity: &EcdsaSigner::from_scalar(&HARNESS_SECRET).expect("valid scalar"), + owner_enc: &kdf::enc_subkey(&HARNESS_SECRET).public(), + scope_id: HARNESS_SCOPE, + root_id: HARNESS_ROOT.0, + children: Vec::new(), + child_scope_index: Vec::new(), + parent_node_seed: None, + owner_write_blob_epoch: Some(OWNER_ROOT_EPOCH), + write_history_link: Vec::new(), + grants: Vec::new(), + }) + .envelope + } + + /// A harness whose snapshot cache holds `cached_root` at the scope-root + /// name, with the head block that record anchors served over the gateway. + /// + /// `None` leaves the cache empty, which is the third arm of + /// [`Drain::load_scope_root`]'s refusal. + fn drain_harness(cached_root: Option) -> DrainHarness { + let write_scope_seed = Zeroizing::new(OWNER_ROOT_WRITE_SCOPE_SEED); + let root_name = derive_write_name(&write_scope_seed, &HARNESS_ROOT.0); + let enc_secret = kdf::enc_subkey(&HARNESS_SECRET); + + let floors = OwnerScopedFloorStore::new(InMemoryFloorStore::default()); + floors.bind(&enc_secret, &kdf::contact_label_seed(&HARNESS_SECRET)); + let snapshot_cache = InMemorySnapshotCache::default(); + let mut blocks = BTreeMap::new(); + if let Some(envelope) = cached_root { + let head_block = encode_envelope(&envelope).expect("the fixture encodes"); + let head_cid = encode_content_cid_str(&compute_cid(DAG_ROOT_CODEC, &head_block)); + let signer = kdf::ipns_keypair( + kdf::write_seed(&OWNER_ROOT_WRITE_SCOPE_SEED, &HARNESS_ROOT.0).as_bytes(), + ); + let record = IpnsRecord::create_v2( + &signer, + format!("/ipfs/{head_cid}").as_bytes(), + HARNESS_SEQUENCE, + HARNESS_TTL_NANOS, + HARNESS_EOL, + ) + .marshal(); + block_on(snapshot_cache.put(root_name.as_str().as_bytes(), &record)) + .expect("the record caches"); + block_on(floor::advance_on_unseal( + &floors, + &HARNESS_ROOT.0, + root_name.as_str().as_bytes(), + HARNESS_SEQUENCE, + OWNER_ROOT_EPOCH, + )) + .expect("the floors seed"); + blocks.insert(head_cid, head_block); + } + + let (events, _event_stream) = mpsc::unbounded(); + DrainHarness { + transport: InMemoryRecordStore::new(vec![EndpointId::new("fake:someguy")]), + api: ApiClient::new( + ScriptedHttp::default(), + InMemoryCredentialStore::default(), + "", + ), + floors, + snapshot_cache, + staging: QueueGenerationStore::new(InMemoryStagingStore::default()), + scheduler: VirtualScheduler::new(), + http: serve(&blocks), + gateway: gateway(), + deadlines: DeadlinePolicy::default(), + placement: Ok(Placement::Hosted), + profile: SyncTimingProfile::CI, + storage_policy: StoragePolicy::CI, + content_profile: ContentProfile::CI, + live_blocks: RefCell::new(LiveBlocks::default()), + entropy: RefCell::new(Box::new(SeededEntropy::new(42))), + base: BaseSnapshot::new(Snapshot::new(HARNESS_ROOT)), + held: RefCell::new(HeldRecords::new()), + blocked: RefCell::new(None), + settings_hold: RefCell::new(None), + pending_reclaim: Cell::new(0), + reclaim_stalls: RefCell::new(Vec::new()), + bookkeeping: RefCell::new(BookkeepingCursors::default()), + orphan_heads: OrphanHeads::default(), + converged_tick: Cell::new(false), + cancels: RefCell::new(UploadCancels::default()), + events, + _event_stream, + bin_keys: BinIndexKeys::derive(&HARNESS_SECRET), + dead_letters: RefCell::new(RetainedDeadLetters::new()), + observed_unlinks: RefCell::new(Vec::new()), + bin_index_record: RefCell::new(None), + bin_index_hold: RefCell::new(None), + pending_scope_exits: RefCell::new(BTreeSet::new()), + root_name, + read_scope_seed: Zeroizing::new(OWNER_ROOT_SCOPE_SEED), + write_scope_seed, + scope_roots: vec![HARNESS_ROOT], + keyless_roots: Vec::new(), + enc_secret, + owner_identity: EcdsaSigner::from_scalar(&HARNESS_SECRET) + .expect("valid scalar") + .verifying_key(), + } + } + + /// The control the two reject arms are read against: the intact fixture + /// reaches the grant-section read and comes back with the ratchet, so a + /// refusal below is the section and not an earlier stage of the load. + #[test] + fn an_intact_scope_root_anchors_the_pass_on_its_carried_ratchet() { + let harness = drain_harness(Some(harness_root_envelope())); + let drain = harness.drain(); + let scope = harness.scope(); + + let loaded = block_on(drain.load_scope_root(&scope.source)).expect("the root loads"); + + assert_eq!(loaded.epoch, OWNER_ROOT_EPOCH); + assert!(loaded.state.is_scope_root); + } + + /// A scope root whose grant section is absent or will not decode carries no + /// backward ratchet, and a pass anchored on it would downgrade every + /// lagging node of the scope to the uncharged epoch-lag hold instead of + /// refusing. An absent cache entry is the same refusal for the same reason: + /// there is no root to anchor on. + #[test] + fn a_scope_root_this_pass_cannot_anchor_a_ratchet_on_is_refused() { + let mut undecodable = harness_root_envelope(); + set_grant_section(&mut undecodable, vec![0xFF; 8]); + let mut absent = harness_root_envelope(); + let without_section: PreservedFields = absent + .unknown + .entries() + .iter() + .filter(|(key, _)| key != "grantSection") + .cloned() + .collect(); + absent.unknown = without_section; + + for (case, cached) in [ + ("a grant section that does not decode", Some(undecodable)), + ("no grant section at all", Some(absent)), + ("no cached record at the root name", None), + ] { + let harness = drain_harness(cached); + let drain = harness.drain(); + let scope = harness.scope(); + + assert_eq!( + block_on(drain.load_scope_root(&scope.source)).err(), + Some(Halt::Unclassified), + "{case}", + ); + } + } } From bab6fe8711a7dcef9444ac1d3a6dab6678e7b4aa Mon Sep 17 00:00:00 2001 From: Michael Yankelev Date: Mon, 14 Sep 2026 09:42:23 +0200 Subject: [PATCH 2/2] fix(engine): charge the drain's unattributed halt on its own budget `Halt::Unclassified` was uncharged, so an op that kept taking it held the strict-FIFO head for ever with no budget and no dead-letter exit. Charge it, but on a second count and a second ceiling: the class carries real availability failures, and the publish attempt budget is also what tells the already-published guard that this device reached the record plane. The op-attempts record now stores a triple per op, under a new format tag, so bytes an older build wrote read as no charges and the op is retried. A spent unattributed budget dead-letters with the staged version preserved and hands back no name, because the class spans both sides of the PUT. The record is now pruned against the pass's own retirements and dropped when it is empty, so the charges an emptied pass leaves do not park a staging row. --- crates/engine/src/sync/drain.rs | 325 ++++++++++++++++++++++++++------ 1 file changed, 268 insertions(+), 57 deletions(-) diff --git a/crates/engine/src/sync/drain.rs b/crates/engine/src/sync/drain.rs index 2f5f914603..fc29e351c5 100644 --- a/crates/engine/src/sync/drain.rs +++ b/crates/engine/src/sync/drain.rs @@ -266,6 +266,17 @@ pub(crate) struct DrainReport { } impl DrainReport { + /// Every op this pass took out of the durable queue, however it left. + fn left_the_queue(&self) -> BTreeSet { + self.published + .iter() + .chain(&self.dropped) + .chain(&self.restore_residue) + .copied() + .chain(self.dead_letters.iter().map(|(op_id, ..)| *op_id)) + .collect() + } + /// Whether the pass left the durable queue exactly as it found it. pub(crate) fn is_empty(&self) -> bool { self.published.is_empty() @@ -275,10 +286,23 @@ impl DrainReport { } } -/// Per-op drain attempt counts, decoded from [`OP_ATTEMPTS_KEY`]. +/// What one op has been charged: publish attempts, and passes that stopped on +/// a halt the valve could not attribute. +/// +/// Two counts rather than one, because two ceilings read them +/// ([`ATTEMPT_BUDGET`], [`UNATTRIBUTED_BUDGET`]) and because +/// [`Drain::create_replays_a_publish`] treats a spent attempt as proof this +/// device already tried to publish — which an unreachable provider is not. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +struct Charges { + attempts: u32, + unattributed: u32, +} + +/// Per-op drain charges, decoded from [`OP_ATTEMPTS_KEY`]. #[derive(Debug, Default)] struct Attempts { - counts: BTreeMap, + counts: BTreeMap, /// Whether this pass changed the counts — an unchanged record is not /// rewritten, so an idle tick makes no staging write. dirty: bool, @@ -290,7 +314,7 @@ impl Attempts { /// direction: an op is retried, never abandoned on unreadable bookkeeping. fn decode(stored: Option>) -> Self { let Some(bytes) = stored.filter(|bytes| { - bytes.first() == Some(&ATTEMPT_FORMAT_V1) && bytes.len() % ATTEMPT_ENTRY_LEN == 1 + bytes.first() == Some(&ATTEMPT_FORMAT_V2) && bytes.len() % ATTEMPT_ENTRY_LEN == 1 }) else { return Self::default(); }; @@ -298,10 +322,16 @@ impl Attempts { counts: bytes[1..] .chunks_exact(ATTEMPT_ENTRY_LEN) .map(|entry| { - let (id, count) = entry.split_at(8); + let (id, counts) = entry.split_at(8); + let (attempts, unattributed) = counts.split_at(4); ( OpId(u64::from_be_bytes(id.try_into().expect("8 bytes"))), - u32::from_be_bytes(count.try_into().expect("4 bytes")), + Charges { + attempts: u32::from_be_bytes(attempts.try_into().expect("4 bytes")), + unattributed: u32::from_be_bytes( + unattributed.try_into().expect("4 bytes"), + ), + }, ) }) .collect(), @@ -311,25 +341,36 @@ impl Attempts { fn encode(&self) -> Vec { let mut bytes = Vec::with_capacity(1 + self.counts.len() * ATTEMPT_ENTRY_LEN); - bytes.push(ATTEMPT_FORMAT_V1); - for (op_id, count) in &self.counts { + bytes.push(ATTEMPT_FORMAT_V2); + for (op_id, charges) in &self.counts { bytes.extend_from_slice(&op_id.0.to_be_bytes()); - bytes.extend_from_slice(&count.to_be_bytes()); + bytes.extend_from_slice(&charges.attempts.to_be_bytes()); + bytes.extend_from_slice(&charges.unattributed.to_be_bytes()); } bytes } - /// What `op_id` has been charged so far. + /// How many publish attempts `op_id` has spent. fn charged_to(&self, op_id: OpId) -> u32 { - self.counts.get(&op_id).copied().unwrap_or(0) + self.counts + .get(&op_id) + .map_or(0, |charges| charges.attempts) } - /// Charge one attempt to `op_id` and return its new count. + /// Charge one publish attempt to `op_id` and return its new count. fn charge(&mut self, op_id: OpId) -> u32 { self.dirty = true; - let count = self.counts.entry(op_id).or_default(); - *count = count.saturating_add(1); - *count + let charges = self.counts.entry(op_id).or_default(); + charges.attempts = charges.attempts.saturating_add(1); + charges.attempts + } + + /// Charge one unattributed halt to `op_id` and return its new count. + fn charge_unattributed(&mut self, op_id: OpId) -> u32 { + self.dirty = true; + let charges = self.counts.entry(op_id).or_default(); + charges.unattributed = charges.unattributed.saturating_add(1); + charges.unattributed } /// Drop every count whose op has left the queue, so the record cannot grow @@ -347,8 +388,12 @@ impl Attempts { #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum Halt { /// A reason the valve does not classify — a seam failure, an unreachable - /// record plane, a load this pass could not do. Charged nothing and retried - /// on the next tick, so an outage never abandons an op. + /// record plane, a load this pass could not do. Retried on the next tick + /// against [`UNATTRIBUTED_BUDGET`] rather than the attempt budget, so an + /// outage does not abandon an op while a halt that never clears still + /// leaves the queue. A spent budget hands back no name: the class spans + /// both sides of the PUT, and cutting a name a live record carries would + /// leave a reference outliving its referent. Unclassified, /// A node this pass must re-author is still behind the scope's epoch, and /// this pass holds no backward ratchet to open it at its own epoch @@ -438,8 +483,9 @@ fn halt_below_another_scope_root( /// /// Held apart from the vault root's own seeds: an owner holding neither vault /// seed runs no vault-root pass, and an op below a keyless scope root would -/// then take [`Halt::Unclassified`] from every pass, spend no attempt budget, -/// and hold the strict-FIFO head for ever with no dead letter (ADR 0012 D6). +/// then take [`Halt::Unclassified`] from every pass and hold the strict-FIFO +/// head on the wide outage budget, which reports a stall no outage explains +/// (ADR 0012 D6). pub(crate) fn charge_the_identity_to_one_pass(scopes: &mut [DrainScope<'_>]) { if let Some(first) = scopes.first_mut() { first.charges_the_identity = true; @@ -532,13 +578,25 @@ const CONTENT_LOST: Halt = Halt::Permanent(DeadLetterReason::ContentUnrecoverabl /// How many non-confirming publish attempts one op gets before it dead-letters. /// -/// Bounds a pathology, not a network outage — only [`Halt::Attempt`] and -/// [`Halt::UploadAttempt`] are charged. +/// Bounds a pathology, not a network outage: it is spent only by a halt the +/// valve could attribute to these bytes. const ATTEMPT_BUDGET: u32 = 5; -/// The staging key holding per-op drain attempt counts: a one-byte format tag -/// followed by `(op_id, attempts)` pairs, big-endian and fixed-width, rewritten -/// each pass over the live queue so a retired op's count leaves with it. +/// How many passes may stop on one op under [`Halt::Unclassified`] before it +/// dead-letters. +/// +/// The class carries real availability failures, so the ceiling is the outage a +/// device may sit through and still publish: at the production poll cadence it +/// is an hour of consecutive halted passes, where [`ATTEMPT_BUDGET`] would be +/// two and a half minutes. It is finite because strict FIFO means the op that +/// keeps halting holds every op behind it, and a queue with no exit is the +/// silent permanent stall the valve exists to remove. +const UNATTRIBUTED_BUDGET: u32 = 120; + +/// The staging key holding per-op drain charges: a one-byte format tag followed +/// by `(op_id, attempts, unattributed)` triples, big-endian and fixed-width, +/// rewritten each pass over the live queue so a retired op's counts leave with +/// it. /// /// It lives in the staging store for the same reason [`DRAINED_OP_MARK_PREFIX`] /// does — the counts and the op ids they name share one durability domain — and @@ -551,10 +609,11 @@ pub const OP_ATTEMPTS_KEY: &[u8] = b"cipherbox/op-attempts"; /// The attempt record's format tag. The staging store is shared with whatever /// build wrote it, so bytes that merely happen to be the right length must not /// parse as counts — a fabricated count would abandon an op early. -const ATTEMPT_FORMAT_V1: u8 = 1; +const ATTEMPT_FORMAT_V2: u8 = 2; -/// One `(op_id, attempts)` pair as [`OP_ATTEMPTS_KEY`] stores it. -const ATTEMPT_ENTRY_LEN: usize = 12; +/// One `(op_id, attempts, unattributed)` triple as [`OP_ATTEMPTS_KEY`] stores +/// it. +const ATTEMPT_ENTRY_LEN: usize = 16; /// One read of the durable queue: this identity's decoded ops, and every id the /// store holds — including other identities' and retained records', which the @@ -912,10 +971,10 @@ fn plane_seals( /// The charged form of a halt, for a refusal no retry of this pass clears. /// -/// [`Halt::Unclassified`] retries free and forever, which is right for a read -/// the next pass may win and wrong for one it will meet again unchanged. A -/// crossing that keeps taking it would hold the FIFO head with nothing reported -/// — the very failure a classified halt exists to prevent. +/// [`Halt::Unclassified`] retries on the wide outage budget, which is right for +/// a read the next pass may win and wrong for one it will meet again unchanged. +/// A crossing that keeps taking it would hold the FIFO head for an hour of +/// passes over a refusal the first pass already settled. fn charge_crossing_read(halt: Halt) -> Halt { match halt { Halt::Unclassified => Halt::UploadAttempt, @@ -1589,6 +1648,12 @@ where let _ = self .pass(scope, exits, &queued, &mut report, &mut attempts) .await; + // Pruned against this pass's own retirements, not just the queue it + // opened on: a count left behind for an op that has gone would park the + // whole record until some later pass happened to read the queue again. + let gone = report.left_the_queue(); + let live: BTreeSet = all_ids.difference(&gone).copied().collect(); + attempts.retain_live(&live); let _ = self.store_attempts(&attempts).await; let _ = self.mark_drained(scope, &queued, &report).await; (report, Some(purges)) @@ -1732,7 +1797,23 @@ where self.clear_bin_index_hold(); } match halt { - Halt::Unclassified | Halt::EpochLagged => {} + Halt::EpochLagged => {} + // Its own budget, its own count: an outage must not spend the + // attempt budget, and a spent attempt is what tells + // [`Self::create_replays_a_publish`] this device already published. + Halt::Unclassified => { + if attempts.charge_unattributed(op_id) < UNATTRIBUTED_BUDGET { + return; + } + self.abandon_keeping_its_name( + scope, + op_id, + op, + DeadLetterReason::AttemptsExhausted, + report, + ) + .await; + } // The facade undid the op against the blocks it could see when the // cancel landed. One more can confirm inside that window — the // upload the drain was already awaiting — and it would be charged @@ -1791,15 +1872,8 @@ where // upload of the version now at the name, and unpinning content a // live record names is loss where leaving rows charged is a leak. Halt::Permanent(reason @ DeadLetterReason::BaseSuperseded) => { - let Ok(preserved) = self.preserve_dead_letter(scope, op_id, reason).await else { - return; - }; - if self.dequeue_op(op_id).await.is_ok() { - self.release_if_refused(preserved, op).await; - report - .dead_letters - .push((op_id, op.target, preserved.observed(reason))); - } + self.abandon_keeping_its_name(scope, op_id, op, reason, report) + .await; } // A replayed create keeps its staged version for the same reason and // hands back the same name a spent budget does: the record standing @@ -5595,6 +5669,29 @@ where } } + /// Abandon one op with its staged version preserved and its name left + /// standing: the halt that took it gives no ground to believe no published + /// record names the op's target, and cutting a name a live record carries + /// would leave a reference outliving its referent. + async fn abandon_keeping_its_name( + &self, + scope: &DrainScope<'_>, + op_id: OpId, + op: &Op, + reason: DeadLetterReason, + report: &mut DrainReport, + ) { + let Ok(preserved) = self.preserve_dead_letter(scope, op_id, reason).await else { + return; + }; + if self.dequeue_op(op_id).await.is_ok() { + self.release_if_refused(preserved, op).await; + report + .dead_letters + .push((op_id, op.target, preserved.observed(reason))); + } + } + /// Abandon one op: retire what its publish registered, then drop it from /// the queue. async fn abandon(&self, scope: &DrainScope<'_>, op_id: OpId, op: &Op) -> Result<(), Halt> { @@ -5775,6 +5872,16 @@ where if !attempts.dirty { return Ok(()); } + // A record the pass emptied is dropped rather than rewritten as a bare + // tag: [`orphan_staging_keys`] holds this key referenced, so a tag-only + // body would park a staging row nothing ever reclaims. + if attempts.counts.is_empty() { + return self + .staging + .remove_staged_bytes(OP_ATTEMPTS_KEY) + .await + .map_err(seam); + } self.staging .put_staged_bytes(OP_ATTEMPTS_KEY, &attempts.encode()) .await @@ -5791,13 +5898,7 @@ where queued: &[(OpId, Op)], report: &DrainReport, ) -> Result<(), Halt> { - let retired: BTreeSet = report - .published - .iter() - .chain(&report.dropped) - .copied() - .chain(report.dead_letters.iter().map(|(op_id, ..)| *op_id)) - .collect(); + let retired = report.left_the_queue(); let Some(mark) = queued .iter() .map_while(|(op_id, _)| retired.contains(op_id).then_some(op_id.0)) @@ -6827,24 +6928,35 @@ mod tests { ); } - fn attempts(pairs: &[(u64, u32)]) -> Attempts { + fn attempts(charges: &[(u64, u32, u32)]) -> Attempts { let mut attempts = Attempts::default(); - for (op_id, count) in pairs { - for _ in 0..*count { + for (op_id, publishes, unattributed) in charges { + for _ in 0..*publishes { attempts.charge(OpId(*op_id)); } + for _ in 0..*unattributed { + attempts.charge_unattributed(OpId(*op_id)); + } } attempts } + fn charges(attempts: u32, unattributed: u32) -> Charges { + Charges { + attempts, + unattributed, + } + } + /// A budget only bounds a pathology if it survives the restart that a - /// half-published op is most likely to hit. + /// half-published op is most likely to hit — and the two counts are read by + /// two ceilings, so both have to survive it apart. #[test] fn the_attempt_record_survives_a_round_trip() { - let stored = attempts(&[(1, 2), (9, 1)]).encode(); + let stored = attempts(&[(1, 2, 3), (9, 0, 1)]).encode(); assert_eq!( Attempts::decode(Some(stored)).counts, - BTreeMap::from([(OpId(1), 2), (OpId(9), 1)]) + BTreeMap::from([(OpId(1), charges(2, 3)), (OpId(9), charges(0, 1))]) ); } @@ -6854,8 +6966,8 @@ mod tests { #[test] fn bytes_this_build_did_not_write_read_as_no_attempts() { let foreign_but_well_sized = { - let mut bytes = attempts(&[(1, 2)]).encode(); - bytes[0] = ATTEMPT_FORMAT_V1.wrapping_add(1); + let mut bytes = attempts(&[(1, 2, 0)]).encode(); + bytes[0] = ATTEMPT_FORMAT_V2.wrapping_add(1); bytes }; for stored in [ @@ -6964,9 +7076,9 @@ mod tests { #[test] fn a_retired_ops_count_leaves_with_it() { - let mut attempts = attempts(&[(1, 3), (2, 1)]); + let mut attempts = attempts(&[(1, 3, 0), (2, 1, 4)]); attempts.retain_live(&BTreeSet::from([OpId(2)])); - assert_eq!(attempts.counts, BTreeMap::from([(OpId(2), 1)])); + assert_eq!(attempts.counts, BTreeMap::from([(OpId(2), charges(1, 4))])); } /// A refusal the API answered with no discriminator stamped on it. @@ -7529,6 +7641,28 @@ mod tests { owner_identity: &self.owner_identity, } } + + /// Queue one op under the same owner secret the pass reads the queue + /// with, and answer the id the store assigned it. + fn queue_an_op(&self, op: &Op) -> OpId { + block_on(stage_op( + &self.staging, + RecordSeal { + owner_enc_secret: &self.enc_secret, + ephemeral_scalar: Zeroizing::new([0x5A; 32]), + }, + op, + )) + .expect("the op queues") + } + + fn queued_op_ids(&self) -> Vec { + block_on(self.staging.queued_ops()) + .expect("the queue reads") + .into_iter() + .map(|(id, _)| id) + .collect() + } } /// The vault root this harness anchors on, as the fixture authors it. @@ -7690,4 +7824,81 @@ mod tests { ); } } + + /// What a run of halted passes over one queued op left behind. + struct Halted { + report: DrainReport, + attempts: Attempts, + op_id: OpId, + still_queued: Vec, + } + + /// Queue one op, then stop `passes` passes on it under `halt`. + fn halted_op(halt: Halt, passes: u32) -> Halted { + let harness = drain_harness(Some(harness_root_envelope())); + let op = Op::rename(NodeId([9; 16]), "renamed.txt", 1, UnixMillis(0)); + let op_id = harness.queue_an_op(&op); + let drain = harness.drain(); + let scope = harness.scope(); + let mut attempts = Attempts::default(); + let mut report = DrainReport::default(); + + for _ in 0..passes { + block_on(drain.apply_valve(&scope, op_id, &op, halt, &mut attempts, &mut report)); + } + Halted { + report, + attempts, + op_id, + still_queued: harness.queued_op_ids(), + } + } + + /// An outage must not abandon an op, so a halt the valve cannot attribute + /// keeps its place at the head of the queue for a whole outage's worth of + /// passes — and then leaves, because strict FIFO means a halt that never + /// clears holds every op behind it with nothing the member can act on. + #[test] + fn an_unattributed_halt_is_charged_on_its_own_budget_and_then_dead_letters() { + let inside = halted_op(Halt::Unclassified, UNATTRIBUTED_BUDGET - 1); + assert!( + inside.report.dead_letters.is_empty(), + "still inside the budget" + ); + assert_eq!( + inside.still_queued, + vec![inside.op_id], + "the op keeps its place at the head", + ); + + let spent = halted_op(Halt::Unclassified, UNATTRIBUTED_BUDGET); + assert_eq!( + spent.report.dead_letters, + vec![( + spent.op_id, + NodeId([9; 16]), + DeadLetterReason::AttemptsExhausted + )], + ); + assert!(spent.still_queued.is_empty(), "the queue moves on"); + } + + /// The two budgets are spent apart. A publish attempt is what tells + /// [`Drain::create_replays_a_publish`] this device already reached the + /// record plane, so an outage that spent one would suppress the + /// already-published verdict a restored data directory depends on. + #[test] + fn an_unattributed_halt_leaves_the_publish_attempt_budget_untouched() { + let outage = halted_op(Halt::Unclassified, ATTEMPT_BUDGET); + assert_eq!(outage.attempts.charged_to(outage.op_id), 0); + assert!(outage.report.dead_letters.is_empty()); + + let refused = halted_op(Halt::UploadAttempt, ATTEMPT_BUDGET); + assert_eq!( + refused.attempts.charged_to(refused.op_id), + ATTEMPT_BUDGET, + "a refusal of these bytes still spends the tighter budget", + ); + assert_eq!(refused.report.dead_letters.len(), 1); + } }