Skip to main content

ironflow_engine/
cancel.rs

1//! Cancellation of a run and of the sub-workflow runs below it.
2//!
3//! A child run executes inside its parent's execution: once the parent stops
4//! (cancelled, failed, retried, abandoned by its worker), nothing drives the
5//! child any more. [`Engine::cancel_descendants`] closes those children so
6//! none stays non-terminal forever, holding its concurrency key.
7//! [`Engine::cancel_run`] cancels a run with its descendants, and wakes the
8//! root of a suspended chain whose child it cancelled so the root observes it.
9
10use std::sync::Arc;
11
12use chrono::Utc;
13use tokio::spawn;
14use tracing::{error, info, warn};
15use uuid::Uuid;
16
17use ironflow_store::error::StoreError;
18use ironflow_store::models::{Run, RunStatus, RunUpdate};
19
20use crate::engine::{Engine, ExecutionMode, chain_root};
21use crate::error::EngineError;
22use crate::notify::{Event, RunStatusChangedEvent};
23
24/// Error recorded on the steps left open by a cancelled run.
25///
26/// # Examples
27///
28/// ```
29/// use ironflow_engine::cancel::RUN_CANCELLED_ERROR;
30///
31/// assert_eq!(RUN_CANCELLED_ERROR, "run cancelled");
32/// ```
33pub const RUN_CANCELLED_ERROR: &str = "run cancelled";
34
35/// Outcome of [`Engine::cancel_run`].
36///
37/// # Examples
38///
39/// ```no_run
40/// use std::sync::Arc;
41/// use ironflow_engine::engine::Engine;
42/// use ironflow_engine::error::EngineError;
43/// use uuid::Uuid;
44///
45/// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
46/// let cancellation = engine.cancel_run(run_id).await?;
47/// println!(
48///     "run {} cancelled with {} sub-runs",
49///     cancellation.run.id,
50///     cancellation.cancelled_descendants.len()
51/// );
52/// # Ok(())
53/// # }
54/// ```
55#[derive(Debug, Clone)]
56pub struct RunCancellation {
57    /// The cancelled run, as stored after the cancellation.
58    pub run: Run,
59    /// The sub-workflow runs below it that this call cancelled, oldest first.
60    pub cancelled_descendants: Vec<Uuid>,
61}
62
63impl Engine {
64    /// Cancel a run and every non-terminal sub-workflow run below it.
65    ///
66    /// The run moves to `Cancelled`, its open steps are closed with
67    /// [`RUN_CANCELLED_ERROR`] and [`Event::RunStatusChanged`] is published.
68    /// Its descendants are then cancelled by
69    /// [`cancel_descendants`](Self::cancel_descendants), which releases their
70    /// concurrency keys.
71    ///
72    /// When the run is a child whose root waits suspended with it
73    /// (`AwaitingApproval` or `Sleeping`), the root is woken so its open
74    /// `Workflow` step observes the cancellation: under
75    /// [`ExecutionMode::Local`] it resumes in a background task, under
76    /// [`ExecutionMode::Workers`] it goes back to `Pending` for a worker. The
77    /// parent then fails, unless its step tolerates the failure with
78    /// `allow_failure`.
79    ///
80    /// Cancelling a run already `Cancelled` is accepted: it publishes nothing
81    /// for the run itself and cancels whatever descendant is still active.
82    ///
83    /// # Errors
84    ///
85    /// Returns [`EngineError::Store`] with [`StoreError::RunNotFound`] for an
86    /// unknown run, with [`StoreError::InvalidTransition`] for a run that
87    /// already finished otherwise (`Completed`, `Failed`, `Warning`), and with
88    /// the store error when the cancellation cannot be persisted.
89    ///
90    /// # Examples
91    ///
92    /// ```no_run
93    /// use std::sync::Arc;
94    /// use ironflow_engine::engine::Engine;
95    /// use ironflow_engine::error::EngineError;
96    /// use uuid::Uuid;
97    ///
98    /// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
99    /// let cancellation = engine.cancel_run(run_id).await?;
100    /// assert!(cancellation.run.status.state.is_terminal());
101    /// # Ok(())
102    /// # }
103    /// ```
104    pub async fn cancel_run(
105        self: &Arc<Self>,
106        run_id: Uuid,
107    ) -> Result<RunCancellation, EngineError> {
108        let run = self.load_run(run_id).await?;
109        let from = run.status.state;
110        if !from.can_transition_to(&RunStatus::Cancelled) {
111            return Err(EngineError::Store(StoreError::InvalidTransition {
112                from,
113                to: RunStatus::Cancelled,
114            }));
115        }
116
117        if from != RunStatus::Cancelled {
118            self.store()
119                .update_run(
120                    run_id,
121                    RunUpdate {
122                        status: Some(RunStatus::Cancelled),
123                        completed_at: Some(Utc::now()),
124                        ..RunUpdate::default()
125                    },
126                )
127                .await?;
128            self.publish_cancelled(&run, None);
129            info!(run_id = %run_id, from = %from, "run cancelled");
130        }
131        self.fail_orphaned_steps(run_id, RUN_CANCELLED_ERROR)
132            .await?;
133
134        let reason = format!("ancestor run {run_id} cancelled");
135        let cancelled_descendants = self.cancel_descendants(run_id, &reason).await?;
136
137        self.wake_root_of_cancelled_child(&run).await;
138
139        Ok(RunCancellation {
140            run: self.load_run(run_id).await?,
141            cancelled_descendants,
142        })
143    }
144
145    /// Cancel every non-terminal sub-workflow run below `run_id`, at any depth.
146    ///
147    /// Each descendant moves to `Cancelled` with `reason` as its error, its
148    /// open steps are closed with the same reason, and
149    /// [`Event::RunStatusChanged`] is published for it. A terminal descendant
150    /// releases its concurrency key. A descendant that finished on its own in
151    /// the meantime is left as it is. Returns the ids of the runs cancelled,
152    /// oldest first: an empty list when nothing was left to cancel.
153    ///
154    /// Called when a run stops for good or ends an attempt
155    /// ([`fail_or_schedule_retry`](Self::fail_or_schedule_retry)), when the
156    /// reaper fails a run whose worker died, and by
157    /// [`cancel_run`](Self::cancel_run).
158    ///
159    /// # Errors
160    ///
161    /// Returns [`EngineError::Store`] when the descendants cannot be listed or
162    /// one of them cannot be updated.
163    ///
164    /// # Examples
165    ///
166    /// ```no_run
167    /// use ironflow_engine::engine::Engine;
168    /// use ironflow_engine::error::EngineError;
169    /// use uuid::Uuid;
170    ///
171    /// # async fn example(engine: &Engine, run_id: Uuid) -> Result<(), EngineError> {
172    /// let cancelled = engine
173    ///     .cancel_descendants(run_id, "parent run stopped")
174    ///     .await?;
175    /// println!("{} sub-runs cancelled", cancelled.len());
176    /// # Ok(())
177    /// # }
178    /// ```
179    pub async fn cancel_descendants(
180        &self,
181        run_id: Uuid,
182        reason: &str,
183    ) -> Result<Vec<Uuid>, EngineError> {
184        let descendants = self.store().list_active_descendants(run_id).await?;
185        let mut cancelled = Vec::with_capacity(descendants.len());
186
187        for descendant in descendants {
188            let update = RunUpdate {
189                status: Some(RunStatus::Cancelled),
190                error: Some(reason.to_string()),
191                completed_at: Some(Utc::now()),
192                ..RunUpdate::default()
193            };
194            match self.store().update_run(descendant.id, update).await {
195                Ok(()) => {}
196                // Finished on its own since it was listed: nothing to cancel.
197                Err(StoreError::InvalidTransition { from, .. }) => {
198                    info!(run_id = %descendant.id, status = %from, "descendant already finished");
199                    continue;
200                }
201                Err(err) => return Err(err.into()),
202            }
203            self.fail_orphaned_steps(descendant.id, reason).await?;
204            self.publish_cancelled(&descendant, Some(reason));
205            cancelled.push(descendant.id);
206        }
207
208        if !cancelled.is_empty() {
209            info!(
210                run_id = %run_id,
211                count = cancelled.len(),
212                reason = %reason,
213                "descendant runs cancelled"
214            );
215        }
216        Ok(cancelled)
217    }
218
219    /// Close the descendants of a run whose attempt just ended, logging any
220    /// failure instead of returning it: the run's own status is already
221    /// persisted and must not be reported as unsaved.
222    pub(crate) async fn cancel_descendants_of_stopped_run(&self, run_id: Uuid, error: &str) {
223        let reason = format!("parent run {run_id} stopped: {error}");
224        if let Err(err) = self.cancel_descendants(run_id, &reason).await {
225            error!(run_id = %run_id, error = %err, "failed to cancel the descendants of a stopped run");
226        }
227    }
228
229    /// Wake the root of `run`'s chain when it waits suspended with it.
230    ///
231    /// Best effort: the cancellation is already persisted, so a failure is
232    /// logged and the root is left for an operator.
233    async fn wake_root_of_cancelled_child(self: &Arc<Self>, run: &Run) {
234        let Some(root_id) = chain_root(run) else {
235            return;
236        };
237        let root = match self.store().get_run(root_id).await {
238            Ok(Some(root)) => root,
239            Ok(None) => return,
240            Err(err) => {
241                warn!(run_id = %run.id, root_run_id = %root_id, error = %err, "cannot read the root of a cancelled child");
242                return;
243            }
244        };
245
246        // A paused root is not woken: it resumes as an active run, so its
247        // replay observes the cancelled child once an operator resumes it.
248        if root.status.state == RunStatus::Paused {
249            if let Err(err) = self.requeue_paused_root(run).await {
250                warn!(run_id = %run.id, root_run_id = %root_id, error = %err, "cannot requeue the paused root of a cancelled child");
251            }
252            return;
253        }
254
255        let woken = match (root.status.state, self.execution_mode()) {
256            (RunStatus::AwaitingApproval | RunStatus::Sleeping, ExecutionMode::Workers) => {
257                self.store()
258                    .update_run_status(root_id, RunStatus::Pending)
259                    .await
260            }
261            (RunStatus::AwaitingApproval, ExecutionMode::Local) => {
262                self.store()
263                    .update_run_status(root_id, RunStatus::Running)
264                    .await
265            }
266            (RunStatus::Sleeping, ExecutionMode::Local) => {
267                match self
268                    .store()
269                    .update_run_status(root_id, RunStatus::Pending)
270                    .await
271                {
272                    Ok(()) => {
273                        self.store()
274                            .update_run_status(root_id, RunStatus::Running)
275                            .await
276                    }
277                    Err(err) => Err(err),
278                }
279            }
280            // Running (the child runs inline), queued or finished: nothing waits.
281            _ => return,
282        };
283        if let Err(err) = woken {
284            warn!(run_id = %run.id, root_run_id = %root_id, error = %err, "cannot wake the root of a cancelled child");
285            return;
286        }
287
288        info!(run_id = %run.id, root_run_id = %root_id, "root run woken to observe its cancelled child");
289        if self.execution_mode() == ExecutionMode::Local {
290            let engine = Arc::clone(self);
291            spawn(async move {
292                if let Err(err) = engine.resume_run(root_id).await {
293                    error!(root_run_id = %root_id, error = %err, "root run stopped after its child was cancelled");
294                }
295            });
296        }
297    }
298
299    /// Publish the move of `run` (as loaded before the update) to `Cancelled`.
300    fn publish_cancelled(&self, run: &Run, error: Option<&str>) {
301        self.event_publisher()
302            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
303                run_id: run.id,
304                workflow_name: run.workflow_name.clone(),
305                from: run.status.state,
306                to: RunStatus::Cancelled,
307                error: error.map(str::to_string),
308                cost_usd: run.cost_usd,
309                duration_ms: run.duration_ms,
310                labels: run.labels.clone(),
311                at: Utc::now(),
312            }));
313    }
314
315    pub(crate) async fn load_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
316        self.store()
317            .get_run(run_id)
318            .await?
319            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))
320    }
321}