use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, Instant};
use anyhow::{Context as _, Result, bail};
use jiff::Timestamp;
use tokio::sync::Semaphore;
use crate::agent::{self, AgentOutput, Invocation, SeatState};
use crate::ask;
use crate::blind;
use crate::bump;
use crate::config::{
AgentSpec, Config, IncompleteReviewPolicy, LeakPolicy, MergeMode, MergeStyle, Prompts,
ResolvedRoles,
};
use crate::git;
use crate::land;
use crate::proc::Quiet as _;
use crate::prompt::{
self, CandidateView, Lens, ReviewPatch, ReviewReconsiderCtx, ReviewSeatReport, Turn,
};
use crate::run::{
BaseSync, Candidate, CommandOutcome, DeliberationRound, DeliberationTurn, FixRecord, Judgement,
MergeOutcome, QuotaLoss, ReviewRecord, ReviewRevoteRecord, ReviewRound, RunState, RunStatus,
Tally, VoteRecord, tail, write_artifact,
};
use crate::verdict::{
self, FinalVote, FixReport, Position, Ranking, Review, ReviewRevote, ReviewVote, Severity,
};
const OUTPUT_TAIL: usize = 8_000;
const EVENT_OUTPUT_TAIL: usize = 2_000;
pub(crate) const STAGNANT_LIMIT: usize = 2;
const BASE_SYNC_ROUNDS: usize = 4;
#[derive(Clone)]
struct SeatJob {
spec: AgentSpec,
seat: SeatState,
cwd: PathBuf,
prompt: String,
timeout: Duration,
allow_write: bool,
sessions: bool,
artifacts: PathBuf,
stem: String,
}
enum AgentOutcome {
Ok(AgentOutput),
Quota(AgentOutput),
Dropped(AgentOutput),
Failed(String),
}
#[derive(Debug, Clone, Default)]
pub struct Pause(Arc<AtomicBool>);
impl Pause {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn park(&self) {
self.0.store(true, Ordering::SeqCst);
}
#[must_use]
pub fn parked(&self) -> bool {
self.0.load(Ordering::SeqCst)
}
}
pub struct Runner {
pub state: RunState,
roles: ResolvedRoles,
sem: Arc<Semaphore>,
pause: Pause,
}
async fn resolve_base(repo: &Path, base_branch: &str, remote: &str) -> Result<String> {
let tracking = format!("{remote}/{base_branch}");
let fetched = git::fetch(repo, remote, base_branch).await;
if let Ok(out) = &fetched
&& out.ok()
&& git::rev_exists(repo, &tracking).await
{
return git::rev_parse(repo, &tracking).await;
}
let why = match &fetched {
Ok(out) if !out.ok() => out.stderr.lines().next().unwrap_or("").to_owned(),
Ok(_) => format!("{remote} has no {base_branch}"),
Err(e) => e.to_string(),
};
tracing::warn!(
"could not read {tracking} ({why}); branching off the local \
{base_branch} instead, which may be behind"
);
git::rev_parse(repo, base_branch).await.with_context(|| {
format!(
"cannot resolve `{base_branch}`; set [merge] base in magi.toml to a \
branch that exists"
)
})
}
impl Runner {
pub async fn start(repo: &Path, instruction: String, config: Config) -> Result<Self> {
let repo = git::toplevel(repo).await?;
let missing = agent::missing_programs(&config.agents);
if !missing.is_empty() {
bail!(
"these agent programs are not on PATH: {}. Fix the roster in \
magi.toml or install them.",
missing.join(", ")
);
}
let base_branch = match config.merge.base.clone() {
Some(b) => b,
None => git::current_branch(&repo)
.await?
.context("HEAD is detached; set [merge] base in magi.toml")?,
};
let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
if !git::is_clean(&repo).await? {
tracing::warn!(
"{} has uncommitted changes; they are not part of this run, \
which branches off {base_branch} ({})",
repo.display(),
&base_commit[..base_commit.len().min(8)]
);
}
let roles = config.resolve_roles()?;
let max_parallel = config.graph.max_parallel.max(1);
let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
state.event("start", format!("run {} created", state.id));
state.save()?;
Ok(Self {
state,
roles,
sem: Arc::new(Semaphore::new(max_parallel)),
pause: Pause::new(),
})
}
pub async fn review(repo: &Path, branch: &str, config: Config) -> Result<Self> {
let repo = git::toplevel(repo).await?;
let missing = agent::missing_programs(&config.agents);
if !missing.is_empty() {
bail!(
"these agent programs are not on PATH: {}. Fix the roster in \
magi.toml or install them.",
missing.join(", ")
);
}
if !git::branch_exists(&repo, branch).await? {
bail!("no branch `{branch}` in {}", repo.display());
}
let base_branch = match config.merge.base.clone() {
Some(b) => b,
None => git::current_branch(&repo)
.await?
.context("HEAD is detached; set [merge] base in magi.toml")?,
};
if base_branch == branch {
bail!("`{branch}` is the base branch; there is nothing to review against");
}
let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
let roles = config.resolve_roles()?;
let max_parallel = config.graph.max_parallel.max(1);
let log = git::log_oneline(&repo, &base_commit, branch)
.await
.unwrap_or_default();
let instruction = format!(
"Review the work already on branch `{branch}`. There is no task \
statement: what the change claims to do is whatever its commits \
say.\n\n{}",
if log.trim().is_empty() {
"(no commit messages)"
} else {
log.trim()
}
);
let mut state = RunState::new(
repo.clone(),
base_branch,
base_commit.clone(),
instruction,
config,
);
let worktree = state.worktree_root().join("under-review");
if let Some(parent) = worktree.parent() {
tokio::fs::create_dir_all(parent).await.ok();
}
let path = worktree.to_string_lossy().to_string();
git::git(&repo, &["worktree", "add", &path, branch])
.await
.with_context(|| {
format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
})?;
let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
.await
.unwrap_or(0);
if commits == 0 {
bail!("`{branch}` has no commits beyond {}", short(&base_commit));
}
let files = git::changed_files(&worktree, &base_commit, "HEAD")
.await
.map(|f| f.len())
.unwrap_or(0);
let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
.await
.unwrap_or_default();
state.candidates.push(Candidate {
index: 0,
label: 'A',
agent: "(existing branch)".to_owned(),
branch: branch.to_owned(),
worktree,
summary: String::new(),
stat,
files,
commits,
empty: false,
failed: None,
duration_ms: 0,
folded: false,
});
state.tally = Some(Tally {
first_choice: BTreeMap::from([('A', 0)]),
borda: BTreeMap::new(),
winner: 'A',
rankings: 0,
unanimous_initial: false,
deliberated: false,
changed_votes: 0,
unanimous_final: false,
tie_break: None,
judges: 0,
present: 0,
quorum: 0,
met_quorum: true,
uncontested: Some("review-only run: nothing competed".to_owned()),
});
state.status = RunStatus::Reviewing;
state.event(
"start",
format!(
"review-only run {} on `{branch}` ({files} files, {commits} commits)",
state.id
),
);
state.save()?;
Ok(Self {
state,
roles,
sem: Arc::new(Semaphore::new(max_parallel)),
pause: Pause::new(),
})
}
pub fn resume(id: &str) -> Result<Self> {
let state = RunState::load(id)?;
let roles = state.config.resolve_roles()?;
let max_parallel = state.config.graph.max_parallel.max(1);
Ok(Self {
state,
roles,
sem: Arc::new(Semaphore::new(max_parallel)),
pause: Pause::new(),
})
}
pub async fn execute(&mut self) -> Result<()> {
self.state.parked = false;
if self.state.clear_active() {
self.state.save()?;
}
if self.state.status == RunStatus::Stalled {
if self.recover_stall().await? {
self.finish_after_tally().await?;
} else {
self.state.save()?;
}
return Ok(());
}
if self.state.status == RunStatus::Landing {
self.run_land().await?;
self.settle_questions();
return Ok(());
}
self.prep().await?;
if self.park_here()? {
return Ok(());
}
self.implement().await?;
if self.park_here()? {
return Ok(());
}
self.judge().await?;
if self.park_here()? {
return Ok(());
}
self.deliberate().await?;
if self.park_here()? {
return Ok(());
}
self.vote().await?;
if self.park_here()? {
return Ok(());
}
self.tally()?;
if self.state.status == RunStatus::Stalled {
self.state.save()?;
return Ok(());
}
self.finish_after_tally().await?;
Ok(())
}
fn park_here(&mut self) -> Result<bool> {
if !self.pause.parked() {
return Ok(false);
}
self.state.event(
"park",
format!(
"parked after `{}` — resume to carry on from here",
self.state.status.as_str()
),
);
self.state.parked = true;
self.state.save()?;
Ok(true)
}
pub fn on_pause(&mut self, pause: Pause) {
self.pause = pause;
}
fn settle_questions(&mut self) {
if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
}
}
async fn finish_after_tally(&mut self) -> Result<()> {
self.fold_losers().await?;
self.sync_to_base().await?;
self.review_loop().await?;
self.sync_to_base().await?;
self.gate().await?;
self.merge().await?;
self.state.save()?;
Ok(())
}
async fn prep(&mut self) -> Result<()> {
if !self.state.candidates.is_empty() {
return Ok(());
}
self.state.status = RunStatus::Prep;
let repo = self.state.repo.clone();
let base = self.state.base_commit.clone();
let root = self.state.worktree_root();
let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
let hooks_dir = self.state.dir().join("hooks");
if self.state.config.blind.commit_msg_hook {
std::fs::create_dir_all(&hooks_dir)
.with_context(|| format!("create {}", hooks_dir.display()))?;
let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
let path = hooks_dir.join("commit-msg");
std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
make_executable(&path)?;
git::acquire_worktree_config(&repo).await?;
self.state.enabled_worktree_config = true;
}
for (index, (spec, label)) in self
.roles
.implementers
.clone()
.into_iter()
.zip(labels)
.enumerate()
{
let branch = self.state.branch_for(label);
let worktree = root.join(format!("cand-{label}"));
git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
if self.state.config.blind.commit_msg_hook {
git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
}
git::local_exclude(&worktree, "/.magi/").await?;
self.state.candidates.push(Candidate {
index,
label,
agent: spec.id.clone(),
branch,
worktree,
summary: String::new(),
stat: String::new(),
files: 0,
commits: 0,
empty: false,
failed: None,
duration_ms: 0,
folded: false,
});
}
for j in 1..=self.roles.judges.len() {
let wt = root.join(format!("judge-{j}"));
if !wt.exists() {
git::worktree_add_detached(&repo, &wt, &base).await?;
}
}
let authors: Vec<&str> = self
.roles
.implementers
.iter()
.map(|a| a.id.as_str())
.collect();
let overlap: Vec<String> = self
.roles
.judges
.iter()
.enumerate()
.filter(|(_, j)| authors.contains(&j.id.as_str()))
.map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
.collect();
if !overlap.is_empty() {
let note = format!(
"{} also authored a candidate; blind, but the panel is less \
independent than {} distinct agents would be",
overlap.join(", "),
self.roles.judges.len()
);
self.state.event("prep", note);
}
self.state.event(
"prep",
format!(
"{} candidates, {} judges, base {} ({})",
self.state.candidates.len(),
self.roles.judges.len(),
&self.state.base_commit[..7.min(self.state.base_commit.len())],
self.state.base_branch
),
);
self.state.status = RunStatus::Implementing;
self.state.save()?;
Ok(())
}
async fn implement(&mut self) -> Result<()> {
let run_id = self.state.id.clone();
let prompts = self.state.config.prompts.clone();
let todo: Vec<usize> = self
.state
.candidates
.iter()
.enumerate()
.filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
.map(|(i, _)| i)
.collect();
if todo.is_empty() {
return self.after_implement();
}
self.state.status = RunStatus::Implementing;
let language = self.state.config.graph.language.clone();
let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
let sessions = self.state.config.graph.sessions;
let artifacts = agent::artifacts_dir(&self.state.dir());
let mut jobs = Vec::new();
for &i in &todo {
let (index, label, worktree) = {
let c = &self.state.candidates[i];
(c.index, c.label, c.worktree.clone())
};
let spec = self.roles.implementers[index].clone();
let seat_key = format!("impl-{label}");
let seat = self.seat(&seat_key, &spec.id);
let instruction = self.state.instruction.clone();
jobs.push(SeatJob {
spec,
seat,
prompt: prompt::implement(&instruction, &worktree.to_string_lossy(), &language),
cwd: worktree,
timeout,
allow_write: true,
sessions,
artifacts: artifacts.clone(),
stem: format!("impl-{label}"),
});
}
self.state.event(
"implement",
format!("{} candidates in parallel", jobs.len()),
);
let sent = jobs.clone();
let cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: &run_id,
node: "implement",
prompts: &prompts,
cache: cache.as_deref(),
};
let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
.await;
for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
let seat_key = seat.key.clone();
self.state.seats.insert(seat.key.clone(), seat);
let label = self.state.candidates[i].label;
let worktree = self.state.candidates[i].worktree.clone();
let base = self.state.base_commit.clone();
let (summary, duration, failed) = match out {
AgentOutcome::Ok(o) => {
let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
let failed = (!o.usable()).then(|| {
if o.timed_out {
"agent timed out".to_owned()
} else {
format!("agent exited with {:?}", o.exit_code)
}
});
(text, o.duration_ms, failed)
}
AgentOutcome::Dropped(o) => {
let why = o
.dropped
.as_ref()
.map(|d| d.why.as_str())
.unwrap_or("the CLI ended the stream without delivering its answer");
(
String::new(),
o.duration_ms,
Some(format!("the CLI dropped the stream ({why})")),
)
}
AgentOutcome::Quota(o) => {
self.state.quota.push(QuotaLoss {
seat: seat_key,
node: "implement".to_owned(),
at: Timestamp::now(),
reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
});
(
String::new(),
o.duration_ms,
Some("rate limited (quota); produced no change".to_owned()),
)
}
AgentOutcome::Failed(e) => (String::new(), 0, Some(e)),
};
let rescued = git::commit_all(
&worktree,
&format!("magi: candidate {label} (uncommitted work)"),
)
.await
.unwrap_or(false);
let commits = git::commits_ahead(&worktree, &base, "HEAD")
.await
.unwrap_or(0);
let patch = git::diff(&worktree, &base, "HEAD")
.await
.unwrap_or_default();
let stat = git::diff_stat(&worktree, &base, "HEAD")
.await
.unwrap_or_default();
let files = git::changed_files(&worktree, &base, "HEAD")
.await
.map(|f| f.len())
.unwrap_or(0);
write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
let c = &mut self.state.candidates[i];
c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
c.stat = stat;
c.files = files;
c.commits = commits;
c.duration_ms = duration;
c.empty = commits == 0 || patch.trim().is_empty();
c.failed = match failed {
Some(_) if c.empty => failed,
_ => None,
};
let note = match (&c.failed, c.empty, rescued) {
(Some(e), _, _) => format!("candidate {label}: {e}"),
(None, true, _) => format!("candidate {label}: no change produced"),
(None, false, true) => {
format!(
"candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
)
}
(None, false, false) => {
format!("candidate {label}: {files} files, {commits} commits")
}
};
self.state.event("implement", note);
self.state.save()?;
}
self.after_implement()
}
async fn resume_undelivered(
&mut self,
results: &mut [(usize, SeatState, AgentOutcome)],
sent: &[SeatJob],
prompts: &Prompts,
run_id: &str,
) {
for (wi, seat, out) in results.iter_mut() {
let Some(dropped) = (match &*out {
AgentOutcome::Dropped(o) => o.dropped.clone(),
_ => None,
}) else {
continue;
};
let Some(job) = sent.get(*wi) else { continue };
if !git::is_clean(&job.cwd).await.unwrap_or(true) {
self.state.event(
"implement",
format!(
"{}: the CLI dropped the stream after {} output tokens ({}), but the \
work is in the tree",
seat.key, dropped.output_tokens, dropped.why
),
);
continue;
}
if !has_context(&job.spec, seat, job.sessions) {
self.state.event(
"implement",
format!(
"{}: the CLI dropped the stream after {} output tokens ({}), but there \
is no session left to resume",
seat.key, dropped.output_tokens, dropped.why
),
);
continue;
}
self.state.event(
"implement",
format!(
"{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
conversation",
seat.key, dropped.output_tokens, dropped.why
),
);
let mut retry = job.clone();
retry.seat = seat.clone();
retry.prompt = prompt::resume_after_drop(&dropped.why);
retry.timeout = retry_budget(job.timeout, true);
retry.stem = format!("{}-resume", job.stem);
let cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: run_id,
node: "implement",
prompts,
cache: cache.as_deref(),
};
let (resumed_seat, resumed) =
run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
*seat = resumed_seat;
*out = resumed;
}
}
fn after_implement(&mut self) -> Result<()> {
if self.state.leaks.is_empty() {
let cfg = self.state.config.blind.clone();
let mut leaks = Vec::new();
for c in &self.state.candidates {
let Some(patch) =
crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
else {
continue;
};
leaks.extend(blind::scan(
&format!("candidate {} patch", c.label),
&patch,
&cfg.vendor_tokens,
));
}
if !leaks.is_empty() {
let summary = leaks
.iter()
.map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
.collect::<Vec<_>>()
.join(", ");
match cfg.on_leak {
LeakPolicy::Fail => {
self.state.status = RunStatus::Failed;
self.state
.event("blind", format!("vendor text in a patch: {summary}"));
self.state.leaks = leaks;
self.state.save()?;
self.settle_questions();
bail!(
"blind.on_leak = \"fail\" and vendor text reached a \
judged patch: {summary}"
);
}
LeakPolicy::Redact => self.state.event(
"blind",
format!("redacting vendor text for judging: {summary}"),
),
LeakPolicy::Warn => self.state.event(
"blind",
format!("vendor text present in a judged patch (shown as-is): {summary}"),
),
}
self.state.leaks = leaks;
}
}
if self.state.viable().is_empty() {
self.state.status = RunStatus::Failed;
self.state.save()?;
self.settle_questions();
bail!("no candidate produced a change; nothing to judge");
}
self.state.status = RunStatus::Judging;
self.state.save()?;
Ok(())
}
async fn judge(&mut self) -> Result<()> {
let run_id = self.state.id.clone();
let prompts = self.state.config.prompts.clone();
if !self.state.judgements.is_empty() || self.state.judge_skipped {
return Ok(());
}
let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
if viable.len() == 1 {
self.state.judge_skipped = true;
self.state.event(
"judge",
format!(
"only candidate {} produced a change; judging skipped",
viable[0].label
),
);
self.state.save()?;
return Ok(());
}
self.state.status = RunStatus::Judging;
let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
let language = self.state.config.graph.language.clone();
let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
let sessions = self.state.config.graph.sessions;
let artifacts = agent::artifacts_dir(&self.state.dir());
let root = self.state.worktree_root();
let base_short = short(&self.state.base_commit);
let mut jobs = Vec::new();
let mut orders = Vec::new();
for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
let order = blind::presentation_order(viable.len(), j, self.state.seed);
let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
let seat_key = format!("judge-{}", j + 1);
let seat = self.seat(&seat_key, &spec.id);
jobs.push(SeatJob {
prompt: prompt::judge(
&self.state.instruction,
&views,
self.roles.judges.len(),
&base_short,
&language,
),
spec,
seat,
cwd: root.join(format!("judge-{}", j + 1)),
timeout,
allow_write: false,
sessions,
artifacts: artifacts.clone(),
stem: format!("judge-{}", j + 1),
});
}
self.state.event(
"judge",
format!(
"{} judges ranking {} candidates blind",
jobs.len(),
viable.len()
),
);
let labels_for_check = labels.clone();
let mut quota_losses = Vec::new();
let cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: &run_id,
node: "judge",
prompts: &prompts,
cache: cache.as_deref(),
};
let results = ask_json_wave::<Ranking>(
jobs,
Arc::clone(&self.sem),
self.state.config.graph.retries,
&ctx,
&mut quota_losses,
&mut self.state,
&move |r: &Ranking| r.validate(&labels_for_check),
)
.await;
self.state.quota.extend(quota_losses);
for (j, (seat, res)) in results.into_iter().enumerate() {
let agent_id = seat.agent.clone();
self.state.seats.insert(seat.key.clone(), seat);
let mut record = Judgement {
judge: j + 1,
seat: format!("judge-{}", j + 1),
agent: agent_id,
ranking: Vec::new(),
reasons: BTreeMap::new(),
confidence: None,
order: orders[j].clone(),
failed: None,
duration_ms: 0,
};
match res {
Ok((ranking, out)) => {
record.ranking = ranking.normalized();
record.reasons = ranking.reasons;
record.confidence = ranking.confidence;
record.duration_ms = out.duration_ms;
self.state.event(
"judge",
format!(
"judge {} ranked {}",
j + 1,
record.ranking.iter().collect::<String>()
),
);
}
Err(e) => {
record.failed = Some(e.to_string());
self.state
.event("judge", format!("judge {} produced no ranking: {e}", j + 1));
}
}
self.state.judgements.push(record);
self.state.save()?;
}
Ok(())
}
async fn deliberate(&mut self) -> Result<()> {
let run_id = self.state.id.clone();
let prompts = self.state.config.prompts.clone();
if !self.state.deliberation.is_empty() {
return Ok(());
}
let tops: Vec<char> = self
.state
.judgements
.iter()
.filter_map(|j| j.ranking.first().copied())
.collect();
let rounds = self.state.config.graph.deliberate_rounds;
if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
self.state.event(
"deliberate",
format!("judges agreed on {} outright; no deliberation", tops[0]),
);
}
self.state.status = RunStatus::Voting;
self.state.save()?;
return Ok(());
}
self.state.status = RunStatus::Deliberating;
self.state.event(
"deliberate",
format!(
"split: first choices were {} — opening {rounds} round(s)",
tops.iter().collect::<String>()
),
);
let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
let language = self.state.config.graph.language.clone();
let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
let sessions = self.state.config.graph.sessions;
let artifacts = agent::artifacts_dir(&self.state.dir());
let root = self.state.worktree_root();
let base_short = short(&self.state.base_commit);
for round in 1..=rounds {
let mut turns: Vec<DeliberationTurn> = Vec::new();
for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
if self.state.judgements[j].failed.is_some() {
continue;
}
let seat_key = format!("judge-{}", j + 1);
let mut seat = self.seat(&seat_key, &spec.id);
let transcript = self.transcript(&turns, j);
let context = if has_context(&spec, &seat, sessions) {
None
} else {
Some(self.candidate_block(&viable, &base_short))
};
let text = prompt::deliberate(
&self.state.instruction,
context.as_deref(),
&transcript,
round,
rounds,
&language,
);
let job = SeatJob {
spec,
seat: seat.clone(),
prompt: text,
cwd: root.join(format!("judge-{}", j + 1)),
timeout,
allow_write: false,
sessions,
artifacts: artifacts.clone(),
stem: format!("delib-{round}-judge-{}", j + 1),
};
let cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: &run_id,
node: "deliberate",
prompts: &prompts,
cache: cache.as_deref(),
};
let (updated, out) =
run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
seat = updated;
let agent_id = seat.agent.clone();
let seat_key = seat.key.clone();
self.state.seats.insert(seat.key.clone(), seat);
let body = match out {
AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
AgentOutcome::Dropped(o) => {
let why =
o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
"the CLI ended the stream without delivering its answer",
);
self.state.event(
"deliberate",
format!(
"judge {} skipped: the CLI dropped the stream ({why})",
j + 1
),
);
continue;
}
AgentOutcome::Quota(o) => {
self.state.quota.push(QuotaLoss {
seat: seat_key,
node: "deliberate".to_owned(),
at: Timestamp::now(),
reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
});
self.state.event(
"deliberate",
format!("judge {} skipped: rate limited (quota)", j + 1),
);
continue;
}
AgentOutcome::Failed(e) => {
self.state
.event("deliberate", format!("judge {} skipped: {e}", j + 1));
continue;
}
};
let tentative = verdict::extract_json::<Position>(&body)
.ok()
.and_then(|p| p.tentative)
.and_then(|s| s.trim().chars().next())
.map(|c| c.to_ascii_uppercase());
self.state.event(
"deliberate",
format!(
"round {round}: judge {} now favours {}",
j + 1,
tentative.map_or("—".to_owned(), |c| c.to_string())
),
);
turns.push(DeliberationTurn {
judge: j + 1,
agent: agent_id,
body: blind::sanitize_prose(&body, &self.state.config.blind),
tentative,
});
}
self.state
.deliberation
.push(DeliberationRound { round, turns });
self.state.save()?;
}
self.state.status = RunStatus::Voting;
self.state.save()?;
Ok(())
}
async fn vote(&mut self) -> Result<()> {
let run_id = self.state.id.clone();
let prompts = self.state.config.prompts.clone();
if !self.state.votes.is_empty() {
return Ok(());
}
let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
if viable.len() == 1 {
return Ok(());
}
self.state.status = RunStatus::Voting;
let language = self.state.config.graph.language.clone();
let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
let sessions = self.state.config.graph.sessions;
let artifacts = agent::artifacts_dir(&self.state.dir());
let root = self.state.worktree_root();
let base_short = short(&self.state.base_commit);
let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
let mut jobs = Vec::new();
let mut seats_at = Vec::new();
for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
if self
.state
.judgements
.get(j)
.is_some_and(|r| r.failed.is_some())
{
continue;
}
let seat_key = format!("judge-{}", j + 1);
let seat = self.seat(&seat_key, &spec.id);
let mut text = prompt::final_vote(&viable, &language);
if !has_context(&spec, &seat, sessions) {
text = format!(
"{}\n\n# Candidates\n\n{}",
text,
self.candidate_block(&candidates, &base_short)
);
}
jobs.push(SeatJob {
spec,
seat,
prompt: text,
cwd: root.join(format!("judge-{}", j + 1)),
timeout,
allow_write: false,
sessions,
artifacts: artifacts.clone(),
stem: format!("vote-judge-{}", j + 1),
});
seats_at.push(j);
}
self.state.event(
"vote",
format!(
"collecting {} final votes one by one, privately",
jobs.len()
),
);
let allowed = viable.clone();
let mut quota_losses = Vec::new();
let cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: &run_id,
node: "vote",
prompts: &prompts,
cache: cache.as_deref(),
};
let results = ask_json_wave::<FinalVote>(
jobs,
Arc::clone(&self.sem),
self.state.config.graph.retries,
&ctx,
&mut quota_losses,
&mut self.state,
&move |v: &FinalVote| match v.label() {
Some(c) if allowed.contains(&c) => Ok(()),
other => bail!("vote {other:?} is not one of {allowed:?}"),
},
)
.await;
self.state.quota.extend(quota_losses);
for (&j, (seat, res)) in seats_at.iter().zip(results) {
let agent_id = seat.agent.clone();
self.state.seats.insert(seat.key.clone(), seat);
let initial = self
.state
.judgements
.get(j)
.and_then(|r| r.ranking.first().copied());
let mut record = VoteRecord {
judge: j + 1,
agent: agent_id,
vote: None,
reason: String::new(),
changed: false,
};
match res {
Ok((v, _)) => {
record.vote = v.label();
record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
self.state.event(
"vote",
format!(
"judge {} voted {}{}",
j + 1,
record.vote.unwrap_or('?'),
if record.changed { " (changed)" } else { "" }
),
);
}
Err(e) => {
self.state
.event("vote", format!("judge {} cast no vote: {e}", j + 1));
}
}
self.state.votes.push(record);
self.state.save()?;
}
Ok(())
}
fn tally(&mut self) -> Result<()> {
if self.state.tally.is_some() {
return Ok(());
}
let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
let tops: Vec<char> = self
.state
.judgements
.iter()
.filter_map(|j| j.ranking.first().copied())
.collect();
let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
let mut cast: Vec<char> = Vec::new();
for (i, j) in self.state.judgements.iter().enumerate() {
let vote = self
.state
.votes
.iter()
.find(|v| v.judge == i + 1)
.and_then(|v| v.vote)
.or_else(|| j.ranking.first().copied());
if let Some(v) = vote {
*first_choice.entry(v).or_insert(0) += 1;
cast.push(v);
}
}
let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
for j in &self.state.judgements {
let n = j.ranking.len();
for (pos, label) in j.ranking.iter().enumerate() {
*borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
}
}
let best = first_choice.values().copied().max().unwrap_or(0);
let mut leaders: Vec<char> = first_choice
.iter()
.filter(|(_, v)| **v == best)
.map(|(k, _)| *k)
.collect();
let mut tie_break = None;
if leaders.len() > 1 {
let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
let borda_leaders: Vec<char> = leaders
.iter()
.copied()
.filter(|l| borda[l] == top_borda)
.collect();
tie_break = Some(if borda_leaders.len() == 1 {
format!(
"{} way tie on first-choice votes, broken by Borda points from the initial rankings",
leaders.len()
)
} else {
format!(
"{} way tie on both first-choice votes and Borda points, broken by label order",
leaders.len()
)
});
leaders = borda_leaders;
leaders.sort_unstable();
}
let winner = *leaders
.first()
.or(viable.first())
.context("no candidate to declare a winner from")?;
let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
let deliberated = !self.state.deliberation.is_empty();
let quota_seats: std::collections::BTreeSet<&str> =
self.state.quota.iter().map(|q| q.seat.as_str()).collect();
let mut present = 0usize;
for (i, j) in self.state.judgements.iter().enumerate() {
if quota_seats.contains(j.seat.as_str()) {
continue;
}
let ranked = !j.ranking.is_empty() && j.failed.is_none();
let voted = self
.state
.votes
.iter()
.any(|v| v.judge == i + 1 && v.vote.is_some());
if ranked || voted {
present += 1;
}
}
let needs_quorum = viable.len() > 1;
let judges_total = if needs_quorum {
self.roles.judges.len()
} else {
0
};
let quorum = if needs_quorum {
judges_total / 2 + 1
} else {
0
};
let met_quorum = !needs_quorum || present >= quorum;
let uncontested = (!needs_quorum).then(|| {
format!("only one candidate ({winner}) produced a usable change; no panel was asked")
});
self.state.event(
"tally",
match &uncontested {
Some(reason) => format!("winner {winner} — {reason}"),
None => format!(
"winner {winner} — votes {} | initial {} | {} changed | \
{present}/{judges_total} judges{}",
first_choice
.iter()
.map(|(k, v)| format!("{k}:{v}"))
.collect::<Vec<_>>()
.join(" "),
if unanimous_initial {
"unanimous"
} else {
"split"
},
changed_votes,
if met_quorum {
String::new()
} else {
format!(" — below quorum ({quorum} required)")
},
),
},
);
if !met_quorum {
self.state.event(
"stall",
format!(
"verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
the run stops here, resumable"
),
);
}
self.state.tally = Some(Tally {
first_choice,
borda,
winner,
rankings: tops.len(),
unanimous_initial,
deliberated,
changed_votes,
unanimous_final,
tie_break,
judges: judges_total,
present,
quorum,
met_quorum,
uncontested,
});
self.state.status = if met_quorum {
RunStatus::Reviewing
} else {
RunStatus::Stalled
};
self.state.save()?;
Ok(())
}
#[allow(clippy::too_many_lines)]
async fn recover_stall(&mut self) -> Result<bool> {
let run_id = self.state.id.clone();
let prompts = self.state.config.prompts.clone();
let quota_seats: BTreeSet<&str> =
self.state.quota.iter().map(|q| q.seat.as_str()).collect();
let absent: Vec<String> = self
.state
.judgements
.iter()
.filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
.map(|j| j.seat.clone())
.collect();
if absent.is_empty() {
return Ok(false);
}
let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
if viable.len() <= 1 {
return Ok(false);
}
let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
let language = self.state.config.graph.language.clone();
let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
let sessions = self.state.config.graph.sessions;
let artifacts = agent::artifacts_dir(&self.state.dir());
let root = self.state.worktree_root();
let base_short = short(&self.state.base_commit);
let candidates: Vec<Candidate> = viable.clone();
let mut positions: Vec<usize> = absent
.iter()
.filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
.collect();
if positions.is_empty() {
return Ok(false);
}
positions.sort_unstable();
positions.dedup();
let mut judge_jobs = Vec::new();
for &j in &positions {
let order = blind::presentation_order(viable.len(), j, self.state.seed);
let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
let seat_key = format!("judge-{}", j + 1);
let spec = self.roles.judges[j].clone();
let seat = self.seat(&seat_key, &spec.id);
judge_jobs.push(SeatJob {
spec,
seat,
prompt: prompt::judge(
&self.state.instruction,
&views,
self.roles.judges.len(),
&base_short,
&language,
),
cwd: root.join(seat_key),
timeout,
allow_write: false,
sessions,
artifacts: artifacts.clone(),
stem: format!("judge-{}-recover", j + 1),
});
}
let labels_for_check = labels.clone();
let mut judge_losses = Vec::new();
let retries = self.state.config.graph.retries;
let cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: &run_id,
node: "judge",
prompts: &prompts,
cache: cache.as_deref(),
};
let results = ask_json_wave::<Ranking>(
judge_jobs,
Arc::clone(&self.sem),
retries,
&ctx,
&mut judge_losses,
&mut self.state,
&move |r: &Ranking| r.validate(&labels_for_check),
)
.await;
let mut recovered: BTreeSet<usize> = BTreeSet::new();
for (&j, (seat, res)) in positions.iter().zip(results) {
self.state.seats.insert(seat.key.clone(), seat);
let record = &mut self.state.judgements[j];
match res {
Ok((ranking, out)) => {
record.ranking = ranking.normalized();
record.reasons = ranking.reasons;
record.confidence = ranking.confidence;
record.failed = None;
record.duration_ms = out.duration_ms;
recovered.insert(j);
self.state.event(
"recover",
format!("judge {} ranked again after the limit", j + 1),
);
}
Err(e) => {
self.state
.event("recover", format!("judge {} still cannot rank: {e}", j + 1));
}
}
}
let mut vote_jobs = Vec::new();
let mut vote_pos: Vec<usize> = Vec::new();
for &j in &recovered {
let seat_key = format!("judge-{}", j + 1);
let spec = self.roles.judges[j].clone();
let seat = self.seat(&seat_key, &spec.id);
let mut text = prompt::final_vote(&labels, &language);
if !has_context(&spec, &seat, sessions) {
text = format!(
"{}\n\n# Candidates\n\n{}",
text,
self.candidate_block(&candidates, &base_short)
);
}
vote_jobs.push(SeatJob {
spec,
seat,
prompt: text,
cwd: root.join(seat_key),
timeout,
allow_write: false,
sessions,
artifacts: artifacts.clone(),
stem: format!("vote-judge-{}-recover", j + 1),
});
vote_pos.push(j);
}
let allowed = labels.clone();
let mut vote_losses = Vec::new();
let vote_retries = self.state.config.graph.retries;
let vote_cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: &run_id,
node: "vote",
prompts: &prompts,
cache: vote_cache.as_deref(),
};
let votes = ask_json_wave::<FinalVote>(
vote_jobs,
Arc::clone(&self.sem),
vote_retries,
&ctx,
&mut vote_losses,
&mut self.state,
&move |v: &FinalVote| match v.label() {
Some(c) if allowed.contains(&c) => Ok(()),
other => bail!("vote {other:?} is not one of {allowed:?}"),
},
)
.await;
for (&j, (seat, res)) in vote_pos.iter().zip(votes) {
let agent_id = seat.agent.clone();
self.state.seats.insert(seat.key.clone(), seat);
match res {
Ok((v, _)) => {
if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
rec.vote = v.label();
rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
} else {
self.state.votes.push(VoteRecord {
judge: j + 1,
agent: agent_id,
vote: v.label(),
reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
changed: false,
});
}
self.state.event(
"recover",
format!("judge {} voted again after the limit", j + 1),
);
}
Err(e) => {
self.state
.event("recover", format!("judge {} still cannot vote: {e}", j + 1));
}
}
}
if !recovered.is_empty() {
let recovered_keys: BTreeSet<String> = recovered
.iter()
.map(|&j| format!("judge-{}", j + 1))
.collect();
self.state
.quota
.retain(|q| !recovered_keys.contains(&q.seat));
}
self.state.tally = None;
self.tally()?;
Ok(self
.state
.tally
.as_ref()
.map(|t| t.met_quorum)
.unwrap_or(false))
}
async fn fold_losers(&mut self) -> Result<()> {
let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
return Ok(());
};
let repo = self.state.repo.clone();
let mut folded = Vec::new();
for i in 0..self.state.candidates.len() {
let c = &self.state.candidates[i];
if c.label == winner || c.folded {
continue;
}
let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
git::worktree_remove(&repo, &wt).await.ok();
git::branch_delete(&repo, &branch).await.ok();
self.state.candidates[i].folded = true;
folded.push(label.to_string());
}
let root = self.state.worktree_root();
for j in 1..=self.roles.judges.len() {
let wt = root.join(format!("judge-{j}"));
if wt.exists() {
git::worktree_remove(&repo, &wt).await.ok();
}
}
if !folded.is_empty() {
self.state
.event("fold", format!("folded candidates {}", folded.join(", ")));
self.state.save()?;
}
Ok(())
}
async fn sync_to_base(&mut self) -> Result<()> {
if self
.state
.base_sync
.as_ref()
.is_some_and(|s| s.conflict.is_some())
{
return Ok(());
}
let Some(winner) = self.state.winner().cloned() else {
return Ok(());
};
let repo = self.state.repo.clone();
let remote = self.state.config.merge.remote.clone();
let base_branch = self.state.base_branch.clone();
let tracking = format!("{remote}/{base_branch}");
git::fetch(&repo, &remote, &base_branch).await.ok();
let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
return Ok(());
};
let head = git::rev_parse(&winner.worktree, "HEAD").await?;
let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
if behind == 0 {
self.state.base_sync = Some(BaseSync {
tip,
behind: 0,
attempts,
conflict: None,
});
self.state.save()?;
return Ok(());
}
if attempts >= BASE_SYNC_ROUNDS {
let why = format!(
"{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
rebase(s); rebasing again would only race it",
winner.branch
);
self.state.status = RunStatus::Blocked;
self.state.base_sync = Some(BaseSync {
tip,
behind,
attempts,
conflict: Some(why.clone()),
});
self.state.event("land", why);
self.state.save()?;
return Ok(());
}
self.state.event(
"land",
format!(
"{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
winner.branch
),
);
self.state.save()?;
let scratch = self.state.dir().join("base-sync");
let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
let attempts = attempts + 1;
match rebased {
Ok(None) => {
git::sync_to_head(&winner.worktree).await?;
self.state.base_sync = Some(BaseSync {
tip: tip.clone(),
behind: 0,
attempts,
conflict: None,
});
self.state
.event("land", format!("rebased {} onto {tracking}", winner.branch));
}
Ok(Some(conflict)) => {
let why = format!(
"{} conflicts with {tracking} and did not rebase: {}",
winner.branch,
conflict.chars().take(600).collect::<String>()
);
self.state.status = RunStatus::Blocked;
self.state.base_sync = Some(BaseSync {
tip,
behind,
attempts,
conflict: Some(why.clone()),
});
self.state.event("land", why);
}
Err(e) => {
let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
self.state.status = RunStatus::Blocked;
self.state.base_sync = Some(BaseSync {
tip,
behind,
attempts,
conflict: Some(why.clone()),
});
self.state.event("land", why);
}
}
self.state.save()?;
Ok(())
}
fn landing_base(&self) -> String {
self.state
.base_sync
.as_ref()
.map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
}
async fn review_loop(&mut self) -> Result<()> {
if self
.state
.base_sync
.as_ref()
.is_some_and(|s| s.conflict.is_some())
{
return Ok(());
}
let run_id = self.state.id.clone();
let prompts = self.state.config.prompts.clone();
let Some(winner) = self.state.winner().cloned() else {
return Ok(());
};
let max_rounds = self.state.config.graph.review_rounds;
if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
self.state.status = status;
self.state.save()?;
return Ok(());
}
self.state.status = RunStatus::Reviewing;
let repo = self.state.repo.clone();
let root = self.state.worktree_root();
let language = self.state.config.graph.language.clone();
let sessions = self.state.config.graph.sessions;
let artifacts = agent::artifacts_dir(&self.state.dir());
let base = self.landing_base();
let base_short = short(&base);
let reviewers = self.roles.reviewers.clone();
let shell = self.state.config.shell();
let mut prev_e2e: Option<String> = None;
for round in (self.state.reviews.len() + 1)..=max_rounds {
let head = git::rev_parse(&winner.worktree, "HEAD").await?;
let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
let mut jobs = Vec::new();
for (r, spec) in reviewers.iter().cloned().enumerate() {
let wt = root.join(format!("review-{}", r + 1));
if wt.exists() {
git::reset_detached(&wt, &head).await?;
} else {
git::worktree_add_detached(&repo, &wt, &head).await?;
}
let seat_key = format!("review-{}", r + 1);
let seat = self.seat(&seat_key, &spec.id);
jobs.push(SeatJob {
prompt: prompt::review(&prompt::ReviewCtx {
instruction: &self.state.instruction,
branch: &winner.branch,
base_short: &base_short,
stat: &stat,
patch: &patch,
e2e: prev_e2e.as_deref(),
reviewers: reviewers.len(),
round,
rounds: max_rounds,
competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
lens: Lens::for_seat(r),
language: &language,
}),
spec,
seat,
cwd: wt,
timeout: Duration::from_secs(self.state.config.graph.timeout_review),
allow_write: false,
sessions,
artifacts: artifacts.clone(),
stem: format!("review-{round}-{}", r + 1),
});
}
self.state.event(
"review",
format!(
"round {round}: {} reviewers on {}",
jobs.len(),
short(&head)
),
);
let mut quota_losses = Vec::new();
let review_retries = self.state.config.graph.retries;
let review_cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: &run_id,
node: "review",
prompts: &prompts,
cache: review_cache.as_deref(),
};
let results = ask_json_wave::<Review>(
jobs,
Arc::clone(&self.sem),
review_retries,
&ctx,
&mut quota_losses,
&mut self.state,
&|_: &Review| Ok(()),
)
.await;
let round_quota_missing = quota_losses.len();
self.state.quota.extend(quota_losses);
let mut records = Vec::new();
let mut all_findings = Vec::new();
for (r, (seat, res)) in results.into_iter().enumerate() {
let agent_id = seat.agent.clone();
self.state.seats.insert(seat.key.clone(), seat);
let mut record = ReviewRecord {
reviewer: r + 1,
agent: agent_id,
summary: String::new(),
findings: Vec::new(),
vote: None,
failed: None,
duration_ms: 0,
};
match res {
Ok((review, out)) => {
record.summary =
blind::sanitize_prose(&review.summary, &self.state.config.blind);
record.vote = Some(review.vote);
record.duration_ms = out.duration_ms;
for (n, mut f) in review.findings.into_iter().enumerate() {
f.id = format!("R{round}-{}-{}", r + 1, n + 1);
f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
f.file = f
.file
.map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
all_findings.push(f.clone());
record.findings.push(f);
}
self.state.event(
"review",
format!(
"round {round}: reviewer {} voted {} with {} finding(s)",
r + 1,
review.vote.label(),
record.findings.len()
),
);
}
Err(e) => {
record.failed = Some(e.to_string());
self.state.event(
"review",
format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
);
}
}
records.push(record);
}
let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
let vote_split =
initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
if vote_split {
self.state.event(
"review",
format!(
"round {round}: votes split ({}) — one round of reconsideration",
initial_votes
.iter()
.map(|v| v.label())
.collect::<Vec<_>>()
.join(", ")
),
);
let panel: Vec<ReviewSeatReport<'_>> = records
.iter()
.filter_map(|r| {
r.vote.map(|vote| ReviewSeatReport {
reviewer: r.reviewer,
vote,
summary: &r.summary,
findings: &r.findings,
})
})
.collect();
let mut jobs = Vec::new();
let mut seats_at = Vec::new();
for (r, spec) in reviewers.iter().cloned().enumerate() {
if records[r].vote.is_none() {
continue;
}
let wt = root.join(format!("review-{}", r + 1));
let seat_key = format!("review-{}", r + 1);
let seat = self.seat(&seat_key, &spec.id);
let patch_ctx = if has_context(&spec, &seat, sessions) {
None
} else {
Some(ReviewPatch {
branch: &winner.branch,
base_short: &base_short,
stat: &stat,
patch: &patch,
})
};
let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
instruction: &self.state.instruction,
reviewer: r + 1,
lens: Lens::for_seat(r),
panel: &panel,
patch: patch_ctx,
round,
rounds: max_rounds,
language: &language,
});
jobs.push(SeatJob {
prompt,
spec,
seat,
cwd: wt,
timeout: Duration::from_secs(self.state.config.graph.timeout_review),
allow_write: false,
sessions,
artifacts: artifacts.clone(),
stem: format!("review-{round}-reconsider-{}", r + 1),
});
seats_at.push(r);
}
let mut recon_quota_losses = Vec::new();
let recon_cache = self.state.config.cache_dir();
let recon_ctx = WaveCtx {
run: &run_id,
node: "review",
prompts: &prompts,
cache: recon_cache.as_deref(),
};
let recon_results = ask_json_wave::<ReviewRevote>(
jobs,
Arc::clone(&self.sem),
review_retries,
&recon_ctx,
&mut recon_quota_losses,
&mut self.state,
&|_: &ReviewRevote| Ok(()),
)
.await;
self.state.quota.extend(recon_quota_losses);
for (&r, (seat, res)) in seats_at.iter().zip(recon_results) {
let agent_id = seat.agent.clone();
self.state.seats.insert(seat.key.clone(), seat);
let mut rec = ReviewRevoteRecord {
reviewer: r + 1,
agent: agent_id,
vote: None,
reason: String::new(),
failed: None,
};
match res {
Ok((rv, _)) => {
rec.vote = Some(rv.vote);
rec.reason =
blind::sanitize_prose(&rv.reason, &self.state.config.blind);
self.state.event(
"review",
format!(
"round {round}: reviewer {} revoted {}",
r + 1,
rv.vote.label()
),
);
}
Err(e) => {
rec.failed = Some(e.to_string());
self.state.event(
"review",
format!("round {round}: reviewer {} did not revote: {e}", r + 1),
);
}
}
reconsideration.push(rec);
}
} else if initial_votes.len() > 1 {
self.state.event(
"review",
format!(
"round {round}: votes agreed ({}) — no reconsideration",
initial_votes[0].label()
),
);
}
let final_votes: Vec<ReviewVote> = records
.iter()
.filter_map(|r| {
reconsideration
.iter()
.find(|rv| rv.reviewer == r.reviewer)
.and_then(|rv| rv.vote)
.or(r.vote)
})
.collect();
let round_verdict = ReviewVote::worst(final_votes);
let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
let defer_e2e =
blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
let reason =
format!("{blocking} blocking finding(s) already required a fix this round");
self.state.event(
"verify",
format!(
"round {round}: {reason} — e2e deferred to the fixer (reviewed head \
{}); it will run once a round has none left",
short(&head)
),
);
(Vec::new(), false, true, Some(reason))
} else {
let e2e_commands = self.state.config.verify.e2e.clone();
let (e2e, verify_retried) = run_e2e_with_retry(
&mut self.state,
&shell,
&e2e_commands,
&winner.worktree,
verify_timeout,
&format!("round {round}"),
)
.await;
(e2e, verify_retried, false, None)
};
let e2e_failures: String = e2e
.iter()
.filter(|o| !o.ok())
.map(|o| format!("$ {}\n{}\n", o.command, o.output_tail))
.collect();
let expected = records.len();
let answered = records.iter().filter(|r| r.failed.is_none()).count();
let incomplete = answered < expected;
let e2e_ok = e2e.iter().all(CommandOutcome::ok);
let policy = self.state.config.graph.incomplete_review;
let clean = round_is_clean(
blocking,
e2e_ok,
answered,
expected,
round_quota_missing,
policy,
);
let mut round_record = ReviewRound {
round,
head: head.clone(),
verified_head: None,
reviews: records,
e2e,
verify_retried,
e2e_deferred,
e2e_defer_reason,
fix: None,
blocking,
answered,
expected,
clean,
progressed: false,
vote_split,
reconsideration,
verdict: round_verdict,
};
if incomplete {
let missing: Vec<String> = round_record
.reviews
.iter()
.filter(|r| r.failed.is_some())
.map(|r| format!("review-{}", r.reviewer))
.collect();
self.state.event(
"review",
format!(
"round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
missing.join(", ")
),
);
}
if clean {
self.state.event(
"review",
if incomplete && policy == IncompleteReviewPolicy::Warn {
format!(
"round {round}: clean (warn policy, incomplete panel) — no \
blocking findings from the seats that answered, verification green"
)
} else if incomplete {
format!(
"round {round}: clean ({} rate-limited reviewer(s) excluded from \
quorum) — no blocking findings from the seats that answered, \
verification green",
expected - answered
)
} else {
format!("round {round}: clean — no blocking findings, verification green")
},
);
self.state.reviews.push(round_record);
self.state.status = RunStatus::Gating;
self.state.save()?;
return Ok(());
}
if incomplete && blocking == 0 && e2e_ok {
self.state.reviews.push(round_record);
self.state.save()?;
if round == max_rounds {
self.state.status = RunStatus::Blocked;
self.state.event(
"review",
format!(
"{} reviewer seat(s) never answered after {max_rounds} rounds; \
refusing to call it clean",
expected - answered
),
);
return Ok(());
}
prev_e2e = None;
continue;
}
if round == max_rounds {
self.state.reviews.push(round_record);
return self
.stop_reviewing(
&format!(
"{blocking} blocking finding(s) still open after {max_rounds} round(s)"
),
&shell,
&winner.worktree,
)
.await;
}
let (fix_spec, fix_seat_key) = match &self.roles.fixer {
Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
_ => (
self.state
.config
.agent(&winner.agent)
.cloned()
.unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
format!("impl-{}", winner.label),
),
};
let seat = self.seat(&fix_seat_key, &fix_spec.id);
let blocking_findings: Vec<_> = all_findings
.iter()
.filter(|f| f.severity.blocks())
.cloned()
.collect();
let job = SeatJob {
prompt: prompt::fix(
&self.state.instruction,
&blocking_findings,
(!e2e_failures.is_empty()).then_some(e2e_failures.as_str()),
e2e_deferred,
round,
max_rounds,
&language,
),
spec: fix_spec.clone(),
seat,
cwd: winner.worktree.clone(),
timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
allow_write: true,
sessions,
artifacts: artifacts.clone(),
stem: format!("fix-{round}"),
};
let before = git::rev_parse(&winner.worktree, "HEAD").await?;
let cache = self.state.config.cache_dir();
let ctx = WaveCtx {
run: &run_id,
node: "fix",
prompts: &prompts,
cache: cache.as_deref(),
};
let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
let agent_id = seat.agent.clone();
let seat_key = seat.key.clone();
self.state.seats.insert(seat.key.clone(), seat);
let mut fix = FixRecord {
agent: agent_id,
addressed: Vec::new(),
rejected: Vec::new(),
notes: String::new(),
committed: false,
failed: None,
duration_ms: 0,
};
match out {
AgentOutcome::Ok(o) => {
fix.duration_ms = o.duration_ms;
match verdict::extract_json::<FixReport>(&o.text) {
Ok(report) => {
fix.addressed = report.addressed;
fix.rejected = report.rejected;
fix.notes =
blind::sanitize_prose(&report.notes, &self.state.config.blind);
}
Err(e) => fix.failed = Some(format!("unparsable fix report: {e}")),
}
}
AgentOutcome::Dropped(o) => {
fix.duration_ms = o.duration_ms;
let why = o
.dropped
.as_ref()
.map(|d| d.why.as_str())
.unwrap_or("the CLI ended the stream without delivering its answer");
fix.failed = Some(format!("the CLI dropped the stream ({why})"));
}
AgentOutcome::Quota(o) => {
self.state.quota.push(QuotaLoss {
seat: seat_key,
node: "fix".to_owned(),
at: Timestamp::now(),
reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
});
fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
}
AgentOutcome::Failed(e) => fix.failed = Some(e),
}
git::commit_all(
&winner.worktree,
&format!("magi: review round {round} fixes (uncommitted work)"),
)
.await
.ok();
let after = git::rev_parse(&winner.worktree, "HEAD").await?;
fix.committed = after != before;
let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
let progressed = diff_after != patch;
let commit_note = if fix.committed {
"committed"
} else {
"NO new commit"
};
let tree_note = if progressed {
"changed vs base"
} else {
"unchanged vs base"
};
self.state.event(
"fix",
match &fix.failed {
Some(reason) => {
format!(
"round {round}: fixer's adoption report was lost ({reason}); \
{commit_note}, tree {tree_note}"
)
}
None => format!(
"round {round}: {} addressed, {} rejected, {commit_note}, tree {tree_note}",
fix.addressed.len(),
fix.rejected.len(),
),
},
);
round_record.fix = Some(fix);
round_record.progressed = progressed;
self.state.reviews.push(round_record);
self.state.save()?;
prev_e2e = (!e2e_failures.is_empty()).then_some(e2e_failures);
let streak = self
.state
.reviews
.iter()
.rev()
.take_while(|r| !r.progressed)
.count();
if streak >= STAGNANT_LIMIT {
return self
.stop_reviewing(
&format!(
"the tree has not moved against base for {streak} round(s) in a row"
),
&shell,
&winner.worktree,
)
.await;
}
}
Ok(())
}
async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
let round_idx = self.state.reviews.len() - 1;
let needs_catchup_run = {
let last = &self.state.reviews[round_idx];
last.e2e.is_empty() && last.e2e_deferred
};
if needs_catchup_run {
let round = self.state.reviews[round_idx].round;
let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
let commands = self.state.config.verify.e2e.clone();
let verified_head = git::rev_parse(worktree, "HEAD").await?;
let (outcomes, verify_retried) = run_e2e_with_retry(
&mut self.state,
shell,
&commands,
worktree,
timeout,
&format!("round {round}: deferred e2e, now catching up before the final decision"),
)
.await;
let last = &mut self.state.reviews[round_idx];
last.e2e = outcomes;
last.verify_retried = verify_retried;
last.e2e_deferred = false;
if verified_head != last.head {
last.verified_head = Some(verified_head);
}
}
let last = &self.state.reviews[round_idx];
let red: Vec<String> = last
.e2e
.iter()
.filter(|o| !o.ok())
.map(|o| {
format!(
"`{}` -> {:?}\n{}",
o.command,
o.code,
tail(&o.output_tail, EVENT_OUTPUT_TAIL)
)
})
.collect();
let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
if red.is_empty() {
self.state.event(
"review",
format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
);
self.state.status = RunStatus::Gating;
} else {
self.state
.event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
self.state.status = RunStatus::Blocked;
}
self.state.save()?;
Ok(())
}
async fn gate(&mut self) -> Result<()> {
if self.state.status == RunStatus::Failed
|| self
.state
.base_sync
.as_ref()
.is_some_and(|s| s.conflict.is_some())
|| review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
!= Some(RunStatus::Gating)
{
return Ok(());
}
if !self.state.gate.is_empty() {
if self.state.gate.iter().any(|outcome| !outcome.ok()) {
self.state.status = RunStatus::Blocked;
self.state.save()?;
}
return Ok(());
}
let Some(winner) = self.state.winner().cloned() else {
return Ok(());
};
self.state.status = RunStatus::Gating;
let shell = self.state.config.shell();
let outcomes = run_commands(
&shell,
&self.state.config.verify.gate,
&winner.worktree,
Duration::from_secs(self.state.config.graph.verify_timeout()),
)
.await;
for o in &outcomes {
self.state.event(
"gate",
format!(
"`{}` -> {}",
o.command,
if o.ok() {
"pass".to_owned()
} else {
format!(
"FAIL ({:?})\n{}",
o.code,
tail(&o.output_tail, EVENT_OUTPUT_TAIL)
)
}
),
);
}
let passed = outcomes.iter().all(CommandOutcome::ok);
self.state.gate = outcomes;
if !passed {
self.state.status = RunStatus::Blocked;
self.state.event("gate", "gate failed; not merging");
}
self.state.save()?;
Ok(())
}
async fn merge(&mut self) -> Result<()> {
if self
.state
.base_sync
.as_ref()
.is_some_and(|s| s.conflict.is_some())
|| review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
!= Some(RunStatus::Gating)
|| self.state.gate.iter().any(|o| !o.ok())
{
return Ok(());
}
if self.state.merge.is_some() {
return Ok(());
}
let Some(winner) = self.state.winner().cloned() else {
return Ok(());
};
let repo = self.state.repo.clone();
let base = self.state.base_branch.clone();
let mode = self.state.config.merge.mode;
let style = self.state.config.merge.style;
let message = pr_body(&self.state, winner.label);
let outcome = match mode {
MergeMode::None => MergeOutcome {
mode,
ok: true,
detail: manual_merge_command(style, &repo, &winner.branch, &message),
},
MergeMode::Local => {
let on = git::current_branch(&repo).await?;
if on.as_deref() != Some(base.as_str()) {
MergeOutcome {
mode,
ok: false,
detail: format!(
"{} has {} checked out, not the base branch {base}",
repo.display(),
on.unwrap_or_else(|| "a detached HEAD".to_owned())
),
}
} else if !git::is_clean(&repo).await? {
MergeOutcome {
mode,
ok: false,
detail: format!("{} is dirty; refusing to merge", repo.display()),
}
} else {
let out = match style {
MergeStyle::Merge => {
git::merge_no_ff(&repo, &winner.branch, &message).await?
}
MergeStyle::Squash => {
git::merge_squash(&repo, &winner.branch, &message).await?
}
MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
};
MergeOutcome {
mode,
ok: out.ok(),
detail: if out.ok() { out.stdout } else { out.stderr },
}
}
}
MergeMode::Pr => {
let remote = self.state.config.merge.remote.clone();
let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
if !pushed.ok() {
MergeOutcome {
mode,
ok: false,
detail: pushed.stderr,
}
} else {
let out = gh_pr_create(&winner.worktree, &base, &winner.branch, &message).await;
match out {
Ok(url) => MergeOutcome {
mode,
ok: true,
detail: url,
},
Err(e) => MergeOutcome {
mode,
ok: false,
detail: e.to_string(),
},
}
}
}
};
self.state.status = match (mode, outcome.ok) {
(MergeMode::None, _) => RunStatus::Ready,
(_, true) => RunStatus::Merged,
(_, false) => RunStatus::Blocked,
};
self.state.event(
"merge",
format!(
"{:?}: {}",
mode,
outcome.detail.lines().next().unwrap_or("")
),
);
self.state.merge = Some(outcome);
self.state.save()?;
if self.state.config.graph.land
&& mode == MergeMode::Pr
&& self.state.status == RunStatus::Merged
{
self.run_land().await?;
}
self.settle_questions();
Ok(())
}
async fn run_land(&mut self) -> Result<()> {
let url = self
.state
.merge
.as_ref()
.map(|m| m.detail.clone())
.unwrap_or_default();
let url = url.lines().next().unwrap_or("").trim().to_owned();
if !url.starts_with("http") {
return Ok(());
}
match land::land(&mut self.state, &url).await {
Ok(pr) if self.state.parked => {
let _ = pr;
}
Ok(pr) => {
self.state.status = match pr.state {
land::PrLifecycle::Merged => RunStatus::Merged,
_ => RunStatus::Blocked,
};
if bump::should_release_bump(self.state.status)
&& let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
{
self.state
.event("bump", format!("release bump skipped: {e:#}"));
}
self.state.save()?;
}
Err(e) => {
self.state.status = RunStatus::Blocked;
self.state.event("land", format!("gave up: {e}"));
self.state.save()?;
}
}
Ok(())
}
fn seat(&mut self, key: &str, agent: &str) -> SeatState {
if let Some(existing) = self.state.seats.get(key)
&& existing.agent == agent
{
return existing.clone();
}
let fresh = SeatState::new(key, agent, self.state.seed);
self.state.seats.insert(key.to_owned(), fresh.clone());
fresh
}
fn view(&self, c: &Candidate) -> CandidateView {
let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
.unwrap_or_default();
let (patch, _) = blind::sanitize_patch(
&format!("candidate {} patch", c.label),
&raw,
&self.state.config.blind,
);
CandidateView {
label: c.label,
branch: c.branch.clone(),
summary: c.summary.clone(),
stat: c.stat.clone(),
patch,
}
}
fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
prompt::judge(
"(see above)",
&views,
self.roles.judges.len(),
base_short,
"en",
)
}
fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
let mut turns = Vec::new();
for j in &self.state.judgements {
if j.ranking.is_empty() {
continue;
}
let reasons = j
.reasons
.iter()
.map(|(k, v)| format!("- {k}: {v}"))
.collect::<Vec<_>>()
.join("\n");
turns.push(Turn {
who: format!("Judge {} (opening ranking)", j.judge),
is_self: j.judge == self_idx + 1,
body: format!(
"Ranked {}{}{reasons}",
j.ranking.iter().collect::<String>(),
if reasons.is_empty() {
""
} else {
", because:\n"
}
),
});
}
for t in self
.state
.deliberation
.iter()
.flat_map(|r| r.turns.iter())
.chain(current)
{
turns.push(Turn {
who: format!("Judge {}", t.judge),
is_self: t.judge == self_idx + 1,
body: t.body.clone(),
});
}
turns
}
}
fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
agent::has_session(spec.kind, seat, sessions)
}
fn short(commit: &str) -> String {
commit.chars().take(7).collect()
}
fn make_executable(path: &Path) -> Result<()> {
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt as _;
let mut perms = std::fs::metadata(path)?.permissions();
perms.set_mode(0o755);
std::fs::set_permissions(path, perms)?;
}
#[cfg(not(unix))]
{
let _ = path;
}
Ok(())
}
struct WaveCtx<'a> {
run: &'a str,
node: &'a str,
prompts: &'a Prompts,
cache: Option<&'a Path>,
}
async fn run_one(
job: SeatJob,
sem: Arc<Semaphore>,
ctx: &WaveCtx<'_>,
state: &mut RunState,
attempt: usize,
) -> (SeatState, AgentOutcome) {
let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
.await
.pop()
.expect("one job in, one result out");
(seat, out)
}
async fn wave(
jobs: Vec<SeatJob>,
sem: Arc<Semaphore>,
ctx: &WaveCtx<'_>,
state: &mut RunState,
attempt: usize,
) -> Vec<(usize, SeatState, AgentOutcome)> {
let WaveCtx {
run,
node,
prompts,
cache,
} = *ctx;
for job in &jobs {
state.seat_started(node, &job.seat.key, job.timeout, attempt);
}
if let Err(e) = state.save() {
tracing::warn!("could not persist in-progress seats: {e:#}");
}
let mut set = tokio::task::JoinSet::new();
let overlay = prompts.overlay(node);
for (i, mut job) in jobs.into_iter().enumerate() {
job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
if cache.is_some() {
job.prompt.push('\n');
job.prompt.push_str(&prompt::build_cache_note(node));
}
let sem = Arc::clone(&sem);
let run = run.to_owned();
let node = node.to_owned();
let cache = cache.map(Path::to_path_buf);
set.spawn(async move {
let _permit = sem.acquire().await;
let mut seat = job.seat;
let out = agent::invoke(
&job.spec,
&mut seat,
&Invocation {
cwd: &job.cwd,
prompt: &job.prompt,
timeout: job.timeout,
allow_write: job.allow_write,
sessions: job.sessions,
artifacts: &job.artifacts,
stem: &job.stem,
run: &run,
node: &node,
cache_dir: cache.as_deref(),
attachments: &[],
},
)
.await;
let out = match out {
Ok(o) if o.usable() => AgentOutcome::Ok(o),
Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
Ok(o) => AgentOutcome::Failed(format!(
"exited with {:?} and no usable output",
o.exit_code
)),
Err(e) => AgentOutcome::Failed(e.to_string()),
};
(i, seat, out)
});
}
let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
while let Some(joined) = set.join_next().await {
let (i, seat, out) = match joined {
Ok(v) => v,
Err(e) => {
tracing::error!("agent task panicked: {e}");
continue;
}
};
state.seat_finished(&seat.key);
if let Err(e) = state.save() {
tracing::warn!("could not persist a seat's completion: {e:#}");
}
if collected.len() <= i {
collected.resize_with(i + 1, || None);
}
collected[i] = Some((i, seat, out));
}
if state
.active
.values()
.any(|a| a.node == node && a.attempt == attempt)
{
state
.active
.retain(|_, a| !(a.node == node && a.attempt == attempt));
if let Err(e) = state.save() {
tracing::warn!("could not persist the end of a wave: {e:#}");
}
}
collected.into_iter().flatten().collect()
}
fn round_is_clean(
blocking: usize,
e2e_ok: bool,
answered: usize,
expected: usize,
quota_missing: usize,
policy: IncompleteReviewPolicy,
) -> bool {
if blocking != 0 || !e2e_ok {
return false;
}
if answered == expected || policy == IncompleteReviewPolicy::Warn {
return true;
}
answered > 0 && expected - answered <= quota_missing
}
fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
return Some(RunStatus::Gating);
}
let last = reviews.last()?;
let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
if reviews.len() < max_rounds && !stagnant {
return None;
}
Some(if last.incomplete() && last.blocking == 0 {
RunStatus::Blocked
} else if last.e2e.iter().all(CommandOutcome::ok) {
RunStatus::Gating
} else {
RunStatus::Blocked
})
}
fn retry_budget(full: Duration, nudged: bool) -> Duration {
if nudged {
(full / 4).max(Duration::from_secs(120)).min(full)
} else {
full
}
}
#[allow(clippy::too_many_arguments)]
async fn ask_json_wave<T>(
jobs: Vec<SeatJob>,
sem: Arc<Semaphore>,
retries: usize,
ctx: &WaveCtx<'_>,
losses: &mut Vec<QuotaLoss>,
state: &mut RunState,
validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
) -> Vec<(SeatState, Result<(T, AgentOutput)>)>
where
T: serde::de::DeserializeOwned + Send + 'static,
{
let n = jobs.len();
let originals: Vec<SeatJob> = jobs;
let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
let mut pending: Vec<usize> = (0..n).collect();
for attempt in 0..=retries {
if pending.is_empty() {
break;
}
let mut batch = Vec::with_capacity(pending.len());
for &i in &pending {
let src = &originals[i];
let (prompt, timeout) = if attempt == 0 {
(src.prompt.clone(), src.timeout)
} else {
let why = done[i]
.as_ref()
.and_then(|r| r.as_ref().err().map(ToString::to_string))
.unwrap_or_else(|| "no parsable answer".to_owned());
let nudge = prompt::nudge(&why);
let nudged = has_context(&src.spec, &seats[i], src.sessions);
let prompt = if nudged {
nudge
} else {
format!("{}\n\n---\n\n{}", src.prompt, nudge)
};
(prompt, retry_budget(src.timeout, nudged))
};
batch.push(SeatJob {
spec: src.spec.clone(),
seat: seats[i].clone(),
cwd: src.cwd.clone(),
prompt,
timeout,
allow_write: src.allow_write,
sessions: src.sessions,
artifacts: src.artifacts.clone(),
stem: if attempt == 0 {
src.stem.clone()
} else {
format!("{}-retry{attempt}", src.stem)
},
});
}
if attempt > 0 {
let seats_out: Vec<&str> = pending
.iter()
.map(|&i| originals[i].seat.key.as_str())
.collect();
state.event(
ctx.node,
format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
);
}
let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
let mut still = Vec::new();
for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
seats[i] = seat;
let (parsed, quota) = match out {
AgentOutcome::Ok(o) => (
match verdict::extract_json::<T>(&o.text) {
Ok(v) => match validate(&v) {
Ok(()) => Ok((v, o)),
Err(e) => Err(e),
},
Err(e) => Err(e),
},
false,
),
AgentOutcome::Quota(o) => {
losses.push(QuotaLoss {
seat: originals[i].seat.key.clone(),
node: ctx.node.to_owned(),
at: Timestamp::now(),
reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
});
(
Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
true,
)
}
AgentOutcome::Dropped(o) => {
let why = o
.dropped
.as_ref()
.map(|d| d.why.as_str())
.unwrap_or("the CLI ended the stream without delivering its answer");
(
Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
false,
)
}
AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
};
let failed = parsed.is_err();
done[i] = Some(parsed);
if failed && !quota {
still.push(i);
}
}
pending = still;
}
seats
.into_iter()
.zip(done)
.map(|(seat, res)| {
(
seat,
res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
)
})
.collect()
}
fn e2e_outcome_label(o: &CommandOutcome) -> String {
if o.ok() {
return "pass".to_owned();
}
let reason = if o.build_failed() {
format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
} else {
format!("FAIL ({:?})", o.code)
};
format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
}
async fn run_e2e_with_retry(
state: &mut RunState,
shell: &[String],
commands: &[String],
worktree: &Path,
timeout: Duration,
context: &str,
) -> (Vec<CommandOutcome>, bool) {
let mut e2e = run_commands(shell, commands, worktree, timeout).await;
for o in &e2e {
state.event(
"verify",
format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
);
}
let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
if verify_retried {
state.event(
"verify",
format!(
"{context}: verify could not build/link, not a test result — retrying once \
before concluding"
),
);
e2e = run_commands(shell, commands, worktree, timeout).await;
for o in &e2e {
state.event(
"verify",
format!(
"{context}: retry `{}` -> {}",
o.command,
e2e_outcome_label(o)
),
);
}
}
(e2e, verify_retried)
}
async fn run_commands(
shell: &[String],
commands: &[String],
cwd: &Path,
timeout: Duration,
) -> Vec<CommandOutcome> {
let mut out = Vec::new();
for command in commands {
let started = Instant::now();
let mut cmd = tokio::process::Command::new(&shell[0]);
cmd.quiet();
cmd.args(&shell[1..])
.arg(command)
.current_dir(cwd)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
let spawned = cmd.spawn();
let (code, body) = match spawned {
Ok(child) => match tokio::time::timeout(timeout, child.wait_with_output()).await {
Ok(Ok(o)) => {
let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
body.push_str(&String::from_utf8_lossy(&o.stderr));
(o.status.code(), body)
}
Ok(Err(e)) => (None, format!("failed to run: {e}")),
Err(_) => (None, format!("timed out after {}s", timeout.as_secs())),
},
Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
};
out.push(CommandOutcome {
command: command.clone(),
code,
output_tail: tail(&body, OUTPUT_TAIL),
duration_ms: started.elapsed().as_millis() as u64,
});
}
out
}
fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
let repo = repo.display();
match style {
MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
MergeStyle::Squash => {
let subject = message.lines().next().unwrap_or(branch);
format!(
"git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
)
}
MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
}
}
fn pr_body(state: &RunState, winner: char) -> String {
let instruction = state.instruction.trim_start();
let mut message = if instruction.is_empty() {
"(empty task)".to_owned()
} else {
instruction.to_owned()
};
let open = state.open_findings();
if !open.is_empty() {
message.push_str("\n\n## Open review findings\n\n");
for f in &open {
message.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
}
}
if let Some(fix) = state.reviews.last().and_then(|r| r.fix.as_ref())
&& !fix.rejected.is_empty()
{
message.push_str("\n## Declined by the fixer\n\n");
for r in &fix.rejected {
message.push_str(&format!("- `{}`: {}\n", r.id, r.why));
}
}
message.push_str(&format!(
"\n\n---\nmagi:run/{} magi:candidate-{}\n",
state.id,
winner.to_ascii_lowercase()
));
message
}
async fn gh_pr_create(cwd: &Path, base: &str, head: &str, body: &str) -> Result<String> {
let title = body.lines().next().unwrap_or("magi run").to_owned();
let out = tokio::process::Command::new("gh")
.args([
"pr", "create", "--base", base, "--head", head, "--title", &title, "--body", body,
])
.current_dir(cwd)
.quiet()
.stdin(std::process::Stdio::null())
.output()
.await
.context("spawn gh")?;
if out.status.success() {
Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
} else {
bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
}
}
pub async fn fold_run(state: &mut RunState, drop_winner: bool) -> Result<Vec<String>> {
let repo = state.repo.clone();
let root = state.worktree_root();
let winner = state.tally.as_ref().map(|t| t.winner);
let mut removed = Vec::new();
for i in 0..state.candidates.len() {
let c = state.candidates[i].clone();
let is_winner = Some(c.label) == winner;
if is_winner && !drop_winner {
continue;
}
if c.worktree.exists() {
git::worktree_remove(&repo, &c.worktree).await.ok();
removed.push(c.worktree.to_string_lossy().into_owned());
}
if git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
git::branch_delete(&repo, &c.branch).await.ok();
removed.push(c.branch.clone());
}
state.candidates[i].folded = true;
}
for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
let path = name.path();
let keep = !drop_winner
&& winner.is_some_and(|w| {
path.file_name()
.is_some_and(|n| n == format!("cand-{w}").as_str())
});
if keep {
continue;
}
git::worktree_remove(&repo, &path).await.ok();
removed.push(path.to_string_lossy().into_owned());
}
remove_if_empty(&root);
if state.enabled_worktree_config && drop_winner {
git::release_worktree_config(&repo).await.ok();
state.enabled_worktree_config = false;
}
state.save()?;
Ok(removed)
}
fn remove_if_empty(dir: &Path) {
if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
std::fs::remove_dir(dir).ok();
}
}
pub fn worst_open(state: &RunState) -> Option<Severity> {
state
.reviews
.last()?
.reviews
.iter()
.flat_map(|r| r.findings.iter())
.map(|f| f.severity)
.max()
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::BTreeMap;
use std::time::Duration;
fn conductor() -> AgentSpec {
AgentSpec {
id: "conductor".to_owned(),
kind: crate::config::AgentKind::Command,
model: None,
command: vec!["true".to_owned()],
extra_args: Vec::new(),
env: BTreeMap::new(),
prompt_delivery: None,
}
}
#[test]
fn remove_if_empty_only_ever_takes_a_bare_directory() {
let dir = tempfile::tempdir().unwrap();
let bay = dir.path().join("ffff");
remove_if_empty(&bay);
assert!(!bay.exists());
std::fs::create_dir_all(bay.join("cand-A")).unwrap();
remove_if_empty(&bay);
assert!(bay.exists(), "non-empty directory must survive");
std::fs::remove_dir(bay.join("cand-A")).unwrap();
remove_if_empty(&bay);
assert!(!bay.exists(), "an empty bay is a leftover, not a record");
}
#[test]
fn a_full_panel_that_found_nothing_is_clean() {
assert!(round_is_clean(
0,
true,
2,
2,
0,
IncompleteReviewPolicy::Block
));
}
#[test]
fn a_missing_seat_is_never_clean_under_the_default_policy() {
assert!(!round_is_clean(
0,
true,
1,
2,
0,
IncompleteReviewPolicy::Block
));
}
#[test]
fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
assert!(!round_is_clean(
1,
true,
1,
2,
0,
IncompleteReviewPolicy::Warn
));
}
#[test]
fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
assert!(round_is_clean(
0,
true,
1,
2,
0,
IncompleteReviewPolicy::Warn
));
}
#[test]
fn a_full_panel_with_an_open_finding_is_not_clean() {
assert!(!round_is_clean(
1,
true,
2,
2,
0,
IncompleteReviewPolicy::Block
));
}
#[test]
fn a_full_panel_with_a_red_e2e_is_not_clean() {
assert!(!round_is_clean(
0,
false,
2,
2,
0,
IncompleteReviewPolicy::Block
));
}
#[test]
fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
assert!(round_is_clean(
0,
true,
1,
2,
1,
IncompleteReviewPolicy::Block
));
}
#[test]
fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
assert!(!round_is_clean(
0,
true,
1,
2,
0,
IncompleteReviewPolicy::Block
));
}
#[test]
fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
assert!(!round_is_clean(
1,
true,
1,
2,
1,
IncompleteReviewPolicy::Block
));
assert!(!round_is_clean(
0,
false,
1,
2,
1,
IncompleteReviewPolicy::Block
));
}
#[test]
fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
assert!(!round_is_clean(
0,
true,
0,
2,
2,
IncompleteReviewPolicy::Block
));
}
fn review_round(
clean: bool,
blocking: usize,
answered: usize,
expected: usize,
progressed: bool,
e2e_ok: bool,
) -> ReviewRound {
ReviewRound {
round: 1,
head: "h".to_owned(),
verified_head: None,
reviews: Vec::new(),
e2e: vec![CommandOutcome {
command: "test".to_owned(),
code: Some(if e2e_ok { 0 } else { 1 }),
output_tail: String::new(),
duration_ms: 0,
}],
verify_retried: false,
e2e_deferred: false,
e2e_defer_reason: None,
fix: None,
blocking,
answered,
expected,
clean,
progressed,
vote_split: false,
reconsideration: Vec::new(),
verdict: None,
}
}
#[test]
fn review_conclusion_is_none_when_nothing_has_run() {
assert_eq!(review_conclusion(&[], 3), None);
}
#[test]
fn review_conclusion_is_none_while_rounds_remain() {
let rounds = vec![review_round(false, 1, 2, 2, true, true)];
assert_eq!(review_conclusion(&rounds, 3), None);
}
#[test]
fn review_conclusion_is_gating_once_a_round_is_clean() {
let rounds = vec![review_round(true, 0, 2, 2, false, true)];
assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
}
#[test]
fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
let rounds = vec![
review_round(false, 1, 2, 2, true, true),
review_round(false, 1, 2, 2, true, true),
];
assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
}
#[test]
fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
let rounds = vec![
review_round(false, 1, 2, 2, true, true),
review_round(false, 1, 2, 2, true, false),
];
assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
}
#[test]
fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
let rounds = vec![review_round(false, 0, 1, 2, false, true)];
assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
}
#[test]
fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
let rounds = vec![
review_round(false, 1, 2, 2, false, true),
review_round(false, 1, 2, 2, false, true),
];
assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
}
fn secs(n: u64) -> Duration {
Duration::from_secs(n)
}
fn init_repo(dir: &Path) {
let run = |args: &[&str]| {
let out = std::process::Command::new("git")
.args(args)
.current_dir(dir)
.quiet()
.output()
.expect("spawn git");
assert!(
out.status.success(),
"git {args:?} failed: {}",
String::from_utf8_lossy(&out.stderr)
);
};
run(&["init", "-b", "main"]);
run(&["config", "user.name", "magi test"]);
run(&["config", "user.email", "magi@example.com"]);
std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
run(&["add", "-A"]);
run(&["commit", "-m", "init"]);
}
fn ask_test_home() {
crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
}
fn runner_at(status: RunStatus) -> Runner {
let mut state = RunState::new(
PathBuf::from("/nonexistent/repo"),
"main".to_owned(),
"deadbeef".to_owned(),
"task".to_owned(),
Config::default(),
);
state.status = status;
Runner {
state,
roles: ResolvedRoles {
implementers: Vec::new(),
judges: Vec::new(),
reviewers: Vec::new(),
fixer: None,
conductor: conductor(),
},
sem: Arc::new(Semaphore::new(1)),
pause: Pause::new(),
}
}
fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
let mut q = ask::Question::new(
run.to_owned(),
"implement".to_owned(),
"impl-A".to_owned(),
"Which storage backend should the cache use?".to_owned(),
String::new(),
vec!["SQLite".to_owned(), "Redis".to_owned()],
);
store.put(&mut q).unwrap();
q
}
#[test]
fn a_failed_runs_open_question_is_abandoned() {
ask_test_home();
let store = ask::Questions::open();
let mut runner = runner_at(RunStatus::Failed);
let run = runner.state.id.clone();
let q = ask_open_question(&store, &run);
runner.settle_questions();
let back = store.get(&q.id).unwrap();
assert!(
!back.status.open(),
"the seat that asked died with the run; nobody is left to read an answer"
);
assert!(
back.detail.contains(&run) && back.detail.contains("failed"),
"the reason names what the run became, not just that it is gone: {}",
back.detail
);
}
#[test]
fn a_merged_runs_open_question_is_abandoned_too() {
ask_test_home();
let store = ask::Questions::open();
for status in [RunStatus::Merged, RunStatus::Ready] {
let mut runner = runner_at(status);
let run = runner.state.id.clone();
let q = ask_open_question(&store, &run);
runner.settle_questions();
let back = store.get(&q.id).unwrap();
assert!(
!back.status.open(),
"{status:?} run's question must not outlive the run"
);
}
}
#[test]
fn a_still_resumable_runs_open_question_is_left_alone() {
ask_test_home();
let store = ask::Questions::open();
for status in [RunStatus::Blocked, RunStatus::Stalled] {
let mut runner = runner_at(status);
let run = runner.state.id.clone();
let q = ask_open_question(&store, &run);
runner.settle_questions();
let back = store.get(&q.id).unwrap();
assert!(
back.status.open(),
"{status:?} is still alive; the question must still be waiting"
);
}
}
#[test]
fn settle_questions_never_touches_an_already_answered_question() {
ask_test_home();
let store = ask::Questions::open();
let mut runner = runner_at(RunStatus::Failed);
let run = runner.state.id.clone();
let mut q = ask_open_question(&store, &run);
q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
.unwrap();
store.put(&mut q).unwrap();
runner.settle_questions();
runner.settle_questions();
let back = store.get(&q.id).unwrap();
assert_eq!(
back.status,
ask::QuestionStatus::Answered,
"a real answer is a decision on record, never overwritten by a sweep"
);
}
#[tokio::test]
async fn merge_does_not_reattempt_once_a_run_has_concluded() {
let tmp = tempfile::tempdir().expect("tempdir");
let repo = tmp.path().join("repo");
std::fs::create_dir_all(&repo).unwrap();
init_repo(&repo);
let mut config = Config::default();
config.merge.mode = MergeMode::Local;
let mut state = RunState::new(
repo.clone(),
"main".to_owned(),
"deadbeef".to_owned(),
"task".to_owned(),
config,
);
state.candidates = vec![Candidate {
index: 0,
label: 'A',
agent: "alpha".to_owned(),
branch: "does-not-exist".to_owned(),
worktree: repo.clone(),
summary: String::new(),
stat: String::new(),
files: 0,
commits: 0,
empty: false,
failed: None,
duration_ms: 0,
folded: false,
}];
state.tally = Some(Tally {
first_choice: BTreeMap::from([('A', 1)]),
borda: BTreeMap::new(),
winner: 'A',
rankings: 1,
unanimous_initial: true,
deliberated: false,
changed_votes: 0,
unanimous_final: true,
tie_break: None,
judges: 0,
present: 0,
quorum: 0,
met_quorum: true,
uncontested: Some("only candidate A produced a change".to_owned()),
});
state.reviews = vec![ReviewRound {
round: 1,
head: "deadbeef".to_owned(),
verified_head: None,
reviews: Vec::new(),
e2e: Vec::new(),
fix: None,
blocking: 0,
answered: 0,
expected: 0,
clean: true,
verify_retried: false,
e2e_deferred: false,
e2e_defer_reason: None,
progressed: false,
vote_split: false,
reconsideration: Vec::new(),
verdict: None,
}];
state.gate = vec![CommandOutcome {
command: "test".to_owned(),
code: Some(0),
output_tail: String::new(),
duration_ms: 0,
}];
state.status = RunStatus::Ready;
state.merge = Some(MergeOutcome {
mode: MergeMode::Local,
ok: false,
detail: "already concluded".to_owned(),
});
let mut runner = Runner {
state,
roles: ResolvedRoles {
implementers: Vec::new(),
judges: Vec::new(),
reviewers: Vec::new(),
fixer: None,
conductor: conductor(),
},
sem: Arc::new(Semaphore::new(1)),
pause: Pause::new(),
};
runner.merge().await.expect("merge");
assert_eq!(
runner.state.status,
RunStatus::Ready,
"a concluded run's status must not change on reentry"
);
assert_eq!(
runner.state.merge.as_ref().map(|m| m.detail.as_str()),
Some("already concluded"),
"merge must not run again once the node already recorded an outcome"
);
}
#[tokio::test]
async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
let tmp = tempfile::tempdir().expect("tempdir");
let repo = tmp.path().join("repo");
std::fs::create_dir_all(&repo).unwrap();
init_repo(&repo);
let mut config = Config::default();
config.merge.mode = MergeMode::Pr;
config.graph.land = true;
config.graph.land_approval = false;
let mut state = RunState::new(
repo.clone(),
"main".to_owned(),
"deadbeef".to_owned(),
"task".to_owned(),
config,
);
state.candidates = vec![Candidate {
index: 0,
label: 'A',
agent: "alpha".to_owned(),
branch: "does-not-exist".to_owned(),
worktree: repo.clone(),
summary: String::new(),
stat: String::new(),
files: 0,
commits: 0,
empty: false,
failed: None,
duration_ms: 0,
folded: false,
}];
state.tally = Some(Tally {
first_choice: BTreeMap::from([('A', 1)]),
borda: BTreeMap::new(),
winner: 'A',
rankings: 1,
unanimous_initial: true,
deliberated: false,
changed_votes: 0,
unanimous_final: true,
tie_break: None,
judges: 0,
present: 0,
quorum: 0,
met_quorum: true,
uncontested: Some("only candidate A produced a change".to_owned()),
});
state.reviews = vec![ReviewRound {
round: 1,
head: "deadbeef".to_owned(),
verified_head: None,
reviews: Vec::new(),
e2e: Vec::new(),
fix: None,
blocking: 0,
answered: 0,
expected: 0,
clean: true,
verify_retried: false,
e2e_deferred: false,
e2e_defer_reason: None,
progressed: false,
vote_split: false,
reconsideration: Vec::new(),
verdict: None,
}];
state.gate = vec![CommandOutcome {
command: "test".to_owned(),
code: Some(0),
output_tail: String::new(),
duration_ms: 0,
}];
state.status = RunStatus::Landing;
state.merge = Some(MergeOutcome {
mode: MergeMode::Pr,
ok: true,
detail: "https://example.invalid/x/y/pull/1".to_owned(),
});
ask_test_home();
let store = ask::Questions::open();
let q = ask_open_question(&store, &state.id);
let mut runner = Runner {
state,
roles: ResolvedRoles {
implementers: Vec::new(),
judges: Vec::new(),
reviewers: Vec::new(),
fixer: None,
conductor: conductor(),
},
sem: Arc::new(Semaphore::new(1)),
pause: Pause::new(),
};
runner.execute().await.expect("execute");
assert_eq!(
runner.state.merge.as_ref().map(|m| m.detail.as_str()),
Some("https://example.invalid/x/y/pull/1"),
"reentry must not push again or open a second pull request over the \
one `land` is already watching"
);
assert_ne!(
runner.state.status,
RunStatus::Landing,
"land could not actually reach the fake pull request, so it must \
have given up rather than left the run silently parked forever"
);
assert_eq!(runner.state.status, RunStatus::Blocked);
assert!(
store.get(&q.id).unwrap().status.open(),
"Blocked is still alive; settle_questions must have been a no-op here"
);
}
fn state_with_round(round: ReviewRound) -> RunState {
let mut s = RunState::new(
PathBuf::from("/repo"),
"main".to_owned(),
"abc1234".to_owned(),
"add retries".to_owned(),
Config::default(),
);
s.reviews = vec![round];
s
}
fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
crate::verdict::Finding {
id: id.to_owned(),
severity,
file: None,
line: None,
title: title.to_owned(),
detail: String::new(),
}
}
#[test]
fn pr_body_names_open_findings_and_declined_ones() {
let round = ReviewRound {
round: 2,
head: "deadbee".to_owned(),
verified_head: None,
reviews: vec![ReviewRecord {
reviewer: 1,
agent: "alpha".to_owned(),
summary: String::new(),
findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
vote: None,
failed: None,
duration_ms: 0,
}],
e2e: vec![CommandOutcome {
command: "cargo test".to_owned(),
code: Some(0),
output_tail: String::new(),
duration_ms: 0,
}],
verify_retried: false,
e2e_deferred: false,
e2e_defer_reason: None,
fix: Some(FixRecord {
agent: "alpha".to_owned(),
addressed: Vec::new(),
rejected: vec![crate::verdict::Rejection {
id: "R1-1-1".to_owned(),
why: "not reachable from any caller".to_owned(),
}],
notes: String::new(),
committed: true,
failed: None,
duration_ms: 0,
}),
blocking: 0,
answered: 1,
expected: 1,
clean: false,
progressed: true,
vote_split: false,
reconsideration: Vec::new(),
verdict: None,
};
let state = state_with_round(round);
let body = pr_body(&state, 'A');
assert!(body.contains("add retries"), "the task must still be there");
assert!(body.contains("R2-1-1"), "{body}");
assert!(body.contains("unused import"), "{body}");
assert!(body.contains("R1-1-1"), "the declined finding: {body}");
assert!(
body.contains("not reachable from any caller"),
"the reason it was declined: {body}"
);
}
#[test]
fn pr_body_says_nothing_extra_when_the_round_was_clean() {
let round = ReviewRound {
round: 1,
head: "deadbee".to_owned(),
verified_head: None,
reviews: vec![ReviewRecord {
reviewer: 1,
agent: "alpha".to_owned(),
summary: String::new(),
findings: Vec::new(),
vote: None,
failed: None,
duration_ms: 0,
}],
e2e: Vec::new(),
verify_retried: false,
e2e_deferred: false,
e2e_defer_reason: None,
fix: None,
blocking: 0,
answered: 1,
expected: 1,
clean: true,
progressed: false,
vote_split: false,
reconsideration: Vec::new(),
verdict: None,
};
let state = state_with_round(round);
let body = pr_body(&state, 'A');
assert!(!body.contains("Open review findings"), "{body}");
assert!(!body.contains("Declined"), "{body}");
}
#[test]
fn pr_body_titles_itself_from_the_task_not_run_or_candidate() {
let state = RunState::new(
PathBuf::from("/repo"),
"main".to_owned(),
"abc1234".to_owned(),
"add retries".to_owned(),
Config::default(),
);
let body = pr_body(&state, 'A');
let title = body.lines().next().unwrap();
assert_eq!(
title, "add retries",
"the title must be the task, not run/candidate bookkeeping: {body}"
);
assert!(
body.contains(&format!("magi:run/{}", state.id)),
"the run id must still be recoverable from the footer: {body}"
);
assert!(
body.contains("magi:candidate-a"),
"the candidate must still be recoverable from the footer: {body}"
);
}
#[test]
fn pr_body_never_titles_itself_off_a_blank_first_line() {
let leading_blank = RunState::new(
PathBuf::from("/repo"),
"main".to_owned(),
"abc1234".to_owned(),
"\n\n \nadd retries\n\ndetails".to_owned(),
Config::default(),
);
let body = pr_body(&leading_blank, 'A');
assert_eq!(
body.lines().next(),
Some("add retries"),
"a leading blank line must not become an empty title: {body}"
);
let whitespace_only = RunState::new(
PathBuf::from("/repo"),
"main".to_owned(),
"abc1234".to_owned(),
" \n \n".to_owned(),
Config::default(),
);
let body = pr_body(&whitespace_only, 'A');
let title = body.lines().next().unwrap_or_default();
assert!(
!title.is_empty(),
"a whitespace-only instruction must still fall back to a non-empty title: {body}"
);
}
#[test]
fn manual_merge_command_matches_the_configured_style() {
let repo = Path::new("/repo");
let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
assert_eq!(
squash,
"git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
\"Merge magi run 0832 (candidate A)\""
);
let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
}
#[test]
fn a_nudge_gets_a_quarter_of_the_budget() {
assert_eq!(retry_budget(secs(1200), true), secs(300));
assert_eq!(retry_budget(secs(3600), true), secs(900));
}
#[test]
fn a_resent_prompt_keeps_the_whole_budget() {
assert_eq!(retry_budget(secs(1200), false), secs(1200));
assert_eq!(retry_budget(secs(60), false), secs(60));
}
#[test]
fn the_floor_never_exceeds_the_original_budget() {
assert_eq!(retry_budget(secs(60), true), secs(60));
assert_eq!(retry_budget(secs(480), true), secs(120));
assert_eq!(retry_budget(secs(0), true), secs(0));
}
}