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}