taquba-workflow 0.13.0

Durable, at-least-once workflow runtime on top of the Taquba task queue. Particularly well-suited for AI agent runs.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
use std::collections::HashMap;
use std::future::Future;
use std::ops::{Deref, DerefMut};
use std::sync::Arc;
use std::time::Duration;

use taquba::object_store::memory::InMemory;
use taquba::{LeaseHandle, PermanentFailure, WorkerError};
use tokio_util::sync::CancellationToken;

use crate::effects::EffectsHandle;
use crate::keys::RunId;
use crate::kv::KvReadHandle;
use crate::memo::{Memo, MemoStore};

/// The delivery a handler runs under: the identity of the run and of
/// the queue job delivering it, the attempt count and the delivery's
/// handles. [`Step`] and [`jobs::JobContext`](crate::jobs::JobContext)
/// dereference to it. It holds handles only; no queue is reachable
/// through it.
///
/// Constructed by the runtime. A test constructs one with
/// [`Delivery::detached`] and assigns the fields it needs.
#[derive(Debug, Clone)]
pub struct Delivery {
    /// Caller-visible run identifier (the value passed to or generated by
    /// [`crate::RunSpec`]).
    pub run_id: RunId,
    /// Submitter-supplied metadata, threaded through every step of the run.
    /// Reserved `workflow.*` headers are stripped before the handler sees them.
    pub headers: HashMap<String, String>,
    /// The Taquba job ID of this delivery, useful for tracing.
    pub job_id: String,
    /// How many times Taquba has attempted this delivery. `1` on the
    /// first attempt; `>1` after a lease expiry / nack retry.
    pub attempts: u32,
    /// The attempt limit: a transient failure on the attempt numbered
    /// `max_attempts` ends the run.
    pub max_attempts: u32,
    /// Cooperative cancellation signal for the run. The runtime cancels
    /// this token when [`crate::WorkflowRuntime::cancel`] is called while
    /// this delivery is in flight, so a long-running handler (e.g. an LLM
    /// call, a slow HTTP request) can short-circuit instead of running to
    /// completion. Typical use:
    ///
    /// ```ignore
    /// tokio::select! {
    ///     out = do_slow_work(step) => out,
    ///     _ = step.cancel_token.cancelled() => {
    ///         Ok(StepOutcome::Cancel { reason: "cooperative".into() })
    ///     }
    /// }
    /// ```
    ///
    /// Handlers that ignore the token remain correct: the runtime still
    /// discards the outcome of a cancelled step and fires the terminal
    /// hook with [`crate::TerminalStatus::Cancelled`]. Watching the token
    /// only reduces cancellation latency for slow steps; it doesn't
    /// change semantics.
    ///
    /// The token is a child of the claim's cancellation token. A
    /// re-delivery of this step observes `is_cancelled() == true`
    /// immediately, because the queue re-fires the claim's cancellation
    /// token from the job's persisted cancellation. Cancelling this token
    /// leaves the claim's token uncancelled, so the runtime does not treat
    /// the step as externally cancelled.
    pub cancel_token: CancellationToken,
    /// The lease handle of this delivery. A long-running handler calls
    /// [`LeaseHandle::ensure_at_least`] at progress points (or once,
    /// with a slow call's timeout, before issuing it) so the step is
    /// not re-queued while it still runs. A detached handle's calls
    /// succeed without effect.
    pub lease: LeaseHandle,
    /// Per-step durable key-value store, scoped to this step's
    /// `(run_id, step_number)`. Use to memoize expensive within-step
    /// side effects (LLM calls, paid APIs) so an at-least-once retry
    /// of this step doesn't re-pay for work the prior attempt already
    /// did:
    ///
    /// ```ignore
    /// let response = match step.memo.get("llm").await? {
    ///     Some(cached) => deserialize(&cached),
    ///     None => {
    ///         let fresh = llm.complete(&prompt).await?;
    ///         step.memo.put("llm", &serialize(&fresh)).await?;
    ///         fresh
    ///     }
    /// };
    /// ```
    ///
    /// See [`Memo`] for the full API.
    pub memo: Memo,
    /// Run-scoped durable key-value store, shared by every step of the
    /// run. Entries live beside the per-step [`Delivery::memo`] entries
    /// and are removed with them when the run's retention expires. Use
    /// it for values a later step reads back, such as an accumulating
    /// journal; the durable channel for the next step's input is
    /// [`StepOutcome::Continue`]'s payload.
    pub run_memo: Memo,
    /// Application KV effects for this step. Writes and deletes staged
    /// here are applied in the same transaction as the settlement that
    /// commits the returned outcome, so application state cannot
    /// diverge from the run's transition on a crash:
    ///
    /// ```ignore
    /// step.effects.put(format!("app/runs/{}", step.run_id), b"done".to_vec())?;
    /// Ok(StepOutcome::Succeed { result })
    /// ```
    ///
    /// See [`EffectsHandle`] for the staging rules.
    pub effects: EffectsHandle,
    /// Read access to the caller KV namespace. A committed value (an earlier
    /// step's applied effect, a [`crate::RunSpec::effects`] write, a direct
    /// [`taquba::Queue::kv_put`]) is readable here. Effects staged by this step
    /// become readable only after it settles. The intended use is a
    /// read-then-stage marker check:
    ///
    /// ```ignore
    /// if step.kv.get(b"app/indexed/doc-1").await?.is_none() {
    ///     index_document(&step.payload).await?; // idempotent
    ///     step.effects.put("app/indexed/doc-1", b"1".to_vec())?;
    /// }
    /// ```
    ///
    /// See [`KvReadHandle`] for the read semantics.
    pub kv: KvReadHandle,
}

