1mod 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#[derive(Debug, Clone)]
43pub struct StepOutput {
44 pub output: Value,
51 pub duration_ms: u64,
53 pub cost_usd: Decimal,
55 pub input_tokens: Option<u64>,
57 pub output_tokens: Option<u64>,
59 pub model: Option<String>,
61 pub debug_messages: Option<Vec<DebugMessage>>,
63}
64
65impl StepOutput {
66 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 pub fn exit_code(&self) -> Option<i64> {
98 self.output.get("exit_code").and_then(Value::as_i64)
99 }
100
101 pub fn stdout(&self) -> &str {
122 self.output
123 .get("stdout")
124 .and_then(Value::as_str)
125 .unwrap_or_default()
126 }
127
128 pub fn stderr(&self) -> &str {
149 self.output
150 .get("stderr")
151 .and_then(Value::as_str)
152 .unwrap_or_default()
153 }
154
155 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 pub fn body(&self) -> &str {
205 self.output
206 .get("body")
207 .and_then(Value::as_str)
208 .unwrap_or_default()
209 }
210
211 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 pub fn json<T: DeserializeOwned>(&self) -> Result<T, EngineError> {
287 from_value(self.output.clone()).map_err(EngineError::Serialization)
288 }
289}
290
291#[derive(Debug, Clone)]
293pub struct ParallelStepResult {
294 pub name: String,
296 pub output: StepOutput,
298 pub step_id: Uuid,
300}
301
302#[derive(Debug, Clone, serde::Serialize)]
330pub struct StepResult {
331 pub trace_id: Uuid,
333 pub name: String,
335 pub status: StepStatus,
337 pub duration_ms: u64,
339 pub cost_usd: Decimal,
341 pub input_tokens: Option<u64>,
343 pub output_tokens: Option<u64>,
345 pub error: Option<String>,
347 pub output_summary: Option<String>,
349}
350
351const OUTPUT_SUMMARY_MAX_LEN: usize = 500;
353
354impl StepResult {
355 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 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
400pub trait StepExecutor: Send + Sync {
405 fn kind(&self) -> StepKind;
411
412 fn execute(
418 &self,
419 provider: &Arc<dyn AgentProvider>,
420 ) -> impl Future<Output = Result<StepOutput, EngineError>> + Send;
421}
422
423pub(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#[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
533pub 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 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 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 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}