ironflow_engine/pause.rs
1//! Operator pause of a run, or of a whole workflow.
2//!
3//! [`Engine::pause_run`] holds a root run, with every sub-workflow run below
4//! it, in `Paused` until [`Engine::resume_paused_run`] puts each one back in
5//! the state it was paused from, or [`Engine::cancel_run`] stops it. A run
6//! executing when it is paused has its step in flight interrupted with
7//! [`STEP_INTERRUPTED_ERROR`](ironflow_store::store::STEP_INTERRUPTED_ERROR),
8//! no further step is started, and the resume replays the run: finished steps
9//! are skipped and the interrupted step is executed again.
10//!
11//! [`Engine::pause_workflow`] holds the queued runs of a workflow instead:
12//! they are created as usual but no worker picks them up until
13//! [`Engine::resume_workflow`].
14
15use std::sync::Arc;
16
17use chrono::Utc;
18use tracing::{debug, info};
19use uuid::Uuid;
20
21use ironflow_store::error::StoreError;
22use ironflow_store::models::{Run, RunStatus, RunUpdate, WorkflowPause};
23
24use crate::engine::{Engine, ExecutionMode, chain_root};
25use crate::error::EngineError;
26use crate::notify::{Event, RunStatusChangedEvent};
27
28/// Outcome of [`Engine::pause_run`].
29///
30/// # Examples
31///
32/// ```no_run
33/// use std::sync::Arc;
34/// use ironflow_engine::engine::Engine;
35/// use ironflow_engine::error::EngineError;
36/// use uuid::Uuid;
37///
38/// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
39/// let pause = engine.pause_run(run_id).await?;
40/// println!(
41/// "run {} paused with {} sub-runs",
42/// pause.run.id,
43/// pause.paused_descendants.len()
44/// );
45/// # Ok(())
46/// # }
47/// ```
48#[derive(Debug, Clone)]
49pub struct RunPause {
50 /// The paused run, as stored after the pause.
51 pub run: Run,
52 /// The sub-workflow runs below it that this call paused, oldest first.
53 pub paused_descendants: Vec<Uuid>,
54}
55
56/// Outcome of [`Engine::resume_paused_run`].
57///
58/// # Examples
59///
60/// ```no_run
61/// use std::sync::Arc;
62/// use ironflow_engine::engine::Engine;
63/// use ironflow_engine::error::EngineError;
64/// use uuid::Uuid;
65///
66/// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
67/// let resume = engine.resume_paused_run(run_id).await?;
68/// println!("run {} is {}", resume.run.id, resume.run.status.state);
69/// # Ok(())
70/// # }
71/// ```
72#[derive(Debug, Clone)]
73pub struct RunResume {
74 /// The resumed run, as stored after the resume.
75 pub run: Run,
76 /// The sub-workflow runs below it that this call resumed, oldest first.
77 pub resumed_descendants: Vec<Uuid>,
78}
79
80impl Engine {
81 /// Pause a root run and every active sub-workflow run below it.
82 ///
83 /// Each run moves to `Paused` and records the state it was paused from in
84 /// [`Run::resume_status`]. A run that was executing has its running steps
85 /// marked `Failed` with
86 /// [`STEP_INTERRUPTED_ERROR`](ironflow_store::store::STEP_INTERRUPTED_ERROR),
87 /// so the resume executes them again. A run held by a worker loses its lease: the worker drops
88 /// the execution at its next renewal, and the reaper never touches a
89 /// paused run. A run executing in-process stops before its next step.
90 /// [`Event::RunStatusChanged`] is published for every run paused.
91 ///
92 /// While the run is paused, an approval, a human input or a signal it
93 /// waits for can still be resolved: the decision is recorded and only
94 /// changes the state the run resumes to.
95 ///
96 /// # Errors
97 ///
98 /// - [`EngineError::Store`] with [`StoreError::RunNotFound`] for an
99 /// unknown run.
100 /// - [`EngineError::ChildRunNotPausable`] for a sub-workflow run: pause
101 /// its root run instead.
102 /// - [`EngineError::Store`] with [`StoreError::InvalidTransition`] for a
103 /// run already paused or finished.
104 /// - [`EngineError::Store`] when the pause cannot be persisted.
105 ///
106 /// # Examples
107 ///
108 /// ```no_run
109 /// use std::sync::Arc;
110 /// use ironflow_engine::engine::Engine;
111 /// use ironflow_engine::error::EngineError;
112 /// use ironflow_store::models::RunStatus;
113 /// use uuid::Uuid;
114 ///
115 /// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
116 /// let pause = engine.pause_run(run_id).await?;
117 /// assert_eq!(pause.run.status.state, RunStatus::Paused);
118 /// # Ok(())
119 /// # }
120 /// ```
121 pub async fn pause_run(&self, run_id: Uuid) -> Result<RunPause, EngineError> {
122 let run = self.load_pausable_root(run_id).await?;
123 let from = run.status.state;
124 if !from.can_transition_to(&RunStatus::Paused) {
125 return Err(EngineError::Store(StoreError::InvalidTransition {
126 from,
127 to: RunStatus::Paused,
128 }));
129 }
130
131 // The root first: an execution in flight checks its own run before
132 // every step, so it stops as soon as possible.
133 self.store()
134 .update_run_status(run_id, RunStatus::Paused)
135 .await?;
136 if from == RunStatus::Running {
137 self.interrupt_running_steps(run_id).await?;
138 }
139 self.publish_transition(&run, RunStatus::Paused);
140 info!(run_id = %run_id, from = %from, "run paused");
141
142 let mut paused_descendants = Vec::new();
143 for descendant in self.store().list_active_descendants(run_id).await? {
144 if !descendant
145 .status
146 .state
147 .can_transition_to(&RunStatus::Paused)
148 {
149 continue;
150 }
151 match self
152 .store()
153 .update_run_status(descendant.id, RunStatus::Paused)
154 .await
155 {
156 Ok(()) => {}
157 // Finished or paused since it was listed: nothing to pause.
158 Err(StoreError::InvalidTransition { .. }) => {
159 debug!(
160 run_id = %descendant.id,
161 "descendant run no longer pausable, skipped"
162 );
163 continue;
164 }
165 Err(err) => return Err(err.into()),
166 }
167 if descendant.status.state == RunStatus::Running {
168 self.interrupt_running_steps(descendant.id).await?;
169 }
170 self.publish_transition(&descendant, RunStatus::Paused);
171 paused_descendants.push(descendant.id);
172 }
173
174 if !paused_descendants.is_empty() {
175 info!(
176 run_id = %run_id,
177 count = paused_descendants.len(),
178 "descendant runs paused"
179 );
180 }
181
182 Ok(RunPause {
183 run: self.load_run(run_id).await?,
184 paused_descendants,
185 })
186 }
187
188 /// Resume a paused root run and the sub-workflow runs paused with it.
189 ///
190 /// Each run goes back to the state recorded in [`Run::resume_status`]:
191 /// a run paused while waiting (`Pending`, `Retrying`, `Sleeping`,
192 /// `AwaitingApproval`) waits again, and a run whose approval, human input
193 /// or signal was resolved during the pause is queued. A root that was
194 /// executing is queued with its interrupted steps, which are executed
195 /// again; finished steps are replayed. A sleeping root whose deadline
196 /// passed during the pause is queued at once; a sleeping sub-workflow run
197 /// in the same case is woken by the [`RunWaker`](crate::wake::RunWaker)
198 /// on its next tick.
199 ///
200 /// Under [`ExecutionMode::Local`] a queued root is resumed in a background
201 /// task, and so is a queued sub-workflow run whose root still waits.
202 /// [`Event::RunStatusChanged`] is published for every run resumed.
203 ///
204 /// # Errors
205 ///
206 /// - [`EngineError::Store`] with [`StoreError::RunNotFound`] for an
207 /// unknown run.
208 /// - [`EngineError::ChildRunNotPausable`] for a sub-workflow run: resume
209 /// its root run instead.
210 /// - [`EngineError::Store`] with [`StoreError::InvalidTransition`] for a
211 /// run that is not paused.
212 /// - [`EngineError::Store`] when the resume cannot be persisted.
213 ///
214 /// # Examples
215 ///
216 /// ```no_run
217 /// use std::sync::Arc;
218 /// use ironflow_engine::engine::Engine;
219 /// use ironflow_engine::error::EngineError;
220 /// use ironflow_store::models::RunStatus;
221 /// use uuid::Uuid;
222 ///
223 /// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
224 /// let resume = engine.resume_paused_run(run_id).await?;
225 /// assert_ne!(resume.run.status.state, RunStatus::Paused);
226 /// # Ok(())
227 /// # }
228 /// ```
229 pub async fn resume_paused_run(
230 self: &Arc<Self>,
231 run_id: Uuid,
232 ) -> Result<RunResume, EngineError> {
233 let run = self.load_pausable_root(run_id).await?;
234 if run.status.state != RunStatus::Paused {
235 return Err(EngineError::Store(StoreError::InvalidTransition {
236 from: run.status.state,
237 to: RunStatus::Pending,
238 }));
239 }
240
241 // The descendants first: the root re-enters them as soon as it runs.
242 let mut resumed_descendants = Vec::new();
243 let mut queued_descendants = Vec::new();
244 for descendant in self.store().list_active_descendants(run_id).await? {
245 if descendant.status.state != RunStatus::Paused {
246 continue;
247 }
248 // A child that was executing stays `Running`: its root's replay
249 // re-enters it and executes its interrupted steps again.
250 let target = descendant.resume_status.unwrap_or(RunStatus::Pending);
251 self.store()
252 .update_run_status(descendant.id, target)
253 .await?;
254 self.publish_transition(&descendant, target);
255 if target == RunStatus::Pending {
256 queued_descendants.push(descendant.id);
257 }
258 resumed_descendants.push(descendant.id);
259 }
260
261 let target = match run.resume_status {
262 // Stopped by the pause in the middle of its execution: queued
263 // again, like a run whose worker lost its lease.
264 Some(RunStatus::Running) | None => {
265 self.interrupt_running_steps(run_id).await?;
266 RunStatus::Pending
267 }
268 // The deadline passed during the pause: due now, like a run the
269 // waker would have claimed. A past `scheduled_at` does not hold
270 // back the pick.
271 Some(RunStatus::Sleeping) if run.scheduled_at.is_some_and(|at| at <= Utc::now()) => {
272 RunStatus::Pending
273 }
274 Some(status) => status,
275 };
276 self.store().update_run_status(run_id, target).await?;
277 self.publish_transition(&run, target);
278 info!(run_id = %run_id, to = %target, "run resumed");
279
280 if self.execution_mode() == ExecutionMode::Local {
281 if target == RunStatus::Pending {
282 self.spawn_local_resume(run_id);
283 } else {
284 // The root waits on its chain: the queued child resumes it.
285 for child_id in queued_descendants {
286 self.spawn_local_resume(child_id);
287 }
288 }
289 }
290
291 Ok(RunResume {
292 run: self.load_run(run_id).await?,
293 resumed_descendants,
294 })
295 }
296
297 /// Pause a registered workflow: its queued runs are no longer picked up.
298 ///
299 /// Runs keep being created; workers skip them until
300 /// [`resume_workflow`](Self::resume_workflow). Runs already executing
301 /// are not affected: pause them with [`pause_run`](Self::pause_run).
302 /// Pausing a paused workflow returns the pause already recorded.
303 ///
304 /// # Errors
305 ///
306 /// Returns [`EngineError::InvalidWorkflow`] when no handler is registered
307 /// under `workflow_name`, and [`EngineError::Store`] when the pause cannot
308 /// be persisted.
309 ///
310 /// # Examples
311 ///
312 /// ```no_run
313 /// use ironflow_engine::engine::Engine;
314 /// use ironflow_engine::error::EngineError;
315 ///
316 /// # async fn example(engine: &Engine) -> Result<(), EngineError> {
317 /// let pause = engine.pause_workflow("deploy", None).await?;
318 /// println!("deploy paused at {}", pause.paused_at);
319 /// # Ok(())
320 /// # }
321 /// ```
322 pub async fn pause_workflow(
323 &self,
324 workflow_name: &str,
325 paused_by: Option<Uuid>,
326 ) -> Result<WorkflowPause, EngineError> {
327 self.require_handler(workflow_name)?;
328 let pause = self
329 .store()
330 .pause_workflow(workflow_name, paused_by)
331 .await?;
332 info!(workflow = %workflow_name, "workflow paused");
333 Ok(pause)
334 }
335
336 /// Resume a paused workflow so its queued runs are picked up again.
337 ///
338 /// Returns `true` when the workflow was paused, `false` when it was not:
339 /// resuming is idempotent.
340 ///
341 /// # Errors
342 ///
343 /// Returns [`EngineError::InvalidWorkflow`] when no handler is registered
344 /// under `workflow_name`, and [`EngineError::Store`] when the resume
345 /// cannot be persisted.
346 ///
347 /// # Examples
348 ///
349 /// ```no_run
350 /// use ironflow_engine::engine::Engine;
351 /// use ironflow_engine::error::EngineError;
352 ///
353 /// # async fn example(engine: &Engine) -> Result<(), EngineError> {
354 /// if engine.resume_workflow("deploy").await? {
355 /// println!("deploy resumed");
356 /// }
357 /// # Ok(())
358 /// # }
359 /// ```
360 pub async fn resume_workflow(&self, workflow_name: &str) -> Result<bool, EngineError> {
361 self.require_handler(workflow_name)?;
362 let was_paused = self.store().resume_workflow(workflow_name).await?;
363 if was_paused {
364 info!(workflow = %workflow_name, "workflow resumed");
365 }
366 Ok(was_paused)
367 }
368
369 fn require_handler(&self, workflow_name: &str) -> Result<(), EngineError> {
370 match self.get_handler(workflow_name) {
371 Some(_) => Ok(()),
372 None => Err(EngineError::InvalidWorkflow(format!(
373 "no handler registered for workflow '{workflow_name}'"
374 ))),
375 }
376 }
377
378 /// Make the paused root of `run`'s chain resume to `Pending` when it
379 /// waits suspended with its child, so that its replay observes what
380 /// happened to the child (cancelled, rejected) once an operator resumes
381 /// it.
382 pub(crate) async fn requeue_paused_root(&self, run: &Run) -> Result<(), EngineError> {
383 let Some(root_id) = chain_root(run) else {
384 return Ok(());
385 };
386 let root = self.load_run(root_id).await?;
387 let suspended = matches!(
388 root.resume_status,
389 Some(RunStatus::AwaitingApproval | RunStatus::Sleeping)
390 );
391 if root.status.state != RunStatus::Paused || !suspended {
392 return Ok(());
393 }
394
395 let update = RunUpdate {
396 resume_status: Some(RunStatus::Pending),
397 ..RunUpdate::default()
398 };
399 self.store().update_run(root_id, update).await?;
400 info!(run_id = %run.id, root_run_id = %root_id, "paused root run will resume to observe its child");
401 Ok(())
402 }
403
404 /// Load `run_id`, refusing a sub-workflow run.
405 async fn load_pausable_root(&self, run_id: Uuid) -> Result<Run, EngineError> {
406 let run = self.load_run(run_id).await?;
407 match chain_root(&run) {
408 Some(root_run_id) => Err(EngineError::ChildRunNotPausable {
409 run_id,
410 root_run_id,
411 }),
412 None => Ok(run),
413 }
414 }
415
416 /// Publish the move of `run` (as loaded before the update) to `to`.
417 fn publish_transition(&self, run: &Run, to: RunStatus) {
418 self.event_publisher()
419 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
420 run_id: run.id,
421 workflow_name: run.workflow_name.clone(),
422 from: run.status.state,
423 to,
424 error: None,
425 cost_usd: run.cost_usd,
426 duration_ms: run.duration_ms,
427 labels: run.labels.clone(),
428 at: Utc::now(),
429 }));
430 }
431}