From e3645106a53393090a1ecc4b8f44f7e0d9eb9737 Mon Sep 17 00:00:00 2001 From: James Ross Date: Mon, 21 Sep 2026 01:24:04 -0700 Subject: [PATCH 1/3] feat(runtime): retain native operation strands across host recovery --- CHANGELOG.md | 5 + crates/warp-core/src/trusted_runtime_host.rs | 28 ++- .../trusted_runtime_host/retained_strands.rs | 209 ++++++++++++++++++ .../tests/trusted_runtime_host_loop_tests.rs | 60 +++++ .../application-contract-hosting.md | 32 +++ xtask/src/main.rs | 4 + xtask/src/run_edict_operation.rs | 1 + xtask/src/run_edict_operation/session.rs | 72 +++++- 8 files changed, 405 insertions(+), 6 deletions(-) create mode 100644 crates/warp-core/src/trusted_runtime_host/retained_strands.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 6a0cb138..c795c914 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ ### Added +- The trusted executable-operation host can retain two native alternative strands, + continue a losing strand under new observations, and recover its fork ancestry + and operation outcomes from the WAL. Application selection records do not + imply merge or promotion. + - 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 diff --git a/crates/warp-core/src/trusted_runtime_host.rs b/crates/warp-core/src/trusted_runtime_host.rs index 93c66ad0..e6681b62 100644 --- a/crates/warp-core/src/trusted_runtime_host.rs +++ b/crates/warp-core/src/trusted_runtime_host.rs @@ -15,6 +15,7 @@ use std::{ use thiserror::Error; mod observed_context; +mod retained_strands; pub use observed_context::EchoOperationContextErrorV1; use crate::causal_anchor::prepare_causal_anchor_admission; @@ -1318,24 +1319,30 @@ impl TrustedRuntimeHost { config: TrustedRuntimeWalConfig, ) -> Result<(), TrustedRuntimeHostError> { let runtime_wal = TrustedRuntimeWal::from_config(config)?; - let recovery = runtime_wal.recover_read_only()?; + let mut recovery = runtime_wal.recover_read_only()?; if let Some(submission_id) = recovery.missing_submission_envelopes.first().copied() { return Err(TrustedRuntimeWalError::SubmissionEnvelopeMissing { submission_id }.into()); } if let Some(receipt_digest) = recovery.missing_runtime_state_deltas.first().copied() { return Err(TrustedRuntimeWalError::RuntimeStateDeltaMissing { receipt_digest }.into()); } - ensure_runtime_authority_is_durable( + let (forked_runtime, forked_provenance) = retained_strands::restore( &self.runtime, &self.provenance, + &runtime_wal, + &mut recovery, + )?; + ensure_runtime_authority_is_durable( + &forked_runtime, + &forked_provenance, &self.engine, &recovery, )?; self.engine .preflight_recovered_echo_operation_packages_v1(&recovery.installed_echo_operations)?; - let mut restored_runtime = self.runtime.clone(); - let mut restored_provenance = self.provenance.clone(); + let mut restored_runtime = forked_runtime; + let mut restored_provenance = forked_provenance; restored_runtime .restore_witnessed_submission_persistence(recovery.witnessed_submissions)?; restore_provenance_entries(&mut restored_provenance, &recovery.provenance_entries)?; @@ -4742,7 +4749,9 @@ fn evaluation_basis_matches_recovered_coordinate( crate::WorldlineTick::from_raw(parent_tick), )) .is_some_and(|parent| { - parent.head_key == Some(basis.writer_head()) + // A native fork preserves the source writer's head in its copied + // prefix. The next writer need not have authored the parent commit. + parent.worldline_id == basis.writer_head().worldline_id && parent.expected.state_root == basis.state_root() && parent.expected.commit_hash == basis.commit_id() && basis.commit_global_tick() == Some(parent.commit_global_tick) @@ -6936,6 +6945,15 @@ mod tests { )); let mut wrong_root = parent.clone(); + let mut other_writer = parent.clone(); + other_writer.head_key = Some(crate::WriterHeadKey { + worldline_id: parent.worldline_id, + head_id: crate::make_head_id("original-fork-writer"), + }); + assert!(evaluation_basis_matches_recovered_coordinate( + basis, + &BTreeMap::from([(parent_coordinate, &other_writer)]) + )); wrong_root.expected.state_root = [39; 32]; let provenance = BTreeMap::from([(parent_coordinate, &wrong_root)]); assert!(validate_operation_receipt_parent_material(basis, &child, &provenance).is_err()); diff --git a/crates/warp-core/src/trusted_runtime_host/retained_strands.rs b/crates/warp-core/src/trusted_runtime_host/retained_strands.rs new file mode 100644 index 00000000..d095fece --- /dev/null +++ b/crates/warp-core/src/trusted_runtime_host/retained_strands.rs @@ -0,0 +1,209 @@ +// SPDX-License-Identifier: Apache-2.0 +// © James Ross Ω FLYING•ROBOTS +//! Bounded trusted-local native strand retention and recovery. +use super::{ + restore_provenance_entries, Hash, ProvenanceService, ProvenanceStore, RuntimeWalActivationGap, + TrustedRuntimeHost, TrustedRuntimeHostError, TrustedRuntimeWal, TrustedRuntimeWalError, + TrustedRuntimeWalRecovery, WalAppendAuthority, WalTransactionId, WalTransactionKind, + WorldlineRuntime, +}; +use crate::causal_wal::{StrandForkRecord, WalRecordKind}; +use crate::playback::SessionId; +use crate::{ + ActorId, AuthorityBinding, AuthorityDomainId, AuthorityDomainRef, CausalPosture, + ForkStrandRequest, InboxPolicy, OriginId, PlaybackMode, RetentionContractId, SealStrength, + SessionContext, StrandId, WorldlineTick, WriterHead, WriterHeadKey, +}; + +fn invalid() -> TrustedRuntimeHostError { + TrustedRuntimeWalError::RuntimeAuthorityNotDurable { + gap: RuntimeWalActivationGap::Provenance, + } + .into() +} +fn digest(domain: &[u8], identity: &[u8]) -> Hash { + let mut hash = blake3::Hasher::new(); + hash.update(domain); + hash.update(identity); + *hash.finalize().as_bytes() +} +fn request(record: &StrandForkRecord) -> Result { + let identity = *record.strand_id.as_bytes(); + let origin = OriginId::from_bytes(identity); + let session = SessionContext::new( + SessionId(identity), + origin, + ActorId::from_bytes(identity), + AuthorityDomainRef::new(origin, AuthorityDomainId::from_bytes(identity)), + AuthorityBinding::LocalUnbound { origin }, + SealStrength::Advisory, + CausalPosture::AuthorOnly, + None, + RetentionContractId::from_bytes(identity), + ) + .map_err(|_| invalid())?; + if record.writer_heads.len() != 1 + || record.writer_heads[0].worldline_id != record.child_worldline_id + || record.child_worldline_id.as_bytes() != &digest(b"echo:local-strand-child:v1", &identity) + || record.writer_heads[0].head_id.as_bytes() + != &digest(b"echo:local-strand-head:v1", &identity) + || record.topology_intent_id != identity + || record.idempotency_key_digest != Some(identity) + || record.retention_posture_digest != digest(b"echo:local-author-only-strand:v1", &identity) + || record.issuer_evidence_digest != identity + { + return Err(invalid()); + } + ForkStrandRequest::from_session_default( + record.strand_id, + record.source_worldline_id, + record.fork_tick, + record.child_worldline_id, + vec![WriterHead::with_routing( + record.writer_heads[0], + PlaybackMode::Play, + InboxPolicy::AcceptAll, + None, + true, + )], + &session, + ) + .map_err(|_| invalid()) +} + +impl TrustedRuntimeHost { + /// Forks a retained AuthorOnly strand under the trusted local host's fixed + /// advisory authority profile. This is not an authenticated admission API. + /// Reusing an existing strand identity returns its original writer head. + /// + /// # Errors + /// Refuses missing WAL/source history, invalid labels, or more than 64 strands. + pub fn fork_local_operation_strand_v1( + &mut self, + source: WriterHeadKey, + label: &str, + ) -> Result { + if label.is_empty() || label.len() > 64 || self.runtime.heads().get(&source).is_none() { + return Err(invalid()); + } + let identity = digest(source.worldline_id.as_bytes(), label.as_bytes()); + let strand_id = StrandId::from_bytes(identity); + if let Some(strand) = self.runtime.strands().get(&strand_id) { + return strand.writer_heads().first().copied().ok_or_else(invalid); + } + if self.runtime.strands().len() >= 64 { + return Err(invalid()); + } + let tip = self + .provenance + .tip_ref(source.worldline_id)? + .ok_or_else(invalid)?; + let entry = self + .provenance + .entry(source.worldline_id, tip.worldline_tick)?; + let head = WriterHeadKey { + worldline_id: crate::WorldlineId::from_bytes(digest( + b"echo:local-strand-child:v1", + &identity, + )), + head_id: crate::head::HeadId::from_bytes(digest( + b"echo:local-strand-head:v1", + &identity, + )), + }; + let record = StrandForkRecord { + topology_intent_id: identity, + strand_id, + source_worldline_id: source.worldline_id, + fork_tick: tip.worldline_tick, + source_commit_hash: tip.commit_hash, + source_boundary_hash: entry.expected.state_root, + child_worldline_id: head.worldline_id, + writer_heads: vec![head], + retention_posture_digest: digest(b"echo:local-author-only-strand:v1", &identity), + issuer_evidence_digest: identity, + idempotency_key_digest: Some(identity), + }; + let mut runtime = self.runtime.clone(); + let mut provenance = self.provenance.clone(); + runtime.fork_strand(&mut provenance, request(&record)?)?; + let wal = self.runtime_wal.as_mut().ok_or_else(invalid)?; + wal.refresh_cursor_from_store_for_writer()?; + let mut builder = wal.builder( + WalTransactionKind::TopologyIntent, + WalAppendAuthority::TrustedScheduler, + WalTransactionId::from_hash(identity), + ); + builder + .push_record( + WalRecordKind::TopologyStrandForkRecorded, + record.to_payload_bytes(), + ) + .map_err(TrustedRuntimeWalError::from)?; + wal.append_transaction( + builder + .commit(Vec::new()) + .map_err(TrustedRuntimeWalError::from)?, + )?; + self.runtime = runtime; + self.provenance = provenance; + Ok(head) + } +} + +pub(super) fn restore( + runtime: &WorldlineRuntime, + provenance: &ProvenanceService, + wal: &TrustedRuntimeWal, + recovery: &mut TrustedRuntimeWalRecovery, +) -> Result<(WorldlineRuntime, ProvenanceService), TrustedRuntimeHostError> { + let mut runtime = runtime.clone(); + let mut provenance = provenance.clone(); + let scan = wal + .store + .recover_read_only() + .map_err(TrustedRuntimeWalError::from)?; + let mut count = 0; + for transaction in scan.transactions { + for frame in transaction.frames { + if frame.header.record_kind != WalRecordKind::TopologyStrandForkRecorded { + continue; + } + count += 1; + if count > 64 { + return Err(invalid()); + } + let record = StrandForkRecord::from_payload_bytes(&frame.payload.canonical_bytes) + .map_err(|_| invalid())?; + let fork_request = request(&record)?; + if runtime.strands().contains(&record.strand_id) { + return Err(invalid()); + } + let source_entries = recovery + .provenance_entries + .iter() + .filter(|entry| entry.worldline_id == record.source_worldline_id) + .cloned() + .collect::>(); + restore_provenance_entries(&mut provenance, &source_entries)?; + runtime.restore_causal_runtime_history(&provenance, &source_entries, &[])?; + let receipt = runtime.fork_strand(&mut provenance, fork_request)?; + if receipt.fork_basis_ref.commit_hash != record.source_commit_hash + || receipt.fork_basis_ref.boundary_hash != record.source_boundary_hash + { + return Err(invalid()); + } + // The copied prefix is derived from the validated native fork and + // retained source entries, not reconstructed from application bytes. + for tick in 0..=record.fork_tick.as_u64() { + recovery.provenance_entries.push( + provenance.entry(record.child_worldline_id, WorldlineTick::from_raw(tick))?, + ); + } + } + } + recovery + .provenance_entries + .sort_by_key(|entry| (entry.worldline_id, entry.worldline_tick)); + Ok((runtime, provenance)) +} 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 52032937..634988ce 100644 --- a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs +++ b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs @@ -2450,3 +2450,63 @@ fn filesystem_causal_anchor_flush_failure_publishes_no_admission() { drop(host); fs::remove_dir_all(&wal_root).expect("failed-anchor WAL fixture should be removable"); } + +#[test] +fn retained_native_fork_recovers_prefix_and_refuses_missing_source() { + let root = temp_runtime_wal_dir("retained-native-fork"); + let (rt, lane) = runtime(); + let head = rt.heads().iter().next().expect("head").0.to_owned(); + let mut host = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); + host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + .expect("wal"); + assert!(host + .fork_local_operation_strand_v1(head, "candidate") + .is_err()); + host.register_contract_package(package()).expect("package"); + let submission = host + .app() + .submit_intent_with_runtime_wal_ack(eint_envelope(lane)) + .expect("submit"); + host.stage_installed_contract_submission(submission.submission_id, &admission_ticket(17)) + .expect("stage"); + host.run_until_idle(4).expect("commit"); + let child = host + .fork_local_operation_strand_v1(head, "candidate") + .expect("native fork"); + let fork = host + .runtime() + .strands() + .find_by_child_worldline(&child.worldline_id) + .expect("strand") + .fork_basis_ref(); + let before = host.runtime_wal().expect("wal").commits().len(); + assert_eq!( + host.fork_local_operation_strand_v1(head, "candidate") + .expect("retry"), + child + ); + assert_eq!(host.runtime_wal().expect("wal").commits().len(), before); + assert!(host.fork_local_operation_strand_v1(head, "").is_err()); + drop(host); + let (rt, _) = runtime(); + let mut restored = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); + restored + .enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + .expect("native topology replay"); + let recovered = restored + .runtime() + .strands() + .find_by_child_worldline(&child.worldline_id) + .expect("retained strand"); + assert_eq!(recovered.fork_basis_ref(), fork); + assert_eq!(recovered.writer_heads(), &[child]); + assert_eq!( + restored + .provenance() + .entry(child.worldline_id, fork.fork_tick) + .expect("prefix") + .expected + .commit_hash, + fork.commit_hash + ); +} diff --git a/docs/architecture/application-contract-hosting.md b/docs/architecture/application-contract-hosting.md index aaf5fe45..e6ad52a3 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -1201,3 +1201,35 @@ successful settlement aperture to equal the admitted request, revalidates the registry grant before settlement, and retains replay bytes before deterministic resumption. This profile adds no application noun, native application callback, provider import, shell, or ambient filesystem capability. + +### Retained alternatives in the trusted local host + +`run-edict-operation --serve --retained-alternatives` starts the parent with one +retained Edict bootstrap operation and forks two native AuthorOnly strands, `a` +and `b`, from that committed basis. `use` changes the selected lane for bounded +observations and operation submission; `ancestry` returns up to 4,096 native +commit identities for that lane. It does not change canonicality or settle work. + +`TrustedRuntimeHost::fork_local_operation_strand_v1` persists Echo's existing +`TopologyStrandForkRecorded` record before exposing the fork. The fixed local +profile derives an advisory, locally unbound AuthorOnly authority and retention +identity from the strand identity. It allows one writer head per strand, labels +of 1–64 bytes and at most 64 retained forks. This is an explicit trusted-host +profile, not caller authentication or a generic authority-policy authoring API. +Reusing a source-worldline/label pair resolves its original retained strand. + +Recovery validates the supported profile and source commit/boundary, replays the +source history, invokes native fork construction at the recorded coordinate, and +then restores child history and Action dispositions. Copied prefix entries are +derived from retained source history and the fork record; application output is +never used to recreate ancestry. The next writer may differ from the writer of +the retained parent commit. Recovery corroborates the parent coordinate, commit, +state root and global tick without requiring those writer identities to match. + +Selection remains an application convention expressed by an Edict operation +that records a chosen worldline and native commit. This driver does not perform +merge, promotion, braid settlement, or cross-worldline atomic validation of a +selection policy. New information on a strand is a new explicit reading; it does +not replace observations bound to an earlier attempt. Recovering a retained +alternative does not establish an efficiency advantage over retained branches +and ordinary artifact retrieval. diff --git a/xtask/src/main.rs b/xtask/src/main.rs index d4066bb7..fbc05179 100644 --- a/xtask/src/main.rs +++ b/xtask/src/main.rs @@ -159,6 +159,9 @@ struct RuntimeCounterDiagnosticArgs { #[derive(Args)] struct RunEdictOperationArgs { + /// Enable two retained native AuthorOnly strands in the trusted host. + #[arg(long, requires = "serve")] + retained_alternatives: bool, /// Serve bounded JSON requests against one persistent worldline. #[arg(long)] serve: bool, @@ -486,6 +489,7 @@ fn main() -> Result<()> { fn run_edict_operation(args: RunEdictOperationArgs) -> Result<()> { let config = run_edict_operation::RunEdictOperationConfig { + retained_alternatives: args.retained_alternatives, package: args.package, verification_report: args.verification_report, lawpack_manifest: args.lawpack_manifest, diff --git a/xtask/src/run_edict_operation.rs b/xtask/src/run_edict_operation.rs index 6b8d7175..6f2e6881 100644 --- a/xtask/src/run_edict_operation.rs +++ b/xtask/src/run_edict_operation.rs @@ -49,6 +49,7 @@ const WARP_ID_SOURCE: &str = "action-lane/v1"; /// Inputs needed to run one exact compiler-produced package. pub struct RunEdictOperationConfig { + pub retained_alternatives: bool, pub package: PathBuf, pub verification_report: PathBuf, pub lawpack_manifest: PathBuf, diff --git a/xtask/src/run_edict_operation/session.rs b/xtask/src/run_edict_operation/session.rs index 92bb93bd..9f9b7abe 100644 --- a/xtask/src/run_edict_operation/session.rs +++ b/xtask/src/run_edict_operation/session.rs @@ -16,6 +16,7 @@ use serde_json::{json, Value}; use std::io::{BufRead, Read, Write}; struct Session { + heads: std::collections::BTreeMap, fixture: HostFixture, package: PackageMetadata, package_id: warp_core::EchoOperationPackageIdV1, @@ -84,6 +85,7 @@ pub fn serve(config: RunEdictOperationConfig) -> Result<()> { ), ); let mut session = Session { + heads: std::collections::BTreeMap::from([("parent".to_owned(), fixture.head)]), fixture, package, package_id, @@ -91,6 +93,35 @@ pub fn serve(config: RunEdictOperationConfig) -> Result<()> { configuration, basis: input.basis, }; + if config.retained_alternatives { + // The initial application operation is retained once. Recovery resolves + // its original binding; it never re-executes it to rebuild history. + if session + .fixture + .host + .echo_operation_observation_v1("host-bootstrap") + .is_err() + { + let result = session + .call(&json!({"op":"observe","attempt":"host-bootstrap","keys":["bootstrap"]}))?; + if result.get("error").is_some() { + bail!("bootstrap observation failed"); + } + } + let result = session.call( + &json!({"op":"submit","attempt":"host-bootstrap","key":"bootstrap","value":""}), + )?; + if result["outcome"] != "committed" { + bail!("bootstrap operation failed: {result}"); + } + for label in ["a", "b"] { + let head = session + .fixture + .host + .fork_local_operation_strand_v1(session.fixture.head, label)?; + session.heads.insert(label.to_owned(), head); + } + } let mut input = std::io::stdin().lock(); loop { let mut line = Vec::new(); @@ -132,7 +163,8 @@ impl Session { fn call(&mut self, request: &Value) -> Result { let allowed: &[&str] = match text(request, "op")? { - "status" => &["op"], + "status" | "ancestry" => &["op"], + "use" => &["op", "lane"], "observe" => &["op", "attempt", "keys"], "changes" | "reading" => &["op", "attempt"], "submit" => &["op", "attempt", "key", "value", "request_id"], @@ -148,6 +180,30 @@ impl Session { bail!("unexpected request field; observation bindings are runtime-owned"); } match text(request, "op")? { + "use" => { + self.fixture.head = *self + .heads + .get(text(request, "lane")?) + .context("unknown lane")?; + self.call(&json!({"op":"status"})) + } + "ancestry" => { + use warp_core::ProvenanceStore; + let provenance = self.fixture.host.provenance(); + let lane = self.fixture.head.worldline_id; + let len = provenance.len(lane)?; + if len > 4096 { + bail!("ancestry aperture exceeded"); + } + let commits = (0..len) + .map(|tick| { + provenance + .entry(lane, warp_core::WorldlineTick::from_raw(tick)) + .map(|entry| hex::encode(entry.expected.commit_hash)) + }) + .collect::, _>>()?; + Ok(json!({"worldline":hex::encode(lane.as_bytes()),"commits":commits})) + } "status" => { let application_basis = echo_operation_anchored_node_creation_application_basis_v1( self.fixture.node, @@ -157,7 +213,21 @@ impl Session { .fixture .host .echo_operation_evaluation_basis_v1(self.fixture.head, application_basis)?; + let strand = self + .fixture + .host + .runtime() + .strands() + .find_by_child_worldline(&self.fixture.head.worldline_id); + let fork = strand.map(|strand| { + let basis = strand.fork_basis_ref(); + json!({ + "source_worldline":hex::encode(basis.source_lane_id.as_bytes()), + "tick":basis.fork_tick.as_u64(),"commit":hex::encode(basis.commit_hash) + }) + }); Ok(json!({ + "fork": fork, "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()), From 24e221b71c9cb33c9e88682cd3408885aa500eb2 Mon Sep 17 00:00:00 2001 From: James Ross Date: Thu, 8 Oct 2026 00:53:38 -0700 Subject: [PATCH 2/3] test: reproduce duplicate retained fork after uncertain commit --- .../tests/trusted_runtime_host_loop_tests.rs | 40 +++++++++++++++++++ 1 file changed, 40 insertions(+) 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 49d239db..fa8b3b6a 100644 --- a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs +++ b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs @@ -2797,3 +2797,43 @@ fn retained_native_fork_recovers_prefix_and_refuses_missing_source() { fork.commit_hash ); } + +#[test] +fn retained_fork_retry_reconciles_uncertain_commit_before_reopening() { + for target in [FilesystemWalFaultTarget::CommitMarkerSynced, + FilesystemWalFaultTarget::AppendFrame, FilesystemWalFaultTarget::FlushCommit] { + let root = temp_runtime_wal_dir("retained-fork-uncertain"); + let (rt, lane) = runtime(); + let head = *rt.heads().iter().next().expect("head").0; + let mut host = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); + host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)).expect("wal"); + host.register_contract_package(package()).expect("package"); + let submission = host.app().submit_intent_with_runtime_wal_ack(eint_envelope(lane)) + .expect("submit"); + host.stage_installed_contract_submission(submission.submission_id, &admission_ticket(17)) + .expect("stage"); + host.run_until_idle(4).expect("source commit"); + let before = host.runtime_wal().expect("wal").commits().len(); + host.inject_runtime_wal_filesystem_fault_for_test(FilesystemWalFaultPlan::fail_next(target)) + .expect("fault"); + assert!(host.fork_local_operation_strand_v1(head, "candidate").is_err()); + let child = host.fork_local_operation_strand_v1(head, "candidate") + .expect("retry original retained fork"); + assert_eq!(host.runtime_wal().expect("wal").commits().len(), before + 1, + "retry must retain exactly one fork transaction"); + let fork = host.runtime().strands().find_by_child_worldline(&child.worldline_id) + .expect("fork").fork_basis_ref(); + drop(host); + let (rt, _) = runtime(); + let mut restored = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); + restored.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + .expect("one retained fork reopens"); + assert_eq!(restored.fork_local_operation_strand_v1(head, "candidate") + .expect("reopened retry"), child); + assert_eq!(restored.runtime().strands().find_by_child_worldline(&child.worldline_id) + .expect("retained fork").fork_basis_ref(), fork); + assert_eq!(restored.runtime_wal().expect("wal").commits().len(), before + 1); + drop(restored); + fs::remove_dir_all(root).expect("owned fixture cleanup"); + } +} From 0fb4b8e2d9d8aaeb7942a7aa96602f9d2f7e28ba Mon Sep 17 00:00:00 2001 From: James Ross Date: Thu, 8 Oct 2026 01:31:36 -0700 Subject: [PATCH 3/3] fix: reconcile retained forks before retries and capacity admission --- CHANGELOG.md | 2 + .../trusted_runtime_host/retained_strands.rs | 90 ++++++- .../tests/trusted_runtime_host_loop_tests.rs | 239 ++++++++++++++++-- .../application-contract-hosting.md | 8 + 4 files changed, 312 insertions(+), 27 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1153929b..5417b9d8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -27,6 +27,8 @@ ### Fixed +- Retained local fork retries reconcile durable topology after uncertain appends, preserving the original basis and preventing duplicate fork records. + - Observation slots preflight their canonical byte bound and charge execution reads before copying Atom payloads. - Private Edict runtime decoding resolves named integer types through their declared width, matching source-function provider acceptance. diff --git a/crates/warp-core/src/trusted_runtime_host/retained_strands.rs b/crates/warp-core/src/trusted_runtime_host/retained_strands.rs index d095fece..fd686c3c 100644 --- a/crates/warp-core/src/trusted_runtime_host/retained_strands.rs +++ b/crates/warp-core/src/trusted_runtime_host/retained_strands.rs @@ -71,6 +71,79 @@ fn request(record: &StrandForkRecord) -> Result Result { + let receipt = runtime.fork_strand(provenance, request(record)?)?; + if receipt.fork_basis_ref.commit_hash != record.source_commit_hash + || receipt.fork_basis_ref.boundary_hash != record.source_boundary_hash + { + return Err(invalid()); + } + record.writer_heads.first().copied().ok_or_else(invalid) +} + +// Recover committed topology before counting capacity or acknowledging retries. +// All candidate changes stay private until every retained record validates. +fn reconcile_pending_forks( + host: &mut TrustedRuntimeHost, +) -> Result, TrustedRuntimeHostError> { + let wal = host.runtime_wal.as_mut().ok_or_else(invalid)?; + wal.refresh_cursor_from_store_for_writer()?; + let scan = wal + .store + .recover_read_only() + .map_err(TrustedRuntimeWalError::from)?; + let mut seen = std::collections::BTreeSet::new(); + let mut pending = Vec::new(); + for transaction in scan.transactions { + for frame in transaction.frames { + if frame.header.record_kind != WalRecordKind::TopologyStrandForkRecorded { + continue; + } + let record = StrandForkRecord::from_payload_bytes(&frame.payload.canonical_bytes) + .map_err(|_| invalid())?; + if !seen.insert(record.strand_id) || seen.len() > 64 { + return Err(invalid()); + } + request(&record)?; + if let Some(strand) = host.runtime.strands().get(&record.strand_id) { + let basis = strand.fork_basis_ref(); + let entry = host + .provenance + .entry(record.source_worldline_id, record.fork_tick)?; + // ForkBasisRef uses legacy lane terminology for this worldline. + let expected_source = record.source_worldline_id; + if basis.source_lane_id != expected_source + || basis.fork_tick != record.fork_tick + || basis.commit_hash != record.source_commit_hash + || basis.boundary_hash != record.source_boundary_hash + || strand.child_worldline_id() != record.child_worldline_id + || strand.writer_heads() != record.writer_heads.as_slice() + || entry.expected.commit_hash != record.source_commit_hash + || entry.expected.state_root != record.source_boundary_hash + { + return Err(invalid()); + } + } else { + pending.push(record); + } + } + } + if !pending.is_empty() { + let mut runtime = host.runtime.clone(); + let mut provenance = host.provenance.clone(); + for record in pending { + apply_record(&mut runtime, &mut provenance, &record)?; + } + host.runtime = runtime; + host.provenance = provenance; + } + Ok(seen) +} + impl TrustedRuntimeHost { /// Forks a retained AuthorOnly strand under the trusted local host's fixed /// advisory authority profile. This is not an authenticated admission API. @@ -88,7 +161,13 @@ impl TrustedRuntimeHost { } let identity = digest(source.worldline_id.as_bytes(), label.as_bytes()); let strand_id = StrandId::from_bytes(identity); + let retained = reconcile_pending_forks(self)?; if let Some(strand) = self.runtime.strands().get(&strand_id) { + if !retained.contains(&strand_id) + || strand.fork_basis_ref().source_lane_id != source.worldline_id + { + return Err(invalid()); + } return strand.writer_heads().first().copied().ok_or_else(invalid); } if self.runtime.strands().len() >= 64 { @@ -126,9 +205,8 @@ impl TrustedRuntimeHost { }; let mut runtime = self.runtime.clone(); let mut provenance = self.provenance.clone(); - runtime.fork_strand(&mut provenance, request(&record)?)?; + apply_record(&mut runtime, &mut provenance, &record)?; let wal = self.runtime_wal.as_mut().ok_or_else(invalid)?; - wal.refresh_cursor_from_store_for_writer()?; let mut builder = wal.builder( WalTransactionKind::TopologyIntent, WalAppendAuthority::TrustedScheduler, @@ -175,7 +253,6 @@ pub(super) fn restore( } let record = StrandForkRecord::from_payload_bytes(&frame.payload.canonical_bytes) .map_err(|_| invalid())?; - let fork_request = request(&record)?; if runtime.strands().contains(&record.strand_id) { return Err(invalid()); } @@ -187,12 +264,7 @@ pub(super) fn restore( .collect::>(); restore_provenance_entries(&mut provenance, &source_entries)?; runtime.restore_causal_runtime_history(&provenance, &source_entries, &[])?; - let receipt = runtime.fork_strand(&mut provenance, fork_request)?; - if receipt.fork_basis_ref.commit_hash != record.source_commit_hash - || receipt.fork_basis_ref.boundary_hash != record.source_boundary_hash - { - return Err(invalid()); - } + apply_record(&mut runtime, &mut provenance, &record)?; // The copied prefix is derived from the validated native fork and // retained source entries, not reconstructed from application bytes. for tick in 0..=record.fork_tick.as_u64() { 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 fa8b3b6a..f8c68bfb 100644 --- a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs +++ b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs @@ -2800,40 +2800,243 @@ fn retained_native_fork_recovers_prefix_and_refuses_missing_source() { #[test] fn retained_fork_retry_reconciles_uncertain_commit_before_reopening() { - for target in [FilesystemWalFaultTarget::CommitMarkerSynced, - FilesystemWalFaultTarget::AppendFrame, FilesystemWalFaultTarget::FlushCommit] { + for target in [ + FilesystemWalFaultTarget::CommitMarkerSynced, + FilesystemWalFaultTarget::AppendFrame, + FilesystemWalFaultTarget::FlushCommit, + ] { let root = temp_runtime_wal_dir("retained-fork-uncertain"); let (rt, lane) = runtime(); let head = *rt.heads().iter().next().expect("head").0; let mut host = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); - host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)).expect("wal"); + host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + .expect("wal"); host.register_contract_package(package()).expect("package"); - let submission = host.app().submit_intent_with_runtime_wal_ack(eint_envelope(lane)) + let submission = host + .app() + .submit_intent_with_runtime_wal_ack(eint_envelope(lane)) .expect("submit"); host.stage_installed_contract_submission(submission.submission_id, &admission_ticket(17)) .expect("stage"); host.run_until_idle(4).expect("source commit"); let before = host.runtime_wal().expect("wal").commits().len(); - host.inject_runtime_wal_filesystem_fault_for_test(FilesystemWalFaultPlan::fail_next(target)) - .expect("fault"); - assert!(host.fork_local_operation_strand_v1(head, "candidate").is_err()); - let child = host.fork_local_operation_strand_v1(head, "candidate") + host.inject_runtime_wal_filesystem_fault_for_test(FilesystemWalFaultPlan::fail_next( + target, + )) + .expect("fault"); + assert!(host + .fork_local_operation_strand_v1(head, "candidate") + .is_err()); + let child = host + .fork_local_operation_strand_v1(head, "candidate") .expect("retry original retained fork"); - assert_eq!(host.runtime_wal().expect("wal").commits().len(), before + 1, - "retry must retain exactly one fork transaction"); - let fork = host.runtime().strands().find_by_child_worldline(&child.worldline_id) - .expect("fork").fork_basis_ref(); + assert_eq!( + host.runtime_wal().expect("wal").commits().len(), + before + 1, + "retry must retain exactly one fork transaction" + ); + let fork = host + .runtime() + .strands() + .find_by_child_worldline(&child.worldline_id) + .expect("fork") + .fork_basis_ref(); drop(host); let (rt, _) = runtime(); let mut restored = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); - restored.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + restored + .enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) .expect("one retained fork reopens"); - assert_eq!(restored.fork_local_operation_strand_v1(head, "candidate") - .expect("reopened retry"), child); - assert_eq!(restored.runtime().strands().find_by_child_worldline(&child.worldline_id) - .expect("retained fork").fork_basis_ref(), fork); - assert_eq!(restored.runtime_wal().expect("wal").commits().len(), before + 1); + assert_eq!( + restored + .fork_local_operation_strand_v1(head, "candidate") + .expect("reopened retry"), + child + ); + assert_eq!( + restored + .runtime() + .strands() + .find_by_child_worldline(&child.worldline_id) + .expect("retained fork") + .fork_basis_ref(), + fork + ); + assert_eq!( + restored.runtime_wal().expect("wal").commits().len(), + before + 1 + ); drop(restored); fs::remove_dir_all(root).expect("owned fixture cleanup"); } } + +fn retained_fork_source_fixture( + root: &std::path::Path, +) -> (TrustedRuntimeHost, WriterHeadKey, WorldlineId) { + let (rt, lane) = runtime(); + let head = *rt.heads().iter().next().expect("head").0; + let mut host = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); + host.enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(root)) + .expect("wal"); + host.register_contract_package(package()).expect("package"); + let submission = host + .app() + .submit_intent_with_runtime_wal_ack(eint_envelope(lane)) + .expect("submit"); + host.stage_installed_contract_submission(submission.submission_id, &admission_ticket(17)) + .expect("stage"); + host.run_until_idle(4).expect("source commit"); + (host, head, lane) +} + +#[test] +fn retained_fork_pending_basis_survives_parent_advance_and_different_label() { + let root = temp_runtime_wal_dir("retained-fork-parent-advance"); + let (mut host, head, lane) = retained_fork_source_fixture(&root); + let original = host + .provenance() + .tip_ref(lane) + .expect("tip") + .expect("source"); + host.inject_runtime_wal_filesystem_fault_for_test(FilesystemWalFaultPlan::fail_next( + FilesystemWalFaultTarget::CommitMarkerSynced, + )) + .expect("fault"); + assert!(host + .fork_local_operation_strand_v1(head, "original") + .is_err()); + // A context write reconciles append permission without exposing topology. + let node = *host + .runtime() + .worldlines() + .get(&lane) + .expect("lane") + .state() + .root(); + host.retain_echo_operation_observation_v1("parent-write-permission", head, &[node]) + .expect("writer prefix reconciled"); + assert_eq!(host.runtime().strands().len(), 0); + let parent_receipt = host + .runtime() + .receipt_correlations() + .next() + .expect("source receipt") + .causal_receipt_ref; + let submission = host + .app() + .submit_intent_with_runtime_wal_ack(causal_eint_envelope(lane, vec![parent_receipt])) + .expect("advance submit"); + host.stage_installed_contract_submission(submission.submission_id, &admission_ticket(18)) + .expect("advance stage"); + host.run_until_idle(4).expect("parent advance"); + let latest = host + .provenance() + .tip_ref(lane) + .expect("tip") + .expect("advanced source"); + assert!(latest.worldline_tick > original.worldline_tick); + let before = host.runtime_wal().expect("wal").commits().len(); + let fresh = host + .fork_local_operation_strand_v1(head, "fresh") + .expect("different label"); + let old = host + .fork_local_operation_strand_v1(head, "original") + .expect("original retry"); + assert_ne!(old, fresh); + assert_eq!(host.runtime().strands().len(), 2); + assert_eq!(host.runtime_wal().expect("wal").commits().len(), before + 1); + let basis = host + .runtime() + .strands() + .find_by_child_worldline(&old.worldline_id) + .expect("old fork") + .fork_basis_ref(); + assert_eq!( + (basis.fork_tick, basis.commit_hash), + (original.worldline_tick, original.commit_hash) + ); + assert_eq!( + host.runtime() + .strands() + .find_by_child_worldline(&fresh.worldline_id) + .expect("fresh fork") + .fork_basis_ref() + .fork_tick, + latest.worldline_tick + ); + drop(host); + let (rt, _) = runtime(); + let mut restored = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); + restored + .enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + .expect("both forks reopen"); + assert_eq!( + restored + .fork_local_operation_strand_v1(head, "original") + .expect("original"), + old + ); + assert_eq!( + restored + .fork_local_operation_strand_v1(head, "fresh") + .expect("fresh"), + fresh + ); + assert_eq!( + restored + .runtime() + .strands() + .find_by_child_worldline(&old.worldline_id) + .expect("old") + .fork_basis_ref(), + basis + ); + drop(restored); + fs::remove_dir_all(root).expect("owned fixture cleanup"); +} + +#[test] +fn retained_fork_pending_commit_counts_against_capacity_before_new_label() { + let root = temp_runtime_wal_dir("retained-fork-capacity"); + let (mut host, head, _) = retained_fork_source_fixture(&root); + for index in 0..63 { + host.fork_local_operation_strand_v1(head, &format!("visible-{index}")) + .expect("visible fork"); + } + let before = host.runtime_wal().expect("wal").commits().len(); + host.inject_runtime_wal_filesystem_fault_for_test(FilesystemWalFaultPlan::fail_next( + FilesystemWalFaultTarget::CommitMarkerSynced, + )) + .expect("fault"); + assert!(host + .fork_local_operation_strand_v1(head, "pending") + .is_err()); + assert_eq!(host.runtime().strands().len(), 63); + assert!(host + .fork_local_operation_strand_v1(head, "sixty-fifth") + .is_err()); + assert_eq!(host.runtime().strands().len(), 64); + let pending = host + .fork_local_operation_strand_v1(head, "pending") + .expect("pending retry"); + assert_eq!(host.runtime_wal().expect("wal").commits().len(), before + 1); + drop(host); + let (rt, _) = runtime(); + let mut restored = TrustedRuntimeHost::new(rt, empty_engine()).expect("host"); + restored + .enable_runtime_wal(TrustedRuntimeWalConfig::filesystem(&root)) + .expect("64 forks reopen"); + assert_eq!(restored.runtime().strands().len(), 64); + assert_eq!( + restored + .fork_local_operation_strand_v1(head, "pending") + .expect("pending"), + pending + ); + assert!(restored + .fork_local_operation_strand_v1(head, "sixty-fifth") + .is_err()); + drop(restored); + fs::remove_dir_all(root).expect("owned fixture cleanup"); +} diff --git a/docs/architecture/application-contract-hosting.md b/docs/architecture/application-contract-hosting.md index a1753eb6..9c8b13fc 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -1582,6 +1582,14 @@ identity from the strand identity. It allows one writer head per strand, labels of 1–64 bytes and at most 64 retained forks. This is an explicit trusted-host profile, not caller authentication or a generic authority-policy authoring API. Reusing a source-worldline/label pair resolves its original retained strand. +Before a retry or fresh fork, the host reconciles its writer cursor and all +pending committed fork records. It validates existing topology bindings and +applies missing forks at their original coordinates on private clones, exposing +them only after the complete pending set validates. A reported append error +therefore cannot cause a second fork record on retry, and pending forks count +against the capacity limit before another label is admitted. Parent advancement +does not change a retained fork's basis. Duplicate records remain a refusal; +this does not migrate an already duplicated history. Recovery validates the supported profile and source commit/boundary, replays the source history, invokes native fork construction at the recorded coordinate, and