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
6#[cfg(feature = "secret-store")]
7use rand::random;
8use serde_json::Value;
9use uuid::Uuid;
10
11use ironflow_core::decision::DecisionProvider;
12use ironflow_core::error::{AgentError, OperationError};
13use ironflow_core::provider::{AgentConfig, AgentOutput, AgentProvider};
14use ironflow_core::providers::record_replay::RecordReplayProvider;
15#[cfg(feature = "secret-store")]
16use ironflow_store::crypto::MasterKey;
17use ironflow_store::error::StoreError;
18use ironflow_store::memory::InMemoryStore;
19use ironflow_store::models::{RunStatus, TriggerKind};
20use ironflow_store::store::{RunStore, Store};
21#[cfg(feature = "secret-store")]
22use ironflow_store::workflow_secrets::ScopedSecretStore;
23
24use crate::config::{HttpConfig, HumanInputConfig, ShellConfig};
25use crate::engine::{Engine, WorkflowResult, chain_root};
26use crate::error::EngineError;
27use crate::executor::{ApprovalOutcome, HumanInputOutcome, SignalOutcome, StepInterceptor};
28use crate::handler::WorkflowHandler;
29use crate::testing::mocks::{
30 MissingAgentProvider, MockAgentProvider, MockHttpResponse, MockInterceptor, MockShellOutput,
31};
32use crate::testing::result::TestResult;
33
34/// Message of the assert guarding every builder method.
35const CONFIGURE_BEFORE_RUN: &str = "configure the TestEngine before its first run";
36
37/// Runs a [`WorkflowHandler`] against an in-memory store with mocked steps.
38///
39/// See the [module documentation](crate::testing) for what the harness replaces
40/// and what it does not.
41///
42/// # Examples
43///
44/// ```no_run
45/// use ironflow_engine::prelude::*;
46/// use ironflow_engine::testing::{MockShellOutput, TestEngine};
47/// use ironflow_store::models::RunStatus;
48/// use serde_json::json;
49///
50/// # struct Deploy;
51/// # impl WorkflowHandler for Deploy {
52/// # fn name(&self) -> &str { "deploy" }
53/// # fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
54/// # Box::pin(async move { ctx.shell("deploy", ShellConfig::new("./deploy.sh")).await?; Ok(()) })
55/// # }
56/// # }
57/// # async fn example() -> Result<(), EngineError> {
58/// let result = TestEngine::new()
59/// .with_handler(Deploy)
60/// .with_mock_shell(|_cfg| Ok(MockShellOutput::ok("deployed")))
61/// .run(json!({"env": "prod"}))
62/// .await?;
63///
64/// assert_eq!(result.status(), RunStatus::Completed);
65/// # Ok(())
66/// # }
67/// ```
68pub struct TestEngine {
69 store: Arc<InMemoryStore>,
70 handlers: Vec<Box<dyn WorkflowHandler>>,
71 primary: Option<String>,
72 provider: Option<Arc<dyn AgentProvider>>,
73 decision_provider: Option<Arc<dyn DecisionProvider>>,
74 mocks: MockInterceptor,
75 engine: Option<Engine>,
76 #[cfg(feature = "secret-store")]
77 secrets: Vec<(String, String)>,
78}
79
80impl Default for TestEngine {
81 fn default() -> Self {
82 Self::new()
83 }
84}
85
86impl fmt::Debug for TestEngine {
87 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
88 f.debug_struct("TestEngine")
89 .field("primary", &self.primary)
90 .field("mocks", &self.mocks)
91 .field("built", &self.engine.is_some())
92 .finish_non_exhaustive()
93 }
94}
95
96impl TestEngine {
97 /// A harness with no handler and no mock.
98 ///
99 /// With the `secret-store` feature the store gets a freshly generated
100 /// master key, so reading a secret that was never set resolves to `None`
101 /// instead of failing with "no master key configured". Seed values with
102 /// `with_secret`.
103 ///
104 /// # Examples
105 ///
106 /// ```
107 /// use ironflow_engine::testing::TestEngine;
108 ///
109 /// let harness = TestEngine::new();
110 /// assert!(format!("{harness:?}").contains("TestEngine"));
111 /// ```
112 pub fn new() -> Self {
113 #[cfg_attr(not(feature = "secret-store"), allow(unused_mut))]
114 let mut store = InMemoryStore::new();
115 #[cfg(feature = "secret-store")]
116 store.set_master_key(
117 MasterKey::from_bytes(&random::<[u8; 32]>())
118 .expect("32 bytes is a valid master key length"),
119 );
120 Self {
121 store: Arc::new(store),
122 handlers: Vec::new(),
123 primary: None,
124 provider: None,
125 decision_provider: None,
126 mocks: MockInterceptor::new(),
127 engine: None,
128 #[cfg(feature = "secret-store")]
129 secrets: Vec::new(),
130 }
131 }
132
133 /// Register a handler. The first one registered is what
134 /// [`run`](Self::run) executes.
135 ///
136 /// # Panics
137 ///
138 /// Panics when called after the first run: the engine is built once, so a
139 /// later registration would be silently ignored.
140 ///
141 /// # Examples
142 ///
143 /// ```no_run
144 /// use ironflow_engine::prelude::*;
145 /// use ironflow_engine::testing::TestEngine;
146 ///
147 /// # struct Deploy;
148 /// # impl WorkflowHandler for Deploy {
149 /// # fn name(&self) -> &str { "deploy" }
150 /// # fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
151 /// # Box::pin(async move { Ok(()) })
152 /// # }
153 /// # }
154 /// let harness = TestEngine::new().with_handler(Deploy);
155 /// # let _ = harness;
156 /// ```
157 pub fn with_handler(mut self, handler: impl WorkflowHandler + 'static) -> Self {
158 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
159 if self.primary.is_none() {
160 self.primary = Some(handler.name().to_string());
161 }
162 self.handlers.push(Box::new(handler));
163 self
164 }
165
166 /// Answer every shell step with `f` instead of spawning a process.
167 ///
168 /// Returning `Err` reproduces a shell failure the same way a non-zero
169 /// [`MockShellOutput::exit_code`] does: a non-zero exit code is an error
170 /// unless the step set `exit_code_as_output`.
171 ///
172 /// # Panics
173 ///
174 /// Panics when called after the first run.
175 ///
176 /// # Examples
177 ///
178 /// ```no_run
179 /// use ironflow_engine::testing::{MockShellOutput, TestEngine};
180 ///
181 /// let harness = TestEngine::new().with_mock_shell(|cfg| {
182 /// if cfg.command.starts_with("git ") {
183 /// Ok(MockShellOutput::ok("abc1234"))
184 /// } else {
185 /// Ok(MockShellOutput::failed(127, "command not found"))
186 /// }
187 /// });
188 /// # let _ = harness;
189 /// ```
190 pub fn with_mock_shell(
191 mut self,
192 f: impl Fn(&ShellConfig) -> Result<MockShellOutput, OperationError> + Send + Sync + 'static,
193 ) -> Self {
194 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
195 self.mocks = self.mocks.shell(f);
196 self
197 }
198
199 /// Answer every HTTP step with `f` instead of sending a request.
200 ///
201 /// A non-2xx [`MockHttpResponse`] is a normal output, like in production.
202 /// Return `Err(OperationError::Http { status: None, .. })` to simulate a
203 /// transport failure.
204 ///
205 /// # Panics
206 ///
207 /// Panics when called after the first run.
208 ///
209 /// # Examples
210 ///
211 /// ```no_run
212 /// use ironflow_engine::testing::{MockHttpResponse, TestEngine};
213 /// use serde_json::json;
214 ///
215 /// let harness = TestEngine::new()
216 /// .with_mock_http(|_cfg| Ok(MockHttpResponse::json(201, &json!({"id": 7}))));
217 /// # let _ = harness;
218 /// ```
219 pub fn with_mock_http(
220 mut self,
221 f: impl Fn(&HttpConfig) -> Result<MockHttpResponse, OperationError> + Send + Sync + 'static,
222 ) -> Self {
223 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
224 self.mocks = self.mocks.http(f);
225 self
226 }
227
228 /// Resolve every approval gate with `outcome` instead of suspending.
229 ///
230 /// Without this, a gated handler ends the run in
231 /// [`RunStatus::AwaitingApproval`] and [`resume`](Self::resume) continues it.
232 ///
233 /// # Panics
234 ///
235 /// Panics when called after the first run.
236 ///
237 /// # Examples
238 ///
239 /// ```no_run
240 /// use ironflow_engine::testing::{ApprovalOutcome, TestEngine};
241 ///
242 /// let harness = TestEngine::new().with_mock_approval(ApprovalOutcome::Approved);
243 /// # let _ = harness;
244 /// ```
245 pub fn with_mock_approval(mut self, outcome: ApprovalOutcome) -> Self {
246 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
247 self.mocks = self.mocks.approval(outcome);
248 self
249 }
250
251 /// Answer every human input step with `f` instead of suspending.
252 ///
253 /// `f` receives the step name and its config. Without this, a handler that
254 /// asks for a human input ends the run in [`RunStatus::AwaitingApproval`];
255 /// write the answer on the step through the store, then call
256 /// [`resume`](Self::resume).
257 ///
258 /// # Panics
259 ///
260 /// Panics when called after the first run.
261 ///
262 /// # Examples
263 ///
264 /// ```no_run
265 /// use ironflow_engine::testing::{HumanInputOutcome, TestEngine};
266 /// use serde_json::json;
267 ///
268 /// let harness = TestEngine::new().with_mock_human_input(|_name, _cfg| {
269 /// HumanInputOutcome::Provided(json!({"answers": ["staging"]}))
270 /// });
271 /// # let _ = harness;
272 /// ```
273 pub fn with_mock_human_input(
274 mut self,
275 f: impl Fn(&str, &HumanInputConfig) -> HumanInputOutcome + Send + Sync + 'static,
276 ) -> Self {
277 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
278 self.mocks = self.mocks.human_input(f);
279 self
280 }
281
282 /// Resolve every signal step with `f` instead of waiting for a signal.
283 ///
284 /// `f` receives the step name, the signal name and the key, and returns
285 /// [`SignalOutcome::Received`] with the payload, or
286 /// [`SignalOutcome::TimedOut`] to make `ctx.wait_for_signal` return `None`.
287 /// Without this, a handler that waits for a signal ends the run in
288 /// [`RunStatus::Sleeping`].
289 ///
290 /// # Panics
291 ///
292 /// Panics when called after the first run.
293 ///
294 /// # Examples
295 ///
296 /// ```no_run
297 /// use ironflow_engine::testing::{SignalOutcome, TestEngine};
298 /// use serde_json::json;
299 ///
300 /// let harness = TestEngine::new().with_mock_signal(|_step, _name, _key| {
301 /// SignalOutcome::Received(json!({"status": "success"}))
302 /// });
303 /// # let _ = harness;
304 /// ```
305 pub fn with_mock_signal(
306 mut self,
307 f: impl Fn(&str, &str, &str) -> SignalOutcome + Send + Sync + 'static,
308 ) -> Self {
309 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
310 self.mocks = self.mocks.signal(f);
311 self
312 }
313
314 /// Answer every agent step with `f` instead of invoking a backend.
315 ///
316 /// # Panics
317 ///
318 /// Panics when called after the first run.
319 ///
320 /// # Examples
321 ///
322 /// ```no_run
323 /// use ironflow_core::provider::AgentOutput;
324 /// use ironflow_engine::testing::TestEngine;
325 /// use serde_json::json;
326 ///
327 /// let harness = TestEngine::new()
328 /// .with_mock_agent(|_cfg| Ok(AgentOutput::new(json!({"score": 9}))));
329 /// # let _ = harness;
330 /// ```
331 pub fn with_mock_agent(
332 mut self,
333 f: impl Fn(&AgentConfig) -> Result<AgentOutput, AgentError> + Send + Sync + 'static,
334 ) -> Self {
335 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
336 self.provider = Some(Arc::new(MockAgentProvider::new(f)));
337 self
338 }
339
340 /// Replay agent steps from fixtures recorded in `fixtures_dir`.
341 ///
342 /// `fixtures_dir` is the **directory**, not a file:
343 /// [`RecordReplayProvider`] keys each fixture by a hash of the
344 /// [`AgentConfig`] and stores it as `<hash>.json` inside it. A missing
345 /// fixture falls back to [`MissingAgentProvider`], so the step fails loudly
346 /// instead of reaching the real Claude CLI.
347 ///
348 /// # Panics
349 ///
350 /// Panics when called after the first run.
351 ///
352 /// # Examples
353 ///
354 /// ```no_run
355 /// use ironflow_engine::testing::TestEngine;
356 ///
357 /// let harness = TestEngine::new().with_recorded_agent("tests/fixtures");
358 /// # let _ = harness;
359 /// ```
360 pub fn with_recorded_agent(mut self, fixtures_dir: &str) -> Self {
361 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
362 self.provider = Some(Arc::new(RecordReplayProvider::replay(
363 MissingAgentProvider,
364 fixtures_dir,
365 )));
366 self
367 }
368
369 /// Use an arbitrary [`AgentProvider`] for agent steps.
370 ///
371 /// The escape hatch for anything the three `with_mock_*` methods do not
372 /// cover, such as recording new fixtures.
373 ///
374 /// # Panics
375 ///
376 /// Panics when called after the first run.
377 ///
378 /// # Examples
379 ///
380 /// ```no_run
381 /// use std::sync::Arc;
382 ///
383 /// use ironflow_core::provider::AgentProvider;
384 /// use ironflow_engine::testing::TestEngine;
385 ///
386 /// # fn example(provider: Arc<dyn AgentProvider>) {
387 /// let harness = TestEngine::new().with_agent_provider(provider);
388 /// # let _ = harness;
389 /// # }
390 /// ```
391 pub fn with_agent_provider(mut self, provider: Arc<dyn AgentProvider>) -> Self {
392 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
393 self.provider = Some(provider);
394 self
395 }
396
397 /// Use a [`DecisionProvider`] for `ctx.decision(...)` steps.
398 ///
399 /// Decision steps are not intercepted: without a provider they fail with
400 /// [`EngineError::NoDecisionProvider`].
401 ///
402 /// # Panics
403 ///
404 /// Panics when called after the first run.
405 ///
406 /// # Examples
407 ///
408 /// ```no_run
409 /// use std::sync::Arc;
410 ///
411 /// use ironflow_core::decision::DecisionProvider;
412 /// use ironflow_engine::testing::TestEngine;
413 ///
414 /// # fn example(provider: Arc<dyn DecisionProvider>) {
415 /// let harness = TestEngine::new().with_decision_provider(provider);
416 /// # let _ = harness;
417 /// # }
418 /// ```
419 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
420 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
421 self.decision_provider = Some(provider);
422 self
423 }
424
425 /// Seed a secret readable by the workflow under test.
426 ///
427 /// The secret is written in the scope of the workflow being run, the same
428 /// scope as `ctx.secrets()` and as the operations run through
429 /// `ctx.operation`. Setting the same key twice keeps the last value.
430 /// Requires the `secret-store` feature.
431 ///
432 /// # Panics
433 ///
434 /// Panics when called after the first run.
435 ///
436 /// # Examples
437 ///
438 /// ```no_run
439 /// use ironflow_engine::testing::TestEngine;
440 ///
441 /// let harness = TestEngine::new().with_secret("git_token", "ghp_test");
442 /// # let _ = harness;
443 /// ```
444 #[cfg(feature = "secret-store")]
445 pub fn with_secret(mut self, key: &str, value: &str) -> Self {
446 assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
447 self.secrets.push((key.to_string(), value.to_string()));
448 self
449 }
450
451 /// The store backing this harness, for assertions the accessors do not
452 /// cover (child runs, step dependencies, logs).
453 ///
454 /// # Examples
455 ///
456 /// ```no_run
457 /// use ironflow_engine::error::EngineError;
458 /// use ironflow_engine::testing::TestEngine;
459 /// use ironflow_store::store::RunStore;
460 /// use uuid::Uuid;
461 ///
462 /// # async fn example(harness: &TestEngine, run_id: Uuid) -> Result<(), EngineError> {
463 /// let steps = harness.store().list_steps(run_id).await?;
464 /// assert!(!steps.is_empty());
465 /// # Ok(())
466 /// # }
467 /// ```
468 pub fn store(&self) -> Arc<InMemoryStore> {
469 self.store.clone()
470 }
471
472 /// Build the underlying [`Engine`] on first use.
473 ///
474 /// Handlers are drained into it, which is why every builder method asserts
475 /// that no run has happened yet.
476 fn ensure_engine(&mut self) -> Result<(), EngineError> {
477 if self.engine.is_some() {
478 return Ok(());
479 }
480
481 let store: Arc<dyn Store> = self.store.clone();
482 let provider = self
483 .provider
484 .clone()
485 .unwrap_or_else(|| Arc::new(MissingAgentProvider));
486 let mocks: Arc<dyn StepInterceptor> = Arc::new(self.mocks.clone());
487 let mut engine = Engine::new(store, provider).with_step_interceptor(mocks);
488 if let Some(decision_provider) = self.decision_provider.clone() {
489 engine = engine.with_decision_provider(decision_provider);
490 }
491 for handler in self.handlers.drain(..) {
492 engine.register_boxed(handler)?;
493 }
494
495 self.engine = Some(engine);
496 Ok(())
497 }
498
499 /// Run the first handler registered with [`with_handler`](Self::with_handler).
500 ///
501 /// A handler that fails is not an error: the returned [`TestResult`] then
502 /// carries [`RunStatus::Failed`] and the message in
503 /// [`TestResult::error`].
504 ///
505 /// # Errors
506 ///
507 /// Returns [`EngineError::InvalidWorkflow`] when no handler was registered
508 /// or two handlers share a name, and [`EngineError::Store`] when the
509 /// in-memory store rejects a write.
510 ///
511 /// # Examples
512 ///
513 /// ```no_run
514 /// use ironflow_engine::prelude::*;
515 /// use ironflow_engine::testing::{MockShellOutput, TestEngine};
516 /// use serde_json::json;
517 ///
518 /// # struct Deploy;
519 /// # impl WorkflowHandler for Deploy {
520 /// # fn name(&self) -> &str { "deploy" }
521 /// # fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
522 /// # Box::pin(async move { Ok(()) })
523 /// # }
524 /// # }
525 /// # async fn example() -> Result<(), EngineError> {
526 /// let result = TestEngine::new()
527 /// .with_handler(Deploy)
528 /// .with_mock_shell(|_cfg| Ok(MockShellOutput::ok("done")))
529 /// .run(json!({}))
530 /// .await?;
531 /// assert!(result.is_completed());
532 /// # Ok(())
533 /// # }
534 /// ```
535 pub async fn run(&mut self, payload: Value) -> Result<TestResult, EngineError> {
536 let name = self.primary.clone().ok_or_else(|| {
537 EngineError::InvalidWorkflow(
538 "TestEngine has no handler: call with_handler(...) first".to_string(),
539 )
540 })?;
541 self.run_workflow(&name, payload).await
542 }
543
544 /// Run a specific registered handler by name.
545 ///
546 /// # Errors
547 ///
548 /// Same as [`run`](Self::run), plus [`EngineError::InvalidWorkflow`] when
549 /// `name` matches no registered handler.
550 ///
551 /// # Examples
552 ///
553 /// ```no_run
554 /// use ironflow_engine::error::EngineError;
555 /// use ironflow_engine::testing::TestEngine;
556 /// use serde_json::json;
557 ///
558 /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
559 /// let result = harness.run_workflow("child", json!({"id": 1})).await?;
560 /// assert!(result.is_completed());
561 /// # Ok(())
562 /// # }
563 /// ```
564 pub async fn run_workflow(
565 &mut self,
566 name: &str,
567 payload: Value,
568 ) -> Result<TestResult, EngineError> {
569 self.ensure_engine()?;
570
571 #[cfg(feature = "secret-store")]
572 {
573 let scoped = ScopedSecretStore::for_workflow(
574 Uuid::new_v5(&Uuid::NAMESPACE_OID, name.as_bytes()),
575 self.store.clone(),
576 );
577 for (key, value) in &self.secrets {
578 scoped.set(key, value).await?;
579 }
580 }
581
582 let engine = self.engine.as_ref().expect("ensure_engine built it");
583
584 // Enqueue then execute, the way the worker does, so the run id is known
585 // before execution and a failed run can still be read back.
586 let trigger = TriggerKind::Manual;
587 let run = engine.enqueue_handler(name, trigger, payload, 0).await?;
588 self.store
589 .update_run_status(run.id, RunStatus::Running)
590 .await?;
591 let execution = engine.execute_handler_run(run.id).await;
592 self.collect(run.id, execution).await
593 }
594
595 /// Resume a run suspended on an approval gate, a human input, a delay or a
596 /// signal wait, the way the API server and the waker do.
597 ///
598 /// A suspended sub-workflow child run is resumed through its root run,
599 /// like in production: the root replays and re-enters the child. The
600 /// returned [`TestResult`] then describes the root run.
601 ///
602 /// # Errors
603 ///
604 /// Returns [`EngineError::Store`] when the run does not exist or is not
605 /// resumable, and [`EngineError::InvalidWorkflow`] when its handler is no
606 /// longer registered.
607 ///
608 /// # Panics
609 ///
610 /// Never in practice: the engine is built by `ensure_engine` just before.
611 ///
612 /// # Examples
613 ///
614 /// ```no_run
615 /// use ironflow_engine::error::EngineError;
616 /// use ironflow_engine::testing::TestEngine;
617 /// use ironflow_store::models::RunStatus;
618 /// use serde_json::json;
619 ///
620 /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
621 /// let suspended = harness.run(json!({})).await?;
622 /// assert_eq!(suspended.status(), RunStatus::AwaitingApproval);
623 ///
624 /// let resumed = harness.resume(suspended.run_id()).await?;
625 /// assert_eq!(resumed.status(), RunStatus::Completed);
626 /// # Ok(())
627 /// # }
628 /// ```
629 pub async fn resume(&mut self, run_id: Uuid) -> Result<TestResult, EngineError> {
630 self.ensure_engine()?;
631 let engine = self.engine.as_ref().expect("ensure_engine built it");
632
633 let run = self
634 .store
635 .get_run(run_id)
636 .await?
637 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
638 // A sleeping run is woken through `Pending`, like the waker does.
639 if run.status.state == RunStatus::Sleeping {
640 self.store
641 .update_run_status(run_id, RunStatus::Pending)
642 .await?;
643 }
644 self.store
645 .update_run_status(run_id, RunStatus::Running)
646 .await?;
647 let execution = engine.resume_run(run_id).await;
648 let reported = chain_root(&run).unwrap_or(run_id);
649 self.collect(reported, execution).await
650 }
651
652 /// Read the run and its steps back from the store.
653 ///
654 /// The engine returns `Err` for a failed run, so the steps are always read
655 /// from the store rather than from the execution result.
656 async fn collect(
657 &self,
658 run_id: Uuid,
659 execution: Result<WorkflowResult, EngineError>,
660 ) -> Result<TestResult, EngineError> {
661 let run = self
662 .store
663 .get_run(run_id)
664 .await?
665 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
666 let steps = self.store.list_steps(run_id).await?;
667
668 let (step_results, error) = match execution {
669 Ok(result) => (result.steps, run.error.clone()),
670 Err(err) => (Vec::new(), Some(err.to_string())),
671 };
672
673 Ok(TestResult::new(run, steps, step_results, error))
674 }
675}
676
677#[cfg(test)]
678mod tests {
679 use std::time::Duration;
680
681 use async_trait::async_trait;
682 use serde_json::json;
683 use tokio::time::timeout;
684
685 use ironflow_core::operation::{Operation, OperationContext};
686
687 use crate::context::WorkflowContext;
688 use crate::handler::HandlerFuture;
689
690 use super::*;
691
692 /// Reads `git_token` through the operation context.
693 struct ReadsToken;
694
695 #[async_trait]
696 impl Operation for ReadsToken {
697 fn kind(&self) -> &str {
698 "reads-token"
699 }
700
701 async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
702 let token = ctx.secrets().get("git_token").await?;
703 Ok(json!({"token": token.map(|secret| secret.value)}))
704 }
705 }
706
707 /// Runs [`ReadsToken`] through `ctx.operation`.
708 struct UsesSecretOp;
709
710 impl WorkflowHandler for UsesSecretOp {
711 fn name(&self) -> &str {
712 "uses-secret-op"
713 }
714
715 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
716 Box::pin(async move {
717 ctx.operation("read-token", &ReadsToken).await?;
718 Ok(())
719 })
720 }
721 }
722
723 #[tokio::test]
724 async fn operation_reading_absent_secret_completes_under_test_engine() {
725 timeout(Duration::from_secs(10), async {
726 let result = TestEngine::new()
727 .with_handler(UsesSecretOp)
728 .run(json!({}))
729 .await
730 .unwrap();
731
732 assert_eq!(
733 result.status(),
734 RunStatus::Completed,
735 "{:?}",
736 result.error()
737 );
738 assert!(result.step("read-token").output()["token"].is_null());
739 })
740 .await
741 .expect("test timed out");
742 }
743
744 #[cfg(feature = "secret-store")]
745 #[tokio::test]
746 async fn with_secret_value_reaches_operation_via_ctx_operation() {
747 timeout(Duration::from_secs(10), async {
748 let result = TestEngine::new()
749 .with_handler(UsesSecretOp)
750 .with_secret("git_token", "ghp_test")
751 .run(json!({}))
752 .await
753 .unwrap();
754
755 assert_eq!(
756 result.status(),
757 RunStatus::Completed,
758 "{:?}",
759 result.error()
760 );
761 assert_eq!(result.step("read-token").output()["token"], "ghp_test");
762 })
763 .await
764 .expect("test timed out");
765 }
766
767 #[cfg(feature = "secret-store")]
768 #[tokio::test]
769 async fn with_secret_last_value_wins() {
770 timeout(Duration::from_secs(10), async {
771 let result = TestEngine::new()
772 .with_handler(UsesSecretOp)
773 .with_secret("git_token", "first")
774 .with_secret("git_token", "second")
775 .run(json!({}))
776 .await
777 .unwrap();
778
779 assert_eq!(result.step("read-token").output()["token"], "second");
780 })
781 .await
782 .expect("test timed out");
783 }
784
785 #[cfg(feature = "secret-store")]
786 #[tokio::test]
787 async fn with_secret_is_readable_through_ctx_secrets() {
788 timeout(Duration::from_secs(10), async {
789 let mut harness = TestEngine::new()
790 .with_handler(UsesSecretOp)
791 .with_secret("k", "v");
792 harness.run(json!({})).await.unwrap();
793
794 let store = harness.store();
795 let scope = |name: &str| {
796 ScopedSecretStore::for_workflow(
797 Uuid::new_v5(&Uuid::NAMESPACE_OID, name.as_bytes()),
798 store.clone(),
799 )
800 };
801 let own = scope("uses-secret-op").get("k").await.unwrap();
802 assert_eq!(own.unwrap().value, "v");
803 assert!(scope("other-workflow").get("k").await.unwrap().is_none());
804 })
805 .await
806 .expect("test timed out");
807 }
808
809 #[cfg(feature = "secret-store")]
810 #[tokio::test]
811 #[should_panic(expected = "configure the TestEngine before its first run")]
812 async fn with_secret_after_first_run_panics() {
813 let mut harness = TestEngine::new().with_handler(UsesSecretOp);
814 harness.run(json!({})).await.unwrap();
815 let _ = harness.with_secret("git_token", "late");
816 }
817}