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 let attachments = self.state.attachments.clone();
1502
1503 let mut jobs = Vec::new();
1504 for &i in &todo {
1505 let (index, label, worktree) = {
1506 let c = &self.state.candidates[i];
1507 (c.index, c.label, c.worktree.clone())
1508 };
1509 let spec = self.roles.implementers[index].clone();
1510 let seat_key = format!("impl-{label}");
1511 let seat = self.seat(&seat_key, &spec.id);
1512 let instruction = seeded_instruction(&self.state);
1513 jobs.push(SeatJob {
1514 spec,
1515 seat,
1516 prompt: prompt::implement(
1517 &instruction,
1518 &worktree.to_string_lossy(),
1519 &language,
1520 brief.as_deref(),
1521 &attachments,
1522 ),
1523 cwd: worktree,
1524 timeout,
1525 allow_write: true,
1526 sessions,
1527 artifacts: artifacts.clone(),
1528 stem: format!("impl-{label}"),
1529 });
1530 }
1531
1532 self.state.event(
1533 "implement",
1534 format!("{} candidates in parallel", jobs.len()),
1535 );
1536 let mut sent = jobs.clone();
1542 let cache = self.state.config.cache_dir();
1543 let ctx = WaveCtx {
1544 run: &run_id,
1545 node: "implement",
1546 prompts: &prompts,
1547 cache: cache.as_deref(),
1548 round: None,
1549 };
1550 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1551 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1552 .await;
1553 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1554 .await;
1555 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1556 .await;
1557
1558 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1559 let seat_key = seat.key.clone();
1560 let agent = seat.agent.clone();
1569 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1570 self.state.seats.insert(seat.key.clone(), seat);
1571 let label = self.state.candidates[i].label;
1572 let worktree = self.state.candidates[i].worktree.clone();
1573 let base = self.state.base_commit.clone();
1574
1575 let (summary, duration, failed, verified_claim) = match out {
1576 AgentOutcome::Ok(o) => {
1577 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1578 let failed = (!o.usable()).then(|| {
1579 if o.timed_out {
1580 "agent timed out".to_owned()
1581 } else {
1582 format!("agent exited with {:?}", o.exit_code)
1583 }
1584 });
1585 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1586 (text, o.duration_ms, failed, verified_claim)
1587 }
1588 AgentOutcome::Dropped(o) => {
1594 let why = o
1595 .dropped
1596 .as_ref()
1597 .map(|d| d.why.as_str())
1598 .unwrap_or("the CLI ended the stream without delivering its answer");
1599 (
1600 String::new(),
1601 o.duration_ms,
1602 Some(format!("the CLI dropped the stream ({why})")),
1603 None,
1604 )
1605 }
1606 AgentOutcome::Quota(o) => {
1607 self.state.quota.push(QuotaLoss {
1608 seat: seat_key,
1609 node: "implement".to_owned(),
1610 at: Timestamp::now(),
1611 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1612 });
1613 (
1614 String::new(),
1615 o.duration_ms,
1616 Some("rate limited (quota); produced no change".to_owned()),
1617 None,
1618 )
1619 }
1620 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1621 };
1622
1623 let rescued = match git::rescue_commit(
1626 &worktree,
1627 &format!("magi: candidate {label} (uncommitted work)"),
1628 )
1629 .await
1630 {
1631 Ok(r) => {
1632 self.state.note_withheld("implement", &r.withheld);
1633 r.committed
1634 }
1635 Err(_) => false,
1636 };
1637 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1638 .await
1639 .unwrap_or(0);
1640 let patch = git::diff(&worktree, &base, "HEAD")
1641 .await
1642 .unwrap_or_default();
1643 let stat = git::diff_stat(&worktree, &base, "HEAD")
1644 .await
1645 .unwrap_or_default();
1646 let files = git::changed_files(&worktree, &base, "HEAD")
1647 .await
1648 .map(|f| f.len())
1649 .unwrap_or(0);
1650 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1651
1652 let c = &mut self.state.candidates[i];
1653 if !exhausted_the_fallback_chain {
1654 c.agent = agent;
1655 }
1656 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1657 c.stat = stat;
1658 c.files = files;
1659 c.commits = commits;
1660 c.duration_ms = duration;
1661 c.empty = commits == 0 || patch.trim().is_empty();
1662 c.failed = match failed {
1665 Some(_) if c.empty => failed,
1666 _ => None,
1667 };
1668 c.verified_noop = if c.empty { verified_claim } else { None };
1673 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1674 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1675 (None, true, Some(_), _) => {
1676 format!("candidate {label}: no change produced (agent-verified no-op)")
1677 }
1678 (None, true, None, _) => format!("candidate {label}: no change produced"),
1679 (None, false, _, true) => {
1680 format!(
1681 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1682 )
1683 }
1684 (None, false, _, false) => {
1685 format!("candidate {label}: {files} files, {commits} commits")
1686 }
1687 };
1688 self.state.event("implement", note);
1689 self.state.save()?;
1690 }
1691
1692 self.after_implement()
1693 }
1694
1695 async fn resume_undelivered(
1723 &mut self,
1724 results: &mut [(usize, SeatState, AgentOutcome)],
1725 sent: &[SeatJob],
1726 prompts: &Prompts,
1727 run_id: &str,
1728 ) {
1729 for (wi, seat, out) in results.iter_mut() {
1730 let Some(dropped) = (match &*out {
1731 AgentOutcome::Dropped(o) => o.dropped.clone(),
1732 _ => None,
1733 }) else {
1734 continue;
1735 };
1736 let Some(job) = sent.get(*wi) else { continue };
1737 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1739 self.state.event(
1740 "implement",
1741 format!(
1742 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1743 work is in the tree",
1744 seat.key, dropped.output_tokens, dropped.why
1745 ),
1746 );
1747 continue;
1748 }
1749 if !has_context(&job.spec, seat, job.sessions) {
1757 self.state.event(
1758 "implement",
1759 format!(
1760 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1761 is no session left to resume",
1762 seat.key, dropped.output_tokens, dropped.why
1763 ),
1764 );
1765 continue;
1766 }
1767 self.state.event(
1768 "implement",
1769 format!(
1770 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1771 conversation",
1772 seat.key, dropped.output_tokens, dropped.why
1773 ),
1774 );
1775 let mut retry = job.clone();
1776 retry.seat = seat.clone();
1777 retry.prompt = prompt::resume_after_drop(&dropped.why);
1778 retry.timeout = retry_budget(job.timeout, true);
1779 retry.stem = format!("{}-resume", job.stem);
1780 let cache = self.state.config.cache_dir();
1781 let ctx = WaveCtx {
1782 run: run_id,
1783 node: "implement",
1784 prompts,
1785 cache: cache.as_deref(),
1786 round: None,
1787 };
1788 let (resumed_seat, resumed) =
1789 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1790 *seat = resumed_seat;
1791 *out = resumed;
1792 }
1793 }
1794
1795 async fn resume_quota_losses(
1857 &mut self,
1858 results: &mut [(usize, SeatState, AgentOutcome)],
1859 sent: &mut [SeatJob],
1860 prompts: &Prompts,
1861 run_id: &str,
1862 ) {
1863 let instruction = seeded_instruction(&self.state);
1864 let language = self.state.config.graph.language.clone();
1865 let brief = self
1866 .state
1867 .advice
1868 .as_ref()
1869 .and_then(|a| a.synthesis.as_deref())
1870 .map(str::to_owned);
1871 let attachments = self.state.attachments.clone();
1872 for (wi, seat, out) in results.iter_mut() {
1873 let Some(job) = sent.get_mut(*wi) else {
1874 continue;
1875 };
1876 let start = self
1881 .roles
1882 .implementer_roster
1883 .iter()
1884 .position(|s| s.id == job.spec.id)
1885 .unwrap_or(0);
1886 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1887 let mut fallback_attempt = 0usize;
1888 while matches!(&*out, AgentOutcome::Quota(_)) {
1889 let Some(next) =
1890 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1891 .cloned()
1892 else {
1893 break;
1894 };
1895 tried.insert(next.id.clone());
1896 fallback_attempt += 1;
1897
1898 if let Ok(r) = git::rescue_commit(
1899 &job.cwd,
1900 &format!(
1901 "magi: candidate {} (uncommitted work before quota fallback)",
1902 seat.key
1903 ),
1904 )
1905 .await
1906 {
1907 self.state.note_withheld("implement", &r.withheld);
1908 }
1909
1910 self.state.event(
1911 "implement",
1912 format!(
1913 "{}: rate limited (quota) on {}; retrying with {}",
1914 seat.key, seat.agent, next.id
1915 ),
1916 );
1917
1918 let new_seat = self.seat(&seat.key, &next.id);
1919 job.spec = next.clone();
1927 let mut retry = job.clone();
1928 retry.seat = new_seat;
1929 retry.prompt = prompt::implement(
1930 &instruction,
1931 &job.cwd.to_string_lossy(),
1932 &language,
1933 brief.as_deref(),
1934 &attachments,
1935 );
1936 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1937 let cache = self.state.config.cache_dir();
1938 let ctx = WaveCtx {
1939 run: run_id,
1940 node: "implement",
1941 prompts,
1942 cache: cache.as_deref(),
1943 round: None,
1944 };
1945 let (fallback_seat, fallback_out) = run_one(
1946 retry,
1947 Arc::clone(&self.sem),
1948 &ctx,
1949 &mut self.state,
1950 fallback_attempt,
1951 )
1952 .await;
1953 *seat = fallback_seat;
1954 *out = fallback_out;
1955 }
1956 }
1957 }
1958
1959 async fn resume_unconfirmed_commands(
1983 &mut self,
1984 results: &mut [(usize, SeatState, AgentOutcome)],
1985 sent: &[SeatJob],
1986 prompts: &Prompts,
1987 run_id: &str,
1988 ) {
1989 for (wi, seat, out) in results.iter_mut() {
1990 let AgentOutcome::Ok(o) = &*out else {
1991 continue;
1992 };
1993 if !has_unconfirmed_command(&o.commands) {
1994 continue;
1995 }
1996 let Some(job) = sent.get(*wi) else { continue };
1997 if !has_context(&job.spec, seat, job.sessions) {
1998 self.state.event(
1999 "implement",
2000 format!(
2001 "{}: the reply named a command whose own CLI never confirmed the exit \
2002 status of, but there is no session left to resume",
2003 seat.key
2004 ),
2005 );
2006 continue;
2007 }
2008 self.state.event(
2009 "implement",
2010 format!(
2011 "{}: the reply named a command whose own CLI never confirmed the exit \
2012 status of; resuming the conversation",
2013 seat.key
2014 ),
2015 );
2016 let mut retry = job.clone();
2017 retry.seat = seat.clone();
2018 retry.prompt = prompt::resume_incomplete(
2019 "a command in your last reply had no confirmed exit status",
2020 );
2021 retry.timeout = retry_budget(job.timeout, true);
2022 retry.stem = format!("{}-confirm", job.stem);
2023 let cache = self.state.config.cache_dir();
2024 let ctx = WaveCtx {
2025 run: run_id,
2026 node: "implement",
2027 prompts,
2028 cache: cache.as_deref(),
2029 round: None,
2030 };
2031 let (resumed_seat, resumed) =
2032 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
2033 *seat = resumed_seat;
2034 *out = resumed;
2035 }
2036 }
2037
2038 async fn continue_fix_report(
2059 &mut self,
2060 mut seat: SeatState,
2061 parse_err: String,
2062 job: &SeatJob,
2063 prompts: &Prompts,
2064 run_id: &str,
2065 round: usize,
2066 ) -> (
2067 SeatState,
2068 Option<FixReport>,
2069 Option<String>,
2070 ContinuationRecord,
2071 ) {
2072 let mut last_err = parse_err;
2073 let mut cumulative_wait_ms = 0u64;
2074 let mut attempts = 0usize;
2075 loop {
2076 if !has_context(&job.spec, &seat, job.sessions) {
2077 self.state.event(
2078 "fix",
2079 format!(
2080 "round {round}: fixer's reply had no adoption report ({last_err}); no \
2081 session left to resume into"
2082 ),
2083 );
2084 let outcome = if attempts == 0 {
2085 ContinuationOutcome::NoSession
2086 } else {
2087 ContinuationOutcome::Exhausted
2088 };
2089 return (
2090 seat,
2091 None,
2092 Some(format!("unparsable fix report: {last_err}")),
2093 ContinuationRecord {
2094 attempts,
2095 cumulative_wait_ms,
2096 outcome,
2097 },
2098 );
2099 }
2100 if attempts >= MAX_FIX_CONTINUATIONS {
2101 self.state.event(
2102 "fix",
2103 format!(
2104 "round {round}: fixer's reply still had no adoption report after \
2105 {attempts} continuation(s) ({last_err}); giving up"
2106 ),
2107 );
2108 return (
2109 seat,
2110 None,
2111 Some(format!(
2112 "unparsable fix report after {attempts} continuation(s): {last_err}"
2113 )),
2114 ContinuationRecord {
2115 attempts,
2116 cumulative_wait_ms,
2117 outcome: ContinuationOutcome::Exhausted,
2118 },
2119 );
2120 }
2121 attempts += 1;
2122 self.state.event(
2123 "fix",
2124 format!(
2125 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
2126 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
2127 ),
2128 );
2129 let mut retry = job.clone();
2130 retry.seat = seat.clone();
2131 retry.prompt = prompt::resume_incomplete(&last_err);
2132 retry.timeout = retry_budget(job.timeout, true);
2133 retry.stem = format!("{}-continue{attempts}", job.stem);
2134 let cache = self.state.config.cache_dir();
2135 let ctx = WaveCtx {
2136 run: run_id,
2137 node: "fix",
2138 prompts,
2139 cache: cache.as_deref(),
2140 round: Some(round),
2141 };
2142 let (resumed_seat, resumed_out) = run_one(
2143 retry,
2144 Arc::clone(&self.sem),
2145 &ctx,
2146 &mut self.state,
2147 attempts,
2148 )
2149 .await;
2150 seat = resumed_seat;
2151 match resumed_out {
2152 AgentOutcome::Ok(o) => {
2153 cumulative_wait_ms += o.duration_ms;
2154 match verdict::extract_json::<FixReport>(&o.text) {
2155 Ok(report) if !has_unconfirmed_command(&o.commands) => {
2156 self.state.event(
2157 "fix",
2158 format!(
2159 "round {round}: fixer's adoption report recovered after \
2160 {attempts} continuation(s)"
2161 ),
2162 );
2163 return (
2164 seat,
2165 Some(report),
2166 None,
2167 ContinuationRecord {
2168 attempts,
2169 cumulative_wait_ms,
2170 outcome: ContinuationOutcome::Resumed,
2171 },
2172 );
2173 }
2174 Ok(_) => {
2182 last_err = "the reply parsed, but it reported a command whose own CLI \
2183 never confirmed an exit status"
2184 .to_owned();
2185 }
2186 Err(e) => last_err = e.to_string(),
2187 }
2188 }
2189 AgentOutcome::Quota(o) => {
2190 cumulative_wait_ms += o.duration_ms;
2191 self.state.quota.push(QuotaLoss {
2192 seat: seat.key.clone(),
2193 node: "fix".to_owned(),
2194 at: Timestamp::now(),
2195 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2196 });
2197 self.state.event(
2198 "fix",
2199 format!(
2200 "round {round}: continuation rate limited (quota); not retrying now"
2201 ),
2202 );
2203 return (
2204 seat,
2205 None,
2206 Some("rate limited (quota) while recovering the fix report".to_owned()),
2207 ContinuationRecord {
2208 attempts,
2209 cumulative_wait_ms,
2210 outcome: ContinuationOutcome::QuotaLost,
2211 },
2212 );
2213 }
2214 AgentOutcome::Dropped(o) => {
2215 cumulative_wait_ms += o.duration_ms;
2216 let why = o
2217 .dropped
2218 .as_ref()
2219 .map(|d| d.why.as_str())
2220 .unwrap_or("the CLI ended the stream without delivering its answer");
2221 last_err = format!("the CLI dropped the stream ({why})");
2222 }
2223 AgentOutcome::Failed(e) => last_err = e,
2224 }
2225 }
2226 }
2227
2228 fn after_implement(&mut self) -> Result<()> {
2229 if self.state.leaks.is_empty() {
2231 let cfg = self.state.config.blind.clone();
2232 let mut leaks = Vec::new();
2233 for c in &self.state.candidates {
2234 let Some(patch) =
2235 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
2236 else {
2237 continue;
2238 };
2239 leaks.extend(blind::scan(
2240 &format!("candidate {} patch", c.label),
2241 &patch,
2242 &cfg.vendor_tokens,
2243 ));
2244 }
2245 if !leaks.is_empty() {
2246 let summary = leaks
2247 .iter()
2248 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2249 .collect::<Vec<_>>()
2250 .join(", ");
2251 match cfg.on_leak {
2252 LeakPolicy::Fail => {
2253 self.state.status = RunStatus::Failed;
2254 self.state
2255 .event("blind", format!("vendor text in a patch: {summary}"));
2256 self.state.leaks = leaks;
2257 self.state.save()?;
2258 self.settle_questions();
2259 bail!(
2260 "blind.on_leak = \"fail\" and vendor text reached a \
2261 judged patch: {summary}"
2262 );
2263 }
2264 LeakPolicy::Redact => self.state.event(
2265 "blind",
2266 format!("redacting vendor text for judging: {summary}"),
2267 ),
2268 LeakPolicy::Warn => self.state.event(
2269 "blind",
2270 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2271 ),
2272 }
2273 self.state.leaks = leaks;
2274 }
2275 }
2276
2277 if self.state.viable().is_empty() {
2278 if self.state.all_candidates_verified_noop() {
2279 self.state.status = RunStatus::VerifiedNoop;
2290 self.state.save()?;
2291 self.settle_questions();
2292 return Ok(());
2293 }
2294 self.state.status = RunStatus::Failed;
2295 self.state.save()?;
2296 self.settle_questions();
2297 bail!("no candidate produced a change; nothing to judge");
2298 }
2299 self.state.status = RunStatus::Judging;
2300 self.state.save()?;
2301 Ok(())
2302 }
2303
2304 async fn judge(&mut self) -> Result<()> {
2307 let run_id = self.state.id.clone();
2312 let prompts = self.state.config.prompts.clone();
2313 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2314 return Ok(());
2315 }
2316 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2317 if viable.len() == 1 {
2318 self.state.judge_skipped = true;
2325 self.state.event(
2326 "judge",
2327 format!(
2328 "only candidate {} produced a change; judging skipped",
2329 viable[0].label
2330 ),
2331 );
2332 self.state.save()?;
2333 return Ok(());
2334 }
2335 self.state.status = RunStatus::Judging;
2336
2337 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2338 let language = self.state.config.graph.language.clone();
2339 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2340 let sessions = self.state.config.graph.sessions;
2341 let artifacts = agent::artifacts_dir(&self.state.dir());
2342 let root = self.state.worktree_root();
2343 let base_short = short(&self.state.base_commit);
2344
2345 let mut jobs = Vec::new();
2346 let mut orders = Vec::new();
2347 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2348 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2349 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2350 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2351 let seat_key = format!("judge-{}", j + 1);
2352 let seat = self.seat(&seat_key, &spec.id);
2353 jobs.push(SeatJob {
2354 prompt: prompt::judge(
2355 &self.state.instruction,
2356 &views,
2357 self.roles.judges.len(),
2358 &base_short,
2359 &language,
2360 ),
2361 spec,
2362 seat,
2363 cwd: root.join(format!("judge-{}", j + 1)),
2364 timeout,
2365 allow_write: false,
2366 sessions,
2367 artifacts: artifacts.clone(),
2368 stem: format!("judge-{}", j + 1),
2369 });
2370 }
2371
2372 self.state.event(
2373 "judge",
2374 format!(
2375 "{} judges ranking {} candidates blind",
2376 jobs.len(),
2377 viable.len()
2378 ),
2379 );
2380 let labels_for_check = labels.clone();
2381 let mut quota_losses = Vec::new();
2382 let cache = self.state.config.cache_dir();
2383 let ctx = WaveCtx {
2384 run: &run_id,
2385 node: "judge",
2386 prompts: &prompts,
2387 cache: cache.as_deref(),
2388 round: None,
2389 };
2390 let results = ask_json_wave::<Ranking>(
2391 jobs,
2392 Arc::clone(&self.sem),
2393 self.state.config.graph.retries,
2394 &ctx,
2395 &mut quota_losses,
2396 &mut self.state,
2397 &move |r: &Ranking| r.validate(&labels_for_check),
2398 )
2399 .await;
2400 self.state.quota.extend(quota_losses);
2401
2402 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2403 let agent_id = seat.agent.clone();
2404 self.state.seats.insert(seat.key.clone(), seat);
2405 let mut record = Judgement {
2406 judge: j + 1,
2407 seat: format!("judge-{}", j + 1),
2408 agent: agent_id,
2409 ranking: Vec::new(),
2410 reasons: BTreeMap::new(),
2411 confidence: None,
2412 order: orders[j].clone(),
2413 failed: None,
2414 duration_ms: 0,
2415 };
2416 match res {
2417 Ok((ranking, out)) => {
2418 record.ranking = ranking.normalized();
2419 record.reasons = ranking.reasons;
2420 record.confidence = ranking.confidence;
2421 record.duration_ms = out.duration_ms;
2422 self.state.event(
2423 "judge",
2424 format!(
2425 "judge {} ranked {}",
2426 j + 1,
2427 record.ranking.iter().collect::<String>()
2428 ),
2429 );
2430 }
2431 Err(e) => {
2432 record.failed = Some(e.to_string());
2433 self.state
2434 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2435 }
2436 }
2437 self.state.judgements.push(record);
2438 self.state.save()?;
2439 }
2440 Ok(())
2441 }
2442
2443 async fn deliberate(&mut self) -> Result<()> {
2446 let run_id = self.state.id.clone();
2451 let prompts = self.state.config.prompts.clone();
2452 if !self.state.deliberation.is_empty() {
2453 return Ok(());
2454 }
2455 let tops: Vec<char> = self
2456 .state
2457 .judgements
2458 .iter()
2459 .filter_map(|j| j.ranking.first().copied())
2460 .collect();
2461 let rounds = self.state.config.graph.deliberate_rounds;
2462 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2463 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2464 self.state.event(
2465 "deliberate",
2466 format!("judges agreed on {} outright; no deliberation", tops[0]),
2467 );
2468 }
2469 self.state.status = RunStatus::Voting;
2470 self.state.save()?;
2471 return Ok(());
2472 }
2473
2474 self.state.status = RunStatus::Deliberating;
2475 self.state.event(
2476 "deliberate",
2477 format!(
2478 "split: first choices were {} — opening {rounds} round(s)",
2479 tops.iter().collect::<String>()
2480 ),
2481 );
2482
2483 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2484 let language = self.state.config.graph.language.clone();
2485 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2486 let sessions = self.state.config.graph.sessions;
2487 let artifacts = agent::artifacts_dir(&self.state.dir());
2488 let root = self.state.worktree_root();
2489 let base_short = short(&self.state.base_commit);
2490
2491 for round in 1..=rounds {
2495 let mut turns: Vec<DeliberationTurn> = Vec::new();
2496 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2497 if self.state.judgements[j].failed.is_some() {
2498 continue;
2499 }
2500 let seat_key = format!("judge-{}", j + 1);
2501 let mut seat = self.seat(&seat_key, &spec.id);
2502 let transcript = self.transcript(&turns, j);
2503 let context = if has_context(&spec, &seat, sessions) {
2504 None
2505 } else {
2506 Some(self.candidate_block(&viable, &base_short))
2507 };
2508 let text = prompt::deliberate(
2509 &self.state.instruction,
2510 context.as_deref(),
2511 &transcript,
2512 round,
2513 rounds,
2514 &language,
2515 );
2516 let job = SeatJob {
2517 spec,
2518 seat: seat.clone(),
2519 prompt: text,
2520 cwd: root.join(format!("judge-{}", j + 1)),
2521 timeout,
2522 allow_write: false,
2523 sessions,
2524 artifacts: artifacts.clone(),
2525 stem: format!("delib-{round}-judge-{}", j + 1),
2526 };
2527 let cache = self.state.config.cache_dir();
2528 let ctx = WaveCtx {
2529 run: &run_id,
2530 node: "deliberate",
2531 prompts: &prompts,
2532 cache: cache.as_deref(),
2533 round: None,
2534 };
2535 let (updated, out) =
2536 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2537 seat = updated;
2538 let agent_id = seat.agent.clone();
2539 let seat_key = seat.key.clone();
2540 self.state.seats.insert(seat.key.clone(), seat);
2541 let body = match out {
2542 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2543 AgentOutcome::Dropped(o) => {
2547 let why =
2548 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2549 "the CLI ended the stream without delivering its answer",
2550 );
2551 self.state.event(
2552 "deliberate",
2553 format!(
2554 "judge {} skipped: the CLI dropped the stream ({why})",
2555 j + 1
2556 ),
2557 );
2558 continue;
2559 }
2560 AgentOutcome::Quota(o) => {
2561 self.state.quota.push(QuotaLoss {
2562 seat: seat_key,
2563 node: "deliberate".to_owned(),
2564 at: Timestamp::now(),
2565 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2566 });
2567 self.state.event(
2568 "deliberate",
2569 format!("judge {} skipped: rate limited (quota)", j + 1),
2570 );
2571 continue;
2572 }
2573 AgentOutcome::Failed(e) => {
2574 self.state
2575 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2576 continue;
2577 }
2578 };
2579 let tentative = verdict::extract_json::<Position>(&body)
2580 .ok()
2581 .and_then(|p| p.tentative)
2582 .and_then(|s| s.trim().chars().next())
2583 .map(|c| c.to_ascii_uppercase());
2584 self.state.event(
2585 "deliberate",
2586 format!(
2587 "round {round}: judge {} now favours {}",
2588 j + 1,
2589 tentative.map_or("—".to_owned(), |c| c.to_string())
2590 ),
2591 );
2592 turns.push(DeliberationTurn {
2593 judge: j + 1,
2594 agent: agent_id,
2595 body: blind::sanitize_prose(&body, &self.state.config.blind),
2596 tentative,
2597 });
2598 }
2599 self.state
2600 .deliberation
2601 .push(DeliberationRound { round, turns });
2602 self.state.save()?;
2603 }
2604
2605 self.state.status = RunStatus::Voting;
2606 self.state.save()?;
2607 Ok(())
2608 }
2609
2610 async fn vote(&mut self) -> Result<()> {
2613 let run_id = self.state.id.clone();
2618 let prompts = self.state.config.prompts.clone();
2619 if !self.state.votes.is_empty() {
2620 return Ok(());
2621 }
2622 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2623 if viable.len() == 1 {
2624 return Ok(());
2625 }
2626 self.state.status = RunStatus::Voting;
2627
2628 let language = self.state.config.graph.language.clone();
2629 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2630 let sessions = self.state.config.graph.sessions;
2631 let artifacts = agent::artifacts_dir(&self.state.dir());
2632 let root = self.state.worktree_root();
2633 let base_short = short(&self.state.base_commit);
2634 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2635
2636 let mut jobs = Vec::new();
2637 let mut seats_at = Vec::new();
2638 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2639 if self
2640 .state
2641 .judgements
2642 .get(j)
2643 .is_some_and(|r| r.failed.is_some())
2644 {
2645 continue;
2646 }
2647 let seat_key = format!("judge-{}", j + 1);
2648 let seat = self.seat(&seat_key, &spec.id);
2649 let mut text = prompt::final_vote(&viable, &language);
2650 if !has_context(&spec, &seat, sessions) {
2651 text = format!(
2652 "{}\n\n# Candidates\n\n{}",
2653 text,
2654 self.candidate_block(&candidates, &base_short)
2655 );
2656 }
2657 jobs.push(SeatJob {
2658 spec,
2659 seat,
2660 prompt: text,
2661 cwd: root.join(format!("judge-{}", j + 1)),
2662 timeout,
2663 allow_write: false,
2664 sessions,
2665 artifacts: artifacts.clone(),
2666 stem: format!("vote-judge-{}", j + 1),
2667 });
2668 seats_at.push(j);
2669 }
2670
2671 self.state.event(
2672 "vote",
2673 format!(
2674 "collecting {} final votes one by one, privately",
2675 jobs.len()
2676 ),
2677 );
2678 let allowed = viable.clone();
2679 let mut quota_losses = Vec::new();
2680 let cache = self.state.config.cache_dir();
2681 let ctx = WaveCtx {
2682 run: &run_id,
2683 node: "vote",
2684 prompts: &prompts,
2685 cache: cache.as_deref(),
2686 round: None,
2687 };
2688 let results = ask_json_wave::<FinalVote>(
2689 jobs,
2690 Arc::clone(&self.sem),
2691 self.state.config.graph.retries,
2692 &ctx,
2693 &mut quota_losses,
2694 &mut self.state,
2695 &move |v: &FinalVote| match v.label() {
2696 Some(c) if allowed.contains(&c) => Ok(()),
2697 other => bail!("vote {other:?} is not one of {allowed:?}"),
2698 },
2699 )
2700 .await;
2701 self.state.quota.extend(quota_losses);
2702
2703 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2704 let agent_id = seat.agent.clone();
2705 self.state.seats.insert(seat.key.clone(), seat);
2706 let initial = self
2707 .state
2708 .judgements
2709 .get(j)
2710 .and_then(|r| r.ranking.first().copied());
2711 let mut record = VoteRecord {
2712 judge: j + 1,
2713 agent: agent_id,
2714 vote: None,
2715 reason: String::new(),
2716 changed: false,
2717 };
2718 match res {
2719 Ok((v, _)) => {
2720 record.vote = v.label();
2721 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2722 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2723 self.state.event(
2724 "vote",
2725 format!(
2726 "judge {} voted {}{}",
2727 j + 1,
2728 record.vote.unwrap_or('?'),
2729 if record.changed { " (changed)" } else { "" }
2730 ),
2731 );
2732 }
2733 Err(e) => {
2734 self.state
2735 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2736 }
2737 }
2738 self.state.votes.push(record);
2739 self.state.save()?;
2740 }
2741 Ok(())
2742 }
2743
2744 fn tally(&mut self) -> Result<()> {
2747 if self.state.tally.is_some() {
2748 return Ok(());
2749 }
2750 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2751 let tops: Vec<char> = self
2752 .state
2753 .judgements
2754 .iter()
2755 .filter_map(|j| j.ranking.first().copied())
2756 .collect();
2757 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2758
2759 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2762 let mut cast: Vec<char> = Vec::new();
2763 for (i, j) in self.state.judgements.iter().enumerate() {
2764 let vote = self
2765 .state
2766 .votes
2767 .iter()
2768 .find(|v| v.judge == i + 1)
2769 .and_then(|v| v.vote)
2770 .or_else(|| j.ranking.first().copied());
2771 if let Some(v) = vote {
2772 *first_choice.entry(v).or_insert(0) += 1;
2773 cast.push(v);
2774 }
2775 }
2776
2777 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2778 for j in &self.state.judgements {
2779 let n = j.ranking.len();
2780 for (pos, label) in j.ranking.iter().enumerate() {
2781 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2782 }
2783 }
2784
2785 let best = first_choice.values().copied().max().unwrap_or(0);
2786 let mut leaders: Vec<char> = first_choice
2787 .iter()
2788 .filter(|(_, v)| **v == best)
2789 .map(|(k, _)| *k)
2790 .collect();
2791 let mut tie_break = None;
2792 if leaders.len() > 1 {
2793 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2794 let borda_leaders: Vec<char> = leaders
2795 .iter()
2796 .copied()
2797 .filter(|l| borda[l] == top_borda)
2798 .collect();
2799 tie_break = Some(if borda_leaders.len() == 1 {
2800 format!(
2801 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2802 leaders.len()
2803 )
2804 } else {
2805 format!(
2806 "{} way tie on both first-choice votes and Borda points, broken by label order",
2807 leaders.len()
2808 )
2809 });
2810 leaders = borda_leaders;
2811 leaders.sort_unstable();
2812 }
2813 let winner = *leaders
2814 .first()
2815 .or(viable.first())
2816 .context("no candidate to declare a winner from")?;
2817
2818 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2819 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2820 let deliberated = !self.state.deliberation.is_empty();
2821
2822 let quota_seats: std::collections::BTreeSet<&str> =
2826 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2827 let mut present = 0usize;
2828 for (i, j) in self.state.judgements.iter().enumerate() {
2829 if quota_seats.contains(j.seat.as_str()) {
2830 continue;
2831 }
2832 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2833 let voted = self
2834 .state
2835 .votes
2836 .iter()
2837 .any(|v| v.judge == i + 1 && v.vote.is_some());
2838 if ranked || voted {
2839 present += 1;
2840 }
2841 }
2842 let needs_quorum = viable.len() > 1;
2848 let judges_total = if needs_quorum {
2849 self.roles.judges.len()
2850 } else {
2851 0
2852 };
2853 let quorum = if needs_quorum {
2854 judges_total / 2 + 1
2855 } else {
2856 0
2857 };
2858 let met_quorum = !needs_quorum || present >= quorum;
2859 let uncontested = (!needs_quorum).then(|| {
2860 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2861 });
2862
2863 self.state.event(
2864 "tally",
2865 match &uncontested {
2866 Some(reason) => format!("winner {winner} — {reason}"),
2867 None => format!(
2868 "winner {winner} — votes {} | initial {} | {} changed | \
2869 {present}/{judges_total} judges{}",
2870 first_choice
2871 .iter()
2872 .map(|(k, v)| format!("{k}:{v}"))
2873 .collect::<Vec<_>>()
2874 .join(" "),
2875 if unanimous_initial {
2876 "unanimous"
2877 } else {
2878 "split"
2879 },
2880 changed_votes,
2881 if met_quorum {
2882 String::new()
2883 } else {
2884 format!(" — below quorum ({quorum} required)")
2885 },
2886 ),
2887 },
2888 );
2889 if !met_quorum {
2890 self.state.event(
2891 "stall",
2892 format!(
2893 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2894 the run stops here, resumable"
2895 ),
2896 );
2897 }
2898 self.state.tally = Some(Tally {
2899 first_choice,
2900 borda,
2901 winner,
2902 rankings: tops.len(),
2903 unanimous_initial,
2904 deliberated,
2905 changed_votes,
2906 unanimous_final,
2907 tie_break,
2908 judges: judges_total,
2909 present,
2910 quorum,
2911 met_quorum,
2912 uncontested,
2913 });
2914 self.state.status = if met_quorum {
2915 RunStatus::Reviewing
2916 } else {
2917 RunStatus::Stalled
2918 };
2919 self.state.save()?;
2920 Ok(())
2921 }
2922
2923 #[allow(clippy::too_many_lines)]
2944 async fn recover_stall(&mut self) -> Result<bool> {
2945 let run_id = self.state.id.clone();
2950 let prompts = self.state.config.prompts.clone();
2951 let quota_seats: BTreeSet<&str> =
2956 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2957 let absent: Vec<String> = self
2958 .state
2959 .judgements
2960 .iter()
2961 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2962 .map(|j| j.seat.clone())
2963 .collect();
2964 if absent.is_empty() {
2965 return Ok(false);
2966 }
2967 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2968 if viable.len() <= 1 {
2969 return Ok(false);
2970 }
2971 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2972 let language = self.state.config.graph.language.clone();
2973 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2974 let sessions = self.state.config.graph.sessions;
2975 let artifacts = agent::artifacts_dir(&self.state.dir());
2976 let root = self.state.worktree_root();
2977 let base_short = short(&self.state.base_commit);
2978 let candidates: Vec<Candidate> = viable.clone();
2979
2980 let mut positions: Vec<usize> = absent
2982 .iter()
2983 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2984 .collect();
2985 if positions.is_empty() {
2986 return Ok(false);
2987 }
2988 positions.sort_unstable();
2989 positions.dedup();
2990
2991 let mut judge_jobs = Vec::new();
2993 for &j in &positions {
2994 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2995 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2996 let seat_key = format!("judge-{}", j + 1);
2997 let spec = self.roles.judges[j].clone();
2998 let seat = self.seat(&seat_key, &spec.id);
2999 judge_jobs.push(SeatJob {
3000 spec,
3001 seat,
3002 prompt: prompt::judge(
3003 &self.state.instruction,
3004 &views,
3005 self.roles.judges.len(),
3006 &base_short,
3007 &language,
3008 ),
3009 cwd: root.join(seat_key),
3010 timeout,
3011 allow_write: false,
3012 sessions,
3013 artifacts: artifacts.clone(),
3014 stem: format!("judge-{}-recover", j + 1),
3015 });
3016 }
3017
3018 let labels_for_check = labels.clone();
3019 let mut judge_losses = Vec::new();
3020 let retries = self.state.config.graph.retries;
3021 let cache = self.state.config.cache_dir();
3022 let ctx = WaveCtx {
3023 run: &run_id,
3024 node: "judge",
3025 prompts: &prompts,
3026 cache: cache.as_deref(),
3027 round: None,
3028 };
3029 let results = ask_json_wave::<Ranking>(
3030 judge_jobs,
3031 Arc::clone(&self.sem),
3032 retries,
3033 &ctx,
3034 &mut judge_losses,
3035 &mut self.state,
3036 &move |r: &Ranking| r.validate(&labels_for_check),
3037 )
3038 .await;
3039
3040 let mut recovered: BTreeSet<usize> = BTreeSet::new();
3042 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
3043 self.state.seats.insert(seat.key.clone(), seat);
3044 let record = &mut self.state.judgements[j];
3045 match res {
3046 Ok((ranking, out)) => {
3047 record.ranking = ranking.normalized();
3048 record.reasons = ranking.reasons;
3049 record.confidence = ranking.confidence;
3050 record.failed = None;
3051 record.duration_ms = out.duration_ms;
3052 recovered.insert(j);
3053 self.state.event(
3054 "recover",
3055 format!("judge {} ranked again after the limit", j + 1),
3056 );
3057 }
3058 Err(e) => {
3059 self.state
3060 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
3061 }
3062 }
3063 }
3064
3065 let mut vote_jobs = Vec::new();
3067 let mut vote_pos: Vec<usize> = Vec::new();
3068 for &j in &recovered {
3069 let seat_key = format!("judge-{}", j + 1);
3070 let spec = self.roles.judges[j].clone();
3071 let seat = self.seat(&seat_key, &spec.id);
3072 let mut text = prompt::final_vote(&labels, &language);
3073 if !has_context(&spec, &seat, sessions) {
3074 text = format!(
3075 "{}\n\n# Candidates\n\n{}",
3076 text,
3077 self.candidate_block(&candidates, &base_short)
3078 );
3079 }
3080 vote_jobs.push(SeatJob {
3081 spec,
3082 seat,
3083 prompt: text,
3084 cwd: root.join(seat_key),
3085 timeout,
3086 allow_write: false,
3087 sessions,
3088 artifacts: artifacts.clone(),
3089 stem: format!("vote-judge-{}-recover", j + 1),
3090 });
3091 vote_pos.push(j);
3092 }
3093 let allowed = labels.clone();
3094 let mut vote_losses = Vec::new();
3095 let vote_retries = self.state.config.graph.retries;
3096 let vote_cache = self.state.config.cache_dir();
3097 let ctx = WaveCtx {
3098 run: &run_id,
3099 node: "vote",
3100 prompts: &prompts,
3101 cache: vote_cache.as_deref(),
3102 round: None,
3103 };
3104 let votes = ask_json_wave::<FinalVote>(
3105 vote_jobs,
3106 Arc::clone(&self.sem),
3107 vote_retries,
3108 &ctx,
3109 &mut vote_losses,
3110 &mut self.state,
3111 &move |v: &FinalVote| match v.label() {
3112 Some(c) if allowed.contains(&c) => Ok(()),
3113 other => bail!("vote {other:?} is not one of {allowed:?}"),
3114 },
3115 )
3116 .await;
3117 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
3118 let agent_id = seat.agent.clone();
3119 self.state.seats.insert(seat.key.clone(), seat);
3120 match res {
3121 Ok((v, _)) => {
3122 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
3123 rec.vote = v.label();
3124 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
3125 } else {
3126 self.state.votes.push(VoteRecord {
3127 judge: j + 1,
3128 agent: agent_id,
3129 vote: v.label(),
3130 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
3131 changed: false,
3132 });
3133 }
3134 self.state.event(
3135 "recover",
3136 format!("judge {} voted again after the limit", j + 1),
3137 );
3138 }
3139 Err(e) => {
3140 self.state
3141 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
3142 }
3143 }
3144 }
3145
3146 let recovered_keys: BTreeSet<String> = recovered
3150 .iter()
3151 .map(|&j| format!("judge-{}", j + 1))
3152 .collect();
3153 self.state
3154 .quota
3155 .retain(|q| !recovered_keys.contains(&q.seat));
3156 for loss in judge_losses.into_iter().chain(vote_losses) {
3160 if recovered_keys.contains(&loss.seat) {
3161 continue;
3162 }
3163 self.state.quota.retain(|q| q.seat != loss.seat);
3164 self.state.quota.push(loss);
3165 }
3166
3167 self.state.tally = None;
3169 self.tally()?;
3170 Ok(self
3171 .state
3172 .tally
3173 .as_ref()
3174 .map(|t| t.met_quorum)
3175 .unwrap_or(false))
3176 }
3177
3178 async fn fold_losers(&mut self) -> Result<()> {
3181 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
3182 return Ok(());
3183 };
3184 let repo = self.state.repo.clone();
3185 let mut folded = Vec::new();
3186 for i in 0..self.state.candidates.len() {
3187 let c = &self.state.candidates[i];
3188 if c.label == winner || c.folded {
3189 continue;
3190 }
3191 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
3192 git::worktree_remove(&repo, &wt).await.ok();
3193 git::branch_delete(&repo, &branch).await.ok();
3194 self.state.candidates[i].folded = true;
3195 folded.push(label.to_string());
3196 }
3197 let root = self.state.worktree_root();
3199 for j in 1..=self.roles.judges.len() {
3200 let wt = root.join(format!("judge-{j}"));
3201 if wt.exists() {
3202 git::worktree_remove(&repo, &wt).await.ok();
3203 }
3204 }
3205 if self.state.config.graph.advise {
3208 for k in 1..=self.state.config.graph.advisors {
3209 let wt = root.join(format!("advisor-{k}"));
3210 if wt.exists() {
3211 git::worktree_remove(&repo, &wt).await.ok();
3212 }
3213 }
3214 }
3215 if !folded.is_empty() {
3216 self.state
3217 .event("fold", format!("folded candidates {}", folded.join(", ")));
3218 self.state.save()?;
3219 }
3220 Ok(())
3221 }
3222
3223 async fn sync_to_base(&mut self) -> Result<()> {
3263 if self
3264 .state
3265 .base_sync
3266 .as_ref()
3267 .is_some_and(|s| s.conflict.is_some())
3268 {
3269 return Ok(());
3270 }
3271 let Some(winner) = self.state.winner().cloned() else {
3272 return Ok(());
3273 };
3274
3275 let repo = self.state.repo.clone();
3276 let remote = self.state.config.merge.remote.clone();
3277 let base_branch = self.state.base_branch.clone();
3278 let tracking = format!("{remote}/{base_branch}");
3279
3280 git::fetch(&repo, &remote, &base_branch).await.ok();
3281 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3285 return Ok(());
3286 };
3287
3288 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3289 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3290 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3291
3292 if behind == 0 {
3293 if let Some(from) = self
3302 .state
3303 .rebase_fixes
3304 .iter()
3305 .rev()
3306 .find_map(|r| r.from.clone())
3307 && from != head
3308 && git::git_raw(&winner.worktree, &["diff", "--quiet", &from])
3309 .await
3310 .is_ok_and(|o| o.ok())
3311 {
3312 git::sync_to_head(&winner.worktree).await?;
3313 }
3314 self.state.base_sync = Some(BaseSync {
3315 tip,
3316 behind: 0,
3317 attempts,
3318 conflict: None,
3319 });
3320 self.state.save()?;
3321 return Ok(());
3322 }
3323
3324 if attempts >= BASE_SYNC_ROUNDS {
3325 let why = format!(
3326 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3327 rebase(s); rebasing again would only race it",
3328 winner.branch
3329 );
3330 self.state.status = RunStatus::Blocked;
3331 self.state.base_sync = Some(BaseSync {
3332 tip,
3333 behind,
3334 attempts,
3335 conflict: Some(why.clone()),
3336 });
3337 self.state.event("land", why);
3338 self.state.save()?;
3339 return Ok(());
3340 }
3341
3342 self.state.event(
3343 "land",
3344 format!(
3345 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3346 winner.branch
3347 ),
3348 );
3349 self.state.save()?;
3350
3351 let branch_tracking = format!("{remote}/{}", winner.branch);
3357 let fetched_branch = git::fetch(&repo, &remote, &winner.branch).await;
3358 let remote_tip = if matches!(&fetched_branch, Ok(o) if o.ok()) {
3359 git::rev_parse(&repo, &branch_tracking).await.ok()
3360 } else {
3361 None
3362 };
3363 if let Some(theirs) = &remote_tip
3367 && !git::is_ancestor(&repo, theirs, &head).await
3368 && !crate::reconcile::origin_missing(&repo, &head, theirs)
3369 .await
3370 .is_ok_and(|missing| missing.is_empty())
3371 {
3372 let why = format!(
3373 "{branch_tracking} ({}) has commits {} does not contain; not rebasing over \
3374 them",
3375 short(theirs),
3376 winner.branch
3377 );
3378 self.state.status = RunStatus::Blocked;
3379 self.state.base_sync = Some(BaseSync {
3380 tip,
3381 behind,
3382 attempts,
3383 conflict: Some(why.clone()),
3384 });
3385 self.state.event("land", why);
3386 self.state.save()?;
3387 return Ok(());
3388 }
3389
3390 let scratch = self.state.dir().join("base-sync");
3391 let rebased = match crate::rebase::rebase_with_fixer(
3392 &mut self.state,
3393 &scratch,
3394 &winner.branch,
3395 &tracking,
3396 )
3397 .await
3398 {
3399 Ok(crate::rebase::Rebased::Applied) => Ok(None),
3400 Ok(crate::rebase::Rebased::Stopped(why)) => Ok(Some(why)),
3401 Err(e) => Err(e),
3402 };
3403 let attempts = attempts + 1;
3404 match rebased {
3405 Ok(None) => {
3406 git::sync_to_head(&winner.worktree).await?;
3410 let mut conflict = None;
3411 if let Some(pinned) = &remote_tip {
3412 let pushed = git::push_pinned(&repo, &remote, &winner.branch, pinned).await;
3413 match pushed {
3414 Ok(o) if o.ok() => self.state.event(
3415 "land",
3416 format!("pushed rebased {} to {remote}", winner.branch),
3417 ),
3418 Ok(o) => {
3419 conflict = Some(format!(
3420 "rebased {} locally but {remote} refused the push (it moved since {}; someone may have pushed): {}",
3421 winner.branch,
3422 short(pinned),
3423 o.stderr.chars().take(600).collect::<String>()
3424 ));
3425 }
3426 Err(e) => {
3427 conflict = Some(format!(
3428 "rebased {} locally but could not push it: {e:#}",
3429 winner.branch
3430 ));
3431 }
3432 }
3433 }
3434 if let Some(why) = &conflict {
3435 self.state.status = RunStatus::Blocked;
3436 self.state.event("land", why.clone());
3437 }
3438 self.state.base_sync = Some(BaseSync {
3439 tip: tip.clone(),
3440 behind: 0,
3441 attempts,
3442 conflict,
3443 });
3444 self.state
3445 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3446 }
3447 Ok(Some(conflict)) => {
3448 let why = format!(
3449 "{} conflicts with {tracking} and did not rebase: {}",
3450 winner.branch,
3451 conflict.chars().take(600).collect::<String>()
3452 );
3453 self.state.status = RunStatus::Blocked;
3454 self.state.base_sync = Some(BaseSync {
3455 tip,
3456 behind,
3457 attempts,
3458 conflict: Some(why.clone()),
3459 });
3460 self.state.event("land", why);
3461 }
3462 Err(e) => {
3463 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3464 self.state.status = RunStatus::Blocked;
3465 self.state.base_sync = Some(BaseSync {
3466 tip,
3467 behind,
3468 attempts,
3469 conflict: Some(why.clone()),
3470 });
3471 self.state.event("land", why);
3472 }
3473 }
3474 self.state.save()?;
3475 Ok(())
3476 }
3477
3478 fn landing_base(&self) -> String {
3488 self.state
3489 .base_sync
3490 .as_ref()
3491 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3492 }
3493
3494 pub async fn fix_selected(
3527 &mut self,
3528 ids: &[String],
3529 reason: &str,
3530 allow_stale: bool,
3531 ) -> Result<()> {
3532 let reason = reason.trim();
3533 if reason.is_empty() {
3534 bail!("a fix request needs a reason — that is the operator's own record of why");
3535 }
3536 if ids.is_empty() {
3537 bail!("no finding id given");
3538 }
3539 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3540 bail!(
3541 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3542 has already concluded — can be given a targeted fix. A run still \
3543 in progress should simply be resumed; a `merged` run's branch has \
3544 already landed, so its answer is a fresh `magi review <branch>`, \
3545 not reopening this run's own record",
3546 self.state.id,
3547 self.state.status.as_str()
3548 );
3549 }
3550 let Some(winner) = self.state.winner().cloned() else {
3551 bail!("run {} has no winning candidate to fix", self.state.id);
3552 };
3553 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3554 bail!(
3555 "branch `{}` no longer exists; this run cannot be extended",
3556 winner.branch
3557 );
3558 }
3559 let home = crate::run::home();
3560 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3561 bail!(
3562 "run {} is currently being worked on by another magi process",
3563 self.state.id
3564 );
3565 }
3566 let _claim = FixClaim::acquire(&self.state.dir())?;
3572
3573 let mut seen = BTreeSet::new();
3577 let mut findings = Vec::new();
3578 let mut missing = Vec::new();
3579 for id in ids {
3580 if !seen.insert(id.clone()) {
3581 continue;
3582 }
3583 match self.state.finding(id) {
3584 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3585 id: f.id.clone(),
3586 severity: f.severity,
3587 reviewer_vote: rec.vote,
3588 round: round.round,
3589 round_head: round.head.clone(),
3590 reviewer: rec.reviewer,
3591 agent: rec.agent.clone(),
3592 file: f.file.clone(),
3593 line: f.line,
3594 title: f.title.clone(),
3595 detail: f.detail.clone(),
3596 outcome: OperatorFixOutcome::Pending,
3597 }),
3598 None => missing.push(id.clone()),
3599 }
3600 }
3601 if !missing.is_empty() {
3602 bail!(
3603 "unknown finding id(s): {}; nothing was changed",
3604 missing.join(", ")
3605 );
3606 }
3607
3608 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3609 let stale_details: Vec<(String, String)> = findings
3610 .iter()
3611 .filter(|f| f.round_head != head_at_request)
3612 .map(|f| (f.id.clone(), f.round_head.clone()))
3613 .collect();
3614 let stale = !stale_details.is_empty();
3615 if stale && !allow_stale {
3616 bail!(
3617 "the branch has moved since some finding(s) were raised — {} — now \
3618 at {}; pass --allow-stale to fix anyway, or re-run review first",
3619 stale_details
3620 .iter()
3621 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3622 .collect::<Vec<_>>()
3623 .join(", "),
3624 short(&head_at_request)
3625 );
3626 }
3627
3628 let request = OperatorFixRequest {
3629 requested_at: Timestamp::now(),
3630 reason: reason.to_owned(),
3631 findings,
3632 head_at_request: head_at_request.clone(),
3633 allow_stale,
3634 stale,
3635 fix: None,
3636 result_head: None,
3637 follow_up_review_run: None,
3638 };
3639 self.state.event(
3640 "fix",
3641 format!(
3642 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3643 request.findings.len(),
3644 request
3645 .findings
3646 .iter()
3647 .map(|f| f.id.as_str())
3648 .collect::<Vec<_>>()
3649 .join(", "),
3650 ),
3651 );
3652 self.state.operator_fixes.push(request);
3659 self.state.save()?;
3660 let request_index = self.state.operator_fixes.len() - 1;
3661
3662 if winner.worktree.exists() {
3671 let dirty = git::git(
3674 &winner.worktree,
3675 &["status", "--porcelain", "--untracked-files=all"],
3676 )
3677 .await?;
3678 let only_withheld = dirty.lines().all(|l| {
3679 l.strip_prefix("?? ")
3680 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3681 });
3682 if !only_withheld {
3683 bail!(
3684 "`{}` has uncommitted changes; refusing to touch it — commit or \
3685 discard them first",
3686 winner.worktree.display()
3687 );
3688 }
3689 git::worktree_remove(&self.state.repo, &winner.worktree)
3690 .await
3691 .ok();
3692 }
3693 let fix_worktree = self.state.worktree_root().join("operator-fix");
3694 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3695 git::git(
3696 &self.state.repo,
3697 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3698 )
3699 .await
3700 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3701 if !git::is_clean(&fix_worktree).await? {
3702 git::worktree_remove(&self.state.repo, &fix_worktree)
3703 .await
3704 .ok();
3705 bail!(
3706 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3707 winner.branch
3708 );
3709 }
3710
3711 let run_id = self.state.id.clone();
3712 let prompts = self.state.config.prompts.clone();
3713 let language = self.state.config.graph.language.clone();
3714 let sessions = self.state.config.graph.sessions;
3715 let artifacts = agent::artifacts_dir(&self.state.dir());
3716 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3717 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3718 _ => (
3719 self.state
3720 .config
3721 .agent(&winner.agent)
3722 .cloned()
3723 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3724 format!("impl-{}", winner.label),
3725 ),
3726 };
3727 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3728 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3729 .findings
3730 .iter()
3731 .map(|f| Finding {
3732 id: f.id.clone(),
3733 severity: f.severity,
3734 file: f.file.clone(),
3735 line: f.line,
3736 title: f.title.clone(),
3737 detail: f.detail.clone(),
3738 })
3739 .collect();
3740 let job = SeatJob {
3741 prompt: prompt::operator_fix(
3742 &self.state.instruction,
3743 &finding_list,
3744 reason,
3745 &stale_details,
3746 &head_at_request,
3747 &language,
3748 ),
3749 spec: fix_spec.clone(),
3750 seat,
3751 cwd: fix_worktree.clone(),
3752 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3753 allow_write: true,
3754 sessions,
3755 artifacts: artifacts.clone(),
3756 stem: "operator-fix".to_owned(),
3757 };
3758 let cache = self.state.config.cache_dir();
3759 let ctx = WaveCtx {
3760 run: &run_id,
3761 node: "fix",
3762 prompts: &prompts,
3763 cache: cache.as_deref(),
3764 round: None,
3765 };
3766 let (seat, out) =
3767 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3768 let agent_id = seat.agent.clone();
3769
3770 let mut fix = FixRecord {
3771 agent: agent_id,
3772 addressed: Vec::new(),
3773 rejected: Vec::new(),
3774 notes: String::new(),
3775 committed: false,
3776 failed: None,
3777 duration_ms: 0,
3778 continuation: None,
3779 };
3780 let mut final_seat = seat.clone();
3781 match out {
3782 AgentOutcome::Ok(o) => {
3783 fix.duration_ms = o.duration_ms;
3784 let parsed = verdict::extract_json::<FixReport>(&o.text);
3785 let incomplete_reason = match &parsed {
3786 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3787 "the reply parsed, but it reported a command whose own CLI \
3788 never confirmed an exit status"
3789 .to_owned(),
3790 ),
3791 Ok(_) => None,
3792 Err(e) => Some(e.to_string()),
3793 };
3794 match incomplete_reason {
3795 None => {
3796 let report = parsed.expect("checked Ok above");
3797 fix.addressed = report.addressed;
3798 fix.rejected = report.rejected;
3799 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3800 }
3801 Some(reason) => {
3802 let (resumed_seat, resolved, failure, cont) = self
3803 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3804 .await;
3805 fix.duration_ms += cont.cumulative_wait_ms;
3806 fix.continuation = Some(cont);
3807 final_seat = resumed_seat;
3808 match resolved {
3809 Some(report) => {
3810 fix.addressed = report.addressed;
3811 fix.rejected = report.rejected;
3812 fix.notes =
3813 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3814 }
3815 None => fix.failed = failure,
3816 }
3817 }
3818 }
3819 }
3820 AgentOutcome::Dropped(o) => {
3821 fix.duration_ms = o.duration_ms;
3822 let why = o
3823 .dropped
3824 .as_ref()
3825 .map(|d| d.why.as_str())
3826 .unwrap_or("the CLI ended the stream without delivering its answer");
3827 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3828 }
3829 AgentOutcome::Quota(o) => {
3830 self.state.quota.push(QuotaLoss {
3831 seat: final_seat.key.clone(),
3832 node: "fix".to_owned(),
3833 at: Timestamp::now(),
3834 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3835 });
3836 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3837 }
3838 AgentOutcome::Failed(e) => fix.failed = Some(e),
3839 }
3840 if fix.continuation.is_none() {
3841 fix.continuation = Some(ContinuationRecord::not_needed());
3842 }
3843 self.state.seats.insert(final_seat.key.clone(), final_seat);
3844
3845 let rescue_message = format!(
3846 "magi: operator-selected fix ({}) (uncommitted work)",
3847 self.state.operator_fixes[request_index]
3848 .findings
3849 .iter()
3850 .map(|f| f.id.as_str())
3851 .collect::<Vec<_>>()
3852 .join(", ")
3853 );
3854 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3855 self.state.note_withheld("fix", &r.withheld);
3856 }
3857 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3858 fix.committed = after != head_at_request;
3859 git::worktree_remove(&self.state.repo, &fix_worktree)
3860 .await
3861 .ok();
3862
3863 self.state.event(
3864 "fix",
3865 match &fix.failed {
3866 Some(reason) => format!(
3867 "operator fix: adoption report was lost ({reason}); {}",
3868 if fix.committed {
3869 "committed"
3870 } else {
3871 "NO new commit"
3872 }
3873 ),
3874 None => format!(
3875 "operator fix: {} addressed, {} rejected, {}",
3876 fix.addressed.len(),
3877 fix.rejected.len(),
3878 if fix.committed {
3879 "committed"
3880 } else {
3881 "NO new commit"
3882 }
3883 ),
3884 },
3885 );
3886
3887 for f in &mut self.state.operator_fixes[request_index].findings {
3894 f.outcome = if fix.failed.is_some() {
3895 OperatorFixOutcome::Unreported
3896 } else if fix.addressed.contains(&f.id) {
3897 OperatorFixOutcome::Addressed
3898 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3899 OperatorFixOutcome::Rejected { why: r.why.clone() }
3900 } else {
3901 OperatorFixOutcome::Unreported
3902 };
3903 }
3904
3905 let committed = fix.committed;
3906 if committed {
3907 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3908 }
3909 self.state.operator_fixes[request_index].fix = Some(fix);
3910 self.state.save()?;
3913
3914 if committed {
3915 self.state.event(
3916 "fix",
3917 format!(
3918 "operator fix committed {}; opening a follow-up review-only run",
3919 short(&after)
3920 ),
3921 );
3922 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3923 Ok(mut follow_up) => {
3924 follow_up.state.event(
3925 "start",
3926 format!(
3927 "requested by an operator fix on run {} for finding(s) {}",
3928 self.state.id,
3929 self.state.operator_fixes[request_index]
3930 .findings
3931 .iter()
3932 .map(|f| f.id.as_str())
3933 .collect::<Vec<_>>()
3934 .join(", "),
3935 ),
3936 );
3937 follow_up.state.save()?;
3938 let follow_up_id = follow_up.state.id.clone();
3939 if let Err(e) = follow_up.execute().await {
3940 self.state.event(
3941 "fix",
3942 format!(
3943 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3944 ),
3945 );
3946 }
3947 self.state.operator_fixes[request_index].follow_up_review_run =
3948 Some(follow_up_id);
3949 }
3950 Err(e) => {
3951 self.state.event(
3952 "fix",
3953 format!("committed the fix but could not open a follow-up review: {e:#}"),
3954 );
3955 }
3956 }
3957 self.state.save()?;
3958 }
3959
3960 Ok(())
3961 }
3962
3963 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3970 match &self.roles.fixer {
3971 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3972 _ => (
3973 self.state
3974 .config
3975 .agent(&winner.agent)
3976 .cloned()
3977 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3978 format!("impl-{}", winner.label),
3979 ),
3980 }
3981 }
3982
3983 async fn review_loop(&mut self) -> Result<()> {
3984 if self
3989 .state
3990 .base_sync
3991 .as_ref()
3992 .is_some_and(|s| s.conflict.is_some())
3993 {
3994 return Ok(());
3995 }
3996 let run_id = self.state.id.clone();
4001 let prompts = self.state.config.prompts.clone();
4002 let Some(winner) = self.state.winner().cloned() else {
4003 return Ok(());
4004 };
4005 let max_rounds = self.state.config.graph.review_rounds;
4006 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
4016 self.state.status = status;
4017 self.state.save()?;
4018 return Ok(());
4019 }
4020 self.state.status = RunStatus::Reviewing;
4021 if self
4031 .state
4032 .reviews
4033 .last()
4034 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
4035 {
4036 let shell = self.state.config.shell();
4037 return self
4038 .stop_reviewing(
4039 "the last round's own verification never resolved",
4040 &shell,
4041 &winner.worktree,
4042 )
4043 .await;
4044 }
4045
4046 let repo = self.state.repo.clone();
4047 let root = self.state.worktree_root();
4048 let language = self.state.config.graph.language.clone();
4049 let sessions = self.state.config.graph.sessions;
4050 let artifacts = agent::artifacts_dir(&self.state.dir());
4051 let base = self.landing_base();
4052 let base_short = short(&base);
4053 let reviewers = self.roles.reviewers.clone();
4054 let shell = self.state.config.shell();
4055
4056 for round in (self.state.reviews.len() + 1)..=max_rounds {
4057 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4058 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
4059 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
4060 let prev_verification = self
4069 .state
4070 .reviews
4071 .last()
4072 .and_then(|r| r.verification_summary(&head));
4073
4074 let mut jobs = Vec::new();
4078 for (r, spec) in reviewers.iter().cloned().enumerate() {
4079 let wt = root.join(format!("review-{}", r + 1));
4080 if wt.exists() {
4081 git::reset_detached(&wt, &head).await?;
4082 } else {
4083 git::worktree_add_detached(&repo, &wt, &head).await?;
4084 }
4085 let seat_key = format!("review-{}", r + 1);
4086 let seat = self.seat(&seat_key, &spec.id);
4087 jobs.push(SeatJob {
4088 prompt: prompt::review(&prompt::ReviewCtx {
4089 instruction: &self.state.instruction,
4090 branch: &winner.branch,
4091 base_short: &base_short,
4092 stat: &stat,
4093 patch: &patch,
4094 verification: prev_verification.as_ref(),
4095 reviewers: reviewers.len(),
4096 round,
4097 rounds: max_rounds,
4098 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
4101 lens: Lens::for_seat(r),
4102 language: &language,
4103 }),
4104 spec,
4105 seat,
4106 cwd: wt,
4107 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4108 allow_write: false,
4109 sessions,
4110 artifacts: artifacts.clone(),
4111 stem: format!("review-{round}-{}", r + 1),
4112 });
4113 }
4114
4115 self.state.event(
4116 "review",
4117 format!(
4118 "round {round}: {} reviewers on {}",
4119 jobs.len(),
4120 short(&head)
4121 ),
4122 );
4123 let mut quota_losses = Vec::new();
4124 let review_retries = self.state.config.graph.retries;
4125 let review_cache = self.state.config.cache_dir();
4126 let ctx = WaveCtx {
4127 run: &run_id,
4128 node: "review",
4129 prompts: &prompts,
4130 cache: review_cache.as_deref(),
4131 round: Some(round),
4132 };
4133 let results = ask_json_wave::<Review>(
4134 jobs,
4135 Arc::clone(&self.sem),
4136 review_retries,
4137 &ctx,
4138 &mut quota_losses,
4139 &mut self.state,
4140 &|_: &Review| Ok(()),
4141 )
4142 .await;
4143 let round_quota_missing = quota_losses.len();
4147 self.state.quota.extend(quota_losses);
4148
4149 let mut records = Vec::new();
4150 let mut all_findings = Vec::new();
4151 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
4152 let agent_id = seat.agent.clone();
4153 self.state.seats.insert(seat.key.clone(), seat);
4154 let mut record = ReviewRecord {
4155 reviewer: r + 1,
4156 agent: agent_id,
4157 summary: String::new(),
4158 findings: Vec::new(),
4159 vote: None,
4160 failed: None,
4161 duration_ms: 0,
4162 attempts,
4168 };
4169 match res {
4170 Ok((review, out)) => {
4171 record.summary =
4179 blind::sanitize_prose(&review.summary, &self.state.config.blind);
4180 record.vote = Some(review.vote);
4181 record.duration_ms = out.duration_ms;
4182 for (n, mut f) in review.findings.into_iter().enumerate() {
4183 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
4186 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
4187 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
4188 f.file = f
4194 .file
4195 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
4196 all_findings.push(f.clone());
4197 record.findings.push(f);
4198 }
4199 self.state.event(
4200 "review",
4201 format!(
4202 "round {round}: reviewer {} voted {} with {} finding(s)",
4203 r + 1,
4204 review.vote.label(),
4205 record.findings.len()
4206 ),
4207 );
4208 }
4209 Err(e) => {
4210 record.failed = Some(e.to_string());
4211 self.state.event(
4212 "review",
4213 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
4214 );
4215 }
4216 }
4217 records.push(record);
4218 }
4219
4220 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
4227 let vote_split =
4228 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
4229 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
4230 if vote_split {
4231 self.state.event(
4232 "review",
4233 format!(
4234 "round {round}: votes split ({}) — one round of reconsideration",
4235 initial_votes
4236 .iter()
4237 .map(|v| v.label())
4238 .collect::<Vec<_>>()
4239 .join(", ")
4240 ),
4241 );
4242 let panel: Vec<ReviewSeatReport<'_>> = records
4245 .iter()
4246 .filter_map(|r| {
4247 r.vote.map(|vote| ReviewSeatReport {
4248 reviewer: r.reviewer,
4249 vote,
4250 summary: &r.summary,
4251 findings: &r.findings,
4252 })
4253 })
4254 .collect();
4255
4256 let mut jobs = Vec::new();
4257 let mut seats_at = Vec::new();
4258 for (r, spec) in reviewers.iter().cloned().enumerate() {
4259 if records[r].vote.is_none() {
4263 continue;
4264 }
4265 let wt = root.join(format!("review-{}", r + 1));
4266 let seat_key = format!("review-{}", r + 1);
4267 let seat = self.seat(&seat_key, &spec.id);
4268 let patch_ctx = if has_context(&spec, &seat, sessions) {
4273 None
4274 } else {
4275 Some(ReviewPatch {
4276 branch: &winner.branch,
4277 base_short: &base_short,
4278 stat: &stat,
4279 patch: &patch,
4280 })
4281 };
4282 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
4283 instruction: &self.state.instruction,
4284 reviewer: r + 1,
4285 lens: Lens::for_seat(r),
4286 panel: &panel,
4287 patch: patch_ctx,
4288 round,
4289 rounds: max_rounds,
4290 language: &language,
4291 });
4292 jobs.push(SeatJob {
4293 prompt,
4294 spec,
4295 seat,
4296 cwd: wt,
4297 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4298 allow_write: false,
4299 sessions,
4300 artifacts: artifacts.clone(),
4301 stem: format!("review-{round}-reconsider-{}", r + 1),
4302 });
4303 seats_at.push(r);
4304 }
4305
4306 let mut recon_quota_losses = Vec::new();
4307 let recon_cache = self.state.config.cache_dir();
4308 let recon_ctx = WaveCtx {
4309 run: &run_id,
4310 node: "review",
4311 prompts: &prompts,
4312 cache: recon_cache.as_deref(),
4313 round: Some(round),
4314 };
4315 let recon_results = ask_json_wave::<ReviewRevote>(
4316 jobs,
4317 Arc::clone(&self.sem),
4318 review_retries,
4319 &recon_ctx,
4320 &mut recon_quota_losses,
4321 &mut self.state,
4322 &|_: &ReviewRevote| Ok(()),
4323 )
4324 .await;
4325 self.state.quota.extend(recon_quota_losses);
4326
4327 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
4328 let agent_id = seat.agent.clone();
4329 self.state.seats.insert(seat.key.clone(), seat);
4330 let mut rec = ReviewRevoteRecord {
4331 reviewer: r + 1,
4332 agent: agent_id,
4333 vote: None,
4334 reason: String::new(),
4335 failed: None,
4336 };
4337 match res {
4338 Ok((rv, _)) => {
4339 rec.vote = Some(rv.vote);
4340 rec.reason =
4341 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
4342 self.state.event(
4343 "review",
4344 format!(
4345 "round {round}: reviewer {} revoted {}",
4346 r + 1,
4347 rv.vote.label()
4348 ),
4349 );
4350 }
4351 Err(e) => {
4352 rec.failed = Some(e.to_string());
4353 self.state.event(
4354 "review",
4355 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4356 );
4357 }
4358 }
4359 reconsideration.push(rec);
4360 }
4361 } else if initial_votes.len() > 1 {
4362 self.state.event(
4363 "review",
4364 format!(
4365 "round {round}: votes agreed ({}) — no reconsideration",
4366 initial_votes[0].label()
4367 ),
4368 );
4369 }
4370
4371 let final_votes: Vec<ReviewVote> = records
4375 .iter()
4376 .filter_map(|r| {
4377 reconsideration
4378 .iter()
4379 .find(|rv| rv.reviewer == r.reviewer)
4380 .and_then(|rv| rv.vote)
4381 .or(r.vote)
4382 })
4383 .collect();
4384 let round_verdict = ReviewVote::worst(final_votes);
4385
4386 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4387 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4388 let defer_e2e =
4399 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4400 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4401 let reason =
4402 format!("{blocking} blocking finding(s) already required a fix this round");
4403 self.state.event(
4404 "verify",
4405 format!(
4406 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4407 {}); it will run once a round has none left",
4408 short(&head)
4409 ),
4410 );
4411 (Vec::new(), false, true, Some(reason))
4412 } else {
4413 let e2e_commands = self.state.config.verify.e2e.clone();
4414 let cache_dir = self.state.config.cache_dir();
4415 let context = format!("round {round}");
4416 let (e2e, verify_retried) = with_cache_lease(
4417 &mut self.state,
4418 cache_dir.as_deref(),
4419 "e2e",
4420 "e2e",
4421 &winner.worktree,
4422 &head,
4423 verify_timeout,
4424 &context,
4425 |state, budget| {
4426 let shell = shell.clone();
4427 let e2e_commands = e2e_commands.clone();
4428 let worktree = winner.worktree.clone();
4429 let context = context.clone();
4430 async move {
4431 run_e2e_with_retry(
4432 state,
4433 &shell,
4434 &e2e_commands,
4435 &worktree,
4436 budget,
4437 &context,
4438 )
4439 .await
4440 }
4441 },
4442 )
4443 .await;
4444 (e2e, verify_retried, false, None)
4445 };
4446
4447 let expected = records.len();
4448 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4449 let incomplete = answered < expected;
4450 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4451 let policy = self.state.config.graph.incomplete_review;
4452 let clean = round_is_clean(
4453 blocking,
4454 e2e_ok,
4455 answered,
4456 expected,
4457 round_quota_missing,
4458 policy,
4459 );
4460
4461 let mut round_record = ReviewRound {
4462 round,
4463 head: head.clone(),
4464 verified_head: None,
4465 verified_at: None,
4466 reviews: records,
4467 e2e,
4468 verify_retried,
4469 e2e_deferred,
4470 e2e_defer_reason,
4471 fix: None,
4472 blocking,
4473 answered,
4474 expected,
4475 clean,
4476 progressed: false,
4477 vote_split,
4478 reconsideration,
4479 verdict: round_verdict,
4480 };
4481 if !matches!(
4492 round_record.e2e_status(),
4493 E2eStatus::Deferred | E2eStatus::NotConfigured
4494 ) {
4495 round_record.verified_head = Some(head.clone());
4496 round_record.verified_at = Some(Timestamp::now());
4497 }
4498 let this_round_verification = round_record.verification_summary(&head);
4499
4500 if incomplete {
4501 let missing: Vec<String> = round_record
4502 .reviews
4503 .iter()
4504 .filter(|r| r.failed.is_some())
4505 .map(|r| format!("review-{}", r.reviewer))
4506 .collect();
4507 self.state.event(
4508 "review",
4509 format!(
4510 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4511 missing.join(", ")
4512 ),
4513 );
4514 }
4515
4516 if clean {
4517 self.state.event(
4518 "review",
4519 if incomplete && policy == IncompleteReviewPolicy::Warn {
4520 format!(
4521 "round {round}: clean (warn policy, incomplete panel) — no \
4522 blocking findings from the seats that answered, verification green"
4523 )
4524 } else if incomplete {
4525 format!(
4526 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4527 quorum) — no blocking findings from the seats that answered, \
4528 verification green",
4529 expected - answered
4530 )
4531 } else {
4532 format!("round {round}: clean — no blocking findings, verification green")
4533 },
4534 );
4535 self.state.reviews.push(round_record);
4536 self.state.status = RunStatus::Gating;
4537 self.state.save()?;
4538 return Ok(());
4539 }
4540
4541 if incomplete && blocking == 0 && e2e_ok {
4549 self.state.reviews.push(round_record);
4550 self.state.save()?;
4551 if round == max_rounds {
4552 self.state.status = RunStatus::Blocked;
4553 self.state.event(
4554 "review",
4555 format!(
4556 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4557 refusing to call it clean",
4558 expected - answered
4559 ),
4560 );
4561 return Ok(());
4562 }
4563 continue;
4564 }
4565
4566 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4578 self.state.reviews.push(round_record);
4579 return self
4580 .stop_reviewing(
4581 "the round's own verification could not run",
4582 &shell,
4583 &winner.worktree,
4584 )
4585 .await;
4586 }
4587
4588 if round == max_rounds {
4589 self.state.reviews.push(round_record);
4590 return self
4591 .stop_reviewing(
4592 &format!(
4593 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4594 ),
4595 &shell,
4596 &winner.worktree,
4597 )
4598 .await;
4599 }
4600
4601 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4604 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4605 let blocking_findings: Vec<_> = all_findings
4606 .iter()
4607 .filter(|f| f.severity.blocks())
4608 .cloned()
4609 .collect();
4610 let job = SeatJob {
4611 prompt: prompt::fix(
4612 &self.state.instruction,
4613 &blocking_findings,
4614 this_round_verification.as_ref(),
4615 round,
4616 max_rounds,
4617 &language,
4618 ),
4619 spec: fix_spec.clone(),
4620 seat,
4621 cwd: winner.worktree.clone(),
4622 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4623 allow_write: true,
4624 sessions,
4625 artifacts: artifacts.clone(),
4626 stem: format!("fix-{round}"),
4627 };
4628 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4629 let cache = self.state.config.cache_dir();
4630 let ctx = WaveCtx {
4631 run: &run_id,
4632 node: "fix",
4633 prompts: &prompts,
4634 cache: cache.as_deref(),
4635 round: Some(round),
4636 };
4637 let (seat, out) =
4638 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4639 let agent_id = seat.agent.clone();
4640
4641 let mut fix = FixRecord {
4642 agent: agent_id,
4643 addressed: Vec::new(),
4644 rejected: Vec::new(),
4645 notes: String::new(),
4646 committed: false,
4647 failed: None,
4648 duration_ms: 0,
4649 continuation: None,
4650 };
4651 let mut continuation = ContinuationRecord::not_needed();
4652 let mut final_seat = seat.clone();
4653 match out {
4654 AgentOutcome::Ok(o) => {
4655 fix.duration_ms = o.duration_ms;
4656 let parsed = verdict::extract_json::<FixReport>(&o.text);
4657 let incomplete_reason = match &parsed {
4664 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4665 "the reply parsed, but it reported a command whose own CLI never \
4666 confirmed an exit status"
4667 .to_owned(),
4668 ),
4669 Ok(_) => None,
4670 Err(e) => Some(e.to_string()),
4671 };
4672 match incomplete_reason {
4673 None => {
4674 let report = parsed.expect("checked Ok above");
4675 fix.addressed = report.addressed;
4676 fix.rejected = report.rejected;
4677 fix.notes =
4678 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4679 }
4680 Some(reason) => {
4681 let (resumed_seat, resolved, failure, cont) = self
4682 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4683 .await;
4684 fix.duration_ms += cont.cumulative_wait_ms;
4685 continuation = cont;
4686 final_seat = resumed_seat;
4687 match resolved {
4688 Some(report) => {
4689 fix.addressed = report.addressed;
4690 fix.rejected = report.rejected;
4691 fix.notes = blind::sanitize_prose(
4692 &report.notes,
4693 &self.state.config.blind,
4694 );
4695 }
4696 None => fix.failed = failure,
4697 }
4698 }
4699 }
4700 }
4701 AgentOutcome::Dropped(o) => {
4703 fix.duration_ms = o.duration_ms;
4704 let why = o
4705 .dropped
4706 .as_ref()
4707 .map(|d| d.why.as_str())
4708 .unwrap_or("the CLI ended the stream without delivering its answer");
4709 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4710 }
4711 AgentOutcome::Quota(o) => {
4712 self.state.quota.push(QuotaLoss {
4713 seat: final_seat.key.clone(),
4714 node: "fix".to_owned(),
4715 at: Timestamp::now(),
4716 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4717 });
4718 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4719 }
4720 AgentOutcome::Failed(e) => fix.failed = Some(e),
4721 }
4722 fix.continuation = Some(continuation);
4723 self.state.seats.insert(final_seat.key.clone(), final_seat);
4724 if let Ok(r) = git::rescue_commit(
4725 &winner.worktree,
4726 &format!("magi: review round {round} fixes (uncommitted work)"),
4727 )
4728 .await
4729 {
4730 self.state.note_withheld("fix", &r.withheld);
4731 }
4732 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4733 fix.committed = after != before;
4734 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4742 let progressed = diff_after != patch;
4743 let commit_note = if fix.committed {
4744 "committed"
4745 } else {
4746 "NO new commit"
4747 };
4748 let tree_note = if progressed {
4749 "changed vs base"
4750 } else {
4751 "unchanged vs base"
4752 };
4753 self.state.event(
4754 "fix",
4755 match &fix.failed {
4756 Some(reason) => {
4762 format!(
4763 "round {round}: fixer's adoption report was lost ({reason}); \
4764 {commit_note}, tree {tree_note}"
4765 )
4766 }
4767 None => format!(
4768 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4769 {tree_note}{}",
4770 fix.addressed.len(),
4771 fix.rejected.len(),
4772 if continuation.outcome == ContinuationOutcome::Resumed {
4773 format!(
4774 " (adoption report recovered after {} continuation(s))",
4775 continuation.attempts
4776 )
4777 } else {
4778 String::new()
4779 },
4780 ),
4781 },
4782 );
4783 round_record.fix = Some(fix);
4784 round_record.progressed = progressed;
4785 self.state.reviews.push(round_record);
4786 self.state.save()?;
4787
4788 if matches!(
4801 continuation.outcome,
4802 ContinuationOutcome::Exhausted
4803 | ContinuationOutcome::QuotaLost
4804 | ContinuationOutcome::NoSession
4805 ) {
4806 return self
4807 .stop_reviewing(
4808 "the fixer's adoption report never came back, even after resuming its \
4809 own seat; refusing to start another round against the same worktree \
4810 while that is unresolved",
4811 &shell,
4812 &winner.worktree,
4813 )
4814 .await;
4815 }
4816
4817 let streak = self
4818 .state
4819 .reviews
4820 .iter()
4821 .rev()
4822 .take_while(|r| !r.progressed)
4823 .count();
4824 if streak >= STAGNANT_LIMIT {
4825 return self
4826 .stop_reviewing(
4827 &format!(
4828 "the tree has not moved against base for {streak} round(s) in a row"
4829 ),
4830 &shell,
4831 &winner.worktree,
4832 )
4833 .await;
4834 }
4835 }
4836 Ok(())
4837 }
4838
4839 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4869 let round_idx = self.state.reviews.len() - 1;
4870 let needs_catchup_run = matches!(
4878 self.state.reviews[round_idx].e2e_status(),
4879 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4880 );
4881 if needs_catchup_run {
4882 let round = self.state.reviews[round_idx].round;
4883 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4884 let commands = self.state.config.verify.e2e.clone();
4885 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4886 let cache_dir = self.state.config.cache_dir();
4887 let context = format!(
4888 "round {round}: verification unresolved, catching up before the final decision"
4889 );
4890 let (outcomes, verify_retried) = with_cache_lease(
4891 &mut self.state,
4892 cache_dir.as_deref(),
4893 "e2e",
4894 "e2e",
4895 worktree,
4896 &attempted_head,
4897 timeout,
4898 &context,
4899 |state, budget| {
4900 let shell = shell.to_vec();
4901 let commands = commands.clone();
4902 let context = context.clone();
4903 async move {
4904 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4905 .await
4906 }
4907 },
4908 )
4909 .await;
4910 let last = &mut self.state.reviews[round_idx];
4911 last.e2e = outcomes;
4912 last.verify_retried = verify_retried;
4913 last.verified_head = Some(attempted_head);
4920 last.verified_at = Some(Timestamp::now());
4921 if verify_inconclusive(&last.e2e) {
4922 self.state.save()?;
4929 return Ok(());
4930 }
4931 last.e2e_deferred = false;
4932 }
4933 let last = &self.state.reviews[round_idx];
4934 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4935
4936 match last.e2e_status() {
4937 E2eStatus::Failed => {
4938 let red: Vec<String> = last
4939 .e2e
4940 .iter()
4941 .filter(|o| !o.ok())
4942 .map(|o| {
4943 format!(
4944 "`{}` -> {:?}\n{}",
4945 o.command,
4946 o.code,
4947 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4948 )
4949 })
4950 .collect();
4951 self.state
4952 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4953 self.state.status = RunStatus::Blocked;
4954 }
4955 E2eStatus::ResourceBlocked => {
4960 self.state.event(
4961 "review",
4962 format!(
4963 "{why}; e2e could not run (shared build cache unavailable); not \
4964 deciding yet"
4965 ),
4966 );
4967 }
4968 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4969 self.state.event(
4970 "review",
4971 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4972 );
4973 self.state.status = RunStatus::Gating;
4974 }
4975 }
4976 self.state.save()?;
4977 Ok(())
4978 }
4979
4980 async fn gate(&mut self) -> Result<()> {
4983 if self.state.status == RunStatus::Failed
4995 || self
4996 .state
4997 .base_sync
4998 .as_ref()
4999 .is_some_and(|s| s.conflict.is_some())
5000 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5001 != Some(RunStatus::Gating)
5002 {
5003 return Ok(());
5004 }
5005 if self.state.gate_ran {
5006 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
5017 self.state.status = RunStatus::Blocked;
5018 self.state.save()?;
5019 }
5020 return Ok(());
5021 }
5022 let Some(winner) = self.state.winner().cloned() else {
5023 return Ok(());
5024 };
5025 self.state.status = RunStatus::Gating;
5026 let mut outcomes = self.run_gate(&winner).await?;
5027 loop {
5028 if verify_inconclusive(&outcomes) {
5039 self.state.save()?;
5040 return Ok(());
5041 }
5042 if outcomes.iter().all(CommandOutcome::ok) {
5043 break;
5044 }
5045 match self.gate_fix_round(&winner, &outcomes).await? {
5046 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
5047 GateFix::Stop => break,
5048 GateFix::Defer => {
5049 self.state.save()?;
5050 return Ok(());
5051 }
5052 }
5053 }
5054 let passed = outcomes.iter().all(CommandOutcome::ok);
5055 self.state.gate = outcomes;
5056 self.state.gate_ran = true;
5057 if !passed {
5058 self.state.status = RunStatus::Blocked;
5059 let spent = self.state.gate_fixes.len();
5060 self.state.event(
5061 "gate",
5062 if spent == 0 {
5063 "gate failed; not merging".to_owned()
5064 } else {
5065 format!("gate failed after {spent} gate-fix round(s); not merging")
5066 },
5067 );
5068 }
5069 self.state.save()?;
5070 Ok(())
5071 }
5072
5073 async fn run_pre_gate(&mut self, winner: &Candidate) {
5083 let commands = self.state.config.verify.pre_gate.clone();
5084 if commands.is_empty() {
5085 return;
5086 }
5087 let shell = self.state.config.shell();
5088 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5089 let (outcomes, _) = run_commands(
5090 &mut self.state,
5091 "pre_gate",
5092 "pre_gate",
5093 0,
5094 &shell,
5095 &commands,
5096 &winner.worktree,
5097 timeout,
5098 )
5099 .await;
5100 for o in &outcomes {
5101 if !o.ok() {
5102 tracing::warn!(
5103 "pre_gate `{}` failed ({:?}); the gate decides",
5104 o.command,
5105 o.code
5106 );
5107 }
5108 self.state.event(
5109 "pre_gate",
5110 format!(
5111 "`{}` -> {}",
5112 o.command,
5113 if o.ok() {
5114 "pass".to_owned()
5115 } else {
5116 format!(
5117 "FAIL ({:?})\n{}",
5118 o.code,
5119 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5120 )
5121 }
5122 ),
5123 );
5124 }
5125 self.state.pre_gate = outcomes;
5126 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
5127 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
5128 Ok(head) => {
5129 self.state
5130 .event("pre_gate", format!("committed mechanical fixes ({head})"));
5131 self.state.pre_gate_commit = Some(head);
5132 }
5133 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
5134 },
5135 Ok(false) => {}
5136 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
5137 }
5138 if let Err(e) = self.state.save() {
5139 tracing::warn!("could not persist the pre_gate record: {e:#}");
5140 }
5141 }
5142
5143 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
5146 self.run_pre_gate(winner).await;
5147 let shell = self.state.config.shell();
5148 let gate_commands = self.state.config.verify.gate.clone();
5149 let outcomes = if gate_commands.is_empty() {
5158 Vec::new()
5159 } else {
5160 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5161 let cache_dir = self.state.config.cache_dir();
5162 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
5163 let (outcomes, _) = with_cache_lease(
5164 &mut self.state,
5165 cache_dir.as_deref(),
5166 "gate",
5167 "gate",
5168 &winner.worktree,
5169 &head,
5170 timeout,
5171 "final gate",
5172 |state, budget| {
5173 let shell = shell.clone();
5174 let gate_commands = gate_commands.clone();
5175 let worktree = winner.worktree.clone();
5176 async move {
5177 let (outcomes, timed_out_pids) = run_commands(
5178 state,
5179 "gate",
5180 "gate",
5181 0,
5182 &shell,
5183 &gate_commands,
5184 &worktree,
5185 budget,
5186 )
5187 .await;
5188 (outcomes, false, timed_out_pids)
5189 }
5190 },
5191 )
5192 .await;
5193 outcomes
5194 };
5195 if outcomes.is_empty() {
5196 self.state.event(
5201 "gate",
5202 "no gate commands configured; nothing to check, passing",
5203 );
5204 }
5205 for o in &outcomes {
5206 self.state.event(
5207 "gate",
5208 format!(
5209 "`{}` -> {}",
5210 o.command,
5211 if o.ok() {
5212 "pass".to_owned()
5213 } else {
5214 format!(
5215 "FAIL ({:?})\n{}",
5216 o.code,
5217 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5218 )
5219 }
5220 ),
5221 );
5222 }
5223 Ok(outcomes)
5224 }
5225
5226 async fn gate_fix_round(
5239 &mut self,
5240 winner: &Candidate,
5241 outcomes: &[CommandOutcome],
5242 ) -> Result<GateFix> {
5243 let cap = self.state.config.graph.gate_fix_rounds;
5244 let spent = self.state.gate_fixes.len();
5245 if spent >= cap {
5246 if cap > 0 {
5247 self.state.event(
5248 "gate",
5249 format!("{spent} gate-fix round(s) spent and the gate still fails"),
5250 );
5251 }
5252 return Ok(GateFix::Stop);
5253 }
5254 if !gate_fixable(outcomes) {
5255 self.state.event(
5256 "gate",
5257 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
5258 command or similar); not spending a fix round on it",
5259 );
5260 return Ok(GateFix::Stop);
5261 }
5262 let min_free = self.state.config.disk.min_free_bytes;
5263 if min_free > 0 {
5264 match crate::disk::free_bytes(&winner.worktree) {
5265 Ok(free) if crate::disk::enough_space(free, min_free) => {}
5266 Ok(free) => {
5267 self.state.event(
5268 "gate",
5269 format!(
5270 "only {free} bytes free ({min_free} required by `[disk] \
5271 min_free_bytes`); not spending a fix round on a failure the disk \
5272 may explain"
5273 ),
5274 );
5275 return Ok(GateFix::Stop);
5276 }
5277 Err(e) => {
5278 self.state.event(
5279 "gate",
5280 format!("free disk space could not be measured ({e:#}); no fix round"),
5281 );
5282 return Ok(GateFix::Stop);
5283 }
5284 }
5285 }
5286
5287 let attempt = spent + 1;
5288 let run_id = self.state.id.clone();
5289 let prompts = self.state.config.prompts.clone();
5290 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
5291 let base = self.landing_base();
5292 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
5293 let seat = self.seat(&fix_seat_key, &fix_spec.id);
5294 let job = SeatJob {
5295 prompt: prompt::gate_fix(
5296 &self.state.instruction,
5297 &failed,
5298 attempt,
5299 cap,
5300 &self.state.config.graph.language,
5301 ),
5302 spec: fix_spec,
5303 seat,
5304 cwd: winner.worktree.clone(),
5305 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
5306 allow_write: true,
5307 sessions: self.state.config.graph.sessions,
5308 artifacts: agent::artifacts_dir(&self.state.dir()),
5309 stem: format!("gate-fix-{attempt}"),
5310 };
5311 self.state.event(
5312 "gate",
5313 format!("gate failed; gate-fix round {attempt} of {cap}"),
5314 );
5315 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
5316 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
5317 let cache = self.state.config.cache_dir();
5318 let ctx = WaveCtx {
5319 run: &run_id,
5320 node: "gate-fix",
5321 prompts: &prompts,
5322 cache: cache.as_deref(),
5323 round: None,
5324 };
5325 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
5326 let mut record = GateFixRecord {
5327 agent: seat.agent.clone(),
5328 failed,
5329 notes: String::new(),
5330 committed: false,
5331 error: None,
5332 };
5333 match out {
5334 AgentOutcome::Ok(o) => {
5335 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
5338 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
5339 }
5340 }
5341 AgentOutcome::Dropped(_) => {
5342 record.error = Some("the CLI dropped the stream".to_owned());
5343 }
5344 AgentOutcome::Quota(o) => {
5345 self.state.quota.push(QuotaLoss {
5346 seat: seat.key.clone(),
5347 node: "gate-fix".to_owned(),
5348 at: Timestamp::now(),
5349 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5350 });
5351 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5352 }
5353 AgentOutcome::Failed(e) => record.error = Some(e),
5354 }
5355 self.state.seats.insert(seat.key.clone(), seat);
5356 if let Ok(r) = git::rescue_commit(
5357 &winner.worktree,
5358 &format!("magi: gate fix {attempt} (uncommitted work)"),
5359 )
5360 .await
5361 {
5362 self.state.note_withheld("gate-fix", &r.withheld);
5363 }
5364 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5365 record.committed = after != before;
5366 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5367 let note = record.error.clone();
5368 self.state.gate_fixes.push(record);
5369 self.state.save()?;
5370 if !changed {
5371 self.state.event(
5372 "gate",
5373 match note {
5374 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5375 None => format!("gate-fix round {attempt}: the tree did not change"),
5376 },
5377 );
5378 return Ok(GateFix::Stop);
5379 }
5380 self.state.event(
5381 "gate",
5382 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5383 );
5384
5385 let commands = self.state.config.verify.e2e.clone();
5386 if !commands.is_empty() {
5387 let shell = self.state.config.shell();
5388 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5389 let cache_dir = self.state.config.cache_dir();
5390 let context = format!("gate-fix round {attempt}");
5391 let (e2e, _) = with_cache_lease(
5392 &mut self.state,
5393 cache_dir.as_deref(),
5394 "e2e",
5395 "e2e",
5396 &winner.worktree,
5397 &after,
5398 timeout,
5399 &context,
5400 |state, budget| {
5401 let shell = shell.clone();
5402 let commands = commands.clone();
5403 let context = context.clone();
5404 let worktree = winner.worktree.clone();
5405 async move {
5406 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5407 .await
5408 }
5409 },
5410 )
5411 .await;
5412 if verify_inconclusive(&e2e) {
5413 return Ok(GateFix::Defer);
5414 }
5415 if e2e.iter().any(|o| !o.ok()) {
5416 self.state.event(
5417 "gate",
5418 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5419 );
5420 return Ok(GateFix::Stop);
5421 }
5422 }
5423 Ok(GateFix::Retry)
5424 }
5425
5426 async fn merge(&mut self) -> Result<()> {
5429 if self
5444 .state
5445 .base_sync
5446 .as_ref()
5447 .is_some_and(|s| s.conflict.is_some())
5448 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5449 != Some(RunStatus::Gating)
5450 || !self.state.gate_status().ok()
5459 {
5460 return Ok(());
5461 }
5462 if self.state.merge.is_some() {
5471 return Ok(());
5472 }
5473 let Some(winner) = self.state.winner().cloned() else {
5474 return Ok(());
5475 };
5476 let repo = self.state.repo.clone();
5477 let base = self.state.base_branch.clone();
5478 let mode = self.state.config.merge.mode;
5479 let style = self.state.config.merge.style;
5480 let pr = pr_message(&self.state, winner.label);
5481 let message = pr.commit_message();
5482
5483 let outcome = match mode {
5484 MergeMode::None => MergeOutcome {
5485 mode,
5486 ok: true,
5487 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5488 empty: false,
5489 },
5490 MergeMode::Pr | MergeMode::Local
5491 if merge_is_empty(&repo, &self.state, &winner.branch, mode).await =>
5492 {
5493 MergeOutcome {
5494 mode,
5495 ok: false,
5496 detail: empty_candidate_detail(&self.state, &base),
5497 empty: true,
5498 }
5499 }
5500 MergeMode::Local => {
5501 let on = git::current_branch(&repo).await?;
5502 if on.as_deref() != Some(base.as_str()) {
5503 MergeOutcome {
5504 mode,
5505 ok: false,
5506 detail: format!(
5507 "{} has {} checked out, not the base branch {base}",
5508 repo.display(),
5509 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5510 ),
5511 empty: false,
5512 }
5513 } else if !git::is_clean(&repo).await? {
5514 MergeOutcome {
5515 mode,
5516 ok: false,
5517 detail: format!("{} is dirty; refusing to merge", repo.display()),
5518 empty: false,
5519 }
5520 } else {
5521 let out = match style {
5522 MergeStyle::Merge => {
5523 git::merge_no_ff(&repo, &winner.branch, &message).await?
5524 }
5525 MergeStyle::Squash => {
5526 git::merge_squash(&repo, &winner.branch, &message).await?
5527 }
5528 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5529 };
5530 MergeOutcome {
5531 mode,
5532 ok: out.ok(),
5533 detail: if out.ok() { out.stdout } else { out.stderr },
5534 empty: false,
5535 }
5536 }
5537 }
5538 MergeMode::Pr => {
5539 let remote = self.state.config.merge.remote.clone();
5540 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5541 if !pushed.ok() {
5542 MergeOutcome {
5543 mode,
5544 ok: false,
5545 detail: pushed.stderr,
5546 empty: false,
5547 }
5548 } else {
5549 let found = land::find_open_pr(&winner.worktree, &winner.branch, &base).await;
5555 let out = match pr_merge_plan(found) {
5556 PrPlan::Create => {
5557 gh_pr_create(
5558 &winner.worktree,
5559 &base,
5560 &winner.branch,
5561 &pr.title,
5562 &pr.body,
5563 )
5564 .await
5565 }
5566 PrPlan::Adopt { url, title } => {
5567 self.state
5568 .event("merge", format!("Pr: adopted open pull request {url}"));
5569 if title != pr.title
5570 && let Err(e) =
5571 land::set_pr_title(&winner.worktree, &url, &pr.title).await
5572 {
5573 tracing::warn!("could not refresh title of {url}: {e:#}");
5574 self.state
5575 .event("merge", format!("Pr: title refresh failed: {e:#}"));
5576 }
5577 Ok(url)
5578 }
5579 PrPlan::Stop(why) => Err(anyhow::anyhow!(why)),
5580 };
5581 match out {
5582 Ok(url) => MergeOutcome {
5583 mode,
5584 ok: true,
5585 detail: url,
5586 empty: false,
5587 },
5588 Err(e) => MergeOutcome {
5589 mode,
5590 ok: false,
5591 detail: e.to_string(),
5592 empty: false,
5593 },
5594 }
5595 }
5596 }
5597 };
5598
5599 self.state.status = match (mode, outcome.ok) {
5600 (MergeMode::None, _) => RunStatus::Ready,
5601 (_, true) => RunStatus::Merged,
5602 (_, false) => RunStatus::Blocked,
5603 };
5604 self.state.event(
5605 "merge",
5606 format!(
5607 "{:?}: {}",
5608 mode,
5609 outcome.detail.lines().next().unwrap_or("")
5610 ),
5611 );
5612 self.state.merge = Some(outcome);
5613 self.state.save()?;
5614
5615 if self.state.config.graph.land
5621 && mode == MergeMode::Pr
5622 && self.state.status == RunStatus::Merged
5623 {
5624 self.run_land().await?;
5625 }
5626 self.settle_questions();
5631 Ok(())
5632 }
5633
5634 async fn run_land(&mut self) -> Result<()> {
5645 let url = self
5646 .state
5647 .merge
5648 .as_ref()
5649 .map(|m| m.detail.clone())
5650 .unwrap_or_default();
5651 let url = url.lines().next().unwrap_or("").trim().to_owned();
5652 if !url.starts_with("http") {
5653 return Ok(());
5654 }
5655 match land::land(&mut self.state, &url).await {
5658 Ok(pr) if self.state.parked => {
5659 let _ = pr;
5663 }
5664 Ok(pr) => {
5665 self.state.status = match pr.state {
5666 land::PrLifecycle::Merged => RunStatus::Merged,
5667 _ => RunStatus::Blocked,
5668 };
5669 if bump::should_release_bump(self.state.status)
5676 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5677 {
5678 self.state
5684 .event("bump", format!("release bump skipped: {e:#}"));
5685 }
5686 self.state.save()?;
5687 }
5688 Err(e) => {
5689 self.state.status = RunStatus::Blocked;
5690 self.state.event("land", format!("gave up: {e}"));
5691 self.state.save()?;
5692 }
5693 }
5694 Ok(())
5695 }
5696
5697 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5701 if let Some(existing) = self.state.seats.get(key)
5702 && existing.agent == agent
5703 {
5704 return existing.clone();
5705 }
5706 let fresh = SeatState::new(key, agent, self.state.seed);
5707 self.state.seats.insert(key.to_owned(), fresh.clone());
5708 fresh
5709 }
5710
5711 fn view(&self, c: &Candidate) -> CandidateView {
5713 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5714 .unwrap_or_default();
5715 let (patch, _) = blind::sanitize_patch(
5716 &format!("candidate {} patch", c.label),
5717 &raw,
5718 &self.state.config.blind,
5719 );
5720 CandidateView {
5721 label: c.label,
5722 branch: c.branch.clone(),
5723 summary: c.summary.clone(),
5724 stat: c.stat.clone(),
5725 patch,
5726 }
5727 }
5728
5729 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5731 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5732 prompt::judge(
5733 "(see above)",
5734 &views,
5735 self.roles.judges.len(),
5736 base_short,
5737 "en",
5738 )
5739 }
5740
5741 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5748 let mut turns = Vec::new();
5749 for j in &self.state.judgements {
5750 if j.ranking.is_empty() {
5751 continue;
5752 }
5753 let reasons = j
5754 .reasons
5755 .iter()
5756 .map(|(k, v)| format!("- {k}: {v}"))
5757 .collect::<Vec<_>>()
5758 .join("\n");
5759 turns.push(Turn {
5760 who: format!("Judge {} (opening ranking)", j.judge),
5761 is_self: j.judge == self_idx + 1,
5762 body: format!(
5763 "Ranked {}{}{reasons}",
5764 j.ranking.iter().collect::<String>(),
5765 if reasons.is_empty() {
5766 ""
5767 } else {
5768 ", because:\n"
5769 }
5770 ),
5771 });
5772 }
5773 for t in self
5774 .state
5775 .deliberation
5776 .iter()
5777 .flat_map(|r| r.turns.iter())
5778 .chain(current)
5779 {
5780 turns.push(Turn {
5781 who: format!("Judge {}", t.judge),
5782 is_self: t.judge == self_idx + 1,
5783 body: t.body.clone(),
5784 });
5785 }
5786 turns
5787 }
5788}
5789
5790fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5792 agent::has_session(spec.kind, seat, sessions)
5793}
5794
5795fn next_untried_implementer<'a>(
5816 roster: &'a [AgentSpec],
5817 start: usize,
5818 tried: &BTreeSet<String>,
5819) -> Option<&'a AgentSpec> {
5820 roster
5821 .get(start + 1..)?
5822 .iter()
5823 .find(|s| !tried.contains(&s.id))
5824}
5825
5826fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5840 commands.iter().any(|c| c.exit_code.is_none())
5841}
5842
5843fn verified_noop_claim(
5856 usable: bool,
5857 commands: &[agent::CommandEvidence],
5858 text: &str,
5859) -> Option<String> {
5860 (usable && !has_unconfirmed_command(commands))
5861 .then(|| verdict::verified_noop(text))
5862 .flatten()
5863}
5864
5865fn short(commit: &str) -> String {
5866 commit.chars().take(7).collect()
5867}
5868
5869fn make_executable(path: &Path) -> Result<()> {
5870 #[cfg(unix)]
5871 {
5872 use std::os::unix::fs::PermissionsExt as _;
5873 let mut perms = std::fs::metadata(path)?.permissions();
5874 perms.set_mode(0o755);
5875 std::fs::set_permissions(path, perms)?;
5876 }
5877 #[cfg(not(unix))]
5878 {
5879 let _ = path;
5880 }
5881 Ok(())
5882}
5883
5884struct WaveCtx<'a> {
5891 run: &'a str,
5894 node: &'a str,
5896 prompts: &'a Prompts,
5897 cache: Option<&'a Path>,
5899 round: Option<usize>,
5902}
5903
5904async fn run_one(
5906 job: SeatJob,
5907 sem: Arc<Semaphore>,
5908 ctx: &WaveCtx<'_>,
5909 state: &mut RunState,
5910 attempt: usize,
5911) -> (SeatState, AgentOutcome) {
5912 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5913 .await
5914 .pop()
5915 .expect("one job in, one result out");
5916 (seat, out)
5917}
5918
5919async fn wave(
5925 jobs: Vec<SeatJob>,
5926 sem: Arc<Semaphore>,
5927 ctx: &WaveCtx<'_>,
5928 state: &mut RunState,
5929 attempt: usize,
5930) -> Vec<(usize, SeatState, AgentOutcome)> {
5931 let WaveCtx {
5932 run,
5933 node,
5934 prompts,
5935 cache,
5936 round,
5937 } = *ctx;
5938 for job in &jobs {
5939 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5940 }
5941 if let Err(e) = state.save() {
5942 tracing::warn!("could not persist in-progress seats: {e:#}");
5947 }
5948 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5964 let wait_started = Instant::now();
5965 let cache_guard = if let Some(cache_dir) = cache {
5966 if jobs_had_a_writer {
5967 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5968 let budget = jobs
5969 .iter()
5970 .map(|j| j.timeout)
5971 .max()
5972 .unwrap_or(Duration::from_secs(60));
5973 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5974 .await
5975 .ok()
5976 } else {
5977 None
5978 }
5979 } else {
5980 None
5981 };
5982 let waited_for_lease = wait_started.elapsed();
5989 let mut set = tokio::task::JoinSet::new();
5990 let overlay = prompts.overlay(node);
5991 for (i, mut job) in jobs.into_iter().enumerate() {
5992 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5993 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
5994 if cache.is_some() {
5995 job.prompt.push('\n');
5996 job.prompt
5997 .push_str(&prompt::build_cache_note(node, job.allow_write));
5998 }
5999 let sem = Arc::clone(&sem);
6000 let run = run.to_owned();
6001 let node = node.to_owned();
6002 let attachments = if node == "implement" {
6005 state.attachments.clone()
6006 } else {
6007 Vec::new()
6008 };
6009 let cache = cache
6020 .filter(|_| job.allow_write && cache_guard.is_some())
6021 .map(Path::to_path_buf);
6022 set.spawn(async move {
6023 let _permit = sem.acquire().await;
6024 let mut seat = job.seat;
6025 let out = agent::invoke(
6026 &job.spec,
6027 &mut seat,
6028 &Invocation {
6029 cwd: &job.cwd,
6030 prompt: &job.prompt,
6031 timeout: job.timeout,
6032 allow_write: job.allow_write,
6033 sessions: job.sessions,
6034 artifacts: &job.artifacts,
6035 stem: &job.stem,
6036 run: &run,
6037 node: &node,
6038 cache_dir: cache.as_deref(),
6039 attachments: &attachments,
6040 },
6041 )
6042 .await;
6043 let out = match out {
6044 Ok(o) if o.usable() => AgentOutcome::Ok(o),
6045 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
6046 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
6054 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
6055 Ok(o) => AgentOutcome::Failed(format!(
6056 "exited with {:?} and no usable output",
6057 o.exit_code
6058 )),
6059 Err(e) => AgentOutcome::Failed(e.to_string()),
6060 };
6061 (i, seat, out)
6062 });
6063 }
6064 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
6065 while let Some(joined) = set.join_next().await {
6066 let (i, seat, out) = match joined {
6067 Ok(v) => v,
6068 Err(e) => {
6072 tracing::error!("agent task panicked: {e}");
6073 continue;
6074 }
6075 };
6076 state.seat_finished(&seat.key);
6077 record_jobs(state, node, round, &seat.key, &out);
6078 if let Err(e) = state.save() {
6079 tracing::warn!("could not persist a seat's completion: {e:#}");
6080 }
6081 if collected.len() <= i {
6082 collected.resize_with(i + 1, || None);
6083 }
6084 collected[i] = Some((i, seat, out));
6085 }
6086 if state
6092 .active
6093 .values()
6094 .any(|a| a.node == node && a.attempt == attempt)
6095 {
6096 state
6097 .active
6098 .retain(|_, a| !(a.node == node && a.attempt == attempt));
6099 if let Err(e) = state.save() {
6100 tracing::warn!("could not persist the end of a wave: {e:#}");
6101 }
6102 }
6103 if let Some(cache_dir) = cache
6110 && jobs_had_a_writer
6111 {
6112 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
6113 }
6114 if let Some(guard) = cache_guard {
6115 guard.release();
6116 }
6117 collected.into_iter().flatten().collect()
6118}
6119
6120fn record_jobs(
6131 state: &mut RunState,
6132 node: &str,
6133 round: Option<usize>,
6134 seat: &str,
6135 out: &AgentOutcome,
6136) {
6137 let commands: &[agent::CommandEvidence] = match out {
6138 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
6139 AgentOutcome::Failed(_) => &[],
6140 };
6141 let checked_at = Timestamp::now();
6142 for c in commands {
6143 state.jobs.push(JobRecord {
6144 node: node.to_owned(),
6145 round,
6146 seat: seat.to_owned(),
6147 id: c.id.clone(),
6148 description: c.description.clone(),
6149 checked_at,
6150 status: match c.exit_code {
6151 Some(0) => JobStatus::Completed,
6152 Some(_) => JobStatus::Failed,
6153 None => JobStatus::Unknown,
6154 },
6155 exit_code: c.exit_code,
6156 result_summary: c.result_summary.clone(),
6157 source: c.source.clone(),
6158 });
6159 }
6160}
6161
6162fn round_is_clean(
6183 blocking: usize,
6184 e2e_ok: bool,
6185 answered: usize,
6186 expected: usize,
6187 quota_missing: usize,
6188 policy: IncompleteReviewPolicy,
6189) -> bool {
6190 if blocking != 0 || !e2e_ok {
6191 return false;
6192 }
6193 if answered == expected || policy == IncompleteReviewPolicy::Warn {
6194 return true;
6195 }
6196 answered > 0 && expected - answered <= quota_missing
6197}
6198
6199fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
6223 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
6224 return Some(RunStatus::Gating);
6225 }
6226 let last = reviews.last()?;
6227 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
6228 if reviews.len() < max_rounds && !stagnant {
6229 return None;
6230 }
6231 if last.incomplete() && last.blocking == 0 {
6232 return Some(RunStatus::Blocked);
6233 }
6234 if last.e2e_status() == E2eStatus::ResourceBlocked {
6235 return None;
6236 }
6237 Some(if last.e2e.iter().all(CommandOutcome::ok) {
6238 RunStatus::Gating
6239 } else {
6240 RunStatus::Blocked
6241 })
6242}
6243
6244fn retry_budget(full: Duration, nudged: bool) -> Duration {
6259 if nudged {
6260 (full / 4).max(Duration::from_secs(120)).min(full)
6261 } else {
6262 full
6263 }
6264}
6265
6266#[allow(clippy::too_many_arguments)]
6279async fn ask_json_wave<T>(
6280 jobs: Vec<SeatJob>,
6281 sem: Arc<Semaphore>,
6282 retries: usize,
6283 ctx: &WaveCtx<'_>,
6284 losses: &mut Vec<QuotaLoss>,
6285 state: &mut RunState,
6286 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
6287) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
6288where
6289 T: serde::de::DeserializeOwned + Send + 'static,
6290{
6291 let n = jobs.len();
6292 let originals: Vec<SeatJob> = jobs;
6293 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
6294 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
6295 let mut attempts_used: Vec<usize> = vec![0; n];
6302 let mut pending: Vec<usize> = (0..n).collect();
6303
6304 for attempt in 0..=retries {
6305 if pending.is_empty() {
6306 break;
6307 }
6308 let mut batch = Vec::with_capacity(pending.len());
6309 for &i in &pending {
6310 let src = &originals[i];
6311 let (prompt, timeout) = if attempt == 0 {
6314 (src.prompt.clone(), src.timeout)
6315 } else {
6316 let why = done[i]
6317 .as_ref()
6318 .and_then(|r| r.as_ref().err().map(ToString::to_string))
6319 .unwrap_or_else(|| "no parsable answer".to_owned());
6320 let nudge = prompt::nudge(&why);
6321 let nudged = has_context(&src.spec, &seats[i], src.sessions);
6322 let prompt = if nudged {
6323 nudge
6324 } else {
6325 format!("{}\n\n---\n\n{}", src.prompt, nudge)
6326 };
6327 (prompt, retry_budget(src.timeout, nudged))
6328 };
6329 batch.push(SeatJob {
6330 spec: src.spec.clone(),
6331 seat: seats[i].clone(),
6332 cwd: src.cwd.clone(),
6333 prompt,
6334 timeout,
6335 allow_write: src.allow_write,
6336 sessions: src.sessions,
6337 artifacts: src.artifacts.clone(),
6338 stem: if attempt == 0 {
6339 src.stem.clone()
6340 } else {
6341 format!("{}-retry{attempt}", src.stem)
6342 },
6343 });
6344 }
6345
6346 if attempt > 0 {
6347 let seats_out: Vec<&str> = pending
6348 .iter()
6349 .map(|&i| originals[i].seat.key.as_str())
6350 .collect();
6351 state.event(
6352 ctx.node,
6353 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
6354 );
6355 }
6356 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
6357 let mut still = Vec::new();
6358 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
6359 seats[i] = seat;
6360 let (parsed, quota) = match out {
6361 AgentOutcome::Ok(o) => (
6362 match verdict::extract_json::<T>(&o.text) {
6363 Ok(v) => match validate(&v) {
6364 Ok(()) => Ok((v, o)),
6365 Err(e) => Err(e),
6366 },
6367 Err(e) => Err(e),
6368 },
6369 false,
6370 ),
6371 AgentOutcome::Quota(o) => {
6372 losses.push(QuotaLoss {
6373 seat: originals[i].seat.key.clone(),
6374 node: ctx.node.to_owned(),
6375 at: Timestamp::now(),
6376 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
6377 });
6378 (
6379 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
6380 true,
6381 )
6382 }
6383 AgentOutcome::Dropped(o) => {
6388 let why = o
6389 .dropped
6390 .as_ref()
6391 .map(|d| d.why.as_str())
6392 .unwrap_or("the CLI ended the stream without delivering its answer");
6393 (
6394 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
6395 false,
6396 )
6397 }
6398 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
6399 };
6400 let failed = parsed.is_err();
6401 done[i] = Some(parsed);
6402 attempts_used[i] = attempt;
6403 if failed && !quota {
6406 still.push(i);
6407 }
6408 }
6409 pending = still;
6410 }
6411
6412 seats
6413 .into_iter()
6414 .zip(done)
6415 .zip(attempts_used)
6416 .map(|((seat, res), attempts)| {
6417 (
6418 seat,
6419 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6420 attempts,
6421 )
6422 })
6423 .collect()
6424}
6425
6426async fn acquire_cache_lease(
6439 state: &mut RunState,
6440 cache_dir: &Path,
6441 owner: &crate::cache::Owner,
6442 budget: Duration,
6443 context: &str,
6444) -> Result<crate::cache::Guard> {
6445 let home = crate::run::home();
6446 let started = Instant::now();
6447 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6448 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6449 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6450 Err(e) => {
6451 state.event(
6452 "verify",
6453 format!("{context}: could not check the shared build cache: {e:#}"),
6454 );
6455 if let Err(e2) = state.save() {
6456 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6457 }
6458 return Err(e);
6459 }
6460 };
6461 state.event(
6462 "verify",
6463 format!(
6464 "{context}: waiting for the shared build cache at {} ({})",
6465 cache_dir.display(),
6466 busy.describe()
6467 ),
6468 );
6469 if let Err(e) = state.save() {
6470 tracing::warn!("could not persist a cache wait: {e:#}");
6471 }
6472 let remaining = budget.saturating_sub(started.elapsed());
6473 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6474 Ok(g) => Ok(g),
6475 Err(e) => {
6476 state.event("verify", format!("{context}: {e:#}"));
6477 if let Err(e2) = state.save() {
6478 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6479 }
6480 Err(e)
6481 }
6482 }
6483}
6484
6485#[allow(clippy::too_many_arguments)]
6506async fn with_cache_lease<'s, F, Fut>(
6507 state: &'s mut RunState,
6508 cache_dir: Option<&Path>,
6509 node: &str,
6510 seat: &str,
6511 worktree: &Path,
6512 head: &str,
6513 budget: Duration,
6514 context: &str,
6515 body: F,
6516) -> (Vec<CommandOutcome>, bool)
6517where
6518 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6519 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6520{
6521 let Some(cache_dir) = cache_dir else {
6522 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6523 return (outcomes, retried);
6524 };
6525 let home = crate::run::home();
6526 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6527 let started = Instant::now();
6528 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6529 Ok(g) => g,
6530 Err(e) => {
6531 return (
6532 vec![CommandOutcome {
6533 command: "(waiting for the shared build cache)".to_owned(),
6534 code: None,
6535 output_tail: e.to_string(),
6536 duration_ms: started.elapsed().as_millis() as u64,
6537 resource_blocked: true,
6538 }],
6539 false,
6540 );
6541 }
6542 };
6543 let identity = crate::cache::Identity::new(worktree, head);
6544 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6545 state.event(
6553 "verify",
6554 format!(
6555 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6556 worktree.display(),
6557 short(head)
6558 ),
6559 );
6560 guard.release();
6561 return (
6562 vec![CommandOutcome {
6563 command: "(confirming the shared build cache is fresh)".to_owned(),
6564 code: None,
6565 output_tail: e.to_string(),
6566 duration_ms: started.elapsed().as_millis() as u64,
6567 resource_blocked: true,
6568 }],
6569 false,
6570 );
6571 }
6572 let remaining = budget.saturating_sub(started.elapsed());
6573 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6574 if !timed_out_pids.is_empty() {
6579 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6580 }
6581 guard.release();
6582 (outcomes, retried)
6583}
6584
6585async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6597 wait_for_pids_with(
6598 pids,
6599 crate::proc::pid_alive,
6600 LEASE_RELEASE_POLL,
6601 LEASE_RELEASE_MAX_WAIT,
6602 )
6603 .await;
6604}
6605
6606async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6612 pids: &[u32],
6613 alive: F,
6614 poll: Duration,
6615 max_wait: Duration,
6616) {
6617 let deadline = Instant::now() + max_wait;
6618 loop {
6619 if pids.iter().all(|&pid| !alive(pid)) {
6620 return;
6621 }
6622 if Instant::now() >= deadline {
6623 return;
6624 }
6625 tokio::time::sleep(poll).await;
6626 }
6627}
6628
6629fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6636 outcomes.iter().any(|o| o.resource_blocked)
6637}
6638
6639enum GateFix {
6641 Retry,
6643 Stop,
6646 Defer,
6649}
6650
6651fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6659 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6660 red.peek().is_some()
6661 && red.all(|o| {
6662 !o.resource_blocked
6663 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6664 && !o.output_tail.trim().is_empty()
6665 })
6666}
6667
6668fn e2e_outcome_label(o: &CommandOutcome) -> String {
6672 if o.ok() {
6673 return "pass".to_owned();
6674 }
6675 let reason = if o.build_failed() {
6676 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6677 } else {
6678 format!("FAIL ({:?})", o.code)
6679 };
6680 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6681}
6682
6683async fn run_e2e_with_retry(
6691 state: &mut RunState,
6692 shell: &[String],
6693 commands: &[String],
6694 worktree: &Path,
6695 timeout: Duration,
6696 context: &str,
6697) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6698 let (mut e2e, mut timed_out_pids) = run_commands(
6699 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6700 )
6701 .await;
6702 for o in &e2e {
6703 state.event(
6704 "verify",
6705 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6706 );
6707 }
6708 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6712 if verify_retried {
6713 state.event(
6714 "verify",
6715 format!(
6716 "{context}: verify could not build/link, not a test result — retrying once \
6717 before concluding"
6718 ),
6719 );
6720 let retried = run_commands(
6721 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6722 )
6723 .await;
6724 e2e = retried.0;
6725 timed_out_pids.extend(retried.1);
6728 for o in &e2e {
6729 state.event(
6730 "verify",
6731 format!(
6732 "{context}: retry `{}` -> {}",
6733 o.command,
6734 e2e_outcome_label(o)
6735 ),
6736 );
6737 }
6738 }
6739 (e2e, verify_retried, timed_out_pids)
6740}
6741
6742#[allow(clippy::too_many_arguments)]
6758async fn run_commands(
6759 state: &mut RunState,
6760 node: &str,
6761 task: &str,
6762 attempt: usize,
6763 shell: &[String],
6764 commands: &[String],
6765 cwd: &Path,
6766 timeout: Duration,
6767) -> (Vec<CommandOutcome>, Vec<u32>) {
6768 if commands.is_empty() {
6769 return (Vec::new(), Vec::new());
6774 }
6775 let mut out = Vec::new();
6776 let mut timed_out_pids = Vec::new();
6777 let total = commands.len();
6778 for (idx, command) in commands.iter().enumerate() {
6779 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6780 if let Err(e) = state.save() {
6781 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6782 }
6783 let started = Instant::now();
6784 let mut cmd = tokio::process::Command::new(&shell[0]);
6785 cmd.quiet();
6786 cmd.args(&shell[1..])
6787 .arg(command)
6788 .current_dir(cwd)
6789 .stdin(std::process::Stdio::null())
6790 .stdout(std::process::Stdio::piped())
6791 .stderr(std::process::Stdio::piped())
6792 .kill_on_drop(true);
6793 let spawned = cmd.spawn();
6794 let (code, body) = match spawned {
6795 Ok(child) => {
6796 let pid = child.id();
6801 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6802 Ok(Ok(o)) => {
6803 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6804 body.push_str(&String::from_utf8_lossy(&o.stderr));
6805 (o.status.code(), body)
6806 }
6807 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6808 Err(_) => {
6809 if let Some(pid) = pid {
6810 timed_out_pids.push(pid);
6811 }
6812 (None, format!("timed out after {}s", timeout.as_secs()))
6813 }
6814 }
6815 }
6816 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6817 };
6818 out.push(CommandOutcome {
6819 command: command.clone(),
6820 code,
6821 output_tail: tail(&body, OUTPUT_TAIL),
6822 duration_ms: started.elapsed().as_millis() as u64,
6823 resource_blocked: false,
6824 });
6825 }
6826 state.task_finished(task);
6827 if let Err(e) = state.save() {
6828 tracing::warn!("could not persist the end of {task}: {e:#}");
6829 }
6830 (out, timed_out_pids)
6831}
6832
6833fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6845 let repo = repo.display();
6846 match style {
6847 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6848 MergeStyle::Squash => {
6849 let subject = message
6852 .lines()
6853 .next()
6854 .unwrap_or(branch)
6855 .replace(['\\', '"', '$', '`'], "");
6856 format!(
6857 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6858 )
6859 }
6860 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6861 }
6862}
6863
6864const PR_TITLE_MAX: usize = 240;
6875
6876struct PrMessage {
6881 title: String,
6882 body: String,
6883}
6884
6885impl PrMessage {
6886 fn commit_message(&self) -> String {
6890 format!("{}\n\n{}", self.title, self.body)
6891 }
6892}
6893
6894fn title_marker(line: &str) -> Option<&str> {
6896 let line = line.trim();
6897 let head = line.get(..6)?;
6898 head.eq_ignore_ascii_case("title:")
6899 .then(|| line[6..].trim())
6900}
6901
6902fn summary_title(summary: &str) -> Option<String> {
6907 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6908 let raw = title_marker(first)?;
6909 if raw.is_empty() {
6910 return None;
6911 }
6912 let title = queue::title_from(raw, PR_TITLE_MAX);
6913 let lower = title.to_ascii_lowercase();
6914 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6915 return None;
6916 }
6917 Some(title)
6918}
6919
6920fn summary_without_title(summary: &str) -> String {
6923 let mut lines = summary.trim().lines().peekable();
6924 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6925 lines.next();
6926 }
6927 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6928}
6929
6930fn pr_message(state: &RunState, winner: char) -> PrMessage {
6944 let summary = state
6945 .candidates
6946 .iter()
6947 .find(|c| c.label == winner)
6948 .map(|c| c.summary.as_str())
6949 .unwrap_or_default();
6950 let title = summary_title(summary).unwrap_or_else(|| {
6953 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6954 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6955 t
6956 } else {
6957 format!(
6958 "chore: land candidate {} of run {}",
6959 winner.to_ascii_uppercase(),
6960 state.id
6961 )
6962 }
6963 });
6964
6965 let mut body = String::new();
6966 let what = summary_without_title(summary);
6967 if !what.is_empty() {
6968 body.push_str("## Summary\n\n");
6969 body.push_str(&what);
6970 body.push_str("\n\n");
6971 }
6972
6973 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6974 if let Some(fix) = fix
6975 && !fix.notes.trim().is_empty()
6976 {
6977 body.push_str("## Review fixes\n\n");
6978 body.push_str(fix.notes.trim());
6979 body.push_str("\n\n");
6980 }
6981
6982 let open = state.open_findings();
6983 if !open.is_empty() {
6984 body.push_str("## Open review findings\n\n");
6985 for f in &open {
6986 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6987 }
6988 body.push('\n');
6989 }
6990
6991 if let Some(fix) = fix
6992 && !fix.rejected.is_empty()
6993 {
6994 body.push_str("## Declined by the fixer\n\n");
6995 for r in &fix.rejected {
6996 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
6997 }
6998 body.push('\n');
6999 }
7000
7001 let task = state.instruction.trim();
7002 let task = if task.is_empty() {
7003 "(empty task)"
7004 } else {
7005 task
7006 };
7007 body.push_str(&format!(
7008 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
7009 task.replace("</details>", "</details>")
7010 ));
7011
7012 body.push_str(&format!(
7013 "\n---\nmagi:run/{} magi:candidate-{}\n",
7014 state.id,
7015 winner.to_ascii_lowercase()
7016 ));
7017
7018 let id = crate::scrub::Identity::current();
7021 PrMessage {
7022 title: crate::scrub::scrub(&title, &id),
7023 body: crate::scrub::scrub(&body, &id),
7024 }
7025}
7026
7027fn seeded_instruction(state: &RunState) -> String {
7031 match refs::describe(&state.seeds) {
7032 Some(facts) => format!(
7033 "{}\n\n# Existing work the task refers to\n\n{facts}\n\n\
7034 Candidates start from the unmerged branch named above, when there \
7035 is one, and carry any unmerged commit named by sha as a \
7036 cherry-pick. Check that this is what the task meant before \
7037 building on it.",
7038 state.instruction
7039 ),
7040 None => state.instruction.clone(),
7041 }
7042}
7043
7044async fn merge_is_empty(repo: &Path, state: &RunState, branch: &str, mode: MergeMode) -> bool {
7048 let base = &state.base_branch;
7049 let mut against = base.clone();
7050 if mode == MergeMode::Pr {
7051 let remote = &state.config.merge.remote;
7052 let tracking = format!("{remote}/{base}");
7053 let fetched = git::fetch(repo, remote, base).await;
7054 if fetched.is_ok_and(|o| o.ok()) && git::rev_exists(repo, &tracking).await {
7055 against = tracking;
7056 }
7057 }
7058 matches!(git::commits_ahead(repo, &against, branch).await, Ok(0))
7059}
7060
7061fn empty_candidate_detail(state: &RunState, base: &str) -> String {
7064 let mut detail = format!(
7065 "empty candidate: the winning branch has 0 commits ahead of {base}, so there is \
7066 nothing to open a pull request for"
7067 );
7068 match refs::describe(&state.seeds) {
7069 Some(facts) => detail.push_str(&format!("\nReferences in the task:\n{facts}")),
7070 None => detail.push_str(
7071 "\nThe task names no existing branch or commit; if it means to land work \
7072 that lives elsewhere, name the branch (magi/<run>/<label>) or the sha.",
7073 ),
7074 }
7075 detail
7076}
7077
7078#[derive(Debug, PartialEq, Eq)]
7081enum PrPlan {
7082 Create,
7083 Adopt { url: String, title: String },
7084 Stop(String),
7085}
7086
7087fn pr_merge_plan(found: Result<land::OpenPr>) -> PrPlan {
7090 match found {
7091 Ok(land::OpenPr::None) => PrPlan::Create,
7092 Ok(land::OpenPr::One { url, title }) => PrPlan::Adopt { url, title },
7093 Ok(land::OpenPr::Many(urls)) => PrPlan::Stop(format!(
7094 "several open pull requests exist for this branch, not picking one: {}",
7095 urls.join(" ")
7096 )),
7097 Err(e) => PrPlan::Stop(format!("could not look up open pull requests: {e:#}")),
7098 }
7099}
7100
7101async fn gh_pr_create(
7103 cwd: &Path,
7104 base: &str,
7105 head: &str,
7106 title: &str,
7107 body: &str,
7108) -> Result<String> {
7109 let out = tokio::process::Command::new("gh")
7110 .args([
7111 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
7112 ])
7113 .current_dir(cwd)
7114 .quiet()
7115 .stdin(std::process::Stdio::null())
7116 .output()
7117 .await
7118 .context("spawn gh")?;
7119 if out.status.success() {
7120 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
7121 } else {
7122 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
7123 }
7124}
7125
7126pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
7135 let repo = state.repo.clone();
7136 let root = state.worktree_root();
7137 let winner = state.tally.as_ref().map(|t| t.winner);
7138 let mut removed = Vec::new();
7139
7140 for i in 0..state.candidates.len() {
7141 let c = state.candidates[i].clone();
7142 let is_winner = Some(c.label) == winner;
7143 if is_winner && !drop_winner {
7144 continue;
7145 }
7146 if c.worktree.exists() {
7147 git::worktree_remove(&repo, &c.worktree).await.ok();
7148 removed.push(c.worktree.to_string_lossy().into_owned());
7149 }
7150 let handed_over = state.released_branches.contains(&c.branch);
7153 if !handed_over && git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
7154 git::branch_delete(&repo, &c.branch).await.ok();
7155 removed.push(c.branch.clone());
7156 }
7157 state.candidates[i].folded = true;
7158 }
7159
7160 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
7161 let path = name.path();
7162 let keep = !drop_winner
7163 && winner.is_some_and(|w| {
7164 path.file_name()
7165 .is_some_and(|n| n == format!("cand-{w}").as_str())
7166 });
7167 if keep {
7168 continue;
7169 }
7170 git::worktree_remove(&repo, &path).await.ok();
7171 removed.push(path.to_string_lossy().into_owned());
7172 }
7173
7174 remove_if_empty(&root);
7183
7184 if state.enabled_worktree_config && drop_winner {
7185 git::release_worktree_config(&repo).await.ok();
7189 state.enabled_worktree_config = false;
7190 }
7191 state.save_under(home)?;
7192 Ok(removed)
7193}
7194
7195fn remove_if_empty(dir: &Path) {
7206 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
7207 std::fs::remove_dir(dir).ok();
7208 }
7209}
7210
7211pub fn worst_open(state: &RunState) -> Option<Severity> {
7213 state
7214 .reviews
7215 .last()?
7216 .reviews
7217 .iter()
7218 .flat_map(|r| r.findings.iter())
7219 .map(|f| f.severity)
7220 .max()
7221}
7222
7223#[cfg(test)]
7224mod tests {
7225 #[test]
7226 fn pr_merge_plan_creates_adopts_or_stops() {
7227 assert_eq!(pr_merge_plan(Ok(land::OpenPr::None)), PrPlan::Create);
7228 assert_eq!(
7229 pr_merge_plan(Ok(land::OpenPr::One {
7230 url: "u".into(),
7231 title: "t".into()
7232 })),
7233 PrPlan::Adopt {
7234 url: "u".into(),
7235 title: "t".into()
7236 }
7237 );
7238 let PrPlan::Stop(many) =
7239 pr_merge_plan(Ok(land::OpenPr::Many(vec!["a".into(), "b".into()])))
7240 else {
7241 panic!("many must stop");
7242 };
7243 assert!(many.contains('a') && many.contains('b'));
7244 let PrPlan::Stop(err) = pr_merge_plan(Err(anyhow::anyhow!("bad token"))) else {
7245 panic!("a failed lookup must stop");
7246 };
7247 assert!(err.contains("bad token"));
7248 }
7249
7250 use super::*;
7251 use crate::run::GateStatus;
7252 use std::collections::BTreeMap;
7253 use std::time::Duration;
7254
7255 fn conductor() -> AgentSpec {
7256 AgentSpec {
7257 id: "conductor".to_owned(),
7258 kind: crate::config::AgentKind::Command,
7259 model: None,
7260 command: vec!["true".to_owned()],
7261 extra_args: Vec::new(),
7262 env: BTreeMap::new(),
7263 prompt_delivery: None,
7264 }
7265 }
7266
7267 fn spec(id: &str) -> AgentSpec {
7268 AgentSpec {
7269 id: id.to_owned(),
7270 kind: crate::config::AgentKind::Command,
7271 model: None,
7272 command: vec!["true".to_owned()],
7273 extra_args: Vec::new(),
7274 env: BTreeMap::new(),
7275 prompt_delivery: None,
7276 }
7277 }
7278
7279 #[test]
7286 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
7287 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7288 let tried = BTreeSet::from(["beta".to_owned()]);
7289 let next = next_untried_implementer(&roster, 1, &tried);
7292 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
7293 }
7294
7295 #[test]
7296 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
7297 let roster = vec![spec("alpha"), spec("beta")];
7298 let tried = BTreeSet::from(["beta".to_owned()]);
7299 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7303 }
7304
7305 #[test]
7306 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
7307 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7308 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
7309 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7313 }
7314
7315 #[test]
7316 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
7317 let roster = vec![spec("a"), spec("a"), spec("b")];
7318 let tried = BTreeSet::from(["a".to_owned()]);
7319 let next = next_untried_implementer(&roster, 0, &tried);
7320 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
7321 }
7322
7323 #[test]
7324 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
7325 let roster = vec![spec("a"), spec("b")];
7326 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
7327 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
7328 }
7329
7330 #[test]
7331 fn remove_if_empty_only_ever_takes_a_bare_directory() {
7332 let dir = tempfile::tempdir().unwrap();
7333 let bay = dir.path().join("ffff");
7334
7335 remove_if_empty(&bay);
7337 assert!(!bay.exists());
7338
7339 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
7342 remove_if_empty(&bay);
7343 assert!(bay.exists(), "non-empty directory must survive");
7344
7345 std::fs::remove_dir(bay.join("cand-A")).unwrap();
7347 remove_if_empty(&bay);
7348 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
7349 }
7350
7351 #[test]
7360 fn a_full_panel_that_found_nothing_is_clean() {
7361 assert!(round_is_clean(
7362 0,
7363 true,
7364 2,
7365 2,
7366 0,
7367 IncompleteReviewPolicy::Block
7368 ));
7369 }
7370
7371 #[test]
7372 fn a_missing_seat_is_never_clean_under_the_default_policy() {
7373 assert!(!round_is_clean(
7374 0,
7375 true,
7376 1,
7377 2,
7378 0,
7379 IncompleteReviewPolicy::Block
7380 ));
7381 }
7382
7383 #[test]
7384 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
7385 assert!(!round_is_clean(
7386 1,
7387 true,
7388 1,
7389 2,
7390 0,
7391 IncompleteReviewPolicy::Warn
7392 ));
7393 }
7394
7395 #[test]
7396 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
7397 assert!(round_is_clean(
7398 0,
7399 true,
7400 1,
7401 2,
7402 0,
7403 IncompleteReviewPolicy::Warn
7404 ));
7405 }
7406
7407 #[test]
7408 fn a_full_panel_with_an_open_finding_is_not_clean() {
7409 assert!(!round_is_clean(
7410 1,
7411 true,
7412 2,
7413 2,
7414 0,
7415 IncompleteReviewPolicy::Block
7416 ));
7417 }
7418
7419 #[test]
7420 fn a_full_panel_with_a_red_e2e_is_not_clean() {
7421 assert!(!round_is_clean(
7422 0,
7423 false,
7424 2,
7425 2,
7426 0,
7427 IncompleteReviewPolicy::Block
7428 ));
7429 }
7430
7431 #[test]
7438 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
7439 assert!(round_is_clean(
7442 0,
7443 true,
7444 1,
7445 2,
7446 1,
7447 IncompleteReviewPolicy::Block
7448 ));
7449 }
7450
7451 #[test]
7452 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
7453 assert!(!round_is_clean(
7456 0,
7457 true,
7458 1,
7459 2,
7460 0,
7461 IncompleteReviewPolicy::Block
7462 ));
7463 }
7464
7465 #[test]
7466 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
7467 assert!(!round_is_clean(
7468 1,
7469 true,
7470 1,
7471 2,
7472 1,
7473 IncompleteReviewPolicy::Block
7474 ));
7475 assert!(!round_is_clean(
7476 0,
7477 false,
7478 1,
7479 2,
7480 1,
7481 IncompleteReviewPolicy::Block
7482 ));
7483 }
7484
7485 #[test]
7486 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
7487 assert!(!round_is_clean(
7491 0,
7492 true,
7493 0,
7494 2,
7495 2,
7496 IncompleteReviewPolicy::Block
7497 ));
7498 }
7499
7500 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
7501 CommandOutcome {
7502 command: "test".to_owned(),
7503 code,
7504 output_tail: String::new(),
7505 duration_ms: 0,
7506 resource_blocked,
7507 }
7508 }
7509
7510 #[test]
7511 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7512 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7513 assert!(
7514 !verify_inconclusive(&[outcome(Some(1), false)]),
7515 "an ordinary failure is still evidence about the patch"
7516 );
7517 assert!(verify_inconclusive(&[outcome(None, true)]));
7518 assert!(
7519 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7520 "one inconclusive outcome taints the whole batch"
7521 );
7522 assert!(!verify_inconclusive(&[]));
7523 }
7524
7525 #[tokio::test]
7526 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7527 let calls = std::sync::atomic::AtomicUsize::new(0);
7531 let started = Instant::now();
7532 wait_for_pids_with(
7533 &[123],
7534 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7535 Duration::from_millis(5),
7536 Duration::from_secs(5),
7537 )
7538 .await;
7539 assert!(
7540 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7541 "must keep checking rather than deciding on the first answer"
7542 );
7543 assert!(
7544 started.elapsed() < Duration::from_secs(1),
7545 "must return the moment it is confirmed dead, not wait out the ceiling"
7546 );
7547 }
7548
7549 #[tokio::test]
7550 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7551 let started = Instant::now();
7552 wait_for_pids_with(
7553 &[123],
7554 |_| true, Duration::from_millis(5),
7556 Duration::from_millis(30),
7557 )
7558 .await;
7559 let elapsed = started.elapsed();
7560 assert!(
7561 elapsed >= Duration::from_millis(30),
7562 "must not give up before its own ceiling: {elapsed:?}"
7563 );
7564 assert!(
7565 elapsed < Duration::from_secs(1),
7566 "must not wait past its own ceiling either: {elapsed:?}"
7567 );
7568 }
7569
7570 #[tokio::test]
7571 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7572 let started = Instant::now();
7573 wait_for_pids_with(
7574 &[],
7575 |_| true,
7576 Duration::from_secs(5),
7577 Duration::from_secs(5),
7578 )
7579 .await;
7580 assert!(
7581 started.elapsed() < Duration::from_millis(200),
7582 "an empty pid list has nothing to confirm"
7583 );
7584 }
7585
7586 fn review_round(
7592 clean: bool,
7593 blocking: usize,
7594 answered: usize,
7595 expected: usize,
7596 progressed: bool,
7597 e2e_ok: bool,
7598 ) -> ReviewRound {
7599 ReviewRound {
7600 round: 1,
7601 head: "h".to_owned(),
7602 verified_head: None,
7603 verified_at: None,
7604 reviews: Vec::new(),
7605 e2e: vec![CommandOutcome {
7606 command: "test".to_owned(),
7607 code: Some(if e2e_ok { 0 } else { 1 }),
7608 output_tail: String::new(),
7609 duration_ms: 0,
7610 resource_blocked: false,
7611 }],
7612 verify_retried: false,
7613 e2e_deferred: false,
7614 e2e_defer_reason: None,
7615 fix: None,
7616 blocking,
7617 answered,
7618 expected,
7619 clean,
7620 progressed,
7621 vote_split: false,
7622 reconsideration: Vec::new(),
7623 verdict: None,
7624 }
7625 }
7626
7627 #[test]
7628 fn review_conclusion_is_none_when_nothing_has_run() {
7629 assert_eq!(review_conclusion(&[], 3), None);
7630 }
7631
7632 #[test]
7633 fn review_conclusion_is_none_while_rounds_remain() {
7634 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7635 assert_eq!(review_conclusion(&rounds, 3), None);
7636 }
7637
7638 #[test]
7639 fn review_conclusion_is_gating_once_a_round_is_clean() {
7640 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7641 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7642 }
7643
7644 #[test]
7645 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7646 let rounds = vec![
7647 review_round(false, 1, 2, 2, true, true),
7648 review_round(false, 1, 2, 2, true, true),
7649 ];
7650 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7651 }
7652
7653 #[test]
7654 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7655 let rounds = vec![
7656 review_round(false, 1, 2, 2, true, true),
7657 review_round(false, 1, 2, 2, true, false),
7658 ];
7659 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7660 }
7661
7662 #[test]
7663 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7664 let mut blocked = review_round(false, 1, 2, 2, true, false);
7671 blocked.e2e[0].resource_blocked = true;
7672 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7673 assert_eq!(review_conclusion(&rounds, 2), None);
7674 }
7675
7676 #[test]
7677 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7678 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7680 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7681 }
7682
7683 #[test]
7684 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7685 let rounds = vec![
7686 review_round(false, 1, 2, 2, false, true),
7687 review_round(false, 1, 2, 2, false, true),
7688 ];
7689 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7690 }
7691
7692 fn secs(n: u64) -> Duration {
7693 Duration::from_secs(n)
7694 }
7695
7696 fn init_repo(dir: &Path) {
7699 let run = |args: &[&str]| {
7700 let out = std::process::Command::new("git")
7701 .args(args)
7702 .current_dir(dir)
7703 .quiet()
7704 .output()
7705 .expect("spawn git");
7706 assert!(
7707 out.status.success(),
7708 "git {args:?} failed: {}",
7709 String::from_utf8_lossy(&out.stderr)
7710 );
7711 };
7712 run(&["init", "-b", "main"]);
7713 run(&["config", "user.name", "magi test"]);
7714 run(&["config", "user.email", "magi@example.com"]);
7715 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7716 run(&["add", "-A"]);
7717 run(&["commit", "-m", "init"]);
7718 }
7719
7720 fn ask_test_home() {
7728 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7729 }
7730
7731 fn runner_at(status: RunStatus) -> Runner {
7734 let mut state = RunState::new(
7735 PathBuf::from("/nonexistent/repo"),
7736 "main".to_owned(),
7737 "deadbeef".to_owned(),
7738 "task".to_owned(),
7739 Config::default(),
7740 );
7741 state.status = status;
7742 Runner {
7743 state,
7744 roles: ResolvedRoles {
7745 implementers: Vec::new(),
7746 judges: Vec::new(),
7747 reviewers: Vec::new(),
7748 fixer: None,
7749 conductor: conductor(),
7750 implementer_roster: Vec::new(),
7751 },
7752 sem: Arc::new(Semaphore::new(1)),
7753 pause: Pause::new(),
7754 interrupt: Pause::new(),
7755 }
7756 }
7757
7758 #[test]
7762 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7763 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7764 let mut runner = runner_at(RunStatus::Implementing);
7765 let interrupt = Pause::new();
7766 runner.watch_interrupt(interrupt.clone());
7767
7768 interrupt.park_because("task a1b2 asked to run first");
7769
7770 assert!(runner.park_here().expect("park_here"));
7771 assert!(runner.state.parked);
7772 let last = runner.state.events.last().expect("a park event");
7773 assert_eq!(last.node, "park");
7774 assert!(
7775 last.message.contains("task a1b2 asked to run first"),
7776 "expected the interrupt reason in {:?}",
7777 last.message
7778 );
7779 }
7780
7781 #[test]
7789 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7790 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7791 let mut runner = runner_at(RunStatus::Implementing);
7792 let shutdown = Pause::new();
7793 runner.on_pause(shutdown.clone());
7794 let interrupt = Pause::new();
7795 runner.watch_interrupt(interrupt.clone());
7796
7797 assert!(!runner.park_here().expect("park_here"));
7799 assert!(!runner.state.parked);
7800
7801 interrupt.park_because("test");
7803 assert!(!shutdown.parked());
7804 assert!(runner.park_here().expect("park_here"));
7805 }
7806
7807 #[tokio::test]
7821 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7822 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7823 let mut runner = runner_at(RunStatus::Implementing);
7824 let interrupt = Pause::new();
7825 runner.watch_interrupt(interrupt.clone());
7826
7827 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7828 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7829
7830 let node = async move {
7834 started_tx.send(()).expect("send started");
7835 finish_rx.await.expect("recv finish");
7836 "node finished"
7837 };
7838
7839 let interrupter = async move {
7840 started_rx.await.expect("recv started");
7841 interrupt.park_because("higher-priority task waiting");
7843 tokio::task::yield_now().await;
7847 finish_tx.send(()).expect("send finish");
7848 };
7849
7850 let (node_result, ()) = tokio::join!(node, interrupter);
7851 assert_eq!(
7852 node_result, "node finished",
7853 "the in-flight call ran to completion"
7854 );
7855
7856 assert!(runner.park_here().expect("park_here"));
7859 assert!(runner.state.parked);
7860 }
7861
7862 #[test]
7868 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7869 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7870 let mut runner = runner_at(RunStatus::Judging);
7871 runner.state.config.agents = vec![conductor()];
7875 runner.state.candidates = vec![Candidate {
7876 index: 0,
7877 label: 'A',
7878 agent: "alpha".to_owned(),
7879 branch: "magi/x/A".to_owned(),
7880 worktree: PathBuf::from("/nonexistent/worktree"),
7881 summary: "did the thing".to_owned(),
7882 stat: "1 file changed".to_owned(),
7883 files: 1,
7884 commits: 1,
7885 empty: false,
7886 failed: None,
7887 verified_noop: None,
7888 duration_ms: 1234,
7889 folded: false,
7890 }];
7891 let run_id = runner.state.id.clone();
7892
7893 let interrupt = Pause::new();
7894 runner.watch_interrupt(interrupt.clone());
7895 interrupt.park_because("task c3d4 asked to run first");
7896 assert!(runner.park_here().expect("park_here"));
7897
7898 let resumed = Runner::resume(&run_id).expect("resume");
7899 assert_eq!(resumed.state.candidates.len(), 1);
7900 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7901 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7902 assert_eq!(resumed.state.status, runner.state.status);
7903 assert!(
7904 resumed.state.parked,
7905 "still parked until `execute` actually walks the graph again"
7906 );
7907 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7908 }
7909
7910 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7912 let mut q = ask::Question::new(
7913 run.to_owned(),
7914 "implement".to_owned(),
7915 "impl-A".to_owned(),
7916 "Which storage backend should the cache use?".to_owned(),
7917 String::new(),
7918 vec!["SQLite".to_owned(), "Redis".to_owned()],
7919 );
7920 store.put(&mut q).unwrap();
7921 q
7922 }
7923
7924 #[test]
7925 fn a_failed_runs_open_question_is_abandoned() {
7926 ask_test_home();
7927 let store = ask::Questions::open();
7928 let mut runner = runner_at(RunStatus::Failed);
7929 let run = runner.state.id.clone();
7930 let q = ask_open_question(&store, &run);
7931
7932 runner.settle_questions();
7933
7934 let back = store.get(&q.id).unwrap();
7935 assert!(
7936 !back.status.open(),
7937 "the seat that asked died with the run; nobody is left to read an answer"
7938 );
7939 assert!(
7940 back.detail.contains(&run) && back.detail.contains("failed"),
7941 "the reason names what the run became, not just that it is gone: {}",
7942 back.detail
7943 );
7944 }
7945
7946 #[test]
7947 fn a_merged_runs_open_question_is_abandoned_too() {
7948 ask_test_home();
7949 let store = ask::Questions::open();
7950 for status in [RunStatus::Merged, RunStatus::Ready] {
7953 let mut runner = runner_at(status);
7954 let run = runner.state.id.clone();
7955 let q = ask_open_question(&store, &run);
7956
7957 runner.settle_questions();
7958
7959 let back = store.get(&q.id).unwrap();
7960 assert!(
7961 !back.status.open(),
7962 "{status:?} run's question must not outlive the run"
7963 );
7964 }
7965 }
7966
7967 #[test]
7968 fn a_still_resumable_runs_open_question_is_left_alone() {
7969 ask_test_home();
7970 let store = ask::Questions::open();
7971 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7977 let mut runner = runner_at(status);
7978 let run = runner.state.id.clone();
7979 let q = ask_open_question(&store, &run);
7980
7981 runner.settle_questions();
7982
7983 let back = store.get(&q.id).unwrap();
7984 assert!(
7985 back.status.open(),
7986 "{status:?} is still alive; the question must still be waiting"
7987 );
7988 }
7989 }
7990
7991 #[test]
7992 fn settle_questions_never_touches_an_already_answered_question() {
7993 ask_test_home();
7994 let store = ask::Questions::open();
7995 let mut runner = runner_at(RunStatus::Failed);
7996 let run = runner.state.id.clone();
7997 let mut q = ask_open_question(&store, &run);
7998 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
7999 .unwrap();
8000 store.put(&mut q).unwrap();
8001
8002 runner.settle_questions();
8007 runner.settle_questions();
8008
8009 let back = store.get(&q.id).unwrap();
8010 assert_eq!(
8011 back.status,
8012 ask::QuestionStatus::Answered,
8013 "a real answer is a decision on record, never overwritten by a sweep"
8014 );
8015 }
8016
8017 #[tokio::test]
8028 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
8029 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
8030 let tmp = tempfile::tempdir().expect("tempdir");
8031 let repo = tmp.path().join("repo");
8032 std::fs::create_dir_all(&repo).unwrap();
8033 init_repo(&repo);
8034
8035 let mut config = Config::default();
8036 config.graph.worktree_root = Some(tmp.path().join("wt"));
8037
8038 let mut state = RunState::new(
8039 repo.clone(),
8040 "main".to_owned(),
8041 "deadbeef".to_owned(),
8042 "task".to_owned(),
8043 config,
8044 );
8045 let root = state.worktree_root();
8046 let wt_a = root.join("cand-A");
8047 let wt_b = root.join("cand-B");
8048 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
8049 .await
8050 .expect("worktree A");
8051 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
8052 .await
8053 .expect("worktree B");
8054
8055 state.candidates = vec![
8056 Candidate {
8057 index: 0,
8058 label: 'A',
8059 agent: "alpha".to_owned(),
8060 branch: "magi/x/A".to_owned(),
8061 worktree: wt_a.clone(),
8062 summary: String::new(),
8063 stat: String::new(),
8064 files: 0,
8065 commits: 0,
8066 empty: false,
8067 failed: None,
8068 verified_noop: None,
8069 duration_ms: 0,
8070 folded: false,
8071 },
8072 Candidate {
8073 index: 1,
8074 label: 'B',
8075 agent: "beta".to_owned(),
8076 branch: "magi/x/B".to_owned(),
8077 worktree: wt_b.clone(),
8078 summary: String::new(),
8079 stat: String::new(),
8080 files: 0,
8081 commits: 0,
8082 empty: false,
8083 failed: None,
8084 verified_noop: None,
8085 duration_ms: 0,
8086 folded: false,
8087 },
8088 ];
8089 state.tally = Some(Tally {
8090 first_choice: BTreeMap::from([('A', 1)]),
8091 borda: BTreeMap::new(),
8092 winner: 'A',
8093 rankings: 1,
8094 unanimous_initial: true,
8095 deliberated: false,
8096 changed_votes: 0,
8097 unanimous_final: true,
8098 tie_break: None,
8099 judges: 1,
8100 present: 1,
8101 quorum: 1,
8102 met_quorum: true,
8103 uncontested: None,
8104 });
8105 state.status = RunStatus::Ready;
8106
8107 fold_run(&mut state, false, &crate::run::home())
8108 .await
8109 .expect("fold_run");
8110
8111 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
8112 assert!(
8113 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
8114 "the unmerged winner's branch survives"
8115 );
8116 assert!(
8117 !state.candidates[0].folded,
8118 "the winner is not marked folded"
8119 );
8120
8121 assert!(!wt_b.exists(), "the loser's worktree is removed");
8122 assert!(
8123 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
8124 "the loser's branch is removed"
8125 );
8126 assert!(state.candidates[1].folded, "the loser is marked folded");
8127 }
8128
8129 #[tokio::test]
8132 async fn fold_run_keeps_a_branch_that_was_handed_to_a_later_run() {
8133 let tmp = tempfile::tempdir().expect("tempdir");
8134 let repo = tmp.path().join("repo");
8135 std::fs::create_dir_all(&repo).unwrap();
8136 init_repo(&repo);
8137 let home = tmp.path().join("home");
8138
8139 let mut config = Config::default();
8140 config.graph.worktree_root = Some(tmp.path().join("wt"));
8141 let mut state = RunState::new(
8142 repo.clone(),
8143 "main".to_owned(),
8144 "deadbeef".to_owned(),
8145 "task".to_owned(),
8146 config,
8147 );
8148 git::git(&repo, &["branch", "magi/x/A", "main"])
8150 .await
8151 .expect("branch");
8152 state.candidates = vec![Candidate {
8153 index: 0,
8154 label: 'A',
8155 agent: "alpha".to_owned(),
8156 branch: "magi/x/A".to_owned(),
8157 worktree: state.worktree_root().join("cand-A"),
8158 summary: String::new(),
8159 stat: String::new(),
8160 files: 0,
8161 commits: 0,
8162 empty: false,
8163 failed: None,
8164 verified_noop: None,
8165 duration_ms: 0,
8166 folded: true,
8167 }];
8168 state.released_to = Some("20260901-000000-new1".to_owned());
8169 state.released_branches = vec!["magi/x/A".to_owned()];
8170
8171 fold_run(&mut state, true, &home).await.expect("fold_run");
8172
8173 assert!(
8174 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
8175 "the handed-over branch survives a fold"
8176 );
8177 }
8178
8179 #[tokio::test]
8183 async fn an_empty_winner_is_detected_before_a_pull_request_is_attempted() {
8184 let tmp = tempfile::tempdir().expect("tempdir");
8185 let repo = tmp.path().join("repo");
8186 std::fs::create_dir_all(&repo).unwrap();
8187 init_repo(&repo);
8188 let run = |args: &[&str]| {
8189 let out = std::process::Command::new("git")
8190 .quiet()
8191 .args(args)
8192 .current_dir(&repo)
8193 .output()
8194 .expect("spawn git");
8195 assert!(out.status.success(), "git {args:?}");
8196 };
8197 run(&["branch", "magi/x/A"]);
8198 run(&["checkout", "-q", "-b", "magi/x/B"]);
8199 std::fs::write(repo.join("f.txt"), "x\n").unwrap();
8200 run(&["add", "-A"]);
8201 run(&["commit", "-q", "-m", "work"]);
8202 run(&["checkout", "-q", "main"]);
8203
8204 let mut state = RunState::new(
8205 repo.clone(),
8206 "main".to_owned(),
8207 "deadbeef".to_owned(),
8208 "task".to_owned(),
8209 Config::default(),
8210 );
8211 state.seeds = vec![refs::Seed {
8212 token: "magi/27b2/A".to_owned(),
8213 kind: refs::SeedKind::Unresolved,
8214 sha: String::new(),
8215 branch: true,
8216 detail: "no branch or commit named magi/27b2/A".to_owned(),
8217 }];
8218
8219 assert!(merge_is_empty(&repo, &state, "magi/x/A", MergeMode::Pr).await);
8220 assert!(merge_is_empty(&repo, &state, "magi/x/A", MergeMode::Local).await);
8221 assert!(!merge_is_empty(&repo, &state, "magi/x/B", MergeMode::Pr).await);
8222 let detail = empty_candidate_detail(&state, "main");
8223 assert!(detail.starts_with("empty candidate"), "{detail}");
8224 assert!(detail.contains("magi/27b2/A"), "{detail}");
8225 }
8226
8227 #[tokio::test]
8236 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
8237 let tmp = tempfile::tempdir().expect("tempdir");
8238 let repo = tmp.path().join("repo");
8239 std::fs::create_dir_all(&repo).unwrap();
8240 init_repo(&repo);
8241
8242 let mut config = Config::default();
8243 config.merge.mode = MergeMode::Local;
8244
8245 let mut state = RunState::new(
8246 repo.clone(),
8247 "main".to_owned(),
8248 "deadbeef".to_owned(),
8249 "task".to_owned(),
8250 config,
8251 );
8252 state.candidates = vec![Candidate {
8253 index: 0,
8254 label: 'A',
8255 agent: "alpha".to_owned(),
8256 branch: "does-not-exist".to_owned(),
8257 worktree: repo.clone(),
8258 summary: String::new(),
8259 stat: String::new(),
8260 files: 0,
8261 commits: 0,
8262 empty: false,
8263 failed: None,
8264 verified_noop: None,
8265 duration_ms: 0,
8266 folded: false,
8267 }];
8268 state.tally = Some(Tally {
8269 first_choice: BTreeMap::from([('A', 1)]),
8270 borda: BTreeMap::new(),
8271 winner: 'A',
8272 rankings: 1,
8273 unanimous_initial: true,
8274 deliberated: false,
8275 changed_votes: 0,
8276 unanimous_final: true,
8277 tie_break: None,
8278 judges: 0,
8279 present: 0,
8280 quorum: 0,
8281 met_quorum: true,
8282 uncontested: Some("only candidate A produced a change".to_owned()),
8283 });
8284 state.reviews = vec![ReviewRound {
8285 round: 1,
8286 head: "deadbeef".to_owned(),
8287 verified_head: None,
8288 verified_at: None,
8289 reviews: Vec::new(),
8290 e2e: Vec::new(),
8291 fix: None,
8292 blocking: 0,
8293 answered: 0,
8294 expected: 0,
8295 clean: true,
8296 verify_retried: false,
8297 e2e_deferred: false,
8298 e2e_defer_reason: None,
8299 progressed: false,
8300 vote_split: false,
8301 reconsideration: Vec::new(),
8302 verdict: None,
8303 }];
8304 state.gate = vec![CommandOutcome {
8305 command: "test".to_owned(),
8306 code: Some(0),
8307 output_tail: String::new(),
8308 duration_ms: 0,
8309 resource_blocked: false,
8310 }];
8311 state.gate_ran = true;
8312 state.status = RunStatus::Ready;
8317 state.merge = Some(MergeOutcome {
8318 mode: MergeMode::Local,
8319 ok: false,
8320 detail: "already concluded".to_owned(),
8321 empty: false,
8322 });
8323
8324 let mut runner = Runner {
8325 state,
8326 roles: ResolvedRoles {
8327 implementers: Vec::new(),
8328 judges: Vec::new(),
8329 reviewers: Vec::new(),
8330 fixer: None,
8331 conductor: conductor(),
8332 implementer_roster: Vec::new(),
8333 },
8334 sem: Arc::new(Semaphore::new(1)),
8335 pause: Pause::new(),
8336 interrupt: Pause::new(),
8337 };
8338
8339 runner.merge().await.expect("merge");
8340
8341 assert_eq!(
8342 runner.state.status,
8343 RunStatus::Ready,
8344 "a concluded run's status must not change on reentry"
8345 );
8346 assert_eq!(
8347 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8348 Some("already concluded"),
8349 "merge must not run again once the node already recorded an outcome"
8350 );
8351 }
8352
8353 #[tokio::test]
8362 async fn merge_refuses_a_gate_that_has_not_actually_run() {
8363 let tmp = tempfile::tempdir().expect("tempdir");
8364 let repo = tmp.path().join("repo");
8365 std::fs::create_dir_all(&repo).unwrap();
8366 init_repo(&repo);
8367
8368 let mut config = Config::default();
8369 config.merge.mode = MergeMode::Local;
8370
8371 let mut state = RunState::new(
8372 repo.clone(),
8373 "main".to_owned(),
8374 "deadbeef".to_owned(),
8375 "task".to_owned(),
8376 config,
8377 );
8378 state.candidates = vec![Candidate {
8379 index: 0,
8380 label: 'A',
8381 agent: "alpha".to_owned(),
8382 branch: "does-not-exist".to_owned(),
8383 worktree: repo.clone(),
8384 summary: String::new(),
8385 stat: String::new(),
8386 files: 0,
8387 commits: 0,
8388 empty: false,
8389 failed: None,
8390 verified_noop: None,
8391 duration_ms: 0,
8392 folded: false,
8393 }];
8394 state.tally = Some(Tally {
8395 first_choice: BTreeMap::from([('A', 1)]),
8396 borda: BTreeMap::new(),
8397 winner: 'A',
8398 rankings: 1,
8399 unanimous_initial: true,
8400 deliberated: false,
8401 changed_votes: 0,
8402 unanimous_final: true,
8403 tie_break: None,
8404 judges: 0,
8405 present: 0,
8406 quorum: 0,
8407 met_quorum: true,
8408 uncontested: Some("only candidate A produced a change".to_owned()),
8409 });
8410 state.reviews = vec![ReviewRound {
8411 round: 1,
8412 head: "deadbeef".to_owned(),
8413 verified_head: None,
8414 verified_at: None,
8415 reviews: Vec::new(),
8416 e2e: Vec::new(),
8417 fix: None,
8418 blocking: 0,
8419 answered: 0,
8420 expected: 0,
8421 clean: true,
8422 verify_retried: false,
8423 e2e_deferred: false,
8424 e2e_defer_reason: None,
8425 progressed: false,
8426 vote_split: false,
8427 reconsideration: Vec::new(),
8428 verdict: None,
8429 }];
8430 state.gate = Vec::new();
8432 state.gate_ran = false;
8433 state.status = RunStatus::Gating;
8434
8435 let mut runner = Runner {
8436 state,
8437 roles: ResolvedRoles {
8438 implementers: Vec::new(),
8439 judges: Vec::new(),
8440 reviewers: Vec::new(),
8441 fixer: None,
8442 conductor: conductor(),
8443 implementer_roster: Vec::new(),
8444 },
8445 sem: Arc::new(Semaphore::new(1)),
8446 pause: Pause::new(),
8447 interrupt: Pause::new(),
8448 };
8449
8450 runner.merge().await.expect("merge");
8451
8452 assert!(
8453 runner.state.merge.is_none(),
8454 "an empty gate must never be read as a passing one: {:?}",
8455 runner.state.merge
8456 );
8457 }
8458
8459 #[tokio::test]
8466 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
8467 let tmp = tempfile::tempdir().expect("tempdir");
8468 let repo = tmp.path().join("repo");
8469 std::fs::create_dir_all(&repo).unwrap();
8470 init_repo(&repo);
8471
8472 let config = Config::default();
8474
8475 let mut state = RunState::new(
8476 repo.clone(),
8477 "main".to_owned(),
8478 "deadbeef".to_owned(),
8479 "task".to_owned(),
8480 config,
8481 );
8482 state.candidates = vec![Candidate {
8483 index: 0,
8484 label: 'A',
8485 agent: "alpha".to_owned(),
8486 branch: "does-not-exist".to_owned(),
8487 worktree: repo.clone(),
8488 summary: String::new(),
8489 stat: String::new(),
8490 files: 0,
8491 commits: 0,
8492 empty: false,
8493 failed: None,
8494 verified_noop: None,
8495 duration_ms: 0,
8496 folded: false,
8497 }];
8498 state.tally = Some(Tally {
8499 first_choice: BTreeMap::from([('A', 1)]),
8500 borda: BTreeMap::new(),
8501 winner: 'A',
8502 rankings: 1,
8503 unanimous_initial: true,
8504 deliberated: false,
8505 changed_votes: 0,
8506 unanimous_final: true,
8507 tie_break: None,
8508 judges: 0,
8509 present: 0,
8510 quorum: 0,
8511 met_quorum: true,
8512 uncontested: Some("only candidate A produced a change".to_owned()),
8513 });
8514 state.reviews = vec![ReviewRound {
8515 round: 1,
8516 head: "deadbeef".to_owned(),
8517 verified_head: None,
8518 verified_at: None,
8519 reviews: Vec::new(),
8520 e2e: Vec::new(),
8521 fix: None,
8522 blocking: 0,
8523 answered: 0,
8524 expected: 0,
8525 clean: true,
8526 verify_retried: false,
8527 e2e_deferred: false,
8528 e2e_defer_reason: None,
8529 progressed: false,
8530 vote_split: false,
8531 reconsideration: Vec::new(),
8532 verdict: None,
8533 }];
8534
8535 let mut runner = Runner {
8536 state,
8537 roles: ResolvedRoles {
8538 implementers: Vec::new(),
8539 judges: Vec::new(),
8540 reviewers: Vec::new(),
8541 fixer: None,
8542 conductor: conductor(),
8543 implementer_roster: Vec::new(),
8544 },
8545 sem: Arc::new(Semaphore::new(1)),
8546 pause: Pause::new(),
8547 interrupt: Pause::new(),
8548 };
8549
8550 runner.gate().await.expect("gate");
8551 assert!(
8552 runner.state.gate_ran,
8553 "zero configured commands is still a real attempt, not an unrun gate"
8554 );
8555 assert!(runner.state.gate.is_empty());
8556 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
8557 assert_ne!(
8558 runner.state.status,
8559 RunStatus::Blocked,
8560 "a gate with nothing to check must not read as failed"
8561 );
8562
8563 runner.merge().await.expect("merge");
8564 assert_eq!(
8565 runner.state.status,
8566 RunStatus::Ready,
8567 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
8568 );
8569 }
8570
8571 #[tokio::test]
8582 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
8583 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8584 let home = crate::run::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 let cache_dir = tmp.path().join("target");
8593
8594 let mut config = Config::default();
8595 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
8596 config.graph.timeout_verify = Some(2);
8599
8600 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8601 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8602 .expect("no io error acquiring directly")
8603 {
8604 crate::cache::AcquireOutcome::Acquired(g) => g,
8605 crate::cache::AcquireOutcome::Busy(b) => {
8606 panic!("expected the direct acquire to win the lease first: {b:?}")
8607 }
8608 };
8609
8610 let mut state = RunState::new(
8611 repo.clone(),
8612 "main".to_owned(),
8613 "deadbeef".to_owned(),
8614 "task".to_owned(),
8615 config,
8616 );
8617 state.candidates = vec![Candidate {
8618 index: 0,
8619 label: 'A',
8620 agent: "alpha".to_owned(),
8621 branch: "does-not-exist".to_owned(),
8622 worktree: repo.clone(),
8623 summary: String::new(),
8624 stat: String::new(),
8625 files: 0,
8626 commits: 0,
8627 empty: false,
8628 failed: None,
8629 verified_noop: None,
8630 duration_ms: 0,
8631 folded: false,
8632 }];
8633 state.tally = Some(Tally {
8634 first_choice: BTreeMap::from([('A', 1)]),
8635 borda: BTreeMap::new(),
8636 winner: 'A',
8637 rankings: 1,
8638 unanimous_initial: true,
8639 deliberated: false,
8640 changed_votes: 0,
8641 unanimous_final: true,
8642 tie_break: None,
8643 judges: 0,
8644 present: 0,
8645 quorum: 0,
8646 met_quorum: true,
8647 uncontested: Some("only candidate A produced a change".to_owned()),
8648 });
8649 state.reviews = vec![ReviewRound {
8650 round: 1,
8651 head: "deadbeef".to_owned(),
8652 verified_head: None,
8653 verified_at: None,
8654 reviews: Vec::new(),
8655 e2e: Vec::new(),
8656 fix: None,
8657 blocking: 0,
8658 answered: 0,
8659 expected: 0,
8660 clean: true,
8661 verify_retried: false,
8662 e2e_deferred: false,
8663 e2e_defer_reason: None,
8664 progressed: false,
8665 vote_split: false,
8666 reconsideration: Vec::new(),
8667 verdict: None,
8668 }];
8669
8670 let mut runner = Runner {
8671 state,
8672 roles: ResolvedRoles {
8673 implementers: Vec::new(),
8674 judges: Vec::new(),
8675 reviewers: Vec::new(),
8676 fixer: None,
8677 conductor: conductor(),
8678 implementer_roster: Vec::new(),
8679 },
8680 sem: Arc::new(Semaphore::new(1)),
8681 pause: Pause::new(),
8682 interrupt: Pause::new(),
8683 };
8684
8685 let started = std::time::Instant::now();
8686 runner.gate().await.expect("gate");
8687 assert!(
8688 started.elapsed() < Duration::from_secs(1),
8689 "a gate with nothing to run must never wait on a lease it never needed"
8690 );
8691 assert!(
8692 runner.state.gate_ran,
8693 "zero commands is still a real, immediate attempt"
8694 );
8695 assert!(runner.state.gate.is_empty());
8696 assert_ne!(
8697 runner.state.status,
8698 RunStatus::Blocked,
8699 "must not read as resource-blocked on a lease it never asked for"
8700 );
8701 }
8702
8703 #[tokio::test]
8713 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8714 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8715
8716 let tmp = tempfile::tempdir().expect("tempdir");
8717 let repo = tmp.path().join("repo");
8718 std::fs::create_dir_all(&repo).unwrap();
8719 init_repo(&repo);
8720
8721 let mut config = Config::default();
8722 config.verify.gate = vec![
8723 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8724 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8725 .to_owned(),
8726 ];
8727
8728 let mut state = RunState::new(
8729 repo.clone(),
8730 "main".to_owned(),
8731 "deadbeef".to_owned(),
8732 "task".to_owned(),
8733 config,
8734 );
8735 let run_id = state.id.clone();
8736 state.candidates = vec![Candidate {
8737 index: 0,
8738 label: 'A',
8739 agent: "alpha".to_owned(),
8740 branch: "does-not-exist".to_owned(),
8741 worktree: repo.clone(),
8742 summary: String::new(),
8743 stat: String::new(),
8744 files: 0,
8745 commits: 0,
8746 empty: false,
8747 failed: None,
8748 verified_noop: None,
8749 duration_ms: 0,
8750 folded: false,
8751 }];
8752 state.tally = Some(Tally {
8753 first_choice: BTreeMap::from([('A', 1)]),
8754 borda: BTreeMap::new(),
8755 winner: 'A',
8756 rankings: 1,
8757 unanimous_initial: true,
8758 deliberated: false,
8759 changed_votes: 0,
8760 unanimous_final: true,
8761 tie_break: None,
8762 judges: 0,
8763 present: 0,
8764 quorum: 0,
8765 met_quorum: true,
8766 uncontested: Some("only candidate A produced a change".to_owned()),
8767 });
8768 state.reviews = vec![ReviewRound {
8769 round: 1,
8770 head: "deadbeef".to_owned(),
8771 verified_head: None,
8772 verified_at: None,
8773 reviews: Vec::new(),
8774 e2e: Vec::new(),
8775 fix: None,
8776 blocking: 0,
8777 answered: 0,
8778 expected: 0,
8779 clean: true,
8780 verify_retried: false,
8781 e2e_deferred: false,
8782 e2e_defer_reason: None,
8783 progressed: false,
8784 vote_split: false,
8785 reconsideration: Vec::new(),
8786 verdict: None,
8787 }];
8788
8789 let mut runner = Runner {
8790 state,
8791 roles: ResolvedRoles {
8792 implementers: Vec::new(),
8793 judges: Vec::new(),
8794 reviewers: Vec::new(),
8795 fixer: None,
8796 conductor: conductor(),
8797 implementer_roster: Vec::new(),
8798 },
8799 sem: Arc::new(Semaphore::new(1)),
8800 pause: Pause::new(),
8801 interrupt: Pause::new(),
8802 };
8803
8804 let started_marker = repo.join("started.marker");
8805 let release_marker = repo.join("release.marker");
8806 let poller = tokio::spawn(async move {
8807 for _ in 0..100 {
8812 if started_marker.exists()
8813 && let Ok(s) = crate::run::RunState::load(&run_id)
8814 && let Some(a) = s.active.get("gate")
8815 {
8816 std::fs::write(&release_marker, b"go").expect("release marker");
8817 return Some(a.clone());
8818 }
8819 tokio::time::sleep(Duration::from_millis(50)).await;
8820 }
8821 None
8822 });
8823
8824 runner.gate().await.expect("gate");
8825 let captured = poller.await.expect("poller task");
8826 let captured = captured.expect(
8827 "the poller never saw a `gate` task entry in run.json while the command was \
8828 still blocked on its own release marker",
8829 );
8830
8831 assert_eq!(captured.task.as_deref(), Some("gate"));
8832 assert_eq!(captured.node, "gate");
8833 assert_eq!(captured.index, Some(1));
8834 assert_eq!(captured.total, Some(1));
8835 assert!(
8836 captured
8837 .command
8838 .as_deref()
8839 .is_some_and(|c| c.contains("started.marker")),
8840 "{captured:?}"
8841 );
8842
8843 assert!(
8844 runner.state.active.is_empty(),
8845 "the entry must be cleared once the command actually finished: {:?}",
8846 runner.state.active
8847 );
8848 assert!(runner.state.gate_ran);
8849 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8850 }
8851
8852 #[tokio::test]
8865 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8866 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8867 let home = crate::run::home();
8868
8869 let tmp = tempfile::tempdir().expect("tempdir");
8870 let repo = tmp.path().join("repo");
8871 std::fs::create_dir_all(&repo).unwrap();
8872 init_repo(&repo);
8873 let head = crate::git::rev_parse(&repo, "HEAD")
8874 .await
8875 .expect("rev-parse");
8876 let cache_dir = tmp.path().join("target");
8879
8880 let mut config = Config::default();
8881 config.verify.e2e = vec![format!(
8882 "CARGO_TARGET_DIR='{}' test -f README.md",
8883 cache_dir.display()
8884 )];
8885 config.graph.review_rounds = 1;
8886 config.graph.timeout_verify = Some(2);
8889
8890 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8891 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8892 .expect("no io error acquiring directly")
8893 {
8894 crate::cache::AcquireOutcome::Acquired(g) => g,
8895 crate::cache::AcquireOutcome::Busy(b) => {
8896 panic!("expected the direct acquire to win the lease first: {b:?}")
8897 }
8898 };
8899
8900 let mut state = RunState::new(
8901 repo.clone(),
8902 "main".to_owned(),
8903 head.clone(),
8904 "task".to_owned(),
8905 config,
8906 );
8907 state.candidates = vec![Candidate {
8908 index: 0,
8909 label: 'A',
8910 agent: "alpha".to_owned(),
8911 branch: "does-not-exist".to_owned(),
8912 worktree: repo.clone(),
8913 summary: String::new(),
8914 stat: String::new(),
8915 files: 0,
8916 commits: 0,
8917 empty: false,
8918 failed: None,
8919 verified_noop: None,
8920 duration_ms: 0,
8921 folded: false,
8922 }];
8923 state.tally = Some(Tally {
8924 first_choice: BTreeMap::from([('A', 1)]),
8925 borda: BTreeMap::new(),
8926 winner: 'A',
8927 rankings: 1,
8928 unanimous_initial: true,
8929 deliberated: false,
8930 changed_votes: 0,
8931 unanimous_final: true,
8932 tie_break: None,
8933 judges: 0,
8934 present: 0,
8935 quorum: 0,
8936 met_quorum: true,
8937 uncontested: Some("only candidate A produced a change".to_owned()),
8938 });
8939 state.reviews = vec![ReviewRound {
8943 round: 1,
8944 head: head.clone(),
8945 verified_head: None,
8946 verified_at: None,
8947 reviews: Vec::new(),
8948 e2e: Vec::new(),
8949 fix: None,
8950 blocking: 1,
8951 answered: 1,
8952 expected: 1,
8953 clean: false,
8954 verify_retried: false,
8955 e2e_deferred: true,
8956 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8957 progressed: false,
8958 vote_split: false,
8959 reconsideration: Vec::new(),
8960 verdict: None,
8961 }];
8962
8963 let mut runner = Runner {
8964 state,
8965 roles: ResolvedRoles {
8966 implementers: Vec::new(),
8967 judges: Vec::new(),
8968 reviewers: Vec::new(),
8969 fixer: None,
8970 conductor: conductor(),
8971 implementer_roster: Vec::new(),
8972 },
8973 sem: Arc::new(Semaphore::new(1)),
8974 pause: Pause::new(),
8975 interrupt: Pause::new(),
8976 };
8977
8978 let shell = runner.state.config.shell();
8979 runner
8980 .stop_reviewing("round budget spent", &shell, &repo)
8981 .await
8982 .expect("stop_reviewing");
8983
8984 let last = runner.state.reviews.last().expect("round record");
8985 assert_eq!(
8986 last.e2e_status(),
8987 E2eStatus::ResourceBlocked,
8988 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8989 failed: {last:?}"
8990 );
8991 assert_eq!(
8992 last.verified_head.as_deref(),
8993 Some(head.as_str()),
8994 "which commit this attempt targeted is known even though nothing finished checking \
8995 it"
8996 );
8997 let first_attempt_at = last
8998 .verified_at
8999 .expect("when this attempt ran is known too");
9000 assert_ne!(
9001 runner.state.status,
9002 RunStatus::Blocked,
9003 "contention is evidence about the machine, not the patch — it must not settle the \
9004 run as blocked: {:?}",
9005 runner.state.status
9006 );
9007 assert!(
9008 !runner
9009 .state
9010 .events
9011 .iter()
9012 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
9013 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
9014 runner.state.events
9015 );
9016
9017 runner
9022 .stop_reviewing("round budget spent", &shell, &repo)
9023 .await
9024 .expect("stop_reviewing retry");
9025 assert_eq!(
9026 runner.state.reviews.len(),
9027 1,
9028 "no new round was started: {:?}",
9029 runner.state.reviews
9030 );
9031 let last = runner.state.reviews.last().expect("round record");
9032 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
9033 assert!(
9034 last.verified_at.expect("still known") > first_attempt_at,
9035 "a second reentry must be a fresh attempt, not a stale copy of the first"
9036 );
9037 assert_ne!(runner.state.status, RunStatus::Blocked);
9038
9039 held.release();
9040 }
9041
9042 #[tokio::test]
9054 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
9055 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
9056 let home = crate::run::home();
9057
9058 let tmp = tempfile::tempdir().expect("tempdir");
9059 let repo = tmp.path().join("repo");
9060 std::fs::create_dir_all(&repo).unwrap();
9061 init_repo(&repo);
9062 let head = crate::git::rev_parse(&repo, "HEAD")
9063 .await
9064 .expect("rev-parse");
9065 let cache_dir = tmp.path().join("target");
9066
9067 let mut config = Config::default();
9068 config.verify.e2e = vec![format!(
9069 "CARGO_TARGET_DIR='{}' test -f README.md",
9070 cache_dir.display()
9071 )];
9072 config.graph.review_rounds = 1;
9073 config.graph.timeout_verify = Some(2);
9074
9075 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
9076 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
9077 .expect("no io error acquiring directly")
9078 {
9079 crate::cache::AcquireOutcome::Acquired(g) => g,
9080 crate::cache::AcquireOutcome::Busy(b) => {
9081 panic!("expected the direct acquire to win the lease first: {b:?}")
9082 }
9083 };
9084
9085 let mut state = RunState::new(
9086 repo.clone(),
9087 "main".to_owned(),
9088 head.clone(),
9089 "task".to_owned(),
9090 config,
9091 );
9092 state.candidates = vec![Candidate {
9093 index: 0,
9094 label: 'A',
9095 agent: "alpha".to_owned(),
9096 branch: "does-not-exist".to_owned(),
9097 worktree: repo.clone(),
9098 summary: String::new(),
9099 stat: String::new(),
9100 files: 0,
9101 commits: 0,
9102 empty: false,
9103 failed: None,
9104 verified_noop: None,
9105 duration_ms: 0,
9106 folded: false,
9107 }];
9108 state.tally = Some(Tally {
9109 first_choice: BTreeMap::from([('A', 1)]),
9110 borda: BTreeMap::new(),
9111 winner: 'A',
9112 rankings: 1,
9113 unanimous_initial: true,
9114 deliberated: false,
9115 changed_votes: 0,
9116 unanimous_final: true,
9117 tie_break: None,
9118 judges: 0,
9119 present: 0,
9120 quorum: 0,
9121 met_quorum: true,
9122 uncontested: Some("only candidate A produced a change".to_owned()),
9123 });
9124 state.reviews = vec![ReviewRound {
9128 round: 1,
9129 head: head.clone(),
9130 verified_head: Some(head.clone()),
9131 verified_at: Some(jiff::Timestamp::now()),
9132 reviews: Vec::new(),
9133 e2e: vec![CommandOutcome {
9134 command: format!(
9135 "CARGO_TARGET_DIR='{}' test -f README.md",
9136 cache_dir.display()
9137 ),
9138 code: None,
9139 output_tail: "waiting for the shared build cache".to_owned(),
9140 duration_ms: 0,
9141 resource_blocked: true,
9142 }],
9143 fix: None,
9144 blocking: 1,
9145 answered: 1,
9146 expected: 1,
9147 clean: false,
9148 verify_retried: false,
9149 e2e_deferred: false,
9150 e2e_defer_reason: None,
9151 progressed: false,
9152 vote_split: false,
9153 reconsideration: Vec::new(),
9154 verdict: None,
9155 }];
9156
9157 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
9158 let mut runner = Runner {
9159 state,
9160 roles: ResolvedRoles {
9161 implementers: Vec::new(),
9162 judges: Vec::new(),
9163 reviewers: Vec::new(),
9164 fixer: None,
9165 conductor: conductor(),
9166 implementer_roster: Vec::new(),
9167 },
9168 sem: Arc::new(Semaphore::new(1)),
9169 pause: Pause::new(),
9170 interrupt: Pause::new(),
9171 };
9172
9173 runner.review_loop().await.expect("review_loop");
9178
9179 assert_eq!(
9180 runner.state.reviews.len(),
9181 1,
9182 "no new round was started on top of the unresolved one: {:?}",
9183 runner.state.reviews
9184 );
9185 let last = &runner.state.reviews[0];
9186 assert_eq!(
9187 last.e2e_status(),
9188 E2eStatus::ResourceBlocked,
9189 "still contended: {last:?}"
9190 );
9191 assert!(
9192 last.verified_at.expect("still known") > first_attempt_at,
9193 "review_loop must have actually retried the check, not left it exactly as found"
9194 );
9195 assert_ne!(
9196 runner.state.status,
9197 RunStatus::Blocked,
9198 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
9199 runner.state.status
9200 );
9201
9202 held.release();
9203 }
9204
9205 #[tokio::test]
9206 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
9207 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
9208 let tmp = tempfile::tempdir().expect("tempdir");
9209 let repo = tmp.path().join("repo");
9210 std::fs::create_dir_all(&repo).unwrap();
9211 init_repo(&repo);
9212
9213 let mut config = Config::default();
9214 config.merge.mode = MergeMode::Pr;
9215 config.graph.land = true;
9216 config.graph.land_approval = false;
9217
9218 let mut state = RunState::new(
9219 repo.clone(),
9220 "main".to_owned(),
9221 "deadbeef".to_owned(),
9222 "task".to_owned(),
9223 config,
9224 );
9225 state.candidates = vec![Candidate {
9226 index: 0,
9227 label: 'A',
9228 agent: "alpha".to_owned(),
9229 branch: "does-not-exist".to_owned(),
9230 worktree: repo.clone(),
9231 summary: String::new(),
9232 stat: String::new(),
9233 files: 0,
9234 commits: 0,
9235 empty: false,
9236 failed: None,
9237 verified_noop: None,
9238 duration_ms: 0,
9239 folded: false,
9240 }];
9241 state.tally = Some(Tally {
9242 first_choice: BTreeMap::from([('A', 1)]),
9243 borda: BTreeMap::new(),
9244 winner: 'A',
9245 rankings: 1,
9246 unanimous_initial: true,
9247 deliberated: false,
9248 changed_votes: 0,
9249 unanimous_final: true,
9250 tie_break: None,
9251 judges: 0,
9252 present: 0,
9253 quorum: 0,
9254 met_quorum: true,
9255 uncontested: Some("only candidate A produced a change".to_owned()),
9256 });
9257 state.reviews = vec![ReviewRound {
9258 round: 1,
9259 head: "deadbeef".to_owned(),
9260 verified_head: None,
9261 verified_at: None,
9262 reviews: Vec::new(),
9263 e2e: Vec::new(),
9264 fix: None,
9265 blocking: 0,
9266 answered: 0,
9267 expected: 0,
9268 clean: true,
9269 verify_retried: false,
9270 e2e_deferred: false,
9271 e2e_defer_reason: None,
9272 progressed: false,
9273 vote_split: false,
9274 reconsideration: Vec::new(),
9275 verdict: None,
9276 }];
9277 state.gate = vec![CommandOutcome {
9278 command: "test".to_owned(),
9279 code: Some(0),
9280 output_tail: String::new(),
9281 duration_ms: 0,
9282 resource_blocked: false,
9283 }];
9284 state.gate_ran = true;
9285 state.status = RunStatus::Landing;
9289 state.merge = Some(MergeOutcome {
9290 mode: MergeMode::Pr,
9291 ok: true,
9292 detail: "https://example.invalid/x/y/pull/1".to_owned(),
9293 empty: false,
9294 });
9295
9296 ask_test_home();
9300 let store = ask::Questions::open();
9301 let q = ask_open_question(&store, &state.id);
9302
9303 let mut runner = Runner {
9304 state,
9305 roles: ResolvedRoles {
9306 implementers: Vec::new(),
9307 judges: Vec::new(),
9308 reviewers: Vec::new(),
9309 fixer: None,
9310 conductor: conductor(),
9311 implementer_roster: Vec::new(),
9312 },
9313 sem: Arc::new(Semaphore::new(1)),
9314 pause: Pause::new(),
9315 interrupt: Pause::new(),
9316 };
9317
9318 runner.execute().await.expect("execute");
9323
9324 assert_eq!(
9325 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
9326 Some("https://example.invalid/x/y/pull/1"),
9327 "reentry must not push again or open a second pull request over the \
9328 one `land` is already watching"
9329 );
9330 assert_ne!(
9331 runner.state.status,
9332 RunStatus::Landing,
9333 "land could not actually reach the fake pull request, so it must \
9334 have given up rather than left the run silently parked forever"
9335 );
9336 assert_eq!(runner.state.status, RunStatus::Blocked);
9340 assert!(
9341 store.get(&q.id).unwrap().status.open(),
9342 "Blocked is still alive; settle_questions must have been a no-op here"
9343 );
9344 }
9345
9346 fn state_with_round(round: ReviewRound) -> RunState {
9347 let mut s = RunState::new(
9348 PathBuf::from("/repo"),
9349 "main".to_owned(),
9350 "abc1234".to_owned(),
9351 "add retries".to_owned(),
9352 Config::default(),
9353 );
9354 s.reviews = vec![round];
9355 s
9356 }
9357
9358 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
9359 crate::verdict::Finding {
9360 id: id.to_owned(),
9361 severity,
9362 file: None,
9363 line: None,
9364 title: title.to_owned(),
9365 detail: String::new(),
9366 }
9367 }
9368
9369 #[test]
9370 fn pr_body_names_open_findings_and_declined_ones() {
9371 let round = ReviewRound {
9372 round: 2,
9373 head: "deadbee".to_owned(),
9374 verified_head: None,
9375 verified_at: None,
9376 reviews: vec![ReviewRecord {
9377 attempts: 0,
9378 reviewer: 1,
9379 agent: "alpha".to_owned(),
9380 summary: String::new(),
9381 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
9382 vote: None,
9383 failed: None,
9384 duration_ms: 0,
9385 }],
9386 e2e: vec![CommandOutcome {
9387 command: "cargo test".to_owned(),
9388 code: Some(0),
9389 output_tail: String::new(),
9390 duration_ms: 0,
9391 resource_blocked: false,
9392 }],
9393 verify_retried: false,
9394 e2e_deferred: false,
9395 e2e_defer_reason: None,
9396 fix: Some(FixRecord {
9397 agent: "alpha".to_owned(),
9398 addressed: Vec::new(),
9399 rejected: vec![crate::verdict::Rejection {
9400 id: "R1-1-1".to_owned(),
9401 why: "not reachable from any caller".to_owned(),
9402 }],
9403 notes: String::new(),
9404 committed: true,
9405 failed: None,
9406 duration_ms: 0,
9407 continuation: None,
9408 }),
9409 blocking: 0,
9410 answered: 1,
9411 expected: 1,
9412 clean: false,
9413 progressed: true,
9414 vote_split: false,
9415 reconsideration: Vec::new(),
9416 verdict: None,
9417 };
9418 let state = state_with_round(round);
9419 let body = pr_message(&state, 'A').body;
9420
9421 assert!(body.contains("add retries"), "the task must still be there");
9422 assert!(body.contains("R2-1-1"), "{body}");
9423 assert!(body.contains("unused import"), "{body}");
9424 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
9425 assert!(
9426 body.contains("not reachable from any caller"),
9427 "the reason it was declined: {body}"
9428 );
9429 }
9430
9431 #[test]
9432 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
9433 let round = ReviewRound {
9434 round: 1,
9435 head: "deadbee".to_owned(),
9436 verified_head: None,
9437 verified_at: None,
9438 reviews: vec![ReviewRecord {
9439 attempts: 0,
9440 reviewer: 1,
9441 agent: "alpha".to_owned(),
9442 summary: String::new(),
9443 findings: Vec::new(),
9444 vote: None,
9445 failed: None,
9446 duration_ms: 0,
9447 }],
9448 e2e: Vec::new(),
9449 verify_retried: false,
9450 e2e_deferred: false,
9451 e2e_defer_reason: None,
9452 fix: None,
9453 blocking: 0,
9454 answered: 1,
9455 expected: 1,
9456 clean: true,
9457 progressed: false,
9458 vote_split: false,
9459 reconsideration: Vec::new(),
9460 verdict: None,
9461 };
9462 let state = state_with_round(round);
9463 let body = pr_message(&state, 'A').body;
9464 assert!(!body.contains("Open review findings"), "{body}");
9465 assert!(!body.contains("Declined"), "{body}");
9466 }
9467
9468 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
9469 let mut state = RunState::new(
9470 PathBuf::from("/repo"),
9471 "main".to_owned(),
9472 "abc1234".to_owned(),
9473 instruction.to_owned(),
9474 Config::default(),
9475 );
9476 state.candidates.push(Candidate {
9477 index: 0,
9478 label: 'A',
9479 agent: "alpha".to_owned(),
9480 branch: "magi/x/A".to_owned(),
9481 worktree: PathBuf::from("/wt"),
9482 summary: summary.to_owned(),
9483 stat: String::new(),
9484 files: 1,
9485 commits: 1,
9486 empty: false,
9487 failed: None,
9488 verified_noop: None,
9489 folded: false,
9490 duration_ms: 0,
9491 });
9492 state
9493 }
9494
9495 #[test]
9496 fn pr_message_describes_the_change_not_the_task() {
9497 let state = state_with_summary(
9498 "今回やってほしいこと: results projector を直す",
9499 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
9500 );
9501 let m = pr_message(&state, 'A');
9502 assert_eq!(m.title, "fix(web): batch the runs list reads");
9503 assert!(
9504 m.body.starts_with("## Summary\n\n- reads run.json once"),
9505 "{}",
9506 m.body
9507 );
9508 assert!(!m.body.contains("TITLE:"), "{}", m.body);
9509 let task_at = m.body.find("今回やってほしいこと").unwrap();
9510 let details_at = m.body.find("<details>").unwrap();
9511 assert!(
9512 details_at < task_at,
9513 "the task lives inside <details>: {}",
9514 m.body
9515 );
9516 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
9517 assert!(m.body.contains("magi:candidate-a"));
9518 }
9519
9520 #[test]
9521 fn pr_message_falls_back_to_the_task_without_a_title_line() {
9522 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
9523 let m = pr_message(&state, 'A');
9524 assert_eq!(m.title, "add retries");
9525 assert!(
9526 m.body.contains("## Summary\n\n- did some things"),
9527 "{}",
9528 m.body
9529 );
9530
9531 let none = RunState::new(
9532 PathBuf::from("/repo"),
9533 "main".to_owned(),
9534 "abc1234".to_owned(),
9535 "add retries".to_owned(),
9536 Config::default(),
9537 );
9538 let m = pr_message(&none, 'A');
9539 assert_eq!(m.title, "add retries");
9540 assert!(!m.body.contains("## Summary"), "{}", m.body);
9541 }
9542
9543 #[test]
9544 fn pr_message_refuses_the_candidate_commit_subject() {
9545 for bad in [
9546 "TITLE: magi: candidate A (uncommitted work)",
9547 "TITLE: chore: stuff (uncommitted work)",
9548 "TITLE: ",
9549 ] {
9550 let state = state_with_summary("add retries", bad);
9551 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
9552 }
9553 }
9554
9555 #[test]
9556 fn pr_message_bounds_a_very_long_task_and_title() {
9557 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
9558 let state = state_with_summary(&long, "- nothing");
9559 let m = pr_message(&state, 'A');
9560 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9561 assert!(!m.title.contains('\n'));
9562
9563 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
9564 let m = pr_message(&state, 'A');
9565 assert!(m.title.starts_with("feat: "));
9566 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9567 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
9568 }
9569
9570 #[test]
9571 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
9572 let mut state = state_with_summary(
9576 "add retries",
9577 "TITLE: fix(web): batch reads\n- reads run.json once",
9578 );
9579 state.config.graph.language = "ja".to_owned();
9580 let m = pr_message(&state, 'A');
9581 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
9582
9583 let task = "今回やってほしいこと: results projector を直す";
9586 let mut state = state_with_summary(task, "- no title line");
9587 state.config.graph.language = "ja".to_owned();
9588 let m = pr_message(&state, 'A');
9589 assert_eq!(
9590 m.title,
9591 format!("chore: land candidate A of run {}", state.id)
9592 );
9593 assert!(
9594 m.body.contains(&format!(
9595 "<summary>Original task</summary>\n\n{task}\n\n</details>"
9596 )),
9597 "{}",
9598 m.body
9599 );
9600 }
9601
9602 #[test]
9603 fn pr_message_scrubs_home_paths_and_addresses() {
9604 let state = state_with_summary(
9605 "fix it in /Users/someone/src/x",
9606 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9607 );
9608 let m = pr_message(&state, 'A');
9609 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9610 assert!(!m.body.contains(leak), "{}", m.body);
9611 }
9612 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9613 }
9614
9615 #[test]
9616 fn pr_message_survives_a_task_that_closes_details() {
9617 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9618 let m = pr_message(&state, 'A');
9619 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9620 }
9621
9622 #[test]
9623 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9624 let cmd = manual_merge_command(
9625 MergeStyle::Squash,
9626 Path::new("/repo"),
9627 "b",
9628 "fix: \"quoted\" $(x) `y`\n\nbody",
9629 );
9630 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9631 }
9632
9633 #[test]
9634 fn manual_merge_command_matches_the_configured_style() {
9635 let repo = Path::new("/repo");
9636 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9637
9638 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9639 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9640
9641 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9642 assert_eq!(
9643 squash,
9644 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9645 \"Merge magi run 0832 (candidate A)\""
9646 );
9647
9648 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9649 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9650 }
9651
9652 #[test]
9653 fn a_nudge_gets_a_quarter_of_the_budget() {
9654 assert_eq!(retry_budget(secs(1200), true), secs(300));
9656 assert_eq!(retry_budget(secs(3600), true), secs(900));
9657 }
9658
9659 #[test]
9660 fn a_resent_prompt_keeps_the_whole_budget() {
9661 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9664 assert_eq!(retry_budget(secs(60), false), secs(60));
9665 }
9666
9667 #[test]
9668 fn the_floor_never_exceeds_the_original_budget() {
9669 assert_eq!(retry_budget(secs(60), true), secs(60));
9673 assert_eq!(retry_budget(secs(480), true), secs(120));
9674 assert_eq!(retry_budget(secs(0), true), secs(0));
9675 }
9676
9677 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9678 agent::CommandEvidence {
9679 id: "item1".to_owned(),
9680 description: "cargo test".to_owned(),
9681 exit_code,
9682 result_summary: String::new(),
9683 source: "codex".to_owned(),
9684 }
9685 }
9686
9687 #[test]
9688 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9689 assert!(!has_unconfirmed_command(&[]));
9693 }
9694
9695 #[test]
9696 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9697 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9701 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9702 assert!(!has_unconfirmed_command(&[
9703 evidence(Some(0)),
9704 evidence(Some(101))
9705 ]));
9706 }
9707
9708 #[test]
9709 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9710 assert!(has_unconfirmed_command(&[
9711 evidence(Some(0)),
9712 evidence(None)
9713 ]));
9714 }
9715
9716 #[test]
9717 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9718 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9719 assert_eq!(
9720 verified_noop_claim(true, &[], text).as_deref(),
9721 Some("already fixed by b32cfc4, on main.")
9722 );
9723 }
9724
9725 #[test]
9726 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9727 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9730 assert!(verified_noop_claim(false, &[], text).is_none());
9731 }
9732
9733 #[test]
9734 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9735 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9736 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9737 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9739 }
9740
9741 #[test]
9742 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9743 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9744 }
9745
9746 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9749 runner.state.candidates = shape
9750 .iter()
9751 .enumerate()
9752 .map(|(i, &(empty, verified))| Candidate {
9753 index: i,
9754 label: (b'A' + i as u8) as char,
9755 agent: "sonnet".to_owned(),
9756 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9757 worktree: PathBuf::from(format!("/wt/{i}")),
9758 summary: String::new(),
9759 stat: String::new(),
9760 files: 0,
9761 commits: 0,
9762 empty,
9763 failed: None,
9764 verified_noop: verified.map(str::to_owned),
9765 duration_ms: 0,
9766 folded: false,
9767 })
9768 .collect();
9769 }
9770
9771 #[test]
9772 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9773 ask_test_home();
9774 let mut runner = runner_at(RunStatus::Implementing);
9775 set_candidates(
9776 &mut runner,
9777 &[
9778 (true, Some("already on main at b32cfc4")),
9779 (true, Some("same fix, see the existing test")),
9780 ],
9781 );
9782
9783 runner
9784 .after_implement()
9785 .expect("a verified no-op is not an error");
9786
9787 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9788 }
9789
9790 #[test]
9791 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9792 ask_test_home();
9793 let mut runner = runner_at(RunStatus::Implementing);
9794 set_candidates(
9798 &mut runner,
9799 &[(true, Some("already on main at b32cfc4")), (true, None)],
9800 );
9801
9802 let err = runner
9803 .after_implement()
9804 .expect_err("an unverified empty candidate must still fail the run");
9805
9806 assert!(
9807 err.to_string().contains("no candidate produced a change"),
9808 "{err}"
9809 );
9810 assert_eq!(runner.state.status, RunStatus::Failed);
9811 }
9812
9813 #[test]
9814 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9815 ask_test_home();
9816 let mut runner = runner_at(RunStatus::Implementing);
9817 set_candidates(&mut runner, &[(true, None), (true, None)]);
9818
9819 let err = runner
9820 .after_implement()
9821 .expect_err("no candidate declared anything; this is an ordinary failure");
9822
9823 assert!(
9824 err.to_string().contains("no candidate produced a change"),
9825 "{err}"
9826 );
9827 assert_eq!(runner.state.status, RunStatus::Failed);
9828 }
9829}