1mod 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, SignalOutcome, StepInterceptor};
42pub use shell::ShellExecutor;
43pub use step_artifacts::StepArtifacts;
44pub use workflow_output::{ConcurrencyConflict, SubWorkflowOutcome, SubWorkflowOutput};
45
46pub(crate) use workflow_output::RecordedWorkflowStep;
47
48pub(crate) const ERROR_KEY: &str = "error";
50
51#[derive(Debug, Clone)]
53pub struct StepOutput {
54 pub output: Value,
61 pub duration_ms: u64,
63 pub cost_usd: Decimal,
65 pub input_tokens: Option<u64>,
67 pub cache_read_input_tokens: Option<u64>,
69 pub cache_creation_input_tokens: Option<u64>,
71 pub output_tokens: Option<u64>,
73 pub model: Option<String>,
75 pub debug_messages: Option<Vec<DebugMessage>>,
77 pub artifacts: StepArtifacts,
80 pub account_id: Option<Uuid>,
83 pub environment_id: Option<String>,
87}
88
89impl StepOutput {
90 pub fn total_tokens(&self) -> u64 {
117 [
118 self.input_tokens,
119 self.cache_read_input_tokens,
120 self.cache_creation_input_tokens,
121 self.output_tokens,
122 ]
123 .into_iter()
124 .map(|t| t.unwrap_or(0))
125 .fold(0u64, u64::saturating_add)
126 }
127
128 pub fn debug_messages_json(&self) -> Option<Value> {
132 self.debug_messages
133 .as_ref()
134 .and_then(|msgs| serde_json::to_value(msgs).ok())
135 }
136
137 pub fn exit_code(&self) -> Option<i64> {
165 self.output.get("exit_code").and_then(Value::as_i64)
166 }
167
168 pub fn stdout(&self) -> &str {
194 self.output
195 .get("stdout")
196 .and_then(Value::as_str)
197 .unwrap_or_default()
198 }
199
200 pub fn stderr(&self) -> &str {
226 self.output
227 .get("stderr")
228 .and_then(Value::as_str)
229 .unwrap_or_default()
230 }
231
232 pub fn status(&self) -> Option<u16> {
260 self.output
261 .get("status")
262 .and_then(Value::as_u64)
263 .and_then(|s| u16::try_from(s).ok())
264 }
265
266 pub fn body(&self) -> &str {
292 self.output
293 .get("body")
294 .and_then(Value::as_str)
295 .unwrap_or_default()
296 }
297
298 pub fn text(&self) -> &str {
326 self.output.as_str().unwrap_or_default()
327 }
328
329 pub fn error(&self) -> Option<&str> {
387 self.output.get(ERROR_KEY).and_then(Value::as_str)
388 }
389
390 pub fn is_success(&self) -> bool {
426 if let Some(code) = self.exit_code() {
427 return code == 0;
428 }
429 if let Some(status) = self.status() {
430 return (200..300).contains(&status);
431 }
432 false
433 }
434
435 pub fn json<T: DeserializeOwned>(&self) -> Result<T, EngineError> {
476 from_value(self.output.clone()).map_err(EngineError::Serialization)
477 }
478}
479
480#[derive(Debug, Clone)]
482pub struct ParallelStepResult {
483 pub name: String,
485 pub output: StepOutput,
487 pub step_id: Uuid,
489}
490
491#[derive(Debug, Clone, serde::Serialize)]
519pub struct StepResult {
520 pub trace_id: Uuid,
522 pub name: String,
524 pub status: StepStatus,
526 pub duration_ms: u64,
528 pub cost_usd: Decimal,
530 pub input_tokens: Option<u64>,
532 pub output_tokens: Option<u64>,
534 pub error: Option<String>,
536 pub output_summary: Option<String>,
538}
539
540const OUTPUT_SUMMARY_MAX_LEN: usize = 500;
542
543impl StepResult {
544 pub fn from_success(trace_id: Uuid, name: &str, output: &StepOutput) -> Self {
546 Self {
547 trace_id,
548 name: name.to_string(),
549 status: StepStatus::Completed,
550 duration_ms: output.duration_ms,
551 cost_usd: output.cost_usd,
552 input_tokens: output.input_tokens,
553 output_tokens: output.output_tokens,
554 error: None,
555 output_summary: summarize_output(&output.output),
556 }
557 }
558
559 pub fn from_failure(
561 trace_id: Uuid,
562 name: &str,
563 error: &str,
564 duration_ms: u64,
565 cost_usd: Decimal,
566 ) -> Self {
567 Self {
568 trace_id,
569 name: name.to_string(),
570 status: StepStatus::Failed,
571 duration_ms,
572 cost_usd,
573 input_tokens: None,
574 output_tokens: None,
575 error: Some(error.to_string()),
576 output_summary: None,
577 }
578 }
579}
580
581fn summarize_output(value: &Value) -> Option<String> {
582 let raw = value.to_string();
583 match raw.char_indices().nth(OUTPUT_SUMMARY_MAX_LEN) {
584 None => Some(raw),
585 Some((byte_idx, _)) => Some(raw[..byte_idx].to_string()),
586 }
587}
588
589pub trait StepExecutor: Send + Sync {
594 fn kind(&self) -> StepKind;
600
601 fn execute(
607 &self,
608 provider: &Arc<dyn AgentProvider>,
609 ) -> impl Future<Output = Result<StepOutput, EngineError>> + Send;
610}
611
612pub(crate) fn step_kind_label(kind: &StepKind) -> Cow<'static, str> {
614 match kind {
615 StepKind::Shell => Cow::Borrowed("shell"),
616 StepKind::Http => Cow::Borrowed("http"),
617 StepKind::Agent => Cow::Borrowed("agent"),
618 StepKind::Workflow => Cow::Borrowed("workflow"),
619 StepKind::Approval => Cow::Borrowed("approval"),
620 StepKind::Decision => Cow::Borrowed("decision"),
621 StepKind::HumanInput => Cow::Borrowed("human_input"),
622 StepKind::Signal => Cow::Borrowed("signal"),
623 StepKind::Custom(name) => Cow::Owned(name.clone()),
624 }
625}
626
627#[tracing::instrument(name = "executor.execute_step", skip_all, fields(step.kind))]
658pub async fn execute_step_config_intercepted(
659 config: &StepConfig,
660 provider: &Arc<dyn AgentProvider>,
661 log_sender: Option<StepLogSender>,
662 interceptor: Option<&Arc<dyn StepInterceptor>>,
663) -> Result<StepOutput, EngineError> {
664 let kind = config.kind();
665 let label = step_kind_label(&kind);
666 Span::current().record("step.kind", label.as_ref());
667
668 let intercepted = interceptor.and_then(|i| i.intercept(config));
669 let result = match intercepted {
670 Some(result) => result,
671 None => match config {
672 StepConfig::Shell(cfg) => {
673 let mut executor = ShellExecutor::new(cfg);
674 if let Some(sender) = log_sender {
675 executor = executor.with_log_sender(sender);
676 }
677 executor.execute(provider).await
678 }
679 StepConfig::Http(cfg) => HttpExecutor::new(cfg).execute(provider).await,
680 StepConfig::Agent(cfg) => {
681 let mut executor = AgentExecutor::new(cfg);
682 if let Some(sender) = log_sender {
683 executor = executor.with_log_sender(sender);
684 }
685 executor.execute(provider).await
686 }
687 StepConfig::Workflow(_) => Err(EngineError::StepConfig(
688 "workflow steps are executed by WorkflowContext, not the executor".to_string(),
689 )),
690 StepConfig::Approval(_) => Err(EngineError::StepConfig(
691 "approval steps are executed by WorkflowContext, not the executor".to_string(),
692 )),
693 StepConfig::Decision(_) => Err(EngineError::StepConfig(
694 "decision steps are executed by WorkflowContext, not the executor".to_string(),
695 )),
696 StepConfig::Delay(_) => Err(EngineError::StepConfig(
697 "delay steps are executed by WorkflowContext, not the executor".to_string(),
698 )),
699 },
700 };
701
702 #[cfg(feature = "prometheus")]
703 {
704 use ironflow_core::metric_names::{
705 STATUS_ERROR, STATUS_SUCCESS, STEP_DURATION_SECONDS, STEPS_TOTAL,
706 };
707 use metrics::{counter, histogram};
708 let status = if result.is_ok() {
709 STATUS_SUCCESS
710 } else {
711 STATUS_ERROR
712 };
713 let kind_label = label.into_owned();
714 counter!(STEPS_TOTAL, "kind" => kind_label.clone(), "status" => status).increment(1);
715 if let Ok(ref output) = result {
716 histogram!(STEP_DURATION_SECONDS, "kind" => kind_label)
717 .record(output.duration_ms as f64 / 1000.0);
718 }
719 }
720
721 result
722}
723
724pub async fn execute_step_config(
750 config: &StepConfig,
751 provider: &Arc<dyn AgentProvider>,
752 log_sender: Option<StepLogSender>,
753) -> Result<StepOutput, EngineError> {
754 execute_step_config_intercepted(config, provider, log_sender, None).await
755}
756
757#[cfg(test)]
758mod tests {
759 use super::*;
760 use ironflow_core::provider::DebugMessage;
761 use ironflow_core::providers::claude::ClaudeCodeProvider;
762 use ironflow_core::providers::record_replay::RecordReplayProvider;
763 use serde_json::json;
764
765 use crate::config::{
766 AgentStepConfig, ApprovalConfig, DecisionConfig, DelayConfig, HttpConfig, ShellConfig,
767 WorkflowStepConfig,
768 };
769
770 fn output_of(value: Value) -> StepOutput {
771 StepOutput {
772 output: value,
773 duration_ms: 1,
774 cost_usd: rust_decimal::Decimal::ZERO,
775 input_tokens: None,
776 cache_read_input_tokens: None,
777 cache_creation_input_tokens: None,
778 output_tokens: None,
779 model: None,
780 debug_messages: None,
781 artifacts: StepArtifacts::default(),
782 account_id: None,
783 environment_id: None,
784 }
785 }
786
787 #[test]
788 fn step_output_error_is_some_for_failure_output() {
789 let output = output_of(json!({"error": "boom"}));
790
791 assert_eq!(output.error(), Some("boom"));
792 assert!(!output.is_success());
793 }
794
795 #[test]
796 fn step_output_error_is_none_for_success_shell_output() {
797 let output = output_of(json!({"stdout": "hi", "stderr": "", "exit_code": 0}));
798
799 assert_eq!(output.error(), None);
800 assert!(output.is_success());
801 }
802
803 #[test]
804 fn step_output_error_is_none_for_http_and_agent_text() {
805 assert_eq!(output_of(json!({"status": 200, "body": ""})).error(), None);
806 assert_eq!(output_of(json!("fine")).error(), None);
807 }
808
809 #[test]
810 fn step_output_error_is_none_when_error_key_not_a_string() {
811 assert_eq!(output_of(json!({"error": 5})).error(), None);
812 assert_eq!(output_of(Value::Null).error(), None);
813 }
814
815 #[test]
816 fn step_output_with_no_debug_messages_returns_none() {
817 let output = StepOutput {
818 output: json!({"result": "ok"}),
819 duration_ms: 100,
820 cost_usd: rust_decimal::Decimal::ZERO,
821 input_tokens: None,
822 cache_read_input_tokens: None,
823 cache_creation_input_tokens: None,
824 output_tokens: None,
825 model: None,
826 debug_messages: None,
827 artifacts: StepArtifacts::default(),
828 account_id: None,
829 environment_id: None,
830 };
831
832 assert_eq!(output.debug_messages_json(), None);
833 }
834
835 #[test]
836 fn step_output_with_empty_debug_messages_returns_some_empty_array() {
837 let output = StepOutput {
838 output: json!({"result": "ok"}),
839 duration_ms: 100,
840 cost_usd: rust_decimal::Decimal::ZERO,
841 input_tokens: None,
842 cache_read_input_tokens: None,
843 cache_creation_input_tokens: None,
844 output_tokens: None,
845 model: None,
846 debug_messages: Some(Vec::new()),
847 artifacts: StepArtifacts::default(),
848 account_id: None,
849 environment_id: None,
850 };
851
852 let json_val = output.debug_messages_json();
853 assert!(json_val.is_some());
854 let arr = json_val.unwrap();
855 assert!(arr.is_array());
856 assert_eq!(arr.as_array().unwrap().len(), 0);
857 }
858
859 #[test]
860 fn step_output_debug_messages_json_serializes_messages() {
861 let json_msgs = json!([
862 {
863 "text": "Hello",
864 "thinking": null,
865 "thinking_redacted": false,
866 "tool_calls": [],
867 "tool_results": [],
868 "stop_reason": "end_turn",
869 "input_tokens": 10,
870 "output_tokens": 20
871 },
872 {
873 "text": "Hi there",
874 "thinking": null,
875 "thinking_redacted": false,
876 "tool_calls": [],
877 "tool_results": [],
878 "stop_reason": "end_turn",
879 "input_tokens": 15,
880 "output_tokens": 25
881 }
882 ]);
883
884 let messages: Vec<DebugMessage> =
885 serde_json::from_value(json_msgs.clone()).expect("deserialize debug messages");
886
887 let output = StepOutput {
888 output: json!({"result": "ok"}),
889 duration_ms: 100,
890 cost_usd: rust_decimal::Decimal::ZERO,
891 input_tokens: None,
892 cache_read_input_tokens: None,
893 cache_creation_input_tokens: None,
894 output_tokens: None,
895 model: None,
896 debug_messages: Some(messages),
897 artifacts: StepArtifacts::default(),
898 account_id: None,
899 environment_id: None,
900 };
901
902 let json_val = output.debug_messages_json();
903 assert!(json_val.is_some());
904
905 let arr = json_val.unwrap();
906 assert!(arr.is_array());
907 let messages_array = arr.as_array().unwrap();
908 assert_eq!(messages_array.len(), 2);
909 assert_eq!(messages_array[0]["text"], "Hello");
910 assert_eq!(messages_array[1]["text"], "Hi there");
911 }
912
913 #[test]
914 fn step_output_contains_all_metrics() {
915 let output = StepOutput {
916 output: json!({"data": "test"}),
917 duration_ms: 5000,
918 cost_usd: rust_decimal::Decimal::new(123, 2),
919 input_tokens: Some(100),
920 cache_read_input_tokens: None,
921 cache_creation_input_tokens: None,
922 output_tokens: Some(200),
923 model: Some("claude-sonnet".to_string()),
924 debug_messages: None,
925 artifacts: StepArtifacts::default(),
926 account_id: None,
927 environment_id: None,
928 };
929
930 assert_eq!(output.duration_ms, 5000);
931 assert_eq!(output.cost_usd, rust_decimal::Decimal::new(123, 2));
932 assert_eq!(output.input_tokens, Some(100));
933 assert_eq!(output.output_tokens, Some(200));
934 assert_eq!(output.model, Some("claude-sonnet".to_string()));
935 }
936
937 #[test]
938 fn step_output_default_tokens_and_model_are_none() {
939 let output = StepOutput {
940 output: json!({}),
941 duration_ms: 0,
942 cost_usd: rust_decimal::Decimal::ZERO,
943 input_tokens: None,
944 cache_read_input_tokens: None,
945 cache_creation_input_tokens: None,
946 output_tokens: None,
947 model: None,
948 debug_messages: None,
949 artifacts: StepArtifacts::default(),
950 account_id: None,
951 environment_id: None,
952 };
953
954 assert!(output.input_tokens.is_none());
955 assert!(output.output_tokens.is_none());
956 assert!(output.model.is_none());
957 }
958
959 #[test]
960 fn parallel_step_result_contains_step_metadata() {
961 let step_id = uuid::Uuid::now_v7();
962 let output = StepOutput {
963 output: json!({"done": true}),
964 duration_ms: 1000,
965 cost_usd: rust_decimal::Decimal::ZERO,
966 input_tokens: None,
967 cache_read_input_tokens: None,
968 cache_creation_input_tokens: None,
969 output_tokens: None,
970 model: None,
971 debug_messages: None,
972 artifacts: StepArtifacts::default(),
973 account_id: None,
974 environment_id: None,
975 };
976
977 let result = ParallelStepResult {
978 name: "build".to_string(),
979 output,
980 step_id,
981 };
982
983 assert_eq!(result.name, "build");
984 assert_eq!(result.step_id, step_id);
985 assert_eq!(result.output.duration_ms, 1000);
986 }
987
988 #[test]
989 fn step_output_serializes_complex_json_output() {
990 let complex_output = json!({
991 "status": "success",
992 "data": {
993 "items": [1, 2, 3],
994 "nested": {
995 "key": "value"
996 }
997 }
998 });
999
1000 let output = StepOutput {
1001 output: complex_output.clone(),
1002 duration_ms: 100,
1003 cost_usd: rust_decimal::Decimal::ZERO,
1004 input_tokens: None,
1005 cache_read_input_tokens: None,
1006 cache_creation_input_tokens: None,
1007 output_tokens: None,
1008 model: None,
1009 debug_messages: None,
1010 artifacts: StepArtifacts::default(),
1011 account_id: None,
1012 environment_id: None,
1013 };
1014
1015 assert_eq!(output.output, complex_output);
1016 assert_eq!(output.output["status"], "success");
1017 assert_eq!(output.output["data"]["items"][0], 1);
1018 assert_eq!(output.output["data"]["nested"]["key"], "value");
1019 }
1020
1021 #[test]
1022 fn step_result_from_success_captures_all_fields() {
1023 let trace_id = Uuid::nil();
1024 let output = StepOutput {
1025 output: json!({"stdout": "ok"}),
1026 duration_ms: 1500,
1027 cost_usd: Decimal::new(42, 2),
1028 input_tokens: Some(100),
1029 cache_read_input_tokens: None,
1030 cache_creation_input_tokens: None,
1031 output_tokens: Some(200),
1032 model: Some("claude-sonnet".to_string()),
1033 debug_messages: None,
1034 artifacts: StepArtifacts::default(),
1035 account_id: None,
1036 environment_id: None,
1037 };
1038
1039 let result = StepResult::from_success(trace_id, "build", &output);
1040
1041 assert_eq!(result.trace_id, trace_id);
1042 assert_eq!(result.name, "build");
1043 assert_eq!(result.status, StepStatus::Completed);
1044 assert_eq!(result.duration_ms, 1500);
1045 assert_eq!(result.cost_usd, Decimal::new(42, 2));
1046 assert_eq!(result.input_tokens, Some(100));
1047 assert_eq!(result.output_tokens, Some(200));
1048 assert!(result.error.is_none());
1049 assert!(result.output_summary.is_some());
1050 assert!(result.output_summary.unwrap().contains("stdout"));
1051 }
1052
1053 #[test]
1054 fn step_result_from_failure_captures_error() {
1055 let trace_id = Uuid::nil();
1056 let result =
1057 StepResult::from_failure(trace_id, "deploy", "connection refused", 500, Decimal::ZERO);
1058
1059 assert_eq!(result.trace_id, trace_id);
1060 assert_eq!(result.name, "deploy");
1061 assert_eq!(result.status, StepStatus::Failed);
1062 assert_eq!(result.duration_ms, 500);
1063 assert_eq!(result.error, Some("connection refused".to_string()));
1064 assert!(result.output_summary.is_none());
1065 }
1066
1067 #[test]
1068 fn step_result_output_summary_truncates_long_output() {
1069 let long_value = json!({"data": "x".repeat(1000)});
1070 let output = StepOutput {
1071 output: long_value,
1072 duration_ms: 0,
1073 cost_usd: Decimal::ZERO,
1074 input_tokens: None,
1075 cache_read_input_tokens: None,
1076 cache_creation_input_tokens: None,
1077 output_tokens: None,
1078 model: None,
1079 debug_messages: None,
1080 artifacts: StepArtifacts::default(),
1081 account_id: None,
1082 environment_id: None,
1083 };
1084
1085 let result = StepResult::from_success(Uuid::nil(), "test", &output);
1086 let summary = result.output_summary.unwrap();
1087 assert_eq!(summary.len(), 500);
1088 }
1089
1090 #[test]
1091 fn step_executor_kind_matches_step_config_kind() {
1092 let shell = ShellConfig::new("echo hi");
1093 let http = HttpConfig::get("https://example.com");
1094 let agent = AgentStepConfig::new("hi");
1095
1096 let shell_kind = ShellExecutor::new(&shell).kind();
1097 let http_kind = HttpExecutor::new(&http).kind();
1098 let agent_kind = AgentExecutor::new(&agent).kind();
1099
1100 assert_eq!(shell_kind, StepConfig::Shell(shell).kind());
1101 assert_eq!(http_kind, StepConfig::Http(http).kind());
1102 assert_eq!(agent_kind, StepConfig::Agent(agent).kind());
1103 }
1104
1105 #[test]
1106 fn step_kind_label_matches_dispatcher_labels() {
1107 let cases: Vec<(StepConfig, &str)> = vec![
1108 (StepConfig::Shell(ShellConfig::new("echo hi")), "shell"),
1109 (
1110 StepConfig::Http(HttpConfig::get("https://example.com")),
1111 "http",
1112 ),
1113 (StepConfig::Agent(AgentStepConfig::new("hi")), "agent"),
1114 (
1115 StepConfig::Workflow(WorkflowStepConfig::new("child", json!({}))),
1116 "workflow",
1117 ),
1118 (
1119 StepConfig::Approval(ApprovalConfig::new("approve?")),
1120 "approval",
1121 ),
1122 (
1123 StepConfig::Decision(DecisionConfig::new(json!({}))),
1124 "decision",
1125 ),
1126 (StepConfig::Delay(DelayConfig::from_secs(1)), "delay"),
1127 ];
1128
1129 for (config, expected) in cases {
1130 assert_eq!(step_kind_label(&config.kind()), expected);
1131 }
1132 }
1133
1134 #[test]
1135 fn step_kind_label_uses_the_custom_kind_name() {
1136 assert_eq!(
1137 step_kind_label(&StepKind::Custom("gitlab".to_string())),
1138 "gitlab"
1139 );
1140 }
1141
1142 struct CannedShell;
1144
1145 impl StepInterceptor for CannedShell {
1146 fn intercept(&self, config: &StepConfig) -> Option<Result<StepOutput, EngineError>> {
1147 match config {
1148 StepConfig::Shell(_) => Some(Ok(StepOutput {
1149 output: json!({"stdout": "canned", "stderr": "", "exit_code": 0}),
1150 duration_ms: 0,
1151 cost_usd: Decimal::ZERO,
1152 input_tokens: None,
1153 cache_read_input_tokens: None,
1154 cache_creation_input_tokens: None,
1155 output_tokens: None,
1156 model: None,
1157 debug_messages: None,
1158 artifacts: StepArtifacts::default(),
1159 account_id: None,
1160 environment_id: None,
1161 })),
1162 _ => None,
1163 }
1164 }
1165 }
1166
1167 fn test_provider() -> Arc<dyn AgentProvider> {
1168 let inner = ClaudeCodeProvider::new();
1169 Arc::new(RecordReplayProvider::replay(
1170 inner,
1171 "/tmp/ironflow-fixtures",
1172 ))
1173 }
1174
1175 #[tokio::test]
1176 async fn intercepted_step_never_reaches_the_shell_executor() {
1177 let interceptor: Arc<dyn StepInterceptor> = Arc::new(CannedShell);
1178 let config = StepConfig::Shell(ShellConfig::new("exit 1"));
1181
1182 let output =
1183 execute_step_config_intercepted(&config, &test_provider(), None, Some(&interceptor))
1184 .await
1185 .expect("the interceptor resolved the step");
1186
1187 assert_eq!(output.stdout(), "canned");
1188 assert_eq!(output.exit_code(), Some(0));
1189 }
1190
1191 #[tokio::test]
1192 async fn a_step_the_interceptor_declines_reaches_the_dispatcher() {
1193 let interceptor: Arc<dyn StepInterceptor> = Arc::new(CannedShell);
1194 let config = StepConfig::Workflow(WorkflowStepConfig::new("child", json!({})));
1197
1198 let err =
1199 execute_step_config_intercepted(&config, &test_provider(), None, Some(&interceptor))
1200 .await
1201 .expect_err("the dispatcher rejects workflow steps");
1202
1203 assert!(matches!(err, EngineError::StepConfig(_)));
1204 }
1205
1206 #[tokio::test]
1207 async fn without_an_interceptor_the_step_runs_for_real() {
1208 let config = StepConfig::Shell(ShellConfig::new("echo hi"));
1209
1210 let output = execute_step_config_intercepted(&config, &test_provider(), None, None)
1211 .await
1212 .expect("echo succeeds");
1213
1214 assert!(output.stdout().contains("hi"));
1215 assert_eq!(output.exit_code(), Some(0));
1216 }
1217}
1218
1219#[cfg(test)]
1220mod output_helper_tests {
1221 use super::*;
1222 use serde::Deserialize;
1223 use serde_json::json;
1224
1225 fn output(value: Value) -> StepOutput {
1226 StepOutput {
1227 output: value,
1228 duration_ms: 1,
1229 cost_usd: Decimal::ZERO,
1230 input_tokens: None,
1231 cache_read_input_tokens: None,
1232 cache_creation_input_tokens: None,
1233 output_tokens: None,
1234 model: None,
1235 debug_messages: None,
1236 artifacts: StepArtifacts::default(),
1237 account_id: None,
1238 environment_id: None,
1239 }
1240 }
1241
1242 #[test]
1243 fn agent_total_tokens_includes_cache_tokens() {
1244 let mut out = output(json!("ok"));
1245 out.input_tokens = Some(100);
1246 out.cache_read_input_tokens = Some(5000);
1247 out.cache_creation_input_tokens = Some(200);
1248 out.output_tokens = Some(50);
1249 assert_eq!(out.total_tokens(), 5350);
1250 }
1251
1252 #[test]
1253 fn agent_total_tokens_all_none_is_zero() {
1254 let out = output(json!("ok"));
1255 assert_eq!(out.total_tokens(), 0);
1256 }
1257
1258 #[test]
1259 fn agent_total_tokens_saturates() {
1260 let mut out = output(json!("ok"));
1261 out.input_tokens = Some(u64::MAX);
1262 out.cache_read_input_tokens = Some(10);
1263 assert_eq!(out.total_tokens(), u64::MAX);
1264 }
1265
1266 #[test]
1267 fn shell_helpers_read_shell_fields() {
1268 let out = output(json!({"stdout": "hi\n", "stderr": "warn", "exit_code": 0}));
1269 assert_eq!(out.exit_code(), Some(0));
1270 assert_eq!(out.stdout(), "hi\n");
1271 assert_eq!(out.stderr(), "warn");
1272 assert!(out.is_success());
1273 assert_eq!(out.status(), None);
1274 assert_eq!(out.body(), "");
1275 }
1276
1277 #[test]
1278 fn shell_non_zero_exit_is_not_success() {
1279 let out = output(json!({"stdout": "", "stderr": "", "exit_code": 127}));
1280 assert_eq!(out.exit_code(), Some(127));
1281 assert!(!out.is_success());
1282 }
1283
1284 #[test]
1285 fn http_helpers_read_http_fields() {
1286 let out = output(json!({"status": 200, "body": "{\"ok\":true}"}));
1287 assert_eq!(out.status(), Some(200));
1288 assert_eq!(out.body(), "{\"ok\":true}");
1289 assert!(out.is_success());
1290 assert_eq!(out.exit_code(), None);
1291 assert_eq!(out.stdout(), "");
1292 }
1293
1294 #[test]
1295 fn http_error_status_is_not_success() {
1296 assert!(!output(json!({"status": 500, "body": ""})).is_success());
1297 assert!(!output(json!({"status": 199, "body": ""})).is_success());
1298 assert!(output(json!({"status": 299, "body": ""})).is_success());
1299 }
1300
1301 #[test]
1302 fn status_out_of_u16_range_is_none() {
1303 assert_eq!(output(json!({"status": 70000})).status(), None);
1304 assert_eq!(output(json!({"status": "200"})).status(), None);
1305 }
1306
1307 #[test]
1308 fn agent_output_without_markers_is_not_success() {
1309 let out = output(json!({"summary": "fine"}));
1310 assert!(!out.is_success());
1311 assert_eq!(out.exit_code(), None);
1312 assert_eq!(out.stdout(), "");
1313 assert_eq!(out.body(), "");
1314 }
1315
1316 #[test]
1317 fn json_deserializes_structured_output() {
1318 #[derive(Deserialize, Debug, PartialEq)]
1319 struct Review {
1320 score: u8,
1321 summary: String,
1322 }
1323 let out = output(json!({"score": 9, "summary": "good"}));
1324 let review: Review = out.json().expect("matches schema");
1325 assert_eq!(
1326 review,
1327 Review {
1328 score: 9,
1329 summary: "good".to_string()
1330 }
1331 );
1332 }
1333
1334 #[test]
1335 fn json_reports_mismatch_as_serialization_error() {
1336 #[derive(Deserialize, Debug)]
1337 struct Review {
1338 #[allow(dead_code)]
1339 score: u8,
1340 }
1341 let out = output(json!({"score": "nine"}));
1342 let err = out.json::<Review>().expect_err("type mismatch");
1343 assert!(matches!(err, EngineError::Serialization(_)));
1344 }
1345
1346 #[test]
1347 fn helpers_tolerate_non_object_output() {
1348 let out = output(json!("plain text"));
1349 assert_eq!(out.exit_code(), None);
1350 assert_eq!(out.status(), None);
1351 assert_eq!(out.stdout(), "");
1352 assert!(!out.is_success());
1353 }
1354
1355 #[test]
1356 fn text_reads_a_plain_agent_answer_only() {
1357 assert_eq!(output(json!("plain text")).text(), "plain text");
1358 assert_eq!(output(json!({"stdout": "x"})).text(), "");
1359 assert_eq!(output(Value::Null).text(), "");
1360 }
1361}