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