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}