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}