ironflow-engine 2.51.7

Workflow orchestration engine for ironflow with FSM-based run lifecycle
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
//! Operator pause of a run, or of a whole workflow.
//!
//! [`Engine::pause_run`] holds a root run, with every sub-workflow run below
//! it, in `Paused` until [`Engine::resume_paused_run`] puts each one back in
//! the state it was paused from, or [`Engine::cancel_run`] stops it. A run
//! executing when it is paused has its step in flight interrupted with
//! [`STEP_INTERRUPTED_ERROR`](ironflow_store::store::STEP_INTERRUPTED_ERROR),
//! no further step is started, and the resume replays the run: finished steps
//! are skipped and the interrupted step is executed again.
//!
//! [`Engine::pause_workflow`] holds the queued runs of a workflow instead:
//! they are created as usual but no worker picks them up until
//! [`Engine::resume_workflow`].

use std::sync::Arc;

use chrono::Utc;
use tracing::{debug, info};
use uuid::Uuid;

use ironflow_store::error::StoreError;
use ironflow_store::models::{Run, RunStatus, RunUpdate, WorkflowPause};

use crate::engine::{Engine, ExecutionMode, chain_root};
use crate::error::EngineError;
use crate::notify::{Event, RunStatusChangedEvent};

/// Outcome of [`Engine::pause_run`].
///
/// # Examples
///
/// ```no_run
/// use std::sync::Arc;
/// use ironflow_engine::engine::Engine;
/// use ironflow_engine::error::EngineError;
/// use uuid::Uuid;
///
/// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
/// let pause = engine.pause_run(run_id).await?;
/// println!(
///     "run {} paused with {} sub-runs",
///     pause.run.id,
///     pause.paused_descendants.len()
/// );
/// # Ok(())
/// # }
/// ```
#[derive(Debug, Clone)]
pub struct RunPause {
    /// The paused run, as stored after the pause.
    pub run: Run,
    /// The sub-workflow runs below it that this call paused, oldest first.
    pub paused_descendants: Vec<Uuid>,
}

/// Outcome of [`Engine::resume_paused_run`].
///
/// # Examples
///
/// ```no_run
/// use std::sync::Arc;
/// use ironflow_engine::engine::Engine;
/// use ironflow_engine::error::EngineError;
/// use uuid::Uuid;
///
/// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
/// let resume = engine.resume_paused_run(run_id).await?;
/// println!("run {} is {}", resume.run.id, resume.run.status.state);
/// # Ok(())
/// # }
/// ```
#[derive(Debug, Clone)]
pub struct RunResume {
    /// The resumed run, as stored after the resume.
    pub run: Run,
    /// The sub-workflow runs below it that this call resumed, oldest first.
    pub resumed_descendants: Vec<Uuid>,
}

impl Engine {
    /// Pause a root run and every active sub-workflow run below it.
    ///
    /// Each run moves to `Paused` and records the state it was paused from in
    /// [`Run::resume_status`]. A run that was executing has its running steps
    /// marked `Failed` with
    /// [`STEP_INTERRUPTED_ERROR`](ironflow_store::store::STEP_INTERRUPTED_ERROR),
    /// so the resume executes them again. A run held by a worker loses its lease: the worker drops
    /// the execution at its next renewal, and the reaper never touches a
    /// paused run. A run executing in-process stops before its next step.
    /// [`Event::RunStatusChanged`] is published for every run paused.
    ///
    /// While the run is paused, an approval, a human input or a signal it
    /// waits for can still be resolved: the decision is recorded and only
    /// changes the state the run resumes to.
    ///
    /// # Errors
    ///
    /// - [`EngineError::Store`] with [`StoreError::RunNotFound`] for an
    ///   unknown run.
    /// - [`EngineError::ChildRunNotPausable`] for a sub-workflow run: pause
    ///   its root run instead.
    /// - [`EngineError::Store`] with [`StoreError::InvalidTransition`] for a
    ///   run already paused or finished.
    /// - [`EngineError::Store`] when the pause cannot be persisted.
    ///
    /// # Examples
    ///
    /// ```no_run
    /// use std::sync::Arc;
    /// use ironflow_engine::engine::Engine;
    /// use ironflow_engine::error::EngineError;
    /// use ironflow_store::models::RunStatus;
    /// use uuid::Uuid;
    ///
    /// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
    /// let pause = engine.pause_run(run_id).await?;
    /// assert_eq!(pause.run.status.state, RunStatus::Paused);
    /// # Ok(())
    /// # }
    /// ```
    pub async fn pause_run(&self, run_id: Uuid) -> Result<RunPause, EngineError> {
        let run = self.load_pausable_root(run_id).await?;
        let from = run.status.state;
        if !from.can_transition_to(&RunStatus::Paused) {
            return Err(EngineError::Store(StoreError::InvalidTransition {
                from,
                to: RunStatus::Paused,
            }));
        }

        // The root first: an execution in flight checks its own run before
        // every step, so it stops as soon as possible.
        self.store()
            .update_run_status(run_id, RunStatus::Paused)
            .await?;
        if from == RunStatus::Running {
            self.interrupt_running_steps(run_id).await?;
        }
        self.publish_transition(&run, RunStatus::Paused);
        info!(run_id = %run_id, from = %from, "run paused");

        let mut paused_descendants = Vec::new();
        for descendant in self.store().list_active_descendants(run_id).await? {
            if !descendant
                .status
                .state
                .can_transition_to(&RunStatus::Paused)
            {
                continue;
            }
            match self
                .store()
                .update_run_status(descendant.id, RunStatus::Paused)
                .await
            {
                Ok(()) => {}
                // Finished or paused since it was listed: nothing to pause.
                Err(StoreError::InvalidTransition { .. }) => {
                    debug!(
                        run_id = %descendant.id,
                        "descendant run no longer pausable, skipped"
                    );
                    continue;
                }
                Err(err) => return Err(err.into()),
            }
            if descendant.status.state == RunStatus::Running {
                self.interrupt_running_steps(descendant.id).await?;
            }
            self.publish_transition(&descendant, RunStatus::Paused);
            paused_descendants.push(descendant.id);
        }

        if !paused_descendants.is_empty() {
            info!(
                run_id = %run_id,
                count = paused_descendants.len(),
                "descendant runs paused"
            );
        }

        Ok(RunPause {
            run: self.load_run(run_id).await?,
            paused_descendants,
        })
    }

