From 834196cfe0244819d452f5a03bf0f0ea59b98c06 Mon Sep 17 00:00:00 2001 From: James Ross Date: Wed, 7 Oct 2026 09:19:36 -0700 Subject: [PATCH 1/3] feat(cas): add a fallible quarantined complete-object port --- CHANGELOG.md | 4 + crates/echo-cas/README.md | 6 + crates/echo-cas/src/disk.rs | 2 +- crates/echo-cas/src/lib.rs | 1 + crates/echo-cas/src/physical_content.rs | 360 ++++++++++++++++++ .../echo-cas/tests/common/physical_content.rs | 129 +++++++ crates/echo-cas/tests/physical_content.rs | 81 ++++ .../echo-keep-physical-content-boundary.md | 11 +- 8 files changed, 591 insertions(+), 3 deletions(-) create mode 100644 crates/echo-cas/src/physical_content.rs create mode 100644 crates/echo-cas/tests/common/physical_content.rs create mode 100644 crates/echo-cas/tests/physical_content.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 571990e8..65fc74fb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,10 @@ ## Unreleased +### Added + +- A fallible complete-object CAS port stages and verifies exact bytes before atomic destination promotion. Memory and disk adapters share conformance checks; existing APIs remain compatible and no durability or authenticated absence is claimed. + ### Fixed - The generic operation runner names typed obstruction kinds, footprint conflicts, and missing outcomes in bounded summaries. It omits raw outcome records and invocation data on both Action error paths. diff --git a/crates/echo-cas/README.md b/crates/echo-cas/README.md index 5db4106d..dcb133a7 100644 --- a/crates/echo-cas/README.md +++ b/crates/echo-cas/README.md @@ -10,3 +10,9 @@ Content-addressed blob store for Echo. that must survive process reconstruction while preserving content-only BLAKE3 hash semantics. Because filesystem writes can fail, `DiskTier` exposes fallible methods directly instead of hiding I/O behind the infallible `BlobStore` trait. + +## Complete-object physical-content port + +`physical_content` adds fallible expected staging, explicit backend publication, borrowed read views, and atomic destination promotion. Both MemoryTier and DiskTier pass one shared suite. Staging requires explicit per-object byte bounds; it does not enforce aggregate MemoryTier capacity. Missing content and unsupported output capabilities return operational errors rather than authenticated absence. Receipts establish no synchronization or crash durability. Existing consumers keep their current APIs. + +[The physical-content boundary](../../docs/architecture/echo-keep-physical-content-boundary.md) owns the evidence and visibility contract. diff --git a/crates/echo-cas/src/disk.rs b/crates/echo-cas/src/disk.rs index 1631c2ce..19c6a9ac 100644 --- a/crates/echo-cas/src/disk.rs +++ b/crates/echo-cas/src/disk.rs @@ -198,7 +198,7 @@ impl DiskTier { self.pins.len() } - fn blob_path(&self, hash: &BlobHash) -> PathBuf { + pub(crate) fn blob_path(&self, hash: &BlobHash) -> PathBuf { let hex = blob_hash_hex(hash); self.blobs_dir.join(&hex[..2]).join(hex) } diff --git a/crates/echo-cas/src/lib.rs b/crates/echo-cas/src/lib.rs index 8c67e448..a6979f6f 100644 --- a/crates/echo-cas/src/lib.rs +++ b/crates/echo-cas/src/lib.rs @@ -24,6 +24,7 @@ mod disk; mod memory; +pub mod physical_content; mod retention; pub use disk::{DiskTier, DiskTierError}; pub use memory::MemoryTier; diff --git a/crates/echo-cas/src/physical_content.rs b/crates/echo-cas/src/physical_content.rs new file mode 100644 index 00000000..c75961a0 --- /dev/null +++ b/crates/echo-cas/src/physical_content.rs @@ -0,0 +1,360 @@ +// SPDX-License-Identifier: Apache-2.0 +// © James Ross Ω FLYING•ROBOTS +//! Fallible complete-object physical-content operations. +//! +//! Private bounded staging precedes atomic destination promotion. Receipts bind +//! exact bytes and length; these adapters establish no authenticated absence, +//! pinned generation, retention, synchronization, or crash-durability evidence. + +use crate::{blob_hash, BlobHash, BlobStore, DiskTier, DiskTierError, MemoryTier}; +use std::io::{self, Read, Write}; + +/// An explicit Echo content identity and exact expected logical length. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct ContentTarget { + /// Raw-content Echo BLAKE3 identity. + pub hash: BlobHash, + /// Exact logical byte count. + pub length: u64, +} + +/// Failure carries no complete content or absence claim. +#[derive(Debug, thiserror::Error)] +pub enum ContentError { + /// Underlying I/O failed; the original cause is retained. + #[error("physical content I/O failed: {0}")] + Io(#[from] io::Error), + /// Existing filesystem CAS publication failed. + #[error(transparent)] + Disk(#[from] DiskTierError), + /// Exact content hash or length did not match. + #[error("physical content identity or length mismatch")] + Mismatch, + /// Explicit byte or allocation budget was exceeded. + #[error("physical content resource limit")] + ResourceLimit, + /// The requested evidence or atomic output capability is unavailable. + #[error("physical content capability unavailable")] + CapabilityUnavailable, + /// An adapter operation failed without a physical-content proposition. + #[error("physical backend operation failed: {0}")] + Backend(#[source] Box), +} + +/// Bytes sealed only after exact Echo identity and length verification. +/// +/// The sealed handle proves bytes, not storage or causal authority. +#[derive(Debug)] +pub struct VerifiedContent { + target: ContentTarget, + bytes: Vec, +} +impl VerifiedContent { + /// Seals an owned bounded buffer after exact verification. + /// + /// # Errors + /// Returns `Mismatch` for wrong length or hash. + pub fn seal(target: ContentTarget, bytes: Vec) -> Result { + if u64::try_from(bytes.len()).ok() != Some(target.length) + || blob_hash(&bytes) != target.hash + { + return Err(ContentError::Mismatch); + } + Ok(Self { target, bytes }) + } + /// Returns the exact verified target. + pub const fn target(&self) -> ContentTarget { + self.target + } + /// Borrows the complete authenticated bytes. + pub fn bytes(&self) -> &[u8] { + &self.bytes + } + /// Transfers the sealed bytes to an atomic destination implementation. + pub fn into_bytes(self) -> Vec { + self.bytes + } +} + +/// A destination capability with all-or-nothing visible publication. +/// +/// Implementations MUST leave prior visible output unchanged on promotion +/// failure and MUST publish the complete sealed handle atomically on success. +/// Ordinary prefix-writing `Write` sinks do not implement this contract. +pub trait TransactionalContentDestination { + /// True only when atomic promotion is supported. + fn supports_atomic_promotion(&self) -> bool; + /// Maximum bytes allowed in private staging and visible output. + fn byte_limit(&self) -> usize; + /// Atomically replaces visible output with verified complete bytes. + /// + /// # Errors + /// Returns an operational failure, preserving prior visible output. + fn promote(&mut self, content: VerifiedContent) -> Result<(), ContentError>; +} + +/// A bounded memory destination whose promotion replaces one owned buffer. +#[derive(Debug)] +pub struct MemoryContentDestination { + visible: Vec, + limit: usize, +} +impl MemoryContentDestination { + /// Creates a destination with an explicit bound and prior visible bytes. + /// + /// # Errors + /// Returns `ResourceLimit` when prior output exceeds the bound. + pub fn new(visible: Vec, limit: usize) -> Result { + if visible.len() > limit { + return Err(ContentError::ResourceLimit); + } + Ok(Self { visible, limit }) + } + /// Borrows only application-visible bytes, never private staging. + pub fn visible(&self) -> &[u8] { + &self.visible + } +} +impl TransactionalContentDestination for MemoryContentDestination { + fn supports_atomic_promotion(&self) -> bool { + true + } + fn byte_limit(&self) -> usize { + self.limit + } + fn promote(&mut self, content: VerifiedContent) -> Result<(), ContentError> { + if content.bytes.len() > self.limit { + return Err(ContentError::ResourceLimit); + } + self.visible = content.into_bytes(); + Ok(()) + } +} + +/// A complete-object success with explicitly unsupported durability evidence. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct ContentReceipt { + target: ContentTarget, +} +impl ContentReceipt { + /// Constructs an exact-byte receipt for an already sealed object. + /// + /// Adapters may emit it only after their named operation succeeds. It grants + /// no storage, generation, absence, retention, or durability evidence. + pub const fn from_verified(content: &VerifiedContent) -> Self { + Self { + target: content.target, + } + } + /// Returns the target bound to the complete successful operation. + pub const fn target(&self) -> ContentTarget { + self.target + } + /// These initial adapters establish no synchronization or crash durability. + pub const fn establishes_durability(&self) -> bool { + false + } + /// These initial adapters establish no authenticated complete-view absence. + pub const fn establishes_complete_view(&self) -> bool { + false + } +} + +/// Invisible, bounded expected-ingestion work; dropping it publishes nothing. +#[derive(Debug)] +pub struct StagedContent(VerifiedContent); +impl StagedContent { + /// Reads and seals a finite source under an explicit allocation byte bound. + /// + /// # Errors + /// Returns source I/O, resource failure, or exact identity/length mismatch. + pub fn read_expected( + target: ContentTarget, + source: &mut dyn Read, + byte_limit: usize, + ) -> Result { + let mut buffer = PrivateContentBuffer::new(byte_limit); + io::copy(source, &mut buffer).map_err(|error| buffer.map_error(error))?; + Ok(Self(VerifiedContent::seal(target, buffer.bytes)?)) + } + /// Returns its exact already-verified target. + pub const fn target(&self) -> ContentTarget { + self.0.target + } + /// Transfers complete staged bytes into an adapter's private publication. + pub fn into_verified(self) -> VerifiedContent { + self.0 + } +} + +/// Borrowed physical view; its concrete capability binds the source backend. +/// +/// A view cannot outlive its store borrow. The existing disk view authenticates +/// bytes at read time; it does not certify a pinned filesystem generation. +pub trait PhysicalContentView { + /// Reconstructs a complete object into a transactional destination. + /// + /// # Errors + /// On any error, no new bytes are visible and no success receipt is emitted. + fn reconstruct( + &self, + target: ContentTarget, + destination: &mut dyn TransactionalContentDestination, + ) -> Result; +} + +/// Expected staged ingestion, explicit publication, and an immutable read view. +pub trait PhysicalContentBackend { + /// Borrowed capability for complete-object reading. + type View<'a>: PhysicalContentView + where + Self: 'a; + /// Creates a view tied to this backend borrow. + fn content_view(&self) -> Self::View<'_>; + /// Publishes already verified staged content under this backend's posture. + /// + /// # Errors + /// Returns backend I/O or resource failure without a publication receipt. + fn publish_content(&mut self, content: StagedContent) -> Result; +} + +/// Bounded private reconstruction and exact verification shared by adapters. +/// +/// The callback receives only a private writer. A complete successful callback +/// is still rehashed before promotion. Its writer may contain an untrusted +/// prefix on failure, but that prefix is dropped without reaching the caller. +/// +/// # Errors +/// Returns failure without changing visible output, including failed promotion. +pub fn reconstruct_quarantined( + target: ContentTarget, + destination: &mut dyn TransactionalContentDestination, + reconstruct: impl FnOnce(&mut dyn Write) -> Result<(), ContentError>, +) -> Result { + if !destination.supports_atomic_promotion() { + return Err(ContentError::CapabilityUnavailable); + } + let limit = destination.byte_limit(); + if usize::try_from(target.length) + .ok() + .is_none_or(|n| n > limit) + { + return Err(ContentError::ResourceLimit); + } + let mut staging = PrivateContentBuffer::new(limit); + if let Err(error) = reconstruct(&mut staging) { + return Err(if staging.resource_failure { + ContentError::ResourceLimit + } else { + error + }); + } + if staging.resource_failure { + return Err(ContentError::ResourceLimit); + } + let verified = VerifiedContent::seal(target, staging.bytes)?; + destination.promote(verified)?; + Ok(ContentReceipt { target }) +} + +struct PrivateContentBuffer { + bytes: Vec, + limit: usize, + resource_failure: bool, +} +impl PrivateContentBuffer { + const fn new(limit: usize) -> Self { + Self { + bytes: Vec::new(), + limit, + resource_failure: false, + } + } + fn map_error(&self, error: io::Error) -> ContentError { + if self.resource_failure { + ContentError::ResourceLimit + } else { + ContentError::Io(error) + } + } +} +impl Write for PrivateContentBuffer { + fn write(&mut self, bytes: &[u8]) -> io::Result { + let length = self.bytes.len().checked_add(bytes.len()); + if length.is_none_or(|n| n > self.limit) + || self.bytes.try_reserve_exact(bytes.len()).is_err() + { + self.resource_failure = true; + return Err(io::Error::other("private content allocation limit")); + } + self.bytes.extend_from_slice(bytes); + Ok(bytes.len()) + } + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } +} + +/// Immutable in-process CAS read capability. +pub struct MemoryContentView<'a>(&'a MemoryTier); +impl PhysicalContentView for MemoryContentView<'_> { + fn reconstruct( + &self, + target: ContentTarget, + destination: &mut dyn TransactionalContentDestination, + ) -> Result { + reconstruct_quarantined(target, destination, |output| { + let bytes = self + .0 + .get(&target.hash) + .ok_or(ContentError::CapabilityUnavailable)?; + output.write_all(&bytes)?; + Ok(()) + }) + } +} +impl PhysicalContentBackend for MemoryTier { + type View<'a> = MemoryContentView<'a>; + fn content_view(&self) -> Self::View<'_> { + MemoryContentView(self) + } + fn publish_content(&mut self, content: StagedContent) -> Result { + let target = content.target(); + self.put_verified(target.hash, content.0.bytes()) + .map_err(|_| ContentError::Mismatch)?; + Ok(ContentReceipt { target }) + } +} + +/// Borrowed disk CAS capability without a pinned-generation assertion. +pub struct DiskContentView<'a>(&'a DiskTier); +impl PhysicalContentView for DiskContentView<'_> { + fn reconstruct( + &self, + target: ContentTarget, + destination: &mut dyn TransactionalContentDestination, + ) -> Result { + reconstruct_quarantined(target, destination, |output| { + let mut file = match std::fs::File::open(self.0.blob_path(&target.hash)) { + Ok(file) => file, + Err(error) if error.kind() == io::ErrorKind::NotFound => { + return Err(ContentError::CapabilityUnavailable) + } + Err(error) => return Err(error.into()), + }; + io::copy(&mut file, output)?; + Ok(()) + }) + } +} +impl PhysicalContentBackend for DiskTier { + type View<'a> = DiskContentView<'a>; + fn content_view(&self) -> Self::View<'_> { + DiskContentView(self) + } + fn publish_content(&mut self, content: StagedContent) -> Result { + let target = content.target(); + self.put_verified(target.hash, content.0.bytes())?; + Ok(ContentReceipt { target }) + } +} diff --git a/crates/echo-cas/tests/common/physical_content.rs b/crates/echo-cas/tests/common/physical_content.rs new file mode 100644 index 00000000..afe9e930 --- /dev/null +++ b/crates/echo-cas/tests/common/physical_content.rs @@ -0,0 +1,129 @@ +// SPDX-License-Identifier: Apache-2.0 +// © James Ross Ω FLYING•ROBOTS +//! Shared backend-neutral complete-object witness, also used by Keep conformance. +use echo_cas::{ + blob_hash, + physical_content::{ + ContentError, ContentTarget, MemoryContentDestination, PhysicalContentBackend, + PhysicalContentView, StagedContent, TransactionalContentDestination, VerifiedContent, + }, +}; +use std::io::{self, Cursor, Read}; + +type TestResult = Result<(), Box>; + +/// Runs the complete-object contract against any admitted backend. +pub fn conformance(backend: &mut impl PhysicalContentBackend) -> TestResult { + for bytes in [ + Vec::new(), + b"exact bytes".to_vec(), + (0..=255).cycle().take(262_145).collect(), + ] { + let target = ContentTarget { + hash: blob_hash(&bytes), + length: u64::try_from(bytes.len())?, + }; + let mut destination = MemoryContentDestination::new(b"prior".to_vec(), bytes.len().max(5))?; + assert!(matches!( + backend.content_view().reconstruct(target, &mut destination), + Err(ContentError::CapabilityUnavailable) + )); + let staged = StagedContent::read_expected(target, &mut Cursor::new(&bytes), bytes.len())?; + // Staging has no observable backend effect. + assert!(backend + .content_view() + .reconstruct(target, &mut destination) + .is_err()); + let receipt = backend.publish_content(staged)?; + assert_eq!(receipt.target(), target); + assert!(!receipt.establishes_durability()); + assert!(!receipt.establishes_complete_view()); + let receipt = backend + .content_view() + .reconstruct(target, &mut destination)?; + assert_eq!(receipt.target(), target); + assert_eq!(destination.visible(), bytes); + assert!(!receipt.establishes_durability()); + let wrong = ContentTarget { + length: target.length + 1, + ..target + }; + let mut unchanged = MemoryContentDestination::new(b"prior".to_vec(), bytes.len() + 6)?; + assert!(matches!( + backend.content_view().reconstruct(wrong, &mut unchanged), + Err(ContentError::Mismatch) + )); + assert_eq!(unchanged.visible(), b"prior"); + let mut failed = FailingDestination { + visible: b"prior".to_vec(), + capable: true, + limit: bytes.len() + 6, + }; + assert!(matches!( + backend.content_view().reconstruct(target, &mut failed), + Err(ContentError::Io(_)) + )); + assert_eq!(failed.visible, b"prior"); + failed.capable = false; + assert!(matches!( + backend.content_view().reconstruct(target, &mut failed), + Err(ContentError::CapabilityUnavailable) + )); + assert_eq!(failed.visible, b"prior"); + if !bytes.is_empty() { + let mut limited = MemoryContentDestination::new(Vec::new(), bytes.len() - 1)?; + assert!(matches!( + backend.content_view().reconstruct(target, &mut limited), + Err(ContentError::ResourceLimit) + )); + assert!(limited.visible().is_empty()); + } + } + let target = ContentTarget { + hash: blob_hash(b"good"), + length: 4, + }; + assert!(matches!( + StagedContent::read_expected(target, &mut Cursor::new(b"bad!"), 4), + Err(ContentError::Mismatch) + )); + assert!(matches!( + StagedContent::read_expected(target, &mut Cursor::new(b"good"), 3), + Err(ContentError::ResourceLimit) + )); + assert!( + matches!(StagedContent::read_expected(target, &mut FailedSource(false), 4), Err(ContentError::Io(e)) if e.kind() == io::ErrorKind::BrokenPipe) + ); + assert!(backend + .content_view() + .reconstruct(target, &mut MemoryContentDestination::new(Vec::new(), 4)?) + .is_err()); + Ok(()) +} +struct FailedSource(bool); +impl Read for FailedSource { + fn read(&mut self, output: &mut [u8]) -> io::Result { + if self.0 { + return Err(io::ErrorKind::BrokenPipe.into()); + } + self.0 = true; + output[0] = 103; + Ok(1) + } +} +struct FailingDestination { + visible: Vec, + capable: bool, + limit: usize, +} +impl TransactionalContentDestination for FailingDestination { + fn supports_atomic_promotion(&self) -> bool { + self.capable + } + fn byte_limit(&self) -> usize { + self.limit + } + fn promote(&mut self, _: VerifiedContent) -> Result<(), ContentError> { + Err(io::Error::other("injected promotion failure").into()) + } +} diff --git a/crates/echo-cas/tests/physical_content.rs b/crates/echo-cas/tests/physical_content.rs new file mode 100644 index 00000000..fb8baaef --- /dev/null +++ b/crates/echo-cas/tests/physical_content.rs @@ -0,0 +1,81 @@ +// SPDX-License-Identifier: Apache-2.0 +// © James Ross Ω FLYING•ROBOTS +//! Complete-object conformance and quarantine faults. +#[path = "common/physical_content.rs"] +mod common; +use echo_cas::{ + blob_hash, + physical_content::{ + reconstruct_quarantined, ContentError, ContentTarget, MemoryContentDestination, + PhysicalContentBackend, PhysicalContentView, + }, + DiskTier, MemoryTier, +}; +use std::fs; +use std::io; +type TestResult = Result<(), Box>; +#[test] +fn memory_complete_object_conformance() -> TestResult { + common::conformance(&mut MemoryTier::new()) +} +#[test] +fn disk_complete_object_conformance_and_corruption_quarantine() -> TestResult { + let root = std::env::temp_dir().join(format!("echo-physical-content-{}", std::process::id())); + let _cleanup = Cleanup(root.clone()); + fs::create_dir_all(&root)?; + let mut backend = DiskTier::open(&root)?; + common::conformance(&mut backend)?; + let bytes = b"corrupt me"; + let hash = backend.put(bytes)?; + let hex = hash.to_string(); + let path = root.join("blobs").join(&hex[..2]).join(hex); + fs::write(&path, b"corrupt xx")?; + let target = ContentTarget { hash, length: 10 }; + let mut destination = MemoryContentDestination::new(b"prior".to_vec(), 16)?; + assert!(matches!( + backend.content_view().reconstruct(target, &mut destination), + Err(ContentError::Mismatch) + )); + assert_eq!(destination.visible(), b"prior"); + // Publication I/O failure cannot produce a publication receipt. + fs::remove_file(&path)?; + fs::create_dir(&path)?; + let staged = echo_cas::physical_content::StagedContent::read_expected( + target, + &mut io::Cursor::new(bytes), + 10, + )?; + assert!(matches!( + backend.publish_content(staged), + Err(ContentError::Disk(_)) + )); + Ok(()) +} +struct Cleanup(std::path::PathBuf); +impl Drop for Cleanup { + fn drop(&mut self) { + let _ = fs::remove_dir_all(&self.0); + } +} +#[test] +fn prefix_and_ignored_writer_failure_never_reach_visible_output() -> TestResult { + let target = ContentTarget { + hash: blob_hash(b"good"), + length: 4, + }; + let mut destination = MemoryContentDestination::new(b"old".to_vec(), 4)?; + let error = reconstruct_quarantined(target, &mut destination, |output| { + output.write_all(b"go")?; + Err(io::Error::from(io::ErrorKind::BrokenPipe).into()) + }); + assert!(matches!(error, Err(ContentError::Io(e)) if e.kind() == io::ErrorKind::BrokenPipe)); + assert_eq!(destination.visible(), b"old"); + let error = reconstruct_quarantined(target, &mut destination, |output| { + output.write_all(b"good")?; + let _ = output.write_all(b"overflow"); + Ok(()) + }); + assert!(matches!(error, Err(ContentError::ResourceLimit))); + assert_eq!(destination.visible(), b"old"); + Ok(()) +} diff --git a/docs/architecture/echo-keep-physical-content-boundary.md b/docs/architecture/echo-keep-physical-content-boundary.md index 0ab8507f..1d74c7d4 100644 --- a/docs/architecture/echo-keep-physical-content-boundary.md +++ b/docs/architecture/echo-keep-physical-content-boundary.md @@ -6,8 +6,7 @@ - **Status:** Accepted for experimental conformance; production adoption is not accepted. - **Decision date:** 2026-08-09 -- **Implementation posture:** No Echo physical-content port or Keep adapter is - implemented on this branch. +- **Implementation posture:** `echo-cas::physical_content` supplies a fallible complete-object port and MemoryTier/DiskTier adapters. Borrowed views authenticate exact bytes but certify no pinned generation, complete-view absence, retention, synchronization, or crash durability. No Keep backend adapter is implemented yet. - **Refines:** [Retained reading storage and proof boundary](../adr/0020-retained-reading-storage-and-proof-boundary.md) - **Depends on:** [Durable external-action settlement](../adr/0026-durable-external-action-settlement.md) - **Related:** [Keep authenticated reconstruction contract](https://github.com/flyingrobots/keep/blob/3bf7b9179db41e90620e6d1875c2d40222a2330b/docs/architecture/authenticated-reconstruction-contract.md) @@ -115,6 +114,14 @@ Echo may retain a bounded materializing helper implemented over this port. The helper is not Keep's foundational contract and must require an explicit byte limit. +## Initial executable port + +`StagedContent::read_expected` reads under an explicit per-object byte limit and seals the Echo hash and exact length. `PhysicalContentBackend::publish_content` publishes that invisible staged object. Existing MemoryTier and DiskTier APIs remain compatible. + +`PhysicalContentView::reconstruct` uses bounded private staging. `TransactionalContentDestination` requires atomic promotion of a sealed `VerifiedContent` handle; an arbitrary `Write` sink is insufficient. `MemoryContentDestination` provides the initial bounded memory implementation. A failed source, integrity check, resource check, or promotion leaves prior visible bytes intact. The shared backend-neutral suite lives in `crates/echo-cas/tests/common/physical_content.rs`. + +The materializing port bounds each staging operation, not aggregate retained MemoryTier capacity or process RSS. Existing MemoryTier budgets remain advisory. Missing content is `CapabilityUnavailable`, with no authenticated absence receipt. Disk views verify bytes at read time without a pinned-generation claim. All initial receipts explicitly report unsupported durability and complete-view evidence. This additive port does not reroute current consumers or adopt Keep in production. + ## Output visibility An ordinary `Write` sink can fail after accepting a prefix. Keep may therefore From e0494292a72cc123b6df0c9cb2f2cc0ab4587c4f Mon Sep 17 00:00:00 2001 From: James Ross Date: Wed, 7 Oct 2026 09:46:32 -0700 Subject: [PATCH 2/3] fix(cas): preflight staging and reject overlong content as corruption --- crates/echo-cas/src/physical_content.rs | 55 ++++++++++++------- .../echo-cas/tests/common/physical_content.rs | 2 + crates/echo-cas/tests/physical_content.rs | 41 +++++++++++++- 3 files changed, 77 insertions(+), 21 deletions(-) diff --git a/crates/echo-cas/src/physical_content.rs b/crates/echo-cas/src/physical_content.rs index c75961a0..b9604aec 100644 --- a/crates/echo-cas/src/physical_content.rs +++ b/crates/echo-cas/src/physical_content.rs @@ -173,7 +173,12 @@ impl StagedContent { source: &mut dyn Read, byte_limit: usize, ) -> Result { - let mut buffer = PrivateContentBuffer::new(byte_limit); + let expected_length = + usize::try_from(target.length).map_err(|_| ContentError::ResourceLimit)?; + if expected_length > byte_limit { + return Err(ContentError::ResourceLimit); + } + let mut buffer = PrivateContentBuffer::new(byte_limit, expected_length); io::copy(source, &mut buffer).map_err(|error| buffer.map_error(error))?; Ok(Self(VerifiedContent::seal(target, buffer.bytes)?)) } @@ -235,23 +240,17 @@ pub fn reconstruct_quarantined( return Err(ContentError::CapabilityUnavailable); } let limit = destination.byte_limit(); - if usize::try_from(target.length) - .ok() - .is_none_or(|n| n > limit) - { + let expected_length = + usize::try_from(target.length).map_err(|_| ContentError::ResourceLimit)?; + if expected_length > limit { return Err(ContentError::ResourceLimit); } - let mut staging = PrivateContentBuffer::new(limit); - if let Err(error) = reconstruct(&mut staging) { - return Err(if staging.resource_failure { - ContentError::ResourceLimit - } else { - error - }); - } - if staging.resource_failure { - return Err(ContentError::ResourceLimit); + let mut staging = PrivateContentBuffer::new(limit, expected_length); + let result = reconstruct(&mut staging); + if let Some(error) = staging.failure() { + return Err(error); } + result?; let verified = VerifiedContent::seal(target, staging.bytes)?; destination.promote(verified)?; Ok(ContentReceipt { target }) @@ -261,26 +260,42 @@ struct PrivateContentBuffer { bytes: Vec, limit: usize, resource_failure: bool, + identity_failure: bool, + expected_length: usize, } impl PrivateContentBuffer { - const fn new(limit: usize) -> Self { + const fn new(limit: usize, expected_length: usize) -> Self { Self { bytes: Vec::new(), limit, resource_failure: false, + identity_failure: false, + expected_length, } } - fn map_error(&self, error: io::Error) -> ContentError { - if self.resource_failure { - ContentError::ResourceLimit + fn failure(&self) -> Option { + if self.identity_failure { + Some(ContentError::Mismatch) + } else if self.resource_failure { + Some(ContentError::ResourceLimit) } else { - ContentError::Io(error) + None } } + fn map_error(&self, error: io::Error) -> ContentError { + self.failure().unwrap_or(ContentError::Io(error)) + } } impl Write for PrivateContentBuffer { fn write(&mut self, bytes: &[u8]) -> io::Result { let length = self.bytes.len().checked_add(bytes.len()); + if length.is_none_or(|n| n > self.expected_length) { + self.identity_failure = true; + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "overlong physical content", + )); + } if length.is_none_or(|n| n > self.limit) || self.bytes.try_reserve_exact(bytes.len()).is_err() { diff --git a/crates/echo-cas/tests/common/physical_content.rs b/crates/echo-cas/tests/common/physical_content.rs index afe9e930..9368bd0d 100644 --- a/crates/echo-cas/tests/common/physical_content.rs +++ b/crates/echo-cas/tests/common/physical_content.rs @@ -36,6 +36,8 @@ pub fn conformance(backend: &mut impl PhysicalContentBackend) -> TestResult { .is_err()); let receipt = backend.publish_content(staged)?; assert_eq!(receipt.target(), target); + let repeated = StagedContent::read_expected(target, &mut Cursor::new(&bytes), bytes.len())?; + assert_eq!(backend.publish_content(repeated)?.target(), target); assert!(!receipt.establishes_durability()); assert!(!receipt.establishes_complete_view()); let receipt = backend diff --git a/crates/echo-cas/tests/physical_content.rs b/crates/echo-cas/tests/physical_content.rs index fb8baaef..c47e700f 100644 --- a/crates/echo-cas/tests/physical_content.rs +++ b/crates/echo-cas/tests/physical_content.rs @@ -75,7 +75,46 @@ fn prefix_and_ignored_writer_failure_never_reach_visible_output() -> TestResult let _ = output.write_all(b"overflow"); Ok(()) }); - assert!(matches!(error, Err(ContentError::ResourceLimit))); + assert!(matches!(error, Err(ContentError::Mismatch))); assert_eq!(destination.visible(), b"old"); Ok(()) } + +#[test] +fn impossible_staging_budget_does_not_read_source() { + struct UnreadSource(bool); + impl std::io::Read for UnreadSource { + fn read(&mut self, _: &mut [u8]) -> io::Result { + self.0 = true; + Err(io::ErrorKind::BrokenPipe.into()) + } + } + let mut source = UnreadSource(false); + let target = ContentTarget { + hash: blob_hash(b"good"), + length: 4, + }; + assert!(matches!( + echo_cas::physical_content::StagedContent::read_expected(target, &mut source, 3), + Err(ContentError::ResourceLimit) + )); + assert!(!source.0); +} + +#[test] +fn overlong_reconstruction_is_mismatch_at_any_sufficient_budget() -> TestResult { + let target = ContentTarget { + hash: blob_hash(b"good"), + length: 4, + }; + for limit in [4, 8] { + let mut destination = MemoryContentDestination::new(b"old".to_vec(), limit)?; + let result = reconstruct_quarantined(target, &mut destination, |output| { + output.write_all(b"goodx")?; + Ok(()) + }); + assert!(matches!(result, Err(ContentError::Mismatch))); + assert_eq!(destination.visible(), b"old"); + } + Ok(()) +} From 7802d898a933e41fdccd7a8e4651a1e295b4e3b8 Mon Sep 17 00:00:00 2001 From: James Ross Date: Wed, 7 Oct 2026 10:07:38 -0700 Subject: [PATCH 3/3] test(cas): allocate disk conformance roots exclusively --- crates/echo-cas/tests/physical_content.rs | 25 +++++++++++++++++++---- 1 file changed, 21 insertions(+), 4 deletions(-) diff --git a/crates/echo-cas/tests/physical_content.rs b/crates/echo-cas/tests/physical_content.rs index c47e700f..8b8fd86e 100644 --- a/crates/echo-cas/tests/physical_content.rs +++ b/crates/echo-cas/tests/physical_content.rs @@ -20,10 +20,9 @@ fn memory_complete_object_conformance() -> TestResult { } #[test] fn disk_complete_object_conformance_and_corruption_quarantine() -> TestResult { - let root = std::env::temp_dir().join(format!("echo-physical-content-{}", std::process::id())); - let _cleanup = Cleanup(root.clone()); - fs::create_dir_all(&root)?; - let mut backend = DiskTier::open(&root)?; + let cleanup = fresh_disk_fixture()?; + let root = &cleanup.0; + let mut backend = DiskTier::open(root)?; common::conformance(&mut backend)?; let bytes = b"corrupt me"; let hash = backend.put(bytes)?; @@ -51,6 +50,24 @@ fn disk_complete_object_conformance_and_corruption_quarantine() -> TestResult { )); Ok(()) } +fn fresh_disk_fixture() -> io::Result { + for slot in 0..1024 { + let root = std::env::temp_dir().join(format!( + "echo-physical-content-{}-{slot}", + std::process::id() + )); + match fs::create_dir(&root) { + Ok(()) => return Ok(Cleanup(root)), + Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {} + Err(error) => return Err(error), + } + } + Err(io::Error::new( + io::ErrorKind::AlreadyExists, + "no fresh disk fixture slot", + )) +} + struct Cleanup(std::path::PathBuf); impl Drop for Cleanup { fn drop(&mut self) {