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