Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions crates/desktop-seams/src/fs_util.rs
Original file line number Diff line number Diff line change
Expand Up @@ -167,12 +167,12 @@ pub(crate) fn keep_first<E>(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)
}

Expand Down
133 changes: 108 additions & 25 deletions crates/desktop-seams/src/staging_store.rs
Original file line number Diff line number Diff line change
@@ -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/<id>.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/<first-id>-<last-id>.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/<hexkey>.bin`).
const SIDECAR_SUFFIX: &str = ".bin";
/// Filename of the durable monotonic op-id counter.
Expand All @@ -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/<first>-<last>.batch` — the marker of a multi-entry enqueue that has
/// not committed ([`BATCH_SUFFIX`]).
/// - `staged/<hexkey>.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
Expand Down Expand Up @@ -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<u64>) -> 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<u64>,
written: impl Iterator<Item = u64>,
) -> 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<u64> {
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<u8>]) -> 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 {
Expand All @@ -103,27 +183,26 @@ impl FileStagingStore {

impl StagingStore for FileStagingStore {
async fn enqueue_op(&self, op: &[u8]) -> SeamResult<OpId> {
// 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<u8>]) -> SeamResult<Vec<OpId>> {
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<Vec<(OpId, Vec<u8>)>> {
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
Expand All @@ -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))
Expand Down Expand Up @@ -221,10 +300,14 @@ fn parse_op_id(name: &str) -> Option<u64> {
name.strip_suffix(OP_SUFFIX)?.parse::<u64>().ok()
}

/// The largest op id currently on disk in `ops_dir`, if any.
fn highest_op_id(ops_dir: &Path) -> std::io::Result<Option<u64>> {
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<RangeInclusive<u64>> {
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()
}
76 changes: 75 additions & 1 deletion crates/desktop-seams/tests/conformance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down Expand Up @@ -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.
Expand Down
57 changes: 31 additions & 26 deletions crates/engine/src/facade.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down Expand Up @@ -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<CommandOutcome, EngineError> {
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<Vec<u8>, 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.
Expand Down
Loading
Loading