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.
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}