    /// Resume a paused root run and the sub-workflow runs paused with it.
    ///
    /// Each run goes back to the state recorded in [`Run::resume_status`]:
    /// a run paused while waiting (`Pending`, `Retrying`, `Sleeping`,
    /// `AwaitingApproval`) waits again, and a run whose approval, human input
    /// or signal was resolved during the pause is queued. A root that was
    /// executing is queued with its interrupted steps, which are executed
    /// again; finished steps are replayed. A sleeping root whose deadline
    /// passed during the pause is queued at once; a sleeping sub-workflow run
    /// in the same case is woken by the [`RunWaker`](crate::wake::RunWaker)
    /// on its next tick.
    ///
    /// Under [`ExecutionMode::Local`] a queued root is resumed in a background
    /// task, and so is a queued sub-workflow run whose root still waits.
    /// [`Event::RunStatusChanged`] is published for every run resumed.
    ///
    /// # Errors
    ///
    /// - [`EngineError::Store`] with [`StoreError::RunNotFound`] for an
    ///   unknown run.
    /// - [`EngineError::ChildRunNotPausable`] for a sub-workflow run: resume
    ///   its root run instead.
    /// - [`EngineError::Store`] with [`StoreError::InvalidTransition`] for a
    ///   run that is not paused.
    /// - [`EngineError::Store`] when the resume cannot be persisted.
    ///
    /// # Examples
    ///
    /// ```no_run
    /// use std::sync::Arc;
    /// use ironflow_engine::engine::Engine;
    /// use ironflow_engine::error::EngineError;
    /// use ironflow_store::models::RunStatus;
    /// use uuid::Uuid;
    ///
    /// # async fn example(engine: Arc<Engine>, run_id: Uuid) -> Result<(), EngineError> {
    /// let resume = engine.resume_paused_run(run_id).await?;
    /// assert_ne!(resume.run.status.state, RunStatus::Paused);
    /// # Ok(())
    /// # }
    /// ```
    pub async fn resume_paused_run(
        self: &Arc<Self>,
        run_id: Uuid,
    ) -> Result<RunResume, EngineError> {
        let run = self.load_pausable_root(run_id).await?;
        if run.status.state != RunStatus::Paused {
            return Err(EngineError::Store(StoreError::InvalidTransition {
                from: run.status.state,
                to: RunStatus::Pending,
            }));
        }

        // The descendants first: the root re-enters them as soon as it runs.
        let mut resumed_descendants = Vec::new();
        let mut queued_descendants = Vec::new();
        for descendant in self.store().list_active_descendants(run_id).await? {
            if descendant.status.state != RunStatus::Paused {
                continue;
            }
            // A child that was executing stays `Running`: its root's replay
            // re-enters it and executes its interrupted steps again.
            let target = descendant.resume_status.unwrap_or(RunStatus::Pending);
            self.store()
                .update_run_status(descendant.id, target)
                .await?;
            self.publish_transition(&descendant, target);
            if target == RunStatus::Pending {
                queued_descendants.push(descendant.id);
            }
            resumed_descendants.push(descendant.id);
        }

        let target = match run.resume_status {
            // Stopped by the pause in the middle of its execution: queued
            // again, like a run whose worker lost its lease.
            Some(RunStatus::Running) | None => {
                self.interrupt_running_steps(run_id).await?;
                RunStatus::Pending
            }
            // The deadline passed during the pause: due now, like a run the
            // waker would have claimed. A past `scheduled_at` does not hold
            // back the pick.
            Some(RunStatus::Sleeping) if run.scheduled_at.is_some_and(|at| at <= Utc::now()) => {
                RunStatus::Pending
            }
            Some(status) => status,
        };
        self.store().update_run_status(run_id, target).await?;
        self.publish_transition(&run, target);
        info!(run_id = %run_id, to = %target, "run resumed");

        if self.execution_mode() == ExecutionMode::Local {
            if target == RunStatus::Pending {
                self.spawn_local_resume(run_id);
            } else {
                // The root waits on its chain: the queued child resumes it.
                for child_id in queued_descendants {
                    self.spawn_local_resume(child_id);
                }
            }
        }

        Ok(RunResume {
            run: self.load_run(run_id).await?,
            resumed_descendants,
        })
    }

