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}