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, ShellConfig};
19use crate::engine::{Engine, WorkflowResult};
20use crate::error::EngineError;
21use crate::executor::{ApprovalOutcome, 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 agent step with `f` instead of invoking a backend.
229 ///
230 /// # Panics
231 ///
232 /// Panics when called after the first run.
233 ///
234 /// # Examples
235 ///
236 /// ```no_run
237 /// use ironflow_core::provider::AgentOutput;
238 /// use ironflow_engine::testing::TestEngine;
239 /// use serde_json::json;
240 ///
241 /// let harness = TestEngine::new()
242 /// .with_mock_agent(|_cfg| Ok(AgentOutput::new(json!({"score": 9}))));
243 /// # let _ = harness;
244 /// ```
245 pub fn with_mock_agent(
246 mut self,
247 f: impl Fn(&AgentConfig) -> Result<AgentOutput, AgentError> + Send + Sync + 'static,
248 ) -> Self {
249 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
250 self.provider = Some(Arc::new(MockAgentProvider::new(f)));
251 self
252 }
253
254 /// Replay agent steps from fixtures recorded in `fixtures_dir`.
255 ///
256 /// `fixtures_dir` is the **directory**, not a file:
257 /// [`RecordReplayProvider`] keys each fixture by a hash of the
258 /// [`AgentConfig`] and stores it as `<hash>.json` inside it. A missing
259 /// fixture falls back to [`MissingAgentProvider`], so the step fails loudly
260 /// instead of reaching the real Claude CLI.
261 ///
262 /// # Panics
263 ///
264 /// Panics when called after the first run.
265 ///
266 /// # Examples
267 ///
268 /// ```no_run
269 /// use ironflow_engine::testing::TestEngine;
270 ///
271 /// let harness = TestEngine::new().with_recorded_agent("tests/fixtures");
272 /// # let _ = harness;
273 /// ```
274 pub fn with_recorded_agent(mut self, fixtures_dir: &str) -> Self {
275 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
276 self.provider = Some(Arc::new(RecordReplayProvider::replay(
277 MissingAgentProvider,
278 fixtures_dir,
279 )));
280 self
281 }
282
283 /// Use an arbitrary [`AgentProvider`] for agent steps.
284 ///
285 /// The escape hatch for anything the three `with_mock_*` methods do not
286 /// cover, such as recording new fixtures.
287 ///
288 /// # Panics
289 ///
290 /// Panics when called after the first run.
291 ///
292 /// # Examples
293 ///
294 /// ```no_run
295 /// use std::sync::Arc;
296 ///
297 /// use ironflow_core::provider::AgentProvider;
298 /// use ironflow_engine::testing::TestEngine;
299 ///
300 /// # fn example(provider: Arc<dyn AgentProvider>) {
301 /// let harness = TestEngine::new().with_agent_provider(provider);
302 /// # let _ = harness;
303 /// # }
304 /// ```
305 pub fn with_agent_provider(mut self, provider: Arc<dyn AgentProvider>) -> Self {
306 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
307 self.provider = Some(provider);
308 self
309 }
310
311 /// Use a [`DecisionProvider`] for `ctx.decision(...)` steps.
312 ///
313 /// Decision steps are not intercepted: without a provider they fail with
314 /// [`EngineError::NoDecisionProvider`].
315 ///
316 /// # Panics
317 ///
318 /// Panics when called after the first run.
319 ///
320 /// # Examples
321 ///
322 /// ```no_run
323 /// use std::sync::Arc;
324 ///
325 /// use ironflow_core::decision::DecisionProvider;
326 /// use ironflow_engine::testing::TestEngine;
327 ///
328 /// # fn example(provider: Arc<dyn DecisionProvider>) {
329 /// let harness = TestEngine::new().with_decision_provider(provider);
330 /// # let _ = harness;
331 /// # }
332 /// ```
333 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
334 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
335 self.decision_provider = Some(provider);
336 self
337 }
338
339 /// The store backing this harness, for assertions the accessors do not
340 /// cover (child runs, step dependencies, logs).
341 ///
342 /// # Examples
343 ///
344 /// ```no_run
345 /// use ironflow_engine::error::EngineError;
346 /// use ironflow_engine::testing::TestEngine;
347 /// use ironflow_store::store::RunStore;
348 /// use uuid::Uuid;
349 ///
350 /// # async fn example(harness: &TestEngine, run_id: Uuid) -> Result<(), EngineError> {
351 /// let steps = harness.store().list_steps(run_id).await?;
352 /// assert!(!steps.is_empty());
353 /// # Ok(())
354 /// # }
355 /// ```
356 pub fn store(&self) -> Arc<InMemoryStore> {
357 self.store.clone()
358 }
359
360 /// Build the underlying [`Engine`] on first use.
361 ///
362 /// Handlers are drained into it, which is why every builder method asserts
363 /// that no run has happened yet.
364 fn ensure_engine(&mut self) -> Result<(), EngineError> {
365 if self.engine.is_some() {
366 return Ok(());
367 }
368
369 let store: Arc<dyn Store> = self.store.clone();
370 let provider = self
371 .provider
372 .clone()
373 .unwrap_or_else(|| Arc::new(MissingAgentProvider));
374 let mocks: Arc<dyn StepInterceptor> = Arc::new(self.mocks.clone());
375 let mut engine = Engine::new(store, provider).with_step_interceptor(mocks);
376 if let Some(decision_provider) = self.decision_provider.clone() {
377 engine = engine.with_decision_provider(decision_provider);
378 }
379 for handler in self.handlers.drain(..) {
380 engine.register_boxed(handler)?;
381 }
382
383 self.engine = Some(engine);
384 Ok(())
385 }
386
387 /// Run the first handler registered with [`with_handler`](Self::with_handler).
388 ///
389 /// A handler that fails is not an error: the returned [`TestResult`] then
390 /// carries [`RunStatus::Failed`] and the message in
391 /// [`TestResult::error`].
392 ///
393 /// # Errors
394 ///
395 /// Returns [`EngineError::InvalidWorkflow`] when no handler was registered
396 /// or two handlers share a name, and [`EngineError::Store`] when the
397 /// in-memory store rejects a write.
398 ///
399 /// # Examples
400 ///
401 /// ```no_run
402 /// use ironflow_engine::prelude::*;
403 /// use ironflow_engine::testing::{MockShellOutput, TestEngine};
404 /// use serde_json::json;
405 ///
406 /// # struct Deploy;
407 /// # impl WorkflowHandler for Deploy {
408 /// # fn name(&self) -> &str { "deploy" }
409 /// # fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
410 /// # Box::pin(async move { Ok(()) })
411 /// # }
412 /// # }
413 /// # async fn example() -> Result<(), EngineError> {
414 /// let result = TestEngine::new()
415 /// .with_handler(Deploy)
416 /// .with_mock_shell(|_cfg| Ok(MockShellOutput::ok("done")))
417 /// .run(json!({}))
418 /// .await?;
419 /// assert!(result.is_completed());
420 /// # Ok(())
421 /// # }
422 /// ```
423 pub async fn run(&mut self, payload: Value) -> Result<TestResult, EngineError> {
424 let name = self.primary.clone().ok_or_else(|| {
425 EngineError::InvalidWorkflow(
426 "TestEngine has no handler: call with_handler(...) first".to_string(),
427 )
428 })?;
429 self.run_workflow(&name, payload).await
430 }
431
432 /// Run a specific registered handler by name.
433 ///
434 /// # Errors
435 ///
436 /// Same as [`run`](Self::run), plus [`EngineError::InvalidWorkflow`] when
437 /// `name` matches no registered handler.
438 ///
439 /// # Examples
440 ///
441 /// ```no_run
442 /// use ironflow_engine::error::EngineError;
443 /// use ironflow_engine::testing::TestEngine;
444 /// use serde_json::json;
445 ///
446 /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
447 /// let result = harness.run_workflow("child", json!({"id": 1})).await?;
448 /// assert!(result.is_completed());
449 /// # Ok(())
450 /// # }
451 /// ```
452 pub async fn run_workflow(
453 &mut self,
454 name: &str,
455 payload: Value,
456 ) -> Result<TestResult, EngineError> {
457 self.ensure_engine()?;
458 let engine = self.engine.as_ref().expect("ensure_engine built it");
459
460 // Enqueue then execute, the way the worker does, so the run id is known
461 // before execution and a failed run can still be read back.
462 let trigger = TriggerKind::Manual;
463 let run = engine.enqueue_handler(name, trigger, payload, 0).await?;
464 self.store
465 .update_run_status(run.id, RunStatus::Running)
466 .await?;
467 let execution = engine.execute_handler_run(run.id).await;
468 self.collect(run.id, execution).await
469 }
470
471 /// Resume a run suspended on an approval gate, the way the API server does.
472 ///
473 /// # Errors
474 ///
475 /// Returns [`EngineError::Store`] when the run does not exist or is not
476 /// resumable, and [`EngineError::InvalidWorkflow`] when its handler is no
477 /// longer registered.
478 ///
479 /// # Examples
480 ///
481 /// ```no_run
482 /// use ironflow_engine::error::EngineError;
483 /// use ironflow_engine::testing::TestEngine;
484 /// use ironflow_store::models::RunStatus;
485 /// use serde_json::json;
486 ///
487 /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
488 /// let suspended = harness.run(json!({})).await?;
489 /// assert_eq!(suspended.status(), RunStatus::AwaitingApproval);
490 ///
491 /// let resumed = harness.resume(suspended.run_id()).await?;
492 /// assert_eq!(resumed.status(), RunStatus::Completed);
493 /// # Ok(())
494 /// # }
495 /// ```
496 pub async fn resume(&mut self, run_id: Uuid) -> Result<TestResult, EngineError> {
497 self.ensure_engine()?;
498 let engine = self.engine.as_ref().expect("ensure_engine built it");
499
500 self.store
501 .update_run_status(run_id, RunStatus::Running)
502 .await?;
503 let execution = engine.resume_run(run_id).await;
504 self.collect(run_id, execution).await
505 }
506
507 /// Read the run and its steps back from the store.
508 ///
509 /// The engine returns `Err` for a failed run, so the steps are always read
510 /// from the store rather than from the execution result.
511 async fn collect(
512 &self,
513 run_id: Uuid,
514 execution: Result<WorkflowResult, EngineError>,
515 ) -> Result<TestResult, EngineError> {
516 let run = self
517 .store
518 .get_run(run_id)
519 .await?
520 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
521 let steps = self.store.list_steps(run_id).await?;
522
523 let (step_results, error) = match execution {
524 Ok(result) => (result.steps, run.error.clone()),
525 Err(err) => (Vec::new(), Some(err.to_string())),
526 };
527
528 Ok(TestResult::new(run, steps, step_results, error))
529 }
530}