Skip to main content

ironflow_engine/testing/
engine.rs

1//! [`TestEngine`] -- the builder that runs a handler against mocked steps.
2
3use std::fmt;
4use std::sync::Arc;
5
6use serde_json::Value;
7use uuid::Uuid;
8
9use ironflow_core::decision::DecisionProvider;
10use ironflow_core::error::{AgentError, OperationError};
11use ironflow_core::provider::{AgentConfig, AgentOutput, AgentProvider};
12use ironflow_core::providers::record_replay::RecordReplayProvider;
13use ironflow_store::error::StoreError;
14use ironflow_store::memory::InMemoryStore;
15use ironflow_store::models::{RunStatus, TriggerKind};
16use ironflow_store::store::{RunStore, Store};
17
18use crate::config::{HttpConfig, HumanInputConfig, ShellConfig};
19use crate::engine::{Engine, WorkflowResult};
20use crate::error::EngineError;
21use crate::executor::{ApprovalOutcome, HumanInputOutcome, StepInterceptor};
22use crate::handler::WorkflowHandler;
23use crate::testing::mocks::{
24    MissingAgentProvider, MockAgentProvider, MockHttpResponse, MockInterceptor, MockShellOutput,
25};
26use crate::testing::result::TestResult;
27
28/// Message of the assert guarding every builder method.
29const CONFIGURE_BEFORE_RUN: &str = "configure the TestEngine before its first run";
30
31/// Runs a [`WorkflowHandler`] against an in-memory store with mocked steps.
32///
33/// See the [module documentation](crate::testing) for what the harness replaces
34/// and what it does not.
35///
36/// # Examples
37///
38/// ```no_run
39/// use ironflow_engine::prelude::*;
40/// use ironflow_engine::testing::{MockShellOutput, TestEngine};
41/// use ironflow_store::models::RunStatus;
42/// use serde_json::json;
43///
44/// # struct Deploy;
45/// # impl WorkflowHandler for Deploy {
46/// #     fn name(&self) -> &str { "deploy" }
47/// #     fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
48/// #         Box::pin(async move { ctx.shell("deploy", ShellConfig::new("./deploy.sh")).await?; Ok(()) })
49/// #     }
50/// # }
51/// # async fn example() -> Result<(), EngineError> {
52/// let result = TestEngine::new()
53///     .with_handler(Deploy)
54///     .with_mock_shell(|_cfg| Ok(MockShellOutput::ok("deployed")))
55///     .run(json!({"env": "prod"}))
56///     .await?;
57///
58/// assert_eq!(result.status(), RunStatus::Completed);
59/// # Ok(())
60/// # }
61/// ```
62pub struct TestEngine {
63    store: Arc<InMemoryStore>,
64    handlers: Vec<Box<dyn WorkflowHandler>>,
65    primary: Option<String>,
66    provider: Option<Arc<dyn AgentProvider>>,
67    decision_provider: Option<Arc<dyn DecisionProvider>>,
68    mocks: MockInterceptor,
69    engine: Option<Engine>,
70}
71
72impl Default for TestEngine {
73    fn default() -> Self {
74        Self::new()
75    }
76}
77
78impl fmt::Debug for TestEngine {
79    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
80        f.debug_struct("TestEngine")
81            .field("primary", &self.primary)
82            .field("mocks", &self.mocks)
83            .field("built", &self.engine.is_some())
84            .finish_non_exhaustive()
85    }
86}
87
88impl TestEngine {
89    /// A harness with no handler and no mock.
90    ///
91    /// # Examples
92    ///
93    /// ```
94    /// use ironflow_engine::testing::TestEngine;
95    ///
96    /// let harness = TestEngine::new();
97    /// assert!(format!("{harness:?}").contains("TestEngine"));
98    /// ```
99    pub fn new() -> Self {
100        Self {
101            store: Arc::new(InMemoryStore::new()),
102            handlers: Vec::new(),
103            primary: None,
104            provider: None,
105            decision_provider: None,
106            mocks: MockInterceptor::new(),
107            engine: None,
108        }
109    }
110
111    /// Register a handler. The first one registered is what
112    /// [`run`](Self::run) executes.
113    ///
114    /// # Panics
115    ///
116    /// Panics when called after the first run: the engine is built once, so a
117    /// later registration would be silently ignored.
118    ///
119    /// # Examples
120    ///
121    /// ```no_run
122    /// use ironflow_engine::prelude::*;
123    /// use ironflow_engine::testing::TestEngine;
124    ///
125    /// # struct Deploy;
126    /// # impl WorkflowHandler for Deploy {
127    /// #     fn name(&self) -> &str { "deploy" }
128    /// #     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
129    /// #         Box::pin(async move { Ok(()) })
130    /// #     }
131    /// # }
132    /// let harness = TestEngine::new().with_handler(Deploy);
133    /// # let _ = harness;
134    /// ```
135    pub fn with_handler(mut self, handler: impl WorkflowHandler + 'static) -> Self {
136        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
137        if self.primary.is_none() {
138            self.primary = Some(handler.name().to_string());
139        }
140        self.handlers.push(Box::new(handler));
141        self
142    }
143
144    /// Answer every shell step with `f` instead of spawning a process.
145    ///
146    /// Returning `Err` reproduces a shell failure the same way a non-zero
147    /// [`MockShellOutput::exit_code`] does.
148    ///
149    /// # Panics
150    ///
151    /// Panics when called after the first run.
152    ///
153    /// # Examples
154    ///
155    /// ```no_run
156    /// use ironflow_engine::testing::{MockShellOutput, TestEngine};
157    ///
158    /// let harness = TestEngine::new().with_mock_shell(|cfg| {
159    ///     if cfg.command.starts_with("git ") {
160    ///         Ok(MockShellOutput::ok("abc1234"))
161    ///     } else {
162    ///         Ok(MockShellOutput::failed(127, "command not found"))
163    ///     }
164    /// });
165    /// # let _ = harness;
166    /// ```
167    pub fn with_mock_shell(
168        mut self,
169        f: impl Fn(&ShellConfig) -> Result<MockShellOutput, OperationError> + Send + Sync + 'static,
170    ) -> Self {
171        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
172        self.mocks = self.mocks.shell(f);
173        self
174    }
175
176    /// Answer every HTTP step with `f` instead of sending a request.
177    ///
178    /// A non-2xx [`MockHttpResponse`] is a normal output, like in production.
179    /// Return `Err(OperationError::Http { status: None, .. })` to simulate a
180    /// transport failure.
181    ///
182    /// # Panics
183    ///
184    /// Panics when called after the first run.
185    ///
186    /// # Examples
187    ///
188    /// ```no_run
189    /// use ironflow_engine::testing::{MockHttpResponse, TestEngine};
190    /// use serde_json::json;
191    ///
192    /// let harness = TestEngine::new()
193    ///     .with_mock_http(|_cfg| Ok(MockHttpResponse::json(201, &json!({"id": 7}))));
194    /// # let _ = harness;
195    /// ```
196    pub fn with_mock_http(
197        mut self,
198        f: impl Fn(&HttpConfig) -> Result<MockHttpResponse, OperationError> + Send + Sync + 'static,
199    ) -> Self {
200        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
201        self.mocks = self.mocks.http(f);
202        self
203    }
204
205    /// Resolve every approval gate with `outcome` instead of suspending.
206    ///
207    /// Without this, a gated handler ends the run in
208    /// [`RunStatus::AwaitingApproval`] and [`resume`](Self::resume) continues it.
209    ///
210    /// # Panics
211    ///
212    /// Panics when called after the first run.
213    ///
214    /// # Examples
215    ///
216    /// ```no_run
217    /// use ironflow_engine::testing::{ApprovalOutcome, TestEngine};
218    ///
219    /// let harness = TestEngine::new().with_mock_approval(ApprovalOutcome::Approved);
220    /// # let _ = harness;
221    /// ```
222    pub fn with_mock_approval(mut self, outcome: ApprovalOutcome) -> Self {
223        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
224        self.mocks = self.mocks.approval(outcome);
225        self
226    }
227
228    /// Answer every human input step with `f` instead of suspending.
229    ///
230    /// `f` receives the step name and its config. Without this, a handler that
231    /// asks for a human input ends the run in [`RunStatus::AwaitingApproval`];
232    /// write the answer on the step through the store, then call
233    /// [`resume`](Self::resume).
234    ///
235    /// # Panics
236    ///
237    /// Panics when called after the first run.
238    ///
239    /// # Examples
240    ///
241    /// ```no_run
242    /// use ironflow_engine::testing::{HumanInputOutcome, TestEngine};
243    /// use serde_json::json;
244    ///
245    /// let harness = TestEngine::new().with_mock_human_input(|_name, _cfg| {
246    ///     HumanInputOutcome::Provided(json!({"answers": ["staging"]}))
247    /// });
248    /// # let _ = harness;
249    /// ```
250    pub fn with_mock_human_input(
251        mut self,
252        f: impl Fn(&str, &HumanInputConfig) -> HumanInputOutcome + Send + Sync + 'static,
253    ) -> Self {
254        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
255        self.mocks = self.mocks.human_input(f);
256        self
257    }
258
259    /// Answer every agent step with `f` instead of invoking a backend.
260    ///
261    /// # Panics
262    ///
263    /// Panics when called after the first run.
264    ///
265    /// # Examples
266    ///
267    /// ```no_run
268    /// use ironflow_core::provider::AgentOutput;
269    /// use ironflow_engine::testing::TestEngine;
270    /// use serde_json::json;
271    ///
272    /// let harness = TestEngine::new()
273    ///     .with_mock_agent(|_cfg| Ok(AgentOutput::new(json!({"score": 9}))));
274    /// # let _ = harness;
275    /// ```
276    pub fn with_mock_agent(
277        mut self,
278        f: impl Fn(&AgentConfig) -> Result<AgentOutput, AgentError> + Send + Sync + 'static,
279    ) -> Self {
280        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
281        self.provider = Some(Arc::new(MockAgentProvider::new(f)));
282        self
283    }
284
285    /// Replay agent steps from fixtures recorded in `fixtures_dir`.
286    ///
287    /// `fixtures_dir` is the **directory**, not a file:
288    /// [`RecordReplayProvider`] keys each fixture by a hash of the
289    /// [`AgentConfig`] and stores it as `<hash>.json` inside it. A missing
290    /// fixture falls back to [`MissingAgentProvider`], so the step fails loudly
291    /// instead of reaching the real Claude CLI.
292    ///
293    /// # Panics
294    ///
295    /// Panics when called after the first run.
296    ///
297    /// # Examples
298    ///
299    /// ```no_run
300    /// use ironflow_engine::testing::TestEngine;
301    ///
302    /// let harness = TestEngine::new().with_recorded_agent("tests/fixtures");
303    /// # let _ = harness;
304    /// ```
305    pub fn with_recorded_agent(mut self, fixtures_dir: &str) -> Self {
306        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
307        self.provider = Some(Arc::new(RecordReplayProvider::replay(
308            MissingAgentProvider,
309            fixtures_dir,
310        )));
311        self
312    }
313
314    /// Use an arbitrary [`AgentProvider`] for agent steps.
315    ///
316    /// The escape hatch for anything the three `with_mock_*` methods do not
317    /// cover, such as recording new fixtures.
318    ///
319    /// # Panics
320    ///
321    /// Panics when called after the first run.
322    ///
323    /// # Examples
324    ///
325    /// ```no_run
326    /// use std::sync::Arc;
327    ///
328    /// use ironflow_core::provider::AgentProvider;
329    /// use ironflow_engine::testing::TestEngine;
330    ///
331    /// # fn example(provider: Arc<dyn AgentProvider>) {
332    /// let harness = TestEngine::new().with_agent_provider(provider);
333    /// # let _ = harness;
334    /// # }
335    /// ```
336    pub fn with_agent_provider(mut self, provider: Arc<dyn AgentProvider>) -> Self {
337        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
338        self.provider = Some(provider);
339        self
340    }
341
342    /// Use a [`DecisionProvider`] for `ctx.decision(...)` steps.
343    ///
344    /// Decision steps are not intercepted: without a provider they fail with
345    /// [`EngineError::NoDecisionProvider`].
346    ///
347    /// # Panics
348    ///
349    /// Panics when called after the first run.
350    ///
351    /// # Examples
352    ///
353    /// ```no_run
354    /// use std::sync::Arc;
355    ///
356    /// use ironflow_core::decision::DecisionProvider;
357    /// use ironflow_engine::testing::TestEngine;
358    ///
359    /// # fn example(provider: Arc<dyn DecisionProvider>) {
360    /// let harness = TestEngine::new().with_decision_provider(provider);
361    /// # let _ = harness;
362    /// # }
363    /// ```
364    pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
365        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
366        self.decision_provider = Some(provider);
367        self
368    }
369
370    /// The store backing this harness, for assertions the accessors do not
371    /// cover (child runs, step dependencies, logs).
372    ///
373    /// # Examples
374    ///
375    /// ```no_run
376    /// use ironflow_engine::error::EngineError;
377    /// use ironflow_engine::testing::TestEngine;
378    /// use ironflow_store::store::RunStore;
379    /// use uuid::Uuid;
380    ///
381    /// # async fn example(harness: &TestEngine, run_id: Uuid) -> Result<(), EngineError> {
382    /// let steps = harness.store().list_steps(run_id).await?;
383    /// assert!(!steps.is_empty());
384    /// # Ok(())
385    /// # }
386    /// ```
387    pub fn store(&self) -> Arc<InMemoryStore> {
388        self.store.clone()
389    }
390
391    /// Build the underlying [`Engine`] on first use.
392    ///
393    /// Handlers are drained into it, which is why every builder method asserts
394    /// that no run has happened yet.
395    fn ensure_engine(&mut self) -> Result<(), EngineError> {
396        if self.engine.is_some() {
397            return Ok(());
398        }
399
400        let store: Arc<dyn Store> = self.store.clone();
401        let provider = self
402            .provider
403            .clone()
404            .unwrap_or_else(|| Arc::new(MissingAgentProvider));
405        let mocks: Arc<dyn StepInterceptor> = Arc::new(self.mocks.clone());
406        let mut engine = Engine::new(store, provider).with_step_interceptor(mocks);
407        if let Some(decision_provider) = self.decision_provider.clone() {
408            engine = engine.with_decision_provider(decision_provider);
409        }
410        for handler in self.handlers.drain(..) {
411            engine.register_boxed(handler)?;
412        }
413
414        self.engine = Some(engine);
415        Ok(())
416    }
417
418    /// Run the first handler registered with [`with_handler`](Self::with_handler).
419    ///
420    /// A handler that fails is not an error: the returned [`TestResult`] then
421    /// carries [`RunStatus::Failed`] and the message in
422    /// [`TestResult::error`].
423    ///
424    /// # Errors
425    ///
426    /// Returns [`EngineError::InvalidWorkflow`] when no handler was registered
427    /// or two handlers share a name, and [`EngineError::Store`] when the
428    /// in-memory store rejects a write.
429    ///
430    /// # Examples
431    ///
432    /// ```no_run
433    /// use ironflow_engine::prelude::*;
434    /// use ironflow_engine::testing::{MockShellOutput, TestEngine};
435    /// use serde_json::json;
436    ///
437    /// # struct Deploy;
438    /// # impl WorkflowHandler for Deploy {
439    /// #     fn name(&self) -> &str { "deploy" }
440    /// #     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
441    /// #         Box::pin(async move { Ok(()) })
442    /// #     }
443    /// # }
444    /// # async fn example() -> Result<(), EngineError> {
445    /// let result = TestEngine::new()
446    ///     .with_handler(Deploy)
447    ///     .with_mock_shell(|_cfg| Ok(MockShellOutput::ok("done")))
448    ///     .run(json!({}))
449    ///     .await?;
450    /// assert!(result.is_completed());
451    /// # Ok(())
452    /// # }
453    /// ```
454    pub async fn run(&mut self, payload: Value) -> Result<TestResult, EngineError> {
455        let name = self.primary.clone().ok_or_else(|| {
456            EngineError::InvalidWorkflow(
457                "TestEngine has no handler: call with_handler(...) first".to_string(),
458            )
459        })?;
460        self.run_workflow(&name, payload).await
461    }
462
463    /// Run a specific registered handler by name.
464    ///
465    /// # Errors
466    ///
467    /// Same as [`run`](Self::run), plus [`EngineError::InvalidWorkflow`] when
468    /// `name` matches no registered handler.
469    ///
470    /// # Examples
471    ///
472    /// ```no_run
473    /// use ironflow_engine::error::EngineError;
474    /// use ironflow_engine::testing::TestEngine;
475    /// use serde_json::json;
476    ///
477    /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
478    /// let result = harness.run_workflow("child", json!({"id": 1})).await?;
479    /// assert!(result.is_completed());
480    /// # Ok(())
481    /// # }
482    /// ```
483    pub async fn run_workflow(
484        &mut self,
485        name: &str,
486        payload: Value,
487    ) -> Result<TestResult, EngineError> {
488        self.ensure_engine()?;
489        let engine = self.engine.as_ref().expect("ensure_engine built it");
490
491        // Enqueue then execute, the way the worker does, so the run id is known
492        // before execution and a failed run can still be read back.
493        let trigger = TriggerKind::Manual;
494        let run = engine.enqueue_handler(name, trigger, payload, 0).await?;
495        self.store
496            .update_run_status(run.id, RunStatus::Running)
497            .await?;
498        let execution = engine.execute_handler_run(run.id).await;
499        self.collect(run.id, execution).await
500    }
501
502    /// Resume a run suspended on an approval gate, the way the API server does.
503    ///
504    /// # Errors
505    ///
506    /// Returns [`EngineError::Store`] when the run does not exist or is not
507    /// resumable, and [`EngineError::InvalidWorkflow`] when its handler is no
508    /// longer registered.
509    ///
510    /// # Examples
511    ///
512    /// ```no_run
513    /// use ironflow_engine::error::EngineError;
514    /// use ironflow_engine::testing::TestEngine;
515    /// use ironflow_store::models::RunStatus;
516    /// use serde_json::json;
517    ///
518    /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
519    /// let suspended = harness.run(json!({})).await?;
520    /// assert_eq!(suspended.status(), RunStatus::AwaitingApproval);
521    ///
522    /// let resumed = harness.resume(suspended.run_id()).await?;
523    /// assert_eq!(resumed.status(), RunStatus::Completed);
524    /// # Ok(())
525    /// # }
526    /// ```
527    pub async fn resume(&mut self, run_id: Uuid) -> Result<TestResult, EngineError> {
528        self.ensure_engine()?;
529        let engine = self.engine.as_ref().expect("ensure_engine built it");
530
531        self.store
532            .update_run_status(run_id, RunStatus::Running)
533            .await?;
534        let execution = engine.resume_run(run_id).await;
535        self.collect(run_id, execution).await
536    }
537
538    /// Read the run and its steps back from the store.
539    ///
540    /// The engine returns `Err` for a failed run, so the steps are always read
541    /// from the store rather than from the execution result.
542    async fn collect(
543        &self,
544        run_id: Uuid,
545        execution: Result<WorkflowResult, EngineError>,
546    ) -> Result<TestResult, EngineError> {
547        let run = self
548            .store
549            .get_run(run_id)
550            .await?
551            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
552        let steps = self.store.list_steps(run_id).await?;
553
554        let (step_results, error) = match execution {
555            Ok(result) => (result.steps, run.error.clone()),
556            Err(err) => (Vec::new(), Some(err.to_string())),
557        };
558
559        Ok(TestResult::new(run, steps, step_results, error))
560    }
561}