Skip to main content

ironflow_engine/
error.rs

1//! Engine error types.
2
3use 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
15/// Business error code carried by [`EngineError::RunBudgetExceeded`].
16pub const RUN_BUDGET_EXCEEDED_CODE: &str = "RUN_BUDGET_EXCEEDED";
17
18/// Business error code carried by [`EngineError::MonthlyBudgetExceeded`].
19pub const MONTHLY_BUDGET_EXCEEDED_CODE: &str = "MONTHLY_BUDGET_EXCEEDED";
20
21/// Business error code carried by [`EngineError::ConcurrencyConflict`].
22///
23/// # Examples
24///
25/// ```
26/// use ironflow_engine::error::CONCURRENCY_CONFLICT_CODE;
27///
28/// assert_eq!(CONCURRENCY_CONFLICT_CODE, "CONCURRENCY_CONFLICT");
29/// ```
30pub const CONCURRENCY_CONFLICT_CODE: &str = "CONCURRENCY_CONFLICT";
31
32/// Business error code for handler-version mismatch on retry.
33pub const HANDLER_VERSION_MISMATCH_CODE: &str = "HANDLER_VERSION_MISMATCH";
34
35/// Errors produced by the workflow engine.
36#[derive(Debug, Error)]
37pub enum EngineError {
38    /// An operation (Shell, Http, Agent) failed during step execution.
39    #[error("operation failed: {0}")]
40    Operation(#[from] OperationError),
41
42    /// The backing store returned an error.
43    #[error("store error: {0}")]
44    Store(StoreError),
45
46    /// Another non-terminal run already holds the requested concurrency key.
47    ///
48    /// Converted from [`StoreError::ConcurrencyConflict`], so every `?` on a
49    /// store call yields this typed variant instead of [`EngineError::Store`].
50    #[error("concurrency key {key:?} is held by active run {run_id}")]
51    ConcurrencyConflict {
52        /// The contested key.
53        key: String,
54        /// The run holding it.
55        run_id: Uuid,
56    },
57
58    /// The requested concurrency limits are invalid: an empty or too long
59    /// group, a zero limit, or the same group listed twice.
60    ///
61    /// Converted from [`StoreError::InvalidConcurrencyLimit`], and returned by
62    /// [`Engine::enqueue_handler_with_options`](crate::engine::Engine::enqueue_handler_with_options)
63    /// before any other check.
64    #[error("invalid concurrency limit: {0}")]
65    InvalidConcurrencyLimit(ConcurrencyLimitError),
66
67    /// The requested run priority is outside
68    /// [`MIN_PRIORITY`](ironflow_store::entities::MIN_PRIORITY)`..=`[`MAX_PRIORITY`](ironflow_store::entities::MAX_PRIORITY).
69    ///
70    /// Returned by
71    /// [`Engine::enqueue_handler_with_options`](crate::engine::Engine::enqueue_handler_with_options)
72    /// before any other check.
73    #[error("invalid priority: {0}")]
74    InvalidPriority(String),
75    /// A requested worker tag is invalid: empty, too long, holding a character
76    /// outside ASCII alphanumerics and `- _ . : / =`, or one tag too many.
77    ///
78    /// Converted from [`StoreError::InvalidWorkerTag`], and returned by
79    /// [`Engine::enqueue_handler_with_options`](crate::engine::Engine::enqueue_handler_with_options)
80    /// before the run is created.
81    #[error("invalid worker tag: {0}")]
82    InvalidWorkerTag(WorkerTagError),
83
84    /// The workflow definition is invalid.
85    #[error("invalid workflow: {0}")]
86    InvalidWorkflow(String),
87
88    /// A step configuration could not be deserialized for execution.
89    #[error("step config error: {0}")]
90    StepConfig(String),
91
92    /// A decision answer was accessed by name with the wrong type or a missing key.
93    #[error("decision error: {0}")]
94    Decision(#[from] ironflow_core::error::DecisionError),
95
96    /// A decision step was reached but no [`DecisionProvider`](ironflow_core::decision::DecisionProvider)
97    /// is wired into the engine.
98    #[error(
99        "decision step '{step}' requires a decision provider; \
100         wire one with Engine::with_decision_provider(...)"
101    )]
102    NoDecisionProvider {
103        /// The decision step that could not run.
104        step: String,
105    },
106
107    /// JSON serialization error.
108    #[error("serialization error: {0}")]
109    Serialization(#[from] serde_json::Error),
110
111    /// The run reached its cumulative cost cap before launching an agent step.
112    ///
113    /// Raised *before* the step is created, so no work and no spend happen.
114    /// The engine transitions the run to
115    /// [`Cancelled`](ironflow_store::entities::RunStatus::Cancelled).
116    #[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        /// The run that hit its cap.
122        run_id: uuid::Uuid,
123        /// The configured cap, in USD.
124        limit_usd: Decimal,
125        /// Cost already accumulated by this run and its ancestors, in USD.
126        spent_usd: Decimal,
127        /// Declared budget of the step that was about to run, in USD.
128        step_budget_usd: Decimal,
129    },
130
131    /// The global monthly cost quota is exhausted; no new run may be created.
132    ///
133    /// Runs already in flight are never interrupted by this error.
134    #[error(
135        "{MONTHLY_BUDGET_EXCEEDED_CODE}: monthly cost quota exhausted \
136         ({spent_usd} USD spent of {limit_usd} USD)"
137    )]
138    MonthlyBudgetExceeded {
139        /// The configured monthly quota, in USD.
140        limit_usd: Decimal,
141        /// Cost already spent during the current calendar month, in USD.
142        spent_usd: Decimal,
143    },
144
145    /// A step declared an output that produced no file.
146    ///
147    /// Raised only when the step itself succeeded: a declared output that never
148    /// materialised is a broken contract, and failing here beats failing later
149    /// in whichever step tried to consume it.
150    #[error("step {step:?} declared output {pattern:?} but no file matched")]
151    MissingArtifact {
152        /// Name of the step that declared the output.
153        step: String,
154        /// The unmatched pattern.
155        pattern: String,
156    },
157
158    /// A handle was asked for an artifact the step never declared.
159    #[error("step {step:?} declares no artifact output named {name:?}")]
160    ArtifactNotDeclared {
161        /// Name of the step the handle was asked from.
162        step: String,
163        /// Name of the artifact that was asked for.
164        name: String,
165    },
166
167    /// A step asked for an artifact that no earlier step produced.
168    #[error("no artifact {name:?} produced by step {step:?} before this point")]
169    ArtifactNotFound {
170        /// Name of the producing step that was searched for.
171        step: String,
172        /// Name of the artifact that was searched for.
173        name: String,
174    },
175
176    /// Artifacts were used but no storage backend is configured.
177    #[error("artifact storage is not configured: {0}")]
178    ArtifactsUnavailable(String),
179
180    /// The artifact storage backend failed.
181    #[error("artifact storage error: {0}")]
182    Artifact(#[from] ArtifactError),
183
184    /// The run requires human approval before continuing.
185    #[error("approval required for run {run_id}, step {step_id}: {message}")]
186    ApprovalRequired {
187        /// The run that is awaiting approval.
188        run_id: uuid::Uuid,
189        /// The approval step that triggered the pause.
190        step_id: uuid::Uuid,
191        /// The approval message.
192        message: String,
193    },
194
195    /// An approval gate was rejected instead of granted.
196    ///
197    /// Raised when a [`StepInterceptor`](crate::executor::StepInterceptor) resolves
198    /// the gate with [`ApprovalOutcome::Rejected`](crate::executor::ApprovalOutcome::Rejected).
199    /// The engine fails the run; the rejection is deterministic, so the run is
200    /// never replayed.
201    #[error("approval rejected for run {run_id}, step {step_id}: {reason}")]
202    ApprovalRejected {
203        /// The run that was stopped.
204        run_id: uuid::Uuid,
205        /// The approval step that was rejected.
206        step_id: uuid::Uuid,
207        /// Why the gate was refused.
208        reason: String,
209    },
210
211    /// The run waits for a typed human input before continuing.
212    ///
213    /// Raised by [`WorkflowContext::human_input`](crate::context::WorkflowContext::human_input)
214    /// when no answer has been given yet. The engine transitions the run to
215    /// [`AwaitingApproval`](ironflow_store::entities::RunStatus::AwaitingApproval).
216    #[error("human input required for run {run_id}, step {step_id}: {message}")]
217    HumanInputRequired {
218        /// The run that is awaiting input.
219        run_id: uuid::Uuid,
220        /// The human input step that triggered the pause.
221        step_id: uuid::Uuid,
222        /// The message displayed to the person answering.
223        message: String,
224    },
225
226    /// A human input request was rejected instead of answered.
227    ///
228    /// Returned to the handler by
229    /// [`WorkflowContext::human_input`](crate::context::WorkflowContext::human_input),
230    /// which decides what happens next. A propagated rejection fails the run,
231    /// and the run is never retried.
232    #[error("human input rejected for run {run_id}, step {step_id}: {reason}")]
233    HumanInputRejected {
234        /// The run the input belongs to.
235        run_id: uuid::Uuid,
236        /// The human input step that was rejected.
237        step_id: uuid::Uuid,
238        /// Why the input was refused.
239        reason: String,
240    },
241
242    /// A delay step suspended the run until the given time.
243    ///
244    /// The engine transitions the run to
245    /// [`Sleeping`](ironflow_store::entities::RunStatus::Sleeping) and sets
246    /// `scheduled_at` so the worker re-queues it automatically.
247    #[error("delay sleeping for run {run_id}, step {step_id}: wake at {wake_at}")]
248    DelaySleeping {
249        /// The run that is sleeping.
250        run_id: uuid::Uuid,
251        /// The delay step.
252        step_id: uuid::Uuid,
253        /// When the run should be woken up.
254        wake_at: chrono::DateTime<chrono::Utc>,
255    },
256
257    /// An agent step found every targeted provider account rate limited and
258    /// suspended the run until capacity comes back.
259    ///
260    /// The engine transitions the run to
261    /// [`Sleeping`](ironflow_store::entities::RunStatus::Sleeping), sets
262    /// `scheduled_at` to `wake_at` and records `kind` in
263    /// `capacity_wait_kind`, so adding or re-enabling an account of that kind
264    /// wakes it early. On wake, the step runs again from zero at the same
265    /// position.
266    ///
267    /// # Examples
268    ///
269    /// ```
270    /// use chrono::Utc;
271    /// use ironflow_engine::error::EngineError;
272    /// use uuid::Uuid;
273    ///
274    /// let err = EngineError::CapacitySleeping {
275    ///     run_id: Uuid::nil(),
276    ///     step_id: Uuid::nil(),
277    ///     kind: "claude".to_string(),
278    ///     wake_at: Utc::now(),
279    /// };
280    /// assert!(err.is_suspension());
281    /// ```
282    #[error("run {run_id} waiting for {kind} capacity at step {step_id}: wake at {wake_at}")]
283    CapacitySleeping {
284        /// The run that is sleeping.
285        run_id: Uuid,
286        /// The agent step that found no capacity.
287        step_id: Uuid,
288        /// Provider kind the run waits for (e.g. `claude`).
289        kind: String,
290        /// When the run should be woken up.
291        wake_at: DateTime<Utc>,
292    },
293
294    /// A signal step suspended the run until a matching signal or its deadline.
295    ///
296    /// Raised by
297    /// [`WorkflowContext::wait_for_signal`](crate::context::WorkflowContext::wait_for_signal)
298    /// when no signal was received yet. The engine transitions the run to
299    /// [`Sleeping`](ironflow_store::entities::RunStatus::Sleeping) with
300    /// `scheduled_at` set to the deadline: a delivery wakes it earlier.
301    #[error("run {run_id} waiting for signal {name:?} with key {key:?} until {deadline_at}")]
302    SignalWaiting {
303        /// The run that is waiting.
304        run_id: Uuid,
305        /// The signal step the run waits on.
306        step_id: Uuid,
307        /// Name of the signal step.
308        step_name: String,
309        /// The awaited signal name.
310        name: String,
311        /// The awaited occurrence key.
312        key: String,
313        /// When the wait times out.
314        deadline_at: DateTime<Utc>,
315    },
316
317    /// A child run started by
318    /// [`WorkflowContext::workflow`](crate::context::WorkflowContext::workflow)
319    /// suspended (approval, human input, delay or signal), and the parent run
320    /// is suspended with it.
321    ///
322    /// The child keeps its own suspension status and the parent's `Workflow`
323    /// step stays open. Resuming the child requeues the root run, which
324    /// replays and re-enters the same child run.
325    ///
326    /// # Examples
327    ///
328    /// ```
329    /// use ironflow_engine::error::EngineError;
330    /// use uuid::Uuid;
331    ///
332    /// let err = EngineError::ChildSuspended {
333    ///     run_id: Uuid::nil(),
334    ///     cause: Box::new(EngineError::HumanInputRequired {
335    ///         run_id: Uuid::nil(),
336    ///         step_id: Uuid::nil(),
337    ///         message: "Answer the questions".to_string(),
338    ///     }),
339    /// };
340    /// assert!(err.is_suspension());
341    /// ```
342    #[error("child run {run_id} suspended: {cause}")]
343    ChildSuspended {
344        /// The direct child run that suspended.
345        run_id: Uuid,
346        /// Why the child suspended: a leaf suspension
347        /// ([`ApprovalRequired`](EngineError::ApprovalRequired),
348        /// [`HumanInputRequired`](EngineError::HumanInputRequired),
349        /// [`DelaySleeping`](EngineError::DelaySleeping),
350        /// [`CapacitySleeping`](EngineError::CapacitySleeping),
351        /// [`SignalWaiting`](EngineError::SignalWaiting)) or a nested
352        /// `ChildSuspended` for a grand-child.
353        cause: Box<EngineError>,
354    },
355
356    /// The child run of a `Workflow` step was cancelled (`POST /runs/:id/cancel`
357    /// on the child itself), while it ran or while its chain was suspended.
358    ///
359    /// The step fails with this error and so does the parent, unless the step
360    /// tolerates the failure with `allow_failure`: it then completes with a
361    /// `Cancelled` [`SubWorkflowOutput`](crate::executor::SubWorkflowOutput).
362    /// Not retried: a new attempt would start a child the user just stopped.
363    ///
364    /// # Examples
365    ///
366    /// ```
367    /// use ironflow_engine::error::EngineError;
368    /// use uuid::Uuid;
369    ///
370    /// let err = EngineError::ChildRunCancelled { run_id: Uuid::nil() };
371    /// assert!(err.to_string().contains("cancelled"));
372    /// ```
373    #[error("child run {run_id} was cancelled")]
374    ChildRunCancelled {
375        /// The cancelled child run.
376        run_id: Uuid,
377    },
378
379    /// An operator paused the run while it executed.
380    ///
381    /// Raised at the next step boundary: the step is not started and the run
382    /// stays `Paused`, neither failed nor retried, until it is resumed or
383    /// cancelled. A `Workflow` step whose child is paused is left as the pause
384    /// recorded it, so the resumed parent re-enters the same child run.
385    ///
386    /// # Examples
387    ///
388    /// ```
389    /// use ironflow_engine::error::EngineError;
390    /// use uuid::Uuid;
391    ///
392    /// let err = EngineError::RunPaused { run_id: Uuid::nil() };
393    /// assert!(err.to_string().contains("paused"));
394    /// ```
395    #[error("run {run_id} was paused")]
396    RunPaused {
397        /// The paused run.
398        run_id: Uuid,
399    },
400
401    /// A pause or a resume targeted a sub-workflow run.
402    ///
403    /// A child run executes inside its root run's execution, so it is paused
404    /// and resumed with its root, never on its own.
405    ///
406    /// # Examples
407    ///
408    /// ```
409    /// use ironflow_engine::error::EngineError;
410    /// use uuid::Uuid;
411    ///
412    /// let err = EngineError::ChildRunNotPausable {
413    ///     run_id: Uuid::nil(),
414    ///     root_run_id: Uuid::nil(),
415    /// };
416    /// assert!(err.to_string().contains("root run"));
417    /// ```
418    #[error("run {run_id} is a sub-workflow run: pause or resume its root run {root_run_id}")]
419    ChildRunNotPausable {
420        /// The sub-workflow run targeted.
421        run_id: Uuid,
422        /// The root run of its chain, the one to pause or resume.
423        root_run_id: Uuid,
424    },
425
426    /// A signal could not be delivered because it is malformed (empty name or
427    /// key).
428    #[error("invalid signal: {0}")]
429    InvalidSignal(String),
430
431    /// A workflow invocation was rejected by the [workflow guard](crate::guard).
432    ///
433    /// The run is transitioned to
434    /// [`Cancelled`](ironflow_store::entities::RunStatus::Cancelled) when this
435    /// error is raised.
436    #[error("{WORKFLOW_GUARD_REJECTED_CODE}: {0}")]
437    WorkflowGuardRejected(#[from] WorkflowRejection),
438
439    /// The step stored at `position` does not match the step the handler just
440    /// called: its name or its `StepKind` differ from what was recorded before
441    /// the run was suspended.
442    ///
443    /// Raised instead of silently serving another step's cached output when the
444    /// handler's code changed while the run was suspended (a deploy during a
445    /// pending approval, human input, decision escalation or delay): step
446    /// positions can shift, and without this check every step after the
447    /// divergence point would silently receive another step's output, vote or
448    /// approval.
449    #[error(
450        "replay divergence at position {position}: handler called '{expected}' but the run \
451         recorded '{recorded}' (handler changed since the run was suspended?)"
452    )]
453    ReplayDivergence {
454        /// The step position where the recorded step and the step just called
455        /// stopped matching.
456        position: u32,
457        /// Identity (`name (kind)`) of the step the handler just called.
458        expected: String,
459        /// Identity (`name (kind)`) of the step recorded at `position`.
460        recorded: String,
461    },
462
463    /// The run's handler changed since the run was created or suspended, and the
464    /// handler's current version is not declared compatible with the version the
465    /// run was created with.
466    ///
467    /// Checked before any step is replayed, on the same rule the manual retry
468    /// endpoint applies (`HANDLER_VERSION_MISMATCH`). Unlike retry, resume has no
469    /// `force` override: replaying an incompatible handler's steps risks serving
470    /// one step's cached output to another (see
471    /// [`ReplayDivergence`](EngineError::ReplayDivergence)).
472    #[error(
473        "{HANDLER_VERSION_MISMATCH_CODE}: run {run_id} for handler '{workflow_name}' was created \
474         with version {run_version}, but the handler is now at version {current_version}; \
475         resume refused (no force override for resume)"
476    )]
477    HandlerVersionMismatch {
478        /// The run that cannot be resumed.
479        run_id: uuid::Uuid,
480        /// The handler's registered name.
481        workflow_name: String,
482        /// The handler version the run was created with.
483        run_version: String,
484        /// The handler's current version.
485        current_version: String,
486    },
487}
488
489impl From<StoreError> for EngineError {
490    fn from(err: StoreError) -> Self {
491        match err {
492            StoreError::ConcurrencyConflict { key, run_id } => {
493                EngineError::ConcurrencyConflict { key, run_id }
494            }
495            StoreError::InvalidConcurrencyLimit(e) => EngineError::InvalidConcurrencyLimit(e),
496            StoreError::InvalidWorkerTag(e) => EngineError::InvalidWorkerTag(e),
497            other => EngineError::Store(other),
498        }
499    }
500}
501
502impl EngineError {
503    /// Whether this error suspends the run instead of failing it.
504    ///
505    /// True for [`ApprovalRequired`](EngineError::ApprovalRequired),
506    /// [`HumanInputRequired`](EngineError::HumanInputRequired),
507    /// [`DelaySleeping`](EngineError::DelaySleeping),
508    /// [`CapacitySleeping`](EngineError::CapacitySleeping),
509    /// [`SignalWaiting`](EngineError::SignalWaiting) and
510    /// [`ChildSuspended`](EngineError::ChildSuspended).
511    ///
512    /// # Examples
513    ///
514    /// ```
515    /// use ironflow_engine::error::EngineError;
516    ///
517    /// assert!(!EngineError::StepConfig("bad".to_string()).is_suspension());
518    /// ```
519    pub fn is_suspension(&self) -> bool {
520        matches!(
521            self,
522            EngineError::ApprovalRequired { .. }
523                | EngineError::HumanInputRequired { .. }
524                | EngineError::DelaySleeping { .. }
525                | EngineError::CapacitySleeping { .. }
526                | EngineError::SignalWaiting { .. }
527                | EngineError::ChildSuspended { .. }
528        )
529    }
530
531    /// The leaf suspension behind a chain of
532    /// [`ChildSuspended`](EngineError::ChildSuspended) errors.
533    ///
534    /// Returns `self` for any error that is not `ChildSuspended`.
535    ///
536    /// # Examples
537    ///
538    /// ```
539    /// use ironflow_engine::error::EngineError;
540    /// use uuid::Uuid;
541    ///
542    /// let err = EngineError::ChildSuspended {
543    ///     run_id: Uuid::nil(),
544    ///     cause: Box::new(EngineError::ApprovalRequired {
545    ///         run_id: Uuid::nil(),
546    ///         step_id: Uuid::nil(),
547    ///         message: "deploy?".to_string(),
548    ///     }),
549    /// };
550    /// assert!(matches!(err.suspension_leaf(), EngineError::ApprovalRequired { .. }));
551    /// ```
552    pub fn suspension_leaf(&self) -> &EngineError {
553        let mut current = self;
554        while let EngineError::ChildSuspended { cause, .. } = current {
555            current = cause;
556        }
557        current
558    }
559
560    /// The status a run suspended by this error takes: `Sleeping` when the
561    /// leaf suspension is a delay, a capacity wait or a signal, `AwaitingApproval` otherwise (a
562    /// gate a human resolves).
563    pub(crate) fn suspension_status(&self) -> RunStatus {
564        match self.suspension_leaf() {
565            EngineError::DelaySleeping { .. }
566            | EngineError::CapacitySleeping { .. }
567            | EngineError::SignalWaiting { .. } => RunStatus::Sleeping,
568            _ => RunStatus::AwaitingApproval,
569        }
570    }
571}
572
573#[cfg(test)]
574mod tests {
575    use super::*;
576    use crate::retry_policy::is_run_retryable;
577
578    #[test]
579    fn invalid_workflow_display() {
580        let err = EngineError::InvalidWorkflow("unknown-handler".to_string());
581        assert!(err.to_string().contains("invalid workflow"));
582        assert!(err.to_string().contains("unknown-handler"));
583    }
584
585    #[test]
586    fn step_config_display() {
587        let err = EngineError::StepConfig("bad shell config".to_string());
588        assert!(err.to_string().contains("step config error"));
589        assert!(err.to_string().contains("bad shell config"));
590    }
591
592    #[test]
593    fn human_input_required_display() {
594        let err = EngineError::HumanInputRequired {
595            run_id: uuid::Uuid::nil(),
596            step_id: uuid::Uuid::nil(),
597            message: "Answer the questions".to_string(),
598        };
599        let text = err.to_string();
600        assert!(text.contains("human input required"));
601        assert!(text.contains("Answer the questions"));
602    }
603
604    #[test]
605    fn human_input_rejected_display() {
606        let err = EngineError::HumanInputRejected {
607            run_id: uuid::Uuid::nil(),
608            step_id: uuid::Uuid::nil(),
609            reason: "not relevant".to_string(),
610        };
611        let text = err.to_string();
612        assert!(text.contains("human input rejected"));
613        assert!(text.contains("not relevant"));
614    }
615
616    #[test]
617    fn store_error_from_conversion() {
618        let store_err = StoreError::RunNotFound(uuid::Uuid::nil());
619        let engine_err = EngineError::from(store_err);
620        assert!(engine_err.to_string().contains("store error"));
621    }
622
623    #[test]
624    fn store_concurrency_conflict_converts_to_engine_variant() {
625        let run_id = Uuid::now_v7();
626        let engine_err = EngineError::from(StoreError::ConcurrencyConflict {
627            key: "issue:12".to_string(),
628            run_id,
629        });
630        match engine_err {
631            EngineError::ConcurrencyConflict {
632                ref key,
633                run_id: holder,
634            } => {
635                assert_eq!(key, "issue:12");
636                assert_eq!(holder, run_id);
637            }
638            ref other => panic!("expected ConcurrencyConflict, got {other:?}"),
639        }
640        assert!(engine_err.to_string().contains("issue:12"));
641        assert!(!is_run_retryable(&engine_err));
642    }
643
644    #[test]
645    fn store_invalid_concurrency_limit_converts_to_engine_variant() {
646        let engine_err = EngineError::from(StoreError::InvalidConcurrencyLimit(
647            ConcurrencyLimitError::ZeroLimit {
648                group: "repo:acme".to_string(),
649            },
650        ));
651        assert!(
652            matches!(
653                engine_err,
654                EngineError::InvalidConcurrencyLimit(ConcurrencyLimitError::ZeroLimit { .. })
655            ),
656            "{engine_err:?}"
657        );
658        assert!(engine_err.to_string().contains("repo:acme"));
659        assert!(!is_run_retryable(&engine_err));
660    }
661
662    #[test]
663    fn invalid_priority_is_not_retryable() {
664        let err = EngineError::InvalidPriority("priority must be between -100 and 100".to_string());
665        assert_eq!(
666            err.to_string(),
667            "invalid priority: priority must be between -100 and 100"
668        );
669        assert!(!is_run_retryable(&err));
670    }
671
672    #[test]
673    fn store_invalid_worker_tag_converts_to_engine_variant() {
674        let engine_err =
675            EngineError::from(StoreError::InvalidWorkerTag(WorkerTagError::InvalidChar {
676                tag: "bad,tag".to_string(),
677            }));
678        assert!(
679            matches!(
680                engine_err,
681                EngineError::InvalidWorkerTag(WorkerTagError::InvalidChar { .. })
682            ),
683            "{engine_err:?}"
684        );
685        assert!(engine_err.to_string().contains("bad,tag"));
686        assert!(!is_run_retryable(&engine_err));
687    }
688
689    #[test]
690    fn run_budget_exceeded_display_carries_code_and_amounts() {
691        let err = EngineError::RunBudgetExceeded {
692            run_id: uuid::Uuid::nil(),
693            limit_usd: Decimal::new(200, 2),
694            spent_usd: Decimal::new(180, 2),
695            step_budget_usd: Decimal::new(50, 2),
696        };
697
698        let msg = err.to_string();
699        assert!(msg.contains(RUN_BUDGET_EXCEEDED_CODE));
700        assert!(msg.contains("2.00"));
701        assert!(msg.contains("1.80"));
702        assert!(msg.contains("0.50"));
703    }
704
705    #[test]
706    fn monthly_budget_exceeded_display_carries_code_and_amounts() {
707        let err = EngineError::MonthlyBudgetExceeded {
708            limit_usd: Decimal::new(10000, 2),
709            spent_usd: Decimal::new(10500, 2),
710        };
711
712        let msg = err.to_string();
713        assert!(msg.contains(MONTHLY_BUDGET_EXCEEDED_CODE));
714        assert!(msg.contains("100.00"));
715        assert!(msg.contains("105.00"));
716    }
717
718    #[test]
719    fn missing_artifact_display_names_the_step_and_pattern() {
720        let err = EngineError::MissingArtifact {
721            step: "build".to_string(),
722            pattern: "target/report.html".to_string(),
723        };
724
725        let msg = err.to_string();
726        assert!(msg.contains("\"build\""));
727        assert!(msg.contains("target/report.html"));
728    }
729
730    #[test]
731    fn artifact_not_declared_display_names_the_step_and_artifact() {
732        let err = EngineError::ArtifactNotDeclared {
733            step: "build".to_string(),
734            name: "report.htm".to_string(),
735        };
736
737        let msg = err.to_string();
738        assert!(msg.contains("\"build\""));
739        assert!(msg.contains("\"report.htm\""));
740    }
741
742    #[test]
743    fn artifact_not_found_display_names_the_producer() {
744        let err = EngineError::ArtifactNotFound {
745            step: "build".to_string(),
746            name: "report.html".to_string(),
747        };
748
749        let msg = err.to_string();
750        assert!(msg.contains("\"build\""));
751        assert!(msg.contains("report.html"));
752    }
753
754    #[test]
755    fn artifacts_unavailable_display() {
756        let err = EngineError::ArtifactsUnavailable("no blob store".to_string());
757        assert!(err.to_string().contains("not configured"));
758    }
759
760    #[test]
761    fn artifact_error_from_conversion() {
762        let engine_err = EngineError::from(ArtifactError::NotFound("a/b".to_string()));
763        assert!(engine_err.to_string().contains("artifact storage error"));
764    }
765
766    #[test]
767    fn serialization_error_from_conversion() {
768        let serde_err = serde_json::from_str::<String>("not json").unwrap_err();
769        let engine_err = EngineError::from(serde_err);
770        assert!(engine_err.to_string().contains("serialization error"));
771    }
772
773    #[test]
774    fn workflow_guard_rejected_display_carries_code_and_detail() {
775        use crate::guard::WorkflowRejection;
776
777        let rejection = WorkflowRejection::MaxDepthExceeded { depth: 6, max: 5 };
778        let err = EngineError::from(rejection);
779
780        let msg = err.to_string();
781        assert!(msg.contains(WORKFLOW_GUARD_REJECTED_CODE));
782        assert!(msg.contains("max call depth exceeded"));
783        assert!(msg.contains("6/5"));
784    }
785
786    #[test]
787    fn workflow_guard_rejected_from_conversion() {
788        use crate::guard::WorkflowRejection;
789
790        let rejection = WorkflowRejection::CycleDetected {
791            target: "wf-b".to_string(),
792            chain: vec!["wf-a".to_string(), "wf-b".to_string()],
793        };
794        let engine_err = EngineError::from(rejection);
795        assert!(engine_err.to_string().contains("cycle detected"));
796    }
797
798    #[test]
799    fn replay_divergence_display_carries_position_and_identities() {
800        let err = EngineError::ReplayDivergence {
801            position: 5,
802            expected: "resolve-base-branch (Shell)".to_string(),
803            recorded: "create-worktree (Shell)".to_string(),
804        };
805
806        let msg = err.to_string();
807        assert!(msg.contains("divergence"));
808        assert!(msg.contains("position 5"));
809        assert!(msg.contains("resolve-base-branch"));
810        assert!(msg.contains("create-worktree"));
811    }
812
813    #[test]
814    fn handler_version_mismatch_display_carries_code_and_versions() {
815        let err = EngineError::HandlerVersionMismatch {
816            run_id: uuid::Uuid::nil(),
817            workflow_name: "deploy".to_string(),
818            run_version: "1.0.0".to_string(),
819            current_version: "2.0.0".to_string(),
820        };
821
822        let msg = err.to_string();
823        assert!(msg.contains(HANDLER_VERSION_MISMATCH_CODE));
824        assert!(msg.contains("1.0.0"));
825        assert!(msg.contains("2.0.0"));
826    }
827
828    fn human_input_required() -> EngineError {
829        EngineError::HumanInputRequired {
830            run_id: Uuid::nil(),
831            step_id: Uuid::nil(),
832            message: "Answer the questions".to_string(),
833        }
834    }
835
836    #[test]
837    fn child_suspended_display_carries_child_and_cause() {
838        let child = Uuid::now_v7();
839        let err = EngineError::ChildSuspended {
840            run_id: child,
841            cause: Box::new(human_input_required()),
842        };
843
844        let msg = err.to_string();
845        assert!(msg.contains(&child.to_string()));
846        assert!(msg.contains("human input required"));
847    }
848
849    #[test]
850    fn leaf_suspensions_and_child_suspended_are_suspensions() {
851        let wake_at = Utc::now();
852        let suspensions = [
853            EngineError::ApprovalRequired {
854                run_id: Uuid::nil(),
855                step_id: Uuid::nil(),
856                message: "deploy?".to_string(),
857            },
858            human_input_required(),
859            EngineError::DelaySleeping {
860                run_id: Uuid::nil(),
861                step_id: Uuid::nil(),
862                wake_at,
863            },
864            EngineError::CapacitySleeping {
865                run_id: Uuid::nil(),
866                step_id: Uuid::nil(),
867                kind: "claude".to_string(),
868                wake_at,
869            },
870            EngineError::SignalWaiting {
871                run_id: Uuid::nil(),
872                step_id: Uuid::nil(),
873                step_name: "wait".to_string(),
874                name: "payment".to_string(),
875                key: "order-1".to_string(),
876                deadline_at: wake_at,
877            },
878            EngineError::ChildSuspended {
879                run_id: Uuid::nil(),
880                cause: Box::new(human_input_required()),
881            },
882        ];
883        for err in &suspensions {
884            assert!(err.is_suspension(), "{err} should be a suspension");
885        }
886    }
887
888    #[test]
889    fn failures_and_rejections_are_not_suspensions() {
890        let failures = [
891            EngineError::InvalidWorkflow("x".to_string()),
892            EngineError::HumanInputRejected {
893                run_id: Uuid::nil(),
894                step_id: Uuid::nil(),
895                reason: "no".to_string(),
896            },
897            EngineError::ApprovalRejected {
898                run_id: Uuid::nil(),
899                step_id: Uuid::nil(),
900                reason: "no".to_string(),
901            },
902        ];
903        for err in &failures {
904            assert!(!err.is_suspension(), "{err} should not be a suspension");
905        }
906    }
907
908    #[test]
909    fn suspension_leaf_unwraps_nested_child_suspensions() {
910        let err = EngineError::ChildSuspended {
911            run_id: Uuid::now_v7(),
912            cause: Box::new(EngineError::ChildSuspended {
913                run_id: Uuid::now_v7(),
914                cause: Box::new(human_input_required()),
915            }),
916        };
917
918        assert!(matches!(
919            err.suspension_leaf(),
920            EngineError::HumanInputRequired { .. }
921        ));
922    }
923
924    #[test]
925    fn suspension_status_follows_the_leaf() {
926        let human = EngineError::ChildSuspended {
927            run_id: Uuid::nil(),
928            cause: Box::new(human_input_required()),
929        };
930        assert_eq!(human.suspension_status(), RunStatus::AwaitingApproval);
931
932        let delay = EngineError::ChildSuspended {
933            run_id: Uuid::nil(),
934            cause: Box::new(EngineError::DelaySleeping {
935                run_id: Uuid::nil(),
936                step_id: Uuid::nil(),
937                wake_at: Utc::now(),
938            }),
939        };
940        assert_eq!(delay.suspension_status(), RunStatus::Sleeping);
941
942        let capacity = EngineError::ChildSuspended {
943            run_id: Uuid::nil(),
944            cause: Box::new(EngineError::CapacitySleeping {
945                run_id: Uuid::nil(),
946                step_id: Uuid::nil(),
947                kind: "claude".to_string(),
948                wake_at: Utc::now(),
949            }),
950        };
951        assert_eq!(capacity.suspension_status(), RunStatus::Sleeping);
952    }
953
954    #[test]
955    fn suspension_leaf_of_a_plain_error_is_itself() {
956        let err = EngineError::StepConfig("bad".to_string());
957        assert!(matches!(err.suspension_leaf(), EngineError::StepConfig(_)));
958    }
959}