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
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
139 changes: 123 additions & 16 deletions crates/worldscript-secure-storage/src/authority.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -425,6 +426,121 @@ pub fn commit_inventory_capture<F: DurableFs, P: KeyProvider>(
capture: InventoryCapture<'_>,
) -> Result<RootCommitted, AuthorityError> {
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<F: DurableFs, P: KeyProvider>(
fs: &mut F,
provider: &mut P,
layout: RootLayout<'_>,
capture: StreamedInventoryCapture<'_>,
) -> Result<RootCommitted, AuthorityError> {
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<F, P, S>(
fs: &mut F,
provider: &mut P,
plan: CapturePlan<'_>,
store_pages: S,
) -> Result<RootCommitted, AuthorityError>
where
F: DurableFs,
P: KeyProvider,
S: FnOnce(
&mut JournalDurableContext<'_, F>,
&LiveMigration,
&MigrationFence,
) -> Result<DirectoryDurability, JournalDurableError>,
{
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
Expand All @@ -439,22 +555,13 @@ pub fn commit_inventory_capture<F: DurableFs, P: KeyProvider>(
// 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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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> {
Expand Down Expand Up @@ -329,7 +329,7 @@ fn sync_chain<F: DurableFs>(
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 {
Expand Down
2 changes: 2 additions & 0 deletions crates/worldscript-secure-storage/src/journal/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ mod page;
mod renewal;
mod state;
mod stream_capture;
mod stream_promote;
mod succession;
mod takeover;
mod wire;
Expand Down Expand Up @@ -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,
Expand Down
44 changes: 29 additions & 15 deletions crates/worldscript-secure-storage/src/journal/stream_capture.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<JournalPageRef>,
Expand Down Expand Up @@ -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<F: DurableFs>(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<F: DurableFs>(&self, fs: &mut F) {
discard_staged(fs, &self.pending, &self.refs, &self.tag);
}
}

Expand Down Expand Up @@ -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)
}
}
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -265,6 +281,7 @@ impl StreamedCapture {
) -> Result<StagedCapture, JournalDurableError> {
let Self {
committed,
journal_dir,
entry_count,
pending,
tag,
Expand All @@ -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)
}
}
Expand All @@ -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<F: DurableFs>(
ctx: &mut JournalDurableContext<'_, F>,
fs: &mut F,
pending: &Path,
refs: &[JournalPageRef],
tag: &WriteOperationId,
Expand All @@ -313,22 +331,18 @@ fn discard_staged<F: DurableFs>(
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<F: DurableFs>(
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<F: DurableFs>(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);
}
}

Expand Down
Loading
Loading