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
//! The worker path: how a claimed job is identified, run and settled.
//! [`StepWorker`] is the [`Worker`] the runtime drives; a step job is
//! parsed into a [`ClaimedStep`], run through the [`StepRunner`] and
//! settled into the [`SettlementEffects`] of its outcome, and a
//! terminal-notification job runs the [`TerminalHook`].

use std::collections::HashMap;
use std::sync::Arc;

use taquba::{
    FailWith, JobRecord, LeaseHandle, PermanentFailure, SettlementEffects, Worker, WorkerError,
};
use tracing::{debug, warn};

use crate::durable::{self, DurableRunOutcome};
use crate::effects::{EffectsHandle, TerminalEffects};
use crate::error::{Error, worker_error};
use crate::group::Membership;
use crate::keys::{HEADER_RUN_ID, HEADER_STEP, HEADER_TERMINAL, RESERVED_HEADER_PREFIX, RunId};
use crate::kv::KvReadHandle;
use crate::runner::{Delivery, Step, StepError, StepErrorKind, StepOutcome, StepRunner, Trigger};
use crate::runtime::{RuntimeInner, StepEnqueueOpts};
use crate::terminal::{RunOutcome, TerminalHook, TerminalStatus};

/// The [`Worker`] of a runtime: every claimed job of the runtime's
/// queue is processed by [`RuntimeInner::process_step`].
pub(crate) struct StepWorker<R, H> {
    pub(crate) inner: Arc<RuntimeInner<R, H>>,
}

impl<R: StepRunner + 'static, H: TerminalHook + 'static> Worker for StepWorker<R, H> {
    async fn process_with_effects(
        &self,
        job: &JobRecord,
        lease: &LeaseHandle,
    ) -> std::result::Result<SettlementEffects, WorkerError> {
        self.inner.process_step(job, lease).await
    }
}

/// One claimed step job as the worker path identifies it: the run and
/// step named by the job's reserved headers, the run's group membership
/// when it has one, the submitter's headers with the reserved ones
/// removed and the queue record itself.
pub(crate) struct ClaimedStep<'a> {
    pub(crate) run_id: RunId,
    pub(crate) step_number: u32,
    pub(crate) membership: Option<Membership>,
    /// Submitter-supplied headers, without the reserved `workflow.` keys.
    pub(crate) headers: HashMap<String, String>,
    pub(crate) job: &'a JobRecord,
}

impl<'a> ClaimedStep<'a> {
    /// Identify the claimed step from `job`'s headers. Fails, permanently,
    /// for a job without the run id header, with a run id or group id
    /// that is not valid or with a step header that is not a `u32`.
    pub(crate) fn parse(job: &'a JobRecord) -> std::result::Result<Self, Error> {
        let run_id = RunId::new(
            job.headers
                .get(HEADER_RUN_ID)
                .ok_or(Error::MissingHeader(HEADER_RUN_ID))?
                .as_str(),
        )?;
        let step_str = job
            .headers
            .get(HEADER_STEP)
            .ok_or(Error::MissingHeader(HEADER_STEP))?;
        let step_number: u32 = step_str.parse().map_err(|_| Error::InvalidStepHeader {
            header: HEADER_STEP,
            value: step_str.clone(),
        })?;
        let headers = job
            .headers
            .iter()
            .filter(|(k, _)| !k.starts_with(RESERVED_HEADER_PREFIX))
            .map(|(k, v)| (k.clone(), v.clone()))
            .collect();
        Ok(Self {
            run_id,
            step_number,
            membership: Membership::from_headers(&job.headers)?,
            headers,
            job,
        })
    }

    /// The enqueue options of the run's next step: this step's priority,
    /// attempt limit and group membership, so a run's per-step settings
    /// and its membership hold across the step boundary.
    pub(crate) fn next_step_opts(&self) -> StepEnqueueOpts {
        StepEnqueueOpts {
            run_at: None,
            priority: Some(self.job.priority),
            max_attempts: Some(self.job.max_attempts),
            reserved_headers: self.reserved_headers_with_none(),
        }
    }

