From 9b6f6bfb14f3efcb67a43bbfce365cd9304a4345 Mon Sep 17 00:00:00 2001 From: James Ross Date: Mon, 21 Sep 2026 00:16:34 -0700 Subject: [PATCH 01/15] feat(runtime): retain observation-bound executable requests across host loss --- CHANGELOG.md | 11 + README.md | 6 +- crates/warp-core/src/causal_wal.rs | 14 +- crates/warp-core/src/echo_operation.rs | 76 ++++- .../warp-core/src/echo_operation/observed.rs | 305 +++++++++++++++++ crates/warp-core/src/lib.rs | 23 +- crates/warp-core/src/trusted_runtime_host.rs | 3 + .../trusted_runtime_host/observed_context.rs | 296 ++++++++++++++++ .../tests/trusted_runtime_host_loop_tests.rs | 92 +++++ .../application-contract-hosting.md | 46 +++ xtask/src/main.rs | 13 +- xtask/src/run_edict_operation.rs | 3 + xtask/src/run_edict_operation/session.rs | 317 ++++++++++++++++++ 13 files changed, 1184 insertions(+), 21 deletions(-) create mode 100644 crates/warp-core/src/echo_operation/observed.rs create mode 100644 crates/warp-core/src/trusted_runtime_host/observed_context.rs create mode 100644 xtask/src/run_edict_operation/session.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 3f16af866..6a0cb1385 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,13 @@ ### Added +- Bounded executable-operation host sessions retain immutable observations and + logical-request bindings in the native WAL. Echo evaluates supplied node/atom + preconditions inside operation preparation, includes their reads in scheduler + footprints, and retains typed refusal or commitment outcomes across reopening. + The generic `run-edict-operation --serve` driver supports repeated operations, + historical outcome lookup, and derived observation-change discovery. + - Strict filesystem WAL stores now persist a checksummed writer-epoch ledger containing the active epoch, its exact latest closed predecessor, and final LSN and commit-digest evidence. Bounded retention keeps ledger writes and @@ -1628,6 +1635,10 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- Reopening a filesystem WAL through an empty writer epoch no longer consumes an + unwritten log position. A second reopen followed by append previously left an + LSN gap and made subsequent recovery fail. + - Generic executable-operation lowering and independent verification now resolve source-local obstruction constructor aliases through the exact digest-locked lawpack import before encoding or comparing the package. diff --git a/README.md b/README.md index 76ad26588..d54331d96 100644 --- a/README.md +++ b/README.md @@ -278,7 +278,11 @@ mutation. A standalone application-owned Edict create-if-absent operation now crosses this route as an exact compiler-produced package and independently accepted verification report through `cargo xtask run-edict-operation`. That witness proves its package-declared typed obstruction and one Action in one -Tick. Its duplicate report exposes equal before/after application-state roots +Tick. Its `--serve` mode continues one worldline through multiple operations, +retains observation-bound requests, and reopens native history; the bounded +trusted-host scope is described in +[application contract hosting](docs/architecture/application-contract-hosting.md#observation-bound-native-host-sessions). +Its duplicate report exposes equal before/after application-state roots and typed target-value digests; the roots commit the reachable graph state, not WAL, history, Receipt, or commit metadata. It does not claim external multi-Action Tick composition, a product-ready application runner, a Jedit diff --git a/crates/warp-core/src/causal_wal.rs b/crates/warp-core/src/causal_wal.rs index 0c574e9be..cb2d327a6 100644 --- a/crates/warp-core/src/causal_wal.rs +++ b/crates/warp-core/src/causal_wal.rs @@ -486,6 +486,8 @@ pub enum WalRecordKind { ExternalActionClaimRecorded, /// Echo admitted one schema-bound external-action settlement. ExternalActionSettlementRecorded, + /// Runtime retained an immutable operation observation or request binding. + ExecutableOperationContextRetained, } impl WalRecordKind { @@ -525,6 +527,7 @@ impl WalRecordKind { Self::ExternalActionRequestRecorded => "ExternalActionRequestRecorded", Self::ExternalActionClaimRecorded => "ExternalActionClaimRecorded", Self::ExternalActionSettlementRecorded => "ExternalActionSettlementRecorded", + Self::ExecutableOperationContextRetained => "ExecutableOperationContextRetained", } } @@ -553,7 +556,8 @@ impl WalRecordKind { Self::SchedulerFaultQuarantined | Self::TrustedRuntimeControlRecorded => { WalAppendAuthority::RuntimeControl } - Self::ExecutableOperationPackageInstalled => WalAppendAuthority::RuntimeControl, + Self::ExecutableOperationPackageInstalled + | Self::ExecutableOperationContextRetained => WalAppendAuthority::RuntimeControl, Self::ExecutableOperationExecutionRecorded | Self::ExecutableOperationStateDeltaRecorded => WalAppendAuthority::ExecutionKernel, Self::CausalAnchorFactRecorded | Self::CausalAnchorAdmissionReceiptRecorded => { @@ -614,6 +618,7 @@ impl WalRecordKind { Self::ExternalActionRequestRecorded => 29, Self::ExternalActionClaimRecorded => 30, Self::ExternalActionSettlementRecorded => 31, + Self::ExecutableOperationContextRetained => 32, } } @@ -650,6 +655,7 @@ impl WalRecordKind { 29 => Ok(Self::ExternalActionRequestRecorded), 30 => Ok(Self::ExternalActionClaimRecorded), 31 => Ok(Self::ExternalActionSettlementRecorded), + 32 => Ok(Self::ExecutableOperationContextRetained), _ => Err(WalDecodeError::UnknownEnumCode { enum_name: "WalRecordKind", code, @@ -1586,7 +1592,7 @@ fn validate_writer_epoch_request( if request.started_at_lsn <= final_lsn { return Err(WalStoreError::WriterEpochLsnRegression); } - } else if request.started_at_lsn <= previous_epoch.started_at_lsn { + } else if request.started_at_lsn < previous_epoch.started_at_lsn { return Err(WalStoreError::WriterEpochLsnRegression); } if request.storage_fencing_token == previous_epoch.storage_fencing_token @@ -5763,8 +5769,10 @@ impl FilesystemWalStore { .unwrap_or_default(); let required_started_at_lsn = previous_closure .final_lsn - .or_else(|| previous_epoch.map(|epoch| epoch.started_at_lsn)) .and_then(Lsn::checked_next) + // An empty epoch reserved but never consumed its first LSN. + // Advancing it would leave a gap in the retained frame sequence. + .or_else(|| previous_epoch.map(|epoch| epoch.started_at_lsn)) .unwrap_or(minimum_started_at_lsn); let started_at_lsn = minimum_started_at_lsn.max(required_started_at_lsn); let ordinal = u64::try_from(self.closed_epochs.len()) diff --git a/crates/warp-core/src/echo_operation.rs b/crates/warp-core/src/echo_operation.rs index fd62bcd33..e145dcd58 100644 --- a/crates/warp-core/src/echo_operation.rs +++ b/crates/warp-core/src/echo_operation.rs @@ -29,6 +29,9 @@ use echo_edict_canonical::{ use sha2::{Digest as _, Sha256}; use thiserror::Error; +mod observed; +pub use observed::EchoOperationObservationV1; + use crate::{ attachment::{AtomPayload, AttachmentKey, AttachmentValue}, clock::{GlobalTick, WorldlineTick}, @@ -2779,6 +2782,7 @@ enum AnchoredNodeOperationModeV1 { /// Canonical basis-bearing invocation emitted by a generated client/helper. #[derive(Clone, Debug, PartialEq, Eq)] pub struct EchoOperationInvocationV1 { + observation: Option, package_id: EchoOperationPackageIdV1, operation_coordinate: String, evaluation_basis: EchoOperationEvaluationBasisV1, @@ -2816,6 +2820,7 @@ impl EchoOperationInvocationV1 { }, replacement_bytes, application_input_bytes: None, + observation: None, } } @@ -2841,6 +2846,7 @@ impl EchoOperationInvocationV1 { kind: EchoOperationInvocationKindV1::AnchoredNodeAttachmentCreateIfAbsent, replacement_bytes, application_input_bytes: None, + observation: None, } } @@ -2868,6 +2874,7 @@ impl EchoOperationInvocationV1 { kind: EchoOperationInvocationKindV1::AnchoredNodeAttachmentCreateIfAbsent, replacement_bytes, application_input_bytes: Some(canonical_application_input_bytes), + observation: None, } } @@ -2893,6 +2900,19 @@ impl EchoOperationInvocationV1 { "invocation delegated step budget must be nonzero", )); } + if let Some(observation) = &self.observation { + let mut inner = self.clone(); + inner.observation = None; + return encode_canonical_cbor_v1(&map_value([ + ("schema", text_value(observed::OBSERVED_SCHEMA)), + ( + "invocation", + CanonicalValueV1::Bytes(inner.to_canonical_bytes()?), + ), + ("observation", observation.to_value()), + ])) + .map_err(canonical_error); + } let common = |schema| { [ ( @@ -2971,7 +2991,7 @@ impl EchoOperationInvocationV1 { encode_canonical_cbor_v1(&value).map_err(canonical_error) } - fn from_canonical_bytes(bytes: &[u8]) -> Result { + pub(crate) fn from_canonical_bytes(bytes: &[u8]) -> Result { let value = decode_canonical_cbor_v1(bytes).map_err(canonical_error)?; let schema = match &value { CanonicalValueV1::Map(entries) => entries @@ -2986,6 +3006,27 @@ impl EchoOperationInvocationV1 { .ok_or_else(|| invalid_structure("invocation schema must be text"))?, _ => return Err(invalid_structure("artifact root must be a map")), }; + if schema == observed::OBSERVED_SCHEMA { + let mut fields = exact_text_map(value, &["schema", "invocation", "observation"])?; + let inner_bytes = take_bytes(&mut fields, "invocation")?; + let inner = decode_canonical_cbor_v1(&inner_bytes).map_err(canonical_error)?; + if let CanonicalValueV1::Map(entries) = &inner { + if entries.iter().any(|(key, value)| { + key == &text_value("schema") && value == &text_value(observed::OBSERVED_SCHEMA) + }) { + return Err(invalid_structure("nested observed invocation is forbidden")); + } + } + let mut invocation = Self::from_canonical_bytes(&inner_bytes)?; + invocation.observation = Some(EchoOperationObservationV1::from_value(take_field( + &mut fields, + "observation", + )?)?); + if invocation.to_canonical_bytes()? != bytes { + return Err(invalid_structure("observed invocation is not canonical")); + } + return Ok(invocation); + } let (expected_fields, create_if_absent, projected) = match schema { INVOCATION_SCHEMA => ( &[ @@ -3063,6 +3104,7 @@ impl EchoOperationInvocationV1 { } }; let invocation = Self { + observation: None, package_id: EchoOperationPackageIdV1(take_hash(&mut fields, "package_id")?), operation_coordinate: take_text(&mut fields, "operation_coordinate")?, evaluation_basis: EchoOperationEvaluationBasisV1::from_value(take_field( @@ -3709,6 +3751,10 @@ pub enum EchoOperationObstructionKindV1 { /// The compiler-owned result projection could not produce one bounded /// canonical result from the exact admitted application input. ResultProjectionInvalid, + /// A supplied observation changed despite a current submission basis. + ObservationChanged, + /// An observed resource is no longer available under the bounded profile. + ObservationUnavailable, } /// Retained runtime policy evidence needed to reproduce one obstruction. @@ -4381,6 +4427,8 @@ fn obstruction_kind_from_code( 11 => Ok(EchoOperationObstructionKindV1::ReplacementTooLarge), 12 => Ok(EchoOperationObstructionKindV1::EvaluationAuthorityMismatch), 13 => Ok(EchoOperationObstructionKindV1::ResultProjectionInvalid), + 14 => Ok(EchoOperationObstructionKindV1::ObservationChanged), + 15 => Ok(EchoOperationObstructionKindV1::ObservationUnavailable), _ => Err(invalid_structure( "unknown executable-operation obstruction kind", )), @@ -4615,6 +4663,17 @@ pub(crate) fn prepare_operation_v1( let node = admitted.invocation.node; let mut actual_footprint = Footprint::default(); let mut budget_meter = EchoOperationBudgetMeterV1::new(admitted.invocation.delegated_budget); + if let Some(observation) = &admitted.invocation.observation { + if let Err(kind) = observation.validate_at_execution( + state, + current_basis, + &mut actual_footprint, + &mut budget_meter, + ) { + return obstruction(kind); + } + } + let observed_footprint = actual_footprint.clone(); let descent_stack = match operation_descent_stack_with_portal_reads(state, node.warp_id, |portal| { if !budget_meter.charge(1, 32, 0) { @@ -4629,7 +4688,7 @@ pub(crate) fn prepare_operation_v1( let Some(store) = state.store(&node.warp_id) else { return obstruction(EchoOperationObstructionKindV1::NodeMissing); }; - let declared_footprint = match mode { + let mut declared_footprint = match mode { AnchoredNodeOperationModeV1::CompareAndSet { .. } => { anchored_node_compare_and_set_footprint(node, &descent_stack) } @@ -4637,6 +4696,7 @@ pub(crate) fn prepare_operation_v1( anchored_node_create_if_absent_footprint(node, &descent_stack) } }; + declared_footprint.union_assign(&observed_footprint); if !budget_meter.charge(1, 32, 0) { return obstruction(EchoOperationObstructionKindV1::BudgetExceeded); } @@ -4739,6 +4799,14 @@ pub(crate) fn prepare_operation_v1( } let mut in_slots = vec![SlotId::Node(node), SlotId::Attachment(slot)]; in_slots.extend(descent_stack.iter().copied().map(SlotId::Attachment)); + in_slots.extend(observed_footprint.n_read.iter().copied().map(SlotId::Node)); + in_slots.extend( + observed_footprint + .a_read + .iter() + .copied() + .map(SlotId::Attachment), + ); let patch = WarpTickPatchV1::new( policy_id, installed.installed_operation_id.as_hash(), @@ -6767,6 +6835,8 @@ fn obstruction_kind_code(kind: EchoOperationObstructionKindV1) -> u8 { EchoOperationObstructionKindV1::ReplacementTooLarge => 11, EchoOperationObstructionKindV1::EvaluationAuthorityMismatch => 12, EchoOperationObstructionKindV1::ResultProjectionInvalid => 13, + EchoOperationObstructionKindV1::ObservationChanged => 14, + EchoOperationObstructionKindV1::ObservationUnavailable => 15, } } @@ -7311,7 +7381,7 @@ mod tests { } } - fn projected_create_fixture( + pub(super) fn projected_create_fixture( max_output_bytes: u64, ) -> ( InstalledEchoOperationV1, diff --git a/crates/warp-core/src/echo_operation/observed.rs b/crates/warp-core/src/echo_operation/observed.rs new file mode 100644 index 000000000..37d7bedf7 --- /dev/null +++ b/crates/warp-core/src/echo_operation/observed.rs @@ -0,0 +1,305 @@ +// SPDX-License-Identifier: Apache-2.0 +// © James Ross Ω FLYING•ROBOTS +//! Bounded observation preconditions for executable actions. +use super::*; + +/// A bounded reading supplied before an executable operation is proposed. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct EchoOperationObservationV1 { + basis: EchoOperationEvaluationBasisV1, + reads: Vec<(NodeKey, Vec)>, +} + +pub(super) const OBSERVED_SCHEMA: &str = "echo.operation-invocation.observed/v1"; +const MAX_READS: usize = 16; +const MAX_BYTES: usize = 4096; + +fn slot_bytes( + state: &WorldlineState, + node: NodeKey, +) -> Result, EchoOperationArtifactErrorV1> { + let store = state + .store(&node.warp_id) + .ok_or_else(|| invalid_structure("observation warp unavailable"))?; + let node_type = store + .node(&node.local_id) + .map_or(CanonicalValueV1::Null, |record| hash_value(record.ty.0)); + let alpha = match store.node_attachment(&node.local_id) { + None => CanonicalValueV1::Null, + Some(AttachmentValue::Atom(atom)) => map_value([ + ("type", hash_value(atom.type_id.0)), + ("bytes", CanonicalValueV1::Bytes(atom.bytes.to_vec())), + ]), + Some(_) => { + return Err(invalid_structure( + "observation requires atom or absent alpha", + )) + } + }; + encode_canonical_cbor_v1(&map_value([("node_type", node_type), ("alpha", alpha)])) + .map_err(canonical_error) +} + +impl EchoOperationObservationV1 { + /// Encodes retained observation evidence for the runtime-owned context log. + pub(crate) fn encode(&self) -> Result, EchoOperationArtifactErrorV1> { + encode_canonical_cbor_v1(&self.to_value()).map_err(canonical_error) + } + + pub(crate) fn decode(bytes: &[u8]) -> Result { + Self::from_value(decode_canonical_cbor_v1(bytes).map_err(canonical_error)?) + } + + /// Derives changed support from this exact retained reading. + pub(crate) fn changed_nodes( + &self, + state: &WorldlineState, + ) -> Result, EchoOperationArtifactErrorV1> { + self.reads + .iter() + .filter_map(|(node, expected)| match slot_bytes(state, *node) { + Ok(actual) if actual == *expected => None, + Ok(_) => Some(Ok(*node)), + Err(error) => Some(Err(error)), + }) + .collect() + } + /// Captures the named node/alpha slots at an explicit observation basis. + pub fn capture( + state: &WorldlineState, + basis: EchoOperationEvaluationBasisV1, + nodes: &[NodeKey], + ) -> Result { + if state.state_root() != basis.state_root() { + return Err(invalid_structure( + "observation basis does not bind supplied state", + )); + } + let nodes = nodes.iter().copied().collect::>(); + if nodes.is_empty() || nodes.len() > MAX_READS { + return Err(invalid_structure("observation needs one to sixteen slots")); + } + let reads = nodes + .into_iter() + .map(|node| Ok((node, slot_bytes(state, node)?))) + .collect::, EchoOperationArtifactErrorV1>>()?; + let observation = Self { basis, reads }; + if observation + .reads + .iter() + .map(|(_, bytes)| bytes.len()) + .sum::() + > MAX_BYTES + { + return Err(invalid_structure("observation exceeds retained byte bound")); + } + Ok(observation) + } + + /// Exact basis of the supplied reading, independent of a later submission. + #[must_use] + pub const fn basis(&self) -> EchoOperationEvaluationBasisV1 { + self.basis + } + + /// Canonical node/alpha values supplied by this reading. + pub fn readings(&self) -> impl Iterator { + self.reads + .iter() + .map(|(node, value)| (*node, value.as_slice())) + } + + pub(super) fn to_value(&self) -> CanonicalValueV1 { + map_value([ + ("basis", self.basis.to_value()), + ( + "reads", + CanonicalValueV1::Array( + self.reads + .iter() + .map(|(node, bytes)| { + map_value([ + ("warp_id", hash_value(node.warp_id.0)), + ("node_id", hash_value(node.local_id.0)), + ("value", CanonicalValueV1::Bytes(bytes.clone())), + ]) + }) + .collect(), + ), + ), + ]) + } + + pub(super) fn from_value( + value: CanonicalValueV1, + ) -> Result { + let mut fields = exact_text_map(value, &["basis", "reads"])?; + let basis = EchoOperationEvaluationBasisV1::from_value(take_field(&mut fields, "basis")?)?; + let CanonicalValueV1::Array(values) = take_field(&mut fields, "reads")? else { + return Err(invalid_structure("observation reads must be an array")); + }; + if values.is_empty() || values.len() > MAX_READS { + return Err(invalid_structure("observation needs one to sixteen slots")); + } + let mut reads = Vec::new(); + let mut total = 0; + for value in values { + let mut fields = exact_text_map(value, &["warp_id", "node_id", "value"])?; + let node = NodeKey { + warp_id: crate::WarpId(take_hash(&mut fields, "warp_id")?), + local_id: crate::NodeId(take_hash(&mut fields, "node_id")?), + }; + let bytes = take_bytes(&mut fields, "value")?; + total += bytes.len(); + if total > MAX_BYTES || reads.last().is_some_and(|(previous, _)| *previous >= node) { + return Err(invalid_structure( + "observation must be bounded and strictly sorted", + )); + } + decode_canonical_cbor_v1(&bytes).map_err(canonical_error)?; + reads.push((node, bytes)); + } + Ok(Self { basis, reads }) + } + + pub(super) fn validate_at_execution( + &self, + state: &WorldlineState, + submission: EchoOperationEvaluationBasisV1, + footprint: &mut Footprint, + meter: &mut EchoOperationBudgetMeterV1, + ) -> Result<(), EchoOperationObstructionKindV1> { + if self.basis.writer_head().worldline_id != submission.writer_head().worldline_id + || self.basis.worldline_tick() > submission.worldline_tick() + { + return Err(EchoOperationObstructionKindV1::ObservationChanged); + } + for (node, expected) in &self.reads { + let portals = + operation_descent_stack_with_portal_reads(state, node.warp_id, |portal| { + footprint.a_read.insert(portal); + meter.charge(1, 32, 0) + })?; + let _ = portals; + let actual = slot_bytes(state, *node) + .map_err(|_| EchoOperationObstructionKindV1::ObservationUnavailable)?; + if !meter.charge(2, 64 + actual.len() as u64, 0) { + return Err(EchoOperationObstructionKindV1::BudgetExceeded); + } + footprint.n_read.insert(*node); + footprint.a_read.insert(AttachmentKey::node_alpha(*node)); + if actual != *expected { + return Err(EchoOperationObstructionKindV1::ObservationChanged); + } + } + Ok(()) + } +} + +impl EchoOperationInvocationV1 { + pub(crate) fn has_observation(&self, observation: &EchoOperationObservationV1) -> bool { + self.observation.as_ref() == Some(observation) + } + pub(crate) fn observed_semantic_identity( + &self, + observation: &EchoOperationObservationV1, + ) -> Result { + let mut normalized = self.clone().with_observation(observation.clone()); + // Transport rebasing does not change the originally proposed meaning. + normalized.evaluation_basis = observation.basis; + Ok(*blake3::hash(&normalized.to_canonical_bytes()?).as_bytes()) + } + /// Binds the operation to inputs observed before its submission basis. + #[must_use] + pub fn with_observation(mut self, observation: EchoOperationObservationV1) -> Self { + self.observation = Some(observation); + self + } +} + +#[cfg(test)] +#[allow(clippy::expect_used, clippy::panic)] +mod tests { + use super::*; + + fn exercise(changed_relevant: bool) -> (EchoOperationPreparationV1, NodeKey) { + let (installed, mut state, old_basis, policy, mut invocation, _) = + super::super::tests::projected_create_fixture(1_024); + let observed = *state.root(); + let observation = EchoOperationObservationV1::capture(&state, old_basis, &[observed]) + .expect("bounded observation"); + let changed = if changed_relevant { + observed + } else { + NodeKey { + warp_id: observed.warp_id, + local_id: crate::make_node_id("unrelated"), + } + }; + let store = state.warp_state.store_mut(&changed.warp_id).expect("store"); + store.insert_node( + changed.local_id, + NodeRecord { + ty: crate::make_type_id("changed"), + }, + ); + let basis = EchoOperationEvaluationBasisV1::new( + old_basis.writer_head(), + WorldlineTick::from_raw(1), + None, + state.state_root(), + old_basis.commit_id, + old_basis.application_basis, + ); + invocation.evaluation_basis = basis; + invocation.delegated_budget = EchoOperationBudgetV1::new(8, 1_024, 1_024); + let invocation = invocation.with_observation(observation); + let authority = EchoOperationEvaluationAuthorityV1::new(); + let admitted = admit_invocation_v1( + Some(&installed), + policy, + &invocation.to_canonical_bytes().expect("encode"), + basis, + &state, + authority.clone(), + ) + .expect("current submission basis admits"); + ( + prepare_operation_v1( + Some(&installed), + admitted, + basis, + &state, + crate::POLICY_ID_NO_POLICY_V0, + &authority, + ), + observed, + ) + } + + #[test] + fn fresh_submission_cannot_erase_changed_observation() { + let (outcome, _) = exercise(true); + assert!( + matches!(outcome, EchoOperationPreparationV1::Obstructed(ref obstruction) if obstruction.kind() == EchoOperationObstructionKindV1::ObservationChanged), + "a fresh submission basis must not launder an old observation" + ); + } + + #[test] + fn unrelated_movement_preserves_support_and_adds_native_reads() { + let (outcome, observed) = exercise(false); + let EchoOperationPreparationV1::Prepared(prepared) = outcome else { + panic!("unrelated movement must not invalidate the observation"); + }; + assert!(prepared + .actual_footprint() + .n_read + .iter() + .any(|node| *node == observed)); + assert!(prepared + .patch() + .in_slots() + .contains(&SlotId::Node(observed))); + } +} diff --git a/crates/warp-core/src/lib.rs b/crates/warp-core/src/lib.rs index ade5e6b74..e76813e52 100644 --- a/crates/warp-core/src/lib.rs +++ b/crates/warp-core/src/lib.rs @@ -257,13 +257,13 @@ pub use echo_operation::{ EchoOperationInstallationErrorV1, EchoOperationInvocationAdmissionErrorKindV1, EchoOperationInvocationAdmissionErrorV1, EchoOperationInvocationAdmissionIdV1, EchoOperationInvocationAdmissionPolicyV1, EchoOperationInvocationIdV1, - EchoOperationInvocationV1, EchoOperationObstructionIdV1, EchoOperationObstructionKindV1, - EchoOperationObstructionV1, EchoOperationPackageAdmissionIdV1, EchoOperationPackageIdV1, - EchoOperationPreparationV1, EchoOperationPrivateEvaluationIdV1, EchoOperationProgramIdV1, - EchoOperationProgramV1, EchoOperationReceiptV1, EchoOperationResultIdV1, - EchoOperationSemanticClosureV1, EchoOperationTerminalPostureV1, ExecutableOperationPackageV1, - InstalledEchoOperationIdV1, InstalledEchoOperationV1, PreparedEchoOperationIdV1, - PreparedEchoOperationV1, ACTION_BATCH_CANDIDATE_LIMIT_V1, + EchoOperationInvocationV1, EchoOperationObservationV1, EchoOperationObstructionIdV1, + EchoOperationObstructionKindV1, EchoOperationObstructionV1, EchoOperationPackageAdmissionIdV1, + EchoOperationPackageIdV1, EchoOperationPreparationV1, EchoOperationPrivateEvaluationIdV1, + EchoOperationProgramIdV1, EchoOperationProgramV1, EchoOperationReceiptV1, + EchoOperationResultIdV1, EchoOperationSemanticClosureV1, EchoOperationTerminalPostureV1, + ExecutableOperationPackageV1, InstalledEchoOperationIdV1, InstalledEchoOperationV1, + PreparedEchoOperationIdV1, PreparedEchoOperationV1, ACTION_BATCH_CANDIDATE_LIMIT_V1, }; pub use edict_target_ir::{ accept_edict_echo_target_ir, execute_accepted_edict_echo_target_ir, AcceptedEdictEchoTargetIr, @@ -484,10 +484,11 @@ pub use tick_patch::{ }; #[cfg(all(feature = "native_rule_bootstrap", feature = "trusted_runtime"))] pub use trusted_runtime_host::{ - EvidenceCatalogPosture, RuntimeWalActivationGap, TrustedRuntimeApp, TrustedRuntimeHost, - TrustedRuntimeHostError, TrustedRuntimeHostParts, TrustedRuntimeHostRunReport, - TrustedRuntimeWal, TrustedRuntimeWalConfig, TrustedRuntimeWalError, TrustedRuntimeWalRecovery, - TrustedRuntimeWalStoreKind, WitnessedCausalAnchorAdmission, + EchoOperationContextErrorV1, EvidenceCatalogPosture, RuntimeWalActivationGap, + TrustedRuntimeApp, TrustedRuntimeHost, TrustedRuntimeHostError, TrustedRuntimeHostParts, + TrustedRuntimeHostRunReport, TrustedRuntimeWal, TrustedRuntimeWalConfig, + TrustedRuntimeWalError, TrustedRuntimeWalRecovery, TrustedRuntimeWalStoreKind, + WitnessedCausalAnchorAdmission, }; pub use tx::TxId; pub use warp_state::{WarpInstance, WarpState}; diff --git a/crates/warp-core/src/trusted_runtime_host.rs b/crates/warp-core/src/trusted_runtime_host.rs index 3acc60ce1..93c66ad0e 100644 --- a/crates/warp-core/src/trusted_runtime_host.rs +++ b/crates/warp-core/src/trusted_runtime_host.rs @@ -14,6 +14,9 @@ use std::{ use thiserror::Error; +mod observed_context; +pub use observed_context::EchoOperationContextErrorV1; + use crate::causal_anchor::prepare_causal_anchor_admission; use crate::{ diff --git a/crates/warp-core/src/trusted_runtime_host/observed_context.rs b/crates/warp-core/src/trusted_runtime_host/observed_context.rs new file mode 100644 index 000000000..165c7f5f6 --- /dev/null +++ b/crates/warp-core/src/trusted_runtime_host/observed_context.rs @@ -0,0 +1,296 @@ +// SPDX-License-Identifier: Apache-2.0 +// © James Ross Ω FLYING•ROBOTS +//! Runtime-owned immutable observation and logical-request bindings. +use super::{ + BTreeMap, Error, Hash, TrustedRuntimeHost, TrustedRuntimeWal, WalAppendAuthority, + WalTransactionId, WalTransactionKind, +}; +use crate::{EchoOperationInvocationV1, EchoOperationObservationV1, NodeKey, WriterHeadKey}; +use echo_edict_canonical::{ + decode_canonical_cbor_v1, encode_canonical_cbor_v1, CanonicalValueV1 as V, +}; + +/// Failure to retain or resolve an immutable operation context. +#[derive(Debug, Error)] +#[error("operation context: {0}")] +pub struct EchoOperationContextErrorV1(String); + +fn error(value: impl std::fmt::Display) -> EchoOperationContextErrorV1 { + EchoOperationContextErrorV1(value.to_string()) +} + +type Result = std::result::Result; +const SCHEMA: &str = "echo.operation-context/v1"; + +fn field_text(value: &V) -> Result<&str> { + if let V::Text(value) = value { + Ok(value) + } else { + Err(error("invalid context text")) + } +} +fn field_bytes(value: &V) -> Result<&[u8]> { + if let V::Bytes(value) = value { + Ok(value) + } else { + Err(error("invalid context bytes")) + } +} +fn check_id(id: &str) -> Result<()> { + if id.is_empty() || id.len() > 128 { + return Err(error("identity must contain 1..128 bytes")); + } + Ok(()) +} + +#[derive(Default)] +struct Contexts { + observations: BTreeMap, + requests: BTreeMap)>, +} + +impl TrustedRuntimeWal { + fn operation_contexts(&self) -> Result { + let report = self.store.recover_read_only().map_err(error)?; + let mut contexts = Contexts::default(); + for transaction in report.transactions { + for frame in transaction.frames { + if frame.header.record_kind + != crate::causal_wal::WalRecordKind::ExecutableOperationContextRetained + { + continue; + } + let V::Array(fields) = + decode_canonical_cbor_v1(&frame.payload.canonical_bytes).map_err(error)? + else { + return Err(error("invalid context record")); + }; + if fields.len() < 4 || field_text(&fields[0])? != SCHEMA { + return Err(error("invalid context schema")); + } + let id = field_text(&fields[2])?.to_owned(); + check_id(&id)?; + match field_text(&fields[1])? { + "observation" if fields.len() == 4 => { + let observation = + EchoOperationObservationV1::decode(field_bytes(&fields[3])?) + .map_err(error)?; + if contexts.observations.insert(id, observation).is_some() { + return Err(error("duplicate observation identity")); + } + } + "request" if fields.len() == 6 => { + let attempt = field_text(&fields[3])?.to_owned(); + let digest: Hash = field_bytes(&fields[4])?.try_into().map_err(error)?; + let invocation_bytes = field_bytes(&fields[5])?.to_vec(); + let invocation = + EchoOperationInvocationV1::from_canonical_bytes(&invocation_bytes) + .map_err(error)?; + let observation = contexts + .observations + .get(&attempt) + .ok_or_else(|| error("request observation unavailable"))?; + if !invocation.has_observation(observation) + || invocation + .observed_semantic_identity(observation) + .map_err(error)? + != digest + { + return Err(error("request identity mismatch")); + } + if contexts + .requests + .insert(id, (attempt, digest, invocation_bytes)) + .is_some() + { + return Err(error("duplicate request identity")); + } + } + _ => return Err(error("invalid context record kind")), + } + } + } + Ok(contexts) + } + + fn retain_operation_context(&mut self, fields: Vec) -> Result<()> { + // A previous flush may have failed after bytes reached storage. Rebuild + // the owned writer cursor and truncate an uncommitted tail before append. + self.refresh_cursor_from_store_for_writer().map_err(error)?; + let bytes = encode_canonical_cbor_v1(&V::Array(fields)).map_err(error)?; + let transaction_id = WalTransactionId::from_hash(*blake3::hash(&bytes).as_bytes()); + let mut builder = self.builder( + WalTransactionKind::RuntimePosture, + WalAppendAuthority::RuntimeControl, + transaction_id, + ); + builder + .push_record( + crate::causal_wal::WalRecordKind::ExecutableOperationContextRetained, + bytes, + ) + .map_err(error)?; + self.append_transaction(builder.commit(Vec::new()).map_err(error)?) + .map_err(error)?; + Ok(()) + } +} + +impl TrustedRuntimeHost { + /// Retains a bounded observation before supplying it to an attempt. + /// + /// The trusted host controls the aperture. Reusing an attempt cannot replace + /// its observations. Missing WAL material obstructs instead of recapturing. + pub fn retain_echo_operation_observation_v1( + &mut self, + attempt: &str, + head: WriterHeadKey, + nodes: &[NodeKey], + ) -> Result { + check_id(attempt)?; + let contexts = self + .runtime_wal + .as_ref() + .ok_or_else(|| error("durable WAL required"))? + .operation_contexts()?; + if contexts.observations.contains_key(attempt) { + return Err(error("attempt observation is immutable")); + } + if contexts.observations.len() >= 1024 { + return Err(error("context admission limit reached")); + } + let first = nodes + .first() + .ok_or_else(|| error("observation needs an aperture"))?; + let application_basis = + crate::echo_operation_anchored_node_absent_application_basis_v1(*first); + let basis = self + .echo_operation_evaluation_basis_v1(head, application_basis) + .map_err(error)?; + let state = self + .runtime + .worldlines() + .get(&head.worldline_id) + .ok_or_else(|| error("worldline unavailable"))? + .state(); + let observation = + EchoOperationObservationV1::capture(state, basis, nodes).map_err(error)?; + self.runtime_wal + .as_mut() + .ok_or_else(|| error("durable WAL required"))? + .retain_operation_context(vec![ + V::Text(SCHEMA.into()), + V::Text("observation".into()), + V::Text(attempt.into()), + V::Bytes(observation.encode().map_err(error)?), + ])?; + Ok(observation) + } + + /// Resolves exactly the retained reading, without refreshing its basis. + pub fn echo_operation_observation_v1( + &self, + attempt: &str, + ) -> Result { + self.runtime_wal + .as_ref() + .ok_or_else(|| error("durable WAL required"))? + .operation_contexts()? + .observations + .remove(attempt) + .ok_or_else(|| error("observation unavailable; re-observation requires a new attempt")) + } + + /// Returns only changed resources in this attempt's retained aperture. + pub fn echo_operation_observation_changes_v1(&self, attempt: &str) -> Result> { + let observation = self.echo_operation_observation_v1(attempt)?; + let state = self + .runtime + .worldlines() + .get(&observation.basis().writer_head().worldline_id) + .ok_or_else(|| error("worldline unavailable"))? + .state(); + observation.changed_nodes(state).map_err(error) + } + + /// Retained commits which wrote the changed support after this reading. + /// Unrelated commits and their payloads are excluded from the result. + pub fn echo_operation_observation_change_commits_v1(&self, attempt: &str) -> Result> { + let observation = self.echo_operation_observation_v1(attempt)?; + let changed = self.echo_operation_observation_changes_v1(attempt)?; + let history = self + .runtime_wal + .as_ref() + .ok_or_else(|| error("durable WAL required"))? + .recover_read_only() + .map_err(error)?; + Ok(history + .provenance_entries + .iter() + .filter(|entry| { + entry.worldline_id == observation.basis().writer_head().worldline_id + && entry.worldline_tick >= observation.basis().worldline_tick() + && entry.patch.as_ref().is_some_and(|patch| { + changed.iter().any(|node| { + patch.out_slots.contains(&crate::SlotId::Node(*node)) + || patch.out_slots.contains(&crate::SlotId::Attachment( + crate::AttachmentKey::node_alpha(*node), + )) + }) + }) + }) + .map(|entry| entry.expected.commit_hash) + .collect()) + } + + /// Binds a logical request before ingress acceptance. Exact retry returns + /// the original invocation, including its original submission basis. + /// Changed semantic input under the same identity is refused. + pub fn bind_echo_operation_request_v1( + &mut self, + request: &str, + attempt: &str, + invocation: EchoOperationInvocationV1, + ) -> Result> { + check_id(request)?; + let contexts = self + .runtime_wal + .as_ref() + .ok_or_else(|| error("durable WAL required"))? + .operation_contexts()?; + let observation = contexts + .observations + .get(attempt) + .ok_or_else(|| error("observation unavailable"))?; + let semantic = invocation + .observed_semantic_identity(observation) + .map_err(error)?; + if let Some((original_attempt, original_semantic, original_bytes)) = + contexts.requests.get(request) + { + if original_attempt != attempt || original_semantic != &semantic { + return Err(error("logical request identity reused with changed input")); + } + return Ok(original_bytes.clone()); + } + if contexts.requests.len() >= 4096 { + return Err(error("request admission limit reached")); + } + let bytes = invocation + .with_observation(observation.clone()) + .to_canonical_bytes() + .map_err(error)?; + self.runtime_wal + .as_mut() + .ok_or_else(|| error("durable WAL required"))? + .retain_operation_context(vec![ + V::Text(SCHEMA.into()), + V::Text("request".into()), + V::Text(request.into()), + V::Text(attempt.into()), + V::Bytes(semantic.to_vec()), + V::Bytes(bytes.clone()), + ])?; + Ok(bytes) + } +} diff --git a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs index c2bdc9591..52032937d 100644 --- a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs +++ b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs @@ -1112,6 +1112,98 @@ fn filesystem_runtime_wal_ack_commits_strict_filesystem_durability() { ); } +#[test] +fn filesystem_runtime_wal_empty_reopen_does_not_consume_a_log_position() { + let root = temp_runtime_wal_dir("empty-reopen-continuation"); + let reopen = || { + let (runtime, first, second) = runtime_pair(); + let mut host = TrustedRuntimeHost::new(runtime, empty_engine()).expect("host"); + host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + .expect("reopen"); + (host, first, second) + }; + let (mut first, lane, _) = reopen(); + first + .app() + .submit_intent_with_runtime_wal_ack(eint_envelope(lane)) + .expect("first append"); + drop(first); + let (empty, _, _) = reopen(); + drop(empty); + let (mut last, _, lane) = reopen(); + last.app() + .submit_intent_with_runtime_wal_ack(eint_envelope(lane)) + .expect("second append"); + let recovered = last + .runtime_wal() + .expect("WAL") + .recover_read_only() + .expect("contiguous history"); + assert_eq!(recovered.certificate.committed_transactions_replayed, 2); + let commits = last.runtime_wal().expect("WAL").commits(); + assert_eq!(commits[1].first_lsn, Lsn::from_raw(3)); +} + +#[test] +fn observed_operation_contexts_recover_without_replacing_original_inputs() { + let root = temp_runtime_wal_dir("observed-request-context"); + let reopen = || { + let (runtime, lane, _) = runtime_pair(); + let node = *runtime + .worldlines() + .get(&lane) + .expect("lane") + .state() + .root(); + let head = WriterHeadKey { + worldline_id: lane, + head_id: make_head_id("default-a"), + }; + let mut host = TrustedRuntimeHost::new(runtime, empty_engine()).expect("host"); + host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + .expect("WAL"); + (host, head, node) + }; + let (mut host, head, node) = reopen(); + let observation = host + .retain_echo_operation_observation_v1("attempt", head, &[node]) + .expect("retain reading"); + let invocation = |value: &[u8]| { + warp_core::EchoOperationInvocationV1::anchored_node_attachment_create_if_absent_with_application_input( + warp_core::echo_operation_package_id_v1(b"context-only-test"), "test.context@1.create", + observation.basis(), [1;32], warp_core::EchoOperationBudgetV1::new(16, 1024, 320), + node, value.to_vec(), vec![0xa0], + ) + }; + // Binding is durable before ingress. This test does not claim that the + // deliberately uninstalled package can be admitted or executed. + let original = host + .bind_echo_operation_request_v1("request", "attempt", invocation(b"old")) + .expect("bind request"); + drop(host); + let (mut recovered, head, node) = reopen(); + assert_eq!( + recovered + .echo_operation_observation_v1("attempt") + .expect("original reading"), + observation + ); + assert_eq!( + recovered + .bind_echo_operation_request_v1("request", "attempt", invocation(b"old")) + .expect("retry"), + original + ); + assert!(recovered + .bind_echo_operation_request_v1("request", "attempt", invocation(b"changed")) + .is_err()); + assert!(recovered + .retain_echo_operation_observation_v1("attempt", head, &[node]) + .is_err()); + assert!(recovered.echo_operation_observation_v1("missing").is_err()); + assert_eq!(recovered.runtime_wal().expect("WAL").commits().len(), 2); +} + #[test] fn filesystem_runtime_wal_ack_recovery_reports_uncommitted_tail_from_root() { let wal_root = temp_runtime_wal_dir("tail-report"); diff --git a/docs/architecture/application-contract-hosting.md b/docs/architecture/application-contract-hosting.md index 0c8939158..aaf5fe458 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -180,6 +180,52 @@ the exact program, after which Echo independently admits each invocation. ## External Edict Provider Artifacts +### Observation-bound native host sessions + +The bounded `cargo xtask run-edict-operation --serve` driver continues one +worldline through multiple compiled create-if-absent operations and reopens its +native WAL. The bootstrap `basis` string selects that worldline once. Each +submission gets a current runtime evaluation basis; it does not replace the +observation basis retained for its attempt. + +The trusted host captures bounded node/atom readings before returning them to a +caller. Immutable observation and logical-request bindings are retained as +`ExecutableOperationContextRetained` WAL records. Request binding precedes +ordinary durable ingress acceptance. Retrying the same semantic request returns +its original canonical invocation and resolves its original native disposition; +different input, package, operation, grant, or observation attempt under that +identity is refused. A crash between binding and ingress can resume acceptance +of that exact invocation. Transport rebasing does not alter its meaning. + +Observation preconditions execute in Echo's operation preparation, against the +state the scheduler will commit from. Their resource reads enter the native +footprint and patch input slots. The operation's admitted budget covers those +reads. Relevant value changes produce `ObservationChanged`; unavailable support +produces an obstruction. Unrelated movement does not invalidate the reading. +The original invocation remains the replay input, including its observation. + +Change discovery is a reading over retained observations and native patches. +It returns changed aperture keys and the retained commits that wrote them, +without persisting a second notification log. It survives loss of the driver's +caches. Unknown observation identities obstruct rather than silently recapture. + +This is a trusted, locally scripted host profile, not an authenticated agent +service or the complete public optics boundary. It covers bounded atomic node +and attachment readings in the caller's granted aperture, not model-internal +influence, arbitrary subtree observations, or speculative strand settlement. +The driver accepts at most 1,024 observations and 4,096 request bindings per WAL; +each observation contains at most 16 nodes and 4,096 retained value bytes. +It currently rebuilds context indexes from the WAL. Production indexing, +retention policy, and authenticated aperture delegation remain separate work. + +Disconnected create-if-absent cells can leave the root-reachable state hash +unchanged. Recovery evidence must therefore include native commit identities +and the recovered cells, not just that root hash. Reopening an empty writer +epoch reuses its unconsumed first LSN under a fresh fenced epoch; it must not +introduce a gap into the retained frame sequence. + +### Provider artifact boundary + Echo also owns the runtime-specific semantics supplied to Edict's generic external provider host. That pipeline has a separate source and output boundary: diff --git a/xtask/src/main.rs b/xtask/src/main.rs index f343b9cd6..d4066bb7d 100644 --- a/xtask/src/main.rs +++ b/xtask/src/main.rs @@ -159,6 +159,9 @@ struct RuntimeCounterDiagnosticArgs { #[derive(Args)] struct RunEdictOperationArgs { + /// Serve bounded JSON requests against one persistent worldline. + #[arg(long)] + serve: bool, /// Exact compiler-produced executable-operation package. #[arg(long)] package: PathBuf, @@ -177,7 +180,7 @@ struct RunEdictOperationArgs { /// Typed operation input encoded as JSON. #[arg(long)] input: PathBuf, - /// Empty directory owned by this run's strict filesystem WAL. + /// Strict filesystem WAL directory; must be empty unless --serve reopens it. #[arg(long)] wal_dir: PathBuf, /// Emit the witness report as JSON. @@ -482,7 +485,7 @@ fn main() -> Result<()> { } fn run_edict_operation(args: RunEdictOperationArgs) -> Result<()> { - let report = run_edict_operation::run(run_edict_operation::RunEdictOperationConfig { + let config = run_edict_operation::RunEdictOperationConfig { package: args.package, verification_report: args.verification_report, lawpack_manifest: args.lawpack_manifest, @@ -490,7 +493,11 @@ fn run_edict_operation(args: RunEdictOperationArgs) -> Result<()> { target_configuration: args.target_configuration, input: args.input, wal_dir: args.wal_dir, - })?; + }; + if args.serve { + return run_edict_operation::serve(config); + } + let report = run_edict_operation::run(config)?; if args.json { println!("{}", serde_json::to_string_pretty(&report)?); diff --git a/xtask/src/run_edict_operation.rs b/xtask/src/run_edict_operation.rs index 1044b0d9f..6b8d71754 100644 --- a/xtask/src/run_edict_operation.rs +++ b/xtask/src/run_edict_operation.rs @@ -193,6 +193,9 @@ struct HostFixture { node: NodeKey, } +mod session; +pub use session::serve; + /// Runs the exact package and returns its durable singleton scheduler witness. pub fn run(config: RunEdictOperationConfig) -> Result { let package_bytes = read_bounded(&config.package, MAX_ARTIFACT_BYTES, "package")?; diff --git a/xtask/src/run_edict_operation/session.rs b/xtask/src/run_edict_operation/session.rs new file mode 100644 index 000000000..92bb93bdc --- /dev/null +++ b/xtask/src/run_edict_operation/session.rs @@ -0,0 +1,317 @@ +// SPDX-License-Identifier: Apache-2.0 +// © James Ross Ω FLYING•ROBOTS +//! Bounded JSONL driver for the native executable-operation host. +use super::{ + application_result_report, bail, build_host, current_state, decode_canonical_cbor_v1, + domain_hash, echo_operation_action_envelope_v1, + echo_operation_anchored_node_creation_application_basis_v1, echo_operation_package_id_v1, fs, + install_package, parse_input, parse_package, read_bounded, validate_closure, + validate_package_configuration, validate_verification_report, Context, Digest, + EchoOperationActionOutcomeV1, EchoOperationAnchoredNodeOccupancyV1, + EchoOperationInvocationAdmissionPolicyV1, EchoOperationInvocationV1, HostFixture, + IngressTarget, NodeId, NodeKey, PackageMetadata, Result, RunEdictOperationConfig, Sha256, + TargetConfiguration, TrustedRuntimeWalConfig, MAX_ARTIFACT_BYTES, MAX_INPUT_BYTES, +}; +use serde_json::{json, Value}; +use std::io::{BufRead, Read, Write}; + +struct Session { + fixture: HostFixture, + package: PackageMetadata, + package_id: warp_core::EchoOperationPackageIdV1, + grant: [u8; 32], + configuration: TargetConfiguration, + basis: String, +} + +pub fn serve(config: RunEdictOperationConfig) -> Result<()> { + let package_bytes = read_bounded(&config.package, MAX_ARTIFACT_BYTES, "package")?; + let package_value = decode_canonical_cbor_v1(&package_bytes)?; + let package = parse_package(&package_value)?; + validate_verification_report( + &read_bounded( + &config.verification_report, + MAX_ARTIFACT_BYTES, + "verification report", + )?, + &package_value, + &package.operation_coordinate, + package.target_ir_identity, + package.result_projection_identity, + )?; + let configuration = validate_closure( + &read_bounded(&config.lawpack_manifest, MAX_ARTIFACT_BYTES, "manifest")?, + &read_bounded(&config.lawpack_adapter, MAX_ARTIFACT_BYTES, "adapter")?, + &read_bounded( + &config.target_configuration, + MAX_ARTIFACT_BYTES, + "configuration", + )?, + &package.lawpack_coordinate, + package.lawpack_identity, + package.target_intrinsic, + )?; + validate_package_configuration(&package, &configuration)?; + let input = parse_input( + &read_bounded(&config.input, MAX_INPUT_BYTES, "bootstrap")?, + &configuration, + )?; + fs::create_dir_all(&config.wal_dir)?; + let mut fixture = build_host(&input.basis, &input.key, false)?; + fixture + .host + .enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&config.wal_dir))?; + let package_id = echo_operation_package_id_v1(&package_bytes); + if fixture + .host + .engine() + .installed_echo_operation_package_v1(package_id) + .is_none() + { + install_package(&mut fixture.host, &package, package_id, package_bytes)?; + } + let grant = domain_hash( + b"echo:edict-operation-runner-authority-grant:v1\0", + &package_id.as_hash(), + ); + fixture + .host + .install_echo_operation_action_admission_policy_v1( + EchoOperationInvocationAdmissionPolicyV1::new( + package.authority_profile_identity, + grant, + package.budget, + ), + ); + let mut session = Session { + fixture, + package, + package_id, + grant, + configuration, + basis: input.basis, + }; + let mut input = std::io::stdin().lock(); + loop { + let mut line = Vec::new(); + if input + .by_ref() + .take(MAX_INPUT_BYTES + 1) + .read_until(b'\n', &mut line)? + == 0 + { + break; + } + if u64::try_from(line.len())? > MAX_INPUT_BYTES { + bail!("request exceeds input limit"); + } + let result = serde_json::from_slice(&line) + .context("invalid JSON request") + .and_then(|request| session.call(&request)); + let response = result.unwrap_or_else(|error| json!({"error": format!("{error:#}")})); + println!("{}", serde_json::to_string(&response)?); + std::io::stdout().flush()?; + } + Ok(()) +} + +fn text<'a>(request: &'a Value, key: &str) -> Result<&'a str> { + request + .get(key) + .and_then(Value::as_str) + .with_context(|| format!("missing string {key}")) +} + +impl Session { + fn node(&self, key: &str) -> NodeKey { + NodeKey { + warp_id: self.fixture.node.warp_id, + local_id: NodeId(Sha256::digest(key.as_bytes()).into()), + } + } + + fn call(&mut self, request: &Value) -> Result { + let allowed: &[&str] = match text(request, "op")? { + "status" => &["op"], + "observe" => &["op", "attempt", "keys"], + "changes" | "reading" => &["op", "attempt"], + "submit" => &["op", "attempt", "key", "value", "request_id"], + "outcome" => &["op", "submission_id"], + _ => bail!("unknown operation"), + }; + if request + .as_object() + .context("request must be an object")? + .keys() + .any(|key| !allowed.contains(&key.as_str())) + { + bail!("unexpected request field; observation bindings are runtime-owned"); + } + match text(request, "op")? { + "status" => { + let application_basis = echo_operation_anchored_node_creation_application_basis_v1( + self.fixture.node, + EchoOperationAnchoredNodeOccupancyV1::Absent, + ); + let basis = self + .fixture + .host + .echo_operation_evaluation_basis_v1(self.fixture.head, application_basis)?; + Ok(json!({ + "worldline": hex::encode(self.fixture.head.worldline_id.as_bytes()), + "state_root": hex::encode(current_state(&self.fixture)?.state_root()), + "commit_id": hex::encode(basis.commit_id()), + "tick": basis.worldline_tick().as_u64(), + })) + } + "observe" => { + let attempt = text(request, "attempt")?; + let keys = request + .get("keys") + .and_then(Value::as_array) + .context("missing keys")?; + let nodes = keys + .iter() + .map(|key| { + key.as_str() + .map(|key| self.node(key)) + .context("key must be a string") + }) + .collect::>>()?; + let observation = self.fixture.host.retain_echo_operation_observation_v1( + attempt, + self.fixture.head, + &nodes, + )?; + let readings = observation + .readings() + .map(|(node, value)| { + json!({ + "node": hex::encode(node.local_id.0), "value_cbor": hex::encode(value) + }) + }) + .collect::>(); + Ok(json!({"attempt":attempt, "readings":readings})) + } + "changes" => { + let nodes = self + .fixture + .host + .echo_operation_observation_changes_v1(text(request, "attempt")?)?; + let commits = self + .fixture + .host + .echo_operation_observation_change_commits_v1(text(request, "attempt")?)?; + Ok( + json!({"changed":!nodes.is_empty(), "nodes":nodes.iter().map(|node| hex::encode(node.local_id.0)).collect::>(), "commits":commits.iter().map(hex::encode).collect::>()}), + ) + } + "reading" => { + let observation = self + .fixture + .host + .echo_operation_observation_v1(text(request, "attempt")?)?; + Ok(json!({"attempt":text(request,"attempt")?, + "observation_commit":hex::encode(observation.basis().commit_id()), + "readings":observation.readings().map(|(node, value)| json!({"node":hex::encode(node.local_id.0),"value_cbor":hex::encode(value)})).collect::>() })) + } + "submit" => { + let attempt = text(request, "attempt")?; + let key = text(request, "key")?; + let input = parse_input( + &serde_json::to_vec( + &json!({"basis":self.basis, "key":key, "value":text(request,"value")?}), + )?, + &self.configuration, + )?; + let node = self.node(key); + let store = current_state(&self.fixture)? + .store(&node.warp_id) + .context("warp unavailable")?; + let occupancy = match ( + store.node(&node.local_id).is_some(), + store.node_attachment(&node.local_id).is_some(), + ) { + (false, false) => EchoOperationAnchoredNodeOccupancyV1::Absent, + (true, false) => EchoOperationAnchoredNodeOccupancyV1::NodeOnly, + (false, true) => EchoOperationAnchoredNodeOccupancyV1::AttachmentOnly, + (true, true) => EchoOperationAnchoredNodeOccupancyV1::NodeAndAttachment, + }; + let application_basis = + echo_operation_anchored_node_creation_application_basis_v1(node, occupancy); + let basis = self + .fixture + .host + .echo_operation_evaluation_basis_v1(self.fixture.head, application_basis)?; + let invocation = EchoOperationInvocationV1::anchored_node_attachment_create_if_absent_with_application_input( + self.package_id, &self.package.operation_coordinate, basis, self.grant, self.package.budget, + node, input.replacement, input.canonical_bytes, + ); + let request_id = request + .get("request_id") + .and_then(Value::as_str) + .unwrap_or(attempt); + let invocation = self + .fixture + .host + .bind_echo_operation_request_v1(request_id, attempt, invocation)?; + let envelope = echo_operation_action_envelope_v1( + IngressTarget::ExactHead { + key: self.fixture.head, + }, + invocation, + )?; + let submission = self + .fixture + .host + .app() + .submit_intent_with_runtime_wal_ack(envelope)? + .submission_id; + if self + .fixture + .host + .echo_operation_action_outcome_v1(&submission) + .is_none() + { + self.fixture.host.tick_once()?; + } + self.outcome(submission) + } + "outcome" => { + let bytes: [u8; 32] = hex::decode(text(request, "submission_id")?)? + .try_into() + .map_err(|_| anyhow::anyhow!("invalid submission identity"))?; + self.outcome(bytes) + } + _ => bail!("unknown operation"), + } + } + + fn outcome(&self, submission: [u8; 32]) -> Result { + let disposition = self + .fixture + .host + .echo_operation_action_outcome_v1(&submission) + .context("outcome unavailable")?; + let mut result = json!({"submission_id":hex::encode(submission)}); + match disposition { + EchoOperationActionOutcomeV1::Committed(receipt) => { + result["outcome"] = json!("committed"); + result["commit_id"] = json!(hex::encode(receipt.commit_id())); + result["receipt_digest"] = json!(hex::encode(receipt.digest())); + result["application_result"] = serde_json::to_value(application_result_report( + receipt + .committed_application_result() + .context("application result unavailable")?, + ))?; + } + EchoOperationActionOutcomeV1::Obstructed(obstruction) => { + result["outcome"] = json!(format!("{:?}", obstruction.kind())); + } + EchoOperationActionOutcomeV1::RejectedFootprintConflict(_) => { + result["outcome"] = json!("FootprintConflict"); + } + } + Ok(result) + } +} From d6703e09b865850f0b70d67b1f8c00b8ec9b5964 Mon Sep 17 00:00:00 2001 From: James Ross Date: Mon, 21 Sep 2026 00:24:18 -0700 Subject: [PATCH 02/15] test(wal): require fresh fencing without consuming empty positions --- crates/warp-core/tests/causal_wal_hardening_tests.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/crates/warp-core/tests/causal_wal_hardening_tests.rs b/crates/warp-core/tests/causal_wal_hardening_tests.rs index 8451c5a41..805690d1e 100644 --- a/crates/warp-core/tests/causal_wal_hardening_tests.rs +++ b/crates/warp-core/tests/causal_wal_hardening_tests.rs @@ -964,7 +964,10 @@ fn filesystem_writer_lease_refuses_overlap_before_takeover() { drop(active); let successor = must_ok(contender.acquire_fresh_writer_epoch(Lsn::from_raw(0))); assert_eq!(successor.previous_epoch_id, Some(epoch_id())); - assert!(successor.started_at_lsn > Lsn::from_raw(0)); + // An empty epoch consumed no record position. Fencing must change the + // writer identity without introducing a gap before the next append. + assert_ne!(successor.epoch_id, epoch_id()); + assert_eq!(successor.started_at_lsn, Lsn::from_raw(0)); drop(contender); must_ok(fs::remove_dir_all(root)); } From 9d92a54a185a4b6399bf3486eb6a50dbbec84e69 Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 08:51:12 -0700 Subject: [PATCH 03/15] fix: bind operation observations to the full writer head --- CHANGELOG.md | 3 +++ .../warp-core/src/echo_operation/observed.rs | 21 ++++++++++++++++++- 2 files changed, 23 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6a0cb1385..09f0154b1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1635,6 +1635,9 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- Observed operations reject observations from another writer head even when + both heads belong to the same worldline. + - Reopening a filesystem WAL through an empty writer epoch no longer consumes an unwritten log position. A second reopen followed by append previously left an LSN gap and made subsequent recovery fail. diff --git a/crates/warp-core/src/echo_operation/observed.rs b/crates/warp-core/src/echo_operation/observed.rs index 37d7bedf7..72670da5f 100644 --- a/crates/warp-core/src/echo_operation/observed.rs +++ b/crates/warp-core/src/echo_operation/observed.rs @@ -169,7 +169,7 @@ impl EchoOperationObservationV1 { footprint: &mut Footprint, meter: &mut EchoOperationBudgetMeterV1, ) -> Result<(), EchoOperationObstructionKindV1> { - if self.basis.writer_head().worldline_id != submission.writer_head().worldline_id + if self.basis.writer_head() != submission.writer_head() || self.basis.worldline_tick() > submission.worldline_tick() { return Err(EchoOperationObstructionKindV1::ObservationChanged); @@ -277,6 +277,25 @@ mod tests { ) } + #[test] + fn observation_cannot_cross_writer_heads_in_one_worldline() { + let (_, state, basis, _, _, _) = super::super::tests::projected_create_fixture(1024); + let observation = EchoOperationObservationV1::capture(&state, basis, &[*state.root()]) + .expect("observation"); + let mut other = basis; + other.writer_head.head_id = crate::make_head_id("another-head"); + let result = observation.validate_at_execution( + &state, + other, + &mut Footprint::default(), + &mut EchoOperationBudgetMeterV1::new(EchoOperationBudgetV1::new(32, 4096, 1024)), + ); + assert_eq!( + result, + Err(EchoOperationObstructionKindV1::ObservationChanged) + ); + } + #[test] fn fresh_submission_cannot_erase_changed_observation() { let (outcome, _) = exercise(true); From 8a74ff7ebbbba9c945597f58259cb13e1f26e0be Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 08:53:14 -0700 Subject: [PATCH 04/15] fix: include observed reads in footprint partition masks --- CHANGELOG.md | 3 +++ crates/warp-core/src/echo_operation/observed.rs | 7 ++++++- 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 09f0154b1..95de316cf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1635,6 +1635,9 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- Observation reads now populate the footprint partition mask as well as the + exact read sets, keeping retained footprint evidence consistent. + - Observed operations reject observations from another writer head even when both heads belong to the same worldline. diff --git a/crates/warp-core/src/echo_operation/observed.rs b/crates/warp-core/src/echo_operation/observed.rs index 72670da5f..e7428e1f1 100644 --- a/crates/warp-core/src/echo_operation/observed.rs +++ b/crates/warp-core/src/echo_operation/observed.rs @@ -186,7 +186,7 @@ impl EchoOperationObservationV1 { if !meter.charge(2, 64 + actual.len() as u64, 0) { return Err(EchoOperationObstructionKindV1::BudgetExceeded); } - footprint.n_read.insert(*node); + record_node_read(footprint, *node); footprint.a_read.insert(AttachmentKey::node_alpha(*node)); if actual != *expected { return Err(EchoOperationObstructionKindV1::ObservationChanged); @@ -316,6 +316,11 @@ mod tests { .n_read .iter() .any(|node| *node == observed)); + assert_ne!( + prepared.actual_footprint().factor_mask & (1_u64 << (observed.local_id.0[0] & 63)), + 0, + "the observation's partition must be represented in the footprint mask" + ); assert!(prepared .patch() .in_slots() From 8521bde89bcd3fd71488a4e93a6a7b3ca8b30715 Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 08:55:34 -0700 Subject: [PATCH 05/15] fix: bind retained observations to actual anchor occupancy --- CHANGELOG.md | 3 +++ .../trusted_runtime_host/observed_context.rs | 23 +++++++++++++++---- .../tests/trusted_runtime_host_loop_tests.rs | 8 +++++++ 3 files changed, 29 insertions(+), 5 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 95de316cf..1d04835a3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1635,6 +1635,9 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- Retained observations bind the anchor's actual occupancy instead of claiming + that every observed anchor was absent. + - Observation reads now populate the footprint partition mask as well as the exact read sets, keeping retained footprint evidence consistent. diff --git a/crates/warp-core/src/trusted_runtime_host/observed_context.rs b/crates/warp-core/src/trusted_runtime_host/observed_context.rs index 165c7f5f6..46b288f63 100644 --- a/crates/warp-core/src/trusted_runtime_host/observed_context.rs +++ b/crates/warp-core/src/trusted_runtime_host/observed_context.rs @@ -162,17 +162,30 @@ impl TrustedRuntimeHost { let first = nodes .first() .ok_or_else(|| error("observation needs an aperture"))?; - let application_basis = - crate::echo_operation_anchored_node_absent_application_basis_v1(*first); - let basis = self - .echo_operation_evaluation_basis_v1(head, application_basis) - .map_err(error)?; let state = self .runtime .worldlines() .get(&head.worldline_id) .ok_or_else(|| error("worldline unavailable"))? .state(); + let store = state + .store(&first.warp_id) + .ok_or_else(|| error("observation warp unavailable"))?; + use crate::EchoOperationAnchoredNodeOccupancyV1 as Occupancy; + let occupancy = match ( + store.node(&first.local_id).is_some(), + store.node_attachment(&first.local_id).is_some(), + ) { + (false, false) => Occupancy::Absent, + (true, false) => Occupancy::NodeOnly, + (false, true) => Occupancy::AttachmentOnly, + (true, true) => Occupancy::NodeAndAttachment, + }; + let application_basis = + crate::echo_operation_anchored_node_creation_application_basis_v1(*first, occupancy); + let basis = self + .echo_operation_evaluation_basis_v1(head, application_basis) + .map_err(error)?; let observation = EchoOperationObservationV1::capture(state, basis, nodes).map_err(error)?; self.runtime_wal diff --git a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs index 52032937d..072f29c6c 100644 --- a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs +++ b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs @@ -1168,6 +1168,14 @@ fn observed_operation_contexts_recover_without_replacing_original_inputs() { let observation = host .retain_echo_operation_observation_v1("attempt", head, &[node]) .expect("retain reading"); + assert_eq!( + observation.basis().application_basis(), + warp_core::echo_operation_anchored_node_creation_application_basis_v1( + node, + warp_core::EchoOperationAnchoredNodeOccupancyV1::NodeOnly, + ), + "the canonical root exists without an alpha attachment" + ); let invocation = |value: &[u8]| { warp_core::EchoOperationInvocationV1::anchored_node_attachment_create_if_absent_with_application_input( warp_core::echo_operation_package_id_v1(b"context-only-test"), "test.context@1.create", From 7305edb2d01189d521cd3db87f104bbde8f0e346 Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 08:56:36 -0700 Subject: [PATCH 06/15] test: preserve request identity across submission occupancy changes --- .../warp-core/src/echo_operation/observed.rs | 31 +++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/crates/warp-core/src/echo_operation/observed.rs b/crates/warp-core/src/echo_operation/observed.rs index e7428e1f1..ce0a479e6 100644 --- a/crates/warp-core/src/echo_operation/observed.rs +++ b/crates/warp-core/src/echo_operation/observed.rs @@ -277,6 +277,37 @@ mod tests { ) } + #[test] + fn retry_identity_survives_submission_occupancy_changes_but_not_input_changes() { + let (_, state, basis, _, invocation, _) = + super::super::tests::projected_create_fixture(1_024); + let observation = EchoOperationObservationV1::capture(&state, basis, &[*state.root()]) + .expect("observation"); + let original = invocation + .observed_semantic_identity(&observation) + .expect("identity"); + let mut retry = invocation.clone(); + retry.evaluation_basis.application_basis = + echo_operation_anchored_node_creation_application_basis_v1( + invocation.node, + EchoOperationAnchoredNodeOccupancyV1::NodeAndAttachment, + ); + retry.evaluation_basis.worldline_tick = WorldlineTick::from_raw(100); + assert_eq!( + retry + .observed_semantic_identity(&observation) + .expect("retry identity"), + original + ); + retry.replacement_bytes.push(1); + assert_ne!( + retry + .observed_semantic_identity(&observation) + .expect("changed identity"), + original + ); + } + #[test] fn observation_cannot_cross_writer_heads_in_one_worldline() { let (_, state, basis, _, _, _) = super::super::tests::projected_create_fixture(1024); From cc7a9fab3089e14a6d62f619b661a9df9ebef847 Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 09:00:38 -0700 Subject: [PATCH 07/15] fix: retain intervening observation changes from native provenance --- CHANGELOG.md | 5 + .../trusted_runtime_host/observed_context.rs | 94 +++++++++----- .../tests/trusted_runtime_host_loop_tests.rs | 120 ++++++++++++++++++ .../application-contract-hosting.md | 9 +- xtask/src/run_edict_operation/session.rs | 8 +- 5 files changed, 195 insertions(+), 41 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1d04835a3..c1705b513 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1635,6 +1635,11 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- Observation change discovery preserves intervening writes after values are + restored, including after reopening. The session driver resolves keys and + commit evidence together from native provenance instead of repeatedly + reconstructing WAL history. + - Retained observations bind the anchor's actual occupancy instead of claiming that every observed anchor was absent. diff --git a/crates/warp-core/src/trusted_runtime_host/observed_context.rs b/crates/warp-core/src/trusted_runtime_host/observed_context.rs index 46b288f63..19f92f5f1 100644 --- a/crates/warp-core/src/trusted_runtime_host/observed_context.rs +++ b/crates/warp-core/src/trusted_runtime_host/observed_context.rs @@ -5,7 +5,9 @@ use super::{ BTreeMap, Error, Hash, TrustedRuntimeHost, TrustedRuntimeWal, WalAppendAuthority, WalTransactionId, WalTransactionKind, }; -use crate::{EchoOperationInvocationV1, EchoOperationObservationV1, NodeKey, WriterHeadKey}; +use crate::{ + EchoOperationInvocationV1, EchoOperationObservationV1, NodeKey, ProvenanceStore, WriterHeadKey, +}; use echo_edict_canonical::{ decode_canonical_cbor_v1, encode_canonical_cbor_v1, CanonicalValueV1 as V, }; @@ -214,46 +216,70 @@ impl TrustedRuntimeHost { .ok_or_else(|| error("observation unavailable; re-observation requires a new attempt")) } - /// Returns only changed resources in this attempt's retained aperture. + /// Returns support changed now or written since this attempt's reading. pub fn echo_operation_observation_changes_v1(&self, attempt: &str) -> Result> { - let observation = self.echo_operation_observation_v1(attempt)?; - let state = self - .runtime - .worldlines() - .get(&observation.basis().writer_head().worldline_id) - .ok_or_else(|| error("worldline unavailable"))? - .state(); - observation.changed_nodes(state).map_err(error) + self.echo_operation_observation_change_evidence_v1(attempt) + .map(|(nodes, _)| nodes) } /// Retained commits which wrote the changed support after this reading. /// Unrelated commits and their payloads are excluded from the result. pub fn echo_operation_observation_change_commits_v1(&self, attempt: &str) -> Result> { + self.echo_operation_observation_change_evidence_v1(attempt) + .map(|(_, commits)| commits) + } + + /// Resolves aperture changes and their commit evidence together. + /// + /// Native retained writes remain discoverable after a value is restored. + /// Current-value differences also remain visible. Missing retained provenance + /// obstructs instead of being interpreted as an unchanged observation. + pub fn echo_operation_observation_change_evidence_v1( + &self, + attempt: &str, + ) -> Result<(Vec, Vec)> { let observation = self.echo_operation_observation_v1(attempt)?; - let changed = self.echo_operation_observation_changes_v1(attempt)?; - let history = self - .runtime_wal - .as_ref() - .ok_or_else(|| error("durable WAL required"))? - .recover_read_only() - .map_err(error)?; - Ok(history - .provenance_entries - .iter() - .filter(|entry| { - entry.worldline_id == observation.basis().writer_head().worldline_id - && entry.worldline_tick >= observation.basis().worldline_tick() - && entry.patch.as_ref().is_some_and(|patch| { - changed.iter().any(|node| { - patch.out_slots.contains(&crate::SlotId::Node(*node)) - || patch.out_slots.contains(&crate::SlotId::Attachment( - crate::AttachmentKey::node_alpha(*node), - )) - }) - }) - }) - .map(|entry| entry.expected.commit_hash) - .collect()) + let lane = observation.basis().writer_head().worldline_id; + let worldline = self + .runtime + .worldlines() + .get(&lane) + .ok_or_else(|| error("worldline unavailable"))?; + let mut changed = observation + .changed_nodes(worldline.state()) + .map_err(error)? + .into_iter() + .collect::>(); + let mut commits = Vec::new(); + // This index is populated by native commitment and reconstructed by WAL + // recovery before the host becomes available. Do not reconstruct it for + // each change query. + for tick in + observation.basis().worldline_tick().as_u64()..worldline.frontier_tick().as_u64() + { + let entry = self + .provenance() + .entry(lane, crate::WorldlineTick::from_raw(tick)) + .map_err(error)?; + let Some(patch) = entry.patch else { + continue; + }; + let mut relevant = false; + for (node, _) in observation.readings() { + if patch.out_slots.contains(&crate::SlotId::Node(node)) + || patch.out_slots.contains(&crate::SlotId::Attachment( + crate::AttachmentKey::node_alpha(node), + )) + { + changed.insert(node); + relevant = true; + } + } + if relevant { + commits.push(entry.expected.commit_hash); + } + } + Ok((changed.into_iter().collect(), commits)) } /// Binds a logical request before ingress acceptance. Exact retry returns diff --git a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs index 072f29c6c..4d24707e5 100644 --- a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs +++ b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs @@ -1144,6 +1144,126 @@ fn filesystem_runtime_wal_empty_reopen_does_not_consume_a_log_position() { assert_eq!(commits[1].first_lsn, Lsn::from_raw(3)); } +#[test] +fn observation_change_discovery_retains_intervening_writes_after_values_are_restored() { + fn toggle_root(view: GraphView<'_>, _: &NodeId, delta: &mut TickDelta) { + let root = make_node_id("root"); + let current = view.node(&root).expect("root").ty; + delta.push(WarpOp::UpsertNode { + node: warp_core::NodeKey { + warp_id: view.warp_id(), + local_id: root, + }, + record: NodeRecord { + ty: if current == make_type_id("world") { + make_type_id("changed") + } else { + make_type_id("world") + }, + }, + }); + } + fn footprint(view: GraphView<'_>, scope: &NodeId) -> warp_core::Footprint { + let mut fp = warp_core::runtime_ingress_eint_read_footprint(view, scope); + fp.n_read + .insert_with_warp(view.warp_id(), make_node_id("root")); + fp.n_write + .insert_with_warp(view.warp_id(), make_node_id("root")); + fp + } + fn matches(view: GraphView<'_>, scope: &NodeId) -> bool { + warp_core::eint_vars_for_op(view, scope, MUTATION_OP_ID).is_some() + } + let path = temp_runtime_wal_dir("observation-aba"); + let (initial, lane) = runtime(); + let root = *initial + .worldlines() + .get(&lane) + .expect("lane") + .state() + .root(); + let head = WriterHeadKey { + worldline_id: lane, + head_id: make_head_id("default"), + }; + let mut host = TrustedRuntimeHost::new(initial, empty_engine()).expect("host"); + host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&path)) + .expect("WAL"); + let mut package = package(); + package.mutation_handlers[0].rule.executor = toggle_root; + package.mutation_handlers[0].rule.matcher = matches; + package.mutation_handlers[0].rule.compute_footprint = footprint; + host.register_contract_package(package).expect("package"); + host.retain_echo_operation_observation_v1("watch", head, &[root]) + .expect("observe"); + let unrelated = warp_core::NodeKey { + warp_id: root.warp_id, + local_id: make_node_id("unrelated"), + }; + host.retain_echo_operation_observation_v1("control", head, &[unrelated]) + .expect("control"); + for index in 0..2 { + let envelope = IngressEnvelope::local_intent( + IngressTarget::DefaultWriter { worldline_id: lane }, + make_intent_kind("echo.intent/eint-v1"), + echo_wasm_abi::pack_intent_v1(MUTATION_OP_ID, &[index]).expect("EINT"), + ); + let submission = host + .app() + .submit_intent_with_runtime_wal_ack(envelope) + .expect("submit"); + host.stage_installed_contract_submission( + submission.submission_id, + &admission_ticket(59 + index), + ) + .expect("stage"); + host.run_until_idle(4).expect("execute"); + } + assert_eq!( + host.runtime() + .worldlines() + .get(&lane) + .expect("lane") + .state() + .store(&root.warp_id) + .expect("store") + .node(&root.local_id) + .expect("root") + .ty, + make_type_id("world") + ); + drop(host); + let (initial, _) = runtime(); + let mut host = TrustedRuntimeHost::new(initial, empty_engine()).expect("fresh host"); + host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&path)) + .expect("reopen"); + host.runtime_wal() + .expect("WAL") + .reset_recover_read_only_call_count_for_test(); + assert_eq!( + host.echo_operation_observation_changes_v1("watch") + .expect("changes"), + vec![root] + ); + assert_eq!( + host.echo_operation_observation_change_commits_v1("watch") + .expect("commits") + .len(), + 2 + ); + assert!(host + .echo_operation_observation_changes_v1("control") + .expect("control changes") + .is_empty()); + assert_eq!( + host.runtime_wal() + .expect("WAL") + .recover_read_only_call_count_for_test(), + 0, + "change discovery must reuse the recovered native provenance index" + ); +} + #[test] fn observed_operation_contexts_recover_without_replacing_original_inputs() { let root = temp_runtime_wal_dir("observed-request-context"); diff --git a/docs/architecture/application-contract-hosting.md b/docs/architecture/application-contract-hosting.md index aaf5fe458..ebf039506 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -205,9 +205,16 @@ produces an obstruction. Unrelated movement does not invalidate the reading. The original invocation remains the replay input, including its observation. Change discovery is a reading over retained observations and native patches. -It returns changed aperture keys and the retained commits that wrote them, +It returns aperture keys whose current values differ or which native patches +wrote after the reading, including a change followed by restoration of the +original value. Notification evidence therefore differs from the value-based +admission precondition: an intervening write can merit attention even when the +original value is valid again. It returns the retained commits that wrote them, without persisting a second notification log. It survives loss of the driver's caches. Unknown observation identities obstruct rather than silently recapture. +The combined change-evidence query resolves the observation once and reads the +host's native provenance index, which is reconstructed during reopening. It does +not repeatedly recover the WAL to derive keys and commit identities separately. This is a trusted, locally scripted host profile, not an authenticated agent service or the complete public optics boundary. It covers bounded atomic node diff --git a/xtask/src/run_edict_operation/session.rs b/xtask/src/run_edict_operation/session.rs index 92bb93bdc..3e5f6450e 100644 --- a/xtask/src/run_edict_operation/session.rs +++ b/xtask/src/run_edict_operation/session.rs @@ -194,14 +194,10 @@ impl Session { Ok(json!({"attempt":attempt, "readings":readings})) } "changes" => { - let nodes = self + let (nodes, commits) = self .fixture .host - .echo_operation_observation_changes_v1(text(request, "attempt")?)?; - let commits = self - .fixture - .host - .echo_operation_observation_change_commits_v1(text(request, "attempt")?)?; + .echo_operation_observation_change_evidence_v1(text(request, "attempt")?)?; Ok( json!({"changed":!nodes.is_empty(), "nodes":nodes.iter().map(|node| hex::encode(node.local_id.0)).collect::>(), "commits":commits.iter().map(hex::encode).collect::>()}), ) From 5b17c7252dd114d6ea600ef6c62ca71062b6da91 Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 09:04:10 -0700 Subject: [PATCH 08/15] fix: refuse writer takeover before WAL tail reconciliation --- CHANGELOG.md | 3 ++ crates/warp-core/src/causal_wal.rs | 15 ++++++ .../tests/causal_wal_hardening_tests.rs | 54 +++++++++++++++++++ 3 files changed, 72 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index c1705b513..3b370efbc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1635,6 +1635,9 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- Filesystem writer takeover refuses an unreconciled WAL tail before reusing + an empty epoch's log position, preventing duplicate physical LSNs. + - Observation change discovery preserves intervening writes after values are restored, including after reopening. The session driver resolves keys and commit evidence together from native provenance instead of repeatedly diff --git a/crates/warp-core/src/causal_wal.rs b/crates/warp-core/src/causal_wal.rs index cb2d327a6..460edc148 100644 --- a/crates/warp-core/src/causal_wal.rs +++ b/crates/warp-core/src/causal_wal.rs @@ -5734,6 +5734,8 @@ impl FilesystemWalStore { /// reread. An active epoch left by a terminated process is closed under /// that lease before the successor is derived and admitted. A concurrently /// live writer retains the lease and prevents takeover. + /// An uncommitted or torn tail must be reconciled by writable recovery + /// before takeover; refusing it preserves the previous epoch ledger. pub fn acquire_fresh_writer_epoch( &mut self, minimum_started_at_lsn: Lsn, @@ -5743,6 +5745,16 @@ impl FilesystemWalStore { } let writer_lock = acquire_writer_epoch_lock(&self.root)?; self.reload_writer_epoch_ledger()?; + let recovery = recover_filesystem_store(&self.root, RecoveryAccessMode::ReadOnly).map_err( + |error| match error { + WalRecoveryError::Store(error) => error, + WalRecoveryError::Validation(error) => WalStoreError::Validation(error), + WalRecoveryError::Index(error) => WalStoreError::RecoveryIndex(error), + }, + )?; + if !matches!(recovery.tail_posture, RecoveryTailPosture::Clean) { + return Err(WalStoreError::SegmentHasUncommittedTail(self.segment_id)); + } if self.active_epoch.is_some() { let previous_ledger = self.writer_epoch_ledger(); @@ -9876,6 +9888,9 @@ pub enum WalValidationError { /// WAL store errors. #[derive(Debug, Error, PartialEq, Eq)] pub enum WalStoreError { + /// The retained prefix could not establish its recovery indexes. + #[error(transparent)] + RecoveryIndex(#[from] WalRecoveryIndexError), /// A writer epoch is already active. #[error("WAL writer epoch already active")] WriterEpochAlreadyActive, diff --git a/crates/warp-core/tests/causal_wal_hardening_tests.rs b/crates/warp-core/tests/causal_wal_hardening_tests.rs index 805690d1e..1f6783815 100644 --- a/crates/warp-core/tests/causal_wal_hardening_tests.rs +++ b/crates/warp-core/tests/causal_wal_hardening_tests.rs @@ -1021,6 +1021,60 @@ fn filesystem_commits_without_writer_epoch_ledger_fail_closed() { must_ok(fs::remove_dir_all(root)); } +#[test] +fn filesystem_takeover_refuses_unreconciled_tail_before_reusing_its_lsn() { + let mut fixture = WalHardeningFixture::new("takeover-uncommitted"); + fixture.append_uncommitted_submission_frame("uncommitted", Lsn::from_raw(0)); + let root = fixture.root.clone(); + drop(fixture); + let mut store = must_ok(FilesystemWalStore::open(&root, WalSegmentId::from_raw(1))); + let error = must_err( + store.acquire_fresh_writer_epoch(Lsn::from_raw(0)), + "takeover must not reuse an occupied uncommitted coordinate", + ); + assert!(matches!(error, WalStoreError::SegmentHasUncommittedTail(_))); + let report = must_ok(recover_filesystem_store( + &root, + RecoveryAccessMode::Writable, + )); + assert_eq!(report.tail_posture, RecoveryTailPosture::TruncatedAll); + let epoch = must_ok(store.acquire_fresh_writer_epoch(Lsn::from_raw(0))); + let builder = WalTransactionBuilder::new( + epoch.epoch_id, + WalSegmentId::from_raw(1), + transaction_id("reconciled"), + WalTransactionKind::SubmissionIntake, + WalAppendAuthority::SubmissionIntake, + Lsn::from_raw(0), + digest("hardening:previous-frame"), + digest("hardening:previous-commit"), + WalDurabilityMode::StrictFilesystem, + PayloadCodecId::from_hash(digest("hardening:codec")), + PayloadSchemaId::from_hash(digest("hardening:schema")), + 1, + 1, + digest("hardening:domain"), + ); + must_ok( + store.append_transaction(must_ok(build_submission_acceptance_transaction( + builder, + submission_acceptance("reconciled"), + vec![frontier( + "reconciled", + AffectedFrontierKind::SubmissionQueue, + )], + ))), + ); + drop(store); + let report = must_ok(recover_filesystem_store( + &root, + RecoveryAccessMode::ReadOnly, + )); + assert_eq!(report.tail_posture, RecoveryTailPosture::Clean); + assert_eq!(report.last_committed_lsn(), Some(Lsn::from_raw(1))); + must_ok(fs::remove_dir_all(root)); +} + #[test] fn duplicate_writer_epoch_id_is_a_chain_gap() { let mut store = InMemoryWalStore::new(); From e823fd1042d1356aa71313152dd2f8904abc151c Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 09:05:58 -0700 Subject: [PATCH 09/15] fix: reject unsupported observation-session budgets before setup --- CHANGELOG.md | 4 ++++ .../application-contract-hosting.md | 10 ++++++++ xtask/src/run_edict_operation/session.rs | 6 +++++ xtask/tests/run_edict_operation.rs | 23 +++++++++++++++++++ 4 files changed, 43 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3b370efbc..176224298 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1635,6 +1635,10 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- The observed-session driver refuses an obviously insufficient compiled + budget before durable setup and explains that observation reads require an + authored, verified allowance rather than silently increasing the grant. + - Filesystem writer takeover refuses an unreconciled WAL tail before reusing an empty epoch's log position, preventing duplicate physical LSNs. diff --git a/docs/architecture/application-contract-hosting.md b/docs/architecture/application-contract-hosting.md index ebf039506..d35f26da1 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -204,6 +204,16 @@ reads. Relevant value changes produce `ObservationChanged`; unavailable support produces an obstruction. Unrelated movement does not invalidate the reading. The original invocation remains the replay input, including its observation. +The pinned Hello Echo fixture has a 64-byte read budget for an unobserved +one-shot create; it is not an observation-session profile. `--serve` refuses +that insufficient budget before creating a WAL. An authored observation profile +must budget for both the operation and its aperture: each node reading costs two +steps and 64 bytes plus its encoded value, with additional descent reads where +applicable. Compile and independently verify that profile before installation; +the driver never increases a package or grant ceiling. Passing the startup +minimum does not guarantee that a larger aperture fits. Runtime metering still +returns `BudgetExceeded` when the actual admitted allowance is exhausted. + Change discovery is a reading over retained observations and native patches. It returns aperture keys whose current values differ or which native patches wrote after the reading, including a change followed by restoration of the diff --git a/xtask/src/run_edict_operation/session.rs b/xtask/src/run_edict_operation/session.rs index 3e5f6450e..121025744 100644 --- a/xtask/src/run_edict_operation/session.rs +++ b/xtask/src/run_edict_operation/session.rs @@ -52,6 +52,12 @@ pub fn serve(config: RunEdictOperationConfig) -> Result<()> { package.target_intrinsic, )?; validate_package_configuration(&package, &configuration)?; + // One root-level observed slot costs two steps and 64 bytes plus its + // nonempty encoded value; the create operation also reads 64 bytes. These + // are necessary lower bounds, not a promise that every aperture fits. + if package.budget.read_bytes() <= 128 || package.budget.steps() < 3 { + bail!("insufficient budget for an observation-bound session: compile an observation-capable profile with more than 128 read bytes and at least 3 steps; the full aperture and operation remain subject to the admitted ceiling"); + } let input = parse_input( &read_bounded(&config.input, MAX_INPUT_BYTES, "bootstrap")?, &configuration, diff --git a/xtask/tests/run_edict_operation.rs b/xtask/tests/run_edict_operation.rs index 80e0a9036..c32ca07dc 100644 --- a/xtask/tests/run_edict_operation.rs +++ b/xtask/tests/run_edict_operation.rs @@ -142,6 +142,29 @@ fn assert_rejected(output: &Output, expected_reason: &str) { ); } +#[test] +fn observed_session_rejects_an_unobserved_only_budget_before_creating_wal() { + let run_dir = TempRunDir::new(); + let wal = run_dir.path().join("wal"); + let output = runner_command( + &fixture_path("executable-operation-package.cbor"), + &fixture_path("verification-report.cbor"), + &fixture_path("input.json"), + &wal, + ) + .arg("--serve") + .output() + .expect("session starts"); + assert_rejected( + &output, + "insufficient budget for an observation-bound session", + ); + assert!( + !wal.exists(), + "unsupported budget must fail before durable setup" + ); +} + #[test] fn compiler_emitted_operation_runs_durably_without_native_callbacks() { let run_dir = TempRunDir::new(); From b8cf973f123f7f3b2fa729713b274190b721e847 Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 09:12:41 -0700 Subject: [PATCH 10/15] fix: keep observed-context imports at module scope --- crates/warp-core/src/trusted_runtime_host/observed_context.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/warp-core/src/trusted_runtime_host/observed_context.rs b/crates/warp-core/src/trusted_runtime_host/observed_context.rs index 19f92f5f1..daf335c62 100644 --- a/crates/warp-core/src/trusted_runtime_host/observed_context.rs +++ b/crates/warp-core/src/trusted_runtime_host/observed_context.rs @@ -5,6 +5,7 @@ use super::{ BTreeMap, Error, Hash, TrustedRuntimeHost, TrustedRuntimeWal, WalAppendAuthority, WalTransactionId, WalTransactionKind, }; +use crate::EchoOperationAnchoredNodeOccupancyV1 as Occupancy; use crate::{ EchoOperationInvocationV1, EchoOperationObservationV1, NodeKey, ProvenanceStore, WriterHeadKey, }; @@ -173,7 +174,6 @@ impl TrustedRuntimeHost { let store = state .store(&first.warp_id) .ok_or_else(|| error("observation warp unavailable"))?; - use crate::EchoOperationAnchoredNodeOccupancyV1 as Occupancy; let occupancy = match ( store.node(&first.local_id).is_some(), store.node_attachment(&first.local_id).is_some(), From 852a51e290bc8f1b6ab309c3bd3595a01aaf60ef Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 09:12:04 -0700 Subject: [PATCH 11/15] fix: release writer leases despite inherited descriptors --- CHANGELOG.md | 4 +++ crates/warp-core/src/causal_wal.rs | 44 ++++++++++++++++++++++++++++-- docs/topics/WAL.md | 5 ++++ 3 files changed, 50 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 176224298..d92e0412b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1635,6 +1635,10 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- Filesystem writer leases explicitly unlock when their owner leaves scope, + so a descriptor briefly inherited by a concurrent child process cannot keep + the departed writer's lease alive and spuriously refuse its successor. + - The observed-session driver refuses an obviously insufficient compiled budget before durable setup and explains that observation reads require an authored, verified allowance rather than silently increasing the grant. diff --git a/crates/warp-core/src/causal_wal.rs b/crates/warp-core/src/causal_wal.rs index 460edc148..691fc6e35 100644 --- a/crates/warp-core/src/causal_wal.rs +++ b/crates/warp-core/src/causal_wal.rs @@ -5585,7 +5585,7 @@ pub struct FilesystemWalStore { active_epoch: Option, closed_epochs: Vec, epoch_closures: BTreeMap, - writer_lock: Option, + writer_lock: Option, manifests: Vec, sync_evidence: Vec, #[cfg(any(test, feature = "host_test"))] @@ -8477,7 +8477,18 @@ fn reconcile_writer_epoch_closures( Ok(()) } -fn acquire_writer_epoch_lock(root: &Path) -> Result { +#[derive(Debug)] +struct WriterEpochLock(File); + +impl Drop for WriterEpochLock { + fn drop(&mut self) { + // A concurrent fork can retain the open file description until exec. + // Relinquish this owner's lock explicitly before closing its descriptor. + let _ = self.0.unlock(); + } +} + +fn acquire_writer_epoch_lock(root: &Path) -> Result { let lock_path = root.join("writer-epoch.lock"); let lock = OpenOptions::new() .create(true) @@ -8491,7 +8502,34 @@ fn acquire_writer_epoch_lock(root: &Path) -> Result { std::fs::TryLockError::Error(error) => Err(error.into()), }; } - Ok(lock) + Ok(WriterEpochLock(lock)) +} + +#[cfg(test)] +#[allow(clippy::expect_used)] +mod writer_lease_tests { + use super::*; + + #[test] + fn dropping_writer_releases_lease_with_an_inherited_descriptor_alive() { + let root = + std::env::temp_dir().join(format!("echo-inherited-lease-{}", std::process::id())); + fs::create_dir_all(&root).expect("scratch directory"); + let lease = acquire_writer_epoch_lock(&root).expect("first writer"); + // A forked process briefly inherits the same open file description + // before exec closes CLOEXEC descriptors. A duplicate models that + // lifetime deterministically, without depending on a process race. + let inherited = lease.0.try_clone().expect("inherited descriptor"); + assert!(matches!( + acquire_writer_epoch_lock(&root), + Err(WalStoreError::WriterEpochLeaseUnavailable) + )); + drop(lease); + let successor = acquire_writer_epoch_lock(&root).expect("successor after owner exits"); + drop(successor); + drop(inherited); + fs::remove_dir_all(root).expect("scratch cleanup"); + } } fn sync_directory_store(path: &Path) -> Result<(), WalStoreError> { diff --git a/docs/topics/WAL.md b/docs/topics/WAL.md index 829624f04..489a56095 100644 --- a/docs/topics/WAL.md +++ b/docs/topics/WAL.md @@ -215,6 +215,11 @@ Duplicate identities, stale or missing predecessor links, reused fencing evidence, LSN regression, corrupted ledgers, and commits without their epoch ledger fail closed before append. +Orderly owner teardown explicitly unlocks the lease before closing its file. +This prevents a descriptor temporarily inherited by a concurrent fork from +extending the departed owner's lock until the child executes its program. +The live owner's lease still excludes every competing writer. + The operating-system lease is the filesystem adapter's exclusion authority. The persisted fencing, process, host, and lease fields are deterministic chain-position markers, not ambient PID or machine measurements and not a From edae3453ba1b52df5083ce44515d22878af412e4 Mon Sep 17 00:00:00 2001 From: James Ross Date: Tue, 22 Sep 2026 09:23:08 -0700 Subject: [PATCH 12/15] test: keep writer-lease scratch paths deterministic --- crates/warp-core/src/causal_wal.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/crates/warp-core/src/causal_wal.rs b/crates/warp-core/src/causal_wal.rs index 691fc6e35..11da163e6 100644 --- a/crates/warp-core/src/causal_wal.rs +++ b/crates/warp-core/src/causal_wal.rs @@ -8512,8 +8512,7 @@ mod writer_lease_tests { #[test] fn dropping_writer_releases_lease_with_an_inherited_descriptor_alive() { - let root = - std::env::temp_dir().join(format!("echo-inherited-lease-{}", std::process::id())); + let root = PathBuf::from("target/warp-core-test-tmp/inherited-writer-lease"); fs::create_dir_all(&root).expect("scratch directory"); let lease = acquire_writer_epoch_lock(&root).expect("first writer"); // A forked process briefly inherits the same open file description From 71b7ff0a7b228de99f5ee933fa0e3bed0a8bff6b Mon Sep 17 00:00:00 2001 From: James Ross Date: Thu, 8 Oct 2026 00:18:05 -0700 Subject: [PATCH 13/15] test: keep one copy of merged takeover regression --- .../tests/causal_wal_hardening_tests.rs | 53 ------------------- 1 file changed, 53 deletions(-) diff --git a/crates/warp-core/tests/causal_wal_hardening_tests.rs b/crates/warp-core/tests/causal_wal_hardening_tests.rs index 73080df99..23675b597 100644 --- a/crates/warp-core/tests/causal_wal_hardening_tests.rs +++ b/crates/warp-core/tests/causal_wal_hardening_tests.rs @@ -3013,56 +3013,3 @@ fn segment_manifest_validation_gate_is_part_of_wal_release_readiness() { ); } -#[test] -fn filesystem_takeover_refuses_unreconciled_tail_before_reusing_its_lsn() { - let mut fixture = WalHardeningFixture::new("takeover-uncommitted"); - fixture.append_uncommitted_submission_frame("uncommitted", Lsn::from_raw(0)); - let root = fixture.root.clone(); - drop(fixture); - let mut store = must_ok(FilesystemWalStore::open(&root, WalSegmentId::from_raw(1))); - let error = must_err( - store.acquire_fresh_writer_epoch(Lsn::from_raw(0)), - "takeover must not reuse an occupied uncommitted coordinate", - ); - assert!(matches!(error, WalStoreError::SegmentHasUncommittedTail(_))); - let report = must_ok(recover_filesystem_store( - &root, - RecoveryAccessMode::Writable, - )); - assert_eq!(report.tail_posture, RecoveryTailPosture::TruncatedAll); - let epoch = must_ok(store.acquire_fresh_writer_epoch(Lsn::from_raw(0))); - let builder = WalTransactionBuilder::new( - epoch.epoch_id, - WalSegmentId::from_raw(1), - transaction_id("reconciled"), - WalTransactionKind::SubmissionIntake, - WalAppendAuthority::SubmissionIntake, - Lsn::from_raw(0), - digest("hardening:previous-frame"), - digest("hardening:previous-commit"), - WalDurabilityMode::StrictFilesystem, - PayloadCodecId::from_hash(digest("hardening:codec")), - PayloadSchemaId::from_hash(digest("hardening:schema")), - 1, - 1, - digest("hardening:domain"), - ); - must_ok( - store.append_transaction(must_ok(build_submission_acceptance_transaction( - builder, - submission_acceptance("reconciled"), - vec![frontier( - "reconciled", - AffectedFrontierKind::SubmissionQueue, - )], - ))), - ); - drop(store); - let report = must_ok(recover_filesystem_store( - &root, - RecoveryAccessMode::ReadOnly, - )); - assert_eq!(report.tail_posture, RecoveryTailPosture::Clean); - assert_eq!(report.last_committed_lsn(), Some(Lsn::from_raw(1))); - must_ok(fs::remove_dir_all(root)); -} From 17c0e8760d03c73882d712a3f0b64908c115a4c4 Mon Sep 17 00:00:00 2001 From: James Ross Date: Thu, 8 Oct 2026 00:20:07 -0700 Subject: [PATCH 14/15] fix: bound observation slots before atom materialization --- CHANGELOG.md | 2 + .../warp-core/src/echo_operation/observed.rs | 123 +++++++++++++++--- .../tests/causal_wal_hardening_tests.rs | 1 - .../application-contract-hosting.md | 3 + 4 files changed, 113 insertions(+), 16 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 345d3bf99..e1a1749b0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,8 @@ ### Fixed +- Observation slots preflight their canonical byte bound and charge execution reads before copying Atom payloads. + - Filesystem writer takeover preserves unused LSNs after empty epochs and refuses unreconciled tails before changing the epoch ledger. - Incremental snapshot state roots now include the existing v1 domain separator and agree with canonical snapshot hashing. Previously emitted incorrect accumulator roots are not migrated. diff --git a/crates/warp-core/src/echo_operation/observed.rs b/crates/warp-core/src/echo_operation/observed.rs index ce0a479e6..035254a16 100644 --- a/crates/warp-core/src/echo_operation/observed.rs +++ b/crates/warp-core/src/echo_operation/observed.rs @@ -14,10 +14,61 @@ pub(super) const OBSERVED_SCHEMA: &str = "echo.operation-invocation.observed/v1" const MAX_READS: usize = 16; const MAX_BYTES: usize = 4096; +// Measure the canonical envelope using empty atom storage, before copying any +// support bytes. The payload bound also makes CBOR byte-string headers <= 3 bytes. +fn slot_encoded_len( + state: &WorldlineState, + node: NodeKey, +) -> Result { + let store = state + .store(&node.warp_id) + .ok_or_else(|| invalid_structure("observation warp unavailable"))?; + let node_type = store + .node(&node.local_id) + .map_or(CanonicalValueV1::Null, |record| hash_value(record.ty.0)); + let (alpha, payload_len) = match store.node_attachment(&node.local_id) { + None => (CanonicalValueV1::Null, None), + Some(AttachmentValue::Atom(atom)) => { + if atom.bytes.len() > MAX_BYTES { + return Err(invalid_structure("observation exceeds retained byte bound")); + } + ( + map_value([ + ("type", hash_value(atom.type_id.0)), + ("bytes", CanonicalValueV1::Bytes(Vec::new())), + ]), + Some(atom.bytes.len()), + ) + } + Some(_) => { + return Err(invalid_structure( + "observation requires atom or absent alpha", + )) + } + }; + let empty_len = + encode_canonical_cbor_v1(&map_value([("node_type", node_type), ("alpha", alpha)])) + .map_err(canonical_error)? + .len(); + let length = payload_len.map_or(empty_len, |length| { + let header_len = match length { + 0..=23 => 1, + 24..=255 => 2, + _ => 3, + }; + empty_len - 1 + header_len + length + }); + if length > MAX_BYTES { + return Err(invalid_structure("observation exceeds retained byte bound")); + } + Ok(length) +} + fn slot_bytes( state: &WorldlineState, node: NodeKey, ) -> Result, EchoOperationArtifactErrorV1> { + slot_encoded_len(state, node)?; let store = state .store(&node.warp_id) .ok_or_else(|| invalid_structure("observation warp unavailable"))?; @@ -79,20 +130,16 @@ impl EchoOperationObservationV1 { if nodes.is_empty() || nodes.len() > MAX_READS { return Err(invalid_structure("observation needs one to sixteen slots")); } - let reads = nodes - .into_iter() - .map(|node| Ok((node, slot_bytes(state, node)?))) - .collect::, EchoOperationArtifactErrorV1>>()?; - let observation = Self { basis, reads }; - if observation - .reads - .iter() - .map(|(_, bytes)| bytes.len()) - .sum::() - > MAX_BYTES - { - return Err(invalid_structure("observation exceeds retained byte bound")); + let mut reads = Vec::new(); + let mut remaining = MAX_BYTES; + for node in nodes { + let length = slot_encoded_len(state, node)?; + remaining = remaining + .checked_sub(length) + .ok_or_else(|| invalid_structure("observation exceeds retained byte bound"))?; + reads.push((node, slot_bytes(state, node)?)); } + let observation = Self { basis, reads }; Ok(observation) } @@ -181,11 +228,13 @@ impl EchoOperationObservationV1 { meter.charge(1, 32, 0) })?; let _ = portals; - let actual = slot_bytes(state, *node) + let length = slot_encoded_len(state, *node) .map_err(|_| EchoOperationObstructionKindV1::ObservationUnavailable)?; - if !meter.charge(2, 64 + actual.len() as u64, 0) { + if !meter.charge(2, 64 + length as u64, 0) { return Err(EchoOperationObstructionKindV1::BudgetExceeded); } + let actual = slot_bytes(state, *node) + .map_err(|_| EchoOperationObstructionKindV1::ObservationUnavailable)?; record_node_read(footprint, *node); footprint.a_read.insert(AttachmentKey::node_alpha(*node)); if actual != *expected { @@ -277,6 +326,50 @@ mod tests { ) } + #[test] + fn slot_preflight_accounts_for_cbor_header_boundaries() { + let (_, mut state, _, _, _, _) = super::super::tests::projected_create_fixture(1_024); + let node = *state.root(); + for length in [0, 23, 24, 255, 256, 512] { + state + .warp_state + .store_mut(&node.warp_id) + .expect("store") + .set_node_attachment( + node.local_id, + Some(AttachmentValue::Atom(crate::AtomPayload::new( + crate::make_type_id("observation"), + vec![7; length].into(), + ))), + ); + assert_eq!( + slot_encoded_len(&state, node).expect("bounded size"), + slot_bytes(&state, node).expect("bounded wire value").len() + ); + } + } + + #[test] + fn oversized_atom_slot_is_refused_before_materializing_observation_bytes() { + let (_, mut state, _, _, _, _) = super::super::tests::projected_create_fixture(1_024); + let node = *state.root(); + state + .warp_state + .store_mut(&node.warp_id) + .expect("store") + .set_node_attachment( + node.local_id, + Some(AttachmentValue::Atom(crate::AtomPayload::new( + crate::make_type_id("oversized-observation"), + vec![7; MAX_BYTES * 2].into(), + ))), + ); + assert!( + slot_bytes(&state, node).is_err(), + "slot encoding must refuse oversized atom support itself" + ); + } + #[test] fn retry_identity_survives_submission_occupancy_changes_but_not_input_changes() { let (_, state, basis, _, invocation, _) = diff --git a/crates/warp-core/tests/causal_wal_hardening_tests.rs b/crates/warp-core/tests/causal_wal_hardening_tests.rs index 23675b597..da172a15a 100644 --- a/crates/warp-core/tests/causal_wal_hardening_tests.rs +++ b/crates/warp-core/tests/causal_wal_hardening_tests.rs @@ -3012,4 +3012,3 @@ fn segment_manifest_validation_gate_is_part_of_wal_release_readiness() { "manifest validation should be an explicit WAL release gate" ); } - diff --git a/docs/architecture/application-contract-hosting.md b/docs/architecture/application-contract-hosting.md index 59e9d8991..adaef0d8a 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -382,6 +382,9 @@ and attachment readings in the caller's granted aperture, not model-internal influence, arbitrary subtree observations, or speculative strand settlement. The driver accepts at most 1,024 observations and 4,096 request bindings per WAL; each observation contains at most 16 nodes and 4,096 retained value bytes. +Capture preflights the aggregate canonical value size before copying Atom bytes. +Execution checks the slot ceiling and charges its admitted read budget before +materializing the value; oversized support obstructs rather than allocating it. It currently rebuilds context indexes from the WAL. Production indexing, retention policy, and authenticated aperture delegation remain separate work. From 49d1ec36419fda71036af2ecd6de66437375e078 Mon Sep 17 00:00:00 2001 From: James Ross Date: Thu, 8 Oct 2026 00:25:36 -0700 Subject: [PATCH 15/15] test: cover aggregate observation and execution budget boundaries --- .../warp-core/src/echo_operation/observed.rs | 59 +++++++++++++++++++ .../application-contract-hosting.md | 3 +- 2 files changed, 61 insertions(+), 1 deletion(-) diff --git a/crates/warp-core/src/echo_operation/observed.rs b/crates/warp-core/src/echo_operation/observed.rs index 035254a16..79ad22044 100644 --- a/crates/warp-core/src/echo_operation/observed.rs +++ b/crates/warp-core/src/echo_operation/observed.rs @@ -326,6 +326,65 @@ mod tests { ) } + #[test] + fn capture_enforces_aggregate_slots_and_execution_meter_boundaries() { + let (_, mut state, mut basis, _, _, _) = + super::super::tests::projected_create_fixture(1_024); + let first = *state.root(); + let second = NodeKey { + warp_id: first.warp_id, + local_id: crate::make_node_id("second-observation"), + }; + for length in [1_900, 2_000] { + let store = state.warp_state.store_mut(&first.warp_id).expect("store"); + store.insert_node( + second.local_id, + NodeRecord { + ty: crate::make_type_id("node"), + }, + ); + for node in [first, second] { + store.set_node_attachment( + node.local_id, + Some(AttachmentValue::Atom(crate::AtomPayload::new( + crate::make_type_id("observation"), + vec![7; length].into(), + ))), + ); + } + basis.state_root = state.state_root(); + let actual = EchoOperationObservationV1::capture(&state, basis, &[first, second]); + if length == 2_000 { + assert!( + actual.is_err(), + "individually bounded slots exceed aggregate allowance" + ); + continue; + } + let observation = actual.expect("two slots fit aggregate allowance"); + let cost = observation + .reads + .iter() + .map(|(_, bytes)| 64 + bytes.len() as u64) + .sum(); + let mut exact = EchoOperationBudgetMeterV1::new(EchoOperationBudgetV1::new(4, cost, 0)); + observation + .validate_at_execution(&state, basis, &mut Footprint::default(), &mut exact) + .expect("exact admitted read budget"); + let mut short = + EchoOperationBudgetMeterV1::new(EchoOperationBudgetV1::new(4, cost - 1, 0)); + assert_eq!( + observation.validate_at_execution( + &state, + basis, + &mut Footprint::default(), + &mut short + ), + Err(EchoOperationObstructionKindV1::BudgetExceeded) + ); + } + } + #[test] fn slot_preflight_accounts_for_cbor_header_boundaries() { let (_, mut state, _, _, _, _) = super::super::tests::projected_create_fixture(1_024); diff --git a/docs/architecture/application-contract-hosting.md b/docs/architecture/application-contract-hosting.md index adaef0d8a..08d425f9e 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -382,7 +382,8 @@ and attachment readings in the caller's granted aperture, not model-internal influence, arbitrary subtree observations, or speculative strand settlement. The driver accepts at most 1,024 observations and 4,096 request bindings per WAL; each observation contains at most 16 nodes and 4,096 retained value bytes. -Capture preflights the aggregate canonical value size before copying Atom bytes. +Capture preflights each slot against the remaining aggregate canonical byte +allowance before copying that slot’s Atom bytes. Execution checks the slot ceiling and charges its admitted read budget before materializing the value; oversized support obstructs rather than allocating it. It currently rebuilds context indexes from the WAL. Production indexing,