1use chrono::Utc;
7use ironflow_store::entities::RunStatus;
8use serde::{Deserialize, Serialize};
9use strum::Display;
10
11use super::{Transition, TransitionError};
12
13#[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 PickedUp,
31 AllStepsCompleted,
33 StepFailed,
35 StepFailedRetryable,
37 RetryStarted,
39 MaxRetriesExceeded,
41 CancelRequested,
43 ApprovalRequested,
45 Approved,
47 Rejected,
49 DelaySleeping,
51 DelayElapsed,
53 SignalReceived,
55 PauseRequested,
57 ResumeRequested,
62}
63
64#[derive(Debug, Clone)]
111pub struct RunFsm {
112 state: RunStatus,
113 history: Vec<Transition<RunStatus, RunEvent>>,
114}
115
116impl RunFsm {
117 pub fn new() -> Self {
129 Self {
130 state: RunStatus::Pending,
131 history: Vec::new(),
132 }
133 }
134
135 pub fn from_state(state: RunStatus) -> Self {
147 Self {
148 state,
149 history: Vec::new(),
150 }
151 }
152
153 pub fn state(&self) -> RunStatus {
155 self.state
156 }
157
158 pub fn history(&self) -> &[Transition<RunStatus, RunEvent>] {
160 &self.history
161 }
162
163 pub fn is_terminal(&self) -> bool {
165 self.state.is_terminal()
166 }
167
168 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 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
231fn next_state(from: RunStatus, event: RunEvent) -> Option<RunStatus> {
234 match (from, event) {
235 (RunStatus::Pending, RunEvent::PickedUp) => Some(RunStatus::Running),
237 (RunStatus::Pending, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
238
239 (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 (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 (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 (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 (
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 _ => None,
276 }
277}
278
279#[cfg(test)]
280mod tests {
281 use super::*;
282
283 #[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 #[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 #[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 #[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 #[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 #[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 #[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}