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 worker tag: {0}")]
74 InvalidWorkerTag(WorkerTagError),
75
76 #[error("invalid workflow: {0}")]
78 InvalidWorkflow(String),
79
80 #[error("step config error: {0}")]
82 StepConfig(String),
83
84 #[error("decision error: {0}")]
86 Decision(#[from] ironflow_core::error::DecisionError),
87
88 #[error(
91 "decision step '{step}' requires a decision provider; \
92 wire one with Engine::with_decision_provider(...)"
93 )]
94 NoDecisionProvider {
95 step: String,
97 },
98
99 #[error("serialization error: {0}")]
101 Serialization(#[from] serde_json::Error),
102
103 #[error(
109 "{RUN_BUDGET_EXCEEDED_CODE}: run {run_id} would exceed its cost cap \
110 (spent {spent_usd} USD + next step {step_budget_usd} USD > cap {limit_usd} USD)"
111 )]
112 RunBudgetExceeded {
113 run_id: uuid::Uuid,
115 limit_usd: Decimal,
117 spent_usd: Decimal,
119 step_budget_usd: Decimal,
121 },
122
123 #[error(
127 "{MONTHLY_BUDGET_EXCEEDED_CODE}: monthly cost quota exhausted \
128 ({spent_usd} USD spent of {limit_usd} USD)"
129 )]
130 MonthlyBudgetExceeded {
131 limit_usd: Decimal,
133 spent_usd: Decimal,
135 },
136
137 #[error("step {step:?} declared output {pattern:?} but no file matched")]
143 MissingArtifact {
144 step: String,
146 pattern: String,
148 },
149
150 #[error("step {step:?} declares no artifact output named {name:?}")]
152 ArtifactNotDeclared {
153 step: String,
155 name: String,
157 },
158
159 #[error("no artifact {name:?} produced by step {step:?} before this point")]
161 ArtifactNotFound {
162 step: String,
164 name: String,
166 },
167
168 #[error("artifact storage is not configured: {0}")]
170 ArtifactsUnavailable(String),
171
172 #[error("artifact storage error: {0}")]
174 Artifact(#[from] ArtifactError),
175
176 #[error("approval required for run {run_id}, step {step_id}: {message}")]
178 ApprovalRequired {
179 run_id: uuid::Uuid,
181 step_id: uuid::Uuid,
183 message: String,
185 },
186
187 #[error("approval rejected for run {run_id}, step {step_id}: {reason}")]
194 ApprovalRejected {
195 run_id: uuid::Uuid,
197 step_id: uuid::Uuid,
199 reason: String,
201 },
202
203 #[error("human input required for run {run_id}, step {step_id}: {message}")]
209 HumanInputRequired {
210 run_id: uuid::Uuid,
212 step_id: uuid::Uuid,
214 message: String,
216 },
217
218 #[error("human input rejected for run {run_id}, step {step_id}: {reason}")]
225 HumanInputRejected {
226 run_id: uuid::Uuid,
228 step_id: uuid::Uuid,
230 reason: String,
232 },
233
234 #[error("delay sleeping for run {run_id}, step {step_id}: wake at {wake_at}")]
240 DelaySleeping {
241 run_id: uuid::Uuid,
243 step_id: uuid::Uuid,
245 wake_at: chrono::DateTime<chrono::Utc>,
247 },
248
249 #[error("run {run_id} waiting for {kind} capacity at step {step_id}: wake at {wake_at}")]
275 CapacitySleeping {
276 run_id: Uuid,
278 step_id: Uuid,
280 kind: String,
282 wake_at: DateTime<Utc>,
284 },
285
286 #[error("run {run_id} waiting for signal {name:?} with key {key:?} until {deadline_at}")]
294 SignalWaiting {
295 run_id: Uuid,
297 step_id: Uuid,
299 step_name: String,
301 name: String,
303 key: String,
305 deadline_at: DateTime<Utc>,
307 },
308
309 #[error("child run {run_id} suspended: {cause}")]
335 ChildSuspended {
336 run_id: Uuid,
338 cause: Box<EngineError>,
346 },
347
348 #[error("child run {run_id} was cancelled")]
366 ChildRunCancelled {
367 run_id: Uuid,
369 },
370
371 #[error("invalid signal: {0}")]
374 InvalidSignal(String),
375
376 #[error("{WORKFLOW_GUARD_REJECTED_CODE}: {0}")]
382 WorkflowGuardRejected(#[from] WorkflowRejection),
383
384 #[error(
395 "replay divergence at position {position}: handler called '{expected}' but the run \
396 recorded '{recorded}' (handler changed since the run was suspended?)"
397 )]
398 ReplayDivergence {
399 position: u32,
402 expected: String,
404 recorded: String,
406 },
407
408 #[error(
418 "{HANDLER_VERSION_MISMATCH_CODE}: run {run_id} for handler '{workflow_name}' was created \
419 with version {run_version}, but the handler is now at version {current_version}; \
420 resume refused (no force override for resume)"
421 )]
422 HandlerVersionMismatch {
423 run_id: uuid::Uuid,
425 workflow_name: String,
427 run_version: String,
429 current_version: String,
431 },
432}
433
434impl From<StoreError> for EngineError {
435 fn from(err: StoreError) -> Self {
436 match err {
437 StoreError::ConcurrencyConflict { key, run_id } => {
438 EngineError::ConcurrencyConflict { key, run_id }
439 }
440 StoreError::InvalidConcurrencyLimit(e) => EngineError::InvalidConcurrencyLimit(e),
441 StoreError::InvalidWorkerTag(e) => EngineError::InvalidWorkerTag(e),
442 other => EngineError::Store(other),
443 }
444 }
445}
446
447impl EngineError {
448 pub fn is_suspension(&self) -> bool {
465 matches!(
466 self,
467 EngineError::ApprovalRequired { .. }
468 | EngineError::HumanInputRequired { .. }
469 | EngineError::DelaySleeping { .. }
470 | EngineError::CapacitySleeping { .. }
471 | EngineError::SignalWaiting { .. }
472 | EngineError::ChildSuspended { .. }
473 )
474 }
475
476 pub fn suspension_leaf(&self) -> &EngineError {
498 let mut current = self;
499 while let EngineError::ChildSuspended { cause, .. } = current {
500 current = cause;
501 }
502 current
503 }
504
505 pub(crate) fn suspension_status(&self) -> RunStatus {
509 match self.suspension_leaf() {
510 EngineError::DelaySleeping { .. }
511 | EngineError::CapacitySleeping { .. }
512 | EngineError::SignalWaiting { .. } => RunStatus::Sleeping,
513 _ => RunStatus::AwaitingApproval,
514 }
515 }
516}
517
518#[cfg(test)]
519mod tests {
520 use super::*;
521 use crate::retry_policy::is_run_retryable;
522
523 #[test]
524 fn invalid_workflow_display() {
525 let err = EngineError::InvalidWorkflow("unknown-handler".to_string());
526 assert!(err.to_string().contains("invalid workflow"));
527 assert!(err.to_string().contains("unknown-handler"));
528 }
529
530 #[test]
531 fn step_config_display() {
532 let err = EngineError::StepConfig("bad shell config".to_string());
533 assert!(err.to_string().contains("step config error"));
534 assert!(err.to_string().contains("bad shell config"));
535 }
536
537 #[test]
538 fn human_input_required_display() {
539 let err = EngineError::HumanInputRequired {
540 run_id: uuid::Uuid::nil(),
541 step_id: uuid::Uuid::nil(),
542 message: "Answer the questions".to_string(),
543 };
544 let text = err.to_string();
545 assert!(text.contains("human input required"));
546 assert!(text.contains("Answer the questions"));
547 }
548
549 #[test]
550 fn human_input_rejected_display() {
551 let err = EngineError::HumanInputRejected {
552 run_id: uuid::Uuid::nil(),
553 step_id: uuid::Uuid::nil(),
554 reason: "not relevant".to_string(),
555 };
556 let text = err.to_string();
557 assert!(text.contains("human input rejected"));
558 assert!(text.contains("not relevant"));
559 }
560
561 #[test]
562 fn store_error_from_conversion() {
563 let store_err = StoreError::RunNotFound(uuid::Uuid::nil());
564 let engine_err = EngineError::from(store_err);
565 assert!(engine_err.to_string().contains("store error"));
566 }
567
568 #[test]
569 fn store_concurrency_conflict_converts_to_engine_variant() {
570 let run_id = Uuid::now_v7();
571 let engine_err = EngineError::from(StoreError::ConcurrencyConflict {
572 key: "issue:12".to_string(),
573 run_id,
574 });
575 match engine_err {
576 EngineError::ConcurrencyConflict {
577 ref key,
578 run_id: holder,
579 } => {
580 assert_eq!(key, "issue:12");
581 assert_eq!(holder, run_id);
582 }
583 ref other => panic!("expected ConcurrencyConflict, got {other:?}"),
584 }
585 assert!(engine_err.to_string().contains("issue:12"));
586 assert!(!is_run_retryable(&engine_err));
587 }
588
589 #[test]
590 fn store_invalid_concurrency_limit_converts_to_engine_variant() {
591 let engine_err = EngineError::from(StoreError::InvalidConcurrencyLimit(
592 ConcurrencyLimitError::ZeroLimit {
593 group: "repo:acme".to_string(),
594 },
595 ));
596 assert!(
597 matches!(
598 engine_err,
599 EngineError::InvalidConcurrencyLimit(ConcurrencyLimitError::ZeroLimit { .. })
600 ),
601 "{engine_err:?}"
602 );
603 assert!(engine_err.to_string().contains("repo:acme"));
604 assert!(!is_run_retryable(&engine_err));
605 }
606
607 #[test]
608 fn store_invalid_worker_tag_converts_to_engine_variant() {
609 let engine_err =
610 EngineError::from(StoreError::InvalidWorkerTag(WorkerTagError::InvalidChar {
611 tag: "bad,tag".to_string(),
612 }));
613 assert!(
614 matches!(
615 engine_err,
616 EngineError::InvalidWorkerTag(WorkerTagError::InvalidChar { .. })
617 ),
618 "{engine_err:?}"
619 );
620 assert!(engine_err.to_string().contains("bad,tag"));
621 assert!(!is_run_retryable(&engine_err));
622 }
623
624 #[test]
625 fn run_budget_exceeded_display_carries_code_and_amounts() {
626 let err = EngineError::RunBudgetExceeded {
627 run_id: uuid::Uuid::nil(),
628 limit_usd: Decimal::new(200, 2),
629 spent_usd: Decimal::new(180, 2),
630 step_budget_usd: Decimal::new(50, 2),
631 };
632
633 let msg = err.to_string();
634 assert!(msg.contains(RUN_BUDGET_EXCEEDED_CODE));
635 assert!(msg.contains("2.00"));
636 assert!(msg.contains("1.80"));
637 assert!(msg.contains("0.50"));
638 }
639
640 #[test]
641 fn monthly_budget_exceeded_display_carries_code_and_amounts() {
642 let err = EngineError::MonthlyBudgetExceeded {
643 limit_usd: Decimal::new(10000, 2),
644 spent_usd: Decimal::new(10500, 2),
645 };
646
647 let msg = err.to_string();
648 assert!(msg.contains(MONTHLY_BUDGET_EXCEEDED_CODE));
649 assert!(msg.contains("100.00"));
650 assert!(msg.contains("105.00"));
651 }
652
653 #[test]
654 fn missing_artifact_display_names_the_step_and_pattern() {
655 let err = EngineError::MissingArtifact {
656 step: "build".to_string(),
657 pattern: "target/report.html".to_string(),
658 };
659
660 let msg = err.to_string();
661 assert!(msg.contains("\"build\""));
662 assert!(msg.contains("target/report.html"));
663 }
664
665 #[test]
666 fn artifact_not_declared_display_names_the_step_and_artifact() {
667 let err = EngineError::ArtifactNotDeclared {
668 step: "build".to_string(),
669 name: "report.htm".to_string(),
670 };
671
672 let msg = err.to_string();
673 assert!(msg.contains("\"build\""));
674 assert!(msg.contains("\"report.htm\""));
675 }
676
677 #[test]
678 fn artifact_not_found_display_names_the_producer() {
679 let err = EngineError::ArtifactNotFound {
680 step: "build".to_string(),
681 name: "report.html".to_string(),
682 };
683
684 let msg = err.to_string();
685 assert!(msg.contains("\"build\""));
686 assert!(msg.contains("report.html"));
687 }
688
689 #[test]
690 fn artifacts_unavailable_display() {
691 let err = EngineError::ArtifactsUnavailable("no blob store".to_string());
692 assert!(err.to_string().contains("not configured"));
693 }
694
695 #[test]
696 fn artifact_error_from_conversion() {
697 let engine_err = EngineError::from(ArtifactError::NotFound("a/b".to_string()));
698 assert!(engine_err.to_string().contains("artifact storage error"));
699 }
700
701 #[test]
702 fn serialization_error_from_conversion() {
703 let serde_err = serde_json::from_str::<String>("not json").unwrap_err();
704 let engine_err = EngineError::from(serde_err);
705 assert!(engine_err.to_string().contains("serialization error"));
706 }
707
708 #[test]
709 fn workflow_guard_rejected_display_carries_code_and_detail() {
710 use crate::guard::WorkflowRejection;
711
712 let rejection = WorkflowRejection::MaxDepthExceeded { depth: 6, max: 5 };
713 let err = EngineError::from(rejection);
714
715 let msg = err.to_string();
716 assert!(msg.contains(WORKFLOW_GUARD_REJECTED_CODE));
717 assert!(msg.contains("max call depth exceeded"));
718 assert!(msg.contains("6/5"));
719 }
720
721 #[test]
722 fn workflow_guard_rejected_from_conversion() {
723 use crate::guard::WorkflowRejection;
724
725 let rejection = WorkflowRejection::CycleDetected {
726 target: "wf-b".to_string(),
727 chain: vec!["wf-a".to_string(), "wf-b".to_string()],
728 };
729 let engine_err = EngineError::from(rejection);
730 assert!(engine_err.to_string().contains("cycle detected"));
731 }
732
733 #[test]
734 fn replay_divergence_display_carries_position_and_identities() {
735 let err = EngineError::ReplayDivergence {
736 position: 5,
737 expected: "resolve-base-branch (Shell)".to_string(),
738 recorded: "create-worktree (Shell)".to_string(),
739 };
740
741 let msg = err.to_string();
742 assert!(msg.contains("divergence"));
743 assert!(msg.contains("position 5"));
744 assert!(msg.contains("resolve-base-branch"));
745 assert!(msg.contains("create-worktree"));
746 }
747
748 #[test]
749 fn handler_version_mismatch_display_carries_code_and_versions() {
750 let err = EngineError::HandlerVersionMismatch {
751 run_id: uuid::Uuid::nil(),
752 workflow_name: "deploy".to_string(),
753 run_version: "1.0.0".to_string(),
754 current_version: "2.0.0".to_string(),
755 };
756
757 let msg = err.to_string();
758 assert!(msg.contains(HANDLER_VERSION_MISMATCH_CODE));
759 assert!(msg.contains("1.0.0"));
760 assert!(msg.contains("2.0.0"));
761 }
762
763 fn human_input_required() -> EngineError {
764 EngineError::HumanInputRequired {
765 run_id: Uuid::nil(),
766 step_id: Uuid::nil(),
767 message: "Answer the questions".to_string(),
768 }
769 }
770
771 #[test]
772 fn child_suspended_display_carries_child_and_cause() {
773 let child = Uuid::now_v7();
774 let err = EngineError::ChildSuspended {
775 run_id: child,
776 cause: Box::new(human_input_required()),
777 };
778
779 let msg = err.to_string();
780 assert!(msg.contains(&child.to_string()));
781 assert!(msg.contains("human input required"));
782 }
783
784 #[test]
785 fn leaf_suspensions_and_child_suspended_are_suspensions() {
786 let wake_at = Utc::now();
787 let suspensions = [
788 EngineError::ApprovalRequired {
789 run_id: Uuid::nil(),
790 step_id: Uuid::nil(),
791 message: "deploy?".to_string(),
792 },
793 human_input_required(),
794 EngineError::DelaySleeping {
795 run_id: Uuid::nil(),
796 step_id: Uuid::nil(),
797 wake_at,
798 },
799 EngineError::CapacitySleeping {
800 run_id: Uuid::nil(),
801 step_id: Uuid::nil(),
802 kind: "claude".to_string(),
803 wake_at,
804 },
805 EngineError::SignalWaiting {
806 run_id: Uuid::nil(),
807 step_id: Uuid::nil(),
808 step_name: "wait".to_string(),
809 name: "payment".to_string(),
810 key: "order-1".to_string(),
811 deadline_at: wake_at,
812 },
813 EngineError::ChildSuspended {
814 run_id: Uuid::nil(),
815 cause: Box::new(human_input_required()),
816 },
817 ];
818 for err in &suspensions {
819 assert!(err.is_suspension(), "{err} should be a suspension");
820 }
821 }
822
823 #[test]
824 fn failures_and_rejections_are_not_suspensions() {
825 let failures = [
826 EngineError::InvalidWorkflow("x".to_string()),
827 EngineError::HumanInputRejected {
828 run_id: Uuid::nil(),
829 step_id: Uuid::nil(),
830 reason: "no".to_string(),
831 },
832 EngineError::ApprovalRejected {
833 run_id: Uuid::nil(),
834 step_id: Uuid::nil(),
835 reason: "no".to_string(),
836 },
837 ];
838 for err in &failures {
839 assert!(!err.is_suspension(), "{err} should not be a suspension");
840 }
841 }
842
843 #[test]
844 fn suspension_leaf_unwraps_nested_child_suspensions() {
845 let err = EngineError::ChildSuspended {
846 run_id: Uuid::now_v7(),
847 cause: Box::new(EngineError::ChildSuspended {
848 run_id: Uuid::now_v7(),
849 cause: Box::new(human_input_required()),
850 }),
851 };
852
853 assert!(matches!(
854 err.suspension_leaf(),
855 EngineError::HumanInputRequired { .. }
856 ));
857 }
858
859 #[test]
860 fn suspension_status_follows_the_leaf() {
861 let human = EngineError::ChildSuspended {
862 run_id: Uuid::nil(),
863 cause: Box::new(human_input_required()),
864 };
865 assert_eq!(human.suspension_status(), RunStatus::AwaitingApproval);
866
867 let delay = EngineError::ChildSuspended {
868 run_id: Uuid::nil(),
869 cause: Box::new(EngineError::DelaySleeping {
870 run_id: Uuid::nil(),
871 step_id: Uuid::nil(),
872 wake_at: Utc::now(),
873 }),
874 };
875 assert_eq!(delay.suspension_status(), RunStatus::Sleeping);
876
877 let capacity = EngineError::ChildSuspended {
878 run_id: Uuid::nil(),
879 cause: Box::new(EngineError::CapacitySleeping {
880 run_id: Uuid::nil(),
881 step_id: Uuid::nil(),
882 kind: "claude".to_string(),
883 wake_at: Utc::now(),
884 }),
885 };
886 assert_eq!(capacity.suspension_status(), RunStatus::Sleeping);
887 }
888
889 #[test]
890 fn suspension_leaf_of_a_plain_error_is_itself() {
891 let err = EngineError::StepConfig("bad".to_string());
892 assert!(matches!(err.suspension_leaf(), EngineError::StepConfig(_)));
893 }
894}