1use std::collections::{BTreeMap, BTreeSet};
20use std::path::{Path, PathBuf};
21use std::sync::atomic::{AtomicBool, Ordering};
22use std::sync::{Arc, Mutex};
23use std::time::{Duration, Instant};
24
25use anyhow::{Context as _, Result, bail};
26use jiff::Timestamp;
27use tokio::sync::Semaphore;
28
29use crate::advise;
30use crate::agent::{self, AgentOutput, Invocation, SeatState};
31use crate::ask;
32use crate::blind;
33use crate::bump;
34use crate::config::{
35 AgentSpec, Config, IncompleteReviewPolicy, LeakPolicy, MergeMode, MergeStyle, Prompts,
36 ResolvedRoles,
37};
38use crate::git;
39use crate::land;
40use crate::proc::Quiet as _;
41use crate::prompt::{
42 self, CandidateView, Lens, ReviewPatch, ReviewReconsiderCtx, ReviewSeatReport, Turn,
43};
44use crate::queue;
45use crate::run::{
46 BaseSync, Candidate, CommandOutcome, ContinuationOutcome, ContinuationRecord,
47 DeliberationRound, DeliberationTurn, E2eStatus, FixRecord, GateFixRecord, JobRecord, JobStatus,
48 Judgement, MergeOutcome, OperatorFixFinding, OperatorFixOutcome, OperatorFixRequest, QuotaLoss,
49 ReviewRecord, ReviewRevoteRecord, ReviewRound, RunState, RunStatus, Tally, VoteRecord, tail,
50 write_artifact,
51};
52use crate::verdict::{
53 self, FinalVote, Finding, FixReport, Position, Proposal, Ranking, Review, ReviewRevote,
54 ReviewVote, Severity,
55};
56
57const OUTPUT_TAIL: usize = 8_000;
59
60const EVENT_OUTPUT_TAIL: usize = 2_000;
63
64const LEASE_RELEASE_POLL: Duration = Duration::from_secs(1);
67
68const LEASE_RELEASE_MAX_WAIT: Duration = Duration::from_secs(30);
85
86pub(crate) const STAGNANT_LIMIT: usize = 2;
100
101const BASE_SYNC_ROUNDS: usize = 4;
114
115const MAX_FIX_CONTINUATIONS: usize = 2;
131
132#[derive(Clone)]
138struct SeatJob {
139 spec: AgentSpec,
140 seat: SeatState,
141 cwd: PathBuf,
142 prompt: String,
143 timeout: Duration,
144 allow_write: bool,
145 sessions: bool,
146 artifacts: PathBuf,
147 stem: String,
148}
149
150enum AgentOutcome {
162 Ok(AgentOutput),
164 Quota(AgentOutput),
166 Dropped(AgentOutput),
169 Failed(String),
171}
172
173#[derive(Debug, Clone, Default)]
200pub struct Pause(Arc<AtomicBool>, Arc<Mutex<Option<String>>>);
201
202impl Pause {
203 #[must_use]
205 pub fn new() -> Self {
206 Self::default()
207 }
208
209 pub fn park(&self) {
211 self.0.store(true, Ordering::SeqCst);
212 }
213
214 pub fn park_because(&self, reason: impl Into<String>) {
220 let mut reason_guard = self
221 .1
222 .lock()
223 .unwrap_or_else(std::sync::PoisonError::into_inner);
224 if reason_guard.is_none() {
225 *reason_guard = Some(reason.into());
226 }
227 drop(reason_guard);
228 self.park();
229 }
230
231 #[must_use]
233 pub fn parked(&self) -> bool {
234 self.0.load(Ordering::SeqCst)
235 }
236
237 #[must_use]
239 pub fn reason(&self) -> Option<String> {
240 self.1
241 .lock()
242 .unwrap_or_else(std::sync::PoisonError::into_inner)
243 .clone()
244 }
245}
246
247pub struct Runner {
249 pub state: RunState,
251 roles: ResolvedRoles,
252 sem: Arc<Semaphore>,
253 pause: Pause,
257 interrupt: Pause,
263}
264
265async fn sync_review_branch(repo: &Path, branch: &str, remote: &str) -> Result<()> {
294 let tracking = format!("{remote}/{branch}");
295 let fetched = git::fetch(repo, remote, branch).await;
296 let fresh = matches!(&fetched, Ok(o) if o.ok()) && git::rev_exists(repo, &tracking).await;
297 let local_exists = git::branch_exists(repo, branch).await?;
298 if !fresh {
299 if !local_exists {
300 bail!("no branch `{branch}` in {} or on {remote}", repo.display());
301 }
302 tracing::warn!(
303 "could not read {tracking}; reviewing the local `{branch}`, which may be stale"
304 );
305 return Ok(());
306 }
307 let remote_sha = git::rev_parse(repo, &tracking).await?;
308 if !local_exists {
309 git::git(repo, &["branch", branch, &tracking]).await?;
310 return Ok(());
311 }
312 let local_sha = git::rev_parse(repo, &format!("refs/heads/{branch}")).await?;
313 if local_sha == remote_sha || git::is_ancestor(repo, &remote_sha, &local_sha).await {
314 return Ok(());
315 }
316 if !git::is_ancestor(repo, &local_sha, &remote_sha).await {
317 let mb = git::git_raw(repo, &["merge-base", &local_sha, &remote_sha]).await?;
321 let placeholder = mb.ok()
322 && git::git_raw(repo, &["diff", "--quiet", &mb.stdout, &local_sha])
323 .await?
324 .ok();
325 if !placeholder {
326 bail!(
327 "local `{branch}` ({}) and {tracking} ({}) have diverged, so it is unclear \
328 which one to review; reconcile them, e.g. `git branch -f {branch} {tracking}` \
329 to review the pushed work, or push the local branch first",
330 short(&local_sha),
331 short(&remote_sha)
332 );
333 }
334 }
335 let out = git::git_raw(repo, &["branch", "-f", branch, &tracking]).await?;
336 if !out.ok() {
337 bail!(
338 "local `{branch}` ({}) is stale against {tracking} ({}) but git will not move it: {}",
339 short(&local_sha),
340 short(&remote_sha),
341 out.stderr
342 );
343 }
344 tracing::warn!(
345 "local `{branch}` was stale: fast-forwarded {} -> {}",
346 short(&local_sha),
347 short(&remote_sha)
348 );
349 Ok(())
350}
351
352async fn resolve_base(repo: &Path, base_branch: &str, remote: &str) -> Result<String> {
353 let tracking = format!("{remote}/{base_branch}");
354 let fetched = git::fetch(repo, remote, base_branch).await;
355 if let Ok(out) = &fetched
356 && out.ok()
357 && git::rev_exists(repo, &tracking).await
358 {
359 return git::rev_parse(repo, &tracking).await;
360 }
361 let why = match &fetched {
362 Ok(out) if !out.ok() => out.stderr.lines().next().unwrap_or("").to_owned(),
363 Ok(_) => format!("{remote} has no {base_branch}"),
364 Err(e) => e.to_string(),
365 };
366 tracing::warn!(
367 "could not read {tracking} ({why}); branching off the local \
368 {base_branch} instead, which may be behind"
369 );
370 git::rev_parse(repo, base_branch).await.with_context(|| {
371 format!(
372 "cannot resolve `{base_branch}`; set [merge] base in magi.toml to a \
373 branch that exists"
374 )
375 })
376}
377
378struct FixClaim {
395 path: PathBuf,
396}
397
398impl FixClaim {
399 fn acquire(dir: &Path) -> Result<Self> {
400 std::fs::create_dir_all(dir).with_context(|| format!("create {}", dir.display()))?;
401 let path = dir.join("fix.lock");
402 match Self::create(&path) {
403 Ok(claim) => Ok(claim),
404 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
405 if Self::reclaim_if_dead(&path) {
406 Self::create(&path).with_context(|| format!("lock {}", path.display()))
407 } else {
408 bail!(
409 "another `magi fix` is already running for this run ({} exists)",
410 path.display()
411 )
412 }
413 }
414 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
415 }
416 }
417
418 fn create(path: &Path) -> std::io::Result<Self> {
419 let mut f = std::fs::OpenOptions::new()
420 .write(true)
421 .create_new(true)
422 .open(path)?;
423 use std::io::Write as _;
424 writeln!(f, "{}", std::process::id())?;
426 Ok(Self {
427 path: path.to_owned(),
428 })
429 }
430
431 fn reclaim_if_dead(path: &Path) -> bool {
435 let dead = std::fs::read_to_string(path)
436 .ok()
437 .and_then(|body| body.trim().parse::<u32>().ok())
438 .is_some_and(|pid| !crate::proc::pid_alive(pid));
439 dead && std::fs::remove_file(path).is_ok()
440 }
441}
442
443impl Drop for FixClaim {
444 fn drop(&mut self) {
445 let _ = std::fs::remove_file(&self.path);
446 }
447}
448
449impl Runner {
450 pub async fn start(repo: &Path, instruction: String, config: Config) -> Result<Self> {
452 let repo = git::toplevel(repo).await?;
453 let missing = agent::missing_programs(&config.agents);
454 if !missing.is_empty() {
455 bail!(
456 "these agent programs are not on PATH: {}. Fix the roster in \
457 magi.toml or install them.",
458 missing.join(", ")
459 );
460 }
461 let base_branch = match config.merge.base.clone() {
462 Some(b) => b,
463 None => git::current_branch(&repo)
464 .await?
465 .context("HEAD is detached; set [merge] base in magi.toml")?,
466 };
467 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
468 if !git::is_clean(&repo).await? {
472 tracing::warn!(
473 "{} has uncommitted changes; they are not part of this run, \
474 which branches off {base_branch} ({})",
475 repo.display(),
476 &base_commit[..base_commit.len().min(8)]
477 );
478 }
479 let roles = config.resolve_roles()?;
480 let max_parallel = config.graph.max_parallel.max(1);
481 let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
482 state.event("start", format!("run {} created", state.id));
483 state.save()?;
484 Ok(Self {
485 state,
486 roles,
487 sem: Arc::new(Semaphore::new(max_parallel)),
488 pause: Pause::new(),
489 interrupt: Pause::new(),
490 })
491 }
492
493 pub async fn review(repo: &Path, branch: &str, config: Config) -> Result<Self> {
507 Self::review_taking_over(repo, branch, config, None).await
508 }
509
510 pub async fn review_taking_over(
516 repo: &Path,
517 branch: &str,
518 config: Config,
519 takeover: Option<crate::handover::Takeover>,
520 ) -> Result<Self> {
521 let repo = git::toplevel(repo).await?;
522 let missing = agent::missing_programs(&config.agents);
523 if !missing.is_empty() {
524 bail!(
525 "these agent programs are not on PATH: {}. Fix the roster in \
526 magi.toml or install them.",
527 missing.join(", ")
528 );
529 }
530 let base_branch = match config.merge.base.clone() {
531 Some(b) => b,
532 None => git::current_branch(&repo)
533 .await?
534 .context("HEAD is detached; set [merge] base in magi.toml")?,
535 };
536 if base_branch == branch {
537 bail!("`{branch}` is the base branch; there is nothing to review against");
538 }
539 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
540
541 let roles = config.resolve_roles()?;
542 let max_parallel = config.graph.max_parallel.max(1);
543 let mut state = RunState::new(
544 repo.clone(),
545 base_branch,
546 base_commit.clone(),
547 String::new(),
548 config,
549 );
550
551 let released = match &takeover {
556 Some(takeover) => crate::handover::release(&repo, branch, &state.id, takeover).await?,
557 None => None,
558 };
559 if let Some(released) = &released {
560 state.event(
561 "release",
562 format!(
563 "took `{branch}` over from run {}: its worktree was released",
564 crate::run::short_of(&released.old_id)
565 ),
566 );
567 }
568 let opened =
569 Self::open_review(&repo, branch, state, roles, max_parallel, base_commit).await;
570 if opened.is_err()
571 && let Some(released) = &released
572 {
573 released.restore(&repo, branch).await;
574 }
575 opened
576 }
577
578 async fn open_review(
581 repo: &Path,
582 branch: &str,
583 mut state: RunState,
584 roles: ResolvedRoles,
585 max_parallel: usize,
586 base_commit: String,
587 ) -> Result<Self> {
588 sync_review_branch(repo, branch, &state.config.merge.remote).await?;
589 let log = git::log_oneline(repo, &base_commit, branch)
592 .await
593 .unwrap_or_default();
594 let instruction = format!(
595 "Review the work already on branch `{branch}`. There is no task \
596 statement: what the change claims to do is whatever its commits \
597 say.\n\n{}",
598 if log.trim().is_empty() {
599 "(no commit messages)"
600 } else {
601 log.trim()
602 }
603 );
604 state.instruction = instruction;
605
606 let worktree = state.worktree_root().join("under-review");
609 if let Some(parent) = worktree.parent() {
610 tokio::fs::create_dir_all(parent).await.ok();
611 }
612 let path = worktree.to_string_lossy().to_string();
613 git::git(repo, &["worktree", "add", &path, branch])
614 .await
615 .with_context(|| {
616 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
617 })?;
618
619 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
620 .await
621 .unwrap_or(0);
622 if commits == 0 {
623 git::worktree_remove(repo, &worktree).await.ok();
624 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
625 }
626 let files = git::changed_files(&worktree, &base_commit, "HEAD")
627 .await
628 .map(|f| f.len())
629 .unwrap_or(0);
630 if files == 0
631 && let (Ok(head_tree), Ok(base_tree)) = (
632 git::tree_of(&worktree, "HEAD").await,
633 git::tree_of(&worktree, &base_commit).await,
634 )
635 && head_tree == base_tree
636 {
637 let head = git::rev_parse(&worktree, "HEAD").await.unwrap_or_default();
638 git::worktree_remove(repo, &worktree).await.ok();
639 bail!(
640 "`{branch}` at {} has a tree identical to base {}; this usually means \
641 the branch ref is stale (check `git rev-parse refs/heads/{branch}` \
642 against `{}/{branch}`) rather than an empty change",
643 short(&head),
644 short(&base_commit),
645 state.config.merge.remote
646 );
647 }
648 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
649 .await
650 .unwrap_or_default();
651
652 state.candidates.push(Candidate {
653 index: 0,
654 label: 'A',
655 agent: "(existing branch)".to_owned(),
658 branch: branch.to_owned(),
659 worktree,
660 summary: String::new(),
661 stat,
662 files,
663 commits,
664 empty: false,
665 failed: None,
666 verified_noop: None,
667 duration_ms: 0,
668 folded: false,
669 });
670 state.tally = Some(Tally {
671 first_choice: BTreeMap::from([('A', 0)]),
672 borda: BTreeMap::new(),
673 winner: 'A',
674 rankings: 0,
675 unanimous_initial: false,
676 deliberated: false,
677 changed_votes: 0,
678 unanimous_final: false,
679 tie_break: None,
680 judges: 0,
684 present: 0,
685 quorum: 0,
686 met_quorum: true,
687 uncontested: Some("review-only run: nothing competed".to_owned()),
688 });
689 state.status = RunStatus::Reviewing;
690 state.event(
691 "start",
692 format!(
693 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
694 state.id
695 ),
696 );
697 state.save()?;
698 Ok(Self {
699 state,
700 roles,
701 sem: Arc::new(Semaphore::new(max_parallel)),
702 pause: Pause::new(),
703 interrupt: Pause::new(),
704 })
705 }
706
707 pub fn resume(id: &str) -> Result<Self> {
709 let state = RunState::load(id)?;
710 if let Some(to) = &state.released_to {
711 bail!(
712 "run {} cannot be resumed: its worktree was released to run {}",
713 state.short(),
714 crate::run::short_of(to)
715 );
716 }
717 let roles = state.config.resolve_roles()?;
718 let max_parallel = state.config.graph.max_parallel.max(1);
719 Ok(Self {
720 state,
721 roles,
722 sem: Arc::new(Semaphore::new(max_parallel)),
723 pause: Pause::new(),
724 interrupt: Pause::new(),
725 })
726 }
727
728 pub async fn execute(&mut self) -> Result<()> {
735 let result = self.execute_graph().await;
736 let ended = if result.is_err() {
737 Some(crate::notices::run_stopped(&self.state.id, &self.state))
738 } else {
739 crate::notices::run_ended(&self.state)
740 };
741 if let Some(notice) = ended {
742 crate::notices::raise(notice);
743 }
744 result
745 }
746
747 async fn execute_graph(&mut self) -> Result<()> {
748 self.state.parked = false;
753 self.state.clear_active();
760 if let Ok(disk) = RunState::load(&self.state.id)
778 && let Some(to) = &disk.released_to
779 {
780 bail!(
781 "run {} cannot continue: its worktree was released to run {}",
782 self.state.short(),
783 crate::run::short_of(to)
784 );
785 }
786 let pid = std::process::id();
787 self.state.driver_pid = Some(pid);
788 self.state.driver_started_at = crate::proc::process_started_at(pid);
789 self.state.save()?;
790 if self.state.status == RunStatus::Stalled {
803 if self.recover_stall().await? {
804 self.finish_after_tally().await?;
805 } else {
806 self.state.save()?;
808 }
809 return Ok(());
810 }
811 if self.state.status == RunStatus::Landing {
821 self.run_land().await?;
822 self.settle_questions();
827 return Ok(());
828 }
829 self.prep().await?;
830 if self.park_here()? {
831 return Ok(());
832 }
833 self.advise().await?;
834 if self.park_here()? {
835 return Ok(());
836 }
837 self.implement().await?;
838 if self.park_here()? {
839 return Ok(());
840 }
841 if self.state.status == RunStatus::VerifiedNoop {
845 return Ok(());
846 }
847 self.judge().await?;
848 if self.park_here()? {
849 return Ok(());
850 }
851 self.deliberate().await?;
852 if self.park_here()? {
853 return Ok(());
854 }
855 self.vote().await?;
856 if self.park_here()? {
857 return Ok(());
858 }
859 self.tally()?;
860 if self.state.status == RunStatus::Stalled {
865 self.state.save()?;
869 return Ok(());
870 }
871 self.finish_after_tally().await?;
872 Ok(())
873 }
874
875 fn park_here(&mut self) -> Result<bool> {
882 if !self.pause.parked() && !self.interrupt.parked() {
887 return Ok(false);
888 }
889 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
890 Some(reason) => format!(
891 "parked after `{}` ({reason}) — resume to carry on from here",
892 self.state.status.as_str()
893 ),
894 None => format!(
895 "parked after `{}` — resume to carry on from here",
896 self.state.status.as_str()
897 ),
898 };
899 self.state.event("park", why);
900 self.state.parked = true;
901 self.state.save()?;
902 Ok(true)
903 }
904
905 pub fn on_pause(&mut self, pause: Pause) {
907 self.pause = pause;
908 }
909
910 pub fn watch_interrupt(&mut self, pause: Pause) {
916 self.interrupt = pause;
917 }
918
919 fn settle_questions(&mut self) {
939 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
940 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
941 }
942 }
943
944 async fn finish_after_tally(&mut self) -> Result<()> {
947 self.fold_losers().await?;
948 self.sync_to_base().await?;
953 self.review_loop().await?;
954 self.sync_to_base().await?;
955 self.gate().await?;
956 self.merge().await?;
957 self.state.save()?;
958 Ok(())
959 }
960
961 async fn prep(&mut self) -> Result<()> {
964 if !self.state.candidates.is_empty() {
965 return Ok(());
966 }
967 self.state.status = RunStatus::Prep;
968 let repo = self.state.repo.clone();
969 let base = self.state.base_commit.clone();
970 let root = self.state.worktree_root();
971 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
972
973 let hooks_dir = self.state.dir().join("hooks");
976 if self.state.config.blind.commit_msg_hook {
977 std::fs::create_dir_all(&hooks_dir)
978 .with_context(|| format!("create {}", hooks_dir.display()))?;
979 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
980 let path = hooks_dir.join("commit-msg");
981 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
982 make_executable(&path)?;
983 git::acquire_worktree_config(&repo).await?;
991 self.state.enabled_worktree_config = true;
992 }
993
994 for (index, (spec, label)) in self
995 .roles
996 .implementers
997 .clone()
998 .into_iter()
999 .zip(labels)
1000 .enumerate()
1001 {
1002 let branch = self.state.branch_for(label);
1003 let worktree = root.join(format!("cand-{label}"));
1004 git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
1005 if self.state.config.blind.commit_msg_hook {
1006 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
1007 }
1008 git::local_exclude(&worktree, "/.magi/").await?;
1009 self.state.candidates.push(Candidate {
1010 index,
1011 label,
1012 agent: spec.id.clone(),
1013 branch,
1014 worktree,
1015 summary: String::new(),
1016 stat: String::new(),
1017 files: 0,
1018 commits: 0,
1019 empty: false,
1020 failed: None,
1021 verified_noop: None,
1022 duration_ms: 0,
1023 folded: false,
1024 });
1025 }
1026
1027 for j in 1..=self.roles.judges.len() {
1028 let wt = root.join(format!("judge-{j}"));
1029 if !wt.exists() {
1030 git::worktree_add_detached(&repo, &wt, &base).await?;
1031 }
1032 }
1033
1034 if self.state.config.graph.advise {
1042 for k in 1..=self.state.config.graph.advisors {
1043 let wt = root.join(format!("advisor-{k}"));
1044 if !wt.exists() {
1045 git::worktree_add_detached(&repo, &wt, &base).await?;
1046 }
1047 }
1048 }
1049
1050 let authors: Vec<&str> = self
1055 .roles
1056 .implementers
1057 .iter()
1058 .map(|a| a.id.as_str())
1059 .collect();
1060 let overlap: Vec<String> = self
1061 .roles
1062 .judges
1063 .iter()
1064 .enumerate()
1065 .filter(|(_, j)| authors.contains(&j.id.as_str()))
1066 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
1067 .collect();
1068 if !overlap.is_empty() {
1069 let note = format!(
1070 "{} also authored a candidate; blind, but the panel is less \
1071 independent than {} distinct agents would be",
1072 overlap.join(", "),
1073 self.roles.judges.len()
1074 );
1075 self.state.event("prep", note);
1076 }
1077
1078 self.state.event(
1079 "prep",
1080 format!(
1081 "{} candidates, {} judges, base {} ({})",
1082 self.state.candidates.len(),
1083 self.roles.judges.len(),
1084 &self.state.base_commit[..7.min(self.state.base_commit.len())],
1085 self.state.base_branch
1086 ),
1087 );
1088 self.state.status = RunStatus::Implementing;
1089 self.state.save()?;
1090 Ok(())
1091 }
1092
1093 async fn advise(&mut self) -> Result<()> {
1126 let implement_untouched = self
1127 .state
1128 .candidates
1129 .iter()
1130 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
1131 if !self.state.config.graph.advise || self.state.advise_attempted {
1132 return Ok(());
1133 }
1134 if !implement_untouched {
1135 self.state.event(
1136 "advise",
1137 "skipping the design-deliberation stage: at least one \
1138 candidate already shows implementation progress, so this \
1139 run is past the point the stage exists to run before"
1140 .to_owned(),
1141 );
1142 self.state.advise_attempted = true;
1143 self.state.save()?;
1144 return Ok(());
1145 }
1146 let run_id = self.state.id.clone();
1147 let prompts = self.state.config.prompts.clone();
1148 let instruction = self.state.instruction.clone();
1149 let language = self.state.config.graph.language.clone();
1150 let root = self.state.worktree_root();
1151 let n = self.state.config.graph.advisors;
1152 let where_recorded = self.state.dir().join("run.json");
1153
1154 let seats = match self.state.config.advisors() {
1155 Ok(seats) if !seats.is_empty() => seats,
1156 Ok(_) => {
1157 self.state.event(
1158 "advise",
1159 format!(
1160 "[graph] advisors is 0; skipping the design-deliberation \
1161 stage and continuing without a synthesis brief (see {})",
1162 where_recorded.display()
1163 ),
1164 );
1165 self.state.advise_attempted = true;
1166 self.state.save()?;
1167 return Ok(());
1168 }
1169 Err(e) => {
1170 self.state.event(
1171 "advise",
1172 format!(
1173 "could not resolve advisor seats ({e:#}); continuing \
1174 without a design-deliberation brief (see {})",
1175 where_recorded.display()
1176 ),
1177 );
1178 self.state.advise_attempted = true;
1179 self.state.save()?;
1180 return Ok(());
1181 }
1182 };
1183
1184 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1185 let artifacts = agent::artifacts_dir(&self.state.dir());
1186 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1187
1188 let mut jobs = Vec::new();
1189 for (i, spec) in seats.iter().cloned().enumerate() {
1190 let seat_key = format!("advisor-{}", i + 1);
1191 let seat = self.seat(&seat_key, &spec.id);
1192 jobs.push(SeatJob {
1193 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1194 spec,
1195 seat,
1196 cwd: worktrees[i % worktrees.len()].clone(),
1197 timeout,
1198 allow_write: false,
1199 sessions: false,
1200 artifacts: artifacts.clone(),
1201 stem: seat_key,
1202 });
1203 }
1204
1205 self.state.event(
1206 "advise",
1207 format!(
1208 "{} advisor seat(s) sketching a design in parallel",
1209 jobs.len()
1210 ),
1211 );
1212 let mut quota_losses = Vec::new();
1213 let cache = self.state.config.cache_dir();
1214 let ctx = WaveCtx {
1215 run: &run_id,
1216 node: "advise",
1217 prompts: &prompts,
1218 cache: cache.as_deref(),
1219 round: None,
1220 };
1221 let results = ask_json_wave::<Proposal>(
1222 jobs,
1223 Arc::clone(&self.sem),
1224 self.state.config.graph.retries,
1225 &ctx,
1226 &mut quota_losses,
1227 &mut self.state,
1228 &|p: &Proposal| p.validate(),
1229 )
1230 .await;
1231 self.state.quota.extend(quota_losses);
1232
1233 let mut records = Vec::with_capacity(results.len());
1234 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1235 let agent_id = seat.agent.clone();
1236 self.state.seats.insert(seat.key.clone(), seat);
1237 match res {
1238 Ok((proposal, out)) => {
1239 self.state
1240 .event("advise", format!("advisor-{} proposed a design", i + 1));
1241 records.push(advise::AdvisorRecord::proposed(
1242 i + 1,
1243 agent_id,
1244 proposal,
1245 out.duration_ms,
1246 ));
1247 }
1248 Err(e) => {
1249 self.state.event(
1250 "advise",
1251 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1252 );
1253 records.push(advise::AdvisorRecord::failed(
1254 i + 1,
1255 agent_id,
1256 e.to_string(),
1257 ));
1258 }
1259 }
1260 }
1261
1262 let mut advice = advise::Advice {
1263 records,
1264 synthesis: None,
1265 };
1266 if advice.proposals().is_empty() {
1267 self.state.event(
1268 "advise",
1269 "no advisor produced a usable proposal; continuing without a \
1270 synthesis brief"
1271 .to_owned(),
1272 );
1273 } else {
1274 match self
1275 .synthesize_brief(
1276 &advice,
1277 &instruction,
1278 &language,
1279 &worktrees[0],
1280 &artifacts,
1281 &run_id,
1282 &prompts,
1283 cache.as_deref(),
1284 )
1285 .await
1286 {
1287 Ok(Some(text)) => {
1288 self.state.event(
1289 "advise",
1290 "synthesized a design brief for the implementer".to_owned(),
1291 );
1292 advice.synthesis = Some(text);
1293 }
1294 Ok(None) => {
1295 self.state.event(
1296 "advise",
1297 "the synthesis seat produced nothing usable; continuing \
1298 without a design brief"
1299 .to_owned(),
1300 );
1301 }
1302 Err(e) => {
1303 self.state.event(
1304 "advise",
1305 format!("could not synthesize a design brief: {e:#}"),
1306 );
1307 }
1308 }
1309 }
1310 advise::apply_reflection(&mut advice);
1311
1312 self.state.advice = Some(advice);
1313 self.state.advise_attempted = true;
1314 self.state.save()?;
1315 Ok(())
1316 }
1317
1318 #[allow(clippy::too_many_arguments)]
1330 async fn synthesize_brief(
1331 &mut self,
1332 advice: &advise::Advice,
1333 instruction: &str,
1334 language: &str,
1335 cwd: &Path,
1336 artifacts: &Path,
1337 run_id: &str,
1338 prompts: &Prompts,
1339 cache: Option<&Path>,
1340 ) -> Result<Option<String>> {
1341 let want = self.state.config.roles.synthesizer.as_deref();
1342 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1343 let mut seat = self.seat("advise-synthesis", &spec.id);
1344 let proposals = advice.proposals();
1345 let mut prompt = prompt::with_overlay(
1346 prompt::synthesize_brief(instruction, &proposals, language),
1347 prompts.overlay("advise"),
1348 );
1349 if cache.is_some() {
1350 prompt.push('\n');
1355 prompt.push_str(&prompt::build_cache_note("advise", false));
1356 }
1357 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1358 let out = agent::invoke(
1359 &spec,
1360 &mut seat,
1361 &Invocation {
1362 cwd,
1363 prompt: &prompt,
1364 timeout,
1365 allow_write: false,
1366 sessions: false,
1367 artifacts,
1368 stem: "advise-synthesis",
1369 run: run_id,
1370 node: "advise",
1371 cache_dir: None,
1372 attachments: &[],
1373 },
1374 )
1375 .await?;
1376 self.state.seats.insert(seat.key.clone(), seat);
1377 if !out.usable() {
1378 return Ok(None);
1379 }
1380 let text =
1381 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1382 Ok((!text.trim().is_empty()).then_some(text))
1383 }
1384
1385 async fn implement(&mut self) -> Result<()> {
1388 let run_id = self.state.id.clone();
1393 let prompts = self.state.config.prompts.clone();
1394 let todo: Vec<usize> = self
1395 .state
1396 .candidates
1397 .iter()
1398 .enumerate()
1399 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1400 .map(|(i, _)| i)
1401 .collect();
1402 if todo.is_empty() {
1403 return self.after_implement();
1404 }
1405 self.state.status = RunStatus::Implementing;
1406
1407 let language = self.state.config.graph.language.clone();
1408 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1409 let sessions = self.state.config.graph.sessions;
1410 let artifacts = agent::artifacts_dir(&self.state.dir());
1411 let brief = self
1415 .state
1416 .advice
1417 .as_ref()
1418 .and_then(|a| a.synthesis.as_deref())
1419 .map(str::to_owned);
1420
1421 let mut jobs = Vec::new();
1422 for &i in &todo {
1423 let (index, label, worktree) = {
1424 let c = &self.state.candidates[i];
1425 (c.index, c.label, c.worktree.clone())
1426 };
1427 let spec = self.roles.implementers[index].clone();
1428 let seat_key = format!("impl-{label}");
1429 let seat = self.seat(&seat_key, &spec.id);
1430 let instruction = self.state.instruction.clone();
1431 jobs.push(SeatJob {
1432 spec,
1433 seat,
1434 prompt: prompt::implement(
1435 &instruction,
1436 &worktree.to_string_lossy(),
1437 &language,
1438 brief.as_deref(),
1439 ),
1440 cwd: worktree,
1441 timeout,
1442 allow_write: true,
1443 sessions,
1444 artifacts: artifacts.clone(),
1445 stem: format!("impl-{label}"),
1446 });
1447 }
1448
1449 self.state.event(
1450 "implement",
1451 format!("{} candidates in parallel", jobs.len()),
1452 );
1453 let mut sent = jobs.clone();
1459 let cache = self.state.config.cache_dir();
1460 let ctx = WaveCtx {
1461 run: &run_id,
1462 node: "implement",
1463 prompts: &prompts,
1464 cache: cache.as_deref(),
1465 round: None,
1466 };
1467 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1468 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1469 .await;
1470 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1471 .await;
1472 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1473 .await;
1474
1475 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1476 let seat_key = seat.key.clone();
1477 let agent = seat.agent.clone();
1486 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1487 self.state.seats.insert(seat.key.clone(), seat);
1488 let label = self.state.candidates[i].label;
1489 let worktree = self.state.candidates[i].worktree.clone();
1490 let base = self.state.base_commit.clone();
1491
1492 let (summary, duration, failed, verified_claim) = match out {
1493 AgentOutcome::Ok(o) => {
1494 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1495 let failed = (!o.usable()).then(|| {
1496 if o.timed_out {
1497 "agent timed out".to_owned()
1498 } else {
1499 format!("agent exited with {:?}", o.exit_code)
1500 }
1501 });
1502 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1503 (text, o.duration_ms, failed, verified_claim)
1504 }
1505 AgentOutcome::Dropped(o) => {
1511 let why = o
1512 .dropped
1513 .as_ref()
1514 .map(|d| d.why.as_str())
1515 .unwrap_or("the CLI ended the stream without delivering its answer");
1516 (
1517 String::new(),
1518 o.duration_ms,
1519 Some(format!("the CLI dropped the stream ({why})")),
1520 None,
1521 )
1522 }
1523 AgentOutcome::Quota(o) => {
1524 self.state.quota.push(QuotaLoss {
1525 seat: seat_key,
1526 node: "implement".to_owned(),
1527 at: Timestamp::now(),
1528 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1529 });
1530 (
1531 String::new(),
1532 o.duration_ms,
1533 Some("rate limited (quota); produced no change".to_owned()),
1534 None,
1535 )
1536 }
1537 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1538 };
1539
1540 let rescued = match git::rescue_commit(
1543 &worktree,
1544 &format!("magi: candidate {label} (uncommitted work)"),
1545 )
1546 .await
1547 {
1548 Ok(r) => {
1549 self.state.note_withheld("implement", &r.withheld);
1550 r.committed
1551 }
1552 Err(_) => false,
1553 };
1554 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1555 .await
1556 .unwrap_or(0);
1557 let patch = git::diff(&worktree, &base, "HEAD")
1558 .await
1559 .unwrap_or_default();
1560 let stat = git::diff_stat(&worktree, &base, "HEAD")
1561 .await
1562 .unwrap_or_default();
1563 let files = git::changed_files(&worktree, &base, "HEAD")
1564 .await
1565 .map(|f| f.len())
1566 .unwrap_or(0);
1567 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1568
1569 let c = &mut self.state.candidates[i];
1570 if !exhausted_the_fallback_chain {
1571 c.agent = agent;
1572 }
1573 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1574 c.stat = stat;
1575 c.files = files;
1576 c.commits = commits;
1577 c.duration_ms = duration;
1578 c.empty = commits == 0 || patch.trim().is_empty();
1579 c.failed = match failed {
1582 Some(_) if c.empty => failed,
1583 _ => None,
1584 };
1585 c.verified_noop = if c.empty { verified_claim } else { None };
1590 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1591 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1592 (None, true, Some(_), _) => {
1593 format!("candidate {label}: no change produced (agent-verified no-op)")
1594 }
1595 (None, true, None, _) => format!("candidate {label}: no change produced"),
1596 (None, false, _, true) => {
1597 format!(
1598 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1599 )
1600 }
1601 (None, false, _, false) => {
1602 format!("candidate {label}: {files} files, {commits} commits")
1603 }
1604 };
1605 self.state.event("implement", note);
1606 self.state.save()?;
1607 }
1608
1609 self.after_implement()
1610 }
1611
1612 async fn resume_undelivered(
1640 &mut self,
1641 results: &mut [(usize, SeatState, AgentOutcome)],
1642 sent: &[SeatJob],
1643 prompts: &Prompts,
1644 run_id: &str,
1645 ) {
1646 for (wi, seat, out) in results.iter_mut() {
1647 let Some(dropped) = (match &*out {
1648 AgentOutcome::Dropped(o) => o.dropped.clone(),
1649 _ => None,
1650 }) else {
1651 continue;
1652 };
1653 let Some(job) = sent.get(*wi) else { continue };
1654 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1656 self.state.event(
1657 "implement",
1658 format!(
1659 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1660 work is in the tree",
1661 seat.key, dropped.output_tokens, dropped.why
1662 ),
1663 );
1664 continue;
1665 }
1666 if !has_context(&job.spec, seat, job.sessions) {
1674 self.state.event(
1675 "implement",
1676 format!(
1677 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1678 is no session left to resume",
1679 seat.key, dropped.output_tokens, dropped.why
1680 ),
1681 );
1682 continue;
1683 }
1684 self.state.event(
1685 "implement",
1686 format!(
1687 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1688 conversation",
1689 seat.key, dropped.output_tokens, dropped.why
1690 ),
1691 );
1692 let mut retry = job.clone();
1693 retry.seat = seat.clone();
1694 retry.prompt = prompt::resume_after_drop(&dropped.why);
1695 retry.timeout = retry_budget(job.timeout, true);
1696 retry.stem = format!("{}-resume", job.stem);
1697 let cache = self.state.config.cache_dir();
1698 let ctx = WaveCtx {
1699 run: run_id,
1700 node: "implement",
1701 prompts,
1702 cache: cache.as_deref(),
1703 round: None,
1704 };
1705 let (resumed_seat, resumed) =
1706 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1707 *seat = resumed_seat;
1708 *out = resumed;
1709 }
1710 }
1711
1712 async fn resume_quota_losses(
1774 &mut self,
1775 results: &mut [(usize, SeatState, AgentOutcome)],
1776 sent: &mut [SeatJob],
1777 prompts: &Prompts,
1778 run_id: &str,
1779 ) {
1780 let instruction = self.state.instruction.clone();
1781 let language = self.state.config.graph.language.clone();
1782 let brief = self
1783 .state
1784 .advice
1785 .as_ref()
1786 .and_then(|a| a.synthesis.as_deref())
1787 .map(str::to_owned);
1788 for (wi, seat, out) in results.iter_mut() {
1789 let Some(job) = sent.get_mut(*wi) else {
1790 continue;
1791 };
1792 let start = self
1797 .roles
1798 .implementer_roster
1799 .iter()
1800 .position(|s| s.id == job.spec.id)
1801 .unwrap_or(0);
1802 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1803 let mut fallback_attempt = 0usize;
1804 while matches!(&*out, AgentOutcome::Quota(_)) {
1805 let Some(next) =
1806 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1807 .cloned()
1808 else {
1809 break;
1810 };
1811 tried.insert(next.id.clone());
1812 fallback_attempt += 1;
1813
1814 if let Ok(r) = git::rescue_commit(
1815 &job.cwd,
1816 &format!(
1817 "magi: candidate {} (uncommitted work before quota fallback)",
1818 seat.key
1819 ),
1820 )
1821 .await
1822 {
1823 self.state.note_withheld("implement", &r.withheld);
1824 }
1825
1826 self.state.event(
1827 "implement",
1828 format!(
1829 "{}: rate limited (quota) on {}; retrying with {}",
1830 seat.key, seat.agent, next.id
1831 ),
1832 );
1833
1834 let new_seat = self.seat(&seat.key, &next.id);
1835 job.spec = next.clone();
1843 let mut retry = job.clone();
1844 retry.seat = new_seat;
1845 retry.prompt = prompt::implement(
1846 &instruction,
1847 &job.cwd.to_string_lossy(),
1848 &language,
1849 brief.as_deref(),
1850 );
1851 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1852 let cache = self.state.config.cache_dir();
1853 let ctx = WaveCtx {
1854 run: run_id,
1855 node: "implement",
1856 prompts,
1857 cache: cache.as_deref(),
1858 round: None,
1859 };
1860 let (fallback_seat, fallback_out) = run_one(
1861 retry,
1862 Arc::clone(&self.sem),
1863 &ctx,
1864 &mut self.state,
1865 fallback_attempt,
1866 )
1867 .await;
1868 *seat = fallback_seat;
1869 *out = fallback_out;
1870 }
1871 }
1872 }
1873
1874 async fn resume_unconfirmed_commands(
1898 &mut self,
1899 results: &mut [(usize, SeatState, AgentOutcome)],
1900 sent: &[SeatJob],
1901 prompts: &Prompts,
1902 run_id: &str,
1903 ) {
1904 for (wi, seat, out) in results.iter_mut() {
1905 let AgentOutcome::Ok(o) = &*out else {
1906 continue;
1907 };
1908 if !has_unconfirmed_command(&o.commands) {
1909 continue;
1910 }
1911 let Some(job) = sent.get(*wi) else { continue };
1912 if !has_context(&job.spec, seat, job.sessions) {
1913 self.state.event(
1914 "implement",
1915 format!(
1916 "{}: the reply named a command whose own CLI never confirmed the exit \
1917 status of, but there is no session left to resume",
1918 seat.key
1919 ),
1920 );
1921 continue;
1922 }
1923 self.state.event(
1924 "implement",
1925 format!(
1926 "{}: the reply named a command whose own CLI never confirmed the exit \
1927 status of; resuming the conversation",
1928 seat.key
1929 ),
1930 );
1931 let mut retry = job.clone();
1932 retry.seat = seat.clone();
1933 retry.prompt = prompt::resume_incomplete(
1934 "a command in your last reply had no confirmed exit status",
1935 );
1936 retry.timeout = retry_budget(job.timeout, true);
1937 retry.stem = format!("{}-confirm", job.stem);
1938 let cache = self.state.config.cache_dir();
1939 let ctx = WaveCtx {
1940 run: run_id,
1941 node: "implement",
1942 prompts,
1943 cache: cache.as_deref(),
1944 round: None,
1945 };
1946 let (resumed_seat, resumed) =
1947 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1948 *seat = resumed_seat;
1949 *out = resumed;
1950 }
1951 }
1952
1953 async fn continue_fix_report(
1974 &mut self,
1975 mut seat: SeatState,
1976 parse_err: String,
1977 job: &SeatJob,
1978 prompts: &Prompts,
1979 run_id: &str,
1980 round: usize,
1981 ) -> (
1982 SeatState,
1983 Option<FixReport>,
1984 Option<String>,
1985 ContinuationRecord,
1986 ) {
1987 let mut last_err = parse_err;
1988 let mut cumulative_wait_ms = 0u64;
1989 let mut attempts = 0usize;
1990 loop {
1991 if !has_context(&job.spec, &seat, job.sessions) {
1992 self.state.event(
1993 "fix",
1994 format!(
1995 "round {round}: fixer's reply had no adoption report ({last_err}); no \
1996 session left to resume into"
1997 ),
1998 );
1999 let outcome = if attempts == 0 {
2000 ContinuationOutcome::NoSession
2001 } else {
2002 ContinuationOutcome::Exhausted
2003 };
2004 return (
2005 seat,
2006 None,
2007 Some(format!("unparsable fix report: {last_err}")),
2008 ContinuationRecord {
2009 attempts,
2010 cumulative_wait_ms,
2011 outcome,
2012 },
2013 );
2014 }
2015 if attempts >= MAX_FIX_CONTINUATIONS {
2016 self.state.event(
2017 "fix",
2018 format!(
2019 "round {round}: fixer's reply still had no adoption report after \
2020 {attempts} continuation(s) ({last_err}); giving up"
2021 ),
2022 );
2023 return (
2024 seat,
2025 None,
2026 Some(format!(
2027 "unparsable fix report after {attempts} continuation(s): {last_err}"
2028 )),
2029 ContinuationRecord {
2030 attempts,
2031 cumulative_wait_ms,
2032 outcome: ContinuationOutcome::Exhausted,
2033 },
2034 );
2035 }
2036 attempts += 1;
2037 self.state.event(
2038 "fix",
2039 format!(
2040 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
2041 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
2042 ),
2043 );
2044 let mut retry = job.clone();
2045 retry.seat = seat.clone();
2046 retry.prompt = prompt::resume_incomplete(&last_err);
2047 retry.timeout = retry_budget(job.timeout, true);
2048 retry.stem = format!("{}-continue{attempts}", job.stem);
2049 let cache = self.state.config.cache_dir();
2050 let ctx = WaveCtx {
2051 run: run_id,
2052 node: "fix",
2053 prompts,
2054 cache: cache.as_deref(),
2055 round: Some(round),
2056 };
2057 let (resumed_seat, resumed_out) = run_one(
2058 retry,
2059 Arc::clone(&self.sem),
2060 &ctx,
2061 &mut self.state,
2062 attempts,
2063 )
2064 .await;
2065 seat = resumed_seat;
2066 match resumed_out {
2067 AgentOutcome::Ok(o) => {
2068 cumulative_wait_ms += o.duration_ms;
2069 match verdict::extract_json::<FixReport>(&o.text) {
2070 Ok(report) if !has_unconfirmed_command(&o.commands) => {
2071 self.state.event(
2072 "fix",
2073 format!(
2074 "round {round}: fixer's adoption report recovered after \
2075 {attempts} continuation(s)"
2076 ),
2077 );
2078 return (
2079 seat,
2080 Some(report),
2081 None,
2082 ContinuationRecord {
2083 attempts,
2084 cumulative_wait_ms,
2085 outcome: ContinuationOutcome::Resumed,
2086 },
2087 );
2088 }
2089 Ok(_) => {
2097 last_err = "the reply parsed, but it reported a command whose own CLI \
2098 never confirmed an exit status"
2099 .to_owned();
2100 }
2101 Err(e) => last_err = e.to_string(),
2102 }
2103 }
2104 AgentOutcome::Quota(o) => {
2105 cumulative_wait_ms += o.duration_ms;
2106 self.state.quota.push(QuotaLoss {
2107 seat: seat.key.clone(),
2108 node: "fix".to_owned(),
2109 at: Timestamp::now(),
2110 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2111 });
2112 self.state.event(
2113 "fix",
2114 format!(
2115 "round {round}: continuation rate limited (quota); not retrying now"
2116 ),
2117 );
2118 return (
2119 seat,
2120 None,
2121 Some("rate limited (quota) while recovering the fix report".to_owned()),
2122 ContinuationRecord {
2123 attempts,
2124 cumulative_wait_ms,
2125 outcome: ContinuationOutcome::QuotaLost,
2126 },
2127 );
2128 }
2129 AgentOutcome::Dropped(o) => {
2130 cumulative_wait_ms += o.duration_ms;
2131 let why = o
2132 .dropped
2133 .as_ref()
2134 .map(|d| d.why.as_str())
2135 .unwrap_or("the CLI ended the stream without delivering its answer");
2136 last_err = format!("the CLI dropped the stream ({why})");
2137 }
2138 AgentOutcome::Failed(e) => last_err = e,
2139 }
2140 }
2141 }
2142
2143 fn after_implement(&mut self) -> Result<()> {
2144 if self.state.leaks.is_empty() {
2146 let cfg = self.state.config.blind.clone();
2147 let mut leaks = Vec::new();
2148 for c in &self.state.candidates {
2149 let Some(patch) =
2150 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
2151 else {
2152 continue;
2153 };
2154 leaks.extend(blind::scan(
2155 &format!("candidate {} patch", c.label),
2156 &patch,
2157 &cfg.vendor_tokens,
2158 ));
2159 }
2160 if !leaks.is_empty() {
2161 let summary = leaks
2162 .iter()
2163 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2164 .collect::<Vec<_>>()
2165 .join(", ");
2166 match cfg.on_leak {
2167 LeakPolicy::Fail => {
2168 self.state.status = RunStatus::Failed;
2169 self.state
2170 .event("blind", format!("vendor text in a patch: {summary}"));
2171 self.state.leaks = leaks;
2172 self.state.save()?;
2173 self.settle_questions();
2174 bail!(
2175 "blind.on_leak = \"fail\" and vendor text reached a \
2176 judged patch: {summary}"
2177 );
2178 }
2179 LeakPolicy::Redact => self.state.event(
2180 "blind",
2181 format!("redacting vendor text for judging: {summary}"),
2182 ),
2183 LeakPolicy::Warn => self.state.event(
2184 "blind",
2185 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2186 ),
2187 }
2188 self.state.leaks = leaks;
2189 }
2190 }
2191
2192 if self.state.viable().is_empty() {
2193 if self.state.all_candidates_verified_noop() {
2194 self.state.status = RunStatus::VerifiedNoop;
2205 self.state.save()?;
2206 self.settle_questions();
2207 return Ok(());
2208 }
2209 self.state.status = RunStatus::Failed;
2210 self.state.save()?;
2211 self.settle_questions();
2212 bail!("no candidate produced a change; nothing to judge");
2213 }
2214 self.state.status = RunStatus::Judging;
2215 self.state.save()?;
2216 Ok(())
2217 }
2218
2219 async fn judge(&mut self) -> Result<()> {
2222 let run_id = self.state.id.clone();
2227 let prompts = self.state.config.prompts.clone();
2228 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2229 return Ok(());
2230 }
2231 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2232 if viable.len() == 1 {
2233 self.state.judge_skipped = true;
2240 self.state.event(
2241 "judge",
2242 format!(
2243 "only candidate {} produced a change; judging skipped",
2244 viable[0].label
2245 ),
2246 );
2247 self.state.save()?;
2248 return Ok(());
2249 }
2250 self.state.status = RunStatus::Judging;
2251
2252 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2253 let language = self.state.config.graph.language.clone();
2254 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2255 let sessions = self.state.config.graph.sessions;
2256 let artifacts = agent::artifacts_dir(&self.state.dir());
2257 let root = self.state.worktree_root();
2258 let base_short = short(&self.state.base_commit);
2259
2260 let mut jobs = Vec::new();
2261 let mut orders = Vec::new();
2262 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2263 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2264 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2265 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2266 let seat_key = format!("judge-{}", j + 1);
2267 let seat = self.seat(&seat_key, &spec.id);
2268 jobs.push(SeatJob {
2269 prompt: prompt::judge(
2270 &self.state.instruction,
2271 &views,
2272 self.roles.judges.len(),
2273 &base_short,
2274 &language,
2275 ),
2276 spec,
2277 seat,
2278 cwd: root.join(format!("judge-{}", j + 1)),
2279 timeout,
2280 allow_write: false,
2281 sessions,
2282 artifacts: artifacts.clone(),
2283 stem: format!("judge-{}", j + 1),
2284 });
2285 }
2286
2287 self.state.event(
2288 "judge",
2289 format!(
2290 "{} judges ranking {} candidates blind",
2291 jobs.len(),
2292 viable.len()
2293 ),
2294 );
2295 let labels_for_check = labels.clone();
2296 let mut quota_losses = Vec::new();
2297 let cache = self.state.config.cache_dir();
2298 let ctx = WaveCtx {
2299 run: &run_id,
2300 node: "judge",
2301 prompts: &prompts,
2302 cache: cache.as_deref(),
2303 round: None,
2304 };
2305 let results = ask_json_wave::<Ranking>(
2306 jobs,
2307 Arc::clone(&self.sem),
2308 self.state.config.graph.retries,
2309 &ctx,
2310 &mut quota_losses,
2311 &mut self.state,
2312 &move |r: &Ranking| r.validate(&labels_for_check),
2313 )
2314 .await;
2315 self.state.quota.extend(quota_losses);
2316
2317 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2318 let agent_id = seat.agent.clone();
2319 self.state.seats.insert(seat.key.clone(), seat);
2320 let mut record = Judgement {
2321 judge: j + 1,
2322 seat: format!("judge-{}", j + 1),
2323 agent: agent_id,
2324 ranking: Vec::new(),
2325 reasons: BTreeMap::new(),
2326 confidence: None,
2327 order: orders[j].clone(),
2328 failed: None,
2329 duration_ms: 0,
2330 };
2331 match res {
2332 Ok((ranking, out)) => {
2333 record.ranking = ranking.normalized();
2334 record.reasons = ranking.reasons;
2335 record.confidence = ranking.confidence;
2336 record.duration_ms = out.duration_ms;
2337 self.state.event(
2338 "judge",
2339 format!(
2340 "judge {} ranked {}",
2341 j + 1,
2342 record.ranking.iter().collect::<String>()
2343 ),
2344 );
2345 }
2346 Err(e) => {
2347 record.failed = Some(e.to_string());
2348 self.state
2349 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2350 }
2351 }
2352 self.state.judgements.push(record);
2353 self.state.save()?;
2354 }
2355 Ok(())
2356 }
2357
2358 async fn deliberate(&mut self) -> Result<()> {
2361 let run_id = self.state.id.clone();
2366 let prompts = self.state.config.prompts.clone();
2367 if !self.state.deliberation.is_empty() {
2368 return Ok(());
2369 }
2370 let tops: Vec<char> = self
2371 .state
2372 .judgements
2373 .iter()
2374 .filter_map(|j| j.ranking.first().copied())
2375 .collect();
2376 let rounds = self.state.config.graph.deliberate_rounds;
2377 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2378 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2379 self.state.event(
2380 "deliberate",
2381 format!("judges agreed on {} outright; no deliberation", tops[0]),
2382 );
2383 }
2384 self.state.status = RunStatus::Voting;
2385 self.state.save()?;
2386 return Ok(());
2387 }
2388
2389 self.state.status = RunStatus::Deliberating;
2390 self.state.event(
2391 "deliberate",
2392 format!(
2393 "split: first choices were {} — opening {rounds} round(s)",
2394 tops.iter().collect::<String>()
2395 ),
2396 );
2397
2398 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2399 let language = self.state.config.graph.language.clone();
2400 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2401 let sessions = self.state.config.graph.sessions;
2402 let artifacts = agent::artifacts_dir(&self.state.dir());
2403 let root = self.state.worktree_root();
2404 let base_short = short(&self.state.base_commit);
2405
2406 for round in 1..=rounds {
2410 let mut turns: Vec<DeliberationTurn> = Vec::new();
2411 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2412 if self.state.judgements[j].failed.is_some() {
2413 continue;
2414 }
2415 let seat_key = format!("judge-{}", j + 1);
2416 let mut seat = self.seat(&seat_key, &spec.id);
2417 let transcript = self.transcript(&turns, j);
2418 let context = if has_context(&spec, &seat, sessions) {
2419 None
2420 } else {
2421 Some(self.candidate_block(&viable, &base_short))
2422 };
2423 let text = prompt::deliberate(
2424 &self.state.instruction,
2425 context.as_deref(),
2426 &transcript,
2427 round,
2428 rounds,
2429 &language,
2430 );
2431 let job = SeatJob {
2432 spec,
2433 seat: seat.clone(),
2434 prompt: text,
2435 cwd: root.join(format!("judge-{}", j + 1)),
2436 timeout,
2437 allow_write: false,
2438 sessions,
2439 artifacts: artifacts.clone(),
2440 stem: format!("delib-{round}-judge-{}", j + 1),
2441 };
2442 let cache = self.state.config.cache_dir();
2443 let ctx = WaveCtx {
2444 run: &run_id,
2445 node: "deliberate",
2446 prompts: &prompts,
2447 cache: cache.as_deref(),
2448 round: None,
2449 };
2450 let (updated, out) =
2451 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2452 seat = updated;
2453 let agent_id = seat.agent.clone();
2454 let seat_key = seat.key.clone();
2455 self.state.seats.insert(seat.key.clone(), seat);
2456 let body = match out {
2457 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2458 AgentOutcome::Dropped(o) => {
2462 let why =
2463 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2464 "the CLI ended the stream without delivering its answer",
2465 );
2466 self.state.event(
2467 "deliberate",
2468 format!(
2469 "judge {} skipped: the CLI dropped the stream ({why})",
2470 j + 1
2471 ),
2472 );
2473 continue;
2474 }
2475 AgentOutcome::Quota(o) => {
2476 self.state.quota.push(QuotaLoss {
2477 seat: seat_key,
2478 node: "deliberate".to_owned(),
2479 at: Timestamp::now(),
2480 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2481 });
2482 self.state.event(
2483 "deliberate",
2484 format!("judge {} skipped: rate limited (quota)", j + 1),
2485 );
2486 continue;
2487 }
2488 AgentOutcome::Failed(e) => {
2489 self.state
2490 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2491 continue;
2492 }
2493 };
2494 let tentative = verdict::extract_json::<Position>(&body)
2495 .ok()
2496 .and_then(|p| p.tentative)
2497 .and_then(|s| s.trim().chars().next())
2498 .map(|c| c.to_ascii_uppercase());
2499 self.state.event(
2500 "deliberate",
2501 format!(
2502 "round {round}: judge {} now favours {}",
2503 j + 1,
2504 tentative.map_or("—".to_owned(), |c| c.to_string())
2505 ),
2506 );
2507 turns.push(DeliberationTurn {
2508 judge: j + 1,
2509 agent: agent_id,
2510 body: blind::sanitize_prose(&body, &self.state.config.blind),
2511 tentative,
2512 });
2513 }
2514 self.state
2515 .deliberation
2516 .push(DeliberationRound { round, turns });
2517 self.state.save()?;
2518 }
2519
2520 self.state.status = RunStatus::Voting;
2521 self.state.save()?;
2522 Ok(())
2523 }
2524
2525 async fn vote(&mut self) -> Result<()> {
2528 let run_id = self.state.id.clone();
2533 let prompts = self.state.config.prompts.clone();
2534 if !self.state.votes.is_empty() {
2535 return Ok(());
2536 }
2537 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2538 if viable.len() == 1 {
2539 return Ok(());
2540 }
2541 self.state.status = RunStatus::Voting;
2542
2543 let language = self.state.config.graph.language.clone();
2544 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2545 let sessions = self.state.config.graph.sessions;
2546 let artifacts = agent::artifacts_dir(&self.state.dir());
2547 let root = self.state.worktree_root();
2548 let base_short = short(&self.state.base_commit);
2549 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2550
2551 let mut jobs = Vec::new();
2552 let mut seats_at = Vec::new();
2553 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2554 if self
2555 .state
2556 .judgements
2557 .get(j)
2558 .is_some_and(|r| r.failed.is_some())
2559 {
2560 continue;
2561 }
2562 let seat_key = format!("judge-{}", j + 1);
2563 let seat = self.seat(&seat_key, &spec.id);
2564 let mut text = prompt::final_vote(&viable, &language);
2565 if !has_context(&spec, &seat, sessions) {
2566 text = format!(
2567 "{}\n\n# Candidates\n\n{}",
2568 text,
2569 self.candidate_block(&candidates, &base_short)
2570 );
2571 }
2572 jobs.push(SeatJob {
2573 spec,
2574 seat,
2575 prompt: text,
2576 cwd: root.join(format!("judge-{}", j + 1)),
2577 timeout,
2578 allow_write: false,
2579 sessions,
2580 artifacts: artifacts.clone(),
2581 stem: format!("vote-judge-{}", j + 1),
2582 });
2583 seats_at.push(j);
2584 }
2585
2586 self.state.event(
2587 "vote",
2588 format!(
2589 "collecting {} final votes one by one, privately",
2590 jobs.len()
2591 ),
2592 );
2593 let allowed = viable.clone();
2594 let mut quota_losses = Vec::new();
2595 let cache = self.state.config.cache_dir();
2596 let ctx = WaveCtx {
2597 run: &run_id,
2598 node: "vote",
2599 prompts: &prompts,
2600 cache: cache.as_deref(),
2601 round: None,
2602 };
2603 let results = ask_json_wave::<FinalVote>(
2604 jobs,
2605 Arc::clone(&self.sem),
2606 self.state.config.graph.retries,
2607 &ctx,
2608 &mut quota_losses,
2609 &mut self.state,
2610 &move |v: &FinalVote| match v.label() {
2611 Some(c) if allowed.contains(&c) => Ok(()),
2612 other => bail!("vote {other:?} is not one of {allowed:?}"),
2613 },
2614 )
2615 .await;
2616 self.state.quota.extend(quota_losses);
2617
2618 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2619 let agent_id = seat.agent.clone();
2620 self.state.seats.insert(seat.key.clone(), seat);
2621 let initial = self
2622 .state
2623 .judgements
2624 .get(j)
2625 .and_then(|r| r.ranking.first().copied());
2626 let mut record = VoteRecord {
2627 judge: j + 1,
2628 agent: agent_id,
2629 vote: None,
2630 reason: String::new(),
2631 changed: false,
2632 };
2633 match res {
2634 Ok((v, _)) => {
2635 record.vote = v.label();
2636 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2637 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2638 self.state.event(
2639 "vote",
2640 format!(
2641 "judge {} voted {}{}",
2642 j + 1,
2643 record.vote.unwrap_or('?'),
2644 if record.changed { " (changed)" } else { "" }
2645 ),
2646 );
2647 }
2648 Err(e) => {
2649 self.state
2650 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2651 }
2652 }
2653 self.state.votes.push(record);
2654 self.state.save()?;
2655 }
2656 Ok(())
2657 }
2658
2659 fn tally(&mut self) -> Result<()> {
2662 if self.state.tally.is_some() {
2663 return Ok(());
2664 }
2665 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2666 let tops: Vec<char> = self
2667 .state
2668 .judgements
2669 .iter()
2670 .filter_map(|j| j.ranking.first().copied())
2671 .collect();
2672 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2673
2674 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2677 let mut cast: Vec<char> = Vec::new();
2678 for (i, j) in self.state.judgements.iter().enumerate() {
2679 let vote = self
2680 .state
2681 .votes
2682 .iter()
2683 .find(|v| v.judge == i + 1)
2684 .and_then(|v| v.vote)
2685 .or_else(|| j.ranking.first().copied());
2686 if let Some(v) = vote {
2687 *first_choice.entry(v).or_insert(0) += 1;
2688 cast.push(v);
2689 }
2690 }
2691
2692 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2693 for j in &self.state.judgements {
2694 let n = j.ranking.len();
2695 for (pos, label) in j.ranking.iter().enumerate() {
2696 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2697 }
2698 }
2699
2700 let best = first_choice.values().copied().max().unwrap_or(0);
2701 let mut leaders: Vec<char> = first_choice
2702 .iter()
2703 .filter(|(_, v)| **v == best)
2704 .map(|(k, _)| *k)
2705 .collect();
2706 let mut tie_break = None;
2707 if leaders.len() > 1 {
2708 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2709 let borda_leaders: Vec<char> = leaders
2710 .iter()
2711 .copied()
2712 .filter(|l| borda[l] == top_borda)
2713 .collect();
2714 tie_break = Some(if borda_leaders.len() == 1 {
2715 format!(
2716 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2717 leaders.len()
2718 )
2719 } else {
2720 format!(
2721 "{} way tie on both first-choice votes and Borda points, broken by label order",
2722 leaders.len()
2723 )
2724 });
2725 leaders = borda_leaders;
2726 leaders.sort_unstable();
2727 }
2728 let winner = *leaders
2729 .first()
2730 .or(viable.first())
2731 .context("no candidate to declare a winner from")?;
2732
2733 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2734 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2735 let deliberated = !self.state.deliberation.is_empty();
2736
2737 let quota_seats: std::collections::BTreeSet<&str> =
2741 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2742 let mut present = 0usize;
2743 for (i, j) in self.state.judgements.iter().enumerate() {
2744 if quota_seats.contains(j.seat.as_str()) {
2745 continue;
2746 }
2747 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2748 let voted = self
2749 .state
2750 .votes
2751 .iter()
2752 .any(|v| v.judge == i + 1 && v.vote.is_some());
2753 if ranked || voted {
2754 present += 1;
2755 }
2756 }
2757 let needs_quorum = viable.len() > 1;
2763 let judges_total = if needs_quorum {
2764 self.roles.judges.len()
2765 } else {
2766 0
2767 };
2768 let quorum = if needs_quorum {
2769 judges_total / 2 + 1
2770 } else {
2771 0
2772 };
2773 let met_quorum = !needs_quorum || present >= quorum;
2774 let uncontested = (!needs_quorum).then(|| {
2775 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2776 });
2777
2778 self.state.event(
2779 "tally",
2780 match &uncontested {
2781 Some(reason) => format!("winner {winner} — {reason}"),
2782 None => format!(
2783 "winner {winner} — votes {} | initial {} | {} changed | \
2784 {present}/{judges_total} judges{}",
2785 first_choice
2786 .iter()
2787 .map(|(k, v)| format!("{k}:{v}"))
2788 .collect::<Vec<_>>()
2789 .join(" "),
2790 if unanimous_initial {
2791 "unanimous"
2792 } else {
2793 "split"
2794 },
2795 changed_votes,
2796 if met_quorum {
2797 String::new()
2798 } else {
2799 format!(" — below quorum ({quorum} required)")
2800 },
2801 ),
2802 },
2803 );
2804 if !met_quorum {
2805 self.state.event(
2806 "stall",
2807 format!(
2808 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2809 the run stops here, resumable"
2810 ),
2811 );
2812 }
2813 self.state.tally = Some(Tally {
2814 first_choice,
2815 borda,
2816 winner,
2817 rankings: tops.len(),
2818 unanimous_initial,
2819 deliberated,
2820 changed_votes,
2821 unanimous_final,
2822 tie_break,
2823 judges: judges_total,
2824 present,
2825 quorum,
2826 met_quorum,
2827 uncontested,
2828 });
2829 self.state.status = if met_quorum {
2830 RunStatus::Reviewing
2831 } else {
2832 RunStatus::Stalled
2833 };
2834 self.state.save()?;
2835 Ok(())
2836 }
2837
2838 #[allow(clippy::too_many_lines)]
2859 async fn recover_stall(&mut self) -> Result<bool> {
2860 let run_id = self.state.id.clone();
2865 let prompts = self.state.config.prompts.clone();
2866 let quota_seats: BTreeSet<&str> =
2871 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2872 let absent: Vec<String> = self
2873 .state
2874 .judgements
2875 .iter()
2876 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2877 .map(|j| j.seat.clone())
2878 .collect();
2879 if absent.is_empty() {
2880 return Ok(false);
2881 }
2882 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2883 if viable.len() <= 1 {
2884 return Ok(false);
2885 }
2886 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2887 let language = self.state.config.graph.language.clone();
2888 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2889 let sessions = self.state.config.graph.sessions;
2890 let artifacts = agent::artifacts_dir(&self.state.dir());
2891 let root = self.state.worktree_root();
2892 let base_short = short(&self.state.base_commit);
2893 let candidates: Vec<Candidate> = viable.clone();
2894
2895 let mut positions: Vec<usize> = absent
2897 .iter()
2898 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2899 .collect();
2900 if positions.is_empty() {
2901 return Ok(false);
2902 }
2903 positions.sort_unstable();
2904 positions.dedup();
2905
2906 let mut judge_jobs = Vec::new();
2908 for &j in &positions {
2909 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2910 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2911 let seat_key = format!("judge-{}", j + 1);
2912 let spec = self.roles.judges[j].clone();
2913 let seat = self.seat(&seat_key, &spec.id);
2914 judge_jobs.push(SeatJob {
2915 spec,
2916 seat,
2917 prompt: prompt::judge(
2918 &self.state.instruction,
2919 &views,
2920 self.roles.judges.len(),
2921 &base_short,
2922 &language,
2923 ),
2924 cwd: root.join(seat_key),
2925 timeout,
2926 allow_write: false,
2927 sessions,
2928 artifacts: artifacts.clone(),
2929 stem: format!("judge-{}-recover", j + 1),
2930 });
2931 }
2932
2933 let labels_for_check = labels.clone();
2934 let mut judge_losses = Vec::new();
2935 let retries = self.state.config.graph.retries;
2936 let cache = self.state.config.cache_dir();
2937 let ctx = WaveCtx {
2938 run: &run_id,
2939 node: "judge",
2940 prompts: &prompts,
2941 cache: cache.as_deref(),
2942 round: None,
2943 };
2944 let results = ask_json_wave::<Ranking>(
2945 judge_jobs,
2946 Arc::clone(&self.sem),
2947 retries,
2948 &ctx,
2949 &mut judge_losses,
2950 &mut self.state,
2951 &move |r: &Ranking| r.validate(&labels_for_check),
2952 )
2953 .await;
2954
2955 let mut recovered: BTreeSet<usize> = BTreeSet::new();
2957 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
2958 self.state.seats.insert(seat.key.clone(), seat);
2959 let record = &mut self.state.judgements[j];
2960 match res {
2961 Ok((ranking, out)) => {
2962 record.ranking = ranking.normalized();
2963 record.reasons = ranking.reasons;
2964 record.confidence = ranking.confidence;
2965 record.failed = None;
2966 record.duration_ms = out.duration_ms;
2967 recovered.insert(j);
2968 self.state.event(
2969 "recover",
2970 format!("judge {} ranked again after the limit", j + 1),
2971 );
2972 }
2973 Err(e) => {
2974 self.state
2975 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
2976 }
2977 }
2978 }
2979
2980 let mut vote_jobs = Vec::new();
2982 let mut vote_pos: Vec<usize> = Vec::new();
2983 for &j in &recovered {
2984 let seat_key = format!("judge-{}", j + 1);
2985 let spec = self.roles.judges[j].clone();
2986 let seat = self.seat(&seat_key, &spec.id);
2987 let mut text = prompt::final_vote(&labels, &language);
2988 if !has_context(&spec, &seat, sessions) {
2989 text = format!(
2990 "{}\n\n# Candidates\n\n{}",
2991 text,
2992 self.candidate_block(&candidates, &base_short)
2993 );
2994 }
2995 vote_jobs.push(SeatJob {
2996 spec,
2997 seat,
2998 prompt: text,
2999 cwd: root.join(seat_key),
3000 timeout,
3001 allow_write: false,
3002 sessions,
3003 artifacts: artifacts.clone(),
3004 stem: format!("vote-judge-{}-recover", j + 1),
3005 });
3006 vote_pos.push(j);
3007 }
3008 let allowed = labels.clone();
3009 let mut vote_losses = Vec::new();
3010 let vote_retries = self.state.config.graph.retries;
3011 let vote_cache = self.state.config.cache_dir();
3012 let ctx = WaveCtx {
3013 run: &run_id,
3014 node: "vote",
3015 prompts: &prompts,
3016 cache: vote_cache.as_deref(),
3017 round: None,
3018 };
3019 let votes = ask_json_wave::<FinalVote>(
3020 vote_jobs,
3021 Arc::clone(&self.sem),
3022 vote_retries,
3023 &ctx,
3024 &mut vote_losses,
3025 &mut self.state,
3026 &move |v: &FinalVote| match v.label() {
3027 Some(c) if allowed.contains(&c) => Ok(()),
3028 other => bail!("vote {other:?} is not one of {allowed:?}"),
3029 },
3030 )
3031 .await;
3032 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
3033 let agent_id = seat.agent.clone();
3034 self.state.seats.insert(seat.key.clone(), seat);
3035 match res {
3036 Ok((v, _)) => {
3037 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
3038 rec.vote = v.label();
3039 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
3040 } else {
3041 self.state.votes.push(VoteRecord {
3042 judge: j + 1,
3043 agent: agent_id,
3044 vote: v.label(),
3045 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
3046 changed: false,
3047 });
3048 }
3049 self.state.event(
3050 "recover",
3051 format!("judge {} voted again after the limit", j + 1),
3052 );
3053 }
3054 Err(e) => {
3055 self.state
3056 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
3057 }
3058 }
3059 }
3060
3061 let recovered_keys: BTreeSet<String> = recovered
3065 .iter()
3066 .map(|&j| format!("judge-{}", j + 1))
3067 .collect();
3068 self.state
3069 .quota
3070 .retain(|q| !recovered_keys.contains(&q.seat));
3071 for loss in judge_losses.into_iter().chain(vote_losses) {
3075 if recovered_keys.contains(&loss.seat) {
3076 continue;
3077 }
3078 self.state.quota.retain(|q| q.seat != loss.seat);
3079 self.state.quota.push(loss);
3080 }
3081
3082 self.state.tally = None;
3084 self.tally()?;
3085 Ok(self
3086 .state
3087 .tally
3088 .as_ref()
3089 .map(|t| t.met_quorum)
3090 .unwrap_or(false))
3091 }
3092
3093 async fn fold_losers(&mut self) -> Result<()> {
3096 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
3097 return Ok(());
3098 };
3099 let repo = self.state.repo.clone();
3100 let mut folded = Vec::new();
3101 for i in 0..self.state.candidates.len() {
3102 let c = &self.state.candidates[i];
3103 if c.label == winner || c.folded {
3104 continue;
3105 }
3106 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
3107 git::worktree_remove(&repo, &wt).await.ok();
3108 git::branch_delete(&repo, &branch).await.ok();
3109 self.state.candidates[i].folded = true;
3110 folded.push(label.to_string());
3111 }
3112 let root = self.state.worktree_root();
3114 for j in 1..=self.roles.judges.len() {
3115 let wt = root.join(format!("judge-{j}"));
3116 if wt.exists() {
3117 git::worktree_remove(&repo, &wt).await.ok();
3118 }
3119 }
3120 if self.state.config.graph.advise {
3123 for k in 1..=self.state.config.graph.advisors {
3124 let wt = root.join(format!("advisor-{k}"));
3125 if wt.exists() {
3126 git::worktree_remove(&repo, &wt).await.ok();
3127 }
3128 }
3129 }
3130 if !folded.is_empty() {
3131 self.state
3132 .event("fold", format!("folded candidates {}", folded.join(", ")));
3133 self.state.save()?;
3134 }
3135 Ok(())
3136 }
3137
3138 async fn sync_to_base(&mut self) -> Result<()> {
3168 if self
3169 .state
3170 .base_sync
3171 .as_ref()
3172 .is_some_and(|s| s.conflict.is_some())
3173 {
3174 return Ok(());
3175 }
3176 let Some(winner) = self.state.winner().cloned() else {
3177 return Ok(());
3178 };
3179
3180 let repo = self.state.repo.clone();
3181 let remote = self.state.config.merge.remote.clone();
3182 let base_branch = self.state.base_branch.clone();
3183 let tracking = format!("{remote}/{base_branch}");
3184
3185 git::fetch(&repo, &remote, &base_branch).await.ok();
3186 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3190 return Ok(());
3191 };
3192
3193 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3194 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3195 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3196
3197 if behind == 0 {
3198 self.state.base_sync = Some(BaseSync {
3199 tip,
3200 behind: 0,
3201 attempts,
3202 conflict: None,
3203 });
3204 self.state.save()?;
3205 return Ok(());
3206 }
3207
3208 if attempts >= BASE_SYNC_ROUNDS {
3209 let why = format!(
3210 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3211 rebase(s); rebasing again would only race it",
3212 winner.branch
3213 );
3214 self.state.status = RunStatus::Blocked;
3215 self.state.base_sync = Some(BaseSync {
3216 tip,
3217 behind,
3218 attempts,
3219 conflict: Some(why.clone()),
3220 });
3221 self.state.event("land", why);
3222 self.state.save()?;
3223 return Ok(());
3224 }
3225
3226 self.state.event(
3227 "land",
3228 format!(
3229 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3230 winner.branch
3231 ),
3232 );
3233 self.state.save()?;
3234
3235 let scratch = self.state.dir().join("base-sync");
3236 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
3237 let attempts = attempts + 1;
3238 match rebased {
3239 Ok(None) => {
3240 git::sync_to_head(&winner.worktree).await?;
3244 self.state.base_sync = Some(BaseSync {
3245 tip: tip.clone(),
3246 behind: 0,
3247 attempts,
3248 conflict: None,
3249 });
3250 self.state
3251 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3252 }
3253 Ok(Some(conflict)) => {
3254 let why = format!(
3255 "{} conflicts with {tracking} and did not rebase: {}",
3256 winner.branch,
3257 conflict.chars().take(600).collect::<String>()
3258 );
3259 self.state.status = RunStatus::Blocked;
3260 self.state.base_sync = Some(BaseSync {
3261 tip,
3262 behind,
3263 attempts,
3264 conflict: Some(why.clone()),
3265 });
3266 self.state.event("land", why);
3267 }
3268 Err(e) => {
3269 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3270 self.state.status = RunStatus::Blocked;
3271 self.state.base_sync = Some(BaseSync {
3272 tip,
3273 behind,
3274 attempts,
3275 conflict: Some(why.clone()),
3276 });
3277 self.state.event("land", why);
3278 }
3279 }
3280 self.state.save()?;
3281 Ok(())
3282 }
3283
3284 fn landing_base(&self) -> String {
3294 self.state
3295 .base_sync
3296 .as_ref()
3297 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3298 }
3299
3300 pub async fn fix_selected(
3333 &mut self,
3334 ids: &[String],
3335 reason: &str,
3336 allow_stale: bool,
3337 ) -> Result<()> {
3338 let reason = reason.trim();
3339 if reason.is_empty() {
3340 bail!("a fix request needs a reason — that is the operator's own record of why");
3341 }
3342 if ids.is_empty() {
3343 bail!("no finding id given");
3344 }
3345 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3346 bail!(
3347 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3348 has already concluded — can be given a targeted fix. A run still \
3349 in progress should simply be resumed; a `merged` run's branch has \
3350 already landed, so its answer is a fresh `magi review <branch>`, \
3351 not reopening this run's own record",
3352 self.state.id,
3353 self.state.status.as_str()
3354 );
3355 }
3356 let Some(winner) = self.state.winner().cloned() else {
3357 bail!("run {} has no winning candidate to fix", self.state.id);
3358 };
3359 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3360 bail!(
3361 "branch `{}` no longer exists; this run cannot be extended",
3362 winner.branch
3363 );
3364 }
3365 let home = crate::run::home();
3366 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3367 bail!(
3368 "run {} is currently being worked on by another magi process",
3369 self.state.id
3370 );
3371 }
3372 let _claim = FixClaim::acquire(&self.state.dir())?;
3378
3379 let mut seen = BTreeSet::new();
3383 let mut findings = Vec::new();
3384 let mut missing = Vec::new();
3385 for id in ids {
3386 if !seen.insert(id.clone()) {
3387 continue;
3388 }
3389 match self.state.finding(id) {
3390 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3391 id: f.id.clone(),
3392 severity: f.severity,
3393 reviewer_vote: rec.vote,
3394 round: round.round,
3395 round_head: round.head.clone(),
3396 reviewer: rec.reviewer,
3397 agent: rec.agent.clone(),
3398 file: f.file.clone(),
3399 line: f.line,
3400 title: f.title.clone(),
3401 detail: f.detail.clone(),
3402 outcome: OperatorFixOutcome::Pending,
3403 }),
3404 None => missing.push(id.clone()),
3405 }
3406 }
3407 if !missing.is_empty() {
3408 bail!(
3409 "unknown finding id(s): {}; nothing was changed",
3410 missing.join(", ")
3411 );
3412 }
3413
3414 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3415 let stale_details: Vec<(String, String)> = findings
3416 .iter()
3417 .filter(|f| f.round_head != head_at_request)
3418 .map(|f| (f.id.clone(), f.round_head.clone()))
3419 .collect();
3420 let stale = !stale_details.is_empty();
3421 if stale && !allow_stale {
3422 bail!(
3423 "the branch has moved since some finding(s) were raised — {} — now \
3424 at {}; pass --allow-stale to fix anyway, or re-run review first",
3425 stale_details
3426 .iter()
3427 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3428 .collect::<Vec<_>>()
3429 .join(", "),
3430 short(&head_at_request)
3431 );
3432 }
3433
3434 let request = OperatorFixRequest {
3435 requested_at: Timestamp::now(),
3436 reason: reason.to_owned(),
3437 findings,
3438 head_at_request: head_at_request.clone(),
3439 allow_stale,
3440 stale,
3441 fix: None,
3442 result_head: None,
3443 follow_up_review_run: None,
3444 };
3445 self.state.event(
3446 "fix",
3447 format!(
3448 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3449 request.findings.len(),
3450 request
3451 .findings
3452 .iter()
3453 .map(|f| f.id.as_str())
3454 .collect::<Vec<_>>()
3455 .join(", "),
3456 ),
3457 );
3458 self.state.operator_fixes.push(request);
3465 self.state.save()?;
3466 let request_index = self.state.operator_fixes.len() - 1;
3467
3468 if winner.worktree.exists() {
3477 let dirty = git::git(
3480 &winner.worktree,
3481 &["status", "--porcelain", "--untracked-files=all"],
3482 )
3483 .await?;
3484 let only_withheld = dirty.lines().all(|l| {
3485 l.strip_prefix("?? ")
3486 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3487 });
3488 if !only_withheld {
3489 bail!(
3490 "`{}` has uncommitted changes; refusing to touch it — commit or \
3491 discard them first",
3492 winner.worktree.display()
3493 );
3494 }
3495 git::worktree_remove(&self.state.repo, &winner.worktree)
3496 .await
3497 .ok();
3498 }
3499 let fix_worktree = self.state.worktree_root().join("operator-fix");
3500 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3501 git::git(
3502 &self.state.repo,
3503 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3504 )
3505 .await
3506 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3507 if !git::is_clean(&fix_worktree).await? {
3508 git::worktree_remove(&self.state.repo, &fix_worktree)
3509 .await
3510 .ok();
3511 bail!(
3512 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3513 winner.branch
3514 );
3515 }
3516
3517 let run_id = self.state.id.clone();
3518 let prompts = self.state.config.prompts.clone();
3519 let language = self.state.config.graph.language.clone();
3520 let sessions = self.state.config.graph.sessions;
3521 let artifacts = agent::artifacts_dir(&self.state.dir());
3522 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3523 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3524 _ => (
3525 self.state
3526 .config
3527 .agent(&winner.agent)
3528 .cloned()
3529 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3530 format!("impl-{}", winner.label),
3531 ),
3532 };
3533 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3534 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3535 .findings
3536 .iter()
3537 .map(|f| Finding {
3538 id: f.id.clone(),
3539 severity: f.severity,
3540 file: f.file.clone(),
3541 line: f.line,
3542 title: f.title.clone(),
3543 detail: f.detail.clone(),
3544 })
3545 .collect();
3546 let job = SeatJob {
3547 prompt: prompt::operator_fix(
3548 &self.state.instruction,
3549 &finding_list,
3550 reason,
3551 &stale_details,
3552 &head_at_request,
3553 &language,
3554 ),
3555 spec: fix_spec.clone(),
3556 seat,
3557 cwd: fix_worktree.clone(),
3558 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3559 allow_write: true,
3560 sessions,
3561 artifacts: artifacts.clone(),
3562 stem: "operator-fix".to_owned(),
3563 };
3564 let cache = self.state.config.cache_dir();
3565 let ctx = WaveCtx {
3566 run: &run_id,
3567 node: "fix",
3568 prompts: &prompts,
3569 cache: cache.as_deref(),
3570 round: None,
3571 };
3572 let (seat, out) =
3573 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3574 let agent_id = seat.agent.clone();
3575
3576 let mut fix = FixRecord {
3577 agent: agent_id,
3578 addressed: Vec::new(),
3579 rejected: Vec::new(),
3580 notes: String::new(),
3581 committed: false,
3582 failed: None,
3583 duration_ms: 0,
3584 continuation: None,
3585 };
3586 let mut final_seat = seat.clone();
3587 match out {
3588 AgentOutcome::Ok(o) => {
3589 fix.duration_ms = o.duration_ms;
3590 let parsed = verdict::extract_json::<FixReport>(&o.text);
3591 let incomplete_reason = match &parsed {
3592 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3593 "the reply parsed, but it reported a command whose own CLI \
3594 never confirmed an exit status"
3595 .to_owned(),
3596 ),
3597 Ok(_) => None,
3598 Err(e) => Some(e.to_string()),
3599 };
3600 match incomplete_reason {
3601 None => {
3602 let report = parsed.expect("checked Ok above");
3603 fix.addressed = report.addressed;
3604 fix.rejected = report.rejected;
3605 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3606 }
3607 Some(reason) => {
3608 let (resumed_seat, resolved, failure, cont) = self
3609 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3610 .await;
3611 fix.duration_ms += cont.cumulative_wait_ms;
3612 fix.continuation = Some(cont);
3613 final_seat = resumed_seat;
3614 match resolved {
3615 Some(report) => {
3616 fix.addressed = report.addressed;
3617 fix.rejected = report.rejected;
3618 fix.notes =
3619 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3620 }
3621 None => fix.failed = failure,
3622 }
3623 }
3624 }
3625 }
3626 AgentOutcome::Dropped(o) => {
3627 fix.duration_ms = o.duration_ms;
3628 let why = o
3629 .dropped
3630 .as_ref()
3631 .map(|d| d.why.as_str())
3632 .unwrap_or("the CLI ended the stream without delivering its answer");
3633 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3634 }
3635 AgentOutcome::Quota(o) => {
3636 self.state.quota.push(QuotaLoss {
3637 seat: final_seat.key.clone(),
3638 node: "fix".to_owned(),
3639 at: Timestamp::now(),
3640 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3641 });
3642 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3643 }
3644 AgentOutcome::Failed(e) => fix.failed = Some(e),
3645 }
3646 if fix.continuation.is_none() {
3647 fix.continuation = Some(ContinuationRecord::not_needed());
3648 }
3649 self.state.seats.insert(final_seat.key.clone(), final_seat);
3650
3651 let rescue_message = format!(
3652 "magi: operator-selected fix ({}) (uncommitted work)",
3653 self.state.operator_fixes[request_index]
3654 .findings
3655 .iter()
3656 .map(|f| f.id.as_str())
3657 .collect::<Vec<_>>()
3658 .join(", ")
3659 );
3660 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3661 self.state.note_withheld("fix", &r.withheld);
3662 }
3663 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3664 fix.committed = after != head_at_request;
3665 git::worktree_remove(&self.state.repo, &fix_worktree)
3666 .await
3667 .ok();
3668
3669 self.state.event(
3670 "fix",
3671 match &fix.failed {
3672 Some(reason) => format!(
3673 "operator fix: adoption report was lost ({reason}); {}",
3674 if fix.committed {
3675 "committed"
3676 } else {
3677 "NO new commit"
3678 }
3679 ),
3680 None => format!(
3681 "operator fix: {} addressed, {} rejected, {}",
3682 fix.addressed.len(),
3683 fix.rejected.len(),
3684 if fix.committed {
3685 "committed"
3686 } else {
3687 "NO new commit"
3688 }
3689 ),
3690 },
3691 );
3692
3693 for f in &mut self.state.operator_fixes[request_index].findings {
3700 f.outcome = if fix.failed.is_some() {
3701 OperatorFixOutcome::Unreported
3702 } else if fix.addressed.contains(&f.id) {
3703 OperatorFixOutcome::Addressed
3704 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3705 OperatorFixOutcome::Rejected { why: r.why.clone() }
3706 } else {
3707 OperatorFixOutcome::Unreported
3708 };
3709 }
3710
3711 let committed = fix.committed;
3712 if committed {
3713 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3714 }
3715 self.state.operator_fixes[request_index].fix = Some(fix);
3716 self.state.save()?;
3719
3720 if committed {
3721 self.state.event(
3722 "fix",
3723 format!(
3724 "operator fix committed {}; opening a follow-up review-only run",
3725 short(&after)
3726 ),
3727 );
3728 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3729 Ok(mut follow_up) => {
3730 follow_up.state.event(
3731 "start",
3732 format!(
3733 "requested by an operator fix on run {} for finding(s) {}",
3734 self.state.id,
3735 self.state.operator_fixes[request_index]
3736 .findings
3737 .iter()
3738 .map(|f| f.id.as_str())
3739 .collect::<Vec<_>>()
3740 .join(", "),
3741 ),
3742 );
3743 follow_up.state.save()?;
3744 let follow_up_id = follow_up.state.id.clone();
3745 if let Err(e) = follow_up.execute().await {
3746 self.state.event(
3747 "fix",
3748 format!(
3749 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3750 ),
3751 );
3752 }
3753 self.state.operator_fixes[request_index].follow_up_review_run =
3754 Some(follow_up_id);
3755 }
3756 Err(e) => {
3757 self.state.event(
3758 "fix",
3759 format!("committed the fix but could not open a follow-up review: {e:#}"),
3760 );
3761 }
3762 }
3763 self.state.save()?;
3764 }
3765
3766 Ok(())
3767 }
3768
3769 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3776 match &self.roles.fixer {
3777 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3778 _ => (
3779 self.state
3780 .config
3781 .agent(&winner.agent)
3782 .cloned()
3783 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3784 format!("impl-{}", winner.label),
3785 ),
3786 }
3787 }
3788
3789 async fn review_loop(&mut self) -> Result<()> {
3790 if self
3795 .state
3796 .base_sync
3797 .as_ref()
3798 .is_some_and(|s| s.conflict.is_some())
3799 {
3800 return Ok(());
3801 }
3802 let run_id = self.state.id.clone();
3807 let prompts = self.state.config.prompts.clone();
3808 let Some(winner) = self.state.winner().cloned() else {
3809 return Ok(());
3810 };
3811 let max_rounds = self.state.config.graph.review_rounds;
3812 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
3822 self.state.status = status;
3823 self.state.save()?;
3824 return Ok(());
3825 }
3826 self.state.status = RunStatus::Reviewing;
3827 if self
3837 .state
3838 .reviews
3839 .last()
3840 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
3841 {
3842 let shell = self.state.config.shell();
3843 return self
3844 .stop_reviewing(
3845 "the last round's own verification never resolved",
3846 &shell,
3847 &winner.worktree,
3848 )
3849 .await;
3850 }
3851
3852 let repo = self.state.repo.clone();
3853 let root = self.state.worktree_root();
3854 let language = self.state.config.graph.language.clone();
3855 let sessions = self.state.config.graph.sessions;
3856 let artifacts = agent::artifacts_dir(&self.state.dir());
3857 let base = self.landing_base();
3858 let base_short = short(&base);
3859 let reviewers = self.roles.reviewers.clone();
3860 let shell = self.state.config.shell();
3861
3862 for round in (self.state.reviews.len() + 1)..=max_rounds {
3863 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3864 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
3865 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
3866 let prev_verification = self
3875 .state
3876 .reviews
3877 .last()
3878 .and_then(|r| r.verification_summary(&head));
3879
3880 let mut jobs = Vec::new();
3884 for (r, spec) in reviewers.iter().cloned().enumerate() {
3885 let wt = root.join(format!("review-{}", r + 1));
3886 if wt.exists() {
3887 git::reset_detached(&wt, &head).await?;
3888 } else {
3889 git::worktree_add_detached(&repo, &wt, &head).await?;
3890 }
3891 let seat_key = format!("review-{}", r + 1);
3892 let seat = self.seat(&seat_key, &spec.id);
3893 jobs.push(SeatJob {
3894 prompt: prompt::review(&prompt::ReviewCtx {
3895 instruction: &self.state.instruction,
3896 branch: &winner.branch,
3897 base_short: &base_short,
3898 stat: &stat,
3899 patch: &patch,
3900 verification: prev_verification.as_ref(),
3901 reviewers: reviewers.len(),
3902 round,
3903 rounds: max_rounds,
3904 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
3907 lens: Lens::for_seat(r),
3908 language: &language,
3909 }),
3910 spec,
3911 seat,
3912 cwd: wt,
3913 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
3914 allow_write: false,
3915 sessions,
3916 artifacts: artifacts.clone(),
3917 stem: format!("review-{round}-{}", r + 1),
3918 });
3919 }
3920
3921 self.state.event(
3922 "review",
3923 format!(
3924 "round {round}: {} reviewers on {}",
3925 jobs.len(),
3926 short(&head)
3927 ),
3928 );
3929 let mut quota_losses = Vec::new();
3930 let review_retries = self.state.config.graph.retries;
3931 let review_cache = self.state.config.cache_dir();
3932 let ctx = WaveCtx {
3933 run: &run_id,
3934 node: "review",
3935 prompts: &prompts,
3936 cache: review_cache.as_deref(),
3937 round: Some(round),
3938 };
3939 let results = ask_json_wave::<Review>(
3940 jobs,
3941 Arc::clone(&self.sem),
3942 review_retries,
3943 &ctx,
3944 &mut quota_losses,
3945 &mut self.state,
3946 &|_: &Review| Ok(()),
3947 )
3948 .await;
3949 let round_quota_missing = quota_losses.len();
3953 self.state.quota.extend(quota_losses);
3954
3955 let mut records = Vec::new();
3956 let mut all_findings = Vec::new();
3957 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
3958 let agent_id = seat.agent.clone();
3959 self.state.seats.insert(seat.key.clone(), seat);
3960 let mut record = ReviewRecord {
3961 reviewer: r + 1,
3962 agent: agent_id,
3963 summary: String::new(),
3964 findings: Vec::new(),
3965 vote: None,
3966 failed: None,
3967 duration_ms: 0,
3968 attempts,
3974 };
3975 match res {
3976 Ok((review, out)) => {
3977 record.summary =
3985 blind::sanitize_prose(&review.summary, &self.state.config.blind);
3986 record.vote = Some(review.vote);
3987 record.duration_ms = out.duration_ms;
3988 for (n, mut f) in review.findings.into_iter().enumerate() {
3989 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
3992 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
3993 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
3994 f.file = f
4000 .file
4001 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
4002 all_findings.push(f.clone());
4003 record.findings.push(f);
4004 }
4005 self.state.event(
4006 "review",
4007 format!(
4008 "round {round}: reviewer {} voted {} with {} finding(s)",
4009 r + 1,
4010 review.vote.label(),
4011 record.findings.len()
4012 ),
4013 );
4014 }
4015 Err(e) => {
4016 record.failed = Some(e.to_string());
4017 self.state.event(
4018 "review",
4019 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
4020 );
4021 }
4022 }
4023 records.push(record);
4024 }
4025
4026 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
4033 let vote_split =
4034 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
4035 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
4036 if vote_split {
4037 self.state.event(
4038 "review",
4039 format!(
4040 "round {round}: votes split ({}) — one round of reconsideration",
4041 initial_votes
4042 .iter()
4043 .map(|v| v.label())
4044 .collect::<Vec<_>>()
4045 .join(", ")
4046 ),
4047 );
4048 let panel: Vec<ReviewSeatReport<'_>> = records
4051 .iter()
4052 .filter_map(|r| {
4053 r.vote.map(|vote| ReviewSeatReport {
4054 reviewer: r.reviewer,
4055 vote,
4056 summary: &r.summary,
4057 findings: &r.findings,
4058 })
4059 })
4060 .collect();
4061
4062 let mut jobs = Vec::new();
4063 let mut seats_at = Vec::new();
4064 for (r, spec) in reviewers.iter().cloned().enumerate() {
4065 if records[r].vote.is_none() {
4069 continue;
4070 }
4071 let wt = root.join(format!("review-{}", r + 1));
4072 let seat_key = format!("review-{}", r + 1);
4073 let seat = self.seat(&seat_key, &spec.id);
4074 let patch_ctx = if has_context(&spec, &seat, sessions) {
4079 None
4080 } else {
4081 Some(ReviewPatch {
4082 branch: &winner.branch,
4083 base_short: &base_short,
4084 stat: &stat,
4085 patch: &patch,
4086 })
4087 };
4088 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
4089 instruction: &self.state.instruction,
4090 reviewer: r + 1,
4091 lens: Lens::for_seat(r),
4092 panel: &panel,
4093 patch: patch_ctx,
4094 round,
4095 rounds: max_rounds,
4096 language: &language,
4097 });
4098 jobs.push(SeatJob {
4099 prompt,
4100 spec,
4101 seat,
4102 cwd: wt,
4103 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4104 allow_write: false,
4105 sessions,
4106 artifacts: artifacts.clone(),
4107 stem: format!("review-{round}-reconsider-{}", r + 1),
4108 });
4109 seats_at.push(r);
4110 }
4111
4112 let mut recon_quota_losses = Vec::new();
4113 let recon_cache = self.state.config.cache_dir();
4114 let recon_ctx = WaveCtx {
4115 run: &run_id,
4116 node: "review",
4117 prompts: &prompts,
4118 cache: recon_cache.as_deref(),
4119 round: Some(round),
4120 };
4121 let recon_results = ask_json_wave::<ReviewRevote>(
4122 jobs,
4123 Arc::clone(&self.sem),
4124 review_retries,
4125 &recon_ctx,
4126 &mut recon_quota_losses,
4127 &mut self.state,
4128 &|_: &ReviewRevote| Ok(()),
4129 )
4130 .await;
4131 self.state.quota.extend(recon_quota_losses);
4132
4133 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
4134 let agent_id = seat.agent.clone();
4135 self.state.seats.insert(seat.key.clone(), seat);
4136 let mut rec = ReviewRevoteRecord {
4137 reviewer: r + 1,
4138 agent: agent_id,
4139 vote: None,
4140 reason: String::new(),
4141 failed: None,
4142 };
4143 match res {
4144 Ok((rv, _)) => {
4145 rec.vote = Some(rv.vote);
4146 rec.reason =
4147 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
4148 self.state.event(
4149 "review",
4150 format!(
4151 "round {round}: reviewer {} revoted {}",
4152 r + 1,
4153 rv.vote.label()
4154 ),
4155 );
4156 }
4157 Err(e) => {
4158 rec.failed = Some(e.to_string());
4159 self.state.event(
4160 "review",
4161 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4162 );
4163 }
4164 }
4165 reconsideration.push(rec);
4166 }
4167 } else if initial_votes.len() > 1 {
4168 self.state.event(
4169 "review",
4170 format!(
4171 "round {round}: votes agreed ({}) — no reconsideration",
4172 initial_votes[0].label()
4173 ),
4174 );
4175 }
4176
4177 let final_votes: Vec<ReviewVote> = records
4181 .iter()
4182 .filter_map(|r| {
4183 reconsideration
4184 .iter()
4185 .find(|rv| rv.reviewer == r.reviewer)
4186 .and_then(|rv| rv.vote)
4187 .or(r.vote)
4188 })
4189 .collect();
4190 let round_verdict = ReviewVote::worst(final_votes);
4191
4192 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4193 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4194 let defer_e2e =
4205 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4206 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4207 let reason =
4208 format!("{blocking} blocking finding(s) already required a fix this round");
4209 self.state.event(
4210 "verify",
4211 format!(
4212 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4213 {}); it will run once a round has none left",
4214 short(&head)
4215 ),
4216 );
4217 (Vec::new(), false, true, Some(reason))
4218 } else {
4219 let e2e_commands = self.state.config.verify.e2e.clone();
4220 let cache_dir = self.state.config.cache_dir();
4221 let context = format!("round {round}");
4222 let (e2e, verify_retried) = with_cache_lease(
4223 &mut self.state,
4224 cache_dir.as_deref(),
4225 "e2e",
4226 "e2e",
4227 &winner.worktree,
4228 &head,
4229 verify_timeout,
4230 &context,
4231 |state, budget| {
4232 let shell = shell.clone();
4233 let e2e_commands = e2e_commands.clone();
4234 let worktree = winner.worktree.clone();
4235 let context = context.clone();
4236 async move {
4237 run_e2e_with_retry(
4238 state,
4239 &shell,
4240 &e2e_commands,
4241 &worktree,
4242 budget,
4243 &context,
4244 )
4245 .await
4246 }
4247 },
4248 )
4249 .await;
4250 (e2e, verify_retried, false, None)
4251 };
4252
4253 let expected = records.len();
4254 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4255 let incomplete = answered < expected;
4256 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4257 let policy = self.state.config.graph.incomplete_review;
4258 let clean = round_is_clean(
4259 blocking,
4260 e2e_ok,
4261 answered,
4262 expected,
4263 round_quota_missing,
4264 policy,
4265 );
4266
4267 let mut round_record = ReviewRound {
4268 round,
4269 head: head.clone(),
4270 verified_head: None,
4271 verified_at: None,
4272 reviews: records,
4273 e2e,
4274 verify_retried,
4275 e2e_deferred,
4276 e2e_defer_reason,
4277 fix: None,
4278 blocking,
4279 answered,
4280 expected,
4281 clean,
4282 progressed: false,
4283 vote_split,
4284 reconsideration,
4285 verdict: round_verdict,
4286 };
4287 if !matches!(
4298 round_record.e2e_status(),
4299 E2eStatus::Deferred | E2eStatus::NotConfigured
4300 ) {
4301 round_record.verified_head = Some(head.clone());
4302 round_record.verified_at = Some(Timestamp::now());
4303 }
4304 let this_round_verification = round_record.verification_summary(&head);
4305
4306 if incomplete {
4307 let missing: Vec<String> = round_record
4308 .reviews
4309 .iter()
4310 .filter(|r| r.failed.is_some())
4311 .map(|r| format!("review-{}", r.reviewer))
4312 .collect();
4313 self.state.event(
4314 "review",
4315 format!(
4316 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4317 missing.join(", ")
4318 ),
4319 );
4320 }
4321
4322 if clean {
4323 self.state.event(
4324 "review",
4325 if incomplete && policy == IncompleteReviewPolicy::Warn {
4326 format!(
4327 "round {round}: clean (warn policy, incomplete panel) — no \
4328 blocking findings from the seats that answered, verification green"
4329 )
4330 } else if incomplete {
4331 format!(
4332 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4333 quorum) — no blocking findings from the seats that answered, \
4334 verification green",
4335 expected - answered
4336 )
4337 } else {
4338 format!("round {round}: clean — no blocking findings, verification green")
4339 },
4340 );
4341 self.state.reviews.push(round_record);
4342 self.state.status = RunStatus::Gating;
4343 self.state.save()?;
4344 return Ok(());
4345 }
4346
4347 if incomplete && blocking == 0 && e2e_ok {
4355 self.state.reviews.push(round_record);
4356 self.state.save()?;
4357 if round == max_rounds {
4358 self.state.status = RunStatus::Blocked;
4359 self.state.event(
4360 "review",
4361 format!(
4362 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4363 refusing to call it clean",
4364 expected - answered
4365 ),
4366 );
4367 return Ok(());
4368 }
4369 continue;
4370 }
4371
4372 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4384 self.state.reviews.push(round_record);
4385 return self
4386 .stop_reviewing(
4387 "the round's own verification could not run",
4388 &shell,
4389 &winner.worktree,
4390 )
4391 .await;
4392 }
4393
4394 if round == max_rounds {
4395 self.state.reviews.push(round_record);
4396 return self
4397 .stop_reviewing(
4398 &format!(
4399 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4400 ),
4401 &shell,
4402 &winner.worktree,
4403 )
4404 .await;
4405 }
4406
4407 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4410 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4411 let blocking_findings: Vec<_> = all_findings
4412 .iter()
4413 .filter(|f| f.severity.blocks())
4414 .cloned()
4415 .collect();
4416 let job = SeatJob {
4417 prompt: prompt::fix(
4418 &self.state.instruction,
4419 &blocking_findings,
4420 this_round_verification.as_ref(),
4421 round,
4422 max_rounds,
4423 &language,
4424 ),
4425 spec: fix_spec.clone(),
4426 seat,
4427 cwd: winner.worktree.clone(),
4428 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4429 allow_write: true,
4430 sessions,
4431 artifacts: artifacts.clone(),
4432 stem: format!("fix-{round}"),
4433 };
4434 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4435 let cache = self.state.config.cache_dir();
4436 let ctx = WaveCtx {
4437 run: &run_id,
4438 node: "fix",
4439 prompts: &prompts,
4440 cache: cache.as_deref(),
4441 round: Some(round),
4442 };
4443 let (seat, out) =
4444 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4445 let agent_id = seat.agent.clone();
4446
4447 let mut fix = FixRecord {
4448 agent: agent_id,
4449 addressed: Vec::new(),
4450 rejected: Vec::new(),
4451 notes: String::new(),
4452 committed: false,
4453 failed: None,
4454 duration_ms: 0,
4455 continuation: None,
4456 };
4457 let mut continuation = ContinuationRecord::not_needed();
4458 let mut final_seat = seat.clone();
4459 match out {
4460 AgentOutcome::Ok(o) => {
4461 fix.duration_ms = o.duration_ms;
4462 let parsed = verdict::extract_json::<FixReport>(&o.text);
4463 let incomplete_reason = match &parsed {
4470 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4471 "the reply parsed, but it reported a command whose own CLI never \
4472 confirmed an exit status"
4473 .to_owned(),
4474 ),
4475 Ok(_) => None,
4476 Err(e) => Some(e.to_string()),
4477 };
4478 match incomplete_reason {
4479 None => {
4480 let report = parsed.expect("checked Ok above");
4481 fix.addressed = report.addressed;
4482 fix.rejected = report.rejected;
4483 fix.notes =
4484 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4485 }
4486 Some(reason) => {
4487 let (resumed_seat, resolved, failure, cont) = self
4488 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4489 .await;
4490 fix.duration_ms += cont.cumulative_wait_ms;
4491 continuation = cont;
4492 final_seat = resumed_seat;
4493 match resolved {
4494 Some(report) => {
4495 fix.addressed = report.addressed;
4496 fix.rejected = report.rejected;
4497 fix.notes = blind::sanitize_prose(
4498 &report.notes,
4499 &self.state.config.blind,
4500 );
4501 }
4502 None => fix.failed = failure,
4503 }
4504 }
4505 }
4506 }
4507 AgentOutcome::Dropped(o) => {
4509 fix.duration_ms = o.duration_ms;
4510 let why = o
4511 .dropped
4512 .as_ref()
4513 .map(|d| d.why.as_str())
4514 .unwrap_or("the CLI ended the stream without delivering its answer");
4515 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4516 }
4517 AgentOutcome::Quota(o) => {
4518 self.state.quota.push(QuotaLoss {
4519 seat: final_seat.key.clone(),
4520 node: "fix".to_owned(),
4521 at: Timestamp::now(),
4522 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4523 });
4524 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4525 }
4526 AgentOutcome::Failed(e) => fix.failed = Some(e),
4527 }
4528 fix.continuation = Some(continuation);
4529 self.state.seats.insert(final_seat.key.clone(), final_seat);
4530 if let Ok(r) = git::rescue_commit(
4531 &winner.worktree,
4532 &format!("magi: review round {round} fixes (uncommitted work)"),
4533 )
4534 .await
4535 {
4536 self.state.note_withheld("fix", &r.withheld);
4537 }
4538 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4539 fix.committed = after != before;
4540 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4548 let progressed = diff_after != patch;
4549 let commit_note = if fix.committed {
4550 "committed"
4551 } else {
4552 "NO new commit"
4553 };
4554 let tree_note = if progressed {
4555 "changed vs base"
4556 } else {
4557 "unchanged vs base"
4558 };
4559 self.state.event(
4560 "fix",
4561 match &fix.failed {
4562 Some(reason) => {
4568 format!(
4569 "round {round}: fixer's adoption report was lost ({reason}); \
4570 {commit_note}, tree {tree_note}"
4571 )
4572 }
4573 None => format!(
4574 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4575 {tree_note}{}",
4576 fix.addressed.len(),
4577 fix.rejected.len(),
4578 if continuation.outcome == ContinuationOutcome::Resumed {
4579 format!(
4580 " (adoption report recovered after {} continuation(s))",
4581 continuation.attempts
4582 )
4583 } else {
4584 String::new()
4585 },
4586 ),
4587 },
4588 );
4589 round_record.fix = Some(fix);
4590 round_record.progressed = progressed;
4591 self.state.reviews.push(round_record);
4592 self.state.save()?;
4593
4594 if matches!(
4607 continuation.outcome,
4608 ContinuationOutcome::Exhausted
4609 | ContinuationOutcome::QuotaLost
4610 | ContinuationOutcome::NoSession
4611 ) {
4612 return self
4613 .stop_reviewing(
4614 "the fixer's adoption report never came back, even after resuming its \
4615 own seat; refusing to start another round against the same worktree \
4616 while that is unresolved",
4617 &shell,
4618 &winner.worktree,
4619 )
4620 .await;
4621 }
4622
4623 let streak = self
4624 .state
4625 .reviews
4626 .iter()
4627 .rev()
4628 .take_while(|r| !r.progressed)
4629 .count();
4630 if streak >= STAGNANT_LIMIT {
4631 return self
4632 .stop_reviewing(
4633 &format!(
4634 "the tree has not moved against base for {streak} round(s) in a row"
4635 ),
4636 &shell,
4637 &winner.worktree,
4638 )
4639 .await;
4640 }
4641 }
4642 Ok(())
4643 }
4644
4645 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4675 let round_idx = self.state.reviews.len() - 1;
4676 let needs_catchup_run = matches!(
4684 self.state.reviews[round_idx].e2e_status(),
4685 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4686 );
4687 if needs_catchup_run {
4688 let round = self.state.reviews[round_idx].round;
4689 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4690 let commands = self.state.config.verify.e2e.clone();
4691 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4692 let cache_dir = self.state.config.cache_dir();
4693 let context = format!(
4694 "round {round}: verification unresolved, catching up before the final decision"
4695 );
4696 let (outcomes, verify_retried) = with_cache_lease(
4697 &mut self.state,
4698 cache_dir.as_deref(),
4699 "e2e",
4700 "e2e",
4701 worktree,
4702 &attempted_head,
4703 timeout,
4704 &context,
4705 |state, budget| {
4706 let shell = shell.to_vec();
4707 let commands = commands.clone();
4708 let context = context.clone();
4709 async move {
4710 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4711 .await
4712 }
4713 },
4714 )
4715 .await;
4716 let last = &mut self.state.reviews[round_idx];
4717 last.e2e = outcomes;
4718 last.verify_retried = verify_retried;
4719 last.verified_head = Some(attempted_head);
4726 last.verified_at = Some(Timestamp::now());
4727 if verify_inconclusive(&last.e2e) {
4728 self.state.save()?;
4735 return Ok(());
4736 }
4737 last.e2e_deferred = false;
4738 }
4739 let last = &self.state.reviews[round_idx];
4740 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4741
4742 match last.e2e_status() {
4743 E2eStatus::Failed => {
4744 let red: Vec<String> = last
4745 .e2e
4746 .iter()
4747 .filter(|o| !o.ok())
4748 .map(|o| {
4749 format!(
4750 "`{}` -> {:?}\n{}",
4751 o.command,
4752 o.code,
4753 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4754 )
4755 })
4756 .collect();
4757 self.state
4758 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4759 self.state.status = RunStatus::Blocked;
4760 }
4761 E2eStatus::ResourceBlocked => {
4766 self.state.event(
4767 "review",
4768 format!(
4769 "{why}; e2e could not run (shared build cache unavailable); not \
4770 deciding yet"
4771 ),
4772 );
4773 }
4774 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4775 self.state.event(
4776 "review",
4777 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4778 );
4779 self.state.status = RunStatus::Gating;
4780 }
4781 }
4782 self.state.save()?;
4783 Ok(())
4784 }
4785
4786 async fn gate(&mut self) -> Result<()> {
4789 if self.state.status == RunStatus::Failed
4801 || self
4802 .state
4803 .base_sync
4804 .as_ref()
4805 .is_some_and(|s| s.conflict.is_some())
4806 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
4807 != Some(RunStatus::Gating)
4808 {
4809 return Ok(());
4810 }
4811 if self.state.gate_ran {
4812 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
4823 self.state.status = RunStatus::Blocked;
4824 self.state.save()?;
4825 }
4826 return Ok(());
4827 }
4828 let Some(winner) = self.state.winner().cloned() else {
4829 return Ok(());
4830 };
4831 self.state.status = RunStatus::Gating;
4832 let mut outcomes = self.run_gate(&winner).await?;
4833 loop {
4834 if verify_inconclusive(&outcomes) {
4845 self.state.save()?;
4846 return Ok(());
4847 }
4848 if outcomes.iter().all(CommandOutcome::ok) {
4849 break;
4850 }
4851 match self.gate_fix_round(&winner, &outcomes).await? {
4852 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
4853 GateFix::Stop => break,
4854 GateFix::Defer => {
4855 self.state.save()?;
4856 return Ok(());
4857 }
4858 }
4859 }
4860 let passed = outcomes.iter().all(CommandOutcome::ok);
4861 self.state.gate = outcomes;
4862 self.state.gate_ran = true;
4863 if !passed {
4864 self.state.status = RunStatus::Blocked;
4865 let spent = self.state.gate_fixes.len();
4866 self.state.event(
4867 "gate",
4868 if spent == 0 {
4869 "gate failed; not merging".to_owned()
4870 } else {
4871 format!("gate failed after {spent} gate-fix round(s); not merging")
4872 },
4873 );
4874 }
4875 self.state.save()?;
4876 Ok(())
4877 }
4878
4879 async fn run_pre_gate(&mut self, winner: &Candidate) {
4889 let commands = self.state.config.verify.pre_gate.clone();
4890 if commands.is_empty() {
4891 return;
4892 }
4893 let shell = self.state.config.shell();
4894 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4895 let (outcomes, _) = run_commands(
4896 &mut self.state,
4897 "pre_gate",
4898 "pre_gate",
4899 0,
4900 &shell,
4901 &commands,
4902 &winner.worktree,
4903 timeout,
4904 )
4905 .await;
4906 for o in &outcomes {
4907 if !o.ok() {
4908 tracing::warn!(
4909 "pre_gate `{}` failed ({:?}); the gate decides",
4910 o.command,
4911 o.code
4912 );
4913 }
4914 self.state.event(
4915 "pre_gate",
4916 format!(
4917 "`{}` -> {}",
4918 o.command,
4919 if o.ok() {
4920 "pass".to_owned()
4921 } else {
4922 format!(
4923 "FAIL ({:?})\n{}",
4924 o.code,
4925 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4926 )
4927 }
4928 ),
4929 );
4930 }
4931 self.state.pre_gate = outcomes;
4932 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
4933 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
4934 Ok(head) => {
4935 self.state
4936 .event("pre_gate", format!("committed mechanical fixes ({head})"));
4937 self.state.pre_gate_commit = Some(head);
4938 }
4939 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
4940 },
4941 Ok(false) => {}
4942 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
4943 }
4944 if let Err(e) = self.state.save() {
4945 tracing::warn!("could not persist the pre_gate record: {e:#}");
4946 }
4947 }
4948
4949 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
4952 self.run_pre_gate(winner).await;
4953 let shell = self.state.config.shell();
4954 let gate_commands = self.state.config.verify.gate.clone();
4955 let outcomes = if gate_commands.is_empty() {
4964 Vec::new()
4965 } else {
4966 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4967 let cache_dir = self.state.config.cache_dir();
4968 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4969 let (outcomes, _) = with_cache_lease(
4970 &mut self.state,
4971 cache_dir.as_deref(),
4972 "gate",
4973 "gate",
4974 &winner.worktree,
4975 &head,
4976 timeout,
4977 "final gate",
4978 |state, budget| {
4979 let shell = shell.clone();
4980 let gate_commands = gate_commands.clone();
4981 let worktree = winner.worktree.clone();
4982 async move {
4983 let (outcomes, timed_out_pids) = run_commands(
4984 state,
4985 "gate",
4986 "gate",
4987 0,
4988 &shell,
4989 &gate_commands,
4990 &worktree,
4991 budget,
4992 )
4993 .await;
4994 (outcomes, false, timed_out_pids)
4995 }
4996 },
4997 )
4998 .await;
4999 outcomes
5000 };
5001 if outcomes.is_empty() {
5002 self.state.event(
5007 "gate",
5008 "no gate commands configured; nothing to check, passing",
5009 );
5010 }
5011 for o in &outcomes {
5012 self.state.event(
5013 "gate",
5014 format!(
5015 "`{}` -> {}",
5016 o.command,
5017 if o.ok() {
5018 "pass".to_owned()
5019 } else {
5020 format!(
5021 "FAIL ({:?})\n{}",
5022 o.code,
5023 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5024 )
5025 }
5026 ),
5027 );
5028 }
5029 Ok(outcomes)
5030 }
5031
5032 async fn gate_fix_round(
5045 &mut self,
5046 winner: &Candidate,
5047 outcomes: &[CommandOutcome],
5048 ) -> Result<GateFix> {
5049 let cap = self.state.config.graph.gate_fix_rounds;
5050 let spent = self.state.gate_fixes.len();
5051 if spent >= cap {
5052 if cap > 0 {
5053 self.state.event(
5054 "gate",
5055 format!("{spent} gate-fix round(s) spent and the gate still fails"),
5056 );
5057 }
5058 return Ok(GateFix::Stop);
5059 }
5060 if !gate_fixable(outcomes) {
5061 self.state.event(
5062 "gate",
5063 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
5064 command or similar); not spending a fix round on it",
5065 );
5066 return Ok(GateFix::Stop);
5067 }
5068 let min_free = self.state.config.disk.min_free_bytes;
5069 if min_free > 0 {
5070 match crate::disk::free_bytes(&winner.worktree) {
5071 Ok(free) if crate::disk::enough_space(free, min_free) => {}
5072 Ok(free) => {
5073 self.state.event(
5074 "gate",
5075 format!(
5076 "only {free} bytes free ({min_free} required by `[disk] \
5077 min_free_bytes`); not spending a fix round on a failure the disk \
5078 may explain"
5079 ),
5080 );
5081 return Ok(GateFix::Stop);
5082 }
5083 Err(e) => {
5084 self.state.event(
5085 "gate",
5086 format!("free disk space could not be measured ({e:#}); no fix round"),
5087 );
5088 return Ok(GateFix::Stop);
5089 }
5090 }
5091 }
5092
5093 let attempt = spent + 1;
5094 let run_id = self.state.id.clone();
5095 let prompts = self.state.config.prompts.clone();
5096 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
5097 let base = self.landing_base();
5098 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
5099 let seat = self.seat(&fix_seat_key, &fix_spec.id);
5100 let job = SeatJob {
5101 prompt: prompt::gate_fix(
5102 &self.state.instruction,
5103 &failed,
5104 attempt,
5105 cap,
5106 &self.state.config.graph.language,
5107 ),
5108 spec: fix_spec,
5109 seat,
5110 cwd: winner.worktree.clone(),
5111 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
5112 allow_write: true,
5113 sessions: self.state.config.graph.sessions,
5114 artifacts: agent::artifacts_dir(&self.state.dir()),
5115 stem: format!("gate-fix-{attempt}"),
5116 };
5117 self.state.event(
5118 "gate",
5119 format!("gate failed; gate-fix round {attempt} of {cap}"),
5120 );
5121 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
5122 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
5123 let cache = self.state.config.cache_dir();
5124 let ctx = WaveCtx {
5125 run: &run_id,
5126 node: "gate-fix",
5127 prompts: &prompts,
5128 cache: cache.as_deref(),
5129 round: None,
5130 };
5131 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
5132 let mut record = GateFixRecord {
5133 agent: seat.agent.clone(),
5134 failed,
5135 notes: String::new(),
5136 committed: false,
5137 error: None,
5138 };
5139 match out {
5140 AgentOutcome::Ok(o) => {
5141 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
5144 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
5145 }
5146 }
5147 AgentOutcome::Dropped(_) => {
5148 record.error = Some("the CLI dropped the stream".to_owned());
5149 }
5150 AgentOutcome::Quota(o) => {
5151 self.state.quota.push(QuotaLoss {
5152 seat: seat.key.clone(),
5153 node: "gate-fix".to_owned(),
5154 at: Timestamp::now(),
5155 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5156 });
5157 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5158 }
5159 AgentOutcome::Failed(e) => record.error = Some(e),
5160 }
5161 self.state.seats.insert(seat.key.clone(), seat);
5162 if let Ok(r) = git::rescue_commit(
5163 &winner.worktree,
5164 &format!("magi: gate fix {attempt} (uncommitted work)"),
5165 )
5166 .await
5167 {
5168 self.state.note_withheld("gate-fix", &r.withheld);
5169 }
5170 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5171 record.committed = after != before;
5172 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5173 let note = record.error.clone();
5174 self.state.gate_fixes.push(record);
5175 self.state.save()?;
5176 if !changed {
5177 self.state.event(
5178 "gate",
5179 match note {
5180 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5181 None => format!("gate-fix round {attempt}: the tree did not change"),
5182 },
5183 );
5184 return Ok(GateFix::Stop);
5185 }
5186 self.state.event(
5187 "gate",
5188 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5189 );
5190
5191 let commands = self.state.config.verify.e2e.clone();
5192 if !commands.is_empty() {
5193 let shell = self.state.config.shell();
5194 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5195 let cache_dir = self.state.config.cache_dir();
5196 let context = format!("gate-fix round {attempt}");
5197 let (e2e, _) = with_cache_lease(
5198 &mut self.state,
5199 cache_dir.as_deref(),
5200 "e2e",
5201 "e2e",
5202 &winner.worktree,
5203 &after,
5204 timeout,
5205 &context,
5206 |state, budget| {
5207 let shell = shell.clone();
5208 let commands = commands.clone();
5209 let context = context.clone();
5210 let worktree = winner.worktree.clone();
5211 async move {
5212 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5213 .await
5214 }
5215 },
5216 )
5217 .await;
5218 if verify_inconclusive(&e2e) {
5219 return Ok(GateFix::Defer);
5220 }
5221 if e2e.iter().any(|o| !o.ok()) {
5222 self.state.event(
5223 "gate",
5224 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5225 );
5226 return Ok(GateFix::Stop);
5227 }
5228 }
5229 Ok(GateFix::Retry)
5230 }
5231
5232 async fn merge(&mut self) -> Result<()> {
5235 if self
5250 .state
5251 .base_sync
5252 .as_ref()
5253 .is_some_and(|s| s.conflict.is_some())
5254 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5255 != Some(RunStatus::Gating)
5256 || !self.state.gate_status().ok()
5265 {
5266 return Ok(());
5267 }
5268 if self.state.merge.is_some() {
5277 return Ok(());
5278 }
5279 let Some(winner) = self.state.winner().cloned() else {
5280 return Ok(());
5281 };
5282 let repo = self.state.repo.clone();
5283 let base = self.state.base_branch.clone();
5284 let mode = self.state.config.merge.mode;
5285 let style = self.state.config.merge.style;
5286 let pr = pr_message(&self.state, winner.label);
5287 let message = pr.commit_message();
5288
5289 let outcome = match mode {
5290 MergeMode::None => MergeOutcome {
5291 mode,
5292 ok: true,
5293 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5294 },
5295 MergeMode::Local => {
5296 let on = git::current_branch(&repo).await?;
5297 if on.as_deref() != Some(base.as_str()) {
5298 MergeOutcome {
5299 mode,
5300 ok: false,
5301 detail: format!(
5302 "{} has {} checked out, not the base branch {base}",
5303 repo.display(),
5304 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5305 ),
5306 }
5307 } else if !git::is_clean(&repo).await? {
5308 MergeOutcome {
5309 mode,
5310 ok: false,
5311 detail: format!("{} is dirty; refusing to merge", repo.display()),
5312 }
5313 } else {
5314 let out = match style {
5315 MergeStyle::Merge => {
5316 git::merge_no_ff(&repo, &winner.branch, &message).await?
5317 }
5318 MergeStyle::Squash => {
5319 git::merge_squash(&repo, &winner.branch, &message).await?
5320 }
5321 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5322 };
5323 MergeOutcome {
5324 mode,
5325 ok: out.ok(),
5326 detail: if out.ok() { out.stdout } else { out.stderr },
5327 }
5328 }
5329 }
5330 MergeMode::Pr => {
5331 let remote = self.state.config.merge.remote.clone();
5332 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5333 if !pushed.ok() {
5334 MergeOutcome {
5335 mode,
5336 ok: false,
5337 detail: pushed.stderr,
5338 }
5339 } else {
5340 let out =
5341 gh_pr_create(&winner.worktree, &base, &winner.branch, &pr.title, &pr.body)
5342 .await;
5343 match out {
5344 Ok(url) => MergeOutcome {
5345 mode,
5346 ok: true,
5347 detail: url,
5348 },
5349 Err(e) => MergeOutcome {
5350 mode,
5351 ok: false,
5352 detail: e.to_string(),
5353 },
5354 }
5355 }
5356 }
5357 };
5358
5359 self.state.status = match (mode, outcome.ok) {
5360 (MergeMode::None, _) => RunStatus::Ready,
5361 (_, true) => RunStatus::Merged,
5362 (_, false) => RunStatus::Blocked,
5363 };
5364 self.state.event(
5365 "merge",
5366 format!(
5367 "{:?}: {}",
5368 mode,
5369 outcome.detail.lines().next().unwrap_or("")
5370 ),
5371 );
5372 self.state.merge = Some(outcome);
5373 self.state.save()?;
5374
5375 if self.state.config.graph.land
5381 && mode == MergeMode::Pr
5382 && self.state.status == RunStatus::Merged
5383 {
5384 self.run_land().await?;
5385 }
5386 self.settle_questions();
5391 Ok(())
5392 }
5393
5394 async fn run_land(&mut self) -> Result<()> {
5405 let url = self
5406 .state
5407 .merge
5408 .as_ref()
5409 .map(|m| m.detail.clone())
5410 .unwrap_or_default();
5411 let url = url.lines().next().unwrap_or("").trim().to_owned();
5412 if !url.starts_with("http") {
5413 return Ok(());
5414 }
5415 match land::land(&mut self.state, &url).await {
5418 Ok(pr) if self.state.parked => {
5419 let _ = pr;
5423 }
5424 Ok(pr) => {
5425 self.state.status = match pr.state {
5426 land::PrLifecycle::Merged => RunStatus::Merged,
5427 _ => RunStatus::Blocked,
5428 };
5429 if bump::should_release_bump(self.state.status)
5436 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5437 {
5438 self.state
5444 .event("bump", format!("release bump skipped: {e:#}"));
5445 }
5446 self.state.save()?;
5447 }
5448 Err(e) => {
5449 self.state.status = RunStatus::Blocked;
5450 self.state.event("land", format!("gave up: {e}"));
5451 self.state.save()?;
5452 }
5453 }
5454 Ok(())
5455 }
5456
5457 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5461 if let Some(existing) = self.state.seats.get(key)
5462 && existing.agent == agent
5463 {
5464 return existing.clone();
5465 }
5466 let fresh = SeatState::new(key, agent, self.state.seed);
5467 self.state.seats.insert(key.to_owned(), fresh.clone());
5468 fresh
5469 }
5470
5471 fn view(&self, c: &Candidate) -> CandidateView {
5473 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5474 .unwrap_or_default();
5475 let (patch, _) = blind::sanitize_patch(
5476 &format!("candidate {} patch", c.label),
5477 &raw,
5478 &self.state.config.blind,
5479 );
5480 CandidateView {
5481 label: c.label,
5482 branch: c.branch.clone(),
5483 summary: c.summary.clone(),
5484 stat: c.stat.clone(),
5485 patch,
5486 }
5487 }
5488
5489 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5491 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5492 prompt::judge(
5493 "(see above)",
5494 &views,
5495 self.roles.judges.len(),
5496 base_short,
5497 "en",
5498 )
5499 }
5500
5501 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5508 let mut turns = Vec::new();
5509 for j in &self.state.judgements {
5510 if j.ranking.is_empty() {
5511 continue;
5512 }
5513 let reasons = j
5514 .reasons
5515 .iter()
5516 .map(|(k, v)| format!("- {k}: {v}"))
5517 .collect::<Vec<_>>()
5518 .join("\n");
5519 turns.push(Turn {
5520 who: format!("Judge {} (opening ranking)", j.judge),
5521 is_self: j.judge == self_idx + 1,
5522 body: format!(
5523 "Ranked {}{}{reasons}",
5524 j.ranking.iter().collect::<String>(),
5525 if reasons.is_empty() {
5526 ""
5527 } else {
5528 ", because:\n"
5529 }
5530 ),
5531 });
5532 }
5533 for t in self
5534 .state
5535 .deliberation
5536 .iter()
5537 .flat_map(|r| r.turns.iter())
5538 .chain(current)
5539 {
5540 turns.push(Turn {
5541 who: format!("Judge {}", t.judge),
5542 is_self: t.judge == self_idx + 1,
5543 body: t.body.clone(),
5544 });
5545 }
5546 turns
5547 }
5548}
5549
5550fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5552 agent::has_session(spec.kind, seat, sessions)
5553}
5554
5555fn next_untried_implementer<'a>(
5576 roster: &'a [AgentSpec],
5577 start: usize,
5578 tried: &BTreeSet<String>,
5579) -> Option<&'a AgentSpec> {
5580 roster
5581 .get(start + 1..)?
5582 .iter()
5583 .find(|s| !tried.contains(&s.id))
5584}
5585
5586fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5600 commands.iter().any(|c| c.exit_code.is_none())
5601}
5602
5603fn verified_noop_claim(
5616 usable: bool,
5617 commands: &[agent::CommandEvidence],
5618 text: &str,
5619) -> Option<String> {
5620 (usable && !has_unconfirmed_command(commands))
5621 .then(|| verdict::verified_noop(text))
5622 .flatten()
5623}
5624
5625fn short(commit: &str) -> String {
5626 commit.chars().take(7).collect()
5627}
5628
5629fn make_executable(path: &Path) -> Result<()> {
5630 #[cfg(unix)]
5631 {
5632 use std::os::unix::fs::PermissionsExt as _;
5633 let mut perms = std::fs::metadata(path)?.permissions();
5634 perms.set_mode(0o755);
5635 std::fs::set_permissions(path, perms)?;
5636 }
5637 #[cfg(not(unix))]
5638 {
5639 let _ = path;
5640 }
5641 Ok(())
5642}
5643
5644struct WaveCtx<'a> {
5651 run: &'a str,
5654 node: &'a str,
5656 prompts: &'a Prompts,
5657 cache: Option<&'a Path>,
5659 round: Option<usize>,
5662}
5663
5664async fn run_one(
5666 job: SeatJob,
5667 sem: Arc<Semaphore>,
5668 ctx: &WaveCtx<'_>,
5669 state: &mut RunState,
5670 attempt: usize,
5671) -> (SeatState, AgentOutcome) {
5672 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5673 .await
5674 .pop()
5675 .expect("one job in, one result out");
5676 (seat, out)
5677}
5678
5679async fn wave(
5685 jobs: Vec<SeatJob>,
5686 sem: Arc<Semaphore>,
5687 ctx: &WaveCtx<'_>,
5688 state: &mut RunState,
5689 attempt: usize,
5690) -> Vec<(usize, SeatState, AgentOutcome)> {
5691 let WaveCtx {
5692 run,
5693 node,
5694 prompts,
5695 cache,
5696 round,
5697 } = *ctx;
5698 for job in &jobs {
5699 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5700 }
5701 if let Err(e) = state.save() {
5702 tracing::warn!("could not persist in-progress seats: {e:#}");
5707 }
5708 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5724 let wait_started = Instant::now();
5725 let cache_guard = if let Some(cache_dir) = cache {
5726 if jobs_had_a_writer {
5727 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5728 let budget = jobs
5729 .iter()
5730 .map(|j| j.timeout)
5731 .max()
5732 .unwrap_or(Duration::from_secs(60));
5733 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5734 .await
5735 .ok()
5736 } else {
5737 None
5738 }
5739 } else {
5740 None
5741 };
5742 let waited_for_lease = wait_started.elapsed();
5749 let mut set = tokio::task::JoinSet::new();
5750 let overlay = prompts.overlay(node);
5751 for (i, mut job) in jobs.into_iter().enumerate() {
5752 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5753 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
5754 if cache.is_some() {
5755 job.prompt.push('\n');
5756 job.prompt
5757 .push_str(&prompt::build_cache_note(node, job.allow_write));
5758 }
5759 let sem = Arc::clone(&sem);
5760 let run = run.to_owned();
5761 let node = node.to_owned();
5762 let cache = cache
5773 .filter(|_| job.allow_write && cache_guard.is_some())
5774 .map(Path::to_path_buf);
5775 set.spawn(async move {
5776 let _permit = sem.acquire().await;
5777 let mut seat = job.seat;
5778 let out = agent::invoke(
5779 &job.spec,
5780 &mut seat,
5781 &Invocation {
5782 cwd: &job.cwd,
5783 prompt: &job.prompt,
5784 timeout: job.timeout,
5785 allow_write: job.allow_write,
5786 sessions: job.sessions,
5787 artifacts: &job.artifacts,
5788 stem: &job.stem,
5789 run: &run,
5790 node: &node,
5791 cache_dir: cache.as_deref(),
5792 attachments: &[],
5793 },
5794 )
5795 .await;
5796 let out = match out {
5797 Ok(o) if o.usable() => AgentOutcome::Ok(o),
5798 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
5799 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
5807 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
5808 Ok(o) => AgentOutcome::Failed(format!(
5809 "exited with {:?} and no usable output",
5810 o.exit_code
5811 )),
5812 Err(e) => AgentOutcome::Failed(e.to_string()),
5813 };
5814 (i, seat, out)
5815 });
5816 }
5817 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
5818 while let Some(joined) = set.join_next().await {
5819 let (i, seat, out) = match joined {
5820 Ok(v) => v,
5821 Err(e) => {
5825 tracing::error!("agent task panicked: {e}");
5826 continue;
5827 }
5828 };
5829 state.seat_finished(&seat.key);
5830 record_jobs(state, node, round, &seat.key, &out);
5831 if let Err(e) = state.save() {
5832 tracing::warn!("could not persist a seat's completion: {e:#}");
5833 }
5834 if collected.len() <= i {
5835 collected.resize_with(i + 1, || None);
5836 }
5837 collected[i] = Some((i, seat, out));
5838 }
5839 if state
5845 .active
5846 .values()
5847 .any(|a| a.node == node && a.attempt == attempt)
5848 {
5849 state
5850 .active
5851 .retain(|_, a| !(a.node == node && a.attempt == attempt));
5852 if let Err(e) = state.save() {
5853 tracing::warn!("could not persist the end of a wave: {e:#}");
5854 }
5855 }
5856 if let Some(cache_dir) = cache
5863 && jobs_had_a_writer
5864 {
5865 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
5866 }
5867 if let Some(guard) = cache_guard {
5868 guard.release();
5869 }
5870 collected.into_iter().flatten().collect()
5871}
5872
5873fn record_jobs(
5884 state: &mut RunState,
5885 node: &str,
5886 round: Option<usize>,
5887 seat: &str,
5888 out: &AgentOutcome,
5889) {
5890 let commands: &[agent::CommandEvidence] = match out {
5891 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
5892 AgentOutcome::Failed(_) => &[],
5893 };
5894 let checked_at = Timestamp::now();
5895 for c in commands {
5896 state.jobs.push(JobRecord {
5897 node: node.to_owned(),
5898 round,
5899 seat: seat.to_owned(),
5900 id: c.id.clone(),
5901 description: c.description.clone(),
5902 checked_at,
5903 status: match c.exit_code {
5904 Some(0) => JobStatus::Completed,
5905 Some(_) => JobStatus::Failed,
5906 None => JobStatus::Unknown,
5907 },
5908 exit_code: c.exit_code,
5909 result_summary: c.result_summary.clone(),
5910 source: c.source.clone(),
5911 });
5912 }
5913}
5914
5915fn round_is_clean(
5936 blocking: usize,
5937 e2e_ok: bool,
5938 answered: usize,
5939 expected: usize,
5940 quota_missing: usize,
5941 policy: IncompleteReviewPolicy,
5942) -> bool {
5943 if blocking != 0 || !e2e_ok {
5944 return false;
5945 }
5946 if answered == expected || policy == IncompleteReviewPolicy::Warn {
5947 return true;
5948 }
5949 answered > 0 && expected - answered <= quota_missing
5950}
5951
5952fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
5976 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
5977 return Some(RunStatus::Gating);
5978 }
5979 let last = reviews.last()?;
5980 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
5981 if reviews.len() < max_rounds && !stagnant {
5982 return None;
5983 }
5984 if last.incomplete() && last.blocking == 0 {
5985 return Some(RunStatus::Blocked);
5986 }
5987 if last.e2e_status() == E2eStatus::ResourceBlocked {
5988 return None;
5989 }
5990 Some(if last.e2e.iter().all(CommandOutcome::ok) {
5991 RunStatus::Gating
5992 } else {
5993 RunStatus::Blocked
5994 })
5995}
5996
5997fn retry_budget(full: Duration, nudged: bool) -> Duration {
6012 if nudged {
6013 (full / 4).max(Duration::from_secs(120)).min(full)
6014 } else {
6015 full
6016 }
6017}
6018
6019#[allow(clippy::too_many_arguments)]
6032async fn ask_json_wave<T>(
6033 jobs: Vec<SeatJob>,
6034 sem: Arc<Semaphore>,
6035 retries: usize,
6036 ctx: &WaveCtx<'_>,
6037 losses: &mut Vec<QuotaLoss>,
6038 state: &mut RunState,
6039 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
6040) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
6041where
6042 T: serde::de::DeserializeOwned + Send + 'static,
6043{
6044 let n = jobs.len();
6045 let originals: Vec<SeatJob> = jobs;
6046 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
6047 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
6048 let mut attempts_used: Vec<usize> = vec![0; n];
6055 let mut pending: Vec<usize> = (0..n).collect();
6056
6057 for attempt in 0..=retries {
6058 if pending.is_empty() {
6059 break;
6060 }
6061 let mut batch = Vec::with_capacity(pending.len());
6062 for &i in &pending {
6063 let src = &originals[i];
6064 let (prompt, timeout) = if attempt == 0 {
6067 (src.prompt.clone(), src.timeout)
6068 } else {
6069 let why = done[i]
6070 .as_ref()
6071 .and_then(|r| r.as_ref().err().map(ToString::to_string))
6072 .unwrap_or_else(|| "no parsable answer".to_owned());
6073 let nudge = prompt::nudge(&why);
6074 let nudged = has_context(&src.spec, &seats[i], src.sessions);
6075 let prompt = if nudged {
6076 nudge
6077 } else {
6078 format!("{}\n\n---\n\n{}", src.prompt, nudge)
6079 };
6080 (prompt, retry_budget(src.timeout, nudged))
6081 };
6082 batch.push(SeatJob {
6083 spec: src.spec.clone(),
6084 seat: seats[i].clone(),
6085 cwd: src.cwd.clone(),
6086 prompt,
6087 timeout,
6088 allow_write: src.allow_write,
6089 sessions: src.sessions,
6090 artifacts: src.artifacts.clone(),
6091 stem: if attempt == 0 {
6092 src.stem.clone()
6093 } else {
6094 format!("{}-retry{attempt}", src.stem)
6095 },
6096 });
6097 }
6098
6099 if attempt > 0 {
6100 let seats_out: Vec<&str> = pending
6101 .iter()
6102 .map(|&i| originals[i].seat.key.as_str())
6103 .collect();
6104 state.event(
6105 ctx.node,
6106 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
6107 );
6108 }
6109 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
6110 let mut still = Vec::new();
6111 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
6112 seats[i] = seat;
6113 let (parsed, quota) = match out {
6114 AgentOutcome::Ok(o) => (
6115 match verdict::extract_json::<T>(&o.text) {
6116 Ok(v) => match validate(&v) {
6117 Ok(()) => Ok((v, o)),
6118 Err(e) => Err(e),
6119 },
6120 Err(e) => Err(e),
6121 },
6122 false,
6123 ),
6124 AgentOutcome::Quota(o) => {
6125 losses.push(QuotaLoss {
6126 seat: originals[i].seat.key.clone(),
6127 node: ctx.node.to_owned(),
6128 at: Timestamp::now(),
6129 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
6130 });
6131 (
6132 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
6133 true,
6134 )
6135 }
6136 AgentOutcome::Dropped(o) => {
6141 let why = o
6142 .dropped
6143 .as_ref()
6144 .map(|d| d.why.as_str())
6145 .unwrap_or("the CLI ended the stream without delivering its answer");
6146 (
6147 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
6148 false,
6149 )
6150 }
6151 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
6152 };
6153 let failed = parsed.is_err();
6154 done[i] = Some(parsed);
6155 attempts_used[i] = attempt;
6156 if failed && !quota {
6159 still.push(i);
6160 }
6161 }
6162 pending = still;
6163 }
6164
6165 seats
6166 .into_iter()
6167 .zip(done)
6168 .zip(attempts_used)
6169 .map(|((seat, res), attempts)| {
6170 (
6171 seat,
6172 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6173 attempts,
6174 )
6175 })
6176 .collect()
6177}
6178
6179async fn acquire_cache_lease(
6192 state: &mut RunState,
6193 cache_dir: &Path,
6194 owner: &crate::cache::Owner,
6195 budget: Duration,
6196 context: &str,
6197) -> Result<crate::cache::Guard> {
6198 let home = crate::run::home();
6199 let started = Instant::now();
6200 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6201 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6202 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6203 Err(e) => {
6204 state.event(
6205 "verify",
6206 format!("{context}: could not check the shared build cache: {e:#}"),
6207 );
6208 if let Err(e2) = state.save() {
6209 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6210 }
6211 return Err(e);
6212 }
6213 };
6214 state.event(
6215 "verify",
6216 format!(
6217 "{context}: waiting for the shared build cache at {} ({})",
6218 cache_dir.display(),
6219 busy.describe()
6220 ),
6221 );
6222 if let Err(e) = state.save() {
6223 tracing::warn!("could not persist a cache wait: {e:#}");
6224 }
6225 let remaining = budget.saturating_sub(started.elapsed());
6226 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6227 Ok(g) => Ok(g),
6228 Err(e) => {
6229 state.event("verify", format!("{context}: {e:#}"));
6230 if let Err(e2) = state.save() {
6231 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6232 }
6233 Err(e)
6234 }
6235 }
6236}
6237
6238#[allow(clippy::too_many_arguments)]
6259async fn with_cache_lease<'s, F, Fut>(
6260 state: &'s mut RunState,
6261 cache_dir: Option<&Path>,
6262 node: &str,
6263 seat: &str,
6264 worktree: &Path,
6265 head: &str,
6266 budget: Duration,
6267 context: &str,
6268 body: F,
6269) -> (Vec<CommandOutcome>, bool)
6270where
6271 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6272 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6273{
6274 let Some(cache_dir) = cache_dir else {
6275 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6276 return (outcomes, retried);
6277 };
6278 let home = crate::run::home();
6279 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6280 let started = Instant::now();
6281 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6282 Ok(g) => g,
6283 Err(e) => {
6284 return (
6285 vec![CommandOutcome {
6286 command: "(waiting for the shared build cache)".to_owned(),
6287 code: None,
6288 output_tail: e.to_string(),
6289 duration_ms: started.elapsed().as_millis() as u64,
6290 resource_blocked: true,
6291 }],
6292 false,
6293 );
6294 }
6295 };
6296 let identity = crate::cache::Identity::new(worktree, head);
6297 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6298 state.event(
6306 "verify",
6307 format!(
6308 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6309 worktree.display(),
6310 short(head)
6311 ),
6312 );
6313 guard.release();
6314 return (
6315 vec![CommandOutcome {
6316 command: "(confirming the shared build cache is fresh)".to_owned(),
6317 code: None,
6318 output_tail: e.to_string(),
6319 duration_ms: started.elapsed().as_millis() as u64,
6320 resource_blocked: true,
6321 }],
6322 false,
6323 );
6324 }
6325 let remaining = budget.saturating_sub(started.elapsed());
6326 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6327 if !timed_out_pids.is_empty() {
6332 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6333 }
6334 guard.release();
6335 (outcomes, retried)
6336}
6337
6338async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6350 wait_for_pids_with(
6351 pids,
6352 crate::proc::pid_alive,
6353 LEASE_RELEASE_POLL,
6354 LEASE_RELEASE_MAX_WAIT,
6355 )
6356 .await;
6357}
6358
6359async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6365 pids: &[u32],
6366 alive: F,
6367 poll: Duration,
6368 max_wait: Duration,
6369) {
6370 let deadline = Instant::now() + max_wait;
6371 loop {
6372 if pids.iter().all(|&pid| !alive(pid)) {
6373 return;
6374 }
6375 if Instant::now() >= deadline {
6376 return;
6377 }
6378 tokio::time::sleep(poll).await;
6379 }
6380}
6381
6382fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6389 outcomes.iter().any(|o| o.resource_blocked)
6390}
6391
6392enum GateFix {
6394 Retry,
6396 Stop,
6399 Defer,
6402}
6403
6404fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6412 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6413 red.peek().is_some()
6414 && red.all(|o| {
6415 !o.resource_blocked
6416 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6417 && !o.output_tail.trim().is_empty()
6418 })
6419}
6420
6421fn e2e_outcome_label(o: &CommandOutcome) -> String {
6425 if o.ok() {
6426 return "pass".to_owned();
6427 }
6428 let reason = if o.build_failed() {
6429 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6430 } else {
6431 format!("FAIL ({:?})", o.code)
6432 };
6433 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6434}
6435
6436async fn run_e2e_with_retry(
6444 state: &mut RunState,
6445 shell: &[String],
6446 commands: &[String],
6447 worktree: &Path,
6448 timeout: Duration,
6449 context: &str,
6450) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6451 let (mut e2e, mut timed_out_pids) = run_commands(
6452 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6453 )
6454 .await;
6455 for o in &e2e {
6456 state.event(
6457 "verify",
6458 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6459 );
6460 }
6461 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6465 if verify_retried {
6466 state.event(
6467 "verify",
6468 format!(
6469 "{context}: verify could not build/link, not a test result — retrying once \
6470 before concluding"
6471 ),
6472 );
6473 let retried = run_commands(
6474 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6475 )
6476 .await;
6477 e2e = retried.0;
6478 timed_out_pids.extend(retried.1);
6481 for o in &e2e {
6482 state.event(
6483 "verify",
6484 format!(
6485 "{context}: retry `{}` -> {}",
6486 o.command,
6487 e2e_outcome_label(o)
6488 ),
6489 );
6490 }
6491 }
6492 (e2e, verify_retried, timed_out_pids)
6493}
6494
6495#[allow(clippy::too_many_arguments)]
6511async fn run_commands(
6512 state: &mut RunState,
6513 node: &str,
6514 task: &str,
6515 attempt: usize,
6516 shell: &[String],
6517 commands: &[String],
6518 cwd: &Path,
6519 timeout: Duration,
6520) -> (Vec<CommandOutcome>, Vec<u32>) {
6521 if commands.is_empty() {
6522 return (Vec::new(), Vec::new());
6527 }
6528 let mut out = Vec::new();
6529 let mut timed_out_pids = Vec::new();
6530 let total = commands.len();
6531 for (idx, command) in commands.iter().enumerate() {
6532 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6533 if let Err(e) = state.save() {
6534 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6535 }
6536 let started = Instant::now();
6537 let mut cmd = tokio::process::Command::new(&shell[0]);
6538 cmd.quiet();
6539 cmd.args(&shell[1..])
6540 .arg(command)
6541 .current_dir(cwd)
6542 .stdin(std::process::Stdio::null())
6543 .stdout(std::process::Stdio::piped())
6544 .stderr(std::process::Stdio::piped())
6545 .kill_on_drop(true);
6546 let spawned = cmd.spawn();
6547 let (code, body) = match spawned {
6548 Ok(child) => {
6549 let pid = child.id();
6554 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6555 Ok(Ok(o)) => {
6556 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6557 body.push_str(&String::from_utf8_lossy(&o.stderr));
6558 (o.status.code(), body)
6559 }
6560 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6561 Err(_) => {
6562 if let Some(pid) = pid {
6563 timed_out_pids.push(pid);
6564 }
6565 (None, format!("timed out after {}s", timeout.as_secs()))
6566 }
6567 }
6568 }
6569 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6570 };
6571 out.push(CommandOutcome {
6572 command: command.clone(),
6573 code,
6574 output_tail: tail(&body, OUTPUT_TAIL),
6575 duration_ms: started.elapsed().as_millis() as u64,
6576 resource_blocked: false,
6577 });
6578 }
6579 state.task_finished(task);
6580 if let Err(e) = state.save() {
6581 tracing::warn!("could not persist the end of {task}: {e:#}");
6582 }
6583 (out, timed_out_pids)
6584}
6585
6586fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6598 let repo = repo.display();
6599 match style {
6600 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6601 MergeStyle::Squash => {
6602 let subject = message
6605 .lines()
6606 .next()
6607 .unwrap_or(branch)
6608 .replace(['\\', '"', '$', '`'], "");
6609 format!(
6610 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6611 )
6612 }
6613 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6614 }
6615}
6616
6617const PR_TITLE_MAX: usize = 240;
6628
6629struct PrMessage {
6634 title: String,
6635 body: String,
6636}
6637
6638impl PrMessage {
6639 fn commit_message(&self) -> String {
6643 format!("{}\n\n{}", self.title, self.body)
6644 }
6645}
6646
6647fn title_marker(line: &str) -> Option<&str> {
6649 let line = line.trim();
6650 let head = line.get(..6)?;
6651 head.eq_ignore_ascii_case("title:")
6652 .then(|| line[6..].trim())
6653}
6654
6655fn summary_title(summary: &str) -> Option<String> {
6660 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6661 let raw = title_marker(first)?;
6662 if raw.is_empty() {
6663 return None;
6664 }
6665 let title = queue::title_from(raw, PR_TITLE_MAX);
6666 let lower = title.to_ascii_lowercase();
6667 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6668 return None;
6669 }
6670 Some(title)
6671}
6672
6673fn summary_without_title(summary: &str) -> String {
6676 let mut lines = summary.trim().lines().peekable();
6677 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6678 lines.next();
6679 }
6680 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6681}
6682
6683fn pr_message(state: &RunState, winner: char) -> PrMessage {
6697 let summary = state
6698 .candidates
6699 .iter()
6700 .find(|c| c.label == winner)
6701 .map(|c| c.summary.as_str())
6702 .unwrap_or_default();
6703 let title = summary_title(summary).unwrap_or_else(|| {
6706 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6707 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6708 t
6709 } else {
6710 format!(
6711 "chore: land candidate {} of run {}",
6712 winner.to_ascii_uppercase(),
6713 state.id
6714 )
6715 }
6716 });
6717
6718 let mut body = String::new();
6719 let what = summary_without_title(summary);
6720 if !what.is_empty() {
6721 body.push_str("## Summary\n\n");
6722 body.push_str(&what);
6723 body.push_str("\n\n");
6724 }
6725
6726 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6727 if let Some(fix) = fix
6728 && !fix.notes.trim().is_empty()
6729 {
6730 body.push_str("## Review fixes\n\n");
6731 body.push_str(fix.notes.trim());
6732 body.push_str("\n\n");
6733 }
6734
6735 let open = state.open_findings();
6736 if !open.is_empty() {
6737 body.push_str("## Open review findings\n\n");
6738 for f in &open {
6739 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6740 }
6741 body.push('\n');
6742 }
6743
6744 if let Some(fix) = fix
6745 && !fix.rejected.is_empty()
6746 {
6747 body.push_str("## Declined by the fixer\n\n");
6748 for r in &fix.rejected {
6749 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
6750 }
6751 body.push('\n');
6752 }
6753
6754 let task = state.instruction.trim();
6755 let task = if task.is_empty() {
6756 "(empty task)"
6757 } else {
6758 task
6759 };
6760 body.push_str(&format!(
6761 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
6762 task.replace("</details>", "</details>")
6763 ));
6764
6765 body.push_str(&format!(
6766 "\n---\nmagi:run/{} magi:candidate-{}\n",
6767 state.id,
6768 winner.to_ascii_lowercase()
6769 ));
6770
6771 let id = crate::scrub::Identity::current();
6774 PrMessage {
6775 title: crate::scrub::scrub(&title, &id),
6776 body: crate::scrub::scrub(&body, &id),
6777 }
6778}
6779
6780async fn gh_pr_create(
6782 cwd: &Path,
6783 base: &str,
6784 head: &str,
6785 title: &str,
6786 body: &str,
6787) -> Result<String> {
6788 let out = tokio::process::Command::new("gh")
6789 .args([
6790 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
6791 ])
6792 .current_dir(cwd)
6793 .quiet()
6794 .stdin(std::process::Stdio::null())
6795 .output()
6796 .await
6797 .context("spawn gh")?;
6798 if out.status.success() {
6799 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
6800 } else {
6801 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
6802 }
6803}
6804
6805pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
6814 let repo = state.repo.clone();
6815 let root = state.worktree_root();
6816 let winner = state.tally.as_ref().map(|t| t.winner);
6817 let mut removed = Vec::new();
6818
6819 for i in 0..state.candidates.len() {
6820 let c = state.candidates[i].clone();
6821 let is_winner = Some(c.label) == winner;
6822 if is_winner && !drop_winner {
6823 continue;
6824 }
6825 if c.worktree.exists() {
6826 git::worktree_remove(&repo, &c.worktree).await.ok();
6827 removed.push(c.worktree.to_string_lossy().into_owned());
6828 }
6829 let handed_over = state.released_branches.contains(&c.branch);
6832 if !handed_over && git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
6833 git::branch_delete(&repo, &c.branch).await.ok();
6834 removed.push(c.branch.clone());
6835 }
6836 state.candidates[i].folded = true;
6837 }
6838
6839 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
6840 let path = name.path();
6841 let keep = !drop_winner
6842 && winner.is_some_and(|w| {
6843 path.file_name()
6844 .is_some_and(|n| n == format!("cand-{w}").as_str())
6845 });
6846 if keep {
6847 continue;
6848 }
6849 git::worktree_remove(&repo, &path).await.ok();
6850 removed.push(path.to_string_lossy().into_owned());
6851 }
6852
6853 remove_if_empty(&root);
6862
6863 if state.enabled_worktree_config && drop_winner {
6864 git::release_worktree_config(&repo).await.ok();
6868 state.enabled_worktree_config = false;
6869 }
6870 state.save_under(home)?;
6871 Ok(removed)
6872}
6873
6874fn remove_if_empty(dir: &Path) {
6885 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
6886 std::fs::remove_dir(dir).ok();
6887 }
6888}
6889
6890pub fn worst_open(state: &RunState) -> Option<Severity> {
6892 state
6893 .reviews
6894 .last()?
6895 .reviews
6896 .iter()
6897 .flat_map(|r| r.findings.iter())
6898 .map(|f| f.severity)
6899 .max()
6900}
6901
6902#[cfg(test)]
6903mod tests {
6904 use super::*;
6905 use crate::run::GateStatus;
6906 use std::collections::BTreeMap;
6907 use std::time::Duration;
6908
6909 fn conductor() -> AgentSpec {
6910 AgentSpec {
6911 id: "conductor".to_owned(),
6912 kind: crate::config::AgentKind::Command,
6913 model: None,
6914 command: vec!["true".to_owned()],
6915 extra_args: Vec::new(),
6916 env: BTreeMap::new(),
6917 prompt_delivery: None,
6918 }
6919 }
6920
6921 fn spec(id: &str) -> AgentSpec {
6922 AgentSpec {
6923 id: id.to_owned(),
6924 kind: crate::config::AgentKind::Command,
6925 model: None,
6926 command: vec!["true".to_owned()],
6927 extra_args: Vec::new(),
6928 env: BTreeMap::new(),
6929 prompt_delivery: None,
6930 }
6931 }
6932
6933 #[test]
6940 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
6941 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6942 let tried = BTreeSet::from(["beta".to_owned()]);
6943 let next = next_untried_implementer(&roster, 1, &tried);
6946 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
6947 }
6948
6949 #[test]
6950 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
6951 let roster = vec![spec("alpha"), spec("beta")];
6952 let tried = BTreeSet::from(["beta".to_owned()]);
6953 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6957 }
6958
6959 #[test]
6960 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
6961 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6962 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
6963 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6967 }
6968
6969 #[test]
6970 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
6971 let roster = vec![spec("a"), spec("a"), spec("b")];
6972 let tried = BTreeSet::from(["a".to_owned()]);
6973 let next = next_untried_implementer(&roster, 0, &tried);
6974 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
6975 }
6976
6977 #[test]
6978 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
6979 let roster = vec![spec("a"), spec("b")];
6980 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
6981 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
6982 }
6983
6984 #[test]
6985 fn remove_if_empty_only_ever_takes_a_bare_directory() {
6986 let dir = tempfile::tempdir().unwrap();
6987 let bay = dir.path().join("ffff");
6988
6989 remove_if_empty(&bay);
6991 assert!(!bay.exists());
6992
6993 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
6996 remove_if_empty(&bay);
6997 assert!(bay.exists(), "non-empty directory must survive");
6998
6999 std::fs::remove_dir(bay.join("cand-A")).unwrap();
7001 remove_if_empty(&bay);
7002 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
7003 }
7004
7005 #[test]
7014 fn a_full_panel_that_found_nothing_is_clean() {
7015 assert!(round_is_clean(
7016 0,
7017 true,
7018 2,
7019 2,
7020 0,
7021 IncompleteReviewPolicy::Block
7022 ));
7023 }
7024
7025 #[test]
7026 fn a_missing_seat_is_never_clean_under_the_default_policy() {
7027 assert!(!round_is_clean(
7028 0,
7029 true,
7030 1,
7031 2,
7032 0,
7033 IncompleteReviewPolicy::Block
7034 ));
7035 }
7036
7037 #[test]
7038 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
7039 assert!(!round_is_clean(
7040 1,
7041 true,
7042 1,
7043 2,
7044 0,
7045 IncompleteReviewPolicy::Warn
7046 ));
7047 }
7048
7049 #[test]
7050 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
7051 assert!(round_is_clean(
7052 0,
7053 true,
7054 1,
7055 2,
7056 0,
7057 IncompleteReviewPolicy::Warn
7058 ));
7059 }
7060
7061 #[test]
7062 fn a_full_panel_with_an_open_finding_is_not_clean() {
7063 assert!(!round_is_clean(
7064 1,
7065 true,
7066 2,
7067 2,
7068 0,
7069 IncompleteReviewPolicy::Block
7070 ));
7071 }
7072
7073 #[test]
7074 fn a_full_panel_with_a_red_e2e_is_not_clean() {
7075 assert!(!round_is_clean(
7076 0,
7077 false,
7078 2,
7079 2,
7080 0,
7081 IncompleteReviewPolicy::Block
7082 ));
7083 }
7084
7085 #[test]
7092 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
7093 assert!(round_is_clean(
7096 0,
7097 true,
7098 1,
7099 2,
7100 1,
7101 IncompleteReviewPolicy::Block
7102 ));
7103 }
7104
7105 #[test]
7106 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
7107 assert!(!round_is_clean(
7110 0,
7111 true,
7112 1,
7113 2,
7114 0,
7115 IncompleteReviewPolicy::Block
7116 ));
7117 }
7118
7119 #[test]
7120 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
7121 assert!(!round_is_clean(
7122 1,
7123 true,
7124 1,
7125 2,
7126 1,
7127 IncompleteReviewPolicy::Block
7128 ));
7129 assert!(!round_is_clean(
7130 0,
7131 false,
7132 1,
7133 2,
7134 1,
7135 IncompleteReviewPolicy::Block
7136 ));
7137 }
7138
7139 #[test]
7140 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
7141 assert!(!round_is_clean(
7145 0,
7146 true,
7147 0,
7148 2,
7149 2,
7150 IncompleteReviewPolicy::Block
7151 ));
7152 }
7153
7154 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
7155 CommandOutcome {
7156 command: "test".to_owned(),
7157 code,
7158 output_tail: String::new(),
7159 duration_ms: 0,
7160 resource_blocked,
7161 }
7162 }
7163
7164 #[test]
7165 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7166 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7167 assert!(
7168 !verify_inconclusive(&[outcome(Some(1), false)]),
7169 "an ordinary failure is still evidence about the patch"
7170 );
7171 assert!(verify_inconclusive(&[outcome(None, true)]));
7172 assert!(
7173 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7174 "one inconclusive outcome taints the whole batch"
7175 );
7176 assert!(!verify_inconclusive(&[]));
7177 }
7178
7179 #[tokio::test]
7180 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7181 let calls = std::sync::atomic::AtomicUsize::new(0);
7185 let started = Instant::now();
7186 wait_for_pids_with(
7187 &[123],
7188 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7189 Duration::from_millis(5),
7190 Duration::from_secs(5),
7191 )
7192 .await;
7193 assert!(
7194 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7195 "must keep checking rather than deciding on the first answer"
7196 );
7197 assert!(
7198 started.elapsed() < Duration::from_secs(1),
7199 "must return the moment it is confirmed dead, not wait out the ceiling"
7200 );
7201 }
7202
7203 #[tokio::test]
7204 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7205 let started = Instant::now();
7206 wait_for_pids_with(
7207 &[123],
7208 |_| true, Duration::from_millis(5),
7210 Duration::from_millis(30),
7211 )
7212 .await;
7213 let elapsed = started.elapsed();
7214 assert!(
7215 elapsed >= Duration::from_millis(30),
7216 "must not give up before its own ceiling: {elapsed:?}"
7217 );
7218 assert!(
7219 elapsed < Duration::from_secs(1),
7220 "must not wait past its own ceiling either: {elapsed:?}"
7221 );
7222 }
7223
7224 #[tokio::test]
7225 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7226 let started = Instant::now();
7227 wait_for_pids_with(
7228 &[],
7229 |_| true,
7230 Duration::from_secs(5),
7231 Duration::from_secs(5),
7232 )
7233 .await;
7234 assert!(
7235 started.elapsed() < Duration::from_millis(200),
7236 "an empty pid list has nothing to confirm"
7237 );
7238 }
7239
7240 fn review_round(
7246 clean: bool,
7247 blocking: usize,
7248 answered: usize,
7249 expected: usize,
7250 progressed: bool,
7251 e2e_ok: bool,
7252 ) -> ReviewRound {
7253 ReviewRound {
7254 round: 1,
7255 head: "h".to_owned(),
7256 verified_head: None,
7257 verified_at: None,
7258 reviews: Vec::new(),
7259 e2e: vec![CommandOutcome {
7260 command: "test".to_owned(),
7261 code: Some(if e2e_ok { 0 } else { 1 }),
7262 output_tail: String::new(),
7263 duration_ms: 0,
7264 resource_blocked: false,
7265 }],
7266 verify_retried: false,
7267 e2e_deferred: false,
7268 e2e_defer_reason: None,
7269 fix: None,
7270 blocking,
7271 answered,
7272 expected,
7273 clean,
7274 progressed,
7275 vote_split: false,
7276 reconsideration: Vec::new(),
7277 verdict: None,
7278 }
7279 }
7280
7281 #[test]
7282 fn review_conclusion_is_none_when_nothing_has_run() {
7283 assert_eq!(review_conclusion(&[], 3), None);
7284 }
7285
7286 #[test]
7287 fn review_conclusion_is_none_while_rounds_remain() {
7288 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7289 assert_eq!(review_conclusion(&rounds, 3), None);
7290 }
7291
7292 #[test]
7293 fn review_conclusion_is_gating_once_a_round_is_clean() {
7294 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7295 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7296 }
7297
7298 #[test]
7299 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7300 let rounds = vec![
7301 review_round(false, 1, 2, 2, true, true),
7302 review_round(false, 1, 2, 2, true, true),
7303 ];
7304 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7305 }
7306
7307 #[test]
7308 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7309 let rounds = vec![
7310 review_round(false, 1, 2, 2, true, true),
7311 review_round(false, 1, 2, 2, true, false),
7312 ];
7313 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7314 }
7315
7316 #[test]
7317 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7318 let mut blocked = review_round(false, 1, 2, 2, true, false);
7325 blocked.e2e[0].resource_blocked = true;
7326 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7327 assert_eq!(review_conclusion(&rounds, 2), None);
7328 }
7329
7330 #[test]
7331 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7332 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7334 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7335 }
7336
7337 #[test]
7338 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7339 let rounds = vec![
7340 review_round(false, 1, 2, 2, false, true),
7341 review_round(false, 1, 2, 2, false, true),
7342 ];
7343 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7344 }
7345
7346 fn secs(n: u64) -> Duration {
7347 Duration::from_secs(n)
7348 }
7349
7350 fn init_repo(dir: &Path) {
7353 let run = |args: &[&str]| {
7354 let out = std::process::Command::new("git")
7355 .args(args)
7356 .current_dir(dir)
7357 .quiet()
7358 .output()
7359 .expect("spawn git");
7360 assert!(
7361 out.status.success(),
7362 "git {args:?} failed: {}",
7363 String::from_utf8_lossy(&out.stderr)
7364 );
7365 };
7366 run(&["init", "-b", "main"]);
7367 run(&["config", "user.name", "magi test"]);
7368 run(&["config", "user.email", "magi@example.com"]);
7369 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7370 run(&["add", "-A"]);
7371 run(&["commit", "-m", "init"]);
7372 }
7373
7374 fn ask_test_home() {
7382 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7383 }
7384
7385 fn runner_at(status: RunStatus) -> Runner {
7388 let mut state = RunState::new(
7389 PathBuf::from("/nonexistent/repo"),
7390 "main".to_owned(),
7391 "deadbeef".to_owned(),
7392 "task".to_owned(),
7393 Config::default(),
7394 );
7395 state.status = status;
7396 Runner {
7397 state,
7398 roles: ResolvedRoles {
7399 implementers: Vec::new(),
7400 judges: Vec::new(),
7401 reviewers: Vec::new(),
7402 fixer: None,
7403 conductor: conductor(),
7404 implementer_roster: Vec::new(),
7405 },
7406 sem: Arc::new(Semaphore::new(1)),
7407 pause: Pause::new(),
7408 interrupt: Pause::new(),
7409 }
7410 }
7411
7412 #[test]
7416 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7417 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7418 let mut runner = runner_at(RunStatus::Implementing);
7419 let interrupt = Pause::new();
7420 runner.watch_interrupt(interrupt.clone());
7421
7422 interrupt.park_because("task a1b2 asked to run first");
7423
7424 assert!(runner.park_here().expect("park_here"));
7425 assert!(runner.state.parked);
7426 let last = runner.state.events.last().expect("a park event");
7427 assert_eq!(last.node, "park");
7428 assert!(
7429 last.message.contains("task a1b2 asked to run first"),
7430 "expected the interrupt reason in {:?}",
7431 last.message
7432 );
7433 }
7434
7435 #[test]
7443 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7444 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7445 let mut runner = runner_at(RunStatus::Implementing);
7446 let shutdown = Pause::new();
7447 runner.on_pause(shutdown.clone());
7448 let interrupt = Pause::new();
7449 runner.watch_interrupt(interrupt.clone());
7450
7451 assert!(!runner.park_here().expect("park_here"));
7453 assert!(!runner.state.parked);
7454
7455 interrupt.park_because("test");
7457 assert!(!shutdown.parked());
7458 assert!(runner.park_here().expect("park_here"));
7459 }
7460
7461 #[tokio::test]
7475 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7476 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7477 let mut runner = runner_at(RunStatus::Implementing);
7478 let interrupt = Pause::new();
7479 runner.watch_interrupt(interrupt.clone());
7480
7481 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7482 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7483
7484 let node = async move {
7488 started_tx.send(()).expect("send started");
7489 finish_rx.await.expect("recv finish");
7490 "node finished"
7491 };
7492
7493 let interrupter = async move {
7494 started_rx.await.expect("recv started");
7495 interrupt.park_because("higher-priority task waiting");
7497 tokio::task::yield_now().await;
7501 finish_tx.send(()).expect("send finish");
7502 };
7503
7504 let (node_result, ()) = tokio::join!(node, interrupter);
7505 assert_eq!(
7506 node_result, "node finished",
7507 "the in-flight call ran to completion"
7508 );
7509
7510 assert!(runner.park_here().expect("park_here"));
7513 assert!(runner.state.parked);
7514 }
7515
7516 #[test]
7522 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7523 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7524 let mut runner = runner_at(RunStatus::Judging);
7525 runner.state.config.agents = vec![conductor()];
7529 runner.state.candidates = vec![Candidate {
7530 index: 0,
7531 label: 'A',
7532 agent: "alpha".to_owned(),
7533 branch: "magi/x/A".to_owned(),
7534 worktree: PathBuf::from("/nonexistent/worktree"),
7535 summary: "did the thing".to_owned(),
7536 stat: "1 file changed".to_owned(),
7537 files: 1,
7538 commits: 1,
7539 empty: false,
7540 failed: None,
7541 verified_noop: None,
7542 duration_ms: 1234,
7543 folded: false,
7544 }];
7545 let run_id = runner.state.id.clone();
7546
7547 let interrupt = Pause::new();
7548 runner.watch_interrupt(interrupt.clone());
7549 interrupt.park_because("task c3d4 asked to run first");
7550 assert!(runner.park_here().expect("park_here"));
7551
7552 let resumed = Runner::resume(&run_id).expect("resume");
7553 assert_eq!(resumed.state.candidates.len(), 1);
7554 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7555 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7556 assert_eq!(resumed.state.status, runner.state.status);
7557 assert!(
7558 resumed.state.parked,
7559 "still parked until `execute` actually walks the graph again"
7560 );
7561 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7562 }
7563
7564 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7566 let mut q = ask::Question::new(
7567 run.to_owned(),
7568 "implement".to_owned(),
7569 "impl-A".to_owned(),
7570 "Which storage backend should the cache use?".to_owned(),
7571 String::new(),
7572 vec!["SQLite".to_owned(), "Redis".to_owned()],
7573 );
7574 store.put(&mut q).unwrap();
7575 q
7576 }
7577
7578 #[test]
7579 fn a_failed_runs_open_question_is_abandoned() {
7580 ask_test_home();
7581 let store = ask::Questions::open();
7582 let mut runner = runner_at(RunStatus::Failed);
7583 let run = runner.state.id.clone();
7584 let q = ask_open_question(&store, &run);
7585
7586 runner.settle_questions();
7587
7588 let back = store.get(&q.id).unwrap();
7589 assert!(
7590 !back.status.open(),
7591 "the seat that asked died with the run; nobody is left to read an answer"
7592 );
7593 assert!(
7594 back.detail.contains(&run) && back.detail.contains("failed"),
7595 "the reason names what the run became, not just that it is gone: {}",
7596 back.detail
7597 );
7598 }
7599
7600 #[test]
7601 fn a_merged_runs_open_question_is_abandoned_too() {
7602 ask_test_home();
7603 let store = ask::Questions::open();
7604 for status in [RunStatus::Merged, RunStatus::Ready] {
7607 let mut runner = runner_at(status);
7608 let run = runner.state.id.clone();
7609 let q = ask_open_question(&store, &run);
7610
7611 runner.settle_questions();
7612
7613 let back = store.get(&q.id).unwrap();
7614 assert!(
7615 !back.status.open(),
7616 "{status:?} run's question must not outlive the run"
7617 );
7618 }
7619 }
7620
7621 #[test]
7622 fn a_still_resumable_runs_open_question_is_left_alone() {
7623 ask_test_home();
7624 let store = ask::Questions::open();
7625 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7631 let mut runner = runner_at(status);
7632 let run = runner.state.id.clone();
7633 let q = ask_open_question(&store, &run);
7634
7635 runner.settle_questions();
7636
7637 let back = store.get(&q.id).unwrap();
7638 assert!(
7639 back.status.open(),
7640 "{status:?} is still alive; the question must still be waiting"
7641 );
7642 }
7643 }
7644
7645 #[test]
7646 fn settle_questions_never_touches_an_already_answered_question() {
7647 ask_test_home();
7648 let store = ask::Questions::open();
7649 let mut runner = runner_at(RunStatus::Failed);
7650 let run = runner.state.id.clone();
7651 let mut q = ask_open_question(&store, &run);
7652 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
7653 .unwrap();
7654 store.put(&mut q).unwrap();
7655
7656 runner.settle_questions();
7661 runner.settle_questions();
7662
7663 let back = store.get(&q.id).unwrap();
7664 assert_eq!(
7665 back.status,
7666 ask::QuestionStatus::Answered,
7667 "a real answer is a decision on record, never overwritten by a sweep"
7668 );
7669 }
7670
7671 #[tokio::test]
7682 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
7683 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
7684 let tmp = tempfile::tempdir().expect("tempdir");
7685 let repo = tmp.path().join("repo");
7686 std::fs::create_dir_all(&repo).unwrap();
7687 init_repo(&repo);
7688
7689 let mut config = Config::default();
7690 config.graph.worktree_root = Some(tmp.path().join("wt"));
7691
7692 let mut state = RunState::new(
7693 repo.clone(),
7694 "main".to_owned(),
7695 "deadbeef".to_owned(),
7696 "task".to_owned(),
7697 config,
7698 );
7699 let root = state.worktree_root();
7700 let wt_a = root.join("cand-A");
7701 let wt_b = root.join("cand-B");
7702 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
7703 .await
7704 .expect("worktree A");
7705 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
7706 .await
7707 .expect("worktree B");
7708
7709 state.candidates = vec![
7710 Candidate {
7711 index: 0,
7712 label: 'A',
7713 agent: "alpha".to_owned(),
7714 branch: "magi/x/A".to_owned(),
7715 worktree: wt_a.clone(),
7716 summary: String::new(),
7717 stat: String::new(),
7718 files: 0,
7719 commits: 0,
7720 empty: false,
7721 failed: None,
7722 verified_noop: None,
7723 duration_ms: 0,
7724 folded: false,
7725 },
7726 Candidate {
7727 index: 1,
7728 label: 'B',
7729 agent: "beta".to_owned(),
7730 branch: "magi/x/B".to_owned(),
7731 worktree: wt_b.clone(),
7732 summary: String::new(),
7733 stat: String::new(),
7734 files: 0,
7735 commits: 0,
7736 empty: false,
7737 failed: None,
7738 verified_noop: None,
7739 duration_ms: 0,
7740 folded: false,
7741 },
7742 ];
7743 state.tally = Some(Tally {
7744 first_choice: BTreeMap::from([('A', 1)]),
7745 borda: BTreeMap::new(),
7746 winner: 'A',
7747 rankings: 1,
7748 unanimous_initial: true,
7749 deliberated: false,
7750 changed_votes: 0,
7751 unanimous_final: true,
7752 tie_break: None,
7753 judges: 1,
7754 present: 1,
7755 quorum: 1,
7756 met_quorum: true,
7757 uncontested: None,
7758 });
7759 state.status = RunStatus::Ready;
7760
7761 fold_run(&mut state, false, &crate::run::home())
7762 .await
7763 .expect("fold_run");
7764
7765 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
7766 assert!(
7767 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7768 "the unmerged winner's branch survives"
7769 );
7770 assert!(
7771 !state.candidates[0].folded,
7772 "the winner is not marked folded"
7773 );
7774
7775 assert!(!wt_b.exists(), "the loser's worktree is removed");
7776 assert!(
7777 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
7778 "the loser's branch is removed"
7779 );
7780 assert!(state.candidates[1].folded, "the loser is marked folded");
7781 }
7782
7783 #[tokio::test]
7786 async fn fold_run_keeps_a_branch_that_was_handed_to_a_later_run() {
7787 let tmp = tempfile::tempdir().expect("tempdir");
7788 let repo = tmp.path().join("repo");
7789 std::fs::create_dir_all(&repo).unwrap();
7790 init_repo(&repo);
7791 let home = tmp.path().join("home");
7792
7793 let mut config = Config::default();
7794 config.graph.worktree_root = Some(tmp.path().join("wt"));
7795 let mut state = RunState::new(
7796 repo.clone(),
7797 "main".to_owned(),
7798 "deadbeef".to_owned(),
7799 "task".to_owned(),
7800 config,
7801 );
7802 git::git(&repo, &["branch", "magi/x/A", "main"])
7804 .await
7805 .expect("branch");
7806 state.candidates = vec![Candidate {
7807 index: 0,
7808 label: 'A',
7809 agent: "alpha".to_owned(),
7810 branch: "magi/x/A".to_owned(),
7811 worktree: state.worktree_root().join("cand-A"),
7812 summary: String::new(),
7813 stat: String::new(),
7814 files: 0,
7815 commits: 0,
7816 empty: false,
7817 failed: None,
7818 verified_noop: None,
7819 duration_ms: 0,
7820 folded: true,
7821 }];
7822 state.released_to = Some("20260901-000000-new1".to_owned());
7823 state.released_branches = vec!["magi/x/A".to_owned()];
7824
7825 fold_run(&mut state, true, &home).await.expect("fold_run");
7826
7827 assert!(
7828 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7829 "the handed-over branch survives a fold"
7830 );
7831 }
7832
7833 #[tokio::test]
7842 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
7843 let tmp = tempfile::tempdir().expect("tempdir");
7844 let repo = tmp.path().join("repo");
7845 std::fs::create_dir_all(&repo).unwrap();
7846 init_repo(&repo);
7847
7848 let mut config = Config::default();
7849 config.merge.mode = MergeMode::Local;
7850
7851 let mut state = RunState::new(
7852 repo.clone(),
7853 "main".to_owned(),
7854 "deadbeef".to_owned(),
7855 "task".to_owned(),
7856 config,
7857 );
7858 state.candidates = vec![Candidate {
7859 index: 0,
7860 label: 'A',
7861 agent: "alpha".to_owned(),
7862 branch: "does-not-exist".to_owned(),
7863 worktree: repo.clone(),
7864 summary: String::new(),
7865 stat: String::new(),
7866 files: 0,
7867 commits: 0,
7868 empty: false,
7869 failed: None,
7870 verified_noop: None,
7871 duration_ms: 0,
7872 folded: false,
7873 }];
7874 state.tally = Some(Tally {
7875 first_choice: BTreeMap::from([('A', 1)]),
7876 borda: BTreeMap::new(),
7877 winner: 'A',
7878 rankings: 1,
7879 unanimous_initial: true,
7880 deliberated: false,
7881 changed_votes: 0,
7882 unanimous_final: true,
7883 tie_break: None,
7884 judges: 0,
7885 present: 0,
7886 quorum: 0,
7887 met_quorum: true,
7888 uncontested: Some("only candidate A produced a change".to_owned()),
7889 });
7890 state.reviews = vec![ReviewRound {
7891 round: 1,
7892 head: "deadbeef".to_owned(),
7893 verified_head: None,
7894 verified_at: None,
7895 reviews: Vec::new(),
7896 e2e: Vec::new(),
7897 fix: None,
7898 blocking: 0,
7899 answered: 0,
7900 expected: 0,
7901 clean: true,
7902 verify_retried: false,
7903 e2e_deferred: false,
7904 e2e_defer_reason: None,
7905 progressed: false,
7906 vote_split: false,
7907 reconsideration: Vec::new(),
7908 verdict: None,
7909 }];
7910 state.gate = vec![CommandOutcome {
7911 command: "test".to_owned(),
7912 code: Some(0),
7913 output_tail: String::new(),
7914 duration_ms: 0,
7915 resource_blocked: false,
7916 }];
7917 state.gate_ran = true;
7918 state.status = RunStatus::Ready;
7923 state.merge = Some(MergeOutcome {
7924 mode: MergeMode::Local,
7925 ok: false,
7926 detail: "already concluded".to_owned(),
7927 });
7928
7929 let mut runner = Runner {
7930 state,
7931 roles: ResolvedRoles {
7932 implementers: Vec::new(),
7933 judges: Vec::new(),
7934 reviewers: Vec::new(),
7935 fixer: None,
7936 conductor: conductor(),
7937 implementer_roster: Vec::new(),
7938 },
7939 sem: Arc::new(Semaphore::new(1)),
7940 pause: Pause::new(),
7941 interrupt: Pause::new(),
7942 };
7943
7944 runner.merge().await.expect("merge");
7945
7946 assert_eq!(
7947 runner.state.status,
7948 RunStatus::Ready,
7949 "a concluded run's status must not change on reentry"
7950 );
7951 assert_eq!(
7952 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
7953 Some("already concluded"),
7954 "merge must not run again once the node already recorded an outcome"
7955 );
7956 }
7957
7958 #[tokio::test]
7967 async fn merge_refuses_a_gate_that_has_not_actually_run() {
7968 let tmp = tempfile::tempdir().expect("tempdir");
7969 let repo = tmp.path().join("repo");
7970 std::fs::create_dir_all(&repo).unwrap();
7971 init_repo(&repo);
7972
7973 let mut config = Config::default();
7974 config.merge.mode = MergeMode::Local;
7975
7976 let mut state = RunState::new(
7977 repo.clone(),
7978 "main".to_owned(),
7979 "deadbeef".to_owned(),
7980 "task".to_owned(),
7981 config,
7982 );
7983 state.candidates = vec![Candidate {
7984 index: 0,
7985 label: 'A',
7986 agent: "alpha".to_owned(),
7987 branch: "does-not-exist".to_owned(),
7988 worktree: repo.clone(),
7989 summary: String::new(),
7990 stat: String::new(),
7991 files: 0,
7992 commits: 0,
7993 empty: false,
7994 failed: None,
7995 verified_noop: None,
7996 duration_ms: 0,
7997 folded: false,
7998 }];
7999 state.tally = Some(Tally {
8000 first_choice: BTreeMap::from([('A', 1)]),
8001 borda: BTreeMap::new(),
8002 winner: 'A',
8003 rankings: 1,
8004 unanimous_initial: true,
8005 deliberated: false,
8006 changed_votes: 0,
8007 unanimous_final: true,
8008 tie_break: None,
8009 judges: 0,
8010 present: 0,
8011 quorum: 0,
8012 met_quorum: true,
8013 uncontested: Some("only candidate A produced a change".to_owned()),
8014 });
8015 state.reviews = vec![ReviewRound {
8016 round: 1,
8017 head: "deadbeef".to_owned(),
8018 verified_head: None,
8019 verified_at: None,
8020 reviews: Vec::new(),
8021 e2e: Vec::new(),
8022 fix: None,
8023 blocking: 0,
8024 answered: 0,
8025 expected: 0,
8026 clean: true,
8027 verify_retried: false,
8028 e2e_deferred: false,
8029 e2e_defer_reason: None,
8030 progressed: false,
8031 vote_split: false,
8032 reconsideration: Vec::new(),
8033 verdict: None,
8034 }];
8035 state.gate = Vec::new();
8037 state.gate_ran = false;
8038 state.status = RunStatus::Gating;
8039
8040 let mut runner = Runner {
8041 state,
8042 roles: ResolvedRoles {
8043 implementers: Vec::new(),
8044 judges: Vec::new(),
8045 reviewers: Vec::new(),
8046 fixer: None,
8047 conductor: conductor(),
8048 implementer_roster: Vec::new(),
8049 },
8050 sem: Arc::new(Semaphore::new(1)),
8051 pause: Pause::new(),
8052 interrupt: Pause::new(),
8053 };
8054
8055 runner.merge().await.expect("merge");
8056
8057 assert!(
8058 runner.state.merge.is_none(),
8059 "an empty gate must never be read as a passing one: {:?}",
8060 runner.state.merge
8061 );
8062 }
8063
8064 #[tokio::test]
8071 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
8072 let tmp = tempfile::tempdir().expect("tempdir");
8073 let repo = tmp.path().join("repo");
8074 std::fs::create_dir_all(&repo).unwrap();
8075 init_repo(&repo);
8076
8077 let config = Config::default();
8079
8080 let mut state = RunState::new(
8081 repo.clone(),
8082 "main".to_owned(),
8083 "deadbeef".to_owned(),
8084 "task".to_owned(),
8085 config,
8086 );
8087 state.candidates = vec![Candidate {
8088 index: 0,
8089 label: 'A',
8090 agent: "alpha".to_owned(),
8091 branch: "does-not-exist".to_owned(),
8092 worktree: repo.clone(),
8093 summary: String::new(),
8094 stat: String::new(),
8095 files: 0,
8096 commits: 0,
8097 empty: false,
8098 failed: None,
8099 verified_noop: None,
8100 duration_ms: 0,
8101 folded: false,
8102 }];
8103 state.tally = Some(Tally {
8104 first_choice: BTreeMap::from([('A', 1)]),
8105 borda: BTreeMap::new(),
8106 winner: 'A',
8107 rankings: 1,
8108 unanimous_initial: true,
8109 deliberated: false,
8110 changed_votes: 0,
8111 unanimous_final: true,
8112 tie_break: None,
8113 judges: 0,
8114 present: 0,
8115 quorum: 0,
8116 met_quorum: true,
8117 uncontested: Some("only candidate A produced a change".to_owned()),
8118 });
8119 state.reviews = vec![ReviewRound {
8120 round: 1,
8121 head: "deadbeef".to_owned(),
8122 verified_head: None,
8123 verified_at: None,
8124 reviews: Vec::new(),
8125 e2e: Vec::new(),
8126 fix: None,
8127 blocking: 0,
8128 answered: 0,
8129 expected: 0,
8130 clean: true,
8131 verify_retried: false,
8132 e2e_deferred: false,
8133 e2e_defer_reason: None,
8134 progressed: false,
8135 vote_split: false,
8136 reconsideration: Vec::new(),
8137 verdict: None,
8138 }];
8139
8140 let mut runner = Runner {
8141 state,
8142 roles: ResolvedRoles {
8143 implementers: Vec::new(),
8144 judges: Vec::new(),
8145 reviewers: Vec::new(),
8146 fixer: None,
8147 conductor: conductor(),
8148 implementer_roster: Vec::new(),
8149 },
8150 sem: Arc::new(Semaphore::new(1)),
8151 pause: Pause::new(),
8152 interrupt: Pause::new(),
8153 };
8154
8155 runner.gate().await.expect("gate");
8156 assert!(
8157 runner.state.gate_ran,
8158 "zero configured commands is still a real attempt, not an unrun gate"
8159 );
8160 assert!(runner.state.gate.is_empty());
8161 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
8162 assert_ne!(
8163 runner.state.status,
8164 RunStatus::Blocked,
8165 "a gate with nothing to check must not read as failed"
8166 );
8167
8168 runner.merge().await.expect("merge");
8169 assert_eq!(
8170 runner.state.status,
8171 RunStatus::Ready,
8172 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
8173 );
8174 }
8175
8176 #[tokio::test]
8187 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
8188 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8189 let home = crate::run::home();
8190
8191 let tmp = tempfile::tempdir().expect("tempdir");
8192 let repo = tmp.path().join("repo");
8193 std::fs::create_dir_all(&repo).unwrap();
8194 init_repo(&repo);
8195 let cache_dir = tmp.path().join("target");
8198
8199 let mut config = Config::default();
8200 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
8201 config.graph.timeout_verify = Some(2);
8204
8205 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8206 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8207 .expect("no io error acquiring directly")
8208 {
8209 crate::cache::AcquireOutcome::Acquired(g) => g,
8210 crate::cache::AcquireOutcome::Busy(b) => {
8211 panic!("expected the direct acquire to win the lease first: {b:?}")
8212 }
8213 };
8214
8215 let mut state = RunState::new(
8216 repo.clone(),
8217 "main".to_owned(),
8218 "deadbeef".to_owned(),
8219 "task".to_owned(),
8220 config,
8221 );
8222 state.candidates = vec![Candidate {
8223 index: 0,
8224 label: 'A',
8225 agent: "alpha".to_owned(),
8226 branch: "does-not-exist".to_owned(),
8227 worktree: repo.clone(),
8228 summary: String::new(),
8229 stat: String::new(),
8230 files: 0,
8231 commits: 0,
8232 empty: false,
8233 failed: None,
8234 verified_noop: None,
8235 duration_ms: 0,
8236 folded: false,
8237 }];
8238 state.tally = Some(Tally {
8239 first_choice: BTreeMap::from([('A', 1)]),
8240 borda: BTreeMap::new(),
8241 winner: 'A',
8242 rankings: 1,
8243 unanimous_initial: true,
8244 deliberated: false,
8245 changed_votes: 0,
8246 unanimous_final: true,
8247 tie_break: None,
8248 judges: 0,
8249 present: 0,
8250 quorum: 0,
8251 met_quorum: true,
8252 uncontested: Some("only candidate A produced a change".to_owned()),
8253 });
8254 state.reviews = vec![ReviewRound {
8255 round: 1,
8256 head: "deadbeef".to_owned(),
8257 verified_head: None,
8258 verified_at: None,
8259 reviews: Vec::new(),
8260 e2e: Vec::new(),
8261 fix: None,
8262 blocking: 0,
8263 answered: 0,
8264 expected: 0,
8265 clean: true,
8266 verify_retried: false,
8267 e2e_deferred: false,
8268 e2e_defer_reason: None,
8269 progressed: false,
8270 vote_split: false,
8271 reconsideration: Vec::new(),
8272 verdict: None,
8273 }];
8274
8275 let mut runner = Runner {
8276 state,
8277 roles: ResolvedRoles {
8278 implementers: Vec::new(),
8279 judges: Vec::new(),
8280 reviewers: Vec::new(),
8281 fixer: None,
8282 conductor: conductor(),
8283 implementer_roster: Vec::new(),
8284 },
8285 sem: Arc::new(Semaphore::new(1)),
8286 pause: Pause::new(),
8287 interrupt: Pause::new(),
8288 };
8289
8290 let started = std::time::Instant::now();
8291 runner.gate().await.expect("gate");
8292 assert!(
8293 started.elapsed() < Duration::from_secs(1),
8294 "a gate with nothing to run must never wait on a lease it never needed"
8295 );
8296 assert!(
8297 runner.state.gate_ran,
8298 "zero commands is still a real, immediate attempt"
8299 );
8300 assert!(runner.state.gate.is_empty());
8301 assert_ne!(
8302 runner.state.status,
8303 RunStatus::Blocked,
8304 "must not read as resource-blocked on a lease it never asked for"
8305 );
8306 }
8307
8308 #[tokio::test]
8318 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8319 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8320
8321 let tmp = tempfile::tempdir().expect("tempdir");
8322 let repo = tmp.path().join("repo");
8323 std::fs::create_dir_all(&repo).unwrap();
8324 init_repo(&repo);
8325
8326 let mut config = Config::default();
8327 config.verify.gate = vec![
8328 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8329 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8330 .to_owned(),
8331 ];
8332
8333 let mut state = RunState::new(
8334 repo.clone(),
8335 "main".to_owned(),
8336 "deadbeef".to_owned(),
8337 "task".to_owned(),
8338 config,
8339 );
8340 let run_id = state.id.clone();
8341 state.candidates = vec![Candidate {
8342 index: 0,
8343 label: 'A',
8344 agent: "alpha".to_owned(),
8345 branch: "does-not-exist".to_owned(),
8346 worktree: repo.clone(),
8347 summary: String::new(),
8348 stat: String::new(),
8349 files: 0,
8350 commits: 0,
8351 empty: false,
8352 failed: None,
8353 verified_noop: None,
8354 duration_ms: 0,
8355 folded: false,
8356 }];
8357 state.tally = Some(Tally {
8358 first_choice: BTreeMap::from([('A', 1)]),
8359 borda: BTreeMap::new(),
8360 winner: 'A',
8361 rankings: 1,
8362 unanimous_initial: true,
8363 deliberated: false,
8364 changed_votes: 0,
8365 unanimous_final: true,
8366 tie_break: None,
8367 judges: 0,
8368 present: 0,
8369 quorum: 0,
8370 met_quorum: true,
8371 uncontested: Some("only candidate A produced a change".to_owned()),
8372 });
8373 state.reviews = vec![ReviewRound {
8374 round: 1,
8375 head: "deadbeef".to_owned(),
8376 verified_head: None,
8377 verified_at: None,
8378 reviews: Vec::new(),
8379 e2e: Vec::new(),
8380 fix: None,
8381 blocking: 0,
8382 answered: 0,
8383 expected: 0,
8384 clean: true,
8385 verify_retried: false,
8386 e2e_deferred: false,
8387 e2e_defer_reason: None,
8388 progressed: false,
8389 vote_split: false,
8390 reconsideration: Vec::new(),
8391 verdict: None,
8392 }];
8393
8394 let mut runner = Runner {
8395 state,
8396 roles: ResolvedRoles {
8397 implementers: Vec::new(),
8398 judges: Vec::new(),
8399 reviewers: Vec::new(),
8400 fixer: None,
8401 conductor: conductor(),
8402 implementer_roster: Vec::new(),
8403 },
8404 sem: Arc::new(Semaphore::new(1)),
8405 pause: Pause::new(),
8406 interrupt: Pause::new(),
8407 };
8408
8409 let started_marker = repo.join("started.marker");
8410 let release_marker = repo.join("release.marker");
8411 let poller = tokio::spawn(async move {
8412 for _ in 0..100 {
8417 if started_marker.exists()
8418 && let Ok(s) = crate::run::RunState::load(&run_id)
8419 && let Some(a) = s.active.get("gate")
8420 {
8421 std::fs::write(&release_marker, b"go").expect("release marker");
8422 return Some(a.clone());
8423 }
8424 tokio::time::sleep(Duration::from_millis(50)).await;
8425 }
8426 None
8427 });
8428
8429 runner.gate().await.expect("gate");
8430 let captured = poller.await.expect("poller task");
8431 let captured = captured.expect(
8432 "the poller never saw a `gate` task entry in run.json while the command was \
8433 still blocked on its own release marker",
8434 );
8435
8436 assert_eq!(captured.task.as_deref(), Some("gate"));
8437 assert_eq!(captured.node, "gate");
8438 assert_eq!(captured.index, Some(1));
8439 assert_eq!(captured.total, Some(1));
8440 assert!(
8441 captured
8442 .command
8443 .as_deref()
8444 .is_some_and(|c| c.contains("started.marker")),
8445 "{captured:?}"
8446 );
8447
8448 assert!(
8449 runner.state.active.is_empty(),
8450 "the entry must be cleared once the command actually finished: {:?}",
8451 runner.state.active
8452 );
8453 assert!(runner.state.gate_ran);
8454 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8455 }
8456
8457 #[tokio::test]
8470 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8471 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8472 let home = crate::run::home();
8473
8474 let tmp = tempfile::tempdir().expect("tempdir");
8475 let repo = tmp.path().join("repo");
8476 std::fs::create_dir_all(&repo).unwrap();
8477 init_repo(&repo);
8478 let head = crate::git::rev_parse(&repo, "HEAD")
8479 .await
8480 .expect("rev-parse");
8481 let cache_dir = tmp.path().join("target");
8484
8485 let mut config = Config::default();
8486 config.verify.e2e = vec![format!(
8487 "CARGO_TARGET_DIR='{}' test -f README.md",
8488 cache_dir.display()
8489 )];
8490 config.graph.review_rounds = 1;
8491 config.graph.timeout_verify = Some(2);
8494
8495 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8496 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8497 .expect("no io error acquiring directly")
8498 {
8499 crate::cache::AcquireOutcome::Acquired(g) => g,
8500 crate::cache::AcquireOutcome::Busy(b) => {
8501 panic!("expected the direct acquire to win the lease first: {b:?}")
8502 }
8503 };
8504
8505 let mut state = RunState::new(
8506 repo.clone(),
8507 "main".to_owned(),
8508 head.clone(),
8509 "task".to_owned(),
8510 config,
8511 );
8512 state.candidates = vec![Candidate {
8513 index: 0,
8514 label: 'A',
8515 agent: "alpha".to_owned(),
8516 branch: "does-not-exist".to_owned(),
8517 worktree: repo.clone(),
8518 summary: String::new(),
8519 stat: String::new(),
8520 files: 0,
8521 commits: 0,
8522 empty: false,
8523 failed: None,
8524 verified_noop: None,
8525 duration_ms: 0,
8526 folded: false,
8527 }];
8528 state.tally = Some(Tally {
8529 first_choice: BTreeMap::from([('A', 1)]),
8530 borda: BTreeMap::new(),
8531 winner: 'A',
8532 rankings: 1,
8533 unanimous_initial: true,
8534 deliberated: false,
8535 changed_votes: 0,
8536 unanimous_final: true,
8537 tie_break: None,
8538 judges: 0,
8539 present: 0,
8540 quorum: 0,
8541 met_quorum: true,
8542 uncontested: Some("only candidate A produced a change".to_owned()),
8543 });
8544 state.reviews = vec![ReviewRound {
8548 round: 1,
8549 head: head.clone(),
8550 verified_head: None,
8551 verified_at: None,
8552 reviews: Vec::new(),
8553 e2e: Vec::new(),
8554 fix: None,
8555 blocking: 1,
8556 answered: 1,
8557 expected: 1,
8558 clean: false,
8559 verify_retried: false,
8560 e2e_deferred: true,
8561 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8562 progressed: false,
8563 vote_split: false,
8564 reconsideration: Vec::new(),
8565 verdict: None,
8566 }];
8567
8568 let mut runner = Runner {
8569 state,
8570 roles: ResolvedRoles {
8571 implementers: Vec::new(),
8572 judges: Vec::new(),
8573 reviewers: Vec::new(),
8574 fixer: None,
8575 conductor: conductor(),
8576 implementer_roster: Vec::new(),
8577 },
8578 sem: Arc::new(Semaphore::new(1)),
8579 pause: Pause::new(),
8580 interrupt: Pause::new(),
8581 };
8582
8583 let shell = runner.state.config.shell();
8584 runner
8585 .stop_reviewing("round budget spent", &shell, &repo)
8586 .await
8587 .expect("stop_reviewing");
8588
8589 let last = runner.state.reviews.last().expect("round record");
8590 assert_eq!(
8591 last.e2e_status(),
8592 E2eStatus::ResourceBlocked,
8593 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8594 failed: {last:?}"
8595 );
8596 assert_eq!(
8597 last.verified_head.as_deref(),
8598 Some(head.as_str()),
8599 "which commit this attempt targeted is known even though nothing finished checking \
8600 it"
8601 );
8602 let first_attempt_at = last
8603 .verified_at
8604 .expect("when this attempt ran is known too");
8605 assert_ne!(
8606 runner.state.status,
8607 RunStatus::Blocked,
8608 "contention is evidence about the machine, not the patch — it must not settle the \
8609 run as blocked: {:?}",
8610 runner.state.status
8611 );
8612 assert!(
8613 !runner
8614 .state
8615 .events
8616 .iter()
8617 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
8618 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
8619 runner.state.events
8620 );
8621
8622 runner
8627 .stop_reviewing("round budget spent", &shell, &repo)
8628 .await
8629 .expect("stop_reviewing retry");
8630 assert_eq!(
8631 runner.state.reviews.len(),
8632 1,
8633 "no new round was started: {:?}",
8634 runner.state.reviews
8635 );
8636 let last = runner.state.reviews.last().expect("round record");
8637 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
8638 assert!(
8639 last.verified_at.expect("still known") > first_attempt_at,
8640 "a second reentry must be a fresh attempt, not a stale copy of the first"
8641 );
8642 assert_ne!(runner.state.status, RunStatus::Blocked);
8643
8644 held.release();
8645 }
8646
8647 #[tokio::test]
8659 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
8660 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8661 let home = crate::run::home();
8662
8663 let tmp = tempfile::tempdir().expect("tempdir");
8664 let repo = tmp.path().join("repo");
8665 std::fs::create_dir_all(&repo).unwrap();
8666 init_repo(&repo);
8667 let head = crate::git::rev_parse(&repo, "HEAD")
8668 .await
8669 .expect("rev-parse");
8670 let cache_dir = tmp.path().join("target");
8671
8672 let mut config = Config::default();
8673 config.verify.e2e = vec![format!(
8674 "CARGO_TARGET_DIR='{}' test -f README.md",
8675 cache_dir.display()
8676 )];
8677 config.graph.review_rounds = 1;
8678 config.graph.timeout_verify = Some(2);
8679
8680 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8681 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8682 .expect("no io error acquiring directly")
8683 {
8684 crate::cache::AcquireOutcome::Acquired(g) => g,
8685 crate::cache::AcquireOutcome::Busy(b) => {
8686 panic!("expected the direct acquire to win the lease first: {b:?}")
8687 }
8688 };
8689
8690 let mut state = RunState::new(
8691 repo.clone(),
8692 "main".to_owned(),
8693 head.clone(),
8694 "task".to_owned(),
8695 config,
8696 );
8697 state.candidates = vec![Candidate {
8698 index: 0,
8699 label: 'A',
8700 agent: "alpha".to_owned(),
8701 branch: "does-not-exist".to_owned(),
8702 worktree: repo.clone(),
8703 summary: String::new(),
8704 stat: String::new(),
8705 files: 0,
8706 commits: 0,
8707 empty: false,
8708 failed: None,
8709 verified_noop: None,
8710 duration_ms: 0,
8711 folded: false,
8712 }];
8713 state.tally = Some(Tally {
8714 first_choice: BTreeMap::from([('A', 1)]),
8715 borda: BTreeMap::new(),
8716 winner: 'A',
8717 rankings: 1,
8718 unanimous_initial: true,
8719 deliberated: false,
8720 changed_votes: 0,
8721 unanimous_final: true,
8722 tie_break: None,
8723 judges: 0,
8724 present: 0,
8725 quorum: 0,
8726 met_quorum: true,
8727 uncontested: Some("only candidate A produced a change".to_owned()),
8728 });
8729 state.reviews = vec![ReviewRound {
8733 round: 1,
8734 head: head.clone(),
8735 verified_head: Some(head.clone()),
8736 verified_at: Some(jiff::Timestamp::now()),
8737 reviews: Vec::new(),
8738 e2e: vec![CommandOutcome {
8739 command: format!(
8740 "CARGO_TARGET_DIR='{}' test -f README.md",
8741 cache_dir.display()
8742 ),
8743 code: None,
8744 output_tail: "waiting for the shared build cache".to_owned(),
8745 duration_ms: 0,
8746 resource_blocked: true,
8747 }],
8748 fix: None,
8749 blocking: 1,
8750 answered: 1,
8751 expected: 1,
8752 clean: false,
8753 verify_retried: false,
8754 e2e_deferred: false,
8755 e2e_defer_reason: None,
8756 progressed: false,
8757 vote_split: false,
8758 reconsideration: Vec::new(),
8759 verdict: None,
8760 }];
8761
8762 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
8763 let mut runner = Runner {
8764 state,
8765 roles: ResolvedRoles {
8766 implementers: Vec::new(),
8767 judges: Vec::new(),
8768 reviewers: Vec::new(),
8769 fixer: None,
8770 conductor: conductor(),
8771 implementer_roster: Vec::new(),
8772 },
8773 sem: Arc::new(Semaphore::new(1)),
8774 pause: Pause::new(),
8775 interrupt: Pause::new(),
8776 };
8777
8778 runner.review_loop().await.expect("review_loop");
8783
8784 assert_eq!(
8785 runner.state.reviews.len(),
8786 1,
8787 "no new round was started on top of the unresolved one: {:?}",
8788 runner.state.reviews
8789 );
8790 let last = &runner.state.reviews[0];
8791 assert_eq!(
8792 last.e2e_status(),
8793 E2eStatus::ResourceBlocked,
8794 "still contended: {last:?}"
8795 );
8796 assert!(
8797 last.verified_at.expect("still known") > first_attempt_at,
8798 "review_loop must have actually retried the check, not left it exactly as found"
8799 );
8800 assert_ne!(
8801 runner.state.status,
8802 RunStatus::Blocked,
8803 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
8804 runner.state.status
8805 );
8806
8807 held.release();
8808 }
8809
8810 #[tokio::test]
8811 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
8812 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8813 let tmp = tempfile::tempdir().expect("tempdir");
8814 let repo = tmp.path().join("repo");
8815 std::fs::create_dir_all(&repo).unwrap();
8816 init_repo(&repo);
8817
8818 let mut config = Config::default();
8819 config.merge.mode = MergeMode::Pr;
8820 config.graph.land = true;
8821 config.graph.land_approval = false;
8822
8823 let mut state = RunState::new(
8824 repo.clone(),
8825 "main".to_owned(),
8826 "deadbeef".to_owned(),
8827 "task".to_owned(),
8828 config,
8829 );
8830 state.candidates = vec![Candidate {
8831 index: 0,
8832 label: 'A',
8833 agent: "alpha".to_owned(),
8834 branch: "does-not-exist".to_owned(),
8835 worktree: repo.clone(),
8836 summary: String::new(),
8837 stat: String::new(),
8838 files: 0,
8839 commits: 0,
8840 empty: false,
8841 failed: None,
8842 verified_noop: None,
8843 duration_ms: 0,
8844 folded: false,
8845 }];
8846 state.tally = Some(Tally {
8847 first_choice: BTreeMap::from([('A', 1)]),
8848 borda: BTreeMap::new(),
8849 winner: 'A',
8850 rankings: 1,
8851 unanimous_initial: true,
8852 deliberated: false,
8853 changed_votes: 0,
8854 unanimous_final: true,
8855 tie_break: None,
8856 judges: 0,
8857 present: 0,
8858 quorum: 0,
8859 met_quorum: true,
8860 uncontested: Some("only candidate A produced a change".to_owned()),
8861 });
8862 state.reviews = vec![ReviewRound {
8863 round: 1,
8864 head: "deadbeef".to_owned(),
8865 verified_head: None,
8866 verified_at: None,
8867 reviews: Vec::new(),
8868 e2e: Vec::new(),
8869 fix: None,
8870 blocking: 0,
8871 answered: 0,
8872 expected: 0,
8873 clean: true,
8874 verify_retried: false,
8875 e2e_deferred: false,
8876 e2e_defer_reason: None,
8877 progressed: false,
8878 vote_split: false,
8879 reconsideration: Vec::new(),
8880 verdict: None,
8881 }];
8882 state.gate = vec![CommandOutcome {
8883 command: "test".to_owned(),
8884 code: Some(0),
8885 output_tail: String::new(),
8886 duration_ms: 0,
8887 resource_blocked: false,
8888 }];
8889 state.gate_ran = true;
8890 state.status = RunStatus::Landing;
8894 state.merge = Some(MergeOutcome {
8895 mode: MergeMode::Pr,
8896 ok: true,
8897 detail: "https://example.invalid/x/y/pull/1".to_owned(),
8898 });
8899
8900 ask_test_home();
8904 let store = ask::Questions::open();
8905 let q = ask_open_question(&store, &state.id);
8906
8907 let mut runner = Runner {
8908 state,
8909 roles: ResolvedRoles {
8910 implementers: Vec::new(),
8911 judges: Vec::new(),
8912 reviewers: Vec::new(),
8913 fixer: None,
8914 conductor: conductor(),
8915 implementer_roster: Vec::new(),
8916 },
8917 sem: Arc::new(Semaphore::new(1)),
8918 pause: Pause::new(),
8919 interrupt: Pause::new(),
8920 };
8921
8922 runner.execute().await.expect("execute");
8927
8928 assert_eq!(
8929 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8930 Some("https://example.invalid/x/y/pull/1"),
8931 "reentry must not push again or open a second pull request over the \
8932 one `land` is already watching"
8933 );
8934 assert_ne!(
8935 runner.state.status,
8936 RunStatus::Landing,
8937 "land could not actually reach the fake pull request, so it must \
8938 have given up rather than left the run silently parked forever"
8939 );
8940 assert_eq!(runner.state.status, RunStatus::Blocked);
8944 assert!(
8945 store.get(&q.id).unwrap().status.open(),
8946 "Blocked is still alive; settle_questions must have been a no-op here"
8947 );
8948 }
8949
8950 fn state_with_round(round: ReviewRound) -> RunState {
8951 let mut s = RunState::new(
8952 PathBuf::from("/repo"),
8953 "main".to_owned(),
8954 "abc1234".to_owned(),
8955 "add retries".to_owned(),
8956 Config::default(),
8957 );
8958 s.reviews = vec![round];
8959 s
8960 }
8961
8962 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
8963 crate::verdict::Finding {
8964 id: id.to_owned(),
8965 severity,
8966 file: None,
8967 line: None,
8968 title: title.to_owned(),
8969 detail: String::new(),
8970 }
8971 }
8972
8973 #[test]
8974 fn pr_body_names_open_findings_and_declined_ones() {
8975 let round = ReviewRound {
8976 round: 2,
8977 head: "deadbee".to_owned(),
8978 verified_head: None,
8979 verified_at: None,
8980 reviews: vec![ReviewRecord {
8981 attempts: 0,
8982 reviewer: 1,
8983 agent: "alpha".to_owned(),
8984 summary: String::new(),
8985 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
8986 vote: None,
8987 failed: None,
8988 duration_ms: 0,
8989 }],
8990 e2e: vec![CommandOutcome {
8991 command: "cargo test".to_owned(),
8992 code: Some(0),
8993 output_tail: String::new(),
8994 duration_ms: 0,
8995 resource_blocked: false,
8996 }],
8997 verify_retried: false,
8998 e2e_deferred: false,
8999 e2e_defer_reason: None,
9000 fix: Some(FixRecord {
9001 agent: "alpha".to_owned(),
9002 addressed: Vec::new(),
9003 rejected: vec![crate::verdict::Rejection {
9004 id: "R1-1-1".to_owned(),
9005 why: "not reachable from any caller".to_owned(),
9006 }],
9007 notes: String::new(),
9008 committed: true,
9009 failed: None,
9010 duration_ms: 0,
9011 continuation: None,
9012 }),
9013 blocking: 0,
9014 answered: 1,
9015 expected: 1,
9016 clean: false,
9017 progressed: true,
9018 vote_split: false,
9019 reconsideration: Vec::new(),
9020 verdict: None,
9021 };
9022 let state = state_with_round(round);
9023 let body = pr_message(&state, 'A').body;
9024
9025 assert!(body.contains("add retries"), "the task must still be there");
9026 assert!(body.contains("R2-1-1"), "{body}");
9027 assert!(body.contains("unused import"), "{body}");
9028 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
9029 assert!(
9030 body.contains("not reachable from any caller"),
9031 "the reason it was declined: {body}"
9032 );
9033 }
9034
9035 #[test]
9036 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
9037 let round = ReviewRound {
9038 round: 1,
9039 head: "deadbee".to_owned(),
9040 verified_head: None,
9041 verified_at: None,
9042 reviews: vec![ReviewRecord {
9043 attempts: 0,
9044 reviewer: 1,
9045 agent: "alpha".to_owned(),
9046 summary: String::new(),
9047 findings: Vec::new(),
9048 vote: None,
9049 failed: None,
9050 duration_ms: 0,
9051 }],
9052 e2e: Vec::new(),
9053 verify_retried: false,
9054 e2e_deferred: false,
9055 e2e_defer_reason: None,
9056 fix: None,
9057 blocking: 0,
9058 answered: 1,
9059 expected: 1,
9060 clean: true,
9061 progressed: false,
9062 vote_split: false,
9063 reconsideration: Vec::new(),
9064 verdict: None,
9065 };
9066 let state = state_with_round(round);
9067 let body = pr_message(&state, 'A').body;
9068 assert!(!body.contains("Open review findings"), "{body}");
9069 assert!(!body.contains("Declined"), "{body}");
9070 }
9071
9072 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
9073 let mut state = RunState::new(
9074 PathBuf::from("/repo"),
9075 "main".to_owned(),
9076 "abc1234".to_owned(),
9077 instruction.to_owned(),
9078 Config::default(),
9079 );
9080 state.candidates.push(Candidate {
9081 index: 0,
9082 label: 'A',
9083 agent: "alpha".to_owned(),
9084 branch: "magi/x/A".to_owned(),
9085 worktree: PathBuf::from("/wt"),
9086 summary: summary.to_owned(),
9087 stat: String::new(),
9088 files: 1,
9089 commits: 1,
9090 empty: false,
9091 failed: None,
9092 verified_noop: None,
9093 folded: false,
9094 duration_ms: 0,
9095 });
9096 state
9097 }
9098
9099 #[test]
9100 fn pr_message_describes_the_change_not_the_task() {
9101 let state = state_with_summary(
9102 "今回やってほしいこと: results projector を直す",
9103 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
9104 );
9105 let m = pr_message(&state, 'A');
9106 assert_eq!(m.title, "fix(web): batch the runs list reads");
9107 assert!(
9108 m.body.starts_with("## Summary\n\n- reads run.json once"),
9109 "{}",
9110 m.body
9111 );
9112 assert!(!m.body.contains("TITLE:"), "{}", m.body);
9113 let task_at = m.body.find("今回やってほしいこと").unwrap();
9114 let details_at = m.body.find("<details>").unwrap();
9115 assert!(
9116 details_at < task_at,
9117 "the task lives inside <details>: {}",
9118 m.body
9119 );
9120 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
9121 assert!(m.body.contains("magi:candidate-a"));
9122 }
9123
9124 #[test]
9125 fn pr_message_falls_back_to_the_task_without_a_title_line() {
9126 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
9127 let m = pr_message(&state, 'A');
9128 assert_eq!(m.title, "add retries");
9129 assert!(
9130 m.body.contains("## Summary\n\n- did some things"),
9131 "{}",
9132 m.body
9133 );
9134
9135 let none = RunState::new(
9136 PathBuf::from("/repo"),
9137 "main".to_owned(),
9138 "abc1234".to_owned(),
9139 "add retries".to_owned(),
9140 Config::default(),
9141 );
9142 let m = pr_message(&none, 'A');
9143 assert_eq!(m.title, "add retries");
9144 assert!(!m.body.contains("## Summary"), "{}", m.body);
9145 }
9146
9147 #[test]
9148 fn pr_message_refuses_the_candidate_commit_subject() {
9149 for bad in [
9150 "TITLE: magi: candidate A (uncommitted work)",
9151 "TITLE: chore: stuff (uncommitted work)",
9152 "TITLE: ",
9153 ] {
9154 let state = state_with_summary("add retries", bad);
9155 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
9156 }
9157 }
9158
9159 #[test]
9160 fn pr_message_bounds_a_very_long_task_and_title() {
9161 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
9162 let state = state_with_summary(&long, "- nothing");
9163 let m = pr_message(&state, 'A');
9164 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9165 assert!(!m.title.contains('\n'));
9166
9167 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
9168 let m = pr_message(&state, 'A');
9169 assert!(m.title.starts_with("feat: "));
9170 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9171 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
9172 }
9173
9174 #[test]
9175 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
9176 let mut state = state_with_summary(
9180 "add retries",
9181 "TITLE: fix(web): batch reads\n- reads run.json once",
9182 );
9183 state.config.graph.language = "ja".to_owned();
9184 let m = pr_message(&state, 'A');
9185 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
9186
9187 let task = "今回やってほしいこと: results projector を直す";
9190 let mut state = state_with_summary(task, "- no title line");
9191 state.config.graph.language = "ja".to_owned();
9192 let m = pr_message(&state, 'A');
9193 assert_eq!(
9194 m.title,
9195 format!("chore: land candidate A of run {}", state.id)
9196 );
9197 assert!(
9198 m.body.contains(&format!(
9199 "<summary>Original task</summary>\n\n{task}\n\n</details>"
9200 )),
9201 "{}",
9202 m.body
9203 );
9204 }
9205
9206 #[test]
9207 fn pr_message_scrubs_home_paths_and_addresses() {
9208 let state = state_with_summary(
9209 "fix it in /Users/someone/src/x",
9210 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9211 );
9212 let m = pr_message(&state, 'A');
9213 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9214 assert!(!m.body.contains(leak), "{}", m.body);
9215 }
9216 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9217 }
9218
9219 #[test]
9220 fn pr_message_survives_a_task_that_closes_details() {
9221 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9222 let m = pr_message(&state, 'A');
9223 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9224 }
9225
9226 #[test]
9227 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9228 let cmd = manual_merge_command(
9229 MergeStyle::Squash,
9230 Path::new("/repo"),
9231 "b",
9232 "fix: \"quoted\" $(x) `y`\n\nbody",
9233 );
9234 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9235 }
9236
9237 #[test]
9238 fn manual_merge_command_matches_the_configured_style() {
9239 let repo = Path::new("/repo");
9240 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9241
9242 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9243 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9244
9245 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9246 assert_eq!(
9247 squash,
9248 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9249 \"Merge magi run 0832 (candidate A)\""
9250 );
9251
9252 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9253 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9254 }
9255
9256 #[test]
9257 fn a_nudge_gets_a_quarter_of_the_budget() {
9258 assert_eq!(retry_budget(secs(1200), true), secs(300));
9260 assert_eq!(retry_budget(secs(3600), true), secs(900));
9261 }
9262
9263 #[test]
9264 fn a_resent_prompt_keeps_the_whole_budget() {
9265 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9268 assert_eq!(retry_budget(secs(60), false), secs(60));
9269 }
9270
9271 #[test]
9272 fn the_floor_never_exceeds_the_original_budget() {
9273 assert_eq!(retry_budget(secs(60), true), secs(60));
9277 assert_eq!(retry_budget(secs(480), true), secs(120));
9278 assert_eq!(retry_budget(secs(0), true), secs(0));
9279 }
9280
9281 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9282 agent::CommandEvidence {
9283 id: "item1".to_owned(),
9284 description: "cargo test".to_owned(),
9285 exit_code,
9286 result_summary: String::new(),
9287 source: "codex".to_owned(),
9288 }
9289 }
9290
9291 #[test]
9292 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9293 assert!(!has_unconfirmed_command(&[]));
9297 }
9298
9299 #[test]
9300 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9301 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9305 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9306 assert!(!has_unconfirmed_command(&[
9307 evidence(Some(0)),
9308 evidence(Some(101))
9309 ]));
9310 }
9311
9312 #[test]
9313 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9314 assert!(has_unconfirmed_command(&[
9315 evidence(Some(0)),
9316 evidence(None)
9317 ]));
9318 }
9319
9320 #[test]
9321 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9322 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9323 assert_eq!(
9324 verified_noop_claim(true, &[], text).as_deref(),
9325 Some("already fixed by b32cfc4, on main.")
9326 );
9327 }
9328
9329 #[test]
9330 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9331 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9334 assert!(verified_noop_claim(false, &[], text).is_none());
9335 }
9336
9337 #[test]
9338 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9339 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9340 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9341 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9343 }
9344
9345 #[test]
9346 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9347 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9348 }
9349
9350 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9353 runner.state.candidates = shape
9354 .iter()
9355 .enumerate()
9356 .map(|(i, &(empty, verified))| Candidate {
9357 index: i,
9358 label: (b'A' + i as u8) as char,
9359 agent: "sonnet".to_owned(),
9360 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9361 worktree: PathBuf::from(format!("/wt/{i}")),
9362 summary: String::new(),
9363 stat: String::new(),
9364 files: 0,
9365 commits: 0,
9366 empty,
9367 failed: None,
9368 verified_noop: verified.map(str::to_owned),
9369 duration_ms: 0,
9370 folded: false,
9371 })
9372 .collect();
9373 }
9374
9375 #[test]
9376 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9377 ask_test_home();
9378 let mut runner = runner_at(RunStatus::Implementing);
9379 set_candidates(
9380 &mut runner,
9381 &[
9382 (true, Some("already on main at b32cfc4")),
9383 (true, Some("same fix, see the existing test")),
9384 ],
9385 );
9386
9387 runner
9388 .after_implement()
9389 .expect("a verified no-op is not an error");
9390
9391 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9392 }
9393
9394 #[test]
9395 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9396 ask_test_home();
9397 let mut runner = runner_at(RunStatus::Implementing);
9398 set_candidates(
9402 &mut runner,
9403 &[(true, Some("already on main at b32cfc4")), (true, None)],
9404 );
9405
9406 let err = runner
9407 .after_implement()
9408 .expect_err("an unverified empty candidate must still fail the run");
9409
9410 assert!(
9411 err.to_string().contains("no candidate produced a change"),
9412 "{err}"
9413 );
9414 assert_eq!(runner.state.status, RunStatus::Failed);
9415 }
9416
9417 #[test]
9418 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9419 ask_test_home();
9420 let mut runner = runner_at(RunStatus::Implementing);
9421 set_candidates(&mut runner, &[(true, None), (true, None)]);
9422
9423 let err = runner
9424 .after_implement()
9425 .expect_err("no candidate declared anything; this is an ordinary failure");
9426
9427 assert!(
9428 err.to_string().contains("no candidate produced a change"),
9429 "{err}"
9430 );
9431 assert_eq!(runner.state.status, RunStatus::Failed);
9432 }
9433}