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}