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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
134 changes: 123 additions & 11 deletions crates/worldscript-secure-storage/src/authority.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -446,16 +446,16 @@ pub fn commit_inventory_capture<F: DurableFs, P: KeyProvider>(
}

/// 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.
Expand All @@ -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<F: DurableFs, P: KeyProvider>(
fs: &mut F,
provider: &mut P,
Expand All @@ -481,21 +483,26 @@ pub fn commit_streamed_inventory_capture<F: DurableFs, P: KeyProvider>(
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,
Expand All @@ -506,6 +513,111 @@ pub fn commit_streamed_inventory_capture<F: DurableFs, P: KeyProvider>(
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<F: DurableFs, P: KeyProvider>(
fs: &mut F,
provider: &P,
layout: RootLayout<'_>,
begin: SessionBegin<'_>,
) -> Result<StagingSession, AuthorityError> {
// 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<F: DurableFs>(
self,
fs: &mut F,
entries: Vec<JournalInventoryEntry>,
) -> Result<Self, AuthorityError> {
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<F: DurableFs>(self, fs: &mut F) -> Result<StagedCapture, AuthorityError> {
Comment thread
qnbs marked this conversation as resolved.
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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<F: DurableFs>(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<F: DurableFs>(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.
Expand Down
12 changes: 6 additions & 6 deletions crates/worldscript-secure-storage/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand All @@ -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,
Expand Down
Loading
Loading