Skip to main content

taquba_workflow/
lib.rs

1//! Durable execution for Rust on object storage: an at-least-once workflow
2//! runtime over the [Taquba] durable task queue. A run is a sequence of
3//! steps whose state (each step's output, memoized side effects and
4//! application KV effects) is stored in a bucket or a local directory, so
5//! a process restarted mid-run resumes at its last committed step.
6//!
7//! The runtime is an embedded library: producers and workers share one
8//! `Arc<Queue>` in one process, and no server, database or control plane
9//! runs beside it. It is built for long-running, expensive, IO- or
10//! API-bound work: LLM agent runs (see `examples/rig_agent.rs` for a Rig
11//! integration), document pipelines, payment flows and command-line tools
12//! whose runs survive an interruption. Implement [`StepRunner`] with
13//! bytes-in, bytes-out per-step logic; the runtime persists everything
14//! else.
15//!
16//! # Scope
17//!
18//! `taquba-workflow` is an imperative step orchestrator: at each step the
19//! runner returns a [`StepOutcome`] (Continue, Succeed, Fail or Cancel)
20//! that decides what happens next, and [`WorkflowRuntime::cancel`] cancels
21//! a run from outside. It is neither a DAG executor (no declarative graph,
22//! no fan-out or fan-in, no dependency-driven scheduling) nor an
23//! event-sourced engine (no event-history replay; a side effect is
24//! recorded only where the runner memoizes it).
25//!
26//! The [`jobs`] module is the typed presentation of the runtime: it runs
27//! one typed async function as a single-step run and returns its result
28//! to an awaiting caller, and a [`jobs::JobGroup`] submits many such
29//! jobs as one durable set. Use a job when the caller awaits a typed
30//! return value and there are no intermediate steps to persist, a job
31//! group when many inputs go through one function and the caller reads
32//! the results as they complete, and a step runner, even for a single
33//! step, when the caller observes the run through cancellation, headers
34//! and a terminal hook.
35//!
36//! # Quick start
37//!
38//! ```no_run
39//! use std::sync::Arc;
40//! use taquba::{Queue, object_store::memory::InMemory};
41//! use taquba_workflow::{
42//!     NoopTerminalHook, RunSpec, Step, StepError, StepOutcome, StepRunner, WorkflowRuntime,
43//! };
44//!
45//! struct EchoRunner;
46//!
47//! impl StepRunner for EchoRunner {
48//!     async fn run_step(&self, step: &Step) -> Result<StepOutcome, StepError> {
49//!         Ok(StepOutcome::Succeed { result: step.payload.clone() })
50//!     }
51//! }
52//!
53//! # async fn run() -> Result<(), Box<dyn std::error::Error>> {
54//! let store = Arc::new(InMemory::new());
55//! let queue = Arc::new(Queue::open(store.clone(), "demo").await?);
56//!
57//! let runtime = WorkflowRuntime::builder(queue, store, EchoRunner, NoopTerminalHook).build();
58//!
59//! let worker = runtime.spawn(std::future::pending::<()>());
60//!
61//! let outcome = runtime.submit(RunSpec {
62//!     input: b"hello".to_vec(),
63//!     ..Default::default()
64//! }).await?;
65//! println!("submitted run {}", outcome.run_id);
66//! worker.shutdown().await?;
67//! # Ok(()) }
68//! ```
69//!
70//! Replace `InMemory` with an S3, GCS or Azure builder in production.
71//! The runtime's queue is named by [`WorkflowRuntimeBuilder::queue_name`];
72//! per-queue settings in [`taquba::OpenOptions::queue_configs`]
73//! (retention, lease duration, attempt limit) are keyed on that name when
74//! the [`taquba::Queue`] is opened.
75//!
76//! # Step outcomes
77//!
78//! | Outcome | Effect |
79//! |---|---|
80//! | [`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`], [`StepOutcome::continue_after`], [`StepOutcome::continue_on_signal`]. |
81//! | [`StepOutcome::Succeed`]` { result }` | Ack; the terminal hook observes [`TerminalStatus::Succeeded`]. |
82//! | [`StepOutcome::Fail`]` { reason }` | Ack; the terminal hook observes [`TerminalStatus::Failed`]. A runner verdict: no dead-letter. |
83//! | [`StepOutcome::Cancel`]` { reason }` | Ack; the terminal hook observes [`TerminalStatus::Cancelled`]. A runner verdict: no dead-letter. |
84//! | `Err(`[`StepError::transient`]`(_))` | Retry per backoff up to `max_attempts`, then dead-letter. |
85//! | `Err(`[`StepError::permanent`]`(_))` | Dead-letter immediately. |
86//!
87//! [`StepOutcome::Fail`] and [`StepOutcome::Cancel`] are runner verdicts
88//! and acknowledge normally; an `Err(`[`StepError::permanent`]`)` is an
89//! infrastructure error and dead-letters, so operators find it through
90//! [`taquba::QueueView::dead_jobs`].
91//!
92//! # The delivery
93//!
94//! A [`Step`] dereferences to its [`Delivery`]: the run id, the
95//! submitter's headers, the queue job id, the attempt count and limit
96//! ([`Delivery::is_last_attempt`] reports whether a transient error from
97//! this attempt dead-letters the step) and the delivery's handles (the
98//! cancellation token, the lease, the per-step and run-scoped memos, the
99//! staged KV effects and committed KV reads). [`jobs::JobContext`]
100//! dereferences to the same type. [`Delivery::detached`] and [`Step::detached`] build instances
101//! bound to no queue, for tests.
102//!
103//! # Submissions
104//!
105//! [`WorkflowRuntime::submit`] takes a [`RunSpec`]: the first step's `input`,
106//! an optional `run_id` (a [`RunId`], validated at construction to 1 to
107//! [`MAX_RUN_ID_LEN`] bytes of `[A-Za-z0-9_-]`, and a ULID is generated when
108//! absent), the [`RunOptions`] of its steps (`headers`, a `priority` and
109//! `max_attempts_per_step` overriding the queue's defaults for every step and a
110//! `run_at` before which the first step is not claimable) and the `effects`
111//! applied with the enqueue. The returned [`SubmitOutcome`] identifies the run
112//! and the queue job that currently represents it.
113//!
114//! `submit` is idempotent on `(run_id, input)`. A re-submission of an active
115//! run with the same input is a no-op and the returned [`SubmitOutcome`] has
116//! `newly_submitted = false`. A re-submission with a different input is
117//! rejected with [`Error::InputMismatch`]. Duplicates are caught by a durable
118//! per-run record written atomically with the step-0 enqueue (via
119//! [`taquba::Queue::enqueue_with_effects`]), so they are caught across process
120//! restarts, even after step 0 is claimed and its dedup key is released. The
121//! record contains a SHA-256 of the original input for the mismatch check. A
122//! current-step pointer under `workflow/steps/` is written with it, rewritten
123//! in the settlement that enqueues each next step and identifies the queue job
124//! that [`SubmitOutcome::job_id`] reports for a duplicate. Both are removed
125//! when the run reaches a terminal state.
126//!
127//! [`WorkflowView::status`] reads the record, the pointer and the step's queue
128//! job into a [`RunStatus`] ([`RunState::Pending`], [`RunState::Running`] or
129//! [`RunState::Cancelling`], with the current step number).
130//! [`WorkflowRuntime::status`] reads through the runtime's view, so the status
131//! is available after a restart and from any runtime over the same queue. A
132//! process without a runtime builds a [`WorkflowView`] from a
133//! `taquba::QueueReader` view and a [`MemoStore`] at the runtime's memo prefix.
134//! A terminated run reports [`RunState::Terminated`] with its status, error,
135//! error kind, final step and time of termination, read from the terminal
136//! record written with the terminating settlement, which
137//! [Memo retention](#memo-retention) removes with the run's memo entries.
138//!
139//! [`WorkflowRuntime::wait`] waits until a run terminates, following its
140//! current step across steps, and reports a [`RunEnd`]: the termination
141//! from the run's terminal record and the committed outcome, when the
142//! worker that terminated the run recorded one. A run
143//! already terminated is reported at once, and
144//! [`WorkflowRuntime::wait_timeout`] bounds the wait. The wait relies on
145//! the queue's in-process completion notification, so it runs in the
146//! process that runs the worker.
147//!
148//! [`WorkflowRuntime::outcome`] returns the committed [`RunOutcome`] of a
149//! terminated run (its result or error, the submitter's headers and the
150//! final step) from the run result record the worker writes to the run's
151//! memo, under the reserved key `workflow.outcome`, before every
152//! terminating settlement it performs. A run terminated without a worker
153//! (a cancellation of a pending step, a step dead-lettered outside the
154//! worker) has no record. A record belongs to the termination the run's
155//! terminal record describes, so a re-submission of a terminated run id
156//! does not report the earlier run's record. The record is removed with
157//! the run's memo entries by the memo sweep.
158//!
159//! # Cancellation
160//!
161//! [`WorkflowRuntime::cancel`] cancels an active run from outside the
162//! runner. The request is recorded on the run's durable record, so it
163//! survives a restart and reaches the run from any runtime over the same
164//! queue; while termination is in flight, `status` reports
165//! [`RunState::Cancelling`]. It returns `Ok(false)` if the run is unknown
166//! or already terminal.
167//!
168//! - If the current step is pending or scheduled, the queued step job is
169//!   removed and the run's notification job is enqueued before the call
170//!   returns.
171//! - If the current step is running, the request is delivered through
172//!   [`Delivery::cancel_token`] (a `tokio_util::sync::CancellationToken`).
173//!   A runner that watches the token returns at once:
174//!
175//!   ```ignore
176//!   tokio::select! {
177//!       out = call_llm(step) => out,
178//!       _ = step.cancel_token.cancelled() => {
179//!           Ok(StepOutcome::Cancel { reason: "cooperative".into() })
180//!       }
181//!   }
182//!   ```
183//!
184//!   A runner that ignores the token runs to completion (a future cannot
185//!   be aborted safely mid-step). In both cases the runner's
186//!   [`StepOutcome`] is discarded, any pending transient retry is
187//!   suppressed and the worker settles the run as
188//!   [`TerminalStatus::Cancelled`] once the step returns. Watching the
189//!   token reduces the latency of cancelling a slow step; the semantics
190//!   are the same.
191//! - A step claimed after the request is settled as cancelled without
192//!   running.
193//!
194//! # Long-running steps
195//!
196//! A step that outlives the queue's lease is re-queued by the reaper
197//! and delivered a second time. A long-running runner avoids this by
198//! extending its lease through [`Delivery::lease`]: call
199//! [`taquba::LeaseHandle::ensure_at_least`] at progress points, or
200//! once, with a slow call's timeout, before issuing the call.
201//!
202//! # Durable signals
203//!
204//! A step can pause the rest of its run until an external event. Returning
205//! [`StepOutcome::continue_on_signal`] (a [`Trigger::OnSignal`]) defers the
206//! next step until a signal for the chosen correlation key arrives via
207//! [`WorkflowRuntime::signal`], or until the timeout elapses. The next
208//! step reads [`Step::signal`]: `Some(payload)` when a signal arrived,
209//! `None` when the timeout fired. The natural fit is a run that waits for
210//! an approval, a webhook callback or another run's completion, with the
211//! timeout as the escalation path.
212//!
213//! ```ignore
214//! // In the runner: pause the run for the payment webhook, or escalate
215//! // after seven days.
216//! Ok(StepOutcome::continue_on_signal(
217//!     order_id.into_bytes(),
218//!     format!("payment:{order_id}"),
219//!     Duration::from_secs(7 * 24 * 3600),
220//! ))
221//!
222//! // In the webhook handler (same process):
223//! match runtime.signal(&format!("payment:{order_id}"), body).await? {
224//!     SignalOutcome::Delivered => { /* a waiting run was woken */ }
225//!     SignalOutcome::Buffered => { /* held for the next waiter */ }
226//! }
227//! ```
228//!
229//! Signals are durable in both directions. The waiting step is a
230//! scheduled job in the store, so the wait survives restarts and
231//! occupies no worker while pending. A signal with no registered waiter
232//! is buffered durably under its correlation key and consumed by the
233//! next waiter registered for it, so a signal that arrives before its
234//! waiter is not lost; [`WorkflowRuntime::clear_signal`] discards a
235//! buffered signal that is no longer wanted. A waiting run resumes
236//! under the runner hosted by the resuming process, so a runner changed
237//! while the run waits is the caller's compatibility concern.
238//!
239//! Delivery follows the crate's at-least-once model: the woken step can
240//! be redelivered and observes the same [`Step::signal`] value on every
241//! attempt. One buffered signal is held per correlation key; a second
242//! signal before consumption replaces the first. One waiter is allowed
243//! per correlation key; registering a second one fails that run, so
244//! choose keys unique to the waiter (include the run id if uniqueness is
245//! uncertain). Signals are scoped to the store: the signaller is the
246//! same process that hosts the runtime, per the single-process design.
247//!
248//! See `examples/durable_approvals.rs` for a runnable approval flow
249//! covering all three delivery paths (signal, timeout and buffered)
250//! across process restarts.
251//!
252//! A signal is also how work on another machine reports back without
253//! an inbound endpoint on this process. The requesting step chooses a
254//! reply key in the bucket, sends it with the request and continues on
255//! a signal. The remote worker writes its reply at that key, and a
256//! watcher task in this process delivers the key as the signal when the
257//! object appears. No lease is held during the wait, because the
258//! requesting step settles before the remote work starts. See `examples/remote_reply.rs`
259//! for a runnable version with the pending-marker layout the watcher
260//! reads.
261//!
262//! # Application KV effects
263//!
264//! Application state that describes a run (a status row, a progress
265//! marker, an outcome record) can be written to Taquba's caller KV
266//! namespace in the same transaction as the run's own transitions, so
267//! a crash cannot leave the two disagreeing. Two surfaces:
268//!
269//! - [`RunSpec::effects`]: a [`taquba::SettlementEffects`] applied
270//!   atomically with the step-0 enqueue. A duplicate submission drops its
271//!   effects.
272//! - [`Delivery::effects`]: an [`EffectsHandle`] that stages writes and
273//!   deletes during a step. Everything staged is applied in the
274//!   settlement transaction that commits the outcome the runner
275//!   returned, whichever outcome that is (`Continue`, `Succeed`,
276//!   `Fail` or `Cancel`).
277//!
278//! ```ignore
279//! // Inside StepRunner::run_step: the outcome record commits with the
280//! // step's own settlement.
281//! step.effects.put(format!("app/runs/{}", step.run_id), summary)?;
282//! Ok(StepOutcome::Succeed { result })
283//! ```
284//!
285//! Semantics:
286//!
287//! - Delivery is at-least-once, so a retried step stages its effects
288//!   again; every staged value must be correct when applied more than
289//!   once (write absolute values).
290//! - No effects are applied when the runner returns a [`StepError`]
291//!   (the step retries or dead-letters) or when an external
292//!   [`WorkflowRuntime::cancel`] overrides the outcome. A runner-issued
293//!   [`StepOutcome::Cancel`] keeps its effects.
294//! - Operations are validated as they are staged: the `workflow/`
295//!   prefix ([`RESERVED_KV_PREFIX`]) is reserved for the runtime,
296//!   values are capped at [`taquba::MAX_KV_VALUE_SIZE`] and a key
297//!   cannot be staged for both a write and a delete within one step.
298//! - With [`WorkflowRuntimeBuilder::step_output_replay`] enabled, the
299//!   replay record stores the staged effects with the outcome, so a
300//!   replayed delivery applies them without invoking the runner.
301//!
302//! The written values are readable inside a step through [`Delivery::kv`]
303//! (a [`KvReadHandle`] exposing `get` only, answering from committed
304//! state, so effects staged by the running step are excluded), through
305//! [`taquba::QueueView::kv_get`] and, from another process, through a
306//! `taquba::QueueReader`.
307//!
308//! See `examples/kv_effects.rs` for a runnable order flow maintaining a
309//! status row through both surfaces.
310//!
311//! # Reserved headers
312//!
313//! Step jobs reserve the `workflow.*` header prefix
314//! ([`RESERVED_HEADER_PREFIX`]); submission rejects user headers starting
315//! with it. [`HEADER_RUN_ID`] and [`HEADER_STEP`] are set by the runtime
316//! on every step. Other headers on [`RunOptions::headers`] thread through
317//! every step and reach the terminal hook on [`RunOutcome::headers`].
318//!
319//! # Run groups
320//!
321//! A [`RunGroup`] is a durable set of runs of one runtime, identified by
322//! a group id: [`WorkflowRuntime::group`] names it, [`RunGroup::submit`]
323//! writes its manifest (the members' keys and inputs) and submits the
324//! members, [`RunGroup::results`] yields each member's [`MemberResult`]
325//! (its termination and, when a worker recorded one, its [`RunOutcome`])
326//! as it terminates, [`RunGroup::cancel`] cancels every active member
327//! and [`RunGroup::status`] counts the members by state. A member's run
328//! id is derived from the group id and its key, so groups never share
329//! run state. A second submission of the same set, or
330//! [`RunGroup::resume`] from the manifest alone, runs again only the
331//! members that did not succeed, which is how a step that fans out
332//! stays safe under a retry and how a batch of inputs is run again
333//! after a partial failure.
334//!
335//! The group's durable state is the manifest in the object store and one
336//! member record per key under `workflow/groups/` in the queue's
337//! key-value namespace, written with the member's submission and
338//! rewritten with its status and error by the settlement that terminates
339//! it. [`RunGroup::forget`] removes it with the members' memo entries
340//! and terminal records, and
341//! [`WorkflowRuntimeBuilder::group_retention`] removes it a window after
342//! a [`RunGroup::results`] consumer observed the last termination,
343//! through a sweep over `workflow/group-terminals/`; a group whose
344//! results are never consumed is retained until it is forgotten.
345//!
346//! # Typed jobs
347//!
348//! The [`jobs`] module runs one typed async function as a single-step run
349//! and returns its result to an awaiting caller: define a [`jobs::Job`]
350//! with typed input fields, an `Output` and an `Error`, register it on a
351//! [`jobs::JobRunner`], submit instances and await the
352//! [`jobs::JobHandle`]. The result is the run result record in the run's
353//! memo, so [`jobs::JobHandle::fetch_result`] reads it after a restart, and a
354//! [`jobs::Job::idempotency_key`] collapses duplicate submissions before
355//! and after completion. A handler that submits further jobs holds a
356//! [`jobs::JobRunner`] in its registered state. The module documentation
357//! covers idempotent submission, retention and the handler context.
358//!
359//! ```ignore
360//! use taquba_workflow::jobs::{Job, JobContext, JobRunner};
361//!
362//! #[derive(serde::Serialize, serde::Deserialize)]
363//! struct SendEmail { to: String }
364//!
365//! impl Job for SendEmail {
366//!     const NAME: &'static str = "email.send";
367//!     type Output = String;
368//!     type Error = EmailError;
369//!
370//!     async fn run(&self, _ctx: JobContext<'_>) -> Result<String, EmailError> {
371//!         Ok(format!("msg-for-{}", self.to))
372//!     }
373//! }
374//!
375//! let runner = JobRunner::builder(queue, store)
376//!     .register::<SendEmail>()
377//!     .build();
378//! let worker = runner.spawn(std::future::pending::<()>());
379//! let message_id = runner.submit(SendEmail { to: "user@example.com".into() }).await?.await?;
380//! worker.shutdown().await?;
381//! ```
382//!
383//! # Job groups
384//!
385//! A [`jobs::JobGroup`] is a [run group](#run-groups) of jobs of one
386//! type: [`jobs::JobRunner::group`] names it, [`jobs::JobGroup::submit`]
387//! takes the jobs, keyed by the job's idempotency key or the positional
388//! `item-{i}`, [`jobs::JobGroup::results`] yields each member's typed
389//! result as it terminates and [`jobs::JobGroup::join`] returns them in
390//! submission order; `resume`, `status`, `cancel`, `forget` and
391//! [`jobs::JobRunnerBuilder::group_retention`] are the run group's.
392//! Streamed output, progress, a failure threshold and a cost rollup are
393//! folds the caller writes over the results;
394//! `examples/group_document_pipeline.rs` shows a per-document pipeline
395//! of memoized stages with counters rolled up by the caller.
396//!
397//! # Idempotency
398//!
399//! Each step is enqueued with [`taquba::EnqueueOptions::dedup_key`] of
400//! `"run:{run_id}:{step_number}"`, so no two pending or scheduled jobs
401//! exist for the same step at the same time. Delivery is at-least-once,
402//! so a step can still be claimed and executed twice if its lease expires
403//! before its acknowledgement: **implementations of [`StepRunner`] must
404//! be idempotent for the same `(run_id, step_number)`.**
405//!
406//! # Memoizing within-step side effects
407//!
408//! Because a retry can re-execute a step, an expensive non-idempotent
409//! side effect (an LLM call, a paid API, one stage of a multi-stage step)
410//! records its result so a retry observes the recorded value and does
411//! not repeat the call. [`Delivery::memo`] is a per-step durable
412//! key-value store scoped to `(run_id, step_number)`:
413//!
414//! ```ignore
415//! // Inside StepRunner::run_step:
416//! if let Some(cached) = step.memo.get("draft").await? {
417//!     return Ok(StepOutcome::Succeed { result: cached });
418//! }
419//! let draft = expensive_call(&step.payload).await?;
420//! step.memo.put("draft", &draft).await?;
421//! Ok(StepOutcome::Succeed { result: draft })
422//! ```
423//!
424//! [`Memo::memoized`] is the typed form: it returns the stored value
425//! when one exists and otherwise runs the computation, stores its value
426//! and returns it. Values are encoded as MessagePack with named fields.
427//! An entry that fails to decode is treated as absent and is
428//! overwritten by the recomputed value; an error from the computation
429//! stores nothing.
430//!
431//! ```ignore
432//! let draft: Draft = step
433//!     .memo
434//!     .memoized("draft", async { expensive_call(&step.payload).await })
435//!     .await?;
436//! ```
437//!
438//! When the natural memo key is the content of an input value,
439//! [`Memo::memoized_by_content`] (and the untyped [`Memo::content_get`]
440//! and [`Memo::content_put`]) serializes that input as MessagePack,
441//! hashes it with SHA-256 and uses the digest as the memo key. The entry
442//! remains scoped to `(run_id, step_number)`; it is not a cross-run
443//! cache. If several logical operations may receive identical inputs,
444//! include an operation name in the serialized input.
445//! [`Memo::content_key`] returns the derived key, for use with
446//! [`Memo::get`] and [`Memo::put`] or for locating an entry from outside
447//! the runtime.
448//!
449//! ```ignore
450//! #[derive(serde::Serialize)]
451//! struct DraftInput<'a> {
452//!     operation: &'static str,
453//!     payload: &'a [u8],
454//! }
455//!
456//! let input = DraftInput { operation: "draft", payload: &step.payload };
457//! let draft: Draft = step
458//!     .memo
459//!     .memoized_by_content(&input, async { expensive_call(&step.payload).await })
460//!     .await?;
461//! ```
462//!
463//! [`Delivery::run_memo`] is the run-scoped variant: one namespace shared
464//! by every step of the run, for values a later step reads back (an
465//! accumulating journal, for example); the key `workflow.outcome` in it
466//! is reserved for the run result record. Its entries are stored beside
467//! the per-step entries and are removed with them when the run's
468//! retention expires.
469//!
470//! Memo entries are stored in the object store passed to
471//! [`WorkflowRuntime::builder`] under the path prefix configured by
472//! [`WorkflowRuntimeBuilder::memo_prefix`] (default `"{queue_name}-memo"`).
473//! A memo is a retry-safety cache whose readers tolerate absence by
474//! re-executing; the durable channel between steps is
475//! [`StepOutcome::Continue`]'s payload.
476//!
477//! # Step-output replay
478//!
479//! [`WorkflowRuntimeBuilder::step_output_replay`] enables a
480//! runtime-managed replay record for every outcome the runner returns,
481//! including `Fail` and `Cancel`. Step errors ([`StepError`]) are not
482//! recorded, so a retry still invokes the runner. The record is keyed by
483//! `(run_id, step_number, SHA-256(step payload))` and is written before
484//! the runtime applies the outcome. If the same step is delivered again
485//! after a crash before its acknowledgement, the stored outcome is
486//! replayed without invoking the runner. The record includes the effects
487//! staged through [`Delivery::effects`], so a replayed outcome applies
488//! them as well. A replayed [`StepOutcome::Continue`] with a
489//! [`Trigger::After`] delay reduces the delay by the time already elapsed
490//! since the outcome was stored, preserving the original schedule.
491//!
492//! Replay is disabled by default because it adds one object-store read
493//! per step delivery (the replay lookup) plus one write per recorded
494//! outcome, and makes that write part of step settlement. The records
495//! are scoped to one run and step, and are removed with the run's memo
496//! entries when memo retention is configured.
497//!
498//! # Memo retention
499//!
500//! By default memo entries are retained indefinitely. To remove them
501//! automatically, configure a retention window via
502//! [`WorkflowRuntimeBuilder::memo_retention`]:
503//!
504//! ```ignore
505//! let runtime = WorkflowRuntime::builder(queue, store, runner, hook)
506//!     .memo_retention(Duration::from_secs(24 * 60 * 60))
507//!     .build();
508//! ```
509//!
510//! Every settlement that commits a terminal outcome (`Succeeded`,
511//! `Failed` or `Cancelled`) writes a terminal record under
512//! `workflow/outcomes/` in the queue's key-value namespace, holding the
513//! status, the error, the final step and the time of termination. When
514//! retention is set, the same transaction writes a terminal marker under
515//! `workflow/terminals/`, so the marker exists exactly when the run's
516//! terminal outcome committed. [`WorkflowRuntime::run`] sweeps the
517//! markers on startup and at every poll interval, removing the memo
518//! entries, the step-output replay entries, the terminal record and the
519//! marker of every run whose marker is at least a window old. A marker
520//! is an entry of a [`taquba::ExpiryIndex`] with the run id as its
521//! suffix, so the sweep reads the expired set from the start of the
522//! range and stops at the first unexpired marker. A pass before a
523//! marker can be expired does not read the index, so a run's state is
524//! removed within a poll interval of the end of its window.
525//!
526//! Because the sweep is keyed on terminal markers and a terminated run
527//! never resumes, the entries of an in-flight run are not removed, with
528//! one exception: entries are addressed by run id, and a terminated run
529//! releases its id, so a second run submitted under that id shares the
530//! first run's entries, and the first run's marker expires against them
531//! while the second run may still be executing. The second run
532//! re-executes the affected steps. Deletion is unguarded because every
533//! reader tolerates absence: a step that finds an entry absent
534//! re-executes the work, as it would under at-least-once delivery
535//! anyway.
536//!
537//! Other cleanup policies (selective retention, externally-driven
538//! sweeps) can be built on [`taquba::QueueView::kv_scan`] over that prefix
539//! and [`MemoStore::clear_memos_for_run`], without configuring
540//! [`WorkflowRuntimeBuilder::memo_retention`].
541//!
542//! # Time injection
543//!
544//! Every timestamp the runtime writes (the `submitted_at_ms` on the
545//! durable per-run record, the `run_at` it computes for a
546//! [`Trigger::After`] delay and the terminal-marker timestamps the
547//! memo-retention sweep consumes) is read through a [`taquba::Clock`].
548//! By default the runtime inherits the clock its [`taquba::Queue`] was
549//! opened with, so passing a [`taquba::MockClock`] to
550//! [`taquba::OpenOptions::clock`] virtualises both in lockstep, and
551//! `MockClock::advance` moves every time-based decision the runtime
552//! makes: [`Trigger::After`] delays, sweep eligibility and
553//! terminal-marker ages.
554//!
555//! ```rust,ignore
556//! let clock = MockClock::new(1_700_000_000_000);
557//! let opts = OpenOptions::default().clock(Arc::new(clock.clone()));
558//! let queue = Queue::open_with_options(store.clone(), "db", opts).await?;
559//! let runtime = WorkflowRuntime::builder(queue, store, runner, hook).build();
560//! // `runtime` reads the same clock as `queue`.
561//! ```
562//!
563//! [`WorkflowRuntimeBuilder::clock`] overrides the inherited default
564//! when a test or specialised setup needs the runtime on a different
565//! time source than the queue.
566//!
567//! # Terminal hook
568//!
569//! [`TerminalHook::on_termination`] processes a run's termination
570//! ([`TerminalStatus::Succeeded`], [`TerminalStatus::Failed`] or
571//! [`TerminalStatus::Cancelled`]), receiving the submitter's headers and
572//! the runner's result or error. Termination is delivered as a queue
573//! job: the settlement that commits a run's terminal outcome atomically
574//! enqueues a notification job, and the hook runs as that job's worker.
575//! The hook therefore observes only outcomes that committed, and
576//! delivery is at-least-once, so implementations must be idempotent. A
577//! transient error ([`StepError::transient`]) retries the notification
578//! per the queue's backoff up to the terminal step's `max_attempts`; a
579//! permanent error dead-letters it.
580//!
581//! The hook stages effects on a [`TerminalEffects`] handle: KV writes
582//! and deletes plus follow-up enqueues, applied in the same transaction
583//! as the notification's acknowledgement. [`TerminalHook::observes`]
584//! (default `true`) is consulted when a run terminates; returning
585//! `false` skips the notification job for that run. [`NoopTerminalHook`]
586//! observes nothing, so runs terminate with no notification cost.
587//!
588//! Runs terminated without an acknowledging settlement (an external
589//! cancellation of a pending step, a step that dead-letters) settle
590//! their notification the same way: the effects are applied by the
591//! dead-letter, by the attempts-exhausting nack or by the
592//! cancellation's removal, so the notification job is created exactly
593//! once on every worker and cancellation path. Two terminations occur
594//! outside any settlement the runtime performs: a job the reaper
595//! dead-letters after its lease expires past the attempt limit, and one
596//! dead-lettered during crash recovery when the queue is opened. The
597//! worker reconciles them: whenever the queue's dead count changes it
598//! terminates every run whose dead step job still has a run record, as
599//! `Failed` with the queue record's last error, enqueueing the
600//! notification in the same transaction.
601//!
602//! `WebhookTerminalHook` (behind the `webhooks` feature) delivers HTTP
603//! callbacks via `taquba-webhooks`, staging the delivery enqueue as a
604//! notification effect so it is created exactly once with the
605//! acknowledgement; set the per-run URL on
606//! [`RunOptions::headers`]`["callback_url"]`. Runs without that header
607//! enqueue no notification.
608//!
609//! [Taquba]: https://docs.rs/taquba
610
611#![warn(missing_docs)]
612
613mod blob;
614mod durable;
615mod effects;
616mod error;
617mod group;
618pub mod jobs;
619mod keys;
620mod kv;
621mod memo;
622mod runner;
623mod runtime;
624mod signal;
625mod sweep;
626mod terminal;
627#[cfg(test)]
628mod test_util;
629mod view;
630mod worker;
631
632pub use effects::{EffectsHandle, TerminalEffects};
633pub use error::{Error, Result};
634pub use group::{GroupMember, GroupStatus, MemberResult, RunGroup};
635pub use keys::{
636    HEADER_RUN_ID, HEADER_SIGNAL_DELIVERED, HEADER_SIGNAL_WAIT, HEADER_STEP, HEADER_TERMINAL,
637    MAX_RUN_ID_LEN, RESERVED_HEADER_PREFIX, RESERVED_KV_PREFIX, RunId,
638};
639pub use kv::KvReadHandle;
640pub use memo::{Memo, MemoStore};
641pub use runner::{Delivery, Step, StepError, StepErrorKind, StepOutcome, StepRunner, Trigger};
642pub use runtime::{
643    RunEnd, RunOptions, RunSpec, RunState, RunStatus, RunTermination, RunnerHandle, SubmitOutcome,
644    WorkflowRuntime, WorkflowRuntimeBuilder,
645};
646pub use signal::SignalOutcome;
647#[cfg(feature = "webhooks")]
648pub use terminal::WebhookTerminalHook;
649pub use terminal::{NoopTerminalHook, RunOutcome, TerminalHook, TerminalStatus};
650pub use view::WorkflowView;