1use std::collections::BTreeMap;
14use std::fmt;
15use std::sync::Arc;
16
17use rust_decimal::Decimal;
18use serde_json::{Value, json, to_string};
19
20use ironflow_core::error::{AgentError, OperationError};
21use ironflow_core::provider::{AgentConfig, AgentOutput, AgentProvider, InvokeFuture};
22
23use crate::config::{ApprovalConfig, HttpConfig, HumanInputConfig, ShellConfig, StepConfig};
24use crate::error::EngineError;
25use crate::executor::{
26 ApprovalOutcome, HumanInputOutcome, SignalOutcome, StepArtifacts, StepInterceptor, StepOutput,
27};
28
29const MISSING_AGENT_PROVIDER: &str = "TestEngine has no agent provider: call with_mock_agent(...), with_recorded_agent(...) or \
31 with_agent_provider(...)";
32
33#[derive(Debug, Clone, PartialEq, Eq, Default)]
47pub struct MockShellOutput {
48 pub stdout: String,
50 pub stderr: String,
52 pub exit_code: i32,
54}
55
56impl MockShellOutput {
57 pub fn ok(stdout: &str) -> Self {
69 Self {
70 stdout: stdout.to_string(),
71 ..Self::default()
72 }
73 }
74
75 pub fn failed(exit_code: i32, stderr: &str) -> Self {
86 Self {
87 stdout: String::new(),
88 stderr: stderr.to_string(),
89 exit_code,
90 }
91 }
92
93 pub(crate) fn into_step_result(
100 self,
101 exit_code_as_output: bool,
102 ) -> Result<StepOutput, EngineError> {
103 if self.exit_code != 0 && !exit_code_as_output {
104 return Err(EngineError::Operation(OperationError::Shell {
105 exit_code: self.exit_code,
106 stderr: self.stderr,
107 }));
108 }
109
110 Ok(StepOutput {
111 output: json!({
112 "stdout": self.stdout,
113 "stderr": self.stderr,
114 "exit_code": self.exit_code,
115 }),
116 duration_ms: 0,
117 cost_usd: Decimal::ZERO,
118 input_tokens: None,
119 cache_read_input_tokens: None,
120 cache_creation_input_tokens: None,
121 output_tokens: None,
122 model: None,
123 debug_messages: None,
124 artifacts: StepArtifacts::default(),
125 account_id: None,
126 environment_id: None,
127 })
128 }
129}
130
131#[derive(Debug, Clone, PartialEq, Eq)]
150pub struct MockHttpResponse {
151 pub status: u16,
153 pub headers: Vec<(String, String)>,
155 pub body: String,
157}
158
159impl Default for MockHttpResponse {
160 fn default() -> Self {
161 Self {
162 status: 200,
163 headers: Vec::new(),
164 body: String::new(),
165 }
166 }
167}
168
169impl MockHttpResponse {
170 pub fn ok(body: &Value) -> Self {
183 Self::json(200, body)
184 }
185
186 pub fn json(status: u16, body: &Value) -> Self {
198 Self {
199 status,
200 headers: Vec::new(),
201 body: to_string(body).unwrap_or_else(|_| body.to_string()),
203 }
204 }
205
206 pub fn text(status: u16, body: &str) -> Self {
217 Self {
218 status,
219 headers: Vec::new(),
220 body: body.to_string(),
221 }
222 }
223
224 pub fn header(mut self, name: &str, value: &str) -> Self {
235 self.headers.push((name.to_string(), value.to_string()));
236 self
237 }
238
239 pub(crate) fn into_step_output(self) -> StepOutput {
242 let headers: BTreeMap<String, String> = self.headers.into_iter().collect();
243 StepOutput {
244 output: json!({
245 "status": self.status,
246 "headers": headers,
247 "body": self.body,
248 }),
249 duration_ms: 0,
250 cost_usd: Decimal::ZERO,
251 input_tokens: None,
252 cache_read_input_tokens: None,
253 cache_creation_input_tokens: None,
254 output_tokens: None,
255 model: None,
256 debug_messages: None,
257 artifacts: StepArtifacts::default(),
258 account_id: None,
259 environment_id: None,
260 }
261 }
262}
263
264pub type ShellMock =
266 Arc<dyn Fn(&ShellConfig) -> Result<MockShellOutput, OperationError> + Send + Sync>;
267
268pub type HttpMock =
270 Arc<dyn Fn(&HttpConfig) -> Result<MockHttpResponse, OperationError> + Send + Sync>;
271
272pub type HumanInputMock = Arc<dyn Fn(&str, &HumanInputConfig) -> HumanInputOutcome + Send + Sync>;
274
275pub type SignalMock = Arc<dyn Fn(&str, &str, &str) -> SignalOutcome + Send + Sync>;
277
278pub type AgentMock = Arc<dyn Fn(&AgentConfig) -> Result<AgentOutput, AgentError> + Send + Sync>;
280
281#[derive(Clone, Default)]
303pub struct MockInterceptor {
304 shell: Option<ShellMock>,
305 http: Option<HttpMock>,
306 approval: Option<ApprovalOutcome>,
307 human_input: Option<HumanInputMock>,
308 signal: Option<SignalMock>,
309}
310
311impl MockInterceptor {
312 pub fn new() -> Self {
326 Self::default()
327 }
328
329 pub fn shell(
341 mut self,
342 f: impl Fn(&ShellConfig) -> Result<MockShellOutput, OperationError> + Send + Sync + 'static,
343 ) -> Self {
344 self.shell = Some(Arc::new(f));
345 self
346 }
347
348 pub fn http(
361 mut self,
362 f: impl Fn(&HttpConfig) -> Result<MockHttpResponse, OperationError> + Send + Sync + 'static,
363 ) -> Self {
364 self.http = Some(Arc::new(f));
365 self
366 }
367
368 pub fn approval(mut self, outcome: ApprovalOutcome) -> Self {
379 self.approval = Some(outcome);
380 self
381 }
382
383 pub fn human_input(
396 mut self,
397 f: impl Fn(&str, &HumanInputConfig) -> HumanInputOutcome + Send + Sync + 'static,
398 ) -> Self {
399 self.human_input = Some(Arc::new(f));
400 self
401 }
402
403 pub fn signal(
418 mut self,
419 f: impl Fn(&str, &str, &str) -> SignalOutcome + Send + Sync + 'static,
420 ) -> Self {
421 self.signal = Some(Arc::new(f));
422 self
423 }
424}
425
426impl fmt::Debug for MockInterceptor {
427 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
428 f.debug_struct("MockInterceptor")
430 .field("shell", &self.shell.is_some())
431 .field("http", &self.http.is_some())
432 .field("approval", &self.approval)
433 .field("human_input", &self.human_input.is_some())
434 .field("signal", &self.signal.is_some())
435 .finish()
436 }
437}
438
439impl StepInterceptor for MockInterceptor {
440 fn intercept(&self, config: &StepConfig) -> Option<Result<StepOutput, EngineError>> {
441 match config {
442 StepConfig::Shell(cfg) => {
443 let mock = self.shell.as_ref()?;
444 Some(match mock(cfg) {
445 Ok(out) => out.into_step_result(cfg.exit_code_as_output),
446 Err(err) => Err(EngineError::Operation(err)),
447 })
448 }
449 StepConfig::Http(cfg) => {
450 let mock = self.http.as_ref()?;
451 Some(match mock(cfg) {
452 Ok(res) => Ok(res.into_step_output()),
453 Err(err) => Err(EngineError::Operation(err)),
454 })
455 }
456 _ => None,
458 }
459 }
460
461 fn intercept_approval(&self, _name: &str, _config: &ApprovalConfig) -> Option<ApprovalOutcome> {
462 self.approval.clone()
463 }
464
465 fn intercept_human_input(
466 &self,
467 name: &str,
468 config: &HumanInputConfig,
469 _schema: &Value,
470 ) -> Option<HumanInputOutcome> {
471 self.human_input.as_ref().map(|f| f(name, config))
472 }
473
474 fn intercept_signal(
475 &self,
476 name: &str,
477 signal_name: &str,
478 key: &str,
479 _schema: &Value,
480 ) -> Option<SignalOutcome> {
481 self.signal.as_ref().map(|f| f(name, signal_name, key))
482 }
483}
484
485pub struct MockAgentProvider {
501 f: AgentMock,
502}
503
504impl MockAgentProvider {
505 pub fn new(
518 f: impl Fn(&AgentConfig) -> Result<AgentOutput, AgentError> + Send + Sync + 'static,
519 ) -> Self {
520 Self { f: Arc::new(f) }
521 }
522}
523
524impl fmt::Debug for MockAgentProvider {
525 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
526 f.debug_struct("MockAgentProvider").finish_non_exhaustive()
527 }
528}
529
530impl AgentProvider for MockAgentProvider {
531 fn invoke<'a>(&'a self, config: &'a AgentConfig) -> InvokeFuture<'a> {
532 let result = (self.f)(config);
533 Box::pin(async move { result })
534 }
535}
536
537#[derive(Debug, Default, Clone, Copy)]
555pub struct MissingAgentProvider;
556
557impl AgentProvider for MissingAgentProvider {
558 fn invoke<'a>(&'a self, _config: &'a AgentConfig) -> InvokeFuture<'a> {
559 Box::pin(async move {
560 Err(AgentError::ProcessFailed {
561 exit_code: -1,
562 stderr: MISSING_AGENT_PROVIDER.to_string(),
563 })
564 })
565 }
566}
567
568#[cfg(test)]
569mod tests {
570 use super::*;
571
572 use crate::config::AgentStepConfig;
573
574 #[test]
575 fn shell_ok_maps_to_the_real_executor_output_shape() {
576 let output = MockShellOutput::ok("hello\n")
577 .into_step_result(false)
578 .expect("exit code 0 succeeds");
579
580 assert_eq!(output.output["stdout"], "hello\n");
581 assert_eq!(output.output["stderr"], "");
582 assert_eq!(output.output["exit_code"], 0);
583 assert_eq!(output.cost_usd, Decimal::ZERO);
584 }
585
586 #[test]
587 fn shell_default_is_an_empty_success() {
588 let default = MockShellOutput::default();
589 assert_eq!(default.exit_code, 0);
590 assert!(default.stdout.is_empty());
591 assert!(default.stderr.is_empty());
592 }
593
594 #[test]
595 fn shell_non_zero_exit_is_an_operation_error() {
596 let err = MockShellOutput::failed(2, "x")
597 .into_step_result(false)
598 .expect_err("a non-zero exit code fails the step");
599
600 match err {
601 EngineError::Operation(OperationError::Shell { exit_code, stderr }) => {
602 assert_eq!(exit_code, 2);
603 assert_eq!(stderr, "x");
604 }
605 other => panic!("expected a shell operation error, got {other}"),
606 }
607 }
608
609 #[test]
610 fn shell_non_zero_exit_with_option_is_an_output() {
611 let output = MockShellOutput::failed(2, "x")
612 .into_step_result(true)
613 .expect("the option turns a non-zero exit into an output");
614
615 assert_eq!(output.output["exit_code"], 2);
616 assert_eq!(output.output["stderr"], "x");
617 assert!(!output.is_success());
618 }
619
620 #[test]
621 fn intercept_shell_non_zero_with_option_completes() {
622 let interceptor = MockInterceptor::new().shell(|_| Ok(MockShellOutput::failed(2, "x")));
623 let config = StepConfig::Shell(ShellConfig::new("x").exit_code_as_output());
624
625 let output = interceptor
626 .intercept(&config)
627 .expect("the shell mock answers")
628 .expect("the step completes");
629 assert_eq!(output.exit_code(), Some(2));
630
631 let plain = StepConfig::Shell(ShellConfig::new("x"));
632 assert!(
633 interceptor
634 .intercept(&plain)
635 .expect("the shell mock answers")
636 .is_err()
637 );
638 }
639
640 #[test]
641 fn http_json_carries_status_body_and_headers() {
642 let output = MockHttpResponse::json(201, &json!({"id": 7}))
643 .header("location", "/things/7")
644 .into_step_output();
645
646 assert_eq!(output.output["status"], 201);
647 assert_eq!(output.output["body"], r#"{"id":7}"#);
648 assert_eq!(output.output["headers"]["location"], "/things/7");
649 }
650
651 #[test]
652 fn http_default_is_an_empty_200() {
653 let default = MockHttpResponse::default();
654 assert_eq!(default.status, 200);
655 assert!(default.body.is_empty());
656 assert!(default.headers.is_empty());
657 }
658
659 #[test]
660 fn http_non_2xx_is_still_an_output() {
661 let output = MockHttpResponse::text(500, "boom").into_step_output();
662 assert_eq!(output.status(), Some(500));
663 assert_eq!(output.body(), "boom");
664 }
665
666 #[test]
667 fn intercept_declines_agent_steps() {
668 let interceptor = MockInterceptor::new().shell(|_| Ok(MockShellOutput::ok("x")));
669 let config = StepConfig::Agent(AgentStepConfig::new("review this"));
670
671 assert!(interceptor.intercept(&config).is_none());
672 }
673
674 #[test]
675 fn intercept_declines_shell_steps_without_a_shell_mock() {
676 let interceptor = MockInterceptor::new();
677 let config = StepConfig::Shell(ShellConfig::new("echo hi"));
678
679 assert!(interceptor.intercept(&config).is_none());
680 }
681
682 #[test]
683 fn intercept_approval_returns_the_configured_outcome() {
684 let interceptor = MockInterceptor::new().approval(ApprovalOutcome::reject("nope"));
685 let config = ApprovalConfig::new("Approve?");
686
687 assert_eq!(
688 interceptor.intercept_approval("gate", &config),
689 Some(ApprovalOutcome::reject("nope"))
690 );
691 assert_eq!(
692 MockInterceptor::new().intercept_approval("gate", &config),
693 None
694 );
695 }
696
697 #[test]
698 fn intercept_human_input_returns_the_mocked_answer() {
699 let interceptor = MockInterceptor::new().human_input(|name, cfg| {
700 HumanInputOutcome::Provided(json!({"step": name, "message": cfg.message()}))
701 });
702 let config = HumanInputConfig::new("Answer?");
703 let answer = json!({"step": "clarify", "message": "Answer?"});
704
705 assert_eq!(
706 interceptor.intercept_human_input("clarify", &config, &json!({})),
707 Some(HumanInputOutcome::Provided(answer))
708 );
709 }
710
711 #[test]
712 fn intercept_human_input_declines_without_a_mock() {
713 let config = HumanInputConfig::new("Answer?");
714
715 assert_eq!(
716 MockInterceptor::new().intercept_human_input("clarify", &config, &json!({})),
717 None
718 );
719 }
720
721 #[test]
722 fn intercept_signal_returns_the_mocked_outcome() {
723 let interceptor = MockInterceptor::new().signal(|step, name, key| {
724 SignalOutcome::Received(json!({"step": step, "name": name, "key": key}))
725 });
726 let expected = json!({"step": "wait-ci", "name": "ci.done", "key": "abc"});
727
728 assert_eq!(
729 interceptor.intercept_signal("wait-ci", "ci.done", "abc", &json!({})),
730 Some(SignalOutcome::Received(expected))
731 );
732 assert_eq!(
733 MockInterceptor::new().intercept_signal("wait-ci", "ci.done", "abc", &json!({})),
734 None
735 );
736 }
737
738 #[test]
739 fn debug_reports_which_seams_are_mocked() {
740 let interceptor = MockInterceptor::new().http(|_| Ok(MockHttpResponse::default()));
741 let rendered = format!("{interceptor:?}");
742
743 assert!(rendered.contains("shell: false"));
744 assert!(rendered.contains("http: true"));
745 }
746
747 #[tokio::test]
748 async fn missing_agent_provider_names_the_three_constructors() {
749 let config = AgentConfig::new("anything");
750 let err = MissingAgentProvider
751 .invoke(&config)
752 .await
753 .expect_err("no agent backend is configured");
754
755 let message = err.to_string();
756 assert!(message.contains("with_mock_agent"));
757 assert!(message.contains("with_recorded_agent"));
758 assert!(message.contains("with_agent_provider"));
759 }
760
761 #[tokio::test]
762 async fn mock_agent_provider_runs_the_closure() {
763 let provider = MockAgentProvider::new(|cfg| {
764 let echoed = json!({"echoed": cfg.prompt.clone()});
765 Ok(AgentOutput::new(echoed))
766 });
767 let config = AgentConfig::new("say hi");
768
769 let output = provider.invoke(&config).await.expect("the mock succeeded");
770
771 assert_eq!(output.value["echoed"], "say hi");
772 }
773}