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;
11
12use crate::guard::{WORKFLOW_GUARD_REJECTED_CODE, WorkflowRejection};
13
14/// Business error code carried by [`EngineError::RunBudgetExceeded`].
15pub const RUN_BUDGET_EXCEEDED_CODE: &str = "RUN_BUDGET_EXCEEDED";
16
17/// Business error code carried by [`EngineError::MonthlyBudgetExceeded`].
18pub const MONTHLY_BUDGET_EXCEEDED_CODE: &str = "MONTHLY_BUDGET_EXCEEDED";
19
20/// Business error code for handler-version mismatch on retry.
21pub const HANDLER_VERSION_MISMATCH_CODE: &str = "HANDLER_VERSION_MISMATCH";
22
23/// Errors produced by the workflow engine.
24#[derive(Debug, Error)]
25pub enum EngineError {
26    /// An operation (Shell, Http, Agent) failed during step execution.
27    #[error("operation failed: {0}")]
28    Operation(#[from] OperationError),
29
30    /// The backing store returned an error.
31    #[error("store error: {0}")]
32    Store(#[from] StoreError),
33
34    /// The workflow definition is invalid.
35    #[error("invalid workflow: {0}")]
36    InvalidWorkflow(String),
37
38    /// A step configuration could not be deserialized for execution.
39    #[error("step config error: {0}")]
40    StepConfig(String),
41
42    /// A decision answer was accessed by name with the wrong type or a missing key.
43    #[error("decision error: {0}")]
44    Decision(#[from] ironflow_core::error::DecisionError),
45
46    /// A decision step was reached but no [`DecisionProvider`](ironflow_core::decision::DecisionProvider)
47    /// is wired into the engine.
48    #[error(
49        "decision step '{step}' requires a decision provider; \
50         wire one with Engine::with_decision_provider(...)"
51    )]
52    NoDecisionProvider {
53        /// The decision step that could not run.
54        step: String,
55    },
56
57    /// JSON serialization error.
58    #[error("serialization error: {0}")]
59    Serialization(#[from] serde_json::Error),
60
61    /// The run reached its cumulative cost cap before launching an agent step.
62    ///
63    /// Raised *before* the step is created, so no work and no spend happen.
64    /// The engine transitions the run to
65    /// [`Cancelled`](ironflow_store::entities::RunStatus::Cancelled).
66    #[error(
67        "{RUN_BUDGET_EXCEEDED_CODE}: run {run_id} would exceed its cost cap \
68         (spent {spent_usd} USD + next step {step_budget_usd} USD > cap {limit_usd} USD)"
69    )]
70    RunBudgetExceeded {
71        /// The run that hit its cap.
72        run_id: uuid::Uuid,
73        /// The configured cap, in USD.
74        limit_usd: Decimal,
75        /// Cost already accumulated by this run and its ancestors, in USD.
76        spent_usd: Decimal,
77        /// Declared budget of the step that was about to run, in USD.
78        step_budget_usd: Decimal,
79    },
80
81    /// The global monthly cost quota is exhausted; no new run may be created.
82    ///
83    /// Runs already in flight are never interrupted by this error.
84    #[error(
85        "{MONTHLY_BUDGET_EXCEEDED_CODE}: monthly cost quota exhausted \
86         ({spent_usd} USD spent of {limit_usd} USD)"
87    )]
88    MonthlyBudgetExceeded {
89        /// The configured monthly quota, in USD.
90        limit_usd: Decimal,
91        /// Cost already spent during the current calendar month, in USD.
92        spent_usd: Decimal,
93    },
94
95    /// A step declared an output that produced no file.
96    ///
97    /// Raised only when the step itself succeeded: a declared output that never
98    /// materialised is a broken contract, and failing here beats failing later
99    /// in whichever step tried to consume it.
100    #[error("step {step:?} declared output {pattern:?} but no file matched")]
101    MissingArtifact {
102        /// Name of the step that declared the output.
103        step: String,
104        /// The unmatched pattern.
105        pattern: String,
106    },
107
108    /// A handle was asked for an artifact the step never declared.
109    #[error("step {step:?} declares no artifact output named {name:?}")]
110    ArtifactNotDeclared {
111        /// Name of the step the handle was asked from.
112        step: String,
113        /// Name of the artifact that was asked for.
114        name: String,
115    },
116
117    /// A step asked for an artifact that no earlier step produced.
118    #[error("no artifact {name:?} produced by step {step:?} before this point")]
119    ArtifactNotFound {
120        /// Name of the producing step that was searched for.
121        step: String,
122        /// Name of the artifact that was searched for.
123        name: String,
124    },
125
126    /// Artifacts were used but no storage backend is configured.
127    #[error("artifact storage is not configured: {0}")]
128    ArtifactsUnavailable(String),
129
130    /// The artifact storage backend failed.
131    #[error("artifact storage error: {0}")]
132    Artifact(#[from] ArtifactError),
133
134    /// The run requires human approval before continuing.
135    #[error("approval required for run {run_id}, step {step_id}: {message}")]
136    ApprovalRequired {
137        /// The run that is awaiting approval.
138        run_id: uuid::Uuid,
139        /// The approval step that triggered the pause.
140        step_id: uuid::Uuid,
141        /// The approval message.
142        message: String,
143    },
144
145    /// An approval gate was rejected instead of granted.
146    ///
147    /// Raised when a [`StepInterceptor`](crate::executor::StepInterceptor) resolves
148    /// the gate with [`ApprovalOutcome::Rejected`](crate::executor::ApprovalOutcome::Rejected).
149    /// The engine fails the run; the rejection is deterministic, so the run is
150    /// never replayed.
151    #[error("approval rejected for run {run_id}, step {step_id}: {reason}")]
152    ApprovalRejected {
153        /// The run that was stopped.
154        run_id: uuid::Uuid,
155        /// The approval step that was rejected.
156        step_id: uuid::Uuid,
157        /// Why the gate was refused.
158        reason: String,
159    },
160
161    /// The run waits for a typed human input before continuing.
162    ///
163    /// Raised by [`WorkflowContext::human_input`](crate::context::WorkflowContext::human_input)
164    /// when no answer has been given yet. The engine transitions the run to
165    /// [`AwaitingApproval`](ironflow_store::entities::RunStatus::AwaitingApproval).
166    #[error("human input required for run {run_id}, step {step_id}: {message}")]
167    HumanInputRequired {
168        /// The run that is awaiting input.
169        run_id: uuid::Uuid,
170        /// The human input step that triggered the pause.
171        step_id: uuid::Uuid,
172        /// The message displayed to the person answering.
173        message: String,
174    },
175
176    /// A human input request was rejected instead of answered.
177    ///
178    /// Returned to the handler by
179    /// [`WorkflowContext::human_input`](crate::context::WorkflowContext::human_input),
180    /// which decides what happens next. A propagated rejection fails the run,
181    /// and the run is never retried.
182    #[error("human input rejected for run {run_id}, step {step_id}: {reason}")]
183    HumanInputRejected {
184        /// The run the input belongs to.
185        run_id: uuid::Uuid,
186        /// The human input step that was rejected.
187        step_id: uuid::Uuid,
188        /// Why the input was refused.
189        reason: String,
190    },
191
192    /// A delay step suspended the run until the given time.
193    ///
194    /// The engine transitions the run to
195    /// [`Sleeping`](ironflow_store::entities::RunStatus::Sleeping) and sets
196    /// `scheduled_at` so the worker re-queues it automatically.
197    #[error("delay sleeping for run {run_id}, step {step_id}: wake at {wake_at}")]
198    DelaySleeping {
199        /// The run that is sleeping.
200        run_id: uuid::Uuid,
201        /// The delay step.
202        step_id: uuid::Uuid,
203        /// When the run should be woken up.
204        wake_at: chrono::DateTime<chrono::Utc>,
205    },
206
207    /// A signal step suspended the run until a matching signal or its deadline.
208    ///
209    /// Raised by
210    /// [`WorkflowContext::wait_for_signal`](crate::context::WorkflowContext::wait_for_signal)
211    /// when no signal was received yet. The engine transitions the run to
212    /// [`Sleeping`](ironflow_store::entities::RunStatus::Sleeping) with
213    /// `scheduled_at` set to the deadline: a delivery wakes it earlier.
214    #[error("run {run_id} waiting for signal {name:?} with key {key:?} until {deadline_at}")]
215    SignalWaiting {
216        /// The run that is waiting.
217        run_id: Uuid,
218        /// The signal step the run waits on.
219        step_id: Uuid,
220        /// Name of the signal step.
221        step_name: String,
222        /// The awaited signal name.
223        name: String,
224        /// The awaited occurrence key.
225        key: String,
226        /// When the wait times out.
227        deadline_at: DateTime<Utc>,
228    },
229
230    /// A signal could not be delivered because it is malformed (empty name or
231    /// key).
232    #[error("invalid signal: {0}")]
233    InvalidSignal(String),
234
235    /// A workflow invocation was rejected by the [workflow guard](crate::guard).
236    ///
237    /// The run is transitioned to
238    /// [`Cancelled`](ironflow_store::entities::RunStatus::Cancelled) when this
239    /// error is raised.
240    #[error("{WORKFLOW_GUARD_REJECTED_CODE}: {0}")]
241    WorkflowGuardRejected(#[from] WorkflowRejection),
242
243    /// The step stored at `position` does not match the step the handler just
244    /// called: its name or its `StepKind` differ from what was recorded before
245    /// the run was suspended.
246    ///
247    /// Raised instead of silently serving another step's cached output when the
248    /// handler's code changed while the run was suspended (a deploy during a
249    /// pending approval, human input, decision escalation or delay): step
250    /// positions can shift, and without this check every step after the
251    /// divergence point would silently receive another step's output, vote or
252    /// approval.
253    #[error(
254        "replay divergence at position {position}: handler called '{expected}' but the run \
255         recorded '{recorded}' (handler changed since the run was suspended?)"
256    )]
257    ReplayDivergence {
258        /// The step position where the recorded step and the step just called
259        /// stopped matching.
260        position: u32,
261        /// Identity (`name (kind)`) of the step the handler just called.
262        expected: String,
263        /// Identity (`name (kind)`) of the step recorded at `position`.
264        recorded: String,
265    },
266
267    /// The run's handler changed since the run was created or suspended, and the
268    /// handler's current version is not declared compatible with the version the
269    /// run was created with.
270    ///
271    /// Checked before any step is replayed, on the same rule the manual retry
272    /// endpoint applies (`HANDLER_VERSION_MISMATCH`). Unlike retry, resume has no
273    /// `force` override: replaying an incompatible handler's steps risks serving
274    /// one step's cached output to another (see
275    /// [`ReplayDivergence`](EngineError::ReplayDivergence)).
276    #[error(
277        "{HANDLER_VERSION_MISMATCH_CODE}: run {run_id} for handler '{workflow_name}' was created \
278         with version {run_version}, but the handler is now at version {current_version}; \
279         resume refused (no force override for resume)"
280    )]
281    HandlerVersionMismatch {
282        /// The run that cannot be resumed.
283        run_id: uuid::Uuid,
284        /// The handler's registered name.
285        workflow_name: String,
286        /// The handler version the run was created with.
287        run_version: String,
288        /// The handler's current version.
289        current_version: String,
290    },
291}
292
293#[cfg(test)]
294mod tests {
295    use super::*;
296
297    #[test]
298    fn invalid_workflow_display() {
299        let err = EngineError::InvalidWorkflow("unknown-handler".to_string());
300        assert!(err.to_string().contains("invalid workflow"));
301        assert!(err.to_string().contains("unknown-handler"));
302    }
303
304    #[test]
305    fn step_config_display() {
306        let err = EngineError::StepConfig("bad shell config".to_string());
307        assert!(err.to_string().contains("step config error"));
308        assert!(err.to_string().contains("bad shell config"));
309    }
310
311    #[test]
312    fn human_input_required_display() {
313        let err = EngineError::HumanInputRequired {
314            run_id: uuid::Uuid::nil(),
315            step_id: uuid::Uuid::nil(),
316            message: "Answer the questions".to_string(),
317        };
318        let text = err.to_string();
319        assert!(text.contains("human input required"));
320        assert!(text.contains("Answer the questions"));
321    }
322
323    #[test]
324    fn human_input_rejected_display() {
325        let err = EngineError::HumanInputRejected {
326            run_id: uuid::Uuid::nil(),
327            step_id: uuid::Uuid::nil(),
328            reason: "not relevant".to_string(),
329        };
330        let text = err.to_string();
331        assert!(text.contains("human input rejected"));
332        assert!(text.contains("not relevant"));
333    }
334
335    #[test]
336    fn store_error_from_conversion() {
337        let store_err = StoreError::RunNotFound(uuid::Uuid::nil());
338        let engine_err = EngineError::from(store_err);
339        assert!(engine_err.to_string().contains("store error"));
340    }
341
342    #[test]
343    fn run_budget_exceeded_display_carries_code_and_amounts() {
344        let err = EngineError::RunBudgetExceeded {
345            run_id: uuid::Uuid::nil(),
346            limit_usd: Decimal::new(200, 2),
347            spent_usd: Decimal::new(180, 2),
348            step_budget_usd: Decimal::new(50, 2),
349        };
350
351        let msg = err.to_string();
352        assert!(msg.contains(RUN_BUDGET_EXCEEDED_CODE));
353        assert!(msg.contains("2.00"));
354        assert!(msg.contains("1.80"));
355        assert!(msg.contains("0.50"));
356    }
357
358    #[test]
359    fn monthly_budget_exceeded_display_carries_code_and_amounts() {
360        let err = EngineError::MonthlyBudgetExceeded {
361            limit_usd: Decimal::new(10000, 2),
362            spent_usd: Decimal::new(10500, 2),
363        };
364
365        let msg = err.to_string();
366        assert!(msg.contains(MONTHLY_BUDGET_EXCEEDED_CODE));
367        assert!(msg.contains("100.00"));
368        assert!(msg.contains("105.00"));
369    }
370
371    #[test]
372    fn missing_artifact_display_names_the_step_and_pattern() {
373        let err = EngineError::MissingArtifact {
374            step: "build".to_string(),
375            pattern: "target/report.html".to_string(),
376        };
377
378        let msg = err.to_string();
379        assert!(msg.contains("\"build\""));
380        assert!(msg.contains("target/report.html"));
381    }
382
383    #[test]
384    fn artifact_not_declared_display_names_the_step_and_artifact() {
385        let err = EngineError::ArtifactNotDeclared {
386            step: "build".to_string(),
387            name: "report.htm".to_string(),
388        };
389
390        let msg = err.to_string();
391        assert!(msg.contains("\"build\""));
392        assert!(msg.contains("\"report.htm\""));
393    }
394
395    #[test]
396    fn artifact_not_found_display_names_the_producer() {
397        let err = EngineError::ArtifactNotFound {
398            step: "build".to_string(),
399            name: "report.html".to_string(),
400        };
401
402        let msg = err.to_string();
403        assert!(msg.contains("\"build\""));
404        assert!(msg.contains("report.html"));
405    }
406
407    #[test]
408    fn artifacts_unavailable_display() {
409        let err = EngineError::ArtifactsUnavailable("no blob store".to_string());
410        assert!(err.to_string().contains("not configured"));
411    }
412
413    #[test]
414    fn artifact_error_from_conversion() {
415        let engine_err = EngineError::from(ArtifactError::NotFound("a/b".to_string()));
416        assert!(engine_err.to_string().contains("artifact storage error"));
417    }
418
419    #[test]
420    fn serialization_error_from_conversion() {
421        let serde_err = serde_json::from_str::<String>("not json").unwrap_err();
422        let engine_err = EngineError::from(serde_err);
423        assert!(engine_err.to_string().contains("serialization error"));
424    }
425
426    #[test]
427    fn workflow_guard_rejected_display_carries_code_and_detail() {
428        use crate::guard::WorkflowRejection;
429
430        let rejection = WorkflowRejection::MaxDepthExceeded { depth: 6, max: 5 };
431        let err = EngineError::from(rejection);
432
433        let msg = err.to_string();
434        assert!(msg.contains(WORKFLOW_GUARD_REJECTED_CODE));
435        assert!(msg.contains("max call depth exceeded"));
436        assert!(msg.contains("6/5"));
437    }
438
439    #[test]
440    fn workflow_guard_rejected_from_conversion() {
441        use crate::guard::WorkflowRejection;
442
443        let rejection = WorkflowRejection::CycleDetected {
444            target: "wf-b".to_string(),
445            chain: vec!["wf-a".to_string(), "wf-b".to_string()],
446        };
447        let engine_err = EngineError::from(rejection);
448        assert!(engine_err.to_string().contains("cycle detected"));
449    }
450
451    #[test]
452    fn replay_divergence_display_carries_position_and_identities() {
453        let err = EngineError::ReplayDivergence {
454            position: 5,
455            expected: "resolve-base-branch (Shell)".to_string(),
456            recorded: "create-worktree (Shell)".to_string(),
457        };
458
459        let msg = err.to_string();
460        assert!(msg.contains("divergence"));
461        assert!(msg.contains("position 5"));
462        assert!(msg.contains("resolve-base-branch"));
463        assert!(msg.contains("create-worktree"));
464    }
465
466    #[test]
467    fn handler_version_mismatch_display_carries_code_and_versions() {
468        let err = EngineError::HandlerVersionMismatch {
469            run_id: uuid::Uuid::nil(),
470            workflow_name: "deploy".to_string(),
471            run_version: "1.0.0".to_string(),
472            current_version: "2.0.0".to_string(),
473        };
474
475        let msg = err.to_string();
476        assert!(msg.contains(HANDLER_VERSION_MISMATCH_CODE));
477        assert!(msg.contains("1.0.0"));
478        assert!(msg.contains("2.0.0"));
479    }
480}