diff --git a/CHANGELOG.md b/CHANGELOG.md index 8e76a122d..271e6e189 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +- **R-15 Gate 4D:** commit a staged inventory (stage two, part a, of the streaming capture). + `commit_streamed_inventory_capture` promotes the pages that `StreamedCapture` staged into the digest + directory of their successor one page at a time (each page confirmed against its reference, its entries + absorbed into the inventory digest, an identical page already there adopted), publishes the successor and + advances the root binding, in the order and under the guarantees of `commit_inventory_capture`, which now + shares one core with it. The staged files are removed only after the root has committed, so a retry after a + failure still has them; a staged page found wrong stops the commit before any manifest is published. The + key-routing staging session is the next part. PR #1011. - **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 diff --git a/crates/worldscript-secure-storage/src/authority.rs b/crates/worldscript-secure-storage/src/authority.rs index 3105603fe..24abb3ee2 100644 --- a/crates/worldscript-secure-storage/src/authority.rs +++ b/crates/worldscript-secure-storage/src/authority.rs @@ -43,10 +43,11 @@ use crate::identity::RecordIdentity; use crate::journal::{ assert_binding_successor, assert_binding_takeover, assert_capture_successor, assert_fence, assert_progress_successor, assert_renewal_successor, assert_takeover_successor, - load_authoritative_manifest, promote_inventory_set_fenced, publish_manifest_fenced, - publish_renewal_fenced, publish_takeover_fenced, CandidateConflict, InventorySetWrite, - JournalDurableContext, JournalDurableError, JournalManifest, JournalTakeover, - MigrationExecutionError, MigrationFence, SealedPage, + load_authoritative_manifest, promote_inventory_set_fenced, promote_staged_inventory_fenced, + publish_manifest_fenced, publish_renewal_fenced, publish_takeover_fenced, CandidateConflict, + InventorySetWrite, JournalDurableContext, JournalDurableError, JournalManifest, + JournalTakeover, MigrationExecutionError, MigrationFence, SealedPage, StagedCapture, + StagedPromotion, }; use crate::journal_route::{resolve_journal_key, JournalRoute, JournalRouteError}; use crate::marker::content_digest; @@ -425,6 +426,121 @@ pub fn commit_inventory_capture( capture: InventoryCapture<'_>, ) -> Result { let checkpoint = capture.checkpoint; + let plan = CapturePlan { + layout, + checkpoint, + committed_revision: capture.committed_manifest.journal_revision, + }; + commit_capture(fs, provider, plan, |journal, committed, fence| { + promote_inventory_set_fenced( + journal, + &InventorySetWrite { + committed_manifest: capture.committed_manifest, + fence, + committed: Some(committed), + successor: checkpoint.manifest, + pages: capture.pages, + }, + ) + }) +} + +/// One capture of the journal's inventory whose pages were staged one at a time (§10.3): the staged +/// capture, the manifest the root binding names, and the owner's token and routes. +#[derive(Clone, Copy)] +pub struct StreamedInventoryCapture<'a> { + /// The finished stage-one capture; its successor is the manifest that is published. + pub staged: &'a StagedCapture, + /// The manifest the root binding names, which the staged capture was begun over. + pub committed_manifest: &'a JournalManifest, + /// The owner's token for the successor. + pub fence: &'a MigrationFence, + pub journal: JournalSource<'a>, + pub root_key_ref: &'a RootKeyRefV1, + pub active_key_epoch: u64, + /// What to do with a durable candidate at the next revision that is not the successor. + pub conflict: CandidateConflict, +} + +/// Captures a staged inventory: promotes its pages into the digest directory, publishes the manifest +/// that names them and advances the root binding to it, in the order and under the guarantees of +/// [`commit_inventory_capture`] (§10.1.1, §10.3). +/// +/// The pages are read back from the staged capture, confirmed against their references and stored one +/// at a time by [`promote_staged_inventory_fenced`], so memory stays one page. The staged files are +/// removed only once the root has committed to the set (best effort), so a retry after a failure +/// before that point still has them and adopts what is already durable; after a success the handle +/// is spent. The journal key is routed from the root here, as for every journal-owner commit. +pub fn commit_streamed_inventory_capture( + fs: &mut F, + provider: &mut P, + layout: RootLayout<'_>, + capture: StreamedInventoryCapture<'_>, +) -> Result { + let staged = capture.staged; + let checkpoint = JournalCheckpoint { + manifest: staged.successor(), + fence: capture.fence, + journal: capture.journal, + root_key_ref: capture.root_key_ref, + active_key_epoch: capture.active_key_epoch, + conflict: capture.conflict, + }; + let plan = CapturePlan { + layout, + checkpoint, + committed_revision: capture.committed_manifest.journal_revision, + }; + let committed = commit_capture(fs, provider, plan, |journal, committed, fence| { + promote_staged_inventory_fenced( + journal, + &StagedPromotion { + committed_manifest: capture.committed_manifest, + fence, + committed, + staged, + }, + ) + })?; + staged.remove_files(fs); + Ok(committed) +} + +/// What a capture commit needs besides the file system, the provider and the step that stores the +/// pages: where the root lives, the checkpoint to publish, and the revision of the committed manifest +/// the pages are written under. +struct CapturePlan<'a> { + layout: RootLayout<'a>, + checkpoint: JournalCheckpoint<'a>, + committed_revision: u64, +} + +/// The commit shared by the capture from pages in memory and the capture of a staged inventory. +/// +/// Everything runs under one `root_commit_mutex`. The committed binding is read from the root and the +/// journal key is routed before any journal write; the pages are stored by `store_pages` before the +/// manifest that names them is published, so the manifest is never written before its pages are +/// durable; and the binding advances as a capture, the one change of the inventory fields it allows. +fn commit_capture( + fs: &mut F, + provider: &mut P, + plan: CapturePlan<'_>, + store_pages: S, +) -> Result +where + F: DurableFs, + P: KeyProvider, + S: FnOnce( + &mut JournalDurableContext<'_, F>, + &LiveMigration, + &MigrationFence, + ) -> Result, +{ + let CapturePlan { + layout, + checkpoint, + committed_revision, + } = plan; check_operation_id(&checkpoint.manifest.operation_id) .map_err(|_| AuthorityError::InvalidOperationId)?; // The caller's token must be the successor's before anything is written, or the pages could be @@ -439,22 +555,13 @@ pub fn commit_inventory_capture( // generation at the committed revision; the successor is published under the caller's own fence. let committed_fence = MigrationFence { fencing_generation: checkpoint.fence.fencing_generation, - journal_revision: capture.committed_manifest.journal_revision, + journal_revision: committed_revision, }; let (published, pages_durability) = { let mut journal = journal_context(&mut *fs, checkpoint.journal, &key, checkpoint.conflict).for_capture(); - let pages_durability = promote_inventory_set_fenced( - &mut journal, - &InventorySetWrite { - committed_manifest: capture.committed_manifest, - fence: &committed_fence, - committed: Some(&committed), - successor: checkpoint.manifest, - pages: capture.pages, - }, - ) - .map_err(AuthorityError::Journal)?; + let pages_durability = store_pages(&mut journal, &committed, &committed_fence) + .map_err(AuthorityError::Journal)?; let published = publish_manifest_fenced( &mut journal, checkpoint.manifest, diff --git a/crates/worldscript-secure-storage/src/journal/inventory_store.rs b/crates/worldscript-secure-storage/src/journal/inventory_store.rs index f8f2d22c5..b8437a64f 100644 --- a/crates/worldscript-secure-storage/src/journal/inventory_store.rs +++ b/crates/worldscript-secure-storage/src/journal/inventory_store.rs @@ -169,7 +169,7 @@ fn verify_set<'a, F: DurableFs>( /// references, which only a verified reader of the stored set can provide (the next slice). Until then /// an older generation is refused, so one page identity and generation can never stand for different /// content across page sets. -fn assert_page_generation( +pub(super) fn assert_page_generation( successor: &JournalManifest, page: &JournalPage, ) -> Result<(), JournalDurableError> { @@ -329,7 +329,7 @@ fn sync_chain( Ok(durability) } -fn both(left: DirectoryDurability, right: DirectoryDurability) -> DirectoryDurability { +pub(super) fn both(left: DirectoryDurability, right: DirectoryDurability) -> DirectoryDurability { if left == DirectoryDurability::Confirmed && right == DirectoryDurability::Confirmed { DirectoryDurability::Confirmed } else { diff --git a/crates/worldscript-secure-storage/src/journal/mod.rs b/crates/worldscript-secure-storage/src/journal/mod.rs index 73b528d6d..1a9e38cae 100644 --- a/crates/worldscript-secure-storage/src/journal/mod.rs +++ b/crates/worldscript-secure-storage/src/journal/mod.rs @@ -17,6 +17,7 @@ mod page; mod renewal; mod state; mod stream_capture; +mod stream_promote; mod succession; mod takeover; mod wire; @@ -150,6 +151,7 @@ pub use state::{ pub use stream_capture::{ load_staged_page, CaptureStart, StagedCapture, StagedPage, StreamedCapture, }; +pub use stream_promote::{promote_staged_inventory_fenced, StagedPromotion}; 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 index 5f7d90052..55527031d 100644 --- a/crates/worldscript-secure-storage/src/journal/stream_capture.rs +++ b/crates/worldscript-secure-storage/src/journal/stream_capture.rs @@ -89,6 +89,8 @@ impl std::fmt::Debug for StreamedCapture { #[derive(Debug, PartialEq, Eq)] pub struct StagedCapture { successor: JournalManifest, + /// The journal directory the pages were staged in, which they may only be promoted into. + journal_dir: PathBuf, pending: PathBuf, tag: WriteOperationId, refs: Vec, @@ -128,10 +130,20 @@ impl StagedCapture { &self.pending } + /// The journal directory the capture was staged in and belongs to. + pub fn journal_dir(&self) -> &Path { + &self.journal_dir + } + /// 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); + self.remove_files(ctx.fs); + } + + /// The same removal for a caller that keeps the handle, such as the commit that has just spent it. + pub(crate) fn remove_files(&self, fs: &mut F) { + discard_staged(fs, &self.pending, &self.refs, &self.tag); } } @@ -200,7 +212,7 @@ impl StreamedCapture { match self.stage_next(ctx, entries) { Ok(()) => Ok(self), Err(error) => { - discard_staged(ctx, &self.pending, &self.refs, &self.tag); + discard_staged(ctx.fs, &self.pending, &self.refs, &self.tag); Err(error) } } @@ -231,9 +243,13 @@ impl StreamedCapture { // 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); + remove_if_staged(ctx.fs, &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); + remove_if_staged( + ctx.fs, + &staging_path(&dir, self.revision, &self.tag), + digest, + ); return Err(error); } self.refs.push(reference); @@ -265,6 +281,7 @@ impl StreamedCapture { ) -> Result { let Self { committed, + journal_dir, entry_count, pending, tag, @@ -286,12 +303,13 @@ impl StreamedCapture { match successor { Ok(successor) => Ok(StagedCapture { successor, + journal_dir, pending, tag, refs, }), Err(error) => { - discard_staged(ctx, &pending, &refs, &tag); + discard_staged(ctx.fs, &pending, &refs, &tag); Err(error) } } @@ -304,7 +322,7 @@ impl StreamedCapture { /// 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>, + fs: &mut F, pending: &Path, refs: &[JournalPageRef], tag: &WriteOperationId, @@ -313,22 +331,18 @@ fn discard_staged( 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); + remove_if_staged(fs, &generation_path(&dir, generation), digest); + remove_if_staged(fs, &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); +fn remove_if_staged(fs: &mut F, path: &Path, digest: &[u8; 32]) { + let found = 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); + let _ = fs.remove_file(path); } } diff --git a/crates/worldscript-secure-storage/src/journal/stream_promote.rs b/crates/worldscript-secure-storage/src/journal/stream_promote.rs new file mode 100644 index 000000000..71e82da8e --- /dev/null +++ b/crates/worldscript-secure-storage/src/journal/stream_promote.rs @@ -0,0 +1,87 @@ +//! Gate 4D streaming capture, stage two: promoting a staged capture into the digest directory +//! (§10.1.1, §10.3). +//! +//! Stage one ([`StreamedCapture`](super::StreamedCapture)) leaves the pages of a finished inventory in +//! a private pending directory. This is the step that makes them reachable: each staged page is read +//! back, confirmed against its reference, and stored under +//! `inventory//page-/`, exactly where +//! [`promote_inventory_set_fenced`](super::promote_inventory_set_fenced) would have put the same +//! bytes, and with the same adoption of an identical page and refusal to replace a different file. +//! +//! Everything that can be refused from metadata alone is refused before the first write: the caller +//! must be the committed owner writing under the committed manifest, that manifest must be the exact +//! root-named generation, the successor must be the capture successor of it, and the staged +//! references must be the page set the successor names. The pages themselves are checked as they are +//! loaded, one at a time, so memory stays one page. That makes this a single pass: a staged page found +//! wrong midway leaves the verified prefix as an inert directory under the digest, as an I/O failure +//! partway through the in-memory store already does, and the inventory digest is confirmed before +//! this returns, so a manifest naming the set is never published on a set that failed it. + +use crate::durable::{DirectoryDurability, DurableFs}; +use crate::root::LiveMigration; + +use super::capture::{assert_capture_successor, assert_page_not_empty, SealedPage}; +use super::digest::InventoryDigestVerifier; +use super::durable::{ + assert_root_named_manifest, with_fence, JournalDurableContext, JournalDurableError, +}; +use super::inventory_store::{assert_page_generation, both, inventory_page_dir, store_page_at}; +use super::manifest::JournalManifest; +use super::state::{assert_page_promote_authority, MigrationExecutionError, MigrationFence}; +use super::stream_capture::{load_staged_page, StagedCapture}; + +/// A finished streamed capture to promote: the manifest the root binding names, the owner's token +/// under it, that binding, and the staged capture built over that manifest. +#[derive(Clone, Copy)] +pub struct StagedPromotion<'a> { + /// The manifest the root binding names: the pages are written under it (R4's page rule). + pub committed_manifest: &'a JournalManifest, + pub fence: &'a MigrationFence, + pub committed: &'a LiveMigration, + pub staged: &'a StagedCapture, +} + +/// Stores the pages of a staged capture under the digest directory of its successor, fenced, and +/// returns whether every directory sync was confirmed. The staged files are left where they are: the +/// caller removes them after the root has committed to the set, so that a retry still has them. +pub fn promote_staged_inventory_fenced( + ctx: &mut JournalDurableContext<'_, F>, + promotion: &StagedPromotion<'_>, +) -> Result { + let committed = promotion.committed_manifest; + let successor = promotion.staged.successor(); + // The capture belongs to the journal it was staged in: its pages are read from there by absolute + // path, so a context for another directory would copy them across. Refused before anything is read. + if ctx.dir != promotion.staged.journal_dir() { + return Err(JournalDurableError::Authority( + MigrationExecutionError::LiveBindingMismatch, + )); + } + with_fence(committed, promotion.fence, || { + assert_page_promote_authority(committed, Some(promotion.committed)) + .map_err(JournalDurableError::Authority)?; + assert_root_named_manifest(ctx, committed, promotion.committed)?; + assert_capture_successor(committed, successor).map_err(JournalDurableError::Authority)?; + // The successor must be a manifest that can be sealed, and the staged references the page + // set it names. + successor.encode()?; + successor.verify_page_set(promotion.staged.page_refs())?; + let mut verifier = + InventoryDigestVerifier::new(successor.inventory_version, successor.entry_count)?; + let mut durability = DirectoryDurability::Confirmed; + for index in 0..successor.page_count { + let staged = load_staged_page(ctx, promotion.staged, index)?; + assert_page_not_empty(&staged.page)?; + assert_page_generation(successor, &staged.page)?; + verifier.absorb_page(&staged.page)?; + let sealed = SealedPage { + page: &staged.page, + envelope: &staged.envelope, + }; + let dir = inventory_page_dir(ctx.dir, &successor.journal_page_set_digest, index); + durability = both(durability, store_page_at(ctx, &dir, committed, &sealed)?); + } + verifier.finish(successor.inventory_digest)?; + Ok(durability) + }) +} diff --git a/crates/worldscript-secure-storage/src/lib.rs b/crates/worldscript-secure-storage/src/lib.rs index b60e5308a..a17d18f09 100644 --- a/crates/worldscript-secure-storage/src/lib.rs +++ b/crates/worldscript-secure-storage/src/lib.rs @@ -61,10 +61,11 @@ pub use admission::{ }; pub use authority::{ advance_live_migration, commit_catalog_change, commit_inventory_capture, - commit_journal_checkpoint, commit_journal_takeover, commit_lease_renewal, AuthorityError, - BindingAdvance, CatalogChange, CatalogCommit, CatalogRecoveryReason, CatalogStep, - CommittedShard, InventoryCapture, JournalCheckpoint, JournalSource, JournalTakeoverCommit, - LoadedCatalog, + commit_journal_checkpoint, commit_journal_takeover, commit_lease_renewal, + commit_streamed_inventory_capture, AuthorityError, BindingAdvance, CatalogChange, + CatalogCommit, CatalogRecoveryReason, CatalogStep, CommittedShard, InventoryCapture, + JournalCheckpoint, JournalSource, JournalTakeoverCommit, LoadedCatalog, + StreamedInventoryCapture, }; #[cfg(feature = "test-support")] pub use authority::{list_records, load_catalog}; @@ -99,8 +100,8 @@ pub use journal::{ 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, + promote_staged_inventory_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, CaptureStart, ForeignInventoryExtension, InventoryDigestVerifier, InventorySetWrite, JournalCheckpointCursor, JournalDurableContext, @@ -108,7 +109,7 @@ pub use journal::{ JournalInventoryExtent, JournalInventorySource, JournalManifest, JournalPage, JournalPageRef, JournalRevision, JournalTakeover, ManifestEnvelopeDigest, ManifestRead, MigrationExecutionError, MigrationFence, MigrationPhase, PublishedManifest, RecoveryReasonCode, - SealedPage, StagedCapture, StagedPage, StreamedCapture, VerifiedInventory, + SealedPage, StagedCapture, StagedPage, StagedPromotion, 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_journal_route_test.rs b/crates/worldscript-secure-storage/tests/gate4d_journal_route_test.rs index fd7d5dee9..2a8c2418d 100644 --- a/crates/worldscript-secure-storage/tests/gate4d_journal_route_test.rs +++ b/crates/worldscript-secure-storage/tests/gate4d_journal_route_test.rs @@ -16,15 +16,17 @@ use worldscript_secure_storage::memory_provider::MemoryKeyProvider; use worldscript_secure_storage::{ advance_live_migration, commit_catalog_change, commit_inventory_capture, commit_journal_checkpoint, commit_journal_takeover, commit_lease_renewal, commit_root, - content_digest, empty_inventory_digest, empty_journal_page_set_digest, generation_path, - load_catalog, operation_type, phase_code, resolve_journal_key, write_key_epoch, AuthorityError, - BindingAdvance, CandidateConflict, CatalogChange, CatalogCommit, InstallationScopeId, - InventoryCapture, JournalCheckpoint, JournalDurableError, JournalError, JournalManifest, + commit_streamed_inventory_capture, content_digest, empty_inventory_digest, + empty_journal_page_set_digest, generation_path, load_catalog, operation_type, phase_code, + resolve_journal_key, write_key_epoch, AuthorityError, BindingAdvance, CandidateConflict, + CaptureStart, CatalogChange, CatalogCommit, InstallationScopeId, InventoryCapture, + JournalCheckpoint, JournalDurableContext, JournalDurableError, JournalError, JournalManifest, JournalRoute, JournalRouteError, JournalSource, JournalTakeoverCommit, Key, KeyEpochCommit, KeyEpochRecord, KeyEpochStatus, KeyProvider, LiveMigration, ManifestRead, MigrationExecutionError, MigrationFence, OpenError, RecordClass, RecordIdentity, RecordMeta, RootBody, RootCommitEvidence, RootCommitGuard, RootCommitRequest, RootCommitState, - RootKeyRefV1, RootLayout, StdFs, WriteOperationId, JOURNAL_MANIFEST_RECORD_SCHEMA, + RootKeyRefV1, RootLayout, StagedCapture, StdFs, StreamedCapture, StreamedInventoryCapture, + WriteOperationId, JOURNAL_MANIFEST_RECORD_SCHEMA, }; const OPERATION: &str = "route-op"; @@ -311,11 +313,12 @@ impl Fixture { commit_root(&mut StdFs, &mut self.provider, layout, request).unwrap(); } - /// The journal in place and a root that binds it. - fn bound(&mut self) { + /// The journal in place and a root that binds it; returns the binding. + fn bound(&mut self) -> LiveMigration { let live = self.store_journal(); self.commit_first_root(); self.bind(&live); + live } /// Whether `key` opens the stored rotation manifest, which only the source epoch's key does. @@ -535,6 +538,51 @@ impl Fixture { } } + /// A finished stage-one capture of an empty inventory over the bound rotation, staged in the journal + /// directory under the journal's own key (staging takes the key from its caller). + fn stage_capture(&self, live: &LiveMigration) -> StagedCapture { + let (key, op) = (source_key(), WriteOperationId::generate().unwrap()); + let journal_dir = self.journal_dir(); + let committed = rotation(); + let start = CaptureStart { + committed_manifest: &committed, + fence: &MigrationFence::from_manifest(&committed), + live, + entry_count: 0, + }; + let mut fs = StdFs; + let mut ctx = JournalDurableContext::new(&mut fs, &key, &journal_dir, &op); + let capture = StreamedCapture::begin(&mut ctx, &start).unwrap(); + capture.finish(&mut ctx).unwrap() + } + + /// The commit of a staged capture over valid inputs, which is expected to be refused: the key is + /// routed from the root, never taken from the caller. + fn streamed_capture_is_refused(&mut self, staged: &StagedCapture) -> AuthorityError { + let root_ref = self.root_ref().clone(); + let (root_dir, journal_dir) = (self.root_dir(), self.journal_dir()); + let committed = rotation(); + let fence = MigrationFence::from_manifest(staged.successor()); + let operation = WriteOperationId::generate().unwrap(); + let capture = StreamedInventoryCapture { + staged, + committed_manifest: &committed, + fence: &fence, + journal: JournalSource { + dir: &journal_dir, + operation: &operation, + }, + root_key_ref: &root_ref, + active_key_epoch: self.root_epoch, + conflict: CandidateConflict::Refuse, + }; + let layout = RootLayout { + root_dir: &root_dir, + }; + commit_streamed_inventory_capture(&mut StdFs, &mut self.provider, layout, capture) + .unwrap_err() + } + /// The five composed journal operations over valid inputs for each (a checkpoint and a capture of /// the next revision, a takeover claim by a new owner at fence plus one, an advance to the stored /// successor, a renewal of the owner's lease): each is expected to be refused, and the error of @@ -622,10 +670,13 @@ fn a_revoked_journal_epoch_refuses_every_journal_operation_before_a_write() { fixture.register(SOURCE_EPOCH, KeyEpochStatus::Active, 1); fixture.register(SOURCE_EPOCH, KeyEpochStatus::Revoked, 2); fixture.register(TARGET_EPOCH, KeyEpochStatus::Active, 1); - fixture.bound(); + let live = fixture.bound(); let stored = fixture.store_successor(); + let staged = fixture.stage_capture(&live); let before = fixture.snapshot(); - for error in fixture.every_operation_is_refused(&stored) { + let mut errors = fixture.every_operation_is_refused(&stored); + errors.push(fixture.streamed_capture_is_refused(&staged)); + for error in errors { assert_eq!( error, AuthorityError::JournalRoute(JournalRouteError::EpochRevoked(SOURCE_EPOCH)) @@ -638,10 +689,13 @@ fn a_revoked_journal_epoch_refuses_every_journal_operation_before_a_write() { fn an_unregistered_journal_epoch_refuses_every_journal_operation_before_a_write() { let mut fixture = Fixture::new(TARGET_EPOCH); fixture.register(TARGET_EPOCH, KeyEpochStatus::Active, 1); - fixture.bound(); + let live = fixture.bound(); let stored = fixture.store_successor(); + let staged = fixture.stage_capture(&live); let before = fixture.snapshot(); - for error in fixture.every_operation_is_refused(&stored) { + let mut errors = fixture.every_operation_is_refused(&stored); + errors.push(fixture.streamed_capture_is_refused(&staged)); + for error in errors { assert_eq!( error, AuthorityError::JournalRoute(JournalRouteError::EpochNotRegistered(SOURCE_EPOCH)) @@ -657,11 +711,14 @@ fn a_route_to_another_key_is_refused_by_the_authenticated_load_before_a_write() let mut fixture = Fixture::with_source_material(TARGET_EPOCH, [0x55; 32]); fixture.register(SOURCE_EPOCH, KeyEpochStatus::RetiredRecoveryOnly, 1); fixture.register(TARGET_EPOCH, KeyEpochStatus::Active, 1); - fixture.bound(); + let live = fixture.bound(); let stored = fixture.store_successor(); + let staged = fixture.stage_capture(&live); let before = fixture.snapshot(); // Every operation has valid inputs, so the key is the only thing left to refuse them. - for error in fixture.every_operation_is_refused(&stored) { + let mut errors = fixture.every_operation_is_refused(&stored); + errors.push(fixture.streamed_capture_is_refused(&staged)); + for error in errors { assert_eq!( error, AuthorityError::Journal(JournalDurableError::Journal(JournalError::Open( diff --git a/crates/worldscript-secure-storage/tests/gate4d_root_binding_test.rs b/crates/worldscript-secure-storage/tests/gate4d_root_binding_test.rs index 05941d66a..371b892a4 100644 --- a/crates/worldscript-secure-storage/tests/gate4d_root_binding_test.rs +++ b/crates/worldscript-secure-storage/tests/gate4d_root_binding_test.rs @@ -13,18 +13,20 @@ use worldscript_secure_storage::memory_provider::{AnchorOp, Fault, MemoryKeyProv use worldscript_secure_storage::{ advance_live_migration, capture_inventory, commit_catalog_change, commit_inventory_capture, commit_journal_checkpoint, commit_journal_takeover, commit_lease_renewal, commit_root, - content_digest, empty_inventory_digest, empty_journal_page_set_digest, generation_path, - inventory_page_dir, load_authoritative_manifest, load_catalog, operation_type, phase_code, + commit_streamed_inventory_capture, content_digest, empty_inventory_digest, + empty_journal_page_set_digest, generation_path, inventory_page_dir, + load_authoritative_manifest, load_catalog, load_staged_page, operation_type, phase_code, promote_manifest_fenced, seal_inventory_pages, source_authority_kind, source_physical_authority_kind, transition_phase, verify_stored_inventory, write_key_epoch, - AuthorityError, BindingAdvance, CandidateConflict, CatalogChange, CatalogCommit, + AuthorityError, BindingAdvance, CandidateConflict, CaptureStart, CatalogChange, CatalogCommit, DirectoryDurability, DurableFs, InstallationScopeId, InventoryCapture, JournalCheckpoint, JournalDurableContext, JournalDurableError, JournalError, JournalInventoryEntry, JournalInventorySource, JournalManifest, JournalPage, JournalSource, JournalTakeoverCommit, KeyEpochCommit, KeyEpochRecord, KeyEpochStatus, KeyProvider, LiveMigration, LoadedCatalog, MigrationExecutionError, MigrationFence, MigrationPhase, RecordClass, RecordIdentity, RootBody, RootCommitEvidence, RootCommitGuard, RootCommitRequest, RootCommitState, RootCommitted, - RootKeyRefV1, RootLayout, SealedPage, StageFailureKind, StdFs, WriteOperationId, + RootKeyRefV1, RootLayout, SealedPage, StageFailureKind, StagedCapture, StagedPage, StdFs, + StreamedCapture, StreamedInventoryCapture, WriteOperationId, }; const OPERATION: &str = "binding-op"; @@ -398,6 +400,51 @@ impl Fixture { ) } + /// A streamed capture of `staged` over `committed`, under the successor's own fence and the + /// committed key route and epoch. + fn capture_streamed( + &mut self, + staged: &StagedCapture, + committed: &JournalManifest, + ) -> Result { + let fence = MigrationFence::from_manifest(staged.successor()); + self.capture_streamed_with(&mut StdFs, (staged, committed), &fence) + } + + /// The same over `fs` and presenting `fence`, which is the successor's own token for an honest + /// caller. `capture` is the staged capture and the manifest the root binds. + fn capture_streamed_with( + &mut self, + fs: &mut F, + capture: (&StagedCapture, &JournalManifest), + fence: &MigrationFence, + ) -> Result { + let op = self + .operation + .clone() + .unwrap_or_else(|| WriteOperationId::generate().unwrap()); + let capture = StreamedInventoryCapture { + staged: capture.0, + committed_manifest: capture.1, + fence, + journal: JournalSource { + dir: &self.journal_dir, + operation: &op, + }, + root_key_ref: &self.key_ref, + active_key_epoch: 1, + conflict: self.conflict, + }; + commit_streamed_inventory_capture( + fs, + &mut self.provider, + RootLayout { + root_dir: &self.root_dir, + }, + capture, + ) + } + /// The relocated candidates in the journal directory, with their bytes. fn rejected_files(&self) -> Vec<(String, Vec)> { self.journal_files() @@ -2380,3 +2427,238 @@ fn a_different_renewal_never_replaces_the_candidate_and_the_original_is_still_ad fixture.renew(&renewal).unwrap(); assert_eq!(fixture.journal_files(), after_failure); } + +// ---- Streaming capture, stage two: the commit of a staged inventory ---- + +/// A finished stage-one capture of `count` entries, two to a page, over `committed` (the manifest the +/// root binds as `bound`), staged in `fixture`'s journal directory. +fn stage_over( + fixture: &Fixture, + committed: &JournalManifest, + bound: &LiveMigration, + count: u32, +) -> StagedCapture { + let (key, op) = (journal_key(), WriteOperationId::generate().unwrap()); + let mut fs = StdFs; + let mut ctx = JournalDurableContext::new(&mut fs, &key, &fixture.journal_dir, &op); + let start = CaptureStart { + committed_manifest: committed, + fence: &MigrationFence::from_manifest(committed), + live: bound, + entry_count: count, + }; + let mut capture = StreamedCapture::begin(&mut ctx, &start).unwrap(); + let all: Vec = (0..count).map(entry).collect(); + for chunk in all.chunks(2) { + capture = capture.push_page(&mut ctx, chunk.to_vec()).unwrap(); + } + capture.finish(&mut ctx).unwrap() +} + +/// Whether any file of `staged` is still on disk. +fn staged_files(staged: &StagedCapture) -> usize { + let mut files = Vec::new(); + collect_files(staged.pending_dir(), &mut files); + files.len() +} + +/// The staged pages read back and confirmed, in page order. +fn staged_pages(fixture: &Fixture, staged: &StagedCapture) -> Vec { + let (key, op) = (journal_key(), WriteOperationId::generate().unwrap()); + let mut fs = StdFs; + let mut ctx = JournalDurableContext::new(&mut fs, &key, &fixture.journal_dir, &op); + (0..staged.page_refs().len() as u32) + .map(|index| load_staged_page(&mut ctx, staged, index).unwrap()) + .collect() +} + +/// The journal tree without the staged files, which a retry is allowed to leave alone or remove. +fn durable_tree(fixture: &Fixture) -> Vec<(PathBuf, Vec)> { + let mut tree = fixture.journal_tree(); + tree.retain(|(path, _)| !path.to_string_lossy().contains("pending-")); + tree +} + +#[test] +fn a_streamed_capture_promotes_the_pages_publishes_the_manifest_and_advances_the_binding() { + let (mut fixture, one, bound) = bound_at_one(); + let staged = stage_over(&fixture, &one, &bound, 5); + let staged_before = staged_files(&staged); + // What the in-memory capture builds from the very same sealed pages. + let pages = staged_pages(&fixture, &staged); + let sealed: Vec<_> = pages + .iter() + .map(|p| SealedPage { + page: &p.page, + envelope: &p.envelope, + }) + .collect(); + let in_memory = capture_inventory(&one, &MigrationFence::from_manifest(&one), &sealed).unwrap(); + let before = fixture.loaded().root; + let committed = fixture.capture_streamed(&staged, &one).unwrap(); + let advanced = binding_of(staged.successor(), digest_of(&fixture.journal_dir, 2)); + let after = fixture.loaded().root; + assert_eq!( + (committed.root_generation, after.live_migration.clone()), + (before.root_generation + 1, Some(advanced.clone())) + ); + // The root names a manifest whose pages are all there and authenticate against it, and the staged + // files are gone once the root has committed to the set. + assert_eq!( + ( + verified_refs(&fixture, &advanced), + staged_before, + staged_files(&staged) + ), + (staged.page_refs().len(), 3, 0) + ); + // Exact parity with the in-memory capture: the committed manifest is the one it would build, and + // every page of the digest directory holds the bytes the staged page held. + let digest = &staged.successor().journal_page_set_digest; + let stored: Vec> = (0..pages.len() as u32) + .map(|index| { + let dir = inventory_page_dir(&fixture.journal_dir, digest, index); + fs::read(generation_path(&dir, 2)).unwrap() + }) + .collect(); + let staged_bytes: Vec> = pages.iter().map(|p| p.envelope.clone()).collect(); + assert_eq!( + ( + resumed_manifest(&fixture, &advanced) == in_memory, + stored == staged_bytes + ), + (true, true) + ); +} + +#[test] +fn a_streamed_capture_that_is_refused_changes_nothing_and_keeps_the_staged_files() { + let (mut fixture, one, bound) = bound_at_one(); + let staged = stage_over(&fixture, &one, &bound, 5); + let before = fixture.loaded().root; + let tree_before = fixture.journal_tree(); + // The right fencing generation at a journal revision the successor does not have. + let wrong_revision = MigrationFence { + fencing_generation: FENCE, + journal_revision: staged.successor().journal_revision + 3, + }; + let refused = fixture.capture_streamed_with(&mut StdFs, (&staged, &one), &wrong_revision); + let mut unbound = Fixture::new(); + unbound.ordinary_commit("bootstrap-root"); + let no_binding = unbound.capture_streamed(&staged, &one); + assert_eq!( + ( + refused, + no_binding, + fixture.loaded().root == before, + fixture.journal_tree() == tree_before + ), + ( + Err(AuthorityError::Journal(JournalDurableError::Fence( + MigrationExecutionError::StaleMigrationOwner + ))), + Err(AuthorityError::NoLiveMigration), + true, + true + ) + ); + assert_eq!(staged_files(&staged), 3); +} + +#[test] +fn a_staged_page_that_changed_refuses_the_commit_before_any_manifest_is_published() { + let (mut fixture, one, bound) = bound_at_one(); + let staged = stage_over(&fixture, &one, &bound, 5); + let before = fixture.loaded().root; + let file = generation_path(&staged.pending_dir().join("page-1"), 2); + let mut flipped = fs::read(&file).unwrap(); + flipped[20] ^= 1; + fs::write(&file, flipped).unwrap(); + let refused = fixture.capture_streamed(&staged, &one); + assert_eq!( + ( + refused, + fixture.loaded().root == before, + generation_path(&fixture.journal_dir, 2).exists(), + staged_files(&staged) + ), + ( + Err(AuthorityError::Journal(JournalDurableError::Journal( + JournalError::PageSetMismatch + ))), + true, + false, + 3 + ) + ); +} + +#[test] +fn a_failure_between_the_staged_pages_and_the_manifest_is_retried_and_completes() { + let (mut fixture, one, bound) = bound_at_one(); + let staged = stage_over(&fixture, &one, &bound, 5); + let before = fixture.loaded().root; + let mut failing = FailManifestFs { + inner: StdFs, + journal_dir: fixture.journal_dir.clone(), + }; + let fence = MigrationFence::from_manifest(staged.successor()); + let error = fixture + .capture_streamed_with(&mut failing, (&staged, &one), &fence) + .unwrap_err(); + // The pages are durable before the manifest; the manifest and the binding are not, and the staged + // files are still there for the retry. + let dir = fixture.journal_dir.clone(); + let digest = staged.successor().journal_page_set_digest; + assert_eq!( + ( + matches!(error, AuthorityError::Journal(_)), + inventory_page_dir(&dir, &digest, 0).is_dir(), + generation_path(&dir, 2).exists(), + fixture.loaded().root == before, + staged_files(&staged) + ), + (true, true, false, true, 3) + ); + fixture.capture_streamed(&staged, &one).unwrap(); + let advanced = binding_of(staged.successor(), digest_of(&dir, 2)); + assert_eq!( + (verified_refs(&fixture, &advanced), staged_files(&staged)), + (staged.page_refs().len(), 0) + ); +} + +#[test] +fn a_failed_root_commit_after_a_streamed_capture_is_retried_with_the_same_staged_capture() { + let (mut fixture, one, bound) = bound_at_one(); + let staged = stage_over(&fixture, &one, &bound, 5); + let before = fixture.loaded().root; + fixture + .provider + .inject(Fault::BeforePersist(AnchorOp::Prepare)); + let failed = fixture.capture_streamed(&staged, &one).unwrap_err(); + // The crash window: pages and manifest are durable, the root still names revision 1, and the + // staged files are still there for the retry. + let durable_after_failure = durable_tree(&fixture); + let kept = staged_files(&staged); + assert_eq!( + ( + matches!(failed, AuthorityError::Root(_)), + fixture.loaded().root == before, + kept + ), + (true, true, 3) + ); + let committed = fixture.capture_streamed(&staged, &one).unwrap(); + let advanced = binding_of(staged.successor(), digest_of(&fixture.journal_dir, 2)); + // The retry wrote nothing new into the journal: it adopted the pages and the manifest. + assert_eq!( + ( + committed.root_generation, + fixture.loaded().root.live_migration, + durable_tree(&fixture) == durable_after_failure, + staged_files(&staged) + ), + (before.root_generation + 1, Some(advanced), true, 0) + ); +} diff --git a/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs b/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs index 3312f92ee..28671cbc6 100644 --- a/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs +++ b/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs @@ -660,3 +660,302 @@ fn a_staged_page_is_trusted_only_after_it_is_confirmed_against_its_reference() { ] ); } + +// ---- Stage two: promoting a staged capture into the digest directory ---- + +/// Promotes `staged` as the committed owner of `as_committed` against the binding `as_live`, over `fs`, +/// in `scenario`'s journal directory. +fn promote_as( + fs: &mut F, + scenario: &Scenario, + staged: &StagedCapture, + as_committed: (&JournalManifest, &LiveMigration), +) -> Result { + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let fence = MigrationFence::from_manifest(as_committed.0); + let mut ctx = JournalDurableContext::new(fs, &key, scenario.journal.path(), &op); + let promotion = StagedPromotion { + committed_manifest: as_committed.0, + fence: &fence, + committed: as_committed.1, + staged, + }; + promote_staged_inventory_fenced(&mut ctx, &promotion) +} + +/// The page file `index` of the successor's digest directory. +fn promoted_page(scenario: &Scenario, staged: &StagedCapture, index: u32) -> std::path::PathBuf { + let digest = &staged.successor().journal_page_set_digest; + generation_path( + &inventory_page_dir(scenario.journal.path(), digest, index), + 4, + ) +} + +#[test] +fn a_promoted_capture_is_a_stored_set_that_verifies_against_its_successor() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let live = &scenario.journal.live; + let promoted = promote_as(&mut StdFs, &scenario, &staged, (&scenario.committed, live)); + // Publishing the successor the way the owner would, and binding it, makes the set verifiable. + let binding = commit_into(scenario.journal.path(), staged.successor()); + let (key, op) = (key(), WriteOperationId::generate().unwrap()); + let mut fs = StdFs; + let mut ctx = JournalDurableContext::new(&mut fs, &key, scenario.journal.path(), &op); + let verified = verify_stored_inventory(&mut ctx, &binding).unwrap(); + assert_eq!( + ( + promoted.is_ok(), + verified.page_refs() == staged.page_refs(), + verified.manifest() == staged.successor() + ), + (true, true, true) + ); +} + +#[test] +fn a_promotion_the_authority_checks_refuse_creates_nothing() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(5), 2).unwrap(); + let mut stale = scenario.journal.live.clone(); + stale.fencing_generation += 1; + let mut not_named = scenario.committed.clone(); + not_named.has_lease_owner = true; + not_named.lease_owner_id = Some("someone".into()); + not_named.lease_expires_unix_ms = Some(1); + let mut ahead = scenario.committed.clone(); + ahead.journal_revision += 1; + // The journal moved on after the capture began: its root-named manifest is now one revision ahead + // of the manifest the capture was built over. + let mut moved_on = scenario.committed.clone(); + moved_on.journal_revision += 1; + let moved_live = commit_into(scenario.journal.path(), &moved_on); + let authority = JournalDurableError::Authority; + let live = &scenario.journal.live; + let cases = [ + ( + "a stale owner", + promote_refusal(&scenario, &staged, (&scenario.committed, &stale)), + authority(MigrationExecutionError::StaleMigrationOwner), + ), + ( + "another revision", + promote_refusal(&scenario, &staged, (&ahead, live)), + authority(MigrationExecutionError::LiveBindingMismatch), + ), + ( + "a manifest the root does not name", + promote_refusal(&scenario, &staged, (¬_named, live)), + authority(MigrationExecutionError::LiveBindingMismatch), + ), + ( + "a journal that moved on after the capture began", + promote_refusal(&scenario, &staged, (&moved_on, &moved_live)), + authority(MigrationExecutionError::StaleJournalRevision), + ), + ]; + let wrong: Vec<_> = cases + .into_iter() + .filter_map(|(name, (refusal, created_nothing), expected)| { + (!(refusal == Some(expected) && created_nothing)).then_some(name) + }) + .collect(); + assert_eq!(wrong, Vec::<&str>::new()); +} + +/// The refusal of a promotion, and whether it created nothing. +fn promote_refusal( + scenario: &Scenario, + staged: &StagedCapture, + as_committed: (&JournalManifest, &LiveMigration), +) -> (Option, bool) { + let mut fs = ObservedFs::new(); + let refusal = promote_as(&mut fs, scenario, staged, as_committed).err(); + (refusal, fs.created_nothing()) +} + +#[test] +fn a_caller_that_is_not_the_committed_owner_is_refused_before_anything_is_read() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(5), 2).unwrap(); + let mut stale = scenario.journal.live.clone(); + stale.fencing_generation += 1; + let mut ahead = scenario.committed.clone(); + ahead.journal_revision += 1; + let reads_before_refusal = |as_committed: (&JournalManifest, &LiveMigration)| { + let mut fs = ObservedFs::new(); + let refusal = promote_as(&mut fs, &scenario, &staged, as_committed).err(); + (refusal.is_some(), fs.reads.len()) + }; + assert_eq!( + [ + reads_before_refusal((&scenario.committed, &stale)), + reads_before_refusal((&ahead, &scenario.journal.live)), + ], + [(true, 0), (true, 0)] + ); +} + +#[test] +fn a_staged_capture_is_promoted_only_into_the_journal_it_was_staged_in() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(5), 2).unwrap(); + // A directory that holds a byte-identical copy of the root-named generation: every check on the + // manifest passes there, 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 live = scenario.journal.live.clone(); + let elsewhere = Scenario { + journal: Journal { dir: copy, live }, + committed: scenario.committed.clone(), + }; + let mut fs = ObservedFs::new(); + let live = &elsewhere.journal.live; + let refused = promote_as(&mut fs, &elsewhere, &staged, (&elsewhere.committed, live)).err(); + assert_eq!( + ( + refused, + fs.created_nothing(), + fs.reads.len(), + files_under(elsewhere.journal.path()).len() + ), + ( + Some(JournalDurableError::Authority( + MigrationExecutionError::LiveBindingMismatch + )), + true, + 0, + 1 + ) + ); +} + +/// A staged page of generation 4, by index. +fn staged_file(staged: &StagedCapture, index: u32) -> std::path::PathBuf { + generation_path(&staged.pending_dir().join(format!("page-{index}")), 4) +} + +fn change_page_1(staged: &StagedCapture) { + let file = staged_file(staged, 1); + let mut bytes = std::fs::read(&file).unwrap(); + bytes[20] ^= 1; + std::fs::write(file, bytes).unwrap(); +} + +fn swap_page_0_for_page_1(staged: &StagedCapture) { + std::fs::write( + staged_file(staged, 0), + std::fs::read(staged_file(staged, 1)).unwrap(), + ) + .unwrap(); +} + +fn remove_page_2(staged: &StagedCapture) { + std::fs::remove_file(staged_file(staged, 2)).unwrap(); +} + +/// A way to make a staged page wrong, the refusal it must meet, and which pages may have been stored. +type WrongPage = ( + &'static str, + fn(&StagedCapture), + JournalDurableError, + [bool; 3], +); + +#[test] +fn a_staged_page_that_is_wrong_stops_the_promotion_and_leaves_only_a_verified_prefix() { + let mismatch = JournalDurableError::Journal(JournalError::PageSetMismatch); + let missing = JournalDurableError::Journal(JournalError::Corrupt("staged page is missing")); + let cases: [WrongPage; 3] = [ + ( + "a changed page", + change_page_1, + mismatch.clone(), + [true, false, false], + ), + ( + "a swapped page", + swap_page_0_for_page_1, + mismatch, + [false, false, false], + ), + ( + "a missing page", + remove_page_2, + missing, + [true, true, false], + ), + ]; + let wrong: Vec<_> = cases + .into_iter() + .filter_map(|(name, tamper, expected, stored)| { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + tamper(&staged); + let live = &scenario.journal.live; + let refused = + promote_as(&mut StdFs, &scenario, &staged, (&scenario.committed, live)).err(); + // Only the pages before the wrong one were stored, and no manifest follows them. + let prefix = [0, 1, 2].map(|index| promoted_page(&scenario, &staged, index).is_file()); + let published = generation_path(scenario.journal.path(), 4).exists(); + (!(refused == Some(expected) && prefix == stored && !published)).then_some(name) + }) + .collect(); + assert_eq!(wrong, Vec::<&str>::new()); +} + +/// Promotes `staged` as the scenario's committed owner, but under `key`. +fn promote_under( + fs: &mut F, + scenario: &Scenario, + staged: &StagedCapture, + key: &Key, +) -> Result { + let op = WriteOperationId::generate().unwrap(); + let fence = scenario.fence(); + let mut ctx = JournalDurableContext::new(fs, key, scenario.journal.path(), &op); + let promotion = StagedPromotion { + committed_manifest: &scenario.committed, + fence: &fence, + committed: &scenario.journal.live, + staged, + }; + promote_staged_inventory_fenced(&mut ctx, &promotion) +} + +#[test] +fn a_promotion_under_another_key_is_refused_before_anything_is_created() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(5), 2).unwrap(); + let mut fs = ObservedFs::new(); + let refused = promote_under(&mut fs, &scenario, &staged, &other_key()).err(); + assert_eq!( + (refused, fs.created_nothing()), + ( + Some(JournalDurableError::Journal(JournalError::Open( + OpenError::Tampered + ))), + true + ) + ); +} + +#[test] +fn a_repeated_promotion_adopts_what_is_already_stored() { + let scenario = Scenario::new(); + let staged = scenario.stream(&mut StdFs, entries(7), 3).unwrap(); + let live = &scenario.journal.live; + let first = promote_as(&mut StdFs, &scenario, &staged, (&scenario.committed, live)); + let after_first = files_under(scenario.journal.path()); + let second = promote_as(&mut StdFs, &scenario, &staged, (&scenario.committed, live)); + assert_eq!( + ( + first.is_ok(), + second.is_ok(), + files_under(scenario.journal.path()) + ), + (true, true, after_first) + ); +} diff --git a/docs/native/R15-SECURE-STORAGE-CONTRACT.md b/docs/native/R15-SECURE-STORAGE-CONTRACT.md index 96445f6b4..f64480532 100644 --- a/docs/native/R15-SECURE-STORAGE-CONTRACT.md +++ b/docs/native/R15-SECURE-STORAGE-CONTRACT.md @@ -2433,8 +2433,15 @@ it is read back by the exact path of its reference and confirmed against that re 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. +holds the bytes that were staged, and an attempt that was killed leaves files that nothing trusts. The staged pages are then promoted to the +digest directory under the root lock and the journal mutex: the caller must be the committed owner of the +exact root-named manifest and the successor must be its capture successor, all before the first write; +each staged page is read back, confirmed against its reference and stored where the in-memory store would +have put the same bytes (an identical page already there is adopted, a different file is never replaced); +and the inventory digest over the entries read must equal the successor's before the manifest naming the +set is published. A staged page found wrong midway leaves the verified prefix as an inert directory under +the digest and publishes nothing. The staged files are removed only after the root has committed to the +set, so that a retry after a failure before that point still has them. **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 e5ee8f563..f5cf97ed1 100644 --- a/docs/native/r15/GATE4D-JOURNAL-DURABLE-EVIDENCE.md +++ b/docs/native/r15/GATE4D-JOURNAL-DURABLE-EVIDENCE.md @@ -665,6 +665,44 @@ 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 two (a) — promoting a staged capture and committing it + +Stage one left a finished inventory in `inventory/pending-*`; this is the step that makes it reachable. + +| Step | Rule | +|---|---| +| `promote_staged_inventory_fenced` | before any read or write: the context is for the journal directory the capture was staged in (a finished capture remembers it and refuses any other, because its pages are read by absolute path), and under the journal mutex the caller is the committed owner writing under the committed manifest (`assert_page_promote_authority`, which refuses before anything is read); then the exact root-named predecessor, the capture successor relation, and the staged references against the page set the successor names; then for each staged page `load_staged_page` (digest first, then open), the entries absorbed into the inventory digest, and `store_page_at` into `inventory//page-/` (an identical page is adopted, a different file never replaced); finally the inventory digest must equal the successor's | +| `commit_streamed_inventory_capture` | `commit_inventory_capture`'s order and guarantees: operation id, fence, root lock, committed binding, journal key routed from the root, the promotion above, `publish_manifest_fenced` as the capture, `BindingStep::Capture`; the staged files are removed (best effort) only after the root has committed, so a retry after a failed root commit still has them and adopts the pages and the manifest candidate; after a success the handle is spent | + +`commit_inventory_capture` and the new commit now share one core (`commit_capture`), which takes the step that +stores the pages as an argument, so the order of operations under the root lock exists once. + +Decisions, disclosed: (a) one pass: pages are verified and stored as they are loaded, so a staged page found +wrong midway leaves an inert prefix under the digest directory (as an I/O failure partway through the in-memory +store already does) and no manifest is published; (b) the commit takes `&StagedCapture` so that a retry has the +files; (c) it is a second entry point beside `commit_inventory_capture`, not a replacement. Two checks cannot +be reached by any caller, because `StreamedCapture::finish` cannot produce a successor that disagrees with +its staged references: the page-set consistency check and the final inventory-digest check are defence in +depth and have no test that fails without them through the public API. + +Proof: the promotion on the journal alone (a promoted capture publishes and binds as a stored set that +`verify_stored_inventory` accepts, with the same references and manifest; a repeated promotion adopts what is +stored; a stale owner and another revision are refused before anything is read; a manifest the root does not +name, a journal that moved on after the capture began, a journal directory other than the one the capture +was staged in (one that holds an identical copy of the root-named manifest) and another key are refused with +nothing created; a staged page that changed, was swapped for another or is missing stops the promotion with +only a verified prefix stored and no manifest published, the missing one as `Corrupt`) and the composed commit with a real root (the root names the successor, which is exactly the manifest the +in-memory capture builds from the same sealed pages, its pages verify and hold the bytes the staged pages +held, and the staged files are gone; a stale token and an unbound migration are refused with the journal tree +and the root unchanged and the staged files kept; a changed staged page refuses the commit before any manifest +is published; a failure between the promoted pages and the manifest, and a failed root commit after it, are each retried +with the same staged capture, which writes nothing new into the journal that is already durable, and then +removes the staged files); the key route refuses the new commit with a revoked or unregistered +epoch and with a route to another key, beside the other journal-owner operations. Mutation-checked: the +root-named, capture-successor, promote-authority and journal-binding checks, the digest directory, the inventory digest +absorption, keeping the staged files after a success, removing them before the root commit, and a fixed +journal key in the shared core each fail the test that owns them. + ## 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 @@ -706,7 +744,7 @@ 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` +Proof (`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 @@ -926,7 +964,7 @@ 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, 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). +- Streaming capture, stage two (b): the authority-level staging session that routes the journal key from the root, so the key is no longer a caller input to staging; until then a caller of `StreamedCapture` supplies the key (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 af4bbfb25..f13db1196 100644 --- a/docs/native/r15/GATE4D-SLICE-B-GAP-MATRIX.md +++ b/docs/native/r15/GATE4D-SLICE-B-GAP-MATRIX.md @@ -89,7 +89,8 @@ 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. +- 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 and the composed commit are stage two (a), the key-routing session stage two (b). +- Streaming capture, stage two (a): `promote_staged_inventory_fenced` stores the staged pages under the digest directory of the successor one page at a time (owner, root-named predecessor and capture successor refused before any write; each page digest-checked, absorbed into the inventory digest, stored with adoption of an identical page; the inventory digest confirmed) and `commit_streamed_inventory_capture` publishes the successor and advances the binding in `commit_inventory_capture`'s order through a shared core, removing the staged files only after the root has committed so a retry still has them. The key-routing staging session is stage two (b). - 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)