taquba_workflow/terminal.rs
1use std::collections::HashMap;
2use std::future::Future;
3
4use crate::effects::TerminalEffects;
5use crate::keys::RunId;
6use crate::runner::StepError;
7
8/// Terminal state of a workflow run, passed to a [`TerminalHook`].
9#[derive(Debug, Clone, Copy, PartialEq, Eq)]
10pub enum TerminalStatus {
11 /// The runner returned [`crate::StepOutcome::Succeed`].
12 Succeeded,
13 /// One of:
14 /// - the runner returned [`crate::StepOutcome::Fail`] (runner verdict);
15 /// - a step returned [`crate::StepError::permanent`];
16 /// - a step exhausted its transient-retry budget; or
17 /// - a step was dead-lettered outside the runner (its lease expired
18 /// past the attempt limit, crash recovery at open, or a permanent
19 /// runtime error before the runner ran) and the worker's
20 /// reconciliation terminated the run with the queue record's
21 /// last error.
22 Failed,
23 /// The run was cancelled. Either:
24 /// - [`crate::WorkflowRuntime::cancel`] was called for this run; or
25 /// - the runner returned [`crate::StepOutcome::Cancel`].
26 ///
27 /// Like [`Self::Failed`] from `StepOutcome::Fail`, this is a clean
28 /// run-level outcome rather than an infrastructure error: the step is
29 /// acked and no dead-letter is produced.
30 Cancelled,
31}
32
33impl TerminalStatus {
34 /// Canonical lowercase identifier for this status, suitable for HTTP
35 /// headers, structured logs, and other wire-format use. Stable across
36 /// minor releases.
37 pub fn as_str(&self) -> &'static str {
38 match self {
39 TerminalStatus::Succeeded => "succeeded",
40 TerminalStatus::Failed => "failed",
41 TerminalStatus::Cancelled => "cancelled",
42 }
43 }
44}
45
46impl std::fmt::Display for TerminalStatus {
47 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
48 f.write_str(self.as_str())
49 }
50}
51
52/// Information passed to a [`TerminalHook`] when a run reaches a terminal
53/// state.
54#[derive(Debug, Clone)]
55pub struct RunOutcome {
56 /// The run's identifier.
57 pub run_id: RunId,
58 /// Whether the run completed successfully or failed.
59 pub status: TerminalStatus,
60 /// Set when `status == Succeeded`: the bytes the runner returned via
61 /// [`crate::StepOutcome::Succeed`].
62 pub result: Option<Vec<u8>>,
63 /// - When `status == Failed`: the human-readable reason recorded on
64 /// the terminal step's `last_error`.
65 /// - When `status == Cancelled`: `Some(reason)` if the runner
66 /// returned [`crate::StepOutcome::Cancel`], or `None` if
67 /// cancellation came from [`crate::WorkflowRuntime::cancel`]
68 /// (which takes no reason at the API level).
69 /// - When `status == Succeeded`: always `None`.
70 pub error: Option<String>,
71 /// Submitter-supplied metadata, threaded through from
72 /// [`crate::RunOptions::headers`].
73 pub headers: HashMap<String, String>,
74 /// Step number of the step that produced the terminal outcome (zero-based).
75 pub final_step: u32,
76}
77
78/// User-implemented hook processing a run's termination.
79///
80/// Termination is delivered as a queue job: the settlement that commits
81/// a run's terminal outcome atomically enqueues a **notification job**
82/// on the same queue, and the hook runs as that job's worker. The
83/// consequences:
84///
85/// - The hook observes only outcomes that committed. A settlement that
86/// loses its claim loses its notification with it, so a redelivered
87/// terminal step notifies only the outcome it actually commits.
88/// - Delivery is at-least-once: a crash after the hook ran but before
89/// the notification job was acknowledged re-delivers it, so
90/// implementations must be idempotent.
91/// - The notification job takes the terminal step's priority and
92/// `max_attempts`. A transient error ([`StepError::transient`]) retries the
93/// notification job per the queue's backoff up to `max_attempts`. A permanent
94/// error dead-letters it, where [`taquba::QueueView::dead_jobs`] finds it.
95/// - Effects staged on the [`TerminalEffects`] handle are applied in
96/// the same transaction as the notification's acknowledgement when
97/// the hook returns `Ok`.
98///
99/// Runs terminated without an acknowledging settlement (an external
100/// cancellation of a pending step, a step that dead-letters) enqueue
101/// the notification job in the transaction of that transition, so it
102/// is created exactly once on every worker and cancellation path. A
103/// step the reaper dead-letters after its lease expires, or one
104/// dead-lettered during crash recovery at open, is reconciled by the
105/// worker, which terminates the run as failed and enqueues the
106/// notification in one transaction.
107pub trait TerminalHook: Send + Sync {
108 /// Process the termination of one run. `outcome` is the committed
109 /// terminal state; effects staged on `effects` commit with this
110 /// notification's acknowledgement.
111 fn on_termination(
112 &self,
113 outcome: &RunOutcome,
114 effects: &TerminalEffects,
115 ) -> impl Future<Output = std::result::Result<(), StepError>> + Send;
116
117 /// Whether a notification job should be enqueued for `outcome`.
118 /// Consulted when the run terminates; returning `false` skips the
119 /// notification entirely, so [`Self::on_termination`] is never
120 /// called for that run. Defaults to `true`.
121 fn observes(&self, outcome: &RunOutcome) -> bool {
122 let _ = outcome;
123 true
124 }
125}
126
127/// A no-op terminal hook. Declares itself unobservant, so runs
128/// terminate without enqueueing a notification job.
129#[derive(Debug, Default, Clone, Copy)]
130pub struct NoopTerminalHook;
131
132impl TerminalHook for NoopTerminalHook {
133 async fn on_termination(
134 &self,
135 _outcome: &RunOutcome,
136 _effects: &TerminalEffects,
137 ) -> std::result::Result<(), StepError> {
138 Ok(())
139 }
140
141 fn observes(&self, _outcome: &RunOutcome) -> bool {
142 false
143 }
144}
145
146#[cfg(feature = "webhooks")]
147mod webhook {
148 use super::{RunOutcome, StepError, TerminalEffects, TerminalHook, TerminalStatus};
149 use std::time::Duration;
150 use taquba_webhooks::{WebhookRequest, webhook_enqueue_request};
151
152 /// Terminal hook that delivers an HTTP webhook via `taquba-webhooks`
153 /// when a run terminates.
154 ///
155 /// The hook reads the target URL from the run's submission headers
156 /// under [`Self::URL_HEADER`] (default `"callback_url"`); runs
157 /// without that header enqueue no notification at all. The default
158 /// key intentionally avoids the reserved `workflow.*` prefix so
159 /// submitters can set it directly via [`crate::RunOptions::headers`].
160 ///
161 /// The webhook enqueue is staged as a notification effect, so the
162 /// delivery job is created exactly once, atomically with the
163 /// notification's acknowledgement.
164 ///
165 /// The webhook body is the raw `result` bytes for succeeded runs, and
166 /// the UTF-8 error message for failed runs. The run identifier and
167 /// terminal status are passed in the `Workflow-Run-Id` and
168 /// `Workflow-Run-Status` HTTP headers respectively.
169 pub struct WebhookTerminalHook {
170 target_queue: String,
171 url_header: String,
172 timeout: Option<Duration>,
173 }
174
175 impl WebhookTerminalHook {
176 /// Default header key the hook looks for on each [`RunOutcome`].
177 /// Deliberately outside the reserved `workflow.*` prefix so submitters
178 /// can set it on [`crate::RunOptions::headers`] without being
179 /// rejected.
180 pub const URL_HEADER: &'static str = "callback_url";
181
182 /// Build a hook that enqueues webhook deliveries onto
183 /// `target_queue`. The submitter sets a callback URL per run via
184 /// the [`Self::URL_HEADER`] header on [`crate::RunOptions::headers`].
185 pub fn new(target_queue: impl Into<String>) -> Self {
186 Self {
187 target_queue: target_queue.into(),
188 url_header: Self::URL_HEADER.to_string(),
189 timeout: None,
190 }
191 }
192
193 /// Override the header key the hook reads. Defaults to
194 /// [`Self::URL_HEADER`].
195 pub fn with_url_header(mut self, header: impl Into<String>) -> Self {
196 self.url_header = header.into();
197 self
198 }
199
200 /// Set a per-delivery timeout passed through to the webhook worker.
201 pub fn with_timeout(mut self, timeout: Duration) -> Self {
202 self.timeout = Some(timeout);
203 self
204 }
205 }
206
207 impl TerminalHook for WebhookTerminalHook {
208 async fn on_termination(
209 &self,
210 outcome: &RunOutcome,
211 effects: &TerminalEffects,
212 ) -> std::result::Result<(), StepError> {
213 let Some(url) = outcome.headers.get(&self.url_header) else {
214 return Ok(());
215 };
216 let mut req = WebhookRequest::new(url)
217 .header("Workflow-Run-Id", outcome.run_id.as_str())
218 .header("Workflow-Run-Status", outcome.status.as_str());
219 if let Some(t) = self.timeout {
220 req = req.timeout(t);
221 }
222 let body = match outcome.status {
223 TerminalStatus::Succeeded => outcome.result.clone().unwrap_or_default(),
224 TerminalStatus::Failed | TerminalStatus::Cancelled => {
225 outcome.error.clone().unwrap_or_default().into_bytes()
226 }
227 };
228 let request = webhook_enqueue_request(&self.target_queue, req, body);
229 effects
230 .enqueue(request)
231 .map_err(|e| StepError::permanent(e.to_string()))?;
232 Ok(())
233 }
234
235 fn observes(&self, outcome: &RunOutcome) -> bool {
236 outcome.headers.contains_key(&self.url_header)
237 }
238 }
239}
240
241#[cfg(feature = "webhooks")]
242pub use webhook::WebhookTerminalHook;