taquba-workflow 0.8.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
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
# taquba-workflow

[![crates.io](https://img.shields.io/crates/v/taquba-workflow.svg)](https://crates.io/crates/taquba-workflow)
[![docs.rs](https://img.shields.io/docsrs/taquba-workflow)](https://docs.rs/taquba-workflow)
[![license](https://img.shields.io/crates/l/taquba-workflow.svg)](#license)

Durable execution on object storage: an at-least-once workflow runtime on
top of the [Taquba](../taquba) durable task queue.

> Part of the [Taquba ecosystem]https://github.com/micllam/taquba; see the
> workspace README for the queue core and the other crates that compose with
> this one.

`taquba-workflow` provides the durable machinery for any multi-step process
that benefits from durable state between steps: idempotent step execution,
retries with backoff, graceful restart, and terminal-state notifications.
Implement `StepRunner` with bytes-in / bytes-out per-step logic; the
runtime persists everything else.

Particularly well-suited for **AI agent runs** (see
[`examples/rig_agent.rs`](examples/rig_agent.rs) for a
[Rig](https://github.com/0xPlaygrounds/rig) integration), but the runtime
itself is framework-neutral and equally usable for ETL pipelines, document
processing, payment flows, etc.

## What this is / isn't

`taquba-workflow` is an **imperative step orchestrator**: at each step
the runner decides what happens next via `StepOutcome` (Continue,
Succeed, Fail, Cancel). External cancellation is supported via
`WorkflowRuntime::cancel`. It is *not*:

- **A DAG executor**. There's no declarative graph, no fan-out / fan-in, no
  dependency-driven scheduling.
- **An event-sourced workflow engine**. There's no event-history replay, no
  per-side-effect recording.

Within the ecosystem, [`taquba-jobs`](https://docs.rs/taquba-jobs) is the
sibling crate for single-shot typed tasks: use it when the caller awaits a
typed return value and there are no intermediate steps to persist; use a
workflow (even a single-step one) when the caller observes the run through
cancellation and a terminal hook rather than awaiting a returned value.
[`taquba-bulk`](https://docs.rs/taquba-bulk) builds on this crate to run
one pipeline over many inputs with batch progress and cost rollup.

## Install

```bash
cargo add taquba-workflow taquba
cargo add tokio --features full
```

Enable the `webhooks` feature for `WebhookTerminalHook`:

```bash
cargo add taquba-workflow --features webhooks
```

## Configuring the queue

Per-queue retention (`QueueConfig::keep_done_jobs` and
`QueueConfig::dead_retention`) is set on the `taquba::Queue` before it's
handed to the runtime. Choose an explicit name via
`WorkflowRuntimeBuilder::queue_name` and key `OpenOptions::queue_configs`
on the same string.

```rust
use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use taquba::{OpenOptions, Queue, QueueConfig, object_store::memory::InMemory};
use taquba_workflow::{NoopTerminalHook, StepError, StepOutcome, StepRunner, WorkflowRuntime, Step};

struct EchoRunner;
impl StepRunner for EchoRunner {
    async fn run_step(&self, step: &Step) -> Result<StepOutcome, StepError> {
        Ok(StepOutcome::Succeed { result: step.payload.clone() })
    }
}

let store = Arc::new(InMemory::new());
let opts = OpenOptions {
    queue_configs: HashMap::from([(
        "agent-runs".to_string(),
        QueueConfig {
            keep_done_jobs: Some(Duration::from_secs(24 * 60 * 60)),
            ..QueueConfig::default()
        },
    )]),
    ..OpenOptions::default()
};
let queue = Arc::new(Queue::open_with_options(store.clone(), "db", opts).await?);
let runtime = WorkflowRuntime::builder(queue, store, EchoRunner, NoopTerminalHook)
    .queue_name("agent-runs") // same string as in queue_configs
    .build();
```

## Quick start

```rust
use std::sync::Arc;
use taquba::{Queue, object_store::memory::InMemory};
use taquba_workflow::{
    NoopTerminalHook, RunSpec, Step, StepError, StepOutcome, StepRunner, WorkflowRuntime,
};

struct EchoRunner;

impl StepRunner for EchoRunner {
    async fn run_step(&self, step: &Step) -> Result<StepOutcome, StepError> {
        Ok(StepOutcome::Succeed { result: step.payload.clone() })
    }
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let store = Arc::new(InMemory::new());
    let queue = Arc::new(Queue::open(store.clone(), "demo").await?);

    let runtime = WorkflowRuntime::builder(queue, store, EchoRunner, NoopTerminalHook).build();

    let worker = runtime.clone();
    tokio::spawn(async move { worker.run(std::future::pending::<()>()).await });

    let handle = runtime.submit(RunSpec {
        input: b"hello".to_vec(),
        ..Default::default()
    }).await?;
    println!("submitted run {}", handle.run_id);
    Ok(())
}
```

## Examples

```bash
cargo run -p taquba-workflow --example single_step
cargo run -p taquba-workflow --example multi_step
cargo run -p taquba-workflow --example crash_resume
cargo run -p taquba-workflow --example fanout_jobs
ANTHROPIC_API_KEY=... cargo run -p taquba-workflow --example rig_agent
OPENAI_API_KEY=...    cargo run -p taquba-workflow --example rig_agent
```

`crash_resume` runs on the local filesystem and is the one example that
demonstrates recovery directly: start it, interrupt it during any stage
and start it again. The second process resumes the same run, skips the
stages already committed and serves the completed units of the
interrupted stage from the memo store.

`rig_agent` is a two-stage AI agent (research, then write) structured
for between-step durability: step 0's research is persisted as queue
state before step 1 begins, so on a persistent store a process that
crashes between the steps resumes at step 1 without re-running the
research. The example itself runs on the in-memory store; substitute a
persistent `object_store` backend, as `crash_resume` does, to observe
recovery across a process restart.

`fanout_jobs` composes the runtime with `taquba-jobs` for fan-out
inside one run: a step submits one typed job per URL to a shared
`JobRunner`, joins the typed results, and memoizes the aggregate so a
step retry does not re-submit the fan-out.

## Step outcomes

| Outcome | Effect |
|---|---|
| `StepOutcome::Continue { payload, when }` | Enqueue the next step; `when` (a `Trigger`) decides when it becomes claimable: `Trigger::Immediate`, `Trigger::After(delay)` or `Trigger::OnSignal { correlation_key, timeout }`. Constructors: `StepOutcome::continue_now(payload)`, `StepOutcome::continue_after(payload, delay)`, `StepOutcome::continue_on_signal(payload, key, timeout)`. |
| `StepOutcome::Succeed { result }` | Ack; terminal hook fires `Succeeded`. |
| `StepOutcome::Fail { reason }` | Ack; terminal hook fires `Failed`. Runner verdict: no dead-letter. |
| `StepOutcome::Cancel { reason }` | Ack; terminal hook fires `Cancelled`. Runner verdict: no dead-letter. |
| `Err(StepError::transient(_))` | Retry per backoff up to `max_attempts`, then dead-letter. |
| `Err(StepError::permanent(_))` | Dead-letter immediately. |

`StepOutcome::Fail` / `StepOutcome::Cancel` vs `Err(StepError::permanent)`:
runner verdicts ack normally; an infrastructure error dead-letters so
operators can find it via `queue.dead_jobs()`.

## Cancellation

Call `WorkflowRuntime::cancel(run_id)` to cancel an active run from
outside the runner:

- If the current step is **pending or scheduled**, the queued step job is
  removed and the terminal hook fires from the `cancel` call before it
  returns.
- If the current step is **running**, cancellation is delivered via
  `Step::cancel_token` (a `tokio_util::sync::CancellationToken`).
  Runners that watch the token can short-circuit immediately:

  ```rust,ignore
  tokio::select! {
      out = call_llm(step) => out,
      _ = step.cancel_token.cancelled() => {
          Ok(StepOutcome::Cancel { reason: "cooperative".into() })
      }
  }
  ```

  Runners that ignore the token are allowed to run to completion (futures
  cannot be safely aborted mid-step). In both cases the runner's
  `StepOutcome` is discarded, any pending transient retry is suppressed,
  and the worker fires the terminal hook with `Cancelled` once the step
  returns. Watching the token only reduces cancellation latency for slow
  steps; it doesn't change semantics.

While termination is in flight, `WorkflowRuntime::status` reports a
`RunState::Cancelling` overlay until the entry is dropped.

Returns `Ok(false)` if the run is unknown or already terminal in this
runtime. `cancel` only reaches runs submitted to this `WorkflowRuntime`
instance; a second runtime in the same process (sharing the queue)
maintains its own registry.

## Durable signals

A step can pause the rest of its run until an external event. Returning
`StepOutcome::continue_on_signal` (a `Trigger::OnSignal`) defers the next
step until a signal for the chosen correlation key arrives via
`WorkflowRuntime::signal`, or until the timeout elapses. The next step
reads `Step::signal`: `Some(payload)` when a signal arrived, `None` when
the timeout fired. The natural fit is a run that waits for an approval, a
webhook callback or another run's completion, with the timeout as the
escalation path.

```rust,ignore
// In the runner: pause the run for the payment webhook, or escalate
// after seven days.
Ok(StepOutcome::continue_on_signal(
    order_id.into_bytes(),
    format!("payment:{order_id}"),
    Duration::from_secs(7 * 24 * 3600),
))

// In the webhook handler (same process):
match runtime.signal(&format!("payment:{order_id}"), body).await? {
    SignalOutcome::Delivered => { /* a waiting run was woken */ }
    SignalOutcome::Buffered => { /* held for the next waiter */ }
}
```

Signals are durable in both directions. The waiting step is a scheduled
job in the store, so the wait survives restarts and costs nothing while
pending. A signal with no registered waiter is buffered durably under its
correlation key and consumed by the next waiter registered for it, so a
signal that arrives before its waiter is not lost;
`WorkflowRuntime::clear_signal` discards a buffered signal that is no
longer wanted.

Semantics: delivery follows the crate's at-least-once model (the woken
step can be redelivered and observes the same `Step::signal` value on
every attempt). One buffered signal is held per correlation key; a second
signal before consumption replaces the first. One waiter is allowed per
correlation key; registering a second one fails that run, so choose keys
unique to the waiter (include the run id if uniqueness is uncertain). Signals are
scoped to the store: the signaller is the same process that hosts the
runtime, per the single-process design.

See [`examples/signals.rs`](examples/signals.rs) for a runnable approval
flow covering all three delivery paths (signal, timeout, buffered).

## Reserved headers

Step jobs reserve the `workflow.*` prefix; submission rejects user
headers starting with it. Other headers on `RunSpec::headers` thread
through every step and reach the terminal hook on `RunOutcome::headers`.

| Key | Meaning |
|---|---|
| `workflow.run_id` | Run identifier. |
| `workflow.step` | Zero-based step number. |

## Idempotency

Each step is enqueued with `dedup_key = "run:{run_id}:{step_number}"`,
preventing concurrent duplicate steps. But Taquba is at-least-once: a
step can be claimed and executed twice if its lease expires before ack.
**`StepRunner` implementations must be idempotent for the same
`(run_id, step_number)`.**

## Memoizing within-step side effects

Because retries can re-execute a step, expensive non-idempotent side
effects (LLM calls, paid APIs, multi-stage processing) need a place to
record their result so retries observe the cached value instead of
paying twice. `Step::memo` is a per-step durable key-value store
scoped to `(run_id, step_number)`:

```rust,ignore
// Inside StepRunner::run_step:
if let Some(cached) = step.memo.get("draft").await? {
    return Ok(StepOutcome::Succeed { result: cached });
}
let draft = expensive_call(&step.payload).await?;
step.memo.put("draft", &draft).await?;
Ok(StepOutcome::Succeed { result: draft })
```

When the natural memo key is the content of an input value,
`Memo::content_get` and `Memo::content_put` serialize that input as
MessagePack, hash it with SHA-256, and use the digest as the memo key:

```rust,ignore
#[derive(serde::Serialize)]
struct DraftInput<'a> {
    operation: &'static str,
    payload: &'a [u8],
}

let input = DraftInput {
    operation: "draft",
    payload: &step.payload,
};
if let Some(cached) = step.memo.content_get(&input).await? {
    return Ok(StepOutcome::Succeed { result: cached });
}
let draft = expensive_call(&step.payload).await?;
step.memo.content_put(&input, &draft).await?;
Ok(StepOutcome::Succeed { result: draft })
```

Content-addressed memo keys remain scoped to `(run_id, step_number)`;
they are not a cross-run cache. If multiple logical operations may
receive identical inputs, include an operation name in the
serialized input.

Memo entries live in the object store passed to
`WorkflowRuntime::builder` under the path prefix configured by
`WorkflowRuntimeBuilder::memo_prefix` (default `"workflow-memo"`).
`Memo` is strictly per-step; the durable channel between steps is
`StepOutcome::Continue`'s payload, not memo.

## Step-output replay

`WorkflowRuntimeBuilder::step_output_replay` enables an additional
runtime-managed replay record for every outcome the runner returns,
including `Fail` and `Cancel`. Step errors (`StepError`) are not recorded,
so retries still invoke the runner. The record is keyed by
`(run_id, step_number, SHA-256(step payload))` and is written before the
runtime applies the outcome. If the same step is delivered again after a
crash before ack, the stored outcome is replayed without invoking the
runner again. A replayed `Continue` with a `Trigger::After` delay reduces
the delay by the time already elapsed since the outcome was stored,
preserving the original schedule.

This is disabled by default because it adds one object-store read per step
delivery (the replay lookup) plus one write per recorded outcome, and makes
that write part of step settlement. The replay records are scoped to one
run and step; they are not a cross-run cache. They are cleared with the
run's memo entries when memo retention is configured.

## Memo retention

By default memo entries are retained indefinitely (appropriate for
short-lived runs or workloads that manage cleanup externally). To
enable automatic cleanup, configure a retention window via
`WorkflowRuntimeBuilder::memo_retention`:

```rust,ignore
let runtime = WorkflowRuntime::builder(queue, store, runner, hook)
    .memo_retention(Duration::from_secs(24 * 60 * 60))
    .build();
```

When retention is set, the runtime writes a small terminal marker for
every terminal state (Succeeded, Failed, Cancelled) and
`WorkflowRuntime::run` spawns a background sweeper that lists those
markers and clears the memo entries, step-output replay entries, and
marker for any run whose marker is older than the retention window. The
first sweep fires on startup so a restarted process catches markers left
behind by an earlier one.

Because the sweep is keyed on those terminal markers, and a terminated
run never resumes, it never deletes the memo or replay entries of an
in-flight run that a resume may still read. A resuming step that finds an
entry absent re-executes the work (delivery is at-least-once
regardless), so a missing entry is always safe to observe rather than a
dangling reference: deletion is left unguarded precisely because every
reader tolerates absence.

Advanced cleanup policies (selective retention, externally-driven
sweeps) can be built directly on `MemoStore::list_terminal_markers`,
`MemoStore::clear_memos_for_run`, and `MemoStore::delete_terminal_marker`
without configuring `WorkflowRuntimeBuilder::memo_retention`.

## Time injection

Every timestamp the runtime writes (the `submitted_at_ms` on the
durable per-run record, the `run_at` it computes when a step
continues with a `Trigger::After` delay, and the terminal-marker timestamps the
memo-retention sweep consumes) is read through a `taquba::Clock`
rather than `SystemTime::now()`. By default the runtime inherits
the clock its `Queue` was opened with, so passing a `MockClock`
to `OpenOptions::clock` virtualises both the queue and the
workflow runtime in lockstep:

```rust,ignore
let clock = MockClock::new(1_700_000_000_000);
let opts = OpenOptions {
    clock: Arc::new(clock.clone()),
    ..OpenOptions::default()
};
let queue = Queue::open_with_options(store.clone(), "db", opts).await?;
let runtime = WorkflowRuntime::builder(queue, store, runner, hook).build();
// `runtime` reads the same clock as `queue`; `clock.advance(...)`
// moves every time-based decision the runtime makes.
```

Override the inherited default via `WorkflowRuntimeBuilder::clock`
when a test or specialised setup needs the runtime on a different
time source than the queue. The common case for production callers
is to leave the default and let the queue's `SystemClock` flow
through.

This makes downstream tests deterministic: `Trigger::After` delays,
memo-retention sweep eligibility, and terminal-marker ages all
advance under explicit `MockClock::advance` calls rather than
wall-clock waits.

## Duplicate submissions

`WorkflowRuntime::submit` is idempotent on `(run_id, spec.input)`. A
re-submission of an active run that carries the same input is a no-op
and the returned `SubmitOutcome` has `newly_submitted = false`. A
re-submission that carries a *different* input is rejected with
`Error::InputMismatch`: reusing a `run_id` with new content is a
programmer error; choose a fresh `run_id` for a new run.

Duplicates are caught from two sources, in order:

1. An in-process registry catches duplicates within the same runtime.
2. A **durable per-run record** written atomically with the step-0
   enqueue (via Taquba's `enqueue_with_kv`) catches duplicates across
   process restarts, even after step 0 has been claimed and its dedup
   key released. The record carries a SHA-256 of the original input so
   the cross-restart mismatch check works even when the in-memory
   registry is empty. The record is cleaned up when the run reaches a
   terminal state.

## Terminal hook

`TerminalHook::on_termination` fires once per run on `Succeeded`,
`Failed`, or `Cancelled`, receiving the submitter's headers and the
runner's result or error. `WebhookTerminalHook` (behind the `webhooks`
feature) fires HTTP callbacks via `taquba-webhooks`; set the per-run URL
on `RunSpec::headers["callback_url"]`.

## License

Licensed under either of

 * Apache License, Version 2.0
   ([LICENSE-APACHE]LICENSE-APACHE or
   <http://www.apache.org/licenses/LICENSE-2.0>)
 * MIT license
   ([LICENSE-MIT]LICENSE-MIT or
   <http://opensource.org/licenses/MIT>)

at your option.

## Contribution

Unless you explicitly state otherwise, any contribution intentionally
submitted for inclusion in the work by you, as defined in the Apache-2.0
license, shall be dual licensed as above, without any additional terms or
conditions.