1use chrono::{DateTime, Utc};
4use rust_decimal::Decimal;
5use thiserror::Error;
6use uuid::Uuid;
7
8use ironflow_artifacts::error::ArtifactError;
9use ironflow_core::error::OperationError;
10use ironflow_store::error::StoreError;
11use ironflow_store::models::{ConcurrencyLimitError, RunStatus};
12
13use crate::guard::{WORKFLOW_GUARD_REJECTED_CODE, WorkflowRejection};
14
15pub const RUN_BUDGET_EXCEEDED_CODE: &str = "RUN_BUDGET_EXCEEDED";
17
18pub const MONTHLY_BUDGET_EXCEEDED_CODE: &str = "MONTHLY_BUDGET_EXCEEDED";
20
21pub const CONCURRENCY_CONFLICT_CODE: &str = "CONCURRENCY_CONFLICT";
31
32pub const HANDLER_VERSION_MISMATCH_CODE: &str = "HANDLER_VERSION_MISMATCH";
34
35#[derive(Debug, Error)]
37pub enum EngineError {
38 #[error("operation failed: {0}")]
40 Operation(#[from] OperationError),
41
42 #[error("store error: {0}")]
44 Store(StoreError),
45
46 #[error("concurrency key {key:?} is held by active run {run_id}")]
51 ConcurrencyConflict {
52 key: String,
54 run_id: Uuid,
56 },
57
58 #[error("invalid concurrency limit: {0}")]
65 InvalidConcurrencyLimit(ConcurrencyLimitError),
66
67 #[error("invalid workflow: {0}")]
69 InvalidWorkflow(String),
70
71 #[error("step config error: {0}")]
73 StepConfig(String),
74
75 #[error("decision error: {0}")]
77 Decision(#[from] ironflow_core::error::DecisionError),
78
79 #[error(
82 "decision step '{step}' requires a decision provider; \
83 wire one with Engine::with_decision_provider(...)"
84 )]
85 NoDecisionProvider {
86 step: String,
88 },
89
90 #[error("serialization error: {0}")]
92 Serialization(#[from] serde_json::Error),
93
94 #[error(
100 "{RUN_BUDGET_EXCEEDED_CODE}: run {run_id} would exceed its cost cap \
101 (spent {spent_usd} USD + next step {step_budget_usd} USD > cap {limit_usd} USD)"
102 )]
103 RunBudgetExceeded {
104 run_id: uuid::Uuid,
106 limit_usd: Decimal,
108 spent_usd: Decimal,
110 step_budget_usd: Decimal,
112 },
113
114 #[error(
118 "{MONTHLY_BUDGET_EXCEEDED_CODE}: monthly cost quota exhausted \
119 ({spent_usd} USD spent of {limit_usd} USD)"
120 )]
121 MonthlyBudgetExceeded {
122 limit_usd: Decimal,
124 spent_usd: Decimal,
126 },
127
128 #[error("step {step:?} declared output {pattern:?} but no file matched")]
134 MissingArtifact {
135 step: String,
137 pattern: String,
139 },
140
141 #[error("step {step:?} declares no artifact output named {name:?}")]
143 ArtifactNotDeclared {
144 step: String,
146 name: String,
148 },
149
150 #[error("no artifact {name:?} produced by step {step:?} before this point")]
152 ArtifactNotFound {
153 step: String,
155 name: String,
157 },
158
159 #[error("artifact storage is not configured: {0}")]
161 ArtifactsUnavailable(String),
162
163 #[error("artifact storage error: {0}")]
165 Artifact(#[from] ArtifactError),
166
167 #[error("approval required for run {run_id}, step {step_id}: {message}")]
169 ApprovalRequired {
170 run_id: uuid::Uuid,
172 step_id: uuid::Uuid,
174 message: String,
176 },
177
178 #[error("approval rejected for run {run_id}, step {step_id}: {reason}")]
185 ApprovalRejected {
186 run_id: uuid::Uuid,
188 step_id: uuid::Uuid,
190 reason: String,
192 },
193
194 #[error("human input required for run {run_id}, step {step_id}: {message}")]
200 HumanInputRequired {
201 run_id: uuid::Uuid,
203 step_id: uuid::Uuid,
205 message: String,
207 },
208
209 #[error("human input rejected for run {run_id}, step {step_id}: {reason}")]
216 HumanInputRejected {
217 run_id: uuid::Uuid,
219 step_id: uuid::Uuid,
221 reason: String,
223 },
224
225 #[error("delay sleeping for run {run_id}, step {step_id}: wake at {wake_at}")]
231 DelaySleeping {
232 run_id: uuid::Uuid,
234 step_id: uuid::Uuid,
236 wake_at: chrono::DateTime<chrono::Utc>,
238 },
239
240 #[error("run {run_id} waiting for signal {name:?} with key {key:?} until {deadline_at}")]
248 SignalWaiting {
249 run_id: Uuid,
251 step_id: Uuid,
253 step_name: String,
255 name: String,
257 key: String,
259 deadline_at: DateTime<Utc>,
261 },
262
263 #[error("child run {run_id} suspended: {cause}")]
289 ChildSuspended {
290 run_id: Uuid,
292 cause: Box<EngineError>,
299 },
300
301 #[error("child run {run_id} was cancelled")]
319 ChildRunCancelled {
320 run_id: Uuid,
322 },
323
324 #[error("invalid signal: {0}")]
327 InvalidSignal(String),
328
329 #[error("{WORKFLOW_GUARD_REJECTED_CODE}: {0}")]
335 WorkflowGuardRejected(#[from] WorkflowRejection),
336
337 #[error(
348 "replay divergence at position {position}: handler called '{expected}' but the run \
349 recorded '{recorded}' (handler changed since the run was suspended?)"
350 )]
351 ReplayDivergence {
352 position: u32,
355 expected: String,
357 recorded: String,
359 },
360
361 #[error(
371 "{HANDLER_VERSION_MISMATCH_CODE}: run {run_id} for handler '{workflow_name}' was created \
372 with version {run_version}, but the handler is now at version {current_version}; \
373 resume refused (no force override for resume)"
374 )]
375 HandlerVersionMismatch {
376 run_id: uuid::Uuid,
378 workflow_name: String,
380 run_version: String,
382 current_version: String,
384 },
385}
386
387impl From<StoreError> for EngineError {
388 fn from(err: StoreError) -> Self {
389 match err {
390 StoreError::ConcurrencyConflict { key, run_id } => {
391 EngineError::ConcurrencyConflict { key, run_id }
392 }
393 StoreError::InvalidConcurrencyLimit(e) => EngineError::InvalidConcurrencyLimit(e),
394 other => EngineError::Store(other),
395 }
396 }
397}
398
399impl EngineError {
400 pub fn is_suspension(&self) -> bool {
416 matches!(
417 self,
418 EngineError::ApprovalRequired { .. }
419 | EngineError::HumanInputRequired { .. }
420 | EngineError::DelaySleeping { .. }
421 | EngineError::SignalWaiting { .. }
422 | EngineError::ChildSuspended { .. }
423 )
424 }
425
426 pub fn suspension_leaf(&self) -> &EngineError {
448 let mut current = self;
449 while let EngineError::ChildSuspended { cause, .. } = current {
450 current = cause;
451 }
452 current
453 }
454
455 pub(crate) fn suspension_status(&self) -> RunStatus {
459 match self.suspension_leaf() {
460 EngineError::DelaySleeping { .. } | EngineError::SignalWaiting { .. } => {
461 RunStatus::Sleeping
462 }
463 _ => RunStatus::AwaitingApproval,
464 }
465 }
466}
467
468#[cfg(test)]
469mod tests {
470 use super::*;
471 use crate::retry_policy::is_run_retryable;
472
473 #[test]
474 fn invalid_workflow_display() {
475 let err = EngineError::InvalidWorkflow("unknown-handler".to_string());
476 assert!(err.to_string().contains("invalid workflow"));
477 assert!(err.to_string().contains("unknown-handler"));
478 }
479
480 #[test]
481 fn step_config_display() {
482 let err = EngineError::StepConfig("bad shell config".to_string());
483 assert!(err.to_string().contains("step config error"));
484 assert!(err.to_string().contains("bad shell config"));
485 }
486
487 #[test]
488 fn human_input_required_display() {
489 let err = EngineError::HumanInputRequired {
490 run_id: uuid::Uuid::nil(),
491 step_id: uuid::Uuid::nil(),
492 message: "Answer the questions".to_string(),
493 };
494 let text = err.to_string();
495 assert!(text.contains("human input required"));
496 assert!(text.contains("Answer the questions"));
497 }
498
499 #[test]
500 fn human_input_rejected_display() {
501 let err = EngineError::HumanInputRejected {
502 run_id: uuid::Uuid::nil(),
503 step_id: uuid::Uuid::nil(),
504 reason: "not relevant".to_string(),
505 };
506 let text = err.to_string();
507 assert!(text.contains("human input rejected"));
508 assert!(text.contains("not relevant"));
509 }
510
511 #[test]
512 fn store_error_from_conversion() {
513 let store_err = StoreError::RunNotFound(uuid::Uuid::nil());
514 let engine_err = EngineError::from(store_err);
515 assert!(engine_err.to_string().contains("store error"));
516 }
517
518 #[test]
519 fn store_concurrency_conflict_converts_to_engine_variant() {
520 let run_id = Uuid::now_v7();
521 let engine_err = EngineError::from(StoreError::ConcurrencyConflict {
522 key: "issue:12".to_string(),
523 run_id,
524 });
525 match engine_err {
526 EngineError::ConcurrencyConflict {
527 ref key,
528 run_id: holder,
529 } => {
530 assert_eq!(key, "issue:12");
531 assert_eq!(holder, run_id);
532 }
533 ref other => panic!("expected ConcurrencyConflict, got {other:?}"),
534 }
535 assert!(engine_err.to_string().contains("issue:12"));
536 assert!(!is_run_retryable(&engine_err));
537 }
538
539 #[test]
540 fn store_invalid_concurrency_limit_converts_to_engine_variant() {
541 let engine_err = EngineError::from(StoreError::InvalidConcurrencyLimit(
542 ConcurrencyLimitError::ZeroLimit {
543 group: "repo:acme".to_string(),
544 },
545 ));
546 assert!(
547 matches!(
548 engine_err,
549 EngineError::InvalidConcurrencyLimit(ConcurrencyLimitError::ZeroLimit { .. })
550 ),
551 "{engine_err:?}"
552 );
553 assert!(engine_err.to_string().contains("repo:acme"));
554 assert!(!is_run_retryable(&engine_err));
555 }
556
557 #[test]
558 fn run_budget_exceeded_display_carries_code_and_amounts() {
559 let err = EngineError::RunBudgetExceeded {
560 run_id: uuid::Uuid::nil(),
561 limit_usd: Decimal::new(200, 2),
562 spent_usd: Decimal::new(180, 2),
563 step_budget_usd: Decimal::new(50, 2),
564 };
565
566 let msg = err.to_string();
567 assert!(msg.contains(RUN_BUDGET_EXCEEDED_CODE));
568 assert!(msg.contains("2.00"));
569 assert!(msg.contains("1.80"));
570 assert!(msg.contains("0.50"));
571 }
572
573 #[test]
574 fn monthly_budget_exceeded_display_carries_code_and_amounts() {
575 let err = EngineError::MonthlyBudgetExceeded {
576 limit_usd: Decimal::new(10000, 2),
577 spent_usd: Decimal::new(10500, 2),
578 };
579
580 let msg = err.to_string();
581 assert!(msg.contains(MONTHLY_BUDGET_EXCEEDED_CODE));
582 assert!(msg.contains("100.00"));
583 assert!(msg.contains("105.00"));
584 }
585
586 #[test]
587 fn missing_artifact_display_names_the_step_and_pattern() {
588 let err = EngineError::MissingArtifact {
589 step: "build".to_string(),
590 pattern: "target/report.html".to_string(),
591 };
592
593 let msg = err.to_string();
594 assert!(msg.contains("\"build\""));
595 assert!(msg.contains("target/report.html"));
596 }
597
598 #[test]
599 fn artifact_not_declared_display_names_the_step_and_artifact() {
600 let err = EngineError::ArtifactNotDeclared {
601 step: "build".to_string(),
602 name: "report.htm".to_string(),
603 };
604
605 let msg = err.to_string();
606 assert!(msg.contains("\"build\""));
607 assert!(msg.contains("\"report.htm\""));
608 }
609
610 #[test]
611 fn artifact_not_found_display_names_the_producer() {
612 let err = EngineError::ArtifactNotFound {
613 step: "build".to_string(),
614 name: "report.html".to_string(),
615 };
616
617 let msg = err.to_string();
618 assert!(msg.contains("\"build\""));
619 assert!(msg.contains("report.html"));
620 }
621
622 #[test]
623 fn artifacts_unavailable_display() {
624 let err = EngineError::ArtifactsUnavailable("no blob store".to_string());
625 assert!(err.to_string().contains("not configured"));
626 }
627
628 #[test]
629 fn artifact_error_from_conversion() {
630 let engine_err = EngineError::from(ArtifactError::NotFound("a/b".to_string()));
631 assert!(engine_err.to_string().contains("artifact storage error"));
632 }
633
634 #[test]
635 fn serialization_error_from_conversion() {
636 let serde_err = serde_json::from_str::<String>("not json").unwrap_err();
637 let engine_err = EngineError::from(serde_err);
638 assert!(engine_err.to_string().contains("serialization error"));
639 }
640
641 #[test]
642 fn workflow_guard_rejected_display_carries_code_and_detail() {
643 use crate::guard::WorkflowRejection;
644
645 let rejection = WorkflowRejection::MaxDepthExceeded { depth: 6, max: 5 };
646 let err = EngineError::from(rejection);
647
648 let msg = err.to_string();
649 assert!(msg.contains(WORKFLOW_GUARD_REJECTED_CODE));
650 assert!(msg.contains("max call depth exceeded"));
651 assert!(msg.contains("6/5"));
652 }
653
654 #[test]
655 fn workflow_guard_rejected_from_conversion() {
656 use crate::guard::WorkflowRejection;
657
658 let rejection = WorkflowRejection::CycleDetected {
659 target: "wf-b".to_string(),
660 chain: vec!["wf-a".to_string(), "wf-b".to_string()],
661 };
662 let engine_err = EngineError::from(rejection);
663 assert!(engine_err.to_string().contains("cycle detected"));
664 }
665
666 #[test]
667 fn replay_divergence_display_carries_position_and_identities() {
668 let err = EngineError::ReplayDivergence {
669 position: 5,
670 expected: "resolve-base-branch (Shell)".to_string(),
671 recorded: "create-worktree (Shell)".to_string(),
672 };
673
674 let msg = err.to_string();
675 assert!(msg.contains("divergence"));
676 assert!(msg.contains("position 5"));
677 assert!(msg.contains("resolve-base-branch"));
678 assert!(msg.contains("create-worktree"));
679 }
680
681 #[test]
682 fn handler_version_mismatch_display_carries_code_and_versions() {
683 let err = EngineError::HandlerVersionMismatch {
684 run_id: uuid::Uuid::nil(),
685 workflow_name: "deploy".to_string(),
686 run_version: "1.0.0".to_string(),
687 current_version: "2.0.0".to_string(),
688 };
689
690 let msg = err.to_string();
691 assert!(msg.contains(HANDLER_VERSION_MISMATCH_CODE));
692 assert!(msg.contains("1.0.0"));
693 assert!(msg.contains("2.0.0"));
694 }
695
696 fn human_input_required() -> EngineError {
697 EngineError::HumanInputRequired {
698 run_id: Uuid::nil(),
699 step_id: Uuid::nil(),
700 message: "Answer the questions".to_string(),
701 }
702 }
703
704 #[test]
705 fn child_suspended_display_carries_child_and_cause() {
706 let child = Uuid::now_v7();
707 let err = EngineError::ChildSuspended {
708 run_id: child,
709 cause: Box::new(human_input_required()),
710 };
711
712 let msg = err.to_string();
713 assert!(msg.contains(&child.to_string()));
714 assert!(msg.contains("human input required"));
715 }
716
717 #[test]
718 fn leaf_suspensions_and_child_suspended_are_suspensions() {
719 let wake_at = Utc::now();
720 let suspensions = [
721 EngineError::ApprovalRequired {
722 run_id: Uuid::nil(),
723 step_id: Uuid::nil(),
724 message: "deploy?".to_string(),
725 },
726 human_input_required(),
727 EngineError::DelaySleeping {
728 run_id: Uuid::nil(),
729 step_id: Uuid::nil(),
730 wake_at,
731 },
732 EngineError::SignalWaiting {
733 run_id: Uuid::nil(),
734 step_id: Uuid::nil(),
735 step_name: "wait".to_string(),
736 name: "payment".to_string(),
737 key: "order-1".to_string(),
738 deadline_at: wake_at,
739 },
740 EngineError::ChildSuspended {
741 run_id: Uuid::nil(),
742 cause: Box::new(human_input_required()),
743 },
744 ];
745 for err in &suspensions {
746 assert!(err.is_suspension(), "{err} should be a suspension");
747 }
748 }
749
750 #[test]
751 fn failures_and_rejections_are_not_suspensions() {
752 let failures = [
753 EngineError::InvalidWorkflow("x".to_string()),
754 EngineError::HumanInputRejected {
755 run_id: Uuid::nil(),
756 step_id: Uuid::nil(),
757 reason: "no".to_string(),
758 },
759 EngineError::ApprovalRejected {
760 run_id: Uuid::nil(),
761 step_id: Uuid::nil(),
762 reason: "no".to_string(),
763 },
764 ];
765 for err in &failures {
766 assert!(!err.is_suspension(), "{err} should not be a suspension");
767 }
768 }
769
770 #[test]
771 fn suspension_leaf_unwraps_nested_child_suspensions() {
772 let err = EngineError::ChildSuspended {
773 run_id: Uuid::now_v7(),
774 cause: Box::new(EngineError::ChildSuspended {
775 run_id: Uuid::now_v7(),
776 cause: Box::new(human_input_required()),
777 }),
778 };
779
780 assert!(matches!(
781 err.suspension_leaf(),
782 EngineError::HumanInputRequired { .. }
783 ));
784 }
785
786 #[test]
787 fn suspension_status_follows_the_leaf() {
788 let human = EngineError::ChildSuspended {
789 run_id: Uuid::nil(),
790 cause: Box::new(human_input_required()),
791 };
792 assert_eq!(human.suspension_status(), RunStatus::AwaitingApproval);
793
794 let delay = EngineError::ChildSuspended {
795 run_id: Uuid::nil(),
796 cause: Box::new(EngineError::DelaySleeping {
797 run_id: Uuid::nil(),
798 step_id: Uuid::nil(),
799 wake_at: Utc::now(),
800 }),
801 };
802 assert_eq!(delay.suspension_status(), RunStatus::Sleeping);
803 }
804
805 #[test]
806 fn suspension_leaf_of_a_plain_error_is_itself() {
807 let err = EngineError::StepConfig("bad".to_string());
808 assert!(matches!(err.suspension_leaf(), EngineError::StepConfig(_)));
809 }
810}