impl Delivery {
    /// Whether this attempt is the last: a transient [`StepError`]
    /// returned from it dead-letters the step and ends the run.
    pub fn is_last_attempt(&self) -> bool {
        self.attempts >= self.max_attempts
    }

    /// A delivery bound to no queue, for tests: run `detached`, attempt
    /// 1 of 3, no headers, a new cancellation token, detached lease,
    /// effects and KV handles, and memos over an in-memory object store.
    pub fn detached() -> Self {
        let run_id = RunId::new("detached").expect("a literal run id");
        let memo_store = MemoStore::new(Arc::new(InMemory::new()), "memo");
        Self {
            run_id: run_id.clone(),
            headers: HashMap::new(),
            job_id: "detached".to_string(),
            attempts: 1,
            max_attempts: 3,
            cancel_token: CancellationToken::new(),
            lease: LeaseHandle::detached(),
            memo: memo_store.new_memo(&run_id, 0),
            run_memo: memo_store.new_run_memo(&run_id),
            effects: EffectsHandle::detached(),
            kv: KvReadHandle::detached(),
        }
    }
}

/// A single step within a workflow run, handed to [`StepRunner::run_step`]:
/// the [`Delivery`] it runs under, which it dereferences to, plus the
/// step number, the step's payload and the signal that reached it.
///
/// Mirrors [`taquba::JobRecord`]: the `payload` is opaque application bytes
/// and `headers` carries user metadata you set at submission (reserved
/// `workflow.*` keys are filtered out before the runner sees them).
///
/// Constructed by the runtime. A test constructs one with
/// [`Step::detached`] and assigns the fields it needs.
#[derive(Debug, Clone)]
pub struct Step {
    /// The delivery this step runs under.
    pub delivery: Delivery,
    /// Zero-based step number. Step 0 is always the first step of a run, with
    /// the original submission input as its `payload`.
    pub step_number: u32,
    /// Application-defined bytes. For step 0 this is the submission `input`;
    /// for later steps it is the bytes returned by the previous step's
    /// [`StepOutcome::Continue`].
    pub payload: Vec<u8>,
    /// The signal payload, when this step was reached through a
    /// [`Trigger::OnSignal`] wait that a signal resolved: the previous
    /// step continued with `OnSignal`, and a
    /// [`crate::WorkflowRuntime::signal`] call for the correlation key
    /// arrived before the timeout. `None` when the timeout elapsed first
    /// and on every step not preceded by an `OnSignal` wait.
    pub signal: Option<Vec<u8>>,
}

impl Deref for Step {
    type Target = Delivery;

    fn deref(&self) -> &Delivery {
        &self.delivery
    }
}

impl DerefMut for Step {
    fn deref_mut(&mut self) -> &mut Delivery {
        &mut self.delivery
    }
}

