diff --git a/.agents/skills/understanding-durable-execution/SKILL.md b/.agents/skills/understanding-durable-execution/SKILL.md index 895cf29f86..face4c1e81 100644 --- a/.agents/skills/understanding-durable-execution/SKILL.md +++ b/.agents/skills/understanding-durable-execution/SKILL.md @@ -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 diff --git a/.agents/skills/understanding-durable-execution/reference/timelines.md b/.agents/skills/understanding-durable-execution/reference/timelines.md index 3f12d74761..6491ac6a45 100644 --- a/.agents/skills/understanding-durable-execution/reference/timelines.md +++ b/.agents/skills/understanding-durable-execution/reference/timelines.md @@ -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 ``` diff --git a/golem-worker-executor/src/durable_host/concurrent/call.rs b/golem-worker-executor/src/durable_host/concurrent/call.rs index 48090370dd..1953fc08c3 100644 --- a/golem-worker-executor/src/durable_host/concurrent/call.rs +++ b/golem-worker-executor/src/durable_host/concurrent/call.rs @@ -3853,10 +3853,47 @@ impl DurableCallSession { /// 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( - mut self, + self, store: &Accessor, get_ctx: fn(&mut T) -> &mut DurableWorkerCtx, ) -> Result, 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( + self, + store: &Accessor, + get_ctx: fn(&mut T) -> &mut DurableWorkerCtx, + span_id: &SpanId, + ) -> Result<(DeferredCallReplayOutcome, 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( + mut self, + store: &Accessor, + get_ctx: fn(&mut T) -> &mut DurableWorkerCtx, + post_end_finish_span: Option, + ) -> Result<(DeferredCallReplayOutcome, bool), WorkerExecutorError> where T: 'static, D: HasData + ?Sized, @@ -3878,24 +3915,41 @@ impl DurableCallSession { .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::(&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 { @@ -3919,10 +3973,20 @@ impl DurableCallSession { .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( @@ -3971,7 +4035,7 @@ impl DurableCallSession { self.abandon_for_trap(); return Err(error); } - Ok(DeferredCallReplayOutcome::Incomplete(self)) + Ok((DeferredCallReplayOutcome::Incomplete(self), false)) } } } @@ -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: &mut DurableWorkerCtx, span_id: &SpanId, diff --git a/golem-worker-executor/src/durable_host/replay_state/cursor.rs b/golem-worker-executor/src/durable_host/replay_state/cursor.rs index 56c08b252a..4710be1ca8 100644 --- a/golem-worker-executor/src/durable_host/replay_state/cursor.rs +++ b/golem-worker-executor/src/durable_host/replay_state/cursor.rs @@ -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; @@ -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, diff --git a/golem-worker-executor/src/durable_host/replay_state/tests.rs b/golem-worker-executor/src/durable_host/replay_state/tests.rs index 763246ed26..602c1cb78a 100644 --- a/golem-worker-executor/src/durable_host/replay_state/tests.rs +++ b/golem-worker-executor/src/durable_host/replay_state/tests.rs @@ -9,7 +9,7 @@ use golem_common::model::entity::{ EntityCallMode, ToolInvocationClaimIdentity, ToolInvocationRejectedIdentity, }; use golem_common::model::environment::EnvironmentId; -use golem_common::model::invocation_context::TraceId; +use golem_common::model::invocation_context::{SpanId, TraceId}; use golem_common::model::oplog::payload::types::{ SerializableP3HttpBodyChunk, SerializableP3HttpConsumeBodyResult, SerializableToolRpcError, }; @@ -2993,6 +2993,132 @@ async fn marked_completion_is_prefetched_without_advancing_past_intervening_entr barrier.acknowledge(); } +#[test] +async fn prefetched_completion_waits_for_interleaved_call_before_consuming_its_span_tail() { + // The outer RPC completion is visible through its delivery marker before the positional + // cursor reaches its End. Its continuation must leave the interleaved timer call untouched, + // then consume only the FinishSpan atomically recorded next to the RPC End. + let span_id = SpanId::generate(); + let rs = replay_state_over(vec![ + noop(), + start_now(), + start_now(), + end_for(3, 43), + end_for(2, 42), + OplogEntry::FinishSpan { + timestamp: Timestamp::now_utc(), + parent_start_index: None, + span_id: span_id.clone(), + }, + delivered_for(2), + ]) + .await; + + let outer = rs + .claim_concurrent_start( + &HostFunctionName::MonotonicClockNow, + &DurableFunctionType::ReadLocal, + ) + .await + .unwrap(); + let outer_end = match rs.await_resolution(outer).await.unwrap() { + Resolution::Completed { + end_idx, + delivery_marker, + .. + } => { + assert_eq!(delivery_marker, Some(OplogIndex::from_u64(7))); + end_idx + } + other => panic!("expected prefetched outer completion, got {other:?}"), + }; + + let mut consume_span = tokio::spawn({ + let rs = rs.clone(); + let span_id = span_id.clone(); + async move { + rs.consume_finish_span_after_terminal(outer_end, span_id) + .await + } + }); + assert!( + tokio::time::timeout(Duration::from_millis(20), &mut consume_span) + .await + .is_err(), + "the outer continuation must wait instead of consuming the timer Start" + ); + assert_eq!(rs.last_replayed_index(), OplogIndex::from_u64(2)); + + let timer = rs + .claim_concurrent_start( + &HostFunctionName::MonotonicClockNow, + &DurableFunctionType::ReadLocal, + ) + .await + .unwrap(); + assert!(matches!( + rs.await_resolution(timer).await.unwrap(), + Resolution::Completed { + end_idx, + delivery_marker: None, + .. + } if end_idx == OplogIndex::from_u64(4) + )); + consume_span + .await + .expect("span-tail task panicked") + .expect("span-tail replay failed"); + assert_eq!(rs.last_replayed_index(), OplogIndex::from_u64(6)); + + let barrier = rs + .await_completion_delivery(OplogIndex::from_u64(2), OplogIndex::from_u64(7)) + .await + .unwrap(); + barrier.acknowledge(); +} + +#[test] +async fn prefetched_completion_rejects_a_different_span_tail_without_consuming_it() { + let recorded_span_id = SpanId::generate(); + let expected_span_id = SpanId::generate(); + let rs = replay_state_over(vec![ + noop(), + start_now(), + end_for(2, 42), + OplogEntry::FinishSpan { + timestamp: Timestamp::now_utc(), + parent_start_index: None, + span_id: recorded_span_id.clone(), + }, + delivered_for(2), + ]) + .await; + let handle = rs + .claim_concurrent_start( + &HostFunctionName::MonotonicClockNow, + &DurableFunctionType::ReadLocal, + ) + .await + .unwrap(); + let end_idx = match rs.await_resolution(handle).await.unwrap() { + Resolution::Completed { end_idx, .. } => end_idx, + other => panic!("expected prefetched completion, got {other:?}"), + }; + + let error = rs + .consume_finish_span_after_terminal(end_idx, expected_span_id) + .await + .expect_err("a different span tail must be rejected"); + assert!(error.to_string().contains("FinishSpan for span")); + assert_eq!(rs.last_replayed_index(), end_idx); + let (index, entry) = rs.get_oplog_entry().await.unwrap(); + assert_eq!(index, end_idx.next()); + assert!(matches!( + entry, + OplogEntry::FinishSpan { span_id, .. } if span_id == recorded_span_id + )); +} + #[test] async fn replay_delivery_marker_holds_cursor_until_guest_boundary() { // A completed before B was started, but A's callback was handed to the guest only after B diff --git a/golem-worker-executor/src/durable_host/wasm_rpc/mod.rs b/golem-worker-executor/src/durable_host/wasm_rpc/mod.rs index c2c8087190..2b4a52fa44 100644 --- a/golem-worker-executor/src/durable_host/wasm_rpc/mod.rs +++ b/golem-worker-executor/src/durable_host/wasm_rpc/mod.rs @@ -3446,7 +3446,7 @@ impl HostFutureInvokeResultWithStore // records the `CompletionDiscarded` marker. The call's durable `FinishSpan` is // appended by the same owned task as the `End` (see `post_end_entry`), so replay // can rely on it unconditionally following the `End` on this path. - let (response, delivery) = if handle.is_live() { + let (response, delivery, replayed_finish_span) = if handle.is_live() { let task = task.expect("a live future-invoke-result must own its background task"); let interrupt_signal = accessor.with(|mut access| { @@ -3485,7 +3485,7 @@ impl HostFutureInvokeResultWithStore .await; } }; - match task_result { + let (response, delivery) = match task_result { Ok(rpc_result) => { let rpc_result = admit_rpc_result_secret_holds(accessor, rpc_result).await?; @@ -3511,15 +3511,20 @@ impl HostFutureInvokeResultWithStore // incomplete for durable-scope recovery, instead of recording an `End`. return Err(handle.trap(anyhow::anyhow!(err.to_string()))); } - } + }; + (response, delivery, false) } else { - match handle - .replay_access_deferred(accessor, accessor.getter()) + let (outcome, replayed_finish_span) = handle + .replay_access_deferred_with_finish_span( + accessor, + accessor.getter(), + &span_id, + ) .await - .map_err(anyhow::Error::from)? - { + .map_err(anyhow::Error::from)?; + match outcome { DeferredCallReplayOutcome::Replayed(response, delivery) => { - (response, delivery) + (response, delivery, replayed_finish_span) } DeferredCallReplayOutcome::Incomplete(mut live) => { // Crash-after-`Start` recovery: the eager `Start` is committed but its @@ -3581,7 +3586,7 @@ impl HostFutureInvokeResultWithStore result = &mut task => Some(result), } }; - match task_result { + let (response, delivery) = match task_result { None => { // Cancelled while re-executing the recovered call: same // handling as the live-path cancellation race above. @@ -3614,7 +3619,8 @@ impl HostFutureInvokeResultWithStore Some(Err(err)) => { return Err(live.trap(anyhow::anyhow!(err.to_string()))); } - } + }; + (response, delivery, false) } } }; @@ -3634,11 +3640,13 @@ impl HostFutureInvokeResultWithStore if delivery.is_replay_discarded() { // The recorded run persisted the `End` (and its `FinishSpan`) but the guest // dropped this future before `get` returned. Mirror the recorded post-`End` - // continuation deterministically — consume the positional `FinishSpan` and - // mark the resource consumed — then park: never return the response, so the - // deterministic guest drops this future at the same point it did live (its - // resource `drop` sees no open handle and writes nothing durable). - finish_span_access(accessor, accessor.getter(), &span_id).await?; + // continuation deterministically — apply the already-consumed `FinishSpan` + // in memory and mark the resource consumed — then park: never return the + // response, so the deterministic guest drops this future at the same point it + // did live (its resource `drop` sees no open handle and writes nothing + // durable). + debug_assert!(replayed_finish_span); + accessor.with(|mut access| finish_span_in_memory(access.get(), &span_id))?; accessor.with(|mut access| { let ctx = access.get(); let entry = ctx @@ -3696,9 +3704,14 @@ impl HostFutureInvokeResultWithStore } } else { // Replay of a delivered completion, or a live unpersisted (snapshotting) - // call: the original span handling applies — replay consumes the positional - // `FinishSpan`, an unpersisted live call appends it here. - finish_span_access(accessor, accessor.getter(), &span_id).await?; + // call: a successful replay already consumed the exact post-`End` tail; + // cancellation replay and an unpersisted live call use positional handling. + if replayed_finish_span { + accessor + .with(|mut access| finish_span_in_memory(access.get(), &span_id))?; + } else { + finish_span_access(accessor, accessor.getter(), &span_id).await?; + } accessor.with(|mut access| { let ctx = access.get(); let entry = ctx