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, 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::Custom(name) => Cow::Owned(name.clone()),
531    }
532}
533
534/// Execute a [`StepConfig`], letting a [`StepInterceptor`] resolve it first.
535///
536/// When `interceptor` returns `Some(result)` for this config, that result is
537/// used as-is and no executor runs. Otherwise the config is dispatched to the
538/// executor matching its [`StepKind`], exactly like [`execute_step_config`].
539///
540/// When a [`StepLogSender`] is provided, executors that support streaming
541/// will emit log lines in real time (e.g. shell stdout/stderr).
542///
543/// # Errors
544///
545/// Returns [`EngineError::Operation`] if the operation fails, or whichever
546/// error the interceptor returned for an intercepted step.
547///
548/// # Examples
549///
550/// ```no_run
551/// use ironflow_engine::config::{StepConfig, ShellConfig};
552/// use ironflow_engine::executor::execute_step_config_intercepted;
553/// use ironflow_core::provider::AgentProvider;
554/// use ironflow_core::providers::claude::ClaudeCodeProvider;
555/// use std::sync::Arc;
556///
557/// # async fn example() -> Result<(), ironflow_engine::error::EngineError> {
558/// let provider: Arc<dyn AgentProvider> = Arc::new(ClaudeCodeProvider::new());
559/// let config = StepConfig::Shell(ShellConfig::new("echo hello"));
560/// let output = execute_step_config_intercepted(&config, &provider, None, None).await?;
561/// # Ok(())
562/// # }
563/// ```
564#[tracing::instrument(name = "executor.execute_step", skip_all, fields(step.kind))]
565pub async fn execute_step_config_intercepted(
566    config: &StepConfig,
567    provider: &Arc<dyn AgentProvider>,
568    log_sender: Option<StepLogSender>,
569    interceptor: Option<&Arc<dyn StepInterceptor>>,
570) -> Result<StepOutput, EngineError> {
571    let kind = config.kind();
572    let label = step_kind_label(&kind);
573    Span::current().record("step.kind", label.as_ref());
574
575    let intercepted = interceptor.and_then(|i| i.intercept(config));
576    let result = match intercepted {
577        Some(result) => result,
578        None => match config {
579            StepConfig::Shell(cfg) => {
580                let mut executor = ShellExecutor::new(cfg);
581                if let Some(sender) = log_sender {
582                    executor = executor.with_log_sender(sender);
583                }
584                executor.execute(provider).await
585            }
586            StepConfig::Http(cfg) => HttpExecutor::new(cfg).execute(provider).await,
587            StepConfig::Agent(cfg) => {
588                let mut executor = AgentExecutor::new(cfg);
589                if let Some(sender) = log_sender {
590                    executor = executor.with_log_sender(sender);
591                }
592                executor.execute(provider).await
593            }
594            StepConfig::Workflow(_) => Err(EngineError::StepConfig(
595                "workflow steps are executed by WorkflowContext, not the executor".to_string(),
596            )),
597            StepConfig::Approval(_) => Err(EngineError::StepConfig(
598                "approval steps are executed by WorkflowContext, not the executor".to_string(),
599            )),
600            StepConfig::Decision(_) => Err(EngineError::StepConfig(
601                "decision steps are executed by WorkflowContext, not the executor".to_string(),
602            )),
603            StepConfig::Delay(_) => Err(EngineError::StepConfig(
604                "delay steps are executed by WorkflowContext, not the executor".to_string(),
605            )),
606        },
607    };
608
609    #[cfg(feature = "prometheus")]
610    {
611        use ironflow_core::metric_names::{
612            STATUS_ERROR, STATUS_SUCCESS, STEP_DURATION_SECONDS, STEPS_TOTAL,
613        };
614        use metrics::{counter, histogram};
615        let status = if result.is_ok() {
616            STATUS_SUCCESS
617        } else {
618            STATUS_ERROR
619        };
620        let kind_label = label.into_owned();
621        counter!(STEPS_TOTAL, "kind" => kind_label.clone(), "status" => status).increment(1);
622        if let Ok(ref output) = result {
623            histogram!(STEP_DURATION_SECONDS, "kind" => kind_label)
624                .record(output.duration_ms as f64 / 1000.0);
625        }
626    }
627
628    result
629}
630
631/// Execute a [`StepConfig`] and return structured output.
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.
639///
640/// # Examples
641///
642/// ```no_run
643/// use ironflow_engine::config::{StepConfig, ShellConfig};
644/// use ironflow_engine::executor::execute_step_config;
645/// use ironflow_core::provider::AgentProvider;
646/// use ironflow_core::providers::claude::ClaudeCodeProvider;
647/// use std::sync::Arc;
648///
649/// # async fn example() -> Result<(), ironflow_engine::error::EngineError> {
650/// let provider: Arc<dyn AgentProvider> = Arc::new(ClaudeCodeProvider::new());
651/// let config = StepConfig::Shell(ShellConfig::new("echo hello"));
652/// let output = execute_step_config(&config, &provider, None).await?;
653/// # Ok(())
654/// # }
655/// ```
656pub async fn execute_step_config(
657    config: &StepConfig,
658    provider: &Arc<dyn AgentProvider>,
659    log_sender: Option<StepLogSender>,
660) -> Result<StepOutput, EngineError> {
661    execute_step_config_intercepted(config, provider, log_sender, None).await
662}
663
664#[cfg(test)]
665mod tests {
666    use super::*;
667    use ironflow_core::provider::DebugMessage;
668    use ironflow_core::providers::claude::ClaudeCodeProvider;
669    use ironflow_core::providers::record_replay::RecordReplayProvider;
670    use serde_json::json;
671
672    use crate::config::{
673        AgentStepConfig, ApprovalConfig, DecisionConfig, DelayConfig, HttpConfig, ShellConfig,
674        WorkflowStepConfig,
675    };
676
677    #[test]
678    fn step_output_with_no_debug_messages_returns_none() {
679        let output = StepOutput {
680            output: json!({"result": "ok"}),
681            duration_ms: 100,
682            cost_usd: rust_decimal::Decimal::ZERO,
683            input_tokens: None,
684            cache_read_input_tokens: None,
685            cache_creation_input_tokens: None,
686            output_tokens: None,
687            model: None,
688            debug_messages: None,
689            artifacts: StepArtifacts::default(),
690        };
691
692        assert_eq!(output.debug_messages_json(), None);
693    }
694
695    #[test]
696    fn step_output_with_empty_debug_messages_returns_some_empty_array() {
697        let output = StepOutput {
698            output: json!({"result": "ok"}),
699            duration_ms: 100,
700            cost_usd: rust_decimal::Decimal::ZERO,
701            input_tokens: None,
702            cache_read_input_tokens: None,
703            cache_creation_input_tokens: None,
704            output_tokens: None,
705            model: None,
706            debug_messages: Some(Vec::new()),
707            artifacts: StepArtifacts::default(),
708        };
709
710        let json_val = output.debug_messages_json();
711        assert!(json_val.is_some());
712        let arr = json_val.unwrap();
713        assert!(arr.is_array());
714        assert_eq!(arr.as_array().unwrap().len(), 0);
715    }
716
717    #[test]
718    fn step_output_debug_messages_json_serializes_messages() {
719        let json_msgs = json!([
720            {
721                "text": "Hello",
722                "thinking": null,
723                "thinking_redacted": false,
724                "tool_calls": [],
725                "tool_results": [],
726                "stop_reason": "end_turn",
727                "input_tokens": 10,
728                "output_tokens": 20
729            },
730            {
731                "text": "Hi there",
732                "thinking": null,
733                "thinking_redacted": false,
734                "tool_calls": [],
735                "tool_results": [],
736                "stop_reason": "end_turn",
737                "input_tokens": 15,
738                "output_tokens": 25
739            }
740        ]);
741
742        let messages: Vec<DebugMessage> =
743            serde_json::from_value(json_msgs.clone()).expect("deserialize debug messages");
744
745        let output = StepOutput {
746            output: json!({"result": "ok"}),
747            duration_ms: 100,
748            cost_usd: rust_decimal::Decimal::ZERO,
749            input_tokens: None,
750            cache_read_input_tokens: None,
751            cache_creation_input_tokens: None,
752            output_tokens: None,
753            model: None,
754            debug_messages: Some(messages),
755            artifacts: StepArtifacts::default(),
756        };
757
758        let json_val = output.debug_messages_json();
759        assert!(json_val.is_some());
760
761        let arr = json_val.unwrap();
762        assert!(arr.is_array());
763        let messages_array = arr.as_array().unwrap();
764        assert_eq!(messages_array.len(), 2);
765        assert_eq!(messages_array[0]["text"], "Hello");
766        assert_eq!(messages_array[1]["text"], "Hi there");
767    }
768
769    #[test]
770    fn step_output_contains_all_metrics() {
771        let output = StepOutput {
772            output: json!({"data": "test"}),
773            duration_ms: 5000,
774            cost_usd: rust_decimal::Decimal::new(123, 2),
775            input_tokens: Some(100),
776            cache_read_input_tokens: None,
777            cache_creation_input_tokens: None,
778            output_tokens: Some(200),
779            model: Some("claude-sonnet".to_string()),
780            debug_messages: None,
781            artifacts: StepArtifacts::default(),
782        };
783
784        assert_eq!(output.duration_ms, 5000);
785        assert_eq!(output.cost_usd, rust_decimal::Decimal::new(123, 2));
786        assert_eq!(output.input_tokens, Some(100));
787        assert_eq!(output.output_tokens, Some(200));
788        assert_eq!(output.model, Some("claude-sonnet".to_string()));
789    }
790
791    #[test]
792    fn step_output_default_tokens_and_model_are_none() {
793        let output = StepOutput {
794            output: json!({}),
795            duration_ms: 0,
796            cost_usd: rust_decimal::Decimal::ZERO,
797            input_tokens: None,
798            cache_read_input_tokens: None,
799            cache_creation_input_tokens: None,
800            output_tokens: None,
801            model: None,
802            debug_messages: None,
803            artifacts: StepArtifacts::default(),
804        };
805
806        assert!(output.input_tokens.is_none());
807        assert!(output.output_tokens.is_none());
808        assert!(output.model.is_none());
809    }
810
811    #[test]
812    fn parallel_step_result_contains_step_metadata() {
813        let step_id = uuid::Uuid::now_v7();
814        let output = StepOutput {
815            output: json!({"done": true}),
816            duration_ms: 1000,
817            cost_usd: rust_decimal::Decimal::ZERO,
818            input_tokens: None,
819            cache_read_input_tokens: None,
820            cache_creation_input_tokens: None,
821            output_tokens: None,
822            model: None,
823            debug_messages: None,
824            artifacts: StepArtifacts::default(),
825        };
826
827        let result = ParallelStepResult {
828            name: "build".to_string(),
829            output,
830            step_id,
831        };
832
833        assert_eq!(result.name, "build");
834        assert_eq!(result.step_id, step_id);
835        assert_eq!(result.output.duration_ms, 1000);
836    }
837
838    #[test]
839    fn step_output_serializes_complex_json_output() {
840        let complex_output = json!({
841            "status": "success",
842            "data": {
843                "items": [1, 2, 3],
844                "nested": {
845                    "key": "value"
846                }
847            }
848        });
849
850        let output = StepOutput {
851            output: complex_output.clone(),
852            duration_ms: 100,
853            cost_usd: rust_decimal::Decimal::ZERO,
854            input_tokens: None,
855            cache_read_input_tokens: None,
856            cache_creation_input_tokens: None,
857            output_tokens: None,
858            model: None,
859            debug_messages: None,
860            artifacts: StepArtifacts::default(),
861        };
862
863        assert_eq!(output.output, complex_output);
864        assert_eq!(output.output["status"], "success");
865        assert_eq!(output.output["data"]["items"][0], 1);
866        assert_eq!(output.output["data"]["nested"]["key"], "value");
867    }
868
869    #[test]
870    fn step_result_from_success_captures_all_fields() {
871        let trace_id = Uuid::nil();
872        let output = StepOutput {
873            output: json!({"stdout": "ok"}),
874            duration_ms: 1500,
875            cost_usd: Decimal::new(42, 2),
876            input_tokens: Some(100),
877            cache_read_input_tokens: None,
878            cache_creation_input_tokens: None,
879            output_tokens: Some(200),
880            model: Some("claude-sonnet".to_string()),
881            debug_messages: None,
882            artifacts: StepArtifacts::default(),
883        };
884
885        let result = StepResult::from_success(trace_id, "build", &output);
886
887        assert_eq!(result.trace_id, trace_id);
888        assert_eq!(result.name, "build");
889        assert_eq!(result.status, StepStatus::Completed);
890        assert_eq!(result.duration_ms, 1500);
891        assert_eq!(result.cost_usd, Decimal::new(42, 2));
892        assert_eq!(result.input_tokens, Some(100));
893        assert_eq!(result.output_tokens, Some(200));
894        assert!(result.error.is_none());
895        assert!(result.output_summary.is_some());
896        assert!(result.output_summary.unwrap().contains("stdout"));
897    }
898
899    #[test]
900    fn step_result_from_failure_captures_error() {
901        let trace_id = Uuid::nil();
902        let result =
903            StepResult::from_failure(trace_id, "deploy", "connection refused", 500, Decimal::ZERO);
904
905        assert_eq!(result.trace_id, trace_id);
906        assert_eq!(result.name, "deploy");
907        assert_eq!(result.status, StepStatus::Failed);
908        assert_eq!(result.duration_ms, 500);
909        assert_eq!(result.error, Some("connection refused".to_string()));
910        assert!(result.output_summary.is_none());
911    }
912
913    #[test]
914    fn step_result_output_summary_truncates_long_output() {
915        let long_value = json!({"data": "x".repeat(1000)});
916        let output = StepOutput {
917            output: long_value,
918            duration_ms: 0,
919            cost_usd: Decimal::ZERO,
920            input_tokens: None,
921            cache_read_input_tokens: None,
922            cache_creation_input_tokens: None,
923            output_tokens: None,
924            model: None,
925            debug_messages: None,
926            artifacts: StepArtifacts::default(),
927        };
928
929        let result = StepResult::from_success(Uuid::nil(), "test", &output);
930        let summary = result.output_summary.unwrap();
931        assert_eq!(summary.len(), 500);
932    }
933
934    #[test]
935    fn step_executor_kind_matches_step_config_kind() {
936        let shell = ShellConfig::new("echo hi");
937        let http = HttpConfig::get("https://example.com");
938        let agent = AgentStepConfig::new("hi");
939
940        let shell_kind = ShellExecutor::new(&shell).kind();
941        let http_kind = HttpExecutor::new(&http).kind();
942        let agent_kind = AgentExecutor::new(&agent).kind();
943
944        assert_eq!(shell_kind, StepConfig::Shell(shell).kind());
945        assert_eq!(http_kind, StepConfig::Http(http).kind());
946        assert_eq!(agent_kind, StepConfig::Agent(agent).kind());
947    }
948
949    #[test]
950    fn step_kind_label_matches_dispatcher_labels() {
951        let cases: Vec<(StepConfig, &str)> = vec![
952            (StepConfig::Shell(ShellConfig::new("echo hi")), "shell"),
953            (
954                StepConfig::Http(HttpConfig::get("https://example.com")),
955                "http",
956            ),
957            (StepConfig::Agent(AgentStepConfig::new("hi")), "agent"),
958            (
959                StepConfig::Workflow(WorkflowStepConfig::new("child", json!({}))),
960                "workflow",
961            ),
962            (
963                StepConfig::Approval(ApprovalConfig::new("approve?")),
964                "approval",
965            ),
966            (
967                StepConfig::Decision(DecisionConfig::new(json!({}))),
968                "decision",
969            ),
970            (StepConfig::Delay(DelayConfig::from_secs(1)), "delay"),
971        ];
972
973        for (config, expected) in cases {
974            assert_eq!(step_kind_label(&config.kind()), expected);
975        }
976    }
977
978    #[test]
979    fn step_kind_label_uses_the_custom_kind_name() {
980        assert_eq!(
981            step_kind_label(&StepKind::Custom("gitlab".to_string())),
982            "gitlab"
983        );
984    }
985
986    /// An interceptor that resolves every shell step with a canned output.
987    struct CannedShell;
988
989    impl StepInterceptor for CannedShell {
990        fn intercept(&self, config: &StepConfig) -> Option<Result<StepOutput, EngineError>> {
991            match config {
992                StepConfig::Shell(_) => Some(Ok(StepOutput {
993                    output: json!({"stdout": "canned", "stderr": "", "exit_code": 0}),
994                    duration_ms: 0,
995                    cost_usd: Decimal::ZERO,
996                    input_tokens: None,
997                    cache_read_input_tokens: None,
998                    cache_creation_input_tokens: None,
999                    output_tokens: None,
1000                    model: None,
1001                    debug_messages: None,
1002                    artifacts: StepArtifacts::default(),
1003                })),
1004                _ => None,
1005            }
1006        }
1007    }
1008
1009    fn test_provider() -> Arc<dyn AgentProvider> {
1010        let inner = ClaudeCodeProvider::new();
1011        Arc::new(RecordReplayProvider::replay(
1012            inner,
1013            "/tmp/ironflow-fixtures",
1014        ))
1015    }
1016
1017    #[tokio::test]
1018    async fn intercepted_step_never_reaches_the_shell_executor() {
1019        let interceptor: Arc<dyn StepInterceptor> = Arc::new(CannedShell);
1020        // A real run of `exit 1` would fail; the canned output proves the
1021        // process was never spawned.
1022        let config = StepConfig::Shell(ShellConfig::new("exit 1"));
1023
1024        let output =
1025            execute_step_config_intercepted(&config, &test_provider(), None, Some(&interceptor))
1026                .await
1027                .expect("the interceptor resolved the step");
1028
1029        assert_eq!(output.stdout(), "canned");
1030        assert_eq!(output.exit_code(), Some(0));
1031    }
1032
1033    #[tokio::test]
1034    async fn a_step_the_interceptor_declines_reaches_the_dispatcher() {
1035        let interceptor: Arc<dyn StepInterceptor> = Arc::new(CannedShell);
1036        // `CannedShell` only answers shell steps, so this one falls through to
1037        // the dispatcher, which refuses workflow configs.
1038        let config = StepConfig::Workflow(WorkflowStepConfig::new("child", json!({})));
1039
1040        let err =
1041            execute_step_config_intercepted(&config, &test_provider(), None, Some(&interceptor))
1042                .await
1043                .expect_err("the dispatcher rejects workflow steps");
1044
1045        assert!(matches!(err, EngineError::StepConfig(_)));
1046    }
1047
1048    #[tokio::test]
1049    async fn without_an_interceptor_the_step_runs_for_real() {
1050        let config = StepConfig::Shell(ShellConfig::new("echo hi"));
1051
1052        let output = execute_step_config_intercepted(&config, &test_provider(), None, None)
1053            .await
1054            .expect("echo succeeds");
1055
1056        assert!(output.stdout().contains("hi"));
1057        assert_eq!(output.exit_code(), Some(0));
1058    }
1059}
1060
1061#[cfg(test)]
1062mod output_helper_tests {
1063    use super::*;
1064    use serde::Deserialize;
1065    use serde_json::json;
1066
1067    fn output(value: Value) -> StepOutput {
1068        StepOutput {
1069            output: value,
1070            duration_ms: 1,
1071            cost_usd: Decimal::ZERO,
1072            input_tokens: None,
1073            cache_read_input_tokens: None,
1074            cache_creation_input_tokens: None,
1075            output_tokens: None,
1076            model: None,
1077            debug_messages: None,
1078            artifacts: StepArtifacts::default(),
1079        }
1080    }
1081
1082    #[test]
1083    fn agent_total_tokens_includes_cache_tokens() {
1084        let mut out = output(json!("ok"));
1085        out.input_tokens = Some(100);
1086        out.cache_read_input_tokens = Some(5000);
1087        out.cache_creation_input_tokens = Some(200);
1088        out.output_tokens = Some(50);
1089        assert_eq!(out.total_tokens(), 5350);
1090    }
1091
1092    #[test]
1093    fn agent_total_tokens_all_none_is_zero() {
1094        let out = output(json!("ok"));
1095        assert_eq!(out.total_tokens(), 0);
1096    }
1097
1098    #[test]
1099    fn agent_total_tokens_saturates() {
1100        let mut out = output(json!("ok"));
1101        out.input_tokens = Some(u64::MAX);
1102        out.cache_read_input_tokens = Some(10);
1103        assert_eq!(out.total_tokens(), u64::MAX);
1104    }
1105
1106    #[test]
1107    fn shell_helpers_read_shell_fields() {
1108        let out = output(json!({"stdout": "hi\n", "stderr": "warn", "exit_code": 0}));
1109        assert_eq!(out.exit_code(), Some(0));
1110        assert_eq!(out.stdout(), "hi\n");
1111        assert_eq!(out.stderr(), "warn");
1112        assert!(out.is_success());
1113        assert_eq!(out.status(), None);
1114        assert_eq!(out.body(), "");
1115    }
1116
1117    #[test]
1118    fn shell_non_zero_exit_is_not_success() {
1119        let out = output(json!({"stdout": "", "stderr": "", "exit_code": 127}));
1120        assert_eq!(out.exit_code(), Some(127));
1121        assert!(!out.is_success());
1122    }
1123
1124    #[test]
1125    fn http_helpers_read_http_fields() {
1126        let out = output(json!({"status": 200, "body": "{\"ok\":true}"}));
1127        assert_eq!(out.status(), Some(200));
1128        assert_eq!(out.body(), "{\"ok\":true}");
1129        assert!(out.is_success());
1130        assert_eq!(out.exit_code(), None);
1131        assert_eq!(out.stdout(), "");
1132    }
1133
1134    #[test]
1135    fn http_error_status_is_not_success() {
1136        assert!(!output(json!({"status": 500, "body": ""})).is_success());
1137        assert!(!output(json!({"status": 199, "body": ""})).is_success());
1138        assert!(output(json!({"status": 299, "body": ""})).is_success());
1139    }
1140
1141    #[test]
1142    fn status_out_of_u16_range_is_none() {
1143        assert_eq!(output(json!({"status": 70000})).status(), None);
1144        assert_eq!(output(json!({"status": "200"})).status(), None);
1145    }
1146
1147    #[test]
1148    fn agent_output_without_markers_is_not_success() {
1149        let out = output(json!({"summary": "fine"}));
1150        assert!(!out.is_success());
1151        assert_eq!(out.exit_code(), None);
1152        assert_eq!(out.stdout(), "");
1153        assert_eq!(out.body(), "");
1154    }
1155
1156    #[test]
1157    fn json_deserializes_structured_output() {
1158        #[derive(Deserialize, Debug, PartialEq)]
1159        struct Review {
1160            score: u8,
1161            summary: String,
1162        }
1163        let out = output(json!({"score": 9, "summary": "good"}));
1164        let review: Review = out.json().expect("matches schema");
1165        assert_eq!(
1166            review,
1167            Review {
1168                score: 9,
1169                summary: "good".to_string()
1170            }
1171        );
1172    }
1173
1174    #[test]
1175    fn json_reports_mismatch_as_serialization_error() {
1176        #[derive(Deserialize, Debug)]
1177        struct Review {
1178            #[allow(dead_code)]
1179            score: u8,
1180        }
1181        let out = output(json!({"score": "nine"}));
1182        let err = out.json::<Review>().expect_err("type mismatch");
1183        assert!(matches!(err, EngineError::Serialization(_)));
1184    }
1185
1186    #[test]
1187    fn helpers_tolerate_non_object_output() {
1188        let out = output(json!("plain text"));
1189        assert_eq!(out.exit_code(), None);
1190        assert_eq!(out.status(), None);
1191        assert_eq!(out.stdout(), "");
1192        assert!(!out.is_success());
1193    }
1194
1195    #[test]
1196    fn text_reads_a_plain_agent_answer_only() {
1197        assert_eq!(output(json!("plain text")).text(), "plain text");
1198        assert_eq!(output(json!({"stdout": "x"})).text(), "");
1199        assert_eq!(output(Value::Null).text(), "");
1200    }
1201}