Skip to content
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

- **R-15 Gate 4D:** stage a captured inventory one page at a time (stage one of the streaming capture).
`StreamedCapture` takes the entries page by page, assigns the page index and generation, checks that they
are valid and strictly ascending across pages, seals each page once and stages its envelope in a private
`inventory/pending-*` directory (never authority, never read by a reader), keeping one page and one
reference per page in memory; the total is announced up front because the inventory digest commits to it
first. `finish` builds the successor with the same function as `capture_inventory`, and
`load_staged_page` confirms a staged page against its reference, a missing one being a broken attempt and
not a recovery state of the journal. A capture is bound to its journal directory and to the key that opened
its manifest, and a failure the caller can handle removes the pages staged so far. Promotion under the root
lock and the composed commit are the next stage. PR #1010.
- **R-15 Gate 4D:** the journal owner can renew its own lease (maintainer decision C, the wiring half).
`commit_lease_renewal` publishes a manifest that moves the lease expiry strictly forward and changes nothing
else (`publish_renewal_fenced`, a publish of its own so the ordinary checkpoint keeps refusing lease changes)
Expand Down
55 changes: 43 additions & 12 deletions crates/worldscript-secure-storage/src/journal/capture.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,23 +48,52 @@ pub fn capture_inventory(
) -> Result<JournalManifest, MigrationExecutionError> {
assert_fence(manifest, fence)?;
assert_inventory_open(manifest)?;
let revision = manifest
.journal_revision
.checked_add(1)
.ok_or(JournalError::InvalidCounter)?;
let revision = next_revision(manifest)?;
let ordered = ordered_pages(pages)?;
let refs = page_refs(&ordered, revision)?;
let entry_count = entry_total(&ordered)?;
let inventory_digest = inventory_digest_of(manifest.inventory_version, &ordered, entry_count)?;
let parts = CapturedInventory {
refs: &refs,
entry_count,
inventory_digest,
};
capture_successor(manifest, &parts)
}

/// What a capture of the inventory consists of, however its pages were produced: one reference per
/// page, the entry total and the inventory digest over the entries.
pub(super) struct CapturedInventory<'a> {
pub refs: &'a [JournalPageRef],
pub entry_count: u32,
pub inventory_digest: [u8; 32],
}

/// The revision a capture of `manifest` is published at: the next one, and the generation of every
/// page the capture writes.
pub(super) fn next_revision(manifest: &JournalManifest) -> Result<u64, MigrationExecutionError> {
manifest
.journal_revision
.checked_add(1)
.ok_or_else(|| JournalError::InvalidCounter.into())
}

