Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
28 changes: 23 additions & 5 deletions crates/warp-core/src/trusted_runtime_host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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)?;
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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());
Expand Down
209 changes: 209 additions & 0 deletions crates/warp-core/src/trusted_runtime_host/retained_strands.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,209 @@
// SPDX-License-Identifier: Apache-2.0
// © James Ross Ω FLYING•ROBOTS <https://github.com/flyingrobots>
//! 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<ForkStrandRequest, TrustedRuntimeHostError> {
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<WriterHeadKey, TrustedRuntimeHostError> {
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::<Vec<_>>();
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))
}
60 changes: 60 additions & 0 deletions crates/warp-core/tests/trusted_runtime_host_loop_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
);
}
32 changes: 32 additions & 0 deletions docs/architecture/application-contract-hosting.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Loading
Loading