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