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, WorkerTagError};
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 priority: {0}")]
74 InvalidPriority(String),
75 #[error("invalid worker tag: {0}")]
82 InvalidWorkerTag(WorkerTagError),
83
84 #[error("invalid workflow: {0}")]
86 InvalidWorkflow(String),
87
88 #[error("step config error: {0}")]
90 StepConfig(String),
91
92 #[error("decision error: {0}")]
94 Decision(#[from] ironflow_core::error::DecisionError),
95
96 #[error(
99 "decision step '{step}' requires a decision provider; \
100 wire one with Engine::with_decision_provider(...)"
101 )]
102 NoDecisionProvider {
103 step: String,
105 },
106
107 #[error("serialization error: {0}")]
109 Serialization(#[from] serde_json::Error),
110
111 #[error(
117 "{RUN_BUDGET_EXCEEDED_CODE}: run {run_id} would exceed its cost cap \
118 (spent {spent_usd} USD + next step {step_budget_usd} USD > cap {limit_usd} USD)"
119 )]
120 RunBudgetExceeded {
121 run_id: uuid::Uuid,
123 limit_usd: Decimal,
125 spent_usd: Decimal,
127 step_budget_usd: Decimal,
129 },
130
131 #[error(
135 "{MONTHLY_BUDGET_EXCEEDED_CODE}: monthly cost quota exhausted \
136 ({spent_usd} USD spent of {limit_usd} USD)"
137 )]
138 MonthlyBudgetExceeded {
139 limit_usd: Decimal,
141 spent_usd: Decimal,
143 },
144
145 #[error("step {step:?} declared output {pattern:?} but no file matched")]
151 MissingArtifact {
152 step: String,
154 pattern: String,
156 },
157
158 #[error("step {step:?} declares no artifact output named {name:?}")]
160 ArtifactNotDeclared {
161 step: String,
163 name: String,
165 },
166
167 #[error("no artifact {name:?} produced by step {step:?} before this point")]
169 ArtifactNotFound {
170 step: String,
172 name: String,
174 },
175
176 #[error("artifact storage is not configured: {0}")]
178 ArtifactsUnavailable(String),
179
180 #[error("artifact storage error: {0}")]
182 Artifact(#[from] ArtifactError),
183
184 #[error("approval required for run {run_id}, step {step_id}: {message}")]
186 ApprovalRequired {
187 run_id: uuid::Uuid,
189 step_id: uuid::Uuid,
191 message: String,
193 },
194
195 #[error("approval rejected for run {run_id}, step {step_id}: {reason}")]
202 ApprovalRejected {
203 run_id: uuid::Uuid,
205 step_id: uuid::Uuid,
207 reason: String,
209 },
210
211 #[error("human input required for run {run_id}, step {step_id}: {message}")]
217 HumanInputRequired {
218 run_id: uuid::Uuid,
220 step_id: uuid::Uuid,
222 message: String,
224 },
225
226 #[error("human input rejected for run {run_id}, step {step_id}: {reason}")]
233 HumanInputRejected {
234 run_id: uuid::Uuid,
236 step_id: uuid::Uuid,
238 reason: String,
240 },
241
242 #[error("delay sleeping for run {run_id}, step {step_id}: wake at {wake_at}")]
248 DelaySleeping {
249 run_id: uuid::Uuid,
251 step_id: uuid::Uuid,
253 wake_at: chrono::DateTime<chrono::Utc>,
255 },
256
257 #[error("run {run_id} waiting for {kind} capacity at step {step_id}: wake at {wake_at}")]
283 CapacitySleeping {
284 run_id: Uuid,
286 step_id: Uuid,
288 kind: String,
290 wake_at: DateTime<Utc>,
292 },
293
294 #[error("run {run_id} waiting for signal {name:?} with key {key:?} until {deadline_at}")]
302 SignalWaiting {
303 run_id: Uuid,
305 step_id: Uuid,
307 step_name: String,
309 name: String,
311 key: String,
313 deadline_at: DateTime<Utc>,
315 },
316
317 #[error("child run {run_id} suspended: {cause}")]
343 ChildSuspended {
344 run_id: Uuid,
346 cause: Box<EngineError>,
354 },
355
356 #[error("child run {run_id} was cancelled")]
374 ChildRunCancelled {
375 run_id: Uuid,
377 },
378
379 #[error("invalid signal: {0}")]
382 InvalidSignal(String),
383
384 #[error("{WORKFLOW_GUARD_REJECTED_CODE}: {0}")]
390 WorkflowGuardRejected(#[from] WorkflowRejection),
391
392 #[error(
403 "replay divergence at position {position}: handler called '{expected}' but the run \
404 recorded '{recorded}' (handler changed since the run was suspended?)"
405 )]
406 ReplayDivergence {
407 position: u32,
410 expected: String,
412 recorded: String,
414 },
415
416 #[error(
426 "{HANDLER_VERSION_MISMATCH_CODE}: run {run_id} for handler '{workflow_name}' was created \
427 with version {run_version}, but the handler is now at version {current_version}; \
428 resume refused (no force override for resume)"
429 )]
430 HandlerVersionMismatch {
431 run_id: uuid::Uuid,
433 workflow_name: String,
435 run_version: String,
437 current_version: String,
439 },
440}
441
442impl From<StoreError> for EngineError {
443 fn from(err: StoreError) -> Self {
444 match err {
445 StoreError::ConcurrencyConflict { key, run_id } => {
446 EngineError::ConcurrencyConflict { key, run_id }
447 }
448 StoreError::InvalidConcurrencyLimit(e) => EngineError::InvalidConcurrencyLimit(e),
449 StoreError::InvalidWorkerTag(e) => EngineError::InvalidWorkerTag(e),
450 other => EngineError::Store(other),
451 }
452 }
453}
454
455impl EngineError {
456 pub fn is_suspension(&self) -> bool {
473 matches!(
474 self,
475 EngineError::ApprovalRequired { .. }
476 | EngineError::HumanInputRequired { .. }
477 | EngineError::DelaySleeping { .. }
478 | EngineError::CapacitySleeping { .. }
479 | EngineError::SignalWaiting { .. }
480 | EngineError::ChildSuspended { .. }
481 )
482 }
483
484 pub fn suspension_leaf(&self) -> &EngineError {
506 let mut current = self;
507 while let EngineError::ChildSuspended { cause, .. } = current {
508 current = cause;
509 }
510 current
511 }
512
513 pub(crate) fn suspension_status(&self) -> RunStatus {
517 match self.suspension_leaf() {
518 EngineError::DelaySleeping { .. }
519 | EngineError::CapacitySleeping { .. }
520 | EngineError::SignalWaiting { .. } => RunStatus::Sleeping,
521 _ => RunStatus::AwaitingApproval,
522 }
523 }
524}
525
526#[cfg(test)]
527mod tests {
528 use super::*;
529 use crate::retry_policy::is_run_retryable;
530
531 #[test]
532 fn invalid_workflow_display() {
533 let err = EngineError::InvalidWorkflow("unknown-handler".to_string());
534 assert!(err.to_string().contains("invalid workflow"));
535 assert!(err.to_string().contains("unknown-handler"));
536 }
537
538 #[test]
539 fn step_config_display() {
540 let err = EngineError::StepConfig("bad shell config".to_string());
541 assert!(err.to_string().contains("step config error"));
542 assert!(err.to_string().contains("bad shell config"));
543 }
544
545 #[test]
546 fn human_input_required_display() {
547 let err = EngineError::HumanInputRequired {
548 run_id: uuid::Uuid::nil(),
549 step_id: uuid::Uuid::nil(),
550 message: "Answer the questions".to_string(),
551 };
552 let text = err.to_string();
553 assert!(text.contains("human input required"));
554 assert!(text.contains("Answer the questions"));
555 }
556
557 #[test]
558 fn human_input_rejected_display() {
559 let err = EngineError::HumanInputRejected {
560 run_id: uuid::Uuid::nil(),
561 step_id: uuid::Uuid::nil(),
562 reason: "not relevant".to_string(),
563 };
564 let text = err.to_string();
565 assert!(text.contains("human input rejected"));
566 assert!(text.contains("not relevant"));
567 }
568
569 #[test]
570 fn store_error_from_conversion() {
571 let store_err = StoreError::RunNotFound(uuid::Uuid::nil());
572 let engine_err = EngineError::from(store_err);
573 assert!(engine_err.to_string().contains("store error"));
574 }
575
576 #[test]
577 fn store_concurrency_conflict_converts_to_engine_variant() {
578 let run_id = Uuid::now_v7();
579 let engine_err = EngineError::from(StoreError::ConcurrencyConflict {
580 key: "issue:12".to_string(),
581 run_id,
582 });
583 match engine_err {
584 EngineError::ConcurrencyConflict {
585 ref key,
586 run_id: holder,
587 } => {
588 assert_eq!(key, "issue:12");
589 assert_eq!(holder, run_id);
590 }
591 ref other => panic!("expected ConcurrencyConflict, got {other:?}"),
592 }
593 assert!(engine_err.to_string().contains("issue:12"));
594 assert!(!is_run_retryable(&engine_err));
595 }
596
597 #[test]
598 fn store_invalid_concurrency_limit_converts_to_engine_variant() {
599 let engine_err = EngineError::from(StoreError::InvalidConcurrencyLimit(
600 ConcurrencyLimitError::ZeroLimit {
601 group: "repo:acme".to_string(),
602 },
603 ));
604 assert!(
605 matches!(
606 engine_err,
607 EngineError::InvalidConcurrencyLimit(ConcurrencyLimitError::ZeroLimit { .. })
608 ),
609 "{engine_err:?}"
610 );
611 assert!(engine_err.to_string().contains("repo:acme"));
612 assert!(!is_run_retryable(&engine_err));
613 }
614
615 #[test]
616 fn invalid_priority_is_not_retryable() {
617 let err = EngineError::InvalidPriority("priority must be between -100 and 100".to_string());
618 assert_eq!(
619 err.to_string(),
620 "invalid priority: priority must be between -100 and 100"
621 );
622 assert!(!is_run_retryable(&err));
623 }
624
625 #[test]
626 fn store_invalid_worker_tag_converts_to_engine_variant() {
627 let engine_err =
628 EngineError::from(StoreError::InvalidWorkerTag(WorkerTagError::InvalidChar {
629 tag: "bad,tag".to_string(),
630 }));
631 assert!(
632 matches!(
633 engine_err,
634 EngineError::InvalidWorkerTag(WorkerTagError::InvalidChar { .. })
635 ),
636 "{engine_err:?}"
637 );
638 assert!(engine_err.to_string().contains("bad,tag"));
639 assert!(!is_run_retryable(&engine_err));
640 }
641
642 #[test]
643 fn run_budget_exceeded_display_carries_code_and_amounts() {
644 let err = EngineError::RunBudgetExceeded {
645 run_id: uuid::Uuid::nil(),
646 limit_usd: Decimal::new(200, 2),
647 spent_usd: Decimal::new(180, 2),
648 step_budget_usd: Decimal::new(50, 2),
649 };
650
651 let msg = err.to_string();
652 assert!(msg.contains(RUN_BUDGET_EXCEEDED_CODE));
653 assert!(msg.contains("2.00"));
654 assert!(msg.contains("1.80"));
655 assert!(msg.contains("0.50"));
656 }
657
658 #[test]
659 fn monthly_budget_exceeded_display_carries_code_and_amounts() {
660 let err = EngineError::MonthlyBudgetExceeded {
661 limit_usd: Decimal::new(10000, 2),
662 spent_usd: Decimal::new(10500, 2),
663 };
664
665 let msg = err.to_string();
666 assert!(msg.contains(MONTHLY_BUDGET_EXCEEDED_CODE));
667 assert!(msg.contains("100.00"));
668 assert!(msg.contains("105.00"));
669 }
670
671 #[test]
672 fn missing_artifact_display_names_the_step_and_pattern() {
673 let err = EngineError::MissingArtifact {
674 step: "build".to_string(),
675 pattern: "target/report.html".to_string(),
676 };
677
678 let msg = err.to_string();
679 assert!(msg.contains("\"build\""));
680 assert!(msg.contains("target/report.html"));
681 }
682
683 #[test]
684 fn artifact_not_declared_display_names_the_step_and_artifact() {
685 let err = EngineError::ArtifactNotDeclared {
686 step: "build".to_string(),
687 name: "report.htm".to_string(),
688 };
689
690 let msg = err.to_string();
691 assert!(msg.contains("\"build\""));
692 assert!(msg.contains("\"report.htm\""));
693 }
694
695 #[test]
696 fn artifact_not_found_display_names_the_producer() {
697 let err = EngineError::ArtifactNotFound {
698 step: "build".to_string(),
699 name: "report.html".to_string(),
700 };
701
702 let msg = err.to_string();
703 assert!(msg.contains("\"build\""));
704 assert!(msg.contains("report.html"));
705 }
706
707 #[test]
708 fn artifacts_unavailable_display() {
709 let err = EngineError::ArtifactsUnavailable("no blob store".to_string());
710 assert!(err.to_string().contains("not configured"));
711 }
712
713 #[test]
714 fn artifact_error_from_conversion() {
715 let engine_err = EngineError::from(ArtifactError::NotFound("a/b".to_string()));
716 assert!(engine_err.to_string().contains("artifact storage error"));
717 }
718
719 #[test]
720 fn serialization_error_from_conversion() {
721 let serde_err = serde_json::from_str::<String>("not json").unwrap_err();
722 let engine_err = EngineError::from(serde_err);
723 assert!(engine_err.to_string().contains("serialization error"));
724 }
725
726 #[test]
727 fn workflow_guard_rejected_display_carries_code_and_detail() {
728 use crate::guard::WorkflowRejection;
729
730 let rejection = WorkflowRejection::MaxDepthExceeded { depth: 6, max: 5 };
731 let err = EngineError::from(rejection);
732
733 let msg = err.to_string();
734 assert!(msg.contains(WORKFLOW_GUARD_REJECTED_CODE));
735 assert!(msg.contains("max call depth exceeded"));
736 assert!(msg.contains("6/5"));
737 }
738
739 #[test]
740 fn workflow_guard_rejected_from_conversion() {
741 use crate::guard::WorkflowRejection;
742
743 let rejection = WorkflowRejection::CycleDetected {
744 target: "wf-b".to_string(),
745 chain: vec!["wf-a".to_string(), "wf-b".to_string()],
746 };
747 let engine_err = EngineError::from(rejection);
748 assert!(engine_err.to_string().contains("cycle detected"));
749 }
750
751 #[test]
752 fn replay_divergence_display_carries_position_and_identities() {
753 let err = EngineError::ReplayDivergence {
754 position: 5,
755 expected: "resolve-base-branch (Shell)".to_string(),
756 recorded: "create-worktree (Shell)".to_string(),
757 };
758
759 let msg = err.to_string();
760 assert!(msg.contains("divergence"));
761 assert!(msg.contains("position 5"));
762 assert!(msg.contains("resolve-base-branch"));
763 assert!(msg.contains("create-worktree"));
764 }
765
766 #[test]
767 fn handler_version_mismatch_display_carries_code_and_versions() {
768 let err = EngineError::HandlerVersionMismatch {
769 run_id: uuid::Uuid::nil(),
770 workflow_name: "deploy".to_string(),
771 run_version: "1.0.0".to_string(),
772 current_version: "2.0.0".to_string(),
773 };
774
775 let msg = err.to_string();
776 assert!(msg.contains(HANDLER_VERSION_MISMATCH_CODE));
777 assert!(msg.contains("1.0.0"));
778 assert!(msg.contains("2.0.0"));
779 }
780
781 fn human_input_required() -> EngineError {
782 EngineError::HumanInputRequired {
783 run_id: Uuid::nil(),
784 step_id: Uuid::nil(),
785 message: "Answer the questions".to_string(),
786 }
787 }
788
789 #[test]
790 fn child_suspended_display_carries_child_and_cause() {
791 let child = Uuid::now_v7();
792 let err = EngineError::ChildSuspended {
793 run_id: child,
794 cause: Box::new(human_input_required()),
795 };
796
797 let msg = err.to_string();
798 assert!(msg.contains(&child.to_string()));
799 assert!(msg.contains("human input required"));
800 }
801
802 #[test]
803 fn leaf_suspensions_and_child_suspended_are_suspensions() {
804 let wake_at = Utc::now();
805 let suspensions = [
806 EngineError::ApprovalRequired {
807 run_id: Uuid::nil(),
808 step_id: Uuid::nil(),
809 message: "deploy?".to_string(),
810 },
811 human_input_required(),
812 EngineError::DelaySleeping {
813 run_id: Uuid::nil(),
814 step_id: Uuid::nil(),
815 wake_at,
816 },
817 EngineError::CapacitySleeping {
818 run_id: Uuid::nil(),
819 step_id: Uuid::nil(),
820 kind: "claude".to_string(),
821 wake_at,
822 },
823 EngineError::SignalWaiting {
824 run_id: Uuid::nil(),
825 step_id: Uuid::nil(),
826 step_name: "wait".to_string(),
827 name: "payment".to_string(),
828 key: "order-1".to_string(),
829 deadline_at: wake_at,
830 },
831 EngineError::ChildSuspended {
832 run_id: Uuid::nil(),
833 cause: Box::new(human_input_required()),
834 },
835 ];
836 for err in &suspensions {
837 assert!(err.is_suspension(), "{err} should be a suspension");
838 }
839 }
840
841 #[test]
842 fn failures_and_rejections_are_not_suspensions() {
843 let failures = [
844 EngineError::InvalidWorkflow("x".to_string()),
845 EngineError::HumanInputRejected {
846 run_id: Uuid::nil(),
847 step_id: Uuid::nil(),
848 reason: "no".to_string(),
849 },
850 EngineError::ApprovalRejected {
851 run_id: Uuid::nil(),
852 step_id: Uuid::nil(),
853 reason: "no".to_string(),
854 },
855 ];
856 for err in &failures {
857 assert!(!err.is_suspension(), "{err} should not be a suspension");
858 }
859 }
860
861 #[test]
862 fn suspension_leaf_unwraps_nested_child_suspensions() {
863 let err = EngineError::ChildSuspended {
864 run_id: Uuid::now_v7(),
865 cause: Box::new(EngineError::ChildSuspended {
866 run_id: Uuid::now_v7(),
867 cause: Box::new(human_input_required()),
868 }),
869 };
870
871 assert!(matches!(
872 err.suspension_leaf(),
873 EngineError::HumanInputRequired { .. }
874 ));
875 }
876
877 #[test]
878 fn suspension_status_follows_the_leaf() {
879 let human = EngineError::ChildSuspended {
880 run_id: Uuid::nil(),
881 cause: Box::new(human_input_required()),
882 };
883 assert_eq!(human.suspension_status(), RunStatus::AwaitingApproval);
884
885 let delay = EngineError::ChildSuspended {
886 run_id: Uuid::nil(),
887 cause: Box::new(EngineError::DelaySleeping {
888 run_id: Uuid::nil(),
889 step_id: Uuid::nil(),
890 wake_at: Utc::now(),
891 }),
892 };
893 assert_eq!(delay.suspension_status(), RunStatus::Sleeping);
894
895 let capacity = EngineError::ChildSuspended {
896 run_id: Uuid::nil(),
897 cause: Box::new(EngineError::CapacitySleeping {
898 run_id: Uuid::nil(),
899 step_id: Uuid::nil(),
900 kind: "claude".to_string(),
901 wake_at: Utc::now(),
902 }),
903 };
904 assert_eq!(capacity.suspension_status(), RunStatus::Sleeping);
905 }
906
907 #[test]
908 fn suspension_leaf_of_a_plain_error_is_itself() {
909 let err = EngineError::StepConfig("bad".to_string());
910 assert!(matches!(err.suspension_leaf(), EngineError::StepConfig(_)));
911 }
912}