use std::fmt;
use std::sync::Arc;
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;
use ironflow_store::error::StoreError;
use ironflow_store::memory::InMemoryStore;
use ironflow_store::models::{RunStatus, TriggerKind};
use ironflow_store::store::{RunStore, Store};
use crate::config::{HttpConfig, ShellConfig};
use crate::engine::{Engine, WorkflowResult};
use crate::error::EngineError;
use crate::executor::{ApprovalOutcome, 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>,
}
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 {
Self {
store: Arc::new(InMemoryStore::new()),
handlers: Vec::new(),
primary: None,
provider: None,
decision_provider: None,
mocks: MockInterceptor::new(),
engine: None,
}
}
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_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
}
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()?;
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))
}
}