diff --git a/crates/desktop-seams/src/fs_util.rs b/crates/desktop-seams/src/fs_util.rs index b6216320c5..e21edc5b33 100644 --- a/crates/desktop-seams/src/fs_util.rs +++ b/crates/desktop-seams/src/fs_util.rs @@ -167,12 +167,12 @@ pub(crate) fn keep_first(kept: Result<(), E>, next: Result<(), E>) -> Result< /// Barriers a directory so a create/rename/unlink inside it is durable /// before the next one is issued. #[cfg(unix)] -fn fsync_dir(dir: &Path) -> io::Result<()> { +pub(crate) fn fsync_dir(dir: &Path) -> io::Result<()> { File::open(dir)?.sync_all() } #[cfg(not(unix))] -fn fsync_dir(dir: &Path) -> io::Result<()> { +pub(crate) fn fsync_dir(dir: &Path) -> io::Result<()> { metadata_log_barrier(dir) } diff --git a/crates/desktop-seams/src/staging_store.rs b/crates/desktop-seams/src/staging_store.rs index 214ea3c589..454072bfe6 100644 --- a/crates/desktop-seams/src/staging_store.rs +++ b/crates/desktop-seams/src/staging_store.rs @@ -1,20 +1,25 @@ //! Desktop [`StagingStore`]: the v1 write journal generalized to every op. use std::collections::BTreeMap; +use std::ops::RangeInclusive; use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex}; use cipherbox_engine::seams::{OpId, SeamResult, StagingStore}; use crate::fs_util::{ - atomic_write, empty_dir, ensure_dir, from_hex, keep_first, list_file_names, read_file_opt, - remove_file_durable, seam_err, to_hex, + atomic_write, empty_dir, ensure_dir, from_hex, fsync_dir, keep_first, list_file_names, + read_file_opt, remove_file_durable, seam_err, to_hex, }; /// Suffix for op-record files (`ops/.op`). The "json record" of the v1 /// journal, generalized: it holds the engine's opaque encoded intent op /// verbatim — the store never parses it. const OP_SUFFIX: &str = ".op"; +/// Suffix for a multi-entry enqueue's marker (`ops/-.batch`, +/// empty). While it stands, the op files in its range are not queued, and +/// a reopen removes them: the marker's removal is the set's commit point. +const BATCH_SUFFIX: &str = ".batch"; /// Suffix for staged-ciphertext sidecar files (`staged/.bin`). const SIDECAR_SUFFIX: &str = ".bin"; /// Filename of the durable monotonic op-id counter. @@ -29,6 +34,8 @@ const COUNTER_FILE: &str = "next_op_id"; /// - `ops/<20-digit-id>.op` — one durable record per queued op, id in the /// filename; enqueue order is id order (FIFO). Each op is opaque engine /// bytes stored verbatim. +/// - `ops/-.batch` — the marker of a multi-entry enqueue that has +/// not committed ([`BATCH_SUFFIX`]). /// - `staged/.bin` — one sidecar per staged-ciphertext key. /// - `next_op_id` — the monotonic id counter, so ids are strictly /// increasing and never reused, even after every op drains and the store @@ -78,17 +85,90 @@ impl FileStagingStore { .and_then(|bytes| <[u8; 8]>::try_from(bytes.as_slice()).ok()) .map(u64::from_le_bytes) .unwrap_or(0); - let highest_op = highest_op_id(&ops_dir) - .map_err(|err| seam_err("staging_store open scan", &err))? + let names = + list_file_names(&ops_dir).map_err(|err| seam_err("staging_store open scan", &err))?; + let highest_op = names + .iter() + .filter_map(|name| parse_op_id(name)) + .max() .map_or(0, |id| id.saturating_add(1)); let next = persisted.max(highest_op).max(1); - Ok(Self { + let store = Self { ops_dir, staged_dir, counter_path, next_op_id: Arc::new(Mutex::new(next)), - }) + }; + for batch in open_batches(&names) { + // Only the ids on disk: a marker's range is not bounded by its set. + let written = names + .iter() + .filter_map(|name| parse_op_id(name)) + .filter(|id| batch.contains(id)); + store + .roll_back_batch(&batch, written) + .map_err(|err| seam_err("staging_store open rollback", &err))?; + } + Ok(store) + } + + fn batch_path(&self, batch: &RangeInclusive) -> PathBuf { + let (first, last) = (batch.start(), batch.end()); + self.ops_dir + .join(format!("{first:020}-{last:020}{BATCH_SUFFIX}")) + } + + /// Removes every op file of an uncommitted set, past a refused removal, + /// then its marker only if all of them went, so a marker that stands hides + /// what is left. An unlink whose barrier refused still takes the file out + /// of the directory, so with no marker the set is gone from this process. + fn roll_back_batch( + &self, + batch: &RangeInclusive, + written: impl Iterator, + ) -> std::io::Result<()> { + let mut removed = Ok(()); + for id in written { + removed = keep_first(removed, remove_file_durable(&self.op_path(id))); + } + removed?; + remove_file_durable(&self.batch_path(batch)) + } + + /// Reserves `count` consecutive ids and durably advances the counter past + /// them before any op file is written: a crash after the bump burns ids, + /// never reuses one. + fn reserve_ids(&self, count: u64) -> SeamResult { + let mut next = self.next_op_id.lock().expect("lock"); + let first = *next; + let advanced = first + count; + atomic_write(&self.counter_path, &advanced.to_le_bytes()) + .map_err(|err| seam_err("staging_store reserve ids", &err))?; + *next = advanced; + Ok(first) + } + + fn write_batch(&self, first: u64, ops: &[Vec]) -> std::io::Result<()> { + let batch = first..=first + ops.len() as u64 - 1; + atomic_write(&self.batch_path(&batch), &[])?; + for (id, op) in (first..).zip(ops) { + if let Err(err) = atomic_write(&self.op_path(id), op) { + // The failed write may have landed its file before a barrier + // refused. A marker left by a failed rollback still hides the + // set, and the next open removes it. + let _ = self.roll_back_batch(&batch, first..=id); + return Err(err); + } + } + if let Err(err) = std::fs::remove_file(self.batch_path(&batch)) { + let _ = self.roll_back_batch(&batch, batch.clone()); + return Err(err); + } + // The unlink commits the whole set; a failed barrier is `Err` with the + // durability unknown, as for `atomic_write`. + fsync_dir(&self.ops_dir) + .map_err(|err| std::io::Error::new(err.kind(), format!("commit barrier: {err}"))) } fn op_path(&self, id: u64) -> PathBuf { @@ -103,27 +183,26 @@ impl FileStagingStore { impl StagingStore for FileStagingStore { async fn enqueue_op(&self, op: &[u8]) -> SeamResult { - // Reserve an id and durably advance the counter *before* writing the - // op file: a crash after the counter bump merely burns an id (ids - // need only be strictly increasing, not contiguous), never reuses - // one. - let id = { - let mut next = self.next_op_id.lock().expect("lock"); - let id = *next; - let advanced = id + 1; - atomic_write(&self.counter_path, &advanced.to_le_bytes()) - .map_err(|err| seam_err("staging_store enqueue_op counter", &err))?; - *next = advanced; - id - }; + let id = self.reserve_ids(1)?; atomic_write(&self.op_path(id), op) .map_err(|err| seam_err("staging_store enqueue_op", &err))?; Ok(OpId(id)) } + async fn enqueue_ops(&self, ops: &[Vec]) -> SeamResult> { + if ops.is_empty() { + return Ok(Vec::new()); + } + let first = self.reserve_ids(ops.len() as u64)?; + self.write_batch(first, ops) + .map_err(|err| seam_err("staging_store enqueue_ops", &err))?; + Ok((first..).take(ops.len()).map(OpId).collect()) + } + async fn queued_ops(&self) -> SeamResult)>> { let names = list_file_names(&self.ops_dir) .map_err(|err| seam_err("staging_store queued_ops", &err))?; + let uncommitted = open_batches(&names); // Keyed by parsed id, not by listed name: `op_path` zero-pads and // `parse_op_id` does not, so `ops/1.op` and `ops/00…01.op` both name op 1. // Reading the re-derived canonical path keeps every returned entry one @@ -134,7 +213,7 @@ impl StagingStore for FileStagingStore { let Some(id) = parse_op_id(&name) else { continue; }; - if ops.contains_key(&id) { + if ops.contains_key(&id) || uncommitted.iter().any(|batch| batch.contains(&id)) { continue; } if let Some(bytes) = read_file_opt(&self.op_path(id)) @@ -221,10 +300,14 @@ fn parse_op_id(name: &str) -> Option { name.strip_suffix(OP_SUFFIX)?.parse::().ok() } -/// The largest op id currently on disk in `ops_dir`, if any. -fn highest_op_id(ops_dir: &Path) -> std::io::Result> { - Ok(list_file_names(ops_dir)? +/// The id range of every multi-entry enqueue whose marker still stands in +/// `names`. +fn open_batches(names: &[String]) -> Vec> { + names .iter() - .filter_map(|name| parse_op_id(name)) - .max()) + .filter_map(|name| { + let (first, last) = name.strip_suffix(BATCH_SUFFIX)?.split_once('-')?; + Some(first.parse().ok()?..=last.parse().ok()?) + }) + .collect() } diff --git a/crates/desktop-seams/tests/conformance.rs b/crates/desktop-seams/tests/conformance.rs index 5206ec53e7..690464b3e4 100644 --- a/crates/desktop-seams/tests/conformance.rs +++ b/crates/desktop-seams/tests/conformance.rs @@ -58,7 +58,10 @@ fn file_staging_store_passes_the_staging_store_kit() { block_on(conformance::staging_store::check( async |backing: Backing| FileStagingStore::open(root.join(backing.label())).unwrap(), async |backing: Backing| match backing { - Backing::Ordering | Backing::FailedReplacement | Backing::Cleared => denial.arm(), + Backing::Ordering + | Backing::FailedReplacement + | Backing::Cleared + | Backing::Batched => denial.arm(), Backing::FailedFirstPut => { std::fs::remove_dir(root.join(backing.label()).join("staged")) .expect("the kit's lever must be armed, or it proves nothing"); @@ -379,6 +382,77 @@ fn staging_store_op_ids_never_reuse_across_a_full_drain() { }); } +/// A multi-entry enqueue whose second entry will not land queues neither, now +/// or after a reopen, and burns the ids it reserved. +#[test] +fn staging_store_a_set_that_fails_part_way_queues_none_of_it() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("staging"); + block_on(async { + let store = FileStagingStore::open(&path).unwrap(); + let single = store.enqueue_op(b"single").await.unwrap(); + // A directory at the second entry's path refuses its atomic rename. + let blocked = path.join("ops").join(format!("{:020}.op", single.0 + 2)); + std::fs::create_dir(&blocked).unwrap(); + + store + .enqueue_ops(&[b"park".to_vec(), b"arrive".to_vec()]) + .await + .expect_err("the second entry cannot land"); + let queued = vec![(single, b"single".to_vec())]; + assert_eq!( + store.queued_ops().await.unwrap(), + queued, + "nor may the first" + ); + + std::fs::remove_dir(&blocked).unwrap(); + let reopened = FileStagingStore::open(&path).unwrap(); + assert_eq!( + reopened.queued_ops().await.unwrap(), + queued, + "not after a reopen either" + ); + let next = reopened.enqueue_op(b"next").await.unwrap(); + assert!( + next.0 > single.0 + 2, + "the set's ids are burnt, never reused" + ); + }); +} + +/// A crash between a set's entries leaves its marker standing: the entries it +/// covers are not queued, and the next open removes them. +#[test] +fn staging_store_a_set_a_crash_interrupted_is_rolled_back_at_open() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("staging"); + let ops = path.join("ops"); + block_on(async { + let store = FileStagingStore::open(&path).unwrap(); + let single = store.enqueue_op(b"single").await.unwrap(); + let (first, last) = (single.0 + 1, single.0 + 2); + let marker = ops.join(format!("{first:020}-{last:020}.batch")); + let parked = ops.join(format!("{first:020}.op")); + std::fs::write(&marker, b"").unwrap(); + std::fs::write(&parked, b"park").unwrap(); + + let queued = vec![(single, b"single".to_vec())]; + assert_eq!( + store.queued_ops().await.unwrap(), + queued, + "an uncommitted set is not queued" + ); + + let reopened = FileStagingStore::open(&path).unwrap(); + assert_eq!(reopened.queued_ops().await.unwrap(), queued); + assert!( + !marker.exists() && !parked.exists(), + "the open removed the set and its marker" + ); + }); +} + /// A clone hands out ids from the same counter. The engine's cold start clones /// this seam into its spawned loops, so two handles that each counted for /// themselves would give one id to two ops. diff --git a/crates/engine/src/facade.rs b/crates/engine/src/facade.rs index 2c8a130c45..67488869b0 100644 --- a/crates/engine/src/facade.rs +++ b/crates/engine/src/facade.rs @@ -173,7 +173,8 @@ use crate::sync::refresh::ManualRefresh; use crate::sync::staging::{ DEAD_LETTER_NOTICES_PREFIX, DroppedVersionDebts, PreservedBounds, PreservedDeadLetter, StagedBlocks, read_dead_letter_notices, read_preserved_dead_letters, reconcile_staging, - release_version_blocks, stage_op, take_dead_letter_notice, take_preserved_dead_letter, + release_version_blocks, stage_op, staged_record, take_dead_letter_notice, + take_preserved_dead_letter, }; use crate::sync::staleness::{Connectivity, classify, next_boundary}; use crate::sync::tick::{ @@ -10626,37 +10627,41 @@ where { /// Journal the legs one command owes, in order, and report the id of the /// last — the op whose publish completes it. /// - /// A leg that will not journal takes the leg before it back off the queue: - /// a staged crossing left half-journaled would park the subtree in the - /// vault-root scope ([`relocation_legs`]), which is neither the move the - /// caller asked for nor the failure they were told about. + /// Two legs go in one atomic write: a staged crossing left half-journaled + /// would park the subtree in the vault-root scope ([`relocation_legs`]), + /// which is neither the move the caller asked for nor the failure they + /// were told about (ADR 0045 D6). async fn stage_legs_and_notify( &mut self, park: Option<&Op>, arrive: &Op, ) -> Result { - let parked = match park { - Some(op) => Some(self.journal(op).await?), - None => None, - }; - match self.journal(arrive).await { - Ok(op_id) => { - // Best-effort push-invalidation trigger; a dropped receiver - // (host torn down) is fine. - let _ = self.events.unbounded_send(Event::SnapshotUpdated); - Ok(CommandOutcome::Queued { op_id }) - } - Err(error) => { - // Best-effort, and the caller still hears why the command - // failed: a leg a tick has already drained is gone from the - // queue, and a cleanup that reported itself instead would name - // the wrong cause. - if let Some(op_id) = parked { - let _ = self.dequeue_op(op_id).await; - } - Err(error) + let op_id = match park { + None => self.journal(arrive).await?, + Some(park) => { + let legs = [self.sealed(park).await?, self.sealed(arrive).await?]; + let ids = self + .seams + .staging_store + .enqueue_ops(&legs) + .await + .map_err(EngineError::from_seam)?; + *ids.last().ok_or_else(|| { + EngineError::from_seam(SeamError::new("enqueue_ops returned no ids")) + })? } - } + }; + // Best-effort push-invalidation trigger; a dropped receiver (host torn + // down) is fine. + let _ = self.events.unbounded_send(Event::SnapshotUpdated); + Ok(CommandOutcome::Queued { op_id }) + } + + /// Seal one op into a durable record, under an ephemeral of its own. + async fn sealed(&self, op: &Op) -> Result, EngineError> { + staged_record(&self.seams.staging_store, self.record_seal()?, op) + .await + .map_err(EngineError::from_seam) } /// Seal one op onto the durable queue, under an ephemeral of its own. diff --git a/crates/engine/src/seams/live.rs b/crates/engine/src/seams/live.rs index 4ab12eaa59..dcdb040d45 100644 --- a/crates/engine/src/seams/live.rs +++ b/crates/engine/src/seams/live.rs @@ -118,6 +118,11 @@ impl StagingStore for LiveSeam { self.seam.enqueue_op(op).await } + async fn enqueue_ops(&self, ops: &[Vec]) -> SeamResult> { + self.writable("staging_store enqueue_ops")?; + self.seam.enqueue_ops(ops).await + } + async fn queued_ops(&self) -> SeamResult)>> { self.seam.queued_ops().await } @@ -226,6 +231,7 @@ mod tests { alive.set(false); seam.enqueue_op(b"op").await.unwrap_err(); + seam.enqueue_ops(&[b"op".to_vec()]).await.unwrap_err(); seam.put_staged_bytes(b"key", b"bytes").await.unwrap_err(); seam.remove_op(op).await.unwrap_err(); seam.remove_staged_bytes(b"key").await.unwrap_err(); diff --git a/crates/engine/src/seams/staging_store.rs b/crates/engine/src/seams/staging_store.rs index 2cc537d6cc..0d6a390a27 100644 --- a/crates/engine/src/seams/staging_store.rs +++ b/crates/engine/src/seams/staging_store.rs @@ -30,9 +30,18 @@ pub struct OpId(pub u64); /// Hosts: IndexedDB + OPFS (web), local journal (desktop). pub trait StagingStore { /// Appends an opaque op record to the durable FIFO queue and returns - /// its id. + /// its id. An `Err` after the commit point can leave the op queued, with + /// its durability unknown. async fn enqueue_op(&self, op: &[u8]) -> SeamResult; + /// Appends several op records to the queue in one atomic write, in order, + /// and returns their ids in that order. No reader, before or after a crash, + /// ever sees part of the set. An `Err` leaves none of them queued, except + /// an `Err` after the commit point, which, as for [`Self::enqueue_op`], can + /// leave the whole set queued with its durability unknown. Each entry is stored as [`Self::enqueue_op`] stores one, so a + /// reader of the queue cannot tell the two apart. + async fn enqueue_ops(&self, ops: &[Vec]) -> SeamResult>; + /// Every queued op in FIFO (ascending-id) order. async fn queued_ops(&self) -> SeamResult)>>; @@ -97,7 +106,7 @@ pub trait StagingStore { /// counted beside the command path's enqueues; the engine wraps the host's /// store once, and every handle it hands out is a clone of that one. /// -/// Only the three methods that change which ops are queued count. Staged bytes +/// Only the methods that change which ops are queued count. Staged bytes /// are not an operand of the state law, and a live write handle churns them. pub struct QueueGenerationStore { seam: S, @@ -158,6 +167,11 @@ impl StagingStore for QueueGenerationStore { self.seam.enqueue_op(op).await } + async fn enqueue_ops(&self, ops: &[Vec]) -> SeamResult> { + self.mutating(); + self.seam.enqueue_ops(ops).await + } + async fn queued_ops(&self) -> SeamResult)>> { self.seam.queued_ops().await } @@ -212,11 +226,14 @@ mod tests { let enqueued = store.generation(); block_on(store.remove_op(op)).expect("remove"); let removed = store.generation(); + block_on(store.enqueue_ops(&[b"a".to_vec(), b"b".to_vec()])).expect("enqueue a set"); + let batched = store.generation(); block_on(store.clear()).expect("clear"); assert_ne!(enqueued, start); assert_ne!(removed, enqueued); - assert_ne!(store.generation(), removed); + assert_ne!(batched, removed); + assert_ne!(store.generation(), batched); } #[test] diff --git a/crates/engine/src/sync/drain.rs b/crates/engine/src/sync/drain.rs index b5c68b6c92..5674e2a5e1 100644 --- a/crates/engine/src/sync/drain.rs +++ b/crates/engine/src/sync/drain.rs @@ -1288,6 +1288,22 @@ struct FolderState { sequence: u64, } +/// The `modified_at` a plan republishes `folder` with: the op's authored time on +/// one of its [authored nodes](Op::authored_nodes), which the overlay stamps +/// from the same set, and the folder's own time otherwise. +fn stamped_modified_at( + pass: &Pass, + op: &Op, + authored: &[NodeId], + folder: NodeId, +) -> Result { + if authored.contains(&folder) { + Ok(op.authored_at.0) + } else { + Ok(pass.folder(folder)?.modified_at) + } +} + /// Where one child ref is going, under what name, and what it displaces — /// rename, relink, and move all reduce to this. struct MovePlan { @@ -2932,15 +2948,11 @@ where // Referent published: only now does the parent gain the ref to it. pass.folder_mut(parent)?.children.push(child.child_ref); - self.publish_folder( - scope, - pass, - parent, - applied.op.authored_at.0, - Some(applied.op_id), - ) - .await - .map_err(Halt::from)?; + let authored = applied.op.authored_nodes(Vec::new); + let modified_at = stamped_modified_at(pass, &applied.op, &authored, parent)?; + self.publish_folder(scope, pass, parent, modified_at, Some(applied.op_id)) + .await + .map_err(Halt::from)?; self.release_staged_blocks(&applied.op).await; self.emit_mirror_shortfall(applied, shortfall); // The parent's repaint lifts the child in without what its own record @@ -2955,6 +2967,8 @@ where 1, Some(&staged.root_cid), ); + } else if let Some(meta) = self.cells.base.borrow_mut().node_mut(child_id) { + applied.op.stamp_authored(meta); } // Held only once the parent names it: a record nothing references is // not one the liveness loop should keep alive. @@ -3077,7 +3091,9 @@ where // leaves the node binned and still linked, which is the residue the // entry-before-unlink order already settles on the retry. let count = unlink_from.len(); + let authored = applied.op.authored_nodes(|| unlink_from.clone()); for (at, parent) in unlink_from.into_iter().enumerate() { + let modified_at = stamped_modified_at(pass, &applied.op, &authored, parent)?; pass.folder_mut(parent)? .children .retain(|entry| entry.id != target.0); @@ -3085,7 +3101,7 @@ where scope, pass, parent, - applied.op.authored_at.0, + modified_at, (at + 1 == count).then_some(applied.op_id), ) .await @@ -3222,11 +3238,13 @@ where Some(existing) => *existing = child, None => into_children.push(child), } + let authored = applied.op.authored_nodes(Vec::new); + let modified_at = stamped_modified_at(pass, &applied.op, &authored, into)?; self.publish_folder( scope, pass, into, - applied.op.authored_at.0, + modified_at, // The entry drop below is this plan's last act, not the relink: a // mark raised here would drop the op on the next pass and leave the // entry standing for a node the vault links again. @@ -4490,8 +4508,6 @@ where if dest == target || self.cells.base.borrow().ancestors(dest).contains(&target) { return Err(Halt::Unclassified); } - let modified_at = applied.op.authored_at.0; - // A crossing this pass carries no second end for is one it cannot // author, and the chain walk below would stall uncharged on the scope // root it cannot load. Charged by the pass holding the tick's @@ -4576,6 +4592,10 @@ where // Only when one folder collapses the plan is the dest-add also its last // record; otherwise the source-remove below is. let single_record = source == dest; + // `source` is the base's winning parent, the one parent a rename or a + // relocation stamps. + let authored = applied.op.authored_nodes(|| vec![source]); + let modified_at = stamped_modified_at(pass, &applied.op, &authored, dest)?; let cas_base = self .publish_folder( scope, @@ -4599,8 +4619,9 @@ where // forever, so a quota refusal, a permanent one, or a spent attempt must // not be flattened into it. Only the undo's own failure is genuinely // unclassified. + let source_modified_at = stamped_modified_at(pass, &applied.op, &authored, source)?; if let Err(failure) = self - .publish_folder(scope, pass, source, modified_at, Some(applied.op_id)) + .publish_folder(scope, pass, source, source_modified_at, Some(applied.op_id)) .await { // A confirmed source-remove is the move complete on the network, so diff --git a/crates/engine/src/sync/op.rs b/crates/engine/src/sync/op.rs index 0ea8f67ce0..5a52c63fc5 100644 --- a/crates/engine/src/sync/op.rs +++ b/crates/engine/src/sync/op.rs @@ -652,14 +652,42 @@ impl Op { self.staged_content().map(|c| &c.root_cid[..]) } - /// Stamp this op's authored facts onto the node it targets — `mtime` - /// **overwriting** the projected time, and a content op's plaintext size. - /// The one function the pending-op overlay and the drain's publish plan - /// share, so a rendered node and the record that will publish it agree - /// (blueprint/engine.md "State law"). + /// The nodes whose next records this op's publish authors at + /// `authored_at`, sorted: the set the pending-op overlay stamps and the + /// drain's publish plan writes, so a rendered node and the record that + /// will publish it agree (ADR 0045 D5). `parents` yields the folders that + /// name the target before the op applies, winner first + /// ([`Snapshot::links_ranked`](crate::sync::model::Snapshot::links_ranked)); + /// only a kind that stamps them calls it. + pub fn authored_nodes(&self, parents: impl FnOnce() -> Vec) -> Vec { + let mut nodes = match &self.kind { + OpKind::Create { parent, .. } => vec![self.target, *parent], + OpKind::UpdateContent { .. } => vec![self.target], + // The name lives in the parent's child ref, so the child's own + // record does not change. A rename and a relocation republish the + // winning parent alone, and a delete unlinks from every parent. + OpKind::Rename { .. } => parents().into_iter().take(1).collect(), + OpKind::Delete { .. } => parents(), + OpKind::Relink { new_parent, .. } | OpKind::Move { new_parent, .. } => { + parents().into_iter().take(1).chain([*new_parent]).collect() + } + OpKind::Restore { into, .. } => vec![*into], + OpKind::Purge { .. } + | OpKind::Prune { .. } + | OpKind::RestoreVersion { .. } + | OpKind::DeleteVersion { .. } => Vec::new(), + }; + nodes.sort(); + nodes.dedup(); + nodes + } + + /// Stamp this op's authored facts onto one of its [authored + /// nodes](Self::authored_nodes): `mtime` **overwriting** the projected + /// time, and, on the target of a content op, its plaintext size. pub fn stamp_authored(&self, meta: &mut NodeMeta) { meta.mtime = Some(self.authored_at.0); - if let Some(content) = self.staged_content() { + if let Some(content) = self.staged_content().filter(|_| meta.id == self.target) { meta.size = Some(content.plaintext_size); } } @@ -956,16 +984,96 @@ mod tests { #[test] fn a_metadata_op_stamps_time_over_a_projection_and_leaves_size_alone() { - let mut node = NodeMeta::new(id(1), "f.txt", NodeKind::File); - node.mtime = Some(999); - node.size = Some(42); - Op::rename(id(1), "g.txt", 1, at(1)).stamp_authored(&mut node); + let mut parent = NodeMeta::new(id(0), "docs", NodeKind::Folder); + parent.mtime = Some(999); + parent.size = Some(42); + Op::rename(id(1), "g.txt", 1, at(1)).stamp_authored(&mut parent); assert_eq!( - node.mtime, + parent.mtime, Some(1), - "the op authors the node's next record, so the projected time is stale" + "the op authors the parent's next record, so the projected time is stale" + ); + assert_eq!(parent.size, Some(42), "a metadata op carries no size"); + } + + #[test] + fn a_content_op_stamps_its_size_on_its_target_alone() { + let mut parent = NodeMeta::new(id(0), "docs", NodeKind::Folder); + let create = NewNode::File { + content: Some(staged(b"root", 9)), + }; + Op::create(id(1), id(0), "f.txt", create, 1, at(5)).stamp_authored(&mut parent); + assert_eq!((parent.mtime, parent.size), (Some(5), None)); + } + + /// ADR 0045 D5. + #[test] + fn each_op_kind_authors_its_own_set_of_nodes() { + let (node, from, to) = (id(1), id(2), id(3)); + let parents = [from]; + for (op, expected, why) in [ + ( + Op::create(node, from, "a", NewNode::Folder, 1, at(1)), + vec![node, from], + "create", + ), + (Op::rename(node, "b", 1, at(1)), vec![from], "rename"), + (Op::delete(node, 1, at(1), 1, false), vec![from], "delete"), + ( + Op::relink(node, from, to, 1, at(1), ScopeCrossing::Intra), + vec![from, to], + "relink", + ), + ( + Op::move_node(node, from, to, "b", None, 1, at(1), ScopeCrossing::Intra), + vec![from, to], + "move", + ), + ( + Op::update_content(node, staged(b"k", 1), None, 1, at(1)), + vec![node], + "update content", + ), + ( + Op::restore(node, to, "a", NodeKind::Folder, 1, at(1)), + vec![to], + "restore", + ), + ] { + assert_eq!(op.authored_nodes(|| parents.to_vec()), expected, "{why}"); + } + assert_eq!( + Op::rename(node, "b", 1, at(1)).authored_nodes(Vec::new), + Vec::new(), + "a node no folder names has no parent to stamp" ); - assert_eq!(node.size, Some(42), "a metadata op carries no size"); + + // Two links, winner first: the winner sorts after the other parent. + let (winner, other) = (id(4), id(2)); + for (op, expected, why) in [ + ( + Op::rename(node, "b", 1, at(1)), + vec![winner], + "a rename keeps the winning parent alone", + ), + ( + Op::move_node(node, winner, to, "b", None, 1, at(1), ScopeCrossing::Intra), + vec![to, winner], + "a move keeps the winning parent and the destination", + ), + ( + Op::relink(node, winner, to, 1, at(1), ScopeCrossing::Intra), + vec![to, winner], + "a relink keeps the winning parent and the destination", + ), + ( + Op::delete(node, 1, at(1), 1, false), + vec![other, winner], + "a delete keeps every parent", + ), + ] { + assert_eq!(op.authored_nodes(|| vec![winner, other]), expected, "{why}"); + } } #[test] diff --git a/crates/engine/src/sync/overlay.rs b/crates/engine/src/sync/overlay.rs index c57602914e..c5c25ea70f 100644 --- a/crates/engine/src/sync/overlay.rs +++ b/crates/engine/src/sync/overlay.rs @@ -24,11 +24,24 @@ pub fn apply_overlay(base: &Snapshot, ops: &[Op]) -> Snapshot { view } -/// Apply one op to the working view optimistically (intent, not rebase). -/// -/// An op that authors a *change* to its target stamps [`Op::stamp_authored`]; -/// the ops that only remove state — a delete, a prune — have nothing to stamp. +/// Apply one op to the working view optimistically (intent, not rebase), then +/// stamp its [authored nodes](Op::authored_nodes). fn apply_one(view: &mut Snapshot, op: &Op) { + let authored = op.authored_nodes(|| { + view.links_ranked(op.target) + .iter() + .map(|link| link.parent) + .collect() + }); + apply_intent(view, op); + for node in authored { + if let Some(meta) = view.node_mut(node) { + op.stamp_authored(meta); + } + } +} + +fn apply_intent(view: &mut Snapshot, op: &Op) { match &op.kind { OpKind::Create { parent, name, node } => { let mut meta = NodeMeta::new(op.target, name.clone(), node.kind()); @@ -36,7 +49,6 @@ fn apply_one(view: &mut Snapshot, op: &Op) { if op.staged_content().is_some() { meta.content_version = Some(1); } - op.stamp_authored(&mut meta); view.upsert_node(meta); view.link_next(*parent, op.target); } @@ -46,7 +58,6 @@ fn apply_one(view: &mut Snapshot, op: &Op) { OpKind::Rename { new_name } => { if let Some(node) = view.node_mut(op.target) { node.rename(new_name.clone()); - op.stamp_authored(node); } } OpKind::Relink { new_parent, .. } => relocate(view, op, *new_parent, None, None), @@ -65,15 +76,12 @@ fn apply_one(view: &mut Snapshot, op: &Op) { OpKind::UpdateContent { .. } => { if let Some(node) = view.node_mut(op.target) { node.content_version = node.content_version.map(|count| count + 1); - op.stamp_authored(node); } } // The node is binned, so nothing in the base renders it: the overlay // materializes it at the destination the command resolved. OpKind::Restore { into, name, kind } => { - let mut meta = NodeMeta::new(op.target, name.clone(), *kind); - op.stamp_authored(&mut meta); - view.upsert_node(meta); + view.upsert_node(NodeMeta::new(op.target, name.clone(), *kind)); view.link_next(*into, op.target); } // A purged node is already absent from the rendered view; the bin is @@ -122,11 +130,8 @@ fn relocate( }) }); view.relocate(op.target, new_parent, vacating); - if let Some(node) = view.node_mut(op.target) { - if let Some(new_name) = new_name { - node.rename(new_name); - } - op.stamp_authored(node); + if let (Some(node), Some(new_name)) = (view.node_mut(op.target), new_name) { + node.rename(new_name); } } diff --git a/crates/engine/src/sync/staging.rs b/crates/engine/src/sync/staging.rs index 4d64024d39..b256115314 100644 --- a/crates/engine/src/sync/staging.rs +++ b/crates/engine/src/sync/staging.rs @@ -95,6 +95,19 @@ pub async fn stage_op( seal: RecordSeal<'_>, op: &Op, ) -> SeamResult { + store + .enqueue_op(&staged_record(store, seal, op).await?) + .await +} + +/// The sealed record [`stage_op`] journals for `op`, after every check it +/// makes, for a caller that journals several records in one atomic write +/// ([`StagingStore::enqueue_ops`]). +pub(crate) async fn staged_record( + store: &S, + seal: RecordSeal<'_>, + op: &Op, +) -> SeamResult> { if !op.crossing_is_coherent() { return Err(SeamError::new( "stage_op: a relocation that keeps its parent cannot claim to leave its scope", @@ -109,7 +122,7 @@ pub async fn stage_op( SeamError::new("stage_op: staged root block does not address to the op's content root") })?; } - store.enqueue_op(&record).await + Ok(record) } /// The staging key holding the **op records** of dead letters whose staged bytes diff --git a/crates/engine/src/testkit/conformance/staging_store.rs b/crates/engine/src/testkit/conformance/staging_store.rs index 183cf3ca07..6f378dc15c 100644 --- a/crates/engine/src/testkit/conformance/staging_store.rs +++ b/crates/engine/src/testkit/conformance/staging_store.rs @@ -1,5 +1,5 @@ //! Conformance kit: [`StagingStore`] FIFO ordering, durability, orphan-GC -//! support, and `put_staged_bytes` failure atomicity. +//! support, multi-entry enqueue, and `put_staged_bytes` failure atomicity. use crate::seams::StagingStore; @@ -25,17 +25,20 @@ pub enum Backing { FailedFirstPut, /// The "forget this device" erase. Cleared, + /// The multi-entry enqueue. + Batched, } impl Backing { /// Every backing the kit asks for, so a host that has to enumerate them /// (one behind a string boundary, say) reads the set off the kit rather /// than transcribing it. - pub const ALL: [Self; 4] = [ + pub const ALL: [Self; 5] = [ Self::Ordering, Self::FailedReplacement, Self::FailedFirstPut, Self::Cleared, + Self::Batched, ]; /// A stable label a host can key a directory, database name, or map entry @@ -46,6 +49,7 @@ impl Backing { Self::FailedReplacement => "failed-replacement", Self::FailedFirstPut => "failed-first-put", Self::Cleared => "cleared", + Self::Batched => "batched", } } } @@ -69,6 +73,70 @@ where failed_replacement_put(&mut open, &mut arm_failed_put).await; failed_first_put(&mut open, &mut arm_failed_put).await; clearing(&mut open).await; + batched_enqueue(&mut open).await; +} + +/// A multi-entry enqueue queues its entries in order, verbatim, under ids that +/// continue the single-entry progression, and they survive reopen like any +/// other entry. +async fn batched_enqueue(open: &mut F) +where + S: StagingStore, + F: AsyncFnMut(Backing) -> S, +{ + let store = open(Backing::Batched).await; + assert!( + store.queued_ops().await.unwrap().is_empty() + && store.staged_keys().await.unwrap().is_empty(), + "every kit backing is its own, and starts empty" + ); + assert_eq!( + store.enqueue_ops(&[]).await.unwrap(), + Vec::new(), + "an empty set queues nothing" + ); + let single = store.enqueue_op(b"op-single").await.unwrap(); + let batch = store + .enqueue_ops(&[b"op-park".to_vec(), b"op-arrive".to_vec()]) + .await + .unwrap(); + let [park, arrive] = batch[..] else { + panic!("a multi-entry enqueue must return one id per entry"); + }; + assert!( + single < park && park < arrive, + "a set's ids continue the progression, in the set's order" + ); + let after = store.enqueue_op(b"op-after").await.unwrap(); + assert!(after > arrive, "and the progression continues past the set"); + + let expected = vec![ + (single, b"op-single".to_vec()), + (park, b"op-park".to_vec()), + (arrive, b"op-arrive".to_vec()), + (after, b"op-after".to_vec()), + ]; + assert_eq!( + store.queued_ops().await.unwrap(), + expected, + "a set's entries queue FIFO with payloads verbatim" + ); + assert_eq!( + open(Backing::Batched).await.queued_ops().await.unwrap(), + expected, + "and survive reopen" + ); + + store.remove_op(park).await.unwrap(); + assert_eq!( + store.queued_ops().await.unwrap(), + vec![ + (single, b"op-single".to_vec()), + (arrive, b"op-arrive".to_vec()), + (after, b"op-after".to_vec()), + ], + "an entry of a set is removed alone" + ); } /// FIFO ordering, id progression, staged-byte accounting, orphan-GC support, diff --git a/crates/engine/src/testkit/fakes/staging_store.rs b/crates/engine/src/testkit/fakes/staging_store.rs index cac7128106..3ecc700b13 100644 --- a/crates/engine/src/testkit/fakes/staging_store.rs +++ b/crates/engine/src/testkit/fakes/staging_store.rs @@ -73,9 +73,10 @@ impl InMemoryStagingStore { self.inner.lock().expect("lock").fail_remove_op = true; } - /// Lets the next `budget` enqueues through and fails every one after, so a - /// test can drop a durable-queue outage in the middle of a multi-op - /// sequence and see what the earlier entries already committed. + /// Lets the next `budget` enqueued entries through and fails every one + /// after, so a test can drop a durable-queue outage in the middle of a + /// multi-op sequence and see what the earlier entries already committed. + /// A multi-entry enqueue the budget cannot cover fails whole. pub fn fail_enqueue_after(&self, budget: u64) { self.inner.lock().expect("lock").enqueue_budget = Some(budget); } @@ -208,6 +209,27 @@ impl StagingStore for InMemoryStagingStore { Ok(op_id) } + async fn enqueue_ops(&self, ops: &[Vec]) -> SeamResult> { + let mut inner = self.inner.lock().expect("lock"); + // The budget counts entries, and a set it cannot cover writes nothing. + let count = ops.len() as u64; + match inner.enqueue_budget { + Some(budget) if budget < count => { + return Err(SeamError::new("enqueue_ops unavailable")); + } + Some(budget) => inner.enqueue_budget = Some(budget - count), + None => {} + } + let mut ids = Vec::with_capacity(ops.len()); + for op in ops { + let op_id = OpId(inner.next_op_id); + inner.next_op_id += 1; + inner.ops.push((op_id, op.clone())); + ids.push(op_id); + } + Ok(ids) + } + async fn queued_ops(&self) -> SeamResult)>> { let mut inner = self.inner.lock().expect("lock"); inner.queue_listings += 1; diff --git a/crates/engine/tests/owner_actions.rs b/crates/engine/tests/owner_actions.rs index 8382dc519e..716def9d40 100644 --- a/crates/engine/tests/owner_actions.rs +++ b/crates/engine/tests/owner_actions.rs @@ -2952,6 +2952,29 @@ fn a_move_between_two_granted_folders_re_seals_into_the_destination_scope() { ); } +/// Both legs of a staged move are journaled or neither is (ADR 0045 D6). With +/// the arriving leg's journal refused and removal refused too, no cleanup can +/// take the parking leg back, so only a set written as one leaves nothing. +#[test] +fn a_staged_move_whose_arriving_leg_will_not_journal_journals_no_leg() { + let mut fx = GrantScenario::new(); + let (holiday, album) = two_granted_folders(&mut fx); + let staging = fx.owner_device.staging_store.inner(); + staging.fail_enqueue_after(1); + staging.fail_remove_op(); + + let refused = block_on(fx.engine.command(Command::Relink { + node: holiday, + new_parent: album, + })); + + assert!(refused.is_err(), "the caller hears that the move failed"); + assert!( + queued_crossings(&fx.owner_device).is_empty(), + "and no leg of it is queued" + ); +} + /// The legs are two durable ops, so a restart between them resumes at the one /// still queued. The cut belongs to the leg that already published, and the /// debt it settled is durable: the resumed leg publishes the arrival and cuts diff --git a/crates/engine/tests/write_plane.rs b/crates/engine/tests/write_plane.rs index 5da0ed07b1..6f7ac7214c 100644 --- a/crates/engine/tests/write_plane.rs +++ b/crates/engine/tests/write_plane.rs @@ -9,7 +9,7 @@ use core::cell::RefCell; use core::num::NonZeroU64; use core::task::{Context, Poll, Waker}; use core::time::Duration; -use std::collections::BTreeSet; +use std::collections::{BTreeMap, BTreeSet}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; @@ -16140,3 +16140,234 @@ fn a_write_over_a_history_retention_cannot_shorten_still_publishes() { "a co-writer's history never parks the member's own write", ); } + +/// The rendered `mtime` of every node in the view's tree. +fn rendered_mtimes(engine: &Engine) -> BTreeMap> { + let view = block_on(engine.view()).expect("a rendered view"); + let mut walk = vec![ROOT]; + let mut mtimes = BTreeMap::new(); + while let Some(node) = walk.pop() { + mtimes.insert(node, view.attrs(node).and_then(|attrs| attrs.mtime)); + walk.extend(view.children(node).into_iter().map(|child| child.id)); + } + mtimes +} + +/// Runs `act` and returns the nodes its overlay stamped: the ones now rendered +/// at its authored time. Every other node keeps its time, and the +/// publish moves no rendered time. +fn stamped_by( + world: &FakeWorld, + engine: &mut Engine, + tasks: &mut [BoxedTask], + act: impl FnOnce(&mut Engine), +) -> BTreeSet { + use cipherbox_engine::seams::Scheduler as _; + + tick(world, engine, tasks); + let before = rendered_mtimes(engine); + let authored_at = Some(world.scheduler.now().0); + act(engine); + let rendered = rendered_mtimes(engine); + for (node, mtime) in &rendered { + if let Some(prior) = before.get(node) { + assert!( + *mtime == authored_at || mtime == prior, + "a node the op does not stamp keeps its time" + ); + } + } + tick(world, engine, tasks); + assert_eq!( + rendered_mtimes(engine), + rendered, + "the publish writes the times the overlay rendered" + ); + rendered + .into_iter() + .filter(|(_, mtime)| *mtime == authored_at) + .map(|(node, _)| node) + .collect() +} + +/// `stamped_by` for one command. +fn stamped_by_command( + world: &FakeWorld, + engine: &mut Engine, + tasks: &mut [BoxedTask], + command: Command, +) -> BTreeSet { + stamped_by(world, engine, tasks, |engine| { + block_on(engine.command(command)).expect("the command journals"); + }) +} + +/// The overlay stamps exactly the nodes whose records the drain republishes at +/// the op's authored time (ADR 0045 D5): a rename or a relocation stamps the +/// folders, not the node. +#[test] +fn an_op_renders_the_times_its_publish_writes() { + let world = FakeWorld::new(); + let blocks = Blocks::default(); + seed_account(&world, &blocks); + let alice = world.device(b"alice"); + let (mut engine, _events, mut tasks) = boot_binning(&world, &blocks, &alice); + for name in ["a", "b"] { + block_on(engine.command(Command::Create { + parent: ROOT, + name: name.into(), + kind: NodeKind::Folder, + })) + .unwrap(); + } + tick(&world, &engine, &mut tasks); + let (a, b) = (child_id(&engine, ROOT, "a"), child_id(&engine, ROOT, "b")); + let created = stamped_by_command( + &world, + &mut engine, + &mut tasks, + Command::Create { + parent: a, + name: "x".into(), + kind: NodeKind::Folder, + }, + ); + let x = child_id(&engine, a, "x"); + assert_eq!( + created, + BTreeSet::from([a, x]), + "a create stamps the node and its parent" + ); + + let rename = Command::Rename { + node: x, + new_name: "y".into(), + }; + let relink = Command::Relink { + node: x, + new_parent: b, + }; + let relocate = Command::Move { + node: x, + new_parent: a, + new_name: "z".into(), + replacing: None, + }; + let restore = Command::Restore { + node: x, + into: None, + }; + for (command, expected, why) in [ + (rename, BTreeSet::from([a]), "a rename stamps the parent"), + ( + relink, + BTreeSet::from([a, b]), + "a relink stamps both parents", + ), + ( + relocate, + BTreeSet::from([a, b]), + "a move stamps both parents", + ), + ( + Command::Delete { node: x }, + BTreeSet::from([a]), + "a delete into the bin stamps the parent", + ), + ( + restore, + BTreeSet::from([a]), + "a restore stamps the folder it restores into", + ), + ] { + assert_eq!( + stamped_by_command(&world, &mut engine, &mut tasks, command), + expected, + "{why}" + ); + } + + write_file( + &mut engine, + WriteTarget::NewFile { + parent: b, + name: "f.bin".into(), + }, + b"first", + ) + .expect("the write commits"); + tick(&world, &engine, &mut tasks); + let file = child_id(&engine, b, "f.bin"); + assert_eq!( + stamped_by(&world, &mut engine, &mut tasks, |engine| { + write_file(engine, version(file), b"second").expect("the edit commits"); + }), + BTreeSet::from([file]), + "a content op stamps the node alone" + ); +} + +/// A dual-linked node's rename or move republishes only the winning parent, +/// which the drain resolves, so the overlay stamps that parent and not the +/// folder of the other link. +#[test] +fn a_dual_linked_node_stamps_only_the_parent_its_publish_writes() { + let world = FakeWorld::new(); + let blocks = Blocks::default(); + seed_account(&world, &blocks); + let alice = world.device(b"alice"); + let (mut engine, _events, mut tasks) = boot(&world, &blocks, &alice, 42); + let (photos, deep) = seed_dual_linked_file(&world, &blocks, &mut engine, &mut tasks); + block_on(engine.command(Command::Create { + parent: ROOT, + name: "albums".into(), + kind: NodeKind::Folder, + })) + .unwrap(); + tick(&world, &engine, &mut tasks); + let albums = child_id(&engine, ROOT, "albums"); + // The winning link: the higher counter, then the lower parent id. + let counter_under = |parent: NodeId| { + published_children(&world.record_store, &blocks, parent) + .iter() + .find(|child| child.id == deep.0) + .map(|child| child.link_counter) + .expect("both folders name the node") + }; + let winner = [ROOT, photos] + .into_iter() + .min_by_key(|parent| (core::cmp::Reverse(counter_under(*parent)), *parent)) + .expect("two links"); + + let renamed = stamped_by_command( + &world, + &mut engine, + &mut tasks, + Command::Rename { + node: deep, + new_name: "renamed.bin".into(), + }, + ); + assert_eq!( + renamed, + BTreeSet::from([winner]), + "a rename stamps the winning parent alone" + ); + + let moved = stamped_by_command( + &world, + &mut engine, + &mut tasks, + Command::Move { + node: deep, + new_parent: albums, + new_name: "moved.bin".into(), + replacing: None, + }, + ); + assert_eq!( + moved, + BTreeSet::from([winner, albums]), + "a move stamps the winning parent and the destination" + ); +} diff --git a/crates/wasm/src/seams_bridge.rs b/crates/wasm/src/seams_bridge.rs index 29d2c0a6c0..f134f562b1 100644 --- a/crates/wasm/src/seams_bridge.rs +++ b/crates/wasm/src/seams_bridge.rs @@ -236,6 +236,8 @@ extern "C" { #[wasm_bindgen(method, catch, js_name = enqueueOp)] async fn enqueue_op(this: &JsStagingStoreSeam, op: &[u8]) -> Result; + #[wasm_bindgen(method, catch, js_name = enqueueOps)] + async fn enqueue_ops(this: &JsStagingStoreSeam, ops: Array) -> Result; #[wasm_bindgen(method, catch, js_name = queuedOps)] async fn queued_ops(this: &JsStagingStoreSeam) -> Result; #[wasm_bindgen(method, catch, js_name = removeOp)] @@ -274,6 +276,21 @@ impl StagingStore for StagingStoreAdapter { required_u64(self.js.enqueue_op(op).await.map_err(seam_error)?).map(OpId) } + async fn enqueue_ops(&self, ops: &[Vec]) -> SeamResult> { + let entries: Array = ops + .iter() + .map(|op| Uint8Array::from(op.as_slice())) + .collect(); + let value = self.js.enqueue_ops(entries).await.map_err(seam_error)?; + let ids: Array = value + .dyn_into() + .map_err(|_| SeamError::new("enqueueOps must return an array"))?; + if ids.length() as usize != ops.len() { + return Err(SeamError::new("enqueueOps must return one id per op")); + } + ids.iter().map(|id| required_u64(id).map(OpId)).collect() + } + async fn queued_ops(&self) -> SeamResult)>> { let value = self.js.queued_ops().await.map_err(seam_error)?; let array: Array = value diff --git a/decisions/0045-a-queued-op-journals-its-crossing-as-a-plan-and-the-drain-decides.md b/decisions/0045-a-queued-op-journals-its-crossing-as-a-plan-and-the-drain-decides.md index 81a6d489c2..faf52a3353 100644 --- a/decisions/0045-a-queued-op-journals-its-crossing-as-a-plan-and-the-drain-decides.md +++ b/decisions/0045-a-queued-op-journals-its-crossing-as-a-plan-and-the-drain-decides.md @@ -160,7 +160,7 @@ record, so the overlay stamps `mtime = authored_at`" (`blueprint/engine.md` "Syn state-law bullet). That target-only form first appeared in the body of FSM1/cipher-box#878, which an agent wrote, and no owner decision stands behind it. It contradicts the resolution of FSM1/cipher-box#830. The owner accepted the FSM1/cipher-box#830 set as the rule on 2026-09-26. -The code does not meet it yet (E3). +`Op::authored_nodes` (`crates/engine/src/sync/op.rs`) is the one function. **D6 — A `move` is a relink and a rename in one entry, and one kernel rename is one command.** The intent op list of `#33` D5 gains `move`. A `move` carries the relink, the rename @@ -338,9 +338,9 @@ Its gain is that the subtree publishes once, not twice. law" gains the citation (ADR 0045) on the rename sentence. `CONTEXT.md` gains the citation (ADR 0045) on the "Cross-scope move" and "Retained record" entries. -8. **The code must follow the reworded blueprint.** FSM1/cipher-box#2013 tracks moving the - overlay and the drain onto the one authored-node set of D5 (E3). FSM1/cipher-box#2014 tracks - journaling the two legs of D6 in one atomic write (E2). +8. **The code follows the reworded blueprint.** FSM1/cipher-box#2103 moved the overlay and the + drain onto the one authored-node set of D5 (E3), and journals the two legs of D6 in one + atomic write (E2). 9. **No wire format, no KDF edge and no op record format changes.** @@ -353,39 +353,7 @@ no member and no attacker can observe the field. While the walk is still dark, t the failed unseal to an uncharged halt, and the op waits at the queue head until the walk proves the boundary. That is an availability cost, not a trust cost. -**E2 — The two legs of D6 are not journaled atomically.** `stage_legs_and_notify` -(`crates/engine/src/facade.rs`) journals the parking leg and the arriving leg in two -`enqueue_op` calls. If the arriving leg fails to journal, the command removes the parking leg -on a best-effort basis (`let _ = self.dequeue_op(op_id)`). That removal cannot undo a parking -leg that a tick already drained, and the code comment allows that case. The source is then cut, -the subtree sits in the vault-root scope, and the caller hears that the command failed. That is -a false failure that performed a scope exit, which D6 forbids. The fix is one atomic journal -write for both legs, which needs a multi-entry enqueue on the `StagingStore` seam, and a test -that fails the arriving leg's journal. FSM1/cipher-box#2014 tracks it. - -**E3 — The overlay and the drain stamp different node sets.** D5 requires one authored-node -set that both call. In the code only the overlay calls `Op::stamp_authored` -(`crates/engine/src/sync/op.rs`), and it stamps the op's target alone. The drain reads -`applied.op.authored_at` at each authoring site in `crates/engine/src/sync/drain.rs` and in -`crates/engine/src/net/author.rs`. The two sets differ: - -- For `rename`, `relink` and `move`, the overlay's `relocate` (`crates/engine/src/sync/overlay.rs`) - stamps the renamed or relinked node. The drain's `publish_ref_move` never republishes that - node's record; it publishes only the parent folders, with `modified_at = authored_at`. So the - overlay stamps a node that the drain does not author, and it does not stamp the parents that - the drain does author. -- For `create`, the drain's `publish_create` stamps the child and the parent, and the overlay - stamps the child only. For `delete`, the drain republishes each parent at `authored_at`, and - the overlay stamps nothing. - -A folder's mtime therefore jumps at publish, which is what section 3 of the resolution of -FSM1/cipher-box#830 exists to prevent. FSM1/cipher-box#2013 tracks the fix. Also, the overlay does not stamp a `prune`, a -`restoreVersion` or a `deleteVersion`, although each authors its target's next record. The -drain keeps the existing `modified_at` for those three kinds (`publish_prune`, -`publish_restore_version` and `publish_delete_version` in `crates/engine/src/sync/drain.rs`), so -the overlay and the drain agree there. Those kinds are outside the six-op list of the -blueprint, and only its sentence "every op but a delete authors its target's next record" -reaches them. +**E2 and E3** were resolved by FSM1/cipher-box#2103 on 2026-09-30. **E4 — `authored_at` is a client-authored time.** A skewed or lying client publishes a skewed `modified_at`, and no gate check catches it. The only reader is the owner's own view, so this @@ -441,10 +409,11 @@ until the arriving leg publishes. FSM1/cipher-box#1765 names this cost on purpos - **D5:** `a_content_op_stamps_its_authored_time_and_plaintext_size` and `a_metadata_op_stamps_time_over_a_projection_and_leaves_size_alone` (`crates/engine/src/sync/op.rs`), and `a_new_child_stamps_the_journaled_time_not_a_clock` - (`crates/engine/src/net/author.rs`) cover the stamp on one node. **Finding:** no test proves - that the overlay and the drain's publish plan stamp the same set of nodes, and no test - enumerates the op kinds, which the resolution of FSM1/cipher-box#830 asked for. The code does - not meet the D5 set today (E3, FSM1/cipher-box#2013). + (`crates/engine/src/net/author.rs`) cover the stamp on one node. + `each_op_kind_authors_its_own_set_of_nodes` (`crates/engine/src/sync/op.rs`) enumerates the + op kinds, and `an_op_renders_the_times_its_publish_writes` + (`crates/engine/tests/write_plane.rs`) proves that the overlay and the drain's publish plan + stamp the same set of nodes. - **D6:** `a_durable_queue_outage_never_destroys_the_destination_a_rename_did_not_replace` and `replacing_a_junk_holding_folder_keeps_the_destination_entry_when_the_queue_fails` (`crates/fuse/tests/fuse_op_core.rs`), `overlay_move_relinks_renames_and_replaces_in_one_step` @@ -457,8 +426,10 @@ until the arriving leg publishes. FSM1/cipher-box#1765 names this cost on purpos `a_move_between_two_granted_folders_re_seals_into_the_destination_scope`, `a_restart_between_the_legs_of_a_staged_move_cuts_the_source_once` and `the_passes_that_cannot_author_a_staged_move_do_not_spend_it` - (`crates/engine/tests/owner_actions.rs`) cover the two legs. **Finding:** no test covers a - two-leg relocation whose arriving leg fails to journal (E2, FSM1/cipher-box#2014). + (`crates/engine/tests/owner_actions.rs`) cover the two legs. + `a_staged_move_whose_arriving_leg_will_not_journal_journals_no_leg` + (`crates/engine/tests/owner_actions.rs`) covers a two-leg relocation whose arriving leg fails + to journal. - **D7:** `conditional_edit_dead_letters_when_another_writer_took_the_head` and `a_second_queued_edit_dead_letters_behind_a_superseded_first` (`crates/engine/src/sync/rebase.rs`), and diff --git a/packages/client/src/seams/stagingStore.ts b/packages/client/src/seams/stagingStore.ts index 2a3839c8f7..e8a777455c 100644 --- a/packages/client/src/seams/stagingStore.ts +++ b/packages/client/src/seams/stagingStore.ts @@ -153,6 +153,28 @@ export class OpfsStagingStore implements StagingStoreSeam { return Number(key); } + async enqueueOps(ops: Uint8Array[]): Promise { + // Copy before the first await, as `enqueueOp` does. + const staged = ops.map((op) => op.slice()); + const db = await this.open(); + // One transaction is the atomic write: a refused add aborts every add. + const tx = db.transaction(STAGING_OPS_STORE, 'readwrite'); + const done = transactionDone(tx); + const store = tx.objectStore(STAGING_OPS_STORE); + const requests: Array> = []; + try { + for (const op of staged) requests.push(requestResult(store.add(op))); + } catch (error) { + // An add that throws, rather than failing its request, would otherwise + // let the transaction commit the adds issued before it. + tx.abort(); + await Promise.allSettled([...requests, done]); + throw error; + } + const [keys] = await Promise.all([Promise.all(requests), done]); + return keys.map(Number); + } + async queuedOps(): Promise> { const db = await this.open(); const tx = db.transaction(STAGING_OPS_STORE, 'readonly'); diff --git a/packages/client/src/seams/types.ts b/packages/client/src/seams/types.ts index 507ea10318..6d5745680c 100644 --- a/packages/client/src/seams/types.ts +++ b/packages/client/src/seams/types.ts @@ -45,6 +45,8 @@ export interface SnapshotCacheSeam { /** Durable op queue (FIFO, strictly increasing never-reused ids) plus staged bytes. */ export interface StagingStoreSeam { enqueueOp(op: Uint8Array): Promise; + /** Appends every op in one atomic write, in order: all or none (an error after commit can leave all). */ + enqueueOps(ops: Uint8Array[]): Promise; queuedOps(): Promise>; removeOp(opId: number): Promise; putStagedBytes(stagingKey: Uint8Array, bytes: Uint8Array): Promise; diff --git a/packages/client/test/browser/conformance.spec.ts b/packages/client/test/browser/conformance.spec.ts index ab2343ed5e..3c248e72aa 100644 --- a/packages/client/test/browser/conformance.spec.ts +++ b/packages/client/test/browser/conformance.spec.ts @@ -56,6 +56,12 @@ test.describe('browser seam conformance', () => { expect(outcome.ok).toBe(true); }); + test('staging store queues a multi-entry enqueue whole or not at all', async ({ page }) => { + const outcome = await runSeam(page, 'stagingStoreBatch'); + expect(outcome.error ?? '', 'staging batch behavioral failure').toBe(''); + expect(outcome.ok).toBe(true); + }); + test('account switching reclaims snapshots and preserves owner-local bookkeeping', async ({ page, }) => { diff --git a/packages/client/test/browser/conformance.worker.ts b/packages/client/test/browser/conformance.worker.ts index 3c8c417931..aec3696d92 100644 --- a/packages/client/test/browser/conformance.worker.ts +++ b/packages/client/test/browser/conformance.worker.ts @@ -130,6 +130,41 @@ async function assertStagedEntryCount(dirName: string, expected: number): Promis } } +/** + * A multi-entry enqueue is one IndexedDB transaction: an add that throws part + * way through leaves none of the set queued, now or after a reopen. + */ +async function runStagingBatchBehavioral(): Promise { + const name = 'conf-staging-batch'; + await deleteDatabase(name); + const store = new OpfsStagingStore(name); + const single = await store.enqueueOp(new Uint8Array([1])); + + const add = IDBObjectStore.prototype.add; + let calls = 0; + IDBObjectStore.prototype.add = function (this: IDBObjectStore, ...args) { + calls += 1; + if (calls === 2) throw new DOMException('refused', 'DataError'); + return add.apply(this, args); + }; + let refused = false; + try { + await store.enqueueOps([new Uint8Array([2]), new Uint8Array([3])]); + } catch { + refused = true; + } finally { + IDBObjectStore.prototype.add = add; + } + if (!refused) throw new Error('enqueueOps: the refused add did not fail the set'); + + for (const reader of [store, new OpfsStagingStore(name)]) { + const ids = (await reader.queuedOps()).map(([id]) => id); + if (ids.length !== 1 || ids[0] !== single) { + throw new Error(`enqueueOps: a failed set left entries queued: [${ids.join(', ')}]`); + } + } +} + /** * In-flight write debris — the temp a killed put leaves behind — is not a * staged record: it stays out of the orphan-GC enumeration and the staging @@ -174,6 +209,7 @@ const STAGED_AFTER_KIT: Record = { 'failed-replacement': 1, 'failed-first-put': 0, cleared: 0, + batched: 0, }; /** Stores are per-fault: the kit asserts on leftover staged counts, so another @@ -510,6 +546,10 @@ async function run(seam: string): Promise { await runStagingDebrisBehavioral(); return; } + case 'stagingStoreBatch': { + await runStagingBatchBehavioral(); + return; + } case 'storeReclaim': { await runStoreReclaimBehavioral(); return;