/// The successor of `manifest` that carries `parts` as its inventory, checked as a valid successor.
/// Both the capture from pages in memory and the streamed capture end here.
pub(super) fn capture_successor(
manifest: &JournalManifest,
parts: &CapturedInventory<'_>,
) -> Result<JournalManifest, MigrationExecutionError> {
let mut next = manifest.clone();
// Only the capture that runs behind `ADMIT`'s barrier is the final one (§10.3).
next.final_inventory_captured = manifest.phase == phase_code::ADMIT;
next.journal_revision = revision;
next.page_count = u32::try_from(ordered.len()).map_err(|_| JournalError::TooManyEntries)?;
next.entry_count = entry_total(&ordered)?;
next.inventory_digest =
inventory_digest_of(manifest.inventory_version, &ordered, next.entry_count)?;
next.journal_page_set_digest = journal_page_set_digest(&refs)?;
next.journal_revision = next_revision(manifest)?;
next.page_count = u32::try_from(parts.refs.len()).map_err(|_| JournalError::TooManyEntries)?;
next.entry_count = parts.entry_count;
next.inventory_digest = parts.inventory_digest;
next.journal_page_set_digest = journal_page_set_digest(parts.refs)?;
next.encode()?;
next.verify_page_set(&refs)?;
next.verify_page_set(parts.refs)?;
assert_manifest_successor(manifest, &next)?;
Ok(next)
}
Expand All @@ -73,7 +102,9 @@ pub fn capture_inventory(
/// starts, while nothing has been converted and, because the final inventory is immutable once
/// captured, before the final capture. `PREPARE` is excluded: ordinary writes are admitted
/// there until `ADMIT`'s barrier, so only the `ADMIT` snapshot can be the commit inventory.
fn assert_inventory_open(manifest: &JournalManifest) -> Result<(), MigrationExecutionError> {
pub(super) fn assert_inventory_open(
manifest: &JournalManifest,
) -> Result<(), MigrationExecutionError> {
let phase = manifest_phase(manifest);
if is_terminal_phase(phase) {
return Err(MigrationExecutionError::TerminalPhase);
Expand Down
55 changes: 50 additions & 5 deletions crates/worldscript-secure-storage/src/journal/durable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -546,18 +546,63 @@ pub fn load_authoritative_manifest<F: DurableFs>(
live: &LiveMigration,
) -> Result<JournalManifest, JournalDurableError> {
let identity = migration_identity(&live.operation_id)?;
let (bytes, digest) = read_root_named_envelope(ctx.fs, ctx.dir, live)?;
let (bytes, _) = read_root_named_envelope(ctx.fs, ctx.dir, live)?;
open_root_named(ctx.key, &identity, live, &bytes)
}

/// Opens the root-named envelope `bytes` under `key` and requires the binding to vouch for it: the
/// part of [`load_authoritative_manifest`] that needs no read. The epoch of the header is compared
/// before the key is used.
fn open_root_named(
key: &Key,
identity: &RecordIdentity,
live: &LiveMigration,
bytes: &[u8],
) -> Result<JournalManifest, JournalDurableError> {
let read = ManifestRead {
record: &identity,
record: identity,
journal_revision: live.journal_revision,
key_epoch: header_key_epoch(&bytes)?,
envelope: &bytes,
key_epoch: header_key_epoch(bytes)?,
envelope: bytes,
};
let manifest = JournalManifest::open(ctx.key, &read)?;
let manifest = JournalManifest::open(key, &read)?;
let digest = ManifestEnvelopeDigest::from_bytes(content_digest(bytes));
assert_live_binding(&manifest, live, digest).map_err(JournalDurableError::Authority)?;
Ok(manifest)
}

/// The exact bytes of the generation the root binding names, after proving that they open under the
/// context's key to exactly `manifest`. A caller that works for a long time keeps the bytes and
/// repeats the key check with [`assert_anchor`], which needs no read.
pub(super) fn root_named_anchor<F: DurableFs>(
ctx: &mut JournalDurableContext<'_, F>,
manifest: &JournalManifest,
live: &LiveMigration,
) -> Result<Vec<u8>, JournalDurableError> {
let (bytes, _) = read_root_named_envelope(ctx.fs, ctx.dir, live)?;
assert_anchor(ctx.key, manifest, live, &bytes)?;
Ok(bytes)
}

/// Requires `bytes`, the root-named generation of `live`, to open under `key` to exactly `manifest`.
/// Pure computation over a few kilobytes: no read, so it can be repeated for every page of a
/// capture.
pub(super) fn assert_anchor(
key: &Key,
manifest: &JournalManifest,
live: &LiveMigration,
bytes: &[u8],
) -> Result<(), JournalDurableError> {
let identity = migration_identity(&live.operation_id)?;
if open_root_named(key, &identity, live, bytes)? == *manifest {
Ok(())
} else {
Err(JournalDurableError::Authority(
MigrationExecutionError::LiveBindingMismatch,
))
}
}

/// Reads the generation the committed root names, bounded, and proves the bytes are what the root
/// committed to: authority first (§6), before any key is used. The binding authenticates these exact
/// bytes, header included, so the epoch in the header is vouched for by the root, never by the
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ use super::state::MigrationExecutionError;
use super::{JournalError, MAX_JOURNAL_PAGE_BYTES};

/// The largest valid sealed page: the page encoding bound plus the envelope header and tag.
const MAX_PAGE_ENVELOPE_BYTES: usize = MAX_JOURNAL_PAGE_BYTES + HEADER_LEN + TAG_LEN;
pub(super) const MAX_PAGE_ENVELOPE_BYTES: usize = MAX_JOURNAL_PAGE_BYTES + HEADER_LEN + TAG_LEN;
/// A page directory holds one generation and, at most, a few staging leftovers; a directory with
/// more entries than this is not read.
const MAX_PAGE_DIRECTORY_ENTRIES: usize = 64;
Expand Down Expand Up @@ -142,7 +142,7 @@ fn hinted_generation<F: DurableFs>(
}

/// The envelope at `path`, never larger than any valid sealed page.
fn read_envelope<F: DurableFs>(
pub(super) fn read_envelope<F: DurableFs>(
ctx: &mut JournalDurableContext<'_, F>,
path: &Path,
) -> Result<Vec<u8>, JournalDurableError> {
Expand All @@ -157,7 +157,7 @@ fn read_envelope<F: DurableFs>(
/// Opens `bytes` as this operation's page `index` at `generation` under the journal key. Read routing
/// is authority-first (§6): the header's epoch is compared with the operation's journal envelope
/// epoch before the key is used, so a misrouted page is never decrypted.
fn open_stored_page<F: DurableFs>(
pub(super) fn open_stored_page<F: DurableFs>(
ctx: &JournalDurableContext<'_, F>,
manifest: &JournalManifest,
index: u32,
Expand Down
34 changes: 19 additions & 15 deletions crates/worldscript-secure-storage/src/journal/inventory_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -238,26 +238,30 @@ fn store_page<F: DurableFs>(
&set.successor.journal_page_set_digest,
sealed.page.page_index(),
);
if holds_envelope(ctx, &dir, sealed)? {
store_page_at(ctx, &dir, set.committed_manifest, sealed)
}

/// Stores one page of `manifest`'s operation into the page directory `dir`, which sits under the
/// journal directory: the same adoption, staging, promotion and directory syncs wherever the page set
/// is being assembled.
pub(super) fn store_page_at<F: DurableFs>(
ctx: &mut JournalDurableContext<'_, F>,
dir: &Path,
manifest: &JournalManifest,
sealed: &SealedPage<'_>,
) -> Result<DirectoryDurability, JournalDurableError> {
if holds_envelope(ctx, dir, sealed)? {
// The earlier attempt may have stopped before its directories were durable, and, under the
// same operation id, may have left its own staging link behind.
let staging = staging_residue(ctx, &dir, sealed.page.page_generation());
return sync_chain(ctx, &dir, staging);
let staging = staging_residue(ctx, dir, sealed.page.page_generation());
return sync_chain(ctx, dir, staging);
}
ctx.fs
.create_dir_all(&dir)
.create_dir_all(dir)
.map_err(|error| JournalDurableError::Stage(stage_io(error)))?;
let identity = migration_page_identity(
&set.committed_manifest.operation_id,
sealed.page.page_index(),
)?;
let epoch = journal_envelope_epoch(set.committed_manifest)?;
let request = stage_request(
&dir,
&identity,
page_meta(sealed.page, epoch),
ctx.operation,
);
let identity = migration_page_identity(&manifest.operation_id, sealed.page.page_index())?;
let epoch = journal_envelope_epoch(manifest)?;
let request = stage_request(dir, &identity, page_meta(sealed.page, epoch), ctx.operation);
let promoted = stage_and_promote_envelope(ctx.fs, ctx.key, &request, sealed.envelope.to_vec())?;
let above = dir.parent().unwrap_or(ctx.dir);
let synced = sync_chain(ctx, above, promoted.staging)?;
Expand Down
4 changes: 4 additions & 0 deletions crates/worldscript-secure-storage/src/journal/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ mod manifest_verify;
mod page;
mod renewal;
mod state;
mod stream_capture;
mod succession;
mod takeover;
mod wire;
Expand Down Expand Up @@ -146,6 +147,9 @@ pub use state::{
JournalInventoryExtent, JournalRevision, ManifestEnvelopeDigest, MigrationExecutionError,
MigrationFence, MigrationPhase, RecoveryReasonCode,
};
pub use stream_capture::{
load_staged_page, CaptureStart, StagedCapture, StagedPage, StreamedCapture,
};
pub use succession::{assert_manifest_successor, assert_progress_successor};
pub use takeover::{
assert_binding_takeover, assert_takeover_promote_authority, assert_takeover_successor,
Expand Down
Loading
Loading