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;
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(&config, &self.provider, step_log_sender).await;
148            let handler_duration = start.elapsed().as_millis() as u64;
149            let completed_at = Utc::now();
150
151            match result {
152                Ok(output) => {
153                    if let Err(store_err) = self
154                        .store
155                        .update_step(
156                            step.id,
157                            StepUpdate {
158                                status: Some(StepStatus::Completed),
159                                output: Some(output.output),
160                                duration_ms: Some(handler_duration),
161                                cost_usd: Some(output.cost_usd),
162                                completed_at: Some(completed_at),
163                                ..StepUpdate::default()
164                            },
165                        )
166                        .await
167                    {
168                        warn!(
169                            run_id = %self.run_id,
170                            handler = %handler.name,
171                            error = %store_err,
172                            "failed to persist error handler completion"
173                        );
174                    }
175
176                    info!(
177                        run_id = %self.run_id,
178                        handler = %handler.name,
179                        duration_ms = handler_duration,
180                        "error handler completed"
181                    );
182                }
183                Err(err) => {
184                    if let Err(store_err) = self
185                        .store
186                        .update_step(
187                            step.id,
188                            StepUpdate {
189                                status: Some(StepStatus::Failed),
190                                error: Some(err.to_string()),
191                                duration_ms: Some(handler_duration),
192                                completed_at: Some(completed_at),
193                                ..StepUpdate::default()
194                            },
195                        )
196                        .await
197                    {
198                        warn!(
199                            run_id = %self.run_id,
200                            handler = %handler.name,
201                            error = %store_err,
202                            "failed to persist error handler failure"
203                        );
204                    }
205
206                    warn!(
207                        run_id = %self.run_id,
208                        handler = %handler.name,
209                        error = %err,
210                        "error handler failed (original error preserved)"
211                    );
212                }
213            }
214        }
215    }
216}
217
218/// Inject error context into a step config before executing it as an error handler.
219fn inject_error_context(
220    config: &mut StepConfig,
221    failed_step: &str,
222    error_msg: &str,
223    duration_ms: u64,
224) {
225    match config {
226        StepConfig::Shell(shell) => {
227            shell
228                .env
229                .push(("IRONFLOW_ERROR_STEP".to_string(), failed_step.to_string()));
230            shell
231                .env
232                .push(("IRONFLOW_ERROR_MESSAGE".to_string(), error_msg.to_string()));
233            shell.env.push((
234                "IRONFLOW_ERROR_DURATION_MS".to_string(),
235                duration_ms.to_string(),
236            ));
237        }
238        StepConfig::Agent(agent) => {
239            agent.prompt = format!(
240                "[Error Context]\nStep \"{}\" failed after {}ms:\n{}\n\n{}",
241                failed_step, duration_ms, error_msg, agent.prompt
242            );
243        }
244        StepConfig::Http(http) => {
245            http.headers
246                .push(("X-Ironflow-Error-Step".to_string(), failed_step.to_string()));
247            http.headers.push((
248                "X-Ironflow-Error-Message".to_string(),
249                error_msg.to_string(),
250            ));
251        }
252        StepConfig::Workflow(_)
253        | StepConfig::Approval(_)
254        | StepConfig::Decision(_)
255        | StepConfig::Delay(_) => {}
256    }
257}