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("invalid signal: {0}")]
304 InvalidSignal(String),
305
306 #[error("{WORKFLOW_GUARD_REJECTED_CODE}: {0}")]
312 WorkflowGuardRejected(#[from] WorkflowRejection),
313
314 #[error(
325 "replay divergence at position {position}: handler called '{expected}' but the run \
326 recorded '{recorded}' (handler changed since the run was suspended?)"
327 )]
328 ReplayDivergence {
329 position: u32,
332 expected: String,
334 recorded: String,
336 },
337
338 #[error(
348 "{HANDLER_VERSION_MISMATCH_CODE}: run {run_id} for handler '{workflow_name}' was created \
349 with version {run_version}, but the handler is now at version {current_version}; \
350 resume refused (no force override for resume)"
351 )]
352 HandlerVersionMismatch {
353 run_id: uuid::Uuid,
355 workflow_name: String,
357 run_version: String,
359 current_version: String,
361 },
362}
363
364impl From<StoreError> for EngineError {
365 fn from(err: StoreError) -> Self {
366 match err {
367 StoreError::ConcurrencyConflict { key, run_id } => {
368 EngineError::ConcurrencyConflict { key, run_id }
369 }
370 StoreError::InvalidConcurrencyLimit(e) => EngineError::InvalidConcurrencyLimit(e),
371 other => EngineError::Store(other),
372 }
373 }
374}
375
376impl EngineError {
377 pub fn is_suspension(&self) -> bool {
393 matches!(
394 self,
395 EngineError::ApprovalRequired { .. }
396 | EngineError::HumanInputRequired { .. }
397 | EngineError::DelaySleeping { .. }
398 | EngineError::SignalWaiting { .. }
399 | EngineError::ChildSuspended { .. }
400 )
401 }
402
403 pub fn suspension_leaf(&self) -> &EngineError {
425 let mut current = self;
426 while let EngineError::ChildSuspended { cause, .. } = current {
427 current = cause;
428 }
429 current
430 }
431
432 pub(crate) fn suspension_status(&self) -> RunStatus {
436 match self.suspension_leaf() {
437 EngineError::DelaySleeping { .. } | EngineError::SignalWaiting { .. } => {
438 RunStatus::Sleeping
439 }
440 _ => RunStatus::AwaitingApproval,
441 }
442 }
443}
444
445#[cfg(test)]
446mod tests {
447 use super::*;
448 use crate::retry_policy::is_run_retryable;
449
450 #[test]
451 fn invalid_workflow_display() {
452 let err = EngineError::InvalidWorkflow("unknown-handler".to_string());
453 assert!(err.to_string().contains("invalid workflow"));
454 assert!(err.to_string().contains("unknown-handler"));
455 }
456
457 #[test]
458 fn step_config_display() {
459 let err = EngineError::StepConfig("bad shell config".to_string());
460 assert!(err.to_string().contains("step config error"));
461 assert!(err.to_string().contains("bad shell config"));
462 }
463
464 #[test]
465 fn human_input_required_display() {
466 let err = EngineError::HumanInputRequired {
467 run_id: uuid::Uuid::nil(),
468 step_id: uuid::Uuid::nil(),
469 message: "Answer the questions".to_string(),
470 };
471 let text = err.to_string();
472 assert!(text.contains("human input required"));
473 assert!(text.contains("Answer the questions"));
474 }
475
476 #[test]
477 fn human_input_rejected_display() {
478 let err = EngineError::HumanInputRejected {
479 run_id: uuid::Uuid::nil(),
480 step_id: uuid::Uuid::nil(),
481 reason: "not relevant".to_string(),
482 };
483 let text = err.to_string();
484 assert!(text.contains("human input rejected"));
485 assert!(text.contains("not relevant"));
486 }
487
488 #[test]
489 fn store_error_from_conversion() {
490 let store_err = StoreError::RunNotFound(uuid::Uuid::nil());
491 let engine_err = EngineError::from(store_err);
492 assert!(engine_err.to_string().contains("store error"));
493 }
494
495 #[test]
496 fn store_concurrency_conflict_converts_to_engine_variant() {
497 let run_id = Uuid::now_v7();
498 let engine_err = EngineError::from(StoreError::ConcurrencyConflict {
499 key: "issue:12".to_string(),
500 run_id,
501 });
502 match engine_err {
503 EngineError::ConcurrencyConflict {
504 ref key,
505 run_id: holder,
506 } => {
507 assert_eq!(key, "issue:12");
508 assert_eq!(holder, run_id);
509 }
510 ref other => panic!("expected ConcurrencyConflict, got {other:?}"),
511 }
512 assert!(engine_err.to_string().contains("issue:12"));
513 assert!(!is_run_retryable(&engine_err));
514 }
515
516 #[test]
517 fn store_invalid_concurrency_limit_converts_to_engine_variant() {
518 let engine_err = EngineError::from(StoreError::InvalidConcurrencyLimit(
519 ConcurrencyLimitError::ZeroLimit {
520 group: "repo:acme".to_string(),
521 },
522 ));
523 assert!(
524 matches!(
525 engine_err,
526 EngineError::InvalidConcurrencyLimit(ConcurrencyLimitError::ZeroLimit { .. })
527 ),
528 "{engine_err:?}"
529 );
530 assert!(engine_err.to_string().contains("repo:acme"));
531 assert!(!is_run_retryable(&engine_err));
532 }
533
534 #[test]
535 fn run_budget_exceeded_display_carries_code_and_amounts() {
536 let err = EngineError::RunBudgetExceeded {
537 run_id: uuid::Uuid::nil(),
538 limit_usd: Decimal::new(200, 2),
539 spent_usd: Decimal::new(180, 2),
540 step_budget_usd: Decimal::new(50, 2),
541 };
542
543 let msg = err.to_string();
544 assert!(msg.contains(RUN_BUDGET_EXCEEDED_CODE));
545 assert!(msg.contains("2.00"));
546 assert!(msg.contains("1.80"));
547 assert!(msg.contains("0.50"));
548 }
549
550 #[test]
551 fn monthly_budget_exceeded_display_carries_code_and_amounts() {
552 let err = EngineError::MonthlyBudgetExceeded {
553 limit_usd: Decimal::new(10000, 2),
554 spent_usd: Decimal::new(10500, 2),
555 };
556
557 let msg = err.to_string();
558 assert!(msg.contains(MONTHLY_BUDGET_EXCEEDED_CODE));
559 assert!(msg.contains("100.00"));
560 assert!(msg.contains("105.00"));
561 }
562
563 #[test]
564 fn missing_artifact_display_names_the_step_and_pattern() {
565 let err = EngineError::MissingArtifact {
566 step: "build".to_string(),
567 pattern: "target/report.html".to_string(),
568 };
569
570 let msg = err.to_string();
571 assert!(msg.contains("\"build\""));
572 assert!(msg.contains("target/report.html"));
573 }
574
575 #[test]
576 fn artifact_not_declared_display_names_the_step_and_artifact() {
577 let err = EngineError::ArtifactNotDeclared {
578 step: "build".to_string(),
579 name: "report.htm".to_string(),
580 };
581
582 let msg = err.to_string();
583 assert!(msg.contains("\"build\""));
584 assert!(msg.contains("\"report.htm\""));
585 }
586
587 #[test]
588 fn artifact_not_found_display_names_the_producer() {
589 let err = EngineError::ArtifactNotFound {
590 step: "build".to_string(),
591 name: "report.html".to_string(),
592 };
593
594 let msg = err.to_string();
595 assert!(msg.contains("\"build\""));
596 assert!(msg.contains("report.html"));
597 }
598
599 #[test]
600 fn artifacts_unavailable_display() {
601 let err = EngineError::ArtifactsUnavailable("no blob store".to_string());
602 assert!(err.to_string().contains("not configured"));
603 }
604
605 #[test]
606 fn artifact_error_from_conversion() {
607 let engine_err = EngineError::from(ArtifactError::NotFound("a/b".to_string()));
608 assert!(engine_err.to_string().contains("artifact storage error"));
609 }
610
611 #[test]
612 fn serialization_error_from_conversion() {
613 let serde_err = serde_json::from_str::<String>("not json").unwrap_err();
614 let engine_err = EngineError::from(serde_err);
615 assert!(engine_err.to_string().contains("serialization error"));
616 }
617
618 #[test]
619 fn workflow_guard_rejected_display_carries_code_and_detail() {
620 use crate::guard::WorkflowRejection;
621
622 let rejection = WorkflowRejection::MaxDepthExceeded { depth: 6, max: 5 };
623 let err = EngineError::from(rejection);
624
625 let msg = err.to_string();
626 assert!(msg.contains(WORKFLOW_GUARD_REJECTED_CODE));
627 assert!(msg.contains("max call depth exceeded"));
628 assert!(msg.contains("6/5"));
629 }
630
631 #[test]
632 fn workflow_guard_rejected_from_conversion() {
633 use crate::guard::WorkflowRejection;
634
635 let rejection = WorkflowRejection::CycleDetected {
636 target: "wf-b".to_string(),
637 chain: vec!["wf-a".to_string(), "wf-b".to_string()],
638 };
639 let engine_err = EngineError::from(rejection);
640 assert!(engine_err.to_string().contains("cycle detected"));
641 }
642
643 #[test]
644 fn replay_divergence_display_carries_position_and_identities() {
645 let err = EngineError::ReplayDivergence {
646 position: 5,
647 expected: "resolve-base-branch (Shell)".to_string(),
648 recorded: "create-worktree (Shell)".to_string(),
649 };
650
651 let msg = err.to_string();
652 assert!(msg.contains("divergence"));
653 assert!(msg.contains("position 5"));
654 assert!(msg.contains("resolve-base-branch"));
655 assert!(msg.contains("create-worktree"));
656 }
657
658 #[test]
659 fn handler_version_mismatch_display_carries_code_and_versions() {
660 let err = EngineError::HandlerVersionMismatch {
661 run_id: uuid::Uuid::nil(),
662 workflow_name: "deploy".to_string(),
663 run_version: "1.0.0".to_string(),
664 current_version: "2.0.0".to_string(),
665 };
666
667 let msg = err.to_string();
668 assert!(msg.contains(HANDLER_VERSION_MISMATCH_CODE));
669 assert!(msg.contains("1.0.0"));
670 assert!(msg.contains("2.0.0"));
671 }
672
673 fn human_input_required() -> EngineError {
674 EngineError::HumanInputRequired {
675 run_id: Uuid::nil(),
676 step_id: Uuid::nil(),
677 message: "Answer the questions".to_string(),
678 }
679 }
680
681 #[test]
682 fn child_suspended_display_carries_child_and_cause() {
683 let child = Uuid::now_v7();
684 let err = EngineError::ChildSuspended {
685 run_id: child,
686 cause: Box::new(human_input_required()),
687 };
688
689 let msg = err.to_string();
690 assert!(msg.contains(&child.to_string()));
691 assert!(msg.contains("human input required"));
692 }
693
694 #[test]
695 fn leaf_suspensions_and_child_suspended_are_suspensions() {
696 let wake_at = Utc::now();
697 let suspensions = [
698 EngineError::ApprovalRequired {
699 run_id: Uuid::nil(),
700 step_id: Uuid::nil(),
701 message: "deploy?".to_string(),
702 },
703 human_input_required(),
704 EngineError::DelaySleeping {
705 run_id: Uuid::nil(),
706 step_id: Uuid::nil(),
707 wake_at,
708 },
709 EngineError::SignalWaiting {
710 run_id: Uuid::nil(),
711 step_id: Uuid::nil(),
712 step_name: "wait".to_string(),
713 name: "payment".to_string(),
714 key: "order-1".to_string(),
715 deadline_at: wake_at,
716 },
717 EngineError::ChildSuspended {
718 run_id: Uuid::nil(),
719 cause: Box::new(human_input_required()),
720 },
721 ];
722 for err in &suspensions {
723 assert!(err.is_suspension(), "{err} should be a suspension");
724 }
725 }
726
727 #[test]
728 fn failures_and_rejections_are_not_suspensions() {
729 let failures = [
730 EngineError::InvalidWorkflow("x".to_string()),
731 EngineError::HumanInputRejected {
732 run_id: Uuid::nil(),
733 step_id: Uuid::nil(),
734 reason: "no".to_string(),
735 },
736 EngineError::ApprovalRejected {
737 run_id: Uuid::nil(),
738 step_id: Uuid::nil(),
739 reason: "no".to_string(),
740 },
741 ];
742 for err in &failures {
743 assert!(!err.is_suspension(), "{err} should not be a suspension");
744 }
745 }
746
747 #[test]
748 fn suspension_leaf_unwraps_nested_child_suspensions() {
749 let err = EngineError::ChildSuspended {
750 run_id: Uuid::now_v7(),
751 cause: Box::new(EngineError::ChildSuspended {
752 run_id: Uuid::now_v7(),
753 cause: Box::new(human_input_required()),
754 }),
755 };
756
757 assert!(matches!(
758 err.suspension_leaf(),
759 EngineError::HumanInputRequired { .. }
760 ));
761 }
762
763 #[test]
764 fn suspension_status_follows_the_leaf() {
765 let human = EngineError::ChildSuspended {
766 run_id: Uuid::nil(),
767 cause: Box::new(human_input_required()),
768 };
769 assert_eq!(human.suspension_status(), RunStatus::AwaitingApproval);
770
771 let delay = EngineError::ChildSuspended {
772 run_id: Uuid::nil(),
773 cause: Box::new(EngineError::DelaySleeping {
774 run_id: Uuid::nil(),
775 step_id: Uuid::nil(),
776 wake_at: Utc::now(),
777 }),
778 };
779 assert_eq!(delay.suspension_status(), RunStatus::Sleeping);
780 }
781
782 #[test]
783 fn suspension_leaf_of_a_plain_error_is_itself() {
784 let err = EngineError::StepConfig("bad".to_string());
785 assert!(matches!(err.suspension_leaf(), EngineError::StepConfig(_)));
786 }
787}