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}