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