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`].
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. The SLA deadline of an approval
95    /// gate keeps running: the [`ApprovalEscalator`](crate::escalation::ApprovalEscalator)
96    /// applies its policy during the pause, with the same effect as a human
97    /// decision.
98    ///
99    /// # Errors
100    ///
101    /// - [`EngineError::Store`] with [`StoreError::RunNotFound`] for an
102    ///   unknown run.
103    /// - [`EngineError::ChildRunNotPausable`] for a sub-workflow run: pause
104    ///   its root run instead.
105    /// - [`EngineError::Store`] with [`StoreError::InvalidTransition`] for a
106    ///   run already paused or finished.
107    /// - [`EngineError::Store`] when the pause cannot be persisted.
108    ///
109    /// # Examples
110    ///
111    /// ```no_run
112    /// use std::sync::Arc;
113    /// use ironflow_engine::engine::Engine;
114    /// use ironflow_engine::error::EngineError;
115    /// use ironflow_store::models::RunStatus;
116    /// use uuid::Uuid;
117    ///
118    /// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
119    /// let pause = engine.pause_run(run_id).await?;
120    /// assert_eq!(pause.run.status.state, RunStatus::Paused);
121    /// # Ok(())
122    /// # }
123    /// ```
124    pub async fn pause_run(&self, run_id: Uuid) -> Result<RunPause, EngineError> {
125        let run = self.load_pausable_root(run_id).await?;
126        let from = run.status.state;
127        if !from.can_transition_to(&RunStatus::Paused) {
128            return Err(EngineError::Store(StoreError::InvalidTransition {
129                from,
130                to: RunStatus::Paused,
131            }));
132        }
133
134        // The root first: an execution in flight checks its own run before
135        // every step, so it stops as soon as possible.
136        self.store()
137            .update_run_status(run_id, RunStatus::Paused)
138            .await?;
139        if from == RunStatus::Running {
140            self.interrupt_running_steps(run_id).await?;
141        }
142        self.publish_transition(&run, RunStatus::Paused);
143        info!(run_id = %run_id, from = %from, "run paused");
144
145        let mut paused_descendants = Vec::new();
146        for descendant in self.store().list_active_descendants(run_id).await? {
147            if !descendant
148                .status
149                .state
150                .can_transition_to(&RunStatus::Paused)
151            {
152                continue;
153            }
154            match self
155                .store()
156                .update_run_status(descendant.id, RunStatus::Paused)
157                .await
158            {
159                Ok(()) => {}
160                // Finished or paused since it was listed: nothing to pause.
161                Err(StoreError::InvalidTransition { .. }) => {
162                    debug!(
163                        run_id = %descendant.id,
164                        "descendant run no longer pausable, skipped"
165                    );
166                    continue;
167                }
168                Err(err) => return Err(err.into()),
169            }
170            if descendant.status.state == RunStatus::Running {
171                self.interrupt_running_steps(descendant.id).await?;
172            }
173            self.publish_transition(&descendant, RunStatus::Paused);
174            paused_descendants.push(descendant.id);
175        }
176
177        if !paused_descendants.is_empty() {
178            info!(
179                run_id = %run_id,
180                count = paused_descendants.len(),
181                "descendant runs paused"
182            );
183        }
184
185        Ok(RunPause {
186            run: self.load_run(run_id).await?,
187            paused_descendants,
188        })
189    }
190
191    /// Resume a paused root run and the sub-workflow runs paused with it.
192    ///
193    /// Each run goes back to the state recorded in [`Run::resume_status`]:
194    /// a run paused while waiting (`Pending`, `Retrying`, `Sleeping`,
195    /// `AwaitingApproval`) waits again, and a run whose approval, human input
196    /// or signal was resolved during the pause is queued. A root that was
197    /// executing is queued with its interrupted steps, which are executed
198    /// again; finished steps are replayed. A sleeping root whose deadline
199    /// passed during the pause is queued at once; a sleeping sub-workflow run
200    /// in the same case is woken by the [`RunWaker`](crate::wake::RunWaker)
201    /// on its next tick.
202    ///
203    /// Under [`ExecutionMode::Local`] a queued root is resumed in a background
204    /// task, and so is a queued sub-workflow run whose root still waits.
205    /// [`Event::RunStatusChanged`] is published for every run resumed.
206    ///
207    /// Under [`ExecutionMode::Local`], an execution still inside its step when
208    /// the run is resumed is not doubled: the run goes back to `Running` and
209    /// that execution carries on, and the background task waits for it to end
210    /// before deciding whether anything is left to restart.
211    ///
212    /// # Errors
213    ///
214    /// - [`EngineError::Store`] with [`StoreError::RunNotFound`] for an
215    ///   unknown run.
216    /// - [`EngineError::ChildRunNotPausable`] for a sub-workflow run: resume
217    ///   its root run instead.
218    /// - [`EngineError::Store`] with [`StoreError::InvalidTransition`] for a
219    ///   run that is not paused.
220    /// - [`EngineError::Store`] when the resume cannot be persisted.
221    ///
222    /// # Examples
223    ///
224    /// ```no_run
225    /// use std::sync::Arc;
226    /// use ironflow_engine::engine::Engine;
227    /// use ironflow_engine::error::EngineError;
228    /// use ironflow_store::models::RunStatus;
229    /// use uuid::Uuid;
230    ///
231    /// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
232    /// let resume = engine.resume_paused_run(run_id).await?;
233    /// assert_ne!(resume.run.status.state, RunStatus::Paused);
234    /// # Ok(())
235    /// # }
236    /// ```
237    pub async fn resume_paused_run(
238        self: &Arc<Self>,
239        run_id: Uuid,
240    ) -> Result<RunResume, EngineError> {
241        let run = self.load_pausable_root(run_id).await?;
242        if run.status.state != RunStatus::Paused {
243            return Err(EngineError::Store(StoreError::InvalidTransition {
244                from: run.status.state,
245                to: RunStatus::Pending,
246            }));
247        }
248
249        // The descendants first: the root re-enters them as soon as it runs.
250        let mut resumed_descendants = Vec::new();
251        let mut queued_descendants = Vec::new();
252        for descendant in self.store().list_active_descendants(run_id).await? {
253            if descendant.status.state != RunStatus::Paused {
254                continue;
255            }
256            // A child that was executing stays `Running`: its root's replay
257            // re-enters it and executes its interrupted steps again.
258            let target = descendant.resume_status.unwrap_or(RunStatus::Pending);
259            self.store()
260                .update_run_status(descendant.id, target)
261                .await?;
262            self.publish_transition(&descendant, target);
263            if target == RunStatus::Pending {
264                queued_descendants.push(descendant.id);
265            }
266            resumed_descendants.push(descendant.id);
267        }
268
269        let target = match run.resume_status {
270            // Stopped by the pause in the middle of its execution: queued
271            // again, like a run whose worker lost its lease.
272            Some(RunStatus::Running) | None => {
273                self.interrupt_running_steps(run_id).await?;
274                RunStatus::Pending
275            }
276            // The deadline passed during the pause: due now, like a run the
277            // waker would have claimed. A past `scheduled_at` does not hold
278            // back the pick.
279            Some(RunStatus::Sleeping) if run.scheduled_at.is_some_and(|at| at <= Utc::now()) => {
280                RunStatus::Pending
281            }
282            Some(status) => status,
283        };
284        self.store().update_run_status(run_id, target).await?;
285        self.publish_transition(&run, target);
286        info!(run_id = %run_id, to = %target, "run resumed");
287
288        if self.execution_mode() == ExecutionMode::Local {
289            if target == RunStatus::Pending {
290                self.continue_in_flight_execution(run_id).await?;
291                self.spawn_local_resume(run_id);
292            } else {
293                // The root waits on its chain: the queued child resumes it.
294                for child_id in queued_descendants {
295                    self.continue_in_flight_execution(child_id).await?;
296                    self.spawn_local_resume(child_id);
297                }
298            }
299        }
300
301        Ok(RunResume {
302            run: self.load_run(run_id).await?,
303            resumed_descendants,
304        })
305    }
306
307    /// Hand a run just queued to `Pending` back to the execution still
308    /// running it in this process, if any.
309    ///
310    /// That execution does not see the pause at a step boundary once the run
311    /// left `Paused`, so the run must be `Running` for it to carry on as a
312    /// legitimate run.
313    async fn continue_in_flight_execution(&self, run_id: Uuid) -> Result<(), EngineError> {
314        if self.is_executing(run_id) {
315            self.store()
316                .update_run_status(run_id, RunStatus::Running)
317                .await?;
318            debug!(run_id = %run_id, "execution still in flight, run continues");
319        }
320        Ok(())
321    }
322
323    /// Pause a registered workflow: its queued runs are no longer picked up.
324    ///
325    /// Runs keep being created; workers skip them until
326    /// [`resume_workflow`](Self::resume_workflow). Runs already executing
327    /// are not affected: pause them with [`pause_run`](Self::pause_run).
328    /// Pausing a paused workflow returns the pause already recorded.
329    ///
330    /// # Errors
331    ///
332    /// Returns [`EngineError::InvalidWorkflow`] when no handler is registered
333    /// under `workflow_name`, and [`EngineError::Store`] when the pause cannot
334    /// be persisted.
335    ///
336    /// # Examples
337    ///
338    /// ```no_run
339    /// use ironflow_engine::engine::Engine;
340    /// use ironflow_engine::error::EngineError;
341    ///
342    /// # async fn example(engine: &Engine) -> Result<(), EngineError> {
343    /// let pause = engine.pause_workflow("deploy", None).await?;
344    /// println!("deploy paused at {}", pause.paused_at);
345    /// # Ok(())
346    /// # }
347    /// ```
348    pub async fn pause_workflow(
349        &self,
350        workflow_name: &str,
351        paused_by: Option<Uuid>,
352    ) -> Result<WorkflowPause, EngineError> {
353        self.require_handler(workflow_name)?;
354        let pause = self
355            .store()
356            .pause_workflow(workflow_name, paused_by)
357            .await?;
358        info!(workflow = %workflow_name, "workflow paused");
359        Ok(pause)
360    }
361
362    /// Resume a paused workflow so its queued runs are picked up again.
363    ///
364    /// Returns `true` when the workflow was paused, `false` when it was not:
365    /// resuming is idempotent.
366    ///
367    /// # Errors
368    ///
369    /// Returns [`EngineError::InvalidWorkflow`] when no handler is registered
370    /// under `workflow_name`, and [`EngineError::Store`] when the resume
371    /// cannot be persisted.
372    ///
373    /// # Examples
374    ///
375    /// ```no_run
376    /// use ironflow_engine::engine::Engine;
377    /// use ironflow_engine::error::EngineError;
378    ///
379    /// # async fn example(engine: &Engine) -> Result<(), EngineError> {
380    /// if engine.resume_workflow("deploy").await? {
381    ///     println!("deploy resumed");
382    /// }
383    /// # Ok(())
384    /// # }
385    /// ```
386    pub async fn resume_workflow(&self, workflow_name: &str) -> Result<bool, EngineError> {
387        self.require_handler(workflow_name)?;
388        let was_paused = self.store().resume_workflow(workflow_name).await?;
389        if was_paused {
390            info!(workflow = %workflow_name, "workflow resumed");
391        }
392        Ok(was_paused)
393    }
394
395    fn require_handler(&self, workflow_name: &str) -> Result<(), EngineError> {
396        match self.get_handler(workflow_name) {
397            Some(_) => Ok(()),
398            None => Err(EngineError::InvalidWorkflow(format!(
399                "no handler registered for workflow '{workflow_name}'"
400            ))),
401        }
402    }
403
404    /// Make the paused root of `run`'s chain resume to `Pending` when it
405    /// waits suspended with its child, so that its replay observes what
406    /// happened to the child (cancelled, rejected) once an operator resumes
407    /// it.
408    pub(crate) async fn requeue_paused_root(&self, run: &Run) -> Result<(), EngineError> {
409        let Some(root_id) = chain_root(run) else {
410            return Ok(());
411        };
412        let root = self.load_run(root_id).await?;
413        let suspended = matches!(
414            root.resume_status,
415            Some(RunStatus::AwaitingApproval | RunStatus::Sleeping)
416        );
417        if root.status.state != RunStatus::Paused || !suspended {
418            return Ok(());
419        }
420
421        let update = RunUpdate {
422            resume_status: Some(RunStatus::Pending),
423            ..RunUpdate::default()
424        };
425        self.store().update_run(root_id, update).await?;
426        info!(run_id = %run.id, root_run_id = %root_id, "paused root run will resume to observe its child");
427        Ok(())
428    }
429
430    /// Load `run_id`, refusing a sub-workflow run.
431    async fn load_pausable_root(&self, run_id: Uuid) -> Result<Run, EngineError> {
432        let run = self.load_run(run_id).await?;
433        match chain_root(&run) {
434            Some(root_run_id) => Err(EngineError::ChildRunNotPausable {
435                run_id,
436                root_run_id,
437            }),
438            None => Ok(run),
439        }
440    }
441
442    /// Publish the move of `run` (as loaded before the update) to `to`.
443    fn publish_transition(&self, run: &Run, to: RunStatus) {
444        self.event_publisher()
445            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
446                run_id: run.id,
447                workflow_name: run.workflow_name.clone(),
448                from: run.status.state,
449                to,
450                error: None,
451                cost_usd: run.cost_usd,
452                duration_ms: run.duration_ms,
453                labels: run.labels.clone(),
454                at: Utc::now(),
455            }));
456    }
457}