impl Step {
    /// Step 0 of a [`Delivery::detached`] delivery with `payload` and no
    /// signal, for tests.
    ///
    /// ```
    /// use taquba_workflow::Step;
    ///
    /// let mut step = Step::detached(b"input");
    /// step.step_number = 2;
    /// step.attempts = 3;
    /// assert_eq!(step.payload, b"input");
    /// assert!(step.is_last_attempt());
    /// ```
    pub fn detached(payload: impl Into<Vec<u8>>) -> Self {
        Self {
            delivery: Delivery::detached(),
            step_number: 0,
            payload: payload.into(),
            signal: None,
        }
    }
}

/// When the next step of a run becomes claimable. Set on the `when`
/// field of [`StepOutcome::Continue`].
#[derive(Debug, Clone)]
pub enum Trigger {
    /// The next step is claimable immediately.
    Immediate,
    /// The next step is claimable `Duration` from now.
    After(Duration),
    /// The next step is claimable when a signal for `correlation_key`
    /// arrives via [`crate::WorkflowRuntime::signal`], or after `timeout`,
    /// whichever comes first. The next step reads [`Step::signal`] to
    /// distinguish the two: `Some(payload)` when a signal arrived, `None`
    /// when the timeout elapsed first. A signal that arrived before this
    /// step settled is consumed at settlement and the next step runs
    /// immediately.
    ///
    /// One waiter per correlation key: registering a second waiter while
    /// one is already waiting fails the step permanently. Choose keys
    /// that are unique per waiter (e.g. include the run id).
    OnSignal {
        /// Caller-chosen key the signaller addresses.
        correlation_key: String,
        /// Upper bound on the wait; the next step runs with
        /// [`Step::signal`] `None` when it elapses first.
        timeout: Duration,
    },
}

/// What the runner wants the runtime to do after this step.
#[derive(Debug, Clone)]
pub enum StepOutcome {
    /// Run is not finished. Enqueue the next step with `payload` as its
    /// bytes; `when` decides when it becomes claimable. The runtime
    /// advances `step_number` by 1. The constructors
    /// [`Self::continue_now`] and [`Self::continue_after`] build the
    /// common forms.
    Continue {
        /// Bytes to hand to the next step's [`Step::payload`].
        payload: Vec<u8>,
        /// When the next step becomes claimable.
        when: Trigger,
    },
    /// The run is finished successfully. The runtime acks the step and fires
    /// the configured terminal hook with
    /// [`crate::TerminalStatus::Succeeded`] and `result` as the body.
    Succeed {
        /// Final result bytes handed to the terminal hook.
        result: Vec<u8>,
    },
    /// The run is finished as failed by the runner's verdict; the runner
    /// ran to completion but the workflow's logical outcome is "no" (e.g.
    /// a validation rule rejected the input, a policy check denied the
    /// request, an agent decided the task can't be fulfilled). The runtime
    /// acks the step and fires the terminal hook with
    /// [`crate::TerminalStatus::Failed`] and `reason` as the error.
    ///
    /// Use this for *workflow-level* failures. For *infrastructure*
    /// failures (network outage, downstream service down, etc.) return
    /// `Err(StepError::transient)` or `Err(StepError::permanent)` instead;
    /// those dead-letter the step so an operator can find it via
    /// [`taquba::QueueView::dead_jobs`]. `Fail` is a successful execution with
    /// a negative outcome and does not dead-letter.
    Fail {
        /// Human-readable reason recorded on [`crate::RunOutcome::error`].
        reason: String,
    },
    /// The run is finished as cancelled by the runner. Use this when the
    /// runner decides on its own that the workflow should stop early
    /// without it being a logical failure (e.g. a downstream cancellation
    /// signal arrived mid-step, the user-supplied input is now obsolete).
    /// The runtime acks the step and fires the terminal hook with
    /// [`crate::TerminalStatus::Cancelled`] and `reason` as the error.
    ///
    /// For *external* cancellation requested by another component in the
    /// process, call [`crate::WorkflowRuntime::cancel`] instead; the
    /// runtime translates that into the same `Cancelled` terminal state.
    Cancel {
        /// Human-readable reason recorded on [`crate::RunOutcome::error`].
        reason: String,
    },
}

