pub mod cli;
mod crash;
mod held;
mod hop;
#[cfg(test)]
mod tests;
use super::{assembler, child_result, driver, terminal};
use crate::prompt::inbox::{self, ExecutorLock};
use crate::prompt::notice::notice;
use crate::prompt::resolve::WorkerConfig;
use crate::prompt::{Deps, Error, retarget};
use crate::workspace::hold;
use brazen::{Content, Message, Role};
use std::path::Path;
#[derive(Debug)]
pub enum AdvanceOutcome {
AlreadyDriven,
NothingToDo,
Terminal,
ToolsPending(ExecutorLock),
Held,
}
#[derive(Debug, PartialEq, Eq)]
enum Warrant {
ModelCallDue,
NothingDue,
Unpaired,
}
fn warrant(messages: &[Message]) -> Warrant {
let mut unanswered = std::collections::HashSet::new();
for m in messages {
for b in &m.content {
match b {
Content::ToolUse { id, .. } => {
unanswered.insert(id.as_str());
}
Content::ToolResult { tool_use_id, .. } => {
unanswered.remove(tool_use_id.as_str());
}
_ => {}
}
}
}
if !unanswered.is_empty() {
return Warrant::Unpaired;
}
match messages.last() {
Some(m) if matches!(m.role, Role::User | Role::Tool) => Warrant::ModelCallDue,
_ => Warrant::NothingDue,
}
}
fn report_retarget(agent_id: &str, outcome: Option<retarget::Outcome>) {
if let Some(retarget::Outcome::Conflicted(paths)) = outcome {
notice!(
"retarget of [{agent_id}] declined — git could not replay {} \
(marked refs/litany/conflicted/{agent_id}, ARCH §2.6); the branch continues \
on its previous config",
paths.join(", "),
);
}
}
pub(in crate::prompt) fn run(
workspace: &Path,
agent_id: &str,
lease: Option<ExecutorLock>,
deps: &Deps<'_>,
resolve: &mut dyn FnMut() -> Result<WorkerConfig, Error>,
) -> Result<AdvanceOutcome, Error> {
let lock = match lease {
Some(lock) => lock,
None => {
let inbox_dir = inbox::inbox_dir(workspace, agent_id);
match inbox::try_acquire(&inbox_dir).map_err(|source| Error::ExecutorLock {
path: inbox_dir.clone(),
source,
})? {
Some(lock) => lock,
None => return Ok(AdvanceOutcome::AlreadyDriven),
}
}
};
let lock = match hold::read(workspace, agent_id, deps.git) {
Some(mark) => match held::resume(workspace, agent_id, &mark, lock, deps, resolve)? {
held::Resumption::Done(outcome) => return Ok(outcome),
held::Resumption::Stale(lock) => lock,
},
None => lock,
};
crash::settle_crashed_window(workspace, agent_id, deps)?;
let delivery = driver::deliver(workspace, agent_id, deps.git)?;
let seen = delivery.left;
let worktree = crate::workspace::agent_worktree(workspace, agent_id);
if !worktree.exists() {
driver::release_then_reprobe(lock, workspace, agent_id, &seen, deps.launcher);
return Ok(AdvanceOutcome::NothingToDo);
}
report_retarget(
agent_id,
retarget::land(workspace, agent_id, &worktree, deps.git)?,
);
let mut cfg = None;
if child_result::has_pending_result(workspace, agent_id)? {
let resolved = resolve()?;
child_result::interpret_pending(workspace, agent_id, &worktree, &resolved.workflow, deps)?;
cfg = Some(resolved);
}
match warrant(&assembler::transcript(&worktree)?) {
Warrant::NothingDue => {
driver::release_then_reprobe(lock, workspace, agent_id, &seen, deps.launcher);
Ok(AdvanceOutcome::NothingToDo)
}
Warrant::Unpaired => Err(Error::UnpairedToolUse {
branch: agent_id.to_string(),
}),
Warrant::ModelCallDue => {
let cfg = match cfg {
Some(cfg) => cfg,
None => resolve()?,
};
match hop::step(workspace, agent_id, &worktree, &cfg, deps)? {
hop::StepOutcome::ToolsRan => Ok(AdvanceOutcome::ToolsPending(lock)),
hop::StepOutcome::Held => {
driver::release_then_reprobe(lock, workspace, agent_id, &seen, deps.launcher);
Ok(AdvanceOutcome::Held)
}
hop::StepOutcome::Terminal(epitaph) => {
terminal::conclude(
workspace,
agent_id,
epitaph,
&cfg.workflow,
lock,
&seen,
deps,
)?;
Ok(AdvanceOutcome::Terminal)
}
}
}
}
}