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, QuotaLoss,
50 ReviewRecord, ReviewRevoteRecord, ReviewRound, RunState, RunStatus, Tally, VoteRecord, tail,
51 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(repo: &Path, instruction: String, config: Config) -> Result<Self> {
455 Self::start_naming(repo, instruction, "", config).await
456 }
457
458 pub async fn start_naming(
462 repo: &Path,
463 instruction: String,
464 also_scan: &str,
465 config: Config,
466 ) -> Result<Self> {
467 let repo = git::toplevel(repo).await?;
468 let missing = agent::missing_programs(&config.agents);
469 if !missing.is_empty() {
470 bail!(
471 "these agent programs are not on PATH: {}. Fix the roster in \
472 magi.toml or install them.",
473 missing.join(", ")
474 );
475 }
476 let base_branch = match config.merge.base.clone() {
477 Some(b) => b,
478 None => git::current_branch(&repo)
479 .await?
480 .context("HEAD is detached; set [merge] base in magi.toml")?,
481 };
482 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
483 if !git::is_clean(&repo).await? {
487 tracing::warn!(
488 "{} has uncommitted changes; they are not part of this run, \
489 which branches off {base_branch} ({})",
490 repo.display(),
491 &base_commit[..base_commit.len().min(8)]
492 );
493 }
494 let roles = config.resolve_roles()?;
495 let max_parallel = config.graph.max_parallel.max(1);
496 let seeds = refs::resolve(
499 &repo,
500 &base_commit,
501 &config.merge.remote,
502 &format!("{also_scan}\n{instruction}"),
503 )
504 .await;
505 refs::plan(&repo, &seeds).await?;
506 let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
507 for seed in &seeds {
508 state.event(
509 "seed",
510 refs::describe(std::slice::from_ref(seed)).unwrap_or_default(),
511 );
512 }
513 state.seeds = seeds;
514 state.event("start", format!("run {} created", state.id));
515 state.save()?;
516 Ok(Self {
517 state,
518 roles,
519 sem: Arc::new(Semaphore::new(max_parallel)),
520 pause: Pause::new(),
521 interrupt: Pause::new(),
522 })
523 }
524
525 pub async fn review(repo: &Path, branch: &str, config: Config) -> Result<Self> {
539 Self::review_taking_over(repo, branch, config, None).await
540 }
541
542 pub async fn review_taking_over(
548 repo: &Path,
549 branch: &str,
550 config: Config,
551 takeover: Option<crate::handover::Takeover>,
552 ) -> Result<Self> {
553 let repo = git::toplevel(repo).await?;
554 let missing = agent::missing_programs(&config.agents);
555 if !missing.is_empty() {
556 bail!(
557 "these agent programs are not on PATH: {}. Fix the roster in \
558 magi.toml or install them.",
559 missing.join(", ")
560 );
561 }
562 let base_branch = match config.merge.base.clone() {
563 Some(b) => b,
564 None => git::current_branch(&repo)
565 .await?
566 .context("HEAD is detached; set [merge] base in magi.toml")?,
567 };
568 if base_branch == branch {
569 bail!("`{branch}` is the base branch; there is nothing to review against");
570 }
571 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
572
573 let roles = config.resolve_roles()?;
574 let max_parallel = config.graph.max_parallel.max(1);
575 let mut state = RunState::new(
576 repo.clone(),
577 base_branch,
578 base_commit.clone(),
579 String::new(),
580 config,
581 );
582
583 let released = match &takeover {
588 Some(takeover) => crate::handover::release(&repo, branch, &state.id, takeover).await?,
589 None => None,
590 };
591 if let Some(released) = &released {
592 state.event(
593 "release",
594 format!(
595 "took `{branch}` over from run {}: its worktree was released",
596 crate::run::short_of(&released.old_id)
597 ),
598 );
599 }
600 if let Some(choice) = takeover.as_ref().and_then(|t| t.choice.as_ref())
604 && let Err(e) =
605 crate::reconcile::apply_choice(&repo, &state.config.merge.remote, branch, choice)
606 .await
607 {
608 if let Some(released) = &released {
609 released.restore(&repo, branch).await;
610 }
611 return Err(e.context("applying the owner's answer about the diverged branch"));
612 }
613 let opened =
614 Self::open_review(&repo, branch, state, roles, max_parallel, base_commit).await;
615 if opened.is_err()
616 && let Some(released) = &released
617 {
618 released.restore(&repo, branch).await;
619 }
620 opened
621 }
622
623 async fn open_review(
626 repo: &Path,
627 branch: &str,
628 mut state: RunState,
629 roles: ResolvedRoles,
630 max_parallel: usize,
631 base_commit: String,
632 ) -> Result<Self> {
633 sync_review_branch(repo, branch, &state.config.merge.remote, &base_commit).await?;
634 let log = git::log_oneline(repo, &base_commit, branch)
637 .await
638 .unwrap_or_default();
639 let instruction = format!(
640 "Review the work already on branch `{branch}`. There is no task \
641 statement: what the change claims to do is whatever its commits \
642 say.\n\n{}",
643 if log.trim().is_empty() {
644 "(no commit messages)"
645 } else {
646 log.trim()
647 }
648 );
649 state.instruction = instruction;
650
651 let worktree = state.worktree_root().join("under-review");
654 if let Some(parent) = worktree.parent() {
655 tokio::fs::create_dir_all(parent).await.ok();
656 }
657 let path = worktree.to_string_lossy().to_string();
658 git::git(repo, &["worktree", "add", &path, branch])
659 .await
660 .with_context(|| {
661 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
662 })?;
663
664 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
665 .await
666 .unwrap_or(0);
667 if commits == 0 {
668 git::worktree_remove(repo, &worktree).await.ok();
669 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
670 }
671 let files = git::changed_files(&worktree, &base_commit, "HEAD")
672 .await
673 .map(|f| f.len())
674 .unwrap_or(0);
675 if files == 0
676 && let (Ok(head_tree), Ok(base_tree)) = (
677 git::tree_of(&worktree, "HEAD").await,
678 git::tree_of(&worktree, &base_commit).await,
679 )
680 && head_tree == base_tree
681 {
682 let head = git::rev_parse(&worktree, "HEAD").await.unwrap_or_default();
683 git::worktree_remove(repo, &worktree).await.ok();
684 bail!(
685 "`{branch}` at {} has a tree identical to base {}; this usually means \
686 the branch ref is stale (check `git rev-parse refs/heads/{branch}` \
687 against `{}/{branch}`) rather than an empty change",
688 short(&head),
689 short(&base_commit),
690 state.config.merge.remote
691 );
692 }
693 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
694 .await
695 .unwrap_or_default();
696
697 state.candidates.push(Candidate {
698 index: 0,
699 label: 'A',
700 agent: "(existing branch)".to_owned(),
703 branch: branch.to_owned(),
704 worktree,
705 summary: String::new(),
706 stat,
707 files,
708 commits,
709 empty: false,
710 failed: None,
711 verified_noop: None,
712 duration_ms: 0,
713 folded: false,
714 });
715 state.tally = Some(Tally {
716 first_choice: BTreeMap::from([('A', 0)]),
717 borda: BTreeMap::new(),
718 winner: 'A',
719 rankings: 0,
720 unanimous_initial: false,
721 deliberated: false,
722 changed_votes: 0,
723 unanimous_final: false,
724 tie_break: None,
725 judges: 0,
729 present: 0,
730 quorum: 0,
731 met_quorum: true,
732 uncontested: Some("review-only run: nothing competed".to_owned()),
733 });
734 state.status = RunStatus::Reviewing;
735 state.event(
736 "start",
737 format!(
738 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
739 state.id
740 ),
741 );
742 state.save()?;
743 Ok(Self {
744 state,
745 roles,
746 sem: Arc::new(Semaphore::new(max_parallel)),
747 pause: Pause::new(),
748 interrupt: Pause::new(),
749 })
750 }
751
752 pub fn resume(id: &str) -> Result<Self> {
754 let state = RunState::load(id)?;
755 if let Some(to) = &state.released_to {
756 bail!(
757 "run {} cannot be resumed: its worktree was released to run {}",
758 state.short(),
759 crate::run::short_of(to)
760 );
761 }
762 let roles = state.config.resolve_roles()?;
763 let max_parallel = state.config.graph.max_parallel.max(1);
764 Ok(Self {
765 state,
766 roles,
767 sem: Arc::new(Semaphore::new(max_parallel)),
768 pause: Pause::new(),
769 interrupt: Pause::new(),
770 })
771 }
772
773 pub async fn execute(&mut self) -> Result<()> {
780 let result = self.execute_graph().await;
781 self.mark_driver_exited();
782 let ended = if result.is_err() {
783 Some(crate::notices::run_stopped(&self.state.id, &self.state))
784 } else {
785 crate::notices::run_ended(&self.state)
786 };
787 if let Some(notice) = ended {
788 crate::notices::raise(notice);
789 }
790 result
791 }
792
793 fn mark_driver_exited(&mut self) {
802 self.state.driver_exited = true;
803 let pid = std::process::id();
804 let Ok(mut disk) = RunState::load(&self.state.id) else {
805 return;
806 };
807 if disk.released_to.is_some() || disk.driver_pid != Some(pid) || disk.driver_exited {
808 return;
809 }
810 disk.driver_exited = true;
811 if let Err(e) = disk.save() {
812 tracing::warn!("could not record that run {} stopped: {e:#}", self.state.id);
813 }
814 }
815
816 async fn execute_graph(&mut self) -> Result<()> {
817 self.state.parked = false;
822 self.state.clear_active();
829 if let Ok(disk) = RunState::load(&self.state.id)
847 && let Some(to) = &disk.released_to
848 {
849 bail!(
850 "run {} cannot continue: its worktree was released to run {}",
851 self.state.short(),
852 crate::run::short_of(to)
853 );
854 }
855 let pid = std::process::id();
856 self.state.driver_pid = Some(pid);
857 self.state.driver_started_at = crate::proc::process_started_at(pid);
858 self.state.driver_exited = false;
859 self.state.save()?;
860 if self.state.status == RunStatus::Stalled {
873 if self.recover_stall().await? {
874 self.finish_after_tally().await?;
875 } else {
876 self.state.save()?;
878 }
879 return Ok(());
880 }
881 if self.state.status == RunStatus::Landing {
891 self.run_land().await?;
892 self.settle_questions();
897 return Ok(());
898 }
899 self.prep().await?;
900 if self.park_here()? {
901 return Ok(());
902 }
903 self.advise().await?;
904 if self.park_here()? {
905 return Ok(());
906 }
907 self.implement().await?;
908 if self.park_here()? {
909 return Ok(());
910 }
911 if self.state.status == RunStatus::VerifiedNoop {
915 return Ok(());
916 }
917 self.judge().await?;
918 if self.park_here()? {
919 return Ok(());
920 }
921 self.deliberate().await?;
922 if self.park_here()? {
923 return Ok(());
924 }
925 self.vote().await?;
926 if self.park_here()? {
927 return Ok(());
928 }
929 self.tally()?;
930 if self.state.status == RunStatus::Stalled {
935 self.state.save()?;
939 return Ok(());
940 }
941 self.finish_after_tally().await?;
942 Ok(())
943 }
944
945 fn park_here(&mut self) -> Result<bool> {
952 if !self.pause.parked() && !self.interrupt.parked() {
957 return Ok(false);
958 }
959 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
960 Some(reason) => format!(
961 "parked after `{}` ({reason}) — resume to carry on from here",
962 self.state.status.as_str()
963 ),
964 None => format!(
965 "parked after `{}` — resume to carry on from here",
966 self.state.status.as_str()
967 ),
968 };
969 self.state.event("park", why);
970 self.state.parked = true;
971 self.state.save()?;
972 Ok(true)
973 }
974
975 pub fn on_pause(&mut self, pause: Pause) {
977 self.pause = pause;
978 }
979
980 pub fn watch_interrupt(&mut self, pause: Pause) {
986 self.interrupt = pause;
987 }
988
989 fn settle_questions(&mut self) {
1009 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
1010 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
1011 }
1012 }
1013
1014 async fn finish_after_tally(&mut self) -> Result<()> {
1017 self.fold_losers().await?;
1018 self.sync_to_base().await?;
1023 self.review_loop().await?;
1024 self.sync_to_base().await?;
1025 self.gate().await?;
1026 self.merge().await?;
1027 self.state.save()?;
1028 Ok(())
1029 }
1030
1031 async fn prep(&mut self) -> Result<()> {
1034 if !self.state.candidates.is_empty() {
1035 return Ok(());
1036 }
1037 self.state.status = RunStatus::Prep;
1038 let repo = self.state.repo.clone();
1039 let base = self.state.base_commit.clone();
1040 let plan = refs::plan(&repo, &self.state.seeds).await?;
1041 let start = plan.start.clone().unwrap_or_else(|| base.clone());
1042 let root = self.state.worktree_root();
1043 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
1044
1045 let hooks_dir = self.state.dir().join("hooks");
1048 if self.state.config.blind.commit_msg_hook {
1049 std::fs::create_dir_all(&hooks_dir)
1050 .with_context(|| format!("create {}", hooks_dir.display()))?;
1051 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
1052 let path = hooks_dir.join("commit-msg");
1053 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
1054 make_executable(&path)?;
1055 git::acquire_worktree_config(&repo).await?;
1063 self.state.enabled_worktree_config = true;
1064 }
1065
1066 for (index, (spec, label)) in self
1067 .roles
1068 .implementers
1069 .clone()
1070 .into_iter()
1071 .zip(labels)
1072 .enumerate()
1073 {
1074 let branch = self.state.branch_for(label);
1075 let worktree = root.join(format!("cand-{label}"));
1076 git::worktree_add_branch(&repo, &worktree, &branch, &start).await?;
1077 if self.state.config.blind.commit_msg_hook {
1078 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
1079 }
1080 git::local_exclude(&worktree, "/.magi/").await?;
1081 for pick in &plan.picks {
1082 if let Err(e) = git::cherry_pick(&worktree, pick).await {
1083 self.state.status = RunStatus::Blocked;
1084 self.state
1085 .event("prep", format!("cannot apply referenced commit: {e}"));
1086 self.state.save()?;
1087 return Err(e);
1088 }
1089 }
1090 self.state.candidates.push(Candidate {
1091 index,
1092 label,
1093 agent: spec.id.clone(),
1094 branch,
1095 worktree,
1096 summary: String::new(),
1097 stat: String::new(),
1098 files: 0,
1099 commits: 0,
1100 empty: false,
1101 failed: None,
1102 verified_noop: None,
1103 duration_ms: 0,
1104 folded: false,
1105 });
1106 }
1107
1108 for j in 1..=self.roles.judges.len() {
1109 let wt = root.join(format!("judge-{j}"));
1110 if !wt.exists() {
1111 git::worktree_add_detached(&repo, &wt, &base).await?;
1112 }
1113 }
1114
1115 if self.state.config.graph.advise {
1123 for k in 1..=self.state.config.graph.advisors {
1124 let wt = root.join(format!("advisor-{k}"));
1125 if !wt.exists() {
1126 git::worktree_add_detached(&repo, &wt, &base).await?;
1127 }
1128 }
1129 }
1130
1131 let authors: Vec<&str> = self
1136 .roles
1137 .implementers
1138 .iter()
1139 .map(|a| a.id.as_str())
1140 .collect();
1141 let overlap: Vec<String> = self
1142 .roles
1143 .judges
1144 .iter()
1145 .enumerate()
1146 .filter(|(_, j)| authors.contains(&j.id.as_str()))
1147 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
1148 .collect();
1149 if !overlap.is_empty() {
1150 let note = format!(
1151 "{} also authored a candidate; blind, but the panel is less \
1152 independent than {} distinct agents would be",
1153 overlap.join(", "),
1154 self.roles.judges.len()
1155 );
1156 self.state.event("prep", note);
1157 }
1158
1159 self.state.event(
1160 "prep",
1161 format!(
1162 "{} candidates, {} judges, base {} ({})",
1163 self.state.candidates.len(),
1164 self.roles.judges.len(),
1165 &self.state.base_commit[..7.min(self.state.base_commit.len())],
1166 self.state.base_branch
1167 ),
1168 );
1169 self.state.status = RunStatus::Implementing;
1170 self.state.save()?;
1171 Ok(())
1172 }
1173
1174 async fn advise(&mut self) -> Result<()> {
1207 let implement_untouched = self
1208 .state
1209 .candidates
1210 .iter()
1211 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
1212 if !self.state.config.graph.advise || self.state.advise_attempted {
1213 return Ok(());
1214 }
1215 if !implement_untouched {
1216 self.state.event(
1217 "advise",
1218 "skipping the design-deliberation stage: at least one \
1219 candidate already shows implementation progress, so this \
1220 run is past the point the stage exists to run before"
1221 .to_owned(),
1222 );
1223 self.state.advise_attempted = true;
1224 self.state.save()?;
1225 return Ok(());
1226 }
1227 let run_id = self.state.id.clone();
1228 let prompts = self.state.config.prompts.clone();
1229 let instruction = self.state.instruction.clone();
1230 let language = self.state.config.graph.language.clone();
1231 let root = self.state.worktree_root();
1232 let n = self.state.config.graph.advisors;
1233 let where_recorded = self.state.dir().join("run.json");
1234
1235 let seats = match self.state.config.advisors() {
1236 Ok(seats) if !seats.is_empty() => seats,
1237 Ok(_) => {
1238 self.state.event(
1239 "advise",
1240 format!(
1241 "[graph] advisors is 0; skipping the design-deliberation \
1242 stage and continuing without a synthesis brief (see {})",
1243 where_recorded.display()
1244 ),
1245 );
1246 self.state.advise_attempted = true;
1247 self.state.save()?;
1248 return Ok(());
1249 }
1250 Err(e) => {
1251 self.state.event(
1252 "advise",
1253 format!(
1254 "could not resolve advisor seats ({e:#}); continuing \
1255 without a design-deliberation brief (see {})",
1256 where_recorded.display()
1257 ),
1258 );
1259 self.state.advise_attempted = true;
1260 self.state.save()?;
1261 return Ok(());
1262 }
1263 };
1264
1265 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1266 let artifacts = agent::artifacts_dir(&self.state.dir());
1267 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1268
1269 let mut jobs = Vec::new();
1270 for (i, spec) in seats.iter().cloned().enumerate() {
1271 let seat_key = format!("advisor-{}", i + 1);
1272 let seat = self.seat(&seat_key, &spec.id);
1273 jobs.push(SeatJob {
1274 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1275 spec,
1276 seat,
1277 cwd: worktrees[i % worktrees.len()].clone(),
1278 timeout,
1279 allow_write: false,
1280 sessions: false,
1281 artifacts: artifacts.clone(),
1282 stem: seat_key,
1283 });
1284 }
1285
1286 self.state.event(
1287 "advise",
1288 format!(
1289 "{} advisor seat(s) sketching a design in parallel",
1290 jobs.len()
1291 ),
1292 );
1293 let mut quota_losses = Vec::new();
1294 let cache = self.state.config.cache_dir();
1295 let ctx = WaveCtx {
1296 run: &run_id,
1297 node: "advise",
1298 prompts: &prompts,
1299 cache: cache.as_deref(),
1300 round: None,
1301 };
1302 let results = ask_json_wave::<Proposal>(
1303 jobs,
1304 Arc::clone(&self.sem),
1305 self.state.config.graph.retries,
1306 &ctx,
1307 &mut quota_losses,
1308 &mut self.state,
1309 &|p: &Proposal| p.validate(),
1310 )
1311 .await;
1312 self.state.quota.extend(quota_losses);
1313
1314 let mut records = Vec::with_capacity(results.len());
1315 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1316 let agent_id = seat.agent.clone();
1317 self.state.seats.insert(seat.key.clone(), seat);
1318 match res {
1319 Ok((proposal, out)) => {
1320 self.state
1321 .event("advise", format!("advisor-{} proposed a design", i + 1));
1322 records.push(advise::AdvisorRecord::proposed(
1323 i + 1,
1324 agent_id,
1325 proposal,
1326 out.duration_ms,
1327 ));
1328 }
1329 Err(e) => {
1330 self.state.event(
1331 "advise",
1332 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1333 );
1334 records.push(advise::AdvisorRecord::failed(
1335 i + 1,
1336 agent_id,
1337 e.to_string(),
1338 ));
1339 }
1340 }
1341 }
1342
1343 let mut advice = advise::Advice {
1344 records,
1345 synthesis: None,
1346 };
1347 if advice.proposals().is_empty() {
1348 self.state.event(
1349 "advise",
1350 "no advisor produced a usable proposal; continuing without a \
1351 synthesis brief"
1352 .to_owned(),
1353 );
1354 } else {
1355 match self
1356 .synthesize_brief(
1357 &advice,
1358 &instruction,
1359 &language,
1360 &worktrees[0],
1361 &artifacts,
1362 &run_id,
1363 &prompts,
1364 cache.as_deref(),
1365 )
1366 .await
1367 {
1368 Ok(Some(text)) => {
1369 self.state.event(
1370 "advise",
1371 "synthesized a design brief for the implementer".to_owned(),
1372 );
1373 advice.synthesis = Some(text);
1374 }
1375 Ok(None) => {
1376 self.state.event(
1377 "advise",
1378 "the synthesis seat produced nothing usable; continuing \
1379 without a design brief"
1380 .to_owned(),
1381 );
1382 }
1383 Err(e) => {
1384 self.state.event(
1385 "advise",
1386 format!("could not synthesize a design brief: {e:#}"),
1387 );
1388 }
1389 }
1390 }
1391 advise::apply_reflection(&mut advice);
1392
1393 self.state.advice = Some(advice);
1394 self.state.advise_attempted = true;
1395 self.state.save()?;
1396 Ok(())
1397 }
1398
1399 #[allow(clippy::too_many_arguments)]
1411 async fn synthesize_brief(
1412 &mut self,
1413 advice: &advise::Advice,
1414 instruction: &str,
1415 language: &str,
1416 cwd: &Path,
1417 artifacts: &Path,
1418 run_id: &str,
1419 prompts: &Prompts,
1420 cache: Option<&Path>,
1421 ) -> Result<Option<String>> {
1422 let want = self.state.config.roles.synthesizer.as_deref();
1423 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1424 let mut seat = self.seat("advise-synthesis", &spec.id);
1425 let proposals = advice.proposals();
1426 let mut prompt = prompt::with_overlay(
1427 prompt::synthesize_brief(instruction, &proposals, language),
1428 prompts.overlay("advise"),
1429 );
1430 if cache.is_some() {
1431 prompt.push('\n');
1436 prompt.push_str(&prompt::build_cache_note("advise", false));
1437 }
1438 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1439 let out = agent::invoke(
1440 &spec,
1441 &mut seat,
1442 &Invocation {
1443 cwd,
1444 prompt: &prompt,
1445 timeout,
1446 allow_write: false,
1447 sessions: false,
1448 artifacts,
1449 stem: "advise-synthesis",
1450 run: run_id,
1451 node: "advise",
1452 cache_dir: None,
1453 attachments: &[],
1454 },
1455 )
1456 .await?;
1457 self.state.seats.insert(seat.key.clone(), seat);
1458 if !out.usable() {
1459 return Ok(None);
1460 }
1461 let text =
1462 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1463 Ok((!text.trim().is_empty()).then_some(text))
1464 }
1465
1466 async fn implement(&mut self) -> Result<()> {
1469 let run_id = self.state.id.clone();
1474 let prompts = self.state.config.prompts.clone();
1475 let todo: Vec<usize> = self
1476 .state
1477 .candidates
1478 .iter()
1479 .enumerate()
1480 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1481 .map(|(i, _)| i)
1482 .collect();
1483 if todo.is_empty() {
1484 return self.after_implement();
1485 }
1486 self.state.status = RunStatus::Implementing;
1487
1488 let language = self.state.config.graph.language.clone();
1489 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1490 let sessions = self.state.config.graph.sessions;
1491 let artifacts = agent::artifacts_dir(&self.state.dir());
1492 let brief = self
1496 .state
1497 .advice
1498 .as_ref()
1499 .and_then(|a| a.synthesis.as_deref())
1500 .map(str::to_owned);
1501
1502 let mut jobs = Vec::new();
1503 for &i in &todo {
1504 let (index, label, worktree) = {
1505 let c = &self.state.candidates[i];
1506 (c.index, c.label, c.worktree.clone())
1507 };
1508 let spec = self.roles.implementers[index].clone();
1509 let seat_key = format!("impl-{label}");
1510 let seat = self.seat(&seat_key, &spec.id);
1511 let instruction = seeded_instruction(&self.state);
1512 jobs.push(SeatJob {
1513 spec,
1514 seat,
1515 prompt: prompt::implement(
1516 &instruction,
1517 &worktree.to_string_lossy(),
1518 &language,
1519 brief.as_deref(),
1520 ),
1521 cwd: worktree,
1522 timeout,
1523 allow_write: true,
1524 sessions,
1525 artifacts: artifacts.clone(),
1526 stem: format!("impl-{label}"),
1527 });
1528 }
1529
1530 self.state.event(
1531 "implement",
1532 format!("{} candidates in parallel", jobs.len()),
1533 );
1534 let mut sent = jobs.clone();
1540 let cache = self.state.config.cache_dir();
1541 let ctx = WaveCtx {
1542 run: &run_id,
1543 node: "implement",
1544 prompts: &prompts,
1545 cache: cache.as_deref(),
1546 round: None,
1547 };
1548 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1549 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1550 .await;
1551 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1552 .await;
1553 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1554 .await;
1555
1556 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1557 let seat_key = seat.key.clone();
1558 let agent = seat.agent.clone();
1567 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1568 self.state.seats.insert(seat.key.clone(), seat);
1569 let label = self.state.candidates[i].label;
1570 let worktree = self.state.candidates[i].worktree.clone();
1571 let base = self.state.base_commit.clone();
1572
1573 let (summary, duration, failed, verified_claim) = match out {
1574 AgentOutcome::Ok(o) => {
1575 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1576 let failed = (!o.usable()).then(|| {
1577 if o.timed_out {
1578 "agent timed out".to_owned()
1579 } else {
1580 format!("agent exited with {:?}", o.exit_code)
1581 }
1582 });
1583 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1584 (text, o.duration_ms, failed, verified_claim)
1585 }
1586 AgentOutcome::Dropped(o) => {
1592 let why = o
1593 .dropped
1594 .as_ref()
1595 .map(|d| d.why.as_str())
1596 .unwrap_or("the CLI ended the stream without delivering its answer");
1597 (
1598 String::new(),
1599 o.duration_ms,
1600 Some(format!("the CLI dropped the stream ({why})")),
1601 None,
1602 )
1603 }
1604 AgentOutcome::Quota(o) => {
1605 self.state.quota.push(QuotaLoss {
1606 seat: seat_key,
1607 node: "implement".to_owned(),
1608 at: Timestamp::now(),
1609 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1610 });
1611 (
1612 String::new(),
1613 o.duration_ms,
1614 Some("rate limited (quota); produced no change".to_owned()),
1615 None,
1616 )
1617 }
1618 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1619 };
1620
1621 let rescued = match git::rescue_commit(
1624 &worktree,
1625 &format!("magi: candidate {label} (uncommitted work)"),
1626 )
1627 .await
1628 {
1629 Ok(r) => {
1630 self.state.note_withheld("implement", &r.withheld);
1631 r.committed
1632 }
1633 Err(_) => false,
1634 };
1635 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1636 .await
1637 .unwrap_or(0);
1638 let patch = git::diff(&worktree, &base, "HEAD")
1639 .await
1640 .unwrap_or_default();
1641 let stat = git::diff_stat(&worktree, &base, "HEAD")
1642 .await
1643 .unwrap_or_default();
1644 let files = git::changed_files(&worktree, &base, "HEAD")
1645 .await
1646 .map(|f| f.len())
1647 .unwrap_or(0);
1648 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1649
1650 let c = &mut self.state.candidates[i];
1651 if !exhausted_the_fallback_chain {
1652 c.agent = agent;
1653 }
1654 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1655 c.stat = stat;
1656 c.files = files;
1657 c.commits = commits;
1658 c.duration_ms = duration;
1659 c.empty = commits == 0 || patch.trim().is_empty();
1660 c.failed = match failed {
1663 Some(_) if c.empty => failed,
1664 _ => None,
1665 };
1666 c.verified_noop = if c.empty { verified_claim } else { None };
1671 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1672 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1673 (None, true, Some(_), _) => {
1674 format!("candidate {label}: no change produced (agent-verified no-op)")
1675 }
1676 (None, true, None, _) => format!("candidate {label}: no change produced"),
1677 (None, false, _, true) => {
1678 format!(
1679 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1680 )
1681 }
1682 (None, false, _, false) => {
1683 format!("candidate {label}: {files} files, {commits} commits")
1684 }
1685 };
1686 self.state.event("implement", note);
1687 self.state.save()?;
1688 }
1689
1690 self.after_implement()
1691 }
1692
1693 async fn resume_undelivered(
1721 &mut self,
1722 results: &mut [(usize, SeatState, AgentOutcome)],
1723 sent: &[SeatJob],
1724 prompts: &Prompts,
1725 run_id: &str,
1726 ) {
1727 for (wi, seat, out) in results.iter_mut() {
1728 let Some(dropped) = (match &*out {
1729 AgentOutcome::Dropped(o) => o.dropped.clone(),
1730 _ => None,
1731 }) else {
1732 continue;
1733 };
1734 let Some(job) = sent.get(*wi) else { continue };
1735 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1737 self.state.event(
1738 "implement",
1739 format!(
1740 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1741 work is in the tree",
1742 seat.key, dropped.output_tokens, dropped.why
1743 ),
1744 );
1745 continue;
1746 }
1747 if !has_context(&job.spec, seat, job.sessions) {
1755 self.state.event(
1756 "implement",
1757 format!(
1758 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1759 is no session left to resume",
1760 seat.key, dropped.output_tokens, dropped.why
1761 ),
1762 );
1763 continue;
1764 }
1765 self.state.event(
1766 "implement",
1767 format!(
1768 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1769 conversation",
1770 seat.key, dropped.output_tokens, dropped.why
1771 ),
1772 );
1773 let mut retry = job.clone();
1774 retry.seat = seat.clone();
1775 retry.prompt = prompt::resume_after_drop(&dropped.why);
1776 retry.timeout = retry_budget(job.timeout, true);
1777 retry.stem = format!("{}-resume", job.stem);
1778 let cache = self.state.config.cache_dir();
1779 let ctx = WaveCtx {
1780 run: run_id,
1781 node: "implement",
1782 prompts,
1783 cache: cache.as_deref(),
1784 round: None,
1785 };
1786 let (resumed_seat, resumed) =
1787 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1788 *seat = resumed_seat;
1789 *out = resumed;
1790 }
1791 }
1792
1793 async fn resume_quota_losses(
1855 &mut self,
1856 results: &mut [(usize, SeatState, AgentOutcome)],
1857 sent: &mut [SeatJob],
1858 prompts: &Prompts,
1859 run_id: &str,
1860 ) {
1861 let instruction = seeded_instruction(&self.state);
1862 let language = self.state.config.graph.language.clone();
1863 let brief = self
1864 .state
1865 .advice
1866 .as_ref()
1867 .and_then(|a| a.synthesis.as_deref())
1868 .map(str::to_owned);
1869 for (wi, seat, out) in results.iter_mut() {
1870 let Some(job) = sent.get_mut(*wi) else {
1871 continue;
1872 };
1873 let start = self
1878 .roles
1879 .implementer_roster
1880 .iter()
1881 .position(|s| s.id == job.spec.id)
1882 .unwrap_or(0);
1883 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1884 let mut fallback_attempt = 0usize;
1885 while matches!(&*out, AgentOutcome::Quota(_)) {
1886 let Some(next) =
1887 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1888 .cloned()
1889 else {
1890 break;
1891 };
1892 tried.insert(next.id.clone());
1893 fallback_attempt += 1;
1894
1895 if let Ok(r) = git::rescue_commit(
1896 &job.cwd,
1897 &format!(
1898 "magi: candidate {} (uncommitted work before quota fallback)",
1899 seat.key
1900 ),
1901 )
1902 .await
1903 {
1904 self.state.note_withheld("implement", &r.withheld);
1905 }
1906
1907 self.state.event(
1908 "implement",
1909 format!(
1910 "{}: rate limited (quota) on {}; retrying with {}",
1911 seat.key, seat.agent, next.id
1912 ),
1913 );
1914
1915 let new_seat = self.seat(&seat.key, &next.id);
1916 job.spec = next.clone();
1924 let mut retry = job.clone();
1925 retry.seat = new_seat;
1926 retry.prompt = prompt::implement(
1927 &instruction,
1928 &job.cwd.to_string_lossy(),
1929 &language,
1930 brief.as_deref(),
1931 );
1932 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1933 let cache = self.state.config.cache_dir();
1934 let ctx = WaveCtx {
1935 run: run_id,
1936 node: "implement",
1937 prompts,
1938 cache: cache.as_deref(),
1939 round: None,
1940 };
1941 let (fallback_seat, fallback_out) = run_one(
1942 retry,
1943 Arc::clone(&self.sem),
1944 &ctx,
1945 &mut self.state,
1946 fallback_attempt,
1947 )
1948 .await;
1949 *seat = fallback_seat;
1950 *out = fallback_out;
1951 }
1952 }
1953 }
1954
1955 async fn resume_unconfirmed_commands(
1979 &mut self,
1980 results: &mut [(usize, SeatState, AgentOutcome)],
1981 sent: &[SeatJob],
1982 prompts: &Prompts,
1983 run_id: &str,
1984 ) {
1985 for (wi, seat, out) in results.iter_mut() {
1986 let AgentOutcome::Ok(o) = &*out else {
1987 continue;
1988 };
1989 if !has_unconfirmed_command(&o.commands) {
1990 continue;
1991 }
1992 let Some(job) = sent.get(*wi) else { continue };
1993 if !has_context(&job.spec, seat, job.sessions) {
1994 self.state.event(
1995 "implement",
1996 format!(
1997 "{}: the reply named a command whose own CLI never confirmed the exit \
1998 status of, but there is no session left to resume",
1999 seat.key
2000 ),
2001 );
2002 continue;
2003 }
2004 self.state.event(
2005 "implement",
2006 format!(
2007 "{}: the reply named a command whose own CLI never confirmed the exit \
2008 status of; resuming the conversation",
2009 seat.key
2010 ),
2011 );
2012 let mut retry = job.clone();
2013 retry.seat = seat.clone();
2014 retry.prompt = prompt::resume_incomplete(
2015 "a command in your last reply had no confirmed exit status",
2016 );
2017 retry.timeout = retry_budget(job.timeout, true);
2018 retry.stem = format!("{}-confirm", job.stem);
2019 let cache = self.state.config.cache_dir();
2020 let ctx = WaveCtx {
2021 run: run_id,
2022 node: "implement",
2023 prompts,
2024 cache: cache.as_deref(),
2025 round: None,
2026 };
2027 let (resumed_seat, resumed) =
2028 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
2029 *seat = resumed_seat;
2030 *out = resumed;
2031 }
2032 }
2033
2034 async fn continue_fix_report(
2055 &mut self,
2056 mut seat: SeatState,
2057 parse_err: String,
2058 job: &SeatJob,
2059 prompts: &Prompts,
2060 run_id: &str,
2061 round: usize,
2062 ) -> (
2063 SeatState,
2064 Option<FixReport>,
2065 Option<String>,
2066 ContinuationRecord,
2067 ) {
2068 let mut last_err = parse_err;
2069 let mut cumulative_wait_ms = 0u64;
2070 let mut attempts = 0usize;
2071 loop {
2072 if !has_context(&job.spec, &seat, job.sessions) {
2073 self.state.event(
2074 "fix",
2075 format!(
2076 "round {round}: fixer's reply had no adoption report ({last_err}); no \
2077 session left to resume into"
2078 ),
2079 );
2080 let outcome = if attempts == 0 {
2081 ContinuationOutcome::NoSession
2082 } else {
2083 ContinuationOutcome::Exhausted
2084 };
2085 return (
2086 seat,
2087 None,
2088 Some(format!("unparsable fix report: {last_err}")),
2089 ContinuationRecord {
2090 attempts,
2091 cumulative_wait_ms,
2092 outcome,
2093 },
2094 );
2095 }
2096 if attempts >= MAX_FIX_CONTINUATIONS {
2097 self.state.event(
2098 "fix",
2099 format!(
2100 "round {round}: fixer's reply still had no adoption report after \
2101 {attempts} continuation(s) ({last_err}); giving up"
2102 ),
2103 );
2104 return (
2105 seat,
2106 None,
2107 Some(format!(
2108 "unparsable fix report after {attempts} continuation(s): {last_err}"
2109 )),
2110 ContinuationRecord {
2111 attempts,
2112 cumulative_wait_ms,
2113 outcome: ContinuationOutcome::Exhausted,
2114 },
2115 );
2116 }
2117 attempts += 1;
2118 self.state.event(
2119 "fix",
2120 format!(
2121 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
2122 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
2123 ),
2124 );
2125 let mut retry = job.clone();
2126 retry.seat = seat.clone();
2127 retry.prompt = prompt::resume_incomplete(&last_err);
2128 retry.timeout = retry_budget(job.timeout, true);
2129 retry.stem = format!("{}-continue{attempts}", job.stem);
2130 let cache = self.state.config.cache_dir();
2131 let ctx = WaveCtx {
2132 run: run_id,
2133 node: "fix",
2134 prompts,
2135 cache: cache.as_deref(),
2136 round: Some(round),
2137 };
2138 let (resumed_seat, resumed_out) = run_one(
2139 retry,
2140 Arc::clone(&self.sem),
2141 &ctx,
2142 &mut self.state,
2143 attempts,
2144 )
2145 .await;
2146 seat = resumed_seat;
2147 match resumed_out {
2148 AgentOutcome::Ok(o) => {
2149 cumulative_wait_ms += o.duration_ms;
2150 match verdict::extract_json::<FixReport>(&o.text) {
2151 Ok(report) if !has_unconfirmed_command(&o.commands) => {
2152 self.state.event(
2153 "fix",
2154 format!(
2155 "round {round}: fixer's adoption report recovered after \
2156 {attempts} continuation(s)"
2157 ),
2158 );
2159 return (
2160 seat,
2161 Some(report),
2162 None,
2163 ContinuationRecord {
2164 attempts,
2165 cumulative_wait_ms,
2166 outcome: ContinuationOutcome::Resumed,
2167 },
2168 );
2169 }
2170 Ok(_) => {
2178 last_err = "the reply parsed, but it reported a command whose own CLI \
2179 never confirmed an exit status"
2180 .to_owned();
2181 }
2182 Err(e) => last_err = e.to_string(),
2183 }
2184 }
2185 AgentOutcome::Quota(o) => {
2186 cumulative_wait_ms += o.duration_ms;
2187 self.state.quota.push(QuotaLoss {
2188 seat: seat.key.clone(),
2189 node: "fix".to_owned(),
2190 at: Timestamp::now(),
2191 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2192 });
2193 self.state.event(
2194 "fix",
2195 format!(
2196 "round {round}: continuation rate limited (quota); not retrying now"
2197 ),
2198 );
2199 return (
2200 seat,
2201 None,
2202 Some("rate limited (quota) while recovering the fix report".to_owned()),
2203 ContinuationRecord {
2204 attempts,
2205 cumulative_wait_ms,
2206 outcome: ContinuationOutcome::QuotaLost,
2207 },
2208 );
2209 }
2210 AgentOutcome::Dropped(o) => {
2211 cumulative_wait_ms += o.duration_ms;
2212 let why = o
2213 .dropped
2214 .as_ref()
2215 .map(|d| d.why.as_str())
2216 .unwrap_or("the CLI ended the stream without delivering its answer");
2217 last_err = format!("the CLI dropped the stream ({why})");
2218 }
2219 AgentOutcome::Failed(e) => last_err = e,
2220 }
2221 }
2222 }
2223
2224 fn after_implement(&mut self) -> Result<()> {
2225 if self.state.leaks.is_empty() {
2227 let cfg = self.state.config.blind.clone();
2228 let mut leaks = Vec::new();
2229 for c in &self.state.candidates {
2230 let Some(patch) =
2231 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
2232 else {
2233 continue;
2234 };
2235 leaks.extend(blind::scan(
2236 &format!("candidate {} patch", c.label),
2237 &patch,
2238 &cfg.vendor_tokens,
2239 ));
2240 }
2241 if !leaks.is_empty() {
2242 let summary = leaks
2243 .iter()
2244 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2245 .collect::<Vec<_>>()
2246 .join(", ");
2247 match cfg.on_leak {
2248 LeakPolicy::Fail => {
2249 self.state.status = RunStatus::Failed;
2250 self.state
2251 .event("blind", format!("vendor text in a patch: {summary}"));
2252 self.state.leaks = leaks;
2253 self.state.save()?;
2254 self.settle_questions();
2255 bail!(
2256 "blind.on_leak = \"fail\" and vendor text reached a \
2257 judged patch: {summary}"
2258 );
2259 }
2260 LeakPolicy::Redact => self.state.event(
2261 "blind",
2262 format!("redacting vendor text for judging: {summary}"),
2263 ),
2264 LeakPolicy::Warn => self.state.event(
2265 "blind",
2266 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2267 ),
2268 }
2269 self.state.leaks = leaks;
2270 }
2271 }
2272
2273 if self.state.viable().is_empty() {
2274 if self.state.all_candidates_verified_noop() {
2275 self.state.status = RunStatus::VerifiedNoop;
2286 self.state.save()?;
2287 self.settle_questions();
2288 return Ok(());
2289 }
2290 self.state.status = RunStatus::Failed;
2291 self.state.save()?;
2292 self.settle_questions();
2293 bail!("no candidate produced a change; nothing to judge");
2294 }
2295 self.state.status = RunStatus::Judging;
2296 self.state.save()?;
2297 Ok(())
2298 }
2299
2300 async fn judge(&mut self) -> Result<()> {
2303 let run_id = self.state.id.clone();
2308 let prompts = self.state.config.prompts.clone();
2309 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2310 return Ok(());
2311 }
2312 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2313 if viable.len() == 1 {
2314 self.state.judge_skipped = true;
2321 self.state.event(
2322 "judge",
2323 format!(
2324 "only candidate {} produced a change; judging skipped",
2325 viable[0].label
2326 ),
2327 );
2328 self.state.save()?;
2329 return Ok(());
2330 }
2331 self.state.status = RunStatus::Judging;
2332
2333 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2334 let language = self.state.config.graph.language.clone();
2335 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2336 let sessions = self.state.config.graph.sessions;
2337 let artifacts = agent::artifacts_dir(&self.state.dir());
2338 let root = self.state.worktree_root();
2339 let base_short = short(&self.state.base_commit);
2340
2341 let mut jobs = Vec::new();
2342 let mut orders = Vec::new();
2343 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2344 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2345 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2346 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2347 let seat_key = format!("judge-{}", j + 1);
2348 let seat = self.seat(&seat_key, &spec.id);
2349 jobs.push(SeatJob {
2350 prompt: prompt::judge(
2351 &self.state.instruction,
2352 &views,
2353 self.roles.judges.len(),
2354 &base_short,
2355 &language,
2356 ),
2357 spec,
2358 seat,
2359 cwd: root.join(format!("judge-{}", j + 1)),
2360 timeout,
2361 allow_write: false,
2362 sessions,
2363 artifacts: artifacts.clone(),
2364 stem: format!("judge-{}", j + 1),
2365 });
2366 }
2367
2368 self.state.event(
2369 "judge",
2370 format!(
2371 "{} judges ranking {} candidates blind",
2372 jobs.len(),
2373 viable.len()
2374 ),
2375 );
2376 let labels_for_check = labels.clone();
2377 let mut quota_losses = Vec::new();
2378 let cache = self.state.config.cache_dir();
2379 let ctx = WaveCtx {
2380 run: &run_id,
2381 node: "judge",
2382 prompts: &prompts,
2383 cache: cache.as_deref(),
2384 round: None,
2385 };
2386 let results = ask_json_wave::<Ranking>(
2387 jobs,
2388 Arc::clone(&self.sem),
2389 self.state.config.graph.retries,
2390 &ctx,
2391 &mut quota_losses,
2392 &mut self.state,
2393 &move |r: &Ranking| r.validate(&labels_for_check),
2394 )
2395 .await;
2396 self.state.quota.extend(quota_losses);
2397
2398 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2399 let agent_id = seat.agent.clone();
2400 self.state.seats.insert(seat.key.clone(), seat);
2401 let mut record = Judgement {
2402 judge: j + 1,
2403 seat: format!("judge-{}", j + 1),
2404 agent: agent_id,
2405 ranking: Vec::new(),
2406 reasons: BTreeMap::new(),
2407 confidence: None,
2408 order: orders[j].clone(),
2409 failed: None,
2410 duration_ms: 0,
2411 };
2412 match res {
2413 Ok((ranking, out)) => {
2414 record.ranking = ranking.normalized();
2415 record.reasons = ranking.reasons;
2416 record.confidence = ranking.confidence;
2417 record.duration_ms = out.duration_ms;
2418 self.state.event(
2419 "judge",
2420 format!(
2421 "judge {} ranked {}",
2422 j + 1,
2423 record.ranking.iter().collect::<String>()
2424 ),
2425 );
2426 }
2427 Err(e) => {
2428 record.failed = Some(e.to_string());
2429 self.state
2430 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2431 }
2432 }
2433 self.state.judgements.push(record);
2434 self.state.save()?;
2435 }
2436 Ok(())
2437 }
2438
2439 async fn deliberate(&mut self) -> Result<()> {
2442 let run_id = self.state.id.clone();
2447 let prompts = self.state.config.prompts.clone();
2448 if !self.state.deliberation.is_empty() {
2449 return Ok(());
2450 }
2451 let tops: Vec<char> = self
2452 .state
2453 .judgements
2454 .iter()
2455 .filter_map(|j| j.ranking.first().copied())
2456 .collect();
2457 let rounds = self.state.config.graph.deliberate_rounds;
2458 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2459 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2460 self.state.event(
2461 "deliberate",
2462 format!("judges agreed on {} outright; no deliberation", tops[0]),
2463 );
2464 }
2465 self.state.status = RunStatus::Voting;
2466 self.state.save()?;
2467 return Ok(());
2468 }
2469
2470 self.state.status = RunStatus::Deliberating;
2471 self.state.event(
2472 "deliberate",
2473 format!(
2474 "split: first choices were {} — opening {rounds} round(s)",
2475 tops.iter().collect::<String>()
2476 ),
2477 );
2478
2479 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2480 let language = self.state.config.graph.language.clone();
2481 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2482 let sessions = self.state.config.graph.sessions;
2483 let artifacts = agent::artifacts_dir(&self.state.dir());
2484 let root = self.state.worktree_root();
2485 let base_short = short(&self.state.base_commit);
2486
2487 for round in 1..=rounds {
2491 let mut turns: Vec<DeliberationTurn> = Vec::new();
2492 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2493 if self.state.judgements[j].failed.is_some() {
2494 continue;
2495 }
2496 let seat_key = format!("judge-{}", j + 1);
2497 let mut seat = self.seat(&seat_key, &spec.id);
2498 let transcript = self.transcript(&turns, j);
2499 let context = if has_context(&spec, &seat, sessions) {
2500 None
2501 } else {
2502 Some(self.candidate_block(&viable, &base_short))
2503 };
2504 let text = prompt::deliberate(
2505 &self.state.instruction,
2506 context.as_deref(),
2507 &transcript,
2508 round,
2509 rounds,
2510 &language,
2511 );
2512 let job = SeatJob {
2513 spec,
2514 seat: seat.clone(),
2515 prompt: text,
2516 cwd: root.join(format!("judge-{}", j + 1)),
2517 timeout,
2518 allow_write: false,
2519 sessions,
2520 artifacts: artifacts.clone(),
2521 stem: format!("delib-{round}-judge-{}", j + 1),
2522 };
2523 let cache = self.state.config.cache_dir();
2524 let ctx = WaveCtx {
2525 run: &run_id,
2526 node: "deliberate",
2527 prompts: &prompts,
2528 cache: cache.as_deref(),
2529 round: None,
2530 };
2531 let (updated, out) =
2532 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2533 seat = updated;
2534 let agent_id = seat.agent.clone();
2535 let seat_key = seat.key.clone();
2536 self.state.seats.insert(seat.key.clone(), seat);
2537 let body = match out {
2538 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2539 AgentOutcome::Dropped(o) => {
2543 let why =
2544 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2545 "the CLI ended the stream without delivering its answer",
2546 );
2547 self.state.event(
2548 "deliberate",
2549 format!(
2550 "judge {} skipped: the CLI dropped the stream ({why})",
2551 j + 1
2552 ),
2553 );
2554 continue;
2555 }
2556 AgentOutcome::Quota(o) => {
2557 self.state.quota.push(QuotaLoss {
2558 seat: seat_key,
2559 node: "deliberate".to_owned(),
2560 at: Timestamp::now(),
2561 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2562 });
2563 self.state.event(
2564 "deliberate",
2565 format!("judge {} skipped: rate limited (quota)", j + 1),
2566 );
2567 continue;
2568 }
2569 AgentOutcome::Failed(e) => {
2570 self.state
2571 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2572 continue;
2573 }
2574 };
2575 let tentative = verdict::extract_json::<Position>(&body)
2576 .ok()
2577 .and_then(|p| p.tentative)
2578 .and_then(|s| s.trim().chars().next())
2579 .map(|c| c.to_ascii_uppercase());
2580 self.state.event(
2581 "deliberate",
2582 format!(
2583 "round {round}: judge {} now favours {}",
2584 j + 1,
2585 tentative.map_or("—".to_owned(), |c| c.to_string())
2586 ),
2587 );
2588 turns.push(DeliberationTurn {
2589 judge: j + 1,
2590 agent: agent_id,
2591 body: blind::sanitize_prose(&body, &self.state.config.blind),
2592 tentative,
2593 });
2594 }
2595 self.state
2596 .deliberation
2597 .push(DeliberationRound { round, turns });
2598 self.state.save()?;
2599 }
2600
2601 self.state.status = RunStatus::Voting;
2602 self.state.save()?;
2603 Ok(())
2604 }
2605
2606 async fn vote(&mut self) -> Result<()> {
2609 let run_id = self.state.id.clone();
2614 let prompts = self.state.config.prompts.clone();
2615 if !self.state.votes.is_empty() {
2616 return Ok(());
2617 }
2618 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2619 if viable.len() == 1 {
2620 return Ok(());
2621 }
2622 self.state.status = RunStatus::Voting;
2623
2624 let language = self.state.config.graph.language.clone();
2625 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2626 let sessions = self.state.config.graph.sessions;
2627 let artifacts = agent::artifacts_dir(&self.state.dir());
2628 let root = self.state.worktree_root();
2629 let base_short = short(&self.state.base_commit);
2630 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2631
2632 let mut jobs = Vec::new();
2633 let mut seats_at = Vec::new();
2634 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2635 if self
2636 .state
2637 .judgements
2638 .get(j)
2639 .is_some_and(|r| r.failed.is_some())
2640 {
2641 continue;
2642 }
2643 let seat_key = format!("judge-{}", j + 1);
2644 let seat = self.seat(&seat_key, &spec.id);
2645 let mut text = prompt::final_vote(&viable, &language);
2646 if !has_context(&spec, &seat, sessions) {
2647 text = format!(
2648 "{}\n\n# Candidates\n\n{}",
2649 text,
2650 self.candidate_block(&candidates, &base_short)
2651 );
2652 }
2653 jobs.push(SeatJob {
2654 spec,
2655 seat,
2656 prompt: text,
2657 cwd: root.join(format!("judge-{}", j + 1)),
2658 timeout,
2659 allow_write: false,
2660 sessions,
2661 artifacts: artifacts.clone(),
2662 stem: format!("vote-judge-{}", j + 1),
2663 });
2664 seats_at.push(j);
2665 }
2666
2667 self.state.event(
2668 "vote",
2669 format!(
2670 "collecting {} final votes one by one, privately",
2671 jobs.len()
2672 ),
2673 );
2674 let allowed = viable.clone();
2675 let mut quota_losses = Vec::new();
2676 let cache = self.state.config.cache_dir();
2677 let ctx = WaveCtx {
2678 run: &run_id,
2679 node: "vote",
2680 prompts: &prompts,
2681 cache: cache.as_deref(),
2682 round: None,
2683 };
2684 let results = ask_json_wave::<FinalVote>(
2685 jobs,
2686 Arc::clone(&self.sem),
2687 self.state.config.graph.retries,
2688 &ctx,
2689 &mut quota_losses,
2690 &mut self.state,
2691 &move |v: &FinalVote| match v.label() {
2692 Some(c) if allowed.contains(&c) => Ok(()),
2693 other => bail!("vote {other:?} is not one of {allowed:?}"),
2694 },
2695 )
2696 .await;
2697 self.state.quota.extend(quota_losses);
2698
2699 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2700 let agent_id = seat.agent.clone();
2701 self.state.seats.insert(seat.key.clone(), seat);
2702 let initial = self
2703 .state
2704 .judgements
2705 .get(j)
2706 .and_then(|r| r.ranking.first().copied());
2707 let mut record = VoteRecord {
2708 judge: j + 1,
2709 agent: agent_id,
2710 vote: None,
2711 reason: String::new(),
2712 changed: false,
2713 };
2714 match res {
2715 Ok((v, _)) => {
2716 record.vote = v.label();
2717 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2718 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2719 self.state.event(
2720 "vote",
2721 format!(
2722 "judge {} voted {}{}",
2723 j + 1,
2724 record.vote.unwrap_or('?'),
2725 if record.changed { " (changed)" } else { "" }
2726 ),
2727 );
2728 }
2729 Err(e) => {
2730 self.state
2731 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2732 }
2733 }
2734 self.state.votes.push(record);
2735 self.state.save()?;
2736 }
2737 Ok(())
2738 }
2739
2740 fn tally(&mut self) -> Result<()> {
2743 if self.state.tally.is_some() {
2744 return Ok(());
2745 }
2746 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2747 let tops: Vec<char> = self
2748 .state
2749 .judgements
2750 .iter()
2751 .filter_map(|j| j.ranking.first().copied())
2752 .collect();
2753 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2754
2755 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2758 let mut cast: Vec<char> = Vec::new();
2759 for (i, j) in self.state.judgements.iter().enumerate() {
2760 let vote = self
2761 .state
2762 .votes
2763 .iter()
2764 .find(|v| v.judge == i + 1)
2765 .and_then(|v| v.vote)
2766 .or_else(|| j.ranking.first().copied());
2767 if let Some(v) = vote {
2768 *first_choice.entry(v).or_insert(0) += 1;
2769 cast.push(v);
2770 }
2771 }
2772
2773 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2774 for j in &self.state.judgements {
2775 let n = j.ranking.len();
2776 for (pos, label) in j.ranking.iter().enumerate() {
2777 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2778 }
2779 }
2780
2781 let best = first_choice.values().copied().max().unwrap_or(0);
2782 let mut leaders: Vec<char> = first_choice
2783 .iter()
2784 .filter(|(_, v)| **v == best)
2785 .map(|(k, _)| *k)
2786 .collect();
2787 let mut tie_break = None;
2788 if leaders.len() > 1 {
2789 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2790 let borda_leaders: Vec<char> = leaders
2791 .iter()
2792 .copied()
2793 .filter(|l| borda[l] == top_borda)
2794 .collect();
2795 tie_break = Some(if borda_leaders.len() == 1 {
2796 format!(
2797 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2798 leaders.len()
2799 )
2800 } else {
2801 format!(
2802 "{} way tie on both first-choice votes and Borda points, broken by label order",
2803 leaders.len()
2804 )
2805 });
2806 leaders = borda_leaders;
2807 leaders.sort_unstable();
2808 }
2809 let winner = *leaders
2810 .first()
2811 .or(viable.first())
2812 .context("no candidate to declare a winner from")?;
2813
2814 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2815 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2816 let deliberated = !self.state.deliberation.is_empty();
2817
2818 let quota_seats: std::collections::BTreeSet<&str> =
2822 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2823 let mut present = 0usize;
2824 for (i, j) in self.state.judgements.iter().enumerate() {
2825 if quota_seats.contains(j.seat.as_str()) {
2826 continue;
2827 }
2828 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2829 let voted = self
2830 .state
2831 .votes
2832 .iter()
2833 .any(|v| v.judge == i + 1 && v.vote.is_some());
2834 if ranked || voted {
2835 present += 1;
2836 }
2837 }
2838 let needs_quorum = viable.len() > 1;
2844 let judges_total = if needs_quorum {
2845 self.roles.judges.len()
2846 } else {
2847 0
2848 };
2849 let quorum = if needs_quorum {
2850 judges_total / 2 + 1
2851 } else {
2852 0
2853 };
2854 let met_quorum = !needs_quorum || present >= quorum;
2855 let uncontested = (!needs_quorum).then(|| {
2856 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2857 });
2858
2859 self.state.event(
2860 "tally",
2861 match &uncontested {
2862 Some(reason) => format!("winner {winner} — {reason}"),
2863 None => format!(
2864 "winner {winner} — votes {} | initial {} | {} changed | \
2865 {present}/{judges_total} judges{}",
2866 first_choice
2867 .iter()
2868 .map(|(k, v)| format!("{k}:{v}"))
2869 .collect::<Vec<_>>()
2870 .join(" "),
2871 if unanimous_initial {
2872 "unanimous"
2873 } else {
2874 "split"
2875 },
2876 changed_votes,
2877 if met_quorum {
2878 String::new()
2879 } else {
2880 format!(" — below quorum ({quorum} required)")
2881 },
2882 ),
2883 },
2884 );
2885 if !met_quorum {
2886 self.state.event(
2887 "stall",
2888 format!(
2889 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2890 the run stops here, resumable"
2891 ),
2892 );
2893 }
2894 self.state.tally = Some(Tally {
2895 first_choice,
2896 borda,
2897 winner,
2898 rankings: tops.len(),
2899 unanimous_initial,
2900 deliberated,
2901 changed_votes,
2902 unanimous_final,
2903 tie_break,
2904 judges: judges_total,
2905 present,
2906 quorum,
2907 met_quorum,
2908 uncontested,
2909 });
2910 self.state.status = if met_quorum {
2911 RunStatus::Reviewing
2912 } else {
2913 RunStatus::Stalled
2914 };
2915 self.state.save()?;
2916 Ok(())
2917 }
2918
2919 #[allow(clippy::too_many_lines)]
2940 async fn recover_stall(&mut self) -> Result<bool> {
2941 let run_id = self.state.id.clone();
2946 let prompts = self.state.config.prompts.clone();
2947 let quota_seats: BTreeSet<&str> =
2952 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2953 let absent: Vec<String> = self
2954 .state
2955 .judgements
2956 .iter()
2957 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2958 .map(|j| j.seat.clone())
2959 .collect();
2960 if absent.is_empty() {
2961 return Ok(false);
2962 }
2963 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2964 if viable.len() <= 1 {
2965 return Ok(false);
2966 }
2967 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2968 let language = self.state.config.graph.language.clone();
2969 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2970 let sessions = self.state.config.graph.sessions;
2971 let artifacts = agent::artifacts_dir(&self.state.dir());
2972 let root = self.state.worktree_root();
2973 let base_short = short(&self.state.base_commit);
2974 let candidates: Vec<Candidate> = viable.clone();
2975
2976 let mut positions: Vec<usize> = absent
2978 .iter()
2979 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2980 .collect();
2981 if positions.is_empty() {
2982 return Ok(false);
2983 }
2984 positions.sort_unstable();
2985 positions.dedup();
2986
2987 let mut judge_jobs = Vec::new();
2989 for &j in &positions {
2990 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2991 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2992 let seat_key = format!("judge-{}", j + 1);
2993 let spec = self.roles.judges[j].clone();
2994 let seat = self.seat(&seat_key, &spec.id);
2995 judge_jobs.push(SeatJob {
2996 spec,
2997 seat,
2998 prompt: prompt::judge(
2999 &self.state.instruction,
3000 &views,
3001 self.roles.judges.len(),
3002 &base_short,
3003 &language,
3004 ),
3005 cwd: root.join(seat_key),
3006 timeout,
3007 allow_write: false,
3008 sessions,
3009 artifacts: artifacts.clone(),
3010 stem: format!("judge-{}-recover", j + 1),
3011 });
3012 }
3013
3014 let labels_for_check = labels.clone();
3015 let mut judge_losses = Vec::new();
3016 let retries = self.state.config.graph.retries;
3017 let cache = self.state.config.cache_dir();
3018 let ctx = WaveCtx {
3019 run: &run_id,
3020 node: "judge",
3021 prompts: &prompts,
3022 cache: cache.as_deref(),
3023 round: None,
3024 };
3025 let results = ask_json_wave::<Ranking>(
3026 judge_jobs,
3027 Arc::clone(&self.sem),
3028 retries,
3029 &ctx,
3030 &mut judge_losses,
3031 &mut self.state,
3032 &move |r: &Ranking| r.validate(&labels_for_check),
3033 )
3034 .await;
3035
3036 let mut recovered: BTreeSet<usize> = BTreeSet::new();
3038 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
3039 self.state.seats.insert(seat.key.clone(), seat);
3040 let record = &mut self.state.judgements[j];
3041 match res {
3042 Ok((ranking, out)) => {
3043 record.ranking = ranking.normalized();
3044 record.reasons = ranking.reasons;
3045 record.confidence = ranking.confidence;
3046 record.failed = None;
3047 record.duration_ms = out.duration_ms;
3048 recovered.insert(j);
3049 self.state.event(
3050 "recover",
3051 format!("judge {} ranked again after the limit", j + 1),
3052 );
3053 }
3054 Err(e) => {
3055 self.state
3056 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
3057 }
3058 }
3059 }
3060
3061 let mut vote_jobs = Vec::new();
3063 let mut vote_pos: Vec<usize> = Vec::new();
3064 for &j in &recovered {
3065 let seat_key = format!("judge-{}", j + 1);
3066 let spec = self.roles.judges[j].clone();
3067 let seat = self.seat(&seat_key, &spec.id);
3068 let mut text = prompt::final_vote(&labels, &language);
3069 if !has_context(&spec, &seat, sessions) {
3070 text = format!(
3071 "{}\n\n# Candidates\n\n{}",
3072 text,
3073 self.candidate_block(&candidates, &base_short)
3074 );
3075 }
3076 vote_jobs.push(SeatJob {
3077 spec,
3078 seat,
3079 prompt: text,
3080 cwd: root.join(seat_key),
3081 timeout,
3082 allow_write: false,
3083 sessions,
3084 artifacts: artifacts.clone(),
3085 stem: format!("vote-judge-{}-recover", j + 1),
3086 });
3087 vote_pos.push(j);
3088 }
3089 let allowed = labels.clone();
3090 let mut vote_losses = Vec::new();
3091 let vote_retries = self.state.config.graph.retries;
3092 let vote_cache = self.state.config.cache_dir();
3093 let ctx = WaveCtx {
3094 run: &run_id,
3095 node: "vote",
3096 prompts: &prompts,
3097 cache: vote_cache.as_deref(),
3098 round: None,
3099 };
3100 let votes = ask_json_wave::<FinalVote>(
3101 vote_jobs,
3102 Arc::clone(&self.sem),
3103 vote_retries,
3104 &ctx,
3105 &mut vote_losses,
3106 &mut self.state,
3107 &move |v: &FinalVote| match v.label() {
3108 Some(c) if allowed.contains(&c) => Ok(()),
3109 other => bail!("vote {other:?} is not one of {allowed:?}"),
3110 },
3111 )
3112 .await;
3113 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
3114 let agent_id = seat.agent.clone();
3115 self.state.seats.insert(seat.key.clone(), seat);
3116 match res {
3117 Ok((v, _)) => {
3118 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
3119 rec.vote = v.label();
3120 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
3121 } else {
3122 self.state.votes.push(VoteRecord {
3123 judge: j + 1,
3124 agent: agent_id,
3125 vote: v.label(),
3126 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
3127 changed: false,
3128 });
3129 }
3130 self.state.event(
3131 "recover",
3132 format!("judge {} voted again after the limit", j + 1),
3133 );
3134 }
3135 Err(e) => {
3136 self.state
3137 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
3138 }
3139 }
3140 }
3141
3142 let recovered_keys: BTreeSet<String> = recovered
3146 .iter()
3147 .map(|&j| format!("judge-{}", j + 1))
3148 .collect();
3149 self.state
3150 .quota
3151 .retain(|q| !recovered_keys.contains(&q.seat));
3152 for loss in judge_losses.into_iter().chain(vote_losses) {
3156 if recovered_keys.contains(&loss.seat) {
3157 continue;
3158 }
3159 self.state.quota.retain(|q| q.seat != loss.seat);
3160 self.state.quota.push(loss);
3161 }
3162
3163 self.state.tally = None;
3165 self.tally()?;
3166 Ok(self
3167 .state
3168 .tally
3169 .as_ref()
3170 .map(|t| t.met_quorum)
3171 .unwrap_or(false))
3172 }
3173
3174 async fn fold_losers(&mut self) -> Result<()> {
3177 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
3178 return Ok(());
3179 };
3180 let repo = self.state.repo.clone();
3181 let mut folded = Vec::new();
3182 for i in 0..self.state.candidates.len() {
3183 let c = &self.state.candidates[i];
3184 if c.label == winner || c.folded {
3185 continue;
3186 }
3187 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
3188 git::worktree_remove(&repo, &wt).await.ok();
3189 git::branch_delete(&repo, &branch).await.ok();
3190 self.state.candidates[i].folded = true;
3191 folded.push(label.to_string());
3192 }
3193 let root = self.state.worktree_root();
3195 for j in 1..=self.roles.judges.len() {
3196 let wt = root.join(format!("judge-{j}"));
3197 if wt.exists() {
3198 git::worktree_remove(&repo, &wt).await.ok();
3199 }
3200 }
3201 if self.state.config.graph.advise {
3204 for k in 1..=self.state.config.graph.advisors {
3205 let wt = root.join(format!("advisor-{k}"));
3206 if wt.exists() {
3207 git::worktree_remove(&repo, &wt).await.ok();
3208 }
3209 }
3210 }
3211 if !folded.is_empty() {
3212 self.state
3213 .event("fold", format!("folded candidates {}", folded.join(", ")));
3214 self.state.save()?;
3215 }
3216 Ok(())
3217 }
3218
3219 async fn sync_to_base(&mut self) -> Result<()> {
3249 if self
3250 .state
3251 .base_sync
3252 .as_ref()
3253 .is_some_and(|s| s.conflict.is_some())
3254 {
3255 return Ok(());
3256 }
3257 let Some(winner) = self.state.winner().cloned() else {
3258 return Ok(());
3259 };
3260
3261 let repo = self.state.repo.clone();
3262 let remote = self.state.config.merge.remote.clone();
3263 let base_branch = self.state.base_branch.clone();
3264 let tracking = format!("{remote}/{base_branch}");
3265
3266 git::fetch(&repo, &remote, &base_branch).await.ok();
3267 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3271 return Ok(());
3272 };
3273
3274 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3275 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3276 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3277
3278 if behind == 0 {
3279 self.state.base_sync = Some(BaseSync {
3280 tip,
3281 behind: 0,
3282 attempts,
3283 conflict: None,
3284 });
3285 self.state.save()?;
3286 return Ok(());
3287 }
3288
3289 if attempts >= BASE_SYNC_ROUNDS {
3290 let why = format!(
3291 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3292 rebase(s); rebasing again would only race it",
3293 winner.branch
3294 );
3295 self.state.status = RunStatus::Blocked;
3296 self.state.base_sync = Some(BaseSync {
3297 tip,
3298 behind,
3299 attempts,
3300 conflict: Some(why.clone()),
3301 });
3302 self.state.event("land", why);
3303 self.state.save()?;
3304 return Ok(());
3305 }
3306
3307 self.state.event(
3308 "land",
3309 format!(
3310 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3311 winner.branch
3312 ),
3313 );
3314 self.state.save()?;
3315
3316 let branch_tracking = format!("{remote}/{}", winner.branch);
3322 let fetched_branch = git::fetch(&repo, &remote, &winner.branch).await;
3323 let remote_tip = if matches!(&fetched_branch, Ok(o) if o.ok()) {
3324 git::rev_parse(&repo, &branch_tracking).await.ok()
3325 } else {
3326 None
3327 };
3328 if let Some(theirs) = &remote_tip
3332 && !git::is_ancestor(&repo, theirs, &head).await
3333 && !crate::reconcile::origin_missing(&repo, &head, theirs)
3334 .await
3335 .is_ok_and(|missing| missing.is_empty())
3336 {
3337 let why = format!(
3338 "{branch_tracking} ({}) has commits {} does not contain; not rebasing over \
3339 them",
3340 short(theirs),
3341 winner.branch
3342 );
3343 self.state.status = RunStatus::Blocked;
3344 self.state.base_sync = Some(BaseSync {
3345 tip,
3346 behind,
3347 attempts,
3348 conflict: Some(why.clone()),
3349 });
3350 self.state.event("land", why);
3351 self.state.save()?;
3352 return Ok(());
3353 }
3354
3355 let scratch = self.state.dir().join("base-sync");
3356 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
3357 let attempts = attempts + 1;
3358 match rebased {
3359 Ok(None) => {
3360 git::sync_to_head(&winner.worktree).await?;
3364 let mut conflict = None;
3365 if let Some(pinned) = &remote_tip {
3366 let pushed = git::push_pinned(&repo, &remote, &winner.branch, pinned).await;
3367 match pushed {
3368 Ok(o) if o.ok() => self.state.event(
3369 "land",
3370 format!("pushed rebased {} to {remote}", winner.branch),
3371 ),
3372 Ok(o) => {
3373 conflict = Some(format!(
3374 "rebased {} locally but {remote} refused the push (it moved since {}; someone may have pushed): {}",
3375 winner.branch,
3376 short(pinned),
3377 o.stderr.chars().take(600).collect::<String>()
3378 ));
3379 }
3380 Err(e) => {
3381 conflict = Some(format!(
3382 "rebased {} locally but could not push it: {e:#}",
3383 winner.branch
3384 ));
3385 }
3386 }
3387 }
3388 if let Some(why) = &conflict {
3389 self.state.status = RunStatus::Blocked;
3390 self.state.event("land", why.clone());
3391 }
3392 self.state.base_sync = Some(BaseSync {
3393 tip: tip.clone(),
3394 behind: 0,
3395 attempts,
3396 conflict,
3397 });
3398 self.state
3399 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3400 }
3401 Ok(Some(conflict)) => {
3402 let why = format!(
3403 "{} conflicts with {tracking} and did not rebase: {}",
3404 winner.branch,
3405 conflict.chars().take(600).collect::<String>()
3406 );
3407 self.state.status = RunStatus::Blocked;
3408 self.state.base_sync = Some(BaseSync {
3409 tip,
3410 behind,
3411 attempts,
3412 conflict: Some(why.clone()),
3413 });
3414 self.state.event("land", why);
3415 }
3416 Err(e) => {
3417 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3418 self.state.status = RunStatus::Blocked;
3419 self.state.base_sync = Some(BaseSync {
3420 tip,
3421 behind,
3422 attempts,
3423 conflict: Some(why.clone()),
3424 });
3425 self.state.event("land", why);
3426 }
3427 }
3428 self.state.save()?;
3429 Ok(())
3430 }
3431
3432 fn landing_base(&self) -> String {
3442 self.state
3443 .base_sync
3444 .as_ref()
3445 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3446 }
3447
3448 pub async fn fix_selected(
3481 &mut self,
3482 ids: &[String],
3483 reason: &str,
3484 allow_stale: bool,
3485 ) -> Result<()> {
3486 let reason = reason.trim();
3487 if reason.is_empty() {
3488 bail!("a fix request needs a reason — that is the operator's own record of why");
3489 }
3490 if ids.is_empty() {
3491 bail!("no finding id given");
3492 }
3493 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3494 bail!(
3495 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3496 has already concluded — can be given a targeted fix. A run still \
3497 in progress should simply be resumed; a `merged` run's branch has \
3498 already landed, so its answer is a fresh `magi review <branch>`, \
3499 not reopening this run's own record",
3500 self.state.id,
3501 self.state.status.as_str()
3502 );
3503 }
3504 let Some(winner) = self.state.winner().cloned() else {
3505 bail!("run {} has no winning candidate to fix", self.state.id);
3506 };
3507 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3508 bail!(
3509 "branch `{}` no longer exists; this run cannot be extended",
3510 winner.branch
3511 );
3512 }
3513 let home = crate::run::home();
3514 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3515 bail!(
3516 "run {} is currently being worked on by another magi process",
3517 self.state.id
3518 );
3519 }
3520 let _claim = FixClaim::acquire(&self.state.dir())?;
3526
3527 let mut seen = BTreeSet::new();
3531 let mut findings = Vec::new();
3532 let mut missing = Vec::new();
3533 for id in ids {
3534 if !seen.insert(id.clone()) {
3535 continue;
3536 }
3537 match self.state.finding(id) {
3538 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3539 id: f.id.clone(),
3540 severity: f.severity,
3541 reviewer_vote: rec.vote,
3542 round: round.round,
3543 round_head: round.head.clone(),
3544 reviewer: rec.reviewer,
3545 agent: rec.agent.clone(),
3546 file: f.file.clone(),
3547 line: f.line,
3548 title: f.title.clone(),
3549 detail: f.detail.clone(),
3550 outcome: OperatorFixOutcome::Pending,
3551 }),
3552 None => missing.push(id.clone()),
3553 }
3554 }
3555 if !missing.is_empty() {
3556 bail!(
3557 "unknown finding id(s): {}; nothing was changed",
3558 missing.join(", ")
3559 );
3560 }
3561
3562 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3563 let stale_details: Vec<(String, String)> = findings
3564 .iter()
3565 .filter(|f| f.round_head != head_at_request)
3566 .map(|f| (f.id.clone(), f.round_head.clone()))
3567 .collect();
3568 let stale = !stale_details.is_empty();
3569 if stale && !allow_stale {
3570 bail!(
3571 "the branch has moved since some finding(s) were raised — {} — now \
3572 at {}; pass --allow-stale to fix anyway, or re-run review first",
3573 stale_details
3574 .iter()
3575 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3576 .collect::<Vec<_>>()
3577 .join(", "),
3578 short(&head_at_request)
3579 );
3580 }
3581
3582 let request = OperatorFixRequest {
3583 requested_at: Timestamp::now(),
3584 reason: reason.to_owned(),
3585 findings,
3586 head_at_request: head_at_request.clone(),
3587 allow_stale,
3588 stale,
3589 fix: None,
3590 result_head: None,
3591 follow_up_review_run: None,
3592 };
3593 self.state.event(
3594 "fix",
3595 format!(
3596 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3597 request.findings.len(),
3598 request
3599 .findings
3600 .iter()
3601 .map(|f| f.id.as_str())
3602 .collect::<Vec<_>>()
3603 .join(", "),
3604 ),
3605 );
3606 self.state.operator_fixes.push(request);
3613 self.state.save()?;
3614 let request_index = self.state.operator_fixes.len() - 1;
3615
3616 if winner.worktree.exists() {
3625 let dirty = git::git(
3628 &winner.worktree,
3629 &["status", "--porcelain", "--untracked-files=all"],
3630 )
3631 .await?;
3632 let only_withheld = dirty.lines().all(|l| {
3633 l.strip_prefix("?? ")
3634 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3635 });
3636 if !only_withheld {
3637 bail!(
3638 "`{}` has uncommitted changes; refusing to touch it — commit or \
3639 discard them first",
3640 winner.worktree.display()
3641 );
3642 }
3643 git::worktree_remove(&self.state.repo, &winner.worktree)
3644 .await
3645 .ok();
3646 }
3647 let fix_worktree = self.state.worktree_root().join("operator-fix");
3648 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3649 git::git(
3650 &self.state.repo,
3651 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3652 )
3653 .await
3654 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3655 if !git::is_clean(&fix_worktree).await? {
3656 git::worktree_remove(&self.state.repo, &fix_worktree)
3657 .await
3658 .ok();
3659 bail!(
3660 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3661 winner.branch
3662 );
3663 }
3664
3665 let run_id = self.state.id.clone();
3666 let prompts = self.state.config.prompts.clone();
3667 let language = self.state.config.graph.language.clone();
3668 let sessions = self.state.config.graph.sessions;
3669 let artifacts = agent::artifacts_dir(&self.state.dir());
3670 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3671 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3672 _ => (
3673 self.state
3674 .config
3675 .agent(&winner.agent)
3676 .cloned()
3677 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3678 format!("impl-{}", winner.label),
3679 ),
3680 };
3681 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3682 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3683 .findings
3684 .iter()
3685 .map(|f| Finding {
3686 id: f.id.clone(),
3687 severity: f.severity,
3688 file: f.file.clone(),
3689 line: f.line,
3690 title: f.title.clone(),
3691 detail: f.detail.clone(),
3692 })
3693 .collect();
3694 let job = SeatJob {
3695 prompt: prompt::operator_fix(
3696 &self.state.instruction,
3697 &finding_list,
3698 reason,
3699 &stale_details,
3700 &head_at_request,
3701 &language,
3702 ),
3703 spec: fix_spec.clone(),
3704 seat,
3705 cwd: fix_worktree.clone(),
3706 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3707 allow_write: true,
3708 sessions,
3709 artifacts: artifacts.clone(),
3710 stem: "operator-fix".to_owned(),
3711 };
3712 let cache = self.state.config.cache_dir();
3713 let ctx = WaveCtx {
3714 run: &run_id,
3715 node: "fix",
3716 prompts: &prompts,
3717 cache: cache.as_deref(),
3718 round: None,
3719 };
3720 let (seat, out) =
3721 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3722 let agent_id = seat.agent.clone();
3723
3724 let mut fix = FixRecord {
3725 agent: agent_id,
3726 addressed: Vec::new(),
3727 rejected: Vec::new(),
3728 notes: String::new(),
3729 committed: false,
3730 failed: None,
3731 duration_ms: 0,
3732 continuation: None,
3733 };
3734 let mut final_seat = seat.clone();
3735 match out {
3736 AgentOutcome::Ok(o) => {
3737 fix.duration_ms = o.duration_ms;
3738 let parsed = verdict::extract_json::<FixReport>(&o.text);
3739 let incomplete_reason = match &parsed {
3740 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3741 "the reply parsed, but it reported a command whose own CLI \
3742 never confirmed an exit status"
3743 .to_owned(),
3744 ),
3745 Ok(_) => None,
3746 Err(e) => Some(e.to_string()),
3747 };
3748 match incomplete_reason {
3749 None => {
3750 let report = parsed.expect("checked Ok above");
3751 fix.addressed = report.addressed;
3752 fix.rejected = report.rejected;
3753 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3754 }
3755 Some(reason) => {
3756 let (resumed_seat, resolved, failure, cont) = self
3757 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3758 .await;
3759 fix.duration_ms += cont.cumulative_wait_ms;
3760 fix.continuation = Some(cont);
3761 final_seat = resumed_seat;
3762 match resolved {
3763 Some(report) => {
3764 fix.addressed = report.addressed;
3765 fix.rejected = report.rejected;
3766 fix.notes =
3767 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3768 }
3769 None => fix.failed = failure,
3770 }
3771 }
3772 }
3773 }
3774 AgentOutcome::Dropped(o) => {
3775 fix.duration_ms = o.duration_ms;
3776 let why = o
3777 .dropped
3778 .as_ref()
3779 .map(|d| d.why.as_str())
3780 .unwrap_or("the CLI ended the stream without delivering its answer");
3781 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3782 }
3783 AgentOutcome::Quota(o) => {
3784 self.state.quota.push(QuotaLoss {
3785 seat: final_seat.key.clone(),
3786 node: "fix".to_owned(),
3787 at: Timestamp::now(),
3788 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3789 });
3790 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3791 }
3792 AgentOutcome::Failed(e) => fix.failed = Some(e),
3793 }
3794 if fix.continuation.is_none() {
3795 fix.continuation = Some(ContinuationRecord::not_needed());
3796 }
3797 self.state.seats.insert(final_seat.key.clone(), final_seat);
3798
3799 let rescue_message = format!(
3800 "magi: operator-selected fix ({}) (uncommitted work)",
3801 self.state.operator_fixes[request_index]
3802 .findings
3803 .iter()
3804 .map(|f| f.id.as_str())
3805 .collect::<Vec<_>>()
3806 .join(", ")
3807 );
3808 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3809 self.state.note_withheld("fix", &r.withheld);
3810 }
3811 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3812 fix.committed = after != head_at_request;
3813 git::worktree_remove(&self.state.repo, &fix_worktree)
3814 .await
3815 .ok();
3816
3817 self.state.event(
3818 "fix",
3819 match &fix.failed {
3820 Some(reason) => format!(
3821 "operator fix: adoption report was lost ({reason}); {}",
3822 if fix.committed {
3823 "committed"
3824 } else {
3825 "NO new commit"
3826 }
3827 ),
3828 None => format!(
3829 "operator fix: {} addressed, {} rejected, {}",
3830 fix.addressed.len(),
3831 fix.rejected.len(),
3832 if fix.committed {
3833 "committed"
3834 } else {
3835 "NO new commit"
3836 }
3837 ),
3838 },
3839 );
3840
3841 for f in &mut self.state.operator_fixes[request_index].findings {
3848 f.outcome = if fix.failed.is_some() {
3849 OperatorFixOutcome::Unreported
3850 } else if fix.addressed.contains(&f.id) {
3851 OperatorFixOutcome::Addressed
3852 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3853 OperatorFixOutcome::Rejected { why: r.why.clone() }
3854 } else {
3855 OperatorFixOutcome::Unreported
3856 };
3857 }
3858
3859 let committed = fix.committed;
3860 if committed {
3861 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3862 }
3863 self.state.operator_fixes[request_index].fix = Some(fix);
3864 self.state.save()?;
3867
3868 if committed {
3869 self.state.event(
3870 "fix",
3871 format!(
3872 "operator fix committed {}; opening a follow-up review-only run",
3873 short(&after)
3874 ),
3875 );
3876 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3877 Ok(mut follow_up) => {
3878 follow_up.state.event(
3879 "start",
3880 format!(
3881 "requested by an operator fix on run {} for finding(s) {}",
3882 self.state.id,
3883 self.state.operator_fixes[request_index]
3884 .findings
3885 .iter()
3886 .map(|f| f.id.as_str())
3887 .collect::<Vec<_>>()
3888 .join(", "),
3889 ),
3890 );
3891 follow_up.state.save()?;
3892 let follow_up_id = follow_up.state.id.clone();
3893 if let Err(e) = follow_up.execute().await {
3894 self.state.event(
3895 "fix",
3896 format!(
3897 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3898 ),
3899 );
3900 }
3901 self.state.operator_fixes[request_index].follow_up_review_run =
3902 Some(follow_up_id);
3903 }
3904 Err(e) => {
3905 self.state.event(
3906 "fix",
3907 format!("committed the fix but could not open a follow-up review: {e:#}"),
3908 );
3909 }
3910 }
3911 self.state.save()?;
3912 }
3913
3914 Ok(())
3915 }
3916
3917 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3924 match &self.roles.fixer {
3925 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3926 _ => (
3927 self.state
3928 .config
3929 .agent(&winner.agent)
3930 .cloned()
3931 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3932 format!("impl-{}", winner.label),
3933 ),
3934 }
3935 }
3936
3937 async fn review_loop(&mut self) -> Result<()> {
3938 if self
3943 .state
3944 .base_sync
3945 .as_ref()
3946 .is_some_and(|s| s.conflict.is_some())
3947 {
3948 return Ok(());
3949 }
3950 let run_id = self.state.id.clone();
3955 let prompts = self.state.config.prompts.clone();
3956 let Some(winner) = self.state.winner().cloned() else {
3957 return Ok(());
3958 };
3959 let max_rounds = self.state.config.graph.review_rounds;
3960 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
3970 self.state.status = status;
3971 self.state.save()?;
3972 return Ok(());
3973 }
3974 self.state.status = RunStatus::Reviewing;
3975 if self
3985 .state
3986 .reviews
3987 .last()
3988 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
3989 {
3990 let shell = self.state.config.shell();
3991 return self
3992 .stop_reviewing(
3993 "the last round's own verification never resolved",
3994 &shell,
3995 &winner.worktree,
3996 )
3997 .await;
3998 }
3999
4000 let repo = self.state.repo.clone();
4001 let root = self.state.worktree_root();
4002 let language = self.state.config.graph.language.clone();
4003 let sessions = self.state.config.graph.sessions;
4004 let artifacts = agent::artifacts_dir(&self.state.dir());
4005 let base = self.landing_base();
4006 let base_short = short(&base);
4007 let reviewers = self.roles.reviewers.clone();
4008 let shell = self.state.config.shell();
4009
4010 for round in (self.state.reviews.len() + 1)..=max_rounds {
4011 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4012 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
4013 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
4014 let prev_verification = self
4023 .state
4024 .reviews
4025 .last()
4026 .and_then(|r| r.verification_summary(&head));
4027
4028 let mut jobs = Vec::new();
4032 for (r, spec) in reviewers.iter().cloned().enumerate() {
4033 let wt = root.join(format!("review-{}", r + 1));
4034 if wt.exists() {
4035 git::reset_detached(&wt, &head).await?;
4036 } else {
4037 git::worktree_add_detached(&repo, &wt, &head).await?;
4038 }
4039 let seat_key = format!("review-{}", r + 1);
4040 let seat = self.seat(&seat_key, &spec.id);
4041 jobs.push(SeatJob {
4042 prompt: prompt::review(&prompt::ReviewCtx {
4043 instruction: &self.state.instruction,
4044 branch: &winner.branch,
4045 base_short: &base_short,
4046 stat: &stat,
4047 patch: &patch,
4048 verification: prev_verification.as_ref(),
4049 reviewers: reviewers.len(),
4050 round,
4051 rounds: max_rounds,
4052 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
4055 lens: Lens::for_seat(r),
4056 language: &language,
4057 }),
4058 spec,
4059 seat,
4060 cwd: wt,
4061 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4062 allow_write: false,
4063 sessions,
4064 artifacts: artifacts.clone(),
4065 stem: format!("review-{round}-{}", r + 1),
4066 });
4067 }
4068
4069 self.state.event(
4070 "review",
4071 format!(
4072 "round {round}: {} reviewers on {}",
4073 jobs.len(),
4074 short(&head)
4075 ),
4076 );
4077 let mut quota_losses = Vec::new();
4078 let review_retries = self.state.config.graph.retries;
4079 let review_cache = self.state.config.cache_dir();
4080 let ctx = WaveCtx {
4081 run: &run_id,
4082 node: "review",
4083 prompts: &prompts,
4084 cache: review_cache.as_deref(),
4085 round: Some(round),
4086 };
4087 let results = ask_json_wave::<Review>(
4088 jobs,
4089 Arc::clone(&self.sem),
4090 review_retries,
4091 &ctx,
4092 &mut quota_losses,
4093 &mut self.state,
4094 &|_: &Review| Ok(()),
4095 )
4096 .await;
4097 let round_quota_missing = quota_losses.len();
4101 self.state.quota.extend(quota_losses);
4102
4103 let mut records = Vec::new();
4104 let mut all_findings = Vec::new();
4105 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
4106 let agent_id = seat.agent.clone();
4107 self.state.seats.insert(seat.key.clone(), seat);
4108 let mut record = ReviewRecord {
4109 reviewer: r + 1,
4110 agent: agent_id,
4111 summary: String::new(),
4112 findings: Vec::new(),
4113 vote: None,
4114 failed: None,
4115 duration_ms: 0,
4116 attempts,
4122 };
4123 match res {
4124 Ok((review, out)) => {
4125 record.summary =
4133 blind::sanitize_prose(&review.summary, &self.state.config.blind);
4134 record.vote = Some(review.vote);
4135 record.duration_ms = out.duration_ms;
4136 for (n, mut f) in review.findings.into_iter().enumerate() {
4137 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
4140 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
4141 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
4142 f.file = f
4148 .file
4149 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
4150 all_findings.push(f.clone());
4151 record.findings.push(f);
4152 }
4153 self.state.event(
4154 "review",
4155 format!(
4156 "round {round}: reviewer {} voted {} with {} finding(s)",
4157 r + 1,
4158 review.vote.label(),
4159 record.findings.len()
4160 ),
4161 );
4162 }
4163 Err(e) => {
4164 record.failed = Some(e.to_string());
4165 self.state.event(
4166 "review",
4167 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
4168 );
4169 }
4170 }
4171 records.push(record);
4172 }
4173
4174 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
4181 let vote_split =
4182 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
4183 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
4184 if vote_split {
4185 self.state.event(
4186 "review",
4187 format!(
4188 "round {round}: votes split ({}) — one round of reconsideration",
4189 initial_votes
4190 .iter()
4191 .map(|v| v.label())
4192 .collect::<Vec<_>>()
4193 .join(", ")
4194 ),
4195 );
4196 let panel: Vec<ReviewSeatReport<'_>> = records
4199 .iter()
4200 .filter_map(|r| {
4201 r.vote.map(|vote| ReviewSeatReport {
4202 reviewer: r.reviewer,
4203 vote,
4204 summary: &r.summary,
4205 findings: &r.findings,
4206 })
4207 })
4208 .collect();
4209
4210 let mut jobs = Vec::new();
4211 let mut seats_at = Vec::new();
4212 for (r, spec) in reviewers.iter().cloned().enumerate() {
4213 if records[r].vote.is_none() {
4217 continue;
4218 }
4219 let wt = root.join(format!("review-{}", r + 1));
4220 let seat_key = format!("review-{}", r + 1);
4221 let seat = self.seat(&seat_key, &spec.id);
4222 let patch_ctx = if has_context(&spec, &seat, sessions) {
4227 None
4228 } else {
4229 Some(ReviewPatch {
4230 branch: &winner.branch,
4231 base_short: &base_short,
4232 stat: &stat,
4233 patch: &patch,
4234 })
4235 };
4236 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
4237 instruction: &self.state.instruction,
4238 reviewer: r + 1,
4239 lens: Lens::for_seat(r),
4240 panel: &panel,
4241 patch: patch_ctx,
4242 round,
4243 rounds: max_rounds,
4244 language: &language,
4245 });
4246 jobs.push(SeatJob {
4247 prompt,
4248 spec,
4249 seat,
4250 cwd: wt,
4251 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4252 allow_write: false,
4253 sessions,
4254 artifacts: artifacts.clone(),
4255 stem: format!("review-{round}-reconsider-{}", r + 1),
4256 });
4257 seats_at.push(r);
4258 }
4259
4260 let mut recon_quota_losses = Vec::new();
4261 let recon_cache = self.state.config.cache_dir();
4262 let recon_ctx = WaveCtx {
4263 run: &run_id,
4264 node: "review",
4265 prompts: &prompts,
4266 cache: recon_cache.as_deref(),
4267 round: Some(round),
4268 };
4269 let recon_results = ask_json_wave::<ReviewRevote>(
4270 jobs,
4271 Arc::clone(&self.sem),
4272 review_retries,
4273 &recon_ctx,
4274 &mut recon_quota_losses,
4275 &mut self.state,
4276 &|_: &ReviewRevote| Ok(()),
4277 )
4278 .await;
4279 self.state.quota.extend(recon_quota_losses);
4280
4281 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
4282 let agent_id = seat.agent.clone();
4283 self.state.seats.insert(seat.key.clone(), seat);
4284 let mut rec = ReviewRevoteRecord {
4285 reviewer: r + 1,
4286 agent: agent_id,
4287 vote: None,
4288 reason: String::new(),
4289 failed: None,
4290 };
4291 match res {
4292 Ok((rv, _)) => {
4293 rec.vote = Some(rv.vote);
4294 rec.reason =
4295 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
4296 self.state.event(
4297 "review",
4298 format!(
4299 "round {round}: reviewer {} revoted {}",
4300 r + 1,
4301 rv.vote.label()
4302 ),
4303 );
4304 }
4305 Err(e) => {
4306 rec.failed = Some(e.to_string());
4307 self.state.event(
4308 "review",
4309 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4310 );
4311 }
4312 }
4313 reconsideration.push(rec);
4314 }
4315 } else if initial_votes.len() > 1 {
4316 self.state.event(
4317 "review",
4318 format!(
4319 "round {round}: votes agreed ({}) — no reconsideration",
4320 initial_votes[0].label()
4321 ),
4322 );
4323 }
4324
4325 let final_votes: Vec<ReviewVote> = records
4329 .iter()
4330 .filter_map(|r| {
4331 reconsideration
4332 .iter()
4333 .find(|rv| rv.reviewer == r.reviewer)
4334 .and_then(|rv| rv.vote)
4335 .or(r.vote)
4336 })
4337 .collect();
4338 let round_verdict = ReviewVote::worst(final_votes);
4339
4340 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4341 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4342 let defer_e2e =
4353 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4354 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4355 let reason =
4356 format!("{blocking} blocking finding(s) already required a fix this round");
4357 self.state.event(
4358 "verify",
4359 format!(
4360 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4361 {}); it will run once a round has none left",
4362 short(&head)
4363 ),
4364 );
4365 (Vec::new(), false, true, Some(reason))
4366 } else {
4367 let e2e_commands = self.state.config.verify.e2e.clone();
4368 let cache_dir = self.state.config.cache_dir();
4369 let context = format!("round {round}");
4370 let (e2e, verify_retried) = with_cache_lease(
4371 &mut self.state,
4372 cache_dir.as_deref(),
4373 "e2e",
4374 "e2e",
4375 &winner.worktree,
4376 &head,
4377 verify_timeout,
4378 &context,
4379 |state, budget| {
4380 let shell = shell.clone();
4381 let e2e_commands = e2e_commands.clone();
4382 let worktree = winner.worktree.clone();
4383 let context = context.clone();
4384 async move {
4385 run_e2e_with_retry(
4386 state,
4387 &shell,
4388 &e2e_commands,
4389 &worktree,
4390 budget,
4391 &context,
4392 )
4393 .await
4394 }
4395 },
4396 )
4397 .await;
4398 (e2e, verify_retried, false, None)
4399 };
4400
4401 let expected = records.len();
4402 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4403 let incomplete = answered < expected;
4404 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4405 let policy = self.state.config.graph.incomplete_review;
4406 let clean = round_is_clean(
4407 blocking,
4408 e2e_ok,
4409 answered,
4410 expected,
4411 round_quota_missing,
4412 policy,
4413 );
4414
4415 let mut round_record = ReviewRound {
4416 round,
4417 head: head.clone(),
4418 verified_head: None,
4419 verified_at: None,
4420 reviews: records,
4421 e2e,
4422 verify_retried,
4423 e2e_deferred,
4424 e2e_defer_reason,
4425 fix: None,
4426 blocking,
4427 answered,
4428 expected,
4429 clean,
4430 progressed: false,
4431 vote_split,
4432 reconsideration,
4433 verdict: round_verdict,
4434 };
4435 if !matches!(
4446 round_record.e2e_status(),
4447 E2eStatus::Deferred | E2eStatus::NotConfigured
4448 ) {
4449 round_record.verified_head = Some(head.clone());
4450 round_record.verified_at = Some(Timestamp::now());
4451 }
4452 let this_round_verification = round_record.verification_summary(&head);
4453
4454 if incomplete {
4455 let missing: Vec<String> = round_record
4456 .reviews
4457 .iter()
4458 .filter(|r| r.failed.is_some())
4459 .map(|r| format!("review-{}", r.reviewer))
4460 .collect();
4461 self.state.event(
4462 "review",
4463 format!(
4464 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4465 missing.join(", ")
4466 ),
4467 );
4468 }
4469
4470 if clean {
4471 self.state.event(
4472 "review",
4473 if incomplete && policy == IncompleteReviewPolicy::Warn {
4474 format!(
4475 "round {round}: clean (warn policy, incomplete panel) — no \
4476 blocking findings from the seats that answered, verification green"
4477 )
4478 } else if incomplete {
4479 format!(
4480 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4481 quorum) — no blocking findings from the seats that answered, \
4482 verification green",
4483 expected - answered
4484 )
4485 } else {
4486 format!("round {round}: clean — no blocking findings, verification green")
4487 },
4488 );
4489 self.state.reviews.push(round_record);
4490 self.state.status = RunStatus::Gating;
4491 self.state.save()?;
4492 return Ok(());
4493 }
4494
4495 if incomplete && blocking == 0 && e2e_ok {
4503 self.state.reviews.push(round_record);
4504 self.state.save()?;
4505 if round == max_rounds {
4506 self.state.status = RunStatus::Blocked;
4507 self.state.event(
4508 "review",
4509 format!(
4510 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4511 refusing to call it clean",
4512 expected - answered
4513 ),
4514 );
4515 return Ok(());
4516 }
4517 continue;
4518 }
4519
4520 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4532 self.state.reviews.push(round_record);
4533 return self
4534 .stop_reviewing(
4535 "the round's own verification could not run",
4536 &shell,
4537 &winner.worktree,
4538 )
4539 .await;
4540 }
4541
4542 if round == max_rounds {
4543 self.state.reviews.push(round_record);
4544 return self
4545 .stop_reviewing(
4546 &format!(
4547 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4548 ),
4549 &shell,
4550 &winner.worktree,
4551 )
4552 .await;
4553 }
4554
4555 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4558 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4559 let blocking_findings: Vec<_> = all_findings
4560 .iter()
4561 .filter(|f| f.severity.blocks())
4562 .cloned()
4563 .collect();
4564 let job = SeatJob {
4565 prompt: prompt::fix(
4566 &self.state.instruction,
4567 &blocking_findings,
4568 this_round_verification.as_ref(),
4569 round,
4570 max_rounds,
4571 &language,
4572 ),
4573 spec: fix_spec.clone(),
4574 seat,
4575 cwd: winner.worktree.clone(),
4576 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4577 allow_write: true,
4578 sessions,
4579 artifacts: artifacts.clone(),
4580 stem: format!("fix-{round}"),
4581 };
4582 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4583 let cache = self.state.config.cache_dir();
4584 let ctx = WaveCtx {
4585 run: &run_id,
4586 node: "fix",
4587 prompts: &prompts,
4588 cache: cache.as_deref(),
4589 round: Some(round),
4590 };
4591 let (seat, out) =
4592 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4593 let agent_id = seat.agent.clone();
4594
4595 let mut fix = FixRecord {
4596 agent: agent_id,
4597 addressed: Vec::new(),
4598 rejected: Vec::new(),
4599 notes: String::new(),
4600 committed: false,
4601 failed: None,
4602 duration_ms: 0,
4603 continuation: None,
4604 };
4605 let mut continuation = ContinuationRecord::not_needed();
4606 let mut final_seat = seat.clone();
4607 match out {
4608 AgentOutcome::Ok(o) => {
4609 fix.duration_ms = o.duration_ms;
4610 let parsed = verdict::extract_json::<FixReport>(&o.text);
4611 let incomplete_reason = match &parsed {
4618 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4619 "the reply parsed, but it reported a command whose own CLI never \
4620 confirmed an exit status"
4621 .to_owned(),
4622 ),
4623 Ok(_) => None,
4624 Err(e) => Some(e.to_string()),
4625 };
4626 match incomplete_reason {
4627 None => {
4628 let report = parsed.expect("checked Ok above");
4629 fix.addressed = report.addressed;
4630 fix.rejected = report.rejected;
4631 fix.notes =
4632 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4633 }
4634 Some(reason) => {
4635 let (resumed_seat, resolved, failure, cont) = self
4636 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4637 .await;
4638 fix.duration_ms += cont.cumulative_wait_ms;
4639 continuation = cont;
4640 final_seat = resumed_seat;
4641 match resolved {
4642 Some(report) => {
4643 fix.addressed = report.addressed;
4644 fix.rejected = report.rejected;
4645 fix.notes = blind::sanitize_prose(
4646 &report.notes,
4647 &self.state.config.blind,
4648 );
4649 }
4650 None => fix.failed = failure,
4651 }
4652 }
4653 }
4654 }
4655 AgentOutcome::Dropped(o) => {
4657 fix.duration_ms = o.duration_ms;
4658 let why = o
4659 .dropped
4660 .as_ref()
4661 .map(|d| d.why.as_str())
4662 .unwrap_or("the CLI ended the stream without delivering its answer");
4663 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4664 }
4665 AgentOutcome::Quota(o) => {
4666 self.state.quota.push(QuotaLoss {
4667 seat: final_seat.key.clone(),
4668 node: "fix".to_owned(),
4669 at: Timestamp::now(),
4670 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4671 });
4672 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4673 }
4674 AgentOutcome::Failed(e) => fix.failed = Some(e),
4675 }
4676 fix.continuation = Some(continuation);
4677 self.state.seats.insert(final_seat.key.clone(), final_seat);
4678 if let Ok(r) = git::rescue_commit(
4679 &winner.worktree,
4680 &format!("magi: review round {round} fixes (uncommitted work)"),
4681 )
4682 .await
4683 {
4684 self.state.note_withheld("fix", &r.withheld);
4685 }
4686 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4687 fix.committed = after != before;
4688 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4696 let progressed = diff_after != patch;
4697 let commit_note = if fix.committed {
4698 "committed"
4699 } else {
4700 "NO new commit"
4701 };
4702 let tree_note = if progressed {
4703 "changed vs base"
4704 } else {
4705 "unchanged vs base"
4706 };
4707 self.state.event(
4708 "fix",
4709 match &fix.failed {
4710 Some(reason) => {
4716 format!(
4717 "round {round}: fixer's adoption report was lost ({reason}); \
4718 {commit_note}, tree {tree_note}"
4719 )
4720 }
4721 None => format!(
4722 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4723 {tree_note}{}",
4724 fix.addressed.len(),
4725 fix.rejected.len(),
4726 if continuation.outcome == ContinuationOutcome::Resumed {
4727 format!(
4728 " (adoption report recovered after {} continuation(s))",
4729 continuation.attempts
4730 )
4731 } else {
4732 String::new()
4733 },
4734 ),
4735 },
4736 );
4737 round_record.fix = Some(fix);
4738 round_record.progressed = progressed;
4739 self.state.reviews.push(round_record);
4740 self.state.save()?;
4741
4742 if matches!(
4755 continuation.outcome,
4756 ContinuationOutcome::Exhausted
4757 | ContinuationOutcome::QuotaLost
4758 | ContinuationOutcome::NoSession
4759 ) {
4760 return self
4761 .stop_reviewing(
4762 "the fixer's adoption report never came back, even after resuming its \
4763 own seat; refusing to start another round against the same worktree \
4764 while that is unresolved",
4765 &shell,
4766 &winner.worktree,
4767 )
4768 .await;
4769 }
4770
4771 let streak = self
4772 .state
4773 .reviews
4774 .iter()
4775 .rev()
4776 .take_while(|r| !r.progressed)
4777 .count();
4778 if streak >= STAGNANT_LIMIT {
4779 return self
4780 .stop_reviewing(
4781 &format!(
4782 "the tree has not moved against base for {streak} round(s) in a row"
4783 ),
4784 &shell,
4785 &winner.worktree,
4786 )
4787 .await;
4788 }
4789 }
4790 Ok(())
4791 }
4792
4793 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4823 let round_idx = self.state.reviews.len() - 1;
4824 let needs_catchup_run = matches!(
4832 self.state.reviews[round_idx].e2e_status(),
4833 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4834 );
4835 if needs_catchup_run {
4836 let round = self.state.reviews[round_idx].round;
4837 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4838 let commands = self.state.config.verify.e2e.clone();
4839 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4840 let cache_dir = self.state.config.cache_dir();
4841 let context = format!(
4842 "round {round}: verification unresolved, catching up before the final decision"
4843 );
4844 let (outcomes, verify_retried) = with_cache_lease(
4845 &mut self.state,
4846 cache_dir.as_deref(),
4847 "e2e",
4848 "e2e",
4849 worktree,
4850 &attempted_head,
4851 timeout,
4852 &context,
4853 |state, budget| {
4854 let shell = shell.to_vec();
4855 let commands = commands.clone();
4856 let context = context.clone();
4857 async move {
4858 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4859 .await
4860 }
4861 },
4862 )
4863 .await;
4864 let last = &mut self.state.reviews[round_idx];
4865 last.e2e = outcomes;
4866 last.verify_retried = verify_retried;
4867 last.verified_head = Some(attempted_head);
4874 last.verified_at = Some(Timestamp::now());
4875 if verify_inconclusive(&last.e2e) {
4876 self.state.save()?;
4883 return Ok(());
4884 }
4885 last.e2e_deferred = false;
4886 }
4887 let last = &self.state.reviews[round_idx];
4888 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4889
4890 match last.e2e_status() {
4891 E2eStatus::Failed => {
4892 let red: Vec<String> = last
4893 .e2e
4894 .iter()
4895 .filter(|o| !o.ok())
4896 .map(|o| {
4897 format!(
4898 "`{}` -> {:?}\n{}",
4899 o.command,
4900 o.code,
4901 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4902 )
4903 })
4904 .collect();
4905 self.state
4906 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4907 self.state.status = RunStatus::Blocked;
4908 }
4909 E2eStatus::ResourceBlocked => {
4914 self.state.event(
4915 "review",
4916 format!(
4917 "{why}; e2e could not run (shared build cache unavailable); not \
4918 deciding yet"
4919 ),
4920 );
4921 }
4922 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4923 self.state.event(
4924 "review",
4925 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4926 );
4927 self.state.status = RunStatus::Gating;
4928 }
4929 }
4930 self.state.save()?;
4931 Ok(())
4932 }
4933
4934 async fn gate(&mut self) -> Result<()> {
4937 if self.state.status == RunStatus::Failed
4949 || self
4950 .state
4951 .base_sync
4952 .as_ref()
4953 .is_some_and(|s| s.conflict.is_some())
4954 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
4955 != Some(RunStatus::Gating)
4956 {
4957 return Ok(());
4958 }
4959 if self.state.gate_ran {
4960 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
4971 self.state.status = RunStatus::Blocked;
4972 self.state.save()?;
4973 }
4974 return Ok(());
4975 }
4976 let Some(winner) = self.state.winner().cloned() else {
4977 return Ok(());
4978 };
4979 self.state.status = RunStatus::Gating;
4980 let mut outcomes = self.run_gate(&winner).await?;
4981 loop {
4982 if verify_inconclusive(&outcomes) {
4993 self.state.save()?;
4994 return Ok(());
4995 }
4996 if outcomes.iter().all(CommandOutcome::ok) {
4997 break;
4998 }
4999 match self.gate_fix_round(&winner, &outcomes).await? {
5000 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
5001 GateFix::Stop => break,
5002 GateFix::Defer => {
5003 self.state.save()?;
5004 return Ok(());
5005 }
5006 }
5007 }
5008 let passed = outcomes.iter().all(CommandOutcome::ok);
5009 self.state.gate = outcomes;
5010 self.state.gate_ran = true;
5011 if !passed {
5012 self.state.status = RunStatus::Blocked;
5013 let spent = self.state.gate_fixes.len();
5014 self.state.event(
5015 "gate",
5016 if spent == 0 {
5017 "gate failed; not merging".to_owned()
5018 } else {
5019 format!("gate failed after {spent} gate-fix round(s); not merging")
5020 },
5021 );
5022 }
5023 self.state.save()?;
5024 Ok(())
5025 }
5026
5027 async fn run_pre_gate(&mut self, winner: &Candidate) {
5037 let commands = self.state.config.verify.pre_gate.clone();
5038 if commands.is_empty() {
5039 return;
5040 }
5041 let shell = self.state.config.shell();
5042 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5043 let (outcomes, _) = run_commands(
5044 &mut self.state,
5045 "pre_gate",
5046 "pre_gate",
5047 0,
5048 &shell,
5049 &commands,
5050 &winner.worktree,
5051 timeout,
5052 )
5053 .await;
5054 for o in &outcomes {
5055 if !o.ok() {
5056 tracing::warn!(
5057 "pre_gate `{}` failed ({:?}); the gate decides",
5058 o.command,
5059 o.code
5060 );
5061 }
5062 self.state.event(
5063 "pre_gate",
5064 format!(
5065 "`{}` -> {}",
5066 o.command,
5067 if o.ok() {
5068 "pass".to_owned()
5069 } else {
5070 format!(
5071 "FAIL ({:?})\n{}",
5072 o.code,
5073 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5074 )
5075 }
5076 ),
5077 );
5078 }
5079 self.state.pre_gate = outcomes;
5080 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
5081 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
5082 Ok(head) => {
5083 self.state
5084 .event("pre_gate", format!("committed mechanical fixes ({head})"));
5085 self.state.pre_gate_commit = Some(head);
5086 }
5087 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
5088 },
5089 Ok(false) => {}
5090 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
5091 }
5092 if let Err(e) = self.state.save() {
5093 tracing::warn!("could not persist the pre_gate record: {e:#}");
5094 }
5095 }
5096
5097 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
5100 self.run_pre_gate(winner).await;
5101 let shell = self.state.config.shell();
5102 let gate_commands = self.state.config.verify.gate.clone();
5103 let outcomes = if gate_commands.is_empty() {
5112 Vec::new()
5113 } else {
5114 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5115 let cache_dir = self.state.config.cache_dir();
5116 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
5117 let (outcomes, _) = with_cache_lease(
5118 &mut self.state,
5119 cache_dir.as_deref(),
5120 "gate",
5121 "gate",
5122 &winner.worktree,
5123 &head,
5124 timeout,
5125 "final gate",
5126 |state, budget| {
5127 let shell = shell.clone();
5128 let gate_commands = gate_commands.clone();
5129 let worktree = winner.worktree.clone();
5130 async move {
5131 let (outcomes, timed_out_pids) = run_commands(
5132 state,
5133 "gate",
5134 "gate",
5135 0,
5136 &shell,
5137 &gate_commands,
5138 &worktree,
5139 budget,
5140 )
5141 .await;
5142 (outcomes, false, timed_out_pids)
5143 }
5144 },
5145 )
5146 .await;
5147 outcomes
5148 };
5149 if outcomes.is_empty() {
5150 self.state.event(
5155 "gate",
5156 "no gate commands configured; nothing to check, passing",
5157 );
5158 }
5159 for o in &outcomes {
5160 self.state.event(
5161 "gate",
5162 format!(
5163 "`{}` -> {}",
5164 o.command,
5165 if o.ok() {
5166 "pass".to_owned()
5167 } else {
5168 format!(
5169 "FAIL ({:?})\n{}",
5170 o.code,
5171 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5172 )
5173 }
5174 ),
5175 );
5176 }
5177 Ok(outcomes)
5178 }
5179
5180 async fn gate_fix_round(
5193 &mut self,
5194 winner: &Candidate,
5195 outcomes: &[CommandOutcome],
5196 ) -> Result<GateFix> {
5197 let cap = self.state.config.graph.gate_fix_rounds;
5198 let spent = self.state.gate_fixes.len();
5199 if spent >= cap {
5200 if cap > 0 {
5201 self.state.event(
5202 "gate",
5203 format!("{spent} gate-fix round(s) spent and the gate still fails"),
5204 );
5205 }
5206 return Ok(GateFix::Stop);
5207 }
5208 if !gate_fixable(outcomes) {
5209 self.state.event(
5210 "gate",
5211 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
5212 command or similar); not spending a fix round on it",
5213 );
5214 return Ok(GateFix::Stop);
5215 }
5216 let min_free = self.state.config.disk.min_free_bytes;
5217 if min_free > 0 {
5218 match crate::disk::free_bytes(&winner.worktree) {
5219 Ok(free) if crate::disk::enough_space(free, min_free) => {}
5220 Ok(free) => {
5221 self.state.event(
5222 "gate",
5223 format!(
5224 "only {free} bytes free ({min_free} required by `[disk] \
5225 min_free_bytes`); not spending a fix round on a failure the disk \
5226 may explain"
5227 ),
5228 );
5229 return Ok(GateFix::Stop);
5230 }
5231 Err(e) => {
5232 self.state.event(
5233 "gate",
5234 format!("free disk space could not be measured ({e:#}); no fix round"),
5235 );
5236 return Ok(GateFix::Stop);
5237 }
5238 }
5239 }
5240
5241 let attempt = spent + 1;
5242 let run_id = self.state.id.clone();
5243 let prompts = self.state.config.prompts.clone();
5244 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
5245 let base = self.landing_base();
5246 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
5247 let seat = self.seat(&fix_seat_key, &fix_spec.id);
5248 let job = SeatJob {
5249 prompt: prompt::gate_fix(
5250 &self.state.instruction,
5251 &failed,
5252 attempt,
5253 cap,
5254 &self.state.config.graph.language,
5255 ),
5256 spec: fix_spec,
5257 seat,
5258 cwd: winner.worktree.clone(),
5259 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
5260 allow_write: true,
5261 sessions: self.state.config.graph.sessions,
5262 artifacts: agent::artifacts_dir(&self.state.dir()),
5263 stem: format!("gate-fix-{attempt}"),
5264 };
5265 self.state.event(
5266 "gate",
5267 format!("gate failed; gate-fix round {attempt} of {cap}"),
5268 );
5269 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
5270 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
5271 let cache = self.state.config.cache_dir();
5272 let ctx = WaveCtx {
5273 run: &run_id,
5274 node: "gate-fix",
5275 prompts: &prompts,
5276 cache: cache.as_deref(),
5277 round: None,
5278 };
5279 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
5280 let mut record = GateFixRecord {
5281 agent: seat.agent.clone(),
5282 failed,
5283 notes: String::new(),
5284 committed: false,
5285 error: None,
5286 };
5287 match out {
5288 AgentOutcome::Ok(o) => {
5289 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
5292 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
5293 }
5294 }
5295 AgentOutcome::Dropped(_) => {
5296 record.error = Some("the CLI dropped the stream".to_owned());
5297 }
5298 AgentOutcome::Quota(o) => {
5299 self.state.quota.push(QuotaLoss {
5300 seat: seat.key.clone(),
5301 node: "gate-fix".to_owned(),
5302 at: Timestamp::now(),
5303 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5304 });
5305 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5306 }
5307 AgentOutcome::Failed(e) => record.error = Some(e),
5308 }
5309 self.state.seats.insert(seat.key.clone(), seat);
5310 if let Ok(r) = git::rescue_commit(
5311 &winner.worktree,
5312 &format!("magi: gate fix {attempt} (uncommitted work)"),
5313 )
5314 .await
5315 {
5316 self.state.note_withheld("gate-fix", &r.withheld);
5317 }
5318 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5319 record.committed = after != before;
5320 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5321 let note = record.error.clone();
5322 self.state.gate_fixes.push(record);
5323 self.state.save()?;
5324 if !changed {
5325 self.state.event(
5326 "gate",
5327 match note {
5328 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5329 None => format!("gate-fix round {attempt}: the tree did not change"),
5330 },
5331 );
5332 return Ok(GateFix::Stop);
5333 }
5334 self.state.event(
5335 "gate",
5336 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5337 );
5338
5339 let commands = self.state.config.verify.e2e.clone();
5340 if !commands.is_empty() {
5341 let shell = self.state.config.shell();
5342 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5343 let cache_dir = self.state.config.cache_dir();
5344 let context = format!("gate-fix round {attempt}");
5345 let (e2e, _) = with_cache_lease(
5346 &mut self.state,
5347 cache_dir.as_deref(),
5348 "e2e",
5349 "e2e",
5350 &winner.worktree,
5351 &after,
5352 timeout,
5353 &context,
5354 |state, budget| {
5355 let shell = shell.clone();
5356 let commands = commands.clone();
5357 let context = context.clone();
5358 let worktree = winner.worktree.clone();
5359 async move {
5360 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5361 .await
5362 }
5363 },
5364 )
5365 .await;
5366 if verify_inconclusive(&e2e) {
5367 return Ok(GateFix::Defer);
5368 }
5369 if e2e.iter().any(|o| !o.ok()) {
5370 self.state.event(
5371 "gate",
5372 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5373 );
5374 return Ok(GateFix::Stop);
5375 }
5376 }
5377 Ok(GateFix::Retry)
5378 }
5379
5380 async fn merge(&mut self) -> Result<()> {
5383 if self
5398 .state
5399 .base_sync
5400 .as_ref()
5401 .is_some_and(|s| s.conflict.is_some())
5402 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5403 != Some(RunStatus::Gating)
5404 || !self.state.gate_status().ok()
5413 {
5414 return Ok(());
5415 }
5416 if self.state.merge.is_some() {
5425 return Ok(());
5426 }
5427 let Some(winner) = self.state.winner().cloned() else {
5428 return Ok(());
5429 };
5430 let repo = self.state.repo.clone();
5431 let base = self.state.base_branch.clone();
5432 let mode = self.state.config.merge.mode;
5433 let style = self.state.config.merge.style;
5434 let pr = pr_message(&self.state, winner.label);
5435 let message = pr.commit_message();
5436
5437 let outcome = match mode {
5438 MergeMode::None => MergeOutcome {
5439 mode,
5440 ok: true,
5441 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5442 empty: false,
5443 },
5444 MergeMode::Pr | MergeMode::Local
5445 if merge_is_empty(&repo, &self.state, &winner.branch, mode).await =>
5446 {
5447 MergeOutcome {
5448 mode,
5449 ok: false,
5450 detail: empty_candidate_detail(&self.state, &base),
5451 empty: true,
5452 }
5453 }
5454 MergeMode::Local => {
5455 let on = git::current_branch(&repo).await?;
5456 if on.as_deref() != Some(base.as_str()) {
5457 MergeOutcome {
5458 mode,
5459 ok: false,
5460 detail: format!(
5461 "{} has {} checked out, not the base branch {base}",
5462 repo.display(),
5463 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5464 ),
5465 empty: false,
5466 }
5467 } else if !git::is_clean(&repo).await? {
5468 MergeOutcome {
5469 mode,
5470 ok: false,
5471 detail: format!("{} is dirty; refusing to merge", repo.display()),
5472 empty: false,
5473 }
5474 } else {
5475 let out = match style {
5476 MergeStyle::Merge => {
5477 git::merge_no_ff(&repo, &winner.branch, &message).await?
5478 }
5479 MergeStyle::Squash => {
5480 git::merge_squash(&repo, &winner.branch, &message).await?
5481 }
5482 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5483 };
5484 MergeOutcome {
5485 mode,
5486 ok: out.ok(),
5487 detail: if out.ok() { out.stdout } else { out.stderr },
5488 empty: false,
5489 }
5490 }
5491 }
5492 MergeMode::Pr => {
5493 let remote = self.state.config.merge.remote.clone();
5494 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5495 if !pushed.ok() {
5496 MergeOutcome {
5497 mode,
5498 ok: false,
5499 detail: pushed.stderr,
5500 empty: false,
5501 }
5502 } else {
5503 let out =
5504 gh_pr_create(&winner.worktree, &base, &winner.branch, &pr.title, &pr.body)
5505 .await;
5506 match out {
5507 Ok(url) => MergeOutcome {
5508 mode,
5509 ok: true,
5510 detail: url,
5511 empty: false,
5512 },
5513 Err(e) => MergeOutcome {
5514 mode,
5515 ok: false,
5516 detail: e.to_string(),
5517 empty: false,
5518 },
5519 }
5520 }
5521 }
5522 };
5523
5524 self.state.status = match (mode, outcome.ok) {
5525 (MergeMode::None, _) => RunStatus::Ready,
5526 (_, true) => RunStatus::Merged,
5527 (_, false) => RunStatus::Blocked,
5528 };
5529 self.state.event(
5530 "merge",
5531 format!(
5532 "{:?}: {}",
5533 mode,
5534 outcome.detail.lines().next().unwrap_or("")
5535 ),
5536 );
5537 self.state.merge = Some(outcome);
5538 self.state.save()?;
5539
5540 if self.state.config.graph.land
5546 && mode == MergeMode::Pr
5547 && self.state.status == RunStatus::Merged
5548 {
5549 self.run_land().await?;
5550 }
5551 self.settle_questions();
5556 Ok(())
5557 }
5558
5559 async fn run_land(&mut self) -> Result<()> {
5570 let url = self
5571 .state
5572 .merge
5573 .as_ref()
5574 .map(|m| m.detail.clone())
5575 .unwrap_or_default();
5576 let url = url.lines().next().unwrap_or("").trim().to_owned();
5577 if !url.starts_with("http") {
5578 return Ok(());
5579 }
5580 match land::land(&mut self.state, &url).await {
5583 Ok(pr) if self.state.parked => {
5584 let _ = pr;
5588 }
5589 Ok(pr) => {
5590 self.state.status = match pr.state {
5591 land::PrLifecycle::Merged => RunStatus::Merged,
5592 _ => RunStatus::Blocked,
5593 };
5594 if bump::should_release_bump(self.state.status)
5601 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5602 {
5603 self.state
5609 .event("bump", format!("release bump skipped: {e:#}"));
5610 }
5611 self.state.save()?;
5612 }
5613 Err(e) => {
5614 self.state.status = RunStatus::Blocked;
5615 self.state.event("land", format!("gave up: {e}"));
5616 self.state.save()?;
5617 }
5618 }
5619 Ok(())
5620 }
5621
5622 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5626 if let Some(existing) = self.state.seats.get(key)
5627 && existing.agent == agent
5628 {
5629 return existing.clone();
5630 }
5631 let fresh = SeatState::new(key, agent, self.state.seed);
5632 self.state.seats.insert(key.to_owned(), fresh.clone());
5633 fresh
5634 }
5635
5636 fn view(&self, c: &Candidate) -> CandidateView {
5638 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5639 .unwrap_or_default();
5640 let (patch, _) = blind::sanitize_patch(
5641 &format!("candidate {} patch", c.label),
5642 &raw,
5643 &self.state.config.blind,
5644 );
5645 CandidateView {
5646 label: c.label,
5647 branch: c.branch.clone(),
5648 summary: c.summary.clone(),
5649 stat: c.stat.clone(),
5650 patch,
5651 }
5652 }
5653
5654 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5656 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5657 prompt::judge(
5658 "(see above)",
5659 &views,
5660 self.roles.judges.len(),
5661 base_short,
5662 "en",
5663 )
5664 }
5665
5666 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5673 let mut turns = Vec::new();
5674 for j in &self.state.judgements {
5675 if j.ranking.is_empty() {
5676 continue;
5677 }
5678 let reasons = j
5679 .reasons
5680 .iter()
5681 .map(|(k, v)| format!("- {k}: {v}"))
5682 .collect::<Vec<_>>()
5683 .join("\n");
5684 turns.push(Turn {
5685 who: format!("Judge {} (opening ranking)", j.judge),
5686 is_self: j.judge == self_idx + 1,
5687 body: format!(
5688 "Ranked {}{}{reasons}",
5689 j.ranking.iter().collect::<String>(),
5690 if reasons.is_empty() {
5691 ""
5692 } else {
5693 ", because:\n"
5694 }
5695 ),
5696 });
5697 }
5698 for t in self
5699 .state
5700 .deliberation
5701 .iter()
5702 .flat_map(|r| r.turns.iter())
5703 .chain(current)
5704 {
5705 turns.push(Turn {
5706 who: format!("Judge {}", t.judge),
5707 is_self: t.judge == self_idx + 1,
5708 body: t.body.clone(),
5709 });
5710 }
5711 turns
5712 }
5713}
5714
5715fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5717 agent::has_session(spec.kind, seat, sessions)
5718}
5719
5720fn next_untried_implementer<'a>(
5741 roster: &'a [AgentSpec],
5742 start: usize,
5743 tried: &BTreeSet<String>,
5744) -> Option<&'a AgentSpec> {
5745 roster
5746 .get(start + 1..)?
5747 .iter()
5748 .find(|s| !tried.contains(&s.id))
5749}
5750
5751fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5765 commands.iter().any(|c| c.exit_code.is_none())
5766}
5767
5768fn verified_noop_claim(
5781 usable: bool,
5782 commands: &[agent::CommandEvidence],
5783 text: &str,
5784) -> Option<String> {
5785 (usable && !has_unconfirmed_command(commands))
5786 .then(|| verdict::verified_noop(text))
5787 .flatten()
5788}
5789
5790fn short(commit: &str) -> String {
5791 commit.chars().take(7).collect()
5792}
5793
5794fn make_executable(path: &Path) -> Result<()> {
5795 #[cfg(unix)]
5796 {
5797 use std::os::unix::fs::PermissionsExt as _;
5798 let mut perms = std::fs::metadata(path)?.permissions();
5799 perms.set_mode(0o755);
5800 std::fs::set_permissions(path, perms)?;
5801 }
5802 #[cfg(not(unix))]
5803 {
5804 let _ = path;
5805 }
5806 Ok(())
5807}
5808
5809struct WaveCtx<'a> {
5816 run: &'a str,
5819 node: &'a str,
5821 prompts: &'a Prompts,
5822 cache: Option<&'a Path>,
5824 round: Option<usize>,
5827}
5828
5829async fn run_one(
5831 job: SeatJob,
5832 sem: Arc<Semaphore>,
5833 ctx: &WaveCtx<'_>,
5834 state: &mut RunState,
5835 attempt: usize,
5836) -> (SeatState, AgentOutcome) {
5837 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5838 .await
5839 .pop()
5840 .expect("one job in, one result out");
5841 (seat, out)
5842}
5843
5844async fn wave(
5850 jobs: Vec<SeatJob>,
5851 sem: Arc<Semaphore>,
5852 ctx: &WaveCtx<'_>,
5853 state: &mut RunState,
5854 attempt: usize,
5855) -> Vec<(usize, SeatState, AgentOutcome)> {
5856 let WaveCtx {
5857 run,
5858 node,
5859 prompts,
5860 cache,
5861 round,
5862 } = *ctx;
5863 for job in &jobs {
5864 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5865 }
5866 if let Err(e) = state.save() {
5867 tracing::warn!("could not persist in-progress seats: {e:#}");
5872 }
5873 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5889 let wait_started = Instant::now();
5890 let cache_guard = if let Some(cache_dir) = cache {
5891 if jobs_had_a_writer {
5892 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5893 let budget = jobs
5894 .iter()
5895 .map(|j| j.timeout)
5896 .max()
5897 .unwrap_or(Duration::from_secs(60));
5898 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5899 .await
5900 .ok()
5901 } else {
5902 None
5903 }
5904 } else {
5905 None
5906 };
5907 let waited_for_lease = wait_started.elapsed();
5914 let mut set = tokio::task::JoinSet::new();
5915 let overlay = prompts.overlay(node);
5916 for (i, mut job) in jobs.into_iter().enumerate() {
5917 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5918 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
5919 if cache.is_some() {
5920 job.prompt.push('\n');
5921 job.prompt
5922 .push_str(&prompt::build_cache_note(node, job.allow_write));
5923 }
5924 let sem = Arc::clone(&sem);
5925 let run = run.to_owned();
5926 let node = node.to_owned();
5927 let cache = cache
5938 .filter(|_| job.allow_write && cache_guard.is_some())
5939 .map(Path::to_path_buf);
5940 set.spawn(async move {
5941 let _permit = sem.acquire().await;
5942 let mut seat = job.seat;
5943 let out = agent::invoke(
5944 &job.spec,
5945 &mut seat,
5946 &Invocation {
5947 cwd: &job.cwd,
5948 prompt: &job.prompt,
5949 timeout: job.timeout,
5950 allow_write: job.allow_write,
5951 sessions: job.sessions,
5952 artifacts: &job.artifacts,
5953 stem: &job.stem,
5954 run: &run,
5955 node: &node,
5956 cache_dir: cache.as_deref(),
5957 attachments: &[],
5958 },
5959 )
5960 .await;
5961 let out = match out {
5962 Ok(o) if o.usable() => AgentOutcome::Ok(o),
5963 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
5964 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
5972 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
5973 Ok(o) => AgentOutcome::Failed(format!(
5974 "exited with {:?} and no usable output",
5975 o.exit_code
5976 )),
5977 Err(e) => AgentOutcome::Failed(e.to_string()),
5978 };
5979 (i, seat, out)
5980 });
5981 }
5982 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
5983 while let Some(joined) = set.join_next().await {
5984 let (i, seat, out) = match joined {
5985 Ok(v) => v,
5986 Err(e) => {
5990 tracing::error!("agent task panicked: {e}");
5991 continue;
5992 }
5993 };
5994 state.seat_finished(&seat.key);
5995 record_jobs(state, node, round, &seat.key, &out);
5996 if let Err(e) = state.save() {
5997 tracing::warn!("could not persist a seat's completion: {e:#}");
5998 }
5999 if collected.len() <= i {
6000 collected.resize_with(i + 1, || None);
6001 }
6002 collected[i] = Some((i, seat, out));
6003 }
6004 if state
6010 .active
6011 .values()
6012 .any(|a| a.node == node && a.attempt == attempt)
6013 {
6014 state
6015 .active
6016 .retain(|_, a| !(a.node == node && a.attempt == attempt));
6017 if let Err(e) = state.save() {
6018 tracing::warn!("could not persist the end of a wave: {e:#}");
6019 }
6020 }
6021 if let Some(cache_dir) = cache
6028 && jobs_had_a_writer
6029 {
6030 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
6031 }
6032 if let Some(guard) = cache_guard {
6033 guard.release();
6034 }
6035 collected.into_iter().flatten().collect()
6036}
6037
6038fn record_jobs(
6049 state: &mut RunState,
6050 node: &str,
6051 round: Option<usize>,
6052 seat: &str,
6053 out: &AgentOutcome,
6054) {
6055 let commands: &[agent::CommandEvidence] = match out {
6056 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
6057 AgentOutcome::Failed(_) => &[],
6058 };
6059 let checked_at = Timestamp::now();
6060 for c in commands {
6061 state.jobs.push(JobRecord {
6062 node: node.to_owned(),
6063 round,
6064 seat: seat.to_owned(),
6065 id: c.id.clone(),
6066 description: c.description.clone(),
6067 checked_at,
6068 status: match c.exit_code {
6069 Some(0) => JobStatus::Completed,
6070 Some(_) => JobStatus::Failed,
6071 None => JobStatus::Unknown,
6072 },
6073 exit_code: c.exit_code,
6074 result_summary: c.result_summary.clone(),
6075 source: c.source.clone(),
6076 });
6077 }
6078}
6079
6080fn round_is_clean(
6101 blocking: usize,
6102 e2e_ok: bool,
6103 answered: usize,
6104 expected: usize,
6105 quota_missing: usize,
6106 policy: IncompleteReviewPolicy,
6107) -> bool {
6108 if blocking != 0 || !e2e_ok {
6109 return false;
6110 }
6111 if answered == expected || policy == IncompleteReviewPolicy::Warn {
6112 return true;
6113 }
6114 answered > 0 && expected - answered <= quota_missing
6115}
6116
6117fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
6141 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
6142 return Some(RunStatus::Gating);
6143 }
6144 let last = reviews.last()?;
6145 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
6146 if reviews.len() < max_rounds && !stagnant {
6147 return None;
6148 }
6149 if last.incomplete() && last.blocking == 0 {
6150 return Some(RunStatus::Blocked);
6151 }
6152 if last.e2e_status() == E2eStatus::ResourceBlocked {
6153 return None;
6154 }
6155 Some(if last.e2e.iter().all(CommandOutcome::ok) {
6156 RunStatus::Gating
6157 } else {
6158 RunStatus::Blocked
6159 })
6160}
6161
6162fn retry_budget(full: Duration, nudged: bool) -> Duration {
6177 if nudged {
6178 (full / 4).max(Duration::from_secs(120)).min(full)
6179 } else {
6180 full
6181 }
6182}
6183
6184#[allow(clippy::too_many_arguments)]
6197async fn ask_json_wave<T>(
6198 jobs: Vec<SeatJob>,
6199 sem: Arc<Semaphore>,
6200 retries: usize,
6201 ctx: &WaveCtx<'_>,
6202 losses: &mut Vec<QuotaLoss>,
6203 state: &mut RunState,
6204 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
6205) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
6206where
6207 T: serde::de::DeserializeOwned + Send + 'static,
6208{
6209 let n = jobs.len();
6210 let originals: Vec<SeatJob> = jobs;
6211 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
6212 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
6213 let mut attempts_used: Vec<usize> = vec![0; n];
6220 let mut pending: Vec<usize> = (0..n).collect();
6221
6222 for attempt in 0..=retries {
6223 if pending.is_empty() {
6224 break;
6225 }
6226 let mut batch = Vec::with_capacity(pending.len());
6227 for &i in &pending {
6228 let src = &originals[i];
6229 let (prompt, timeout) = if attempt == 0 {
6232 (src.prompt.clone(), src.timeout)
6233 } else {
6234 let why = done[i]
6235 .as_ref()
6236 .and_then(|r| r.as_ref().err().map(ToString::to_string))
6237 .unwrap_or_else(|| "no parsable answer".to_owned());
6238 let nudge = prompt::nudge(&why);
6239 let nudged = has_context(&src.spec, &seats[i], src.sessions);
6240 let prompt = if nudged {
6241 nudge
6242 } else {
6243 format!("{}\n\n---\n\n{}", src.prompt, nudge)
6244 };
6245 (prompt, retry_budget(src.timeout, nudged))
6246 };
6247 batch.push(SeatJob {
6248 spec: src.spec.clone(),
6249 seat: seats[i].clone(),
6250 cwd: src.cwd.clone(),
6251 prompt,
6252 timeout,
6253 allow_write: src.allow_write,
6254 sessions: src.sessions,
6255 artifacts: src.artifacts.clone(),
6256 stem: if attempt == 0 {
6257 src.stem.clone()
6258 } else {
6259 format!("{}-retry{attempt}", src.stem)
6260 },
6261 });
6262 }
6263
6264 if attempt > 0 {
6265 let seats_out: Vec<&str> = pending
6266 .iter()
6267 .map(|&i| originals[i].seat.key.as_str())
6268 .collect();
6269 state.event(
6270 ctx.node,
6271 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
6272 );
6273 }
6274 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
6275 let mut still = Vec::new();
6276 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
6277 seats[i] = seat;
6278 let (parsed, quota) = match out {
6279 AgentOutcome::Ok(o) => (
6280 match verdict::extract_json::<T>(&o.text) {
6281 Ok(v) => match validate(&v) {
6282 Ok(()) => Ok((v, o)),
6283 Err(e) => Err(e),
6284 },
6285 Err(e) => Err(e),
6286 },
6287 false,
6288 ),
6289 AgentOutcome::Quota(o) => {
6290 losses.push(QuotaLoss {
6291 seat: originals[i].seat.key.clone(),
6292 node: ctx.node.to_owned(),
6293 at: Timestamp::now(),
6294 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
6295 });
6296 (
6297 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
6298 true,
6299 )
6300 }
6301 AgentOutcome::Dropped(o) => {
6306 let why = o
6307 .dropped
6308 .as_ref()
6309 .map(|d| d.why.as_str())
6310 .unwrap_or("the CLI ended the stream without delivering its answer");
6311 (
6312 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
6313 false,
6314 )
6315 }
6316 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
6317 };
6318 let failed = parsed.is_err();
6319 done[i] = Some(parsed);
6320 attempts_used[i] = attempt;
6321 if failed && !quota {
6324 still.push(i);
6325 }
6326 }
6327 pending = still;
6328 }
6329
6330 seats
6331 .into_iter()
6332 .zip(done)
6333 .zip(attempts_used)
6334 .map(|((seat, res), attempts)| {
6335 (
6336 seat,
6337 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6338 attempts,
6339 )
6340 })
6341 .collect()
6342}
6343
6344async fn acquire_cache_lease(
6357 state: &mut RunState,
6358 cache_dir: &Path,
6359 owner: &crate::cache::Owner,
6360 budget: Duration,
6361 context: &str,
6362) -> Result<crate::cache::Guard> {
6363 let home = crate::run::home();
6364 let started = Instant::now();
6365 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6366 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6367 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6368 Err(e) => {
6369 state.event(
6370 "verify",
6371 format!("{context}: could not check the shared build cache: {e:#}"),
6372 );
6373 if let Err(e2) = state.save() {
6374 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6375 }
6376 return Err(e);
6377 }
6378 };
6379 state.event(
6380 "verify",
6381 format!(
6382 "{context}: waiting for the shared build cache at {} ({})",
6383 cache_dir.display(),
6384 busy.describe()
6385 ),
6386 );
6387 if let Err(e) = state.save() {
6388 tracing::warn!("could not persist a cache wait: {e:#}");
6389 }
6390 let remaining = budget.saturating_sub(started.elapsed());
6391 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6392 Ok(g) => Ok(g),
6393 Err(e) => {
6394 state.event("verify", format!("{context}: {e:#}"));
6395 if let Err(e2) = state.save() {
6396 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6397 }
6398 Err(e)
6399 }
6400 }
6401}
6402
6403#[allow(clippy::too_many_arguments)]
6424async fn with_cache_lease<'s, F, Fut>(
6425 state: &'s mut RunState,
6426 cache_dir: Option<&Path>,
6427 node: &str,
6428 seat: &str,
6429 worktree: &Path,
6430 head: &str,
6431 budget: Duration,
6432 context: &str,
6433 body: F,
6434) -> (Vec<CommandOutcome>, bool)
6435where
6436 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6437 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6438{
6439 let Some(cache_dir) = cache_dir else {
6440 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6441 return (outcomes, retried);
6442 };
6443 let home = crate::run::home();
6444 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6445 let started = Instant::now();
6446 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6447 Ok(g) => g,
6448 Err(e) => {
6449 return (
6450 vec![CommandOutcome {
6451 command: "(waiting for the shared build cache)".to_owned(),
6452 code: None,
6453 output_tail: e.to_string(),
6454 duration_ms: started.elapsed().as_millis() as u64,
6455 resource_blocked: true,
6456 }],
6457 false,
6458 );
6459 }
6460 };
6461 let identity = crate::cache::Identity::new(worktree, head);
6462 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6463 state.event(
6471 "verify",
6472 format!(
6473 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6474 worktree.display(),
6475 short(head)
6476 ),
6477 );
6478 guard.release();
6479 return (
6480 vec![CommandOutcome {
6481 command: "(confirming the shared build cache is fresh)".to_owned(),
6482 code: None,
6483 output_tail: e.to_string(),
6484 duration_ms: started.elapsed().as_millis() as u64,
6485 resource_blocked: true,
6486 }],
6487 false,
6488 );
6489 }
6490 let remaining = budget.saturating_sub(started.elapsed());
6491 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6492 if !timed_out_pids.is_empty() {
6497 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6498 }
6499 guard.release();
6500 (outcomes, retried)
6501}
6502
6503async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6515 wait_for_pids_with(
6516 pids,
6517 crate::proc::pid_alive,
6518 LEASE_RELEASE_POLL,
6519 LEASE_RELEASE_MAX_WAIT,
6520 )
6521 .await;
6522}
6523
6524async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6530 pids: &[u32],
6531 alive: F,
6532 poll: Duration,
6533 max_wait: Duration,
6534) {
6535 let deadline = Instant::now() + max_wait;
6536 loop {
6537 if pids.iter().all(|&pid| !alive(pid)) {
6538 return;
6539 }
6540 if Instant::now() >= deadline {
6541 return;
6542 }
6543 tokio::time::sleep(poll).await;
6544 }
6545}
6546
6547fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6554 outcomes.iter().any(|o| o.resource_blocked)
6555}
6556
6557enum GateFix {
6559 Retry,
6561 Stop,
6564 Defer,
6567}
6568
6569fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6577 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6578 red.peek().is_some()
6579 && red.all(|o| {
6580 !o.resource_blocked
6581 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6582 && !o.output_tail.trim().is_empty()
6583 })
6584}
6585
6586fn e2e_outcome_label(o: &CommandOutcome) -> String {
6590 if o.ok() {
6591 return "pass".to_owned();
6592 }
6593 let reason = if o.build_failed() {
6594 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6595 } else {
6596 format!("FAIL ({:?})", o.code)
6597 };
6598 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6599}
6600
6601async fn run_e2e_with_retry(
6609 state: &mut RunState,
6610 shell: &[String],
6611 commands: &[String],
6612 worktree: &Path,
6613 timeout: Duration,
6614 context: &str,
6615) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6616 let (mut e2e, mut timed_out_pids) = run_commands(
6617 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6618 )
6619 .await;
6620 for o in &e2e {
6621 state.event(
6622 "verify",
6623 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6624 );
6625 }
6626 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6630 if verify_retried {
6631 state.event(
6632 "verify",
6633 format!(
6634 "{context}: verify could not build/link, not a test result — retrying once \
6635 before concluding"
6636 ),
6637 );
6638 let retried = run_commands(
6639 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6640 )
6641 .await;
6642 e2e = retried.0;
6643 timed_out_pids.extend(retried.1);
6646 for o in &e2e {
6647 state.event(
6648 "verify",
6649 format!(
6650 "{context}: retry `{}` -> {}",
6651 o.command,
6652 e2e_outcome_label(o)
6653 ),
6654 );
6655 }
6656 }
6657 (e2e, verify_retried, timed_out_pids)
6658}
6659
6660#[allow(clippy::too_many_arguments)]
6676async fn run_commands(
6677 state: &mut RunState,
6678 node: &str,
6679 task: &str,
6680 attempt: usize,
6681 shell: &[String],
6682 commands: &[String],
6683 cwd: &Path,
6684 timeout: Duration,
6685) -> (Vec<CommandOutcome>, Vec<u32>) {
6686 if commands.is_empty() {
6687 return (Vec::new(), Vec::new());
6692 }
6693 let mut out = Vec::new();
6694 let mut timed_out_pids = Vec::new();
6695 let total = commands.len();
6696 for (idx, command) in commands.iter().enumerate() {
6697 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6698 if let Err(e) = state.save() {
6699 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6700 }
6701 let started = Instant::now();
6702 let mut cmd = tokio::process::Command::new(&shell[0]);
6703 cmd.quiet();
6704 cmd.args(&shell[1..])
6705 .arg(command)
6706 .current_dir(cwd)
6707 .stdin(std::process::Stdio::null())
6708 .stdout(std::process::Stdio::piped())
6709 .stderr(std::process::Stdio::piped())
6710 .kill_on_drop(true);
6711 let spawned = cmd.spawn();
6712 let (code, body) = match spawned {
6713 Ok(child) => {
6714 let pid = child.id();
6719 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6720 Ok(Ok(o)) => {
6721 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6722 body.push_str(&String::from_utf8_lossy(&o.stderr));
6723 (o.status.code(), body)
6724 }
6725 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6726 Err(_) => {
6727 if let Some(pid) = pid {
6728 timed_out_pids.push(pid);
6729 }
6730 (None, format!("timed out after {}s", timeout.as_secs()))
6731 }
6732 }
6733 }
6734 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6735 };
6736 out.push(CommandOutcome {
6737 command: command.clone(),
6738 code,
6739 output_tail: tail(&body, OUTPUT_TAIL),
6740 duration_ms: started.elapsed().as_millis() as u64,
6741 resource_blocked: false,
6742 });
6743 }
6744 state.task_finished(task);
6745 if let Err(e) = state.save() {
6746 tracing::warn!("could not persist the end of {task}: {e:#}");
6747 }
6748 (out, timed_out_pids)
6749}
6750
6751fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6763 let repo = repo.display();
6764 match style {
6765 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6766 MergeStyle::Squash => {
6767 let subject = message
6770 .lines()
6771 .next()
6772 .unwrap_or(branch)
6773 .replace(['\\', '"', '$', '`'], "");
6774 format!(
6775 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6776 )
6777 }
6778 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6779 }
6780}
6781
6782const PR_TITLE_MAX: usize = 240;
6793
6794struct PrMessage {
6799 title: String,
6800 body: String,
6801}
6802
6803impl PrMessage {
6804 fn commit_message(&self) -> String {
6808 format!("{}\n\n{}", self.title, self.body)
6809 }
6810}
6811
6812fn title_marker(line: &str) -> Option<&str> {
6814 let line = line.trim();
6815 let head = line.get(..6)?;
6816 head.eq_ignore_ascii_case("title:")
6817 .then(|| line[6..].trim())
6818}
6819
6820fn summary_title(summary: &str) -> Option<String> {
6825 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6826 let raw = title_marker(first)?;
6827 if raw.is_empty() {
6828 return None;
6829 }
6830 let title = queue::title_from(raw, PR_TITLE_MAX);
6831 let lower = title.to_ascii_lowercase();
6832 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6833 return None;
6834 }
6835 Some(title)
6836}
6837
6838fn summary_without_title(summary: &str) -> String {
6841 let mut lines = summary.trim().lines().peekable();
6842 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6843 lines.next();
6844 }
6845 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6846}
6847
6848fn pr_message(state: &RunState, winner: char) -> PrMessage {
6862 let summary = state
6863 .candidates
6864 .iter()
6865 .find(|c| c.label == winner)
6866 .map(|c| c.summary.as_str())
6867 .unwrap_or_default();
6868 let title = summary_title(summary).unwrap_or_else(|| {
6871 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6872 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6873 t
6874 } else {
6875 format!(
6876 "chore: land candidate {} of run {}",
6877 winner.to_ascii_uppercase(),
6878 state.id
6879 )
6880 }
6881 });
6882
6883 let mut body = String::new();
6884 let what = summary_without_title(summary);
6885 if !what.is_empty() {
6886 body.push_str("## Summary\n\n");
6887 body.push_str(&what);
6888 body.push_str("\n\n");
6889 }
6890
6891 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6892 if let Some(fix) = fix
6893 && !fix.notes.trim().is_empty()
6894 {
6895 body.push_str("## Review fixes\n\n");
6896 body.push_str(fix.notes.trim());
6897 body.push_str("\n\n");
6898 }
6899
6900 let open = state.open_findings();
6901 if !open.is_empty() {
6902 body.push_str("## Open review findings\n\n");
6903 for f in &open {
6904 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6905 }
6906 body.push('\n');
6907 }
6908
6909 if let Some(fix) = fix
6910 && !fix.rejected.is_empty()
6911 {
6912 body.push_str("## Declined by the fixer\n\n");
6913 for r in &fix.rejected {
6914 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
6915 }
6916 body.push('\n');
6917 }
6918
6919 let task = state.instruction.trim();
6920 let task = if task.is_empty() {
6921 "(empty task)"
6922 } else {
6923 task
6924 };
6925 body.push_str(&format!(
6926 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
6927 task.replace("</details>", "</details>")
6928 ));
6929
6930 body.push_str(&format!(
6931 "\n---\nmagi:run/{} magi:candidate-{}\n",
6932 state.id,
6933 winner.to_ascii_lowercase()
6934 ));
6935
6936 let id = crate::scrub::Identity::current();
6939 PrMessage {
6940 title: crate::scrub::scrub(&title, &id),
6941 body: crate::scrub::scrub(&body, &id),
6942 }
6943}
6944
6945fn seeded_instruction(state: &RunState) -> String {
6949 match refs::describe(&state.seeds) {
6950 Some(facts) => format!(
6951 "{}\n\n# Existing work the task refers to\n\n{facts}\n\n\
6952 Candidates start from the unmerged branch named above, when there \
6953 is one, and carry any unmerged commit named by sha as a \
6954 cherry-pick. Check that this is what the task meant before \
6955 building on it.",
6956 state.instruction
6957 ),
6958 None => state.instruction.clone(),
6959 }
6960}
6961
6962async fn merge_is_empty(repo: &Path, state: &RunState, branch: &str, mode: MergeMode) -> bool {
6966 let base = &state.base_branch;
6967 let mut against = base.clone();
6968 if mode == MergeMode::Pr {
6969 let remote = &state.config.merge.remote;
6970 let tracking = format!("{remote}/{base}");
6971 let fetched = git::fetch(repo, remote, base).await;
6972 if fetched.is_ok_and(|o| o.ok()) && git::rev_exists(repo, &tracking).await {
6973 against = tracking;
6974 }
6975 }
6976 matches!(git::commits_ahead(repo, &against, branch).await, Ok(0))
6977}
6978
6979fn empty_candidate_detail(state: &RunState, base: &str) -> String {
6982 let mut detail = format!(
6983 "empty candidate: the winning branch has 0 commits ahead of {base}, so there is \
6984 nothing to open a pull request for"
6985 );
6986 match refs::describe(&state.seeds) {
6987 Some(facts) => detail.push_str(&format!("\nReferences in the task:\n{facts}")),
6988 None => detail.push_str(
6989 "\nThe task names no existing branch or commit; if it means to land work \
6990 that lives elsewhere, name the branch (magi/<run>/<label>) or the sha.",
6991 ),
6992 }
6993 detail
6994}
6995
6996async fn gh_pr_create(
6998 cwd: &Path,
6999 base: &str,
7000 head: &str,
7001 title: &str,
7002 body: &str,
7003) -> Result<String> {
7004 let out = tokio::process::Command::new("gh")
7005 .args([
7006 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
7007 ])
7008 .current_dir(cwd)
7009 .quiet()
7010 .stdin(std::process::Stdio::null())
7011 .output()
7012 .await
7013 .context("spawn gh")?;
7014 if out.status.success() {
7015 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
7016 } else {
7017 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
7018 }
7019}
7020
7021pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
7030 let repo = state.repo.clone();
7031 let root = state.worktree_root();
7032 let winner = state.tally.as_ref().map(|t| t.winner);
7033 let mut removed = Vec::new();
7034
7035 for i in 0..state.candidates.len() {
7036 let c = state.candidates[i].clone();
7037 let is_winner = Some(c.label) == winner;
7038 if is_winner && !drop_winner {
7039 continue;
7040 }
7041 if c.worktree.exists() {
7042 git::worktree_remove(&repo, &c.worktree).await.ok();
7043 removed.push(c.worktree.to_string_lossy().into_owned());
7044 }
7045 let handed_over = state.released_branches.contains(&c.branch);
7048 if !handed_over && git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
7049 git::branch_delete(&repo, &c.branch).await.ok();
7050 removed.push(c.branch.clone());
7051 }
7052 state.candidates[i].folded = true;
7053 }
7054
7055 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
7056 let path = name.path();
7057 let keep = !drop_winner
7058 && winner.is_some_and(|w| {
7059 path.file_name()
7060 .is_some_and(|n| n == format!("cand-{w}").as_str())
7061 });
7062 if keep {
7063 continue;
7064 }
7065 git::worktree_remove(&repo, &path).await.ok();
7066 removed.push(path.to_string_lossy().into_owned());
7067 }
7068
7069 remove_if_empty(&root);
7078
7079 if state.enabled_worktree_config && drop_winner {
7080 git::release_worktree_config(&repo).await.ok();
7084 state.enabled_worktree_config = false;
7085 }
7086 state.save_under(home)?;
7087 Ok(removed)
7088}
7089
7090fn remove_if_empty(dir: &Path) {
7101 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
7102 std::fs::remove_dir(dir).ok();
7103 }
7104}
7105
7106pub fn worst_open(state: &RunState) -> Option<Severity> {
7108 state
7109 .reviews
7110 .last()?
7111 .reviews
7112 .iter()
7113 .flat_map(|r| r.findings.iter())
7114 .map(|f| f.severity)
7115 .max()
7116}
7117
7118#[cfg(test)]
7119mod tests {
7120 use super::*;
7121 use crate::run::GateStatus;
7122 use std::collections::BTreeMap;
7123 use std::time::Duration;
7124
7125 fn conductor() -> AgentSpec {
7126 AgentSpec {
7127 id: "conductor".to_owned(),
7128 kind: crate::config::AgentKind::Command,
7129 model: None,
7130 command: vec!["true".to_owned()],
7131 extra_args: Vec::new(),
7132 env: BTreeMap::new(),
7133 prompt_delivery: None,
7134 }
7135 }
7136
7137 fn spec(id: &str) -> AgentSpec {
7138 AgentSpec {
7139 id: id.to_owned(),
7140 kind: crate::config::AgentKind::Command,
7141 model: None,
7142 command: vec!["true".to_owned()],
7143 extra_args: Vec::new(),
7144 env: BTreeMap::new(),
7145 prompt_delivery: None,
7146 }
7147 }
7148
7149 #[test]
7156 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
7157 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7158 let tried = BTreeSet::from(["beta".to_owned()]);
7159 let next = next_untried_implementer(&roster, 1, &tried);
7162 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
7163 }
7164
7165 #[test]
7166 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
7167 let roster = vec![spec("alpha"), spec("beta")];
7168 let tried = BTreeSet::from(["beta".to_owned()]);
7169 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7173 }
7174
7175 #[test]
7176 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
7177 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7178 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
7179 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7183 }
7184
7185 #[test]
7186 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
7187 let roster = vec![spec("a"), spec("a"), spec("b")];
7188 let tried = BTreeSet::from(["a".to_owned()]);
7189 let next = next_untried_implementer(&roster, 0, &tried);
7190 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
7191 }
7192
7193 #[test]
7194 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
7195 let roster = vec![spec("a"), spec("b")];
7196 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
7197 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
7198 }
7199
7200 #[test]
7201 fn remove_if_empty_only_ever_takes_a_bare_directory() {
7202 let dir = tempfile::tempdir().unwrap();
7203 let bay = dir.path().join("ffff");
7204
7205 remove_if_empty(&bay);
7207 assert!(!bay.exists());
7208
7209 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
7212 remove_if_empty(&bay);
7213 assert!(bay.exists(), "non-empty directory must survive");
7214
7215 std::fs::remove_dir(bay.join("cand-A")).unwrap();
7217 remove_if_empty(&bay);
7218 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
7219 }
7220
7221 #[test]
7230 fn a_full_panel_that_found_nothing_is_clean() {
7231 assert!(round_is_clean(
7232 0,
7233 true,
7234 2,
7235 2,
7236 0,
7237 IncompleteReviewPolicy::Block
7238 ));
7239 }
7240
7241 #[test]
7242 fn a_missing_seat_is_never_clean_under_the_default_policy() {
7243 assert!(!round_is_clean(
7244 0,
7245 true,
7246 1,
7247 2,
7248 0,
7249 IncompleteReviewPolicy::Block
7250 ));
7251 }
7252
7253 #[test]
7254 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
7255 assert!(!round_is_clean(
7256 1,
7257 true,
7258 1,
7259 2,
7260 0,
7261 IncompleteReviewPolicy::Warn
7262 ));
7263 }
7264
7265 #[test]
7266 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
7267 assert!(round_is_clean(
7268 0,
7269 true,
7270 1,
7271 2,
7272 0,
7273 IncompleteReviewPolicy::Warn
7274 ));
7275 }
7276
7277 #[test]
7278 fn a_full_panel_with_an_open_finding_is_not_clean() {
7279 assert!(!round_is_clean(
7280 1,
7281 true,
7282 2,
7283 2,
7284 0,
7285 IncompleteReviewPolicy::Block
7286 ));
7287 }
7288
7289 #[test]
7290 fn a_full_panel_with_a_red_e2e_is_not_clean() {
7291 assert!(!round_is_clean(
7292 0,
7293 false,
7294 2,
7295 2,
7296 0,
7297 IncompleteReviewPolicy::Block
7298 ));
7299 }
7300
7301 #[test]
7308 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
7309 assert!(round_is_clean(
7312 0,
7313 true,
7314 1,
7315 2,
7316 1,
7317 IncompleteReviewPolicy::Block
7318 ));
7319 }
7320
7321 #[test]
7322 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
7323 assert!(!round_is_clean(
7326 0,
7327 true,
7328 1,
7329 2,
7330 0,
7331 IncompleteReviewPolicy::Block
7332 ));
7333 }
7334
7335 #[test]
7336 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
7337 assert!(!round_is_clean(
7338 1,
7339 true,
7340 1,
7341 2,
7342 1,
7343 IncompleteReviewPolicy::Block
7344 ));
7345 assert!(!round_is_clean(
7346 0,
7347 false,
7348 1,
7349 2,
7350 1,
7351 IncompleteReviewPolicy::Block
7352 ));
7353 }
7354
7355 #[test]
7356 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
7357 assert!(!round_is_clean(
7361 0,
7362 true,
7363 0,
7364 2,
7365 2,
7366 IncompleteReviewPolicy::Block
7367 ));
7368 }
7369
7370 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
7371 CommandOutcome {
7372 command: "test".to_owned(),
7373 code,
7374 output_tail: String::new(),
7375 duration_ms: 0,
7376 resource_blocked,
7377 }
7378 }
7379
7380 #[test]
7381 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7382 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7383 assert!(
7384 !verify_inconclusive(&[outcome(Some(1), false)]),
7385 "an ordinary failure is still evidence about the patch"
7386 );
7387 assert!(verify_inconclusive(&[outcome(None, true)]));
7388 assert!(
7389 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7390 "one inconclusive outcome taints the whole batch"
7391 );
7392 assert!(!verify_inconclusive(&[]));
7393 }
7394
7395 #[tokio::test]
7396 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7397 let calls = std::sync::atomic::AtomicUsize::new(0);
7401 let started = Instant::now();
7402 wait_for_pids_with(
7403 &[123],
7404 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7405 Duration::from_millis(5),
7406 Duration::from_secs(5),
7407 )
7408 .await;
7409 assert!(
7410 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7411 "must keep checking rather than deciding on the first answer"
7412 );
7413 assert!(
7414 started.elapsed() < Duration::from_secs(1),
7415 "must return the moment it is confirmed dead, not wait out the ceiling"
7416 );
7417 }
7418
7419 #[tokio::test]
7420 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7421 let started = Instant::now();
7422 wait_for_pids_with(
7423 &[123],
7424 |_| true, Duration::from_millis(5),
7426 Duration::from_millis(30),
7427 )
7428 .await;
7429 let elapsed = started.elapsed();
7430 assert!(
7431 elapsed >= Duration::from_millis(30),
7432 "must not give up before its own ceiling: {elapsed:?}"
7433 );
7434 assert!(
7435 elapsed < Duration::from_secs(1),
7436 "must not wait past its own ceiling either: {elapsed:?}"
7437 );
7438 }
7439
7440 #[tokio::test]
7441 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7442 let started = Instant::now();
7443 wait_for_pids_with(
7444 &[],
7445 |_| true,
7446 Duration::from_secs(5),
7447 Duration::from_secs(5),
7448 )
7449 .await;
7450 assert!(
7451 started.elapsed() < Duration::from_millis(200),
7452 "an empty pid list has nothing to confirm"
7453 );
7454 }
7455
7456 fn review_round(
7462 clean: bool,
7463 blocking: usize,
7464 answered: usize,
7465 expected: usize,
7466 progressed: bool,
7467 e2e_ok: bool,
7468 ) -> ReviewRound {
7469 ReviewRound {
7470 round: 1,
7471 head: "h".to_owned(),
7472 verified_head: None,
7473 verified_at: None,
7474 reviews: Vec::new(),
7475 e2e: vec![CommandOutcome {
7476 command: "test".to_owned(),
7477 code: Some(if e2e_ok { 0 } else { 1 }),
7478 output_tail: String::new(),
7479 duration_ms: 0,
7480 resource_blocked: false,
7481 }],
7482 verify_retried: false,
7483 e2e_deferred: false,
7484 e2e_defer_reason: None,
7485 fix: None,
7486 blocking,
7487 answered,
7488 expected,
7489 clean,
7490 progressed,
7491 vote_split: false,
7492 reconsideration: Vec::new(),
7493 verdict: None,
7494 }
7495 }
7496
7497 #[test]
7498 fn review_conclusion_is_none_when_nothing_has_run() {
7499 assert_eq!(review_conclusion(&[], 3), None);
7500 }
7501
7502 #[test]
7503 fn review_conclusion_is_none_while_rounds_remain() {
7504 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7505 assert_eq!(review_conclusion(&rounds, 3), None);
7506 }
7507
7508 #[test]
7509 fn review_conclusion_is_gating_once_a_round_is_clean() {
7510 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7511 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7512 }
7513
7514 #[test]
7515 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7516 let rounds = vec![
7517 review_round(false, 1, 2, 2, true, true),
7518 review_round(false, 1, 2, 2, true, true),
7519 ];
7520 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7521 }
7522
7523 #[test]
7524 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7525 let rounds = vec![
7526 review_round(false, 1, 2, 2, true, true),
7527 review_round(false, 1, 2, 2, true, false),
7528 ];
7529 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7530 }
7531
7532 #[test]
7533 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7534 let mut blocked = review_round(false, 1, 2, 2, true, false);
7541 blocked.e2e[0].resource_blocked = true;
7542 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7543 assert_eq!(review_conclusion(&rounds, 2), None);
7544 }
7545
7546 #[test]
7547 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7548 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7550 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7551 }
7552
7553 #[test]
7554 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7555 let rounds = vec![
7556 review_round(false, 1, 2, 2, false, true),
7557 review_round(false, 1, 2, 2, false, true),
7558 ];
7559 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7560 }
7561
7562 fn secs(n: u64) -> Duration {
7563 Duration::from_secs(n)
7564 }
7565
7566 fn init_repo(dir: &Path) {
7569 let run = |args: &[&str]| {
7570 let out = std::process::Command::new("git")
7571 .args(args)
7572 .current_dir(dir)
7573 .quiet()
7574 .output()
7575 .expect("spawn git");
7576 assert!(
7577 out.status.success(),
7578 "git {args:?} failed: {}",
7579 String::from_utf8_lossy(&out.stderr)
7580 );
7581 };
7582 run(&["init", "-b", "main"]);
7583 run(&["config", "user.name", "magi test"]);
7584 run(&["config", "user.email", "magi@example.com"]);
7585 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7586 run(&["add", "-A"]);
7587 run(&["commit", "-m", "init"]);
7588 }
7589
7590 fn ask_test_home() {
7598 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7599 }
7600
7601 fn runner_at(status: RunStatus) -> Runner {
7604 let mut state = RunState::new(
7605 PathBuf::from("/nonexistent/repo"),
7606 "main".to_owned(),
7607 "deadbeef".to_owned(),
7608 "task".to_owned(),
7609 Config::default(),
7610 );
7611 state.status = status;
7612 Runner {
7613 state,
7614 roles: ResolvedRoles {
7615 implementers: Vec::new(),
7616 judges: Vec::new(),
7617 reviewers: Vec::new(),
7618 fixer: None,
7619 conductor: conductor(),
7620 implementer_roster: Vec::new(),
7621 },
7622 sem: Arc::new(Semaphore::new(1)),
7623 pause: Pause::new(),
7624 interrupt: Pause::new(),
7625 }
7626 }
7627
7628 #[test]
7632 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7633 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7634 let mut runner = runner_at(RunStatus::Implementing);
7635 let interrupt = Pause::new();
7636 runner.watch_interrupt(interrupt.clone());
7637
7638 interrupt.park_because("task a1b2 asked to run first");
7639
7640 assert!(runner.park_here().expect("park_here"));
7641 assert!(runner.state.parked);
7642 let last = runner.state.events.last().expect("a park event");
7643 assert_eq!(last.node, "park");
7644 assert!(
7645 last.message.contains("task a1b2 asked to run first"),
7646 "expected the interrupt reason in {:?}",
7647 last.message
7648 );
7649 }
7650
7651 #[test]
7659 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7660 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7661 let mut runner = runner_at(RunStatus::Implementing);
7662 let shutdown = Pause::new();
7663 runner.on_pause(shutdown.clone());
7664 let interrupt = Pause::new();
7665 runner.watch_interrupt(interrupt.clone());
7666
7667 assert!(!runner.park_here().expect("park_here"));
7669 assert!(!runner.state.parked);
7670
7671 interrupt.park_because("test");
7673 assert!(!shutdown.parked());
7674 assert!(runner.park_here().expect("park_here"));
7675 }
7676
7677 #[tokio::test]
7691 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7692 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7693 let mut runner = runner_at(RunStatus::Implementing);
7694 let interrupt = Pause::new();
7695 runner.watch_interrupt(interrupt.clone());
7696
7697 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7698 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7699
7700 let node = async move {
7704 started_tx.send(()).expect("send started");
7705 finish_rx.await.expect("recv finish");
7706 "node finished"
7707 };
7708
7709 let interrupter = async move {
7710 started_rx.await.expect("recv started");
7711 interrupt.park_because("higher-priority task waiting");
7713 tokio::task::yield_now().await;
7717 finish_tx.send(()).expect("send finish");
7718 };
7719
7720 let (node_result, ()) = tokio::join!(node, interrupter);
7721 assert_eq!(
7722 node_result, "node finished",
7723 "the in-flight call ran to completion"
7724 );
7725
7726 assert!(runner.park_here().expect("park_here"));
7729 assert!(runner.state.parked);
7730 }
7731
7732 #[test]
7738 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7739 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7740 let mut runner = runner_at(RunStatus::Judging);
7741 runner.state.config.agents = vec![conductor()];
7745 runner.state.candidates = vec![Candidate {
7746 index: 0,
7747 label: 'A',
7748 agent: "alpha".to_owned(),
7749 branch: "magi/x/A".to_owned(),
7750 worktree: PathBuf::from("/nonexistent/worktree"),
7751 summary: "did the thing".to_owned(),
7752 stat: "1 file changed".to_owned(),
7753 files: 1,
7754 commits: 1,
7755 empty: false,
7756 failed: None,
7757 verified_noop: None,
7758 duration_ms: 1234,
7759 folded: false,
7760 }];
7761 let run_id = runner.state.id.clone();
7762
7763 let interrupt = Pause::new();
7764 runner.watch_interrupt(interrupt.clone());
7765 interrupt.park_because("task c3d4 asked to run first");
7766 assert!(runner.park_here().expect("park_here"));
7767
7768 let resumed = Runner::resume(&run_id).expect("resume");
7769 assert_eq!(resumed.state.candidates.len(), 1);
7770 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7771 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7772 assert_eq!(resumed.state.status, runner.state.status);
7773 assert!(
7774 resumed.state.parked,
7775 "still parked until `execute` actually walks the graph again"
7776 );
7777 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7778 }
7779
7780 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7782 let mut q = ask::Question::new(
7783 run.to_owned(),
7784 "implement".to_owned(),
7785 "impl-A".to_owned(),
7786 "Which storage backend should the cache use?".to_owned(),
7787 String::new(),
7788 vec!["SQLite".to_owned(), "Redis".to_owned()],
7789 );
7790 store.put(&mut q).unwrap();
7791 q
7792 }
7793
7794 #[test]
7795 fn a_failed_runs_open_question_is_abandoned() {
7796 ask_test_home();
7797 let store = ask::Questions::open();
7798 let mut runner = runner_at(RunStatus::Failed);
7799 let run = runner.state.id.clone();
7800 let q = ask_open_question(&store, &run);
7801
7802 runner.settle_questions();
7803
7804 let back = store.get(&q.id).unwrap();
7805 assert!(
7806 !back.status.open(),
7807 "the seat that asked died with the run; nobody is left to read an answer"
7808 );
7809 assert!(
7810 back.detail.contains(&run) && back.detail.contains("failed"),
7811 "the reason names what the run became, not just that it is gone: {}",
7812 back.detail
7813 );
7814 }
7815
7816 #[test]
7817 fn a_merged_runs_open_question_is_abandoned_too() {
7818 ask_test_home();
7819 let store = ask::Questions::open();
7820 for status in [RunStatus::Merged, RunStatus::Ready] {
7823 let mut runner = runner_at(status);
7824 let run = runner.state.id.clone();
7825 let q = ask_open_question(&store, &run);
7826
7827 runner.settle_questions();
7828
7829 let back = store.get(&q.id).unwrap();
7830 assert!(
7831 !back.status.open(),
7832 "{status:?} run's question must not outlive the run"
7833 );
7834 }
7835 }
7836
7837 #[test]
7838 fn a_still_resumable_runs_open_question_is_left_alone() {
7839 ask_test_home();
7840 let store = ask::Questions::open();
7841 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7847 let mut runner = runner_at(status);
7848 let run = runner.state.id.clone();
7849 let q = ask_open_question(&store, &run);
7850
7851 runner.settle_questions();
7852
7853 let back = store.get(&q.id).unwrap();
7854 assert!(
7855 back.status.open(),
7856 "{status:?} is still alive; the question must still be waiting"
7857 );
7858 }
7859 }
7860
7861 #[test]
7862 fn settle_questions_never_touches_an_already_answered_question() {
7863 ask_test_home();
7864 let store = ask::Questions::open();
7865 let mut runner = runner_at(RunStatus::Failed);
7866 let run = runner.state.id.clone();
7867 let mut q = ask_open_question(&store, &run);
7868 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
7869 .unwrap();
7870 store.put(&mut q).unwrap();
7871
7872 runner.settle_questions();
7877 runner.settle_questions();
7878
7879 let back = store.get(&q.id).unwrap();
7880 assert_eq!(
7881 back.status,
7882 ask::QuestionStatus::Answered,
7883 "a real answer is a decision on record, never overwritten by a sweep"
7884 );
7885 }
7886
7887 #[tokio::test]
7898 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
7899 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
7900 let tmp = tempfile::tempdir().expect("tempdir");
7901 let repo = tmp.path().join("repo");
7902 std::fs::create_dir_all(&repo).unwrap();
7903 init_repo(&repo);
7904
7905 let mut config = Config::default();
7906 config.graph.worktree_root = Some(tmp.path().join("wt"));
7907
7908 let mut state = RunState::new(
7909 repo.clone(),
7910 "main".to_owned(),
7911 "deadbeef".to_owned(),
7912 "task".to_owned(),
7913 config,
7914 );
7915 let root = state.worktree_root();
7916 let wt_a = root.join("cand-A");
7917 let wt_b = root.join("cand-B");
7918 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
7919 .await
7920 .expect("worktree A");
7921 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
7922 .await
7923 .expect("worktree B");
7924
7925 state.candidates = vec![
7926 Candidate {
7927 index: 0,
7928 label: 'A',
7929 agent: "alpha".to_owned(),
7930 branch: "magi/x/A".to_owned(),
7931 worktree: wt_a.clone(),
7932 summary: String::new(),
7933 stat: String::new(),
7934 files: 0,
7935 commits: 0,
7936 empty: false,
7937 failed: None,
7938 verified_noop: None,
7939 duration_ms: 0,
7940 folded: false,
7941 },
7942 Candidate {
7943 index: 1,
7944 label: 'B',
7945 agent: "beta".to_owned(),
7946 branch: "magi/x/B".to_owned(),
7947 worktree: wt_b.clone(),
7948 summary: String::new(),
7949 stat: String::new(),
7950 files: 0,
7951 commits: 0,
7952 empty: false,
7953 failed: None,
7954 verified_noop: None,
7955 duration_ms: 0,
7956 folded: false,
7957 },
7958 ];
7959 state.tally = Some(Tally {
7960 first_choice: BTreeMap::from([('A', 1)]),
7961 borda: BTreeMap::new(),
7962 winner: 'A',
7963 rankings: 1,
7964 unanimous_initial: true,
7965 deliberated: false,
7966 changed_votes: 0,
7967 unanimous_final: true,
7968 tie_break: None,
7969 judges: 1,
7970 present: 1,
7971 quorum: 1,
7972 met_quorum: true,
7973 uncontested: None,
7974 });
7975 state.status = RunStatus::Ready;
7976
7977 fold_run(&mut state, false, &crate::run::home())
7978 .await
7979 .expect("fold_run");
7980
7981 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
7982 assert!(
7983 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7984 "the unmerged winner's branch survives"
7985 );
7986 assert!(
7987 !state.candidates[0].folded,
7988 "the winner is not marked folded"
7989 );
7990
7991 assert!(!wt_b.exists(), "the loser's worktree is removed");
7992 assert!(
7993 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
7994 "the loser's branch is removed"
7995 );
7996 assert!(state.candidates[1].folded, "the loser is marked folded");
7997 }
7998
7999 #[tokio::test]
8002 async fn fold_run_keeps_a_branch_that_was_handed_to_a_later_run() {
8003 let tmp = tempfile::tempdir().expect("tempdir");
8004 let repo = tmp.path().join("repo");
8005 std::fs::create_dir_all(&repo).unwrap();
8006 init_repo(&repo);
8007 let home = tmp.path().join("home");
8008
8009 let mut config = Config::default();
8010 config.graph.worktree_root = Some(tmp.path().join("wt"));
8011 let mut state = RunState::new(
8012 repo.clone(),
8013 "main".to_owned(),
8014 "deadbeef".to_owned(),
8015 "task".to_owned(),
8016 config,
8017 );
8018 git::git(&repo, &["branch", "magi/x/A", "main"])
8020 .await
8021 .expect("branch");
8022 state.candidates = vec![Candidate {
8023 index: 0,
8024 label: 'A',
8025 agent: "alpha".to_owned(),
8026 branch: "magi/x/A".to_owned(),
8027 worktree: state.worktree_root().join("cand-A"),
8028 summary: String::new(),
8029 stat: String::new(),
8030 files: 0,
8031 commits: 0,
8032 empty: false,
8033 failed: None,
8034 verified_noop: None,
8035 duration_ms: 0,
8036 folded: true,
8037 }];
8038 state.released_to = Some("20260901-000000-new1".to_owned());
8039 state.released_branches = vec!["magi/x/A".to_owned()];
8040
8041 fold_run(&mut state, true, &home).await.expect("fold_run");
8042
8043 assert!(
8044 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
8045 "the handed-over branch survives a fold"
8046 );
8047 }
8048
8049 #[tokio::test]
8053 async fn an_empty_winner_is_detected_before_a_pull_request_is_attempted() {
8054 let tmp = tempfile::tempdir().expect("tempdir");
8055 let repo = tmp.path().join("repo");
8056 std::fs::create_dir_all(&repo).unwrap();
8057 init_repo(&repo);
8058 let run = |args: &[&str]| {
8059 let out = std::process::Command::new("git")
8060 .quiet()
8061 .args(args)
8062 .current_dir(&repo)
8063 .output()
8064 .expect("spawn git");
8065 assert!(out.status.success(), "git {args:?}");
8066 };
8067 run(&["branch", "magi/x/A"]);
8068 run(&["checkout", "-q", "-b", "magi/x/B"]);
8069 std::fs::write(repo.join("f.txt"), "x\n").unwrap();
8070 run(&["add", "-A"]);
8071 run(&["commit", "-q", "-m", "work"]);
8072 run(&["checkout", "-q", "main"]);
8073
8074 let mut state = RunState::new(
8075 repo.clone(),
8076 "main".to_owned(),
8077 "deadbeef".to_owned(),
8078 "task".to_owned(),
8079 Config::default(),
8080 );
8081 state.seeds = vec![refs::Seed {
8082 token: "magi/27b2/A".to_owned(),
8083 kind: refs::SeedKind::Unresolved,
8084 sha: String::new(),
8085 branch: true,
8086 detail: "no branch or commit named magi/27b2/A".to_owned(),
8087 }];
8088
8089 assert!(merge_is_empty(&repo, &state, "magi/x/A", MergeMode::Pr).await);
8090 assert!(merge_is_empty(&repo, &state, "magi/x/A", MergeMode::Local).await);
8091 assert!(!merge_is_empty(&repo, &state, "magi/x/B", MergeMode::Pr).await);
8092 let detail = empty_candidate_detail(&state, "main");
8093 assert!(detail.starts_with("empty candidate"), "{detail}");
8094 assert!(detail.contains("magi/27b2/A"), "{detail}");
8095 }
8096
8097 #[tokio::test]
8106 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
8107 let tmp = tempfile::tempdir().expect("tempdir");
8108 let repo = tmp.path().join("repo");
8109 std::fs::create_dir_all(&repo).unwrap();
8110 init_repo(&repo);
8111
8112 let mut config = Config::default();
8113 config.merge.mode = MergeMode::Local;
8114
8115 let mut state = RunState::new(
8116 repo.clone(),
8117 "main".to_owned(),
8118 "deadbeef".to_owned(),
8119 "task".to_owned(),
8120 config,
8121 );
8122 state.candidates = vec![Candidate {
8123 index: 0,
8124 label: 'A',
8125 agent: "alpha".to_owned(),
8126 branch: "does-not-exist".to_owned(),
8127 worktree: repo.clone(),
8128 summary: String::new(),
8129 stat: String::new(),
8130 files: 0,
8131 commits: 0,
8132 empty: false,
8133 failed: None,
8134 verified_noop: None,
8135 duration_ms: 0,
8136 folded: false,
8137 }];
8138 state.tally = Some(Tally {
8139 first_choice: BTreeMap::from([('A', 1)]),
8140 borda: BTreeMap::new(),
8141 winner: 'A',
8142 rankings: 1,
8143 unanimous_initial: true,
8144 deliberated: false,
8145 changed_votes: 0,
8146 unanimous_final: true,
8147 tie_break: None,
8148 judges: 0,
8149 present: 0,
8150 quorum: 0,
8151 met_quorum: true,
8152 uncontested: Some("only candidate A produced a change".to_owned()),
8153 });
8154 state.reviews = vec![ReviewRound {
8155 round: 1,
8156 head: "deadbeef".to_owned(),
8157 verified_head: None,
8158 verified_at: None,
8159 reviews: Vec::new(),
8160 e2e: Vec::new(),
8161 fix: None,
8162 blocking: 0,
8163 answered: 0,
8164 expected: 0,
8165 clean: true,
8166 verify_retried: false,
8167 e2e_deferred: false,
8168 e2e_defer_reason: None,
8169 progressed: false,
8170 vote_split: false,
8171 reconsideration: Vec::new(),
8172 verdict: None,
8173 }];
8174 state.gate = vec![CommandOutcome {
8175 command: "test".to_owned(),
8176 code: Some(0),
8177 output_tail: String::new(),
8178 duration_ms: 0,
8179 resource_blocked: false,
8180 }];
8181 state.gate_ran = true;
8182 state.status = RunStatus::Ready;
8187 state.merge = Some(MergeOutcome {
8188 mode: MergeMode::Local,
8189 ok: false,
8190 detail: "already concluded".to_owned(),
8191 empty: false,
8192 });
8193
8194 let mut runner = Runner {
8195 state,
8196 roles: ResolvedRoles {
8197 implementers: Vec::new(),
8198 judges: Vec::new(),
8199 reviewers: Vec::new(),
8200 fixer: None,
8201 conductor: conductor(),
8202 implementer_roster: Vec::new(),
8203 },
8204 sem: Arc::new(Semaphore::new(1)),
8205 pause: Pause::new(),
8206 interrupt: Pause::new(),
8207 };
8208
8209 runner.merge().await.expect("merge");
8210
8211 assert_eq!(
8212 runner.state.status,
8213 RunStatus::Ready,
8214 "a concluded run's status must not change on reentry"
8215 );
8216 assert_eq!(
8217 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8218 Some("already concluded"),
8219 "merge must not run again once the node already recorded an outcome"
8220 );
8221 }
8222
8223 #[tokio::test]
8232 async fn merge_refuses_a_gate_that_has_not_actually_run() {
8233 let tmp = tempfile::tempdir().expect("tempdir");
8234 let repo = tmp.path().join("repo");
8235 std::fs::create_dir_all(&repo).unwrap();
8236 init_repo(&repo);
8237
8238 let mut config = Config::default();
8239 config.merge.mode = MergeMode::Local;
8240
8241 let mut state = RunState::new(
8242 repo.clone(),
8243 "main".to_owned(),
8244 "deadbeef".to_owned(),
8245 "task".to_owned(),
8246 config,
8247 );
8248 state.candidates = vec![Candidate {
8249 index: 0,
8250 label: 'A',
8251 agent: "alpha".to_owned(),
8252 branch: "does-not-exist".to_owned(),
8253 worktree: repo.clone(),
8254 summary: String::new(),
8255 stat: String::new(),
8256 files: 0,
8257 commits: 0,
8258 empty: false,
8259 failed: None,
8260 verified_noop: None,
8261 duration_ms: 0,
8262 folded: false,
8263 }];
8264 state.tally = Some(Tally {
8265 first_choice: BTreeMap::from([('A', 1)]),
8266 borda: BTreeMap::new(),
8267 winner: 'A',
8268 rankings: 1,
8269 unanimous_initial: true,
8270 deliberated: false,
8271 changed_votes: 0,
8272 unanimous_final: true,
8273 tie_break: None,
8274 judges: 0,
8275 present: 0,
8276 quorum: 0,
8277 met_quorum: true,
8278 uncontested: Some("only candidate A produced a change".to_owned()),
8279 });
8280 state.reviews = vec![ReviewRound {
8281 round: 1,
8282 head: "deadbeef".to_owned(),
8283 verified_head: None,
8284 verified_at: None,
8285 reviews: Vec::new(),
8286 e2e: Vec::new(),
8287 fix: None,
8288 blocking: 0,
8289 answered: 0,
8290 expected: 0,
8291 clean: true,
8292 verify_retried: false,
8293 e2e_deferred: false,
8294 e2e_defer_reason: None,
8295 progressed: false,
8296 vote_split: false,
8297 reconsideration: Vec::new(),
8298 verdict: None,
8299 }];
8300 state.gate = Vec::new();
8302 state.gate_ran = false;
8303 state.status = RunStatus::Gating;
8304
8305 let mut runner = Runner {
8306 state,
8307 roles: ResolvedRoles {
8308 implementers: Vec::new(),
8309 judges: Vec::new(),
8310 reviewers: Vec::new(),
8311 fixer: None,
8312 conductor: conductor(),
8313 implementer_roster: Vec::new(),
8314 },
8315 sem: Arc::new(Semaphore::new(1)),
8316 pause: Pause::new(),
8317 interrupt: Pause::new(),
8318 };
8319
8320 runner.merge().await.expect("merge");
8321
8322 assert!(
8323 runner.state.merge.is_none(),
8324 "an empty gate must never be read as a passing one: {:?}",
8325 runner.state.merge
8326 );
8327 }
8328
8329 #[tokio::test]
8336 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
8337 let tmp = tempfile::tempdir().expect("tempdir");
8338 let repo = tmp.path().join("repo");
8339 std::fs::create_dir_all(&repo).unwrap();
8340 init_repo(&repo);
8341
8342 let config = Config::default();
8344
8345 let mut state = RunState::new(
8346 repo.clone(),
8347 "main".to_owned(),
8348 "deadbeef".to_owned(),
8349 "task".to_owned(),
8350 config,
8351 );
8352 state.candidates = vec![Candidate {
8353 index: 0,
8354 label: 'A',
8355 agent: "alpha".to_owned(),
8356 branch: "does-not-exist".to_owned(),
8357 worktree: repo.clone(),
8358 summary: String::new(),
8359 stat: String::new(),
8360 files: 0,
8361 commits: 0,
8362 empty: false,
8363 failed: None,
8364 verified_noop: None,
8365 duration_ms: 0,
8366 folded: false,
8367 }];
8368 state.tally = Some(Tally {
8369 first_choice: BTreeMap::from([('A', 1)]),
8370 borda: BTreeMap::new(),
8371 winner: 'A',
8372 rankings: 1,
8373 unanimous_initial: true,
8374 deliberated: false,
8375 changed_votes: 0,
8376 unanimous_final: true,
8377 tie_break: None,
8378 judges: 0,
8379 present: 0,
8380 quorum: 0,
8381 met_quorum: true,
8382 uncontested: Some("only candidate A produced a change".to_owned()),
8383 });
8384 state.reviews = vec![ReviewRound {
8385 round: 1,
8386 head: "deadbeef".to_owned(),
8387 verified_head: None,
8388 verified_at: None,
8389 reviews: Vec::new(),
8390 e2e: Vec::new(),
8391 fix: None,
8392 blocking: 0,
8393 answered: 0,
8394 expected: 0,
8395 clean: true,
8396 verify_retried: false,
8397 e2e_deferred: false,
8398 e2e_defer_reason: None,
8399 progressed: false,
8400 vote_split: false,
8401 reconsideration: Vec::new(),
8402 verdict: None,
8403 }];
8404
8405 let mut runner = Runner {
8406 state,
8407 roles: ResolvedRoles {
8408 implementers: Vec::new(),
8409 judges: Vec::new(),
8410 reviewers: Vec::new(),
8411 fixer: None,
8412 conductor: conductor(),
8413 implementer_roster: Vec::new(),
8414 },
8415 sem: Arc::new(Semaphore::new(1)),
8416 pause: Pause::new(),
8417 interrupt: Pause::new(),
8418 };
8419
8420 runner.gate().await.expect("gate");
8421 assert!(
8422 runner.state.gate_ran,
8423 "zero configured commands is still a real attempt, not an unrun gate"
8424 );
8425 assert!(runner.state.gate.is_empty());
8426 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
8427 assert_ne!(
8428 runner.state.status,
8429 RunStatus::Blocked,
8430 "a gate with nothing to check must not read as failed"
8431 );
8432
8433 runner.merge().await.expect("merge");
8434 assert_eq!(
8435 runner.state.status,
8436 RunStatus::Ready,
8437 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
8438 );
8439 }
8440
8441 #[tokio::test]
8452 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
8453 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8454 let home = crate::run::home();
8455
8456 let tmp = tempfile::tempdir().expect("tempdir");
8457 let repo = tmp.path().join("repo");
8458 std::fs::create_dir_all(&repo).unwrap();
8459 init_repo(&repo);
8460 let cache_dir = tmp.path().join("target");
8463
8464 let mut config = Config::default();
8465 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
8466 config.graph.timeout_verify = Some(2);
8469
8470 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8471 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8472 .expect("no io error acquiring directly")
8473 {
8474 crate::cache::AcquireOutcome::Acquired(g) => g,
8475 crate::cache::AcquireOutcome::Busy(b) => {
8476 panic!("expected the direct acquire to win the lease first: {b:?}")
8477 }
8478 };
8479
8480 let mut state = RunState::new(
8481 repo.clone(),
8482 "main".to_owned(),
8483 "deadbeef".to_owned(),
8484 "task".to_owned(),
8485 config,
8486 );
8487 state.candidates = vec![Candidate {
8488 index: 0,
8489 label: 'A',
8490 agent: "alpha".to_owned(),
8491 branch: "does-not-exist".to_owned(),
8492 worktree: repo.clone(),
8493 summary: String::new(),
8494 stat: String::new(),
8495 files: 0,
8496 commits: 0,
8497 empty: false,
8498 failed: None,
8499 verified_noop: None,
8500 duration_ms: 0,
8501 folded: false,
8502 }];
8503 state.tally = Some(Tally {
8504 first_choice: BTreeMap::from([('A', 1)]),
8505 borda: BTreeMap::new(),
8506 winner: 'A',
8507 rankings: 1,
8508 unanimous_initial: true,
8509 deliberated: false,
8510 changed_votes: 0,
8511 unanimous_final: true,
8512 tie_break: None,
8513 judges: 0,
8514 present: 0,
8515 quorum: 0,
8516 met_quorum: true,
8517 uncontested: Some("only candidate A produced a change".to_owned()),
8518 });
8519 state.reviews = vec![ReviewRound {
8520 round: 1,
8521 head: "deadbeef".to_owned(),
8522 verified_head: None,
8523 verified_at: None,
8524 reviews: Vec::new(),
8525 e2e: Vec::new(),
8526 fix: None,
8527 blocking: 0,
8528 answered: 0,
8529 expected: 0,
8530 clean: true,
8531 verify_retried: false,
8532 e2e_deferred: false,
8533 e2e_defer_reason: None,
8534 progressed: false,
8535 vote_split: false,
8536 reconsideration: Vec::new(),
8537 verdict: None,
8538 }];
8539
8540 let mut runner = Runner {
8541 state,
8542 roles: ResolvedRoles {
8543 implementers: Vec::new(),
8544 judges: Vec::new(),
8545 reviewers: Vec::new(),
8546 fixer: None,
8547 conductor: conductor(),
8548 implementer_roster: Vec::new(),
8549 },
8550 sem: Arc::new(Semaphore::new(1)),
8551 pause: Pause::new(),
8552 interrupt: Pause::new(),
8553 };
8554
8555 let started = std::time::Instant::now();
8556 runner.gate().await.expect("gate");
8557 assert!(
8558 started.elapsed() < Duration::from_secs(1),
8559 "a gate with nothing to run must never wait on a lease it never needed"
8560 );
8561 assert!(
8562 runner.state.gate_ran,
8563 "zero commands is still a real, immediate attempt"
8564 );
8565 assert!(runner.state.gate.is_empty());
8566 assert_ne!(
8567 runner.state.status,
8568 RunStatus::Blocked,
8569 "must not read as resource-blocked on a lease it never asked for"
8570 );
8571 }
8572
8573 #[tokio::test]
8583 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8584 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8585
8586 let tmp = tempfile::tempdir().expect("tempdir");
8587 let repo = tmp.path().join("repo");
8588 std::fs::create_dir_all(&repo).unwrap();
8589 init_repo(&repo);
8590
8591 let mut config = Config::default();
8592 config.verify.gate = vec![
8593 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8594 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8595 .to_owned(),
8596 ];
8597
8598 let mut state = RunState::new(
8599 repo.clone(),
8600 "main".to_owned(),
8601 "deadbeef".to_owned(),
8602 "task".to_owned(),
8603 config,
8604 );
8605 let run_id = state.id.clone();
8606 state.candidates = vec![Candidate {
8607 index: 0,
8608 label: 'A',
8609 agent: "alpha".to_owned(),
8610 branch: "does-not-exist".to_owned(),
8611 worktree: repo.clone(),
8612 summary: String::new(),
8613 stat: String::new(),
8614 files: 0,
8615 commits: 0,
8616 empty: false,
8617 failed: None,
8618 verified_noop: None,
8619 duration_ms: 0,
8620 folded: false,
8621 }];
8622 state.tally = Some(Tally {
8623 first_choice: BTreeMap::from([('A', 1)]),
8624 borda: BTreeMap::new(),
8625 winner: 'A',
8626 rankings: 1,
8627 unanimous_initial: true,
8628 deliberated: false,
8629 changed_votes: 0,
8630 unanimous_final: true,
8631 tie_break: None,
8632 judges: 0,
8633 present: 0,
8634 quorum: 0,
8635 met_quorum: true,
8636 uncontested: Some("only candidate A produced a change".to_owned()),
8637 });
8638 state.reviews = vec![ReviewRound {
8639 round: 1,
8640 head: "deadbeef".to_owned(),
8641 verified_head: None,
8642 verified_at: None,
8643 reviews: Vec::new(),
8644 e2e: Vec::new(),
8645 fix: None,
8646 blocking: 0,
8647 answered: 0,
8648 expected: 0,
8649 clean: true,
8650 verify_retried: false,
8651 e2e_deferred: false,
8652 e2e_defer_reason: None,
8653 progressed: false,
8654 vote_split: false,
8655 reconsideration: Vec::new(),
8656 verdict: None,
8657 }];
8658
8659 let mut runner = Runner {
8660 state,
8661 roles: ResolvedRoles {
8662 implementers: Vec::new(),
8663 judges: Vec::new(),
8664 reviewers: Vec::new(),
8665 fixer: None,
8666 conductor: conductor(),
8667 implementer_roster: Vec::new(),
8668 },
8669 sem: Arc::new(Semaphore::new(1)),
8670 pause: Pause::new(),
8671 interrupt: Pause::new(),
8672 };
8673
8674 let started_marker = repo.join("started.marker");
8675 let release_marker = repo.join("release.marker");
8676 let poller = tokio::spawn(async move {
8677 for _ in 0..100 {
8682 if started_marker.exists()
8683 && let Ok(s) = crate::run::RunState::load(&run_id)
8684 && let Some(a) = s.active.get("gate")
8685 {
8686 std::fs::write(&release_marker, b"go").expect("release marker");
8687 return Some(a.clone());
8688 }
8689 tokio::time::sleep(Duration::from_millis(50)).await;
8690 }
8691 None
8692 });
8693
8694 runner.gate().await.expect("gate");
8695 let captured = poller.await.expect("poller task");
8696 let captured = captured.expect(
8697 "the poller never saw a `gate` task entry in run.json while the command was \
8698 still blocked on its own release marker",
8699 );
8700
8701 assert_eq!(captured.task.as_deref(), Some("gate"));
8702 assert_eq!(captured.node, "gate");
8703 assert_eq!(captured.index, Some(1));
8704 assert_eq!(captured.total, Some(1));
8705 assert!(
8706 captured
8707 .command
8708 .as_deref()
8709 .is_some_and(|c| c.contains("started.marker")),
8710 "{captured:?}"
8711 );
8712
8713 assert!(
8714 runner.state.active.is_empty(),
8715 "the entry must be cleared once the command actually finished: {:?}",
8716 runner.state.active
8717 );
8718 assert!(runner.state.gate_ran);
8719 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8720 }
8721
8722 #[tokio::test]
8735 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8736 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8737 let home = crate::run::home();
8738
8739 let tmp = tempfile::tempdir().expect("tempdir");
8740 let repo = tmp.path().join("repo");
8741 std::fs::create_dir_all(&repo).unwrap();
8742 init_repo(&repo);
8743 let head = crate::git::rev_parse(&repo, "HEAD")
8744 .await
8745 .expect("rev-parse");
8746 let cache_dir = tmp.path().join("target");
8749
8750 let mut config = Config::default();
8751 config.verify.e2e = vec![format!(
8752 "CARGO_TARGET_DIR='{}' test -f README.md",
8753 cache_dir.display()
8754 )];
8755 config.graph.review_rounds = 1;
8756 config.graph.timeout_verify = Some(2);
8759
8760 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8761 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8762 .expect("no io error acquiring directly")
8763 {
8764 crate::cache::AcquireOutcome::Acquired(g) => g,
8765 crate::cache::AcquireOutcome::Busy(b) => {
8766 panic!("expected the direct acquire to win the lease first: {b:?}")
8767 }
8768 };
8769
8770 let mut state = RunState::new(
8771 repo.clone(),
8772 "main".to_owned(),
8773 head.clone(),
8774 "task".to_owned(),
8775 config,
8776 );
8777 state.candidates = vec![Candidate {
8778 index: 0,
8779 label: 'A',
8780 agent: "alpha".to_owned(),
8781 branch: "does-not-exist".to_owned(),
8782 worktree: repo.clone(),
8783 summary: String::new(),
8784 stat: String::new(),
8785 files: 0,
8786 commits: 0,
8787 empty: false,
8788 failed: None,
8789 verified_noop: None,
8790 duration_ms: 0,
8791 folded: false,
8792 }];
8793 state.tally = Some(Tally {
8794 first_choice: BTreeMap::from([('A', 1)]),
8795 borda: BTreeMap::new(),
8796 winner: 'A',
8797 rankings: 1,
8798 unanimous_initial: true,
8799 deliberated: false,
8800 changed_votes: 0,
8801 unanimous_final: true,
8802 tie_break: None,
8803 judges: 0,
8804 present: 0,
8805 quorum: 0,
8806 met_quorum: true,
8807 uncontested: Some("only candidate A produced a change".to_owned()),
8808 });
8809 state.reviews = vec![ReviewRound {
8813 round: 1,
8814 head: head.clone(),
8815 verified_head: None,
8816 verified_at: None,
8817 reviews: Vec::new(),
8818 e2e: Vec::new(),
8819 fix: None,
8820 blocking: 1,
8821 answered: 1,
8822 expected: 1,
8823 clean: false,
8824 verify_retried: false,
8825 e2e_deferred: true,
8826 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8827 progressed: false,
8828 vote_split: false,
8829 reconsideration: Vec::new(),
8830 verdict: None,
8831 }];
8832
8833 let mut runner = Runner {
8834 state,
8835 roles: ResolvedRoles {
8836 implementers: Vec::new(),
8837 judges: Vec::new(),
8838 reviewers: Vec::new(),
8839 fixer: None,
8840 conductor: conductor(),
8841 implementer_roster: Vec::new(),
8842 },
8843 sem: Arc::new(Semaphore::new(1)),
8844 pause: Pause::new(),
8845 interrupt: Pause::new(),
8846 };
8847
8848 let shell = runner.state.config.shell();
8849 runner
8850 .stop_reviewing("round budget spent", &shell, &repo)
8851 .await
8852 .expect("stop_reviewing");
8853
8854 let last = runner.state.reviews.last().expect("round record");
8855 assert_eq!(
8856 last.e2e_status(),
8857 E2eStatus::ResourceBlocked,
8858 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8859 failed: {last:?}"
8860 );
8861 assert_eq!(
8862 last.verified_head.as_deref(),
8863 Some(head.as_str()),
8864 "which commit this attempt targeted is known even though nothing finished checking \
8865 it"
8866 );
8867 let first_attempt_at = last
8868 .verified_at
8869 .expect("when this attempt ran is known too");
8870 assert_ne!(
8871 runner.state.status,
8872 RunStatus::Blocked,
8873 "contention is evidence about the machine, not the patch — it must not settle the \
8874 run as blocked: {:?}",
8875 runner.state.status
8876 );
8877 assert!(
8878 !runner
8879 .state
8880 .events
8881 .iter()
8882 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
8883 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
8884 runner.state.events
8885 );
8886
8887 runner
8892 .stop_reviewing("round budget spent", &shell, &repo)
8893 .await
8894 .expect("stop_reviewing retry");
8895 assert_eq!(
8896 runner.state.reviews.len(),
8897 1,
8898 "no new round was started: {:?}",
8899 runner.state.reviews
8900 );
8901 let last = runner.state.reviews.last().expect("round record");
8902 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
8903 assert!(
8904 last.verified_at.expect("still known") > first_attempt_at,
8905 "a second reentry must be a fresh attempt, not a stale copy of the first"
8906 );
8907 assert_ne!(runner.state.status, RunStatus::Blocked);
8908
8909 held.release();
8910 }
8911
8912 #[tokio::test]
8924 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
8925 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8926 let home = crate::run::home();
8927
8928 let tmp = tempfile::tempdir().expect("tempdir");
8929 let repo = tmp.path().join("repo");
8930 std::fs::create_dir_all(&repo).unwrap();
8931 init_repo(&repo);
8932 let head = crate::git::rev_parse(&repo, "HEAD")
8933 .await
8934 .expect("rev-parse");
8935 let cache_dir = tmp.path().join("target");
8936
8937 let mut config = Config::default();
8938 config.verify.e2e = vec![format!(
8939 "CARGO_TARGET_DIR='{}' test -f README.md",
8940 cache_dir.display()
8941 )];
8942 config.graph.review_rounds = 1;
8943 config.graph.timeout_verify = Some(2);
8944
8945 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8946 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8947 .expect("no io error acquiring directly")
8948 {
8949 crate::cache::AcquireOutcome::Acquired(g) => g,
8950 crate::cache::AcquireOutcome::Busy(b) => {
8951 panic!("expected the direct acquire to win the lease first: {b:?}")
8952 }
8953 };
8954
8955 let mut state = RunState::new(
8956 repo.clone(),
8957 "main".to_owned(),
8958 head.clone(),
8959 "task".to_owned(),
8960 config,
8961 );
8962 state.candidates = vec![Candidate {
8963 index: 0,
8964 label: 'A',
8965 agent: "alpha".to_owned(),
8966 branch: "does-not-exist".to_owned(),
8967 worktree: repo.clone(),
8968 summary: String::new(),
8969 stat: String::new(),
8970 files: 0,
8971 commits: 0,
8972 empty: false,
8973 failed: None,
8974 verified_noop: None,
8975 duration_ms: 0,
8976 folded: false,
8977 }];
8978 state.tally = Some(Tally {
8979 first_choice: BTreeMap::from([('A', 1)]),
8980 borda: BTreeMap::new(),
8981 winner: 'A',
8982 rankings: 1,
8983 unanimous_initial: true,
8984 deliberated: false,
8985 changed_votes: 0,
8986 unanimous_final: true,
8987 tie_break: None,
8988 judges: 0,
8989 present: 0,
8990 quorum: 0,
8991 met_quorum: true,
8992 uncontested: Some("only candidate A produced a change".to_owned()),
8993 });
8994 state.reviews = vec![ReviewRound {
8998 round: 1,
8999 head: head.clone(),
9000 verified_head: Some(head.clone()),
9001 verified_at: Some(jiff::Timestamp::now()),
9002 reviews: Vec::new(),
9003 e2e: vec![CommandOutcome {
9004 command: format!(
9005 "CARGO_TARGET_DIR='{}' test -f README.md",
9006 cache_dir.display()
9007 ),
9008 code: None,
9009 output_tail: "waiting for the shared build cache".to_owned(),
9010 duration_ms: 0,
9011 resource_blocked: true,
9012 }],
9013 fix: None,
9014 blocking: 1,
9015 answered: 1,
9016 expected: 1,
9017 clean: false,
9018 verify_retried: false,
9019 e2e_deferred: false,
9020 e2e_defer_reason: None,
9021 progressed: false,
9022 vote_split: false,
9023 reconsideration: Vec::new(),
9024 verdict: None,
9025 }];
9026
9027 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
9028 let mut runner = Runner {
9029 state,
9030 roles: ResolvedRoles {
9031 implementers: Vec::new(),
9032 judges: Vec::new(),
9033 reviewers: Vec::new(),
9034 fixer: None,
9035 conductor: conductor(),
9036 implementer_roster: Vec::new(),
9037 },
9038 sem: Arc::new(Semaphore::new(1)),
9039 pause: Pause::new(),
9040 interrupt: Pause::new(),
9041 };
9042
9043 runner.review_loop().await.expect("review_loop");
9048
9049 assert_eq!(
9050 runner.state.reviews.len(),
9051 1,
9052 "no new round was started on top of the unresolved one: {:?}",
9053 runner.state.reviews
9054 );
9055 let last = &runner.state.reviews[0];
9056 assert_eq!(
9057 last.e2e_status(),
9058 E2eStatus::ResourceBlocked,
9059 "still contended: {last:?}"
9060 );
9061 assert!(
9062 last.verified_at.expect("still known") > first_attempt_at,
9063 "review_loop must have actually retried the check, not left it exactly as found"
9064 );
9065 assert_ne!(
9066 runner.state.status,
9067 RunStatus::Blocked,
9068 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
9069 runner.state.status
9070 );
9071
9072 held.release();
9073 }
9074
9075 #[tokio::test]
9076 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
9077 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
9078 let tmp = tempfile::tempdir().expect("tempdir");
9079 let repo = tmp.path().join("repo");
9080 std::fs::create_dir_all(&repo).unwrap();
9081 init_repo(&repo);
9082
9083 let mut config = Config::default();
9084 config.merge.mode = MergeMode::Pr;
9085 config.graph.land = true;
9086 config.graph.land_approval = false;
9087
9088 let mut state = RunState::new(
9089 repo.clone(),
9090 "main".to_owned(),
9091 "deadbeef".to_owned(),
9092 "task".to_owned(),
9093 config,
9094 );
9095 state.candidates = vec![Candidate {
9096 index: 0,
9097 label: 'A',
9098 agent: "alpha".to_owned(),
9099 branch: "does-not-exist".to_owned(),
9100 worktree: repo.clone(),
9101 summary: String::new(),
9102 stat: String::new(),
9103 files: 0,
9104 commits: 0,
9105 empty: false,
9106 failed: None,
9107 verified_noop: None,
9108 duration_ms: 0,
9109 folded: false,
9110 }];
9111 state.tally = Some(Tally {
9112 first_choice: BTreeMap::from([('A', 1)]),
9113 borda: BTreeMap::new(),
9114 winner: 'A',
9115 rankings: 1,
9116 unanimous_initial: true,
9117 deliberated: false,
9118 changed_votes: 0,
9119 unanimous_final: true,
9120 tie_break: None,
9121 judges: 0,
9122 present: 0,
9123 quorum: 0,
9124 met_quorum: true,
9125 uncontested: Some("only candidate A produced a change".to_owned()),
9126 });
9127 state.reviews = vec![ReviewRound {
9128 round: 1,
9129 head: "deadbeef".to_owned(),
9130 verified_head: None,
9131 verified_at: None,
9132 reviews: Vec::new(),
9133 e2e: Vec::new(),
9134 fix: None,
9135 blocking: 0,
9136 answered: 0,
9137 expected: 0,
9138 clean: true,
9139 verify_retried: false,
9140 e2e_deferred: false,
9141 e2e_defer_reason: None,
9142 progressed: false,
9143 vote_split: false,
9144 reconsideration: Vec::new(),
9145 verdict: None,
9146 }];
9147 state.gate = vec![CommandOutcome {
9148 command: "test".to_owned(),
9149 code: Some(0),
9150 output_tail: String::new(),
9151 duration_ms: 0,
9152 resource_blocked: false,
9153 }];
9154 state.gate_ran = true;
9155 state.status = RunStatus::Landing;
9159 state.merge = Some(MergeOutcome {
9160 mode: MergeMode::Pr,
9161 ok: true,
9162 detail: "https://example.invalid/x/y/pull/1".to_owned(),
9163 empty: false,
9164 });
9165
9166 ask_test_home();
9170 let store = ask::Questions::open();
9171 let q = ask_open_question(&store, &state.id);
9172
9173 let mut runner = Runner {
9174 state,
9175 roles: ResolvedRoles {
9176 implementers: Vec::new(),
9177 judges: Vec::new(),
9178 reviewers: Vec::new(),
9179 fixer: None,
9180 conductor: conductor(),
9181 implementer_roster: Vec::new(),
9182 },
9183 sem: Arc::new(Semaphore::new(1)),
9184 pause: Pause::new(),
9185 interrupt: Pause::new(),
9186 };
9187
9188 runner.execute().await.expect("execute");
9193
9194 assert_eq!(
9195 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
9196 Some("https://example.invalid/x/y/pull/1"),
9197 "reentry must not push again or open a second pull request over the \
9198 one `land` is already watching"
9199 );
9200 assert_ne!(
9201 runner.state.status,
9202 RunStatus::Landing,
9203 "land could not actually reach the fake pull request, so it must \
9204 have given up rather than left the run silently parked forever"
9205 );
9206 assert_eq!(runner.state.status, RunStatus::Blocked);
9210 assert!(
9211 store.get(&q.id).unwrap().status.open(),
9212 "Blocked is still alive; settle_questions must have been a no-op here"
9213 );
9214 }
9215
9216 fn state_with_round(round: ReviewRound) -> RunState {
9217 let mut s = RunState::new(
9218 PathBuf::from("/repo"),
9219 "main".to_owned(),
9220 "abc1234".to_owned(),
9221 "add retries".to_owned(),
9222 Config::default(),
9223 );
9224 s.reviews = vec![round];
9225 s
9226 }
9227
9228 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
9229 crate::verdict::Finding {
9230 id: id.to_owned(),
9231 severity,
9232 file: None,
9233 line: None,
9234 title: title.to_owned(),
9235 detail: String::new(),
9236 }
9237 }
9238
9239 #[test]
9240 fn pr_body_names_open_findings_and_declined_ones() {
9241 let round = ReviewRound {
9242 round: 2,
9243 head: "deadbee".to_owned(),
9244 verified_head: None,
9245 verified_at: None,
9246 reviews: vec![ReviewRecord {
9247 attempts: 0,
9248 reviewer: 1,
9249 agent: "alpha".to_owned(),
9250 summary: String::new(),
9251 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
9252 vote: None,
9253 failed: None,
9254 duration_ms: 0,
9255 }],
9256 e2e: vec![CommandOutcome {
9257 command: "cargo test".to_owned(),
9258 code: Some(0),
9259 output_tail: String::new(),
9260 duration_ms: 0,
9261 resource_blocked: false,
9262 }],
9263 verify_retried: false,
9264 e2e_deferred: false,
9265 e2e_defer_reason: None,
9266 fix: Some(FixRecord {
9267 agent: "alpha".to_owned(),
9268 addressed: Vec::new(),
9269 rejected: vec![crate::verdict::Rejection {
9270 id: "R1-1-1".to_owned(),
9271 why: "not reachable from any caller".to_owned(),
9272 }],
9273 notes: String::new(),
9274 committed: true,
9275 failed: None,
9276 duration_ms: 0,
9277 continuation: None,
9278 }),
9279 blocking: 0,
9280 answered: 1,
9281 expected: 1,
9282 clean: false,
9283 progressed: true,
9284 vote_split: false,
9285 reconsideration: Vec::new(),
9286 verdict: None,
9287 };
9288 let state = state_with_round(round);
9289 let body = pr_message(&state, 'A').body;
9290
9291 assert!(body.contains("add retries"), "the task must still be there");
9292 assert!(body.contains("R2-1-1"), "{body}");
9293 assert!(body.contains("unused import"), "{body}");
9294 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
9295 assert!(
9296 body.contains("not reachable from any caller"),
9297 "the reason it was declined: {body}"
9298 );
9299 }
9300
9301 #[test]
9302 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
9303 let round = ReviewRound {
9304 round: 1,
9305 head: "deadbee".to_owned(),
9306 verified_head: None,
9307 verified_at: None,
9308 reviews: vec![ReviewRecord {
9309 attempts: 0,
9310 reviewer: 1,
9311 agent: "alpha".to_owned(),
9312 summary: String::new(),
9313 findings: Vec::new(),
9314 vote: None,
9315 failed: None,
9316 duration_ms: 0,
9317 }],
9318 e2e: Vec::new(),
9319 verify_retried: false,
9320 e2e_deferred: false,
9321 e2e_defer_reason: None,
9322 fix: None,
9323 blocking: 0,
9324 answered: 1,
9325 expected: 1,
9326 clean: true,
9327 progressed: false,
9328 vote_split: false,
9329 reconsideration: Vec::new(),
9330 verdict: None,
9331 };
9332 let state = state_with_round(round);
9333 let body = pr_message(&state, 'A').body;
9334 assert!(!body.contains("Open review findings"), "{body}");
9335 assert!(!body.contains("Declined"), "{body}");
9336 }
9337
9338 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
9339 let mut state = RunState::new(
9340 PathBuf::from("/repo"),
9341 "main".to_owned(),
9342 "abc1234".to_owned(),
9343 instruction.to_owned(),
9344 Config::default(),
9345 );
9346 state.candidates.push(Candidate {
9347 index: 0,
9348 label: 'A',
9349 agent: "alpha".to_owned(),
9350 branch: "magi/x/A".to_owned(),
9351 worktree: PathBuf::from("/wt"),
9352 summary: summary.to_owned(),
9353 stat: String::new(),
9354 files: 1,
9355 commits: 1,
9356 empty: false,
9357 failed: None,
9358 verified_noop: None,
9359 folded: false,
9360 duration_ms: 0,
9361 });
9362 state
9363 }
9364
9365 #[test]
9366 fn pr_message_describes_the_change_not_the_task() {
9367 let state = state_with_summary(
9368 "今回やってほしいこと: results projector を直す",
9369 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
9370 );
9371 let m = pr_message(&state, 'A');
9372 assert_eq!(m.title, "fix(web): batch the runs list reads");
9373 assert!(
9374 m.body.starts_with("## Summary\n\n- reads run.json once"),
9375 "{}",
9376 m.body
9377 );
9378 assert!(!m.body.contains("TITLE:"), "{}", m.body);
9379 let task_at = m.body.find("今回やってほしいこと").unwrap();
9380 let details_at = m.body.find("<details>").unwrap();
9381 assert!(
9382 details_at < task_at,
9383 "the task lives inside <details>: {}",
9384 m.body
9385 );
9386 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
9387 assert!(m.body.contains("magi:candidate-a"));
9388 }
9389
9390 #[test]
9391 fn pr_message_falls_back_to_the_task_without_a_title_line() {
9392 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
9393 let m = pr_message(&state, 'A');
9394 assert_eq!(m.title, "add retries");
9395 assert!(
9396 m.body.contains("## Summary\n\n- did some things"),
9397 "{}",
9398 m.body
9399 );
9400
9401 let none = RunState::new(
9402 PathBuf::from("/repo"),
9403 "main".to_owned(),
9404 "abc1234".to_owned(),
9405 "add retries".to_owned(),
9406 Config::default(),
9407 );
9408 let m = pr_message(&none, 'A');
9409 assert_eq!(m.title, "add retries");
9410 assert!(!m.body.contains("## Summary"), "{}", m.body);
9411 }
9412
9413 #[test]
9414 fn pr_message_refuses_the_candidate_commit_subject() {
9415 for bad in [
9416 "TITLE: magi: candidate A (uncommitted work)",
9417 "TITLE: chore: stuff (uncommitted work)",
9418 "TITLE: ",
9419 ] {
9420 let state = state_with_summary("add retries", bad);
9421 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
9422 }
9423 }
9424
9425 #[test]
9426 fn pr_message_bounds_a_very_long_task_and_title() {
9427 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
9428 let state = state_with_summary(&long, "- nothing");
9429 let m = pr_message(&state, 'A');
9430 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9431 assert!(!m.title.contains('\n'));
9432
9433 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
9434 let m = pr_message(&state, 'A');
9435 assert!(m.title.starts_with("feat: "));
9436 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9437 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
9438 }
9439
9440 #[test]
9441 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
9442 let mut state = state_with_summary(
9446 "add retries",
9447 "TITLE: fix(web): batch reads\n- reads run.json once",
9448 );
9449 state.config.graph.language = "ja".to_owned();
9450 let m = pr_message(&state, 'A');
9451 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
9452
9453 let task = "今回やってほしいこと: results projector を直す";
9456 let mut state = state_with_summary(task, "- no title line");
9457 state.config.graph.language = "ja".to_owned();
9458 let m = pr_message(&state, 'A');
9459 assert_eq!(
9460 m.title,
9461 format!("chore: land candidate A of run {}", state.id)
9462 );
9463 assert!(
9464 m.body.contains(&format!(
9465 "<summary>Original task</summary>\n\n{task}\n\n</details>"
9466 )),
9467 "{}",
9468 m.body
9469 );
9470 }
9471
9472 #[test]
9473 fn pr_message_scrubs_home_paths_and_addresses() {
9474 let state = state_with_summary(
9475 "fix it in /Users/someone/src/x",
9476 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9477 );
9478 let m = pr_message(&state, 'A');
9479 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9480 assert!(!m.body.contains(leak), "{}", m.body);
9481 }
9482 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9483 }
9484
9485 #[test]
9486 fn pr_message_survives_a_task_that_closes_details() {
9487 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9488 let m = pr_message(&state, 'A');
9489 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9490 }
9491
9492 #[test]
9493 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9494 let cmd = manual_merge_command(
9495 MergeStyle::Squash,
9496 Path::new("/repo"),
9497 "b",
9498 "fix: \"quoted\" $(x) `y`\n\nbody",
9499 );
9500 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9501 }
9502
9503 #[test]
9504 fn manual_merge_command_matches_the_configured_style() {
9505 let repo = Path::new("/repo");
9506 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9507
9508 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9509 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9510
9511 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9512 assert_eq!(
9513 squash,
9514 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9515 \"Merge magi run 0832 (candidate A)\""
9516 );
9517
9518 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9519 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9520 }
9521
9522 #[test]
9523 fn a_nudge_gets_a_quarter_of_the_budget() {
9524 assert_eq!(retry_budget(secs(1200), true), secs(300));
9526 assert_eq!(retry_budget(secs(3600), true), secs(900));
9527 }
9528
9529 #[test]
9530 fn a_resent_prompt_keeps_the_whole_budget() {
9531 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9534 assert_eq!(retry_budget(secs(60), false), secs(60));
9535 }
9536
9537 #[test]
9538 fn the_floor_never_exceeds_the_original_budget() {
9539 assert_eq!(retry_budget(secs(60), true), secs(60));
9543 assert_eq!(retry_budget(secs(480), true), secs(120));
9544 assert_eq!(retry_budget(secs(0), true), secs(0));
9545 }
9546
9547 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9548 agent::CommandEvidence {
9549 id: "item1".to_owned(),
9550 description: "cargo test".to_owned(),
9551 exit_code,
9552 result_summary: String::new(),
9553 source: "codex".to_owned(),
9554 }
9555 }
9556
9557 #[test]
9558 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9559 assert!(!has_unconfirmed_command(&[]));
9563 }
9564
9565 #[test]
9566 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9567 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9571 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9572 assert!(!has_unconfirmed_command(&[
9573 evidence(Some(0)),
9574 evidence(Some(101))
9575 ]));
9576 }
9577
9578 #[test]
9579 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9580 assert!(has_unconfirmed_command(&[
9581 evidence(Some(0)),
9582 evidence(None)
9583 ]));
9584 }
9585
9586 #[test]
9587 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9588 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9589 assert_eq!(
9590 verified_noop_claim(true, &[], text).as_deref(),
9591 Some("already fixed by b32cfc4, on main.")
9592 );
9593 }
9594
9595 #[test]
9596 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9597 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9600 assert!(verified_noop_claim(false, &[], text).is_none());
9601 }
9602
9603 #[test]
9604 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9605 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9606 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9607 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9609 }
9610
9611 #[test]
9612 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9613 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9614 }
9615
9616 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9619 runner.state.candidates = shape
9620 .iter()
9621 .enumerate()
9622 .map(|(i, &(empty, verified))| Candidate {
9623 index: i,
9624 label: (b'A' + i as u8) as char,
9625 agent: "sonnet".to_owned(),
9626 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9627 worktree: PathBuf::from(format!("/wt/{i}")),
9628 summary: String::new(),
9629 stat: String::new(),
9630 files: 0,
9631 commits: 0,
9632 empty,
9633 failed: None,
9634 verified_noop: verified.map(str::to_owned),
9635 duration_ms: 0,
9636 folded: false,
9637 })
9638 .collect();
9639 }
9640
9641 #[test]
9642 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9643 ask_test_home();
9644 let mut runner = runner_at(RunStatus::Implementing);
9645 set_candidates(
9646 &mut runner,
9647 &[
9648 (true, Some("already on main at b32cfc4")),
9649 (true, Some("same fix, see the existing test")),
9650 ],
9651 );
9652
9653 runner
9654 .after_implement()
9655 .expect("a verified no-op is not an error");
9656
9657 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9658 }
9659
9660 #[test]
9661 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9662 ask_test_home();
9663 let mut runner = runner_at(RunStatus::Implementing);
9664 set_candidates(
9668 &mut runner,
9669 &[(true, Some("already on main at b32cfc4")), (true, None)],
9670 );
9671
9672 let err = runner
9673 .after_implement()
9674 .expect_err("an unverified empty candidate must still fail the run");
9675
9676 assert!(
9677 err.to_string().contains("no candidate produced a change"),
9678 "{err}"
9679 );
9680 assert_eq!(runner.state.status, RunStatus::Failed);
9681 }
9682
9683 #[test]
9684 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9685 ask_test_home();
9686 let mut runner = runner_at(RunStatus::Implementing);
9687 set_candidates(&mut runner, &[(true, None), (true, None)]);
9688
9689 let err = runner
9690 .after_implement()
9691 .expect_err("no candidate declared anything; this is an ordinary failure");
9692
9693 assert!(
9694 err.to_string().contains("no candidate produced a change"),
9695 "{err}"
9696 );
9697 assert_eq!(runner.state.status, RunStatus::Failed);
9698 }
9699}