From 11be73dd07becca98c88f4c69ed69536acacb447 Mon Sep 17 00:00:00 2001 From: James Ross Date: Fri, 2 Oct 2026 12:30:48 -0700 Subject: [PATCH 1/3] Test: reject writable conversion of sealed receipts (#146) --- src/adapters/sealed_segment.rs | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/src/adapters/sealed_segment.rs b/src/adapters/sealed_segment.rs index 70b5b988..1393f739 100644 --- a/src/adapters/sealed_segment.rs +++ b/src/adapters/sealed_segment.rs @@ -7,6 +7,19 @@ use super::{ClosedSegment, SegmentDigest, SegmentStage}; /// /// The stage remains unpublished. This type exposes no mutable stage handle, /// makes no directory-durability claim, and is not a catalog reference. +/// +/// Sealing cannot hand writable authority to a conversion callback, including +/// with `repository-tasks` enabled: +/// +/// ```compile_fail +/// use keep::{SealedSegment, SegmentStage}; +/// fn rewrite(sealed: SealedSegment) { +/// let _ = sealed.map_stage(|mut stage| { +/// let _ = std::io::Write::write_all(&mut stage, b"!"); +/// stage +/// }); +/// } +/// ``` #[must_use] pub struct SealedSegment where From 038b283df9375424b44b1b68036def2e1ade927c Mon Sep 17 00:00:00 2001 From: James Ross Date: Fri, 2 Oct 2026 12:41:57 -0700 Subject: [PATCH 2/3] Fix: retain sealed stage authority through crash observation (#146) --- CHANGELOG.md | 2 + docs/formats/segment-store-v1/rationale.md | 4 + docs/formats/segment-store-v1/requirements.md | 2 +- .../sealed-stage-observation.md | 30 ++++ src/adapters/exports.rs | 4 + src/adapters/mod.rs | 4 + src/adapters/observed_segment_stage.rs | 84 ++++++++++ src/adapters/sealed_segment.rs | 16 -- src/adapters/segment_stage_observer.rs | 41 +++++ src/lib.rs | 9 +- tests/observed_segment_stage.rs | 147 ++++++++++++++++++ .../byte_equivalence.rs | 49 ++++++ .../durability_failure.rs | 62 ++++++++ .../production_protocol/publication.rs | 9 +- .../production_protocol/segment_stage.rs | 83 +++++----- 15 files changed, 477 insertions(+), 69 deletions(-) create mode 100644 docs/testing-evidence/sealed-stage-observation.md create mode 100644 src/adapters/observed_segment_stage.rs create mode 100644 src/adapters/segment_stage_observer.rs create mode 100644 tests/observed_segment_stage.rs create mode 100644 tests/observed_segment_stage/byte_equivalence.rs create mode 100644 tests/observed_segment_stage/durability_failure.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index c23e7c27..31401267 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,8 @@ after its public API and format compatibility policies are established. ## [Unreleased] +- Sealed segment receipts no longer expose writable stages through `map_stage`. Repository crash injection uses a private-stage observation wrapper whose sealed conversion preserves stage identity and publisher authority without handing storage to callbacks (#146). + - Retention recovery execution errors report the exact failed boundary, original typed cause, known namespace effects and uncertain effect/durability; retries freshly observe the store. Observed stage identity remains binding across reopening, and cleanup preserves verified pool evidence rather than promising the removed pathname survives (#99). - Retention recovery now preserves incomplete stages and requires explicit disposition before any recovery mutation or publication retry; automatic incomplete-stage disposal is deferred by maintainer decision (#99). diff --git a/docs/formats/segment-store-v1/rationale.md b/docs/formats/segment-store-v1/rationale.md index 7c03e0e1..878d40d3 100644 --- a/docs/formats/segment-store-v1/rationale.md +++ b/docs/formats/segment-store-v1/rationale.md @@ -176,6 +176,10 @@ filesystem whose individual operations returned success. Each would let an unsupported platform manufacture the authority that the proof is meant to represent. +## Sealed-stage observation + +Repository segment observation keeps writable stage authority private (#146). The former sealed `map_stage` callback could mutate an already synchronized stage while keeping the original receipt metadata; publication-time revalidation did not make that receipt truthful. `ObservedSegmentStage` now owns the stage and performs its actual write, flush and synchronization operations, exposing only lengths and before/after durability events to its observer. Its specialized sealed `without_observer` conversion drops the observer and preserves the same hidden stage and metadata together. There is no arbitrary user conversion of a sealed stage. Returning an error from an after-write hook cannot undo bytes; an `Interrupted` cause is retained inside a non-retryable I/O error so ordinary write retry cannot duplicate those effects. A new arbitrary unwrap trait or caller-supplied closure would reopen the capability escape and was rejected. + ## Observation before recovery Store opening performs no repair. It produces either one verified reader diff --git a/docs/formats/segment-store-v1/requirements.md b/docs/formats/segment-store-v1/requirements.md index 526fd25c..b0f69b46 100644 --- a/docs/formats/segment-store-v1/requirements.md +++ b/docs/formats/segment-store-v1/requirements.md @@ -47,7 +47,7 @@ retention, or garbage collection. | `KEEP-SEGMENT-002` | Staged and sealed writer states are distinct, consuming types | `tests/segment_writer.rs` | Implemented in #15 | | `KEEP-SEGMENT-003` | Short, interrupted, zero-progress, invalid-count, storage, permission, flush, and synchronization failures retain exact phases and offsets | `tests/segment_writer/write_contract_laws.rs`, `tests/segment_writer/refusal_laws.rs`, `tests/segment_writer/durability_laws.rs` | Implemented in #15 | | `KEEP-SEGMENT-004` | Complete-segment admission verifies bounds, framing, checksums, logical identities, duplicate refusal, terminal state, and physical digest before exposure | `tests/segment.rs`, `tests/segment/identity_laws.rs`, `tests/segment/framing_laws.rs` | Implemented in #15 | -| `KEEP-SEGMENT-005` | The public sealed receipt exposes no mutable stage handle | `src/adapters/sealed_segment.rs` | Production API implemented in #15; repository-task escape tracked in #146 | +| `KEEP-SEGMENT-005` | The public sealed receipt exposes no mutable stage handle, including with repository-task observation enabled | `src/adapters/sealed_segment.rs` compile-fail law; `tests/observed_segment_stage.rs`; [capability evidence](../../testing-evidence/sealed-stage-observation.md) | Production API implemented in #15; repository-task escape removed for #146 | | `KEEP-SEGMENT-006` | Malformed, unsupported, partial, conflicting, and corrupt input returns boundary-typed errors | `tests/segment_header/mutation_laws.rs`, `tests/segment_record_header/framing_laws.rs`, `tests/segment_seal/framing_laws.rs`, `tests/segment/identity_laws.rs` | Implemented in #15 | | `KEEP-SEGMENT-007` | Record, nested-layout, segment-length, and temporary identity-index allocation remain explicitly bounded | `tests/segment_memory.rs`, `tests/segment_record_memory.rs`, `tests/segment_seal_memory.rs` | Implemented in #15 | | `KEEP-SEGMENT-008` | Filesystem staging uses exclusive fixed-name creation and never enumerates storage as a content index | `src/adapters/filesystem_segment_stage_tests.rs`, `src/adapters/filesystem_segment_stage.rs` | Implemented in #15 | diff --git a/docs/testing-evidence/sealed-stage-observation.md b/docs/testing-evidence/sealed-stage-observation.md new file mode 100644 index 00000000..9861b48b --- /dev/null +++ b/docs/testing-evidence/sealed-stage-observation.md @@ -0,0 +1,30 @@ +# Sealed-stage observation capability + +Change kind: bug fix for #146, with a cohesive replacement observation boundary for its crash-harness consumer. Owner: `@flyingrobots`. The branch starts at main `6051abb25a9fd33ae7ee0de5614514b709a4d82a`; it is independent of the public-stage and catalog-admission follow-up branches. No on-disk format, identity, publication phase or new dependency is introduced. Removing the feature-gated public `map_stage` method intentionally breaks callers that requested writable authority from an already sealed receipt. + +## RED and protected boundary + +The permanent compile-fail example in `SealedSegment` attempts to call `map_stage` and write another byte through its callback. On the unfixed parent, `cargo test --doc --all-features sealed_segment --locked` failed with “Test compiled successfully, but it's marked compile_fail.” Regression commit `11be73d` preserves that RED. This is compiler-enforced capability evidence; it is not represented as a filesystem failure. With the arbitrary conversion removed, the same example passes by refusing access to that API. + +The replacement stores both stage and observer in private fields. Observation receives only requested/actual write lengths and explicit before/after flush/synchronize events. `without_observer` is implemented only for the built-in wrapper and returns another sealed receipt over the same stage; no callback or public extractor receives writable storage. Generic `close` still drops the stage. Filesystem selection still checks its original publisher authority and admitted segment coordinates. The usual `SegmentStage` exclusive-ownership precondition remains; a caller's independently duplicated storage handle is not made exclusive by a wrapper. + +## Runtime evidence and oracles + +| Claim | Evidence and oracle | +| --- | --- | +| Removing observation preserves bytes and private publisher provenance | Actual admitted ext4 stage equals the independently specified empty-segment golden; selection by its original publisher succeeds after observer removal. | +| Short-write observation preserves ordinary writer output over the tested input family | Generated payload lengths 1 through 64, with a fixed repeating byte pattern, compare observed seven-byte-prefix writes against plain filesystem writing. Every case owns fresh storage. Ascending exhaustive exploration reports the smallest failing length in this family; the independent golden complements this shared-implementation differential oracle. | +| Excessive observer limits refuse before writes | Exact header `SegmentWriteError::Write`, zero prior completed bytes and `InvalidInput`; the actual stage remains empty. | +| A post-write interruption cannot retry completed effects | Exact header write refusal retains the nested original `Interrupted` cause, with a non-retryable outer kind; actual bytes equal exactly one golden header. The prior-completed-call offset is not a rollback claim for this failed call. | +| Observation cannot fabricate successful durability | Faulting stage ports drive the production writer through the observation wrapper; exact prefix flush/synchronize failures preserve errno 5 and produce no sealed receipt. These are port-level fault simulations, not physical disk failure evidence. | +| Crash injection remains connected to production operations | Complete existing debug and optimized process-death matrices pass after replacing the writable decorator with `CrashSegmentObserver`. Header, record and seal interruption offsets, before/during/after choices, flush/sync ordering, restart bytes and immutable admission oracles remain unchanged. | + +Separate copied mutations demonstrated runtime RED for skipped underlying flush, skipped underlying synchronization, clamped excessive limits, retryable post-effect interruption, and zeroed writes. Each used a dedicated build directory. Zeroed writes also failed the generated differential assertion at its first payload length. These are distinct outcome calibrations, not a mutation score or harness-case-count claim. The mutation sources were not committed. + +## Execution and limits + +All Rust execution uses copied Docker source, pinned Rust 1.96.0 and dedicated output on Linux aarch64. Filesystem laws use actual production platform admission on private ext4 scratch; no writable host checkout mount or fabricated admission proof is used. The capability regression is run with `--all-features`, since that is what exposed the old method. Both feature configurations are checked with warnings-denied Clippy. The final PR records exact candidate and hosted-check SHAs alongside full debug/release, formatting, source-structure, documentation and existing golden/corruption/memory evidence. + +Filesystem laws are medium; the faulting durability-port laws and capability example are small/static evidence. No random seeds, sleeps or uncontrolled scheduling drive the oracles. The differential input family is bounded and deterministic; it does not establish equivalence for arbitrary payloads or malicious storage implementations. Replay with `cargo test --test observed_segment_stage --all-features --locked`, its `--release` variant, and the all-feature doctest command above. Re-run the production campaigns with `cargo xtask durability-crash-matrix` and its optimized invocation. + +Existing resource-enforcement and suite-budget gaps remain those disclosed by the binding testing enforcement profile; this change claims neither a new sandbox policy nor measured latency SLOs. Process death and simulated errno failures do not establish physical power-loss behavior. No performance optimization is claimed. Retire the laws only if the capability is removed or stronger boundary evidence demonstrably subsumes the risk. Original roadmap checkboxes and unrelated recovery/admission scope remain unchanged. diff --git a/src/adapters/exports.rs b/src/adapters/exports.rs index e2c8c51f..2a6badf1 100644 --- a/src/adapters/exports.rs +++ b/src/adapters/exports.rs @@ -64,6 +64,8 @@ pub use super::layout_encode_error::LayoutEncodeError; pub use super::layout_id_binary_error::LayoutIdBinaryParseError; pub use super::layout_id_text_error::LayoutIdTextParseError; pub use super::layout_record::CanonicalLayoutRecord; +#[cfg(feature = "repository-tasks")] +pub use super::observed_segment_stage::ObservedSegmentStage; pub use super::opened_reusable_segment::OpenedReusableSegment; pub use super::publication_head_decode_error::PublicationHeadDecodeError; pub use super::recovery::*; @@ -92,6 +94,8 @@ pub use super::segment_seal::SegmentSeal; pub use super::segment_seal_error::SegmentSealError; pub use super::segment_stage::SegmentStage; pub use super::segment_stage_create_error::SegmentStageCreateError; +#[cfg(feature = "repository-tasks")] +pub use super::segment_stage_observer::{SegmentStageDurabilityEvent, SegmentStageObserver}; pub use super::segment_write_error::SegmentWriteError; pub use super::segment_write_phase::{SegmentDurabilityPhase, SegmentWritePhase}; pub use super::staged_segment::StagedSegment; diff --git a/src/adapters/mod.rs b/src/adapters/mod.rs index 39627605..db478fee 100644 --- a/src/adapters/mod.rs +++ b/src/adapters/mod.rs @@ -151,6 +151,8 @@ mod layout_record_format; mod layout_record_framing; mod loaded_segment; mod lower_hex; +#[cfg(feature = "repository-tasks")] +mod observed_segment_stage; mod opened_reusable_segment; mod physical_pool_name; mod publication_head_decode_error; @@ -213,6 +215,8 @@ mod segment_seal_hash; mod segment_stage; mod segment_stage_create_error; mod segment_stage_create_error_display; +#[cfg(feature = "repository-tasks")] +mod segment_stage_observer; mod segment_stage_write; mod segment_write_error; mod segment_write_error_display; diff --git a/src/adapters/observed_segment_stage.rs b/src/adapters/observed_segment_stage.rs new file mode 100644 index 00000000..8245358e --- /dev/null +++ b/src/adapters/observed_segment_stage.rs @@ -0,0 +1,84 @@ +//! This module owns transparent stage observation with private writable authority. + +use std::io::{self, Write}; + +use super::{SealedSegment, SegmentStage, SegmentStageDurabilityEvent, SegmentStageObserver}; + +/// A repository stage decorator whose observer never receives its inner stage. +/// +/// Owns the stage until close or sealing. Construction allocates nothing; actual +/// writes and durability operations block according to the supplied stage and +/// observer. The supplied stage must satisfy [`SegmentStage`]'s ownership contract. +pub struct ObservedSegmentStage { + stage: S, + observer: O, +} + +impl ObservedSegmentStage { + /// Installs observation before the stage is consumed by the writer. + pub const fn new(stage: S, observer: O) -> Self { + Self { stage, observer } + } +} + +impl Write for ObservedSegmentStage { + fn write(&mut self, bytes: &[u8]) -> io::Result { + let allowed = self.observer.before_write(bytes.len())?; + let prefix = bytes.get(..allowed).ok_or_else(|| { + io::Error::new( + io::ErrorKind::InvalidInput, + "observer write limit exceeds input", + ) + })?; + let written = self.stage.write(prefix)?; + if written > prefix.len() { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "stage exceeded observed write limit", + )); + } + self.observer.after_write(written).map_err(|source| { + // Interrupted means no write happened to Write::write_all. The + // observer runs after real effects, so preserve the cause without + // allowing a caller to silently repeat those bytes. + if source.kind() == io::ErrorKind::Interrupted { + io::Error::other(source) + } else { + source + } + })?; + Ok(written) + } + + fn flush(&mut self) -> io::Result<()> { + self.observer + .durability(SegmentStageDurabilityEvent::BeforeFlush)?; + self.stage.flush()?; + self.observer + .durability(SegmentStageDurabilityEvent::AfterFlush) + } +} + +impl SegmentStage for ObservedSegmentStage { + fn synchronize(&mut self) -> io::Result<()> { + self.observer + .durability(SegmentStageDurabilityEvent::BeforeSynchronize)?; + self.stage.synchronize()?; + self.observer + .durability(SegmentStageDurabilityEvent::AfterSynchronize) + } +} + +impl SealedSegment> { + /// Removes observation while keeping the sealed stage inaccessible to callers. + /// + /// Drops only the observer. The original stage, exact metadata and completed + /// durability evidence stay together; no user callback receives the stage. + /// Filesystem publication still checks actual bytes and publisher authority. + pub fn without_observer(self) -> SealedSegment { + let (wrapped, count, length, digest) = self.into_parts(); + let ObservedSegmentStage { stage, observer } = wrapped; + drop(observer); + SealedSegment::admitted(stage, count, length, digest) + } +} diff --git a/src/adapters/sealed_segment.rs b/src/adapters/sealed_segment.rs index 1393f739..ac0c802a 100644 --- a/src/adapters/sealed_segment.rs +++ b/src/adapters/sealed_segment.rs @@ -66,22 +66,6 @@ where self.digest } - /// Replaces a repository-only storage decorator while preserving sealed - /// metadata. - /// - /// This supports transparent storage-port decorators. The mapping does not - /// change, revalidate, or publish the sealed bytes; authority-bound - /// adapters still validate the returned stage before publication. - #[cfg(feature = "repository-tasks")] - #[doc(hidden)] - pub fn map_stage(self, map: impl FnOnce(S) -> T) -> SealedSegment - where - T: SegmentStage, - { - let (stage, record_count, segment_length, digest) = self.into_parts(); - SealedSegment::admitted(map(stage), record_count, segment_length, digest) - } - pub(super) fn into_parts(self) -> (S, u32, u64, SegmentDigest) { let Self { _stage: stage, diff --git a/src/adapters/segment_stage_observer.rs b/src/adapters/segment_stage_observer.rs new file mode 100644 index 00000000..756a992e --- /dev/null +++ b/src/adapters/segment_stage_observer.rs @@ -0,0 +1,41 @@ +//! This module owns repository fault-observation hooks without stage authority. + +use std::io; + +/// An observation boundary around a real stage durability operation. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum SegmentStageDurabilityEvent { + /// Before flushing buffered bytes. + BeforeFlush, + /// After a successful flush. + AfterFlush, + /// Before synchronizing persisted bytes. + BeforeSynchronize, + /// After a successful synchronization. + AfterSynchronize, +} + +/// Repository fault hooks that can interrupt operations but cannot access storage. +/// +/// The wrapper always performs actual writes, flushes and synchronization itself. +/// Observers can block or fail at a boundary, or limit a write to a strict prefix. +/// No hook receives a writable stage or can authorize a sealed receipt. +pub trait SegmentStageObserver { + /// Selects how many of the requested bytes may reach the next actual write. + /// + /// # Errors + /// Returns the injected pre-write failure. Limits beyond `requested` are refused. + fn before_write(&mut self, requested: usize) -> io::Result; + + /// Observes the actual successful write count before the caller continues. + /// + /// # Errors + /// Returns an injected failure after the write; this does not undo its effects. + fn after_write(&mut self, written: usize) -> io::Result<()>; + + /// Observes a boundary before or after an actual durability operation. + /// + /// # Errors + /// Returns the injected failure at this boundary, without rolling back effects. + fn durability(&mut self, event: SegmentStageDurabilityEvent) -> io::Result<()>; +} diff --git a/src/lib.rs b/src/lib.rs index 1127d095..7c1be6a7 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -55,9 +55,6 @@ mod profile; mod reference; mod retention; -#[cfg(feature = "repository-tasks")] -#[doc(hidden)] -pub use adapters::RepositoryInitializationStorage; pub use adapters::{ AdmittedCatalog, AdmittedRecoveryStageBytes, AdmittedSegment, AdmittedSegmentRecord, AdmittedStoreFormatMarker, AdmittedStoreMigrationIntent, AdmittedStoreMigrationReceipt, @@ -167,6 +164,12 @@ pub use adapters::{ StoreMigrationRecoveryReceipt, StoreMigrationRecoveryStorage, StoreMigrationResidue, StoreMigrationStageDecodeError, plan_store_migration_recovery, recover_store_migration, }; +#[cfg(feature = "repository-tasks")] +#[doc(hidden)] +pub use adapters::{ + ObservedSegmentStage, RepositoryInitializationStorage, SegmentStageDurabilityEvent, + SegmentStageObserver, +}; pub use blob::{ BlobHashError, BlobHasher, BlobId, BlobLength, BlobReadError, ByteLength, ByteOffset, ByteRange, ByteRangeError, diff --git a/tests/observed_segment_stage.rs b/tests/observed_segment_stage.rs new file mode 100644 index 00000000..66f0eaef --- /dev/null +++ b/tests/observed_segment_stage.rs @@ -0,0 +1,147 @@ +//! Medium public filesystem laws for sealed observation; oracle: canonical bytes or typed refusal. +//! Delete only when the observation capability is removed or stronger boundary laws replace these. + +#![cfg(all(target_os = "linux", feature = "repository-tasks"))] + +#[path = "observed_segment_stage/byte_equivalence.rs"] +mod byte_equivalence; +#[path = "observed_segment_stage/durability_failure.rs"] +mod durability_failure; +#[path = "segment_filesystem_stage/sandbox.rs"] +pub mod sandbox; +mod support; + +use std::error::Error; +use std::fs; +use std::io; + +use keep::{ + AdmittedSegment, CatalogRestartByteLimit, CatalogRestartPolicy, FilesystemCatalogPublisher, + FilesystemPlatformAdmission, LayoutEntryLimit, ObservedSegmentStage, SegmentReadPolicy, + SegmentRecordLimit, SegmentStageDurabilityEvent, SegmentStageObserver, SegmentWriteError, + SegmentWritePhase, StagedSegment, +}; + +#[test] +fn removing_observation_preserves_canonical_bytes_and_publisher_authority() +-> Result<(), Box> { + let directory = sandbox::TestDirectory::create("sealed-observation")?; + let publisher = publisher(&directory)?; + let stage = ObservedSegmentStage::new(publisher.create_segment_stage()?, Observer::Prefixes); + let sealed = StagedSegment::begin(stage, SegmentRecordLimit::MAXIMUM)? + .seal()? + .without_observer(); + let canonical = support::decode_hex( + include_str!("../conformance/segment-store/v1/empty-segment.hex").trim_end(), + )?; + assert_eq!( + fs::read(directory.path().join("staging/current.seg"))?, + canonical + ); + let admitted = AdmittedSegment::decode(&canonical, read_policy())?; + let _selection = publisher.select_segment(sealed, &admitted)?; + drop(publisher); + directory.remove()?; + Ok(()) +} + +#[test] +fn an_excessive_observation_limit_refuses_before_writing() -> Result<(), Box> { + let directory = sandbox::TestDirectory::create("invalid-observation-limit")?; + let publisher = publisher(&directory)?; + let stage = ObservedSegmentStage::new(publisher.create_segment_stage()?, Observer::Excessive); + let Err(error) = StagedSegment::begin(stage, SegmentRecordLimit::MAXIMUM) else { + return Err("excessive observation limit was admitted".into()); + }; + assert!( + matches!(error, SegmentWriteError::Write { phase: SegmentWritePhase::Header, bytes_written: 0, ref source } + if source.kind() == io::ErrorKind::InvalidInput) + ); + assert_eq!(fs::read(directory.path().join("staging/current.seg"))?, []); + drop(publisher); + directory.remove()?; + Ok(()) +} + +#[test] +fn a_post_write_interruption_cannot_retry_already_written_bytes() -> Result<(), Box> { + let directory = sandbox::TestDirectory::create("post-write-observation-interruption")?; + let publisher = publisher(&directory)?; + let stage = ObservedSegmentStage::new( + publisher.create_segment_stage()?, + Observer::InterruptAfterWrite, + ); + let Err(error) = StagedSegment::begin(stage, SegmentRecordLimit::MAXIMUM) else { + return Err("post-effect interruption was retried into successful admission".into()); + }; + let SegmentWriteError::Write { + phase: SegmentWritePhase::Header, + bytes_written: 0, + source, + } = error + else { + return Err(format!("unexpected interruption failure: {error}").into()); + }; + assert_eq!(source.kind(), io::ErrorKind::Other); + let original = source + .get_ref() + .and_then(|error| error.downcast_ref::()) + .ok_or("original interruption cause lost")?; + assert_eq!(original.kind(), io::ErrorKind::Interrupted); + let canonical = support::decode_hex( + include_str!("../conformance/segment-store/v1/empty-segment.hex").trim_end(), + )?; + assert_eq!( + fs::read(directory.path().join("staging/current.seg"))?, + canonical.get(..64).ok_or("missing golden header")? + ); + drop(publisher); + directory.remove()?; + Ok(()) +} + +enum Observer { + Prefixes, + Excessive, + InterruptAfterWrite, +} + +impl SegmentStageObserver for Observer { + fn before_write(&mut self, requested: usize) -> io::Result { + match self { + Self::Prefixes => Ok(requested.min(7)), + Self::Excessive => requested + .checked_add(1) + .ok_or_else(|| io::Error::other("limit overflow")), + Self::InterruptAfterWrite => Ok(requested), + } + } + + fn after_write(&mut self, _written: usize) -> io::Result<()> { + match self { + Self::InterruptAfterWrite => { + *self = Self::Prefixes; + Err(io::Error::from(io::ErrorKind::Interrupted)) + } + Self::Prefixes | Self::Excessive => Ok(()), + } + } + + fn durability(&mut self, _event: SegmentStageDurabilityEvent) -> io::Result<()> { + Ok(()) + } +} + +fn publisher( + directory: &sandbox::TestDirectory, +) -> Result> { + let admission = FilesystemPlatformAdmission::initialize(directory.path())?; + Ok(FilesystemCatalogPublisher::open( + admission, + CatalogRestartPolicy::new(read_policy(), CatalogRestartByteLimit::new(1_048_576)?), + )?) +} + +const fn read_policy() -> SegmentReadPolicy { + SegmentReadPolicy::new(SegmentRecordLimit::MAXIMUM, LayoutEntryLimit::MAXIMUM) +} diff --git a/tests/observed_segment_stage/byte_equivalence.rs b/tests/observed_segment_stage/byte_equivalence.rs new file mode 100644 index 00000000..75c72f9f --- /dev/null +++ b/tests/observed_segment_stage/byte_equivalence.rs @@ -0,0 +1,49 @@ +//! Medium filesystem differential evidence over an explicitly bounded payload space. + +use std::error::Error; +use std::fs; + +use keep::{AdmittedSegmentRecord, ObservedSegmentStage, SegmentRecordLimit, StagedSegment}; + +use super::{Observer, publisher, sandbox}; + +#[test] +fn observed_short_writes_match_plain_writes_across_generated_payload_lengths() +-> Result<(), Box> { + // Ascending exhaustive length exploration reports the smallest failing + // member of this input family. Empty segments have an independent golden. + for length in 1..=64 { + let payload: Vec<_> = b"generated bytes" + .iter() + .copied() + .cycle() + .take(length) + .collect(); + let plain_dir = sandbox::TestDirectory::create(&format!("plain-stage-{length}"))?; + let observed_dir = sandbox::TestDirectory::create(&format!("observed-stage-{length}"))?; + let plain = publisher(&plain_dir)?; + let observed = publisher(&observed_dir)?; + let plain_sealed = + StagedSegment::begin(plain.create_segment_stage()?, SegmentRecordLimit::MAXIMUM)? + .append(AdmittedSegmentRecord::for_chunk(&payload)?)? + .seal()?; + let decorated = + ObservedSegmentStage::new(observed.create_segment_stage()?, Observer::Prefixes); + let observed_sealed = StagedSegment::begin(decorated, SegmentRecordLimit::MAXIMUM)? + .append(AdmittedSegmentRecord::for_chunk(&payload)?)? + .seal()? + .without_observer(); + assert_eq!( + fs::read(observed_dir.path().join("staging/current.seg"))?, + fs::read(plain_dir.path().join("staging/current.seg"))?, + "observed/plain storage differs for generated payload length {length}" + ); + drop(plain_sealed); + drop(observed_sealed); + drop(plain); + drop(observed); + plain_dir.remove()?; + observed_dir.remove()?; + } + Ok(()) +} diff --git a/tests/observed_segment_stage/durability_failure.rs b/tests/observed_segment_stage/durability_failure.rs new file mode 100644 index 00000000..6820ff01 --- /dev/null +++ b/tests/observed_segment_stage/durability_failure.rs @@ -0,0 +1,62 @@ +//! Small port-level laws: observation cannot conceal real stage durability failure. +//! Oracle: `SegmentStage`'s fallible flush/synchronize contract, not physical disk evidence. + +use std::error::Error; +use std::io::{self, Write}; + +use keep::{ + ObservedSegmentStage, SegmentDurabilityPhase, SegmentRecordLimit, SegmentStage, + SegmentWriteError, StagedSegment, +}; + +use super::Observer; + +#[derive(Clone, Copy)] +enum Failure { + Flush, + Synchronize, +} + +struct FailingStage(Failure); + +impl Write for FailingStage { + fn write(&mut self, bytes: &[u8]) -> io::Result { + Ok(bytes.len()) + } + fn flush(&mut self) -> io::Result<()> { + match self.0 { + Failure::Flush => Err(io::Error::from_raw_os_error(5)), + Failure::Synchronize => Ok(()), + } + } +} + +impl SegmentStage for FailingStage { + fn synchronize(&mut self) -> io::Result<()> { + Err(io::Error::from_raw_os_error(5)) + } +} + +#[test] +fn observation_preserves_the_underlying_flush_refusal() -> Result<(), Box> { + let stage = ObservedSegmentStage::new(FailingStage(Failure::Flush), Observer::Prefixes); + let Err(error) = StagedSegment::begin(stage, SegmentRecordLimit::MAXIMUM)?.seal() else { + return Err("observation manufactured sealing despite failed flush".into()); + }; + assert!( + matches!(error, SegmentWriteError::Flush { phase: SegmentDurabilityPhase::RecordPrefix, ref source } if source.raw_os_error() == Some(5)) + ); + Ok(()) +} + +#[test] +fn observation_preserves_the_underlying_synchronization_refusal() -> Result<(), Box> { + let stage = ObservedSegmentStage::new(FailingStage(Failure::Synchronize), Observer::Prefixes); + let Err(error) = StagedSegment::begin(stage, SegmentRecordLimit::MAXIMUM)?.seal() else { + return Err("observation manufactured sealing despite failed synchronization".into()); + }; + assert!( + matches!(error, SegmentWriteError::Synchronize { phase: SegmentDurabilityPhase::RecordPrefix, ref source } if source.raw_os_error() == Some(5)) + ); + Ok(()) +} diff --git a/xtask/src/durability_crash_matrix/production_protocol/publication.rs b/xtask/src/durability_crash_matrix/production_protocol/publication.rs index 0f2a4f65..b692ab68 100644 --- a/xtask/src/durability_crash_matrix/production_protocol/publication.rs +++ b/xtask/src/durability_crash_matrix/production_protocol/publication.rs @@ -5,14 +5,15 @@ use std::path::Path; use keep::{ AdmittedSegment, AdmittedSegmentRecord, CanonicalCatalog, CatalogGeneration, - CatalogPublicationExpectation, SegmentRecordLimit, StagedSegment, publish_catalog_generation, + CatalogPublicationExpectation, ObservedSegmentStage, SegmentRecordLimit, StagedSegment, + publish_catalog_generation, }; use xtask::DurabilityCrashPoint; use super::control::{CrashControl, DuringTiming}; use super::initialization; use super::publication_storage::CrashPublicationStorage; -use super::segment_stage::CrashSegmentStage; +use super::segment_stage::CrashSegmentObserver; use super::{DurabilityCrashMatrixError, verification}; pub(super) fn run( @@ -31,7 +32,7 @@ pub(super) fn run( .after(point, DuringTiming::After) .map_err(crash_gate)?; - let stage = CrashSegmentStage::new(stage, control); + let stage = ObservedSegmentStage::new(stage, CrashSegmentObserver::new(control)); let record = AdmittedSegmentRecord::for_chunk(&[0]) .map_err(|source| verification("admit crash segment record", source))?; let sealed = StagedSegment::begin(stage, SegmentRecordLimit::MAXIMUM) @@ -40,7 +41,7 @@ pub(super) fn run( .map_err(|source| verification("append production segment record", source))? .seal() .map_err(|source| verification("seal production segment", source))? - .map_stage(CrashSegmentStage::into_inner); + .without_observer(); let segment_bytes = fs::read(store_root.join("staging/current.seg")) .map_err(|source| DurabilityCrashMatrixError::io("read sealed crash segment", source))?; diff --git a/xtask/src/durability_crash_matrix/production_protocol/segment_stage.rs b/xtask/src/durability_crash_matrix/production_protocol/segment_stage.rs index 68cf9d9e..6889b81a 100644 --- a/xtask/src/durability_crash_matrix/production_protocol/segment_stage.rs +++ b/xtask/src/durability_crash_matrix/production_protocol/segment_stage.rs @@ -1,8 +1,8 @@ //! This module owns crash injection around the production segment-stage port. -use std::io::{self, Write}; +use std::io; -use keep::SegmentStage; +use keep::{SegmentStageDurabilityEvent, SegmentStageObserver}; use xtask::{DurabilityCrashPoint, DurabilityCrashPosition}; use super::control::{CrashControl, DuringTiming}; @@ -14,25 +14,19 @@ const RECORD_INTERRUPTION: usize = 136; const SEAL_END: usize = 337; const SEAL_INTERRUPTION: usize = 273; -pub(super) struct CrashSegmentStage<'control, S> { - inner: S, +pub(super) struct CrashSegmentObserver<'control> { control: &'control mut CrashControl, bytes_written: usize, } -impl<'control, S> CrashSegmentStage<'control, S> { - pub(super) const fn new(inner: S, control: &'control mut CrashControl) -> Self { +impl<'control> CrashSegmentObserver<'control> { + pub(super) const fn new(control: &'control mut CrashControl) -> Self { Self { - inner, control, bytes_written: 0, } } - pub(super) fn into_inner(self) -> S { - self.inner - } - fn write_boundary(&self) -> io::Result<(DurabilityCrashPoint, usize, usize)> { match self.bytes_written { 0..HEADER_END => Ok(( @@ -71,11 +65,8 @@ impl<'control, S> CrashSegmentStage<'control, S> { } } -impl Write for CrashSegmentStage<'_, S> -where - S: SegmentStage, -{ - fn write(&mut self, bytes: &[u8]) -> io::Result { +impl SegmentStageObserver for CrashSegmentObserver<'_> { + fn before_write(&mut self, requested: usize) -> io::Result { let (point, interruption, end) = self.write_boundary()?; let position = self.control.position(point); if position == Some(DurabilityCrashPosition::Before) { @@ -89,11 +80,12 @@ where let remaining = limit .checked_sub(self.bytes_written) .ok_or_else(|| io::Error::other("segment crash boundary moved backward"))?; - let allowed = remaining.min(bytes.len()); - let prefix = bytes - .get(..allowed) - .ok_or_else(|| io::Error::other("segment write prefix exceeded input"))?; - let written = self.inner.write(prefix)?; + Ok(remaining.min(requested)) + } + + fn after_write(&mut self, written: usize) -> io::Result<()> { + let (point, interruption, end) = self.write_boundary()?; + let position = self.control.position(point); self.bytes_written = self .bytes_written .checked_add(written) @@ -103,31 +95,32 @@ where { self.control.await_process_death()?; } - Ok(written) + Ok(()) } - fn flush(&mut self) -> io::Result<()> { - let point = self.durability_point( - DurabilityCrashPoint::FlushSegmentRecordPrefix, - DurabilityCrashPoint::FlushSealedSegment, - )?; - self.control.before(point, DuringTiming::Before)?; - self.inner.flush()?; - self.control.after(point, DuringTiming::Before) - } -} - -impl SegmentStage for CrashSegmentStage<'_, S> -where - S: SegmentStage, -{ - fn synchronize(&mut self) -> io::Result<()> { - let point = self.durability_point( - DurabilityCrashPoint::SynchronizeSegmentRecordPrefix, - DurabilityCrashPoint::SynchronizeSealedSegment, - )?; - self.control.before(point, DuringTiming::Before)?; - self.inner.synchronize()?; - self.control.after(point, DuringTiming::Before) + fn durability(&mut self, event: SegmentStageDurabilityEvent) -> io::Result<()> { + let point = match event { + SegmentStageDurabilityEvent::BeforeFlush | SegmentStageDurabilityEvent::AfterFlush => { + self.durability_point( + DurabilityCrashPoint::FlushSegmentRecordPrefix, + DurabilityCrashPoint::FlushSealedSegment, + )? + } + SegmentStageDurabilityEvent::BeforeSynchronize + | SegmentStageDurabilityEvent::AfterSynchronize => self.durability_point( + DurabilityCrashPoint::SynchronizeSegmentRecordPrefix, + DurabilityCrashPoint::SynchronizeSealedSegment, + )?, + }; + match event { + SegmentStageDurabilityEvent::BeforeFlush + | SegmentStageDurabilityEvent::BeforeSynchronize => { + self.control.before(point, DuringTiming::Before) + } + SegmentStageDurabilityEvent::AfterFlush + | SegmentStageDurabilityEvent::AfterSynchronize => { + self.control.after(point, DuringTiming::Before) + } + } } } From 15aa976a77d5b65a4b6e8f8a2f9e01d2ba3b0d58 Mon Sep 17 00:00:00 2001 From: James Ross Date: Fri, 2 Oct 2026 12:42:49 -0700 Subject: [PATCH 3/3] Fix: require explicit use of observed stage authority (#146) --- src/adapters/observed_segment_stage.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/adapters/observed_segment_stage.rs b/src/adapters/observed_segment_stage.rs index 8245358e..6af557c0 100644 --- a/src/adapters/observed_segment_stage.rs +++ b/src/adapters/observed_segment_stage.rs @@ -9,6 +9,7 @@ use super::{SealedSegment, SegmentStage, SegmentStageDurabilityEvent, SegmentSta /// Owns the stage until close or sealing. Construction allocates nothing; actual /// writes and durability operations block according to the supplied stage and /// observer. The supplied stage must satisfy [`SegmentStage`]'s ownership contract. +#[must_use] pub struct ObservedSegmentStage { stage: S, observer: O,