    /// The reserved headers of the next step: the membership's, plus
    /// `header`.
    pub(crate) fn reserved_headers_with(
        &self,
        header: (&'static str, String),
    ) -> Vec<(&'static str, String)> {
        let mut headers = self.reserved_headers_with_none();
        headers.push(header);
        headers
    }

    fn reserved_headers_with_none(&self) -> Vec<(&'static str, String)> {
        self.membership
            .as_ref()
            .map(Membership::reserved_headers)
            .unwrap_or_default()
    }

    /// The outcome of the run terminating at this step with `status`.
    fn outcome(
        &self,
        status: TerminalStatus,
        result: Option<Vec<u8>>,
        error: Option<String>,
    ) -> RunOutcome {
        RunOutcome {
            run_id: self.run_id.clone(),
            status,
            result,
            error,
            headers: self.headers.clone(),
            final_step: self.step_number,
        }
    }

    /// A `Succeeded` outcome of the run at this step.
    pub(crate) fn succeeded(&self, result: Vec<u8>) -> RunOutcome {
        self.outcome(TerminalStatus::Succeeded, Some(result), None)
    }

    /// A `Failed` outcome of the run at this step.
    pub(crate) fn failed(&self, error: String) -> RunOutcome {
        self.outcome(TerminalStatus::Failed, None, Some(error))
    }

    /// A `Cancelled` outcome of the run at this step; `reason` is
    /// `None` for an external cancellation.
    pub(crate) fn cancelled(&self, reason: Option<String>) -> RunOutcome {
        self.outcome(TerminalStatus::Cancelled, None, reason)
    }
}

impl<R: StepRunner, H: TerminalHook> RuntimeInner<R, H> {
    /// The error a step returns for a failure that terminates its run:
    /// the effects of the `Failed` termination on a [`FailWith`], so the
    /// core applies them only with the dead-lettering settlement. The
    /// run result record is written first; a failed write is logged.
    pub(crate) async fn terminating_failure(
        &self,
        claimed: &ClaimedStep<'_>,
        error: StepError,
        input_hash: [u8; 32],
    ) -> WorkerError {
        let outcome = claimed.failed(error.message.clone());
        let termination = self
            .core
            .termination(&outcome, Some(error.kind), input_hash);
        if let Err(err) = self.core.store_run_result(&outcome, &termination).await {
            warn!(run_id = %claimed.run_id, "failed to write the run result record: {err}");
        }
        let effects = self
            .core
            .terminate_collecting_effects(&outcome, claimed, termination);
        FailWith::new(error.into_worker_error(), effects).into()
    }

    /// Terminate the run of `claimed` with `outcome` from an acking
    /// settlement: write the run result record, then build the
    /// settlement effects. A failed write is a transient error, so the
    /// step is delivered again. `error_kind` classifies a `Failed`
    /// outcome.
    async fn terminate_recorded(
        &self,
        claimed: &ClaimedStep<'_>,
        outcome: RunOutcome,
        input_hash: [u8; 32],
        error_kind: Option<StepErrorKind>,
    ) -> std::result::Result<SettlementEffects, WorkerError> {
        let termination = self.core.termination(&outcome, error_kind, input_hash);
        self.core
            .store_run_result(&outcome, &termination)
            .await
            .map_err(worker_error)?;
        Ok(self
            .core
            .terminate_collecting_effects(&outcome, claimed, termination))
    }

