diff --git a/CHANGELOG.md b/CHANGELOG.md index 7a71cb8c3..8e76a122d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +- **R-15 Gate 4D:** stage a captured inventory one page at a time (stage one of the streaming capture). + `StreamedCapture` takes the entries page by page, assigns the page index and generation, checks that they + are valid and strictly ascending across pages, seals each page once and stages its envelope in a private + `inventory/pending-*` directory (never authority, never read by a reader), keeping one page and one + reference per page in memory; the total is announced up front because the inventory digest commits to it + first. `finish` builds the successor with the same function as `capture_inventory`, and + `load_staged_page` confirms a staged page against its reference, a missing one being a broken attempt and + not a recovery state of the journal. A capture is bound to its journal directory and to the key that opened + its manifest, and a failure the caller can handle removes the pages staged so far. Promotion under the root + lock and the composed commit are the next stage. PR #1010. - **R-15 Gate 4D:** the journal owner can renew its own lease (maintainer decision C, the wiring half). `commit_lease_renewal` publishes a manifest that moves the lease expiry strictly forward and changes nothing else (`publish_renewal_fenced`, a publish of its own so the ordinary checkpoint keeps refusing lease changes) diff --git a/crates/worldscript-secure-storage/src/journal/capture.rs b/crates/worldscript-secure-storage/src/journal/capture.rs index 176e71277..1879db609 100644 --- a/crates/worldscript-secure-storage/src/journal/capture.rs +++ b/crates/worldscript-secure-storage/src/journal/capture.rs @@ -48,23 +48,52 @@ pub fn capture_inventory( ) -> Result { assert_fence(manifest, fence)?; assert_inventory_open(manifest)?; - let revision = manifest - .journal_revision - .checked_add(1) - .ok_or(JournalError::InvalidCounter)?; + let revision = next_revision(manifest)?; let ordered = ordered_pages(pages)?; let refs = page_refs(&ordered, revision)?; + let entry_count = entry_total(&ordered)?; + let inventory_digest = inventory_digest_of(manifest.inventory_version, &ordered, entry_count)?; + let parts = CapturedInventory { + refs: &refs, + entry_count, + inventory_digest, + }; + capture_successor(manifest, &parts) +} + +/// What a capture of the inventory consists of, however its pages were produced: one reference per +/// page, the entry total and the inventory digest over the entries. +pub(super) struct CapturedInventory<'a> { + pub refs: &'a [JournalPageRef], + pub entry_count: u32, + pub inventory_digest: [u8; 32], +} + +/// The revision a capture of `manifest` is published at: the next one, and the generation of every +/// page the capture writes. +pub(super) fn next_revision(manifest: &JournalManifest) -> Result { + manifest + .journal_revision + .checked_add(1) + .ok_or_else(|| JournalError::InvalidCounter.into()) +} + +/// The successor of `manifest` that carries `parts` as its inventory, checked as a valid successor. +/// Both the capture from pages in memory and the streamed capture end here. +pub(super) fn capture_successor( + manifest: &JournalManifest, + parts: &CapturedInventory<'_>, +) -> Result { let mut next = manifest.clone(); // Only the capture that runs behind `ADMIT`'s barrier is the final one (§10.3). next.final_inventory_captured = manifest.phase == phase_code::ADMIT; - next.journal_revision = revision; - next.page_count = u32::try_from(ordered.len()).map_err(|_| JournalError::TooManyEntries)?; - next.entry_count = entry_total(&ordered)?; - next.inventory_digest = - inventory_digest_of(manifest.inventory_version, &ordered, next.entry_count)?; - next.journal_page_set_digest = journal_page_set_digest(&refs)?; + next.journal_revision = next_revision(manifest)?; + next.page_count = u32::try_from(parts.refs.len()).map_err(|_| JournalError::TooManyEntries)?; + next.entry_count = parts.entry_count; + next.inventory_digest = parts.inventory_digest; + next.journal_page_set_digest = journal_page_set_digest(parts.refs)?; next.encode()?; - next.verify_page_set(&refs)?; + next.verify_page_set(parts.refs)?; assert_manifest_successor(manifest, &next)?; Ok(next) } @@ -73,7 +102,9 @@ pub fn capture_inventory( /// starts, while nothing has been converted and, because the final inventory is immutable once /// captured, before the final capture. `PREPARE` is excluded: ordinary writes are admitted /// there until `ADMIT`'s barrier, so only the `ADMIT` snapshot can be the commit inventory. -fn assert_inventory_open(manifest: &JournalManifest) -> Result<(), MigrationExecutionError> { +pub(super) fn assert_inventory_open( + manifest: &JournalManifest, +) -> Result<(), MigrationExecutionError> { let phase = manifest_phase(manifest); if is_terminal_phase(phase) { return Err(MigrationExecutionError::TerminalPhase); diff --git a/crates/worldscript-secure-storage/src/journal/durable.rs b/crates/worldscript-secure-storage/src/journal/durable.rs index 7bbca9481..a897fc24f 100644 --- a/crates/worldscript-secure-storage/src/journal/durable.rs +++ b/crates/worldscript-secure-storage/src/journal/durable.rs @@ -546,18 +546,63 @@ pub fn load_authoritative_manifest( live: &LiveMigration, ) -> Result { let identity = migration_identity(&live.operation_id)?; - let (bytes, digest) = read_root_named_envelope(ctx.fs, ctx.dir, live)?; + let (bytes, _) = read_root_named_envelope(ctx.fs, ctx.dir, live)?; + open_root_named(ctx.key, &identity, live, &bytes) +} + +/// Opens the root-named envelope `bytes` under `key` and requires the binding to vouch for it: the +/// part of [`load_authoritative_manifest`] that needs no read. The epoch of the header is compared +/// before the key is used. +fn open_root_named( + key: &Key, + identity: &RecordIdentity, + live: &LiveMigration, + bytes: &[u8], +) -> Result { let read = ManifestRead { - record: &identity, + record: identity, journal_revision: live.journal_revision, - key_epoch: header_key_epoch(&bytes)?, - envelope: &bytes, + key_epoch: header_key_epoch(bytes)?, + envelope: bytes, }; - let manifest = JournalManifest::open(ctx.key, &read)?; + let manifest = JournalManifest::open(key, &read)?; + let digest = ManifestEnvelopeDigest::from_bytes(content_digest(bytes)); assert_live_binding(&manifest, live, digest).map_err(JournalDurableError::Authority)?; Ok(manifest) } +/// The exact bytes of the generation the root binding names, after proving that they open under the +/// context's key to exactly `manifest`. A caller that works for a long time keeps the bytes and +/// repeats the key check with [`assert_anchor`], which needs no read. +pub(super) fn root_named_anchor( + ctx: &mut JournalDurableContext<'_, F>, + manifest: &JournalManifest, + live: &LiveMigration, +) -> Result, JournalDurableError> { + let (bytes, _) = read_root_named_envelope(ctx.fs, ctx.dir, live)?; + assert_anchor(ctx.key, manifest, live, &bytes)?; + Ok(bytes) +} + +/// Requires `bytes`, the root-named generation of `live`, to open under `key` to exactly `manifest`. +/// Pure computation over a few kilobytes: no read, so it can be repeated for every page of a +/// capture. +pub(super) fn assert_anchor( + key: &Key, + manifest: &JournalManifest, + live: &LiveMigration, + bytes: &[u8], +) -> Result<(), JournalDurableError> { + let identity = migration_identity(&live.operation_id)?; + if open_root_named(key, &identity, live, bytes)? == *manifest { + Ok(()) + } else { + Err(JournalDurableError::Authority( + MigrationExecutionError::LiveBindingMismatch, + )) + } +} + /// Reads the generation the committed root names, bounded, and proves the bytes are what the root /// committed to: authority first (§6), before any key is used. The binding authenticates these exact /// bytes, header included, so the epoch in the header is vouched for by the root, never by the diff --git a/crates/worldscript-secure-storage/src/journal/inventory_read.rs b/crates/worldscript-secure-storage/src/journal/inventory_read.rs index fccdf2c88..38442dc13 100644 --- a/crates/worldscript-secure-storage/src/journal/inventory_read.rs +++ b/crates/worldscript-secure-storage/src/journal/inventory_read.rs @@ -38,7 +38,7 @@ use super::state::MigrationExecutionError; use super::{JournalError, MAX_JOURNAL_PAGE_BYTES}; /// The largest valid sealed page: the page encoding bound plus the envelope header and tag. -const MAX_PAGE_ENVELOPE_BYTES: usize = MAX_JOURNAL_PAGE_BYTES + HEADER_LEN + TAG_LEN; +pub(super) const MAX_PAGE_ENVELOPE_BYTES: usize = MAX_JOURNAL_PAGE_BYTES + HEADER_LEN + TAG_LEN; /// A page directory holds one generation and, at most, a few staging leftovers; a directory with /// more entries than this is not read. const MAX_PAGE_DIRECTORY_ENTRIES: usize = 64; @@ -142,7 +142,7 @@ fn hinted_generation( } /// The envelope at `path`, never larger than any valid sealed page. -fn read_envelope( +pub(super) fn read_envelope( ctx: &mut JournalDurableContext<'_, F>, path: &Path, ) -> Result, JournalDurableError> { @@ -157,7 +157,7 @@ fn read_envelope( /// Opens `bytes` as this operation's page `index` at `generation` under the journal key. Read routing /// is authority-first (§6): the header's epoch is compared with the operation's journal envelope /// epoch before the key is used, so a misrouted page is never decrypted. -fn open_stored_page( +pub(super) fn open_stored_page( ctx: &JournalDurableContext<'_, F>, manifest: &JournalManifest, index: u32, diff --git a/crates/worldscript-secure-storage/src/journal/inventory_store.rs b/crates/worldscript-secure-storage/src/journal/inventory_store.rs index 2a718f2ee..f8f2d22c5 100644 --- a/crates/worldscript-secure-storage/src/journal/inventory_store.rs +++ b/crates/worldscript-secure-storage/src/journal/inventory_store.rs @@ -238,26 +238,30 @@ fn store_page( &set.successor.journal_page_set_digest, sealed.page.page_index(), ); - if holds_envelope(ctx, &dir, sealed)? { + store_page_at(ctx, &dir, set.committed_manifest, sealed) +} + +/// Stores one page of `manifest`'s operation into the page directory `dir`, which sits under the +/// journal directory: the same adoption, staging, promotion and directory syncs wherever the page set +/// is being assembled. +pub(super) fn store_page_at( + ctx: &mut JournalDurableContext<'_, F>, + dir: &Path, + manifest: &JournalManifest, + sealed: &SealedPage<'_>, +) -> Result { + if holds_envelope(ctx, dir, sealed)? { // The earlier attempt may have stopped before its directories were durable, and, under the // same operation id, may have left its own staging link behind. - let staging = staging_residue(ctx, &dir, sealed.page.page_generation()); - return sync_chain(ctx, &dir, staging); + let staging = staging_residue(ctx, dir, sealed.page.page_generation()); + return sync_chain(ctx, dir, staging); } ctx.fs - .create_dir_all(&dir) + .create_dir_all(dir) .map_err(|error| JournalDurableError::Stage(stage_io(error)))?; - let identity = migration_page_identity( - &set.committed_manifest.operation_id, - sealed.page.page_index(), - )?; - let epoch = journal_envelope_epoch(set.committed_manifest)?; - let request = stage_request( - &dir, - &identity, - page_meta(sealed.page, epoch), - ctx.operation, - ); + let identity = migration_page_identity(&manifest.operation_id, sealed.page.page_index())?; + let epoch = journal_envelope_epoch(manifest)?; + let request = stage_request(dir, &identity, page_meta(sealed.page, epoch), ctx.operation); let promoted = stage_and_promote_envelope(ctx.fs, ctx.key, &request, sealed.envelope.to_vec())?; let above = dir.parent().unwrap_or(ctx.dir); let synced = sync_chain(ctx, above, promoted.staging)?; diff --git a/crates/worldscript-secure-storage/src/journal/mod.rs b/crates/worldscript-secure-storage/src/journal/mod.rs index 49fef19d0..73b528d6d 100644 --- a/crates/worldscript-secure-storage/src/journal/mod.rs +++ b/crates/worldscript-secure-storage/src/journal/mod.rs @@ -16,6 +16,7 @@ mod manifest_verify; mod page; mod renewal; mod state; +mod stream_capture; mod succession; mod takeover; mod wire; @@ -146,6 +147,9 @@ pub use state::{ JournalInventoryExtent, JournalRevision, ManifestEnvelopeDigest, MigrationExecutionError, MigrationFence, MigrationPhase, RecoveryReasonCode, }; +pub use stream_capture::{ + load_staged_page, CaptureStart, StagedCapture, StagedPage, StreamedCapture, +}; pub use succession::{assert_manifest_successor, assert_progress_successor}; pub use takeover::{ assert_binding_takeover, assert_takeover_promote_authority, assert_takeover_successor, diff --git a/crates/worldscript-secure-storage/src/journal/stream_capture.rs b/crates/worldscript-secure-storage/src/journal/stream_capture.rs new file mode 100644 index 000000000..5f7d90052 --- /dev/null +++ b/crates/worldscript-secure-storage/src/journal/stream_capture.rs @@ -0,0 +1,383 @@ +//! Gate 4D streaming capture, stage one: building and staging a captured inventory one page at a time +//! (§10.1.1, §10.3). +//! +//! [`capture_inventory`](super::capture::capture_inventory) and the inventory store hold every sealed +//! page in memory, which the format does not bound: it allows a million entries of up to a kilobyte in +//! pages of up to 4096 entries. A page cannot be written to its final place as it is sealed, because +//! that place is keyed by the page-set digest, which needs the envelope digest of every page, and a +//! second sealing pass would use fresh nonces and give different bytes. So each envelope is kept on +//! disk, in a private pending directory, between "sealed and hashed" and "digest known": +//! +//! ```text +//! /inventory/pending--/page-/generation-.wsr1 +//! ``` +//! +//! The pending directory has the layout of a digest directory, so the final promotion (stage two) is +//! the same store of the same bytes. Nothing in it is authority: the manifest never names it, no +//! reader looks for it, and a staged page is a candidate whose absence is a broken attempt, never a +//! recovery state of the journal. A failure that the caller can still handle removes the staged page +//! files on the way out (best effort), and [`StagedCapture::discard`] does the same for a finished +//! capture that is not going to be promoted. The empty directories stay, because the file system +//! abstraction has no directory removal, and so do the files of an attempt that was killed: both are +//! inert and are reclaimed by a separate cleanup (an acceptance criterion on #359). +//! +//! The inventory digest commits to the total entry count before any entry (§5.4), so the total is +//! given up front and checked at the end. Memory is one page and one reference per page; the builder +//! holds no page and no envelope after a push returns. No lock is taken while pages are staged; the +//! authority checks made by [`StreamedCapture::begin`] only refuse early, and stage two repeats them +//! under the root lock and the journal mutex. + +use std::io::ErrorKind; +use std::path::{Path, PathBuf}; + +use crate::durable::{generation_path, staging_path, DurableFs, WriteOperationId}; +use crate::marker::content_digest; +use crate::root::LiveMigration; + +use super::capture::{ + assert_inventory_open, assert_page_not_empty, capture_successor, next_revision, + CapturedInventory, SealedPage, +}; +use super::digest::{page_ref_for, InventoryDigestVerifier}; +use super::durable::{ + assert_anchor, root_named_anchor, stage_io, with_fence, JournalDurableContext, + JournalDurableError, +}; +use super::inventory::JournalInventoryEntry; +use super::inventory_read::{open_stored_page, MAX_PAGE_ENVELOPE_BYTES}; +use super::inventory_store::{seal_inventory_pages, store_page_at}; +use super::manifest::{JournalManifest, JournalPageRef}; +use super::page::JournalPage; +use super::state::{assert_page_promote_authority, MigrationExecutionError, MigrationFence}; +use super::JournalError; + +/// A capture whose pages are being staged. Every push consumes the builder and returns it, so a +/// builder that failed partway cannot be pushed to again. It belongs to one journal directory and to +/// the key that authenticated the manifest it started from; a push under another context is refused. +pub struct StreamedCapture { + committed: JournalManifest, + live: LiveMigration, + /// The exact bytes of the root-named generation, which the capture's key must keep opening to + /// `committed`: a few kilobytes, kept so that no push reads the journal for it. + anchor: Vec, + journal_dir: PathBuf, + revision: u64, + entry_count: u32, + pending: PathBuf, + /// The attempt's own operation identity: it names the pending directory and every staging link, so + /// the names a failed promotion can leave behind are known to a later cleanup. + tag: WriteOperationId, + refs: Vec, + verifier: InventoryDigestVerifier, +} + +// The verifier is deliberately left out: it is a running hash, not state worth printing. +impl std::fmt::Debug for StreamedCapture { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("StreamedCapture") + .field("revision", &self.revision) + .field("entry_count", &self.entry_count) + .field("pages_staged", &self.refs.len()) + .field("pending", &self.pending) + .finish_non_exhaustive() + } +} + +/// A finished capture: the successor manifest that names the page set, and where the staged pages +/// are. Nothing is published yet. It is a handle to files on disk, so it is deliberately not +/// `Clone`: whoever holds it owns the decision to promote or to [`discard`](Self::discard) them. +#[derive(Debug, PartialEq, Eq)] +pub struct StagedCapture { + successor: JournalManifest, + pending: PathBuf, + tag: WriteOperationId, + refs: Vec, +} + +/// A staged page read back and confirmed against its reference: the decoded page and the exact +/// envelope bytes the page set binds. +#[derive(Clone, PartialEq, Eq)] +pub struct StagedPage { + pub page: JournalPage, + pub envelope: Vec, +} + +/// What a streamed capture starts from: the manifest the root binding names, the owner's fence, that +/// binding, and the announced total of entries. +#[derive(Clone, Copy)] +pub struct CaptureStart<'a> { + pub committed_manifest: &'a JournalManifest, + pub fence: &'a MigrationFence, + pub live: &'a LiveMigration, + pub entry_count: u32, +} + +impl StagedCapture { + /// The successor manifest that captures the staged inventory. + pub fn successor(&self) -> &JournalManifest { + &self.successor + } + + /// The reference of every staged page, in page-index order. + pub fn page_refs(&self) -> &[JournalPageRef] { + &self.refs + } + + /// The private pending directory the pages are staged in. + pub fn pending_dir(&self) -> &Path { + &self.pending + } + + /// Abandons the capture: the staged page files are removed, best effort. The empty directories + /// stay, because the file system abstraction cannot remove a directory. + pub fn discard(self, ctx: &mut JournalDurableContext<'_, F>) { + discard_staged(ctx, &self.pending, &self.refs, &self.tag); + } +} + +impl StreamedCapture { + /// Starts a capture of `start.entry_count` entries over the manifest the root binding names. + /// + /// Refused before anything is created: a stale fence, a caller that is not the committed owner + /// writing under the committed manifest, a manifest that is not the exact root-named generation, + /// an inventory that is not open (`DISCOVER` or `ADMIT`, not yet captured as final, nothing + /// converted) and a total above the format's bound. These are early refusals only: whoever + /// promotes the staged pages repeats them under the locks. + pub fn begin( + ctx: &mut JournalDurableContext<'_, F>, + start: &CaptureStart<'_>, + ) -> Result { + let committed = start.committed_manifest; + let (verifier, anchor) = with_fence(committed, start.fence, || { + assert_page_promote_authority(committed, Some(start.live)) + .map_err(JournalDurableError::Authority)?; + let anchor = root_named_anchor(ctx, committed, start.live)?; + assert_inventory_open(committed).map_err(JournalDurableError::Authority)?; + let verifier = + InventoryDigestVerifier::new(committed.inventory_version, start.entry_count)?; + Ok((verifier, anchor)) + })?; + let revision = next_revision(committed).map_err(JournalDurableError::Authority)?; + // A fresh random tag, never the caller's operation id: two attempts must never share a + // pending directory, because each seals its pages with its own nonces. + let tag = WriteOperationId::generate() + .map_err(|error| JournalDurableError::Journal(JournalError::Seal(error)))?; + let pending = ctx + .dir + .join("inventory") + .join(format!("pending-{revision}-{}", tag.as_str())); + Ok(Self { + committed: committed.clone(), + live: start.live.clone(), + anchor, + journal_dir: ctx.dir.to_path_buf(), + revision, + entry_count: start.entry_count, + pending, + tag, + refs: Vec::new(), + verifier, + }) + } + + /// The private pending directory the pages are being staged in. + pub fn pending_dir(&self) -> &Path { + &self.pending + } + + /// Stages the next page. The page takes the next index and the capture's generation; its + /// entries must be valid and strictly ascending after every earlier entry, and the running total + /// must stay within the count given at the start. The sealed envelope is staged, synced, and + /// dropped. The context must be the one the capture started under: the same journal directory, + /// and a key that authenticates the manifest the capture started from. + /// + /// A failure removes the staged pages with the builder, best effort, because nobody else can. + pub fn push_page( + mut self, + ctx: &mut JournalDurableContext<'_, F>, + entries: Vec, + ) -> Result { + match self.stage_next(ctx, entries) { + Ok(()) => Ok(self), + Err(error) => { + discard_staged(ctx, &self.pending, &self.refs, &self.tag); + Err(error) + } + } + } + + fn stage_next( + &mut self, + ctx: &mut JournalDurableContext<'_, F>, + entries: Vec, + ) -> Result<(), JournalDurableError> { + self.assert_context(ctx)?; + let index = u32::try_from(self.refs.len()).map_err(|_| JournalError::InvalidPageIndex)?; + let page = JournalPage::new(index, self.revision, entries)?; + assert_page_not_empty(&page)?; + self.verifier.absorb_page(&page)?; + let envelope = + seal_inventory_pages(ctx.key, &self.committed, std::slice::from_ref(&page))?.remove(0); + let reference = page_ref_for(&page, &envelope)?; + let dir = self.pending.join(format!("page-{index}")); + let sealed = SealedPage { + page: &page, + envelope: &envelope, + }; + // Staged under the attempt's own operation identity, not the caller's, so that the name of + // a staging link that a promotion leaves behind is known to the cleanup. + let mut staging = JournalDurableContext::new(&mut *ctx.fs, ctx.key, ctx.dir, &self.tag); + if let Err(error) = store_page_at(&mut staging, &dir, &self.committed, &sealed) { + // The page in flight may be on disk already, or the failure may be that a different file + // sits in its slot: only a file that holds exactly the bytes staged here is ours to remove. + let digest = &reference.page_content_digest; + remove_if_staged(ctx, &generation_path(&dir, self.revision), digest); + // The attempt's own staging link, which a failed promotion reports and leaves behind. + remove_if_staged(ctx, &staging_path(&dir, self.revision, &self.tag), digest); + return Err(error); + } + self.refs.push(reference); + Ok(()) + } + + /// The capture belongs to one journal directory and one key. A context for another journal, or + /// whose key does not open the manifest the capture started from, is refused before anything is + /// sealed or written, so a page can never be staged that the journal key cannot open. + fn assert_context( + &self, + ctx: &JournalDurableContext<'_, F>, + ) -> Result<(), JournalDurableError> { + if ctx.dir != self.journal_dir.as_path() { + return Err(JournalDurableError::Authority( + MigrationExecutionError::LiveBindingMismatch, + )); + } + assert_anchor(ctx.key, &self.committed, &self.live, &self.anchor) + } + + /// Ends the capture: exactly the announced number of entries must have been staged. Returns the + /// successor manifest, built and checked exactly as + /// [`capture_inventory`](super::capture::capture_inventory) builds it from the same pages. A + /// capture that does not end here is abandoned: its staged pages are removed, best effort. + pub fn finish( + self, + ctx: &mut JournalDurableContext<'_, F>, + ) -> Result { + let Self { + committed, + entry_count, + pending, + tag, + refs, + verifier, + .. + } = self; + let successor = verifier + .finish_digest() + .map_err(JournalDurableError::from) + .and_then(|inventory_digest| { + let parts = CapturedInventory { + refs: &refs, + entry_count, + inventory_digest, + }; + capture_successor(&committed, &parts).map_err(JournalDurableError::Authority) + }); + match successor { + Ok(successor) => Ok(StagedCapture { + successor, + pending, + tag, + refs, + }), + Err(error) => { + discard_staged(ctx, &pending, &refs, &tag); + Err(error) + } + } + } +} + +/// Removes the staged page files that `refs` describe, and the staging link a promotion may have left +/// next to each (named by the attempt's `tag`, whether or not the promotion reported success). +/// Best effort: a file that cannot be removed stays as inert residue, and the error that led here is +/// the one worth reporting, so removal failures are not. A file is removed only if it still holds +/// exactly the bytes this capture staged (see [`remove_if_staged`]). +fn discard_staged( + ctx: &mut JournalDurableContext<'_, F>, + pending: &Path, + refs: &[JournalPageRef], + tag: &WriteOperationId, +) { + for (index, reference) in refs.iter().enumerate() { + let dir = pending.join(format!("page-{index}")); + let digest = &reference.page_content_digest; + let generation = reference.page_generation; + remove_if_staged(ctx, &generation_path(&dir, generation), digest); + remove_if_staged(ctx, &staging_path(&dir, generation, tag), digest); + } +} + +/// Removes the file at `path` only if its bytes hash to `digest`, the digest of what this capture +/// staged there, so a file that was replaced, or that a failure found already sitting in the slot, is +/// never deleted. The read is bounded by the largest valid page envelope. +fn remove_if_staged( + ctx: &mut JournalDurableContext<'_, F>, + path: &Path, + digest: &[u8; 32], +) { + let found = ctx.fs.read_at_most(path, MAX_PAGE_ENVELOPE_BYTES); + if matches!(found, Ok(Some(bytes)) if content_digest(&bytes) == *digest) { + let _ = ctx.fs.remove_file(path); + } +} + +/// The staged envelope at `path`, never larger than any valid sealed page. A staged file is a +/// candidate, not authority: a missing one is a broken attempt to abandon, never a recovery state of +/// the journal, so it is not reported as `RecoveryRequired` the way a missing authoritative page is. +fn read_staged_envelope( + ctx: &mut JournalDurableContext<'_, F>, + path: &Path, +) -> Result, JournalDurableError> { + match ctx.fs.read_at_most(path, MAX_PAGE_ENVELOPE_BYTES) { + Ok(Some(bytes)) => Ok(bytes), + Ok(None) => Err(JournalError::Corrupt("staged page exceeds the envelope bound").into()), + Err(error) if error.kind() == ErrorKind::NotFound => { + Err(JournalError::Corrupt("staged page is missing").into()) + } + Err(error) => Err(JournalDurableError::Stage(stage_io(error))), + } +} + +/// Reads staged page `page_index` back by the exact path of its reference, bounded. +/// +/// The envelope must hash to the reference before it is opened, and it opens only as this +/// operation's page at this index and generation, under the journal key and at the operation's +/// journal envelope epoch (compared before the key is used), so a staged file that changed after it +/// was written is refused rather than trusted. +pub fn load_staged_page( + ctx: &mut JournalDurableContext<'_, F>, + staged: &StagedCapture, + page_index: u32, +) -> Result { + let reference = usize::try_from(page_index) + .ok() + .and_then(|index| staged.refs.get(index)) + .ok_or(JournalError::InvalidPageIndex)?; + let dir = staged.pending.join(format!("page-{page_index}")); + let envelope = read_staged_envelope(ctx, &generation_path(&dir, reference.page_generation))?; + if content_digest(&envelope) != reference.page_content_digest { + return Err(JournalError::PageSetMismatch.into()); + } + let page = open_stored_page( + ctx, + &staged.successor, + page_index, + reference.page_generation, + &envelope, + )?; + if page.entries().len() as u64 != u64::from(reference.page_entry_count) { + return Err(JournalError::EntryCountMismatch.into()); + } + Ok(StagedPage { page, envelope }) +} diff --git a/crates/worldscript-secure-storage/src/lib.rs b/crates/worldscript-secure-storage/src/lib.rs index b3362bd95..b60e5308a 100644 --- a/crates/worldscript-secure-storage/src/lib.rs +++ b/crates/worldscript-secure-storage/src/lib.rs @@ -96,18 +96,19 @@ pub use journal::{ assert_takeover_successor, authoritative_manifest_revision, capture_inventory, checkpoint_progress, empty_inventory_digest, empty_journal_page_set_digest, inventory_digest, inventory_page_dir, is_terminal_phase, journal_envelope_epoch, journal_page_set_digest, - load_authoritative_manifest, load_inventory_page, load_manifest_generation, mark_done, - mark_recovery, operation_type, ordinary_mutating_writes_admitted, page_ref_for, phase_code, - promote_inventory_set_fenced, promote_manifest_fenced, promote_page_fenced, + load_authoritative_manifest, load_inventory_page, load_manifest_generation, load_staged_page, + mark_done, mark_recovery, operation_type, ordinary_mutating_writes_admitted, page_ref_for, + phase_code, promote_inventory_set_fenced, promote_manifest_fenced, promote_page_fenced, publish_manifest_fenced, publish_renewal_fenced, publish_takeover_fenced, root_named_journal_epoch, seal_inventory_pages, source_authority_kind, source_physical_authority_kind, source_scheme_id, transition_phase, verify_stored_inventory, - with_fence, CandidateConflict, ForeignInventoryExtension, InventoryDigestVerifier, - InventorySetWrite, JournalCheckpointCursor, JournalDurableContext, JournalDurableError, - JournalDurableGuard, JournalError, JournalInventoryEntry, JournalInventoryExtent, - JournalInventorySource, JournalManifest, JournalPage, JournalPageRef, JournalRevision, - JournalTakeover, ManifestEnvelopeDigest, ManifestRead, MigrationExecutionError, MigrationFence, - MigrationPhase, PublishedManifest, RecoveryReasonCode, SealedPage, VerifiedInventory, + with_fence, CandidateConflict, CaptureStart, ForeignInventoryExtension, + InventoryDigestVerifier, InventorySetWrite, JournalCheckpointCursor, JournalDurableContext, + JournalDurableError, JournalDurableGuard, JournalError, JournalInventoryEntry, + JournalInventoryExtent, JournalInventorySource, JournalManifest, JournalPage, JournalPageRef, + JournalRevision, JournalTakeover, ManifestEnvelopeDigest, ManifestRead, + MigrationExecutionError, MigrationFence, MigrationPhase, PublishedManifest, RecoveryReasonCode, + SealedPage, StagedCapture, StagedPage, StreamedCapture, VerifiedInventory, JOURNAL_MANIFEST_FORMAT_VERSION, JOURNAL_MANIFEST_RECORD_SCHEMA, JOURNAL_PAGE_FORMAT_VERSION, JOURNAL_PAGE_RECORD_SCHEMA, MAX_JOURNAL_ENTRY_BYTES, MAX_JOURNAL_INVENTORY_ENTRIES, MAX_JOURNAL_MANIFEST_ENVELOPE_BYTES, MAX_JOURNAL_PAGE_BYTES, MAX_JOURNAL_PAGE_DESCRIPTORS, diff --git a/crates/worldscript-secure-storage/tests/gate4d_inventory_read_test.rs b/crates/worldscript-secure-storage/tests/gate4d_inventory_read_test.rs index 9b4ab52a6..9c5d23a5f 100644 --- a/crates/worldscript-secure-storage/tests/gate4d_inventory_read_test.rs +++ b/crates/worldscript-secure-storage/tests/gate4d_inventory_read_test.rs @@ -1,10 +1,13 @@ //! Gate 4D Slice C1b-2: reading back a stored inventory page set (§10.1.1). +#[path = "support/journal_fixture.rs"] +mod journal_fixture; #[path = "support/inventory.rs"] mod support; use std::path::PathBuf; +use journal_fixture::*; use support::*; use worldscript_secure_storage::*; diff --git a/crates/worldscript-secure-storage/tests/gate4d_inventory_store_test.rs b/crates/worldscript-secure-storage/tests/gate4d_inventory_store_test.rs index 8c4981fbb..3ff01ddfc 100644 --- a/crates/worldscript-secure-storage/tests/gate4d_inventory_store_test.rs +++ b/crates/worldscript-secure-storage/tests/gate4d_inventory_store_test.rs @@ -1,10 +1,13 @@ //! Gate 4D Slice C1b-1: where the pages of a captured inventory live, and writing them (§10.1.1). +#[path = "support/journal_fixture.rs"] +mod journal_fixture; #[path = "support/inventory.rs"] mod support; use std::path::Path; +use journal_fixture::*; use support::*; use worldscript_secure_storage::*; diff --git a/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs b/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs new file mode 100644 index 000000000..3312f92ee --- /dev/null +++ b/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs @@ -0,0 +1,662 @@ +//! Gate 4D streaming capture, stage one: an inventory built and staged one page at a time (§10.1.1, +//! §10.3). The streamed successor must be exactly what `capture_inventory` builds from the pages that +//! were staged, the staged files are inert, and a staged page is trusted only after it is confirmed +//! against its reference. Every test is a table of named cases with one assertion over all of them. + +#[path = "support/journal_fixture.rs"] +mod journal_fixture; + +use std::collections::BTreeSet; +use std::path::Path; + +use journal_fixture::*; +use worldscript_secure_storage::*; + +/// `count` ascending entries (six digits, so the order is the numeric one for any count). +fn entries(count: u32) -> Vec { + (0..count).map(entry_n).collect() +} + +fn entry_n(n: u32) -> JournalInventoryEntry { + let record = RecordIdentity::new(RecordClass::Codex, &[&format!("p{n:06}")]).unwrap(); + let source = JournalInventorySource { + authority_kind: source_authority_kind::LEGACY_PLAINTEXT, + physical_authority_kind: source_physical_authority_kind::TAURI_FILESYSTEM, + generation: None, + evidence_digest: Some([n as u8; 32]), + foreign: None, + }; + JournalInventoryEntry::new(record, source).unwrap() +} + +/// A change to the committed manifest a scenario starts from. +type Change = fn(&mut JournalManifest); + +/// A real journal whose root-named generation is `committed`. +struct Scenario { + journal: Journal, + committed: JournalManifest, +} + +impl Scenario { + fn new() -> Self { + Self::with(|_| {}) + } + + /// The committed manifest after `change`, committed as the root-named generation. + fn with(change: Change) -> Self { + let mut committed = manifest_at(COMMITTED_REVISION); + change(&mut committed); + let dir = TempDir::new(); + let live = commit_into(&dir.0, &committed); + Scenario { + journal: Journal { dir, live }, + committed, + } + } + + fn fence(&self) -> MigrationFence { + MigrationFence::from_manifest(&self.committed) + } + + /// Starts a capture of `count` entries over `fs`. + fn begin( + &self, + fs: &mut F, + count: u32, + ) -> Result { + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut ctx = JournalDurableContext::new(fs, &key, self.journal.path(), &op); + let start = CaptureStart { + committed_manifest: &self.committed, + fence: &self.fence(), + live: &self.journal.live, + entry_count: count, + }; + StreamedCapture::begin(&mut ctx, &start) + } + + /// Stages `entries` cut into pages of `per_page`. + fn stream( + &self, + fs: &mut F, + entries: Vec, + per_page: usize, + ) -> Result { + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut capture = self.begin(fs, entries.len() as u32)?; + let mut ctx = JournalDurableContext::new(fs, &key, self.journal.path(), &op); + for chunk in entries.chunks(per_page) { + capture = capture.push_page(&mut ctx, chunk.to_vec())?; + } + capture.finish(&mut ctx) + } + + /// Abandons a finished capture. + fn discard(&self, staged: StagedCapture) { + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = ObservedFs::new(); + let mut ctx = JournalDurableContext::new(&mut fs, &key, self.journal.path(), &op); + staged.discard(&mut ctx); + } + + /// Staged page `index`, read back under `key`. + fn read_under( + &self, + key: &Key, + staged: &StagedCapture, + index: u32, + ) -> Result { + let op = WriteOperationId::generate().unwrap(); + let mut fs = StdFs; + let mut ctx = JournalDurableContext::new(&mut fs, key, self.journal.path(), &op); + load_staged_page(&mut ctx, staged, index) + } + + fn read(&self, staged: &StagedCapture, index: u32) -> Result { + self.read_under(&key(), staged, index) + } + + /// Every staged page, read back and confirmed. + fn staged_pages(&self, staged: &StagedCapture) -> Vec { + (0..staged.page_refs().len() as u32) + .map(|index| self.read(staged, index).unwrap()) + .collect() + } + + /// What `capture_inventory` builds from the pages the stream staged. + fn replayed(&self, staged: &StagedCapture) -> JournalManifest { + let pages = self.staged_pages(staged); + let sealed: Vec<_> = pages + .iter() + .map(|p| SealedPage { + page: &p.page, + envelope: &p.envelope, + }) + .collect(); + capture_inventory(&self.committed, &self.fence(), &sealed).unwrap() + } +} + +/// Every file under `root`, as a path relative to it. +fn files_under(root: &Path) -> BTreeSet { + let mut found = BTreeSet::new(); + let mut stack = vec![root.to_path_buf()]; + while let Some(dir) = stack.pop() { + for item in std::fs::read_dir(&dir).unwrap() { + let path = item.unwrap().path(); + if path.is_dir() { + stack.push(path); + } else { + let relative = path.strip_prefix(root).unwrap(); + found.insert(relative.to_string_lossy().replace('\\', "/")); + } + } + } + found +} + +fn admit(manifest: &mut JournalManifest) { + manifest.phase = phase_code::ADMIT; +} + +fn already_final(manifest: &mut JournalManifest) { + manifest.phase = phase_code::ADMIT; + manifest.final_inventory_captured = true; +} + +fn converting(manifest: &mut JournalManifest) { + manifest.phase = phase_code::CONVERT; + manifest.final_inventory_captured = true; +} + +fn preparing(manifest: &mut JournalManifest) { + manifest.phase = phase_code::PREPARE; +} + +fn finished(manifest: &mut JournalManifest) { + manifest.phase = phase_code::DONE; + manifest.final_inventory_captured = true; +} + +#[test] +fn the_streamed_successor_is_what_capture_inventory_builds_from_the_staged_pages() { + let cases: [(&str, u32, usize, Change); 6] = [ + ("one entry", 1, 1, |_| {}), + ("an uneven split", 7, 3, |_| {}), + ("an even split", 8, 4, |_| {}), + ("one page", 5, 5, |_| {}), + ("the final capture in ADMIT", 5, 2, admit), + ("a full page and one more", 4097, 4096, |_| {}), + ]; + let wrong: Vec<_> = cases + .into_iter() + .filter_map(|(name, count, per_page, change)| { + let scenario = Scenario::with(change); + let staged = scenario + .stream(&mut StdFs, entries(count), per_page) + .unwrap(); + let same = *staged.successor() == scenario.replayed(&staged); + let counted = staged.successor().entry_count == count; + (!(same && counted)).then_some(name) + }) + .collect(); + assert_eq!(wrong, Vec::<&str>::new()); +} + +#[test] +fn an_empty_inventory_stages_nothing_and_matches_the_empty_capture() { + let scenario = Scenario::new(); + let before = files_under(scenario.journal.path()); + let staged = scenario.stream(&mut StdFs, Vec::new(), 1).unwrap(); + let empty = capture_inventory(&scenario.committed, &scenario.fence(), &[]).unwrap(); + assert_eq!( + ( + *staged.successor() == empty, + files_under(scenario.journal.path()) + ), + (true, before) + ); +} + +#[test] +fn staged_pages_are_inert_files_in_a_private_pending_directory() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let pending = staged + .pending_dir() + .strip_prefix(scenario.journal.path()) + .unwrap(); + let pending = pending.to_string_lossy().replace('\\', "/"); + let mut expected = BTreeSet::from(["generation-3.wsr1".to_owned()]); + expected.extend((0..3).map(|i| format!("{pending}/page-{i}/generation-4.wsr1"))); + // The directory is named for what it is, never like a digest directory a reader would resolve. + assert_eq!( + ( + pending.starts_with("inventory/pending-4-"), + files_under(scenario.journal.path()) + ), + (true, expected) + ); +} + +#[test] +fn two_attempts_under_one_operation_id_never_share_a_pending_directory() { + let scenario = Scenario::new(); + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = StdFs; + let mut ctx = JournalDurableContext::new(&mut fs, &key, scenario.journal.path(), &op); + let start = CaptureStart { + committed_manifest: &scenario.committed, + fence: &scenario.fence(), + live: &scenario.journal.live, + entry_count: 0, + }; + let mut attempt = || { + let capture = StreamedCapture::begin(&mut ctx, &start).unwrap(); + capture.finish(&mut ctx).unwrap() + }; + let (first, second) = (attempt(), attempt()); + assert_ne!(first.pending_dir(), second.pending_dir()); +} + +#[test] +fn the_staged_pages_are_accepted_by_the_existing_store() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let pages = scenario.staged_pages(&staged); + let sealed: Vec<_> = pages + .iter() + .map(|p| SealedPage { + page: &p.page, + envelope: &p.envelope, + }) + .collect(); + let set = InventorySetWrite { + committed_manifest: &scenario.committed, + fence: &scenario.fence(), + committed: Some(&scenario.journal.live), + successor: staged.successor(), + pages: &sealed, + }; + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = StdFs; + let mut ctx = JournalDurableContext::new(&mut fs, &key, scenario.journal.path(), &op); + assert!(promote_inventory_set_fenced(&mut ctx, &set).is_ok()); +} + +#[test] +fn a_capture_that_cannot_start_creates_nothing() { + let refused = JournalDurableError::Authority; + let cases: [(&str, Change, u32, JournalDurableError); 5] = [ + ( + "a capture already final", + already_final, + 1, + refused(MigrationExecutionError::FrozenFieldChanged), + ), + ( + "conversion already started", + converting, + 1, + refused(MigrationExecutionError::FrozenFieldChanged), + ), + ( + "a phase without a snapshot", + preparing, + 1, + refused(MigrationExecutionError::InvalidPhaseTransition), + ), + ( + "a finished journal", + finished, + 1, + refused(MigrationExecutionError::TerminalPhase), + ), + ( + "a total above the bound", + |_| {}, + MAX_JOURNAL_INVENTORY_ENTRIES + 1, + JournalDurableError::Journal(JournalError::TooManyEntries), + ), + ]; + let wrong: Vec<_> = cases + .into_iter() + .filter_map(|(name, change, count, error)| { + let scenario = Scenario::with(change); + let mut fs = ObservedFs::new(); + let refusal = scenario.begin(&mut fs, count).err(); + (!(refusal == Some(error) && fs.created_nothing())).then_some(name) + }) + .collect(); + assert_eq!(wrong, Vec::<&str>::new()); +} + +/// The refusal of a `begin` over `committed` against the binding `live`, which must create nothing. +fn begin_refusal( + scenario: &Scenario, + committed: &JournalManifest, + live: &LiveMigration, +) -> (Option, bool) { + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = ObservedFs::new(); + let mut ctx = JournalDurableContext::new(&mut fs, &key, scenario.journal.path(), &op); + let start = CaptureStart { + committed_manifest: committed, + fence: &MigrationFence::from_manifest(committed), + live, + entry_count: 1, + }; + let refusal = StreamedCapture::begin(&mut ctx, &start).err(); + (refusal, fs.created_nothing()) +} + +#[test] +fn a_caller_that_is_not_the_committed_owner_of_the_root_named_manifest_is_refused() { + let scenario = Scenario::new(); + let mut moved = scenario.journal.live.clone(); + moved.fencing_generation += 1; + let mut other = scenario.committed.clone(); + other.has_lease_owner = true; + other.lease_owner_id = Some("someone".into()); + other.lease_expires_unix_ms = Some(1); + let authority = JournalDurableError::Authority; + assert_eq!( + [ + begin_refusal(&scenario, &scenario.committed, &moved), + begin_refusal(&scenario, &other, &scenario.journal.live), + ], + [ + ( + Some(authority(MigrationExecutionError::StaleMigrationOwner)), + true + ), + ( + Some(authority(MigrationExecutionError::LiveBindingMismatch)), + true + ), + ] + ); +} + +#[test] +fn entries_arrive_in_order_and_the_announced_total_is_enforced_at_both_ends() { + let scenario = Scenario::new(); + let too_many = JournalDurableError::Journal(JournalError::TooManyEntries); + let short = JournalDurableError::Journal(JournalError::EntryCountMismatch); + let unordered = JournalDurableError::Journal(JournalError::NotStrictlyAscending); + let empty_page = JournalDurableError::Journal(JournalError::InvalidDescriptorCount); + let run = |count: u32, pages: Vec>| -> Option { + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = StdFs; + let mut capture = scenario.begin(&mut fs, count).unwrap(); + let mut ctx = JournalDurableContext::new(&mut fs, &key, scenario.journal.path(), &op); + for page in pages { + capture = match capture.push_page(&mut ctx, page) { + Ok(next) => next, + Err(error) => return Some(error), + }; + } + capture.finish(&mut ctx).err() + }; + let all = entries(6); + let outcomes = [ + run(6, vec![all[3..6].to_vec(), all[0..3].to_vec()]), + run(4, vec![all[0..5].to_vec()]), + run(4, vec![all[0..3].to_vec()]), + run(1, vec![Vec::new()]), + run(6, vec![all[0..3].to_vec(), all[2..6].to_vec()]), + ]; + assert_eq!( + outcomes, + [ + Some(unordered.clone()), + Some(too_many), + Some(short), + Some(empty_page), + Some(unordered) + ] + ); +} + +#[test] +fn a_failure_partway_removes_the_pages_it_staged_and_a_new_attempt_is_independent() { + let scenario = Scenario::new(); + let committed_only = files_under(scenario.journal.path()); + let mut failing = ObservedFs::new(); + failing.fail_create_at = Some(3); + let failed = scenario.stream(&mut failing, entries(7), 3); + let after_failure = files_under(scenario.journal.path()); + let second = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + assert_eq!( + (failed.is_err(), after_failure, second.page_refs().len()), + (true, committed_only, 3) + ); +} + +/// The pages a handled failure of any kind takes with it: the capture that ends short of its total +/// and the finished capture that is discarded, each followed by the files left on disk. +#[test] +fn an_abandoned_capture_removes_its_staged_pages() { + let scenario = Scenario::new(); + let committed_only = files_under(scenario.journal.path()); + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = StdFs; + let capture = scenario.begin(&mut fs, 4).unwrap(); + let mut ctx = JournalDurableContext::new(&mut fs, &key, scenario.journal.path(), &op); + let capture = capture.push_page(&mut ctx, entries(3)).unwrap(); + let staged_then_short = ( + files_under(scenario.journal.path()).len(), + capture.finish(&mut ctx).err(), + ); + let short = JournalDurableError::Journal(JournalError::EntryCountMismatch); + let after_short = files_under(scenario.journal.path()); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let staged_files = files_under(scenario.journal.path()).len(); + scenario.discard(staged); + assert_eq!( + ( + staged_then_short, + after_short, + staged_files, + files_under(scenario.journal.path()) + ), + ((2, Some(short)), committed_only.clone(), 4, committed_only) + ); +} + +#[test] +fn a_cleanup_that_cannot_remove_a_file_leaves_inert_residue_and_does_not_fail() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = ObservedFs::new(); + fs.fail_remove = true; + let mut ctx = JournalDurableContext::new(&mut fs, &key, scenario.journal.path(), &op); + staged.discard(&mut ctx); + assert_eq!(files_under(scenario.journal.path()).len(), 4); +} + +/// A push under `key` in `dir`, over a capture that began in `scenario`'s journal. +fn push_under(scenario: &Scenario, key: &Key, dir: &Path) -> Option { + let op = WriteOperationId::generate().unwrap(); + let mut fs = StdFs; + let capture = scenario.begin(&mut fs, 3).unwrap(); + let mut ctx = JournalDurableContext::new(&mut fs, key, dir, &op); + capture.push_page(&mut ctx, entries(3)).err() +} + +#[test] +fn a_push_under_another_key_or_another_journal_is_refused_before_anything_is_staged() { + let scenario = Scenario::new(); + // A directory that holds a byte-identical copy of the root-named generation: the manifest check + // alone would accept it, so only the capture's own binding to its journal refuses it. + let copy = TempDir::new(); + let generation = generation_path(scenario.journal.path(), COMMITTED_REVISION); + std::fs::copy(&generation, generation_path(©.0, COMMITTED_REVISION)).unwrap(); + let committed_only = files_under(scenario.journal.path()); + let outcomes = [ + push_under(&scenario, &other_key(), scenario.journal.path()), + push_under(&scenario, &key(), ©.0), + ]; + assert_eq!( + ( + outcomes, + files_under(scenario.journal.path()), + files_under(©.0) + ), + ( + [ + Some(JournalDurableError::Journal(JournalError::Open( + OpenError::Tampered + ))), + Some(JournalDurableError::Authority( + MigrationExecutionError::LiveBindingMismatch + )), + ], + committed_only.clone(), + committed_only + ) + ); +} + +#[test] +fn a_failure_after_a_page_was_promoted_removes_that_page_too() { + let scenario = Scenario::new(); + let committed_only = files_under(scenario.journal.path()); + let mut failing = ObservedFs::new(); + failing.fail_sync_matching = Some("page-1"); + let failed = scenario.stream(&mut failing, entries(7), 3); + assert_eq!( + ( + matches!(failed, Err(JournalDurableError::Stage(_))), + files_under(scenario.journal.path()) + ), + (true, committed_only) + ); +} + +#[test] +fn pushing_a_page_does_not_read_the_root_named_manifest_again() { + let scenario = Scenario::new(); + let mut fs = ObservedFs::new(); + scenario.stream(&mut fs, entries(12), 3).unwrap(); + let generation = generation_path(scenario.journal.path(), COMMITTED_REVISION); + // Four pages were pushed; the key check of each is made over the bytes `begin` kept. + assert_eq!( + fs.reads.iter().filter(|path| **path == generation).count(), + 1 + ); +} + +#[test] +fn a_failure_never_removes_a_file_it_did_not_write() { + let scenario = Scenario::new(); + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = StdFs; + let capture = scenario.begin(&mut fs, 6).unwrap(); + let mut ctx = JournalDurableContext::new(&mut fs, &key, scenario.journal.path(), &op); + let all = entries(6); + let capture = capture.push_page(&mut ctx, all[0..3].to_vec()).unwrap(); + let pending = capture.pending_dir().to_path_buf(); + // A different file already sits in the slot of the next page: the push fails because of it. + let slot = generation_path(&pending.join("page-1"), 4); + std::fs::create_dir_all(slot.parent().unwrap()).unwrap(); + std::fs::write(&slot, b"bytes of someone else").unwrap(); + let failed = capture.push_page(&mut ctx, all[3..6].to_vec()).err(); + let first = generation_path(&pending.join("page-0"), 4); + // Nothing of the failed attempt is left but the file it did not write: no page, and not the + // staging link a failed promotion reports either. + assert_eq!( + ( + matches!(failed, Some(JournalDurableError::Stage(_))), + std::fs::read(&slot).unwrap(), + first.exists(), + files_under(&pending), + ), + ( + true, + b"bytes of someone else".to_vec(), + false, + BTreeSet::from(["page-1/generation-4.wsr1".to_owned()]) + ) + ); +} + +#[test] +fn a_cleanup_never_removes_a_staged_page_that_was_replaced() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let replaced = generation_path(&staged.pending_dir().join("page-1"), 4); + std::fs::write(&replaced, b"a file that is not the staged page").unwrap(); + scenario.discard(staged); + let committed = "generation-3.wsr1".to_owned(); + let remaining: Vec<_> = files_under(scenario.journal.path()).into_iter().collect(); + assert_eq!( + ( + remaining.len(), + remaining.contains(&committed), + std::fs::read(&replaced).unwrap() + ), + (2, true, b"a file that is not the staged page".to_vec()) + ); +} + +#[test] +fn a_staging_link_left_after_a_successful_promotion_is_found_by_the_cleanup() { + let scenario = Scenario::new(); + let committed_only = files_under(scenario.journal.path()); + // Promotion succeeds, but the staging link cannot be unlinked: the page is staged, the link stays. + let mut failing = ObservedFs::new(); + failing.fail_remove = true; + let staged = scenario.stream(&mut failing, entries(7), 3).unwrap(); + let with_links = files_under(scenario.journal.path()).len(); + scenario.discard(staged); + assert_eq!( + (with_links, files_under(scenario.journal.path())), + (1 + 3 * 2, committed_only) + ); +} + +#[test] +fn a_missing_staged_page_is_a_broken_attempt_not_a_recovery_state_of_the_journal() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let file = generation_path(&staged.pending_dir().join("page-1"), 4); + std::fs::remove_file(file).unwrap(); + assert_eq!( + scenario.read(&staged, 1).err(), + Some(JournalDurableError::Journal(JournalError::Corrupt( + "staged page is missing" + ))) + ); +} + +#[test] +fn a_staged_page_is_trusted_only_after_it_is_confirmed_against_its_reference() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let file = |index: u32| generation_path(&staged.pending_dir().join(format!("page-{index}")), 4); + let mut flipped = std::fs::read(file(0)).unwrap(); + flipped[20] ^= 1; + std::fs::write(file(0), flipped).unwrap(); + std::fs::write(file(1), std::fs::read(file(2)).unwrap()).unwrap(); + let outcomes = [ + scenario.read(&staged, 0).err(), + scenario.read(&staged, 1).err(), + scenario.read(&staged, 7).err(), + scenario.read_under(&other_key(), &staged, 2).err(), + ]; + let mismatch = JournalDurableError::Journal(JournalError::PageSetMismatch); + assert_eq!( + outcomes, + [ + Some(mismatch.clone()), + Some(mismatch), + Some(JournalDurableError::Journal(JournalError::InvalidPageIndex)), + Some(JournalDurableError::Journal(JournalError::Open( + OpenError::Tampered + ))), + ] + ); +} diff --git a/crates/worldscript-secure-storage/tests/support/inventory.rs b/crates/worldscript-secure-storage/tests/support/inventory.rs index c07fab4ba..eea7217bd 100644 --- a/crates/worldscript-secure-storage/tests/support/inventory.rs +++ b/crates/worldscript-secure-storage/tests/support/inventory.rs @@ -1,164 +1,18 @@ -//! Shared Gate 4D inventory fixtures: a real committed journal, a captured and sealed inventory, a -//! file system that observes what a store creates, and the fenced store drivers. +//! Shared Gate 4D inventory fixtures: a captured and sealed inventory over a committed journal (see +//! `journal_fixture`) and the fenced store drivers. -use std::ffi::OsString; -use std::fs::File; -use std::io; use std::path::{Path, PathBuf}; -use std::sync::atomic::{AtomicU32, Ordering}; +use crate::journal_fixture::*; use worldscript_secure_storage::{ - capture_inventory, content_digest, empty_inventory_digest, empty_journal_page_set_digest, - generation_path, inventory_page_dir, journal_page_set_digest, operation_type, page_ref_for, - parse_envelope, phase_code, promote_inventory_set_fenced, promote_manifest_fenced, - seal_inventory_pages, source_authority_kind, source_physical_authority_kind, - DirectoryDurability, DurableFs, InventorySetWrite, JournalDurableContext, JournalDurableError, - JournalInventoryEntry, JournalInventorySource, JournalManifest, JournalPage, Key, - LiveMigration, MigrationFence, RecordClass, RecordIdentity, RecordMeta, SealedPage, StdFs, - WriteOperationId, + capture_inventory, inventory_page_dir, journal_page_set_digest, page_ref_for, parse_envelope, + promote_inventory_set_fenced, seal_inventory_pages, source_authority_kind, + source_physical_authority_kind, DirectoryDurability, DurableFs, InventorySetWrite, + JournalDurableContext, JournalDurableError, JournalInventoryEntry, JournalInventorySource, + JournalManifest, JournalPage, Key, LiveMigration, MigrationFence, RecordClass, RecordIdentity, + RecordMeta, SealedPage, WriteOperationId, }; -pub const OPERATION: &str = "store-op"; -pub const COMMITTED_REVISION: u64 = 3; - -pub struct TempDir(pub PathBuf); - -impl TempDir { - pub fn new() -> Self { - static NEXT: AtomicU32 = AtomicU32::new(0); - let path = std::env::temp_dir().join(format!( - "wss-gate4d-store-{}-{}", - std::process::id(), - NEXT.fetch_add(1, Ordering::Relaxed) - )); - std::fs::create_dir_all(&path).unwrap(); - TempDir(path) - } -} - -impl Drop for TempDir { - fn drop(&mut self) { - let _ = std::fs::remove_dir_all(&self.0); - } -} - -/// A real file system that counts what a refused store must never do, logs directory syncs and can -/// fail the sync of one directory or the n-th file creation. -pub struct ObservedFs { - pub inner: StdFs, - pub creates: u32, - pub dirs_created: u32, - pub synced: Vec, - pub fail_sync_of: Option, - pub fail_create_at: Option, - pub fail_remove: bool, -} - -impl ObservedFs { - pub fn new() -> Self { - Self { - inner: StdFs, - creates: 0, - dirs_created: 0, - synced: Vec::new(), - fail_sync_of: None, - fail_create_at: None, - fail_remove: false, - } - } - - pub fn created_nothing(&self) -> bool { - (self.creates, self.dirs_created) == (0, 0) - } -} - -impl DurableFs for ObservedFs { - type File = File; - - fn create_new(&mut self, path: &Path) -> io::Result { - self.creates += 1; - if self.fail_create_at == Some(self.creates) { - return Err(io::Error::other("injected file creation failure")); - } - self.inner.create_new(path) - } - - fn sync_file(&mut self, file: &mut File) -> io::Result<()> { - self.inner.sync_file(file) - } - - fn read(&mut self, path: &Path) -> io::Result> { - self.inner.read(path) - } - - fn link_no_replace(&mut self, from: &Path, to: &Path) -> io::Result<()> { - self.inner.link_no_replace(from, to) - } - - fn remove_file(&mut self, path: &Path) -> io::Result<()> { - if self.fail_remove { - return Err(io::Error::other("injected file removal failure")); - } - self.inner.remove_file(path) - } - - fn sync_dir(&mut self, dir: &Path) -> io::Result { - self.synced.push(dir.to_path_buf()); - if self.fail_sync_of.as_deref() == Some(dir) { - return Err(io::Error::other("injected directory sync failure")); - } - self.inner.sync_dir(dir) - } - - fn list_dir(&mut self, dir: &Path) -> io::Result> { - self.inner.list_dir(dir) - } - - fn rename_replace(&mut self, from: &Path, to: &Path) -> io::Result<()> { - self.inner.rename_replace(from, to) - } - - fn create_dir_all(&mut self, dir: &Path) -> io::Result<()> { - self.dirs_created += 1; - self.inner.create_dir_all(dir) - } -} - -pub fn key() -> worldscript_secure_storage::Key { - worldscript_secure_storage::Key::from_bytes(&mut [9u8; 32]) -} - -/// A key that is not the journal key: what another epoch's pages are sealed under. -pub fn other_key() -> worldscript_secure_storage::Key { - worldscript_secure_storage::Key::from_bytes(&mut [10u8; 32]) -} - -pub fn manifest_at(revision: u64) -> JournalManifest { - JournalManifest { - operation_id: OPERATION.into(), - journal_revision: revision, - operation_type: operation_type::ROTATE, - phase: phase_code::DISCOVER, - source_epoch: 1, - target_epoch: 2, - has_target_root_key_ref: true, - target_root_key_ref_digest: Some([0x42; 32]), - fencing_generation: 7, - inventory_version: 1, - inventory_digest: empty_inventory_digest(1), - page_count: 0, - entry_count: 0, - journal_page_set_digest: empty_journal_page_set_digest(), - final_inventory_captured: false, - cursor_page_index: 0, - cursor_entry_index: 0, - has_lease_owner: false, - lease_owner_id: None, - lease_expires_unix_ms: None, - recovery_reason_code: 0, - } -} - pub fn entry(n: u32) -> JournalInventoryEntry { let record = RecordIdentity::new(RecordClass::Codex, &[&format!("p{n:03}")]).unwrap(); JournalInventoryEntry::new( @@ -246,18 +100,6 @@ pub fn seal_all<'a>(pages: &'a [JournalPage], envelopes: &'a [Vec]) -> Vec &Path { - &self.dir.0 - } -} - /// Commits `captured`'s committed manifest to a fresh directory and returns the binding naming it. pub fn journal_of(captured: &Captured) -> Journal { let dir = TempDir::new(); @@ -265,28 +107,6 @@ pub fn journal_of(captured: &Captured) -> Journal { Journal { dir, live } } -/// Commits `manifest` as the root-named generation of `dir`; returns the binding that names it. -pub fn commit_into(dir: &Path, manifest: &JournalManifest) -> LiveMigration { - let previous = LiveMigration { - operation_id: manifest.operation_id.clone(), - fencing_generation: manifest.fencing_generation, - journal_revision: manifest.journal_revision - 1, - manifest_digest: [0x11; 32], - }; - let op = WriteOperationId::generate().unwrap(); - let key = key(); - let mut fs = StdFs; - let mut ctx = JournalDurableContext::new(&mut fs, &key, dir, &op); - let fence = MigrationFence::from_manifest(manifest); - promote_manifest_fenced(&mut ctx, manifest, &fence, Some(&previous)).unwrap(); - let bytes = std::fs::read(generation_path(dir, manifest.journal_revision)).unwrap(); - LiveMigration { - manifest_digest: content_digest(&bytes), - journal_revision: manifest.journal_revision, - ..previous - } -} - /// The committed owner stores `captured` into `journal`. pub fn store( fs: &mut F, diff --git a/crates/worldscript-secure-storage/tests/support/journal_fixture.rs b/crates/worldscript-secure-storage/tests/support/journal_fixture.rs new file mode 100644 index 000000000..247f7822f --- /dev/null +++ b/crates/worldscript-secure-storage/tests/support/journal_fixture.rs @@ -0,0 +1,207 @@ +//! Shared Gate 4D journal fixtures: a real committed journal directory, a file system that observes +//! what a store creates, and the keys and manifests the journal tests start from. + +use std::ffi::OsString; +use std::fs::File; +use std::io; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicU32, Ordering}; + +use worldscript_secure_storage::{ + content_digest, empty_inventory_digest, empty_journal_page_set_digest, generation_path, + operation_type, phase_code, promote_manifest_fenced, DirectoryDurability, DurableFs, + JournalDurableContext, JournalManifest, LiveMigration, MigrationFence, StdFs, WriteOperationId, +}; + +pub const OPERATION: &str = "store-op"; +pub const COMMITTED_REVISION: u64 = 3; + +pub struct TempDir(pub PathBuf); + +impl TempDir { + /// A fresh directory that this call created: a name that already exists (planted, or left by an + /// earlier run) is skipped, never reused, so `Drop` only ever removes what was created here. + pub fn new() -> Self { + static NEXT: AtomicU32 = AtomicU32::new(0); + loop { + let path = std::env::temp_dir().join(format!( + "wss-gate4d-store-{}-{}", + std::process::id(), + NEXT.fetch_add(1, Ordering::Relaxed) + )); + match std::fs::create_dir(&path) { + Ok(()) => return TempDir(path), + Err(error) if error.kind() == io::ErrorKind::AlreadyExists => continue, + Err(error) => panic!("cannot create the test directory {path:?}: {error}"), + } + } + } +} + +impl Drop for TempDir { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.0); + } +} + +/// A real file system that counts what a refused store must never do, logs directory syncs and can +/// fail the sync of one directory or the n-th file creation. +pub struct ObservedFs { + pub inner: StdFs, + pub creates: u32, + pub dirs_created: u32, + pub synced: Vec, + pub fail_sync_of: Option, + /// Fails the sync of every directory whose path contains this fragment. + pub fail_sync_matching: Option<&'static str>, + pub fail_create_at: Option, + pub fail_remove: bool, + /// Every path read, in order. + pub reads: Vec, +} + +impl ObservedFs { + pub fn new() -> Self { + Self { + inner: StdFs, + creates: 0, + dirs_created: 0, + synced: Vec::new(), + fail_sync_of: None, + fail_sync_matching: None, + fail_create_at: None, + fail_remove: false, + reads: Vec::new(), + } + } + + pub fn created_nothing(&self) -> bool { + (self.creates, self.dirs_created) == (0, 0) + } +} + +impl DurableFs for ObservedFs { + type File = File; + + fn create_new(&mut self, path: &Path) -> io::Result { + self.creates += 1; + if self.fail_create_at == Some(self.creates) { + return Err(io::Error::other("injected file creation failure")); + } + self.inner.create_new(path) + } + + fn sync_file(&mut self, file: &mut File) -> io::Result<()> { + self.inner.sync_file(file) + } + + fn read(&mut self, path: &Path) -> io::Result> { + self.reads.push(path.to_path_buf()); + self.inner.read(path) + } + + fn link_no_replace(&mut self, from: &Path, to: &Path) -> io::Result<()> { + self.inner.link_no_replace(from, to) + } + + fn remove_file(&mut self, path: &Path) -> io::Result<()> { + if self.fail_remove { + return Err(io::Error::other("injected file removal failure")); + } + self.inner.remove_file(path) + } + + fn sync_dir(&mut self, dir: &Path) -> io::Result { + self.synced.push(dir.to_path_buf()); + let matching = self.fail_sync_matching; + if matching.is_some_and(|fragment| dir.to_string_lossy().contains(fragment)) { + return Err(io::Error::other("injected directory sync failure")); + } + if self.fail_sync_of.as_deref() == Some(dir) { + return Err(io::Error::other("injected directory sync failure")); + } + self.inner.sync_dir(dir) + } + + fn list_dir(&mut self, dir: &Path) -> io::Result> { + self.inner.list_dir(dir) + } + + fn rename_replace(&mut self, from: &Path, to: &Path) -> io::Result<()> { + self.inner.rename_replace(from, to) + } + + fn create_dir_all(&mut self, dir: &Path) -> io::Result<()> { + self.dirs_created += 1; + self.inner.create_dir_all(dir) + } +} + +pub fn key() -> worldscript_secure_storage::Key { + worldscript_secure_storage::Key::from_bytes(&mut [9u8; 32]) +} + +/// A key that is not the journal key: what another epoch's pages are sealed under. +pub fn other_key() -> worldscript_secure_storage::Key { + worldscript_secure_storage::Key::from_bytes(&mut [10u8; 32]) +} + +pub fn manifest_at(revision: u64) -> JournalManifest { + JournalManifest { + operation_id: OPERATION.into(), + journal_revision: revision, + operation_type: operation_type::ROTATE, + phase: phase_code::DISCOVER, + source_epoch: 1, + target_epoch: 2, + has_target_root_key_ref: true, + target_root_key_ref_digest: Some([0x42; 32]), + fencing_generation: 7, + inventory_version: 1, + inventory_digest: empty_inventory_digest(1), + page_count: 0, + entry_count: 0, + journal_page_set_digest: empty_journal_page_set_digest(), + final_inventory_captured: false, + cursor_page_index: 0, + cursor_entry_index: 0, + has_lease_owner: false, + lease_owner_id: None, + lease_expires_unix_ms: None, + recovery_reason_code: 0, + } +} + +/// A journal directory whose root-named generation is a committed manifest. +pub struct Journal { + pub dir: TempDir, + pub live: LiveMigration, +} + +impl Journal { + pub fn path(&self) -> &Path { + &self.dir.0 + } +} + +/// Commits `manifest` as the root-named generation of `dir`; returns the binding that names it. +pub fn commit_into(dir: &Path, manifest: &JournalManifest) -> LiveMigration { + let previous = LiveMigration { + operation_id: manifest.operation_id.clone(), + fencing_generation: manifest.fencing_generation, + journal_revision: manifest.journal_revision - 1, + manifest_digest: [0x11; 32], + }; + let op = WriteOperationId::generate().unwrap(); + let key = key(); + let mut fs = StdFs; + let mut ctx = JournalDurableContext::new(&mut fs, &key, dir, &op); + let fence = MigrationFence::from_manifest(manifest); + promote_manifest_fenced(&mut ctx, manifest, &fence, Some(&previous)).unwrap(); + let bytes = std::fs::read(generation_path(dir, manifest.journal_revision)).unwrap(); + LiveMigration { + manifest_digest: content_digest(&bytes), + journal_revision: manifest.journal_revision, + ..previous + } +} diff --git a/docs/native/R15-SECURE-STORAGE-CONTRACT.md b/docs/native/R15-SECURE-STORAGE-CONTRACT.md index 99f980715..96445f6b4 100644 --- a/docs/native/R15-SECURE-STORAGE-CONTRACT.md +++ b/docs/native/R15-SECURE-STORAGE-CONTRACT.md @@ -2420,6 +2420,21 @@ with. Entering `ADMIT` does not set it, so after a crash or restart a journal ca inventory carried forward from a final one actually captured behind the barrier. The final inventory is the commit inventory and immutable once captured, so it is captured once. +**Capturing an inventory too large for memory.** The pages of a capture need not be held at once. A page +is sealed once (a second sealing would use another nonce and bind a different digest), so its envelope is +staged on disk in a pending directory of its own, named for what it is and never as a digest directory, +while the page references and the inventory digest accumulate page by page. The inventory digest commits +to the total entry count before any entry (§5.4), so the total is announced when the capture starts and +the capture ends only when exactly that many entries were staged. The builder assigns the consecutive +page indexes and the capture's generation, entries must be valid and strictly ascending across pages, and +the successor manifest is built and checked exactly as for pages held in memory. Staged files are not +authority: the manifest never names them, no reader looks for them, a staged page is trusted only after +it is read back by the exact path of its reference and confirmed against that reference, and a staged +page that is missing is a broken attempt to abandon, never a recovery state of the journal. A capture +belongs to one journal directory and to the key that opened the manifest it started from; a page is never +staged under another context. A failure the caller can still handle removes the pages staged so far, each only if it still +holds the bytes that were staged, and an attempt that was killed leaves files that nothing trusts. Promoting the staged pages to the digest +directory, under the root lock, is a separate step. **Journal envelope epoch.** The manifest and every page of one operation are sealed under a single key epoch that is stable for the whole operation and derived from the authenticated manifest alone, never from the authority root's `active_key_epoch`, which moves at cutover: `ENABLE` seals under diff --git a/docs/native/r15/GATE4D-JOURNAL-DURABLE-EVIDENCE.md b/docs/native/r15/GATE4D-JOURNAL-DURABLE-EVIDENCE.md index 981920e25..e5ee8f563 100644 --- a/docs/native/r15/GATE4D-JOURNAL-DURABLE-EVIDENCE.md +++ b/docs/native/r15/GATE4D-JOURNAL-DURABLE-EVIDENCE.md @@ -665,6 +665,66 @@ after authentication (D2b-2), and resolving the journal key through the authenti readability proof for a root already at the target epoch (D2b-3a), used by every journal-owner operation (D2b-3b). Each is described in its own section below. +## Streaming capture, stage one — building and staging one page at a time + +The store and `capture_inventory` hold every sealed page in memory, and the format allows a million entries +of up to a kilobyte in pages of up to 4096 entries. A page cannot be written to its final place as it is +sealed: that place is keyed by the page-set digest, which needs every page's envelope digest, and sealing +uses a fresh nonce, so a second pass would bind different bytes. Stage one therefore keeps each envelope on +disk in a private pending directory, `inventory/pending--/page-/generation-.wsr1` +(the layout of a digest directory, so stage two stores the same bytes), and keeps in memory one page and one +reference per page. + +| Step | Rule | +|---|---| +| `StreamedCapture::begin` | fence; the caller is the committed owner writing under the committed manifest (`StaleMigrationOwner`); the manifest is the exact root-named generation (`LiveBindingMismatch`); the inventory is open (`DISCOVER` or `ADMIT`, not yet final, nothing converted); the announced total is within the format bound. Nothing is created. These are early refusals: stage two repeats them under the locks | +| `push_page` | the context is the one the capture started under: the same journal directory (`LiveBindingMismatch`) and a key that opens the manifest the capture started from (`Open(Tampered)`), checked before anything is sealed over the few kilobytes of root-named manifest that `begin` kept, so a push reads nothing for it; the page takes the next index and the capture's generation; it is non-empty, its entries are valid and strictly ascending after every earlier entry, the running total stays within the announced one; it is sealed once under the operation's journal envelope epoch, its reference recorded, its envelope staged and synced, and the bytes dropped. A failure removes the pages staged so far and the attempt's own staging link, each only if it still holds exactly the bytes that were staged (checked against its reference), so a file this attempt did not write, or a staged page that was replaced, is never deleted | +| `finish` | exactly the announced total was staged; the successor is built by the same function `capture_inventory` ends in, and checked as a valid successor; a failure removes the staged pages | +| `StagedCapture::discard` | removes the staged page files of a capture that will not be promoted | +| `load_staged_page` | the exact path of the reference, bounded; a missing file is `Corrupt`, never `RecoveryRequired`, because a staged page is a candidate and not authority; the envelope hashes to the reference before it is opened; it opens only as this operation's page at this index and generation, with the epoch compared before the key is used | + +Differences from the admission, disclosed: (1) the inventory digest commits to the entry total before any +entry, so `begin` takes the total up front and the caller counts before it streams (the admission did not +say so); (2) `begin` takes a journal context that carries the key, and the authority-level wrapper that +routes the key from the root arrives with stage two, because the routed key has to be resolved again at +promotion anyway; (3) the admission planned a counting file system to bound the page size, but peak memory is +not measured: it is bounded by construction, because `StreamedCapture` has no field that can hold a page +or an envelope, and a push drops both before it returns; (4) the tests are a file of their own, +`gate4d_stream_capture_test`, over the journal fixtures that were split out of the inventory fixtures +(`support/journal_fixture.rs`), so a test that needs only a committed journal does not carry the page +fixtures; putting them into the store test made that file lose its cohesion. + +Cleanup is best effort and removes files only: the file system abstraction cannot remove a directory, so the +empty page directories stay, and the files of an attempt that was killed stay as well. The pages are staged under +the attempt's own random identity, not the caller's operation id, so the name of a staging link that a +promotion leaves behind (after a failed promotion, or after a successful one whose link could not be +unlinked) is known to the cleanup, which removes it under the same rule. The check of a file against its digest and its +removal are two operations (the abstraction has no unlink by handle), which is acceptable because the pending +directory is private, named by a fresh random tag and never authority: the check guards against files a +failure found in place or that were replaced by accident, not against a writer with access to the journal +directory, who could alter the authoritative files directly. Both are inert and +are reclaimed by the separate cleanup recorded below. `begin` and `push_page` carry the capture's binding to its +journal and key because a context is supplied on every call and the key is never compared directly. + +Proof (eighteen tests in `gate4d_stream_capture_test`): the streamed successor equals what `capture_inventory` +builds from the pages that were staged, for one entry, uneven and even splits, one page, the final capture +in `ADMIT` and a full page plus one (4097 entries); an empty inventory stages nothing and matches the empty +capture; the staged files are exactly the pending layout under a `pending-` name; the staged pages are +accepted by the existing store; every way a capture cannot start (a capture already final, conversion +started, a phase without a snapshot, a finished journal, a total above the bound, a stale owner, a manifest +that is not the root-named one) is refused with nothing created; an out-of-order page, a total exceeded, a +total not reached and an empty page are refused; a failure partway removes the pages it staged, also when +the failing page was already promoted, and a new attempt is independent; a capture that ends short and a +discarded capture leave no page file, and a removal that fails is not an error; a push under another key +or into a journal directory that holds an identical copy of the root-named generation is refused before +anything is staged and a push reads the root-named manifest no more than `begin` did; a failure whose cause is a different file in the slot of the page in flight leaves that file alone, and so does the cleanup of a staged page that was replaced, and the staging link of the failed promotion is removed, as is one that a successful promotion left behind; two attempts under one operation id never share a directory; a staged page that +changed, was swapped, is missing, is out of range or is read under another key is refused. Mutation-checked: +no digest absorption, no root-named check, no open-inventory check, no `pending-` prefix, a pending +directory named by the caller's operation id, a shifted page directory, no digest check on read, no +context check (the directory and the key each on their own), no cleanup on a failed push, on a failed +finish or in `discard`, no slot for the page in flight, an unconditional removal of a staged file, no removal of the staging link, staging links named by the caller's operation id, a manifest read on every push, and a missing staged page reported as +`RecoveryRequired` each fail the test that owns them. Existing capture, store, reader and durable tests are unchanged and pass. + ## Slice D3b — the owner renews its own lease After D3a every ordinary successor and the plain binding advance refuse a lease change, so a live owner could @@ -866,7 +926,8 @@ this API's reach; the journal key route of D2b-3 fails closed on a `Revoked` or - Conversion (C2+) over the verified page set: it must require exclusive admission by construction and re-read the root-bound manifest with `final_inventory_captured = 1` (maintainer decision D). - Write barrier of the final capture: `commit_inventory_capture` takes no admission guard; the barrier is the durable `ADMIT` phase the orchestrator establishes by draining writers, and the write path must refuse ordinary mutating writes by that phase (`ordinary_mutating_writes_admitted`) before the final capture has a caller (Gate 4E/5; acceptance criterion on #359). - Inheriting unchanged pages: the C1b-2 reader now returns the authenticated page references, so the store may accept a page that keeps an earlier generation if those references name exactly its bytes (acceptance criterion on #359, a follow-up slice). Until then every page of a capture is rewritten at the new revision. -- Streaming capture: `promote_inventory_set_fenced` verifies the set in memory; a one-page-at-a-time seal, digest and promote is needed before very large inventories (acceptance criterion on #359). +- Streaming capture, stage two: the promotion of the staged pages to the digest directory under the root lock and the journal mutex (re-reading and re-verifying each page, adopting on retry), the authority-level wrapper that routes the key, and the composed commit; until then `commit_inventory_capture` still holds the set in memory (acceptance criterion on #359). +- Reclaiming abandoned pending directories: an attempt that was killed, or a finished capture dropped without `discard`, leaves inert files under `inventory/pending-*`, and every attempt leaves its empty page directories; none is authority or ever read. Reclaiming them needs a directory-removal primitive and a sweep that knows no live attempt owns them, as for the orphaned digest directories of a discarded capture (acceptance criterion on #359). - The cross-process lease CAS. - Mixed-key conversion, Gate 4E/5/6/7, production authority switch. diff --git a/docs/native/r15/GATE4D-SLICE-B-GAP-MATRIX.md b/docs/native/r15/GATE4D-SLICE-B-GAP-MATRIX.md index e10943ab6..af4bbfb25 100644 --- a/docs/native/r15/GATE4D-SLICE-B-GAP-MATRIX.md +++ b/docs/native/r15/GATE4D-SLICE-B-GAP-MATRIX.md @@ -89,6 +89,7 @@ load of the newer generation yet). - Epoch relation: `journal_envelope_epoch` refuses a `ROTATE` whose target is not above its source (§8.3 item 2) and an `ENVELOPE_MIGRATION` whose target differs from its source (§10.4, stated by this slice: it keeps the key epoch, a newer epoch is a rotation), so such a manifest neither encodes nor decodes; closes the pre-caller validation criterion before any caller writes a rotation journal. - Slice D3a: the successor relation constrains the cursor and the lease across manifests (maintainer decision C): a forward phase change enters the new phase at `(0, 0)` (`CursorNotReset`), entering `RECOVERY_REQUIRED` keeps the last cursor, and an ordinary successor leaves the lease alone, so the owner and the fence change only by takeover; `assert_renewal_successor` is the predicate of the same-owner renewal (expiry strictly forward), wired in D3b. - Slice D3b: `commit_lease_renewal` is the composed operation of the same-owner renewal: its own fenced publish (`publish_renewal_fenced`, so the ordinary checkpoint and the plain binding advance still refuse a lease change), the journal key through the registry route, `BindingStep::Renewal` re-proving the renewal relation at the root under the lock, adoption of an identical candidate on retry; the root lock settles a renewal racing a takeover. No orchestrator calls it yet. +- Streaming capture, stage one: `StreamedCapture` builds and stages a captured inventory one page at a time (the total announced up front, because the inventory digest commits to it first), each page sealed once and its envelope staged in a private `inventory/pending-*` directory that is never authority; the successor is built by the same function as `capture_inventory`; `load_staged_page` confirms a staged page against its reference and treats a missing one as a broken attempt, not a recovery state. A capture is bound to its journal directory and key, and a handled failure removes the pages staged so far. Promotion under the root lock, the key-routing wrapper and the composed commit are stage two. - Slice C: record conversion / mixed-key inventory execution as live truth requires - Gate 4E: first enable/disable closure - Root two-phase commit coupling with step F when journal + root must advance together (after B2)