Skip to main content

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;