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::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 workflow: {0}")]
60 InvalidWorkflow(String),
61
62 #[error("step config error: {0}")]
64 StepConfig(String),
65
66 #[error("decision error: {0}")]
68 Decision(#[from] ironflow_core::error::DecisionError),
69
70 #[error(
73 "decision step '{step}' requires a decision provider; \
74 wire one with Engine::with_decision_provider(...)"
75 )]
76 NoDecisionProvider {
77 step: String,
79 },
80
81 #[error("serialization error: {0}")]
83 Serialization(#[from] serde_json::Error),
84
85 #[error(
91 "{RUN_BUDGET_EXCEEDED_CODE}: run {run_id} would exceed its cost cap \
92 (spent {spent_usd} USD + next step {step_budget_usd} USD > cap {limit_usd} USD)"
93 )]
94 RunBudgetExceeded {
95 run_id: uuid::Uuid,
97 limit_usd: Decimal,
99 spent_usd: Decimal,
101 step_budget_usd: Decimal,
103 },
104
105 #[error(
109 "{MONTHLY_BUDGET_EXCEEDED_CODE}: monthly cost quota exhausted \
110 ({spent_usd} USD spent of {limit_usd} USD)"
111 )]
112 MonthlyBudgetExceeded {
113 limit_usd: Decimal,
115 spent_usd: Decimal,
117 },
118
119 #[error("step {step:?} declared output {pattern:?} but no file matched")]
125 MissingArtifact {
126 step: String,
128 pattern: String,
130 },
131
132 #[error("step {step:?} declares no artifact output named {name:?}")]
134 ArtifactNotDeclared {
135 step: String,
137 name: String,
139 },
140
141 #[error("no artifact {name:?} produced by step {step:?} before this point")]
143 ArtifactNotFound {
144 step: String,
146 name: String,
148 },
149
150 #[error("artifact storage is not configured: {0}")]
152 ArtifactsUnavailable(String),
153
154 #[error("artifact storage error: {0}")]
156 Artifact(#[from] ArtifactError),
157
158 #[error("approval required for run {run_id}, step {step_id}: {message}")]
160 ApprovalRequired {
161 run_id: uuid::Uuid,
163 step_id: uuid::Uuid,
165 message: String,
167 },
168
169 #[error("approval rejected for run {run_id}, step {step_id}: {reason}")]
176 ApprovalRejected {
177 run_id: uuid::Uuid,
179 step_id: uuid::Uuid,
181 reason: String,
183 },
184
185 #[error("human input required for run {run_id}, step {step_id}: {message}")]
191 HumanInputRequired {
192 run_id: uuid::Uuid,
194 step_id: uuid::Uuid,
196 message: String,
198 },
199
200 #[error("human input rejected for run {run_id}, step {step_id}: {reason}")]
207 HumanInputRejected {
208 run_id: uuid::Uuid,
210 step_id: uuid::Uuid,
212 reason: String,
214 },
215
216 #[error("delay sleeping for run {run_id}, step {step_id}: wake at {wake_at}")]
222 DelaySleeping {
223 run_id: uuid::Uuid,
225 step_id: uuid::Uuid,
227 wake_at: chrono::DateTime<chrono::Utc>,
229 },
230
231 #[error("run {run_id} waiting for signal {name:?} with key {key:?} until {deadline_at}")]
239 SignalWaiting {
240 run_id: Uuid,
242 step_id: Uuid,
244 step_name: String,
246 name: String,
248 key: String,
250 deadline_at: DateTime<Utc>,
252 },
253
254 #[error("child run {run_id} suspended: {cause}")]
280 ChildSuspended {
281 run_id: Uuid,
283 cause: Box<EngineError>,
290 },
291
292 #[error("invalid signal: {0}")]
295 InvalidSignal(String),
296
297 #[error("{WORKFLOW_GUARD_REJECTED_CODE}: {0}")]
303 WorkflowGuardRejected(#[from] WorkflowRejection),
304
305 #[error(
316 "replay divergence at position {position}: handler called '{expected}' but the run \
317 recorded '{recorded}' (handler changed since the run was suspended?)"
318 )]
319 ReplayDivergence {
320 position: u32,
323 expected: String,
325 recorded: String,
327 },
328
329 #[error(
339 "{HANDLER_VERSION_MISMATCH_CODE}: run {run_id} for handler '{workflow_name}' was created \
340 with version {run_version}, but the handler is now at version {current_version}; \
341 resume refused (no force override for resume)"
342 )]
343 HandlerVersionMismatch {
344 run_id: uuid::Uuid,
346 workflow_name: String,
348 run_version: String,
350 current_version: String,
352 },
353}
354
355impl From<StoreError> for EngineError {
356 fn from(err: StoreError) -> Self {
357 match err {
358 StoreError::ConcurrencyConflict { key, run_id } => {
359 EngineError::ConcurrencyConflict { key, run_id }
360 }
361 other => EngineError::Store(other),
362 }
363 }
364}
365
366impl EngineError {
367 pub fn is_suspension(&self) -> bool {
383 matches!(
384 self,
385 EngineError::ApprovalRequired { .. }
386 | EngineError::HumanInputRequired { .. }
387 | EngineError::DelaySleeping { .. }
388 | EngineError::SignalWaiting { .. }
389 | EngineError::ChildSuspended { .. }
390 )
391 }
392
393 pub fn suspension_leaf(&self) -> &EngineError {
415 let mut current = self;
416 while let EngineError::ChildSuspended { cause, .. } = current {
417 current = cause;
418 }
419 current
420 }
421
422 pub(crate) fn suspension_status(&self) -> RunStatus {
426 match self.suspension_leaf() {
427 EngineError::DelaySleeping { .. } | EngineError::SignalWaiting { .. } => {
428 RunStatus::Sleeping
429 }
430 _ => RunStatus::AwaitingApproval,
431 }
432 }
433}
434
435#[cfg(test)]
436mod tests {
437 use super::*;
438 use crate::retry_policy::is_run_retryable;
439
440 #[test]
441 fn invalid_workflow_display() {
442 let err = EngineError::InvalidWorkflow("unknown-handler".to_string());
443 assert!(err.to_string().contains("invalid workflow"));
444 assert!(err.to_string().contains("unknown-handler"));
445 }
446
447 #[test]
448 fn step_config_display() {
449 let err = EngineError::StepConfig("bad shell config".to_string());
450 assert!(err.to_string().contains("step config error"));
451 assert!(err.to_string().contains("bad shell config"));
452 }
453
454 #[test]
455 fn human_input_required_display() {
456 let err = EngineError::HumanInputRequired {
457 run_id: uuid::Uuid::nil(),
458 step_id: uuid::Uuid::nil(),
459 message: "Answer the questions".to_string(),
460 };
461 let text = err.to_string();
462 assert!(text.contains("human input required"));
463 assert!(text.contains("Answer the questions"));
464 }
465
466 #[test]
467 fn human_input_rejected_display() {
468 let err = EngineError::HumanInputRejected {
469 run_id: uuid::Uuid::nil(),
470 step_id: uuid::Uuid::nil(),
471 reason: "not relevant".to_string(),
472 };
473 let text = err.to_string();
474 assert!(text.contains("human input rejected"));
475 assert!(text.contains("not relevant"));
476 }
477
478 #[test]
479 fn store_error_from_conversion() {
480 let store_err = StoreError::RunNotFound(uuid::Uuid::nil());
481 let engine_err = EngineError::from(store_err);
482 assert!(engine_err.to_string().contains("store error"));
483 }
484
485 #[test]
486 fn store_concurrency_conflict_converts_to_engine_variant() {
487 let run_id = Uuid::now_v7();
488 let engine_err = EngineError::from(StoreError::ConcurrencyConflict {
489 key: "issue:12".to_string(),
490 run_id,
491 });
492 match engine_err {
493 EngineError::ConcurrencyConflict {
494 ref key,
495 run_id: holder,
496 } => {
497 assert_eq!(key, "issue:12");
498 assert_eq!(holder, run_id);
499 }
500 ref other => panic!("expected ConcurrencyConflict, got {other:?}"),
501 }
502 assert!(engine_err.to_string().contains("issue:12"));
503 assert!(!is_run_retryable(&engine_err));
504 }
505
506 #[test]
507 fn run_budget_exceeded_display_carries_code_and_amounts() {
508 let err = EngineError::RunBudgetExceeded {
509 run_id: uuid::Uuid::nil(),
510 limit_usd: Decimal::new(200, 2),
511 spent_usd: Decimal::new(180, 2),
512 step_budget_usd: Decimal::new(50, 2),
513 };
514
515 let msg = err.to_string();
516 assert!(msg.contains(RUN_BUDGET_EXCEEDED_CODE));
517 assert!(msg.contains("2.00"));
518 assert!(msg.contains("1.80"));
519 assert!(msg.contains("0.50"));
520 }
521
522 #[test]
523 fn monthly_budget_exceeded_display_carries_code_and_amounts() {
524 let err = EngineError::MonthlyBudgetExceeded {
525 limit_usd: Decimal::new(10000, 2),
526 spent_usd: Decimal::new(10500, 2),
527 };
528
529 let msg = err.to_string();
530 assert!(msg.contains(MONTHLY_BUDGET_EXCEEDED_CODE));
531 assert!(msg.contains("100.00"));
532 assert!(msg.contains("105.00"));
533 }
534
535 #[test]
536 fn missing_artifact_display_names_the_step_and_pattern() {
537 let err = EngineError::MissingArtifact {
538 step: "build".to_string(),
539 pattern: "target/report.html".to_string(),
540 };
541
542 let msg = err.to_string();
543 assert!(msg.contains("\"build\""));
544 assert!(msg.contains("target/report.html"));
545 }
546
547 #[test]
548 fn artifact_not_declared_display_names_the_step_and_artifact() {
549 let err = EngineError::ArtifactNotDeclared {
550 step: "build".to_string(),
551 name: "report.htm".to_string(),
552 };
553
554 let msg = err.to_string();
555 assert!(msg.contains("\"build\""));
556 assert!(msg.contains("\"report.htm\""));
557 }
558
559 #[test]
560 fn artifact_not_found_display_names_the_producer() {
561 let err = EngineError::ArtifactNotFound {
562 step: "build".to_string(),
563 name: "report.html".to_string(),
564 };
565
566 let msg = err.to_string();
567 assert!(msg.contains("\"build\""));
568 assert!(msg.contains("report.html"));
569 }
570
571 #[test]
572 fn artifacts_unavailable_display() {
573 let err = EngineError::ArtifactsUnavailable("no blob store".to_string());
574 assert!(err.to_string().contains("not configured"));
575 }
576
577 #[test]
578 fn artifact_error_from_conversion() {
579 let engine_err = EngineError::from(ArtifactError::NotFound("a/b".to_string()));
580 assert!(engine_err.to_string().contains("artifact storage error"));
581 }
582
583 #[test]
584 fn serialization_error_from_conversion() {
585 let serde_err = serde_json::from_str::<String>("not json").unwrap_err();
586 let engine_err = EngineError::from(serde_err);
587 assert!(engine_err.to_string().contains("serialization error"));
588 }
589
590 #[test]
591 fn workflow_guard_rejected_display_carries_code_and_detail() {
592 use crate::guard::WorkflowRejection;
593
594 let rejection = WorkflowRejection::MaxDepthExceeded { depth: 6, max: 5 };
595 let err = EngineError::from(rejection);
596
597 let msg = err.to_string();
598 assert!(msg.contains(WORKFLOW_GUARD_REJECTED_CODE));
599 assert!(msg.contains("max call depth exceeded"));
600 assert!(msg.contains("6/5"));
601 }
602
603 #[test]
604 fn workflow_guard_rejected_from_conversion() {
605 use crate::guard::WorkflowRejection;
606
607 let rejection = WorkflowRejection::CycleDetected {
608 target: "wf-b".to_string(),
609 chain: vec!["wf-a".to_string(), "wf-b".to_string()],
610 };
611 let engine_err = EngineError::from(rejection);
612 assert!(engine_err.to_string().contains("cycle detected"));
613 }
614
615 #[test]
616 fn replay_divergence_display_carries_position_and_identities() {
617 let err = EngineError::ReplayDivergence {
618 position: 5,
619 expected: "resolve-base-branch (Shell)".to_string(),
620 recorded: "create-worktree (Shell)".to_string(),
621 };
622
623 let msg = err.to_string();
624 assert!(msg.contains("divergence"));
625 assert!(msg.contains("position 5"));
626 assert!(msg.contains("resolve-base-branch"));
627 assert!(msg.contains("create-worktree"));
628 }
629
630 #[test]
631 fn handler_version_mismatch_display_carries_code_and_versions() {
632 let err = EngineError::HandlerVersionMismatch {
633 run_id: uuid::Uuid::nil(),
634 workflow_name: "deploy".to_string(),
635 run_version: "1.0.0".to_string(),
636 current_version: "2.0.0".to_string(),
637 };
638
639 let msg = err.to_string();
640 assert!(msg.contains(HANDLER_VERSION_MISMATCH_CODE));
641 assert!(msg.contains("1.0.0"));
642 assert!(msg.contains("2.0.0"));
643 }
644
645 fn human_input_required() -> EngineError {
646 EngineError::HumanInputRequired {
647 run_id: Uuid::nil(),
648 step_id: Uuid::nil(),
649 message: "Answer the questions".to_string(),
650 }
651 }
652
653 #[test]
654 fn child_suspended_display_carries_child_and_cause() {
655 let child = Uuid::now_v7();
656 let err = EngineError::ChildSuspended {
657 run_id: child,
658 cause: Box::new(human_input_required()),
659 };
660
661 let msg = err.to_string();
662 assert!(msg.contains(&child.to_string()));
663 assert!(msg.contains("human input required"));
664 }
665
666 #[test]
667 fn leaf_suspensions_and_child_suspended_are_suspensions() {
668 let wake_at = Utc::now();
669 let suspensions = [
670 EngineError::ApprovalRequired {
671 run_id: Uuid::nil(),
672 step_id: Uuid::nil(),
673 message: "deploy?".to_string(),
674 },
675 human_input_required(),
676 EngineError::DelaySleeping {
677 run_id: Uuid::nil(),
678 step_id: Uuid::nil(),
679 wake_at,
680 },
681 EngineError::SignalWaiting {
682 run_id: Uuid::nil(),
683 step_id: Uuid::nil(),
684 step_name: "wait".to_string(),
685 name: "payment".to_string(),
686 key: "order-1".to_string(),
687 deadline_at: wake_at,
688 },
689 EngineError::ChildSuspended {
690 run_id: Uuid::nil(),
691 cause: Box::new(human_input_required()),
692 },
693 ];
694 for err in &suspensions {
695 assert!(err.is_suspension(), "{err} should be a suspension");
696 }
697 }
698
699 #[test]
700 fn failures_and_rejections_are_not_suspensions() {
701 let failures = [
702 EngineError::InvalidWorkflow("x".to_string()),
703 EngineError::HumanInputRejected {
704 run_id: Uuid::nil(),
705 step_id: Uuid::nil(),
706 reason: "no".to_string(),
707 },
708 EngineError::ApprovalRejected {
709 run_id: Uuid::nil(),
710 step_id: Uuid::nil(),
711 reason: "no".to_string(),
712 },
713 ];
714 for err in &failures {
715 assert!(!err.is_suspension(), "{err} should not be a suspension");
716 }
717 }
718
719 #[test]
720 fn suspension_leaf_unwraps_nested_child_suspensions() {
721 let err = EngineError::ChildSuspended {
722 run_id: Uuid::now_v7(),
723 cause: Box::new(EngineError::ChildSuspended {
724 run_id: Uuid::now_v7(),
725 cause: Box::new(human_input_required()),
726 }),
727 };
728
729 assert!(matches!(
730 err.suspension_leaf(),
731 EngineError::HumanInputRequired { .. }
732 ));
733 }
734
735 #[test]
736 fn suspension_status_follows_the_leaf() {
737 let human = EngineError::ChildSuspended {
738 run_id: Uuid::nil(),
739 cause: Box::new(human_input_required()),
740 };
741 assert_eq!(human.suspension_status(), RunStatus::AwaitingApproval);
742
743 let delay = EngineError::ChildSuspended {
744 run_id: Uuid::nil(),
745 cause: Box::new(EngineError::DelaySleeping {
746 run_id: Uuid::nil(),
747 step_id: Uuid::nil(),
748 wake_at: Utc::now(),
749 }),
750 };
751 assert_eq!(delay.suspension_status(), RunStatus::Sleeping);
752 }
753
754 #[test]
755 fn suspension_leaf_of_a_plain_error_is_itself() {
756 let err = EngineError::StepConfig("bad".to_string());
757 assert!(matches!(err.suspension_leaf(), EngineError::StepConfig(_)));
758 }
759}