Skip to content
Open
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 @@ -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.
Expand Down
20 changes: 20 additions & 0 deletions crates/warp-core/src/causal_wal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Path>,
mode: RecoveryAccessMode,
) -> Result<RecoveryScanReport, WalRecoveryError> {
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 {
Expand Down
29 changes: 29 additions & 0 deletions crates/warp-core/src/trusted_runtime_host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,12 @@ impl From<RuntimeError> 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),
Expand Down Expand Up @@ -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<crate::evidence::CausalSegmentCatalog>,
evidence_catalog_posture: EvidenceCatalogPosture,
Expand Down Expand Up @@ -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()
Expand All @@ -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),
Expand Down Expand Up @@ -2566,6 +2579,7 @@ impl TrustedRuntimeWal {
fn in_memory_rollback_snapshot(&self) -> Option<Self> {
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"))]
Expand Down Expand Up @@ -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)?
Expand Down Expand Up @@ -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(())
}

Expand Down Expand Up @@ -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;
Expand Down
171 changes: 100 additions & 71 deletions crates/warp-core/src/trusted_runtime_host/observed_context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down Expand Up @@ -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<String, EchoOperationObservationV1>,
requests: BTreeMap<String, (String, Hash, Vec<u8>)>,
}

impl TrustedRuntimeWal {
fn operation_contexts(&self) -> Result<Contexts> {
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<Self> {
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<std::borrow::Cow<'_, Contexts>> {
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<V>) -> 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(
Expand Down Expand Up @@ -151,11 +177,12 @@ impl TrustedRuntimeHost {
nodes: &[NodeKey],
) -> Result<EchoOperationObservationV1> {
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"));
}
Expand Down Expand Up @@ -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"))
}

Expand Down Expand Up @@ -292,11 +320,12 @@ impl TrustedRuntimeHost {
invocation: EchoOperationInvocationV1,
) -> Result<Vec<u8>> {
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)
Expand Down
Loading
Loading