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};
26use crate::error::EngineError;
27use crate::executor::{ApprovalOutcome, HumanInputOutcome, 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    /// Answer every agent step with `f` instead of invoking a backend.
283    ///
284    /// # Panics
285    ///
286    /// Panics when called after the first run.
287    ///
288    /// # Examples
289    ///
290    /// ```no_run
291    /// use ironflow_core::provider::AgentOutput;
292    /// use ironflow_engine::testing::TestEngine;
293    /// use serde_json::json;
294    ///
295    /// let harness = TestEngine::new()
296    ///     .with_mock_agent(|_cfg| Ok(AgentOutput::new(json!({"score": 9}))));
297    /// # let _ = harness;
298    /// ```
299    pub fn with_mock_agent(
300        mut self,
301        f: impl Fn(&AgentConfig) -> Result<AgentOutput, AgentError> + Send + Sync + 'static,
302    ) -> Self {
303        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
304        self.provider = Some(Arc::new(MockAgentProvider::new(f)));
305        self
306    }
307
308    /// Replay agent steps from fixtures recorded in `fixtures_dir`.
309    ///
310    /// `fixtures_dir` is the **directory**, not a file:
311    /// [`RecordReplayProvider`] keys each fixture by a hash of the
312    /// [`AgentConfig`] and stores it as `<hash>.json` inside it. A missing
313    /// fixture falls back to [`MissingAgentProvider`], so the step fails loudly
314    /// instead of reaching the real Claude CLI.
315    ///
316    /// # Panics
317    ///
318    /// Panics when called after the first run.
319    ///
320    /// # Examples
321    ///
322    /// ```no_run
323    /// use ironflow_engine::testing::TestEngine;
324    ///
325    /// let harness = TestEngine::new().with_recorded_agent("tests/fixtures");
326    /// # let _ = harness;
327    /// ```
328    pub fn with_recorded_agent(mut self, fixtures_dir: &str) -> Self {
329        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
330        self.provider = Some(Arc::new(RecordReplayProvider::replay(
331            MissingAgentProvider,
332            fixtures_dir,
333        )));
334        self
335    }
336
337    /// Use an arbitrary [`AgentProvider`] for agent steps.
338    ///
339    /// The escape hatch for anything the three `with_mock_*` methods do not
340    /// cover, such as recording new fixtures.
341    ///
342    /// # Panics
343    ///
344    /// Panics when called after the first run.
345    ///
346    /// # Examples
347    ///
348    /// ```no_run
349    /// use std::sync::Arc;
350    ///
351    /// use ironflow_core::provider::AgentProvider;
352    /// use ironflow_engine::testing::TestEngine;
353    ///
354    /// # fn example(provider: Arc<dyn AgentProvider>) {
355    /// let harness = TestEngine::new().with_agent_provider(provider);
356    /// # let _ = harness;
357    /// # }
358    /// ```
359    pub fn with_agent_provider(mut self, provider: Arc<dyn AgentProvider>) -> Self {
360        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
361        self.provider = Some(provider);
362        self
363    }
364
365    /// Use a [`DecisionProvider`] for `ctx.decision(...)` steps.
366    ///
367    /// Decision steps are not intercepted: without a provider they fail with
368    /// [`EngineError::NoDecisionProvider`].
369    ///
370    /// # Panics
371    ///
372    /// Panics when called after the first run.
373    ///
374    /// # Examples
375    ///
376    /// ```no_run
377    /// use std::sync::Arc;
378    ///
379    /// use ironflow_core::decision::DecisionProvider;
380    /// use ironflow_engine::testing::TestEngine;
381    ///
382    /// # fn example(provider: Arc<dyn DecisionProvider>) {
383    /// let harness = TestEngine::new().with_decision_provider(provider);
384    /// # let _ = harness;
385    /// # }
386    /// ```
387    pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
388        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
389        self.decision_provider = Some(provider);
390        self
391    }
392
393    /// Seed a secret readable by the workflow under test.
394    ///
395    /// The secret is written in the scope of the workflow being run, the same
396    /// scope as `ctx.secrets()` and as the operations run through
397    /// `ctx.operation`. Setting the same key twice keeps the last value.
398    /// Requires the `secret-store` feature.
399    ///
400    /// # Panics
401    ///
402    /// Panics when called after the first run.
403    ///
404    /// # Examples
405    ///
406    /// ```no_run
407    /// use ironflow_engine::testing::TestEngine;
408    ///
409    /// let harness = TestEngine::new().with_secret("git_token", "ghp_test");
410    /// # let _ = harness;
411    /// ```
412    #[cfg(feature = "secret-store")]
413    pub fn with_secret(mut self, key: &str, value: &str) -> Self {
414        assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
415        self.secrets.push((key.to_string(), value.to_string()));
416        self
417    }
418
419    /// The store backing this harness, for assertions the accessors do not
420    /// cover (child runs, step dependencies, logs).
421    ///
422    /// # Examples
423    ///
424    /// ```no_run
425    /// use ironflow_engine::error::EngineError;
426    /// use ironflow_engine::testing::TestEngine;
427    /// use ironflow_store::store::RunStore;
428    /// use uuid::Uuid;
429    ///
430    /// # async fn example(harness: &TestEngine, run_id: Uuid) -> Result<(), EngineError> {
431    /// let steps = harness.store().list_steps(run_id).await?;
432    /// assert!(!steps.is_empty());
433    /// # Ok(())
434    /// # }
435    /// ```
436    pub fn store(&self) -> Arc<InMemoryStore> {
437        self.store.clone()
438    }
439
440    /// Build the underlying [`Engine`] on first use.
441    ///
442    /// Handlers are drained into it, which is why every builder method asserts
443    /// that no run has happened yet.
444    fn ensure_engine(&mut self) -> Result<(), EngineError> {
445        if self.engine.is_some() {
446            return Ok(());
447        }
448
449        let store: Arc<dyn Store> = self.store.clone();
450        let provider = self
451            .provider
452            .clone()
453            .unwrap_or_else(|| Arc::new(MissingAgentProvider));
454        let mocks: Arc<dyn StepInterceptor> = Arc::new(self.mocks.clone());
455        let mut engine = Engine::new(store, provider).with_step_interceptor(mocks);
456        if let Some(decision_provider) = self.decision_provider.clone() {
457            engine = engine.with_decision_provider(decision_provider);
458        }
459        for handler in self.handlers.drain(..) {
460            engine.register_boxed(handler)?;
461        }
462
463        self.engine = Some(engine);
464        Ok(())
465    }
466
467    /// Run the first handler registered with [`with_handler`](Self::with_handler).
468    ///
469    /// A handler that fails is not an error: the returned [`TestResult`] then
470    /// carries [`RunStatus::Failed`] and the message in
471    /// [`TestResult::error`].
472    ///
473    /// # Errors
474    ///
475    /// Returns [`EngineError::InvalidWorkflow`] when no handler was registered
476    /// or two handlers share a name, and [`EngineError::Store`] when the
477    /// in-memory store rejects a write.
478    ///
479    /// # Examples
480    ///
481    /// ```no_run
482    /// use ironflow_engine::prelude::*;
483    /// use ironflow_engine::testing::{MockShellOutput, TestEngine};
484    /// use serde_json::json;
485    ///
486    /// # struct Deploy;
487    /// # impl WorkflowHandler for Deploy {
488    /// #     fn name(&self) -> &str { "deploy" }
489    /// #     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
490    /// #         Box::pin(async move { Ok(()) })
491    /// #     }
492    /// # }
493    /// # async fn example() -> Result<(), EngineError> {
494    /// let result = TestEngine::new()
495    ///     .with_handler(Deploy)
496    ///     .with_mock_shell(|_cfg| Ok(MockShellOutput::ok("done")))
497    ///     .run(json!({}))
498    ///     .await?;
499    /// assert!(result.is_completed());
500    /// # Ok(())
501    /// # }
502    /// ```
503    pub async fn run(&mut self, payload: Value) -> Result<TestResult, EngineError> {
504        let name = self.primary.clone().ok_or_else(|| {
505            EngineError::InvalidWorkflow(
506                "TestEngine has no handler: call with_handler(...) first".to_string(),
507            )
508        })?;
509        self.run_workflow(&name, payload).await
510    }
511
512    /// Run a specific registered handler by name.
513    ///
514    /// # Errors
515    ///
516    /// Same as [`run`](Self::run), plus [`EngineError::InvalidWorkflow`] when
517    /// `name` matches no registered handler.
518    ///
519    /// # Examples
520    ///
521    /// ```no_run
522    /// use ironflow_engine::error::EngineError;
523    /// use ironflow_engine::testing::TestEngine;
524    /// use serde_json::json;
525    ///
526    /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
527    /// let result = harness.run_workflow("child", json!({"id": 1})).await?;
528    /// assert!(result.is_completed());
529    /// # Ok(())
530    /// # }
531    /// ```
532    pub async fn run_workflow(
533        &mut self,
534        name: &str,
535        payload: Value,
536    ) -> Result<TestResult, EngineError> {
537        self.ensure_engine()?;
538
539        #[cfg(feature = "secret-store")]
540        {
541            let scoped = ScopedSecretStore::for_workflow(
542                Uuid::new_v5(&Uuid::NAMESPACE_OID, name.as_bytes()),
543                self.store.clone(),
544            );
545            for (key, value) in &self.secrets {
546                scoped.set(key, value).await?;
547            }
548        }
549
550        let engine = self.engine.as_ref().expect("ensure_engine built it");
551
552        // Enqueue then execute, the way the worker does, so the run id is known
553        // before execution and a failed run can still be read back.
554        let trigger = TriggerKind::Manual;
555        let run = engine.enqueue_handler(name, trigger, payload, 0).await?;
556        self.store
557            .update_run_status(run.id, RunStatus::Running)
558            .await?;
559        let execution = engine.execute_handler_run(run.id).await;
560        self.collect(run.id, execution).await
561    }
562
563    /// Resume a run suspended on an approval gate, the way the API server does.
564    ///
565    /// # Errors
566    ///
567    /// Returns [`EngineError::Store`] when the run does not exist or is not
568    /// resumable, and [`EngineError::InvalidWorkflow`] when its handler is no
569    /// longer registered.
570    ///
571    /// # Examples
572    ///
573    /// ```no_run
574    /// use ironflow_engine::error::EngineError;
575    /// use ironflow_engine::testing::TestEngine;
576    /// use ironflow_store::models::RunStatus;
577    /// use serde_json::json;
578    ///
579    /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
580    /// let suspended = harness.run(json!({})).await?;
581    /// assert_eq!(suspended.status(), RunStatus::AwaitingApproval);
582    ///
583    /// let resumed = harness.resume(suspended.run_id()).await?;
584    /// assert_eq!(resumed.status(), RunStatus::Completed);
585    /// # Ok(())
586    /// # }
587    /// ```
588    pub async fn resume(&mut self, run_id: Uuid) -> Result<TestResult, EngineError> {
589        self.ensure_engine()?;
590        let engine = self.engine.as_ref().expect("ensure_engine built it");
591
592        self.store
593            .update_run_status(run_id, RunStatus::Running)
594            .await?;
595        let execution = engine.resume_run(run_id).await;
596        self.collect(run_id, execution).await
597    }
598
599    /// Read the run and its steps back from the store.
600    ///
601    /// The engine returns `Err` for a failed run, so the steps are always read
602    /// from the store rather than from the execution result.
603    async fn collect(
604        &self,
605        run_id: Uuid,
606        execution: Result<WorkflowResult, EngineError>,
607    ) -> Result<TestResult, EngineError> {
608        let run = self
609            .store
610            .get_run(run_id)
611            .await?
612            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
613        let steps = self.store.list_steps(run_id).await?;
614
615        let (step_results, error) = match execution {
616            Ok(result) => (result.steps, run.error.clone()),
617            Err(err) => (Vec::new(), Some(err.to_string())),
618        };
619
620        Ok(TestResult::new(run, steps, step_results, error))
621    }
622}
623
624#[cfg(test)]
625mod tests {
626    use std::time::Duration;
627
628    use async_trait::async_trait;
629    use serde_json::json;
630    use tokio::time::timeout;
631
632    use ironflow_core::operation::{Operation, OperationContext};
633
634    use crate::context::WorkflowContext;
635    use crate::handler::HandlerFuture;
636
637    use super::*;
638
639    /// Reads `git_token` through the operation context.
640    struct ReadsToken;
641
642    #[async_trait]
643    impl Operation for ReadsToken {
644        fn kind(&self) -> &str {
645            "reads-token"
646        }
647
648        async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
649            let token = ctx.secrets().get("git_token").await?;
650            Ok(json!({"token": token.map(|secret| secret.value)}))
651        }
652    }
653
654    /// Runs [`ReadsToken`] through `ctx.operation`.
655    struct UsesSecretOp;
656
657    impl WorkflowHandler for UsesSecretOp {
658        fn name(&self) -> &str {
659            "uses-secret-op"
660        }
661
662        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
663            Box::pin(async move {
664                ctx.operation("read-token", &ReadsToken).await?;
665                Ok(())
666            })
667        }
668    }
669
670    #[tokio::test]
671    async fn operation_reading_absent_secret_completes_under_test_engine() {
672        timeout(Duration::from_secs(10), async {
673            let result = TestEngine::new()
674                .with_handler(UsesSecretOp)
675                .run(json!({}))
676                .await
677                .unwrap();
678
679            assert_eq!(
680                result.status(),
681                RunStatus::Completed,
682                "{:?}",
683                result.error()
684            );
685            assert!(result.step("read-token").output()["token"].is_null());
686        })
687        .await
688        .expect("test timed out");
689    }
690
691    #[cfg(feature = "secret-store")]
692    #[tokio::test]
693    async fn with_secret_value_reaches_operation_via_ctx_operation() {
694        timeout(Duration::from_secs(10), async {
695            let result = TestEngine::new()
696                .with_handler(UsesSecretOp)
697                .with_secret("git_token", "ghp_test")
698                .run(json!({}))
699                .await
700                .unwrap();
701
702            assert_eq!(
703                result.status(),
704                RunStatus::Completed,
705                "{:?}",
706                result.error()
707            );
708            assert_eq!(result.step("read-token").output()["token"], "ghp_test");
709        })
710        .await
711        .expect("test timed out");
712    }
713
714    #[cfg(feature = "secret-store")]
715    #[tokio::test]
716    async fn with_secret_last_value_wins() {
717        timeout(Duration::from_secs(10), async {
718            let result = TestEngine::new()
719                .with_handler(UsesSecretOp)
720                .with_secret("git_token", "first")
721                .with_secret("git_token", "second")
722                .run(json!({}))
723                .await
724                .unwrap();
725
726            assert_eq!(result.step("read-token").output()["token"], "second");
727        })
728        .await
729        .expect("test timed out");
730    }
731
732    #[cfg(feature = "secret-store")]
733    #[tokio::test]
734    async fn with_secret_is_readable_through_ctx_secrets() {
735        timeout(Duration::from_secs(10), async {
736            let mut harness = TestEngine::new()
737                .with_handler(UsesSecretOp)
738                .with_secret("k", "v");
739            harness.run(json!({})).await.unwrap();
740
741            let store = harness.store();
742            let scope = |name: &str| {
743                ScopedSecretStore::for_workflow(
744                    Uuid::new_v5(&Uuid::NAMESPACE_OID, name.as_bytes()),
745                    store.clone(),
746                )
747            };
748            let own = scope("uses-secret-op").get("k").await.unwrap();
749            assert_eq!(own.unwrap().value, "v");
750            assert!(scope("other-workflow").get("k").await.unwrap().is_none());
751        })
752        .await
753        .expect("test timed out");
754    }
755
756    #[cfg(feature = "secret-store")]
757    #[tokio::test]
758    #[should_panic(expected = "configure the TestEngine before its first run")]
759    async fn with_secret_after_first_run_panics() {
760        let mut harness = TestEngine::new().with_handler(UsesSecretOp);
761        harness.run(json!({})).await.unwrap();
762        let _ = harness.with_secret("git_token", "late");
763    }
764}