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: a non-zero exit code is an error
148 /// unless the step set `exit_code_as_output`.
149 ///
150 /// # Panics
151 ///
152 /// Panics when called after the first run.
153 ///
154 /// # Examples
155 ///
156 /// ```no_run
157 /// use ironflow_engine::testing::{MockShellOutput, TestEngine};
158 ///
159 /// let harness = TestEngine::new().with_mock_shell(|cfg| {
160 /// if cfg.command.starts_with("git ") {
161 /// Ok(MockShellOutput::ok("abc1234"))
162 /// } else {
163 /// Ok(MockShellOutput::failed(127, "command not found"))
164 /// }
165 /// });
166 /// # let _ = harness;
167 /// ```
168 pub fn with_mock_shell(
169 mut self,
170 f: impl Fn(&ShellConfig) -> Result<MockShellOutput, OperationError> + Send + Sync + 'static,
171 ) -> Self {
172 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
173 self.mocks = self.mocks.shell(f);
174 self
175 }
176
177 /// Answer every HTTP step with `f` instead of sending a request.
178 ///
179 /// A non-2xx [`MockHttpResponse`] is a normal output, like in production.
180 /// Return `Err(OperationError::Http { status: None, .. })` to simulate a
181 /// transport failure.
182 ///
183 /// # Panics
184 ///
185 /// Panics when called after the first run.
186 ///
187 /// # Examples
188 ///
189 /// ```no_run
190 /// use ironflow_engine::testing::{MockHttpResponse, TestEngine};
191 /// use serde_json::json;
192 ///
193 /// let harness = TestEngine::new()
194 /// .with_mock_http(|_cfg| Ok(MockHttpResponse::json(201, &json!({"id": 7}))));
195 /// # let _ = harness;
196 /// ```
197 pub fn with_mock_http(
198 mut self,
199 f: impl Fn(&HttpConfig) -> Result<MockHttpResponse, OperationError> + Send + Sync + 'static,
200 ) -> Self {
201 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
202 self.mocks = self.mocks.http(f);
203 self
204 }
205
206 /// Resolve every approval gate with `outcome` instead of suspending.
207 ///
208 /// Without this, a gated handler ends the run in
209 /// [`RunStatus::AwaitingApproval`] and [`resume`](Self::resume) continues it.
210 ///
211 /// # Panics
212 ///
213 /// Panics when called after the first run.
214 ///
215 /// # Examples
216 ///
217 /// ```no_run
218 /// use ironflow_engine::testing::{ApprovalOutcome, TestEngine};
219 ///
220 /// let harness = TestEngine::new().with_mock_approval(ApprovalOutcome::Approved);
221 /// # let _ = harness;
222 /// ```
223 pub fn with_mock_approval(mut self, outcome: ApprovalOutcome) -> Self {
224 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
225 self.mocks = self.mocks.approval(outcome);
226 self
227 }
228
229 /// Answer every human input step with `f` instead of suspending.
230 ///
231 /// `f` receives the step name and its config. Without this, a handler that
232 /// asks for a human input ends the run in [`RunStatus::AwaitingApproval`];
233 /// write the answer on the step through the store, then call
234 /// [`resume`](Self::resume).
235 ///
236 /// # Panics
237 ///
238 /// Panics when called after the first run.
239 ///
240 /// # Examples
241 ///
242 /// ```no_run
243 /// use ironflow_engine::testing::{HumanInputOutcome, TestEngine};
244 /// use serde_json::json;
245 ///
246 /// let harness = TestEngine::new().with_mock_human_input(|_name, _cfg| {
247 /// HumanInputOutcome::Provided(json!({"answers": ["staging"]}))
248 /// });
249 /// # let _ = harness;
250 /// ```
251 pub fn with_mock_human_input(
252 mut self,
253 f: impl Fn(&str, &HumanInputConfig) -> HumanInputOutcome + Send + Sync + 'static,
254 ) -> Self {
255 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
256 self.mocks = self.mocks.human_input(f);
257 self
258 }
259
260 /// Answer every agent step with `f` instead of invoking a backend.
261 ///
262 /// # Panics
263 ///
264 /// Panics when called after the first run.
265 ///
266 /// # Examples
267 ///
268 /// ```no_run
269 /// use ironflow_core::provider::AgentOutput;
270 /// use ironflow_engine::testing::TestEngine;
271 /// use serde_json::json;
272 ///
273 /// let harness = TestEngine::new()
274 /// .with_mock_agent(|_cfg| Ok(AgentOutput::new(json!({"score": 9}))));
275 /// # let _ = harness;
276 /// ```
277 pub fn with_mock_agent(
278 mut self,
279 f: impl Fn(&AgentConfig) -> Result<AgentOutput, AgentError> + Send + Sync + 'static,
280 ) -> Self {
281 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
282 self.provider = Some(Arc::new(MockAgentProvider::new(f)));
283 self
284 }
285
286 /// Replay agent steps from fixtures recorded in `fixtures_dir`.
287 ///
288 /// `fixtures_dir` is the **directory**, not a file:
289 /// [`RecordReplayProvider`] keys each fixture by a hash of the
290 /// [`AgentConfig`] and stores it as `<hash>.json` inside it. A missing
291 /// fixture falls back to [`MissingAgentProvider`], so the step fails loudly
292 /// instead of reaching the real Claude CLI.
293 ///
294 /// # Panics
295 ///
296 /// Panics when called after the first run.
297 ///
298 /// # Examples
299 ///
300 /// ```no_run
301 /// use ironflow_engine::testing::TestEngine;
302 ///
303 /// let harness = TestEngine::new().with_recorded_agent("tests/fixtures");
304 /// # let _ = harness;
305 /// ```
306 pub fn with_recorded_agent(mut self, fixtures_dir: &str) -> Self {
307 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
308 self.provider = Some(Arc::new(RecordReplayProvider::replay(
309 MissingAgentProvider,
310 fixtures_dir,
311 )));
312 self
313 }
314
315 /// Use an arbitrary [`AgentProvider`] for agent steps.
316 ///
317 /// The escape hatch for anything the three `with_mock_*` methods do not
318 /// cover, such as recording new fixtures.
319 ///
320 /// # Panics
321 ///
322 /// Panics when called after the first run.
323 ///
324 /// # Examples
325 ///
326 /// ```no_run
327 /// use std::sync::Arc;
328 ///
329 /// use ironflow_core::provider::AgentProvider;
330 /// use ironflow_engine::testing::TestEngine;
331 ///
332 /// # fn example(provider: Arc<dyn AgentProvider>) {
333 /// let harness = TestEngine::new().with_agent_provider(provider);
334 /// # let _ = harness;
335 /// # }
336 /// ```
337 pub fn with_agent_provider(mut self, provider: Arc<dyn AgentProvider>) -> Self {
338 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
339 self.provider = Some(provider);
340 self
341 }
342
343 /// Use a [`DecisionProvider`] for `ctx.decision(...)` steps.
344 ///
345 /// Decision steps are not intercepted: without a provider they fail with
346 /// [`EngineError::NoDecisionProvider`].
347 ///
348 /// # Panics
349 ///
350 /// Panics when called after the first run.
351 ///
352 /// # Examples
353 ///
354 /// ```no_run
355 /// use std::sync::Arc;
356 ///
357 /// use ironflow_core::decision::DecisionProvider;
358 /// use ironflow_engine::testing::TestEngine;
359 ///
360 /// # fn example(provider: Arc<dyn DecisionProvider>) {
361 /// let harness = TestEngine::new().with_decision_provider(provider);
362 /// # let _ = harness;
363 /// # }
364 /// ```
365 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
366 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
367 self.decision_provider = Some(provider);
368 self
369 }
370
371 /// The store backing this harness, for assertions the accessors do not
372 /// cover (child runs, step dependencies, logs).
373 ///
374 /// # Examples
375 ///
376 /// ```no_run
377 /// use ironflow_engine::error::EngineError;
378 /// use ironflow_engine::testing::TestEngine;
379 /// use ironflow_store::store::RunStore;
380 /// use uuid::Uuid;
381 ///
382 /// # async fn example(harness: &TestEngine, run_id: Uuid) -> Result<(), EngineError> {
383 /// let steps = harness.store().list_steps(run_id).await?;
384 /// assert!(!steps.is_empty());
385 /// # Ok(())
386 /// # }
387 /// ```
388 pub fn store(&self) -> Arc<InMemoryStore> {
389 self.store.clone()
390 }
391
392 /// Build the underlying [`Engine`] on first use.
393 ///
394 /// Handlers are drained into it, which is why every builder method asserts
395 /// that no run has happened yet.
396 fn ensure_engine(&mut self) -> Result<(), EngineError> {
397 if self.engine.is_some() {
398 return Ok(());
399 }
400
401 let store: Arc<dyn Store> = self.store.clone();
402 let provider = self
403 .provider
404 .clone()
405 .unwrap_or_else(|| Arc::new(MissingAgentProvider));
406 let mocks: Arc<dyn StepInterceptor> = Arc::new(self.mocks.clone());
407 let mut engine = Engine::new(store, provider).with_step_interceptor(mocks);
408 if let Some(decision_provider) = self.decision_provider.clone() {
409 engine = engine.with_decision_provider(decision_provider);
410 }
411 for handler in self.handlers.drain(..) {
412 engine.register_boxed(handler)?;
413 }
414
415 self.engine = Some(engine);
416 Ok(())
417 }
418
419 /// Run the first handler registered with [`with_handler`](Self::with_handler).
420 ///
421 /// A handler that fails is not an error: the returned [`TestResult`] then
422 /// carries [`RunStatus::Failed`] and the message in
423 /// [`TestResult::error`].
424 ///
425 /// # Errors
426 ///
427 /// Returns [`EngineError::InvalidWorkflow`] when no handler was registered
428 /// or two handlers share a name, and [`EngineError::Store`] when the
429 /// in-memory store rejects a write.
430 ///
431 /// # Examples
432 ///
433 /// ```no_run
434 /// use ironflow_engine::prelude::*;
435 /// use ironflow_engine::testing::{MockShellOutput, TestEngine};
436 /// use serde_json::json;
437 ///
438 /// # struct Deploy;
439 /// # impl WorkflowHandler for Deploy {
440 /// # fn name(&self) -> &str { "deploy" }
441 /// # fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
442 /// # Box::pin(async move { Ok(()) })
443 /// # }
444 /// # }
445 /// # async fn example() -> Result<(), EngineError> {
446 /// let result = TestEngine::new()
447 /// .with_handler(Deploy)
448 /// .with_mock_shell(|_cfg| Ok(MockShellOutput::ok("done")))
449 /// .run(json!({}))
450 /// .await?;
451 /// assert!(result.is_completed());
452 /// # Ok(())
453 /// # }
454 /// ```
455 pub async fn run(&mut self, payload: Value) -> Result<TestResult, EngineError> {
456 let name = self.primary.clone().ok_or_else(|| {
457 EngineError::InvalidWorkflow(
458 "TestEngine has no handler: call with_handler(...) first".to_string(),
459 )
460 })?;
461 self.run_workflow(&name, payload).await
462 }
463
464 /// Run a specific registered handler by name.
465 ///
466 /// # Errors
467 ///
468 /// Same as [`run`](Self::run), plus [`EngineError::InvalidWorkflow`] when
469 /// `name` matches no registered handler.
470 ///
471 /// # Examples
472 ///
473 /// ```no_run
474 /// use ironflow_engine::error::EngineError;
475 /// use ironflow_engine::testing::TestEngine;
476 /// use serde_json::json;
477 ///
478 /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
479 /// let result = harness.run_workflow("child", json!({"id": 1})).await?;
480 /// assert!(result.is_completed());
481 /// # Ok(())
482 /// # }
483 /// ```
484 pub async fn run_workflow(
485 &mut self,
486 name: &str,
487 payload: Value,
488 ) -> Result<TestResult, EngineError> {
489 self.ensure_engine()?;
490 let engine = self.engine.as_ref().expect("ensure_engine built it");
491
492 // Enqueue then execute, the way the worker does, so the run id is known
493 // before execution and a failed run can still be read back.
494 let trigger = TriggerKind::Manual;
495 let run = engine.enqueue_handler(name, trigger, payload, 0).await?;
496 self.store
497 .update_run_status(run.id, RunStatus::Running)
498 .await?;
499 let execution = engine.execute_handler_run(run.id).await;
500 self.collect(run.id, execution).await
501 }
502
503 /// Resume a run suspended on an approval gate, the way the API server does.
504 ///
505 /// # Errors
506 ///
507 /// Returns [`EngineError::Store`] when the run does not exist or is not
508 /// resumable, and [`EngineError::InvalidWorkflow`] when its handler is no
509 /// longer registered.
510 ///
511 /// # Examples
512 ///
513 /// ```no_run
514 /// use ironflow_engine::error::EngineError;
515 /// use ironflow_engine::testing::TestEngine;
516 /// use ironflow_store::models::RunStatus;
517 /// use serde_json::json;
518 ///
519 /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
520 /// let suspended = harness.run(json!({})).await?;
521 /// assert_eq!(suspended.status(), RunStatus::AwaitingApproval);
522 ///
523 /// let resumed = harness.resume(suspended.run_id()).await?;
524 /// assert_eq!(resumed.status(), RunStatus::Completed);
525 /// # Ok(())
526 /// # }
527 /// ```
528 pub async fn resume(&mut self, run_id: Uuid) -> Result<TestResult, EngineError> {
529 self.ensure_engine()?;
530 let engine = self.engine.as_ref().expect("ensure_engine built it");
531
532 self.store
533 .update_run_status(run_id, RunStatus::Running)
534 .await?;
535 let execution = engine.resume_run(run_id).await;
536 self.collect(run_id, execution).await
537 }
538
539 /// Read the run and its steps back from the store.
540 ///
541 /// The engine returns `Err` for a failed run, so the steps are always read
542 /// from the store rather than from the execution result.
543 async fn collect(
544 &self,
545 run_id: Uuid,
546 execution: Result<WorkflowResult, EngineError>,
547 ) -> Result<TestResult, EngineError> {
548 let run = self
549 .store
550 .get_run(run_id)
551 .await?
552 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
553 let steps = self.store.list_steps(run_id).await?;
554
555 let (step_results, error) = match execution {
556 Ok(result) => (result.steps, run.error.clone()),
557 Err(err) => (Vec::new(), Some(err.to_string())),
558 };
559
560 Ok(TestResult::new(run, steps, step_results, error))
561 }
562}