diff --git a/CHANGELOG.md b/CHANGELOG.md index d92e0412b..fd0d48cf9 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 93c66ad0e..e6681b62e 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 000000000..d095fecee --- /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 4d24707e5..00452d149 100644 --- a/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs +++ b/crates/warp-core/tests/trusted_runtime_host_loop_tests.rs @@ -2578,3 +2578,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 d35f26da1..8f6419e07 100644 --- a/docs/architecture/application-contract-hosting.md +++ b/docs/architecture/application-contract-hosting.md @@ -1218,3 +1218,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 d4066bb7d..fbc051796 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 6b8d71754..6f2e6881f 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 121025744..25d78993a 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, @@ -90,6 +91,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, @@ -97,6 +99,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(); @@ -138,7 +169,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"], @@ -154,6 +186,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, @@ -163,7 +219,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()),