//! Workflow engine: interprets `ScriptPrimitive` values (agent, parallel, //! pipeline, phase) by spawning subagents, collecting results, and //! managing concurrency. use std::collections::HashMap; use std::sync::{Arc, Mutex}; use serde::{Deserialize, Serialize}; use super::script::{ScriptPrimitive, WorkflowScript}; static FINDINGS: Mutex> = Mutex::new(Vec::new()); /// The lifecycle state of an agent within a workflow run. #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] pub enum AgentState { Idle, Running, Completed, Failed, } /// Timestamped status of one workflow agent. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct AgentStatus { pub state: AgentState, pub started_at: Option, pub completed_at: Option, pub error: Option, } /// A single agent tracked within a workflow run. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct WorkflowAgent { pub id: String, pub name: String, pub status: AgentStatus, } /// Orchestrator for running workflow scripts: holds agent roster and a /// shared finding accumulator visible to all pipeline stages. #[derive(Debug, Clone)] pub struct WorkflowEngine { pub agents: Vec, pub findings: Vec, } impl WorkflowEngine { /// Create an empty workflow engine with no agents or findings. pub fn new() -> Self { WorkflowEngine { agents: Vec::new(), findings: Vec::new(), } } } /// Spawn a single synchronous subagent with the given prompt, passing it /// any findings from earlier sibling agents. /// /// Flow: build an `AgentDefinition` -> build a `SubagentContext` -> /// inject findings into the system prompt -> call `run_subagent` on a /// dedicated mpsc channel. /// /// Return: the agent's text output, or an error on failure. fn spawn_single_agent(prompt: &str, findings_snapshot: Vec) -> anyhow::Result { use crate::app::subagent::context::build_subagent_context; use crate::app::subagent::engine::run_subagent; use crate::app::subagent::spawn::AgentDefinition; let def = AgentDefinition::new("workflow-agent".to_string(), "coder".to_string()) .with_max_steps(usize::MAX); let mut ctx = build_subagent_context(def); let findings_section = if findings_snapshot.is_empty() { String::new() } else { format!( "\n\nFindings from sibling agents in this workflow run:\n{}", findings_snapshot .iter() .enumerate() .map(|(i, f)| format!("{}. {}", i + 1, f)) .collect::>() .join("\n") ) }; ctx.system_prompt = format!("{}{}", prompt, findings_section); let (tx, _rx) = tokio::sync::mpsc::channel(32); run_subagent(ctx, tx) } type ParallelResult = (usize, anyhow::Result>); /// Recursively execute a `ScriptPrimitive` tree, respecting an overall /// concurrency cap for parallel branches. /// /// Flow: match the primitive -> /// `Agent` -> `spawn_single_agent` /// `Parallel` -> spawn threads up to `concurrency_cap`, join /// `Pipeline` -> spawn threads sequentially, collect in order /// `Phase` -> recurse (pass-through wrapper) /// /// Why: parallelism is implemented with `std::thread::spawn` and a /// counting semaphore so the main async event loop remains unblocked. /// /// Return: a `Vec` of all agent outputs (or error strings) in /// the order they were submitted. pub fn execute_primitive( primitive: &ScriptPrimitive, args: &HashMap, concurrency_cap: usize, ) -> anyhow::Result> { match primitive { ScriptPrimitive::Agent(prompt) => { let resolved = resolve_template(prompt, args); let findings_snapshot = FINDINGS.lock().map(|f| f.clone()).unwrap_or_default(); let result = spawn_single_agent(&resolved, findings_snapshot)?; Ok(vec![result]) } ScriptPrimitive::Parallel(scripts) => { let semaphore = Arc::new(Semaphore::new(concurrency_cap.max(1))); let results: Arc>> = Arc::new(Mutex::new(Vec::new())); let handles: Vec<_> = scripts .iter() .enumerate() .map(|(idx, script)| { let script = script.clone(); let args = args.clone(); let sem = Arc::clone(&semaphore); let results = Arc::clone(&results); let cap = concurrency_cap; std::thread::spawn(move || { let _permit = sem.acquire(); let result = execute_primitive(&script, &args, cap); if let Ok(mut locked) = results.lock() { locked.push((idx, result)); } }) }) .collect(); for handle in handles { let _ = handle.join(); } let mut locked = results.lock().map_err(|_| anyhow::anyhow!("parallel results lock poisoned"))?; locked.sort_by_key(|(idx, _)| *idx); let mut all = Vec::new(); for (_, res) in locked.drain(..) { match res { Ok(outputs) => all.extend(outputs), Err(e) => all.push(format!("agent error: {}", e)), } } Ok(all) } ScriptPrimitive::Pipeline(scripts) => { let results_store: Arc>>>> = Arc::new(Mutex::new(vec![None; scripts.len()])); let args_arc = Arc::new(args.clone()); let handles: Vec<_> = scripts .iter() .enumerate() .map(|(idx, script)| { let script = script.clone(); let args = Arc::clone(&args_arc); let store = Arc::clone(&results_store); let cap = concurrency_cap; std::thread::spawn(move || { let result = execute_primitive(&script, &args, cap); if let Ok(mut locked) = store.lock() { locked[idx] = Some(result.unwrap_or_else(|e| vec![format!("pipeline stage {} error: {}", idx, e)])); } }) }) .collect(); for handle in handles { let _ = handle.join(); } let locked = results_store.lock().map_err(|_| anyhow::anyhow!("pipeline results lock poisoned"))?; let mut all = Vec::new(); for outputs in locked.iter().flatten() { all.extend(outputs.iter().cloned()); } Ok(all) } ScriptPrimitive::Phase { name: _name, script } => { execute_primitive(script, args, concurrency_cap) } } } /// Run a `WorkflowScript` with the given template arguments and produce a /// summary string. /// /// Flow: clear the global finding store -> cap concurrency to 5 -> call /// `execute_primitive` on the script's root primitive -> format results /// into a one-line-per-agent summary. /// /// Return: a human-readable summary string. pub fn run_workflow(script: &WorkflowScript, args: &HashMap) -> anyhow::Result { if let Ok(mut findings) = FINDINGS.lock() { findings.clear(); } let concurrency_cap = if script.options.max_concurrency > 0 { script.options.max_concurrency.min(5) } else { 5 }; let results = execute_primitive(&script.script, args, concurrency_cap)?; let summary = if results.is_empty() { "workflow completed with no output".to_string() } else { format!( "workflow '{}' completed. {} agent result(s):\n{}", script.name, results.len(), results .iter() .enumerate() .map(|(i, r)| format!("[{}] {}", i + 1, r.lines().next().unwrap_or(r))) .collect::>() .join("\n") ) }; Ok(summary) } /// Add a finding text to the global workflow findings list, making it /// visible to sibling agents spawned later in the same run. pub fn note_finding(text: &str) { if let Ok(mut findings) = FINDINGS.lock() { findings.push(text.to_string()); } } /// Simple template engine: replace `{{key}}` placeholders with values /// from `args`. /// /// Why: a structed template engine is unnecessary for the limited /// use-case; this is intentionally simple and safe. fn resolve_template(template: &str, args: &HashMap) -> String { let mut result = template.to_string(); for (key, value) in args { result = result.replace(&format!("{{{{{}}}}}", key), value); } result } /// A counting semaphore built from a `Mutex` + `Condvar`. /// /// Used by `execute_primitive` to cap concurrent parallel branches. struct Semaphore { count: Mutex, condvar: std::sync::Condvar, } impl Semaphore { fn new(count: usize) -> Self { Semaphore { count: Mutex::new(count), condvar: std::sync::Condvar::new(), } } fn acquire(&self) -> SemaphoreGuard<'_> { let mut count = self.count.lock().unwrap(); while *count == 0 { count = self.condvar.wait(count).unwrap(); } *count -= 1; SemaphoreGuard { sem: self } } } struct SemaphoreGuard<'a> { sem: &'a Semaphore, } impl<'a> Drop for SemaphoreGuard<'a> { fn drop(&mut self) { let mut count = self.sem.count.lock().unwrap(); *count += 1; self.sem.condvar.notify_one(); } }