    /// Process a terminal-notification job: decode the committed
    /// outcome and run the configured [`TerminalHook`] as the job's
    /// worker. Effects the hook stages join this job's acknowledgement.
    /// A transient hook error retries the job per the queue's backoff;
    /// a permanent one dead-letters it.
    async fn process_notification(
        &self,
        job: &JobRecord,
    ) -> std::result::Result<SettlementEffects, WorkerError> {
        let outcome: RunOutcome = durable::decode::<DurableRunOutcome>(&job.payload)
            .map_err(|err| {
                warn!(job_id = %job.id, error = %err, "terminal notification has a malformed payload");
                PermanentFailure::new(err.to_string())
            })?
            .into();
        let effects = TerminalEffects::for_delivery();
        let result = self.terminal_hook.on_termination(&outcome, &effects).await;
        let (staged, enqueues) = effects.seal_and_take();
        match result {
            Ok(()) => Ok(SettlementEffects::default()
                .enqueues(enqueues)
                .kv_writes(staged.writes)
                .kv_deletes(staged.deletes.into_iter().collect())),
            Err(err) => Err(err.into_worker_error()),
        }
    }

    /// Process one claimed job of the runtime's queue: a notification
    /// job runs the terminal hook; a step job is run through the
    /// runner, or its stored outcome replayed, and settled with the
    /// effects of the outcome.
    pub(crate) async fn process_step(
        &self,
        job: &JobRecord,
        lease: &LeaseHandle,
    ) -> std::result::Result<SettlementEffects, WorkerError> {
        if job.headers.contains_key(HEADER_TERMINAL) {
            return self.process_notification(job).await;
        }

        let claimed = match ClaimedStep::parse(job) {
            Ok(claimed) => claimed,
            Err(err) => {
                warn!(job_id = %job.id, error = %err, "workflow step has malformed headers");
                return Err(worker_error(err));
            }
        };
        let run_id = &claimed.run_id;
        let step_number = claimed.step_number;

        let (step_signal, signal_kv_deletes) = self
            .core
            .resolve_step_signal(job, run_id, step_number)
            .await
            .map_err(worker_error)?;

        // A cancellation requested before this claim is recorded on the
        // run record; the step is settled as cancelled without running.
        let record = self
            .core
            .view
            .run_record(run_id)
            .await
            .map_err(worker_error)?
            .ok_or_else(|| worker_error(Error::InconsistentRunState(run_id.clone())))?;
        if record.cancel_requested {
            let mut effects = self
                .terminate_recorded(&claimed, claimed.cancelled(None), record.input_hash, None)
                .await?;
            effects.kv_deletes.extend(signal_kv_deletes);
            return Ok(effects);
        }

        // A cancellation requested during the delivery fires the claim's
        // token (`Queue::cancel`, and a re-claim from the job's persisted
        // `cancel_requested`). The runner receives a child, so a runner
        // firing its own token is not treated as an external
        // cancellation below.
        let claim_cancel = lease.cancel_token().clone();

        let effects_handle = EffectsHandle::for_delivery();
        let step = Step {
            delivery: Delivery {
                run_id: run_id.clone(),
                headers: claimed.headers.clone(),
                job_id: job.id.clone(),
                attempts: job.attempts,
                max_attempts: job.max_attempts,
                cancel_token: claim_cancel.child_token(),
                lease: lease.clone(),
                memo: self.core.memo_store.new_memo(run_id, step_number),
                run_memo: self.core.memo_store.new_run_memo(run_id),
                effects: effects_handle.clone(),
                kv: KvReadHandle::for_delivery(self.core.queue.clone()),
            },
            step_number,
            payload: job.payload.clone(),
            signal: step_signal,
        };

        let replayed = if self.core.step_output_replay {
            self.core
                .load_step_output(run_id, step_number, &job.payload)
                .await
                .map_err(worker_error)?
        } else {
            None
        };
        let (outcome, replayed) = match replayed {
            Some((outcome, effects)) => {
                debug!(run_id = %run_id, step_number, "replaying stored step outcome");
                (Ok(outcome), Some(effects))
            }
            None => (self.runner.run_step(&step).await, None),
        };

        // Sealed as soon as the runner has returned: an effect staged
        // through a retained handle clone after this point could not
        // join the settlement, so staging it errors.
        let staged = effects_handle.seal_and_take();
        let replayed_step_output = replayed.is_some();
        let caller_effects = replayed.unwrap_or(staged);

        // A request that lands after this read is recorded on the run
        // record and settles the next step before it runs.
        let external_cancel = claim_cancel.is_cancelled();

        if self.core.step_output_replay
            && !replayed_step_output
            && !external_cancel
            && let Ok(ref outcome) = outcome
        {
            self.core
                .store_step_output(run_id, step_number, &job.payload, outcome, &caller_effects)
                .await
                .map_err(worker_error)?;
        }

        let runner_cancelled = matches!(outcome, Ok(StepOutcome::Cancel { .. }));
        let settled = self
            .settle_outcome(&claimed, outcome, external_cancel, record.input_hash)
            .await;

        // Durable signal entries scheduled for cleanup are deleted with the
        // step's settlement; on retry paths they survive for redelivery.
        settled.map(|mut effects| {
            effects.kv_deletes.extend(signal_kv_deletes);
            // The external-cancel override commits Cancelled in place of
            // the runner's outcome; the staged effects describe that
            // outcome and are discarded with it. A runner-issued Cancel
            // keeps its effects.
            if runner_cancelled || !external_cancel {
                effects.kv_writes.extend(caller_effects.writes);
                effects.kv_deletes.extend(caller_effects.deletes);
            }
            effects
        })
    }

