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}
56
57/// Finite state machine for a workflow run.
58///
59/// Wraps a [`RunStatus`] and enforces valid transitions via typed
60/// [`RunEvent`]s. Records every transition in a history log.
61///
62/// # Transition table
63///
64/// | From | Event | To |
65/// |------|-------|----|
66/// | Pending | PickedUp | Running |
67/// | Pending | CancelRequested | Cancelled |
68/// | Running | AllStepsCompleted | Completed |
69/// | Running | StepFailed | Failed |
70/// | Running | StepFailedRetryable | Retrying |
71/// | Running | CancelRequested | Cancelled |
72/// | Retrying | RetryStarted | Running |
73/// | Retrying | MaxRetriesExceeded | Failed |
74/// | Retrying | CancelRequested | Cancelled |
75/// | Running | ApprovalRequested | AwaitingApproval |
76/// | AwaitingApproval | Approved | Running |
77/// | AwaitingApproval | Rejected | Failed |
78/// | AwaitingApproval | CancelRequested | Cancelled |
79/// | Running | DelaySleeping | Sleeping |
80/// | Sleeping | DelayElapsed | Pending |
81/// | Sleeping | SignalReceived | Pending |
82/// | Sleeping | CancelRequested | Cancelled |
83///
84/// # Examples
85///
86/// ```
87/// use ironflow_engine::fsm::{RunFsm, RunEvent};
88/// use ironflow_store::entities::RunStatus;
89///
90/// let mut fsm = RunFsm::new();
91/// assert_eq!(fsm.state(), RunStatus::Pending);
92///
93/// fsm.apply(RunEvent::PickedUp).unwrap();
94/// assert_eq!(fsm.state(), RunStatus::Running);
95///
96/// fsm.apply(RunEvent::AllStepsCompleted).unwrap();
97/// assert_eq!(fsm.state(), RunStatus::Completed);
98/// assert_eq!(fsm.history().len(), 2);
99/// ```
100#[derive(Debug, Clone)]
101pub struct RunFsm {
102    state: RunStatus,
103    history: Vec<Transition<RunStatus, RunEvent>>,
104}
105
106impl RunFsm {
107    /// Create a new FSM in `Pending` state.
108    ///
109    /// # Examples
110    ///
111    /// ```
112    /// use ironflow_engine::fsm::RunFsm;
113    /// use ironflow_store::entities::RunStatus;
114    ///
115    /// let fsm = RunFsm::new();
116    /// assert_eq!(fsm.state(), RunStatus::Pending);
117    /// ```
118    pub fn new() -> Self {
119        Self {
120            state: RunStatus::Pending,
121            history: Vec::new(),
122        }
123    }
124
125    /// Create a FSM from an existing state (e.g. loaded from DB).
126    ///
127    /// # Examples
128    ///
129    /// ```
130    /// use ironflow_engine::fsm::RunFsm;
131    /// use ironflow_store::entities::RunStatus;
132    ///
133    /// let fsm = RunFsm::from_state(RunStatus::Running);
134    /// assert_eq!(fsm.state(), RunStatus::Running);
135    /// ```
136    pub fn from_state(state: RunStatus) -> Self {
137        Self {
138            state,
139            history: Vec::new(),
140        }
141    }
142
143    /// Returns the current state.
144    pub fn state(&self) -> RunStatus {
145        self.state
146    }
147
148    /// Returns the full transition history.
149    pub fn history(&self) -> &[Transition<RunStatus, RunEvent>] {
150        &self.history
151    }
152
153    /// Returns `true` if the FSM is in a terminal state.
154    pub fn is_terminal(&self) -> bool {
155        self.state.is_terminal()
156    }
157
158    /// Apply an event, transitioning to a new state if valid.
159    ///
160    /// # Errors
161    ///
162    /// Returns [`TransitionError`] if the event is not allowed in the current state.
163    ///
164    /// # Examples
165    ///
166    /// ```
167    /// use ironflow_engine::fsm::{RunFsm, RunEvent};
168    /// use ironflow_store::entities::RunStatus;
169    ///
170    /// let mut fsm = RunFsm::new();
171    ///
172    /// // Valid transition
173    /// assert!(fsm.apply(RunEvent::PickedUp).is_ok());
174    ///
175    /// // Invalid: can't pick up a running run
176    /// assert!(fsm.apply(RunEvent::PickedUp).is_err());
177    /// ```
178    pub fn apply(
179        &mut self,
180        event: RunEvent,
181    ) -> Result<RunStatus, TransitionError<RunStatus, RunEvent>> {
182        let next = next_state(self.state, event).ok_or(TransitionError {
183            from: self.state,
184            event,
185        })?;
186
187        let transition = Transition {
188            from: self.state,
189            to: next,
190            event,
191            at: Utc::now(),
192        };
193
194        self.history.push(transition);
195        self.state = next;
196        Ok(next)
197    }
198
199    /// Check if an event would be accepted without applying it.
200    ///
201    /// # Examples
202    ///
203    /// ```
204    /// use ironflow_engine::fsm::{RunFsm, RunEvent};
205    ///
206    /// let fsm = RunFsm::new();
207    /// assert!(fsm.can_apply(RunEvent::PickedUp));
208    /// assert!(!fsm.can_apply(RunEvent::AllStepsCompleted));
209    /// ```
210    pub fn can_apply(&self, event: RunEvent) -> bool {
211        next_state(self.state, event).is_some()
212    }
213}
214
215impl Default for RunFsm {
216    fn default() -> Self {
217        Self::new()
218    }
219}
220
221/// Pure transition function — returns the next state for a given (state, event)
222/// pair, or `None` if the transition is invalid.
223fn next_state(from: RunStatus, event: RunEvent) -> Option<RunStatus> {
224    match (from, event) {
225        // Pending
226        (RunStatus::Pending, RunEvent::PickedUp) => Some(RunStatus::Running),
227        (RunStatus::Pending, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
228
229        // Running
230        (RunStatus::Running, RunEvent::AllStepsCompleted) => Some(RunStatus::Completed),
231        (RunStatus::Running, RunEvent::StepFailed) => Some(RunStatus::Failed),
232        (RunStatus::Running, RunEvent::StepFailedRetryable) => Some(RunStatus::Retrying),
233        (RunStatus::Running, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
234
235        // Retrying
236        (RunStatus::Retrying, RunEvent::RetryStarted) => Some(RunStatus::Running),
237        (RunStatus::Retrying, RunEvent::MaxRetriesExceeded) => Some(RunStatus::Failed),
238        (RunStatus::Retrying, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
239
240        // Approval
241        (RunStatus::Running, RunEvent::ApprovalRequested) => Some(RunStatus::AwaitingApproval),
242        (RunStatus::AwaitingApproval, RunEvent::Approved) => Some(RunStatus::Running),
243        (RunStatus::AwaitingApproval, RunEvent::Rejected) => Some(RunStatus::Failed),
244        (RunStatus::AwaitingApproval, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
245
246        // Delay
247        (RunStatus::Running, RunEvent::DelaySleeping) => Some(RunStatus::Sleeping),
248        (RunStatus::Sleeping, RunEvent::DelayElapsed) => Some(RunStatus::Pending),
249        (RunStatus::Sleeping, RunEvent::SignalReceived) => Some(RunStatus::Pending),
250        (RunStatus::Sleeping, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
251
252        // Terminal states and all other combos → invalid
253        _ => None,
254    }
255}
256
257#[cfg(test)]
258mod tests {
259    use super::*;
260
261    // ---- Happy paths ----
262
263    #[test]
264    fn pending_to_running() {
265        let mut fsm = RunFsm::new();
266        let result = fsm.apply(RunEvent::PickedUp);
267        assert!(result.is_ok());
268        assert_eq!(fsm.state(), RunStatus::Running);
269    }
270
271    #[test]
272    fn full_success_path() {
273        let mut fsm = RunFsm::new();
274        fsm.apply(RunEvent::PickedUp).unwrap();
275        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
276        assert_eq!(fsm.state(), RunStatus::Completed);
277        assert!(fsm.is_terminal());
278        assert_eq!(fsm.history().len(), 2);
279    }
280
281    #[test]
282    fn full_failure_path() {
283        let mut fsm = RunFsm::new();
284        fsm.apply(RunEvent::PickedUp).unwrap();
285        fsm.apply(RunEvent::StepFailed).unwrap();
286        assert_eq!(fsm.state(), RunStatus::Failed);
287        assert!(fsm.is_terminal());
288    }
289
290    #[test]
291    fn retry_then_success() {
292        let mut fsm = RunFsm::new();
293        fsm.apply(RunEvent::PickedUp).unwrap();
294        fsm.apply(RunEvent::StepFailedRetryable).unwrap();
295        assert_eq!(fsm.state(), RunStatus::Retrying);
296
297        fsm.apply(RunEvent::RetryStarted).unwrap();
298        assert_eq!(fsm.state(), RunStatus::Running);
299
300        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
301        assert_eq!(fsm.state(), RunStatus::Completed);
302        assert_eq!(fsm.history().len(), 4);
303    }
304
305    #[test]
306    fn retry_then_max_retries_exceeded() {
307        let mut fsm = RunFsm::new();
308        fsm.apply(RunEvent::PickedUp).unwrap();
309        fsm.apply(RunEvent::StepFailedRetryable).unwrap();
310        fsm.apply(RunEvent::MaxRetriesExceeded).unwrap();
311        assert_eq!(fsm.state(), RunStatus::Failed);
312    }
313
314    #[test]
315    fn cancel_from_pending() {
316        let mut fsm = RunFsm::new();
317        fsm.apply(RunEvent::CancelRequested).unwrap();
318        assert_eq!(fsm.state(), RunStatus::Cancelled);
319        assert!(fsm.is_terminal());
320    }
321
322    #[test]
323    fn cancel_from_running() {
324        let mut fsm = RunFsm::new();
325        fsm.apply(RunEvent::PickedUp).unwrap();
326        fsm.apply(RunEvent::CancelRequested).unwrap();
327        assert_eq!(fsm.state(), RunStatus::Cancelled);
328    }
329
330    #[test]
331    fn cancel_from_retrying() {
332        let mut fsm = RunFsm::new();
333        fsm.apply(RunEvent::PickedUp).unwrap();
334        fsm.apply(RunEvent::StepFailedRetryable).unwrap();
335        fsm.apply(RunEvent::CancelRequested).unwrap();
336        assert_eq!(fsm.state(), RunStatus::Cancelled);
337    }
338
339    // ---- Invalid transitions ----
340
341    #[test]
342    fn cannot_complete_from_pending() {
343        let mut fsm = RunFsm::new();
344        let result = fsm.apply(RunEvent::AllStepsCompleted);
345        assert!(result.is_err());
346        assert_eq!(fsm.state(), RunStatus::Pending);
347    }
348
349    #[test]
350    fn cannot_pick_up_running() {
351        let mut fsm = RunFsm::new();
352        fsm.apply(RunEvent::PickedUp).unwrap();
353        let result = fsm.apply(RunEvent::PickedUp);
354        assert!(result.is_err());
355    }
356
357    #[test]
358    fn cannot_transition_from_terminal() {
359        let mut fsm = RunFsm::new();
360        fsm.apply(RunEvent::PickedUp).unwrap();
361        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
362
363        assert!(fsm.apply(RunEvent::PickedUp).is_err());
364        assert!(fsm.apply(RunEvent::CancelRequested).is_err());
365        assert!(fsm.apply(RunEvent::StepFailed).is_err());
366    }
367
368    // ---- can_apply ----
369
370    #[test]
371    fn can_apply_checks_without_mutation() {
372        let fsm = RunFsm::new();
373        assert!(fsm.can_apply(RunEvent::PickedUp));
374        assert!(fsm.can_apply(RunEvent::CancelRequested));
375        assert!(!fsm.can_apply(RunEvent::AllStepsCompleted));
376        assert!(!fsm.can_apply(RunEvent::StepFailed));
377        assert_eq!(fsm.state(), RunStatus::Pending);
378    }
379
380    // ---- from_state ----
381
382    #[test]
383    fn from_state_resumes_at_given_state() {
384        let mut fsm = RunFsm::from_state(RunStatus::Running);
385        assert_eq!(fsm.state(), RunStatus::Running);
386        assert!(fsm.history().is_empty());
387
388        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
389        assert_eq!(fsm.state(), RunStatus::Completed);
390    }
391
392    // ---- History ----
393
394    #[test]
395    fn history_records_transitions() {
396        let mut fsm = RunFsm::new();
397        fsm.apply(RunEvent::PickedUp).unwrap();
398        fsm.apply(RunEvent::StepFailedRetryable).unwrap();
399        fsm.apply(RunEvent::RetryStarted).unwrap();
400
401        let history = fsm.history();
402        assert_eq!(history.len(), 3);
403
404        assert_eq!(history[0].from, RunStatus::Pending);
405        assert_eq!(history[0].to, RunStatus::Running);
406        assert_eq!(history[0].event, RunEvent::PickedUp);
407
408        assert_eq!(history[1].from, RunStatus::Running);
409        assert_eq!(history[1].to, RunStatus::Retrying);
410        assert_eq!(history[1].event, RunEvent::StepFailedRetryable);
411
412        assert_eq!(history[2].from, RunStatus::Retrying);
413        assert_eq!(history[2].to, RunStatus::Running);
414        assert_eq!(history[2].event, RunEvent::RetryStarted);
415    }
416
417    // ---- Approval transitions ----
418
419    #[test]
420    fn running_to_awaiting_approval() {
421        let mut fsm = RunFsm::new();
422        fsm.apply(RunEvent::PickedUp).unwrap();
423        fsm.apply(RunEvent::ApprovalRequested).unwrap();
424        assert_eq!(fsm.state(), RunStatus::AwaitingApproval);
425        assert!(!fsm.is_terminal());
426    }
427
428    #[test]
429    fn awaiting_approval_approved_resumes_running() {
430        let mut fsm = RunFsm::new();
431        fsm.apply(RunEvent::PickedUp).unwrap();
432        fsm.apply(RunEvent::ApprovalRequested).unwrap();
433        fsm.apply(RunEvent::Approved).unwrap();
434        assert_eq!(fsm.state(), RunStatus::Running);
435    }
436
437    #[test]
438    fn awaiting_approval_rejected_fails() {
439        let mut fsm = RunFsm::new();
440        fsm.apply(RunEvent::PickedUp).unwrap();
441        fsm.apply(RunEvent::ApprovalRequested).unwrap();
442        fsm.apply(RunEvent::Rejected).unwrap();
443        assert_eq!(fsm.state(), RunStatus::Failed);
444        assert!(fsm.is_terminal());
445    }
446
447    #[test]
448    fn awaiting_approval_cancel() {
449        let mut fsm = RunFsm::new();
450        fsm.apply(RunEvent::PickedUp).unwrap();
451        fsm.apply(RunEvent::ApprovalRequested).unwrap();
452        fsm.apply(RunEvent::CancelRequested).unwrap();
453        assert_eq!(fsm.state(), RunStatus::Cancelled);
454        assert!(fsm.is_terminal());
455    }
456
457    #[test]
458    fn cannot_approve_from_pending() {
459        let mut fsm = RunFsm::new();
460        assert!(fsm.apply(RunEvent::Approved).is_err());
461    }
462
463    #[test]
464    fn approval_then_complete() {
465        let mut fsm = RunFsm::new();
466        fsm.apply(RunEvent::PickedUp).unwrap();
467        fsm.apply(RunEvent::ApprovalRequested).unwrap();
468        fsm.apply(RunEvent::Approved).unwrap();
469        fsm.apply(RunEvent::AllStepsCompleted).unwrap();
470        assert_eq!(fsm.state(), RunStatus::Completed);
471        assert_eq!(fsm.history().len(), 4);
472    }
473
474    #[test]
475    fn sleeping_signal_received_goes_pending() {
476        let mut fsm = RunFsm::new();
477        fsm.apply(RunEvent::PickedUp).unwrap();
478        fsm.apply(RunEvent::DelaySleeping).unwrap();
479        fsm.apply(RunEvent::SignalReceived).unwrap();
480        assert_eq!(fsm.state(), RunStatus::Pending);
481        assert_eq!(RunEvent::SignalReceived.to_string(), "signal_received");
482    }
483
484    #[test]
485    fn cannot_receive_signal_while_running() {
486        let mut fsm = RunFsm::new();
487        fsm.apply(RunEvent::PickedUp).unwrap();
488        assert!(fsm.apply(RunEvent::SignalReceived).is_err());
489    }
490
491    // ---- TransitionError Display ----
492
493    #[test]
494    fn transition_error_display() {
495        let mut fsm = RunFsm::new();
496        let err = fsm.apply(RunEvent::AllStepsCompleted).unwrap_err();
497        let msg = err.to_string();
498        assert!(msg.contains("all_steps_completed"));
499        assert!(msg.contains("Pending"));
500    }
501}