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