Skip to main content

ironflow_engine/fsm/
run_fsm.rs

1//! [`RunFsm`] — Finite state machine for the run lifecycle.
2//!
3//! Events drive transitions; the FSM rejects invalid ones and keeps
4//! a full history of state changes.
5
6use chrono::Utc;
7use ironflow_store::entities::RunStatus;
8use serde::{Deserialize, Serialize};
9use strum::Display;
10
11use super::{Transition, TransitionError};
12
13/// Events that drive [`RunFsm`] transitions.
14///
15/// Each event represents something that happened during execution.
16///
17/// # Examples
18///
19/// ```
20/// use ironflow_engine::fsm::RunEvent;
21///
22/// let event = RunEvent::PickedUp;
23/// assert_eq!(event.to_string(), "picked_up");
24/// ```
25#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Display)]
26#[serde(rename_all = "snake_case")]
27#[strum(serialize_all = "snake_case")]
28pub enum RunEvent {
29    /// Worker or inline executor picked up the run.
30    PickedUp,
31    /// All steps completed successfully.
32    AllStepsCompleted,
33    /// A step failed and the error is not retryable, or max retries exhausted.
34    StepFailed,
35    /// A step failed but a retry is possible.
36    StepFailedRetryable,
37    /// Retry attempt started.
38    RetryStarted,
39    /// Maximum retries exhausted after a retryable failure.
40    MaxRetriesExceeded,
41    /// User or API requested cancellation.
42    CancelRequested,
43    /// A step requires human approval before continuing.
44    ApprovalRequested,
45    /// Human approved the run to continue.
46    Approved,
47    /// Human rejected the run.
48    Rejected,
49    /// A delay step suspended the run until a scheduled time.
50    DelaySleeping,
51    /// The delay elapsed and the run is re-queued.
52    DelayElapsed,
53    /// A signal resumed the run.
54    SignalReceived,
55    /// An operator paused the run.
56    PauseRequested,
57    /// An operator resumed a paused run, which goes back to the queue.
58    ///
59    /// The engine restores the exact state the run was paused in (kept in
60    /// `Run::resume_status`); the FSM models the common case, a requeue.
61    ResumeRequested,
62}
63
64/// Finite state machine for a workflow run.
65///
66/// Wraps a [`RunStatus`] and enforces valid transitions via typed
67/// [`RunEvent`]s. Records every transition in a history log.
68///
69/// # Transition table
70///
71/// | From | Event | To |
72/// |------|-------|----|
73/// | Pending | PickedUp | Running |
74/// | Pending | CancelRequested | Cancelled |
75/// | Running | AllStepsCompleted | Completed |
76/// | Running | StepFailed | Failed |
77/// | Running | StepFailedRetryable | Retrying |
78/// | Running | CancelRequested | Cancelled |
79/// | Retrying | RetryStarted | Running |
80/// | Retrying | MaxRetriesExceeded | Failed |
81/// | Retrying | CancelRequested | Cancelled |
82/// | Running | ApprovalRequested | AwaitingApproval |
83/// | AwaitingApproval | Approved | Running |
84/// | AwaitingApproval | Rejected | Failed |
85/// | AwaitingApproval | CancelRequested | Cancelled |
86/// | Running | DelaySleeping | Sleeping |
87/// | Sleeping | DelayElapsed | Pending |
88/// | Sleeping | SignalReceived | Pending |
89/// | Sleeping | CancelRequested | Cancelled |
90/// | Pending, Retrying, Sleeping, AwaitingApproval, Running | PauseRequested | Paused |
91/// | Paused | ResumeRequested | Pending |
92/// | Paused | CancelRequested | Cancelled |
93///
94/// # Examples
95///
96/// ```
97/// use ironflow_engine::fsm::{RunFsm, RunEvent};
98/// use ironflow_store::entities::RunStatus;
99///
100/// let mut fsm = RunFsm::new();
101/// assert_eq!(fsm.state(), RunStatus::Pending);
102///
103/// fsm.apply(RunEvent::PickedUp).unwrap();
104/// assert_eq!(fsm.state(), RunStatus::Running);
105///
106/// fsm.apply(RunEvent::AllStepsCompleted).unwrap();
107/// assert_eq!(fsm.state(), RunStatus::Completed);
108/// assert_eq!(fsm.history().len(), 2);
109/// ```
110#[derive(Debug, Clone)]
111pub struct RunFsm {
112    state: RunStatus,
113    history: Vec<Transition<RunStatus, RunEvent>>,
114}
115
116impl RunFsm {
117    /// Create a new FSM in `Pending` state.
118    ///
119    /// # Examples
120    ///
121    /// ```
122    /// use ironflow_engine::fsm::RunFsm;
123    /// use ironflow_store::entities::RunStatus;
124    ///
125    /// let fsm = RunFsm::new();
126    /// assert_eq!(fsm.state(), RunStatus::Pending);
127    /// ```
128    pub fn new() -> Self {
129        Self {
130            state: RunStatus::Pending,
131            history: Vec::new(),
132        }
133    }
134
135    /// Create a FSM from an existing state (e.g. loaded from DB).
136    ///
137    /// # Examples
138    ///
139    /// ```
140    /// use ironflow_engine::fsm::RunFsm;
141    /// use ironflow_store::entities::RunStatus;
142    ///
143    /// let fsm = RunFsm::from_state(RunStatus::Running);
144    /// assert_eq!(fsm.state(), RunStatus::Running);
145    /// ```
146    pub fn from_state(state: RunStatus) -> Self {
147        Self {
148            state,
149            history: Vec::new(),
150        }
151    }
152
153    /// Returns the current state.
154    pub fn state(&self) -> RunStatus {
155        self.state
156    }
157
158    /// Returns the full transition history.
159    pub fn history(&self) -> &[Transition<RunStatus, RunEvent>] {
160        &self.history
161    }
162
163    /// Returns `true` if the FSM is in a terminal state.
164    pub fn is_terminal(&self) -> bool {
165        self.state.is_terminal()
166    }
167
168    /// Apply an event, transitioning to a new state if valid.
169    ///
170    /// # Errors
171    ///
172    /// Returns [`TransitionError`] if the event is not allowed in the current state.
173    ///
174    /// # Examples
175    ///
176    /// ```
177    /// use ironflow_engine::fsm::{RunFsm, RunEvent};
178    /// use ironflow_store::entities::RunStatus;
179    ///
180    /// let mut fsm = RunFsm::new();
181    ///
182    /// // Valid transition
183    /// assert!(fsm.apply(RunEvent::PickedUp).is_ok());
184    ///
185    /// // Invalid: can't pick up a running run
186    /// assert!(fsm.apply(RunEvent::PickedUp).is_err());
187    /// ```
188    pub fn apply(
189        &mut self,
190        event: RunEvent,
191    ) -> Result<RunStatus, TransitionError<RunStatus, RunEvent>> {
192        let next = next_state(self.state, event).ok_or(TransitionError {
193            from: self.state,
194            event,
195        })?;
196
197        let transition = Transition {
198            from: self.state,
199            to: next,
200            event,
201            at: Utc::now(),
202        };
203
204        self.history.push(transition);
205        self.state = next;
206        Ok(next)
207    }
208
209    /// Check if an event would be accepted without applying it.
210    ///
211    /// # Examples
212    ///
213    /// ```
214    /// use ironflow_engine::fsm::{RunFsm, RunEvent};
215    ///
216    /// let fsm = RunFsm::new();
217    /// assert!(fsm.can_apply(RunEvent::PickedUp));
218    /// assert!(!fsm.can_apply(RunEvent::AllStepsCompleted));
219    /// ```
220    pub fn can_apply(&self, event: RunEvent) -> bool {
221        next_state(self.state, event).is_some()
222    }
223}
224
225impl Default for RunFsm {
226    fn default() -> Self {
227        Self::new()
228    }
229}
230
231/// Pure transition function — returns the next state for a given (state, event)
232/// pair, or `None` if the transition is invalid.
233fn next_state(from: RunStatus, event: RunEvent) -> Option<RunStatus> {
234    match (from, event) {
235        // Pending
236        (RunStatus::Pending, RunEvent::PickedUp) => Some(RunStatus::Running),
237        (RunStatus::Pending, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
238
239        // Running
240        (RunStatus::Running, RunEvent::AllStepsCompleted) => Some(RunStatus::Completed),
241        (RunStatus::Running, RunEvent::StepFailed) => Some(RunStatus::Failed),
242        (RunStatus::Running, RunEvent::StepFailedRetryable) => Some(RunStatus::Retrying),
243        (RunStatus::Running, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
244
245        // Retrying
246        (RunStatus::Retrying, RunEvent::RetryStarted) => Some(RunStatus::Running),
247        (RunStatus::Retrying, RunEvent::MaxRetriesExceeded) => Some(RunStatus::Failed),
248        (RunStatus::Retrying, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
249
250        // Approval
251        (RunStatus::Running, RunEvent::ApprovalRequested) => Some(RunStatus::AwaitingApproval),
252        (RunStatus::AwaitingApproval, RunEvent::Approved) => Some(RunStatus::Running),
253        (RunStatus::AwaitingApproval, RunEvent::Rejected) => Some(RunStatus::Failed),
254        (RunStatus::AwaitingApproval, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
255
256        // Delay
257        (RunStatus::Running, RunEvent::DelaySleeping) => Some(RunStatus::Sleeping),
258        (RunStatus::Sleeping, RunEvent::DelayElapsed) => Some(RunStatus::Pending),
259        (RunStatus::Sleeping, RunEvent::SignalReceived) => Some(RunStatus::Pending),
260        (RunStatus::Sleeping, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
261
262        // Pause
263        (
264            RunStatus::Pending
265            | RunStatus::Retrying
266            | RunStatus::Sleeping
267            | RunStatus::AwaitingApproval
268            | RunStatus::Running,
269            RunEvent::PauseRequested,
270        ) => Some(RunStatus::Paused),
271        (RunStatus::Paused, RunEvent::ResumeRequested) => Some(RunStatus::Pending),
272        (RunStatus::Paused, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
273
274        // Terminal states and all other combos → invalid
275        _ => None,
276    }
277}
278
279#[cfg(test)]
280mod tests {
281    use super::*;
282
283    // ---- Happy paths ----
284
285    #[test]
286    fn pending_to_running() {
287        let mut fsm = RunFsm::new();
288        let result = fsm.apply(RunEvent::PickedUp);
289        assert!(result.is_ok());
290        assert_eq!(fsm.state(), RunStatus::Running);
291    }
292
293    #[test]
294    fn full_success_path() {
295        let mut fsm = RunFsm::new();
296        fsm.apply(RunEvent::PickedUp).unwrap();
297        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
298        assert_eq!(fsm.state(), RunStatus::Completed);
299        assert!(fsm.is_terminal());
300        assert_eq!(fsm.history().len(), 2);
301    }
302
303    #[test]
304    fn full_failure_path() {
305        let mut fsm = RunFsm::new();
306        fsm.apply(RunEvent::PickedUp).unwrap();
307        fsm.apply(RunEvent::StepFailed).unwrap();
308        assert_eq!(fsm.state(), RunStatus::Failed);
309        assert!(fsm.is_terminal());
310    }
311
312    #[test]
313    fn retry_then_success() {
314        let mut fsm = RunFsm::new();
315        fsm.apply(RunEvent::PickedUp).unwrap();
316        fsm.apply(RunEvent::StepFailedRetryable).unwrap();
317        assert_eq!(fsm.state(), RunStatus::Retrying);
318
319        fsm.apply(RunEvent::RetryStarted).unwrap();
320        assert_eq!(fsm.state(), RunStatus::Running);
321
322        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
323        assert_eq!(fsm.state(), RunStatus::Completed);
324        assert_eq!(fsm.history().len(), 4);
325    }
326
327    #[test]
328    fn retry_then_max_retries_exceeded() {
329        let mut fsm = RunFsm::new();
330        fsm.apply(RunEvent::PickedUp).unwrap();
331        fsm.apply(RunEvent::StepFailedRetryable).unwrap();
332        fsm.apply(RunEvent::MaxRetriesExceeded).unwrap();
333        assert_eq!(fsm.state(), RunStatus::Failed);
334    }
335
336    #[test]
337    fn cancel_from_pending() {
338        let mut fsm = RunFsm::new();
339        fsm.apply(RunEvent::CancelRequested).unwrap();
340        assert_eq!(fsm.state(), RunStatus::Cancelled);
341        assert!(fsm.is_terminal());
342    }
343
344    #[test]
345    fn cancel_from_running() {
346        let mut fsm = RunFsm::new();
347        fsm.apply(RunEvent::PickedUp).unwrap();
348        fsm.apply(RunEvent::CancelRequested).unwrap();
349        assert_eq!(fsm.state(), RunStatus::Cancelled);
350    }
351
352    #[test]
353    fn cancel_from_retrying() {
354        let mut fsm = RunFsm::new();
355        fsm.apply(RunEvent::PickedUp).unwrap();
356        fsm.apply(RunEvent::StepFailedRetryable).unwrap();
357        fsm.apply(RunEvent::CancelRequested).unwrap();
358        assert_eq!(fsm.state(), RunStatus::Cancelled);
359    }
360
361    // ---- Invalid transitions ----
362
363    #[test]
364    fn cannot_complete_from_pending() {
365        let mut fsm = RunFsm::new();
366        let result = fsm.apply(RunEvent::AllStepsCompleted);
367        assert!(result.is_err());
368        assert_eq!(fsm.state(), RunStatus::Pending);
369    }
370
371    #[test]
372    fn cannot_pick_up_running() {
373        let mut fsm = RunFsm::new();
374        fsm.apply(RunEvent::PickedUp).unwrap();
375        let result = fsm.apply(RunEvent::PickedUp);
376        assert!(result.is_err());
377    }
378
379    #[test]
380    fn cannot_transition_from_terminal() {
381        let mut fsm = RunFsm::new();
382        fsm.apply(RunEvent::PickedUp).unwrap();
383        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
384
385        assert!(fsm.apply(RunEvent::PickedUp).is_err());
386        assert!(fsm.apply(RunEvent::CancelRequested).is_err());
387        assert!(fsm.apply(RunEvent::StepFailed).is_err());
388    }
389
390    // ---- can_apply ----
391
392    #[test]
393    fn can_apply_checks_without_mutation() {
394        let fsm = RunFsm::new();
395        assert!(fsm.can_apply(RunEvent::PickedUp));
396        assert!(fsm.can_apply(RunEvent::CancelRequested));
397        assert!(!fsm.can_apply(RunEvent::AllStepsCompleted));
398        assert!(!fsm.can_apply(RunEvent::StepFailed));
399        assert_eq!(fsm.state(), RunStatus::Pending);
400    }
401
402    // ---- from_state ----
403
404    #[test]
405    fn from_state_resumes_at_given_state() {
406        let mut fsm = RunFsm::from_state(RunStatus::Running);
407        assert_eq!(fsm.state(), RunStatus::Running);
408        assert!(fsm.history().is_empty());
409
410        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
411        assert_eq!(fsm.state(), RunStatus::Completed);
412    }
413
414    // ---- History ----
415
416    #[test]
417    fn history_records_transitions() {
418        let mut fsm = RunFsm::new();
419        fsm.apply(RunEvent::PickedUp).unwrap();
420        fsm.apply(RunEvent::StepFailedRetryable).unwrap();
421        fsm.apply(RunEvent::RetryStarted).unwrap();
422
423        let history = fsm.history();
424        assert_eq!(history.len(), 3);
425
426        assert_eq!(history[0].from, RunStatus::Pending);
427        assert_eq!(history[0].to, RunStatus::Running);
428        assert_eq!(history[0].event, RunEvent::PickedUp);
429
430        assert_eq!(history[1].from, RunStatus::Running);
431        assert_eq!(history[1].to, RunStatus::Retrying);
432        assert_eq!(history[1].event, RunEvent::StepFailedRetryable);
433
434        assert_eq!(history[2].from, RunStatus::Retrying);
435        assert_eq!(history[2].to, RunStatus::Running);
436        assert_eq!(history[2].event, RunEvent::RetryStarted);
437    }
438
439    // ---- Approval transitions ----
440
441    #[test]
442    fn running_to_awaiting_approval() {
443        let mut fsm = RunFsm::new();
444        fsm.apply(RunEvent::PickedUp).unwrap();
445        fsm.apply(RunEvent::ApprovalRequested).unwrap();
446        assert_eq!(fsm.state(), RunStatus::AwaitingApproval);
447        assert!(!fsm.is_terminal());
448    }
449
450    #[test]
451    fn awaiting_approval_approved_resumes_running() {
452        let mut fsm = RunFsm::new();
453        fsm.apply(RunEvent::PickedUp).unwrap();
454        fsm.apply(RunEvent::ApprovalRequested).unwrap();
455        fsm.apply(RunEvent::Approved).unwrap();
456        assert_eq!(fsm.state(), RunStatus::Running);
457    }
458
459    #[test]
460    fn awaiting_approval_rejected_fails() {
461        let mut fsm = RunFsm::new();
462        fsm.apply(RunEvent::PickedUp).unwrap();
463        fsm.apply(RunEvent::ApprovalRequested).unwrap();
464        fsm.apply(RunEvent::Rejected).unwrap();
465        assert_eq!(fsm.state(), RunStatus::Failed);
466        assert!(fsm.is_terminal());
467    }
468
469    #[test]
470    fn awaiting_approval_cancel() {
471        let mut fsm = RunFsm::new();
472        fsm.apply(RunEvent::PickedUp).unwrap();
473        fsm.apply(RunEvent::ApprovalRequested).unwrap();
474        fsm.apply(RunEvent::CancelRequested).unwrap();
475        assert_eq!(fsm.state(), RunStatus::Cancelled);
476        assert!(fsm.is_terminal());
477    }
478
479    #[test]
480    fn cannot_approve_from_pending() {
481        let mut fsm = RunFsm::new();
482        assert!(fsm.apply(RunEvent::Approved).is_err());
483    }
484
485    #[test]
486    fn approval_then_complete() {
487        let mut fsm = RunFsm::new();
488        fsm.apply(RunEvent::PickedUp).unwrap();
489        fsm.apply(RunEvent::ApprovalRequested).unwrap();
490        fsm.apply(RunEvent::Approved).unwrap();
491        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
492        assert_eq!(fsm.state(), RunStatus::Completed);
493        assert_eq!(fsm.history().len(), 4);
494    }
495
496    #[test]
497    fn sleeping_signal_received_goes_pending() {
498        let mut fsm = RunFsm::new();
499        fsm.apply(RunEvent::PickedUp).unwrap();
500        fsm.apply(RunEvent::DelaySleeping).unwrap();
501        fsm.apply(RunEvent::SignalReceived).unwrap();
502        assert_eq!(fsm.state(), RunStatus::Pending);
503        assert_eq!(RunEvent::SignalReceived.to_string(), "signal_received");
504    }
505
506    #[test]
507    fn cannot_receive_signal_while_running() {
508        let mut fsm = RunFsm::new();
509        fsm.apply(RunEvent::PickedUp).unwrap();
510        assert!(fsm.apply(RunEvent::SignalReceived).is_err());
511    }
512
513    #[test]
514    fn every_active_state_can_be_paused() {
515        for state in [
516            RunStatus::Pending,
517            RunStatus::Retrying,
518            RunStatus::Sleeping,
519            RunStatus::AwaitingApproval,
520            RunStatus::Running,
521        ] {
522            let mut fsm = RunFsm::from_state(state);
523            assert_eq!(
524                fsm.apply(RunEvent::PauseRequested).unwrap(),
525                RunStatus::Paused
526            );
527        }
528        assert_eq!(RunEvent::PauseRequested.to_string(), "pause_requested");
529    }
530
531    #[test]
532    fn paused_resume_goes_pending() {
533        let mut fsm = RunFsm::new();
534        fsm.apply(RunEvent::PauseRequested).unwrap();
535        fsm.apply(RunEvent::ResumeRequested).unwrap();
536        assert_eq!(fsm.state(), RunStatus::Pending);
537        assert_eq!(RunEvent::ResumeRequested.to_string(), "resume_requested");
538    }
539
540    #[test]
541    fn paused_can_be_cancelled() {
542        let mut fsm = RunFsm::from_state(RunStatus::Paused);
543        assert_eq!(
544            fsm.apply(RunEvent::CancelRequested).unwrap(),
545            RunStatus::Cancelled
546        );
547    }
548
549    #[test]
550    fn cannot_pause_terminal_or_paused_run() {
551        for state in [
552            RunStatus::Completed,
553            RunStatus::Failed,
554            RunStatus::Cancelled,
555            RunStatus::Warning,
556            RunStatus::Paused,
557        ] {
558            assert!(!RunFsm::from_state(state).can_apply(RunEvent::PauseRequested));
559        }
560    }
561
562    #[test]
563    fn cannot_resume_a_run_that_is_not_paused() {
564        let mut fsm = RunFsm::new();
565        assert!(fsm.apply(RunEvent::ResumeRequested).is_err());
566    }
567
568    // ---- TransitionError Display ----
569
570    #[test]
571    fn transition_error_display() {
572        let mut fsm = RunFsm::new();
573        let err = fsm.apply(RunEvent::AllStepsCompleted).unwrap_err();
574        let msg = err.to_string();
575        assert!(msg.contains("all_steps_completed"));
576        assert!(msg.contains("Pending"));
577    }
578}