mod agent;
mod http;
mod shell;
use std::future::Future;
use std::sync::Arc;
use rust_decimal::Decimal;
use serde::de::DeserializeOwned;
use serde_json::{Value, from_value};
use uuid::Uuid;
use ironflow_core::provider::{AgentProvider, DebugMessage};
use ironflow_store::entities::StepStatus;
use crate::config::StepConfig;
use crate::error::EngineError;
use crate::log_sender::StepLogSender;
pub use agent::AgentExecutor;
pub use http::HttpExecutor;
pub use shell::ShellExecutor;
#[derive(Debug, Clone)]
pub struct StepOutput {
pub output: Value,
pub duration_ms: u64,
pub cost_usd: Decimal,
pub input_tokens: Option<u64>,
pub output_tokens: Option<u64>,
pub model: Option<String>,
pub debug_messages: Option<Vec<DebugMessage>>,
}
impl StepOutput {
pub fn debug_messages_json(&self) -> Option<Value> {
self.debug_messages
.as_ref()
.and_then(|msgs| serde_json::to_value(msgs).ok())
}
pub fn exit_code(&self) -> Option<i64> {
self.output.get("exit_code").and_then(Value::as_i64)
}
pub fn stdout(&self) -> &str {
self.output
.get("stdout")
.and_then(Value::as_str)
.unwrap_or_default()
}
pub fn stderr(&self) -> &str {
self.output
.get("stderr")
.and_then(Value::as_str)
.unwrap_or_default()
}
pub fn status(&self) -> Option<u16> {
self.output
.get("status")
.and_then(Value::as_u64)
.and_then(|s| u16::try_from(s).ok())
}
pub fn body(&self) -> &str {
self.output
.get("body")
.and_then(Value::as_str)
.unwrap_or_default()
}
pub fn is_success(&self) -> bool {
if let Some(code) = self.exit_code() {
return code == 0;
}
if let Some(status) = self.status() {
return (200..300).contains(&status);
}
false
}
pub fn json<T: DeserializeOwned>(&self) -> Result<T, EngineError> {
from_value(self.output.clone()).map_err(EngineError::Serialization)
}
}
#[derive(Debug, Clone)]
pub struct ParallelStepResult {
pub name: String,
pub output: StepOutput,
pub step_id: Uuid,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct StepResult {
pub trace_id: Uuid,
pub name: String,
pub status: StepStatus,
pub duration_ms: u64,
pub cost_usd: Decimal,
pub input_tokens: Option<u64>,
pub output_tokens: Option<u64>,
pub error: Option<String>,
pub output_summary: Option<String>,
}
const OUTPUT_SUMMARY_MAX_LEN: usize = 500;
impl StepResult {
pub fn from_success(trace_id: Uuid, name: &str, output: &StepOutput) -> Self {
Self {
trace_id,
name: name.to_string(),
status: StepStatus::Completed,
duration_ms: output.duration_ms,
cost_usd: output.cost_usd,
input_tokens: output.input_tokens,
output_tokens: output.output_tokens,
error: None,
output_summary: summarize_output(&output.output),
}
}
pub fn from_failure(
trace_id: Uuid,
name: &str,
error: &str,
duration_ms: u64,
cost_usd: Decimal,
) -> Self {
Self {
trace_id,
name: name.to_string(),
status: StepStatus::Failed,
duration_ms,
cost_usd,
input_tokens: None,
output_tokens: None,
error: Some(error.to_string()),
output_summary: None,
}
}
}
fn summarize_output(value: &Value) -> Option<String> {
let raw = value.to_string();
match raw.char_indices().nth(OUTPUT_SUMMARY_MAX_LEN) {
None => Some(raw),
Some((byte_idx, _)) => Some(raw[..byte_idx].to_string()),
}
}
pub trait StepExecutor: Send + Sync {
fn execute(
&self,
provider: &Arc<dyn AgentProvider>,
) -> impl Future<Output = Result<StepOutput, EngineError>> + Send;
}
#[tracing::instrument(name = "executor.execute_step", skip_all, fields(step.kind))]
pub async fn execute_step_config(
config: &StepConfig,
provider: &Arc<dyn AgentProvider>,
log_sender: Option<StepLogSender>,
) -> Result<StepOutput, EngineError> {
let kind = match config {
StepConfig::Shell(_) => "shell",
StepConfig::Http(_) => "http",
StepConfig::Agent(_) => "agent",
StepConfig::Workflow(_) => "workflow",
StepConfig::Approval(_) => "approval",
};
tracing::Span::current().record("step.kind", kind);
let result = match config {
StepConfig::Shell(cfg) => {
let mut executor = ShellExecutor::new(cfg);
if let Some(sender) = log_sender {
executor = executor.with_log_sender(sender);
}
executor.execute(provider).await
}
StepConfig::Http(cfg) => HttpExecutor::new(cfg).execute(provider).await,
StepConfig::Agent(cfg) => {
let mut executor = AgentExecutor::new(cfg);
if let Some(sender) = log_sender {
executor = executor.with_log_sender(sender);
}
executor.execute(provider).await
}
StepConfig::Workflow(_) => Err(EngineError::StepConfig(
"workflow steps are executed by WorkflowContext, not the executor".to_string(),
)),
StepConfig::Approval(_) => Err(EngineError::StepConfig(
"approval steps are executed by WorkflowContext, not the executor".to_string(),
)),
};
#[cfg(feature = "prometheus")]
{
use ironflow_core::metric_names::{
STATUS_ERROR, STATUS_SUCCESS, STEP_DURATION_SECONDS, STEPS_TOTAL,
};
use metrics::{counter, histogram};
let status = if result.is_ok() {
STATUS_SUCCESS
} else {
STATUS_ERROR
};
counter!(STEPS_TOTAL, "kind" => kind, "status" => status).increment(1);
if let Ok(ref output) = result {
histogram!(STEP_DURATION_SECONDS, "kind" => kind)
.record(output.duration_ms as f64 / 1000.0);
}
}
result
}
#[cfg(test)]
mod tests {
use super::*;
use ironflow_core::provider::DebugMessage;
use serde_json::json;
#[test]
fn step_output_with_no_debug_messages_returns_none() {
let output = StepOutput {
output: json!({"result": "ok"}),
duration_ms: 100,
cost_usd: rust_decimal::Decimal::ZERO,
input_tokens: None,
output_tokens: None,
model: None,
debug_messages: None,
};
assert_eq!(output.debug_messages_json(), None);
}
#[test]
fn step_output_with_empty_debug_messages_returns_some_empty_array() {
let output = StepOutput {
output: json!({"result": "ok"}),
duration_ms: 100,
cost_usd: rust_decimal::Decimal::ZERO,
input_tokens: None,
output_tokens: None,
model: None,
debug_messages: Some(Vec::new()),
};
let json_val = output.debug_messages_json();
assert!(json_val.is_some());
let arr = json_val.unwrap();
assert!(arr.is_array());
assert_eq!(arr.as_array().unwrap().len(), 0);
}
#[test]
fn step_output_debug_messages_json_serializes_messages() {
let json_msgs = json!([
{
"text": "Hello",
"thinking": null,
"thinking_redacted": false,
"tool_calls": [],
"tool_results": [],
"stop_reason": "end_turn",
"input_tokens": 10,
"output_tokens": 20
},
{
"text": "Hi there",
"thinking": null,
"thinking_redacted": false,
"tool_calls": [],
"tool_results": [],
"stop_reason": "end_turn",
"input_tokens": 15,
"output_tokens": 25
}
]);
let messages: Vec<DebugMessage> =
serde_json::from_value(json_msgs.clone()).expect("deserialize debug messages");
let output = StepOutput {
output: json!({"result": "ok"}),
duration_ms: 100,
cost_usd: rust_decimal::Decimal::ZERO,
input_tokens: None,
output_tokens: None,
model: None,
debug_messages: Some(messages),
};
let json_val = output.debug_messages_json();
assert!(json_val.is_some());
let arr = json_val.unwrap();
assert!(arr.is_array());
let messages_array = arr.as_array().unwrap();
assert_eq!(messages_array.len(), 2);
assert_eq!(messages_array[0]["text"], "Hello");
assert_eq!(messages_array[1]["text"], "Hi there");
}
#[test]
fn step_output_contains_all_metrics() {
let output = StepOutput {
output: json!({"data": "test"}),
duration_ms: 5000,
cost_usd: rust_decimal::Decimal::new(123, 2),
input_tokens: Some(100),
output_tokens: Some(200),
model: Some("claude-sonnet".to_string()),
debug_messages: None,
};
assert_eq!(output.duration_ms, 5000);
assert_eq!(output.cost_usd, rust_decimal::Decimal::new(123, 2));
assert_eq!(output.input_tokens, Some(100));
assert_eq!(output.output_tokens, Some(200));
assert_eq!(output.model, Some("claude-sonnet".to_string()));
}
#[test]
fn step_output_default_tokens_and_model_are_none() {
let output = StepOutput {
output: json!({}),
duration_ms: 0,
cost_usd: rust_decimal::Decimal::ZERO,
input_tokens: None,
output_tokens: None,
model: None,
debug_messages: None,
};
assert!(output.input_tokens.is_none());
assert!(output.output_tokens.is_none());
assert!(output.model.is_none());
}
#[test]
fn parallel_step_result_contains_step_metadata() {
let step_id = uuid::Uuid::now_v7();
let output = StepOutput {
output: json!({"done": true}),
duration_ms: 1000,
cost_usd: rust_decimal::Decimal::ZERO,
input_tokens: None,
output_tokens: None,
model: None,
debug_messages: None,
};
let result = ParallelStepResult {
name: "build".to_string(),
output,
step_id,
};
assert_eq!(result.name, "build");
assert_eq!(result.step_id, step_id);
assert_eq!(result.output.duration_ms, 1000);
}
#[test]
fn step_output_serializes_complex_json_output() {
let complex_output = json!({
"status": "success",
"data": {
"items": [1, 2, 3],
"nested": {
"key": "value"
}
}
});
let output = StepOutput {
output: complex_output.clone(),
duration_ms: 100,
cost_usd: rust_decimal::Decimal::ZERO,
input_tokens: None,
output_tokens: None,
model: None,
debug_messages: None,
};
assert_eq!(output.output, complex_output);
assert_eq!(output.output["status"], "success");
assert_eq!(output.output["data"]["items"][0], 1);
assert_eq!(output.output["data"]["nested"]["key"], "value");
}
#[test]
fn step_result_from_success_captures_all_fields() {
let trace_id = Uuid::nil();
let output = StepOutput {
output: json!({"stdout": "ok"}),
duration_ms: 1500,
cost_usd: Decimal::new(42, 2),
input_tokens: Some(100),
output_tokens: Some(200),
model: Some("claude-sonnet".to_string()),
debug_messages: None,
};
let result = StepResult::from_success(trace_id, "build", &output);
assert_eq!(result.trace_id, trace_id);
assert_eq!(result.name, "build");
assert_eq!(result.status, StepStatus::Completed);
assert_eq!(result.duration_ms, 1500);
assert_eq!(result.cost_usd, Decimal::new(42, 2));
assert_eq!(result.input_tokens, Some(100));
assert_eq!(result.output_tokens, Some(200));
assert!(result.error.is_none());
assert!(result.output_summary.is_some());
assert!(result.output_summary.unwrap().contains("stdout"));
}
#[test]
fn step_result_from_failure_captures_error() {
let trace_id = Uuid::nil();
let result =
StepResult::from_failure(trace_id, "deploy", "connection refused", 500, Decimal::ZERO);
assert_eq!(result.trace_id, trace_id);
assert_eq!(result.name, "deploy");
assert_eq!(result.status, StepStatus::Failed);
assert_eq!(result.duration_ms, 500);
assert_eq!(result.error, Some("connection refused".to_string()));
assert!(result.output_summary.is_none());
}
#[test]
fn step_result_output_summary_truncates_long_output() {
let long_value = json!({"data": "x".repeat(1000)});
let output = StepOutput {
output: long_value,
duration_ms: 0,
cost_usd: Decimal::ZERO,
input_tokens: None,
output_tokens: None,
model: None,
debug_messages: None,
};
let result = StepResult::from_success(Uuid::nil(), "test", &output);
let summary = result.output_summary.unwrap();
assert_eq!(summary.len(), 500);
}
}
#[cfg(test)]
mod output_helper_tests {
use super::*;
use serde::Deserialize;
use serde_json::json;
fn output(value: Value) -> StepOutput {
StepOutput {
output: value,
duration_ms: 1,
cost_usd: Decimal::ZERO,
input_tokens: None,
output_tokens: None,
model: None,
debug_messages: None,
}
}
#[test]
fn shell_helpers_read_shell_fields() {
let out = output(json!({"stdout": "hi\n", "stderr": "warn", "exit_code": 0}));
assert_eq!(out.exit_code(), Some(0));
assert_eq!(out.stdout(), "hi\n");
assert_eq!(out.stderr(), "warn");
assert!(out.is_success());
assert_eq!(out.status(), None);
assert_eq!(out.body(), "");
}
#[test]
fn shell_non_zero_exit_is_not_success() {
let out = output(json!({"stdout": "", "stderr": "", "exit_code": 127}));
assert_eq!(out.exit_code(), Some(127));
assert!(!out.is_success());
}
#[test]
fn http_helpers_read_http_fields() {
let out = output(json!({"status": 200, "body": "{\"ok\":true}"}));
assert_eq!(out.status(), Some(200));
assert_eq!(out.body(), "{\"ok\":true}");
assert!(out.is_success());
assert_eq!(out.exit_code(), None);
assert_eq!(out.stdout(), "");
}
#[test]
fn http_error_status_is_not_success() {
assert!(!output(json!({"status": 500, "body": ""})).is_success());
assert!(!output(json!({"status": 199, "body": ""})).is_success());
assert!(output(json!({"status": 299, "body": ""})).is_success());
}
#[test]
fn status_out_of_u16_range_is_none() {
assert_eq!(output(json!({"status": 70000})).status(), None);
assert_eq!(output(json!({"status": "200"})).status(), None);
}
#[test]
fn agent_output_without_markers_is_not_success() {
let out = output(json!({"summary": "fine"}));
assert!(!out.is_success());
assert_eq!(out.exit_code(), None);
assert_eq!(out.stdout(), "");
assert_eq!(out.body(), "");
}
#[test]
fn json_deserializes_structured_output() {
#[derive(Deserialize, Debug, PartialEq)]
struct Review {
score: u8,
summary: String,
}
let out = output(json!({"score": 9, "summary": "good"}));
let review: Review = out.json().expect("matches schema");
assert_eq!(
review,
Review {
score: 9,
summary: "good".to_string()
}
);
}
#[test]
fn json_reports_mismatch_as_serialization_error() {
#[derive(Deserialize, Debug)]
struct Review {
#[allow(dead_code)]
score: u8,
}
let out = output(json!({"score": "nine"}));
let err = out.json::<Review>().expect_err("type mismatch");
assert!(matches!(err, EngineError::Serialization(_)));
}
#[test]
fn helpers_tolerate_non_object_output() {
let out = output(json!("plain text"));
assert_eq!(out.exit_code(), None);
assert_eq!(out.status(), None);
assert_eq!(out.stdout(), "");
assert!(!out.is_success());
}
}