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}
56
57#[derive(Debug, Clone)]
101pub struct RunFsm {
102 state: RunStatus,
103 history: Vec<Transition<RunStatus, RunEvent>>,
104}
105
106impl RunFsm {
107 pub fn new() -> Self {
119 Self {
120 state: RunStatus::Pending,
121 history: Vec::new(),
122 }
123 }
124
125 pub fn from_state(state: RunStatus) -> Self {
137 Self {
138 state,
139 history: Vec::new(),
140 }
141 }
142
143 pub fn state(&self) -> RunStatus {
145 self.state
146 }
147
148 pub fn history(&self) -> &[Transition<RunStatus, RunEvent>] {
150 &self.history
151 }
152
153 pub fn is_terminal(&self) -> bool {
155 self.state.is_terminal()
156 }
157
158 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 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
221fn next_state(from: RunStatus, event: RunEvent) -> Option<RunStatus> {
224 match (from, event) {
225 (RunStatus::Pending, RunEvent::PickedUp) => Some(RunStatus::Running),
227 (RunStatus::Pending, RunEvent::CancelRequested) => Some(RunStatus::Cancelled),
228
229 (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 (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 (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 (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 _ => None,
254 }
255}
256
257#[cfg(test)]
258mod tests {
259 use super::*;
260
261 #[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 #[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 #[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 #[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 #[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 #[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 #[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}