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}