    /// Pause a registered workflow: its queued runs are no longer picked up.
    ///
    /// Runs keep being created; workers skip them until
    /// [`resume_workflow`](Self::resume_workflow). Runs already executing
    /// are not affected: pause them with [`pause_run`](Self::pause_run).
    /// Pausing a paused workflow returns the pause already recorded.
    ///
    /// # Errors
    ///
    /// Returns [`EngineError::InvalidWorkflow`] when no handler is registered
    /// under `workflow_name`, and [`EngineError::Store`] when the pause cannot
    /// be persisted.
    ///
    /// # Examples
    ///
    /// ```no_run
    /// use ironflow_engine::engine::Engine;
    /// use ironflow_engine::error::EngineError;
    ///
    /// # async fn example(engine: &Engine) -> Result<(), EngineError> {
    /// let pause = engine.pause_workflow("deploy", None).await?;
    /// println!("deploy paused at {}", pause.paused_at);
    /// # Ok(())
    /// # }
    /// ```
    pub async fn pause_workflow(
        &self,
        workflow_name: &str,
        paused_by: Option<Uuid>,
    ) -> Result<WorkflowPause, EngineError> {
        self.require_handler(workflow_name)?;
        let pause = self
            .store()
            .pause_workflow(workflow_name, paused_by)
            .await?;
        info!(workflow = %workflow_name, "workflow paused");
        Ok(pause)
    }

    /// Resume a paused workflow so its queued runs are picked up again.
    ///
    /// Returns `true` when the workflow was paused, `false` when it was not:
    /// resuming is idempotent.
    ///
    /// # Errors
    ///
    /// Returns [`EngineError::InvalidWorkflow`] when no handler is registered
    /// under `workflow_name`, and [`EngineError::Store`] when the resume
    /// cannot be persisted.
    ///
    /// # Examples
    ///
    /// ```no_run
    /// use ironflow_engine::engine::Engine;
    /// use ironflow_engine::error::EngineError;
    ///
    /// # async fn example(engine: &Engine) -> Result<(), EngineError> {
    /// if engine.resume_workflow("deploy").await? {
    ///     println!("deploy resumed");
    /// }
    /// # Ok(())
    /// # }
    /// ```
    pub async fn resume_workflow(&self, workflow_name: &str) -> Result<bool, EngineError> {
        self.require_handler(workflow_name)?;
        let was_paused = self.store().resume_workflow(workflow_name).await?;
        if was_paused {
            info!(workflow = %workflow_name, "workflow resumed");
        }
        Ok(was_paused)
    }

    fn require_handler(&self, workflow_name: &str) -> Result<(), EngineError> {
        match self.get_handler(workflow_name) {
            Some(_) => Ok(()),
            None => Err(EngineError::InvalidWorkflow(format!(
                "no handler registered for workflow '{workflow_name}'"
            ))),
        }
    }

    /// Make the paused root of `run`'s chain resume to `Pending` when it
    /// waits suspended with its child, so that its replay observes what
    /// happened to the child (cancelled, rejected) once an operator resumes
    /// it.
    pub(crate) async fn requeue_paused_root(&self, run: &Run) -> Result<(), EngineError> {
        let Some(root_id) = chain_root(run) else {
            return Ok(());
        };
        let root = self.load_run(root_id).await?;
        let suspended = matches!(
            root.resume_status,
            Some(RunStatus::AwaitingApproval | RunStatus::Sleeping)
        );
        if root.status.state != RunStatus::Paused || !suspended {
            return Ok(());
        }

        let update = RunUpdate {
            resume_status: Some(RunStatus::Pending),
            ..RunUpdate::default()
        };
        self.store().update_run(root_id, update).await?;
        info!(run_id = %run.id, root_run_id = %root_id, "paused root run will resume to observe its child");
        Ok(())
    }

    /// Load `run_id`, refusing a sub-workflow run.
    async fn load_pausable_root(&self, run_id: Uuid) -> Result<Run, EngineError> {
        let run = self.load_run(run_id).await?;
        match chain_root(&run) {
            Some(root_run_id) => Err(EngineError::ChildRunNotPausable {
                run_id,
                root_run_id,
            }),
            None => Ok(run),
        }
    }

    /// Publish the move of `run` (as loaded before the update) to `to`.
    fn publish_transition(&self, run: &Run, to: RunStatus) {
        self.event_publisher()
            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
                run_id: run.id,
                workflow_name: run.workflow_name.clone(),
                from: run.status.state,
                to,
                error: None,
                cost_usd: run.cost_usd,
                duration_ms: run.duration_ms,
                labels: run.labels.clone(),
                at: Utc::now(),
            }));
    }
}