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
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}