1mod agent;
13mod decision;
14mod http;
15mod shell;
16
17use std::borrow::Cow;
18use std::future::Future;
19use std::sync::Arc;
20
21use rust_decimal::Decimal;
22use serde::de::DeserializeOwned;
23use serde_json::{Value, from_value};
24use tracing::Span;
25use uuid::Uuid;
26
27use ironflow_core::provider::{AgentProvider, DebugMessage};
28use ironflow_store::entities::{StepKind, StepStatus};
29
30use crate::config::StepConfig;
31use crate::error::EngineError;
32use crate::log_sender::StepLogSender;
33
34pub use agent::AgentExecutor;
35pub use decision::{DecisionExecution, execute_decision};
36pub use http::HttpExecutor;
37pub use shell::ShellExecutor;
38
39#[derive(Debug, Clone)]
41pub struct StepOutput {
42 pub output: Value,
49 pub duration_ms: u64,
51 pub cost_usd: Decimal,
53 pub input_tokens: Option<u64>,
55 pub output_tokens: Option<u64>,
57 pub model: Option<String>,
59 pub debug_messages: Option<Vec<DebugMessage>>,
61}
62
63impl StepOutput {
64 pub fn debug_messages_json(&self) -> Option<Value> {
68 self.debug_messages
69 .as_ref()
70 .and_then(|msgs| serde_json::to_value(msgs).ok())
71 }
72
73 pub fn exit_code(&self) -> Option<i64> {
96 self.output.get("exit_code").and_then(Value::as_i64)
97 }
98
99 pub fn stdout(&self) -> &str {
120 self.output
121 .get("stdout")
122 .and_then(Value::as_str)
123 .unwrap_or_default()
124 }
125
126 pub fn stderr(&self) -> &str {
147 self.output
148 .get("stderr")
149 .and_then(Value::as_str)
150 .unwrap_or_default()
151 }
152
153 pub fn status(&self) -> Option<u16> {
176 self.output
177 .get("status")
178 .and_then(Value::as_u64)
179 .and_then(|s| u16::try_from(s).ok())
180 }
181
182 pub fn body(&self) -> &str {
203 self.output
204 .get("body")
205 .and_then(Value::as_str)
206 .unwrap_or_default()
207 }
208
209 pub fn is_success(&self) -> bool {
240 if let Some(code) = self.exit_code() {
241 return code == 0;
242 }
243 if let Some(status) = self.status() {
244 return (200..300).contains(&status);
245 }
246 false
247 }
248
249 pub fn json<T: DeserializeOwned>(&self) -> Result<T, EngineError> {
285 from_value(self.output.clone()).map_err(EngineError::Serialization)
286 }
287}
288
289#[derive(Debug, Clone)]
291pub struct ParallelStepResult {
292 pub name: String,
294 pub output: StepOutput,
296 pub step_id: Uuid,
298}
299
300#[derive(Debug, Clone, serde::Serialize)]
328pub struct StepResult {
329 pub trace_id: Uuid,
331 pub name: String,
333 pub status: StepStatus,
335 pub duration_ms: u64,
337 pub cost_usd: Decimal,
339 pub input_tokens: Option<u64>,
341 pub output_tokens: Option<u64>,
343 pub error: Option<String>,
345 pub output_summary: Option<String>,
347}
348
349const OUTPUT_SUMMARY_MAX_LEN: usize = 500;
351
352impl StepResult {
353 pub fn from_success(trace_id: Uuid, name: &str, output: &StepOutput) -> Self {
355 Self {
356 trace_id,
357 name: name.to_string(),
358 status: StepStatus::Completed,
359 duration_ms: output.duration_ms,
360 cost_usd: output.cost_usd,
361 input_tokens: output.input_tokens,
362 output_tokens: output.output_tokens,
363 error: None,
364 output_summary: summarize_output(&output.output),
365 }
366 }
367
368 pub fn from_failure(
370 trace_id: Uuid,
371 name: &str,
372 error: &str,
373 duration_ms: u64,
374 cost_usd: Decimal,
375 ) -> Self {
376 Self {
377 trace_id,
378 name: name.to_string(),
379 status: StepStatus::Failed,
380 duration_ms,
381 cost_usd,
382 input_tokens: None,
383 output_tokens: None,
384 error: Some(error.to_string()),
385 output_summary: None,
386 }
387 }
388}
389
390fn summarize_output(value: &Value) -> Option<String> {
391 let raw = value.to_string();
392 match raw.char_indices().nth(OUTPUT_SUMMARY_MAX_LEN) {
393 None => Some(raw),
394 Some((byte_idx, _)) => Some(raw[..byte_idx].to_string()),
395 }
396}
397
398pub trait StepExecutor: Send + Sync {
403 fn kind(&self) -> StepKind;
409
410 fn execute(
416 &self,
417 provider: &Arc<dyn AgentProvider>,
418 ) -> impl Future<Output = Result<StepOutput, EngineError>> + Send;
419}
420
421pub(crate) fn step_kind_label(kind: &StepKind) -> Cow<'static, str> {
423 match kind {
424 StepKind::Shell => Cow::Borrowed("shell"),
425 StepKind::Http => Cow::Borrowed("http"),
426 StepKind::Agent => Cow::Borrowed("agent"),
427 StepKind::Workflow => Cow::Borrowed("workflow"),
428 StepKind::Approval => Cow::Borrowed("approval"),
429 StepKind::Decision => Cow::Borrowed("decision"),
430 StepKind::Custom(name) => Cow::Owned(name.clone()),
431 }
432}
433
434#[tracing::instrument(name = "executor.execute_step", skip_all, fields(step.kind))]
460pub async fn execute_step_config(
461 config: &StepConfig,
462 provider: &Arc<dyn AgentProvider>,
463 log_sender: Option<StepLogSender>,
464) -> Result<StepOutput, EngineError> {
465 let kind = config.kind();
466 let label = step_kind_label(&kind);
467 Span::current().record("step.kind", label.as_ref());
468
469 let result = match config {
470 StepConfig::Shell(cfg) => {
471 let mut executor = ShellExecutor::new(cfg);
472 if let Some(sender) = log_sender {
473 executor = executor.with_log_sender(sender);
474 }
475 executor.execute(provider).await
476 }
477 StepConfig::Http(cfg) => HttpExecutor::new(cfg).execute(provider).await,
478 StepConfig::Agent(cfg) => {
479 let mut executor = AgentExecutor::new(cfg);
480 if let Some(sender) = log_sender {
481 executor = executor.with_log_sender(sender);
482 }
483 executor.execute(provider).await
484 }
485 StepConfig::Workflow(_) => Err(EngineError::StepConfig(
486 "workflow steps are executed by WorkflowContext, not the executor".to_string(),
487 )),
488 StepConfig::Approval(_) => Err(EngineError::StepConfig(
489 "approval steps are executed by WorkflowContext, not the executor".to_string(),
490 )),
491 StepConfig::Decision(_) => Err(EngineError::StepConfig(
492 "decision steps are executed by WorkflowContext, not the executor".to_string(),
493 )),
494 StepConfig::Delay(_) => Err(EngineError::StepConfig(
495 "delay steps are executed by WorkflowContext, not the executor".to_string(),
496 )),
497 };
498
499 #[cfg(feature = "prometheus")]
500 {
501 use ironflow_core::metric_names::{
502 STATUS_ERROR, STATUS_SUCCESS, STEP_DURATION_SECONDS, STEPS_TOTAL,
503 };
504 use metrics::{counter, histogram};
505 let status = if result.is_ok() {
506 STATUS_SUCCESS
507 } else {
508 STATUS_ERROR
509 };
510 let kind_label = label.into_owned();
511 counter!(STEPS_TOTAL, "kind" => kind_label.clone(), "status" => status).increment(1);
512 if let Ok(ref output) = result {
513 histogram!(STEP_DURATION_SECONDS, "kind" => kind_label)
514 .record(output.duration_ms as f64 / 1000.0);
515 }
516 }
517
518 result
519}
520
521#[cfg(test)]
522mod tests {
523 use super::*;
524 use ironflow_core::provider::DebugMessage;
525 use serde_json::json;
526
527 use crate::config::{
528 AgentStepConfig, ApprovalConfig, DecisionConfig, DelayConfig, HttpConfig, ShellConfig,
529 WorkflowStepConfig,
530 };
531
532 #[test]
533 fn step_output_with_no_debug_messages_returns_none() {
534 let output = StepOutput {
535 output: json!({"result": "ok"}),
536 duration_ms: 100,
537 cost_usd: rust_decimal::Decimal::ZERO,
538 input_tokens: None,
539 output_tokens: None,
540 model: None,
541 debug_messages: None,
542 };
543
544 assert_eq!(output.debug_messages_json(), None);
545 }
546
547 #[test]
548 fn step_output_with_empty_debug_messages_returns_some_empty_array() {
549 let output = StepOutput {
550 output: json!({"result": "ok"}),
551 duration_ms: 100,
552 cost_usd: rust_decimal::Decimal::ZERO,
553 input_tokens: None,
554 output_tokens: None,
555 model: None,
556 debug_messages: Some(Vec::new()),
557 };
558
559 let json_val = output.debug_messages_json();
560 assert!(json_val.is_some());
561 let arr = json_val.unwrap();
562 assert!(arr.is_array());
563 assert_eq!(arr.as_array().unwrap().len(), 0);
564 }
565
566 #[test]
567 fn step_output_debug_messages_json_serializes_messages() {
568 let json_msgs = json!([
569 {
570 "text": "Hello",
571 "thinking": null,
572 "thinking_redacted": false,
573 "tool_calls": [],
574 "tool_results": [],
575 "stop_reason": "end_turn",
576 "input_tokens": 10,
577 "output_tokens": 20
578 },
579 {
580 "text": "Hi there",
581 "thinking": null,
582 "thinking_redacted": false,
583 "tool_calls": [],
584 "tool_results": [],
585 "stop_reason": "end_turn",
586 "input_tokens": 15,
587 "output_tokens": 25
588 }
589 ]);
590
591 let messages: Vec<DebugMessage> =
592 serde_json::from_value(json_msgs.clone()).expect("deserialize debug messages");
593
594 let output = StepOutput {
595 output: json!({"result": "ok"}),
596 duration_ms: 100,
597 cost_usd: rust_decimal::Decimal::ZERO,
598 input_tokens: None,
599 output_tokens: None,
600 model: None,
601 debug_messages: Some(messages),
602 };
603
604 let json_val = output.debug_messages_json();
605 assert!(json_val.is_some());
606
607 let arr = json_val.unwrap();
608 assert!(arr.is_array());
609 let messages_array = arr.as_array().unwrap();
610 assert_eq!(messages_array.len(), 2);
611 assert_eq!(messages_array[0]["text"], "Hello");
612 assert_eq!(messages_array[1]["text"], "Hi there");
613 }
614
615 #[test]
616 fn step_output_contains_all_metrics() {
617 let output = StepOutput {
618 output: json!({"data": "test"}),
619 duration_ms: 5000,
620 cost_usd: rust_decimal::Decimal::new(123, 2),
621 input_tokens: Some(100),
622 output_tokens: Some(200),
623 model: Some("claude-sonnet".to_string()),
624 debug_messages: None,
625 };
626
627 assert_eq!(output.duration_ms, 5000);
628 assert_eq!(output.cost_usd, rust_decimal::Decimal::new(123, 2));
629 assert_eq!(output.input_tokens, Some(100));
630 assert_eq!(output.output_tokens, Some(200));
631 assert_eq!(output.model, Some("claude-sonnet".to_string()));
632 }
633
634 #[test]
635 fn step_output_default_tokens_and_model_are_none() {
636 let output = StepOutput {
637 output: json!({}),
638 duration_ms: 0,
639 cost_usd: rust_decimal::Decimal::ZERO,
640 input_tokens: None,
641 output_tokens: None,
642 model: None,
643 debug_messages: None,
644 };
645
646 assert!(output.input_tokens.is_none());
647 assert!(output.output_tokens.is_none());
648 assert!(output.model.is_none());
649 }
650
651 #[test]
652 fn parallel_step_result_contains_step_metadata() {
653 let step_id = uuid::Uuid::now_v7();
654 let output = StepOutput {
655 output: json!({"done": true}),
656 duration_ms: 1000,
657 cost_usd: rust_decimal::Decimal::ZERO,
658 input_tokens: None,
659 output_tokens: None,
660 model: None,
661 debug_messages: None,
662 };
663
664 let result = ParallelStepResult {
665 name: "build".to_string(),
666 output,
667 step_id,
668 };
669
670 assert_eq!(result.name, "build");
671 assert_eq!(result.step_id, step_id);
672 assert_eq!(result.output.duration_ms, 1000);
673 }
674
675 #[test]
676 fn step_output_serializes_complex_json_output() {
677 let complex_output = json!({
678 "status": "success",
679 "data": {
680 "items": [1, 2, 3],
681 "nested": {
682 "key": "value"
683 }
684 }
685 });
686
687 let output = StepOutput {
688 output: complex_output.clone(),
689 duration_ms: 100,
690 cost_usd: rust_decimal::Decimal::ZERO,
691 input_tokens: None,
692 output_tokens: None,
693 model: None,
694 debug_messages: None,
695 };
696
697 assert_eq!(output.output, complex_output);
698 assert_eq!(output.output["status"], "success");
699 assert_eq!(output.output["data"]["items"][0], 1);
700 assert_eq!(output.output["data"]["nested"]["key"], "value");
701 }
702
703 #[test]
704 fn step_result_from_success_captures_all_fields() {
705 let trace_id = Uuid::nil();
706 let output = StepOutput {
707 output: json!({"stdout": "ok"}),
708 duration_ms: 1500,
709 cost_usd: Decimal::new(42, 2),
710 input_tokens: Some(100),
711 output_tokens: Some(200),
712 model: Some("claude-sonnet".to_string()),
713 debug_messages: None,
714 };
715
716 let result = StepResult::from_success(trace_id, "build", &output);
717
718 assert_eq!(result.trace_id, trace_id);
719 assert_eq!(result.name, "build");
720 assert_eq!(result.status, StepStatus::Completed);
721 assert_eq!(result.duration_ms, 1500);
722 assert_eq!(result.cost_usd, Decimal::new(42, 2));
723 assert_eq!(result.input_tokens, Some(100));
724 assert_eq!(result.output_tokens, Some(200));
725 assert!(result.error.is_none());
726 assert!(result.output_summary.is_some());
727 assert!(result.output_summary.unwrap().contains("stdout"));
728 }
729
730 #[test]
731 fn step_result_from_failure_captures_error() {
732 let trace_id = Uuid::nil();
733 let result =
734 StepResult::from_failure(trace_id, "deploy", "connection refused", 500, Decimal::ZERO);
735
736 assert_eq!(result.trace_id, trace_id);
737 assert_eq!(result.name, "deploy");
738 assert_eq!(result.status, StepStatus::Failed);
739 assert_eq!(result.duration_ms, 500);
740 assert_eq!(result.error, Some("connection refused".to_string()));
741 assert!(result.output_summary.is_none());
742 }
743
744 #[test]
745 fn step_result_output_summary_truncates_long_output() {
746 let long_value = json!({"data": "x".repeat(1000)});
747 let output = StepOutput {
748 output: long_value,
749 duration_ms: 0,
750 cost_usd: Decimal::ZERO,
751 input_tokens: None,
752 output_tokens: None,
753 model: None,
754 debug_messages: None,
755 };
756
757 let result = StepResult::from_success(Uuid::nil(), "test", &output);
758 let summary = result.output_summary.unwrap();
759 assert_eq!(summary.len(), 500);
760 }
761
762 #[test]
763 fn step_executor_kind_matches_step_config_kind() {
764 let shell = ShellConfig::new("echo hi");
765 let http = HttpConfig::get("https://example.com");
766 let agent = AgentStepConfig::new("hi");
767
768 let shell_kind = ShellExecutor::new(&shell).kind();
769 let http_kind = HttpExecutor::new(&http).kind();
770 let agent_kind = AgentExecutor::new(&agent).kind();
771
772 assert_eq!(shell_kind, StepConfig::Shell(shell).kind());
773 assert_eq!(http_kind, StepConfig::Http(http).kind());
774 assert_eq!(agent_kind, StepConfig::Agent(agent).kind());
775 }
776
777 #[test]
778 fn step_kind_label_matches_dispatcher_labels() {
779 let cases: Vec<(StepConfig, &str)> = vec![
780 (StepConfig::Shell(ShellConfig::new("echo hi")), "shell"),
781 (
782 StepConfig::Http(HttpConfig::get("https://example.com")),
783 "http",
784 ),
785 (StepConfig::Agent(AgentStepConfig::new("hi")), "agent"),
786 (
787 StepConfig::Workflow(WorkflowStepConfig::new("child", json!({}))),
788 "workflow",
789 ),
790 (
791 StepConfig::Approval(ApprovalConfig::new("approve?")),
792 "approval",
793 ),
794 (
795 StepConfig::Decision(DecisionConfig::new(json!({}))),
796 "decision",
797 ),
798 (StepConfig::Delay(DelayConfig::from_secs(1)), "delay"),
799 ];
800
801 for (config, expected) in cases {
802 assert_eq!(step_kind_label(&config.kind()), expected);
803 }
804 }
805
806 #[test]
807 fn step_kind_label_uses_the_custom_kind_name() {
808 assert_eq!(
809 step_kind_label(&StepKind::Custom("gitlab".to_string())),
810 "gitlab"
811 );
812 }
813}
814
815#[cfg(test)]
816mod output_helper_tests {
817 use super::*;
818 use serde::Deserialize;
819 use serde_json::json;
820
821 fn output(value: Value) -> StepOutput {
822 StepOutput {
823 output: value,
824 duration_ms: 1,
825 cost_usd: Decimal::ZERO,
826 input_tokens: None,
827 output_tokens: None,
828 model: None,
829 debug_messages: None,
830 }
831 }
832
833 #[test]
834 fn shell_helpers_read_shell_fields() {
835 let out = output(json!({"stdout": "hi\n", "stderr": "warn", "exit_code": 0}));
836 assert_eq!(out.exit_code(), Some(0));
837 assert_eq!(out.stdout(), "hi\n");
838 assert_eq!(out.stderr(), "warn");
839 assert!(out.is_success());
840 assert_eq!(out.status(), None);
841 assert_eq!(out.body(), "");
842 }
843
844 #[test]
845 fn shell_non_zero_exit_is_not_success() {
846 let out = output(json!({"stdout": "", "stderr": "", "exit_code": 127}));
847 assert_eq!(out.exit_code(), Some(127));
848 assert!(!out.is_success());
849 }
850
851 #[test]
852 fn http_helpers_read_http_fields() {
853 let out = output(json!({"status": 200, "body": "{\"ok\":true}"}));
854 assert_eq!(out.status(), Some(200));
855 assert_eq!(out.body(), "{\"ok\":true}");
856 assert!(out.is_success());
857 assert_eq!(out.exit_code(), None);
858 assert_eq!(out.stdout(), "");
859 }
860
861 #[test]
862 fn http_error_status_is_not_success() {
863 assert!(!output(json!({"status": 500, "body": ""})).is_success());
864 assert!(!output(json!({"status": 199, "body": ""})).is_success());
865 assert!(output(json!({"status": 299, "body": ""})).is_success());
866 }
867
868 #[test]
869 fn status_out_of_u16_range_is_none() {
870 assert_eq!(output(json!({"status": 70000})).status(), None);
871 assert_eq!(output(json!({"status": "200"})).status(), None);
872 }
873
874 #[test]
875 fn agent_output_without_markers_is_not_success() {
876 let out = output(json!({"summary": "fine"}));
877 assert!(!out.is_success());
878 assert_eq!(out.exit_code(), None);
879 assert_eq!(out.stdout(), "");
880 assert_eq!(out.body(), "");
881 }
882
883 #[test]
884 fn json_deserializes_structured_output() {
885 #[derive(Deserialize, Debug, PartialEq)]
886 struct Review {
887 score: u8,
888 summary: String,
889 }
890 let out = output(json!({"score": 9, "summary": "good"}));
891 let review: Review = out.json().expect("matches schema");
892 assert_eq!(
893 review,
894 Review {
895 score: 9,
896 summary: "good".to_string()
897 }
898 );
899 }
900
901 #[test]
902 fn json_reports_mismatch_as_serialization_error() {
903 #[derive(Deserialize, Debug)]
904 struct Review {
905 #[allow(dead_code)]
906 score: u8,
907 }
908 let out = output(json!({"score": "nine"}));
909 let err = out.json::<Review>().expect_err("type mismatch");
910 assert!(matches!(err, EngineError::Serialization(_)));
911 }
912
913 #[test]
914 fn helpers_tolerate_non_object_output() {
915 let out = output(json!("plain text"));
916 assert_eq!(out.exit_code(), None);
917 assert_eq!(out.status(), None);
918 assert_eq!(out.stdout(), "");
919 assert!(!out.is_success());
920 }
921}