From 6b6dbf36bbe7c77ee1f831c6d17f128f6b6fc3d9 Mon Sep 17 00:00:00 2001 From: Michael Yankelev Date: Tue, 15 Sep 2026 02:03:38 +0200 Subject: [PATCH 1/2] refactor: hold the drain queue head in one cell with a reason The drain held the queue head in three parallel cells, one for the account quota, one for a settings refusal and one for a bin index the pass could not read. All three carried the same head and were mutually exclusive, but nothing made the exclusion structural: each halt arm cleared the other arms by hand. One `QueueHold { op_id, node, reason }` replaces them, with the reason as `Quota`, `Settings` or `BinIndex`. One pre-pass gate dispatches on the reason and is the only place a hold lets go. The three public `SnapshotView` and `SessionStatus` fields become `queue_hold`, the three wasm-bound classes become one `QueueHold` that names its reason, and the client protocol carries one discriminated union in place of three nullable fields. The drain test harness also regains the two `OwnerRootSpec` fields that a semantic merge conflict between two merged pull requests left off its owner-root fixture, which broke the engine test build on main. --- .../file-browser/FileBrowserActions.test.tsx | 6 +- .../file-browser/QueueHoldNotice.test.tsx | 21 +- .../file-browser/QueueHoldNotice.tsx | 38 +-- .../file-browser/UploadPanel.test.tsx | 12 +- apps/web/src/engine/introspection.ts | 2 +- apps/web/src/engine/snapshotStore.test.ts | 2 +- apps/web/src/engine/snapshotStore.ts | 6 +- apps/web/src/engine/testFakes.ts | 4 +- .../src/vault/useFolderNavigation.test.tsx | 4 +- crates/engine/src/facade.rs | 78 ++---- crates/engine/src/lib.rs | 12 +- crates/engine/src/sync/drain.rs | 244 ++++++++---------- crates/engine/src/sync/mod.rs | 2 +- crates/engine/tests/write_plane.rs | 115 ++++++--- crates/fuse/tests/fuse_op_core.rs | 4 +- crates/wasm/src/lib.rs | 171 +++++------- crates/wasm/tests/boundary.rs | 115 +++++---- .../client/src/broadcastTransport.test.ts | 4 +- packages/client/src/index.ts | 3 +- packages/client/src/testkit.ts | 4 +- .../client/src/worker/commandCodec.test.ts | 80 +++--- packages/client/src/worker/commandCodec.ts | 49 ++-- packages/client/src/worker/engineWasm.ts | 21 +- packages/client/src/worker/protocol.ts | 55 ++-- 24 files changed, 502 insertions(+), 550 deletions(-) diff --git a/apps/web/src/components/file-browser/FileBrowserActions.test.tsx b/apps/web/src/components/file-browser/FileBrowserActions.test.tsx index 88f328f8a9..192def0f1c 100644 --- a/apps/web/src/components/file-browser/FileBrowserActions.test.tsx +++ b/apps/web/src/components/file-browser/FileBrowserActions.test.tsx @@ -59,9 +59,7 @@ function folderView(overrides: Partial = {}): SnapshotDescri children: [], ancestors: [], deadLetters: [], - blocked: null, - settingsHold: null, - binIndexHold: null, + queueHold: null, retainedRecords: 0, staleness: 'fresh', ...overrides, @@ -1364,7 +1362,7 @@ describe('the queue overlay', () => { engine, folderView({ children: [file(NOTE, 'notes.txt')], - binIndexHold: { opId: 8n, node: NOTE, check: 'timed-out' }, + queueHold: { reason: 'bin-index', opId: 8n, node: NOTE, check: 'timed-out' }, }) ); diff --git a/apps/web/src/components/file-browser/QueueHoldNotice.test.tsx b/apps/web/src/components/file-browser/QueueHoldNotice.test.tsx index 795aafcbfa..796796242e 100644 --- a/apps/web/src/components/file-browser/QueueHoldNotice.test.tsx +++ b/apps/web/src/components/file-browser/QueueHoldNotice.test.tsx @@ -21,7 +21,7 @@ describe('the queue hold notice', () => { render( ); @@ -34,7 +34,9 @@ describe('the queue hold notice', () => { it('names why the bin index did not resolve, and clears when the hold clears', () => { const { rerender } = render( ); expect(screen.getByTestId('queue-hold-notice').textContent).toContain( @@ -49,7 +51,12 @@ describe('the queue hold notice', () => { render( ); @@ -59,17 +66,15 @@ describe('the queue hold notice', () => { expect(notice.textContent).not.toContain('child-0'); }); - it('reports both holds at once', () => { + it('leaves the over-quota hold to the upload panel that renders its figure', () => { render( ); - expect(screen.getByTestId('queue-hold-notice').textContent).toContain('2 changes are waiting'); - expect(screen.getAllByRole('listitem')).toHaveLength(2); + expect(screen.queryByTestId('queue-hold-notice')).toBeNull(); }); }); diff --git a/apps/web/src/components/file-browser/QueueHoldNotice.tsx b/apps/web/src/components/file-browser/QueueHoldNotice.tsx index 9143a0c80f..54dd89895c 100644 --- a/apps/web/src/components/file-browser/QueueHoldNotice.tsx +++ b/apps/web/src/components/file-browser/QueueHoldNotice.tsx @@ -24,38 +24,24 @@ const BIN_INDEX_CAUSES: Record = { }; /** - * The two held queue heads the engine reports beside the over-quota hold: the - * member's own settings refused the head, or the owner's bin index did not - * resolve for it. Both clear, so the notice follows the snapshot and goes when - * the hold does. + * The held queue head, when the member's own settings refused it or the owner's + * bin index did not resolve for it. The over-quota hold is the upload panel's, + * which renders the figure it carries. A hold clears, so the notice follows the + * snapshot and goes when the hold does. */ export function QueueHoldNotice({ view }: { view: SnapshotDescriptor | null }) { - const holds: { key: string; text: string }[] = []; - if (view?.settingsHold != null) { - const { node, check } = view.settingsHold; - holds.push({ - key: 'settings', - text: `${held(view, node)} waits on your settings: ${SETTINGS_CAUSES[check]}.`, - }); - } - if (view?.binIndexHold != null) { - const { node, check } = view.binIndexHold; - holds.push({ - key: 'bin-index', - text: `${held(view, node)} waits on your bin: ${BIN_INDEX_CAUSES[check]}.`, - }); - } - if (holds.length === 0) return null; + const hold = view?.queueHold ?? null; + if (view == null || hold === null || hold.reason === 'quota') return null; + const text = + hold.reason === 'settings' + ? `${held(view, hold.node)} waits on your settings: ${SETTINGS_CAUSES[hold.check]}.` + : `${held(view, hold.node)} waits on your bin: ${BIN_INDEX_CAUSES[hold.check]}.`; return (
-

- {`[!] ${holds.length === 1 ? 'a change is' : `${holds.length} changes are`} waiting`} -

+

[!] a change is waiting

    - {holds.map((hold) => ( -
  • {hold.text}
  • - ))} +
  • {text}
); diff --git a/apps/web/src/components/file-browser/UploadPanel.test.tsx b/apps/web/src/components/file-browser/UploadPanel.test.tsx index 9a1d14fa15..d1f4f2c2d9 100644 --- a/apps/web/src/components/file-browser/UploadPanel.test.tsx +++ b/apps/web/src/components/file-browser/UploadPanel.test.tsx @@ -126,7 +126,7 @@ describe('the upload panel', () => { await act(async () => { engine.publish({ ...view(), - blocked: { opId: 1n, node: ROOT_ID, neededBytes: 900n }, + queueHold: { reason: 'quota', opId: 1n, node: ROOT_ID, neededBytes: 900n }, }); }); await waitFor(() => @@ -146,14 +146,20 @@ describe('the upload panel', () => { draw(engine.client, FOLDER); await queueOne(engine); await act(async () => { - engine.publish({ ...view(), blocked: { opId: 1n, node: ROOT_ID, neededBytes: 900n } }); + engine.publish({ + ...view(), + queueHold: { reason: 'quota', opId: 1n, node: ROOT_ID, neededBytes: 900n }, + }); }); await waitFor(() => expect(screen.getByTestId('upload-row-hold')).toBeTruthy()); // A hold on another session's op charges the same budget but is not this // row's business. await act(async () => { - engine.publish({ ...view(), blocked: { opId: 99n, node: ROOT_ID, neededBytes: 900n } }); + engine.publish({ + ...view(), + queueHold: { reason: 'quota', opId: 99n, node: ROOT_ID, neededBytes: 900n }, + }); }); await waitFor(() => expect(screen.queryByTestId('upload-row-hold')).toBeNull()); diff --git a/apps/web/src/engine/introspection.ts b/apps/web/src/engine/introspection.ts index d6531b0f10..358688e645 100644 --- a/apps/web/src/engine/introspection.ts +++ b/apps/web/src/engine/introspection.ts @@ -323,7 +323,7 @@ async function digest(bytes: Uint8Array): Promise { function settled(view: SnapshotDescriptor): boolean { return ( view.staleness === 'fresh' && - view.blocked === null && + view.queueHold?.reason !== 'quota' && view.children.every((child) => child.pending === 'none') ); } diff --git a/apps/web/src/engine/snapshotStore.test.ts b/apps/web/src/engine/snapshotStore.test.ts index 96d2556973..2e7ab6c25b 100644 --- a/apps/web/src/engine/snapshotStore.test.ts +++ b/apps/web/src/engine/snapshotStore.test.ts @@ -398,7 +398,7 @@ describe('failure classification', () => { describe("the drain's over-budget hold", () => { const HELD = 7n; const held = (opId: bigint, neededBytes: bigint) => ({ - view: { ...view(), blocked: { opId, node: ROOT_ID, neededBytes } }, + view: { ...view(), queueHold: { reason: 'quota' as const, opId, node: ROOT_ID, neededBytes } }, error: null, }); diff --git a/apps/web/src/engine/snapshotStore.ts b/apps/web/src/engine/snapshotStore.ts index c7677b8a68..89bac59fca 100644 --- a/apps/web/src/engine/snapshotStore.ts +++ b/apps/web/src/engine/snapshotStore.ts @@ -45,9 +45,9 @@ export function isRecoverable(error: SnapshotError): boolean { * it a snapshot field rather than an event. */ export function heldBytes(state: SnapshotState, opId: bigint | null): bigint | null { - const blocked = state.view?.blocked; - if (blocked == null || opId === null || blocked.opId !== opId) return null; - return blocked.neededBytes; + const hold = state.view?.queueHold; + if (hold == null || hold.reason !== 'quota' || opId === null || hold.opId !== opId) return null; + return hold.neededBytes; } /** Durable queue entries this session cannot read but whose bytes it is charged for. */ diff --git a/apps/web/src/engine/testFakes.ts b/apps/web/src/engine/testFakes.ts index eabca40f7a..15c11fe9b1 100644 --- a/apps/web/src/engine/testFakes.ts +++ b/apps/web/src/engine/testFakes.ts @@ -37,9 +37,7 @@ export function view( })), ancestors: [], deadLetters: [], - blocked: null, - settingsHold: null, - binIndexHold: null, + queueHold: null, retainedRecords: 0, staleness, }; diff --git a/apps/web/src/vault/useFolderNavigation.test.tsx b/apps/web/src/vault/useFolderNavigation.test.tsx index d3c2d6fa52..74de6ce3ab 100644 --- a/apps/web/src/vault/useFolderNavigation.test.tsx +++ b/apps/web/src/vault/useFolderNavigation.test.tsx @@ -20,9 +20,7 @@ function folderView(overrides: Partial = {}): SnapshotDescri children: [], ancestors: [], deadLetters: [], - blocked: null, - settingsHold: null, - binIndexHold: null, + queueHold: null, retainedRecords: 0, staleness: 'fresh', ...overrides, diff --git a/crates/engine/src/facade.rs b/crates/engine/src/facade.rs index b34a18a971..3ef79fb021 100644 --- a/crates/engine/src/facade.rs +++ b/crates/engine/src/facade.rs @@ -142,7 +142,7 @@ use crate::sync::render::{BaseSnapshot, RenderKey, RenderMemo}; use crate::sync::scope_exit_debt::SCOPE_EXIT_DEBT_PREFIX; use cipherbox_core::hex::lower as hex_lower; -pub use crate::sync::drain::{BinIndexHold, BlockedOp, SettingsHold}; +pub use crate::sync::drain::{QueueHold, QueueHoldReason}; pub use crate::sync::rebase::DeadLetterReason; use crate::sync::record::{RecordReader, RecordSeal}; pub use crate::sync::refresh::ForcedPass; @@ -710,18 +710,11 @@ pub struct DeadLetter { pub struct SessionStatus { /// Every retained dead-lettered op, with its reason. pub dead_letters: Vec, - /// The over-quota hold, if the drain has one. Read rather than evented: + /// The held queue head, if the drain has one. Read rather than evented: /// this is a state that *clears*, and a lost "resumed" would strand a host - /// on a blockage that is gone. - pub blocked: Option, - /// The settings-refused hold, if the drain has one. Read for the same - /// reason as `blocked`, and it names the rule so a host can tell the member - /// which part of their own provider config to fix. - pub settings_hold: Option, - /// The bin-index-refused hold, if the drain has one. Read for the same - /// reason as `blocked`, and it names the reason so a withheld record is a - /// cause the member can see rather than a silently stalled queue. - pub bin_index_hold: Option, + /// on a hold that is gone. It names its reason, so a stall is a cause the + /// member can see rather than a silent queue. + pub queue_hold: Option, /// How many durable queue entries this session holds but cannot read /// (CONTEXT.md "Retained record"). Deliberately unattributed — it says the /// device is not empty, never whose work it holds — and it exists so an @@ -751,12 +744,8 @@ pub struct SnapshotView { pub ancestors: Vec, /// See [`SessionStatus::dead_letters`]. pub dead_letters: Vec, - /// See [`SessionStatus::blocked`]. - pub blocked: Option, - /// See [`SessionStatus::settings_hold`]. - pub settings_hold: Option, - /// See [`SessionStatus::bin_index_hold`]. - pub bin_index_hold: Option, + /// See [`SessionStatus::queue_hold`]. + pub queue_hold: Option, /// See [`SessionStatus::retained_records`]. pub retained_records: usize, /// See [`SessionStatus::staleness`]. @@ -772,9 +761,7 @@ impl fmt::Debug for SnapshotView { .field("children", &self.children) .field("ancestors", &self.ancestors) .field("dead_letters", &self.dead_letters) - .field("blocked", &self.blocked) - .field("settings_hold", &self.settings_hold) - .field("bin_index_hold", &self.bin_index_hold) + .field("queue_hold", &self.queue_hold) .field("retained_records", &self.retained_records) .field("staleness", &self.staleness) .finish() @@ -4613,17 +4600,10 @@ pub struct Engine { /// Memo of the durable queue scan every read renders through /// ([`scan_queue`](Self::scan_queue)). queue_scan: Rc>, - /// The drain's over-quota hold, written by the drain tick and read by + /// The drain's held queue head, written by the drain tick and read by /// [`snapshot`](Self::snapshot). In-memory: a restart re-derives it from the - /// next drain attempt's own 413 rather than trusting a stale verdict. - blocked: Rc>>, - /// The drain's settings-refused hold, on the same in-memory terms as - /// [`blocked`](Self::blocked): a restart re-derives it from the next drain - /// attempt's own verdict. - settings_hold: Rc>>, - /// The drain's bin-index-refused hold, on the same in-memory terms as - /// [`blocked`](Self::blocked). - bin_index_hold: Rc>>, + /// next drain attempt's own verdict rather than trusting a stale one. + queue_hold: Rc>>, /// Pinned bytes a published prune still owes the registry, written by the /// drain tick and read by [`pending_reclaim_bytes`](Self::pending_reclaim_bytes). /// In-memory: the durable record is the retire ledger, which every pass re-reads. @@ -4794,9 +4774,7 @@ impl Engine { focus_hinted: Cell::new(None), dead_letters: Rc::new(RefCell::new(BTreeMap::new())), queue_scan: Rc::new(RefCell::new(QueueScanMemo::default())), - blocked: Rc::new(RefCell::new(None)), - settings_hold: Rc::new(RefCell::new(None)), - bin_index_hold: Rc::new(RefCell::new(None)), + queue_hold: Rc::new(RefCell::new(None)), pending_reclaim: Rc::new(Cell::new(0)), reclaim_stalls: Rc::new(RefCell::new(Vec::new())), bookkeeping: Rc::new(RefCell::new(BookkeepingCursors::default())), @@ -5873,9 +5851,7 @@ where { let entropy = self.entropy.clone(); let scope_write_seeds = self.scope_write_seeds.clone(); let dead_letters = self.dead_letters.clone(); - let blocked = self.blocked.clone(); - let settings_hold = self.settings_hold.clone(); - let bin_index_hold = self.bin_index_hold.clone(); + let queue_hold = self.queue_hold.clone(); let pending_reclaim = self.pending_reclaim.clone(); let reclaim_stalls = self.reclaim_stalls.clone(); let bookkeeping = self.bookkeeping.clone(); @@ -6507,8 +6483,7 @@ where { entropy: &entropy, base: &base, held: &held, - blocked: &blocked, - settings_hold: &settings_hold, + hold: &queue_hold, pending_reclaim: &pending_reclaim, reclaim_stalls: &reclaim_stalls, bookkeeping: &bookkeeping, @@ -6520,7 +6495,6 @@ where { bin_retention_days: owner_bin_retention_days(&tick_settings), dead_letters: &dead_letters, bin_index_record: &bin_index_record, - bin_index_hold: &bin_index_hold, established_bin_index: RefCell::new(None), observed_unlinks: &observed_unlinks, pending_scope_exits: &pending_scope_exits, @@ -9470,9 +9444,7 @@ where { let retained_records = self.scan_queue().await?.retained; Ok(SessionStatus { dead_letters: self.retained_dead_letters(), - blocked: *self.blocked.borrow(), - settings_hold: *self.settings_hold.borrow(), - bin_index_hold: *self.bin_index_hold.borrow(), + queue_hold: *self.queue_hold.borrow(), retained_records, staleness: self.staleness_now(), }) @@ -9563,9 +9535,7 @@ where { children, ancestors, dead_letters: self.retained_dead_letters(), - blocked: *self.blocked.borrow(), - settings_hold: *self.settings_hold.borrow(), - bin_index_hold: *self.bin_index_hold.borrow(), + queue_hold: *self.queue_hold.borrow(), retained_records: scan.retained, staleness: self.staleness_now(), }) @@ -11466,9 +11436,7 @@ mod tests { name: FOLDER.to_string(), }], dead_letters: Vec::new(), - blocked: None, - settings_hold: None, - bin_index_hold: None, + queue_hold: None, retained_records: 0, staleness: Staleness::Fresh, }; @@ -13597,20 +13565,20 @@ mod tests { fn a_settings_refused_hold_reaches_both_read_surfaces() { let (engine, _events) = started(); let root = engine.root(); - let hold = SettingsHold { + let hold = QueueHold { op_id: OpId(1), node: root, - refusal: crate::settings::SettingsRefusal::Byo( + reason: QueueHoldReason::Settings(crate::settings::SettingsRefusal::Byo( crate::content::ProviderError::InsecureTransport, - ), + )), }; - *engine.settings_hold.borrow_mut() = Some(hold); + *engine.queue_hold.borrow_mut() = Some(hold); assert_eq!( - block_on(engine.snapshot(root)).unwrap().settings_hold, + block_on(engine.snapshot(root)).unwrap().queue_hold, Some(hold) ); - assert_eq!(block_on(engine.status()).unwrap().settings_hold, Some(hold)); + assert_eq!(block_on(engine.status()).unwrap().queue_hold, Some(hold)); } #[test] diff --git a/crates/engine/src/lib.rs b/crates/engine/src/lib.rs index b99279db88..7db1757c38 100644 --- a/crates/engine/src/lib.rs +++ b/crates/engine/src/lib.rs @@ -126,13 +126,13 @@ pub use settings::{ }; pub use storage_policy::{Headroom, StoragePlatform, StoragePolicy}; pub use sync::{ - AppliedOp, BlockedOp, Connectivity, DeadLetterReason, DropReason, FocusTarget, FocusWindow, + AppliedOp, Connectivity, DeadLetterReason, DropReason, FocusTarget, FocusWindow, HeadReconciliation, Link, NewNode, NodeMeta, Op, OpKind, OpRecordError, OpResolution, - PointerError, PointerFetch, PointerRecord, RecordClass, RecordReader, RecordSeal, Repair, - Replaced, ReplayReport, ScopeCrossing, SessionRole, SettingsHold, Snapshot, StagedContent, - TickCause, TickControl, VaultPointerAdoption, apply_overlay, apply_repairs, classify, - decode_queue, encode_op_record, focus_set, observed_repair, rebase_one, reconcile_head, - record_content_root_cid, replay, resolve_vault_pointer, stage_op, + PointerError, PointerFetch, PointerRecord, QueueHold, QueueHoldReason, RecordClass, + RecordReader, RecordSeal, Repair, Replaced, ReplayReport, ScopeCrossing, SessionRole, Snapshot, + StagedContent, TickCause, TickControl, VaultPointerAdoption, apply_overlay, apply_repairs, + classify, decode_queue, encode_op_record, focus_set, observed_repair, rebase_one, + reconcile_head, record_content_root_cid, replay, resolve_vault_pointer, stage_op, }; /// Placeholder identity item; kept for the sibling crate stubs' dependency diff --git a/crates/engine/src/sync/drain.rs b/crates/engine/src/sync/drain.rs index 30ebd06b53..eee74f2286 100644 --- a/crates/engine/src/sync/drain.rs +++ b/crates/engine/src/sync/drain.rs @@ -194,7 +194,7 @@ fn names_this_scope(end: &ScopeEnd<'_>, child: &ChildRef) -> bool { /// index cannot hold the queue head for good. A plane this pass could not read /// is availability and waits uncharged — but it waits as a *reported* hold, so /// a party who withholds one record does not stall the queue in silence -/// ([`BinIndexHold`]). +/// ([`QueueHoldReason::BinIndex`]). /// /// A stranded mint is neither. The hold's only exit is the record resolving, /// and on a single-device account nothing is left to publish it — so the op @@ -428,7 +428,7 @@ enum Halt { Permanent(DeadLetterReason), /// The member's own settings were refused before any request was built. /// Not a failure of the op — it holds the head and its staging reservation - /// until those settings change ([`SettingsHold`]). + /// until those settings change ([`QueueHoldReason::Settings`]). HeldBySettings(SettingsRefusal), /// Over the account quota. Not a failure of the op — it holds the head and /// its staging reservation until a quota probe reports room. @@ -440,7 +440,7 @@ enum Halt { /// The bin index plane did not establish the current index, and the reason /// is availability rather than a verdict on bytes it served. Not a failure /// of the op — it holds the head and its staging reservation until the - /// record resolves ([`BinIndexHold`]). + /// record resolves ([`QueueHoldReason::BinIndex`]). HeldByBinIndex(DefaultsReason), /// The user cancelled the upload. The facade has already undone it, so the /// valve does nothing but stop the pass. @@ -622,50 +622,57 @@ struct Queue { all_ids: BTreeSet, } -/// The queue head is held over rather than failed: the account quota refused -/// it, and it keeps its place and its staging reservation until a probe on a -/// later drain tick reports room. +/// Why the queue head is held over rather than failed. Each reason names its +/// own exit, and [`Drain::hold_admits_the_head`] is the one gate that tries it. #[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct BlockedOp { - /// The held op. - pub op_id: OpId, - /// The node the op targets, so a host can point at it. - pub node: NodeId, - /// The byte count the resume probe must find room for. - pub needed_bytes: u64, +pub enum QueueHoldReason { + /// The account quota refused the upload. The exit is a probe on a later + /// drain tick reporting room for `needed_bytes`. + Quota { + /// The byte count the resume probe must find room for. + needed_bytes: u64, + }, + /// The member's own settings were refused before any request was built, so + /// every retry reaches the same verdict and charging one would spend the + /// version's budget and then release its staged blocks. The exit is + /// settings that name a placement this rule no longer refuses. + /// + /// Render it through [`SettingsRefusal::check`], which names the rule and + /// never the endpoint or the bearer the settings carry. + Settings(SettingsRefusal), + /// The bin index plane did not establish the current index. The exit is a + /// load that establishes it. + /// + /// Reported, because a party who withholds the record — or one head block + /// of it — otherwise stops every queued operation for the account with no + /// cause the member can see (blueprint/engine.md "Bin index record"). + BinIndex(DefaultsReason), } -/// The queue head is held over rather than failed: the member's own settings -/// were refused before any request was built, so every retry reaches the same -/// verdict and charging one would spend the version's budget and then release -/// its staged blocks. It keeps its place and its staging reservation until the -/// settings name a placement that clears the refusal. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct SettingsHold { - /// The held op. - pub op_id: OpId, - /// The node the op targets, so a host can point at it. - pub node: NodeId, - /// Which rule refused. Render it through [`SettingsRefusal::check`], which - /// names the rule and never the endpoint or the bearer the settings carry. - pub refusal: SettingsRefusal, +impl QueueHoldReason { + /// The stable name a host dispatches on. + pub fn name(&self) -> &'static str { + match self { + Self::Quota { .. } => "quota", + Self::Settings(_) => "settings", + Self::BinIndex(_) => "bin-index", + } + } } -/// The queue head is held over rather than failed: the bin index plane did not -/// establish the current index, so the op that needs it keeps its place and its -/// staging reservation until the record resolves. +/// The queue head is held over rather than failed: it keeps its place and its +/// staging reservation until its reason's own exit comes. /// -/// Reported, because a party who withholds the record — or one head block of it -/// — otherwise stops every queued operation for the account with no cause the -/// member can see (blueprint/engine.md "Bin index record"). +/// One cell, so one head cannot be claimed for two reasons at once and no arm +/// has another arm's state to clear. #[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct BinIndexHold { +pub struct QueueHold { /// The held op. pub op_id: OpId, /// The node the op targets, so a host can point at it. pub node: NodeId, - /// Why the load did not establish the index. - pub reason: DefaultsReason, + /// What it waits on. + pub reason: QueueHoldReason, } /// The captures one pass adopts into the bin. A peer chooses both the trigger @@ -1046,13 +1053,9 @@ pub(crate) struct Drain<'a, T, H: Http, C: CredentialStore, F, S, St, Sch> { pub(crate) base: &'a BaseSnapshot, /// The live held-record set the liveness loop keeps alive. pub(crate) held: &'a RefCell, - /// The over-quota hold, shared with the facade's read surface. It clears - /// only here, on a quota probe reporting room. - pub(crate) blocked: &'a RefCell>, - /// The settings-refused hold, shared with the facade's read surface. It - /// clears only here, once `placement` no longer carries the refusal that - /// took it. - pub(crate) settings_hold: &'a RefCell>, + /// The held queue head, shared with the facade's read surface. It clears + /// only here, when its reason's own exit comes. + pub(crate) hold: &'a RefCell>, /// Pinned bytes the retire ledger still owes, shared with the facade's read /// surface. Rewritten at the end of every pass from the ledger itself. pub(crate) pending_reclaim: &'a Cell, @@ -1095,9 +1098,6 @@ pub(crate) struct Drain<'a, T, H: Http, C: CredentialStore, F, S, St, Sch> { /// with the facade's renewal slot. A load fills it too, so the sub-EOL /// renewal keeps the record alive on a session that publishes nothing. pub(crate) bin_index_record: &'a RefCell>, - /// The bin-index-refused hold, shared with the facade's read surface. It - /// clears only here, on a load that establishes the index. - pub(crate) bin_index_hold: &'a RefCell>, /// The bin index this pass has established: the one it resolved, or the one /// its last confirmed publish left standing. Carried so a bulk soft delete /// costs one resolve rather than one per operation; the publish stays per @@ -1620,22 +1620,12 @@ where }; let queued = mine; if queued.is_empty() { - self.clear_block(); - self.clear_settings_hold(); - self.clear_bin_index_hold(); + self.release_hold(); // A debt outlives the op that owed it, so an empty queue is still a // pass that drives the cuts this device owes. self.cut_exited_scopes(scope, exits).await; return (report, Some(Vec::new())); } - // The hold names one op, so it goes as soon as that op does. - if self - .bin_index_hold - .borrow() - .is_some_and(|hold| !still_queued(&queued, hold.op_id)) - { - self.clear_bin_index_hold(); - } let purges = queued .iter() .filter(|(_, op)| matches!(op.kind, OpKind::Purge { .. })) @@ -1682,10 +1672,7 @@ where report: &mut DrainReport, attempts: &mut Attempts, ) -> Result<(), Halt> { - if !self.quota_admits_the_held_head(queued).await { - return Ok(()); - } - if !self.settings_admit_the_held_head(queued) { + if !self.hold_admits_the_head(queued).await { return Ok(()); } @@ -1791,9 +1778,9 @@ where ) { // The bin plane has no probe of its own — the load is the only one — so // its hold goes here, on the first halt that is not it. Every other - // hold's own pre-pass gate is what lets go of it. + // reason has an exit the pre-pass gate can try. if !matches!(halt, Halt::HeldByBinIndex(_)) { - self.clear_bin_index_hold(); + self.release_bin_index_hold(); } match halt { Halt::EpochLagged => {} @@ -1896,48 +1883,50 @@ where Halt::Permanent(reason) => { self.dead_letter(scope, op_id, op, reason, report).await; } - // One pass raises one halt, and each hold's own gate is what lets - // go of it — so taking one drops the others rather than leaving two - // cells claiming the same head for different reasons. Halt::Blocked { needed_bytes } => { - self.clear_settings_hold(); - *self.blocked.borrow_mut() = Some(BlockedOp { - op_id, - node: op.target, - needed_bytes, - }); + self.hold_head(op_id, op, QueueHoldReason::Quota { needed_bytes }); } Halt::HeldBySettings(refusal) => { - self.clear_block(); - *self.settings_hold.borrow_mut() = Some(SettingsHold { - op_id, - node: op.target, - refusal, - }); + self.hold_head(op_id, op, QueueHoldReason::Settings(refusal)); } Halt::HeldByBinIndex(reason) => { - self.clear_block(); - self.clear_settings_hold(); - *self.bin_index_hold.borrow_mut() = Some(BinIndexHold { - op_id, - node: op.target, - reason, - }); + self.hold_head(op_id, op, QueueHoldReason::BinIndex(reason)); } } } - /// Whether a held head may be tried again this tick. A `GET /account/quota` - /// probe reporting room is the hold's only exit, so an unanswered probe - /// leaves it in place. - async fn quota_admits_the_held_head(&self, queued: &[(OpId, Op)]) -> bool { - let Some(blocked) = *self.blocked.borrow() else { + /// Whether the held head may be tried again this tick: the one gate over + /// the one hold cell. A hold whose op has left the queue, or whose reason's + /// own exit has come, lets go of the cell. + async fn hold_admits_the_head(&self, queued: &[(OpId, Op)]) -> bool { + let Some(hold) = *self.hold.borrow() else { return true; }; - if !still_queued(queued, blocked.op_id) { - self.clear_block(); - return true; + if still_queued(queued, hold.op_id) { + match hold.reason { + QueueHoldReason::Quota { needed_bytes } => { + if !self.quota_admits(needed_bytes).await { + return false; + } + } + QueueHoldReason::Settings(refusal) => { + if settings_refusal(self.placement) == Some(refusal) { + return false; + } + } + // The bin index load is its own probe, so this reason neither + // stops a pass nor clears before one: [`Self::apply_valve`] and + // [`Self::establish_bin_index`] are its exits. + QueueHoldReason::BinIndex(_) => return true, + } } + self.release_hold(); + true + } + + /// Whether a `GET /account/quota` probe reports room for a held head. The + /// probe is the hold's only exit, so an unanswered one leaves it in place. + async fn quota_admits(&self, needed_bytes: u64) -> bool { let Ok(placement) = self.placement.as_ref() else { return false; }; @@ -1945,45 +1934,34 @@ where // could give bears on a hold under a placement without one — and an // endpoint that will not answer would park the head on that question. if !placement.has_hosted_leg() { - self.clear_block(); return true; } let Ok(quota) = self.api.quota().await else { return false; }; - if pre_flight_quota_check(blocked.needed_bytes, "a, true).is_err() { - return false; - } - self.clear_block(); - true + pre_flight_quota_check(needed_bytes, "a, true).is_ok() } - fn clear_block(&self) { - *self.blocked.borrow_mut() = None; + fn hold_head(&self, op_id: OpId, op: &Op, reason: QueueHoldReason) { + *self.hold.borrow_mut() = Some(QueueHold { + op_id, + node: op.target, + reason, + }); } - fn clear_bin_index_hold(&self) { - *self.bin_index_hold.borrow_mut() = None; + fn release_hold(&self) { + *self.hold.borrow_mut() = None; } - /// Whether a settings-held head may be tried again this tick: only once the - /// placement this pass runs under stops reaching the verdict that took the - /// hold. - fn settings_admit_the_held_head(&self, queued: &[(OpId, Op)]) -> bool { - let Some(hold) = *self.settings_hold.borrow() else { - return true; - }; - if still_queued(queued, hold.op_id) - && settings_refusal(self.placement) == Some(hold.refusal) - { - return false; + fn release_bin_index_hold(&self) { + let held_by_bin_index = matches!( + self.hold.borrow().map(|hold| hold.reason), + Some(QueueHoldReason::BinIndex(_)) + ); + if held_by_bin_index { + self.release_hold(); } - self.clear_settings_hold(); - true - } - - fn clear_settings_hold(&self) { - *self.settings_hold.borrow_mut() = None; } /// This identity's queued ops, minus restore residue: an op at or below the @@ -3589,7 +3567,7 @@ where /// hold a refused load took. fn establish_bin_index(&self, index: BinIndex) { *self.established_bin_index.borrow_mut() = Some(index); - self.clear_bin_index_hold(); + self.release_bin_index_hold(); } /// Publish the bin index and hold the confirmed record for renewal. @@ -6152,7 +6130,7 @@ fn classify_upload(error: ApiError, refused_bytes: u64) -> Halt { /// that repairs itself. Everything the provider *answered* is charged. /// /// A policy verdict is neither: it is deterministic, so it holds the op rather -/// than charging it ([`SettingsHold`]). +/// than charging it ([`QueueHoldReason::Settings`]). fn classify_placement(error: ProviderError) -> Halt { if error.is_deterministic() { return Halt::HeldBySettings(SettingsRefusal::Byo(error)); @@ -6189,8 +6167,7 @@ fn blocks(count: usize) -> u32 { /// The key-free classification an [`OpPhase::UploadFailed`] carries, or `None` /// where the halt is not a failed attempt: a hold keeps the op and its -/// reservation, and the host reads them from `SnapshotView::blocked`, -/// `SnapshotView::settings_hold` and `SnapshotView::bin_index_hold`. +/// reservation, and the host reads them from `SnapshotView::queue_hold`. fn upload_failure(halt: Halt) -> Option<&'static str> { match halt { // A cancel reports `UploadCancelled` from the facade that ordered it. @@ -7509,8 +7486,9 @@ mod tests { 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, + OWNER_ROOT_EPOCH, OWNER_ROOT_POINTER_READ_KEY, OWNER_ROOT_SCOPE_SEED, + OWNER_ROOT_WRITE_SCOPE_SEED, OwnerRootSpec, SeededEntropy, block_on, gateway, + owner_root_fixture, owner_root_pseudonym, serve, }; /// The login secret the harness's owner identity, enc secret and bin keys @@ -7557,8 +7535,7 @@ mod tests { entropy: RefCell>, base: BaseSnapshot, held: RefCell, - blocked: RefCell>, - settings_hold: RefCell>, + hold: RefCell>, pending_reclaim: Cell, reclaim_stalls: RefCell>, bookkeeping: RefCell, @@ -7572,7 +7549,6 @@ mod tests { 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]>, @@ -7603,8 +7579,7 @@ mod tests { entropy: &self.entropy, base: &self.base, held: &self.held, - blocked: &self.blocked, - settings_hold: &self.settings_hold, + hold: &self.hold, pending_reclaim: &self.pending_reclaim, reclaim_stalls: &self.reclaim_stalls, bookkeeping: &self.bookkeeping, @@ -7616,7 +7591,6 @@ mod tests { 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, @@ -7669,6 +7643,8 @@ mod tests { owner_root_fixture(OwnerRootSpec { owner_identity: &EcdsaSigner::from_scalar(&HARNESS_SECRET).expect("valid scalar"), owner_enc: &kdf::enc_subkey(&HARNESS_SECRET).public(), + writer_pseudonym: &owner_root_pseudonym(), + pointer_read_key: OWNER_ROOT_POINTER_READ_KEY, scope_id: HARNESS_SCOPE, root_id: HARNESS_ROOT.0, children: Vec::new(), @@ -7745,8 +7721,7 @@ mod tests { 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), + hold: RefCell::new(None), pending_reclaim: Cell::new(0), reclaim_stalls: RefCell::new(Vec::new()), bookkeeping: RefCell::new(BookkeepingCursors::default()), @@ -7759,7 +7734,6 @@ mod tests { 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), diff --git a/crates/engine/src/sync/mod.rs b/crates/engine/src/sync/mod.rs index 8a39691661..db670a23f4 100644 --- a/crates/engine/src/sync/mod.rs +++ b/crates/engine/src/sync/mod.rs @@ -44,7 +44,7 @@ pub use doomed::{ MAX_BOOKKEEPING_OPENS, MAX_JOURNAL_REPLAYS, MAX_QUARANTINE_ATTEMPTS, doomed_journal_key, }; pub use drain::{ - BlockedOp, DRAINED_OP_MARK_PREFIX, OP_ATTEMPTS_KEY, PUBLISHED_OP_MARK_PREFIX, SettingsHold, + DRAINED_OP_MARK_PREFIX, OP_ATTEMPTS_KEY, PUBLISHED_OP_MARK_PREFIX, QueueHold, QueueHoldReason, owner_scoped_key, owner_tag, }; pub use model::{Link, NodeMeta, Snapshot, case_fold, collation_key, suffix_name}; diff --git a/crates/engine/tests/write_plane.rs b/crates/engine/tests/write_plane.rs index b0c7a72146..ded122b6e4 100644 --- a/crates/engine/tests/write_plane.rs +++ b/crates/engine/tests/write_plane.rs @@ -34,7 +34,7 @@ use cipherbox_engine::content::{ ByoIpfsConfig, ByoKind, DAG_ROOT_CODEC, PinMode, RetentionPolicy, SealedChunk, SessionBearer, assemble, decode_root, }; -use cipherbox_engine::facade::{BinOrigin, PendingClass}; +use cipherbox_engine::facade::{BinOrigin, PendingClass, SnapshotView}; use cipherbox_engine::net::OrphanHeads; use cipherbox_engine::net::author::{ AuthoredHead, ENVELOPE_V, EnvelopeAuthoring, author_child_envelope, @@ -50,8 +50,8 @@ use cipherbox_engine::seams::{ SnapshotCache, StagingStore, UnixMillis, }; use cipherbox_engine::settings::{ - Destinations, SettingsOrigin, SettingsPublishError, VaultSettings, publish_settings, - settings_name, + Destinations, SettingsOrigin, SettingsPublishError, SettingsRefusal, VaultSettings, + publish_settings, settings_name, }; use cipherbox_engine::sync::pointer::{open_repoint, vault_pointer_name}; use cipherbox_engine::sync::{ @@ -77,10 +77,10 @@ use cipherbox_engine::{ CommittedSet, ContentProfile, DEFAULT_BIN_RETENTION_DAYS, DeadLetter, DeadLetterReason, DefaultsReason, Engine, EngineError, Entropy, EntropyError, Event, EventStream, GatewayConfig, LoginSecret, MAX_FOCUS_FILES, MAX_FOLDER_CHILDREN, MAX_OPEN_STREAMS, NodeId, NodeKind, Op, - OpKind, OpPhase, OverBudgetCause, Placement, PlacementRefusal, PrevEpochSeed, RecordReader, - RecordSeal, ResealSeeds, ScopeCrossing, ScopeRootIdentity, StoragePolicy, SyncTimingProfile, - WriteHistory, WriteTarget, decode_queue, load_bin_index, publish_bin_index, reseal_scope_root, - stage_op, + OpKind, OpPhase, OverBudgetCause, Placement, PlacementRefusal, PrevEpochSeed, QueueHold, + QueueHoldReason, RecordReader, RecordSeal, ResealSeeds, ScopeCrossing, ScopeRootIdentity, + StoragePolicy, SyncTimingProfile, WriteHistory, WriteTarget, decode_queue, load_bin_index, + publish_bin_index, reseal_scope_root, stage_op, }; /// The override seed a rotation mints for `SCOPE`'s second read epoch. @@ -88,6 +88,48 @@ const ROTATED_READ_SCOPE_SEED: [u8; 32] = [0xA5; 32]; /// The stable per-scope pointer read key the owner-root fixture's grant blobs /// carry. const POINTER_READ_KEY: [u8; 32] = [0x88; 32]; + +/// The held queue head when the account quota is what holds it, with the byte +/// count the resume probe must find room for. +fn quota_hold(view: &SnapshotView) -> Option<(QueueHold, u64)> { + match view.queue_hold { + Some( + hold @ QueueHold { + reason: QueueHoldReason::Quota { needed_bytes }, + .. + }, + ) => Some((hold, needed_bytes)), + _ => None, + } +} + +/// The held queue head when the member's own settings are what hold it, with +/// the rule that refused. +fn settings_hold(view: &SnapshotView) -> Option<(QueueHold, SettingsRefusal)> { + match view.queue_hold { + Some( + hold @ QueueHold { + reason: QueueHoldReason::Settings(refusal), + .. + }, + ) => Some((hold, refusal)), + _ => None, + } +} + +/// The held queue head when the owner's bin index is what holds it, with the +/// reason the load did not establish the index. +fn bin_index_hold(view: &SnapshotView) -> Option<(QueueHold, DefaultsReason)> { + match view.queue_hold { + Some( + hold @ QueueHold { + reason: QueueHoldReason::BinIndex(reason), + .. + }, + ) => Some((hold, reason)), + _ => None, + } +} /// The destination set the upload mark opens on. const DESTINATIONS_LEN: usize = Destinations::LEN; @@ -3368,15 +3410,9 @@ fn an_over_quota_upload_holds_the_op_without_reporting_a_failure() { )), "a full account is not a failed upload attempt" ); - assert_eq!( - block_on(engine.snapshot(ROOT)) - .expect("a snapshot") - .blocked - .expect("the over-quota head is held") - .op_id, - op_id, - "the hold is what the host acts on" - ); + let (hold, _) = quota_hold(&block_on(engine.snapshot(ROOT)).expect("a snapshot")) + .expect("the over-quota head is held"); + assert_eq!(hold.op_id, op_id, "the hold is what the host acts on"); } /// An op the completion record already covers is restore residue: a data dir @@ -4726,11 +4762,10 @@ fn a_withheld_bin_index_reports_a_named_hold_that_clears_when_it_resolves() { tick(&world, &engine, &mut tasks); let view = block_on(engine.snapshot(ROOT)).expect("a snapshot"); - let hold = view - .bin_index_hold - .expect("a bin index the pass cannot read holds the head"); + let (hold, reason) = + bin_index_hold(&view).expect("a bin index the pass cannot read holds the head"); assert_eq!(hold.node, second); - assert_eq!(hold.reason.check(), "suppressed"); + assert_eq!(reason.check(), "suppressed"); assert!( view.dead_letters.is_empty(), "a withheld record is not a failed op" @@ -4745,7 +4780,7 @@ fn a_withheld_bin_index_reports_a_named_hold_that_clears_when_it_resolves() { let view = block_on(engine.snapshot(ROOT)).expect("a snapshot"); assert!( - view.bin_index_hold.is_none(), + bin_index_hold(&view).is_none(), "the hold clears when the record resolves" ); assert!( @@ -7374,13 +7409,11 @@ fn an_over_quota_source_remove_holds_the_op_with_its_needed_bytes() { }); tick(&world, &engine, &mut tasks); - let blocked = block_on(engine.snapshot(ROOT)) - .expect("a snapshot") - .blocked + let (hold, needed_bytes) = quota_hold(&block_on(engine.snapshot(ROOT)).expect("a snapshot")) .expect("the over-quota source-remove is held"); - assert_eq!(blocked.op_id, op_id); + assert_eq!(hold.op_id, op_id); assert!( - blocked.needed_bytes > 0, + needed_bytes > 0, "the figure the resume probe must find room for survives the leg" ); } @@ -7870,11 +7903,11 @@ fn an_over_quota_413_holds_the_head_and_a_quota_probe_with_room_resumes_it() { tick(&world, &engine, &mut tasks); let view = block_on(engine.snapshot(ROOT)).expect("a snapshot"); - let blocked = view.blocked.expect("the over-quota head is held"); - assert_eq!(blocked.op_id, held); - assert_eq!(blocked.node, photos); + let (hold, needed_bytes) = quota_hold(&view).expect("the over-quota head is held"); + assert_eq!(hold.op_id, held); + assert_eq!(hold.node, photos); assert!( - blocked.needed_bytes > 0, + needed_bytes > 0, "the hold records what the refused upload asked for" ); assert!( @@ -7900,7 +7933,7 @@ fn an_over_quota_413_holds_the_head_and_a_quota_probe_with_room_resumes_it() { before, "a held head re-probes the quota, never the upload" ); - assert!(block_on(engine.snapshot(ROOT)).unwrap().blocked.is_some()); + assert!(quota_hold(&block_on(engine.snapshot(ROOT)).unwrap()).is_some()); // Room appears. blocks.accept_uploads(); @@ -7908,7 +7941,7 @@ fn an_over_quota_413_holds_the_head_and_a_quota_probe_with_room_resumes_it() { tick(&world, &engine, &mut tasks); let view = block_on(engine.snapshot(ROOT)).expect("a snapshot"); - assert!(view.blocked.is_none(), "the probe cleared the hold"); + assert!(quota_hold(&view).is_none(), "the probe cleared the hold"); assert_eq!( published_names(&world.record_store, &blocks, ROOT), vec!["notes".to_owned(), "photos".to_owned()], @@ -7930,7 +7963,7 @@ fn an_over_cap_413_is_permanent_and_its_reason_reaches_the_host() { let view = block_on(engine.snapshot(ROOT)).expect("a snapshot"); assert!( - view.blocked.is_none(), + quota_hold(&view).is_none(), "the transport cap is not the account-quota gate" ); assert_eq!( @@ -7974,7 +8007,7 @@ fn a_413_the_api_did_not_stamp_neither_blocks_nor_abandons_the_op() { let view = block_on(engine.snapshot(ROOT)).expect("a snapshot"); assert!( - view.blocked.is_none(), + quota_hold(&view).is_none(), "no positive quota evidence, no hold" ); assert!( @@ -11031,10 +11064,7 @@ fn a_hold_under_an_external_placement_clears_without_a_quota_probe() { create(&mut engine, "photos"); tick(&world, &engine, &mut tasks); assert!( - block_on(engine.snapshot(ROOT)) - .expect("a snapshot") - .blocked - .is_some(), + quota_hold(&block_on(engine.snapshot(ROOT)).expect("a snapshot")).is_some(), "the refused head block held the op" ); @@ -11053,7 +11083,7 @@ fn a_hold_under_an_external_placement_clears_without_a_quota_probe() { let view = block_on(engine.snapshot(ROOT)).expect("a snapshot"); assert!( - view.blocked.is_none(), + quota_hold(&view).is_none(), "an unreachable quota endpoint never gates a placement it does not cover" ); assert_eq!( @@ -11834,11 +11864,11 @@ fn a_deterministic_placement_refusal_holds_the_queued_write_rather_than_charging .contains(&root_cid), "and its staged version with it" ); - let hold = view.settings_hold.expect("the pass names what it waits on"); + let (hold, refusal) = settings_hold(&view).expect("the pass names what it waits on"); assert_eq!(hold.op_id, op_id); assert_eq!(hold.node, photo); assert_eq!( - hold.refusal.check(), + refusal.check(), "byo-provider-missing", "the rule that refused, never the settings it read" ); @@ -11875,7 +11905,8 @@ fn a_degraded_settings_load_retries_the_queued_write_and_takes_no_hold() { "an outage this pass could not resolve never spends the budget" ); assert_eq!( - view.settings_hold, None, + settings_hold(&view), + None, "no settings change is what this head is waiting for" ); assert_eq!( diff --git a/crates/fuse/tests/fuse_op_core.rs b/crates/fuse/tests/fuse_op_core.rs index fa77bb08b5..2f1315e708 100644 --- a/crates/fuse/tests/fuse_op_core.rs +++ b/crates/fuse/tests/fuse_op_core.rs @@ -2316,7 +2316,7 @@ fn a_dead_lettered_op_reaches_the_mount_status_with_its_reason() { assert!( block_on(core.status()) .expect("the status reads again") - .blocked + .queue_hold .is_none(), "a dead letter is not a drain hold" ); @@ -2357,7 +2357,7 @@ fn the_mount_status_reports_a_quiet_mount_as_quiet() { let status = block_on(core.status()).expect("the status reads"); assert!(status.dead_letters.is_empty()); - assert!(status.blocked.is_none()); + assert!(status.queue_hold.is_none()); assert_eq!(status.retained_records, 0); assert_eq!(status.staleness, Staleness::Fresh); } diff --git a/crates/wasm/src/lib.rs b/crates/wasm/src/lib.rs index 58373b0f48..1a3f4dea8c 100644 --- a/crates/wasm/src/lib.rs +++ b/crates/wasm/src/lib.rs @@ -736,15 +736,15 @@ impl DeadLetter { } } -/// The queue head held over the account quota, keeping its place and its -/// staging reservation until a quota probe reports room. +/// The queue head held over rather than failed, keeping its place and its +/// staging reservation until its reason's own exit comes. #[wasm_bindgen] -pub struct BlockedOp { - inner: facade::BlockedOp, +pub struct QueueHold { + inner: facade::QueueHold, } #[wasm_bindgen] -impl BlockedOp { +impl QueueHold { /// The held op id (a `u64`, crossing as a `bigint`). #[wasm_bindgen(getter, js_name = opId)] pub fn op_id(&self) -> u64 { @@ -757,67 +757,31 @@ impl BlockedOp { self.inner.node.0.to_vec() } - /// The byte count the resume probe must find room for. - #[wasm_bindgen(getter, js_name = neededBytes)] - pub fn needed_bytes(&self) -> u64 { - self.inner.needed_bytes - } -} - -/// The queue head held over the member's own settings, keeping its place and -/// its staging reservation until those settings change. -#[wasm_bindgen] -pub struct SettingsHold { - inner: facade::SettingsHold, -} - -#[wasm_bindgen] -impl SettingsHold { - /// The held op id (a `u64`, crossing as a `bigint`). - #[wasm_bindgen(getter, js_name = opId)] - pub fn op_id(&self) -> u64 { - self.inner.op_id.0 - } - - /// The 16 raw bytes of the node the held op targets. - #[wasm_bindgen(getter)] - pub fn node(&self) -> Vec { - self.inner.node.0.to_vec() - } - - /// The stable check name of the rule that refused. Never the endpoint or - /// the bearer those settings carry. + /// What the hold waits on: `quota`, `settings` or `bin-index`. #[wasm_bindgen(getter)] - pub fn check(&self) -> String { - self.inner.refusal.check().to_owned() + pub fn reason(&self) -> String { + self.inner.reason.name().to_owned() } -} -/// The queue head held over the owner's bin index, keeping its place and its -/// staging reservation until that record resolves. -#[wasm_bindgen] -pub struct BinIndexHold { - inner: facade::BinIndexHold, -} - -#[wasm_bindgen] -impl BinIndexHold { - /// The held op id (a `u64`, crossing as a `bigint`). - #[wasm_bindgen(getter, js_name = opId)] - pub fn op_id(&self) -> u64 { - self.inner.op_id.0 - } - - /// The 16 raw bytes of the node the held op targets. - #[wasm_bindgen(getter)] - pub fn node(&self) -> Vec { - self.inner.node.0.to_vec() + /// The byte count the resume probe must find room for, on a quota hold and + /// nowhere else. + #[wasm_bindgen(getter, js_name = neededBytes)] + pub fn needed_bytes(&self) -> Option { + match self.inner.reason { + facade::QueueHoldReason::Quota { needed_bytes } => Some(needed_bytes), + _ => None, + } } - /// The stable check name of what the bin index load produced. + /// The stable check name of what refused, on a settings or a bin index hold. + /// Never the endpoint or the bearer the settings carry. #[wasm_bindgen(getter)] - pub fn check(&self) -> String { - self.inner.reason.check().to_owned() + pub fn check(&self) -> Option { + match self.inner.reason { + facade::QueueHoldReason::Quota { .. } => None, + facade::QueueHoldReason::Settings(refusal) => Some(refusal.check().to_owned()), + facade::QueueHoldReason::BinIndex(reason) => Some(reason.check().to_owned()), + } } } @@ -912,24 +876,10 @@ impl SnapshotView { .collect() } - /// The drain's over-quota hold, or `undefined`. - #[wasm_bindgen(getter)] - pub fn blocked(&self) -> Option { - self.inner.blocked.map(|inner| BlockedOp { inner }) - } - - /// The drain's settings-refused hold, or `undefined`. - #[wasm_bindgen(getter, js_name = settingsHold)] - pub fn settings_hold(&self) -> Option { - self.inner.settings_hold.map(|inner| SettingsHold { inner }) - } - - /// The drain's bin-index-refused hold, or `undefined`. - #[wasm_bindgen(getter, js_name = binIndexHold)] - pub fn bin_index_hold(&self) -> Option { - self.inner - .bin_index_hold - .map(|inner| BinIndexHold { inner }) + /// The drain's held queue head, or `undefined`. + #[wasm_bindgen(getter, js_name = queueHold)] + pub fn queue_hold(&self) -> Option { + self.inner.queue_hold.map(|inner| QueueHold { inner }) } /// Durable queue entries this session holds but cannot read — another @@ -2602,22 +2552,10 @@ mod tests { reason: facade::DeadLetterReason::AttemptsExhausted, }, ], - blocked: Some(facade::BlockedOp { + queue_hold: Some(facade::QueueHold { op_id: OpId(12), node: facade::NodeId([5u8; 16]), - needed_bytes: 4096, - }), - settings_hold: Some(facade::SettingsHold { - op_id: OpId(13), - node: facade::NodeId([6u8; 16]), - refusal: cipherbox_engine::SettingsRefusal::Byo( - cipherbox_engine::ProviderError::InsecureTransport, - ), - }), - bin_index_hold: Some(facade::BinIndexHold { - op_id: OpId(14), - node: facade::NodeId([7u8; 16]), - reason: cipherbox_engine::DefaultsReason::Suppressed, + reason: facade::QueueHoldReason::Quota { needed_bytes: 4096 }, }), retained_records: 3, staleness: facade::Staleness::Reconciling, @@ -2637,18 +2575,12 @@ mod tests { (11, DeadLetterReason::AttemptsExhausted), ] ); - let blocked = view.blocked().expect("the view carries the hold"); - assert_eq!(blocked.op_id(), 12); - assert_eq!(blocked.node(), vec![5u8; 16]); - assert_eq!(blocked.needed_bytes(), 4096); - let held = view.settings_hold().expect("the view carries the hold"); - assert_eq!(held.op_id(), 13); - assert_eq!(held.node(), vec![6u8; 16]); - assert_eq!(held.check(), "byo-endpoint-insecure"); - let bin_held = view.bin_index_hold().expect("the view carries the hold"); - assert_eq!(bin_held.op_id(), 14); - assert_eq!(bin_held.node(), vec![7u8; 16]); - assert_eq!(bin_held.check(), "suppressed"); + let hold = view.queue_hold().expect("the view carries the hold"); + assert_eq!(hold.op_id(), 12); + assert_eq!(hold.node(), vec![5u8; 16]); + assert_eq!(hold.reason(), "quota"); + assert_eq!(hold.needed_bytes(), Some(4096)); + assert_eq!(hold.check(), None); assert_eq!(view.retained_records(), 3); assert_eq!(view.staleness(), Staleness::Reconciling); @@ -2674,4 +2606,35 @@ mod tests { assert_eq!(ancestors[0].id(), vec![1u8; 16]); assert_eq!(ancestors[0].name(), ""); } + + /// A host dispatches on the reason, and each one carries exactly the figure + /// its own notice renders. + #[test] + fn a_queue_hold_names_its_reason_and_carries_only_that_reasons_figure() { + let settings = QueueHold { + inner: facade::QueueHold { + op_id: OpId(13), + node: facade::NodeId([6u8; 16]), + reason: facade::QueueHoldReason::Settings(cipherbox_engine::SettingsRefusal::Byo( + cipherbox_engine::ProviderError::InsecureTransport, + )), + }, + }; + assert_eq!(settings.reason(), "settings"); + assert_eq!(settings.check().as_deref(), Some("byo-endpoint-insecure")); + assert_eq!(settings.needed_bytes(), None); + + let bin_index = QueueHold { + inner: facade::QueueHold { + op_id: OpId(14), + node: facade::NodeId([7u8; 16]), + reason: facade::QueueHoldReason::BinIndex( + cipherbox_engine::DefaultsReason::Suppressed, + ), + }, + }; + assert_eq!(bin_index.reason(), "bin-index"); + assert_eq!(bin_index.check().as_deref(), Some("suppressed")); + assert_eq!(bin_index.needed_bytes(), None); + } } diff --git a/crates/wasm/tests/boundary.rs b/crates/wasm/tests/boundary.rs index 9b3123c3e2..c883eaceef 100644 --- a/crates/wasm/tests/boundary.rs +++ b/crates/wasm/tests/boundary.rs @@ -343,22 +343,12 @@ fn snapshot_view_getters_cross_with_boundary_shapes() { op_id: OpId(9), reason: facade::DeadLetterReason::SuffixExhausted, }], - blocked: Some(facade::BlockedOp { + queue_hold: Some(facade::QueueHold { op_id: OpId(12), node: facade::NodeId([6u8; 16]), - needed_bytes: u64::MAX, - }), - settings_hold: Some(facade::SettingsHold { - op_id: OpId(13), - node: facade::NodeId([7u8; 16]), - refusal: cipherbox_engine::SettingsRefusal::Byo( - cipherbox_engine::ProviderError::BlockedAddress, - ), - }), - bin_index_hold: Some(facade::BinIndexHold { - op_id: OpId(14), - node: facade::NodeId([8u8; 16]), - reason: cipherbox_engine::DefaultsReason::Suppressed, + reason: facade::QueueHoldReason::Quota { + needed_bytes: u64::MAX, + }, }), retained_records: 0, staleness: facade::Staleness::Fresh, @@ -416,8 +406,18 @@ fn snapshot_view_getters_cross_with_boundary_shapes() { "the reason crosses as its mirror-enum ordinal" ); - let blocked = get(&view, "blocked"); - let needed = get(&blocked, "neededBytes"); + let hold = get(&view, "queueHold"); + assert_eq!( + get(&hold, "opId").js_typeof(), + JsValue::from_str("bigint"), + "a held op's opId must cross as a JS bigint, never a number" + ); + assert_eq!( + get(&hold, "reason"), + JsValue::from_str("quota"), + "the host dispatches on the reason name" + ); + let needed = get(&hold, "neededBytes"); assert_eq!( needed.js_typeof(), JsValue::from_str("bigint"), @@ -432,45 +432,13 @@ fn snapshot_view_getters_cross_with_boundary_shapes() { ), u64::MAX.to_string() ); - assert_eq!( - get(&blocked, "node") - .unchecked_into::() - .to_vec(), - vec![6u8; 16] - ); - - let held = get(&view, "settingsHold"); - assert_eq!( - get(&held, "opId").js_typeof(), - JsValue::from_str("bigint"), - "a held op's opId must cross as a JS bigint, never a number" - ); - assert_eq!( - get(&held, "node").unchecked_into::().to_vec(), - vec![7u8; 16] - ); - assert_eq!( - get(&held, "check"), - JsValue::from_str("byo-endpoint-blocked"), - "the refusing rule crosses by its stable check name" - ); - - let bin_held = get(&view, "binIndexHold"); - assert_eq!( - get(&bin_held, "opId").js_typeof(), - JsValue::from_str("bigint"), - "a held op's opId must cross as a JS bigint, never a number" - ); - assert_eq!( - get(&bin_held, "node") - .unchecked_into::() - .to_vec(), - vec![8u8; 16] + assert!( + get(&hold, "check").is_undefined(), + "a quota hold carries no check name" ); assert_eq!( - get(&bin_held, "check"), - JsValue::from_str("suppressed"), - "the load outcome crosses by its stable check name, carrying no figures" + get(&hold, "node").unchecked_into::().to_vec(), + vec![6u8; 16] ); let children = get(&view, "children"); @@ -539,6 +507,47 @@ fn snapshot_view_getters_cross_with_boundary_shapes() { ); } +/// A settings hold crosses with its check name and no byte figure: the host +/// renders the rule that refused, and nothing a quota hold would carry. +#[wasm_bindgen_test] +fn a_settings_queue_hold_crosses_with_its_check_and_no_byte_figure() { + let view: JsValue = SnapshotView::from_facade(facade::SnapshotView { + root: facade::NodeId([1u8; 16]), + folder: facade::NodeId([1u8; 16]), + folder_name: String::new(), + children: Vec::new(), + ancestors: Vec::new(), + dead_letters: Vec::new(), + queue_hold: Some(facade::QueueHold { + op_id: OpId(13), + node: facade::NodeId([7u8; 16]), + reason: facade::QueueHoldReason::Settings(cipherbox_engine::SettingsRefusal::Byo( + cipherbox_engine::ProviderError::BlockedAddress, + )), + }), + retained_records: 0, + staleness: facade::Staleness::Fresh, + }) + .into(); + + let hold = Reflect::get(&view, &JsValue::from_str("queueHold")).expect("getter is readable"); + let get = |key: &str| Reflect::get(&hold, &JsValue::from_str(key)).expect("getter is readable"); + assert_eq!(get("reason"), JsValue::from_str("settings")); + assert_eq!( + get("check"), + JsValue::from_str("byo-endpoint-blocked"), + "the refusing rule crosses by its stable check name" + ); + assert!( + get("neededBytes").is_undefined(), + "only a quota hold carries a byte figure" + ); + assert_eq!( + get("node").unchecked_into::().to_vec(), + vec![7u8; 16] + ); +} + /// The refusal builds a `JsError`, so it is only reachable on this target. #[wasm_bindgen_test] fn a_zero_retention_cap_is_refused_rather_than_defaulted() { diff --git a/packages/client/src/broadcastTransport.test.ts b/packages/client/src/broadcastTransport.test.ts index b41aeaf715..977874a4bf 100644 --- a/packages/client/src/broadcastTransport.test.ts +++ b/packages/client/src/broadcastTransport.test.ts @@ -625,9 +625,7 @@ describe('broadcast transport ↔ leader relay', () => { ], ancestors: [{ id: new Uint8Array(16).fill(1), name: '' }], deadLetters: [{ opId: 7n, reason: 'suffixExhausted' }], - blocked: null, - settingsHold: null, - binIndexHold: null, + queueHold: null, retainedRecords: 0, staleness: 'reconciling', }; diff --git a/packages/client/src/index.ts b/packages/client/src/index.ts index 911a8d7225..2e2467faa9 100644 --- a/packages/client/src/index.ts +++ b/packages/client/src/index.ts @@ -64,11 +64,12 @@ export type { OpProgressPhase, DeadLetterReason, DeadLetterDescriptor, - BlockedOpDescriptor, SettingsHoldCheck, SettingsHoldDescriptor, BinIndexHoldCheck, BinIndexHoldDescriptor, + QuotaHoldDescriptor, + QueueHoldDescriptor, SnapshotDescriptor, SnapshotChildDescriptor, BreadcrumbDescriptor, diff --git a/packages/client/src/testkit.ts b/packages/client/src/testkit.ts index 9404661c71..f46130c20b 100644 --- a/packages/client/src/testkit.ts +++ b/packages/client/src/testkit.ts @@ -117,9 +117,7 @@ export function emptySnapshot(folder: Uint8Array = new Uint8Array(16)): Snapshot children: [], ancestors: [], deadLetters: [], - blocked: null, - settingsHold: null, - binIndexHold: null, + queueHold: null, retainedRecords: 0, staleness: 'fresh', }; diff --git a/packages/client/src/worker/commandCodec.test.ts b/packages/client/src/worker/commandCodec.test.ts index fc7d4df640..7bbb6104f6 100644 --- a/packages/client/src/worker/commandCodec.test.ts +++ b/packages/client/src/worker/commandCodec.test.ts @@ -1099,21 +1099,12 @@ describe('readSnapshot', () => { { opId: 9n, reason: 4 }, { opId: 9_007_199_254_740_993n, reason: 6 }, ], - blocked: { + queueHold: { opId: 12n, node: new Uint8Array(16).fill(6), + reason: 'quota', neededBytes: 9_007_199_254_740_993n, }, - settingsHold: { - opId: 13n, - node: new Uint8Array(16).fill(7), - check: 'byo-provider-missing', - }, - binIndexHold: { - opId: 14n, - node: new Uint8Array(16).fill(8), - check: 'suppressed', - }, retainedRecords: 2, staleness: 1, }; @@ -1162,40 +1153,64 @@ describe('readSnapshot', () => { { opId: 9n, reason: 'undecodable' }, { opId: 9_007_199_254_740_993n, reason: 'attemptsExhausted' }, ], - blocked: { + queueHold: { opId: 12n, node: new Uint8Array(16).fill(6), + reason: 'quota', neededBytes: 9_007_199_254_740_993n, }, - settingsHold: { + retainedRecords: 2, + staleness: 'reconciling', + }); + }); + + it('maps an absent hold to null', () => { + expect(readSnapshot(fakeWasm, baseView()).queueHold).toBeNull(); + }); + + it('reads a held head by its reason', () => { + const view = { + ...baseView(), + queueHold: { opId: 13n, node: new Uint8Array(16).fill(7), + reason: 'settings', check: 'byo-provider-missing', }, - binIndexHold: { - opId: 14n, - node: new Uint8Array(16).fill(8), - check: 'suppressed', - }, - retainedRecords: 2, - staleness: 'reconciling', + }; + expect(readSnapshot(fakeWasm, view).queueHold).toEqual({ + opId: 13n, + node: new Uint8Array(16).fill(7), + reason: 'settings', + check: 'byo-provider-missing', }); }); - it('maps an absent over-budget hold to null', () => { - expect(readSnapshot(fakeWasm, baseView()).blocked).toBeNull(); + it('fails closed on a hold reason this build cannot name', () => { + const view = { + ...baseView(), + queueHold: { opId: 1n, node: new Uint8Array(16), reason: 'weather' }, + }; + expect(() => readSnapshot(fakeWasm, view)).toThrow('unknown WASM queue hold reason: weather'); }); - it('maps an absent settings hold and bin index hold to null', () => { - const view = readSnapshot(fakeWasm, baseView()); - expect(view.settingsHold).toBeNull(); - expect(view.binIndexHold).toBeNull(); + it('fails closed on a quota hold that carries no byte count', () => { + const view = { + ...baseView(), + queueHold: { opId: 1n, node: new Uint8Array(16), reason: 'quota' }, + }; + expect(() => readSnapshot(fakeWasm, view)).toThrow('WASM quota hold carries no byte count'); }); it('fails closed on a hold check this build cannot name', () => { const settings = { ...baseView(), - settingsHold: { opId: 1n, node: new Uint8Array(16), check: 'byo-unreachable' }, + queueHold: { + opId: 1n, + node: new Uint8Array(16), + reason: 'settings', + check: 'byo-unreachable', + }, }; expect(() => readSnapshot(fakeWasm, settings)).toThrow( 'unknown WASM settings hold check: byo-unreachable' @@ -1203,17 +1218,22 @@ describe('readSnapshot', () => { const bin = { ...baseView(), - binIndexHold: { opId: 1n, node: new Uint8Array(16), check: 'stranded-mint' }, + queueHold: { + opId: 1n, + node: new Uint8Array(16), + reason: 'bin-index', + check: 'stranded-mint', + }, }; expect(() => readSnapshot(fakeWasm, bin)).toThrow( - 'unknown WASM bin index hold check: stranded-mint' + 'unknown WASM bin-index hold check: stranded-mint' ); }); it('holds each check vocabulary apart', () => { const crossed = { ...baseView(), - settingsHold: { opId: 1n, node: new Uint8Array(16), check: 'suppressed' }, + queueHold: { opId: 1n, node: new Uint8Array(16), reason: 'settings', check: 'suppressed' }, }; expect(() => readSnapshot(fakeWasm, crossed)).toThrow('unknown WASM settings hold check'); }); diff --git a/packages/client/src/worker/commandCodec.ts b/packages/client/src/worker/commandCodec.ts index b588815f1b..bc9d5a442c 100644 --- a/packages/client/src/worker/commandCodec.ts +++ b/packages/client/src/worker/commandCodec.ts @@ -13,7 +13,6 @@ import type { AuthMethodKind, BinDescriptor, BinOriginDescriptor, - BlockedOpDescriptor, ByoKind, CommandDescriptor, DeadLetterReason, @@ -30,6 +29,7 @@ import type { RegisteredDeviceDescriptor, SettingsOrigin, SharingDescriptor, + QueueHoldDescriptor, SnapshotDescriptor, Staleness, VaultStorageDescriptor, @@ -39,7 +39,6 @@ import type { WasmAuthMethod, WasmBinRow, WasmBinView, - WasmBlockedOp, WasmCommand, WasmEvent, WasmByoIpfsConfig, @@ -487,31 +486,35 @@ function deadLetterReason(wasm: EngineWasm, reason: number | undefined): DeadLet } } -function blockedHold(blocked: WasmBlockedOp | undefined): BlockedOpDescriptor | null { - if (blocked === undefined) return null; - return { - opId: blocked.opId, - node: blocked.node, - neededBytes: blocked.neededBytes, - }; -} - /** - * Reads a held queue head, refusing a check name this build does not know. A - * hold whose cause cannot be named would render as an unexplained stall, which - * is the state the hold exists to remove. + * Reads the held queue head, refusing a reason or a check name this build does + * not know. A hold whose cause cannot be named would render as an unexplained + * stall, which is the state the hold exists to remove. */ -function queueHold( - hold: WasmQueueHold | undefined, - checks: readonly TCheck[], - held: string -): { opId: bigint; node: Uint8Array; check: TCheck } | null { +function queueHold(hold: WasmQueueHold | undefined): QueueHoldDescriptor | null { if (hold === undefined) return null; + const head = { opId: hold.opId, node: hold.node }; + switch (hold.reason) { + case 'quota': + if (hold.neededBytes === undefined) { + throw new Error('WASM quota hold carries no byte count'); + } + return { ...head, reason: 'quota', neededBytes: hold.neededBytes }; + case 'settings': + return { ...head, reason: 'settings', check: holdCheck(hold, SETTINGS_HOLD_CHECKS) }; + case 'bin-index': + return { ...head, reason: 'bin-index', check: holdCheck(hold, BIN_INDEX_HOLD_CHECKS) }; + default: + throw new Error(`unknown WASM queue hold reason: ${hold.reason}`); + } +} + +function holdCheck(hold: WasmQueueHold, checks: readonly TCheck[]): TCheck { const check = checks.find((known) => known === hold.check); if (check === undefined) { - throw new Error(`unknown WASM ${held} hold check: ${hold.check}`); + throw new Error(`unknown WASM ${hold.reason} hold check: ${hold.check}`); } - return { opId: hold.opId, node: hold.node, check }; + return check; } function nodeKindFrom(wasm: EngineWasm, kind: number): NodeKind { @@ -606,9 +609,7 @@ export function readSnapshot(wasm: EngineWasm, view: WasmSnapshotView): Snapshot opId: dead.opId, reason: deadLetterReason(wasm, dead.reason), })), - blocked: blockedHold(view.blocked), - settingsHold: queueHold(view.settingsHold, SETTINGS_HOLD_CHECKS, 'settings'), - binIndexHold: queueHold(view.binIndexHold, BIN_INDEX_HOLD_CHECKS, 'bin index'), + queueHold: queueHold(view.queueHold), retainedRecords: view.retainedRecords, staleness: staleness(wasm, view.staleness), }; diff --git a/packages/client/src/worker/engineWasm.ts b/packages/client/src/worker/engineWasm.ts index 7bd4ad55d5..f5314bb21e 100644 --- a/packages/client/src/worker/engineWasm.ts +++ b/packages/client/src/worker/engineWasm.ts @@ -90,22 +90,17 @@ export interface WasmDeadLetter { readonly reason: number; } -/** wasm-bindgen `BlockedOp` — the drain's over-budget hold. */ -export interface WasmBlockedOp { - readonly opId: bigint; - readonly node: Uint8Array; - readonly neededBytes: bigint; -} - /** - * wasm-bindgen `SettingsHold` / `BinIndexHold` — a held queue head and the - * stable check name of what refused it. The two carry different check - * vocabularies, which `protocol.ts` maps apart. + * wasm-bindgen `QueueHold` — the held queue head and the reason it is held. + * Each reason carries exactly one figure, so the other is `undefined`: + * `neededBytes` on a quota hold, `check` on the two the host renders by name. */ export interface WasmQueueHold { readonly opId: bigint; readonly node: Uint8Array; - readonly check: string; + readonly reason: string; + readonly neededBytes?: bigint; + readonly check?: string; } /** @@ -127,9 +122,7 @@ export interface WasmSnapshotView { readonly children: WasmSnapshotChild[]; readonly ancestors: WasmBreadcrumb[]; readonly deadLetters: readonly WasmDeadLetter[]; - readonly blocked?: WasmBlockedOp; - readonly settingsHold?: WasmQueueHold; - readonly binIndexHold?: WasmQueueHold; + readonly queueHold?: WasmQueueHold; readonly retainedRecords: number; readonly staleness: number; } diff --git a/packages/client/src/worker/protocol.ts b/packages/client/src/worker/protocol.ts index 4ed5c070fa..6b1a9aa720 100644 --- a/packages/client/src/worker/protocol.ts +++ b/packages/client/src/worker/protocol.ts @@ -88,13 +88,6 @@ export interface DeadLetterDescriptor { reason: DeadLetterReason; } -/** The drain's over-budget hold, as data (mirrors the facade `BlockedOp`). */ -export interface BlockedOpDescriptor { - opId: bigint; - node: Uint8Array; - neededBytes: bigint; -} - /** * The rule that refused the member's own settings, as the engine's stable check * names. Only the verdicts a settings hold can carry: a hold waits on the member @@ -126,27 +119,43 @@ export const BIN_INDEX_HOLD_CHECKS = [ export type BinIndexHoldCheck = (typeof BIN_INDEX_HOLD_CHECKS)[number]; -/** - * The queue head held over the member's own settings, as data (mirrors the - * facade `SettingsHold`). The check names the rule, never the endpoint or the - * bearer those settings carry. - */ -export interface SettingsHoldDescriptor { +/** The held op and the node it targets, which every hold reason carries. */ +interface HeldQueueHead { opId: bigint; node: Uint8Array; - check: SettingsHoldCheck; +} + +/** The queue head held over the account quota. */ +export interface QuotaHoldDescriptor extends HeldQueueHead { + reason: 'quota'; + neededBytes: bigint; } /** - * The queue head held over the owner's bin index, as data (mirrors the facade - * `BinIndexHold`). + * The queue head held over the member's own settings. The check names the rule, + * never the endpoint or the bearer those settings carry. */ -export interface BinIndexHoldDescriptor { - opId: bigint; - node: Uint8Array; +export interface SettingsHoldDescriptor extends HeldQueueHead { + reason: 'settings'; + check: SettingsHoldCheck; +} + +/** The queue head held over the owner's bin index. */ +export interface BinIndexHoldDescriptor extends HeldQueueHead { + reason: 'bin-index'; check: BinIndexHoldCheck; } +/** + * The one held queue head, as data (mirrors the facade `QueueHold`). One head + * is held for one reason, so a host dispatches on `reason` rather than reading + * parallel fields. + */ +export type QueueHoldDescriptor = + | QuotaHoldDescriptor + | SettingsHoldDescriptor + | BinIndexHoldDescriptor; + /** * One direct child in a snapshot, as data. `size`/`mtime`/`contentVersion` are * `null` until projected. @@ -177,12 +186,8 @@ export interface SnapshotDescriptor { children: SnapshotChildDescriptor[]; ancestors: BreadcrumbDescriptor[]; deadLetters: DeadLetterDescriptor[]; - /** The drain's over-budget hold, or `null` when nothing is held. */ - blocked: BlockedOpDescriptor | null; - /** The drain's settings-refused hold, or `null` when nothing is held. */ - settingsHold: SettingsHoldDescriptor | null; - /** The drain's bin-index-refused hold, or `null` when nothing is held. */ - binIndexHold: BinIndexHoldDescriptor | null; + /** The drain's held queue head, or `null` when nothing is held. */ + queueHold: QueueHoldDescriptor | null; /** * Durable queue entries this session holds but cannot read — another * identity's, or written by a newer build. They occupy staged bytes against From a84b9550b924223f0179b4e00ddf4be5c2cfc543 Mon Sep 17 00:00:00 2001 From: Michael Yankelev Date: Tue, 15 Sep 2026 03:11:36 +0200 Subject: [PATCH 2/2] fix: let a quota hold go when the settings refuse the placement A head held over the account quota re-probes the quota on every tick, and a placement the session could not decide made that probe answer "no room" without ever asking. The member then read "over quota" while the cause was their own settings, and the label stood until the settings were fixed. The probe now separates the two placements it cannot use. A refusal of the member's own settings is a verdict, so the hold goes and the pass re-takes the hold its own rule names. An outage is not a verdict: the head stays held rather than running a pass that would spend the unattributed budget on a placement no pass can decide. --- crates/engine/src/sync/drain.rs | 46 ++++++++++++++++++++++++++++++++- 1 file changed, 45 insertions(+), 1 deletion(-) diff --git a/crates/engine/src/sync/drain.rs b/crates/engine/src/sync/drain.rs index eee74f2286..30e129c514 100644 --- a/crates/engine/src/sync/drain.rs +++ b/crates/engine/src/sync/drain.rs @@ -1928,7 +1928,12 @@ where /// probe is the hold's only exit, so an unanswered one leaves it in place. async fn quota_admits(&self, needed_bytes: u64) -> bool { let Ok(placement) = self.placement.as_ref() else { - return false; + // A placement the settings themselves refuse is a verdict the pass + // re-takes as its own hold, so the head stops waiting under a cause + // the member cannot act on. An outage is not a verdict: it keeps + // the head where it is rather than spending the unattributed budget + // on a placement no pass can decide. + return settings_refusal(self.placement).is_some(); }; // Only the hosted leg is quota-gated, so no answer the quota endpoint // could give bears on a hold under a placement without one — and an @@ -7798,6 +7803,45 @@ mod tests { } } + /// A quota hold's exit is a probe, and a placement the session cannot use + /// answers that probe in two different ways. A refusal of the member's own + /// settings is a verdict: the pass re-takes the hold its own rule names, so + /// the member is not left reading "over quota" over a cause they cannot act + /// on. An outage is not a verdict, and letting the head run would spend the + /// unattributed budget on a placement no pass can decide. + #[test] + fn a_quota_hold_lets_go_of_a_refusing_placement_and_waits_out_an_undecided_one() { + let node = NodeId([9; 16]); + let queued = vec![(OpId(1), Op::rename(node, "renamed.txt", 1, UnixMillis(0)))]; + let over_quota = QueueHold { + op_id: OpId(1), + node, + reason: QueueHoldReason::Quota { needed_bytes: 4096 }, + }; + + let mut refusing = drain_harness(None); + refusing.placement = Err(PlacementRefusal::NoProvider); + *refusing.hold.borrow_mut() = Some(over_quota); + assert!(block_on(refusing.drain().hold_admits_the_head(&queued))); + assert_eq!( + *refusing.hold.borrow(), + None, + "the settings refusal is what the next pass names", + ); + + let mut undecided = drain_harness(None); + undecided.placement = Err(PlacementRefusal::SettingsUnavailable( + DefaultsReason::Suppressed, + )); + *undecided.hold.borrow_mut() = Some(over_quota); + assert!(!block_on(undecided.drain().hold_admits_the_head(&queued))); + assert_eq!( + *undecided.hold.borrow(), + Some(over_quota), + "an outage leaves the head held rather than charging it", + ); + } + /// What a run of halted passes over one queued op left behind. struct Halted { report: DrainReport,