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 sync_review_branch(repo: &Path, branch: &str, remote: &str) -> Result<()> {
294 let tracking = format!("{remote}/{branch}");
295 let fetched = git::fetch(repo, remote, branch).await;
296 let fresh = matches!(&fetched, Ok(o) if o.ok()) && git::rev_exists(repo, &tracking).await;
297 let local_exists = git::branch_exists(repo, branch).await?;
298 if !fresh {
299 if !local_exists {
300 bail!("no branch `{branch}` in {} or on {remote}", repo.display());
301 }
302 tracing::warn!(
303 "could not read {tracking}; reviewing the local `{branch}`, which may be stale"
304 );
305 return Ok(());
306 }
307 let remote_sha = git::rev_parse(repo, &tracking).await?;
308 if !local_exists {
309 git::git(repo, &["branch", branch, &tracking]).await?;
310 return Ok(());
311 }
312 let local_sha = git::rev_parse(repo, &format!("refs/heads/{branch}")).await?;
313 if local_sha == remote_sha || git::is_ancestor(repo, &remote_sha, &local_sha).await {
314 return Ok(());
315 }
316 if !git::is_ancestor(repo, &local_sha, &remote_sha).await {
317 let mb = git::git_raw(repo, &["merge-base", &local_sha, &remote_sha]).await?;
321 let placeholder = mb.ok()
322 && git::git_raw(repo, &["diff", "--quiet", &mb.stdout, &local_sha])
323 .await?
324 .ok();
325 if !placeholder {
326 bail!(
327 "local `{branch}` ({}) and {tracking} ({}) have diverged, so it is unclear \
328 which one to review; reconcile them, e.g. `git branch -f {branch} {tracking}` \
329 to review the pushed work, or push the local branch first",
330 short(&local_sha),
331 short(&remote_sha)
332 );
333 }
334 }
335 let out = git::git_raw(repo, &["branch", "-f", branch, &tracking]).await?;
336 if !out.ok() {
337 bail!(
338 "local `{branch}` ({}) is stale against {tracking} ({}) but git will not move it: {}",
339 short(&local_sha),
340 short(&remote_sha),
341 out.stderr
342 );
343 }
344 tracing::warn!(
345 "local `{branch}` was stale: fast-forwarded {} -> {}",
346 short(&local_sha),
347 short(&remote_sha)
348 );
349 Ok(())
350}
351
352async fn resolve_base(repo: &Path, base_branch: &str, remote: &str) -> Result<String> {
353 let tracking = format!("{remote}/{base_branch}");
354 let fetched = git::fetch(repo, remote, base_branch).await;
355 if let Ok(out) = &fetched
356 && out.ok()
357 && git::rev_exists(repo, &tracking).await
358 {
359 return git::rev_parse(repo, &tracking).await;
360 }
361 let why = match &fetched {
362 Ok(out) if !out.ok() => out.stderr.lines().next().unwrap_or("").to_owned(),
363 Ok(_) => format!("{remote} has no {base_branch}"),
364 Err(e) => e.to_string(),
365 };
366 tracing::warn!(
367 "could not read {tracking} ({why}); branching off the local \
368 {base_branch} instead, which may be behind"
369 );
370 git::rev_parse(repo, base_branch).await.with_context(|| {
371 format!(
372 "cannot resolve `{base_branch}`; set [merge] base in magi.toml to a \
373 branch that exists"
374 )
375 })
376}
377
378struct FixClaim {
395 path: PathBuf,
396}
397
398impl FixClaim {
399 fn acquire(dir: &Path) -> Result<Self> {
400 std::fs::create_dir_all(dir).with_context(|| format!("create {}", dir.display()))?;
401 let path = dir.join("fix.lock");
402 match Self::create(&path) {
403 Ok(claim) => Ok(claim),
404 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
405 if Self::reclaim_if_dead(&path) {
406 Self::create(&path).with_context(|| format!("lock {}", path.display()))
407 } else {
408 bail!(
409 "another `magi fix` is already running for this run ({} exists)",
410 path.display()
411 )
412 }
413 }
414 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
415 }
416 }
417
418 fn create(path: &Path) -> std::io::Result<Self> {
419 let mut f = std::fs::OpenOptions::new()
420 .write(true)
421 .create_new(true)
422 .open(path)?;
423 use std::io::Write as _;
424 writeln!(f, "{}", std::process::id())?;
426 Ok(Self {
427 path: path.to_owned(),
428 })
429 }
430
431 fn reclaim_if_dead(path: &Path) -> bool {
435 let dead = std::fs::read_to_string(path)
436 .ok()
437 .and_then(|body| body.trim().parse::<u32>().ok())
438 .is_some_and(|pid| !crate::proc::pid_alive(pid));
439 dead && std::fs::remove_file(path).is_ok()
440 }
441}
442
443impl Drop for FixClaim {
444 fn drop(&mut self) {
445 let _ = std::fs::remove_file(&self.path);
446 }
447}
448
449impl Runner {
450 pub async fn start(repo: &Path, instruction: String, config: Config) -> Result<Self> {
452 let repo = git::toplevel(repo).await?;
453 let missing = agent::missing_programs(&config.agents);
454 if !missing.is_empty() {
455 bail!(
456 "these agent programs are not on PATH: {}. Fix the roster in \
457 magi.toml or install them.",
458 missing.join(", ")
459 );
460 }
461 let base_branch = match config.merge.base.clone() {
462 Some(b) => b,
463 None => git::current_branch(&repo)
464 .await?
465 .context("HEAD is detached; set [merge] base in magi.toml")?,
466 };
467 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
468 if !git::is_clean(&repo).await? {
472 tracing::warn!(
473 "{} has uncommitted changes; they are not part of this run, \
474 which branches off {base_branch} ({})",
475 repo.display(),
476 &base_commit[..base_commit.len().min(8)]
477 );
478 }
479 let roles = config.resolve_roles()?;
480 let max_parallel = config.graph.max_parallel.max(1);
481 let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
482 state.event("start", format!("run {} created", state.id));
483 state.save()?;
484 Ok(Self {
485 state,
486 roles,
487 sem: Arc::new(Semaphore::new(max_parallel)),
488 pause: Pause::new(),
489 interrupt: Pause::new(),
490 })
491 }
492
493 pub async fn review(repo: &Path, branch: &str, config: Config) -> Result<Self> {
507 let repo = git::toplevel(repo).await?;
508 let missing = agent::missing_programs(&config.agents);
509 if !missing.is_empty() {
510 bail!(
511 "these agent programs are not on PATH: {}. Fix the roster in \
512 magi.toml or install them.",
513 missing.join(", ")
514 );
515 }
516 sync_review_branch(&repo, branch, &config.merge.remote).await?;
517 let base_branch = match config.merge.base.clone() {
518 Some(b) => b,
519 None => git::current_branch(&repo)
520 .await?
521 .context("HEAD is detached; set [merge] base in magi.toml")?,
522 };
523 if base_branch == branch {
524 bail!("`{branch}` is the base branch; there is nothing to review against");
525 }
526 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
527
528 let roles = config.resolve_roles()?;
529 let max_parallel = config.graph.max_parallel.max(1);
530 let log = git::log_oneline(&repo, &base_commit, branch)
533 .await
534 .unwrap_or_default();
535 let instruction = format!(
536 "Review the work already on branch `{branch}`. There is no task \
537 statement: what the change claims to do is whatever its commits \
538 say.\n\n{}",
539 if log.trim().is_empty() {
540 "(no commit messages)"
541 } else {
542 log.trim()
543 }
544 );
545 let mut state = RunState::new(
546 repo.clone(),
547 base_branch,
548 base_commit.clone(),
549 instruction,
550 config,
551 );
552
553 let worktree = state.worktree_root().join("under-review");
556 if let Some(parent) = worktree.parent() {
557 tokio::fs::create_dir_all(parent).await.ok();
558 }
559 let path = worktree.to_string_lossy().to_string();
560 git::git(&repo, &["worktree", "add", &path, branch])
561 .await
562 .with_context(|| {
563 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
564 })?;
565
566 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
567 .await
568 .unwrap_or(0);
569 if commits == 0 {
570 git::worktree_remove(&repo, &worktree).await.ok();
571 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
572 }
573 let files = git::changed_files(&worktree, &base_commit, "HEAD")
574 .await
575 .map(|f| f.len())
576 .unwrap_or(0);
577 if files == 0
578 && let (Ok(head_tree), Ok(base_tree)) = (
579 git::tree_of(&worktree, "HEAD").await,
580 git::tree_of(&worktree, &base_commit).await,
581 )
582 && head_tree == base_tree
583 {
584 let head = git::rev_parse(&worktree, "HEAD").await.unwrap_or_default();
585 git::worktree_remove(&repo, &worktree).await.ok();
586 bail!(
587 "`{branch}` at {} has a tree identical to base {}; this usually means \
588 the branch ref is stale (check `git rev-parse refs/heads/{branch}` \
589 against `{}/{branch}`) rather than an empty change",
590 short(&head),
591 short(&base_commit),
592 state.config.merge.remote
593 );
594 }
595 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
596 .await
597 .unwrap_or_default();
598
599 state.candidates.push(Candidate {
600 index: 0,
601 label: 'A',
602 agent: "(existing branch)".to_owned(),
605 branch: branch.to_owned(),
606 worktree,
607 summary: String::new(),
608 stat,
609 files,
610 commits,
611 empty: false,
612 failed: None,
613 verified_noop: None,
614 duration_ms: 0,
615 folded: false,
616 });
617 state.tally = Some(Tally {
618 first_choice: BTreeMap::from([('A', 0)]),
619 borda: BTreeMap::new(),
620 winner: 'A',
621 rankings: 0,
622 unanimous_initial: false,
623 deliberated: false,
624 changed_votes: 0,
625 unanimous_final: false,
626 tie_break: None,
627 judges: 0,
631 present: 0,
632 quorum: 0,
633 met_quorum: true,
634 uncontested: Some("review-only run: nothing competed".to_owned()),
635 });
636 state.status = RunStatus::Reviewing;
637 state.event(
638 "start",
639 format!(
640 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
641 state.id
642 ),
643 );
644 state.save()?;
645 Ok(Self {
646 state,
647 roles,
648 sem: Arc::new(Semaphore::new(max_parallel)),
649 pause: Pause::new(),
650 interrupt: Pause::new(),
651 })
652 }
653
654 pub fn resume(id: &str) -> Result<Self> {
656 let state = RunState::load(id)?;
657 let roles = state.config.resolve_roles()?;
658 let max_parallel = state.config.graph.max_parallel.max(1);
659 Ok(Self {
660 state,
661 roles,
662 sem: Arc::new(Semaphore::new(max_parallel)),
663 pause: Pause::new(),
664 interrupt: Pause::new(),
665 })
666 }
667
668 pub async fn execute(&mut self) -> Result<()> {
675 let result = self.execute_graph().await;
676 let ended = if result.is_err() {
677 Some(crate::notices::run_stopped(&self.state.id, &self.state))
678 } else {
679 crate::notices::run_ended(&self.state)
680 };
681 if let Some(notice) = ended {
682 crate::notices::raise(notice);
683 }
684 result
685 }
686
687 async fn execute_graph(&mut self) -> Result<()> {
688 self.state.parked = false;
693 self.state.clear_active();
700 let pid = std::process::id();
715 self.state.driver_pid = Some(pid);
716 self.state.driver_started_at = crate::proc::process_started_at(pid);
717 self.state.save()?;
718 if self.state.status == RunStatus::Stalled {
731 if self.recover_stall().await? {
732 self.finish_after_tally().await?;
733 } else {
734 self.state.save()?;
736 }
737 return Ok(());
738 }
739 if self.state.status == RunStatus::Landing {
749 self.run_land().await?;
750 self.settle_questions();
755 return Ok(());
756 }
757 self.prep().await?;
758 if self.park_here()? {
759 return Ok(());
760 }
761 self.advise().await?;
762 if self.park_here()? {
763 return Ok(());
764 }
765 self.implement().await?;
766 if self.park_here()? {
767 return Ok(());
768 }
769 if self.state.status == RunStatus::VerifiedNoop {
773 return Ok(());
774 }
775 self.judge().await?;
776 if self.park_here()? {
777 return Ok(());
778 }
779 self.deliberate().await?;
780 if self.park_here()? {
781 return Ok(());
782 }
783 self.vote().await?;
784 if self.park_here()? {
785 return Ok(());
786 }
787 self.tally()?;
788 if self.state.status == RunStatus::Stalled {
793 self.state.save()?;
797 return Ok(());
798 }
799 self.finish_after_tally().await?;
800 Ok(())
801 }
802
803 fn park_here(&mut self) -> Result<bool> {
810 if !self.pause.parked() && !self.interrupt.parked() {
815 return Ok(false);
816 }
817 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
818 Some(reason) => format!(
819 "parked after `{}` ({reason}) — resume to carry on from here",
820 self.state.status.as_str()
821 ),
822 None => format!(
823 "parked after `{}` — resume to carry on from here",
824 self.state.status.as_str()
825 ),
826 };
827 self.state.event("park", why);
828 self.state.parked = true;
829 self.state.save()?;
830 Ok(true)
831 }
832
833 pub fn on_pause(&mut self, pause: Pause) {
835 self.pause = pause;
836 }
837
838 pub fn watch_interrupt(&mut self, pause: Pause) {
844 self.interrupt = pause;
845 }
846
847 fn settle_questions(&mut self) {
867 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
868 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
869 }
870 }
871
872 async fn finish_after_tally(&mut self) -> Result<()> {
875 self.fold_losers().await?;
876 self.sync_to_base().await?;
881 self.review_loop().await?;
882 self.sync_to_base().await?;
883 self.gate().await?;
884 self.merge().await?;
885 self.state.save()?;
886 Ok(())
887 }
888
889 async fn prep(&mut self) -> Result<()> {
892 if !self.state.candidates.is_empty() {
893 return Ok(());
894 }
895 self.state.status = RunStatus::Prep;
896 let repo = self.state.repo.clone();
897 let base = self.state.base_commit.clone();
898 let root = self.state.worktree_root();
899 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
900
901 let hooks_dir = self.state.dir().join("hooks");
904 if self.state.config.blind.commit_msg_hook {
905 std::fs::create_dir_all(&hooks_dir)
906 .with_context(|| format!("create {}", hooks_dir.display()))?;
907 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
908 let path = hooks_dir.join("commit-msg");
909 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
910 make_executable(&path)?;
911 git::acquire_worktree_config(&repo).await?;
919 self.state.enabled_worktree_config = true;
920 }
921
922 for (index, (spec, label)) in self
923 .roles
924 .implementers
925 .clone()
926 .into_iter()
927 .zip(labels)
928 .enumerate()
929 {
930 let branch = self.state.branch_for(label);
931 let worktree = root.join(format!("cand-{label}"));
932 git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
933 if self.state.config.blind.commit_msg_hook {
934 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
935 }
936 git::local_exclude(&worktree, "/.magi/").await?;
937 self.state.candidates.push(Candidate {
938 index,
939 label,
940 agent: spec.id.clone(),
941 branch,
942 worktree,
943 summary: String::new(),
944 stat: String::new(),
945 files: 0,
946 commits: 0,
947 empty: false,
948 failed: None,
949 verified_noop: None,
950 duration_ms: 0,
951 folded: false,
952 });
953 }
954
955 for j in 1..=self.roles.judges.len() {
956 let wt = root.join(format!("judge-{j}"));
957 if !wt.exists() {
958 git::worktree_add_detached(&repo, &wt, &base).await?;
959 }
960 }
961
962 if self.state.config.graph.advise {
970 for k in 1..=self.state.config.graph.advisors {
971 let wt = root.join(format!("advisor-{k}"));
972 if !wt.exists() {
973 git::worktree_add_detached(&repo, &wt, &base).await?;
974 }
975 }
976 }
977
978 let authors: Vec<&str> = self
983 .roles
984 .implementers
985 .iter()
986 .map(|a| a.id.as_str())
987 .collect();
988 let overlap: Vec<String> = self
989 .roles
990 .judges
991 .iter()
992 .enumerate()
993 .filter(|(_, j)| authors.contains(&j.id.as_str()))
994 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
995 .collect();
996 if !overlap.is_empty() {
997 let note = format!(
998 "{} also authored a candidate; blind, but the panel is less \
999 independent than {} distinct agents would be",
1000 overlap.join(", "),
1001 self.roles.judges.len()
1002 );
1003 self.state.event("prep", note);
1004 }
1005
1006 self.state.event(
1007 "prep",
1008 format!(
1009 "{} candidates, {} judges, base {} ({})",
1010 self.state.candidates.len(),
1011 self.roles.judges.len(),
1012 &self.state.base_commit[..7.min(self.state.base_commit.len())],
1013 self.state.base_branch
1014 ),
1015 );
1016 self.state.status = RunStatus::Implementing;
1017 self.state.save()?;
1018 Ok(())
1019 }
1020
1021 async fn advise(&mut self) -> Result<()> {
1054 let implement_untouched = self
1055 .state
1056 .candidates
1057 .iter()
1058 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
1059 if !self.state.config.graph.advise || self.state.advise_attempted {
1060 return Ok(());
1061 }
1062 if !implement_untouched {
1063 self.state.event(
1064 "advise",
1065 "skipping the design-deliberation stage: at least one \
1066 candidate already shows implementation progress, so this \
1067 run is past the point the stage exists to run before"
1068 .to_owned(),
1069 );
1070 self.state.advise_attempted = true;
1071 self.state.save()?;
1072 return Ok(());
1073 }
1074 let run_id = self.state.id.clone();
1075 let prompts = self.state.config.prompts.clone();
1076 let instruction = self.state.instruction.clone();
1077 let language = self.state.config.graph.language.clone();
1078 let root = self.state.worktree_root();
1079 let n = self.state.config.graph.advisors;
1080 let where_recorded = self.state.dir().join("run.json");
1081
1082 let seats = match self.state.config.advisors() {
1083 Ok(seats) if !seats.is_empty() => seats,
1084 Ok(_) => {
1085 self.state.event(
1086 "advise",
1087 format!(
1088 "[graph] advisors is 0; skipping the design-deliberation \
1089 stage and continuing without a synthesis brief (see {})",
1090 where_recorded.display()
1091 ),
1092 );
1093 self.state.advise_attempted = true;
1094 self.state.save()?;
1095 return Ok(());
1096 }
1097 Err(e) => {
1098 self.state.event(
1099 "advise",
1100 format!(
1101 "could not resolve advisor seats ({e:#}); continuing \
1102 without a design-deliberation brief (see {})",
1103 where_recorded.display()
1104 ),
1105 );
1106 self.state.advise_attempted = true;
1107 self.state.save()?;
1108 return Ok(());
1109 }
1110 };
1111
1112 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1113 let artifacts = agent::artifacts_dir(&self.state.dir());
1114 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1115
1116 let mut jobs = Vec::new();
1117 for (i, spec) in seats.iter().cloned().enumerate() {
1118 let seat_key = format!("advisor-{}", i + 1);
1119 let seat = self.seat(&seat_key, &spec.id);
1120 jobs.push(SeatJob {
1121 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1122 spec,
1123 seat,
1124 cwd: worktrees[i % worktrees.len()].clone(),
1125 timeout,
1126 allow_write: false,
1127 sessions: false,
1128 artifacts: artifacts.clone(),
1129 stem: seat_key,
1130 });
1131 }
1132
1133 self.state.event(
1134 "advise",
1135 format!(
1136 "{} advisor seat(s) sketching a design in parallel",
1137 jobs.len()
1138 ),
1139 );
1140 let mut quota_losses = Vec::new();
1141 let cache = self.state.config.cache_dir();
1142 let ctx = WaveCtx {
1143 run: &run_id,
1144 node: "advise",
1145 prompts: &prompts,
1146 cache: cache.as_deref(),
1147 round: None,
1148 };
1149 let results = ask_json_wave::<Proposal>(
1150 jobs,
1151 Arc::clone(&self.sem),
1152 self.state.config.graph.retries,
1153 &ctx,
1154 &mut quota_losses,
1155 &mut self.state,
1156 &|p: &Proposal| p.validate(),
1157 )
1158 .await;
1159 self.state.quota.extend(quota_losses);
1160
1161 let mut records = Vec::with_capacity(results.len());
1162 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1163 let agent_id = seat.agent.clone();
1164 self.state.seats.insert(seat.key.clone(), seat);
1165 match res {
1166 Ok((proposal, out)) => {
1167 self.state
1168 .event("advise", format!("advisor-{} proposed a design", i + 1));
1169 records.push(advise::AdvisorRecord::proposed(
1170 i + 1,
1171 agent_id,
1172 proposal,
1173 out.duration_ms,
1174 ));
1175 }
1176 Err(e) => {
1177 self.state.event(
1178 "advise",
1179 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1180 );
1181 records.push(advise::AdvisorRecord::failed(
1182 i + 1,
1183 agent_id,
1184 e.to_string(),
1185 ));
1186 }
1187 }
1188 }
1189
1190 let mut advice = advise::Advice {
1191 records,
1192 synthesis: None,
1193 };
1194 if advice.proposals().is_empty() {
1195 self.state.event(
1196 "advise",
1197 "no advisor produced a usable proposal; continuing without a \
1198 synthesis brief"
1199 .to_owned(),
1200 );
1201 } else {
1202 match self
1203 .synthesize_brief(
1204 &advice,
1205 &instruction,
1206 &language,
1207 &worktrees[0],
1208 &artifacts,
1209 &run_id,
1210 &prompts,
1211 cache.as_deref(),
1212 )
1213 .await
1214 {
1215 Ok(Some(text)) => {
1216 self.state.event(
1217 "advise",
1218 "synthesized a design brief for the implementer".to_owned(),
1219 );
1220 advice.synthesis = Some(text);
1221 }
1222 Ok(None) => {
1223 self.state.event(
1224 "advise",
1225 "the synthesis seat produced nothing usable; continuing \
1226 without a design brief"
1227 .to_owned(),
1228 );
1229 }
1230 Err(e) => {
1231 self.state.event(
1232 "advise",
1233 format!("could not synthesize a design brief: {e:#}"),
1234 );
1235 }
1236 }
1237 }
1238 advise::apply_reflection(&mut advice);
1239
1240 self.state.advice = Some(advice);
1241 self.state.advise_attempted = true;
1242 self.state.save()?;
1243 Ok(())
1244 }
1245
1246 #[allow(clippy::too_many_arguments)]
1258 async fn synthesize_brief(
1259 &mut self,
1260 advice: &advise::Advice,
1261 instruction: &str,
1262 language: &str,
1263 cwd: &Path,
1264 artifacts: &Path,
1265 run_id: &str,
1266 prompts: &Prompts,
1267 cache: Option<&Path>,
1268 ) -> Result<Option<String>> {
1269 let want = self.state.config.roles.synthesizer.as_deref();
1270 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1271 let mut seat = self.seat("advise-synthesis", &spec.id);
1272 let proposals = advice.proposals();
1273 let mut prompt = prompt::with_overlay(
1274 prompt::synthesize_brief(instruction, &proposals, language),
1275 prompts.overlay("advise"),
1276 );
1277 if cache.is_some() {
1278 prompt.push('\n');
1283 prompt.push_str(&prompt::build_cache_note("advise", false));
1284 }
1285 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1286 let out = agent::invoke(
1287 &spec,
1288 &mut seat,
1289 &Invocation {
1290 cwd,
1291 prompt: &prompt,
1292 timeout,
1293 allow_write: false,
1294 sessions: false,
1295 artifacts,
1296 stem: "advise-synthesis",
1297 run: run_id,
1298 node: "advise",
1299 cache_dir: None,
1300 attachments: &[],
1301 },
1302 )
1303 .await?;
1304 self.state.seats.insert(seat.key.clone(), seat);
1305 if !out.usable() {
1306 return Ok(None);
1307 }
1308 let text =
1309 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1310 Ok((!text.trim().is_empty()).then_some(text))
1311 }
1312
1313 async fn implement(&mut self) -> Result<()> {
1316 let run_id = self.state.id.clone();
1321 let prompts = self.state.config.prompts.clone();
1322 let todo: Vec<usize> = self
1323 .state
1324 .candidates
1325 .iter()
1326 .enumerate()
1327 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1328 .map(|(i, _)| i)
1329 .collect();
1330 if todo.is_empty() {
1331 return self.after_implement();
1332 }
1333 self.state.status = RunStatus::Implementing;
1334
1335 let language = self.state.config.graph.language.clone();
1336 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1337 let sessions = self.state.config.graph.sessions;
1338 let artifacts = agent::artifacts_dir(&self.state.dir());
1339 let brief = self
1343 .state
1344 .advice
1345 .as_ref()
1346 .and_then(|a| a.synthesis.as_deref())
1347 .map(str::to_owned);
1348
1349 let mut jobs = Vec::new();
1350 for &i in &todo {
1351 let (index, label, worktree) = {
1352 let c = &self.state.candidates[i];
1353 (c.index, c.label, c.worktree.clone())
1354 };
1355 let spec = self.roles.implementers[index].clone();
1356 let seat_key = format!("impl-{label}");
1357 let seat = self.seat(&seat_key, &spec.id);
1358 let instruction = self.state.instruction.clone();
1359 jobs.push(SeatJob {
1360 spec,
1361 seat,
1362 prompt: prompt::implement(
1363 &instruction,
1364 &worktree.to_string_lossy(),
1365 &language,
1366 brief.as_deref(),
1367 ),
1368 cwd: worktree,
1369 timeout,
1370 allow_write: true,
1371 sessions,
1372 artifacts: artifacts.clone(),
1373 stem: format!("impl-{label}"),
1374 });
1375 }
1376
1377 self.state.event(
1378 "implement",
1379 format!("{} candidates in parallel", jobs.len()),
1380 );
1381 let mut sent = jobs.clone();
1387 let cache = self.state.config.cache_dir();
1388 let ctx = WaveCtx {
1389 run: &run_id,
1390 node: "implement",
1391 prompts: &prompts,
1392 cache: cache.as_deref(),
1393 round: None,
1394 };
1395 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1396 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1397 .await;
1398 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1399 .await;
1400 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1401 .await;
1402
1403 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1404 let seat_key = seat.key.clone();
1405 let agent = seat.agent.clone();
1414 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1415 self.state.seats.insert(seat.key.clone(), seat);
1416 let label = self.state.candidates[i].label;
1417 let worktree = self.state.candidates[i].worktree.clone();
1418 let base = self.state.base_commit.clone();
1419
1420 let (summary, duration, failed, verified_claim) = match out {
1421 AgentOutcome::Ok(o) => {
1422 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1423 let failed = (!o.usable()).then(|| {
1424 if o.timed_out {
1425 "agent timed out".to_owned()
1426 } else {
1427 format!("agent exited with {:?}", o.exit_code)
1428 }
1429 });
1430 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1431 (text, o.duration_ms, failed, verified_claim)
1432 }
1433 AgentOutcome::Dropped(o) => {
1439 let why = o
1440 .dropped
1441 .as_ref()
1442 .map(|d| d.why.as_str())
1443 .unwrap_or("the CLI ended the stream without delivering its answer");
1444 (
1445 String::new(),
1446 o.duration_ms,
1447 Some(format!("the CLI dropped the stream ({why})")),
1448 None,
1449 )
1450 }
1451 AgentOutcome::Quota(o) => {
1452 self.state.quota.push(QuotaLoss {
1453 seat: seat_key,
1454 node: "implement".to_owned(),
1455 at: Timestamp::now(),
1456 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1457 });
1458 (
1459 String::new(),
1460 o.duration_ms,
1461 Some("rate limited (quota); produced no change".to_owned()),
1462 None,
1463 )
1464 }
1465 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1466 };
1467
1468 let rescued = match git::rescue_commit(
1471 &worktree,
1472 &format!("magi: candidate {label} (uncommitted work)"),
1473 )
1474 .await
1475 {
1476 Ok(r) => {
1477 self.state.note_withheld("implement", &r.withheld);
1478 r.committed
1479 }
1480 Err(_) => false,
1481 };
1482 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1483 .await
1484 .unwrap_or(0);
1485 let patch = git::diff(&worktree, &base, "HEAD")
1486 .await
1487 .unwrap_or_default();
1488 let stat = git::diff_stat(&worktree, &base, "HEAD")
1489 .await
1490 .unwrap_or_default();
1491 let files = git::changed_files(&worktree, &base, "HEAD")
1492 .await
1493 .map(|f| f.len())
1494 .unwrap_or(0);
1495 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1496
1497 let c = &mut self.state.candidates[i];
1498 if !exhausted_the_fallback_chain {
1499 c.agent = agent;
1500 }
1501 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1502 c.stat = stat;
1503 c.files = files;
1504 c.commits = commits;
1505 c.duration_ms = duration;
1506 c.empty = commits == 0 || patch.trim().is_empty();
1507 c.failed = match failed {
1510 Some(_) if c.empty => failed,
1511 _ => None,
1512 };
1513 c.verified_noop = if c.empty { verified_claim } else { None };
1518 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1519 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1520 (None, true, Some(_), _) => {
1521 format!("candidate {label}: no change produced (agent-verified no-op)")
1522 }
1523 (None, true, None, _) => format!("candidate {label}: no change produced"),
1524 (None, false, _, true) => {
1525 format!(
1526 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1527 )
1528 }
1529 (None, false, _, false) => {
1530 format!("candidate {label}: {files} files, {commits} commits")
1531 }
1532 };
1533 self.state.event("implement", note);
1534 self.state.save()?;
1535 }
1536
1537 self.after_implement()
1538 }
1539
1540 async fn resume_undelivered(
1568 &mut self,
1569 results: &mut [(usize, SeatState, AgentOutcome)],
1570 sent: &[SeatJob],
1571 prompts: &Prompts,
1572 run_id: &str,
1573 ) {
1574 for (wi, seat, out) in results.iter_mut() {
1575 let Some(dropped) = (match &*out {
1576 AgentOutcome::Dropped(o) => o.dropped.clone(),
1577 _ => None,
1578 }) else {
1579 continue;
1580 };
1581 let Some(job) = sent.get(*wi) else { continue };
1582 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1584 self.state.event(
1585 "implement",
1586 format!(
1587 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1588 work is in the tree",
1589 seat.key, dropped.output_tokens, dropped.why
1590 ),
1591 );
1592 continue;
1593 }
1594 if !has_context(&job.spec, seat, job.sessions) {
1602 self.state.event(
1603 "implement",
1604 format!(
1605 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1606 is no session left to resume",
1607 seat.key, dropped.output_tokens, dropped.why
1608 ),
1609 );
1610 continue;
1611 }
1612 self.state.event(
1613 "implement",
1614 format!(
1615 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1616 conversation",
1617 seat.key, dropped.output_tokens, dropped.why
1618 ),
1619 );
1620 let mut retry = job.clone();
1621 retry.seat = seat.clone();
1622 retry.prompt = prompt::resume_after_drop(&dropped.why);
1623 retry.timeout = retry_budget(job.timeout, true);
1624 retry.stem = format!("{}-resume", job.stem);
1625 let cache = self.state.config.cache_dir();
1626 let ctx = WaveCtx {
1627 run: run_id,
1628 node: "implement",
1629 prompts,
1630 cache: cache.as_deref(),
1631 round: None,
1632 };
1633 let (resumed_seat, resumed) =
1634 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1635 *seat = resumed_seat;
1636 *out = resumed;
1637 }
1638 }
1639
1640 async fn resume_quota_losses(
1702 &mut self,
1703 results: &mut [(usize, SeatState, AgentOutcome)],
1704 sent: &mut [SeatJob],
1705 prompts: &Prompts,
1706 run_id: &str,
1707 ) {
1708 let instruction = self.state.instruction.clone();
1709 let language = self.state.config.graph.language.clone();
1710 let brief = self
1711 .state
1712 .advice
1713 .as_ref()
1714 .and_then(|a| a.synthesis.as_deref())
1715 .map(str::to_owned);
1716 for (wi, seat, out) in results.iter_mut() {
1717 let Some(job) = sent.get_mut(*wi) else {
1718 continue;
1719 };
1720 let start = self
1725 .roles
1726 .implementer_roster
1727 .iter()
1728 .position(|s| s.id == job.spec.id)
1729 .unwrap_or(0);
1730 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1731 let mut fallback_attempt = 0usize;
1732 while matches!(&*out, AgentOutcome::Quota(_)) {
1733 let Some(next) =
1734 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1735 .cloned()
1736 else {
1737 break;
1738 };
1739 tried.insert(next.id.clone());
1740 fallback_attempt += 1;
1741
1742 if let Ok(r) = git::rescue_commit(
1743 &job.cwd,
1744 &format!(
1745 "magi: candidate {} (uncommitted work before quota fallback)",
1746 seat.key
1747 ),
1748 )
1749 .await
1750 {
1751 self.state.note_withheld("implement", &r.withheld);
1752 }
1753
1754 self.state.event(
1755 "implement",
1756 format!(
1757 "{}: rate limited (quota) on {}; retrying with {}",
1758 seat.key, seat.agent, next.id
1759 ),
1760 );
1761
1762 let new_seat = self.seat(&seat.key, &next.id);
1763 job.spec = next.clone();
1771 let mut retry = job.clone();
1772 retry.seat = new_seat;
1773 retry.prompt = prompt::implement(
1774 &instruction,
1775 &job.cwd.to_string_lossy(),
1776 &language,
1777 brief.as_deref(),
1778 );
1779 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1780 let cache = self.state.config.cache_dir();
1781 let ctx = WaveCtx {
1782 run: run_id,
1783 node: "implement",
1784 prompts,
1785 cache: cache.as_deref(),
1786 round: None,
1787 };
1788 let (fallback_seat, fallback_out) = run_one(
1789 retry,
1790 Arc::clone(&self.sem),
1791 &ctx,
1792 &mut self.state,
1793 fallback_attempt,
1794 )
1795 .await;
1796 *seat = fallback_seat;
1797 *out = fallback_out;
1798 }
1799 }
1800 }
1801
1802 async fn resume_unconfirmed_commands(
1826 &mut self,
1827 results: &mut [(usize, SeatState, AgentOutcome)],
1828 sent: &[SeatJob],
1829 prompts: &Prompts,
1830 run_id: &str,
1831 ) {
1832 for (wi, seat, out) in results.iter_mut() {
1833 let AgentOutcome::Ok(o) = &*out else {
1834 continue;
1835 };
1836 if !has_unconfirmed_command(&o.commands) {
1837 continue;
1838 }
1839 let Some(job) = sent.get(*wi) else { continue };
1840 if !has_context(&job.spec, seat, job.sessions) {
1841 self.state.event(
1842 "implement",
1843 format!(
1844 "{}: the reply named a command whose own CLI never confirmed the exit \
1845 status of, but there is no session left to resume",
1846 seat.key
1847 ),
1848 );
1849 continue;
1850 }
1851 self.state.event(
1852 "implement",
1853 format!(
1854 "{}: the reply named a command whose own CLI never confirmed the exit \
1855 status of; resuming the conversation",
1856 seat.key
1857 ),
1858 );
1859 let mut retry = job.clone();
1860 retry.seat = seat.clone();
1861 retry.prompt = prompt::resume_incomplete(
1862 "a command in your last reply had no confirmed exit status",
1863 );
1864 retry.timeout = retry_budget(job.timeout, true);
1865 retry.stem = format!("{}-confirm", job.stem);
1866 let cache = self.state.config.cache_dir();
1867 let ctx = WaveCtx {
1868 run: run_id,
1869 node: "implement",
1870 prompts,
1871 cache: cache.as_deref(),
1872 round: None,
1873 };
1874 let (resumed_seat, resumed) =
1875 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1876 *seat = resumed_seat;
1877 *out = resumed;
1878 }
1879 }
1880
1881 async fn continue_fix_report(
1902 &mut self,
1903 mut seat: SeatState,
1904 parse_err: String,
1905 job: &SeatJob,
1906 prompts: &Prompts,
1907 run_id: &str,
1908 round: usize,
1909 ) -> (
1910 SeatState,
1911 Option<FixReport>,
1912 Option<String>,
1913 ContinuationRecord,
1914 ) {
1915 let mut last_err = parse_err;
1916 let mut cumulative_wait_ms = 0u64;
1917 let mut attempts = 0usize;
1918 loop {
1919 if !has_context(&job.spec, &seat, job.sessions) {
1920 self.state.event(
1921 "fix",
1922 format!(
1923 "round {round}: fixer's reply had no adoption report ({last_err}); no \
1924 session left to resume into"
1925 ),
1926 );
1927 let outcome = if attempts == 0 {
1928 ContinuationOutcome::NoSession
1929 } else {
1930 ContinuationOutcome::Exhausted
1931 };
1932 return (
1933 seat,
1934 None,
1935 Some(format!("unparsable fix report: {last_err}")),
1936 ContinuationRecord {
1937 attempts,
1938 cumulative_wait_ms,
1939 outcome,
1940 },
1941 );
1942 }
1943 if attempts >= MAX_FIX_CONTINUATIONS {
1944 self.state.event(
1945 "fix",
1946 format!(
1947 "round {round}: fixer's reply still had no adoption report after \
1948 {attempts} continuation(s) ({last_err}); giving up"
1949 ),
1950 );
1951 return (
1952 seat,
1953 None,
1954 Some(format!(
1955 "unparsable fix report after {attempts} continuation(s): {last_err}"
1956 )),
1957 ContinuationRecord {
1958 attempts,
1959 cumulative_wait_ms,
1960 outcome: ContinuationOutcome::Exhausted,
1961 },
1962 );
1963 }
1964 attempts += 1;
1965 self.state.event(
1966 "fix",
1967 format!(
1968 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
1969 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
1970 ),
1971 );
1972 let mut retry = job.clone();
1973 retry.seat = seat.clone();
1974 retry.prompt = prompt::resume_incomplete(&last_err);
1975 retry.timeout = retry_budget(job.timeout, true);
1976 retry.stem = format!("{}-continue{attempts}", job.stem);
1977 let cache = self.state.config.cache_dir();
1978 let ctx = WaveCtx {
1979 run: run_id,
1980 node: "fix",
1981 prompts,
1982 cache: cache.as_deref(),
1983 round: Some(round),
1984 };
1985 let (resumed_seat, resumed_out) = run_one(
1986 retry,
1987 Arc::clone(&self.sem),
1988 &ctx,
1989 &mut self.state,
1990 attempts,
1991 )
1992 .await;
1993 seat = resumed_seat;
1994 match resumed_out {
1995 AgentOutcome::Ok(o) => {
1996 cumulative_wait_ms += o.duration_ms;
1997 match verdict::extract_json::<FixReport>(&o.text) {
1998 Ok(report) if !has_unconfirmed_command(&o.commands) => {
1999 self.state.event(
2000 "fix",
2001 format!(
2002 "round {round}: fixer's adoption report recovered after \
2003 {attempts} continuation(s)"
2004 ),
2005 );
2006 return (
2007 seat,
2008 Some(report),
2009 None,
2010 ContinuationRecord {
2011 attempts,
2012 cumulative_wait_ms,
2013 outcome: ContinuationOutcome::Resumed,
2014 },
2015 );
2016 }
2017 Ok(_) => {
2025 last_err = "the reply parsed, but it reported a command whose own CLI \
2026 never confirmed an exit status"
2027 .to_owned();
2028 }
2029 Err(e) => last_err = e.to_string(),
2030 }
2031 }
2032 AgentOutcome::Quota(o) => {
2033 cumulative_wait_ms += o.duration_ms;
2034 self.state.quota.push(QuotaLoss {
2035 seat: seat.key.clone(),
2036 node: "fix".to_owned(),
2037 at: Timestamp::now(),
2038 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2039 });
2040 self.state.event(
2041 "fix",
2042 format!(
2043 "round {round}: continuation rate limited (quota); not retrying now"
2044 ),
2045 );
2046 return (
2047 seat,
2048 None,
2049 Some("rate limited (quota) while recovering the fix report".to_owned()),
2050 ContinuationRecord {
2051 attempts,
2052 cumulative_wait_ms,
2053 outcome: ContinuationOutcome::QuotaLost,
2054 },
2055 );
2056 }
2057 AgentOutcome::Dropped(o) => {
2058 cumulative_wait_ms += o.duration_ms;
2059 let why = o
2060 .dropped
2061 .as_ref()
2062 .map(|d| d.why.as_str())
2063 .unwrap_or("the CLI ended the stream without delivering its answer");
2064 last_err = format!("the CLI dropped the stream ({why})");
2065 }
2066 AgentOutcome::Failed(e) => last_err = e,
2067 }
2068 }
2069 }
2070
2071 fn after_implement(&mut self) -> Result<()> {
2072 if self.state.leaks.is_empty() {
2074 let cfg = self.state.config.blind.clone();
2075 let mut leaks = Vec::new();
2076 for c in &self.state.candidates {
2077 let Some(patch) =
2078 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
2079 else {
2080 continue;
2081 };
2082 leaks.extend(blind::scan(
2083 &format!("candidate {} patch", c.label),
2084 &patch,
2085 &cfg.vendor_tokens,
2086 ));
2087 }
2088 if !leaks.is_empty() {
2089 let summary = leaks
2090 .iter()
2091 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2092 .collect::<Vec<_>>()
2093 .join(", ");
2094 match cfg.on_leak {
2095 LeakPolicy::Fail => {
2096 self.state.status = RunStatus::Failed;
2097 self.state
2098 .event("blind", format!("vendor text in a patch: {summary}"));
2099 self.state.leaks = leaks;
2100 self.state.save()?;
2101 self.settle_questions();
2102 bail!(
2103 "blind.on_leak = \"fail\" and vendor text reached a \
2104 judged patch: {summary}"
2105 );
2106 }
2107 LeakPolicy::Redact => self.state.event(
2108 "blind",
2109 format!("redacting vendor text for judging: {summary}"),
2110 ),
2111 LeakPolicy::Warn => self.state.event(
2112 "blind",
2113 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2114 ),
2115 }
2116 self.state.leaks = leaks;
2117 }
2118 }
2119
2120 if self.state.viable().is_empty() {
2121 if self.state.all_candidates_verified_noop() {
2122 self.state.status = RunStatus::VerifiedNoop;
2133 self.state.save()?;
2134 self.settle_questions();
2135 return Ok(());
2136 }
2137 self.state.status = RunStatus::Failed;
2138 self.state.save()?;
2139 self.settle_questions();
2140 bail!("no candidate produced a change; nothing to judge");
2141 }
2142 self.state.status = RunStatus::Judging;
2143 self.state.save()?;
2144 Ok(())
2145 }
2146
2147 async fn judge(&mut self) -> Result<()> {
2150 let run_id = self.state.id.clone();
2155 let prompts = self.state.config.prompts.clone();
2156 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2157 return Ok(());
2158 }
2159 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2160 if viable.len() == 1 {
2161 self.state.judge_skipped = true;
2168 self.state.event(
2169 "judge",
2170 format!(
2171 "only candidate {} produced a change; judging skipped",
2172 viable[0].label
2173 ),
2174 );
2175 self.state.save()?;
2176 return Ok(());
2177 }
2178 self.state.status = RunStatus::Judging;
2179
2180 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2181 let language = self.state.config.graph.language.clone();
2182 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2183 let sessions = self.state.config.graph.sessions;
2184 let artifacts = agent::artifacts_dir(&self.state.dir());
2185 let root = self.state.worktree_root();
2186 let base_short = short(&self.state.base_commit);
2187
2188 let mut jobs = Vec::new();
2189 let mut orders = Vec::new();
2190 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2191 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2192 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2193 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2194 let seat_key = format!("judge-{}", j + 1);
2195 let seat = self.seat(&seat_key, &spec.id);
2196 jobs.push(SeatJob {
2197 prompt: prompt::judge(
2198 &self.state.instruction,
2199 &views,
2200 self.roles.judges.len(),
2201 &base_short,
2202 &language,
2203 ),
2204 spec,
2205 seat,
2206 cwd: root.join(format!("judge-{}", j + 1)),
2207 timeout,
2208 allow_write: false,
2209 sessions,
2210 artifacts: artifacts.clone(),
2211 stem: format!("judge-{}", j + 1),
2212 });
2213 }
2214
2215 self.state.event(
2216 "judge",
2217 format!(
2218 "{} judges ranking {} candidates blind",
2219 jobs.len(),
2220 viable.len()
2221 ),
2222 );
2223 let labels_for_check = labels.clone();
2224 let mut quota_losses = Vec::new();
2225 let cache = self.state.config.cache_dir();
2226 let ctx = WaveCtx {
2227 run: &run_id,
2228 node: "judge",
2229 prompts: &prompts,
2230 cache: cache.as_deref(),
2231 round: None,
2232 };
2233 let results = ask_json_wave::<Ranking>(
2234 jobs,
2235 Arc::clone(&self.sem),
2236 self.state.config.graph.retries,
2237 &ctx,
2238 &mut quota_losses,
2239 &mut self.state,
2240 &move |r: &Ranking| r.validate(&labels_for_check),
2241 )
2242 .await;
2243 self.state.quota.extend(quota_losses);
2244
2245 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2246 let agent_id = seat.agent.clone();
2247 self.state.seats.insert(seat.key.clone(), seat);
2248 let mut record = Judgement {
2249 judge: j + 1,
2250 seat: format!("judge-{}", j + 1),
2251 agent: agent_id,
2252 ranking: Vec::new(),
2253 reasons: BTreeMap::new(),
2254 confidence: None,
2255 order: orders[j].clone(),
2256 failed: None,
2257 duration_ms: 0,
2258 };
2259 match res {
2260 Ok((ranking, out)) => {
2261 record.ranking = ranking.normalized();
2262 record.reasons = ranking.reasons;
2263 record.confidence = ranking.confidence;
2264 record.duration_ms = out.duration_ms;
2265 self.state.event(
2266 "judge",
2267 format!(
2268 "judge {} ranked {}",
2269 j + 1,
2270 record.ranking.iter().collect::<String>()
2271 ),
2272 );
2273 }
2274 Err(e) => {
2275 record.failed = Some(e.to_string());
2276 self.state
2277 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2278 }
2279 }
2280 self.state.judgements.push(record);
2281 self.state.save()?;
2282 }
2283 Ok(())
2284 }
2285
2286 async fn deliberate(&mut self) -> Result<()> {
2289 let run_id = self.state.id.clone();
2294 let prompts = self.state.config.prompts.clone();
2295 if !self.state.deliberation.is_empty() {
2296 return Ok(());
2297 }
2298 let tops: Vec<char> = self
2299 .state
2300 .judgements
2301 .iter()
2302 .filter_map(|j| j.ranking.first().copied())
2303 .collect();
2304 let rounds = self.state.config.graph.deliberate_rounds;
2305 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2306 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2307 self.state.event(
2308 "deliberate",
2309 format!("judges agreed on {} outright; no deliberation", tops[0]),
2310 );
2311 }
2312 self.state.status = RunStatus::Voting;
2313 self.state.save()?;
2314 return Ok(());
2315 }
2316
2317 self.state.status = RunStatus::Deliberating;
2318 self.state.event(
2319 "deliberate",
2320 format!(
2321 "split: first choices were {} — opening {rounds} round(s)",
2322 tops.iter().collect::<String>()
2323 ),
2324 );
2325
2326 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2327 let language = self.state.config.graph.language.clone();
2328 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2329 let sessions = self.state.config.graph.sessions;
2330 let artifacts = agent::artifacts_dir(&self.state.dir());
2331 let root = self.state.worktree_root();
2332 let base_short = short(&self.state.base_commit);
2333
2334 for round in 1..=rounds {
2338 let mut turns: Vec<DeliberationTurn> = Vec::new();
2339 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2340 if self.state.judgements[j].failed.is_some() {
2341 continue;
2342 }
2343 let seat_key = format!("judge-{}", j + 1);
2344 let mut seat = self.seat(&seat_key, &spec.id);
2345 let transcript = self.transcript(&turns, j);
2346 let context = if has_context(&spec, &seat, sessions) {
2347 None
2348 } else {
2349 Some(self.candidate_block(&viable, &base_short))
2350 };
2351 let text = prompt::deliberate(
2352 &self.state.instruction,
2353 context.as_deref(),
2354 &transcript,
2355 round,
2356 rounds,
2357 &language,
2358 );
2359 let job = SeatJob {
2360 spec,
2361 seat: seat.clone(),
2362 prompt: text,
2363 cwd: root.join(format!("judge-{}", j + 1)),
2364 timeout,
2365 allow_write: false,
2366 sessions,
2367 artifacts: artifacts.clone(),
2368 stem: format!("delib-{round}-judge-{}", j + 1),
2369 };
2370 let cache = self.state.config.cache_dir();
2371 let ctx = WaveCtx {
2372 run: &run_id,
2373 node: "deliberate",
2374 prompts: &prompts,
2375 cache: cache.as_deref(),
2376 round: None,
2377 };
2378 let (updated, out) =
2379 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2380 seat = updated;
2381 let agent_id = seat.agent.clone();
2382 let seat_key = seat.key.clone();
2383 self.state.seats.insert(seat.key.clone(), seat);
2384 let body = match out {
2385 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2386 AgentOutcome::Dropped(o) => {
2390 let why =
2391 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2392 "the CLI ended the stream without delivering its answer",
2393 );
2394 self.state.event(
2395 "deliberate",
2396 format!(
2397 "judge {} skipped: the CLI dropped the stream ({why})",
2398 j + 1
2399 ),
2400 );
2401 continue;
2402 }
2403 AgentOutcome::Quota(o) => {
2404 self.state.quota.push(QuotaLoss {
2405 seat: seat_key,
2406 node: "deliberate".to_owned(),
2407 at: Timestamp::now(),
2408 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2409 });
2410 self.state.event(
2411 "deliberate",
2412 format!("judge {} skipped: rate limited (quota)", j + 1),
2413 );
2414 continue;
2415 }
2416 AgentOutcome::Failed(e) => {
2417 self.state
2418 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2419 continue;
2420 }
2421 };
2422 let tentative = verdict::extract_json::<Position>(&body)
2423 .ok()
2424 .and_then(|p| p.tentative)
2425 .and_then(|s| s.trim().chars().next())
2426 .map(|c| c.to_ascii_uppercase());
2427 self.state.event(
2428 "deliberate",
2429 format!(
2430 "round {round}: judge {} now favours {}",
2431 j + 1,
2432 tentative.map_or("—".to_owned(), |c| c.to_string())
2433 ),
2434 );
2435 turns.push(DeliberationTurn {
2436 judge: j + 1,
2437 agent: agent_id,
2438 body: blind::sanitize_prose(&body, &self.state.config.blind),
2439 tentative,
2440 });
2441 }
2442 self.state
2443 .deliberation
2444 .push(DeliberationRound { round, turns });
2445 self.state.save()?;
2446 }
2447
2448 self.state.status = RunStatus::Voting;
2449 self.state.save()?;
2450 Ok(())
2451 }
2452
2453 async fn vote(&mut self) -> Result<()> {
2456 let run_id = self.state.id.clone();
2461 let prompts = self.state.config.prompts.clone();
2462 if !self.state.votes.is_empty() {
2463 return Ok(());
2464 }
2465 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2466 if viable.len() == 1 {
2467 return Ok(());
2468 }
2469 self.state.status = RunStatus::Voting;
2470
2471 let language = self.state.config.graph.language.clone();
2472 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2473 let sessions = self.state.config.graph.sessions;
2474 let artifacts = agent::artifacts_dir(&self.state.dir());
2475 let root = self.state.worktree_root();
2476 let base_short = short(&self.state.base_commit);
2477 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2478
2479 let mut jobs = Vec::new();
2480 let mut seats_at = Vec::new();
2481 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2482 if self
2483 .state
2484 .judgements
2485 .get(j)
2486 .is_some_and(|r| r.failed.is_some())
2487 {
2488 continue;
2489 }
2490 let seat_key = format!("judge-{}", j + 1);
2491 let seat = self.seat(&seat_key, &spec.id);
2492 let mut text = prompt::final_vote(&viable, &language);
2493 if !has_context(&spec, &seat, sessions) {
2494 text = format!(
2495 "{}\n\n# Candidates\n\n{}",
2496 text,
2497 self.candidate_block(&candidates, &base_short)
2498 );
2499 }
2500 jobs.push(SeatJob {
2501 spec,
2502 seat,
2503 prompt: text,
2504 cwd: root.join(format!("judge-{}", j + 1)),
2505 timeout,
2506 allow_write: false,
2507 sessions,
2508 artifacts: artifacts.clone(),
2509 stem: format!("vote-judge-{}", j + 1),
2510 });
2511 seats_at.push(j);
2512 }
2513
2514 self.state.event(
2515 "vote",
2516 format!(
2517 "collecting {} final votes one by one, privately",
2518 jobs.len()
2519 ),
2520 );
2521 let allowed = viable.clone();
2522 let mut quota_losses = Vec::new();
2523 let cache = self.state.config.cache_dir();
2524 let ctx = WaveCtx {
2525 run: &run_id,
2526 node: "vote",
2527 prompts: &prompts,
2528 cache: cache.as_deref(),
2529 round: None,
2530 };
2531 let results = ask_json_wave::<FinalVote>(
2532 jobs,
2533 Arc::clone(&self.sem),
2534 self.state.config.graph.retries,
2535 &ctx,
2536 &mut quota_losses,
2537 &mut self.state,
2538 &move |v: &FinalVote| match v.label() {
2539 Some(c) if allowed.contains(&c) => Ok(()),
2540 other => bail!("vote {other:?} is not one of {allowed:?}"),
2541 },
2542 )
2543 .await;
2544 self.state.quota.extend(quota_losses);
2545
2546 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2547 let agent_id = seat.agent.clone();
2548 self.state.seats.insert(seat.key.clone(), seat);
2549 let initial = self
2550 .state
2551 .judgements
2552 .get(j)
2553 .and_then(|r| r.ranking.first().copied());
2554 let mut record = VoteRecord {
2555 judge: j + 1,
2556 agent: agent_id,
2557 vote: None,
2558 reason: String::new(),
2559 changed: false,
2560 };
2561 match res {
2562 Ok((v, _)) => {
2563 record.vote = v.label();
2564 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2565 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2566 self.state.event(
2567 "vote",
2568 format!(
2569 "judge {} voted {}{}",
2570 j + 1,
2571 record.vote.unwrap_or('?'),
2572 if record.changed { " (changed)" } else { "" }
2573 ),
2574 );
2575 }
2576 Err(e) => {
2577 self.state
2578 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2579 }
2580 }
2581 self.state.votes.push(record);
2582 self.state.save()?;
2583 }
2584 Ok(())
2585 }
2586
2587 fn tally(&mut self) -> Result<()> {
2590 if self.state.tally.is_some() {
2591 return Ok(());
2592 }
2593 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2594 let tops: Vec<char> = self
2595 .state
2596 .judgements
2597 .iter()
2598 .filter_map(|j| j.ranking.first().copied())
2599 .collect();
2600 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2601
2602 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2605 let mut cast: Vec<char> = Vec::new();
2606 for (i, j) in self.state.judgements.iter().enumerate() {
2607 let vote = self
2608 .state
2609 .votes
2610 .iter()
2611 .find(|v| v.judge == i + 1)
2612 .and_then(|v| v.vote)
2613 .or_else(|| j.ranking.first().copied());
2614 if let Some(v) = vote {
2615 *first_choice.entry(v).or_insert(0) += 1;
2616 cast.push(v);
2617 }
2618 }
2619
2620 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2621 for j in &self.state.judgements {
2622 let n = j.ranking.len();
2623 for (pos, label) in j.ranking.iter().enumerate() {
2624 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2625 }
2626 }
2627
2628 let best = first_choice.values().copied().max().unwrap_or(0);
2629 let mut leaders: Vec<char> = first_choice
2630 .iter()
2631 .filter(|(_, v)| **v == best)
2632 .map(|(k, _)| *k)
2633 .collect();
2634 let mut tie_break = None;
2635 if leaders.len() > 1 {
2636 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2637 let borda_leaders: Vec<char> = leaders
2638 .iter()
2639 .copied()
2640 .filter(|l| borda[l] == top_borda)
2641 .collect();
2642 tie_break = Some(if borda_leaders.len() == 1 {
2643 format!(
2644 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2645 leaders.len()
2646 )
2647 } else {
2648 format!(
2649 "{} way tie on both first-choice votes and Borda points, broken by label order",
2650 leaders.len()
2651 )
2652 });
2653 leaders = borda_leaders;
2654 leaders.sort_unstable();
2655 }
2656 let winner = *leaders
2657 .first()
2658 .or(viable.first())
2659 .context("no candidate to declare a winner from")?;
2660
2661 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2662 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2663 let deliberated = !self.state.deliberation.is_empty();
2664
2665 let quota_seats: std::collections::BTreeSet<&str> =
2669 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2670 let mut present = 0usize;
2671 for (i, j) in self.state.judgements.iter().enumerate() {
2672 if quota_seats.contains(j.seat.as_str()) {
2673 continue;
2674 }
2675 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2676 let voted = self
2677 .state
2678 .votes
2679 .iter()
2680 .any(|v| v.judge == i + 1 && v.vote.is_some());
2681 if ranked || voted {
2682 present += 1;
2683 }
2684 }
2685 let needs_quorum = viable.len() > 1;
2691 let judges_total = if needs_quorum {
2692 self.roles.judges.len()
2693 } else {
2694 0
2695 };
2696 let quorum = if needs_quorum {
2697 judges_total / 2 + 1
2698 } else {
2699 0
2700 };
2701 let met_quorum = !needs_quorum || present >= quorum;
2702 let uncontested = (!needs_quorum).then(|| {
2703 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2704 });
2705
2706 self.state.event(
2707 "tally",
2708 match &uncontested {
2709 Some(reason) => format!("winner {winner} — {reason}"),
2710 None => format!(
2711 "winner {winner} — votes {} | initial {} | {} changed | \
2712 {present}/{judges_total} judges{}",
2713 first_choice
2714 .iter()
2715 .map(|(k, v)| format!("{k}:{v}"))
2716 .collect::<Vec<_>>()
2717 .join(" "),
2718 if unanimous_initial {
2719 "unanimous"
2720 } else {
2721 "split"
2722 },
2723 changed_votes,
2724 if met_quorum {
2725 String::new()
2726 } else {
2727 format!(" — below quorum ({quorum} required)")
2728 },
2729 ),
2730 },
2731 );
2732 if !met_quorum {
2733 self.state.event(
2734 "stall",
2735 format!(
2736 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2737 the run stops here, resumable"
2738 ),
2739 );
2740 }
2741 self.state.tally = Some(Tally {
2742 first_choice,
2743 borda,
2744 winner,
2745 rankings: tops.len(),
2746 unanimous_initial,
2747 deliberated,
2748 changed_votes,
2749 unanimous_final,
2750 tie_break,
2751 judges: judges_total,
2752 present,
2753 quorum,
2754 met_quorum,
2755 uncontested,
2756 });
2757 self.state.status = if met_quorum {
2758 RunStatus::Reviewing
2759 } else {
2760 RunStatus::Stalled
2761 };
2762 self.state.save()?;
2763 Ok(())
2764 }
2765
2766 #[allow(clippy::too_many_lines)]
2787 async fn recover_stall(&mut self) -> Result<bool> {
2788 let run_id = self.state.id.clone();
2793 let prompts = self.state.config.prompts.clone();
2794 let quota_seats: BTreeSet<&str> =
2799 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2800 let absent: Vec<String> = self
2801 .state
2802 .judgements
2803 .iter()
2804 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2805 .map(|j| j.seat.clone())
2806 .collect();
2807 if absent.is_empty() {
2808 return Ok(false);
2809 }
2810 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2811 if viable.len() <= 1 {
2812 return Ok(false);
2813 }
2814 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2815 let language = self.state.config.graph.language.clone();
2816 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2817 let sessions = self.state.config.graph.sessions;
2818 let artifacts = agent::artifacts_dir(&self.state.dir());
2819 let root = self.state.worktree_root();
2820 let base_short = short(&self.state.base_commit);
2821 let candidates: Vec<Candidate> = viable.clone();
2822
2823 let mut positions: Vec<usize> = absent
2825 .iter()
2826 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2827 .collect();
2828 if positions.is_empty() {
2829 return Ok(false);
2830 }
2831 positions.sort_unstable();
2832 positions.dedup();
2833
2834 let mut judge_jobs = Vec::new();
2836 for &j in &positions {
2837 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2838 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2839 let seat_key = format!("judge-{}", j + 1);
2840 let spec = self.roles.judges[j].clone();
2841 let seat = self.seat(&seat_key, &spec.id);
2842 judge_jobs.push(SeatJob {
2843 spec,
2844 seat,
2845 prompt: prompt::judge(
2846 &self.state.instruction,
2847 &views,
2848 self.roles.judges.len(),
2849 &base_short,
2850 &language,
2851 ),
2852 cwd: root.join(seat_key),
2853 timeout,
2854 allow_write: false,
2855 sessions,
2856 artifacts: artifacts.clone(),
2857 stem: format!("judge-{}-recover", j + 1),
2858 });
2859 }
2860
2861 let labels_for_check = labels.clone();
2862 let mut judge_losses = Vec::new();
2863 let retries = self.state.config.graph.retries;
2864 let cache = self.state.config.cache_dir();
2865 let ctx = WaveCtx {
2866 run: &run_id,
2867 node: "judge",
2868 prompts: &prompts,
2869 cache: cache.as_deref(),
2870 round: None,
2871 };
2872 let results = ask_json_wave::<Ranking>(
2873 judge_jobs,
2874 Arc::clone(&self.sem),
2875 retries,
2876 &ctx,
2877 &mut judge_losses,
2878 &mut self.state,
2879 &move |r: &Ranking| r.validate(&labels_for_check),
2880 )
2881 .await;
2882
2883 let mut recovered: BTreeSet<usize> = BTreeSet::new();
2885 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
2886 self.state.seats.insert(seat.key.clone(), seat);
2887 let record = &mut self.state.judgements[j];
2888 match res {
2889 Ok((ranking, out)) => {
2890 record.ranking = ranking.normalized();
2891 record.reasons = ranking.reasons;
2892 record.confidence = ranking.confidence;
2893 record.failed = None;
2894 record.duration_ms = out.duration_ms;
2895 recovered.insert(j);
2896 self.state.event(
2897 "recover",
2898 format!("judge {} ranked again after the limit", j + 1),
2899 );
2900 }
2901 Err(e) => {
2902 self.state
2903 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
2904 }
2905 }
2906 }
2907
2908 let mut vote_jobs = Vec::new();
2910 let mut vote_pos: Vec<usize> = Vec::new();
2911 for &j in &recovered {
2912 let seat_key = format!("judge-{}", j + 1);
2913 let spec = self.roles.judges[j].clone();
2914 let seat = self.seat(&seat_key, &spec.id);
2915 let mut text = prompt::final_vote(&labels, &language);
2916 if !has_context(&spec, &seat, sessions) {
2917 text = format!(
2918 "{}\n\n# Candidates\n\n{}",
2919 text,
2920 self.candidate_block(&candidates, &base_short)
2921 );
2922 }
2923 vote_jobs.push(SeatJob {
2924 spec,
2925 seat,
2926 prompt: text,
2927 cwd: root.join(seat_key),
2928 timeout,
2929 allow_write: false,
2930 sessions,
2931 artifacts: artifacts.clone(),
2932 stem: format!("vote-judge-{}-recover", j + 1),
2933 });
2934 vote_pos.push(j);
2935 }
2936 let allowed = labels.clone();
2937 let mut vote_losses = Vec::new();
2938 let vote_retries = self.state.config.graph.retries;
2939 let vote_cache = self.state.config.cache_dir();
2940 let ctx = WaveCtx {
2941 run: &run_id,
2942 node: "vote",
2943 prompts: &prompts,
2944 cache: vote_cache.as_deref(),
2945 round: None,
2946 };
2947 let votes = ask_json_wave::<FinalVote>(
2948 vote_jobs,
2949 Arc::clone(&self.sem),
2950 vote_retries,
2951 &ctx,
2952 &mut vote_losses,
2953 &mut self.state,
2954 &move |v: &FinalVote| match v.label() {
2955 Some(c) if allowed.contains(&c) => Ok(()),
2956 other => bail!("vote {other:?} is not one of {allowed:?}"),
2957 },
2958 )
2959 .await;
2960 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
2961 let agent_id = seat.agent.clone();
2962 self.state.seats.insert(seat.key.clone(), seat);
2963 match res {
2964 Ok((v, _)) => {
2965 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
2966 rec.vote = v.label();
2967 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2968 } else {
2969 self.state.votes.push(VoteRecord {
2970 judge: j + 1,
2971 agent: agent_id,
2972 vote: v.label(),
2973 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
2974 changed: false,
2975 });
2976 }
2977 self.state.event(
2978 "recover",
2979 format!("judge {} voted again after the limit", j + 1),
2980 );
2981 }
2982 Err(e) => {
2983 self.state
2984 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
2985 }
2986 }
2987 }
2988
2989 let recovered_keys: BTreeSet<String> = recovered
2993 .iter()
2994 .map(|&j| format!("judge-{}", j + 1))
2995 .collect();
2996 self.state
2997 .quota
2998 .retain(|q| !recovered_keys.contains(&q.seat));
2999 for loss in judge_losses.into_iter().chain(vote_losses) {
3003 if recovered_keys.contains(&loss.seat) {
3004 continue;
3005 }
3006 self.state.quota.retain(|q| q.seat != loss.seat);
3007 self.state.quota.push(loss);
3008 }
3009
3010 self.state.tally = None;
3012 self.tally()?;
3013 Ok(self
3014 .state
3015 .tally
3016 .as_ref()
3017 .map(|t| t.met_quorum)
3018 .unwrap_or(false))
3019 }
3020
3021 async fn fold_losers(&mut self) -> Result<()> {
3024 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
3025 return Ok(());
3026 };
3027 let repo = self.state.repo.clone();
3028 let mut folded = Vec::new();
3029 for i in 0..self.state.candidates.len() {
3030 let c = &self.state.candidates[i];
3031 if c.label == winner || c.folded {
3032 continue;
3033 }
3034 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
3035 git::worktree_remove(&repo, &wt).await.ok();
3036 git::branch_delete(&repo, &branch).await.ok();
3037 self.state.candidates[i].folded = true;
3038 folded.push(label.to_string());
3039 }
3040 let root = self.state.worktree_root();
3042 for j in 1..=self.roles.judges.len() {
3043 let wt = root.join(format!("judge-{j}"));
3044 if wt.exists() {
3045 git::worktree_remove(&repo, &wt).await.ok();
3046 }
3047 }
3048 if self.state.config.graph.advise {
3051 for k in 1..=self.state.config.graph.advisors {
3052 let wt = root.join(format!("advisor-{k}"));
3053 if wt.exists() {
3054 git::worktree_remove(&repo, &wt).await.ok();
3055 }
3056 }
3057 }
3058 if !folded.is_empty() {
3059 self.state
3060 .event("fold", format!("folded candidates {}", folded.join(", ")));
3061 self.state.save()?;
3062 }
3063 Ok(())
3064 }
3065
3066 async fn sync_to_base(&mut self) -> Result<()> {
3096 if self
3097 .state
3098 .base_sync
3099 .as_ref()
3100 .is_some_and(|s| s.conflict.is_some())
3101 {
3102 return Ok(());
3103 }
3104 let Some(winner) = self.state.winner().cloned() else {
3105 return Ok(());
3106 };
3107
3108 let repo = self.state.repo.clone();
3109 let remote = self.state.config.merge.remote.clone();
3110 let base_branch = self.state.base_branch.clone();
3111 let tracking = format!("{remote}/{base_branch}");
3112
3113 git::fetch(&repo, &remote, &base_branch).await.ok();
3114 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3118 return Ok(());
3119 };
3120
3121 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3122 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3123 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3124
3125 if behind == 0 {
3126 self.state.base_sync = Some(BaseSync {
3127 tip,
3128 behind: 0,
3129 attempts,
3130 conflict: None,
3131 });
3132 self.state.save()?;
3133 return Ok(());
3134 }
3135
3136 if attempts >= BASE_SYNC_ROUNDS {
3137 let why = format!(
3138 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3139 rebase(s); rebasing again would only race it",
3140 winner.branch
3141 );
3142 self.state.status = RunStatus::Blocked;
3143 self.state.base_sync = Some(BaseSync {
3144 tip,
3145 behind,
3146 attempts,
3147 conflict: Some(why.clone()),
3148 });
3149 self.state.event("land", why);
3150 self.state.save()?;
3151 return Ok(());
3152 }
3153
3154 self.state.event(
3155 "land",
3156 format!(
3157 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3158 winner.branch
3159 ),
3160 );
3161 self.state.save()?;
3162
3163 let scratch = self.state.dir().join("base-sync");
3164 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
3165 let attempts = attempts + 1;
3166 match rebased {
3167 Ok(None) => {
3168 git::sync_to_head(&winner.worktree).await?;
3172 self.state.base_sync = Some(BaseSync {
3173 tip: tip.clone(),
3174 behind: 0,
3175 attempts,
3176 conflict: None,
3177 });
3178 self.state
3179 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3180 }
3181 Ok(Some(conflict)) => {
3182 let why = format!(
3183 "{} conflicts with {tracking} and did not rebase: {}",
3184 winner.branch,
3185 conflict.chars().take(600).collect::<String>()
3186 );
3187 self.state.status = RunStatus::Blocked;
3188 self.state.base_sync = Some(BaseSync {
3189 tip,
3190 behind,
3191 attempts,
3192 conflict: Some(why.clone()),
3193 });
3194 self.state.event("land", why);
3195 }
3196 Err(e) => {
3197 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3198 self.state.status = RunStatus::Blocked;
3199 self.state.base_sync = Some(BaseSync {
3200 tip,
3201 behind,
3202 attempts,
3203 conflict: Some(why.clone()),
3204 });
3205 self.state.event("land", why);
3206 }
3207 }
3208 self.state.save()?;
3209 Ok(())
3210 }
3211
3212 fn landing_base(&self) -> String {
3222 self.state
3223 .base_sync
3224 .as_ref()
3225 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3226 }
3227
3228 pub async fn fix_selected(
3261 &mut self,
3262 ids: &[String],
3263 reason: &str,
3264 allow_stale: bool,
3265 ) -> Result<()> {
3266 let reason = reason.trim();
3267 if reason.is_empty() {
3268 bail!("a fix request needs a reason — that is the operator's own record of why");
3269 }
3270 if ids.is_empty() {
3271 bail!("no finding id given");
3272 }
3273 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3274 bail!(
3275 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3276 has already concluded — can be given a targeted fix. A run still \
3277 in progress should simply be resumed; a `merged` run's branch has \
3278 already landed, so its answer is a fresh `magi review <branch>`, \
3279 not reopening this run's own record",
3280 self.state.id,
3281 self.state.status.as_str()
3282 );
3283 }
3284 let Some(winner) = self.state.winner().cloned() else {
3285 bail!("run {} has no winning candidate to fix", self.state.id);
3286 };
3287 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3288 bail!(
3289 "branch `{}` no longer exists; this run cannot be extended",
3290 winner.branch
3291 );
3292 }
3293 let home = crate::run::home();
3294 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3295 bail!(
3296 "run {} is currently being worked on by another magi process",
3297 self.state.id
3298 );
3299 }
3300 let _claim = FixClaim::acquire(&self.state.dir())?;
3306
3307 let mut seen = BTreeSet::new();
3311 let mut findings = Vec::new();
3312 let mut missing = Vec::new();
3313 for id in ids {
3314 if !seen.insert(id.clone()) {
3315 continue;
3316 }
3317 match self.state.finding(id) {
3318 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3319 id: f.id.clone(),
3320 severity: f.severity,
3321 reviewer_vote: rec.vote,
3322 round: round.round,
3323 round_head: round.head.clone(),
3324 reviewer: rec.reviewer,
3325 agent: rec.agent.clone(),
3326 file: f.file.clone(),
3327 line: f.line,
3328 title: f.title.clone(),
3329 detail: f.detail.clone(),
3330 outcome: OperatorFixOutcome::Pending,
3331 }),
3332 None => missing.push(id.clone()),
3333 }
3334 }
3335 if !missing.is_empty() {
3336 bail!(
3337 "unknown finding id(s): {}; nothing was changed",
3338 missing.join(", ")
3339 );
3340 }
3341
3342 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3343 let stale_details: Vec<(String, String)> = findings
3344 .iter()
3345 .filter(|f| f.round_head != head_at_request)
3346 .map(|f| (f.id.clone(), f.round_head.clone()))
3347 .collect();
3348 let stale = !stale_details.is_empty();
3349 if stale && !allow_stale {
3350 bail!(
3351 "the branch has moved since some finding(s) were raised — {} — now \
3352 at {}; pass --allow-stale to fix anyway, or re-run review first",
3353 stale_details
3354 .iter()
3355 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3356 .collect::<Vec<_>>()
3357 .join(", "),
3358 short(&head_at_request)
3359 );
3360 }
3361
3362 let request = OperatorFixRequest {
3363 requested_at: Timestamp::now(),
3364 reason: reason.to_owned(),
3365 findings,
3366 head_at_request: head_at_request.clone(),
3367 allow_stale,
3368 stale,
3369 fix: None,
3370 result_head: None,
3371 follow_up_review_run: None,
3372 };
3373 self.state.event(
3374 "fix",
3375 format!(
3376 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3377 request.findings.len(),
3378 request
3379 .findings
3380 .iter()
3381 .map(|f| f.id.as_str())
3382 .collect::<Vec<_>>()
3383 .join(", "),
3384 ),
3385 );
3386 self.state.operator_fixes.push(request);
3393 self.state.save()?;
3394 let request_index = self.state.operator_fixes.len() - 1;
3395
3396 if winner.worktree.exists() {
3405 let dirty = git::git(
3408 &winner.worktree,
3409 &["status", "--porcelain", "--untracked-files=all"],
3410 )
3411 .await?;
3412 let only_withheld = dirty.lines().all(|l| {
3413 l.strip_prefix("?? ")
3414 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3415 });
3416 if !only_withheld {
3417 bail!(
3418 "`{}` has uncommitted changes; refusing to touch it — commit or \
3419 discard them first",
3420 winner.worktree.display()
3421 );
3422 }
3423 git::worktree_remove(&self.state.repo, &winner.worktree)
3424 .await
3425 .ok();
3426 }
3427 let fix_worktree = self.state.worktree_root().join("operator-fix");
3428 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3429 git::git(
3430 &self.state.repo,
3431 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3432 )
3433 .await
3434 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3435 if !git::is_clean(&fix_worktree).await? {
3436 git::worktree_remove(&self.state.repo, &fix_worktree)
3437 .await
3438 .ok();
3439 bail!(
3440 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3441 winner.branch
3442 );
3443 }
3444
3445 let run_id = self.state.id.clone();
3446 let prompts = self.state.config.prompts.clone();
3447 let language = self.state.config.graph.language.clone();
3448 let sessions = self.state.config.graph.sessions;
3449 let artifacts = agent::artifacts_dir(&self.state.dir());
3450 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3451 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3452 _ => (
3453 self.state
3454 .config
3455 .agent(&winner.agent)
3456 .cloned()
3457 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3458 format!("impl-{}", winner.label),
3459 ),
3460 };
3461 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3462 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3463 .findings
3464 .iter()
3465 .map(|f| Finding {
3466 id: f.id.clone(),
3467 severity: f.severity,
3468 file: f.file.clone(),
3469 line: f.line,
3470 title: f.title.clone(),
3471 detail: f.detail.clone(),
3472 })
3473 .collect();
3474 let job = SeatJob {
3475 prompt: prompt::operator_fix(
3476 &self.state.instruction,
3477 &finding_list,
3478 reason,
3479 &stale_details,
3480 &head_at_request,
3481 &language,
3482 ),
3483 spec: fix_spec.clone(),
3484 seat,
3485 cwd: fix_worktree.clone(),
3486 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3487 allow_write: true,
3488 sessions,
3489 artifacts: artifacts.clone(),
3490 stem: "operator-fix".to_owned(),
3491 };
3492 let cache = self.state.config.cache_dir();
3493 let ctx = WaveCtx {
3494 run: &run_id,
3495 node: "fix",
3496 prompts: &prompts,
3497 cache: cache.as_deref(),
3498 round: None,
3499 };
3500 let (seat, out) =
3501 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3502 let agent_id = seat.agent.clone();
3503
3504 let mut fix = FixRecord {
3505 agent: agent_id,
3506 addressed: Vec::new(),
3507 rejected: Vec::new(),
3508 notes: String::new(),
3509 committed: false,
3510 failed: None,
3511 duration_ms: 0,
3512 continuation: None,
3513 };
3514 let mut final_seat = seat.clone();
3515 match out {
3516 AgentOutcome::Ok(o) => {
3517 fix.duration_ms = o.duration_ms;
3518 let parsed = verdict::extract_json::<FixReport>(&o.text);
3519 let incomplete_reason = match &parsed {
3520 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3521 "the reply parsed, but it reported a command whose own CLI \
3522 never confirmed an exit status"
3523 .to_owned(),
3524 ),
3525 Ok(_) => None,
3526 Err(e) => Some(e.to_string()),
3527 };
3528 match incomplete_reason {
3529 None => {
3530 let report = parsed.expect("checked Ok above");
3531 fix.addressed = report.addressed;
3532 fix.rejected = report.rejected;
3533 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3534 }
3535 Some(reason) => {
3536 let (resumed_seat, resolved, failure, cont) = self
3537 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3538 .await;
3539 fix.duration_ms += cont.cumulative_wait_ms;
3540 fix.continuation = Some(cont);
3541 final_seat = resumed_seat;
3542 match resolved {
3543 Some(report) => {
3544 fix.addressed = report.addressed;
3545 fix.rejected = report.rejected;
3546 fix.notes =
3547 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3548 }
3549 None => fix.failed = failure,
3550 }
3551 }
3552 }
3553 }
3554 AgentOutcome::Dropped(o) => {
3555 fix.duration_ms = o.duration_ms;
3556 let why = o
3557 .dropped
3558 .as_ref()
3559 .map(|d| d.why.as_str())
3560 .unwrap_or("the CLI ended the stream without delivering its answer");
3561 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3562 }
3563 AgentOutcome::Quota(o) => {
3564 self.state.quota.push(QuotaLoss {
3565 seat: final_seat.key.clone(),
3566 node: "fix".to_owned(),
3567 at: Timestamp::now(),
3568 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3569 });
3570 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3571 }
3572 AgentOutcome::Failed(e) => fix.failed = Some(e),
3573 }
3574 if fix.continuation.is_none() {
3575 fix.continuation = Some(ContinuationRecord::not_needed());
3576 }
3577 self.state.seats.insert(final_seat.key.clone(), final_seat);
3578
3579 let rescue_message = format!(
3580 "magi: operator-selected fix ({}) (uncommitted work)",
3581 self.state.operator_fixes[request_index]
3582 .findings
3583 .iter()
3584 .map(|f| f.id.as_str())
3585 .collect::<Vec<_>>()
3586 .join(", ")
3587 );
3588 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3589 self.state.note_withheld("fix", &r.withheld);
3590 }
3591 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3592 fix.committed = after != head_at_request;
3593 git::worktree_remove(&self.state.repo, &fix_worktree)
3594 .await
3595 .ok();
3596
3597 self.state.event(
3598 "fix",
3599 match &fix.failed {
3600 Some(reason) => format!(
3601 "operator fix: adoption report was lost ({reason}); {}",
3602 if fix.committed {
3603 "committed"
3604 } else {
3605 "NO new commit"
3606 }
3607 ),
3608 None => format!(
3609 "operator fix: {} addressed, {} rejected, {}",
3610 fix.addressed.len(),
3611 fix.rejected.len(),
3612 if fix.committed {
3613 "committed"
3614 } else {
3615 "NO new commit"
3616 }
3617 ),
3618 },
3619 );
3620
3621 for f in &mut self.state.operator_fixes[request_index].findings {
3628 f.outcome = if fix.failed.is_some() {
3629 OperatorFixOutcome::Unreported
3630 } else if fix.addressed.contains(&f.id) {
3631 OperatorFixOutcome::Addressed
3632 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3633 OperatorFixOutcome::Rejected { why: r.why.clone() }
3634 } else {
3635 OperatorFixOutcome::Unreported
3636 };
3637 }
3638
3639 let committed = fix.committed;
3640 if committed {
3641 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3642 }
3643 self.state.operator_fixes[request_index].fix = Some(fix);
3644 self.state.save()?;
3647
3648 if committed {
3649 self.state.event(
3650 "fix",
3651 format!(
3652 "operator fix committed {}; opening a follow-up review-only run",
3653 short(&after)
3654 ),
3655 );
3656 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3657 Ok(mut follow_up) => {
3658 follow_up.state.event(
3659 "start",
3660 format!(
3661 "requested by an operator fix on run {} for finding(s) {}",
3662 self.state.id,
3663 self.state.operator_fixes[request_index]
3664 .findings
3665 .iter()
3666 .map(|f| f.id.as_str())
3667 .collect::<Vec<_>>()
3668 .join(", "),
3669 ),
3670 );
3671 follow_up.state.save()?;
3672 let follow_up_id = follow_up.state.id.clone();
3673 if let Err(e) = follow_up.execute().await {
3674 self.state.event(
3675 "fix",
3676 format!(
3677 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3678 ),
3679 );
3680 }
3681 self.state.operator_fixes[request_index].follow_up_review_run =
3682 Some(follow_up_id);
3683 }
3684 Err(e) => {
3685 self.state.event(
3686 "fix",
3687 format!("committed the fix but could not open a follow-up review: {e:#}"),
3688 );
3689 }
3690 }
3691 self.state.save()?;
3692 }
3693
3694 Ok(())
3695 }
3696
3697 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3704 match &self.roles.fixer {
3705 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3706 _ => (
3707 self.state
3708 .config
3709 .agent(&winner.agent)
3710 .cloned()
3711 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3712 format!("impl-{}", winner.label),
3713 ),
3714 }
3715 }
3716
3717 async fn review_loop(&mut self) -> Result<()> {
3718 if self
3723 .state
3724 .base_sync
3725 .as_ref()
3726 .is_some_and(|s| s.conflict.is_some())
3727 {
3728 return Ok(());
3729 }
3730 let run_id = self.state.id.clone();
3735 let prompts = self.state.config.prompts.clone();
3736 let Some(winner) = self.state.winner().cloned() else {
3737 return Ok(());
3738 };
3739 let max_rounds = self.state.config.graph.review_rounds;
3740 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
3750 self.state.status = status;
3751 self.state.save()?;
3752 return Ok(());
3753 }
3754 self.state.status = RunStatus::Reviewing;
3755 if self
3765 .state
3766 .reviews
3767 .last()
3768 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
3769 {
3770 let shell = self.state.config.shell();
3771 return self
3772 .stop_reviewing(
3773 "the last round's own verification never resolved",
3774 &shell,
3775 &winner.worktree,
3776 )
3777 .await;
3778 }
3779
3780 let repo = self.state.repo.clone();
3781 let root = self.state.worktree_root();
3782 let language = self.state.config.graph.language.clone();
3783 let sessions = self.state.config.graph.sessions;
3784 let artifacts = agent::artifacts_dir(&self.state.dir());
3785 let base = self.landing_base();
3786 let base_short = short(&base);
3787 let reviewers = self.roles.reviewers.clone();
3788 let shell = self.state.config.shell();
3789
3790 for round in (self.state.reviews.len() + 1)..=max_rounds {
3791 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3792 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
3793 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
3794 let prev_verification = self
3803 .state
3804 .reviews
3805 .last()
3806 .and_then(|r| r.verification_summary(&head));
3807
3808 let mut jobs = Vec::new();
3812 for (r, spec) in reviewers.iter().cloned().enumerate() {
3813 let wt = root.join(format!("review-{}", r + 1));
3814 if wt.exists() {
3815 git::reset_detached(&wt, &head).await?;
3816 } else {
3817 git::worktree_add_detached(&repo, &wt, &head).await?;
3818 }
3819 let seat_key = format!("review-{}", r + 1);
3820 let seat = self.seat(&seat_key, &spec.id);
3821 jobs.push(SeatJob {
3822 prompt: prompt::review(&prompt::ReviewCtx {
3823 instruction: &self.state.instruction,
3824 branch: &winner.branch,
3825 base_short: &base_short,
3826 stat: &stat,
3827 patch: &patch,
3828 verification: prev_verification.as_ref(),
3829 reviewers: reviewers.len(),
3830 round,
3831 rounds: max_rounds,
3832 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
3835 lens: Lens::for_seat(r),
3836 language: &language,
3837 }),
3838 spec,
3839 seat,
3840 cwd: wt,
3841 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
3842 allow_write: false,
3843 sessions,
3844 artifacts: artifacts.clone(),
3845 stem: format!("review-{round}-{}", r + 1),
3846 });
3847 }
3848
3849 self.state.event(
3850 "review",
3851 format!(
3852 "round {round}: {} reviewers on {}",
3853 jobs.len(),
3854 short(&head)
3855 ),
3856 );
3857 let mut quota_losses = Vec::new();
3858 let review_retries = self.state.config.graph.retries;
3859 let review_cache = self.state.config.cache_dir();
3860 let ctx = WaveCtx {
3861 run: &run_id,
3862 node: "review",
3863 prompts: &prompts,
3864 cache: review_cache.as_deref(),
3865 round: Some(round),
3866 };
3867 let results = ask_json_wave::<Review>(
3868 jobs,
3869 Arc::clone(&self.sem),
3870 review_retries,
3871 &ctx,
3872 &mut quota_losses,
3873 &mut self.state,
3874 &|_: &Review| Ok(()),
3875 )
3876 .await;
3877 let round_quota_missing = quota_losses.len();
3881 self.state.quota.extend(quota_losses);
3882
3883 let mut records = Vec::new();
3884 let mut all_findings = Vec::new();
3885 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
3886 let agent_id = seat.agent.clone();
3887 self.state.seats.insert(seat.key.clone(), seat);
3888 let mut record = ReviewRecord {
3889 reviewer: r + 1,
3890 agent: agent_id,
3891 summary: String::new(),
3892 findings: Vec::new(),
3893 vote: None,
3894 failed: None,
3895 duration_ms: 0,
3896 attempts,
3902 };
3903 match res {
3904 Ok((review, out)) => {
3905 record.summary =
3913 blind::sanitize_prose(&review.summary, &self.state.config.blind);
3914 record.vote = Some(review.vote);
3915 record.duration_ms = out.duration_ms;
3916 for (n, mut f) in review.findings.into_iter().enumerate() {
3917 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
3920 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
3921 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
3922 f.file = f
3928 .file
3929 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
3930 all_findings.push(f.clone());
3931 record.findings.push(f);
3932 }
3933 self.state.event(
3934 "review",
3935 format!(
3936 "round {round}: reviewer {} voted {} with {} finding(s)",
3937 r + 1,
3938 review.vote.label(),
3939 record.findings.len()
3940 ),
3941 );
3942 }
3943 Err(e) => {
3944 record.failed = Some(e.to_string());
3945 self.state.event(
3946 "review",
3947 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
3948 );
3949 }
3950 }
3951 records.push(record);
3952 }
3953
3954 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
3961 let vote_split =
3962 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
3963 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
3964 if vote_split {
3965 self.state.event(
3966 "review",
3967 format!(
3968 "round {round}: votes split ({}) — one round of reconsideration",
3969 initial_votes
3970 .iter()
3971 .map(|v| v.label())
3972 .collect::<Vec<_>>()
3973 .join(", ")
3974 ),
3975 );
3976 let panel: Vec<ReviewSeatReport<'_>> = records
3979 .iter()
3980 .filter_map(|r| {
3981 r.vote.map(|vote| ReviewSeatReport {
3982 reviewer: r.reviewer,
3983 vote,
3984 summary: &r.summary,
3985 findings: &r.findings,
3986 })
3987 })
3988 .collect();
3989
3990 let mut jobs = Vec::new();
3991 let mut seats_at = Vec::new();
3992 for (r, spec) in reviewers.iter().cloned().enumerate() {
3993 if records[r].vote.is_none() {
3997 continue;
3998 }
3999 let wt = root.join(format!("review-{}", r + 1));
4000 let seat_key = format!("review-{}", r + 1);
4001 let seat = self.seat(&seat_key, &spec.id);
4002 let patch_ctx = if has_context(&spec, &seat, sessions) {
4007 None
4008 } else {
4009 Some(ReviewPatch {
4010 branch: &winner.branch,
4011 base_short: &base_short,
4012 stat: &stat,
4013 patch: &patch,
4014 })
4015 };
4016 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
4017 instruction: &self.state.instruction,
4018 reviewer: r + 1,
4019 lens: Lens::for_seat(r),
4020 panel: &panel,
4021 patch: patch_ctx,
4022 round,
4023 rounds: max_rounds,
4024 language: &language,
4025 });
4026 jobs.push(SeatJob {
4027 prompt,
4028 spec,
4029 seat,
4030 cwd: wt,
4031 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4032 allow_write: false,
4033 sessions,
4034 artifacts: artifacts.clone(),
4035 stem: format!("review-{round}-reconsider-{}", r + 1),
4036 });
4037 seats_at.push(r);
4038 }
4039
4040 let mut recon_quota_losses = Vec::new();
4041 let recon_cache = self.state.config.cache_dir();
4042 let recon_ctx = WaveCtx {
4043 run: &run_id,
4044 node: "review",
4045 prompts: &prompts,
4046 cache: recon_cache.as_deref(),
4047 round: Some(round),
4048 };
4049 let recon_results = ask_json_wave::<ReviewRevote>(
4050 jobs,
4051 Arc::clone(&self.sem),
4052 review_retries,
4053 &recon_ctx,
4054 &mut recon_quota_losses,
4055 &mut self.state,
4056 &|_: &ReviewRevote| Ok(()),
4057 )
4058 .await;
4059 self.state.quota.extend(recon_quota_losses);
4060
4061 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
4062 let agent_id = seat.agent.clone();
4063 self.state.seats.insert(seat.key.clone(), seat);
4064 let mut rec = ReviewRevoteRecord {
4065 reviewer: r + 1,
4066 agent: agent_id,
4067 vote: None,
4068 reason: String::new(),
4069 failed: None,
4070 };
4071 match res {
4072 Ok((rv, _)) => {
4073 rec.vote = Some(rv.vote);
4074 rec.reason =
4075 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
4076 self.state.event(
4077 "review",
4078 format!(
4079 "round {round}: reviewer {} revoted {}",
4080 r + 1,
4081 rv.vote.label()
4082 ),
4083 );
4084 }
4085 Err(e) => {
4086 rec.failed = Some(e.to_string());
4087 self.state.event(
4088 "review",
4089 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4090 );
4091 }
4092 }
4093 reconsideration.push(rec);
4094 }
4095 } else if initial_votes.len() > 1 {
4096 self.state.event(
4097 "review",
4098 format!(
4099 "round {round}: votes agreed ({}) — no reconsideration",
4100 initial_votes[0].label()
4101 ),
4102 );
4103 }
4104
4105 let final_votes: Vec<ReviewVote> = records
4109 .iter()
4110 .filter_map(|r| {
4111 reconsideration
4112 .iter()
4113 .find(|rv| rv.reviewer == r.reviewer)
4114 .and_then(|rv| rv.vote)
4115 .or(r.vote)
4116 })
4117 .collect();
4118 let round_verdict = ReviewVote::worst(final_votes);
4119
4120 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4121 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4122 let defer_e2e =
4133 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4134 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4135 let reason =
4136 format!("{blocking} blocking finding(s) already required a fix this round");
4137 self.state.event(
4138 "verify",
4139 format!(
4140 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4141 {}); it will run once a round has none left",
4142 short(&head)
4143 ),
4144 );
4145 (Vec::new(), false, true, Some(reason))
4146 } else {
4147 let e2e_commands = self.state.config.verify.e2e.clone();
4148 let cache_dir = self.state.config.cache_dir();
4149 let context = format!("round {round}");
4150 let (e2e, verify_retried) = with_cache_lease(
4151 &mut self.state,
4152 cache_dir.as_deref(),
4153 "e2e",
4154 "e2e",
4155 &winner.worktree,
4156 &head,
4157 verify_timeout,
4158 &context,
4159 |state, budget| {
4160 let shell = shell.clone();
4161 let e2e_commands = e2e_commands.clone();
4162 let worktree = winner.worktree.clone();
4163 let context = context.clone();
4164 async move {
4165 run_e2e_with_retry(
4166 state,
4167 &shell,
4168 &e2e_commands,
4169 &worktree,
4170 budget,
4171 &context,
4172 )
4173 .await
4174 }
4175 },
4176 )
4177 .await;
4178 (e2e, verify_retried, false, None)
4179 };
4180
4181 let expected = records.len();
4182 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4183 let incomplete = answered < expected;
4184 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4185 let policy = self.state.config.graph.incomplete_review;
4186 let clean = round_is_clean(
4187 blocking,
4188 e2e_ok,
4189 answered,
4190 expected,
4191 round_quota_missing,
4192 policy,
4193 );
4194
4195 let mut round_record = ReviewRound {
4196 round,
4197 head: head.clone(),
4198 verified_head: None,
4199 verified_at: None,
4200 reviews: records,
4201 e2e,
4202 verify_retried,
4203 e2e_deferred,
4204 e2e_defer_reason,
4205 fix: None,
4206 blocking,
4207 answered,
4208 expected,
4209 clean,
4210 progressed: false,
4211 vote_split,
4212 reconsideration,
4213 verdict: round_verdict,
4214 };
4215 if !matches!(
4226 round_record.e2e_status(),
4227 E2eStatus::Deferred | E2eStatus::NotConfigured
4228 ) {
4229 round_record.verified_head = Some(head.clone());
4230 round_record.verified_at = Some(Timestamp::now());
4231 }
4232 let this_round_verification = round_record.verification_summary(&head);
4233
4234 if incomplete {
4235 let missing: Vec<String> = round_record
4236 .reviews
4237 .iter()
4238 .filter(|r| r.failed.is_some())
4239 .map(|r| format!("review-{}", r.reviewer))
4240 .collect();
4241 self.state.event(
4242 "review",
4243 format!(
4244 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4245 missing.join(", ")
4246 ),
4247 );
4248 }
4249
4250 if clean {
4251 self.state.event(
4252 "review",
4253 if incomplete && policy == IncompleteReviewPolicy::Warn {
4254 format!(
4255 "round {round}: clean (warn policy, incomplete panel) — no \
4256 blocking findings from the seats that answered, verification green"
4257 )
4258 } else if incomplete {
4259 format!(
4260 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4261 quorum) — no blocking findings from the seats that answered, \
4262 verification green",
4263 expected - answered
4264 )
4265 } else {
4266 format!("round {round}: clean — no blocking findings, verification green")
4267 },
4268 );
4269 self.state.reviews.push(round_record);
4270 self.state.status = RunStatus::Gating;
4271 self.state.save()?;
4272 return Ok(());
4273 }
4274
4275 if incomplete && blocking == 0 && e2e_ok {
4283 self.state.reviews.push(round_record);
4284 self.state.save()?;
4285 if round == max_rounds {
4286 self.state.status = RunStatus::Blocked;
4287 self.state.event(
4288 "review",
4289 format!(
4290 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4291 refusing to call it clean",
4292 expected - answered
4293 ),
4294 );
4295 return Ok(());
4296 }
4297 continue;
4298 }
4299
4300 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4312 self.state.reviews.push(round_record);
4313 return self
4314 .stop_reviewing(
4315 "the round's own verification could not run",
4316 &shell,
4317 &winner.worktree,
4318 )
4319 .await;
4320 }
4321
4322 if round == max_rounds {
4323 self.state.reviews.push(round_record);
4324 return self
4325 .stop_reviewing(
4326 &format!(
4327 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4328 ),
4329 &shell,
4330 &winner.worktree,
4331 )
4332 .await;
4333 }
4334
4335 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4338 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4339 let blocking_findings: Vec<_> = all_findings
4340 .iter()
4341 .filter(|f| f.severity.blocks())
4342 .cloned()
4343 .collect();
4344 let job = SeatJob {
4345 prompt: prompt::fix(
4346 &self.state.instruction,
4347 &blocking_findings,
4348 this_round_verification.as_ref(),
4349 round,
4350 max_rounds,
4351 &language,
4352 ),
4353 spec: fix_spec.clone(),
4354 seat,
4355 cwd: winner.worktree.clone(),
4356 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4357 allow_write: true,
4358 sessions,
4359 artifacts: artifacts.clone(),
4360 stem: format!("fix-{round}"),
4361 };
4362 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4363 let cache = self.state.config.cache_dir();
4364 let ctx = WaveCtx {
4365 run: &run_id,
4366 node: "fix",
4367 prompts: &prompts,
4368 cache: cache.as_deref(),
4369 round: Some(round),
4370 };
4371 let (seat, out) =
4372 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4373 let agent_id = seat.agent.clone();
4374
4375 let mut fix = FixRecord {
4376 agent: agent_id,
4377 addressed: Vec::new(),
4378 rejected: Vec::new(),
4379 notes: String::new(),
4380 committed: false,
4381 failed: None,
4382 duration_ms: 0,
4383 continuation: None,
4384 };
4385 let mut continuation = ContinuationRecord::not_needed();
4386 let mut final_seat = seat.clone();
4387 match out {
4388 AgentOutcome::Ok(o) => {
4389 fix.duration_ms = o.duration_ms;
4390 let parsed = verdict::extract_json::<FixReport>(&o.text);
4391 let incomplete_reason = match &parsed {
4398 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4399 "the reply parsed, but it reported a command whose own CLI never \
4400 confirmed an exit status"
4401 .to_owned(),
4402 ),
4403 Ok(_) => None,
4404 Err(e) => Some(e.to_string()),
4405 };
4406 match incomplete_reason {
4407 None => {
4408 let report = parsed.expect("checked Ok above");
4409 fix.addressed = report.addressed;
4410 fix.rejected = report.rejected;
4411 fix.notes =
4412 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4413 }
4414 Some(reason) => {
4415 let (resumed_seat, resolved, failure, cont) = self
4416 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4417 .await;
4418 fix.duration_ms += cont.cumulative_wait_ms;
4419 continuation = cont;
4420 final_seat = resumed_seat;
4421 match resolved {
4422 Some(report) => {
4423 fix.addressed = report.addressed;
4424 fix.rejected = report.rejected;
4425 fix.notes = blind::sanitize_prose(
4426 &report.notes,
4427 &self.state.config.blind,
4428 );
4429 }
4430 None => fix.failed = failure,
4431 }
4432 }
4433 }
4434 }
4435 AgentOutcome::Dropped(o) => {
4437 fix.duration_ms = o.duration_ms;
4438 let why = o
4439 .dropped
4440 .as_ref()
4441 .map(|d| d.why.as_str())
4442 .unwrap_or("the CLI ended the stream without delivering its answer");
4443 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4444 }
4445 AgentOutcome::Quota(o) => {
4446 self.state.quota.push(QuotaLoss {
4447 seat: final_seat.key.clone(),
4448 node: "fix".to_owned(),
4449 at: Timestamp::now(),
4450 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4451 });
4452 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4453 }
4454 AgentOutcome::Failed(e) => fix.failed = Some(e),
4455 }
4456 fix.continuation = Some(continuation);
4457 self.state.seats.insert(final_seat.key.clone(), final_seat);
4458 if let Ok(r) = git::rescue_commit(
4459 &winner.worktree,
4460 &format!("magi: review round {round} fixes (uncommitted work)"),
4461 )
4462 .await
4463 {
4464 self.state.note_withheld("fix", &r.withheld);
4465 }
4466 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4467 fix.committed = after != before;
4468 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4476 let progressed = diff_after != patch;
4477 let commit_note = if fix.committed {
4478 "committed"
4479 } else {
4480 "NO new commit"
4481 };
4482 let tree_note = if progressed {
4483 "changed vs base"
4484 } else {
4485 "unchanged vs base"
4486 };
4487 self.state.event(
4488 "fix",
4489 match &fix.failed {
4490 Some(reason) => {
4496 format!(
4497 "round {round}: fixer's adoption report was lost ({reason}); \
4498 {commit_note}, tree {tree_note}"
4499 )
4500 }
4501 None => format!(
4502 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4503 {tree_note}{}",
4504 fix.addressed.len(),
4505 fix.rejected.len(),
4506 if continuation.outcome == ContinuationOutcome::Resumed {
4507 format!(
4508 " (adoption report recovered after {} continuation(s))",
4509 continuation.attempts
4510 )
4511 } else {
4512 String::new()
4513 },
4514 ),
4515 },
4516 );
4517 round_record.fix = Some(fix);
4518 round_record.progressed = progressed;
4519 self.state.reviews.push(round_record);
4520 self.state.save()?;
4521
4522 if matches!(
4535 continuation.outcome,
4536 ContinuationOutcome::Exhausted
4537 | ContinuationOutcome::QuotaLost
4538 | ContinuationOutcome::NoSession
4539 ) {
4540 return self
4541 .stop_reviewing(
4542 "the fixer's adoption report never came back, even after resuming its \
4543 own seat; refusing to start another round against the same worktree \
4544 while that is unresolved",
4545 &shell,
4546 &winner.worktree,
4547 )
4548 .await;
4549 }
4550
4551 let streak = self
4552 .state
4553 .reviews
4554 .iter()
4555 .rev()
4556 .take_while(|r| !r.progressed)
4557 .count();
4558 if streak >= STAGNANT_LIMIT {
4559 return self
4560 .stop_reviewing(
4561 &format!(
4562 "the tree has not moved against base for {streak} round(s) in a row"
4563 ),
4564 &shell,
4565 &winner.worktree,
4566 )
4567 .await;
4568 }
4569 }
4570 Ok(())
4571 }
4572
4573 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4603 let round_idx = self.state.reviews.len() - 1;
4604 let needs_catchup_run = matches!(
4612 self.state.reviews[round_idx].e2e_status(),
4613 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4614 );
4615 if needs_catchup_run {
4616 let round = self.state.reviews[round_idx].round;
4617 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4618 let commands = self.state.config.verify.e2e.clone();
4619 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4620 let cache_dir = self.state.config.cache_dir();
4621 let context = format!(
4622 "round {round}: verification unresolved, catching up before the final decision"
4623 );
4624 let (outcomes, verify_retried) = with_cache_lease(
4625 &mut self.state,
4626 cache_dir.as_deref(),
4627 "e2e",
4628 "e2e",
4629 worktree,
4630 &attempted_head,
4631 timeout,
4632 &context,
4633 |state, budget| {
4634 let shell = shell.to_vec();
4635 let commands = commands.clone();
4636 let context = context.clone();
4637 async move {
4638 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4639 .await
4640 }
4641 },
4642 )
4643 .await;
4644 let last = &mut self.state.reviews[round_idx];
4645 last.e2e = outcomes;
4646 last.verify_retried = verify_retried;
4647 last.verified_head = Some(attempted_head);
4654 last.verified_at = Some(Timestamp::now());
4655 if verify_inconclusive(&last.e2e) {
4656 self.state.save()?;
4663 return Ok(());
4664 }
4665 last.e2e_deferred = false;
4666 }
4667 let last = &self.state.reviews[round_idx];
4668 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4669
4670 match last.e2e_status() {
4671 E2eStatus::Failed => {
4672 let red: Vec<String> = last
4673 .e2e
4674 .iter()
4675 .filter(|o| !o.ok())
4676 .map(|o| {
4677 format!(
4678 "`{}` -> {:?}\n{}",
4679 o.command,
4680 o.code,
4681 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4682 )
4683 })
4684 .collect();
4685 self.state
4686 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4687 self.state.status = RunStatus::Blocked;
4688 }
4689 E2eStatus::ResourceBlocked => {
4694 self.state.event(
4695 "review",
4696 format!(
4697 "{why}; e2e could not run (shared build cache unavailable); not \
4698 deciding yet"
4699 ),
4700 );
4701 }
4702 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4703 self.state.event(
4704 "review",
4705 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4706 );
4707 self.state.status = RunStatus::Gating;
4708 }
4709 }
4710 self.state.save()?;
4711 Ok(())
4712 }
4713
4714 async fn gate(&mut self) -> Result<()> {
4717 if self.state.status == RunStatus::Failed
4729 || self
4730 .state
4731 .base_sync
4732 .as_ref()
4733 .is_some_and(|s| s.conflict.is_some())
4734 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
4735 != Some(RunStatus::Gating)
4736 {
4737 return Ok(());
4738 }
4739 if self.state.gate_ran {
4740 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
4751 self.state.status = RunStatus::Blocked;
4752 self.state.save()?;
4753 }
4754 return Ok(());
4755 }
4756 let Some(winner) = self.state.winner().cloned() else {
4757 return Ok(());
4758 };
4759 self.state.status = RunStatus::Gating;
4760 let mut outcomes = self.run_gate(&winner).await?;
4761 loop {
4762 if verify_inconclusive(&outcomes) {
4773 self.state.save()?;
4774 return Ok(());
4775 }
4776 if outcomes.iter().all(CommandOutcome::ok) {
4777 break;
4778 }
4779 match self.gate_fix_round(&winner, &outcomes).await? {
4780 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
4781 GateFix::Stop => break,
4782 GateFix::Defer => {
4783 self.state.save()?;
4784 return Ok(());
4785 }
4786 }
4787 }
4788 let passed = outcomes.iter().all(CommandOutcome::ok);
4789 self.state.gate = outcomes;
4790 self.state.gate_ran = true;
4791 if !passed {
4792 self.state.status = RunStatus::Blocked;
4793 let spent = self.state.gate_fixes.len();
4794 self.state.event(
4795 "gate",
4796 if spent == 0 {
4797 "gate failed; not merging".to_owned()
4798 } else {
4799 format!("gate failed after {spent} gate-fix round(s); not merging")
4800 },
4801 );
4802 }
4803 self.state.save()?;
4804 Ok(())
4805 }
4806
4807 async fn run_pre_gate(&mut self, winner: &Candidate) {
4817 let commands = self.state.config.verify.pre_gate.clone();
4818 if commands.is_empty() {
4819 return;
4820 }
4821 let shell = self.state.config.shell();
4822 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4823 let (outcomes, _) = run_commands(
4824 &mut self.state,
4825 "pre_gate",
4826 "pre_gate",
4827 0,
4828 &shell,
4829 &commands,
4830 &winner.worktree,
4831 timeout,
4832 )
4833 .await;
4834 for o in &outcomes {
4835 if !o.ok() {
4836 tracing::warn!(
4837 "pre_gate `{}` failed ({:?}); the gate decides",
4838 o.command,
4839 o.code
4840 );
4841 }
4842 self.state.event(
4843 "pre_gate",
4844 format!(
4845 "`{}` -> {}",
4846 o.command,
4847 if o.ok() {
4848 "pass".to_owned()
4849 } else {
4850 format!(
4851 "FAIL ({:?})\n{}",
4852 o.code,
4853 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4854 )
4855 }
4856 ),
4857 );
4858 }
4859 self.state.pre_gate = outcomes;
4860 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
4861 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
4862 Ok(head) => {
4863 self.state
4864 .event("pre_gate", format!("committed mechanical fixes ({head})"));
4865 self.state.pre_gate_commit = Some(head);
4866 }
4867 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
4868 },
4869 Ok(false) => {}
4870 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
4871 }
4872 if let Err(e) = self.state.save() {
4873 tracing::warn!("could not persist the pre_gate record: {e:#}");
4874 }
4875 }
4876
4877 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
4880 self.run_pre_gate(winner).await;
4881 let shell = self.state.config.shell();
4882 let gate_commands = self.state.config.verify.gate.clone();
4883 let outcomes = if gate_commands.is_empty() {
4892 Vec::new()
4893 } else {
4894 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4895 let cache_dir = self.state.config.cache_dir();
4896 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4897 let (outcomes, _) = with_cache_lease(
4898 &mut self.state,
4899 cache_dir.as_deref(),
4900 "gate",
4901 "gate",
4902 &winner.worktree,
4903 &head,
4904 timeout,
4905 "final gate",
4906 |state, budget| {
4907 let shell = shell.clone();
4908 let gate_commands = gate_commands.clone();
4909 let worktree = winner.worktree.clone();
4910 async move {
4911 let (outcomes, timed_out_pids) = run_commands(
4912 state,
4913 "gate",
4914 "gate",
4915 0,
4916 &shell,
4917 &gate_commands,
4918 &worktree,
4919 budget,
4920 )
4921 .await;
4922 (outcomes, false, timed_out_pids)
4923 }
4924 },
4925 )
4926 .await;
4927 outcomes
4928 };
4929 if outcomes.is_empty() {
4930 self.state.event(
4935 "gate",
4936 "no gate commands configured; nothing to check, passing",
4937 );
4938 }
4939 for o in &outcomes {
4940 self.state.event(
4941 "gate",
4942 format!(
4943 "`{}` -> {}",
4944 o.command,
4945 if o.ok() {
4946 "pass".to_owned()
4947 } else {
4948 format!(
4949 "FAIL ({:?})\n{}",
4950 o.code,
4951 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4952 )
4953 }
4954 ),
4955 );
4956 }
4957 Ok(outcomes)
4958 }
4959
4960 async fn gate_fix_round(
4973 &mut self,
4974 winner: &Candidate,
4975 outcomes: &[CommandOutcome],
4976 ) -> Result<GateFix> {
4977 let cap = self.state.config.graph.gate_fix_rounds;
4978 let spent = self.state.gate_fixes.len();
4979 if spent >= cap {
4980 if cap > 0 {
4981 self.state.event(
4982 "gate",
4983 format!("{spent} gate-fix round(s) spent and the gate still fails"),
4984 );
4985 }
4986 return Ok(GateFix::Stop);
4987 }
4988 if !gate_fixable(outcomes) {
4989 self.state.event(
4990 "gate",
4991 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
4992 command or similar); not spending a fix round on it",
4993 );
4994 return Ok(GateFix::Stop);
4995 }
4996 let min_free = self.state.config.disk.min_free_bytes;
4997 if min_free > 0 {
4998 match crate::disk::free_bytes(&winner.worktree) {
4999 Ok(free) if crate::disk::enough_space(free, min_free) => {}
5000 Ok(free) => {
5001 self.state.event(
5002 "gate",
5003 format!(
5004 "only {free} bytes free ({min_free} required by `[disk] \
5005 min_free_bytes`); not spending a fix round on a failure the disk \
5006 may explain"
5007 ),
5008 );
5009 return Ok(GateFix::Stop);
5010 }
5011 Err(e) => {
5012 self.state.event(
5013 "gate",
5014 format!("free disk space could not be measured ({e:#}); no fix round"),
5015 );
5016 return Ok(GateFix::Stop);
5017 }
5018 }
5019 }
5020
5021 let attempt = spent + 1;
5022 let run_id = self.state.id.clone();
5023 let prompts = self.state.config.prompts.clone();
5024 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
5025 let base = self.landing_base();
5026 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
5027 let seat = self.seat(&fix_seat_key, &fix_spec.id);
5028 let job = SeatJob {
5029 prompt: prompt::gate_fix(
5030 &self.state.instruction,
5031 &failed,
5032 attempt,
5033 cap,
5034 &self.state.config.graph.language,
5035 ),
5036 spec: fix_spec,
5037 seat,
5038 cwd: winner.worktree.clone(),
5039 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
5040 allow_write: true,
5041 sessions: self.state.config.graph.sessions,
5042 artifacts: agent::artifacts_dir(&self.state.dir()),
5043 stem: format!("gate-fix-{attempt}"),
5044 };
5045 self.state.event(
5046 "gate",
5047 format!("gate failed; gate-fix round {attempt} of {cap}"),
5048 );
5049 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
5050 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
5051 let cache = self.state.config.cache_dir();
5052 let ctx = WaveCtx {
5053 run: &run_id,
5054 node: "gate-fix",
5055 prompts: &prompts,
5056 cache: cache.as_deref(),
5057 round: None,
5058 };
5059 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
5060 let mut record = GateFixRecord {
5061 agent: seat.agent.clone(),
5062 failed,
5063 notes: String::new(),
5064 committed: false,
5065 error: None,
5066 };
5067 match out {
5068 AgentOutcome::Ok(o) => {
5069 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
5072 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
5073 }
5074 }
5075 AgentOutcome::Dropped(_) => {
5076 record.error = Some("the CLI dropped the stream".to_owned());
5077 }
5078 AgentOutcome::Quota(o) => {
5079 self.state.quota.push(QuotaLoss {
5080 seat: seat.key.clone(),
5081 node: "gate-fix".to_owned(),
5082 at: Timestamp::now(),
5083 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5084 });
5085 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5086 }
5087 AgentOutcome::Failed(e) => record.error = Some(e),
5088 }
5089 self.state.seats.insert(seat.key.clone(), seat);
5090 if let Ok(r) = git::rescue_commit(
5091 &winner.worktree,
5092 &format!("magi: gate fix {attempt} (uncommitted work)"),
5093 )
5094 .await
5095 {
5096 self.state.note_withheld("gate-fix", &r.withheld);
5097 }
5098 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5099 record.committed = after != before;
5100 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5101 let note = record.error.clone();
5102 self.state.gate_fixes.push(record);
5103 self.state.save()?;
5104 if !changed {
5105 self.state.event(
5106 "gate",
5107 match note {
5108 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5109 None => format!("gate-fix round {attempt}: the tree did not change"),
5110 },
5111 );
5112 return Ok(GateFix::Stop);
5113 }
5114 self.state.event(
5115 "gate",
5116 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5117 );
5118
5119 let commands = self.state.config.verify.e2e.clone();
5120 if !commands.is_empty() {
5121 let shell = self.state.config.shell();
5122 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5123 let cache_dir = self.state.config.cache_dir();
5124 let context = format!("gate-fix round {attempt}");
5125 let (e2e, _) = with_cache_lease(
5126 &mut self.state,
5127 cache_dir.as_deref(),
5128 "e2e",
5129 "e2e",
5130 &winner.worktree,
5131 &after,
5132 timeout,
5133 &context,
5134 |state, budget| {
5135 let shell = shell.clone();
5136 let commands = commands.clone();
5137 let context = context.clone();
5138 let worktree = winner.worktree.clone();
5139 async move {
5140 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5141 .await
5142 }
5143 },
5144 )
5145 .await;
5146 if verify_inconclusive(&e2e) {
5147 return Ok(GateFix::Defer);
5148 }
5149 if e2e.iter().any(|o| !o.ok()) {
5150 self.state.event(
5151 "gate",
5152 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5153 );
5154 return Ok(GateFix::Stop);
5155 }
5156 }
5157 Ok(GateFix::Retry)
5158 }
5159
5160 async fn merge(&mut self) -> Result<()> {
5163 if self
5178 .state
5179 .base_sync
5180 .as_ref()
5181 .is_some_and(|s| s.conflict.is_some())
5182 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5183 != Some(RunStatus::Gating)
5184 || !self.state.gate_status().ok()
5193 {
5194 return Ok(());
5195 }
5196 if self.state.merge.is_some() {
5205 return Ok(());
5206 }
5207 let Some(winner) = self.state.winner().cloned() else {
5208 return Ok(());
5209 };
5210 let repo = self.state.repo.clone();
5211 let base = self.state.base_branch.clone();
5212 let mode = self.state.config.merge.mode;
5213 let style = self.state.config.merge.style;
5214 let pr = pr_message(&self.state, winner.label);
5215 let message = pr.commit_message();
5216
5217 let outcome = match mode {
5218 MergeMode::None => MergeOutcome {
5219 mode,
5220 ok: true,
5221 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5222 },
5223 MergeMode::Local => {
5224 let on = git::current_branch(&repo).await?;
5225 if on.as_deref() != Some(base.as_str()) {
5226 MergeOutcome {
5227 mode,
5228 ok: false,
5229 detail: format!(
5230 "{} has {} checked out, not the base branch {base}",
5231 repo.display(),
5232 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5233 ),
5234 }
5235 } else if !git::is_clean(&repo).await? {
5236 MergeOutcome {
5237 mode,
5238 ok: false,
5239 detail: format!("{} is dirty; refusing to merge", repo.display()),
5240 }
5241 } else {
5242 let out = match style {
5243 MergeStyle::Merge => {
5244 git::merge_no_ff(&repo, &winner.branch, &message).await?
5245 }
5246 MergeStyle::Squash => {
5247 git::merge_squash(&repo, &winner.branch, &message).await?
5248 }
5249 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5250 };
5251 MergeOutcome {
5252 mode,
5253 ok: out.ok(),
5254 detail: if out.ok() { out.stdout } else { out.stderr },
5255 }
5256 }
5257 }
5258 MergeMode::Pr => {
5259 let remote = self.state.config.merge.remote.clone();
5260 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5261 if !pushed.ok() {
5262 MergeOutcome {
5263 mode,
5264 ok: false,
5265 detail: pushed.stderr,
5266 }
5267 } else {
5268 let out =
5269 gh_pr_create(&winner.worktree, &base, &winner.branch, &pr.title, &pr.body)
5270 .await;
5271 match out {
5272 Ok(url) => MergeOutcome {
5273 mode,
5274 ok: true,
5275 detail: url,
5276 },
5277 Err(e) => MergeOutcome {
5278 mode,
5279 ok: false,
5280 detail: e.to_string(),
5281 },
5282 }
5283 }
5284 }
5285 };
5286
5287 self.state.status = match (mode, outcome.ok) {
5288 (MergeMode::None, _) => RunStatus::Ready,
5289 (_, true) => RunStatus::Merged,
5290 (_, false) => RunStatus::Blocked,
5291 };
5292 self.state.event(
5293 "merge",
5294 format!(
5295 "{:?}: {}",
5296 mode,
5297 outcome.detail.lines().next().unwrap_or("")
5298 ),
5299 );
5300 self.state.merge = Some(outcome);
5301 self.state.save()?;
5302
5303 if self.state.config.graph.land
5309 && mode == MergeMode::Pr
5310 && self.state.status == RunStatus::Merged
5311 {
5312 self.run_land().await?;
5313 }
5314 self.settle_questions();
5319 Ok(())
5320 }
5321
5322 async fn run_land(&mut self) -> Result<()> {
5333 let url = self
5334 .state
5335 .merge
5336 .as_ref()
5337 .map(|m| m.detail.clone())
5338 .unwrap_or_default();
5339 let url = url.lines().next().unwrap_or("").trim().to_owned();
5340 if !url.starts_with("http") {
5341 return Ok(());
5342 }
5343 match land::land(&mut self.state, &url).await {
5346 Ok(pr) if self.state.parked => {
5347 let _ = pr;
5351 }
5352 Ok(pr) => {
5353 self.state.status = match pr.state {
5354 land::PrLifecycle::Merged => RunStatus::Merged,
5355 _ => RunStatus::Blocked,
5356 };
5357 if bump::should_release_bump(self.state.status)
5364 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5365 {
5366 self.state
5372 .event("bump", format!("release bump skipped: {e:#}"));
5373 }
5374 self.state.save()?;
5375 }
5376 Err(e) => {
5377 self.state.status = RunStatus::Blocked;
5378 self.state.event("land", format!("gave up: {e}"));
5379 self.state.save()?;
5380 }
5381 }
5382 Ok(())
5383 }
5384
5385 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5389 if let Some(existing) = self.state.seats.get(key)
5390 && existing.agent == agent
5391 {
5392 return existing.clone();
5393 }
5394 let fresh = SeatState::new(key, agent, self.state.seed);
5395 self.state.seats.insert(key.to_owned(), fresh.clone());
5396 fresh
5397 }
5398
5399 fn view(&self, c: &Candidate) -> CandidateView {
5401 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5402 .unwrap_or_default();
5403 let (patch, _) = blind::sanitize_patch(
5404 &format!("candidate {} patch", c.label),
5405 &raw,
5406 &self.state.config.blind,
5407 );
5408 CandidateView {
5409 label: c.label,
5410 branch: c.branch.clone(),
5411 summary: c.summary.clone(),
5412 stat: c.stat.clone(),
5413 patch,
5414 }
5415 }
5416
5417 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5419 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5420 prompt::judge(
5421 "(see above)",
5422 &views,
5423 self.roles.judges.len(),
5424 base_short,
5425 "en",
5426 )
5427 }
5428
5429 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5436 let mut turns = Vec::new();
5437 for j in &self.state.judgements {
5438 if j.ranking.is_empty() {
5439 continue;
5440 }
5441 let reasons = j
5442 .reasons
5443 .iter()
5444 .map(|(k, v)| format!("- {k}: {v}"))
5445 .collect::<Vec<_>>()
5446 .join("\n");
5447 turns.push(Turn {
5448 who: format!("Judge {} (opening ranking)", j.judge),
5449 is_self: j.judge == self_idx + 1,
5450 body: format!(
5451 "Ranked {}{}{reasons}",
5452 j.ranking.iter().collect::<String>(),
5453 if reasons.is_empty() {
5454 ""
5455 } else {
5456 ", because:\n"
5457 }
5458 ),
5459 });
5460 }
5461 for t in self
5462 .state
5463 .deliberation
5464 .iter()
5465 .flat_map(|r| r.turns.iter())
5466 .chain(current)
5467 {
5468 turns.push(Turn {
5469 who: format!("Judge {}", t.judge),
5470 is_self: t.judge == self_idx + 1,
5471 body: t.body.clone(),
5472 });
5473 }
5474 turns
5475 }
5476}
5477
5478fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5480 agent::has_session(spec.kind, seat, sessions)
5481}
5482
5483fn next_untried_implementer<'a>(
5504 roster: &'a [AgentSpec],
5505 start: usize,
5506 tried: &BTreeSet<String>,
5507) -> Option<&'a AgentSpec> {
5508 roster
5509 .get(start + 1..)?
5510 .iter()
5511 .find(|s| !tried.contains(&s.id))
5512}
5513
5514fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5528 commands.iter().any(|c| c.exit_code.is_none())
5529}
5530
5531fn verified_noop_claim(
5544 usable: bool,
5545 commands: &[agent::CommandEvidence],
5546 text: &str,
5547) -> Option<String> {
5548 (usable && !has_unconfirmed_command(commands))
5549 .then(|| verdict::verified_noop(text))
5550 .flatten()
5551}
5552
5553fn short(commit: &str) -> String {
5554 commit.chars().take(7).collect()
5555}
5556
5557fn make_executable(path: &Path) -> Result<()> {
5558 #[cfg(unix)]
5559 {
5560 use std::os::unix::fs::PermissionsExt as _;
5561 let mut perms = std::fs::metadata(path)?.permissions();
5562 perms.set_mode(0o755);
5563 std::fs::set_permissions(path, perms)?;
5564 }
5565 #[cfg(not(unix))]
5566 {
5567 let _ = path;
5568 }
5569 Ok(())
5570}
5571
5572struct WaveCtx<'a> {
5579 run: &'a str,
5582 node: &'a str,
5584 prompts: &'a Prompts,
5585 cache: Option<&'a Path>,
5587 round: Option<usize>,
5590}
5591
5592async fn run_one(
5594 job: SeatJob,
5595 sem: Arc<Semaphore>,
5596 ctx: &WaveCtx<'_>,
5597 state: &mut RunState,
5598 attempt: usize,
5599) -> (SeatState, AgentOutcome) {
5600 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5601 .await
5602 .pop()
5603 .expect("one job in, one result out");
5604 (seat, out)
5605}
5606
5607async fn wave(
5613 jobs: Vec<SeatJob>,
5614 sem: Arc<Semaphore>,
5615 ctx: &WaveCtx<'_>,
5616 state: &mut RunState,
5617 attempt: usize,
5618) -> Vec<(usize, SeatState, AgentOutcome)> {
5619 let WaveCtx {
5620 run,
5621 node,
5622 prompts,
5623 cache,
5624 round,
5625 } = *ctx;
5626 for job in &jobs {
5627 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5628 }
5629 if let Err(e) = state.save() {
5630 tracing::warn!("could not persist in-progress seats: {e:#}");
5635 }
5636 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5652 let wait_started = Instant::now();
5653 let cache_guard = if let Some(cache_dir) = cache {
5654 if jobs_had_a_writer {
5655 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5656 let budget = jobs
5657 .iter()
5658 .map(|j| j.timeout)
5659 .max()
5660 .unwrap_or(Duration::from_secs(60));
5661 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5662 .await
5663 .ok()
5664 } else {
5665 None
5666 }
5667 } else {
5668 None
5669 };
5670 let waited_for_lease = wait_started.elapsed();
5677 let mut set = tokio::task::JoinSet::new();
5678 let overlay = prompts.overlay(node);
5679 for (i, mut job) in jobs.into_iter().enumerate() {
5680 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5681 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
5682 if cache.is_some() {
5683 job.prompt.push('\n');
5684 job.prompt
5685 .push_str(&prompt::build_cache_note(node, job.allow_write));
5686 }
5687 let sem = Arc::clone(&sem);
5688 let run = run.to_owned();
5689 let node = node.to_owned();
5690 let cache = cache
5701 .filter(|_| job.allow_write && cache_guard.is_some())
5702 .map(Path::to_path_buf);
5703 set.spawn(async move {
5704 let _permit = sem.acquire().await;
5705 let mut seat = job.seat;
5706 let out = agent::invoke(
5707 &job.spec,
5708 &mut seat,
5709 &Invocation {
5710 cwd: &job.cwd,
5711 prompt: &job.prompt,
5712 timeout: job.timeout,
5713 allow_write: job.allow_write,
5714 sessions: job.sessions,
5715 artifacts: &job.artifacts,
5716 stem: &job.stem,
5717 run: &run,
5718 node: &node,
5719 cache_dir: cache.as_deref(),
5720 attachments: &[],
5721 },
5722 )
5723 .await;
5724 let out = match out {
5725 Ok(o) if o.usable() => AgentOutcome::Ok(o),
5726 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
5727 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
5735 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
5736 Ok(o) => AgentOutcome::Failed(format!(
5737 "exited with {:?} and no usable output",
5738 o.exit_code
5739 )),
5740 Err(e) => AgentOutcome::Failed(e.to_string()),
5741 };
5742 (i, seat, out)
5743 });
5744 }
5745 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
5746 while let Some(joined) = set.join_next().await {
5747 let (i, seat, out) = match joined {
5748 Ok(v) => v,
5749 Err(e) => {
5753 tracing::error!("agent task panicked: {e}");
5754 continue;
5755 }
5756 };
5757 state.seat_finished(&seat.key);
5758 record_jobs(state, node, round, &seat.key, &out);
5759 if let Err(e) = state.save() {
5760 tracing::warn!("could not persist a seat's completion: {e:#}");
5761 }
5762 if collected.len() <= i {
5763 collected.resize_with(i + 1, || None);
5764 }
5765 collected[i] = Some((i, seat, out));
5766 }
5767 if state
5773 .active
5774 .values()
5775 .any(|a| a.node == node && a.attempt == attempt)
5776 {
5777 state
5778 .active
5779 .retain(|_, a| !(a.node == node && a.attempt == attempt));
5780 if let Err(e) = state.save() {
5781 tracing::warn!("could not persist the end of a wave: {e:#}");
5782 }
5783 }
5784 if let Some(cache_dir) = cache
5791 && jobs_had_a_writer
5792 {
5793 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
5794 }
5795 if let Some(guard) = cache_guard {
5796 guard.release();
5797 }
5798 collected.into_iter().flatten().collect()
5799}
5800
5801fn record_jobs(
5812 state: &mut RunState,
5813 node: &str,
5814 round: Option<usize>,
5815 seat: &str,
5816 out: &AgentOutcome,
5817) {
5818 let commands: &[agent::CommandEvidence] = match out {
5819 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
5820 AgentOutcome::Failed(_) => &[],
5821 };
5822 let checked_at = Timestamp::now();
5823 for c in commands {
5824 state.jobs.push(JobRecord {
5825 node: node.to_owned(),
5826 round,
5827 seat: seat.to_owned(),
5828 id: c.id.clone(),
5829 description: c.description.clone(),
5830 checked_at,
5831 status: match c.exit_code {
5832 Some(0) => JobStatus::Completed,
5833 Some(_) => JobStatus::Failed,
5834 None => JobStatus::Unknown,
5835 },
5836 exit_code: c.exit_code,
5837 result_summary: c.result_summary.clone(),
5838 source: c.source.clone(),
5839 });
5840 }
5841}
5842
5843fn round_is_clean(
5864 blocking: usize,
5865 e2e_ok: bool,
5866 answered: usize,
5867 expected: usize,
5868 quota_missing: usize,
5869 policy: IncompleteReviewPolicy,
5870) -> bool {
5871 if blocking != 0 || !e2e_ok {
5872 return false;
5873 }
5874 if answered == expected || policy == IncompleteReviewPolicy::Warn {
5875 return true;
5876 }
5877 answered > 0 && expected - answered <= quota_missing
5878}
5879
5880fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
5904 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
5905 return Some(RunStatus::Gating);
5906 }
5907 let last = reviews.last()?;
5908 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
5909 if reviews.len() < max_rounds && !stagnant {
5910 return None;
5911 }
5912 if last.incomplete() && last.blocking == 0 {
5913 return Some(RunStatus::Blocked);
5914 }
5915 if last.e2e_status() == E2eStatus::ResourceBlocked {
5916 return None;
5917 }
5918 Some(if last.e2e.iter().all(CommandOutcome::ok) {
5919 RunStatus::Gating
5920 } else {
5921 RunStatus::Blocked
5922 })
5923}
5924
5925fn retry_budget(full: Duration, nudged: bool) -> Duration {
5940 if nudged {
5941 (full / 4).max(Duration::from_secs(120)).min(full)
5942 } else {
5943 full
5944 }
5945}
5946
5947#[allow(clippy::too_many_arguments)]
5960async fn ask_json_wave<T>(
5961 jobs: Vec<SeatJob>,
5962 sem: Arc<Semaphore>,
5963 retries: usize,
5964 ctx: &WaveCtx<'_>,
5965 losses: &mut Vec<QuotaLoss>,
5966 state: &mut RunState,
5967 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
5968) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
5969where
5970 T: serde::de::DeserializeOwned + Send + 'static,
5971{
5972 let n = jobs.len();
5973 let originals: Vec<SeatJob> = jobs;
5974 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
5975 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
5976 let mut attempts_used: Vec<usize> = vec![0; n];
5983 let mut pending: Vec<usize> = (0..n).collect();
5984
5985 for attempt in 0..=retries {
5986 if pending.is_empty() {
5987 break;
5988 }
5989 let mut batch = Vec::with_capacity(pending.len());
5990 for &i in &pending {
5991 let src = &originals[i];
5992 let (prompt, timeout) = if attempt == 0 {
5995 (src.prompt.clone(), src.timeout)
5996 } else {
5997 let why = done[i]
5998 .as_ref()
5999 .and_then(|r| r.as_ref().err().map(ToString::to_string))
6000 .unwrap_or_else(|| "no parsable answer".to_owned());
6001 let nudge = prompt::nudge(&why);
6002 let nudged = has_context(&src.spec, &seats[i], src.sessions);
6003 let prompt = if nudged {
6004 nudge
6005 } else {
6006 format!("{}\n\n---\n\n{}", src.prompt, nudge)
6007 };
6008 (prompt, retry_budget(src.timeout, nudged))
6009 };
6010 batch.push(SeatJob {
6011 spec: src.spec.clone(),
6012 seat: seats[i].clone(),
6013 cwd: src.cwd.clone(),
6014 prompt,
6015 timeout,
6016 allow_write: src.allow_write,
6017 sessions: src.sessions,
6018 artifacts: src.artifacts.clone(),
6019 stem: if attempt == 0 {
6020 src.stem.clone()
6021 } else {
6022 format!("{}-retry{attempt}", src.stem)
6023 },
6024 });
6025 }
6026
6027 if attempt > 0 {
6028 let seats_out: Vec<&str> = pending
6029 .iter()
6030 .map(|&i| originals[i].seat.key.as_str())
6031 .collect();
6032 state.event(
6033 ctx.node,
6034 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
6035 );
6036 }
6037 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
6038 let mut still = Vec::new();
6039 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
6040 seats[i] = seat;
6041 let (parsed, quota) = match out {
6042 AgentOutcome::Ok(o) => (
6043 match verdict::extract_json::<T>(&o.text) {
6044 Ok(v) => match validate(&v) {
6045 Ok(()) => Ok((v, o)),
6046 Err(e) => Err(e),
6047 },
6048 Err(e) => Err(e),
6049 },
6050 false,
6051 ),
6052 AgentOutcome::Quota(o) => {
6053 losses.push(QuotaLoss {
6054 seat: originals[i].seat.key.clone(),
6055 node: ctx.node.to_owned(),
6056 at: Timestamp::now(),
6057 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
6058 });
6059 (
6060 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
6061 true,
6062 )
6063 }
6064 AgentOutcome::Dropped(o) => {
6069 let why = o
6070 .dropped
6071 .as_ref()
6072 .map(|d| d.why.as_str())
6073 .unwrap_or("the CLI ended the stream without delivering its answer");
6074 (
6075 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
6076 false,
6077 )
6078 }
6079 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
6080 };
6081 let failed = parsed.is_err();
6082 done[i] = Some(parsed);
6083 attempts_used[i] = attempt;
6084 if failed && !quota {
6087 still.push(i);
6088 }
6089 }
6090 pending = still;
6091 }
6092
6093 seats
6094 .into_iter()
6095 .zip(done)
6096 .zip(attempts_used)
6097 .map(|((seat, res), attempts)| {
6098 (
6099 seat,
6100 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6101 attempts,
6102 )
6103 })
6104 .collect()
6105}
6106
6107async fn acquire_cache_lease(
6120 state: &mut RunState,
6121 cache_dir: &Path,
6122 owner: &crate::cache::Owner,
6123 budget: Duration,
6124 context: &str,
6125) -> Result<crate::cache::Guard> {
6126 let home = crate::run::home();
6127 let started = Instant::now();
6128 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6129 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6130 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6131 Err(e) => {
6132 state.event(
6133 "verify",
6134 format!("{context}: could not check the shared build cache: {e:#}"),
6135 );
6136 if let Err(e2) = state.save() {
6137 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6138 }
6139 return Err(e);
6140 }
6141 };
6142 state.event(
6143 "verify",
6144 format!(
6145 "{context}: waiting for the shared build cache at {} ({})",
6146 cache_dir.display(),
6147 busy.describe()
6148 ),
6149 );
6150 if let Err(e) = state.save() {
6151 tracing::warn!("could not persist a cache wait: {e:#}");
6152 }
6153 let remaining = budget.saturating_sub(started.elapsed());
6154 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6155 Ok(g) => Ok(g),
6156 Err(e) => {
6157 state.event("verify", format!("{context}: {e:#}"));
6158 if let Err(e2) = state.save() {
6159 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6160 }
6161 Err(e)
6162 }
6163 }
6164}
6165
6166#[allow(clippy::too_many_arguments)]
6187async fn with_cache_lease<'s, F, Fut>(
6188 state: &'s mut RunState,
6189 cache_dir: Option<&Path>,
6190 node: &str,
6191 seat: &str,
6192 worktree: &Path,
6193 head: &str,
6194 budget: Duration,
6195 context: &str,
6196 body: F,
6197) -> (Vec<CommandOutcome>, bool)
6198where
6199 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6200 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6201{
6202 let Some(cache_dir) = cache_dir else {
6203 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6204 return (outcomes, retried);
6205 };
6206 let home = crate::run::home();
6207 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6208 let started = Instant::now();
6209 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6210 Ok(g) => g,
6211 Err(e) => {
6212 return (
6213 vec![CommandOutcome {
6214 command: "(waiting for the shared build cache)".to_owned(),
6215 code: None,
6216 output_tail: e.to_string(),
6217 duration_ms: started.elapsed().as_millis() as u64,
6218 resource_blocked: true,
6219 }],
6220 false,
6221 );
6222 }
6223 };
6224 let identity = crate::cache::Identity::new(worktree, head);
6225 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6226 state.event(
6234 "verify",
6235 format!(
6236 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6237 worktree.display(),
6238 short(head)
6239 ),
6240 );
6241 guard.release();
6242 return (
6243 vec![CommandOutcome {
6244 command: "(confirming the shared build cache is fresh)".to_owned(),
6245 code: None,
6246 output_tail: e.to_string(),
6247 duration_ms: started.elapsed().as_millis() as u64,
6248 resource_blocked: true,
6249 }],
6250 false,
6251 );
6252 }
6253 let remaining = budget.saturating_sub(started.elapsed());
6254 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6255 if !timed_out_pids.is_empty() {
6260 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6261 }
6262 guard.release();
6263 (outcomes, retried)
6264}
6265
6266async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6278 wait_for_pids_with(
6279 pids,
6280 crate::proc::pid_alive,
6281 LEASE_RELEASE_POLL,
6282 LEASE_RELEASE_MAX_WAIT,
6283 )
6284 .await;
6285}
6286
6287async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6293 pids: &[u32],
6294 alive: F,
6295 poll: Duration,
6296 max_wait: Duration,
6297) {
6298 let deadline = Instant::now() + max_wait;
6299 loop {
6300 if pids.iter().all(|&pid| !alive(pid)) {
6301 return;
6302 }
6303 if Instant::now() >= deadline {
6304 return;
6305 }
6306 tokio::time::sleep(poll).await;
6307 }
6308}
6309
6310fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6317 outcomes.iter().any(|o| o.resource_blocked)
6318}
6319
6320enum GateFix {
6322 Retry,
6324 Stop,
6327 Defer,
6330}
6331
6332fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6340 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6341 red.peek().is_some()
6342 && red.all(|o| {
6343 !o.resource_blocked
6344 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6345 && !o.output_tail.trim().is_empty()
6346 })
6347}
6348
6349fn e2e_outcome_label(o: &CommandOutcome) -> String {
6353 if o.ok() {
6354 return "pass".to_owned();
6355 }
6356 let reason = if o.build_failed() {
6357 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6358 } else {
6359 format!("FAIL ({:?})", o.code)
6360 };
6361 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6362}
6363
6364async fn run_e2e_with_retry(
6372 state: &mut RunState,
6373 shell: &[String],
6374 commands: &[String],
6375 worktree: &Path,
6376 timeout: Duration,
6377 context: &str,
6378) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6379 let (mut e2e, mut timed_out_pids) = run_commands(
6380 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6381 )
6382 .await;
6383 for o in &e2e {
6384 state.event(
6385 "verify",
6386 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6387 );
6388 }
6389 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6393 if verify_retried {
6394 state.event(
6395 "verify",
6396 format!(
6397 "{context}: verify could not build/link, not a test result — retrying once \
6398 before concluding"
6399 ),
6400 );
6401 let retried = run_commands(
6402 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6403 )
6404 .await;
6405 e2e = retried.0;
6406 timed_out_pids.extend(retried.1);
6409 for o in &e2e {
6410 state.event(
6411 "verify",
6412 format!(
6413 "{context}: retry `{}` -> {}",
6414 o.command,
6415 e2e_outcome_label(o)
6416 ),
6417 );
6418 }
6419 }
6420 (e2e, verify_retried, timed_out_pids)
6421}
6422
6423#[allow(clippy::too_many_arguments)]
6439async fn run_commands(
6440 state: &mut RunState,
6441 node: &str,
6442 task: &str,
6443 attempt: usize,
6444 shell: &[String],
6445 commands: &[String],
6446 cwd: &Path,
6447 timeout: Duration,
6448) -> (Vec<CommandOutcome>, Vec<u32>) {
6449 if commands.is_empty() {
6450 return (Vec::new(), Vec::new());
6455 }
6456 let mut out = Vec::new();
6457 let mut timed_out_pids = Vec::new();
6458 let total = commands.len();
6459 for (idx, command) in commands.iter().enumerate() {
6460 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6461 if let Err(e) = state.save() {
6462 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6463 }
6464 let started = Instant::now();
6465 let mut cmd = tokio::process::Command::new(&shell[0]);
6466 cmd.quiet();
6467 cmd.args(&shell[1..])
6468 .arg(command)
6469 .current_dir(cwd)
6470 .stdin(std::process::Stdio::null())
6471 .stdout(std::process::Stdio::piped())
6472 .stderr(std::process::Stdio::piped())
6473 .kill_on_drop(true);
6474 let spawned = cmd.spawn();
6475 let (code, body) = match spawned {
6476 Ok(child) => {
6477 let pid = child.id();
6482 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6483 Ok(Ok(o)) => {
6484 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6485 body.push_str(&String::from_utf8_lossy(&o.stderr));
6486 (o.status.code(), body)
6487 }
6488 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6489 Err(_) => {
6490 if let Some(pid) = pid {
6491 timed_out_pids.push(pid);
6492 }
6493 (None, format!("timed out after {}s", timeout.as_secs()))
6494 }
6495 }
6496 }
6497 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6498 };
6499 out.push(CommandOutcome {
6500 command: command.clone(),
6501 code,
6502 output_tail: tail(&body, OUTPUT_TAIL),
6503 duration_ms: started.elapsed().as_millis() as u64,
6504 resource_blocked: false,
6505 });
6506 }
6507 state.task_finished(task);
6508 if let Err(e) = state.save() {
6509 tracing::warn!("could not persist the end of {task}: {e:#}");
6510 }
6511 (out, timed_out_pids)
6512}
6513
6514fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6526 let repo = repo.display();
6527 match style {
6528 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6529 MergeStyle::Squash => {
6530 let subject = message
6533 .lines()
6534 .next()
6535 .unwrap_or(branch)
6536 .replace(['\\', '"', '$', '`'], "");
6537 format!(
6538 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6539 )
6540 }
6541 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6542 }
6543}
6544
6545const PR_TITLE_MAX: usize = 240;
6556
6557struct PrMessage {
6562 title: String,
6563 body: String,
6564}
6565
6566impl PrMessage {
6567 fn commit_message(&self) -> String {
6571 format!("{}\n\n{}", self.title, self.body)
6572 }
6573}
6574
6575fn title_marker(line: &str) -> Option<&str> {
6577 let line = line.trim();
6578 let head = line.get(..6)?;
6579 head.eq_ignore_ascii_case("title:")
6580 .then(|| line[6..].trim())
6581}
6582
6583fn summary_title(summary: &str) -> Option<String> {
6588 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6589 let raw = title_marker(first)?;
6590 if raw.is_empty() {
6591 return None;
6592 }
6593 let title = queue::title_from(raw, PR_TITLE_MAX);
6594 let lower = title.to_ascii_lowercase();
6595 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6596 return None;
6597 }
6598 Some(title)
6599}
6600
6601fn summary_without_title(summary: &str) -> String {
6604 let mut lines = summary.trim().lines().peekable();
6605 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6606 lines.next();
6607 }
6608 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6609}
6610
6611fn pr_message(state: &RunState, winner: char) -> PrMessage {
6625 let summary = state
6626 .candidates
6627 .iter()
6628 .find(|c| c.label == winner)
6629 .map(|c| c.summary.as_str())
6630 .unwrap_or_default();
6631 let title = summary_title(summary).unwrap_or_else(|| {
6634 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6635 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6636 t
6637 } else {
6638 format!(
6639 "chore: land candidate {} of run {}",
6640 winner.to_ascii_uppercase(),
6641 state.id
6642 )
6643 }
6644 });
6645
6646 let mut body = String::new();
6647 let what = summary_without_title(summary);
6648 if !what.is_empty() {
6649 body.push_str("## Summary\n\n");
6650 body.push_str(&what);
6651 body.push_str("\n\n");
6652 }
6653
6654 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6655 if let Some(fix) = fix
6656 && !fix.notes.trim().is_empty()
6657 {
6658 body.push_str("## Review fixes\n\n");
6659 body.push_str(fix.notes.trim());
6660 body.push_str("\n\n");
6661 }
6662
6663 let open = state.open_findings();
6664 if !open.is_empty() {
6665 body.push_str("## Open review findings\n\n");
6666 for f in &open {
6667 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6668 }
6669 body.push('\n');
6670 }
6671
6672 if let Some(fix) = fix
6673 && !fix.rejected.is_empty()
6674 {
6675 body.push_str("## Declined by the fixer\n\n");
6676 for r in &fix.rejected {
6677 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
6678 }
6679 body.push('\n');
6680 }
6681
6682 let task = state.instruction.trim();
6683 let task = if task.is_empty() {
6684 "(empty task)"
6685 } else {
6686 task
6687 };
6688 body.push_str(&format!(
6689 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
6690 task.replace("</details>", "</details>")
6691 ));
6692
6693 body.push_str(&format!(
6694 "\n---\nmagi:run/{} magi:candidate-{}\n",
6695 state.id,
6696 winner.to_ascii_lowercase()
6697 ));
6698
6699 let id = crate::scrub::Identity::current();
6702 PrMessage {
6703 title: crate::scrub::scrub(&title, &id),
6704 body: crate::scrub::scrub(&body, &id),
6705 }
6706}
6707
6708async fn gh_pr_create(
6710 cwd: &Path,
6711 base: &str,
6712 head: &str,
6713 title: &str,
6714 body: &str,
6715) -> Result<String> {
6716 let out = tokio::process::Command::new("gh")
6717 .args([
6718 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
6719 ])
6720 .current_dir(cwd)
6721 .quiet()
6722 .stdin(std::process::Stdio::null())
6723 .output()
6724 .await
6725 .context("spawn gh")?;
6726 if out.status.success() {
6727 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
6728 } else {
6729 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
6730 }
6731}
6732
6733pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
6742 let repo = state.repo.clone();
6743 let root = state.worktree_root();
6744 let winner = state.tally.as_ref().map(|t| t.winner);
6745 let mut removed = Vec::new();
6746
6747 for i in 0..state.candidates.len() {
6748 let c = state.candidates[i].clone();
6749 let is_winner = Some(c.label) == winner;
6750 if is_winner && !drop_winner {
6751 continue;
6752 }
6753 if c.worktree.exists() {
6754 git::worktree_remove(&repo, &c.worktree).await.ok();
6755 removed.push(c.worktree.to_string_lossy().into_owned());
6756 }
6757 if git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
6758 git::branch_delete(&repo, &c.branch).await.ok();
6759 removed.push(c.branch.clone());
6760 }
6761 state.candidates[i].folded = true;
6762 }
6763
6764 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
6765 let path = name.path();
6766 let keep = !drop_winner
6767 && winner.is_some_and(|w| {
6768 path.file_name()
6769 .is_some_and(|n| n == format!("cand-{w}").as_str())
6770 });
6771 if keep {
6772 continue;
6773 }
6774 git::worktree_remove(&repo, &path).await.ok();
6775 removed.push(path.to_string_lossy().into_owned());
6776 }
6777
6778 remove_if_empty(&root);
6787
6788 if state.enabled_worktree_config && drop_winner {
6789 git::release_worktree_config(&repo).await.ok();
6793 state.enabled_worktree_config = false;
6794 }
6795 state.save_under(home)?;
6796 Ok(removed)
6797}
6798
6799fn remove_if_empty(dir: &Path) {
6810 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
6811 std::fs::remove_dir(dir).ok();
6812 }
6813}
6814
6815pub fn worst_open(state: &RunState) -> Option<Severity> {
6817 state
6818 .reviews
6819 .last()?
6820 .reviews
6821 .iter()
6822 .flat_map(|r| r.findings.iter())
6823 .map(|f| f.severity)
6824 .max()
6825}
6826
6827#[cfg(test)]
6828mod tests {
6829 use super::*;
6830 use crate::run::GateStatus;
6831 use std::collections::BTreeMap;
6832 use std::time::Duration;
6833
6834 fn conductor() -> AgentSpec {
6835 AgentSpec {
6836 id: "conductor".to_owned(),
6837 kind: crate::config::AgentKind::Command,
6838 model: None,
6839 command: vec!["true".to_owned()],
6840 extra_args: Vec::new(),
6841 env: BTreeMap::new(),
6842 prompt_delivery: None,
6843 }
6844 }
6845
6846 fn spec(id: &str) -> AgentSpec {
6847 AgentSpec {
6848 id: id.to_owned(),
6849 kind: crate::config::AgentKind::Command,
6850 model: None,
6851 command: vec!["true".to_owned()],
6852 extra_args: Vec::new(),
6853 env: BTreeMap::new(),
6854 prompt_delivery: None,
6855 }
6856 }
6857
6858 #[test]
6865 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
6866 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6867 let tried = BTreeSet::from(["beta".to_owned()]);
6868 let next = next_untried_implementer(&roster, 1, &tried);
6871 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
6872 }
6873
6874 #[test]
6875 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
6876 let roster = vec![spec("alpha"), spec("beta")];
6877 let tried = BTreeSet::from(["beta".to_owned()]);
6878 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6882 }
6883
6884 #[test]
6885 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
6886 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6887 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
6888 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6892 }
6893
6894 #[test]
6895 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
6896 let roster = vec![spec("a"), spec("a"), spec("b")];
6897 let tried = BTreeSet::from(["a".to_owned()]);
6898 let next = next_untried_implementer(&roster, 0, &tried);
6899 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
6900 }
6901
6902 #[test]
6903 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
6904 let roster = vec![spec("a"), spec("b")];
6905 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
6906 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
6907 }
6908
6909 #[test]
6910 fn remove_if_empty_only_ever_takes_a_bare_directory() {
6911 let dir = tempfile::tempdir().unwrap();
6912 let bay = dir.path().join("ffff");
6913
6914 remove_if_empty(&bay);
6916 assert!(!bay.exists());
6917
6918 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
6921 remove_if_empty(&bay);
6922 assert!(bay.exists(), "non-empty directory must survive");
6923
6924 std::fs::remove_dir(bay.join("cand-A")).unwrap();
6926 remove_if_empty(&bay);
6927 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
6928 }
6929
6930 #[test]
6939 fn a_full_panel_that_found_nothing_is_clean() {
6940 assert!(round_is_clean(
6941 0,
6942 true,
6943 2,
6944 2,
6945 0,
6946 IncompleteReviewPolicy::Block
6947 ));
6948 }
6949
6950 #[test]
6951 fn a_missing_seat_is_never_clean_under_the_default_policy() {
6952 assert!(!round_is_clean(
6953 0,
6954 true,
6955 1,
6956 2,
6957 0,
6958 IncompleteReviewPolicy::Block
6959 ));
6960 }
6961
6962 #[test]
6963 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
6964 assert!(!round_is_clean(
6965 1,
6966 true,
6967 1,
6968 2,
6969 0,
6970 IncompleteReviewPolicy::Warn
6971 ));
6972 }
6973
6974 #[test]
6975 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
6976 assert!(round_is_clean(
6977 0,
6978 true,
6979 1,
6980 2,
6981 0,
6982 IncompleteReviewPolicy::Warn
6983 ));
6984 }
6985
6986 #[test]
6987 fn a_full_panel_with_an_open_finding_is_not_clean() {
6988 assert!(!round_is_clean(
6989 1,
6990 true,
6991 2,
6992 2,
6993 0,
6994 IncompleteReviewPolicy::Block
6995 ));
6996 }
6997
6998 #[test]
6999 fn a_full_panel_with_a_red_e2e_is_not_clean() {
7000 assert!(!round_is_clean(
7001 0,
7002 false,
7003 2,
7004 2,
7005 0,
7006 IncompleteReviewPolicy::Block
7007 ));
7008 }
7009
7010 #[test]
7017 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
7018 assert!(round_is_clean(
7021 0,
7022 true,
7023 1,
7024 2,
7025 1,
7026 IncompleteReviewPolicy::Block
7027 ));
7028 }
7029
7030 #[test]
7031 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
7032 assert!(!round_is_clean(
7035 0,
7036 true,
7037 1,
7038 2,
7039 0,
7040 IncompleteReviewPolicy::Block
7041 ));
7042 }
7043
7044 #[test]
7045 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
7046 assert!(!round_is_clean(
7047 1,
7048 true,
7049 1,
7050 2,
7051 1,
7052 IncompleteReviewPolicy::Block
7053 ));
7054 assert!(!round_is_clean(
7055 0,
7056 false,
7057 1,
7058 2,
7059 1,
7060 IncompleteReviewPolicy::Block
7061 ));
7062 }
7063
7064 #[test]
7065 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
7066 assert!(!round_is_clean(
7070 0,
7071 true,
7072 0,
7073 2,
7074 2,
7075 IncompleteReviewPolicy::Block
7076 ));
7077 }
7078
7079 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
7080 CommandOutcome {
7081 command: "test".to_owned(),
7082 code,
7083 output_tail: String::new(),
7084 duration_ms: 0,
7085 resource_blocked,
7086 }
7087 }
7088
7089 #[test]
7090 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7091 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7092 assert!(
7093 !verify_inconclusive(&[outcome(Some(1), false)]),
7094 "an ordinary failure is still evidence about the patch"
7095 );
7096 assert!(verify_inconclusive(&[outcome(None, true)]));
7097 assert!(
7098 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7099 "one inconclusive outcome taints the whole batch"
7100 );
7101 assert!(!verify_inconclusive(&[]));
7102 }
7103
7104 #[tokio::test]
7105 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7106 let calls = std::sync::atomic::AtomicUsize::new(0);
7110 let started = Instant::now();
7111 wait_for_pids_with(
7112 &[123],
7113 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7114 Duration::from_millis(5),
7115 Duration::from_secs(5),
7116 )
7117 .await;
7118 assert!(
7119 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7120 "must keep checking rather than deciding on the first answer"
7121 );
7122 assert!(
7123 started.elapsed() < Duration::from_secs(1),
7124 "must return the moment it is confirmed dead, not wait out the ceiling"
7125 );
7126 }
7127
7128 #[tokio::test]
7129 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7130 let started = Instant::now();
7131 wait_for_pids_with(
7132 &[123],
7133 |_| true, Duration::from_millis(5),
7135 Duration::from_millis(30),
7136 )
7137 .await;
7138 let elapsed = started.elapsed();
7139 assert!(
7140 elapsed >= Duration::from_millis(30),
7141 "must not give up before its own ceiling: {elapsed:?}"
7142 );
7143 assert!(
7144 elapsed < Duration::from_secs(1),
7145 "must not wait past its own ceiling either: {elapsed:?}"
7146 );
7147 }
7148
7149 #[tokio::test]
7150 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7151 let started = Instant::now();
7152 wait_for_pids_with(
7153 &[],
7154 |_| true,
7155 Duration::from_secs(5),
7156 Duration::from_secs(5),
7157 )
7158 .await;
7159 assert!(
7160 started.elapsed() < Duration::from_millis(200),
7161 "an empty pid list has nothing to confirm"
7162 );
7163 }
7164
7165 fn review_round(
7171 clean: bool,
7172 blocking: usize,
7173 answered: usize,
7174 expected: usize,
7175 progressed: bool,
7176 e2e_ok: bool,
7177 ) -> ReviewRound {
7178 ReviewRound {
7179 round: 1,
7180 head: "h".to_owned(),
7181 verified_head: None,
7182 verified_at: None,
7183 reviews: Vec::new(),
7184 e2e: vec![CommandOutcome {
7185 command: "test".to_owned(),
7186 code: Some(if e2e_ok { 0 } else { 1 }),
7187 output_tail: String::new(),
7188 duration_ms: 0,
7189 resource_blocked: false,
7190 }],
7191 verify_retried: false,
7192 e2e_deferred: false,
7193 e2e_defer_reason: None,
7194 fix: None,
7195 blocking,
7196 answered,
7197 expected,
7198 clean,
7199 progressed,
7200 vote_split: false,
7201 reconsideration: Vec::new(),
7202 verdict: None,
7203 }
7204 }
7205
7206 #[test]
7207 fn review_conclusion_is_none_when_nothing_has_run() {
7208 assert_eq!(review_conclusion(&[], 3), None);
7209 }
7210
7211 #[test]
7212 fn review_conclusion_is_none_while_rounds_remain() {
7213 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7214 assert_eq!(review_conclusion(&rounds, 3), None);
7215 }
7216
7217 #[test]
7218 fn review_conclusion_is_gating_once_a_round_is_clean() {
7219 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7220 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7221 }
7222
7223 #[test]
7224 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7225 let rounds = vec![
7226 review_round(false, 1, 2, 2, true, true),
7227 review_round(false, 1, 2, 2, true, true),
7228 ];
7229 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7230 }
7231
7232 #[test]
7233 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7234 let rounds = vec![
7235 review_round(false, 1, 2, 2, true, true),
7236 review_round(false, 1, 2, 2, true, false),
7237 ];
7238 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7239 }
7240
7241 #[test]
7242 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7243 let mut blocked = review_round(false, 1, 2, 2, true, false);
7250 blocked.e2e[0].resource_blocked = true;
7251 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7252 assert_eq!(review_conclusion(&rounds, 2), None);
7253 }
7254
7255 #[test]
7256 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7257 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7259 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7260 }
7261
7262 #[test]
7263 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7264 let rounds = vec![
7265 review_round(false, 1, 2, 2, false, true),
7266 review_round(false, 1, 2, 2, false, true),
7267 ];
7268 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7269 }
7270
7271 fn secs(n: u64) -> Duration {
7272 Duration::from_secs(n)
7273 }
7274
7275 fn init_repo(dir: &Path) {
7278 let run = |args: &[&str]| {
7279 let out = std::process::Command::new("git")
7280 .args(args)
7281 .current_dir(dir)
7282 .quiet()
7283 .output()
7284 .expect("spawn git");
7285 assert!(
7286 out.status.success(),
7287 "git {args:?} failed: {}",
7288 String::from_utf8_lossy(&out.stderr)
7289 );
7290 };
7291 run(&["init", "-b", "main"]);
7292 run(&["config", "user.name", "magi test"]);
7293 run(&["config", "user.email", "magi@example.com"]);
7294 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7295 run(&["add", "-A"]);
7296 run(&["commit", "-m", "init"]);
7297 }
7298
7299 fn ask_test_home() {
7307 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7308 }
7309
7310 fn runner_at(status: RunStatus) -> Runner {
7313 let mut state = RunState::new(
7314 PathBuf::from("/nonexistent/repo"),
7315 "main".to_owned(),
7316 "deadbeef".to_owned(),
7317 "task".to_owned(),
7318 Config::default(),
7319 );
7320 state.status = status;
7321 Runner {
7322 state,
7323 roles: ResolvedRoles {
7324 implementers: Vec::new(),
7325 judges: Vec::new(),
7326 reviewers: Vec::new(),
7327 fixer: None,
7328 conductor: conductor(),
7329 implementer_roster: Vec::new(),
7330 },
7331 sem: Arc::new(Semaphore::new(1)),
7332 pause: Pause::new(),
7333 interrupt: Pause::new(),
7334 }
7335 }
7336
7337 #[test]
7341 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7342 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7343 let mut runner = runner_at(RunStatus::Implementing);
7344 let interrupt = Pause::new();
7345 runner.watch_interrupt(interrupt.clone());
7346
7347 interrupt.park_because("task a1b2 asked to run first");
7348
7349 assert!(runner.park_here().expect("park_here"));
7350 assert!(runner.state.parked);
7351 let last = runner.state.events.last().expect("a park event");
7352 assert_eq!(last.node, "park");
7353 assert!(
7354 last.message.contains("task a1b2 asked to run first"),
7355 "expected the interrupt reason in {:?}",
7356 last.message
7357 );
7358 }
7359
7360 #[test]
7368 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7369 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7370 let mut runner = runner_at(RunStatus::Implementing);
7371 let shutdown = Pause::new();
7372 runner.on_pause(shutdown.clone());
7373 let interrupt = Pause::new();
7374 runner.watch_interrupt(interrupt.clone());
7375
7376 assert!(!runner.park_here().expect("park_here"));
7378 assert!(!runner.state.parked);
7379
7380 interrupt.park_because("test");
7382 assert!(!shutdown.parked());
7383 assert!(runner.park_here().expect("park_here"));
7384 }
7385
7386 #[tokio::test]
7400 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7401 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7402 let mut runner = runner_at(RunStatus::Implementing);
7403 let interrupt = Pause::new();
7404 runner.watch_interrupt(interrupt.clone());
7405
7406 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7407 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7408
7409 let node = async move {
7413 started_tx.send(()).expect("send started");
7414 finish_rx.await.expect("recv finish");
7415 "node finished"
7416 };
7417
7418 let interrupter = async move {
7419 started_rx.await.expect("recv started");
7420 interrupt.park_because("higher-priority task waiting");
7422 tokio::task::yield_now().await;
7426 finish_tx.send(()).expect("send finish");
7427 };
7428
7429 let (node_result, ()) = tokio::join!(node, interrupter);
7430 assert_eq!(
7431 node_result, "node finished",
7432 "the in-flight call ran to completion"
7433 );
7434
7435 assert!(runner.park_here().expect("park_here"));
7438 assert!(runner.state.parked);
7439 }
7440
7441 #[test]
7447 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7448 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7449 let mut runner = runner_at(RunStatus::Judging);
7450 runner.state.config.agents = vec![conductor()];
7454 runner.state.candidates = vec![Candidate {
7455 index: 0,
7456 label: 'A',
7457 agent: "alpha".to_owned(),
7458 branch: "magi/x/A".to_owned(),
7459 worktree: PathBuf::from("/nonexistent/worktree"),
7460 summary: "did the thing".to_owned(),
7461 stat: "1 file changed".to_owned(),
7462 files: 1,
7463 commits: 1,
7464 empty: false,
7465 failed: None,
7466 verified_noop: None,
7467 duration_ms: 1234,
7468 folded: false,
7469 }];
7470 let run_id = runner.state.id.clone();
7471
7472 let interrupt = Pause::new();
7473 runner.watch_interrupt(interrupt.clone());
7474 interrupt.park_because("task c3d4 asked to run first");
7475 assert!(runner.park_here().expect("park_here"));
7476
7477 let resumed = Runner::resume(&run_id).expect("resume");
7478 assert_eq!(resumed.state.candidates.len(), 1);
7479 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7480 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7481 assert_eq!(resumed.state.status, runner.state.status);
7482 assert!(
7483 resumed.state.parked,
7484 "still parked until `execute` actually walks the graph again"
7485 );
7486 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7487 }
7488
7489 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7491 let mut q = ask::Question::new(
7492 run.to_owned(),
7493 "implement".to_owned(),
7494 "impl-A".to_owned(),
7495 "Which storage backend should the cache use?".to_owned(),
7496 String::new(),
7497 vec!["SQLite".to_owned(), "Redis".to_owned()],
7498 );
7499 store.put(&mut q).unwrap();
7500 q
7501 }
7502
7503 #[test]
7504 fn a_failed_runs_open_question_is_abandoned() {
7505 ask_test_home();
7506 let store = ask::Questions::open();
7507 let mut runner = runner_at(RunStatus::Failed);
7508 let run = runner.state.id.clone();
7509 let q = ask_open_question(&store, &run);
7510
7511 runner.settle_questions();
7512
7513 let back = store.get(&q.id).unwrap();
7514 assert!(
7515 !back.status.open(),
7516 "the seat that asked died with the run; nobody is left to read an answer"
7517 );
7518 assert!(
7519 back.detail.contains(&run) && back.detail.contains("failed"),
7520 "the reason names what the run became, not just that it is gone: {}",
7521 back.detail
7522 );
7523 }
7524
7525 #[test]
7526 fn a_merged_runs_open_question_is_abandoned_too() {
7527 ask_test_home();
7528 let store = ask::Questions::open();
7529 for status in [RunStatus::Merged, RunStatus::Ready] {
7532 let mut runner = runner_at(status);
7533 let run = runner.state.id.clone();
7534 let q = ask_open_question(&store, &run);
7535
7536 runner.settle_questions();
7537
7538 let back = store.get(&q.id).unwrap();
7539 assert!(
7540 !back.status.open(),
7541 "{status:?} run's question must not outlive the run"
7542 );
7543 }
7544 }
7545
7546 #[test]
7547 fn a_still_resumable_runs_open_question_is_left_alone() {
7548 ask_test_home();
7549 let store = ask::Questions::open();
7550 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7556 let mut runner = runner_at(status);
7557 let run = runner.state.id.clone();
7558 let q = ask_open_question(&store, &run);
7559
7560 runner.settle_questions();
7561
7562 let back = store.get(&q.id).unwrap();
7563 assert!(
7564 back.status.open(),
7565 "{status:?} is still alive; the question must still be waiting"
7566 );
7567 }
7568 }
7569
7570 #[test]
7571 fn settle_questions_never_touches_an_already_answered_question() {
7572 ask_test_home();
7573 let store = ask::Questions::open();
7574 let mut runner = runner_at(RunStatus::Failed);
7575 let run = runner.state.id.clone();
7576 let mut q = ask_open_question(&store, &run);
7577 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
7578 .unwrap();
7579 store.put(&mut q).unwrap();
7580
7581 runner.settle_questions();
7586 runner.settle_questions();
7587
7588 let back = store.get(&q.id).unwrap();
7589 assert_eq!(
7590 back.status,
7591 ask::QuestionStatus::Answered,
7592 "a real answer is a decision on record, never overwritten by a sweep"
7593 );
7594 }
7595
7596 #[tokio::test]
7607 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
7608 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
7609 let tmp = tempfile::tempdir().expect("tempdir");
7610 let repo = tmp.path().join("repo");
7611 std::fs::create_dir_all(&repo).unwrap();
7612 init_repo(&repo);
7613
7614 let mut config = Config::default();
7615 config.graph.worktree_root = Some(tmp.path().join("wt"));
7616
7617 let mut state = RunState::new(
7618 repo.clone(),
7619 "main".to_owned(),
7620 "deadbeef".to_owned(),
7621 "task".to_owned(),
7622 config,
7623 );
7624 let root = state.worktree_root();
7625 let wt_a = root.join("cand-A");
7626 let wt_b = root.join("cand-B");
7627 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
7628 .await
7629 .expect("worktree A");
7630 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
7631 .await
7632 .expect("worktree B");
7633
7634 state.candidates = vec![
7635 Candidate {
7636 index: 0,
7637 label: 'A',
7638 agent: "alpha".to_owned(),
7639 branch: "magi/x/A".to_owned(),
7640 worktree: wt_a.clone(),
7641 summary: String::new(),
7642 stat: String::new(),
7643 files: 0,
7644 commits: 0,
7645 empty: false,
7646 failed: None,
7647 verified_noop: None,
7648 duration_ms: 0,
7649 folded: false,
7650 },
7651 Candidate {
7652 index: 1,
7653 label: 'B',
7654 agent: "beta".to_owned(),
7655 branch: "magi/x/B".to_owned(),
7656 worktree: wt_b.clone(),
7657 summary: String::new(),
7658 stat: String::new(),
7659 files: 0,
7660 commits: 0,
7661 empty: false,
7662 failed: None,
7663 verified_noop: None,
7664 duration_ms: 0,
7665 folded: false,
7666 },
7667 ];
7668 state.tally = Some(Tally {
7669 first_choice: BTreeMap::from([('A', 1)]),
7670 borda: BTreeMap::new(),
7671 winner: 'A',
7672 rankings: 1,
7673 unanimous_initial: true,
7674 deliberated: false,
7675 changed_votes: 0,
7676 unanimous_final: true,
7677 tie_break: None,
7678 judges: 1,
7679 present: 1,
7680 quorum: 1,
7681 met_quorum: true,
7682 uncontested: None,
7683 });
7684 state.status = RunStatus::Ready;
7685
7686 fold_run(&mut state, false, &crate::run::home())
7687 .await
7688 .expect("fold_run");
7689
7690 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
7691 assert!(
7692 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7693 "the unmerged winner's branch survives"
7694 );
7695 assert!(
7696 !state.candidates[0].folded,
7697 "the winner is not marked folded"
7698 );
7699
7700 assert!(!wt_b.exists(), "the loser's worktree is removed");
7701 assert!(
7702 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
7703 "the loser's branch is removed"
7704 );
7705 assert!(state.candidates[1].folded, "the loser is marked folded");
7706 }
7707
7708 #[tokio::test]
7717 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
7718 let tmp = tempfile::tempdir().expect("tempdir");
7719 let repo = tmp.path().join("repo");
7720 std::fs::create_dir_all(&repo).unwrap();
7721 init_repo(&repo);
7722
7723 let mut config = Config::default();
7724 config.merge.mode = MergeMode::Local;
7725
7726 let mut state = RunState::new(
7727 repo.clone(),
7728 "main".to_owned(),
7729 "deadbeef".to_owned(),
7730 "task".to_owned(),
7731 config,
7732 );
7733 state.candidates = vec![Candidate {
7734 index: 0,
7735 label: 'A',
7736 agent: "alpha".to_owned(),
7737 branch: "does-not-exist".to_owned(),
7738 worktree: repo.clone(),
7739 summary: String::new(),
7740 stat: String::new(),
7741 files: 0,
7742 commits: 0,
7743 empty: false,
7744 failed: None,
7745 verified_noop: None,
7746 duration_ms: 0,
7747 folded: false,
7748 }];
7749 state.tally = Some(Tally {
7750 first_choice: BTreeMap::from([('A', 1)]),
7751 borda: BTreeMap::new(),
7752 winner: 'A',
7753 rankings: 1,
7754 unanimous_initial: true,
7755 deliberated: false,
7756 changed_votes: 0,
7757 unanimous_final: true,
7758 tie_break: None,
7759 judges: 0,
7760 present: 0,
7761 quorum: 0,
7762 met_quorum: true,
7763 uncontested: Some("only candidate A produced a change".to_owned()),
7764 });
7765 state.reviews = vec![ReviewRound {
7766 round: 1,
7767 head: "deadbeef".to_owned(),
7768 verified_head: None,
7769 verified_at: None,
7770 reviews: Vec::new(),
7771 e2e: Vec::new(),
7772 fix: None,
7773 blocking: 0,
7774 answered: 0,
7775 expected: 0,
7776 clean: true,
7777 verify_retried: false,
7778 e2e_deferred: false,
7779 e2e_defer_reason: None,
7780 progressed: false,
7781 vote_split: false,
7782 reconsideration: Vec::new(),
7783 verdict: None,
7784 }];
7785 state.gate = vec![CommandOutcome {
7786 command: "test".to_owned(),
7787 code: Some(0),
7788 output_tail: String::new(),
7789 duration_ms: 0,
7790 resource_blocked: false,
7791 }];
7792 state.gate_ran = true;
7793 state.status = RunStatus::Ready;
7798 state.merge = Some(MergeOutcome {
7799 mode: MergeMode::Local,
7800 ok: false,
7801 detail: "already concluded".to_owned(),
7802 });
7803
7804 let mut runner = Runner {
7805 state,
7806 roles: ResolvedRoles {
7807 implementers: Vec::new(),
7808 judges: Vec::new(),
7809 reviewers: Vec::new(),
7810 fixer: None,
7811 conductor: conductor(),
7812 implementer_roster: Vec::new(),
7813 },
7814 sem: Arc::new(Semaphore::new(1)),
7815 pause: Pause::new(),
7816 interrupt: Pause::new(),
7817 };
7818
7819 runner.merge().await.expect("merge");
7820
7821 assert_eq!(
7822 runner.state.status,
7823 RunStatus::Ready,
7824 "a concluded run's status must not change on reentry"
7825 );
7826 assert_eq!(
7827 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
7828 Some("already concluded"),
7829 "merge must not run again once the node already recorded an outcome"
7830 );
7831 }
7832
7833 #[tokio::test]
7842 async fn merge_refuses_a_gate_that_has_not_actually_run() {
7843 let tmp = tempfile::tempdir().expect("tempdir");
7844 let repo = tmp.path().join("repo");
7845 std::fs::create_dir_all(&repo).unwrap();
7846 init_repo(&repo);
7847
7848 let mut config = Config::default();
7849 config.merge.mode = MergeMode::Local;
7850
7851 let mut state = RunState::new(
7852 repo.clone(),
7853 "main".to_owned(),
7854 "deadbeef".to_owned(),
7855 "task".to_owned(),
7856 config,
7857 );
7858 state.candidates = vec![Candidate {
7859 index: 0,
7860 label: 'A',
7861 agent: "alpha".to_owned(),
7862 branch: "does-not-exist".to_owned(),
7863 worktree: repo.clone(),
7864 summary: String::new(),
7865 stat: String::new(),
7866 files: 0,
7867 commits: 0,
7868 empty: false,
7869 failed: None,
7870 verified_noop: None,
7871 duration_ms: 0,
7872 folded: false,
7873 }];
7874 state.tally = Some(Tally {
7875 first_choice: BTreeMap::from([('A', 1)]),
7876 borda: BTreeMap::new(),
7877 winner: 'A',
7878 rankings: 1,
7879 unanimous_initial: true,
7880 deliberated: false,
7881 changed_votes: 0,
7882 unanimous_final: true,
7883 tie_break: None,
7884 judges: 0,
7885 present: 0,
7886 quorum: 0,
7887 met_quorum: true,
7888 uncontested: Some("only candidate A produced a change".to_owned()),
7889 });
7890 state.reviews = vec![ReviewRound {
7891 round: 1,
7892 head: "deadbeef".to_owned(),
7893 verified_head: None,
7894 verified_at: None,
7895 reviews: Vec::new(),
7896 e2e: Vec::new(),
7897 fix: None,
7898 blocking: 0,
7899 answered: 0,
7900 expected: 0,
7901 clean: true,
7902 verify_retried: false,
7903 e2e_deferred: false,
7904 e2e_defer_reason: None,
7905 progressed: false,
7906 vote_split: false,
7907 reconsideration: Vec::new(),
7908 verdict: None,
7909 }];
7910 state.gate = Vec::new();
7912 state.gate_ran = false;
7913 state.status = RunStatus::Gating;
7914
7915 let mut runner = Runner {
7916 state,
7917 roles: ResolvedRoles {
7918 implementers: Vec::new(),
7919 judges: Vec::new(),
7920 reviewers: Vec::new(),
7921 fixer: None,
7922 conductor: conductor(),
7923 implementer_roster: Vec::new(),
7924 },
7925 sem: Arc::new(Semaphore::new(1)),
7926 pause: Pause::new(),
7927 interrupt: Pause::new(),
7928 };
7929
7930 runner.merge().await.expect("merge");
7931
7932 assert!(
7933 runner.state.merge.is_none(),
7934 "an empty gate must never be read as a passing one: {:?}",
7935 runner.state.merge
7936 );
7937 }
7938
7939 #[tokio::test]
7946 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
7947 let tmp = tempfile::tempdir().expect("tempdir");
7948 let repo = tmp.path().join("repo");
7949 std::fs::create_dir_all(&repo).unwrap();
7950 init_repo(&repo);
7951
7952 let config = Config::default();
7954
7955 let mut state = RunState::new(
7956 repo.clone(),
7957 "main".to_owned(),
7958 "deadbeef".to_owned(),
7959 "task".to_owned(),
7960 config,
7961 );
7962 state.candidates = vec![Candidate {
7963 index: 0,
7964 label: 'A',
7965 agent: "alpha".to_owned(),
7966 branch: "does-not-exist".to_owned(),
7967 worktree: repo.clone(),
7968 summary: String::new(),
7969 stat: String::new(),
7970 files: 0,
7971 commits: 0,
7972 empty: false,
7973 failed: None,
7974 verified_noop: None,
7975 duration_ms: 0,
7976 folded: false,
7977 }];
7978 state.tally = Some(Tally {
7979 first_choice: BTreeMap::from([('A', 1)]),
7980 borda: BTreeMap::new(),
7981 winner: 'A',
7982 rankings: 1,
7983 unanimous_initial: true,
7984 deliberated: false,
7985 changed_votes: 0,
7986 unanimous_final: true,
7987 tie_break: None,
7988 judges: 0,
7989 present: 0,
7990 quorum: 0,
7991 met_quorum: true,
7992 uncontested: Some("only candidate A produced a change".to_owned()),
7993 });
7994 state.reviews = vec![ReviewRound {
7995 round: 1,
7996 head: "deadbeef".to_owned(),
7997 verified_head: None,
7998 verified_at: None,
7999 reviews: Vec::new(),
8000 e2e: Vec::new(),
8001 fix: None,
8002 blocking: 0,
8003 answered: 0,
8004 expected: 0,
8005 clean: true,
8006 verify_retried: false,
8007 e2e_deferred: false,
8008 e2e_defer_reason: None,
8009 progressed: false,
8010 vote_split: false,
8011 reconsideration: Vec::new(),
8012 verdict: None,
8013 }];
8014
8015 let mut runner = Runner {
8016 state,
8017 roles: ResolvedRoles {
8018 implementers: Vec::new(),
8019 judges: Vec::new(),
8020 reviewers: Vec::new(),
8021 fixer: None,
8022 conductor: conductor(),
8023 implementer_roster: Vec::new(),
8024 },
8025 sem: Arc::new(Semaphore::new(1)),
8026 pause: Pause::new(),
8027 interrupt: Pause::new(),
8028 };
8029
8030 runner.gate().await.expect("gate");
8031 assert!(
8032 runner.state.gate_ran,
8033 "zero configured commands is still a real attempt, not an unrun gate"
8034 );
8035 assert!(runner.state.gate.is_empty());
8036 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
8037 assert_ne!(
8038 runner.state.status,
8039 RunStatus::Blocked,
8040 "a gate with nothing to check must not read as failed"
8041 );
8042
8043 runner.merge().await.expect("merge");
8044 assert_eq!(
8045 runner.state.status,
8046 RunStatus::Ready,
8047 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
8048 );
8049 }
8050
8051 #[tokio::test]
8062 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
8063 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8064 let home = crate::run::home();
8065
8066 let tmp = tempfile::tempdir().expect("tempdir");
8067 let repo = tmp.path().join("repo");
8068 std::fs::create_dir_all(&repo).unwrap();
8069 init_repo(&repo);
8070 let cache_dir = tmp.path().join("target");
8073
8074 let mut config = Config::default();
8075 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
8076 config.graph.timeout_verify = Some(2);
8079
8080 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8081 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8082 .expect("no io error acquiring directly")
8083 {
8084 crate::cache::AcquireOutcome::Acquired(g) => g,
8085 crate::cache::AcquireOutcome::Busy(b) => {
8086 panic!("expected the direct acquire to win the lease first: {b:?}")
8087 }
8088 };
8089
8090 let mut state = RunState::new(
8091 repo.clone(),
8092 "main".to_owned(),
8093 "deadbeef".to_owned(),
8094 "task".to_owned(),
8095 config,
8096 );
8097 state.candidates = vec![Candidate {
8098 index: 0,
8099 label: 'A',
8100 agent: "alpha".to_owned(),
8101 branch: "does-not-exist".to_owned(),
8102 worktree: repo.clone(),
8103 summary: String::new(),
8104 stat: String::new(),
8105 files: 0,
8106 commits: 0,
8107 empty: false,
8108 failed: None,
8109 verified_noop: None,
8110 duration_ms: 0,
8111 folded: false,
8112 }];
8113 state.tally = Some(Tally {
8114 first_choice: BTreeMap::from([('A', 1)]),
8115 borda: BTreeMap::new(),
8116 winner: 'A',
8117 rankings: 1,
8118 unanimous_initial: true,
8119 deliberated: false,
8120 changed_votes: 0,
8121 unanimous_final: true,
8122 tie_break: None,
8123 judges: 0,
8124 present: 0,
8125 quorum: 0,
8126 met_quorum: true,
8127 uncontested: Some("only candidate A produced a change".to_owned()),
8128 });
8129 state.reviews = vec![ReviewRound {
8130 round: 1,
8131 head: "deadbeef".to_owned(),
8132 verified_head: None,
8133 verified_at: None,
8134 reviews: Vec::new(),
8135 e2e: Vec::new(),
8136 fix: None,
8137 blocking: 0,
8138 answered: 0,
8139 expected: 0,
8140 clean: true,
8141 verify_retried: false,
8142 e2e_deferred: false,
8143 e2e_defer_reason: None,
8144 progressed: false,
8145 vote_split: false,
8146 reconsideration: Vec::new(),
8147 verdict: None,
8148 }];
8149
8150 let mut runner = Runner {
8151 state,
8152 roles: ResolvedRoles {
8153 implementers: Vec::new(),
8154 judges: Vec::new(),
8155 reviewers: Vec::new(),
8156 fixer: None,
8157 conductor: conductor(),
8158 implementer_roster: Vec::new(),
8159 },
8160 sem: Arc::new(Semaphore::new(1)),
8161 pause: Pause::new(),
8162 interrupt: Pause::new(),
8163 };
8164
8165 let started = std::time::Instant::now();
8166 runner.gate().await.expect("gate");
8167 assert!(
8168 started.elapsed() < Duration::from_secs(1),
8169 "a gate with nothing to run must never wait on a lease it never needed"
8170 );
8171 assert!(
8172 runner.state.gate_ran,
8173 "zero commands is still a real, immediate attempt"
8174 );
8175 assert!(runner.state.gate.is_empty());
8176 assert_ne!(
8177 runner.state.status,
8178 RunStatus::Blocked,
8179 "must not read as resource-blocked on a lease it never asked for"
8180 );
8181 }
8182
8183 #[tokio::test]
8193 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8194 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8195
8196 let tmp = tempfile::tempdir().expect("tempdir");
8197 let repo = tmp.path().join("repo");
8198 std::fs::create_dir_all(&repo).unwrap();
8199 init_repo(&repo);
8200
8201 let mut config = Config::default();
8202 config.verify.gate = vec![
8203 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8204 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8205 .to_owned(),
8206 ];
8207
8208 let mut state = RunState::new(
8209 repo.clone(),
8210 "main".to_owned(),
8211 "deadbeef".to_owned(),
8212 "task".to_owned(),
8213 config,
8214 );
8215 let run_id = state.id.clone();
8216 state.candidates = vec![Candidate {
8217 index: 0,
8218 label: 'A',
8219 agent: "alpha".to_owned(),
8220 branch: "does-not-exist".to_owned(),
8221 worktree: repo.clone(),
8222 summary: String::new(),
8223 stat: String::new(),
8224 files: 0,
8225 commits: 0,
8226 empty: false,
8227 failed: None,
8228 verified_noop: None,
8229 duration_ms: 0,
8230 folded: false,
8231 }];
8232 state.tally = Some(Tally {
8233 first_choice: BTreeMap::from([('A', 1)]),
8234 borda: BTreeMap::new(),
8235 winner: 'A',
8236 rankings: 1,
8237 unanimous_initial: true,
8238 deliberated: false,
8239 changed_votes: 0,
8240 unanimous_final: true,
8241 tie_break: None,
8242 judges: 0,
8243 present: 0,
8244 quorum: 0,
8245 met_quorum: true,
8246 uncontested: Some("only candidate A produced a change".to_owned()),
8247 });
8248 state.reviews = vec![ReviewRound {
8249 round: 1,
8250 head: "deadbeef".to_owned(),
8251 verified_head: None,
8252 verified_at: None,
8253 reviews: Vec::new(),
8254 e2e: Vec::new(),
8255 fix: None,
8256 blocking: 0,
8257 answered: 0,
8258 expected: 0,
8259 clean: true,
8260 verify_retried: false,
8261 e2e_deferred: false,
8262 e2e_defer_reason: None,
8263 progressed: false,
8264 vote_split: false,
8265 reconsideration: Vec::new(),
8266 verdict: None,
8267 }];
8268
8269 let mut runner = Runner {
8270 state,
8271 roles: ResolvedRoles {
8272 implementers: Vec::new(),
8273 judges: Vec::new(),
8274 reviewers: Vec::new(),
8275 fixer: None,
8276 conductor: conductor(),
8277 implementer_roster: Vec::new(),
8278 },
8279 sem: Arc::new(Semaphore::new(1)),
8280 pause: Pause::new(),
8281 interrupt: Pause::new(),
8282 };
8283
8284 let started_marker = repo.join("started.marker");
8285 let release_marker = repo.join("release.marker");
8286 let poller = tokio::spawn(async move {
8287 for _ in 0..100 {
8292 if started_marker.exists()
8293 && let Ok(s) = crate::run::RunState::load(&run_id)
8294 && let Some(a) = s.active.get("gate")
8295 {
8296 std::fs::write(&release_marker, b"go").expect("release marker");
8297 return Some(a.clone());
8298 }
8299 tokio::time::sleep(Duration::from_millis(50)).await;
8300 }
8301 None
8302 });
8303
8304 runner.gate().await.expect("gate");
8305 let captured = poller.await.expect("poller task");
8306 let captured = captured.expect(
8307 "the poller never saw a `gate` task entry in run.json while the command was \
8308 still blocked on its own release marker",
8309 );
8310
8311 assert_eq!(captured.task.as_deref(), Some("gate"));
8312 assert_eq!(captured.node, "gate");
8313 assert_eq!(captured.index, Some(1));
8314 assert_eq!(captured.total, Some(1));
8315 assert!(
8316 captured
8317 .command
8318 .as_deref()
8319 .is_some_and(|c| c.contains("started.marker")),
8320 "{captured:?}"
8321 );
8322
8323 assert!(
8324 runner.state.active.is_empty(),
8325 "the entry must be cleared once the command actually finished: {:?}",
8326 runner.state.active
8327 );
8328 assert!(runner.state.gate_ran);
8329 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8330 }
8331
8332 #[tokio::test]
8345 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8346 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8347 let home = crate::run::home();
8348
8349 let tmp = tempfile::tempdir().expect("tempdir");
8350 let repo = tmp.path().join("repo");
8351 std::fs::create_dir_all(&repo).unwrap();
8352 init_repo(&repo);
8353 let head = crate::git::rev_parse(&repo, "HEAD")
8354 .await
8355 .expect("rev-parse");
8356 let cache_dir = tmp.path().join("target");
8359
8360 let mut config = Config::default();
8361 config.verify.e2e = vec![format!(
8362 "CARGO_TARGET_DIR='{}' test -f README.md",
8363 cache_dir.display()
8364 )];
8365 config.graph.review_rounds = 1;
8366 config.graph.timeout_verify = Some(2);
8369
8370 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8371 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8372 .expect("no io error acquiring directly")
8373 {
8374 crate::cache::AcquireOutcome::Acquired(g) => g,
8375 crate::cache::AcquireOutcome::Busy(b) => {
8376 panic!("expected the direct acquire to win the lease first: {b:?}")
8377 }
8378 };
8379
8380 let mut state = RunState::new(
8381 repo.clone(),
8382 "main".to_owned(),
8383 head.clone(),
8384 "task".to_owned(),
8385 config,
8386 );
8387 state.candidates = vec![Candidate {
8388 index: 0,
8389 label: 'A',
8390 agent: "alpha".to_owned(),
8391 branch: "does-not-exist".to_owned(),
8392 worktree: repo.clone(),
8393 summary: String::new(),
8394 stat: String::new(),
8395 files: 0,
8396 commits: 0,
8397 empty: false,
8398 failed: None,
8399 verified_noop: None,
8400 duration_ms: 0,
8401 folded: false,
8402 }];
8403 state.tally = Some(Tally {
8404 first_choice: BTreeMap::from([('A', 1)]),
8405 borda: BTreeMap::new(),
8406 winner: 'A',
8407 rankings: 1,
8408 unanimous_initial: true,
8409 deliberated: false,
8410 changed_votes: 0,
8411 unanimous_final: true,
8412 tie_break: None,
8413 judges: 0,
8414 present: 0,
8415 quorum: 0,
8416 met_quorum: true,
8417 uncontested: Some("only candidate A produced a change".to_owned()),
8418 });
8419 state.reviews = vec![ReviewRound {
8423 round: 1,
8424 head: head.clone(),
8425 verified_head: None,
8426 verified_at: None,
8427 reviews: Vec::new(),
8428 e2e: Vec::new(),
8429 fix: None,
8430 blocking: 1,
8431 answered: 1,
8432 expected: 1,
8433 clean: false,
8434 verify_retried: false,
8435 e2e_deferred: true,
8436 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8437 progressed: false,
8438 vote_split: false,
8439 reconsideration: Vec::new(),
8440 verdict: None,
8441 }];
8442
8443 let mut runner = Runner {
8444 state,
8445 roles: ResolvedRoles {
8446 implementers: Vec::new(),
8447 judges: Vec::new(),
8448 reviewers: Vec::new(),
8449 fixer: None,
8450 conductor: conductor(),
8451 implementer_roster: Vec::new(),
8452 },
8453 sem: Arc::new(Semaphore::new(1)),
8454 pause: Pause::new(),
8455 interrupt: Pause::new(),
8456 };
8457
8458 let shell = runner.state.config.shell();
8459 runner
8460 .stop_reviewing("round budget spent", &shell, &repo)
8461 .await
8462 .expect("stop_reviewing");
8463
8464 let last = runner.state.reviews.last().expect("round record");
8465 assert_eq!(
8466 last.e2e_status(),
8467 E2eStatus::ResourceBlocked,
8468 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8469 failed: {last:?}"
8470 );
8471 assert_eq!(
8472 last.verified_head.as_deref(),
8473 Some(head.as_str()),
8474 "which commit this attempt targeted is known even though nothing finished checking \
8475 it"
8476 );
8477 let first_attempt_at = last
8478 .verified_at
8479 .expect("when this attempt ran is known too");
8480 assert_ne!(
8481 runner.state.status,
8482 RunStatus::Blocked,
8483 "contention is evidence about the machine, not the patch — it must not settle the \
8484 run as blocked: {:?}",
8485 runner.state.status
8486 );
8487 assert!(
8488 !runner
8489 .state
8490 .events
8491 .iter()
8492 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
8493 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
8494 runner.state.events
8495 );
8496
8497 runner
8502 .stop_reviewing("round budget spent", &shell, &repo)
8503 .await
8504 .expect("stop_reviewing retry");
8505 assert_eq!(
8506 runner.state.reviews.len(),
8507 1,
8508 "no new round was started: {:?}",
8509 runner.state.reviews
8510 );
8511 let last = runner.state.reviews.last().expect("round record");
8512 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
8513 assert!(
8514 last.verified_at.expect("still known") > first_attempt_at,
8515 "a second reentry must be a fresh attempt, not a stale copy of the first"
8516 );
8517 assert_ne!(runner.state.status, RunStatus::Blocked);
8518
8519 held.release();
8520 }
8521
8522 #[tokio::test]
8534 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
8535 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8536 let home = crate::run::home();
8537
8538 let tmp = tempfile::tempdir().expect("tempdir");
8539 let repo = tmp.path().join("repo");
8540 std::fs::create_dir_all(&repo).unwrap();
8541 init_repo(&repo);
8542 let head = crate::git::rev_parse(&repo, "HEAD")
8543 .await
8544 .expect("rev-parse");
8545 let cache_dir = tmp.path().join("target");
8546
8547 let mut config = Config::default();
8548 config.verify.e2e = vec![format!(
8549 "CARGO_TARGET_DIR='{}' test -f README.md",
8550 cache_dir.display()
8551 )];
8552 config.graph.review_rounds = 1;
8553 config.graph.timeout_verify = Some(2);
8554
8555 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8556 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8557 .expect("no io error acquiring directly")
8558 {
8559 crate::cache::AcquireOutcome::Acquired(g) => g,
8560 crate::cache::AcquireOutcome::Busy(b) => {
8561 panic!("expected the direct acquire to win the lease first: {b:?}")
8562 }
8563 };
8564
8565 let mut state = RunState::new(
8566 repo.clone(),
8567 "main".to_owned(),
8568 head.clone(),
8569 "task".to_owned(),
8570 config,
8571 );
8572 state.candidates = vec![Candidate {
8573 index: 0,
8574 label: 'A',
8575 agent: "alpha".to_owned(),
8576 branch: "does-not-exist".to_owned(),
8577 worktree: repo.clone(),
8578 summary: String::new(),
8579 stat: String::new(),
8580 files: 0,
8581 commits: 0,
8582 empty: false,
8583 failed: None,
8584 verified_noop: None,
8585 duration_ms: 0,
8586 folded: false,
8587 }];
8588 state.tally = Some(Tally {
8589 first_choice: BTreeMap::from([('A', 1)]),
8590 borda: BTreeMap::new(),
8591 winner: 'A',
8592 rankings: 1,
8593 unanimous_initial: true,
8594 deliberated: false,
8595 changed_votes: 0,
8596 unanimous_final: true,
8597 tie_break: None,
8598 judges: 0,
8599 present: 0,
8600 quorum: 0,
8601 met_quorum: true,
8602 uncontested: Some("only candidate A produced a change".to_owned()),
8603 });
8604 state.reviews = vec![ReviewRound {
8608 round: 1,
8609 head: head.clone(),
8610 verified_head: Some(head.clone()),
8611 verified_at: Some(jiff::Timestamp::now()),
8612 reviews: Vec::new(),
8613 e2e: vec![CommandOutcome {
8614 command: format!(
8615 "CARGO_TARGET_DIR='{}' test -f README.md",
8616 cache_dir.display()
8617 ),
8618 code: None,
8619 output_tail: "waiting for the shared build cache".to_owned(),
8620 duration_ms: 0,
8621 resource_blocked: true,
8622 }],
8623 fix: None,
8624 blocking: 1,
8625 answered: 1,
8626 expected: 1,
8627 clean: false,
8628 verify_retried: false,
8629 e2e_deferred: false,
8630 e2e_defer_reason: None,
8631 progressed: false,
8632 vote_split: false,
8633 reconsideration: Vec::new(),
8634 verdict: None,
8635 }];
8636
8637 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
8638 let mut runner = Runner {
8639 state,
8640 roles: ResolvedRoles {
8641 implementers: Vec::new(),
8642 judges: Vec::new(),
8643 reviewers: Vec::new(),
8644 fixer: None,
8645 conductor: conductor(),
8646 implementer_roster: Vec::new(),
8647 },
8648 sem: Arc::new(Semaphore::new(1)),
8649 pause: Pause::new(),
8650 interrupt: Pause::new(),
8651 };
8652
8653 runner.review_loop().await.expect("review_loop");
8658
8659 assert_eq!(
8660 runner.state.reviews.len(),
8661 1,
8662 "no new round was started on top of the unresolved one: {:?}",
8663 runner.state.reviews
8664 );
8665 let last = &runner.state.reviews[0];
8666 assert_eq!(
8667 last.e2e_status(),
8668 E2eStatus::ResourceBlocked,
8669 "still contended: {last:?}"
8670 );
8671 assert!(
8672 last.verified_at.expect("still known") > first_attempt_at,
8673 "review_loop must have actually retried the check, not left it exactly as found"
8674 );
8675 assert_ne!(
8676 runner.state.status,
8677 RunStatus::Blocked,
8678 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
8679 runner.state.status
8680 );
8681
8682 held.release();
8683 }
8684
8685 #[tokio::test]
8686 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
8687 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8688 let tmp = tempfile::tempdir().expect("tempdir");
8689 let repo = tmp.path().join("repo");
8690 std::fs::create_dir_all(&repo).unwrap();
8691 init_repo(&repo);
8692
8693 let mut config = Config::default();
8694 config.merge.mode = MergeMode::Pr;
8695 config.graph.land = true;
8696 config.graph.land_approval = false;
8697
8698 let mut state = RunState::new(
8699 repo.clone(),
8700 "main".to_owned(),
8701 "deadbeef".to_owned(),
8702 "task".to_owned(),
8703 config,
8704 );
8705 state.candidates = vec![Candidate {
8706 index: 0,
8707 label: 'A',
8708 agent: "alpha".to_owned(),
8709 branch: "does-not-exist".to_owned(),
8710 worktree: repo.clone(),
8711 summary: String::new(),
8712 stat: String::new(),
8713 files: 0,
8714 commits: 0,
8715 empty: false,
8716 failed: None,
8717 verified_noop: None,
8718 duration_ms: 0,
8719 folded: false,
8720 }];
8721 state.tally = Some(Tally {
8722 first_choice: BTreeMap::from([('A', 1)]),
8723 borda: BTreeMap::new(),
8724 winner: 'A',
8725 rankings: 1,
8726 unanimous_initial: true,
8727 deliberated: false,
8728 changed_votes: 0,
8729 unanimous_final: true,
8730 tie_break: None,
8731 judges: 0,
8732 present: 0,
8733 quorum: 0,
8734 met_quorum: true,
8735 uncontested: Some("only candidate A produced a change".to_owned()),
8736 });
8737 state.reviews = vec![ReviewRound {
8738 round: 1,
8739 head: "deadbeef".to_owned(),
8740 verified_head: None,
8741 verified_at: None,
8742 reviews: Vec::new(),
8743 e2e: Vec::new(),
8744 fix: None,
8745 blocking: 0,
8746 answered: 0,
8747 expected: 0,
8748 clean: true,
8749 verify_retried: false,
8750 e2e_deferred: false,
8751 e2e_defer_reason: None,
8752 progressed: false,
8753 vote_split: false,
8754 reconsideration: Vec::new(),
8755 verdict: None,
8756 }];
8757 state.gate = vec![CommandOutcome {
8758 command: "test".to_owned(),
8759 code: Some(0),
8760 output_tail: String::new(),
8761 duration_ms: 0,
8762 resource_blocked: false,
8763 }];
8764 state.gate_ran = true;
8765 state.status = RunStatus::Landing;
8769 state.merge = Some(MergeOutcome {
8770 mode: MergeMode::Pr,
8771 ok: true,
8772 detail: "https://example.invalid/x/y/pull/1".to_owned(),
8773 });
8774
8775 ask_test_home();
8779 let store = ask::Questions::open();
8780 let q = ask_open_question(&store, &state.id);
8781
8782 let mut runner = Runner {
8783 state,
8784 roles: ResolvedRoles {
8785 implementers: Vec::new(),
8786 judges: Vec::new(),
8787 reviewers: Vec::new(),
8788 fixer: None,
8789 conductor: conductor(),
8790 implementer_roster: Vec::new(),
8791 },
8792 sem: Arc::new(Semaphore::new(1)),
8793 pause: Pause::new(),
8794 interrupt: Pause::new(),
8795 };
8796
8797 runner.execute().await.expect("execute");
8802
8803 assert_eq!(
8804 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8805 Some("https://example.invalid/x/y/pull/1"),
8806 "reentry must not push again or open a second pull request over the \
8807 one `land` is already watching"
8808 );
8809 assert_ne!(
8810 runner.state.status,
8811 RunStatus::Landing,
8812 "land could not actually reach the fake pull request, so it must \
8813 have given up rather than left the run silently parked forever"
8814 );
8815 assert_eq!(runner.state.status, RunStatus::Blocked);
8819 assert!(
8820 store.get(&q.id).unwrap().status.open(),
8821 "Blocked is still alive; settle_questions must have been a no-op here"
8822 );
8823 }
8824
8825 fn state_with_round(round: ReviewRound) -> RunState {
8826 let mut s = RunState::new(
8827 PathBuf::from("/repo"),
8828 "main".to_owned(),
8829 "abc1234".to_owned(),
8830 "add retries".to_owned(),
8831 Config::default(),
8832 );
8833 s.reviews = vec![round];
8834 s
8835 }
8836
8837 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
8838 crate::verdict::Finding {
8839 id: id.to_owned(),
8840 severity,
8841 file: None,
8842 line: None,
8843 title: title.to_owned(),
8844 detail: String::new(),
8845 }
8846 }
8847
8848 #[test]
8849 fn pr_body_names_open_findings_and_declined_ones() {
8850 let round = ReviewRound {
8851 round: 2,
8852 head: "deadbee".to_owned(),
8853 verified_head: None,
8854 verified_at: None,
8855 reviews: vec![ReviewRecord {
8856 attempts: 0,
8857 reviewer: 1,
8858 agent: "alpha".to_owned(),
8859 summary: String::new(),
8860 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
8861 vote: None,
8862 failed: None,
8863 duration_ms: 0,
8864 }],
8865 e2e: vec![CommandOutcome {
8866 command: "cargo test".to_owned(),
8867 code: Some(0),
8868 output_tail: String::new(),
8869 duration_ms: 0,
8870 resource_blocked: false,
8871 }],
8872 verify_retried: false,
8873 e2e_deferred: false,
8874 e2e_defer_reason: None,
8875 fix: Some(FixRecord {
8876 agent: "alpha".to_owned(),
8877 addressed: Vec::new(),
8878 rejected: vec![crate::verdict::Rejection {
8879 id: "R1-1-1".to_owned(),
8880 why: "not reachable from any caller".to_owned(),
8881 }],
8882 notes: String::new(),
8883 committed: true,
8884 failed: None,
8885 duration_ms: 0,
8886 continuation: None,
8887 }),
8888 blocking: 0,
8889 answered: 1,
8890 expected: 1,
8891 clean: false,
8892 progressed: true,
8893 vote_split: false,
8894 reconsideration: Vec::new(),
8895 verdict: None,
8896 };
8897 let state = state_with_round(round);
8898 let body = pr_message(&state, 'A').body;
8899
8900 assert!(body.contains("add retries"), "the task must still be there");
8901 assert!(body.contains("R2-1-1"), "{body}");
8902 assert!(body.contains("unused import"), "{body}");
8903 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
8904 assert!(
8905 body.contains("not reachable from any caller"),
8906 "the reason it was declined: {body}"
8907 );
8908 }
8909
8910 #[test]
8911 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
8912 let round = ReviewRound {
8913 round: 1,
8914 head: "deadbee".to_owned(),
8915 verified_head: None,
8916 verified_at: None,
8917 reviews: vec![ReviewRecord {
8918 attempts: 0,
8919 reviewer: 1,
8920 agent: "alpha".to_owned(),
8921 summary: String::new(),
8922 findings: Vec::new(),
8923 vote: None,
8924 failed: None,
8925 duration_ms: 0,
8926 }],
8927 e2e: Vec::new(),
8928 verify_retried: false,
8929 e2e_deferred: false,
8930 e2e_defer_reason: None,
8931 fix: None,
8932 blocking: 0,
8933 answered: 1,
8934 expected: 1,
8935 clean: true,
8936 progressed: false,
8937 vote_split: false,
8938 reconsideration: Vec::new(),
8939 verdict: None,
8940 };
8941 let state = state_with_round(round);
8942 let body = pr_message(&state, 'A').body;
8943 assert!(!body.contains("Open review findings"), "{body}");
8944 assert!(!body.contains("Declined"), "{body}");
8945 }
8946
8947 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
8948 let mut state = RunState::new(
8949 PathBuf::from("/repo"),
8950 "main".to_owned(),
8951 "abc1234".to_owned(),
8952 instruction.to_owned(),
8953 Config::default(),
8954 );
8955 state.candidates.push(Candidate {
8956 index: 0,
8957 label: 'A',
8958 agent: "alpha".to_owned(),
8959 branch: "magi/x/A".to_owned(),
8960 worktree: PathBuf::from("/wt"),
8961 summary: summary.to_owned(),
8962 stat: String::new(),
8963 files: 1,
8964 commits: 1,
8965 empty: false,
8966 failed: None,
8967 verified_noop: None,
8968 folded: false,
8969 duration_ms: 0,
8970 });
8971 state
8972 }
8973
8974 #[test]
8975 fn pr_message_describes_the_change_not_the_task() {
8976 let state = state_with_summary(
8977 "今回やってほしいこと: results projector を直す",
8978 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
8979 );
8980 let m = pr_message(&state, 'A');
8981 assert_eq!(m.title, "fix(web): batch the runs list reads");
8982 assert!(
8983 m.body.starts_with("## Summary\n\n- reads run.json once"),
8984 "{}",
8985 m.body
8986 );
8987 assert!(!m.body.contains("TITLE:"), "{}", m.body);
8988 let task_at = m.body.find("今回やってほしいこと").unwrap();
8989 let details_at = m.body.find("<details>").unwrap();
8990 assert!(
8991 details_at < task_at,
8992 "the task lives inside <details>: {}",
8993 m.body
8994 );
8995 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
8996 assert!(m.body.contains("magi:candidate-a"));
8997 }
8998
8999 #[test]
9000 fn pr_message_falls_back_to_the_task_without_a_title_line() {
9001 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
9002 let m = pr_message(&state, 'A');
9003 assert_eq!(m.title, "add retries");
9004 assert!(
9005 m.body.contains("## Summary\n\n- did some things"),
9006 "{}",
9007 m.body
9008 );
9009
9010 let none = RunState::new(
9011 PathBuf::from("/repo"),
9012 "main".to_owned(),
9013 "abc1234".to_owned(),
9014 "add retries".to_owned(),
9015 Config::default(),
9016 );
9017 let m = pr_message(&none, 'A');
9018 assert_eq!(m.title, "add retries");
9019 assert!(!m.body.contains("## Summary"), "{}", m.body);
9020 }
9021
9022 #[test]
9023 fn pr_message_refuses_the_candidate_commit_subject() {
9024 for bad in [
9025 "TITLE: magi: candidate A (uncommitted work)",
9026 "TITLE: chore: stuff (uncommitted work)",
9027 "TITLE: ",
9028 ] {
9029 let state = state_with_summary("add retries", bad);
9030 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
9031 }
9032 }
9033
9034 #[test]
9035 fn pr_message_bounds_a_very_long_task_and_title() {
9036 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
9037 let state = state_with_summary(&long, "- nothing");
9038 let m = pr_message(&state, 'A');
9039 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9040 assert!(!m.title.contains('\n'));
9041
9042 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
9043 let m = pr_message(&state, 'A');
9044 assert!(m.title.starts_with("feat: "));
9045 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9046 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
9047 }
9048
9049 #[test]
9050 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
9051 let mut state = state_with_summary(
9055 "add retries",
9056 "TITLE: fix(web): batch reads\n- reads run.json once",
9057 );
9058 state.config.graph.language = "ja".to_owned();
9059 let m = pr_message(&state, 'A');
9060 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
9061
9062 let task = "今回やってほしいこと: results projector を直す";
9065 let mut state = state_with_summary(task, "- no title line");
9066 state.config.graph.language = "ja".to_owned();
9067 let m = pr_message(&state, 'A');
9068 assert_eq!(
9069 m.title,
9070 format!("chore: land candidate A of run {}", state.id)
9071 );
9072 assert!(
9073 m.body.contains(&format!(
9074 "<summary>Original task</summary>\n\n{task}\n\n</details>"
9075 )),
9076 "{}",
9077 m.body
9078 );
9079 }
9080
9081 #[test]
9082 fn pr_message_scrubs_home_paths_and_addresses() {
9083 let state = state_with_summary(
9084 "fix it in /Users/someone/src/x",
9085 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9086 );
9087 let m = pr_message(&state, 'A');
9088 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9089 assert!(!m.body.contains(leak), "{}", m.body);
9090 }
9091 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9092 }
9093
9094 #[test]
9095 fn pr_message_survives_a_task_that_closes_details() {
9096 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9097 let m = pr_message(&state, 'A');
9098 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9099 }
9100
9101 #[test]
9102 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9103 let cmd = manual_merge_command(
9104 MergeStyle::Squash,
9105 Path::new("/repo"),
9106 "b",
9107 "fix: \"quoted\" $(x) `y`\n\nbody",
9108 );
9109 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9110 }
9111
9112 #[test]
9113 fn manual_merge_command_matches_the_configured_style() {
9114 let repo = Path::new("/repo");
9115 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9116
9117 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9118 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9119
9120 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9121 assert_eq!(
9122 squash,
9123 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9124 \"Merge magi run 0832 (candidate A)\""
9125 );
9126
9127 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9128 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9129 }
9130
9131 #[test]
9132 fn a_nudge_gets_a_quarter_of_the_budget() {
9133 assert_eq!(retry_budget(secs(1200), true), secs(300));
9135 assert_eq!(retry_budget(secs(3600), true), secs(900));
9136 }
9137
9138 #[test]
9139 fn a_resent_prompt_keeps_the_whole_budget() {
9140 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9143 assert_eq!(retry_budget(secs(60), false), secs(60));
9144 }
9145
9146 #[test]
9147 fn the_floor_never_exceeds_the_original_budget() {
9148 assert_eq!(retry_budget(secs(60), true), secs(60));
9152 assert_eq!(retry_budget(secs(480), true), secs(120));
9153 assert_eq!(retry_budget(secs(0), true), secs(0));
9154 }
9155
9156 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9157 agent::CommandEvidence {
9158 id: "item1".to_owned(),
9159 description: "cargo test".to_owned(),
9160 exit_code,
9161 result_summary: String::new(),
9162 source: "codex".to_owned(),
9163 }
9164 }
9165
9166 #[test]
9167 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9168 assert!(!has_unconfirmed_command(&[]));
9172 }
9173
9174 #[test]
9175 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9176 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9180 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9181 assert!(!has_unconfirmed_command(&[
9182 evidence(Some(0)),
9183 evidence(Some(101))
9184 ]));
9185 }
9186
9187 #[test]
9188 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9189 assert!(has_unconfirmed_command(&[
9190 evidence(Some(0)),
9191 evidence(None)
9192 ]));
9193 }
9194
9195 #[test]
9196 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9197 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9198 assert_eq!(
9199 verified_noop_claim(true, &[], text).as_deref(),
9200 Some("already fixed by b32cfc4, on main.")
9201 );
9202 }
9203
9204 #[test]
9205 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9206 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9209 assert!(verified_noop_claim(false, &[], text).is_none());
9210 }
9211
9212 #[test]
9213 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9214 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9215 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9216 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9218 }
9219
9220 #[test]
9221 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9222 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9223 }
9224
9225 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9228 runner.state.candidates = shape
9229 .iter()
9230 .enumerate()
9231 .map(|(i, &(empty, verified))| Candidate {
9232 index: i,
9233 label: (b'A' + i as u8) as char,
9234 agent: "sonnet".to_owned(),
9235 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9236 worktree: PathBuf::from(format!("/wt/{i}")),
9237 summary: String::new(),
9238 stat: String::new(),
9239 files: 0,
9240 commits: 0,
9241 empty,
9242 failed: None,
9243 verified_noop: verified.map(str::to_owned),
9244 duration_ms: 0,
9245 folded: false,
9246 })
9247 .collect();
9248 }
9249
9250 #[test]
9251 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9252 ask_test_home();
9253 let mut runner = runner_at(RunStatus::Implementing);
9254 set_candidates(
9255 &mut runner,
9256 &[
9257 (true, Some("already on main at b32cfc4")),
9258 (true, Some("same fix, see the existing test")),
9259 ],
9260 );
9261
9262 runner
9263 .after_implement()
9264 .expect("a verified no-op is not an error");
9265
9266 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9267 }
9268
9269 #[test]
9270 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9271 ask_test_home();
9272 let mut runner = runner_at(RunStatus::Implementing);
9273 set_candidates(
9277 &mut runner,
9278 &[(true, Some("already on main at b32cfc4")), (true, None)],
9279 );
9280
9281 let err = runner
9282 .after_implement()
9283 .expect_err("an unverified empty candidate must still fail the run");
9284
9285 assert!(
9286 err.to_string().contains("no candidate produced a change"),
9287 "{err}"
9288 );
9289 assert_eq!(runner.state.status, RunStatus::Failed);
9290 }
9291
9292 #[test]
9293 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9294 ask_test_home();
9295 let mut runner = runner_at(RunStatus::Implementing);
9296 set_candidates(&mut runner, &[(true, None), (true, None)]);
9297
9298 let err = runner
9299 .after_implement()
9300 .expect_err("no candidate declared anything; this is an ordinary failure");
9301
9302 assert!(
9303 err.to_string().contains("no candidate produced a change"),
9304 "{err}"
9305 );
9306 assert_eq!(runner.state.status, RunStatus::Failed);
9307 }
9308}