diff --git a/CHANGELOG.md b/CHANGELOG.md index 271e6e189..f557675b2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +- **R-15 Gate 4D:** stage a streamed inventory without handling a journal key (stage two, part b, of the + streaming capture). `begin_streamed_capture` routes the journal key from the authenticated root and + registry, reads the binding and the committed manifest from the root and returns a `StagingSession` that + stages page by page under that key and one journal directory value; `commit_streamed_inventory_capture` no + longer takes a committed manifest or a journal directory, promoting into the directory the capture was + staged in and loading the manifest from the root under the lock. PR #1012. - **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 diff --git a/crates/worldscript-secure-storage/src/authority.rs b/crates/worldscript-secure-storage/src/authority.rs index 24abb3ee2..66b6bb99a 100644 --- a/crates/worldscript-secure-storage/src/authority.rs +++ b/crates/worldscript-secure-storage/src/authority.rs @@ -45,9 +45,9 @@ use crate::journal::{ assert_progress_successor, assert_renewal_successor, assert_takeover_successor, 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, + CaptureStart, InventorySetWrite, JournalDurableContext, JournalDurableError, + JournalInventoryEntry, JournalManifest, JournalTakeover, MigrationExecutionError, + MigrationFence, SealedPage, StagedCapture, StagedPromotion, StreamedCapture, }; use crate::journal_route::{resolve_journal_key, JournalRoute, JournalRouteError}; use crate::marker::content_digest; @@ -446,16 +446,16 @@ pub fn commit_inventory_capture( } /// 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. +/// capture and the owner's token and routes. The journal directory is the one the capture was staged +/// in, and the committed manifest is read from the root, so neither is a caller input. #[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>, + /// The write operation of this commit. + pub operation: &'a WriteOperationId, 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. @@ -470,7 +470,9 @@ pub struct StreamedInventoryCapture<'a> { /// 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. +/// is spent. The journal key is routed from the root here, as for every journal-owner commit, the +/// committed manifest is the one the root names (loaded under that key while the lock is held), and +/// the journal directory is the one the capture was staged in. pub fn commit_streamed_inventory_capture( fs: &mut F, provider: &mut P, @@ -481,21 +483,26 @@ pub fn commit_streamed_inventory_capture( let checkpoint = JournalCheckpoint { manifest: staged.successor(), fence: capture.fence, - journal: capture.journal, + journal: JournalSource { + dir: staged.journal_dir(), + operation: capture.operation, + }, root_key_ref: capture.root_key_ref, active_key_epoch: capture.active_key_epoch, conflict: capture.conflict, }; + // The successor is the revision after the one the pages are written under. let plan = CapturePlan { layout, checkpoint, - committed_revision: capture.committed_manifest.journal_revision, + committed_revision: staged.successor().journal_revision.saturating_sub(1), }; let committed = commit_capture(fs, provider, plan, |journal, committed, fence| { + let manifest = load_authoritative_manifest(journal, committed)?; promote_staged_inventory_fenced( journal, &StagedPromotion { - committed_manifest: capture.committed_manifest, + committed_manifest: &manifest, fence, committed, staged, @@ -506,6 +513,111 @@ pub fn commit_streamed_inventory_capture( Ok(committed) } +/// What a staging session needs from its caller: where the journal lives and the write operation that +/// stages it, the owner's token at the committed revision, and the total number of entries the +/// inventory will have (the inventory digest commits to it first). +#[derive(Clone, Copy)] +pub struct SessionBegin<'a> { + pub journal: JournalSource<'a>, + pub fence: &'a MigrationFence, + pub entry_count: u32, +} + +/// A streamed capture in progress together with the journal key and directory it stages under, so a +/// caller of the staging stage never handles a key and never spells the journal directory twice. +pub struct StagingSession { + capture: StreamedCapture, + key: Key, + journal_dir: PathBuf, + operation: WriteOperationId, +} + +// The key is deliberately left out of the debug output. +impl std::fmt::Debug for StagingSession { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("StagingSession") + .field("capture", &self.capture) + .field("journal_dir", &self.journal_dir) + .finish_non_exhaustive() + } +} + +/// Starts staging an inventory one page at a time, with the journal key routed from the root. +/// +/// The binding is read from the committed root, the journal key is resolved through the key-epoch +/// registry ([`resolve_journal_key`]), the committed manifest is loaded by the exact root-named path +/// under that key, and the stage-one capture begins over it; the caller supplies no key, manifest or +/// binding. Nothing is created and no lock is taken: these are early refusals, and the commit repeats +/// every authority check under the root lock. +pub fn begin_streamed_capture( + fs: &mut F, + provider: &P, + layout: RootLayout<'_>, + begin: SessionBegin<'_>, +) -> Result { + // The root alone names the binding: the catalog pages are not read, so starting a capture of a + // very large inventory does not first materialise a very large catalog. + let live = load_committed_root(fs, provider, layout) + .map_err(AuthorityError::Root)? + .and_then(|view| view.root.live_migration) + .ok_or(AuthorityError::NoLiveMigration)?; + let key = route_journal_key(fs, provider, layout, begin.journal)?; + let capture = { + let mut ctx = journal_context(&mut *fs, begin.journal, &key, CandidateConflict::Refuse); + let committed = + load_authoritative_manifest(&mut ctx, &live).map_err(AuthorityError::Journal)?; + let start = CaptureStart { + committed_manifest: &committed, + fence: begin.fence, + live: &live, + entry_count: begin.entry_count, + }; + StreamedCapture::begin(&mut ctx, &start).map_err(AuthorityError::Journal)? + }; + Ok(StagingSession { + capture, + key, + journal_dir: begin.journal.dir.to_path_buf(), + operation: begin.journal.operation.clone(), + }) +} + +impl StagingSession { + /// Stages the next page under the session's key and directory; see [`StreamedCapture::push_page`]. + pub fn push_page( + self, + fs: &mut F, + entries: Vec, + ) -> Result { + let Self { + capture, + key, + journal_dir, + operation, + } = self; + let capture = { + let mut ctx = JournalDurableContext::new(fs, &key, &journal_dir, &operation); + capture + .push_page(&mut ctx, entries) + .map_err(AuthorityError::Journal)? + }; + Ok(Self { + capture, + key, + journal_dir, + operation, + }) + } + + /// Ends the capture; see [`StreamedCapture::finish`]. + pub fn finish(self, fs: &mut F) -> Result { + let mut ctx = JournalDurableContext::new(fs, &self.key, &self.journal_dir, &self.operation); + self.capture + .finish(&mut ctx) + .map_err(AuthorityError::Journal) + } +} + /// 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. diff --git a/crates/worldscript-secure-storage/src/journal/stream_capture.rs b/crates/worldscript-secure-storage/src/journal/stream_capture.rs index 55527031d..f6d98ae20 100644 --- a/crates/worldscript-secure-storage/src/journal/stream_capture.rs +++ b/crates/worldscript-secure-storage/src/journal/stream_capture.rs @@ -136,9 +136,11 @@ impl StagedCapture { } /// 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>) { - self.remove_files(ctx.fs); + /// stay, because the file system abstraction cannot remove a directory. Needs no key: a caller + /// that staged through a session never held one, and removing a file the capture recorded by + /// its digest authenticates nothing. + pub fn discard(self, fs: &mut F) { + self.remove_files(fs); } /// The same removal for a caller that keeps the handle, such as the commit that has just spent it. diff --git a/crates/worldscript-secure-storage/src/lib.rs b/crates/worldscript-secure-storage/src/lib.rs index a17d18f09..88961ea43 100644 --- a/crates/worldscript-secure-storage/src/lib.rs +++ b/crates/worldscript-secure-storage/src/lib.rs @@ -60,12 +60,12 @@ pub use admission::{ SharedAdmissionGuard, OPERATION_ADMISSION_LOCK_FILE, }; pub use authority::{ - advance_live_migration, commit_catalog_change, commit_inventory_capture, - 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, + advance_live_migration, begin_streamed_capture, commit_catalog_change, + commit_inventory_capture, 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, + SessionBegin, StagingSession, StreamedInventoryCapture, }; #[cfg(feature = "test-support")] pub use authority::{list_records, load_catalog}; 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 2a8c2418d..267baaae9 100644 --- a/crates/worldscript-secure-storage/tests/gate4d_journal_route_test.rs +++ b/crates/worldscript-secure-storage/tests/gate4d_journal_route_test.rs @@ -14,19 +14,20 @@ use std::sync::atomic::{AtomicU32, Ordering}; 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, - 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, StagedCapture, StdFs, StreamedCapture, StreamedInventoryCapture, - WriteOperationId, JOURNAL_MANIFEST_RECORD_SCHEMA, + advance_live_migration, begin_streamed_capture, commit_catalog_change, + commit_inventory_capture, commit_journal_checkpoint, commit_journal_takeover, + commit_lease_renewal, commit_root, 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, SessionBegin, StagedCapture, + StdFs, StreamedCapture, StreamedInventoryCapture, WriteOperationId, + JOURNAL_MANIFEST_RECORD_SCHEMA, }; const OPERATION: &str = "route-op"; @@ -560,18 +561,13 @@ impl Fixture { /// 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 root_dir = self.root_dir(); 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, - }, + operation: &operation, root_key_ref: &root_ref, active_key_epoch: self.root_epoch, conflict: CandidateConflict::Refuse, @@ -583,6 +579,26 @@ impl Fixture { .unwrap_err() } + /// The start of a staging session over valid inputs, which is expected to be refused: the key is + /// routed from the root before anything is loaded or created. + fn session_is_refused(&self) -> AuthorityError { + let (root_dir, journal_dir) = (self.root_dir(), self.journal_dir()); + let fence = MigrationFence::from_manifest(&rotation()); + let operation = WriteOperationId::generate().unwrap(); + let begin = SessionBegin { + journal: JournalSource { + dir: &journal_dir, + operation: &operation, + }, + fence: &fence, + entry_count: 0, + }; + let layout = RootLayout { + root_dir: &root_dir, + }; + begin_streamed_capture(&mut StdFs, &self.provider, layout, begin).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 @@ -676,6 +692,7 @@ fn a_revoked_journal_epoch_refuses_every_journal_operation_before_a_write() { let before = fixture.snapshot(); let mut errors = fixture.every_operation_is_refused(&stored); errors.push(fixture.streamed_capture_is_refused(&staged)); + errors.push(fixture.session_is_refused()); for error in errors { assert_eq!( error, @@ -695,6 +712,7 @@ fn an_unregistered_journal_epoch_refuses_every_journal_operation_before_a_write( let before = fixture.snapshot(); let mut errors = fixture.every_operation_is_refused(&stored); errors.push(fixture.streamed_capture_is_refused(&staged)); + errors.push(fixture.session_is_refused()); for error in errors { assert_eq!( error, @@ -718,6 +736,7 @@ fn a_route_to_another_key_is_refused_by_the_authenticated_load_before_a_write() // Every operation has valid inputs, so the key is the only thing left to refuse them. let mut errors = fixture.every_operation_is_refused(&stored); errors.push(fixture.streamed_capture_is_refused(&staged)); + errors.push(fixture.session_is_refused()); for error in errors { assert_eq!( error, 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 371b892a4..0eeb98efa 100644 --- a/crates/worldscript-secure-storage/tests/gate4d_root_binding_test.rs +++ b/crates/worldscript-secure-storage/tests/gate4d_root_binding_test.rs @@ -11,22 +11,23 @@ use std::sync::atomic::{AtomicU32, Ordering}; use worldscript_secure_storage::memory_provider::{AnchorOp, Fault, MemoryKeyProvider}; 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, - commit_streamed_inventory_capture, content_digest, empty_inventory_digest, - empty_journal_page_set_digest, generation_path, inventory_page_dir, + advance_live_migration, begin_streamed_capture, capture_inventory, commit_catalog_change, + commit_inventory_capture, commit_journal_checkpoint, commit_journal_takeover, + commit_lease_renewal, commit_root, 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, CaptureStart, CatalogChange, CatalogCommit, - DirectoryDurability, DurableFs, InstallationScopeId, InventoryCapture, JournalCheckpoint, - JournalDurableContext, JournalDurableError, JournalError, JournalInventoryEntry, - JournalInventorySource, JournalManifest, JournalPage, JournalSource, JournalTakeoverCommit, - KeyEpochCommit, KeyEpochRecord, KeyEpochStatus, KeyProvider, LiveMigration, LoadedCatalog, + AuthorityError, BindingAdvance, CandidateConflict, CatalogChange, CatalogCommit, + CatalogDescriptor, CommitMarker, CommittedGeneration, DirectoryDurability, DurableFs, + InstallationScopeId, InventoryCapture, JournalCheckpoint, JournalDurableContext, + JournalDurableError, JournalError, JournalInventoryEntry, JournalInventorySource, + JournalManifest, JournalPage, JournalSource, JournalTakeoverCommit, KeyEpochCommit, + KeyEpochRecord, KeyEpochStatus, KeyProvider, LiveMigration, LoadedCatalog, MarkerBody, MigrationExecutionError, MigrationFence, MigrationPhase, RecordClass, RecordIdentity, RootBody, RootCommitEvidence, RootCommitGuard, RootCommitRequest, RootCommitState, RootCommitted, - RootKeyRefV1, RootLayout, SealedPage, StageFailureKind, StagedCapture, StagedPage, StdFs, - StreamedCapture, StreamedInventoryCapture, WriteOperationId, + RootKeyRefV1, RootLayout, SealedPage, SessionBegin, StageFailureKind, StagedCapture, + StagedPage, StagingSession, StdFs, StreamedInventoryCapture, WriteOperationId, }; const OPERATION: &str = "binding-op"; @@ -194,6 +195,38 @@ impl Fixture { .root_generation } + /// An ordinary catalog commit that adds one descriptor, so that the catalog has a page. + fn commit_descriptor(&mut self, operation_id: &str) { + let record = RecordIdentity::new(RecordClass::Codex, &["a-record"]).unwrap(); + let body = MarkerBody::Active { + committed_generation: 1, + committed_epoch: 1, + content_digest: [0xab; 32], + }; + let marker = CommitMarker::new(&record, 1, body).unwrap(); + let readable = CommittedGeneration { + generation: 1, + epoch: 1, + content_digest: [0xab; 32], + }; + let descriptor = + CatalogDescriptor::new_unverified(&record, &marker, Some(readable)).unwrap(); + let root_dir = self.root_dir.clone(); + let commit = CatalogCommit { + change: CatalogChange { + upsert: &[descriptor], + remove: &[], + }, + root_key_ref: &self.key_ref, + active_key_epoch: 1, + operation_id, + }; + let layout = RootLayout { + root_dir: &root_dir, + }; + commit_catalog_change(&mut StdFs, &mut self.provider, layout, commit).unwrap(); + } + /// Commits the first ordinary root, then a root that binds `binding` (no producer exists in /// the crate yet, so the test writes the root directly). fn bind(&mut self, binding: &LiveMigration) { @@ -400,23 +433,22 @@ impl Fixture { ) } - /// A streamed capture of `staged` over `committed`, under the successor's own fence and the - /// committed key route and epoch. + /// A streamed capture of `staged`, under the successor's own fence and the committed key route and + /// epoch. The journal directory and the committed manifest are not inputs. 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) + self.capture_streamed_with(&mut StdFs, staged, &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. + /// caller. fn capture_streamed_with( &mut self, fs: &mut F, - capture: (&StagedCapture, &JournalManifest), + staged: &StagedCapture, fence: &MigrationFence, ) -> Result { let op = self @@ -424,13 +456,9 @@ impl Fixture { .clone() .unwrap_or_else(|| WriteOperationId::generate().unwrap()); let capture = StreamedInventoryCapture { - staged: capture.0, - committed_manifest: capture.1, + staged, fence, - journal: JournalSource { - dir: &self.journal_dir, - operation: &op, - }, + operation: &op, root_key_ref: &self.key_ref, active_key_epoch: 1, conflict: self.conflict, @@ -2430,29 +2458,38 @@ fn a_different_renewal_never_replaces_the_candidate_and_the_original_is_still_ad // ---- 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( +/// The start of a staging session in `fixture`'s journal directory under `fence`. The caller holds no +/// journal key, no committed manifest and no binding. +fn begin_session( fixture: &Fixture, - committed: &JournalManifest, - bound: &LiveMigration, + fence: &MigrationFence, 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, +) -> Result { + let op = WriteOperationId::generate().unwrap(); + let begin = SessionBegin { + journal: JournalSource { + dir: &fixture.journal_dir, + operation: &op, + }, + fence, entry_count: count, }; - let mut capture = StreamedCapture::begin(&mut ctx, &start).unwrap(); + let layout = RootLayout { + root_dir: &fixture.root_dir, + }; + begin_streamed_capture(&mut StdFs, &fixture.provider, layout, begin) +} + +/// A finished capture of `count` entries, two to a page, staged through a session in `fixture`'s +/// journal directory over `committed`, the manifest the root binds. +fn stage_over(fixture: &mut Fixture, committed: &JournalManifest, count: u32) -> StagedCapture { + let fence = MigrationFence::from_manifest(committed); + let mut session = begin_session(fixture, &fence, count).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(); + session = session.push_page(&mut StdFs, chunk.to_vec()).unwrap(); } - capture.finish(&mut ctx).unwrap() + session.finish(&mut StdFs).unwrap() } /// Whether any file of `staged` is still on disk. @@ -2481,8 +2518,8 @@ fn durable_tree(fixture: &Fixture) -> Vec<(PathBuf, Vec)> { #[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 (mut fixture, one, _) = bound_at_one(); + let staged = stage_over(&mut fixture, &one, 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); @@ -2495,7 +2532,7 @@ fn a_streamed_capture_promotes_the_pages_publishes_the_manifest_and_advances_the .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 committed = fixture.capture_streamed(&staged).unwrap(); let advanced = binding_of(staged.successor(), digest_of(&fixture.journal_dir, 2)); let after = fixture.loaded().root; assert_eq!( @@ -2533,8 +2570,8 @@ fn a_streamed_capture_promotes_the_pages_publishes_the_manifest_and_advances_the #[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 (mut fixture, one, _) = bound_at_one(); + let staged = stage_over(&mut fixture, &one, 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. @@ -2542,10 +2579,10 @@ fn a_streamed_capture_that_is_refused_changes_nothing_and_keeps_the_staged_files fencing_generation: FENCE, journal_revision: staged.successor().journal_revision + 3, }; - let refused = fixture.capture_streamed_with(&mut StdFs, (&staged, &one), &wrong_revision); + let refused = fixture.capture_streamed_with(&mut StdFs, &staged, &wrong_revision); let mut unbound = Fixture::new(); unbound.ordinary_commit("bootstrap-root"); - let no_binding = unbound.capture_streamed(&staged, &one); + let no_binding = unbound.capture_streamed(&staged); assert_eq!( ( refused, @@ -2567,14 +2604,14 @@ fn a_streamed_capture_that_is_refused_changes_nothing_and_keeps_the_staged_files #[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 (mut fixture, one, _) = bound_at_one(); + let staged = stage_over(&mut fixture, &one, 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); + let refused = fixture.capture_streamed(&staged); assert_eq!( ( refused, @@ -2595,8 +2632,8 @@ fn a_staged_page_that_changed_refuses_the_commit_before_any_manifest_is_publishe #[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 (mut fixture, one, _) = bound_at_one(); + let staged = stage_over(&mut fixture, &one, 5); let before = fixture.loaded().root; let mut failing = FailManifestFs { inner: StdFs, @@ -2604,7 +2641,7 @@ fn a_failure_between_the_staged_pages_and_the_manifest_is_retried_and_completes( }; let fence = MigrationFence::from_manifest(staged.successor()); let error = fixture - .capture_streamed_with(&mut failing, (&staged, &one), &fence) + .capture_streamed_with(&mut failing, &staged, &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. @@ -2620,7 +2657,7 @@ fn a_failure_between_the_staged_pages_and_the_manifest_is_retried_and_completes( ), (true, true, false, true, 3) ); - fixture.capture_streamed(&staged, &one).unwrap(); + fixture.capture_streamed(&staged).unwrap(); let advanced = binding_of(staged.successor(), digest_of(&dir, 2)); assert_eq!( (verified_refs(&fixture, &advanced), staged_files(&staged)), @@ -2630,13 +2667,13 @@ fn a_failure_between_the_staged_pages_and_the_manifest_is_retried_and_completes( #[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 (mut fixture, one, _) = bound_at_one(); + let staged = stage_over(&mut fixture, &one, 5); let before = fixture.loaded().root; fixture .provider .inject(Fault::BeforePersist(AnchorOp::Prepare)); - let failed = fixture.capture_streamed(&staged, &one).unwrap_err(); + let failed = fixture.capture_streamed(&staged).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); @@ -2649,7 +2686,7 @@ fn a_failed_root_commit_after_a_streamed_capture_is_retried_with_the_same_staged ), (true, true, 3) ); - let committed = fixture.capture_streamed(&staged, &one).unwrap(); + let committed = fixture.capture_streamed(&staged).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!( @@ -2662,3 +2699,97 @@ fn a_failed_root_commit_after_a_streamed_capture_is_retried_with_the_same_staged (before.root_generation + 1, Some(advanced), true, 0) ); } + +#[test] +fn a_staging_session_is_refused_before_anything_is_created() { + let (fixture, one, _) = bound_at_one(); + let tree_before = fixture.journal_tree(); + let stale = MigrationFence { + fencing_generation: FENCE - 1, + journal_revision: one.journal_revision, + }; + let mut unbound = Fixture::new(); + unbound.ordinary_commit("bootstrap-root"); + let outcomes = [ + begin_session(&fixture, &stale, 3).err(), + begin_session(&unbound, &MigrationFence::from_manifest(&one), 3).err(), + ]; + assert_eq!( + ( + outcomes, + fixture.journal_tree() == tree_before, + unbound.journal_tree().is_empty() + ), + ( + [ + Some(AuthorityError::Journal(JournalDurableError::Fence( + MigrationExecutionError::StaleMigrationOwner + ))), + Some(AuthorityError::NoLiveMigration) + ], + true, + true + ) + ); +} + +#[test] +fn a_commit_after_the_journal_moved_on_is_refused_before_any_write_and_keeps_the_staged_files() { + let (mut fixture, one, _) = bound_at_one(); + let staged = stage_over(&mut fixture, &one, 5); + // The owner moves the journal on while the capture waits: the manifest the root names is no + // longer the one the pages were staged over. + fixture + .checkpoint(&manifest_at(OPERATION, 2, FENCE)) + .unwrap(); + let before = fixture.loaded().root; + let tree_before = fixture.journal_tree(); + let refused = fixture.capture_streamed(&staged); + assert_eq!( + ( + refused, + fixture.loaded().root == before, + fixture.journal_tree() == tree_before, + staged_files(&staged) + ), + ( + Err(AuthorityError::Journal(JournalDurableError::Fence( + MigrationExecutionError::StaleMigrationOwner + ))), + true, + true, + 3 + ) + ); +} + +#[test] +fn a_session_starts_from_the_root_alone_and_does_not_read_the_catalog_pages() { + let (mut fixture, one, _) = bound_at_one(); + fixture.commit_descriptor("add-a-record"); + // A catalog page that no longer verifies: the catalog cannot be loaded, the committed root can. + let mut pages = Vec::new(); + collect_files(&fixture.root_dir.join("catalog"), &mut pages); + let (page, mut bytes) = pages.into_iter().next().unwrap(); + bytes[20] ^= 1; + fs::write(&page, bytes).unwrap(); + let layout = RootLayout { + root_dir: &fixture.root_dir, + }; + let catalog_loads = load_catalog(&mut StdFs, &fixture.provider, layout).is_ok(); + let staged = stage_over(&mut fixture, &one, 3); + assert_eq!((catalog_loads, staged.page_refs().len()), (false, 2)); +} + +#[test] +fn a_finished_session_capture_is_discarded_without_a_key() { + let (mut fixture, one, _) = bound_at_one(); + let staged = stage_over(&mut fixture, &one, 5); + let before = staged_files(&staged); + let pending = staged.pending_dir().to_path_buf(); + // The caller of a session never held the journal key, and discarding needs only the file system. + staged.discard(&mut StdFs); + let mut left = Vec::new(); + collect_files(&pending, &mut left); + assert_eq!((before, left.len()), (3, 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 28671cbc6..a2a87e9ec 100644 --- a/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs +++ b/crates/worldscript-secure-storage/tests/gate4d_stream_capture_test.rs @@ -94,10 +94,7 @@ impl Scenario { /// 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.discard(&mut ObservedFs::new()); } /// Staged page `index`, read back under `key`. @@ -469,11 +466,9 @@ fn an_abandoned_capture_removes_its_staged_pages() { 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); + staged.discard(&mut fs); assert_eq!(files_under(scenario.journal.path()).len(), 4); } diff --git a/docs/native/R15-SECURE-STORAGE-CONTRACT.md b/docs/native/R15-SECURE-STORAGE-CONTRACT.md index f64480532..aa8388070 100644 --- a/docs/native/R15-SECURE-STORAGE-CONTRACT.md +++ b/docs/native/R15-SECURE-STORAGE-CONTRACT.md @@ -2432,16 +2432,23 @@ authority: the manifest never names them, no reader looks for them, a staged pag 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. 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 +staged under another context. While a capture is being staged, 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. Once a capture is finished, its staged files are kept +through the promotion and the commit unless the caller discards the capture. 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. +the digest and publishes nothing. A failure of the promotion or of the commit does not remove the +staged files of a finished capture: they are removed only after the root has committed to the set, so that +a retry after such a failure still has them. A caller of this path handles no +journal key, no committed manifest and no binding: the key is routed from the authenticated root and +registry, both the start of the capture and the commit read the committed manifest and the binding from +the root, and the commit promotes into the journal directory the capture was staged in, so there is one +directory value per capture. **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 f5cf97ed1..e0fe37dd5 100644 --- a/docs/native/r15/GATE4D-JOURNAL-DURABLE-EVIDENCE.md +++ b/docs/native/r15/GATE4D-JOURNAL-DURABLE-EVIDENCE.md @@ -665,6 +665,32 @@ 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 (b) — the staging session routes the key + +After stage two (a) the stage-one builder still took the journal key, the committed manifest and the binding +from its caller, and the commit took the committed manifest and a journal directory, compared with the staged +one by spelling (a deferred finding on #1011). The authority-level session removes those inputs. + +| Step | Rule | +|---|---| +| `begin_streamed_capture` | reads the binding from the committed root alone (`NoLiveMigration` if none; the catalog pages are not read, so starting a capture does not first materialise a catalog that can hold a million records), routes the journal key through the key-epoch registry (`JournalRoute` errors; a route to a key the journal was not sealed under is `Open(Tampered)` from the authenticated load of the root-named manifest), loads the committed manifest by the exact root-named path, and starts the stage-one capture over it (`StreamedCapture::begin`, so the token, the open inventory and the announced total are checked as before). The caller supplies the journal source, the owner's token and the announced total. Nothing is created and no lock is taken | +| `StagingSession::push_page` / `finish` | the stage-one operations under the key and the directory the session holds, so staging never handles a key and spells the journal directory once | +| `commit_streamed_inventory_capture` | no longer takes a committed manifest or a journal directory: it promotes into the directory the capture was staged in, and loads the committed manifest from the root under the routed key while the lock is held, inside the shared core | + +Decisions, disclosed: (a) the session reads the manifest and the binding itself, which goes beyond the QNB-11 +text (it names only the key); (b) `StreamedInventoryCapture` changes shape (stage two (a) had no caller); +(c) the session holds the routed key for its lifetime and the commit routes again, so a key that went stale in +between is refused at promotion; (d) staging stays lock-free, so every check made at the start is repeated by +the commit under the root lock. + +Proof: a whole capture through a session with no key, manifest or binding in the caller's hands is committed, +verifies, and equals the in-memory capture of the same pages; a session is refused before anything is created +with a stale token and with no bound migration, and it starts with an unreadable catalog page because it reads the root alone; a finished capture is discarded with no key; a commit after the owner moved the journal on is refused before +any write with the staged files kept; the route refuses the session with a revoked or unregistered epoch and +with a route to another key, beside the other journal-owner operations. Mutation-checked: a fixed journal key +in the session, no binding requirement, the whole catalog loaded at the start, and a wrong committed revision in +the commit each fail the test that owns them. The pre-existing stage-two (a) cases pass unchanged against the smaller input. + ## 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. @@ -718,7 +744,7 @@ reference per page. | `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 | +| `StagedCapture::discard` | removes the staged page files of a capture that will not be promoted; needs only the file system, no key, because a caller that staged through a session never held one | | `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 @@ -964,7 +990,6 @@ 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 (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 f13db1196..9083de979 100644 --- a/docs/native/r15/GATE4D-SLICE-B-GAP-MATRIX.md +++ b/docs/native/r15/GATE4D-SLICE-B-GAP-MATRIX.md @@ -91,6 +91,7 @@ load of the newer generation yet). - 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 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). +- Streaming capture, stage two (b): `begin_streamed_capture` returns a `StagingSession` that routes the journal key from the root, reads the binding and the committed manifest from it and holds the journal directory once, so staging handles no key; `commit_streamed_inventory_capture` loses its committed-manifest and journal-directory inputs (it promotes into the staged directory and loads the manifest from the root under the lock). - 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)