Skip to main content

ironflow_engine/executor/
mod.rs

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