Skip to main content

taquba_workflow/
runner.rs

1use std::collections::HashMap;
2use std::future::Future;
3use std::time::Duration;
4
5use taquba::LeaseHandle;
6use tokio_util::sync::CancellationToken;
7
8use crate::effects::EffectsHandle;
9use crate::kv::KvReadHandle;
10use crate::memo::Memo;
11
12/// A single step within a workflow run, handed to [`StepRunner::run_step`].
13///
14/// Mirrors [`taquba::JobRecord`]: the `payload` is opaque application bytes
15/// and `headers` carries user metadata you set at submission (reserved
16/// `workflow.*` keys are filtered out before the runner sees them).
17#[derive(Debug, Clone)]
18pub struct Step {
19    /// Caller-visible run identifier (the value passed to or generated by
20    /// [`crate::RunSpec`]).
21    pub run_id: String,
22    /// Zero-based step number. Step 0 is always the first step of a run, with
23    /// the original submission input as its `payload`.
24    pub step_number: u32,
25    /// Application-defined bytes. For step 0 this is the submission `input`;
26    /// for later steps it is the bytes returned by the previous step's
27    /// [`StepOutcome::Continue`].
28    pub payload: Vec<u8>,
29    /// Submitter-supplied metadata, threaded through every step of the run.
30    /// Reserved `workflow.*` headers are stripped before the runner sees them.
31    pub headers: HashMap<String, String>,
32    /// The Taquba job ID for this step, useful for tracing.
33    pub job_id: String,
34    /// How many times Taquba has attempted to deliver this step. `1` on the
35    /// first attempt; `>1` after a lease expiry / nack retry.
36    pub attempts: u32,
37    /// The attempt limit for this step: a transient failure on the attempt
38    /// numbered `max_attempts` ends the run.
39    pub max_attempts: u32,
40    /// Cooperative cancellation signal for the run. The runtime cancels
41    /// this token when [`crate::WorkflowRuntime::cancel`] is called while
42    /// this step is in flight, so a long-running runner (e.g. an LLM call,
43    /// a slow HTTP request) can short-circuit instead of running to
44    /// completion. Typical use:
45    ///
46    /// ```ignore
47    /// tokio::select! {
48    ///     out = do_slow_work(step) => out,
49    ///     _ = step.cancel_token.cancelled() => {
50    ///         Ok(StepOutcome::Cancel { reason: "cooperative".into() })
51    ///     }
52    /// }
53    /// ```
54    ///
55    /// Runners that ignore the token remain correct: the runtime still
56    /// discards the outcome of a cancelled step and fires the terminal
57    /// hook with [`crate::TerminalStatus::Cancelled`]. Watching the token
58    /// only reduces cancellation latency for slow steps; it doesn't
59    /// change semantics.
60    ///
61    /// The token is a child of the delivery's claim token. A re-delivery
62    /// of this step observes `is_cancelled() == true` immediately,
63    /// because the queue re-fires the claim's token from the job's
64    /// persisted cancellation. Cancelling this token leaves the claim's
65    /// token uncancelled, so the runtime does not treat the step as
66    /// externally cancelled.
67    pub cancel_token: CancellationToken,
68    /// The lease handle for this step's delivery. A long-running runner
69    /// calls [`LeaseHandle::ensure_at_least`] at progress points (or
70    /// once, with a slow call's timeout, before issuing it) so the step
71    /// is not re-queued while it still runs. Use
72    /// [`LeaseHandle::detached`] when constructing a `Step` in tests; a
73    /// detached handle's calls succeed without effect.
74    pub lease: LeaseHandle,
75    /// Per-step durable key-value store, scoped to this step's
76    /// `(run_id, step_number)`. Use to memoize expensive within-step
77    /// side effects (LLM calls, paid APIs) so an at-least-once retry
78    /// of this step doesn't re-pay for work the prior attempt already
79    /// did:
80    ///
81    /// ```ignore
82    /// let response = match step.memo.get("llm").await? {
83    ///     Some(cached) => deserialize(&cached),
84    ///     None => {
85    ///         let fresh = llm.complete(&prompt).await?;
86    ///         step.memo.put("llm", &serialize(&fresh)).await?;
87    ///         fresh
88    ///     }
89    /// };
90    /// ```
91    ///
92    /// See [`Memo`] for the full API.
93    pub memo: Memo,
94    /// Run-scoped durable key-value store, shared by every step of the
95    /// run. Entries live beside the per-step [`Step::memo`] entries and
96    /// are removed with them when the run's retention expires. Use it
97    /// for values a later step reads back, such as an accumulating
98    /// journal; the durable channel for the next step's input is
99    /// [`StepOutcome::Continue`]'s payload.
100    pub run_memo: Memo,
101    /// Application KV effects for this step. Writes and deletes staged
102    /// here are applied in the same transaction as the settlement that
103    /// commits the returned outcome, so application state cannot
104    /// diverge from the run's transition on a crash:
105    ///
106    /// ```ignore
107    /// step.effects.put(format!("app/runs/{}", step.run_id), b"done".to_vec())?;
108    /// Ok(StepOutcome::Succeed { result })
109    /// ```
110    ///
111    /// See [`EffectsHandle`] for the staging rules. Use
112    /// [`EffectsHandle::detached`] when constructing a `Step` in tests.
113    pub effects: EffectsHandle,
114    /// Read access to the caller KV namespace. A committed value (an
115    /// earlier step's applied effect, a [`crate::RunSpec::kv_writes`]
116    /// entry, a direct [`taquba::Queue::kv_put`]) is readable here;
117    /// effects staged by this step become readable only after it
118    /// settles. The intended use is a read-then-stage marker check:
119    ///
120    /// ```ignore
121    /// if step.kv.get(b"app/indexed/doc-1").await?.is_none() {
122    ///     index_document(&step.payload).await?; // idempotent
123    ///     step.effects.put("app/indexed/doc-1", b"1".to_vec())?;
124    /// }
125    /// ```
126    ///
127    /// See [`KvReadHandle`] for the read semantics. Use
128    /// [`KvReadHandle::detached`] when constructing a `Step` in tests.
129    pub kv: KvReadHandle,
130    /// The signal payload, when this step was reached through a
131    /// [`Trigger::OnSignal`] wait that a signal resolved: the previous
132    /// step continued with `OnSignal`, and a
133    /// [`crate::WorkflowRuntime::signal`] call for the correlation key
134    /// arrived before the timeout. `None` when the timeout elapsed first
135    /// and on every step not preceded by an `OnSignal` wait.
136    pub signal: Option<Vec<u8>>,
137}
138
139/// When the next step of a run becomes claimable. Set on the `when`
140/// field of [`StepOutcome::Continue`].
141#[derive(Debug, Clone)]
142#[non_exhaustive]
143pub enum Trigger {
144    /// The next step is claimable immediately.
145    Immediate,
146    /// The next step is claimable `Duration` from now.
147    After(Duration),
148    /// The next step is claimable when a signal for `correlation_key`
149    /// arrives via [`crate::WorkflowRuntime::signal`], or after `timeout`,
150    /// whichever comes first. The next step reads [`Step::signal`] to
151    /// distinguish the two: `Some(payload)` when a signal arrived, `None`
152    /// when the timeout elapsed first. A signal that arrived before this
153    /// step settled is consumed at settlement and the next step runs
154    /// immediately.
155    ///
156    /// One waiter per correlation key: registering a second waiter while
157    /// one is already waiting fails the step permanently. Choose keys
158    /// that are unique per waiter (e.g. include the run id).
159    OnSignal {
160        /// Caller-chosen key the signaller addresses.
161        correlation_key: String,
162        /// Upper bound on the wait; the next step runs with
163        /// [`Step::signal`] `None` when it elapses first.
164        timeout: Duration,
165    },
166}
167
168/// What the runner wants the runtime to do after this step.
169#[derive(Debug, Clone)]
170#[non_exhaustive]
171pub enum StepOutcome {
172    /// Run is not finished. Enqueue the next step with `payload` as its
173    /// bytes; `when` decides when it becomes claimable. The runtime
174    /// advances `step_number` by 1. The constructors
175    /// [`Self::continue_now`] and [`Self::continue_after`] build the
176    /// common forms.
177    Continue {
178        /// Bytes to hand to the next step's [`Step::payload`].
179        payload: Vec<u8>,
180        /// When the next step becomes claimable.
181        when: Trigger,
182    },
183    /// The run is finished successfully. The runtime acks the step and fires
184    /// the configured terminal hook with
185    /// [`crate::TerminalStatus::Succeeded`] and `result` as the body.
186    Succeed {
187        /// Final result bytes handed to the terminal hook.
188        result: Vec<u8>,
189    },
190    /// The run is finished as failed by the runner's verdict; the runner
191    /// ran to completion but the workflow's logical outcome is "no" (e.g.
192    /// a validation rule rejected the input, a policy check denied the
193    /// request, an agent decided the task can't be fulfilled). The runtime
194    /// acks the step and fires the terminal hook with
195    /// [`crate::TerminalStatus::Failed`] and `reason` as the error.
196    ///
197    /// Use this for *workflow-level* failures. For *infrastructure*
198    /// failures (network outage, downstream service down, etc.) return
199    /// `Err(StepError::transient)` or `Err(StepError::permanent)` instead;
200    /// those dead-letter the step so an operator can find it via
201    /// [`taquba::Queue::dead_jobs`]. `Fail` is a successful execution with
202    /// a negative outcome and does not dead-letter.
203    Fail {
204        /// Human-readable reason recorded on [`crate::RunOutcome::error`].
205        reason: String,
206    },
207    /// The run is finished as cancelled by the runner. Use this when the
208    /// runner decides on its own that the workflow should stop early
209    /// without it being a logical failure (e.g. a downstream cancellation
210    /// signal arrived mid-step, the user-supplied input is now obsolete).
211    /// The runtime acks the step and fires the terminal hook with
212    /// [`crate::TerminalStatus::Cancelled`] and `reason` as the error.
213    ///
214    /// For *external* cancellation requested by another component in the
215    /// process, call [`crate::WorkflowRuntime::cancel`] instead; the
216    /// runtime translates that into the same `Cancelled` terminal state.
217    Cancel {
218        /// Human-readable reason recorded on [`crate::RunOutcome::error`].
219        reason: String,
220    },
221}
222
223impl StepOutcome {
224    /// Continue the run; the next step is claimable immediately.
225    pub fn continue_now(payload: Vec<u8>) -> Self {
226        Self::Continue {
227            payload,
228            when: Trigger::Immediate,
229        }
230    }
231
232    /// Continue the run; the next step is claimable `delay` from now.
233    pub fn continue_after(payload: Vec<u8>, delay: Duration) -> Self {
234        Self::Continue {
235            payload,
236            when: Trigger::After(delay),
237        }
238    }
239
240    /// Continue the run; the next step is claimable when a signal for
241    /// `correlation_key` arrives, or after `timeout` at the latest.
242    pub fn continue_on_signal(
243        payload: Vec<u8>,
244        correlation_key: impl Into<String>,
245        timeout: Duration,
246    ) -> Self {
247        Self::Continue {
248            payload,
249            when: Trigger::OnSignal {
250                correlation_key: correlation_key.into(),
251                timeout,
252            },
253        }
254    }
255}
256
257/// Failure outcomes the runner can return.
258#[derive(Debug)]
259pub struct StepError {
260    /// Human-readable message recorded on the underlying job's `last_error`.
261    pub message: String,
262    /// Whether to retry the step or fail the run immediately.
263    pub kind: StepErrorKind,
264}
265
266impl StepError {
267    /// Build a transient error: Taquba retries the step per the queue's
268    /// backoff/`max_attempts`. Once `max_attempts` is exhausted, the step is
269    /// dead-lettered and the run terminates as failed.
270    pub fn transient(message: impl Into<String>) -> Self {
271        Self {
272            message: message.into(),
273            kind: StepErrorKind::Transient,
274        }
275    }
276
277    /// Build a permanent error: the step is dead-lettered immediately and the
278    /// run terminates as failed.
279    pub fn permanent(message: impl Into<String>) -> Self {
280        Self {
281            message: message.into(),
282            kind: StepErrorKind::Permanent,
283        }
284    }
285}
286
287impl std::fmt::Display for StepError {
288    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
289        f.write_str(&self.message)
290    }
291}
292
293impl std::error::Error for StepError {}
294
295impl From<crate::Error> for StepError {
296    fn from(err: crate::Error) -> Self {
297        let permanent = err.is_permanent();
298        let message = err.to_string();
299        if permanent {
300            Self::permanent(message)
301        } else {
302            Self::transient(message)
303        }
304    }
305}
306
307/// Whether a [`StepError`] should retry or fail the run.
308#[derive(Debug, Clone, Copy, PartialEq, Eq)]
309#[non_exhaustive]
310pub enum StepErrorKind {
311    /// Retry per the queue's backoff policy until `max_attempts` is reached.
312    Transient,
313    /// Dead-letter the step immediately; terminate the run as failed.
314    Permanent,
315}
316
317/// User-implemented logic that advances a single workflow step.
318///
319/// Implementations must be idempotent for the same `(run_id, step_number)`:
320/// Taquba is at-least-once, so a step can be claimed and processed more than
321/// once if a lease expires before the worker acks. Returning the same
322/// `StepOutcome` for the same input is the easiest way to satisfy this.
323pub trait StepRunner: Send + Sync {
324    /// Process a single step of a workflow run. Return [`StepOutcome::Continue`]
325    /// to enqueue the next step, [`StepOutcome::Succeed`] to finish the run
326    /// successfully, [`StepOutcome::Fail`] to terminate the run as Failed by
327    /// runner verdict, [`StepOutcome::Cancel`] to terminate the run as
328    /// Cancelled by runner verdict, or `Err(StepError)` to retry /
329    /// dead-letter on infrastructure errors.
330    fn run_step(
331        &self,
332        step: &Step,
333    ) -> impl Future<Output = std::result::Result<StepOutcome, StepError>> + Send;
334}
335
336#[cfg(test)]
337mod tests {
338    use super::*;
339
340    #[test]
341    fn from_workflow_error_maps_via_is_permanent() {
342        let permanent: StepError = crate::Error::InputMismatch("run-1".into()).into();
343        assert_eq!(permanent.kind, StepErrorKind::Permanent);
344
345        let store_err = taquba::object_store::Error::NotFound {
346            path: "x".into(),
347            source: "missing".into(),
348        };
349        let transient: StepError = crate::Error::Store(store_err).into();
350        assert_eq!(transient.kind, StepErrorKind::Transient);
351    }
352}