Skip to main content

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}