mod flush;
mod landing;
mod proposal;
mod verifier;
pub(super) use flush::run_flush;
use super::{drain, transcript, transfer};
use crate::config::{Action, Event, Workflow};
use crate::prompt::{Deps, Error, compactor, inbox, reviewer, role};
use crate::template::GitRunner;
use std::path::{Path, PathBuf};
pub(super) struct ChildResult {
pub(super) child_id: String,
pub(super) terminal_ref: String,
pub(super) epitaph: String,
pub(super) response: Option<String>,
pub(super) path: PathBuf,
}
pub(super) fn own_result_ref(agent_id: &str, sender: &str, body: &str) -> Option<String> {
if inbox::parent_of(sender).as_deref() != Some(agent_id) {
return None;
}
transfer::terminal_ref_of(body)
}
pub(super) fn has_pending_result(workspace: &Path, agent_id: &str) -> Result<bool, Error> {
let dir = inbox::inbox_dir(workspace, agent_id);
for msg in drain::pending(&dir)? {
if own_result_ref(agent_id, &msg.sender, &read_body(&msg.path)?).is_some() {
return Ok(true);
}
}
Ok(false)
}
pub(super) fn interpret_pending(
workspace: &Path,
agent_id: &str,
worktree: &Path,
workflow: &Workflow,
deps: &Deps<'_>,
) -> Result<(), Error> {
let dir = inbox::inbox_dir(workspace, agent_id);
let results = load_results(agent_id, &dir)?;
let events: Vec<Event> = results
.iter()
.map(|cr| child_event(worktree, cr, deps.git))
.collect::<Result<_, _>>()?;
for (cr, &event) in results.iter().zip(&events) {
if matches!(event, Event::VerifierApprove | Event::VerifierReject) {
for action in child_actions(workflow, event) {
verifier::execute(
&action, event, workspace, agent_id, worktree, cr, &results, deps,
)?;
}
}
}
for (cr, &event) in results.iter().zip(&events) {
if matches!(event, Event::VerifierApprove | Event::VerifierReject) || !cr.path.exists() {
continue;
}
for action in child_actions(workflow, event) {
execute_child(
&action, event, workspace, agent_id, worktree, cr, workflow, deps,
)?;
}
}
Ok(())
}
fn load_results(agent_id: &str, dir: &Path) -> Result<Vec<ChildResult>, Error> {
let mut out = Vec::new();
for msg in drain::pending(dir)? {
let body = read_body(&msg.path)?;
let Some(terminal_ref) = own_result_ref(agent_id, &msg.sender, &body) else {
continue;
};
let (epitaph, response) = split_frontmatter(&body);
out.push(ChildResult {
child_id: msg.sender,
terminal_ref,
epitaph,
response,
path: msg.path,
});
}
Ok(out)
}
fn read_body(path: &Path) -> Result<String, Error> {
std::fs::read_to_string(path).map_err(Error::Io)
}
pub(super) fn split_frontmatter(body: &str) -> (String, Option<String>) {
let mut lines = body.lines();
if lines.next() != Some("---") {
return (String::new(), None);
}
let mut epitaph = String::new();
for line in lines.by_ref() {
if line == "---" {
break;
}
if let Some(v) = line.strip_prefix("epitaph:") {
epitaph = v.trim().to_string();
}
}
let rest = lines.collect::<Vec<_>>().join("\n");
let response = (!rest.trim().is_empty()).then_some(rest);
(epitaph, response)
}
fn child_event(worktree: &Path, cr: &ChildResult, git: &dyn GitRunner) -> Result<Event, Error> {
let derived = role::derive(worktree, &cr.terminal_ref, &cr.child_id, git)?;
Ok(match derived.as_deref() {
Some(compactor::COMPACTOR_ROLE) => Event::CompactorReturn,
Some(reviewer::REVIEWER_ROLE) => Event::ReviewerReturn,
Some(verifier::VERIFIER_ROLE) => verifier::verdict(cr),
_ => Event::WorkerReturn,
})
}
fn child_actions(workflow: &Workflow, event: Event) -> Vec<Action> {
let bound = workflow.actions_for(event);
if !bound.is_empty() {
return bound;
}
match event {
Event::CompactorReturn => vec![Action::LandCompaction],
Event::ReviewerReturn => vec![Action::StageProposal],
Event::VerifierReject => vec![Action::Dispatch {
role: crate::prompt::WORKER_ROLE.to_string(),
with: Some(verifier::FEEDBACK.to_string()),
mode: None,
}],
_ => vec![Action::DeliverResult],
}
}
#[allow(clippy::too_many_arguments)] fn execute_child(
action: &Action,
event: Event,
workspace: &Path,
agent_id: &str,
worktree: &Path,
cr: &ChildResult,
workflow: &Workflow,
deps: &Deps<'_>,
) -> Result<(), Error> {
match action {
Action::Dispatch { role, .. } if role == verifier::VERIFIER_ROLE => {
verifier::dispatch(workspace, agent_id, worktree, cr, deps)
}
Action::GateReturnOn { .. } => Ok(()),
Action::DeliverResult => deliver_result(worktree, agent_id, cr, deps.git),
Action::LandCompaction if landing::qualifies(cr) => {
landing::land(worktree, agent_id, cr, workflow, deps.git)
}
Action::StageProposal if landing::qualifies(cr) => {
proposal::stage(workspace, worktree, cr, deps)
}
Action::LandCompaction | Action::StageProposal => {
deliver_result(worktree, agent_id, cr, deps.git)
}
other => Err(Error::ActionUnsupported {
action: format!("{other:?}"),
event: event.as_str(),
}),
}
}
pub(super) fn deliver_result(
worktree: &Path,
agent_id: &str,
cr: &ChildResult,
git: &dyn GitRunner,
) -> Result<(), Error> {
transfer::apply(worktree, &cr.child_id, &cr.terminal_ref, git)?;
transcript::deliver_message(worktree, agent_id, &cr.child_id, &cr.path, git)
}
#[cfg(test)]
mod tests;