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