Skip to main content

taquba_workflow/
runtime.rs

1use std::collections::HashMap;
2use std::future::Future;
3use std::sync::Arc;
4use std::time::{Duration, SystemTime, UNIX_EPOCH};
5
6use futures_util::TryStreamExt;
7use taquba::object_store::ObjectStore;
8use taquba::{
9    Clock, EnqueueOptions, EnqueueRequest, EnqueueResult, JobRecord, JobStatus, Queue,
10    SettlementEffects, WaitOutcome, WorkerHandle,
11};
12use tokio_util::sync::CancellationToken;
13use tracing::{debug, instrument, warn};
14
15use crate::durable::{
16    self, DurableCurrentStep, DurableErrorKind, DurableRunOutcome, DurableRunRecord,
17    DurableRunResult, DurableStepOutcome, DurableStepOutcomeRecord, DurableTermination,
18};
19use crate::effects::StagedEffects;
20use crate::error::{Error, Result};
21use crate::group::{GroupStore, Membership, RunGroup, pending_member, terminated_member};
22use crate::keys::{
23    DEDUP_PREFIX, GROUP_TERMINAL_KV_PREFIX, HEADER_RUN_ID, HEADER_STEP, HEADER_TERMINAL,
24    RESERVED_HEADER_PREFIX, RESERVED_KV_PREFIX, RunId, TERMINAL_KV_PREFIX, hash_input,
25    outcome_kv_key, run_kv_key, step_kv_key,
26};
27use crate::memo::{MemoStore, RUN_RESULT_MEMO_KEY};
28use crate::runner::{StepErrorKind, StepOutcome, StepRunner, Trigger};
29use crate::sweep::{Clearable, Sweep, run_periodically};
30use crate::terminal::{RunOutcome, TerminalHook, TerminalStatus};
31use crate::view::WorkflowView;
32use crate::worker::{ClaimedStep, StepWorker};
33
34/// The encoded current-step pointer for `job_id` at `step_number`.
35fn current_step_bytes(step_number: u32, job_id: &str) -> Vec<u8> {
36    durable::encode(&DurableCurrentStep {
37        step_number,
38        job_id: job_id.to_string(),
39    })
40}
41
42/// The portion of `delay` still ahead of `now_ms`, measured from
43/// `stored_at_ms`. Saturates to the full delay if the clock reads
44/// earlier than the stored timestamp.
45fn remaining_delay(stored_at_ms: u64, now_ms: u64, delay: Duration) -> Duration {
46    let elapsed = Duration::from_millis(now_ms.saturating_sub(stored_at_ms));
47    delay.saturating_sub(elapsed)
48}
49
50/// Per-step enqueue options the runtime forwards through to Taquba. The
51/// runtime always owns `headers` (it injects [`HEADER_RUN_ID`] and
52/// [`HEADER_STEP`]) and `dedup_key` (it derives one from
53/// `(run_id, step_number)`), so callers only pick the three fields below.
54#[derive(Debug, Default)]
55pub(crate) struct StepEnqueueOpts {
56    /// Earliest claimable time for the step. `None` means immediate.
57    pub(crate) run_at: Option<SystemTime>,
58    /// Per-step priority override.
59    pub(crate) priority: Option<u32>,
60    /// Per-step `max_attempts` override.
61    pub(crate) max_attempts: Option<u32>,
62    /// Additional runtime-owned reserved headers for the step job.
63    pub(crate) reserved_headers: Vec<(&'static str, String)>,
64}
65
66/// The settings of a run's steps, applied by [`WorkflowRuntime::submit`]
67/// through [`RunSpec::options`], by [`RunGroup::submit`] and
68/// [`RunGroup::resume`] to every member of a group and by
69/// [`jobs::JobRunner::submit_with`](crate::jobs::JobRunner::submit_with)
70/// to a job.
71#[derive(Debug, Clone, Default)]
72pub struct RunOptions {
73    /// Submitter-supplied metadata, threaded through every step of the
74    /// run and surfaced to the terminal hook. Reserved `workflow.*` keys
75    /// are rejected at submission with [`Error::ReservedHeaderInSubmit`].
76    pub headers: HashMap<String, String>,
77    /// Priority of every step; the queue's default when `None`.
78    pub priority: Option<u32>,
79    /// Attempt limit of every step; the queue's `max_attempts` when
80    /// `None`.
81    pub max_attempts_per_step: Option<u32>,
82    /// Earliest time the first step may run. The step-0 job waits in
83    /// the queue's scheduled state until the queue's clock passes this
84    /// time; `None` makes it claimable at once.
85    pub run_at: Option<SystemTime>,
86}
87
88/// Spec passed to [`WorkflowRuntime::submit`].
89#[derive(Debug, Clone, Default)]
90pub struct RunSpec {
91    /// Caller-supplied run identifier. If `None`, the runtime generates
92    /// a ULID. The dedup key for the first step job is `run:{run_id}:0`, so
93    /// re-submitting the same `run_id` while the run is active returns the
94    /// existing job rather than creating a duplicate.
95    ///
96    /// A terminated run releases its id for re-submission. The second
97    /// run shares the first run's memo and step-output entries, which is
98    /// what makes a re-submission resume from them, and under
99    /// [`WorkflowRuntimeBuilder::memo_retention`] the first run's marker
100    /// expires against those shared entries even while the second run is
101    /// executing. The second run then re-executes the affected steps.
102    pub run_id: Option<RunId>,
103    /// Bytes handed to the runner as the first step's payload.
104    pub input: Vec<u8>,
105    /// The settings of the run's steps.
106    pub options: RunOptions,
107    /// Effects applied in the same transaction as the step-0 enqueue, as
108    /// [`taquba::Queue::enqueue_with_effects`] applies them: the enqueues, the
109    /// KV writes, the KV deletes and the expiry entries. Applied only when the
110    /// submission is new: a duplicate submission's effects are dropped, and the
111    /// effects do not participate in the duplicate-submission input check. A KV
112    /// key written or deleted must not start with the reserved `workflow/`
113    /// prefix ([`RESERVED_KV_PREFIX`]). Values are capped at
114    /// [`taquba::MAX_KV_VALUE_SIZE`].
115    pub effects: SettlementEffects,
116}
117
118/// Outcome of [`WorkflowRuntime::submit`].
119///
120/// `submit` is idempotent on `run_id`: re-submitting an active run is a
121/// no-op and the returned `SubmitOutcome` carries `newly_submitted = false`.
122#[derive(Debug, Clone)]
123pub struct SubmitOutcome {
124    /// The run's identifier (generated if the spec didn't carry one).
125    pub run_id: RunId,
126    /// `true` if this call enqueued a new run; `false` if a run with this
127    /// id was already active (its durable run record exists) and this
128    /// call was a no-op. Call
129    /// [`WorkflowRuntime::status`] for the run's current state when
130    /// needed.
131    pub newly_submitted: bool,
132    /// The id of the queue job currently representing the run: its
133    /// first step for a new submission, and the step the run has reached
134    /// for a duplicate, read from the run's durable current-step pointer.
135    pub job_id: String,
136}
137
138/// Status snapshot of a run, read from its durable state by
139/// [`WorkflowView::status`].
140#[derive(Debug, Clone)]
141pub struct RunStatus {
142    /// The run's identifier.
143    pub run_id: RunId,
144    /// Lifecycle state of the run's current step, or its termination.
145    pub state: RunState,
146    /// Step number of the run's current step; the final step of a
147    /// terminated run.
148    pub current_step: u32,
149}
150
151/// Lifecycle state tracked in [`RunStatus::state`].
152#[derive(Debug, Clone, PartialEq, Eq)]
153pub enum RunState {
154    /// A step job exists in the queue but has not yet been claimed.
155    Pending,
156    /// A step is currently being processed by a worker.
157    Running,
158    /// [`WorkflowRuntime::cancel`] was called for this run and the run
159    /// has not yet terminated. Reported until the in-flight step
160    /// returns and the runtime settles the run as
161    /// [`crate::TerminalStatus::Cancelled`]; after that,
162    /// [`WorkflowRuntime::status`] returns `None`.
163    ///
164    /// Only set by external cancellation. A pure runner-issued
165    /// [`crate::StepOutcome::Cancel`] (with no external `cancel()`
166    /// call) terminates as `Cancelled` without ever transitioning
167    /// through `Cancelling`: a runner-issued cancel is observed when
168    /// `run_step` returns, and the run terminates at that point.
169    Cancelling,
170    /// The run reached a terminal state. Reported from the run's terminal
171    /// record, written with the terminating settlement and removed with
172    /// the run's memo entries by the memo sweep under
173    /// [`WorkflowRuntimeBuilder::memo_retention`].
174    Terminated(RunTermination),
175}
176
177/// A run result record as read back: the committed outcome and the
178/// termination the record belongs to.
179#[derive(Debug, Clone)]
180pub(crate) struct RunResult {
181    pub(crate) termination: RunTermination,
182    pub(crate) outcome: RunOutcome,
183}
184
185/// The committed terminal outcome of a run, as
186/// [`RunState::Terminated`] reports it.
187#[derive(Debug, Clone, PartialEq, Eq)]
188pub struct RunTermination {
189    /// How the run ended.
190    pub status: TerminalStatus,
191    /// The failure reason, or the runner's reason for a
192    /// [`StepOutcome::Cancel`]; `None` for a success and for an external
193    /// cancellation.
194    pub error: Option<String>,
195    /// The classification of a failure: the kind of the
196    /// [`StepError`](crate::StepError) that dead-lettered the run, or
197    /// [`StepErrorKind::Permanent`] for a [`StepOutcome::Fail`] verdict.
198    /// `None` for a success, a cancellation and a termination outside
199    /// the worker.
200    pub error_kind: Option<StepErrorKind>,
201    /// The number of the step whose settlement terminated the run.
202    pub final_step: u32,
203    /// The runtime clock's time at the terminating settlement, in
204    /// milliseconds since the Unix epoch.
205    pub terminated_at_ms: u64,
206}
207
208impl From<DurableTermination> for RunTermination {
209    fn from(record: DurableTermination) -> Self {
210        Self {
211            status: record.status.into(),
212            error: record.error,
213            error_kind: record.error_kind.map(Into::into),
214            final_step: record.final_step,
215            terminated_at_ms: record.terminated_at_ms,
216        }
217    }
218}
219
220/// The end of a run, as [`WorkflowRuntime::wait`] reports it.
221#[derive(Debug, Clone)]
222pub struct RunEnd {
223    /// The run's termination, from its terminal record.
224    pub termination: RunTermination,
225    /// The committed outcome, when the worker that terminated the run
226    /// recorded one; see [`WorkflowRuntime::outcome`].
227    pub outcome: Option<RunOutcome>,
228}
229
230/// Builder for [`WorkflowRuntime`].
231///
232/// Construct via [`WorkflowRuntime::builder`].
233pub struct WorkflowRuntimeBuilder<R, H> {
234    queue: Arc<Queue>,
235    object_store: Arc<dyn ObjectStore>,
236    queue_name: String,
237    memo_prefix: Option<String>,
238    runner: R,
239    terminal_hook: H,
240    max_concurrent_steps: usize,
241    poll_interval: Duration,
242    memo_retention: Option<Duration>,
243    group_retention: Option<Duration>,
244    step_output_replay: bool,
245    clock: Arc<dyn Clock>,
246}
247
248impl<R: StepRunner, H: TerminalHook> WorkflowRuntimeBuilder<R, H> {
249    /// The Taquba queue name that step jobs are enqueued onto. Defaults to
250    /// `"workflow-steps"`. Multiple runtimes can share a `Queue` handle by
251    /// using distinct queue names.
252    pub fn queue_name(mut self, name: impl Into<String>) -> Self {
253        self.queue_name = name.into();
254        self
255    }
256
257    /// The object-store path prefix of the [`Delivery::memo`](crate::Delivery::memo)
258    /// entries. Defaults to `"{queue_name}-memo"`, so runtimes with
259    /// distinct queue names on one object store have distinct memo
260    /// namespaces.
261    pub fn memo_prefix(mut self, prefix: impl Into<String>) -> Self {
262        self.memo_prefix = Some(prefix.into());
263        self
264    }
265
266    /// Maximum number of steps processed concurrently in [`WorkflowRuntime::run`].
267    /// Defaults to 16.
268    pub fn max_concurrent_steps(mut self, n: usize) -> Self {
269        assert!(n > 0, "max_concurrent_steps must be at least 1");
270        self.max_concurrent_steps = n;
271        self
272    }
273
274    /// Maximum time the worker loop waits on an empty queue before re-checking.
275    /// Defaults to 250ms. The retention sweeps run at this interval.
276    pub fn poll_interval(mut self, interval: Duration) -> Self {
277        self.poll_interval = interval;
278        self
279    }
280
281    /// Enable memo retention with the given window. When set, the
282    /// runtime writes a terminal marker for every run that reaches a
283    /// terminal state, and the in-process sweeper clears that run's memo
284    /// entries and terminal record `retention` after termination. When
285    /// unset (default), no marker is written and memo entries and
286    /// terminal records are retained indefinitely.
287    pub fn memo_retention(mut self, retention: Duration) -> Self {
288        self.memo_retention = Some(retention);
289        self
290    }
291
292    /// Enable content-addressed replay of runner-returned step outcomes.
293    ///
294    /// When enabled, the runtime writes every [`StepOutcome`] the runner
295    /// returns, including `Fail` and `Cancel`, to object storage before
296    /// applying it. Step errors ([`StepError`](crate::StepError)) are not
297    /// recorded, so retries still invoke the runner. The replay key is
298    /// scoped to `(run_id, step_number, SHA-256(step payload))`. If the
299    /// same step is delivered again after a crash before ack, the stored
300    /// outcome is replayed without invoking the runner again. The record
301    /// includes the effects staged through
302    /// [`Delivery::effects`](crate::Delivery::effects), so a
303    /// replayed outcome applies them as well. A replayed
304    /// [`StepOutcome::Continue`] with a [`Trigger::After`] delay reduces
305    /// the delay by the time already elapsed since the outcome was
306    /// stored, preserving the original schedule.
307    ///
308    /// This is disabled by default because it adds one object-store read
309    /// per step delivery (the replay lookup) plus one write per recorded
310    /// outcome, and makes that write part of step settlement.
311    pub fn step_output_replay(mut self) -> Self {
312        self.step_output_replay = true;
313        self
314    }
315
316    /// Override the [`Clock`] the runtime reads its timestamps from.
317    /// Defaults to the same clock the [`Queue`] was opened with (via
318    /// [`Queue::clock`]).
319    pub fn clock(mut self, clock: Arc<dyn Clock>) -> Self {
320        self.clock = clock;
321        self
322    }
323
324    /// Remove a [`RunGroup`]'s state (its manifest, member records and
325    /// the memo entries and terminal records of its members) `retention` after a
326    /// [`RunGroup::results`] consumer observed the last member's
327    /// termination, through a sweep when the worker starts and at every
328    /// poll interval after that. When unset (default), no group
329    /// terminal marker is written. A group whose results are never
330    /// consumed is retained until [`RunGroup::forget`] in either case.
331    pub fn group_retention(mut self, retention: Duration) -> Self {
332        self.group_retention = Some(retention);
333        self
334    }
335
336    /// Finalize the builder.
337    pub fn build(self) -> WorkflowRuntime<R, H>
338    where
339        H: 'static,
340    {
341        let terminal_hook = Arc::new(self.terminal_hook);
342        let observes: Arc<dyn Fn(&RunOutcome) -> bool + Send + Sync> = {
343            let hook = terminal_hook.clone();
344            Arc::new(move |outcome| hook.observes(outcome))
345        };
346        let memo_prefix = self
347            .memo_prefix
348            .unwrap_or_else(|| format!("{}-memo", self.queue_name));
349        let memo_store = MemoStore::new(self.object_store.clone(), memo_prefix.clone());
350        let group_store = GroupStore::new(
351            self.object_store,
352            memo_prefix,
353            memo_store.clone(),
354            self.queue.clone(),
355        );
356        let memo_sweep = self.memo_retention.map(|retention| {
357            Arc::new(Sweep::new(
358                TERMINAL_KV_PREFIX,
359                retention,
360                RunStore {
361                    memo_store: memo_store.clone(),
362                },
363            ))
364        });
365        let group_sweep = self.group_retention.map(|retention| {
366            Arc::new(Sweep::new(
367                GROUP_TERMINAL_KV_PREFIX,
368                retention,
369                group_store.clone(),
370            ))
371        });
372        let view = WorkflowView::new(self.queue.view().clone(), memo_store.clone());
373        let core = RuntimeCore {
374            queue: self.queue,
375            view,
376            queue_name: self.queue_name,
377            max_concurrent_steps: self.max_concurrent_steps,
378            poll_interval: self.poll_interval,
379            memo_store,
380            group_store,
381            memo_sweep,
382            group_sweep,
383            step_output_replay: self.step_output_replay,
384            clock: self.clock,
385            observes,
386        };
387        let inner = RuntimeInner {
388            runner: self.runner,
389            terminal_hook,
390            core: Arc::new(core),
391        };
392        WorkflowRuntime {
393            inner: Arc::new(inner),
394        }
395    }
396}
397
398/// The retained state of a terminated run: its memo and step-output
399/// entries in the object store and its terminal record in the queue's
400/// KV namespace, removed together by the memo sweep.
401struct RunStore {
402    memo_store: MemoStore,
403}
404
405impl Clearable for RunStore {
406    type Error = Error;
407
408    async fn clear(&self, run_id: &RunId) -> Result<Vec<Vec<u8>>> {
409        self.memo_store.clear_memos_for_run(run_id).await?;
410        Ok(vec![outcome_kv_key(run_id)])
411    }
412}
413
414/// Durable runtime for workflow runs. Cheap to clone (internally `Arc`).
415pub struct WorkflowRuntime<R, H> {
416    pub(crate) inner: Arc<RuntimeInner<R, H>>,
417}
418
419impl<R, H> Clone for WorkflowRuntime<R, H> {
420    fn clone(&self) -> Self {
421        Self {
422            inner: self.inner.clone(),
423        }
424    }
425}
426
427/// A runtime's [`StepRunner`] and [`TerminalHook`], with the shared
428/// [`RuntimeCore`] they operate on: the executing half of a runtime,
429/// held by the worker.
430pub(crate) struct RuntimeInner<R, H> {
431    pub(crate) runner: R,
432    pub(crate) terminal_hook: Arc<H>,
433    pub(crate) core: Arc<RuntimeCore>,
434}
435
436/// The control half of a runtime: the queue, the stores and the settings,
437/// with every operation that reads or settles run state without invoking the
438/// runner or the hook. It is not generic over the runner and hook types, so the
439/// worker, [`WorkflowRuntime`] and [`RunGroup`](crate::RunGroup) share it.
440pub(crate) struct RuntimeCore {
441    pub(crate) queue: Arc<Queue>,
442    pub(crate) view: WorkflowView,
443    queue_name: String,
444    max_concurrent_steps: usize,
445    poll_interval: Duration,
446    pub(crate) memo_store: MemoStore,
447    pub(crate) group_store: GroupStore,
448    /// The sweep that removes the memo entries and the terminal record
449    /// of a run a window after its termination, when
450    /// [`WorkflowRuntimeBuilder::memo_retention`] is set. Without it no
451    /// terminal marker is written.
452    pub(crate) memo_sweep: Option<Arc<Sweep>>,
453    /// The sweep that removes the state of a group a window after its
454    /// last member terminated, when
455    /// [`WorkflowRuntimeBuilder::group_retention`] is set. Without it no
456    /// group marker is written.
457    pub(crate) group_sweep: Option<Arc<Sweep>>,
458    /// Whether runner-returned step outcomes are persisted and replayed
459    /// by `(run_id, step_number, SHA-256(step payload))`.
460    pub(crate) step_output_replay: bool,
461    /// Time source. Defaults to the queue's clock; tests can substitute
462    /// a [`MockClock`](taquba::MockClock) to virtualise time.
463    pub(crate) clock: Arc<dyn Clock>,
464    /// [`TerminalHook::observes`] of the runtime's hook, which decides
465    /// whether a termination enqueues a notification job.
466    observes: Arc<dyn Fn(&RunOutcome) -> bool + Send + Sync>,
467}
468
469impl<R: StepRunner, H: TerminalHook> WorkflowRuntime<R, H> {
470    /// Start configuring a runtime. Takes the four required dependencies
471    /// (Taquba queue, object store, [`StepRunner`], [`TerminalHook`]); optional
472    /// fields are set via [`WorkflowRuntimeBuilder`] methods before [`build`].
473    ///
474    /// The object store backs [`Delivery::memo`]; it does **not** need to be
475    /// the same store the [`Queue`] was opened with, though sharing one store
476    /// is the common case (just clone the `Arc`). Use a distinct
477    /// [`WorkflowRuntimeBuilder::memo_prefix`] when multiple runtimes share one
478    /// store.
479    ///
480    /// Use [`crate::NoopTerminalHook`] if you don't need terminal callbacks.
481    ///
482    /// [`Delivery::memo`]: crate::Delivery::memo
483    /// [`build`]: WorkflowRuntimeBuilder::build
484    pub fn builder(
485        queue: Arc<Queue>,
486        object_store: Arc<dyn ObjectStore>,
487        runner: R,
488        terminal_hook: H,
489    ) -> WorkflowRuntimeBuilder<R, H> {
490        let clock = queue.clock();
491        WorkflowRuntimeBuilder {
492            queue,
493            object_store,
494            queue_name: "workflow-steps".to_string(),
495            memo_prefix: None,
496            runner,
497            terminal_hook,
498            max_concurrent_steps: 16,
499            poll_interval: Duration::from_millis(250),
500            memo_retention: None,
501            group_retention: None,
502            step_output_replay: false,
503            clock,
504        }
505    }
506
507    /// Submit a new run. Enqueues step 0 with payload `spec.input`.
508    ///
509    /// Idempotent on `(run_id, spec.input)`: if a run with the same id is
510    /// already active (its durable run record in Taquba's user KV
511    /// namespace exists, whichever process submitted it) and
512    /// `spec.input` matches the original submission, this
513    /// call is a no-op and the returned [`SubmitOutcome`] has
514    /// `newly_submitted = false`. A re-submission of an active `run_id`
515    /// with a *different* input is rejected with [`Error::InputMismatch`];
516    /// pick a fresh `run_id` for a new run.
517    pub async fn submit(&self, spec: RunSpec) -> Result<SubmitOutcome> {
518        self.inner.core.submit(spec).await
519    }
520
521    /// The runtime's view of its store: the read-only queries over the queue
522    /// and the memo store the runtime writes to.
523    pub fn view(&self) -> &WorkflowView {
524        &self.inner.core.view
525    }
526
527    /// [`WorkflowView::status`] through the runtime's view, so the status is
528    /// available after a restart and from any runtime over the same queue.
529    pub async fn status(&self, run_id: &RunId) -> Result<Option<RunStatus>> {
530        self.inner.core.view.status(run_id).await
531    }
532
533    /// [`WorkflowView::outcome`] through the runtime's view.
534    pub async fn outcome(&self, run_id: &RunId) -> Result<Option<RunOutcome>> {
535        self.inner.core.view.outcome(run_id).await
536    }
537
538    /// Wait until the run `run_id` terminates and report its end. The
539    /// wait follows the run's current step across its steps, so it
540    /// answers for a run of any length, after a restart and from any
541    /// runtime over the same queue; a run already terminated is
542    /// reported at once from its records. A step dead-lettered outside
543    /// the worker is terminated by the worker's dead-step
544    /// reconciliation, which the wait polls for at the poll interval.
545    ///
546    /// Returns [`Error::RunNotFound`] for a run the runtime has no
547    /// record of: never submitted, or terminated and swept.
548    pub async fn wait(&self, run_id: &RunId) -> Result<RunEnd> {
549        self.inner.core.wait(run_id).await
550    }
551
552    /// [`Self::wait`] bounded by `timeout`; `Ok(None)` when the timeout
553    /// elapses first.
554    pub async fn wait_timeout(&self, run_id: &RunId, timeout: Duration) -> Result<Option<RunEnd>> {
555        self.inner.core.wait_timeout(run_id, timeout).await
556    }
557
558    /// Request cancellation of an active run.
559    ///
560    /// Returns `Ok(true)` once the request is recorded on the run's
561    /// durable record, or `Ok(false)` if the run is unknown or already
562    /// terminal, including a run whose current step the queue
563    /// dead-lettered outside the worker, which the worker's dead-step
564    /// reconciliation terminates as failed. The request reaches a run
565    /// after a restart and from any runtime over the same queue.
566    ///
567    /// The run terminates as [`TerminalStatus::Cancelled`](crate::TerminalStatus::Cancelled) and its
568    /// notification job is enqueued for the terminal hook:
569    ///
570    /// - **Pending / scheduled step**: the queued step job is removed
571    ///   and the notification enqueued in one transaction before this
572    ///   call returns; the hook runs from a worker afterwards.
573    /// - **Running step**: cancellation is delivered to the runner via
574    ///   [`Delivery::cancel_token`](crate::Delivery::cancel_token); runners
575    ///   that watch the token short-circuit immediately. Runners that ignore
576    ///   the token are allowed to run to completion (futures cannot be safely
577    ///   aborted mid-step). In both cases the runner's [`StepOutcome`] /
578    ///   [`StepError`](crate::StepError) is discarded and the worker settles
579    ///   the run once the step returns, with any pending transient retry
580    ///   suppressed and the step acked rather than nacked.
581    /// - A step claimed after the request is settled as cancelled
582    ///   without running.
583    ///
584    /// Cancellation is best-effort: a run whose terminal step settles while the
585    /// request is recorded keeps the outcome it committed. A request applies
586    /// only to the run it is recorded on, and a later submission of the same
587    /// run id starts without it.
588    pub async fn cancel(&self, run_id: &RunId) -> Result<bool> {
589        self.inner.core.cancel(run_id).await
590    }
591
592    /// The group named `id`.
593    pub fn group(&self, id: RunId) -> RunGroup {
594        RunGroup::new(self.inner.core.clone(), id)
595    }
596
597    /// A group with a generated id.
598    pub fn new_group(&self) -> RunGroup {
599        RunGroup::new(self.inner.core.clone(), RunId::generate())
600    }
601
602    /// Spawn [`Self::run`] as a Tokio task and return a handle for
603    /// graceful shutdown. The worker runs until `shutdown` resolves or
604    /// [`RunnerHandle::shutdown`] is called; in-flight steps finish
605    /// either way.
606    pub fn spawn<F>(&self, shutdown: F) -> RunnerHandle
607    where
608        F: Future<Output = ()> + Send + 'static,
609        R: 'static,
610        H: 'static,
611    {
612        let runtime = self.clone();
613        WorkerHandle::spawn(shutdown, |stop| async move { runtime.run_with(stop).await })
614    }
615
616    /// Drive the step worker loop until `shutdown` resolves. Spawns up
617    /// to `max_concurrent_steps` step processors, the dead-step
618    /// reconciliation that terminates runs whose step the queue
619    /// dead-lettered outside the worker and, when
620    /// [`WorkflowRuntimeBuilder::memo_retention`] is set, a
621    /// memo-retention sweeper, all running in parallel. All halt cleanly
622    /// when `shutdown` resolves or the worker errors.
623    pub async fn run<F>(&self, shutdown: F) -> Result<()>
624    where
625        F: Future<Output = ()>,
626        R: 'static,
627        H: 'static,
628    {
629        let stop = CancellationToken::new();
630        let mut worker = std::pin::pin!(self.run_with(stop.clone()));
631        tokio::select! {
632            res = &mut worker => res,
633            () = shutdown => {
634                stop.cancel();
635                worker.await
636            }
637        }
638    }
639
640    /// [`Self::run`] over a token: the worker and the background loops
641    /// stop when `stop` is cancelled, and `stop` is cancelled when the
642    /// worker returns on its own, so the background loops halt with it.
643    async fn run_with(&self, stop: CancellationToken) -> Result<()>
644    where
645        R: 'static,
646        H: 'static,
647    {
648        let mut background: Vec<_> = self
649            .inner
650            .core
651            .sweeps()
652            .map(|sweep| {
653                let sweep = sweep.clone();
654                let core = self.inner.core.clone();
655                let token = stop.clone();
656                tokio::spawn(async move {
657                    sweep
658                        .run(&core.queue, &*core.clock, core.poll_interval, token)
659                        .await;
660                })
661            })
662            .collect();
663        background.push({
664            let core = self.inner.core.clone();
665            let token = stop.clone();
666            tokio::spawn(async move { core.run_dead_step_reconciliation(token).await })
667        });
668
669        let worker = Arc::new(StepWorker {
670            inner: self.inner.clone(),
671        });
672        let result = taquba::run_worker_concurrent(
673            &self.inner.core.queue,
674            &self.inner.core.queue_name,
675            worker,
676            self.inner.core.max_concurrent_steps,
677            self.inner.core.poll_interval,
678            stop.clone().cancelled_owned(),
679        )
680        .await;
681        stop.cancel();
682
683        for handle in background {
684            let _ = handle.await;
685        }
686
687        result?;
688        Ok(())
689    }
690}
691
692/// A handle to a worker task spawned by [`WorkflowRuntime::spawn`].
693///
694/// Dropping a `RunnerHandle` does not stop the worker: the task
695/// continues until the `shutdown` future passed to `spawn` resolves.
696/// Call [`shutdown`](WorkerHandle::shutdown) or
697/// [`wait`](WorkerHandle::wait) to stop or join the worker explicitly.
698pub type RunnerHandle = WorkerHandle<Result<()>>;
699
700impl RuntimeCore {
701    /// [`WorkflowRuntime::submit`].
702    #[instrument(skip(self, spec), fields(run_id))]
703    pub(crate) async fn submit(&self, spec: RunSpec) -> Result<SubmitOutcome> {
704        let run_id = Self::validate_spec(&spec)?;
705        tracing::Span::current().record("run_id", run_id.as_str());
706        self.enqueue_run(&run_id, spec, None).await
707    }
708
709    /// [`Self::submit`] of a member of a group: the run's step jobs
710    /// hold the membership and its member record is written with the
711    /// enqueue.
712    pub(crate) async fn submit_member(
713        &self,
714        membership: &Membership,
715        spec: RunSpec,
716    ) -> Result<SubmitOutcome> {
717        let run_id = Self::validate_spec(&spec)?;
718        self.enqueue_run(&run_id, spec, Some(membership)).await
719    }
720
721    /// Check `spec`'s headers and KV keys, and return the run id,
722    /// generated when the spec names none.
723    fn validate_spec(spec: &RunSpec) -> Result<RunId> {
724        for k in spec.options.headers.keys() {
725            if k.starts_with(RESERVED_HEADER_PREFIX) {
726                return Err(Error::ReservedHeaderInSubmit(k.clone()));
727            }
728        }
729        let deletes = spec.effects.kv_deletes.iter();
730        for key in spec.effects.kv_writes.keys().chain(deletes) {
731            if key.starts_with(RESERVED_KV_PREFIX.as_bytes()) {
732                return Err(Error::ReservedKvKey(
733                    String::from_utf8_lossy(key).into_owned(),
734                ));
735            }
736        }
737        Ok(spec.run_id.clone().unwrap_or_else(RunId::generate))
738    }
739
740    /// Enqueue step 0 of `run_id` unless the run is active. The run
741    /// record is written with the step-0 enqueue and deleted with the
742    /// termination, so it identifies an active run whichever process
743    /// submitted it, and the step's dedup key serialises two
744    /// submissions that both find no record. A `membership` is set on
745    /// the step job and its pending member record is written with the
746    /// enqueue.
747    async fn enqueue_run(
748        &self,
749        run_id: &RunId,
750        spec: RunSpec,
751        membership: Option<&Membership>,
752    ) -> Result<SubmitOutcome> {
753        let input_hash = hash_input(&spec.input);
754        let duplicate = |job_id: String| SubmitOutcome {
755            run_id: run_id.clone(),
756            newly_submitted: false,
757            job_id,
758        };
759        let check_input = |existing: DurableRunRecord| {
760            if existing.input_hash == input_hash {
761                Ok(())
762            } else {
763                Err(Error::InputMismatch(run_id.clone()))
764            }
765        };
766
767        if let Some(existing) = self.view.run_record(run_id).await? {
768            check_input(existing)?;
769            let current = self.current_step(run_id).await?;
770            return Ok(duplicate(current.job_id));
771        }
772
773        let opts = StepEnqueueOpts {
774            run_at: spec.options.run_at,
775            priority: spec.options.priority,
776            max_attempts: spec.options.max_attempts_per_step,
777            reserved_headers: membership
778                .map(Membership::reserved_headers)
779                .unwrap_or_default(),
780        };
781        let (request, job_id) =
782            self.step_enqueue_request(run_id, 0, spec.input, &spec.options.headers, opts);
783
784        let record_bytes = durable::encode(&DurableRunRecord {
785            run_id: run_id.clone(),
786            submitted_at_ms: self.clock.now_ms(),
787            input_hash,
788            cancel_requested: false,
789        });
790        let mut effects = spec
791            .effects
792            .kv_put(run_kv_key(run_id), record_bytes)
793            .kv_put(step_kv_key(run_id), current_step_bytes(0, &job_id));
794        if let Some(membership) = membership {
795            effects = effects.kv_put(
796                membership.kv_key(),
797                durable::encode(&pending_member(run_id)),
798            );
799        }
800
801        let job_id = match self
802            .queue
803            .enqueue_with_effects(&request.queue, request.payload, request.options, effects)
804            .await?
805            .0
806        {
807            EnqueueResult::New(id) => id,
808            // A concurrent submission committed first; its record holds
809            // the input to check against. A dedup hit without a record
810            // is a store this runtime did not write, reported as a
811            // duplicate.
812            EnqueueResult::AlreadyEnqueued(existing) => {
813                if let Some(record) = self.view.run_record(run_id).await? {
814                    check_input(record)?;
815                }
816                return Ok(duplicate(existing));
817            }
818        };
819
820        debug!(run_id = %run_id, job_id = %job_id, "run submitted");
821        Ok(SubmitOutcome {
822            run_id: run_id.clone(),
823            newly_submitted: true,
824            job_id,
825        })
826    }
827
828    /// [`WorkflowRuntime::wait`].
829    pub(crate) async fn wait(&self, run_id: &RunId) -> Result<RunEnd> {
830        self.wait_run(run_id)
831            .await?
832            .ok_or_else(|| Error::RunNotFound(run_id.clone()))
833    }
834
835    /// [`WorkflowRuntime::wait_timeout`].
836    pub(crate) async fn wait_timeout(
837        &self,
838        run_id: &RunId,
839        timeout: Duration,
840    ) -> Result<Option<RunEnd>> {
841        match tokio::time::timeout(timeout, self.wait(run_id)).await {
842            Ok(end) => end.map(Some),
843            Err(_) => Ok(None),
844        }
845    }
846
847    /// [`WorkflowRuntime::cancel`].
848    pub(crate) async fn cancel(&self, run_id: &RunId) -> Result<bool> {
849        let Some(input_hash) = self.request_cancel(run_id).await? else {
850            return Ok(false);
851        };
852        // Settle the current step now: remove it while it is queued, or
853        // fire the claim's cancellation token, the parent of
854        // `Delivery::cancel_token`, while it runs. A step that settles in
855        // between is followed to its successor; a step claimed after the
856        // request terminates the run on its own.
857        loop {
858            let Some((_, job)) = self.view.current_job(run_id).await? else {
859                // Terminated on its own after the request was recorded.
860                return Ok(false);
861            };
862            if job.status == JobStatus::Dead {
863                // Dead-lettered by the queue outside the worker;
864                // reconciliation terminates the run and the request is
865                // not honoured.
866                return Ok(false);
867            }
868            let claimed = ClaimedStep::parse(&job)?;
869            // `error` is `None`: external cancellation supplies no reason
870            // at the API level. The effects are built before the outcome
871            // is known; the queue applies them only on `Removed`.
872            let outcome = claimed.cancelled(None);
873            let termination = self.termination(&outcome, None, input_hash);
874            let effects = self.terminate_collecting_effects(&outcome, &claimed, termination);
875            match self.queue.cancel_with(&job.id, effects).await?.0 {
876                taquba::CancelOutcome::Removed | taquba::CancelOutcome::Requested => {
877                    return Ok(true);
878                }
879                taquba::CancelOutcome::NotFound => continue,
880            }
881        }
882    }
883
884    /// Settle a run into its terminal state: return the deletes of the
885    /// durable run record and the current-step pointer, the writes of
886    /// the terminal record and the terminal marker (when memo retention
887    /// is enabled), the write of the member record (when the run is a
888    /// group member) and the terminal-notification enqueue
889    /// (when the hook observes this outcome) as [`SettlementEffects`]
890    /// for the settlement transaction. The notification job's payload
891    /// is the committed outcome and the configured [`TerminalHook`]
892    /// runs as its worker; `terminal_step` is the step that produced
893    /// the outcome, or the pending step a cancellation removes, and
894    /// `termination` the record of the outcome the settlement writes.
895    ///
896    /// The effects are pure: nothing is written and no state is
897    /// mutated here, so a caller that builds them and then commits a
898    /// non-terminal outcome leaves no trace. A settlement that fails
899    /// redelivers the step, which re-terminates and rebuilds the same
900    /// effects.
901    pub(crate) fn terminate_collecting_effects(
902        &self,
903        outcome: &RunOutcome,
904        terminal_step: &ClaimedStep<'_>,
905        termination: DurableTermination,
906    ) -> SettlementEffects {
907        let terminated_at_ms = termination.terminated_at_ms;
908        let kv_deletes = vec![run_kv_key(&outcome.run_id), step_kv_key(&outcome.run_id)];
909        let mut kv_writes = HashMap::new();
910        kv_writes.insert(
911            outcome_kv_key(&outcome.run_id),
912            durable::encode(&termination),
913        );
914        if let Some(membership) = &terminal_step.membership {
915            kv_writes.insert(
916                membership.kv_key(),
917                durable::encode(&terminated_member(&outcome.run_id, termination)),
918            );
919        }
920        let enqueues = if (self.observes)(outcome) {
921            vec![self.notification_enqueue_request(outcome, Some(terminal_step.job))]
922        } else {
923            Vec::new()
924        };
925        let effects = SettlementEffects::default()
926            .enqueues(enqueues)
927            .kv_writes(kv_writes)
928            .kv_deletes(kv_deletes);
929        match &self.memo_sweep {
930            Some(sweep) => sweep.mark(effects, &outcome.run_id, terminated_at_ms),
931            None => effects,
932        }
933    }
934
935    /// Terminate every run whose step job the queue dead-lettered
936    /// outside the worker path: a lease that expired past the attempt
937    /// limit, or a claim dead-lettered by crash recovery when the queue
938    /// was opened. Such a settlement runs no workflow code, so the run
939    /// record and the current-step pointer survive it and no
940    /// notification is enqueued. A dead step job that the run's
941    /// current-step pointer names identifies the case exactly, because
942    /// every worker-path dead-letter deletes the pointer in its own
943    /// transaction and a re-submission of the run id names a new job.
944    /// The run terminates as [`TerminalStatus::Failed`] with the queue record's
945    /// last error, through the same effects as a worker-path
946    /// termination committed as one transaction with no transition of
947    /// their own; the runner's failure writes cannot apply, since no
948    /// runner returned. Returns the number of runs terminated.
949    pub(crate) async fn reconcile_dead_steps(&self) -> Result<usize> {
950        const PAGE: usize = 256;
951        let mut terminated = 0usize;
952        let mut dead = std::pin::pin!(self.queue.view().jobs(
953            &self.queue_name,
954            JobStatus::Dead,
955            PAGE
956        ));
957        while let Some(job) = dead.try_next().await? {
958            if job.headers.contains_key(HEADER_TERMINAL) {
959                continue;
960            }
961            let Ok(claimed) = ClaimedStep::parse(&job) else {
962                continue;
963            };
964            let run_id = &claimed.run_id;
965            let current = self.view.current_step_if_active(run_id).await?;
966            if current.is_none_or(|current| current.job_id != job.id) {
967                continue;
968            }
969            // The record is written and deleted with the pointer; a
970            // pointer without one is a store the runtime did not write,
971            // left for the worker to report.
972            let Some(record) = self.view.run_record(run_id).await? else {
973                warn!(run_id = %run_id, job_id = %job.id, "dead step has a current-step pointer but no run record");
974                continue;
975            };
976            let error = job
977                .last_error
978                .clone()
979                .unwrap_or_else(|| "step dead-lettered outside the worker".to_string());
980            let outcome = claimed.failed(error);
981            let termination = self.termination(&outcome, None, record.input_hash);
982            let effects = self.terminate_collecting_effects(&outcome, &claimed, termination);
983            self.queue.commit_effects(effects).await?;
984            warn!(run_id = %run_id, step_number = claimed.step_number, job_id = %job.id, "terminated a run whose step was dead-lettered outside the worker");
985            terminated += 1;
986        }
987        Ok(terminated)
988    }
989
990    /// The reconciliation loop: a pass when the worker starts, then a
991    /// pass whenever the queue's dead count has changed since the last
992    /// successful pass, checked every poll interval, until `stop` is
993    /// cancelled.
994    async fn run_dead_step_reconciliation(&self, stop: CancellationToken) {
995        run_periodically(
996            self.poll_interval,
997            &stop,
998            None,
999            |reconciled_at: Option<i64>| async move {
1000                match self.queue.view().stats(&self.queue_name).await {
1001                    Ok(stats) if reconciled_at != Some(stats.dead) => {
1002                        match self.reconcile_dead_steps().await {
1003                            Ok(_) => Some(stats.dead),
1004                            Err(err) => {
1005                                warn!("dead-step reconciliation failed: {err}");
1006                                reconciled_at
1007                            }
1008                        }
1009                    }
1010                    Ok(_) => reconciled_at,
1011                    Err(err) => {
1012                        warn!("dead-step reconciliation could not read queue stats: {err}");
1013                        reconciled_at
1014                    }
1015                }
1016            },
1017        )
1018        .await;
1019    }
1020
1021    /// The retention sweeps [`WorkflowRuntime::run`] runs.
1022    fn sweeps(&self) -> impl Iterator<Item = &Arc<Sweep>> {
1023        self.memo_sweep.iter().chain(self.group_sweep.iter())
1024    }
1025
1026    /// One pass of every retention sweep. Returns the number of markers
1027    /// removed.
1028    #[cfg(test)]
1029    pub(crate) async fn sweep_once(&self) -> Result<usize> {
1030        let mut removed = 0;
1031        for sweep in self.sweeps() {
1032            removed += sweep.pass(&self.queue, &*self.clock).await?;
1033        }
1034        Ok(removed)
1035    }
1036
1037    /// The current-step pointer of a run whose durable record exists.
1038    pub(crate) async fn current_step(&self, run_id: &RunId) -> Result<DurableCurrentStep> {
1039        self.view
1040            .current_step_if_active(run_id)
1041            .await?
1042            .ok_or_else(|| Error::InconsistentRunState(run_id.clone()))
1043    }
1044
1045    /// Wait until `run_id` terminates, following its current step
1046    /// across steps, and report its end; `None` for a run with no
1047    /// current step and no record. A current step the queue
1048    /// dead-lettered outside the worker is polled at the poll interval
1049    /// until reconciliation terminates the run.
1050    pub(crate) async fn wait_run(&self, run_id: &RunId) -> Result<Option<RunEnd>> {
1051        loop {
1052            let Some((current, _)) = self.view.current_job(run_id).await? else {
1053                return self.run_end(run_id).await;
1054            };
1055            match self.queue.wait_for_completion(&current.job_id).await? {
1056                // `NotFound`: the job settled between the two reads; the
1057                // next read follows the pointer.
1058                WaitOutcome::Done(_) | WaitOutcome::Cancelled | WaitOutcome::NotFound => {}
1059                WaitOutcome::Dead(_) => {
1060                    // A worker-path dead-letter deletes the pointer with
1061                    // its settlement; a pointer that still names the dead
1062                    // job identifies a dead-letter outside the worker,
1063                    // which reconciliation terminates.
1064                    let unreconciled = self
1065                        .view
1066                        .current_step_if_active(run_id)
1067                        .await?
1068                        .is_some_and(|step| step.job_id == current.job_id);
1069                    if unreconciled {
1070                        tokio::time::sleep(self.poll_interval).await;
1071                    }
1072                }
1073            }
1074        }
1075    }
1076
1077    /// The end of the terminated run `run_id` from its terminal record
1078    /// and the run result record of that termination; `None` when no
1079    /// terminal record remains.
1080    async fn run_end(&self, run_id: &RunId) -> Result<Option<RunEnd>> {
1081        let Some(termination) = self.view.terminal_record(run_id).await? else {
1082            return Ok(None);
1083        };
1084        let termination = RunTermination::from(termination);
1085        let outcome = self
1086            .view
1087            .run_result_of(run_id, &termination)
1088            .await?
1089            .map(|result| result.outcome);
1090        Ok(Some(RunEnd {
1091            termination,
1092            outcome,
1093        }))
1094    }
1095
1096    /// The termination of `outcome`'s run at the clock's current time;
1097    /// `input_hash` is the run record's.
1098    pub(crate) fn termination(
1099        &self,
1100        outcome: &RunOutcome,
1101        error_kind: Option<StepErrorKind>,
1102        input_hash: [u8; 32],
1103    ) -> DurableTermination {
1104        DurableTermination {
1105            status: outcome.status.into(),
1106            error: outcome.error.clone(),
1107            error_kind: error_kind.map(DurableErrorKind::from),
1108            final_step: outcome.final_step,
1109            terminated_at_ms: self.clock.now_ms(),
1110            input_hash,
1111        }
1112    }
1113
1114    /// Write the run result record of `outcome`'s run, belonging to
1115    /// `termination`.
1116    pub(crate) async fn store_run_result(
1117        &self,
1118        outcome: &RunOutcome,
1119        termination: &DurableTermination,
1120    ) -> Result<()> {
1121        let record = DurableRunResult {
1122            termination: termination.clone(),
1123            outcome: DurableRunOutcome::from(outcome),
1124        };
1125        self.memo_store
1126            .new_run_memo(&outcome.run_id)
1127            .put(RUN_RESULT_MEMO_KEY, &durable::encode(&record))
1128            .await
1129    }
1130
1131    /// Record a cancellation request on the run record of `run_id`.
1132    /// Returns the record's input hash when the run is active, `None`
1133    /// otherwise; a request already recorded counts as recorded again.
1134    async fn request_cancel(&self, run_id: &RunId) -> Result<Option<[u8; 32]>> {
1135        let key = run_kv_key(run_id);
1136        loop {
1137            let Some(current) = self.queue.view().kv_get(&key).await? else {
1138                return Ok(None);
1139            };
1140            let mut record: DurableRunRecord = durable::decode(&current)?;
1141            if record.cancel_requested {
1142                return Ok(Some(record.input_hash));
1143            }
1144            record.cancel_requested = true;
1145            if self
1146                .queue
1147                .kv_compare_put(&key, Some(&current), &durable::encode(&record))
1148                .await?
1149            {
1150                return Ok(Some(record.input_hash));
1151            }
1152        }
1153    }
1154
1155    /// Build the enqueue request for one step of a run, with a
1156    /// pre-assigned job id so the current-step pointer written with
1157    /// the enqueue can name it. Returns the request and the assigned
1158    /// id.
1159    fn step_enqueue_request(
1160        &self,
1161        run_id: &RunId,
1162        step_number: u32,
1163        payload: Vec<u8>,
1164        user_headers: &HashMap<String, String>,
1165        opts: StepEnqueueOpts,
1166    ) -> (EnqueueRequest, String) {
1167        let job_id = self.queue.next_job_id();
1168        let mut headers = user_headers.clone();
1169        headers.insert(HEADER_RUN_ID.to_string(), run_id.to_string());
1170        headers.insert(HEADER_STEP.to_string(), step_number.to_string());
1171        for (key, value) in &opts.reserved_headers {
1172            headers.insert((*key).to_string(), value.clone());
1173        }
1174
1175        let request = EnqueueRequest {
1176            queue: self.queue_name.clone(),
1177            payload,
1178            options: EnqueueOptions::default()
1179                .headers(headers)
1180                .run_at(opts.run_at)
1181                .priority(opts.priority)
1182                .max_attempts(opts.max_attempts)
1183                .dedup_key(Some(format!("{DEDUP_PREFIX}{run_id}:{step_number}")))
1184                .id_override(Some(job_id.clone())),
1185        };
1186        (request, job_id)
1187    }
1188
1189    /// Build the enqueue request for a run's terminal-notification job.
1190    /// The priority and attempt limit are inherited from `terminal_step`,
1191    /// the job of the step that produced the outcome, when one exists.
1192    fn notification_enqueue_request(
1193        &self,
1194        outcome: &RunOutcome,
1195        terminal_step: Option<&JobRecord>,
1196    ) -> EnqueueRequest {
1197        let payload = durable::encode(&DurableRunOutcome::from(outcome));
1198        let mut headers = HashMap::new();
1199        headers.insert(HEADER_RUN_ID.to_string(), outcome.run_id.to_string());
1200        headers.insert(HEADER_TERMINAL.to_string(), "1".to_string());
1201        EnqueueRequest {
1202            queue: self.queue_name.clone(),
1203            payload,
1204            options: EnqueueOptions::default()
1205                .headers(headers)
1206                .priority(terminal_step.map(|job| job.priority))
1207                .max_attempts(terminal_step.map(|job| job.max_attempts))
1208                .dedup_key(Some(format!("{DEDUP_PREFIX}{}:terminal", outcome.run_id))),
1209        }
1210    }
1211
1212    /// The instant `delay` after the runtime's clock now, as a
1213    /// [`SystemTime`] for an enqueue's `run_at`.
1214    pub(crate) fn run_at_after(&self, delay: Duration) -> SystemTime {
1215        UNIX_EPOCH + Duration::from_millis(self.clock.now_ms()) + delay
1216    }
1217
1218    pub(crate) async fn load_step_output(
1219        &self,
1220        run_id: &RunId,
1221        step_number: u32,
1222        step_payload: &[u8],
1223    ) -> Result<Option<(StepOutcome, StagedEffects)>> {
1224        let Some(bytes) = self
1225            .memo_store
1226            .get_step_output(run_id, step_number, step_payload)
1227            .await?
1228        else {
1229            return Ok(None);
1230        };
1231        let Some(record) = durable::decode_or_absent::<DurableStepOutcomeRecord>(
1232            &bytes,
1233            "step-output replay record",
1234            &format_args!("{run_id}/{step_number}"),
1235        ) else {
1236            return Ok(None);
1237        };
1238        let mut outcome = StepOutcome::from(record.outcome);
1239        match &mut outcome {
1240            StepOutcome::Continue {
1241                when: Trigger::After(delay),
1242                ..
1243            } => {
1244                *delay = remaining_delay(record.stored_at_ms, self.clock.now_ms(), *delay);
1245            }
1246            StepOutcome::Continue {
1247                when: Trigger::OnSignal { timeout, .. },
1248                ..
1249            } => {
1250                *timeout = remaining_delay(record.stored_at_ms, self.clock.now_ms(), *timeout);
1251            }
1252            _ => {}
1253        }
1254        Ok(Some((outcome, record.effects)))
1255    }
1256
1257    pub(crate) async fn store_step_output(
1258        &self,
1259        run_id: &RunId,
1260        step_number: u32,
1261        step_payload: &[u8],
1262        outcome: &StepOutcome,
1263        effects: &StagedEffects,
1264    ) -> Result<()> {
1265        let record = DurableStepOutcomeRecord {
1266            stored_at_ms: self.clock.now_ms(),
1267            outcome: DurableStepOutcome::from(outcome),
1268            effects: effects.clone(),
1269        };
1270        let bytes = rmp_serde::to_vec_named(&record)?;
1271        self.memo_store
1272            .put_step_output(run_id, step_number, step_payload, &bytes)
1273            .await
1274    }
1275
1276    /// Build the effects that advance the run of `claimed` to its next
1277    /// step: the next step's enqueue joins the current step's
1278    /// acknowledgement transaction, so the transition is atomic.
1279    pub(crate) async fn advance(
1280        &self,
1281        claimed: &ClaimedStep<'_>,
1282        payload: Vec<u8>,
1283        opts: StepEnqueueOpts,
1284    ) -> SettlementEffects {
1285        self.advance_with_kv(claimed, payload, opts, |_| HashMap::new())
1286            .await
1287    }
1288
1289    /// [`Self::advance`] with caller KV writes joined to the same
1290    /// acknowledgement transaction. `kv_writes` receives the next step's
1291    /// pre-assigned job id so the writes can reference it.
1292    pub(crate) async fn advance_with_kv(
1293        &self,
1294        claimed: &ClaimedStep<'_>,
1295        payload: Vec<u8>,
1296        opts: StepEnqueueOpts,
1297        kv_writes: impl FnOnce(&str) -> HashMap<Vec<u8>, Vec<u8>>,
1298    ) -> SettlementEffects {
1299        let run_id = &claimed.run_id;
1300        let next_step = claimed.step_number + 1;
1301        let (request, next_job_id) =
1302            self.step_enqueue_request(run_id, next_step, payload, &claimed.headers, opts);
1303        let mut kv_writes = kv_writes(&next_job_id);
1304        kv_writes.insert(
1305            step_kv_key(run_id),
1306            current_step_bytes(next_step, &next_job_id),
1307        );
1308        SettlementEffects::default()
1309            .enqueues(vec![request])
1310            .kv_writes(kv_writes)
1311    }
1312}
1313
1314#[cfg(test)]
1315mod tests {
1316    use super::*;
1317    use crate::durable::DurableMember;
1318    use crate::effects::{EffectsHandle, TerminalEffects};
1319    use crate::group::GroupMember;
1320    use crate::keys::group_member_kv_key;
1321    use crate::keys::{TERMINAL_KV_PREFIX, signal_buf_kv_key, signal_wait_kv_key};
1322    use crate::runner::{Step, StepError};
1323    use crate::signal::SignalOutcome;
1324    use crate::terminal::NoopTerminalHook;
1325    use crate::terminal::TerminalStatus;
1326    use crate::test_util::{
1327        advance, fast_options, open_queue, open_queue_at, open_queue_at_with, open_queue_with, rid,
1328    };
1329    use crate::view::WorkflowView;
1330    use std::sync::Mutex as StdMutex;
1331    use std::sync::atomic::{AtomicU32, Ordering};
1332    use taquba::object_store::ObjectStoreExt;
1333    use taquba::object_store::memory::InMemory;
1334    use taquba::{Expired, ExpiryIndex};
1335    use taquba::{LeaseHandle, MockClock, OpenOptions, QueueConfig, QueueReader};
1336    use tokio::sync::oneshot;
1337
1338    /// Recording terminal hook backed by an mpsc channel.
1339    struct ChannelHook {
1340        tx: tokio::sync::mpsc::UnboundedSender<RunOutcome>,
1341    }
1342
1343    impl TerminalHook for ChannelHook {
1344        async fn on_termination(
1345            &self,
1346            outcome: &RunOutcome,
1347            _effects: &TerminalEffects,
1348        ) -> std::result::Result<(), StepError> {
1349            let _ = self.tx.send(outcome.clone());
1350            Ok(())
1351        }
1352    }
1353
1354    /// Runner that executes a fixed list of step outcomes in order.
1355    struct ScriptedRunner {
1356        script: Arc<StdMutex<Vec<StepOutcome>>>,
1357    }
1358
1359    impl ScriptedRunner {
1360        fn new(steps: Vec<StepOutcome>) -> Self {
1361            Self {
1362                script: Arc::new(StdMutex::new(steps)),
1363            }
1364        }
1365    }
1366
1367    impl StepRunner for ScriptedRunner {
1368        async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1369            let next = self.script.lock().unwrap().remove(0);
1370            Ok(next)
1371        }
1372    }
1373
1374    /// Runner that returns a clone of `result` on every step and counts
1375    /// its calls.
1376    struct FixedRunner {
1377        result: std::result::Result<StepOutcome, StepError>,
1378        calls: Arc<AtomicU32>,
1379    }
1380
1381    impl FixedRunner {
1382        fn new(result: std::result::Result<StepOutcome, StepError>) -> Self {
1383            Self {
1384                result,
1385                calls: Arc::new(AtomicU32::new(0)),
1386            }
1387        }
1388    }
1389
1390    impl StepRunner for FixedRunner {
1391        async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1392            self.calls.fetch_add(1, Ordering::SeqCst);
1393            self.result.clone()
1394        }
1395    }
1396
1397    /// Runner whose step never returns.
1398    struct PauseRunner;
1399
1400    impl StepRunner for PauseRunner {
1401        async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1402            std::future::pending().await
1403        }
1404    }
1405
1406    /// Runner whose step must not be claimed.
1407    struct UnreachableRunner;
1408
1409    impl StepRunner for UnreachableRunner {
1410        async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1411            unreachable!("worker must not claim the step");
1412        }
1413    }
1414
1415    /// Runner that reports the claim of its step, holds the step until
1416    /// its [`Gate`] is released and then returns a clone of `result`.
1417    struct GatedRunner {
1418        claimed: Arc<tokio::sync::Notify>,
1419        release: tokio::sync::Mutex<Option<oneshot::Receiver<()>>>,
1420        result: std::result::Result<StepOutcome, StepError>,
1421        calls: Arc<AtomicU32>,
1422    }
1423
1424    /// The test's side of a [`GatedRunner`].
1425    struct Gate {
1426        claimed: Arc<tokio::sync::Notify>,
1427        release: StdMutex<Option<oneshot::Sender<()>>>,
1428        calls: Arc<AtomicU32>,
1429    }
1430
1431    impl GatedRunner {
1432        fn new(result: std::result::Result<StepOutcome, StepError>) -> (Self, Gate) {
1433            let claimed = Arc::new(tokio::sync::Notify::new());
1434            let calls = Arc::new(AtomicU32::new(0));
1435            let (release_tx, release_rx) = oneshot::channel();
1436            let runner = Self {
1437                claimed: claimed.clone(),
1438                release: tokio::sync::Mutex::new(Some(release_rx)),
1439                result,
1440                calls: calls.clone(),
1441            };
1442            let gate = Gate {
1443                claimed,
1444                release: StdMutex::new(Some(release_tx)),
1445                calls,
1446            };
1447            (runner, gate)
1448        }
1449    }
1450
1451    impl StepRunner for GatedRunner {
1452        async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
1453            self.calls.fetch_add(1, Ordering::SeqCst);
1454            self.claimed.notify_one();
1455            let rx = self
1456                .release
1457                .lock()
1458                .await
1459                .take()
1460                .expect("gate consumed twice");
1461            let _ = rx.await;
1462            self.result.clone()
1463        }
1464    }
1465
1466    impl Gate {
1467        /// Wait until the runner holds the step.
1468        async fn claimed(&self) {
1469            tokio::time::timeout(Duration::from_secs(2), self.claimed.notified())
1470                .await
1471                .expect("runner reached gate");
1472        }
1473
1474        fn release(&self) {
1475            if let Some(tx) = self.release.lock().unwrap().take() {
1476                let _ = tx.send(());
1477            }
1478        }
1479    }
1480
1481    /// Every terminal marker in the queue's KV namespace, as
1482    /// `(run_id, terminal_at_ms)` pairs in key order (oldest first).
1483    async fn terminal_markers(queue: &Queue) -> Vec<(RunId, u64)> {
1484        let page = queue
1485            .view()
1486            .kv_scan(TERMINAL_KV_PREFIX, .., 1_000)
1487            .await
1488            .unwrap();
1489        let index = ExpiryIndex::new(TERMINAL_KV_PREFIX);
1490        page.entries
1491            .iter()
1492            .map(|(key, _)| {
1493                let (at_ms, suffix) = index.parse(key).expect("well-formed marker key");
1494                let id = std::str::from_utf8(suffix).expect("run id");
1495                (RunId::new(id).expect("run id"), at_ms)
1496            })
1497            .collect()
1498    }
1499
1500    /// The terminal status `status` reports for `run_id`; `None` while
1501    /// the run is active or unknown.
1502    async fn terminal_status_of<R: StepRunner, H: TerminalHook>(
1503        runtime: &WorkflowRuntime<R, H>,
1504        run_id: &RunId,
1505    ) -> Option<TerminalStatus> {
1506        match runtime.status(run_id).await.unwrap().map(|s| s.state) {
1507            Some(RunState::Terminated(termination)) => Some(termination.status),
1508            _ => None,
1509        }
1510    }
1511
1512    fn spawn_runtime<R, H>(runtime: WorkflowRuntime<R, H>) -> oneshot::Sender<()>
1513    where
1514        R: StepRunner + 'static,
1515        H: TerminalHook + 'static,
1516    {
1517        let (tx, rx) = oneshot::channel::<()>();
1518        tokio::spawn(async move {
1519            let _ = runtime
1520                .run(async move {
1521                    let _ = rx.await;
1522                })
1523                .await;
1524        });
1525        tx
1526    }
1527
1528    /// A runner that extends its lease through the step and reports the
1529    /// resulting expiry.
1530    struct RenewingRunner {
1531        queue: Arc<Queue>,
1532        tx: tokio::sync::mpsc::UnboundedSender<u64>,
1533    }
1534
1535    impl StepRunner for RenewingRunner {
1536        async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
1537            step.lease
1538                .ensure_at_least(Duration::from_secs(600))
1539                .map_err(|e| StepError::transient(e.to_string()))?;
1540            let expiry = self
1541                .queue
1542                .lease_expiry("workflow-steps", &step.job_id)
1543                .expect("a running step holds a lease");
1544            let _ = self.tx.send(expiry);
1545            Ok(StepOutcome::Succeed { result: Vec::new() })
1546        }
1547    }
1548
1549    #[tokio::test(start_paused = true)]
1550    async fn a_step_runner_extends_its_lease_through_the_step() {
1551        let base = 1_700_000_000_000;
1552        let (queue, store, _clock) = open_queue_at(base).await;
1553        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1554        let (hook_tx, mut hook_rx) = tokio::sync::mpsc::unbounded_channel();
1555        let runtime = WorkflowRuntime::builder(
1556            queue.clone(),
1557            store.clone(),
1558            RenewingRunner { queue, tx },
1559            ChannelHook { tx: hook_tx },
1560        )
1561        .build();
1562        let shutdown = spawn_runtime(runtime.clone());
1563
1564        runtime
1565            .submit(RunSpec {
1566                input: Vec::new(),
1567                ..Default::default()
1568            })
1569            .await
1570            .unwrap();
1571        let expiry = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1572            .await
1573            .unwrap()
1574            .unwrap();
1575        assert!(
1576            expiry >= base + 600_000,
1577            "the extension must reach the lease registry",
1578        );
1579        let outcome = tokio::time::timeout(Duration::from_secs(2), hook_rx.recv())
1580            .await
1581            .unwrap()
1582            .unwrap();
1583        assert_eq!(outcome.status, TerminalStatus::Succeeded);
1584
1585        let _ = shutdown.send(());
1586    }
1587
1588    #[tokio::test(start_paused = true)]
1589    async fn single_step_succeeds_and_writes_no_marker_without_retention() {
1590        let (queue, store) = open_queue().await;
1591        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1592        let runtime = WorkflowRuntime::builder(
1593            queue,
1594            store.clone(),
1595            ScriptedRunner::new(vec![StepOutcome::Succeed {
1596                result: b"done".to_vec(),
1597            }]),
1598            ChannelHook { tx },
1599        )
1600        .build();
1601        let shutdown = spawn_runtime(runtime.clone());
1602
1603        let handle = runtime
1604            .submit(RunSpec {
1605                input: b"in".to_vec(),
1606                ..Default::default()
1607            })
1608            .await
1609            .unwrap();
1610        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1611            .await
1612            .unwrap()
1613            .unwrap();
1614
1615        assert_eq!(outcome.run_id, handle.run_id);
1616        assert_eq!(outcome.status, TerminalStatus::Succeeded);
1617        assert_eq!(outcome.result.as_deref(), Some(b"done".as_slice()));
1618        assert_eq!(outcome.final_step, 0);
1619        assert_eq!(
1620            terminal_status_of(&runtime, &handle.run_id).await,
1621            Some(TerminalStatus::Succeeded),
1622            "the terminal record is written without retention",
1623        );
1624        assert!(terminal_markers(&runtime.inner.core.queue).await.is_empty());
1625
1626        let _ = shutdown.send(());
1627    }
1628
1629    #[tokio::test(start_paused = true)]
1630    async fn multi_step_run_advances_through_continue_with_its_headers() {
1631        let (queue, store) = open_queue().await;
1632        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1633        let runtime = WorkflowRuntime::builder(
1634            queue,
1635            store.clone(),
1636            ScriptedRunner::new(vec![
1637                StepOutcome::continue_now(b"step1".to_vec()),
1638                StepOutcome::continue_now(b"step2".to_vec()),
1639                StepOutcome::Succeed {
1640                    result: b"final".to_vec(),
1641                },
1642            ]),
1643            ChannelHook { tx },
1644        )
1645        .build();
1646        let shutdown = spawn_runtime(runtime.clone());
1647
1648        let handle = runtime
1649            .submit(RunSpec {
1650                input: b"start".to_vec(),
1651                options: RunOptions {
1652                    headers: HashMap::from([
1653                        ("trace_id".to_string(), "abc-123".to_string()),
1654                        ("tenant".to_string(), "acme".to_string()),
1655                    ]),
1656                    ..Default::default()
1657                },
1658                ..Default::default()
1659            })
1660            .await
1661            .unwrap();
1662        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1663            .await
1664            .unwrap()
1665            .unwrap();
1666
1667        assert_eq!(outcome.run_id, handle.run_id);
1668        assert_eq!(outcome.final_step, 2);
1669        assert_eq!(outcome.status, TerminalStatus::Succeeded);
1670        assert_eq!(outcome.result.as_deref(), Some(b"final".as_slice()));
1671        assert_eq!(outcome.headers.get("trace_id").unwrap(), "abc-123");
1672        assert_eq!(outcome.headers.get("tenant").unwrap(), "acme");
1673        assert!(!outcome.headers.contains_key(HEADER_RUN_ID));
1674        assert!(!outcome.headers.contains_key(HEADER_STEP));
1675
1676        let recorded =
1677            runtime.outcome(&handle.run_id).await.unwrap().expect(
1678                "the worker writes the run result record before the terminating settlement",
1679            );
1680        assert_eq!(recorded.status, TerminalStatus::Succeeded);
1681        assert_eq!(recorded.final_step, 2);
1682        assert_eq!(recorded.result.as_deref(), Some(b"final".as_slice()));
1683        assert_eq!(recorded.headers, outcome.headers);
1684
1685        let _ = shutdown.send(());
1686    }
1687
1688    #[tokio::test(start_paused = true)]
1689    async fn continue_after_delays_next_step_until_promotion() {
1690        let initial = 1_700_000_000_000u64;
1691        let (queue, store, clock) = open_queue_at(initial).await;
1692        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1693        let runtime = WorkflowRuntime::builder(
1694            queue.clone(),
1695            store.clone(),
1696            ScriptedRunner::new(vec![
1697                StepOutcome::continue_after(b"step1".to_vec(), Duration::from_secs(60)),
1698                StepOutcome::Succeed {
1699                    result: b"final".to_vec(),
1700                },
1701            ]),
1702            ChannelHook { tx },
1703        )
1704        .build();
1705        let shutdown = spawn_runtime(runtime.clone());
1706
1707        let handle = runtime
1708            .submit(RunSpec {
1709                input: b"start".to_vec(),
1710                ..Default::default()
1711            })
1712            .await
1713            .unwrap();
1714
1715        // The delayed step is held in `scheduled` and the run must not
1716        // terminate while the delay is pending.
1717        assert!(
1718            tokio::time::timeout(Duration::from_millis(500), rx.recv())
1719                .await
1720                .is_err()
1721        );
1722        let stats = queue.view().stats("workflow-steps").await.unwrap();
1723        assert_eq!(stats.scheduled, 1);
1724
1725        advance(&clock, Duration::from_secs(61)).await;
1726        queue.promote_scheduled_now().await.unwrap();
1727
1728        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1729            .await
1730            .unwrap()
1731            .unwrap();
1732        assert_eq!(outcome.run_id, handle.run_id);
1733        assert_eq!(outcome.final_step, 1);
1734        assert_eq!(outcome.status, TerminalStatus::Succeeded);
1735
1736        let _ = shutdown.send(());
1737    }
1738
1739    type ObservedSignals = Arc<StdMutex<Vec<Option<Vec<u8>>>>>;
1740
1741    /// Two-step runner: step 0 continues with a signal wait on
1742    /// `correlation_key`, step 1 records [`Step::signal`] and succeeds.
1743    struct SignalProbe {
1744        correlation_key: String,
1745        timeout: Duration,
1746        observed: ObservedSignals,
1747    }
1748
1749    impl StepRunner for SignalProbe {
1750        async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
1751            if step.step_number == 0 {
1752                Ok(StepOutcome::continue_on_signal(
1753                    Vec::new(),
1754                    self.correlation_key.clone(),
1755                    self.timeout,
1756                ))
1757            } else {
1758                self.observed.lock().unwrap().push(step.signal.clone());
1759                Ok(StepOutcome::Succeed { result: Vec::new() })
1760            }
1761        }
1762    }
1763
1764    async fn wait_for_scheduled(queue: &Queue, count: i64) {
1765        for _ in 0..200 {
1766            if queue
1767                .view()
1768                .stats("workflow-steps")
1769                .await
1770                .unwrap()
1771                .scheduled
1772                == count
1773            {
1774                return;
1775            }
1776            tokio::time::sleep(Duration::from_millis(10)).await;
1777        }
1778        panic!("scheduled count never reached {count}");
1779    }
1780
1781    fn signal_probe_runtime(
1782        queue: Arc<Queue>,
1783        store: Arc<dyn taquba::object_store::ObjectStore>,
1784        correlation_key: &str,
1785        timeout: Duration,
1786    ) -> (
1787        WorkflowRuntime<SignalProbe, ChannelHook>,
1788        ObservedSignals,
1789        tokio::sync::mpsc::UnboundedReceiver<RunOutcome>,
1790    ) {
1791        let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
1792        let observed = Arc::new(StdMutex::new(Vec::new()));
1793        let runtime = WorkflowRuntime::builder(
1794            queue,
1795            store,
1796            SignalProbe {
1797                correlation_key: correlation_key.to_string(),
1798                timeout,
1799                observed: observed.clone(),
1800            },
1801            ChannelHook { tx },
1802        )
1803        .build();
1804        (runtime, observed, rx)
1805    }
1806
1807    #[tokio::test(start_paused = true)]
1808    async fn signal_wakes_waiting_run_early_with_payload() {
1809        let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
1810        let (runtime, observed, mut rx) =
1811            signal_probe_runtime(queue.clone(), store, "order-1", Duration::from_secs(3600));
1812        let shutdown = spawn_runtime(runtime.clone());
1813
1814        runtime
1815            .submit(RunSpec {
1816                input: Vec::new(),
1817                ..Default::default()
1818            })
1819            .await
1820            .unwrap();
1821        wait_for_scheduled(&queue, 1).await;
1822
1823        let outcome = runtime.signal("order-1", b"paid".to_vec()).await.unwrap();
1824        assert_eq!(outcome, SignalOutcome::Delivered);
1825
1826        // The run completes without the timeout elapsing on the mock clock.
1827        let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1828            .await
1829            .unwrap()
1830            .unwrap();
1831        assert_eq!(terminal.status, TerminalStatus::Succeeded);
1832        assert_eq!(
1833            observed.lock().unwrap().as_slice(),
1834            &[Some(b"paid".to_vec())]
1835        );
1836
1837        // Both durable signal entries are consumed.
1838        assert!(
1839            queue
1840                .view()
1841                .kv_get(&signal_wait_kv_key("order-1"))
1842                .await
1843                .unwrap()
1844                .is_none()
1845        );
1846        assert!(
1847            queue
1848                .view()
1849                .kv_get(&signal_buf_kv_key("order-1"))
1850                .await
1851                .unwrap()
1852                .is_none()
1853        );
1854
1855        let _ = shutdown.send(());
1856    }
1857
1858    #[tokio::test(start_paused = true)]
1859    async fn a_run_waiting_on_a_signal_survives_a_close_and_reopen() {
1860        let store: Arc<dyn taquba::object_store::ObjectStore> = Arc::new(InMemory::new());
1861        let open = |store: Arc<dyn taquba::object_store::ObjectStore>| async move {
1862            Arc::new(
1863                Queue::open_with_options(
1864                    store,
1865                    "test",
1866                    OpenOptions::default().clock(Arc::new(MockClock::new(1_700_000_000_000))),
1867                )
1868                .await
1869                .unwrap(),
1870            )
1871        };
1872
1873        let queue = open(store.clone()).await;
1874        let (runtime, _observed, _rx) = signal_probe_runtime(
1875            queue.clone(),
1876            store.clone(),
1877            "approval",
1878            Duration::from_secs(3600),
1879        );
1880        let (stop_tx, stop_rx) = oneshot::channel::<()>();
1881        let worker = tokio::spawn({
1882            let runtime = runtime.clone();
1883            async move {
1884                runtime
1885                    .run(async move {
1886                        let _ = stop_rx.await;
1887                    })
1888                    .await
1889            }
1890        });
1891        runtime
1892            .submit(RunSpec {
1893                input: Vec::new(),
1894                ..Default::default()
1895            })
1896            .await
1897            .unwrap();
1898        wait_for_scheduled(&queue, 1).await;
1899        let _ = stop_tx.send(());
1900        worker.await.unwrap().unwrap();
1901        drop(runtime);
1902        Arc::into_inner(queue)
1903            .expect("no other queue references at close")
1904            .close()
1905            .await
1906            .unwrap();
1907
1908        let queue = open(store.clone()).await;
1909        assert_eq!(
1910            queue
1911                .view()
1912                .stats("workflow-steps")
1913                .await
1914                .unwrap()
1915                .scheduled,
1916            1
1917        );
1918
1919        let (runtime, observed, mut rx) =
1920            signal_probe_runtime(queue.clone(), store, "approval", Duration::from_secs(3600));
1921        let shutdown = spawn_runtime(runtime.clone());
1922
1923        let delivery = runtime
1924            .signal("approval", b"approved".to_vec())
1925            .await
1926            .unwrap();
1927        assert_eq!(delivery, SignalOutcome::Delivered);
1928
1929        let terminal = tokio::time::timeout(Duration::from_secs(5), rx.recv())
1930            .await
1931            .unwrap()
1932            .unwrap();
1933        assert_eq!(terminal.status, TerminalStatus::Succeeded);
1934        assert_eq!(
1935            observed.lock().unwrap().as_slice(),
1936            &[Some(b"approved".to_vec())]
1937        );
1938        let _ = shutdown.send(());
1939    }
1940
1941    #[tokio::test(start_paused = true)]
1942    async fn signal_timeout_delivers_none() {
1943        let (queue, store, clock) = open_queue_at(1_700_000_000_000).await;
1944        let (runtime, observed, mut rx) =
1945            signal_probe_runtime(queue.clone(), store, "order-2", Duration::from_secs(60));
1946        let shutdown = spawn_runtime(runtime.clone());
1947
1948        runtime
1949            .submit(RunSpec {
1950                input: Vec::new(),
1951                ..Default::default()
1952            })
1953            .await
1954            .unwrap();
1955        wait_for_scheduled(&queue, 1).await;
1956
1957        advance(&clock, Duration::from_secs(61)).await;
1958        queue.promote_scheduled_now().await.unwrap();
1959
1960        let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
1961            .await
1962            .unwrap()
1963            .unwrap();
1964        assert_eq!(terminal.status, TerminalStatus::Succeeded);
1965        assert_eq!(observed.lock().unwrap().as_slice(), &[None]);
1966
1967        assert!(
1968            queue
1969                .view()
1970                .kv_get(&signal_wait_kv_key("order-2"))
1971                .await
1972                .unwrap()
1973                .is_none()
1974        );
1975
1976        let _ = shutdown.send(());
1977    }
1978
1979    #[tokio::test(start_paused = true)]
1980    async fn a_buffered_signal_is_consumed_at_registration_and_a_later_one_replaces_it() {
1981        let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
1982        let (runtime, observed, mut rx) =
1983            signal_probe_runtime(queue.clone(), store, "order-3", Duration::from_secs(3600));
1984        let shutdown = spawn_runtime(runtime.clone());
1985
1986        assert_eq!(
1987            runtime.signal("order-3", b"first".to_vec()).await.unwrap(),
1988            SignalOutcome::Buffered
1989        );
1990        assert_eq!(
1991            runtime.signal("order-3", b"second".to_vec()).await.unwrap(),
1992            SignalOutcome::Buffered
1993        );
1994
1995        runtime
1996            .submit(RunSpec {
1997                input: Vec::new(),
1998                ..Default::default()
1999            })
2000            .await
2001            .unwrap();
2002
2003        // The buffered signal is consumed at registration: the run
2004        // completes without any waiting and without the timeout elapsing.
2005        let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
2006            .await
2007            .unwrap()
2008            .unwrap();
2009        assert_eq!(terminal.status, TerminalStatus::Succeeded);
2010        assert_eq!(
2011            observed.lock().unwrap().as_slice(),
2012            &[Some(b"second".to_vec())]
2013        );
2014
2015        assert!(
2016            queue
2017                .view()
2018                .kv_get(&signal_buf_kv_key("order-3"))
2019                .await
2020                .unwrap()
2021                .is_none()
2022        );
2023
2024        let _ = shutdown.send(());
2025    }
2026
2027    #[tokio::test(start_paused = true)]
2028    async fn clear_signal_discards_buffered_signal() {
2029        let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
2030        let (runtime, _observed, _rx) =
2031            signal_probe_runtime(queue.clone(), store, "order-5", Duration::from_secs(60));
2032
2033        assert_eq!(
2034            runtime.signal("order-5", b"stale".to_vec()).await.unwrap(),
2035            SignalOutcome::Buffered
2036        );
2037        assert!(runtime.clear_signal("order-5").await.unwrap());
2038        assert!(!runtime.clear_signal("order-5").await.unwrap());
2039        assert!(
2040            queue
2041                .view()
2042                .kv_get(&signal_buf_kv_key("order-5"))
2043                .await
2044                .unwrap()
2045                .is_none()
2046        );
2047    }
2048
2049    #[tokio::test(start_paused = true)]
2050    async fn duplicate_waiter_registration_fails_the_run() {
2051        let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
2052        let (runtime, _observed, mut rx) =
2053            signal_probe_runtime(queue.clone(), store, "order-6", Duration::from_secs(3600));
2054        let shutdown = spawn_runtime(runtime.clone());
2055
2056        runtime
2057            .submit(RunSpec {
2058                run_id: Some(rid("run-a")),
2059                input: Vec::new(),
2060                ..Default::default()
2061            })
2062            .await
2063            .unwrap();
2064        wait_for_scheduled(&queue, 1).await;
2065
2066        runtime
2067            .submit(RunSpec {
2068                run_id: Some(rid("run-b")),
2069                input: Vec::new(),
2070                ..Default::default()
2071            })
2072            .await
2073            .unwrap();
2074
2075        let terminal = tokio::time::timeout(Duration::from_secs(5), rx.recv())
2076            .await
2077            .unwrap()
2078            .unwrap();
2079        assert_eq!(terminal.run_id, "run-b");
2080        assert_eq!(terminal.status, TerminalStatus::Failed);
2081        assert!(
2082            terminal
2083                .error
2084                .as_deref()
2085                .is_some_and(|e| e.contains("already registered"))
2086        );
2087        assert_eq!(
2088            queue.view().stats("workflow-steps").await.unwrap().dead,
2089            1,
2090            "the rejected registration dead-letters run-b's step",
2091        );
2092        assert!(
2093            queue
2094                .view()
2095                .kv_get(&run_kv_key(&rid("run-b")))
2096                .await
2097                .unwrap()
2098                .is_none(),
2099            "the run record delete rides the dead-letter",
2100        );
2101        assert!(
2102            queue
2103                .view()
2104                .kv_get(&run_kv_key(&rid("run-a")))
2105                .await
2106                .unwrap()
2107                .is_some(),
2108            "the waiting run keeps its record",
2109        );
2110
2111        let _ = shutdown.send(());
2112    }
2113
2114    #[tokio::test(start_paused = true)]
2115    async fn buffered_signal_missed_by_the_wake_is_delivered_at_timeout() {
2116        let (queue, store, clock) = open_queue_at(1_700_000_000_000).await;
2117        let (runtime, observed, mut rx) =
2118            signal_probe_runtime(queue.clone(), store, "order-7", Duration::from_secs(60));
2119        let shutdown = spawn_runtime(runtime.clone());
2120
2121        runtime
2122            .submit(RunSpec {
2123                input: Vec::new(),
2124                ..Default::default()
2125            })
2126            .await
2127            .unwrap();
2128        wait_for_scheduled(&queue, 1).await;
2129
2130        // A signal that was buffered without winning the wake (written
2131        // directly to simulate the settling-registration window).
2132        queue
2133            .kv_put(&signal_buf_kv_key("order-7"), b"late")
2134            .await
2135            .unwrap();
2136
2137        advance(&clock, Duration::from_secs(61)).await;
2138        queue.promote_scheduled_now().await.unwrap();
2139
2140        let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
2141            .await
2142            .unwrap()
2143            .unwrap();
2144        assert_eq!(terminal.status, TerminalStatus::Succeeded);
2145        assert_eq!(
2146            observed.lock().unwrap().as_slice(),
2147            &[Some(b"late".to_vec())]
2148        );
2149
2150        let _ = shutdown.send(());
2151    }
2152
2153    #[tokio::test(start_paused = true)]
2154    async fn cancelled_waiter_leaves_no_live_index_for_the_next_signal() {
2155        let (queue, store, _clock) = open_queue_at(1_700_000_000_000).await;
2156        let (runtime, _observed, mut rx) =
2157            signal_probe_runtime(queue.clone(), store, "order-8", Duration::from_secs(3600));
2158        let shutdown = spawn_runtime(runtime.clone());
2159
2160        let handle = runtime
2161            .submit(RunSpec {
2162                input: Vec::new(),
2163                ..Default::default()
2164            })
2165            .await
2166            .unwrap();
2167        wait_for_scheduled(&queue, 1).await;
2168
2169        assert!(runtime.cancel(&handle.run_id).await.unwrap());
2170        let terminal = tokio::time::timeout(Duration::from_secs(2), rx.recv())
2171            .await
2172            .unwrap()
2173            .unwrap();
2174        assert_eq!(terminal.status, TerminalStatus::Cancelled);
2175
2176        // The signal finds only a stale index entry, cleans it and buffers.
2177        assert_eq!(
2178            runtime.signal("order-8", b"orphan".to_vec()).await.unwrap(),
2179            SignalOutcome::Buffered
2180        );
2181        assert!(
2182            queue
2183                .view()
2184                .kv_get(&signal_wait_kv_key("order-8"))
2185                .await
2186                .unwrap()
2187                .is_none()
2188        );
2189
2190        let _ = shutdown.send(());
2191    }
2192
2193    #[tokio::test(start_paused = true)]
2194    async fn a_failure_notification_inherits_the_step_limits() {
2195        let (queue, store) = open_queue().await;
2196        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
2197        let runtime = WorkflowRuntime::builder(
2198            queue.clone(),
2199            store,
2200            FixedRunner::new(Err(StepError::permanent("nope"))),
2201            ChannelHook { tx },
2202        )
2203        .build();
2204        runtime
2205            .submit(RunSpec {
2206                input: b"x".to_vec(),
2207                options: RunOptions {
2208                    priority: Some(3),
2209                    max_attempts_per_step: Some(5),
2210                    ..Default::default()
2211                },
2212                ..Default::default()
2213            })
2214            .await
2215            .unwrap();
2216
2217        let job = queue
2218            .claim("workflow-steps", Duration::from_secs(30))
2219            .await
2220            .unwrap()
2221            .unwrap();
2222        let err = runtime
2223            .inner
2224            .process_step(&job, &LeaseHandle::detached())
2225            .await
2226            .unwrap_err();
2227        let failure = err
2228            .downcast_ref::<taquba::FailWith>()
2229            .expect("a terminating failure carries its effects");
2230        let notification = &failure.effects.enqueues[0];
2231        assert_eq!(notification.options.priority, Some(3));
2232        assert_eq!(notification.options.max_attempts, Some(5));
2233    }
2234
2235    #[tokio::test(start_paused = true)]
2236    async fn a_duplicate_submit_is_idempotent_drops_its_effects_and_rejects_a_changed_input() {
2237        let (queue, store) = open_queue().await;
2238        let runtime = WorkflowRuntime::builder(
2239            queue.clone(),
2240            store,
2241            ScriptedRunner::new(vec![]),
2242            NoopTerminalHook,
2243        )
2244        .build();
2245        // No worker loop runs, so the step stays queued and the run is
2246        // active for every later submit.
2247        let spec = |input: &[u8], key: &[u8]| RunSpec {
2248            run_id: Some(rid("fixed-id")),
2249            input: input.to_vec(),
2250            effects: SettlementEffects::default().kv_put(key, b"1"),
2251            ..Default::default()
2252        };
2253
2254        let first = runtime.submit(spec(b"x", b"app/first")).await.unwrap();
2255        assert!(first.newly_submitted);
2256        assert!(runtime.status(&rid("fixed-id")).await.unwrap().is_some());
2257        assert_eq!(
2258            queue.view().kv_get(b"app/first").await.unwrap().as_deref(),
2259            Some(b"1".as_slice())
2260        );
2261
2262        let duplicate = runtime.submit(spec(b"x", b"app/second")).await.unwrap();
2263        assert_eq!(duplicate.run_id, "fixed-id");
2264        assert!(!duplicate.newly_submitted);
2265        assert_eq!(duplicate.job_id, first.job_id);
2266        assert!(queue.view().kv_get(b"app/second").await.unwrap().is_none());
2267
2268        let err = runtime.submit(spec(b"y", b"app/third")).await.unwrap_err();
2269        assert!(matches!(&err, Error::InputMismatch(id) if id == "fixed-id"));
2270        assert!(err.is_permanent());
2271        assert!(queue.view().kv_get(b"app/third").await.unwrap().is_none());
2272    }
2273
2274    #[tokio::test(start_paused = true)]
2275    async fn a_duplicate_known_only_from_the_durable_record_reports_the_current_job() {
2276        let (queue, store) = open_queue().await;
2277        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
2278        // Step 0 continues into a step scheduled an hour out, so the run
2279        // rests at step 1 with a job the second runtime never saw.
2280        let first = WorkflowRuntime::builder(
2281            queue.clone(),
2282            store.clone(),
2283            ScriptedRunner::new(vec![StepOutcome::continue_after(
2284                b"next".to_vec(),
2285                Duration::from_secs(3600),
2286            )]),
2287            ChannelHook { tx },
2288        )
2289        .build();
2290        let shutdown = spawn_runtime(first.clone());
2291        let submitted = first
2292            .submit(RunSpec {
2293                run_id: Some(rid("durable")),
2294                input: b"x".to_vec(),
2295                ..Default::default()
2296            })
2297            .await
2298            .unwrap();
2299        for _ in 0..200 {
2300            if queue
2301                .view()
2302                .stats("workflow-steps")
2303                .await
2304                .unwrap()
2305                .scheduled
2306                == 1
2307            {
2308                break;
2309            }
2310            tokio::time::sleep(Duration::from_millis(10)).await;
2311        }
2312        let _ = shutdown.send(());
2313
2314        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
2315        let second = WorkflowRuntime::builder(
2316            queue.clone(),
2317            store,
2318            ScriptedRunner::new(vec![]),
2319            ChannelHook { tx },
2320        )
2321        .build();
2322        let duplicate = second
2323            .submit(RunSpec {
2324                run_id: Some(rid("durable")),
2325                input: b"x".to_vec(),
2326                ..Default::default()
2327            })
2328            .await
2329            .unwrap();
2330        assert!(!duplicate.newly_submitted);
2331        assert_ne!(
2332            duplicate.job_id, submitted.job_id,
2333            "the pointer moved to step 1"
2334        );
2335        let step_1 = queue
2336            .view()
2337            .get_job(&duplicate.job_id)
2338            .await
2339            .unwrap()
2340            .unwrap();
2341        assert_eq!(step_1.status, taquba::JobStatus::Scheduled);
2342        assert_eq!(
2343            step_1.headers.get(HEADER_STEP).map(String::as_str),
2344            Some("1")
2345        );
2346    }
2347
2348    #[tokio::test(start_paused = true)]
2349    async fn a_run_submitted_with_run_at_stays_scheduled_until_then() {
2350        struct Echo;
2351        impl StepRunner for Echo {
2352            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
2353                Ok(StepOutcome::Succeed {
2354                    result: step.payload.clone(),
2355                })
2356            }
2357        }
2358
2359        let t0 = 1_700_000_000_000;
2360        let (queue, store, clock) = open_queue_at(t0).await;
2361        let runtime =
2362            WorkflowRuntime::builder(queue.clone(), store, Echo, NoopTerminalHook).build();
2363        let shutdown = spawn_runtime(runtime.clone());
2364
2365        runtime
2366            .submit(RunSpec {
2367                input: b"x".to_vec(),
2368                options: RunOptions {
2369                    run_at: Some(UNIX_EPOCH + Duration::from_millis(t0 + 60_000)),
2370                    ..Default::default()
2371                },
2372                ..Default::default()
2373            })
2374            .await
2375            .unwrap();
2376        let scheduled = queue
2377            .view()
2378            .list_jobs("workflow-steps", taquba::JobStatus::Scheduled, None, 10)
2379            .await
2380            .unwrap()
2381            .jobs;
2382        assert_eq!(scheduled.len(), 1);
2383        let job_id = scheduled[0].id.clone();
2384
2385        let waiter = tokio::spawn({
2386            let queue = queue.clone();
2387            async move { queue.wait_for_completion(&job_id).await }
2388        });
2389        advance(&clock, Duration::from_secs(120)).await;
2390        assert!(matches!(
2391            waiter.await.unwrap().unwrap(),
2392            taquba::WaitOutcome::Done(_),
2393        ));
2394
2395        let _ = shutdown.send(());
2396    }
2397
2398    #[tokio::test(start_paused = true)]
2399    async fn a_step_reports_its_attempt_limit() {
2400        struct Recording(Arc<std::sync::Mutex<Option<u32>>>);
2401        impl StepRunner for Recording {
2402            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
2403                *self.0.lock().unwrap() = Some(step.max_attempts);
2404                Ok(StepOutcome::Succeed { result: Vec::new() })
2405            }
2406        }
2407
2408        let seen = Arc::new(std::sync::Mutex::new(None));
2409        let (queue, store) = open_queue().await;
2410        let runtime = WorkflowRuntime::builder(
2411            queue.clone(),
2412            store,
2413            Recording(seen.clone()),
2414            NoopTerminalHook,
2415        )
2416        .build();
2417        let shutdown = spawn_runtime(runtime.clone());
2418
2419        let outcome = runtime
2420            .submit(RunSpec {
2421                input: b"x".to_vec(),
2422                options: RunOptions {
2423                    max_attempts_per_step: Some(7),
2424                    ..Default::default()
2425                },
2426                ..Default::default()
2427            })
2428            .await
2429            .unwrap();
2430        let job_id = outcome.job_id;
2431        queue.wait_for_completion(&job_id).await.unwrap();
2432        assert_eq!(*seen.lock().unwrap(), Some(7));
2433
2434        let _ = shutdown.send(());
2435    }
2436
2437    #[tokio::test(start_paused = true)]
2438    async fn concurrent_submits_of_one_run_admit_one_and_reject_a_changed_input() {
2439        let (queue, store) = open_queue().await;
2440        let runtime =
2441            WorkflowRuntime::builder(queue, store.clone(), PauseRunner, NoopTerminalHook).build();
2442        let spec = |input: &[u8]| RunSpec {
2443            run_id: Some(rid("raced")),
2444            input: input.to_vec(),
2445            ..Default::default()
2446        };
2447
2448        let (first, same, changed) = tokio::join!(
2449            runtime.submit(spec(b"x")),
2450            runtime.submit(spec(b"x")),
2451            runtime.submit(spec(b"y")),
2452        );
2453        let first = first.unwrap();
2454        let same = same.unwrap();
2455        assert!(first.newly_submitted);
2456        assert!(!same.newly_submitted);
2457        assert_eq!(same.job_id, first.job_id);
2458        assert!(matches!(changed, Err(Error::InputMismatch(id)) if id == "raced"));
2459    }
2460
2461    #[tokio::test(start_paused = true)]
2462    async fn restart_resumes_at_next_step() {
2463        // Headline durability test: after step 0 has acked and step 1 is in
2464        // the queue, kill runtime A entirely and spawn runtime B on the same
2465        // Queue handle. B should claim and complete step 1 without re-running
2466        // step 0.
2467        //
2468        // To make this race-free we gate step 0's runner: the test holds the
2469        // gate while signalling shutdown to A so A enters drain mode without
2470        // ever claiming step 1. Then the gate is opened, A's spawned step-0
2471        // task finishes (enqueueing step 1 + acking step 0) and A exits.
2472        struct CompleteOnStep1;
2473        impl StepRunner for CompleteOnStep1 {
2474            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
2475                assert_eq!(step.step_number, 1, "runtime B should only ever see step 1");
2476                assert_eq!(step.payload.as_slice(), b"step1-payload");
2477                Ok(StepOutcome::Succeed {
2478                    result: b"resumed".to_vec(),
2479                })
2480            }
2481        }
2482
2483        let (queue, store) = open_queue().await;
2484
2485        let (runner, gate) =
2486            GatedRunner::new(Ok(StepOutcome::continue_now(b"step1-payload".to_vec())));
2487        let runtime_a =
2488            WorkflowRuntime::builder(queue.clone(), store.clone(), runner, NoopTerminalHook)
2489                .max_concurrent_steps(1)
2490                .build();
2491
2492        let (shutdown_a_tx, shutdown_a_rx) = oneshot::channel::<()>();
2493        let worker_a = {
2494            let runtime_a = runtime_a.clone();
2495            tokio::spawn(async move {
2496                let _ = runtime_a
2497                    .run(async move {
2498                        let _ = shutdown_a_rx.await;
2499                    })
2500                    .await;
2501            })
2502        };
2503
2504        let handle = runtime_a
2505            .submit(RunSpec {
2506                input: b"input".to_vec(),
2507                ..Default::default()
2508            })
2509            .await
2510            .unwrap();
2511
2512        gate.claimed().await;
2513        let s = runtime_a
2514            .status(&handle.run_id)
2515            .await
2516            .unwrap()
2517            .expect("status");
2518        assert_eq!(s.state, RunState::Running);
2519        assert_eq!(s.current_step, 0);
2520
2521        // A's worker is in the at-capacity select-loop. Signal shutdown
2522        // first, then open the gate so step 0 finishes processing inside
2523        // drain mode (A will not claim step 1).
2524        let _ = shutdown_a_tx.send(());
2525        gate.release();
2526
2527        worker_a.await.expect("runtime A drained cleanly");
2528
2529        // Bring up runtime B on the same Queue handle. It should pick up
2530        // step 1 from where A left off.
2531        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2532        let runtime_b =
2533            WorkflowRuntime::builder(queue, store.clone(), CompleteOnStep1, ChannelHook { tx })
2534                .build();
2535        let shutdown_b = spawn_runtime(runtime_b.clone());
2536
2537        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
2538            .await
2539            .expect("hook fired in time")
2540            .expect("hook channel open");
2541
2542        assert_eq!(outcome.run_id, handle.run_id);
2543        assert_eq!(outcome.status, TerminalStatus::Succeeded);
2544        assert_eq!(outcome.result.as_deref(), Some(b"resumed".as_slice()));
2545        assert_eq!(outcome.final_step, 1);
2546
2547        let _ = shutdown_b.send(());
2548    }
2549
2550    #[test]
2551    fn remaining_delay_measures_from_stored_timestamp() {
2552        let delay = Duration::from_secs(10);
2553        assert_eq!(remaining_delay(1_000, 4_000, delay), Duration::from_secs(7));
2554        assert_eq!(remaining_delay(1_000, 20_000, delay), Duration::ZERO);
2555        assert_eq!(remaining_delay(5_000, 1_000, delay), delay);
2556    }
2557
2558    #[tokio::test(start_paused = true)]
2559    async fn step_output_replay_skips_runner_after_crash_before_ack() {
2560        let (queue, store) = open_queue().await;
2561        let calls = Arc::new(AtomicU32::new(0));
2562        let runtime = WorkflowRuntime::builder(
2563            queue.clone(),
2564            store.clone(),
2565            FixedRunner {
2566                result: Ok(StepOutcome::continue_now(b"step1-payload".to_vec())),
2567                calls: calls.clone(),
2568            },
2569            NoopTerminalHook,
2570        )
2571        .step_output_replay()
2572        .build();
2573
2574        runtime
2575            .submit(RunSpec {
2576                run_id: Some(rid("replay-run")),
2577                input: b"input".to_vec(),
2578                ..Default::default()
2579            })
2580            .await
2581            .unwrap();
2582
2583        let job = queue
2584            .claim("workflow-steps", Duration::from_secs(30))
2585            .await
2586            .unwrap()
2587            .unwrap();
2588
2589        // Simulate a crash after the step output was stored but before
2590        // the settlement committed: discard the effects of the first
2591        // delivery so nothing is enqueued.
2592        let _ = runtime
2593            .inner
2594            .process_step(&job, &LeaseHandle::detached())
2595            .await
2596            .unwrap();
2597        assert_eq!(calls.load(Ordering::SeqCst), 1);
2598
2599        // Re-processing the same claimed record replays the stored step
2600        // outcome without invoking the runner a second time.
2601        let effects = runtime
2602            .inner
2603            .process_step(&job, &LeaseHandle::detached())
2604            .await
2605            .unwrap();
2606        assert_eq!(calls.load(Ordering::SeqCst), 1);
2607        queue.ack_with(&job, effects).await.unwrap();
2608
2609        let next = queue
2610            .claim("workflow-steps", Duration::from_secs(30))
2611            .await
2612            .unwrap()
2613            .unwrap();
2614        assert_eq!(next.payload.as_slice(), b"step1-payload");
2615        assert_eq!(next.headers.get(HEADER_RUN_ID).unwrap(), "replay-run");
2616        assert_eq!(next.headers.get(HEADER_STEP).unwrap(), "1");
2617        assert!(
2618            queue
2619                .claim("workflow-steps", Duration::from_secs(30))
2620                .await
2621                .unwrap()
2622                .is_none(),
2623            "the replayed continue must enqueue step 1 exactly once",
2624        );
2625    }
2626
2627    #[tokio::test(start_paused = true)]
2628    async fn corrupt_step_output_replay_entry_falls_back_to_runner() {
2629        let (queue, store) = open_queue().await;
2630        let calls = Arc::new(AtomicU32::new(0));
2631        let runtime = WorkflowRuntime::builder(
2632            queue.clone(),
2633            store.clone(),
2634            FixedRunner {
2635                result: Ok(StepOutcome::continue_now(b"step1-payload".to_vec())),
2636                calls: calls.clone(),
2637            },
2638            NoopTerminalHook,
2639        )
2640        .step_output_replay()
2641        .build();
2642
2643        runtime
2644            .submit(RunSpec {
2645                run_id: Some(rid("corrupt-run")),
2646                input: b"input".to_vec(),
2647                ..Default::default()
2648            })
2649            .await
2650            .unwrap();
2651
2652        let job = queue
2653            .claim("workflow-steps", Duration::from_secs(30))
2654            .await
2655            .unwrap()
2656            .unwrap();
2657        runtime
2658            .inner
2659            .core
2660            .memo_store
2661            .put_step_output(&rid("corrupt-run"), 0, &job.payload, b"not msgpack")
2662            .await
2663            .unwrap();
2664
2665        runtime
2666            .inner
2667            .process_step(&job, &LeaseHandle::detached())
2668            .await
2669            .unwrap();
2670        assert_eq!(
2671            calls.load(Ordering::SeqCst),
2672            1,
2673            "corrupt entry is treated as a miss",
2674        );
2675
2676        // The recomputed outcome overwrites the corrupt entry, so a
2677        // second delivery replays it without invoking the runner again.
2678        runtime
2679            .inner
2680            .process_step(&job, &LeaseHandle::detached())
2681            .await
2682            .unwrap();
2683        assert_eq!(calls.load(Ordering::SeqCst), 1);
2684    }
2685
2686    #[tokio::test(start_paused = true)]
2687    async fn step_output_replay_of_terminal_outcome_skips_runner() {
2688        let (queue, store) = open_queue().await;
2689        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2690        let calls = Arc::new(AtomicU32::new(0));
2691        let runtime = WorkflowRuntime::builder(
2692            queue.clone(),
2693            store.clone(),
2694            FixedRunner {
2695                result: Ok(StepOutcome::Succeed {
2696                    result: b"final".to_vec(),
2697                }),
2698                calls: calls.clone(),
2699            },
2700            ChannelHook { tx },
2701        )
2702        .step_output_replay()
2703        .build();
2704
2705        runtime
2706            .submit(RunSpec {
2707                run_id: Some(rid("terminal-replay")),
2708                input: b"input".to_vec(),
2709                ..Default::default()
2710            })
2711            .await
2712            .unwrap();
2713        let job = queue
2714            .claim("workflow-steps", Duration::from_secs(30))
2715            .await
2716            .unwrap()
2717            .unwrap();
2718
2719        runtime
2720            .inner
2721            .process_step(&job, &LeaseHandle::detached())
2722            .await
2723            .unwrap();
2724        assert_eq!(calls.load(Ordering::SeqCst), 1);
2725
2726        // Redelivery after a crash before ack: the stored outcome
2727        // settles the run without invoking the runner again.
2728        let effects = runtime
2729            .inner
2730            .process_step(&job, &LeaseHandle::detached())
2731            .await
2732            .unwrap();
2733        assert_eq!(calls.load(Ordering::SeqCst), 1);
2734        queue.ack_with(&job, effects).await.unwrap();
2735
2736        // The committed settlement enqueued one notification; the hook
2737        // observes the replayed outcome when it is processed.
2738        let notification = queue
2739            .claim("workflow-steps", Duration::from_secs(30))
2740            .await
2741            .unwrap()
2742            .unwrap();
2743        let effects = runtime
2744            .inner
2745            .process_step(&notification, &LeaseHandle::detached())
2746            .await
2747            .unwrap();
2748        queue.ack_with(&notification, effects).await.unwrap();
2749        let outcome = rx.recv().await.unwrap();
2750        assert_eq!(outcome.status, TerminalStatus::Succeeded);
2751        assert_eq!(outcome.result.as_deref(), Some(b"final".as_slice()));
2752        assert_eq!(calls.load(Ordering::SeqCst), 1);
2753    }
2754
2755    /// Submits a run whose runner always returns
2756    /// [`StepError::transient`], capped at `max_attempts`. Asserts the
2757    /// runner is invoked exactly `max_attempts` times (per-step max-attempts
2758    /// propagation) and that the terminal hook fires Failed exactly once on
2759    /// the final attempt (fire-once-on-last-attempt logic).
2760    async fn assert_transient_retries_until_max(max_attempts: u32) {
2761        let (queue, store) = open_queue_with(fast_options()).await;
2762        let calls = Arc::new(AtomicU32::new(0));
2763        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2764        let runtime = WorkflowRuntime::builder(
2765            queue.clone(),
2766            store.clone(),
2767            FixedRunner {
2768                result: Err(StepError::transient("flaky")),
2769                calls: calls.clone(),
2770            },
2771            ChannelHook { tx },
2772        )
2773        .memo_retention(Duration::from_secs(60))
2774        .build();
2775        let shutdown = spawn_runtime(runtime.clone());
2776
2777        let handle = runtime
2778            .submit(RunSpec {
2779                input: b"x".to_vec(),
2780                options: RunOptions {
2781                    max_attempts_per_step: Some(max_attempts),
2782                    ..Default::default()
2783                },
2784                ..Default::default()
2785            })
2786            .await
2787            .unwrap();
2788
2789        let outcome = tokio::time::timeout(Duration::from_secs(3), rx.recv())
2790            .await
2791            .expect("hook fired in time")
2792            .expect("hook channel open");
2793
2794        assert_eq!(outcome.status, TerminalStatus::Failed);
2795        assert_eq!(outcome.error.as_deref(), Some("flaky"));
2796        assert_eq!(
2797            calls.load(Ordering::SeqCst),
2798            max_attempts,
2799            "runner called once per attempt up to max_attempts"
2800        );
2801
2802        // Settle window: assert no duplicate hook fires after the terminal one.
2803        tokio::time::sleep(Duration::from_millis(50)).await;
2804        assert!(rx.try_recv().is_err(), "hook fired more than once");
2805
2806        // The notification job was enqueued with the exhausted nack, so
2807        // its effects are committed once the hook fires.
2808        assert_eq!(queue.view().stats("workflow-steps").await.unwrap().dead, 1);
2809        assert!(
2810            queue
2811                .view()
2812                .kv_get(&run_kv_key(&handle.run_id))
2813                .await
2814                .unwrap()
2815                .is_none(),
2816            "the run record delete rides the exhausted nack",
2817        );
2818        assert_eq!(
2819            terminal_markers(&queue)
2820                .await
2821                .iter()
2822                .filter(|(run_id, _)| *run_id == handle.run_id)
2823                .count(),
2824            1,
2825            "the terminal marker rides the exhausted nack",
2826        );
2827
2828        let _ = shutdown.send(());
2829    }
2830
2831    #[tokio::test(start_paused = true)]
2832    async fn a_cancellation_survives_a_restart() {
2833        // Models a restart: the request recorded on the run record and
2834        // the job's persisted `cancel_requested` survive while a fresh
2835        // runtime starts with no process state. The runner returns
2836        // Succeed, so a Cancelled outcome shows the request was read.
2837        let (queue, store, _clock) = open_queue_at(10_000).await;
2838        let (tx_a, _rx_a) = tokio::sync::mpsc::unbounded_channel();
2839        let before = WorkflowRuntime::builder(
2840            queue.clone(),
2841            store.clone(),
2842            ScriptedRunner::new(vec![StepOutcome::Succeed {
2843                result: b"done".to_vec(),
2844            }]),
2845            ChannelHook { tx: tx_a },
2846        )
2847        .build();
2848
2849        let handle = before
2850            .submit(RunSpec {
2851                input: b"x".to_vec(),
2852                ..Default::default()
2853            })
2854            .await
2855            .unwrap();
2856        let claim = queue
2857            .claim("workflow-steps", Duration::from_secs(30))
2858            .await
2859            .unwrap()
2860            .expect("step 0 is claimable");
2861        assert!(before.cancel(&handle.run_id).await.unwrap());
2862
2863        let (tx_b, mut rx_b) = tokio::sync::mpsc::unbounded_channel();
2864        let after = WorkflowRuntime::builder(
2865            queue.clone(),
2866            store.clone(),
2867            ScriptedRunner::new(vec![StepOutcome::Succeed {
2868                result: b"done".to_vec(),
2869            }]),
2870            ChannelHook { tx: tx_b },
2871        )
2872        .build();
2873        assert_eq!(
2874            after.status(&handle.run_id).await.unwrap().map(|s| s.state),
2875            Some(RunState::Cancelling),
2876            "the fresh runtime reads the request from the run record",
2877        );
2878
2879        let effects = after
2880            .inner
2881            .process_step(&claim, &queue.lease_handle(&claim))
2882            .await
2883            .unwrap();
2884        queue.ack_with(&claim, effects).await.unwrap();
2885
2886        let notification = queue
2887            .claim("workflow-steps", Duration::from_secs(30))
2888            .await
2889            .unwrap()
2890            .expect("the terminal notification is claimable");
2891        let effects = after
2892            .inner
2893            .process_step(&notification, &LeaseHandle::detached())
2894            .await
2895            .unwrap();
2896        queue.ack_with(&notification, effects).await.unwrap();
2897
2898        let outcome = rx_b.recv().await.unwrap();
2899        assert_eq!(outcome.status, TerminalStatus::Cancelled);
2900        assert!(outcome.result.is_none(), "the succeed payload is discarded");
2901    }
2902
2903    #[tokio::test(start_paused = true)]
2904    async fn a_cancellation_after_the_settlement_read_reaches_the_next_step() {
2905        // The worker reads the claim's token once the runner has returned.
2906        // A request recorded after that read does not affect the
2907        // advancing settlement; it is read from the run record when the
2908        // next step is claimed, which is then settled without running.
2909        let (queue, store, _clock) = open_queue_at(10_000).await;
2910        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2911        let runtime = WorkflowRuntime::builder(
2912            queue.clone(),
2913            store.clone(),
2914            ScriptedRunner::new(vec![
2915                StepOutcome::Continue {
2916                    payload: b"next".to_vec(),
2917                    when: Trigger::Immediate,
2918                },
2919                StepOutcome::Succeed {
2920                    result: b"done".to_vec(),
2921                },
2922            ]),
2923            ChannelHook { tx },
2924        )
2925        .build();
2926
2927        let handle = runtime
2928            .submit(RunSpec {
2929                input: b"x".to_vec(),
2930                ..Default::default()
2931            })
2932            .await
2933            .unwrap();
2934        let step0 = queue
2935            .claim("workflow-steps", Duration::from_secs(30))
2936            .await
2937            .unwrap()
2938            .expect("step 0 is claimable");
2939        let effects = runtime
2940            .inner
2941            .process_step(&step0, &queue.lease_handle(&step0))
2942            .await
2943            .unwrap();
2944        assert!(runtime.cancel(&handle.run_id).await.unwrap());
2945        queue.ack_with(&step0, effects).await.unwrap();
2946
2947        let step1 = queue
2948            .claim("workflow-steps", Duration::from_secs(30))
2949            .await
2950            .unwrap()
2951            .expect("step 1 is claimable");
2952        let effects = runtime
2953            .inner
2954            .process_step(&step1, &queue.lease_handle(&step1))
2955            .await
2956            .unwrap();
2957        queue.ack_with(&step1, effects).await.unwrap();
2958
2959        let notification = queue
2960            .claim("workflow-steps", Duration::from_secs(30))
2961            .await
2962            .unwrap()
2963            .expect("the terminal notification is claimable");
2964        let effects = runtime
2965            .inner
2966            .process_step(&notification, &LeaseHandle::detached())
2967            .await
2968            .unwrap();
2969        queue.ack_with(&notification, effects).await.unwrap();
2970
2971        let outcome = rx.recv().await.unwrap();
2972        assert_eq!(outcome.status, TerminalStatus::Cancelled);
2973        assert_eq!(outcome.final_step, 1);
2974        assert_eq!(
2975            terminal_status_of(&runtime, &handle.run_id).await,
2976            Some(outcome.status)
2977        );
2978    }
2979
2980    #[tokio::test(start_paused = true)]
2981    async fn cancelling_a_pending_run_commits_its_marker_and_fires_the_hook_once() {
2982        // Pending case: a run sits in the queue, we call `cancel()` before
2983        // any worker claims it. `cancel` removes the step job and enqueues
2984        // the notification before returning.
2985
2986        let (queue, store, _clock) = open_queue_at(10_000).await;
2987        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2988        let runtime = WorkflowRuntime::builder(
2989            queue.clone(),
2990            store.clone(),
2991            UnreachableRunner,
2992            ChannelHook { tx },
2993        )
2994        .memo_retention(Duration::from_secs(60))
2995        .build();
2996        // Note: deliberately do NOT spawn the worker loop, so the submitted
2997        // step stays Pending in the queue while we cancel it.
2998
2999        let mut headers = HashMap::new();
3000        headers.insert("tenant".to_string(), "acme".to_string());
3001
3002        let handle = runtime
3003            .submit(RunSpec {
3004                input: b"x".to_vec(),
3005                options: RunOptions {
3006                    headers,
3007                    ..Default::default()
3008                },
3009                ..Default::default()
3010            })
3011            .await
3012            .unwrap();
3013        let status = runtime
3014            .status(&handle.run_id)
3015            .await
3016            .unwrap()
3017            .expect("active");
3018        assert_eq!(status.state, RunState::Pending);
3019
3020        let was_cancelled = runtime.cancel(&handle.run_id).await.unwrap();
3021        assert!(was_cancelled);
3022        let status = runtime.status(&handle.run_id).await.unwrap().unwrap();
3023        assert_eq!(
3024            status.state,
3025            RunState::Terminated(RunTermination {
3026                status: TerminalStatus::Cancelled,
3027                error: None,
3028                error_kind: None,
3029                final_step: 0,
3030                terminated_at_ms: 10_000,
3031            }),
3032            "the terminal record commits with the removal",
3033        );
3034        assert_eq!(status.current_step, 0);
3035        assert!(
3036            runtime.outcome(&handle.run_id).await.unwrap().is_none(),
3037            "no worker terminated the run, so no run result record exists",
3038        );
3039        assert!(
3040            !runtime.cancel(&handle.run_id).await.unwrap(),
3041            "a second cancel finds no run record",
3042        );
3043
3044        // The marker and the run record's delete commit with the removal,
3045        // before the notification is processed.
3046        let markers = terminal_markers(&queue).await;
3047        assert_eq!(markers, vec![(handle.run_id.clone(), 10_000)]);
3048        assert_eq!(
3049            queue
3050                .view()
3051                .kv_get(&run_kv_key(&handle.run_id))
3052                .await
3053                .unwrap(),
3054            None,
3055        );
3056
3057        // The step job is gone; the one claimable job is the notification.
3058        let notification = queue
3059            .claim("workflow-steps", Duration::from_secs(30))
3060            .await
3061            .unwrap()
3062            .unwrap();
3063        let effects = runtime
3064            .inner
3065            .process_step(&notification, &LeaseHandle::detached())
3066            .await
3067            .unwrap();
3068        queue.ack_with(&notification, effects).await.unwrap();
3069        let outcome = rx.recv().await.unwrap();
3070        assert_eq!(outcome.run_id, handle.run_id);
3071        assert_eq!(outcome.status, TerminalStatus::Cancelled);
3072        // External cancellation carries no reason: `error` is `None`.
3073        assert!(outcome.error.is_none());
3074        assert_eq!(outcome.headers.get("tenant").unwrap(), "acme");
3075        assert!(
3076            queue
3077                .claim("workflow-steps", Duration::from_secs(30))
3078                .await
3079                .unwrap()
3080                .is_none(),
3081            "the cancel enqueues one notification",
3082        );
3083        assert!(rx.try_recv().is_err());
3084
3085        let stats = queue.view().stats("workflow-steps").await.unwrap();
3086        assert_eq!(stats.dead, 0, "cancel must not dead-letter");
3087        assert_eq!(stats.pending, 0, "cancelled job must be removed");
3088    }
3089
3090    #[tokio::test(start_paused = true)]
3091    async fn a_view_reads_the_same_state_through_the_runtime_and_a_reader() {
3092        let (queue, store, _clock) = open_queue_at(10_000).await;
3093        let runtime = WorkflowRuntime::builder(
3094            queue.clone(),
3095            store.clone(),
3096            UnreachableRunner,
3097            NoopTerminalHook,
3098        )
3099        .memo_prefix("memo")
3100        .build();
3101        // The worker loop is not spawned, so the step stays pending.
3102        let handle = runtime.submit(RunSpec::default()).await.unwrap();
3103
3104        let reader = QueueReader::open(store.clone(), "test").await.unwrap();
3105        let view = WorkflowView::new(reader.view().clone(), MemoStore::new(store.clone(), "memo"));
3106        let through_reader = view.status(&handle.run_id).await.unwrap().expect("active");
3107        let through_runtime = runtime.status(&handle.run_id).await.unwrap().unwrap();
3108        assert_eq!(through_reader.run_id, handle.run_id);
3109        assert_eq!(through_reader.state, RunState::Pending);
3110        assert_eq!(through_reader.state, through_runtime.state);
3111        assert_eq!(through_reader.current_step, through_runtime.current_step);
3112        assert!(view.outcome(&handle.run_id).await.unwrap().is_none());
3113        assert!(view.status(&rid("unknown")).await.unwrap().is_none());
3114
3115        assert!(runtime.cancel(&handle.run_id).await.unwrap());
3116        let reader = QueueReader::open(store.clone(), "test").await.unwrap();
3117        let view = WorkflowView::new(reader.view().clone(), MemoStore::new(store, "memo"));
3118        let status = view
3119            .status(&handle.run_id)
3120            .await
3121            .unwrap()
3122            .expect("terminal record");
3123        assert_eq!(
3124            status.state,
3125            RunState::Terminated(RunTermination {
3126                status: TerminalStatus::Cancelled,
3127                error: None,
3128                error_kind: None,
3129                final_step: 0,
3130                terminated_at_ms: 10_000,
3131            }),
3132        );
3133    }
3134
3135    #[tokio::test(start_paused = true)]
3136    async fn the_status_of_a_run_does_not_read_its_offloaded_step_payload() {
3137        let payloads: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
3138        let (queue, store, _clock) = open_queue_at_with(
3139            10_000,
3140            OpenOptions::default()
3141                .payload_offload_threshold(64)
3142                .payload_store(payloads.clone()),
3143        )
3144        .await;
3145        let runtime = WorkflowRuntime::builder(queue, store, UnreachableRunner, NoopTerminalHook)
3146            .memo_prefix("memo")
3147            .build();
3148        let handle = runtime
3149            .submit(RunSpec {
3150                input: vec![7u8; 512],
3151                ..Default::default()
3152            })
3153            .await
3154            .unwrap();
3155
3156        // A read that fetches the step's payload object fails once the object
3157        // is deleted.
3158        let objects: Vec<_> = payloads.list(None).try_collect().await.unwrap();
3159        assert!(!objects.is_empty(), "the step payload is offloaded");
3160        for object in objects {
3161            payloads.delete(&object.location).await.unwrap();
3162        }
3163        let status = runtime.status(&handle.run_id).await.unwrap().unwrap();
3164        assert_eq!(status.state, RunState::Pending);
3165    }
3166
3167    /// Drive a single step that blocks on a gate, calls `cancel(run_id)`
3168    /// while the step is in-flight, and then has the runner return the
3169    /// supplied error. Asserts that external cancellation suppresses the
3170    /// error path entirely: the hook fires `Cancelled` (not `Failed`),
3171    /// no dead-letter is produced regardless of `permanent`/`transient`,
3172    /// and the worker returns `Ok` (no retry, no PermanentFailure
3173    /// propagation).
3174    async fn assert_cancel_suppresses_runner_error(error: StepError) {
3175        let (queue, store) = open_queue_with(fast_options()).await;
3176        let (runner, gate) = GatedRunner::new(Err(error));
3177        let (hook_tx, mut hook_rx) = tokio::sync::mpsc::unbounded_channel();
3178        let runtime = WorkflowRuntime::builder(
3179            queue.clone(),
3180            store.clone(),
3181            runner,
3182            ChannelHook { tx: hook_tx },
3183        )
3184        .build();
3185        let shutdown = spawn_runtime(runtime.clone());
3186
3187        let handle = runtime
3188            .submit(RunSpec {
3189                input: b"x".to_vec(),
3190                ..Default::default()
3191            })
3192            .await
3193            .unwrap();
3194        gate.claimed().await;
3195
3196        let was_cancelled = runtime.cancel(&handle.run_id).await.unwrap();
3197        assert!(was_cancelled);
3198
3199        // Release the runner. It returns Err; without cancellation this
3200        // would either dead-letter (permanent) or nack for retry
3201        // (transient). Cancellation must suppress both.
3202        gate.release();
3203
3204        let outcome = tokio::time::timeout(Duration::from_secs(2), hook_rx.recv())
3205            .await
3206            .expect("hook fired")
3207            .expect("hook channel open");
3208        assert_eq!(outcome.status, TerminalStatus::Cancelled);
3209        assert!(
3210            outcome.error.is_none(),
3211            "external cancel must carry no reason (Some(_) would imply runner-issued StepOutcome::Cancel)",
3212        );
3213        assert_eq!(
3214            terminal_status_of(&runtime, &handle.run_id).await,
3215            Some(outcome.status)
3216        );
3217
3218        // Settle window: assert no retry attempt and no dead-letter or
3219        // duplicate hook fires after the terminal one.
3220        tokio::time::sleep(Duration::from_millis(100)).await;
3221        assert_eq!(
3222            gate.calls.load(Ordering::SeqCst),
3223            1,
3224            "cancellation must suppress retries",
3225        );
3226        let stats = queue.view().stats("workflow-steps").await.unwrap();
3227        assert_eq!(stats.dead, 0, "cancellation must suppress dead-letter");
3228        assert!(
3229            hook_rx.try_recv().is_err(),
3230            "hook must fire exactly once for the cancelled run",
3231        );
3232
3233        let _ = shutdown.send(());
3234    }
3235
3236    #[tokio::test(start_paused = true)]
3237    async fn cancel_suppresses_a_runner_error() {
3238        // Without cancellation, `StepError::permanent` dead-letters the
3239        // step and `StepError::transient` nacks for retry. With an
3240        // external cancel in flight, the worker must ack and fire
3241        // `Cancelled` instead, without re-invoking the runner.
3242        assert_cancel_suppresses_runner_error(StepError::permanent("would-dead-letter")).await;
3243        assert_cancel_suppresses_runner_error(StepError::transient("would-retry")).await;
3244    }
3245
3246    #[tokio::test(start_paused = true)]
3247    async fn cancel_signals_step_token_for_cooperative_short_circuit() {
3248        // A runner that watches `step.cancel_token` should short-circuit
3249        // long after-claim work as soon as `WorkflowRuntime::cancel` is
3250        // called. Without the token, cancellation latency is bounded by
3251        // step duration; with it, the runner returns essentially
3252        // immediately. The test pins this by using a step that would
3253        // otherwise sleep for 30 seconds; if the token didn't fire, the
3254        // test would time out.
3255        struct CooperativeRunner {
3256            claimed: Arc<tokio::sync::Notify>,
3257        }
3258        impl StepRunner for CooperativeRunner {
3259            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
3260                self.claimed.notify_one();
3261                tokio::select! {
3262                    _ = tokio::time::sleep(Duration::from_secs(30)) => {
3263                        Ok(StepOutcome::Succeed { result: b"slow".to_vec() })
3264                    }
3265                    _ = step.cancel_token.cancelled() => {
3266                        Ok(StepOutcome::Cancel { reason: "cooperative".to_string() })
3267                    }
3268                }
3269            }
3270        }
3271
3272        let (queue, store) = open_queue().await;
3273        let claimed = Arc::new(tokio::sync::Notify::new());
3274        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3275        let runtime = WorkflowRuntime::builder(
3276            queue.clone(),
3277            store.clone(),
3278            CooperativeRunner {
3279                claimed: claimed.clone(),
3280            },
3281            ChannelHook { tx },
3282        )
3283        .build();
3284        let shutdown = spawn_runtime(runtime.clone());
3285
3286        let handle = runtime
3287            .submit(RunSpec {
3288                input: b"x".to_vec(),
3289                ..Default::default()
3290            })
3291            .await
3292            .unwrap();
3293        tokio::time::timeout(Duration::from_secs(2), claimed.notified())
3294            .await
3295            .expect("runner observed token");
3296
3297        let was_cancelled = runtime.cancel(&handle.run_id).await.unwrap();
3298        assert!(was_cancelled);
3299
3300        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3301            .await
3302            .expect("hook fired well before the 30s sleep would have")
3303            .expect("hook channel open");
3304
3305        assert_eq!(outcome.status, TerminalStatus::Cancelled);
3306        // Runner-issued Cancel wins precedence over external cancel, so
3307        // the runner's reason surfaces.
3308        assert_eq!(outcome.error.as_deref(), Some("cooperative"));
3309        assert_eq!(
3310            terminal_status_of(&runtime, &handle.run_id).await,
3311            Some(outcome.status)
3312        );
3313
3314        let stats = queue.view().stats("workflow-steps").await.unwrap();
3315        assert_eq!(stats.dead, 0);
3316
3317        let _ = shutdown.send(());
3318    }
3319
3320    #[tokio::test(start_paused = true)]
3321    async fn cancel_returns_false_for_a_terminated_or_unknown_run() {
3322        // Submit a run that succeeds normally, wait for the terminal
3323        // hook, then call `cancel`. The run record was deleted with the
3324        // success, so `cancel` must report `Ok(false)` and must not fire
3325        // a second hook.
3326        let (queue, store) = open_queue().await;
3327        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3328        let runtime = WorkflowRuntime::builder(
3329            queue,
3330            store.clone(),
3331            ScriptedRunner::new(vec![StepOutcome::Succeed {
3332                result: b"done".to_vec(),
3333            }]),
3334            ChannelHook { tx },
3335        )
3336        .build();
3337        let shutdown = spawn_runtime(runtime.clone());
3338
3339        let handle = runtime
3340            .submit(RunSpec {
3341                input: b"x".to_vec(),
3342                ..Default::default()
3343            })
3344            .await
3345            .unwrap();
3346
3347        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3348            .await
3349            .expect("Succeeded hook fired")
3350            .expect("hook channel open");
3351        assert_eq!(outcome.status, TerminalStatus::Succeeded);
3352        assert_eq!(
3353            terminal_status_of(&runtime, &handle.run_id).await,
3354            Some(outcome.status)
3355        );
3356
3357        let was_cancelled = runtime.cancel(&handle.run_id).await.unwrap();
3358        assert!(
3359            !was_cancelled,
3360            "cancel on an already-terminated run must report Ok(false)",
3361        );
3362
3363        tokio::time::sleep(Duration::from_millis(50)).await;
3364        assert!(
3365            rx.try_recv().is_err(),
3366            "no Cancelled hook may fire after the run already terminated as Succeeded",
3367        );
3368
3369        assert!(!runtime.cancel(&rid("never-submitted")).await.unwrap());
3370
3371        let _ = shutdown.send(());
3372    }
3373
3374    #[tokio::test(start_paused = true)]
3375    async fn transient_retries_until_max_attempts() {
3376        assert_transient_retries_until_max(1).await;
3377        assert_transient_retries_until_max(3).await;
3378    }
3379
3380    #[tokio::test(start_paused = true)]
3381    async fn step_memo_survives_across_attempts_of_the_same_step() {
3382        // First attempt writes the memo entry and returns a transient
3383        // error so the runtime retries the step. The second attempt
3384        // reads the same key back and succeeds. This exercises the
3385        // central use case: at-least-once retries of one step should
3386        // short-circuit work the prior attempt already did.
3387        struct MemoRetryRunner;
3388        impl StepRunner for MemoRetryRunner {
3389            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
3390                if step.attempts == 1 {
3391                    step.memo
3392                        .put("cached", b"first-attempt-value")
3393                        .await
3394                        .map_err(|e| StepError::transient(e.to_string()))?;
3395                    return Err(StepError::transient("force a retry"));
3396                }
3397                let got = step
3398                    .memo
3399                    .get("cached")
3400                    .await
3401                    .map_err(|e| StepError::transient(e.to_string()))?;
3402                assert_eq!(got, Some(b"first-attempt-value".to_vec()));
3403                Ok(StepOutcome::Succeed {
3404                    result: got.unwrap_or_default(),
3405                })
3406            }
3407        }
3408
3409        let (queue, store) = open_queue_with(fast_options()).await;
3410        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3411        let runtime =
3412            WorkflowRuntime::builder(queue, store, MemoRetryRunner, ChannelHook { tx }).build();
3413        let shutdown = spawn_runtime(runtime.clone());
3414
3415        runtime
3416            .submit(RunSpec {
3417                input: b"start".to_vec(),
3418                options: RunOptions {
3419                    max_attempts_per_step: Some(3),
3420                    ..Default::default()
3421                },
3422                ..Default::default()
3423            })
3424            .await
3425            .unwrap();
3426        let outcome = tokio::time::timeout(Duration::from_secs(3), rx.recv())
3427            .await
3428            .expect("hook fired in time")
3429            .expect("hook channel open");
3430        assert_eq!(outcome.status, TerminalStatus::Succeeded);
3431        assert_eq!(
3432            outcome.result.as_deref(),
3433            Some(b"first-attempt-value".as_slice())
3434        );
3435
3436        let _ = shutdown.send(());
3437    }
3438
3439    #[tokio::test(start_paused = true)]
3440    async fn terminal_marker_is_written_at_the_runtime_clock() {
3441        // The queue's MockClock is shared into the runtime by default
3442        // (via Queue::clock()), so a `clock.advance` between submit and
3443        // terminate is visible in the marker's terminal_at_ms.
3444        let (queue, store, clock) = open_queue_at(10_000).await;
3445        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3446        let runtime = WorkflowRuntime::builder(
3447            queue,
3448            store.clone(),
3449            ScriptedRunner::new(vec![StepOutcome::Succeed {
3450                result: b"done".to_vec(),
3451            }]),
3452            ChannelHook { tx },
3453        )
3454        .memo_retention(Duration::from_secs(60))
3455        .build();
3456        let shutdown = spawn_runtime(runtime.clone());
3457
3458        let handle = runtime
3459            .submit(RunSpec {
3460                input: b"in".to_vec(),
3461                ..Default::default()
3462            })
3463            .await
3464            .unwrap();
3465        advance(&clock, Duration::from_secs(30)).await;
3466        let _ = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3467            .await
3468            .unwrap()
3469            .unwrap();
3470
3471        let markers = terminal_markers(&runtime.inner.core.queue).await;
3472        assert_eq!(markers.len(), 1);
3473        assert_eq!(markers[0].0, handle.run_id);
3474        // MockClock only moves on explicit advance/set, so the value the
3475        // effects builder reads is exactly the post-advance clock.
3476        assert_eq!(markers[0].1, 10_000 + 30_000);
3477
3478        let _ = shutdown.send(());
3479    }
3480
3481    #[tokio::test(start_paused = true)]
3482    async fn submit_rejects_reserved_headers_and_reserved_kv_keys() {
3483        let (queue, store, _clock) = open_queue_at(10_000).await;
3484        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
3485        let runtime = WorkflowRuntime::builder(
3486            queue,
3487            store,
3488            ScriptedRunner::new(vec![]),
3489            ChannelHook { tx },
3490        )
3491        .build();
3492
3493        let err = runtime
3494            .submit(RunSpec {
3495                input: b"x".to_vec(),
3496                options: RunOptions {
3497                    headers: HashMap::from([("workflow.run_id".to_string(), "evil".to_string())]),
3498                    ..Default::default()
3499                },
3500                ..Default::default()
3501            })
3502            .await
3503            .unwrap_err();
3504        assert!(
3505            matches!(&err, Error::ReservedHeaderInSubmit(k) if k == "workflow.run_id"),
3506            "got: {err:?}"
3507        );
3508
3509        let err = runtime
3510            .submit(RunSpec {
3511                input: Vec::new(),
3512                effects: SettlementEffects::default().kv_put(b"workflow/x", b"v"),
3513                ..Default::default()
3514            })
3515            .await
3516            .unwrap_err();
3517        assert!(matches!(err, Error::ReservedKvKey(_)));
3518
3519        let err = runtime
3520            .submit(RunSpec {
3521                input: Vec::new(),
3522                effects: SettlementEffects::default().kv_delete(b"workflow/x"),
3523                ..Default::default()
3524            })
3525            .await
3526            .unwrap_err();
3527        assert!(matches!(err, Error::ReservedKvKey(_)));
3528    }
3529
3530    #[tokio::test(start_paused = true)]
3531    async fn submit_applies_the_deletes_and_expiry_entries_of_the_spec() {
3532        let (queue, store, _clock) = open_queue_at(10_000).await;
3533        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
3534        let runtime = WorkflowRuntime::builder(
3535            queue.clone(),
3536            store,
3537            ScriptedRunner::new(vec![]),
3538            ChannelHook { tx },
3539        )
3540        .build();
3541        let index = ExpiryIndex::new(b"app/expiry/".to_vec());
3542        queue.kv_put(b"app/stale", b"1").await.unwrap();
3543        // An entry at 9_000 sets the bound of the index, so a pass reads the
3544        // entry at 8_000 only if the submit records it.
3545        queue
3546            .commit_effects(SettlementEffects::default().expiry_entry(&index, 9_000, b"later"))
3547            .await
3548            .unwrap();
3549
3550        runtime
3551            .submit(RunSpec {
3552                input: b"x".to_vec(),
3553                effects: SettlementEffects::default()
3554                    .kv_delete(b"app/stale")
3555                    .expiry_entry(&index, 8_000, b"run"),
3556                ..Default::default()
3557            })
3558            .await
3559            .unwrap();
3560
3561        assert!(queue.view().kv_get(b"app/stale").await.unwrap().is_none());
3562        let mut seen = Vec::new();
3563        index
3564            .pass(&queue, 10_000, Duration::ZERO, |at_ms, suffix| {
3565                seen.push((at_ms, suffix));
3566                std::future::ready(Expired::Delete(SettlementEffects::default()))
3567            })
3568            .await
3569            .unwrap();
3570        assert_eq!(seen, [(8_000, b"run".to_vec()), (9_000, b"later".to_vec())]);
3571    }
3572
3573    #[tokio::test(start_paused = true)]
3574    async fn a_malformed_terminal_marker_is_deleted_without_clearing_memos() {
3575        let (queue, store, clock) = open_queue_at(10_000).await;
3576        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
3577        let runtime = WorkflowRuntime::builder(
3578            queue.clone(),
3579            store.clone(),
3580            ScriptedRunner::new(vec![]),
3581            ChannelHook { tx },
3582        )
3583        .memo_retention(Duration::from_secs(60))
3584        .build();
3585
3586        let memos = MemoStore::new(store, "workflow-steps-memo");
3587        memos
3588            .new_memo(&rid("bystander"), 0)
3589            .put("k", b"expensive")
3590            .await
3591            .unwrap();
3592        // A marker with an empty id, which `RunId` rejects.
3593        let marker = ExpiryIndex::new(TERMINAL_KV_PREFIX).entry_key(0, b"");
3594        queue.kv_put(&marker, b"").await.unwrap();
3595        // A key too short for a time.
3596        let unparseable = [TERMINAL_KV_PREFIX, b"short"].concat();
3597        queue.kv_put(&unparseable, b"").await.unwrap();
3598
3599        advance(&clock, Duration::from_secs(3_600)).await;
3600        runtime.inner.core.sweep_once().await.unwrap();
3601        assert_eq!(
3602            memos.new_memo(&rid("bystander"), 0).get("k").await.unwrap(),
3603            Some(b"expensive".to_vec()),
3604            "an unrelated run's memo entries must survive",
3605        );
3606        assert!(
3607            queue.view().kv_get(&marker).await.unwrap().is_none(),
3608            "the marker is removed and not retried on every sweep",
3609        );
3610        assert!(queue.view().kv_get(&unparseable).await.unwrap().is_none());
3611    }
3612
3613    #[tokio::test(start_paused = true)]
3614    async fn cancelling_a_running_step_overrides_its_outcome_and_writes_its_marker_at_settlement() {
3615        // Between the request and the settlement the run reports
3616        // Cancelling and holds its record with no marker (the queue
3617        // discards the effects on the `Requested` arm); the settlement
3618        // then commits Cancelled in place of the runner's outcome.
3619        let (queue, store, _clock) = open_queue_at(10_000).await;
3620        let (runner, gate) = GatedRunner::new(Ok(StepOutcome::Succeed {
3621            result: b"would-have-succeeded".to_vec(),
3622        }));
3623        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3624        let runtime =
3625            WorkflowRuntime::builder(queue.clone(), store.clone(), runner, ChannelHook { tx })
3626                .memo_retention(Duration::from_secs(60))
3627                .build();
3628        let shutdown = spawn_runtime(runtime.clone());
3629
3630        let handle = runtime
3631            .submit(RunSpec {
3632                input: b"x".to_vec(),
3633                ..Default::default()
3634            })
3635            .await
3636            .unwrap();
3637        gate.claimed().await;
3638        assert_eq!(
3639            runtime
3640                .status(&handle.run_id)
3641                .await
3642                .unwrap()
3643                .expect("active")
3644                .state,
3645            RunState::Running
3646        );
3647
3648        assert!(runtime.cancel(&handle.run_id).await.unwrap());
3649        assert_eq!(
3650            runtime
3651                .status(&handle.run_id)
3652                .await
3653                .unwrap()
3654                .expect("entry retained while termination is in flight")
3655                .state,
3656            RunState::Cancelling
3657        );
3658        assert!(
3659            terminal_markers(&queue).await.is_empty(),
3660            "a run still executing its step must have no terminal marker",
3661        );
3662        assert!(
3663            queue
3664                .view()
3665                .kv_get(&run_kv_key(&handle.run_id))
3666                .await
3667                .unwrap()
3668                .is_some(),
3669            "the run record must survive a cancel the worker has to finish",
3670        );
3671
3672        gate.release();
3673        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3674            .await
3675            .expect("hook fired")
3676            .expect("hook channel open");
3677        assert_eq!(outcome.status, TerminalStatus::Cancelled);
3678        assert!(
3679            outcome.result.is_none(),
3680            "succeed payload must be discarded"
3681        );
3682        let status = runtime.status(&handle.run_id).await.unwrap().unwrap();
3683        assert!(matches!(
3684            status.state,
3685            RunState::Terminated(RunTermination {
3686                status: TerminalStatus::Cancelled,
3687                error: None,
3688                ..
3689            })
3690        ));
3691        assert_eq!(
3692            terminal_markers(&queue).await,
3693            vec![(handle.run_id.clone(), 10_000)],
3694            "the worker's settlement writes it",
3695        );
3696        assert_eq!(queue.view().stats("workflow-steps").await.unwrap().dead, 0);
3697
3698        let _ = shutdown.send(());
3699    }
3700
3701    #[tokio::test(start_paused = true)]
3702    async fn a_permanent_step_error_dead_letters_with_its_marker_and_no_staged_effects() {
3703        let (queue, store, _clock) = open_queue_at(10_000).await;
3704        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3705        let runtime = WorkflowRuntime::builder(
3706            queue.clone(),
3707            store,
3708            EffectStagingRunner::new(vec![Err(StepError::permanent("nope"))]),
3709            ChannelHook { tx },
3710        )
3711        .memo_retention(Duration::from_secs(60))
3712        .build();
3713        let shutdown = spawn_runtime(runtime.clone());
3714
3715        let handle = runtime
3716            .submit(RunSpec {
3717                input: b"x".to_vec(),
3718                ..Default::default()
3719            })
3720            .await
3721            .unwrap();
3722        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3723            .await
3724            .unwrap()
3725            .unwrap();
3726        assert_eq!(outcome.run_id, handle.run_id);
3727        assert_eq!(outcome.status, TerminalStatus::Failed);
3728        assert_eq!(outcome.error.as_deref(), Some("nope"));
3729        let status = runtime.status(&handle.run_id).await.unwrap().unwrap();
3730        assert!(
3731            matches!(
3732                status.state,
3733                RunState::Terminated(RunTermination {
3734                    status: TerminalStatus::Failed,
3735                    error: Some(ref error),
3736                    error_kind: Some(StepErrorKind::Permanent),
3737                    ..
3738                }) if error == "nope"
3739            ),
3740            "the terminal record commits with the dead-letter and carries the error kind",
3741        );
3742        let recorded = runtime.outcome(&handle.run_id).await.unwrap().unwrap();
3743        assert_eq!(recorded.status, TerminalStatus::Failed);
3744        assert_eq!(recorded.error.as_deref(), Some("nope"));
3745
3746        // The notification was enqueued by the dead-letter transaction, so
3747        // the dead job and the marker are already visible, and the staged
3748        // effect was discarded with the failure.
3749        assert_eq!(queue.view().stats("workflow-steps").await.unwrap().dead, 1);
3750        assert_eq!(
3751            terminal_markers(&queue).await,
3752            vec![(handle.run_id.clone(), 10_000)],
3753        );
3754        assert!(queue.view().kv_get(b"app/step-0").await.unwrap().is_none());
3755
3756        let _ = shutdown.send(());
3757    }
3758
3759    #[tokio::test(start_paused = true)]
3760    async fn a_retrying_step_error_commits_no_terminal_marker() {
3761        // A transient failure with attempts left nacks for retry, and
3762        // `nack_with` discards the effects on that branch.
3763        struct FlakyRunner {
3764            attempts: Arc<std::sync::atomic::AtomicUsize>,
3765            clock: MockClock,
3766        }
3767        impl StepRunner for FlakyRunner {
3768            async fn run_step(&self, _step: &Step) -> std::result::Result<StepOutcome, StepError> {
3769                if self
3770                    .attempts
3771                    .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
3772                    == 0
3773                {
3774                    return Err(StepError::transient("flaky"));
3775                }
3776                // Separate the two settlements in clock time before the
3777                // succeeding one: a marker wrongly written for the retry
3778                // would otherwise share this one's key and go unobserved.
3779                self.clock.advance(Duration::from_secs(1));
3780                Ok(StepOutcome::Succeed {
3781                    result: b"done".to_vec(),
3782                })
3783            }
3784        }
3785
3786        let (queue, store, clock) = open_queue_at_with(10_000, fast_options()).await;
3787        let attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
3788        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3789        let runtime = WorkflowRuntime::builder(
3790            queue.clone(),
3791            store.clone(),
3792            FlakyRunner {
3793                attempts: attempts.clone(),
3794                clock: clock.clone(),
3795            },
3796            ChannelHook { tx },
3797        )
3798        .memo_retention(Duration::from_secs(60))
3799        .build();
3800        let shutdown = spawn_runtime(runtime.clone());
3801
3802        let handle = runtime
3803            .submit(RunSpec {
3804                input: b"x".to_vec(),
3805                options: RunOptions {
3806                    max_attempts_per_step: Some(3),
3807                    ..Default::default()
3808                },
3809                ..Default::default()
3810            })
3811            .await
3812            .unwrap();
3813        let outcome = tokio::time::timeout(Duration::from_secs(5), rx.recv())
3814            .await
3815            .unwrap()
3816            .unwrap();
3817        assert_eq!(outcome.status, TerminalStatus::Succeeded);
3818        assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 2);
3819
3820        // Exactly one marker, stamped after the retry's clock advance:
3821        // the retry produced none, the success did.
3822        let markers = terminal_markers(&queue).await;
3823        assert_eq!(markers.len(), 1);
3824        assert_eq!(markers[0].0, handle.run_id);
3825        assert_eq!(markers[0].1, 11_000);
3826
3827        let _ = shutdown.send(());
3828    }
3829
3830    #[tokio::test(start_paused = true)]
3831    async fn the_sweep_clears_only_markers_older_than_the_cutoff() {
3832        // Markers sort by timestamp, so the sweep scans from the start
3833        // of the range and returns at the first unexpired marker.
3834        let (queue, store, _clock) = open_queue_at(10_000).await;
3835        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
3836        let runtime = WorkflowRuntime::builder(
3837            queue.clone(),
3838            store.clone(),
3839            ScriptedRunner::new(vec![]),
3840            ChannelHook { tx },
3841        )
3842        .memo_retention(Duration::from_secs(1))
3843        .build();
3844
3845        let memos = MemoStore::new(store, "workflow-steps-memo");
3846        let sweep = runtime.inner.core.memo_sweep.as_ref().unwrap();
3847        for (run_id, at_ms) in [("old", 1_000u64), ("young", 9_500u64)] {
3848            let run_id = rid(run_id);
3849            memos.new_memo(&run_id, 0).put("k", b"v").await.unwrap();
3850            queue
3851                .commit_effects(sweep.mark(SettlementEffects::default(), &run_id, at_ms))
3852                .await
3853                .unwrap();
3854        }
3855
3856        // Clock at 10_000 with 1s retention: cutoff 9_000.
3857        let cleared = runtime.inner.core.sweep_once().await.unwrap();
3858        assert_eq!(cleared, 1);
3859
3860        let remaining = terminal_markers(&queue).await;
3861        assert_eq!(remaining.len(), 1);
3862        assert_eq!(remaining[0].0, "young");
3863        assert_eq!(memos.new_memo(&rid("old"), 0).get("k").await.unwrap(), None);
3864        assert_eq!(
3865            memos.new_memo(&rid("young"), 0).get("k").await.unwrap(),
3866            Some(b"v".to_vec()),
3867        );
3868    }
3869
3870    #[tokio::test(start_paused = true)]
3871    async fn a_re_submitted_run_shares_its_entries_until_the_first_run_expires() {
3872        let (queue, store, clock) = open_queue_at(10_000).await;
3873        let runtime = WorkflowRuntime::builder(
3874            queue.clone(),
3875            store.clone(),
3876            FixedRunner::new(Ok(StepOutcome::Succeed {
3877                result: b"done".to_vec(),
3878            })),
3879            NoopTerminalHook,
3880        )
3881        .memo_retention(Duration::from_secs(1))
3882        .build();
3883        let spec = RunSpec {
3884            run_id: Some(rid("shared")),
3885            input: b"x".to_vec(),
3886            ..Default::default()
3887        };
3888        runtime.submit(spec.clone()).await.unwrap();
3889        let claim = queue
3890            .claim("workflow-steps", Duration::from_secs(30))
3891            .await
3892            .unwrap()
3893            .unwrap();
3894        let effects = runtime
3895            .inner
3896            .process_step(&claim, &LeaseHandle::detached())
3897            .await
3898            .unwrap();
3899        queue.ack_with(&claim, effects).await.unwrap();
3900        let memos = MemoStore::new(store, "workflow-steps-memo");
3901        memos
3902            .new_memo(&rid("shared"), 0)
3903            .put("k", b"v")
3904            .await
3905            .unwrap();
3906
3907        assert!(runtime.submit(spec).await.unwrap().newly_submitted);
3908        assert_eq!(
3909            memos.new_memo(&rid("shared"), 0).get("k").await.unwrap(),
3910            Some(b"v".to_vec()),
3911            "the second run reads the first run's entry",
3912        );
3913
3914        // The first run's marker expires while the second run is active.
3915        clock.advance(Duration::from_secs(2));
3916        assert_eq!(runtime.inner.core.sweep_once().await.unwrap(), 1);
3917        assert_eq!(
3918            memos.new_memo(&rid("shared"), 0).get("k").await.unwrap(),
3919            None
3920        );
3921        assert_eq!(
3922            runtime
3923                .status(&rid("shared"))
3924                .await
3925                .unwrap()
3926                .map(|s| s.state),
3927            Some(RunState::Pending),
3928            "the second run is still active",
3929        );
3930    }
3931
3932    /// Yield up to `iters` times waiting for `cond` to become true.
3933    /// Used in sweeper tests to let the spawned sweep task make
3934    /// progress between `tokio::time::advance` and the assertion;
3935    /// returns true if the condition held within the budget.
3936    async fn yield_until<F, Fut>(iters: usize, mut cond: F) -> bool
3937    where
3938        F: FnMut() -> Fut,
3939        Fut: Future<Output = bool>,
3940    {
3941        for _ in 0..iters {
3942            if cond().await {
3943                return true;
3944            }
3945            tokio::task::yield_now().await;
3946        }
3947        false
3948    }
3949
3950    #[tokio::test(start_paused = true)]
3951    async fn the_sweeper_clears_a_marker_only_after_retention_elapses() {
3952        // Retention 200ms and a 10ms poll interval. A pass 199ms after
3953        // the marker is written leaves it. Within a poll interval of
3954        // the boundary the sweep loop clears the marker and the run's
3955        // memo entries.
3956        let (queue, store, clock) = open_queue_at(10_000).await;
3957        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3958        let runtime = WorkflowRuntime::builder(
3959            queue,
3960            store.clone(),
3961            ScriptedRunner::new(vec![StepOutcome::Succeed {
3962                result: b"done".to_vec(),
3963            }]),
3964            ChannelHook { tx },
3965        )
3966        .memo_retention(Duration::from_millis(200))
3967        .poll_interval(Duration::from_millis(10))
3968        .build();
3969        let shutdown = spawn_runtime(runtime.clone());
3970
3971        let handle = runtime
3972            .submit(RunSpec {
3973                input: b"in".to_vec(),
3974                ..Default::default()
3975            })
3976            .await
3977            .unwrap();
3978        let _ = tokio::time::timeout(Duration::from_secs(2), rx.recv())
3979            .await
3980            .unwrap()
3981            .unwrap();
3982        let memos = MemoStore::new(store.clone(), "workflow-steps-memo");
3983        memos
3984            .new_memo(&handle.run_id, 0)
3985            .put("k", b"cached")
3986            .await
3987            .unwrap();
3988
3989        advance(&clock, Duration::from_millis(199)).await;
3990        assert_eq!(runtime.inner.core.sweep_once().await.unwrap(), 0);
3991        let markers = terminal_markers(&runtime.inner.core.queue).await;
3992        assert_eq!(
3993            markers.len(),
3994            1,
3995            "a marker within the window must not be swept"
3996        );
3997
3998        advance(&clock, Duration::from_millis(1)).await;
3999        advance(&clock, Duration::from_millis(10)).await;
4000        let cleared = yield_until(50, || async {
4001            terminal_markers(&runtime.inner.core.queue).await.is_empty()
4002        })
4003        .await;
4004        assert!(cleared, "sweeper did not clear the expired marker");
4005        assert_eq!(
4006            memos.new_memo(&handle.run_id, 0).get("k").await.unwrap(),
4007            None,
4008            "sweeper did not clear the run's memo entries",
4009        );
4010        assert!(
4011            runtime.status(&handle.run_id).await.unwrap().is_none(),
4012            "sweeper did not clear the run's terminal record",
4013        );
4014
4015        let _ = shutdown.send(());
4016    }
4017
4018    #[tokio::test(start_paused = true)]
4019    async fn sweeper_keeps_memos_of_runs_without_a_terminal_marker() {
4020        // A run gets a terminal marker only once it terminates, and a
4021        // terminated run never resumes. The sweep is keyed on those
4022        // markers, so an in-flight run's memo entries are never deleted
4023        // out from under a resume, even past the retention window. Here
4024        // a memo entry exists for a run with no terminal marker;
4025        // advancing well past retention must leave it in place.
4026        let (queue, store, clock) = open_queue_at(10_000).await;
4027        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
4028        let runtime = WorkflowRuntime::builder(
4029            queue,
4030            store.clone(),
4031            ScriptedRunner::new(vec![]),
4032            ChannelHook { tx },
4033        )
4034        .memo_retention(Duration::from_millis(100))
4035        .build();
4036        let shutdown = spawn_runtime(runtime.clone());
4037
4038        let memos = MemoStore::new(store.clone(), "workflow-steps-memo");
4039        memos
4040            .new_memo(&rid("in-flight-run"), 0)
4041            .put("k", b"cached")
4042            .await
4043            .unwrap();
4044
4045        advance(&clock, Duration::from_millis(500)).await;
4046        // Give the sweeper several ticks to run against the advanced clock.
4047        for _ in 0..50 {
4048            tokio::task::yield_now().await;
4049        }
4050
4051        assert_eq!(
4052            memos
4053                .new_memo(&rid("in-flight-run"), 0)
4054                .get("k")
4055                .await
4056                .unwrap(),
4057            Some(b"cached".to_vec()),
4058            "sweep must not remove memos of a run with no terminal marker",
4059        );
4060
4061        let _ = shutdown.send(());
4062    }
4063
4064    async fn wait_for_kv(queue: &Queue, key: &[u8]) -> Vec<u8> {
4065        for _ in 0..200 {
4066            if let Some(v) = queue.view().kv_get(key).await.unwrap() {
4067                return v.to_vec();
4068            }
4069            tokio::time::sleep(Duration::from_millis(10)).await;
4070        }
4071        panic!(
4072            "kv key `{}` was never written",
4073            String::from_utf8_lossy(key)
4074        );
4075    }
4076
4077    async fn wait_for_drained(queue: &Queue) {
4078        for _ in 0..200 {
4079            let stats = queue.view().stats("workflow-steps").await.unwrap();
4080            if stats.pending == 0 && stats.claimed == 0 && stats.scheduled == 0 {
4081                return;
4082            }
4083            tokio::time::sleep(Duration::from_millis(10)).await;
4084        }
4085        panic!("the queue never drained");
4086    }
4087
4088    /// Runner that stages one `app/step-{n}` write per step and returns
4089    /// the next scripted result.
4090    struct EffectStagingRunner {
4091        script: Arc<StdMutex<Vec<std::result::Result<StepOutcome, StepError>>>>,
4092    }
4093
4094    impl EffectStagingRunner {
4095        fn new(script: Vec<std::result::Result<StepOutcome, StepError>>) -> Self {
4096            Self {
4097                script: Arc::new(StdMutex::new(script)),
4098            }
4099        }
4100    }
4101
4102    impl StepRunner for EffectStagingRunner {
4103        async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4104            step.effects
4105                .put(format!("app/step-{}", step.step_number), b"done".to_vec())
4106                .map_err(|e| StepError::permanent(e.to_string()))?;
4107            self.script.lock().unwrap().remove(0)
4108        }
4109    }
4110
4111    #[tokio::test(start_paused = true)]
4112    async fn step_effects_commit_with_the_acking_settlement() {
4113        let (queue, store) = open_queue().await;
4114        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4115        let runtime = WorkflowRuntime::builder(
4116            queue.clone(),
4117            store,
4118            EffectStagingRunner::new(vec![
4119                Ok(StepOutcome::continue_now(b"next".to_vec())),
4120                Ok(StepOutcome::Succeed { result: Vec::new() }),
4121            ]),
4122            ChannelHook { tx },
4123        )
4124        .build();
4125        let shutdown = spawn_runtime(runtime.clone());
4126
4127        runtime
4128            .submit(RunSpec {
4129                input: Vec::new(),
4130                ..Default::default()
4131            })
4132            .await
4133            .unwrap();
4134        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4135            .await
4136            .unwrap()
4137            .unwrap();
4138        assert_eq!(outcome.status, TerminalStatus::Succeeded);
4139        assert_eq!(wait_for_kv(&queue, b"app/step-0").await, b"done");
4140        assert_eq!(wait_for_kv(&queue, b"app/step-1").await, b"done");
4141
4142        let _ = shutdown.send(());
4143    }
4144
4145    #[tokio::test(start_paused = true)]
4146    async fn a_step_effect_is_readable_by_the_next_step() {
4147        struct ReadingRunner {
4148            read_under_staging: Arc<StdMutex<Option<Option<Vec<u8>>>>>,
4149        }
4150        impl StepRunner for ReadingRunner {
4151            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4152                let read = step
4153                    .kv
4154                    .get(b"app/marker")
4155                    .await
4156                    .map_err(|e| StepError::permanent(e.to_string()))?
4157                    .map(|b| b.to_vec());
4158                if step.step_number == 0 {
4159                    step.effects
4160                        .put("app/marker", b"v".to_vec())
4161                        .map_err(|e| StepError::permanent(e.to_string()))?;
4162                    let staged_read = step
4163                        .kv
4164                        .get(b"app/marker")
4165                        .await
4166                        .map_err(|e| StepError::permanent(e.to_string()))?
4167                        .map(|b| b.to_vec());
4168                    *self.read_under_staging.lock().unwrap() = Some(staged_read);
4169                    return Ok(StepOutcome::continue_now(Vec::new()));
4170                }
4171                Ok(StepOutcome::Succeed {
4172                    result: read.unwrap_or_default(),
4173                })
4174            }
4175        }
4176
4177        let (queue, store) = open_queue().await;
4178        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4179        let read_under_staging = Arc::new(StdMutex::new(None));
4180        let runtime = WorkflowRuntime::builder(
4181            queue.clone(),
4182            store,
4183            ReadingRunner {
4184                read_under_staging: read_under_staging.clone(),
4185            },
4186            ChannelHook { tx },
4187        )
4188        .build();
4189        let shutdown = spawn_runtime(runtime.clone());
4190
4191        runtime
4192            .submit(RunSpec {
4193                input: Vec::new(),
4194                ..Default::default()
4195            })
4196            .await
4197            .unwrap();
4198        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4199            .await
4200            .unwrap()
4201            .unwrap();
4202        assert_eq!(outcome.status, TerminalStatus::Succeeded);
4203        assert_eq!(outcome.result.as_deref(), Some(b"v".as_slice()));
4204        assert_eq!(*read_under_staging.lock().unwrap(), Some(None));
4205
4206        let _ = shutdown.send(());
4207    }
4208
4209    #[tokio::test(start_paused = true)]
4210    async fn a_run_memo_written_in_one_step_is_readable_in_the_next() {
4211        struct JournalRunner;
4212        impl StepRunner for JournalRunner {
4213            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4214                if step.step_number == 0 {
4215                    step.run_memo.put("journal", b"entry").await?;
4216                    return Ok(StepOutcome::continue_now(Vec::new()));
4217                }
4218                let value = step.run_memo.get("journal").await?;
4219                Ok(StepOutcome::Succeed {
4220                    result: value.unwrap_or_default(),
4221                })
4222            }
4223        }
4224
4225        let (queue, store) = open_queue().await;
4226        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4227        let runtime =
4228            WorkflowRuntime::builder(queue, store, JournalRunner, ChannelHook { tx }).build();
4229        let shutdown = spawn_runtime(runtime.clone());
4230
4231        runtime
4232            .submit(RunSpec {
4233                input: Vec::new(),
4234                ..Default::default()
4235            })
4236            .await
4237            .unwrap();
4238        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4239            .await
4240            .unwrap()
4241            .unwrap();
4242        assert_eq!(outcome.status, TerminalStatus::Succeeded);
4243        assert_eq!(outcome.result.as_deref(), Some(b"entry".as_slice()));
4244
4245        let _ = shutdown.send(());
4246    }
4247
4248    #[tokio::test(start_paused = true)]
4249    async fn a_fail_verdict_acks_with_its_effects_and_no_dead_letter() {
4250        let (queue, store) = open_queue().await;
4251        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4252        let runtime = WorkflowRuntime::builder(
4253            queue.clone(),
4254            store,
4255            EffectStagingRunner::new(vec![Ok(StepOutcome::Fail {
4256                reason: "denied".to_string(),
4257            })]),
4258            ChannelHook { tx },
4259        )
4260        .build();
4261        let shutdown = spawn_runtime(runtime.clone());
4262
4263        let handle = runtime
4264            .submit(RunSpec {
4265                input: Vec::new(),
4266                ..Default::default()
4267            })
4268            .await
4269            .unwrap();
4270        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4271            .await
4272            .unwrap()
4273            .unwrap();
4274        assert_eq!(outcome.run_id, handle.run_id);
4275        assert_eq!(outcome.status, TerminalStatus::Failed);
4276        assert_eq!(outcome.error.as_deref(), Some("denied"));
4277        assert_eq!(
4278            terminal_status_of(&runtime, &handle.run_id).await,
4279            Some(outcome.status)
4280        );
4281        assert_eq!(wait_for_kv(&queue, b"app/step-0").await, b"done");
4282        assert_eq!(
4283            queue.view().stats("workflow-steps").await.unwrap().dead,
4284            0,
4285            "a Fail verdict must not dead-letter"
4286        );
4287
4288        let _ = shutdown.send(());
4289    }
4290
4291    #[tokio::test(start_paused = true)]
4292    async fn a_runner_cancelling_its_own_token_is_not_an_external_cancel() {
4293        // The runner receives a child of the claim's token, so firing it
4294        // leaves the parent uncancelled and the step's staged effects are
4295        // applied.
4296        struct SelfCancellingRunner;
4297        impl StepRunner for SelfCancellingRunner {
4298            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4299                step.cancel_token.cancel();
4300                step.effects
4301                    .put("app/step-0", b"done")
4302                    .map_err(|e| StepError::permanent(e.to_string()))?;
4303                Ok(StepOutcome::Succeed {
4304                    result: b"finished".to_vec(),
4305                })
4306            }
4307        }
4308
4309        let (queue, store, _clock) = open_queue_at(10_000).await;
4310        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4311        let runtime = WorkflowRuntime::builder(
4312            queue.clone(),
4313            store,
4314            SelfCancellingRunner,
4315            ChannelHook { tx },
4316        )
4317        .build();
4318        let shutdown = spawn_runtime(runtime.clone());
4319
4320        runtime
4321            .submit(RunSpec {
4322                input: Vec::new(),
4323                ..Default::default()
4324            })
4325            .await
4326            .unwrap();
4327        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4328            .await
4329            .unwrap()
4330            .unwrap();
4331        assert_eq!(outcome.status, TerminalStatus::Succeeded);
4332        assert_eq!(outcome.result.as_deref(), Some(b"finished".as_slice()));
4333        assert_eq!(wait_for_kv(&queue, b"app/step-0").await, b"done");
4334
4335        let _ = shutdown.send(());
4336    }
4337
4338    #[tokio::test(start_paused = true)]
4339    async fn a_cancel_verdict_acks_with_its_effects_and_no_dead_letter() {
4340        let (queue, store) = open_queue().await;
4341        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4342        let runtime = WorkflowRuntime::builder(
4343            queue.clone(),
4344            store,
4345            EffectStagingRunner::new(vec![Ok(StepOutcome::Cancel {
4346                reason: "obsolete".to_string(),
4347            })]),
4348            ChannelHook { tx },
4349        )
4350        .build();
4351        let shutdown = spawn_runtime(runtime.clone());
4352
4353        let handle = runtime
4354            .submit(RunSpec {
4355                input: Vec::new(),
4356                ..Default::default()
4357            })
4358            .await
4359            .unwrap();
4360        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4361            .await
4362            .unwrap()
4363            .unwrap();
4364        assert_eq!(outcome.run_id, handle.run_id);
4365        assert_eq!(outcome.status, TerminalStatus::Cancelled);
4366        assert_eq!(outcome.error.as_deref(), Some("obsolete"));
4367        assert_eq!(
4368            terminal_status_of(&runtime, &handle.run_id).await,
4369            Some(outcome.status)
4370        );
4371        assert_eq!(wait_for_kv(&queue, b"app/step-0").await, b"done");
4372        assert_eq!(
4373            queue.view().stats("workflow-steps").await.unwrap().dead,
4374            0,
4375            "a Cancel verdict must not dead-letter"
4376        );
4377
4378        let _ = shutdown.send(());
4379    }
4380
4381    #[tokio::test(start_paused = true)]
4382    async fn an_external_cancel_discards_staged_effects() {
4383        struct StageThenAwaitCancel {
4384            started: tokio::sync::mpsc::UnboundedSender<()>,
4385        }
4386
4387        impl StepRunner for StageThenAwaitCancel {
4388            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4389                step.effects
4390                    .put(b"app/override".to_vec(), b"staged".to_vec())
4391                    .map_err(|e| StepError::permanent(e.to_string()))?;
4392                let _ = self.started.send(());
4393                step.cancel_token.cancelled().await;
4394                Ok(StepOutcome::continue_now(Vec::new()))
4395            }
4396        }
4397
4398        let (queue, store) = open_queue().await;
4399        let (started_tx, mut started_rx) = tokio::sync::mpsc::unbounded_channel();
4400        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4401        let runtime = WorkflowRuntime::builder(
4402            queue.clone(),
4403            store,
4404            StageThenAwaitCancel {
4405                started: started_tx,
4406            },
4407            ChannelHook { tx },
4408        )
4409        .build();
4410        let shutdown = spawn_runtime(runtime.clone());
4411
4412        let handle = runtime
4413            .submit(RunSpec {
4414                input: Vec::new(),
4415                ..Default::default()
4416            })
4417            .await
4418            .unwrap();
4419        tokio::time::timeout(Duration::from_secs(2), started_rx.recv())
4420            .await
4421            .unwrap()
4422            .unwrap();
4423        assert!(runtime.cancel(&handle.run_id).await.unwrap());
4424
4425        let outcome = tokio::time::timeout(Duration::from_secs(2), rx.recv())
4426            .await
4427            .unwrap()
4428            .unwrap();
4429        assert_eq!(outcome.status, TerminalStatus::Cancelled);
4430        assert_eq!(outcome.error, None);
4431
4432        wait_for_drained(&queue).await;
4433        assert!(
4434            queue
4435                .view()
4436                .kv_get(b"app/override")
4437                .await
4438                .unwrap()
4439                .is_none()
4440        );
4441
4442        let _ = shutdown.send(());
4443    }
4444
4445    #[tokio::test(start_paused = true)]
4446    async fn a_run_dead_lettered_by_the_reaper_is_terminated_by_reconciliation() {
4447        let (queue, store, clock) = open_queue_at_with(
4448            1_700_000_000_000,
4449            fast_options().default_queue_config(
4450                QueueConfig::default()
4451                    .retry_backoff_base(Duration::ZERO)
4452                    .lease_duration(Duration::from_secs(1)),
4453            ),
4454        )
4455        .await;
4456        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4457        let runtime =
4458            WorkflowRuntime::builder(queue.clone(), store, PauseRunner, ChannelHook { tx })
4459                .poll_interval(Duration::from_millis(10))
4460                .build();
4461        let shutdown = spawn_runtime(runtime.clone());
4462
4463        let submitted = runtime
4464            .submit(RunSpec {
4465                run_id: Some(rid("hung")),
4466                input: Vec::new(),
4467                options: RunOptions {
4468                    max_attempts_per_step: Some(1),
4469                    ..Default::default()
4470                },
4471                ..Default::default()
4472            })
4473            .await
4474            .unwrap();
4475        for _ in 0..200 {
4476            if queue.view().stats("workflow-steps").await.unwrap().claimed == 1 {
4477                break;
4478            }
4479            tokio::time::sleep(Duration::from_millis(10)).await;
4480        }
4481        assert_eq!(
4482            runtime.status(&rid("hung")).await.unwrap().map(|s| s.state),
4483            Some(RunState::Running)
4484        );
4485
4486        // The lease expires past the attempt limit: the reaper dead-letters
4487        // the step outside the worker, with no worker in the loop.
4488        advance(&clock, Duration::from_secs(2)).await;
4489        let outcome = tokio::time::timeout(Duration::from_secs(5), rx.recv())
4490            .await
4491            .unwrap()
4492            .unwrap();
4493        assert_eq!(outcome.run_id, "hung");
4494        assert_eq!(outcome.status, TerminalStatus::Failed);
4495        assert_eq!(outcome.final_step, 0);
4496        assert_eq!(
4497            queue
4498                .view()
4499                .get_job(&submitted.job_id)
4500                .await
4501                .unwrap()
4502                .unwrap()
4503                .status,
4504            JobStatus::Dead
4505        );
4506        assert!(
4507            queue
4508                .view()
4509                .kv_get(&run_kv_key(&rid("hung")))
4510                .await
4511                .unwrap()
4512                .is_none()
4513        );
4514        assert!(
4515            queue
4516                .view()
4517                .kv_get(&step_kv_key(&rid("hung")))
4518                .await
4519                .unwrap()
4520                .is_none()
4521        );
4522        assert_eq!(
4523            terminal_status_of(&runtime, &rid("hung")).await,
4524            Some(TerminalStatus::Failed)
4525        );
4526        let _ = shutdown.send(());
4527
4528        // The run id is submitted again while the dead job is retained:
4529        // the pass identifies a dead-letter outside the worker by the
4530        // current step's job, so the new run is left alone.
4531        let again = runtime
4532            .submit(RunSpec {
4533                run_id: Some(rid("hung")),
4534                input: Vec::new(),
4535                ..Default::default()
4536            })
4537            .await
4538            .unwrap();
4539        assert!(again.newly_submitted);
4540        assert_eq!(runtime.inner.core.reconcile_dead_steps().await.unwrap(), 0);
4541        assert!(
4542            runtime.status(&rid("hung")).await.unwrap().is_some(),
4543            "the re-submitted run is active"
4544        );
4545    }
4546
4547    #[tokio::test(start_paused = true)]
4548    async fn a_wait_follows_the_run_across_its_steps() {
4549        let (queue, store) = open_queue().await;
4550        let runtime = WorkflowRuntime::builder(
4551            queue,
4552            store,
4553            ScriptedRunner::new(vec![
4554                StepOutcome::continue_now(b"next".to_vec()),
4555                StepOutcome::Succeed {
4556                    result: b"done".to_vec(),
4557                },
4558            ]),
4559            NoopTerminalHook,
4560        )
4561        .build();
4562        assert!(matches!(
4563            runtime.wait(&rid("absent")).await,
4564            Err(Error::RunNotFound(id)) if id == "absent"
4565        ));
4566        let shutdown = spawn_runtime(runtime.clone());
4567
4568        runtime
4569            .submit(RunSpec {
4570                run_id: Some(rid("two")),
4571                input: b"x".to_vec(),
4572                ..Default::default()
4573            })
4574            .await
4575            .unwrap();
4576        let end = tokio::time::timeout(Duration::from_secs(5), runtime.wait(&rid("two")))
4577            .await
4578            .expect("the wait resolved")
4579            .unwrap();
4580        let outcome = end.outcome.expect("the worker recorded the outcome");
4581        assert_eq!(
4582            (outcome.final_step, outcome.result.as_deref()),
4583            (1, Some(b"done".as_slice()))
4584        );
4585        assert_eq!(
4586            (end.termination.status, end.termination.final_step),
4587            (TerminalStatus::Succeeded, 1)
4588        );
4589
4590        // A run already terminated is reported at once.
4591        let again = runtime.wait(&rid("two")).await.unwrap();
4592        assert_eq!(again.termination, end.termination);
4593        let _ = shutdown.send(());
4594    }
4595
4596    #[tokio::test(start_paused = true)]
4597    async fn a_result_record_of_an_earlier_run_is_not_reported_for_a_re_submitted_run_id() {
4598        let (queue, store, clock) = open_queue_at(10_000).await;
4599        let runtime = WorkflowRuntime::builder(
4600            queue.clone(),
4601            store,
4602            FixedRunner::new(Ok(StepOutcome::Succeed {
4603                result: b"done".to_vec(),
4604            })),
4605            NoopTerminalHook,
4606        )
4607        .build();
4608        let spec = RunSpec {
4609            run_id: Some(rid("again")),
4610            input: b"x".to_vec(),
4611            ..Default::default()
4612        };
4613        runtime.submit(spec.clone()).await.unwrap();
4614        let claim = queue
4615            .claim("workflow-steps", Duration::from_secs(30))
4616            .await
4617            .unwrap()
4618            .unwrap();
4619        let effects = runtime
4620            .inner
4621            .process_step(&claim, &LeaseHandle::detached())
4622            .await
4623            .unwrap();
4624        queue.ack_with(&claim, effects).await.unwrap();
4625        let first = runtime.wait(&rid("again")).await.unwrap();
4626        assert_eq!(first.termination.status, TerminalStatus::Succeeded);
4627        assert!(first.outcome.is_some());
4628
4629        // The run id is submitted again: the earlier run's records
4630        // remain until the new run's termination overwrites them.
4631        clock.advance(Duration::from_secs(1));
4632        assert!(runtime.submit(spec).await.unwrap().newly_submitted);
4633        assert!(
4634            runtime.outcome(&rid("again")).await.unwrap().is_none(),
4635            "the run is active"
4636        );
4637        assert!(runtime.cancel(&rid("again")).await.unwrap());
4638        let end = runtime.wait(&rid("again")).await.unwrap();
4639        assert_eq!(
4640            (end.termination.status, end.termination.terminated_at_ms),
4641            (TerminalStatus::Cancelled, 11_000)
4642        );
4643        assert!(
4644            end.outcome.is_none(),
4645            "no worker terminated the new run, so the earlier run's record is not its outcome"
4646        );
4647        assert!(runtime.outcome(&rid("again")).await.unwrap().is_none());
4648    }
4649
4650    #[tokio::test(start_paused = true)]
4651    async fn a_cancel_request_does_not_reach_a_re_submission_of_the_run_id() {
4652        let (queue, store, _clock) = open_queue_at(10_000).await;
4653        let runtime =
4654            WorkflowRuntime::builder(queue.clone(), store, UnreachableRunner, NoopTerminalHook)
4655                .build();
4656        let spec = RunSpec {
4657            run_id: Some(rid("again")),
4658            input: b"x".to_vec(),
4659            ..Default::default()
4660        };
4661        runtime.submit(spec.clone()).await.unwrap();
4662        assert!(runtime.cancel(&rid("again")).await.unwrap());
4663        let end = runtime.wait(&rid("again")).await.unwrap();
4664        assert_eq!(end.termination.status, TerminalStatus::Cancelled);
4665
4666        assert!(runtime.submit(spec).await.unwrap().newly_submitted);
4667        let record = runtime
4668            .view()
4669            .run_record(&rid("again"))
4670            .await
4671            .unwrap()
4672            .unwrap();
4673        assert!(!record.cancel_requested);
4674        let status = runtime.status(&rid("again")).await.unwrap().unwrap();
4675        assert_eq!(status.state, RunState::Pending);
4676    }
4677
4678    #[tokio::test(start_paused = true)]
4679    async fn a_step_dead_lettered_outside_the_worker_is_waited_for_until_reconciliation() {
4680        let (queue, store, _clock) = open_queue_at(10_000).await;
4681        let runtime =
4682            WorkflowRuntime::builder(queue.clone(), store, UnreachableRunner, NoopTerminalHook)
4683                .poll_interval(Duration::from_millis(10))
4684                .build();
4685        runtime
4686            .submit(RunSpec {
4687                run_id: Some(rid("hung")),
4688                input: Vec::new(),
4689                ..Default::default()
4690            })
4691            .await
4692            .unwrap();
4693        let claim = queue
4694            .claim("workflow-steps", Duration::from_secs(60))
4695            .await
4696            .unwrap()
4697            .unwrap();
4698        queue.dead_letter(&claim, "hung").await.unwrap();
4699
4700        let waiting = tokio::spawn({
4701            let runtime = runtime.clone();
4702            async move { runtime.wait(&rid("hung")).await }
4703        });
4704        assert!(
4705            !runtime.cancel(&rid("hung")).await.unwrap(),
4706            "the request is not honoured"
4707        );
4708        assert!(
4709            runtime
4710                .wait_timeout(&rid("hung"), Duration::from_secs(1))
4711                .await
4712                .unwrap()
4713                .is_none(),
4714            "nothing terminates the run without a worker"
4715        );
4716        assert_eq!(runtime.inner.core.reconcile_dead_steps().await.unwrap(), 1);
4717        let end = tokio::time::timeout(Duration::from_secs(5), waiting)
4718            .await
4719            .expect("the wait resolved")
4720            .unwrap()
4721            .unwrap();
4722        assert_eq!(
4723            end.termination,
4724            RunTermination {
4725                status: TerminalStatus::Failed,
4726                error: Some("hung".into()),
4727                error_kind: None,
4728                final_step: 0,
4729                terminated_at_ms: 10_000,
4730            },
4731            "the termination is read from the terminal record",
4732        );
4733        assert!(end.outcome.is_none());
4734        assert_eq!(
4735            runtime.wait(&rid("hung")).await.unwrap().termination,
4736            end.termination,
4737            "the record is retained without memo retention",
4738        );
4739    }
4740
4741    #[tokio::test(start_paused = true)]
4742    async fn a_pointer_over_a_missing_job_is_an_inconsistent_run_state() {
4743        let (queue, store) = open_queue().await;
4744        let runtime =
4745            WorkflowRuntime::builder(queue.clone(), store, UnreachableRunner, NoopTerminalHook)
4746                .build();
4747        let submitted = runtime
4748            .submit(RunSpec {
4749                run_id: Some(rid("torn")),
4750                input: Vec::new(),
4751                ..Default::default()
4752            })
4753            .await
4754            .unwrap();
4755        // The step job is removed without the runtime, so the pointer
4756        // outlives it.
4757        queue.cancel(&submitted.job_id).await.unwrap();
4758        assert!(matches!(
4759            runtime.status(&rid("torn")).await,
4760            Err(Error::InconsistentRunState(id)) if id == "torn"
4761        ));
4762        assert!(matches!(
4763            runtime.wait(&rid("torn")).await,
4764            Err(Error::InconsistentRunState(id)) if id == "torn"
4765        ));
4766        assert!(matches!(
4767            runtime.cancel(&rid("torn")).await,
4768            Err(Error::InconsistentRunState(id)) if id == "torn"
4769        ));
4770    }
4771
4772    #[tokio::test(start_paused = true)]
4773    async fn a_member_record_is_rewritten_only_by_the_terminating_settlement() {
4774        struct RecordReadingRunner {
4775            pending_seen: Arc<StdMutex<Vec<bool>>>,
4776        }
4777
4778        impl StepRunner for RecordReadingRunner {
4779            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4780                let record = step
4781                    .kv
4782                    .get(&group_member_kv_key(&rid("g"), "m"))
4783                    .await?
4784                    .expect("the member record is written with the submission");
4785                let member: DurableMember = rmp_serde::from_slice(&record).unwrap();
4786                self.pending_seen
4787                    .lock()
4788                    .unwrap()
4789                    .push(member.terminated.is_none());
4790                Err(StepError::transient("still failing"))
4791            }
4792        }
4793
4794        let (queue, store) = open_queue_with(fast_options()).await;
4795        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4796        let pending_seen = Arc::new(StdMutex::new(Vec::new()));
4797        let runtime = WorkflowRuntime::builder(
4798            queue.clone(),
4799            store,
4800            RecordReadingRunner {
4801                pending_seen: pending_seen.clone(),
4802            },
4803            ChannelHook { tx },
4804        )
4805        .build();
4806        let shutdown = spawn_runtime(runtime.clone());
4807
4808        let group = runtime.group(rid("g"));
4809        group
4810            .submit(
4811                vec![GroupMember {
4812                    key: "m".to_string(),
4813                    input: Vec::new(),
4814                }],
4815                &RunOptions {
4816                    max_attempts_per_step: Some(2),
4817                    ..RunOptions::default()
4818                },
4819            )
4820            .await
4821            .unwrap();
4822        let outcome = tokio::time::timeout(Duration::from_secs(5), rx.recv())
4823            .await
4824            .unwrap()
4825            .unwrap();
4826        assert_eq!(outcome.status, TerminalStatus::Failed);
4827        for _ in 0..200 {
4828            if queue.view().stats("workflow-steps").await.unwrap().dead == 1 {
4829                break;
4830            }
4831            tokio::time::sleep(Duration::from_millis(10)).await;
4832        }
4833        assert_eq!(queue.view().stats("workflow-steps").await.unwrap().dead, 1);
4834
4835        // The retried first attempt left the record pending; the
4836        // exhausted second attempt's termination commits with the
4837        // dead-letter.
4838        assert_eq!(*pending_seen.lock().unwrap(), vec![true, true]);
4839        let members = group.members().await.unwrap();
4840        assert_eq!(members.len(), 1);
4841        assert_eq!(members[0].key, "m");
4842        assert_eq!(members[0].status(), Some(TerminalStatus::Failed));
4843        assert_eq!(
4844            members[0]
4845                .record
4846                .terminated
4847                .as_ref()
4848                .unwrap()
4849                .error
4850                .as_deref(),
4851            Some("still failing")
4852        );
4853
4854        let _ = shutdown.send(());
4855    }
4856
4857    #[tokio::test(start_paused = true)]
4858    async fn a_replayed_step_outcome_restores_its_staged_effects() {
4859        struct StagingContinueRunner {
4860            calls: Arc<AtomicU32>,
4861        }
4862
4863        impl StepRunner for StagingContinueRunner {
4864            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4865                self.calls.fetch_add(1, Ordering::SeqCst);
4866                step.effects
4867                    .put(b"app/replayed".to_vec(), b"v".to_vec())
4868                    .map_err(|e| StepError::transient(e.to_string()))?;
4869                step.effects
4870                    .delete(b"app/stale".to_vec())
4871                    .map_err(|e| StepError::transient(e.to_string()))?;
4872                Ok(StepOutcome::continue_now(b"step1".to_vec()))
4873            }
4874        }
4875
4876        let (queue, store) = open_queue().await;
4877        let calls = Arc::new(AtomicU32::new(0));
4878        let runtime = WorkflowRuntime::builder(
4879            queue.clone(),
4880            store,
4881            StagingContinueRunner {
4882                calls: calls.clone(),
4883            },
4884            NoopTerminalHook,
4885        )
4886        .step_output_replay()
4887        .build();
4888
4889        queue.kv_put(b"app/stale", b"old").await.unwrap();
4890        runtime
4891            .submit(RunSpec {
4892                run_id: Some(rid("replay-effects")),
4893                input: b"input".to_vec(),
4894                ..Default::default()
4895            })
4896            .await
4897            .unwrap();
4898
4899        let job = queue
4900            .claim("workflow-steps", Duration::from_secs(30))
4901            .await
4902            .unwrap()
4903            .unwrap();
4904
4905        // First delivery: the returned effects are dropped, simulating a
4906        // crash between the replay-record write and the settlement.
4907        let _ = runtime
4908            .inner
4909            .process_step(&job, &LeaseHandle::detached())
4910            .await
4911            .unwrap();
4912        assert_eq!(calls.load(Ordering::SeqCst), 1);
4913        assert!(
4914            queue
4915                .view()
4916                .kv_get(b"app/replayed")
4917                .await
4918                .unwrap()
4919                .is_none()
4920        );
4921
4922        // Redelivery replays the stored outcome and restores the staged
4923        // effects into the settlement without invoking the runner.
4924        let effects = runtime
4925            .inner
4926            .process_step(&job, &LeaseHandle::detached())
4927            .await
4928            .unwrap();
4929        assert_eq!(calls.load(Ordering::SeqCst), 1);
4930        assert_eq!(
4931            effects.kv_writes.get(b"app/replayed".as_slice()),
4932            Some(&b"v".to_vec())
4933        );
4934        assert!(effects.kv_deletes.contains(&b"app/stale".to_vec()));
4935        queue.ack_with(&job, effects).await.unwrap();
4936        assert_eq!(
4937            queue
4938                .view()
4939                .kv_get(b"app/replayed")
4940                .await
4941                .unwrap()
4942                .as_deref(),
4943            Some(b"v".as_slice())
4944        );
4945        assert!(queue.view().kv_get(b"app/stale").await.unwrap().is_none());
4946    }
4947
4948    #[tokio::test(start_paused = true)]
4949    async fn only_the_committed_outcome_produces_a_notification() {
4950        struct GatedSecondAttempt {
4951            calls: Arc<AtomicU32>,
4952            running: tokio::sync::mpsc::UnboundedSender<()>,
4953        }
4954
4955        impl StepRunner for GatedSecondAttempt {
4956            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
4957                if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
4958                    return Ok(StepOutcome::Succeed {
4959                        result: b"done".to_vec(),
4960                    });
4961                }
4962                let _ = self.running.send(());
4963                step.cancel_token.cancelled().await;
4964                Ok(StepOutcome::Succeed {
4965                    result: b"done".to_vec(),
4966                })
4967            }
4968        }
4969
4970        let (queue, store) = open_queue().await;
4971        let (running_tx, mut running_rx) = tokio::sync::mpsc::unbounded_channel();
4972        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
4973        let runtime = WorkflowRuntime::builder(
4974            queue.clone(),
4975            store,
4976            GatedSecondAttempt {
4977                calls: Arc::new(AtomicU32::new(0)),
4978                running: running_tx,
4979            },
4980            ChannelHook { tx },
4981        )
4982        .build();
4983
4984        runtime
4985            .submit(RunSpec {
4986                run_id: Some(rid("phantom")),
4987                input: Vec::new(),
4988                ..Default::default()
4989            })
4990            .await
4991            .unwrap();
4992        let job = queue
4993            .claim("workflow-steps", Duration::from_secs(30))
4994            .await
4995            .unwrap()
4996            .unwrap();
4997
4998        // First settlement attempt: the runner succeeds, but the effects
4999        // are dropped, as when the settlement loses the claim. The
5000        // Succeeded notification is dropped with them.
5001        let _ = runtime
5002            .inner
5003            .process_step(&job, &queue.lease_handle(&job))
5004            .await
5005            .unwrap();
5006
5007        // The redelivered attempt observes an external cancel and
5008        // commits Cancelled.
5009        let worker = {
5010            let inner = runtime.inner.clone();
5011            let queue = queue.clone();
5012            tokio::spawn(async move {
5013                let effects = inner
5014                    .process_step(&job, &queue.lease_handle(&job))
5015                    .await
5016                    .unwrap();
5017                queue.ack_with(&job, effects).await.unwrap();
5018            })
5019        };
5020        running_rx.recv().await.unwrap();
5021        assert!(runtime.cancel(&rid("phantom")).await.unwrap());
5022        worker.await.unwrap();
5023
5024        let notification = queue
5025            .claim("workflow-steps", Duration::from_secs(30))
5026            .await
5027            .unwrap()
5028            .expect("the committed settlement enqueued its notification");
5029        let effects = runtime
5030            .inner
5031            .process_step(&notification, &LeaseHandle::detached())
5032            .await
5033            .unwrap();
5034        queue.ack_with(&notification, effects).await.unwrap();
5035        let outcome = rx.recv().await.unwrap();
5036        assert_eq!(outcome.status, TerminalStatus::Cancelled);
5037        assert!(
5038            queue
5039                .claim("workflow-steps", Duration::from_secs(30))
5040                .await
5041                .unwrap()
5042                .is_none(),
5043            "the outcome that never committed must produce no notification",
5044        );
5045        assert!(rx.try_recv().is_err());
5046    }
5047
5048    #[tokio::test(start_paused = true)]
5049    async fn hook_effects_commit_with_the_notification_ack() {
5050        struct EffectHook;
5051
5052        impl TerminalHook for EffectHook {
5053            async fn on_termination(
5054                &self,
5055                outcome: &RunOutcome,
5056                effects: &TerminalEffects,
5057            ) -> std::result::Result<(), StepError> {
5058                effects
5059                    .put(
5060                        format!("app/outcomes/{}", outcome.run_id),
5061                        outcome.status.as_str(),
5062                    )
5063                    .map_err(|e| StepError::permanent(e.to_string()))?;
5064                effects
5065                    .enqueue(EnqueueRequest {
5066                        queue: "side-effects".to_string(),
5067                        payload: outcome.run_id.to_string().into_bytes(),
5068                        options: EnqueueOptions::default(),
5069                    })
5070                    .map_err(|e| StepError::permanent(e.to_string()))?;
5071                Ok(())
5072            }
5073        }
5074
5075        let (queue, store) = open_queue().await;
5076        let runtime = WorkflowRuntime::builder(
5077            queue.clone(),
5078            store,
5079            ScriptedRunner::new(vec![StepOutcome::Succeed { result: Vec::new() }]),
5080            EffectHook,
5081        )
5082        .build();
5083        let shutdown = spawn_runtime(runtime.clone());
5084
5085        runtime
5086            .submit(RunSpec {
5087                run_id: Some(rid("hooked")),
5088                input: Vec::new(),
5089                ..Default::default()
5090            })
5091            .await
5092            .unwrap();
5093
5094        assert_eq!(
5095            wait_for_kv(&queue, b"app/outcomes/hooked").await,
5096            b"succeeded"
5097        );
5098        let side = queue
5099            .claim("side-effects", Duration::from_secs(30))
5100            .await
5101            .unwrap()
5102            .expect("the staged enqueue committed with the notification ack");
5103        assert_eq!(side.payload.as_slice(), b"hooked");
5104
5105        let _ = shutdown.send(());
5106    }
5107
5108    #[tokio::test(start_paused = true)]
5109    async fn a_transiently_failing_hook_retries_the_notification() {
5110        struct FlakyHook {
5111            calls: Arc<AtomicU32>,
5112        }
5113
5114        impl TerminalHook for FlakyHook {
5115            async fn on_termination(
5116                &self,
5117                outcome: &RunOutcome,
5118                effects: &TerminalEffects,
5119            ) -> std::result::Result<(), StepError> {
5120                if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
5121                    return Err(StepError::transient("first attempt fails"));
5122                }
5123                effects
5124                    .put(format!("app/notified/{}", outcome.run_id), b"1".to_vec())
5125                    .map_err(|e| StepError::permanent(e.to_string()))?;
5126                Ok(())
5127            }
5128        }
5129
5130        let (queue, store) = open_queue_with(fast_options()).await;
5131        let calls = Arc::new(AtomicU32::new(0));
5132        let runtime = WorkflowRuntime::builder(
5133            queue.clone(),
5134            store,
5135            ScriptedRunner::new(vec![StepOutcome::Succeed { result: Vec::new() }]),
5136            FlakyHook {
5137                calls: calls.clone(),
5138            },
5139        )
5140        .build();
5141        let shutdown = spawn_runtime(runtime.clone());
5142
5143        runtime
5144            .submit(RunSpec {
5145                run_id: Some(rid("flaky")),
5146                input: Vec::new(),
5147                ..Default::default()
5148            })
5149            .await
5150            .unwrap();
5151
5152        assert_eq!(wait_for_kv(&queue, b"app/notified/flaky").await, b"1");
5153        assert_eq!(calls.load(Ordering::SeqCst), 2);
5154
5155        let _ = shutdown.send(());
5156    }
5157
5158    #[tokio::test(start_paused = true)]
5159    async fn a_noop_hook_enqueues_no_notification() {
5160        let (queue, store) = open_queue().await;
5161        let runtime = WorkflowRuntime::builder(
5162            queue.clone(),
5163            store,
5164            ScriptedRunner::new(vec![StepOutcome::Succeed { result: Vec::new() }]),
5165            NoopTerminalHook,
5166        )
5167        .build();
5168        runtime
5169            .submit(RunSpec {
5170                input: Vec::new(),
5171                ..Default::default()
5172            })
5173            .await
5174            .unwrap();
5175        let job = queue
5176            .claim("workflow-steps", Duration::from_secs(30))
5177            .await
5178            .unwrap()
5179            .unwrap();
5180        let effects = runtime
5181            .inner
5182            .process_step(&job, &LeaseHandle::detached())
5183            .await
5184            .unwrap();
5185        assert!(effects.enqueues.is_empty());
5186        queue.ack_with(&job, effects).await.unwrap();
5187        assert!(
5188            queue
5189                .claim("workflow-steps", Duration::from_secs(30))
5190                .await
5191                .unwrap()
5192                .is_none()
5193        );
5194    }
5195
5196    #[cfg(feature = "webhooks")]
5197    #[tokio::test(start_paused = true)]
5198    async fn the_webhook_hook_stages_its_delivery_as_a_notification_effect() {
5199        use crate::terminal::WebhookTerminalHook;
5200
5201        let (queue, store) = open_queue().await;
5202        let runtime = WorkflowRuntime::builder(
5203            queue.clone(),
5204            store,
5205            ScriptedRunner::new(vec![
5206                StepOutcome::Succeed {
5207                    result: b"payload".to_vec(),
5208                },
5209                StepOutcome::Succeed { result: Vec::new() },
5210            ]),
5211            WebhookTerminalHook::new("callbacks"),
5212        )
5213        .build();
5214        let shutdown = spawn_runtime(runtime.clone());
5215
5216        runtime
5217            .submit(RunSpec {
5218                run_id: Some(rid("with-callback")),
5219                input: Vec::new(),
5220                options: RunOptions {
5221                    headers: HashMap::from([(
5222                        "callback_url".to_string(),
5223                        "https://example.com/done".to_string(),
5224                    )]),
5225                    ..Default::default()
5226                },
5227                ..Default::default()
5228            })
5229            .await
5230            .unwrap();
5231
5232        let webhook = loop {
5233            if let Some(job) = queue
5234                .claim("callbacks", Duration::from_secs(30))
5235                .await
5236                .unwrap()
5237            {
5238                break job;
5239            }
5240            tokio::time::sleep(Duration::from_millis(10)).await;
5241        };
5242        assert_eq!(webhook.payload.as_slice(), b"payload");
5243        assert_eq!(
5244            webhook.headers.get("webhook.url").unwrap(),
5245            "https://example.com/done"
5246        );
5247        assert_eq!(
5248            webhook.headers.get("http.Workflow-Run-Status").unwrap(),
5249            "succeeded"
5250        );
5251
5252        // A run without a callback header enqueues no notification.
5253        runtime
5254            .submit(RunSpec {
5255                run_id: Some(rid("without-callback")),
5256                input: Vec::new(),
5257                ..Default::default()
5258            })
5259            .await
5260            .unwrap();
5261        wait_for_drained(&queue).await;
5262        assert!(
5263            queue
5264                .claim("callbacks", Duration::from_secs(30))
5265                .await
5266                .unwrap()
5267                .is_none()
5268        );
5269
5270        let _ = shutdown.send(());
5271    }
5272
5273    #[tokio::test(start_paused = true)]
5274    async fn a_late_write_through_an_escaped_handle_is_refused() {
5275        struct EscapingRunner {
5276            escaped: Arc<StdMutex<Option<EffectsHandle>>>,
5277        }
5278
5279        impl StepRunner for EscapingRunner {
5280            async fn run_step(&self, step: &Step) -> std::result::Result<StepOutcome, StepError> {
5281                *self.escaped.lock().unwrap() = Some(step.effects.clone());
5282                Ok(StepOutcome::Succeed { result: Vec::new() })
5283            }
5284        }
5285
5286        let (queue, store) = open_queue().await;
5287        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
5288        let escaped = Arc::new(StdMutex::new(None));
5289        let runtime = WorkflowRuntime::builder(
5290            queue,
5291            store,
5292            EscapingRunner {
5293                escaped: escaped.clone(),
5294            },
5295            ChannelHook { tx },
5296        )
5297        .build();
5298        let shutdown = spawn_runtime(runtime.clone());
5299
5300        runtime
5301            .submit(RunSpec {
5302                input: Vec::new(),
5303                ..Default::default()
5304            })
5305            .await
5306            .unwrap();
5307        tokio::time::timeout(Duration::from_secs(2), rx.recv())
5308            .await
5309            .unwrap()
5310            .unwrap();
5311
5312        let handle = escaped.lock().unwrap().take().unwrap();
5313        assert!(matches!(
5314            handle.put("app/late", "v"),
5315            Err(Error::EffectsSealed)
5316        ));
5317
5318        let _ = shutdown.send(());
5319    }
5320}