diff --git a/CHANGELOG.md b/CHANGELOG.md index fd0d48cf..f67e268b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1640,6 +1640,11 @@ Applied, Rejected, Obstructed}` with receipt evidence and typed contract ### Fixed +- Trusted-host operation observations and logical-request bindings now advance a + derived index from acknowledged WAL appends instead of reconstructing retained + history on each call. Append errors invalidate the fast path; reconciliation + still resolves commits that became durable before the caller observed failure. + - 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. diff --git a/crates/warp-core/src/causal_wal.rs b/crates/warp-core/src/causal_wal.rs index 11da163e..f4f86b4e 100644 --- a/crates/warp-core/src/causal_wal.rs +++ b/crates/warp-core/src/causal_wal.rs @@ -6163,13 +6163,33 @@ impl WalStorePort for FilesystemWalStore { } } +thread_local! { + static FILESYSTEM_RECOVERY_WORK: std::cell::Cell<(u64, u64)> = const { std::cell::Cell::new((0, 0)) }; +} + +/// Cumulative filesystem recovery calls and decoded frames on the calling thread. +/// Failed scans count as calls; frames count after segment decoding succeeds. +/// Diagnostic work counters are not retained evidence or admission authority. +#[must_use] +pub fn filesystem_recovery_work() -> (u64, u64) { + FILESYSTEM_RECOVERY_WORK.get() +} + /// Recovers committed transactions from filesystem WAL segments. pub fn recover_filesystem_store( root: impl AsRef, mode: RecoveryAccessMode, ) -> Result { + FILESYSTEM_RECOVERY_WORK.set(( + filesystem_recovery_work().0 + 1, + filesystem_recovery_work().1, + )); let root = root.as_ref(); let (frames, commits, torn_tail) = read_filesystem_segments(root)?; + FILESYSTEM_RECOVERY_WORK.set(( + filesystem_recovery_work().0, + filesystem_recovery_work().1 + frames.len() as u64, + )); let mut report = recover_from_frames_and_commits(&frames, &commits, mode)?; if torn_tail && matches!(report.tail_posture, RecoveryTailPosture::Clean) { report.tail_posture = match mode { diff --git a/crates/warp-core/src/trusted_runtime_host.rs b/crates/warp-core/src/trusted_runtime_host.rs index e6681b62..32d018f1 100644 --- a/crates/warp-core/src/trusted_runtime_host.rs +++ b/crates/warp-core/src/trusted_runtime_host.rs @@ -166,6 +166,12 @@ impl From for TrustedRuntimeHostError { /// Error returned by the trusted runtime WAL adapter. #[derive(Debug, Error, PartialEq, Eq)] pub enum TrustedRuntimeWalError { + /// Retained operation contexts could not be validated. + #[error("trusted runtime WAL context error: {0}")] + OperationContext(#[from] EchoOperationContextErrorV1), + /// An earlier append must be reconciled before another can reach storage. + #[error("trusted runtime WAL writer prefix requires reconciliation")] + UncertainWriterPrefix, /// WAL transaction construction failed before storage append. #[error("trusted runtime WAL transaction build error: {0}")] Build(#[from] WalBuildError), @@ -2440,6 +2446,8 @@ impl CausalAnchorClaimProjection { /// Minimal trusted-runtime WAL adapter for ACK-boundary integration tests. #[derive(Debug)] pub struct TrustedRuntimeWal { + // Derived from a validated committed prefix; None after any uncertain append. + operation_context_index: Option<(Hash, observed_context::Contexts)>, store: TrustedRuntimeWalStore, evidence_catalog: Option, evidence_catalog_posture: EvidenceCatalogPosture, @@ -2484,6 +2492,7 @@ impl TrustedRuntimeWal { let mut store = TrustedRuntimeWalStore::open(store)?; let recovery_report = store.recover_for_writer()?; let recovered_cursor = TrustedRuntimeWalCursor::from_recovery(&recovery_report)?; + let contexts = observed_context::Contexts::from_recovery(&recovery_report)?; let durable_submission_acceptances = recover_submission_index(&recovery_report) .map_err(WalRecoveryError::from)? .entries() @@ -2503,6 +2512,10 @@ impl TrustedRuntimeWal { let writer_epoch = writer_epoch.epoch_id; let durability_mode = store.durability_mode(); Ok(Self { + operation_context_index: Some(( + recovered_cursor.previous_committed_transaction_digest, + contexts, + )), store, writer_epoch, segment_id: WalSegmentId::from_raw(1), @@ -2566,6 +2579,7 @@ impl TrustedRuntimeWal { fn in_memory_rollback_snapshot(&self) -> Option { Some(Self { store: TrustedRuntimeWalStore::InMemory(self.store.cloned_in_memory_store()?), + operation_context_index: self.operation_context_index.clone(), evidence_catalog: self.evidence_catalog.clone(), evidence_catalog_posture: self.evidence_catalog_posture.clone(), #[cfg(any(test, feature = "host_test"))] @@ -2805,7 +2819,9 @@ impl TrustedRuntimeWal { } fn refresh_cursor_from_store_for_writer(&mut self) -> Result<(), TrustedRuntimeWalError> { + self.operation_context_index = None; let report = self.store.recover_for_writer()?; + let contexts = observed_context::Contexts::from_recovery(&report)?; let cursor = TrustedRuntimeWalCursor::from_recovery(&report)?; let durable_submission_acceptances = recover_submission_index(&report) .map_err(WalRecoveryError::from)? @@ -2840,6 +2856,7 @@ impl TrustedRuntimeWal { self.causal_history_frontier_digest = cursor.causal_history_frontier_digest; self.causal_anchor_claim_projection = cursor.causal_anchor_claim_projection; self.durable_submission_acceptances = durable_submission_acceptances; + self.operation_context_index = Some((self.previous_committed_transaction_digest, contexts)); Ok(()) } @@ -3281,7 +3298,19 @@ impl TrustedRuntimeWal { let last_good_commit = self.previous_committed_transaction_digest; let commit = transaction.commit.clone(); let frames = transaction.frames.clone(); + // Take the candidate out of circulation before storage can fail. Apply + // only the new records, then publish it after durable acknowledgement. + // A transaction built from an uncertain cursor never reaches storage. + let (prefix, mut contexts) = self + .operation_context_index + .take() + .ok_or(TrustedRuntimeWalError::UncertainWriterPrefix)?; + if prefix != self.previous_committed_transaction_digest { + return Err(TrustedRuntimeWalError::UncertainWriterPrefix); + } + contexts.apply_frames(&frames)?; self.store.append_transaction(transaction)?; + self.operation_context_index = Some((commit.commit_digest, contexts)); self.next_lsn = next_lsn; self.previous_frame_digest = last_frame_digest; self.previous_committed_transaction_digest = commit.commit_digest; 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 daf335c6..9882c7ea 100644 --- a/crates/warp-core/src/trusted_runtime_host/observed_context.rs +++ b/crates/warp-core/src/trusted_runtime_host/observed_context.rs @@ -14,7 +14,7 @@ use echo_edict_canonical::{ }; /// Failure to retain or resolve an immutable operation context. -#[derive(Debug, Error)] +#[derive(Debug, Error, PartialEq, Eq)] #[error("operation context: {0}")] pub struct EchoOperationContextErrorV1(String); @@ -46,80 +46,106 @@ fn check_id(id: &str) -> Result<()> { Ok(()) } -#[derive(Default)] -struct Contexts { +#[derive(Debug, Clone, Default)] +pub(super) 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")); +impl Contexts { + pub(super) fn from_recovery(report: &crate::causal_wal::RecoveryScanReport) -> Result { + let mut contexts = Self::default(); + for transaction in &report.transactions { + contexts.apply_frames(&transaction.frames)?; + } + Ok(contexts) + } + + pub(super) fn apply_frames(&mut self, frames: &[crate::causal_wal::WalFrame]) -> Result<()> { + for frame in 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 self.observations.insert(id, observation).is_some() { + return Err(error("duplicate observation identity")); + } } - 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 = self + .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")); } - "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")); - } + if self + .requests + .insert(id, (attempt, digest, invocation_bytes)) + .is_some() + { + return Err(error("duplicate request identity")); } - _ => return Err(error("invalid context record kind")), } + _ => return Err(error("invalid context record kind")), } } - Ok(contexts) + Ok(()) + } +} + +impl TrustedRuntimeWal { + fn operation_contexts(&self) -> Result> { + if let Some((prefix, contexts)) = &self.operation_context_index { + if *prefix == self.previous_committed_transaction_digest { + return Ok(std::borrow::Cow::Borrowed(contexts)); + } + } + // A read after an uncertain append resolves retained truth, but does not + // restore permission to append from the old writer cursor. + let report = self.store.recover_read_only().map_err(error)?; + Ok(std::borrow::Cow::Owned(Contexts::from_recovery(&report)?)) + } + + fn reconcile_operation_context_prefix(&mut self) -> Result<()> { + if self + .operation_context_index + .as_ref() + .is_none_or(|(prefix, _)| *prefix != self.previous_committed_transaction_digest) + { + self.refresh_cursor_from_store_for_writer().map_err(error)?; + } + Ok(()) } 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)?; + self.reconcile_operation_context_prefix()?; 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( @@ -151,11 +177,12 @@ impl TrustedRuntimeHost { nodes: &[NodeKey], ) -> Result { check_id(attempt)?; - let contexts = self + let wal = self .runtime_wal - .as_ref() - .ok_or_else(|| error("durable WAL required"))? - .operation_contexts()?; + .as_mut() + .ok_or_else(|| error("durable WAL required"))?; + wal.reconcile_operation_context_prefix()?; + let contexts = wal.operation_contexts()?; if contexts.observations.contains_key(attempt) { return Err(error("attempt observation is immutable")); } @@ -212,7 +239,8 @@ impl TrustedRuntimeHost { .ok_or_else(|| error("durable WAL required"))? .operation_contexts()? .observations - .remove(attempt) + .get(attempt) + .cloned() .ok_or_else(|| error("observation unavailable; re-observation requires a new attempt")) } @@ -292,11 +320,12 @@ impl TrustedRuntimeHost { invocation: EchoOperationInvocationV1, ) -> Result> { check_id(request)?; - let contexts = self + let wal = self .runtime_wal - .as_ref() - .ok_or_else(|| error("durable WAL required"))? - .operation_contexts()?; + .as_mut() + .ok_or_else(|| error("durable WAL required"))?; + wal.reconcile_operation_context_prefix()?; + let contexts = wal.operation_contexts()?; let observation = contexts .observations .get(attempt) 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 00452d14..49d239db 100644 --- a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs +++ b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs @@ -1332,6 +1332,165 @@ fn observed_operation_contexts_recover_without_replacing_original_inputs() { assert_eq!(recovered.runtime_wal().expect("WAL").commits().len(), 2); } +#[test] +fn healthy_contexts_advance_without_reconstructing_the_wal() { + let root = temp_runtime_wal_dir("context-steady-prefix"); + 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"); + let before = warp_core::causal_wal::filesystem_recovery_work(); + for i in 0..8 { + let attempt = format!("attempt-{i}"); + let observation = host + .retain_echo_operation_observation_v1(&attempt, head, &[node]) + .expect("capture"); + assert_eq!( + host.echo_operation_observation_v1(&attempt) + .expect("reading"), + observation + ); + let invocation = 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, b"value".to_vec(), vec![0xa0]); + let original = host + .bind_echo_operation_request_v1(&attempt, &attempt, invocation.clone()) + .expect("bind"); + assert_eq!( + host.bind_echo_operation_request_v1(&attempt, &attempt, invocation) + .expect("retry"), + original + ); + } + let after = warp_core::causal_wal::filesystem_recovery_work(); + assert_eq!( + after, before, + "healthy operations rescanned committed records" + ); +} + +#[test] +fn context_append_failure_reconciles_the_actual_committed_prefix() { + for target in [ + FilesystemWalFaultTarget::AppendFrame, + FilesystemWalFaultTarget::FlushCommit, + FilesystemWalFaultTarget::CommitMarkerSynced, + ] { + let root = temp_runtime_wal_dir("context-uncertain-prefix"); + 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.inject_runtime_wal_filesystem_fault_for_test(FilesystemWalFaultPlan::fail_next( + target, + )) + .expect("fault"); + assert!(host + .retain_echo_operation_observation_v1("uncertain", head, &[node]) + .is_err()); + let committed = target == FilesystemWalFaultTarget::CommitMarkerSynced; + assert_eq!( + host.echo_operation_observation_v1("uncertain").is_ok(), + committed + ); + let retry = host.retain_echo_operation_observation_v1("uncertain", head, &[node]); + assert_eq!( + retry.is_err(), + committed, + "durable observation cannot be replaced" + ); + host.retain_echo_operation_observation_v1("next", head, &[node]) + .expect("reconciled append"); + let report = host + .runtime_wal() + .expect("WAL") + .recover_read_only() + .expect("valid retained prefix"); + assert_eq!(report.certificate.committed_transactions_replayed, 2); + assert_eq!(host.runtime_wal().expect("WAL").commits().len(), 2); + } +} + +#[test] +fn durable_request_error_resolves_original_bytes_before_retrying() { + let root = temp_runtime_wal_dir("request-uncertain-prefix"); + 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("observation"); + 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]) + }; + host.inject_runtime_wal_filesystem_fault_for_test(FilesystemWalFaultPlan::fail_next( + FilesystemWalFaultTarget::CommitMarkerSynced, + )) + .expect("fault"); + assert!(host + .bind_echo_operation_request_v1("request", "attempt", invocation(b"original")) + .is_err()); + // The caller did not see success, but the semantic identity is already fixed. + assert!(host + .bind_echo_operation_request_v1("request", "attempt", invocation(b"replacement")) + .is_err()); + let scans = warp_core::causal_wal::filesystem_recovery_work(); + let bytes = host + .bind_echo_operation_request_v1("request", "attempt", invocation(b"original")) + .expect("resolve durable request"); + assert_eq!(scans, warp_core::causal_wal::filesystem_recovery_work()); + assert_eq!(host.runtime_wal().expect("WAL").commits().len(), 2); + drop(host); + let (mut recovered, _, _) = reopen(); + assert_eq!( + recovered + .bind_echo_operation_request_v1("request", "attempt", invocation(b"original")) + .expect("original after reopen"), + bytes + ); + assert!(recovered + .bind_echo_operation_request_v1("request", "attempt", invocation(b"replacement")) + .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 8f6419e0..b0a12c5a 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -232,8 +232,23 @@ 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. +The context index and writer cursor are derived from one validated committed WAL +prefix on opening or reconciliation. Healthy appends apply only their new context +records and advance the index's commit binding after durable acknowledgement. +Reads borrow that index without replaying history. Any append error invalidates +the binding, including an error reported after the commit marker was synced. +Reads in that posture resolve retained records; context writes reconcile the +writer cursor before checking identity or constructing another transaction. No +transaction from an uncertain cursor can reach storage. Exact retries recover +the original request rather than replacing its semantic input. +Retention policy and authenticated aperture delegation remain separate work. + +For diagnostic workload accounting, `ECHO_WAL_PROFILE=1` makes the session driver +emit cumulative thread-local recovery calls and decoded frame counts on stderr +after each RPC. Counters restart with each process. They count all filesystem +recovery calls, including required opening and explicit recovery, and do not +change stdout responses or retained evidence. They are work counts, not CPU-time +or physical-I/O measurements. Disconnected create-if-absent cells can leave the root-reachable state hash unchanged. Recovery evidence must therefore include native commit identities diff --git a/xtask/src/run_edict_operation/session.rs b/xtask/src/run_edict_operation/session.rs index 25d78993..14215187 100644 --- a/xtask/src/run_edict_operation/session.rs +++ b/xtask/src/run_edict_operation/session.rs @@ -148,6 +148,10 @@ pub fn serve(config: RunEdictOperationConfig) -> Result<()> { let response = result.unwrap_or_else(|error| json!({"error": format!("{error:#}")})); println!("{}", serde_json::to_string(&response)?); std::io::stdout().flush()?; + if std::env::var_os("ECHO_WAL_PROFILE").is_some() { + let (scans, frames) = warp_core::causal_wal::filesystem_recovery_work(); + eprintln!("{}", json!({"wal_profile":{"scans":scans,"frames":frames}})); + } } Ok(()) }