#![allow(dead_code)]
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use beam_core::{
HostExecutorPrepareResult, WorkflowDispatchOutcome, WorkflowDispatchRun,
WorkflowExecutionHooks,
workflow_definition::{HostExecutorNode, SubagentNode},
};
use serde_json::Value;
pub fn temp_run_dir(label: &str) -> std::path::PathBuf {
std::env::temp_dir().join(format!(
"beam-regression-{label}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos()
))
}
#[derive(Clone)]
pub struct FakeHooks;
#[async_trait]
impl WorkflowExecutionHooks for FakeHooks {
async fn execute_subagent(
&mut self,
_ctx: WorkflowDispatchRun<'_>,
_node: &SubagentNode,
resolved_prompt: String,
) -> anyhow::Result<WorkflowDispatchOutcome> {
Ok(WorkflowDispatchOutcome::Succeeded {
output: Value::String(resolved_prompt),
session: None,
})
}
async fn execute_host_executor(
&mut self,
_ctx: WorkflowDispatchRun<'_>,
_node: &HostExecutorNode,
resolved_input: Value,
) -> anyhow::Result<WorkflowDispatchOutcome> {
Ok(WorkflowDispatchOutcome::Succeeded {
output: resolved_input,
session: None,
})
}
}
#[derive(Clone)]
pub struct SpyHooks {
pub prepare_called: Arc<Mutex<bool>>,
pub execute_called: Arc<Mutex<bool>>,
}
impl SpyHooks {
pub fn new() -> Self {
Self {
prepare_called: Arc::new(Mutex::new(false)),
execute_called: Arc::new(Mutex::new(false)),
}
}
}
#[async_trait]
impl WorkflowExecutionHooks for SpyHooks {
async fn execute_subagent(
&mut self,
_ctx: WorkflowDispatchRun<'_>,
_node: &SubagentNode,
resolved_prompt: String,
) -> anyhow::Result<WorkflowDispatchOutcome> {
Ok(WorkflowDispatchOutcome::Succeeded {
output: Value::String(resolved_prompt),
session: None,
})
}
async fn execute_host_executor(
&mut self,
_ctx: WorkflowDispatchRun<'_>,
_node: &HostExecutorNode,
parsed_input: Value,
) -> anyhow::Result<WorkflowDispatchOutcome> {
*self.execute_called.lock().unwrap() = true;
Ok(WorkflowDispatchOutcome::Succeeded {
output: parsed_input,
session: None,
})
}
fn prepare_host_executor(
&self,
_executor_name: &str,
resolved_input: &Value,
) -> anyhow::Result<HostExecutorPrepareResult> {
*self.prepare_called.lock().unwrap() = true;
Ok(HostExecutorPrepareResult {
parsed_input: resolved_input.clone(),
canonical_input: resolved_input.clone(),
provider: "test-provider".to_string(),
idempotency_ttl_ms: 42_000,
})
}
}
#[derive(Clone)]
pub struct FailingPrepareHooks {
pub execute_called: Arc<Mutex<bool>>,
}
impl FailingPrepareHooks {
pub fn new() -> Self {
Self {
execute_called: Arc::new(Mutex::new(false)),
}
}
}
#[async_trait]
impl WorkflowExecutionHooks for FailingPrepareHooks {
async fn execute_subagent(
&mut self,
_ctx: WorkflowDispatchRun<'_>,
_node: &SubagentNode,
resolved_prompt: String,
) -> anyhow::Result<WorkflowDispatchOutcome> {
Ok(WorkflowDispatchOutcome::Succeeded {
output: Value::String(resolved_prompt),
session: None,
})
}
async fn execute_host_executor(
&mut self,
_ctx: WorkflowDispatchRun<'_>,
_node: &HostExecutorNode,
parsed_input: Value,
) -> anyhow::Result<WorkflowDispatchOutcome> {
*self.execute_called.lock().unwrap() = true;
Ok(WorkflowDispatchOutcome::Succeeded {
output: parsed_input,
session: None,
})
}
fn prepare_host_executor(
&self,
_executor_name: &str,
_resolved_input: &Value,
) -> anyhow::Result<HostExecutorPrepareResult> {
anyhow::bail!("prepare_host_executor forced failure")
}
}
#[derive(Clone)]
pub struct PanicHooks;
#[async_trait]
impl WorkflowExecutionHooks for PanicHooks {
async fn execute_subagent(
&mut self,
_ctx: WorkflowDispatchRun<'_>,
_node: &SubagentNode,
resolved_prompt: String,
) -> anyhow::Result<WorkflowDispatchOutcome> {
Ok(WorkflowDispatchOutcome::Succeeded {
output: Value::String(resolved_prompt),
session: None,
})
}
async fn execute_host_executor(
&mut self,
_ctx: WorkflowDispatchRun<'_>,
_node: &HostExecutorNode,
_resolved_input: Value,
) -> anyhow::Result<WorkflowDispatchOutcome> {
anyhow::bail!("simulated executor failure")
}
}