use std::fmt;
use std::sync::Arc;
#[cfg(feature = "secret-store")]
use rand::random;
use serde_json::Value;
use uuid::Uuid;
use ironflow_core::decision::DecisionProvider;
use ironflow_core::error::{AgentError, OperationError};
use ironflow_core::provider::{AgentConfig, AgentOutput, AgentProvider};
use ironflow_core::providers::record_replay::RecordReplayProvider;
#[cfg(feature = "secret-store")]
use ironflow_store::crypto::MasterKey;
use ironflow_store::error::StoreError;
use ironflow_store::memory::InMemoryStore;
use ironflow_store::models::{RunStatus, TriggerKind};
use ironflow_store::store::{RunStore, Store};
#[cfg(feature = "secret-store")]
use ironflow_store::workflow_secrets::ScopedSecretStore;
use crate::config::{HttpConfig, HumanInputConfig, ShellConfig};
use crate::engine::{Engine, WorkflowResult};
use crate::error::EngineError;
use crate::executor::{ApprovalOutcome, HumanInputOutcome, StepInterceptor};
use crate::handler::WorkflowHandler;
use crate::testing::mocks::{
MissingAgentProvider, MockAgentProvider, MockHttpResponse, MockInterceptor, MockShellOutput,
};
use crate::testing::result::TestResult;
const CONFIGURE_BEFORE_RUN: &str = "configure the TestEngine before its first run";
pub struct TestEngine {
store: Arc<InMemoryStore>,
handlers: Vec<Box<dyn WorkflowHandler>>,
primary: Option<String>,
provider: Option<Arc<dyn AgentProvider>>,
decision_provider: Option<Arc<dyn DecisionProvider>>,
mocks: MockInterceptor,
engine: Option<Engine>,
#[cfg(feature = "secret-store")]
secrets: Vec<(String, String)>,
}
impl Default for TestEngine {
fn default() -> Self {
Self::new()
}
}
impl fmt::Debug for TestEngine {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("TestEngine")
.field("primary", &self.primary)
.field("mocks", &self.mocks)
.field("built", &self.engine.is_some())
.finish_non_exhaustive()
}
}
impl TestEngine {
pub fn new() -> Self {
#[cfg_attr(not(feature = "secret-store"), allow(unused_mut))]
let mut store = InMemoryStore::new();
#[cfg(feature = "secret-store")]
store.set_master_key(
MasterKey::from_bytes(&random::<[u8; 32]>())
.expect("32 bytes is a valid master key length"),
);
Self {
store: Arc::new(store),
handlers: Vec::new(),
primary: None,
provider: None,
decision_provider: None,
mocks: MockInterceptor::new(),
engine: None,
#[cfg(feature = "secret-store")]
secrets: Vec::new(),
}
}
pub fn with_handler(mut self, handler: impl WorkflowHandler + 'static) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
if self.primary.is_none() {
self.primary = Some(handler.name().to_string());
}
self.handlers.push(Box::new(handler));
self
}
pub fn with_mock_shell(
mut self,
f: impl Fn(&ShellConfig) -> Result<MockShellOutput, OperationError> + Send + Sync + 'static,
) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.mocks = self.mocks.shell(f);
self
}
pub fn with_mock_http(
mut self,
f: impl Fn(&HttpConfig) -> Result<MockHttpResponse, OperationError> + Send + Sync + 'static,
) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.mocks = self.mocks.http(f);
self
}
pub fn with_mock_approval(mut self, outcome: ApprovalOutcome) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.mocks = self.mocks.approval(outcome);
self
}
pub fn with_mock_human_input(
mut self,
f: impl Fn(&str, &HumanInputConfig) -> HumanInputOutcome + Send + Sync + 'static,
) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.mocks = self.mocks.human_input(f);
self
}
pub fn with_mock_agent(
mut self,
f: impl Fn(&AgentConfig) -> Result<AgentOutput, AgentError> + Send + Sync + 'static,
) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.provider = Some(Arc::new(MockAgentProvider::new(f)));
self
}
pub fn with_recorded_agent(mut self, fixtures_dir: &str) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.provider = Some(Arc::new(RecordReplayProvider::replay(
MissingAgentProvider,
fixtures_dir,
)));
self
}
pub fn with_agent_provider(mut self, provider: Arc<dyn AgentProvider>) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.provider = Some(provider);
self
}
pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.decision_provider = Some(provider);
self
}
#[cfg(feature = "secret-store")]
pub fn with_secret(mut self, key: &str, value: &str) -> Self {
assert!(self.engine.is_none(), "{CONFIGURE_BEFORE_RUN}");
self.secrets.push((key.to_string(), value.to_string()));
self
}
pub fn store(&self) -> Arc<InMemoryStore> {
self.store.clone()
}
fn ensure_engine(&mut self) -> Result<(), EngineError> {
if self.engine.is_some() {
return Ok(());
}
let store: Arc<dyn Store> = self.store.clone();
let provider = self
.provider
.clone()
.unwrap_or_else(|| Arc::new(MissingAgentProvider));
let mocks: Arc<dyn StepInterceptor> = Arc::new(self.mocks.clone());
let mut engine = Engine::new(store, provider).with_step_interceptor(mocks);
if let Some(decision_provider) = self.decision_provider.clone() {
engine = engine.with_decision_provider(decision_provider);
}
for handler in self.handlers.drain(..) {
engine.register_boxed(handler)?;
}
self.engine = Some(engine);
Ok(())
}
pub async fn run(&mut self, payload: Value) -> Result<TestResult, EngineError> {
let name = self.primary.clone().ok_or_else(|| {
EngineError::InvalidWorkflow(
"TestEngine has no handler: call with_handler(...) first".to_string(),
)
})?;
self.run_workflow(&name, payload).await
}
pub async fn run_workflow(
&mut self,
name: &str,
payload: Value,
) -> Result<TestResult, EngineError> {
self.ensure_engine()?;
#[cfg(feature = "secret-store")]
{
let scoped = ScopedSecretStore::for_workflow(
Uuid::new_v5(&Uuid::NAMESPACE_OID, name.as_bytes()),
self.store.clone(),
);
for (key, value) in &self.secrets {
scoped.set(key, value).await?;
}
}
let engine = self.engine.as_ref().expect("ensure_engine built it");
let trigger = TriggerKind::Manual;
let run = engine.enqueue_handler(name, trigger, payload, 0).await?;
self.store
.update_run_status(run.id, RunStatus::Running)
.await?;
let execution = engine.execute_handler_run(run.id).await;
self.collect(run.id, execution).await
}
pub async fn resume(&mut self, run_id: Uuid) -> Result<TestResult, EngineError> {
self.ensure_engine()?;
let engine = self.engine.as_ref().expect("ensure_engine built it");
self.store
.update_run_status(run_id, RunStatus::Running)
.await?;
let execution = engine.resume_run(run_id).await;
self.collect(run_id, execution).await
}
async fn collect(
&self,
run_id: Uuid,
execution: Result<WorkflowResult, EngineError>,
) -> Result<TestResult, EngineError> {
let run = self
.store
.get_run(run_id)
.await?
.ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
let steps = self.store.list_steps(run_id).await?;
let (step_results, error) = match execution {
Ok(result) => (result.steps, run.error.clone()),
Err(err) => (Vec::new(), Some(err.to_string())),
};
Ok(TestResult::new(run, steps, step_results, error))
}
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use async_trait::async_trait;
use serde_json::json;
use tokio::time::timeout;
use ironflow_core::operation::{Operation, OperationContext};
use crate::context::WorkflowContext;
use crate::handler::HandlerFuture;
use super::*;
struct ReadsToken;
#[async_trait]
impl Operation for ReadsToken {
fn kind(&self) -> &str {
"reads-token"
}
async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
let token = ctx.secrets().get("git_token").await?;
Ok(json!({"token": token.map(|secret| secret.value)}))
}
}
struct UsesSecretOp;
impl WorkflowHandler for UsesSecretOp {
fn name(&self) -> &str {
"uses-secret-op"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
ctx.operation("read-token", &ReadsToken).await?;
Ok(())
})
}
}
#[tokio::test]
async fn operation_reading_absent_secret_completes_under_test_engine() {
timeout(Duration::from_secs(10), async {
let result = TestEngine::new()
.with_handler(UsesSecretOp)
.run(json!({}))
.await
.unwrap();
assert_eq!(
result.status(),
RunStatus::Completed,
"{:?}",
result.error()
);
assert!(result.step("read-token").output()["token"].is_null());
})
.await
.expect("test timed out");
}
#[cfg(feature = "secret-store")]
#[tokio::test]
async fn with_secret_value_reaches_operation_via_ctx_operation() {
timeout(Duration::from_secs(10), async {
let result = TestEngine::new()
.with_handler(UsesSecretOp)
.with_secret("git_token", "ghp_test")
.run(json!({}))
.await
.unwrap();
assert_eq!(
result.status(),
RunStatus::Completed,
"{:?}",
result.error()
);
assert_eq!(result.step("read-token").output()["token"], "ghp_test");
})
.await
.expect("test timed out");
}
#[cfg(feature = "secret-store")]
#[tokio::test]
async fn with_secret_last_value_wins() {
timeout(Duration::from_secs(10), async {
let result = TestEngine::new()
.with_handler(UsesSecretOp)
.with_secret("git_token", "first")
.with_secret("git_token", "second")
.run(json!({}))
.await
.unwrap();
assert_eq!(result.step("read-token").output()["token"], "second");
})
.await
.expect("test timed out");
}
#[cfg(feature = "secret-store")]
#[tokio::test]
async fn with_secret_is_readable_through_ctx_secrets() {
timeout(Duration::from_secs(10), async {
let mut harness = TestEngine::new()
.with_handler(UsesSecretOp)
.with_secret("k", "v");
harness.run(json!({})).await.unwrap();
let store = harness.store();
let scope = |name: &str| {
ScopedSecretStore::for_workflow(
Uuid::new_v5(&Uuid::NAMESPACE_OID, name.as_bytes()),
store.clone(),
)
};
let own = scope("uses-secret-op").get("k").await.unwrap();
assert_eq!(own.unwrap().value, "v");
assert!(scope("other-workflow").get("k").await.unwrap().is_none());
})
.await
.expect("test timed out");
}
#[cfg(feature = "secret-store")]
#[tokio::test]
#[should_panic(expected = "configure the TestEngine before its first run")]
async fn with_secret_after_first_run_panics() {
let mut harness = TestEngine::new().with_handler(UsesSecretOp);
harness.run(json!({})).await.unwrap();
let _ = harness.with_secret("git_token", "late");
}
}