2026-07-12 11:28:39 +07:00
|
|
|
//! Workflow engine: interprets `ScriptPrimitive` values (agent, parallel,
|
|
|
|
|
//! pipeline, phase) by spawning subagents, collecting results, and
|
|
|
|
|
//! managing concurrency.
|
|
|
|
|
|
2026-07-11 13:16:10 +07:00
|
|
|
use std::collections::HashMap;
|
2026-07-11 23:45:13 +07:00
|
|
|
use std::sync::{Arc, Mutex};
|
2026-07-11 13:16:10 +07:00
|
|
|
use serde::{Deserialize, Serialize};
|
|
|
|
|
use super::script::{ScriptPrimitive, WorkflowScript};
|
|
|
|
|
|
2026-07-11 23:45:13 +07:00
|
|
|
static FINDINGS: Mutex<Vec<String>> = Mutex::new(Vec::new());
|
2026-07-11 18:23:01 +07:00
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// The lifecycle state of an agent within a workflow run.
|
2026-07-11 13:16:10 +07:00
|
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
|
|
|
|
pub enum AgentState {
|
|
|
|
|
Idle,
|
|
|
|
|
Running,
|
|
|
|
|
Completed,
|
|
|
|
|
Failed,
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// Timestamped status of one workflow agent.
|
2026-07-11 13:16:10 +07:00
|
|
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
|
|
|
|
pub struct AgentStatus {
|
|
|
|
|
pub state: AgentState,
|
|
|
|
|
pub started_at: Option<i64>,
|
|
|
|
|
pub completed_at: Option<i64>,
|
|
|
|
|
pub error: Option<String>,
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// A single agent tracked within a workflow run.
|
2026-07-11 13:16:10 +07:00
|
|
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
|
|
|
|
pub struct WorkflowAgent {
|
|
|
|
|
pub id: String,
|
|
|
|
|
pub name: String,
|
|
|
|
|
pub status: AgentStatus,
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// Orchestrator for running workflow scripts: holds agent roster and a
|
|
|
|
|
/// shared finding accumulator visible to all pipeline stages.
|
2026-07-11 13:16:10 +07:00
|
|
|
#[derive(Debug, Clone)]
|
|
|
|
|
pub struct WorkflowEngine {
|
|
|
|
|
pub agents: Vec<WorkflowAgent>,
|
|
|
|
|
pub findings: Vec<String>,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl WorkflowEngine {
|
2026-07-12 11:28:39 +07:00
|
|
|
/// Create an empty workflow engine with no agents or findings.
|
2026-07-11 13:16:10 +07:00
|
|
|
pub fn new() -> Self {
|
|
|
|
|
WorkflowEngine {
|
|
|
|
|
agents: Vec::new(),
|
|
|
|
|
findings: Vec::new(),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// 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.
|
2026-07-11 23:45:13 +07:00
|
|
|
fn spawn_single_agent(prompt: &str, findings_snapshot: Vec<String>) -> anyhow::Result<String> {
|
|
|
|
|
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())
|
2026-07-12 04:01:10 +07:00
|
|
|
.with_max_steps(usize::MAX);
|
2026-07-11 23:45:13 +07:00
|
|
|
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::<Vec<_>>()
|
|
|
|
|
.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<Vec<String>>);
|
|
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// 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<String>` of all agent outputs (or error strings) in
|
|
|
|
|
/// the order they were submitted.
|
2026-07-11 23:45:13 +07:00
|
|
|
pub fn execute_primitive(
|
|
|
|
|
primitive: &ScriptPrimitive,
|
|
|
|
|
args: &HashMap<String, String>,
|
|
|
|
|
concurrency_cap: usize,
|
2026-07-12 11:55:02 +07:00
|
|
|
continue_on_error: bool,
|
2026-07-11 23:45:13 +07:00
|
|
|
) -> anyhow::Result<Vec<String>> {
|
2026-07-11 13:16:10 +07:00
|
|
|
match primitive {
|
2026-07-11 23:45:13 +07:00
|
|
|
ScriptPrimitive::Agent(prompt) => {
|
|
|
|
|
let resolved = resolve_template(prompt, args);
|
|
|
|
|
let findings_snapshot = FINDINGS.lock().map(|f| f.clone()).unwrap_or_default();
|
2026-07-12 11:55:02 +07:00
|
|
|
match spawn_single_agent(&resolved, findings_snapshot) {
|
|
|
|
|
Ok(text) => Ok(vec![text]),
|
|
|
|
|
Err(e) => {
|
|
|
|
|
if continue_on_error {
|
|
|
|
|
Ok(vec![format!("agent error: {}", e)])
|
|
|
|
|
} else {
|
|
|
|
|
Err(e)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-07-11 13:16:10 +07:00
|
|
|
}
|
2026-07-11 23:45:13 +07:00
|
|
|
|
2026-07-11 13:16:10 +07:00
|
|
|
ScriptPrimitive::Parallel(scripts) => {
|
2026-07-11 23:45:13 +07:00
|
|
|
let semaphore = Arc::new(Semaphore::new(concurrency_cap.max(1)));
|
|
|
|
|
let results: Arc<Mutex<Vec<ParallelResult>>> =
|
|
|
|
|
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();
|
2026-07-12 11:55:02 +07:00
|
|
|
let result = execute_primitive(&script, &args, cap, continue_on_error);
|
2026-07-11 23:45:13 +07:00
|
|
|
if let Ok(mut locked) = results.lock() {
|
|
|
|
|
locked.push((idx, result));
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
})
|
|
|
|
|
.collect();
|
|
|
|
|
|
|
|
|
|
for handle in handles {
|
|
|
|
|
let _ = handle.join();
|
2026-07-11 13:16:10 +07:00
|
|
|
}
|
2026-07-11 23:45:13 +07:00
|
|
|
|
|
|
|
|
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)
|
2026-07-11 13:16:10 +07:00
|
|
|
}
|
2026-07-11 23:45:13 +07:00
|
|
|
|
2026-07-11 13:16:10 +07:00
|
|
|
ScriptPrimitive::Pipeline(scripts) => {
|
2026-07-11 23:45:13 +07:00
|
|
|
let results_store: Arc<Mutex<Vec<Option<Vec<String>>>>> =
|
|
|
|
|
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 || {
|
2026-07-12 11:55:02 +07:00
|
|
|
let result = execute_primitive(&script, &args, cap, continue_on_error);
|
2026-07-11 23:45:13 +07:00
|
|
|
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();
|
2026-07-11 13:16:10 +07:00
|
|
|
}
|
2026-07-11 23:45:13 +07:00
|
|
|
|
|
|
|
|
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)
|
2026-07-11 13:16:10 +07:00
|
|
|
}
|
2026-07-11 23:45:13 +07:00
|
|
|
|
2026-07-11 13:16:10 +07:00
|
|
|
ScriptPrimitive::Phase { name: _name, script } => {
|
2026-07-12 11:55:02 +07:00
|
|
|
execute_primitive(script, args, concurrency_cap, continue_on_error)
|
2026-07-11 13:16:10 +07:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// 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.
|
2026-07-11 13:16:10 +07:00
|
|
|
pub fn run_workflow(script: &WorkflowScript, args: &HashMap<String, String>) -> anyhow::Result<String> {
|
2026-07-11 23:45:13 +07:00
|
|
|
if let Ok(mut findings) = FINDINGS.lock() {
|
|
|
|
|
findings.clear();
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-11 13:16:10 +07:00
|
|
|
let concurrency_cap = if script.options.max_concurrency > 0 {
|
|
|
|
|
script.options.max_concurrency.min(5)
|
|
|
|
|
} else {
|
|
|
|
|
5
|
|
|
|
|
};
|
2026-07-11 23:45:13 +07:00
|
|
|
|
2026-07-12 11:55:02 +07:00
|
|
|
let results = execute_primitive(&script.script, args, concurrency_cap, script.options.continue_on_error)?;
|
2026-07-11 23:45:13 +07:00
|
|
|
|
|
|
|
|
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::<Vec<_>>()
|
|
|
|
|
.join("\n")
|
|
|
|
|
)
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
Ok(summary)
|
2026-07-11 13:16:10 +07:00
|
|
|
}
|
|
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// Add a finding text to the global workflow findings list, making it
|
|
|
|
|
/// visible to sibling agents spawned later in the same run.
|
2026-07-11 13:16:10 +07:00
|
|
|
pub fn note_finding(text: &str) {
|
2026-07-11 18:23:01 +07:00
|
|
|
if let Ok(mut findings) = FINDINGS.lock() {
|
|
|
|
|
findings.push(text.to_string());
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-07-11 23:45:13 +07:00
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// 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.
|
2026-07-11 23:45:13 +07:00
|
|
|
fn resolve_template(template: &str, args: &HashMap<String, String>) -> String {
|
|
|
|
|
let mut result = template.to_string();
|
|
|
|
|
for (key, value) in args {
|
|
|
|
|
result = result.replace(&format!("{{{{{}}}}}", key), value);
|
|
|
|
|
}
|
|
|
|
|
result
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-12 11:28:39 +07:00
|
|
|
/// A counting semaphore built from a `Mutex` + `Condvar`.
|
|
|
|
|
///
|
|
|
|
|
/// Used by `execute_primitive` to cap concurrent parallel branches.
|
2026-07-11 23:45:13 +07:00
|
|
|
struct Semaphore {
|
|
|
|
|
count: Mutex<usize>,
|
|
|
|
|
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();
|
|
|
|
|
}
|
|
|
|
|
}
|