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::refs;
46use crate::run::{
47 BaseSync, Candidate, CommandOutcome, ContinuationOutcome, ContinuationRecord,
48 DeliberationRound, DeliberationTurn, E2eStatus, FixRecord, GateFixRecord, JobRecord, JobStatus,
49 Judgement, MergeOutcome, OperatorFixFinding, OperatorFixOutcome, OperatorFixRequest, Origin,
50 QuotaLoss, ReviewRecord, ReviewRevoteRecord, ReviewRound, RunState, RunStatus, Tally,
51 VoteRecord, tail, write_artifact,
52};
53use crate::verdict::{
54 self, FinalVote, Finding, FixReport, Position, Proposal, Ranking, Review, ReviewRevote,
55 ReviewVote, Severity,
56};
57
58const OUTPUT_TAIL: usize = 8_000;
60
61const EVENT_OUTPUT_TAIL: usize = 2_000;
64
65const LEASE_RELEASE_POLL: Duration = Duration::from_secs(1);
68
69const LEASE_RELEASE_MAX_WAIT: Duration = Duration::from_secs(30);
86
87pub(crate) const STAGNANT_LIMIT: usize = 2;
101
102const BASE_SYNC_ROUNDS: usize = 4;
115
116const MAX_FIX_CONTINUATIONS: usize = 2;
132
133#[derive(Clone)]
139struct SeatJob {
140 spec: AgentSpec,
141 seat: SeatState,
142 cwd: PathBuf,
143 prompt: String,
144 timeout: Duration,
145 allow_write: bool,
146 sessions: bool,
147 artifacts: PathBuf,
148 stem: String,
149}
150
151enum AgentOutcome {
163 Ok(AgentOutput),
165 Quota(AgentOutput),
167 Dropped(AgentOutput),
170 Failed(String),
172}
173
174#[derive(Debug, Clone, Default)]
201pub struct Pause(Arc<AtomicBool>, Arc<Mutex<Option<String>>>);
202
203impl Pause {
204 #[must_use]
206 pub fn new() -> Self {
207 Self::default()
208 }
209
210 pub fn park(&self) {
212 self.0.store(true, Ordering::SeqCst);
213 }
214
215 pub fn park_because(&self, reason: impl Into<String>) {
221 let mut reason_guard = self
222 .1
223 .lock()
224 .unwrap_or_else(std::sync::PoisonError::into_inner);
225 if reason_guard.is_none() {
226 *reason_guard = Some(reason.into());
227 }
228 drop(reason_guard);
229 self.park();
230 }
231
232 #[must_use]
234 pub fn parked(&self) -> bool {
235 self.0.load(Ordering::SeqCst)
236 }
237
238 #[must_use]
240 pub fn reason(&self) -> Option<String> {
241 self.1
242 .lock()
243 .unwrap_or_else(std::sync::PoisonError::into_inner)
244 .clone()
245 }
246}
247
248pub struct Runner {
250 pub state: RunState,
252 roles: ResolvedRoles,
253 sem: Arc<Semaphore>,
254 pause: Pause,
258 interrupt: Pause,
264}
265
266async fn sync_review_branch(repo: &Path, branch: &str, remote: &str, base: &str) -> Result<()> {
295 let tracking = format!("{remote}/{branch}");
296 let fetched = git::fetch(repo, remote, branch).await;
297 let fresh = matches!(&fetched, Ok(o) if o.ok()) && git::rev_exists(repo, &tracking).await;
298 let local_exists = git::branch_exists(repo, branch).await?;
299 if !fresh {
300 if !local_exists {
301 bail!("no branch `{branch}` in {} or on {remote}", repo.display());
302 }
303 tracing::warn!(
304 "could not read {tracking}; reviewing the local `{branch}`, which may be stale"
305 );
306 return Ok(());
307 }
308 let remote_sha = git::rev_parse(repo, &tracking).await?;
309 if !local_exists {
310 git::git(repo, &["branch", branch, &tracking]).await?;
311 return Ok(());
312 }
313 let local_sha = git::rev_parse(repo, &format!("refs/heads/{branch}")).await?;
314 if local_sha == remote_sha || git::is_ancestor(repo, &remote_sha, &local_sha).await {
315 return Ok(());
316 }
317 if !git::is_ancestor(repo, &local_sha, &remote_sha).await {
318 match crate::reconcile::reconcile(repo, remote, branch, &local_sha, &remote_sha, base)
324 .await?
325 {
326 crate::reconcile::Reconciliation::Pushed => {
327 tracing::warn!(
328 "local `{branch}` ({}) is {tracking} ({}) rebased; pushed it over",
329 short(&local_sha),
330 short(&remote_sha)
331 );
332 return Ok(());
333 }
334 crate::reconcile::Reconciliation::Placeholder => {}
335 crate::reconcile::Reconciliation::Genuine(d) => return Err((*d).into()),
336 }
337 }
338 let out = git::git_raw(repo, &["branch", "-f", branch, &tracking]).await?;
339 if !out.ok() {
340 bail!(
341 "local `{branch}` ({}) is stale against {tracking} ({}) but git will not move it: {}",
342 short(&local_sha),
343 short(&remote_sha),
344 out.stderr
345 );
346 }
347 tracing::warn!(
348 "local `{branch}` was stale: fast-forwarded {} -> {}",
349 short(&local_sha),
350 short(&remote_sha)
351 );
352 Ok(())
353}
354
355async fn resolve_base(repo: &Path, base_branch: &str, remote: &str) -> Result<String> {
356 let tracking = format!("{remote}/{base_branch}");
357 let fetched = git::fetch(repo, remote, base_branch).await;
358 if let Ok(out) = &fetched
359 && out.ok()
360 && git::rev_exists(repo, &tracking).await
361 {
362 return git::rev_parse(repo, &tracking).await;
363 }
364 let why = match &fetched {
365 Ok(out) if !out.ok() => out.stderr.lines().next().unwrap_or("").to_owned(),
366 Ok(_) => format!("{remote} has no {base_branch}"),
367 Err(e) => e.to_string(),
368 };
369 tracing::warn!(
370 "could not read {tracking} ({why}); branching off the local \
371 {base_branch} instead, which may be behind"
372 );
373 git::rev_parse(repo, base_branch).await.with_context(|| {
374 format!(
375 "cannot resolve `{base_branch}`; set [merge] base in magi.toml to a \
376 branch that exists"
377 )
378 })
379}
380
381struct FixClaim {
398 path: PathBuf,
399}
400
401impl FixClaim {
402 fn acquire(dir: &Path) -> Result<Self> {
403 std::fs::create_dir_all(dir).with_context(|| format!("create {}", dir.display()))?;
404 let path = dir.join("fix.lock");
405 match Self::create(&path) {
406 Ok(claim) => Ok(claim),
407 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
408 if Self::reclaim_if_dead(&path) {
409 Self::create(&path).with_context(|| format!("lock {}", path.display()))
410 } else {
411 bail!(
412 "another `magi fix` is already running for this run ({} exists)",
413 path.display()
414 )
415 }
416 }
417 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
418 }
419 }
420
421 fn create(path: &Path) -> std::io::Result<Self> {
422 let mut f = std::fs::OpenOptions::new()
423 .write(true)
424 .create_new(true)
425 .open(path)?;
426 use std::io::Write as _;
427 writeln!(f, "{}", std::process::id())?;
429 Ok(Self {
430 path: path.to_owned(),
431 })
432 }
433
434 fn reclaim_if_dead(path: &Path) -> bool {
438 let dead = std::fs::read_to_string(path)
439 .ok()
440 .and_then(|body| body.trim().parse::<u32>().ok())
441 .is_some_and(|pid| !crate::proc::pid_alive(pid));
442 dead && std::fs::remove_file(path).is_ok()
443 }
444}
445
446impl Drop for FixClaim {
447 fn drop(&mut self) {
448 let _ = std::fs::remove_file(&self.path);
449 }
450}
451
452impl Runner {
453 pub async fn start(
455 repo: &Path,
456 instruction: String,
457 config: Config,
458 origin: Origin,
459 ) -> Result<Self> {
460 Self::start_naming(repo, instruction, "", config, origin).await
461 }
462
463 pub async fn start_naming(
467 repo: &Path,
468 instruction: String,
469 also_scan: &str,
470 config: Config,
471 origin: Origin,
472 ) -> Result<Self> {
473 let repo = git::toplevel(repo).await?;
474 let missing = agent::missing_programs(&config.agents);
475 if !missing.is_empty() {
476 bail!(
477 "these agent programs are not on PATH: {}. Fix the roster in \
478 magi.toml or install them.",
479 missing.join(", ")
480 );
481 }
482 let base_branch = match config.merge.base.clone() {
483 Some(b) => b,
484 None => git::current_branch(&repo)
485 .await?
486 .context("HEAD is detached; set [merge] base in magi.toml")?,
487 };
488 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
489 if !git::is_clean(&repo).await? {
493 tracing::warn!(
494 "{} has uncommitted changes; they are not part of this run, \
495 which branches off {base_branch} ({})",
496 repo.display(),
497 &base_commit[..base_commit.len().min(8)]
498 );
499 }
500 let roles = config.resolve_roles()?;
501 let max_parallel = config.graph.max_parallel.max(1);
502 let seeds = refs::resolve(
505 &repo,
506 &base_commit,
507 &config.merge.remote,
508 &format!("{also_scan}\n{instruction}"),
509 )
510 .await;
511 refs::plan(&repo, &seeds).await?;
512 let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
513 state.origin = Some(origin);
517 for seed in &seeds {
518 state.event(
519 "seed",
520 refs::describe(std::slice::from_ref(seed)).unwrap_or_default(),
521 );
522 }
523 state.seeds = seeds;
524 state.event("start", format!("run {} created", state.id));
525 state.save()?;
526 Ok(Self {
527 state,
528 roles,
529 sem: Arc::new(Semaphore::new(max_parallel)),
530 pause: Pause::new(),
531 interrupt: Pause::new(),
532 })
533 }
534
535 pub async fn review(repo: &Path, branch: &str, config: Config, origin: Origin) -> Result<Self> {
549 Self::review_taking_over(repo, branch, config, None, origin).await
550 }
551
552 pub async fn review_taking_over(
559 repo: &Path,
560 branch: &str,
561 config: Config,
562 takeover: Option<crate::handover::Takeover>,
563 origin: Origin,
564 ) -> Result<Self> {
565 let repo = git::toplevel(repo).await?;
566 let missing = agent::missing_programs(&config.agents);
567 if !missing.is_empty() {
568 bail!(
569 "these agent programs are not on PATH: {}. Fix the roster in \
570 magi.toml or install them.",
571 missing.join(", ")
572 );
573 }
574 let base_branch = match config.merge.base.clone() {
575 Some(b) => b,
576 None => git::current_branch(&repo)
577 .await?
578 .context("HEAD is detached; set [merge] base in magi.toml")?,
579 };
580 if base_branch == branch {
581 bail!("`{branch}` is the base branch; there is nothing to review against");
582 }
583 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
584
585 let roles = config.resolve_roles()?;
586 let max_parallel = config.graph.max_parallel.max(1);
587 let mut state = RunState::new(
588 repo.clone(),
589 base_branch,
590 base_commit.clone(),
591 String::new(),
592 config,
593 );
594 state.origin = Some(origin);
595
596 let takeover = takeover.unwrap_or_else(|| crate::handover::Takeover {
603 earlier: Vec::new(),
604 home: crate::run::home(),
605 choice: None,
606 });
607 let released = crate::handover::release(&repo, branch, &state.id, &takeover).await?;
608 if let Some(released) = &released {
609 state.event(
610 "release",
611 format!(
612 "took `{branch}` over from run {}: its worktree was released: {}",
613 crate::run::short_of(&released.old_id),
614 released.audit
615 ),
616 );
617 }
618 if let Some(choice) = takeover.choice.as_ref()
622 && let Err(e) =
623 crate::reconcile::apply_choice(&repo, &state.config.merge.remote, branch, choice)
624 .await
625 {
626 if let Some(released) = &released {
627 released.restore(&repo, branch).await;
628 }
629 return Err(e.context("applying the owner's answer about the diverged branch"));
630 }
631 let opened =
632 Self::open_review(&repo, branch, state, roles, max_parallel, base_commit).await;
633 if opened.is_err()
634 && let Some(released) = &released
635 {
636 released.restore(&repo, branch).await;
637 }
638 opened
639 }
640
641 async fn open_review(
644 repo: &Path,
645 branch: &str,
646 mut state: RunState,
647 roles: ResolvedRoles,
648 max_parallel: usize,
649 base_commit: String,
650 ) -> Result<Self> {
651 sync_review_branch(repo, branch, &state.config.merge.remote, &base_commit).await?;
652 let log = git::log_oneline(repo, &base_commit, branch)
655 .await
656 .unwrap_or_default();
657 let instruction = format!(
658 "Review the work already on branch `{branch}`. There is no task \
659 statement: what the change claims to do is whatever its commits \
660 say.\n\n{}",
661 if log.trim().is_empty() {
662 "(no commit messages)"
663 } else {
664 log.trim()
665 }
666 );
667 state.instruction = instruction;
668
669 let worktree = state.worktree_root().join("under-review");
672 if let Some(parent) = worktree.parent() {
673 tokio::fs::create_dir_all(parent).await.ok();
674 }
675 let path = worktree.to_string_lossy().to_string();
676 git::git(repo, &["worktree", "add", &path, branch])
677 .await
678 .with_context(|| {
679 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
680 })?;
681
682 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
683 .await
684 .unwrap_or(0);
685 if commits == 0 {
686 git::worktree_remove(repo, &worktree).await.ok();
687 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
688 }
689 let files = git::changed_files(&worktree, &base_commit, "HEAD")
690 .await
691 .map(|f| f.len())
692 .unwrap_or(0);
693 if files == 0
694 && let (Ok(head_tree), Ok(base_tree)) = (
695 git::tree_of(&worktree, "HEAD").await,
696 git::tree_of(&worktree, &base_commit).await,
697 )
698 && head_tree == base_tree
699 {
700 let head = git::rev_parse(&worktree, "HEAD").await.unwrap_or_default();
701 git::worktree_remove(repo, &worktree).await.ok();
702 bail!(
703 "`{branch}` at {} has a tree identical to base {}; this usually means \
704 the branch ref is stale (check `git rev-parse refs/heads/{branch}` \
705 against `{}/{branch}`) rather than an empty change",
706 short(&head),
707 short(&base_commit),
708 state.config.merge.remote
709 );
710 }
711 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
712 .await
713 .unwrap_or_default();
714
715 state.candidates.push(Candidate {
716 index: 0,
717 label: 'A',
718 agent: "(existing branch)".to_owned(),
721 branch: branch.to_owned(),
722 worktree,
723 summary: String::new(),
724 stat,
725 files,
726 commits,
727 empty: false,
728 failed: None,
729 verified_noop: None,
730 duration_ms: 0,
731 folded: false,
732 });
733 state.tally = Some(Tally {
734 first_choice: BTreeMap::from([('A', 0)]),
735 borda: BTreeMap::new(),
736 winner: 'A',
737 rankings: 0,
738 unanimous_initial: false,
739 deliberated: false,
740 changed_votes: 0,
741 unanimous_final: false,
742 tie_break: None,
743 judges: 0,
747 present: 0,
748 quorum: 0,
749 met_quorum: true,
750 uncontested: Some("review-only run: nothing competed".to_owned()),
751 });
752 state.status = RunStatus::Reviewing;
753 state.event(
754 "start",
755 format!(
756 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
757 state.id
758 ),
759 );
760 state.save()?;
761 Ok(Self {
762 state,
763 roles,
764 sem: Arc::new(Semaphore::new(max_parallel)),
765 pause: Pause::new(),
766 interrupt: Pause::new(),
767 })
768 }
769
770 pub fn resume(id: &str) -> Result<Self> {
772 let state = RunState::load(id)?;
773 if let Some(to) = &state.released_to {
774 bail!(
775 "run {} cannot be resumed: its worktree was released to run {}",
776 state.short(),
777 crate::run::short_of(to)
778 );
779 }
780 let roles = state.config.resolve_roles()?;
781 let max_parallel = state.config.graph.max_parallel.max(1);
782 Ok(Self {
783 state,
784 roles,
785 sem: Arc::new(Semaphore::new(max_parallel)),
786 pause: Pause::new(),
787 interrupt: Pause::new(),
788 })
789 }
790
791 pub async fn execute(&mut self) -> Result<()> {
798 let result = self.execute_graph().await;
799 self.mark_driver_exited();
800 let ended = if result.is_err() {
801 Some(crate::notices::run_stopped(&self.state.id, &self.state))
802 } else {
803 crate::notices::run_ended(&self.state)
804 };
805 if let Some(notice) = ended {
806 crate::notices::raise(notice);
807 }
808 result
809 }
810
811 fn mark_driver_exited(&mut self) {
820 self.state.driver_exited = true;
821 let pid = std::process::id();
822 let Ok(mut disk) = RunState::load(&self.state.id) else {
823 return;
824 };
825 if disk.released_to.is_some() || disk.driver_pid != Some(pid) || disk.driver_exited {
826 return;
827 }
828 disk.driver_exited = true;
829 if let Err(e) = disk.save() {
830 tracing::warn!("could not record that run {} stopped: {e:#}", self.state.id);
831 }
832 }
833
834 async fn execute_graph(&mut self) -> Result<()> {
835 self.state.parked = false;
840 self.state.clear_active();
847 if let Ok(disk) = RunState::load(&self.state.id)
865 && let Some(to) = &disk.released_to
866 {
867 bail!(
868 "run {} cannot continue: its worktree was released to run {}",
869 self.state.short(),
870 crate::run::short_of(to)
871 );
872 }
873 let pid = std::process::id();
874 self.state.driver_pid = Some(pid);
875 self.state.driver_started_at = crate::proc::process_started_at(pid);
876 self.state.driver_exited = false;
877 self.state.save()?;
878 if self.state.status == RunStatus::Stalled {
891 if self.recover_stall().await? {
892 self.finish_after_tally().await?;
893 } else {
894 self.state.save()?;
896 }
897 return Ok(());
898 }
899 if self.state.status == RunStatus::Landing {
909 self.run_land().await?;
910 self.settle_questions();
915 return Ok(());
916 }
917 self.prep().await?;
918 if self.park_here()? {
919 return Ok(());
920 }
921 self.advise().await?;
922 if self.park_here()? {
923 return Ok(());
924 }
925 self.implement().await?;
926 if self.park_here()? {
927 return Ok(());
928 }
929 if self.state.status == RunStatus::VerifiedNoop {
933 return Ok(());
934 }
935 self.judge().await?;
936 if self.park_here()? {
937 return Ok(());
938 }
939 self.deliberate().await?;
940 if self.park_here()? {
941 return Ok(());
942 }
943 self.vote().await?;
944 if self.park_here()? {
945 return Ok(());
946 }
947 self.tally()?;
948 if self.state.status == RunStatus::Stalled {
953 self.state.save()?;
957 return Ok(());
958 }
959 self.finish_after_tally().await?;
960 Ok(())
961 }
962
963 fn park_here(&mut self) -> Result<bool> {
970 if !self.pause.parked() && !self.interrupt.parked() {
975 return Ok(false);
976 }
977 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
978 Some(reason) => format!(
979 "parked after `{}` ({reason}) — resume to carry on from here",
980 self.state.status.as_str()
981 ),
982 None => format!(
983 "parked after `{}` — resume to carry on from here",
984 self.state.status.as_str()
985 ),
986 };
987 self.state.event("park", why);
988 self.state.parked = true;
989 self.state.save()?;
990 Ok(true)
991 }
992
993 pub fn on_pause(&mut self, pause: Pause) {
995 self.pause = pause;
996 }
997
998 pub fn watch_interrupt(&mut self, pause: Pause) {
1004 self.interrupt = pause;
1005 }
1006
1007 fn settle_questions(&mut self) {
1027 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
1028 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
1029 }
1030 }
1031
1032 async fn finish_after_tally(&mut self) -> Result<()> {
1035 self.fold_losers().await?;
1036 self.sync_to_base().await?;
1041 if self.state.status == RunStatus::AlreadyInBase {
1042 return Ok(());
1043 }
1044 self.review_loop().await?;
1045 self.sync_to_base().await?;
1046 if self.state.status == RunStatus::AlreadyInBase {
1047 return Ok(());
1048 }
1049 self.gate().await?;
1050 self.merge().await?;
1051 self.state.save()?;
1052 Ok(())
1053 }
1054
1055 async fn prep(&mut self) -> Result<()> {
1058 if !self.state.candidates.is_empty() {
1059 return Ok(());
1060 }
1061 self.state.status = RunStatus::Prep;
1062 let repo = self.state.repo.clone();
1063 let base = self.state.base_commit.clone();
1064 let plan = refs::plan(&repo, &self.state.seeds).await?;
1065 let start = plan.start.clone().unwrap_or_else(|| base.clone());
1066 let root = self.state.worktree_root();
1067 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
1068
1069 let hooks_dir = self.state.dir().join("hooks");
1072 if self.state.config.blind.commit_msg_hook {
1073 std::fs::create_dir_all(&hooks_dir)
1074 .with_context(|| format!("create {}", hooks_dir.display()))?;
1075 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
1076 let path = hooks_dir.join("commit-msg");
1077 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
1078 make_executable(&path)?;
1079 git::acquire_worktree_config(&repo).await?;
1087 self.state.enabled_worktree_config = true;
1088 }
1089
1090 for (index, (spec, label)) in self
1091 .roles
1092 .implementers
1093 .clone()
1094 .into_iter()
1095 .zip(labels)
1096 .enumerate()
1097 {
1098 let branch = self.state.branch_for(label);
1099 let worktree = root.join(format!("cand-{label}"));
1100 git::worktree_add_branch(&repo, &worktree, &branch, &start).await?;
1101 if self.state.config.blind.commit_msg_hook {
1102 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
1103 }
1104 git::local_exclude(&worktree, "/.magi/").await?;
1105 for pick in &plan.picks {
1106 if let Err(e) = git::cherry_pick(&worktree, pick).await {
1107 self.state.status = RunStatus::Blocked;
1108 self.state
1109 .event("prep", format!("cannot apply referenced commit: {e}"));
1110 self.state.save()?;
1111 return Err(e);
1112 }
1113 }
1114 self.state.candidates.push(Candidate {
1115 index,
1116 label,
1117 agent: spec.id.clone(),
1118 branch,
1119 worktree,
1120 summary: String::new(),
1121 stat: String::new(),
1122 files: 0,
1123 commits: 0,
1124 empty: false,
1125 failed: None,
1126 verified_noop: None,
1127 duration_ms: 0,
1128 folded: false,
1129 });
1130 }
1131
1132 for j in 1..=self.roles.judges.len() {
1133 let wt = root.join(format!("judge-{j}"));
1134 if !wt.exists() {
1135 git::worktree_add_detached(&repo, &wt, &base).await?;
1136 }
1137 }
1138
1139 if self.state.config.graph.advise {
1147 for k in 1..=self.state.config.graph.advisors {
1148 let wt = root.join(format!("advisor-{k}"));
1149 if !wt.exists() {
1150 git::worktree_add_detached(&repo, &wt, &base).await?;
1151 }
1152 }
1153 }
1154
1155 let authors: Vec<&str> = self
1160 .roles
1161 .implementers
1162 .iter()
1163 .map(|a| a.id.as_str())
1164 .collect();
1165 let overlap: Vec<String> = self
1166 .roles
1167 .judges
1168 .iter()
1169 .enumerate()
1170 .filter(|(_, j)| authors.contains(&j.id.as_str()))
1171 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
1172 .collect();
1173 if !overlap.is_empty() {
1174 let note = format!(
1175 "{} also authored a candidate; blind, but the panel is less \
1176 independent than {} distinct agents would be",
1177 overlap.join(", "),
1178 self.roles.judges.len()
1179 );
1180 self.state.event("prep", note);
1181 }
1182
1183 self.state.event(
1184 "prep",
1185 format!(
1186 "{} candidates, {} judges, base {} ({})",
1187 self.state.candidates.len(),
1188 self.roles.judges.len(),
1189 &self.state.base_commit[..7.min(self.state.base_commit.len())],
1190 self.state.base_branch
1191 ),
1192 );
1193 self.state.status = RunStatus::Implementing;
1194 self.state.save()?;
1195 Ok(())
1196 }
1197
1198 async fn advise(&mut self) -> Result<()> {
1231 let implement_untouched = self
1232 .state
1233 .candidates
1234 .iter()
1235 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
1236 if !self.state.config.graph.advise || self.state.advise_attempted {
1237 return Ok(());
1238 }
1239 if !implement_untouched {
1240 self.state.event(
1241 "advise",
1242 "skipping the design-deliberation stage: at least one \
1243 candidate already shows implementation progress, so this \
1244 run is past the point the stage exists to run before"
1245 .to_owned(),
1246 );
1247 self.state.advise_attempted = true;
1248 self.state.save()?;
1249 return Ok(());
1250 }
1251 let run_id = self.state.id.clone();
1252 let prompts = self.state.config.prompts.clone();
1253 let instruction = self.state.instruction.clone();
1254 let language = self.state.config.graph.language.clone();
1255 let root = self.state.worktree_root();
1256 let n = self.state.config.graph.advisors;
1257 let where_recorded = self.state.dir().join("run.json");
1258
1259 let seats = match self.state.config.advisors() {
1260 Ok(seats) if !seats.is_empty() => seats,
1261 Ok(_) => {
1262 self.state.event(
1263 "advise",
1264 format!(
1265 "[graph] advisors is 0; skipping the design-deliberation \
1266 stage and continuing without a synthesis brief (see {})",
1267 where_recorded.display()
1268 ),
1269 );
1270 self.state.advise_attempted = true;
1271 self.state.save()?;
1272 return Ok(());
1273 }
1274 Err(e) => {
1275 self.state.event(
1276 "advise",
1277 format!(
1278 "could not resolve advisor seats ({e:#}); continuing \
1279 without a design-deliberation brief (see {})",
1280 where_recorded.display()
1281 ),
1282 );
1283 self.state.advise_attempted = true;
1284 self.state.save()?;
1285 return Ok(());
1286 }
1287 };
1288
1289 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1290 let artifacts = agent::artifacts_dir(&self.state.dir());
1291 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1292
1293 let mut jobs = Vec::new();
1294 for (i, spec) in seats.iter().cloned().enumerate() {
1295 let seat_key = format!("advisor-{}", i + 1);
1296 let seat = self.seat(&seat_key, &spec.id);
1297 jobs.push(SeatJob {
1298 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1299 spec,
1300 seat,
1301 cwd: worktrees[i % worktrees.len()].clone(),
1302 timeout,
1303 allow_write: false,
1304 sessions: false,
1305 artifacts: artifacts.clone(),
1306 stem: seat_key,
1307 });
1308 }
1309
1310 self.state.event(
1311 "advise",
1312 format!(
1313 "{} advisor seat(s) sketching a design in parallel",
1314 jobs.len()
1315 ),
1316 );
1317 let mut quota_losses = Vec::new();
1318 let cache = self.state.config.cache_dir();
1319 let ctx = WaveCtx {
1320 run: &run_id,
1321 node: "advise",
1322 prompts: &prompts,
1323 cache: cache.as_deref(),
1324 round: None,
1325 };
1326 let results = ask_json_wave::<Proposal>(
1327 jobs,
1328 Arc::clone(&self.sem),
1329 self.state.config.graph.retries,
1330 &ctx,
1331 &mut quota_losses,
1332 &mut self.state,
1333 &|p: &Proposal| p.validate(),
1334 )
1335 .await;
1336 self.state.quota.extend(quota_losses);
1337
1338 let mut records = Vec::with_capacity(results.len());
1339 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1340 let agent_id = seat.agent.clone();
1341 self.state.seats.insert(seat.key.clone(), seat);
1342 match res {
1343 Ok((proposal, out)) => {
1344 self.state
1345 .event("advise", format!("advisor-{} proposed a design", i + 1));
1346 records.push(advise::AdvisorRecord::proposed(
1347 i + 1,
1348 agent_id,
1349 proposal,
1350 out.duration_ms,
1351 ));
1352 }
1353 Err(e) => {
1354 self.state.event(
1355 "advise",
1356 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1357 );
1358 records.push(advise::AdvisorRecord::failed(
1359 i + 1,
1360 agent_id,
1361 e.to_string(),
1362 ));
1363 }
1364 }
1365 }
1366
1367 let mut advice = advise::Advice {
1368 records,
1369 synthesis: None,
1370 };
1371 if advice.proposals().is_empty() {
1372 self.state.event(
1373 "advise",
1374 "no advisor produced a usable proposal; continuing without a \
1375 synthesis brief"
1376 .to_owned(),
1377 );
1378 } else {
1379 match self
1380 .synthesize_brief(
1381 &advice,
1382 &instruction,
1383 &language,
1384 &worktrees[0],
1385 &artifacts,
1386 &run_id,
1387 &prompts,
1388 cache.as_deref(),
1389 )
1390 .await
1391 {
1392 Ok(Some(text)) => {
1393 self.state.event(
1394 "advise",
1395 "synthesized a design brief for the implementer".to_owned(),
1396 );
1397 advice.synthesis = Some(text);
1398 }
1399 Ok(None) => {
1400 self.state.event(
1401 "advise",
1402 "the synthesis seat produced nothing usable; continuing \
1403 without a design brief"
1404 .to_owned(),
1405 );
1406 }
1407 Err(e) => {
1408 self.state.event(
1409 "advise",
1410 format!("could not synthesize a design brief: {e:#}"),
1411 );
1412 }
1413 }
1414 }
1415 advise::apply_reflection(&mut advice);
1416
1417 self.state.advice = Some(advice);
1418 self.state.advise_attempted = true;
1419 self.state.save()?;
1420 Ok(())
1421 }
1422
1423 #[allow(clippy::too_many_arguments)]
1435 async fn synthesize_brief(
1436 &mut self,
1437 advice: &advise::Advice,
1438 instruction: &str,
1439 language: &str,
1440 cwd: &Path,
1441 artifacts: &Path,
1442 run_id: &str,
1443 prompts: &Prompts,
1444 cache: Option<&Path>,
1445 ) -> Result<Option<String>> {
1446 let want = self.state.config.roles.synthesizer.as_deref();
1447 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1448 let mut seat = self.seat("advise-synthesis", &spec.id);
1449 let proposals = advice.proposals();
1450 let mut prompt = prompt::with_overlay(
1451 prompt::synthesize_brief(instruction, &proposals, language),
1452 prompts.overlay("advise"),
1453 );
1454 if cache.is_some() {
1455 prompt.push('\n');
1460 prompt.push_str(&prompt::build_cache_note("advise", false));
1461 }
1462 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1463 let out = agent::invoke(
1464 &spec,
1465 &mut seat,
1466 &Invocation {
1467 cwd,
1468 prompt: &prompt,
1469 timeout,
1470 allow_write: false,
1471 sessions: false,
1472 artifacts,
1473 stem: "advise-synthesis",
1474 run: run_id,
1475 node: "advise",
1476 cache_dir: None,
1477 attachments: &[],
1478 },
1479 )
1480 .await?;
1481 self.state.seats.insert(seat.key.clone(), seat);
1482 if !out.usable() {
1483 return Ok(None);
1484 }
1485 let text =
1486 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1487 Ok((!text.trim().is_empty()).then_some(text))
1488 }
1489
1490 async fn implement(&mut self) -> Result<()> {
1493 let run_id = self.state.id.clone();
1498 let prompts = self.state.config.prompts.clone();
1499 let todo: Vec<usize> = self
1500 .state
1501 .candidates
1502 .iter()
1503 .enumerate()
1504 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1505 .map(|(i, _)| i)
1506 .collect();
1507 if todo.is_empty() {
1508 return self.after_implement();
1509 }
1510 self.state.status = RunStatus::Implementing;
1511
1512 let language = self.state.config.graph.language.clone();
1513 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1514 let sessions = self.state.config.graph.sessions;
1515 let artifacts = agent::artifacts_dir(&self.state.dir());
1516 let brief = self
1520 .state
1521 .advice
1522 .as_ref()
1523 .and_then(|a| a.synthesis.as_deref())
1524 .map(str::to_owned);
1525 let attachments = self.state.attachments.clone();
1526
1527 let mut jobs = Vec::new();
1528 for &i in &todo {
1529 let (index, label, worktree) = {
1530 let c = &self.state.candidates[i];
1531 (c.index, c.label, c.worktree.clone())
1532 };
1533 let spec = self.roles.implementers[index].clone();
1534 let seat_key = format!("impl-{label}");
1535 let seat = self.seat(&seat_key, &spec.id);
1536 let instruction = seeded_instruction(&self.state);
1537 jobs.push(SeatJob {
1538 spec,
1539 seat,
1540 prompt: prompt::implement(
1541 &instruction,
1542 &worktree.to_string_lossy(),
1543 &language,
1544 brief.as_deref(),
1545 &attachments,
1546 ),
1547 cwd: worktree,
1548 timeout,
1549 allow_write: true,
1550 sessions,
1551 artifacts: artifacts.clone(),
1552 stem: format!("impl-{label}"),
1553 });
1554 }
1555
1556 self.state.event(
1557 "implement",
1558 format!("{} candidates in parallel", jobs.len()),
1559 );
1560 let mut sent = jobs.clone();
1566 let cache = self.state.config.cache_dir();
1567 let ctx = WaveCtx {
1568 run: &run_id,
1569 node: "implement",
1570 prompts: &prompts,
1571 cache: cache.as_deref(),
1572 round: None,
1573 };
1574 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1575 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1576 .await;
1577 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1578 .await;
1579 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1580 .await;
1581
1582 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1583 let seat_key = seat.key.clone();
1584 let agent = seat.agent.clone();
1593 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1594 self.state.seats.insert(seat.key.clone(), seat);
1595 let label = self.state.candidates[i].label;
1596 let worktree = self.state.candidates[i].worktree.clone();
1597 let base = self.state.base_commit.clone();
1598
1599 let (summary, duration, failed, verified_claim) = match out {
1600 AgentOutcome::Ok(o) => {
1601 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1602 let failed = (!o.usable()).then(|| {
1603 if o.timed_out {
1604 "agent timed out".to_owned()
1605 } else {
1606 format!("agent exited with {:?}", o.exit_code)
1607 }
1608 });
1609 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1610 (text, o.duration_ms, failed, verified_claim)
1611 }
1612 AgentOutcome::Dropped(o) => {
1618 let why = o
1619 .dropped
1620 .as_ref()
1621 .map(|d| d.why.as_str())
1622 .unwrap_or("the CLI ended the stream without delivering its answer");
1623 (
1624 String::new(),
1625 o.duration_ms,
1626 Some(format!("the CLI dropped the stream ({why})")),
1627 None,
1628 )
1629 }
1630 AgentOutcome::Quota(o) => {
1631 self.state.quota.push(QuotaLoss {
1632 seat: seat_key,
1633 node: "implement".to_owned(),
1634 at: Timestamp::now(),
1635 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1636 });
1637 (
1638 String::new(),
1639 o.duration_ms,
1640 Some("rate limited (quota); produced no change".to_owned()),
1641 None,
1642 )
1643 }
1644 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1645 };
1646
1647 let rescued = match git::rescue_commit(
1650 &worktree,
1651 &format!("magi: candidate {label} (uncommitted work)"),
1652 )
1653 .await
1654 {
1655 Ok(r) => {
1656 self.state.note_withheld("implement", &r.withheld);
1657 r.committed
1658 }
1659 Err(_) => false,
1660 };
1661 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1662 .await
1663 .unwrap_or(0);
1664 let patch = git::diff(&worktree, &base, "HEAD")
1665 .await
1666 .unwrap_or_default();
1667 let stat = git::diff_stat(&worktree, &base, "HEAD")
1668 .await
1669 .unwrap_or_default();
1670 let files = git::changed_files(&worktree, &base, "HEAD")
1671 .await
1672 .map(|f| f.len())
1673 .unwrap_or(0);
1674 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1675
1676 let c = &mut self.state.candidates[i];
1677 if !exhausted_the_fallback_chain {
1678 c.agent = agent;
1679 }
1680 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1681 c.stat = stat;
1682 c.files = files;
1683 c.commits = commits;
1684 c.duration_ms = duration;
1685 c.empty = commits == 0 || patch.trim().is_empty();
1686 c.failed = match failed {
1689 Some(_) if c.empty => failed,
1690 _ => None,
1691 };
1692 c.verified_noop = if c.empty { verified_claim } else { None };
1697 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1698 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1699 (None, true, Some(_), _) => {
1700 format!("candidate {label}: no change produced (agent-verified no-op)")
1701 }
1702 (None, true, None, _) => format!("candidate {label}: no change produced"),
1703 (None, false, _, true) => {
1704 format!(
1705 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1706 )
1707 }
1708 (None, false, _, false) => {
1709 format!("candidate {label}: {files} files, {commits} commits")
1710 }
1711 };
1712 self.state.event("implement", note);
1713 self.state.save()?;
1714 }
1715
1716 self.after_implement()
1717 }
1718
1719 async fn resume_undelivered(
1747 &mut self,
1748 results: &mut [(usize, SeatState, AgentOutcome)],
1749 sent: &[SeatJob],
1750 prompts: &Prompts,
1751 run_id: &str,
1752 ) {
1753 for (wi, seat, out) in results.iter_mut() {
1754 let Some(dropped) = (match &*out {
1755 AgentOutcome::Dropped(o) => o.dropped.clone(),
1756 _ => None,
1757 }) else {
1758 continue;
1759 };
1760 let Some(job) = sent.get(*wi) else { continue };
1761 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1763 self.state.event(
1764 "implement",
1765 format!(
1766 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1767 work is in the tree",
1768 seat.key, dropped.output_tokens, dropped.why
1769 ),
1770 );
1771 continue;
1772 }
1773 if !has_context(&job.spec, seat, job.sessions) {
1781 self.state.event(
1782 "implement",
1783 format!(
1784 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1785 is no session left to resume",
1786 seat.key, dropped.output_tokens, dropped.why
1787 ),
1788 );
1789 continue;
1790 }
1791 self.state.event(
1792 "implement",
1793 format!(
1794 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1795 conversation",
1796 seat.key, dropped.output_tokens, dropped.why
1797 ),
1798 );
1799 let mut retry = job.clone();
1800 retry.seat = seat.clone();
1801 retry.prompt = prompt::resume_after_drop(&dropped.why);
1802 retry.timeout = retry_budget(job.timeout, true);
1803 retry.stem = format!("{}-resume", job.stem);
1804 let cache = self.state.config.cache_dir();
1805 let ctx = WaveCtx {
1806 run: run_id,
1807 node: "implement",
1808 prompts,
1809 cache: cache.as_deref(),
1810 round: None,
1811 };
1812 let (resumed_seat, resumed) =
1813 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1814 *seat = resumed_seat;
1815 *out = resumed;
1816 }
1817 }
1818
1819 async fn resume_quota_losses(
1881 &mut self,
1882 results: &mut [(usize, SeatState, AgentOutcome)],
1883 sent: &mut [SeatJob],
1884 prompts: &Prompts,
1885 run_id: &str,
1886 ) {
1887 let instruction = seeded_instruction(&self.state);
1888 let language = self.state.config.graph.language.clone();
1889 let brief = self
1890 .state
1891 .advice
1892 .as_ref()
1893 .and_then(|a| a.synthesis.as_deref())
1894 .map(str::to_owned);
1895 let attachments = self.state.attachments.clone();
1896 for (wi, seat, out) in results.iter_mut() {
1897 let Some(job) = sent.get_mut(*wi) else {
1898 continue;
1899 };
1900 let start = self
1905 .roles
1906 .implementer_roster
1907 .iter()
1908 .position(|s| s.id == job.spec.id)
1909 .unwrap_or(0);
1910 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1911 let mut fallback_attempt = 0usize;
1912 while matches!(&*out, AgentOutcome::Quota(_)) {
1913 let Some(next) =
1914 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1915 .cloned()
1916 else {
1917 break;
1918 };
1919 tried.insert(next.id.clone());
1920 fallback_attempt += 1;
1921
1922 if let Ok(r) = git::rescue_commit(
1923 &job.cwd,
1924 &format!(
1925 "magi: candidate {} (uncommitted work before quota fallback)",
1926 seat.key
1927 ),
1928 )
1929 .await
1930 {
1931 self.state.note_withheld("implement", &r.withheld);
1932 }
1933
1934 self.state.event(
1935 "implement",
1936 format!(
1937 "{}: rate limited (quota) on {}; retrying with {}",
1938 seat.key, seat.agent, next.id
1939 ),
1940 );
1941
1942 let new_seat = self.seat(&seat.key, &next.id);
1943 job.spec = next.clone();
1951 let mut retry = job.clone();
1952 retry.seat = new_seat;
1953 retry.prompt = prompt::implement(
1954 &instruction,
1955 &job.cwd.to_string_lossy(),
1956 &language,
1957 brief.as_deref(),
1958 &attachments,
1959 );
1960 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1961 let cache = self.state.config.cache_dir();
1962 let ctx = WaveCtx {
1963 run: run_id,
1964 node: "implement",
1965 prompts,
1966 cache: cache.as_deref(),
1967 round: None,
1968 };
1969 let (fallback_seat, fallback_out) = run_one(
1970 retry,
1971 Arc::clone(&self.sem),
1972 &ctx,
1973 &mut self.state,
1974 fallback_attempt,
1975 )
1976 .await;
1977 *seat = fallback_seat;
1978 *out = fallback_out;
1979 }
1980 }
1981 }
1982
1983 async fn resume_unconfirmed_commands(
2007 &mut self,
2008 results: &mut [(usize, SeatState, AgentOutcome)],
2009 sent: &[SeatJob],
2010 prompts: &Prompts,
2011 run_id: &str,
2012 ) {
2013 for (wi, seat, out) in results.iter_mut() {
2014 let AgentOutcome::Ok(o) = &*out else {
2015 continue;
2016 };
2017 if !has_unconfirmed_command(&o.commands) {
2018 continue;
2019 }
2020 let Some(job) = sent.get(*wi) else { continue };
2021 if !has_context(&job.spec, seat, job.sessions) {
2022 self.state.event(
2023 "implement",
2024 format!(
2025 "{}: the reply named a command whose own CLI never confirmed the exit \
2026 status of, but there is no session left to resume",
2027 seat.key
2028 ),
2029 );
2030 continue;
2031 }
2032 self.state.event(
2033 "implement",
2034 format!(
2035 "{}: the reply named a command whose own CLI never confirmed the exit \
2036 status of; resuming the conversation",
2037 seat.key
2038 ),
2039 );
2040 let mut retry = job.clone();
2041 retry.seat = seat.clone();
2042 retry.prompt = prompt::resume_incomplete(
2043 "a command in your last reply had no confirmed exit status",
2044 );
2045 retry.timeout = retry_budget(job.timeout, true);
2046 retry.stem = format!("{}-confirm", job.stem);
2047 let cache = self.state.config.cache_dir();
2048 let ctx = WaveCtx {
2049 run: run_id,
2050 node: "implement",
2051 prompts,
2052 cache: cache.as_deref(),
2053 round: None,
2054 };
2055 let (resumed_seat, resumed) =
2056 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
2057 *seat = resumed_seat;
2058 *out = resumed;
2059 }
2060 }
2061
2062 async fn continue_fix_report(
2083 &mut self,
2084 mut seat: SeatState,
2085 parse_err: String,
2086 job: &SeatJob,
2087 prompts: &Prompts,
2088 run_id: &str,
2089 round: usize,
2090 ) -> (
2091 SeatState,
2092 Option<FixReport>,
2093 Option<String>,
2094 ContinuationRecord,
2095 ) {
2096 let mut last_err = parse_err;
2097 let mut cumulative_wait_ms = 0u64;
2098 let mut attempts = 0usize;
2099 loop {
2100 if !has_context(&job.spec, &seat, job.sessions) {
2101 self.state.event(
2102 "fix",
2103 format!(
2104 "round {round}: fixer's reply had no adoption report ({last_err}); no \
2105 session left to resume into"
2106 ),
2107 );
2108 let outcome = if attempts == 0 {
2109 ContinuationOutcome::NoSession
2110 } else {
2111 ContinuationOutcome::Exhausted
2112 };
2113 return (
2114 seat,
2115 None,
2116 Some(format!("unparsable fix report: {last_err}")),
2117 ContinuationRecord {
2118 attempts,
2119 cumulative_wait_ms,
2120 outcome,
2121 },
2122 );
2123 }
2124 if attempts >= MAX_FIX_CONTINUATIONS {
2125 self.state.event(
2126 "fix",
2127 format!(
2128 "round {round}: fixer's reply still had no adoption report after \
2129 {attempts} continuation(s) ({last_err}); giving up"
2130 ),
2131 );
2132 return (
2133 seat,
2134 None,
2135 Some(format!(
2136 "unparsable fix report after {attempts} continuation(s): {last_err}"
2137 )),
2138 ContinuationRecord {
2139 attempts,
2140 cumulative_wait_ms,
2141 outcome: ContinuationOutcome::Exhausted,
2142 },
2143 );
2144 }
2145 attempts += 1;
2146 self.state.event(
2147 "fix",
2148 format!(
2149 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
2150 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
2151 ),
2152 );
2153 let mut retry = job.clone();
2154 retry.seat = seat.clone();
2155 retry.prompt = prompt::resume_incomplete(&last_err);
2156 retry.timeout = retry_budget(job.timeout, true);
2157 retry.stem = format!("{}-continue{attempts}", job.stem);
2158 let cache = self.state.config.cache_dir();
2159 let ctx = WaveCtx {
2160 run: run_id,
2161 node: "fix",
2162 prompts,
2163 cache: cache.as_deref(),
2164 round: Some(round),
2165 };
2166 let (resumed_seat, resumed_out) = run_one(
2167 retry,
2168 Arc::clone(&self.sem),
2169 &ctx,
2170 &mut self.state,
2171 attempts,
2172 )
2173 .await;
2174 seat = resumed_seat;
2175 match resumed_out {
2176 AgentOutcome::Ok(o) => {
2177 cumulative_wait_ms += o.duration_ms;
2178 match verdict::extract_json::<FixReport>(&o.text) {
2179 Ok(report) if !has_unconfirmed_command(&o.commands) => {
2180 self.state.event(
2181 "fix",
2182 format!(
2183 "round {round}: fixer's adoption report recovered after \
2184 {attempts} continuation(s)"
2185 ),
2186 );
2187 return (
2188 seat,
2189 Some(report),
2190 None,
2191 ContinuationRecord {
2192 attempts,
2193 cumulative_wait_ms,
2194 outcome: ContinuationOutcome::Resumed,
2195 },
2196 );
2197 }
2198 Ok(_) => {
2206 last_err = "the reply parsed, but it reported a command whose own CLI \
2207 never confirmed an exit status"
2208 .to_owned();
2209 }
2210 Err(e) => last_err = e.to_string(),
2211 }
2212 }
2213 AgentOutcome::Quota(o) => {
2214 cumulative_wait_ms += o.duration_ms;
2215 self.state.quota.push(QuotaLoss {
2216 seat: seat.key.clone(),
2217 node: "fix".to_owned(),
2218 at: Timestamp::now(),
2219 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2220 });
2221 self.state.event(
2222 "fix",
2223 format!(
2224 "round {round}: continuation rate limited (quota); not retrying now"
2225 ),
2226 );
2227 return (
2228 seat,
2229 None,
2230 Some("rate limited (quota) while recovering the fix report".to_owned()),
2231 ContinuationRecord {
2232 attempts,
2233 cumulative_wait_ms,
2234 outcome: ContinuationOutcome::QuotaLost,
2235 },
2236 );
2237 }
2238 AgentOutcome::Dropped(o) => {
2239 cumulative_wait_ms += o.duration_ms;
2240 let why = o
2241 .dropped
2242 .as_ref()
2243 .map(|d| d.why.as_str())
2244 .unwrap_or("the CLI ended the stream without delivering its answer");
2245 last_err = format!("the CLI dropped the stream ({why})");
2246 }
2247 AgentOutcome::Failed(e) => last_err = e,
2248 }
2249 }
2250 }
2251
2252 fn after_implement(&mut self) -> Result<()> {
2253 if self.state.leaks.is_empty() {
2255 let cfg = self.state.config.blind.clone();
2256 let mut leaks = Vec::new();
2257 for c in &self.state.candidates {
2258 let Some(patch) =
2259 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
2260 else {
2261 continue;
2262 };
2263 leaks.extend(blind::scan(
2264 &format!("candidate {} patch", c.label),
2265 &patch,
2266 &cfg.vendor_tokens,
2267 ));
2268 }
2269 if !leaks.is_empty() {
2270 let summary = leaks
2271 .iter()
2272 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2273 .collect::<Vec<_>>()
2274 .join(", ");
2275 match cfg.on_leak {
2276 LeakPolicy::Fail => {
2277 self.state.status = RunStatus::Failed;
2278 self.state
2279 .event("blind", format!("vendor text in a patch: {summary}"));
2280 self.state.leaks = leaks;
2281 self.state.save()?;
2282 self.settle_questions();
2283 bail!(
2284 "blind.on_leak = \"fail\" and vendor text reached a \
2285 judged patch: {summary}"
2286 );
2287 }
2288 LeakPolicy::Redact => self.state.event(
2289 "blind",
2290 format!("redacting vendor text for judging: {summary}"),
2291 ),
2292 LeakPolicy::Warn => self.state.event(
2293 "blind",
2294 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2295 ),
2296 }
2297 self.state.leaks = leaks;
2298 }
2299 }
2300
2301 if self.state.viable().is_empty() {
2302 if self.state.all_candidates_verified_noop() {
2303 self.state.status = RunStatus::VerifiedNoop;
2314 self.state.save()?;
2315 self.settle_questions();
2316 return Ok(());
2317 }
2318 self.state.status = RunStatus::Failed;
2319 self.state.save()?;
2320 self.settle_questions();
2321 bail!("no candidate produced a change; nothing to judge");
2322 }
2323 self.state.status = RunStatus::Judging;
2324 self.state.save()?;
2325 Ok(())
2326 }
2327
2328 async fn judge(&mut self) -> Result<()> {
2331 let run_id = self.state.id.clone();
2336 let prompts = self.state.config.prompts.clone();
2337 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2338 return Ok(());
2339 }
2340 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2341 if viable.len() == 1 {
2342 self.state.judge_skipped = true;
2349 self.state.event(
2350 "judge",
2351 format!(
2352 "only candidate {} produced a change; judging skipped",
2353 viable[0].label
2354 ),
2355 );
2356 self.state.save()?;
2357 return Ok(());
2358 }
2359 self.state.status = RunStatus::Judging;
2360
2361 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2362 let language = self.state.config.graph.language.clone();
2363 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2364 let sessions = self.state.config.graph.sessions;
2365 let artifacts = agent::artifacts_dir(&self.state.dir());
2366 let root = self.state.worktree_root();
2367 let base_short = short(&self.state.base_commit);
2368
2369 let mut jobs = Vec::new();
2370 let mut orders = Vec::new();
2371 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2372 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2373 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2374 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2375 let seat_key = format!("judge-{}", j + 1);
2376 let seat = self.seat(&seat_key, &spec.id);
2377 jobs.push(SeatJob {
2378 prompt: prompt::judge(
2379 &self.state.instruction,
2380 &views,
2381 self.roles.judges.len(),
2382 &base_short,
2383 &language,
2384 ),
2385 spec,
2386 seat,
2387 cwd: root.join(format!("judge-{}", j + 1)),
2388 timeout,
2389 allow_write: false,
2390 sessions,
2391 artifacts: artifacts.clone(),
2392 stem: format!("judge-{}", j + 1),
2393 });
2394 }
2395
2396 self.state.event(
2397 "judge",
2398 format!(
2399 "{} judges ranking {} candidates blind",
2400 jobs.len(),
2401 viable.len()
2402 ),
2403 );
2404 let labels_for_check = labels.clone();
2405 let mut quota_losses = Vec::new();
2406 let cache = self.state.config.cache_dir();
2407 let ctx = WaveCtx {
2408 run: &run_id,
2409 node: "judge",
2410 prompts: &prompts,
2411 cache: cache.as_deref(),
2412 round: None,
2413 };
2414 let results = ask_json_wave::<Ranking>(
2415 jobs,
2416 Arc::clone(&self.sem),
2417 self.state.config.graph.retries,
2418 &ctx,
2419 &mut quota_losses,
2420 &mut self.state,
2421 &move |r: &Ranking| r.validate(&labels_for_check),
2422 )
2423 .await;
2424 self.state.quota.extend(quota_losses);
2425
2426 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2427 let agent_id = seat.agent.clone();
2428 self.state.seats.insert(seat.key.clone(), seat);
2429 let mut record = Judgement {
2430 judge: j + 1,
2431 seat: format!("judge-{}", j + 1),
2432 agent: agent_id,
2433 ranking: Vec::new(),
2434 reasons: BTreeMap::new(),
2435 confidence: None,
2436 order: orders[j].clone(),
2437 failed: None,
2438 duration_ms: 0,
2439 };
2440 match res {
2441 Ok((ranking, out)) => {
2442 record.ranking = ranking.normalized();
2443 record.reasons = ranking.reasons;
2444 record.confidence = ranking.confidence;
2445 record.duration_ms = out.duration_ms;
2446 self.state.event(
2447 "judge",
2448 format!(
2449 "judge {} ranked {}",
2450 j + 1,
2451 record.ranking.iter().collect::<String>()
2452 ),
2453 );
2454 }
2455 Err(e) => {
2456 record.failed = Some(e.to_string());
2457 self.state
2458 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2459 }
2460 }
2461 self.state.judgements.push(record);
2462 self.state.save()?;
2463 }
2464 Ok(())
2465 }
2466
2467 async fn deliberate(&mut self) -> Result<()> {
2470 let run_id = self.state.id.clone();
2475 let prompts = self.state.config.prompts.clone();
2476 if !self.state.deliberation.is_empty() {
2477 return Ok(());
2478 }
2479 let tops: Vec<char> = self
2480 .state
2481 .judgements
2482 .iter()
2483 .filter_map(|j| j.ranking.first().copied())
2484 .collect();
2485 let rounds = self.state.config.graph.deliberate_rounds;
2486 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2487 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2488 self.state.event(
2489 "deliberate",
2490 format!("judges agreed on {} outright; no deliberation", tops[0]),
2491 );
2492 }
2493 self.state.status = RunStatus::Voting;
2494 self.state.save()?;
2495 return Ok(());
2496 }
2497
2498 self.state.status = RunStatus::Deliberating;
2499 self.state.event(
2500 "deliberate",
2501 format!(
2502 "split: first choices were {} — opening {rounds} round(s)",
2503 tops.iter().collect::<String>()
2504 ),
2505 );
2506
2507 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2508 let language = self.state.config.graph.language.clone();
2509 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2510 let sessions = self.state.config.graph.sessions;
2511 let artifacts = agent::artifacts_dir(&self.state.dir());
2512 let root = self.state.worktree_root();
2513 let base_short = short(&self.state.base_commit);
2514
2515 for round in 1..=rounds {
2519 let mut turns: Vec<DeliberationTurn> = Vec::new();
2520 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2521 if self.state.judgements[j].failed.is_some() {
2522 continue;
2523 }
2524 let seat_key = format!("judge-{}", j + 1);
2525 let mut seat = self.seat(&seat_key, &spec.id);
2526 let transcript = self.transcript(&turns, j);
2527 let context = if has_context(&spec, &seat, sessions) {
2528 None
2529 } else {
2530 Some(self.candidate_block(&viable, &base_short))
2531 };
2532 let text = prompt::deliberate(
2533 &self.state.instruction,
2534 context.as_deref(),
2535 &transcript,
2536 round,
2537 rounds,
2538 &language,
2539 );
2540 let job = SeatJob {
2541 spec,
2542 seat: seat.clone(),
2543 prompt: text,
2544 cwd: root.join(format!("judge-{}", j + 1)),
2545 timeout,
2546 allow_write: false,
2547 sessions,
2548 artifacts: artifacts.clone(),
2549 stem: format!("delib-{round}-judge-{}", j + 1),
2550 };
2551 let cache = self.state.config.cache_dir();
2552 let ctx = WaveCtx {
2553 run: &run_id,
2554 node: "deliberate",
2555 prompts: &prompts,
2556 cache: cache.as_deref(),
2557 round: None,
2558 };
2559 let (updated, out) =
2560 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2561 seat = updated;
2562 let agent_id = seat.agent.clone();
2563 let seat_key = seat.key.clone();
2564 self.state.seats.insert(seat.key.clone(), seat);
2565 let body = match out {
2566 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2567 AgentOutcome::Dropped(o) => {
2571 let why =
2572 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2573 "the CLI ended the stream without delivering its answer",
2574 );
2575 self.state.event(
2576 "deliberate",
2577 format!(
2578 "judge {} skipped: the CLI dropped the stream ({why})",
2579 j + 1
2580 ),
2581 );
2582 continue;
2583 }
2584 AgentOutcome::Quota(o) => {
2585 self.state.quota.push(QuotaLoss {
2586 seat: seat_key,
2587 node: "deliberate".to_owned(),
2588 at: Timestamp::now(),
2589 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2590 });
2591 self.state.event(
2592 "deliberate",
2593 format!("judge {} skipped: rate limited (quota)", j + 1),
2594 );
2595 continue;
2596 }
2597 AgentOutcome::Failed(e) => {
2598 self.state
2599 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2600 continue;
2601 }
2602 };
2603 let tentative = verdict::extract_json::<Position>(&body)
2604 .ok()
2605 .and_then(|p| p.tentative)
2606 .and_then(|s| s.trim().chars().next())
2607 .map(|c| c.to_ascii_uppercase());
2608 self.state.event(
2609 "deliberate",
2610 format!(
2611 "round {round}: judge {} now favours {}",
2612 j + 1,
2613 tentative.map_or("—".to_owned(), |c| c.to_string())
2614 ),
2615 );
2616 turns.push(DeliberationTurn {
2617 judge: j + 1,
2618 agent: agent_id,
2619 body: blind::sanitize_prose(&body, &self.state.config.blind),
2620 tentative,
2621 });
2622 }
2623 self.state
2624 .deliberation
2625 .push(DeliberationRound { round, turns });
2626 self.state.save()?;
2627 }
2628
2629 self.state.status = RunStatus::Voting;
2630 self.state.save()?;
2631 Ok(())
2632 }
2633
2634 async fn vote(&mut self) -> Result<()> {
2637 let run_id = self.state.id.clone();
2642 let prompts = self.state.config.prompts.clone();
2643 if !self.state.votes.is_empty() {
2644 return Ok(());
2645 }
2646 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2647 if viable.len() == 1 {
2648 return Ok(());
2649 }
2650 self.state.status = RunStatus::Voting;
2651
2652 let language = self.state.config.graph.language.clone();
2653 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2654 let sessions = self.state.config.graph.sessions;
2655 let artifacts = agent::artifacts_dir(&self.state.dir());
2656 let root = self.state.worktree_root();
2657 let base_short = short(&self.state.base_commit);
2658 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2659
2660 let mut jobs = Vec::new();
2661 let mut seats_at = Vec::new();
2662 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2663 if self
2664 .state
2665 .judgements
2666 .get(j)
2667 .is_some_and(|r| r.failed.is_some())
2668 {
2669 continue;
2670 }
2671 let seat_key = format!("judge-{}", j + 1);
2672 let seat = self.seat(&seat_key, &spec.id);
2673 let mut text = prompt::final_vote(&viable, &language);
2674 if !has_context(&spec, &seat, sessions) {
2675 text = format!(
2676 "{}\n\n# Candidates\n\n{}",
2677 text,
2678 self.candidate_block(&candidates, &base_short)
2679 );
2680 }
2681 jobs.push(SeatJob {
2682 spec,
2683 seat,
2684 prompt: text,
2685 cwd: root.join(format!("judge-{}", j + 1)),
2686 timeout,
2687 allow_write: false,
2688 sessions,
2689 artifacts: artifacts.clone(),
2690 stem: format!("vote-judge-{}", j + 1),
2691 });
2692 seats_at.push(j);
2693 }
2694
2695 self.state.event(
2696 "vote",
2697 format!(
2698 "collecting {} final votes one by one, privately",
2699 jobs.len()
2700 ),
2701 );
2702 let allowed = viable.clone();
2703 let mut quota_losses = Vec::new();
2704 let cache = self.state.config.cache_dir();
2705 let ctx = WaveCtx {
2706 run: &run_id,
2707 node: "vote",
2708 prompts: &prompts,
2709 cache: cache.as_deref(),
2710 round: None,
2711 };
2712 let results = ask_json_wave::<FinalVote>(
2713 jobs,
2714 Arc::clone(&self.sem),
2715 self.state.config.graph.retries,
2716 &ctx,
2717 &mut quota_losses,
2718 &mut self.state,
2719 &move |v: &FinalVote| match v.label() {
2720 Some(c) if allowed.contains(&c) => Ok(()),
2721 other => bail!("vote {other:?} is not one of {allowed:?}"),
2722 },
2723 )
2724 .await;
2725 self.state.quota.extend(quota_losses);
2726
2727 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2728 let agent_id = seat.agent.clone();
2729 self.state.seats.insert(seat.key.clone(), seat);
2730 let initial = self
2731 .state
2732 .judgements
2733 .get(j)
2734 .and_then(|r| r.ranking.first().copied());
2735 let mut record = VoteRecord {
2736 judge: j + 1,
2737 agent: agent_id,
2738 vote: None,
2739 reason: String::new(),
2740 changed: false,
2741 };
2742 match res {
2743 Ok((v, _)) => {
2744 record.vote = v.label();
2745 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2746 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2747 self.state.event(
2748 "vote",
2749 format!(
2750 "judge {} voted {}{}",
2751 j + 1,
2752 record.vote.unwrap_or('?'),
2753 if record.changed { " (changed)" } else { "" }
2754 ),
2755 );
2756 }
2757 Err(e) => {
2758 self.state
2759 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2760 }
2761 }
2762 self.state.votes.push(record);
2763 self.state.save()?;
2764 }
2765 Ok(())
2766 }
2767
2768 fn tally(&mut self) -> Result<()> {
2771 if self.state.tally.is_some() {
2772 return Ok(());
2773 }
2774 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2775 let tops: Vec<char> = self
2776 .state
2777 .judgements
2778 .iter()
2779 .filter_map(|j| j.ranking.first().copied())
2780 .collect();
2781 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2782
2783 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2786 let mut cast: Vec<char> = Vec::new();
2787 for (i, j) in self.state.judgements.iter().enumerate() {
2788 let vote = self
2789 .state
2790 .votes
2791 .iter()
2792 .find(|v| v.judge == i + 1)
2793 .and_then(|v| v.vote)
2794 .or_else(|| j.ranking.first().copied());
2795 if let Some(v) = vote {
2796 *first_choice.entry(v).or_insert(0) += 1;
2797 cast.push(v);
2798 }
2799 }
2800
2801 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2802 for j in &self.state.judgements {
2803 let n = j.ranking.len();
2804 for (pos, label) in j.ranking.iter().enumerate() {
2805 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2806 }
2807 }
2808
2809 let best = first_choice.values().copied().max().unwrap_or(0);
2810 let mut leaders: Vec<char> = first_choice
2811 .iter()
2812 .filter(|(_, v)| **v == best)
2813 .map(|(k, _)| *k)
2814 .collect();
2815 let mut tie_break = None;
2816 if leaders.len() > 1 {
2817 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2818 let borda_leaders: Vec<char> = leaders
2819 .iter()
2820 .copied()
2821 .filter(|l| borda[l] == top_borda)
2822 .collect();
2823 tie_break = Some(if borda_leaders.len() == 1 {
2824 format!(
2825 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2826 leaders.len()
2827 )
2828 } else {
2829 format!(
2830 "{} way tie on both first-choice votes and Borda points, broken by label order",
2831 leaders.len()
2832 )
2833 });
2834 leaders = borda_leaders;
2835 leaders.sort_unstable();
2836 }
2837 let winner = *leaders
2838 .first()
2839 .or(viable.first())
2840 .context("no candidate to declare a winner from")?;
2841
2842 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2843 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2844 let deliberated = !self.state.deliberation.is_empty();
2845
2846 let quota_seats: std::collections::BTreeSet<&str> =
2850 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2851 let mut present = 0usize;
2852 for (i, j) in self.state.judgements.iter().enumerate() {
2853 if quota_seats.contains(j.seat.as_str()) {
2854 continue;
2855 }
2856 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2857 let voted = self
2858 .state
2859 .votes
2860 .iter()
2861 .any(|v| v.judge == i + 1 && v.vote.is_some());
2862 if ranked || voted {
2863 present += 1;
2864 }
2865 }
2866 let needs_quorum = viable.len() > 1;
2872 let judges_total = if needs_quorum {
2873 self.roles.judges.len()
2874 } else {
2875 0
2876 };
2877 let quorum = if needs_quorum {
2878 judges_total / 2 + 1
2879 } else {
2880 0
2881 };
2882 let met_quorum = !needs_quorum || present >= quorum;
2883 let uncontested = (!needs_quorum).then(|| {
2884 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2885 });
2886
2887 self.state.event(
2888 "tally",
2889 match &uncontested {
2890 Some(reason) => format!("winner {winner} — {reason}"),
2891 None => format!(
2892 "winner {winner} — votes {} | initial {} | {} changed | \
2893 {present}/{judges_total} judges{}",
2894 first_choice
2895 .iter()
2896 .map(|(k, v)| format!("{k}:{v}"))
2897 .collect::<Vec<_>>()
2898 .join(" "),
2899 if unanimous_initial {
2900 "unanimous"
2901 } else {
2902 "split"
2903 },
2904 changed_votes,
2905 if met_quorum {
2906 String::new()
2907 } else {
2908 format!(" — below quorum ({quorum} required)")
2909 },
2910 ),
2911 },
2912 );
2913 if !met_quorum {
2914 self.state.event(
2915 "stall",
2916 format!(
2917 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2918 the run stops here, resumable"
2919 ),
2920 );
2921 }
2922 self.state.tally = Some(Tally {
2923 first_choice,
2924 borda,
2925 winner,
2926 rankings: tops.len(),
2927 unanimous_initial,
2928 deliberated,
2929 changed_votes,
2930 unanimous_final,
2931 tie_break,
2932 judges: judges_total,
2933 present,
2934 quorum,
2935 met_quorum,
2936 uncontested,
2937 });
2938 self.state.status = if met_quorum {
2939 RunStatus::Reviewing
2940 } else {
2941 RunStatus::Stalled
2942 };
2943 self.state.save()?;
2944 Ok(())
2945 }
2946
2947 #[allow(clippy::too_many_lines)]
2968 async fn recover_stall(&mut self) -> Result<bool> {
2969 let run_id = self.state.id.clone();
2974 let prompts = self.state.config.prompts.clone();
2975 let quota_seats: BTreeSet<&str> =
2980 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2981 let absent: Vec<String> = self
2982 .state
2983 .judgements
2984 .iter()
2985 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2986 .map(|j| j.seat.clone())
2987 .collect();
2988 if absent.is_empty() {
2989 return Ok(false);
2990 }
2991 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2992 if viable.len() <= 1 {
2993 return Ok(false);
2994 }
2995 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2996 let language = self.state.config.graph.language.clone();
2997 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2998 let sessions = self.state.config.graph.sessions;
2999 let artifacts = agent::artifacts_dir(&self.state.dir());
3000 let root = self.state.worktree_root();
3001 let base_short = short(&self.state.base_commit);
3002 let candidates: Vec<Candidate> = viable.clone();
3003
3004 let mut positions: Vec<usize> = absent
3006 .iter()
3007 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
3008 .collect();
3009 if positions.is_empty() {
3010 return Ok(false);
3011 }
3012 positions.sort_unstable();
3013 positions.dedup();
3014
3015 let mut judge_jobs = Vec::new();
3017 for &j in &positions {
3018 let order = blind::presentation_order(viable.len(), j, self.state.seed);
3019 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
3020 let seat_key = format!("judge-{}", j + 1);
3021 let spec = self.roles.judges[j].clone();
3022 let seat = self.seat(&seat_key, &spec.id);
3023 judge_jobs.push(SeatJob {
3024 spec,
3025 seat,
3026 prompt: prompt::judge(
3027 &self.state.instruction,
3028 &views,
3029 self.roles.judges.len(),
3030 &base_short,
3031 &language,
3032 ),
3033 cwd: root.join(seat_key),
3034 timeout,
3035 allow_write: false,
3036 sessions,
3037 artifacts: artifacts.clone(),
3038 stem: format!("judge-{}-recover", j + 1),
3039 });
3040 }
3041
3042 let labels_for_check = labels.clone();
3043 let mut judge_losses = Vec::new();
3044 let retries = self.state.config.graph.retries;
3045 let cache = self.state.config.cache_dir();
3046 let ctx = WaveCtx {
3047 run: &run_id,
3048 node: "judge",
3049 prompts: &prompts,
3050 cache: cache.as_deref(),
3051 round: None,
3052 };
3053 let results = ask_json_wave::<Ranking>(
3054 judge_jobs,
3055 Arc::clone(&self.sem),
3056 retries,
3057 &ctx,
3058 &mut judge_losses,
3059 &mut self.state,
3060 &move |r: &Ranking| r.validate(&labels_for_check),
3061 )
3062 .await;
3063
3064 let mut recovered: BTreeSet<usize> = BTreeSet::new();
3066 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
3067 self.state.seats.insert(seat.key.clone(), seat);
3068 let record = &mut self.state.judgements[j];
3069 match res {
3070 Ok((ranking, out)) => {
3071 record.ranking = ranking.normalized();
3072 record.reasons = ranking.reasons;
3073 record.confidence = ranking.confidence;
3074 record.failed = None;
3075 record.duration_ms = out.duration_ms;
3076 recovered.insert(j);
3077 self.state.event(
3078 "recover",
3079 format!("judge {} ranked again after the limit", j + 1),
3080 );
3081 }
3082 Err(e) => {
3083 self.state
3084 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
3085 }
3086 }
3087 }
3088
3089 let mut vote_jobs = Vec::new();
3091 let mut vote_pos: Vec<usize> = Vec::new();
3092 for &j in &recovered {
3093 let seat_key = format!("judge-{}", j + 1);
3094 let spec = self.roles.judges[j].clone();
3095 let seat = self.seat(&seat_key, &spec.id);
3096 let mut text = prompt::final_vote(&labels, &language);
3097 if !has_context(&spec, &seat, sessions) {
3098 text = format!(
3099 "{}\n\n# Candidates\n\n{}",
3100 text,
3101 self.candidate_block(&candidates, &base_short)
3102 );
3103 }
3104 vote_jobs.push(SeatJob {
3105 spec,
3106 seat,
3107 prompt: text,
3108 cwd: root.join(seat_key),
3109 timeout,
3110 allow_write: false,
3111 sessions,
3112 artifacts: artifacts.clone(),
3113 stem: format!("vote-judge-{}-recover", j + 1),
3114 });
3115 vote_pos.push(j);
3116 }
3117 let allowed = labels.clone();
3118 let mut vote_losses = Vec::new();
3119 let vote_retries = self.state.config.graph.retries;
3120 let vote_cache = self.state.config.cache_dir();
3121 let ctx = WaveCtx {
3122 run: &run_id,
3123 node: "vote",
3124 prompts: &prompts,
3125 cache: vote_cache.as_deref(),
3126 round: None,
3127 };
3128 let votes = ask_json_wave::<FinalVote>(
3129 vote_jobs,
3130 Arc::clone(&self.sem),
3131 vote_retries,
3132 &ctx,
3133 &mut vote_losses,
3134 &mut self.state,
3135 &move |v: &FinalVote| match v.label() {
3136 Some(c) if allowed.contains(&c) => Ok(()),
3137 other => bail!("vote {other:?} is not one of {allowed:?}"),
3138 },
3139 )
3140 .await;
3141 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
3142 let agent_id = seat.agent.clone();
3143 self.state.seats.insert(seat.key.clone(), seat);
3144 match res {
3145 Ok((v, _)) => {
3146 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
3147 rec.vote = v.label();
3148 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
3149 } else {
3150 self.state.votes.push(VoteRecord {
3151 judge: j + 1,
3152 agent: agent_id,
3153 vote: v.label(),
3154 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
3155 changed: false,
3156 });
3157 }
3158 self.state.event(
3159 "recover",
3160 format!("judge {} voted again after the limit", j + 1),
3161 );
3162 }
3163 Err(e) => {
3164 self.state
3165 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
3166 }
3167 }
3168 }
3169
3170 let recovered_keys: BTreeSet<String> = recovered
3174 .iter()
3175 .map(|&j| format!("judge-{}", j + 1))
3176 .collect();
3177 self.state
3178 .quota
3179 .retain(|q| !recovered_keys.contains(&q.seat));
3180 for loss in judge_losses.into_iter().chain(vote_losses) {
3184 if recovered_keys.contains(&loss.seat) {
3185 continue;
3186 }
3187 self.state.quota.retain(|q| q.seat != loss.seat);
3188 self.state.quota.push(loss);
3189 }
3190
3191 self.state.tally = None;
3193 self.tally()?;
3194 Ok(self
3195 .state
3196 .tally
3197 .as_ref()
3198 .map(|t| t.met_quorum)
3199 .unwrap_or(false))
3200 }
3201
3202 async fn fold_losers(&mut self) -> Result<()> {
3205 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
3206 return Ok(());
3207 };
3208 let repo = self.state.repo.clone();
3209 let mut folded = Vec::new();
3210 for i in 0..self.state.candidates.len() {
3211 let c = &self.state.candidates[i];
3212 if c.label == winner || c.folded {
3213 continue;
3214 }
3215 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
3216 git::worktree_remove(&repo, &wt).await.ok();
3217 git::branch_delete(&repo, &branch).await.ok();
3218 self.state.candidates[i].folded = true;
3219 folded.push(label.to_string());
3220 }
3221 let root = self.state.worktree_root();
3223 for j in 1..=self.roles.judges.len() {
3224 let wt = root.join(format!("judge-{j}"));
3225 if wt.exists() {
3226 git::worktree_remove(&repo, &wt).await.ok();
3227 }
3228 }
3229 if self.state.config.graph.advise {
3232 for k in 1..=self.state.config.graph.advisors {
3233 let wt = root.join(format!("advisor-{k}"));
3234 if wt.exists() {
3235 git::worktree_remove(&repo, &wt).await.ok();
3236 }
3237 }
3238 }
3239 if !folded.is_empty() {
3240 self.state
3241 .event("fold", format!("folded candidates {}", folded.join(", ")));
3242 self.state.save()?;
3243 }
3244 Ok(())
3245 }
3246
3247 async fn sync_to_base(&mut self) -> Result<()> {
3287 if self.state.status == RunStatus::AlreadyInBase {
3288 return Ok(());
3289 }
3290 let conflicted = self
3291 .state
3292 .base_sync
3293 .as_ref()
3294 .is_some_and(|s| s.conflict.is_some());
3295 let Some(winner) = self.state.winner().cloned() else {
3296 return Ok(());
3297 };
3298
3299 let repo = self.state.repo.clone();
3300 let remote = self.state.config.merge.remote.clone();
3301 let base_branch = self.state.base_branch.clone();
3302 let tracking = format!("{remote}/{base_branch}");
3303
3304 git::fetch(&repo, &remote, &base_branch).await.ok();
3305 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3309 return Ok(());
3310 };
3311
3312 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3313 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3314 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3315
3316 if (behind > 0 || conflicted || head == tip)
3325 && self
3326 .settle_already_in(&winner.branch, &tip, &head, attempts, behind)
3327 .await?
3328 {
3329 return Ok(());
3330 }
3331 if conflicted {
3332 return Ok(());
3333 }
3334
3335 if behind == 0 {
3336 if let Some(from) = self
3345 .state
3346 .rebase_fixes
3347 .iter()
3348 .rev()
3349 .find_map(|r| r.from.clone())
3350 && from != head
3351 && git::git_raw(&winner.worktree, &["diff", "--quiet", &from])
3352 .await
3353 .is_ok_and(|o| o.ok())
3354 {
3355 git::sync_to_head(&winner.worktree).await?;
3356 }
3357 self.state.base_sync = Some(BaseSync {
3358 tip,
3359 behind: 0,
3360 attempts,
3361 conflict: None,
3362 already_in: None,
3363 });
3364 self.state.save()?;
3365 return Ok(());
3366 }
3367
3368 if attempts >= BASE_SYNC_ROUNDS {
3369 let why = format!(
3370 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3371 rebase(s); rebasing again would only race it",
3372 winner.branch
3373 );
3374 self.state.status = RunStatus::Blocked;
3375 self.state.base_sync = Some(BaseSync {
3376 tip,
3377 behind,
3378 attempts,
3379 conflict: Some(why.clone()),
3380 already_in: None,
3381 });
3382 self.state.event("land", why);
3383 self.state.save()?;
3384 return Ok(());
3385 }
3386
3387 self.state.event(
3388 "land",
3389 format!(
3390 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3391 winner.branch
3392 ),
3393 );
3394 self.state.save()?;
3395
3396 let branch_tracking = format!("{remote}/{}", winner.branch);
3402 let fetched_branch = git::fetch(&repo, &remote, &winner.branch).await;
3403 let remote_tip = if matches!(&fetched_branch, Ok(o) if o.ok()) {
3404 git::rev_parse(&repo, &branch_tracking).await.ok()
3405 } else {
3406 None
3407 };
3408 if let Some(theirs) = &remote_tip
3412 && !git::is_ancestor(&repo, theirs, &head).await
3413 && !crate::reconcile::origin_missing(&repo, &head, theirs)
3414 .await
3415 .is_ok_and(|missing| missing.is_empty())
3416 {
3417 let why = format!(
3418 "{branch_tracking} ({}) has commits {} does not contain; not rebasing over \
3419 them",
3420 short(theirs),
3421 winner.branch
3422 );
3423 self.state.status = RunStatus::Blocked;
3424 self.state.base_sync = Some(BaseSync {
3425 tip,
3426 behind,
3427 attempts,
3428 conflict: Some(why.clone()),
3429 already_in: None,
3430 });
3431 self.state.event("land", why);
3432 self.state.save()?;
3433 return Ok(());
3434 }
3435
3436 let scratch = self.state.dir().join("base-sync");
3437 let rebased = match crate::rebase::rebase_with_fixer(
3438 &mut self.state,
3439 &scratch,
3440 &winner.branch,
3441 &tracking,
3442 )
3443 .await
3444 {
3445 Ok(crate::rebase::Rebased::Applied) => Ok(None),
3446 Ok(crate::rebase::Rebased::Stopped(why)) => Ok(Some(why)),
3447 Err(e) => Err(e),
3448 };
3449 let attempts = attempts + 1;
3450 match rebased {
3451 Ok(None) => {
3452 git::sync_to_head(&winner.worktree).await?;
3456 let mut conflict = None;
3457 if let Some(pinned) = &remote_tip {
3458 let pushed = git::push_pinned(&repo, &remote, &winner.branch, pinned).await;
3459 match pushed {
3460 Ok(o) if o.ok() => self.state.event(
3461 "land",
3462 format!("pushed rebased {} to {remote}", winner.branch),
3463 ),
3464 Ok(o) => {
3465 conflict = Some(format!(
3466 "rebased {} locally but {remote} refused the push (it moved since {}; someone may have pushed): {}",
3467 winner.branch,
3468 short(pinned),
3469 o.stderr.chars().take(600).collect::<String>()
3470 ));
3471 }
3472 Err(e) => {
3473 conflict = Some(format!(
3474 "rebased {} locally but could not push it: {e:#}",
3475 winner.branch
3476 ));
3477 }
3478 }
3479 }
3480 if let Some(why) = &conflict {
3481 self.state.status = RunStatus::Blocked;
3482 self.state.event("land", why.clone());
3483 }
3484 self.state.base_sync = Some(BaseSync {
3485 tip: tip.clone(),
3486 behind: 0,
3487 attempts,
3488 conflict,
3489 already_in: None,
3490 });
3491 self.state
3492 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3493 }
3494 Ok(Some(conflict)) => {
3495 let why = format!(
3496 "{} conflicts with {tracking} and did not rebase: {}",
3497 winner.branch,
3498 conflict.chars().take(600).collect::<String>()
3499 );
3500 self.state.status = RunStatus::Blocked;
3501 self.state.base_sync = Some(BaseSync {
3502 tip,
3503 behind,
3504 attempts,
3505 conflict: Some(why.clone()),
3506 already_in: None,
3507 });
3508 self.state.event("land", why);
3509 }
3510 Err(e) => {
3511 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3512 self.state.status = RunStatus::Blocked;
3513 self.state.base_sync = Some(BaseSync {
3514 tip,
3515 behind,
3516 attempts,
3517 conflict: Some(why.clone()),
3518 already_in: None,
3519 });
3520 self.state.event("land", why);
3521 }
3522 }
3523 self.state.save()?;
3524 Ok(())
3525 }
3526
3527 async fn settle_already_in(
3544 &mut self,
3545 branch: &str,
3546 tip: &str,
3547 head: &str,
3548 attempts: usize,
3549 behind: usize,
3550 ) -> Result<bool> {
3551 let repo = self.state.repo.clone();
3552 let remote = self.state.config.merge.remote.clone();
3553 let start = self.state.base_commit.clone();
3554 let evidence = match crate::already::classify(&repo, tip, head, Some(&start)).await {
3555 Ok(Some(e)) => e,
3556 Ok(None) => return Ok(false),
3557 Err(e) => {
3558 tracing::warn!("already-in-base check for {branch}: {e:#}");
3559 return Ok(false);
3560 }
3561 };
3562 let mut verified = vec![head.to_owned()];
3563 let fetched = git::fetch(&repo, &remote, branch).await;
3564 if matches!(&fetched, Ok(o) if o.ok())
3565 && let Ok(theirs) = git::rev_parse(&repo, &format!("{remote}/{branch}")).await
3566 && theirs != head
3567 {
3568 match crate::already::classify(&repo, tip, &theirs, Some(&start)).await {
3569 Ok(Some(_)) => verified.push(theirs),
3570 _ => return Ok(false),
3571 }
3572 }
3573 let base_branch = self.state.base_branch.clone();
3574 let message = format!(
3575 "{branch} is already in {remote}/{base_branch} as {} ({} match); nothing left to \
3576 land",
3577 evidence.names(),
3578 evidence.proof.as_str()
3579 );
3580 let closed =
3581 crate::land::close_superseded_pr(&mut self.state, branch, &evidence, &verified).await;
3582 match closed {
3583 Ok(Ok(url)) => self
3584 .state
3585 .event("land", format!("closed {url}: superseded on {base_branch}")),
3586 Ok(Err(why)) => self
3587 .state
3588 .event("land", format!("did not close a pull request: {why}")),
3589 Err(e) => {
3590 let why = format!(
3591 "{branch} is already in {remote}/{base_branch}, but its pull request could \
3592 not be closed ({e:#}); resume to retry"
3593 );
3594 self.state.status = RunStatus::Blocked;
3595 self.state.base_sync = Some(BaseSync {
3596 tip: tip.to_owned(),
3597 behind,
3598 attempts,
3599 conflict: Some(why.clone()),
3600 already_in: None,
3601 });
3602 self.state.event("land", why);
3603 self.state.save()?;
3604 return Ok(true);
3605 }
3606 }
3607 self.state.status = RunStatus::AlreadyInBase;
3608 self.state.base_sync = Some(BaseSync {
3609 tip: tip.to_owned(),
3610 behind,
3611 attempts,
3612 conflict: None,
3613 already_in: Some(evidence),
3614 });
3615 self.state.event("land", message);
3616 self.state.save()?;
3617 self.settle_questions();
3618 Ok(true)
3619 }
3620
3621 fn landing_base(&self) -> String {
3631 self.state
3632 .base_sync
3633 .as_ref()
3634 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3635 }
3636
3637 pub async fn fix_selected(
3670 &mut self,
3671 ids: &[String],
3672 reason: &str,
3673 allow_stale: bool,
3674 ) -> Result<()> {
3675 let reason = reason.trim();
3676 if reason.is_empty() {
3677 bail!("a fix request needs a reason — that is the operator's own record of why");
3678 }
3679 if ids.is_empty() {
3680 bail!("no finding id given");
3681 }
3682 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3683 bail!(
3684 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3685 has already concluded — can be given a targeted fix. A run still \
3686 in progress should simply be resumed; a `merged` run's branch has \
3687 already landed, so its answer is a fresh `magi review <branch>`, \
3688 not reopening this run's own record",
3689 self.state.id,
3690 self.state.status.as_str()
3691 );
3692 }
3693 let Some(winner) = self.state.winner().cloned() else {
3694 bail!("run {} has no winning candidate to fix", self.state.id);
3695 };
3696 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3697 bail!(
3698 "branch `{}` no longer exists; this run cannot be extended",
3699 winner.branch
3700 );
3701 }
3702 let home = crate::run::home();
3703 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3704 bail!(
3705 "run {} is currently being worked on by another magi process",
3706 self.state.id
3707 );
3708 }
3709 let _claim = FixClaim::acquire(&self.state.dir())?;
3715
3716 let mut seen = BTreeSet::new();
3720 let mut findings = Vec::new();
3721 let mut missing = Vec::new();
3722 for id in ids {
3723 if !seen.insert(id.clone()) {
3724 continue;
3725 }
3726 match self.state.finding(id) {
3727 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3728 id: f.id.clone(),
3729 severity: f.severity,
3730 reviewer_vote: rec.vote,
3731 round: round.round,
3732 round_head: round.head.clone(),
3733 reviewer: rec.reviewer,
3734 agent: rec.agent.clone(),
3735 file: f.file.clone(),
3736 line: f.line,
3737 title: f.title.clone(),
3738 detail: f.detail.clone(),
3739 outcome: OperatorFixOutcome::Pending,
3740 }),
3741 None => missing.push(id.clone()),
3742 }
3743 }
3744 if !missing.is_empty() {
3745 bail!(
3746 "unknown finding id(s): {}; nothing was changed",
3747 missing.join(", ")
3748 );
3749 }
3750
3751 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3752 let stale_details: Vec<(String, String)> = findings
3753 .iter()
3754 .filter(|f| f.round_head != head_at_request)
3755 .map(|f| (f.id.clone(), f.round_head.clone()))
3756 .collect();
3757 let stale = !stale_details.is_empty();
3758 if stale && !allow_stale {
3759 bail!(
3760 "the branch has moved since some finding(s) were raised — {} — now \
3761 at {}; pass --allow-stale to fix anyway, or re-run review first",
3762 stale_details
3763 .iter()
3764 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3765 .collect::<Vec<_>>()
3766 .join(", "),
3767 short(&head_at_request)
3768 );
3769 }
3770
3771 let request = OperatorFixRequest {
3772 requested_at: Timestamp::now(),
3773 reason: reason.to_owned(),
3774 findings,
3775 head_at_request: head_at_request.clone(),
3776 allow_stale,
3777 stale,
3778 fix: None,
3779 result_head: None,
3780 follow_up_review_run: None,
3781 };
3782 self.state.event(
3783 "fix",
3784 format!(
3785 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3786 request.findings.len(),
3787 request
3788 .findings
3789 .iter()
3790 .map(|f| f.id.as_str())
3791 .collect::<Vec<_>>()
3792 .join(", "),
3793 ),
3794 );
3795 self.state.operator_fixes.push(request);
3802 self.state.save()?;
3803 let request_index = self.state.operator_fixes.len() - 1;
3804
3805 if winner.worktree.exists() {
3814 let dirty = git::git(
3817 &winner.worktree,
3818 &["status", "--porcelain", "--untracked-files=all"],
3819 )
3820 .await?;
3821 let only_withheld = dirty.lines().all(|l| {
3822 l.strip_prefix("?? ")
3823 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3824 });
3825 if !only_withheld {
3826 bail!(
3827 "`{}` has uncommitted changes; refusing to touch it — commit or \
3828 discard them first",
3829 winner.worktree.display()
3830 );
3831 }
3832 git::worktree_remove(&self.state.repo, &winner.worktree)
3833 .await
3834 .ok();
3835 }
3836 let fix_worktree = self.state.worktree_root().join("operator-fix");
3837 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3838 git::git(
3839 &self.state.repo,
3840 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3841 )
3842 .await
3843 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3844 if !git::is_clean(&fix_worktree).await? {
3845 git::worktree_remove(&self.state.repo, &fix_worktree)
3846 .await
3847 .ok();
3848 bail!(
3849 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3850 winner.branch
3851 );
3852 }
3853
3854 let run_id = self.state.id.clone();
3855 let prompts = self.state.config.prompts.clone();
3856 let language = self.state.config.graph.language.clone();
3857 let sessions = self.state.config.graph.sessions;
3858 let artifacts = agent::artifacts_dir(&self.state.dir());
3859 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3860 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3861 _ => (
3862 self.state
3863 .config
3864 .agent(&winner.agent)
3865 .cloned()
3866 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3867 format!("impl-{}", winner.label),
3868 ),
3869 };
3870 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3871 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3872 .findings
3873 .iter()
3874 .map(|f| Finding {
3875 id: f.id.clone(),
3876 severity: f.severity,
3877 file: f.file.clone(),
3878 line: f.line,
3879 title: f.title.clone(),
3880 detail: f.detail.clone(),
3881 })
3882 .collect();
3883 let job = SeatJob {
3884 prompt: prompt::operator_fix(
3885 &self.state.instruction,
3886 &finding_list,
3887 reason,
3888 &stale_details,
3889 &head_at_request,
3890 &language,
3891 ),
3892 spec: fix_spec.clone(),
3893 seat,
3894 cwd: fix_worktree.clone(),
3895 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3896 allow_write: true,
3897 sessions,
3898 artifacts: artifacts.clone(),
3899 stem: "operator-fix".to_owned(),
3900 };
3901 let cache = self.state.config.cache_dir();
3902 let ctx = WaveCtx {
3903 run: &run_id,
3904 node: "fix",
3905 prompts: &prompts,
3906 cache: cache.as_deref(),
3907 round: None,
3908 };
3909 let (seat, out) =
3910 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3911 let agent_id = seat.agent.clone();
3912
3913 let mut fix = FixRecord {
3914 agent: agent_id,
3915 addressed: Vec::new(),
3916 rejected: Vec::new(),
3917 notes: String::new(),
3918 committed: false,
3919 failed: None,
3920 duration_ms: 0,
3921 continuation: None,
3922 };
3923 let mut final_seat = seat.clone();
3924 match out {
3925 AgentOutcome::Ok(o) => {
3926 fix.duration_ms = o.duration_ms;
3927 let parsed = verdict::extract_json::<FixReport>(&o.text);
3928 let incomplete_reason = match &parsed {
3929 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3930 "the reply parsed, but it reported a command whose own CLI \
3931 never confirmed an exit status"
3932 .to_owned(),
3933 ),
3934 Ok(_) => None,
3935 Err(e) => Some(e.to_string()),
3936 };
3937 match incomplete_reason {
3938 None => {
3939 let report = parsed.expect("checked Ok above");
3940 fix.addressed = report.addressed;
3941 fix.rejected = report.rejected;
3942 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3943 }
3944 Some(reason) => {
3945 let (resumed_seat, resolved, failure, cont) = self
3946 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3947 .await;
3948 fix.duration_ms += cont.cumulative_wait_ms;
3949 fix.continuation = Some(cont);
3950 final_seat = resumed_seat;
3951 match resolved {
3952 Some(report) => {
3953 fix.addressed = report.addressed;
3954 fix.rejected = report.rejected;
3955 fix.notes =
3956 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3957 }
3958 None => fix.failed = failure,
3959 }
3960 }
3961 }
3962 }
3963 AgentOutcome::Dropped(o) => {
3964 fix.duration_ms = o.duration_ms;
3965 let why = o
3966 .dropped
3967 .as_ref()
3968 .map(|d| d.why.as_str())
3969 .unwrap_or("the CLI ended the stream without delivering its answer");
3970 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3971 }
3972 AgentOutcome::Quota(o) => {
3973 self.state.quota.push(QuotaLoss {
3974 seat: final_seat.key.clone(),
3975 node: "fix".to_owned(),
3976 at: Timestamp::now(),
3977 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3978 });
3979 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3980 }
3981 AgentOutcome::Failed(e) => fix.failed = Some(e),
3982 }
3983 if fix.continuation.is_none() {
3984 fix.continuation = Some(ContinuationRecord::not_needed());
3985 }
3986 self.state.seats.insert(final_seat.key.clone(), final_seat);
3987
3988 let rescue_message = format!(
3989 "magi: operator-selected fix ({}) (uncommitted work)",
3990 self.state.operator_fixes[request_index]
3991 .findings
3992 .iter()
3993 .map(|f| f.id.as_str())
3994 .collect::<Vec<_>>()
3995 .join(", ")
3996 );
3997 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3998 self.state.note_withheld("fix", &r.withheld);
3999 }
4000 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
4001 fix.committed = after != head_at_request;
4002 git::worktree_remove(&self.state.repo, &fix_worktree)
4003 .await
4004 .ok();
4005
4006 self.state.event(
4007 "fix",
4008 match &fix.failed {
4009 Some(reason) => format!(
4010 "operator fix: adoption report was lost ({reason}); {}",
4011 if fix.committed {
4012 "committed"
4013 } else {
4014 "NO new commit"
4015 }
4016 ),
4017 None => format!(
4018 "operator fix: {} addressed, {} rejected, {}",
4019 fix.addressed.len(),
4020 fix.rejected.len(),
4021 if fix.committed {
4022 "committed"
4023 } else {
4024 "NO new commit"
4025 }
4026 ),
4027 },
4028 );
4029
4030 for f in &mut self.state.operator_fixes[request_index].findings {
4037 f.outcome = if fix.failed.is_some() {
4038 OperatorFixOutcome::Unreported
4039 } else if fix.addressed.contains(&f.id) {
4040 OperatorFixOutcome::Addressed
4041 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
4042 OperatorFixOutcome::Rejected { why: r.why.clone() }
4043 } else {
4044 OperatorFixOutcome::Unreported
4045 };
4046 }
4047
4048 let committed = fix.committed;
4049 if committed {
4050 self.state.operator_fixes[request_index].result_head = Some(after.clone());
4051 }
4052 self.state.operator_fixes[request_index].fix = Some(fix);
4053 self.state.save()?;
4056
4057 if committed {
4058 self.state.event(
4059 "fix",
4060 format!(
4061 "operator fix committed {}; opening a follow-up review-only run",
4062 short(&after)
4063 ),
4064 );
4065 let origin =
4068 Origin::operator().serving(self.state.origin.as_ref().and_then(|o| o.task.clone()));
4069 match Self::review(
4070 &self.state.repo,
4071 &winner.branch,
4072 self.state.config.clone(),
4073 origin,
4074 )
4075 .await
4076 {
4077 Ok(mut follow_up) => {
4078 follow_up.state.event(
4079 "start",
4080 format!(
4081 "requested by an operator fix on run {} for finding(s) {}",
4082 self.state.id,
4083 self.state.operator_fixes[request_index]
4084 .findings
4085 .iter()
4086 .map(|f| f.id.as_str())
4087 .collect::<Vec<_>>()
4088 .join(", "),
4089 ),
4090 );
4091 follow_up.state.save()?;
4092 let follow_up_id = follow_up.state.id.clone();
4093 if let Err(e) = follow_up.execute().await {
4094 self.state.event(
4095 "fix",
4096 format!(
4097 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
4098 ),
4099 );
4100 }
4101 self.state.operator_fixes[request_index].follow_up_review_run =
4102 Some(follow_up_id);
4103 }
4104 Err(e) => {
4105 self.state.event(
4106 "fix",
4107 format!("committed the fix but could not open a follow-up review: {e:#}"),
4108 );
4109 }
4110 }
4111 self.state.save()?;
4112 }
4113
4114 Ok(())
4115 }
4116
4117 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
4124 match &self.roles.fixer {
4125 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
4126 _ => (
4127 self.state
4128 .config
4129 .agent(&winner.agent)
4130 .cloned()
4131 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
4132 format!("impl-{}", winner.label),
4133 ),
4134 }
4135 }
4136
4137 async fn review_loop(&mut self) -> Result<()> {
4138 if self
4143 .state
4144 .base_sync
4145 .as_ref()
4146 .is_some_and(|s| s.conflict.is_some())
4147 {
4148 return Ok(());
4149 }
4150 let run_id = self.state.id.clone();
4155 let prompts = self.state.config.prompts.clone();
4156 let Some(winner) = self.state.winner().cloned() else {
4157 return Ok(());
4158 };
4159 let max_rounds = self.state.config.graph.review_rounds;
4160 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
4170 self.state.status = status;
4171 self.state.save()?;
4172 return Ok(());
4173 }
4174 self.state.status = RunStatus::Reviewing;
4175 if self
4185 .state
4186 .reviews
4187 .last()
4188 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
4189 {
4190 let shell = self.state.config.shell();
4191 return self
4192 .stop_reviewing(
4193 "the last round's own verification never resolved",
4194 &shell,
4195 &winner.worktree,
4196 )
4197 .await;
4198 }
4199
4200 let repo = self.state.repo.clone();
4201 let root = self.state.worktree_root();
4202 let language = self.state.config.graph.language.clone();
4203 let sessions = self.state.config.graph.sessions;
4204 let artifacts = agent::artifacts_dir(&self.state.dir());
4205 let base = self.landing_base();
4206 let base_short = short(&base);
4207 let reviewers = self.roles.reviewers.clone();
4208 let shell = self.state.config.shell();
4209
4210 for round in (self.state.reviews.len() + 1)..=max_rounds {
4211 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4212 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
4213 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
4214 let prev_verification = self
4223 .state
4224 .reviews
4225 .last()
4226 .and_then(|r| r.verification_summary(&head));
4227
4228 let mut jobs = Vec::new();
4232 for (r, spec) in reviewers.iter().cloned().enumerate() {
4233 let wt = root.join(format!("review-{}", r + 1));
4234 if wt.exists() {
4235 git::reset_detached(&wt, &head).await?;
4236 } else {
4237 git::worktree_add_detached(&repo, &wt, &head).await?;
4238 }
4239 let seat_key = format!("review-{}", r + 1);
4240 let seat = self.seat(&seat_key, &spec.id);
4241 jobs.push(SeatJob {
4242 prompt: prompt::review(&prompt::ReviewCtx {
4243 instruction: &self.state.instruction,
4244 branch: &winner.branch,
4245 base_short: &base_short,
4246 stat: &stat,
4247 patch: &patch,
4248 verification: prev_verification.as_ref(),
4249 reviewers: reviewers.len(),
4250 round,
4251 rounds: max_rounds,
4252 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
4255 lens: Lens::for_seat(r),
4256 language: &language,
4257 }),
4258 spec,
4259 seat,
4260 cwd: wt,
4261 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4262 allow_write: false,
4263 sessions,
4264 artifacts: artifacts.clone(),
4265 stem: format!("review-{round}-{}", r + 1),
4266 });
4267 }
4268
4269 self.state.event(
4270 "review",
4271 format!(
4272 "round {round}: {} reviewers on {}",
4273 jobs.len(),
4274 short(&head)
4275 ),
4276 );
4277 let mut quota_losses = Vec::new();
4278 let review_retries = self.state.config.graph.retries;
4279 let review_cache = self.state.config.cache_dir();
4280 let ctx = WaveCtx {
4281 run: &run_id,
4282 node: "review",
4283 prompts: &prompts,
4284 cache: review_cache.as_deref(),
4285 round: Some(round),
4286 };
4287 let results = ask_json_wave::<Review>(
4288 jobs,
4289 Arc::clone(&self.sem),
4290 review_retries,
4291 &ctx,
4292 &mut quota_losses,
4293 &mut self.state,
4294 &|_: &Review| Ok(()),
4295 )
4296 .await;
4297 let round_quota_missing = quota_losses.len();
4301 self.state.quota.extend(quota_losses);
4302
4303 let mut records = Vec::new();
4304 let mut all_findings = Vec::new();
4305 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
4306 let agent_id = seat.agent.clone();
4307 self.state.seats.insert(seat.key.clone(), seat);
4308 let mut record = ReviewRecord {
4309 reviewer: r + 1,
4310 agent: agent_id,
4311 summary: String::new(),
4312 findings: Vec::new(),
4313 vote: None,
4314 failed: None,
4315 duration_ms: 0,
4316 attempts,
4322 };
4323 match res {
4324 Ok((review, out)) => {
4325 record.summary =
4333 blind::sanitize_prose(&review.summary, &self.state.config.blind);
4334 record.vote = Some(review.vote);
4335 record.duration_ms = out.duration_ms;
4336 for (n, mut f) in review.findings.into_iter().enumerate() {
4337 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
4340 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
4341 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
4342 f.file = f
4348 .file
4349 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
4350 all_findings.push(f.clone());
4351 record.findings.push(f);
4352 }
4353 self.state.event(
4354 "review",
4355 format!(
4356 "round {round}: reviewer {} voted {} with {} finding(s)",
4357 r + 1,
4358 review.vote.label(),
4359 record.findings.len()
4360 ),
4361 );
4362 }
4363 Err(e) => {
4364 record.failed = Some(e.to_string());
4365 self.state.event(
4366 "review",
4367 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
4368 );
4369 }
4370 }
4371 records.push(record);
4372 }
4373
4374 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
4381 let vote_split =
4382 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
4383 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
4384 if vote_split {
4385 self.state.event(
4386 "review",
4387 format!(
4388 "round {round}: votes split ({}) — one round of reconsideration",
4389 initial_votes
4390 .iter()
4391 .map(|v| v.label())
4392 .collect::<Vec<_>>()
4393 .join(", ")
4394 ),
4395 );
4396 let panel: Vec<ReviewSeatReport<'_>> = records
4399 .iter()
4400 .filter_map(|r| {
4401 r.vote.map(|vote| ReviewSeatReport {
4402 reviewer: r.reviewer,
4403 vote,
4404 summary: &r.summary,
4405 findings: &r.findings,
4406 })
4407 })
4408 .collect();
4409
4410 let mut jobs = Vec::new();
4411 let mut seats_at = Vec::new();
4412 for (r, spec) in reviewers.iter().cloned().enumerate() {
4413 if records[r].vote.is_none() {
4417 continue;
4418 }
4419 let wt = root.join(format!("review-{}", r + 1));
4420 let seat_key = format!("review-{}", r + 1);
4421 let seat = self.seat(&seat_key, &spec.id);
4422 let patch_ctx = if has_context(&spec, &seat, sessions) {
4427 None
4428 } else {
4429 Some(ReviewPatch {
4430 branch: &winner.branch,
4431 base_short: &base_short,
4432 stat: &stat,
4433 patch: &patch,
4434 })
4435 };
4436 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
4437 instruction: &self.state.instruction,
4438 reviewer: r + 1,
4439 lens: Lens::for_seat(r),
4440 panel: &panel,
4441 patch: patch_ctx,
4442 round,
4443 rounds: max_rounds,
4444 language: &language,
4445 });
4446 jobs.push(SeatJob {
4447 prompt,
4448 spec,
4449 seat,
4450 cwd: wt,
4451 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4452 allow_write: false,
4453 sessions,
4454 artifacts: artifacts.clone(),
4455 stem: format!("review-{round}-reconsider-{}", r + 1),
4456 });
4457 seats_at.push(r);
4458 }
4459
4460 let mut recon_quota_losses = Vec::new();
4461 let recon_cache = self.state.config.cache_dir();
4462 let recon_ctx = WaveCtx {
4463 run: &run_id,
4464 node: "review",
4465 prompts: &prompts,
4466 cache: recon_cache.as_deref(),
4467 round: Some(round),
4468 };
4469 let recon_results = ask_json_wave::<ReviewRevote>(
4470 jobs,
4471 Arc::clone(&self.sem),
4472 review_retries,
4473 &recon_ctx,
4474 &mut recon_quota_losses,
4475 &mut self.state,
4476 &|_: &ReviewRevote| Ok(()),
4477 )
4478 .await;
4479 self.state.quota.extend(recon_quota_losses);
4480
4481 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
4482 let agent_id = seat.agent.clone();
4483 self.state.seats.insert(seat.key.clone(), seat);
4484 let mut rec = ReviewRevoteRecord {
4485 reviewer: r + 1,
4486 agent: agent_id,
4487 vote: None,
4488 reason: String::new(),
4489 failed: None,
4490 };
4491 match res {
4492 Ok((rv, _)) => {
4493 rec.vote = Some(rv.vote);
4494 rec.reason =
4495 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
4496 self.state.event(
4497 "review",
4498 format!(
4499 "round {round}: reviewer {} revoted {}",
4500 r + 1,
4501 rv.vote.label()
4502 ),
4503 );
4504 }
4505 Err(e) => {
4506 rec.failed = Some(e.to_string());
4507 self.state.event(
4508 "review",
4509 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4510 );
4511 }
4512 }
4513 reconsideration.push(rec);
4514 }
4515 } else if initial_votes.len() > 1 {
4516 self.state.event(
4517 "review",
4518 format!(
4519 "round {round}: votes agreed ({}) — no reconsideration",
4520 initial_votes[0].label()
4521 ),
4522 );
4523 }
4524
4525 let final_votes: Vec<ReviewVote> = records
4529 .iter()
4530 .filter_map(|r| {
4531 reconsideration
4532 .iter()
4533 .find(|rv| rv.reviewer == r.reviewer)
4534 .and_then(|rv| rv.vote)
4535 .or(r.vote)
4536 })
4537 .collect();
4538 let round_verdict = ReviewVote::worst(final_votes);
4539
4540 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4541 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4542 let defer_e2e =
4553 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4554 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4555 let reason =
4556 format!("{blocking} blocking finding(s) already required a fix this round");
4557 self.state.event(
4558 "verify",
4559 format!(
4560 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4561 {}); it will run once a round has none left",
4562 short(&head)
4563 ),
4564 );
4565 (Vec::new(), false, true, Some(reason))
4566 } else {
4567 let e2e_commands = self.state.config.verify.e2e.clone();
4568 let cache_dir = self.state.config.cache_dir();
4569 let context = format!("round {round}");
4570 let (e2e, verify_retried) = with_cache_lease(
4571 &mut self.state,
4572 cache_dir.as_deref(),
4573 "e2e",
4574 "e2e",
4575 &winner.worktree,
4576 &head,
4577 verify_timeout,
4578 &context,
4579 |state, budget| {
4580 let shell = shell.clone();
4581 let e2e_commands = e2e_commands.clone();
4582 let worktree = winner.worktree.clone();
4583 let context = context.clone();
4584 async move {
4585 run_e2e_with_retry(
4586 state,
4587 &shell,
4588 &e2e_commands,
4589 &worktree,
4590 budget,
4591 &context,
4592 )
4593 .await
4594 }
4595 },
4596 )
4597 .await;
4598 (e2e, verify_retried, false, None)
4599 };
4600
4601 let expected = records.len();
4602 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4603 let incomplete = answered < expected;
4604 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4605 let policy = self.state.config.graph.incomplete_review;
4606 let clean = round_is_clean(
4607 blocking,
4608 e2e_ok,
4609 answered,
4610 expected,
4611 round_quota_missing,
4612 policy,
4613 );
4614
4615 let mut round_record = ReviewRound {
4616 round,
4617 head: head.clone(),
4618 verified_head: None,
4619 verified_at: None,
4620 reviews: records,
4621 e2e,
4622 verify_retried,
4623 e2e_deferred,
4624 e2e_defer_reason,
4625 fix: None,
4626 blocking,
4627 answered,
4628 expected,
4629 clean,
4630 progressed: false,
4631 vote_split,
4632 reconsideration,
4633 verdict: round_verdict,
4634 };
4635 if !matches!(
4646 round_record.e2e_status(),
4647 E2eStatus::Deferred | E2eStatus::NotConfigured
4648 ) {
4649 round_record.verified_head = Some(head.clone());
4650 round_record.verified_at = Some(Timestamp::now());
4651 }
4652 let this_round_verification = round_record.verification_summary(&head);
4653
4654 if incomplete {
4655 let missing: Vec<String> = round_record
4656 .reviews
4657 .iter()
4658 .filter(|r| r.failed.is_some())
4659 .map(|r| format!("review-{}", r.reviewer))
4660 .collect();
4661 self.state.event(
4662 "review",
4663 format!(
4664 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4665 missing.join(", ")
4666 ),
4667 );
4668 }
4669
4670 if clean {
4671 self.state.event(
4672 "review",
4673 if incomplete && policy == IncompleteReviewPolicy::Warn {
4674 format!(
4675 "round {round}: clean (warn policy, incomplete panel) — no \
4676 blocking findings from the seats that answered, verification green"
4677 )
4678 } else if incomplete {
4679 format!(
4680 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4681 quorum) — no blocking findings from the seats that answered, \
4682 verification green",
4683 expected - answered
4684 )
4685 } else {
4686 format!("round {round}: clean — no blocking findings, verification green")
4687 },
4688 );
4689 self.state.reviews.push(round_record);
4690 self.state.status = RunStatus::Gating;
4691 self.state.save()?;
4692 return Ok(());
4693 }
4694
4695 if incomplete && blocking == 0 && e2e_ok {
4703 self.state.reviews.push(round_record);
4704 self.state.save()?;
4705 if round == max_rounds {
4706 self.state.status = RunStatus::Blocked;
4707 self.state.event(
4708 "review",
4709 format!(
4710 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4711 refusing to call it clean",
4712 expected - answered
4713 ),
4714 );
4715 return Ok(());
4716 }
4717 continue;
4718 }
4719
4720 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4732 self.state.reviews.push(round_record);
4733 return self
4734 .stop_reviewing(
4735 "the round's own verification could not run",
4736 &shell,
4737 &winner.worktree,
4738 )
4739 .await;
4740 }
4741
4742 if round == max_rounds {
4743 self.state.reviews.push(round_record);
4744 return self
4745 .stop_reviewing(
4746 &format!(
4747 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4748 ),
4749 &shell,
4750 &winner.worktree,
4751 )
4752 .await;
4753 }
4754
4755 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4758 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4759 let blocking_findings: Vec<_> = all_findings
4760 .iter()
4761 .filter(|f| f.severity.blocks())
4762 .cloned()
4763 .collect();
4764 let job = SeatJob {
4765 prompt: prompt::fix(
4766 &self.state.instruction,
4767 &blocking_findings,
4768 this_round_verification.as_ref(),
4769 round,
4770 max_rounds,
4771 &language,
4772 ),
4773 spec: fix_spec.clone(),
4774 seat,
4775 cwd: winner.worktree.clone(),
4776 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4777 allow_write: true,
4778 sessions,
4779 artifacts: artifacts.clone(),
4780 stem: format!("fix-{round}"),
4781 };
4782 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4783 let cache = self.state.config.cache_dir();
4784 let ctx = WaveCtx {
4785 run: &run_id,
4786 node: "fix",
4787 prompts: &prompts,
4788 cache: cache.as_deref(),
4789 round: Some(round),
4790 };
4791 let (seat, out) =
4792 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4793 let agent_id = seat.agent.clone();
4794
4795 let mut fix = FixRecord {
4796 agent: agent_id,
4797 addressed: Vec::new(),
4798 rejected: Vec::new(),
4799 notes: String::new(),
4800 committed: false,
4801 failed: None,
4802 duration_ms: 0,
4803 continuation: None,
4804 };
4805 let mut continuation = ContinuationRecord::not_needed();
4806 let mut final_seat = seat.clone();
4807 match out {
4808 AgentOutcome::Ok(o) => {
4809 fix.duration_ms = o.duration_ms;
4810 let parsed = verdict::extract_json::<FixReport>(&o.text);
4811 let incomplete_reason = match &parsed {
4818 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4819 "the reply parsed, but it reported a command whose own CLI never \
4820 confirmed an exit status"
4821 .to_owned(),
4822 ),
4823 Ok(_) => None,
4824 Err(e) => Some(e.to_string()),
4825 };
4826 match incomplete_reason {
4827 None => {
4828 let report = parsed.expect("checked Ok above");
4829 fix.addressed = report.addressed;
4830 fix.rejected = report.rejected;
4831 fix.notes =
4832 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4833 }
4834 Some(reason) => {
4835 let (resumed_seat, resolved, failure, cont) = self
4836 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4837 .await;
4838 fix.duration_ms += cont.cumulative_wait_ms;
4839 continuation = cont;
4840 final_seat = resumed_seat;
4841 match resolved {
4842 Some(report) => {
4843 fix.addressed = report.addressed;
4844 fix.rejected = report.rejected;
4845 fix.notes = blind::sanitize_prose(
4846 &report.notes,
4847 &self.state.config.blind,
4848 );
4849 }
4850 None => fix.failed = failure,
4851 }
4852 }
4853 }
4854 }
4855 AgentOutcome::Dropped(o) => {
4857 fix.duration_ms = o.duration_ms;
4858 let why = o
4859 .dropped
4860 .as_ref()
4861 .map(|d| d.why.as_str())
4862 .unwrap_or("the CLI ended the stream without delivering its answer");
4863 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4864 }
4865 AgentOutcome::Quota(o) => {
4866 self.state.quota.push(QuotaLoss {
4867 seat: final_seat.key.clone(),
4868 node: "fix".to_owned(),
4869 at: Timestamp::now(),
4870 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4871 });
4872 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4873 }
4874 AgentOutcome::Failed(e) => fix.failed = Some(e),
4875 }
4876 fix.continuation = Some(continuation);
4877 self.state.seats.insert(final_seat.key.clone(), final_seat);
4878 if let Ok(r) = git::rescue_commit(
4879 &winner.worktree,
4880 &format!("magi: review round {round} fixes (uncommitted work)"),
4881 )
4882 .await
4883 {
4884 self.state.note_withheld("fix", &r.withheld);
4885 }
4886 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4887 fix.committed = after != before;
4888 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4896 let progressed = diff_after != patch;
4897 let commit_note = if fix.committed {
4898 "committed"
4899 } else {
4900 "NO new commit"
4901 };
4902 let tree_note = if progressed {
4903 "changed vs base"
4904 } else {
4905 "unchanged vs base"
4906 };
4907 self.state.event(
4908 "fix",
4909 match &fix.failed {
4910 Some(reason) => {
4916 format!(
4917 "round {round}: fixer's adoption report was lost ({reason}); \
4918 {commit_note}, tree {tree_note}"
4919 )
4920 }
4921 None => format!(
4922 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4923 {tree_note}{}",
4924 fix.addressed.len(),
4925 fix.rejected.len(),
4926 if continuation.outcome == ContinuationOutcome::Resumed {
4927 format!(
4928 " (adoption report recovered after {} continuation(s))",
4929 continuation.attempts
4930 )
4931 } else {
4932 String::new()
4933 },
4934 ),
4935 },
4936 );
4937 round_record.fix = Some(fix);
4938 round_record.progressed = progressed;
4939 self.state.reviews.push(round_record);
4940 self.state.save()?;
4941
4942 if matches!(
4955 continuation.outcome,
4956 ContinuationOutcome::Exhausted
4957 | ContinuationOutcome::QuotaLost
4958 | ContinuationOutcome::NoSession
4959 ) {
4960 return self
4961 .stop_reviewing(
4962 "the fixer's adoption report never came back, even after resuming its \
4963 own seat; refusing to start another round against the same worktree \
4964 while that is unresolved",
4965 &shell,
4966 &winner.worktree,
4967 )
4968 .await;
4969 }
4970
4971 let streak = self
4972 .state
4973 .reviews
4974 .iter()
4975 .rev()
4976 .take_while(|r| !r.progressed)
4977 .count();
4978 if streak >= STAGNANT_LIMIT {
4979 return self
4980 .stop_reviewing(
4981 &format!(
4982 "the tree has not moved against base for {streak} round(s) in a row"
4983 ),
4984 &shell,
4985 &winner.worktree,
4986 )
4987 .await;
4988 }
4989 }
4990 Ok(())
4991 }
4992
4993 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
5023 let round_idx = self.state.reviews.len() - 1;
5024 let needs_catchup_run = matches!(
5032 self.state.reviews[round_idx].e2e_status(),
5033 E2eStatus::Deferred | E2eStatus::ResourceBlocked
5034 );
5035 if needs_catchup_run {
5036 let round = self.state.reviews[round_idx].round;
5037 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5038 let commands = self.state.config.verify.e2e.clone();
5039 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
5040 let cache_dir = self.state.config.cache_dir();
5041 let context = format!(
5042 "round {round}: verification unresolved, catching up before the final decision"
5043 );
5044 let (outcomes, verify_retried) = with_cache_lease(
5045 &mut self.state,
5046 cache_dir.as_deref(),
5047 "e2e",
5048 "e2e",
5049 worktree,
5050 &attempted_head,
5051 timeout,
5052 &context,
5053 |state, budget| {
5054 let shell = shell.to_vec();
5055 let commands = commands.clone();
5056 let context = context.clone();
5057 async move {
5058 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
5059 .await
5060 }
5061 },
5062 )
5063 .await;
5064 let last = &mut self.state.reviews[round_idx];
5065 last.e2e = outcomes;
5066 last.verify_retried = verify_retried;
5067 last.verified_head = Some(attempted_head);
5074 last.verified_at = Some(Timestamp::now());
5075 if verify_inconclusive(&last.e2e) {
5076 self.state.save()?;
5083 return Ok(());
5084 }
5085 last.e2e_deferred = false;
5086 }
5087 let last = &self.state.reviews[round_idx];
5088 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
5089
5090 match last.e2e_status() {
5091 E2eStatus::Failed => {
5092 let red: Vec<String> = last
5093 .e2e
5094 .iter()
5095 .filter(|o| !o.ok())
5096 .map(|o| {
5097 format!(
5098 "`{}` -> {:?}\n{}",
5099 o.command,
5100 o.code,
5101 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5102 )
5103 })
5104 .collect();
5105 self.state
5106 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
5107 self.state.status = RunStatus::Blocked;
5108 }
5109 E2eStatus::ResourceBlocked => {
5114 self.state.event(
5115 "review",
5116 format!(
5117 "{why}; e2e could not run (shared build cache unavailable); not \
5118 deciding yet"
5119 ),
5120 );
5121 }
5122 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
5123 self.state.event(
5124 "review",
5125 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
5126 );
5127 self.state.status = RunStatus::Gating;
5128 }
5129 }
5130 self.state.save()?;
5131 Ok(())
5132 }
5133
5134 async fn gate(&mut self) -> Result<()> {
5137 if self.state.status == RunStatus::Failed
5149 || self
5150 .state
5151 .base_sync
5152 .as_ref()
5153 .is_some_and(|s| s.conflict.is_some())
5154 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5155 != Some(RunStatus::Gating)
5156 {
5157 return Ok(());
5158 }
5159 if self.state.gate_ran {
5160 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
5171 self.state.status = RunStatus::Blocked;
5172 self.state.save()?;
5173 }
5174 return Ok(());
5175 }
5176 let Some(winner) = self.state.winner().cloned() else {
5177 return Ok(());
5178 };
5179 self.state.status = RunStatus::Gating;
5180 let mut outcomes = self.run_gate(&winner).await?;
5181 loop {
5182 if verify_inconclusive(&outcomes) {
5193 self.state.save()?;
5194 return Ok(());
5195 }
5196 if outcomes.iter().all(CommandOutcome::ok) {
5197 break;
5198 }
5199 match self.gate_fix_round(&winner, &outcomes).await? {
5200 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
5201 GateFix::Stop => break,
5202 GateFix::Defer => {
5203 self.state.save()?;
5204 return Ok(());
5205 }
5206 }
5207 }
5208 let passed = outcomes.iter().all(CommandOutcome::ok);
5209 self.state.gate = outcomes;
5210 self.state.gate_ran = true;
5211 if !passed {
5212 self.state.status = RunStatus::Blocked;
5213 let spent = self.state.gate_fixes.len();
5214 self.state.event(
5215 "gate",
5216 if spent == 0 {
5217 "gate failed; not merging".to_owned()
5218 } else {
5219 format!("gate failed after {spent} gate-fix round(s); not merging")
5220 },
5221 );
5222 }
5223 self.state.save()?;
5224 Ok(())
5225 }
5226
5227 async fn run_pre_gate(&mut self, winner: &Candidate) {
5237 let commands = self.state.config.verify.pre_gate.clone();
5238 if commands.is_empty() {
5239 return;
5240 }
5241 let shell = self.state.config.shell();
5242 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5243 let (outcomes, _) = run_commands(
5244 &mut self.state,
5245 "pre_gate",
5246 "pre_gate",
5247 0,
5248 &shell,
5249 &commands,
5250 &winner.worktree,
5251 timeout,
5252 )
5253 .await;
5254 for o in &outcomes {
5255 if !o.ok() {
5256 tracing::warn!(
5257 "pre_gate `{}` failed ({:?}); the gate decides",
5258 o.command,
5259 o.code
5260 );
5261 }
5262 self.state.event(
5263 "pre_gate",
5264 format!(
5265 "`{}` -> {}",
5266 o.command,
5267 if o.ok() {
5268 "pass".to_owned()
5269 } else {
5270 format!(
5271 "FAIL ({:?})\n{}",
5272 o.code,
5273 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5274 )
5275 }
5276 ),
5277 );
5278 }
5279 self.state.pre_gate = outcomes;
5280 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
5281 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
5282 Ok(head) => {
5283 self.state
5284 .event("pre_gate", format!("committed mechanical fixes ({head})"));
5285 self.state.pre_gate_commit = Some(head);
5286 }
5287 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
5288 },
5289 Ok(false) => {}
5290 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
5291 }
5292 if let Err(e) = self.state.save() {
5293 tracing::warn!("could not persist the pre_gate record: {e:#}");
5294 }
5295 }
5296
5297 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
5300 self.run_pre_gate(winner).await;
5301 let shell = self.state.config.shell();
5302 let gate_commands = self.state.config.verify.gate.clone();
5303 let outcomes = if gate_commands.is_empty() {
5312 Vec::new()
5313 } else {
5314 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5315 let cache_dir = self.state.config.cache_dir();
5316 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
5317 let (outcomes, _) = with_cache_lease(
5318 &mut self.state,
5319 cache_dir.as_deref(),
5320 "gate",
5321 "gate",
5322 &winner.worktree,
5323 &head,
5324 timeout,
5325 "final gate",
5326 |state, budget| {
5327 let shell = shell.clone();
5328 let gate_commands = gate_commands.clone();
5329 let worktree = winner.worktree.clone();
5330 async move {
5331 let (outcomes, timed_out_pids) = run_commands(
5332 state,
5333 "gate",
5334 "gate",
5335 0,
5336 &shell,
5337 &gate_commands,
5338 &worktree,
5339 budget,
5340 )
5341 .await;
5342 (outcomes, false, timed_out_pids)
5343 }
5344 },
5345 )
5346 .await;
5347 outcomes
5348 };
5349 if outcomes.is_empty() {
5350 self.state.event(
5355 "gate",
5356 "no gate commands configured; nothing to check, passing",
5357 );
5358 }
5359 for o in &outcomes {
5360 self.state.event(
5361 "gate",
5362 format!(
5363 "`{}` -> {}",
5364 o.command,
5365 if o.ok() {
5366 "pass".to_owned()
5367 } else {
5368 format!(
5369 "FAIL ({:?})\n{}",
5370 o.code,
5371 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5372 )
5373 }
5374 ),
5375 );
5376 }
5377 Ok(outcomes)
5378 }
5379
5380 async fn gate_fix_round(
5393 &mut self,
5394 winner: &Candidate,
5395 outcomes: &[CommandOutcome],
5396 ) -> Result<GateFix> {
5397 let cap = self.state.config.graph.gate_fix_rounds;
5398 let spent = self.state.gate_fixes.len();
5399 if spent >= cap {
5400 if cap > 0 {
5401 self.state.event(
5402 "gate",
5403 format!("{spent} gate-fix round(s) spent and the gate still fails"),
5404 );
5405 }
5406 return Ok(GateFix::Stop);
5407 }
5408 if !gate_fixable(outcomes) {
5409 self.state.event(
5410 "gate",
5411 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
5412 command or similar); not spending a fix round on it",
5413 );
5414 return Ok(GateFix::Stop);
5415 }
5416 let min_free = self.state.config.disk.min_free_bytes;
5417 if min_free > 0 {
5418 match crate::disk::free_bytes(&winner.worktree) {
5419 Ok(free) if crate::disk::enough_space(free, min_free) => {}
5420 Ok(free) => {
5421 self.state.event(
5422 "gate",
5423 format!(
5424 "only {free} bytes free ({min_free} required by `[disk] \
5425 min_free_bytes`); not spending a fix round on a failure the disk \
5426 may explain"
5427 ),
5428 );
5429 return Ok(GateFix::Stop);
5430 }
5431 Err(e) => {
5432 self.state.event(
5433 "gate",
5434 format!("free disk space could not be measured ({e:#}); no fix round"),
5435 );
5436 return Ok(GateFix::Stop);
5437 }
5438 }
5439 }
5440
5441 let attempt = spent + 1;
5442 let run_id = self.state.id.clone();
5443 let prompts = self.state.config.prompts.clone();
5444 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
5445 let base = self.landing_base();
5446 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
5447 let seat = self.seat(&fix_seat_key, &fix_spec.id);
5448 let job = SeatJob {
5449 prompt: prompt::gate_fix(
5450 &self.state.instruction,
5451 &failed,
5452 attempt,
5453 cap,
5454 &self.state.config.graph.language,
5455 ),
5456 spec: fix_spec,
5457 seat,
5458 cwd: winner.worktree.clone(),
5459 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
5460 allow_write: true,
5461 sessions: self.state.config.graph.sessions,
5462 artifacts: agent::artifacts_dir(&self.state.dir()),
5463 stem: format!("gate-fix-{attempt}"),
5464 };
5465 self.state.event(
5466 "gate",
5467 format!("gate failed; gate-fix round {attempt} of {cap}"),
5468 );
5469 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
5470 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
5471 let cache = self.state.config.cache_dir();
5472 let ctx = WaveCtx {
5473 run: &run_id,
5474 node: "gate-fix",
5475 prompts: &prompts,
5476 cache: cache.as_deref(),
5477 round: None,
5478 };
5479 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
5480 let mut record = GateFixRecord {
5481 agent: seat.agent.clone(),
5482 failed,
5483 notes: String::new(),
5484 committed: false,
5485 error: None,
5486 };
5487 match out {
5488 AgentOutcome::Ok(o) => {
5489 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
5492 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
5493 }
5494 }
5495 AgentOutcome::Dropped(_) => {
5496 record.error = Some("the CLI dropped the stream".to_owned());
5497 }
5498 AgentOutcome::Quota(o) => {
5499 self.state.quota.push(QuotaLoss {
5500 seat: seat.key.clone(),
5501 node: "gate-fix".to_owned(),
5502 at: Timestamp::now(),
5503 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5504 });
5505 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5506 }
5507 AgentOutcome::Failed(e) => record.error = Some(e),
5508 }
5509 self.state.seats.insert(seat.key.clone(), seat);
5510 if let Ok(r) = git::rescue_commit(
5511 &winner.worktree,
5512 &format!("magi: gate fix {attempt} (uncommitted work)"),
5513 )
5514 .await
5515 {
5516 self.state.note_withheld("gate-fix", &r.withheld);
5517 }
5518 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5519 record.committed = after != before;
5520 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5521 let note = record.error.clone();
5522 self.state.gate_fixes.push(record);
5523 self.state.save()?;
5524 if !changed {
5525 self.state.event(
5526 "gate",
5527 match note {
5528 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5529 None => format!("gate-fix round {attempt}: the tree did not change"),
5530 },
5531 );
5532 return Ok(GateFix::Stop);
5533 }
5534 self.state.event(
5535 "gate",
5536 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5537 );
5538
5539 let commands = self.state.config.verify.e2e.clone();
5540 if !commands.is_empty() {
5541 let shell = self.state.config.shell();
5542 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5543 let cache_dir = self.state.config.cache_dir();
5544 let context = format!("gate-fix round {attempt}");
5545 let (e2e, _) = with_cache_lease(
5546 &mut self.state,
5547 cache_dir.as_deref(),
5548 "e2e",
5549 "e2e",
5550 &winner.worktree,
5551 &after,
5552 timeout,
5553 &context,
5554 |state, budget| {
5555 let shell = shell.clone();
5556 let commands = commands.clone();
5557 let context = context.clone();
5558 let worktree = winner.worktree.clone();
5559 async move {
5560 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5561 .await
5562 }
5563 },
5564 )
5565 .await;
5566 if verify_inconclusive(&e2e) {
5567 return Ok(GateFix::Defer);
5568 }
5569 if e2e.iter().any(|o| !o.ok()) {
5570 self.state.event(
5571 "gate",
5572 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5573 );
5574 return Ok(GateFix::Stop);
5575 }
5576 }
5577 Ok(GateFix::Retry)
5578 }
5579
5580 async fn merge(&mut self) -> Result<()> {
5583 if self
5598 .state
5599 .base_sync
5600 .as_ref()
5601 .is_some_and(|s| s.conflict.is_some())
5602 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5603 != Some(RunStatus::Gating)
5604 || !self.state.gate_status().ok()
5613 {
5614 return Ok(());
5615 }
5616 if self.state.merge.is_some() {
5625 return Ok(());
5626 }
5627 let Some(winner) = self.state.winner().cloned() else {
5628 return Ok(());
5629 };
5630 let repo = self.state.repo.clone();
5631 let base = self.state.base_branch.clone();
5632 let mode = self.state.config.merge.mode;
5633 let style = self.state.config.merge.style;
5634 let pr = pr_message(&self.state, winner.label);
5635 let message = pr.commit_message();
5636
5637 let outcome = match mode {
5638 MergeMode::None => MergeOutcome {
5639 mode,
5640 ok: true,
5641 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5642 empty: false,
5643 },
5644 MergeMode::Pr | MergeMode::Local
5645 if merge_is_empty(&repo, &self.state, &winner.branch, mode).await =>
5646 {
5647 MergeOutcome {
5648 mode,
5649 ok: false,
5650 detail: empty_candidate_detail(&self.state, &base),
5651 empty: true,
5652 }
5653 }
5654 MergeMode::Local => {
5655 let on = git::current_branch(&repo).await?;
5656 if on.as_deref() != Some(base.as_str()) {
5657 MergeOutcome {
5658 mode,
5659 ok: false,
5660 detail: format!(
5661 "{} has {} checked out, not the base branch {base}",
5662 repo.display(),
5663 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5664 ),
5665 empty: false,
5666 }
5667 } else if !git::is_clean(&repo).await? {
5668 MergeOutcome {
5669 mode,
5670 ok: false,
5671 detail: format!("{} is dirty; refusing to merge", repo.display()),
5672 empty: false,
5673 }
5674 } else {
5675 let out = match style {
5676 MergeStyle::Merge => {
5677 git::merge_no_ff(&repo, &winner.branch, &message).await?
5678 }
5679 MergeStyle::Squash => {
5680 git::merge_squash(&repo, &winner.branch, &message).await?
5681 }
5682 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5683 };
5684 MergeOutcome {
5685 mode,
5686 ok: out.ok(),
5687 detail: if out.ok() { out.stdout } else { out.stderr },
5688 empty: false,
5689 }
5690 }
5691 }
5692 MergeMode::Pr => {
5693 let remote = self.state.config.merge.remote.clone();
5694 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5695 if !pushed.ok() {
5696 MergeOutcome {
5697 mode,
5698 ok: false,
5699 detail: pushed.stderr,
5700 empty: false,
5701 }
5702 } else {
5703 let found = land::find_open_pr(&winner.worktree, &winner.branch, &base).await;
5709 let out = match pr_merge_plan(found) {
5710 PrPlan::Create => {
5711 gh_pr_create(
5712 &winner.worktree,
5713 &base,
5714 &winner.branch,
5715 &pr.title,
5716 &pr.body,
5717 )
5718 .await
5719 }
5720 PrPlan::Adopt { url, title } => {
5721 self.state
5722 .event("merge", format!("Pr: adopted open pull request {url}"));
5723 if title != pr.title
5724 && let Err(e) =
5725 land::set_pr_title(&winner.worktree, &url, &pr.title).await
5726 {
5727 tracing::warn!("could not refresh title of {url}: {e:#}");
5728 self.state
5729 .event("merge", format!("Pr: title refresh failed: {e:#}"));
5730 }
5731 Ok(url)
5732 }
5733 PrPlan::Stop(why) => Err(anyhow::anyhow!(why)),
5734 };
5735 match out {
5736 Ok(url) => MergeOutcome {
5737 mode,
5738 ok: true,
5739 detail: url,
5740 empty: false,
5741 },
5742 Err(e) => MergeOutcome {
5743 mode,
5744 ok: false,
5745 detail: e.to_string(),
5746 empty: false,
5747 },
5748 }
5749 }
5750 }
5751 };
5752
5753 self.state.status = match (mode, outcome.ok) {
5754 (MergeMode::None, _) => RunStatus::Ready,
5755 (_, true) => RunStatus::Merged,
5756 (_, false) => RunStatus::Blocked,
5757 };
5758 self.state.event(
5759 "merge",
5760 format!(
5761 "{:?}: {}",
5762 mode,
5763 outcome.detail.lines().next().unwrap_or("")
5764 ),
5765 );
5766 self.state.merge = Some(outcome);
5767 self.state.save()?;
5768
5769 if self.state.config.graph.land
5775 && mode == MergeMode::Pr
5776 && self.state.status == RunStatus::Merged
5777 {
5778 self.run_land().await?;
5779 }
5780 self.settle_questions();
5785 Ok(())
5786 }
5787
5788 async fn run_land(&mut self) -> Result<()> {
5799 let url = self
5800 .state
5801 .merge
5802 .as_ref()
5803 .map(|m| m.detail.clone())
5804 .unwrap_or_default();
5805 let url = url.lines().next().unwrap_or("").trim().to_owned();
5806 if !url.starts_with("http") {
5807 return Ok(());
5808 }
5809 match land::land(&mut self.state, &url).await {
5812 Ok(pr) if self.state.parked => {
5813 let _ = pr;
5817 }
5818 Ok(pr) => {
5819 self.state.status = match pr.state {
5820 land::PrLifecycle::Merged => RunStatus::Merged,
5821 _ => RunStatus::Blocked,
5822 };
5823 if bump::should_release_bump(self.state.status)
5830 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5831 {
5832 self.state
5838 .event("bump", format!("release bump skipped: {e:#}"));
5839 }
5840 self.state.save()?;
5841 }
5842 Err(e) => {
5843 self.state.status = RunStatus::Blocked;
5844 self.state.event("land", format!("gave up: {e}"));
5845 self.state.save()?;
5846 }
5847 }
5848 Ok(())
5849 }
5850
5851 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5855 if let Some(existing) = self.state.seats.get(key)
5856 && existing.agent == agent
5857 {
5858 return existing.clone();
5859 }
5860 let fresh = SeatState::new(key, agent, self.state.seed);
5861 self.state.seats.insert(key.to_owned(), fresh.clone());
5862 fresh
5863 }
5864
5865 fn view(&self, c: &Candidate) -> CandidateView {
5867 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5868 .unwrap_or_default();
5869 let (patch, _) = blind::sanitize_patch(
5870 &format!("candidate {} patch", c.label),
5871 &raw,
5872 &self.state.config.blind,
5873 );
5874 CandidateView {
5875 label: c.label,
5876 branch: c.branch.clone(),
5877 summary: c.summary.clone(),
5878 stat: c.stat.clone(),
5879 patch,
5880 }
5881 }
5882
5883 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5885 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5886 prompt::judge(
5887 "(see above)",
5888 &views,
5889 self.roles.judges.len(),
5890 base_short,
5891 "en",
5892 )
5893 }
5894
5895 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5902 let mut turns = Vec::new();
5903 for j in &self.state.judgements {
5904 if j.ranking.is_empty() {
5905 continue;
5906 }
5907 let reasons = j
5908 .reasons
5909 .iter()
5910 .map(|(k, v)| format!("- {k}: {v}"))
5911 .collect::<Vec<_>>()
5912 .join("\n");
5913 turns.push(Turn {
5914 who: format!("Judge {} (opening ranking)", j.judge),
5915 is_self: j.judge == self_idx + 1,
5916 body: format!(
5917 "Ranked {}{}{reasons}",
5918 j.ranking.iter().collect::<String>(),
5919 if reasons.is_empty() {
5920 ""
5921 } else {
5922 ", because:\n"
5923 }
5924 ),
5925 });
5926 }
5927 for t in self
5928 .state
5929 .deliberation
5930 .iter()
5931 .flat_map(|r| r.turns.iter())
5932 .chain(current)
5933 {
5934 turns.push(Turn {
5935 who: format!("Judge {}", t.judge),
5936 is_self: t.judge == self_idx + 1,
5937 body: t.body.clone(),
5938 });
5939 }
5940 turns
5941 }
5942}
5943
5944fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5946 agent::has_session(spec.kind, seat, sessions)
5947}
5948
5949fn next_untried_implementer<'a>(
5970 roster: &'a [AgentSpec],
5971 start: usize,
5972 tried: &BTreeSet<String>,
5973) -> Option<&'a AgentSpec> {
5974 roster
5975 .get(start + 1..)?
5976 .iter()
5977 .find(|s| !tried.contains(&s.id))
5978}
5979
5980fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5994 commands.iter().any(|c| c.exit_code.is_none())
5995}
5996
5997fn verified_noop_claim(
6010 usable: bool,
6011 commands: &[agent::CommandEvidence],
6012 text: &str,
6013) -> Option<String> {
6014 (usable && !has_unconfirmed_command(commands))
6015 .then(|| verdict::verified_noop(text))
6016 .flatten()
6017}
6018
6019fn short(commit: &str) -> String {
6020 commit.chars().take(7).collect()
6021}
6022
6023fn make_executable(path: &Path) -> Result<()> {
6024 #[cfg(unix)]
6025 {
6026 use std::os::unix::fs::PermissionsExt as _;
6027 let mut perms = std::fs::metadata(path)?.permissions();
6028 perms.set_mode(0o755);
6029 std::fs::set_permissions(path, perms)?;
6030 }
6031 #[cfg(not(unix))]
6032 {
6033 let _ = path;
6034 }
6035 Ok(())
6036}
6037
6038struct WaveCtx<'a> {
6045 run: &'a str,
6048 node: &'a str,
6050 prompts: &'a Prompts,
6051 cache: Option<&'a Path>,
6053 round: Option<usize>,
6056}
6057
6058async fn run_one(
6060 job: SeatJob,
6061 sem: Arc<Semaphore>,
6062 ctx: &WaveCtx<'_>,
6063 state: &mut RunState,
6064 attempt: usize,
6065) -> (SeatState, AgentOutcome) {
6066 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
6067 .await
6068 .pop()
6069 .expect("one job in, one result out");
6070 (seat, out)
6071}
6072
6073async fn wave(
6079 jobs: Vec<SeatJob>,
6080 sem: Arc<Semaphore>,
6081 ctx: &WaveCtx<'_>,
6082 state: &mut RunState,
6083 attempt: usize,
6084) -> Vec<(usize, SeatState, AgentOutcome)> {
6085 let WaveCtx {
6086 run,
6087 node,
6088 prompts,
6089 cache,
6090 round,
6091 } = *ctx;
6092 for job in &jobs {
6093 state.seat_started(node, &job.seat.key, job.timeout, attempt);
6094 }
6095 if let Err(e) = state.save() {
6096 tracing::warn!("could not persist in-progress seats: {e:#}");
6101 }
6102 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
6118 let wait_started = Instant::now();
6119 let cache_guard = if let Some(cache_dir) = cache {
6120 if jobs_had_a_writer {
6121 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
6122 let budget = jobs
6123 .iter()
6124 .map(|j| j.timeout)
6125 .max()
6126 .unwrap_or(Duration::from_secs(60));
6127 acquire_cache_lease(state, cache_dir, &owner, budget, node)
6128 .await
6129 .ok()
6130 } else {
6131 None
6132 }
6133 } else {
6134 None
6135 };
6136 let waited_for_lease = wait_started.elapsed();
6143 let mut set = tokio::task::JoinSet::new();
6144 let overlay = prompts.overlay(node);
6145 for (i, mut job) in jobs.into_iter().enumerate() {
6146 job.timeout = job.timeout.saturating_sub(waited_for_lease);
6147 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
6148 if cache.is_some() {
6149 job.prompt.push('\n');
6150 job.prompt
6151 .push_str(&prompt::build_cache_note(node, job.allow_write));
6152 }
6153 let sem = Arc::clone(&sem);
6154 let run = run.to_owned();
6155 let node = node.to_owned();
6156 let attachments = if node == "implement" {
6159 state.attachments.clone()
6160 } else {
6161 Vec::new()
6162 };
6163 let cache = cache
6174 .filter(|_| job.allow_write && cache_guard.is_some())
6175 .map(Path::to_path_buf);
6176 set.spawn(async move {
6177 let _permit = sem.acquire().await;
6178 let mut seat = job.seat;
6179 let out = agent::invoke(
6180 &job.spec,
6181 &mut seat,
6182 &Invocation {
6183 cwd: &job.cwd,
6184 prompt: &job.prompt,
6185 timeout: job.timeout,
6186 allow_write: job.allow_write,
6187 sessions: job.sessions,
6188 artifacts: &job.artifacts,
6189 stem: &job.stem,
6190 run: &run,
6191 node: &node,
6192 cache_dir: cache.as_deref(),
6193 attachments: &attachments,
6194 },
6195 )
6196 .await;
6197 let out = match out {
6198 Ok(o) if o.usable() => AgentOutcome::Ok(o),
6199 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
6200 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
6208 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
6209 Ok(o) => AgentOutcome::Failed(format!(
6210 "exited with {:?} and no usable output",
6211 o.exit_code
6212 )),
6213 Err(e) => AgentOutcome::Failed(e.to_string()),
6214 };
6215 (i, seat, out)
6216 });
6217 }
6218 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
6219 while let Some(joined) = set.join_next().await {
6220 let (i, seat, out) = match joined {
6221 Ok(v) => v,
6222 Err(e) => {
6226 tracing::error!("agent task panicked: {e}");
6227 continue;
6228 }
6229 };
6230 state.seat_finished(&seat.key);
6231 record_jobs(state, node, round, &seat.key, &out);
6232 if let Err(e) = state.save() {
6233 tracing::warn!("could not persist a seat's completion: {e:#}");
6234 }
6235 if collected.len() <= i {
6236 collected.resize_with(i + 1, || None);
6237 }
6238 collected[i] = Some((i, seat, out));
6239 }
6240 if state
6246 .active
6247 .values()
6248 .any(|a| a.node == node && a.attempt == attempt)
6249 {
6250 state
6251 .active
6252 .retain(|_, a| !(a.node == node && a.attempt == attempt));
6253 if let Err(e) = state.save() {
6254 tracing::warn!("could not persist the end of a wave: {e:#}");
6255 }
6256 }
6257 if let Some(cache_dir) = cache
6264 && jobs_had_a_writer
6265 {
6266 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
6267 }
6268 if let Some(guard) = cache_guard {
6269 guard.release();
6270 }
6271 collected.into_iter().flatten().collect()
6272}
6273
6274fn record_jobs(
6285 state: &mut RunState,
6286 node: &str,
6287 round: Option<usize>,
6288 seat: &str,
6289 out: &AgentOutcome,
6290) {
6291 let commands: &[agent::CommandEvidence] = match out {
6292 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
6293 AgentOutcome::Failed(_) => &[],
6294 };
6295 let checked_at = Timestamp::now();
6296 for c in commands {
6297 state.jobs.push(JobRecord {
6298 node: node.to_owned(),
6299 round,
6300 seat: seat.to_owned(),
6301 id: c.id.clone(),
6302 description: c.description.clone(),
6303 checked_at,
6304 status: match c.exit_code {
6305 Some(0) => JobStatus::Completed,
6306 Some(_) => JobStatus::Failed,
6307 None => JobStatus::Unknown,
6308 },
6309 exit_code: c.exit_code,
6310 result_summary: c.result_summary.clone(),
6311 source: c.source.clone(),
6312 });
6313 }
6314}
6315
6316fn round_is_clean(
6337 blocking: usize,
6338 e2e_ok: bool,
6339 answered: usize,
6340 expected: usize,
6341 quota_missing: usize,
6342 policy: IncompleteReviewPolicy,
6343) -> bool {
6344 if blocking != 0 || !e2e_ok {
6345 return false;
6346 }
6347 if answered == expected || policy == IncompleteReviewPolicy::Warn {
6348 return true;
6349 }
6350 answered > 0 && expected - answered <= quota_missing
6351}
6352
6353fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
6377 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
6378 return Some(RunStatus::Gating);
6379 }
6380 let last = reviews.last()?;
6381 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
6382 if reviews.len() < max_rounds && !stagnant {
6383 return None;
6384 }
6385 if last.incomplete() && last.blocking == 0 {
6386 return Some(RunStatus::Blocked);
6387 }
6388 if last.e2e_status() == E2eStatus::ResourceBlocked {
6389 return None;
6390 }
6391 Some(if last.e2e.iter().all(CommandOutcome::ok) {
6392 RunStatus::Gating
6393 } else {
6394 RunStatus::Blocked
6395 })
6396}
6397
6398fn retry_budget(full: Duration, nudged: bool) -> Duration {
6413 if nudged {
6414 (full / 4).max(Duration::from_secs(120)).min(full)
6415 } else {
6416 full
6417 }
6418}
6419
6420#[allow(clippy::too_many_arguments)]
6433async fn ask_json_wave<T>(
6434 jobs: Vec<SeatJob>,
6435 sem: Arc<Semaphore>,
6436 retries: usize,
6437 ctx: &WaveCtx<'_>,
6438 losses: &mut Vec<QuotaLoss>,
6439 state: &mut RunState,
6440 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
6441) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
6442where
6443 T: serde::de::DeserializeOwned + Send + 'static,
6444{
6445 let n = jobs.len();
6446 let originals: Vec<SeatJob> = jobs;
6447 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
6448 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
6449 let mut attempts_used: Vec<usize> = vec![0; n];
6456 let mut pending: Vec<usize> = (0..n).collect();
6457
6458 for attempt in 0..=retries {
6459 if pending.is_empty() {
6460 break;
6461 }
6462 let mut batch = Vec::with_capacity(pending.len());
6463 for &i in &pending {
6464 let src = &originals[i];
6465 let (prompt, timeout) = if attempt == 0 {
6468 (src.prompt.clone(), src.timeout)
6469 } else {
6470 let why = done[i]
6471 .as_ref()
6472 .and_then(|r| r.as_ref().err().map(ToString::to_string))
6473 .unwrap_or_else(|| "no parsable answer".to_owned());
6474 let nudge = prompt::nudge(&why);
6475 let nudged = has_context(&src.spec, &seats[i], src.sessions);
6476 let prompt = if nudged {
6477 nudge
6478 } else {
6479 format!("{}\n\n---\n\n{}", src.prompt, nudge)
6480 };
6481 (prompt, retry_budget(src.timeout, nudged))
6482 };
6483 batch.push(SeatJob {
6484 spec: src.spec.clone(),
6485 seat: seats[i].clone(),
6486 cwd: src.cwd.clone(),
6487 prompt,
6488 timeout,
6489 allow_write: src.allow_write,
6490 sessions: src.sessions,
6491 artifacts: src.artifacts.clone(),
6492 stem: if attempt == 0 {
6493 src.stem.clone()
6494 } else {
6495 format!("{}-retry{attempt}", src.stem)
6496 },
6497 });
6498 }
6499
6500 if attempt > 0 {
6501 let seats_out: Vec<&str> = pending
6502 .iter()
6503 .map(|&i| originals[i].seat.key.as_str())
6504 .collect();
6505 state.event(
6506 ctx.node,
6507 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
6508 );
6509 }
6510 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
6511 let mut still = Vec::new();
6512 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
6513 seats[i] = seat;
6514 let (parsed, quota) = match out {
6515 AgentOutcome::Ok(o) => (
6516 match verdict::extract_json::<T>(&o.text) {
6517 Ok(v) => match validate(&v) {
6518 Ok(()) => Ok((v, o)),
6519 Err(e) => Err(e),
6520 },
6521 Err(e) => Err(e),
6522 },
6523 false,
6524 ),
6525 AgentOutcome::Quota(o) => {
6526 losses.push(QuotaLoss {
6527 seat: originals[i].seat.key.clone(),
6528 node: ctx.node.to_owned(),
6529 at: Timestamp::now(),
6530 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
6531 });
6532 (
6533 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
6534 true,
6535 )
6536 }
6537 AgentOutcome::Dropped(o) => {
6542 let why = o
6543 .dropped
6544 .as_ref()
6545 .map(|d| d.why.as_str())
6546 .unwrap_or("the CLI ended the stream without delivering its answer");
6547 (
6548 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
6549 false,
6550 )
6551 }
6552 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
6553 };
6554 let failed = parsed.is_err();
6555 done[i] = Some(parsed);
6556 attempts_used[i] = attempt;
6557 if failed && !quota {
6560 still.push(i);
6561 }
6562 }
6563 pending = still;
6564 }
6565
6566 seats
6567 .into_iter()
6568 .zip(done)
6569 .zip(attempts_used)
6570 .map(|((seat, res), attempts)| {
6571 (
6572 seat,
6573 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6574 attempts,
6575 )
6576 })
6577 .collect()
6578}
6579
6580async fn acquire_cache_lease(
6593 state: &mut RunState,
6594 cache_dir: &Path,
6595 owner: &crate::cache::Owner,
6596 budget: Duration,
6597 context: &str,
6598) -> Result<crate::cache::Guard> {
6599 let home = crate::run::home();
6600 let started = Instant::now();
6601 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6602 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6603 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6604 Err(e) => {
6605 state.event(
6606 "verify",
6607 format!("{context}: could not check the shared build cache: {e:#}"),
6608 );
6609 if let Err(e2) = state.save() {
6610 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6611 }
6612 return Err(e);
6613 }
6614 };
6615 state.event(
6616 "verify",
6617 format!(
6618 "{context}: waiting for the shared build cache at {} ({})",
6619 cache_dir.display(),
6620 busy.describe()
6621 ),
6622 );
6623 if let Err(e) = state.save() {
6624 tracing::warn!("could not persist a cache wait: {e:#}");
6625 }
6626 let remaining = budget.saturating_sub(started.elapsed());
6627 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6628 Ok(g) => Ok(g),
6629 Err(e) => {
6630 state.event("verify", format!("{context}: {e:#}"));
6631 if let Err(e2) = state.save() {
6632 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6633 }
6634 Err(e)
6635 }
6636 }
6637}
6638
6639#[allow(clippy::too_many_arguments)]
6660async fn with_cache_lease<'s, F, Fut>(
6661 state: &'s mut RunState,
6662 cache_dir: Option<&Path>,
6663 node: &str,
6664 seat: &str,
6665 worktree: &Path,
6666 head: &str,
6667 budget: Duration,
6668 context: &str,
6669 body: F,
6670) -> (Vec<CommandOutcome>, bool)
6671where
6672 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6673 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6674{
6675 let Some(cache_dir) = cache_dir else {
6676 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6677 return (outcomes, retried);
6678 };
6679 let home = crate::run::home();
6680 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6681 let started = Instant::now();
6682 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6683 Ok(g) => g,
6684 Err(e) => {
6685 return (
6686 vec![CommandOutcome {
6687 command: "(waiting for the shared build cache)".to_owned(),
6688 code: None,
6689 output_tail: e.to_string(),
6690 duration_ms: started.elapsed().as_millis() as u64,
6691 resource_blocked: true,
6692 }],
6693 false,
6694 );
6695 }
6696 };
6697 let identity = crate::cache::Identity::new(worktree, head);
6698 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6699 state.event(
6707 "verify",
6708 format!(
6709 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6710 worktree.display(),
6711 short(head)
6712 ),
6713 );
6714 guard.release();
6715 return (
6716 vec![CommandOutcome {
6717 command: "(confirming the shared build cache is fresh)".to_owned(),
6718 code: None,
6719 output_tail: e.to_string(),
6720 duration_ms: started.elapsed().as_millis() as u64,
6721 resource_blocked: true,
6722 }],
6723 false,
6724 );
6725 }
6726 let remaining = budget.saturating_sub(started.elapsed());
6727 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6728 if !timed_out_pids.is_empty() {
6733 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6734 }
6735 guard.release();
6736 (outcomes, retried)
6737}
6738
6739async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6751 wait_for_pids_with(
6752 pids,
6753 crate::proc::pid_alive,
6754 LEASE_RELEASE_POLL,
6755 LEASE_RELEASE_MAX_WAIT,
6756 )
6757 .await;
6758}
6759
6760async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6766 pids: &[u32],
6767 alive: F,
6768 poll: Duration,
6769 max_wait: Duration,
6770) {
6771 let deadline = Instant::now() + max_wait;
6772 loop {
6773 if pids.iter().all(|&pid| !alive(pid)) {
6774 return;
6775 }
6776 if Instant::now() >= deadline {
6777 return;
6778 }
6779 tokio::time::sleep(poll).await;
6780 }
6781}
6782
6783fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6790 outcomes.iter().any(|o| o.resource_blocked)
6791}
6792
6793enum GateFix {
6795 Retry,
6797 Stop,
6800 Defer,
6803}
6804
6805fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6813 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6814 red.peek().is_some()
6815 && red.all(|o| {
6816 !o.resource_blocked
6817 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6818 && !o.output_tail.trim().is_empty()
6819 })
6820}
6821
6822fn e2e_outcome_label(o: &CommandOutcome) -> String {
6826 if o.ok() {
6827 return "pass".to_owned();
6828 }
6829 let reason = if o.build_failed() {
6830 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6831 } else {
6832 format!("FAIL ({:?})", o.code)
6833 };
6834 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6835}
6836
6837async fn run_e2e_with_retry(
6845 state: &mut RunState,
6846 shell: &[String],
6847 commands: &[String],
6848 worktree: &Path,
6849 timeout: Duration,
6850 context: &str,
6851) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6852 let (mut e2e, mut timed_out_pids) = run_commands(
6853 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6854 )
6855 .await;
6856 for o in &e2e {
6857 state.event(
6858 "verify",
6859 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6860 );
6861 }
6862 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6866 if verify_retried {
6867 state.event(
6868 "verify",
6869 format!(
6870 "{context}: verify could not build/link, not a test result — retrying once \
6871 before concluding"
6872 ),
6873 );
6874 let retried = run_commands(
6875 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6876 )
6877 .await;
6878 e2e = retried.0;
6879 timed_out_pids.extend(retried.1);
6882 for o in &e2e {
6883 state.event(
6884 "verify",
6885 format!(
6886 "{context}: retry `{}` -> {}",
6887 o.command,
6888 e2e_outcome_label(o)
6889 ),
6890 );
6891 }
6892 }
6893 (e2e, verify_retried, timed_out_pids)
6894}
6895
6896#[allow(clippy::too_many_arguments)]
6912async fn run_commands(
6913 state: &mut RunState,
6914 node: &str,
6915 task: &str,
6916 attempt: usize,
6917 shell: &[String],
6918 commands: &[String],
6919 cwd: &Path,
6920 timeout: Duration,
6921) -> (Vec<CommandOutcome>, Vec<u32>) {
6922 if commands.is_empty() {
6923 return (Vec::new(), Vec::new());
6928 }
6929 let mut out = Vec::new();
6930 let mut timed_out_pids = Vec::new();
6931 let total = commands.len();
6932 for (idx, command) in commands.iter().enumerate() {
6933 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6934 if let Err(e) = state.save() {
6935 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6936 }
6937 let started = Instant::now();
6938 let mut cmd = tokio::process::Command::new(&shell[0]);
6939 cmd.quiet();
6940 cmd.args(&shell[1..])
6941 .arg(command)
6942 .current_dir(cwd)
6943 .stdin(std::process::Stdio::null())
6944 .stdout(std::process::Stdio::piped())
6945 .stderr(std::process::Stdio::piped())
6946 .kill_on_drop(true);
6947 let spawned = cmd.spawn();
6948 let (code, body) = match spawned {
6949 Ok(child) => {
6950 let pid = child.id();
6955 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6956 Ok(Ok(o)) => {
6957 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6958 body.push_str(&String::from_utf8_lossy(&o.stderr));
6959 (o.status.code(), body)
6960 }
6961 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6962 Err(_) => {
6963 if let Some(pid) = pid {
6964 timed_out_pids.push(pid);
6965 }
6966 (None, format!("timed out after {}s", timeout.as_secs()))
6967 }
6968 }
6969 }
6970 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6971 };
6972 out.push(CommandOutcome {
6973 command: command.clone(),
6974 code,
6975 output_tail: tail(&body, OUTPUT_TAIL),
6976 duration_ms: started.elapsed().as_millis() as u64,
6977 resource_blocked: false,
6978 });
6979 }
6980 state.task_finished(task);
6981 if let Err(e) = state.save() {
6982 tracing::warn!("could not persist the end of {task}: {e:#}");
6983 }
6984 (out, timed_out_pids)
6985}
6986
6987fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6999 let repo = repo.display();
7000 match style {
7001 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
7002 MergeStyle::Squash => {
7003 let subject = message
7006 .lines()
7007 .next()
7008 .unwrap_or(branch)
7009 .replace(['\\', '"', '$', '`'], "");
7010 format!(
7011 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
7012 )
7013 }
7014 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
7015 }
7016}
7017
7018const PR_TITLE_MAX: usize = 240;
7029
7030struct PrMessage {
7035 title: String,
7036 body: String,
7037}
7038
7039impl PrMessage {
7040 fn commit_message(&self) -> String {
7044 format!("{}\n\n{}", self.title, self.body)
7045 }
7046}
7047
7048fn title_marker(line: &str) -> Option<&str> {
7050 let line = line.trim();
7051 let head = line.get(..6)?;
7052 head.eq_ignore_ascii_case("title:")
7053 .then(|| line[6..].trim())
7054}
7055
7056fn summary_title(summary: &str) -> Option<String> {
7061 let first = summary.lines().find(|l| !l.trim().is_empty())?;
7062 let raw = title_marker(first)?;
7063 if raw.is_empty() {
7064 return None;
7065 }
7066 let title = queue::title_from(raw, PR_TITLE_MAX);
7067 let lower = title.to_ascii_lowercase();
7068 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
7069 return None;
7070 }
7071 Some(title)
7072}
7073
7074fn summary_without_title(summary: &str) -> String {
7077 let mut lines = summary.trim().lines().peekable();
7078 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
7079 lines.next();
7080 }
7081 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
7082}
7083
7084fn pr_message(state: &RunState, winner: char) -> PrMessage {
7098 let summary = state
7099 .candidates
7100 .iter()
7101 .find(|c| c.label == winner)
7102 .map(|c| c.summary.as_str())
7103 .unwrap_or_default();
7104 let title = summary_title(summary).unwrap_or_else(|| {
7107 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
7108 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
7109 t
7110 } else {
7111 format!(
7112 "chore: land candidate {} of run {}",
7113 winner.to_ascii_uppercase(),
7114 state.id
7115 )
7116 }
7117 });
7118
7119 let mut body = String::new();
7120 let what = summary_without_title(summary);
7121 if !what.is_empty() {
7122 body.push_str("## Summary\n\n");
7123 body.push_str(&what);
7124 body.push_str("\n\n");
7125 }
7126
7127 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
7128 if let Some(fix) = fix
7129 && !fix.notes.trim().is_empty()
7130 {
7131 body.push_str("## Review fixes\n\n");
7132 body.push_str(fix.notes.trim());
7133 body.push_str("\n\n");
7134 }
7135
7136 let open = state.open_findings();
7137 if !open.is_empty() {
7138 body.push_str("## Open review findings\n\n");
7139 for f in &open {
7140 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
7141 }
7142 body.push('\n');
7143 }
7144
7145 if let Some(fix) = fix
7146 && !fix.rejected.is_empty()
7147 {
7148 body.push_str("## Declined by the fixer\n\n");
7149 for r in &fix.rejected {
7150 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
7151 }
7152 body.push('\n');
7153 }
7154
7155 let task = state.instruction.trim();
7156 let task = if task.is_empty() {
7157 "(empty task)"
7158 } else {
7159 task
7160 };
7161 body.push_str(&format!(
7162 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
7163 task.replace("</details>", "</details>")
7164 ));
7165
7166 body.push_str(&format!(
7167 "\n---\nmagi:run/{} magi:candidate-{}\n",
7168 state.id,
7169 winner.to_ascii_lowercase()
7170 ));
7171
7172 let id = crate::scrub::Identity::current();
7175 PrMessage {
7176 title: crate::scrub::scrub(&title, &id),
7177 body: crate::scrub::scrub(&body, &id),
7178 }
7179}
7180
7181fn seeded_instruction(state: &RunState) -> String {
7185 match refs::describe(&state.seeds) {
7186 Some(facts) => format!(
7187 "{}\n\n# Existing work the task refers to\n\n{facts}\n\n\
7188 Candidates start from the unmerged branch named above, when there \
7189 is one, and carry any unmerged commit named by sha as a \
7190 cherry-pick. Check that this is what the task meant before \
7191 building on it.",
7192 state.instruction
7193 ),
7194 None => state.instruction.clone(),
7195 }
7196}
7197
7198async fn merge_is_empty(repo: &Path, state: &RunState, branch: &str, mode: MergeMode) -> bool {
7202 let base = &state.base_branch;
7203 let mut against = base.clone();
7204 if mode == MergeMode::Pr {
7205 let remote = &state.config.merge.remote;
7206 let tracking = format!("{remote}/{base}");
7207 let fetched = git::fetch(repo, remote, base).await;
7208 if fetched.is_ok_and(|o| o.ok()) && git::rev_exists(repo, &tracking).await {
7209 against = tracking;
7210 }
7211 }
7212 matches!(git::commits_ahead(repo, &against, branch).await, Ok(0))
7213}
7214
7215fn empty_candidate_detail(state: &RunState, base: &str) -> String {
7218 let mut detail = format!(
7219 "empty candidate: the winning branch has 0 commits ahead of {base}, so there is \
7220 nothing to open a pull request for"
7221 );
7222 match refs::describe(&state.seeds) {
7223 Some(facts) => detail.push_str(&format!("\nReferences in the task:\n{facts}")),
7224 None => detail.push_str(
7225 "\nThe task names no existing branch or commit; if it means to land work \
7226 that lives elsewhere, name the branch (magi/<run>/<label>) or the sha.",
7227 ),
7228 }
7229 detail
7230}
7231
7232#[derive(Debug, PartialEq, Eq)]
7235enum PrPlan {
7236 Create,
7237 Adopt { url: String, title: String },
7238 Stop(String),
7239}
7240
7241fn pr_merge_plan(found: Result<land::OpenPr>) -> PrPlan {
7244 match found {
7245 Ok(land::OpenPr::None) => PrPlan::Create,
7246 Ok(land::OpenPr::One { url, title }) => PrPlan::Adopt { url, title },
7247 Ok(land::OpenPr::Many(urls)) => PrPlan::Stop(format!(
7248 "several open pull requests exist for this branch, not picking one: {}",
7249 urls.join(" ")
7250 )),
7251 Err(e) => PrPlan::Stop(format!("could not look up open pull requests: {e:#}")),
7252 }
7253}
7254
7255async fn gh_pr_create(
7257 cwd: &Path,
7258 base: &str,
7259 head: &str,
7260 title: &str,
7261 body: &str,
7262) -> Result<String> {
7263 let out = tokio::process::Command::new("gh")
7264 .args([
7265 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
7266 ])
7267 .current_dir(cwd)
7268 .quiet()
7269 .stdin(std::process::Stdio::null())
7270 .output()
7271 .await
7272 .context("spawn gh")?;
7273 if out.status.success() {
7274 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
7275 } else {
7276 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
7277 }
7278}
7279
7280pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
7289 let repo = state.repo.clone();
7290 let root = state.worktree_root();
7291 let winner = state.tally.as_ref().map(|t| t.winner);
7292 let mut removed = Vec::new();
7293
7294 for i in 0..state.candidates.len() {
7295 let c = state.candidates[i].clone();
7296 let is_winner = Some(c.label) == winner;
7297 if is_winner && !drop_winner {
7298 continue;
7299 }
7300 if c.worktree.exists() {
7301 git::worktree_remove(&repo, &c.worktree).await.ok();
7302 removed.push(c.worktree.to_string_lossy().into_owned());
7303 }
7304 let handed_over = state.released_branches.contains(&c.branch);
7307 if !handed_over && git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
7308 git::branch_delete(&repo, &c.branch).await.ok();
7309 removed.push(c.branch.clone());
7310 }
7311 state.candidates[i].folded = true;
7312 }
7313
7314 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
7315 let path = name.path();
7316 let keep = !drop_winner
7317 && winner.is_some_and(|w| {
7318 path.file_name()
7319 .is_some_and(|n| n == format!("cand-{w}").as_str())
7320 });
7321 if keep {
7322 continue;
7323 }
7324 git::worktree_remove(&repo, &path).await.ok();
7325 removed.push(path.to_string_lossy().into_owned());
7326 }
7327
7328 remove_if_empty(&root);
7337
7338 if state.enabled_worktree_config && drop_winner {
7339 git::release_worktree_config(&repo).await.ok();
7343 state.enabled_worktree_config = false;
7344 }
7345 state.save_under(home)?;
7346 Ok(removed)
7347}
7348
7349fn remove_if_empty(dir: &Path) {
7360 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
7361 std::fs::remove_dir(dir).ok();
7362 }
7363}
7364
7365pub fn worst_open(state: &RunState) -> Option<Severity> {
7367 state
7368 .reviews
7369 .last()?
7370 .reviews
7371 .iter()
7372 .flat_map(|r| r.findings.iter())
7373 .map(|f| f.severity)
7374 .max()
7375}
7376
7377#[cfg(test)]
7378mod tests {
7379 #[test]
7380 fn pr_merge_plan_creates_adopts_or_stops() {
7381 assert_eq!(pr_merge_plan(Ok(land::OpenPr::None)), PrPlan::Create);
7382 assert_eq!(
7383 pr_merge_plan(Ok(land::OpenPr::One {
7384 url: "u".into(),
7385 title: "t".into()
7386 })),
7387 PrPlan::Adopt {
7388 url: "u".into(),
7389 title: "t".into()
7390 }
7391 );
7392 let PrPlan::Stop(many) =
7393 pr_merge_plan(Ok(land::OpenPr::Many(vec!["a".into(), "b".into()])))
7394 else {
7395 panic!("many must stop");
7396 };
7397 assert!(many.contains('a') && many.contains('b'));
7398 let PrPlan::Stop(err) = pr_merge_plan(Err(anyhow::anyhow!("bad token"))) else {
7399 panic!("a failed lookup must stop");
7400 };
7401 assert!(err.contains("bad token"));
7402 }
7403
7404 use super::*;
7405 use crate::run::GateStatus;
7406 use std::collections::BTreeMap;
7407 use std::time::Duration;
7408
7409 fn conductor() -> AgentSpec {
7410 AgentSpec {
7411 id: "conductor".to_owned(),
7412 kind: crate::config::AgentKind::Command,
7413 model: None,
7414 command: vec!["true".to_owned()],
7415 extra_args: Vec::new(),
7416 env: BTreeMap::new(),
7417 prompt_delivery: None,
7418 }
7419 }
7420
7421 fn spec(id: &str) -> AgentSpec {
7422 AgentSpec {
7423 id: id.to_owned(),
7424 kind: crate::config::AgentKind::Command,
7425 model: None,
7426 command: vec!["true".to_owned()],
7427 extra_args: Vec::new(),
7428 env: BTreeMap::new(),
7429 prompt_delivery: None,
7430 }
7431 }
7432
7433 #[test]
7440 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
7441 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7442 let tried = BTreeSet::from(["beta".to_owned()]);
7443 let next = next_untried_implementer(&roster, 1, &tried);
7446 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
7447 }
7448
7449 #[test]
7450 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
7451 let roster = vec![spec("alpha"), spec("beta")];
7452 let tried = BTreeSet::from(["beta".to_owned()]);
7453 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7457 }
7458
7459 #[test]
7460 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
7461 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7462 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
7463 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7467 }
7468
7469 #[test]
7470 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
7471 let roster = vec![spec("a"), spec("a"), spec("b")];
7472 let tried = BTreeSet::from(["a".to_owned()]);
7473 let next = next_untried_implementer(&roster, 0, &tried);
7474 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
7475 }
7476
7477 #[test]
7478 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
7479 let roster = vec![spec("a"), spec("b")];
7480 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
7481 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
7482 }
7483
7484 #[test]
7485 fn remove_if_empty_only_ever_takes_a_bare_directory() {
7486 let dir = tempfile::tempdir().unwrap();
7487 let bay = dir.path().join("ffff");
7488
7489 remove_if_empty(&bay);
7491 assert!(!bay.exists());
7492
7493 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
7496 remove_if_empty(&bay);
7497 assert!(bay.exists(), "non-empty directory must survive");
7498
7499 std::fs::remove_dir(bay.join("cand-A")).unwrap();
7501 remove_if_empty(&bay);
7502 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
7503 }
7504
7505 #[test]
7514 fn a_full_panel_that_found_nothing_is_clean() {
7515 assert!(round_is_clean(
7516 0,
7517 true,
7518 2,
7519 2,
7520 0,
7521 IncompleteReviewPolicy::Block
7522 ));
7523 }
7524
7525 #[test]
7526 fn a_missing_seat_is_never_clean_under_the_default_policy() {
7527 assert!(!round_is_clean(
7528 0,
7529 true,
7530 1,
7531 2,
7532 0,
7533 IncompleteReviewPolicy::Block
7534 ));
7535 }
7536
7537 #[test]
7538 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
7539 assert!(!round_is_clean(
7540 1,
7541 true,
7542 1,
7543 2,
7544 0,
7545 IncompleteReviewPolicy::Warn
7546 ));
7547 }
7548
7549 #[test]
7550 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
7551 assert!(round_is_clean(
7552 0,
7553 true,
7554 1,
7555 2,
7556 0,
7557 IncompleteReviewPolicy::Warn
7558 ));
7559 }
7560
7561 #[test]
7562 fn a_full_panel_with_an_open_finding_is_not_clean() {
7563 assert!(!round_is_clean(
7564 1,
7565 true,
7566 2,
7567 2,
7568 0,
7569 IncompleteReviewPolicy::Block
7570 ));
7571 }
7572
7573 #[test]
7574 fn a_full_panel_with_a_red_e2e_is_not_clean() {
7575 assert!(!round_is_clean(
7576 0,
7577 false,
7578 2,
7579 2,
7580 0,
7581 IncompleteReviewPolicy::Block
7582 ));
7583 }
7584
7585 #[test]
7592 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
7593 assert!(round_is_clean(
7596 0,
7597 true,
7598 1,
7599 2,
7600 1,
7601 IncompleteReviewPolicy::Block
7602 ));
7603 }
7604
7605 #[test]
7606 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
7607 assert!(!round_is_clean(
7610 0,
7611 true,
7612 1,
7613 2,
7614 0,
7615 IncompleteReviewPolicy::Block
7616 ));
7617 }
7618
7619 #[test]
7620 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
7621 assert!(!round_is_clean(
7622 1,
7623 true,
7624 1,
7625 2,
7626 1,
7627 IncompleteReviewPolicy::Block
7628 ));
7629 assert!(!round_is_clean(
7630 0,
7631 false,
7632 1,
7633 2,
7634 1,
7635 IncompleteReviewPolicy::Block
7636 ));
7637 }
7638
7639 #[test]
7640 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
7641 assert!(!round_is_clean(
7645 0,
7646 true,
7647 0,
7648 2,
7649 2,
7650 IncompleteReviewPolicy::Block
7651 ));
7652 }
7653
7654 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
7655 CommandOutcome {
7656 command: "test".to_owned(),
7657 code,
7658 output_tail: String::new(),
7659 duration_ms: 0,
7660 resource_blocked,
7661 }
7662 }
7663
7664 #[test]
7665 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7666 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7667 assert!(
7668 !verify_inconclusive(&[outcome(Some(1), false)]),
7669 "an ordinary failure is still evidence about the patch"
7670 );
7671 assert!(verify_inconclusive(&[outcome(None, true)]));
7672 assert!(
7673 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7674 "one inconclusive outcome taints the whole batch"
7675 );
7676 assert!(!verify_inconclusive(&[]));
7677 }
7678
7679 #[tokio::test]
7680 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7681 let calls = std::sync::atomic::AtomicUsize::new(0);
7685 let started = Instant::now();
7686 wait_for_pids_with(
7687 &[123],
7688 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7689 Duration::from_millis(5),
7690 Duration::from_secs(5),
7691 )
7692 .await;
7693 assert!(
7694 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7695 "must keep checking rather than deciding on the first answer"
7696 );
7697 assert!(
7698 started.elapsed() < Duration::from_secs(1),
7699 "must return the moment it is confirmed dead, not wait out the ceiling"
7700 );
7701 }
7702
7703 #[tokio::test]
7704 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7705 let started = Instant::now();
7706 wait_for_pids_with(
7707 &[123],
7708 |_| true, Duration::from_millis(5),
7710 Duration::from_millis(30),
7711 )
7712 .await;
7713 let elapsed = started.elapsed();
7714 assert!(
7715 elapsed >= Duration::from_millis(30),
7716 "must not give up before its own ceiling: {elapsed:?}"
7717 );
7718 assert!(
7719 elapsed < Duration::from_secs(1),
7720 "must not wait past its own ceiling either: {elapsed:?}"
7721 );
7722 }
7723
7724 #[tokio::test]
7725 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7726 let started = Instant::now();
7727 wait_for_pids_with(
7728 &[],
7729 |_| true,
7730 Duration::from_secs(5),
7731 Duration::from_secs(5),
7732 )
7733 .await;
7734 assert!(
7735 started.elapsed() < Duration::from_millis(200),
7736 "an empty pid list has nothing to confirm"
7737 );
7738 }
7739
7740 fn review_round(
7746 clean: bool,
7747 blocking: usize,
7748 answered: usize,
7749 expected: usize,
7750 progressed: bool,
7751 e2e_ok: bool,
7752 ) -> ReviewRound {
7753 ReviewRound {
7754 round: 1,
7755 head: "h".to_owned(),
7756 verified_head: None,
7757 verified_at: None,
7758 reviews: Vec::new(),
7759 e2e: vec![CommandOutcome {
7760 command: "test".to_owned(),
7761 code: Some(if e2e_ok { 0 } else { 1 }),
7762 output_tail: String::new(),
7763 duration_ms: 0,
7764 resource_blocked: false,
7765 }],
7766 verify_retried: false,
7767 e2e_deferred: false,
7768 e2e_defer_reason: None,
7769 fix: None,
7770 blocking,
7771 answered,
7772 expected,
7773 clean,
7774 progressed,
7775 vote_split: false,
7776 reconsideration: Vec::new(),
7777 verdict: None,
7778 }
7779 }
7780
7781 #[test]
7782 fn review_conclusion_is_none_when_nothing_has_run() {
7783 assert_eq!(review_conclusion(&[], 3), None);
7784 }
7785
7786 #[test]
7787 fn review_conclusion_is_none_while_rounds_remain() {
7788 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7789 assert_eq!(review_conclusion(&rounds, 3), None);
7790 }
7791
7792 #[test]
7793 fn review_conclusion_is_gating_once_a_round_is_clean() {
7794 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7795 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7796 }
7797
7798 #[test]
7799 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7800 let rounds = vec![
7801 review_round(false, 1, 2, 2, true, true),
7802 review_round(false, 1, 2, 2, true, true),
7803 ];
7804 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7805 }
7806
7807 #[test]
7808 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7809 let rounds = vec![
7810 review_round(false, 1, 2, 2, true, true),
7811 review_round(false, 1, 2, 2, true, false),
7812 ];
7813 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7814 }
7815
7816 #[test]
7817 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7818 let mut blocked = review_round(false, 1, 2, 2, true, false);
7825 blocked.e2e[0].resource_blocked = true;
7826 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7827 assert_eq!(review_conclusion(&rounds, 2), None);
7828 }
7829
7830 #[test]
7831 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7832 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7834 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7835 }
7836
7837 #[test]
7838 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7839 let rounds = vec![
7840 review_round(false, 1, 2, 2, false, true),
7841 review_round(false, 1, 2, 2, false, true),
7842 ];
7843 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7844 }
7845
7846 fn secs(n: u64) -> Duration {
7847 Duration::from_secs(n)
7848 }
7849
7850 fn init_repo(dir: &Path) {
7853 let run = |args: &[&str]| {
7854 let out = std::process::Command::new("git")
7855 .args(args)
7856 .current_dir(dir)
7857 .quiet()
7858 .output()
7859 .expect("spawn git");
7860 assert!(
7861 out.status.success(),
7862 "git {args:?} failed: {}",
7863 String::from_utf8_lossy(&out.stderr)
7864 );
7865 };
7866 run(&["init", "-b", "main"]);
7867 run(&["config", "user.name", "magi test"]);
7868 run(&["config", "user.email", "magi@example.com"]);
7869 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7870 run(&["add", "-A"]);
7871 run(&["commit", "-m", "init"]);
7872 }
7873
7874 fn ask_test_home() {
7882 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7883 }
7884
7885 fn runner_at(status: RunStatus) -> Runner {
7888 let mut state = RunState::new(
7889 PathBuf::from("/nonexistent/repo"),
7890 "main".to_owned(),
7891 "deadbeef".to_owned(),
7892 "task".to_owned(),
7893 Config::default(),
7894 );
7895 state.status = status;
7896 Runner {
7897 state,
7898 roles: ResolvedRoles {
7899 implementers: Vec::new(),
7900 judges: Vec::new(),
7901 reviewers: Vec::new(),
7902 fixer: None,
7903 conductor: conductor(),
7904 implementer_roster: Vec::new(),
7905 },
7906 sem: Arc::new(Semaphore::new(1)),
7907 pause: Pause::new(),
7908 interrupt: Pause::new(),
7909 }
7910 }
7911
7912 #[test]
7916 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7917 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7918 let mut runner = runner_at(RunStatus::Implementing);
7919 let interrupt = Pause::new();
7920 runner.watch_interrupt(interrupt.clone());
7921
7922 interrupt.park_because("task a1b2 asked to run first");
7923
7924 assert!(runner.park_here().expect("park_here"));
7925 assert!(runner.state.parked);
7926 let last = runner.state.events.last().expect("a park event");
7927 assert_eq!(last.node, "park");
7928 assert!(
7929 last.message.contains("task a1b2 asked to run first"),
7930 "expected the interrupt reason in {:?}",
7931 last.message
7932 );
7933 }
7934
7935 #[test]
7943 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7944 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7945 let mut runner = runner_at(RunStatus::Implementing);
7946 let shutdown = Pause::new();
7947 runner.on_pause(shutdown.clone());
7948 let interrupt = Pause::new();
7949 runner.watch_interrupt(interrupt.clone());
7950
7951 assert!(!runner.park_here().expect("park_here"));
7953 assert!(!runner.state.parked);
7954
7955 interrupt.park_because("test");
7957 assert!(!shutdown.parked());
7958 assert!(runner.park_here().expect("park_here"));
7959 }
7960
7961 #[tokio::test]
7975 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7976 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7977 let mut runner = runner_at(RunStatus::Implementing);
7978 let interrupt = Pause::new();
7979 runner.watch_interrupt(interrupt.clone());
7980
7981 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7982 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7983
7984 let node = async move {
7988 started_tx.send(()).expect("send started");
7989 finish_rx.await.expect("recv finish");
7990 "node finished"
7991 };
7992
7993 let interrupter = async move {
7994 started_rx.await.expect("recv started");
7995 interrupt.park_because("higher-priority task waiting");
7997 tokio::task::yield_now().await;
8001 finish_tx.send(()).expect("send finish");
8002 };
8003
8004 let (node_result, ()) = tokio::join!(node, interrupter);
8005 assert_eq!(
8006 node_result, "node finished",
8007 "the in-flight call ran to completion"
8008 );
8009
8010 assert!(runner.park_here().expect("park_here"));
8013 assert!(runner.state.parked);
8014 }
8015
8016 #[test]
8022 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
8023 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
8024 let mut runner = runner_at(RunStatus::Judging);
8025 runner.state.config.agents = vec![conductor()];
8029 runner.state.candidates = vec![Candidate {
8030 index: 0,
8031 label: 'A',
8032 agent: "alpha".to_owned(),
8033 branch: "magi/x/A".to_owned(),
8034 worktree: PathBuf::from("/nonexistent/worktree"),
8035 summary: "did the thing".to_owned(),
8036 stat: "1 file changed".to_owned(),
8037 files: 1,
8038 commits: 1,
8039 empty: false,
8040 failed: None,
8041 verified_noop: None,
8042 duration_ms: 1234,
8043 folded: false,
8044 }];
8045 let run_id = runner.state.id.clone();
8046
8047 let interrupt = Pause::new();
8048 runner.watch_interrupt(interrupt.clone());
8049 interrupt.park_because("task c3d4 asked to run first");
8050 assert!(runner.park_here().expect("park_here"));
8051
8052 let resumed = Runner::resume(&run_id).expect("resume");
8053 assert_eq!(resumed.state.candidates.len(), 1);
8054 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
8055 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
8056 assert_eq!(resumed.state.status, runner.state.status);
8057 assert!(
8058 resumed.state.parked,
8059 "still parked until `execute` actually walks the graph again"
8060 );
8061 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
8062 }
8063
8064 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
8066 let mut q = ask::Question::new(
8067 run.to_owned(),
8068 "implement".to_owned(),
8069 "impl-A".to_owned(),
8070 "Which storage backend should the cache use?".to_owned(),
8071 String::new(),
8072 vec!["SQLite".to_owned(), "Redis".to_owned()],
8073 );
8074 store.put(&mut q).unwrap();
8075 q
8076 }
8077
8078 #[test]
8079 fn a_failed_runs_open_question_is_abandoned() {
8080 ask_test_home();
8081 let store = ask::Questions::open();
8082 let mut runner = runner_at(RunStatus::Failed);
8083 let run = runner.state.id.clone();
8084 let q = ask_open_question(&store, &run);
8085
8086 runner.settle_questions();
8087
8088 let back = store.get(&q.id).unwrap();
8089 assert!(
8090 !back.status.open(),
8091 "the seat that asked died with the run; nobody is left to read an answer"
8092 );
8093 assert!(
8094 back.detail.contains(&run) && back.detail.contains("failed"),
8095 "the reason names what the run became, not just that it is gone: {}",
8096 back.detail
8097 );
8098 }
8099
8100 #[test]
8101 fn a_merged_runs_open_question_is_abandoned_too() {
8102 ask_test_home();
8103 let store = ask::Questions::open();
8104 for status in [RunStatus::Merged, RunStatus::Ready] {
8107 let mut runner = runner_at(status);
8108 let run = runner.state.id.clone();
8109 let q = ask_open_question(&store, &run);
8110
8111 runner.settle_questions();
8112
8113 let back = store.get(&q.id).unwrap();
8114 assert!(
8115 !back.status.open(),
8116 "{status:?} run's question must not outlive the run"
8117 );
8118 }
8119 }
8120
8121 #[test]
8122 fn a_still_resumable_runs_open_question_is_left_alone() {
8123 ask_test_home();
8124 let store = ask::Questions::open();
8125 for status in [RunStatus::Blocked, RunStatus::Stalled] {
8131 let mut runner = runner_at(status);
8132 let run = runner.state.id.clone();
8133 let q = ask_open_question(&store, &run);
8134
8135 runner.settle_questions();
8136
8137 let back = store.get(&q.id).unwrap();
8138 assert!(
8139 back.status.open(),
8140 "{status:?} is still alive; the question must still be waiting"
8141 );
8142 }
8143 }
8144
8145 #[test]
8146 fn settle_questions_never_touches_an_already_answered_question() {
8147 ask_test_home();
8148 let store = ask::Questions::open();
8149 let mut runner = runner_at(RunStatus::Failed);
8150 let run = runner.state.id.clone();
8151 let mut q = ask_open_question(&store, &run);
8152 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
8153 .unwrap();
8154 store.put(&mut q).unwrap();
8155
8156 runner.settle_questions();
8161 runner.settle_questions();
8162
8163 let back = store.get(&q.id).unwrap();
8164 assert_eq!(
8165 back.status,
8166 ask::QuestionStatus::Answered,
8167 "a real answer is a decision on record, never overwritten by a sweep"
8168 );
8169 }
8170
8171 #[tokio::test]
8182 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
8183 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
8184 let tmp = tempfile::tempdir().expect("tempdir");
8185 let repo = tmp.path().join("repo");
8186 std::fs::create_dir_all(&repo).unwrap();
8187 init_repo(&repo);
8188
8189 let mut config = Config::default();
8190 config.graph.worktree_root = Some(tmp.path().join("wt"));
8191
8192 let mut state = RunState::new(
8193 repo.clone(),
8194 "main".to_owned(),
8195 "deadbeef".to_owned(),
8196 "task".to_owned(),
8197 config,
8198 );
8199 let root = state.worktree_root();
8200 let wt_a = root.join("cand-A");
8201 let wt_b = root.join("cand-B");
8202 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
8203 .await
8204 .expect("worktree A");
8205 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
8206 .await
8207 .expect("worktree B");
8208
8209 state.candidates = vec![
8210 Candidate {
8211 index: 0,
8212 label: 'A',
8213 agent: "alpha".to_owned(),
8214 branch: "magi/x/A".to_owned(),
8215 worktree: wt_a.clone(),
8216 summary: String::new(),
8217 stat: String::new(),
8218 files: 0,
8219 commits: 0,
8220 empty: false,
8221 failed: None,
8222 verified_noop: None,
8223 duration_ms: 0,
8224 folded: false,
8225 },
8226 Candidate {
8227 index: 1,
8228 label: 'B',
8229 agent: "beta".to_owned(),
8230 branch: "magi/x/B".to_owned(),
8231 worktree: wt_b.clone(),
8232 summary: String::new(),
8233 stat: String::new(),
8234 files: 0,
8235 commits: 0,
8236 empty: false,
8237 failed: None,
8238 verified_noop: None,
8239 duration_ms: 0,
8240 folded: false,
8241 },
8242 ];
8243 state.tally = Some(Tally {
8244 first_choice: BTreeMap::from([('A', 1)]),
8245 borda: BTreeMap::new(),
8246 winner: 'A',
8247 rankings: 1,
8248 unanimous_initial: true,
8249 deliberated: false,
8250 changed_votes: 0,
8251 unanimous_final: true,
8252 tie_break: None,
8253 judges: 1,
8254 present: 1,
8255 quorum: 1,
8256 met_quorum: true,
8257 uncontested: None,
8258 });
8259 state.status = RunStatus::Ready;
8260
8261 fold_run(&mut state, false, &crate::run::home())
8262 .await
8263 .expect("fold_run");
8264
8265 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
8266 assert!(
8267 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
8268 "the unmerged winner's branch survives"
8269 );
8270 assert!(
8271 !state.candidates[0].folded,
8272 "the winner is not marked folded"
8273 );
8274
8275 assert!(!wt_b.exists(), "the loser's worktree is removed");
8276 assert!(
8277 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
8278 "the loser's branch is removed"
8279 );
8280 assert!(state.candidates[1].folded, "the loser is marked folded");
8281 }
8282
8283 #[tokio::test]
8286 async fn fold_run_keeps_a_branch_that_was_handed_to_a_later_run() {
8287 let tmp = tempfile::tempdir().expect("tempdir");
8288 let repo = tmp.path().join("repo");
8289 std::fs::create_dir_all(&repo).unwrap();
8290 init_repo(&repo);
8291 let home = tmp.path().join("home");
8292
8293 let mut config = Config::default();
8294 config.graph.worktree_root = Some(tmp.path().join("wt"));
8295 let mut state = RunState::new(
8296 repo.clone(),
8297 "main".to_owned(),
8298 "deadbeef".to_owned(),
8299 "task".to_owned(),
8300 config,
8301 );
8302 git::git(&repo, &["branch", "magi/x/A", "main"])
8304 .await
8305 .expect("branch");
8306 state.candidates = vec![Candidate {
8307 index: 0,
8308 label: 'A',
8309 agent: "alpha".to_owned(),
8310 branch: "magi/x/A".to_owned(),
8311 worktree: state.worktree_root().join("cand-A"),
8312 summary: String::new(),
8313 stat: String::new(),
8314 files: 0,
8315 commits: 0,
8316 empty: false,
8317 failed: None,
8318 verified_noop: None,
8319 duration_ms: 0,
8320 folded: true,
8321 }];
8322 state.released_to = Some("20260901-000000-new1".to_owned());
8323 state.released_branches = vec!["magi/x/A".to_owned()];
8324
8325 fold_run(&mut state, true, &home).await.expect("fold_run");
8326
8327 assert!(
8328 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
8329 "the handed-over branch survives a fold"
8330 );
8331 }
8332
8333 #[tokio::test]
8337 async fn an_empty_winner_is_detected_before_a_pull_request_is_attempted() {
8338 let tmp = tempfile::tempdir().expect("tempdir");
8339 let repo = tmp.path().join("repo");
8340 std::fs::create_dir_all(&repo).unwrap();
8341 init_repo(&repo);
8342 let run = |args: &[&str]| {
8343 let out = std::process::Command::new("git")
8344 .quiet()
8345 .args(args)
8346 .current_dir(&repo)
8347 .output()
8348 .expect("spawn git");
8349 assert!(out.status.success(), "git {args:?}");
8350 };
8351 run(&["branch", "magi/x/A"]);
8352 run(&["checkout", "-q", "-b", "magi/x/B"]);
8353 std::fs::write(repo.join("f.txt"), "x\n").unwrap();
8354 run(&["add", "-A"]);
8355 run(&["commit", "-q", "-m", "work"]);
8356 run(&["checkout", "-q", "main"]);
8357
8358 let mut state = RunState::new(
8359 repo.clone(),
8360 "main".to_owned(),
8361 "deadbeef".to_owned(),
8362 "task".to_owned(),
8363 Config::default(),
8364 );
8365 state.seeds = vec![refs::Seed {
8366 token: "magi/27b2/A".to_owned(),
8367 kind: refs::SeedKind::Unresolved,
8368 sha: String::new(),
8369 branch: true,
8370 detail: "no branch or commit named magi/27b2/A".to_owned(),
8371 }];
8372
8373 assert!(merge_is_empty(&repo, &state, "magi/x/A", MergeMode::Pr).await);
8374 assert!(merge_is_empty(&repo, &state, "magi/x/A", MergeMode::Local).await);
8375 assert!(!merge_is_empty(&repo, &state, "magi/x/B", MergeMode::Pr).await);
8376 let detail = empty_candidate_detail(&state, "main");
8377 assert!(detail.starts_with("empty candidate"), "{detail}");
8378 assert!(detail.contains("magi/27b2/A"), "{detail}");
8379 }
8380
8381 #[tokio::test]
8390 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
8391 let tmp = tempfile::tempdir().expect("tempdir");
8392 let repo = tmp.path().join("repo");
8393 std::fs::create_dir_all(&repo).unwrap();
8394 init_repo(&repo);
8395
8396 let mut config = Config::default();
8397 config.merge.mode = MergeMode::Local;
8398
8399 let mut state = RunState::new(
8400 repo.clone(),
8401 "main".to_owned(),
8402 "deadbeef".to_owned(),
8403 "task".to_owned(),
8404 config,
8405 );
8406 state.candidates = vec![Candidate {
8407 index: 0,
8408 label: 'A',
8409 agent: "alpha".to_owned(),
8410 branch: "does-not-exist".to_owned(),
8411 worktree: repo.clone(),
8412 summary: String::new(),
8413 stat: String::new(),
8414 files: 0,
8415 commits: 0,
8416 empty: false,
8417 failed: None,
8418 verified_noop: None,
8419 duration_ms: 0,
8420 folded: false,
8421 }];
8422 state.tally = Some(Tally {
8423 first_choice: BTreeMap::from([('A', 1)]),
8424 borda: BTreeMap::new(),
8425 winner: 'A',
8426 rankings: 1,
8427 unanimous_initial: true,
8428 deliberated: false,
8429 changed_votes: 0,
8430 unanimous_final: true,
8431 tie_break: None,
8432 judges: 0,
8433 present: 0,
8434 quorum: 0,
8435 met_quorum: true,
8436 uncontested: Some("only candidate A produced a change".to_owned()),
8437 });
8438 state.reviews = vec![ReviewRound {
8439 round: 1,
8440 head: "deadbeef".to_owned(),
8441 verified_head: None,
8442 verified_at: None,
8443 reviews: Vec::new(),
8444 e2e: Vec::new(),
8445 fix: None,
8446 blocking: 0,
8447 answered: 0,
8448 expected: 0,
8449 clean: true,
8450 verify_retried: false,
8451 e2e_deferred: false,
8452 e2e_defer_reason: None,
8453 progressed: false,
8454 vote_split: false,
8455 reconsideration: Vec::new(),
8456 verdict: None,
8457 }];
8458 state.gate = vec![CommandOutcome {
8459 command: "test".to_owned(),
8460 code: Some(0),
8461 output_tail: String::new(),
8462 duration_ms: 0,
8463 resource_blocked: false,
8464 }];
8465 state.gate_ran = true;
8466 state.status = RunStatus::Ready;
8471 state.merge = Some(MergeOutcome {
8472 mode: MergeMode::Local,
8473 ok: false,
8474 detail: "already concluded".to_owned(),
8475 empty: false,
8476 });
8477
8478 let mut runner = Runner {
8479 state,
8480 roles: ResolvedRoles {
8481 implementers: Vec::new(),
8482 judges: Vec::new(),
8483 reviewers: Vec::new(),
8484 fixer: None,
8485 conductor: conductor(),
8486 implementer_roster: Vec::new(),
8487 },
8488 sem: Arc::new(Semaphore::new(1)),
8489 pause: Pause::new(),
8490 interrupt: Pause::new(),
8491 };
8492
8493 runner.merge().await.expect("merge");
8494
8495 assert_eq!(
8496 runner.state.status,
8497 RunStatus::Ready,
8498 "a concluded run's status must not change on reentry"
8499 );
8500 assert_eq!(
8501 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8502 Some("already concluded"),
8503 "merge must not run again once the node already recorded an outcome"
8504 );
8505 }
8506
8507 #[tokio::test]
8516 async fn merge_refuses_a_gate_that_has_not_actually_run() {
8517 let tmp = tempfile::tempdir().expect("tempdir");
8518 let repo = tmp.path().join("repo");
8519 std::fs::create_dir_all(&repo).unwrap();
8520 init_repo(&repo);
8521
8522 let mut config = Config::default();
8523 config.merge.mode = MergeMode::Local;
8524
8525 let mut state = RunState::new(
8526 repo.clone(),
8527 "main".to_owned(),
8528 "deadbeef".to_owned(),
8529 "task".to_owned(),
8530 config,
8531 );
8532 state.candidates = vec![Candidate {
8533 index: 0,
8534 label: 'A',
8535 agent: "alpha".to_owned(),
8536 branch: "does-not-exist".to_owned(),
8537 worktree: repo.clone(),
8538 summary: String::new(),
8539 stat: String::new(),
8540 files: 0,
8541 commits: 0,
8542 empty: false,
8543 failed: None,
8544 verified_noop: None,
8545 duration_ms: 0,
8546 folded: false,
8547 }];
8548 state.tally = Some(Tally {
8549 first_choice: BTreeMap::from([('A', 1)]),
8550 borda: BTreeMap::new(),
8551 winner: 'A',
8552 rankings: 1,
8553 unanimous_initial: true,
8554 deliberated: false,
8555 changed_votes: 0,
8556 unanimous_final: true,
8557 tie_break: None,
8558 judges: 0,
8559 present: 0,
8560 quorum: 0,
8561 met_quorum: true,
8562 uncontested: Some("only candidate A produced a change".to_owned()),
8563 });
8564 state.reviews = vec![ReviewRound {
8565 round: 1,
8566 head: "deadbeef".to_owned(),
8567 verified_head: None,
8568 verified_at: None,
8569 reviews: Vec::new(),
8570 e2e: Vec::new(),
8571 fix: None,
8572 blocking: 0,
8573 answered: 0,
8574 expected: 0,
8575 clean: true,
8576 verify_retried: false,
8577 e2e_deferred: false,
8578 e2e_defer_reason: None,
8579 progressed: false,
8580 vote_split: false,
8581 reconsideration: Vec::new(),
8582 verdict: None,
8583 }];
8584 state.gate = Vec::new();
8586 state.gate_ran = false;
8587 state.status = RunStatus::Gating;
8588
8589 let mut runner = Runner {
8590 state,
8591 roles: ResolvedRoles {
8592 implementers: Vec::new(),
8593 judges: Vec::new(),
8594 reviewers: Vec::new(),
8595 fixer: None,
8596 conductor: conductor(),
8597 implementer_roster: Vec::new(),
8598 },
8599 sem: Arc::new(Semaphore::new(1)),
8600 pause: Pause::new(),
8601 interrupt: Pause::new(),
8602 };
8603
8604 runner.merge().await.expect("merge");
8605
8606 assert!(
8607 runner.state.merge.is_none(),
8608 "an empty gate must never be read as a passing one: {:?}",
8609 runner.state.merge
8610 );
8611 }
8612
8613 #[tokio::test]
8620 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
8621 let tmp = tempfile::tempdir().expect("tempdir");
8622 let repo = tmp.path().join("repo");
8623 std::fs::create_dir_all(&repo).unwrap();
8624 init_repo(&repo);
8625
8626 let config = Config::default();
8628
8629 let mut state = RunState::new(
8630 repo.clone(),
8631 "main".to_owned(),
8632 "deadbeef".to_owned(),
8633 "task".to_owned(),
8634 config,
8635 );
8636 state.candidates = vec![Candidate {
8637 index: 0,
8638 label: 'A',
8639 agent: "alpha".to_owned(),
8640 branch: "does-not-exist".to_owned(),
8641 worktree: repo.clone(),
8642 summary: String::new(),
8643 stat: String::new(),
8644 files: 0,
8645 commits: 0,
8646 empty: false,
8647 failed: None,
8648 verified_noop: None,
8649 duration_ms: 0,
8650 folded: false,
8651 }];
8652 state.tally = Some(Tally {
8653 first_choice: BTreeMap::from([('A', 1)]),
8654 borda: BTreeMap::new(),
8655 winner: 'A',
8656 rankings: 1,
8657 unanimous_initial: true,
8658 deliberated: false,
8659 changed_votes: 0,
8660 unanimous_final: true,
8661 tie_break: None,
8662 judges: 0,
8663 present: 0,
8664 quorum: 0,
8665 met_quorum: true,
8666 uncontested: Some("only candidate A produced a change".to_owned()),
8667 });
8668 state.reviews = vec![ReviewRound {
8669 round: 1,
8670 head: "deadbeef".to_owned(),
8671 verified_head: None,
8672 verified_at: None,
8673 reviews: Vec::new(),
8674 e2e: Vec::new(),
8675 fix: None,
8676 blocking: 0,
8677 answered: 0,
8678 expected: 0,
8679 clean: true,
8680 verify_retried: false,
8681 e2e_deferred: false,
8682 e2e_defer_reason: None,
8683 progressed: false,
8684 vote_split: false,
8685 reconsideration: Vec::new(),
8686 verdict: None,
8687 }];
8688
8689 let mut runner = Runner {
8690 state,
8691 roles: ResolvedRoles {
8692 implementers: Vec::new(),
8693 judges: Vec::new(),
8694 reviewers: Vec::new(),
8695 fixer: None,
8696 conductor: conductor(),
8697 implementer_roster: Vec::new(),
8698 },
8699 sem: Arc::new(Semaphore::new(1)),
8700 pause: Pause::new(),
8701 interrupt: Pause::new(),
8702 };
8703
8704 runner.gate().await.expect("gate");
8705 assert!(
8706 runner.state.gate_ran,
8707 "zero configured commands is still a real attempt, not an unrun gate"
8708 );
8709 assert!(runner.state.gate.is_empty());
8710 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
8711 assert_ne!(
8712 runner.state.status,
8713 RunStatus::Blocked,
8714 "a gate with nothing to check must not read as failed"
8715 );
8716
8717 runner.merge().await.expect("merge");
8718 assert_eq!(
8719 runner.state.status,
8720 RunStatus::Ready,
8721 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
8722 );
8723 }
8724
8725 #[tokio::test]
8736 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
8737 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8738 let home = crate::run::home();
8739
8740 let tmp = tempfile::tempdir().expect("tempdir");
8741 let repo = tmp.path().join("repo");
8742 std::fs::create_dir_all(&repo).unwrap();
8743 init_repo(&repo);
8744 let cache_dir = tmp.path().join("target");
8747
8748 let mut config = Config::default();
8749 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
8750 config.graph.timeout_verify = Some(2);
8753
8754 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8755 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8756 .expect("no io error acquiring directly")
8757 {
8758 crate::cache::AcquireOutcome::Acquired(g) => g,
8759 crate::cache::AcquireOutcome::Busy(b) => {
8760 panic!("expected the direct acquire to win the lease first: {b:?}")
8761 }
8762 };
8763
8764 let mut state = RunState::new(
8765 repo.clone(),
8766 "main".to_owned(),
8767 "deadbeef".to_owned(),
8768 "task".to_owned(),
8769 config,
8770 );
8771 state.candidates = vec![Candidate {
8772 index: 0,
8773 label: 'A',
8774 agent: "alpha".to_owned(),
8775 branch: "does-not-exist".to_owned(),
8776 worktree: repo.clone(),
8777 summary: String::new(),
8778 stat: String::new(),
8779 files: 0,
8780 commits: 0,
8781 empty: false,
8782 failed: None,
8783 verified_noop: None,
8784 duration_ms: 0,
8785 folded: false,
8786 }];
8787 state.tally = Some(Tally {
8788 first_choice: BTreeMap::from([('A', 1)]),
8789 borda: BTreeMap::new(),
8790 winner: 'A',
8791 rankings: 1,
8792 unanimous_initial: true,
8793 deliberated: false,
8794 changed_votes: 0,
8795 unanimous_final: true,
8796 tie_break: None,
8797 judges: 0,
8798 present: 0,
8799 quorum: 0,
8800 met_quorum: true,
8801 uncontested: Some("only candidate A produced a change".to_owned()),
8802 });
8803 state.reviews = vec![ReviewRound {
8804 round: 1,
8805 head: "deadbeef".to_owned(),
8806 verified_head: None,
8807 verified_at: None,
8808 reviews: Vec::new(),
8809 e2e: Vec::new(),
8810 fix: None,
8811 blocking: 0,
8812 answered: 0,
8813 expected: 0,
8814 clean: true,
8815 verify_retried: false,
8816 e2e_deferred: false,
8817 e2e_defer_reason: None,
8818 progressed: false,
8819 vote_split: false,
8820 reconsideration: Vec::new(),
8821 verdict: None,
8822 }];
8823
8824 let mut runner = Runner {
8825 state,
8826 roles: ResolvedRoles {
8827 implementers: Vec::new(),
8828 judges: Vec::new(),
8829 reviewers: Vec::new(),
8830 fixer: None,
8831 conductor: conductor(),
8832 implementer_roster: Vec::new(),
8833 },
8834 sem: Arc::new(Semaphore::new(1)),
8835 pause: Pause::new(),
8836 interrupt: Pause::new(),
8837 };
8838
8839 let started = std::time::Instant::now();
8840 runner.gate().await.expect("gate");
8841 assert!(
8842 started.elapsed() < Duration::from_secs(1),
8843 "a gate with nothing to run must never wait on a lease it never needed"
8844 );
8845 assert!(
8846 runner.state.gate_ran,
8847 "zero commands is still a real, immediate attempt"
8848 );
8849 assert!(runner.state.gate.is_empty());
8850 assert_ne!(
8851 runner.state.status,
8852 RunStatus::Blocked,
8853 "must not read as resource-blocked on a lease it never asked for"
8854 );
8855 }
8856
8857 #[tokio::test]
8867 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8868 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8869
8870 let tmp = tempfile::tempdir().expect("tempdir");
8871 let repo = tmp.path().join("repo");
8872 std::fs::create_dir_all(&repo).unwrap();
8873 init_repo(&repo);
8874
8875 let mut config = Config::default();
8876 config.verify.gate = vec![
8877 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8878 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8879 .to_owned(),
8880 ];
8881
8882 let mut state = RunState::new(
8883 repo.clone(),
8884 "main".to_owned(),
8885 "deadbeef".to_owned(),
8886 "task".to_owned(),
8887 config,
8888 );
8889 let run_id = state.id.clone();
8890 state.candidates = vec![Candidate {
8891 index: 0,
8892 label: 'A',
8893 agent: "alpha".to_owned(),
8894 branch: "does-not-exist".to_owned(),
8895 worktree: repo.clone(),
8896 summary: String::new(),
8897 stat: String::new(),
8898 files: 0,
8899 commits: 0,
8900 empty: false,
8901 failed: None,
8902 verified_noop: None,
8903 duration_ms: 0,
8904 folded: false,
8905 }];
8906 state.tally = Some(Tally {
8907 first_choice: BTreeMap::from([('A', 1)]),
8908 borda: BTreeMap::new(),
8909 winner: 'A',
8910 rankings: 1,
8911 unanimous_initial: true,
8912 deliberated: false,
8913 changed_votes: 0,
8914 unanimous_final: true,
8915 tie_break: None,
8916 judges: 0,
8917 present: 0,
8918 quorum: 0,
8919 met_quorum: true,
8920 uncontested: Some("only candidate A produced a change".to_owned()),
8921 });
8922 state.reviews = vec![ReviewRound {
8923 round: 1,
8924 head: "deadbeef".to_owned(),
8925 verified_head: None,
8926 verified_at: None,
8927 reviews: Vec::new(),
8928 e2e: Vec::new(),
8929 fix: None,
8930 blocking: 0,
8931 answered: 0,
8932 expected: 0,
8933 clean: true,
8934 verify_retried: false,
8935 e2e_deferred: false,
8936 e2e_defer_reason: None,
8937 progressed: false,
8938 vote_split: false,
8939 reconsideration: Vec::new(),
8940 verdict: None,
8941 }];
8942
8943 let mut runner = Runner {
8944 state,
8945 roles: ResolvedRoles {
8946 implementers: Vec::new(),
8947 judges: Vec::new(),
8948 reviewers: Vec::new(),
8949 fixer: None,
8950 conductor: conductor(),
8951 implementer_roster: Vec::new(),
8952 },
8953 sem: Arc::new(Semaphore::new(1)),
8954 pause: Pause::new(),
8955 interrupt: Pause::new(),
8956 };
8957
8958 let started_marker = repo.join("started.marker");
8959 let release_marker = repo.join("release.marker");
8960 let poller = tokio::spawn(async move {
8961 for _ in 0..100 {
8966 if started_marker.exists()
8967 && let Ok(s) = crate::run::RunState::load(&run_id)
8968 && let Some(a) = s.active.get("gate")
8969 {
8970 std::fs::write(&release_marker, b"go").expect("release marker");
8971 return Some(a.clone());
8972 }
8973 tokio::time::sleep(Duration::from_millis(50)).await;
8974 }
8975 None
8976 });
8977
8978 runner.gate().await.expect("gate");
8979 let captured = poller.await.expect("poller task");
8980 let captured = captured.expect(
8981 "the poller never saw a `gate` task entry in run.json while the command was \
8982 still blocked on its own release marker",
8983 );
8984
8985 assert_eq!(captured.task.as_deref(), Some("gate"));
8986 assert_eq!(captured.node, "gate");
8987 assert_eq!(captured.index, Some(1));
8988 assert_eq!(captured.total, Some(1));
8989 assert!(
8990 captured
8991 .command
8992 .as_deref()
8993 .is_some_and(|c| c.contains("started.marker")),
8994 "{captured:?}"
8995 );
8996
8997 assert!(
8998 runner.state.active.is_empty(),
8999 "the entry must be cleared once the command actually finished: {:?}",
9000 runner.state.active
9001 );
9002 assert!(runner.state.gate_ran);
9003 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
9004 }
9005
9006 #[tokio::test]
9019 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
9020 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
9021 let home = crate::run::home();
9022
9023 let tmp = tempfile::tempdir().expect("tempdir");
9024 let repo = tmp.path().join("repo");
9025 std::fs::create_dir_all(&repo).unwrap();
9026 init_repo(&repo);
9027 let head = crate::git::rev_parse(&repo, "HEAD")
9028 .await
9029 .expect("rev-parse");
9030 let cache_dir = tmp.path().join("target");
9033
9034 let mut config = Config::default();
9035 config.verify.e2e = vec![format!(
9036 "CARGO_TARGET_DIR='{}' test -f README.md",
9037 cache_dir.display()
9038 )];
9039 config.graph.review_rounds = 1;
9040 config.graph.timeout_verify = Some(2);
9043
9044 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
9045 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
9046 .expect("no io error acquiring directly")
9047 {
9048 crate::cache::AcquireOutcome::Acquired(g) => g,
9049 crate::cache::AcquireOutcome::Busy(b) => {
9050 panic!("expected the direct acquire to win the lease first: {b:?}")
9051 }
9052 };
9053
9054 let mut state = RunState::new(
9055 repo.clone(),
9056 "main".to_owned(),
9057 head.clone(),
9058 "task".to_owned(),
9059 config,
9060 );
9061 state.candidates = vec![Candidate {
9062 index: 0,
9063 label: 'A',
9064 agent: "alpha".to_owned(),
9065 branch: "does-not-exist".to_owned(),
9066 worktree: repo.clone(),
9067 summary: String::new(),
9068 stat: String::new(),
9069 files: 0,
9070 commits: 0,
9071 empty: false,
9072 failed: None,
9073 verified_noop: None,
9074 duration_ms: 0,
9075 folded: false,
9076 }];
9077 state.tally = Some(Tally {
9078 first_choice: BTreeMap::from([('A', 1)]),
9079 borda: BTreeMap::new(),
9080 winner: 'A',
9081 rankings: 1,
9082 unanimous_initial: true,
9083 deliberated: false,
9084 changed_votes: 0,
9085 unanimous_final: true,
9086 tie_break: None,
9087 judges: 0,
9088 present: 0,
9089 quorum: 0,
9090 met_quorum: true,
9091 uncontested: Some("only candidate A produced a change".to_owned()),
9092 });
9093 state.reviews = vec![ReviewRound {
9097 round: 1,
9098 head: head.clone(),
9099 verified_head: None,
9100 verified_at: None,
9101 reviews: Vec::new(),
9102 e2e: Vec::new(),
9103 fix: None,
9104 blocking: 1,
9105 answered: 1,
9106 expected: 1,
9107 clean: false,
9108 verify_retried: false,
9109 e2e_deferred: true,
9110 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
9111 progressed: false,
9112 vote_split: false,
9113 reconsideration: Vec::new(),
9114 verdict: None,
9115 }];
9116
9117 let mut runner = Runner {
9118 state,
9119 roles: ResolvedRoles {
9120 implementers: Vec::new(),
9121 judges: Vec::new(),
9122 reviewers: Vec::new(),
9123 fixer: None,
9124 conductor: conductor(),
9125 implementer_roster: Vec::new(),
9126 },
9127 sem: Arc::new(Semaphore::new(1)),
9128 pause: Pause::new(),
9129 interrupt: Pause::new(),
9130 };
9131
9132 let shell = runner.state.config.shell();
9133 runner
9134 .stop_reviewing("round budget spent", &shell, &repo)
9135 .await
9136 .expect("stop_reviewing");
9137
9138 let last = runner.state.reviews.last().expect("round record");
9139 assert_eq!(
9140 last.e2e_status(),
9141 E2eStatus::ResourceBlocked,
9142 "the shared cache is still held; the attempt must read as blocked, not deferred or \
9143 failed: {last:?}"
9144 );
9145 assert_eq!(
9146 last.verified_head.as_deref(),
9147 Some(head.as_str()),
9148 "which commit this attempt targeted is known even though nothing finished checking \
9149 it"
9150 );
9151 let first_attempt_at = last
9152 .verified_at
9153 .expect("when this attempt ran is known too");
9154 assert_ne!(
9155 runner.state.status,
9156 RunStatus::Blocked,
9157 "contention is evidence about the machine, not the patch — it must not settle the \
9158 run as blocked: {:?}",
9159 runner.state.status
9160 );
9161 assert!(
9162 !runner
9163 .state
9164 .events
9165 .iter()
9166 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
9167 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
9168 runner.state.events
9169 );
9170
9171 runner
9176 .stop_reviewing("round budget spent", &shell, &repo)
9177 .await
9178 .expect("stop_reviewing retry");
9179 assert_eq!(
9180 runner.state.reviews.len(),
9181 1,
9182 "no new round was started: {:?}",
9183 runner.state.reviews
9184 );
9185 let last = runner.state.reviews.last().expect("round record");
9186 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
9187 assert!(
9188 last.verified_at.expect("still known") > first_attempt_at,
9189 "a second reentry must be a fresh attempt, not a stale copy of the first"
9190 );
9191 assert_ne!(runner.state.status, RunStatus::Blocked);
9192
9193 held.release();
9194 }
9195
9196 #[tokio::test]
9208 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
9209 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
9210 let home = crate::run::home();
9211
9212 let tmp = tempfile::tempdir().expect("tempdir");
9213 let repo = tmp.path().join("repo");
9214 std::fs::create_dir_all(&repo).unwrap();
9215 init_repo(&repo);
9216 let head = crate::git::rev_parse(&repo, "HEAD")
9217 .await
9218 .expect("rev-parse");
9219 let cache_dir = tmp.path().join("target");
9220
9221 let mut config = Config::default();
9222 config.verify.e2e = vec![format!(
9223 "CARGO_TARGET_DIR='{}' test -f README.md",
9224 cache_dir.display()
9225 )];
9226 config.graph.review_rounds = 1;
9227 config.graph.timeout_verify = Some(2);
9228
9229 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
9230 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
9231 .expect("no io error acquiring directly")
9232 {
9233 crate::cache::AcquireOutcome::Acquired(g) => g,
9234 crate::cache::AcquireOutcome::Busy(b) => {
9235 panic!("expected the direct acquire to win the lease first: {b:?}")
9236 }
9237 };
9238
9239 let mut state = RunState::new(
9240 repo.clone(),
9241 "main".to_owned(),
9242 head.clone(),
9243 "task".to_owned(),
9244 config,
9245 );
9246 state.candidates = vec![Candidate {
9247 index: 0,
9248 label: 'A',
9249 agent: "alpha".to_owned(),
9250 branch: "does-not-exist".to_owned(),
9251 worktree: repo.clone(),
9252 summary: String::new(),
9253 stat: String::new(),
9254 files: 0,
9255 commits: 0,
9256 empty: false,
9257 failed: None,
9258 verified_noop: None,
9259 duration_ms: 0,
9260 folded: false,
9261 }];
9262 state.tally = Some(Tally {
9263 first_choice: BTreeMap::from([('A', 1)]),
9264 borda: BTreeMap::new(),
9265 winner: 'A',
9266 rankings: 1,
9267 unanimous_initial: true,
9268 deliberated: false,
9269 changed_votes: 0,
9270 unanimous_final: true,
9271 tie_break: None,
9272 judges: 0,
9273 present: 0,
9274 quorum: 0,
9275 met_quorum: true,
9276 uncontested: Some("only candidate A produced a change".to_owned()),
9277 });
9278 state.reviews = vec![ReviewRound {
9282 round: 1,
9283 head: head.clone(),
9284 verified_head: Some(head.clone()),
9285 verified_at: Some(jiff::Timestamp::now()),
9286 reviews: Vec::new(),
9287 e2e: vec![CommandOutcome {
9288 command: format!(
9289 "CARGO_TARGET_DIR='{}' test -f README.md",
9290 cache_dir.display()
9291 ),
9292 code: None,
9293 output_tail: "waiting for the shared build cache".to_owned(),
9294 duration_ms: 0,
9295 resource_blocked: true,
9296 }],
9297 fix: None,
9298 blocking: 1,
9299 answered: 1,
9300 expected: 1,
9301 clean: false,
9302 verify_retried: false,
9303 e2e_deferred: false,
9304 e2e_defer_reason: None,
9305 progressed: false,
9306 vote_split: false,
9307 reconsideration: Vec::new(),
9308 verdict: None,
9309 }];
9310
9311 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
9312 let mut runner = Runner {
9313 state,
9314 roles: ResolvedRoles {
9315 implementers: Vec::new(),
9316 judges: Vec::new(),
9317 reviewers: Vec::new(),
9318 fixer: None,
9319 conductor: conductor(),
9320 implementer_roster: Vec::new(),
9321 },
9322 sem: Arc::new(Semaphore::new(1)),
9323 pause: Pause::new(),
9324 interrupt: Pause::new(),
9325 };
9326
9327 runner.review_loop().await.expect("review_loop");
9332
9333 assert_eq!(
9334 runner.state.reviews.len(),
9335 1,
9336 "no new round was started on top of the unresolved one: {:?}",
9337 runner.state.reviews
9338 );
9339 let last = &runner.state.reviews[0];
9340 assert_eq!(
9341 last.e2e_status(),
9342 E2eStatus::ResourceBlocked,
9343 "still contended: {last:?}"
9344 );
9345 assert!(
9346 last.verified_at.expect("still known") > first_attempt_at,
9347 "review_loop must have actually retried the check, not left it exactly as found"
9348 );
9349 assert_ne!(
9350 runner.state.status,
9351 RunStatus::Blocked,
9352 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
9353 runner.state.status
9354 );
9355
9356 held.release();
9357 }
9358
9359 #[tokio::test]
9360 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
9361 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
9362 let tmp = tempfile::tempdir().expect("tempdir");
9363 let repo = tmp.path().join("repo");
9364 std::fs::create_dir_all(&repo).unwrap();
9365 init_repo(&repo);
9366
9367 let mut config = Config::default();
9368 config.merge.mode = MergeMode::Pr;
9369 config.graph.land = true;
9370 config.graph.land_approval = false;
9371
9372 let mut state = RunState::new(
9373 repo.clone(),
9374 "main".to_owned(),
9375 "deadbeef".to_owned(),
9376 "task".to_owned(),
9377 config,
9378 );
9379 state.candidates = vec![Candidate {
9380 index: 0,
9381 label: 'A',
9382 agent: "alpha".to_owned(),
9383 branch: "does-not-exist".to_owned(),
9384 worktree: repo.clone(),
9385 summary: String::new(),
9386 stat: String::new(),
9387 files: 0,
9388 commits: 0,
9389 empty: false,
9390 failed: None,
9391 verified_noop: None,
9392 duration_ms: 0,
9393 folded: false,
9394 }];
9395 state.tally = Some(Tally {
9396 first_choice: BTreeMap::from([('A', 1)]),
9397 borda: BTreeMap::new(),
9398 winner: 'A',
9399 rankings: 1,
9400 unanimous_initial: true,
9401 deliberated: false,
9402 changed_votes: 0,
9403 unanimous_final: true,
9404 tie_break: None,
9405 judges: 0,
9406 present: 0,
9407 quorum: 0,
9408 met_quorum: true,
9409 uncontested: Some("only candidate A produced a change".to_owned()),
9410 });
9411 state.reviews = vec![ReviewRound {
9412 round: 1,
9413 head: "deadbeef".to_owned(),
9414 verified_head: None,
9415 verified_at: None,
9416 reviews: Vec::new(),
9417 e2e: Vec::new(),
9418 fix: None,
9419 blocking: 0,
9420 answered: 0,
9421 expected: 0,
9422 clean: true,
9423 verify_retried: false,
9424 e2e_deferred: false,
9425 e2e_defer_reason: None,
9426 progressed: false,
9427 vote_split: false,
9428 reconsideration: Vec::new(),
9429 verdict: None,
9430 }];
9431 state.gate = vec![CommandOutcome {
9432 command: "test".to_owned(),
9433 code: Some(0),
9434 output_tail: String::new(),
9435 duration_ms: 0,
9436 resource_blocked: false,
9437 }];
9438 state.gate_ran = true;
9439 state.status = RunStatus::Landing;
9443 state.merge = Some(MergeOutcome {
9444 mode: MergeMode::Pr,
9445 ok: true,
9446 detail: "https://example.invalid/x/y/pull/1".to_owned(),
9447 empty: false,
9448 });
9449
9450 ask_test_home();
9454 let store = ask::Questions::open();
9455 let q = ask_open_question(&store, &state.id);
9456
9457 let mut runner = Runner {
9458 state,
9459 roles: ResolvedRoles {
9460 implementers: Vec::new(),
9461 judges: Vec::new(),
9462 reviewers: Vec::new(),
9463 fixer: None,
9464 conductor: conductor(),
9465 implementer_roster: Vec::new(),
9466 },
9467 sem: Arc::new(Semaphore::new(1)),
9468 pause: Pause::new(),
9469 interrupt: Pause::new(),
9470 };
9471
9472 runner.execute().await.expect("execute");
9477
9478 assert_eq!(
9479 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
9480 Some("https://example.invalid/x/y/pull/1"),
9481 "reentry must not push again or open a second pull request over the \
9482 one `land` is already watching"
9483 );
9484 assert_ne!(
9485 runner.state.status,
9486 RunStatus::Landing,
9487 "land could not actually reach the fake pull request, so it must \
9488 have given up rather than left the run silently parked forever"
9489 );
9490 assert_eq!(runner.state.status, RunStatus::Blocked);
9494 assert!(
9495 store.get(&q.id).unwrap().status.open(),
9496 "Blocked is still alive; settle_questions must have been a no-op here"
9497 );
9498 }
9499
9500 fn state_with_round(round: ReviewRound) -> RunState {
9501 let mut s = RunState::new(
9502 PathBuf::from("/repo"),
9503 "main".to_owned(),
9504 "abc1234".to_owned(),
9505 "add retries".to_owned(),
9506 Config::default(),
9507 );
9508 s.reviews = vec![round];
9509 s
9510 }
9511
9512 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
9513 crate::verdict::Finding {
9514 id: id.to_owned(),
9515 severity,
9516 file: None,
9517 line: None,
9518 title: title.to_owned(),
9519 detail: String::new(),
9520 }
9521 }
9522
9523 #[test]
9524 fn pr_body_names_open_findings_and_declined_ones() {
9525 let round = ReviewRound {
9526 round: 2,
9527 head: "deadbee".to_owned(),
9528 verified_head: None,
9529 verified_at: None,
9530 reviews: vec![ReviewRecord {
9531 attempts: 0,
9532 reviewer: 1,
9533 agent: "alpha".to_owned(),
9534 summary: String::new(),
9535 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
9536 vote: None,
9537 failed: None,
9538 duration_ms: 0,
9539 }],
9540 e2e: vec![CommandOutcome {
9541 command: "cargo test".to_owned(),
9542 code: Some(0),
9543 output_tail: String::new(),
9544 duration_ms: 0,
9545 resource_blocked: false,
9546 }],
9547 verify_retried: false,
9548 e2e_deferred: false,
9549 e2e_defer_reason: None,
9550 fix: Some(FixRecord {
9551 agent: "alpha".to_owned(),
9552 addressed: Vec::new(),
9553 rejected: vec![crate::verdict::Rejection {
9554 id: "R1-1-1".to_owned(),
9555 why: "not reachable from any caller".to_owned(),
9556 }],
9557 notes: String::new(),
9558 committed: true,
9559 failed: None,
9560 duration_ms: 0,
9561 continuation: None,
9562 }),
9563 blocking: 0,
9564 answered: 1,
9565 expected: 1,
9566 clean: false,
9567 progressed: true,
9568 vote_split: false,
9569 reconsideration: Vec::new(),
9570 verdict: None,
9571 };
9572 let state = state_with_round(round);
9573 let body = pr_message(&state, 'A').body;
9574
9575 assert!(body.contains("add retries"), "the task must still be there");
9576 assert!(body.contains("R2-1-1"), "{body}");
9577 assert!(body.contains("unused import"), "{body}");
9578 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
9579 assert!(
9580 body.contains("not reachable from any caller"),
9581 "the reason it was declined: {body}"
9582 );
9583 }
9584
9585 #[test]
9586 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
9587 let round = ReviewRound {
9588 round: 1,
9589 head: "deadbee".to_owned(),
9590 verified_head: None,
9591 verified_at: None,
9592 reviews: vec![ReviewRecord {
9593 attempts: 0,
9594 reviewer: 1,
9595 agent: "alpha".to_owned(),
9596 summary: String::new(),
9597 findings: Vec::new(),
9598 vote: None,
9599 failed: None,
9600 duration_ms: 0,
9601 }],
9602 e2e: Vec::new(),
9603 verify_retried: false,
9604 e2e_deferred: false,
9605 e2e_defer_reason: None,
9606 fix: None,
9607 blocking: 0,
9608 answered: 1,
9609 expected: 1,
9610 clean: true,
9611 progressed: false,
9612 vote_split: false,
9613 reconsideration: Vec::new(),
9614 verdict: None,
9615 };
9616 let state = state_with_round(round);
9617 let body = pr_message(&state, 'A').body;
9618 assert!(!body.contains("Open review findings"), "{body}");
9619 assert!(!body.contains("Declined"), "{body}");
9620 }
9621
9622 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
9623 let mut state = RunState::new(
9624 PathBuf::from("/repo"),
9625 "main".to_owned(),
9626 "abc1234".to_owned(),
9627 instruction.to_owned(),
9628 Config::default(),
9629 );
9630 state.candidates.push(Candidate {
9631 index: 0,
9632 label: 'A',
9633 agent: "alpha".to_owned(),
9634 branch: "magi/x/A".to_owned(),
9635 worktree: PathBuf::from("/wt"),
9636 summary: summary.to_owned(),
9637 stat: String::new(),
9638 files: 1,
9639 commits: 1,
9640 empty: false,
9641 failed: None,
9642 verified_noop: None,
9643 folded: false,
9644 duration_ms: 0,
9645 });
9646 state
9647 }
9648
9649 #[test]
9650 fn pr_message_describes_the_change_not_the_task() {
9651 let state = state_with_summary(
9652 "今回やってほしいこと: results projector を直す",
9653 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
9654 );
9655 let m = pr_message(&state, 'A');
9656 assert_eq!(m.title, "fix(web): batch the runs list reads");
9657 assert!(
9658 m.body.starts_with("## Summary\n\n- reads run.json once"),
9659 "{}",
9660 m.body
9661 );
9662 assert!(!m.body.contains("TITLE:"), "{}", m.body);
9663 let task_at = m.body.find("今回やってほしいこと").unwrap();
9664 let details_at = m.body.find("<details>").unwrap();
9665 assert!(
9666 details_at < task_at,
9667 "the task lives inside <details>: {}",
9668 m.body
9669 );
9670 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
9671 assert!(m.body.contains("magi:candidate-a"));
9672 }
9673
9674 #[test]
9675 fn pr_message_falls_back_to_the_task_without_a_title_line() {
9676 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
9677 let m = pr_message(&state, 'A');
9678 assert_eq!(m.title, "add retries");
9679 assert!(
9680 m.body.contains("## Summary\n\n- did some things"),
9681 "{}",
9682 m.body
9683 );
9684
9685 let none = RunState::new(
9686 PathBuf::from("/repo"),
9687 "main".to_owned(),
9688 "abc1234".to_owned(),
9689 "add retries".to_owned(),
9690 Config::default(),
9691 );
9692 let m = pr_message(&none, 'A');
9693 assert_eq!(m.title, "add retries");
9694 assert!(!m.body.contains("## Summary"), "{}", m.body);
9695 }
9696
9697 #[test]
9698 fn pr_message_refuses_the_candidate_commit_subject() {
9699 for bad in [
9700 "TITLE: magi: candidate A (uncommitted work)",
9701 "TITLE: chore: stuff (uncommitted work)",
9702 "TITLE: ",
9703 ] {
9704 let state = state_with_summary("add retries", bad);
9705 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
9706 }
9707 }
9708
9709 #[test]
9710 fn pr_message_bounds_a_very_long_task_and_title() {
9711 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
9712 let state = state_with_summary(&long, "- nothing");
9713 let m = pr_message(&state, 'A');
9714 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9715 assert!(!m.title.contains('\n'));
9716
9717 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
9718 let m = pr_message(&state, 'A');
9719 assert!(m.title.starts_with("feat: "));
9720 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9721 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
9722 }
9723
9724 #[test]
9725 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
9726 let mut state = state_with_summary(
9730 "add retries",
9731 "TITLE: fix(web): batch reads\n- reads run.json once",
9732 );
9733 state.config.graph.language = "ja".to_owned();
9734 let m = pr_message(&state, 'A');
9735 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
9736
9737 let task = "今回やってほしいこと: results projector を直す";
9740 let mut state = state_with_summary(task, "- no title line");
9741 state.config.graph.language = "ja".to_owned();
9742 let m = pr_message(&state, 'A');
9743 assert_eq!(
9744 m.title,
9745 format!("chore: land candidate A of run {}", state.id)
9746 );
9747 assert!(
9748 m.body.contains(&format!(
9749 "<summary>Original task</summary>\n\n{task}\n\n</details>"
9750 )),
9751 "{}",
9752 m.body
9753 );
9754 }
9755
9756 #[test]
9757 fn pr_message_scrubs_home_paths_and_addresses() {
9758 let state = state_with_summary(
9759 "fix it in /Users/someone/src/x",
9760 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9761 );
9762 let m = pr_message(&state, 'A');
9763 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9764 assert!(!m.body.contains(leak), "{}", m.body);
9765 }
9766 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9767 }
9768
9769 #[test]
9770 fn pr_message_survives_a_task_that_closes_details() {
9771 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9772 let m = pr_message(&state, 'A');
9773 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9774 }
9775
9776 #[test]
9777 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9778 let cmd = manual_merge_command(
9779 MergeStyle::Squash,
9780 Path::new("/repo"),
9781 "b",
9782 "fix: \"quoted\" $(x) `y`\n\nbody",
9783 );
9784 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9785 }
9786
9787 #[test]
9788 fn manual_merge_command_matches_the_configured_style() {
9789 let repo = Path::new("/repo");
9790 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9791
9792 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9793 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9794
9795 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9796 assert_eq!(
9797 squash,
9798 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9799 \"Merge magi run 0832 (candidate A)\""
9800 );
9801
9802 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9803 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9804 }
9805
9806 #[test]
9807 fn a_nudge_gets_a_quarter_of_the_budget() {
9808 assert_eq!(retry_budget(secs(1200), true), secs(300));
9810 assert_eq!(retry_budget(secs(3600), true), secs(900));
9811 }
9812
9813 #[test]
9814 fn a_resent_prompt_keeps_the_whole_budget() {
9815 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9818 assert_eq!(retry_budget(secs(60), false), secs(60));
9819 }
9820
9821 #[test]
9822 fn the_floor_never_exceeds_the_original_budget() {
9823 assert_eq!(retry_budget(secs(60), true), secs(60));
9827 assert_eq!(retry_budget(secs(480), true), secs(120));
9828 assert_eq!(retry_budget(secs(0), true), secs(0));
9829 }
9830
9831 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9832 agent::CommandEvidence {
9833 id: "item1".to_owned(),
9834 description: "cargo test".to_owned(),
9835 exit_code,
9836 result_summary: String::new(),
9837 source: "codex".to_owned(),
9838 }
9839 }
9840
9841 #[test]
9842 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9843 assert!(!has_unconfirmed_command(&[]));
9847 }
9848
9849 #[test]
9850 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9851 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9855 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9856 assert!(!has_unconfirmed_command(&[
9857 evidence(Some(0)),
9858 evidence(Some(101))
9859 ]));
9860 }
9861
9862 #[test]
9863 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9864 assert!(has_unconfirmed_command(&[
9865 evidence(Some(0)),
9866 evidence(None)
9867 ]));
9868 }
9869
9870 #[test]
9871 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9872 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9873 assert_eq!(
9874 verified_noop_claim(true, &[], text).as_deref(),
9875 Some("already fixed by b32cfc4, on main.")
9876 );
9877 }
9878
9879 #[test]
9880 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9881 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9884 assert!(verified_noop_claim(false, &[], text).is_none());
9885 }
9886
9887 #[test]
9888 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9889 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9890 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9891 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9893 }
9894
9895 #[test]
9896 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9897 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9898 }
9899
9900 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9903 runner.state.candidates = shape
9904 .iter()
9905 .enumerate()
9906 .map(|(i, &(empty, verified))| Candidate {
9907 index: i,
9908 label: (b'A' + i as u8) as char,
9909 agent: "sonnet".to_owned(),
9910 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9911 worktree: PathBuf::from(format!("/wt/{i}")),
9912 summary: String::new(),
9913 stat: String::new(),
9914 files: 0,
9915 commits: 0,
9916 empty,
9917 failed: None,
9918 verified_noop: verified.map(str::to_owned),
9919 duration_ms: 0,
9920 folded: false,
9921 })
9922 .collect();
9923 }
9924
9925 #[test]
9926 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9927 ask_test_home();
9928 let mut runner = runner_at(RunStatus::Implementing);
9929 set_candidates(
9930 &mut runner,
9931 &[
9932 (true, Some("already on main at b32cfc4")),
9933 (true, Some("same fix, see the existing test")),
9934 ],
9935 );
9936
9937 runner
9938 .after_implement()
9939 .expect("a verified no-op is not an error");
9940
9941 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9942 }
9943
9944 #[test]
9945 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9946 ask_test_home();
9947 let mut runner = runner_at(RunStatus::Implementing);
9948 set_candidates(
9952 &mut runner,
9953 &[(true, Some("already on main at b32cfc4")), (true, None)],
9954 );
9955
9956 let err = runner
9957 .after_implement()
9958 .expect_err("an unverified empty candidate must still fail the run");
9959
9960 assert!(
9961 err.to_string().contains("no candidate produced a change"),
9962 "{err}"
9963 );
9964 assert_eq!(runner.state.status, RunStatus::Failed);
9965 }
9966
9967 #[test]
9968 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9969 ask_test_home();
9970 let mut runner = runner_at(RunStatus::Implementing);
9971 set_candidates(&mut runner, &[(true, None), (true, None)]);
9972
9973 let err = runner
9974 .after_implement()
9975 .expect_err("no candidate declared anything; this is an ordinary failure");
9976
9977 assert!(
9978 err.to_string().contains("no candidate produced a change"),
9979 "{err}"
9980 );
9981 assert_eq!(runner.state.status, RunStatus::Failed);
9982 }
9983}