1use std::collections::{BTreeMap, BTreeSet};
20use std::path::{Path, PathBuf};
21use std::sync::atomic::{AtomicBool, Ordering};
22use std::sync::{Arc, Mutex};
23use std::time::{Duration, Instant};
24
25use anyhow::{Context as _, Result, bail};
26use jiff::Timestamp;
27use tokio::sync::Semaphore;
28
29use crate::advise;
30use crate::agent::{self, AgentOutput, Invocation, SeatState};
31use crate::ask;
32use crate::blind;
33use crate::bump;
34use crate::config::{
35 AgentSpec, Config, IncompleteReviewPolicy, LeakPolicy, MergeMode, MergeStyle, Prompts,
36 ResolvedRoles,
37};
38use crate::git;
39use crate::land;
40use crate::proc::Quiet as _;
41use crate::prompt::{
42 self, CandidateView, Lens, ReviewPatch, ReviewReconsiderCtx, ReviewSeatReport, Turn,
43};
44use crate::queue;
45use crate::run::{
46 BaseSync, Candidate, CommandOutcome, ContinuationOutcome, ContinuationRecord,
47 DeliberationRound, DeliberationTurn, E2eStatus, FixRecord, GateFixRecord, JobRecord, JobStatus,
48 Judgement, MergeOutcome, OperatorFixFinding, OperatorFixOutcome, OperatorFixRequest, QuotaLoss,
49 ReviewRecord, ReviewRevoteRecord, ReviewRound, RunState, RunStatus, Tally, VoteRecord, tail,
50 write_artifact,
51};
52use crate::verdict::{
53 self, FinalVote, Finding, FixReport, Position, Proposal, Ranking, Review, ReviewRevote,
54 ReviewVote, Severity,
55};
56
57const OUTPUT_TAIL: usize = 8_000;
59
60const EVENT_OUTPUT_TAIL: usize = 2_000;
63
64const LEASE_RELEASE_POLL: Duration = Duration::from_secs(1);
67
68const LEASE_RELEASE_MAX_WAIT: Duration = Duration::from_secs(30);
85
86pub(crate) const STAGNANT_LIMIT: usize = 2;
100
101const BASE_SYNC_ROUNDS: usize = 4;
114
115const MAX_FIX_CONTINUATIONS: usize = 2;
131
132#[derive(Clone)]
138struct SeatJob {
139 spec: AgentSpec,
140 seat: SeatState,
141 cwd: PathBuf,
142 prompt: String,
143 timeout: Duration,
144 allow_write: bool,
145 sessions: bool,
146 artifacts: PathBuf,
147 stem: String,
148}
149
150enum AgentOutcome {
162 Ok(AgentOutput),
164 Quota(AgentOutput),
166 Dropped(AgentOutput),
169 Failed(String),
171}
172
173#[derive(Debug, Clone, Default)]
200pub struct Pause(Arc<AtomicBool>, Arc<Mutex<Option<String>>>);
201
202impl Pause {
203 #[must_use]
205 pub fn new() -> Self {
206 Self::default()
207 }
208
209 pub fn park(&self) {
211 self.0.store(true, Ordering::SeqCst);
212 }
213
214 pub fn park_because(&self, reason: impl Into<String>) {
220 let mut reason_guard = self
221 .1
222 .lock()
223 .unwrap_or_else(std::sync::PoisonError::into_inner);
224 if reason_guard.is_none() {
225 *reason_guard = Some(reason.into());
226 }
227 drop(reason_guard);
228 self.park();
229 }
230
231 #[must_use]
233 pub fn parked(&self) -> bool {
234 self.0.load(Ordering::SeqCst)
235 }
236
237 #[must_use]
239 pub fn reason(&self) -> Option<String> {
240 self.1
241 .lock()
242 .unwrap_or_else(std::sync::PoisonError::into_inner)
243 .clone()
244 }
245}
246
247pub struct Runner {
249 pub state: RunState,
251 roles: ResolvedRoles,
252 sem: Arc<Semaphore>,
253 pause: Pause,
257 interrupt: Pause,
263}
264
265async fn resolve_base(repo: &Path, base_branch: &str, remote: &str) -> Result<String> {
285 let tracking = format!("{remote}/{base_branch}");
286 let fetched = git::fetch(repo, remote, base_branch).await;
287 if let Ok(out) = &fetched
288 && out.ok()
289 && git::rev_exists(repo, &tracking).await
290 {
291 return git::rev_parse(repo, &tracking).await;
292 }
293 let why = match &fetched {
294 Ok(out) if !out.ok() => out.stderr.lines().next().unwrap_or("").to_owned(),
295 Ok(_) => format!("{remote} has no {base_branch}"),
296 Err(e) => e.to_string(),
297 };
298 tracing::warn!(
299 "could not read {tracking} ({why}); branching off the local \
300 {base_branch} instead, which may be behind"
301 );
302 git::rev_parse(repo, base_branch).await.with_context(|| {
303 format!(
304 "cannot resolve `{base_branch}`; set [merge] base in magi.toml to a \
305 branch that exists"
306 )
307 })
308}
309
310struct FixClaim {
327 path: PathBuf,
328}
329
330impl FixClaim {
331 fn acquire(dir: &Path) -> Result<Self> {
332 std::fs::create_dir_all(dir).with_context(|| format!("create {}", dir.display()))?;
333 let path = dir.join("fix.lock");
334 match Self::create(&path) {
335 Ok(claim) => Ok(claim),
336 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
337 if Self::reclaim_if_dead(&path) {
338 Self::create(&path).with_context(|| format!("lock {}", path.display()))
339 } else {
340 bail!(
341 "another `magi fix` is already running for this run ({} exists)",
342 path.display()
343 )
344 }
345 }
346 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
347 }
348 }
349
350 fn create(path: &Path) -> std::io::Result<Self> {
351 let mut f = std::fs::OpenOptions::new()
352 .write(true)
353 .create_new(true)
354 .open(path)?;
355 use std::io::Write as _;
356 writeln!(f, "{}", std::process::id())?;
358 Ok(Self {
359 path: path.to_owned(),
360 })
361 }
362
363 fn reclaim_if_dead(path: &Path) -> bool {
367 let dead = std::fs::read_to_string(path)
368 .ok()
369 .and_then(|body| body.trim().parse::<u32>().ok())
370 .is_some_and(|pid| !crate::proc::pid_alive(pid));
371 dead && std::fs::remove_file(path).is_ok()
372 }
373}
374
375impl Drop for FixClaim {
376 fn drop(&mut self) {
377 let _ = std::fs::remove_file(&self.path);
378 }
379}
380
381impl Runner {
382 pub async fn start(repo: &Path, instruction: String, config: Config) -> Result<Self> {
384 let repo = git::toplevel(repo).await?;
385 let missing = agent::missing_programs(&config.agents);
386 if !missing.is_empty() {
387 bail!(
388 "these agent programs are not on PATH: {}. Fix the roster in \
389 magi.toml or install them.",
390 missing.join(", ")
391 );
392 }
393 let base_branch = match config.merge.base.clone() {
394 Some(b) => b,
395 None => git::current_branch(&repo)
396 .await?
397 .context("HEAD is detached; set [merge] base in magi.toml")?,
398 };
399 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
400 if !git::is_clean(&repo).await? {
404 tracing::warn!(
405 "{} has uncommitted changes; they are not part of this run, \
406 which branches off {base_branch} ({})",
407 repo.display(),
408 &base_commit[..base_commit.len().min(8)]
409 );
410 }
411 let roles = config.resolve_roles()?;
412 let max_parallel = config.graph.max_parallel.max(1);
413 let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
414 state.event("start", format!("run {} created", state.id));
415 state.save()?;
416 Ok(Self {
417 state,
418 roles,
419 sem: Arc::new(Semaphore::new(max_parallel)),
420 pause: Pause::new(),
421 interrupt: Pause::new(),
422 })
423 }
424
425 pub async fn review(repo: &Path, branch: &str, config: Config) -> Result<Self> {
439 let repo = git::toplevel(repo).await?;
440 let missing = agent::missing_programs(&config.agents);
441 if !missing.is_empty() {
442 bail!(
443 "these agent programs are not on PATH: {}. Fix the roster in \
444 magi.toml or install them.",
445 missing.join(", ")
446 );
447 }
448 if !git::branch_exists(&repo, branch).await? {
449 bail!("no branch `{branch}` in {}", repo.display());
450 }
451 let base_branch = match config.merge.base.clone() {
452 Some(b) => b,
453 None => git::current_branch(&repo)
454 .await?
455 .context("HEAD is detached; set [merge] base in magi.toml")?,
456 };
457 if base_branch == branch {
458 bail!("`{branch}` is the base branch; there is nothing to review against");
459 }
460 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
461
462 let roles = config.resolve_roles()?;
463 let max_parallel = config.graph.max_parallel.max(1);
464 let log = git::log_oneline(&repo, &base_commit, branch)
467 .await
468 .unwrap_or_default();
469 let instruction = format!(
470 "Review the work already on branch `{branch}`. There is no task \
471 statement: what the change claims to do is whatever its commits \
472 say.\n\n{}",
473 if log.trim().is_empty() {
474 "(no commit messages)"
475 } else {
476 log.trim()
477 }
478 );
479 let mut state = RunState::new(
480 repo.clone(),
481 base_branch,
482 base_commit.clone(),
483 instruction,
484 config,
485 );
486
487 let worktree = state.worktree_root().join("under-review");
490 if let Some(parent) = worktree.parent() {
491 tokio::fs::create_dir_all(parent).await.ok();
492 }
493 let path = worktree.to_string_lossy().to_string();
494 git::git(&repo, &["worktree", "add", &path, branch])
495 .await
496 .with_context(|| {
497 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
498 })?;
499
500 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
501 .await
502 .unwrap_or(0);
503 if commits == 0 {
504 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
505 }
506 let files = git::changed_files(&worktree, &base_commit, "HEAD")
507 .await
508 .map(|f| f.len())
509 .unwrap_or(0);
510 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
511 .await
512 .unwrap_or_default();
513
514 state.candidates.push(Candidate {
515 index: 0,
516 label: 'A',
517 agent: "(existing branch)".to_owned(),
520 branch: branch.to_owned(),
521 worktree,
522 summary: String::new(),
523 stat,
524 files,
525 commits,
526 empty: false,
527 failed: None,
528 verified_noop: None,
529 duration_ms: 0,
530 folded: false,
531 });
532 state.tally = Some(Tally {
533 first_choice: BTreeMap::from([('A', 0)]),
534 borda: BTreeMap::new(),
535 winner: 'A',
536 rankings: 0,
537 unanimous_initial: false,
538 deliberated: false,
539 changed_votes: 0,
540 unanimous_final: false,
541 tie_break: None,
542 judges: 0,
546 present: 0,
547 quorum: 0,
548 met_quorum: true,
549 uncontested: Some("review-only run: nothing competed".to_owned()),
550 });
551 state.status = RunStatus::Reviewing;
552 state.event(
553 "start",
554 format!(
555 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
556 state.id
557 ),
558 );
559 state.save()?;
560 Ok(Self {
561 state,
562 roles,
563 sem: Arc::new(Semaphore::new(max_parallel)),
564 pause: Pause::new(),
565 interrupt: Pause::new(),
566 })
567 }
568
569 pub fn resume(id: &str) -> Result<Self> {
571 let state = RunState::load(id)?;
572 let roles = state.config.resolve_roles()?;
573 let max_parallel = state.config.graph.max_parallel.max(1);
574 Ok(Self {
575 state,
576 roles,
577 sem: Arc::new(Semaphore::new(max_parallel)),
578 pause: Pause::new(),
579 interrupt: Pause::new(),
580 })
581 }
582
583 pub async fn execute(&mut self) -> Result<()> {
590 let result = self.execute_graph().await;
591 let ended = if result.is_err() {
592 Some(crate::notices::run_stopped(&self.state.id, &self.state))
593 } else {
594 crate::notices::run_ended(&self.state)
595 };
596 if let Some(notice) = ended {
597 crate::notices::raise(notice);
598 }
599 result
600 }
601
602 async fn execute_graph(&mut self) -> Result<()> {
603 self.state.parked = false;
608 self.state.clear_active();
615 let pid = std::process::id();
630 self.state.driver_pid = Some(pid);
631 self.state.driver_started_at = crate::proc::process_started_at(pid);
632 self.state.save()?;
633 if self.state.status == RunStatus::Stalled {
646 if self.recover_stall().await? {
647 self.finish_after_tally().await?;
648 } else {
649 self.state.save()?;
651 }
652 return Ok(());
653 }
654 if self.state.status == RunStatus::Landing {
664 self.run_land().await?;
665 self.settle_questions();
670 return Ok(());
671 }
672 self.prep().await?;
673 if self.park_here()? {
674 return Ok(());
675 }
676 self.advise().await?;
677 if self.park_here()? {
678 return Ok(());
679 }
680 self.implement().await?;
681 if self.park_here()? {
682 return Ok(());
683 }
684 if self.state.status == RunStatus::VerifiedNoop {
688 return Ok(());
689 }
690 self.judge().await?;
691 if self.park_here()? {
692 return Ok(());
693 }
694 self.deliberate().await?;
695 if self.park_here()? {
696 return Ok(());
697 }
698 self.vote().await?;
699 if self.park_here()? {
700 return Ok(());
701 }
702 self.tally()?;
703 if self.state.status == RunStatus::Stalled {
708 self.state.save()?;
712 return Ok(());
713 }
714 self.finish_after_tally().await?;
715 Ok(())
716 }
717
718 fn park_here(&mut self) -> Result<bool> {
725 if !self.pause.parked() && !self.interrupt.parked() {
730 return Ok(false);
731 }
732 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
733 Some(reason) => format!(
734 "parked after `{}` ({reason}) — resume to carry on from here",
735 self.state.status.as_str()
736 ),
737 None => format!(
738 "parked after `{}` — resume to carry on from here",
739 self.state.status.as_str()
740 ),
741 };
742 self.state.event("park", why);
743 self.state.parked = true;
744 self.state.save()?;
745 Ok(true)
746 }
747
748 pub fn on_pause(&mut self, pause: Pause) {
750 self.pause = pause;
751 }
752
753 pub fn watch_interrupt(&mut self, pause: Pause) {
759 self.interrupt = pause;
760 }
761
762 fn settle_questions(&mut self) {
782 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
783 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
784 }
785 }
786
787 async fn finish_after_tally(&mut self) -> Result<()> {
790 self.fold_losers().await?;
791 self.sync_to_base().await?;
796 self.review_loop().await?;
797 self.sync_to_base().await?;
798 self.gate().await?;
799 self.merge().await?;
800 self.state.save()?;
801 Ok(())
802 }
803
804 async fn prep(&mut self) -> Result<()> {
807 if !self.state.candidates.is_empty() {
808 return Ok(());
809 }
810 self.state.status = RunStatus::Prep;
811 let repo = self.state.repo.clone();
812 let base = self.state.base_commit.clone();
813 let root = self.state.worktree_root();
814 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
815
816 let hooks_dir = self.state.dir().join("hooks");
819 if self.state.config.blind.commit_msg_hook {
820 std::fs::create_dir_all(&hooks_dir)
821 .with_context(|| format!("create {}", hooks_dir.display()))?;
822 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
823 let path = hooks_dir.join("commit-msg");
824 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
825 make_executable(&path)?;
826 git::acquire_worktree_config(&repo).await?;
834 self.state.enabled_worktree_config = true;
835 }
836
837 for (index, (spec, label)) in self
838 .roles
839 .implementers
840 .clone()
841 .into_iter()
842 .zip(labels)
843 .enumerate()
844 {
845 let branch = self.state.branch_for(label);
846 let worktree = root.join(format!("cand-{label}"));
847 git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
848 if self.state.config.blind.commit_msg_hook {
849 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
850 }
851 git::local_exclude(&worktree, "/.magi/").await?;
852 self.state.candidates.push(Candidate {
853 index,
854 label,
855 agent: spec.id.clone(),
856 branch,
857 worktree,
858 summary: String::new(),
859 stat: String::new(),
860 files: 0,
861 commits: 0,
862 empty: false,
863 failed: None,
864 verified_noop: None,
865 duration_ms: 0,
866 folded: false,
867 });
868 }
869
870 for j in 1..=self.roles.judges.len() {
871 let wt = root.join(format!("judge-{j}"));
872 if !wt.exists() {
873 git::worktree_add_detached(&repo, &wt, &base).await?;
874 }
875 }
876
877 if self.state.config.graph.advise {
885 for k in 1..=self.state.config.graph.advisors {
886 let wt = root.join(format!("advisor-{k}"));
887 if !wt.exists() {
888 git::worktree_add_detached(&repo, &wt, &base).await?;
889 }
890 }
891 }
892
893 let authors: Vec<&str> = self
898 .roles
899 .implementers
900 .iter()
901 .map(|a| a.id.as_str())
902 .collect();
903 let overlap: Vec<String> = self
904 .roles
905 .judges
906 .iter()
907 .enumerate()
908 .filter(|(_, j)| authors.contains(&j.id.as_str()))
909 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
910 .collect();
911 if !overlap.is_empty() {
912 let note = format!(
913 "{} also authored a candidate; blind, but the panel is less \
914 independent than {} distinct agents would be",
915 overlap.join(", "),
916 self.roles.judges.len()
917 );
918 self.state.event("prep", note);
919 }
920
921 self.state.event(
922 "prep",
923 format!(
924 "{} candidates, {} judges, base {} ({})",
925 self.state.candidates.len(),
926 self.roles.judges.len(),
927 &self.state.base_commit[..7.min(self.state.base_commit.len())],
928 self.state.base_branch
929 ),
930 );
931 self.state.status = RunStatus::Implementing;
932 self.state.save()?;
933 Ok(())
934 }
935
936 async fn advise(&mut self) -> Result<()> {
969 let implement_untouched = self
970 .state
971 .candidates
972 .iter()
973 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
974 if !self.state.config.graph.advise || self.state.advise_attempted {
975 return Ok(());
976 }
977 if !implement_untouched {
978 self.state.event(
979 "advise",
980 "skipping the design-deliberation stage: at least one \
981 candidate already shows implementation progress, so this \
982 run is past the point the stage exists to run before"
983 .to_owned(),
984 );
985 self.state.advise_attempted = true;
986 self.state.save()?;
987 return Ok(());
988 }
989 let run_id = self.state.id.clone();
990 let prompts = self.state.config.prompts.clone();
991 let instruction = self.state.instruction.clone();
992 let language = self.state.config.graph.language.clone();
993 let root = self.state.worktree_root();
994 let n = self.state.config.graph.advisors;
995 let where_recorded = self.state.dir().join("run.json");
996
997 let seats = match self.state.config.advisors() {
998 Ok(seats) if !seats.is_empty() => seats,
999 Ok(_) => {
1000 self.state.event(
1001 "advise",
1002 format!(
1003 "[graph] advisors is 0; skipping the design-deliberation \
1004 stage and continuing without a synthesis brief (see {})",
1005 where_recorded.display()
1006 ),
1007 );
1008 self.state.advise_attempted = true;
1009 self.state.save()?;
1010 return Ok(());
1011 }
1012 Err(e) => {
1013 self.state.event(
1014 "advise",
1015 format!(
1016 "could not resolve advisor seats ({e:#}); continuing \
1017 without a design-deliberation brief (see {})",
1018 where_recorded.display()
1019 ),
1020 );
1021 self.state.advise_attempted = true;
1022 self.state.save()?;
1023 return Ok(());
1024 }
1025 };
1026
1027 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1028 let artifacts = agent::artifacts_dir(&self.state.dir());
1029 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1030
1031 let mut jobs = Vec::new();
1032 for (i, spec) in seats.iter().cloned().enumerate() {
1033 let seat_key = format!("advisor-{}", i + 1);
1034 let seat = self.seat(&seat_key, &spec.id);
1035 jobs.push(SeatJob {
1036 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1037 spec,
1038 seat,
1039 cwd: worktrees[i % worktrees.len()].clone(),
1040 timeout,
1041 allow_write: false,
1042 sessions: false,
1043 artifacts: artifacts.clone(),
1044 stem: seat_key,
1045 });
1046 }
1047
1048 self.state.event(
1049 "advise",
1050 format!(
1051 "{} advisor seat(s) sketching a design in parallel",
1052 jobs.len()
1053 ),
1054 );
1055 let mut quota_losses = Vec::new();
1056 let cache = self.state.config.cache_dir();
1057 let ctx = WaveCtx {
1058 run: &run_id,
1059 node: "advise",
1060 prompts: &prompts,
1061 cache: cache.as_deref(),
1062 round: None,
1063 };
1064 let results = ask_json_wave::<Proposal>(
1065 jobs,
1066 Arc::clone(&self.sem),
1067 self.state.config.graph.retries,
1068 &ctx,
1069 &mut quota_losses,
1070 &mut self.state,
1071 &|p: &Proposal| p.validate(),
1072 )
1073 .await;
1074 self.state.quota.extend(quota_losses);
1075
1076 let mut records = Vec::with_capacity(results.len());
1077 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1078 let agent_id = seat.agent.clone();
1079 self.state.seats.insert(seat.key.clone(), seat);
1080 match res {
1081 Ok((proposal, out)) => {
1082 self.state
1083 .event("advise", format!("advisor-{} proposed a design", i + 1));
1084 records.push(advise::AdvisorRecord::proposed(
1085 i + 1,
1086 agent_id,
1087 proposal,
1088 out.duration_ms,
1089 ));
1090 }
1091 Err(e) => {
1092 self.state.event(
1093 "advise",
1094 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1095 );
1096 records.push(advise::AdvisorRecord::failed(
1097 i + 1,
1098 agent_id,
1099 e.to_string(),
1100 ));
1101 }
1102 }
1103 }
1104
1105 let mut advice = advise::Advice {
1106 records,
1107 synthesis: None,
1108 };
1109 if advice.proposals().is_empty() {
1110 self.state.event(
1111 "advise",
1112 "no advisor produced a usable proposal; continuing without a \
1113 synthesis brief"
1114 .to_owned(),
1115 );
1116 } else {
1117 match self
1118 .synthesize_brief(
1119 &advice,
1120 &instruction,
1121 &language,
1122 &worktrees[0],
1123 &artifacts,
1124 &run_id,
1125 &prompts,
1126 cache.as_deref(),
1127 )
1128 .await
1129 {
1130 Ok(Some(text)) => {
1131 self.state.event(
1132 "advise",
1133 "synthesized a design brief for the implementer".to_owned(),
1134 );
1135 advice.synthesis = Some(text);
1136 }
1137 Ok(None) => {
1138 self.state.event(
1139 "advise",
1140 "the synthesis seat produced nothing usable; continuing \
1141 without a design brief"
1142 .to_owned(),
1143 );
1144 }
1145 Err(e) => {
1146 self.state.event(
1147 "advise",
1148 format!("could not synthesize a design brief: {e:#}"),
1149 );
1150 }
1151 }
1152 }
1153 advise::apply_reflection(&mut advice);
1154
1155 self.state.advice = Some(advice);
1156 self.state.advise_attempted = true;
1157 self.state.save()?;
1158 Ok(())
1159 }
1160
1161 #[allow(clippy::too_many_arguments)]
1173 async fn synthesize_brief(
1174 &mut self,
1175 advice: &advise::Advice,
1176 instruction: &str,
1177 language: &str,
1178 cwd: &Path,
1179 artifacts: &Path,
1180 run_id: &str,
1181 prompts: &Prompts,
1182 cache: Option<&Path>,
1183 ) -> Result<Option<String>> {
1184 let want = self.state.config.roles.synthesizer.as_deref();
1185 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1186 let mut seat = self.seat("advise-synthesis", &spec.id);
1187 let proposals = advice.proposals();
1188 let mut prompt = prompt::with_overlay(
1189 prompt::synthesize_brief(instruction, &proposals, language),
1190 prompts.overlay("advise"),
1191 );
1192 if cache.is_some() {
1193 prompt.push('\n');
1198 prompt.push_str(&prompt::build_cache_note("advise", false));
1199 }
1200 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1201 let out = agent::invoke(
1202 &spec,
1203 &mut seat,
1204 &Invocation {
1205 cwd,
1206 prompt: &prompt,
1207 timeout,
1208 allow_write: false,
1209 sessions: false,
1210 artifacts,
1211 stem: "advise-synthesis",
1212 run: run_id,
1213 node: "advise",
1214 cache_dir: None,
1215 attachments: &[],
1216 },
1217 )
1218 .await?;
1219 self.state.seats.insert(seat.key.clone(), seat);
1220 if !out.usable() {
1221 return Ok(None);
1222 }
1223 let text =
1224 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1225 Ok((!text.trim().is_empty()).then_some(text))
1226 }
1227
1228 async fn implement(&mut self) -> Result<()> {
1231 let run_id = self.state.id.clone();
1236 let prompts = self.state.config.prompts.clone();
1237 let todo: Vec<usize> = self
1238 .state
1239 .candidates
1240 .iter()
1241 .enumerate()
1242 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1243 .map(|(i, _)| i)
1244 .collect();
1245 if todo.is_empty() {
1246 return self.after_implement();
1247 }
1248 self.state.status = RunStatus::Implementing;
1249
1250 let language = self.state.config.graph.language.clone();
1251 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1252 let sessions = self.state.config.graph.sessions;
1253 let artifacts = agent::artifacts_dir(&self.state.dir());
1254 let brief = self
1258 .state
1259 .advice
1260 .as_ref()
1261 .and_then(|a| a.synthesis.as_deref())
1262 .map(str::to_owned);
1263
1264 let mut jobs = Vec::new();
1265 for &i in &todo {
1266 let (index, label, worktree) = {
1267 let c = &self.state.candidates[i];
1268 (c.index, c.label, c.worktree.clone())
1269 };
1270 let spec = self.roles.implementers[index].clone();
1271 let seat_key = format!("impl-{label}");
1272 let seat = self.seat(&seat_key, &spec.id);
1273 let instruction = self.state.instruction.clone();
1274 jobs.push(SeatJob {
1275 spec,
1276 seat,
1277 prompt: prompt::implement(
1278 &instruction,
1279 &worktree.to_string_lossy(),
1280 &language,
1281 brief.as_deref(),
1282 ),
1283 cwd: worktree,
1284 timeout,
1285 allow_write: true,
1286 sessions,
1287 artifacts: artifacts.clone(),
1288 stem: format!("impl-{label}"),
1289 });
1290 }
1291
1292 self.state.event(
1293 "implement",
1294 format!("{} candidates in parallel", jobs.len()),
1295 );
1296 let mut sent = jobs.clone();
1302 let cache = self.state.config.cache_dir();
1303 let ctx = WaveCtx {
1304 run: &run_id,
1305 node: "implement",
1306 prompts: &prompts,
1307 cache: cache.as_deref(),
1308 round: None,
1309 };
1310 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1311 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1312 .await;
1313 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1314 .await;
1315 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1316 .await;
1317
1318 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1319 let seat_key = seat.key.clone();
1320 let agent = seat.agent.clone();
1329 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1330 self.state.seats.insert(seat.key.clone(), seat);
1331 let label = self.state.candidates[i].label;
1332 let worktree = self.state.candidates[i].worktree.clone();
1333 let base = self.state.base_commit.clone();
1334
1335 let (summary, duration, failed, verified_claim) = match out {
1336 AgentOutcome::Ok(o) => {
1337 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1338 let failed = (!o.usable()).then(|| {
1339 if o.timed_out {
1340 "agent timed out".to_owned()
1341 } else {
1342 format!("agent exited with {:?}", o.exit_code)
1343 }
1344 });
1345 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1346 (text, o.duration_ms, failed, verified_claim)
1347 }
1348 AgentOutcome::Dropped(o) => {
1354 let why = o
1355 .dropped
1356 .as_ref()
1357 .map(|d| d.why.as_str())
1358 .unwrap_or("the CLI ended the stream without delivering its answer");
1359 (
1360 String::new(),
1361 o.duration_ms,
1362 Some(format!("the CLI dropped the stream ({why})")),
1363 None,
1364 )
1365 }
1366 AgentOutcome::Quota(o) => {
1367 self.state.quota.push(QuotaLoss {
1368 seat: seat_key,
1369 node: "implement".to_owned(),
1370 at: Timestamp::now(),
1371 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1372 });
1373 (
1374 String::new(),
1375 o.duration_ms,
1376 Some("rate limited (quota); produced no change".to_owned()),
1377 None,
1378 )
1379 }
1380 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1381 };
1382
1383 let rescued = match git::rescue_commit(
1386 &worktree,
1387 &format!("magi: candidate {label} (uncommitted work)"),
1388 )
1389 .await
1390 {
1391 Ok(r) => {
1392 self.state.note_withheld("implement", &r.withheld);
1393 r.committed
1394 }
1395 Err(_) => false,
1396 };
1397 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1398 .await
1399 .unwrap_or(0);
1400 let patch = git::diff(&worktree, &base, "HEAD")
1401 .await
1402 .unwrap_or_default();
1403 let stat = git::diff_stat(&worktree, &base, "HEAD")
1404 .await
1405 .unwrap_or_default();
1406 let files = git::changed_files(&worktree, &base, "HEAD")
1407 .await
1408 .map(|f| f.len())
1409 .unwrap_or(0);
1410 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1411
1412 let c = &mut self.state.candidates[i];
1413 if !exhausted_the_fallback_chain {
1414 c.agent = agent;
1415 }
1416 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1417 c.stat = stat;
1418 c.files = files;
1419 c.commits = commits;
1420 c.duration_ms = duration;
1421 c.empty = commits == 0 || patch.trim().is_empty();
1422 c.failed = match failed {
1425 Some(_) if c.empty => failed,
1426 _ => None,
1427 };
1428 c.verified_noop = if c.empty { verified_claim } else { None };
1433 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1434 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1435 (None, true, Some(_), _) => {
1436 format!("candidate {label}: no change produced (agent-verified no-op)")
1437 }
1438 (None, true, None, _) => format!("candidate {label}: no change produced"),
1439 (None, false, _, true) => {
1440 format!(
1441 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1442 )
1443 }
1444 (None, false, _, false) => {
1445 format!("candidate {label}: {files} files, {commits} commits")
1446 }
1447 };
1448 self.state.event("implement", note);
1449 self.state.save()?;
1450 }
1451
1452 self.after_implement()
1453 }
1454
1455 async fn resume_undelivered(
1483 &mut self,
1484 results: &mut [(usize, SeatState, AgentOutcome)],
1485 sent: &[SeatJob],
1486 prompts: &Prompts,
1487 run_id: &str,
1488 ) {
1489 for (wi, seat, out) in results.iter_mut() {
1490 let Some(dropped) = (match &*out {
1491 AgentOutcome::Dropped(o) => o.dropped.clone(),
1492 _ => None,
1493 }) else {
1494 continue;
1495 };
1496 let Some(job) = sent.get(*wi) else { continue };
1497 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1499 self.state.event(
1500 "implement",
1501 format!(
1502 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1503 work is in the tree",
1504 seat.key, dropped.output_tokens, dropped.why
1505 ),
1506 );
1507 continue;
1508 }
1509 if !has_context(&job.spec, seat, job.sessions) {
1517 self.state.event(
1518 "implement",
1519 format!(
1520 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1521 is no session left to resume",
1522 seat.key, dropped.output_tokens, dropped.why
1523 ),
1524 );
1525 continue;
1526 }
1527 self.state.event(
1528 "implement",
1529 format!(
1530 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1531 conversation",
1532 seat.key, dropped.output_tokens, dropped.why
1533 ),
1534 );
1535 let mut retry = job.clone();
1536 retry.seat = seat.clone();
1537 retry.prompt = prompt::resume_after_drop(&dropped.why);
1538 retry.timeout = retry_budget(job.timeout, true);
1539 retry.stem = format!("{}-resume", job.stem);
1540 let cache = self.state.config.cache_dir();
1541 let ctx = WaveCtx {
1542 run: run_id,
1543 node: "implement",
1544 prompts,
1545 cache: cache.as_deref(),
1546 round: None,
1547 };
1548 let (resumed_seat, resumed) =
1549 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1550 *seat = resumed_seat;
1551 *out = resumed;
1552 }
1553 }
1554
1555 async fn resume_quota_losses(
1617 &mut self,
1618 results: &mut [(usize, SeatState, AgentOutcome)],
1619 sent: &mut [SeatJob],
1620 prompts: &Prompts,
1621 run_id: &str,
1622 ) {
1623 let instruction = self.state.instruction.clone();
1624 let language = self.state.config.graph.language.clone();
1625 let brief = self
1626 .state
1627 .advice
1628 .as_ref()
1629 .and_then(|a| a.synthesis.as_deref())
1630 .map(str::to_owned);
1631 for (wi, seat, out) in results.iter_mut() {
1632 let Some(job) = sent.get_mut(*wi) else {
1633 continue;
1634 };
1635 let start = self
1640 .roles
1641 .implementer_roster
1642 .iter()
1643 .position(|s| s.id == job.spec.id)
1644 .unwrap_or(0);
1645 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1646 let mut fallback_attempt = 0usize;
1647 while matches!(&*out, AgentOutcome::Quota(_)) {
1648 let Some(next) =
1649 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1650 .cloned()
1651 else {
1652 break;
1653 };
1654 tried.insert(next.id.clone());
1655 fallback_attempt += 1;
1656
1657 if let Ok(r) = git::rescue_commit(
1658 &job.cwd,
1659 &format!(
1660 "magi: candidate {} (uncommitted work before quota fallback)",
1661 seat.key
1662 ),
1663 )
1664 .await
1665 {
1666 self.state.note_withheld("implement", &r.withheld);
1667 }
1668
1669 self.state.event(
1670 "implement",
1671 format!(
1672 "{}: rate limited (quota) on {}; retrying with {}",
1673 seat.key, seat.agent, next.id
1674 ),
1675 );
1676
1677 let new_seat = self.seat(&seat.key, &next.id);
1678 job.spec = next.clone();
1686 let mut retry = job.clone();
1687 retry.seat = new_seat;
1688 retry.prompt = prompt::implement(
1689 &instruction,
1690 &job.cwd.to_string_lossy(),
1691 &language,
1692 brief.as_deref(),
1693 );
1694 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1695 let cache = self.state.config.cache_dir();
1696 let ctx = WaveCtx {
1697 run: run_id,
1698 node: "implement",
1699 prompts,
1700 cache: cache.as_deref(),
1701 round: None,
1702 };
1703 let (fallback_seat, fallback_out) = run_one(
1704 retry,
1705 Arc::clone(&self.sem),
1706 &ctx,
1707 &mut self.state,
1708 fallback_attempt,
1709 )
1710 .await;
1711 *seat = fallback_seat;
1712 *out = fallback_out;
1713 }
1714 }
1715 }
1716
1717 async fn resume_unconfirmed_commands(
1741 &mut self,
1742 results: &mut [(usize, SeatState, AgentOutcome)],
1743 sent: &[SeatJob],
1744 prompts: &Prompts,
1745 run_id: &str,
1746 ) {
1747 for (wi, seat, out) in results.iter_mut() {
1748 let AgentOutcome::Ok(o) = &*out else {
1749 continue;
1750 };
1751 if !has_unconfirmed_command(&o.commands) {
1752 continue;
1753 }
1754 let Some(job) = sent.get(*wi) else { continue };
1755 if !has_context(&job.spec, seat, job.sessions) {
1756 self.state.event(
1757 "implement",
1758 format!(
1759 "{}: the reply named a command whose own CLI never confirmed the exit \
1760 status of, but there is no session left to resume",
1761 seat.key
1762 ),
1763 );
1764 continue;
1765 }
1766 self.state.event(
1767 "implement",
1768 format!(
1769 "{}: the reply named a command whose own CLI never confirmed the exit \
1770 status of; resuming the conversation",
1771 seat.key
1772 ),
1773 );
1774 let mut retry = job.clone();
1775 retry.seat = seat.clone();
1776 retry.prompt = prompt::resume_incomplete(
1777 "a command in your last reply had no confirmed exit status",
1778 );
1779 retry.timeout = retry_budget(job.timeout, true);
1780 retry.stem = format!("{}-confirm", job.stem);
1781 let cache = self.state.config.cache_dir();
1782 let ctx = WaveCtx {
1783 run: run_id,
1784 node: "implement",
1785 prompts,
1786 cache: cache.as_deref(),
1787 round: None,
1788 };
1789 let (resumed_seat, resumed) =
1790 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1791 *seat = resumed_seat;
1792 *out = resumed;
1793 }
1794 }
1795
1796 async fn continue_fix_report(
1817 &mut self,
1818 mut seat: SeatState,
1819 parse_err: String,
1820 job: &SeatJob,
1821 prompts: &Prompts,
1822 run_id: &str,
1823 round: usize,
1824 ) -> (
1825 SeatState,
1826 Option<FixReport>,
1827 Option<String>,
1828 ContinuationRecord,
1829 ) {
1830 let mut last_err = parse_err;
1831 let mut cumulative_wait_ms = 0u64;
1832 let mut attempts = 0usize;
1833 loop {
1834 if !has_context(&job.spec, &seat, job.sessions) {
1835 self.state.event(
1836 "fix",
1837 format!(
1838 "round {round}: fixer's reply had no adoption report ({last_err}); no \
1839 session left to resume into"
1840 ),
1841 );
1842 let outcome = if attempts == 0 {
1843 ContinuationOutcome::NoSession
1844 } else {
1845 ContinuationOutcome::Exhausted
1846 };
1847 return (
1848 seat,
1849 None,
1850 Some(format!("unparsable fix report: {last_err}")),
1851 ContinuationRecord {
1852 attempts,
1853 cumulative_wait_ms,
1854 outcome,
1855 },
1856 );
1857 }
1858 if attempts >= MAX_FIX_CONTINUATIONS {
1859 self.state.event(
1860 "fix",
1861 format!(
1862 "round {round}: fixer's reply still had no adoption report after \
1863 {attempts} continuation(s) ({last_err}); giving up"
1864 ),
1865 );
1866 return (
1867 seat,
1868 None,
1869 Some(format!(
1870 "unparsable fix report after {attempts} continuation(s): {last_err}"
1871 )),
1872 ContinuationRecord {
1873 attempts,
1874 cumulative_wait_ms,
1875 outcome: ContinuationOutcome::Exhausted,
1876 },
1877 );
1878 }
1879 attempts += 1;
1880 self.state.event(
1881 "fix",
1882 format!(
1883 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
1884 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
1885 ),
1886 );
1887 let mut retry = job.clone();
1888 retry.seat = seat.clone();
1889 retry.prompt = prompt::resume_incomplete(&last_err);
1890 retry.timeout = retry_budget(job.timeout, true);
1891 retry.stem = format!("{}-continue{attempts}", job.stem);
1892 let cache = self.state.config.cache_dir();
1893 let ctx = WaveCtx {
1894 run: run_id,
1895 node: "fix",
1896 prompts,
1897 cache: cache.as_deref(),
1898 round: Some(round),
1899 };
1900 let (resumed_seat, resumed_out) = run_one(
1901 retry,
1902 Arc::clone(&self.sem),
1903 &ctx,
1904 &mut self.state,
1905 attempts,
1906 )
1907 .await;
1908 seat = resumed_seat;
1909 match resumed_out {
1910 AgentOutcome::Ok(o) => {
1911 cumulative_wait_ms += o.duration_ms;
1912 match verdict::extract_json::<FixReport>(&o.text) {
1913 Ok(report) if !has_unconfirmed_command(&o.commands) => {
1914 self.state.event(
1915 "fix",
1916 format!(
1917 "round {round}: fixer's adoption report recovered after \
1918 {attempts} continuation(s)"
1919 ),
1920 );
1921 return (
1922 seat,
1923 Some(report),
1924 None,
1925 ContinuationRecord {
1926 attempts,
1927 cumulative_wait_ms,
1928 outcome: ContinuationOutcome::Resumed,
1929 },
1930 );
1931 }
1932 Ok(_) => {
1940 last_err = "the reply parsed, but it reported a command whose own CLI \
1941 never confirmed an exit status"
1942 .to_owned();
1943 }
1944 Err(e) => last_err = e.to_string(),
1945 }
1946 }
1947 AgentOutcome::Quota(o) => {
1948 cumulative_wait_ms += o.duration_ms;
1949 self.state.quota.push(QuotaLoss {
1950 seat: seat.key.clone(),
1951 node: "fix".to_owned(),
1952 at: Timestamp::now(),
1953 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1954 });
1955 self.state.event(
1956 "fix",
1957 format!(
1958 "round {round}: continuation rate limited (quota); not retrying now"
1959 ),
1960 );
1961 return (
1962 seat,
1963 None,
1964 Some("rate limited (quota) while recovering the fix report".to_owned()),
1965 ContinuationRecord {
1966 attempts,
1967 cumulative_wait_ms,
1968 outcome: ContinuationOutcome::QuotaLost,
1969 },
1970 );
1971 }
1972 AgentOutcome::Dropped(o) => {
1973 cumulative_wait_ms += o.duration_ms;
1974 let why = o
1975 .dropped
1976 .as_ref()
1977 .map(|d| d.why.as_str())
1978 .unwrap_or("the CLI ended the stream without delivering its answer");
1979 last_err = format!("the CLI dropped the stream ({why})");
1980 }
1981 AgentOutcome::Failed(e) => last_err = e,
1982 }
1983 }
1984 }
1985
1986 fn after_implement(&mut self) -> Result<()> {
1987 if self.state.leaks.is_empty() {
1989 let cfg = self.state.config.blind.clone();
1990 let mut leaks = Vec::new();
1991 for c in &self.state.candidates {
1992 let Some(patch) =
1993 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
1994 else {
1995 continue;
1996 };
1997 leaks.extend(blind::scan(
1998 &format!("candidate {} patch", c.label),
1999 &patch,
2000 &cfg.vendor_tokens,
2001 ));
2002 }
2003 if !leaks.is_empty() {
2004 let summary = leaks
2005 .iter()
2006 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2007 .collect::<Vec<_>>()
2008 .join(", ");
2009 match cfg.on_leak {
2010 LeakPolicy::Fail => {
2011 self.state.status = RunStatus::Failed;
2012 self.state
2013 .event("blind", format!("vendor text in a patch: {summary}"));
2014 self.state.leaks = leaks;
2015 self.state.save()?;
2016 self.settle_questions();
2017 bail!(
2018 "blind.on_leak = \"fail\" and vendor text reached a \
2019 judged patch: {summary}"
2020 );
2021 }
2022 LeakPolicy::Redact => self.state.event(
2023 "blind",
2024 format!("redacting vendor text for judging: {summary}"),
2025 ),
2026 LeakPolicy::Warn => self.state.event(
2027 "blind",
2028 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2029 ),
2030 }
2031 self.state.leaks = leaks;
2032 }
2033 }
2034
2035 if self.state.viable().is_empty() {
2036 if self.state.all_candidates_verified_noop() {
2037 self.state.status = RunStatus::VerifiedNoop;
2048 self.state.save()?;
2049 self.settle_questions();
2050 return Ok(());
2051 }
2052 self.state.status = RunStatus::Failed;
2053 self.state.save()?;
2054 self.settle_questions();
2055 bail!("no candidate produced a change; nothing to judge");
2056 }
2057 self.state.status = RunStatus::Judging;
2058 self.state.save()?;
2059 Ok(())
2060 }
2061
2062 async fn judge(&mut self) -> Result<()> {
2065 let run_id = self.state.id.clone();
2070 let prompts = self.state.config.prompts.clone();
2071 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2072 return Ok(());
2073 }
2074 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2075 if viable.len() == 1 {
2076 self.state.judge_skipped = true;
2083 self.state.event(
2084 "judge",
2085 format!(
2086 "only candidate {} produced a change; judging skipped",
2087 viable[0].label
2088 ),
2089 );
2090 self.state.save()?;
2091 return Ok(());
2092 }
2093 self.state.status = RunStatus::Judging;
2094
2095 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2096 let language = self.state.config.graph.language.clone();
2097 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2098 let sessions = self.state.config.graph.sessions;
2099 let artifacts = agent::artifacts_dir(&self.state.dir());
2100 let root = self.state.worktree_root();
2101 let base_short = short(&self.state.base_commit);
2102
2103 let mut jobs = Vec::new();
2104 let mut orders = Vec::new();
2105 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2106 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2107 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2108 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2109 let seat_key = format!("judge-{}", j + 1);
2110 let seat = self.seat(&seat_key, &spec.id);
2111 jobs.push(SeatJob {
2112 prompt: prompt::judge(
2113 &self.state.instruction,
2114 &views,
2115 self.roles.judges.len(),
2116 &base_short,
2117 &language,
2118 ),
2119 spec,
2120 seat,
2121 cwd: root.join(format!("judge-{}", j + 1)),
2122 timeout,
2123 allow_write: false,
2124 sessions,
2125 artifacts: artifacts.clone(),
2126 stem: format!("judge-{}", j + 1),
2127 });
2128 }
2129
2130 self.state.event(
2131 "judge",
2132 format!(
2133 "{} judges ranking {} candidates blind",
2134 jobs.len(),
2135 viable.len()
2136 ),
2137 );
2138 let labels_for_check = labels.clone();
2139 let mut quota_losses = Vec::new();
2140 let cache = self.state.config.cache_dir();
2141 let ctx = WaveCtx {
2142 run: &run_id,
2143 node: "judge",
2144 prompts: &prompts,
2145 cache: cache.as_deref(),
2146 round: None,
2147 };
2148 let results = ask_json_wave::<Ranking>(
2149 jobs,
2150 Arc::clone(&self.sem),
2151 self.state.config.graph.retries,
2152 &ctx,
2153 &mut quota_losses,
2154 &mut self.state,
2155 &move |r: &Ranking| r.validate(&labels_for_check),
2156 )
2157 .await;
2158 self.state.quota.extend(quota_losses);
2159
2160 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2161 let agent_id = seat.agent.clone();
2162 self.state.seats.insert(seat.key.clone(), seat);
2163 let mut record = Judgement {
2164 judge: j + 1,
2165 seat: format!("judge-{}", j + 1),
2166 agent: agent_id,
2167 ranking: Vec::new(),
2168 reasons: BTreeMap::new(),
2169 confidence: None,
2170 order: orders[j].clone(),
2171 failed: None,
2172 duration_ms: 0,
2173 };
2174 match res {
2175 Ok((ranking, out)) => {
2176 record.ranking = ranking.normalized();
2177 record.reasons = ranking.reasons;
2178 record.confidence = ranking.confidence;
2179 record.duration_ms = out.duration_ms;
2180 self.state.event(
2181 "judge",
2182 format!(
2183 "judge {} ranked {}",
2184 j + 1,
2185 record.ranking.iter().collect::<String>()
2186 ),
2187 );
2188 }
2189 Err(e) => {
2190 record.failed = Some(e.to_string());
2191 self.state
2192 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2193 }
2194 }
2195 self.state.judgements.push(record);
2196 self.state.save()?;
2197 }
2198 Ok(())
2199 }
2200
2201 async fn deliberate(&mut self) -> Result<()> {
2204 let run_id = self.state.id.clone();
2209 let prompts = self.state.config.prompts.clone();
2210 if !self.state.deliberation.is_empty() {
2211 return Ok(());
2212 }
2213 let tops: Vec<char> = self
2214 .state
2215 .judgements
2216 .iter()
2217 .filter_map(|j| j.ranking.first().copied())
2218 .collect();
2219 let rounds = self.state.config.graph.deliberate_rounds;
2220 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2221 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2222 self.state.event(
2223 "deliberate",
2224 format!("judges agreed on {} outright; no deliberation", tops[0]),
2225 );
2226 }
2227 self.state.status = RunStatus::Voting;
2228 self.state.save()?;
2229 return Ok(());
2230 }
2231
2232 self.state.status = RunStatus::Deliberating;
2233 self.state.event(
2234 "deliberate",
2235 format!(
2236 "split: first choices were {} — opening {rounds} round(s)",
2237 tops.iter().collect::<String>()
2238 ),
2239 );
2240
2241 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2242 let language = self.state.config.graph.language.clone();
2243 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2244 let sessions = self.state.config.graph.sessions;
2245 let artifacts = agent::artifacts_dir(&self.state.dir());
2246 let root = self.state.worktree_root();
2247 let base_short = short(&self.state.base_commit);
2248
2249 for round in 1..=rounds {
2253 let mut turns: Vec<DeliberationTurn> = Vec::new();
2254 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2255 if self.state.judgements[j].failed.is_some() {
2256 continue;
2257 }
2258 let seat_key = format!("judge-{}", j + 1);
2259 let mut seat = self.seat(&seat_key, &spec.id);
2260 let transcript = self.transcript(&turns, j);
2261 let context = if has_context(&spec, &seat, sessions) {
2262 None
2263 } else {
2264 Some(self.candidate_block(&viable, &base_short))
2265 };
2266 let text = prompt::deliberate(
2267 &self.state.instruction,
2268 context.as_deref(),
2269 &transcript,
2270 round,
2271 rounds,
2272 &language,
2273 );
2274 let job = SeatJob {
2275 spec,
2276 seat: seat.clone(),
2277 prompt: text,
2278 cwd: root.join(format!("judge-{}", j + 1)),
2279 timeout,
2280 allow_write: false,
2281 sessions,
2282 artifacts: artifacts.clone(),
2283 stem: format!("delib-{round}-judge-{}", j + 1),
2284 };
2285 let cache = self.state.config.cache_dir();
2286 let ctx = WaveCtx {
2287 run: &run_id,
2288 node: "deliberate",
2289 prompts: &prompts,
2290 cache: cache.as_deref(),
2291 round: None,
2292 };
2293 let (updated, out) =
2294 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2295 seat = updated;
2296 let agent_id = seat.agent.clone();
2297 let seat_key = seat.key.clone();
2298 self.state.seats.insert(seat.key.clone(), seat);
2299 let body = match out {
2300 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2301 AgentOutcome::Dropped(o) => {
2305 let why =
2306 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2307 "the CLI ended the stream without delivering its answer",
2308 );
2309 self.state.event(
2310 "deliberate",
2311 format!(
2312 "judge {} skipped: the CLI dropped the stream ({why})",
2313 j + 1
2314 ),
2315 );
2316 continue;
2317 }
2318 AgentOutcome::Quota(o) => {
2319 self.state.quota.push(QuotaLoss {
2320 seat: seat_key,
2321 node: "deliberate".to_owned(),
2322 at: Timestamp::now(),
2323 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2324 });
2325 self.state.event(
2326 "deliberate",
2327 format!("judge {} skipped: rate limited (quota)", j + 1),
2328 );
2329 continue;
2330 }
2331 AgentOutcome::Failed(e) => {
2332 self.state
2333 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2334 continue;
2335 }
2336 };
2337 let tentative = verdict::extract_json::<Position>(&body)
2338 .ok()
2339 .and_then(|p| p.tentative)
2340 .and_then(|s| s.trim().chars().next())
2341 .map(|c| c.to_ascii_uppercase());
2342 self.state.event(
2343 "deliberate",
2344 format!(
2345 "round {round}: judge {} now favours {}",
2346 j + 1,
2347 tentative.map_or("—".to_owned(), |c| c.to_string())
2348 ),
2349 );
2350 turns.push(DeliberationTurn {
2351 judge: j + 1,
2352 agent: agent_id,
2353 body: blind::sanitize_prose(&body, &self.state.config.blind),
2354 tentative,
2355 });
2356 }
2357 self.state
2358 .deliberation
2359 .push(DeliberationRound { round, turns });
2360 self.state.save()?;
2361 }
2362
2363 self.state.status = RunStatus::Voting;
2364 self.state.save()?;
2365 Ok(())
2366 }
2367
2368 async fn vote(&mut self) -> Result<()> {
2371 let run_id = self.state.id.clone();
2376 let prompts = self.state.config.prompts.clone();
2377 if !self.state.votes.is_empty() {
2378 return Ok(());
2379 }
2380 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2381 if viable.len() == 1 {
2382 return Ok(());
2383 }
2384 self.state.status = RunStatus::Voting;
2385
2386 let language = self.state.config.graph.language.clone();
2387 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2388 let sessions = self.state.config.graph.sessions;
2389 let artifacts = agent::artifacts_dir(&self.state.dir());
2390 let root = self.state.worktree_root();
2391 let base_short = short(&self.state.base_commit);
2392 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2393
2394 let mut jobs = Vec::new();
2395 let mut seats_at = Vec::new();
2396 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2397 if self
2398 .state
2399 .judgements
2400 .get(j)
2401 .is_some_and(|r| r.failed.is_some())
2402 {
2403 continue;
2404 }
2405 let seat_key = format!("judge-{}", j + 1);
2406 let seat = self.seat(&seat_key, &spec.id);
2407 let mut text = prompt::final_vote(&viable, &language);
2408 if !has_context(&spec, &seat, sessions) {
2409 text = format!(
2410 "{}\n\n# Candidates\n\n{}",
2411 text,
2412 self.candidate_block(&candidates, &base_short)
2413 );
2414 }
2415 jobs.push(SeatJob {
2416 spec,
2417 seat,
2418 prompt: text,
2419 cwd: root.join(format!("judge-{}", j + 1)),
2420 timeout,
2421 allow_write: false,
2422 sessions,
2423 artifacts: artifacts.clone(),
2424 stem: format!("vote-judge-{}", j + 1),
2425 });
2426 seats_at.push(j);
2427 }
2428
2429 self.state.event(
2430 "vote",
2431 format!(
2432 "collecting {} final votes one by one, privately",
2433 jobs.len()
2434 ),
2435 );
2436 let allowed = viable.clone();
2437 let mut quota_losses = Vec::new();
2438 let cache = self.state.config.cache_dir();
2439 let ctx = WaveCtx {
2440 run: &run_id,
2441 node: "vote",
2442 prompts: &prompts,
2443 cache: cache.as_deref(),
2444 round: None,
2445 };
2446 let results = ask_json_wave::<FinalVote>(
2447 jobs,
2448 Arc::clone(&self.sem),
2449 self.state.config.graph.retries,
2450 &ctx,
2451 &mut quota_losses,
2452 &mut self.state,
2453 &move |v: &FinalVote| match v.label() {
2454 Some(c) if allowed.contains(&c) => Ok(()),
2455 other => bail!("vote {other:?} is not one of {allowed:?}"),
2456 },
2457 )
2458 .await;
2459 self.state.quota.extend(quota_losses);
2460
2461 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2462 let agent_id = seat.agent.clone();
2463 self.state.seats.insert(seat.key.clone(), seat);
2464 let initial = self
2465 .state
2466 .judgements
2467 .get(j)
2468 .and_then(|r| r.ranking.first().copied());
2469 let mut record = VoteRecord {
2470 judge: j + 1,
2471 agent: agent_id,
2472 vote: None,
2473 reason: String::new(),
2474 changed: false,
2475 };
2476 match res {
2477 Ok((v, _)) => {
2478 record.vote = v.label();
2479 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2480 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2481 self.state.event(
2482 "vote",
2483 format!(
2484 "judge {} voted {}{}",
2485 j + 1,
2486 record.vote.unwrap_or('?'),
2487 if record.changed { " (changed)" } else { "" }
2488 ),
2489 );
2490 }
2491 Err(e) => {
2492 self.state
2493 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2494 }
2495 }
2496 self.state.votes.push(record);
2497 self.state.save()?;
2498 }
2499 Ok(())
2500 }
2501
2502 fn tally(&mut self) -> Result<()> {
2505 if self.state.tally.is_some() {
2506 return Ok(());
2507 }
2508 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2509 let tops: Vec<char> = self
2510 .state
2511 .judgements
2512 .iter()
2513 .filter_map(|j| j.ranking.first().copied())
2514 .collect();
2515 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2516
2517 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2520 let mut cast: Vec<char> = Vec::new();
2521 for (i, j) in self.state.judgements.iter().enumerate() {
2522 let vote = self
2523 .state
2524 .votes
2525 .iter()
2526 .find(|v| v.judge == i + 1)
2527 .and_then(|v| v.vote)
2528 .or_else(|| j.ranking.first().copied());
2529 if let Some(v) = vote {
2530 *first_choice.entry(v).or_insert(0) += 1;
2531 cast.push(v);
2532 }
2533 }
2534
2535 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2536 for j in &self.state.judgements {
2537 let n = j.ranking.len();
2538 for (pos, label) in j.ranking.iter().enumerate() {
2539 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2540 }
2541 }
2542
2543 let best = first_choice.values().copied().max().unwrap_or(0);
2544 let mut leaders: Vec<char> = first_choice
2545 .iter()
2546 .filter(|(_, v)| **v == best)
2547 .map(|(k, _)| *k)
2548 .collect();
2549 let mut tie_break = None;
2550 if leaders.len() > 1 {
2551 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2552 let borda_leaders: Vec<char> = leaders
2553 .iter()
2554 .copied()
2555 .filter(|l| borda[l] == top_borda)
2556 .collect();
2557 tie_break = Some(if borda_leaders.len() == 1 {
2558 format!(
2559 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2560 leaders.len()
2561 )
2562 } else {
2563 format!(
2564 "{} way tie on both first-choice votes and Borda points, broken by label order",
2565 leaders.len()
2566 )
2567 });
2568 leaders = borda_leaders;
2569 leaders.sort_unstable();
2570 }
2571 let winner = *leaders
2572 .first()
2573 .or(viable.first())
2574 .context("no candidate to declare a winner from")?;
2575
2576 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2577 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2578 let deliberated = !self.state.deliberation.is_empty();
2579
2580 let quota_seats: std::collections::BTreeSet<&str> =
2584 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2585 let mut present = 0usize;
2586 for (i, j) in self.state.judgements.iter().enumerate() {
2587 if quota_seats.contains(j.seat.as_str()) {
2588 continue;
2589 }
2590 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2591 let voted = self
2592 .state
2593 .votes
2594 .iter()
2595 .any(|v| v.judge == i + 1 && v.vote.is_some());
2596 if ranked || voted {
2597 present += 1;
2598 }
2599 }
2600 let needs_quorum = viable.len() > 1;
2606 let judges_total = if needs_quorum {
2607 self.roles.judges.len()
2608 } else {
2609 0
2610 };
2611 let quorum = if needs_quorum {
2612 judges_total / 2 + 1
2613 } else {
2614 0
2615 };
2616 let met_quorum = !needs_quorum || present >= quorum;
2617 let uncontested = (!needs_quorum).then(|| {
2618 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2619 });
2620
2621 self.state.event(
2622 "tally",
2623 match &uncontested {
2624 Some(reason) => format!("winner {winner} — {reason}"),
2625 None => format!(
2626 "winner {winner} — votes {} | initial {} | {} changed | \
2627 {present}/{judges_total} judges{}",
2628 first_choice
2629 .iter()
2630 .map(|(k, v)| format!("{k}:{v}"))
2631 .collect::<Vec<_>>()
2632 .join(" "),
2633 if unanimous_initial {
2634 "unanimous"
2635 } else {
2636 "split"
2637 },
2638 changed_votes,
2639 if met_quorum {
2640 String::new()
2641 } else {
2642 format!(" — below quorum ({quorum} required)")
2643 },
2644 ),
2645 },
2646 );
2647 if !met_quorum {
2648 self.state.event(
2649 "stall",
2650 format!(
2651 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2652 the run stops here, resumable"
2653 ),
2654 );
2655 }
2656 self.state.tally = Some(Tally {
2657 first_choice,
2658 borda,
2659 winner,
2660 rankings: tops.len(),
2661 unanimous_initial,
2662 deliberated,
2663 changed_votes,
2664 unanimous_final,
2665 tie_break,
2666 judges: judges_total,
2667 present,
2668 quorum,
2669 met_quorum,
2670 uncontested,
2671 });
2672 self.state.status = if met_quorum {
2673 RunStatus::Reviewing
2674 } else {
2675 RunStatus::Stalled
2676 };
2677 self.state.save()?;
2678 Ok(())
2679 }
2680
2681 #[allow(clippy::too_many_lines)]
2702 async fn recover_stall(&mut self) -> Result<bool> {
2703 let run_id = self.state.id.clone();
2708 let prompts = self.state.config.prompts.clone();
2709 let quota_seats: BTreeSet<&str> =
2714 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2715 let absent: Vec<String> = self
2716 .state
2717 .judgements
2718 .iter()
2719 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2720 .map(|j| j.seat.clone())
2721 .collect();
2722 if absent.is_empty() {
2723 return Ok(false);
2724 }
2725 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2726 if viable.len() <= 1 {
2727 return Ok(false);
2728 }
2729 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2730 let language = self.state.config.graph.language.clone();
2731 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2732 let sessions = self.state.config.graph.sessions;
2733 let artifacts = agent::artifacts_dir(&self.state.dir());
2734 let root = self.state.worktree_root();
2735 let base_short = short(&self.state.base_commit);
2736 let candidates: Vec<Candidate> = viable.clone();
2737
2738 let mut positions: Vec<usize> = absent
2740 .iter()
2741 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2742 .collect();
2743 if positions.is_empty() {
2744 return Ok(false);
2745 }
2746 positions.sort_unstable();
2747 positions.dedup();
2748
2749 let mut judge_jobs = Vec::new();
2751 for &j in &positions {
2752 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2753 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2754 let seat_key = format!("judge-{}", j + 1);
2755 let spec = self.roles.judges[j].clone();
2756 let seat = self.seat(&seat_key, &spec.id);
2757 judge_jobs.push(SeatJob {
2758 spec,
2759 seat,
2760 prompt: prompt::judge(
2761 &self.state.instruction,
2762 &views,
2763 self.roles.judges.len(),
2764 &base_short,
2765 &language,
2766 ),
2767 cwd: root.join(seat_key),
2768 timeout,
2769 allow_write: false,
2770 sessions,
2771 artifacts: artifacts.clone(),
2772 stem: format!("judge-{}-recover", j + 1),
2773 });
2774 }
2775
2776 let labels_for_check = labels.clone();
2777 let mut judge_losses = Vec::new();
2778 let retries = self.state.config.graph.retries;
2779 let cache = self.state.config.cache_dir();
2780 let ctx = WaveCtx {
2781 run: &run_id,
2782 node: "judge",
2783 prompts: &prompts,
2784 cache: cache.as_deref(),
2785 round: None,
2786 };
2787 let results = ask_json_wave::<Ranking>(
2788 judge_jobs,
2789 Arc::clone(&self.sem),
2790 retries,
2791 &ctx,
2792 &mut judge_losses,
2793 &mut self.state,
2794 &move |r: &Ranking| r.validate(&labels_for_check),
2795 )
2796 .await;
2797
2798 let mut recovered: BTreeSet<usize> = BTreeSet::new();
2800 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
2801 self.state.seats.insert(seat.key.clone(), seat);
2802 let record = &mut self.state.judgements[j];
2803 match res {
2804 Ok((ranking, out)) => {
2805 record.ranking = ranking.normalized();
2806 record.reasons = ranking.reasons;
2807 record.confidence = ranking.confidence;
2808 record.failed = None;
2809 record.duration_ms = out.duration_ms;
2810 recovered.insert(j);
2811 self.state.event(
2812 "recover",
2813 format!("judge {} ranked again after the limit", j + 1),
2814 );
2815 }
2816 Err(e) => {
2817 self.state
2818 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
2819 }
2820 }
2821 }
2822
2823 let mut vote_jobs = Vec::new();
2825 let mut vote_pos: Vec<usize> = Vec::new();
2826 for &j in &recovered {
2827 let seat_key = format!("judge-{}", j + 1);
2828 let spec = self.roles.judges[j].clone();
2829 let seat = self.seat(&seat_key, &spec.id);
2830 let mut text = prompt::final_vote(&labels, &language);
2831 if !has_context(&spec, &seat, sessions) {
2832 text = format!(
2833 "{}\n\n# Candidates\n\n{}",
2834 text,
2835 self.candidate_block(&candidates, &base_short)
2836 );
2837 }
2838 vote_jobs.push(SeatJob {
2839 spec,
2840 seat,
2841 prompt: text,
2842 cwd: root.join(seat_key),
2843 timeout,
2844 allow_write: false,
2845 sessions,
2846 artifacts: artifacts.clone(),
2847 stem: format!("vote-judge-{}-recover", j + 1),
2848 });
2849 vote_pos.push(j);
2850 }
2851 let allowed = labels.clone();
2852 let mut vote_losses = Vec::new();
2853 let vote_retries = self.state.config.graph.retries;
2854 let vote_cache = self.state.config.cache_dir();
2855 let ctx = WaveCtx {
2856 run: &run_id,
2857 node: "vote",
2858 prompts: &prompts,
2859 cache: vote_cache.as_deref(),
2860 round: None,
2861 };
2862 let votes = ask_json_wave::<FinalVote>(
2863 vote_jobs,
2864 Arc::clone(&self.sem),
2865 vote_retries,
2866 &ctx,
2867 &mut vote_losses,
2868 &mut self.state,
2869 &move |v: &FinalVote| match v.label() {
2870 Some(c) if allowed.contains(&c) => Ok(()),
2871 other => bail!("vote {other:?} is not one of {allowed:?}"),
2872 },
2873 )
2874 .await;
2875 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
2876 let agent_id = seat.agent.clone();
2877 self.state.seats.insert(seat.key.clone(), seat);
2878 match res {
2879 Ok((v, _)) => {
2880 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
2881 rec.vote = v.label();
2882 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2883 } else {
2884 self.state.votes.push(VoteRecord {
2885 judge: j + 1,
2886 agent: agent_id,
2887 vote: v.label(),
2888 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
2889 changed: false,
2890 });
2891 }
2892 self.state.event(
2893 "recover",
2894 format!("judge {} voted again after the limit", j + 1),
2895 );
2896 }
2897 Err(e) => {
2898 self.state
2899 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
2900 }
2901 }
2902 }
2903
2904 let recovered_keys: BTreeSet<String> = recovered
2908 .iter()
2909 .map(|&j| format!("judge-{}", j + 1))
2910 .collect();
2911 self.state
2912 .quota
2913 .retain(|q| !recovered_keys.contains(&q.seat));
2914 for loss in judge_losses.into_iter().chain(vote_losses) {
2918 if recovered_keys.contains(&loss.seat) {
2919 continue;
2920 }
2921 self.state.quota.retain(|q| q.seat != loss.seat);
2922 self.state.quota.push(loss);
2923 }
2924
2925 self.state.tally = None;
2927 self.tally()?;
2928 Ok(self
2929 .state
2930 .tally
2931 .as_ref()
2932 .map(|t| t.met_quorum)
2933 .unwrap_or(false))
2934 }
2935
2936 async fn fold_losers(&mut self) -> Result<()> {
2939 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
2940 return Ok(());
2941 };
2942 let repo = self.state.repo.clone();
2943 let mut folded = Vec::new();
2944 for i in 0..self.state.candidates.len() {
2945 let c = &self.state.candidates[i];
2946 if c.label == winner || c.folded {
2947 continue;
2948 }
2949 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
2950 git::worktree_remove(&repo, &wt).await.ok();
2951 git::branch_delete(&repo, &branch).await.ok();
2952 self.state.candidates[i].folded = true;
2953 folded.push(label.to_string());
2954 }
2955 let root = self.state.worktree_root();
2957 for j in 1..=self.roles.judges.len() {
2958 let wt = root.join(format!("judge-{j}"));
2959 if wt.exists() {
2960 git::worktree_remove(&repo, &wt).await.ok();
2961 }
2962 }
2963 if self.state.config.graph.advise {
2966 for k in 1..=self.state.config.graph.advisors {
2967 let wt = root.join(format!("advisor-{k}"));
2968 if wt.exists() {
2969 git::worktree_remove(&repo, &wt).await.ok();
2970 }
2971 }
2972 }
2973 if !folded.is_empty() {
2974 self.state
2975 .event("fold", format!("folded candidates {}", folded.join(", ")));
2976 self.state.save()?;
2977 }
2978 Ok(())
2979 }
2980
2981 async fn sync_to_base(&mut self) -> Result<()> {
3011 if self
3012 .state
3013 .base_sync
3014 .as_ref()
3015 .is_some_and(|s| s.conflict.is_some())
3016 {
3017 return Ok(());
3018 }
3019 let Some(winner) = self.state.winner().cloned() else {
3020 return Ok(());
3021 };
3022
3023 let repo = self.state.repo.clone();
3024 let remote = self.state.config.merge.remote.clone();
3025 let base_branch = self.state.base_branch.clone();
3026 let tracking = format!("{remote}/{base_branch}");
3027
3028 git::fetch(&repo, &remote, &base_branch).await.ok();
3029 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3033 return Ok(());
3034 };
3035
3036 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3037 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3038 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3039
3040 if behind == 0 {
3041 self.state.base_sync = Some(BaseSync {
3042 tip,
3043 behind: 0,
3044 attempts,
3045 conflict: None,
3046 });
3047 self.state.save()?;
3048 return Ok(());
3049 }
3050
3051 if attempts >= BASE_SYNC_ROUNDS {
3052 let why = format!(
3053 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3054 rebase(s); rebasing again would only race it",
3055 winner.branch
3056 );
3057 self.state.status = RunStatus::Blocked;
3058 self.state.base_sync = Some(BaseSync {
3059 tip,
3060 behind,
3061 attempts,
3062 conflict: Some(why.clone()),
3063 });
3064 self.state.event("land", why);
3065 self.state.save()?;
3066 return Ok(());
3067 }
3068
3069 self.state.event(
3070 "land",
3071 format!(
3072 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3073 winner.branch
3074 ),
3075 );
3076 self.state.save()?;
3077
3078 let scratch = self.state.dir().join("base-sync");
3079 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
3080 let attempts = attempts + 1;
3081 match rebased {
3082 Ok(None) => {
3083 git::sync_to_head(&winner.worktree).await?;
3087 self.state.base_sync = Some(BaseSync {
3088 tip: tip.clone(),
3089 behind: 0,
3090 attempts,
3091 conflict: None,
3092 });
3093 self.state
3094 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3095 }
3096 Ok(Some(conflict)) => {
3097 let why = format!(
3098 "{} conflicts with {tracking} and did not rebase: {}",
3099 winner.branch,
3100 conflict.chars().take(600).collect::<String>()
3101 );
3102 self.state.status = RunStatus::Blocked;
3103 self.state.base_sync = Some(BaseSync {
3104 tip,
3105 behind,
3106 attempts,
3107 conflict: Some(why.clone()),
3108 });
3109 self.state.event("land", why);
3110 }
3111 Err(e) => {
3112 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3113 self.state.status = RunStatus::Blocked;
3114 self.state.base_sync = Some(BaseSync {
3115 tip,
3116 behind,
3117 attempts,
3118 conflict: Some(why.clone()),
3119 });
3120 self.state.event("land", why);
3121 }
3122 }
3123 self.state.save()?;
3124 Ok(())
3125 }
3126
3127 fn landing_base(&self) -> String {
3137 self.state
3138 .base_sync
3139 .as_ref()
3140 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3141 }
3142
3143 pub async fn fix_selected(
3176 &mut self,
3177 ids: &[String],
3178 reason: &str,
3179 allow_stale: bool,
3180 ) -> Result<()> {
3181 let reason = reason.trim();
3182 if reason.is_empty() {
3183 bail!("a fix request needs a reason — that is the operator's own record of why");
3184 }
3185 if ids.is_empty() {
3186 bail!("no finding id given");
3187 }
3188 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3189 bail!(
3190 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3191 has already concluded — can be given a targeted fix. A run still \
3192 in progress should simply be resumed; a `merged` run's branch has \
3193 already landed, so its answer is a fresh `magi review <branch>`, \
3194 not reopening this run's own record",
3195 self.state.id,
3196 self.state.status.as_str()
3197 );
3198 }
3199 let Some(winner) = self.state.winner().cloned() else {
3200 bail!("run {} has no winning candidate to fix", self.state.id);
3201 };
3202 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3203 bail!(
3204 "branch `{}` no longer exists; this run cannot be extended",
3205 winner.branch
3206 );
3207 }
3208 let home = crate::run::home();
3209 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3210 bail!(
3211 "run {} is currently being worked on by another magi process",
3212 self.state.id
3213 );
3214 }
3215 let _claim = FixClaim::acquire(&self.state.dir())?;
3221
3222 let mut seen = BTreeSet::new();
3226 let mut findings = Vec::new();
3227 let mut missing = Vec::new();
3228 for id in ids {
3229 if !seen.insert(id.clone()) {
3230 continue;
3231 }
3232 match self.state.finding(id) {
3233 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3234 id: f.id.clone(),
3235 severity: f.severity,
3236 reviewer_vote: rec.vote,
3237 round: round.round,
3238 round_head: round.head.clone(),
3239 reviewer: rec.reviewer,
3240 agent: rec.agent.clone(),
3241 file: f.file.clone(),
3242 line: f.line,
3243 title: f.title.clone(),
3244 detail: f.detail.clone(),
3245 outcome: OperatorFixOutcome::Pending,
3246 }),
3247 None => missing.push(id.clone()),
3248 }
3249 }
3250 if !missing.is_empty() {
3251 bail!(
3252 "unknown finding id(s): {}; nothing was changed",
3253 missing.join(", ")
3254 );
3255 }
3256
3257 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3258 let stale_details: Vec<(String, String)> = findings
3259 .iter()
3260 .filter(|f| f.round_head != head_at_request)
3261 .map(|f| (f.id.clone(), f.round_head.clone()))
3262 .collect();
3263 let stale = !stale_details.is_empty();
3264 if stale && !allow_stale {
3265 bail!(
3266 "the branch has moved since some finding(s) were raised — {} — now \
3267 at {}; pass --allow-stale to fix anyway, or re-run review first",
3268 stale_details
3269 .iter()
3270 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3271 .collect::<Vec<_>>()
3272 .join(", "),
3273 short(&head_at_request)
3274 );
3275 }
3276
3277 let request = OperatorFixRequest {
3278 requested_at: Timestamp::now(),
3279 reason: reason.to_owned(),
3280 findings,
3281 head_at_request: head_at_request.clone(),
3282 allow_stale,
3283 stale,
3284 fix: None,
3285 result_head: None,
3286 follow_up_review_run: None,
3287 };
3288 self.state.event(
3289 "fix",
3290 format!(
3291 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3292 request.findings.len(),
3293 request
3294 .findings
3295 .iter()
3296 .map(|f| f.id.as_str())
3297 .collect::<Vec<_>>()
3298 .join(", "),
3299 ),
3300 );
3301 self.state.operator_fixes.push(request);
3308 self.state.save()?;
3309 let request_index = self.state.operator_fixes.len() - 1;
3310
3311 if winner.worktree.exists() {
3320 let dirty = git::git(
3323 &winner.worktree,
3324 &["status", "--porcelain", "--untracked-files=all"],
3325 )
3326 .await?;
3327 let only_withheld = dirty.lines().all(|l| {
3328 l.strip_prefix("?? ")
3329 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3330 });
3331 if !only_withheld {
3332 bail!(
3333 "`{}` has uncommitted changes; refusing to touch it — commit or \
3334 discard them first",
3335 winner.worktree.display()
3336 );
3337 }
3338 git::worktree_remove(&self.state.repo, &winner.worktree)
3339 .await
3340 .ok();
3341 }
3342 let fix_worktree = self.state.worktree_root().join("operator-fix");
3343 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3344 git::git(
3345 &self.state.repo,
3346 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3347 )
3348 .await
3349 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3350 if !git::is_clean(&fix_worktree).await? {
3351 git::worktree_remove(&self.state.repo, &fix_worktree)
3352 .await
3353 .ok();
3354 bail!(
3355 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3356 winner.branch
3357 );
3358 }
3359
3360 let run_id = self.state.id.clone();
3361 let prompts = self.state.config.prompts.clone();
3362 let language = self.state.config.graph.language.clone();
3363 let sessions = self.state.config.graph.sessions;
3364 let artifacts = agent::artifacts_dir(&self.state.dir());
3365 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3366 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3367 _ => (
3368 self.state
3369 .config
3370 .agent(&winner.agent)
3371 .cloned()
3372 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3373 format!("impl-{}", winner.label),
3374 ),
3375 };
3376 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3377 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3378 .findings
3379 .iter()
3380 .map(|f| Finding {
3381 id: f.id.clone(),
3382 severity: f.severity,
3383 file: f.file.clone(),
3384 line: f.line,
3385 title: f.title.clone(),
3386 detail: f.detail.clone(),
3387 })
3388 .collect();
3389 let job = SeatJob {
3390 prompt: prompt::operator_fix(
3391 &self.state.instruction,
3392 &finding_list,
3393 reason,
3394 &stale_details,
3395 &head_at_request,
3396 &language,
3397 ),
3398 spec: fix_spec.clone(),
3399 seat,
3400 cwd: fix_worktree.clone(),
3401 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3402 allow_write: true,
3403 sessions,
3404 artifacts: artifacts.clone(),
3405 stem: "operator-fix".to_owned(),
3406 };
3407 let cache = self.state.config.cache_dir();
3408 let ctx = WaveCtx {
3409 run: &run_id,
3410 node: "fix",
3411 prompts: &prompts,
3412 cache: cache.as_deref(),
3413 round: None,
3414 };
3415 let (seat, out) =
3416 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3417 let agent_id = seat.agent.clone();
3418
3419 let mut fix = FixRecord {
3420 agent: agent_id,
3421 addressed: Vec::new(),
3422 rejected: Vec::new(),
3423 notes: String::new(),
3424 committed: false,
3425 failed: None,
3426 duration_ms: 0,
3427 continuation: None,
3428 };
3429 let mut final_seat = seat.clone();
3430 match out {
3431 AgentOutcome::Ok(o) => {
3432 fix.duration_ms = o.duration_ms;
3433 let parsed = verdict::extract_json::<FixReport>(&o.text);
3434 let incomplete_reason = match &parsed {
3435 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3436 "the reply parsed, but it reported a command whose own CLI \
3437 never confirmed an exit status"
3438 .to_owned(),
3439 ),
3440 Ok(_) => None,
3441 Err(e) => Some(e.to_string()),
3442 };
3443 match incomplete_reason {
3444 None => {
3445 let report = parsed.expect("checked Ok above");
3446 fix.addressed = report.addressed;
3447 fix.rejected = report.rejected;
3448 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3449 }
3450 Some(reason) => {
3451 let (resumed_seat, resolved, failure, cont) = self
3452 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3453 .await;
3454 fix.duration_ms += cont.cumulative_wait_ms;
3455 fix.continuation = Some(cont);
3456 final_seat = resumed_seat;
3457 match resolved {
3458 Some(report) => {
3459 fix.addressed = report.addressed;
3460 fix.rejected = report.rejected;
3461 fix.notes =
3462 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3463 }
3464 None => fix.failed = failure,
3465 }
3466 }
3467 }
3468 }
3469 AgentOutcome::Dropped(o) => {
3470 fix.duration_ms = o.duration_ms;
3471 let why = o
3472 .dropped
3473 .as_ref()
3474 .map(|d| d.why.as_str())
3475 .unwrap_or("the CLI ended the stream without delivering its answer");
3476 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3477 }
3478 AgentOutcome::Quota(o) => {
3479 self.state.quota.push(QuotaLoss {
3480 seat: final_seat.key.clone(),
3481 node: "fix".to_owned(),
3482 at: Timestamp::now(),
3483 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3484 });
3485 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3486 }
3487 AgentOutcome::Failed(e) => fix.failed = Some(e),
3488 }
3489 if fix.continuation.is_none() {
3490 fix.continuation = Some(ContinuationRecord::not_needed());
3491 }
3492 self.state.seats.insert(final_seat.key.clone(), final_seat);
3493
3494 let rescue_message = format!(
3495 "magi: operator-selected fix ({}) (uncommitted work)",
3496 self.state.operator_fixes[request_index]
3497 .findings
3498 .iter()
3499 .map(|f| f.id.as_str())
3500 .collect::<Vec<_>>()
3501 .join(", ")
3502 );
3503 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3504 self.state.note_withheld("fix", &r.withheld);
3505 }
3506 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3507 fix.committed = after != head_at_request;
3508 git::worktree_remove(&self.state.repo, &fix_worktree)
3509 .await
3510 .ok();
3511
3512 self.state.event(
3513 "fix",
3514 match &fix.failed {
3515 Some(reason) => format!(
3516 "operator fix: adoption report was lost ({reason}); {}",
3517 if fix.committed {
3518 "committed"
3519 } else {
3520 "NO new commit"
3521 }
3522 ),
3523 None => format!(
3524 "operator fix: {} addressed, {} rejected, {}",
3525 fix.addressed.len(),
3526 fix.rejected.len(),
3527 if fix.committed {
3528 "committed"
3529 } else {
3530 "NO new commit"
3531 }
3532 ),
3533 },
3534 );
3535
3536 for f in &mut self.state.operator_fixes[request_index].findings {
3543 f.outcome = if fix.failed.is_some() {
3544 OperatorFixOutcome::Unreported
3545 } else if fix.addressed.contains(&f.id) {
3546 OperatorFixOutcome::Addressed
3547 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3548 OperatorFixOutcome::Rejected { why: r.why.clone() }
3549 } else {
3550 OperatorFixOutcome::Unreported
3551 };
3552 }
3553
3554 let committed = fix.committed;
3555 if committed {
3556 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3557 }
3558 self.state.operator_fixes[request_index].fix = Some(fix);
3559 self.state.save()?;
3562
3563 if committed {
3564 self.state.event(
3565 "fix",
3566 format!(
3567 "operator fix committed {}; opening a follow-up review-only run",
3568 short(&after)
3569 ),
3570 );
3571 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3572 Ok(mut follow_up) => {
3573 follow_up.state.event(
3574 "start",
3575 format!(
3576 "requested by an operator fix on run {} for finding(s) {}",
3577 self.state.id,
3578 self.state.operator_fixes[request_index]
3579 .findings
3580 .iter()
3581 .map(|f| f.id.as_str())
3582 .collect::<Vec<_>>()
3583 .join(", "),
3584 ),
3585 );
3586 follow_up.state.save()?;
3587 let follow_up_id = follow_up.state.id.clone();
3588 if let Err(e) = follow_up.execute().await {
3589 self.state.event(
3590 "fix",
3591 format!(
3592 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3593 ),
3594 );
3595 }
3596 self.state.operator_fixes[request_index].follow_up_review_run =
3597 Some(follow_up_id);
3598 }
3599 Err(e) => {
3600 self.state.event(
3601 "fix",
3602 format!("committed the fix but could not open a follow-up review: {e:#}"),
3603 );
3604 }
3605 }
3606 self.state.save()?;
3607 }
3608
3609 Ok(())
3610 }
3611
3612 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3619 match &self.roles.fixer {
3620 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3621 _ => (
3622 self.state
3623 .config
3624 .agent(&winner.agent)
3625 .cloned()
3626 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3627 format!("impl-{}", winner.label),
3628 ),
3629 }
3630 }
3631
3632 async fn review_loop(&mut self) -> Result<()> {
3633 if self
3638 .state
3639 .base_sync
3640 .as_ref()
3641 .is_some_and(|s| s.conflict.is_some())
3642 {
3643 return Ok(());
3644 }
3645 let run_id = self.state.id.clone();
3650 let prompts = self.state.config.prompts.clone();
3651 let Some(winner) = self.state.winner().cloned() else {
3652 return Ok(());
3653 };
3654 let max_rounds = self.state.config.graph.review_rounds;
3655 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
3665 self.state.status = status;
3666 self.state.save()?;
3667 return Ok(());
3668 }
3669 self.state.status = RunStatus::Reviewing;
3670 if self
3680 .state
3681 .reviews
3682 .last()
3683 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
3684 {
3685 let shell = self.state.config.shell();
3686 return self
3687 .stop_reviewing(
3688 "the last round's own verification never resolved",
3689 &shell,
3690 &winner.worktree,
3691 )
3692 .await;
3693 }
3694
3695 let repo = self.state.repo.clone();
3696 let root = self.state.worktree_root();
3697 let language = self.state.config.graph.language.clone();
3698 let sessions = self.state.config.graph.sessions;
3699 let artifacts = agent::artifacts_dir(&self.state.dir());
3700 let base = self.landing_base();
3701 let base_short = short(&base);
3702 let reviewers = self.roles.reviewers.clone();
3703 let shell = self.state.config.shell();
3704
3705 for round in (self.state.reviews.len() + 1)..=max_rounds {
3706 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3707 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
3708 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
3709 let prev_verification = self
3718 .state
3719 .reviews
3720 .last()
3721 .and_then(|r| r.verification_summary(&head));
3722
3723 let mut jobs = Vec::new();
3727 for (r, spec) in reviewers.iter().cloned().enumerate() {
3728 let wt = root.join(format!("review-{}", r + 1));
3729 if wt.exists() {
3730 git::reset_detached(&wt, &head).await?;
3731 } else {
3732 git::worktree_add_detached(&repo, &wt, &head).await?;
3733 }
3734 let seat_key = format!("review-{}", r + 1);
3735 let seat = self.seat(&seat_key, &spec.id);
3736 jobs.push(SeatJob {
3737 prompt: prompt::review(&prompt::ReviewCtx {
3738 instruction: &self.state.instruction,
3739 branch: &winner.branch,
3740 base_short: &base_short,
3741 stat: &stat,
3742 patch: &patch,
3743 verification: prev_verification.as_ref(),
3744 reviewers: reviewers.len(),
3745 round,
3746 rounds: max_rounds,
3747 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
3750 lens: Lens::for_seat(r),
3751 language: &language,
3752 }),
3753 spec,
3754 seat,
3755 cwd: wt,
3756 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
3757 allow_write: false,
3758 sessions,
3759 artifacts: artifacts.clone(),
3760 stem: format!("review-{round}-{}", r + 1),
3761 });
3762 }
3763
3764 self.state.event(
3765 "review",
3766 format!(
3767 "round {round}: {} reviewers on {}",
3768 jobs.len(),
3769 short(&head)
3770 ),
3771 );
3772 let mut quota_losses = Vec::new();
3773 let review_retries = self.state.config.graph.retries;
3774 let review_cache = self.state.config.cache_dir();
3775 let ctx = WaveCtx {
3776 run: &run_id,
3777 node: "review",
3778 prompts: &prompts,
3779 cache: review_cache.as_deref(),
3780 round: Some(round),
3781 };
3782 let results = ask_json_wave::<Review>(
3783 jobs,
3784 Arc::clone(&self.sem),
3785 review_retries,
3786 &ctx,
3787 &mut quota_losses,
3788 &mut self.state,
3789 &|_: &Review| Ok(()),
3790 )
3791 .await;
3792 let round_quota_missing = quota_losses.len();
3796 self.state.quota.extend(quota_losses);
3797
3798 let mut records = Vec::new();
3799 let mut all_findings = Vec::new();
3800 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
3801 let agent_id = seat.agent.clone();
3802 self.state.seats.insert(seat.key.clone(), seat);
3803 let mut record = ReviewRecord {
3804 reviewer: r + 1,
3805 agent: agent_id,
3806 summary: String::new(),
3807 findings: Vec::new(),
3808 vote: None,
3809 failed: None,
3810 duration_ms: 0,
3811 attempts,
3817 };
3818 match res {
3819 Ok((review, out)) => {
3820 record.summary =
3828 blind::sanitize_prose(&review.summary, &self.state.config.blind);
3829 record.vote = Some(review.vote);
3830 record.duration_ms = out.duration_ms;
3831 for (n, mut f) in review.findings.into_iter().enumerate() {
3832 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
3835 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
3836 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
3837 f.file = f
3843 .file
3844 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
3845 all_findings.push(f.clone());
3846 record.findings.push(f);
3847 }
3848 self.state.event(
3849 "review",
3850 format!(
3851 "round {round}: reviewer {} voted {} with {} finding(s)",
3852 r + 1,
3853 review.vote.label(),
3854 record.findings.len()
3855 ),
3856 );
3857 }
3858 Err(e) => {
3859 record.failed = Some(e.to_string());
3860 self.state.event(
3861 "review",
3862 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
3863 );
3864 }
3865 }
3866 records.push(record);
3867 }
3868
3869 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
3876 let vote_split =
3877 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
3878 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
3879 if vote_split {
3880 self.state.event(
3881 "review",
3882 format!(
3883 "round {round}: votes split ({}) — one round of reconsideration",
3884 initial_votes
3885 .iter()
3886 .map(|v| v.label())
3887 .collect::<Vec<_>>()
3888 .join(", ")
3889 ),
3890 );
3891 let panel: Vec<ReviewSeatReport<'_>> = records
3894 .iter()
3895 .filter_map(|r| {
3896 r.vote.map(|vote| ReviewSeatReport {
3897 reviewer: r.reviewer,
3898 vote,
3899 summary: &r.summary,
3900 findings: &r.findings,
3901 })
3902 })
3903 .collect();
3904
3905 let mut jobs = Vec::new();
3906 let mut seats_at = Vec::new();
3907 for (r, spec) in reviewers.iter().cloned().enumerate() {
3908 if records[r].vote.is_none() {
3912 continue;
3913 }
3914 let wt = root.join(format!("review-{}", r + 1));
3915 let seat_key = format!("review-{}", r + 1);
3916 let seat = self.seat(&seat_key, &spec.id);
3917 let patch_ctx = if has_context(&spec, &seat, sessions) {
3922 None
3923 } else {
3924 Some(ReviewPatch {
3925 branch: &winner.branch,
3926 base_short: &base_short,
3927 stat: &stat,
3928 patch: &patch,
3929 })
3930 };
3931 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
3932 instruction: &self.state.instruction,
3933 reviewer: r + 1,
3934 lens: Lens::for_seat(r),
3935 panel: &panel,
3936 patch: patch_ctx,
3937 round,
3938 rounds: max_rounds,
3939 language: &language,
3940 });
3941 jobs.push(SeatJob {
3942 prompt,
3943 spec,
3944 seat,
3945 cwd: wt,
3946 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
3947 allow_write: false,
3948 sessions,
3949 artifacts: artifacts.clone(),
3950 stem: format!("review-{round}-reconsider-{}", r + 1),
3951 });
3952 seats_at.push(r);
3953 }
3954
3955 let mut recon_quota_losses = Vec::new();
3956 let recon_cache = self.state.config.cache_dir();
3957 let recon_ctx = WaveCtx {
3958 run: &run_id,
3959 node: "review",
3960 prompts: &prompts,
3961 cache: recon_cache.as_deref(),
3962 round: Some(round),
3963 };
3964 let recon_results = ask_json_wave::<ReviewRevote>(
3965 jobs,
3966 Arc::clone(&self.sem),
3967 review_retries,
3968 &recon_ctx,
3969 &mut recon_quota_losses,
3970 &mut self.state,
3971 &|_: &ReviewRevote| Ok(()),
3972 )
3973 .await;
3974 self.state.quota.extend(recon_quota_losses);
3975
3976 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
3977 let agent_id = seat.agent.clone();
3978 self.state.seats.insert(seat.key.clone(), seat);
3979 let mut rec = ReviewRevoteRecord {
3980 reviewer: r + 1,
3981 agent: agent_id,
3982 vote: None,
3983 reason: String::new(),
3984 failed: None,
3985 };
3986 match res {
3987 Ok((rv, _)) => {
3988 rec.vote = Some(rv.vote);
3989 rec.reason =
3990 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
3991 self.state.event(
3992 "review",
3993 format!(
3994 "round {round}: reviewer {} revoted {}",
3995 r + 1,
3996 rv.vote.label()
3997 ),
3998 );
3999 }
4000 Err(e) => {
4001 rec.failed = Some(e.to_string());
4002 self.state.event(
4003 "review",
4004 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4005 );
4006 }
4007 }
4008 reconsideration.push(rec);
4009 }
4010 } else if initial_votes.len() > 1 {
4011 self.state.event(
4012 "review",
4013 format!(
4014 "round {round}: votes agreed ({}) — no reconsideration",
4015 initial_votes[0].label()
4016 ),
4017 );
4018 }
4019
4020 let final_votes: Vec<ReviewVote> = records
4024 .iter()
4025 .filter_map(|r| {
4026 reconsideration
4027 .iter()
4028 .find(|rv| rv.reviewer == r.reviewer)
4029 .and_then(|rv| rv.vote)
4030 .or(r.vote)
4031 })
4032 .collect();
4033 let round_verdict = ReviewVote::worst(final_votes);
4034
4035 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4036 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4037 let defer_e2e =
4048 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4049 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4050 let reason =
4051 format!("{blocking} blocking finding(s) already required a fix this round");
4052 self.state.event(
4053 "verify",
4054 format!(
4055 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4056 {}); it will run once a round has none left",
4057 short(&head)
4058 ),
4059 );
4060 (Vec::new(), false, true, Some(reason))
4061 } else {
4062 let e2e_commands = self.state.config.verify.e2e.clone();
4063 let cache_dir = self.state.config.cache_dir();
4064 let context = format!("round {round}");
4065 let (e2e, verify_retried) = with_cache_lease(
4066 &mut self.state,
4067 cache_dir.as_deref(),
4068 "e2e",
4069 "e2e",
4070 &winner.worktree,
4071 &head,
4072 verify_timeout,
4073 &context,
4074 |state, budget| {
4075 let shell = shell.clone();
4076 let e2e_commands = e2e_commands.clone();
4077 let worktree = winner.worktree.clone();
4078 let context = context.clone();
4079 async move {
4080 run_e2e_with_retry(
4081 state,
4082 &shell,
4083 &e2e_commands,
4084 &worktree,
4085 budget,
4086 &context,
4087 )
4088 .await
4089 }
4090 },
4091 )
4092 .await;
4093 (e2e, verify_retried, false, None)
4094 };
4095
4096 let expected = records.len();
4097 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4098 let incomplete = answered < expected;
4099 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4100 let policy = self.state.config.graph.incomplete_review;
4101 let clean = round_is_clean(
4102 blocking,
4103 e2e_ok,
4104 answered,
4105 expected,
4106 round_quota_missing,
4107 policy,
4108 );
4109
4110 let mut round_record = ReviewRound {
4111 round,
4112 head: head.clone(),
4113 verified_head: None,
4114 verified_at: None,
4115 reviews: records,
4116 e2e,
4117 verify_retried,
4118 e2e_deferred,
4119 e2e_defer_reason,
4120 fix: None,
4121 blocking,
4122 answered,
4123 expected,
4124 clean,
4125 progressed: false,
4126 vote_split,
4127 reconsideration,
4128 verdict: round_verdict,
4129 };
4130 if !matches!(
4141 round_record.e2e_status(),
4142 E2eStatus::Deferred | E2eStatus::NotConfigured
4143 ) {
4144 round_record.verified_head = Some(head.clone());
4145 round_record.verified_at = Some(Timestamp::now());
4146 }
4147 let this_round_verification = round_record.verification_summary(&head);
4148
4149 if incomplete {
4150 let missing: Vec<String> = round_record
4151 .reviews
4152 .iter()
4153 .filter(|r| r.failed.is_some())
4154 .map(|r| format!("review-{}", r.reviewer))
4155 .collect();
4156 self.state.event(
4157 "review",
4158 format!(
4159 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4160 missing.join(", ")
4161 ),
4162 );
4163 }
4164
4165 if clean {
4166 self.state.event(
4167 "review",
4168 if incomplete && policy == IncompleteReviewPolicy::Warn {
4169 format!(
4170 "round {round}: clean (warn policy, incomplete panel) — no \
4171 blocking findings from the seats that answered, verification green"
4172 )
4173 } else if incomplete {
4174 format!(
4175 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4176 quorum) — no blocking findings from the seats that answered, \
4177 verification green",
4178 expected - answered
4179 )
4180 } else {
4181 format!("round {round}: clean — no blocking findings, verification green")
4182 },
4183 );
4184 self.state.reviews.push(round_record);
4185 self.state.status = RunStatus::Gating;
4186 self.state.save()?;
4187 return Ok(());
4188 }
4189
4190 if incomplete && blocking == 0 && e2e_ok {
4198 self.state.reviews.push(round_record);
4199 self.state.save()?;
4200 if round == max_rounds {
4201 self.state.status = RunStatus::Blocked;
4202 self.state.event(
4203 "review",
4204 format!(
4205 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4206 refusing to call it clean",
4207 expected - answered
4208 ),
4209 );
4210 return Ok(());
4211 }
4212 continue;
4213 }
4214
4215 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4227 self.state.reviews.push(round_record);
4228 return self
4229 .stop_reviewing(
4230 "the round's own verification could not run",
4231 &shell,
4232 &winner.worktree,
4233 )
4234 .await;
4235 }
4236
4237 if round == max_rounds {
4238 self.state.reviews.push(round_record);
4239 return self
4240 .stop_reviewing(
4241 &format!(
4242 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4243 ),
4244 &shell,
4245 &winner.worktree,
4246 )
4247 .await;
4248 }
4249
4250 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4253 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4254 let blocking_findings: Vec<_> = all_findings
4255 .iter()
4256 .filter(|f| f.severity.blocks())
4257 .cloned()
4258 .collect();
4259 let job = SeatJob {
4260 prompt: prompt::fix(
4261 &self.state.instruction,
4262 &blocking_findings,
4263 this_round_verification.as_ref(),
4264 round,
4265 max_rounds,
4266 &language,
4267 ),
4268 spec: fix_spec.clone(),
4269 seat,
4270 cwd: winner.worktree.clone(),
4271 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4272 allow_write: true,
4273 sessions,
4274 artifacts: artifacts.clone(),
4275 stem: format!("fix-{round}"),
4276 };
4277 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4278 let cache = self.state.config.cache_dir();
4279 let ctx = WaveCtx {
4280 run: &run_id,
4281 node: "fix",
4282 prompts: &prompts,
4283 cache: cache.as_deref(),
4284 round: Some(round),
4285 };
4286 let (seat, out) =
4287 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4288 let agent_id = seat.agent.clone();
4289
4290 let mut fix = FixRecord {
4291 agent: agent_id,
4292 addressed: Vec::new(),
4293 rejected: Vec::new(),
4294 notes: String::new(),
4295 committed: false,
4296 failed: None,
4297 duration_ms: 0,
4298 continuation: None,
4299 };
4300 let mut continuation = ContinuationRecord::not_needed();
4301 let mut final_seat = seat.clone();
4302 match out {
4303 AgentOutcome::Ok(o) => {
4304 fix.duration_ms = o.duration_ms;
4305 let parsed = verdict::extract_json::<FixReport>(&o.text);
4306 let incomplete_reason = match &parsed {
4313 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4314 "the reply parsed, but it reported a command whose own CLI never \
4315 confirmed an exit status"
4316 .to_owned(),
4317 ),
4318 Ok(_) => None,
4319 Err(e) => Some(e.to_string()),
4320 };
4321 match incomplete_reason {
4322 None => {
4323 let report = parsed.expect("checked Ok above");
4324 fix.addressed = report.addressed;
4325 fix.rejected = report.rejected;
4326 fix.notes =
4327 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4328 }
4329 Some(reason) => {
4330 let (resumed_seat, resolved, failure, cont) = self
4331 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4332 .await;
4333 fix.duration_ms += cont.cumulative_wait_ms;
4334 continuation = cont;
4335 final_seat = resumed_seat;
4336 match resolved {
4337 Some(report) => {
4338 fix.addressed = report.addressed;
4339 fix.rejected = report.rejected;
4340 fix.notes = blind::sanitize_prose(
4341 &report.notes,
4342 &self.state.config.blind,
4343 );
4344 }
4345 None => fix.failed = failure,
4346 }
4347 }
4348 }
4349 }
4350 AgentOutcome::Dropped(o) => {
4352 fix.duration_ms = o.duration_ms;
4353 let why = o
4354 .dropped
4355 .as_ref()
4356 .map(|d| d.why.as_str())
4357 .unwrap_or("the CLI ended the stream without delivering its answer");
4358 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4359 }
4360 AgentOutcome::Quota(o) => {
4361 self.state.quota.push(QuotaLoss {
4362 seat: final_seat.key.clone(),
4363 node: "fix".to_owned(),
4364 at: Timestamp::now(),
4365 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4366 });
4367 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4368 }
4369 AgentOutcome::Failed(e) => fix.failed = Some(e),
4370 }
4371 fix.continuation = Some(continuation);
4372 self.state.seats.insert(final_seat.key.clone(), final_seat);
4373 if let Ok(r) = git::rescue_commit(
4374 &winner.worktree,
4375 &format!("magi: review round {round} fixes (uncommitted work)"),
4376 )
4377 .await
4378 {
4379 self.state.note_withheld("fix", &r.withheld);
4380 }
4381 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4382 fix.committed = after != before;
4383 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4391 let progressed = diff_after != patch;
4392 let commit_note = if fix.committed {
4393 "committed"
4394 } else {
4395 "NO new commit"
4396 };
4397 let tree_note = if progressed {
4398 "changed vs base"
4399 } else {
4400 "unchanged vs base"
4401 };
4402 self.state.event(
4403 "fix",
4404 match &fix.failed {
4405 Some(reason) => {
4411 format!(
4412 "round {round}: fixer's adoption report was lost ({reason}); \
4413 {commit_note}, tree {tree_note}"
4414 )
4415 }
4416 None => format!(
4417 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4418 {tree_note}{}",
4419 fix.addressed.len(),
4420 fix.rejected.len(),
4421 if continuation.outcome == ContinuationOutcome::Resumed {
4422 format!(
4423 " (adoption report recovered after {} continuation(s))",
4424 continuation.attempts
4425 )
4426 } else {
4427 String::new()
4428 },
4429 ),
4430 },
4431 );
4432 round_record.fix = Some(fix);
4433 round_record.progressed = progressed;
4434 self.state.reviews.push(round_record);
4435 self.state.save()?;
4436
4437 if matches!(
4450 continuation.outcome,
4451 ContinuationOutcome::Exhausted
4452 | ContinuationOutcome::QuotaLost
4453 | ContinuationOutcome::NoSession
4454 ) {
4455 return self
4456 .stop_reviewing(
4457 "the fixer's adoption report never came back, even after resuming its \
4458 own seat; refusing to start another round against the same worktree \
4459 while that is unresolved",
4460 &shell,
4461 &winner.worktree,
4462 )
4463 .await;
4464 }
4465
4466 let streak = self
4467 .state
4468 .reviews
4469 .iter()
4470 .rev()
4471 .take_while(|r| !r.progressed)
4472 .count();
4473 if streak >= STAGNANT_LIMIT {
4474 return self
4475 .stop_reviewing(
4476 &format!(
4477 "the tree has not moved against base for {streak} round(s) in a row"
4478 ),
4479 &shell,
4480 &winner.worktree,
4481 )
4482 .await;
4483 }
4484 }
4485 Ok(())
4486 }
4487
4488 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4518 let round_idx = self.state.reviews.len() - 1;
4519 let needs_catchup_run = matches!(
4527 self.state.reviews[round_idx].e2e_status(),
4528 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4529 );
4530 if needs_catchup_run {
4531 let round = self.state.reviews[round_idx].round;
4532 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4533 let commands = self.state.config.verify.e2e.clone();
4534 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4535 let cache_dir = self.state.config.cache_dir();
4536 let context = format!(
4537 "round {round}: verification unresolved, catching up before the final decision"
4538 );
4539 let (outcomes, verify_retried) = with_cache_lease(
4540 &mut self.state,
4541 cache_dir.as_deref(),
4542 "e2e",
4543 "e2e",
4544 worktree,
4545 &attempted_head,
4546 timeout,
4547 &context,
4548 |state, budget| {
4549 let shell = shell.to_vec();
4550 let commands = commands.clone();
4551 let context = context.clone();
4552 async move {
4553 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4554 .await
4555 }
4556 },
4557 )
4558 .await;
4559 let last = &mut self.state.reviews[round_idx];
4560 last.e2e = outcomes;
4561 last.verify_retried = verify_retried;
4562 last.verified_head = Some(attempted_head);
4569 last.verified_at = Some(Timestamp::now());
4570 if verify_inconclusive(&last.e2e) {
4571 self.state.save()?;
4578 return Ok(());
4579 }
4580 last.e2e_deferred = false;
4581 }
4582 let last = &self.state.reviews[round_idx];
4583 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4584
4585 match last.e2e_status() {
4586 E2eStatus::Failed => {
4587 let red: Vec<String> = last
4588 .e2e
4589 .iter()
4590 .filter(|o| !o.ok())
4591 .map(|o| {
4592 format!(
4593 "`{}` -> {:?}\n{}",
4594 o.command,
4595 o.code,
4596 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4597 )
4598 })
4599 .collect();
4600 self.state
4601 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4602 self.state.status = RunStatus::Blocked;
4603 }
4604 E2eStatus::ResourceBlocked => {
4609 self.state.event(
4610 "review",
4611 format!(
4612 "{why}; e2e could not run (shared build cache unavailable); not \
4613 deciding yet"
4614 ),
4615 );
4616 }
4617 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4618 self.state.event(
4619 "review",
4620 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4621 );
4622 self.state.status = RunStatus::Gating;
4623 }
4624 }
4625 self.state.save()?;
4626 Ok(())
4627 }
4628
4629 async fn gate(&mut self) -> Result<()> {
4632 if self.state.status == RunStatus::Failed
4644 || self
4645 .state
4646 .base_sync
4647 .as_ref()
4648 .is_some_and(|s| s.conflict.is_some())
4649 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
4650 != Some(RunStatus::Gating)
4651 {
4652 return Ok(());
4653 }
4654 if self.state.gate_ran {
4655 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
4666 self.state.status = RunStatus::Blocked;
4667 self.state.save()?;
4668 }
4669 return Ok(());
4670 }
4671 let Some(winner) = self.state.winner().cloned() else {
4672 return Ok(());
4673 };
4674 self.state.status = RunStatus::Gating;
4675 let mut outcomes = self.run_gate(&winner).await?;
4676 loop {
4677 if verify_inconclusive(&outcomes) {
4688 self.state.save()?;
4689 return Ok(());
4690 }
4691 if outcomes.iter().all(CommandOutcome::ok) {
4692 break;
4693 }
4694 match self.gate_fix_round(&winner, &outcomes).await? {
4695 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
4696 GateFix::Stop => break,
4697 GateFix::Defer => {
4698 self.state.save()?;
4699 return Ok(());
4700 }
4701 }
4702 }
4703 let passed = outcomes.iter().all(CommandOutcome::ok);
4704 self.state.gate = outcomes;
4705 self.state.gate_ran = true;
4706 if !passed {
4707 self.state.status = RunStatus::Blocked;
4708 let spent = self.state.gate_fixes.len();
4709 self.state.event(
4710 "gate",
4711 if spent == 0 {
4712 "gate failed; not merging".to_owned()
4713 } else {
4714 format!("gate failed after {spent} gate-fix round(s); not merging")
4715 },
4716 );
4717 }
4718 self.state.save()?;
4719 Ok(())
4720 }
4721
4722 async fn run_pre_gate(&mut self, winner: &Candidate) {
4732 let commands = self.state.config.verify.pre_gate.clone();
4733 if commands.is_empty() {
4734 return;
4735 }
4736 let shell = self.state.config.shell();
4737 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4738 let (outcomes, _) = run_commands(
4739 &mut self.state,
4740 "pre_gate",
4741 "pre_gate",
4742 0,
4743 &shell,
4744 &commands,
4745 &winner.worktree,
4746 timeout,
4747 )
4748 .await;
4749 for o in &outcomes {
4750 if !o.ok() {
4751 tracing::warn!(
4752 "pre_gate `{}` failed ({:?}); the gate decides",
4753 o.command,
4754 o.code
4755 );
4756 }
4757 self.state.event(
4758 "pre_gate",
4759 format!(
4760 "`{}` -> {}",
4761 o.command,
4762 if o.ok() {
4763 "pass".to_owned()
4764 } else {
4765 format!(
4766 "FAIL ({:?})\n{}",
4767 o.code,
4768 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4769 )
4770 }
4771 ),
4772 );
4773 }
4774 self.state.pre_gate = outcomes;
4775 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
4776 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
4777 Ok(head) => {
4778 self.state
4779 .event("pre_gate", format!("committed mechanical fixes ({head})"));
4780 self.state.pre_gate_commit = Some(head);
4781 }
4782 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
4783 },
4784 Ok(false) => {}
4785 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
4786 }
4787 if let Err(e) = self.state.save() {
4788 tracing::warn!("could not persist the pre_gate record: {e:#}");
4789 }
4790 }
4791
4792 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
4795 self.run_pre_gate(winner).await;
4796 let shell = self.state.config.shell();
4797 let gate_commands = self.state.config.verify.gate.clone();
4798 let outcomes = if gate_commands.is_empty() {
4807 Vec::new()
4808 } else {
4809 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4810 let cache_dir = self.state.config.cache_dir();
4811 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4812 let (outcomes, _) = with_cache_lease(
4813 &mut self.state,
4814 cache_dir.as_deref(),
4815 "gate",
4816 "gate",
4817 &winner.worktree,
4818 &head,
4819 timeout,
4820 "final gate",
4821 |state, budget| {
4822 let shell = shell.clone();
4823 let gate_commands = gate_commands.clone();
4824 let worktree = winner.worktree.clone();
4825 async move {
4826 let (outcomes, timed_out_pids) = run_commands(
4827 state,
4828 "gate",
4829 "gate",
4830 0,
4831 &shell,
4832 &gate_commands,
4833 &worktree,
4834 budget,
4835 )
4836 .await;
4837 (outcomes, false, timed_out_pids)
4838 }
4839 },
4840 )
4841 .await;
4842 outcomes
4843 };
4844 if outcomes.is_empty() {
4845 self.state.event(
4850 "gate",
4851 "no gate commands configured; nothing to check, passing",
4852 );
4853 }
4854 for o in &outcomes {
4855 self.state.event(
4856 "gate",
4857 format!(
4858 "`{}` -> {}",
4859 o.command,
4860 if o.ok() {
4861 "pass".to_owned()
4862 } else {
4863 format!(
4864 "FAIL ({:?})\n{}",
4865 o.code,
4866 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4867 )
4868 }
4869 ),
4870 );
4871 }
4872 Ok(outcomes)
4873 }
4874
4875 async fn gate_fix_round(
4888 &mut self,
4889 winner: &Candidate,
4890 outcomes: &[CommandOutcome],
4891 ) -> Result<GateFix> {
4892 let cap = self.state.config.graph.gate_fix_rounds;
4893 let spent = self.state.gate_fixes.len();
4894 if spent >= cap {
4895 if cap > 0 {
4896 self.state.event(
4897 "gate",
4898 format!("{spent} gate-fix round(s) spent and the gate still fails"),
4899 );
4900 }
4901 return Ok(GateFix::Stop);
4902 }
4903 if !gate_fixable(outcomes) {
4904 self.state.event(
4905 "gate",
4906 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
4907 command or similar); not spending a fix round on it",
4908 );
4909 return Ok(GateFix::Stop);
4910 }
4911 let min_free = self.state.config.disk.min_free_bytes;
4912 if min_free > 0 {
4913 match crate::disk::free_bytes(&winner.worktree) {
4914 Ok(free) if crate::disk::enough_space(free, min_free) => {}
4915 Ok(free) => {
4916 self.state.event(
4917 "gate",
4918 format!(
4919 "only {free} bytes free ({min_free} required by `[disk] \
4920 min_free_bytes`); not spending a fix round on a failure the disk \
4921 may explain"
4922 ),
4923 );
4924 return Ok(GateFix::Stop);
4925 }
4926 Err(e) => {
4927 self.state.event(
4928 "gate",
4929 format!("free disk space could not be measured ({e:#}); no fix round"),
4930 );
4931 return Ok(GateFix::Stop);
4932 }
4933 }
4934 }
4935
4936 let attempt = spent + 1;
4937 let run_id = self.state.id.clone();
4938 let prompts = self.state.config.prompts.clone();
4939 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
4940 let base = self.landing_base();
4941 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
4942 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4943 let job = SeatJob {
4944 prompt: prompt::gate_fix(
4945 &self.state.instruction,
4946 &failed,
4947 attempt,
4948 cap,
4949 &self.state.config.graph.language,
4950 ),
4951 spec: fix_spec,
4952 seat,
4953 cwd: winner.worktree.clone(),
4954 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4955 allow_write: true,
4956 sessions: self.state.config.graph.sessions,
4957 artifacts: agent::artifacts_dir(&self.state.dir()),
4958 stem: format!("gate-fix-{attempt}"),
4959 };
4960 self.state.event(
4961 "gate",
4962 format!("gate failed; gate-fix round {attempt} of {cap}"),
4963 );
4964 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4965 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
4966 let cache = self.state.config.cache_dir();
4967 let ctx = WaveCtx {
4968 run: &run_id,
4969 node: "gate-fix",
4970 prompts: &prompts,
4971 cache: cache.as_deref(),
4972 round: None,
4973 };
4974 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4975 let mut record = GateFixRecord {
4976 agent: seat.agent.clone(),
4977 failed,
4978 notes: String::new(),
4979 committed: false,
4980 error: None,
4981 };
4982 match out {
4983 AgentOutcome::Ok(o) => {
4984 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
4987 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
4988 }
4989 }
4990 AgentOutcome::Dropped(_) => {
4991 record.error = Some("the CLI dropped the stream".to_owned());
4992 }
4993 AgentOutcome::Quota(o) => {
4994 self.state.quota.push(QuotaLoss {
4995 seat: seat.key.clone(),
4996 node: "gate-fix".to_owned(),
4997 at: Timestamp::now(),
4998 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4999 });
5000 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5001 }
5002 AgentOutcome::Failed(e) => record.error = Some(e),
5003 }
5004 self.state.seats.insert(seat.key.clone(), seat);
5005 if let Ok(r) = git::rescue_commit(
5006 &winner.worktree,
5007 &format!("magi: gate fix {attempt} (uncommitted work)"),
5008 )
5009 .await
5010 {
5011 self.state.note_withheld("gate-fix", &r.withheld);
5012 }
5013 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5014 record.committed = after != before;
5015 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5016 let note = record.error.clone();
5017 self.state.gate_fixes.push(record);
5018 self.state.save()?;
5019 if !changed {
5020 self.state.event(
5021 "gate",
5022 match note {
5023 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5024 None => format!("gate-fix round {attempt}: the tree did not change"),
5025 },
5026 );
5027 return Ok(GateFix::Stop);
5028 }
5029 self.state.event(
5030 "gate",
5031 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5032 );
5033
5034 let commands = self.state.config.verify.e2e.clone();
5035 if !commands.is_empty() {
5036 let shell = self.state.config.shell();
5037 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5038 let cache_dir = self.state.config.cache_dir();
5039 let context = format!("gate-fix round {attempt}");
5040 let (e2e, _) = with_cache_lease(
5041 &mut self.state,
5042 cache_dir.as_deref(),
5043 "e2e",
5044 "e2e",
5045 &winner.worktree,
5046 &after,
5047 timeout,
5048 &context,
5049 |state, budget| {
5050 let shell = shell.clone();
5051 let commands = commands.clone();
5052 let context = context.clone();
5053 let worktree = winner.worktree.clone();
5054 async move {
5055 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5056 .await
5057 }
5058 },
5059 )
5060 .await;
5061 if verify_inconclusive(&e2e) {
5062 return Ok(GateFix::Defer);
5063 }
5064 if e2e.iter().any(|o| !o.ok()) {
5065 self.state.event(
5066 "gate",
5067 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5068 );
5069 return Ok(GateFix::Stop);
5070 }
5071 }
5072 Ok(GateFix::Retry)
5073 }
5074
5075 async fn merge(&mut self) -> Result<()> {
5078 if self
5093 .state
5094 .base_sync
5095 .as_ref()
5096 .is_some_and(|s| s.conflict.is_some())
5097 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5098 != Some(RunStatus::Gating)
5099 || !self.state.gate_status().ok()
5108 {
5109 return Ok(());
5110 }
5111 if self.state.merge.is_some() {
5120 return Ok(());
5121 }
5122 let Some(winner) = self.state.winner().cloned() else {
5123 return Ok(());
5124 };
5125 let repo = self.state.repo.clone();
5126 let base = self.state.base_branch.clone();
5127 let mode = self.state.config.merge.mode;
5128 let style = self.state.config.merge.style;
5129 let pr = pr_message(&self.state, winner.label);
5130 let message = pr.commit_message();
5131
5132 let outcome = match mode {
5133 MergeMode::None => MergeOutcome {
5134 mode,
5135 ok: true,
5136 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5137 },
5138 MergeMode::Local => {
5139 let on = git::current_branch(&repo).await?;
5140 if on.as_deref() != Some(base.as_str()) {
5141 MergeOutcome {
5142 mode,
5143 ok: false,
5144 detail: format!(
5145 "{} has {} checked out, not the base branch {base}",
5146 repo.display(),
5147 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5148 ),
5149 }
5150 } else if !git::is_clean(&repo).await? {
5151 MergeOutcome {
5152 mode,
5153 ok: false,
5154 detail: format!("{} is dirty; refusing to merge", repo.display()),
5155 }
5156 } else {
5157 let out = match style {
5158 MergeStyle::Merge => {
5159 git::merge_no_ff(&repo, &winner.branch, &message).await?
5160 }
5161 MergeStyle::Squash => {
5162 git::merge_squash(&repo, &winner.branch, &message).await?
5163 }
5164 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5165 };
5166 MergeOutcome {
5167 mode,
5168 ok: out.ok(),
5169 detail: if out.ok() { out.stdout } else { out.stderr },
5170 }
5171 }
5172 }
5173 MergeMode::Pr => {
5174 let remote = self.state.config.merge.remote.clone();
5175 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5176 if !pushed.ok() {
5177 MergeOutcome {
5178 mode,
5179 ok: false,
5180 detail: pushed.stderr,
5181 }
5182 } else {
5183 let out =
5184 gh_pr_create(&winner.worktree, &base, &winner.branch, &pr.title, &pr.body)
5185 .await;
5186 match out {
5187 Ok(url) => MergeOutcome {
5188 mode,
5189 ok: true,
5190 detail: url,
5191 },
5192 Err(e) => MergeOutcome {
5193 mode,
5194 ok: false,
5195 detail: e.to_string(),
5196 },
5197 }
5198 }
5199 }
5200 };
5201
5202 self.state.status = match (mode, outcome.ok) {
5203 (MergeMode::None, _) => RunStatus::Ready,
5204 (_, true) => RunStatus::Merged,
5205 (_, false) => RunStatus::Blocked,
5206 };
5207 self.state.event(
5208 "merge",
5209 format!(
5210 "{:?}: {}",
5211 mode,
5212 outcome.detail.lines().next().unwrap_or("")
5213 ),
5214 );
5215 self.state.merge = Some(outcome);
5216 self.state.save()?;
5217
5218 if self.state.config.graph.land
5224 && mode == MergeMode::Pr
5225 && self.state.status == RunStatus::Merged
5226 {
5227 self.run_land().await?;
5228 }
5229 self.settle_questions();
5234 Ok(())
5235 }
5236
5237 async fn run_land(&mut self) -> Result<()> {
5248 let url = self
5249 .state
5250 .merge
5251 .as_ref()
5252 .map(|m| m.detail.clone())
5253 .unwrap_or_default();
5254 let url = url.lines().next().unwrap_or("").trim().to_owned();
5255 if !url.starts_with("http") {
5256 return Ok(());
5257 }
5258 match land::land(&mut self.state, &url).await {
5261 Ok(pr) if self.state.parked => {
5262 let _ = pr;
5266 }
5267 Ok(pr) => {
5268 self.state.status = match pr.state {
5269 land::PrLifecycle::Merged => RunStatus::Merged,
5270 _ => RunStatus::Blocked,
5271 };
5272 if bump::should_release_bump(self.state.status)
5279 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5280 {
5281 self.state
5287 .event("bump", format!("release bump skipped: {e:#}"));
5288 }
5289 self.state.save()?;
5290 }
5291 Err(e) => {
5292 self.state.status = RunStatus::Blocked;
5293 self.state.event("land", format!("gave up: {e}"));
5294 self.state.save()?;
5295 }
5296 }
5297 Ok(())
5298 }
5299
5300 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5304 if let Some(existing) = self.state.seats.get(key)
5305 && existing.agent == agent
5306 {
5307 return existing.clone();
5308 }
5309 let fresh = SeatState::new(key, agent, self.state.seed);
5310 self.state.seats.insert(key.to_owned(), fresh.clone());
5311 fresh
5312 }
5313
5314 fn view(&self, c: &Candidate) -> CandidateView {
5316 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5317 .unwrap_or_default();
5318 let (patch, _) = blind::sanitize_patch(
5319 &format!("candidate {} patch", c.label),
5320 &raw,
5321 &self.state.config.blind,
5322 );
5323 CandidateView {
5324 label: c.label,
5325 branch: c.branch.clone(),
5326 summary: c.summary.clone(),
5327 stat: c.stat.clone(),
5328 patch,
5329 }
5330 }
5331
5332 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5334 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5335 prompt::judge(
5336 "(see above)",
5337 &views,
5338 self.roles.judges.len(),
5339 base_short,
5340 "en",
5341 )
5342 }
5343
5344 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5351 let mut turns = Vec::new();
5352 for j in &self.state.judgements {
5353 if j.ranking.is_empty() {
5354 continue;
5355 }
5356 let reasons = j
5357 .reasons
5358 .iter()
5359 .map(|(k, v)| format!("- {k}: {v}"))
5360 .collect::<Vec<_>>()
5361 .join("\n");
5362 turns.push(Turn {
5363 who: format!("Judge {} (opening ranking)", j.judge),
5364 is_self: j.judge == self_idx + 1,
5365 body: format!(
5366 "Ranked {}{}{reasons}",
5367 j.ranking.iter().collect::<String>(),
5368 if reasons.is_empty() {
5369 ""
5370 } else {
5371 ", because:\n"
5372 }
5373 ),
5374 });
5375 }
5376 for t in self
5377 .state
5378 .deliberation
5379 .iter()
5380 .flat_map(|r| r.turns.iter())
5381 .chain(current)
5382 {
5383 turns.push(Turn {
5384 who: format!("Judge {}", t.judge),
5385 is_self: t.judge == self_idx + 1,
5386 body: t.body.clone(),
5387 });
5388 }
5389 turns
5390 }
5391}
5392
5393fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5395 agent::has_session(spec.kind, seat, sessions)
5396}
5397
5398fn next_untried_implementer<'a>(
5419 roster: &'a [AgentSpec],
5420 start: usize,
5421 tried: &BTreeSet<String>,
5422) -> Option<&'a AgentSpec> {
5423 roster
5424 .get(start + 1..)?
5425 .iter()
5426 .find(|s| !tried.contains(&s.id))
5427}
5428
5429fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5443 commands.iter().any(|c| c.exit_code.is_none())
5444}
5445
5446fn verified_noop_claim(
5459 usable: bool,
5460 commands: &[agent::CommandEvidence],
5461 text: &str,
5462) -> Option<String> {
5463 (usable && !has_unconfirmed_command(commands))
5464 .then(|| verdict::verified_noop(text))
5465 .flatten()
5466}
5467
5468fn short(commit: &str) -> String {
5469 commit.chars().take(7).collect()
5470}
5471
5472fn make_executable(path: &Path) -> Result<()> {
5473 #[cfg(unix)]
5474 {
5475 use std::os::unix::fs::PermissionsExt as _;
5476 let mut perms = std::fs::metadata(path)?.permissions();
5477 perms.set_mode(0o755);
5478 std::fs::set_permissions(path, perms)?;
5479 }
5480 #[cfg(not(unix))]
5481 {
5482 let _ = path;
5483 }
5484 Ok(())
5485}
5486
5487struct WaveCtx<'a> {
5494 run: &'a str,
5497 node: &'a str,
5499 prompts: &'a Prompts,
5500 cache: Option<&'a Path>,
5502 round: Option<usize>,
5505}
5506
5507async fn run_one(
5509 job: SeatJob,
5510 sem: Arc<Semaphore>,
5511 ctx: &WaveCtx<'_>,
5512 state: &mut RunState,
5513 attempt: usize,
5514) -> (SeatState, AgentOutcome) {
5515 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5516 .await
5517 .pop()
5518 .expect("one job in, one result out");
5519 (seat, out)
5520}
5521
5522async fn wave(
5528 jobs: Vec<SeatJob>,
5529 sem: Arc<Semaphore>,
5530 ctx: &WaveCtx<'_>,
5531 state: &mut RunState,
5532 attempt: usize,
5533) -> Vec<(usize, SeatState, AgentOutcome)> {
5534 let WaveCtx {
5535 run,
5536 node,
5537 prompts,
5538 cache,
5539 round,
5540 } = *ctx;
5541 for job in &jobs {
5542 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5543 }
5544 if let Err(e) = state.save() {
5545 tracing::warn!("could not persist in-progress seats: {e:#}");
5550 }
5551 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5567 let wait_started = Instant::now();
5568 let cache_guard = if let Some(cache_dir) = cache {
5569 if jobs_had_a_writer {
5570 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5571 let budget = jobs
5572 .iter()
5573 .map(|j| j.timeout)
5574 .max()
5575 .unwrap_or(Duration::from_secs(60));
5576 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5577 .await
5578 .ok()
5579 } else {
5580 None
5581 }
5582 } else {
5583 None
5584 };
5585 let waited_for_lease = wait_started.elapsed();
5592 let mut set = tokio::task::JoinSet::new();
5593 let overlay = prompts.overlay(node);
5594 for (i, mut job) in jobs.into_iter().enumerate() {
5595 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5596 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
5597 if cache.is_some() {
5598 job.prompt.push('\n');
5599 job.prompt
5600 .push_str(&prompt::build_cache_note(node, job.allow_write));
5601 }
5602 let sem = Arc::clone(&sem);
5603 let run = run.to_owned();
5604 let node = node.to_owned();
5605 let cache = cache
5616 .filter(|_| job.allow_write && cache_guard.is_some())
5617 .map(Path::to_path_buf);
5618 set.spawn(async move {
5619 let _permit = sem.acquire().await;
5620 let mut seat = job.seat;
5621 let out = agent::invoke(
5622 &job.spec,
5623 &mut seat,
5624 &Invocation {
5625 cwd: &job.cwd,
5626 prompt: &job.prompt,
5627 timeout: job.timeout,
5628 allow_write: job.allow_write,
5629 sessions: job.sessions,
5630 artifacts: &job.artifacts,
5631 stem: &job.stem,
5632 run: &run,
5633 node: &node,
5634 cache_dir: cache.as_deref(),
5635 attachments: &[],
5636 },
5637 )
5638 .await;
5639 let out = match out {
5640 Ok(o) if o.usable() => AgentOutcome::Ok(o),
5641 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
5642 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
5650 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
5651 Ok(o) => AgentOutcome::Failed(format!(
5652 "exited with {:?} and no usable output",
5653 o.exit_code
5654 )),
5655 Err(e) => AgentOutcome::Failed(e.to_string()),
5656 };
5657 (i, seat, out)
5658 });
5659 }
5660 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
5661 while let Some(joined) = set.join_next().await {
5662 let (i, seat, out) = match joined {
5663 Ok(v) => v,
5664 Err(e) => {
5668 tracing::error!("agent task panicked: {e}");
5669 continue;
5670 }
5671 };
5672 state.seat_finished(&seat.key);
5673 record_jobs(state, node, round, &seat.key, &out);
5674 if let Err(e) = state.save() {
5675 tracing::warn!("could not persist a seat's completion: {e:#}");
5676 }
5677 if collected.len() <= i {
5678 collected.resize_with(i + 1, || None);
5679 }
5680 collected[i] = Some((i, seat, out));
5681 }
5682 if state
5688 .active
5689 .values()
5690 .any(|a| a.node == node && a.attempt == attempt)
5691 {
5692 state
5693 .active
5694 .retain(|_, a| !(a.node == node && a.attempt == attempt));
5695 if let Err(e) = state.save() {
5696 tracing::warn!("could not persist the end of a wave: {e:#}");
5697 }
5698 }
5699 if let Some(cache_dir) = cache
5706 && jobs_had_a_writer
5707 {
5708 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
5709 }
5710 if let Some(guard) = cache_guard {
5711 guard.release();
5712 }
5713 collected.into_iter().flatten().collect()
5714}
5715
5716fn record_jobs(
5727 state: &mut RunState,
5728 node: &str,
5729 round: Option<usize>,
5730 seat: &str,
5731 out: &AgentOutcome,
5732) {
5733 let commands: &[agent::CommandEvidence] = match out {
5734 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
5735 AgentOutcome::Failed(_) => &[],
5736 };
5737 let checked_at = Timestamp::now();
5738 for c in commands {
5739 state.jobs.push(JobRecord {
5740 node: node.to_owned(),
5741 round,
5742 seat: seat.to_owned(),
5743 id: c.id.clone(),
5744 description: c.description.clone(),
5745 checked_at,
5746 status: match c.exit_code {
5747 Some(0) => JobStatus::Completed,
5748 Some(_) => JobStatus::Failed,
5749 None => JobStatus::Unknown,
5750 },
5751 exit_code: c.exit_code,
5752 result_summary: c.result_summary.clone(),
5753 source: c.source.clone(),
5754 });
5755 }
5756}
5757
5758fn round_is_clean(
5779 blocking: usize,
5780 e2e_ok: bool,
5781 answered: usize,
5782 expected: usize,
5783 quota_missing: usize,
5784 policy: IncompleteReviewPolicy,
5785) -> bool {
5786 if blocking != 0 || !e2e_ok {
5787 return false;
5788 }
5789 if answered == expected || policy == IncompleteReviewPolicy::Warn {
5790 return true;
5791 }
5792 answered > 0 && expected - answered <= quota_missing
5793}
5794
5795fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
5819 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
5820 return Some(RunStatus::Gating);
5821 }
5822 let last = reviews.last()?;
5823 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
5824 if reviews.len() < max_rounds && !stagnant {
5825 return None;
5826 }
5827 if last.incomplete() && last.blocking == 0 {
5828 return Some(RunStatus::Blocked);
5829 }
5830 if last.e2e_status() == E2eStatus::ResourceBlocked {
5831 return None;
5832 }
5833 Some(if last.e2e.iter().all(CommandOutcome::ok) {
5834 RunStatus::Gating
5835 } else {
5836 RunStatus::Blocked
5837 })
5838}
5839
5840fn retry_budget(full: Duration, nudged: bool) -> Duration {
5855 if nudged {
5856 (full / 4).max(Duration::from_secs(120)).min(full)
5857 } else {
5858 full
5859 }
5860}
5861
5862#[allow(clippy::too_many_arguments)]
5875async fn ask_json_wave<T>(
5876 jobs: Vec<SeatJob>,
5877 sem: Arc<Semaphore>,
5878 retries: usize,
5879 ctx: &WaveCtx<'_>,
5880 losses: &mut Vec<QuotaLoss>,
5881 state: &mut RunState,
5882 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
5883) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
5884where
5885 T: serde::de::DeserializeOwned + Send + 'static,
5886{
5887 let n = jobs.len();
5888 let originals: Vec<SeatJob> = jobs;
5889 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
5890 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
5891 let mut attempts_used: Vec<usize> = vec![0; n];
5898 let mut pending: Vec<usize> = (0..n).collect();
5899
5900 for attempt in 0..=retries {
5901 if pending.is_empty() {
5902 break;
5903 }
5904 let mut batch = Vec::with_capacity(pending.len());
5905 for &i in &pending {
5906 let src = &originals[i];
5907 let (prompt, timeout) = if attempt == 0 {
5910 (src.prompt.clone(), src.timeout)
5911 } else {
5912 let why = done[i]
5913 .as_ref()
5914 .and_then(|r| r.as_ref().err().map(ToString::to_string))
5915 .unwrap_or_else(|| "no parsable answer".to_owned());
5916 let nudge = prompt::nudge(&why);
5917 let nudged = has_context(&src.spec, &seats[i], src.sessions);
5918 let prompt = if nudged {
5919 nudge
5920 } else {
5921 format!("{}\n\n---\n\n{}", src.prompt, nudge)
5922 };
5923 (prompt, retry_budget(src.timeout, nudged))
5924 };
5925 batch.push(SeatJob {
5926 spec: src.spec.clone(),
5927 seat: seats[i].clone(),
5928 cwd: src.cwd.clone(),
5929 prompt,
5930 timeout,
5931 allow_write: src.allow_write,
5932 sessions: src.sessions,
5933 artifacts: src.artifacts.clone(),
5934 stem: if attempt == 0 {
5935 src.stem.clone()
5936 } else {
5937 format!("{}-retry{attempt}", src.stem)
5938 },
5939 });
5940 }
5941
5942 if attempt > 0 {
5943 let seats_out: Vec<&str> = pending
5944 .iter()
5945 .map(|&i| originals[i].seat.key.as_str())
5946 .collect();
5947 state.event(
5948 ctx.node,
5949 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
5950 );
5951 }
5952 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
5953 let mut still = Vec::new();
5954 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
5955 seats[i] = seat;
5956 let (parsed, quota) = match out {
5957 AgentOutcome::Ok(o) => (
5958 match verdict::extract_json::<T>(&o.text) {
5959 Ok(v) => match validate(&v) {
5960 Ok(()) => Ok((v, o)),
5961 Err(e) => Err(e),
5962 },
5963 Err(e) => Err(e),
5964 },
5965 false,
5966 ),
5967 AgentOutcome::Quota(o) => {
5968 losses.push(QuotaLoss {
5969 seat: originals[i].seat.key.clone(),
5970 node: ctx.node.to_owned(),
5971 at: Timestamp::now(),
5972 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5973 });
5974 (
5975 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
5976 true,
5977 )
5978 }
5979 AgentOutcome::Dropped(o) => {
5984 let why = o
5985 .dropped
5986 .as_ref()
5987 .map(|d| d.why.as_str())
5988 .unwrap_or("the CLI ended the stream without delivering its answer");
5989 (
5990 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
5991 false,
5992 )
5993 }
5994 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
5995 };
5996 let failed = parsed.is_err();
5997 done[i] = Some(parsed);
5998 attempts_used[i] = attempt;
5999 if failed && !quota {
6002 still.push(i);
6003 }
6004 }
6005 pending = still;
6006 }
6007
6008 seats
6009 .into_iter()
6010 .zip(done)
6011 .zip(attempts_used)
6012 .map(|((seat, res), attempts)| {
6013 (
6014 seat,
6015 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6016 attempts,
6017 )
6018 })
6019 .collect()
6020}
6021
6022async fn acquire_cache_lease(
6035 state: &mut RunState,
6036 cache_dir: &Path,
6037 owner: &crate::cache::Owner,
6038 budget: Duration,
6039 context: &str,
6040) -> Result<crate::cache::Guard> {
6041 let home = crate::run::home();
6042 let started = Instant::now();
6043 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6044 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6045 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6046 Err(e) => {
6047 state.event(
6048 "verify",
6049 format!("{context}: could not check the shared build cache: {e:#}"),
6050 );
6051 if let Err(e2) = state.save() {
6052 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6053 }
6054 return Err(e);
6055 }
6056 };
6057 state.event(
6058 "verify",
6059 format!(
6060 "{context}: waiting for the shared build cache at {} ({})",
6061 cache_dir.display(),
6062 busy.describe()
6063 ),
6064 );
6065 if let Err(e) = state.save() {
6066 tracing::warn!("could not persist a cache wait: {e:#}");
6067 }
6068 let remaining = budget.saturating_sub(started.elapsed());
6069 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6070 Ok(g) => Ok(g),
6071 Err(e) => {
6072 state.event("verify", format!("{context}: {e:#}"));
6073 if let Err(e2) = state.save() {
6074 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6075 }
6076 Err(e)
6077 }
6078 }
6079}
6080
6081#[allow(clippy::too_many_arguments)]
6102async fn with_cache_lease<'s, F, Fut>(
6103 state: &'s mut RunState,
6104 cache_dir: Option<&Path>,
6105 node: &str,
6106 seat: &str,
6107 worktree: &Path,
6108 head: &str,
6109 budget: Duration,
6110 context: &str,
6111 body: F,
6112) -> (Vec<CommandOutcome>, bool)
6113where
6114 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6115 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6116{
6117 let Some(cache_dir) = cache_dir else {
6118 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6119 return (outcomes, retried);
6120 };
6121 let home = crate::run::home();
6122 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6123 let started = Instant::now();
6124 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6125 Ok(g) => g,
6126 Err(e) => {
6127 return (
6128 vec![CommandOutcome {
6129 command: "(waiting for the shared build cache)".to_owned(),
6130 code: None,
6131 output_tail: e.to_string(),
6132 duration_ms: started.elapsed().as_millis() as u64,
6133 resource_blocked: true,
6134 }],
6135 false,
6136 );
6137 }
6138 };
6139 let identity = crate::cache::Identity::new(worktree, head);
6140 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6141 state.event(
6149 "verify",
6150 format!(
6151 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6152 worktree.display(),
6153 short(head)
6154 ),
6155 );
6156 guard.release();
6157 return (
6158 vec![CommandOutcome {
6159 command: "(confirming the shared build cache is fresh)".to_owned(),
6160 code: None,
6161 output_tail: e.to_string(),
6162 duration_ms: started.elapsed().as_millis() as u64,
6163 resource_blocked: true,
6164 }],
6165 false,
6166 );
6167 }
6168 let remaining = budget.saturating_sub(started.elapsed());
6169 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6170 if !timed_out_pids.is_empty() {
6175 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6176 }
6177 guard.release();
6178 (outcomes, retried)
6179}
6180
6181async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6193 wait_for_pids_with(
6194 pids,
6195 crate::proc::pid_alive,
6196 LEASE_RELEASE_POLL,
6197 LEASE_RELEASE_MAX_WAIT,
6198 )
6199 .await;
6200}
6201
6202async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6208 pids: &[u32],
6209 alive: F,
6210 poll: Duration,
6211 max_wait: Duration,
6212) {
6213 let deadline = Instant::now() + max_wait;
6214 loop {
6215 if pids.iter().all(|&pid| !alive(pid)) {
6216 return;
6217 }
6218 if Instant::now() >= deadline {
6219 return;
6220 }
6221 tokio::time::sleep(poll).await;
6222 }
6223}
6224
6225fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6232 outcomes.iter().any(|o| o.resource_blocked)
6233}
6234
6235enum GateFix {
6237 Retry,
6239 Stop,
6242 Defer,
6245}
6246
6247fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6255 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6256 red.peek().is_some()
6257 && red.all(|o| {
6258 !o.resource_blocked
6259 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6260 && !o.output_tail.trim().is_empty()
6261 })
6262}
6263
6264fn e2e_outcome_label(o: &CommandOutcome) -> String {
6268 if o.ok() {
6269 return "pass".to_owned();
6270 }
6271 let reason = if o.build_failed() {
6272 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6273 } else {
6274 format!("FAIL ({:?})", o.code)
6275 };
6276 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6277}
6278
6279async fn run_e2e_with_retry(
6287 state: &mut RunState,
6288 shell: &[String],
6289 commands: &[String],
6290 worktree: &Path,
6291 timeout: Duration,
6292 context: &str,
6293) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6294 let (mut e2e, mut timed_out_pids) = run_commands(
6295 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6296 )
6297 .await;
6298 for o in &e2e {
6299 state.event(
6300 "verify",
6301 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6302 );
6303 }
6304 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6308 if verify_retried {
6309 state.event(
6310 "verify",
6311 format!(
6312 "{context}: verify could not build/link, not a test result — retrying once \
6313 before concluding"
6314 ),
6315 );
6316 let retried = run_commands(
6317 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6318 )
6319 .await;
6320 e2e = retried.0;
6321 timed_out_pids.extend(retried.1);
6324 for o in &e2e {
6325 state.event(
6326 "verify",
6327 format!(
6328 "{context}: retry `{}` -> {}",
6329 o.command,
6330 e2e_outcome_label(o)
6331 ),
6332 );
6333 }
6334 }
6335 (e2e, verify_retried, timed_out_pids)
6336}
6337
6338#[allow(clippy::too_many_arguments)]
6354async fn run_commands(
6355 state: &mut RunState,
6356 node: &str,
6357 task: &str,
6358 attempt: usize,
6359 shell: &[String],
6360 commands: &[String],
6361 cwd: &Path,
6362 timeout: Duration,
6363) -> (Vec<CommandOutcome>, Vec<u32>) {
6364 if commands.is_empty() {
6365 return (Vec::new(), Vec::new());
6370 }
6371 let mut out = Vec::new();
6372 let mut timed_out_pids = Vec::new();
6373 let total = commands.len();
6374 for (idx, command) in commands.iter().enumerate() {
6375 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6376 if let Err(e) = state.save() {
6377 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6378 }
6379 let started = Instant::now();
6380 let mut cmd = tokio::process::Command::new(&shell[0]);
6381 cmd.quiet();
6382 cmd.args(&shell[1..])
6383 .arg(command)
6384 .current_dir(cwd)
6385 .stdin(std::process::Stdio::null())
6386 .stdout(std::process::Stdio::piped())
6387 .stderr(std::process::Stdio::piped())
6388 .kill_on_drop(true);
6389 let spawned = cmd.spawn();
6390 let (code, body) = match spawned {
6391 Ok(child) => {
6392 let pid = child.id();
6397 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6398 Ok(Ok(o)) => {
6399 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6400 body.push_str(&String::from_utf8_lossy(&o.stderr));
6401 (o.status.code(), body)
6402 }
6403 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6404 Err(_) => {
6405 if let Some(pid) = pid {
6406 timed_out_pids.push(pid);
6407 }
6408 (None, format!("timed out after {}s", timeout.as_secs()))
6409 }
6410 }
6411 }
6412 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6413 };
6414 out.push(CommandOutcome {
6415 command: command.clone(),
6416 code,
6417 output_tail: tail(&body, OUTPUT_TAIL),
6418 duration_ms: started.elapsed().as_millis() as u64,
6419 resource_blocked: false,
6420 });
6421 }
6422 state.task_finished(task);
6423 if let Err(e) = state.save() {
6424 tracing::warn!("could not persist the end of {task}: {e:#}");
6425 }
6426 (out, timed_out_pids)
6427}
6428
6429fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6441 let repo = repo.display();
6442 match style {
6443 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6444 MergeStyle::Squash => {
6445 let subject = message
6448 .lines()
6449 .next()
6450 .unwrap_or(branch)
6451 .replace(['\\', '"', '$', '`'], "");
6452 format!(
6453 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6454 )
6455 }
6456 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6457 }
6458}
6459
6460const PR_TITLE_MAX: usize = 240;
6471
6472struct PrMessage {
6477 title: String,
6478 body: String,
6479}
6480
6481impl PrMessage {
6482 fn commit_message(&self) -> String {
6486 format!("{}\n\n{}", self.title, self.body)
6487 }
6488}
6489
6490fn title_marker(line: &str) -> Option<&str> {
6492 let line = line.trim();
6493 let head = line.get(..6)?;
6494 head.eq_ignore_ascii_case("title:")
6495 .then(|| line[6..].trim())
6496}
6497
6498fn summary_title(summary: &str) -> Option<String> {
6503 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6504 let raw = title_marker(first)?;
6505 if raw.is_empty() {
6506 return None;
6507 }
6508 let title = queue::title_from(raw, PR_TITLE_MAX);
6509 let lower = title.to_ascii_lowercase();
6510 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6511 return None;
6512 }
6513 Some(title)
6514}
6515
6516fn summary_without_title(summary: &str) -> String {
6519 let mut lines = summary.trim().lines().peekable();
6520 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6521 lines.next();
6522 }
6523 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6524}
6525
6526fn pr_message(state: &RunState, winner: char) -> PrMessage {
6540 let summary = state
6541 .candidates
6542 .iter()
6543 .find(|c| c.label == winner)
6544 .map(|c| c.summary.as_str())
6545 .unwrap_or_default();
6546 let title = summary_title(summary).unwrap_or_else(|| {
6549 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6550 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6551 t
6552 } else {
6553 format!(
6554 "chore: land candidate {} of run {}",
6555 winner.to_ascii_uppercase(),
6556 state.id
6557 )
6558 }
6559 });
6560
6561 let mut body = String::new();
6562 let what = summary_without_title(summary);
6563 if !what.is_empty() {
6564 body.push_str("## Summary\n\n");
6565 body.push_str(&what);
6566 body.push_str("\n\n");
6567 }
6568
6569 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6570 if let Some(fix) = fix
6571 && !fix.notes.trim().is_empty()
6572 {
6573 body.push_str("## Review fixes\n\n");
6574 body.push_str(fix.notes.trim());
6575 body.push_str("\n\n");
6576 }
6577
6578 let open = state.open_findings();
6579 if !open.is_empty() {
6580 body.push_str("## Open review findings\n\n");
6581 for f in &open {
6582 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6583 }
6584 body.push('\n');
6585 }
6586
6587 if let Some(fix) = fix
6588 && !fix.rejected.is_empty()
6589 {
6590 body.push_str("## Declined by the fixer\n\n");
6591 for r in &fix.rejected {
6592 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
6593 }
6594 body.push('\n');
6595 }
6596
6597 let task = state.instruction.trim();
6598 let task = if task.is_empty() {
6599 "(empty task)"
6600 } else {
6601 task
6602 };
6603 body.push_str(&format!(
6604 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
6605 task.replace("</details>", "</details>")
6606 ));
6607
6608 body.push_str(&format!(
6609 "\n---\nmagi:run/{} magi:candidate-{}\n",
6610 state.id,
6611 winner.to_ascii_lowercase()
6612 ));
6613
6614 let id = crate::scrub::Identity::current();
6617 PrMessage {
6618 title: crate::scrub::scrub(&title, &id),
6619 body: crate::scrub::scrub(&body, &id),
6620 }
6621}
6622
6623async fn gh_pr_create(
6625 cwd: &Path,
6626 base: &str,
6627 head: &str,
6628 title: &str,
6629 body: &str,
6630) -> Result<String> {
6631 let out = tokio::process::Command::new("gh")
6632 .args([
6633 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
6634 ])
6635 .current_dir(cwd)
6636 .quiet()
6637 .stdin(std::process::Stdio::null())
6638 .output()
6639 .await
6640 .context("spawn gh")?;
6641 if out.status.success() {
6642 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
6643 } else {
6644 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
6645 }
6646}
6647
6648pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
6657 let repo = state.repo.clone();
6658 let root = state.worktree_root();
6659 let winner = state.tally.as_ref().map(|t| t.winner);
6660 let mut removed = Vec::new();
6661
6662 for i in 0..state.candidates.len() {
6663 let c = state.candidates[i].clone();
6664 let is_winner = Some(c.label) == winner;
6665 if is_winner && !drop_winner {
6666 continue;
6667 }
6668 if c.worktree.exists() {
6669 git::worktree_remove(&repo, &c.worktree).await.ok();
6670 removed.push(c.worktree.to_string_lossy().into_owned());
6671 }
6672 if git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
6673 git::branch_delete(&repo, &c.branch).await.ok();
6674 removed.push(c.branch.clone());
6675 }
6676 state.candidates[i].folded = true;
6677 }
6678
6679 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
6680 let path = name.path();
6681 let keep = !drop_winner
6682 && winner.is_some_and(|w| {
6683 path.file_name()
6684 .is_some_and(|n| n == format!("cand-{w}").as_str())
6685 });
6686 if keep {
6687 continue;
6688 }
6689 git::worktree_remove(&repo, &path).await.ok();
6690 removed.push(path.to_string_lossy().into_owned());
6691 }
6692
6693 remove_if_empty(&root);
6702
6703 if state.enabled_worktree_config && drop_winner {
6704 git::release_worktree_config(&repo).await.ok();
6708 state.enabled_worktree_config = false;
6709 }
6710 state.save_under(home)?;
6711 Ok(removed)
6712}
6713
6714fn remove_if_empty(dir: &Path) {
6725 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
6726 std::fs::remove_dir(dir).ok();
6727 }
6728}
6729
6730pub fn worst_open(state: &RunState) -> Option<Severity> {
6732 state
6733 .reviews
6734 .last()?
6735 .reviews
6736 .iter()
6737 .flat_map(|r| r.findings.iter())
6738 .map(|f| f.severity)
6739 .max()
6740}
6741
6742#[cfg(test)]
6743mod tests {
6744 use super::*;
6745 use crate::run::GateStatus;
6746 use std::collections::BTreeMap;
6747 use std::time::Duration;
6748
6749 fn conductor() -> AgentSpec {
6750 AgentSpec {
6751 id: "conductor".to_owned(),
6752 kind: crate::config::AgentKind::Command,
6753 model: None,
6754 command: vec!["true".to_owned()],
6755 extra_args: Vec::new(),
6756 env: BTreeMap::new(),
6757 prompt_delivery: None,
6758 }
6759 }
6760
6761 fn spec(id: &str) -> AgentSpec {
6762 AgentSpec {
6763 id: id.to_owned(),
6764 kind: crate::config::AgentKind::Command,
6765 model: None,
6766 command: vec!["true".to_owned()],
6767 extra_args: Vec::new(),
6768 env: BTreeMap::new(),
6769 prompt_delivery: None,
6770 }
6771 }
6772
6773 #[test]
6780 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
6781 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6782 let tried = BTreeSet::from(["beta".to_owned()]);
6783 let next = next_untried_implementer(&roster, 1, &tried);
6786 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
6787 }
6788
6789 #[test]
6790 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
6791 let roster = vec![spec("alpha"), spec("beta")];
6792 let tried = BTreeSet::from(["beta".to_owned()]);
6793 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6797 }
6798
6799 #[test]
6800 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
6801 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6802 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
6803 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6807 }
6808
6809 #[test]
6810 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
6811 let roster = vec![spec("a"), spec("a"), spec("b")];
6812 let tried = BTreeSet::from(["a".to_owned()]);
6813 let next = next_untried_implementer(&roster, 0, &tried);
6814 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
6815 }
6816
6817 #[test]
6818 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
6819 let roster = vec![spec("a"), spec("b")];
6820 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
6821 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
6822 }
6823
6824 #[test]
6825 fn remove_if_empty_only_ever_takes_a_bare_directory() {
6826 let dir = tempfile::tempdir().unwrap();
6827 let bay = dir.path().join("ffff");
6828
6829 remove_if_empty(&bay);
6831 assert!(!bay.exists());
6832
6833 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
6836 remove_if_empty(&bay);
6837 assert!(bay.exists(), "non-empty directory must survive");
6838
6839 std::fs::remove_dir(bay.join("cand-A")).unwrap();
6841 remove_if_empty(&bay);
6842 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
6843 }
6844
6845 #[test]
6854 fn a_full_panel_that_found_nothing_is_clean() {
6855 assert!(round_is_clean(
6856 0,
6857 true,
6858 2,
6859 2,
6860 0,
6861 IncompleteReviewPolicy::Block
6862 ));
6863 }
6864
6865 #[test]
6866 fn a_missing_seat_is_never_clean_under_the_default_policy() {
6867 assert!(!round_is_clean(
6868 0,
6869 true,
6870 1,
6871 2,
6872 0,
6873 IncompleteReviewPolicy::Block
6874 ));
6875 }
6876
6877 #[test]
6878 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
6879 assert!(!round_is_clean(
6880 1,
6881 true,
6882 1,
6883 2,
6884 0,
6885 IncompleteReviewPolicy::Warn
6886 ));
6887 }
6888
6889 #[test]
6890 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
6891 assert!(round_is_clean(
6892 0,
6893 true,
6894 1,
6895 2,
6896 0,
6897 IncompleteReviewPolicy::Warn
6898 ));
6899 }
6900
6901 #[test]
6902 fn a_full_panel_with_an_open_finding_is_not_clean() {
6903 assert!(!round_is_clean(
6904 1,
6905 true,
6906 2,
6907 2,
6908 0,
6909 IncompleteReviewPolicy::Block
6910 ));
6911 }
6912
6913 #[test]
6914 fn a_full_panel_with_a_red_e2e_is_not_clean() {
6915 assert!(!round_is_clean(
6916 0,
6917 false,
6918 2,
6919 2,
6920 0,
6921 IncompleteReviewPolicy::Block
6922 ));
6923 }
6924
6925 #[test]
6932 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
6933 assert!(round_is_clean(
6936 0,
6937 true,
6938 1,
6939 2,
6940 1,
6941 IncompleteReviewPolicy::Block
6942 ));
6943 }
6944
6945 #[test]
6946 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
6947 assert!(!round_is_clean(
6950 0,
6951 true,
6952 1,
6953 2,
6954 0,
6955 IncompleteReviewPolicy::Block
6956 ));
6957 }
6958
6959 #[test]
6960 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
6961 assert!(!round_is_clean(
6962 1,
6963 true,
6964 1,
6965 2,
6966 1,
6967 IncompleteReviewPolicy::Block
6968 ));
6969 assert!(!round_is_clean(
6970 0,
6971 false,
6972 1,
6973 2,
6974 1,
6975 IncompleteReviewPolicy::Block
6976 ));
6977 }
6978
6979 #[test]
6980 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
6981 assert!(!round_is_clean(
6985 0,
6986 true,
6987 0,
6988 2,
6989 2,
6990 IncompleteReviewPolicy::Block
6991 ));
6992 }
6993
6994 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
6995 CommandOutcome {
6996 command: "test".to_owned(),
6997 code,
6998 output_tail: String::new(),
6999 duration_ms: 0,
7000 resource_blocked,
7001 }
7002 }
7003
7004 #[test]
7005 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7006 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7007 assert!(
7008 !verify_inconclusive(&[outcome(Some(1), false)]),
7009 "an ordinary failure is still evidence about the patch"
7010 );
7011 assert!(verify_inconclusive(&[outcome(None, true)]));
7012 assert!(
7013 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7014 "one inconclusive outcome taints the whole batch"
7015 );
7016 assert!(!verify_inconclusive(&[]));
7017 }
7018
7019 #[tokio::test]
7020 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7021 let calls = std::sync::atomic::AtomicUsize::new(0);
7025 let started = Instant::now();
7026 wait_for_pids_with(
7027 &[123],
7028 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7029 Duration::from_millis(5),
7030 Duration::from_secs(5),
7031 )
7032 .await;
7033 assert!(
7034 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7035 "must keep checking rather than deciding on the first answer"
7036 );
7037 assert!(
7038 started.elapsed() < Duration::from_secs(1),
7039 "must return the moment it is confirmed dead, not wait out the ceiling"
7040 );
7041 }
7042
7043 #[tokio::test]
7044 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7045 let started = Instant::now();
7046 wait_for_pids_with(
7047 &[123],
7048 |_| true, Duration::from_millis(5),
7050 Duration::from_millis(30),
7051 )
7052 .await;
7053 let elapsed = started.elapsed();
7054 assert!(
7055 elapsed >= Duration::from_millis(30),
7056 "must not give up before its own ceiling: {elapsed:?}"
7057 );
7058 assert!(
7059 elapsed < Duration::from_secs(1),
7060 "must not wait past its own ceiling either: {elapsed:?}"
7061 );
7062 }
7063
7064 #[tokio::test]
7065 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7066 let started = Instant::now();
7067 wait_for_pids_with(
7068 &[],
7069 |_| true,
7070 Duration::from_secs(5),
7071 Duration::from_secs(5),
7072 )
7073 .await;
7074 assert!(
7075 started.elapsed() < Duration::from_millis(200),
7076 "an empty pid list has nothing to confirm"
7077 );
7078 }
7079
7080 fn review_round(
7086 clean: bool,
7087 blocking: usize,
7088 answered: usize,
7089 expected: usize,
7090 progressed: bool,
7091 e2e_ok: bool,
7092 ) -> ReviewRound {
7093 ReviewRound {
7094 round: 1,
7095 head: "h".to_owned(),
7096 verified_head: None,
7097 verified_at: None,
7098 reviews: Vec::new(),
7099 e2e: vec![CommandOutcome {
7100 command: "test".to_owned(),
7101 code: Some(if e2e_ok { 0 } else { 1 }),
7102 output_tail: String::new(),
7103 duration_ms: 0,
7104 resource_blocked: false,
7105 }],
7106 verify_retried: false,
7107 e2e_deferred: false,
7108 e2e_defer_reason: None,
7109 fix: None,
7110 blocking,
7111 answered,
7112 expected,
7113 clean,
7114 progressed,
7115 vote_split: false,
7116 reconsideration: Vec::new(),
7117 verdict: None,
7118 }
7119 }
7120
7121 #[test]
7122 fn review_conclusion_is_none_when_nothing_has_run() {
7123 assert_eq!(review_conclusion(&[], 3), None);
7124 }
7125
7126 #[test]
7127 fn review_conclusion_is_none_while_rounds_remain() {
7128 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7129 assert_eq!(review_conclusion(&rounds, 3), None);
7130 }
7131
7132 #[test]
7133 fn review_conclusion_is_gating_once_a_round_is_clean() {
7134 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7135 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7136 }
7137
7138 #[test]
7139 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7140 let rounds = vec![
7141 review_round(false, 1, 2, 2, true, true),
7142 review_round(false, 1, 2, 2, true, true),
7143 ];
7144 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7145 }
7146
7147 #[test]
7148 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7149 let rounds = vec![
7150 review_round(false, 1, 2, 2, true, true),
7151 review_round(false, 1, 2, 2, true, false),
7152 ];
7153 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7154 }
7155
7156 #[test]
7157 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7158 let mut blocked = review_round(false, 1, 2, 2, true, false);
7165 blocked.e2e[0].resource_blocked = true;
7166 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7167 assert_eq!(review_conclusion(&rounds, 2), None);
7168 }
7169
7170 #[test]
7171 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7172 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7174 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7175 }
7176
7177 #[test]
7178 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7179 let rounds = vec![
7180 review_round(false, 1, 2, 2, false, true),
7181 review_round(false, 1, 2, 2, false, true),
7182 ];
7183 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7184 }
7185
7186 fn secs(n: u64) -> Duration {
7187 Duration::from_secs(n)
7188 }
7189
7190 fn init_repo(dir: &Path) {
7193 let run = |args: &[&str]| {
7194 let out = std::process::Command::new("git")
7195 .args(args)
7196 .current_dir(dir)
7197 .quiet()
7198 .output()
7199 .expect("spawn git");
7200 assert!(
7201 out.status.success(),
7202 "git {args:?} failed: {}",
7203 String::from_utf8_lossy(&out.stderr)
7204 );
7205 };
7206 run(&["init", "-b", "main"]);
7207 run(&["config", "user.name", "magi test"]);
7208 run(&["config", "user.email", "magi@example.com"]);
7209 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7210 run(&["add", "-A"]);
7211 run(&["commit", "-m", "init"]);
7212 }
7213
7214 fn ask_test_home() {
7222 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7223 }
7224
7225 fn runner_at(status: RunStatus) -> Runner {
7228 let mut state = RunState::new(
7229 PathBuf::from("/nonexistent/repo"),
7230 "main".to_owned(),
7231 "deadbeef".to_owned(),
7232 "task".to_owned(),
7233 Config::default(),
7234 );
7235 state.status = status;
7236 Runner {
7237 state,
7238 roles: ResolvedRoles {
7239 implementers: Vec::new(),
7240 judges: Vec::new(),
7241 reviewers: Vec::new(),
7242 fixer: None,
7243 conductor: conductor(),
7244 implementer_roster: Vec::new(),
7245 },
7246 sem: Arc::new(Semaphore::new(1)),
7247 pause: Pause::new(),
7248 interrupt: Pause::new(),
7249 }
7250 }
7251
7252 #[test]
7256 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7257 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7258 let mut runner = runner_at(RunStatus::Implementing);
7259 let interrupt = Pause::new();
7260 runner.watch_interrupt(interrupt.clone());
7261
7262 interrupt.park_because("task a1b2 asked to run first");
7263
7264 assert!(runner.park_here().expect("park_here"));
7265 assert!(runner.state.parked);
7266 let last = runner.state.events.last().expect("a park event");
7267 assert_eq!(last.node, "park");
7268 assert!(
7269 last.message.contains("task a1b2 asked to run first"),
7270 "expected the interrupt reason in {:?}",
7271 last.message
7272 );
7273 }
7274
7275 #[test]
7283 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7284 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7285 let mut runner = runner_at(RunStatus::Implementing);
7286 let shutdown = Pause::new();
7287 runner.on_pause(shutdown.clone());
7288 let interrupt = Pause::new();
7289 runner.watch_interrupt(interrupt.clone());
7290
7291 assert!(!runner.park_here().expect("park_here"));
7293 assert!(!runner.state.parked);
7294
7295 interrupt.park_because("test");
7297 assert!(!shutdown.parked());
7298 assert!(runner.park_here().expect("park_here"));
7299 }
7300
7301 #[tokio::test]
7315 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7316 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7317 let mut runner = runner_at(RunStatus::Implementing);
7318 let interrupt = Pause::new();
7319 runner.watch_interrupt(interrupt.clone());
7320
7321 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7322 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7323
7324 let node = async move {
7328 started_tx.send(()).expect("send started");
7329 finish_rx.await.expect("recv finish");
7330 "node finished"
7331 };
7332
7333 let interrupter = async move {
7334 started_rx.await.expect("recv started");
7335 interrupt.park_because("higher-priority task waiting");
7337 tokio::task::yield_now().await;
7341 finish_tx.send(()).expect("send finish");
7342 };
7343
7344 let (node_result, ()) = tokio::join!(node, interrupter);
7345 assert_eq!(
7346 node_result, "node finished",
7347 "the in-flight call ran to completion"
7348 );
7349
7350 assert!(runner.park_here().expect("park_here"));
7353 assert!(runner.state.parked);
7354 }
7355
7356 #[test]
7362 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7363 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7364 let mut runner = runner_at(RunStatus::Judging);
7365 runner.state.config.agents = vec![conductor()];
7369 runner.state.candidates = vec![Candidate {
7370 index: 0,
7371 label: 'A',
7372 agent: "alpha".to_owned(),
7373 branch: "magi/x/A".to_owned(),
7374 worktree: PathBuf::from("/nonexistent/worktree"),
7375 summary: "did the thing".to_owned(),
7376 stat: "1 file changed".to_owned(),
7377 files: 1,
7378 commits: 1,
7379 empty: false,
7380 failed: None,
7381 verified_noop: None,
7382 duration_ms: 1234,
7383 folded: false,
7384 }];
7385 let run_id = runner.state.id.clone();
7386
7387 let interrupt = Pause::new();
7388 runner.watch_interrupt(interrupt.clone());
7389 interrupt.park_because("task c3d4 asked to run first");
7390 assert!(runner.park_here().expect("park_here"));
7391
7392 let resumed = Runner::resume(&run_id).expect("resume");
7393 assert_eq!(resumed.state.candidates.len(), 1);
7394 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7395 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7396 assert_eq!(resumed.state.status, runner.state.status);
7397 assert!(
7398 resumed.state.parked,
7399 "still parked until `execute` actually walks the graph again"
7400 );
7401 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7402 }
7403
7404 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7406 let mut q = ask::Question::new(
7407 run.to_owned(),
7408 "implement".to_owned(),
7409 "impl-A".to_owned(),
7410 "Which storage backend should the cache use?".to_owned(),
7411 String::new(),
7412 vec!["SQLite".to_owned(), "Redis".to_owned()],
7413 );
7414 store.put(&mut q).unwrap();
7415 q
7416 }
7417
7418 #[test]
7419 fn a_failed_runs_open_question_is_abandoned() {
7420 ask_test_home();
7421 let store = ask::Questions::open();
7422 let mut runner = runner_at(RunStatus::Failed);
7423 let run = runner.state.id.clone();
7424 let q = ask_open_question(&store, &run);
7425
7426 runner.settle_questions();
7427
7428 let back = store.get(&q.id).unwrap();
7429 assert!(
7430 !back.status.open(),
7431 "the seat that asked died with the run; nobody is left to read an answer"
7432 );
7433 assert!(
7434 back.detail.contains(&run) && back.detail.contains("failed"),
7435 "the reason names what the run became, not just that it is gone: {}",
7436 back.detail
7437 );
7438 }
7439
7440 #[test]
7441 fn a_merged_runs_open_question_is_abandoned_too() {
7442 ask_test_home();
7443 let store = ask::Questions::open();
7444 for status in [RunStatus::Merged, RunStatus::Ready] {
7447 let mut runner = runner_at(status);
7448 let run = runner.state.id.clone();
7449 let q = ask_open_question(&store, &run);
7450
7451 runner.settle_questions();
7452
7453 let back = store.get(&q.id).unwrap();
7454 assert!(
7455 !back.status.open(),
7456 "{status:?} run's question must not outlive the run"
7457 );
7458 }
7459 }
7460
7461 #[test]
7462 fn a_still_resumable_runs_open_question_is_left_alone() {
7463 ask_test_home();
7464 let store = ask::Questions::open();
7465 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7471 let mut runner = runner_at(status);
7472 let run = runner.state.id.clone();
7473 let q = ask_open_question(&store, &run);
7474
7475 runner.settle_questions();
7476
7477 let back = store.get(&q.id).unwrap();
7478 assert!(
7479 back.status.open(),
7480 "{status:?} is still alive; the question must still be waiting"
7481 );
7482 }
7483 }
7484
7485 #[test]
7486 fn settle_questions_never_touches_an_already_answered_question() {
7487 ask_test_home();
7488 let store = ask::Questions::open();
7489 let mut runner = runner_at(RunStatus::Failed);
7490 let run = runner.state.id.clone();
7491 let mut q = ask_open_question(&store, &run);
7492 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
7493 .unwrap();
7494 store.put(&mut q).unwrap();
7495
7496 runner.settle_questions();
7501 runner.settle_questions();
7502
7503 let back = store.get(&q.id).unwrap();
7504 assert_eq!(
7505 back.status,
7506 ask::QuestionStatus::Answered,
7507 "a real answer is a decision on record, never overwritten by a sweep"
7508 );
7509 }
7510
7511 #[tokio::test]
7522 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
7523 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
7524 let tmp = tempfile::tempdir().expect("tempdir");
7525 let repo = tmp.path().join("repo");
7526 std::fs::create_dir_all(&repo).unwrap();
7527 init_repo(&repo);
7528
7529 let mut config = Config::default();
7530 config.graph.worktree_root = Some(tmp.path().join("wt"));
7531
7532 let mut state = RunState::new(
7533 repo.clone(),
7534 "main".to_owned(),
7535 "deadbeef".to_owned(),
7536 "task".to_owned(),
7537 config,
7538 );
7539 let root = state.worktree_root();
7540 let wt_a = root.join("cand-A");
7541 let wt_b = root.join("cand-B");
7542 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
7543 .await
7544 .expect("worktree A");
7545 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
7546 .await
7547 .expect("worktree B");
7548
7549 state.candidates = vec![
7550 Candidate {
7551 index: 0,
7552 label: 'A',
7553 agent: "alpha".to_owned(),
7554 branch: "magi/x/A".to_owned(),
7555 worktree: wt_a.clone(),
7556 summary: String::new(),
7557 stat: String::new(),
7558 files: 0,
7559 commits: 0,
7560 empty: false,
7561 failed: None,
7562 verified_noop: None,
7563 duration_ms: 0,
7564 folded: false,
7565 },
7566 Candidate {
7567 index: 1,
7568 label: 'B',
7569 agent: "beta".to_owned(),
7570 branch: "magi/x/B".to_owned(),
7571 worktree: wt_b.clone(),
7572 summary: String::new(),
7573 stat: String::new(),
7574 files: 0,
7575 commits: 0,
7576 empty: false,
7577 failed: None,
7578 verified_noop: None,
7579 duration_ms: 0,
7580 folded: false,
7581 },
7582 ];
7583 state.tally = Some(Tally {
7584 first_choice: BTreeMap::from([('A', 1)]),
7585 borda: BTreeMap::new(),
7586 winner: 'A',
7587 rankings: 1,
7588 unanimous_initial: true,
7589 deliberated: false,
7590 changed_votes: 0,
7591 unanimous_final: true,
7592 tie_break: None,
7593 judges: 1,
7594 present: 1,
7595 quorum: 1,
7596 met_quorum: true,
7597 uncontested: None,
7598 });
7599 state.status = RunStatus::Ready;
7600
7601 fold_run(&mut state, false, &crate::run::home())
7602 .await
7603 .expect("fold_run");
7604
7605 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
7606 assert!(
7607 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7608 "the unmerged winner's branch survives"
7609 );
7610 assert!(
7611 !state.candidates[0].folded,
7612 "the winner is not marked folded"
7613 );
7614
7615 assert!(!wt_b.exists(), "the loser's worktree is removed");
7616 assert!(
7617 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
7618 "the loser's branch is removed"
7619 );
7620 assert!(state.candidates[1].folded, "the loser is marked folded");
7621 }
7622
7623 #[tokio::test]
7632 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
7633 let tmp = tempfile::tempdir().expect("tempdir");
7634 let repo = tmp.path().join("repo");
7635 std::fs::create_dir_all(&repo).unwrap();
7636 init_repo(&repo);
7637
7638 let mut config = Config::default();
7639 config.merge.mode = MergeMode::Local;
7640
7641 let mut state = RunState::new(
7642 repo.clone(),
7643 "main".to_owned(),
7644 "deadbeef".to_owned(),
7645 "task".to_owned(),
7646 config,
7647 );
7648 state.candidates = vec![Candidate {
7649 index: 0,
7650 label: 'A',
7651 agent: "alpha".to_owned(),
7652 branch: "does-not-exist".to_owned(),
7653 worktree: repo.clone(),
7654 summary: String::new(),
7655 stat: String::new(),
7656 files: 0,
7657 commits: 0,
7658 empty: false,
7659 failed: None,
7660 verified_noop: None,
7661 duration_ms: 0,
7662 folded: false,
7663 }];
7664 state.tally = Some(Tally {
7665 first_choice: BTreeMap::from([('A', 1)]),
7666 borda: BTreeMap::new(),
7667 winner: 'A',
7668 rankings: 1,
7669 unanimous_initial: true,
7670 deliberated: false,
7671 changed_votes: 0,
7672 unanimous_final: true,
7673 tie_break: None,
7674 judges: 0,
7675 present: 0,
7676 quorum: 0,
7677 met_quorum: true,
7678 uncontested: Some("only candidate A produced a change".to_owned()),
7679 });
7680 state.reviews = vec![ReviewRound {
7681 round: 1,
7682 head: "deadbeef".to_owned(),
7683 verified_head: None,
7684 verified_at: None,
7685 reviews: Vec::new(),
7686 e2e: Vec::new(),
7687 fix: None,
7688 blocking: 0,
7689 answered: 0,
7690 expected: 0,
7691 clean: true,
7692 verify_retried: false,
7693 e2e_deferred: false,
7694 e2e_defer_reason: None,
7695 progressed: false,
7696 vote_split: false,
7697 reconsideration: Vec::new(),
7698 verdict: None,
7699 }];
7700 state.gate = vec![CommandOutcome {
7701 command: "test".to_owned(),
7702 code: Some(0),
7703 output_tail: String::new(),
7704 duration_ms: 0,
7705 resource_blocked: false,
7706 }];
7707 state.gate_ran = true;
7708 state.status = RunStatus::Ready;
7713 state.merge = Some(MergeOutcome {
7714 mode: MergeMode::Local,
7715 ok: false,
7716 detail: "already concluded".to_owned(),
7717 });
7718
7719 let mut runner = Runner {
7720 state,
7721 roles: ResolvedRoles {
7722 implementers: Vec::new(),
7723 judges: Vec::new(),
7724 reviewers: Vec::new(),
7725 fixer: None,
7726 conductor: conductor(),
7727 implementer_roster: Vec::new(),
7728 },
7729 sem: Arc::new(Semaphore::new(1)),
7730 pause: Pause::new(),
7731 interrupt: Pause::new(),
7732 };
7733
7734 runner.merge().await.expect("merge");
7735
7736 assert_eq!(
7737 runner.state.status,
7738 RunStatus::Ready,
7739 "a concluded run's status must not change on reentry"
7740 );
7741 assert_eq!(
7742 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
7743 Some("already concluded"),
7744 "merge must not run again once the node already recorded an outcome"
7745 );
7746 }
7747
7748 #[tokio::test]
7757 async fn merge_refuses_a_gate_that_has_not_actually_run() {
7758 let tmp = tempfile::tempdir().expect("tempdir");
7759 let repo = tmp.path().join("repo");
7760 std::fs::create_dir_all(&repo).unwrap();
7761 init_repo(&repo);
7762
7763 let mut config = Config::default();
7764 config.merge.mode = MergeMode::Local;
7765
7766 let mut state = RunState::new(
7767 repo.clone(),
7768 "main".to_owned(),
7769 "deadbeef".to_owned(),
7770 "task".to_owned(),
7771 config,
7772 );
7773 state.candidates = vec![Candidate {
7774 index: 0,
7775 label: 'A',
7776 agent: "alpha".to_owned(),
7777 branch: "does-not-exist".to_owned(),
7778 worktree: repo.clone(),
7779 summary: String::new(),
7780 stat: String::new(),
7781 files: 0,
7782 commits: 0,
7783 empty: false,
7784 failed: None,
7785 verified_noop: None,
7786 duration_ms: 0,
7787 folded: false,
7788 }];
7789 state.tally = Some(Tally {
7790 first_choice: BTreeMap::from([('A', 1)]),
7791 borda: BTreeMap::new(),
7792 winner: 'A',
7793 rankings: 1,
7794 unanimous_initial: true,
7795 deliberated: false,
7796 changed_votes: 0,
7797 unanimous_final: true,
7798 tie_break: None,
7799 judges: 0,
7800 present: 0,
7801 quorum: 0,
7802 met_quorum: true,
7803 uncontested: Some("only candidate A produced a change".to_owned()),
7804 });
7805 state.reviews = vec![ReviewRound {
7806 round: 1,
7807 head: "deadbeef".to_owned(),
7808 verified_head: None,
7809 verified_at: None,
7810 reviews: Vec::new(),
7811 e2e: Vec::new(),
7812 fix: None,
7813 blocking: 0,
7814 answered: 0,
7815 expected: 0,
7816 clean: true,
7817 verify_retried: false,
7818 e2e_deferred: false,
7819 e2e_defer_reason: None,
7820 progressed: false,
7821 vote_split: false,
7822 reconsideration: Vec::new(),
7823 verdict: None,
7824 }];
7825 state.gate = Vec::new();
7827 state.gate_ran = false;
7828 state.status = RunStatus::Gating;
7829
7830 let mut runner = Runner {
7831 state,
7832 roles: ResolvedRoles {
7833 implementers: Vec::new(),
7834 judges: Vec::new(),
7835 reviewers: Vec::new(),
7836 fixer: None,
7837 conductor: conductor(),
7838 implementer_roster: Vec::new(),
7839 },
7840 sem: Arc::new(Semaphore::new(1)),
7841 pause: Pause::new(),
7842 interrupt: Pause::new(),
7843 };
7844
7845 runner.merge().await.expect("merge");
7846
7847 assert!(
7848 runner.state.merge.is_none(),
7849 "an empty gate must never be read as a passing one: {:?}",
7850 runner.state.merge
7851 );
7852 }
7853
7854 #[tokio::test]
7861 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
7862 let tmp = tempfile::tempdir().expect("tempdir");
7863 let repo = tmp.path().join("repo");
7864 std::fs::create_dir_all(&repo).unwrap();
7865 init_repo(&repo);
7866
7867 let config = Config::default();
7869
7870 let mut state = RunState::new(
7871 repo.clone(),
7872 "main".to_owned(),
7873 "deadbeef".to_owned(),
7874 "task".to_owned(),
7875 config,
7876 );
7877 state.candidates = vec![Candidate {
7878 index: 0,
7879 label: 'A',
7880 agent: "alpha".to_owned(),
7881 branch: "does-not-exist".to_owned(),
7882 worktree: repo.clone(),
7883 summary: String::new(),
7884 stat: String::new(),
7885 files: 0,
7886 commits: 0,
7887 empty: false,
7888 failed: None,
7889 verified_noop: None,
7890 duration_ms: 0,
7891 folded: false,
7892 }];
7893 state.tally = Some(Tally {
7894 first_choice: BTreeMap::from([('A', 1)]),
7895 borda: BTreeMap::new(),
7896 winner: 'A',
7897 rankings: 1,
7898 unanimous_initial: true,
7899 deliberated: false,
7900 changed_votes: 0,
7901 unanimous_final: true,
7902 tie_break: None,
7903 judges: 0,
7904 present: 0,
7905 quorum: 0,
7906 met_quorum: true,
7907 uncontested: Some("only candidate A produced a change".to_owned()),
7908 });
7909 state.reviews = vec![ReviewRound {
7910 round: 1,
7911 head: "deadbeef".to_owned(),
7912 verified_head: None,
7913 verified_at: None,
7914 reviews: Vec::new(),
7915 e2e: Vec::new(),
7916 fix: None,
7917 blocking: 0,
7918 answered: 0,
7919 expected: 0,
7920 clean: true,
7921 verify_retried: false,
7922 e2e_deferred: false,
7923 e2e_defer_reason: None,
7924 progressed: false,
7925 vote_split: false,
7926 reconsideration: Vec::new(),
7927 verdict: None,
7928 }];
7929
7930 let mut runner = Runner {
7931 state,
7932 roles: ResolvedRoles {
7933 implementers: Vec::new(),
7934 judges: Vec::new(),
7935 reviewers: Vec::new(),
7936 fixer: None,
7937 conductor: conductor(),
7938 implementer_roster: Vec::new(),
7939 },
7940 sem: Arc::new(Semaphore::new(1)),
7941 pause: Pause::new(),
7942 interrupt: Pause::new(),
7943 };
7944
7945 runner.gate().await.expect("gate");
7946 assert!(
7947 runner.state.gate_ran,
7948 "zero configured commands is still a real attempt, not an unrun gate"
7949 );
7950 assert!(runner.state.gate.is_empty());
7951 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
7952 assert_ne!(
7953 runner.state.status,
7954 RunStatus::Blocked,
7955 "a gate with nothing to check must not read as failed"
7956 );
7957
7958 runner.merge().await.expect("merge");
7959 assert_eq!(
7960 runner.state.status,
7961 RunStatus::Ready,
7962 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
7963 );
7964 }
7965
7966 #[tokio::test]
7977 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
7978 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
7979 let home = crate::run::home();
7980
7981 let tmp = tempfile::tempdir().expect("tempdir");
7982 let repo = tmp.path().join("repo");
7983 std::fs::create_dir_all(&repo).unwrap();
7984 init_repo(&repo);
7985 let cache_dir = tmp.path().join("target");
7988
7989 let mut config = Config::default();
7990 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
7991 config.graph.timeout_verify = Some(2);
7994
7995 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
7996 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
7997 .expect("no io error acquiring directly")
7998 {
7999 crate::cache::AcquireOutcome::Acquired(g) => g,
8000 crate::cache::AcquireOutcome::Busy(b) => {
8001 panic!("expected the direct acquire to win the lease first: {b:?}")
8002 }
8003 };
8004
8005 let mut state = RunState::new(
8006 repo.clone(),
8007 "main".to_owned(),
8008 "deadbeef".to_owned(),
8009 "task".to_owned(),
8010 config,
8011 );
8012 state.candidates = vec![Candidate {
8013 index: 0,
8014 label: 'A',
8015 agent: "alpha".to_owned(),
8016 branch: "does-not-exist".to_owned(),
8017 worktree: repo.clone(),
8018 summary: String::new(),
8019 stat: String::new(),
8020 files: 0,
8021 commits: 0,
8022 empty: false,
8023 failed: None,
8024 verified_noop: None,
8025 duration_ms: 0,
8026 folded: false,
8027 }];
8028 state.tally = Some(Tally {
8029 first_choice: BTreeMap::from([('A', 1)]),
8030 borda: BTreeMap::new(),
8031 winner: 'A',
8032 rankings: 1,
8033 unanimous_initial: true,
8034 deliberated: false,
8035 changed_votes: 0,
8036 unanimous_final: true,
8037 tie_break: None,
8038 judges: 0,
8039 present: 0,
8040 quorum: 0,
8041 met_quorum: true,
8042 uncontested: Some("only candidate A produced a change".to_owned()),
8043 });
8044 state.reviews = vec![ReviewRound {
8045 round: 1,
8046 head: "deadbeef".to_owned(),
8047 verified_head: None,
8048 verified_at: None,
8049 reviews: Vec::new(),
8050 e2e: Vec::new(),
8051 fix: None,
8052 blocking: 0,
8053 answered: 0,
8054 expected: 0,
8055 clean: true,
8056 verify_retried: false,
8057 e2e_deferred: false,
8058 e2e_defer_reason: None,
8059 progressed: false,
8060 vote_split: false,
8061 reconsideration: Vec::new(),
8062 verdict: None,
8063 }];
8064
8065 let mut runner = Runner {
8066 state,
8067 roles: ResolvedRoles {
8068 implementers: Vec::new(),
8069 judges: Vec::new(),
8070 reviewers: Vec::new(),
8071 fixer: None,
8072 conductor: conductor(),
8073 implementer_roster: Vec::new(),
8074 },
8075 sem: Arc::new(Semaphore::new(1)),
8076 pause: Pause::new(),
8077 interrupt: Pause::new(),
8078 };
8079
8080 let started = std::time::Instant::now();
8081 runner.gate().await.expect("gate");
8082 assert!(
8083 started.elapsed() < Duration::from_secs(1),
8084 "a gate with nothing to run must never wait on a lease it never needed"
8085 );
8086 assert!(
8087 runner.state.gate_ran,
8088 "zero commands is still a real, immediate attempt"
8089 );
8090 assert!(runner.state.gate.is_empty());
8091 assert_ne!(
8092 runner.state.status,
8093 RunStatus::Blocked,
8094 "must not read as resource-blocked on a lease it never asked for"
8095 );
8096 }
8097
8098 #[tokio::test]
8108 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8109 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8110
8111 let tmp = tempfile::tempdir().expect("tempdir");
8112 let repo = tmp.path().join("repo");
8113 std::fs::create_dir_all(&repo).unwrap();
8114 init_repo(&repo);
8115
8116 let mut config = Config::default();
8117 config.verify.gate = vec![
8118 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8119 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8120 .to_owned(),
8121 ];
8122
8123 let mut state = RunState::new(
8124 repo.clone(),
8125 "main".to_owned(),
8126 "deadbeef".to_owned(),
8127 "task".to_owned(),
8128 config,
8129 );
8130 let run_id = state.id.clone();
8131 state.candidates = vec![Candidate {
8132 index: 0,
8133 label: 'A',
8134 agent: "alpha".to_owned(),
8135 branch: "does-not-exist".to_owned(),
8136 worktree: repo.clone(),
8137 summary: String::new(),
8138 stat: String::new(),
8139 files: 0,
8140 commits: 0,
8141 empty: false,
8142 failed: None,
8143 verified_noop: None,
8144 duration_ms: 0,
8145 folded: false,
8146 }];
8147 state.tally = Some(Tally {
8148 first_choice: BTreeMap::from([('A', 1)]),
8149 borda: BTreeMap::new(),
8150 winner: 'A',
8151 rankings: 1,
8152 unanimous_initial: true,
8153 deliberated: false,
8154 changed_votes: 0,
8155 unanimous_final: true,
8156 tie_break: None,
8157 judges: 0,
8158 present: 0,
8159 quorum: 0,
8160 met_quorum: true,
8161 uncontested: Some("only candidate A produced a change".to_owned()),
8162 });
8163 state.reviews = vec![ReviewRound {
8164 round: 1,
8165 head: "deadbeef".to_owned(),
8166 verified_head: None,
8167 verified_at: None,
8168 reviews: Vec::new(),
8169 e2e: Vec::new(),
8170 fix: None,
8171 blocking: 0,
8172 answered: 0,
8173 expected: 0,
8174 clean: true,
8175 verify_retried: false,
8176 e2e_deferred: false,
8177 e2e_defer_reason: None,
8178 progressed: false,
8179 vote_split: false,
8180 reconsideration: Vec::new(),
8181 verdict: None,
8182 }];
8183
8184 let mut runner = Runner {
8185 state,
8186 roles: ResolvedRoles {
8187 implementers: Vec::new(),
8188 judges: Vec::new(),
8189 reviewers: Vec::new(),
8190 fixer: None,
8191 conductor: conductor(),
8192 implementer_roster: Vec::new(),
8193 },
8194 sem: Arc::new(Semaphore::new(1)),
8195 pause: Pause::new(),
8196 interrupt: Pause::new(),
8197 };
8198
8199 let started_marker = repo.join("started.marker");
8200 let release_marker = repo.join("release.marker");
8201 let poller = tokio::spawn(async move {
8202 for _ in 0..100 {
8207 if started_marker.exists()
8208 && let Ok(s) = crate::run::RunState::load(&run_id)
8209 && let Some(a) = s.active.get("gate")
8210 {
8211 std::fs::write(&release_marker, b"go").expect("release marker");
8212 return Some(a.clone());
8213 }
8214 tokio::time::sleep(Duration::from_millis(50)).await;
8215 }
8216 None
8217 });
8218
8219 runner.gate().await.expect("gate");
8220 let captured = poller.await.expect("poller task");
8221 let captured = captured.expect(
8222 "the poller never saw a `gate` task entry in run.json while the command was \
8223 still blocked on its own release marker",
8224 );
8225
8226 assert_eq!(captured.task.as_deref(), Some("gate"));
8227 assert_eq!(captured.node, "gate");
8228 assert_eq!(captured.index, Some(1));
8229 assert_eq!(captured.total, Some(1));
8230 assert!(
8231 captured
8232 .command
8233 .as_deref()
8234 .is_some_and(|c| c.contains("started.marker")),
8235 "{captured:?}"
8236 );
8237
8238 assert!(
8239 runner.state.active.is_empty(),
8240 "the entry must be cleared once the command actually finished: {:?}",
8241 runner.state.active
8242 );
8243 assert!(runner.state.gate_ran);
8244 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8245 }
8246
8247 #[tokio::test]
8260 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8261 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8262 let home = crate::run::home();
8263
8264 let tmp = tempfile::tempdir().expect("tempdir");
8265 let repo = tmp.path().join("repo");
8266 std::fs::create_dir_all(&repo).unwrap();
8267 init_repo(&repo);
8268 let head = crate::git::rev_parse(&repo, "HEAD")
8269 .await
8270 .expect("rev-parse");
8271 let cache_dir = tmp.path().join("target");
8274
8275 let mut config = Config::default();
8276 config.verify.e2e = vec![format!(
8277 "CARGO_TARGET_DIR='{}' test -f README.md",
8278 cache_dir.display()
8279 )];
8280 config.graph.review_rounds = 1;
8281 config.graph.timeout_verify = Some(2);
8284
8285 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8286 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8287 .expect("no io error acquiring directly")
8288 {
8289 crate::cache::AcquireOutcome::Acquired(g) => g,
8290 crate::cache::AcquireOutcome::Busy(b) => {
8291 panic!("expected the direct acquire to win the lease first: {b:?}")
8292 }
8293 };
8294
8295 let mut state = RunState::new(
8296 repo.clone(),
8297 "main".to_owned(),
8298 head.clone(),
8299 "task".to_owned(),
8300 config,
8301 );
8302 state.candidates = vec![Candidate {
8303 index: 0,
8304 label: 'A',
8305 agent: "alpha".to_owned(),
8306 branch: "does-not-exist".to_owned(),
8307 worktree: repo.clone(),
8308 summary: String::new(),
8309 stat: String::new(),
8310 files: 0,
8311 commits: 0,
8312 empty: false,
8313 failed: None,
8314 verified_noop: None,
8315 duration_ms: 0,
8316 folded: false,
8317 }];
8318 state.tally = Some(Tally {
8319 first_choice: BTreeMap::from([('A', 1)]),
8320 borda: BTreeMap::new(),
8321 winner: 'A',
8322 rankings: 1,
8323 unanimous_initial: true,
8324 deliberated: false,
8325 changed_votes: 0,
8326 unanimous_final: true,
8327 tie_break: None,
8328 judges: 0,
8329 present: 0,
8330 quorum: 0,
8331 met_quorum: true,
8332 uncontested: Some("only candidate A produced a change".to_owned()),
8333 });
8334 state.reviews = vec![ReviewRound {
8338 round: 1,
8339 head: head.clone(),
8340 verified_head: None,
8341 verified_at: None,
8342 reviews: Vec::new(),
8343 e2e: Vec::new(),
8344 fix: None,
8345 blocking: 1,
8346 answered: 1,
8347 expected: 1,
8348 clean: false,
8349 verify_retried: false,
8350 e2e_deferred: true,
8351 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8352 progressed: false,
8353 vote_split: false,
8354 reconsideration: Vec::new(),
8355 verdict: None,
8356 }];
8357
8358 let mut runner = Runner {
8359 state,
8360 roles: ResolvedRoles {
8361 implementers: Vec::new(),
8362 judges: Vec::new(),
8363 reviewers: Vec::new(),
8364 fixer: None,
8365 conductor: conductor(),
8366 implementer_roster: Vec::new(),
8367 },
8368 sem: Arc::new(Semaphore::new(1)),
8369 pause: Pause::new(),
8370 interrupt: Pause::new(),
8371 };
8372
8373 let shell = runner.state.config.shell();
8374 runner
8375 .stop_reviewing("round budget spent", &shell, &repo)
8376 .await
8377 .expect("stop_reviewing");
8378
8379 let last = runner.state.reviews.last().expect("round record");
8380 assert_eq!(
8381 last.e2e_status(),
8382 E2eStatus::ResourceBlocked,
8383 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8384 failed: {last:?}"
8385 );
8386 assert_eq!(
8387 last.verified_head.as_deref(),
8388 Some(head.as_str()),
8389 "which commit this attempt targeted is known even though nothing finished checking \
8390 it"
8391 );
8392 let first_attempt_at = last
8393 .verified_at
8394 .expect("when this attempt ran is known too");
8395 assert_ne!(
8396 runner.state.status,
8397 RunStatus::Blocked,
8398 "contention is evidence about the machine, not the patch — it must not settle the \
8399 run as blocked: {:?}",
8400 runner.state.status
8401 );
8402 assert!(
8403 !runner
8404 .state
8405 .events
8406 .iter()
8407 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
8408 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
8409 runner.state.events
8410 );
8411
8412 runner
8417 .stop_reviewing("round budget spent", &shell, &repo)
8418 .await
8419 .expect("stop_reviewing retry");
8420 assert_eq!(
8421 runner.state.reviews.len(),
8422 1,
8423 "no new round was started: {:?}",
8424 runner.state.reviews
8425 );
8426 let last = runner.state.reviews.last().expect("round record");
8427 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
8428 assert!(
8429 last.verified_at.expect("still known") > first_attempt_at,
8430 "a second reentry must be a fresh attempt, not a stale copy of the first"
8431 );
8432 assert_ne!(runner.state.status, RunStatus::Blocked);
8433
8434 held.release();
8435 }
8436
8437 #[tokio::test]
8449 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
8450 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8451 let home = crate::run::home();
8452
8453 let tmp = tempfile::tempdir().expect("tempdir");
8454 let repo = tmp.path().join("repo");
8455 std::fs::create_dir_all(&repo).unwrap();
8456 init_repo(&repo);
8457 let head = crate::git::rev_parse(&repo, "HEAD")
8458 .await
8459 .expect("rev-parse");
8460 let cache_dir = tmp.path().join("target");
8461
8462 let mut config = Config::default();
8463 config.verify.e2e = vec![format!(
8464 "CARGO_TARGET_DIR='{}' test -f README.md",
8465 cache_dir.display()
8466 )];
8467 config.graph.review_rounds = 1;
8468 config.graph.timeout_verify = Some(2);
8469
8470 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8471 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8472 .expect("no io error acquiring directly")
8473 {
8474 crate::cache::AcquireOutcome::Acquired(g) => g,
8475 crate::cache::AcquireOutcome::Busy(b) => {
8476 panic!("expected the direct acquire to win the lease first: {b:?}")
8477 }
8478 };
8479
8480 let mut state = RunState::new(
8481 repo.clone(),
8482 "main".to_owned(),
8483 head.clone(),
8484 "task".to_owned(),
8485 config,
8486 );
8487 state.candidates = vec![Candidate {
8488 index: 0,
8489 label: 'A',
8490 agent: "alpha".to_owned(),
8491 branch: "does-not-exist".to_owned(),
8492 worktree: repo.clone(),
8493 summary: String::new(),
8494 stat: String::new(),
8495 files: 0,
8496 commits: 0,
8497 empty: false,
8498 failed: None,
8499 verified_noop: None,
8500 duration_ms: 0,
8501 folded: false,
8502 }];
8503 state.tally = Some(Tally {
8504 first_choice: BTreeMap::from([('A', 1)]),
8505 borda: BTreeMap::new(),
8506 winner: 'A',
8507 rankings: 1,
8508 unanimous_initial: true,
8509 deliberated: false,
8510 changed_votes: 0,
8511 unanimous_final: true,
8512 tie_break: None,
8513 judges: 0,
8514 present: 0,
8515 quorum: 0,
8516 met_quorum: true,
8517 uncontested: Some("only candidate A produced a change".to_owned()),
8518 });
8519 state.reviews = vec![ReviewRound {
8523 round: 1,
8524 head: head.clone(),
8525 verified_head: Some(head.clone()),
8526 verified_at: Some(jiff::Timestamp::now()),
8527 reviews: Vec::new(),
8528 e2e: vec![CommandOutcome {
8529 command: format!(
8530 "CARGO_TARGET_DIR='{}' test -f README.md",
8531 cache_dir.display()
8532 ),
8533 code: None,
8534 output_tail: "waiting for the shared build cache".to_owned(),
8535 duration_ms: 0,
8536 resource_blocked: true,
8537 }],
8538 fix: None,
8539 blocking: 1,
8540 answered: 1,
8541 expected: 1,
8542 clean: false,
8543 verify_retried: false,
8544 e2e_deferred: false,
8545 e2e_defer_reason: None,
8546 progressed: false,
8547 vote_split: false,
8548 reconsideration: Vec::new(),
8549 verdict: None,
8550 }];
8551
8552 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
8553 let mut runner = Runner {
8554 state,
8555 roles: ResolvedRoles {
8556 implementers: Vec::new(),
8557 judges: Vec::new(),
8558 reviewers: Vec::new(),
8559 fixer: None,
8560 conductor: conductor(),
8561 implementer_roster: Vec::new(),
8562 },
8563 sem: Arc::new(Semaphore::new(1)),
8564 pause: Pause::new(),
8565 interrupt: Pause::new(),
8566 };
8567
8568 runner.review_loop().await.expect("review_loop");
8573
8574 assert_eq!(
8575 runner.state.reviews.len(),
8576 1,
8577 "no new round was started on top of the unresolved one: {:?}",
8578 runner.state.reviews
8579 );
8580 let last = &runner.state.reviews[0];
8581 assert_eq!(
8582 last.e2e_status(),
8583 E2eStatus::ResourceBlocked,
8584 "still contended: {last:?}"
8585 );
8586 assert!(
8587 last.verified_at.expect("still known") > first_attempt_at,
8588 "review_loop must have actually retried the check, not left it exactly as found"
8589 );
8590 assert_ne!(
8591 runner.state.status,
8592 RunStatus::Blocked,
8593 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
8594 runner.state.status
8595 );
8596
8597 held.release();
8598 }
8599
8600 #[tokio::test]
8601 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
8602 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8603 let tmp = tempfile::tempdir().expect("tempdir");
8604 let repo = tmp.path().join("repo");
8605 std::fs::create_dir_all(&repo).unwrap();
8606 init_repo(&repo);
8607
8608 let mut config = Config::default();
8609 config.merge.mode = MergeMode::Pr;
8610 config.graph.land = true;
8611 config.graph.land_approval = false;
8612
8613 let mut state = RunState::new(
8614 repo.clone(),
8615 "main".to_owned(),
8616 "deadbeef".to_owned(),
8617 "task".to_owned(),
8618 config,
8619 );
8620 state.candidates = vec![Candidate {
8621 index: 0,
8622 label: 'A',
8623 agent: "alpha".to_owned(),
8624 branch: "does-not-exist".to_owned(),
8625 worktree: repo.clone(),
8626 summary: String::new(),
8627 stat: String::new(),
8628 files: 0,
8629 commits: 0,
8630 empty: false,
8631 failed: None,
8632 verified_noop: None,
8633 duration_ms: 0,
8634 folded: false,
8635 }];
8636 state.tally = Some(Tally {
8637 first_choice: BTreeMap::from([('A', 1)]),
8638 borda: BTreeMap::new(),
8639 winner: 'A',
8640 rankings: 1,
8641 unanimous_initial: true,
8642 deliberated: false,
8643 changed_votes: 0,
8644 unanimous_final: true,
8645 tie_break: None,
8646 judges: 0,
8647 present: 0,
8648 quorum: 0,
8649 met_quorum: true,
8650 uncontested: Some("only candidate A produced a change".to_owned()),
8651 });
8652 state.reviews = vec![ReviewRound {
8653 round: 1,
8654 head: "deadbeef".to_owned(),
8655 verified_head: None,
8656 verified_at: None,
8657 reviews: Vec::new(),
8658 e2e: Vec::new(),
8659 fix: None,
8660 blocking: 0,
8661 answered: 0,
8662 expected: 0,
8663 clean: true,
8664 verify_retried: false,
8665 e2e_deferred: false,
8666 e2e_defer_reason: None,
8667 progressed: false,
8668 vote_split: false,
8669 reconsideration: Vec::new(),
8670 verdict: None,
8671 }];
8672 state.gate = vec![CommandOutcome {
8673 command: "test".to_owned(),
8674 code: Some(0),
8675 output_tail: String::new(),
8676 duration_ms: 0,
8677 resource_blocked: false,
8678 }];
8679 state.gate_ran = true;
8680 state.status = RunStatus::Landing;
8684 state.merge = Some(MergeOutcome {
8685 mode: MergeMode::Pr,
8686 ok: true,
8687 detail: "https://example.invalid/x/y/pull/1".to_owned(),
8688 });
8689
8690 ask_test_home();
8694 let store = ask::Questions::open();
8695 let q = ask_open_question(&store, &state.id);
8696
8697 let mut runner = Runner {
8698 state,
8699 roles: ResolvedRoles {
8700 implementers: Vec::new(),
8701 judges: Vec::new(),
8702 reviewers: Vec::new(),
8703 fixer: None,
8704 conductor: conductor(),
8705 implementer_roster: Vec::new(),
8706 },
8707 sem: Arc::new(Semaphore::new(1)),
8708 pause: Pause::new(),
8709 interrupt: Pause::new(),
8710 };
8711
8712 runner.execute().await.expect("execute");
8717
8718 assert_eq!(
8719 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8720 Some("https://example.invalid/x/y/pull/1"),
8721 "reentry must not push again or open a second pull request over the \
8722 one `land` is already watching"
8723 );
8724 assert_ne!(
8725 runner.state.status,
8726 RunStatus::Landing,
8727 "land could not actually reach the fake pull request, so it must \
8728 have given up rather than left the run silently parked forever"
8729 );
8730 assert_eq!(runner.state.status, RunStatus::Blocked);
8734 assert!(
8735 store.get(&q.id).unwrap().status.open(),
8736 "Blocked is still alive; settle_questions must have been a no-op here"
8737 );
8738 }
8739
8740 fn state_with_round(round: ReviewRound) -> RunState {
8741 let mut s = RunState::new(
8742 PathBuf::from("/repo"),
8743 "main".to_owned(),
8744 "abc1234".to_owned(),
8745 "add retries".to_owned(),
8746 Config::default(),
8747 );
8748 s.reviews = vec![round];
8749 s
8750 }
8751
8752 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
8753 crate::verdict::Finding {
8754 id: id.to_owned(),
8755 severity,
8756 file: None,
8757 line: None,
8758 title: title.to_owned(),
8759 detail: String::new(),
8760 }
8761 }
8762
8763 #[test]
8764 fn pr_body_names_open_findings_and_declined_ones() {
8765 let round = ReviewRound {
8766 round: 2,
8767 head: "deadbee".to_owned(),
8768 verified_head: None,
8769 verified_at: None,
8770 reviews: vec![ReviewRecord {
8771 attempts: 0,
8772 reviewer: 1,
8773 agent: "alpha".to_owned(),
8774 summary: String::new(),
8775 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
8776 vote: None,
8777 failed: None,
8778 duration_ms: 0,
8779 }],
8780 e2e: vec![CommandOutcome {
8781 command: "cargo test".to_owned(),
8782 code: Some(0),
8783 output_tail: String::new(),
8784 duration_ms: 0,
8785 resource_blocked: false,
8786 }],
8787 verify_retried: false,
8788 e2e_deferred: false,
8789 e2e_defer_reason: None,
8790 fix: Some(FixRecord {
8791 agent: "alpha".to_owned(),
8792 addressed: Vec::new(),
8793 rejected: vec![crate::verdict::Rejection {
8794 id: "R1-1-1".to_owned(),
8795 why: "not reachable from any caller".to_owned(),
8796 }],
8797 notes: String::new(),
8798 committed: true,
8799 failed: None,
8800 duration_ms: 0,
8801 continuation: None,
8802 }),
8803 blocking: 0,
8804 answered: 1,
8805 expected: 1,
8806 clean: false,
8807 progressed: true,
8808 vote_split: false,
8809 reconsideration: Vec::new(),
8810 verdict: None,
8811 };
8812 let state = state_with_round(round);
8813 let body = pr_message(&state, 'A').body;
8814
8815 assert!(body.contains("add retries"), "the task must still be there");
8816 assert!(body.contains("R2-1-1"), "{body}");
8817 assert!(body.contains("unused import"), "{body}");
8818 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
8819 assert!(
8820 body.contains("not reachable from any caller"),
8821 "the reason it was declined: {body}"
8822 );
8823 }
8824
8825 #[test]
8826 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
8827 let round = ReviewRound {
8828 round: 1,
8829 head: "deadbee".to_owned(),
8830 verified_head: None,
8831 verified_at: None,
8832 reviews: vec![ReviewRecord {
8833 attempts: 0,
8834 reviewer: 1,
8835 agent: "alpha".to_owned(),
8836 summary: String::new(),
8837 findings: Vec::new(),
8838 vote: None,
8839 failed: None,
8840 duration_ms: 0,
8841 }],
8842 e2e: Vec::new(),
8843 verify_retried: false,
8844 e2e_deferred: false,
8845 e2e_defer_reason: None,
8846 fix: None,
8847 blocking: 0,
8848 answered: 1,
8849 expected: 1,
8850 clean: true,
8851 progressed: false,
8852 vote_split: false,
8853 reconsideration: Vec::new(),
8854 verdict: None,
8855 };
8856 let state = state_with_round(round);
8857 let body = pr_message(&state, 'A').body;
8858 assert!(!body.contains("Open review findings"), "{body}");
8859 assert!(!body.contains("Declined"), "{body}");
8860 }
8861
8862 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
8863 let mut state = RunState::new(
8864 PathBuf::from("/repo"),
8865 "main".to_owned(),
8866 "abc1234".to_owned(),
8867 instruction.to_owned(),
8868 Config::default(),
8869 );
8870 state.candidates.push(Candidate {
8871 index: 0,
8872 label: 'A',
8873 agent: "alpha".to_owned(),
8874 branch: "magi/x/A".to_owned(),
8875 worktree: PathBuf::from("/wt"),
8876 summary: summary.to_owned(),
8877 stat: String::new(),
8878 files: 1,
8879 commits: 1,
8880 empty: false,
8881 failed: None,
8882 verified_noop: None,
8883 folded: false,
8884 duration_ms: 0,
8885 });
8886 state
8887 }
8888
8889 #[test]
8890 fn pr_message_describes_the_change_not_the_task() {
8891 let state = state_with_summary(
8892 "今回やってほしいこと: results projector を直す",
8893 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
8894 );
8895 let m = pr_message(&state, 'A');
8896 assert_eq!(m.title, "fix(web): batch the runs list reads");
8897 assert!(
8898 m.body.starts_with("## Summary\n\n- reads run.json once"),
8899 "{}",
8900 m.body
8901 );
8902 assert!(!m.body.contains("TITLE:"), "{}", m.body);
8903 let task_at = m.body.find("今回やってほしいこと").unwrap();
8904 let details_at = m.body.find("<details>").unwrap();
8905 assert!(
8906 details_at < task_at,
8907 "the task lives inside <details>: {}",
8908 m.body
8909 );
8910 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
8911 assert!(m.body.contains("magi:candidate-a"));
8912 }
8913
8914 #[test]
8915 fn pr_message_falls_back_to_the_task_without_a_title_line() {
8916 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
8917 let m = pr_message(&state, 'A');
8918 assert_eq!(m.title, "add retries");
8919 assert!(
8920 m.body.contains("## Summary\n\n- did some things"),
8921 "{}",
8922 m.body
8923 );
8924
8925 let none = RunState::new(
8926 PathBuf::from("/repo"),
8927 "main".to_owned(),
8928 "abc1234".to_owned(),
8929 "add retries".to_owned(),
8930 Config::default(),
8931 );
8932 let m = pr_message(&none, 'A');
8933 assert_eq!(m.title, "add retries");
8934 assert!(!m.body.contains("## Summary"), "{}", m.body);
8935 }
8936
8937 #[test]
8938 fn pr_message_refuses_the_candidate_commit_subject() {
8939 for bad in [
8940 "TITLE: magi: candidate A (uncommitted work)",
8941 "TITLE: chore: stuff (uncommitted work)",
8942 "TITLE: ",
8943 ] {
8944 let state = state_with_summary("add retries", bad);
8945 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
8946 }
8947 }
8948
8949 #[test]
8950 fn pr_message_bounds_a_very_long_task_and_title() {
8951 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
8952 let state = state_with_summary(&long, "- nothing");
8953 let m = pr_message(&state, 'A');
8954 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
8955 assert!(!m.title.contains('\n'));
8956
8957 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
8958 let m = pr_message(&state, 'A');
8959 assert!(m.title.starts_with("feat: "));
8960 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
8961 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
8962 }
8963
8964 #[test]
8965 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
8966 let mut state = state_with_summary(
8970 "add retries",
8971 "TITLE: fix(web): batch reads\n- reads run.json once",
8972 );
8973 state.config.graph.language = "ja".to_owned();
8974 let m = pr_message(&state, 'A');
8975 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
8976
8977 let task = "今回やってほしいこと: results projector を直す";
8980 let mut state = state_with_summary(task, "- no title line");
8981 state.config.graph.language = "ja".to_owned();
8982 let m = pr_message(&state, 'A');
8983 assert_eq!(
8984 m.title,
8985 format!("chore: land candidate A of run {}", state.id)
8986 );
8987 assert!(
8988 m.body.contains(&format!(
8989 "<summary>Original task</summary>\n\n{task}\n\n</details>"
8990 )),
8991 "{}",
8992 m.body
8993 );
8994 }
8995
8996 #[test]
8997 fn pr_message_scrubs_home_paths_and_addresses() {
8998 let state = state_with_summary(
8999 "fix it in /Users/someone/src/x",
9000 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9001 );
9002 let m = pr_message(&state, 'A');
9003 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9004 assert!(!m.body.contains(leak), "{}", m.body);
9005 }
9006 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9007 }
9008
9009 #[test]
9010 fn pr_message_survives_a_task_that_closes_details() {
9011 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9012 let m = pr_message(&state, 'A');
9013 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9014 }
9015
9016 #[test]
9017 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9018 let cmd = manual_merge_command(
9019 MergeStyle::Squash,
9020 Path::new("/repo"),
9021 "b",
9022 "fix: \"quoted\" $(x) `y`\n\nbody",
9023 );
9024 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9025 }
9026
9027 #[test]
9028 fn manual_merge_command_matches_the_configured_style() {
9029 let repo = Path::new("/repo");
9030 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9031
9032 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9033 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9034
9035 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9036 assert_eq!(
9037 squash,
9038 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9039 \"Merge magi run 0832 (candidate A)\""
9040 );
9041
9042 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9043 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9044 }
9045
9046 #[test]
9047 fn a_nudge_gets_a_quarter_of_the_budget() {
9048 assert_eq!(retry_budget(secs(1200), true), secs(300));
9050 assert_eq!(retry_budget(secs(3600), true), secs(900));
9051 }
9052
9053 #[test]
9054 fn a_resent_prompt_keeps_the_whole_budget() {
9055 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9058 assert_eq!(retry_budget(secs(60), false), secs(60));
9059 }
9060
9061 #[test]
9062 fn the_floor_never_exceeds_the_original_budget() {
9063 assert_eq!(retry_budget(secs(60), true), secs(60));
9067 assert_eq!(retry_budget(secs(480), true), secs(120));
9068 assert_eq!(retry_budget(secs(0), true), secs(0));
9069 }
9070
9071 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9072 agent::CommandEvidence {
9073 id: "item1".to_owned(),
9074 description: "cargo test".to_owned(),
9075 exit_code,
9076 result_summary: String::new(),
9077 source: "codex".to_owned(),
9078 }
9079 }
9080
9081 #[test]
9082 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9083 assert!(!has_unconfirmed_command(&[]));
9087 }
9088
9089 #[test]
9090 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9091 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9095 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9096 assert!(!has_unconfirmed_command(&[
9097 evidence(Some(0)),
9098 evidence(Some(101))
9099 ]));
9100 }
9101
9102 #[test]
9103 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9104 assert!(has_unconfirmed_command(&[
9105 evidence(Some(0)),
9106 evidence(None)
9107 ]));
9108 }
9109
9110 #[test]
9111 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9112 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9113 assert_eq!(
9114 verified_noop_claim(true, &[], text).as_deref(),
9115 Some("already fixed by b32cfc4, on main.")
9116 );
9117 }
9118
9119 #[test]
9120 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9121 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9124 assert!(verified_noop_claim(false, &[], text).is_none());
9125 }
9126
9127 #[test]
9128 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9129 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9130 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9131 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9133 }
9134
9135 #[test]
9136 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9137 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9138 }
9139
9140 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9143 runner.state.candidates = shape
9144 .iter()
9145 .enumerate()
9146 .map(|(i, &(empty, verified))| Candidate {
9147 index: i,
9148 label: (b'A' + i as u8) as char,
9149 agent: "sonnet".to_owned(),
9150 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9151 worktree: PathBuf::from(format!("/wt/{i}")),
9152 summary: String::new(),
9153 stat: String::new(),
9154 files: 0,
9155 commits: 0,
9156 empty,
9157 failed: None,
9158 verified_noop: verified.map(str::to_owned),
9159 duration_ms: 0,
9160 folded: false,
9161 })
9162 .collect();
9163 }
9164
9165 #[test]
9166 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9167 ask_test_home();
9168 let mut runner = runner_at(RunStatus::Implementing);
9169 set_candidates(
9170 &mut runner,
9171 &[
9172 (true, Some("already on main at b32cfc4")),
9173 (true, Some("same fix, see the existing test")),
9174 ],
9175 );
9176
9177 runner
9178 .after_implement()
9179 .expect("a verified no-op is not an error");
9180
9181 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9182 }
9183
9184 #[test]
9185 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9186 ask_test_home();
9187 let mut runner = runner_at(RunStatus::Implementing);
9188 set_candidates(
9192 &mut runner,
9193 &[(true, Some("already on main at b32cfc4")), (true, None)],
9194 );
9195
9196 let err = runner
9197 .after_implement()
9198 .expect_err("an unverified empty candidate must still fail the run");
9199
9200 assert!(
9201 err.to_string().contains("no candidate produced a change"),
9202 "{err}"
9203 );
9204 assert_eq!(runner.state.status, RunStatus::Failed);
9205 }
9206
9207 #[test]
9208 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9209 ask_test_home();
9210 let mut runner = runner_at(RunStatus::Implementing);
9211 set_candidates(&mut runner, &[(true, None), (true, None)]);
9212
9213 let err = runner
9214 .after_implement()
9215 .expect_err("no candidate declared anything; this is an ordinary failure");
9216
9217 assert!(
9218 err.to_string().contains("no candidate produced a change"),
9219 "{err}"
9220 );
9221 assert_eq!(runner.state.status, RunStatus::Failed);
9222 }
9223}