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(
549 repo: &Path,
550 branch: &str,
551 config: Config,
552 takeover: Option<crate::handover::Takeover>,
553 ) -> Result<Self> {
554 let repo = git::toplevel(repo).await?;
555 let missing = agent::missing_programs(&config.agents);
556 if !missing.is_empty() {
557 bail!(
558 "these agent programs are not on PATH: {}. Fix the roster in \
559 magi.toml or install them.",
560 missing.join(", ")
561 );
562 }
563 let base_branch = match config.merge.base.clone() {
564 Some(b) => b,
565 None => git::current_branch(&repo)
566 .await?
567 .context("HEAD is detached; set [merge] base in magi.toml")?,
568 };
569 if base_branch == branch {
570 bail!("`{branch}` is the base branch; there is nothing to review against");
571 }
572 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
573
574 let roles = config.resolve_roles()?;
575 let max_parallel = config.graph.max_parallel.max(1);
576 let mut state = RunState::new(
577 repo.clone(),
578 base_branch,
579 base_commit.clone(),
580 String::new(),
581 config,
582 );
583
584 let takeover = takeover.unwrap_or_else(|| crate::handover::Takeover {
591 earlier: Vec::new(),
592 home: crate::run::home(),
593 choice: None,
594 });
595 let released = crate::handover::release(&repo, branch, &state.id, &takeover).await?;
596 if let Some(released) = &released {
597 state.event(
598 "release",
599 format!(
600 "took `{branch}` over from run {}: its worktree was released: {}",
601 crate::run::short_of(&released.old_id),
602 released.audit
603 ),
604 );
605 }
606 if let Some(choice) = takeover.choice.as_ref()
610 && let Err(e) =
611 crate::reconcile::apply_choice(&repo, &state.config.merge.remote, branch, choice)
612 .await
613 {
614 if let Some(released) = &released {
615 released.restore(&repo, branch).await;
616 }
617 return Err(e.context("applying the owner's answer about the diverged branch"));
618 }
619 let opened =
620 Self::open_review(&repo, branch, state, roles, max_parallel, base_commit).await;
621 if opened.is_err()
622 && let Some(released) = &released
623 {
624 released.restore(&repo, branch).await;
625 }
626 opened
627 }
628
629 async fn open_review(
632 repo: &Path,
633 branch: &str,
634 mut state: RunState,
635 roles: ResolvedRoles,
636 max_parallel: usize,
637 base_commit: String,
638 ) -> Result<Self> {
639 sync_review_branch(repo, branch, &state.config.merge.remote, &base_commit).await?;
640 let log = git::log_oneline(repo, &base_commit, branch)
643 .await
644 .unwrap_or_default();
645 let instruction = format!(
646 "Review the work already on branch `{branch}`. There is no task \
647 statement: what the change claims to do is whatever its commits \
648 say.\n\n{}",
649 if log.trim().is_empty() {
650 "(no commit messages)"
651 } else {
652 log.trim()
653 }
654 );
655 state.instruction = instruction;
656
657 let worktree = state.worktree_root().join("under-review");
660 if let Some(parent) = worktree.parent() {
661 tokio::fs::create_dir_all(parent).await.ok();
662 }
663 let path = worktree.to_string_lossy().to_string();
664 git::git(repo, &["worktree", "add", &path, branch])
665 .await
666 .with_context(|| {
667 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
668 })?;
669
670 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
671 .await
672 .unwrap_or(0);
673 if commits == 0 {
674 git::worktree_remove(repo, &worktree).await.ok();
675 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
676 }
677 let files = git::changed_files(&worktree, &base_commit, "HEAD")
678 .await
679 .map(|f| f.len())
680 .unwrap_or(0);
681 if files == 0
682 && let (Ok(head_tree), Ok(base_tree)) = (
683 git::tree_of(&worktree, "HEAD").await,
684 git::tree_of(&worktree, &base_commit).await,
685 )
686 && head_tree == base_tree
687 {
688 let head = git::rev_parse(&worktree, "HEAD").await.unwrap_or_default();
689 git::worktree_remove(repo, &worktree).await.ok();
690 bail!(
691 "`{branch}` at {} has a tree identical to base {}; this usually means \
692 the branch ref is stale (check `git rev-parse refs/heads/{branch}` \
693 against `{}/{branch}`) rather than an empty change",
694 short(&head),
695 short(&base_commit),
696 state.config.merge.remote
697 );
698 }
699 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
700 .await
701 .unwrap_or_default();
702
703 state.candidates.push(Candidate {
704 index: 0,
705 label: 'A',
706 agent: "(existing branch)".to_owned(),
709 branch: branch.to_owned(),
710 worktree,
711 summary: String::new(),
712 stat,
713 files,
714 commits,
715 empty: false,
716 failed: None,
717 verified_noop: None,
718 duration_ms: 0,
719 folded: false,
720 });
721 state.tally = Some(Tally {
722 first_choice: BTreeMap::from([('A', 0)]),
723 borda: BTreeMap::new(),
724 winner: 'A',
725 rankings: 0,
726 unanimous_initial: false,
727 deliberated: false,
728 changed_votes: 0,
729 unanimous_final: false,
730 tie_break: None,
731 judges: 0,
735 present: 0,
736 quorum: 0,
737 met_quorum: true,
738 uncontested: Some("review-only run: nothing competed".to_owned()),
739 });
740 state.status = RunStatus::Reviewing;
741 state.event(
742 "start",
743 format!(
744 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
745 state.id
746 ),
747 );
748 state.save()?;
749 Ok(Self {
750 state,
751 roles,
752 sem: Arc::new(Semaphore::new(max_parallel)),
753 pause: Pause::new(),
754 interrupt: Pause::new(),
755 })
756 }
757
758 pub fn resume(id: &str) -> Result<Self> {
760 let state = RunState::load(id)?;
761 if let Some(to) = &state.released_to {
762 bail!(
763 "run {} cannot be resumed: its worktree was released to run {}",
764 state.short(),
765 crate::run::short_of(to)
766 );
767 }
768 let roles = state.config.resolve_roles()?;
769 let max_parallel = state.config.graph.max_parallel.max(1);
770 Ok(Self {
771 state,
772 roles,
773 sem: Arc::new(Semaphore::new(max_parallel)),
774 pause: Pause::new(),
775 interrupt: Pause::new(),
776 })
777 }
778
779 pub async fn execute(&mut self) -> Result<()> {
786 let result = self.execute_graph().await;
787 self.mark_driver_exited();
788 let ended = if result.is_err() {
789 Some(crate::notices::run_stopped(&self.state.id, &self.state))
790 } else {
791 crate::notices::run_ended(&self.state)
792 };
793 if let Some(notice) = ended {
794 crate::notices::raise(notice);
795 }
796 result
797 }
798
799 fn mark_driver_exited(&mut self) {
808 self.state.driver_exited = true;
809 let pid = std::process::id();
810 let Ok(mut disk) = RunState::load(&self.state.id) else {
811 return;
812 };
813 if disk.released_to.is_some() || disk.driver_pid != Some(pid) || disk.driver_exited {
814 return;
815 }
816 disk.driver_exited = true;
817 if let Err(e) = disk.save() {
818 tracing::warn!("could not record that run {} stopped: {e:#}", self.state.id);
819 }
820 }
821
822 async fn execute_graph(&mut self) -> Result<()> {
823 self.state.parked = false;
828 self.state.clear_active();
835 if let Ok(disk) = RunState::load(&self.state.id)
853 && let Some(to) = &disk.released_to
854 {
855 bail!(
856 "run {} cannot continue: its worktree was released to run {}",
857 self.state.short(),
858 crate::run::short_of(to)
859 );
860 }
861 let pid = std::process::id();
862 self.state.driver_pid = Some(pid);
863 self.state.driver_started_at = crate::proc::process_started_at(pid);
864 self.state.driver_exited = false;
865 self.state.save()?;
866 if self.state.status == RunStatus::Stalled {
879 if self.recover_stall().await? {
880 self.finish_after_tally().await?;
881 } else {
882 self.state.save()?;
884 }
885 return Ok(());
886 }
887 if self.state.status == RunStatus::Landing {
897 self.run_land().await?;
898 self.settle_questions();
903 return Ok(());
904 }
905 self.prep().await?;
906 if self.park_here()? {
907 return Ok(());
908 }
909 self.advise().await?;
910 if self.park_here()? {
911 return Ok(());
912 }
913 self.implement().await?;
914 if self.park_here()? {
915 return Ok(());
916 }
917 if self.state.status == RunStatus::VerifiedNoop {
921 return Ok(());
922 }
923 self.judge().await?;
924 if self.park_here()? {
925 return Ok(());
926 }
927 self.deliberate().await?;
928 if self.park_here()? {
929 return Ok(());
930 }
931 self.vote().await?;
932 if self.park_here()? {
933 return Ok(());
934 }
935 self.tally()?;
936 if self.state.status == RunStatus::Stalled {
941 self.state.save()?;
945 return Ok(());
946 }
947 self.finish_after_tally().await?;
948 Ok(())
949 }
950
951 fn park_here(&mut self) -> Result<bool> {
958 if !self.pause.parked() && !self.interrupt.parked() {
963 return Ok(false);
964 }
965 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
966 Some(reason) => format!(
967 "parked after `{}` ({reason}) — resume to carry on from here",
968 self.state.status.as_str()
969 ),
970 None => format!(
971 "parked after `{}` — resume to carry on from here",
972 self.state.status.as_str()
973 ),
974 };
975 self.state.event("park", why);
976 self.state.parked = true;
977 self.state.save()?;
978 Ok(true)
979 }
980
981 pub fn on_pause(&mut self, pause: Pause) {
983 self.pause = pause;
984 }
985
986 pub fn watch_interrupt(&mut self, pause: Pause) {
992 self.interrupt = pause;
993 }
994
995 fn settle_questions(&mut self) {
1015 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
1016 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
1017 }
1018 }
1019
1020 async fn finish_after_tally(&mut self) -> Result<()> {
1023 self.fold_losers().await?;
1024 self.sync_to_base().await?;
1029 self.review_loop().await?;
1030 self.sync_to_base().await?;
1031 self.gate().await?;
1032 self.merge().await?;
1033 self.state.save()?;
1034 Ok(())
1035 }
1036
1037 async fn prep(&mut self) -> Result<()> {
1040 if !self.state.candidates.is_empty() {
1041 return Ok(());
1042 }
1043 self.state.status = RunStatus::Prep;
1044 let repo = self.state.repo.clone();
1045 let base = self.state.base_commit.clone();
1046 let plan = refs::plan(&repo, &self.state.seeds).await?;
1047 let start = plan.start.clone().unwrap_or_else(|| base.clone());
1048 let root = self.state.worktree_root();
1049 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
1050
1051 let hooks_dir = self.state.dir().join("hooks");
1054 if self.state.config.blind.commit_msg_hook {
1055 std::fs::create_dir_all(&hooks_dir)
1056 .with_context(|| format!("create {}", hooks_dir.display()))?;
1057 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
1058 let path = hooks_dir.join("commit-msg");
1059 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
1060 make_executable(&path)?;
1061 git::acquire_worktree_config(&repo).await?;
1069 self.state.enabled_worktree_config = true;
1070 }
1071
1072 for (index, (spec, label)) in self
1073 .roles
1074 .implementers
1075 .clone()
1076 .into_iter()
1077 .zip(labels)
1078 .enumerate()
1079 {
1080 let branch = self.state.branch_for(label);
1081 let worktree = root.join(format!("cand-{label}"));
1082 git::worktree_add_branch(&repo, &worktree, &branch, &start).await?;
1083 if self.state.config.blind.commit_msg_hook {
1084 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
1085 }
1086 git::local_exclude(&worktree, "/.magi/").await?;
1087 for pick in &plan.picks {
1088 if let Err(e) = git::cherry_pick(&worktree, pick).await {
1089 self.state.status = RunStatus::Blocked;
1090 self.state
1091 .event("prep", format!("cannot apply referenced commit: {e}"));
1092 self.state.save()?;
1093 return Err(e);
1094 }
1095 }
1096 self.state.candidates.push(Candidate {
1097 index,
1098 label,
1099 agent: spec.id.clone(),
1100 branch,
1101 worktree,
1102 summary: String::new(),
1103 stat: String::new(),
1104 files: 0,
1105 commits: 0,
1106 empty: false,
1107 failed: None,
1108 verified_noop: None,
1109 duration_ms: 0,
1110 folded: false,
1111 });
1112 }
1113
1114 for j in 1..=self.roles.judges.len() {
1115 let wt = root.join(format!("judge-{j}"));
1116 if !wt.exists() {
1117 git::worktree_add_detached(&repo, &wt, &base).await?;
1118 }
1119 }
1120
1121 if self.state.config.graph.advise {
1129 for k in 1..=self.state.config.graph.advisors {
1130 let wt = root.join(format!("advisor-{k}"));
1131 if !wt.exists() {
1132 git::worktree_add_detached(&repo, &wt, &base).await?;
1133 }
1134 }
1135 }
1136
1137 let authors: Vec<&str> = self
1142 .roles
1143 .implementers
1144 .iter()
1145 .map(|a| a.id.as_str())
1146 .collect();
1147 let overlap: Vec<String> = self
1148 .roles
1149 .judges
1150 .iter()
1151 .enumerate()
1152 .filter(|(_, j)| authors.contains(&j.id.as_str()))
1153 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
1154 .collect();
1155 if !overlap.is_empty() {
1156 let note = format!(
1157 "{} also authored a candidate; blind, but the panel is less \
1158 independent than {} distinct agents would be",
1159 overlap.join(", "),
1160 self.roles.judges.len()
1161 );
1162 self.state.event("prep", note);
1163 }
1164
1165 self.state.event(
1166 "prep",
1167 format!(
1168 "{} candidates, {} judges, base {} ({})",
1169 self.state.candidates.len(),
1170 self.roles.judges.len(),
1171 &self.state.base_commit[..7.min(self.state.base_commit.len())],
1172 self.state.base_branch
1173 ),
1174 );
1175 self.state.status = RunStatus::Implementing;
1176 self.state.save()?;
1177 Ok(())
1178 }
1179
1180 async fn advise(&mut self) -> Result<()> {
1213 let implement_untouched = self
1214 .state
1215 .candidates
1216 .iter()
1217 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
1218 if !self.state.config.graph.advise || self.state.advise_attempted {
1219 return Ok(());
1220 }
1221 if !implement_untouched {
1222 self.state.event(
1223 "advise",
1224 "skipping the design-deliberation stage: at least one \
1225 candidate already shows implementation progress, so this \
1226 run is past the point the stage exists to run before"
1227 .to_owned(),
1228 );
1229 self.state.advise_attempted = true;
1230 self.state.save()?;
1231 return Ok(());
1232 }
1233 let run_id = self.state.id.clone();
1234 let prompts = self.state.config.prompts.clone();
1235 let instruction = self.state.instruction.clone();
1236 let language = self.state.config.graph.language.clone();
1237 let root = self.state.worktree_root();
1238 let n = self.state.config.graph.advisors;
1239 let where_recorded = self.state.dir().join("run.json");
1240
1241 let seats = match self.state.config.advisors() {
1242 Ok(seats) if !seats.is_empty() => seats,
1243 Ok(_) => {
1244 self.state.event(
1245 "advise",
1246 format!(
1247 "[graph] advisors is 0; skipping the design-deliberation \
1248 stage and continuing without a synthesis brief (see {})",
1249 where_recorded.display()
1250 ),
1251 );
1252 self.state.advise_attempted = true;
1253 self.state.save()?;
1254 return Ok(());
1255 }
1256 Err(e) => {
1257 self.state.event(
1258 "advise",
1259 format!(
1260 "could not resolve advisor seats ({e:#}); continuing \
1261 without a design-deliberation brief (see {})",
1262 where_recorded.display()
1263 ),
1264 );
1265 self.state.advise_attempted = true;
1266 self.state.save()?;
1267 return Ok(());
1268 }
1269 };
1270
1271 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1272 let artifacts = agent::artifacts_dir(&self.state.dir());
1273 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1274
1275 let mut jobs = Vec::new();
1276 for (i, spec) in seats.iter().cloned().enumerate() {
1277 let seat_key = format!("advisor-{}", i + 1);
1278 let seat = self.seat(&seat_key, &spec.id);
1279 jobs.push(SeatJob {
1280 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1281 spec,
1282 seat,
1283 cwd: worktrees[i % worktrees.len()].clone(),
1284 timeout,
1285 allow_write: false,
1286 sessions: false,
1287 artifacts: artifacts.clone(),
1288 stem: seat_key,
1289 });
1290 }
1291
1292 self.state.event(
1293 "advise",
1294 format!(
1295 "{} advisor seat(s) sketching a design in parallel",
1296 jobs.len()
1297 ),
1298 );
1299 let mut quota_losses = Vec::new();
1300 let cache = self.state.config.cache_dir();
1301 let ctx = WaveCtx {
1302 run: &run_id,
1303 node: "advise",
1304 prompts: &prompts,
1305 cache: cache.as_deref(),
1306 round: None,
1307 };
1308 let results = ask_json_wave::<Proposal>(
1309 jobs,
1310 Arc::clone(&self.sem),
1311 self.state.config.graph.retries,
1312 &ctx,
1313 &mut quota_losses,
1314 &mut self.state,
1315 &|p: &Proposal| p.validate(),
1316 )
1317 .await;
1318 self.state.quota.extend(quota_losses);
1319
1320 let mut records = Vec::with_capacity(results.len());
1321 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1322 let agent_id = seat.agent.clone();
1323 self.state.seats.insert(seat.key.clone(), seat);
1324 match res {
1325 Ok((proposal, out)) => {
1326 self.state
1327 .event("advise", format!("advisor-{} proposed a design", i + 1));
1328 records.push(advise::AdvisorRecord::proposed(
1329 i + 1,
1330 agent_id,
1331 proposal,
1332 out.duration_ms,
1333 ));
1334 }
1335 Err(e) => {
1336 self.state.event(
1337 "advise",
1338 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1339 );
1340 records.push(advise::AdvisorRecord::failed(
1341 i + 1,
1342 agent_id,
1343 e.to_string(),
1344 ));
1345 }
1346 }
1347 }
1348
1349 let mut advice = advise::Advice {
1350 records,
1351 synthesis: None,
1352 };
1353 if advice.proposals().is_empty() {
1354 self.state.event(
1355 "advise",
1356 "no advisor produced a usable proposal; continuing without a \
1357 synthesis brief"
1358 .to_owned(),
1359 );
1360 } else {
1361 match self
1362 .synthesize_brief(
1363 &advice,
1364 &instruction,
1365 &language,
1366 &worktrees[0],
1367 &artifacts,
1368 &run_id,
1369 &prompts,
1370 cache.as_deref(),
1371 )
1372 .await
1373 {
1374 Ok(Some(text)) => {
1375 self.state.event(
1376 "advise",
1377 "synthesized a design brief for the implementer".to_owned(),
1378 );
1379 advice.synthesis = Some(text);
1380 }
1381 Ok(None) => {
1382 self.state.event(
1383 "advise",
1384 "the synthesis seat produced nothing usable; continuing \
1385 without a design brief"
1386 .to_owned(),
1387 );
1388 }
1389 Err(e) => {
1390 self.state.event(
1391 "advise",
1392 format!("could not synthesize a design brief: {e:#}"),
1393 );
1394 }
1395 }
1396 }
1397 advise::apply_reflection(&mut advice);
1398
1399 self.state.advice = Some(advice);
1400 self.state.advise_attempted = true;
1401 self.state.save()?;
1402 Ok(())
1403 }
1404
1405 #[allow(clippy::too_many_arguments)]
1417 async fn synthesize_brief(
1418 &mut self,
1419 advice: &advise::Advice,
1420 instruction: &str,
1421 language: &str,
1422 cwd: &Path,
1423 artifacts: &Path,
1424 run_id: &str,
1425 prompts: &Prompts,
1426 cache: Option<&Path>,
1427 ) -> Result<Option<String>> {
1428 let want = self.state.config.roles.synthesizer.as_deref();
1429 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1430 let mut seat = self.seat("advise-synthesis", &spec.id);
1431 let proposals = advice.proposals();
1432 let mut prompt = prompt::with_overlay(
1433 prompt::synthesize_brief(instruction, &proposals, language),
1434 prompts.overlay("advise"),
1435 );
1436 if cache.is_some() {
1437 prompt.push('\n');
1442 prompt.push_str(&prompt::build_cache_note("advise", false));
1443 }
1444 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1445 let out = agent::invoke(
1446 &spec,
1447 &mut seat,
1448 &Invocation {
1449 cwd,
1450 prompt: &prompt,
1451 timeout,
1452 allow_write: false,
1453 sessions: false,
1454 artifacts,
1455 stem: "advise-synthesis",
1456 run: run_id,
1457 node: "advise",
1458 cache_dir: None,
1459 attachments: &[],
1460 },
1461 )
1462 .await?;
1463 self.state.seats.insert(seat.key.clone(), seat);
1464 if !out.usable() {
1465 return Ok(None);
1466 }
1467 let text =
1468 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1469 Ok((!text.trim().is_empty()).then_some(text))
1470 }
1471
1472 async fn implement(&mut self) -> Result<()> {
1475 let run_id = self.state.id.clone();
1480 let prompts = self.state.config.prompts.clone();
1481 let todo: Vec<usize> = self
1482 .state
1483 .candidates
1484 .iter()
1485 .enumerate()
1486 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1487 .map(|(i, _)| i)
1488 .collect();
1489 if todo.is_empty() {
1490 return self.after_implement();
1491 }
1492 self.state.status = RunStatus::Implementing;
1493
1494 let language = self.state.config.graph.language.clone();
1495 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1496 let sessions = self.state.config.graph.sessions;
1497 let artifacts = agent::artifacts_dir(&self.state.dir());
1498 let brief = self
1502 .state
1503 .advice
1504 .as_ref()
1505 .and_then(|a| a.synthesis.as_deref())
1506 .map(str::to_owned);
1507 let attachments = self.state.attachments.clone();
1508
1509 let mut jobs = Vec::new();
1510 for &i in &todo {
1511 let (index, label, worktree) = {
1512 let c = &self.state.candidates[i];
1513 (c.index, c.label, c.worktree.clone())
1514 };
1515 let spec = self.roles.implementers[index].clone();
1516 let seat_key = format!("impl-{label}");
1517 let seat = self.seat(&seat_key, &spec.id);
1518 let instruction = seeded_instruction(&self.state);
1519 jobs.push(SeatJob {
1520 spec,
1521 seat,
1522 prompt: prompt::implement(
1523 &instruction,
1524 &worktree.to_string_lossy(),
1525 &language,
1526 brief.as_deref(),
1527 &attachments,
1528 ),
1529 cwd: worktree,
1530 timeout,
1531 allow_write: true,
1532 sessions,
1533 artifacts: artifacts.clone(),
1534 stem: format!("impl-{label}"),
1535 });
1536 }
1537
1538 self.state.event(
1539 "implement",
1540 format!("{} candidates in parallel", jobs.len()),
1541 );
1542 let mut sent = jobs.clone();
1548 let cache = self.state.config.cache_dir();
1549 let ctx = WaveCtx {
1550 run: &run_id,
1551 node: "implement",
1552 prompts: &prompts,
1553 cache: cache.as_deref(),
1554 round: None,
1555 };
1556 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1557 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1558 .await;
1559 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1560 .await;
1561 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1562 .await;
1563
1564 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1565 let seat_key = seat.key.clone();
1566 let agent = seat.agent.clone();
1575 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1576 self.state.seats.insert(seat.key.clone(), seat);
1577 let label = self.state.candidates[i].label;
1578 let worktree = self.state.candidates[i].worktree.clone();
1579 let base = self.state.base_commit.clone();
1580
1581 let (summary, duration, failed, verified_claim) = match out {
1582 AgentOutcome::Ok(o) => {
1583 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1584 let failed = (!o.usable()).then(|| {
1585 if o.timed_out {
1586 "agent timed out".to_owned()
1587 } else {
1588 format!("agent exited with {:?}", o.exit_code)
1589 }
1590 });
1591 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1592 (text, o.duration_ms, failed, verified_claim)
1593 }
1594 AgentOutcome::Dropped(o) => {
1600 let why = o
1601 .dropped
1602 .as_ref()
1603 .map(|d| d.why.as_str())
1604 .unwrap_or("the CLI ended the stream without delivering its answer");
1605 (
1606 String::new(),
1607 o.duration_ms,
1608 Some(format!("the CLI dropped the stream ({why})")),
1609 None,
1610 )
1611 }
1612 AgentOutcome::Quota(o) => {
1613 self.state.quota.push(QuotaLoss {
1614 seat: seat_key,
1615 node: "implement".to_owned(),
1616 at: Timestamp::now(),
1617 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1618 });
1619 (
1620 String::new(),
1621 o.duration_ms,
1622 Some("rate limited (quota); produced no change".to_owned()),
1623 None,
1624 )
1625 }
1626 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1627 };
1628
1629 let rescued = match git::rescue_commit(
1632 &worktree,
1633 &format!("magi: candidate {label} (uncommitted work)"),
1634 )
1635 .await
1636 {
1637 Ok(r) => {
1638 self.state.note_withheld("implement", &r.withheld);
1639 r.committed
1640 }
1641 Err(_) => false,
1642 };
1643 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1644 .await
1645 .unwrap_or(0);
1646 let patch = git::diff(&worktree, &base, "HEAD")
1647 .await
1648 .unwrap_or_default();
1649 let stat = git::diff_stat(&worktree, &base, "HEAD")
1650 .await
1651 .unwrap_or_default();
1652 let files = git::changed_files(&worktree, &base, "HEAD")
1653 .await
1654 .map(|f| f.len())
1655 .unwrap_or(0);
1656 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1657
1658 let c = &mut self.state.candidates[i];
1659 if !exhausted_the_fallback_chain {
1660 c.agent = agent;
1661 }
1662 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1663 c.stat = stat;
1664 c.files = files;
1665 c.commits = commits;
1666 c.duration_ms = duration;
1667 c.empty = commits == 0 || patch.trim().is_empty();
1668 c.failed = match failed {
1671 Some(_) if c.empty => failed,
1672 _ => None,
1673 };
1674 c.verified_noop = if c.empty { verified_claim } else { None };
1679 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1680 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1681 (None, true, Some(_), _) => {
1682 format!("candidate {label}: no change produced (agent-verified no-op)")
1683 }
1684 (None, true, None, _) => format!("candidate {label}: no change produced"),
1685 (None, false, _, true) => {
1686 format!(
1687 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1688 )
1689 }
1690 (None, false, _, false) => {
1691 format!("candidate {label}: {files} files, {commits} commits")
1692 }
1693 };
1694 self.state.event("implement", note);
1695 self.state.save()?;
1696 }
1697
1698 self.after_implement()
1699 }
1700
1701 async fn resume_undelivered(
1729 &mut self,
1730 results: &mut [(usize, SeatState, AgentOutcome)],
1731 sent: &[SeatJob],
1732 prompts: &Prompts,
1733 run_id: &str,
1734 ) {
1735 for (wi, seat, out) in results.iter_mut() {
1736 let Some(dropped) = (match &*out {
1737 AgentOutcome::Dropped(o) => o.dropped.clone(),
1738 _ => None,
1739 }) else {
1740 continue;
1741 };
1742 let Some(job) = sent.get(*wi) else { continue };
1743 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1745 self.state.event(
1746 "implement",
1747 format!(
1748 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1749 work is in the tree",
1750 seat.key, dropped.output_tokens, dropped.why
1751 ),
1752 );
1753 continue;
1754 }
1755 if !has_context(&job.spec, seat, job.sessions) {
1763 self.state.event(
1764 "implement",
1765 format!(
1766 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1767 is no session left to resume",
1768 seat.key, dropped.output_tokens, dropped.why
1769 ),
1770 );
1771 continue;
1772 }
1773 self.state.event(
1774 "implement",
1775 format!(
1776 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1777 conversation",
1778 seat.key, dropped.output_tokens, dropped.why
1779 ),
1780 );
1781 let mut retry = job.clone();
1782 retry.seat = seat.clone();
1783 retry.prompt = prompt::resume_after_drop(&dropped.why);
1784 retry.timeout = retry_budget(job.timeout, true);
1785 retry.stem = format!("{}-resume", job.stem);
1786 let cache = self.state.config.cache_dir();
1787 let ctx = WaveCtx {
1788 run: run_id,
1789 node: "implement",
1790 prompts,
1791 cache: cache.as_deref(),
1792 round: None,
1793 };
1794 let (resumed_seat, resumed) =
1795 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1796 *seat = resumed_seat;
1797 *out = resumed;
1798 }
1799 }
1800
1801 async fn resume_quota_losses(
1863 &mut self,
1864 results: &mut [(usize, SeatState, AgentOutcome)],
1865 sent: &mut [SeatJob],
1866 prompts: &Prompts,
1867 run_id: &str,
1868 ) {
1869 let instruction = seeded_instruction(&self.state);
1870 let language = self.state.config.graph.language.clone();
1871 let brief = self
1872 .state
1873 .advice
1874 .as_ref()
1875 .and_then(|a| a.synthesis.as_deref())
1876 .map(str::to_owned);
1877 let attachments = self.state.attachments.clone();
1878 for (wi, seat, out) in results.iter_mut() {
1879 let Some(job) = sent.get_mut(*wi) else {
1880 continue;
1881 };
1882 let start = self
1887 .roles
1888 .implementer_roster
1889 .iter()
1890 .position(|s| s.id == job.spec.id)
1891 .unwrap_or(0);
1892 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1893 let mut fallback_attempt = 0usize;
1894 while matches!(&*out, AgentOutcome::Quota(_)) {
1895 let Some(next) =
1896 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1897 .cloned()
1898 else {
1899 break;
1900 };
1901 tried.insert(next.id.clone());
1902 fallback_attempt += 1;
1903
1904 if let Ok(r) = git::rescue_commit(
1905 &job.cwd,
1906 &format!(
1907 "magi: candidate {} (uncommitted work before quota fallback)",
1908 seat.key
1909 ),
1910 )
1911 .await
1912 {
1913 self.state.note_withheld("implement", &r.withheld);
1914 }
1915
1916 self.state.event(
1917 "implement",
1918 format!(
1919 "{}: rate limited (quota) on {}; retrying with {}",
1920 seat.key, seat.agent, next.id
1921 ),
1922 );
1923
1924 let new_seat = self.seat(&seat.key, &next.id);
1925 job.spec = next.clone();
1933 let mut retry = job.clone();
1934 retry.seat = new_seat;
1935 retry.prompt = prompt::implement(
1936 &instruction,
1937 &job.cwd.to_string_lossy(),
1938 &language,
1939 brief.as_deref(),
1940 &attachments,
1941 );
1942 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1943 let cache = self.state.config.cache_dir();
1944 let ctx = WaveCtx {
1945 run: run_id,
1946 node: "implement",
1947 prompts,
1948 cache: cache.as_deref(),
1949 round: None,
1950 };
1951 let (fallback_seat, fallback_out) = run_one(
1952 retry,
1953 Arc::clone(&self.sem),
1954 &ctx,
1955 &mut self.state,
1956 fallback_attempt,
1957 )
1958 .await;
1959 *seat = fallback_seat;
1960 *out = fallback_out;
1961 }
1962 }
1963 }
1964
1965 async fn resume_unconfirmed_commands(
1989 &mut self,
1990 results: &mut [(usize, SeatState, AgentOutcome)],
1991 sent: &[SeatJob],
1992 prompts: &Prompts,
1993 run_id: &str,
1994 ) {
1995 for (wi, seat, out) in results.iter_mut() {
1996 let AgentOutcome::Ok(o) = &*out else {
1997 continue;
1998 };
1999 if !has_unconfirmed_command(&o.commands) {
2000 continue;
2001 }
2002 let Some(job) = sent.get(*wi) else { continue };
2003 if !has_context(&job.spec, seat, job.sessions) {
2004 self.state.event(
2005 "implement",
2006 format!(
2007 "{}: the reply named a command whose own CLI never confirmed the exit \
2008 status of, but there is no session left to resume",
2009 seat.key
2010 ),
2011 );
2012 continue;
2013 }
2014 self.state.event(
2015 "implement",
2016 format!(
2017 "{}: the reply named a command whose own CLI never confirmed the exit \
2018 status of; resuming the conversation",
2019 seat.key
2020 ),
2021 );
2022 let mut retry = job.clone();
2023 retry.seat = seat.clone();
2024 retry.prompt = prompt::resume_incomplete(
2025 "a command in your last reply had no confirmed exit status",
2026 );
2027 retry.timeout = retry_budget(job.timeout, true);
2028 retry.stem = format!("{}-confirm", job.stem);
2029 let cache = self.state.config.cache_dir();
2030 let ctx = WaveCtx {
2031 run: run_id,
2032 node: "implement",
2033 prompts,
2034 cache: cache.as_deref(),
2035 round: None,
2036 };
2037 let (resumed_seat, resumed) =
2038 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
2039 *seat = resumed_seat;
2040 *out = resumed;
2041 }
2042 }
2043
2044 async fn continue_fix_report(
2065 &mut self,
2066 mut seat: SeatState,
2067 parse_err: String,
2068 job: &SeatJob,
2069 prompts: &Prompts,
2070 run_id: &str,
2071 round: usize,
2072 ) -> (
2073 SeatState,
2074 Option<FixReport>,
2075 Option<String>,
2076 ContinuationRecord,
2077 ) {
2078 let mut last_err = parse_err;
2079 let mut cumulative_wait_ms = 0u64;
2080 let mut attempts = 0usize;
2081 loop {
2082 if !has_context(&job.spec, &seat, job.sessions) {
2083 self.state.event(
2084 "fix",
2085 format!(
2086 "round {round}: fixer's reply had no adoption report ({last_err}); no \
2087 session left to resume into"
2088 ),
2089 );
2090 let outcome = if attempts == 0 {
2091 ContinuationOutcome::NoSession
2092 } else {
2093 ContinuationOutcome::Exhausted
2094 };
2095 return (
2096 seat,
2097 None,
2098 Some(format!("unparsable fix report: {last_err}")),
2099 ContinuationRecord {
2100 attempts,
2101 cumulative_wait_ms,
2102 outcome,
2103 },
2104 );
2105 }
2106 if attempts >= MAX_FIX_CONTINUATIONS {
2107 self.state.event(
2108 "fix",
2109 format!(
2110 "round {round}: fixer's reply still had no adoption report after \
2111 {attempts} continuation(s) ({last_err}); giving up"
2112 ),
2113 );
2114 return (
2115 seat,
2116 None,
2117 Some(format!(
2118 "unparsable fix report after {attempts} continuation(s): {last_err}"
2119 )),
2120 ContinuationRecord {
2121 attempts,
2122 cumulative_wait_ms,
2123 outcome: ContinuationOutcome::Exhausted,
2124 },
2125 );
2126 }
2127 attempts += 1;
2128 self.state.event(
2129 "fix",
2130 format!(
2131 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
2132 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
2133 ),
2134 );
2135 let mut retry = job.clone();
2136 retry.seat = seat.clone();
2137 retry.prompt = prompt::resume_incomplete(&last_err);
2138 retry.timeout = retry_budget(job.timeout, true);
2139 retry.stem = format!("{}-continue{attempts}", job.stem);
2140 let cache = self.state.config.cache_dir();
2141 let ctx = WaveCtx {
2142 run: run_id,
2143 node: "fix",
2144 prompts,
2145 cache: cache.as_deref(),
2146 round: Some(round),
2147 };
2148 let (resumed_seat, resumed_out) = run_one(
2149 retry,
2150 Arc::clone(&self.sem),
2151 &ctx,
2152 &mut self.state,
2153 attempts,
2154 )
2155 .await;
2156 seat = resumed_seat;
2157 match resumed_out {
2158 AgentOutcome::Ok(o) => {
2159 cumulative_wait_ms += o.duration_ms;
2160 match verdict::extract_json::<FixReport>(&o.text) {
2161 Ok(report) if !has_unconfirmed_command(&o.commands) => {
2162 self.state.event(
2163 "fix",
2164 format!(
2165 "round {round}: fixer's adoption report recovered after \
2166 {attempts} continuation(s)"
2167 ),
2168 );
2169 return (
2170 seat,
2171 Some(report),
2172 None,
2173 ContinuationRecord {
2174 attempts,
2175 cumulative_wait_ms,
2176 outcome: ContinuationOutcome::Resumed,
2177 },
2178 );
2179 }
2180 Ok(_) => {
2188 last_err = "the reply parsed, but it reported a command whose own CLI \
2189 never confirmed an exit status"
2190 .to_owned();
2191 }
2192 Err(e) => last_err = e.to_string(),
2193 }
2194 }
2195 AgentOutcome::Quota(o) => {
2196 cumulative_wait_ms += o.duration_ms;
2197 self.state.quota.push(QuotaLoss {
2198 seat: seat.key.clone(),
2199 node: "fix".to_owned(),
2200 at: Timestamp::now(),
2201 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2202 });
2203 self.state.event(
2204 "fix",
2205 format!(
2206 "round {round}: continuation rate limited (quota); not retrying now"
2207 ),
2208 );
2209 return (
2210 seat,
2211 None,
2212 Some("rate limited (quota) while recovering the fix report".to_owned()),
2213 ContinuationRecord {
2214 attempts,
2215 cumulative_wait_ms,
2216 outcome: ContinuationOutcome::QuotaLost,
2217 },
2218 );
2219 }
2220 AgentOutcome::Dropped(o) => {
2221 cumulative_wait_ms += o.duration_ms;
2222 let why = o
2223 .dropped
2224 .as_ref()
2225 .map(|d| d.why.as_str())
2226 .unwrap_or("the CLI ended the stream without delivering its answer");
2227 last_err = format!("the CLI dropped the stream ({why})");
2228 }
2229 AgentOutcome::Failed(e) => last_err = e,
2230 }
2231 }
2232 }
2233
2234 fn after_implement(&mut self) -> Result<()> {
2235 if self.state.leaks.is_empty() {
2237 let cfg = self.state.config.blind.clone();
2238 let mut leaks = Vec::new();
2239 for c in &self.state.candidates {
2240 let Some(patch) =
2241 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
2242 else {
2243 continue;
2244 };
2245 leaks.extend(blind::scan(
2246 &format!("candidate {} patch", c.label),
2247 &patch,
2248 &cfg.vendor_tokens,
2249 ));
2250 }
2251 if !leaks.is_empty() {
2252 let summary = leaks
2253 .iter()
2254 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2255 .collect::<Vec<_>>()
2256 .join(", ");
2257 match cfg.on_leak {
2258 LeakPolicy::Fail => {
2259 self.state.status = RunStatus::Failed;
2260 self.state
2261 .event("blind", format!("vendor text in a patch: {summary}"));
2262 self.state.leaks = leaks;
2263 self.state.save()?;
2264 self.settle_questions();
2265 bail!(
2266 "blind.on_leak = \"fail\" and vendor text reached a \
2267 judged patch: {summary}"
2268 );
2269 }
2270 LeakPolicy::Redact => self.state.event(
2271 "blind",
2272 format!("redacting vendor text for judging: {summary}"),
2273 ),
2274 LeakPolicy::Warn => self.state.event(
2275 "blind",
2276 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2277 ),
2278 }
2279 self.state.leaks = leaks;
2280 }
2281 }
2282
2283 if self.state.viable().is_empty() {
2284 if self.state.all_candidates_verified_noop() {
2285 self.state.status = RunStatus::VerifiedNoop;
2296 self.state.save()?;
2297 self.settle_questions();
2298 return Ok(());
2299 }
2300 self.state.status = RunStatus::Failed;
2301 self.state.save()?;
2302 self.settle_questions();
2303 bail!("no candidate produced a change; nothing to judge");
2304 }
2305 self.state.status = RunStatus::Judging;
2306 self.state.save()?;
2307 Ok(())
2308 }
2309
2310 async fn judge(&mut self) -> Result<()> {
2313 let run_id = self.state.id.clone();
2318 let prompts = self.state.config.prompts.clone();
2319 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2320 return Ok(());
2321 }
2322 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2323 if viable.len() == 1 {
2324 self.state.judge_skipped = true;
2331 self.state.event(
2332 "judge",
2333 format!(
2334 "only candidate {} produced a change; judging skipped",
2335 viable[0].label
2336 ),
2337 );
2338 self.state.save()?;
2339 return Ok(());
2340 }
2341 self.state.status = RunStatus::Judging;
2342
2343 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2344 let language = self.state.config.graph.language.clone();
2345 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2346 let sessions = self.state.config.graph.sessions;
2347 let artifacts = agent::artifacts_dir(&self.state.dir());
2348 let root = self.state.worktree_root();
2349 let base_short = short(&self.state.base_commit);
2350
2351 let mut jobs = Vec::new();
2352 let mut orders = Vec::new();
2353 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2354 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2355 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2356 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2357 let seat_key = format!("judge-{}", j + 1);
2358 let seat = self.seat(&seat_key, &spec.id);
2359 jobs.push(SeatJob {
2360 prompt: prompt::judge(
2361 &self.state.instruction,
2362 &views,
2363 self.roles.judges.len(),
2364 &base_short,
2365 &language,
2366 ),
2367 spec,
2368 seat,
2369 cwd: root.join(format!("judge-{}", j + 1)),
2370 timeout,
2371 allow_write: false,
2372 sessions,
2373 artifacts: artifacts.clone(),
2374 stem: format!("judge-{}", j + 1),
2375 });
2376 }
2377
2378 self.state.event(
2379 "judge",
2380 format!(
2381 "{} judges ranking {} candidates blind",
2382 jobs.len(),
2383 viable.len()
2384 ),
2385 );
2386 let labels_for_check = labels.clone();
2387 let mut quota_losses = Vec::new();
2388 let cache = self.state.config.cache_dir();
2389 let ctx = WaveCtx {
2390 run: &run_id,
2391 node: "judge",
2392 prompts: &prompts,
2393 cache: cache.as_deref(),
2394 round: None,
2395 };
2396 let results = ask_json_wave::<Ranking>(
2397 jobs,
2398 Arc::clone(&self.sem),
2399 self.state.config.graph.retries,
2400 &ctx,
2401 &mut quota_losses,
2402 &mut self.state,
2403 &move |r: &Ranking| r.validate(&labels_for_check),
2404 )
2405 .await;
2406 self.state.quota.extend(quota_losses);
2407
2408 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2409 let agent_id = seat.agent.clone();
2410 self.state.seats.insert(seat.key.clone(), seat);
2411 let mut record = Judgement {
2412 judge: j + 1,
2413 seat: format!("judge-{}", j + 1),
2414 agent: agent_id,
2415 ranking: Vec::new(),
2416 reasons: BTreeMap::new(),
2417 confidence: None,
2418 order: orders[j].clone(),
2419 failed: None,
2420 duration_ms: 0,
2421 };
2422 match res {
2423 Ok((ranking, out)) => {
2424 record.ranking = ranking.normalized();
2425 record.reasons = ranking.reasons;
2426 record.confidence = ranking.confidence;
2427 record.duration_ms = out.duration_ms;
2428 self.state.event(
2429 "judge",
2430 format!(
2431 "judge {} ranked {}",
2432 j + 1,
2433 record.ranking.iter().collect::<String>()
2434 ),
2435 );
2436 }
2437 Err(e) => {
2438 record.failed = Some(e.to_string());
2439 self.state
2440 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2441 }
2442 }
2443 self.state.judgements.push(record);
2444 self.state.save()?;
2445 }
2446 Ok(())
2447 }
2448
2449 async fn deliberate(&mut self) -> Result<()> {
2452 let run_id = self.state.id.clone();
2457 let prompts = self.state.config.prompts.clone();
2458 if !self.state.deliberation.is_empty() {
2459 return Ok(());
2460 }
2461 let tops: Vec<char> = self
2462 .state
2463 .judgements
2464 .iter()
2465 .filter_map(|j| j.ranking.first().copied())
2466 .collect();
2467 let rounds = self.state.config.graph.deliberate_rounds;
2468 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2469 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2470 self.state.event(
2471 "deliberate",
2472 format!("judges agreed on {} outright; no deliberation", tops[0]),
2473 );
2474 }
2475 self.state.status = RunStatus::Voting;
2476 self.state.save()?;
2477 return Ok(());
2478 }
2479
2480 self.state.status = RunStatus::Deliberating;
2481 self.state.event(
2482 "deliberate",
2483 format!(
2484 "split: first choices were {} — opening {rounds} round(s)",
2485 tops.iter().collect::<String>()
2486 ),
2487 );
2488
2489 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2490 let language = self.state.config.graph.language.clone();
2491 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2492 let sessions = self.state.config.graph.sessions;
2493 let artifacts = agent::artifacts_dir(&self.state.dir());
2494 let root = self.state.worktree_root();
2495 let base_short = short(&self.state.base_commit);
2496
2497 for round in 1..=rounds {
2501 let mut turns: Vec<DeliberationTurn> = Vec::new();
2502 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2503 if self.state.judgements[j].failed.is_some() {
2504 continue;
2505 }
2506 let seat_key = format!("judge-{}", j + 1);
2507 let mut seat = self.seat(&seat_key, &spec.id);
2508 let transcript = self.transcript(&turns, j);
2509 let context = if has_context(&spec, &seat, sessions) {
2510 None
2511 } else {
2512 Some(self.candidate_block(&viable, &base_short))
2513 };
2514 let text = prompt::deliberate(
2515 &self.state.instruction,
2516 context.as_deref(),
2517 &transcript,
2518 round,
2519 rounds,
2520 &language,
2521 );
2522 let job = SeatJob {
2523 spec,
2524 seat: seat.clone(),
2525 prompt: text,
2526 cwd: root.join(format!("judge-{}", j + 1)),
2527 timeout,
2528 allow_write: false,
2529 sessions,
2530 artifacts: artifacts.clone(),
2531 stem: format!("delib-{round}-judge-{}", j + 1),
2532 };
2533 let cache = self.state.config.cache_dir();
2534 let ctx = WaveCtx {
2535 run: &run_id,
2536 node: "deliberate",
2537 prompts: &prompts,
2538 cache: cache.as_deref(),
2539 round: None,
2540 };
2541 let (updated, out) =
2542 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2543 seat = updated;
2544 let agent_id = seat.agent.clone();
2545 let seat_key = seat.key.clone();
2546 self.state.seats.insert(seat.key.clone(), seat);
2547 let body = match out {
2548 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2549 AgentOutcome::Dropped(o) => {
2553 let why =
2554 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2555 "the CLI ended the stream without delivering its answer",
2556 );
2557 self.state.event(
2558 "deliberate",
2559 format!(
2560 "judge {} skipped: the CLI dropped the stream ({why})",
2561 j + 1
2562 ),
2563 );
2564 continue;
2565 }
2566 AgentOutcome::Quota(o) => {
2567 self.state.quota.push(QuotaLoss {
2568 seat: seat_key,
2569 node: "deliberate".to_owned(),
2570 at: Timestamp::now(),
2571 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2572 });
2573 self.state.event(
2574 "deliberate",
2575 format!("judge {} skipped: rate limited (quota)", j + 1),
2576 );
2577 continue;
2578 }
2579 AgentOutcome::Failed(e) => {
2580 self.state
2581 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2582 continue;
2583 }
2584 };
2585 let tentative = verdict::extract_json::<Position>(&body)
2586 .ok()
2587 .and_then(|p| p.tentative)
2588 .and_then(|s| s.trim().chars().next())
2589 .map(|c| c.to_ascii_uppercase());
2590 self.state.event(
2591 "deliberate",
2592 format!(
2593 "round {round}: judge {} now favours {}",
2594 j + 1,
2595 tentative.map_or("—".to_owned(), |c| c.to_string())
2596 ),
2597 );
2598 turns.push(DeliberationTurn {
2599 judge: j + 1,
2600 agent: agent_id,
2601 body: blind::sanitize_prose(&body, &self.state.config.blind),
2602 tentative,
2603 });
2604 }
2605 self.state
2606 .deliberation
2607 .push(DeliberationRound { round, turns });
2608 self.state.save()?;
2609 }
2610
2611 self.state.status = RunStatus::Voting;
2612 self.state.save()?;
2613 Ok(())
2614 }
2615
2616 async fn vote(&mut self) -> Result<()> {
2619 let run_id = self.state.id.clone();
2624 let prompts = self.state.config.prompts.clone();
2625 if !self.state.votes.is_empty() {
2626 return Ok(());
2627 }
2628 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2629 if viable.len() == 1 {
2630 return Ok(());
2631 }
2632 self.state.status = RunStatus::Voting;
2633
2634 let language = self.state.config.graph.language.clone();
2635 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2636 let sessions = self.state.config.graph.sessions;
2637 let artifacts = agent::artifacts_dir(&self.state.dir());
2638 let root = self.state.worktree_root();
2639 let base_short = short(&self.state.base_commit);
2640 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2641
2642 let mut jobs = Vec::new();
2643 let mut seats_at = Vec::new();
2644 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2645 if self
2646 .state
2647 .judgements
2648 .get(j)
2649 .is_some_and(|r| r.failed.is_some())
2650 {
2651 continue;
2652 }
2653 let seat_key = format!("judge-{}", j + 1);
2654 let seat = self.seat(&seat_key, &spec.id);
2655 let mut text = prompt::final_vote(&viable, &language);
2656 if !has_context(&spec, &seat, sessions) {
2657 text = format!(
2658 "{}\n\n# Candidates\n\n{}",
2659 text,
2660 self.candidate_block(&candidates, &base_short)
2661 );
2662 }
2663 jobs.push(SeatJob {
2664 spec,
2665 seat,
2666 prompt: text,
2667 cwd: root.join(format!("judge-{}", j + 1)),
2668 timeout,
2669 allow_write: false,
2670 sessions,
2671 artifacts: artifacts.clone(),
2672 stem: format!("vote-judge-{}", j + 1),
2673 });
2674 seats_at.push(j);
2675 }
2676
2677 self.state.event(
2678 "vote",
2679 format!(
2680 "collecting {} final votes one by one, privately",
2681 jobs.len()
2682 ),
2683 );
2684 let allowed = viable.clone();
2685 let mut quota_losses = Vec::new();
2686 let cache = self.state.config.cache_dir();
2687 let ctx = WaveCtx {
2688 run: &run_id,
2689 node: "vote",
2690 prompts: &prompts,
2691 cache: cache.as_deref(),
2692 round: None,
2693 };
2694 let results = ask_json_wave::<FinalVote>(
2695 jobs,
2696 Arc::clone(&self.sem),
2697 self.state.config.graph.retries,
2698 &ctx,
2699 &mut quota_losses,
2700 &mut self.state,
2701 &move |v: &FinalVote| match v.label() {
2702 Some(c) if allowed.contains(&c) => Ok(()),
2703 other => bail!("vote {other:?} is not one of {allowed:?}"),
2704 },
2705 )
2706 .await;
2707 self.state.quota.extend(quota_losses);
2708
2709 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2710 let agent_id = seat.agent.clone();
2711 self.state.seats.insert(seat.key.clone(), seat);
2712 let initial = self
2713 .state
2714 .judgements
2715 .get(j)
2716 .and_then(|r| r.ranking.first().copied());
2717 let mut record = VoteRecord {
2718 judge: j + 1,
2719 agent: agent_id,
2720 vote: None,
2721 reason: String::new(),
2722 changed: false,
2723 };
2724 match res {
2725 Ok((v, _)) => {
2726 record.vote = v.label();
2727 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2728 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2729 self.state.event(
2730 "vote",
2731 format!(
2732 "judge {} voted {}{}",
2733 j + 1,
2734 record.vote.unwrap_or('?'),
2735 if record.changed { " (changed)" } else { "" }
2736 ),
2737 );
2738 }
2739 Err(e) => {
2740 self.state
2741 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2742 }
2743 }
2744 self.state.votes.push(record);
2745 self.state.save()?;
2746 }
2747 Ok(())
2748 }
2749
2750 fn tally(&mut self) -> Result<()> {
2753 if self.state.tally.is_some() {
2754 return Ok(());
2755 }
2756 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2757 let tops: Vec<char> = self
2758 .state
2759 .judgements
2760 .iter()
2761 .filter_map(|j| j.ranking.first().copied())
2762 .collect();
2763 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2764
2765 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2768 let mut cast: Vec<char> = Vec::new();
2769 for (i, j) in self.state.judgements.iter().enumerate() {
2770 let vote = self
2771 .state
2772 .votes
2773 .iter()
2774 .find(|v| v.judge == i + 1)
2775 .and_then(|v| v.vote)
2776 .or_else(|| j.ranking.first().copied());
2777 if let Some(v) = vote {
2778 *first_choice.entry(v).or_insert(0) += 1;
2779 cast.push(v);
2780 }
2781 }
2782
2783 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2784 for j in &self.state.judgements {
2785 let n = j.ranking.len();
2786 for (pos, label) in j.ranking.iter().enumerate() {
2787 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2788 }
2789 }
2790
2791 let best = first_choice.values().copied().max().unwrap_or(0);
2792 let mut leaders: Vec<char> = first_choice
2793 .iter()
2794 .filter(|(_, v)| **v == best)
2795 .map(|(k, _)| *k)
2796 .collect();
2797 let mut tie_break = None;
2798 if leaders.len() > 1 {
2799 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2800 let borda_leaders: Vec<char> = leaders
2801 .iter()
2802 .copied()
2803 .filter(|l| borda[l] == top_borda)
2804 .collect();
2805 tie_break = Some(if borda_leaders.len() == 1 {
2806 format!(
2807 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2808 leaders.len()
2809 )
2810 } else {
2811 format!(
2812 "{} way tie on both first-choice votes and Borda points, broken by label order",
2813 leaders.len()
2814 )
2815 });
2816 leaders = borda_leaders;
2817 leaders.sort_unstable();
2818 }
2819 let winner = *leaders
2820 .first()
2821 .or(viable.first())
2822 .context("no candidate to declare a winner from")?;
2823
2824 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2825 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2826 let deliberated = !self.state.deliberation.is_empty();
2827
2828 let quota_seats: std::collections::BTreeSet<&str> =
2832 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2833 let mut present = 0usize;
2834 for (i, j) in self.state.judgements.iter().enumerate() {
2835 if quota_seats.contains(j.seat.as_str()) {
2836 continue;
2837 }
2838 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2839 let voted = self
2840 .state
2841 .votes
2842 .iter()
2843 .any(|v| v.judge == i + 1 && v.vote.is_some());
2844 if ranked || voted {
2845 present += 1;
2846 }
2847 }
2848 let needs_quorum = viable.len() > 1;
2854 let judges_total = if needs_quorum {
2855 self.roles.judges.len()
2856 } else {
2857 0
2858 };
2859 let quorum = if needs_quorum {
2860 judges_total / 2 + 1
2861 } else {
2862 0
2863 };
2864 let met_quorum = !needs_quorum || present >= quorum;
2865 let uncontested = (!needs_quorum).then(|| {
2866 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2867 });
2868
2869 self.state.event(
2870 "tally",
2871 match &uncontested {
2872 Some(reason) => format!("winner {winner} — {reason}"),
2873 None => format!(
2874 "winner {winner} — votes {} | initial {} | {} changed | \
2875 {present}/{judges_total} judges{}",
2876 first_choice
2877 .iter()
2878 .map(|(k, v)| format!("{k}:{v}"))
2879 .collect::<Vec<_>>()
2880 .join(" "),
2881 if unanimous_initial {
2882 "unanimous"
2883 } else {
2884 "split"
2885 },
2886 changed_votes,
2887 if met_quorum {
2888 String::new()
2889 } else {
2890 format!(" — below quorum ({quorum} required)")
2891 },
2892 ),
2893 },
2894 );
2895 if !met_quorum {
2896 self.state.event(
2897 "stall",
2898 format!(
2899 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2900 the run stops here, resumable"
2901 ),
2902 );
2903 }
2904 self.state.tally = Some(Tally {
2905 first_choice,
2906 borda,
2907 winner,
2908 rankings: tops.len(),
2909 unanimous_initial,
2910 deliberated,
2911 changed_votes,
2912 unanimous_final,
2913 tie_break,
2914 judges: judges_total,
2915 present,
2916 quorum,
2917 met_quorum,
2918 uncontested,
2919 });
2920 self.state.status = if met_quorum {
2921 RunStatus::Reviewing
2922 } else {
2923 RunStatus::Stalled
2924 };
2925 self.state.save()?;
2926 Ok(())
2927 }
2928
2929 #[allow(clippy::too_many_lines)]
2950 async fn recover_stall(&mut self) -> Result<bool> {
2951 let run_id = self.state.id.clone();
2956 let prompts = self.state.config.prompts.clone();
2957 let quota_seats: BTreeSet<&str> =
2962 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2963 let absent: Vec<String> = self
2964 .state
2965 .judgements
2966 .iter()
2967 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2968 .map(|j| j.seat.clone())
2969 .collect();
2970 if absent.is_empty() {
2971 return Ok(false);
2972 }
2973 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2974 if viable.len() <= 1 {
2975 return Ok(false);
2976 }
2977 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2978 let language = self.state.config.graph.language.clone();
2979 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2980 let sessions = self.state.config.graph.sessions;
2981 let artifacts = agent::artifacts_dir(&self.state.dir());
2982 let root = self.state.worktree_root();
2983 let base_short = short(&self.state.base_commit);
2984 let candidates: Vec<Candidate> = viable.clone();
2985
2986 let mut positions: Vec<usize> = absent
2988 .iter()
2989 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2990 .collect();
2991 if positions.is_empty() {
2992 return Ok(false);
2993 }
2994 positions.sort_unstable();
2995 positions.dedup();
2996
2997 let mut judge_jobs = Vec::new();
2999 for &j in &positions {
3000 let order = blind::presentation_order(viable.len(), j, self.state.seed);
3001 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
3002 let seat_key = format!("judge-{}", j + 1);
3003 let spec = self.roles.judges[j].clone();
3004 let seat = self.seat(&seat_key, &spec.id);
3005 judge_jobs.push(SeatJob {
3006 spec,
3007 seat,
3008 prompt: prompt::judge(
3009 &self.state.instruction,
3010 &views,
3011 self.roles.judges.len(),
3012 &base_short,
3013 &language,
3014 ),
3015 cwd: root.join(seat_key),
3016 timeout,
3017 allow_write: false,
3018 sessions,
3019 artifacts: artifacts.clone(),
3020 stem: format!("judge-{}-recover", j + 1),
3021 });
3022 }
3023
3024 let labels_for_check = labels.clone();
3025 let mut judge_losses = Vec::new();
3026 let retries = self.state.config.graph.retries;
3027 let cache = self.state.config.cache_dir();
3028 let ctx = WaveCtx {
3029 run: &run_id,
3030 node: "judge",
3031 prompts: &prompts,
3032 cache: cache.as_deref(),
3033 round: None,
3034 };
3035 let results = ask_json_wave::<Ranking>(
3036 judge_jobs,
3037 Arc::clone(&self.sem),
3038 retries,
3039 &ctx,
3040 &mut judge_losses,
3041 &mut self.state,
3042 &move |r: &Ranking| r.validate(&labels_for_check),
3043 )
3044 .await;
3045
3046 let mut recovered: BTreeSet<usize> = BTreeSet::new();
3048 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
3049 self.state.seats.insert(seat.key.clone(), seat);
3050 let record = &mut self.state.judgements[j];
3051 match res {
3052 Ok((ranking, out)) => {
3053 record.ranking = ranking.normalized();
3054 record.reasons = ranking.reasons;
3055 record.confidence = ranking.confidence;
3056 record.failed = None;
3057 record.duration_ms = out.duration_ms;
3058 recovered.insert(j);
3059 self.state.event(
3060 "recover",
3061 format!("judge {} ranked again after the limit", j + 1),
3062 );
3063 }
3064 Err(e) => {
3065 self.state
3066 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
3067 }
3068 }
3069 }
3070
3071 let mut vote_jobs = Vec::new();
3073 let mut vote_pos: Vec<usize> = Vec::new();
3074 for &j in &recovered {
3075 let seat_key = format!("judge-{}", j + 1);
3076 let spec = self.roles.judges[j].clone();
3077 let seat = self.seat(&seat_key, &spec.id);
3078 let mut text = prompt::final_vote(&labels, &language);
3079 if !has_context(&spec, &seat, sessions) {
3080 text = format!(
3081 "{}\n\n# Candidates\n\n{}",
3082 text,
3083 self.candidate_block(&candidates, &base_short)
3084 );
3085 }
3086 vote_jobs.push(SeatJob {
3087 spec,
3088 seat,
3089 prompt: text,
3090 cwd: root.join(seat_key),
3091 timeout,
3092 allow_write: false,
3093 sessions,
3094 artifacts: artifacts.clone(),
3095 stem: format!("vote-judge-{}-recover", j + 1),
3096 });
3097 vote_pos.push(j);
3098 }
3099 let allowed = labels.clone();
3100 let mut vote_losses = Vec::new();
3101 let vote_retries = self.state.config.graph.retries;
3102 let vote_cache = self.state.config.cache_dir();
3103 let ctx = WaveCtx {
3104 run: &run_id,
3105 node: "vote",
3106 prompts: &prompts,
3107 cache: vote_cache.as_deref(),
3108 round: None,
3109 };
3110 let votes = ask_json_wave::<FinalVote>(
3111 vote_jobs,
3112 Arc::clone(&self.sem),
3113 vote_retries,
3114 &ctx,
3115 &mut vote_losses,
3116 &mut self.state,
3117 &move |v: &FinalVote| match v.label() {
3118 Some(c) if allowed.contains(&c) => Ok(()),
3119 other => bail!("vote {other:?} is not one of {allowed:?}"),
3120 },
3121 )
3122 .await;
3123 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
3124 let agent_id = seat.agent.clone();
3125 self.state.seats.insert(seat.key.clone(), seat);
3126 match res {
3127 Ok((v, _)) => {
3128 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
3129 rec.vote = v.label();
3130 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
3131 } else {
3132 self.state.votes.push(VoteRecord {
3133 judge: j + 1,
3134 agent: agent_id,
3135 vote: v.label(),
3136 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
3137 changed: false,
3138 });
3139 }
3140 self.state.event(
3141 "recover",
3142 format!("judge {} voted again after the limit", j + 1),
3143 );
3144 }
3145 Err(e) => {
3146 self.state
3147 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
3148 }
3149 }
3150 }
3151
3152 let recovered_keys: BTreeSet<String> = recovered
3156 .iter()
3157 .map(|&j| format!("judge-{}", j + 1))
3158 .collect();
3159 self.state
3160 .quota
3161 .retain(|q| !recovered_keys.contains(&q.seat));
3162 for loss in judge_losses.into_iter().chain(vote_losses) {
3166 if recovered_keys.contains(&loss.seat) {
3167 continue;
3168 }
3169 self.state.quota.retain(|q| q.seat != loss.seat);
3170 self.state.quota.push(loss);
3171 }
3172
3173 self.state.tally = None;
3175 self.tally()?;
3176 Ok(self
3177 .state
3178 .tally
3179 .as_ref()
3180 .map(|t| t.met_quorum)
3181 .unwrap_or(false))
3182 }
3183
3184 async fn fold_losers(&mut self) -> Result<()> {
3187 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
3188 return Ok(());
3189 };
3190 let repo = self.state.repo.clone();
3191 let mut folded = Vec::new();
3192 for i in 0..self.state.candidates.len() {
3193 let c = &self.state.candidates[i];
3194 if c.label == winner || c.folded {
3195 continue;
3196 }
3197 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
3198 git::worktree_remove(&repo, &wt).await.ok();
3199 git::branch_delete(&repo, &branch).await.ok();
3200 self.state.candidates[i].folded = true;
3201 folded.push(label.to_string());
3202 }
3203 let root = self.state.worktree_root();
3205 for j in 1..=self.roles.judges.len() {
3206 let wt = root.join(format!("judge-{j}"));
3207 if wt.exists() {
3208 git::worktree_remove(&repo, &wt).await.ok();
3209 }
3210 }
3211 if self.state.config.graph.advise {
3214 for k in 1..=self.state.config.graph.advisors {
3215 let wt = root.join(format!("advisor-{k}"));
3216 if wt.exists() {
3217 git::worktree_remove(&repo, &wt).await.ok();
3218 }
3219 }
3220 }
3221 if !folded.is_empty() {
3222 self.state
3223 .event("fold", format!("folded candidates {}", folded.join(", ")));
3224 self.state.save()?;
3225 }
3226 Ok(())
3227 }
3228
3229 async fn sync_to_base(&mut self) -> Result<()> {
3269 if self
3270 .state
3271 .base_sync
3272 .as_ref()
3273 .is_some_and(|s| s.conflict.is_some())
3274 {
3275 return Ok(());
3276 }
3277 let Some(winner) = self.state.winner().cloned() else {
3278 return Ok(());
3279 };
3280
3281 let repo = self.state.repo.clone();
3282 let remote = self.state.config.merge.remote.clone();
3283 let base_branch = self.state.base_branch.clone();
3284 let tracking = format!("{remote}/{base_branch}");
3285
3286 git::fetch(&repo, &remote, &base_branch).await.ok();
3287 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3291 return Ok(());
3292 };
3293
3294 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3295 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3296 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3297
3298 if behind == 0 {
3299 if let Some(from) = self
3308 .state
3309 .rebase_fixes
3310 .iter()
3311 .rev()
3312 .find_map(|r| r.from.clone())
3313 && from != head
3314 && git::git_raw(&winner.worktree, &["diff", "--quiet", &from])
3315 .await
3316 .is_ok_and(|o| o.ok())
3317 {
3318 git::sync_to_head(&winner.worktree).await?;
3319 }
3320 self.state.base_sync = Some(BaseSync {
3321 tip,
3322 behind: 0,
3323 attempts,
3324 conflict: None,
3325 });
3326 self.state.save()?;
3327 return Ok(());
3328 }
3329
3330 if attempts >= BASE_SYNC_ROUNDS {
3331 let why = format!(
3332 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3333 rebase(s); rebasing again would only race it",
3334 winner.branch
3335 );
3336 self.state.status = RunStatus::Blocked;
3337 self.state.base_sync = Some(BaseSync {
3338 tip,
3339 behind,
3340 attempts,
3341 conflict: Some(why.clone()),
3342 });
3343 self.state.event("land", why);
3344 self.state.save()?;
3345 return Ok(());
3346 }
3347
3348 self.state.event(
3349 "land",
3350 format!(
3351 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3352 winner.branch
3353 ),
3354 );
3355 self.state.save()?;
3356
3357 let branch_tracking = format!("{remote}/{}", winner.branch);
3363 let fetched_branch = git::fetch(&repo, &remote, &winner.branch).await;
3364 let remote_tip = if matches!(&fetched_branch, Ok(o) if o.ok()) {
3365 git::rev_parse(&repo, &branch_tracking).await.ok()
3366 } else {
3367 None
3368 };
3369 if let Some(theirs) = &remote_tip
3373 && !git::is_ancestor(&repo, theirs, &head).await
3374 && !crate::reconcile::origin_missing(&repo, &head, theirs)
3375 .await
3376 .is_ok_and(|missing| missing.is_empty())
3377 {
3378 let why = format!(
3379 "{branch_tracking} ({}) has commits {} does not contain; not rebasing over \
3380 them",
3381 short(theirs),
3382 winner.branch
3383 );
3384 self.state.status = RunStatus::Blocked;
3385 self.state.base_sync = Some(BaseSync {
3386 tip,
3387 behind,
3388 attempts,
3389 conflict: Some(why.clone()),
3390 });
3391 self.state.event("land", why);
3392 self.state.save()?;
3393 return Ok(());
3394 }
3395
3396 let scratch = self.state.dir().join("base-sync");
3397 let rebased = match crate::rebase::rebase_with_fixer(
3398 &mut self.state,
3399 &scratch,
3400 &winner.branch,
3401 &tracking,
3402 )
3403 .await
3404 {
3405 Ok(crate::rebase::Rebased::Applied) => Ok(None),
3406 Ok(crate::rebase::Rebased::Stopped(why)) => Ok(Some(why)),
3407 Err(e) => Err(e),
3408 };
3409 let attempts = attempts + 1;
3410 match rebased {
3411 Ok(None) => {
3412 git::sync_to_head(&winner.worktree).await?;
3416 let mut conflict = None;
3417 if let Some(pinned) = &remote_tip {
3418 let pushed = git::push_pinned(&repo, &remote, &winner.branch, pinned).await;
3419 match pushed {
3420 Ok(o) if o.ok() => self.state.event(
3421 "land",
3422 format!("pushed rebased {} to {remote}", winner.branch),
3423 ),
3424 Ok(o) => {
3425 conflict = Some(format!(
3426 "rebased {} locally but {remote} refused the push (it moved since {}; someone may have pushed): {}",
3427 winner.branch,
3428 short(pinned),
3429 o.stderr.chars().take(600).collect::<String>()
3430 ));
3431 }
3432 Err(e) => {
3433 conflict = Some(format!(
3434 "rebased {} locally but could not push it: {e:#}",
3435 winner.branch
3436 ));
3437 }
3438 }
3439 }
3440 if let Some(why) = &conflict {
3441 self.state.status = RunStatus::Blocked;
3442 self.state.event("land", why.clone());
3443 }
3444 self.state.base_sync = Some(BaseSync {
3445 tip: tip.clone(),
3446 behind: 0,
3447 attempts,
3448 conflict,
3449 });
3450 self.state
3451 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3452 }
3453 Ok(Some(conflict)) => {
3454 let why = format!(
3455 "{} conflicts with {tracking} and did not rebase: {}",
3456 winner.branch,
3457 conflict.chars().take(600).collect::<String>()
3458 );
3459 self.state.status = RunStatus::Blocked;
3460 self.state.base_sync = Some(BaseSync {
3461 tip,
3462 behind,
3463 attempts,
3464 conflict: Some(why.clone()),
3465 });
3466 self.state.event("land", why);
3467 }
3468 Err(e) => {
3469 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3470 self.state.status = RunStatus::Blocked;
3471 self.state.base_sync = Some(BaseSync {
3472 tip,
3473 behind,
3474 attempts,
3475 conflict: Some(why.clone()),
3476 });
3477 self.state.event("land", why);
3478 }
3479 }
3480 self.state.save()?;
3481 Ok(())
3482 }
3483
3484 fn landing_base(&self) -> String {
3494 self.state
3495 .base_sync
3496 .as_ref()
3497 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3498 }
3499
3500 pub async fn fix_selected(
3533 &mut self,
3534 ids: &[String],
3535 reason: &str,
3536 allow_stale: bool,
3537 ) -> Result<()> {
3538 let reason = reason.trim();
3539 if reason.is_empty() {
3540 bail!("a fix request needs a reason — that is the operator's own record of why");
3541 }
3542 if ids.is_empty() {
3543 bail!("no finding id given");
3544 }
3545 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3546 bail!(
3547 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3548 has already concluded — can be given a targeted fix. A run still \
3549 in progress should simply be resumed; a `merged` run's branch has \
3550 already landed, so its answer is a fresh `magi review <branch>`, \
3551 not reopening this run's own record",
3552 self.state.id,
3553 self.state.status.as_str()
3554 );
3555 }
3556 let Some(winner) = self.state.winner().cloned() else {
3557 bail!("run {} has no winning candidate to fix", self.state.id);
3558 };
3559 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3560 bail!(
3561 "branch `{}` no longer exists; this run cannot be extended",
3562 winner.branch
3563 );
3564 }
3565 let home = crate::run::home();
3566 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3567 bail!(
3568 "run {} is currently being worked on by another magi process",
3569 self.state.id
3570 );
3571 }
3572 let _claim = FixClaim::acquire(&self.state.dir())?;
3578
3579 let mut seen = BTreeSet::new();
3583 let mut findings = Vec::new();
3584 let mut missing = Vec::new();
3585 for id in ids {
3586 if !seen.insert(id.clone()) {
3587 continue;
3588 }
3589 match self.state.finding(id) {
3590 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3591 id: f.id.clone(),
3592 severity: f.severity,
3593 reviewer_vote: rec.vote,
3594 round: round.round,
3595 round_head: round.head.clone(),
3596 reviewer: rec.reviewer,
3597 agent: rec.agent.clone(),
3598 file: f.file.clone(),
3599 line: f.line,
3600 title: f.title.clone(),
3601 detail: f.detail.clone(),
3602 outcome: OperatorFixOutcome::Pending,
3603 }),
3604 None => missing.push(id.clone()),
3605 }
3606 }
3607 if !missing.is_empty() {
3608 bail!(
3609 "unknown finding id(s): {}; nothing was changed",
3610 missing.join(", ")
3611 );
3612 }
3613
3614 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3615 let stale_details: Vec<(String, String)> = findings
3616 .iter()
3617 .filter(|f| f.round_head != head_at_request)
3618 .map(|f| (f.id.clone(), f.round_head.clone()))
3619 .collect();
3620 let stale = !stale_details.is_empty();
3621 if stale && !allow_stale {
3622 bail!(
3623 "the branch has moved since some finding(s) were raised — {} — now \
3624 at {}; pass --allow-stale to fix anyway, or re-run review first",
3625 stale_details
3626 .iter()
3627 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3628 .collect::<Vec<_>>()
3629 .join(", "),
3630 short(&head_at_request)
3631 );
3632 }
3633
3634 let request = OperatorFixRequest {
3635 requested_at: Timestamp::now(),
3636 reason: reason.to_owned(),
3637 findings,
3638 head_at_request: head_at_request.clone(),
3639 allow_stale,
3640 stale,
3641 fix: None,
3642 result_head: None,
3643 follow_up_review_run: None,
3644 };
3645 self.state.event(
3646 "fix",
3647 format!(
3648 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3649 request.findings.len(),
3650 request
3651 .findings
3652 .iter()
3653 .map(|f| f.id.as_str())
3654 .collect::<Vec<_>>()
3655 .join(", "),
3656 ),
3657 );
3658 self.state.operator_fixes.push(request);
3665 self.state.save()?;
3666 let request_index = self.state.operator_fixes.len() - 1;
3667
3668 if winner.worktree.exists() {
3677 let dirty = git::git(
3680 &winner.worktree,
3681 &["status", "--porcelain", "--untracked-files=all"],
3682 )
3683 .await?;
3684 let only_withheld = dirty.lines().all(|l| {
3685 l.strip_prefix("?? ")
3686 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3687 });
3688 if !only_withheld {
3689 bail!(
3690 "`{}` has uncommitted changes; refusing to touch it — commit or \
3691 discard them first",
3692 winner.worktree.display()
3693 );
3694 }
3695 git::worktree_remove(&self.state.repo, &winner.worktree)
3696 .await
3697 .ok();
3698 }
3699 let fix_worktree = self.state.worktree_root().join("operator-fix");
3700 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3701 git::git(
3702 &self.state.repo,
3703 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3704 )
3705 .await
3706 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3707 if !git::is_clean(&fix_worktree).await? {
3708 git::worktree_remove(&self.state.repo, &fix_worktree)
3709 .await
3710 .ok();
3711 bail!(
3712 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3713 winner.branch
3714 );
3715 }
3716
3717 let run_id = self.state.id.clone();
3718 let prompts = self.state.config.prompts.clone();
3719 let language = self.state.config.graph.language.clone();
3720 let sessions = self.state.config.graph.sessions;
3721 let artifacts = agent::artifacts_dir(&self.state.dir());
3722 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3723 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3724 _ => (
3725 self.state
3726 .config
3727 .agent(&winner.agent)
3728 .cloned()
3729 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3730 format!("impl-{}", winner.label),
3731 ),
3732 };
3733 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3734 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3735 .findings
3736 .iter()
3737 .map(|f| Finding {
3738 id: f.id.clone(),
3739 severity: f.severity,
3740 file: f.file.clone(),
3741 line: f.line,
3742 title: f.title.clone(),
3743 detail: f.detail.clone(),
3744 })
3745 .collect();
3746 let job = SeatJob {
3747 prompt: prompt::operator_fix(
3748 &self.state.instruction,
3749 &finding_list,
3750 reason,
3751 &stale_details,
3752 &head_at_request,
3753 &language,
3754 ),
3755 spec: fix_spec.clone(),
3756 seat,
3757 cwd: fix_worktree.clone(),
3758 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3759 allow_write: true,
3760 sessions,
3761 artifacts: artifacts.clone(),
3762 stem: "operator-fix".to_owned(),
3763 };
3764 let cache = self.state.config.cache_dir();
3765 let ctx = WaveCtx {
3766 run: &run_id,
3767 node: "fix",
3768 prompts: &prompts,
3769 cache: cache.as_deref(),
3770 round: None,
3771 };
3772 let (seat, out) =
3773 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3774 let agent_id = seat.agent.clone();
3775
3776 let mut fix = FixRecord {
3777 agent: agent_id,
3778 addressed: Vec::new(),
3779 rejected: Vec::new(),
3780 notes: String::new(),
3781 committed: false,
3782 failed: None,
3783 duration_ms: 0,
3784 continuation: None,
3785 };
3786 let mut final_seat = seat.clone();
3787 match out {
3788 AgentOutcome::Ok(o) => {
3789 fix.duration_ms = o.duration_ms;
3790 let parsed = verdict::extract_json::<FixReport>(&o.text);
3791 let incomplete_reason = match &parsed {
3792 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3793 "the reply parsed, but it reported a command whose own CLI \
3794 never confirmed an exit status"
3795 .to_owned(),
3796 ),
3797 Ok(_) => None,
3798 Err(e) => Some(e.to_string()),
3799 };
3800 match incomplete_reason {
3801 None => {
3802 let report = parsed.expect("checked Ok above");
3803 fix.addressed = report.addressed;
3804 fix.rejected = report.rejected;
3805 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3806 }
3807 Some(reason) => {
3808 let (resumed_seat, resolved, failure, cont) = self
3809 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3810 .await;
3811 fix.duration_ms += cont.cumulative_wait_ms;
3812 fix.continuation = Some(cont);
3813 final_seat = resumed_seat;
3814 match resolved {
3815 Some(report) => {
3816 fix.addressed = report.addressed;
3817 fix.rejected = report.rejected;
3818 fix.notes =
3819 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3820 }
3821 None => fix.failed = failure,
3822 }
3823 }
3824 }
3825 }
3826 AgentOutcome::Dropped(o) => {
3827 fix.duration_ms = o.duration_ms;
3828 let why = o
3829 .dropped
3830 .as_ref()
3831 .map(|d| d.why.as_str())
3832 .unwrap_or("the CLI ended the stream without delivering its answer");
3833 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3834 }
3835 AgentOutcome::Quota(o) => {
3836 self.state.quota.push(QuotaLoss {
3837 seat: final_seat.key.clone(),
3838 node: "fix".to_owned(),
3839 at: Timestamp::now(),
3840 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3841 });
3842 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3843 }
3844 AgentOutcome::Failed(e) => fix.failed = Some(e),
3845 }
3846 if fix.continuation.is_none() {
3847 fix.continuation = Some(ContinuationRecord::not_needed());
3848 }
3849 self.state.seats.insert(final_seat.key.clone(), final_seat);
3850
3851 let rescue_message = format!(
3852 "magi: operator-selected fix ({}) (uncommitted work)",
3853 self.state.operator_fixes[request_index]
3854 .findings
3855 .iter()
3856 .map(|f| f.id.as_str())
3857 .collect::<Vec<_>>()
3858 .join(", ")
3859 );
3860 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3861 self.state.note_withheld("fix", &r.withheld);
3862 }
3863 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3864 fix.committed = after != head_at_request;
3865 git::worktree_remove(&self.state.repo, &fix_worktree)
3866 .await
3867 .ok();
3868
3869 self.state.event(
3870 "fix",
3871 match &fix.failed {
3872 Some(reason) => format!(
3873 "operator fix: adoption report was lost ({reason}); {}",
3874 if fix.committed {
3875 "committed"
3876 } else {
3877 "NO new commit"
3878 }
3879 ),
3880 None => format!(
3881 "operator fix: {} addressed, {} rejected, {}",
3882 fix.addressed.len(),
3883 fix.rejected.len(),
3884 if fix.committed {
3885 "committed"
3886 } else {
3887 "NO new commit"
3888 }
3889 ),
3890 },
3891 );
3892
3893 for f in &mut self.state.operator_fixes[request_index].findings {
3900 f.outcome = if fix.failed.is_some() {
3901 OperatorFixOutcome::Unreported
3902 } else if fix.addressed.contains(&f.id) {
3903 OperatorFixOutcome::Addressed
3904 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3905 OperatorFixOutcome::Rejected { why: r.why.clone() }
3906 } else {
3907 OperatorFixOutcome::Unreported
3908 };
3909 }
3910
3911 let committed = fix.committed;
3912 if committed {
3913 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3914 }
3915 self.state.operator_fixes[request_index].fix = Some(fix);
3916 self.state.save()?;
3919
3920 if committed {
3921 self.state.event(
3922 "fix",
3923 format!(
3924 "operator fix committed {}; opening a follow-up review-only run",
3925 short(&after)
3926 ),
3927 );
3928 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3929 Ok(mut follow_up) => {
3930 follow_up.state.event(
3931 "start",
3932 format!(
3933 "requested by an operator fix on run {} for finding(s) {}",
3934 self.state.id,
3935 self.state.operator_fixes[request_index]
3936 .findings
3937 .iter()
3938 .map(|f| f.id.as_str())
3939 .collect::<Vec<_>>()
3940 .join(", "),
3941 ),
3942 );
3943 follow_up.state.save()?;
3944 let follow_up_id = follow_up.state.id.clone();
3945 if let Err(e) = follow_up.execute().await {
3946 self.state.event(
3947 "fix",
3948 format!(
3949 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3950 ),
3951 );
3952 }
3953 self.state.operator_fixes[request_index].follow_up_review_run =
3954 Some(follow_up_id);
3955 }
3956 Err(e) => {
3957 self.state.event(
3958 "fix",
3959 format!("committed the fix but could not open a follow-up review: {e:#}"),
3960 );
3961 }
3962 }
3963 self.state.save()?;
3964 }
3965
3966 Ok(())
3967 }
3968
3969 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3976 match &self.roles.fixer {
3977 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3978 _ => (
3979 self.state
3980 .config
3981 .agent(&winner.agent)
3982 .cloned()
3983 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3984 format!("impl-{}", winner.label),
3985 ),
3986 }
3987 }
3988
3989 async fn review_loop(&mut self) -> Result<()> {
3990 if self
3995 .state
3996 .base_sync
3997 .as_ref()
3998 .is_some_and(|s| s.conflict.is_some())
3999 {
4000 return Ok(());
4001 }
4002 let run_id = self.state.id.clone();
4007 let prompts = self.state.config.prompts.clone();
4008 let Some(winner) = self.state.winner().cloned() else {
4009 return Ok(());
4010 };
4011 let max_rounds = self.state.config.graph.review_rounds;
4012 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
4022 self.state.status = status;
4023 self.state.save()?;
4024 return Ok(());
4025 }
4026 self.state.status = RunStatus::Reviewing;
4027 if self
4037 .state
4038 .reviews
4039 .last()
4040 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
4041 {
4042 let shell = self.state.config.shell();
4043 return self
4044 .stop_reviewing(
4045 "the last round's own verification never resolved",
4046 &shell,
4047 &winner.worktree,
4048 )
4049 .await;
4050 }
4051
4052 let repo = self.state.repo.clone();
4053 let root = self.state.worktree_root();
4054 let language = self.state.config.graph.language.clone();
4055 let sessions = self.state.config.graph.sessions;
4056 let artifacts = agent::artifacts_dir(&self.state.dir());
4057 let base = self.landing_base();
4058 let base_short = short(&base);
4059 let reviewers = self.roles.reviewers.clone();
4060 let shell = self.state.config.shell();
4061
4062 for round in (self.state.reviews.len() + 1)..=max_rounds {
4063 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4064 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
4065 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
4066 let prev_verification = self
4075 .state
4076 .reviews
4077 .last()
4078 .and_then(|r| r.verification_summary(&head));
4079
4080 let mut jobs = Vec::new();
4084 for (r, spec) in reviewers.iter().cloned().enumerate() {
4085 let wt = root.join(format!("review-{}", r + 1));
4086 if wt.exists() {
4087 git::reset_detached(&wt, &head).await?;
4088 } else {
4089 git::worktree_add_detached(&repo, &wt, &head).await?;
4090 }
4091 let seat_key = format!("review-{}", r + 1);
4092 let seat = self.seat(&seat_key, &spec.id);
4093 jobs.push(SeatJob {
4094 prompt: prompt::review(&prompt::ReviewCtx {
4095 instruction: &self.state.instruction,
4096 branch: &winner.branch,
4097 base_short: &base_short,
4098 stat: &stat,
4099 patch: &patch,
4100 verification: prev_verification.as_ref(),
4101 reviewers: reviewers.len(),
4102 round,
4103 rounds: max_rounds,
4104 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
4107 lens: Lens::for_seat(r),
4108 language: &language,
4109 }),
4110 spec,
4111 seat,
4112 cwd: wt,
4113 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4114 allow_write: false,
4115 sessions,
4116 artifacts: artifacts.clone(),
4117 stem: format!("review-{round}-{}", r + 1),
4118 });
4119 }
4120
4121 self.state.event(
4122 "review",
4123 format!(
4124 "round {round}: {} reviewers on {}",
4125 jobs.len(),
4126 short(&head)
4127 ),
4128 );
4129 let mut quota_losses = Vec::new();
4130 let review_retries = self.state.config.graph.retries;
4131 let review_cache = self.state.config.cache_dir();
4132 let ctx = WaveCtx {
4133 run: &run_id,
4134 node: "review",
4135 prompts: &prompts,
4136 cache: review_cache.as_deref(),
4137 round: Some(round),
4138 };
4139 let results = ask_json_wave::<Review>(
4140 jobs,
4141 Arc::clone(&self.sem),
4142 review_retries,
4143 &ctx,
4144 &mut quota_losses,
4145 &mut self.state,
4146 &|_: &Review| Ok(()),
4147 )
4148 .await;
4149 let round_quota_missing = quota_losses.len();
4153 self.state.quota.extend(quota_losses);
4154
4155 let mut records = Vec::new();
4156 let mut all_findings = Vec::new();
4157 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
4158 let agent_id = seat.agent.clone();
4159 self.state.seats.insert(seat.key.clone(), seat);
4160 let mut record = ReviewRecord {
4161 reviewer: r + 1,
4162 agent: agent_id,
4163 summary: String::new(),
4164 findings: Vec::new(),
4165 vote: None,
4166 failed: None,
4167 duration_ms: 0,
4168 attempts,
4174 };
4175 match res {
4176 Ok((review, out)) => {
4177 record.summary =
4185 blind::sanitize_prose(&review.summary, &self.state.config.blind);
4186 record.vote = Some(review.vote);
4187 record.duration_ms = out.duration_ms;
4188 for (n, mut f) in review.findings.into_iter().enumerate() {
4189 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
4192 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
4193 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
4194 f.file = f
4200 .file
4201 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
4202 all_findings.push(f.clone());
4203 record.findings.push(f);
4204 }
4205 self.state.event(
4206 "review",
4207 format!(
4208 "round {round}: reviewer {} voted {} with {} finding(s)",
4209 r + 1,
4210 review.vote.label(),
4211 record.findings.len()
4212 ),
4213 );
4214 }
4215 Err(e) => {
4216 record.failed = Some(e.to_string());
4217 self.state.event(
4218 "review",
4219 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
4220 );
4221 }
4222 }
4223 records.push(record);
4224 }
4225
4226 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
4233 let vote_split =
4234 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
4235 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
4236 if vote_split {
4237 self.state.event(
4238 "review",
4239 format!(
4240 "round {round}: votes split ({}) — one round of reconsideration",
4241 initial_votes
4242 .iter()
4243 .map(|v| v.label())
4244 .collect::<Vec<_>>()
4245 .join(", ")
4246 ),
4247 );
4248 let panel: Vec<ReviewSeatReport<'_>> = records
4251 .iter()
4252 .filter_map(|r| {
4253 r.vote.map(|vote| ReviewSeatReport {
4254 reviewer: r.reviewer,
4255 vote,
4256 summary: &r.summary,
4257 findings: &r.findings,
4258 })
4259 })
4260 .collect();
4261
4262 let mut jobs = Vec::new();
4263 let mut seats_at = Vec::new();
4264 for (r, spec) in reviewers.iter().cloned().enumerate() {
4265 if records[r].vote.is_none() {
4269 continue;
4270 }
4271 let wt = root.join(format!("review-{}", r + 1));
4272 let seat_key = format!("review-{}", r + 1);
4273 let seat = self.seat(&seat_key, &spec.id);
4274 let patch_ctx = if has_context(&spec, &seat, sessions) {
4279 None
4280 } else {
4281 Some(ReviewPatch {
4282 branch: &winner.branch,
4283 base_short: &base_short,
4284 stat: &stat,
4285 patch: &patch,
4286 })
4287 };
4288 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
4289 instruction: &self.state.instruction,
4290 reviewer: r + 1,
4291 lens: Lens::for_seat(r),
4292 panel: &panel,
4293 patch: patch_ctx,
4294 round,
4295 rounds: max_rounds,
4296 language: &language,
4297 });
4298 jobs.push(SeatJob {
4299 prompt,
4300 spec,
4301 seat,
4302 cwd: wt,
4303 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4304 allow_write: false,
4305 sessions,
4306 artifacts: artifacts.clone(),
4307 stem: format!("review-{round}-reconsider-{}", r + 1),
4308 });
4309 seats_at.push(r);
4310 }
4311
4312 let mut recon_quota_losses = Vec::new();
4313 let recon_cache = self.state.config.cache_dir();
4314 let recon_ctx = WaveCtx {
4315 run: &run_id,
4316 node: "review",
4317 prompts: &prompts,
4318 cache: recon_cache.as_deref(),
4319 round: Some(round),
4320 };
4321 let recon_results = ask_json_wave::<ReviewRevote>(
4322 jobs,
4323 Arc::clone(&self.sem),
4324 review_retries,
4325 &recon_ctx,
4326 &mut recon_quota_losses,
4327 &mut self.state,
4328 &|_: &ReviewRevote| Ok(()),
4329 )
4330 .await;
4331 self.state.quota.extend(recon_quota_losses);
4332
4333 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
4334 let agent_id = seat.agent.clone();
4335 self.state.seats.insert(seat.key.clone(), seat);
4336 let mut rec = ReviewRevoteRecord {
4337 reviewer: r + 1,
4338 agent: agent_id,
4339 vote: None,
4340 reason: String::new(),
4341 failed: None,
4342 };
4343 match res {
4344 Ok((rv, _)) => {
4345 rec.vote = Some(rv.vote);
4346 rec.reason =
4347 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
4348 self.state.event(
4349 "review",
4350 format!(
4351 "round {round}: reviewer {} revoted {}",
4352 r + 1,
4353 rv.vote.label()
4354 ),
4355 );
4356 }
4357 Err(e) => {
4358 rec.failed = Some(e.to_string());
4359 self.state.event(
4360 "review",
4361 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4362 );
4363 }
4364 }
4365 reconsideration.push(rec);
4366 }
4367 } else if initial_votes.len() > 1 {
4368 self.state.event(
4369 "review",
4370 format!(
4371 "round {round}: votes agreed ({}) — no reconsideration",
4372 initial_votes[0].label()
4373 ),
4374 );
4375 }
4376
4377 let final_votes: Vec<ReviewVote> = records
4381 .iter()
4382 .filter_map(|r| {
4383 reconsideration
4384 .iter()
4385 .find(|rv| rv.reviewer == r.reviewer)
4386 .and_then(|rv| rv.vote)
4387 .or(r.vote)
4388 })
4389 .collect();
4390 let round_verdict = ReviewVote::worst(final_votes);
4391
4392 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4393 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4394 let defer_e2e =
4405 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4406 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4407 let reason =
4408 format!("{blocking} blocking finding(s) already required a fix this round");
4409 self.state.event(
4410 "verify",
4411 format!(
4412 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4413 {}); it will run once a round has none left",
4414 short(&head)
4415 ),
4416 );
4417 (Vec::new(), false, true, Some(reason))
4418 } else {
4419 let e2e_commands = self.state.config.verify.e2e.clone();
4420 let cache_dir = self.state.config.cache_dir();
4421 let context = format!("round {round}");
4422 let (e2e, verify_retried) = with_cache_lease(
4423 &mut self.state,
4424 cache_dir.as_deref(),
4425 "e2e",
4426 "e2e",
4427 &winner.worktree,
4428 &head,
4429 verify_timeout,
4430 &context,
4431 |state, budget| {
4432 let shell = shell.clone();
4433 let e2e_commands = e2e_commands.clone();
4434 let worktree = winner.worktree.clone();
4435 let context = context.clone();
4436 async move {
4437 run_e2e_with_retry(
4438 state,
4439 &shell,
4440 &e2e_commands,
4441 &worktree,
4442 budget,
4443 &context,
4444 )
4445 .await
4446 }
4447 },
4448 )
4449 .await;
4450 (e2e, verify_retried, false, None)
4451 };
4452
4453 let expected = records.len();
4454 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4455 let incomplete = answered < expected;
4456 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4457 let policy = self.state.config.graph.incomplete_review;
4458 let clean = round_is_clean(
4459 blocking,
4460 e2e_ok,
4461 answered,
4462 expected,
4463 round_quota_missing,
4464 policy,
4465 );
4466
4467 let mut round_record = ReviewRound {
4468 round,
4469 head: head.clone(),
4470 verified_head: None,
4471 verified_at: None,
4472 reviews: records,
4473 e2e,
4474 verify_retried,
4475 e2e_deferred,
4476 e2e_defer_reason,
4477 fix: None,
4478 blocking,
4479 answered,
4480 expected,
4481 clean,
4482 progressed: false,
4483 vote_split,
4484 reconsideration,
4485 verdict: round_verdict,
4486 };
4487 if !matches!(
4498 round_record.e2e_status(),
4499 E2eStatus::Deferred | E2eStatus::NotConfigured
4500 ) {
4501 round_record.verified_head = Some(head.clone());
4502 round_record.verified_at = Some(Timestamp::now());
4503 }
4504 let this_round_verification = round_record.verification_summary(&head);
4505
4506 if incomplete {
4507 let missing: Vec<String> = round_record
4508 .reviews
4509 .iter()
4510 .filter(|r| r.failed.is_some())
4511 .map(|r| format!("review-{}", r.reviewer))
4512 .collect();
4513 self.state.event(
4514 "review",
4515 format!(
4516 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4517 missing.join(", ")
4518 ),
4519 );
4520 }
4521
4522 if clean {
4523 self.state.event(
4524 "review",
4525 if incomplete && policy == IncompleteReviewPolicy::Warn {
4526 format!(
4527 "round {round}: clean (warn policy, incomplete panel) — no \
4528 blocking findings from the seats that answered, verification green"
4529 )
4530 } else if incomplete {
4531 format!(
4532 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4533 quorum) — no blocking findings from the seats that answered, \
4534 verification green",
4535 expected - answered
4536 )
4537 } else {
4538 format!("round {round}: clean — no blocking findings, verification green")
4539 },
4540 );
4541 self.state.reviews.push(round_record);
4542 self.state.status = RunStatus::Gating;
4543 self.state.save()?;
4544 return Ok(());
4545 }
4546
4547 if incomplete && blocking == 0 && e2e_ok {
4555 self.state.reviews.push(round_record);
4556 self.state.save()?;
4557 if round == max_rounds {
4558 self.state.status = RunStatus::Blocked;
4559 self.state.event(
4560 "review",
4561 format!(
4562 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4563 refusing to call it clean",
4564 expected - answered
4565 ),
4566 );
4567 return Ok(());
4568 }
4569 continue;
4570 }
4571
4572 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4584 self.state.reviews.push(round_record);
4585 return self
4586 .stop_reviewing(
4587 "the round's own verification could not run",
4588 &shell,
4589 &winner.worktree,
4590 )
4591 .await;
4592 }
4593
4594 if round == max_rounds {
4595 self.state.reviews.push(round_record);
4596 return self
4597 .stop_reviewing(
4598 &format!(
4599 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4600 ),
4601 &shell,
4602 &winner.worktree,
4603 )
4604 .await;
4605 }
4606
4607 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4610 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4611 let blocking_findings: Vec<_> = all_findings
4612 .iter()
4613 .filter(|f| f.severity.blocks())
4614 .cloned()
4615 .collect();
4616 let job = SeatJob {
4617 prompt: prompt::fix(
4618 &self.state.instruction,
4619 &blocking_findings,
4620 this_round_verification.as_ref(),
4621 round,
4622 max_rounds,
4623 &language,
4624 ),
4625 spec: fix_spec.clone(),
4626 seat,
4627 cwd: winner.worktree.clone(),
4628 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4629 allow_write: true,
4630 sessions,
4631 artifacts: artifacts.clone(),
4632 stem: format!("fix-{round}"),
4633 };
4634 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4635 let cache = self.state.config.cache_dir();
4636 let ctx = WaveCtx {
4637 run: &run_id,
4638 node: "fix",
4639 prompts: &prompts,
4640 cache: cache.as_deref(),
4641 round: Some(round),
4642 };
4643 let (seat, out) =
4644 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4645 let agent_id = seat.agent.clone();
4646
4647 let mut fix = FixRecord {
4648 agent: agent_id,
4649 addressed: Vec::new(),
4650 rejected: Vec::new(),
4651 notes: String::new(),
4652 committed: false,
4653 failed: None,
4654 duration_ms: 0,
4655 continuation: None,
4656 };
4657 let mut continuation = ContinuationRecord::not_needed();
4658 let mut final_seat = seat.clone();
4659 match out {
4660 AgentOutcome::Ok(o) => {
4661 fix.duration_ms = o.duration_ms;
4662 let parsed = verdict::extract_json::<FixReport>(&o.text);
4663 let incomplete_reason = match &parsed {
4670 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4671 "the reply parsed, but it reported a command whose own CLI never \
4672 confirmed an exit status"
4673 .to_owned(),
4674 ),
4675 Ok(_) => None,
4676 Err(e) => Some(e.to_string()),
4677 };
4678 match incomplete_reason {
4679 None => {
4680 let report = parsed.expect("checked Ok above");
4681 fix.addressed = report.addressed;
4682 fix.rejected = report.rejected;
4683 fix.notes =
4684 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4685 }
4686 Some(reason) => {
4687 let (resumed_seat, resolved, failure, cont) = self
4688 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4689 .await;
4690 fix.duration_ms += cont.cumulative_wait_ms;
4691 continuation = cont;
4692 final_seat = resumed_seat;
4693 match resolved {
4694 Some(report) => {
4695 fix.addressed = report.addressed;
4696 fix.rejected = report.rejected;
4697 fix.notes = blind::sanitize_prose(
4698 &report.notes,
4699 &self.state.config.blind,
4700 );
4701 }
4702 None => fix.failed = failure,
4703 }
4704 }
4705 }
4706 }
4707 AgentOutcome::Dropped(o) => {
4709 fix.duration_ms = o.duration_ms;
4710 let why = o
4711 .dropped
4712 .as_ref()
4713 .map(|d| d.why.as_str())
4714 .unwrap_or("the CLI ended the stream without delivering its answer");
4715 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4716 }
4717 AgentOutcome::Quota(o) => {
4718 self.state.quota.push(QuotaLoss {
4719 seat: final_seat.key.clone(),
4720 node: "fix".to_owned(),
4721 at: Timestamp::now(),
4722 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4723 });
4724 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4725 }
4726 AgentOutcome::Failed(e) => fix.failed = Some(e),
4727 }
4728 fix.continuation = Some(continuation);
4729 self.state.seats.insert(final_seat.key.clone(), final_seat);
4730 if let Ok(r) = git::rescue_commit(
4731 &winner.worktree,
4732 &format!("magi: review round {round} fixes (uncommitted work)"),
4733 )
4734 .await
4735 {
4736 self.state.note_withheld("fix", &r.withheld);
4737 }
4738 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4739 fix.committed = after != before;
4740 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4748 let progressed = diff_after != patch;
4749 let commit_note = if fix.committed {
4750 "committed"
4751 } else {
4752 "NO new commit"
4753 };
4754 let tree_note = if progressed {
4755 "changed vs base"
4756 } else {
4757 "unchanged vs base"
4758 };
4759 self.state.event(
4760 "fix",
4761 match &fix.failed {
4762 Some(reason) => {
4768 format!(
4769 "round {round}: fixer's adoption report was lost ({reason}); \
4770 {commit_note}, tree {tree_note}"
4771 )
4772 }
4773 None => format!(
4774 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4775 {tree_note}{}",
4776 fix.addressed.len(),
4777 fix.rejected.len(),
4778 if continuation.outcome == ContinuationOutcome::Resumed {
4779 format!(
4780 " (adoption report recovered after {} continuation(s))",
4781 continuation.attempts
4782 )
4783 } else {
4784 String::new()
4785 },
4786 ),
4787 },
4788 );
4789 round_record.fix = Some(fix);
4790 round_record.progressed = progressed;
4791 self.state.reviews.push(round_record);
4792 self.state.save()?;
4793
4794 if matches!(
4807 continuation.outcome,
4808 ContinuationOutcome::Exhausted
4809 | ContinuationOutcome::QuotaLost
4810 | ContinuationOutcome::NoSession
4811 ) {
4812 return self
4813 .stop_reviewing(
4814 "the fixer's adoption report never came back, even after resuming its \
4815 own seat; refusing to start another round against the same worktree \
4816 while that is unresolved",
4817 &shell,
4818 &winner.worktree,
4819 )
4820 .await;
4821 }
4822
4823 let streak = self
4824 .state
4825 .reviews
4826 .iter()
4827 .rev()
4828 .take_while(|r| !r.progressed)
4829 .count();
4830 if streak >= STAGNANT_LIMIT {
4831 return self
4832 .stop_reviewing(
4833 &format!(
4834 "the tree has not moved against base for {streak} round(s) in a row"
4835 ),
4836 &shell,
4837 &winner.worktree,
4838 )
4839 .await;
4840 }
4841 }
4842 Ok(())
4843 }
4844
4845 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4875 let round_idx = self.state.reviews.len() - 1;
4876 let needs_catchup_run = matches!(
4884 self.state.reviews[round_idx].e2e_status(),
4885 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4886 );
4887 if needs_catchup_run {
4888 let round = self.state.reviews[round_idx].round;
4889 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4890 let commands = self.state.config.verify.e2e.clone();
4891 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4892 let cache_dir = self.state.config.cache_dir();
4893 let context = format!(
4894 "round {round}: verification unresolved, catching up before the final decision"
4895 );
4896 let (outcomes, verify_retried) = with_cache_lease(
4897 &mut self.state,
4898 cache_dir.as_deref(),
4899 "e2e",
4900 "e2e",
4901 worktree,
4902 &attempted_head,
4903 timeout,
4904 &context,
4905 |state, budget| {
4906 let shell = shell.to_vec();
4907 let commands = commands.clone();
4908 let context = context.clone();
4909 async move {
4910 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4911 .await
4912 }
4913 },
4914 )
4915 .await;
4916 let last = &mut self.state.reviews[round_idx];
4917 last.e2e = outcomes;
4918 last.verify_retried = verify_retried;
4919 last.verified_head = Some(attempted_head);
4926 last.verified_at = Some(Timestamp::now());
4927 if verify_inconclusive(&last.e2e) {
4928 self.state.save()?;
4935 return Ok(());
4936 }
4937 last.e2e_deferred = false;
4938 }
4939 let last = &self.state.reviews[round_idx];
4940 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4941
4942 match last.e2e_status() {
4943 E2eStatus::Failed => {
4944 let red: Vec<String> = last
4945 .e2e
4946 .iter()
4947 .filter(|o| !o.ok())
4948 .map(|o| {
4949 format!(
4950 "`{}` -> {:?}\n{}",
4951 o.command,
4952 o.code,
4953 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4954 )
4955 })
4956 .collect();
4957 self.state
4958 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4959 self.state.status = RunStatus::Blocked;
4960 }
4961 E2eStatus::ResourceBlocked => {
4966 self.state.event(
4967 "review",
4968 format!(
4969 "{why}; e2e could not run (shared build cache unavailable); not \
4970 deciding yet"
4971 ),
4972 );
4973 }
4974 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4975 self.state.event(
4976 "review",
4977 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4978 );
4979 self.state.status = RunStatus::Gating;
4980 }
4981 }
4982 self.state.save()?;
4983 Ok(())
4984 }
4985
4986 async fn gate(&mut self) -> Result<()> {
4989 if self.state.status == RunStatus::Failed
5001 || self
5002 .state
5003 .base_sync
5004 .as_ref()
5005 .is_some_and(|s| s.conflict.is_some())
5006 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5007 != Some(RunStatus::Gating)
5008 {
5009 return Ok(());
5010 }
5011 if self.state.gate_ran {
5012 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
5023 self.state.status = RunStatus::Blocked;
5024 self.state.save()?;
5025 }
5026 return Ok(());
5027 }
5028 let Some(winner) = self.state.winner().cloned() else {
5029 return Ok(());
5030 };
5031 self.state.status = RunStatus::Gating;
5032 let mut outcomes = self.run_gate(&winner).await?;
5033 loop {
5034 if verify_inconclusive(&outcomes) {
5045 self.state.save()?;
5046 return Ok(());
5047 }
5048 if outcomes.iter().all(CommandOutcome::ok) {
5049 break;
5050 }
5051 match self.gate_fix_round(&winner, &outcomes).await? {
5052 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
5053 GateFix::Stop => break,
5054 GateFix::Defer => {
5055 self.state.save()?;
5056 return Ok(());
5057 }
5058 }
5059 }
5060 let passed = outcomes.iter().all(CommandOutcome::ok);
5061 self.state.gate = outcomes;
5062 self.state.gate_ran = true;
5063 if !passed {
5064 self.state.status = RunStatus::Blocked;
5065 let spent = self.state.gate_fixes.len();
5066 self.state.event(
5067 "gate",
5068 if spent == 0 {
5069 "gate failed; not merging".to_owned()
5070 } else {
5071 format!("gate failed after {spent} gate-fix round(s); not merging")
5072 },
5073 );
5074 }
5075 self.state.save()?;
5076 Ok(())
5077 }
5078
5079 async fn run_pre_gate(&mut self, winner: &Candidate) {
5089 let commands = self.state.config.verify.pre_gate.clone();
5090 if commands.is_empty() {
5091 return;
5092 }
5093 let shell = self.state.config.shell();
5094 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5095 let (outcomes, _) = run_commands(
5096 &mut self.state,
5097 "pre_gate",
5098 "pre_gate",
5099 0,
5100 &shell,
5101 &commands,
5102 &winner.worktree,
5103 timeout,
5104 )
5105 .await;
5106 for o in &outcomes {
5107 if !o.ok() {
5108 tracing::warn!(
5109 "pre_gate `{}` failed ({:?}); the gate decides",
5110 o.command,
5111 o.code
5112 );
5113 }
5114 self.state.event(
5115 "pre_gate",
5116 format!(
5117 "`{}` -> {}",
5118 o.command,
5119 if o.ok() {
5120 "pass".to_owned()
5121 } else {
5122 format!(
5123 "FAIL ({:?})\n{}",
5124 o.code,
5125 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5126 )
5127 }
5128 ),
5129 );
5130 }
5131 self.state.pre_gate = outcomes;
5132 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
5133 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
5134 Ok(head) => {
5135 self.state
5136 .event("pre_gate", format!("committed mechanical fixes ({head})"));
5137 self.state.pre_gate_commit = Some(head);
5138 }
5139 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
5140 },
5141 Ok(false) => {}
5142 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
5143 }
5144 if let Err(e) = self.state.save() {
5145 tracing::warn!("could not persist the pre_gate record: {e:#}");
5146 }
5147 }
5148
5149 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
5152 self.run_pre_gate(winner).await;
5153 let shell = self.state.config.shell();
5154 let gate_commands = self.state.config.verify.gate.clone();
5155 let outcomes = if gate_commands.is_empty() {
5164 Vec::new()
5165 } else {
5166 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5167 let cache_dir = self.state.config.cache_dir();
5168 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
5169 let (outcomes, _) = with_cache_lease(
5170 &mut self.state,
5171 cache_dir.as_deref(),
5172 "gate",
5173 "gate",
5174 &winner.worktree,
5175 &head,
5176 timeout,
5177 "final gate",
5178 |state, budget| {
5179 let shell = shell.clone();
5180 let gate_commands = gate_commands.clone();
5181 let worktree = winner.worktree.clone();
5182 async move {
5183 let (outcomes, timed_out_pids) = run_commands(
5184 state,
5185 "gate",
5186 "gate",
5187 0,
5188 &shell,
5189 &gate_commands,
5190 &worktree,
5191 budget,
5192 )
5193 .await;
5194 (outcomes, false, timed_out_pids)
5195 }
5196 },
5197 )
5198 .await;
5199 outcomes
5200 };
5201 if outcomes.is_empty() {
5202 self.state.event(
5207 "gate",
5208 "no gate commands configured; nothing to check, passing",
5209 );
5210 }
5211 for o in &outcomes {
5212 self.state.event(
5213 "gate",
5214 format!(
5215 "`{}` -> {}",
5216 o.command,
5217 if o.ok() {
5218 "pass".to_owned()
5219 } else {
5220 format!(
5221 "FAIL ({:?})\n{}",
5222 o.code,
5223 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5224 )
5225 }
5226 ),
5227 );
5228 }
5229 Ok(outcomes)
5230 }
5231
5232 async fn gate_fix_round(
5245 &mut self,
5246 winner: &Candidate,
5247 outcomes: &[CommandOutcome],
5248 ) -> Result<GateFix> {
5249 let cap = self.state.config.graph.gate_fix_rounds;
5250 let spent = self.state.gate_fixes.len();
5251 if spent >= cap {
5252 if cap > 0 {
5253 self.state.event(
5254 "gate",
5255 format!("{spent} gate-fix round(s) spent and the gate still fails"),
5256 );
5257 }
5258 return Ok(GateFix::Stop);
5259 }
5260 if !gate_fixable(outcomes) {
5261 self.state.event(
5262 "gate",
5263 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
5264 command or similar); not spending a fix round on it",
5265 );
5266 return Ok(GateFix::Stop);
5267 }
5268 let min_free = self.state.config.disk.min_free_bytes;
5269 if min_free > 0 {
5270 match crate::disk::free_bytes(&winner.worktree) {
5271 Ok(free) if crate::disk::enough_space(free, min_free) => {}
5272 Ok(free) => {
5273 self.state.event(
5274 "gate",
5275 format!(
5276 "only {free} bytes free ({min_free} required by `[disk] \
5277 min_free_bytes`); not spending a fix round on a failure the disk \
5278 may explain"
5279 ),
5280 );
5281 return Ok(GateFix::Stop);
5282 }
5283 Err(e) => {
5284 self.state.event(
5285 "gate",
5286 format!("free disk space could not be measured ({e:#}); no fix round"),
5287 );
5288 return Ok(GateFix::Stop);
5289 }
5290 }
5291 }
5292
5293 let attempt = spent + 1;
5294 let run_id = self.state.id.clone();
5295 let prompts = self.state.config.prompts.clone();
5296 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
5297 let base = self.landing_base();
5298 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
5299 let seat = self.seat(&fix_seat_key, &fix_spec.id);
5300 let job = SeatJob {
5301 prompt: prompt::gate_fix(
5302 &self.state.instruction,
5303 &failed,
5304 attempt,
5305 cap,
5306 &self.state.config.graph.language,
5307 ),
5308 spec: fix_spec,
5309 seat,
5310 cwd: winner.worktree.clone(),
5311 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
5312 allow_write: true,
5313 sessions: self.state.config.graph.sessions,
5314 artifacts: agent::artifacts_dir(&self.state.dir()),
5315 stem: format!("gate-fix-{attempt}"),
5316 };
5317 self.state.event(
5318 "gate",
5319 format!("gate failed; gate-fix round {attempt} of {cap}"),
5320 );
5321 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
5322 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
5323 let cache = self.state.config.cache_dir();
5324 let ctx = WaveCtx {
5325 run: &run_id,
5326 node: "gate-fix",
5327 prompts: &prompts,
5328 cache: cache.as_deref(),
5329 round: None,
5330 };
5331 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
5332 let mut record = GateFixRecord {
5333 agent: seat.agent.clone(),
5334 failed,
5335 notes: String::new(),
5336 committed: false,
5337 error: None,
5338 };
5339 match out {
5340 AgentOutcome::Ok(o) => {
5341 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
5344 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
5345 }
5346 }
5347 AgentOutcome::Dropped(_) => {
5348 record.error = Some("the CLI dropped the stream".to_owned());
5349 }
5350 AgentOutcome::Quota(o) => {
5351 self.state.quota.push(QuotaLoss {
5352 seat: seat.key.clone(),
5353 node: "gate-fix".to_owned(),
5354 at: Timestamp::now(),
5355 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5356 });
5357 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5358 }
5359 AgentOutcome::Failed(e) => record.error = Some(e),
5360 }
5361 self.state.seats.insert(seat.key.clone(), seat);
5362 if let Ok(r) = git::rescue_commit(
5363 &winner.worktree,
5364 &format!("magi: gate fix {attempt} (uncommitted work)"),
5365 )
5366 .await
5367 {
5368 self.state.note_withheld("gate-fix", &r.withheld);
5369 }
5370 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5371 record.committed = after != before;
5372 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5373 let note = record.error.clone();
5374 self.state.gate_fixes.push(record);
5375 self.state.save()?;
5376 if !changed {
5377 self.state.event(
5378 "gate",
5379 match note {
5380 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5381 None => format!("gate-fix round {attempt}: the tree did not change"),
5382 },
5383 );
5384 return Ok(GateFix::Stop);
5385 }
5386 self.state.event(
5387 "gate",
5388 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5389 );
5390
5391 let commands = self.state.config.verify.e2e.clone();
5392 if !commands.is_empty() {
5393 let shell = self.state.config.shell();
5394 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5395 let cache_dir = self.state.config.cache_dir();
5396 let context = format!("gate-fix round {attempt}");
5397 let (e2e, _) = with_cache_lease(
5398 &mut self.state,
5399 cache_dir.as_deref(),
5400 "e2e",
5401 "e2e",
5402 &winner.worktree,
5403 &after,
5404 timeout,
5405 &context,
5406 |state, budget| {
5407 let shell = shell.clone();
5408 let commands = commands.clone();
5409 let context = context.clone();
5410 let worktree = winner.worktree.clone();
5411 async move {
5412 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5413 .await
5414 }
5415 },
5416 )
5417 .await;
5418 if verify_inconclusive(&e2e) {
5419 return Ok(GateFix::Defer);
5420 }
5421 if e2e.iter().any(|o| !o.ok()) {
5422 self.state.event(
5423 "gate",
5424 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5425 );
5426 return Ok(GateFix::Stop);
5427 }
5428 }
5429 Ok(GateFix::Retry)
5430 }
5431
5432 async fn merge(&mut self) -> Result<()> {
5435 if self
5450 .state
5451 .base_sync
5452 .as_ref()
5453 .is_some_and(|s| s.conflict.is_some())
5454 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5455 != Some(RunStatus::Gating)
5456 || !self.state.gate_status().ok()
5465 {
5466 return Ok(());
5467 }
5468 if self.state.merge.is_some() {
5477 return Ok(());
5478 }
5479 let Some(winner) = self.state.winner().cloned() else {
5480 return Ok(());
5481 };
5482 let repo = self.state.repo.clone();
5483 let base = self.state.base_branch.clone();
5484 let mode = self.state.config.merge.mode;
5485 let style = self.state.config.merge.style;
5486 let pr = pr_message(&self.state, winner.label);
5487 let message = pr.commit_message();
5488
5489 let outcome = match mode {
5490 MergeMode::None => MergeOutcome {
5491 mode,
5492 ok: true,
5493 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5494 empty: false,
5495 },
5496 MergeMode::Pr | MergeMode::Local
5497 if merge_is_empty(&repo, &self.state, &winner.branch, mode).await =>
5498 {
5499 MergeOutcome {
5500 mode,
5501 ok: false,
5502 detail: empty_candidate_detail(&self.state, &base),
5503 empty: true,
5504 }
5505 }
5506 MergeMode::Local => {
5507 let on = git::current_branch(&repo).await?;
5508 if on.as_deref() != Some(base.as_str()) {
5509 MergeOutcome {
5510 mode,
5511 ok: false,
5512 detail: format!(
5513 "{} has {} checked out, not the base branch {base}",
5514 repo.display(),
5515 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5516 ),
5517 empty: false,
5518 }
5519 } else if !git::is_clean(&repo).await? {
5520 MergeOutcome {
5521 mode,
5522 ok: false,
5523 detail: format!("{} is dirty; refusing to merge", repo.display()),
5524 empty: false,
5525 }
5526 } else {
5527 let out = match style {
5528 MergeStyle::Merge => {
5529 git::merge_no_ff(&repo, &winner.branch, &message).await?
5530 }
5531 MergeStyle::Squash => {
5532 git::merge_squash(&repo, &winner.branch, &message).await?
5533 }
5534 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5535 };
5536 MergeOutcome {
5537 mode,
5538 ok: out.ok(),
5539 detail: if out.ok() { out.stdout } else { out.stderr },
5540 empty: false,
5541 }
5542 }
5543 }
5544 MergeMode::Pr => {
5545 let remote = self.state.config.merge.remote.clone();
5546 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5547 if !pushed.ok() {
5548 MergeOutcome {
5549 mode,
5550 ok: false,
5551 detail: pushed.stderr,
5552 empty: false,
5553 }
5554 } else {
5555 let found = land::find_open_pr(&winner.worktree, &winner.branch, &base).await;
5561 let out = match pr_merge_plan(found) {
5562 PrPlan::Create => {
5563 gh_pr_create(
5564 &winner.worktree,
5565 &base,
5566 &winner.branch,
5567 &pr.title,
5568 &pr.body,
5569 )
5570 .await
5571 }
5572 PrPlan::Adopt { url, title } => {
5573 self.state
5574 .event("merge", format!("Pr: adopted open pull request {url}"));
5575 if title != pr.title
5576 && let Err(e) =
5577 land::set_pr_title(&winner.worktree, &url, &pr.title).await
5578 {
5579 tracing::warn!("could not refresh title of {url}: {e:#}");
5580 self.state
5581 .event("merge", format!("Pr: title refresh failed: {e:#}"));
5582 }
5583 Ok(url)
5584 }
5585 PrPlan::Stop(why) => Err(anyhow::anyhow!(why)),
5586 };
5587 match out {
5588 Ok(url) => MergeOutcome {
5589 mode,
5590 ok: true,
5591 detail: url,
5592 empty: false,
5593 },
5594 Err(e) => MergeOutcome {
5595 mode,
5596 ok: false,
5597 detail: e.to_string(),
5598 empty: false,
5599 },
5600 }
5601 }
5602 }
5603 };
5604
5605 self.state.status = match (mode, outcome.ok) {
5606 (MergeMode::None, _) => RunStatus::Ready,
5607 (_, true) => RunStatus::Merged,
5608 (_, false) => RunStatus::Blocked,
5609 };
5610 self.state.event(
5611 "merge",
5612 format!(
5613 "{:?}: {}",
5614 mode,
5615 outcome.detail.lines().next().unwrap_or("")
5616 ),
5617 );
5618 self.state.merge = Some(outcome);
5619 self.state.save()?;
5620
5621 if self.state.config.graph.land
5627 && mode == MergeMode::Pr
5628 && self.state.status == RunStatus::Merged
5629 {
5630 self.run_land().await?;
5631 }
5632 self.settle_questions();
5637 Ok(())
5638 }
5639
5640 async fn run_land(&mut self) -> Result<()> {
5651 let url = self
5652 .state
5653 .merge
5654 .as_ref()
5655 .map(|m| m.detail.clone())
5656 .unwrap_or_default();
5657 let url = url.lines().next().unwrap_or("").trim().to_owned();
5658 if !url.starts_with("http") {
5659 return Ok(());
5660 }
5661 match land::land(&mut self.state, &url).await {
5664 Ok(pr) if self.state.parked => {
5665 let _ = pr;
5669 }
5670 Ok(pr) => {
5671 self.state.status = match pr.state {
5672 land::PrLifecycle::Merged => RunStatus::Merged,
5673 _ => RunStatus::Blocked,
5674 };
5675 if bump::should_release_bump(self.state.status)
5682 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5683 {
5684 self.state
5690 .event("bump", format!("release bump skipped: {e:#}"));
5691 }
5692 self.state.save()?;
5693 }
5694 Err(e) => {
5695 self.state.status = RunStatus::Blocked;
5696 self.state.event("land", format!("gave up: {e}"));
5697 self.state.save()?;
5698 }
5699 }
5700 Ok(())
5701 }
5702
5703 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5707 if let Some(existing) = self.state.seats.get(key)
5708 && existing.agent == agent
5709 {
5710 return existing.clone();
5711 }
5712 let fresh = SeatState::new(key, agent, self.state.seed);
5713 self.state.seats.insert(key.to_owned(), fresh.clone());
5714 fresh
5715 }
5716
5717 fn view(&self, c: &Candidate) -> CandidateView {
5719 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5720 .unwrap_or_default();
5721 let (patch, _) = blind::sanitize_patch(
5722 &format!("candidate {} patch", c.label),
5723 &raw,
5724 &self.state.config.blind,
5725 );
5726 CandidateView {
5727 label: c.label,
5728 branch: c.branch.clone(),
5729 summary: c.summary.clone(),
5730 stat: c.stat.clone(),
5731 patch,
5732 }
5733 }
5734
5735 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5737 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5738 prompt::judge(
5739 "(see above)",
5740 &views,
5741 self.roles.judges.len(),
5742 base_short,
5743 "en",
5744 )
5745 }
5746
5747 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5754 let mut turns = Vec::new();
5755 for j in &self.state.judgements {
5756 if j.ranking.is_empty() {
5757 continue;
5758 }
5759 let reasons = j
5760 .reasons
5761 .iter()
5762 .map(|(k, v)| format!("- {k}: {v}"))
5763 .collect::<Vec<_>>()
5764 .join("\n");
5765 turns.push(Turn {
5766 who: format!("Judge {} (opening ranking)", j.judge),
5767 is_self: j.judge == self_idx + 1,
5768 body: format!(
5769 "Ranked {}{}{reasons}",
5770 j.ranking.iter().collect::<String>(),
5771 if reasons.is_empty() {
5772 ""
5773 } else {
5774 ", because:\n"
5775 }
5776 ),
5777 });
5778 }
5779 for t in self
5780 .state
5781 .deliberation
5782 .iter()
5783 .flat_map(|r| r.turns.iter())
5784 .chain(current)
5785 {
5786 turns.push(Turn {
5787 who: format!("Judge {}", t.judge),
5788 is_self: t.judge == self_idx + 1,
5789 body: t.body.clone(),
5790 });
5791 }
5792 turns
5793 }
5794}
5795
5796fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5798 agent::has_session(spec.kind, seat, sessions)
5799}
5800
5801fn next_untried_implementer<'a>(
5822 roster: &'a [AgentSpec],
5823 start: usize,
5824 tried: &BTreeSet<String>,
5825) -> Option<&'a AgentSpec> {
5826 roster
5827 .get(start + 1..)?
5828 .iter()
5829 .find(|s| !tried.contains(&s.id))
5830}
5831
5832fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5846 commands.iter().any(|c| c.exit_code.is_none())
5847}
5848
5849fn verified_noop_claim(
5862 usable: bool,
5863 commands: &[agent::CommandEvidence],
5864 text: &str,
5865) -> Option<String> {
5866 (usable && !has_unconfirmed_command(commands))
5867 .then(|| verdict::verified_noop(text))
5868 .flatten()
5869}
5870
5871fn short(commit: &str) -> String {
5872 commit.chars().take(7).collect()
5873}
5874
5875fn make_executable(path: &Path) -> Result<()> {
5876 #[cfg(unix)]
5877 {
5878 use std::os::unix::fs::PermissionsExt as _;
5879 let mut perms = std::fs::metadata(path)?.permissions();
5880 perms.set_mode(0o755);
5881 std::fs::set_permissions(path, perms)?;
5882 }
5883 #[cfg(not(unix))]
5884 {
5885 let _ = path;
5886 }
5887 Ok(())
5888}
5889
5890struct WaveCtx<'a> {
5897 run: &'a str,
5900 node: &'a str,
5902 prompts: &'a Prompts,
5903 cache: Option<&'a Path>,
5905 round: Option<usize>,
5908}
5909
5910async fn run_one(
5912 job: SeatJob,
5913 sem: Arc<Semaphore>,
5914 ctx: &WaveCtx<'_>,
5915 state: &mut RunState,
5916 attempt: usize,
5917) -> (SeatState, AgentOutcome) {
5918 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5919 .await
5920 .pop()
5921 .expect("one job in, one result out");
5922 (seat, out)
5923}
5924
5925async fn wave(
5931 jobs: Vec<SeatJob>,
5932 sem: Arc<Semaphore>,
5933 ctx: &WaveCtx<'_>,
5934 state: &mut RunState,
5935 attempt: usize,
5936) -> Vec<(usize, SeatState, AgentOutcome)> {
5937 let WaveCtx {
5938 run,
5939 node,
5940 prompts,
5941 cache,
5942 round,
5943 } = *ctx;
5944 for job in &jobs {
5945 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5946 }
5947 if let Err(e) = state.save() {
5948 tracing::warn!("could not persist in-progress seats: {e:#}");
5953 }
5954 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5970 let wait_started = Instant::now();
5971 let cache_guard = if let Some(cache_dir) = cache {
5972 if jobs_had_a_writer {
5973 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5974 let budget = jobs
5975 .iter()
5976 .map(|j| j.timeout)
5977 .max()
5978 .unwrap_or(Duration::from_secs(60));
5979 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5980 .await
5981 .ok()
5982 } else {
5983 None
5984 }
5985 } else {
5986 None
5987 };
5988 let waited_for_lease = wait_started.elapsed();
5995 let mut set = tokio::task::JoinSet::new();
5996 let overlay = prompts.overlay(node);
5997 for (i, mut job) in jobs.into_iter().enumerate() {
5998 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5999 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
6000 if cache.is_some() {
6001 job.prompt.push('\n');
6002 job.prompt
6003 .push_str(&prompt::build_cache_note(node, job.allow_write));
6004 }
6005 let sem = Arc::clone(&sem);
6006 let run = run.to_owned();
6007 let node = node.to_owned();
6008 let attachments = if node == "implement" {
6011 state.attachments.clone()
6012 } else {
6013 Vec::new()
6014 };
6015 let cache = cache
6026 .filter(|_| job.allow_write && cache_guard.is_some())
6027 .map(Path::to_path_buf);
6028 set.spawn(async move {
6029 let _permit = sem.acquire().await;
6030 let mut seat = job.seat;
6031 let out = agent::invoke(
6032 &job.spec,
6033 &mut seat,
6034 &Invocation {
6035 cwd: &job.cwd,
6036 prompt: &job.prompt,
6037 timeout: job.timeout,
6038 allow_write: job.allow_write,
6039 sessions: job.sessions,
6040 artifacts: &job.artifacts,
6041 stem: &job.stem,
6042 run: &run,
6043 node: &node,
6044 cache_dir: cache.as_deref(),
6045 attachments: &attachments,
6046 },
6047 )
6048 .await;
6049 let out = match out {
6050 Ok(o) if o.usable() => AgentOutcome::Ok(o),
6051 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
6052 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
6060 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
6061 Ok(o) => AgentOutcome::Failed(format!(
6062 "exited with {:?} and no usable output",
6063 o.exit_code
6064 )),
6065 Err(e) => AgentOutcome::Failed(e.to_string()),
6066 };
6067 (i, seat, out)
6068 });
6069 }
6070 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
6071 while let Some(joined) = set.join_next().await {
6072 let (i, seat, out) = match joined {
6073 Ok(v) => v,
6074 Err(e) => {
6078 tracing::error!("agent task panicked: {e}");
6079 continue;
6080 }
6081 };
6082 state.seat_finished(&seat.key);
6083 record_jobs(state, node, round, &seat.key, &out);
6084 if let Err(e) = state.save() {
6085 tracing::warn!("could not persist a seat's completion: {e:#}");
6086 }
6087 if collected.len() <= i {
6088 collected.resize_with(i + 1, || None);
6089 }
6090 collected[i] = Some((i, seat, out));
6091 }
6092 if state
6098 .active
6099 .values()
6100 .any(|a| a.node == node && a.attempt == attempt)
6101 {
6102 state
6103 .active
6104 .retain(|_, a| !(a.node == node && a.attempt == attempt));
6105 if let Err(e) = state.save() {
6106 tracing::warn!("could not persist the end of a wave: {e:#}");
6107 }
6108 }
6109 if let Some(cache_dir) = cache
6116 && jobs_had_a_writer
6117 {
6118 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
6119 }
6120 if let Some(guard) = cache_guard {
6121 guard.release();
6122 }
6123 collected.into_iter().flatten().collect()
6124}
6125
6126fn record_jobs(
6137 state: &mut RunState,
6138 node: &str,
6139 round: Option<usize>,
6140 seat: &str,
6141 out: &AgentOutcome,
6142) {
6143 let commands: &[agent::CommandEvidence] = match out {
6144 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
6145 AgentOutcome::Failed(_) => &[],
6146 };
6147 let checked_at = Timestamp::now();
6148 for c in commands {
6149 state.jobs.push(JobRecord {
6150 node: node.to_owned(),
6151 round,
6152 seat: seat.to_owned(),
6153 id: c.id.clone(),
6154 description: c.description.clone(),
6155 checked_at,
6156 status: match c.exit_code {
6157 Some(0) => JobStatus::Completed,
6158 Some(_) => JobStatus::Failed,
6159 None => JobStatus::Unknown,
6160 },
6161 exit_code: c.exit_code,
6162 result_summary: c.result_summary.clone(),
6163 source: c.source.clone(),
6164 });
6165 }
6166}
6167
6168fn round_is_clean(
6189 blocking: usize,
6190 e2e_ok: bool,
6191 answered: usize,
6192 expected: usize,
6193 quota_missing: usize,
6194 policy: IncompleteReviewPolicy,
6195) -> bool {
6196 if blocking != 0 || !e2e_ok {
6197 return false;
6198 }
6199 if answered == expected || policy == IncompleteReviewPolicy::Warn {
6200 return true;
6201 }
6202 answered > 0 && expected - answered <= quota_missing
6203}
6204
6205fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
6229 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
6230 return Some(RunStatus::Gating);
6231 }
6232 let last = reviews.last()?;
6233 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
6234 if reviews.len() < max_rounds && !stagnant {
6235 return None;
6236 }
6237 if last.incomplete() && last.blocking == 0 {
6238 return Some(RunStatus::Blocked);
6239 }
6240 if last.e2e_status() == E2eStatus::ResourceBlocked {
6241 return None;
6242 }
6243 Some(if last.e2e.iter().all(CommandOutcome::ok) {
6244 RunStatus::Gating
6245 } else {
6246 RunStatus::Blocked
6247 })
6248}
6249
6250fn retry_budget(full: Duration, nudged: bool) -> Duration {
6265 if nudged {
6266 (full / 4).max(Duration::from_secs(120)).min(full)
6267 } else {
6268 full
6269 }
6270}
6271
6272#[allow(clippy::too_many_arguments)]
6285async fn ask_json_wave<T>(
6286 jobs: Vec<SeatJob>,
6287 sem: Arc<Semaphore>,
6288 retries: usize,
6289 ctx: &WaveCtx<'_>,
6290 losses: &mut Vec<QuotaLoss>,
6291 state: &mut RunState,
6292 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
6293) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
6294where
6295 T: serde::de::DeserializeOwned + Send + 'static,
6296{
6297 let n = jobs.len();
6298 let originals: Vec<SeatJob> = jobs;
6299 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
6300 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
6301 let mut attempts_used: Vec<usize> = vec![0; n];
6308 let mut pending: Vec<usize> = (0..n).collect();
6309
6310 for attempt in 0..=retries {
6311 if pending.is_empty() {
6312 break;
6313 }
6314 let mut batch = Vec::with_capacity(pending.len());
6315 for &i in &pending {
6316 let src = &originals[i];
6317 let (prompt, timeout) = if attempt == 0 {
6320 (src.prompt.clone(), src.timeout)
6321 } else {
6322 let why = done[i]
6323 .as_ref()
6324 .and_then(|r| r.as_ref().err().map(ToString::to_string))
6325 .unwrap_or_else(|| "no parsable answer".to_owned());
6326 let nudge = prompt::nudge(&why);
6327 let nudged = has_context(&src.spec, &seats[i], src.sessions);
6328 let prompt = if nudged {
6329 nudge
6330 } else {
6331 format!("{}\n\n---\n\n{}", src.prompt, nudge)
6332 };
6333 (prompt, retry_budget(src.timeout, nudged))
6334 };
6335 batch.push(SeatJob {
6336 spec: src.spec.clone(),
6337 seat: seats[i].clone(),
6338 cwd: src.cwd.clone(),
6339 prompt,
6340 timeout,
6341 allow_write: src.allow_write,
6342 sessions: src.sessions,
6343 artifacts: src.artifacts.clone(),
6344 stem: if attempt == 0 {
6345 src.stem.clone()
6346 } else {
6347 format!("{}-retry{attempt}", src.stem)
6348 },
6349 });
6350 }
6351
6352 if attempt > 0 {
6353 let seats_out: Vec<&str> = pending
6354 .iter()
6355 .map(|&i| originals[i].seat.key.as_str())
6356 .collect();
6357 state.event(
6358 ctx.node,
6359 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
6360 );
6361 }
6362 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
6363 let mut still = Vec::new();
6364 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
6365 seats[i] = seat;
6366 let (parsed, quota) = match out {
6367 AgentOutcome::Ok(o) => (
6368 match verdict::extract_json::<T>(&o.text) {
6369 Ok(v) => match validate(&v) {
6370 Ok(()) => Ok((v, o)),
6371 Err(e) => Err(e),
6372 },
6373 Err(e) => Err(e),
6374 },
6375 false,
6376 ),
6377 AgentOutcome::Quota(o) => {
6378 losses.push(QuotaLoss {
6379 seat: originals[i].seat.key.clone(),
6380 node: ctx.node.to_owned(),
6381 at: Timestamp::now(),
6382 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
6383 });
6384 (
6385 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
6386 true,
6387 )
6388 }
6389 AgentOutcome::Dropped(o) => {
6394 let why = o
6395 .dropped
6396 .as_ref()
6397 .map(|d| d.why.as_str())
6398 .unwrap_or("the CLI ended the stream without delivering its answer");
6399 (
6400 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
6401 false,
6402 )
6403 }
6404 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
6405 };
6406 let failed = parsed.is_err();
6407 done[i] = Some(parsed);
6408 attempts_used[i] = attempt;
6409 if failed && !quota {
6412 still.push(i);
6413 }
6414 }
6415 pending = still;
6416 }
6417
6418 seats
6419 .into_iter()
6420 .zip(done)
6421 .zip(attempts_used)
6422 .map(|((seat, res), attempts)| {
6423 (
6424 seat,
6425 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6426 attempts,
6427 )
6428 })
6429 .collect()
6430}
6431
6432async fn acquire_cache_lease(
6445 state: &mut RunState,
6446 cache_dir: &Path,
6447 owner: &crate::cache::Owner,
6448 budget: Duration,
6449 context: &str,
6450) -> Result<crate::cache::Guard> {
6451 let home = crate::run::home();
6452 let started = Instant::now();
6453 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6454 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6455 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6456 Err(e) => {
6457 state.event(
6458 "verify",
6459 format!("{context}: could not check the shared build cache: {e:#}"),
6460 );
6461 if let Err(e2) = state.save() {
6462 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6463 }
6464 return Err(e);
6465 }
6466 };
6467 state.event(
6468 "verify",
6469 format!(
6470 "{context}: waiting for the shared build cache at {} ({})",
6471 cache_dir.display(),
6472 busy.describe()
6473 ),
6474 );
6475 if let Err(e) = state.save() {
6476 tracing::warn!("could not persist a cache wait: {e:#}");
6477 }
6478 let remaining = budget.saturating_sub(started.elapsed());
6479 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6480 Ok(g) => Ok(g),
6481 Err(e) => {
6482 state.event("verify", format!("{context}: {e:#}"));
6483 if let Err(e2) = state.save() {
6484 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6485 }
6486 Err(e)
6487 }
6488 }
6489}
6490
6491#[allow(clippy::too_many_arguments)]
6512async fn with_cache_lease<'s, F, Fut>(
6513 state: &'s mut RunState,
6514 cache_dir: Option<&Path>,
6515 node: &str,
6516 seat: &str,
6517 worktree: &Path,
6518 head: &str,
6519 budget: Duration,
6520 context: &str,
6521 body: F,
6522) -> (Vec<CommandOutcome>, bool)
6523where
6524 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6525 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6526{
6527 let Some(cache_dir) = cache_dir else {
6528 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6529 return (outcomes, retried);
6530 };
6531 let home = crate::run::home();
6532 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6533 let started = Instant::now();
6534 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6535 Ok(g) => g,
6536 Err(e) => {
6537 return (
6538 vec![CommandOutcome {
6539 command: "(waiting for the shared build cache)".to_owned(),
6540 code: None,
6541 output_tail: e.to_string(),
6542 duration_ms: started.elapsed().as_millis() as u64,
6543 resource_blocked: true,
6544 }],
6545 false,
6546 );
6547 }
6548 };
6549 let identity = crate::cache::Identity::new(worktree, head);
6550 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6551 state.event(
6559 "verify",
6560 format!(
6561 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6562 worktree.display(),
6563 short(head)
6564 ),
6565 );
6566 guard.release();
6567 return (
6568 vec![CommandOutcome {
6569 command: "(confirming the shared build cache is fresh)".to_owned(),
6570 code: None,
6571 output_tail: e.to_string(),
6572 duration_ms: started.elapsed().as_millis() as u64,
6573 resource_blocked: true,
6574 }],
6575 false,
6576 );
6577 }
6578 let remaining = budget.saturating_sub(started.elapsed());
6579 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6580 if !timed_out_pids.is_empty() {
6585 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6586 }
6587 guard.release();
6588 (outcomes, retried)
6589}
6590
6591async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6603 wait_for_pids_with(
6604 pids,
6605 crate::proc::pid_alive,
6606 LEASE_RELEASE_POLL,
6607 LEASE_RELEASE_MAX_WAIT,
6608 )
6609 .await;
6610}
6611
6612async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6618 pids: &[u32],
6619 alive: F,
6620 poll: Duration,
6621 max_wait: Duration,
6622) {
6623 let deadline = Instant::now() + max_wait;
6624 loop {
6625 if pids.iter().all(|&pid| !alive(pid)) {
6626 return;
6627 }
6628 if Instant::now() >= deadline {
6629 return;
6630 }
6631 tokio::time::sleep(poll).await;
6632 }
6633}
6634
6635fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6642 outcomes.iter().any(|o| o.resource_blocked)
6643}
6644
6645enum GateFix {
6647 Retry,
6649 Stop,
6652 Defer,
6655}
6656
6657fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6665 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6666 red.peek().is_some()
6667 && red.all(|o| {
6668 !o.resource_blocked
6669 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6670 && !o.output_tail.trim().is_empty()
6671 })
6672}
6673
6674fn e2e_outcome_label(o: &CommandOutcome) -> String {
6678 if o.ok() {
6679 return "pass".to_owned();
6680 }
6681 let reason = if o.build_failed() {
6682 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6683 } else {
6684 format!("FAIL ({:?})", o.code)
6685 };
6686 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6687}
6688
6689async fn run_e2e_with_retry(
6697 state: &mut RunState,
6698 shell: &[String],
6699 commands: &[String],
6700 worktree: &Path,
6701 timeout: Duration,
6702 context: &str,
6703) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6704 let (mut e2e, mut timed_out_pids) = run_commands(
6705 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6706 )
6707 .await;
6708 for o in &e2e {
6709 state.event(
6710 "verify",
6711 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6712 );
6713 }
6714 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6718 if verify_retried {
6719 state.event(
6720 "verify",
6721 format!(
6722 "{context}: verify could not build/link, not a test result — retrying once \
6723 before concluding"
6724 ),
6725 );
6726 let retried = run_commands(
6727 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6728 )
6729 .await;
6730 e2e = retried.0;
6731 timed_out_pids.extend(retried.1);
6734 for o in &e2e {
6735 state.event(
6736 "verify",
6737 format!(
6738 "{context}: retry `{}` -> {}",
6739 o.command,
6740 e2e_outcome_label(o)
6741 ),
6742 );
6743 }
6744 }
6745 (e2e, verify_retried, timed_out_pids)
6746}
6747
6748#[allow(clippy::too_many_arguments)]
6764async fn run_commands(
6765 state: &mut RunState,
6766 node: &str,
6767 task: &str,
6768 attempt: usize,
6769 shell: &[String],
6770 commands: &[String],
6771 cwd: &Path,
6772 timeout: Duration,
6773) -> (Vec<CommandOutcome>, Vec<u32>) {
6774 if commands.is_empty() {
6775 return (Vec::new(), Vec::new());
6780 }
6781 let mut out = Vec::new();
6782 let mut timed_out_pids = Vec::new();
6783 let total = commands.len();
6784 for (idx, command) in commands.iter().enumerate() {
6785 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6786 if let Err(e) = state.save() {
6787 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6788 }
6789 let started = Instant::now();
6790 let mut cmd = tokio::process::Command::new(&shell[0]);
6791 cmd.quiet();
6792 cmd.args(&shell[1..])
6793 .arg(command)
6794 .current_dir(cwd)
6795 .stdin(std::process::Stdio::null())
6796 .stdout(std::process::Stdio::piped())
6797 .stderr(std::process::Stdio::piped())
6798 .kill_on_drop(true);
6799 let spawned = cmd.spawn();
6800 let (code, body) = match spawned {
6801 Ok(child) => {
6802 let pid = child.id();
6807 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6808 Ok(Ok(o)) => {
6809 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6810 body.push_str(&String::from_utf8_lossy(&o.stderr));
6811 (o.status.code(), body)
6812 }
6813 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6814 Err(_) => {
6815 if let Some(pid) = pid {
6816 timed_out_pids.push(pid);
6817 }
6818 (None, format!("timed out after {}s", timeout.as_secs()))
6819 }
6820 }
6821 }
6822 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6823 };
6824 out.push(CommandOutcome {
6825 command: command.clone(),
6826 code,
6827 output_tail: tail(&body, OUTPUT_TAIL),
6828 duration_ms: started.elapsed().as_millis() as u64,
6829 resource_blocked: false,
6830 });
6831 }
6832 state.task_finished(task);
6833 if let Err(e) = state.save() {
6834 tracing::warn!("could not persist the end of {task}: {e:#}");
6835 }
6836 (out, timed_out_pids)
6837}
6838
6839fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6851 let repo = repo.display();
6852 match style {
6853 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6854 MergeStyle::Squash => {
6855 let subject = message
6858 .lines()
6859 .next()
6860 .unwrap_or(branch)
6861 .replace(['\\', '"', '$', '`'], "");
6862 format!(
6863 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6864 )
6865 }
6866 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6867 }
6868}
6869
6870const PR_TITLE_MAX: usize = 240;
6881
6882struct PrMessage {
6887 title: String,
6888 body: String,
6889}
6890
6891impl PrMessage {
6892 fn commit_message(&self) -> String {
6896 format!("{}\n\n{}", self.title, self.body)
6897 }
6898}
6899
6900fn title_marker(line: &str) -> Option<&str> {
6902 let line = line.trim();
6903 let head = line.get(..6)?;
6904 head.eq_ignore_ascii_case("title:")
6905 .then(|| line[6..].trim())
6906}
6907
6908fn summary_title(summary: &str) -> Option<String> {
6913 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6914 let raw = title_marker(first)?;
6915 if raw.is_empty() {
6916 return None;
6917 }
6918 let title = queue::title_from(raw, PR_TITLE_MAX);
6919 let lower = title.to_ascii_lowercase();
6920 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6921 return None;
6922 }
6923 Some(title)
6924}
6925
6926fn summary_without_title(summary: &str) -> String {
6929 let mut lines = summary.trim().lines().peekable();
6930 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6931 lines.next();
6932 }
6933 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6934}
6935
6936fn pr_message(state: &RunState, winner: char) -> PrMessage {
6950 let summary = state
6951 .candidates
6952 .iter()
6953 .find(|c| c.label == winner)
6954 .map(|c| c.summary.as_str())
6955 .unwrap_or_default();
6956 let title = summary_title(summary).unwrap_or_else(|| {
6959 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6960 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6961 t
6962 } else {
6963 format!(
6964 "chore: land candidate {} of run {}",
6965 winner.to_ascii_uppercase(),
6966 state.id
6967 )
6968 }
6969 });
6970
6971 let mut body = String::new();
6972 let what = summary_without_title(summary);
6973 if !what.is_empty() {
6974 body.push_str("## Summary\n\n");
6975 body.push_str(&what);
6976 body.push_str("\n\n");
6977 }
6978
6979 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6980 if let Some(fix) = fix
6981 && !fix.notes.trim().is_empty()
6982 {
6983 body.push_str("## Review fixes\n\n");
6984 body.push_str(fix.notes.trim());
6985 body.push_str("\n\n");
6986 }
6987
6988 let open = state.open_findings();
6989 if !open.is_empty() {
6990 body.push_str("## Open review findings\n\n");
6991 for f in &open {
6992 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6993 }
6994 body.push('\n');
6995 }
6996
6997 if let Some(fix) = fix
6998 && !fix.rejected.is_empty()
6999 {
7000 body.push_str("## Declined by the fixer\n\n");
7001 for r in &fix.rejected {
7002 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
7003 }
7004 body.push('\n');
7005 }
7006
7007 let task = state.instruction.trim();
7008 let task = if task.is_empty() {
7009 "(empty task)"
7010 } else {
7011 task
7012 };
7013 body.push_str(&format!(
7014 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
7015 task.replace("</details>", "</details>")
7016 ));
7017
7018 body.push_str(&format!(
7019 "\n---\nmagi:run/{} magi:candidate-{}\n",
7020 state.id,
7021 winner.to_ascii_lowercase()
7022 ));
7023
7024 let id = crate::scrub::Identity::current();
7027 PrMessage {
7028 title: crate::scrub::scrub(&title, &id),
7029 body: crate::scrub::scrub(&body, &id),
7030 }
7031}
7032
7033fn seeded_instruction(state: &RunState) -> String {
7037 match refs::describe(&state.seeds) {
7038 Some(facts) => format!(
7039 "{}\n\n# Existing work the task refers to\n\n{facts}\n\n\
7040 Candidates start from the unmerged branch named above, when there \
7041 is one, and carry any unmerged commit named by sha as a \
7042 cherry-pick. Check that this is what the task meant before \
7043 building on it.",
7044 state.instruction
7045 ),
7046 None => state.instruction.clone(),
7047 }
7048}
7049
7050async fn merge_is_empty(repo: &Path, state: &RunState, branch: &str, mode: MergeMode) -> bool {
7054 let base = &state.base_branch;
7055 let mut against = base.clone();
7056 if mode == MergeMode::Pr {
7057 let remote = &state.config.merge.remote;
7058 let tracking = format!("{remote}/{base}");
7059 let fetched = git::fetch(repo, remote, base).await;
7060 if fetched.is_ok_and(|o| o.ok()) && git::rev_exists(repo, &tracking).await {
7061 against = tracking;
7062 }
7063 }
7064 matches!(git::commits_ahead(repo, &against, branch).await, Ok(0))
7065}
7066
7067fn empty_candidate_detail(state: &RunState, base: &str) -> String {
7070 let mut detail = format!(
7071 "empty candidate: the winning branch has 0 commits ahead of {base}, so there is \
7072 nothing to open a pull request for"
7073 );
7074 match refs::describe(&state.seeds) {
7075 Some(facts) => detail.push_str(&format!("\nReferences in the task:\n{facts}")),
7076 None => detail.push_str(
7077 "\nThe task names no existing branch or commit; if it means to land work \
7078 that lives elsewhere, name the branch (magi/<run>/<label>) or the sha.",
7079 ),
7080 }
7081 detail
7082}
7083
7084#[derive(Debug, PartialEq, Eq)]
7087enum PrPlan {
7088 Create,
7089 Adopt { url: String, title: String },
7090 Stop(String),
7091}
7092
7093fn pr_merge_plan(found: Result<land::OpenPr>) -> PrPlan {
7096 match found {
7097 Ok(land::OpenPr::None) => PrPlan::Create,
7098 Ok(land::OpenPr::One { url, title }) => PrPlan::Adopt { url, title },
7099 Ok(land::OpenPr::Many(urls)) => PrPlan::Stop(format!(
7100 "several open pull requests exist for this branch, not picking one: {}",
7101 urls.join(" ")
7102 )),
7103 Err(e) => PrPlan::Stop(format!("could not look up open pull requests: {e:#}")),
7104 }
7105}
7106
7107async fn gh_pr_create(
7109 cwd: &Path,
7110 base: &str,
7111 head: &str,
7112 title: &str,
7113 body: &str,
7114) -> Result<String> {
7115 let out = tokio::process::Command::new("gh")
7116 .args([
7117 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
7118 ])
7119 .current_dir(cwd)
7120 .quiet()
7121 .stdin(std::process::Stdio::null())
7122 .output()
7123 .await
7124 .context("spawn gh")?;
7125 if out.status.success() {
7126 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
7127 } else {
7128 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
7129 }
7130}
7131
7132pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
7141 let repo = state.repo.clone();
7142 let root = state.worktree_root();
7143 let winner = state.tally.as_ref().map(|t| t.winner);
7144 let mut removed = Vec::new();
7145
7146 for i in 0..state.candidates.len() {
7147 let c = state.candidates[i].clone();
7148 let is_winner = Some(c.label) == winner;
7149 if is_winner && !drop_winner {
7150 continue;
7151 }
7152 if c.worktree.exists() {
7153 git::worktree_remove(&repo, &c.worktree).await.ok();
7154 removed.push(c.worktree.to_string_lossy().into_owned());
7155 }
7156 let handed_over = state.released_branches.contains(&c.branch);
7159 if !handed_over && git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
7160 git::branch_delete(&repo, &c.branch).await.ok();
7161 removed.push(c.branch.clone());
7162 }
7163 state.candidates[i].folded = true;
7164 }
7165
7166 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
7167 let path = name.path();
7168 let keep = !drop_winner
7169 && winner.is_some_and(|w| {
7170 path.file_name()
7171 .is_some_and(|n| n == format!("cand-{w}").as_str())
7172 });
7173 if keep {
7174 continue;
7175 }
7176 git::worktree_remove(&repo, &path).await.ok();
7177 removed.push(path.to_string_lossy().into_owned());
7178 }
7179
7180 remove_if_empty(&root);
7189
7190 if state.enabled_worktree_config && drop_winner {
7191 git::release_worktree_config(&repo).await.ok();
7195 state.enabled_worktree_config = false;
7196 }
7197 state.save_under(home)?;
7198 Ok(removed)
7199}
7200
7201fn remove_if_empty(dir: &Path) {
7212 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
7213 std::fs::remove_dir(dir).ok();
7214 }
7215}
7216
7217pub fn worst_open(state: &RunState) -> Option<Severity> {
7219 state
7220 .reviews
7221 .last()?
7222 .reviews
7223 .iter()
7224 .flat_map(|r| r.findings.iter())
7225 .map(|f| f.severity)
7226 .max()
7227}
7228
7229#[cfg(test)]
7230mod tests {
7231 #[test]
7232 fn pr_merge_plan_creates_adopts_or_stops() {
7233 assert_eq!(pr_merge_plan(Ok(land::OpenPr::None)), PrPlan::Create);
7234 assert_eq!(
7235 pr_merge_plan(Ok(land::OpenPr::One {
7236 url: "u".into(),
7237 title: "t".into()
7238 })),
7239 PrPlan::Adopt {
7240 url: "u".into(),
7241 title: "t".into()
7242 }
7243 );
7244 let PrPlan::Stop(many) =
7245 pr_merge_plan(Ok(land::OpenPr::Many(vec!["a".into(), "b".into()])))
7246 else {
7247 panic!("many must stop");
7248 };
7249 assert!(many.contains('a') && many.contains('b'));
7250 let PrPlan::Stop(err) = pr_merge_plan(Err(anyhow::anyhow!("bad token"))) else {
7251 panic!("a failed lookup must stop");
7252 };
7253 assert!(err.contains("bad token"));
7254 }
7255
7256 use super::*;
7257 use crate::run::GateStatus;
7258 use std::collections::BTreeMap;
7259 use std::time::Duration;
7260
7261 fn conductor() -> AgentSpec {
7262 AgentSpec {
7263 id: "conductor".to_owned(),
7264 kind: crate::config::AgentKind::Command,
7265 model: None,
7266 command: vec!["true".to_owned()],
7267 extra_args: Vec::new(),
7268 env: BTreeMap::new(),
7269 prompt_delivery: None,
7270 }
7271 }
7272
7273 fn spec(id: &str) -> AgentSpec {
7274 AgentSpec {
7275 id: id.to_owned(),
7276 kind: crate::config::AgentKind::Command,
7277 model: None,
7278 command: vec!["true".to_owned()],
7279 extra_args: Vec::new(),
7280 env: BTreeMap::new(),
7281 prompt_delivery: None,
7282 }
7283 }
7284
7285 #[test]
7292 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
7293 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7294 let tried = BTreeSet::from(["beta".to_owned()]);
7295 let next = next_untried_implementer(&roster, 1, &tried);
7298 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
7299 }
7300
7301 #[test]
7302 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
7303 let roster = vec![spec("alpha"), spec("beta")];
7304 let tried = BTreeSet::from(["beta".to_owned()]);
7305 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7309 }
7310
7311 #[test]
7312 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
7313 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7314 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
7315 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7319 }
7320
7321 #[test]
7322 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
7323 let roster = vec![spec("a"), spec("a"), spec("b")];
7324 let tried = BTreeSet::from(["a".to_owned()]);
7325 let next = next_untried_implementer(&roster, 0, &tried);
7326 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
7327 }
7328
7329 #[test]
7330 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
7331 let roster = vec![spec("a"), spec("b")];
7332 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
7333 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
7334 }
7335
7336 #[test]
7337 fn remove_if_empty_only_ever_takes_a_bare_directory() {
7338 let dir = tempfile::tempdir().unwrap();
7339 let bay = dir.path().join("ffff");
7340
7341 remove_if_empty(&bay);
7343 assert!(!bay.exists());
7344
7345 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
7348 remove_if_empty(&bay);
7349 assert!(bay.exists(), "non-empty directory must survive");
7350
7351 std::fs::remove_dir(bay.join("cand-A")).unwrap();
7353 remove_if_empty(&bay);
7354 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
7355 }
7356
7357 #[test]
7366 fn a_full_panel_that_found_nothing_is_clean() {
7367 assert!(round_is_clean(
7368 0,
7369 true,
7370 2,
7371 2,
7372 0,
7373 IncompleteReviewPolicy::Block
7374 ));
7375 }
7376
7377 #[test]
7378 fn a_missing_seat_is_never_clean_under_the_default_policy() {
7379 assert!(!round_is_clean(
7380 0,
7381 true,
7382 1,
7383 2,
7384 0,
7385 IncompleteReviewPolicy::Block
7386 ));
7387 }
7388
7389 #[test]
7390 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
7391 assert!(!round_is_clean(
7392 1,
7393 true,
7394 1,
7395 2,
7396 0,
7397 IncompleteReviewPolicy::Warn
7398 ));
7399 }
7400
7401 #[test]
7402 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
7403 assert!(round_is_clean(
7404 0,
7405 true,
7406 1,
7407 2,
7408 0,
7409 IncompleteReviewPolicy::Warn
7410 ));
7411 }
7412
7413 #[test]
7414 fn a_full_panel_with_an_open_finding_is_not_clean() {
7415 assert!(!round_is_clean(
7416 1,
7417 true,
7418 2,
7419 2,
7420 0,
7421 IncompleteReviewPolicy::Block
7422 ));
7423 }
7424
7425 #[test]
7426 fn a_full_panel_with_a_red_e2e_is_not_clean() {
7427 assert!(!round_is_clean(
7428 0,
7429 false,
7430 2,
7431 2,
7432 0,
7433 IncompleteReviewPolicy::Block
7434 ));
7435 }
7436
7437 #[test]
7444 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
7445 assert!(round_is_clean(
7448 0,
7449 true,
7450 1,
7451 2,
7452 1,
7453 IncompleteReviewPolicy::Block
7454 ));
7455 }
7456
7457 #[test]
7458 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
7459 assert!(!round_is_clean(
7462 0,
7463 true,
7464 1,
7465 2,
7466 0,
7467 IncompleteReviewPolicy::Block
7468 ));
7469 }
7470
7471 #[test]
7472 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
7473 assert!(!round_is_clean(
7474 1,
7475 true,
7476 1,
7477 2,
7478 1,
7479 IncompleteReviewPolicy::Block
7480 ));
7481 assert!(!round_is_clean(
7482 0,
7483 false,
7484 1,
7485 2,
7486 1,
7487 IncompleteReviewPolicy::Block
7488 ));
7489 }
7490
7491 #[test]
7492 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
7493 assert!(!round_is_clean(
7497 0,
7498 true,
7499 0,
7500 2,
7501 2,
7502 IncompleteReviewPolicy::Block
7503 ));
7504 }
7505
7506 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
7507 CommandOutcome {
7508 command: "test".to_owned(),
7509 code,
7510 output_tail: String::new(),
7511 duration_ms: 0,
7512 resource_blocked,
7513 }
7514 }
7515
7516 #[test]
7517 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7518 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7519 assert!(
7520 !verify_inconclusive(&[outcome(Some(1), false)]),
7521 "an ordinary failure is still evidence about the patch"
7522 );
7523 assert!(verify_inconclusive(&[outcome(None, true)]));
7524 assert!(
7525 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7526 "one inconclusive outcome taints the whole batch"
7527 );
7528 assert!(!verify_inconclusive(&[]));
7529 }
7530
7531 #[tokio::test]
7532 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7533 let calls = std::sync::atomic::AtomicUsize::new(0);
7537 let started = Instant::now();
7538 wait_for_pids_with(
7539 &[123],
7540 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7541 Duration::from_millis(5),
7542 Duration::from_secs(5),
7543 )
7544 .await;
7545 assert!(
7546 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7547 "must keep checking rather than deciding on the first answer"
7548 );
7549 assert!(
7550 started.elapsed() < Duration::from_secs(1),
7551 "must return the moment it is confirmed dead, not wait out the ceiling"
7552 );
7553 }
7554
7555 #[tokio::test]
7556 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7557 let started = Instant::now();
7558 wait_for_pids_with(
7559 &[123],
7560 |_| true, Duration::from_millis(5),
7562 Duration::from_millis(30),
7563 )
7564 .await;
7565 let elapsed = started.elapsed();
7566 assert!(
7567 elapsed >= Duration::from_millis(30),
7568 "must not give up before its own ceiling: {elapsed:?}"
7569 );
7570 assert!(
7571 elapsed < Duration::from_secs(1),
7572 "must not wait past its own ceiling either: {elapsed:?}"
7573 );
7574 }
7575
7576 #[tokio::test]
7577 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7578 let started = Instant::now();
7579 wait_for_pids_with(
7580 &[],
7581 |_| true,
7582 Duration::from_secs(5),
7583 Duration::from_secs(5),
7584 )
7585 .await;
7586 assert!(
7587 started.elapsed() < Duration::from_millis(200),
7588 "an empty pid list has nothing to confirm"
7589 );
7590 }
7591
7592 fn review_round(
7598 clean: bool,
7599 blocking: usize,
7600 answered: usize,
7601 expected: usize,
7602 progressed: bool,
7603 e2e_ok: bool,
7604 ) -> ReviewRound {
7605 ReviewRound {
7606 round: 1,
7607 head: "h".to_owned(),
7608 verified_head: None,
7609 verified_at: None,
7610 reviews: Vec::new(),
7611 e2e: vec![CommandOutcome {
7612 command: "test".to_owned(),
7613 code: Some(if e2e_ok { 0 } else { 1 }),
7614 output_tail: String::new(),
7615 duration_ms: 0,
7616 resource_blocked: false,
7617 }],
7618 verify_retried: false,
7619 e2e_deferred: false,
7620 e2e_defer_reason: None,
7621 fix: None,
7622 blocking,
7623 answered,
7624 expected,
7625 clean,
7626 progressed,
7627 vote_split: false,
7628 reconsideration: Vec::new(),
7629 verdict: None,
7630 }
7631 }
7632
7633 #[test]
7634 fn review_conclusion_is_none_when_nothing_has_run() {
7635 assert_eq!(review_conclusion(&[], 3), None);
7636 }
7637
7638 #[test]
7639 fn review_conclusion_is_none_while_rounds_remain() {
7640 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7641 assert_eq!(review_conclusion(&rounds, 3), None);
7642 }
7643
7644 #[test]
7645 fn review_conclusion_is_gating_once_a_round_is_clean() {
7646 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7647 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7648 }
7649
7650 #[test]
7651 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7652 let rounds = vec![
7653 review_round(false, 1, 2, 2, true, true),
7654 review_round(false, 1, 2, 2, true, true),
7655 ];
7656 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7657 }
7658
7659 #[test]
7660 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7661 let rounds = vec![
7662 review_round(false, 1, 2, 2, true, true),
7663 review_round(false, 1, 2, 2, true, false),
7664 ];
7665 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7666 }
7667
7668 #[test]
7669 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7670 let mut blocked = review_round(false, 1, 2, 2, true, false);
7677 blocked.e2e[0].resource_blocked = true;
7678 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7679 assert_eq!(review_conclusion(&rounds, 2), None);
7680 }
7681
7682 #[test]
7683 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7684 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7686 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7687 }
7688
7689 #[test]
7690 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7691 let rounds = vec![
7692 review_round(false, 1, 2, 2, false, true),
7693 review_round(false, 1, 2, 2, false, true),
7694 ];
7695 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7696 }
7697
7698 fn secs(n: u64) -> Duration {
7699 Duration::from_secs(n)
7700 }
7701
7702 fn init_repo(dir: &Path) {
7705 let run = |args: &[&str]| {
7706 let out = std::process::Command::new("git")
7707 .args(args)
7708 .current_dir(dir)
7709 .quiet()
7710 .output()
7711 .expect("spawn git");
7712 assert!(
7713 out.status.success(),
7714 "git {args:?} failed: {}",
7715 String::from_utf8_lossy(&out.stderr)
7716 );
7717 };
7718 run(&["init", "-b", "main"]);
7719 run(&["config", "user.name", "magi test"]);
7720 run(&["config", "user.email", "magi@example.com"]);
7721 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7722 run(&["add", "-A"]);
7723 run(&["commit", "-m", "init"]);
7724 }
7725
7726 fn ask_test_home() {
7734 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7735 }
7736
7737 fn runner_at(status: RunStatus) -> Runner {
7740 let mut state = RunState::new(
7741 PathBuf::from("/nonexistent/repo"),
7742 "main".to_owned(),
7743 "deadbeef".to_owned(),
7744 "task".to_owned(),
7745 Config::default(),
7746 );
7747 state.status = status;
7748 Runner {
7749 state,
7750 roles: ResolvedRoles {
7751 implementers: Vec::new(),
7752 judges: Vec::new(),
7753 reviewers: Vec::new(),
7754 fixer: None,
7755 conductor: conductor(),
7756 implementer_roster: Vec::new(),
7757 },
7758 sem: Arc::new(Semaphore::new(1)),
7759 pause: Pause::new(),
7760 interrupt: Pause::new(),
7761 }
7762 }
7763
7764 #[test]
7768 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7769 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7770 let mut runner = runner_at(RunStatus::Implementing);
7771 let interrupt = Pause::new();
7772 runner.watch_interrupt(interrupt.clone());
7773
7774 interrupt.park_because("task a1b2 asked to run first");
7775
7776 assert!(runner.park_here().expect("park_here"));
7777 assert!(runner.state.parked);
7778 let last = runner.state.events.last().expect("a park event");
7779 assert_eq!(last.node, "park");
7780 assert!(
7781 last.message.contains("task a1b2 asked to run first"),
7782 "expected the interrupt reason in {:?}",
7783 last.message
7784 );
7785 }
7786
7787 #[test]
7795 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7796 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7797 let mut runner = runner_at(RunStatus::Implementing);
7798 let shutdown = Pause::new();
7799 runner.on_pause(shutdown.clone());
7800 let interrupt = Pause::new();
7801 runner.watch_interrupt(interrupt.clone());
7802
7803 assert!(!runner.park_here().expect("park_here"));
7805 assert!(!runner.state.parked);
7806
7807 interrupt.park_because("test");
7809 assert!(!shutdown.parked());
7810 assert!(runner.park_here().expect("park_here"));
7811 }
7812
7813 #[tokio::test]
7827 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7828 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7829 let mut runner = runner_at(RunStatus::Implementing);
7830 let interrupt = Pause::new();
7831 runner.watch_interrupt(interrupt.clone());
7832
7833 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7834 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7835
7836 let node = async move {
7840 started_tx.send(()).expect("send started");
7841 finish_rx.await.expect("recv finish");
7842 "node finished"
7843 };
7844
7845 let interrupter = async move {
7846 started_rx.await.expect("recv started");
7847 interrupt.park_because("higher-priority task waiting");
7849 tokio::task::yield_now().await;
7853 finish_tx.send(()).expect("send finish");
7854 };
7855
7856 let (node_result, ()) = tokio::join!(node, interrupter);
7857 assert_eq!(
7858 node_result, "node finished",
7859 "the in-flight call ran to completion"
7860 );
7861
7862 assert!(runner.park_here().expect("park_here"));
7865 assert!(runner.state.parked);
7866 }
7867
7868 #[test]
7874 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7875 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7876 let mut runner = runner_at(RunStatus::Judging);
7877 runner.state.config.agents = vec![conductor()];
7881 runner.state.candidates = vec![Candidate {
7882 index: 0,
7883 label: 'A',
7884 agent: "alpha".to_owned(),
7885 branch: "magi/x/A".to_owned(),
7886 worktree: PathBuf::from("/nonexistent/worktree"),
7887 summary: "did the thing".to_owned(),
7888 stat: "1 file changed".to_owned(),
7889 files: 1,
7890 commits: 1,
7891 empty: false,
7892 failed: None,
7893 verified_noop: None,
7894 duration_ms: 1234,
7895 folded: false,
7896 }];
7897 let run_id = runner.state.id.clone();
7898
7899 let interrupt = Pause::new();
7900 runner.watch_interrupt(interrupt.clone());
7901 interrupt.park_because("task c3d4 asked to run first");
7902 assert!(runner.park_here().expect("park_here"));
7903
7904 let resumed = Runner::resume(&run_id).expect("resume");
7905 assert_eq!(resumed.state.candidates.len(), 1);
7906 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7907 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7908 assert_eq!(resumed.state.status, runner.state.status);
7909 assert!(
7910 resumed.state.parked,
7911 "still parked until `execute` actually walks the graph again"
7912 );
7913 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7914 }
7915
7916 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7918 let mut q = ask::Question::new(
7919 run.to_owned(),
7920 "implement".to_owned(),
7921 "impl-A".to_owned(),
7922 "Which storage backend should the cache use?".to_owned(),
7923 String::new(),
7924 vec!["SQLite".to_owned(), "Redis".to_owned()],
7925 );
7926 store.put(&mut q).unwrap();
7927 q
7928 }
7929
7930 #[test]
7931 fn a_failed_runs_open_question_is_abandoned() {
7932 ask_test_home();
7933 let store = ask::Questions::open();
7934 let mut runner = runner_at(RunStatus::Failed);
7935 let run = runner.state.id.clone();
7936 let q = ask_open_question(&store, &run);
7937
7938 runner.settle_questions();
7939
7940 let back = store.get(&q.id).unwrap();
7941 assert!(
7942 !back.status.open(),
7943 "the seat that asked died with the run; nobody is left to read an answer"
7944 );
7945 assert!(
7946 back.detail.contains(&run) && back.detail.contains("failed"),
7947 "the reason names what the run became, not just that it is gone: {}",
7948 back.detail
7949 );
7950 }
7951
7952 #[test]
7953 fn a_merged_runs_open_question_is_abandoned_too() {
7954 ask_test_home();
7955 let store = ask::Questions::open();
7956 for status in [RunStatus::Merged, RunStatus::Ready] {
7959 let mut runner = runner_at(status);
7960 let run = runner.state.id.clone();
7961 let q = ask_open_question(&store, &run);
7962
7963 runner.settle_questions();
7964
7965 let back = store.get(&q.id).unwrap();
7966 assert!(
7967 !back.status.open(),
7968 "{status:?} run's question must not outlive the run"
7969 );
7970 }
7971 }
7972
7973 #[test]
7974 fn a_still_resumable_runs_open_question_is_left_alone() {
7975 ask_test_home();
7976 let store = ask::Questions::open();
7977 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7983 let mut runner = runner_at(status);
7984 let run = runner.state.id.clone();
7985 let q = ask_open_question(&store, &run);
7986
7987 runner.settle_questions();
7988
7989 let back = store.get(&q.id).unwrap();
7990 assert!(
7991 back.status.open(),
7992 "{status:?} is still alive; the question must still be waiting"
7993 );
7994 }
7995 }
7996
7997 #[test]
7998 fn settle_questions_never_touches_an_already_answered_question() {
7999 ask_test_home();
8000 let store = ask::Questions::open();
8001 let mut runner = runner_at(RunStatus::Failed);
8002 let run = runner.state.id.clone();
8003 let mut q = ask_open_question(&store, &run);
8004 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
8005 .unwrap();
8006 store.put(&mut q).unwrap();
8007
8008 runner.settle_questions();
8013 runner.settle_questions();
8014
8015 let back = store.get(&q.id).unwrap();
8016 assert_eq!(
8017 back.status,
8018 ask::QuestionStatus::Answered,
8019 "a real answer is a decision on record, never overwritten by a sweep"
8020 );
8021 }
8022
8023 #[tokio::test]
8034 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
8035 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
8036 let tmp = tempfile::tempdir().expect("tempdir");
8037 let repo = tmp.path().join("repo");
8038 std::fs::create_dir_all(&repo).unwrap();
8039 init_repo(&repo);
8040
8041 let mut config = Config::default();
8042 config.graph.worktree_root = Some(tmp.path().join("wt"));
8043
8044 let mut state = RunState::new(
8045 repo.clone(),
8046 "main".to_owned(),
8047 "deadbeef".to_owned(),
8048 "task".to_owned(),
8049 config,
8050 );
8051 let root = state.worktree_root();
8052 let wt_a = root.join("cand-A");
8053 let wt_b = root.join("cand-B");
8054 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
8055 .await
8056 .expect("worktree A");
8057 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
8058 .await
8059 .expect("worktree B");
8060
8061 state.candidates = vec![
8062 Candidate {
8063 index: 0,
8064 label: 'A',
8065 agent: "alpha".to_owned(),
8066 branch: "magi/x/A".to_owned(),
8067 worktree: wt_a.clone(),
8068 summary: String::new(),
8069 stat: String::new(),
8070 files: 0,
8071 commits: 0,
8072 empty: false,
8073 failed: None,
8074 verified_noop: None,
8075 duration_ms: 0,
8076 folded: false,
8077 },
8078 Candidate {
8079 index: 1,
8080 label: 'B',
8081 agent: "beta".to_owned(),
8082 branch: "magi/x/B".to_owned(),
8083 worktree: wt_b.clone(),
8084 summary: String::new(),
8085 stat: String::new(),
8086 files: 0,
8087 commits: 0,
8088 empty: false,
8089 failed: None,
8090 verified_noop: None,
8091 duration_ms: 0,
8092 folded: false,
8093 },
8094 ];
8095 state.tally = Some(Tally {
8096 first_choice: BTreeMap::from([('A', 1)]),
8097 borda: BTreeMap::new(),
8098 winner: 'A',
8099 rankings: 1,
8100 unanimous_initial: true,
8101 deliberated: false,
8102 changed_votes: 0,
8103 unanimous_final: true,
8104 tie_break: None,
8105 judges: 1,
8106 present: 1,
8107 quorum: 1,
8108 met_quorum: true,
8109 uncontested: None,
8110 });
8111 state.status = RunStatus::Ready;
8112
8113 fold_run(&mut state, false, &crate::run::home())
8114 .await
8115 .expect("fold_run");
8116
8117 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
8118 assert!(
8119 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
8120 "the unmerged winner's branch survives"
8121 );
8122 assert!(
8123 !state.candidates[0].folded,
8124 "the winner is not marked folded"
8125 );
8126
8127 assert!(!wt_b.exists(), "the loser's worktree is removed");
8128 assert!(
8129 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
8130 "the loser's branch is removed"
8131 );
8132 assert!(state.candidates[1].folded, "the loser is marked folded");
8133 }
8134
8135 #[tokio::test]
8138 async fn fold_run_keeps_a_branch_that_was_handed_to_a_later_run() {
8139 let tmp = tempfile::tempdir().expect("tempdir");
8140 let repo = tmp.path().join("repo");
8141 std::fs::create_dir_all(&repo).unwrap();
8142 init_repo(&repo);
8143 let home = tmp.path().join("home");
8144
8145 let mut config = Config::default();
8146 config.graph.worktree_root = Some(tmp.path().join("wt"));
8147 let mut state = RunState::new(
8148 repo.clone(),
8149 "main".to_owned(),
8150 "deadbeef".to_owned(),
8151 "task".to_owned(),
8152 config,
8153 );
8154 git::git(&repo, &["branch", "magi/x/A", "main"])
8156 .await
8157 .expect("branch");
8158 state.candidates = vec![Candidate {
8159 index: 0,
8160 label: 'A',
8161 agent: "alpha".to_owned(),
8162 branch: "magi/x/A".to_owned(),
8163 worktree: state.worktree_root().join("cand-A"),
8164 summary: String::new(),
8165 stat: String::new(),
8166 files: 0,
8167 commits: 0,
8168 empty: false,
8169 failed: None,
8170 verified_noop: None,
8171 duration_ms: 0,
8172 folded: true,
8173 }];
8174 state.released_to = Some("20260901-000000-new1".to_owned());
8175 state.released_branches = vec!["magi/x/A".to_owned()];
8176
8177 fold_run(&mut state, true, &home).await.expect("fold_run");
8178
8179 assert!(
8180 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
8181 "the handed-over branch survives a fold"
8182 );
8183 }
8184
8185 #[tokio::test]
8189 async fn an_empty_winner_is_detected_before_a_pull_request_is_attempted() {
8190 let tmp = tempfile::tempdir().expect("tempdir");
8191 let repo = tmp.path().join("repo");
8192 std::fs::create_dir_all(&repo).unwrap();
8193 init_repo(&repo);
8194 let run = |args: &[&str]| {
8195 let out = std::process::Command::new("git")
8196 .quiet()
8197 .args(args)
8198 .current_dir(&repo)
8199 .output()
8200 .expect("spawn git");
8201 assert!(out.status.success(), "git {args:?}");
8202 };
8203 run(&["branch", "magi/x/A"]);
8204 run(&["checkout", "-q", "-b", "magi/x/B"]);
8205 std::fs::write(repo.join("f.txt"), "x\n").unwrap();
8206 run(&["add", "-A"]);
8207 run(&["commit", "-q", "-m", "work"]);
8208 run(&["checkout", "-q", "main"]);
8209
8210 let mut state = RunState::new(
8211 repo.clone(),
8212 "main".to_owned(),
8213 "deadbeef".to_owned(),
8214 "task".to_owned(),
8215 Config::default(),
8216 );
8217 state.seeds = vec![refs::Seed {
8218 token: "magi/27b2/A".to_owned(),
8219 kind: refs::SeedKind::Unresolved,
8220 sha: String::new(),
8221 branch: true,
8222 detail: "no branch or commit named magi/27b2/A".to_owned(),
8223 }];
8224
8225 assert!(merge_is_empty(&repo, &state, "magi/x/A", MergeMode::Pr).await);
8226 assert!(merge_is_empty(&repo, &state, "magi/x/A", MergeMode::Local).await);
8227 assert!(!merge_is_empty(&repo, &state, "magi/x/B", MergeMode::Pr).await);
8228 let detail = empty_candidate_detail(&state, "main");
8229 assert!(detail.starts_with("empty candidate"), "{detail}");
8230 assert!(detail.contains("magi/27b2/A"), "{detail}");
8231 }
8232
8233 #[tokio::test]
8242 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
8243 let tmp = tempfile::tempdir().expect("tempdir");
8244 let repo = tmp.path().join("repo");
8245 std::fs::create_dir_all(&repo).unwrap();
8246 init_repo(&repo);
8247
8248 let mut config = Config::default();
8249 config.merge.mode = MergeMode::Local;
8250
8251 let mut state = RunState::new(
8252 repo.clone(),
8253 "main".to_owned(),
8254 "deadbeef".to_owned(),
8255 "task".to_owned(),
8256 config,
8257 );
8258 state.candidates = vec![Candidate {
8259 index: 0,
8260 label: 'A',
8261 agent: "alpha".to_owned(),
8262 branch: "does-not-exist".to_owned(),
8263 worktree: repo.clone(),
8264 summary: String::new(),
8265 stat: String::new(),
8266 files: 0,
8267 commits: 0,
8268 empty: false,
8269 failed: None,
8270 verified_noop: None,
8271 duration_ms: 0,
8272 folded: false,
8273 }];
8274 state.tally = Some(Tally {
8275 first_choice: BTreeMap::from([('A', 1)]),
8276 borda: BTreeMap::new(),
8277 winner: 'A',
8278 rankings: 1,
8279 unanimous_initial: true,
8280 deliberated: false,
8281 changed_votes: 0,
8282 unanimous_final: true,
8283 tie_break: None,
8284 judges: 0,
8285 present: 0,
8286 quorum: 0,
8287 met_quorum: true,
8288 uncontested: Some("only candidate A produced a change".to_owned()),
8289 });
8290 state.reviews = vec![ReviewRound {
8291 round: 1,
8292 head: "deadbeef".to_owned(),
8293 verified_head: None,
8294 verified_at: None,
8295 reviews: Vec::new(),
8296 e2e: Vec::new(),
8297 fix: None,
8298 blocking: 0,
8299 answered: 0,
8300 expected: 0,
8301 clean: true,
8302 verify_retried: false,
8303 e2e_deferred: false,
8304 e2e_defer_reason: None,
8305 progressed: false,
8306 vote_split: false,
8307 reconsideration: Vec::new(),
8308 verdict: None,
8309 }];
8310 state.gate = vec![CommandOutcome {
8311 command: "test".to_owned(),
8312 code: Some(0),
8313 output_tail: String::new(),
8314 duration_ms: 0,
8315 resource_blocked: false,
8316 }];
8317 state.gate_ran = true;
8318 state.status = RunStatus::Ready;
8323 state.merge = Some(MergeOutcome {
8324 mode: MergeMode::Local,
8325 ok: false,
8326 detail: "already concluded".to_owned(),
8327 empty: false,
8328 });
8329
8330 let mut runner = Runner {
8331 state,
8332 roles: ResolvedRoles {
8333 implementers: Vec::new(),
8334 judges: Vec::new(),
8335 reviewers: Vec::new(),
8336 fixer: None,
8337 conductor: conductor(),
8338 implementer_roster: Vec::new(),
8339 },
8340 sem: Arc::new(Semaphore::new(1)),
8341 pause: Pause::new(),
8342 interrupt: Pause::new(),
8343 };
8344
8345 runner.merge().await.expect("merge");
8346
8347 assert_eq!(
8348 runner.state.status,
8349 RunStatus::Ready,
8350 "a concluded run's status must not change on reentry"
8351 );
8352 assert_eq!(
8353 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8354 Some("already concluded"),
8355 "merge must not run again once the node already recorded an outcome"
8356 );
8357 }
8358
8359 #[tokio::test]
8368 async fn merge_refuses_a_gate_that_has_not_actually_run() {
8369 let tmp = tempfile::tempdir().expect("tempdir");
8370 let repo = tmp.path().join("repo");
8371 std::fs::create_dir_all(&repo).unwrap();
8372 init_repo(&repo);
8373
8374 let mut config = Config::default();
8375 config.merge.mode = MergeMode::Local;
8376
8377 let mut state = RunState::new(
8378 repo.clone(),
8379 "main".to_owned(),
8380 "deadbeef".to_owned(),
8381 "task".to_owned(),
8382 config,
8383 );
8384 state.candidates = vec![Candidate {
8385 index: 0,
8386 label: 'A',
8387 agent: "alpha".to_owned(),
8388 branch: "does-not-exist".to_owned(),
8389 worktree: repo.clone(),
8390 summary: String::new(),
8391 stat: String::new(),
8392 files: 0,
8393 commits: 0,
8394 empty: false,
8395 failed: None,
8396 verified_noop: None,
8397 duration_ms: 0,
8398 folded: false,
8399 }];
8400 state.tally = Some(Tally {
8401 first_choice: BTreeMap::from([('A', 1)]),
8402 borda: BTreeMap::new(),
8403 winner: 'A',
8404 rankings: 1,
8405 unanimous_initial: true,
8406 deliberated: false,
8407 changed_votes: 0,
8408 unanimous_final: true,
8409 tie_break: None,
8410 judges: 0,
8411 present: 0,
8412 quorum: 0,
8413 met_quorum: true,
8414 uncontested: Some("only candidate A produced a change".to_owned()),
8415 });
8416 state.reviews = vec![ReviewRound {
8417 round: 1,
8418 head: "deadbeef".to_owned(),
8419 verified_head: None,
8420 verified_at: None,
8421 reviews: Vec::new(),
8422 e2e: Vec::new(),
8423 fix: None,
8424 blocking: 0,
8425 answered: 0,
8426 expected: 0,
8427 clean: true,
8428 verify_retried: false,
8429 e2e_deferred: false,
8430 e2e_defer_reason: None,
8431 progressed: false,
8432 vote_split: false,
8433 reconsideration: Vec::new(),
8434 verdict: None,
8435 }];
8436 state.gate = Vec::new();
8438 state.gate_ran = false;
8439 state.status = RunStatus::Gating;
8440
8441 let mut runner = Runner {
8442 state,
8443 roles: ResolvedRoles {
8444 implementers: Vec::new(),
8445 judges: Vec::new(),
8446 reviewers: Vec::new(),
8447 fixer: None,
8448 conductor: conductor(),
8449 implementer_roster: Vec::new(),
8450 },
8451 sem: Arc::new(Semaphore::new(1)),
8452 pause: Pause::new(),
8453 interrupt: Pause::new(),
8454 };
8455
8456 runner.merge().await.expect("merge");
8457
8458 assert!(
8459 runner.state.merge.is_none(),
8460 "an empty gate must never be read as a passing one: {:?}",
8461 runner.state.merge
8462 );
8463 }
8464
8465 #[tokio::test]
8472 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
8473 let tmp = tempfile::tempdir().expect("tempdir");
8474 let repo = tmp.path().join("repo");
8475 std::fs::create_dir_all(&repo).unwrap();
8476 init_repo(&repo);
8477
8478 let config = Config::default();
8480
8481 let mut state = RunState::new(
8482 repo.clone(),
8483 "main".to_owned(),
8484 "deadbeef".to_owned(),
8485 "task".to_owned(),
8486 config,
8487 );
8488 state.candidates = vec![Candidate {
8489 index: 0,
8490 label: 'A',
8491 agent: "alpha".to_owned(),
8492 branch: "does-not-exist".to_owned(),
8493 worktree: repo.clone(),
8494 summary: String::new(),
8495 stat: String::new(),
8496 files: 0,
8497 commits: 0,
8498 empty: false,
8499 failed: None,
8500 verified_noop: None,
8501 duration_ms: 0,
8502 folded: false,
8503 }];
8504 state.tally = Some(Tally {
8505 first_choice: BTreeMap::from([('A', 1)]),
8506 borda: BTreeMap::new(),
8507 winner: 'A',
8508 rankings: 1,
8509 unanimous_initial: true,
8510 deliberated: false,
8511 changed_votes: 0,
8512 unanimous_final: true,
8513 tie_break: None,
8514 judges: 0,
8515 present: 0,
8516 quorum: 0,
8517 met_quorum: true,
8518 uncontested: Some("only candidate A produced a change".to_owned()),
8519 });
8520 state.reviews = vec![ReviewRound {
8521 round: 1,
8522 head: "deadbeef".to_owned(),
8523 verified_head: None,
8524 verified_at: None,
8525 reviews: Vec::new(),
8526 e2e: Vec::new(),
8527 fix: None,
8528 blocking: 0,
8529 answered: 0,
8530 expected: 0,
8531 clean: true,
8532 verify_retried: false,
8533 e2e_deferred: false,
8534 e2e_defer_reason: None,
8535 progressed: false,
8536 vote_split: false,
8537 reconsideration: Vec::new(),
8538 verdict: None,
8539 }];
8540
8541 let mut runner = Runner {
8542 state,
8543 roles: ResolvedRoles {
8544 implementers: Vec::new(),
8545 judges: Vec::new(),
8546 reviewers: Vec::new(),
8547 fixer: None,
8548 conductor: conductor(),
8549 implementer_roster: Vec::new(),
8550 },
8551 sem: Arc::new(Semaphore::new(1)),
8552 pause: Pause::new(),
8553 interrupt: Pause::new(),
8554 };
8555
8556 runner.gate().await.expect("gate");
8557 assert!(
8558 runner.state.gate_ran,
8559 "zero configured commands is still a real attempt, not an unrun gate"
8560 );
8561 assert!(runner.state.gate.is_empty());
8562 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
8563 assert_ne!(
8564 runner.state.status,
8565 RunStatus::Blocked,
8566 "a gate with nothing to check must not read as failed"
8567 );
8568
8569 runner.merge().await.expect("merge");
8570 assert_eq!(
8571 runner.state.status,
8572 RunStatus::Ready,
8573 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
8574 );
8575 }
8576
8577 #[tokio::test]
8588 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
8589 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8590 let home = crate::run::home();
8591
8592 let tmp = tempfile::tempdir().expect("tempdir");
8593 let repo = tmp.path().join("repo");
8594 std::fs::create_dir_all(&repo).unwrap();
8595 init_repo(&repo);
8596 let cache_dir = tmp.path().join("target");
8599
8600 let mut config = Config::default();
8601 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
8602 config.graph.timeout_verify = Some(2);
8605
8606 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8607 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8608 .expect("no io error acquiring directly")
8609 {
8610 crate::cache::AcquireOutcome::Acquired(g) => g,
8611 crate::cache::AcquireOutcome::Busy(b) => {
8612 panic!("expected the direct acquire to win the lease first: {b:?}")
8613 }
8614 };
8615
8616 let mut state = RunState::new(
8617 repo.clone(),
8618 "main".to_owned(),
8619 "deadbeef".to_owned(),
8620 "task".to_owned(),
8621 config,
8622 );
8623 state.candidates = vec![Candidate {
8624 index: 0,
8625 label: 'A',
8626 agent: "alpha".to_owned(),
8627 branch: "does-not-exist".to_owned(),
8628 worktree: repo.clone(),
8629 summary: String::new(),
8630 stat: String::new(),
8631 files: 0,
8632 commits: 0,
8633 empty: false,
8634 failed: None,
8635 verified_noop: None,
8636 duration_ms: 0,
8637 folded: false,
8638 }];
8639 state.tally = Some(Tally {
8640 first_choice: BTreeMap::from([('A', 1)]),
8641 borda: BTreeMap::new(),
8642 winner: 'A',
8643 rankings: 1,
8644 unanimous_initial: true,
8645 deliberated: false,
8646 changed_votes: 0,
8647 unanimous_final: true,
8648 tie_break: None,
8649 judges: 0,
8650 present: 0,
8651 quorum: 0,
8652 met_quorum: true,
8653 uncontested: Some("only candidate A produced a change".to_owned()),
8654 });
8655 state.reviews = vec![ReviewRound {
8656 round: 1,
8657 head: "deadbeef".to_owned(),
8658 verified_head: None,
8659 verified_at: None,
8660 reviews: Vec::new(),
8661 e2e: Vec::new(),
8662 fix: None,
8663 blocking: 0,
8664 answered: 0,
8665 expected: 0,
8666 clean: true,
8667 verify_retried: false,
8668 e2e_deferred: false,
8669 e2e_defer_reason: None,
8670 progressed: false,
8671 vote_split: false,
8672 reconsideration: Vec::new(),
8673 verdict: None,
8674 }];
8675
8676 let mut runner = Runner {
8677 state,
8678 roles: ResolvedRoles {
8679 implementers: Vec::new(),
8680 judges: Vec::new(),
8681 reviewers: Vec::new(),
8682 fixer: None,
8683 conductor: conductor(),
8684 implementer_roster: Vec::new(),
8685 },
8686 sem: Arc::new(Semaphore::new(1)),
8687 pause: Pause::new(),
8688 interrupt: Pause::new(),
8689 };
8690
8691 let started = std::time::Instant::now();
8692 runner.gate().await.expect("gate");
8693 assert!(
8694 started.elapsed() < Duration::from_secs(1),
8695 "a gate with nothing to run must never wait on a lease it never needed"
8696 );
8697 assert!(
8698 runner.state.gate_ran,
8699 "zero commands is still a real, immediate attempt"
8700 );
8701 assert!(runner.state.gate.is_empty());
8702 assert_ne!(
8703 runner.state.status,
8704 RunStatus::Blocked,
8705 "must not read as resource-blocked on a lease it never asked for"
8706 );
8707 }
8708
8709 #[tokio::test]
8719 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8720 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8721
8722 let tmp = tempfile::tempdir().expect("tempdir");
8723 let repo = tmp.path().join("repo");
8724 std::fs::create_dir_all(&repo).unwrap();
8725 init_repo(&repo);
8726
8727 let mut config = Config::default();
8728 config.verify.gate = vec![
8729 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8730 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8731 .to_owned(),
8732 ];
8733
8734 let mut state = RunState::new(
8735 repo.clone(),
8736 "main".to_owned(),
8737 "deadbeef".to_owned(),
8738 "task".to_owned(),
8739 config,
8740 );
8741 let run_id = state.id.clone();
8742 state.candidates = vec![Candidate {
8743 index: 0,
8744 label: 'A',
8745 agent: "alpha".to_owned(),
8746 branch: "does-not-exist".to_owned(),
8747 worktree: repo.clone(),
8748 summary: String::new(),
8749 stat: String::new(),
8750 files: 0,
8751 commits: 0,
8752 empty: false,
8753 failed: None,
8754 verified_noop: None,
8755 duration_ms: 0,
8756 folded: false,
8757 }];
8758 state.tally = Some(Tally {
8759 first_choice: BTreeMap::from([('A', 1)]),
8760 borda: BTreeMap::new(),
8761 winner: 'A',
8762 rankings: 1,
8763 unanimous_initial: true,
8764 deliberated: false,
8765 changed_votes: 0,
8766 unanimous_final: true,
8767 tie_break: None,
8768 judges: 0,
8769 present: 0,
8770 quorum: 0,
8771 met_quorum: true,
8772 uncontested: Some("only candidate A produced a change".to_owned()),
8773 });
8774 state.reviews = vec![ReviewRound {
8775 round: 1,
8776 head: "deadbeef".to_owned(),
8777 verified_head: None,
8778 verified_at: None,
8779 reviews: Vec::new(),
8780 e2e: Vec::new(),
8781 fix: None,
8782 blocking: 0,
8783 answered: 0,
8784 expected: 0,
8785 clean: true,
8786 verify_retried: false,
8787 e2e_deferred: false,
8788 e2e_defer_reason: None,
8789 progressed: false,
8790 vote_split: false,
8791 reconsideration: Vec::new(),
8792 verdict: None,
8793 }];
8794
8795 let mut runner = Runner {
8796 state,
8797 roles: ResolvedRoles {
8798 implementers: Vec::new(),
8799 judges: Vec::new(),
8800 reviewers: Vec::new(),
8801 fixer: None,
8802 conductor: conductor(),
8803 implementer_roster: Vec::new(),
8804 },
8805 sem: Arc::new(Semaphore::new(1)),
8806 pause: Pause::new(),
8807 interrupt: Pause::new(),
8808 };
8809
8810 let started_marker = repo.join("started.marker");
8811 let release_marker = repo.join("release.marker");
8812 let poller = tokio::spawn(async move {
8813 for _ in 0..100 {
8818 if started_marker.exists()
8819 && let Ok(s) = crate::run::RunState::load(&run_id)
8820 && let Some(a) = s.active.get("gate")
8821 {
8822 std::fs::write(&release_marker, b"go").expect("release marker");
8823 return Some(a.clone());
8824 }
8825 tokio::time::sleep(Duration::from_millis(50)).await;
8826 }
8827 None
8828 });
8829
8830 runner.gate().await.expect("gate");
8831 let captured = poller.await.expect("poller task");
8832 let captured = captured.expect(
8833 "the poller never saw a `gate` task entry in run.json while the command was \
8834 still blocked on its own release marker",
8835 );
8836
8837 assert_eq!(captured.task.as_deref(), Some("gate"));
8838 assert_eq!(captured.node, "gate");
8839 assert_eq!(captured.index, Some(1));
8840 assert_eq!(captured.total, Some(1));
8841 assert!(
8842 captured
8843 .command
8844 .as_deref()
8845 .is_some_and(|c| c.contains("started.marker")),
8846 "{captured:?}"
8847 );
8848
8849 assert!(
8850 runner.state.active.is_empty(),
8851 "the entry must be cleared once the command actually finished: {:?}",
8852 runner.state.active
8853 );
8854 assert!(runner.state.gate_ran);
8855 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8856 }
8857
8858 #[tokio::test]
8871 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8872 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8873 let home = crate::run::home();
8874
8875 let tmp = tempfile::tempdir().expect("tempdir");
8876 let repo = tmp.path().join("repo");
8877 std::fs::create_dir_all(&repo).unwrap();
8878 init_repo(&repo);
8879 let head = crate::git::rev_parse(&repo, "HEAD")
8880 .await
8881 .expect("rev-parse");
8882 let cache_dir = tmp.path().join("target");
8885
8886 let mut config = Config::default();
8887 config.verify.e2e = vec![format!(
8888 "CARGO_TARGET_DIR='{}' test -f README.md",
8889 cache_dir.display()
8890 )];
8891 config.graph.review_rounds = 1;
8892 config.graph.timeout_verify = Some(2);
8895
8896 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8897 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8898 .expect("no io error acquiring directly")
8899 {
8900 crate::cache::AcquireOutcome::Acquired(g) => g,
8901 crate::cache::AcquireOutcome::Busy(b) => {
8902 panic!("expected the direct acquire to win the lease first: {b:?}")
8903 }
8904 };
8905
8906 let mut state = RunState::new(
8907 repo.clone(),
8908 "main".to_owned(),
8909 head.clone(),
8910 "task".to_owned(),
8911 config,
8912 );
8913 state.candidates = vec![Candidate {
8914 index: 0,
8915 label: 'A',
8916 agent: "alpha".to_owned(),
8917 branch: "does-not-exist".to_owned(),
8918 worktree: repo.clone(),
8919 summary: String::new(),
8920 stat: String::new(),
8921 files: 0,
8922 commits: 0,
8923 empty: false,
8924 failed: None,
8925 verified_noop: None,
8926 duration_ms: 0,
8927 folded: false,
8928 }];
8929 state.tally = Some(Tally {
8930 first_choice: BTreeMap::from([('A', 1)]),
8931 borda: BTreeMap::new(),
8932 winner: 'A',
8933 rankings: 1,
8934 unanimous_initial: true,
8935 deliberated: false,
8936 changed_votes: 0,
8937 unanimous_final: true,
8938 tie_break: None,
8939 judges: 0,
8940 present: 0,
8941 quorum: 0,
8942 met_quorum: true,
8943 uncontested: Some("only candidate A produced a change".to_owned()),
8944 });
8945 state.reviews = vec![ReviewRound {
8949 round: 1,
8950 head: head.clone(),
8951 verified_head: None,
8952 verified_at: None,
8953 reviews: Vec::new(),
8954 e2e: Vec::new(),
8955 fix: None,
8956 blocking: 1,
8957 answered: 1,
8958 expected: 1,
8959 clean: false,
8960 verify_retried: false,
8961 e2e_deferred: true,
8962 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8963 progressed: false,
8964 vote_split: false,
8965 reconsideration: Vec::new(),
8966 verdict: None,
8967 }];
8968
8969 let mut runner = Runner {
8970 state,
8971 roles: ResolvedRoles {
8972 implementers: Vec::new(),
8973 judges: Vec::new(),
8974 reviewers: Vec::new(),
8975 fixer: None,
8976 conductor: conductor(),
8977 implementer_roster: Vec::new(),
8978 },
8979 sem: Arc::new(Semaphore::new(1)),
8980 pause: Pause::new(),
8981 interrupt: Pause::new(),
8982 };
8983
8984 let shell = runner.state.config.shell();
8985 runner
8986 .stop_reviewing("round budget spent", &shell, &repo)
8987 .await
8988 .expect("stop_reviewing");
8989
8990 let last = runner.state.reviews.last().expect("round record");
8991 assert_eq!(
8992 last.e2e_status(),
8993 E2eStatus::ResourceBlocked,
8994 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8995 failed: {last:?}"
8996 );
8997 assert_eq!(
8998 last.verified_head.as_deref(),
8999 Some(head.as_str()),
9000 "which commit this attempt targeted is known even though nothing finished checking \
9001 it"
9002 );
9003 let first_attempt_at = last
9004 .verified_at
9005 .expect("when this attempt ran is known too");
9006 assert_ne!(
9007 runner.state.status,
9008 RunStatus::Blocked,
9009 "contention is evidence about the machine, not the patch — it must not settle the \
9010 run as blocked: {:?}",
9011 runner.state.status
9012 );
9013 assert!(
9014 !runner
9015 .state
9016 .events
9017 .iter()
9018 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
9019 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
9020 runner.state.events
9021 );
9022
9023 runner
9028 .stop_reviewing("round budget spent", &shell, &repo)
9029 .await
9030 .expect("stop_reviewing retry");
9031 assert_eq!(
9032 runner.state.reviews.len(),
9033 1,
9034 "no new round was started: {:?}",
9035 runner.state.reviews
9036 );
9037 let last = runner.state.reviews.last().expect("round record");
9038 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
9039 assert!(
9040 last.verified_at.expect("still known") > first_attempt_at,
9041 "a second reentry must be a fresh attempt, not a stale copy of the first"
9042 );
9043 assert_ne!(runner.state.status, RunStatus::Blocked);
9044
9045 held.release();
9046 }
9047
9048 #[tokio::test]
9060 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
9061 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
9062 let home = crate::run::home();
9063
9064 let tmp = tempfile::tempdir().expect("tempdir");
9065 let repo = tmp.path().join("repo");
9066 std::fs::create_dir_all(&repo).unwrap();
9067 init_repo(&repo);
9068 let head = crate::git::rev_parse(&repo, "HEAD")
9069 .await
9070 .expect("rev-parse");
9071 let cache_dir = tmp.path().join("target");
9072
9073 let mut config = Config::default();
9074 config.verify.e2e = vec![format!(
9075 "CARGO_TARGET_DIR='{}' test -f README.md",
9076 cache_dir.display()
9077 )];
9078 config.graph.review_rounds = 1;
9079 config.graph.timeout_verify = Some(2);
9080
9081 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
9082 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
9083 .expect("no io error acquiring directly")
9084 {
9085 crate::cache::AcquireOutcome::Acquired(g) => g,
9086 crate::cache::AcquireOutcome::Busy(b) => {
9087 panic!("expected the direct acquire to win the lease first: {b:?}")
9088 }
9089 };
9090
9091 let mut state = RunState::new(
9092 repo.clone(),
9093 "main".to_owned(),
9094 head.clone(),
9095 "task".to_owned(),
9096 config,
9097 );
9098 state.candidates = vec![Candidate {
9099 index: 0,
9100 label: 'A',
9101 agent: "alpha".to_owned(),
9102 branch: "does-not-exist".to_owned(),
9103 worktree: repo.clone(),
9104 summary: String::new(),
9105 stat: String::new(),
9106 files: 0,
9107 commits: 0,
9108 empty: false,
9109 failed: None,
9110 verified_noop: None,
9111 duration_ms: 0,
9112 folded: false,
9113 }];
9114 state.tally = Some(Tally {
9115 first_choice: BTreeMap::from([('A', 1)]),
9116 borda: BTreeMap::new(),
9117 winner: 'A',
9118 rankings: 1,
9119 unanimous_initial: true,
9120 deliberated: false,
9121 changed_votes: 0,
9122 unanimous_final: true,
9123 tie_break: None,
9124 judges: 0,
9125 present: 0,
9126 quorum: 0,
9127 met_quorum: true,
9128 uncontested: Some("only candidate A produced a change".to_owned()),
9129 });
9130 state.reviews = vec![ReviewRound {
9134 round: 1,
9135 head: head.clone(),
9136 verified_head: Some(head.clone()),
9137 verified_at: Some(jiff::Timestamp::now()),
9138 reviews: Vec::new(),
9139 e2e: vec![CommandOutcome {
9140 command: format!(
9141 "CARGO_TARGET_DIR='{}' test -f README.md",
9142 cache_dir.display()
9143 ),
9144 code: None,
9145 output_tail: "waiting for the shared build cache".to_owned(),
9146 duration_ms: 0,
9147 resource_blocked: true,
9148 }],
9149 fix: None,
9150 blocking: 1,
9151 answered: 1,
9152 expected: 1,
9153 clean: false,
9154 verify_retried: false,
9155 e2e_deferred: false,
9156 e2e_defer_reason: None,
9157 progressed: false,
9158 vote_split: false,
9159 reconsideration: Vec::new(),
9160 verdict: None,
9161 }];
9162
9163 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
9164 let mut runner = Runner {
9165 state,
9166 roles: ResolvedRoles {
9167 implementers: Vec::new(),
9168 judges: Vec::new(),
9169 reviewers: Vec::new(),
9170 fixer: None,
9171 conductor: conductor(),
9172 implementer_roster: Vec::new(),
9173 },
9174 sem: Arc::new(Semaphore::new(1)),
9175 pause: Pause::new(),
9176 interrupt: Pause::new(),
9177 };
9178
9179 runner.review_loop().await.expect("review_loop");
9184
9185 assert_eq!(
9186 runner.state.reviews.len(),
9187 1,
9188 "no new round was started on top of the unresolved one: {:?}",
9189 runner.state.reviews
9190 );
9191 let last = &runner.state.reviews[0];
9192 assert_eq!(
9193 last.e2e_status(),
9194 E2eStatus::ResourceBlocked,
9195 "still contended: {last:?}"
9196 );
9197 assert!(
9198 last.verified_at.expect("still known") > first_attempt_at,
9199 "review_loop must have actually retried the check, not left it exactly as found"
9200 );
9201 assert_ne!(
9202 runner.state.status,
9203 RunStatus::Blocked,
9204 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
9205 runner.state.status
9206 );
9207
9208 held.release();
9209 }
9210
9211 #[tokio::test]
9212 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
9213 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
9214 let tmp = tempfile::tempdir().expect("tempdir");
9215 let repo = tmp.path().join("repo");
9216 std::fs::create_dir_all(&repo).unwrap();
9217 init_repo(&repo);
9218
9219 let mut config = Config::default();
9220 config.merge.mode = MergeMode::Pr;
9221 config.graph.land = true;
9222 config.graph.land_approval = false;
9223
9224 let mut state = RunState::new(
9225 repo.clone(),
9226 "main".to_owned(),
9227 "deadbeef".to_owned(),
9228 "task".to_owned(),
9229 config,
9230 );
9231 state.candidates = vec![Candidate {
9232 index: 0,
9233 label: 'A',
9234 agent: "alpha".to_owned(),
9235 branch: "does-not-exist".to_owned(),
9236 worktree: repo.clone(),
9237 summary: String::new(),
9238 stat: String::new(),
9239 files: 0,
9240 commits: 0,
9241 empty: false,
9242 failed: None,
9243 verified_noop: None,
9244 duration_ms: 0,
9245 folded: false,
9246 }];
9247 state.tally = Some(Tally {
9248 first_choice: BTreeMap::from([('A', 1)]),
9249 borda: BTreeMap::new(),
9250 winner: 'A',
9251 rankings: 1,
9252 unanimous_initial: true,
9253 deliberated: false,
9254 changed_votes: 0,
9255 unanimous_final: true,
9256 tie_break: None,
9257 judges: 0,
9258 present: 0,
9259 quorum: 0,
9260 met_quorum: true,
9261 uncontested: Some("only candidate A produced a change".to_owned()),
9262 });
9263 state.reviews = vec![ReviewRound {
9264 round: 1,
9265 head: "deadbeef".to_owned(),
9266 verified_head: None,
9267 verified_at: None,
9268 reviews: Vec::new(),
9269 e2e: Vec::new(),
9270 fix: None,
9271 blocking: 0,
9272 answered: 0,
9273 expected: 0,
9274 clean: true,
9275 verify_retried: false,
9276 e2e_deferred: false,
9277 e2e_defer_reason: None,
9278 progressed: false,
9279 vote_split: false,
9280 reconsideration: Vec::new(),
9281 verdict: None,
9282 }];
9283 state.gate = vec![CommandOutcome {
9284 command: "test".to_owned(),
9285 code: Some(0),
9286 output_tail: String::new(),
9287 duration_ms: 0,
9288 resource_blocked: false,
9289 }];
9290 state.gate_ran = true;
9291 state.status = RunStatus::Landing;
9295 state.merge = Some(MergeOutcome {
9296 mode: MergeMode::Pr,
9297 ok: true,
9298 detail: "https://example.invalid/x/y/pull/1".to_owned(),
9299 empty: false,
9300 });
9301
9302 ask_test_home();
9306 let store = ask::Questions::open();
9307 let q = ask_open_question(&store, &state.id);
9308
9309 let mut runner = Runner {
9310 state,
9311 roles: ResolvedRoles {
9312 implementers: Vec::new(),
9313 judges: Vec::new(),
9314 reviewers: Vec::new(),
9315 fixer: None,
9316 conductor: conductor(),
9317 implementer_roster: Vec::new(),
9318 },
9319 sem: Arc::new(Semaphore::new(1)),
9320 pause: Pause::new(),
9321 interrupt: Pause::new(),
9322 };
9323
9324 runner.execute().await.expect("execute");
9329
9330 assert_eq!(
9331 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
9332 Some("https://example.invalid/x/y/pull/1"),
9333 "reentry must not push again or open a second pull request over the \
9334 one `land` is already watching"
9335 );
9336 assert_ne!(
9337 runner.state.status,
9338 RunStatus::Landing,
9339 "land could not actually reach the fake pull request, so it must \
9340 have given up rather than left the run silently parked forever"
9341 );
9342 assert_eq!(runner.state.status, RunStatus::Blocked);
9346 assert!(
9347 store.get(&q.id).unwrap().status.open(),
9348 "Blocked is still alive; settle_questions must have been a no-op here"
9349 );
9350 }
9351
9352 fn state_with_round(round: ReviewRound) -> RunState {
9353 let mut s = RunState::new(
9354 PathBuf::from("/repo"),
9355 "main".to_owned(),
9356 "abc1234".to_owned(),
9357 "add retries".to_owned(),
9358 Config::default(),
9359 );
9360 s.reviews = vec![round];
9361 s
9362 }
9363
9364 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
9365 crate::verdict::Finding {
9366 id: id.to_owned(),
9367 severity,
9368 file: None,
9369 line: None,
9370 title: title.to_owned(),
9371 detail: String::new(),
9372 }
9373 }
9374
9375 #[test]
9376 fn pr_body_names_open_findings_and_declined_ones() {
9377 let round = ReviewRound {
9378 round: 2,
9379 head: "deadbee".to_owned(),
9380 verified_head: None,
9381 verified_at: None,
9382 reviews: vec![ReviewRecord {
9383 attempts: 0,
9384 reviewer: 1,
9385 agent: "alpha".to_owned(),
9386 summary: String::new(),
9387 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
9388 vote: None,
9389 failed: None,
9390 duration_ms: 0,
9391 }],
9392 e2e: vec![CommandOutcome {
9393 command: "cargo test".to_owned(),
9394 code: Some(0),
9395 output_tail: String::new(),
9396 duration_ms: 0,
9397 resource_blocked: false,
9398 }],
9399 verify_retried: false,
9400 e2e_deferred: false,
9401 e2e_defer_reason: None,
9402 fix: Some(FixRecord {
9403 agent: "alpha".to_owned(),
9404 addressed: Vec::new(),
9405 rejected: vec![crate::verdict::Rejection {
9406 id: "R1-1-1".to_owned(),
9407 why: "not reachable from any caller".to_owned(),
9408 }],
9409 notes: String::new(),
9410 committed: true,
9411 failed: None,
9412 duration_ms: 0,
9413 continuation: None,
9414 }),
9415 blocking: 0,
9416 answered: 1,
9417 expected: 1,
9418 clean: false,
9419 progressed: true,
9420 vote_split: false,
9421 reconsideration: Vec::new(),
9422 verdict: None,
9423 };
9424 let state = state_with_round(round);
9425 let body = pr_message(&state, 'A').body;
9426
9427 assert!(body.contains("add retries"), "the task must still be there");
9428 assert!(body.contains("R2-1-1"), "{body}");
9429 assert!(body.contains("unused import"), "{body}");
9430 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
9431 assert!(
9432 body.contains("not reachable from any caller"),
9433 "the reason it was declined: {body}"
9434 );
9435 }
9436
9437 #[test]
9438 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
9439 let round = ReviewRound {
9440 round: 1,
9441 head: "deadbee".to_owned(),
9442 verified_head: None,
9443 verified_at: None,
9444 reviews: vec![ReviewRecord {
9445 attempts: 0,
9446 reviewer: 1,
9447 agent: "alpha".to_owned(),
9448 summary: String::new(),
9449 findings: Vec::new(),
9450 vote: None,
9451 failed: None,
9452 duration_ms: 0,
9453 }],
9454 e2e: Vec::new(),
9455 verify_retried: false,
9456 e2e_deferred: false,
9457 e2e_defer_reason: None,
9458 fix: None,
9459 blocking: 0,
9460 answered: 1,
9461 expected: 1,
9462 clean: true,
9463 progressed: false,
9464 vote_split: false,
9465 reconsideration: Vec::new(),
9466 verdict: None,
9467 };
9468 let state = state_with_round(round);
9469 let body = pr_message(&state, 'A').body;
9470 assert!(!body.contains("Open review findings"), "{body}");
9471 assert!(!body.contains("Declined"), "{body}");
9472 }
9473
9474 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
9475 let mut state = RunState::new(
9476 PathBuf::from("/repo"),
9477 "main".to_owned(),
9478 "abc1234".to_owned(),
9479 instruction.to_owned(),
9480 Config::default(),
9481 );
9482 state.candidates.push(Candidate {
9483 index: 0,
9484 label: 'A',
9485 agent: "alpha".to_owned(),
9486 branch: "magi/x/A".to_owned(),
9487 worktree: PathBuf::from("/wt"),
9488 summary: summary.to_owned(),
9489 stat: String::new(),
9490 files: 1,
9491 commits: 1,
9492 empty: false,
9493 failed: None,
9494 verified_noop: None,
9495 folded: false,
9496 duration_ms: 0,
9497 });
9498 state
9499 }
9500
9501 #[test]
9502 fn pr_message_describes_the_change_not_the_task() {
9503 let state = state_with_summary(
9504 "今回やってほしいこと: results projector を直す",
9505 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
9506 );
9507 let m = pr_message(&state, 'A');
9508 assert_eq!(m.title, "fix(web): batch the runs list reads");
9509 assert!(
9510 m.body.starts_with("## Summary\n\n- reads run.json once"),
9511 "{}",
9512 m.body
9513 );
9514 assert!(!m.body.contains("TITLE:"), "{}", m.body);
9515 let task_at = m.body.find("今回やってほしいこと").unwrap();
9516 let details_at = m.body.find("<details>").unwrap();
9517 assert!(
9518 details_at < task_at,
9519 "the task lives inside <details>: {}",
9520 m.body
9521 );
9522 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
9523 assert!(m.body.contains("magi:candidate-a"));
9524 }
9525
9526 #[test]
9527 fn pr_message_falls_back_to_the_task_without_a_title_line() {
9528 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
9529 let m = pr_message(&state, 'A');
9530 assert_eq!(m.title, "add retries");
9531 assert!(
9532 m.body.contains("## Summary\n\n- did some things"),
9533 "{}",
9534 m.body
9535 );
9536
9537 let none = RunState::new(
9538 PathBuf::from("/repo"),
9539 "main".to_owned(),
9540 "abc1234".to_owned(),
9541 "add retries".to_owned(),
9542 Config::default(),
9543 );
9544 let m = pr_message(&none, 'A');
9545 assert_eq!(m.title, "add retries");
9546 assert!(!m.body.contains("## Summary"), "{}", m.body);
9547 }
9548
9549 #[test]
9550 fn pr_message_refuses_the_candidate_commit_subject() {
9551 for bad in [
9552 "TITLE: magi: candidate A (uncommitted work)",
9553 "TITLE: chore: stuff (uncommitted work)",
9554 "TITLE: ",
9555 ] {
9556 let state = state_with_summary("add retries", bad);
9557 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
9558 }
9559 }
9560
9561 #[test]
9562 fn pr_message_bounds_a_very_long_task_and_title() {
9563 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
9564 let state = state_with_summary(&long, "- nothing");
9565 let m = pr_message(&state, 'A');
9566 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9567 assert!(!m.title.contains('\n'));
9568
9569 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
9570 let m = pr_message(&state, 'A');
9571 assert!(m.title.starts_with("feat: "));
9572 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9573 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
9574 }
9575
9576 #[test]
9577 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
9578 let mut state = state_with_summary(
9582 "add retries",
9583 "TITLE: fix(web): batch reads\n- reads run.json once",
9584 );
9585 state.config.graph.language = "ja".to_owned();
9586 let m = pr_message(&state, 'A');
9587 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
9588
9589 let task = "今回やってほしいこと: results projector を直す";
9592 let mut state = state_with_summary(task, "- no title line");
9593 state.config.graph.language = "ja".to_owned();
9594 let m = pr_message(&state, 'A');
9595 assert_eq!(
9596 m.title,
9597 format!("chore: land candidate A of run {}", state.id)
9598 );
9599 assert!(
9600 m.body.contains(&format!(
9601 "<summary>Original task</summary>\n\n{task}\n\n</details>"
9602 )),
9603 "{}",
9604 m.body
9605 );
9606 }
9607
9608 #[test]
9609 fn pr_message_scrubs_home_paths_and_addresses() {
9610 let state = state_with_summary(
9611 "fix it in /Users/someone/src/x",
9612 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9613 );
9614 let m = pr_message(&state, 'A');
9615 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9616 assert!(!m.body.contains(leak), "{}", m.body);
9617 }
9618 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9619 }
9620
9621 #[test]
9622 fn pr_message_survives_a_task_that_closes_details() {
9623 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9624 let m = pr_message(&state, 'A');
9625 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9626 }
9627
9628 #[test]
9629 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9630 let cmd = manual_merge_command(
9631 MergeStyle::Squash,
9632 Path::new("/repo"),
9633 "b",
9634 "fix: \"quoted\" $(x) `y`\n\nbody",
9635 );
9636 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9637 }
9638
9639 #[test]
9640 fn manual_merge_command_matches_the_configured_style() {
9641 let repo = Path::new("/repo");
9642 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9643
9644 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9645 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9646
9647 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9648 assert_eq!(
9649 squash,
9650 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9651 \"Merge magi run 0832 (candidate A)\""
9652 );
9653
9654 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9655 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9656 }
9657
9658 #[test]
9659 fn a_nudge_gets_a_quarter_of_the_budget() {
9660 assert_eq!(retry_budget(secs(1200), true), secs(300));
9662 assert_eq!(retry_budget(secs(3600), true), secs(900));
9663 }
9664
9665 #[test]
9666 fn a_resent_prompt_keeps_the_whole_budget() {
9667 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9670 assert_eq!(retry_budget(secs(60), false), secs(60));
9671 }
9672
9673 #[test]
9674 fn the_floor_never_exceeds_the_original_budget() {
9675 assert_eq!(retry_budget(secs(60), true), secs(60));
9679 assert_eq!(retry_budget(secs(480), true), secs(120));
9680 assert_eq!(retry_budget(secs(0), true), secs(0));
9681 }
9682
9683 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9684 agent::CommandEvidence {
9685 id: "item1".to_owned(),
9686 description: "cargo test".to_owned(),
9687 exit_code,
9688 result_summary: String::new(),
9689 source: "codex".to_owned(),
9690 }
9691 }
9692
9693 #[test]
9694 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9695 assert!(!has_unconfirmed_command(&[]));
9699 }
9700
9701 #[test]
9702 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9703 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9707 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9708 assert!(!has_unconfirmed_command(&[
9709 evidence(Some(0)),
9710 evidence(Some(101))
9711 ]));
9712 }
9713
9714 #[test]
9715 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9716 assert!(has_unconfirmed_command(&[
9717 evidence(Some(0)),
9718 evidence(None)
9719 ]));
9720 }
9721
9722 #[test]
9723 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9724 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9725 assert_eq!(
9726 verified_noop_claim(true, &[], text).as_deref(),
9727 Some("already fixed by b32cfc4, on main.")
9728 );
9729 }
9730
9731 #[test]
9732 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9733 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9736 assert!(verified_noop_claim(false, &[], text).is_none());
9737 }
9738
9739 #[test]
9740 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9741 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9742 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9743 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9745 }
9746
9747 #[test]
9748 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9749 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9750 }
9751
9752 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9755 runner.state.candidates = shape
9756 .iter()
9757 .enumerate()
9758 .map(|(i, &(empty, verified))| Candidate {
9759 index: i,
9760 label: (b'A' + i as u8) as char,
9761 agent: "sonnet".to_owned(),
9762 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9763 worktree: PathBuf::from(format!("/wt/{i}")),
9764 summary: String::new(),
9765 stat: String::new(),
9766 files: 0,
9767 commits: 0,
9768 empty,
9769 failed: None,
9770 verified_noop: verified.map(str::to_owned),
9771 duration_ms: 0,
9772 folded: false,
9773 })
9774 .collect();
9775 }
9776
9777 #[test]
9778 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9779 ask_test_home();
9780 let mut runner = runner_at(RunStatus::Implementing);
9781 set_candidates(
9782 &mut runner,
9783 &[
9784 (true, Some("already on main at b32cfc4")),
9785 (true, Some("same fix, see the existing test")),
9786 ],
9787 );
9788
9789 runner
9790 .after_implement()
9791 .expect("a verified no-op is not an error");
9792
9793 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9794 }
9795
9796 #[test]
9797 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9798 ask_test_home();
9799 let mut runner = runner_at(RunStatus::Implementing);
9800 set_candidates(
9804 &mut runner,
9805 &[(true, Some("already on main at b32cfc4")), (true, None)],
9806 );
9807
9808 let err = runner
9809 .after_implement()
9810 .expect_err("an unverified empty candidate must still fail the run");
9811
9812 assert!(
9813 err.to_string().contains("no candidate produced a change"),
9814 "{err}"
9815 );
9816 assert_eq!(runner.state.status, RunStatus::Failed);
9817 }
9818
9819 #[test]
9820 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9821 ask_test_home();
9822 let mut runner = runner_at(RunStatus::Implementing);
9823 set_candidates(&mut runner, &[(true, None), (true, None)]);
9824
9825 let err = runner
9826 .after_implement()
9827 .expect_err("no candidate declared anything; this is an ordinary failure");
9828
9829 assert!(
9830 err.to_string().contains("no candidate produced a change"),
9831 "{err}"
9832 );
9833 assert_eq!(runner.state.status, RunStatus::Failed);
9834 }
9835}