    /// The settlement of `outcome` for `claimed`: the effects of the
    /// run's transition on an acknowledging outcome, or the error that
    /// retries or dead-letters the step.
    ///
    /// Cancellation precedence: a runner-issued [`StepOutcome::Cancel`]
    /// wins and carries its reason on [`RunOutcome::error`]; otherwise
    /// an external cancellation overrides whatever the runner returned,
    /// including a transient retry and a permanent failure, with
    /// `error: None` so consumers can distinguish the two.
    async fn settle_outcome(
        &self,
        claimed: &ClaimedStep<'_>,
        outcome: std::result::Result<StepOutcome, StepError>,
        external_cancel: bool,
        input_hash: [u8; 32],
    ) -> std::result::Result<SettlementEffects, WorkerError> {
        match outcome {
            Ok(StepOutcome::Cancel { reason }) => {
                self.terminate_recorded(claimed, claimed.cancelled(Some(reason)), input_hash, None)
                    .await
            }
            _ if external_cancel => {
                self.terminate_recorded(claimed, claimed.cancelled(None), input_hash, None)
                    .await
            }
            Ok(StepOutcome::Continue { payload, when }) => match when {
                Trigger::Immediate => Ok(self
                    .core
                    .advance(claimed, payload, claimed.next_step_opts())
                    .await),
                Trigger::After(delay) => {
                    let opts = StepEnqueueOpts {
                        run_at: Some(self.core.run_at_after(delay)),
                        ..claimed.next_step_opts()
                    };
                    Ok(self.core.advance(claimed, payload, opts).await)
                }
                Trigger::OnSignal {
                    correlation_key,
                    timeout,
                } => {
                    self.advance_on_signal(claimed, payload, &correlation_key, timeout, input_hash)
                        .await
                }
            },
            Ok(StepOutcome::Succeed { result }) => {
                self.terminate_recorded(claimed, claimed.succeeded(result), input_hash, None)
                    .await
            }
            // A runner verdict: the step ran cleanly and is acknowledged;
            // the run terminates as `Failed` without a dead-letter.
            Ok(StepOutcome::Fail { reason }) => {
                self.terminate_recorded(
                    claimed,
                    claimed.failed(reason),
                    input_hash,
                    Some(StepErrorKind::Permanent),
                )
                .await
            }
            // A permanent error, or a transient one on the last attempt,
            // dead-letters the step and terminates the run.
            Err(err) if err.kind == StepErrorKind::Permanent || claimed.job.is_last_attempt() => {
                Err(self.terminating_failure(claimed, err, input_hash).await)
            }
            Err(err) => Err(err.into_worker_error()),
        }
    }
}