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, 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, the way the API server does.
596    ///
597    /// # Errors
598    ///
599    /// Returns [`EngineError::Store`] when the run does not exist or is not
600    /// resumable, and [`EngineError::InvalidWorkflow`] when its handler is no
601    /// longer registered.
602    ///
603    /// # Examples
604    ///
605    /// ```no_run
606    /// use ironflow_engine::error::EngineError;
607    /// use ironflow_engine::testing::TestEngine;
608    /// use ironflow_store::models::RunStatus;
609    /// use serde_json::json;
610    ///
611    /// # async fn example(harness: &mut TestEngine) -> Result<(), EngineError> {
612    /// let suspended = harness.run(json!({})).await?;
613    /// assert_eq!(suspended.status(), RunStatus::AwaitingApproval);
614    ///
615    /// let resumed = harness.resume(suspended.run_id()).await?;
616    /// assert_eq!(resumed.status(), RunStatus::Completed);
617    /// # Ok(())
618    /// # }
619    /// ```
620    pub async fn resume(&mut self, run_id: Uuid) -> Result<TestResult, EngineError> {
621        self.ensure_engine()?;
622        let engine = self.engine.as_ref().expect("ensure_engine built it");
623
624        self.store
625            .update_run_status(run_id, RunStatus::Running)
626            .await?;
627        let execution = engine.resume_run(run_id).await;
628        self.collect(run_id, execution).await
629    }
630
631    /// Read the run and its steps back from the store.
632    ///
633    /// The engine returns `Err` for a failed run, so the steps are always read
634    /// from the store rather than from the execution result.
635    async fn collect(
636        &self,
637        run_id: Uuid,
638        execution: Result<WorkflowResult, EngineError>,
639    ) -> Result<TestResult, EngineError> {
640        let run = self
641            .store
642            .get_run(run_id)
643            .await?
644            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
645        let steps = self.store.list_steps(run_id).await?;
646
647        let (step_results, error) = match execution {
648            Ok(result) => (result.steps, run.error.clone()),
649            Err(err) => (Vec::new(), Some(err.to_string())),
650        };
651
652        Ok(TestResult::new(run, steps, step_results, error))
653    }
654}
655
656#[cfg(test)]
657mod tests {
658    use std::time::Duration;
659
660    use async_trait::async_trait;
661    use serde_json::json;
662    use tokio::time::timeout;
663
664    use ironflow_core::operation::{Operation, OperationContext};
665
666    use crate::context::WorkflowContext;
667    use crate::handler::HandlerFuture;
668
669    use super::*;
670
671    /// Reads `git_token` through the operation context.
672    struct ReadsToken;
673
674    #[async_trait]
675    impl Operation for ReadsToken {
676        fn kind(&self) -> &str {
677            "reads-token"
678        }
679
680        async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
681            let token = ctx.secrets().get("git_token").await?;
682            Ok(json!({"token": token.map(|secret| secret.value)}))
683        }
684    }
685
686    /// Runs [`ReadsToken`] through `ctx.operation`.
687    struct UsesSecretOp;
688
689    impl WorkflowHandler for UsesSecretOp {
690        fn name(&self) -> &str {
691            "uses-secret-op"
692        }
693
694        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
695            Box::pin(async move {
696                ctx.operation("read-token", &ReadsToken).await?;
697                Ok(())
698            })
699        }
700    }
701
702    #[tokio::test]
703    async fn operation_reading_absent_secret_completes_under_test_engine() {
704        timeout(Duration::from_secs(10), async {
705            let result = TestEngine::new()
706                .with_handler(UsesSecretOp)
707                .run(json!({}))
708                .await
709                .unwrap();
710
711            assert_eq!(
712                result.status(),
713                RunStatus::Completed,
714                "{:?}",
715                result.error()
716            );
717            assert!(result.step("read-token").output()["token"].is_null());
718        })
719        .await
720        .expect("test timed out");
721    }
722
723    #[cfg(feature = "secret-store")]
724    #[tokio::test]
725    async fn with_secret_value_reaches_operation_via_ctx_operation() {
726        timeout(Duration::from_secs(10), async {
727            let result = TestEngine::new()
728                .with_handler(UsesSecretOp)
729                .with_secret("git_token", "ghp_test")
730                .run(json!({}))
731                .await
732                .unwrap();
733
734            assert_eq!(
735                result.status(),
736                RunStatus::Completed,
737                "{:?}",
738                result.error()
739            );
740            assert_eq!(result.step("read-token").output()["token"], "ghp_test");
741        })
742        .await
743        .expect("test timed out");
744    }
745
746    #[cfg(feature = "secret-store")]
747    #[tokio::test]
748    async fn with_secret_last_value_wins() {
749        timeout(Duration::from_secs(10), async {
750            let result = TestEngine::new()
751                .with_handler(UsesSecretOp)
752                .with_secret("git_token", "first")
753                .with_secret("git_token", "second")
754                .run(json!({}))
755                .await
756                .unwrap();
757
758            assert_eq!(result.step("read-token").output()["token"], "second");
759        })
760        .await
761        .expect("test timed out");
762    }
763
764    #[cfg(feature = "secret-store")]
765    #[tokio::test]
766    async fn with_secret_is_readable_through_ctx_secrets() {
767        timeout(Duration::from_secs(10), async {
768            let mut harness = TestEngine::new()
769                .with_handler(UsesSecretOp)
770                .with_secret("k", "v");
771            harness.run(json!({})).await.unwrap();
772
773            let store = harness.store();
774            let scope = |name: &str| {
775                ScopedSecretStore::for_workflow(
776                    Uuid::new_v5(&Uuid::NAMESPACE_OID, name.as_bytes()),
777                    store.clone(),
778                )
779            };
780            let own = scope("uses-secret-op").get("k").await.unwrap();
781            assert_eq!(own.unwrap().value, "v");
782            assert!(scope("other-workflow").get("k").await.unwrap().is_none());
783        })
784        .await
785        .expect("test timed out");
786    }
787
788    #[cfg(feature = "secret-store")]
789    #[tokio::test]
790    #[should_panic(expected = "configure the TestEngine before its first run")]
791    async fn with_secret_after_first_run_panics() {
792        let mut harness = TestEngine::new().with_handler(UsesSecretOp);
793        harness.run(json!({})).await.unwrap();
794        let _ = harness.with_secret("git_token", "late");
795    }
796}