Skip to content
Closed
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
8 changes: 8 additions & 0 deletions .agents/skills/understanding-durable-execution/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -387,6 +387,14 @@ possible; a hard error by itself is not a recovery.
Host completion time and guest observation time are different facts. Only the second is a guest
input, and only the second must recur exactly; the first may vary between runs.

A deferred accessor may append a mandatory positional tail entry, such as `FinishSpan`, in the
same owned task immediately after its `End`. Completion-marker lookahead can resolve that `End`
before the positional cursor reaches it. The replaying continuation must therefore wait for the
owners of any interleaved entries to advance normally, auto-drain its own terminal, and then
atomically validate and consume the exact adjacent tail entry by identity. It must not issue an
ordinary positional read early, scan past unrelated work, or apply the tail's in-memory effect
before the durable entry is consumed.

Fork and revert retain the exact inclusive prefix. A cut between `Start` and its terminal uses
ordinary incomplete-call recovery; a cut between an accessor `End` and its delivery/discard
marker uses `AtReplayTail`. No outcome beyond the cut is inherited. Revert leaves deleted entries
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -186,6 +186,13 @@ not a guest input; delivery order (#64 before #65) is. Owner: `ReplayDeliveryBar
property this relies on (bare Wasmtime, no oplog); the marker mechanics are covered by
`replay_state/tests.rs` and `concurrent/tests.rs`.

If A atomically recorded a positional tail such as `FinishSpan` immediately after `End A`, a
prefetched resolution does not make that tail positionally available. A's continuation waits while
the owners of entries before `End A` advance them, lets the cursor auto-drain `End A`, then validates
and consumes the exact adjacent `FinishSpan` before applying its in-memory span transition. Reading
`FinishSpan` through the ordinary positional API as soon as A resolves would steal whichever
interleaved entry is still at the cursor head.

## 9. Automatic update with snapshot

```
Expand Down
94 changes: 79 additions & 15 deletions golem-worker-executor/src/durable_host/concurrent/call.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3853,10 +3853,47 @@ impl<Pair: HostPayloadPair, P: DropPolicy> DurableCallSession<Pair, P> {
/// was armed — and the *caller* parks at its own delivery boundary after performing its
/// deterministic post-`End` continuation (e.g. consuming a positional `FinishSpan`).
pub async fn replay_access_deferred<T, D, Ctx>(
mut self,
self,
store: &Accessor<T, D>,
get_ctx: fn(&mut T) -> &mut DurableWorkerCtx<Ctx>,
) -> Result<DeferredCallReplayOutcome<Pair, P>, WorkerExecutorError>
where
T: 'static,
D: HasData + ?Sized,
Ctx: WorkerCtx,
{
self.replay_access_deferred_impl(store, get_ctx, None)
.await
.map(|(outcome, _)| outcome)
}

/// Replays a deferred accessor call whose successful live terminal atomically appended a
/// `FinishSpan` immediately after its `End`.
///
/// The boolean reports whether that durable span tail was consumed here. A cancellation with
/// a partial response has no post-`End` tail and still requires the caller's ordinary
/// positional span handling.
pub async fn replay_access_deferred_with_finish_span<T, D, Ctx>(
self,
store: &Accessor<T, D>,
get_ctx: fn(&mut T) -> &mut DurableWorkerCtx<Ctx>,
span_id: &SpanId,
) -> Result<(DeferredCallReplayOutcome<Pair, P>, bool), WorkerExecutorError>
where
T: 'static,
D: HasData + ?Sized,
Ctx: WorkerCtx,
{
self.replay_access_deferred_impl(store, get_ctx, Some(span_id.clone()))
.await
}

async fn replay_access_deferred_impl<T, D, Ctx>(
mut self,
store: &Accessor<T, D>,
get_ctx: fn(&mut T) -> &mut DurableWorkerCtx<Ctx>,
post_end_finish_span: Option<SpanId>,
) -> Result<(DeferredCallReplayOutcome<Pair, P>, bool), WorkerExecutorError>
where
T: 'static,
D: HasData + ?Sized,
Expand All @@ -3878,24 +3915,41 @@ impl<Pair: HostPayloadPair, P: DropPolicy> DurableCallSession<Pair, P> {
.take()
.expect("replay_access_deferred() called on a live handle");
let outcome = replay_state.await_resolution_outcome(replay).await?;
let finish_span_tail = post_end_finish_span.and_then(|span_id| match &outcome {
ResolutionOutcome::Resolved(Resolution::Completed { end_idx, .. })
| ResolutionOutcome::Resolved(Resolution::CompletedButDiscarded { end_idx, .. }) => {
Some((*end_idx, span_id))
}
ResolutionOutcome::Resolved(Resolution::Cancelled { .. })
| ResolutionOutcome::Incomplete => None,
});
match classify_replay_resolution(outcome) {
ReplayedResolution::Delivered(payload, disposition) => {
self.finished = true;
let response = decode_replayed_payload::<Pair>(&oplog, payload).await?;
end_durable_function_access(store, get_ctx, function_type, begin_index, false)
.await?;
let consumed_finish_span = finish_span_tail.is_some();
if let Some((end_idx, span_id)) = finish_span_tail {
replay_state
.consume_finish_span_after_terminal(end_idx, span_id)
.await?;
}
// The delivery token is constructed only after the fallible decode / scope close
// succeeded: a token dropped on the error path would log a spurious
// unconsumed-token warning.
Ok(DeferredCallReplayOutcome::Replayed(
response,
CompletionDelivery::replay_delivered(
disposition,
self.start_idx,
completion_marker_recorder,
self.trap_context(),
self.cleanup_sink.clone(),
Ok((
DeferredCallReplayOutcome::Replayed(
response,
CompletionDelivery::replay_delivered(
disposition,
self.start_idx,
completion_marker_recorder,
self.trap_context(),
self.cleanup_sink.clone(),
),
),
consumed_finish_span,
))
}
ReplayedResolution::Undelivered(UndeliveredTerminal::CompletionDiscarded {
Expand All @@ -3919,10 +3973,20 @@ impl<Pair: HostPayloadPair, P: DropPolicy> DurableCallSession<Pair, P> {
.await?;
end_durable_function_access(store, get_ctx, function_type, begin_index, false)
.await?;
let consumed_finish_span = finish_span_tail.is_some();
if let Some((tail_end_idx, span_id)) = finish_span_tail {
debug_assert_eq!(tail_end_idx, end_idx);
replay_state
.consume_finish_span_after_terminal(tail_end_idx, span_id)
.await?;
}
// As above, the token is constructed only after the fallible operations.
Ok(DeferredCallReplayOutcome::Replayed(
response,
CompletionDelivery::replay_discarded(),
Ok((
DeferredCallReplayOutcome::Replayed(
response,
CompletionDelivery::replay_discarded(),
),
consumed_finish_span,
))
}
ReplayedResolution::Undelivered(
Expand Down Expand Up @@ -3971,7 +4035,7 @@ impl<Pair: HostPayloadPair, P: DropPolicy> DurableCallSession<Pair, P> {
self.abandon_for_trap();
return Err(error);
}
Ok(DeferredCallReplayOutcome::Incomplete(self))
Ok((DeferredCallReplayOutcome::Incomplete(self), false))
}
}
}
Expand Down Expand Up @@ -4923,8 +4987,8 @@ where

/// Finishes a span in the in-memory invocation context only, without writing or consuming any
/// oplog entry: pops the current-span pointer if it points at this span, then marks the span
/// finished. This is the non-durable half of [`finish_span_access`], used directly for spans
/// whose identity is derived from durable records (no `StartSpan`/`FinishSpan` entries exist).
/// finished. This is the in-memory half of [`finish_span_access`], used directly when the durable
/// entry either does not exist or has already been handled separately.
pub(crate) fn finish_span_in_memory<Ctx: WorkerCtx>(
ctx: &mut DurableWorkerCtx<Ctx>,
span_id: &SpanId,
Expand Down
79 changes: 79 additions & 0 deletions golem-worker-executor/src/durable_host/replay_state/cursor.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
use super::claims::{RequestClaimIdentity, StartClaim, recorded_request_payload_matches};
use super::*;
use crate::durable_host::PositionalRead;
use golem_common::model::invocation_context::SpanId;
#[cfg(feature = "test-utils")]
use std::pin::Pin;

Expand Down Expand Up @@ -2826,6 +2827,84 @@ impl ReplayState {
}
}

/// Waits until the positional cursor reaches a prefetched call's terminal, then consumes the
/// exact `FinishSpan` recorded immediately after it.
///
/// A completion marker lets a concurrent call resolve from lookahead before its `End` reaches
/// the cursor. Its host continuation may therefore run while earlier interleaved calls still
/// own the cursor head. Those calls must advance normally; once they do, this method
/// auto-drains the resolver-owned terminal and atomically validates the adjacent span tail.
pub(crate) async fn consume_finish_span_after_terminal(
&self,
terminal_index: OplogIndex,
span_id: SpanId,
) -> Result<(), WorkerExecutorError> {
self.run_owned_cursor_op(move |state| async move {
loop {
let progress = state.cursor.progress.notified();
tokio::pin!(progress);
progress.as_mut().enable();

let consumed = state
.with_tx(async |tx| {
if tx.cursor.last_replayed_index() < terminal_index {
// Drive only entries already owned by claimed calls. An unrelated
// positional entry remains untouched for its own replaying task.
tx.try_get_oplog_entry(|_| false).await?;
}

let cursor_index = tx.cursor.last_replayed_index();
if cursor_index < terminal_index {
return Ok(false);
}
if cursor_index > terminal_index {
return Err(WorkerExecutorError::unexpected_oplog_entry(
format!("FinishSpan immediately after terminal {terminal_index}"),
format!(
"replay cursor already advanced to {cursor_index} for {}",
tx.cursor.owned_agent_id
),
));
}
if tx.cursor.is_live() {
return Err(WorkerExecutorError::unexpected_oplog_entry(
format!("FinishSpan immediately after terminal {terminal_index}"),
format!(
"end of replay for {} at terminal {terminal_index}",
tx.cursor.owned_agent_id
),
));
}

let (read_index, entry) = tx.raw_read_next_oplog_entry().await?;
let expected_index = terminal_index.next();
match &entry {
OplogEntry::FinishSpan {
span_id: recorded_span_id,
..
} if read_index == expected_index && recorded_span_id == &span_id => {
tx.commit_consumed_entry(read_index, &entry).await?;
Ok(true)
}
_ => Err(WorkerExecutorError::unexpected_oplog_entry(
format!(
"FinishSpan for span {span_id} at {expected_index} immediately after terminal {terminal_index}"
),
format!("{entry:?} at {read_index}"),
)),
}
})
.await?;

if consumed {
return Ok(());
}
progress.await;
}
})
.await
}

/// Owned-task variant of [`Self::get_oplog_entry_or_replay_end`].
pub async fn get_oplog_entry_or_replay_end_owned(
&self,
Expand Down
Loading
Loading