Skip to main content

ironflow_engine/executor/
mod.rs

1//! Step executor — reconstructs operations from configs and runs them.
2//!
3//! Each step type (shell, HTTP, agent) has its own executor implementing
4//! the [`StepExecutor`] trait. The [`execute_step_config`] function dispatches
5//! to the appropriate executor based on the [`StepConfig`] variant.
6//!
7//! Every executor declares the [`StepKind`] it handles via
8//! [`StepExecutor::kind`], and [`execute_step_config`] derives the span and
9//! metric label from that kind. A new step type therefore only has to declare
10//! its kind instead of extending a match here.
11
12mod agent;
13mod decision;
14mod http;
15mod interceptor;
16mod shell;
17mod step_artifacts;
18mod stored;
19mod workflow_output;
20
21use std::borrow::Cow;
22use std::future::Future;
23use std::sync::Arc;
24
25use rust_decimal::Decimal;
26use serde::de::DeserializeOwned;
27use serde_json::{Value, from_value};
28use tracing::Span;
29use uuid::Uuid;
30
31use ironflow_core::provider::{AgentProvider, DebugMessage};
32use ironflow_store::entities::{StepKind, StepStatus};
33
34use crate::config::StepConfig;
35use crate::error::EngineError;
36use crate::log_sender::StepLogSender;
37
38pub use agent::AgentExecutor;
39pub use decision::{DecisionExecution, execute_decision};
40pub use http::HttpExecutor;
41pub use interceptor::{ApprovalOutcome, HumanInputOutcome, SignalOutcome, StepInterceptor};
42pub use shell::ShellExecutor;
43pub use step_artifacts::StepArtifacts;
44pub use workflow_output::SubWorkflowOutput;
45
46/// Key under which a tolerated (`allow_failure`) step failure stores its message.
47pub(crate) const ERROR_KEY: &str = "error";
48
49/// Result of executing a single step.
50#[derive(Debug, Clone)]
51pub struct StepOutput {
52    /// Serialized output (stdout for shell, body for http, value for agent).
53    ///
54    /// For agent steps with a JSON schema, the value may not strictly conform
55    /// to the schema: Claude CLI can flatten wrapper objects with a single
56    /// array field, returning a bare array instead of `{"items": [...]}`.
57    /// Callers should handle both the expected wrapper and a bare value.
58    pub output: Value,
59    /// Wall-clock duration in milliseconds.
60    pub duration_ms: u64,
61    /// Cost in USD (agent steps only).
62    pub cost_usd: Decimal,
63    /// Uncached input token count (agent steps only).
64    pub input_tokens: Option<u64>,
65    /// Input tokens served from the prompt cache (agent steps only).
66    pub cache_read_input_tokens: Option<u64>,
67    /// Input tokens written to the prompt cache (agent steps only).
68    pub cache_creation_input_tokens: Option<u64>,
69    /// Output token count (agent steps only).
70    pub output_tokens: Option<u64>,
71    /// Model identifier used for agent steps (e.g. `"claude-sonnet-4-20250514"`).
72    pub model: Option<String>,
73    /// Conversation trace from verbose agent invocations.
74    pub debug_messages: Option<Vec<DebugMessage>>,
75    /// Artifacts the step can hand out through [`artifact`](Self::artifact).
76    /// Filled by the workflow context; executors leave the default.
77    pub artifacts: StepArtifacts,
78    /// Provider Account the agent step ran under, `None` for other steps and
79    /// agent steps that used the worker environment.
80    pub account_id: Option<Uuid>,
81}
82
83impl StepOutput {
84    /// Total tokens consumed by the step: uncached input, cache reads, cache
85    /// writes and output. Missing counts are treated as 0 and the sum saturates.
86    ///
87    /// # Examples
88    ///
89    /// ```
90    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
91    /// use rust_decimal::Decimal;
92    /// use serde_json::json;
93    ///
94    /// let output = StepOutput {
95    ///     output: json!("ok"),
96    ///     duration_ms: 10,
97    ///     cost_usd: Decimal::ZERO,
98    ///     input_tokens: Some(100),
99    ///     cache_read_input_tokens: Some(5000),
100    ///     cache_creation_input_tokens: Some(200),
101    ///     output_tokens: Some(50),
102    ///     model: None,
103    ///     debug_messages: None,
104    ///     artifacts: StepArtifacts::default(),
105    ///     account_id: None,
106    /// };
107    /// assert_eq!(output.total_tokens(), 5350);
108    /// ```
109    pub fn total_tokens(&self) -> u64 {
110        [
111            self.input_tokens,
112            self.cache_read_input_tokens,
113            self.cache_creation_input_tokens,
114            self.output_tokens,
115        ]
116        .into_iter()
117        .map(|t| t.unwrap_or(0))
118        .fold(0u64, u64::saturating_add)
119    }
120
121    /// Serialize debug messages to a JSON [`Value`] for store persistence.
122    ///
123    /// Returns `None` when verbose mode was off (no messages captured).
124    pub fn debug_messages_json(&self) -> Option<Value> {
125        self.debug_messages
126            .as_ref()
127            .and_then(|msgs| serde_json::to_value(msgs).ok())
128    }
129
130    /// Exit code of a shell step.
131    ///
132    /// Returns `None` for non-shell steps or when the field is absent.
133    ///
134    /// # Examples
135    ///
136    /// ```
137    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
138    /// use rust_decimal::Decimal;
139    /// use serde_json::json;
140    ///
141    /// let output = StepOutput {
142    ///     output: json!({"stdout": "ok\n", "stderr": "", "exit_code": 0}),
143    ///     duration_ms: 3,
144    ///     cost_usd: Decimal::ZERO,
145    ///     input_tokens: None,
146    ///     cache_read_input_tokens: None,
147    ///     cache_creation_input_tokens: None,
148    ///     output_tokens: None,
149    ///     model: None,
150    ///     debug_messages: None,
151    ///     artifacts: StepArtifacts::default(),
152    ///     account_id: None,
153    /// };
154    /// assert_eq!(output.exit_code(), Some(0));
155    /// ```
156    pub fn exit_code(&self) -> Option<i64> {
157        self.output.get("exit_code").and_then(Value::as_i64)
158    }
159
160    /// Standard output of a shell step, or an empty string for other kinds.
161    ///
162    /// # Examples
163    ///
164    /// ```
165    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
166    /// use rust_decimal::Decimal;
167    /// use serde_json::json;
168    ///
169    /// let output = StepOutput {
170    ///     output: json!({"stdout": "42 tests passed\n", "stderr": "", "exit_code": 0}),
171    ///     duration_ms: 3,
172    ///     cost_usd: Decimal::ZERO,
173    ///     input_tokens: None,
174    ///     cache_read_input_tokens: None,
175    ///     cache_creation_input_tokens: None,
176    ///     output_tokens: None,
177    ///     model: None,
178    ///     debug_messages: None,
179    ///     artifacts: StepArtifacts::default(),
180    ///     account_id: None,
181    /// };
182    /// assert!(output.stdout().contains("42 tests"));
183    /// ```
184    pub fn stdout(&self) -> &str {
185        self.output
186            .get("stdout")
187            .and_then(Value::as_str)
188            .unwrap_or_default()
189    }
190
191    /// Standard error of a shell step, or an empty string for other kinds.
192    ///
193    /// # Examples
194    ///
195    /// ```
196    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
197    /// use rust_decimal::Decimal;
198    /// use serde_json::json;
199    ///
200    /// let output = StepOutput {
201    ///     output: json!({"stdout": "", "stderr": "warning: unused", "exit_code": 0}),
202    ///     duration_ms: 3,
203    ///     cost_usd: Decimal::ZERO,
204    ///     input_tokens: None,
205    ///     cache_read_input_tokens: None,
206    ///     cache_creation_input_tokens: None,
207    ///     output_tokens: None,
208    ///     model: None,
209    ///     debug_messages: None,
210    ///     artifacts: StepArtifacts::default(),
211    ///     account_id: None,
212    /// };
213    /// assert_eq!(output.stderr(), "warning: unused");
214    /// ```
215    pub fn stderr(&self) -> &str {
216        self.output
217            .get("stderr")
218            .and_then(Value::as_str)
219            .unwrap_or_default()
220    }
221
222    /// HTTP status code of an HTTP step.
223    ///
224    /// Returns `None` for non-HTTP steps or when the field is absent.
225    ///
226    /// # Examples
227    ///
228    /// ```
229    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
230    /// use rust_decimal::Decimal;
231    /// use serde_json::json;
232    ///
233    /// let output = StepOutput {
234    ///     output: json!({"status": 204, "body": ""}),
235    ///     duration_ms: 3,
236    ///     cost_usd: Decimal::ZERO,
237    ///     input_tokens: None,
238    ///     cache_read_input_tokens: None,
239    ///     cache_creation_input_tokens: None,
240    ///     output_tokens: None,
241    ///     model: None,
242    ///     debug_messages: None,
243    ///     artifacts: StepArtifacts::default(),
244    ///     account_id: None,
245    /// };
246    /// assert_eq!(output.status(), Some(204));
247    /// ```
248    pub fn status(&self) -> Option<u16> {
249        self.output
250            .get("status")
251            .and_then(Value::as_u64)
252            .and_then(|s| u16::try_from(s).ok())
253    }
254
255    /// Response body of an HTTP step, or an empty string for other kinds.
256    ///
257    /// # Examples
258    ///
259    /// ```
260    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
261    /// use rust_decimal::Decimal;
262    /// use serde_json::json;
263    ///
264    /// let output = StepOutput {
265    ///     output: json!({"status": 200, "body": "{\"ok\":true}"}),
266    ///     duration_ms: 3,
267    ///     cost_usd: Decimal::ZERO,
268    ///     input_tokens: None,
269    ///     cache_read_input_tokens: None,
270    ///     cache_creation_input_tokens: None,
271    ///     output_tokens: None,
272    ///     model: None,
273    ///     debug_messages: None,
274    ///     artifacts: StepArtifacts::default(),
275    ///     account_id: None,
276    /// };
277    /// assert_eq!(output.body(), "{\"ok\":true}");
278    /// ```
279    pub fn body(&self) -> &str {
280        self.output
281            .get("body")
282            .and_then(Value::as_str)
283            .unwrap_or_default()
284    }
285
286    /// Text answer of an agent step without structured output, or an empty
287    /// string for other kinds.
288    ///
289    /// # Examples
290    ///
291    /// ```
292    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
293    /// use rust_decimal::Decimal;
294    /// use serde_json::json;
295    ///
296    /// let answer = StepOutput {
297    ///     output: json!("Looks good."),
298    ///     duration_ms: 3,
299    ///     cost_usd: Decimal::ZERO,
300    ///     input_tokens: None,
301    ///     cache_read_input_tokens: None,
302    ///     cache_creation_input_tokens: None,
303    ///     output_tokens: None,
304    ///     model: None,
305    ///     debug_messages: None,
306    ///     artifacts: StepArtifacts::default(),
307    ///     account_id: None,
308    /// };
309    /// assert_eq!(answer.text(), "Looks good.");
310    /// assert_eq!(StepOutput { output: json!({"stdout": "x"}), ..answer }.text(), "");
311    /// ```
312    pub fn text(&self) -> &str {
313        self.output.as_str().unwrap_or_default()
314    }
315
316    /// Error message of a step tolerated with `allow_failure()` that failed.
317    ///
318    /// An `allow_failure` step that failed (timeout, spawn failure, non-zero
319    /// exit of a shell step without `exit_code_as_output()`, failed input
320    /// preparation) has the output `{"error": "<message>"}`. This accessor
321    /// returns that message. It returns `None` for a successful step, and for
322    /// a shell step configured with `exit_code_as_output()`: the exit code is
323    /// data there, read it with [`exit_code`](Self::exit_code).
324    ///
325    /// It complements [`stdout`](Self::stdout), [`stderr`](Self::stderr),
326    /// [`exit_code`](Self::exit_code), [`status`](Self::status),
327    /// [`body`](Self::body) and [`text`](Self::text). It reads the `"error"`
328    /// key of the output, so it is meant for `allow_failure` steps: an agent or
329    /// operation whose own JSON has an `"error"` key is read the same way.
330    ///
331    /// # Examples
332    ///
333    /// ```
334    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
335    /// use rust_decimal::Decimal;
336    /// use serde_json::json;
337    ///
338    /// let failed = StepOutput {
339    ///     output: json!({"error": "timed out"}),
340    ///     duration_ms: 3,
341    ///     cost_usd: Decimal::ZERO,
342    ///     input_tokens: None,
343    ///     cache_read_input_tokens: None,
344    ///     cache_creation_input_tokens: None,
345    ///     output_tokens: None,
346    ///     model: None,
347    ///     debug_messages: None,
348    ///     artifacts: StepArtifacts::default(),
349    ///     account_id: None,
350    /// };
351    /// assert_eq!(failed.error(), Some("timed out"));
352    /// assert!(!failed.is_success());
353    ///
354    /// let ok = StepOutput { output: json!({"stdout": "", "stderr": "", "exit_code": 0}), ..failed };
355    /// assert_eq!(ok.error(), None);
356    /// ```
357    ///
358    /// ```no_run
359    /// use ironflow_engine::config::ShellConfig;
360    /// use ironflow_engine::context::WorkflowContext;
361    /// use ironflow_engine::error::EngineError;
362    ///
363    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
364    /// let out = ctx
365    ///     .shell("fetch", ShellConfig::new("sleep 30").timeout_secs(1).allow_failure())
366    ///     .await?;
367    /// assert!(out.error().is_some());
368    /// assert!(!out.is_success());
369    /// # Ok(())
370    /// # }
371    /// ```
372    pub fn error(&self) -> Option<&str> {
373        self.output.get(ERROR_KEY).and_then(Value::as_str)
374    }
375
376    /// Whether the step succeeded from the point of view of its own kind.
377    ///
378    /// - Shell step: the exit code is `0`.
379    /// - HTTP step: the status is in the `2xx` range.
380    /// - Any other kind: `false`, since no success marker is recorded.
381    ///
382    /// Mostly useful after a step configured with `allow_failure()`, since a
383    /// failing step otherwise returns an error from the context method.
384    ///
385    /// # Examples
386    ///
387    /// ```
388    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
389    /// use rust_decimal::Decimal;
390    /// use serde_json::json;
391    ///
392    /// let shell = StepOutput {
393    ///     output: json!({"stdout": "", "stderr": "", "exit_code": 1}),
394    ///     duration_ms: 3,
395    ///     cost_usd: Decimal::ZERO,
396    ///     input_tokens: None,
397    ///     cache_read_input_tokens: None,
398    ///     cache_creation_input_tokens: None,
399    ///     output_tokens: None,
400    ///     model: None,
401    ///     debug_messages: None,
402    ///     artifacts: StepArtifacts::default(),
403    ///     account_id: None,
404    /// };
405    /// assert!(!shell.is_success());
406    ///
407    /// let http = StepOutput { output: json!({"status": 201, "body": ""}), ..shell.clone() };
408    /// assert!(http.is_success());
409    /// ```
410    pub fn is_success(&self) -> bool {
411        if let Some(code) = self.exit_code() {
412            return code == 0;
413        }
414        if let Some(status) = self.status() {
415            return (200..300).contains(&status);
416        }
417        false
418    }
419
420    /// Deserialize the step output into `T`.
421    ///
422    /// Intended for agent steps constrained by a JSON schema, and for custom
423    /// operations that return structured JSON.
424    ///
425    /// # Errors
426    ///
427    /// Returns [`EngineError::Serialization`] when the output does not match `T`.
428    ///
429    /// # Examples
430    ///
431    /// ```
432    /// use ironflow_engine::executor::{StepArtifacts, StepOutput};
433    /// use rust_decimal::Decimal;
434    /// use serde::Deserialize;
435    /// use serde_json::json;
436    ///
437    /// #[derive(Deserialize)]
438    /// struct Review {
439    ///     score: u8,
440    /// }
441    ///
442    /// let output = StepOutput {
443    ///     output: json!({"score": 8}),
444    ///     duration_ms: 3,
445    ///     cost_usd: Decimal::ZERO,
446    ///     input_tokens: None,
447    ///     cache_read_input_tokens: None,
448    ///     cache_creation_input_tokens: None,
449    ///     output_tokens: None,
450    ///     model: None,
451    ///     debug_messages: None,
452    ///     artifacts: StepArtifacts::default(),
453    ///     account_id: None,
454    /// };
455    /// let review: Review = output.json()?;
456    /// assert_eq!(review.score, 8);
457    /// # Ok::<(), ironflow_engine::error::EngineError>(())
458    /// ```
459    pub fn json<T: DeserializeOwned>(&self) -> Result<T, EngineError> {
460        from_value(self.output.clone()).map_err(EngineError::Serialization)
461    }
462}
463
464/// Result of a single step within a [`parallel`](crate::context::WorkflowContext::parallel) batch.
465#[derive(Debug, Clone)]
466pub struct ParallelStepResult {
467    /// The step name (same as provided to `parallel()`).
468    pub name: String,
469    /// The step execution output.
470    pub output: StepOutput,
471    /// The step ID in the store (for dependency tracking).
472    pub step_id: Uuid,
473}
474
475/// Enriched result of a completed step, for post-execution inspection.
476///
477/// Collects the step's trace ID, status, metrics, and a truncated output
478/// summary into a single struct that the [`WorkflowContext`](crate::context::WorkflowContext)
479/// accumulates over the run.
480///
481/// # Examples
482///
483/// ```
484/// use ironflow_engine::executor::StepResult;
485/// use ironflow_store::entities::StepStatus;
486/// use rust_decimal::Decimal;
487/// use uuid::Uuid;
488///
489/// let result = StepResult {
490///     trace_id: Uuid::nil(),
491///     name: "build".to_string(),
492///     status: StepStatus::Completed,
493///     duration_ms: 1200,
494///     cost_usd: Decimal::ZERO,
495///     input_tokens: None,
496///     output_tokens: None,
497///     error: None,
498///     output_summary: Some("ok".to_string()),
499/// };
500/// assert_eq!(result.status, StepStatus::Completed);
501/// ```
502#[derive(Debug, Clone, serde::Serialize)]
503pub struct StepResult {
504    /// Deterministic trace ID for log correlation.
505    pub trace_id: Uuid,
506    /// Step name.
507    pub name: String,
508    /// Terminal status.
509    pub status: StepStatus,
510    /// Wall-clock duration in milliseconds.
511    pub duration_ms: u64,
512    /// Cost in USD.
513    pub cost_usd: Decimal,
514    /// Input token count (agent steps only).
515    pub input_tokens: Option<u64>,
516    /// Output token count (agent steps only).
517    pub output_tokens: Option<u64>,
518    /// Error message if the step failed.
519    pub error: Option<String>,
520    /// First 500 characters of the serialized output.
521    pub output_summary: Option<String>,
522}
523
524/// Maximum length of [`StepResult::output_summary`].
525const OUTPUT_SUMMARY_MAX_LEN: usize = 500;
526
527impl StepResult {
528    /// Build from a completed step's output.
529    pub fn from_success(trace_id: Uuid, name: &str, output: &StepOutput) -> Self {
530        Self {
531            trace_id,
532            name: name.to_string(),
533            status: StepStatus::Completed,
534            duration_ms: output.duration_ms,
535            cost_usd: output.cost_usd,
536            input_tokens: output.input_tokens,
537            output_tokens: output.output_tokens,
538            error: None,
539            output_summary: summarize_output(&output.output),
540        }
541    }
542
543    /// Build from a failed step.
544    pub fn from_failure(
545        trace_id: Uuid,
546        name: &str,
547        error: &str,
548        duration_ms: u64,
549        cost_usd: Decimal,
550    ) -> Self {
551        Self {
552            trace_id,
553            name: name.to_string(),
554            status: StepStatus::Failed,
555            duration_ms,
556            cost_usd,
557            input_tokens: None,
558            output_tokens: None,
559            error: Some(error.to_string()),
560            output_summary: None,
561        }
562    }
563}
564
565fn summarize_output(value: &Value) -> Option<String> {
566    let raw = value.to_string();
567    match raw.char_indices().nth(OUTPUT_SUMMARY_MAX_LEN) {
568        None => Some(raw),
569        Some((byte_idx, _)) => Some(raw[..byte_idx].to_string()),
570    }
571}
572
573/// Trait for step executors.
574///
575/// Each step type implements this trait to execute its specific operation
576/// and return a [`StepOutput`].
577pub trait StepExecutor: Send + Sync {
578    /// The [`StepKind`] this executor handles.
579    ///
580    /// The dispatcher uses it to label spans and metrics, so a new step type
581    /// only has to declare its kind here instead of extending a match in
582    /// [`execute_step_config`].
583    fn kind(&self) -> StepKind;
584
585    /// Execute the step and return structured output.
586    ///
587    /// # Errors
588    ///
589    /// Returns [`EngineError`] if the operation fails.
590    fn execute(
591        &self,
592        provider: &Arc<dyn AgentProvider>,
593    ) -> impl Future<Output = Result<StepOutput, EngineError>> + Send;
594}
595
596/// Span and metric label for a step kind.
597pub(crate) fn step_kind_label(kind: &StepKind) -> Cow<'static, str> {
598    match kind {
599        StepKind::Shell => Cow::Borrowed("shell"),
600        StepKind::Http => Cow::Borrowed("http"),
601        StepKind::Agent => Cow::Borrowed("agent"),
602        StepKind::Workflow => Cow::Borrowed("workflow"),
603        StepKind::Approval => Cow::Borrowed("approval"),
604        StepKind::Decision => Cow::Borrowed("decision"),
605        StepKind::HumanInput => Cow::Borrowed("human_input"),
606        StepKind::Signal => Cow::Borrowed("signal"),
607        StepKind::Custom(name) => Cow::Owned(name.clone()),
608    }
609}
610
611/// Execute a [`StepConfig`], letting a [`StepInterceptor`] resolve it first.
612///
613/// When `interceptor` returns `Some(result)` for this config, that result is
614/// used as-is and no executor runs. Otherwise the config is dispatched to the
615/// executor matching its [`StepKind`], exactly like [`execute_step_config`].
616///
617/// When a [`StepLogSender`] is provided, executors that support streaming
618/// will emit log lines in real time (e.g. shell stdout/stderr).
619///
620/// # Errors
621///
622/// Returns [`EngineError::Operation`] if the operation fails, or whichever
623/// error the interceptor returned for an intercepted step.
624///
625/// # Examples
626///
627/// ```no_run
628/// use ironflow_engine::config::{StepConfig, ShellConfig};
629/// use ironflow_engine::executor::execute_step_config_intercepted;
630/// use ironflow_core::provider::AgentProvider;
631/// use ironflow_core::providers::claude::ClaudeCodeProvider;
632/// use std::sync::Arc;
633///
634/// # async fn example() -> Result<(), ironflow_engine::error::EngineError> {
635/// let provider: Arc<dyn AgentProvider> = Arc::new(ClaudeCodeProvider::new());
636/// let config = StepConfig::Shell(ShellConfig::new("echo hello"));
637/// let output = execute_step_config_intercepted(&config, &provider, None, None).await?;
638/// # Ok(())
639/// # }
640/// ```
641#[tracing::instrument(name = "executor.execute_step", skip_all, fields(step.kind))]
642pub async fn execute_step_config_intercepted(
643    config: &StepConfig,
644    provider: &Arc<dyn AgentProvider>,
645    log_sender: Option<StepLogSender>,
646    interceptor: Option<&Arc<dyn StepInterceptor>>,
647) -> Result<StepOutput, EngineError> {
648    let kind = config.kind();
649    let label = step_kind_label(&kind);
650    Span::current().record("step.kind", label.as_ref());
651
652    let intercepted = interceptor.and_then(|i| i.intercept(config));
653    let result = match intercepted {
654        Some(result) => result,
655        None => match config {
656            StepConfig::Shell(cfg) => {
657                let mut executor = ShellExecutor::new(cfg);
658                if let Some(sender) = log_sender {
659                    executor = executor.with_log_sender(sender);
660                }
661                executor.execute(provider).await
662            }
663            StepConfig::Http(cfg) => HttpExecutor::new(cfg).execute(provider).await,
664            StepConfig::Agent(cfg) => {
665                let mut executor = AgentExecutor::new(cfg);
666                if let Some(sender) = log_sender {
667                    executor = executor.with_log_sender(sender);
668                }
669                executor.execute(provider).await
670            }
671            StepConfig::Workflow(_) => Err(EngineError::StepConfig(
672                "workflow steps are executed by WorkflowContext, not the executor".to_string(),
673            )),
674            StepConfig::Approval(_) => Err(EngineError::StepConfig(
675                "approval steps are executed by WorkflowContext, not the executor".to_string(),
676            )),
677            StepConfig::Decision(_) => Err(EngineError::StepConfig(
678                "decision steps are executed by WorkflowContext, not the executor".to_string(),
679            )),
680            StepConfig::Delay(_) => Err(EngineError::StepConfig(
681                "delay steps are executed by WorkflowContext, not the executor".to_string(),
682            )),
683        },
684    };
685
686    #[cfg(feature = "prometheus")]
687    {
688        use ironflow_core::metric_names::{
689            STATUS_ERROR, STATUS_SUCCESS, STEP_DURATION_SECONDS, STEPS_TOTAL,
690        };
691        use metrics::{counter, histogram};
692        let status = if result.is_ok() {
693            STATUS_SUCCESS
694        } else {
695            STATUS_ERROR
696        };
697        let kind_label = label.into_owned();
698        counter!(STEPS_TOTAL, "kind" => kind_label.clone(), "status" => status).increment(1);
699        if let Ok(ref output) = result {
700            histogram!(STEP_DURATION_SECONDS, "kind" => kind_label)
701                .record(output.duration_ms as f64 / 1000.0);
702        }
703    }
704
705    result
706}
707
708/// Execute a [`StepConfig`] and return structured output.
709///
710/// When a [`StepLogSender`] is provided, executors that support streaming
711/// will emit log lines in real time (e.g. shell stdout/stderr).
712///
713/// # Errors
714///
715/// Returns [`EngineError::Operation`] if the operation fails.
716///
717/// # Examples
718///
719/// ```no_run
720/// use ironflow_engine::config::{StepConfig, ShellConfig};
721/// use ironflow_engine::executor::execute_step_config;
722/// use ironflow_core::provider::AgentProvider;
723/// use ironflow_core::providers::claude::ClaudeCodeProvider;
724/// use std::sync::Arc;
725///
726/// # async fn example() -> Result<(), ironflow_engine::error::EngineError> {
727/// let provider: Arc<dyn AgentProvider> = Arc::new(ClaudeCodeProvider::new());
728/// let config = StepConfig::Shell(ShellConfig::new("echo hello"));
729/// let output = execute_step_config(&config, &provider, None).await?;
730/// # Ok(())
731/// # }
732/// ```
733pub async fn execute_step_config(
734    config: &StepConfig,
735    provider: &Arc<dyn AgentProvider>,
736    log_sender: Option<StepLogSender>,
737) -> Result<StepOutput, EngineError> {
738    execute_step_config_intercepted(config, provider, log_sender, None).await
739}
740
741#[cfg(test)]
742mod tests {
743    use super::*;
744    use ironflow_core::provider::DebugMessage;
745    use ironflow_core::providers::claude::ClaudeCodeProvider;
746    use ironflow_core::providers::record_replay::RecordReplayProvider;
747    use serde_json::json;
748
749    use crate::config::{
750        AgentStepConfig, ApprovalConfig, DecisionConfig, DelayConfig, HttpConfig, ShellConfig,
751        WorkflowStepConfig,
752    };
753
754    fn output_of(value: Value) -> StepOutput {
755        StepOutput {
756            output: value,
757            duration_ms: 1,
758            cost_usd: rust_decimal::Decimal::ZERO,
759            input_tokens: None,
760            cache_read_input_tokens: None,
761            cache_creation_input_tokens: None,
762            output_tokens: None,
763            model: None,
764            debug_messages: None,
765            artifacts: StepArtifacts::default(),
766            account_id: None,
767        }
768    }
769
770    #[test]
771    fn step_output_error_is_some_for_failure_output() {
772        let output = output_of(json!({"error": "boom"}));
773
774        assert_eq!(output.error(), Some("boom"));
775        assert!(!output.is_success());
776    }
777
778    #[test]
779    fn step_output_error_is_none_for_success_shell_output() {
780        let output = output_of(json!({"stdout": "hi", "stderr": "", "exit_code": 0}));
781
782        assert_eq!(output.error(), None);
783        assert!(output.is_success());
784    }
785
786    #[test]
787    fn step_output_error_is_none_for_http_and_agent_text() {
788        assert_eq!(output_of(json!({"status": 200, "body": ""})).error(), None);
789        assert_eq!(output_of(json!("fine")).error(), None);
790    }
791
792    #[test]
793    fn step_output_error_is_none_when_error_key_not_a_string() {
794        assert_eq!(output_of(json!({"error": 5})).error(), None);
795        assert_eq!(output_of(Value::Null).error(), None);
796    }
797
798    #[test]
799    fn step_output_with_no_debug_messages_returns_none() {
800        let output = StepOutput {
801            output: json!({"result": "ok"}),
802            duration_ms: 100,
803            cost_usd: rust_decimal::Decimal::ZERO,
804            input_tokens: None,
805            cache_read_input_tokens: None,
806            cache_creation_input_tokens: None,
807            output_tokens: None,
808            model: None,
809            debug_messages: None,
810            artifacts: StepArtifacts::default(),
811            account_id: None,
812        };
813
814        assert_eq!(output.debug_messages_json(), None);
815    }
816
817    #[test]
818    fn step_output_with_empty_debug_messages_returns_some_empty_array() {
819        let output = StepOutput {
820            output: json!({"result": "ok"}),
821            duration_ms: 100,
822            cost_usd: rust_decimal::Decimal::ZERO,
823            input_tokens: None,
824            cache_read_input_tokens: None,
825            cache_creation_input_tokens: None,
826            output_tokens: None,
827            model: None,
828            debug_messages: Some(Vec::new()),
829            artifacts: StepArtifacts::default(),
830            account_id: None,
831        };
832
833        let json_val = output.debug_messages_json();
834        assert!(json_val.is_some());
835        let arr = json_val.unwrap();
836        assert!(arr.is_array());
837        assert_eq!(arr.as_array().unwrap().len(), 0);
838    }
839
840    #[test]
841    fn step_output_debug_messages_json_serializes_messages() {
842        let json_msgs = json!([
843            {
844                "text": "Hello",
845                "thinking": null,
846                "thinking_redacted": false,
847                "tool_calls": [],
848                "tool_results": [],
849                "stop_reason": "end_turn",
850                "input_tokens": 10,
851                "output_tokens": 20
852            },
853            {
854                "text": "Hi there",
855                "thinking": null,
856                "thinking_redacted": false,
857                "tool_calls": [],
858                "tool_results": [],
859                "stop_reason": "end_turn",
860                "input_tokens": 15,
861                "output_tokens": 25
862            }
863        ]);
864
865        let messages: Vec<DebugMessage> =
866            serde_json::from_value(json_msgs.clone()).expect("deserialize debug messages");
867
868        let output = StepOutput {
869            output: json!({"result": "ok"}),
870            duration_ms: 100,
871            cost_usd: rust_decimal::Decimal::ZERO,
872            input_tokens: None,
873            cache_read_input_tokens: None,
874            cache_creation_input_tokens: None,
875            output_tokens: None,
876            model: None,
877            debug_messages: Some(messages),
878            artifacts: StepArtifacts::default(),
879            account_id: None,
880        };
881
882        let json_val = output.debug_messages_json();
883        assert!(json_val.is_some());
884
885        let arr = json_val.unwrap();
886        assert!(arr.is_array());
887        let messages_array = arr.as_array().unwrap();
888        assert_eq!(messages_array.len(), 2);
889        assert_eq!(messages_array[0]["text"], "Hello");
890        assert_eq!(messages_array[1]["text"], "Hi there");
891    }
892
893    #[test]
894    fn step_output_contains_all_metrics() {
895        let output = StepOutput {
896            output: json!({"data": "test"}),
897            duration_ms: 5000,
898            cost_usd: rust_decimal::Decimal::new(123, 2),
899            input_tokens: Some(100),
900            cache_read_input_tokens: None,
901            cache_creation_input_tokens: None,
902            output_tokens: Some(200),
903            model: Some("claude-sonnet".to_string()),
904            debug_messages: None,
905            artifacts: StepArtifacts::default(),
906            account_id: None,
907        };
908
909        assert_eq!(output.duration_ms, 5000);
910        assert_eq!(output.cost_usd, rust_decimal::Decimal::new(123, 2));
911        assert_eq!(output.input_tokens, Some(100));
912        assert_eq!(output.output_tokens, Some(200));
913        assert_eq!(output.model, Some("claude-sonnet".to_string()));
914    }
915
916    #[test]
917    fn step_output_default_tokens_and_model_are_none() {
918        let output = StepOutput {
919            output: json!({}),
920            duration_ms: 0,
921            cost_usd: rust_decimal::Decimal::ZERO,
922            input_tokens: None,
923            cache_read_input_tokens: None,
924            cache_creation_input_tokens: None,
925            output_tokens: None,
926            model: None,
927            debug_messages: None,
928            artifacts: StepArtifacts::default(),
929            account_id: None,
930        };
931
932        assert!(output.input_tokens.is_none());
933        assert!(output.output_tokens.is_none());
934        assert!(output.model.is_none());
935    }
936
937    #[test]
938    fn parallel_step_result_contains_step_metadata() {
939        let step_id = uuid::Uuid::now_v7();
940        let output = StepOutput {
941            output: json!({"done": true}),
942            duration_ms: 1000,
943            cost_usd: rust_decimal::Decimal::ZERO,
944            input_tokens: None,
945            cache_read_input_tokens: None,
946            cache_creation_input_tokens: None,
947            output_tokens: None,
948            model: None,
949            debug_messages: None,
950            artifacts: StepArtifacts::default(),
951            account_id: None,
952        };
953
954        let result = ParallelStepResult {
955            name: "build".to_string(),
956            output,
957            step_id,
958        };
959
960        assert_eq!(result.name, "build");
961        assert_eq!(result.step_id, step_id);
962        assert_eq!(result.output.duration_ms, 1000);
963    }
964
965    #[test]
966    fn step_output_serializes_complex_json_output() {
967        let complex_output = json!({
968            "status": "success",
969            "data": {
970                "items": [1, 2, 3],
971                "nested": {
972                    "key": "value"
973                }
974            }
975        });
976
977        let output = StepOutput {
978            output: complex_output.clone(),
979            duration_ms: 100,
980            cost_usd: rust_decimal::Decimal::ZERO,
981            input_tokens: None,
982            cache_read_input_tokens: None,
983            cache_creation_input_tokens: None,
984            output_tokens: None,
985            model: None,
986            debug_messages: None,
987            artifacts: StepArtifacts::default(),
988            account_id: None,
989        };
990
991        assert_eq!(output.output, complex_output);
992        assert_eq!(output.output["status"], "success");
993        assert_eq!(output.output["data"]["items"][0], 1);
994        assert_eq!(output.output["data"]["nested"]["key"], "value");
995    }
996
997    #[test]
998    fn step_result_from_success_captures_all_fields() {
999        let trace_id = Uuid::nil();
1000        let output = StepOutput {
1001            output: json!({"stdout": "ok"}),
1002            duration_ms: 1500,
1003            cost_usd: Decimal::new(42, 2),
1004            input_tokens: Some(100),
1005            cache_read_input_tokens: None,
1006            cache_creation_input_tokens: None,
1007            output_tokens: Some(200),
1008            model: Some("claude-sonnet".to_string()),
1009            debug_messages: None,
1010            artifacts: StepArtifacts::default(),
1011            account_id: None,
1012        };
1013
1014        let result = StepResult::from_success(trace_id, "build", &output);
1015
1016        assert_eq!(result.trace_id, trace_id);
1017        assert_eq!(result.name, "build");
1018        assert_eq!(result.status, StepStatus::Completed);
1019        assert_eq!(result.duration_ms, 1500);
1020        assert_eq!(result.cost_usd, Decimal::new(42, 2));
1021        assert_eq!(result.input_tokens, Some(100));
1022        assert_eq!(result.output_tokens, Some(200));
1023        assert!(result.error.is_none());
1024        assert!(result.output_summary.is_some());
1025        assert!(result.output_summary.unwrap().contains("stdout"));
1026    }
1027
1028    #[test]
1029    fn step_result_from_failure_captures_error() {
1030        let trace_id = Uuid::nil();
1031        let result =
1032            StepResult::from_failure(trace_id, "deploy", "connection refused", 500, Decimal::ZERO);
1033
1034        assert_eq!(result.trace_id, trace_id);
1035        assert_eq!(result.name, "deploy");
1036        assert_eq!(result.status, StepStatus::Failed);
1037        assert_eq!(result.duration_ms, 500);
1038        assert_eq!(result.error, Some("connection refused".to_string()));
1039        assert!(result.output_summary.is_none());
1040    }
1041
1042    #[test]
1043    fn step_result_output_summary_truncates_long_output() {
1044        let long_value = json!({"data": "x".repeat(1000)});
1045        let output = StepOutput {
1046            output: long_value,
1047            duration_ms: 0,
1048            cost_usd: Decimal::ZERO,
1049            input_tokens: None,
1050            cache_read_input_tokens: None,
1051            cache_creation_input_tokens: None,
1052            output_tokens: None,
1053            model: None,
1054            debug_messages: None,
1055            artifacts: StepArtifacts::default(),
1056            account_id: None,
1057        };
1058
1059        let result = StepResult::from_success(Uuid::nil(), "test", &output);
1060        let summary = result.output_summary.unwrap();
1061        assert_eq!(summary.len(), 500);
1062    }
1063
1064    #[test]
1065    fn step_executor_kind_matches_step_config_kind() {
1066        let shell = ShellConfig::new("echo hi");
1067        let http = HttpConfig::get("https://example.com");
1068        let agent = AgentStepConfig::new("hi");
1069
1070        let shell_kind = ShellExecutor::new(&shell).kind();
1071        let http_kind = HttpExecutor::new(&http).kind();
1072        let agent_kind = AgentExecutor::new(&agent).kind();
1073
1074        assert_eq!(shell_kind, StepConfig::Shell(shell).kind());
1075        assert_eq!(http_kind, StepConfig::Http(http).kind());
1076        assert_eq!(agent_kind, StepConfig::Agent(agent).kind());
1077    }
1078
1079    #[test]
1080    fn step_kind_label_matches_dispatcher_labels() {
1081        let cases: Vec<(StepConfig, &str)> = vec![
1082            (StepConfig::Shell(ShellConfig::new("echo hi")), "shell"),
1083            (
1084                StepConfig::Http(HttpConfig::get("https://example.com")),
1085                "http",
1086            ),
1087            (StepConfig::Agent(AgentStepConfig::new("hi")), "agent"),
1088            (
1089                StepConfig::Workflow(WorkflowStepConfig::new("child", json!({}))),
1090                "workflow",
1091            ),
1092            (
1093                StepConfig::Approval(ApprovalConfig::new("approve?")),
1094                "approval",
1095            ),
1096            (
1097                StepConfig::Decision(DecisionConfig::new(json!({}))),
1098                "decision",
1099            ),
1100            (StepConfig::Delay(DelayConfig::from_secs(1)), "delay"),
1101        ];
1102
1103        for (config, expected) in cases {
1104            assert_eq!(step_kind_label(&config.kind()), expected);
1105        }
1106    }
1107
1108    #[test]
1109    fn step_kind_label_uses_the_custom_kind_name() {
1110        assert_eq!(
1111            step_kind_label(&StepKind::Custom("gitlab".to_string())),
1112            "gitlab"
1113        );
1114    }
1115
1116    /// An interceptor that resolves every shell step with a canned output.
1117    struct CannedShell;
1118
1119    impl StepInterceptor for CannedShell {
1120        fn intercept(&self, config: &StepConfig) -> Option<Result<StepOutput, EngineError>> {
1121            match config {
1122                StepConfig::Shell(_) => Some(Ok(StepOutput {
1123                    output: json!({"stdout": "canned", "stderr": "", "exit_code": 0}),
1124                    duration_ms: 0,
1125                    cost_usd: Decimal::ZERO,
1126                    input_tokens: None,
1127                    cache_read_input_tokens: None,
1128                    cache_creation_input_tokens: None,
1129                    output_tokens: None,
1130                    model: None,
1131                    debug_messages: None,
1132                    artifacts: StepArtifacts::default(),
1133                    account_id: None,
1134                })),
1135                _ => None,
1136            }
1137        }
1138    }
1139
1140    fn test_provider() -> Arc<dyn AgentProvider> {
1141        let inner = ClaudeCodeProvider::new();
1142        Arc::new(RecordReplayProvider::replay(
1143            inner,
1144            "/tmp/ironflow-fixtures",
1145        ))
1146    }
1147
1148    #[tokio::test]
1149    async fn intercepted_step_never_reaches_the_shell_executor() {
1150        let interceptor: Arc<dyn StepInterceptor> = Arc::new(CannedShell);
1151        // A real run of `exit 1` would fail; the canned output proves the
1152        // process was never spawned.
1153        let config = StepConfig::Shell(ShellConfig::new("exit 1"));
1154
1155        let output =
1156            execute_step_config_intercepted(&config, &test_provider(), None, Some(&interceptor))
1157                .await
1158                .expect("the interceptor resolved the step");
1159
1160        assert_eq!(output.stdout(), "canned");
1161        assert_eq!(output.exit_code(), Some(0));
1162    }
1163
1164    #[tokio::test]
1165    async fn a_step_the_interceptor_declines_reaches_the_dispatcher() {
1166        let interceptor: Arc<dyn StepInterceptor> = Arc::new(CannedShell);
1167        // `CannedShell` only answers shell steps, so this one falls through to
1168        // the dispatcher, which refuses workflow configs.
1169        let config = StepConfig::Workflow(WorkflowStepConfig::new("child", json!({})));
1170
1171        let err =
1172            execute_step_config_intercepted(&config, &test_provider(), None, Some(&interceptor))
1173                .await
1174                .expect_err("the dispatcher rejects workflow steps");
1175
1176        assert!(matches!(err, EngineError::StepConfig(_)));
1177    }
1178
1179    #[tokio::test]
1180    async fn without_an_interceptor_the_step_runs_for_real() {
1181        let config = StepConfig::Shell(ShellConfig::new("echo hi"));
1182
1183        let output = execute_step_config_intercepted(&config, &test_provider(), None, None)
1184            .await
1185            .expect("echo succeeds");
1186
1187        assert!(output.stdout().contains("hi"));
1188        assert_eq!(output.exit_code(), Some(0));
1189    }
1190}
1191
1192#[cfg(test)]
1193mod output_helper_tests {
1194    use super::*;
1195    use serde::Deserialize;
1196    use serde_json::json;
1197
1198    fn output(value: Value) -> StepOutput {
1199        StepOutput {
1200            output: value,
1201            duration_ms: 1,
1202            cost_usd: Decimal::ZERO,
1203            input_tokens: None,
1204            cache_read_input_tokens: None,
1205            cache_creation_input_tokens: None,
1206            output_tokens: None,
1207            model: None,
1208            debug_messages: None,
1209            artifacts: StepArtifacts::default(),
1210            account_id: None,
1211        }
1212    }
1213
1214    #[test]
1215    fn agent_total_tokens_includes_cache_tokens() {
1216        let mut out = output(json!("ok"));
1217        out.input_tokens = Some(100);
1218        out.cache_read_input_tokens = Some(5000);
1219        out.cache_creation_input_tokens = Some(200);
1220        out.output_tokens = Some(50);
1221        assert_eq!(out.total_tokens(), 5350);
1222    }
1223
1224    #[test]
1225    fn agent_total_tokens_all_none_is_zero() {
1226        let out = output(json!("ok"));
1227        assert_eq!(out.total_tokens(), 0);
1228    }
1229
1230    #[test]
1231    fn agent_total_tokens_saturates() {
1232        let mut out = output(json!("ok"));
1233        out.input_tokens = Some(u64::MAX);
1234        out.cache_read_input_tokens = Some(10);
1235        assert_eq!(out.total_tokens(), u64::MAX);
1236    }
1237
1238    #[test]
1239    fn shell_helpers_read_shell_fields() {
1240        let out = output(json!({"stdout": "hi\n", "stderr": "warn", "exit_code": 0}));
1241        assert_eq!(out.exit_code(), Some(0));
1242        assert_eq!(out.stdout(), "hi\n");
1243        assert_eq!(out.stderr(), "warn");
1244        assert!(out.is_success());
1245        assert_eq!(out.status(), None);
1246        assert_eq!(out.body(), "");
1247    }
1248
1249    #[test]
1250    fn shell_non_zero_exit_is_not_success() {
1251        let out = output(json!({"stdout": "", "stderr": "", "exit_code": 127}));
1252        assert_eq!(out.exit_code(), Some(127));
1253        assert!(!out.is_success());
1254    }
1255
1256    #[test]
1257    fn http_helpers_read_http_fields() {
1258        let out = output(json!({"status": 200, "body": "{\"ok\":true}"}));
1259        assert_eq!(out.status(), Some(200));
1260        assert_eq!(out.body(), "{\"ok\":true}");
1261        assert!(out.is_success());
1262        assert_eq!(out.exit_code(), None);
1263        assert_eq!(out.stdout(), "");
1264    }
1265
1266    #[test]
1267    fn http_error_status_is_not_success() {
1268        assert!(!output(json!({"status": 500, "body": ""})).is_success());
1269        assert!(!output(json!({"status": 199, "body": ""})).is_success());
1270        assert!(output(json!({"status": 299, "body": ""})).is_success());
1271    }
1272
1273    #[test]
1274    fn status_out_of_u16_range_is_none() {
1275        assert_eq!(output(json!({"status": 70000})).status(), None);
1276        assert_eq!(output(json!({"status": "200"})).status(), None);
1277    }
1278
1279    #[test]
1280    fn agent_output_without_markers_is_not_success() {
1281        let out = output(json!({"summary": "fine"}));
1282        assert!(!out.is_success());
1283        assert_eq!(out.exit_code(), None);
1284        assert_eq!(out.stdout(), "");
1285        assert_eq!(out.body(), "");
1286    }
1287
1288    #[test]
1289    fn json_deserializes_structured_output() {
1290        #[derive(Deserialize, Debug, PartialEq)]
1291        struct Review {
1292            score: u8,
1293            summary: String,
1294        }
1295        let out = output(json!({"score": 9, "summary": "good"}));
1296        let review: Review = out.json().expect("matches schema");
1297        assert_eq!(
1298            review,
1299            Review {
1300                score: 9,
1301                summary: "good".to_string()
1302            }
1303        );
1304    }
1305
1306    #[test]
1307    fn json_reports_mismatch_as_serialization_error() {
1308        #[derive(Deserialize, Debug)]
1309        struct Review {
1310            #[allow(dead_code)]
1311            score: u8,
1312        }
1313        let out = output(json!({"score": "nine"}));
1314        let err = out.json::<Review>().expect_err("type mismatch");
1315        assert!(matches!(err, EngineError::Serialization(_)));
1316    }
1317
1318    #[test]
1319    fn helpers_tolerate_non_object_output() {
1320        let out = output(json!("plain text"));
1321        assert_eq!(out.exit_code(), None);
1322        assert_eq!(out.status(), None);
1323        assert_eq!(out.stdout(), "");
1324        assert!(!out.is_success());
1325    }
1326
1327    #[test]
1328    fn text_reads_a_plain_agent_answer_only() {
1329        assert_eq!(output(json!("plain text")).text(), "plain text");
1330        assert_eq!(output(json!({"stdout": "x"})).text(), "");
1331        assert_eq!(output(Value::Null).text(), "");
1332    }
1333}