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<()> {
585 self.state.parked = false;
590 self.state.clear_active();
597 let pid = std::process::id();
612 self.state.driver_pid = Some(pid);
613 self.state.driver_started_at = crate::proc::process_started_at(pid);
614 self.state.save()?;
615 if self.state.status == RunStatus::Stalled {
628 if self.recover_stall().await? {
629 self.finish_after_tally().await?;
630 } else {
631 self.state.save()?;
633 }
634 return Ok(());
635 }
636 if self.state.status == RunStatus::Landing {
646 self.run_land().await?;
647 self.settle_questions();
652 return Ok(());
653 }
654 self.prep().await?;
655 if self.park_here()? {
656 return Ok(());
657 }
658 self.advise().await?;
659 if self.park_here()? {
660 return Ok(());
661 }
662 self.implement().await?;
663 if self.park_here()? {
664 return Ok(());
665 }
666 if self.state.status == RunStatus::VerifiedNoop {
670 return Ok(());
671 }
672 self.judge().await?;
673 if self.park_here()? {
674 return Ok(());
675 }
676 self.deliberate().await?;
677 if self.park_here()? {
678 return Ok(());
679 }
680 self.vote().await?;
681 if self.park_here()? {
682 return Ok(());
683 }
684 self.tally()?;
685 if self.state.status == RunStatus::Stalled {
690 self.state.save()?;
694 return Ok(());
695 }
696 self.finish_after_tally().await?;
697 Ok(())
698 }
699
700 fn park_here(&mut self) -> Result<bool> {
707 if !self.pause.parked() && !self.interrupt.parked() {
712 return Ok(false);
713 }
714 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
715 Some(reason) => format!(
716 "parked after `{}` ({reason}) — resume to carry on from here",
717 self.state.status.as_str()
718 ),
719 None => format!(
720 "parked after `{}` — resume to carry on from here",
721 self.state.status.as_str()
722 ),
723 };
724 self.state.event("park", why);
725 self.state.parked = true;
726 self.state.save()?;
727 Ok(true)
728 }
729
730 pub fn on_pause(&mut self, pause: Pause) {
732 self.pause = pause;
733 }
734
735 pub fn watch_interrupt(&mut self, pause: Pause) {
741 self.interrupt = pause;
742 }
743
744 fn settle_questions(&mut self) {
764 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
765 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
766 }
767 }
768
769 async fn finish_after_tally(&mut self) -> Result<()> {
772 self.fold_losers().await?;
773 self.sync_to_base().await?;
778 self.review_loop().await?;
779 self.sync_to_base().await?;
780 self.gate().await?;
781 self.merge().await?;
782 self.state.save()?;
783 Ok(())
784 }
785
786 async fn prep(&mut self) -> Result<()> {
789 if !self.state.candidates.is_empty() {
790 return Ok(());
791 }
792 self.state.status = RunStatus::Prep;
793 let repo = self.state.repo.clone();
794 let base = self.state.base_commit.clone();
795 let root = self.state.worktree_root();
796 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
797
798 let hooks_dir = self.state.dir().join("hooks");
801 if self.state.config.blind.commit_msg_hook {
802 std::fs::create_dir_all(&hooks_dir)
803 .with_context(|| format!("create {}", hooks_dir.display()))?;
804 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
805 let path = hooks_dir.join("commit-msg");
806 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
807 make_executable(&path)?;
808 git::acquire_worktree_config(&repo).await?;
816 self.state.enabled_worktree_config = true;
817 }
818
819 for (index, (spec, label)) in self
820 .roles
821 .implementers
822 .clone()
823 .into_iter()
824 .zip(labels)
825 .enumerate()
826 {
827 let branch = self.state.branch_for(label);
828 let worktree = root.join(format!("cand-{label}"));
829 git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
830 if self.state.config.blind.commit_msg_hook {
831 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
832 }
833 git::local_exclude(&worktree, "/.magi/").await?;
834 self.state.candidates.push(Candidate {
835 index,
836 label,
837 agent: spec.id.clone(),
838 branch,
839 worktree,
840 summary: String::new(),
841 stat: String::new(),
842 files: 0,
843 commits: 0,
844 empty: false,
845 failed: None,
846 verified_noop: None,
847 duration_ms: 0,
848 folded: false,
849 });
850 }
851
852 for j in 1..=self.roles.judges.len() {
853 let wt = root.join(format!("judge-{j}"));
854 if !wt.exists() {
855 git::worktree_add_detached(&repo, &wt, &base).await?;
856 }
857 }
858
859 if self.state.config.graph.advise {
867 for k in 1..=self.state.config.graph.advisors {
868 let wt = root.join(format!("advisor-{k}"));
869 if !wt.exists() {
870 git::worktree_add_detached(&repo, &wt, &base).await?;
871 }
872 }
873 }
874
875 let authors: Vec<&str> = self
880 .roles
881 .implementers
882 .iter()
883 .map(|a| a.id.as_str())
884 .collect();
885 let overlap: Vec<String> = self
886 .roles
887 .judges
888 .iter()
889 .enumerate()
890 .filter(|(_, j)| authors.contains(&j.id.as_str()))
891 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
892 .collect();
893 if !overlap.is_empty() {
894 let note = format!(
895 "{} also authored a candidate; blind, but the panel is less \
896 independent than {} distinct agents would be",
897 overlap.join(", "),
898 self.roles.judges.len()
899 );
900 self.state.event("prep", note);
901 }
902
903 self.state.event(
904 "prep",
905 format!(
906 "{} candidates, {} judges, base {} ({})",
907 self.state.candidates.len(),
908 self.roles.judges.len(),
909 &self.state.base_commit[..7.min(self.state.base_commit.len())],
910 self.state.base_branch
911 ),
912 );
913 self.state.status = RunStatus::Implementing;
914 self.state.save()?;
915 Ok(())
916 }
917
918 async fn advise(&mut self) -> Result<()> {
951 let implement_untouched = self
952 .state
953 .candidates
954 .iter()
955 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
956 if !self.state.config.graph.advise || self.state.advise_attempted {
957 return Ok(());
958 }
959 if !implement_untouched {
960 self.state.event(
961 "advise",
962 "skipping the design-deliberation stage: at least one \
963 candidate already shows implementation progress, so this \
964 run is past the point the stage exists to run before"
965 .to_owned(),
966 );
967 self.state.advise_attempted = true;
968 self.state.save()?;
969 return Ok(());
970 }
971 let run_id = self.state.id.clone();
972 let prompts = self.state.config.prompts.clone();
973 let instruction = self.state.instruction.clone();
974 let language = self.state.config.graph.language.clone();
975 let root = self.state.worktree_root();
976 let n = self.state.config.graph.advisors;
977 let where_recorded = self.state.dir().join("run.json");
978
979 let seats = match self.state.config.advisors() {
980 Ok(seats) if !seats.is_empty() => seats,
981 Ok(_) => {
982 self.state.event(
983 "advise",
984 format!(
985 "[graph] advisors is 0; skipping the design-deliberation \
986 stage and continuing without a synthesis brief (see {})",
987 where_recorded.display()
988 ),
989 );
990 self.state.advise_attempted = true;
991 self.state.save()?;
992 return Ok(());
993 }
994 Err(e) => {
995 self.state.event(
996 "advise",
997 format!(
998 "could not resolve advisor seats ({e:#}); continuing \
999 without a design-deliberation brief (see {})",
1000 where_recorded.display()
1001 ),
1002 );
1003 self.state.advise_attempted = true;
1004 self.state.save()?;
1005 return Ok(());
1006 }
1007 };
1008
1009 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1010 let artifacts = agent::artifacts_dir(&self.state.dir());
1011 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1012
1013 let mut jobs = Vec::new();
1014 for (i, spec) in seats.iter().cloned().enumerate() {
1015 let seat_key = format!("advisor-{}", i + 1);
1016 let seat = self.seat(&seat_key, &spec.id);
1017 jobs.push(SeatJob {
1018 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1019 spec,
1020 seat,
1021 cwd: worktrees[i % worktrees.len()].clone(),
1022 timeout,
1023 allow_write: false,
1024 sessions: false,
1025 artifacts: artifacts.clone(),
1026 stem: seat_key,
1027 });
1028 }
1029
1030 self.state.event(
1031 "advise",
1032 format!(
1033 "{} advisor seat(s) sketching a design in parallel",
1034 jobs.len()
1035 ),
1036 );
1037 let mut quota_losses = Vec::new();
1038 let cache = self.state.config.cache_dir();
1039 let ctx = WaveCtx {
1040 run: &run_id,
1041 node: "advise",
1042 prompts: &prompts,
1043 cache: cache.as_deref(),
1044 round: None,
1045 };
1046 let results = ask_json_wave::<Proposal>(
1047 jobs,
1048 Arc::clone(&self.sem),
1049 self.state.config.graph.retries,
1050 &ctx,
1051 &mut quota_losses,
1052 &mut self.state,
1053 &|p: &Proposal| p.validate(),
1054 )
1055 .await;
1056 self.state.quota.extend(quota_losses);
1057
1058 let mut records = Vec::with_capacity(results.len());
1059 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1060 let agent_id = seat.agent.clone();
1061 self.state.seats.insert(seat.key.clone(), seat);
1062 match res {
1063 Ok((proposal, out)) => {
1064 self.state
1065 .event("advise", format!("advisor-{} proposed a design", i + 1));
1066 records.push(advise::AdvisorRecord::proposed(
1067 i + 1,
1068 agent_id,
1069 proposal,
1070 out.duration_ms,
1071 ));
1072 }
1073 Err(e) => {
1074 self.state.event(
1075 "advise",
1076 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1077 );
1078 records.push(advise::AdvisorRecord::failed(
1079 i + 1,
1080 agent_id,
1081 e.to_string(),
1082 ));
1083 }
1084 }
1085 }
1086
1087 let mut advice = advise::Advice {
1088 records,
1089 synthesis: None,
1090 };
1091 if advice.proposals().is_empty() {
1092 self.state.event(
1093 "advise",
1094 "no advisor produced a usable proposal; continuing without a \
1095 synthesis brief"
1096 .to_owned(),
1097 );
1098 } else {
1099 match self
1100 .synthesize_brief(
1101 &advice,
1102 &instruction,
1103 &language,
1104 &worktrees[0],
1105 &artifacts,
1106 &run_id,
1107 &prompts,
1108 cache.as_deref(),
1109 )
1110 .await
1111 {
1112 Ok(Some(text)) => {
1113 self.state.event(
1114 "advise",
1115 "synthesized a design brief for the implementer".to_owned(),
1116 );
1117 advice.synthesis = Some(text);
1118 }
1119 Ok(None) => {
1120 self.state.event(
1121 "advise",
1122 "the synthesis seat produced nothing usable; continuing \
1123 without a design brief"
1124 .to_owned(),
1125 );
1126 }
1127 Err(e) => {
1128 self.state.event(
1129 "advise",
1130 format!("could not synthesize a design brief: {e:#}"),
1131 );
1132 }
1133 }
1134 }
1135 advise::apply_reflection(&mut advice);
1136
1137 self.state.advice = Some(advice);
1138 self.state.advise_attempted = true;
1139 self.state.save()?;
1140 Ok(())
1141 }
1142
1143 #[allow(clippy::too_many_arguments)]
1155 async fn synthesize_brief(
1156 &mut self,
1157 advice: &advise::Advice,
1158 instruction: &str,
1159 language: &str,
1160 cwd: &Path,
1161 artifacts: &Path,
1162 run_id: &str,
1163 prompts: &Prompts,
1164 cache: Option<&Path>,
1165 ) -> Result<Option<String>> {
1166 let want = self.state.config.roles.synthesizer.as_deref();
1167 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1168 let mut seat = self.seat("advise-synthesis", &spec.id);
1169 let proposals = advice.proposals();
1170 let mut prompt = prompt::with_overlay(
1171 prompt::synthesize_brief(instruction, &proposals, language),
1172 prompts.overlay("advise"),
1173 );
1174 if cache.is_some() {
1175 prompt.push('\n');
1180 prompt.push_str(&prompt::build_cache_note("advise", false));
1181 }
1182 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1183 let out = agent::invoke(
1184 &spec,
1185 &mut seat,
1186 &Invocation {
1187 cwd,
1188 prompt: &prompt,
1189 timeout,
1190 allow_write: false,
1191 sessions: false,
1192 artifacts,
1193 stem: "advise-synthesis",
1194 run: run_id,
1195 node: "advise",
1196 cache_dir: None,
1197 attachments: &[],
1198 },
1199 )
1200 .await?;
1201 self.state.seats.insert(seat.key.clone(), seat);
1202 if !out.usable() {
1203 return Ok(None);
1204 }
1205 let text =
1206 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1207 Ok((!text.trim().is_empty()).then_some(text))
1208 }
1209
1210 async fn implement(&mut self) -> Result<()> {
1213 let run_id = self.state.id.clone();
1218 let prompts = self.state.config.prompts.clone();
1219 let todo: Vec<usize> = self
1220 .state
1221 .candidates
1222 .iter()
1223 .enumerate()
1224 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1225 .map(|(i, _)| i)
1226 .collect();
1227 if todo.is_empty() {
1228 return self.after_implement();
1229 }
1230 self.state.status = RunStatus::Implementing;
1231
1232 let language = self.state.config.graph.language.clone();
1233 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1234 let sessions = self.state.config.graph.sessions;
1235 let artifacts = agent::artifacts_dir(&self.state.dir());
1236 let brief = self
1240 .state
1241 .advice
1242 .as_ref()
1243 .and_then(|a| a.synthesis.as_deref())
1244 .map(str::to_owned);
1245
1246 let mut jobs = Vec::new();
1247 for &i in &todo {
1248 let (index, label, worktree) = {
1249 let c = &self.state.candidates[i];
1250 (c.index, c.label, c.worktree.clone())
1251 };
1252 let spec = self.roles.implementers[index].clone();
1253 let seat_key = format!("impl-{label}");
1254 let seat = self.seat(&seat_key, &spec.id);
1255 let instruction = self.state.instruction.clone();
1256 jobs.push(SeatJob {
1257 spec,
1258 seat,
1259 prompt: prompt::implement(
1260 &instruction,
1261 &worktree.to_string_lossy(),
1262 &language,
1263 brief.as_deref(),
1264 ),
1265 cwd: worktree,
1266 timeout,
1267 allow_write: true,
1268 sessions,
1269 artifacts: artifacts.clone(),
1270 stem: format!("impl-{label}"),
1271 });
1272 }
1273
1274 self.state.event(
1275 "implement",
1276 format!("{} candidates in parallel", jobs.len()),
1277 );
1278 let mut sent = jobs.clone();
1284 let cache = self.state.config.cache_dir();
1285 let ctx = WaveCtx {
1286 run: &run_id,
1287 node: "implement",
1288 prompts: &prompts,
1289 cache: cache.as_deref(),
1290 round: None,
1291 };
1292 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1293 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1294 .await;
1295 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1296 .await;
1297 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1298 .await;
1299
1300 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1301 let seat_key = seat.key.clone();
1302 let agent = seat.agent.clone();
1311 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1312 self.state.seats.insert(seat.key.clone(), seat);
1313 let label = self.state.candidates[i].label;
1314 let worktree = self.state.candidates[i].worktree.clone();
1315 let base = self.state.base_commit.clone();
1316
1317 let (summary, duration, failed, verified_claim) = match out {
1318 AgentOutcome::Ok(o) => {
1319 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1320 let failed = (!o.usable()).then(|| {
1321 if o.timed_out {
1322 "agent timed out".to_owned()
1323 } else {
1324 format!("agent exited with {:?}", o.exit_code)
1325 }
1326 });
1327 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1328 (text, o.duration_ms, failed, verified_claim)
1329 }
1330 AgentOutcome::Dropped(o) => {
1336 let why = o
1337 .dropped
1338 .as_ref()
1339 .map(|d| d.why.as_str())
1340 .unwrap_or("the CLI ended the stream without delivering its answer");
1341 (
1342 String::new(),
1343 o.duration_ms,
1344 Some(format!("the CLI dropped the stream ({why})")),
1345 None,
1346 )
1347 }
1348 AgentOutcome::Quota(o) => {
1349 self.state.quota.push(QuotaLoss {
1350 seat: seat_key,
1351 node: "implement".to_owned(),
1352 at: Timestamp::now(),
1353 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1354 });
1355 (
1356 String::new(),
1357 o.duration_ms,
1358 Some("rate limited (quota); produced no change".to_owned()),
1359 None,
1360 )
1361 }
1362 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1363 };
1364
1365 let rescued = match git::rescue_commit(
1368 &worktree,
1369 &format!("magi: candidate {label} (uncommitted work)"),
1370 )
1371 .await
1372 {
1373 Ok(r) => {
1374 self.state.note_withheld("implement", &r.withheld);
1375 r.committed
1376 }
1377 Err(_) => false,
1378 };
1379 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1380 .await
1381 .unwrap_or(0);
1382 let patch = git::diff(&worktree, &base, "HEAD")
1383 .await
1384 .unwrap_or_default();
1385 let stat = git::diff_stat(&worktree, &base, "HEAD")
1386 .await
1387 .unwrap_or_default();
1388 let files = git::changed_files(&worktree, &base, "HEAD")
1389 .await
1390 .map(|f| f.len())
1391 .unwrap_or(0);
1392 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1393
1394 let c = &mut self.state.candidates[i];
1395 if !exhausted_the_fallback_chain {
1396 c.agent = agent;
1397 }
1398 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1399 c.stat = stat;
1400 c.files = files;
1401 c.commits = commits;
1402 c.duration_ms = duration;
1403 c.empty = commits == 0 || patch.trim().is_empty();
1404 c.failed = match failed {
1407 Some(_) if c.empty => failed,
1408 _ => None,
1409 };
1410 c.verified_noop = if c.empty { verified_claim } else { None };
1415 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1416 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1417 (None, true, Some(_), _) => {
1418 format!("candidate {label}: no change produced (agent-verified no-op)")
1419 }
1420 (None, true, None, _) => format!("candidate {label}: no change produced"),
1421 (None, false, _, true) => {
1422 format!(
1423 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1424 )
1425 }
1426 (None, false, _, false) => {
1427 format!("candidate {label}: {files} files, {commits} commits")
1428 }
1429 };
1430 self.state.event("implement", note);
1431 self.state.save()?;
1432 }
1433
1434 self.after_implement()
1435 }
1436
1437 async fn resume_undelivered(
1465 &mut self,
1466 results: &mut [(usize, SeatState, AgentOutcome)],
1467 sent: &[SeatJob],
1468 prompts: &Prompts,
1469 run_id: &str,
1470 ) {
1471 for (wi, seat, out) in results.iter_mut() {
1472 let Some(dropped) = (match &*out {
1473 AgentOutcome::Dropped(o) => o.dropped.clone(),
1474 _ => None,
1475 }) else {
1476 continue;
1477 };
1478 let Some(job) = sent.get(*wi) else { continue };
1479 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1481 self.state.event(
1482 "implement",
1483 format!(
1484 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1485 work is in the tree",
1486 seat.key, dropped.output_tokens, dropped.why
1487 ),
1488 );
1489 continue;
1490 }
1491 if !has_context(&job.spec, seat, job.sessions) {
1499 self.state.event(
1500 "implement",
1501 format!(
1502 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1503 is no session left to resume",
1504 seat.key, dropped.output_tokens, dropped.why
1505 ),
1506 );
1507 continue;
1508 }
1509 self.state.event(
1510 "implement",
1511 format!(
1512 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1513 conversation",
1514 seat.key, dropped.output_tokens, dropped.why
1515 ),
1516 );
1517 let mut retry = job.clone();
1518 retry.seat = seat.clone();
1519 retry.prompt = prompt::resume_after_drop(&dropped.why);
1520 retry.timeout = retry_budget(job.timeout, true);
1521 retry.stem = format!("{}-resume", job.stem);
1522 let cache = self.state.config.cache_dir();
1523 let ctx = WaveCtx {
1524 run: run_id,
1525 node: "implement",
1526 prompts,
1527 cache: cache.as_deref(),
1528 round: None,
1529 };
1530 let (resumed_seat, resumed) =
1531 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1532 *seat = resumed_seat;
1533 *out = resumed;
1534 }
1535 }
1536
1537 async fn resume_quota_losses(
1599 &mut self,
1600 results: &mut [(usize, SeatState, AgentOutcome)],
1601 sent: &mut [SeatJob],
1602 prompts: &Prompts,
1603 run_id: &str,
1604 ) {
1605 let instruction = self.state.instruction.clone();
1606 let language = self.state.config.graph.language.clone();
1607 let brief = self
1608 .state
1609 .advice
1610 .as_ref()
1611 .and_then(|a| a.synthesis.as_deref())
1612 .map(str::to_owned);
1613 for (wi, seat, out) in results.iter_mut() {
1614 let Some(job) = sent.get_mut(*wi) else {
1615 continue;
1616 };
1617 let start = self
1622 .roles
1623 .implementer_roster
1624 .iter()
1625 .position(|s| s.id == job.spec.id)
1626 .unwrap_or(0);
1627 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1628 let mut fallback_attempt = 0usize;
1629 while matches!(&*out, AgentOutcome::Quota(_)) {
1630 let Some(next) =
1631 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1632 .cloned()
1633 else {
1634 break;
1635 };
1636 tried.insert(next.id.clone());
1637 fallback_attempt += 1;
1638
1639 if let Ok(r) = git::rescue_commit(
1640 &job.cwd,
1641 &format!(
1642 "magi: candidate {} (uncommitted work before quota fallback)",
1643 seat.key
1644 ),
1645 )
1646 .await
1647 {
1648 self.state.note_withheld("implement", &r.withheld);
1649 }
1650
1651 self.state.event(
1652 "implement",
1653 format!(
1654 "{}: rate limited (quota) on {}; retrying with {}",
1655 seat.key, seat.agent, next.id
1656 ),
1657 );
1658
1659 let new_seat = self.seat(&seat.key, &next.id);
1660 job.spec = next.clone();
1668 let mut retry = job.clone();
1669 retry.seat = new_seat;
1670 retry.prompt = prompt::implement(
1671 &instruction,
1672 &job.cwd.to_string_lossy(),
1673 &language,
1674 brief.as_deref(),
1675 );
1676 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1677 let cache = self.state.config.cache_dir();
1678 let ctx = WaveCtx {
1679 run: run_id,
1680 node: "implement",
1681 prompts,
1682 cache: cache.as_deref(),
1683 round: None,
1684 };
1685 let (fallback_seat, fallback_out) = run_one(
1686 retry,
1687 Arc::clone(&self.sem),
1688 &ctx,
1689 &mut self.state,
1690 fallback_attempt,
1691 )
1692 .await;
1693 *seat = fallback_seat;
1694 *out = fallback_out;
1695 }
1696 }
1697 }
1698
1699 async fn resume_unconfirmed_commands(
1723 &mut self,
1724 results: &mut [(usize, SeatState, AgentOutcome)],
1725 sent: &[SeatJob],
1726 prompts: &Prompts,
1727 run_id: &str,
1728 ) {
1729 for (wi, seat, out) in results.iter_mut() {
1730 let AgentOutcome::Ok(o) = &*out else {
1731 continue;
1732 };
1733 if !has_unconfirmed_command(&o.commands) {
1734 continue;
1735 }
1736 let Some(job) = sent.get(*wi) else { continue };
1737 if !has_context(&job.spec, seat, job.sessions) {
1738 self.state.event(
1739 "implement",
1740 format!(
1741 "{}: the reply named a command whose own CLI never confirmed the exit \
1742 status of, but there is no session left to resume",
1743 seat.key
1744 ),
1745 );
1746 continue;
1747 }
1748 self.state.event(
1749 "implement",
1750 format!(
1751 "{}: the reply named a command whose own CLI never confirmed the exit \
1752 status of; resuming the conversation",
1753 seat.key
1754 ),
1755 );
1756 let mut retry = job.clone();
1757 retry.seat = seat.clone();
1758 retry.prompt = prompt::resume_incomplete(
1759 "a command in your last reply had no confirmed exit status",
1760 );
1761 retry.timeout = retry_budget(job.timeout, true);
1762 retry.stem = format!("{}-confirm", job.stem);
1763 let cache = self.state.config.cache_dir();
1764 let ctx = WaveCtx {
1765 run: run_id,
1766 node: "implement",
1767 prompts,
1768 cache: cache.as_deref(),
1769 round: None,
1770 };
1771 let (resumed_seat, resumed) =
1772 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1773 *seat = resumed_seat;
1774 *out = resumed;
1775 }
1776 }
1777
1778 async fn continue_fix_report(
1799 &mut self,
1800 mut seat: SeatState,
1801 parse_err: String,
1802 job: &SeatJob,
1803 prompts: &Prompts,
1804 run_id: &str,
1805 round: usize,
1806 ) -> (
1807 SeatState,
1808 Option<FixReport>,
1809 Option<String>,
1810 ContinuationRecord,
1811 ) {
1812 let mut last_err = parse_err;
1813 let mut cumulative_wait_ms = 0u64;
1814 let mut attempts = 0usize;
1815 loop {
1816 if !has_context(&job.spec, &seat, job.sessions) {
1817 self.state.event(
1818 "fix",
1819 format!(
1820 "round {round}: fixer's reply had no adoption report ({last_err}); no \
1821 session left to resume into"
1822 ),
1823 );
1824 let outcome = if attempts == 0 {
1825 ContinuationOutcome::NoSession
1826 } else {
1827 ContinuationOutcome::Exhausted
1828 };
1829 return (
1830 seat,
1831 None,
1832 Some(format!("unparsable fix report: {last_err}")),
1833 ContinuationRecord {
1834 attempts,
1835 cumulative_wait_ms,
1836 outcome,
1837 },
1838 );
1839 }
1840 if attempts >= MAX_FIX_CONTINUATIONS {
1841 self.state.event(
1842 "fix",
1843 format!(
1844 "round {round}: fixer's reply still had no adoption report after \
1845 {attempts} continuation(s) ({last_err}); giving up"
1846 ),
1847 );
1848 return (
1849 seat,
1850 None,
1851 Some(format!(
1852 "unparsable fix report after {attempts} continuation(s): {last_err}"
1853 )),
1854 ContinuationRecord {
1855 attempts,
1856 cumulative_wait_ms,
1857 outcome: ContinuationOutcome::Exhausted,
1858 },
1859 );
1860 }
1861 attempts += 1;
1862 self.state.event(
1863 "fix",
1864 format!(
1865 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
1866 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
1867 ),
1868 );
1869 let mut retry = job.clone();
1870 retry.seat = seat.clone();
1871 retry.prompt = prompt::resume_incomplete(&last_err);
1872 retry.timeout = retry_budget(job.timeout, true);
1873 retry.stem = format!("{}-continue{attempts}", job.stem);
1874 let cache = self.state.config.cache_dir();
1875 let ctx = WaveCtx {
1876 run: run_id,
1877 node: "fix",
1878 prompts,
1879 cache: cache.as_deref(),
1880 round: Some(round),
1881 };
1882 let (resumed_seat, resumed_out) = run_one(
1883 retry,
1884 Arc::clone(&self.sem),
1885 &ctx,
1886 &mut self.state,
1887 attempts,
1888 )
1889 .await;
1890 seat = resumed_seat;
1891 match resumed_out {
1892 AgentOutcome::Ok(o) => {
1893 cumulative_wait_ms += o.duration_ms;
1894 match verdict::extract_json::<FixReport>(&o.text) {
1895 Ok(report) if !has_unconfirmed_command(&o.commands) => {
1896 self.state.event(
1897 "fix",
1898 format!(
1899 "round {round}: fixer's adoption report recovered after \
1900 {attempts} continuation(s)"
1901 ),
1902 );
1903 return (
1904 seat,
1905 Some(report),
1906 None,
1907 ContinuationRecord {
1908 attempts,
1909 cumulative_wait_ms,
1910 outcome: ContinuationOutcome::Resumed,
1911 },
1912 );
1913 }
1914 Ok(_) => {
1922 last_err = "the reply parsed, but it reported a command whose own CLI \
1923 never confirmed an exit status"
1924 .to_owned();
1925 }
1926 Err(e) => last_err = e.to_string(),
1927 }
1928 }
1929 AgentOutcome::Quota(o) => {
1930 cumulative_wait_ms += o.duration_ms;
1931 self.state.quota.push(QuotaLoss {
1932 seat: seat.key.clone(),
1933 node: "fix".to_owned(),
1934 at: Timestamp::now(),
1935 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1936 });
1937 self.state.event(
1938 "fix",
1939 format!(
1940 "round {round}: continuation rate limited (quota); not retrying now"
1941 ),
1942 );
1943 return (
1944 seat,
1945 None,
1946 Some("rate limited (quota) while recovering the fix report".to_owned()),
1947 ContinuationRecord {
1948 attempts,
1949 cumulative_wait_ms,
1950 outcome: ContinuationOutcome::QuotaLost,
1951 },
1952 );
1953 }
1954 AgentOutcome::Dropped(o) => {
1955 cumulative_wait_ms += o.duration_ms;
1956 let why = o
1957 .dropped
1958 .as_ref()
1959 .map(|d| d.why.as_str())
1960 .unwrap_or("the CLI ended the stream without delivering its answer");
1961 last_err = format!("the CLI dropped the stream ({why})");
1962 }
1963 AgentOutcome::Failed(e) => last_err = e,
1964 }
1965 }
1966 }
1967
1968 fn after_implement(&mut self) -> Result<()> {
1969 if self.state.leaks.is_empty() {
1971 let cfg = self.state.config.blind.clone();
1972 let mut leaks = Vec::new();
1973 for c in &self.state.candidates {
1974 let Some(patch) =
1975 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
1976 else {
1977 continue;
1978 };
1979 leaks.extend(blind::scan(
1980 &format!("candidate {} patch", c.label),
1981 &patch,
1982 &cfg.vendor_tokens,
1983 ));
1984 }
1985 if !leaks.is_empty() {
1986 let summary = leaks
1987 .iter()
1988 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
1989 .collect::<Vec<_>>()
1990 .join(", ");
1991 match cfg.on_leak {
1992 LeakPolicy::Fail => {
1993 self.state.status = RunStatus::Failed;
1994 self.state
1995 .event("blind", format!("vendor text in a patch: {summary}"));
1996 self.state.leaks = leaks;
1997 self.state.save()?;
1998 self.settle_questions();
1999 bail!(
2000 "blind.on_leak = \"fail\" and vendor text reached a \
2001 judged patch: {summary}"
2002 );
2003 }
2004 LeakPolicy::Redact => self.state.event(
2005 "blind",
2006 format!("redacting vendor text for judging: {summary}"),
2007 ),
2008 LeakPolicy::Warn => self.state.event(
2009 "blind",
2010 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2011 ),
2012 }
2013 self.state.leaks = leaks;
2014 }
2015 }
2016
2017 if self.state.viable().is_empty() {
2018 if self.state.all_candidates_verified_noop() {
2019 self.state.status = RunStatus::VerifiedNoop;
2030 self.state.save()?;
2031 self.settle_questions();
2032 return Ok(());
2033 }
2034 self.state.status = RunStatus::Failed;
2035 self.state.save()?;
2036 self.settle_questions();
2037 bail!("no candidate produced a change; nothing to judge");
2038 }
2039 self.state.status = RunStatus::Judging;
2040 self.state.save()?;
2041 Ok(())
2042 }
2043
2044 async fn judge(&mut self) -> Result<()> {
2047 let run_id = self.state.id.clone();
2052 let prompts = self.state.config.prompts.clone();
2053 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2054 return Ok(());
2055 }
2056 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2057 if viable.len() == 1 {
2058 self.state.judge_skipped = true;
2065 self.state.event(
2066 "judge",
2067 format!(
2068 "only candidate {} produced a change; judging skipped",
2069 viable[0].label
2070 ),
2071 );
2072 self.state.save()?;
2073 return Ok(());
2074 }
2075 self.state.status = RunStatus::Judging;
2076
2077 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2078 let language = self.state.config.graph.language.clone();
2079 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2080 let sessions = self.state.config.graph.sessions;
2081 let artifacts = agent::artifacts_dir(&self.state.dir());
2082 let root = self.state.worktree_root();
2083 let base_short = short(&self.state.base_commit);
2084
2085 let mut jobs = Vec::new();
2086 let mut orders = Vec::new();
2087 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2088 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2089 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2090 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2091 let seat_key = format!("judge-{}", j + 1);
2092 let seat = self.seat(&seat_key, &spec.id);
2093 jobs.push(SeatJob {
2094 prompt: prompt::judge(
2095 &self.state.instruction,
2096 &views,
2097 self.roles.judges.len(),
2098 &base_short,
2099 &language,
2100 ),
2101 spec,
2102 seat,
2103 cwd: root.join(format!("judge-{}", j + 1)),
2104 timeout,
2105 allow_write: false,
2106 sessions,
2107 artifacts: artifacts.clone(),
2108 stem: format!("judge-{}", j + 1),
2109 });
2110 }
2111
2112 self.state.event(
2113 "judge",
2114 format!(
2115 "{} judges ranking {} candidates blind",
2116 jobs.len(),
2117 viable.len()
2118 ),
2119 );
2120 let labels_for_check = labels.clone();
2121 let mut quota_losses = Vec::new();
2122 let cache = self.state.config.cache_dir();
2123 let ctx = WaveCtx {
2124 run: &run_id,
2125 node: "judge",
2126 prompts: &prompts,
2127 cache: cache.as_deref(),
2128 round: None,
2129 };
2130 let results = ask_json_wave::<Ranking>(
2131 jobs,
2132 Arc::clone(&self.sem),
2133 self.state.config.graph.retries,
2134 &ctx,
2135 &mut quota_losses,
2136 &mut self.state,
2137 &move |r: &Ranking| r.validate(&labels_for_check),
2138 )
2139 .await;
2140 self.state.quota.extend(quota_losses);
2141
2142 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2143 let agent_id = seat.agent.clone();
2144 self.state.seats.insert(seat.key.clone(), seat);
2145 let mut record = Judgement {
2146 judge: j + 1,
2147 seat: format!("judge-{}", j + 1),
2148 agent: agent_id,
2149 ranking: Vec::new(),
2150 reasons: BTreeMap::new(),
2151 confidence: None,
2152 order: orders[j].clone(),
2153 failed: None,
2154 duration_ms: 0,
2155 };
2156 match res {
2157 Ok((ranking, out)) => {
2158 record.ranking = ranking.normalized();
2159 record.reasons = ranking.reasons;
2160 record.confidence = ranking.confidence;
2161 record.duration_ms = out.duration_ms;
2162 self.state.event(
2163 "judge",
2164 format!(
2165 "judge {} ranked {}",
2166 j + 1,
2167 record.ranking.iter().collect::<String>()
2168 ),
2169 );
2170 }
2171 Err(e) => {
2172 record.failed = Some(e.to_string());
2173 self.state
2174 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2175 }
2176 }
2177 self.state.judgements.push(record);
2178 self.state.save()?;
2179 }
2180 Ok(())
2181 }
2182
2183 async fn deliberate(&mut self) -> Result<()> {
2186 let run_id = self.state.id.clone();
2191 let prompts = self.state.config.prompts.clone();
2192 if !self.state.deliberation.is_empty() {
2193 return Ok(());
2194 }
2195 let tops: Vec<char> = self
2196 .state
2197 .judgements
2198 .iter()
2199 .filter_map(|j| j.ranking.first().copied())
2200 .collect();
2201 let rounds = self.state.config.graph.deliberate_rounds;
2202 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2203 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2204 self.state.event(
2205 "deliberate",
2206 format!("judges agreed on {} outright; no deliberation", tops[0]),
2207 );
2208 }
2209 self.state.status = RunStatus::Voting;
2210 self.state.save()?;
2211 return Ok(());
2212 }
2213
2214 self.state.status = RunStatus::Deliberating;
2215 self.state.event(
2216 "deliberate",
2217 format!(
2218 "split: first choices were {} — opening {rounds} round(s)",
2219 tops.iter().collect::<String>()
2220 ),
2221 );
2222
2223 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2224 let language = self.state.config.graph.language.clone();
2225 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2226 let sessions = self.state.config.graph.sessions;
2227 let artifacts = agent::artifacts_dir(&self.state.dir());
2228 let root = self.state.worktree_root();
2229 let base_short = short(&self.state.base_commit);
2230
2231 for round in 1..=rounds {
2235 let mut turns: Vec<DeliberationTurn> = Vec::new();
2236 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2237 if self.state.judgements[j].failed.is_some() {
2238 continue;
2239 }
2240 let seat_key = format!("judge-{}", j + 1);
2241 let mut seat = self.seat(&seat_key, &spec.id);
2242 let transcript = self.transcript(&turns, j);
2243 let context = if has_context(&spec, &seat, sessions) {
2244 None
2245 } else {
2246 Some(self.candidate_block(&viable, &base_short))
2247 };
2248 let text = prompt::deliberate(
2249 &self.state.instruction,
2250 context.as_deref(),
2251 &transcript,
2252 round,
2253 rounds,
2254 &language,
2255 );
2256 let job = SeatJob {
2257 spec,
2258 seat: seat.clone(),
2259 prompt: text,
2260 cwd: root.join(format!("judge-{}", j + 1)),
2261 timeout,
2262 allow_write: false,
2263 sessions,
2264 artifacts: artifacts.clone(),
2265 stem: format!("delib-{round}-judge-{}", j + 1),
2266 };
2267 let cache = self.state.config.cache_dir();
2268 let ctx = WaveCtx {
2269 run: &run_id,
2270 node: "deliberate",
2271 prompts: &prompts,
2272 cache: cache.as_deref(),
2273 round: None,
2274 };
2275 let (updated, out) =
2276 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2277 seat = updated;
2278 let agent_id = seat.agent.clone();
2279 let seat_key = seat.key.clone();
2280 self.state.seats.insert(seat.key.clone(), seat);
2281 let body = match out {
2282 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2283 AgentOutcome::Dropped(o) => {
2287 let why =
2288 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2289 "the CLI ended the stream without delivering its answer",
2290 );
2291 self.state.event(
2292 "deliberate",
2293 format!(
2294 "judge {} skipped: the CLI dropped the stream ({why})",
2295 j + 1
2296 ),
2297 );
2298 continue;
2299 }
2300 AgentOutcome::Quota(o) => {
2301 self.state.quota.push(QuotaLoss {
2302 seat: seat_key,
2303 node: "deliberate".to_owned(),
2304 at: Timestamp::now(),
2305 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2306 });
2307 self.state.event(
2308 "deliberate",
2309 format!("judge {} skipped: rate limited (quota)", j + 1),
2310 );
2311 continue;
2312 }
2313 AgentOutcome::Failed(e) => {
2314 self.state
2315 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2316 continue;
2317 }
2318 };
2319 let tentative = verdict::extract_json::<Position>(&body)
2320 .ok()
2321 .and_then(|p| p.tentative)
2322 .and_then(|s| s.trim().chars().next())
2323 .map(|c| c.to_ascii_uppercase());
2324 self.state.event(
2325 "deliberate",
2326 format!(
2327 "round {round}: judge {} now favours {}",
2328 j + 1,
2329 tentative.map_or("—".to_owned(), |c| c.to_string())
2330 ),
2331 );
2332 turns.push(DeliberationTurn {
2333 judge: j + 1,
2334 agent: agent_id,
2335 body: blind::sanitize_prose(&body, &self.state.config.blind),
2336 tentative,
2337 });
2338 }
2339 self.state
2340 .deliberation
2341 .push(DeliberationRound { round, turns });
2342 self.state.save()?;
2343 }
2344
2345 self.state.status = RunStatus::Voting;
2346 self.state.save()?;
2347 Ok(())
2348 }
2349
2350 async fn vote(&mut self) -> Result<()> {
2353 let run_id = self.state.id.clone();
2358 let prompts = self.state.config.prompts.clone();
2359 if !self.state.votes.is_empty() {
2360 return Ok(());
2361 }
2362 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2363 if viable.len() == 1 {
2364 return Ok(());
2365 }
2366 self.state.status = RunStatus::Voting;
2367
2368 let language = self.state.config.graph.language.clone();
2369 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2370 let sessions = self.state.config.graph.sessions;
2371 let artifacts = agent::artifacts_dir(&self.state.dir());
2372 let root = self.state.worktree_root();
2373 let base_short = short(&self.state.base_commit);
2374 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2375
2376 let mut jobs = Vec::new();
2377 let mut seats_at = Vec::new();
2378 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2379 if self
2380 .state
2381 .judgements
2382 .get(j)
2383 .is_some_and(|r| r.failed.is_some())
2384 {
2385 continue;
2386 }
2387 let seat_key = format!("judge-{}", j + 1);
2388 let seat = self.seat(&seat_key, &spec.id);
2389 let mut text = prompt::final_vote(&viable, &language);
2390 if !has_context(&spec, &seat, sessions) {
2391 text = format!(
2392 "{}\n\n# Candidates\n\n{}",
2393 text,
2394 self.candidate_block(&candidates, &base_short)
2395 );
2396 }
2397 jobs.push(SeatJob {
2398 spec,
2399 seat,
2400 prompt: text,
2401 cwd: root.join(format!("judge-{}", j + 1)),
2402 timeout,
2403 allow_write: false,
2404 sessions,
2405 artifacts: artifacts.clone(),
2406 stem: format!("vote-judge-{}", j + 1),
2407 });
2408 seats_at.push(j);
2409 }
2410
2411 self.state.event(
2412 "vote",
2413 format!(
2414 "collecting {} final votes one by one, privately",
2415 jobs.len()
2416 ),
2417 );
2418 let allowed = viable.clone();
2419 let mut quota_losses = Vec::new();
2420 let cache = self.state.config.cache_dir();
2421 let ctx = WaveCtx {
2422 run: &run_id,
2423 node: "vote",
2424 prompts: &prompts,
2425 cache: cache.as_deref(),
2426 round: None,
2427 };
2428 let results = ask_json_wave::<FinalVote>(
2429 jobs,
2430 Arc::clone(&self.sem),
2431 self.state.config.graph.retries,
2432 &ctx,
2433 &mut quota_losses,
2434 &mut self.state,
2435 &move |v: &FinalVote| match v.label() {
2436 Some(c) if allowed.contains(&c) => Ok(()),
2437 other => bail!("vote {other:?} is not one of {allowed:?}"),
2438 },
2439 )
2440 .await;
2441 self.state.quota.extend(quota_losses);
2442
2443 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2444 let agent_id = seat.agent.clone();
2445 self.state.seats.insert(seat.key.clone(), seat);
2446 let initial = self
2447 .state
2448 .judgements
2449 .get(j)
2450 .and_then(|r| r.ranking.first().copied());
2451 let mut record = VoteRecord {
2452 judge: j + 1,
2453 agent: agent_id,
2454 vote: None,
2455 reason: String::new(),
2456 changed: false,
2457 };
2458 match res {
2459 Ok((v, _)) => {
2460 record.vote = v.label();
2461 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2462 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2463 self.state.event(
2464 "vote",
2465 format!(
2466 "judge {} voted {}{}",
2467 j + 1,
2468 record.vote.unwrap_or('?'),
2469 if record.changed { " (changed)" } else { "" }
2470 ),
2471 );
2472 }
2473 Err(e) => {
2474 self.state
2475 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2476 }
2477 }
2478 self.state.votes.push(record);
2479 self.state.save()?;
2480 }
2481 Ok(())
2482 }
2483
2484 fn tally(&mut self) -> Result<()> {
2487 if self.state.tally.is_some() {
2488 return Ok(());
2489 }
2490 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2491 let tops: Vec<char> = self
2492 .state
2493 .judgements
2494 .iter()
2495 .filter_map(|j| j.ranking.first().copied())
2496 .collect();
2497 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2498
2499 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2502 let mut cast: Vec<char> = Vec::new();
2503 for (i, j) in self.state.judgements.iter().enumerate() {
2504 let vote = self
2505 .state
2506 .votes
2507 .iter()
2508 .find(|v| v.judge == i + 1)
2509 .and_then(|v| v.vote)
2510 .or_else(|| j.ranking.first().copied());
2511 if let Some(v) = vote {
2512 *first_choice.entry(v).or_insert(0) += 1;
2513 cast.push(v);
2514 }
2515 }
2516
2517 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2518 for j in &self.state.judgements {
2519 let n = j.ranking.len();
2520 for (pos, label) in j.ranking.iter().enumerate() {
2521 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2522 }
2523 }
2524
2525 let best = first_choice.values().copied().max().unwrap_or(0);
2526 let mut leaders: Vec<char> = first_choice
2527 .iter()
2528 .filter(|(_, v)| **v == best)
2529 .map(|(k, _)| *k)
2530 .collect();
2531 let mut tie_break = None;
2532 if leaders.len() > 1 {
2533 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2534 let borda_leaders: Vec<char> = leaders
2535 .iter()
2536 .copied()
2537 .filter(|l| borda[l] == top_borda)
2538 .collect();
2539 tie_break = Some(if borda_leaders.len() == 1 {
2540 format!(
2541 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2542 leaders.len()
2543 )
2544 } else {
2545 format!(
2546 "{} way tie on both first-choice votes and Borda points, broken by label order",
2547 leaders.len()
2548 )
2549 });
2550 leaders = borda_leaders;
2551 leaders.sort_unstable();
2552 }
2553 let winner = *leaders
2554 .first()
2555 .or(viable.first())
2556 .context("no candidate to declare a winner from")?;
2557
2558 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2559 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2560 let deliberated = !self.state.deliberation.is_empty();
2561
2562 let quota_seats: std::collections::BTreeSet<&str> =
2566 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2567 let mut present = 0usize;
2568 for (i, j) in self.state.judgements.iter().enumerate() {
2569 if quota_seats.contains(j.seat.as_str()) {
2570 continue;
2571 }
2572 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2573 let voted = self
2574 .state
2575 .votes
2576 .iter()
2577 .any(|v| v.judge == i + 1 && v.vote.is_some());
2578 if ranked || voted {
2579 present += 1;
2580 }
2581 }
2582 let needs_quorum = viable.len() > 1;
2588 let judges_total = if needs_quorum {
2589 self.roles.judges.len()
2590 } else {
2591 0
2592 };
2593 let quorum = if needs_quorum {
2594 judges_total / 2 + 1
2595 } else {
2596 0
2597 };
2598 let met_quorum = !needs_quorum || present >= quorum;
2599 let uncontested = (!needs_quorum).then(|| {
2600 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2601 });
2602
2603 self.state.event(
2604 "tally",
2605 match &uncontested {
2606 Some(reason) => format!("winner {winner} — {reason}"),
2607 None => format!(
2608 "winner {winner} — votes {} | initial {} | {} changed | \
2609 {present}/{judges_total} judges{}",
2610 first_choice
2611 .iter()
2612 .map(|(k, v)| format!("{k}:{v}"))
2613 .collect::<Vec<_>>()
2614 .join(" "),
2615 if unanimous_initial {
2616 "unanimous"
2617 } else {
2618 "split"
2619 },
2620 changed_votes,
2621 if met_quorum {
2622 String::new()
2623 } else {
2624 format!(" — below quorum ({quorum} required)")
2625 },
2626 ),
2627 },
2628 );
2629 if !met_quorum {
2630 self.state.event(
2631 "stall",
2632 format!(
2633 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2634 the run stops here, resumable"
2635 ),
2636 );
2637 }
2638 self.state.tally = Some(Tally {
2639 first_choice,
2640 borda,
2641 winner,
2642 rankings: tops.len(),
2643 unanimous_initial,
2644 deliberated,
2645 changed_votes,
2646 unanimous_final,
2647 tie_break,
2648 judges: judges_total,
2649 present,
2650 quorum,
2651 met_quorum,
2652 uncontested,
2653 });
2654 self.state.status = if met_quorum {
2655 RunStatus::Reviewing
2656 } else {
2657 RunStatus::Stalled
2658 };
2659 self.state.save()?;
2660 Ok(())
2661 }
2662
2663 #[allow(clippy::too_many_lines)]
2684 async fn recover_stall(&mut self) -> Result<bool> {
2685 let run_id = self.state.id.clone();
2690 let prompts = self.state.config.prompts.clone();
2691 let quota_seats: BTreeSet<&str> =
2696 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2697 let absent: Vec<String> = self
2698 .state
2699 .judgements
2700 .iter()
2701 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2702 .map(|j| j.seat.clone())
2703 .collect();
2704 if absent.is_empty() {
2705 return Ok(false);
2706 }
2707 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2708 if viable.len() <= 1 {
2709 return Ok(false);
2710 }
2711 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2712 let language = self.state.config.graph.language.clone();
2713 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2714 let sessions = self.state.config.graph.sessions;
2715 let artifacts = agent::artifacts_dir(&self.state.dir());
2716 let root = self.state.worktree_root();
2717 let base_short = short(&self.state.base_commit);
2718 let candidates: Vec<Candidate> = viable.clone();
2719
2720 let mut positions: Vec<usize> = absent
2722 .iter()
2723 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2724 .collect();
2725 if positions.is_empty() {
2726 return Ok(false);
2727 }
2728 positions.sort_unstable();
2729 positions.dedup();
2730
2731 let mut judge_jobs = Vec::new();
2733 for &j in &positions {
2734 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2735 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2736 let seat_key = format!("judge-{}", j + 1);
2737 let spec = self.roles.judges[j].clone();
2738 let seat = self.seat(&seat_key, &spec.id);
2739 judge_jobs.push(SeatJob {
2740 spec,
2741 seat,
2742 prompt: prompt::judge(
2743 &self.state.instruction,
2744 &views,
2745 self.roles.judges.len(),
2746 &base_short,
2747 &language,
2748 ),
2749 cwd: root.join(seat_key),
2750 timeout,
2751 allow_write: false,
2752 sessions,
2753 artifacts: artifacts.clone(),
2754 stem: format!("judge-{}-recover", j + 1),
2755 });
2756 }
2757
2758 let labels_for_check = labels.clone();
2759 let mut judge_losses = Vec::new();
2760 let retries = self.state.config.graph.retries;
2761 let cache = self.state.config.cache_dir();
2762 let ctx = WaveCtx {
2763 run: &run_id,
2764 node: "judge",
2765 prompts: &prompts,
2766 cache: cache.as_deref(),
2767 round: None,
2768 };
2769 let results = ask_json_wave::<Ranking>(
2770 judge_jobs,
2771 Arc::clone(&self.sem),
2772 retries,
2773 &ctx,
2774 &mut judge_losses,
2775 &mut self.state,
2776 &move |r: &Ranking| r.validate(&labels_for_check),
2777 )
2778 .await;
2779
2780 let mut recovered: BTreeSet<usize> = BTreeSet::new();
2782 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
2783 self.state.seats.insert(seat.key.clone(), seat);
2784 let record = &mut self.state.judgements[j];
2785 match res {
2786 Ok((ranking, out)) => {
2787 record.ranking = ranking.normalized();
2788 record.reasons = ranking.reasons;
2789 record.confidence = ranking.confidence;
2790 record.failed = None;
2791 record.duration_ms = out.duration_ms;
2792 recovered.insert(j);
2793 self.state.event(
2794 "recover",
2795 format!("judge {} ranked again after the limit", j + 1),
2796 );
2797 }
2798 Err(e) => {
2799 self.state
2800 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
2801 }
2802 }
2803 }
2804
2805 let mut vote_jobs = Vec::new();
2807 let mut vote_pos: Vec<usize> = Vec::new();
2808 for &j in &recovered {
2809 let seat_key = format!("judge-{}", j + 1);
2810 let spec = self.roles.judges[j].clone();
2811 let seat = self.seat(&seat_key, &spec.id);
2812 let mut text = prompt::final_vote(&labels, &language);
2813 if !has_context(&spec, &seat, sessions) {
2814 text = format!(
2815 "{}\n\n# Candidates\n\n{}",
2816 text,
2817 self.candidate_block(&candidates, &base_short)
2818 );
2819 }
2820 vote_jobs.push(SeatJob {
2821 spec,
2822 seat,
2823 prompt: text,
2824 cwd: root.join(seat_key),
2825 timeout,
2826 allow_write: false,
2827 sessions,
2828 artifacts: artifacts.clone(),
2829 stem: format!("vote-judge-{}-recover", j + 1),
2830 });
2831 vote_pos.push(j);
2832 }
2833 let allowed = labels.clone();
2834 let mut vote_losses = Vec::new();
2835 let vote_retries = self.state.config.graph.retries;
2836 let vote_cache = self.state.config.cache_dir();
2837 let ctx = WaveCtx {
2838 run: &run_id,
2839 node: "vote",
2840 prompts: &prompts,
2841 cache: vote_cache.as_deref(),
2842 round: None,
2843 };
2844 let votes = ask_json_wave::<FinalVote>(
2845 vote_jobs,
2846 Arc::clone(&self.sem),
2847 vote_retries,
2848 &ctx,
2849 &mut vote_losses,
2850 &mut self.state,
2851 &move |v: &FinalVote| match v.label() {
2852 Some(c) if allowed.contains(&c) => Ok(()),
2853 other => bail!("vote {other:?} is not one of {allowed:?}"),
2854 },
2855 )
2856 .await;
2857 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
2858 let agent_id = seat.agent.clone();
2859 self.state.seats.insert(seat.key.clone(), seat);
2860 match res {
2861 Ok((v, _)) => {
2862 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
2863 rec.vote = v.label();
2864 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2865 } else {
2866 self.state.votes.push(VoteRecord {
2867 judge: j + 1,
2868 agent: agent_id,
2869 vote: v.label(),
2870 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
2871 changed: false,
2872 });
2873 }
2874 self.state.event(
2875 "recover",
2876 format!("judge {} voted again after the limit", j + 1),
2877 );
2878 }
2879 Err(e) => {
2880 self.state
2881 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
2882 }
2883 }
2884 }
2885
2886 let recovered_keys: BTreeSet<String> = recovered
2890 .iter()
2891 .map(|&j| format!("judge-{}", j + 1))
2892 .collect();
2893 self.state
2894 .quota
2895 .retain(|q| !recovered_keys.contains(&q.seat));
2896 for loss in judge_losses.into_iter().chain(vote_losses) {
2900 if recovered_keys.contains(&loss.seat) {
2901 continue;
2902 }
2903 self.state.quota.retain(|q| q.seat != loss.seat);
2904 self.state.quota.push(loss);
2905 }
2906
2907 self.state.tally = None;
2909 self.tally()?;
2910 Ok(self
2911 .state
2912 .tally
2913 .as_ref()
2914 .map(|t| t.met_quorum)
2915 .unwrap_or(false))
2916 }
2917
2918 async fn fold_losers(&mut self) -> Result<()> {
2921 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
2922 return Ok(());
2923 };
2924 let repo = self.state.repo.clone();
2925 let mut folded = Vec::new();
2926 for i in 0..self.state.candidates.len() {
2927 let c = &self.state.candidates[i];
2928 if c.label == winner || c.folded {
2929 continue;
2930 }
2931 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
2932 git::worktree_remove(&repo, &wt).await.ok();
2933 git::branch_delete(&repo, &branch).await.ok();
2934 self.state.candidates[i].folded = true;
2935 folded.push(label.to_string());
2936 }
2937 let root = self.state.worktree_root();
2939 for j in 1..=self.roles.judges.len() {
2940 let wt = root.join(format!("judge-{j}"));
2941 if wt.exists() {
2942 git::worktree_remove(&repo, &wt).await.ok();
2943 }
2944 }
2945 if self.state.config.graph.advise {
2948 for k in 1..=self.state.config.graph.advisors {
2949 let wt = root.join(format!("advisor-{k}"));
2950 if wt.exists() {
2951 git::worktree_remove(&repo, &wt).await.ok();
2952 }
2953 }
2954 }
2955 if !folded.is_empty() {
2956 self.state
2957 .event("fold", format!("folded candidates {}", folded.join(", ")));
2958 self.state.save()?;
2959 }
2960 Ok(())
2961 }
2962
2963 async fn sync_to_base(&mut self) -> Result<()> {
2993 if self
2994 .state
2995 .base_sync
2996 .as_ref()
2997 .is_some_and(|s| s.conflict.is_some())
2998 {
2999 return Ok(());
3000 }
3001 let Some(winner) = self.state.winner().cloned() else {
3002 return Ok(());
3003 };
3004
3005 let repo = self.state.repo.clone();
3006 let remote = self.state.config.merge.remote.clone();
3007 let base_branch = self.state.base_branch.clone();
3008 let tracking = format!("{remote}/{base_branch}");
3009
3010 git::fetch(&repo, &remote, &base_branch).await.ok();
3011 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3015 return Ok(());
3016 };
3017
3018 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3019 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3020 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3021
3022 if behind == 0 {
3023 self.state.base_sync = Some(BaseSync {
3024 tip,
3025 behind: 0,
3026 attempts,
3027 conflict: None,
3028 });
3029 self.state.save()?;
3030 return Ok(());
3031 }
3032
3033 if attempts >= BASE_SYNC_ROUNDS {
3034 let why = format!(
3035 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3036 rebase(s); rebasing again would only race it",
3037 winner.branch
3038 );
3039 self.state.status = RunStatus::Blocked;
3040 self.state.base_sync = Some(BaseSync {
3041 tip,
3042 behind,
3043 attempts,
3044 conflict: Some(why.clone()),
3045 });
3046 self.state.event("land", why);
3047 self.state.save()?;
3048 return Ok(());
3049 }
3050
3051 self.state.event(
3052 "land",
3053 format!(
3054 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3055 winner.branch
3056 ),
3057 );
3058 self.state.save()?;
3059
3060 let scratch = self.state.dir().join("base-sync");
3061 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
3062 let attempts = attempts + 1;
3063 match rebased {
3064 Ok(None) => {
3065 git::sync_to_head(&winner.worktree).await?;
3069 self.state.base_sync = Some(BaseSync {
3070 tip: tip.clone(),
3071 behind: 0,
3072 attempts,
3073 conflict: None,
3074 });
3075 self.state
3076 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3077 }
3078 Ok(Some(conflict)) => {
3079 let why = format!(
3080 "{} conflicts with {tracking} and did not rebase: {}",
3081 winner.branch,
3082 conflict.chars().take(600).collect::<String>()
3083 );
3084 self.state.status = RunStatus::Blocked;
3085 self.state.base_sync = Some(BaseSync {
3086 tip,
3087 behind,
3088 attempts,
3089 conflict: Some(why.clone()),
3090 });
3091 self.state.event("land", why);
3092 }
3093 Err(e) => {
3094 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3095 self.state.status = RunStatus::Blocked;
3096 self.state.base_sync = Some(BaseSync {
3097 tip,
3098 behind,
3099 attempts,
3100 conflict: Some(why.clone()),
3101 });
3102 self.state.event("land", why);
3103 }
3104 }
3105 self.state.save()?;
3106 Ok(())
3107 }
3108
3109 fn landing_base(&self) -> String {
3119 self.state
3120 .base_sync
3121 .as_ref()
3122 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3123 }
3124
3125 pub async fn fix_selected(
3158 &mut self,
3159 ids: &[String],
3160 reason: &str,
3161 allow_stale: bool,
3162 ) -> Result<()> {
3163 let reason = reason.trim();
3164 if reason.is_empty() {
3165 bail!("a fix request needs a reason — that is the operator's own record of why");
3166 }
3167 if ids.is_empty() {
3168 bail!("no finding id given");
3169 }
3170 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3171 bail!(
3172 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3173 has already concluded — can be given a targeted fix. A run still \
3174 in progress should simply be resumed; a `merged` run's branch has \
3175 already landed, so its answer is a fresh `magi review <branch>`, \
3176 not reopening this run's own record",
3177 self.state.id,
3178 self.state.status.as_str()
3179 );
3180 }
3181 let Some(winner) = self.state.winner().cloned() else {
3182 bail!("run {} has no winning candidate to fix", self.state.id);
3183 };
3184 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3185 bail!(
3186 "branch `{}` no longer exists; this run cannot be extended",
3187 winner.branch
3188 );
3189 }
3190 let home = crate::run::home();
3191 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3192 bail!(
3193 "run {} is currently being worked on by another magi process",
3194 self.state.id
3195 );
3196 }
3197 let _claim = FixClaim::acquire(&self.state.dir())?;
3203
3204 let mut seen = BTreeSet::new();
3208 let mut findings = Vec::new();
3209 let mut missing = Vec::new();
3210 for id in ids {
3211 if !seen.insert(id.clone()) {
3212 continue;
3213 }
3214 match self.state.finding(id) {
3215 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3216 id: f.id.clone(),
3217 severity: f.severity,
3218 reviewer_vote: rec.vote,
3219 round: round.round,
3220 round_head: round.head.clone(),
3221 reviewer: rec.reviewer,
3222 agent: rec.agent.clone(),
3223 file: f.file.clone(),
3224 line: f.line,
3225 title: f.title.clone(),
3226 detail: f.detail.clone(),
3227 outcome: OperatorFixOutcome::Pending,
3228 }),
3229 None => missing.push(id.clone()),
3230 }
3231 }
3232 if !missing.is_empty() {
3233 bail!(
3234 "unknown finding id(s): {}; nothing was changed",
3235 missing.join(", ")
3236 );
3237 }
3238
3239 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3240 let stale_details: Vec<(String, String)> = findings
3241 .iter()
3242 .filter(|f| f.round_head != head_at_request)
3243 .map(|f| (f.id.clone(), f.round_head.clone()))
3244 .collect();
3245 let stale = !stale_details.is_empty();
3246 if stale && !allow_stale {
3247 bail!(
3248 "the branch has moved since some finding(s) were raised — {} — now \
3249 at {}; pass --allow-stale to fix anyway, or re-run review first",
3250 stale_details
3251 .iter()
3252 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3253 .collect::<Vec<_>>()
3254 .join(", "),
3255 short(&head_at_request)
3256 );
3257 }
3258
3259 let request = OperatorFixRequest {
3260 requested_at: Timestamp::now(),
3261 reason: reason.to_owned(),
3262 findings,
3263 head_at_request: head_at_request.clone(),
3264 allow_stale,
3265 stale,
3266 fix: None,
3267 result_head: None,
3268 follow_up_review_run: None,
3269 };
3270 self.state.event(
3271 "fix",
3272 format!(
3273 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3274 request.findings.len(),
3275 request
3276 .findings
3277 .iter()
3278 .map(|f| f.id.as_str())
3279 .collect::<Vec<_>>()
3280 .join(", "),
3281 ),
3282 );
3283 self.state.operator_fixes.push(request);
3290 self.state.save()?;
3291 let request_index = self.state.operator_fixes.len() - 1;
3292
3293 if winner.worktree.exists() {
3302 let dirty = git::git(
3305 &winner.worktree,
3306 &["status", "--porcelain", "--untracked-files=all"],
3307 )
3308 .await?;
3309 let only_withheld = dirty.lines().all(|l| {
3310 l.strip_prefix("?? ")
3311 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3312 });
3313 if !only_withheld {
3314 bail!(
3315 "`{}` has uncommitted changes; refusing to touch it — commit or \
3316 discard them first",
3317 winner.worktree.display()
3318 );
3319 }
3320 git::worktree_remove(&self.state.repo, &winner.worktree)
3321 .await
3322 .ok();
3323 }
3324 let fix_worktree = self.state.worktree_root().join("operator-fix");
3325 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3326 git::git(
3327 &self.state.repo,
3328 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3329 )
3330 .await
3331 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3332 if !git::is_clean(&fix_worktree).await? {
3333 git::worktree_remove(&self.state.repo, &fix_worktree)
3334 .await
3335 .ok();
3336 bail!(
3337 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3338 winner.branch
3339 );
3340 }
3341
3342 let run_id = self.state.id.clone();
3343 let prompts = self.state.config.prompts.clone();
3344 let language = self.state.config.graph.language.clone();
3345 let sessions = self.state.config.graph.sessions;
3346 let artifacts = agent::artifacts_dir(&self.state.dir());
3347 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3348 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3349 _ => (
3350 self.state
3351 .config
3352 .agent(&winner.agent)
3353 .cloned()
3354 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3355 format!("impl-{}", winner.label),
3356 ),
3357 };
3358 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3359 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3360 .findings
3361 .iter()
3362 .map(|f| Finding {
3363 id: f.id.clone(),
3364 severity: f.severity,
3365 file: f.file.clone(),
3366 line: f.line,
3367 title: f.title.clone(),
3368 detail: f.detail.clone(),
3369 })
3370 .collect();
3371 let job = SeatJob {
3372 prompt: prompt::operator_fix(
3373 &self.state.instruction,
3374 &finding_list,
3375 reason,
3376 &stale_details,
3377 &head_at_request,
3378 &language,
3379 ),
3380 spec: fix_spec.clone(),
3381 seat,
3382 cwd: fix_worktree.clone(),
3383 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3384 allow_write: true,
3385 sessions,
3386 artifacts: artifacts.clone(),
3387 stem: "operator-fix".to_owned(),
3388 };
3389 let cache = self.state.config.cache_dir();
3390 let ctx = WaveCtx {
3391 run: &run_id,
3392 node: "fix",
3393 prompts: &prompts,
3394 cache: cache.as_deref(),
3395 round: None,
3396 };
3397 let (seat, out) =
3398 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3399 let agent_id = seat.agent.clone();
3400
3401 let mut fix = FixRecord {
3402 agent: agent_id,
3403 addressed: Vec::new(),
3404 rejected: Vec::new(),
3405 notes: String::new(),
3406 committed: false,
3407 failed: None,
3408 duration_ms: 0,
3409 continuation: None,
3410 };
3411 let mut final_seat = seat.clone();
3412 match out {
3413 AgentOutcome::Ok(o) => {
3414 fix.duration_ms = o.duration_ms;
3415 let parsed = verdict::extract_json::<FixReport>(&o.text);
3416 let incomplete_reason = match &parsed {
3417 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3418 "the reply parsed, but it reported a command whose own CLI \
3419 never confirmed an exit status"
3420 .to_owned(),
3421 ),
3422 Ok(_) => None,
3423 Err(e) => Some(e.to_string()),
3424 };
3425 match incomplete_reason {
3426 None => {
3427 let report = parsed.expect("checked Ok above");
3428 fix.addressed = report.addressed;
3429 fix.rejected = report.rejected;
3430 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3431 }
3432 Some(reason) => {
3433 let (resumed_seat, resolved, failure, cont) = self
3434 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3435 .await;
3436 fix.duration_ms += cont.cumulative_wait_ms;
3437 fix.continuation = Some(cont);
3438 final_seat = resumed_seat;
3439 match resolved {
3440 Some(report) => {
3441 fix.addressed = report.addressed;
3442 fix.rejected = report.rejected;
3443 fix.notes =
3444 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3445 }
3446 None => fix.failed = failure,
3447 }
3448 }
3449 }
3450 }
3451 AgentOutcome::Dropped(o) => {
3452 fix.duration_ms = o.duration_ms;
3453 let why = o
3454 .dropped
3455 .as_ref()
3456 .map(|d| d.why.as_str())
3457 .unwrap_or("the CLI ended the stream without delivering its answer");
3458 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3459 }
3460 AgentOutcome::Quota(o) => {
3461 self.state.quota.push(QuotaLoss {
3462 seat: final_seat.key.clone(),
3463 node: "fix".to_owned(),
3464 at: Timestamp::now(),
3465 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3466 });
3467 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3468 }
3469 AgentOutcome::Failed(e) => fix.failed = Some(e),
3470 }
3471 if fix.continuation.is_none() {
3472 fix.continuation = Some(ContinuationRecord::not_needed());
3473 }
3474 self.state.seats.insert(final_seat.key.clone(), final_seat);
3475
3476 let rescue_message = format!(
3477 "magi: operator-selected fix ({}) (uncommitted work)",
3478 self.state.operator_fixes[request_index]
3479 .findings
3480 .iter()
3481 .map(|f| f.id.as_str())
3482 .collect::<Vec<_>>()
3483 .join(", ")
3484 );
3485 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3486 self.state.note_withheld("fix", &r.withheld);
3487 }
3488 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3489 fix.committed = after != head_at_request;
3490 git::worktree_remove(&self.state.repo, &fix_worktree)
3491 .await
3492 .ok();
3493
3494 self.state.event(
3495 "fix",
3496 match &fix.failed {
3497 Some(reason) => format!(
3498 "operator fix: adoption report was lost ({reason}); {}",
3499 if fix.committed {
3500 "committed"
3501 } else {
3502 "NO new commit"
3503 }
3504 ),
3505 None => format!(
3506 "operator fix: {} addressed, {} rejected, {}",
3507 fix.addressed.len(),
3508 fix.rejected.len(),
3509 if fix.committed {
3510 "committed"
3511 } else {
3512 "NO new commit"
3513 }
3514 ),
3515 },
3516 );
3517
3518 for f in &mut self.state.operator_fixes[request_index].findings {
3525 f.outcome = if fix.failed.is_some() {
3526 OperatorFixOutcome::Unreported
3527 } else if fix.addressed.contains(&f.id) {
3528 OperatorFixOutcome::Addressed
3529 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3530 OperatorFixOutcome::Rejected { why: r.why.clone() }
3531 } else {
3532 OperatorFixOutcome::Unreported
3533 };
3534 }
3535
3536 let committed = fix.committed;
3537 if committed {
3538 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3539 }
3540 self.state.operator_fixes[request_index].fix = Some(fix);
3541 self.state.save()?;
3544
3545 if committed {
3546 self.state.event(
3547 "fix",
3548 format!(
3549 "operator fix committed {}; opening a follow-up review-only run",
3550 short(&after)
3551 ),
3552 );
3553 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3554 Ok(mut follow_up) => {
3555 follow_up.state.event(
3556 "start",
3557 format!(
3558 "requested by an operator fix on run {} for finding(s) {}",
3559 self.state.id,
3560 self.state.operator_fixes[request_index]
3561 .findings
3562 .iter()
3563 .map(|f| f.id.as_str())
3564 .collect::<Vec<_>>()
3565 .join(", "),
3566 ),
3567 );
3568 follow_up.state.save()?;
3569 let follow_up_id = follow_up.state.id.clone();
3570 if let Err(e) = follow_up.execute().await {
3571 self.state.event(
3572 "fix",
3573 format!(
3574 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3575 ),
3576 );
3577 }
3578 self.state.operator_fixes[request_index].follow_up_review_run =
3579 Some(follow_up_id);
3580 }
3581 Err(e) => {
3582 self.state.event(
3583 "fix",
3584 format!("committed the fix but could not open a follow-up review: {e:#}"),
3585 );
3586 }
3587 }
3588 self.state.save()?;
3589 }
3590
3591 Ok(())
3592 }
3593
3594 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3601 match &self.roles.fixer {
3602 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3603 _ => (
3604 self.state
3605 .config
3606 .agent(&winner.agent)
3607 .cloned()
3608 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3609 format!("impl-{}", winner.label),
3610 ),
3611 }
3612 }
3613
3614 async fn review_loop(&mut self) -> Result<()> {
3615 if self
3620 .state
3621 .base_sync
3622 .as_ref()
3623 .is_some_and(|s| s.conflict.is_some())
3624 {
3625 return Ok(());
3626 }
3627 let run_id = self.state.id.clone();
3632 let prompts = self.state.config.prompts.clone();
3633 let Some(winner) = self.state.winner().cloned() else {
3634 return Ok(());
3635 };
3636 let max_rounds = self.state.config.graph.review_rounds;
3637 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
3647 self.state.status = status;
3648 self.state.save()?;
3649 return Ok(());
3650 }
3651 self.state.status = RunStatus::Reviewing;
3652 if self
3662 .state
3663 .reviews
3664 .last()
3665 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
3666 {
3667 let shell = self.state.config.shell();
3668 return self
3669 .stop_reviewing(
3670 "the last round's own verification never resolved",
3671 &shell,
3672 &winner.worktree,
3673 )
3674 .await;
3675 }
3676
3677 let repo = self.state.repo.clone();
3678 let root = self.state.worktree_root();
3679 let language = self.state.config.graph.language.clone();
3680 let sessions = self.state.config.graph.sessions;
3681 let artifacts = agent::artifacts_dir(&self.state.dir());
3682 let base = self.landing_base();
3683 let base_short = short(&base);
3684 let reviewers = self.roles.reviewers.clone();
3685 let shell = self.state.config.shell();
3686
3687 for round in (self.state.reviews.len() + 1)..=max_rounds {
3688 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3689 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
3690 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
3691 let prev_verification = self
3700 .state
3701 .reviews
3702 .last()
3703 .and_then(|r| r.verification_summary(&head));
3704
3705 let mut jobs = Vec::new();
3709 for (r, spec) in reviewers.iter().cloned().enumerate() {
3710 let wt = root.join(format!("review-{}", r + 1));
3711 if wt.exists() {
3712 git::reset_detached(&wt, &head).await?;
3713 } else {
3714 git::worktree_add_detached(&repo, &wt, &head).await?;
3715 }
3716 let seat_key = format!("review-{}", r + 1);
3717 let seat = self.seat(&seat_key, &spec.id);
3718 jobs.push(SeatJob {
3719 prompt: prompt::review(&prompt::ReviewCtx {
3720 instruction: &self.state.instruction,
3721 branch: &winner.branch,
3722 base_short: &base_short,
3723 stat: &stat,
3724 patch: &patch,
3725 verification: prev_verification.as_ref(),
3726 reviewers: reviewers.len(),
3727 round,
3728 rounds: max_rounds,
3729 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
3732 lens: Lens::for_seat(r),
3733 language: &language,
3734 }),
3735 spec,
3736 seat,
3737 cwd: wt,
3738 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
3739 allow_write: false,
3740 sessions,
3741 artifacts: artifacts.clone(),
3742 stem: format!("review-{round}-{}", r + 1),
3743 });
3744 }
3745
3746 self.state.event(
3747 "review",
3748 format!(
3749 "round {round}: {} reviewers on {}",
3750 jobs.len(),
3751 short(&head)
3752 ),
3753 );
3754 let mut quota_losses = Vec::new();
3755 let review_retries = self.state.config.graph.retries;
3756 let review_cache = self.state.config.cache_dir();
3757 let ctx = WaveCtx {
3758 run: &run_id,
3759 node: "review",
3760 prompts: &prompts,
3761 cache: review_cache.as_deref(),
3762 round: Some(round),
3763 };
3764 let results = ask_json_wave::<Review>(
3765 jobs,
3766 Arc::clone(&self.sem),
3767 review_retries,
3768 &ctx,
3769 &mut quota_losses,
3770 &mut self.state,
3771 &|_: &Review| Ok(()),
3772 )
3773 .await;
3774 let round_quota_missing = quota_losses.len();
3778 self.state.quota.extend(quota_losses);
3779
3780 let mut records = Vec::new();
3781 let mut all_findings = Vec::new();
3782 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
3783 let agent_id = seat.agent.clone();
3784 self.state.seats.insert(seat.key.clone(), seat);
3785 let mut record = ReviewRecord {
3786 reviewer: r + 1,
3787 agent: agent_id,
3788 summary: String::new(),
3789 findings: Vec::new(),
3790 vote: None,
3791 failed: None,
3792 duration_ms: 0,
3793 attempts,
3799 };
3800 match res {
3801 Ok((review, out)) => {
3802 record.summary =
3810 blind::sanitize_prose(&review.summary, &self.state.config.blind);
3811 record.vote = Some(review.vote);
3812 record.duration_ms = out.duration_ms;
3813 for (n, mut f) in review.findings.into_iter().enumerate() {
3814 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
3817 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
3818 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
3819 f.file = f
3825 .file
3826 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
3827 all_findings.push(f.clone());
3828 record.findings.push(f);
3829 }
3830 self.state.event(
3831 "review",
3832 format!(
3833 "round {round}: reviewer {} voted {} with {} finding(s)",
3834 r + 1,
3835 review.vote.label(),
3836 record.findings.len()
3837 ),
3838 );
3839 }
3840 Err(e) => {
3841 record.failed = Some(e.to_string());
3842 self.state.event(
3843 "review",
3844 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
3845 );
3846 }
3847 }
3848 records.push(record);
3849 }
3850
3851 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
3858 let vote_split =
3859 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
3860 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
3861 if vote_split {
3862 self.state.event(
3863 "review",
3864 format!(
3865 "round {round}: votes split ({}) — one round of reconsideration",
3866 initial_votes
3867 .iter()
3868 .map(|v| v.label())
3869 .collect::<Vec<_>>()
3870 .join(", ")
3871 ),
3872 );
3873 let panel: Vec<ReviewSeatReport<'_>> = records
3876 .iter()
3877 .filter_map(|r| {
3878 r.vote.map(|vote| ReviewSeatReport {
3879 reviewer: r.reviewer,
3880 vote,
3881 summary: &r.summary,
3882 findings: &r.findings,
3883 })
3884 })
3885 .collect();
3886
3887 let mut jobs = Vec::new();
3888 let mut seats_at = Vec::new();
3889 for (r, spec) in reviewers.iter().cloned().enumerate() {
3890 if records[r].vote.is_none() {
3894 continue;
3895 }
3896 let wt = root.join(format!("review-{}", r + 1));
3897 let seat_key = format!("review-{}", r + 1);
3898 let seat = self.seat(&seat_key, &spec.id);
3899 let patch_ctx = if has_context(&spec, &seat, sessions) {
3904 None
3905 } else {
3906 Some(ReviewPatch {
3907 branch: &winner.branch,
3908 base_short: &base_short,
3909 stat: &stat,
3910 patch: &patch,
3911 })
3912 };
3913 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
3914 instruction: &self.state.instruction,
3915 reviewer: r + 1,
3916 lens: Lens::for_seat(r),
3917 panel: &panel,
3918 patch: patch_ctx,
3919 round,
3920 rounds: max_rounds,
3921 language: &language,
3922 });
3923 jobs.push(SeatJob {
3924 prompt,
3925 spec,
3926 seat,
3927 cwd: wt,
3928 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
3929 allow_write: false,
3930 sessions,
3931 artifacts: artifacts.clone(),
3932 stem: format!("review-{round}-reconsider-{}", r + 1),
3933 });
3934 seats_at.push(r);
3935 }
3936
3937 let mut recon_quota_losses = Vec::new();
3938 let recon_cache = self.state.config.cache_dir();
3939 let recon_ctx = WaveCtx {
3940 run: &run_id,
3941 node: "review",
3942 prompts: &prompts,
3943 cache: recon_cache.as_deref(),
3944 round: Some(round),
3945 };
3946 let recon_results = ask_json_wave::<ReviewRevote>(
3947 jobs,
3948 Arc::clone(&self.sem),
3949 review_retries,
3950 &recon_ctx,
3951 &mut recon_quota_losses,
3952 &mut self.state,
3953 &|_: &ReviewRevote| Ok(()),
3954 )
3955 .await;
3956 self.state.quota.extend(recon_quota_losses);
3957
3958 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
3959 let agent_id = seat.agent.clone();
3960 self.state.seats.insert(seat.key.clone(), seat);
3961 let mut rec = ReviewRevoteRecord {
3962 reviewer: r + 1,
3963 agent: agent_id,
3964 vote: None,
3965 reason: String::new(),
3966 failed: None,
3967 };
3968 match res {
3969 Ok((rv, _)) => {
3970 rec.vote = Some(rv.vote);
3971 rec.reason =
3972 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
3973 self.state.event(
3974 "review",
3975 format!(
3976 "round {round}: reviewer {} revoted {}",
3977 r + 1,
3978 rv.vote.label()
3979 ),
3980 );
3981 }
3982 Err(e) => {
3983 rec.failed = Some(e.to_string());
3984 self.state.event(
3985 "review",
3986 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
3987 );
3988 }
3989 }
3990 reconsideration.push(rec);
3991 }
3992 } else if initial_votes.len() > 1 {
3993 self.state.event(
3994 "review",
3995 format!(
3996 "round {round}: votes agreed ({}) — no reconsideration",
3997 initial_votes[0].label()
3998 ),
3999 );
4000 }
4001
4002 let final_votes: Vec<ReviewVote> = records
4006 .iter()
4007 .filter_map(|r| {
4008 reconsideration
4009 .iter()
4010 .find(|rv| rv.reviewer == r.reviewer)
4011 .and_then(|rv| rv.vote)
4012 .or(r.vote)
4013 })
4014 .collect();
4015 let round_verdict = ReviewVote::worst(final_votes);
4016
4017 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4018 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4019 let defer_e2e =
4030 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4031 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4032 let reason =
4033 format!("{blocking} blocking finding(s) already required a fix this round");
4034 self.state.event(
4035 "verify",
4036 format!(
4037 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4038 {}); it will run once a round has none left",
4039 short(&head)
4040 ),
4041 );
4042 (Vec::new(), false, true, Some(reason))
4043 } else {
4044 let e2e_commands = self.state.config.verify.e2e.clone();
4045 let cache_dir = self.state.config.cache_dir();
4046 let context = format!("round {round}");
4047 let (e2e, verify_retried) = with_cache_lease(
4048 &mut self.state,
4049 cache_dir.as_deref(),
4050 "e2e",
4051 "e2e",
4052 &winner.worktree,
4053 &head,
4054 verify_timeout,
4055 &context,
4056 |state, budget| {
4057 let shell = shell.clone();
4058 let e2e_commands = e2e_commands.clone();
4059 let worktree = winner.worktree.clone();
4060 let context = context.clone();
4061 async move {
4062 run_e2e_with_retry(
4063 state,
4064 &shell,
4065 &e2e_commands,
4066 &worktree,
4067 budget,
4068 &context,
4069 )
4070 .await
4071 }
4072 },
4073 )
4074 .await;
4075 (e2e, verify_retried, false, None)
4076 };
4077
4078 let expected = records.len();
4079 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4080 let incomplete = answered < expected;
4081 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4082 let policy = self.state.config.graph.incomplete_review;
4083 let clean = round_is_clean(
4084 blocking,
4085 e2e_ok,
4086 answered,
4087 expected,
4088 round_quota_missing,
4089 policy,
4090 );
4091
4092 let mut round_record = ReviewRound {
4093 round,
4094 head: head.clone(),
4095 verified_head: None,
4096 verified_at: None,
4097 reviews: records,
4098 e2e,
4099 verify_retried,
4100 e2e_deferred,
4101 e2e_defer_reason,
4102 fix: None,
4103 blocking,
4104 answered,
4105 expected,
4106 clean,
4107 progressed: false,
4108 vote_split,
4109 reconsideration,
4110 verdict: round_verdict,
4111 };
4112 if !matches!(
4123 round_record.e2e_status(),
4124 E2eStatus::Deferred | E2eStatus::NotConfigured
4125 ) {
4126 round_record.verified_head = Some(head.clone());
4127 round_record.verified_at = Some(Timestamp::now());
4128 }
4129 let this_round_verification = round_record.verification_summary(&head);
4130
4131 if incomplete {
4132 let missing: Vec<String> = round_record
4133 .reviews
4134 .iter()
4135 .filter(|r| r.failed.is_some())
4136 .map(|r| format!("review-{}", r.reviewer))
4137 .collect();
4138 self.state.event(
4139 "review",
4140 format!(
4141 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4142 missing.join(", ")
4143 ),
4144 );
4145 }
4146
4147 if clean {
4148 self.state.event(
4149 "review",
4150 if incomplete && policy == IncompleteReviewPolicy::Warn {
4151 format!(
4152 "round {round}: clean (warn policy, incomplete panel) — no \
4153 blocking findings from the seats that answered, verification green"
4154 )
4155 } else if incomplete {
4156 format!(
4157 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4158 quorum) — no blocking findings from the seats that answered, \
4159 verification green",
4160 expected - answered
4161 )
4162 } else {
4163 format!("round {round}: clean — no blocking findings, verification green")
4164 },
4165 );
4166 self.state.reviews.push(round_record);
4167 self.state.status = RunStatus::Gating;
4168 self.state.save()?;
4169 return Ok(());
4170 }
4171
4172 if incomplete && blocking == 0 && e2e_ok {
4180 self.state.reviews.push(round_record);
4181 self.state.save()?;
4182 if round == max_rounds {
4183 self.state.status = RunStatus::Blocked;
4184 self.state.event(
4185 "review",
4186 format!(
4187 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4188 refusing to call it clean",
4189 expected - answered
4190 ),
4191 );
4192 return Ok(());
4193 }
4194 continue;
4195 }
4196
4197 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4209 self.state.reviews.push(round_record);
4210 return self
4211 .stop_reviewing(
4212 "the round's own verification could not run",
4213 &shell,
4214 &winner.worktree,
4215 )
4216 .await;
4217 }
4218
4219 if round == max_rounds {
4220 self.state.reviews.push(round_record);
4221 return self
4222 .stop_reviewing(
4223 &format!(
4224 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4225 ),
4226 &shell,
4227 &winner.worktree,
4228 )
4229 .await;
4230 }
4231
4232 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4235 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4236 let blocking_findings: Vec<_> = all_findings
4237 .iter()
4238 .filter(|f| f.severity.blocks())
4239 .cloned()
4240 .collect();
4241 let job = SeatJob {
4242 prompt: prompt::fix(
4243 &self.state.instruction,
4244 &blocking_findings,
4245 this_round_verification.as_ref(),
4246 round,
4247 max_rounds,
4248 &language,
4249 ),
4250 spec: fix_spec.clone(),
4251 seat,
4252 cwd: winner.worktree.clone(),
4253 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4254 allow_write: true,
4255 sessions,
4256 artifacts: artifacts.clone(),
4257 stem: format!("fix-{round}"),
4258 };
4259 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4260 let cache = self.state.config.cache_dir();
4261 let ctx = WaveCtx {
4262 run: &run_id,
4263 node: "fix",
4264 prompts: &prompts,
4265 cache: cache.as_deref(),
4266 round: Some(round),
4267 };
4268 let (seat, out) =
4269 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4270 let agent_id = seat.agent.clone();
4271
4272 let mut fix = FixRecord {
4273 agent: agent_id,
4274 addressed: Vec::new(),
4275 rejected: Vec::new(),
4276 notes: String::new(),
4277 committed: false,
4278 failed: None,
4279 duration_ms: 0,
4280 continuation: None,
4281 };
4282 let mut continuation = ContinuationRecord::not_needed();
4283 let mut final_seat = seat.clone();
4284 match out {
4285 AgentOutcome::Ok(o) => {
4286 fix.duration_ms = o.duration_ms;
4287 let parsed = verdict::extract_json::<FixReport>(&o.text);
4288 let incomplete_reason = match &parsed {
4295 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4296 "the reply parsed, but it reported a command whose own CLI never \
4297 confirmed an exit status"
4298 .to_owned(),
4299 ),
4300 Ok(_) => None,
4301 Err(e) => Some(e.to_string()),
4302 };
4303 match incomplete_reason {
4304 None => {
4305 let report = parsed.expect("checked Ok above");
4306 fix.addressed = report.addressed;
4307 fix.rejected = report.rejected;
4308 fix.notes =
4309 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4310 }
4311 Some(reason) => {
4312 let (resumed_seat, resolved, failure, cont) = self
4313 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4314 .await;
4315 fix.duration_ms += cont.cumulative_wait_ms;
4316 continuation = cont;
4317 final_seat = resumed_seat;
4318 match resolved {
4319 Some(report) => {
4320 fix.addressed = report.addressed;
4321 fix.rejected = report.rejected;
4322 fix.notes = blind::sanitize_prose(
4323 &report.notes,
4324 &self.state.config.blind,
4325 );
4326 }
4327 None => fix.failed = failure,
4328 }
4329 }
4330 }
4331 }
4332 AgentOutcome::Dropped(o) => {
4334 fix.duration_ms = o.duration_ms;
4335 let why = o
4336 .dropped
4337 .as_ref()
4338 .map(|d| d.why.as_str())
4339 .unwrap_or("the CLI ended the stream without delivering its answer");
4340 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4341 }
4342 AgentOutcome::Quota(o) => {
4343 self.state.quota.push(QuotaLoss {
4344 seat: final_seat.key.clone(),
4345 node: "fix".to_owned(),
4346 at: Timestamp::now(),
4347 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4348 });
4349 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4350 }
4351 AgentOutcome::Failed(e) => fix.failed = Some(e),
4352 }
4353 fix.continuation = Some(continuation);
4354 self.state.seats.insert(final_seat.key.clone(), final_seat);
4355 if let Ok(r) = git::rescue_commit(
4356 &winner.worktree,
4357 &format!("magi: review round {round} fixes (uncommitted work)"),
4358 )
4359 .await
4360 {
4361 self.state.note_withheld("fix", &r.withheld);
4362 }
4363 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4364 fix.committed = after != before;
4365 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4373 let progressed = diff_after != patch;
4374 let commit_note = if fix.committed {
4375 "committed"
4376 } else {
4377 "NO new commit"
4378 };
4379 let tree_note = if progressed {
4380 "changed vs base"
4381 } else {
4382 "unchanged vs base"
4383 };
4384 self.state.event(
4385 "fix",
4386 match &fix.failed {
4387 Some(reason) => {
4393 format!(
4394 "round {round}: fixer's adoption report was lost ({reason}); \
4395 {commit_note}, tree {tree_note}"
4396 )
4397 }
4398 None => format!(
4399 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4400 {tree_note}{}",
4401 fix.addressed.len(),
4402 fix.rejected.len(),
4403 if continuation.outcome == ContinuationOutcome::Resumed {
4404 format!(
4405 " (adoption report recovered after {} continuation(s))",
4406 continuation.attempts
4407 )
4408 } else {
4409 String::new()
4410 },
4411 ),
4412 },
4413 );
4414 round_record.fix = Some(fix);
4415 round_record.progressed = progressed;
4416 self.state.reviews.push(round_record);
4417 self.state.save()?;
4418
4419 if matches!(
4432 continuation.outcome,
4433 ContinuationOutcome::Exhausted
4434 | ContinuationOutcome::QuotaLost
4435 | ContinuationOutcome::NoSession
4436 ) {
4437 return self
4438 .stop_reviewing(
4439 "the fixer's adoption report never came back, even after resuming its \
4440 own seat; refusing to start another round against the same worktree \
4441 while that is unresolved",
4442 &shell,
4443 &winner.worktree,
4444 )
4445 .await;
4446 }
4447
4448 let streak = self
4449 .state
4450 .reviews
4451 .iter()
4452 .rev()
4453 .take_while(|r| !r.progressed)
4454 .count();
4455 if streak >= STAGNANT_LIMIT {
4456 return self
4457 .stop_reviewing(
4458 &format!(
4459 "the tree has not moved against base for {streak} round(s) in a row"
4460 ),
4461 &shell,
4462 &winner.worktree,
4463 )
4464 .await;
4465 }
4466 }
4467 Ok(())
4468 }
4469
4470 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4500 let round_idx = self.state.reviews.len() - 1;
4501 let needs_catchup_run = matches!(
4509 self.state.reviews[round_idx].e2e_status(),
4510 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4511 );
4512 if needs_catchup_run {
4513 let round = self.state.reviews[round_idx].round;
4514 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4515 let commands = self.state.config.verify.e2e.clone();
4516 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4517 let cache_dir = self.state.config.cache_dir();
4518 let context = format!(
4519 "round {round}: verification unresolved, catching up before the final decision"
4520 );
4521 let (outcomes, verify_retried) = with_cache_lease(
4522 &mut self.state,
4523 cache_dir.as_deref(),
4524 "e2e",
4525 "e2e",
4526 worktree,
4527 &attempted_head,
4528 timeout,
4529 &context,
4530 |state, budget| {
4531 let shell = shell.to_vec();
4532 let commands = commands.clone();
4533 let context = context.clone();
4534 async move {
4535 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4536 .await
4537 }
4538 },
4539 )
4540 .await;
4541 let last = &mut self.state.reviews[round_idx];
4542 last.e2e = outcomes;
4543 last.verify_retried = verify_retried;
4544 last.verified_head = Some(attempted_head);
4551 last.verified_at = Some(Timestamp::now());
4552 if verify_inconclusive(&last.e2e) {
4553 self.state.save()?;
4560 return Ok(());
4561 }
4562 last.e2e_deferred = false;
4563 }
4564 let last = &self.state.reviews[round_idx];
4565 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4566
4567 match last.e2e_status() {
4568 E2eStatus::Failed => {
4569 let red: Vec<String> = last
4570 .e2e
4571 .iter()
4572 .filter(|o| !o.ok())
4573 .map(|o| {
4574 format!(
4575 "`{}` -> {:?}\n{}",
4576 o.command,
4577 o.code,
4578 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4579 )
4580 })
4581 .collect();
4582 self.state
4583 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4584 self.state.status = RunStatus::Blocked;
4585 }
4586 E2eStatus::ResourceBlocked => {
4591 self.state.event(
4592 "review",
4593 format!(
4594 "{why}; e2e could not run (shared build cache unavailable); not \
4595 deciding yet"
4596 ),
4597 );
4598 }
4599 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4600 self.state.event(
4601 "review",
4602 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4603 );
4604 self.state.status = RunStatus::Gating;
4605 }
4606 }
4607 self.state.save()?;
4608 Ok(())
4609 }
4610
4611 async fn gate(&mut self) -> Result<()> {
4614 if self.state.status == RunStatus::Failed
4626 || self
4627 .state
4628 .base_sync
4629 .as_ref()
4630 .is_some_and(|s| s.conflict.is_some())
4631 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
4632 != Some(RunStatus::Gating)
4633 {
4634 return Ok(());
4635 }
4636 if self.state.gate_ran {
4637 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
4648 self.state.status = RunStatus::Blocked;
4649 self.state.save()?;
4650 }
4651 return Ok(());
4652 }
4653 let Some(winner) = self.state.winner().cloned() else {
4654 return Ok(());
4655 };
4656 self.state.status = RunStatus::Gating;
4657 let mut outcomes = self.run_gate(&winner).await?;
4658 loop {
4659 if verify_inconclusive(&outcomes) {
4670 self.state.save()?;
4671 return Ok(());
4672 }
4673 if outcomes.iter().all(CommandOutcome::ok) {
4674 break;
4675 }
4676 match self.gate_fix_round(&winner, &outcomes).await? {
4677 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
4678 GateFix::Stop => break,
4679 GateFix::Defer => {
4680 self.state.save()?;
4681 return Ok(());
4682 }
4683 }
4684 }
4685 let passed = outcomes.iter().all(CommandOutcome::ok);
4686 self.state.gate = outcomes;
4687 self.state.gate_ran = true;
4688 if !passed {
4689 self.state.status = RunStatus::Blocked;
4690 let spent = self.state.gate_fixes.len();
4691 self.state.event(
4692 "gate",
4693 if spent == 0 {
4694 "gate failed; not merging".to_owned()
4695 } else {
4696 format!("gate failed after {spent} gate-fix round(s); not merging")
4697 },
4698 );
4699 }
4700 self.state.save()?;
4701 Ok(())
4702 }
4703
4704 async fn run_pre_gate(&mut self, winner: &Candidate) {
4714 let commands = self.state.config.verify.pre_gate.clone();
4715 if commands.is_empty() {
4716 return;
4717 }
4718 let shell = self.state.config.shell();
4719 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4720 let (outcomes, _) = run_commands(
4721 &mut self.state,
4722 "pre_gate",
4723 "pre_gate",
4724 0,
4725 &shell,
4726 &commands,
4727 &winner.worktree,
4728 timeout,
4729 )
4730 .await;
4731 for o in &outcomes {
4732 if !o.ok() {
4733 tracing::warn!(
4734 "pre_gate `{}` failed ({:?}); the gate decides",
4735 o.command,
4736 o.code
4737 );
4738 }
4739 self.state.event(
4740 "pre_gate",
4741 format!(
4742 "`{}` -> {}",
4743 o.command,
4744 if o.ok() {
4745 "pass".to_owned()
4746 } else {
4747 format!(
4748 "FAIL ({:?})\n{}",
4749 o.code,
4750 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4751 )
4752 }
4753 ),
4754 );
4755 }
4756 self.state.pre_gate = outcomes;
4757 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
4758 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
4759 Ok(head) => {
4760 self.state
4761 .event("pre_gate", format!("committed mechanical fixes ({head})"));
4762 self.state.pre_gate_commit = Some(head);
4763 }
4764 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
4765 },
4766 Ok(false) => {}
4767 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
4768 }
4769 if let Err(e) = self.state.save() {
4770 tracing::warn!("could not persist the pre_gate record: {e:#}");
4771 }
4772 }
4773
4774 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
4777 self.run_pre_gate(winner).await;
4778 let shell = self.state.config.shell();
4779 let gate_commands = self.state.config.verify.gate.clone();
4780 let outcomes = if gate_commands.is_empty() {
4789 Vec::new()
4790 } else {
4791 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4792 let cache_dir = self.state.config.cache_dir();
4793 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4794 let (outcomes, _) = with_cache_lease(
4795 &mut self.state,
4796 cache_dir.as_deref(),
4797 "gate",
4798 "gate",
4799 &winner.worktree,
4800 &head,
4801 timeout,
4802 "final gate",
4803 |state, budget| {
4804 let shell = shell.clone();
4805 let gate_commands = gate_commands.clone();
4806 let worktree = winner.worktree.clone();
4807 async move {
4808 let (outcomes, timed_out_pids) = run_commands(
4809 state,
4810 "gate",
4811 "gate",
4812 0,
4813 &shell,
4814 &gate_commands,
4815 &worktree,
4816 budget,
4817 )
4818 .await;
4819 (outcomes, false, timed_out_pids)
4820 }
4821 },
4822 )
4823 .await;
4824 outcomes
4825 };
4826 if outcomes.is_empty() {
4827 self.state.event(
4832 "gate",
4833 "no gate commands configured; nothing to check, passing",
4834 );
4835 }
4836 for o in &outcomes {
4837 self.state.event(
4838 "gate",
4839 format!(
4840 "`{}` -> {}",
4841 o.command,
4842 if o.ok() {
4843 "pass".to_owned()
4844 } else {
4845 format!(
4846 "FAIL ({:?})\n{}",
4847 o.code,
4848 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4849 )
4850 }
4851 ),
4852 );
4853 }
4854 Ok(outcomes)
4855 }
4856
4857 async fn gate_fix_round(
4870 &mut self,
4871 winner: &Candidate,
4872 outcomes: &[CommandOutcome],
4873 ) -> Result<GateFix> {
4874 let cap = self.state.config.graph.gate_fix_rounds;
4875 let spent = self.state.gate_fixes.len();
4876 if spent >= cap {
4877 if cap > 0 {
4878 self.state.event(
4879 "gate",
4880 format!("{spent} gate-fix round(s) spent and the gate still fails"),
4881 );
4882 }
4883 return Ok(GateFix::Stop);
4884 }
4885 if !gate_fixable(outcomes) {
4886 self.state.event(
4887 "gate",
4888 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
4889 command or similar); not spending a fix round on it",
4890 );
4891 return Ok(GateFix::Stop);
4892 }
4893 let min_free = self.state.config.disk.min_free_bytes;
4894 if min_free > 0 {
4895 match crate::disk::free_bytes(&winner.worktree) {
4896 Ok(free) if crate::disk::enough_space(free, min_free) => {}
4897 Ok(free) => {
4898 self.state.event(
4899 "gate",
4900 format!(
4901 "only {free} bytes free ({min_free} required by `[disk] \
4902 min_free_bytes`); not spending a fix round on a failure the disk \
4903 may explain"
4904 ),
4905 );
4906 return Ok(GateFix::Stop);
4907 }
4908 Err(e) => {
4909 self.state.event(
4910 "gate",
4911 format!("free disk space could not be measured ({e:#}); no fix round"),
4912 );
4913 return Ok(GateFix::Stop);
4914 }
4915 }
4916 }
4917
4918 let attempt = spent + 1;
4919 let run_id = self.state.id.clone();
4920 let prompts = self.state.config.prompts.clone();
4921 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
4922 let base = self.landing_base();
4923 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
4924 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4925 let job = SeatJob {
4926 prompt: prompt::gate_fix(
4927 &self.state.instruction,
4928 &failed,
4929 attempt,
4930 cap,
4931 &self.state.config.graph.language,
4932 ),
4933 spec: fix_spec,
4934 seat,
4935 cwd: winner.worktree.clone(),
4936 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4937 allow_write: true,
4938 sessions: self.state.config.graph.sessions,
4939 artifacts: agent::artifacts_dir(&self.state.dir()),
4940 stem: format!("gate-fix-{attempt}"),
4941 };
4942 self.state.event(
4943 "gate",
4944 format!("gate failed; gate-fix round {attempt} of {cap}"),
4945 );
4946 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4947 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
4948 let cache = self.state.config.cache_dir();
4949 let ctx = WaveCtx {
4950 run: &run_id,
4951 node: "gate-fix",
4952 prompts: &prompts,
4953 cache: cache.as_deref(),
4954 round: None,
4955 };
4956 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4957 let mut record = GateFixRecord {
4958 agent: seat.agent.clone(),
4959 failed,
4960 notes: String::new(),
4961 committed: false,
4962 error: None,
4963 };
4964 match out {
4965 AgentOutcome::Ok(o) => {
4966 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
4969 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
4970 }
4971 }
4972 AgentOutcome::Dropped(_) => {
4973 record.error = Some("the CLI dropped the stream".to_owned());
4974 }
4975 AgentOutcome::Quota(o) => {
4976 self.state.quota.push(QuotaLoss {
4977 seat: seat.key.clone(),
4978 node: "gate-fix".to_owned(),
4979 at: Timestamp::now(),
4980 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4981 });
4982 record.error = Some("rate limited (quota); fixer could not run".to_owned());
4983 }
4984 AgentOutcome::Failed(e) => record.error = Some(e),
4985 }
4986 self.state.seats.insert(seat.key.clone(), seat);
4987 if let Ok(r) = git::rescue_commit(
4988 &winner.worktree,
4989 &format!("magi: gate fix {attempt} (uncommitted work)"),
4990 )
4991 .await
4992 {
4993 self.state.note_withheld("gate-fix", &r.withheld);
4994 }
4995 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4996 record.committed = after != before;
4997 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
4998 let note = record.error.clone();
4999 self.state.gate_fixes.push(record);
5000 self.state.save()?;
5001 if !changed {
5002 self.state.event(
5003 "gate",
5004 match note {
5005 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5006 None => format!("gate-fix round {attempt}: the tree did not change"),
5007 },
5008 );
5009 return Ok(GateFix::Stop);
5010 }
5011 self.state.event(
5012 "gate",
5013 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5014 );
5015
5016 let commands = self.state.config.verify.e2e.clone();
5017 if !commands.is_empty() {
5018 let shell = self.state.config.shell();
5019 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5020 let cache_dir = self.state.config.cache_dir();
5021 let context = format!("gate-fix round {attempt}");
5022 let (e2e, _) = with_cache_lease(
5023 &mut self.state,
5024 cache_dir.as_deref(),
5025 "e2e",
5026 "e2e",
5027 &winner.worktree,
5028 &after,
5029 timeout,
5030 &context,
5031 |state, budget| {
5032 let shell = shell.clone();
5033 let commands = commands.clone();
5034 let context = context.clone();
5035 let worktree = winner.worktree.clone();
5036 async move {
5037 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5038 .await
5039 }
5040 },
5041 )
5042 .await;
5043 if verify_inconclusive(&e2e) {
5044 return Ok(GateFix::Defer);
5045 }
5046 if e2e.iter().any(|o| !o.ok()) {
5047 self.state.event(
5048 "gate",
5049 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5050 );
5051 return Ok(GateFix::Stop);
5052 }
5053 }
5054 Ok(GateFix::Retry)
5055 }
5056
5057 async fn merge(&mut self) -> Result<()> {
5060 if self
5075 .state
5076 .base_sync
5077 .as_ref()
5078 .is_some_and(|s| s.conflict.is_some())
5079 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5080 != Some(RunStatus::Gating)
5081 || !self.state.gate_status().ok()
5090 {
5091 return Ok(());
5092 }
5093 if self.state.merge.is_some() {
5102 return Ok(());
5103 }
5104 let Some(winner) = self.state.winner().cloned() else {
5105 return Ok(());
5106 };
5107 let repo = self.state.repo.clone();
5108 let base = self.state.base_branch.clone();
5109 let mode = self.state.config.merge.mode;
5110 let style = self.state.config.merge.style;
5111 let pr = pr_message(&self.state, winner.label);
5112 let message = pr.commit_message();
5113
5114 let outcome = match mode {
5115 MergeMode::None => MergeOutcome {
5116 mode,
5117 ok: true,
5118 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5119 },
5120 MergeMode::Local => {
5121 let on = git::current_branch(&repo).await?;
5122 if on.as_deref() != Some(base.as_str()) {
5123 MergeOutcome {
5124 mode,
5125 ok: false,
5126 detail: format!(
5127 "{} has {} checked out, not the base branch {base}",
5128 repo.display(),
5129 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5130 ),
5131 }
5132 } else if !git::is_clean(&repo).await? {
5133 MergeOutcome {
5134 mode,
5135 ok: false,
5136 detail: format!("{} is dirty; refusing to merge", repo.display()),
5137 }
5138 } else {
5139 let out = match style {
5140 MergeStyle::Merge => {
5141 git::merge_no_ff(&repo, &winner.branch, &message).await?
5142 }
5143 MergeStyle::Squash => {
5144 git::merge_squash(&repo, &winner.branch, &message).await?
5145 }
5146 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5147 };
5148 MergeOutcome {
5149 mode,
5150 ok: out.ok(),
5151 detail: if out.ok() { out.stdout } else { out.stderr },
5152 }
5153 }
5154 }
5155 MergeMode::Pr => {
5156 let remote = self.state.config.merge.remote.clone();
5157 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5158 if !pushed.ok() {
5159 MergeOutcome {
5160 mode,
5161 ok: false,
5162 detail: pushed.stderr,
5163 }
5164 } else {
5165 let out =
5166 gh_pr_create(&winner.worktree, &base, &winner.branch, &pr.title, &pr.body)
5167 .await;
5168 match out {
5169 Ok(url) => MergeOutcome {
5170 mode,
5171 ok: true,
5172 detail: url,
5173 },
5174 Err(e) => MergeOutcome {
5175 mode,
5176 ok: false,
5177 detail: e.to_string(),
5178 },
5179 }
5180 }
5181 }
5182 };
5183
5184 self.state.status = match (mode, outcome.ok) {
5185 (MergeMode::None, _) => RunStatus::Ready,
5186 (_, true) => RunStatus::Merged,
5187 (_, false) => RunStatus::Blocked,
5188 };
5189 self.state.event(
5190 "merge",
5191 format!(
5192 "{:?}: {}",
5193 mode,
5194 outcome.detail.lines().next().unwrap_or("")
5195 ),
5196 );
5197 self.state.merge = Some(outcome);
5198 self.state.save()?;
5199
5200 if self.state.config.graph.land
5206 && mode == MergeMode::Pr
5207 && self.state.status == RunStatus::Merged
5208 {
5209 self.run_land().await?;
5210 }
5211 self.settle_questions();
5216 Ok(())
5217 }
5218
5219 async fn run_land(&mut self) -> Result<()> {
5230 let url = self
5231 .state
5232 .merge
5233 .as_ref()
5234 .map(|m| m.detail.clone())
5235 .unwrap_or_default();
5236 let url = url.lines().next().unwrap_or("").trim().to_owned();
5237 if !url.starts_with("http") {
5238 return Ok(());
5239 }
5240 match land::land(&mut self.state, &url).await {
5243 Ok(pr) if self.state.parked => {
5244 let _ = pr;
5248 }
5249 Ok(pr) => {
5250 self.state.status = match pr.state {
5251 land::PrLifecycle::Merged => RunStatus::Merged,
5252 _ => RunStatus::Blocked,
5253 };
5254 if bump::should_release_bump(self.state.status)
5261 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5262 {
5263 self.state
5269 .event("bump", format!("release bump skipped: {e:#}"));
5270 }
5271 self.state.save()?;
5272 }
5273 Err(e) => {
5274 self.state.status = RunStatus::Blocked;
5275 self.state.event("land", format!("gave up: {e}"));
5276 self.state.save()?;
5277 }
5278 }
5279 Ok(())
5280 }
5281
5282 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5286 if let Some(existing) = self.state.seats.get(key)
5287 && existing.agent == agent
5288 {
5289 return existing.clone();
5290 }
5291 let fresh = SeatState::new(key, agent, self.state.seed);
5292 self.state.seats.insert(key.to_owned(), fresh.clone());
5293 fresh
5294 }
5295
5296 fn view(&self, c: &Candidate) -> CandidateView {
5298 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5299 .unwrap_or_default();
5300 let (patch, _) = blind::sanitize_patch(
5301 &format!("candidate {} patch", c.label),
5302 &raw,
5303 &self.state.config.blind,
5304 );
5305 CandidateView {
5306 label: c.label,
5307 branch: c.branch.clone(),
5308 summary: c.summary.clone(),
5309 stat: c.stat.clone(),
5310 patch,
5311 }
5312 }
5313
5314 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5316 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5317 prompt::judge(
5318 "(see above)",
5319 &views,
5320 self.roles.judges.len(),
5321 base_short,
5322 "en",
5323 )
5324 }
5325
5326 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5333 let mut turns = Vec::new();
5334 for j in &self.state.judgements {
5335 if j.ranking.is_empty() {
5336 continue;
5337 }
5338 let reasons = j
5339 .reasons
5340 .iter()
5341 .map(|(k, v)| format!("- {k}: {v}"))
5342 .collect::<Vec<_>>()
5343 .join("\n");
5344 turns.push(Turn {
5345 who: format!("Judge {} (opening ranking)", j.judge),
5346 is_self: j.judge == self_idx + 1,
5347 body: format!(
5348 "Ranked {}{}{reasons}",
5349 j.ranking.iter().collect::<String>(),
5350 if reasons.is_empty() {
5351 ""
5352 } else {
5353 ", because:\n"
5354 }
5355 ),
5356 });
5357 }
5358 for t in self
5359 .state
5360 .deliberation
5361 .iter()
5362 .flat_map(|r| r.turns.iter())
5363 .chain(current)
5364 {
5365 turns.push(Turn {
5366 who: format!("Judge {}", t.judge),
5367 is_self: t.judge == self_idx + 1,
5368 body: t.body.clone(),
5369 });
5370 }
5371 turns
5372 }
5373}
5374
5375fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5377 agent::has_session(spec.kind, seat, sessions)
5378}
5379
5380fn next_untried_implementer<'a>(
5401 roster: &'a [AgentSpec],
5402 start: usize,
5403 tried: &BTreeSet<String>,
5404) -> Option<&'a AgentSpec> {
5405 roster
5406 .get(start + 1..)?
5407 .iter()
5408 .find(|s| !tried.contains(&s.id))
5409}
5410
5411fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5425 commands.iter().any(|c| c.exit_code.is_none())
5426}
5427
5428fn verified_noop_claim(
5441 usable: bool,
5442 commands: &[agent::CommandEvidence],
5443 text: &str,
5444) -> Option<String> {
5445 (usable && !has_unconfirmed_command(commands))
5446 .then(|| verdict::verified_noop(text))
5447 .flatten()
5448}
5449
5450fn short(commit: &str) -> String {
5451 commit.chars().take(7).collect()
5452}
5453
5454fn make_executable(path: &Path) -> Result<()> {
5455 #[cfg(unix)]
5456 {
5457 use std::os::unix::fs::PermissionsExt as _;
5458 let mut perms = std::fs::metadata(path)?.permissions();
5459 perms.set_mode(0o755);
5460 std::fs::set_permissions(path, perms)?;
5461 }
5462 #[cfg(not(unix))]
5463 {
5464 let _ = path;
5465 }
5466 Ok(())
5467}
5468
5469struct WaveCtx<'a> {
5476 run: &'a str,
5479 node: &'a str,
5481 prompts: &'a Prompts,
5482 cache: Option<&'a Path>,
5484 round: Option<usize>,
5487}
5488
5489async fn run_one(
5491 job: SeatJob,
5492 sem: Arc<Semaphore>,
5493 ctx: &WaveCtx<'_>,
5494 state: &mut RunState,
5495 attempt: usize,
5496) -> (SeatState, AgentOutcome) {
5497 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5498 .await
5499 .pop()
5500 .expect("one job in, one result out");
5501 (seat, out)
5502}
5503
5504async fn wave(
5510 jobs: Vec<SeatJob>,
5511 sem: Arc<Semaphore>,
5512 ctx: &WaveCtx<'_>,
5513 state: &mut RunState,
5514 attempt: usize,
5515) -> Vec<(usize, SeatState, AgentOutcome)> {
5516 let WaveCtx {
5517 run,
5518 node,
5519 prompts,
5520 cache,
5521 round,
5522 } = *ctx;
5523 for job in &jobs {
5524 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5525 }
5526 if let Err(e) = state.save() {
5527 tracing::warn!("could not persist in-progress seats: {e:#}");
5532 }
5533 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5549 let wait_started = Instant::now();
5550 let cache_guard = if let Some(cache_dir) = cache {
5551 if jobs_had_a_writer {
5552 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5553 let budget = jobs
5554 .iter()
5555 .map(|j| j.timeout)
5556 .max()
5557 .unwrap_or(Duration::from_secs(60));
5558 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5559 .await
5560 .ok()
5561 } else {
5562 None
5563 }
5564 } else {
5565 None
5566 };
5567 let waited_for_lease = wait_started.elapsed();
5574 let mut set = tokio::task::JoinSet::new();
5575 let overlay = prompts.overlay(node);
5576 for (i, mut job) in jobs.into_iter().enumerate() {
5577 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5578 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
5579 if cache.is_some() {
5580 job.prompt.push('\n');
5581 job.prompt
5582 .push_str(&prompt::build_cache_note(node, job.allow_write));
5583 }
5584 let sem = Arc::clone(&sem);
5585 let run = run.to_owned();
5586 let node = node.to_owned();
5587 let cache = cache
5598 .filter(|_| job.allow_write && cache_guard.is_some())
5599 .map(Path::to_path_buf);
5600 set.spawn(async move {
5601 let _permit = sem.acquire().await;
5602 let mut seat = job.seat;
5603 let out = agent::invoke(
5604 &job.spec,
5605 &mut seat,
5606 &Invocation {
5607 cwd: &job.cwd,
5608 prompt: &job.prompt,
5609 timeout: job.timeout,
5610 allow_write: job.allow_write,
5611 sessions: job.sessions,
5612 artifacts: &job.artifacts,
5613 stem: &job.stem,
5614 run: &run,
5615 node: &node,
5616 cache_dir: cache.as_deref(),
5617 attachments: &[],
5618 },
5619 )
5620 .await;
5621 let out = match out {
5622 Ok(o) if o.usable() => AgentOutcome::Ok(o),
5623 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
5624 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
5632 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
5633 Ok(o) => AgentOutcome::Failed(format!(
5634 "exited with {:?} and no usable output",
5635 o.exit_code
5636 )),
5637 Err(e) => AgentOutcome::Failed(e.to_string()),
5638 };
5639 (i, seat, out)
5640 });
5641 }
5642 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
5643 while let Some(joined) = set.join_next().await {
5644 let (i, seat, out) = match joined {
5645 Ok(v) => v,
5646 Err(e) => {
5650 tracing::error!("agent task panicked: {e}");
5651 continue;
5652 }
5653 };
5654 state.seat_finished(&seat.key);
5655 record_jobs(state, node, round, &seat.key, &out);
5656 if let Err(e) = state.save() {
5657 tracing::warn!("could not persist a seat's completion: {e:#}");
5658 }
5659 if collected.len() <= i {
5660 collected.resize_with(i + 1, || None);
5661 }
5662 collected[i] = Some((i, seat, out));
5663 }
5664 if state
5670 .active
5671 .values()
5672 .any(|a| a.node == node && a.attempt == attempt)
5673 {
5674 state
5675 .active
5676 .retain(|_, a| !(a.node == node && a.attempt == attempt));
5677 if let Err(e) = state.save() {
5678 tracing::warn!("could not persist the end of a wave: {e:#}");
5679 }
5680 }
5681 if let Some(cache_dir) = cache
5688 && jobs_had_a_writer
5689 {
5690 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
5691 }
5692 if let Some(guard) = cache_guard {
5693 guard.release();
5694 }
5695 collected.into_iter().flatten().collect()
5696}
5697
5698fn record_jobs(
5709 state: &mut RunState,
5710 node: &str,
5711 round: Option<usize>,
5712 seat: &str,
5713 out: &AgentOutcome,
5714) {
5715 let commands: &[agent::CommandEvidence] = match out {
5716 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
5717 AgentOutcome::Failed(_) => &[],
5718 };
5719 let checked_at = Timestamp::now();
5720 for c in commands {
5721 state.jobs.push(JobRecord {
5722 node: node.to_owned(),
5723 round,
5724 seat: seat.to_owned(),
5725 id: c.id.clone(),
5726 description: c.description.clone(),
5727 checked_at,
5728 status: match c.exit_code {
5729 Some(0) => JobStatus::Completed,
5730 Some(_) => JobStatus::Failed,
5731 None => JobStatus::Unknown,
5732 },
5733 exit_code: c.exit_code,
5734 result_summary: c.result_summary.clone(),
5735 source: c.source.clone(),
5736 });
5737 }
5738}
5739
5740fn round_is_clean(
5761 blocking: usize,
5762 e2e_ok: bool,
5763 answered: usize,
5764 expected: usize,
5765 quota_missing: usize,
5766 policy: IncompleteReviewPolicy,
5767) -> bool {
5768 if blocking != 0 || !e2e_ok {
5769 return false;
5770 }
5771 if answered == expected || policy == IncompleteReviewPolicy::Warn {
5772 return true;
5773 }
5774 answered > 0 && expected - answered <= quota_missing
5775}
5776
5777fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
5801 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
5802 return Some(RunStatus::Gating);
5803 }
5804 let last = reviews.last()?;
5805 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
5806 if reviews.len() < max_rounds && !stagnant {
5807 return None;
5808 }
5809 if last.incomplete() && last.blocking == 0 {
5810 return Some(RunStatus::Blocked);
5811 }
5812 if last.e2e_status() == E2eStatus::ResourceBlocked {
5813 return None;
5814 }
5815 Some(if last.e2e.iter().all(CommandOutcome::ok) {
5816 RunStatus::Gating
5817 } else {
5818 RunStatus::Blocked
5819 })
5820}
5821
5822fn retry_budget(full: Duration, nudged: bool) -> Duration {
5837 if nudged {
5838 (full / 4).max(Duration::from_secs(120)).min(full)
5839 } else {
5840 full
5841 }
5842}
5843
5844#[allow(clippy::too_many_arguments)]
5857async fn ask_json_wave<T>(
5858 jobs: Vec<SeatJob>,
5859 sem: Arc<Semaphore>,
5860 retries: usize,
5861 ctx: &WaveCtx<'_>,
5862 losses: &mut Vec<QuotaLoss>,
5863 state: &mut RunState,
5864 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
5865) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
5866where
5867 T: serde::de::DeserializeOwned + Send + 'static,
5868{
5869 let n = jobs.len();
5870 let originals: Vec<SeatJob> = jobs;
5871 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
5872 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
5873 let mut attempts_used: Vec<usize> = vec![0; n];
5880 let mut pending: Vec<usize> = (0..n).collect();
5881
5882 for attempt in 0..=retries {
5883 if pending.is_empty() {
5884 break;
5885 }
5886 let mut batch = Vec::with_capacity(pending.len());
5887 for &i in &pending {
5888 let src = &originals[i];
5889 let (prompt, timeout) = if attempt == 0 {
5892 (src.prompt.clone(), src.timeout)
5893 } else {
5894 let why = done[i]
5895 .as_ref()
5896 .and_then(|r| r.as_ref().err().map(ToString::to_string))
5897 .unwrap_or_else(|| "no parsable answer".to_owned());
5898 let nudge = prompt::nudge(&why);
5899 let nudged = has_context(&src.spec, &seats[i], src.sessions);
5900 let prompt = if nudged {
5901 nudge
5902 } else {
5903 format!("{}\n\n---\n\n{}", src.prompt, nudge)
5904 };
5905 (prompt, retry_budget(src.timeout, nudged))
5906 };
5907 batch.push(SeatJob {
5908 spec: src.spec.clone(),
5909 seat: seats[i].clone(),
5910 cwd: src.cwd.clone(),
5911 prompt,
5912 timeout,
5913 allow_write: src.allow_write,
5914 sessions: src.sessions,
5915 artifacts: src.artifacts.clone(),
5916 stem: if attempt == 0 {
5917 src.stem.clone()
5918 } else {
5919 format!("{}-retry{attempt}", src.stem)
5920 },
5921 });
5922 }
5923
5924 if attempt > 0 {
5925 let seats_out: Vec<&str> = pending
5926 .iter()
5927 .map(|&i| originals[i].seat.key.as_str())
5928 .collect();
5929 state.event(
5930 ctx.node,
5931 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
5932 );
5933 }
5934 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
5935 let mut still = Vec::new();
5936 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
5937 seats[i] = seat;
5938 let (parsed, quota) = match out {
5939 AgentOutcome::Ok(o) => (
5940 match verdict::extract_json::<T>(&o.text) {
5941 Ok(v) => match validate(&v) {
5942 Ok(()) => Ok((v, o)),
5943 Err(e) => Err(e),
5944 },
5945 Err(e) => Err(e),
5946 },
5947 false,
5948 ),
5949 AgentOutcome::Quota(o) => {
5950 losses.push(QuotaLoss {
5951 seat: originals[i].seat.key.clone(),
5952 node: ctx.node.to_owned(),
5953 at: Timestamp::now(),
5954 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5955 });
5956 (
5957 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
5958 true,
5959 )
5960 }
5961 AgentOutcome::Dropped(o) => {
5966 let why = o
5967 .dropped
5968 .as_ref()
5969 .map(|d| d.why.as_str())
5970 .unwrap_or("the CLI ended the stream without delivering its answer");
5971 (
5972 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
5973 false,
5974 )
5975 }
5976 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
5977 };
5978 let failed = parsed.is_err();
5979 done[i] = Some(parsed);
5980 attempts_used[i] = attempt;
5981 if failed && !quota {
5984 still.push(i);
5985 }
5986 }
5987 pending = still;
5988 }
5989
5990 seats
5991 .into_iter()
5992 .zip(done)
5993 .zip(attempts_used)
5994 .map(|((seat, res), attempts)| {
5995 (
5996 seat,
5997 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
5998 attempts,
5999 )
6000 })
6001 .collect()
6002}
6003
6004async fn acquire_cache_lease(
6017 state: &mut RunState,
6018 cache_dir: &Path,
6019 owner: &crate::cache::Owner,
6020 budget: Duration,
6021 context: &str,
6022) -> Result<crate::cache::Guard> {
6023 let home = crate::run::home();
6024 let started = Instant::now();
6025 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6026 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6027 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6028 Err(e) => {
6029 state.event(
6030 "verify",
6031 format!("{context}: could not check the shared build cache: {e:#}"),
6032 );
6033 if let Err(e2) = state.save() {
6034 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6035 }
6036 return Err(e);
6037 }
6038 };
6039 state.event(
6040 "verify",
6041 format!(
6042 "{context}: waiting for the shared build cache at {} ({})",
6043 cache_dir.display(),
6044 busy.describe()
6045 ),
6046 );
6047 if let Err(e) = state.save() {
6048 tracing::warn!("could not persist a cache wait: {e:#}");
6049 }
6050 let remaining = budget.saturating_sub(started.elapsed());
6051 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6052 Ok(g) => Ok(g),
6053 Err(e) => {
6054 state.event("verify", format!("{context}: {e:#}"));
6055 if let Err(e2) = state.save() {
6056 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6057 }
6058 Err(e)
6059 }
6060 }
6061}
6062
6063#[allow(clippy::too_many_arguments)]
6084async fn with_cache_lease<'s, F, Fut>(
6085 state: &'s mut RunState,
6086 cache_dir: Option<&Path>,
6087 node: &str,
6088 seat: &str,
6089 worktree: &Path,
6090 head: &str,
6091 budget: Duration,
6092 context: &str,
6093 body: F,
6094) -> (Vec<CommandOutcome>, bool)
6095where
6096 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6097 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6098{
6099 let Some(cache_dir) = cache_dir else {
6100 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6101 return (outcomes, retried);
6102 };
6103 let home = crate::run::home();
6104 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6105 let started = Instant::now();
6106 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6107 Ok(g) => g,
6108 Err(e) => {
6109 return (
6110 vec![CommandOutcome {
6111 command: "(waiting for the shared build cache)".to_owned(),
6112 code: None,
6113 output_tail: e.to_string(),
6114 duration_ms: started.elapsed().as_millis() as u64,
6115 resource_blocked: true,
6116 }],
6117 false,
6118 );
6119 }
6120 };
6121 let identity = crate::cache::Identity::new(worktree, head);
6122 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6123 state.event(
6131 "verify",
6132 format!(
6133 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6134 worktree.display(),
6135 short(head)
6136 ),
6137 );
6138 guard.release();
6139 return (
6140 vec![CommandOutcome {
6141 command: "(confirming the shared build cache is fresh)".to_owned(),
6142 code: None,
6143 output_tail: e.to_string(),
6144 duration_ms: started.elapsed().as_millis() as u64,
6145 resource_blocked: true,
6146 }],
6147 false,
6148 );
6149 }
6150 let remaining = budget.saturating_sub(started.elapsed());
6151 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6152 if !timed_out_pids.is_empty() {
6157 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6158 }
6159 guard.release();
6160 (outcomes, retried)
6161}
6162
6163async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6175 wait_for_pids_with(
6176 pids,
6177 crate::proc::pid_alive,
6178 LEASE_RELEASE_POLL,
6179 LEASE_RELEASE_MAX_WAIT,
6180 )
6181 .await;
6182}
6183
6184async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6190 pids: &[u32],
6191 alive: F,
6192 poll: Duration,
6193 max_wait: Duration,
6194) {
6195 let deadline = Instant::now() + max_wait;
6196 loop {
6197 if pids.iter().all(|&pid| !alive(pid)) {
6198 return;
6199 }
6200 if Instant::now() >= deadline {
6201 return;
6202 }
6203 tokio::time::sleep(poll).await;
6204 }
6205}
6206
6207fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6214 outcomes.iter().any(|o| o.resource_blocked)
6215}
6216
6217enum GateFix {
6219 Retry,
6221 Stop,
6224 Defer,
6227}
6228
6229fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6237 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6238 red.peek().is_some()
6239 && red.all(|o| {
6240 !o.resource_blocked
6241 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6242 && !o.output_tail.trim().is_empty()
6243 })
6244}
6245
6246fn e2e_outcome_label(o: &CommandOutcome) -> String {
6250 if o.ok() {
6251 return "pass".to_owned();
6252 }
6253 let reason = if o.build_failed() {
6254 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6255 } else {
6256 format!("FAIL ({:?})", o.code)
6257 };
6258 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6259}
6260
6261async fn run_e2e_with_retry(
6269 state: &mut RunState,
6270 shell: &[String],
6271 commands: &[String],
6272 worktree: &Path,
6273 timeout: Duration,
6274 context: &str,
6275) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6276 let (mut e2e, mut timed_out_pids) = run_commands(
6277 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6278 )
6279 .await;
6280 for o in &e2e {
6281 state.event(
6282 "verify",
6283 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6284 );
6285 }
6286 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6290 if verify_retried {
6291 state.event(
6292 "verify",
6293 format!(
6294 "{context}: verify could not build/link, not a test result — retrying once \
6295 before concluding"
6296 ),
6297 );
6298 let retried = run_commands(
6299 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6300 )
6301 .await;
6302 e2e = retried.0;
6303 timed_out_pids.extend(retried.1);
6306 for o in &e2e {
6307 state.event(
6308 "verify",
6309 format!(
6310 "{context}: retry `{}` -> {}",
6311 o.command,
6312 e2e_outcome_label(o)
6313 ),
6314 );
6315 }
6316 }
6317 (e2e, verify_retried, timed_out_pids)
6318}
6319
6320#[allow(clippy::too_many_arguments)]
6336async fn run_commands(
6337 state: &mut RunState,
6338 node: &str,
6339 task: &str,
6340 attempt: usize,
6341 shell: &[String],
6342 commands: &[String],
6343 cwd: &Path,
6344 timeout: Duration,
6345) -> (Vec<CommandOutcome>, Vec<u32>) {
6346 if commands.is_empty() {
6347 return (Vec::new(), Vec::new());
6352 }
6353 let mut out = Vec::new();
6354 let mut timed_out_pids = Vec::new();
6355 let total = commands.len();
6356 for (idx, command) in commands.iter().enumerate() {
6357 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6358 if let Err(e) = state.save() {
6359 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6360 }
6361 let started = Instant::now();
6362 let mut cmd = tokio::process::Command::new(&shell[0]);
6363 cmd.quiet();
6364 cmd.args(&shell[1..])
6365 .arg(command)
6366 .current_dir(cwd)
6367 .stdin(std::process::Stdio::null())
6368 .stdout(std::process::Stdio::piped())
6369 .stderr(std::process::Stdio::piped())
6370 .kill_on_drop(true);
6371 let spawned = cmd.spawn();
6372 let (code, body) = match spawned {
6373 Ok(child) => {
6374 let pid = child.id();
6379 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6380 Ok(Ok(o)) => {
6381 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6382 body.push_str(&String::from_utf8_lossy(&o.stderr));
6383 (o.status.code(), body)
6384 }
6385 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6386 Err(_) => {
6387 if let Some(pid) = pid {
6388 timed_out_pids.push(pid);
6389 }
6390 (None, format!("timed out after {}s", timeout.as_secs()))
6391 }
6392 }
6393 }
6394 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6395 };
6396 out.push(CommandOutcome {
6397 command: command.clone(),
6398 code,
6399 output_tail: tail(&body, OUTPUT_TAIL),
6400 duration_ms: started.elapsed().as_millis() as u64,
6401 resource_blocked: false,
6402 });
6403 }
6404 state.task_finished(task);
6405 if let Err(e) = state.save() {
6406 tracing::warn!("could not persist the end of {task}: {e:#}");
6407 }
6408 (out, timed_out_pids)
6409}
6410
6411fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6423 let repo = repo.display();
6424 match style {
6425 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6426 MergeStyle::Squash => {
6427 let subject = message
6430 .lines()
6431 .next()
6432 .unwrap_or(branch)
6433 .replace(['\\', '"', '$', '`'], "");
6434 format!(
6435 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6436 )
6437 }
6438 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6439 }
6440}
6441
6442const PR_TITLE_MAX: usize = 240;
6453
6454struct PrMessage {
6459 title: String,
6460 body: String,
6461}
6462
6463impl PrMessage {
6464 fn commit_message(&self) -> String {
6468 format!("{}\n\n{}", self.title, self.body)
6469 }
6470}
6471
6472fn title_marker(line: &str) -> Option<&str> {
6474 let line = line.trim();
6475 let head = line.get(..6)?;
6476 head.eq_ignore_ascii_case("title:")
6477 .then(|| line[6..].trim())
6478}
6479
6480fn summary_title(summary: &str) -> Option<String> {
6485 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6486 let raw = title_marker(first)?;
6487 if raw.is_empty() {
6488 return None;
6489 }
6490 let title = queue::title_from(raw, PR_TITLE_MAX);
6491 let lower = title.to_ascii_lowercase();
6492 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6493 return None;
6494 }
6495 Some(title)
6496}
6497
6498fn summary_without_title(summary: &str) -> String {
6501 let mut lines = summary.trim().lines().peekable();
6502 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6503 lines.next();
6504 }
6505 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6506}
6507
6508fn pr_message(state: &RunState, winner: char) -> PrMessage {
6522 let summary = state
6523 .candidates
6524 .iter()
6525 .find(|c| c.label == winner)
6526 .map(|c| c.summary.as_str())
6527 .unwrap_or_default();
6528 let title = summary_title(summary).unwrap_or_else(|| {
6531 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6532 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6533 t
6534 } else {
6535 format!(
6536 "chore: land candidate {} of run {}",
6537 winner.to_ascii_uppercase(),
6538 state.id
6539 )
6540 }
6541 });
6542
6543 let mut body = String::new();
6544 let what = summary_without_title(summary);
6545 if !what.is_empty() {
6546 body.push_str("## Summary\n\n");
6547 body.push_str(&what);
6548 body.push_str("\n\n");
6549 }
6550
6551 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6552 if let Some(fix) = fix
6553 && !fix.notes.trim().is_empty()
6554 {
6555 body.push_str("## Review fixes\n\n");
6556 body.push_str(fix.notes.trim());
6557 body.push_str("\n\n");
6558 }
6559
6560 let open = state.open_findings();
6561 if !open.is_empty() {
6562 body.push_str("## Open review findings\n\n");
6563 for f in &open {
6564 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6565 }
6566 body.push('\n');
6567 }
6568
6569 if let Some(fix) = fix
6570 && !fix.rejected.is_empty()
6571 {
6572 body.push_str("## Declined by the fixer\n\n");
6573 for r in &fix.rejected {
6574 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
6575 }
6576 body.push('\n');
6577 }
6578
6579 let task = state.instruction.trim();
6580 let task = if task.is_empty() {
6581 "(empty task)"
6582 } else {
6583 task
6584 };
6585 body.push_str(&format!(
6586 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
6587 task.replace("</details>", "</details>")
6588 ));
6589
6590 body.push_str(&format!(
6591 "\n---\nmagi:run/{} magi:candidate-{}\n",
6592 state.id,
6593 winner.to_ascii_lowercase()
6594 ));
6595
6596 let id = crate::scrub::Identity::current();
6599 PrMessage {
6600 title: crate::scrub::scrub(&title, &id),
6601 body: crate::scrub::scrub(&body, &id),
6602 }
6603}
6604
6605async fn gh_pr_create(
6607 cwd: &Path,
6608 base: &str,
6609 head: &str,
6610 title: &str,
6611 body: &str,
6612) -> Result<String> {
6613 let out = tokio::process::Command::new("gh")
6614 .args([
6615 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
6616 ])
6617 .current_dir(cwd)
6618 .quiet()
6619 .stdin(std::process::Stdio::null())
6620 .output()
6621 .await
6622 .context("spawn gh")?;
6623 if out.status.success() {
6624 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
6625 } else {
6626 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
6627 }
6628}
6629
6630pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
6639 let repo = state.repo.clone();
6640 let root = state.worktree_root();
6641 let winner = state.tally.as_ref().map(|t| t.winner);
6642 let mut removed = Vec::new();
6643
6644 for i in 0..state.candidates.len() {
6645 let c = state.candidates[i].clone();
6646 let is_winner = Some(c.label) == winner;
6647 if is_winner && !drop_winner {
6648 continue;
6649 }
6650 if c.worktree.exists() {
6651 git::worktree_remove(&repo, &c.worktree).await.ok();
6652 removed.push(c.worktree.to_string_lossy().into_owned());
6653 }
6654 if git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
6655 git::branch_delete(&repo, &c.branch).await.ok();
6656 removed.push(c.branch.clone());
6657 }
6658 state.candidates[i].folded = true;
6659 }
6660
6661 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
6662 let path = name.path();
6663 let keep = !drop_winner
6664 && winner.is_some_and(|w| {
6665 path.file_name()
6666 .is_some_and(|n| n == format!("cand-{w}").as_str())
6667 });
6668 if keep {
6669 continue;
6670 }
6671 git::worktree_remove(&repo, &path).await.ok();
6672 removed.push(path.to_string_lossy().into_owned());
6673 }
6674
6675 remove_if_empty(&root);
6684
6685 if state.enabled_worktree_config && drop_winner {
6686 git::release_worktree_config(&repo).await.ok();
6690 state.enabled_worktree_config = false;
6691 }
6692 state.save_under(home)?;
6693 Ok(removed)
6694}
6695
6696fn remove_if_empty(dir: &Path) {
6707 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
6708 std::fs::remove_dir(dir).ok();
6709 }
6710}
6711
6712pub fn worst_open(state: &RunState) -> Option<Severity> {
6714 state
6715 .reviews
6716 .last()?
6717 .reviews
6718 .iter()
6719 .flat_map(|r| r.findings.iter())
6720 .map(|f| f.severity)
6721 .max()
6722}
6723
6724#[cfg(test)]
6725mod tests {
6726 use super::*;
6727 use crate::run::GateStatus;
6728 use std::collections::BTreeMap;
6729 use std::time::Duration;
6730
6731 fn conductor() -> AgentSpec {
6732 AgentSpec {
6733 id: "conductor".to_owned(),
6734 kind: crate::config::AgentKind::Command,
6735 model: None,
6736 command: vec!["true".to_owned()],
6737 extra_args: Vec::new(),
6738 env: BTreeMap::new(),
6739 prompt_delivery: None,
6740 }
6741 }
6742
6743 fn spec(id: &str) -> AgentSpec {
6744 AgentSpec {
6745 id: id.to_owned(),
6746 kind: crate::config::AgentKind::Command,
6747 model: None,
6748 command: vec!["true".to_owned()],
6749 extra_args: Vec::new(),
6750 env: BTreeMap::new(),
6751 prompt_delivery: None,
6752 }
6753 }
6754
6755 #[test]
6762 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
6763 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6764 let tried = BTreeSet::from(["beta".to_owned()]);
6765 let next = next_untried_implementer(&roster, 1, &tried);
6768 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
6769 }
6770
6771 #[test]
6772 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
6773 let roster = vec![spec("alpha"), spec("beta")];
6774 let tried = BTreeSet::from(["beta".to_owned()]);
6775 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6779 }
6780
6781 #[test]
6782 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
6783 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6784 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
6785 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6789 }
6790
6791 #[test]
6792 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
6793 let roster = vec![spec("a"), spec("a"), spec("b")];
6794 let tried = BTreeSet::from(["a".to_owned()]);
6795 let next = next_untried_implementer(&roster, 0, &tried);
6796 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
6797 }
6798
6799 #[test]
6800 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
6801 let roster = vec![spec("a"), spec("b")];
6802 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
6803 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
6804 }
6805
6806 #[test]
6807 fn remove_if_empty_only_ever_takes_a_bare_directory() {
6808 let dir = tempfile::tempdir().unwrap();
6809 let bay = dir.path().join("ffff");
6810
6811 remove_if_empty(&bay);
6813 assert!(!bay.exists());
6814
6815 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
6818 remove_if_empty(&bay);
6819 assert!(bay.exists(), "non-empty directory must survive");
6820
6821 std::fs::remove_dir(bay.join("cand-A")).unwrap();
6823 remove_if_empty(&bay);
6824 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
6825 }
6826
6827 #[test]
6836 fn a_full_panel_that_found_nothing_is_clean() {
6837 assert!(round_is_clean(
6838 0,
6839 true,
6840 2,
6841 2,
6842 0,
6843 IncompleteReviewPolicy::Block
6844 ));
6845 }
6846
6847 #[test]
6848 fn a_missing_seat_is_never_clean_under_the_default_policy() {
6849 assert!(!round_is_clean(
6850 0,
6851 true,
6852 1,
6853 2,
6854 0,
6855 IncompleteReviewPolicy::Block
6856 ));
6857 }
6858
6859 #[test]
6860 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
6861 assert!(!round_is_clean(
6862 1,
6863 true,
6864 1,
6865 2,
6866 0,
6867 IncompleteReviewPolicy::Warn
6868 ));
6869 }
6870
6871 #[test]
6872 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
6873 assert!(round_is_clean(
6874 0,
6875 true,
6876 1,
6877 2,
6878 0,
6879 IncompleteReviewPolicy::Warn
6880 ));
6881 }
6882
6883 #[test]
6884 fn a_full_panel_with_an_open_finding_is_not_clean() {
6885 assert!(!round_is_clean(
6886 1,
6887 true,
6888 2,
6889 2,
6890 0,
6891 IncompleteReviewPolicy::Block
6892 ));
6893 }
6894
6895 #[test]
6896 fn a_full_panel_with_a_red_e2e_is_not_clean() {
6897 assert!(!round_is_clean(
6898 0,
6899 false,
6900 2,
6901 2,
6902 0,
6903 IncompleteReviewPolicy::Block
6904 ));
6905 }
6906
6907 #[test]
6914 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
6915 assert!(round_is_clean(
6918 0,
6919 true,
6920 1,
6921 2,
6922 1,
6923 IncompleteReviewPolicy::Block
6924 ));
6925 }
6926
6927 #[test]
6928 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
6929 assert!(!round_is_clean(
6932 0,
6933 true,
6934 1,
6935 2,
6936 0,
6937 IncompleteReviewPolicy::Block
6938 ));
6939 }
6940
6941 #[test]
6942 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
6943 assert!(!round_is_clean(
6944 1,
6945 true,
6946 1,
6947 2,
6948 1,
6949 IncompleteReviewPolicy::Block
6950 ));
6951 assert!(!round_is_clean(
6952 0,
6953 false,
6954 1,
6955 2,
6956 1,
6957 IncompleteReviewPolicy::Block
6958 ));
6959 }
6960
6961 #[test]
6962 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
6963 assert!(!round_is_clean(
6967 0,
6968 true,
6969 0,
6970 2,
6971 2,
6972 IncompleteReviewPolicy::Block
6973 ));
6974 }
6975
6976 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
6977 CommandOutcome {
6978 command: "test".to_owned(),
6979 code,
6980 output_tail: String::new(),
6981 duration_ms: 0,
6982 resource_blocked,
6983 }
6984 }
6985
6986 #[test]
6987 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
6988 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
6989 assert!(
6990 !verify_inconclusive(&[outcome(Some(1), false)]),
6991 "an ordinary failure is still evidence about the patch"
6992 );
6993 assert!(verify_inconclusive(&[outcome(None, true)]));
6994 assert!(
6995 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
6996 "one inconclusive outcome taints the whole batch"
6997 );
6998 assert!(!verify_inconclusive(&[]));
6999 }
7000
7001 #[tokio::test]
7002 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7003 let calls = std::sync::atomic::AtomicUsize::new(0);
7007 let started = Instant::now();
7008 wait_for_pids_with(
7009 &[123],
7010 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7011 Duration::from_millis(5),
7012 Duration::from_secs(5),
7013 )
7014 .await;
7015 assert!(
7016 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7017 "must keep checking rather than deciding on the first answer"
7018 );
7019 assert!(
7020 started.elapsed() < Duration::from_secs(1),
7021 "must return the moment it is confirmed dead, not wait out the ceiling"
7022 );
7023 }
7024
7025 #[tokio::test]
7026 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7027 let started = Instant::now();
7028 wait_for_pids_with(
7029 &[123],
7030 |_| true, Duration::from_millis(5),
7032 Duration::from_millis(30),
7033 )
7034 .await;
7035 let elapsed = started.elapsed();
7036 assert!(
7037 elapsed >= Duration::from_millis(30),
7038 "must not give up before its own ceiling: {elapsed:?}"
7039 );
7040 assert!(
7041 elapsed < Duration::from_secs(1),
7042 "must not wait past its own ceiling either: {elapsed:?}"
7043 );
7044 }
7045
7046 #[tokio::test]
7047 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7048 let started = Instant::now();
7049 wait_for_pids_with(
7050 &[],
7051 |_| true,
7052 Duration::from_secs(5),
7053 Duration::from_secs(5),
7054 )
7055 .await;
7056 assert!(
7057 started.elapsed() < Duration::from_millis(200),
7058 "an empty pid list has nothing to confirm"
7059 );
7060 }
7061
7062 fn review_round(
7068 clean: bool,
7069 blocking: usize,
7070 answered: usize,
7071 expected: usize,
7072 progressed: bool,
7073 e2e_ok: bool,
7074 ) -> ReviewRound {
7075 ReviewRound {
7076 round: 1,
7077 head: "h".to_owned(),
7078 verified_head: None,
7079 verified_at: None,
7080 reviews: Vec::new(),
7081 e2e: vec![CommandOutcome {
7082 command: "test".to_owned(),
7083 code: Some(if e2e_ok { 0 } else { 1 }),
7084 output_tail: String::new(),
7085 duration_ms: 0,
7086 resource_blocked: false,
7087 }],
7088 verify_retried: false,
7089 e2e_deferred: false,
7090 e2e_defer_reason: None,
7091 fix: None,
7092 blocking,
7093 answered,
7094 expected,
7095 clean,
7096 progressed,
7097 vote_split: false,
7098 reconsideration: Vec::new(),
7099 verdict: None,
7100 }
7101 }
7102
7103 #[test]
7104 fn review_conclusion_is_none_when_nothing_has_run() {
7105 assert_eq!(review_conclusion(&[], 3), None);
7106 }
7107
7108 #[test]
7109 fn review_conclusion_is_none_while_rounds_remain() {
7110 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7111 assert_eq!(review_conclusion(&rounds, 3), None);
7112 }
7113
7114 #[test]
7115 fn review_conclusion_is_gating_once_a_round_is_clean() {
7116 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7117 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7118 }
7119
7120 #[test]
7121 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7122 let rounds = vec![
7123 review_round(false, 1, 2, 2, true, true),
7124 review_round(false, 1, 2, 2, true, true),
7125 ];
7126 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7127 }
7128
7129 #[test]
7130 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7131 let rounds = vec![
7132 review_round(false, 1, 2, 2, true, true),
7133 review_round(false, 1, 2, 2, true, false),
7134 ];
7135 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7136 }
7137
7138 #[test]
7139 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7140 let mut blocked = review_round(false, 1, 2, 2, true, false);
7147 blocked.e2e[0].resource_blocked = true;
7148 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7149 assert_eq!(review_conclusion(&rounds, 2), None);
7150 }
7151
7152 #[test]
7153 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7154 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7156 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7157 }
7158
7159 #[test]
7160 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7161 let rounds = vec![
7162 review_round(false, 1, 2, 2, false, true),
7163 review_round(false, 1, 2, 2, false, true),
7164 ];
7165 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7166 }
7167
7168 fn secs(n: u64) -> Duration {
7169 Duration::from_secs(n)
7170 }
7171
7172 fn init_repo(dir: &Path) {
7175 let run = |args: &[&str]| {
7176 let out = std::process::Command::new("git")
7177 .args(args)
7178 .current_dir(dir)
7179 .quiet()
7180 .output()
7181 .expect("spawn git");
7182 assert!(
7183 out.status.success(),
7184 "git {args:?} failed: {}",
7185 String::from_utf8_lossy(&out.stderr)
7186 );
7187 };
7188 run(&["init", "-b", "main"]);
7189 run(&["config", "user.name", "magi test"]);
7190 run(&["config", "user.email", "magi@example.com"]);
7191 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7192 run(&["add", "-A"]);
7193 run(&["commit", "-m", "init"]);
7194 }
7195
7196 fn ask_test_home() {
7204 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7205 }
7206
7207 fn runner_at(status: RunStatus) -> Runner {
7210 let mut state = RunState::new(
7211 PathBuf::from("/nonexistent/repo"),
7212 "main".to_owned(),
7213 "deadbeef".to_owned(),
7214 "task".to_owned(),
7215 Config::default(),
7216 );
7217 state.status = status;
7218 Runner {
7219 state,
7220 roles: ResolvedRoles {
7221 implementers: Vec::new(),
7222 judges: Vec::new(),
7223 reviewers: Vec::new(),
7224 fixer: None,
7225 conductor: conductor(),
7226 implementer_roster: Vec::new(),
7227 },
7228 sem: Arc::new(Semaphore::new(1)),
7229 pause: Pause::new(),
7230 interrupt: Pause::new(),
7231 }
7232 }
7233
7234 #[test]
7238 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7239 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7240 let mut runner = runner_at(RunStatus::Implementing);
7241 let interrupt = Pause::new();
7242 runner.watch_interrupt(interrupt.clone());
7243
7244 interrupt.park_because("task a1b2 asked to run first");
7245
7246 assert!(runner.park_here().expect("park_here"));
7247 assert!(runner.state.parked);
7248 let last = runner.state.events.last().expect("a park event");
7249 assert_eq!(last.node, "park");
7250 assert!(
7251 last.message.contains("task a1b2 asked to run first"),
7252 "expected the interrupt reason in {:?}",
7253 last.message
7254 );
7255 }
7256
7257 #[test]
7265 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7266 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7267 let mut runner = runner_at(RunStatus::Implementing);
7268 let shutdown = Pause::new();
7269 runner.on_pause(shutdown.clone());
7270 let interrupt = Pause::new();
7271 runner.watch_interrupt(interrupt.clone());
7272
7273 assert!(!runner.park_here().expect("park_here"));
7275 assert!(!runner.state.parked);
7276
7277 interrupt.park_because("test");
7279 assert!(!shutdown.parked());
7280 assert!(runner.park_here().expect("park_here"));
7281 }
7282
7283 #[tokio::test]
7297 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7298 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7299 let mut runner = runner_at(RunStatus::Implementing);
7300 let interrupt = Pause::new();
7301 runner.watch_interrupt(interrupt.clone());
7302
7303 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7304 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7305
7306 let node = async move {
7310 started_tx.send(()).expect("send started");
7311 finish_rx.await.expect("recv finish");
7312 "node finished"
7313 };
7314
7315 let interrupter = async move {
7316 started_rx.await.expect("recv started");
7317 interrupt.park_because("higher-priority task waiting");
7319 tokio::task::yield_now().await;
7323 finish_tx.send(()).expect("send finish");
7324 };
7325
7326 let (node_result, ()) = tokio::join!(node, interrupter);
7327 assert_eq!(
7328 node_result, "node finished",
7329 "the in-flight call ran to completion"
7330 );
7331
7332 assert!(runner.park_here().expect("park_here"));
7335 assert!(runner.state.parked);
7336 }
7337
7338 #[test]
7344 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7345 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7346 let mut runner = runner_at(RunStatus::Judging);
7347 runner.state.config.agents = vec![conductor()];
7351 runner.state.candidates = vec![Candidate {
7352 index: 0,
7353 label: 'A',
7354 agent: "alpha".to_owned(),
7355 branch: "magi/x/A".to_owned(),
7356 worktree: PathBuf::from("/nonexistent/worktree"),
7357 summary: "did the thing".to_owned(),
7358 stat: "1 file changed".to_owned(),
7359 files: 1,
7360 commits: 1,
7361 empty: false,
7362 failed: None,
7363 verified_noop: None,
7364 duration_ms: 1234,
7365 folded: false,
7366 }];
7367 let run_id = runner.state.id.clone();
7368
7369 let interrupt = Pause::new();
7370 runner.watch_interrupt(interrupt.clone());
7371 interrupt.park_because("task c3d4 asked to run first");
7372 assert!(runner.park_here().expect("park_here"));
7373
7374 let resumed = Runner::resume(&run_id).expect("resume");
7375 assert_eq!(resumed.state.candidates.len(), 1);
7376 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7377 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7378 assert_eq!(resumed.state.status, runner.state.status);
7379 assert!(
7380 resumed.state.parked,
7381 "still parked until `execute` actually walks the graph again"
7382 );
7383 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7384 }
7385
7386 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7388 let mut q = ask::Question::new(
7389 run.to_owned(),
7390 "implement".to_owned(),
7391 "impl-A".to_owned(),
7392 "Which storage backend should the cache use?".to_owned(),
7393 String::new(),
7394 vec!["SQLite".to_owned(), "Redis".to_owned()],
7395 );
7396 store.put(&mut q).unwrap();
7397 q
7398 }
7399
7400 #[test]
7401 fn a_failed_runs_open_question_is_abandoned() {
7402 ask_test_home();
7403 let store = ask::Questions::open();
7404 let mut runner = runner_at(RunStatus::Failed);
7405 let run = runner.state.id.clone();
7406 let q = ask_open_question(&store, &run);
7407
7408 runner.settle_questions();
7409
7410 let back = store.get(&q.id).unwrap();
7411 assert!(
7412 !back.status.open(),
7413 "the seat that asked died with the run; nobody is left to read an answer"
7414 );
7415 assert!(
7416 back.detail.contains(&run) && back.detail.contains("failed"),
7417 "the reason names what the run became, not just that it is gone: {}",
7418 back.detail
7419 );
7420 }
7421
7422 #[test]
7423 fn a_merged_runs_open_question_is_abandoned_too() {
7424 ask_test_home();
7425 let store = ask::Questions::open();
7426 for status in [RunStatus::Merged, RunStatus::Ready] {
7429 let mut runner = runner_at(status);
7430 let run = runner.state.id.clone();
7431 let q = ask_open_question(&store, &run);
7432
7433 runner.settle_questions();
7434
7435 let back = store.get(&q.id).unwrap();
7436 assert!(
7437 !back.status.open(),
7438 "{status:?} run's question must not outlive the run"
7439 );
7440 }
7441 }
7442
7443 #[test]
7444 fn a_still_resumable_runs_open_question_is_left_alone() {
7445 ask_test_home();
7446 let store = ask::Questions::open();
7447 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7453 let mut runner = runner_at(status);
7454 let run = runner.state.id.clone();
7455 let q = ask_open_question(&store, &run);
7456
7457 runner.settle_questions();
7458
7459 let back = store.get(&q.id).unwrap();
7460 assert!(
7461 back.status.open(),
7462 "{status:?} is still alive; the question must still be waiting"
7463 );
7464 }
7465 }
7466
7467 #[test]
7468 fn settle_questions_never_touches_an_already_answered_question() {
7469 ask_test_home();
7470 let store = ask::Questions::open();
7471 let mut runner = runner_at(RunStatus::Failed);
7472 let run = runner.state.id.clone();
7473 let mut q = ask_open_question(&store, &run);
7474 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
7475 .unwrap();
7476 store.put(&mut q).unwrap();
7477
7478 runner.settle_questions();
7483 runner.settle_questions();
7484
7485 let back = store.get(&q.id).unwrap();
7486 assert_eq!(
7487 back.status,
7488 ask::QuestionStatus::Answered,
7489 "a real answer is a decision on record, never overwritten by a sweep"
7490 );
7491 }
7492
7493 #[tokio::test]
7504 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
7505 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
7506 let tmp = tempfile::tempdir().expect("tempdir");
7507 let repo = tmp.path().join("repo");
7508 std::fs::create_dir_all(&repo).unwrap();
7509 init_repo(&repo);
7510
7511 let mut config = Config::default();
7512 config.graph.worktree_root = Some(tmp.path().join("wt"));
7513
7514 let mut state = RunState::new(
7515 repo.clone(),
7516 "main".to_owned(),
7517 "deadbeef".to_owned(),
7518 "task".to_owned(),
7519 config,
7520 );
7521 let root = state.worktree_root();
7522 let wt_a = root.join("cand-A");
7523 let wt_b = root.join("cand-B");
7524 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
7525 .await
7526 .expect("worktree A");
7527 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
7528 .await
7529 .expect("worktree B");
7530
7531 state.candidates = vec![
7532 Candidate {
7533 index: 0,
7534 label: 'A',
7535 agent: "alpha".to_owned(),
7536 branch: "magi/x/A".to_owned(),
7537 worktree: wt_a.clone(),
7538 summary: String::new(),
7539 stat: String::new(),
7540 files: 0,
7541 commits: 0,
7542 empty: false,
7543 failed: None,
7544 verified_noop: None,
7545 duration_ms: 0,
7546 folded: false,
7547 },
7548 Candidate {
7549 index: 1,
7550 label: 'B',
7551 agent: "beta".to_owned(),
7552 branch: "magi/x/B".to_owned(),
7553 worktree: wt_b.clone(),
7554 summary: String::new(),
7555 stat: String::new(),
7556 files: 0,
7557 commits: 0,
7558 empty: false,
7559 failed: None,
7560 verified_noop: None,
7561 duration_ms: 0,
7562 folded: false,
7563 },
7564 ];
7565 state.tally = Some(Tally {
7566 first_choice: BTreeMap::from([('A', 1)]),
7567 borda: BTreeMap::new(),
7568 winner: 'A',
7569 rankings: 1,
7570 unanimous_initial: true,
7571 deliberated: false,
7572 changed_votes: 0,
7573 unanimous_final: true,
7574 tie_break: None,
7575 judges: 1,
7576 present: 1,
7577 quorum: 1,
7578 met_quorum: true,
7579 uncontested: None,
7580 });
7581 state.status = RunStatus::Ready;
7582
7583 fold_run(&mut state, false, &crate::run::home())
7584 .await
7585 .expect("fold_run");
7586
7587 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
7588 assert!(
7589 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7590 "the unmerged winner's branch survives"
7591 );
7592 assert!(
7593 !state.candidates[0].folded,
7594 "the winner is not marked folded"
7595 );
7596
7597 assert!(!wt_b.exists(), "the loser's worktree is removed");
7598 assert!(
7599 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
7600 "the loser's branch is removed"
7601 );
7602 assert!(state.candidates[1].folded, "the loser is marked folded");
7603 }
7604
7605 #[tokio::test]
7614 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
7615 let tmp = tempfile::tempdir().expect("tempdir");
7616 let repo = tmp.path().join("repo");
7617 std::fs::create_dir_all(&repo).unwrap();
7618 init_repo(&repo);
7619
7620 let mut config = Config::default();
7621 config.merge.mode = MergeMode::Local;
7622
7623 let mut state = RunState::new(
7624 repo.clone(),
7625 "main".to_owned(),
7626 "deadbeef".to_owned(),
7627 "task".to_owned(),
7628 config,
7629 );
7630 state.candidates = vec![Candidate {
7631 index: 0,
7632 label: 'A',
7633 agent: "alpha".to_owned(),
7634 branch: "does-not-exist".to_owned(),
7635 worktree: repo.clone(),
7636 summary: String::new(),
7637 stat: String::new(),
7638 files: 0,
7639 commits: 0,
7640 empty: false,
7641 failed: None,
7642 verified_noop: None,
7643 duration_ms: 0,
7644 folded: false,
7645 }];
7646 state.tally = Some(Tally {
7647 first_choice: BTreeMap::from([('A', 1)]),
7648 borda: BTreeMap::new(),
7649 winner: 'A',
7650 rankings: 1,
7651 unanimous_initial: true,
7652 deliberated: false,
7653 changed_votes: 0,
7654 unanimous_final: true,
7655 tie_break: None,
7656 judges: 0,
7657 present: 0,
7658 quorum: 0,
7659 met_quorum: true,
7660 uncontested: Some("only candidate A produced a change".to_owned()),
7661 });
7662 state.reviews = vec![ReviewRound {
7663 round: 1,
7664 head: "deadbeef".to_owned(),
7665 verified_head: None,
7666 verified_at: None,
7667 reviews: Vec::new(),
7668 e2e: Vec::new(),
7669 fix: None,
7670 blocking: 0,
7671 answered: 0,
7672 expected: 0,
7673 clean: true,
7674 verify_retried: false,
7675 e2e_deferred: false,
7676 e2e_defer_reason: None,
7677 progressed: false,
7678 vote_split: false,
7679 reconsideration: Vec::new(),
7680 verdict: None,
7681 }];
7682 state.gate = vec![CommandOutcome {
7683 command: "test".to_owned(),
7684 code: Some(0),
7685 output_tail: String::new(),
7686 duration_ms: 0,
7687 resource_blocked: false,
7688 }];
7689 state.gate_ran = true;
7690 state.status = RunStatus::Ready;
7695 state.merge = Some(MergeOutcome {
7696 mode: MergeMode::Local,
7697 ok: false,
7698 detail: "already concluded".to_owned(),
7699 });
7700
7701 let mut runner = Runner {
7702 state,
7703 roles: ResolvedRoles {
7704 implementers: Vec::new(),
7705 judges: Vec::new(),
7706 reviewers: Vec::new(),
7707 fixer: None,
7708 conductor: conductor(),
7709 implementer_roster: Vec::new(),
7710 },
7711 sem: Arc::new(Semaphore::new(1)),
7712 pause: Pause::new(),
7713 interrupt: Pause::new(),
7714 };
7715
7716 runner.merge().await.expect("merge");
7717
7718 assert_eq!(
7719 runner.state.status,
7720 RunStatus::Ready,
7721 "a concluded run's status must not change on reentry"
7722 );
7723 assert_eq!(
7724 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
7725 Some("already concluded"),
7726 "merge must not run again once the node already recorded an outcome"
7727 );
7728 }
7729
7730 #[tokio::test]
7739 async fn merge_refuses_a_gate_that_has_not_actually_run() {
7740 let tmp = tempfile::tempdir().expect("tempdir");
7741 let repo = tmp.path().join("repo");
7742 std::fs::create_dir_all(&repo).unwrap();
7743 init_repo(&repo);
7744
7745 let mut config = Config::default();
7746 config.merge.mode = MergeMode::Local;
7747
7748 let mut state = RunState::new(
7749 repo.clone(),
7750 "main".to_owned(),
7751 "deadbeef".to_owned(),
7752 "task".to_owned(),
7753 config,
7754 );
7755 state.candidates = vec![Candidate {
7756 index: 0,
7757 label: 'A',
7758 agent: "alpha".to_owned(),
7759 branch: "does-not-exist".to_owned(),
7760 worktree: repo.clone(),
7761 summary: String::new(),
7762 stat: String::new(),
7763 files: 0,
7764 commits: 0,
7765 empty: false,
7766 failed: None,
7767 verified_noop: None,
7768 duration_ms: 0,
7769 folded: false,
7770 }];
7771 state.tally = Some(Tally {
7772 first_choice: BTreeMap::from([('A', 1)]),
7773 borda: BTreeMap::new(),
7774 winner: 'A',
7775 rankings: 1,
7776 unanimous_initial: true,
7777 deliberated: false,
7778 changed_votes: 0,
7779 unanimous_final: true,
7780 tie_break: None,
7781 judges: 0,
7782 present: 0,
7783 quorum: 0,
7784 met_quorum: true,
7785 uncontested: Some("only candidate A produced a change".to_owned()),
7786 });
7787 state.reviews = vec![ReviewRound {
7788 round: 1,
7789 head: "deadbeef".to_owned(),
7790 verified_head: None,
7791 verified_at: None,
7792 reviews: Vec::new(),
7793 e2e: Vec::new(),
7794 fix: None,
7795 blocking: 0,
7796 answered: 0,
7797 expected: 0,
7798 clean: true,
7799 verify_retried: false,
7800 e2e_deferred: false,
7801 e2e_defer_reason: None,
7802 progressed: false,
7803 vote_split: false,
7804 reconsideration: Vec::new(),
7805 verdict: None,
7806 }];
7807 state.gate = Vec::new();
7809 state.gate_ran = false;
7810 state.status = RunStatus::Gating;
7811
7812 let mut runner = Runner {
7813 state,
7814 roles: ResolvedRoles {
7815 implementers: Vec::new(),
7816 judges: Vec::new(),
7817 reviewers: Vec::new(),
7818 fixer: None,
7819 conductor: conductor(),
7820 implementer_roster: Vec::new(),
7821 },
7822 sem: Arc::new(Semaphore::new(1)),
7823 pause: Pause::new(),
7824 interrupt: Pause::new(),
7825 };
7826
7827 runner.merge().await.expect("merge");
7828
7829 assert!(
7830 runner.state.merge.is_none(),
7831 "an empty gate must never be read as a passing one: {:?}",
7832 runner.state.merge
7833 );
7834 }
7835
7836 #[tokio::test]
7843 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
7844 let tmp = tempfile::tempdir().expect("tempdir");
7845 let repo = tmp.path().join("repo");
7846 std::fs::create_dir_all(&repo).unwrap();
7847 init_repo(&repo);
7848
7849 let config = Config::default();
7851
7852 let mut state = RunState::new(
7853 repo.clone(),
7854 "main".to_owned(),
7855 "deadbeef".to_owned(),
7856 "task".to_owned(),
7857 config,
7858 );
7859 state.candidates = vec![Candidate {
7860 index: 0,
7861 label: 'A',
7862 agent: "alpha".to_owned(),
7863 branch: "does-not-exist".to_owned(),
7864 worktree: repo.clone(),
7865 summary: String::new(),
7866 stat: String::new(),
7867 files: 0,
7868 commits: 0,
7869 empty: false,
7870 failed: None,
7871 verified_noop: None,
7872 duration_ms: 0,
7873 folded: false,
7874 }];
7875 state.tally = Some(Tally {
7876 first_choice: BTreeMap::from([('A', 1)]),
7877 borda: BTreeMap::new(),
7878 winner: 'A',
7879 rankings: 1,
7880 unanimous_initial: true,
7881 deliberated: false,
7882 changed_votes: 0,
7883 unanimous_final: true,
7884 tie_break: None,
7885 judges: 0,
7886 present: 0,
7887 quorum: 0,
7888 met_quorum: true,
7889 uncontested: Some("only candidate A produced a change".to_owned()),
7890 });
7891 state.reviews = vec![ReviewRound {
7892 round: 1,
7893 head: "deadbeef".to_owned(),
7894 verified_head: None,
7895 verified_at: None,
7896 reviews: Vec::new(),
7897 e2e: Vec::new(),
7898 fix: None,
7899 blocking: 0,
7900 answered: 0,
7901 expected: 0,
7902 clean: true,
7903 verify_retried: false,
7904 e2e_deferred: false,
7905 e2e_defer_reason: None,
7906 progressed: false,
7907 vote_split: false,
7908 reconsideration: Vec::new(),
7909 verdict: None,
7910 }];
7911
7912 let mut runner = Runner {
7913 state,
7914 roles: ResolvedRoles {
7915 implementers: Vec::new(),
7916 judges: Vec::new(),
7917 reviewers: Vec::new(),
7918 fixer: None,
7919 conductor: conductor(),
7920 implementer_roster: Vec::new(),
7921 },
7922 sem: Arc::new(Semaphore::new(1)),
7923 pause: Pause::new(),
7924 interrupt: Pause::new(),
7925 };
7926
7927 runner.gate().await.expect("gate");
7928 assert!(
7929 runner.state.gate_ran,
7930 "zero configured commands is still a real attempt, not an unrun gate"
7931 );
7932 assert!(runner.state.gate.is_empty());
7933 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
7934 assert_ne!(
7935 runner.state.status,
7936 RunStatus::Blocked,
7937 "a gate with nothing to check must not read as failed"
7938 );
7939
7940 runner.merge().await.expect("merge");
7941 assert_eq!(
7942 runner.state.status,
7943 RunStatus::Ready,
7944 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
7945 );
7946 }
7947
7948 #[tokio::test]
7959 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
7960 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
7961 let home = crate::run::home();
7962
7963 let tmp = tempfile::tempdir().expect("tempdir");
7964 let repo = tmp.path().join("repo");
7965 std::fs::create_dir_all(&repo).unwrap();
7966 init_repo(&repo);
7967 let cache_dir = tmp.path().join("target");
7970
7971 let mut config = Config::default();
7972 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
7973 config.graph.timeout_verify = Some(2);
7976
7977 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
7978 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
7979 .expect("no io error acquiring directly")
7980 {
7981 crate::cache::AcquireOutcome::Acquired(g) => g,
7982 crate::cache::AcquireOutcome::Busy(b) => {
7983 panic!("expected the direct acquire to win the lease first: {b:?}")
7984 }
7985 };
7986
7987 let mut state = RunState::new(
7988 repo.clone(),
7989 "main".to_owned(),
7990 "deadbeef".to_owned(),
7991 "task".to_owned(),
7992 config,
7993 );
7994 state.candidates = vec![Candidate {
7995 index: 0,
7996 label: 'A',
7997 agent: "alpha".to_owned(),
7998 branch: "does-not-exist".to_owned(),
7999 worktree: repo.clone(),
8000 summary: String::new(),
8001 stat: String::new(),
8002 files: 0,
8003 commits: 0,
8004 empty: false,
8005 failed: None,
8006 verified_noop: None,
8007 duration_ms: 0,
8008 folded: false,
8009 }];
8010 state.tally = Some(Tally {
8011 first_choice: BTreeMap::from([('A', 1)]),
8012 borda: BTreeMap::new(),
8013 winner: 'A',
8014 rankings: 1,
8015 unanimous_initial: true,
8016 deliberated: false,
8017 changed_votes: 0,
8018 unanimous_final: true,
8019 tie_break: None,
8020 judges: 0,
8021 present: 0,
8022 quorum: 0,
8023 met_quorum: true,
8024 uncontested: Some("only candidate A produced a change".to_owned()),
8025 });
8026 state.reviews = vec![ReviewRound {
8027 round: 1,
8028 head: "deadbeef".to_owned(),
8029 verified_head: None,
8030 verified_at: None,
8031 reviews: Vec::new(),
8032 e2e: Vec::new(),
8033 fix: None,
8034 blocking: 0,
8035 answered: 0,
8036 expected: 0,
8037 clean: true,
8038 verify_retried: false,
8039 e2e_deferred: false,
8040 e2e_defer_reason: None,
8041 progressed: false,
8042 vote_split: false,
8043 reconsideration: Vec::new(),
8044 verdict: None,
8045 }];
8046
8047 let mut runner = Runner {
8048 state,
8049 roles: ResolvedRoles {
8050 implementers: Vec::new(),
8051 judges: Vec::new(),
8052 reviewers: Vec::new(),
8053 fixer: None,
8054 conductor: conductor(),
8055 implementer_roster: Vec::new(),
8056 },
8057 sem: Arc::new(Semaphore::new(1)),
8058 pause: Pause::new(),
8059 interrupt: Pause::new(),
8060 };
8061
8062 let started = std::time::Instant::now();
8063 runner.gate().await.expect("gate");
8064 assert!(
8065 started.elapsed() < Duration::from_secs(1),
8066 "a gate with nothing to run must never wait on a lease it never needed"
8067 );
8068 assert!(
8069 runner.state.gate_ran,
8070 "zero commands is still a real, immediate attempt"
8071 );
8072 assert!(runner.state.gate.is_empty());
8073 assert_ne!(
8074 runner.state.status,
8075 RunStatus::Blocked,
8076 "must not read as resource-blocked on a lease it never asked for"
8077 );
8078 }
8079
8080 #[tokio::test]
8090 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8091 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8092
8093 let tmp = tempfile::tempdir().expect("tempdir");
8094 let repo = tmp.path().join("repo");
8095 std::fs::create_dir_all(&repo).unwrap();
8096 init_repo(&repo);
8097
8098 let mut config = Config::default();
8099 config.verify.gate = vec![
8100 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8101 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8102 .to_owned(),
8103 ];
8104
8105 let mut state = RunState::new(
8106 repo.clone(),
8107 "main".to_owned(),
8108 "deadbeef".to_owned(),
8109 "task".to_owned(),
8110 config,
8111 );
8112 let run_id = state.id.clone();
8113 state.candidates = vec![Candidate {
8114 index: 0,
8115 label: 'A',
8116 agent: "alpha".to_owned(),
8117 branch: "does-not-exist".to_owned(),
8118 worktree: repo.clone(),
8119 summary: String::new(),
8120 stat: String::new(),
8121 files: 0,
8122 commits: 0,
8123 empty: false,
8124 failed: None,
8125 verified_noop: None,
8126 duration_ms: 0,
8127 folded: false,
8128 }];
8129 state.tally = Some(Tally {
8130 first_choice: BTreeMap::from([('A', 1)]),
8131 borda: BTreeMap::new(),
8132 winner: 'A',
8133 rankings: 1,
8134 unanimous_initial: true,
8135 deliberated: false,
8136 changed_votes: 0,
8137 unanimous_final: true,
8138 tie_break: None,
8139 judges: 0,
8140 present: 0,
8141 quorum: 0,
8142 met_quorum: true,
8143 uncontested: Some("only candidate A produced a change".to_owned()),
8144 });
8145 state.reviews = vec![ReviewRound {
8146 round: 1,
8147 head: "deadbeef".to_owned(),
8148 verified_head: None,
8149 verified_at: None,
8150 reviews: Vec::new(),
8151 e2e: Vec::new(),
8152 fix: None,
8153 blocking: 0,
8154 answered: 0,
8155 expected: 0,
8156 clean: true,
8157 verify_retried: false,
8158 e2e_deferred: false,
8159 e2e_defer_reason: None,
8160 progressed: false,
8161 vote_split: false,
8162 reconsideration: Vec::new(),
8163 verdict: None,
8164 }];
8165
8166 let mut runner = Runner {
8167 state,
8168 roles: ResolvedRoles {
8169 implementers: Vec::new(),
8170 judges: Vec::new(),
8171 reviewers: Vec::new(),
8172 fixer: None,
8173 conductor: conductor(),
8174 implementer_roster: Vec::new(),
8175 },
8176 sem: Arc::new(Semaphore::new(1)),
8177 pause: Pause::new(),
8178 interrupt: Pause::new(),
8179 };
8180
8181 let started_marker = repo.join("started.marker");
8182 let release_marker = repo.join("release.marker");
8183 let poller = tokio::spawn(async move {
8184 for _ in 0..100 {
8189 if started_marker.exists()
8190 && let Ok(s) = crate::run::RunState::load(&run_id)
8191 && let Some(a) = s.active.get("gate")
8192 {
8193 std::fs::write(&release_marker, b"go").expect("release marker");
8194 return Some(a.clone());
8195 }
8196 tokio::time::sleep(Duration::from_millis(50)).await;
8197 }
8198 None
8199 });
8200
8201 runner.gate().await.expect("gate");
8202 let captured = poller.await.expect("poller task");
8203 let captured = captured.expect(
8204 "the poller never saw a `gate` task entry in run.json while the command was \
8205 still blocked on its own release marker",
8206 );
8207
8208 assert_eq!(captured.task.as_deref(), Some("gate"));
8209 assert_eq!(captured.node, "gate");
8210 assert_eq!(captured.index, Some(1));
8211 assert_eq!(captured.total, Some(1));
8212 assert!(
8213 captured
8214 .command
8215 .as_deref()
8216 .is_some_and(|c| c.contains("started.marker")),
8217 "{captured:?}"
8218 );
8219
8220 assert!(
8221 runner.state.active.is_empty(),
8222 "the entry must be cleared once the command actually finished: {:?}",
8223 runner.state.active
8224 );
8225 assert!(runner.state.gate_ran);
8226 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8227 }
8228
8229 #[tokio::test]
8242 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8243 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8244 let home = crate::run::home();
8245
8246 let tmp = tempfile::tempdir().expect("tempdir");
8247 let repo = tmp.path().join("repo");
8248 std::fs::create_dir_all(&repo).unwrap();
8249 init_repo(&repo);
8250 let head = crate::git::rev_parse(&repo, "HEAD")
8251 .await
8252 .expect("rev-parse");
8253 let cache_dir = tmp.path().join("target");
8256
8257 let mut config = Config::default();
8258 config.verify.e2e = vec![format!(
8259 "CARGO_TARGET_DIR='{}' test -f README.md",
8260 cache_dir.display()
8261 )];
8262 config.graph.review_rounds = 1;
8263 config.graph.timeout_verify = Some(2);
8266
8267 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8268 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8269 .expect("no io error acquiring directly")
8270 {
8271 crate::cache::AcquireOutcome::Acquired(g) => g,
8272 crate::cache::AcquireOutcome::Busy(b) => {
8273 panic!("expected the direct acquire to win the lease first: {b:?}")
8274 }
8275 };
8276
8277 let mut state = RunState::new(
8278 repo.clone(),
8279 "main".to_owned(),
8280 head.clone(),
8281 "task".to_owned(),
8282 config,
8283 );
8284 state.candidates = vec![Candidate {
8285 index: 0,
8286 label: 'A',
8287 agent: "alpha".to_owned(),
8288 branch: "does-not-exist".to_owned(),
8289 worktree: repo.clone(),
8290 summary: String::new(),
8291 stat: String::new(),
8292 files: 0,
8293 commits: 0,
8294 empty: false,
8295 failed: None,
8296 verified_noop: None,
8297 duration_ms: 0,
8298 folded: false,
8299 }];
8300 state.tally = Some(Tally {
8301 first_choice: BTreeMap::from([('A', 1)]),
8302 borda: BTreeMap::new(),
8303 winner: 'A',
8304 rankings: 1,
8305 unanimous_initial: true,
8306 deliberated: false,
8307 changed_votes: 0,
8308 unanimous_final: true,
8309 tie_break: None,
8310 judges: 0,
8311 present: 0,
8312 quorum: 0,
8313 met_quorum: true,
8314 uncontested: Some("only candidate A produced a change".to_owned()),
8315 });
8316 state.reviews = vec![ReviewRound {
8320 round: 1,
8321 head: head.clone(),
8322 verified_head: None,
8323 verified_at: None,
8324 reviews: Vec::new(),
8325 e2e: Vec::new(),
8326 fix: None,
8327 blocking: 1,
8328 answered: 1,
8329 expected: 1,
8330 clean: false,
8331 verify_retried: false,
8332 e2e_deferred: true,
8333 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8334 progressed: false,
8335 vote_split: false,
8336 reconsideration: Vec::new(),
8337 verdict: None,
8338 }];
8339
8340 let mut runner = Runner {
8341 state,
8342 roles: ResolvedRoles {
8343 implementers: Vec::new(),
8344 judges: Vec::new(),
8345 reviewers: Vec::new(),
8346 fixer: None,
8347 conductor: conductor(),
8348 implementer_roster: Vec::new(),
8349 },
8350 sem: Arc::new(Semaphore::new(1)),
8351 pause: Pause::new(),
8352 interrupt: Pause::new(),
8353 };
8354
8355 let shell = runner.state.config.shell();
8356 runner
8357 .stop_reviewing("round budget spent", &shell, &repo)
8358 .await
8359 .expect("stop_reviewing");
8360
8361 let last = runner.state.reviews.last().expect("round record");
8362 assert_eq!(
8363 last.e2e_status(),
8364 E2eStatus::ResourceBlocked,
8365 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8366 failed: {last:?}"
8367 );
8368 assert_eq!(
8369 last.verified_head.as_deref(),
8370 Some(head.as_str()),
8371 "which commit this attempt targeted is known even though nothing finished checking \
8372 it"
8373 );
8374 let first_attempt_at = last
8375 .verified_at
8376 .expect("when this attempt ran is known too");
8377 assert_ne!(
8378 runner.state.status,
8379 RunStatus::Blocked,
8380 "contention is evidence about the machine, not the patch — it must not settle the \
8381 run as blocked: {:?}",
8382 runner.state.status
8383 );
8384 assert!(
8385 !runner
8386 .state
8387 .events
8388 .iter()
8389 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
8390 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
8391 runner.state.events
8392 );
8393
8394 runner
8399 .stop_reviewing("round budget spent", &shell, &repo)
8400 .await
8401 .expect("stop_reviewing retry");
8402 assert_eq!(
8403 runner.state.reviews.len(),
8404 1,
8405 "no new round was started: {:?}",
8406 runner.state.reviews
8407 );
8408 let last = runner.state.reviews.last().expect("round record");
8409 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
8410 assert!(
8411 last.verified_at.expect("still known") > first_attempt_at,
8412 "a second reentry must be a fresh attempt, not a stale copy of the first"
8413 );
8414 assert_ne!(runner.state.status, RunStatus::Blocked);
8415
8416 held.release();
8417 }
8418
8419 #[tokio::test]
8431 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
8432 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8433 let home = crate::run::home();
8434
8435 let tmp = tempfile::tempdir().expect("tempdir");
8436 let repo = tmp.path().join("repo");
8437 std::fs::create_dir_all(&repo).unwrap();
8438 init_repo(&repo);
8439 let head = crate::git::rev_parse(&repo, "HEAD")
8440 .await
8441 .expect("rev-parse");
8442 let cache_dir = tmp.path().join("target");
8443
8444 let mut config = Config::default();
8445 config.verify.e2e = vec![format!(
8446 "CARGO_TARGET_DIR='{}' test -f README.md",
8447 cache_dir.display()
8448 )];
8449 config.graph.review_rounds = 1;
8450 config.graph.timeout_verify = Some(2);
8451
8452 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8453 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8454 .expect("no io error acquiring directly")
8455 {
8456 crate::cache::AcquireOutcome::Acquired(g) => g,
8457 crate::cache::AcquireOutcome::Busy(b) => {
8458 panic!("expected the direct acquire to win the lease first: {b:?}")
8459 }
8460 };
8461
8462 let mut state = RunState::new(
8463 repo.clone(),
8464 "main".to_owned(),
8465 head.clone(),
8466 "task".to_owned(),
8467 config,
8468 );
8469 state.candidates = vec![Candidate {
8470 index: 0,
8471 label: 'A',
8472 agent: "alpha".to_owned(),
8473 branch: "does-not-exist".to_owned(),
8474 worktree: repo.clone(),
8475 summary: String::new(),
8476 stat: String::new(),
8477 files: 0,
8478 commits: 0,
8479 empty: false,
8480 failed: None,
8481 verified_noop: None,
8482 duration_ms: 0,
8483 folded: false,
8484 }];
8485 state.tally = Some(Tally {
8486 first_choice: BTreeMap::from([('A', 1)]),
8487 borda: BTreeMap::new(),
8488 winner: 'A',
8489 rankings: 1,
8490 unanimous_initial: true,
8491 deliberated: false,
8492 changed_votes: 0,
8493 unanimous_final: true,
8494 tie_break: None,
8495 judges: 0,
8496 present: 0,
8497 quorum: 0,
8498 met_quorum: true,
8499 uncontested: Some("only candidate A produced a change".to_owned()),
8500 });
8501 state.reviews = vec![ReviewRound {
8505 round: 1,
8506 head: head.clone(),
8507 verified_head: Some(head.clone()),
8508 verified_at: Some(jiff::Timestamp::now()),
8509 reviews: Vec::new(),
8510 e2e: vec![CommandOutcome {
8511 command: format!(
8512 "CARGO_TARGET_DIR='{}' test -f README.md",
8513 cache_dir.display()
8514 ),
8515 code: None,
8516 output_tail: "waiting for the shared build cache".to_owned(),
8517 duration_ms: 0,
8518 resource_blocked: true,
8519 }],
8520 fix: None,
8521 blocking: 1,
8522 answered: 1,
8523 expected: 1,
8524 clean: false,
8525 verify_retried: false,
8526 e2e_deferred: false,
8527 e2e_defer_reason: None,
8528 progressed: false,
8529 vote_split: false,
8530 reconsideration: Vec::new(),
8531 verdict: None,
8532 }];
8533
8534 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
8535 let mut runner = Runner {
8536 state,
8537 roles: ResolvedRoles {
8538 implementers: Vec::new(),
8539 judges: Vec::new(),
8540 reviewers: Vec::new(),
8541 fixer: None,
8542 conductor: conductor(),
8543 implementer_roster: Vec::new(),
8544 },
8545 sem: Arc::new(Semaphore::new(1)),
8546 pause: Pause::new(),
8547 interrupt: Pause::new(),
8548 };
8549
8550 runner.review_loop().await.expect("review_loop");
8555
8556 assert_eq!(
8557 runner.state.reviews.len(),
8558 1,
8559 "no new round was started on top of the unresolved one: {:?}",
8560 runner.state.reviews
8561 );
8562 let last = &runner.state.reviews[0];
8563 assert_eq!(
8564 last.e2e_status(),
8565 E2eStatus::ResourceBlocked,
8566 "still contended: {last:?}"
8567 );
8568 assert!(
8569 last.verified_at.expect("still known") > first_attempt_at,
8570 "review_loop must have actually retried the check, not left it exactly as found"
8571 );
8572 assert_ne!(
8573 runner.state.status,
8574 RunStatus::Blocked,
8575 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
8576 runner.state.status
8577 );
8578
8579 held.release();
8580 }
8581
8582 #[tokio::test]
8583 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
8584 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8585 let tmp = tempfile::tempdir().expect("tempdir");
8586 let repo = tmp.path().join("repo");
8587 std::fs::create_dir_all(&repo).unwrap();
8588 init_repo(&repo);
8589
8590 let mut config = Config::default();
8591 config.merge.mode = MergeMode::Pr;
8592 config.graph.land = true;
8593 config.graph.land_approval = false;
8594
8595 let mut state = RunState::new(
8596 repo.clone(),
8597 "main".to_owned(),
8598 "deadbeef".to_owned(),
8599 "task".to_owned(),
8600 config,
8601 );
8602 state.candidates = vec![Candidate {
8603 index: 0,
8604 label: 'A',
8605 agent: "alpha".to_owned(),
8606 branch: "does-not-exist".to_owned(),
8607 worktree: repo.clone(),
8608 summary: String::new(),
8609 stat: String::new(),
8610 files: 0,
8611 commits: 0,
8612 empty: false,
8613 failed: None,
8614 verified_noop: None,
8615 duration_ms: 0,
8616 folded: false,
8617 }];
8618 state.tally = Some(Tally {
8619 first_choice: BTreeMap::from([('A', 1)]),
8620 borda: BTreeMap::new(),
8621 winner: 'A',
8622 rankings: 1,
8623 unanimous_initial: true,
8624 deliberated: false,
8625 changed_votes: 0,
8626 unanimous_final: true,
8627 tie_break: None,
8628 judges: 0,
8629 present: 0,
8630 quorum: 0,
8631 met_quorum: true,
8632 uncontested: Some("only candidate A produced a change".to_owned()),
8633 });
8634 state.reviews = vec![ReviewRound {
8635 round: 1,
8636 head: "deadbeef".to_owned(),
8637 verified_head: None,
8638 verified_at: None,
8639 reviews: Vec::new(),
8640 e2e: Vec::new(),
8641 fix: None,
8642 blocking: 0,
8643 answered: 0,
8644 expected: 0,
8645 clean: true,
8646 verify_retried: false,
8647 e2e_deferred: false,
8648 e2e_defer_reason: None,
8649 progressed: false,
8650 vote_split: false,
8651 reconsideration: Vec::new(),
8652 verdict: None,
8653 }];
8654 state.gate = vec![CommandOutcome {
8655 command: "test".to_owned(),
8656 code: Some(0),
8657 output_tail: String::new(),
8658 duration_ms: 0,
8659 resource_blocked: false,
8660 }];
8661 state.gate_ran = true;
8662 state.status = RunStatus::Landing;
8666 state.merge = Some(MergeOutcome {
8667 mode: MergeMode::Pr,
8668 ok: true,
8669 detail: "https://example.invalid/x/y/pull/1".to_owned(),
8670 });
8671
8672 ask_test_home();
8676 let store = ask::Questions::open();
8677 let q = ask_open_question(&store, &state.id);
8678
8679 let mut runner = Runner {
8680 state,
8681 roles: ResolvedRoles {
8682 implementers: Vec::new(),
8683 judges: Vec::new(),
8684 reviewers: Vec::new(),
8685 fixer: None,
8686 conductor: conductor(),
8687 implementer_roster: Vec::new(),
8688 },
8689 sem: Arc::new(Semaphore::new(1)),
8690 pause: Pause::new(),
8691 interrupt: Pause::new(),
8692 };
8693
8694 runner.execute().await.expect("execute");
8699
8700 assert_eq!(
8701 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8702 Some("https://example.invalid/x/y/pull/1"),
8703 "reentry must not push again or open a second pull request over the \
8704 one `land` is already watching"
8705 );
8706 assert_ne!(
8707 runner.state.status,
8708 RunStatus::Landing,
8709 "land could not actually reach the fake pull request, so it must \
8710 have given up rather than left the run silently parked forever"
8711 );
8712 assert_eq!(runner.state.status, RunStatus::Blocked);
8716 assert!(
8717 store.get(&q.id).unwrap().status.open(),
8718 "Blocked is still alive; settle_questions must have been a no-op here"
8719 );
8720 }
8721
8722 fn state_with_round(round: ReviewRound) -> RunState {
8723 let mut s = RunState::new(
8724 PathBuf::from("/repo"),
8725 "main".to_owned(),
8726 "abc1234".to_owned(),
8727 "add retries".to_owned(),
8728 Config::default(),
8729 );
8730 s.reviews = vec![round];
8731 s
8732 }
8733
8734 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
8735 crate::verdict::Finding {
8736 id: id.to_owned(),
8737 severity,
8738 file: None,
8739 line: None,
8740 title: title.to_owned(),
8741 detail: String::new(),
8742 }
8743 }
8744
8745 #[test]
8746 fn pr_body_names_open_findings_and_declined_ones() {
8747 let round = ReviewRound {
8748 round: 2,
8749 head: "deadbee".to_owned(),
8750 verified_head: None,
8751 verified_at: None,
8752 reviews: vec![ReviewRecord {
8753 attempts: 0,
8754 reviewer: 1,
8755 agent: "alpha".to_owned(),
8756 summary: String::new(),
8757 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
8758 vote: None,
8759 failed: None,
8760 duration_ms: 0,
8761 }],
8762 e2e: vec![CommandOutcome {
8763 command: "cargo test".to_owned(),
8764 code: Some(0),
8765 output_tail: String::new(),
8766 duration_ms: 0,
8767 resource_blocked: false,
8768 }],
8769 verify_retried: false,
8770 e2e_deferred: false,
8771 e2e_defer_reason: None,
8772 fix: Some(FixRecord {
8773 agent: "alpha".to_owned(),
8774 addressed: Vec::new(),
8775 rejected: vec![crate::verdict::Rejection {
8776 id: "R1-1-1".to_owned(),
8777 why: "not reachable from any caller".to_owned(),
8778 }],
8779 notes: String::new(),
8780 committed: true,
8781 failed: None,
8782 duration_ms: 0,
8783 continuation: None,
8784 }),
8785 blocking: 0,
8786 answered: 1,
8787 expected: 1,
8788 clean: false,
8789 progressed: true,
8790 vote_split: false,
8791 reconsideration: Vec::new(),
8792 verdict: None,
8793 };
8794 let state = state_with_round(round);
8795 let body = pr_message(&state, 'A').body;
8796
8797 assert!(body.contains("add retries"), "the task must still be there");
8798 assert!(body.contains("R2-1-1"), "{body}");
8799 assert!(body.contains("unused import"), "{body}");
8800 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
8801 assert!(
8802 body.contains("not reachable from any caller"),
8803 "the reason it was declined: {body}"
8804 );
8805 }
8806
8807 #[test]
8808 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
8809 let round = ReviewRound {
8810 round: 1,
8811 head: "deadbee".to_owned(),
8812 verified_head: None,
8813 verified_at: None,
8814 reviews: vec![ReviewRecord {
8815 attempts: 0,
8816 reviewer: 1,
8817 agent: "alpha".to_owned(),
8818 summary: String::new(),
8819 findings: Vec::new(),
8820 vote: None,
8821 failed: None,
8822 duration_ms: 0,
8823 }],
8824 e2e: Vec::new(),
8825 verify_retried: false,
8826 e2e_deferred: false,
8827 e2e_defer_reason: None,
8828 fix: None,
8829 blocking: 0,
8830 answered: 1,
8831 expected: 1,
8832 clean: true,
8833 progressed: false,
8834 vote_split: false,
8835 reconsideration: Vec::new(),
8836 verdict: None,
8837 };
8838 let state = state_with_round(round);
8839 let body = pr_message(&state, 'A').body;
8840 assert!(!body.contains("Open review findings"), "{body}");
8841 assert!(!body.contains("Declined"), "{body}");
8842 }
8843
8844 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
8845 let mut state = RunState::new(
8846 PathBuf::from("/repo"),
8847 "main".to_owned(),
8848 "abc1234".to_owned(),
8849 instruction.to_owned(),
8850 Config::default(),
8851 );
8852 state.candidates.push(Candidate {
8853 index: 0,
8854 label: 'A',
8855 agent: "alpha".to_owned(),
8856 branch: "magi/x/A".to_owned(),
8857 worktree: PathBuf::from("/wt"),
8858 summary: summary.to_owned(),
8859 stat: String::new(),
8860 files: 1,
8861 commits: 1,
8862 empty: false,
8863 failed: None,
8864 verified_noop: None,
8865 folded: false,
8866 duration_ms: 0,
8867 });
8868 state
8869 }
8870
8871 #[test]
8872 fn pr_message_describes_the_change_not_the_task() {
8873 let state = state_with_summary(
8874 "今回やってほしいこと: results projector を直す",
8875 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
8876 );
8877 let m = pr_message(&state, 'A');
8878 assert_eq!(m.title, "fix(web): batch the runs list reads");
8879 assert!(
8880 m.body.starts_with("## Summary\n\n- reads run.json once"),
8881 "{}",
8882 m.body
8883 );
8884 assert!(!m.body.contains("TITLE:"), "{}", m.body);
8885 let task_at = m.body.find("今回やってほしいこと").unwrap();
8886 let details_at = m.body.find("<details>").unwrap();
8887 assert!(
8888 details_at < task_at,
8889 "the task lives inside <details>: {}",
8890 m.body
8891 );
8892 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
8893 assert!(m.body.contains("magi:candidate-a"));
8894 }
8895
8896 #[test]
8897 fn pr_message_falls_back_to_the_task_without_a_title_line() {
8898 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
8899 let m = pr_message(&state, 'A');
8900 assert_eq!(m.title, "add retries");
8901 assert!(
8902 m.body.contains("## Summary\n\n- did some things"),
8903 "{}",
8904 m.body
8905 );
8906
8907 let none = RunState::new(
8908 PathBuf::from("/repo"),
8909 "main".to_owned(),
8910 "abc1234".to_owned(),
8911 "add retries".to_owned(),
8912 Config::default(),
8913 );
8914 let m = pr_message(&none, 'A');
8915 assert_eq!(m.title, "add retries");
8916 assert!(!m.body.contains("## Summary"), "{}", m.body);
8917 }
8918
8919 #[test]
8920 fn pr_message_refuses_the_candidate_commit_subject() {
8921 for bad in [
8922 "TITLE: magi: candidate A (uncommitted work)",
8923 "TITLE: chore: stuff (uncommitted work)",
8924 "TITLE: ",
8925 ] {
8926 let state = state_with_summary("add retries", bad);
8927 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
8928 }
8929 }
8930
8931 #[test]
8932 fn pr_message_bounds_a_very_long_task_and_title() {
8933 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
8934 let state = state_with_summary(&long, "- nothing");
8935 let m = pr_message(&state, 'A');
8936 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
8937 assert!(!m.title.contains('\n'));
8938
8939 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
8940 let m = pr_message(&state, 'A');
8941 assert!(m.title.starts_with("feat: "));
8942 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
8943 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
8944 }
8945
8946 #[test]
8947 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
8948 let mut state = state_with_summary(
8952 "add retries",
8953 "TITLE: fix(web): batch reads\n- reads run.json once",
8954 );
8955 state.config.graph.language = "ja".to_owned();
8956 let m = pr_message(&state, 'A');
8957 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
8958
8959 let task = "今回やってほしいこと: results projector を直す";
8962 let mut state = state_with_summary(task, "- no title line");
8963 state.config.graph.language = "ja".to_owned();
8964 let m = pr_message(&state, 'A');
8965 assert_eq!(
8966 m.title,
8967 format!("chore: land candidate A of run {}", state.id)
8968 );
8969 assert!(
8970 m.body.contains(&format!(
8971 "<summary>Original task</summary>\n\n{task}\n\n</details>"
8972 )),
8973 "{}",
8974 m.body
8975 );
8976 }
8977
8978 #[test]
8979 fn pr_message_scrubs_home_paths_and_addresses() {
8980 let state = state_with_summary(
8981 "fix it in /Users/someone/src/x",
8982 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
8983 );
8984 let m = pr_message(&state, 'A');
8985 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
8986 assert!(!m.body.contains(leak), "{}", m.body);
8987 }
8988 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
8989 }
8990
8991 #[test]
8992 fn pr_message_survives_a_task_that_closes_details() {
8993 let state = state_with_summary("a </details> b", "TITLE: fix: x");
8994 let m = pr_message(&state, 'A');
8995 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
8996 }
8997
8998 #[test]
8999 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9000 let cmd = manual_merge_command(
9001 MergeStyle::Squash,
9002 Path::new("/repo"),
9003 "b",
9004 "fix: \"quoted\" $(x) `y`\n\nbody",
9005 );
9006 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9007 }
9008
9009 #[test]
9010 fn manual_merge_command_matches_the_configured_style() {
9011 let repo = Path::new("/repo");
9012 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9013
9014 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9015 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9016
9017 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9018 assert_eq!(
9019 squash,
9020 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9021 \"Merge magi run 0832 (candidate A)\""
9022 );
9023
9024 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9025 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9026 }
9027
9028 #[test]
9029 fn a_nudge_gets_a_quarter_of_the_budget() {
9030 assert_eq!(retry_budget(secs(1200), true), secs(300));
9032 assert_eq!(retry_budget(secs(3600), true), secs(900));
9033 }
9034
9035 #[test]
9036 fn a_resent_prompt_keeps_the_whole_budget() {
9037 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9040 assert_eq!(retry_budget(secs(60), false), secs(60));
9041 }
9042
9043 #[test]
9044 fn the_floor_never_exceeds_the_original_budget() {
9045 assert_eq!(retry_budget(secs(60), true), secs(60));
9049 assert_eq!(retry_budget(secs(480), true), secs(120));
9050 assert_eq!(retry_budget(secs(0), true), secs(0));
9051 }
9052
9053 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9054 agent::CommandEvidence {
9055 id: "item1".to_owned(),
9056 description: "cargo test".to_owned(),
9057 exit_code,
9058 result_summary: String::new(),
9059 source: "codex".to_owned(),
9060 }
9061 }
9062
9063 #[test]
9064 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9065 assert!(!has_unconfirmed_command(&[]));
9069 }
9070
9071 #[test]
9072 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9073 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9077 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9078 assert!(!has_unconfirmed_command(&[
9079 evidence(Some(0)),
9080 evidence(Some(101))
9081 ]));
9082 }
9083
9084 #[test]
9085 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9086 assert!(has_unconfirmed_command(&[
9087 evidence(Some(0)),
9088 evidence(None)
9089 ]));
9090 }
9091
9092 #[test]
9093 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9094 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9095 assert_eq!(
9096 verified_noop_claim(true, &[], text).as_deref(),
9097 Some("already fixed by b32cfc4, on main.")
9098 );
9099 }
9100
9101 #[test]
9102 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9103 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9106 assert!(verified_noop_claim(false, &[], text).is_none());
9107 }
9108
9109 #[test]
9110 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9111 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9112 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9113 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9115 }
9116
9117 #[test]
9118 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9119 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9120 }
9121
9122 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9125 runner.state.candidates = shape
9126 .iter()
9127 .enumerate()
9128 .map(|(i, &(empty, verified))| Candidate {
9129 index: i,
9130 label: (b'A' + i as u8) as char,
9131 agent: "sonnet".to_owned(),
9132 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9133 worktree: PathBuf::from(format!("/wt/{i}")),
9134 summary: String::new(),
9135 stat: String::new(),
9136 files: 0,
9137 commits: 0,
9138 empty,
9139 failed: None,
9140 verified_noop: verified.map(str::to_owned),
9141 duration_ms: 0,
9142 folded: false,
9143 })
9144 .collect();
9145 }
9146
9147 #[test]
9148 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9149 ask_test_home();
9150 let mut runner = runner_at(RunStatus::Implementing);
9151 set_candidates(
9152 &mut runner,
9153 &[
9154 (true, Some("already on main at b32cfc4")),
9155 (true, Some("same fix, see the existing test")),
9156 ],
9157 );
9158
9159 runner
9160 .after_implement()
9161 .expect("a verified no-op is not an error");
9162
9163 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9164 }
9165
9166 #[test]
9167 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9168 ask_test_home();
9169 let mut runner = runner_at(RunStatus::Implementing);
9170 set_candidates(
9174 &mut runner,
9175 &[(true, Some("already on main at b32cfc4")), (true, None)],
9176 );
9177
9178 let err = runner
9179 .after_implement()
9180 .expect_err("an unverified empty candidate must still fail the run");
9181
9182 assert!(
9183 err.to_string().contains("no candidate produced a change"),
9184 "{err}"
9185 );
9186 assert_eq!(runner.state.status, RunStatus::Failed);
9187 }
9188
9189 #[test]
9190 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9191 ask_test_home();
9192 let mut runner = runner_at(RunStatus::Implementing);
9193 set_candidates(&mut runner, &[(true, None), (true, None)]);
9194
9195 let err = runner
9196 .after_implement()
9197 .expect_err("no candidate declared anything; this is an ordinary failure");
9198
9199 assert!(
9200 err.to_string().contains("no candidate produced a change"),
9201 "{err}"
9202 );
9203 assert_eq!(runner.state.status, RunStatus::Failed);
9204 }
9205}