impl StepOutcome {
    /// Continue the run; the next step is claimable immediately.
    pub fn continue_now(payload: Vec<u8>) -> Self {
        Self::Continue {
            payload,
            when: Trigger::Immediate,
        }
    }

    /// Continue the run; the next step is claimable `delay` from now.
    pub fn continue_after(payload: Vec<u8>, delay: Duration) -> Self {
        Self::Continue {
            payload,
            when: Trigger::After(delay),
        }
    }

    /// Continue the run; the next step is claimable when a signal for
    /// `correlation_key` arrives, or after `timeout` at the latest.
    pub fn continue_on_signal(
        payload: Vec<u8>,
        correlation_key: impl Into<String>,
        timeout: Duration,
    ) -> Self {
        Self::Continue {
            payload,
            when: Trigger::OnSignal {
                correlation_key: correlation_key.into(),
                timeout,
            },
        }
    }
}

/// Failure outcomes the runner can return.
#[derive(Debug, Clone)]
pub struct StepError {
    /// Human-readable message recorded on the underlying job's `last_error`.
    pub message: String,
    /// Whether to retry the step or fail the run immediately.
    pub kind: StepErrorKind,
}

impl StepError {
    /// Build a transient error: Taquba retries the step per the queue's
    /// backoff/`max_attempts`. Once `max_attempts` is exhausted, the step is
    /// dead-lettered and the run terminates as failed.
    pub fn transient(message: impl Into<String>) -> Self {
        Self {
            message: message.into(),
            kind: StepErrorKind::Transient,
        }
    }

    /// Build a permanent error: the step is dead-lettered immediately and the
    /// run terminates as failed.
    pub fn permanent(message: impl Into<String>) -> Self {
        Self {
            message: message.into(),
            kind: StepErrorKind::Permanent,
        }
    }

    /// The worker error reporting this failure: a [`PermanentFailure`]
    /// for a permanent one, a retrying error otherwise.
    pub(crate) fn into_worker_error(self) -> WorkerError {
        match self.kind {
            StepErrorKind::Permanent => PermanentFailure::new(self.message).into(),
            StepErrorKind::Transient => self.message.into(),
        }
    }
}

impl std::fmt::Display for StepError {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.write_str(&self.message)
    }
}

impl std::error::Error for StepError {}

impl From<crate::Error> for StepError {
    fn from(err: crate::Error) -> Self {
        let permanent = err.is_permanent();
        let message = err.to_string();
        if permanent {
            Self::permanent(message)
        } else {
            Self::transient(message)
        }
    }
}

/// Whether a [`StepError`] should retry or fail the run.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StepErrorKind {
    /// Retry per the queue's backoff policy until `max_attempts` is reached.
    Transient,
    /// Dead-letter the step immediately; terminate the run as failed.
    Permanent,
}

/// User-implemented logic that advances a single workflow step.
///
/// Implementations must be idempotent for the same `(run_id, step_number)`:
/// Taquba is at-least-once, so a step can be claimed and processed more than
/// once if a lease expires before the worker acks. Returning the same
/// `StepOutcome` for the same input is the easiest way to satisfy this.
pub trait StepRunner: Send + Sync {
    /// Process a single step of a workflow run. Return [`StepOutcome::Continue`]
    /// to enqueue the next step, [`StepOutcome::Succeed`] to finish the run
    /// successfully, [`StepOutcome::Fail`] to terminate the run as Failed by
    /// runner verdict, [`StepOutcome::Cancel`] to terminate the run as
    /// Cancelled by runner verdict, or `Err(StepError)` to retry /
    /// dead-letter on infrastructure errors.
    fn run_step(
        &self,
        step: &Step,
    ) -> impl Future<Output = std::result::Result<StepOutcome, StepError>> + Send;
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::test_util::rid;

    #[test]
    fn from_workflow_error_maps_via_is_permanent() {
        let permanent: StepError = crate::Error::InputMismatch(rid("run-1")).into();
        assert_eq!(permanent.kind, StepErrorKind::Permanent);

        let store_err = taquba::object_store::Error::NotFound {
            path: "x".into(),
            source: "missing".into(),
        };
        let transient: StepError = crate::Error::Store(store_err).into();
        assert_eq!(transient.kind, StepErrorKind::Transient);
    }
}