Skip to main content

ironflow_engine/context/
error_handlers.rs

1//! Error handlers registered on [`WorkflowContext`].
2//!
3//! A handler registered with [`on_error`](WorkflowContext::on_error) fires once
4//! when a later step fails. Execution is best-effort: a handler failure is
5//! logged and never replaces the error that triggered it.
6
7use std::mem::take;
8use std::time::Instant;
9
10use chrono::Utc;
11use serde_json::json;
12use tracing::{info, warn};
13
14use ironflow_store::models::{NewStep, StepStatus, StepUpdate, step_trace_id};
15
16use crate::config::StepConfig;
17use crate::executor::execute_step_config_intercepted;
18use crate::log_sender::StepLogSender;
19
20use super::{OnErrorHandler, WorkflowContext};
21
22impl WorkflowContext {
23    /// Register an error handler that fires when any subsequent step fails.
24    ///
25    /// The handler is consumed after firing (fire-once). Multiple handlers
26    /// can be registered; they fire in registration order.
27    ///
28    /// Error handler execution is best-effort: if a handler fails, the error
29    /// is logged but the original step error is preserved. Error handler steps
30    /// appear in the run timeline with
31    /// [`Step::is_error_handler`](ironflow_store::models::Step::is_error_handler)
32    /// set to `true`.
33    ///
34    /// # Examples
35    ///
36    /// ```no_run
37    /// use ironflow_engine::context::WorkflowContext;
38    /// use ironflow_engine::config::ShellConfig;
39    /// use ironflow_engine::error::EngineError;
40    ///
41    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
42    /// ctx.on_error("cleanup", ShellConfig::new("rm -rf /tmp/build"));
43    /// ctx.shell("build", ShellConfig::new("cargo build")).await?;
44    /// # Ok(())
45    /// # }
46    /// ```
47    pub fn on_error(&mut self, name: &str, config: impl Into<StepConfig>) {
48        self.error_handlers.push(OnErrorHandler {
49            name: name.to_string(),
50            config: config.into(),
51        });
52    }
53
54    /// Remove all registered error handlers.
55    ///
56    /// # Examples
57    ///
58    /// ```no_run
59    /// use ironflow_engine::context::WorkflowContext;
60    /// use ironflow_engine::config::ShellConfig;
61    /// use ironflow_engine::error::EngineError;
62    ///
63    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
64    /// ctx.on_error("cleanup", ShellConfig::new("rm -rf /tmp/build"));
65    /// ctx.shell("build", ShellConfig::new("cargo build")).await?;
66    /// ctx.clear_error_handlers();
67    /// // cleanup will NOT fire if deploy fails
68    /// ctx.shell("deploy", ShellConfig::new("./deploy.sh")).await?;
69    /// # Ok(())
70    /// # }
71    /// ```
72    pub fn clear_error_handlers(&mut self) {
73        self.error_handlers.clear();
74    }
75
76    /// Execute all registered error handlers after a step failure.
77    ///
78    /// Drains the handler list (fire-once). Each handler creates its own
79    /// step record with `is_error_handler = true`. Handler failures are
80    /// logged but never propagated.
81    pub(super) async fn fire_error_handlers(
82        &mut self,
83        failed_step_name: &str,
84        error_msg: &str,
85        duration_ms: u64,
86    ) {
87        let handlers = take(&mut self.error_handlers);
88        if handlers.is_empty() {
89            return;
90        }
91
92        let error_context = json!({
93            "failed_step": failed_step_name,
94            "error": error_msg,
95            "duration_ms": duration_ms,
96        });
97
98        for handler in handlers {
99            let mut config = handler.config.clone();
100            inject_error_context(&mut config, failed_step_name, error_msg, duration_ms);
101
102            let position = self.position;
103            self.position += 1;
104
105            let trace_id = step_trace_id(self.run_id, &handler.name, position);
106            let step = match self
107                .store
108                .create_step(NewStep {
109                    run_id: self.run_id,
110                    trace_id,
111                    name: handler.name.clone(),
112                    kind: config.kind(),
113                    position,
114                    input: Some(error_context.clone()),
115                    is_error_handler: true,
116                })
117                .await
118            {
119                Ok(step) => step,
120                Err(err) => {
121                    warn!(
122                        run_id = %self.run_id,
123                        handler = %handler.name,
124                        error = %err,
125                        "failed to create error handler step"
126                    );
127                    continue;
128                }
129            };
130
131            if let Err(err) = self.start_step(step.id, Utc::now()).await {
132                warn!(
133                    run_id = %self.run_id,
134                    handler = %handler.name,
135                    error = %err,
136                    "failed to start error handler step"
137                );
138                continue;
139            }
140
141            let step_log_sender = self
142                .log_sender
143                .as_ref()
144                .map(|s| StepLogSender::new(s.clone(), self.run_id, step.id, handler.name.clone()));
145
146            let start = Instant::now();
147            let result = execute_step_config_intercepted(
148                &config,
149                &self.provider,
150                step_log_sender,
151                self.step_interceptor(),
152            )
153            .await;
154            let handler_duration = start.elapsed().as_millis() as u64;
155            let completed_at = Utc::now();
156
157            match result {
158                Ok(output) => {
159                    if let Err(store_err) = self
160                        .store
161                        .update_step(
162                            step.id,
163                            StepUpdate {
164                                status: Some(StepStatus::Completed),
165                                output: Some(output.output),
166                                duration_ms: Some(handler_duration),
167                                cost_usd: Some(output.cost_usd),
168                                completed_at: Some(completed_at),
169                                ..StepUpdate::default()
170                            },
171                        )
172                        .await
173                    {
174                        warn!(
175                            run_id = %self.run_id,
176                            handler = %handler.name,
177                            error = %store_err,
178                            "failed to persist error handler completion"
179                        );
180                    }
181
182                    info!(
183                        run_id = %self.run_id,
184                        handler = %handler.name,
185                        duration_ms = handler_duration,
186                        "error handler completed"
187                    );
188                }
189                Err(err) => {
190                    if let Err(store_err) = self
191                        .store
192                        .update_step(
193                            step.id,
194                            StepUpdate {
195                                status: Some(StepStatus::Failed),
196                                error: Some(err.to_string()),
197                                duration_ms: Some(handler_duration),
198                                completed_at: Some(completed_at),
199                                ..StepUpdate::default()
200                            },
201                        )
202                        .await
203                    {
204                        warn!(
205                            run_id = %self.run_id,
206                            handler = %handler.name,
207                            error = %store_err,
208                            "failed to persist error handler failure"
209                        );
210                    }
211
212                    warn!(
213                        run_id = %self.run_id,
214                        handler = %handler.name,
215                        error = %err,
216                        "error handler failed (original error preserved)"
217                    );
218                }
219            }
220        }
221    }
222}
223
224/// Inject error context into a step config before executing it as an error handler.
225fn inject_error_context(
226    config: &mut StepConfig,
227    failed_step: &str,
228    error_msg: &str,
229    duration_ms: u64,
230) {
231    match config {
232        StepConfig::Shell(shell) => {
233            shell
234                .env
235                .push(("IRONFLOW_ERROR_STEP".to_string(), failed_step.to_string()));
236            shell
237                .env
238                .push(("IRONFLOW_ERROR_MESSAGE".to_string(), error_msg.to_string()));
239            shell.env.push((
240                "IRONFLOW_ERROR_DURATION_MS".to_string(),
241                duration_ms.to_string(),
242            ));
243        }
244        StepConfig::Agent(agent) => {
245            agent.prompt = format!(
246                "[Error Context]\nStep \"{}\" failed after {}ms:\n{}\n\n{}",
247                failed_step, duration_ms, error_msg, agent.prompt
248            );
249        }
250        StepConfig::Http(http) => {
251            http.headers
252                .push(("X-Ironflow-Error-Step".to_string(), failed_step.to_string()));
253            http.headers.push((
254                "X-Ironflow-Error-Message".to_string(),
255                error_msg.to_string(),
256            ));
257        }
258        StepConfig::Workflow(_)
259        | StepConfig::Approval(_)
260        | StepConfig::Decision(_)
261        | StepConfig::Delay(_) => {}
262    }
263}