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 self.mark_driver_exited();
737 let ended = if result.is_err() {
738 Some(crate::notices::run_stopped(&self.state.id, &self.state))
739 } else {
740 crate::notices::run_ended(&self.state)
741 };
742 if let Some(notice) = ended {
743 crate::notices::raise(notice);
744 }
745 result
746 }
747
748 fn mark_driver_exited(&mut self) {
757 self.state.driver_exited = true;
758 let pid = std::process::id();
759 let Ok(mut disk) = RunState::load(&self.state.id) else {
760 return;
761 };
762 if disk.released_to.is_some() || disk.driver_pid != Some(pid) || disk.driver_exited {
763 return;
764 }
765 disk.driver_exited = true;
766 if let Err(e) = disk.save() {
767 tracing::warn!("could not record that run {} stopped: {e:#}", self.state.id);
768 }
769 }
770
771 async fn execute_graph(&mut self) -> Result<()> {
772 self.state.parked = false;
777 self.state.clear_active();
784 if let Ok(disk) = RunState::load(&self.state.id)
802 && let Some(to) = &disk.released_to
803 {
804 bail!(
805 "run {} cannot continue: its worktree was released to run {}",
806 self.state.short(),
807 crate::run::short_of(to)
808 );
809 }
810 let pid = std::process::id();
811 self.state.driver_pid = Some(pid);
812 self.state.driver_started_at = crate::proc::process_started_at(pid);
813 self.state.driver_exited = false;
814 self.state.save()?;
815 if self.state.status == RunStatus::Stalled {
828 if self.recover_stall().await? {
829 self.finish_after_tally().await?;
830 } else {
831 self.state.save()?;
833 }
834 return Ok(());
835 }
836 if self.state.status == RunStatus::Landing {
846 self.run_land().await?;
847 self.settle_questions();
852 return Ok(());
853 }
854 self.prep().await?;
855 if self.park_here()? {
856 return Ok(());
857 }
858 self.advise().await?;
859 if self.park_here()? {
860 return Ok(());
861 }
862 self.implement().await?;
863 if self.park_here()? {
864 return Ok(());
865 }
866 if self.state.status == RunStatus::VerifiedNoop {
870 return Ok(());
871 }
872 self.judge().await?;
873 if self.park_here()? {
874 return Ok(());
875 }
876 self.deliberate().await?;
877 if self.park_here()? {
878 return Ok(());
879 }
880 self.vote().await?;
881 if self.park_here()? {
882 return Ok(());
883 }
884 self.tally()?;
885 if self.state.status == RunStatus::Stalled {
890 self.state.save()?;
894 return Ok(());
895 }
896 self.finish_after_tally().await?;
897 Ok(())
898 }
899
900 fn park_here(&mut self) -> Result<bool> {
907 if !self.pause.parked() && !self.interrupt.parked() {
912 return Ok(false);
913 }
914 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
915 Some(reason) => format!(
916 "parked after `{}` ({reason}) — resume to carry on from here",
917 self.state.status.as_str()
918 ),
919 None => format!(
920 "parked after `{}` — resume to carry on from here",
921 self.state.status.as_str()
922 ),
923 };
924 self.state.event("park", why);
925 self.state.parked = true;
926 self.state.save()?;
927 Ok(true)
928 }
929
930 pub fn on_pause(&mut self, pause: Pause) {
932 self.pause = pause;
933 }
934
935 pub fn watch_interrupt(&mut self, pause: Pause) {
941 self.interrupt = pause;
942 }
943
944 fn settle_questions(&mut self) {
964 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
965 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
966 }
967 }
968
969 async fn finish_after_tally(&mut self) -> Result<()> {
972 self.fold_losers().await?;
973 self.sync_to_base().await?;
978 self.review_loop().await?;
979 self.sync_to_base().await?;
980 self.gate().await?;
981 self.merge().await?;
982 self.state.save()?;
983 Ok(())
984 }
985
986 async fn prep(&mut self) -> Result<()> {
989 if !self.state.candidates.is_empty() {
990 return Ok(());
991 }
992 self.state.status = RunStatus::Prep;
993 let repo = self.state.repo.clone();
994 let base = self.state.base_commit.clone();
995 let root = self.state.worktree_root();
996 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
997
998 let hooks_dir = self.state.dir().join("hooks");
1001 if self.state.config.blind.commit_msg_hook {
1002 std::fs::create_dir_all(&hooks_dir)
1003 .with_context(|| format!("create {}", hooks_dir.display()))?;
1004 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
1005 let path = hooks_dir.join("commit-msg");
1006 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
1007 make_executable(&path)?;
1008 git::acquire_worktree_config(&repo).await?;
1016 self.state.enabled_worktree_config = true;
1017 }
1018
1019 for (index, (spec, label)) in self
1020 .roles
1021 .implementers
1022 .clone()
1023 .into_iter()
1024 .zip(labels)
1025 .enumerate()
1026 {
1027 let branch = self.state.branch_for(label);
1028 let worktree = root.join(format!("cand-{label}"));
1029 git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
1030 if self.state.config.blind.commit_msg_hook {
1031 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
1032 }
1033 git::local_exclude(&worktree, "/.magi/").await?;
1034 self.state.candidates.push(Candidate {
1035 index,
1036 label,
1037 agent: spec.id.clone(),
1038 branch,
1039 worktree,
1040 summary: String::new(),
1041 stat: String::new(),
1042 files: 0,
1043 commits: 0,
1044 empty: false,
1045 failed: None,
1046 verified_noop: None,
1047 duration_ms: 0,
1048 folded: false,
1049 });
1050 }
1051
1052 for j in 1..=self.roles.judges.len() {
1053 let wt = root.join(format!("judge-{j}"));
1054 if !wt.exists() {
1055 git::worktree_add_detached(&repo, &wt, &base).await?;
1056 }
1057 }
1058
1059 if self.state.config.graph.advise {
1067 for k in 1..=self.state.config.graph.advisors {
1068 let wt = root.join(format!("advisor-{k}"));
1069 if !wt.exists() {
1070 git::worktree_add_detached(&repo, &wt, &base).await?;
1071 }
1072 }
1073 }
1074
1075 let authors: Vec<&str> = self
1080 .roles
1081 .implementers
1082 .iter()
1083 .map(|a| a.id.as_str())
1084 .collect();
1085 let overlap: Vec<String> = self
1086 .roles
1087 .judges
1088 .iter()
1089 .enumerate()
1090 .filter(|(_, j)| authors.contains(&j.id.as_str()))
1091 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
1092 .collect();
1093 if !overlap.is_empty() {
1094 let note = format!(
1095 "{} also authored a candidate; blind, but the panel is less \
1096 independent than {} distinct agents would be",
1097 overlap.join(", "),
1098 self.roles.judges.len()
1099 );
1100 self.state.event("prep", note);
1101 }
1102
1103 self.state.event(
1104 "prep",
1105 format!(
1106 "{} candidates, {} judges, base {} ({})",
1107 self.state.candidates.len(),
1108 self.roles.judges.len(),
1109 &self.state.base_commit[..7.min(self.state.base_commit.len())],
1110 self.state.base_branch
1111 ),
1112 );
1113 self.state.status = RunStatus::Implementing;
1114 self.state.save()?;
1115 Ok(())
1116 }
1117
1118 async fn advise(&mut self) -> Result<()> {
1151 let implement_untouched = self
1152 .state
1153 .candidates
1154 .iter()
1155 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
1156 if !self.state.config.graph.advise || self.state.advise_attempted {
1157 return Ok(());
1158 }
1159 if !implement_untouched {
1160 self.state.event(
1161 "advise",
1162 "skipping the design-deliberation stage: at least one \
1163 candidate already shows implementation progress, so this \
1164 run is past the point the stage exists to run before"
1165 .to_owned(),
1166 );
1167 self.state.advise_attempted = true;
1168 self.state.save()?;
1169 return Ok(());
1170 }
1171 let run_id = self.state.id.clone();
1172 let prompts = self.state.config.prompts.clone();
1173 let instruction = self.state.instruction.clone();
1174 let language = self.state.config.graph.language.clone();
1175 let root = self.state.worktree_root();
1176 let n = self.state.config.graph.advisors;
1177 let where_recorded = self.state.dir().join("run.json");
1178
1179 let seats = match self.state.config.advisors() {
1180 Ok(seats) if !seats.is_empty() => seats,
1181 Ok(_) => {
1182 self.state.event(
1183 "advise",
1184 format!(
1185 "[graph] advisors is 0; skipping the design-deliberation \
1186 stage and continuing without a synthesis brief (see {})",
1187 where_recorded.display()
1188 ),
1189 );
1190 self.state.advise_attempted = true;
1191 self.state.save()?;
1192 return Ok(());
1193 }
1194 Err(e) => {
1195 self.state.event(
1196 "advise",
1197 format!(
1198 "could not resolve advisor seats ({e:#}); continuing \
1199 without a design-deliberation brief (see {})",
1200 where_recorded.display()
1201 ),
1202 );
1203 self.state.advise_attempted = true;
1204 self.state.save()?;
1205 return Ok(());
1206 }
1207 };
1208
1209 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1210 let artifacts = agent::artifacts_dir(&self.state.dir());
1211 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1212
1213 let mut jobs = Vec::new();
1214 for (i, spec) in seats.iter().cloned().enumerate() {
1215 let seat_key = format!("advisor-{}", i + 1);
1216 let seat = self.seat(&seat_key, &spec.id);
1217 jobs.push(SeatJob {
1218 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1219 spec,
1220 seat,
1221 cwd: worktrees[i % worktrees.len()].clone(),
1222 timeout,
1223 allow_write: false,
1224 sessions: false,
1225 artifacts: artifacts.clone(),
1226 stem: seat_key,
1227 });
1228 }
1229
1230 self.state.event(
1231 "advise",
1232 format!(
1233 "{} advisor seat(s) sketching a design in parallel",
1234 jobs.len()
1235 ),
1236 );
1237 let mut quota_losses = Vec::new();
1238 let cache = self.state.config.cache_dir();
1239 let ctx = WaveCtx {
1240 run: &run_id,
1241 node: "advise",
1242 prompts: &prompts,
1243 cache: cache.as_deref(),
1244 round: None,
1245 };
1246 let results = ask_json_wave::<Proposal>(
1247 jobs,
1248 Arc::clone(&self.sem),
1249 self.state.config.graph.retries,
1250 &ctx,
1251 &mut quota_losses,
1252 &mut self.state,
1253 &|p: &Proposal| p.validate(),
1254 )
1255 .await;
1256 self.state.quota.extend(quota_losses);
1257
1258 let mut records = Vec::with_capacity(results.len());
1259 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1260 let agent_id = seat.agent.clone();
1261 self.state.seats.insert(seat.key.clone(), seat);
1262 match res {
1263 Ok((proposal, out)) => {
1264 self.state
1265 .event("advise", format!("advisor-{} proposed a design", i + 1));
1266 records.push(advise::AdvisorRecord::proposed(
1267 i + 1,
1268 agent_id,
1269 proposal,
1270 out.duration_ms,
1271 ));
1272 }
1273 Err(e) => {
1274 self.state.event(
1275 "advise",
1276 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1277 );
1278 records.push(advise::AdvisorRecord::failed(
1279 i + 1,
1280 agent_id,
1281 e.to_string(),
1282 ));
1283 }
1284 }
1285 }
1286
1287 let mut advice = advise::Advice {
1288 records,
1289 synthesis: None,
1290 };
1291 if advice.proposals().is_empty() {
1292 self.state.event(
1293 "advise",
1294 "no advisor produced a usable proposal; continuing without a \
1295 synthesis brief"
1296 .to_owned(),
1297 );
1298 } else {
1299 match self
1300 .synthesize_brief(
1301 &advice,
1302 &instruction,
1303 &language,
1304 &worktrees[0],
1305 &artifacts,
1306 &run_id,
1307 &prompts,
1308 cache.as_deref(),
1309 )
1310 .await
1311 {
1312 Ok(Some(text)) => {
1313 self.state.event(
1314 "advise",
1315 "synthesized a design brief for the implementer".to_owned(),
1316 );
1317 advice.synthesis = Some(text);
1318 }
1319 Ok(None) => {
1320 self.state.event(
1321 "advise",
1322 "the synthesis seat produced nothing usable; continuing \
1323 without a design brief"
1324 .to_owned(),
1325 );
1326 }
1327 Err(e) => {
1328 self.state.event(
1329 "advise",
1330 format!("could not synthesize a design brief: {e:#}"),
1331 );
1332 }
1333 }
1334 }
1335 advise::apply_reflection(&mut advice);
1336
1337 self.state.advice = Some(advice);
1338 self.state.advise_attempted = true;
1339 self.state.save()?;
1340 Ok(())
1341 }
1342
1343 #[allow(clippy::too_many_arguments)]
1355 async fn synthesize_brief(
1356 &mut self,
1357 advice: &advise::Advice,
1358 instruction: &str,
1359 language: &str,
1360 cwd: &Path,
1361 artifacts: &Path,
1362 run_id: &str,
1363 prompts: &Prompts,
1364 cache: Option<&Path>,
1365 ) -> Result<Option<String>> {
1366 let want = self.state.config.roles.synthesizer.as_deref();
1367 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1368 let mut seat = self.seat("advise-synthesis", &spec.id);
1369 let proposals = advice.proposals();
1370 let mut prompt = prompt::with_overlay(
1371 prompt::synthesize_brief(instruction, &proposals, language),
1372 prompts.overlay("advise"),
1373 );
1374 if cache.is_some() {
1375 prompt.push('\n');
1380 prompt.push_str(&prompt::build_cache_note("advise", false));
1381 }
1382 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1383 let out = agent::invoke(
1384 &spec,
1385 &mut seat,
1386 &Invocation {
1387 cwd,
1388 prompt: &prompt,
1389 timeout,
1390 allow_write: false,
1391 sessions: false,
1392 artifacts,
1393 stem: "advise-synthesis",
1394 run: run_id,
1395 node: "advise",
1396 cache_dir: None,
1397 attachments: &[],
1398 },
1399 )
1400 .await?;
1401 self.state.seats.insert(seat.key.clone(), seat);
1402 if !out.usable() {
1403 return Ok(None);
1404 }
1405 let text =
1406 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1407 Ok((!text.trim().is_empty()).then_some(text))
1408 }
1409
1410 async fn implement(&mut self) -> Result<()> {
1413 let run_id = self.state.id.clone();
1418 let prompts = self.state.config.prompts.clone();
1419 let todo: Vec<usize> = self
1420 .state
1421 .candidates
1422 .iter()
1423 .enumerate()
1424 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1425 .map(|(i, _)| i)
1426 .collect();
1427 if todo.is_empty() {
1428 return self.after_implement();
1429 }
1430 self.state.status = RunStatus::Implementing;
1431
1432 let language = self.state.config.graph.language.clone();
1433 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1434 let sessions = self.state.config.graph.sessions;
1435 let artifacts = agent::artifacts_dir(&self.state.dir());
1436 let brief = self
1440 .state
1441 .advice
1442 .as_ref()
1443 .and_then(|a| a.synthesis.as_deref())
1444 .map(str::to_owned);
1445
1446 let mut jobs = Vec::new();
1447 for &i in &todo {
1448 let (index, label, worktree) = {
1449 let c = &self.state.candidates[i];
1450 (c.index, c.label, c.worktree.clone())
1451 };
1452 let spec = self.roles.implementers[index].clone();
1453 let seat_key = format!("impl-{label}");
1454 let seat = self.seat(&seat_key, &spec.id);
1455 let instruction = self.state.instruction.clone();
1456 jobs.push(SeatJob {
1457 spec,
1458 seat,
1459 prompt: prompt::implement(
1460 &instruction,
1461 &worktree.to_string_lossy(),
1462 &language,
1463 brief.as_deref(),
1464 ),
1465 cwd: worktree,
1466 timeout,
1467 allow_write: true,
1468 sessions,
1469 artifacts: artifacts.clone(),
1470 stem: format!("impl-{label}"),
1471 });
1472 }
1473
1474 self.state.event(
1475 "implement",
1476 format!("{} candidates in parallel", jobs.len()),
1477 );
1478 let mut sent = jobs.clone();
1484 let cache = self.state.config.cache_dir();
1485 let ctx = WaveCtx {
1486 run: &run_id,
1487 node: "implement",
1488 prompts: &prompts,
1489 cache: cache.as_deref(),
1490 round: None,
1491 };
1492 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1493 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1494 .await;
1495 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1496 .await;
1497 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1498 .await;
1499
1500 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1501 let seat_key = seat.key.clone();
1502 let agent = seat.agent.clone();
1511 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1512 self.state.seats.insert(seat.key.clone(), seat);
1513 let label = self.state.candidates[i].label;
1514 let worktree = self.state.candidates[i].worktree.clone();
1515 let base = self.state.base_commit.clone();
1516
1517 let (summary, duration, failed, verified_claim) = match out {
1518 AgentOutcome::Ok(o) => {
1519 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1520 let failed = (!o.usable()).then(|| {
1521 if o.timed_out {
1522 "agent timed out".to_owned()
1523 } else {
1524 format!("agent exited with {:?}", o.exit_code)
1525 }
1526 });
1527 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1528 (text, o.duration_ms, failed, verified_claim)
1529 }
1530 AgentOutcome::Dropped(o) => {
1536 let why = o
1537 .dropped
1538 .as_ref()
1539 .map(|d| d.why.as_str())
1540 .unwrap_or("the CLI ended the stream without delivering its answer");
1541 (
1542 String::new(),
1543 o.duration_ms,
1544 Some(format!("the CLI dropped the stream ({why})")),
1545 None,
1546 )
1547 }
1548 AgentOutcome::Quota(o) => {
1549 self.state.quota.push(QuotaLoss {
1550 seat: seat_key,
1551 node: "implement".to_owned(),
1552 at: Timestamp::now(),
1553 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1554 });
1555 (
1556 String::new(),
1557 o.duration_ms,
1558 Some("rate limited (quota); produced no change".to_owned()),
1559 None,
1560 )
1561 }
1562 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1563 };
1564
1565 let rescued = match git::rescue_commit(
1568 &worktree,
1569 &format!("magi: candidate {label} (uncommitted work)"),
1570 )
1571 .await
1572 {
1573 Ok(r) => {
1574 self.state.note_withheld("implement", &r.withheld);
1575 r.committed
1576 }
1577 Err(_) => false,
1578 };
1579 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1580 .await
1581 .unwrap_or(0);
1582 let patch = git::diff(&worktree, &base, "HEAD")
1583 .await
1584 .unwrap_or_default();
1585 let stat = git::diff_stat(&worktree, &base, "HEAD")
1586 .await
1587 .unwrap_or_default();
1588 let files = git::changed_files(&worktree, &base, "HEAD")
1589 .await
1590 .map(|f| f.len())
1591 .unwrap_or(0);
1592 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1593
1594 let c = &mut self.state.candidates[i];
1595 if !exhausted_the_fallback_chain {
1596 c.agent = agent;
1597 }
1598 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1599 c.stat = stat;
1600 c.files = files;
1601 c.commits = commits;
1602 c.duration_ms = duration;
1603 c.empty = commits == 0 || patch.trim().is_empty();
1604 c.failed = match failed {
1607 Some(_) if c.empty => failed,
1608 _ => None,
1609 };
1610 c.verified_noop = if c.empty { verified_claim } else { None };
1615 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1616 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1617 (None, true, Some(_), _) => {
1618 format!("candidate {label}: no change produced (agent-verified no-op)")
1619 }
1620 (None, true, None, _) => format!("candidate {label}: no change produced"),
1621 (None, false, _, true) => {
1622 format!(
1623 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1624 )
1625 }
1626 (None, false, _, false) => {
1627 format!("candidate {label}: {files} files, {commits} commits")
1628 }
1629 };
1630 self.state.event("implement", note);
1631 self.state.save()?;
1632 }
1633
1634 self.after_implement()
1635 }
1636
1637 async fn resume_undelivered(
1665 &mut self,
1666 results: &mut [(usize, SeatState, AgentOutcome)],
1667 sent: &[SeatJob],
1668 prompts: &Prompts,
1669 run_id: &str,
1670 ) {
1671 for (wi, seat, out) in results.iter_mut() {
1672 let Some(dropped) = (match &*out {
1673 AgentOutcome::Dropped(o) => o.dropped.clone(),
1674 _ => None,
1675 }) else {
1676 continue;
1677 };
1678 let Some(job) = sent.get(*wi) else { continue };
1679 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1681 self.state.event(
1682 "implement",
1683 format!(
1684 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1685 work is in the tree",
1686 seat.key, dropped.output_tokens, dropped.why
1687 ),
1688 );
1689 continue;
1690 }
1691 if !has_context(&job.spec, seat, job.sessions) {
1699 self.state.event(
1700 "implement",
1701 format!(
1702 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1703 is no session left to resume",
1704 seat.key, dropped.output_tokens, dropped.why
1705 ),
1706 );
1707 continue;
1708 }
1709 self.state.event(
1710 "implement",
1711 format!(
1712 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1713 conversation",
1714 seat.key, dropped.output_tokens, dropped.why
1715 ),
1716 );
1717 let mut retry = job.clone();
1718 retry.seat = seat.clone();
1719 retry.prompt = prompt::resume_after_drop(&dropped.why);
1720 retry.timeout = retry_budget(job.timeout, true);
1721 retry.stem = format!("{}-resume", job.stem);
1722 let cache = self.state.config.cache_dir();
1723 let ctx = WaveCtx {
1724 run: run_id,
1725 node: "implement",
1726 prompts,
1727 cache: cache.as_deref(),
1728 round: None,
1729 };
1730 let (resumed_seat, resumed) =
1731 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1732 *seat = resumed_seat;
1733 *out = resumed;
1734 }
1735 }
1736
1737 async fn resume_quota_losses(
1799 &mut self,
1800 results: &mut [(usize, SeatState, AgentOutcome)],
1801 sent: &mut [SeatJob],
1802 prompts: &Prompts,
1803 run_id: &str,
1804 ) {
1805 let instruction = self.state.instruction.clone();
1806 let language = self.state.config.graph.language.clone();
1807 let brief = self
1808 .state
1809 .advice
1810 .as_ref()
1811 .and_then(|a| a.synthesis.as_deref())
1812 .map(str::to_owned);
1813 for (wi, seat, out) in results.iter_mut() {
1814 let Some(job) = sent.get_mut(*wi) else {
1815 continue;
1816 };
1817 let start = self
1822 .roles
1823 .implementer_roster
1824 .iter()
1825 .position(|s| s.id == job.spec.id)
1826 .unwrap_or(0);
1827 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1828 let mut fallback_attempt = 0usize;
1829 while matches!(&*out, AgentOutcome::Quota(_)) {
1830 let Some(next) =
1831 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1832 .cloned()
1833 else {
1834 break;
1835 };
1836 tried.insert(next.id.clone());
1837 fallback_attempt += 1;
1838
1839 if let Ok(r) = git::rescue_commit(
1840 &job.cwd,
1841 &format!(
1842 "magi: candidate {} (uncommitted work before quota fallback)",
1843 seat.key
1844 ),
1845 )
1846 .await
1847 {
1848 self.state.note_withheld("implement", &r.withheld);
1849 }
1850
1851 self.state.event(
1852 "implement",
1853 format!(
1854 "{}: rate limited (quota) on {}; retrying with {}",
1855 seat.key, seat.agent, next.id
1856 ),
1857 );
1858
1859 let new_seat = self.seat(&seat.key, &next.id);
1860 job.spec = next.clone();
1868 let mut retry = job.clone();
1869 retry.seat = new_seat;
1870 retry.prompt = prompt::implement(
1871 &instruction,
1872 &job.cwd.to_string_lossy(),
1873 &language,
1874 brief.as_deref(),
1875 );
1876 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1877 let cache = self.state.config.cache_dir();
1878 let ctx = WaveCtx {
1879 run: run_id,
1880 node: "implement",
1881 prompts,
1882 cache: cache.as_deref(),
1883 round: None,
1884 };
1885 let (fallback_seat, fallback_out) = run_one(
1886 retry,
1887 Arc::clone(&self.sem),
1888 &ctx,
1889 &mut self.state,
1890 fallback_attempt,
1891 )
1892 .await;
1893 *seat = fallback_seat;
1894 *out = fallback_out;
1895 }
1896 }
1897 }
1898
1899 async fn resume_unconfirmed_commands(
1923 &mut self,
1924 results: &mut [(usize, SeatState, AgentOutcome)],
1925 sent: &[SeatJob],
1926 prompts: &Prompts,
1927 run_id: &str,
1928 ) {
1929 for (wi, seat, out) in results.iter_mut() {
1930 let AgentOutcome::Ok(o) = &*out else {
1931 continue;
1932 };
1933 if !has_unconfirmed_command(&o.commands) {
1934 continue;
1935 }
1936 let Some(job) = sent.get(*wi) else { continue };
1937 if !has_context(&job.spec, seat, job.sessions) {
1938 self.state.event(
1939 "implement",
1940 format!(
1941 "{}: the reply named a command whose own CLI never confirmed the exit \
1942 status of, but there is no session left to resume",
1943 seat.key
1944 ),
1945 );
1946 continue;
1947 }
1948 self.state.event(
1949 "implement",
1950 format!(
1951 "{}: the reply named a command whose own CLI never confirmed the exit \
1952 status of; resuming the conversation",
1953 seat.key
1954 ),
1955 );
1956 let mut retry = job.clone();
1957 retry.seat = seat.clone();
1958 retry.prompt = prompt::resume_incomplete(
1959 "a command in your last reply had no confirmed exit status",
1960 );
1961 retry.timeout = retry_budget(job.timeout, true);
1962 retry.stem = format!("{}-confirm", job.stem);
1963 let cache = self.state.config.cache_dir();
1964 let ctx = WaveCtx {
1965 run: run_id,
1966 node: "implement",
1967 prompts,
1968 cache: cache.as_deref(),
1969 round: None,
1970 };
1971 let (resumed_seat, resumed) =
1972 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1973 *seat = resumed_seat;
1974 *out = resumed;
1975 }
1976 }
1977
1978 async fn continue_fix_report(
1999 &mut self,
2000 mut seat: SeatState,
2001 parse_err: String,
2002 job: &SeatJob,
2003 prompts: &Prompts,
2004 run_id: &str,
2005 round: usize,
2006 ) -> (
2007 SeatState,
2008 Option<FixReport>,
2009 Option<String>,
2010 ContinuationRecord,
2011 ) {
2012 let mut last_err = parse_err;
2013 let mut cumulative_wait_ms = 0u64;
2014 let mut attempts = 0usize;
2015 loop {
2016 if !has_context(&job.spec, &seat, job.sessions) {
2017 self.state.event(
2018 "fix",
2019 format!(
2020 "round {round}: fixer's reply had no adoption report ({last_err}); no \
2021 session left to resume into"
2022 ),
2023 );
2024 let outcome = if attempts == 0 {
2025 ContinuationOutcome::NoSession
2026 } else {
2027 ContinuationOutcome::Exhausted
2028 };
2029 return (
2030 seat,
2031 None,
2032 Some(format!("unparsable fix report: {last_err}")),
2033 ContinuationRecord {
2034 attempts,
2035 cumulative_wait_ms,
2036 outcome,
2037 },
2038 );
2039 }
2040 if attempts >= MAX_FIX_CONTINUATIONS {
2041 self.state.event(
2042 "fix",
2043 format!(
2044 "round {round}: fixer's reply still had no adoption report after \
2045 {attempts} continuation(s) ({last_err}); giving up"
2046 ),
2047 );
2048 return (
2049 seat,
2050 None,
2051 Some(format!(
2052 "unparsable fix report after {attempts} continuation(s): {last_err}"
2053 )),
2054 ContinuationRecord {
2055 attempts,
2056 cumulative_wait_ms,
2057 outcome: ContinuationOutcome::Exhausted,
2058 },
2059 );
2060 }
2061 attempts += 1;
2062 self.state.event(
2063 "fix",
2064 format!(
2065 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
2066 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
2067 ),
2068 );
2069 let mut retry = job.clone();
2070 retry.seat = seat.clone();
2071 retry.prompt = prompt::resume_incomplete(&last_err);
2072 retry.timeout = retry_budget(job.timeout, true);
2073 retry.stem = format!("{}-continue{attempts}", job.stem);
2074 let cache = self.state.config.cache_dir();
2075 let ctx = WaveCtx {
2076 run: run_id,
2077 node: "fix",
2078 prompts,
2079 cache: cache.as_deref(),
2080 round: Some(round),
2081 };
2082 let (resumed_seat, resumed_out) = run_one(
2083 retry,
2084 Arc::clone(&self.sem),
2085 &ctx,
2086 &mut self.state,
2087 attempts,
2088 )
2089 .await;
2090 seat = resumed_seat;
2091 match resumed_out {
2092 AgentOutcome::Ok(o) => {
2093 cumulative_wait_ms += o.duration_ms;
2094 match verdict::extract_json::<FixReport>(&o.text) {
2095 Ok(report) if !has_unconfirmed_command(&o.commands) => {
2096 self.state.event(
2097 "fix",
2098 format!(
2099 "round {round}: fixer's adoption report recovered after \
2100 {attempts} continuation(s)"
2101 ),
2102 );
2103 return (
2104 seat,
2105 Some(report),
2106 None,
2107 ContinuationRecord {
2108 attempts,
2109 cumulative_wait_ms,
2110 outcome: ContinuationOutcome::Resumed,
2111 },
2112 );
2113 }
2114 Ok(_) => {
2122 last_err = "the reply parsed, but it reported a command whose own CLI \
2123 never confirmed an exit status"
2124 .to_owned();
2125 }
2126 Err(e) => last_err = e.to_string(),
2127 }
2128 }
2129 AgentOutcome::Quota(o) => {
2130 cumulative_wait_ms += o.duration_ms;
2131 self.state.quota.push(QuotaLoss {
2132 seat: seat.key.clone(),
2133 node: "fix".to_owned(),
2134 at: Timestamp::now(),
2135 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2136 });
2137 self.state.event(
2138 "fix",
2139 format!(
2140 "round {round}: continuation rate limited (quota); not retrying now"
2141 ),
2142 );
2143 return (
2144 seat,
2145 None,
2146 Some("rate limited (quota) while recovering the fix report".to_owned()),
2147 ContinuationRecord {
2148 attempts,
2149 cumulative_wait_ms,
2150 outcome: ContinuationOutcome::QuotaLost,
2151 },
2152 );
2153 }
2154 AgentOutcome::Dropped(o) => {
2155 cumulative_wait_ms += o.duration_ms;
2156 let why = o
2157 .dropped
2158 .as_ref()
2159 .map(|d| d.why.as_str())
2160 .unwrap_or("the CLI ended the stream without delivering its answer");
2161 last_err = format!("the CLI dropped the stream ({why})");
2162 }
2163 AgentOutcome::Failed(e) => last_err = e,
2164 }
2165 }
2166 }
2167
2168 fn after_implement(&mut self) -> Result<()> {
2169 if self.state.leaks.is_empty() {
2171 let cfg = self.state.config.blind.clone();
2172 let mut leaks = Vec::new();
2173 for c in &self.state.candidates {
2174 let Some(patch) =
2175 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
2176 else {
2177 continue;
2178 };
2179 leaks.extend(blind::scan(
2180 &format!("candidate {} patch", c.label),
2181 &patch,
2182 &cfg.vendor_tokens,
2183 ));
2184 }
2185 if !leaks.is_empty() {
2186 let summary = leaks
2187 .iter()
2188 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2189 .collect::<Vec<_>>()
2190 .join(", ");
2191 match cfg.on_leak {
2192 LeakPolicy::Fail => {
2193 self.state.status = RunStatus::Failed;
2194 self.state
2195 .event("blind", format!("vendor text in a patch: {summary}"));
2196 self.state.leaks = leaks;
2197 self.state.save()?;
2198 self.settle_questions();
2199 bail!(
2200 "blind.on_leak = \"fail\" and vendor text reached a \
2201 judged patch: {summary}"
2202 );
2203 }
2204 LeakPolicy::Redact => self.state.event(
2205 "blind",
2206 format!("redacting vendor text for judging: {summary}"),
2207 ),
2208 LeakPolicy::Warn => self.state.event(
2209 "blind",
2210 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2211 ),
2212 }
2213 self.state.leaks = leaks;
2214 }
2215 }
2216
2217 if self.state.viable().is_empty() {
2218 if self.state.all_candidates_verified_noop() {
2219 self.state.status = RunStatus::VerifiedNoop;
2230 self.state.save()?;
2231 self.settle_questions();
2232 return Ok(());
2233 }
2234 self.state.status = RunStatus::Failed;
2235 self.state.save()?;
2236 self.settle_questions();
2237 bail!("no candidate produced a change; nothing to judge");
2238 }
2239 self.state.status = RunStatus::Judging;
2240 self.state.save()?;
2241 Ok(())
2242 }
2243
2244 async fn judge(&mut self) -> Result<()> {
2247 let run_id = self.state.id.clone();
2252 let prompts = self.state.config.prompts.clone();
2253 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2254 return Ok(());
2255 }
2256 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2257 if viable.len() == 1 {
2258 self.state.judge_skipped = true;
2265 self.state.event(
2266 "judge",
2267 format!(
2268 "only candidate {} produced a change; judging skipped",
2269 viable[0].label
2270 ),
2271 );
2272 self.state.save()?;
2273 return Ok(());
2274 }
2275 self.state.status = RunStatus::Judging;
2276
2277 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2278 let language = self.state.config.graph.language.clone();
2279 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2280 let sessions = self.state.config.graph.sessions;
2281 let artifacts = agent::artifacts_dir(&self.state.dir());
2282 let root = self.state.worktree_root();
2283 let base_short = short(&self.state.base_commit);
2284
2285 let mut jobs = Vec::new();
2286 let mut orders = Vec::new();
2287 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2288 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2289 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2290 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2291 let seat_key = format!("judge-{}", j + 1);
2292 let seat = self.seat(&seat_key, &spec.id);
2293 jobs.push(SeatJob {
2294 prompt: prompt::judge(
2295 &self.state.instruction,
2296 &views,
2297 self.roles.judges.len(),
2298 &base_short,
2299 &language,
2300 ),
2301 spec,
2302 seat,
2303 cwd: root.join(format!("judge-{}", j + 1)),
2304 timeout,
2305 allow_write: false,
2306 sessions,
2307 artifacts: artifacts.clone(),
2308 stem: format!("judge-{}", j + 1),
2309 });
2310 }
2311
2312 self.state.event(
2313 "judge",
2314 format!(
2315 "{} judges ranking {} candidates blind",
2316 jobs.len(),
2317 viable.len()
2318 ),
2319 );
2320 let labels_for_check = labels.clone();
2321 let mut quota_losses = Vec::new();
2322 let cache = self.state.config.cache_dir();
2323 let ctx = WaveCtx {
2324 run: &run_id,
2325 node: "judge",
2326 prompts: &prompts,
2327 cache: cache.as_deref(),
2328 round: None,
2329 };
2330 let results = ask_json_wave::<Ranking>(
2331 jobs,
2332 Arc::clone(&self.sem),
2333 self.state.config.graph.retries,
2334 &ctx,
2335 &mut quota_losses,
2336 &mut self.state,
2337 &move |r: &Ranking| r.validate(&labels_for_check),
2338 )
2339 .await;
2340 self.state.quota.extend(quota_losses);
2341
2342 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2343 let agent_id = seat.agent.clone();
2344 self.state.seats.insert(seat.key.clone(), seat);
2345 let mut record = Judgement {
2346 judge: j + 1,
2347 seat: format!("judge-{}", j + 1),
2348 agent: agent_id,
2349 ranking: Vec::new(),
2350 reasons: BTreeMap::new(),
2351 confidence: None,
2352 order: orders[j].clone(),
2353 failed: None,
2354 duration_ms: 0,
2355 };
2356 match res {
2357 Ok((ranking, out)) => {
2358 record.ranking = ranking.normalized();
2359 record.reasons = ranking.reasons;
2360 record.confidence = ranking.confidence;
2361 record.duration_ms = out.duration_ms;
2362 self.state.event(
2363 "judge",
2364 format!(
2365 "judge {} ranked {}",
2366 j + 1,
2367 record.ranking.iter().collect::<String>()
2368 ),
2369 );
2370 }
2371 Err(e) => {
2372 record.failed = Some(e.to_string());
2373 self.state
2374 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2375 }
2376 }
2377 self.state.judgements.push(record);
2378 self.state.save()?;
2379 }
2380 Ok(())
2381 }
2382
2383 async fn deliberate(&mut self) -> Result<()> {
2386 let run_id = self.state.id.clone();
2391 let prompts = self.state.config.prompts.clone();
2392 if !self.state.deliberation.is_empty() {
2393 return Ok(());
2394 }
2395 let tops: Vec<char> = self
2396 .state
2397 .judgements
2398 .iter()
2399 .filter_map(|j| j.ranking.first().copied())
2400 .collect();
2401 let rounds = self.state.config.graph.deliberate_rounds;
2402 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2403 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2404 self.state.event(
2405 "deliberate",
2406 format!("judges agreed on {} outright; no deliberation", tops[0]),
2407 );
2408 }
2409 self.state.status = RunStatus::Voting;
2410 self.state.save()?;
2411 return Ok(());
2412 }
2413
2414 self.state.status = RunStatus::Deliberating;
2415 self.state.event(
2416 "deliberate",
2417 format!(
2418 "split: first choices were {} — opening {rounds} round(s)",
2419 tops.iter().collect::<String>()
2420 ),
2421 );
2422
2423 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2424 let language = self.state.config.graph.language.clone();
2425 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2426 let sessions = self.state.config.graph.sessions;
2427 let artifacts = agent::artifacts_dir(&self.state.dir());
2428 let root = self.state.worktree_root();
2429 let base_short = short(&self.state.base_commit);
2430
2431 for round in 1..=rounds {
2435 let mut turns: Vec<DeliberationTurn> = Vec::new();
2436 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2437 if self.state.judgements[j].failed.is_some() {
2438 continue;
2439 }
2440 let seat_key = format!("judge-{}", j + 1);
2441 let mut seat = self.seat(&seat_key, &spec.id);
2442 let transcript = self.transcript(&turns, j);
2443 let context = if has_context(&spec, &seat, sessions) {
2444 None
2445 } else {
2446 Some(self.candidate_block(&viable, &base_short))
2447 };
2448 let text = prompt::deliberate(
2449 &self.state.instruction,
2450 context.as_deref(),
2451 &transcript,
2452 round,
2453 rounds,
2454 &language,
2455 );
2456 let job = SeatJob {
2457 spec,
2458 seat: seat.clone(),
2459 prompt: text,
2460 cwd: root.join(format!("judge-{}", j + 1)),
2461 timeout,
2462 allow_write: false,
2463 sessions,
2464 artifacts: artifacts.clone(),
2465 stem: format!("delib-{round}-judge-{}", j + 1),
2466 };
2467 let cache = self.state.config.cache_dir();
2468 let ctx = WaveCtx {
2469 run: &run_id,
2470 node: "deliberate",
2471 prompts: &prompts,
2472 cache: cache.as_deref(),
2473 round: None,
2474 };
2475 let (updated, out) =
2476 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2477 seat = updated;
2478 let agent_id = seat.agent.clone();
2479 let seat_key = seat.key.clone();
2480 self.state.seats.insert(seat.key.clone(), seat);
2481 let body = match out {
2482 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2483 AgentOutcome::Dropped(o) => {
2487 let why =
2488 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2489 "the CLI ended the stream without delivering its answer",
2490 );
2491 self.state.event(
2492 "deliberate",
2493 format!(
2494 "judge {} skipped: the CLI dropped the stream ({why})",
2495 j + 1
2496 ),
2497 );
2498 continue;
2499 }
2500 AgentOutcome::Quota(o) => {
2501 self.state.quota.push(QuotaLoss {
2502 seat: seat_key,
2503 node: "deliberate".to_owned(),
2504 at: Timestamp::now(),
2505 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2506 });
2507 self.state.event(
2508 "deliberate",
2509 format!("judge {} skipped: rate limited (quota)", j + 1),
2510 );
2511 continue;
2512 }
2513 AgentOutcome::Failed(e) => {
2514 self.state
2515 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2516 continue;
2517 }
2518 };
2519 let tentative = verdict::extract_json::<Position>(&body)
2520 .ok()
2521 .and_then(|p| p.tentative)
2522 .and_then(|s| s.trim().chars().next())
2523 .map(|c| c.to_ascii_uppercase());
2524 self.state.event(
2525 "deliberate",
2526 format!(
2527 "round {round}: judge {} now favours {}",
2528 j + 1,
2529 tentative.map_or("—".to_owned(), |c| c.to_string())
2530 ),
2531 );
2532 turns.push(DeliberationTurn {
2533 judge: j + 1,
2534 agent: agent_id,
2535 body: blind::sanitize_prose(&body, &self.state.config.blind),
2536 tentative,
2537 });
2538 }
2539 self.state
2540 .deliberation
2541 .push(DeliberationRound { round, turns });
2542 self.state.save()?;
2543 }
2544
2545 self.state.status = RunStatus::Voting;
2546 self.state.save()?;
2547 Ok(())
2548 }
2549
2550 async fn vote(&mut self) -> Result<()> {
2553 let run_id = self.state.id.clone();
2558 let prompts = self.state.config.prompts.clone();
2559 if !self.state.votes.is_empty() {
2560 return Ok(());
2561 }
2562 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2563 if viable.len() == 1 {
2564 return Ok(());
2565 }
2566 self.state.status = RunStatus::Voting;
2567
2568 let language = self.state.config.graph.language.clone();
2569 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2570 let sessions = self.state.config.graph.sessions;
2571 let artifacts = agent::artifacts_dir(&self.state.dir());
2572 let root = self.state.worktree_root();
2573 let base_short = short(&self.state.base_commit);
2574 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2575
2576 let mut jobs = Vec::new();
2577 let mut seats_at = Vec::new();
2578 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2579 if self
2580 .state
2581 .judgements
2582 .get(j)
2583 .is_some_and(|r| r.failed.is_some())
2584 {
2585 continue;
2586 }
2587 let seat_key = format!("judge-{}", j + 1);
2588 let seat = self.seat(&seat_key, &spec.id);
2589 let mut text = prompt::final_vote(&viable, &language);
2590 if !has_context(&spec, &seat, sessions) {
2591 text = format!(
2592 "{}\n\n# Candidates\n\n{}",
2593 text,
2594 self.candidate_block(&candidates, &base_short)
2595 );
2596 }
2597 jobs.push(SeatJob {
2598 spec,
2599 seat,
2600 prompt: text,
2601 cwd: root.join(format!("judge-{}", j + 1)),
2602 timeout,
2603 allow_write: false,
2604 sessions,
2605 artifacts: artifacts.clone(),
2606 stem: format!("vote-judge-{}", j + 1),
2607 });
2608 seats_at.push(j);
2609 }
2610
2611 self.state.event(
2612 "vote",
2613 format!(
2614 "collecting {} final votes one by one, privately",
2615 jobs.len()
2616 ),
2617 );
2618 let allowed = viable.clone();
2619 let mut quota_losses = Vec::new();
2620 let cache = self.state.config.cache_dir();
2621 let ctx = WaveCtx {
2622 run: &run_id,
2623 node: "vote",
2624 prompts: &prompts,
2625 cache: cache.as_deref(),
2626 round: None,
2627 };
2628 let results = ask_json_wave::<FinalVote>(
2629 jobs,
2630 Arc::clone(&self.sem),
2631 self.state.config.graph.retries,
2632 &ctx,
2633 &mut quota_losses,
2634 &mut self.state,
2635 &move |v: &FinalVote| match v.label() {
2636 Some(c) if allowed.contains(&c) => Ok(()),
2637 other => bail!("vote {other:?} is not one of {allowed:?}"),
2638 },
2639 )
2640 .await;
2641 self.state.quota.extend(quota_losses);
2642
2643 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2644 let agent_id = seat.agent.clone();
2645 self.state.seats.insert(seat.key.clone(), seat);
2646 let initial = self
2647 .state
2648 .judgements
2649 .get(j)
2650 .and_then(|r| r.ranking.first().copied());
2651 let mut record = VoteRecord {
2652 judge: j + 1,
2653 agent: agent_id,
2654 vote: None,
2655 reason: String::new(),
2656 changed: false,
2657 };
2658 match res {
2659 Ok((v, _)) => {
2660 record.vote = v.label();
2661 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2662 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2663 self.state.event(
2664 "vote",
2665 format!(
2666 "judge {} voted {}{}",
2667 j + 1,
2668 record.vote.unwrap_or('?'),
2669 if record.changed { " (changed)" } else { "" }
2670 ),
2671 );
2672 }
2673 Err(e) => {
2674 self.state
2675 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2676 }
2677 }
2678 self.state.votes.push(record);
2679 self.state.save()?;
2680 }
2681 Ok(())
2682 }
2683
2684 fn tally(&mut self) -> Result<()> {
2687 if self.state.tally.is_some() {
2688 return Ok(());
2689 }
2690 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2691 let tops: Vec<char> = self
2692 .state
2693 .judgements
2694 .iter()
2695 .filter_map(|j| j.ranking.first().copied())
2696 .collect();
2697 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2698
2699 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2702 let mut cast: Vec<char> = Vec::new();
2703 for (i, j) in self.state.judgements.iter().enumerate() {
2704 let vote = self
2705 .state
2706 .votes
2707 .iter()
2708 .find(|v| v.judge == i + 1)
2709 .and_then(|v| v.vote)
2710 .or_else(|| j.ranking.first().copied());
2711 if let Some(v) = vote {
2712 *first_choice.entry(v).or_insert(0) += 1;
2713 cast.push(v);
2714 }
2715 }
2716
2717 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2718 for j in &self.state.judgements {
2719 let n = j.ranking.len();
2720 for (pos, label) in j.ranking.iter().enumerate() {
2721 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2722 }
2723 }
2724
2725 let best = first_choice.values().copied().max().unwrap_or(0);
2726 let mut leaders: Vec<char> = first_choice
2727 .iter()
2728 .filter(|(_, v)| **v == best)
2729 .map(|(k, _)| *k)
2730 .collect();
2731 let mut tie_break = None;
2732 if leaders.len() > 1 {
2733 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2734 let borda_leaders: Vec<char> = leaders
2735 .iter()
2736 .copied()
2737 .filter(|l| borda[l] == top_borda)
2738 .collect();
2739 tie_break = Some(if borda_leaders.len() == 1 {
2740 format!(
2741 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2742 leaders.len()
2743 )
2744 } else {
2745 format!(
2746 "{} way tie on both first-choice votes and Borda points, broken by label order",
2747 leaders.len()
2748 )
2749 });
2750 leaders = borda_leaders;
2751 leaders.sort_unstable();
2752 }
2753 let winner = *leaders
2754 .first()
2755 .or(viable.first())
2756 .context("no candidate to declare a winner from")?;
2757
2758 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2759 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2760 let deliberated = !self.state.deliberation.is_empty();
2761
2762 let quota_seats: std::collections::BTreeSet<&str> =
2766 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2767 let mut present = 0usize;
2768 for (i, j) in self.state.judgements.iter().enumerate() {
2769 if quota_seats.contains(j.seat.as_str()) {
2770 continue;
2771 }
2772 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2773 let voted = self
2774 .state
2775 .votes
2776 .iter()
2777 .any(|v| v.judge == i + 1 && v.vote.is_some());
2778 if ranked || voted {
2779 present += 1;
2780 }
2781 }
2782 let needs_quorum = viable.len() > 1;
2788 let judges_total = if needs_quorum {
2789 self.roles.judges.len()
2790 } else {
2791 0
2792 };
2793 let quorum = if needs_quorum {
2794 judges_total / 2 + 1
2795 } else {
2796 0
2797 };
2798 let met_quorum = !needs_quorum || present >= quorum;
2799 let uncontested = (!needs_quorum).then(|| {
2800 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2801 });
2802
2803 self.state.event(
2804 "tally",
2805 match &uncontested {
2806 Some(reason) => format!("winner {winner} — {reason}"),
2807 None => format!(
2808 "winner {winner} — votes {} | initial {} | {} changed | \
2809 {present}/{judges_total} judges{}",
2810 first_choice
2811 .iter()
2812 .map(|(k, v)| format!("{k}:{v}"))
2813 .collect::<Vec<_>>()
2814 .join(" "),
2815 if unanimous_initial {
2816 "unanimous"
2817 } else {
2818 "split"
2819 },
2820 changed_votes,
2821 if met_quorum {
2822 String::new()
2823 } else {
2824 format!(" — below quorum ({quorum} required)")
2825 },
2826 ),
2827 },
2828 );
2829 if !met_quorum {
2830 self.state.event(
2831 "stall",
2832 format!(
2833 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2834 the run stops here, resumable"
2835 ),
2836 );
2837 }
2838 self.state.tally = Some(Tally {
2839 first_choice,
2840 borda,
2841 winner,
2842 rankings: tops.len(),
2843 unanimous_initial,
2844 deliberated,
2845 changed_votes,
2846 unanimous_final,
2847 tie_break,
2848 judges: judges_total,
2849 present,
2850 quorum,
2851 met_quorum,
2852 uncontested,
2853 });
2854 self.state.status = if met_quorum {
2855 RunStatus::Reviewing
2856 } else {
2857 RunStatus::Stalled
2858 };
2859 self.state.save()?;
2860 Ok(())
2861 }
2862
2863 #[allow(clippy::too_many_lines)]
2884 async fn recover_stall(&mut self) -> Result<bool> {
2885 let run_id = self.state.id.clone();
2890 let prompts = self.state.config.prompts.clone();
2891 let quota_seats: BTreeSet<&str> =
2896 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2897 let absent: Vec<String> = self
2898 .state
2899 .judgements
2900 .iter()
2901 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2902 .map(|j| j.seat.clone())
2903 .collect();
2904 if absent.is_empty() {
2905 return Ok(false);
2906 }
2907 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2908 if viable.len() <= 1 {
2909 return Ok(false);
2910 }
2911 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2912 let language = self.state.config.graph.language.clone();
2913 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2914 let sessions = self.state.config.graph.sessions;
2915 let artifacts = agent::artifacts_dir(&self.state.dir());
2916 let root = self.state.worktree_root();
2917 let base_short = short(&self.state.base_commit);
2918 let candidates: Vec<Candidate> = viable.clone();
2919
2920 let mut positions: Vec<usize> = absent
2922 .iter()
2923 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2924 .collect();
2925 if positions.is_empty() {
2926 return Ok(false);
2927 }
2928 positions.sort_unstable();
2929 positions.dedup();
2930
2931 let mut judge_jobs = Vec::new();
2933 for &j in &positions {
2934 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2935 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2936 let seat_key = format!("judge-{}", j + 1);
2937 let spec = self.roles.judges[j].clone();
2938 let seat = self.seat(&seat_key, &spec.id);
2939 judge_jobs.push(SeatJob {
2940 spec,
2941 seat,
2942 prompt: prompt::judge(
2943 &self.state.instruction,
2944 &views,
2945 self.roles.judges.len(),
2946 &base_short,
2947 &language,
2948 ),
2949 cwd: root.join(seat_key),
2950 timeout,
2951 allow_write: false,
2952 sessions,
2953 artifacts: artifacts.clone(),
2954 stem: format!("judge-{}-recover", j + 1),
2955 });
2956 }
2957
2958 let labels_for_check = labels.clone();
2959 let mut judge_losses = Vec::new();
2960 let retries = self.state.config.graph.retries;
2961 let cache = self.state.config.cache_dir();
2962 let ctx = WaveCtx {
2963 run: &run_id,
2964 node: "judge",
2965 prompts: &prompts,
2966 cache: cache.as_deref(),
2967 round: None,
2968 };
2969 let results = ask_json_wave::<Ranking>(
2970 judge_jobs,
2971 Arc::clone(&self.sem),
2972 retries,
2973 &ctx,
2974 &mut judge_losses,
2975 &mut self.state,
2976 &move |r: &Ranking| r.validate(&labels_for_check),
2977 )
2978 .await;
2979
2980 let mut recovered: BTreeSet<usize> = BTreeSet::new();
2982 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
2983 self.state.seats.insert(seat.key.clone(), seat);
2984 let record = &mut self.state.judgements[j];
2985 match res {
2986 Ok((ranking, out)) => {
2987 record.ranking = ranking.normalized();
2988 record.reasons = ranking.reasons;
2989 record.confidence = ranking.confidence;
2990 record.failed = None;
2991 record.duration_ms = out.duration_ms;
2992 recovered.insert(j);
2993 self.state.event(
2994 "recover",
2995 format!("judge {} ranked again after the limit", j + 1),
2996 );
2997 }
2998 Err(e) => {
2999 self.state
3000 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
3001 }
3002 }
3003 }
3004
3005 let mut vote_jobs = Vec::new();
3007 let mut vote_pos: Vec<usize> = Vec::new();
3008 for &j in &recovered {
3009 let seat_key = format!("judge-{}", j + 1);
3010 let spec = self.roles.judges[j].clone();
3011 let seat = self.seat(&seat_key, &spec.id);
3012 let mut text = prompt::final_vote(&labels, &language);
3013 if !has_context(&spec, &seat, sessions) {
3014 text = format!(
3015 "{}\n\n# Candidates\n\n{}",
3016 text,
3017 self.candidate_block(&candidates, &base_short)
3018 );
3019 }
3020 vote_jobs.push(SeatJob {
3021 spec,
3022 seat,
3023 prompt: text,
3024 cwd: root.join(seat_key),
3025 timeout,
3026 allow_write: false,
3027 sessions,
3028 artifacts: artifacts.clone(),
3029 stem: format!("vote-judge-{}-recover", j + 1),
3030 });
3031 vote_pos.push(j);
3032 }
3033 let allowed = labels.clone();
3034 let mut vote_losses = Vec::new();
3035 let vote_retries = self.state.config.graph.retries;
3036 let vote_cache = self.state.config.cache_dir();
3037 let ctx = WaveCtx {
3038 run: &run_id,
3039 node: "vote",
3040 prompts: &prompts,
3041 cache: vote_cache.as_deref(),
3042 round: None,
3043 };
3044 let votes = ask_json_wave::<FinalVote>(
3045 vote_jobs,
3046 Arc::clone(&self.sem),
3047 vote_retries,
3048 &ctx,
3049 &mut vote_losses,
3050 &mut self.state,
3051 &move |v: &FinalVote| match v.label() {
3052 Some(c) if allowed.contains(&c) => Ok(()),
3053 other => bail!("vote {other:?} is not one of {allowed:?}"),
3054 },
3055 )
3056 .await;
3057 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
3058 let agent_id = seat.agent.clone();
3059 self.state.seats.insert(seat.key.clone(), seat);
3060 match res {
3061 Ok((v, _)) => {
3062 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
3063 rec.vote = v.label();
3064 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
3065 } else {
3066 self.state.votes.push(VoteRecord {
3067 judge: j + 1,
3068 agent: agent_id,
3069 vote: v.label(),
3070 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
3071 changed: false,
3072 });
3073 }
3074 self.state.event(
3075 "recover",
3076 format!("judge {} voted again after the limit", j + 1),
3077 );
3078 }
3079 Err(e) => {
3080 self.state
3081 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
3082 }
3083 }
3084 }
3085
3086 let recovered_keys: BTreeSet<String> = recovered
3090 .iter()
3091 .map(|&j| format!("judge-{}", j + 1))
3092 .collect();
3093 self.state
3094 .quota
3095 .retain(|q| !recovered_keys.contains(&q.seat));
3096 for loss in judge_losses.into_iter().chain(vote_losses) {
3100 if recovered_keys.contains(&loss.seat) {
3101 continue;
3102 }
3103 self.state.quota.retain(|q| q.seat != loss.seat);
3104 self.state.quota.push(loss);
3105 }
3106
3107 self.state.tally = None;
3109 self.tally()?;
3110 Ok(self
3111 .state
3112 .tally
3113 .as_ref()
3114 .map(|t| t.met_quorum)
3115 .unwrap_or(false))
3116 }
3117
3118 async fn fold_losers(&mut self) -> Result<()> {
3121 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
3122 return Ok(());
3123 };
3124 let repo = self.state.repo.clone();
3125 let mut folded = Vec::new();
3126 for i in 0..self.state.candidates.len() {
3127 let c = &self.state.candidates[i];
3128 if c.label == winner || c.folded {
3129 continue;
3130 }
3131 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
3132 git::worktree_remove(&repo, &wt).await.ok();
3133 git::branch_delete(&repo, &branch).await.ok();
3134 self.state.candidates[i].folded = true;
3135 folded.push(label.to_string());
3136 }
3137 let root = self.state.worktree_root();
3139 for j in 1..=self.roles.judges.len() {
3140 let wt = root.join(format!("judge-{j}"));
3141 if wt.exists() {
3142 git::worktree_remove(&repo, &wt).await.ok();
3143 }
3144 }
3145 if self.state.config.graph.advise {
3148 for k in 1..=self.state.config.graph.advisors {
3149 let wt = root.join(format!("advisor-{k}"));
3150 if wt.exists() {
3151 git::worktree_remove(&repo, &wt).await.ok();
3152 }
3153 }
3154 }
3155 if !folded.is_empty() {
3156 self.state
3157 .event("fold", format!("folded candidates {}", folded.join(", ")));
3158 self.state.save()?;
3159 }
3160 Ok(())
3161 }
3162
3163 async fn sync_to_base(&mut self) -> Result<()> {
3193 if self
3194 .state
3195 .base_sync
3196 .as_ref()
3197 .is_some_and(|s| s.conflict.is_some())
3198 {
3199 return Ok(());
3200 }
3201 let Some(winner) = self.state.winner().cloned() else {
3202 return Ok(());
3203 };
3204
3205 let repo = self.state.repo.clone();
3206 let remote = self.state.config.merge.remote.clone();
3207 let base_branch = self.state.base_branch.clone();
3208 let tracking = format!("{remote}/{base_branch}");
3209
3210 git::fetch(&repo, &remote, &base_branch).await.ok();
3211 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3215 return Ok(());
3216 };
3217
3218 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3219 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3220 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3221
3222 if behind == 0 {
3223 self.state.base_sync = Some(BaseSync {
3224 tip,
3225 behind: 0,
3226 attempts,
3227 conflict: None,
3228 });
3229 self.state.save()?;
3230 return Ok(());
3231 }
3232
3233 if attempts >= BASE_SYNC_ROUNDS {
3234 let why = format!(
3235 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3236 rebase(s); rebasing again would only race it",
3237 winner.branch
3238 );
3239 self.state.status = RunStatus::Blocked;
3240 self.state.base_sync = Some(BaseSync {
3241 tip,
3242 behind,
3243 attempts,
3244 conflict: Some(why.clone()),
3245 });
3246 self.state.event("land", why);
3247 self.state.save()?;
3248 return Ok(());
3249 }
3250
3251 self.state.event(
3252 "land",
3253 format!(
3254 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3255 winner.branch
3256 ),
3257 );
3258 self.state.save()?;
3259
3260 let scratch = self.state.dir().join("base-sync");
3261 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
3262 let attempts = attempts + 1;
3263 match rebased {
3264 Ok(None) => {
3265 git::sync_to_head(&winner.worktree).await?;
3269 self.state.base_sync = Some(BaseSync {
3270 tip: tip.clone(),
3271 behind: 0,
3272 attempts,
3273 conflict: None,
3274 });
3275 self.state
3276 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3277 }
3278 Ok(Some(conflict)) => {
3279 let why = format!(
3280 "{} conflicts with {tracking} and did not rebase: {}",
3281 winner.branch,
3282 conflict.chars().take(600).collect::<String>()
3283 );
3284 self.state.status = RunStatus::Blocked;
3285 self.state.base_sync = Some(BaseSync {
3286 tip,
3287 behind,
3288 attempts,
3289 conflict: Some(why.clone()),
3290 });
3291 self.state.event("land", why);
3292 }
3293 Err(e) => {
3294 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3295 self.state.status = RunStatus::Blocked;
3296 self.state.base_sync = Some(BaseSync {
3297 tip,
3298 behind,
3299 attempts,
3300 conflict: Some(why.clone()),
3301 });
3302 self.state.event("land", why);
3303 }
3304 }
3305 self.state.save()?;
3306 Ok(())
3307 }
3308
3309 fn landing_base(&self) -> String {
3319 self.state
3320 .base_sync
3321 .as_ref()
3322 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3323 }
3324
3325 pub async fn fix_selected(
3358 &mut self,
3359 ids: &[String],
3360 reason: &str,
3361 allow_stale: bool,
3362 ) -> Result<()> {
3363 let reason = reason.trim();
3364 if reason.is_empty() {
3365 bail!("a fix request needs a reason — that is the operator's own record of why");
3366 }
3367 if ids.is_empty() {
3368 bail!("no finding id given");
3369 }
3370 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3371 bail!(
3372 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3373 has already concluded — can be given a targeted fix. A run still \
3374 in progress should simply be resumed; a `merged` run's branch has \
3375 already landed, so its answer is a fresh `magi review <branch>`, \
3376 not reopening this run's own record",
3377 self.state.id,
3378 self.state.status.as_str()
3379 );
3380 }
3381 let Some(winner) = self.state.winner().cloned() else {
3382 bail!("run {} has no winning candidate to fix", self.state.id);
3383 };
3384 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3385 bail!(
3386 "branch `{}` no longer exists; this run cannot be extended",
3387 winner.branch
3388 );
3389 }
3390 let home = crate::run::home();
3391 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3392 bail!(
3393 "run {} is currently being worked on by another magi process",
3394 self.state.id
3395 );
3396 }
3397 let _claim = FixClaim::acquire(&self.state.dir())?;
3403
3404 let mut seen = BTreeSet::new();
3408 let mut findings = Vec::new();
3409 let mut missing = Vec::new();
3410 for id in ids {
3411 if !seen.insert(id.clone()) {
3412 continue;
3413 }
3414 match self.state.finding(id) {
3415 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3416 id: f.id.clone(),
3417 severity: f.severity,
3418 reviewer_vote: rec.vote,
3419 round: round.round,
3420 round_head: round.head.clone(),
3421 reviewer: rec.reviewer,
3422 agent: rec.agent.clone(),
3423 file: f.file.clone(),
3424 line: f.line,
3425 title: f.title.clone(),
3426 detail: f.detail.clone(),
3427 outcome: OperatorFixOutcome::Pending,
3428 }),
3429 None => missing.push(id.clone()),
3430 }
3431 }
3432 if !missing.is_empty() {
3433 bail!(
3434 "unknown finding id(s): {}; nothing was changed",
3435 missing.join(", ")
3436 );
3437 }
3438
3439 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3440 let stale_details: Vec<(String, String)> = findings
3441 .iter()
3442 .filter(|f| f.round_head != head_at_request)
3443 .map(|f| (f.id.clone(), f.round_head.clone()))
3444 .collect();
3445 let stale = !stale_details.is_empty();
3446 if stale && !allow_stale {
3447 bail!(
3448 "the branch has moved since some finding(s) were raised — {} — now \
3449 at {}; pass --allow-stale to fix anyway, or re-run review first",
3450 stale_details
3451 .iter()
3452 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3453 .collect::<Vec<_>>()
3454 .join(", "),
3455 short(&head_at_request)
3456 );
3457 }
3458
3459 let request = OperatorFixRequest {
3460 requested_at: Timestamp::now(),
3461 reason: reason.to_owned(),
3462 findings,
3463 head_at_request: head_at_request.clone(),
3464 allow_stale,
3465 stale,
3466 fix: None,
3467 result_head: None,
3468 follow_up_review_run: None,
3469 };
3470 self.state.event(
3471 "fix",
3472 format!(
3473 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3474 request.findings.len(),
3475 request
3476 .findings
3477 .iter()
3478 .map(|f| f.id.as_str())
3479 .collect::<Vec<_>>()
3480 .join(", "),
3481 ),
3482 );
3483 self.state.operator_fixes.push(request);
3490 self.state.save()?;
3491 let request_index = self.state.operator_fixes.len() - 1;
3492
3493 if winner.worktree.exists() {
3502 let dirty = git::git(
3505 &winner.worktree,
3506 &["status", "--porcelain", "--untracked-files=all"],
3507 )
3508 .await?;
3509 let only_withheld = dirty.lines().all(|l| {
3510 l.strip_prefix("?? ")
3511 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3512 });
3513 if !only_withheld {
3514 bail!(
3515 "`{}` has uncommitted changes; refusing to touch it — commit or \
3516 discard them first",
3517 winner.worktree.display()
3518 );
3519 }
3520 git::worktree_remove(&self.state.repo, &winner.worktree)
3521 .await
3522 .ok();
3523 }
3524 let fix_worktree = self.state.worktree_root().join("operator-fix");
3525 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3526 git::git(
3527 &self.state.repo,
3528 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3529 )
3530 .await
3531 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3532 if !git::is_clean(&fix_worktree).await? {
3533 git::worktree_remove(&self.state.repo, &fix_worktree)
3534 .await
3535 .ok();
3536 bail!(
3537 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3538 winner.branch
3539 );
3540 }
3541
3542 let run_id = self.state.id.clone();
3543 let prompts = self.state.config.prompts.clone();
3544 let language = self.state.config.graph.language.clone();
3545 let sessions = self.state.config.graph.sessions;
3546 let artifacts = agent::artifacts_dir(&self.state.dir());
3547 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3548 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3549 _ => (
3550 self.state
3551 .config
3552 .agent(&winner.agent)
3553 .cloned()
3554 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3555 format!("impl-{}", winner.label),
3556 ),
3557 };
3558 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3559 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3560 .findings
3561 .iter()
3562 .map(|f| Finding {
3563 id: f.id.clone(),
3564 severity: f.severity,
3565 file: f.file.clone(),
3566 line: f.line,
3567 title: f.title.clone(),
3568 detail: f.detail.clone(),
3569 })
3570 .collect();
3571 let job = SeatJob {
3572 prompt: prompt::operator_fix(
3573 &self.state.instruction,
3574 &finding_list,
3575 reason,
3576 &stale_details,
3577 &head_at_request,
3578 &language,
3579 ),
3580 spec: fix_spec.clone(),
3581 seat,
3582 cwd: fix_worktree.clone(),
3583 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3584 allow_write: true,
3585 sessions,
3586 artifacts: artifacts.clone(),
3587 stem: "operator-fix".to_owned(),
3588 };
3589 let cache = self.state.config.cache_dir();
3590 let ctx = WaveCtx {
3591 run: &run_id,
3592 node: "fix",
3593 prompts: &prompts,
3594 cache: cache.as_deref(),
3595 round: None,
3596 };
3597 let (seat, out) =
3598 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3599 let agent_id = seat.agent.clone();
3600
3601 let mut fix = FixRecord {
3602 agent: agent_id,
3603 addressed: Vec::new(),
3604 rejected: Vec::new(),
3605 notes: String::new(),
3606 committed: false,
3607 failed: None,
3608 duration_ms: 0,
3609 continuation: None,
3610 };
3611 let mut final_seat = seat.clone();
3612 match out {
3613 AgentOutcome::Ok(o) => {
3614 fix.duration_ms = o.duration_ms;
3615 let parsed = verdict::extract_json::<FixReport>(&o.text);
3616 let incomplete_reason = match &parsed {
3617 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3618 "the reply parsed, but it reported a command whose own CLI \
3619 never confirmed an exit status"
3620 .to_owned(),
3621 ),
3622 Ok(_) => None,
3623 Err(e) => Some(e.to_string()),
3624 };
3625 match incomplete_reason {
3626 None => {
3627 let report = parsed.expect("checked Ok above");
3628 fix.addressed = report.addressed;
3629 fix.rejected = report.rejected;
3630 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3631 }
3632 Some(reason) => {
3633 let (resumed_seat, resolved, failure, cont) = self
3634 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3635 .await;
3636 fix.duration_ms += cont.cumulative_wait_ms;
3637 fix.continuation = Some(cont);
3638 final_seat = resumed_seat;
3639 match resolved {
3640 Some(report) => {
3641 fix.addressed = report.addressed;
3642 fix.rejected = report.rejected;
3643 fix.notes =
3644 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3645 }
3646 None => fix.failed = failure,
3647 }
3648 }
3649 }
3650 }
3651 AgentOutcome::Dropped(o) => {
3652 fix.duration_ms = o.duration_ms;
3653 let why = o
3654 .dropped
3655 .as_ref()
3656 .map(|d| d.why.as_str())
3657 .unwrap_or("the CLI ended the stream without delivering its answer");
3658 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3659 }
3660 AgentOutcome::Quota(o) => {
3661 self.state.quota.push(QuotaLoss {
3662 seat: final_seat.key.clone(),
3663 node: "fix".to_owned(),
3664 at: Timestamp::now(),
3665 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3666 });
3667 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3668 }
3669 AgentOutcome::Failed(e) => fix.failed = Some(e),
3670 }
3671 if fix.continuation.is_none() {
3672 fix.continuation = Some(ContinuationRecord::not_needed());
3673 }
3674 self.state.seats.insert(final_seat.key.clone(), final_seat);
3675
3676 let rescue_message = format!(
3677 "magi: operator-selected fix ({}) (uncommitted work)",
3678 self.state.operator_fixes[request_index]
3679 .findings
3680 .iter()
3681 .map(|f| f.id.as_str())
3682 .collect::<Vec<_>>()
3683 .join(", ")
3684 );
3685 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3686 self.state.note_withheld("fix", &r.withheld);
3687 }
3688 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3689 fix.committed = after != head_at_request;
3690 git::worktree_remove(&self.state.repo, &fix_worktree)
3691 .await
3692 .ok();
3693
3694 self.state.event(
3695 "fix",
3696 match &fix.failed {
3697 Some(reason) => format!(
3698 "operator fix: adoption report was lost ({reason}); {}",
3699 if fix.committed {
3700 "committed"
3701 } else {
3702 "NO new commit"
3703 }
3704 ),
3705 None => format!(
3706 "operator fix: {} addressed, {} rejected, {}",
3707 fix.addressed.len(),
3708 fix.rejected.len(),
3709 if fix.committed {
3710 "committed"
3711 } else {
3712 "NO new commit"
3713 }
3714 ),
3715 },
3716 );
3717
3718 for f in &mut self.state.operator_fixes[request_index].findings {
3725 f.outcome = if fix.failed.is_some() {
3726 OperatorFixOutcome::Unreported
3727 } else if fix.addressed.contains(&f.id) {
3728 OperatorFixOutcome::Addressed
3729 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3730 OperatorFixOutcome::Rejected { why: r.why.clone() }
3731 } else {
3732 OperatorFixOutcome::Unreported
3733 };
3734 }
3735
3736 let committed = fix.committed;
3737 if committed {
3738 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3739 }
3740 self.state.operator_fixes[request_index].fix = Some(fix);
3741 self.state.save()?;
3744
3745 if committed {
3746 self.state.event(
3747 "fix",
3748 format!(
3749 "operator fix committed {}; opening a follow-up review-only run",
3750 short(&after)
3751 ),
3752 );
3753 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3754 Ok(mut follow_up) => {
3755 follow_up.state.event(
3756 "start",
3757 format!(
3758 "requested by an operator fix on run {} for finding(s) {}",
3759 self.state.id,
3760 self.state.operator_fixes[request_index]
3761 .findings
3762 .iter()
3763 .map(|f| f.id.as_str())
3764 .collect::<Vec<_>>()
3765 .join(", "),
3766 ),
3767 );
3768 follow_up.state.save()?;
3769 let follow_up_id = follow_up.state.id.clone();
3770 if let Err(e) = follow_up.execute().await {
3771 self.state.event(
3772 "fix",
3773 format!(
3774 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3775 ),
3776 );
3777 }
3778 self.state.operator_fixes[request_index].follow_up_review_run =
3779 Some(follow_up_id);
3780 }
3781 Err(e) => {
3782 self.state.event(
3783 "fix",
3784 format!("committed the fix but could not open a follow-up review: {e:#}"),
3785 );
3786 }
3787 }
3788 self.state.save()?;
3789 }
3790
3791 Ok(())
3792 }
3793
3794 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3801 match &self.roles.fixer {
3802 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3803 _ => (
3804 self.state
3805 .config
3806 .agent(&winner.agent)
3807 .cloned()
3808 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3809 format!("impl-{}", winner.label),
3810 ),
3811 }
3812 }
3813
3814 async fn review_loop(&mut self) -> Result<()> {
3815 if self
3820 .state
3821 .base_sync
3822 .as_ref()
3823 .is_some_and(|s| s.conflict.is_some())
3824 {
3825 return Ok(());
3826 }
3827 let run_id = self.state.id.clone();
3832 let prompts = self.state.config.prompts.clone();
3833 let Some(winner) = self.state.winner().cloned() else {
3834 return Ok(());
3835 };
3836 let max_rounds = self.state.config.graph.review_rounds;
3837 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
3847 self.state.status = status;
3848 self.state.save()?;
3849 return Ok(());
3850 }
3851 self.state.status = RunStatus::Reviewing;
3852 if self
3862 .state
3863 .reviews
3864 .last()
3865 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
3866 {
3867 let shell = self.state.config.shell();
3868 return self
3869 .stop_reviewing(
3870 "the last round's own verification never resolved",
3871 &shell,
3872 &winner.worktree,
3873 )
3874 .await;
3875 }
3876
3877 let repo = self.state.repo.clone();
3878 let root = self.state.worktree_root();
3879 let language = self.state.config.graph.language.clone();
3880 let sessions = self.state.config.graph.sessions;
3881 let artifacts = agent::artifacts_dir(&self.state.dir());
3882 let base = self.landing_base();
3883 let base_short = short(&base);
3884 let reviewers = self.roles.reviewers.clone();
3885 let shell = self.state.config.shell();
3886
3887 for round in (self.state.reviews.len() + 1)..=max_rounds {
3888 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3889 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
3890 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
3891 let prev_verification = self
3900 .state
3901 .reviews
3902 .last()
3903 .and_then(|r| r.verification_summary(&head));
3904
3905 let mut jobs = Vec::new();
3909 for (r, spec) in reviewers.iter().cloned().enumerate() {
3910 let wt = root.join(format!("review-{}", r + 1));
3911 if wt.exists() {
3912 git::reset_detached(&wt, &head).await?;
3913 } else {
3914 git::worktree_add_detached(&repo, &wt, &head).await?;
3915 }
3916 let seat_key = format!("review-{}", r + 1);
3917 let seat = self.seat(&seat_key, &spec.id);
3918 jobs.push(SeatJob {
3919 prompt: prompt::review(&prompt::ReviewCtx {
3920 instruction: &self.state.instruction,
3921 branch: &winner.branch,
3922 base_short: &base_short,
3923 stat: &stat,
3924 patch: &patch,
3925 verification: prev_verification.as_ref(),
3926 reviewers: reviewers.len(),
3927 round,
3928 rounds: max_rounds,
3929 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
3932 lens: Lens::for_seat(r),
3933 language: &language,
3934 }),
3935 spec,
3936 seat,
3937 cwd: wt,
3938 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
3939 allow_write: false,
3940 sessions,
3941 artifacts: artifacts.clone(),
3942 stem: format!("review-{round}-{}", r + 1),
3943 });
3944 }
3945
3946 self.state.event(
3947 "review",
3948 format!(
3949 "round {round}: {} reviewers on {}",
3950 jobs.len(),
3951 short(&head)
3952 ),
3953 );
3954 let mut quota_losses = Vec::new();
3955 let review_retries = self.state.config.graph.retries;
3956 let review_cache = self.state.config.cache_dir();
3957 let ctx = WaveCtx {
3958 run: &run_id,
3959 node: "review",
3960 prompts: &prompts,
3961 cache: review_cache.as_deref(),
3962 round: Some(round),
3963 };
3964 let results = ask_json_wave::<Review>(
3965 jobs,
3966 Arc::clone(&self.sem),
3967 review_retries,
3968 &ctx,
3969 &mut quota_losses,
3970 &mut self.state,
3971 &|_: &Review| Ok(()),
3972 )
3973 .await;
3974 let round_quota_missing = quota_losses.len();
3978 self.state.quota.extend(quota_losses);
3979
3980 let mut records = Vec::new();
3981 let mut all_findings = Vec::new();
3982 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
3983 let agent_id = seat.agent.clone();
3984 self.state.seats.insert(seat.key.clone(), seat);
3985 let mut record = ReviewRecord {
3986 reviewer: r + 1,
3987 agent: agent_id,
3988 summary: String::new(),
3989 findings: Vec::new(),
3990 vote: None,
3991 failed: None,
3992 duration_ms: 0,
3993 attempts,
3999 };
4000 match res {
4001 Ok((review, out)) => {
4002 record.summary =
4010 blind::sanitize_prose(&review.summary, &self.state.config.blind);
4011 record.vote = Some(review.vote);
4012 record.duration_ms = out.duration_ms;
4013 for (n, mut f) in review.findings.into_iter().enumerate() {
4014 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
4017 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
4018 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
4019 f.file = f
4025 .file
4026 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
4027 all_findings.push(f.clone());
4028 record.findings.push(f);
4029 }
4030 self.state.event(
4031 "review",
4032 format!(
4033 "round {round}: reviewer {} voted {} with {} finding(s)",
4034 r + 1,
4035 review.vote.label(),
4036 record.findings.len()
4037 ),
4038 );
4039 }
4040 Err(e) => {
4041 record.failed = Some(e.to_string());
4042 self.state.event(
4043 "review",
4044 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
4045 );
4046 }
4047 }
4048 records.push(record);
4049 }
4050
4051 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
4058 let vote_split =
4059 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
4060 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
4061 if vote_split {
4062 self.state.event(
4063 "review",
4064 format!(
4065 "round {round}: votes split ({}) — one round of reconsideration",
4066 initial_votes
4067 .iter()
4068 .map(|v| v.label())
4069 .collect::<Vec<_>>()
4070 .join(", ")
4071 ),
4072 );
4073 let panel: Vec<ReviewSeatReport<'_>> = records
4076 .iter()
4077 .filter_map(|r| {
4078 r.vote.map(|vote| ReviewSeatReport {
4079 reviewer: r.reviewer,
4080 vote,
4081 summary: &r.summary,
4082 findings: &r.findings,
4083 })
4084 })
4085 .collect();
4086
4087 let mut jobs = Vec::new();
4088 let mut seats_at = Vec::new();
4089 for (r, spec) in reviewers.iter().cloned().enumerate() {
4090 if records[r].vote.is_none() {
4094 continue;
4095 }
4096 let wt = root.join(format!("review-{}", r + 1));
4097 let seat_key = format!("review-{}", r + 1);
4098 let seat = self.seat(&seat_key, &spec.id);
4099 let patch_ctx = if has_context(&spec, &seat, sessions) {
4104 None
4105 } else {
4106 Some(ReviewPatch {
4107 branch: &winner.branch,
4108 base_short: &base_short,
4109 stat: &stat,
4110 patch: &patch,
4111 })
4112 };
4113 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
4114 instruction: &self.state.instruction,
4115 reviewer: r + 1,
4116 lens: Lens::for_seat(r),
4117 panel: &panel,
4118 patch: patch_ctx,
4119 round,
4120 rounds: max_rounds,
4121 language: &language,
4122 });
4123 jobs.push(SeatJob {
4124 prompt,
4125 spec,
4126 seat,
4127 cwd: wt,
4128 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4129 allow_write: false,
4130 sessions,
4131 artifacts: artifacts.clone(),
4132 stem: format!("review-{round}-reconsider-{}", r + 1),
4133 });
4134 seats_at.push(r);
4135 }
4136
4137 let mut recon_quota_losses = Vec::new();
4138 let recon_cache = self.state.config.cache_dir();
4139 let recon_ctx = WaveCtx {
4140 run: &run_id,
4141 node: "review",
4142 prompts: &prompts,
4143 cache: recon_cache.as_deref(),
4144 round: Some(round),
4145 };
4146 let recon_results = ask_json_wave::<ReviewRevote>(
4147 jobs,
4148 Arc::clone(&self.sem),
4149 review_retries,
4150 &recon_ctx,
4151 &mut recon_quota_losses,
4152 &mut self.state,
4153 &|_: &ReviewRevote| Ok(()),
4154 )
4155 .await;
4156 self.state.quota.extend(recon_quota_losses);
4157
4158 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
4159 let agent_id = seat.agent.clone();
4160 self.state.seats.insert(seat.key.clone(), seat);
4161 let mut rec = ReviewRevoteRecord {
4162 reviewer: r + 1,
4163 agent: agent_id,
4164 vote: None,
4165 reason: String::new(),
4166 failed: None,
4167 };
4168 match res {
4169 Ok((rv, _)) => {
4170 rec.vote = Some(rv.vote);
4171 rec.reason =
4172 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
4173 self.state.event(
4174 "review",
4175 format!(
4176 "round {round}: reviewer {} revoted {}",
4177 r + 1,
4178 rv.vote.label()
4179 ),
4180 );
4181 }
4182 Err(e) => {
4183 rec.failed = Some(e.to_string());
4184 self.state.event(
4185 "review",
4186 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4187 );
4188 }
4189 }
4190 reconsideration.push(rec);
4191 }
4192 } else if initial_votes.len() > 1 {
4193 self.state.event(
4194 "review",
4195 format!(
4196 "round {round}: votes agreed ({}) — no reconsideration",
4197 initial_votes[0].label()
4198 ),
4199 );
4200 }
4201
4202 let final_votes: Vec<ReviewVote> = records
4206 .iter()
4207 .filter_map(|r| {
4208 reconsideration
4209 .iter()
4210 .find(|rv| rv.reviewer == r.reviewer)
4211 .and_then(|rv| rv.vote)
4212 .or(r.vote)
4213 })
4214 .collect();
4215 let round_verdict = ReviewVote::worst(final_votes);
4216
4217 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4218 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4219 let defer_e2e =
4230 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4231 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4232 let reason =
4233 format!("{blocking} blocking finding(s) already required a fix this round");
4234 self.state.event(
4235 "verify",
4236 format!(
4237 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4238 {}); it will run once a round has none left",
4239 short(&head)
4240 ),
4241 );
4242 (Vec::new(), false, true, Some(reason))
4243 } else {
4244 let e2e_commands = self.state.config.verify.e2e.clone();
4245 let cache_dir = self.state.config.cache_dir();
4246 let context = format!("round {round}");
4247 let (e2e, verify_retried) = with_cache_lease(
4248 &mut self.state,
4249 cache_dir.as_deref(),
4250 "e2e",
4251 "e2e",
4252 &winner.worktree,
4253 &head,
4254 verify_timeout,
4255 &context,
4256 |state, budget| {
4257 let shell = shell.clone();
4258 let e2e_commands = e2e_commands.clone();
4259 let worktree = winner.worktree.clone();
4260 let context = context.clone();
4261 async move {
4262 run_e2e_with_retry(
4263 state,
4264 &shell,
4265 &e2e_commands,
4266 &worktree,
4267 budget,
4268 &context,
4269 )
4270 .await
4271 }
4272 },
4273 )
4274 .await;
4275 (e2e, verify_retried, false, None)
4276 };
4277
4278 let expected = records.len();
4279 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4280 let incomplete = answered < expected;
4281 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4282 let policy = self.state.config.graph.incomplete_review;
4283 let clean = round_is_clean(
4284 blocking,
4285 e2e_ok,
4286 answered,
4287 expected,
4288 round_quota_missing,
4289 policy,
4290 );
4291
4292 let mut round_record = ReviewRound {
4293 round,
4294 head: head.clone(),
4295 verified_head: None,
4296 verified_at: None,
4297 reviews: records,
4298 e2e,
4299 verify_retried,
4300 e2e_deferred,
4301 e2e_defer_reason,
4302 fix: None,
4303 blocking,
4304 answered,
4305 expected,
4306 clean,
4307 progressed: false,
4308 vote_split,
4309 reconsideration,
4310 verdict: round_verdict,
4311 };
4312 if !matches!(
4323 round_record.e2e_status(),
4324 E2eStatus::Deferred | E2eStatus::NotConfigured
4325 ) {
4326 round_record.verified_head = Some(head.clone());
4327 round_record.verified_at = Some(Timestamp::now());
4328 }
4329 let this_round_verification = round_record.verification_summary(&head);
4330
4331 if incomplete {
4332 let missing: Vec<String> = round_record
4333 .reviews
4334 .iter()
4335 .filter(|r| r.failed.is_some())
4336 .map(|r| format!("review-{}", r.reviewer))
4337 .collect();
4338 self.state.event(
4339 "review",
4340 format!(
4341 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4342 missing.join(", ")
4343 ),
4344 );
4345 }
4346
4347 if clean {
4348 self.state.event(
4349 "review",
4350 if incomplete && policy == IncompleteReviewPolicy::Warn {
4351 format!(
4352 "round {round}: clean (warn policy, incomplete panel) — no \
4353 blocking findings from the seats that answered, verification green"
4354 )
4355 } else if incomplete {
4356 format!(
4357 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4358 quorum) — no blocking findings from the seats that answered, \
4359 verification green",
4360 expected - answered
4361 )
4362 } else {
4363 format!("round {round}: clean — no blocking findings, verification green")
4364 },
4365 );
4366 self.state.reviews.push(round_record);
4367 self.state.status = RunStatus::Gating;
4368 self.state.save()?;
4369 return Ok(());
4370 }
4371
4372 if incomplete && blocking == 0 && e2e_ok {
4380 self.state.reviews.push(round_record);
4381 self.state.save()?;
4382 if round == max_rounds {
4383 self.state.status = RunStatus::Blocked;
4384 self.state.event(
4385 "review",
4386 format!(
4387 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4388 refusing to call it clean",
4389 expected - answered
4390 ),
4391 );
4392 return Ok(());
4393 }
4394 continue;
4395 }
4396
4397 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4409 self.state.reviews.push(round_record);
4410 return self
4411 .stop_reviewing(
4412 "the round's own verification could not run",
4413 &shell,
4414 &winner.worktree,
4415 )
4416 .await;
4417 }
4418
4419 if round == max_rounds {
4420 self.state.reviews.push(round_record);
4421 return self
4422 .stop_reviewing(
4423 &format!(
4424 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4425 ),
4426 &shell,
4427 &winner.worktree,
4428 )
4429 .await;
4430 }
4431
4432 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4435 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4436 let blocking_findings: Vec<_> = all_findings
4437 .iter()
4438 .filter(|f| f.severity.blocks())
4439 .cloned()
4440 .collect();
4441 let job = SeatJob {
4442 prompt: prompt::fix(
4443 &self.state.instruction,
4444 &blocking_findings,
4445 this_round_verification.as_ref(),
4446 round,
4447 max_rounds,
4448 &language,
4449 ),
4450 spec: fix_spec.clone(),
4451 seat,
4452 cwd: winner.worktree.clone(),
4453 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4454 allow_write: true,
4455 sessions,
4456 artifacts: artifacts.clone(),
4457 stem: format!("fix-{round}"),
4458 };
4459 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4460 let cache = self.state.config.cache_dir();
4461 let ctx = WaveCtx {
4462 run: &run_id,
4463 node: "fix",
4464 prompts: &prompts,
4465 cache: cache.as_deref(),
4466 round: Some(round),
4467 };
4468 let (seat, out) =
4469 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4470 let agent_id = seat.agent.clone();
4471
4472 let mut fix = FixRecord {
4473 agent: agent_id,
4474 addressed: Vec::new(),
4475 rejected: Vec::new(),
4476 notes: String::new(),
4477 committed: false,
4478 failed: None,
4479 duration_ms: 0,
4480 continuation: None,
4481 };
4482 let mut continuation = ContinuationRecord::not_needed();
4483 let mut final_seat = seat.clone();
4484 match out {
4485 AgentOutcome::Ok(o) => {
4486 fix.duration_ms = o.duration_ms;
4487 let parsed = verdict::extract_json::<FixReport>(&o.text);
4488 let incomplete_reason = match &parsed {
4495 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4496 "the reply parsed, but it reported a command whose own CLI never \
4497 confirmed an exit status"
4498 .to_owned(),
4499 ),
4500 Ok(_) => None,
4501 Err(e) => Some(e.to_string()),
4502 };
4503 match incomplete_reason {
4504 None => {
4505 let report = parsed.expect("checked Ok above");
4506 fix.addressed = report.addressed;
4507 fix.rejected = report.rejected;
4508 fix.notes =
4509 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4510 }
4511 Some(reason) => {
4512 let (resumed_seat, resolved, failure, cont) = self
4513 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4514 .await;
4515 fix.duration_ms += cont.cumulative_wait_ms;
4516 continuation = cont;
4517 final_seat = resumed_seat;
4518 match resolved {
4519 Some(report) => {
4520 fix.addressed = report.addressed;
4521 fix.rejected = report.rejected;
4522 fix.notes = blind::sanitize_prose(
4523 &report.notes,
4524 &self.state.config.blind,
4525 );
4526 }
4527 None => fix.failed = failure,
4528 }
4529 }
4530 }
4531 }
4532 AgentOutcome::Dropped(o) => {
4534 fix.duration_ms = o.duration_ms;
4535 let why = o
4536 .dropped
4537 .as_ref()
4538 .map(|d| d.why.as_str())
4539 .unwrap_or("the CLI ended the stream without delivering its answer");
4540 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4541 }
4542 AgentOutcome::Quota(o) => {
4543 self.state.quota.push(QuotaLoss {
4544 seat: final_seat.key.clone(),
4545 node: "fix".to_owned(),
4546 at: Timestamp::now(),
4547 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4548 });
4549 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4550 }
4551 AgentOutcome::Failed(e) => fix.failed = Some(e),
4552 }
4553 fix.continuation = Some(continuation);
4554 self.state.seats.insert(final_seat.key.clone(), final_seat);
4555 if let Ok(r) = git::rescue_commit(
4556 &winner.worktree,
4557 &format!("magi: review round {round} fixes (uncommitted work)"),
4558 )
4559 .await
4560 {
4561 self.state.note_withheld("fix", &r.withheld);
4562 }
4563 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4564 fix.committed = after != before;
4565 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4573 let progressed = diff_after != patch;
4574 let commit_note = if fix.committed {
4575 "committed"
4576 } else {
4577 "NO new commit"
4578 };
4579 let tree_note = if progressed {
4580 "changed vs base"
4581 } else {
4582 "unchanged vs base"
4583 };
4584 self.state.event(
4585 "fix",
4586 match &fix.failed {
4587 Some(reason) => {
4593 format!(
4594 "round {round}: fixer's adoption report was lost ({reason}); \
4595 {commit_note}, tree {tree_note}"
4596 )
4597 }
4598 None => format!(
4599 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4600 {tree_note}{}",
4601 fix.addressed.len(),
4602 fix.rejected.len(),
4603 if continuation.outcome == ContinuationOutcome::Resumed {
4604 format!(
4605 " (adoption report recovered after {} continuation(s))",
4606 continuation.attempts
4607 )
4608 } else {
4609 String::new()
4610 },
4611 ),
4612 },
4613 );
4614 round_record.fix = Some(fix);
4615 round_record.progressed = progressed;
4616 self.state.reviews.push(round_record);
4617 self.state.save()?;
4618
4619 if matches!(
4632 continuation.outcome,
4633 ContinuationOutcome::Exhausted
4634 | ContinuationOutcome::QuotaLost
4635 | ContinuationOutcome::NoSession
4636 ) {
4637 return self
4638 .stop_reviewing(
4639 "the fixer's adoption report never came back, even after resuming its \
4640 own seat; refusing to start another round against the same worktree \
4641 while that is unresolved",
4642 &shell,
4643 &winner.worktree,
4644 )
4645 .await;
4646 }
4647
4648 let streak = self
4649 .state
4650 .reviews
4651 .iter()
4652 .rev()
4653 .take_while(|r| !r.progressed)
4654 .count();
4655 if streak >= STAGNANT_LIMIT {
4656 return self
4657 .stop_reviewing(
4658 &format!(
4659 "the tree has not moved against base for {streak} round(s) in a row"
4660 ),
4661 &shell,
4662 &winner.worktree,
4663 )
4664 .await;
4665 }
4666 }
4667 Ok(())
4668 }
4669
4670 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4700 let round_idx = self.state.reviews.len() - 1;
4701 let needs_catchup_run = matches!(
4709 self.state.reviews[round_idx].e2e_status(),
4710 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4711 );
4712 if needs_catchup_run {
4713 let round = self.state.reviews[round_idx].round;
4714 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4715 let commands = self.state.config.verify.e2e.clone();
4716 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4717 let cache_dir = self.state.config.cache_dir();
4718 let context = format!(
4719 "round {round}: verification unresolved, catching up before the final decision"
4720 );
4721 let (outcomes, verify_retried) = with_cache_lease(
4722 &mut self.state,
4723 cache_dir.as_deref(),
4724 "e2e",
4725 "e2e",
4726 worktree,
4727 &attempted_head,
4728 timeout,
4729 &context,
4730 |state, budget| {
4731 let shell = shell.to_vec();
4732 let commands = commands.clone();
4733 let context = context.clone();
4734 async move {
4735 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4736 .await
4737 }
4738 },
4739 )
4740 .await;
4741 let last = &mut self.state.reviews[round_idx];
4742 last.e2e = outcomes;
4743 last.verify_retried = verify_retried;
4744 last.verified_head = Some(attempted_head);
4751 last.verified_at = Some(Timestamp::now());
4752 if verify_inconclusive(&last.e2e) {
4753 self.state.save()?;
4760 return Ok(());
4761 }
4762 last.e2e_deferred = false;
4763 }
4764 let last = &self.state.reviews[round_idx];
4765 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4766
4767 match last.e2e_status() {
4768 E2eStatus::Failed => {
4769 let red: Vec<String> = last
4770 .e2e
4771 .iter()
4772 .filter(|o| !o.ok())
4773 .map(|o| {
4774 format!(
4775 "`{}` -> {:?}\n{}",
4776 o.command,
4777 o.code,
4778 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4779 )
4780 })
4781 .collect();
4782 self.state
4783 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4784 self.state.status = RunStatus::Blocked;
4785 }
4786 E2eStatus::ResourceBlocked => {
4791 self.state.event(
4792 "review",
4793 format!(
4794 "{why}; e2e could not run (shared build cache unavailable); not \
4795 deciding yet"
4796 ),
4797 );
4798 }
4799 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4800 self.state.event(
4801 "review",
4802 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4803 );
4804 self.state.status = RunStatus::Gating;
4805 }
4806 }
4807 self.state.save()?;
4808 Ok(())
4809 }
4810
4811 async fn gate(&mut self) -> Result<()> {
4814 if self.state.status == RunStatus::Failed
4826 || self
4827 .state
4828 .base_sync
4829 .as_ref()
4830 .is_some_and(|s| s.conflict.is_some())
4831 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
4832 != Some(RunStatus::Gating)
4833 {
4834 return Ok(());
4835 }
4836 if self.state.gate_ran {
4837 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
4848 self.state.status = RunStatus::Blocked;
4849 self.state.save()?;
4850 }
4851 return Ok(());
4852 }
4853 let Some(winner) = self.state.winner().cloned() else {
4854 return Ok(());
4855 };
4856 self.state.status = RunStatus::Gating;
4857 let mut outcomes = self.run_gate(&winner).await?;
4858 loop {
4859 if verify_inconclusive(&outcomes) {
4870 self.state.save()?;
4871 return Ok(());
4872 }
4873 if outcomes.iter().all(CommandOutcome::ok) {
4874 break;
4875 }
4876 match self.gate_fix_round(&winner, &outcomes).await? {
4877 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
4878 GateFix::Stop => break,
4879 GateFix::Defer => {
4880 self.state.save()?;
4881 return Ok(());
4882 }
4883 }
4884 }
4885 let passed = outcomes.iter().all(CommandOutcome::ok);
4886 self.state.gate = outcomes;
4887 self.state.gate_ran = true;
4888 if !passed {
4889 self.state.status = RunStatus::Blocked;
4890 let spent = self.state.gate_fixes.len();
4891 self.state.event(
4892 "gate",
4893 if spent == 0 {
4894 "gate failed; not merging".to_owned()
4895 } else {
4896 format!("gate failed after {spent} gate-fix round(s); not merging")
4897 },
4898 );
4899 }
4900 self.state.save()?;
4901 Ok(())
4902 }
4903
4904 async fn run_pre_gate(&mut self, winner: &Candidate) {
4914 let commands = self.state.config.verify.pre_gate.clone();
4915 if commands.is_empty() {
4916 return;
4917 }
4918 let shell = self.state.config.shell();
4919 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4920 let (outcomes, _) = run_commands(
4921 &mut self.state,
4922 "pre_gate",
4923 "pre_gate",
4924 0,
4925 &shell,
4926 &commands,
4927 &winner.worktree,
4928 timeout,
4929 )
4930 .await;
4931 for o in &outcomes {
4932 if !o.ok() {
4933 tracing::warn!(
4934 "pre_gate `{}` failed ({:?}); the gate decides",
4935 o.command,
4936 o.code
4937 );
4938 }
4939 self.state.event(
4940 "pre_gate",
4941 format!(
4942 "`{}` -> {}",
4943 o.command,
4944 if o.ok() {
4945 "pass".to_owned()
4946 } else {
4947 format!(
4948 "FAIL ({:?})\n{}",
4949 o.code,
4950 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4951 )
4952 }
4953 ),
4954 );
4955 }
4956 self.state.pre_gate = outcomes;
4957 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
4958 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
4959 Ok(head) => {
4960 self.state
4961 .event("pre_gate", format!("committed mechanical fixes ({head})"));
4962 self.state.pre_gate_commit = Some(head);
4963 }
4964 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
4965 },
4966 Ok(false) => {}
4967 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
4968 }
4969 if let Err(e) = self.state.save() {
4970 tracing::warn!("could not persist the pre_gate record: {e:#}");
4971 }
4972 }
4973
4974 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
4977 self.run_pre_gate(winner).await;
4978 let shell = self.state.config.shell();
4979 let gate_commands = self.state.config.verify.gate.clone();
4980 let outcomes = if gate_commands.is_empty() {
4989 Vec::new()
4990 } else {
4991 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4992 let cache_dir = self.state.config.cache_dir();
4993 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
4994 let (outcomes, _) = with_cache_lease(
4995 &mut self.state,
4996 cache_dir.as_deref(),
4997 "gate",
4998 "gate",
4999 &winner.worktree,
5000 &head,
5001 timeout,
5002 "final gate",
5003 |state, budget| {
5004 let shell = shell.clone();
5005 let gate_commands = gate_commands.clone();
5006 let worktree = winner.worktree.clone();
5007 async move {
5008 let (outcomes, timed_out_pids) = run_commands(
5009 state,
5010 "gate",
5011 "gate",
5012 0,
5013 &shell,
5014 &gate_commands,
5015 &worktree,
5016 budget,
5017 )
5018 .await;
5019 (outcomes, false, timed_out_pids)
5020 }
5021 },
5022 )
5023 .await;
5024 outcomes
5025 };
5026 if outcomes.is_empty() {
5027 self.state.event(
5032 "gate",
5033 "no gate commands configured; nothing to check, passing",
5034 );
5035 }
5036 for o in &outcomes {
5037 self.state.event(
5038 "gate",
5039 format!(
5040 "`{}` -> {}",
5041 o.command,
5042 if o.ok() {
5043 "pass".to_owned()
5044 } else {
5045 format!(
5046 "FAIL ({:?})\n{}",
5047 o.code,
5048 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5049 )
5050 }
5051 ),
5052 );
5053 }
5054 Ok(outcomes)
5055 }
5056
5057 async fn gate_fix_round(
5070 &mut self,
5071 winner: &Candidate,
5072 outcomes: &[CommandOutcome],
5073 ) -> Result<GateFix> {
5074 let cap = self.state.config.graph.gate_fix_rounds;
5075 let spent = self.state.gate_fixes.len();
5076 if spent >= cap {
5077 if cap > 0 {
5078 self.state.event(
5079 "gate",
5080 format!("{spent} gate-fix round(s) spent and the gate still fails"),
5081 );
5082 }
5083 return Ok(GateFix::Stop);
5084 }
5085 if !gate_fixable(outcomes) {
5086 self.state.event(
5087 "gate",
5088 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
5089 command or similar); not spending a fix round on it",
5090 );
5091 return Ok(GateFix::Stop);
5092 }
5093 let min_free = self.state.config.disk.min_free_bytes;
5094 if min_free > 0 {
5095 match crate::disk::free_bytes(&winner.worktree) {
5096 Ok(free) if crate::disk::enough_space(free, min_free) => {}
5097 Ok(free) => {
5098 self.state.event(
5099 "gate",
5100 format!(
5101 "only {free} bytes free ({min_free} required by `[disk] \
5102 min_free_bytes`); not spending a fix round on a failure the disk \
5103 may explain"
5104 ),
5105 );
5106 return Ok(GateFix::Stop);
5107 }
5108 Err(e) => {
5109 self.state.event(
5110 "gate",
5111 format!("free disk space could not be measured ({e:#}); no fix round"),
5112 );
5113 return Ok(GateFix::Stop);
5114 }
5115 }
5116 }
5117
5118 let attempt = spent + 1;
5119 let run_id = self.state.id.clone();
5120 let prompts = self.state.config.prompts.clone();
5121 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
5122 let base = self.landing_base();
5123 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
5124 let seat = self.seat(&fix_seat_key, &fix_spec.id);
5125 let job = SeatJob {
5126 prompt: prompt::gate_fix(
5127 &self.state.instruction,
5128 &failed,
5129 attempt,
5130 cap,
5131 &self.state.config.graph.language,
5132 ),
5133 spec: fix_spec,
5134 seat,
5135 cwd: winner.worktree.clone(),
5136 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
5137 allow_write: true,
5138 sessions: self.state.config.graph.sessions,
5139 artifacts: agent::artifacts_dir(&self.state.dir()),
5140 stem: format!("gate-fix-{attempt}"),
5141 };
5142 self.state.event(
5143 "gate",
5144 format!("gate failed; gate-fix round {attempt} of {cap}"),
5145 );
5146 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
5147 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
5148 let cache = self.state.config.cache_dir();
5149 let ctx = WaveCtx {
5150 run: &run_id,
5151 node: "gate-fix",
5152 prompts: &prompts,
5153 cache: cache.as_deref(),
5154 round: None,
5155 };
5156 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
5157 let mut record = GateFixRecord {
5158 agent: seat.agent.clone(),
5159 failed,
5160 notes: String::new(),
5161 committed: false,
5162 error: None,
5163 };
5164 match out {
5165 AgentOutcome::Ok(o) => {
5166 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
5169 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
5170 }
5171 }
5172 AgentOutcome::Dropped(_) => {
5173 record.error = Some("the CLI dropped the stream".to_owned());
5174 }
5175 AgentOutcome::Quota(o) => {
5176 self.state.quota.push(QuotaLoss {
5177 seat: seat.key.clone(),
5178 node: "gate-fix".to_owned(),
5179 at: Timestamp::now(),
5180 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5181 });
5182 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5183 }
5184 AgentOutcome::Failed(e) => record.error = Some(e),
5185 }
5186 self.state.seats.insert(seat.key.clone(), seat);
5187 if let Ok(r) = git::rescue_commit(
5188 &winner.worktree,
5189 &format!("magi: gate fix {attempt} (uncommitted work)"),
5190 )
5191 .await
5192 {
5193 self.state.note_withheld("gate-fix", &r.withheld);
5194 }
5195 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5196 record.committed = after != before;
5197 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5198 let note = record.error.clone();
5199 self.state.gate_fixes.push(record);
5200 self.state.save()?;
5201 if !changed {
5202 self.state.event(
5203 "gate",
5204 match note {
5205 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5206 None => format!("gate-fix round {attempt}: the tree did not change"),
5207 },
5208 );
5209 return Ok(GateFix::Stop);
5210 }
5211 self.state.event(
5212 "gate",
5213 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5214 );
5215
5216 let commands = self.state.config.verify.e2e.clone();
5217 if !commands.is_empty() {
5218 let shell = self.state.config.shell();
5219 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5220 let cache_dir = self.state.config.cache_dir();
5221 let context = format!("gate-fix round {attempt}");
5222 let (e2e, _) = with_cache_lease(
5223 &mut self.state,
5224 cache_dir.as_deref(),
5225 "e2e",
5226 "e2e",
5227 &winner.worktree,
5228 &after,
5229 timeout,
5230 &context,
5231 |state, budget| {
5232 let shell = shell.clone();
5233 let commands = commands.clone();
5234 let context = context.clone();
5235 let worktree = winner.worktree.clone();
5236 async move {
5237 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5238 .await
5239 }
5240 },
5241 )
5242 .await;
5243 if verify_inconclusive(&e2e) {
5244 return Ok(GateFix::Defer);
5245 }
5246 if e2e.iter().any(|o| !o.ok()) {
5247 self.state.event(
5248 "gate",
5249 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5250 );
5251 return Ok(GateFix::Stop);
5252 }
5253 }
5254 Ok(GateFix::Retry)
5255 }
5256
5257 async fn merge(&mut self) -> Result<()> {
5260 if self
5275 .state
5276 .base_sync
5277 .as_ref()
5278 .is_some_and(|s| s.conflict.is_some())
5279 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5280 != Some(RunStatus::Gating)
5281 || !self.state.gate_status().ok()
5290 {
5291 return Ok(());
5292 }
5293 if self.state.merge.is_some() {
5302 return Ok(());
5303 }
5304 let Some(winner) = self.state.winner().cloned() else {
5305 return Ok(());
5306 };
5307 let repo = self.state.repo.clone();
5308 let base = self.state.base_branch.clone();
5309 let mode = self.state.config.merge.mode;
5310 let style = self.state.config.merge.style;
5311 let pr = pr_message(&self.state, winner.label);
5312 let message = pr.commit_message();
5313
5314 let outcome = match mode {
5315 MergeMode::None => MergeOutcome {
5316 mode,
5317 ok: true,
5318 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5319 },
5320 MergeMode::Local => {
5321 let on = git::current_branch(&repo).await?;
5322 if on.as_deref() != Some(base.as_str()) {
5323 MergeOutcome {
5324 mode,
5325 ok: false,
5326 detail: format!(
5327 "{} has {} checked out, not the base branch {base}",
5328 repo.display(),
5329 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5330 ),
5331 }
5332 } else if !git::is_clean(&repo).await? {
5333 MergeOutcome {
5334 mode,
5335 ok: false,
5336 detail: format!("{} is dirty; refusing to merge", repo.display()),
5337 }
5338 } else {
5339 let out = match style {
5340 MergeStyle::Merge => {
5341 git::merge_no_ff(&repo, &winner.branch, &message).await?
5342 }
5343 MergeStyle::Squash => {
5344 git::merge_squash(&repo, &winner.branch, &message).await?
5345 }
5346 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5347 };
5348 MergeOutcome {
5349 mode,
5350 ok: out.ok(),
5351 detail: if out.ok() { out.stdout } else { out.stderr },
5352 }
5353 }
5354 }
5355 MergeMode::Pr => {
5356 let remote = self.state.config.merge.remote.clone();
5357 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5358 if !pushed.ok() {
5359 MergeOutcome {
5360 mode,
5361 ok: false,
5362 detail: pushed.stderr,
5363 }
5364 } else {
5365 let out =
5366 gh_pr_create(&winner.worktree, &base, &winner.branch, &pr.title, &pr.body)
5367 .await;
5368 match out {
5369 Ok(url) => MergeOutcome {
5370 mode,
5371 ok: true,
5372 detail: url,
5373 },
5374 Err(e) => MergeOutcome {
5375 mode,
5376 ok: false,
5377 detail: e.to_string(),
5378 },
5379 }
5380 }
5381 }
5382 };
5383
5384 self.state.status = match (mode, outcome.ok) {
5385 (MergeMode::None, _) => RunStatus::Ready,
5386 (_, true) => RunStatus::Merged,
5387 (_, false) => RunStatus::Blocked,
5388 };
5389 self.state.event(
5390 "merge",
5391 format!(
5392 "{:?}: {}",
5393 mode,
5394 outcome.detail.lines().next().unwrap_or("")
5395 ),
5396 );
5397 self.state.merge = Some(outcome);
5398 self.state.save()?;
5399
5400 if self.state.config.graph.land
5406 && mode == MergeMode::Pr
5407 && self.state.status == RunStatus::Merged
5408 {
5409 self.run_land().await?;
5410 }
5411 self.settle_questions();
5416 Ok(())
5417 }
5418
5419 async fn run_land(&mut self) -> Result<()> {
5430 let url = self
5431 .state
5432 .merge
5433 .as_ref()
5434 .map(|m| m.detail.clone())
5435 .unwrap_or_default();
5436 let url = url.lines().next().unwrap_or("").trim().to_owned();
5437 if !url.starts_with("http") {
5438 return Ok(());
5439 }
5440 match land::land(&mut self.state, &url).await {
5443 Ok(pr) if self.state.parked => {
5444 let _ = pr;
5448 }
5449 Ok(pr) => {
5450 self.state.status = match pr.state {
5451 land::PrLifecycle::Merged => RunStatus::Merged,
5452 _ => RunStatus::Blocked,
5453 };
5454 if bump::should_release_bump(self.state.status)
5461 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5462 {
5463 self.state
5469 .event("bump", format!("release bump skipped: {e:#}"));
5470 }
5471 self.state.save()?;
5472 }
5473 Err(e) => {
5474 self.state.status = RunStatus::Blocked;
5475 self.state.event("land", format!("gave up: {e}"));
5476 self.state.save()?;
5477 }
5478 }
5479 Ok(())
5480 }
5481
5482 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5486 if let Some(existing) = self.state.seats.get(key)
5487 && existing.agent == agent
5488 {
5489 return existing.clone();
5490 }
5491 let fresh = SeatState::new(key, agent, self.state.seed);
5492 self.state.seats.insert(key.to_owned(), fresh.clone());
5493 fresh
5494 }
5495
5496 fn view(&self, c: &Candidate) -> CandidateView {
5498 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5499 .unwrap_or_default();
5500 let (patch, _) = blind::sanitize_patch(
5501 &format!("candidate {} patch", c.label),
5502 &raw,
5503 &self.state.config.blind,
5504 );
5505 CandidateView {
5506 label: c.label,
5507 branch: c.branch.clone(),
5508 summary: c.summary.clone(),
5509 stat: c.stat.clone(),
5510 patch,
5511 }
5512 }
5513
5514 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5516 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5517 prompt::judge(
5518 "(see above)",
5519 &views,
5520 self.roles.judges.len(),
5521 base_short,
5522 "en",
5523 )
5524 }
5525
5526 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5533 let mut turns = Vec::new();
5534 for j in &self.state.judgements {
5535 if j.ranking.is_empty() {
5536 continue;
5537 }
5538 let reasons = j
5539 .reasons
5540 .iter()
5541 .map(|(k, v)| format!("- {k}: {v}"))
5542 .collect::<Vec<_>>()
5543 .join("\n");
5544 turns.push(Turn {
5545 who: format!("Judge {} (opening ranking)", j.judge),
5546 is_self: j.judge == self_idx + 1,
5547 body: format!(
5548 "Ranked {}{}{reasons}",
5549 j.ranking.iter().collect::<String>(),
5550 if reasons.is_empty() {
5551 ""
5552 } else {
5553 ", because:\n"
5554 }
5555 ),
5556 });
5557 }
5558 for t in self
5559 .state
5560 .deliberation
5561 .iter()
5562 .flat_map(|r| r.turns.iter())
5563 .chain(current)
5564 {
5565 turns.push(Turn {
5566 who: format!("Judge {}", t.judge),
5567 is_self: t.judge == self_idx + 1,
5568 body: t.body.clone(),
5569 });
5570 }
5571 turns
5572 }
5573}
5574
5575fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5577 agent::has_session(spec.kind, seat, sessions)
5578}
5579
5580fn next_untried_implementer<'a>(
5601 roster: &'a [AgentSpec],
5602 start: usize,
5603 tried: &BTreeSet<String>,
5604) -> Option<&'a AgentSpec> {
5605 roster
5606 .get(start + 1..)?
5607 .iter()
5608 .find(|s| !tried.contains(&s.id))
5609}
5610
5611fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5625 commands.iter().any(|c| c.exit_code.is_none())
5626}
5627
5628fn verified_noop_claim(
5641 usable: bool,
5642 commands: &[agent::CommandEvidence],
5643 text: &str,
5644) -> Option<String> {
5645 (usable && !has_unconfirmed_command(commands))
5646 .then(|| verdict::verified_noop(text))
5647 .flatten()
5648}
5649
5650fn short(commit: &str) -> String {
5651 commit.chars().take(7).collect()
5652}
5653
5654fn make_executable(path: &Path) -> Result<()> {
5655 #[cfg(unix)]
5656 {
5657 use std::os::unix::fs::PermissionsExt as _;
5658 let mut perms = std::fs::metadata(path)?.permissions();
5659 perms.set_mode(0o755);
5660 std::fs::set_permissions(path, perms)?;
5661 }
5662 #[cfg(not(unix))]
5663 {
5664 let _ = path;
5665 }
5666 Ok(())
5667}
5668
5669struct WaveCtx<'a> {
5676 run: &'a str,
5679 node: &'a str,
5681 prompts: &'a Prompts,
5682 cache: Option<&'a Path>,
5684 round: Option<usize>,
5687}
5688
5689async fn run_one(
5691 job: SeatJob,
5692 sem: Arc<Semaphore>,
5693 ctx: &WaveCtx<'_>,
5694 state: &mut RunState,
5695 attempt: usize,
5696) -> (SeatState, AgentOutcome) {
5697 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5698 .await
5699 .pop()
5700 .expect("one job in, one result out");
5701 (seat, out)
5702}
5703
5704async fn wave(
5710 jobs: Vec<SeatJob>,
5711 sem: Arc<Semaphore>,
5712 ctx: &WaveCtx<'_>,
5713 state: &mut RunState,
5714 attempt: usize,
5715) -> Vec<(usize, SeatState, AgentOutcome)> {
5716 let WaveCtx {
5717 run,
5718 node,
5719 prompts,
5720 cache,
5721 round,
5722 } = *ctx;
5723 for job in &jobs {
5724 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5725 }
5726 if let Err(e) = state.save() {
5727 tracing::warn!("could not persist in-progress seats: {e:#}");
5732 }
5733 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5749 let wait_started = Instant::now();
5750 let cache_guard = if let Some(cache_dir) = cache {
5751 if jobs_had_a_writer {
5752 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5753 let budget = jobs
5754 .iter()
5755 .map(|j| j.timeout)
5756 .max()
5757 .unwrap_or(Duration::from_secs(60));
5758 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5759 .await
5760 .ok()
5761 } else {
5762 None
5763 }
5764 } else {
5765 None
5766 };
5767 let waited_for_lease = wait_started.elapsed();
5774 let mut set = tokio::task::JoinSet::new();
5775 let overlay = prompts.overlay(node);
5776 for (i, mut job) in jobs.into_iter().enumerate() {
5777 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5778 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
5779 if cache.is_some() {
5780 job.prompt.push('\n');
5781 job.prompt
5782 .push_str(&prompt::build_cache_note(node, job.allow_write));
5783 }
5784 let sem = Arc::clone(&sem);
5785 let run = run.to_owned();
5786 let node = node.to_owned();
5787 let cache = cache
5798 .filter(|_| job.allow_write && cache_guard.is_some())
5799 .map(Path::to_path_buf);
5800 set.spawn(async move {
5801 let _permit = sem.acquire().await;
5802 let mut seat = job.seat;
5803 let out = agent::invoke(
5804 &job.spec,
5805 &mut seat,
5806 &Invocation {
5807 cwd: &job.cwd,
5808 prompt: &job.prompt,
5809 timeout: job.timeout,
5810 allow_write: job.allow_write,
5811 sessions: job.sessions,
5812 artifacts: &job.artifacts,
5813 stem: &job.stem,
5814 run: &run,
5815 node: &node,
5816 cache_dir: cache.as_deref(),
5817 attachments: &[],
5818 },
5819 )
5820 .await;
5821 let out = match out {
5822 Ok(o) if o.usable() => AgentOutcome::Ok(o),
5823 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
5824 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
5832 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
5833 Ok(o) => AgentOutcome::Failed(format!(
5834 "exited with {:?} and no usable output",
5835 o.exit_code
5836 )),
5837 Err(e) => AgentOutcome::Failed(e.to_string()),
5838 };
5839 (i, seat, out)
5840 });
5841 }
5842 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
5843 while let Some(joined) = set.join_next().await {
5844 let (i, seat, out) = match joined {
5845 Ok(v) => v,
5846 Err(e) => {
5850 tracing::error!("agent task panicked: {e}");
5851 continue;
5852 }
5853 };
5854 state.seat_finished(&seat.key);
5855 record_jobs(state, node, round, &seat.key, &out);
5856 if let Err(e) = state.save() {
5857 tracing::warn!("could not persist a seat's completion: {e:#}");
5858 }
5859 if collected.len() <= i {
5860 collected.resize_with(i + 1, || None);
5861 }
5862 collected[i] = Some((i, seat, out));
5863 }
5864 if state
5870 .active
5871 .values()
5872 .any(|a| a.node == node && a.attempt == attempt)
5873 {
5874 state
5875 .active
5876 .retain(|_, a| !(a.node == node && a.attempt == attempt));
5877 if let Err(e) = state.save() {
5878 tracing::warn!("could not persist the end of a wave: {e:#}");
5879 }
5880 }
5881 if let Some(cache_dir) = cache
5888 && jobs_had_a_writer
5889 {
5890 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
5891 }
5892 if let Some(guard) = cache_guard {
5893 guard.release();
5894 }
5895 collected.into_iter().flatten().collect()
5896}
5897
5898fn record_jobs(
5909 state: &mut RunState,
5910 node: &str,
5911 round: Option<usize>,
5912 seat: &str,
5913 out: &AgentOutcome,
5914) {
5915 let commands: &[agent::CommandEvidence] = match out {
5916 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
5917 AgentOutcome::Failed(_) => &[],
5918 };
5919 let checked_at = Timestamp::now();
5920 for c in commands {
5921 state.jobs.push(JobRecord {
5922 node: node.to_owned(),
5923 round,
5924 seat: seat.to_owned(),
5925 id: c.id.clone(),
5926 description: c.description.clone(),
5927 checked_at,
5928 status: match c.exit_code {
5929 Some(0) => JobStatus::Completed,
5930 Some(_) => JobStatus::Failed,
5931 None => JobStatus::Unknown,
5932 },
5933 exit_code: c.exit_code,
5934 result_summary: c.result_summary.clone(),
5935 source: c.source.clone(),
5936 });
5937 }
5938}
5939
5940fn round_is_clean(
5961 blocking: usize,
5962 e2e_ok: bool,
5963 answered: usize,
5964 expected: usize,
5965 quota_missing: usize,
5966 policy: IncompleteReviewPolicy,
5967) -> bool {
5968 if blocking != 0 || !e2e_ok {
5969 return false;
5970 }
5971 if answered == expected || policy == IncompleteReviewPolicy::Warn {
5972 return true;
5973 }
5974 answered > 0 && expected - answered <= quota_missing
5975}
5976
5977fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
6001 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
6002 return Some(RunStatus::Gating);
6003 }
6004 let last = reviews.last()?;
6005 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
6006 if reviews.len() < max_rounds && !stagnant {
6007 return None;
6008 }
6009 if last.incomplete() && last.blocking == 0 {
6010 return Some(RunStatus::Blocked);
6011 }
6012 if last.e2e_status() == E2eStatus::ResourceBlocked {
6013 return None;
6014 }
6015 Some(if last.e2e.iter().all(CommandOutcome::ok) {
6016 RunStatus::Gating
6017 } else {
6018 RunStatus::Blocked
6019 })
6020}
6021
6022fn retry_budget(full: Duration, nudged: bool) -> Duration {
6037 if nudged {
6038 (full / 4).max(Duration::from_secs(120)).min(full)
6039 } else {
6040 full
6041 }
6042}
6043
6044#[allow(clippy::too_many_arguments)]
6057async fn ask_json_wave<T>(
6058 jobs: Vec<SeatJob>,
6059 sem: Arc<Semaphore>,
6060 retries: usize,
6061 ctx: &WaveCtx<'_>,
6062 losses: &mut Vec<QuotaLoss>,
6063 state: &mut RunState,
6064 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
6065) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
6066where
6067 T: serde::de::DeserializeOwned + Send + 'static,
6068{
6069 let n = jobs.len();
6070 let originals: Vec<SeatJob> = jobs;
6071 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
6072 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
6073 let mut attempts_used: Vec<usize> = vec![0; n];
6080 let mut pending: Vec<usize> = (0..n).collect();
6081
6082 for attempt in 0..=retries {
6083 if pending.is_empty() {
6084 break;
6085 }
6086 let mut batch = Vec::with_capacity(pending.len());
6087 for &i in &pending {
6088 let src = &originals[i];
6089 let (prompt, timeout) = if attempt == 0 {
6092 (src.prompt.clone(), src.timeout)
6093 } else {
6094 let why = done[i]
6095 .as_ref()
6096 .and_then(|r| r.as_ref().err().map(ToString::to_string))
6097 .unwrap_or_else(|| "no parsable answer".to_owned());
6098 let nudge = prompt::nudge(&why);
6099 let nudged = has_context(&src.spec, &seats[i], src.sessions);
6100 let prompt = if nudged {
6101 nudge
6102 } else {
6103 format!("{}\n\n---\n\n{}", src.prompt, nudge)
6104 };
6105 (prompt, retry_budget(src.timeout, nudged))
6106 };
6107 batch.push(SeatJob {
6108 spec: src.spec.clone(),
6109 seat: seats[i].clone(),
6110 cwd: src.cwd.clone(),
6111 prompt,
6112 timeout,
6113 allow_write: src.allow_write,
6114 sessions: src.sessions,
6115 artifacts: src.artifacts.clone(),
6116 stem: if attempt == 0 {
6117 src.stem.clone()
6118 } else {
6119 format!("{}-retry{attempt}", src.stem)
6120 },
6121 });
6122 }
6123
6124 if attempt > 0 {
6125 let seats_out: Vec<&str> = pending
6126 .iter()
6127 .map(|&i| originals[i].seat.key.as_str())
6128 .collect();
6129 state.event(
6130 ctx.node,
6131 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
6132 );
6133 }
6134 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
6135 let mut still = Vec::new();
6136 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
6137 seats[i] = seat;
6138 let (parsed, quota) = match out {
6139 AgentOutcome::Ok(o) => (
6140 match verdict::extract_json::<T>(&o.text) {
6141 Ok(v) => match validate(&v) {
6142 Ok(()) => Ok((v, o)),
6143 Err(e) => Err(e),
6144 },
6145 Err(e) => Err(e),
6146 },
6147 false,
6148 ),
6149 AgentOutcome::Quota(o) => {
6150 losses.push(QuotaLoss {
6151 seat: originals[i].seat.key.clone(),
6152 node: ctx.node.to_owned(),
6153 at: Timestamp::now(),
6154 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
6155 });
6156 (
6157 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
6158 true,
6159 )
6160 }
6161 AgentOutcome::Dropped(o) => {
6166 let why = o
6167 .dropped
6168 .as_ref()
6169 .map(|d| d.why.as_str())
6170 .unwrap_or("the CLI ended the stream without delivering its answer");
6171 (
6172 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
6173 false,
6174 )
6175 }
6176 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
6177 };
6178 let failed = parsed.is_err();
6179 done[i] = Some(parsed);
6180 attempts_used[i] = attempt;
6181 if failed && !quota {
6184 still.push(i);
6185 }
6186 }
6187 pending = still;
6188 }
6189
6190 seats
6191 .into_iter()
6192 .zip(done)
6193 .zip(attempts_used)
6194 .map(|((seat, res), attempts)| {
6195 (
6196 seat,
6197 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6198 attempts,
6199 )
6200 })
6201 .collect()
6202}
6203
6204async fn acquire_cache_lease(
6217 state: &mut RunState,
6218 cache_dir: &Path,
6219 owner: &crate::cache::Owner,
6220 budget: Duration,
6221 context: &str,
6222) -> Result<crate::cache::Guard> {
6223 let home = crate::run::home();
6224 let started = Instant::now();
6225 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6226 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6227 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6228 Err(e) => {
6229 state.event(
6230 "verify",
6231 format!("{context}: could not check the shared build cache: {e:#}"),
6232 );
6233 if let Err(e2) = state.save() {
6234 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6235 }
6236 return Err(e);
6237 }
6238 };
6239 state.event(
6240 "verify",
6241 format!(
6242 "{context}: waiting for the shared build cache at {} ({})",
6243 cache_dir.display(),
6244 busy.describe()
6245 ),
6246 );
6247 if let Err(e) = state.save() {
6248 tracing::warn!("could not persist a cache wait: {e:#}");
6249 }
6250 let remaining = budget.saturating_sub(started.elapsed());
6251 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6252 Ok(g) => Ok(g),
6253 Err(e) => {
6254 state.event("verify", format!("{context}: {e:#}"));
6255 if let Err(e2) = state.save() {
6256 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6257 }
6258 Err(e)
6259 }
6260 }
6261}
6262
6263#[allow(clippy::too_many_arguments)]
6284async fn with_cache_lease<'s, F, Fut>(
6285 state: &'s mut RunState,
6286 cache_dir: Option<&Path>,
6287 node: &str,
6288 seat: &str,
6289 worktree: &Path,
6290 head: &str,
6291 budget: Duration,
6292 context: &str,
6293 body: F,
6294) -> (Vec<CommandOutcome>, bool)
6295where
6296 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6297 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6298{
6299 let Some(cache_dir) = cache_dir else {
6300 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6301 return (outcomes, retried);
6302 };
6303 let home = crate::run::home();
6304 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6305 let started = Instant::now();
6306 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6307 Ok(g) => g,
6308 Err(e) => {
6309 return (
6310 vec![CommandOutcome {
6311 command: "(waiting for the shared build cache)".to_owned(),
6312 code: None,
6313 output_tail: e.to_string(),
6314 duration_ms: started.elapsed().as_millis() as u64,
6315 resource_blocked: true,
6316 }],
6317 false,
6318 );
6319 }
6320 };
6321 let identity = crate::cache::Identity::new(worktree, head);
6322 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6323 state.event(
6331 "verify",
6332 format!(
6333 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6334 worktree.display(),
6335 short(head)
6336 ),
6337 );
6338 guard.release();
6339 return (
6340 vec![CommandOutcome {
6341 command: "(confirming the shared build cache is fresh)".to_owned(),
6342 code: None,
6343 output_tail: e.to_string(),
6344 duration_ms: started.elapsed().as_millis() as u64,
6345 resource_blocked: true,
6346 }],
6347 false,
6348 );
6349 }
6350 let remaining = budget.saturating_sub(started.elapsed());
6351 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6352 if !timed_out_pids.is_empty() {
6357 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6358 }
6359 guard.release();
6360 (outcomes, retried)
6361}
6362
6363async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6375 wait_for_pids_with(
6376 pids,
6377 crate::proc::pid_alive,
6378 LEASE_RELEASE_POLL,
6379 LEASE_RELEASE_MAX_WAIT,
6380 )
6381 .await;
6382}
6383
6384async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6390 pids: &[u32],
6391 alive: F,
6392 poll: Duration,
6393 max_wait: Duration,
6394) {
6395 let deadline = Instant::now() + max_wait;
6396 loop {
6397 if pids.iter().all(|&pid| !alive(pid)) {
6398 return;
6399 }
6400 if Instant::now() >= deadline {
6401 return;
6402 }
6403 tokio::time::sleep(poll).await;
6404 }
6405}
6406
6407fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6414 outcomes.iter().any(|o| o.resource_blocked)
6415}
6416
6417enum GateFix {
6419 Retry,
6421 Stop,
6424 Defer,
6427}
6428
6429fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6437 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6438 red.peek().is_some()
6439 && red.all(|o| {
6440 !o.resource_blocked
6441 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6442 && !o.output_tail.trim().is_empty()
6443 })
6444}
6445
6446fn e2e_outcome_label(o: &CommandOutcome) -> String {
6450 if o.ok() {
6451 return "pass".to_owned();
6452 }
6453 let reason = if o.build_failed() {
6454 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6455 } else {
6456 format!("FAIL ({:?})", o.code)
6457 };
6458 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6459}
6460
6461async fn run_e2e_with_retry(
6469 state: &mut RunState,
6470 shell: &[String],
6471 commands: &[String],
6472 worktree: &Path,
6473 timeout: Duration,
6474 context: &str,
6475) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6476 let (mut e2e, mut timed_out_pids) = run_commands(
6477 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6478 )
6479 .await;
6480 for o in &e2e {
6481 state.event(
6482 "verify",
6483 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6484 );
6485 }
6486 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6490 if verify_retried {
6491 state.event(
6492 "verify",
6493 format!(
6494 "{context}: verify could not build/link, not a test result — retrying once \
6495 before concluding"
6496 ),
6497 );
6498 let retried = run_commands(
6499 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6500 )
6501 .await;
6502 e2e = retried.0;
6503 timed_out_pids.extend(retried.1);
6506 for o in &e2e {
6507 state.event(
6508 "verify",
6509 format!(
6510 "{context}: retry `{}` -> {}",
6511 o.command,
6512 e2e_outcome_label(o)
6513 ),
6514 );
6515 }
6516 }
6517 (e2e, verify_retried, timed_out_pids)
6518}
6519
6520#[allow(clippy::too_many_arguments)]
6536async fn run_commands(
6537 state: &mut RunState,
6538 node: &str,
6539 task: &str,
6540 attempt: usize,
6541 shell: &[String],
6542 commands: &[String],
6543 cwd: &Path,
6544 timeout: Duration,
6545) -> (Vec<CommandOutcome>, Vec<u32>) {
6546 if commands.is_empty() {
6547 return (Vec::new(), Vec::new());
6552 }
6553 let mut out = Vec::new();
6554 let mut timed_out_pids = Vec::new();
6555 let total = commands.len();
6556 for (idx, command) in commands.iter().enumerate() {
6557 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6558 if let Err(e) = state.save() {
6559 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6560 }
6561 let started = Instant::now();
6562 let mut cmd = tokio::process::Command::new(&shell[0]);
6563 cmd.quiet();
6564 cmd.args(&shell[1..])
6565 .arg(command)
6566 .current_dir(cwd)
6567 .stdin(std::process::Stdio::null())
6568 .stdout(std::process::Stdio::piped())
6569 .stderr(std::process::Stdio::piped())
6570 .kill_on_drop(true);
6571 let spawned = cmd.spawn();
6572 let (code, body) = match spawned {
6573 Ok(child) => {
6574 let pid = child.id();
6579 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6580 Ok(Ok(o)) => {
6581 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6582 body.push_str(&String::from_utf8_lossy(&o.stderr));
6583 (o.status.code(), body)
6584 }
6585 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6586 Err(_) => {
6587 if let Some(pid) = pid {
6588 timed_out_pids.push(pid);
6589 }
6590 (None, format!("timed out after {}s", timeout.as_secs()))
6591 }
6592 }
6593 }
6594 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6595 };
6596 out.push(CommandOutcome {
6597 command: command.clone(),
6598 code,
6599 output_tail: tail(&body, OUTPUT_TAIL),
6600 duration_ms: started.elapsed().as_millis() as u64,
6601 resource_blocked: false,
6602 });
6603 }
6604 state.task_finished(task);
6605 if let Err(e) = state.save() {
6606 tracing::warn!("could not persist the end of {task}: {e:#}");
6607 }
6608 (out, timed_out_pids)
6609}
6610
6611fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6623 let repo = repo.display();
6624 match style {
6625 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6626 MergeStyle::Squash => {
6627 let subject = message
6630 .lines()
6631 .next()
6632 .unwrap_or(branch)
6633 .replace(['\\', '"', '$', '`'], "");
6634 format!(
6635 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6636 )
6637 }
6638 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6639 }
6640}
6641
6642const PR_TITLE_MAX: usize = 240;
6653
6654struct PrMessage {
6659 title: String,
6660 body: String,
6661}
6662
6663impl PrMessage {
6664 fn commit_message(&self) -> String {
6668 format!("{}\n\n{}", self.title, self.body)
6669 }
6670}
6671
6672fn title_marker(line: &str) -> Option<&str> {
6674 let line = line.trim();
6675 let head = line.get(..6)?;
6676 head.eq_ignore_ascii_case("title:")
6677 .then(|| line[6..].trim())
6678}
6679
6680fn summary_title(summary: &str) -> Option<String> {
6685 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6686 let raw = title_marker(first)?;
6687 if raw.is_empty() {
6688 return None;
6689 }
6690 let title = queue::title_from(raw, PR_TITLE_MAX);
6691 let lower = title.to_ascii_lowercase();
6692 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6693 return None;
6694 }
6695 Some(title)
6696}
6697
6698fn summary_without_title(summary: &str) -> String {
6701 let mut lines = summary.trim().lines().peekable();
6702 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6703 lines.next();
6704 }
6705 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6706}
6707
6708fn pr_message(state: &RunState, winner: char) -> PrMessage {
6722 let summary = state
6723 .candidates
6724 .iter()
6725 .find(|c| c.label == winner)
6726 .map(|c| c.summary.as_str())
6727 .unwrap_or_default();
6728 let title = summary_title(summary).unwrap_or_else(|| {
6731 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6732 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6733 t
6734 } else {
6735 format!(
6736 "chore: land candidate {} of run {}",
6737 winner.to_ascii_uppercase(),
6738 state.id
6739 )
6740 }
6741 });
6742
6743 let mut body = String::new();
6744 let what = summary_without_title(summary);
6745 if !what.is_empty() {
6746 body.push_str("## Summary\n\n");
6747 body.push_str(&what);
6748 body.push_str("\n\n");
6749 }
6750
6751 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6752 if let Some(fix) = fix
6753 && !fix.notes.trim().is_empty()
6754 {
6755 body.push_str("## Review fixes\n\n");
6756 body.push_str(fix.notes.trim());
6757 body.push_str("\n\n");
6758 }
6759
6760 let open = state.open_findings();
6761 if !open.is_empty() {
6762 body.push_str("## Open review findings\n\n");
6763 for f in &open {
6764 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6765 }
6766 body.push('\n');
6767 }
6768
6769 if let Some(fix) = fix
6770 && !fix.rejected.is_empty()
6771 {
6772 body.push_str("## Declined by the fixer\n\n");
6773 for r in &fix.rejected {
6774 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
6775 }
6776 body.push('\n');
6777 }
6778
6779 let task = state.instruction.trim();
6780 let task = if task.is_empty() {
6781 "(empty task)"
6782 } else {
6783 task
6784 };
6785 body.push_str(&format!(
6786 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
6787 task.replace("</details>", "</details>")
6788 ));
6789
6790 body.push_str(&format!(
6791 "\n---\nmagi:run/{} magi:candidate-{}\n",
6792 state.id,
6793 winner.to_ascii_lowercase()
6794 ));
6795
6796 let id = crate::scrub::Identity::current();
6799 PrMessage {
6800 title: crate::scrub::scrub(&title, &id),
6801 body: crate::scrub::scrub(&body, &id),
6802 }
6803}
6804
6805async fn gh_pr_create(
6807 cwd: &Path,
6808 base: &str,
6809 head: &str,
6810 title: &str,
6811 body: &str,
6812) -> Result<String> {
6813 let out = tokio::process::Command::new("gh")
6814 .args([
6815 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
6816 ])
6817 .current_dir(cwd)
6818 .quiet()
6819 .stdin(std::process::Stdio::null())
6820 .output()
6821 .await
6822 .context("spawn gh")?;
6823 if out.status.success() {
6824 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
6825 } else {
6826 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
6827 }
6828}
6829
6830pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
6839 let repo = state.repo.clone();
6840 let root = state.worktree_root();
6841 let winner = state.tally.as_ref().map(|t| t.winner);
6842 let mut removed = Vec::new();
6843
6844 for i in 0..state.candidates.len() {
6845 let c = state.candidates[i].clone();
6846 let is_winner = Some(c.label) == winner;
6847 if is_winner && !drop_winner {
6848 continue;
6849 }
6850 if c.worktree.exists() {
6851 git::worktree_remove(&repo, &c.worktree).await.ok();
6852 removed.push(c.worktree.to_string_lossy().into_owned());
6853 }
6854 let handed_over = state.released_branches.contains(&c.branch);
6857 if !handed_over && git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
6858 git::branch_delete(&repo, &c.branch).await.ok();
6859 removed.push(c.branch.clone());
6860 }
6861 state.candidates[i].folded = true;
6862 }
6863
6864 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
6865 let path = name.path();
6866 let keep = !drop_winner
6867 && winner.is_some_and(|w| {
6868 path.file_name()
6869 .is_some_and(|n| n == format!("cand-{w}").as_str())
6870 });
6871 if keep {
6872 continue;
6873 }
6874 git::worktree_remove(&repo, &path).await.ok();
6875 removed.push(path.to_string_lossy().into_owned());
6876 }
6877
6878 remove_if_empty(&root);
6887
6888 if state.enabled_worktree_config && drop_winner {
6889 git::release_worktree_config(&repo).await.ok();
6893 state.enabled_worktree_config = false;
6894 }
6895 state.save_under(home)?;
6896 Ok(removed)
6897}
6898
6899fn remove_if_empty(dir: &Path) {
6910 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
6911 std::fs::remove_dir(dir).ok();
6912 }
6913}
6914
6915pub fn worst_open(state: &RunState) -> Option<Severity> {
6917 state
6918 .reviews
6919 .last()?
6920 .reviews
6921 .iter()
6922 .flat_map(|r| r.findings.iter())
6923 .map(|f| f.severity)
6924 .max()
6925}
6926
6927#[cfg(test)]
6928mod tests {
6929 use super::*;
6930 use crate::run::GateStatus;
6931 use std::collections::BTreeMap;
6932 use std::time::Duration;
6933
6934 fn conductor() -> AgentSpec {
6935 AgentSpec {
6936 id: "conductor".to_owned(),
6937 kind: crate::config::AgentKind::Command,
6938 model: None,
6939 command: vec!["true".to_owned()],
6940 extra_args: Vec::new(),
6941 env: BTreeMap::new(),
6942 prompt_delivery: None,
6943 }
6944 }
6945
6946 fn spec(id: &str) -> AgentSpec {
6947 AgentSpec {
6948 id: id.to_owned(),
6949 kind: crate::config::AgentKind::Command,
6950 model: None,
6951 command: vec!["true".to_owned()],
6952 extra_args: Vec::new(),
6953 env: BTreeMap::new(),
6954 prompt_delivery: None,
6955 }
6956 }
6957
6958 #[test]
6965 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
6966 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6967 let tried = BTreeSet::from(["beta".to_owned()]);
6968 let next = next_untried_implementer(&roster, 1, &tried);
6971 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
6972 }
6973
6974 #[test]
6975 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
6976 let roster = vec![spec("alpha"), spec("beta")];
6977 let tried = BTreeSet::from(["beta".to_owned()]);
6978 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6982 }
6983
6984 #[test]
6985 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
6986 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
6987 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
6988 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
6992 }
6993
6994 #[test]
6995 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
6996 let roster = vec![spec("a"), spec("a"), spec("b")];
6997 let tried = BTreeSet::from(["a".to_owned()]);
6998 let next = next_untried_implementer(&roster, 0, &tried);
6999 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
7000 }
7001
7002 #[test]
7003 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
7004 let roster = vec![spec("a"), spec("b")];
7005 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
7006 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
7007 }
7008
7009 #[test]
7010 fn remove_if_empty_only_ever_takes_a_bare_directory() {
7011 let dir = tempfile::tempdir().unwrap();
7012 let bay = dir.path().join("ffff");
7013
7014 remove_if_empty(&bay);
7016 assert!(!bay.exists());
7017
7018 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
7021 remove_if_empty(&bay);
7022 assert!(bay.exists(), "non-empty directory must survive");
7023
7024 std::fs::remove_dir(bay.join("cand-A")).unwrap();
7026 remove_if_empty(&bay);
7027 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
7028 }
7029
7030 #[test]
7039 fn a_full_panel_that_found_nothing_is_clean() {
7040 assert!(round_is_clean(
7041 0,
7042 true,
7043 2,
7044 2,
7045 0,
7046 IncompleteReviewPolicy::Block
7047 ));
7048 }
7049
7050 #[test]
7051 fn a_missing_seat_is_never_clean_under_the_default_policy() {
7052 assert!(!round_is_clean(
7053 0,
7054 true,
7055 1,
7056 2,
7057 0,
7058 IncompleteReviewPolicy::Block
7059 ));
7060 }
7061
7062 #[test]
7063 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
7064 assert!(!round_is_clean(
7065 1,
7066 true,
7067 1,
7068 2,
7069 0,
7070 IncompleteReviewPolicy::Warn
7071 ));
7072 }
7073
7074 #[test]
7075 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
7076 assert!(round_is_clean(
7077 0,
7078 true,
7079 1,
7080 2,
7081 0,
7082 IncompleteReviewPolicy::Warn
7083 ));
7084 }
7085
7086 #[test]
7087 fn a_full_panel_with_an_open_finding_is_not_clean() {
7088 assert!(!round_is_clean(
7089 1,
7090 true,
7091 2,
7092 2,
7093 0,
7094 IncompleteReviewPolicy::Block
7095 ));
7096 }
7097
7098 #[test]
7099 fn a_full_panel_with_a_red_e2e_is_not_clean() {
7100 assert!(!round_is_clean(
7101 0,
7102 false,
7103 2,
7104 2,
7105 0,
7106 IncompleteReviewPolicy::Block
7107 ));
7108 }
7109
7110 #[test]
7117 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
7118 assert!(round_is_clean(
7121 0,
7122 true,
7123 1,
7124 2,
7125 1,
7126 IncompleteReviewPolicy::Block
7127 ));
7128 }
7129
7130 #[test]
7131 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
7132 assert!(!round_is_clean(
7135 0,
7136 true,
7137 1,
7138 2,
7139 0,
7140 IncompleteReviewPolicy::Block
7141 ));
7142 }
7143
7144 #[test]
7145 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
7146 assert!(!round_is_clean(
7147 1,
7148 true,
7149 1,
7150 2,
7151 1,
7152 IncompleteReviewPolicy::Block
7153 ));
7154 assert!(!round_is_clean(
7155 0,
7156 false,
7157 1,
7158 2,
7159 1,
7160 IncompleteReviewPolicy::Block
7161 ));
7162 }
7163
7164 #[test]
7165 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
7166 assert!(!round_is_clean(
7170 0,
7171 true,
7172 0,
7173 2,
7174 2,
7175 IncompleteReviewPolicy::Block
7176 ));
7177 }
7178
7179 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
7180 CommandOutcome {
7181 command: "test".to_owned(),
7182 code,
7183 output_tail: String::new(),
7184 duration_ms: 0,
7185 resource_blocked,
7186 }
7187 }
7188
7189 #[test]
7190 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7191 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7192 assert!(
7193 !verify_inconclusive(&[outcome(Some(1), false)]),
7194 "an ordinary failure is still evidence about the patch"
7195 );
7196 assert!(verify_inconclusive(&[outcome(None, true)]));
7197 assert!(
7198 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7199 "one inconclusive outcome taints the whole batch"
7200 );
7201 assert!(!verify_inconclusive(&[]));
7202 }
7203
7204 #[tokio::test]
7205 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7206 let calls = std::sync::atomic::AtomicUsize::new(0);
7210 let started = Instant::now();
7211 wait_for_pids_with(
7212 &[123],
7213 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7214 Duration::from_millis(5),
7215 Duration::from_secs(5),
7216 )
7217 .await;
7218 assert!(
7219 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7220 "must keep checking rather than deciding on the first answer"
7221 );
7222 assert!(
7223 started.elapsed() < Duration::from_secs(1),
7224 "must return the moment it is confirmed dead, not wait out the ceiling"
7225 );
7226 }
7227
7228 #[tokio::test]
7229 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7230 let started = Instant::now();
7231 wait_for_pids_with(
7232 &[123],
7233 |_| true, Duration::from_millis(5),
7235 Duration::from_millis(30),
7236 )
7237 .await;
7238 let elapsed = started.elapsed();
7239 assert!(
7240 elapsed >= Duration::from_millis(30),
7241 "must not give up before its own ceiling: {elapsed:?}"
7242 );
7243 assert!(
7244 elapsed < Duration::from_secs(1),
7245 "must not wait past its own ceiling either: {elapsed:?}"
7246 );
7247 }
7248
7249 #[tokio::test]
7250 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7251 let started = Instant::now();
7252 wait_for_pids_with(
7253 &[],
7254 |_| true,
7255 Duration::from_secs(5),
7256 Duration::from_secs(5),
7257 )
7258 .await;
7259 assert!(
7260 started.elapsed() < Duration::from_millis(200),
7261 "an empty pid list has nothing to confirm"
7262 );
7263 }
7264
7265 fn review_round(
7271 clean: bool,
7272 blocking: usize,
7273 answered: usize,
7274 expected: usize,
7275 progressed: bool,
7276 e2e_ok: bool,
7277 ) -> ReviewRound {
7278 ReviewRound {
7279 round: 1,
7280 head: "h".to_owned(),
7281 verified_head: None,
7282 verified_at: None,
7283 reviews: Vec::new(),
7284 e2e: vec![CommandOutcome {
7285 command: "test".to_owned(),
7286 code: Some(if e2e_ok { 0 } else { 1 }),
7287 output_tail: String::new(),
7288 duration_ms: 0,
7289 resource_blocked: false,
7290 }],
7291 verify_retried: false,
7292 e2e_deferred: false,
7293 e2e_defer_reason: None,
7294 fix: None,
7295 blocking,
7296 answered,
7297 expected,
7298 clean,
7299 progressed,
7300 vote_split: false,
7301 reconsideration: Vec::new(),
7302 verdict: None,
7303 }
7304 }
7305
7306 #[test]
7307 fn review_conclusion_is_none_when_nothing_has_run() {
7308 assert_eq!(review_conclusion(&[], 3), None);
7309 }
7310
7311 #[test]
7312 fn review_conclusion_is_none_while_rounds_remain() {
7313 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7314 assert_eq!(review_conclusion(&rounds, 3), None);
7315 }
7316
7317 #[test]
7318 fn review_conclusion_is_gating_once_a_round_is_clean() {
7319 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7320 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7321 }
7322
7323 #[test]
7324 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7325 let rounds = vec![
7326 review_round(false, 1, 2, 2, true, true),
7327 review_round(false, 1, 2, 2, true, true),
7328 ];
7329 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7330 }
7331
7332 #[test]
7333 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7334 let rounds = vec![
7335 review_round(false, 1, 2, 2, true, true),
7336 review_round(false, 1, 2, 2, true, false),
7337 ];
7338 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7339 }
7340
7341 #[test]
7342 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7343 let mut blocked = review_round(false, 1, 2, 2, true, false);
7350 blocked.e2e[0].resource_blocked = true;
7351 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7352 assert_eq!(review_conclusion(&rounds, 2), None);
7353 }
7354
7355 #[test]
7356 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7357 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7359 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7360 }
7361
7362 #[test]
7363 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7364 let rounds = vec![
7365 review_round(false, 1, 2, 2, false, true),
7366 review_round(false, 1, 2, 2, false, true),
7367 ];
7368 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7369 }
7370
7371 fn secs(n: u64) -> Duration {
7372 Duration::from_secs(n)
7373 }
7374
7375 fn init_repo(dir: &Path) {
7378 let run = |args: &[&str]| {
7379 let out = std::process::Command::new("git")
7380 .args(args)
7381 .current_dir(dir)
7382 .quiet()
7383 .output()
7384 .expect("spawn git");
7385 assert!(
7386 out.status.success(),
7387 "git {args:?} failed: {}",
7388 String::from_utf8_lossy(&out.stderr)
7389 );
7390 };
7391 run(&["init", "-b", "main"]);
7392 run(&["config", "user.name", "magi test"]);
7393 run(&["config", "user.email", "magi@example.com"]);
7394 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7395 run(&["add", "-A"]);
7396 run(&["commit", "-m", "init"]);
7397 }
7398
7399 fn ask_test_home() {
7407 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7408 }
7409
7410 fn runner_at(status: RunStatus) -> Runner {
7413 let mut state = RunState::new(
7414 PathBuf::from("/nonexistent/repo"),
7415 "main".to_owned(),
7416 "deadbeef".to_owned(),
7417 "task".to_owned(),
7418 Config::default(),
7419 );
7420 state.status = status;
7421 Runner {
7422 state,
7423 roles: ResolvedRoles {
7424 implementers: Vec::new(),
7425 judges: Vec::new(),
7426 reviewers: Vec::new(),
7427 fixer: None,
7428 conductor: conductor(),
7429 implementer_roster: Vec::new(),
7430 },
7431 sem: Arc::new(Semaphore::new(1)),
7432 pause: Pause::new(),
7433 interrupt: Pause::new(),
7434 }
7435 }
7436
7437 #[test]
7441 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7442 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7443 let mut runner = runner_at(RunStatus::Implementing);
7444 let interrupt = Pause::new();
7445 runner.watch_interrupt(interrupt.clone());
7446
7447 interrupt.park_because("task a1b2 asked to run first");
7448
7449 assert!(runner.park_here().expect("park_here"));
7450 assert!(runner.state.parked);
7451 let last = runner.state.events.last().expect("a park event");
7452 assert_eq!(last.node, "park");
7453 assert!(
7454 last.message.contains("task a1b2 asked to run first"),
7455 "expected the interrupt reason in {:?}",
7456 last.message
7457 );
7458 }
7459
7460 #[test]
7468 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7469 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7470 let mut runner = runner_at(RunStatus::Implementing);
7471 let shutdown = Pause::new();
7472 runner.on_pause(shutdown.clone());
7473 let interrupt = Pause::new();
7474 runner.watch_interrupt(interrupt.clone());
7475
7476 assert!(!runner.park_here().expect("park_here"));
7478 assert!(!runner.state.parked);
7479
7480 interrupt.park_because("test");
7482 assert!(!shutdown.parked());
7483 assert!(runner.park_here().expect("park_here"));
7484 }
7485
7486 #[tokio::test]
7500 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7501 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7502 let mut runner = runner_at(RunStatus::Implementing);
7503 let interrupt = Pause::new();
7504 runner.watch_interrupt(interrupt.clone());
7505
7506 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7507 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7508
7509 let node = async move {
7513 started_tx.send(()).expect("send started");
7514 finish_rx.await.expect("recv finish");
7515 "node finished"
7516 };
7517
7518 let interrupter = async move {
7519 started_rx.await.expect("recv started");
7520 interrupt.park_because("higher-priority task waiting");
7522 tokio::task::yield_now().await;
7526 finish_tx.send(()).expect("send finish");
7527 };
7528
7529 let (node_result, ()) = tokio::join!(node, interrupter);
7530 assert_eq!(
7531 node_result, "node finished",
7532 "the in-flight call ran to completion"
7533 );
7534
7535 assert!(runner.park_here().expect("park_here"));
7538 assert!(runner.state.parked);
7539 }
7540
7541 #[test]
7547 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7548 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7549 let mut runner = runner_at(RunStatus::Judging);
7550 runner.state.config.agents = vec![conductor()];
7554 runner.state.candidates = vec![Candidate {
7555 index: 0,
7556 label: 'A',
7557 agent: "alpha".to_owned(),
7558 branch: "magi/x/A".to_owned(),
7559 worktree: PathBuf::from("/nonexistent/worktree"),
7560 summary: "did the thing".to_owned(),
7561 stat: "1 file changed".to_owned(),
7562 files: 1,
7563 commits: 1,
7564 empty: false,
7565 failed: None,
7566 verified_noop: None,
7567 duration_ms: 1234,
7568 folded: false,
7569 }];
7570 let run_id = runner.state.id.clone();
7571
7572 let interrupt = Pause::new();
7573 runner.watch_interrupt(interrupt.clone());
7574 interrupt.park_because("task c3d4 asked to run first");
7575 assert!(runner.park_here().expect("park_here"));
7576
7577 let resumed = Runner::resume(&run_id).expect("resume");
7578 assert_eq!(resumed.state.candidates.len(), 1);
7579 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7580 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7581 assert_eq!(resumed.state.status, runner.state.status);
7582 assert!(
7583 resumed.state.parked,
7584 "still parked until `execute` actually walks the graph again"
7585 );
7586 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7587 }
7588
7589 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7591 let mut q = ask::Question::new(
7592 run.to_owned(),
7593 "implement".to_owned(),
7594 "impl-A".to_owned(),
7595 "Which storage backend should the cache use?".to_owned(),
7596 String::new(),
7597 vec!["SQLite".to_owned(), "Redis".to_owned()],
7598 );
7599 store.put(&mut q).unwrap();
7600 q
7601 }
7602
7603 #[test]
7604 fn a_failed_runs_open_question_is_abandoned() {
7605 ask_test_home();
7606 let store = ask::Questions::open();
7607 let mut runner = runner_at(RunStatus::Failed);
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 "the seat that asked died with the run; nobody is left to read an answer"
7617 );
7618 assert!(
7619 back.detail.contains(&run) && back.detail.contains("failed"),
7620 "the reason names what the run became, not just that it is gone: {}",
7621 back.detail
7622 );
7623 }
7624
7625 #[test]
7626 fn a_merged_runs_open_question_is_abandoned_too() {
7627 ask_test_home();
7628 let store = ask::Questions::open();
7629 for status in [RunStatus::Merged, RunStatus::Ready] {
7632 let mut runner = runner_at(status);
7633 let run = runner.state.id.clone();
7634 let q = ask_open_question(&store, &run);
7635
7636 runner.settle_questions();
7637
7638 let back = store.get(&q.id).unwrap();
7639 assert!(
7640 !back.status.open(),
7641 "{status:?} run's question must not outlive the run"
7642 );
7643 }
7644 }
7645
7646 #[test]
7647 fn a_still_resumable_runs_open_question_is_left_alone() {
7648 ask_test_home();
7649 let store = ask::Questions::open();
7650 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7656 let mut runner = runner_at(status);
7657 let run = runner.state.id.clone();
7658 let q = ask_open_question(&store, &run);
7659
7660 runner.settle_questions();
7661
7662 let back = store.get(&q.id).unwrap();
7663 assert!(
7664 back.status.open(),
7665 "{status:?} is still alive; the question must still be waiting"
7666 );
7667 }
7668 }
7669
7670 #[test]
7671 fn settle_questions_never_touches_an_already_answered_question() {
7672 ask_test_home();
7673 let store = ask::Questions::open();
7674 let mut runner = runner_at(RunStatus::Failed);
7675 let run = runner.state.id.clone();
7676 let mut q = ask_open_question(&store, &run);
7677 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
7678 .unwrap();
7679 store.put(&mut q).unwrap();
7680
7681 runner.settle_questions();
7686 runner.settle_questions();
7687
7688 let back = store.get(&q.id).unwrap();
7689 assert_eq!(
7690 back.status,
7691 ask::QuestionStatus::Answered,
7692 "a real answer is a decision on record, never overwritten by a sweep"
7693 );
7694 }
7695
7696 #[tokio::test]
7707 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
7708 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
7709 let tmp = tempfile::tempdir().expect("tempdir");
7710 let repo = tmp.path().join("repo");
7711 std::fs::create_dir_all(&repo).unwrap();
7712 init_repo(&repo);
7713
7714 let mut config = Config::default();
7715 config.graph.worktree_root = Some(tmp.path().join("wt"));
7716
7717 let mut state = RunState::new(
7718 repo.clone(),
7719 "main".to_owned(),
7720 "deadbeef".to_owned(),
7721 "task".to_owned(),
7722 config,
7723 );
7724 let root = state.worktree_root();
7725 let wt_a = root.join("cand-A");
7726 let wt_b = root.join("cand-B");
7727 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
7728 .await
7729 .expect("worktree A");
7730 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
7731 .await
7732 .expect("worktree B");
7733
7734 state.candidates = vec![
7735 Candidate {
7736 index: 0,
7737 label: 'A',
7738 agent: "alpha".to_owned(),
7739 branch: "magi/x/A".to_owned(),
7740 worktree: wt_a.clone(),
7741 summary: String::new(),
7742 stat: String::new(),
7743 files: 0,
7744 commits: 0,
7745 empty: false,
7746 failed: None,
7747 verified_noop: None,
7748 duration_ms: 0,
7749 folded: false,
7750 },
7751 Candidate {
7752 index: 1,
7753 label: 'B',
7754 agent: "beta".to_owned(),
7755 branch: "magi/x/B".to_owned(),
7756 worktree: wt_b.clone(),
7757 summary: String::new(),
7758 stat: String::new(),
7759 files: 0,
7760 commits: 0,
7761 empty: false,
7762 failed: None,
7763 verified_noop: None,
7764 duration_ms: 0,
7765 folded: false,
7766 },
7767 ];
7768 state.tally = Some(Tally {
7769 first_choice: BTreeMap::from([('A', 1)]),
7770 borda: BTreeMap::new(),
7771 winner: 'A',
7772 rankings: 1,
7773 unanimous_initial: true,
7774 deliberated: false,
7775 changed_votes: 0,
7776 unanimous_final: true,
7777 tie_break: None,
7778 judges: 1,
7779 present: 1,
7780 quorum: 1,
7781 met_quorum: true,
7782 uncontested: None,
7783 });
7784 state.status = RunStatus::Ready;
7785
7786 fold_run(&mut state, false, &crate::run::home())
7787 .await
7788 .expect("fold_run");
7789
7790 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
7791 assert!(
7792 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7793 "the unmerged winner's branch survives"
7794 );
7795 assert!(
7796 !state.candidates[0].folded,
7797 "the winner is not marked folded"
7798 );
7799
7800 assert!(!wt_b.exists(), "the loser's worktree is removed");
7801 assert!(
7802 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
7803 "the loser's branch is removed"
7804 );
7805 assert!(state.candidates[1].folded, "the loser is marked folded");
7806 }
7807
7808 #[tokio::test]
7811 async fn fold_run_keeps_a_branch_that_was_handed_to_a_later_run() {
7812 let tmp = tempfile::tempdir().expect("tempdir");
7813 let repo = tmp.path().join("repo");
7814 std::fs::create_dir_all(&repo).unwrap();
7815 init_repo(&repo);
7816 let home = tmp.path().join("home");
7817
7818 let mut config = Config::default();
7819 config.graph.worktree_root = Some(tmp.path().join("wt"));
7820 let mut state = RunState::new(
7821 repo.clone(),
7822 "main".to_owned(),
7823 "deadbeef".to_owned(),
7824 "task".to_owned(),
7825 config,
7826 );
7827 git::git(&repo, &["branch", "magi/x/A", "main"])
7829 .await
7830 .expect("branch");
7831 state.candidates = vec![Candidate {
7832 index: 0,
7833 label: 'A',
7834 agent: "alpha".to_owned(),
7835 branch: "magi/x/A".to_owned(),
7836 worktree: state.worktree_root().join("cand-A"),
7837 summary: String::new(),
7838 stat: String::new(),
7839 files: 0,
7840 commits: 0,
7841 empty: false,
7842 failed: None,
7843 verified_noop: None,
7844 duration_ms: 0,
7845 folded: true,
7846 }];
7847 state.released_to = Some("20260901-000000-new1".to_owned());
7848 state.released_branches = vec!["magi/x/A".to_owned()];
7849
7850 fold_run(&mut state, true, &home).await.expect("fold_run");
7851
7852 assert!(
7853 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7854 "the handed-over branch survives a fold"
7855 );
7856 }
7857
7858 #[tokio::test]
7867 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
7868 let tmp = tempfile::tempdir().expect("tempdir");
7869 let repo = tmp.path().join("repo");
7870 std::fs::create_dir_all(&repo).unwrap();
7871 init_repo(&repo);
7872
7873 let mut config = Config::default();
7874 config.merge.mode = MergeMode::Local;
7875
7876 let mut state = RunState::new(
7877 repo.clone(),
7878 "main".to_owned(),
7879 "deadbeef".to_owned(),
7880 "task".to_owned(),
7881 config,
7882 );
7883 state.candidates = vec![Candidate {
7884 index: 0,
7885 label: 'A',
7886 agent: "alpha".to_owned(),
7887 branch: "does-not-exist".to_owned(),
7888 worktree: repo.clone(),
7889 summary: String::new(),
7890 stat: String::new(),
7891 files: 0,
7892 commits: 0,
7893 empty: false,
7894 failed: None,
7895 verified_noop: None,
7896 duration_ms: 0,
7897 folded: false,
7898 }];
7899 state.tally = Some(Tally {
7900 first_choice: BTreeMap::from([('A', 1)]),
7901 borda: BTreeMap::new(),
7902 winner: 'A',
7903 rankings: 1,
7904 unanimous_initial: true,
7905 deliberated: false,
7906 changed_votes: 0,
7907 unanimous_final: true,
7908 tie_break: None,
7909 judges: 0,
7910 present: 0,
7911 quorum: 0,
7912 met_quorum: true,
7913 uncontested: Some("only candidate A produced a change".to_owned()),
7914 });
7915 state.reviews = vec![ReviewRound {
7916 round: 1,
7917 head: "deadbeef".to_owned(),
7918 verified_head: None,
7919 verified_at: None,
7920 reviews: Vec::new(),
7921 e2e: Vec::new(),
7922 fix: None,
7923 blocking: 0,
7924 answered: 0,
7925 expected: 0,
7926 clean: true,
7927 verify_retried: false,
7928 e2e_deferred: false,
7929 e2e_defer_reason: None,
7930 progressed: false,
7931 vote_split: false,
7932 reconsideration: Vec::new(),
7933 verdict: None,
7934 }];
7935 state.gate = vec![CommandOutcome {
7936 command: "test".to_owned(),
7937 code: Some(0),
7938 output_tail: String::new(),
7939 duration_ms: 0,
7940 resource_blocked: false,
7941 }];
7942 state.gate_ran = true;
7943 state.status = RunStatus::Ready;
7948 state.merge = Some(MergeOutcome {
7949 mode: MergeMode::Local,
7950 ok: false,
7951 detail: "already concluded".to_owned(),
7952 });
7953
7954 let mut runner = Runner {
7955 state,
7956 roles: ResolvedRoles {
7957 implementers: Vec::new(),
7958 judges: Vec::new(),
7959 reviewers: Vec::new(),
7960 fixer: None,
7961 conductor: conductor(),
7962 implementer_roster: Vec::new(),
7963 },
7964 sem: Arc::new(Semaphore::new(1)),
7965 pause: Pause::new(),
7966 interrupt: Pause::new(),
7967 };
7968
7969 runner.merge().await.expect("merge");
7970
7971 assert_eq!(
7972 runner.state.status,
7973 RunStatus::Ready,
7974 "a concluded run's status must not change on reentry"
7975 );
7976 assert_eq!(
7977 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
7978 Some("already concluded"),
7979 "merge must not run again once the node already recorded an outcome"
7980 );
7981 }
7982
7983 #[tokio::test]
7992 async fn merge_refuses_a_gate_that_has_not_actually_run() {
7993 let tmp = tempfile::tempdir().expect("tempdir");
7994 let repo = tmp.path().join("repo");
7995 std::fs::create_dir_all(&repo).unwrap();
7996 init_repo(&repo);
7997
7998 let mut config = Config::default();
7999 config.merge.mode = MergeMode::Local;
8000
8001 let mut state = RunState::new(
8002 repo.clone(),
8003 "main".to_owned(),
8004 "deadbeef".to_owned(),
8005 "task".to_owned(),
8006 config,
8007 );
8008 state.candidates = vec![Candidate {
8009 index: 0,
8010 label: 'A',
8011 agent: "alpha".to_owned(),
8012 branch: "does-not-exist".to_owned(),
8013 worktree: repo.clone(),
8014 summary: String::new(),
8015 stat: String::new(),
8016 files: 0,
8017 commits: 0,
8018 empty: false,
8019 failed: None,
8020 verified_noop: None,
8021 duration_ms: 0,
8022 folded: false,
8023 }];
8024 state.tally = Some(Tally {
8025 first_choice: BTreeMap::from([('A', 1)]),
8026 borda: BTreeMap::new(),
8027 winner: 'A',
8028 rankings: 1,
8029 unanimous_initial: true,
8030 deliberated: false,
8031 changed_votes: 0,
8032 unanimous_final: true,
8033 tie_break: None,
8034 judges: 0,
8035 present: 0,
8036 quorum: 0,
8037 met_quorum: true,
8038 uncontested: Some("only candidate A produced a change".to_owned()),
8039 });
8040 state.reviews = vec![ReviewRound {
8041 round: 1,
8042 head: "deadbeef".to_owned(),
8043 verified_head: None,
8044 verified_at: None,
8045 reviews: Vec::new(),
8046 e2e: Vec::new(),
8047 fix: None,
8048 blocking: 0,
8049 answered: 0,
8050 expected: 0,
8051 clean: true,
8052 verify_retried: false,
8053 e2e_deferred: false,
8054 e2e_defer_reason: None,
8055 progressed: false,
8056 vote_split: false,
8057 reconsideration: Vec::new(),
8058 verdict: None,
8059 }];
8060 state.gate = Vec::new();
8062 state.gate_ran = false;
8063 state.status = RunStatus::Gating;
8064
8065 let mut runner = Runner {
8066 state,
8067 roles: ResolvedRoles {
8068 implementers: Vec::new(),
8069 judges: Vec::new(),
8070 reviewers: Vec::new(),
8071 fixer: None,
8072 conductor: conductor(),
8073 implementer_roster: Vec::new(),
8074 },
8075 sem: Arc::new(Semaphore::new(1)),
8076 pause: Pause::new(),
8077 interrupt: Pause::new(),
8078 };
8079
8080 runner.merge().await.expect("merge");
8081
8082 assert!(
8083 runner.state.merge.is_none(),
8084 "an empty gate must never be read as a passing one: {:?}",
8085 runner.state.merge
8086 );
8087 }
8088
8089 #[tokio::test]
8096 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
8097 let tmp = tempfile::tempdir().expect("tempdir");
8098 let repo = tmp.path().join("repo");
8099 std::fs::create_dir_all(&repo).unwrap();
8100 init_repo(&repo);
8101
8102 let config = Config::default();
8104
8105 let mut state = RunState::new(
8106 repo.clone(),
8107 "main".to_owned(),
8108 "deadbeef".to_owned(),
8109 "task".to_owned(),
8110 config,
8111 );
8112 state.candidates = vec![Candidate {
8113 index: 0,
8114 label: 'A',
8115 agent: "alpha".to_owned(),
8116 branch: "does-not-exist".to_owned(),
8117 worktree: repo.clone(),
8118 summary: String::new(),
8119 stat: String::new(),
8120 files: 0,
8121 commits: 0,
8122 empty: false,
8123 failed: None,
8124 verified_noop: None,
8125 duration_ms: 0,
8126 folded: false,
8127 }];
8128 state.tally = Some(Tally {
8129 first_choice: BTreeMap::from([('A', 1)]),
8130 borda: BTreeMap::new(),
8131 winner: 'A',
8132 rankings: 1,
8133 unanimous_initial: true,
8134 deliberated: false,
8135 changed_votes: 0,
8136 unanimous_final: true,
8137 tie_break: None,
8138 judges: 0,
8139 present: 0,
8140 quorum: 0,
8141 met_quorum: true,
8142 uncontested: Some("only candidate A produced a change".to_owned()),
8143 });
8144 state.reviews = vec![ReviewRound {
8145 round: 1,
8146 head: "deadbeef".to_owned(),
8147 verified_head: None,
8148 verified_at: None,
8149 reviews: Vec::new(),
8150 e2e: Vec::new(),
8151 fix: None,
8152 blocking: 0,
8153 answered: 0,
8154 expected: 0,
8155 clean: true,
8156 verify_retried: false,
8157 e2e_deferred: false,
8158 e2e_defer_reason: None,
8159 progressed: false,
8160 vote_split: false,
8161 reconsideration: Vec::new(),
8162 verdict: None,
8163 }];
8164
8165 let mut runner = Runner {
8166 state,
8167 roles: ResolvedRoles {
8168 implementers: Vec::new(),
8169 judges: Vec::new(),
8170 reviewers: Vec::new(),
8171 fixer: None,
8172 conductor: conductor(),
8173 implementer_roster: Vec::new(),
8174 },
8175 sem: Arc::new(Semaphore::new(1)),
8176 pause: Pause::new(),
8177 interrupt: Pause::new(),
8178 };
8179
8180 runner.gate().await.expect("gate");
8181 assert!(
8182 runner.state.gate_ran,
8183 "zero configured commands is still a real attempt, not an unrun gate"
8184 );
8185 assert!(runner.state.gate.is_empty());
8186 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
8187 assert_ne!(
8188 runner.state.status,
8189 RunStatus::Blocked,
8190 "a gate with nothing to check must not read as failed"
8191 );
8192
8193 runner.merge().await.expect("merge");
8194 assert_eq!(
8195 runner.state.status,
8196 RunStatus::Ready,
8197 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
8198 );
8199 }
8200
8201 #[tokio::test]
8212 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
8213 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8214 let home = crate::run::home();
8215
8216 let tmp = tempfile::tempdir().expect("tempdir");
8217 let repo = tmp.path().join("repo");
8218 std::fs::create_dir_all(&repo).unwrap();
8219 init_repo(&repo);
8220 let cache_dir = tmp.path().join("target");
8223
8224 let mut config = Config::default();
8225 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
8226 config.graph.timeout_verify = Some(2);
8229
8230 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8231 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8232 .expect("no io error acquiring directly")
8233 {
8234 crate::cache::AcquireOutcome::Acquired(g) => g,
8235 crate::cache::AcquireOutcome::Busy(b) => {
8236 panic!("expected the direct acquire to win the lease first: {b:?}")
8237 }
8238 };
8239
8240 let mut state = RunState::new(
8241 repo.clone(),
8242 "main".to_owned(),
8243 "deadbeef".to_owned(),
8244 "task".to_owned(),
8245 config,
8246 );
8247 state.candidates = vec![Candidate {
8248 index: 0,
8249 label: 'A',
8250 agent: "alpha".to_owned(),
8251 branch: "does-not-exist".to_owned(),
8252 worktree: repo.clone(),
8253 summary: String::new(),
8254 stat: String::new(),
8255 files: 0,
8256 commits: 0,
8257 empty: false,
8258 failed: None,
8259 verified_noop: None,
8260 duration_ms: 0,
8261 folded: false,
8262 }];
8263 state.tally = Some(Tally {
8264 first_choice: BTreeMap::from([('A', 1)]),
8265 borda: BTreeMap::new(),
8266 winner: 'A',
8267 rankings: 1,
8268 unanimous_initial: true,
8269 deliberated: false,
8270 changed_votes: 0,
8271 unanimous_final: true,
8272 tie_break: None,
8273 judges: 0,
8274 present: 0,
8275 quorum: 0,
8276 met_quorum: true,
8277 uncontested: Some("only candidate A produced a change".to_owned()),
8278 });
8279 state.reviews = vec![ReviewRound {
8280 round: 1,
8281 head: "deadbeef".to_owned(),
8282 verified_head: None,
8283 verified_at: None,
8284 reviews: Vec::new(),
8285 e2e: Vec::new(),
8286 fix: None,
8287 blocking: 0,
8288 answered: 0,
8289 expected: 0,
8290 clean: true,
8291 verify_retried: false,
8292 e2e_deferred: false,
8293 e2e_defer_reason: None,
8294 progressed: false,
8295 vote_split: false,
8296 reconsideration: Vec::new(),
8297 verdict: None,
8298 }];
8299
8300 let mut runner = Runner {
8301 state,
8302 roles: ResolvedRoles {
8303 implementers: Vec::new(),
8304 judges: Vec::new(),
8305 reviewers: Vec::new(),
8306 fixer: None,
8307 conductor: conductor(),
8308 implementer_roster: Vec::new(),
8309 },
8310 sem: Arc::new(Semaphore::new(1)),
8311 pause: Pause::new(),
8312 interrupt: Pause::new(),
8313 };
8314
8315 let started = std::time::Instant::now();
8316 runner.gate().await.expect("gate");
8317 assert!(
8318 started.elapsed() < Duration::from_secs(1),
8319 "a gate with nothing to run must never wait on a lease it never needed"
8320 );
8321 assert!(
8322 runner.state.gate_ran,
8323 "zero commands is still a real, immediate attempt"
8324 );
8325 assert!(runner.state.gate.is_empty());
8326 assert_ne!(
8327 runner.state.status,
8328 RunStatus::Blocked,
8329 "must not read as resource-blocked on a lease it never asked for"
8330 );
8331 }
8332
8333 #[tokio::test]
8343 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8344 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8345
8346 let tmp = tempfile::tempdir().expect("tempdir");
8347 let repo = tmp.path().join("repo");
8348 std::fs::create_dir_all(&repo).unwrap();
8349 init_repo(&repo);
8350
8351 let mut config = Config::default();
8352 config.verify.gate = vec![
8353 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8354 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8355 .to_owned(),
8356 ];
8357
8358 let mut state = RunState::new(
8359 repo.clone(),
8360 "main".to_owned(),
8361 "deadbeef".to_owned(),
8362 "task".to_owned(),
8363 config,
8364 );
8365 let run_id = state.id.clone();
8366 state.candidates = vec![Candidate {
8367 index: 0,
8368 label: 'A',
8369 agent: "alpha".to_owned(),
8370 branch: "does-not-exist".to_owned(),
8371 worktree: repo.clone(),
8372 summary: String::new(),
8373 stat: String::new(),
8374 files: 0,
8375 commits: 0,
8376 empty: false,
8377 failed: None,
8378 verified_noop: None,
8379 duration_ms: 0,
8380 folded: false,
8381 }];
8382 state.tally = Some(Tally {
8383 first_choice: BTreeMap::from([('A', 1)]),
8384 borda: BTreeMap::new(),
8385 winner: 'A',
8386 rankings: 1,
8387 unanimous_initial: true,
8388 deliberated: false,
8389 changed_votes: 0,
8390 unanimous_final: true,
8391 tie_break: None,
8392 judges: 0,
8393 present: 0,
8394 quorum: 0,
8395 met_quorum: true,
8396 uncontested: Some("only candidate A produced a change".to_owned()),
8397 });
8398 state.reviews = vec![ReviewRound {
8399 round: 1,
8400 head: "deadbeef".to_owned(),
8401 verified_head: None,
8402 verified_at: None,
8403 reviews: Vec::new(),
8404 e2e: Vec::new(),
8405 fix: None,
8406 blocking: 0,
8407 answered: 0,
8408 expected: 0,
8409 clean: true,
8410 verify_retried: false,
8411 e2e_deferred: false,
8412 e2e_defer_reason: None,
8413 progressed: false,
8414 vote_split: false,
8415 reconsideration: Vec::new(),
8416 verdict: None,
8417 }];
8418
8419 let mut runner = Runner {
8420 state,
8421 roles: ResolvedRoles {
8422 implementers: Vec::new(),
8423 judges: Vec::new(),
8424 reviewers: Vec::new(),
8425 fixer: None,
8426 conductor: conductor(),
8427 implementer_roster: Vec::new(),
8428 },
8429 sem: Arc::new(Semaphore::new(1)),
8430 pause: Pause::new(),
8431 interrupt: Pause::new(),
8432 };
8433
8434 let started_marker = repo.join("started.marker");
8435 let release_marker = repo.join("release.marker");
8436 let poller = tokio::spawn(async move {
8437 for _ in 0..100 {
8442 if started_marker.exists()
8443 && let Ok(s) = crate::run::RunState::load(&run_id)
8444 && let Some(a) = s.active.get("gate")
8445 {
8446 std::fs::write(&release_marker, b"go").expect("release marker");
8447 return Some(a.clone());
8448 }
8449 tokio::time::sleep(Duration::from_millis(50)).await;
8450 }
8451 None
8452 });
8453
8454 runner.gate().await.expect("gate");
8455 let captured = poller.await.expect("poller task");
8456 let captured = captured.expect(
8457 "the poller never saw a `gate` task entry in run.json while the command was \
8458 still blocked on its own release marker",
8459 );
8460
8461 assert_eq!(captured.task.as_deref(), Some("gate"));
8462 assert_eq!(captured.node, "gate");
8463 assert_eq!(captured.index, Some(1));
8464 assert_eq!(captured.total, Some(1));
8465 assert!(
8466 captured
8467 .command
8468 .as_deref()
8469 .is_some_and(|c| c.contains("started.marker")),
8470 "{captured:?}"
8471 );
8472
8473 assert!(
8474 runner.state.active.is_empty(),
8475 "the entry must be cleared once the command actually finished: {:?}",
8476 runner.state.active
8477 );
8478 assert!(runner.state.gate_ran);
8479 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8480 }
8481
8482 #[tokio::test]
8495 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8496 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8497 let home = crate::run::home();
8498
8499 let tmp = tempfile::tempdir().expect("tempdir");
8500 let repo = tmp.path().join("repo");
8501 std::fs::create_dir_all(&repo).unwrap();
8502 init_repo(&repo);
8503 let head = crate::git::rev_parse(&repo, "HEAD")
8504 .await
8505 .expect("rev-parse");
8506 let cache_dir = tmp.path().join("target");
8509
8510 let mut config = Config::default();
8511 config.verify.e2e = vec![format!(
8512 "CARGO_TARGET_DIR='{}' test -f README.md",
8513 cache_dir.display()
8514 )];
8515 config.graph.review_rounds = 1;
8516 config.graph.timeout_verify = Some(2);
8519
8520 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8521 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8522 .expect("no io error acquiring directly")
8523 {
8524 crate::cache::AcquireOutcome::Acquired(g) => g,
8525 crate::cache::AcquireOutcome::Busy(b) => {
8526 panic!("expected the direct acquire to win the lease first: {b:?}")
8527 }
8528 };
8529
8530 let mut state = RunState::new(
8531 repo.clone(),
8532 "main".to_owned(),
8533 head.clone(),
8534 "task".to_owned(),
8535 config,
8536 );
8537 state.candidates = vec![Candidate {
8538 index: 0,
8539 label: 'A',
8540 agent: "alpha".to_owned(),
8541 branch: "does-not-exist".to_owned(),
8542 worktree: repo.clone(),
8543 summary: String::new(),
8544 stat: String::new(),
8545 files: 0,
8546 commits: 0,
8547 empty: false,
8548 failed: None,
8549 verified_noop: None,
8550 duration_ms: 0,
8551 folded: false,
8552 }];
8553 state.tally = Some(Tally {
8554 first_choice: BTreeMap::from([('A', 1)]),
8555 borda: BTreeMap::new(),
8556 winner: 'A',
8557 rankings: 1,
8558 unanimous_initial: true,
8559 deliberated: false,
8560 changed_votes: 0,
8561 unanimous_final: true,
8562 tie_break: None,
8563 judges: 0,
8564 present: 0,
8565 quorum: 0,
8566 met_quorum: true,
8567 uncontested: Some("only candidate A produced a change".to_owned()),
8568 });
8569 state.reviews = vec![ReviewRound {
8573 round: 1,
8574 head: head.clone(),
8575 verified_head: None,
8576 verified_at: None,
8577 reviews: Vec::new(),
8578 e2e: Vec::new(),
8579 fix: None,
8580 blocking: 1,
8581 answered: 1,
8582 expected: 1,
8583 clean: false,
8584 verify_retried: false,
8585 e2e_deferred: true,
8586 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8587 progressed: false,
8588 vote_split: false,
8589 reconsideration: Vec::new(),
8590 verdict: None,
8591 }];
8592
8593 let mut runner = Runner {
8594 state,
8595 roles: ResolvedRoles {
8596 implementers: Vec::new(),
8597 judges: Vec::new(),
8598 reviewers: Vec::new(),
8599 fixer: None,
8600 conductor: conductor(),
8601 implementer_roster: Vec::new(),
8602 },
8603 sem: Arc::new(Semaphore::new(1)),
8604 pause: Pause::new(),
8605 interrupt: Pause::new(),
8606 };
8607
8608 let shell = runner.state.config.shell();
8609 runner
8610 .stop_reviewing("round budget spent", &shell, &repo)
8611 .await
8612 .expect("stop_reviewing");
8613
8614 let last = runner.state.reviews.last().expect("round record");
8615 assert_eq!(
8616 last.e2e_status(),
8617 E2eStatus::ResourceBlocked,
8618 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8619 failed: {last:?}"
8620 );
8621 assert_eq!(
8622 last.verified_head.as_deref(),
8623 Some(head.as_str()),
8624 "which commit this attempt targeted is known even though nothing finished checking \
8625 it"
8626 );
8627 let first_attempt_at = last
8628 .verified_at
8629 .expect("when this attempt ran is known too");
8630 assert_ne!(
8631 runner.state.status,
8632 RunStatus::Blocked,
8633 "contention is evidence about the machine, not the patch — it must not settle the \
8634 run as blocked: {:?}",
8635 runner.state.status
8636 );
8637 assert!(
8638 !runner
8639 .state
8640 .events
8641 .iter()
8642 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
8643 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
8644 runner.state.events
8645 );
8646
8647 runner
8652 .stop_reviewing("round budget spent", &shell, &repo)
8653 .await
8654 .expect("stop_reviewing retry");
8655 assert_eq!(
8656 runner.state.reviews.len(),
8657 1,
8658 "no new round was started: {:?}",
8659 runner.state.reviews
8660 );
8661 let last = runner.state.reviews.last().expect("round record");
8662 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
8663 assert!(
8664 last.verified_at.expect("still known") > first_attempt_at,
8665 "a second reentry must be a fresh attempt, not a stale copy of the first"
8666 );
8667 assert_ne!(runner.state.status, RunStatus::Blocked);
8668
8669 held.release();
8670 }
8671
8672 #[tokio::test]
8684 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
8685 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8686 let home = crate::run::home();
8687
8688 let tmp = tempfile::tempdir().expect("tempdir");
8689 let repo = tmp.path().join("repo");
8690 std::fs::create_dir_all(&repo).unwrap();
8691 init_repo(&repo);
8692 let head = crate::git::rev_parse(&repo, "HEAD")
8693 .await
8694 .expect("rev-parse");
8695 let cache_dir = tmp.path().join("target");
8696
8697 let mut config = Config::default();
8698 config.verify.e2e = vec![format!(
8699 "CARGO_TARGET_DIR='{}' test -f README.md",
8700 cache_dir.display()
8701 )];
8702 config.graph.review_rounds = 1;
8703 config.graph.timeout_verify = Some(2);
8704
8705 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8706 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8707 .expect("no io error acquiring directly")
8708 {
8709 crate::cache::AcquireOutcome::Acquired(g) => g,
8710 crate::cache::AcquireOutcome::Busy(b) => {
8711 panic!("expected the direct acquire to win the lease first: {b:?}")
8712 }
8713 };
8714
8715 let mut state = RunState::new(
8716 repo.clone(),
8717 "main".to_owned(),
8718 head.clone(),
8719 "task".to_owned(),
8720 config,
8721 );
8722 state.candidates = vec![Candidate {
8723 index: 0,
8724 label: 'A',
8725 agent: "alpha".to_owned(),
8726 branch: "does-not-exist".to_owned(),
8727 worktree: repo.clone(),
8728 summary: String::new(),
8729 stat: String::new(),
8730 files: 0,
8731 commits: 0,
8732 empty: false,
8733 failed: None,
8734 verified_noop: None,
8735 duration_ms: 0,
8736 folded: false,
8737 }];
8738 state.tally = Some(Tally {
8739 first_choice: BTreeMap::from([('A', 1)]),
8740 borda: BTreeMap::new(),
8741 winner: 'A',
8742 rankings: 1,
8743 unanimous_initial: true,
8744 deliberated: false,
8745 changed_votes: 0,
8746 unanimous_final: true,
8747 tie_break: None,
8748 judges: 0,
8749 present: 0,
8750 quorum: 0,
8751 met_quorum: true,
8752 uncontested: Some("only candidate A produced a change".to_owned()),
8753 });
8754 state.reviews = vec![ReviewRound {
8758 round: 1,
8759 head: head.clone(),
8760 verified_head: Some(head.clone()),
8761 verified_at: Some(jiff::Timestamp::now()),
8762 reviews: Vec::new(),
8763 e2e: vec![CommandOutcome {
8764 command: format!(
8765 "CARGO_TARGET_DIR='{}' test -f README.md",
8766 cache_dir.display()
8767 ),
8768 code: None,
8769 output_tail: "waiting for the shared build cache".to_owned(),
8770 duration_ms: 0,
8771 resource_blocked: true,
8772 }],
8773 fix: None,
8774 blocking: 1,
8775 answered: 1,
8776 expected: 1,
8777 clean: false,
8778 verify_retried: false,
8779 e2e_deferred: false,
8780 e2e_defer_reason: None,
8781 progressed: false,
8782 vote_split: false,
8783 reconsideration: Vec::new(),
8784 verdict: None,
8785 }];
8786
8787 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
8788 let mut runner = Runner {
8789 state,
8790 roles: ResolvedRoles {
8791 implementers: Vec::new(),
8792 judges: Vec::new(),
8793 reviewers: Vec::new(),
8794 fixer: None,
8795 conductor: conductor(),
8796 implementer_roster: Vec::new(),
8797 },
8798 sem: Arc::new(Semaphore::new(1)),
8799 pause: Pause::new(),
8800 interrupt: Pause::new(),
8801 };
8802
8803 runner.review_loop().await.expect("review_loop");
8808
8809 assert_eq!(
8810 runner.state.reviews.len(),
8811 1,
8812 "no new round was started on top of the unresolved one: {:?}",
8813 runner.state.reviews
8814 );
8815 let last = &runner.state.reviews[0];
8816 assert_eq!(
8817 last.e2e_status(),
8818 E2eStatus::ResourceBlocked,
8819 "still contended: {last:?}"
8820 );
8821 assert!(
8822 last.verified_at.expect("still known") > first_attempt_at,
8823 "review_loop must have actually retried the check, not left it exactly as found"
8824 );
8825 assert_ne!(
8826 runner.state.status,
8827 RunStatus::Blocked,
8828 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
8829 runner.state.status
8830 );
8831
8832 held.release();
8833 }
8834
8835 #[tokio::test]
8836 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
8837 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8838 let tmp = tempfile::tempdir().expect("tempdir");
8839 let repo = tmp.path().join("repo");
8840 std::fs::create_dir_all(&repo).unwrap();
8841 init_repo(&repo);
8842
8843 let mut config = Config::default();
8844 config.merge.mode = MergeMode::Pr;
8845 config.graph.land = true;
8846 config.graph.land_approval = false;
8847
8848 let mut state = RunState::new(
8849 repo.clone(),
8850 "main".to_owned(),
8851 "deadbeef".to_owned(),
8852 "task".to_owned(),
8853 config,
8854 );
8855 state.candidates = vec![Candidate {
8856 index: 0,
8857 label: 'A',
8858 agent: "alpha".to_owned(),
8859 branch: "does-not-exist".to_owned(),
8860 worktree: repo.clone(),
8861 summary: String::new(),
8862 stat: String::new(),
8863 files: 0,
8864 commits: 0,
8865 empty: false,
8866 failed: None,
8867 verified_noop: None,
8868 duration_ms: 0,
8869 folded: false,
8870 }];
8871 state.tally = Some(Tally {
8872 first_choice: BTreeMap::from([('A', 1)]),
8873 borda: BTreeMap::new(),
8874 winner: 'A',
8875 rankings: 1,
8876 unanimous_initial: true,
8877 deliberated: false,
8878 changed_votes: 0,
8879 unanimous_final: true,
8880 tie_break: None,
8881 judges: 0,
8882 present: 0,
8883 quorum: 0,
8884 met_quorum: true,
8885 uncontested: Some("only candidate A produced a change".to_owned()),
8886 });
8887 state.reviews = vec![ReviewRound {
8888 round: 1,
8889 head: "deadbeef".to_owned(),
8890 verified_head: None,
8891 verified_at: None,
8892 reviews: Vec::new(),
8893 e2e: Vec::new(),
8894 fix: None,
8895 blocking: 0,
8896 answered: 0,
8897 expected: 0,
8898 clean: true,
8899 verify_retried: false,
8900 e2e_deferred: false,
8901 e2e_defer_reason: None,
8902 progressed: false,
8903 vote_split: false,
8904 reconsideration: Vec::new(),
8905 verdict: None,
8906 }];
8907 state.gate = vec![CommandOutcome {
8908 command: "test".to_owned(),
8909 code: Some(0),
8910 output_tail: String::new(),
8911 duration_ms: 0,
8912 resource_blocked: false,
8913 }];
8914 state.gate_ran = true;
8915 state.status = RunStatus::Landing;
8919 state.merge = Some(MergeOutcome {
8920 mode: MergeMode::Pr,
8921 ok: true,
8922 detail: "https://example.invalid/x/y/pull/1".to_owned(),
8923 });
8924
8925 ask_test_home();
8929 let store = ask::Questions::open();
8930 let q = ask_open_question(&store, &state.id);
8931
8932 let mut runner = Runner {
8933 state,
8934 roles: ResolvedRoles {
8935 implementers: Vec::new(),
8936 judges: Vec::new(),
8937 reviewers: Vec::new(),
8938 fixer: None,
8939 conductor: conductor(),
8940 implementer_roster: Vec::new(),
8941 },
8942 sem: Arc::new(Semaphore::new(1)),
8943 pause: Pause::new(),
8944 interrupt: Pause::new(),
8945 };
8946
8947 runner.execute().await.expect("execute");
8952
8953 assert_eq!(
8954 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8955 Some("https://example.invalid/x/y/pull/1"),
8956 "reentry must not push again or open a second pull request over the \
8957 one `land` is already watching"
8958 );
8959 assert_ne!(
8960 runner.state.status,
8961 RunStatus::Landing,
8962 "land could not actually reach the fake pull request, so it must \
8963 have given up rather than left the run silently parked forever"
8964 );
8965 assert_eq!(runner.state.status, RunStatus::Blocked);
8969 assert!(
8970 store.get(&q.id).unwrap().status.open(),
8971 "Blocked is still alive; settle_questions must have been a no-op here"
8972 );
8973 }
8974
8975 fn state_with_round(round: ReviewRound) -> RunState {
8976 let mut s = RunState::new(
8977 PathBuf::from("/repo"),
8978 "main".to_owned(),
8979 "abc1234".to_owned(),
8980 "add retries".to_owned(),
8981 Config::default(),
8982 );
8983 s.reviews = vec![round];
8984 s
8985 }
8986
8987 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
8988 crate::verdict::Finding {
8989 id: id.to_owned(),
8990 severity,
8991 file: None,
8992 line: None,
8993 title: title.to_owned(),
8994 detail: String::new(),
8995 }
8996 }
8997
8998 #[test]
8999 fn pr_body_names_open_findings_and_declined_ones() {
9000 let round = ReviewRound {
9001 round: 2,
9002 head: "deadbee".to_owned(),
9003 verified_head: None,
9004 verified_at: None,
9005 reviews: vec![ReviewRecord {
9006 attempts: 0,
9007 reviewer: 1,
9008 agent: "alpha".to_owned(),
9009 summary: String::new(),
9010 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
9011 vote: None,
9012 failed: None,
9013 duration_ms: 0,
9014 }],
9015 e2e: vec![CommandOutcome {
9016 command: "cargo test".to_owned(),
9017 code: Some(0),
9018 output_tail: String::new(),
9019 duration_ms: 0,
9020 resource_blocked: false,
9021 }],
9022 verify_retried: false,
9023 e2e_deferred: false,
9024 e2e_defer_reason: None,
9025 fix: Some(FixRecord {
9026 agent: "alpha".to_owned(),
9027 addressed: Vec::new(),
9028 rejected: vec![crate::verdict::Rejection {
9029 id: "R1-1-1".to_owned(),
9030 why: "not reachable from any caller".to_owned(),
9031 }],
9032 notes: String::new(),
9033 committed: true,
9034 failed: None,
9035 duration_ms: 0,
9036 continuation: None,
9037 }),
9038 blocking: 0,
9039 answered: 1,
9040 expected: 1,
9041 clean: false,
9042 progressed: true,
9043 vote_split: false,
9044 reconsideration: Vec::new(),
9045 verdict: None,
9046 };
9047 let state = state_with_round(round);
9048 let body = pr_message(&state, 'A').body;
9049
9050 assert!(body.contains("add retries"), "the task must still be there");
9051 assert!(body.contains("R2-1-1"), "{body}");
9052 assert!(body.contains("unused import"), "{body}");
9053 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
9054 assert!(
9055 body.contains("not reachable from any caller"),
9056 "the reason it was declined: {body}"
9057 );
9058 }
9059
9060 #[test]
9061 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
9062 let round = ReviewRound {
9063 round: 1,
9064 head: "deadbee".to_owned(),
9065 verified_head: None,
9066 verified_at: None,
9067 reviews: vec![ReviewRecord {
9068 attempts: 0,
9069 reviewer: 1,
9070 agent: "alpha".to_owned(),
9071 summary: String::new(),
9072 findings: Vec::new(),
9073 vote: None,
9074 failed: None,
9075 duration_ms: 0,
9076 }],
9077 e2e: Vec::new(),
9078 verify_retried: false,
9079 e2e_deferred: false,
9080 e2e_defer_reason: None,
9081 fix: None,
9082 blocking: 0,
9083 answered: 1,
9084 expected: 1,
9085 clean: true,
9086 progressed: false,
9087 vote_split: false,
9088 reconsideration: Vec::new(),
9089 verdict: None,
9090 };
9091 let state = state_with_round(round);
9092 let body = pr_message(&state, 'A').body;
9093 assert!(!body.contains("Open review findings"), "{body}");
9094 assert!(!body.contains("Declined"), "{body}");
9095 }
9096
9097 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
9098 let mut state = RunState::new(
9099 PathBuf::from("/repo"),
9100 "main".to_owned(),
9101 "abc1234".to_owned(),
9102 instruction.to_owned(),
9103 Config::default(),
9104 );
9105 state.candidates.push(Candidate {
9106 index: 0,
9107 label: 'A',
9108 agent: "alpha".to_owned(),
9109 branch: "magi/x/A".to_owned(),
9110 worktree: PathBuf::from("/wt"),
9111 summary: summary.to_owned(),
9112 stat: String::new(),
9113 files: 1,
9114 commits: 1,
9115 empty: false,
9116 failed: None,
9117 verified_noop: None,
9118 folded: false,
9119 duration_ms: 0,
9120 });
9121 state
9122 }
9123
9124 #[test]
9125 fn pr_message_describes_the_change_not_the_task() {
9126 let state = state_with_summary(
9127 "今回やってほしいこと: results projector を直す",
9128 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
9129 );
9130 let m = pr_message(&state, 'A');
9131 assert_eq!(m.title, "fix(web): batch the runs list reads");
9132 assert!(
9133 m.body.starts_with("## Summary\n\n- reads run.json once"),
9134 "{}",
9135 m.body
9136 );
9137 assert!(!m.body.contains("TITLE:"), "{}", m.body);
9138 let task_at = m.body.find("今回やってほしいこと").unwrap();
9139 let details_at = m.body.find("<details>").unwrap();
9140 assert!(
9141 details_at < task_at,
9142 "the task lives inside <details>: {}",
9143 m.body
9144 );
9145 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
9146 assert!(m.body.contains("magi:candidate-a"));
9147 }
9148
9149 #[test]
9150 fn pr_message_falls_back_to_the_task_without_a_title_line() {
9151 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
9152 let m = pr_message(&state, 'A');
9153 assert_eq!(m.title, "add retries");
9154 assert!(
9155 m.body.contains("## Summary\n\n- did some things"),
9156 "{}",
9157 m.body
9158 );
9159
9160 let none = RunState::new(
9161 PathBuf::from("/repo"),
9162 "main".to_owned(),
9163 "abc1234".to_owned(),
9164 "add retries".to_owned(),
9165 Config::default(),
9166 );
9167 let m = pr_message(&none, 'A');
9168 assert_eq!(m.title, "add retries");
9169 assert!(!m.body.contains("## Summary"), "{}", m.body);
9170 }
9171
9172 #[test]
9173 fn pr_message_refuses_the_candidate_commit_subject() {
9174 for bad in [
9175 "TITLE: magi: candidate A (uncommitted work)",
9176 "TITLE: chore: stuff (uncommitted work)",
9177 "TITLE: ",
9178 ] {
9179 let state = state_with_summary("add retries", bad);
9180 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
9181 }
9182 }
9183
9184 #[test]
9185 fn pr_message_bounds_a_very_long_task_and_title() {
9186 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
9187 let state = state_with_summary(&long, "- nothing");
9188 let m = pr_message(&state, 'A');
9189 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9190 assert!(!m.title.contains('\n'));
9191
9192 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
9193 let m = pr_message(&state, 'A');
9194 assert!(m.title.starts_with("feat: "));
9195 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9196 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
9197 }
9198
9199 #[test]
9200 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
9201 let mut state = state_with_summary(
9205 "add retries",
9206 "TITLE: fix(web): batch reads\n- reads run.json once",
9207 );
9208 state.config.graph.language = "ja".to_owned();
9209 let m = pr_message(&state, 'A');
9210 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
9211
9212 let task = "今回やってほしいこと: results projector を直す";
9215 let mut state = state_with_summary(task, "- no title line");
9216 state.config.graph.language = "ja".to_owned();
9217 let m = pr_message(&state, 'A');
9218 assert_eq!(
9219 m.title,
9220 format!("chore: land candidate A of run {}", state.id)
9221 );
9222 assert!(
9223 m.body.contains(&format!(
9224 "<summary>Original task</summary>\n\n{task}\n\n</details>"
9225 )),
9226 "{}",
9227 m.body
9228 );
9229 }
9230
9231 #[test]
9232 fn pr_message_scrubs_home_paths_and_addresses() {
9233 let state = state_with_summary(
9234 "fix it in /Users/someone/src/x",
9235 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9236 );
9237 let m = pr_message(&state, 'A');
9238 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9239 assert!(!m.body.contains(leak), "{}", m.body);
9240 }
9241 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9242 }
9243
9244 #[test]
9245 fn pr_message_survives_a_task_that_closes_details() {
9246 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9247 let m = pr_message(&state, 'A');
9248 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9249 }
9250
9251 #[test]
9252 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9253 let cmd = manual_merge_command(
9254 MergeStyle::Squash,
9255 Path::new("/repo"),
9256 "b",
9257 "fix: \"quoted\" $(x) `y`\n\nbody",
9258 );
9259 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9260 }
9261
9262 #[test]
9263 fn manual_merge_command_matches_the_configured_style() {
9264 let repo = Path::new("/repo");
9265 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9266
9267 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9268 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9269
9270 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9271 assert_eq!(
9272 squash,
9273 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9274 \"Merge magi run 0832 (candidate A)\""
9275 );
9276
9277 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9278 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9279 }
9280
9281 #[test]
9282 fn a_nudge_gets_a_quarter_of_the_budget() {
9283 assert_eq!(retry_budget(secs(1200), true), secs(300));
9285 assert_eq!(retry_budget(secs(3600), true), secs(900));
9286 }
9287
9288 #[test]
9289 fn a_resent_prompt_keeps_the_whole_budget() {
9290 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9293 assert_eq!(retry_budget(secs(60), false), secs(60));
9294 }
9295
9296 #[test]
9297 fn the_floor_never_exceeds_the_original_budget() {
9298 assert_eq!(retry_budget(secs(60), true), secs(60));
9302 assert_eq!(retry_budget(secs(480), true), secs(120));
9303 assert_eq!(retry_budget(secs(0), true), secs(0));
9304 }
9305
9306 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9307 agent::CommandEvidence {
9308 id: "item1".to_owned(),
9309 description: "cargo test".to_owned(),
9310 exit_code,
9311 result_summary: String::new(),
9312 source: "codex".to_owned(),
9313 }
9314 }
9315
9316 #[test]
9317 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9318 assert!(!has_unconfirmed_command(&[]));
9322 }
9323
9324 #[test]
9325 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9326 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9330 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9331 assert!(!has_unconfirmed_command(&[
9332 evidence(Some(0)),
9333 evidence(Some(101))
9334 ]));
9335 }
9336
9337 #[test]
9338 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9339 assert!(has_unconfirmed_command(&[
9340 evidence(Some(0)),
9341 evidence(None)
9342 ]));
9343 }
9344
9345 #[test]
9346 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9347 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9348 assert_eq!(
9349 verified_noop_claim(true, &[], text).as_deref(),
9350 Some("already fixed by b32cfc4, on main.")
9351 );
9352 }
9353
9354 #[test]
9355 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9356 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9359 assert!(verified_noop_claim(false, &[], text).is_none());
9360 }
9361
9362 #[test]
9363 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9364 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9365 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9366 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9368 }
9369
9370 #[test]
9371 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9372 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9373 }
9374
9375 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9378 runner.state.candidates = shape
9379 .iter()
9380 .enumerate()
9381 .map(|(i, &(empty, verified))| Candidate {
9382 index: i,
9383 label: (b'A' + i as u8) as char,
9384 agent: "sonnet".to_owned(),
9385 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9386 worktree: PathBuf::from(format!("/wt/{i}")),
9387 summary: String::new(),
9388 stat: String::new(),
9389 files: 0,
9390 commits: 0,
9391 empty,
9392 failed: None,
9393 verified_noop: verified.map(str::to_owned),
9394 duration_ms: 0,
9395 folded: false,
9396 })
9397 .collect();
9398 }
9399
9400 #[test]
9401 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9402 ask_test_home();
9403 let mut runner = runner_at(RunStatus::Implementing);
9404 set_candidates(
9405 &mut runner,
9406 &[
9407 (true, Some("already on main at b32cfc4")),
9408 (true, Some("same fix, see the existing test")),
9409 ],
9410 );
9411
9412 runner
9413 .after_implement()
9414 .expect("a verified no-op is not an error");
9415
9416 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9417 }
9418
9419 #[test]
9420 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9421 ask_test_home();
9422 let mut runner = runner_at(RunStatus::Implementing);
9423 set_candidates(
9427 &mut runner,
9428 &[(true, Some("already on main at b32cfc4")), (true, None)],
9429 );
9430
9431 let err = runner
9432 .after_implement()
9433 .expect_err("an unverified empty candidate must still fail the run");
9434
9435 assert!(
9436 err.to_string().contains("no candidate produced a change"),
9437 "{err}"
9438 );
9439 assert_eq!(runner.state.status, RunStatus::Failed);
9440 }
9441
9442 #[test]
9443 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9444 ask_test_home();
9445 let mut runner = runner_at(RunStatus::Implementing);
9446 set_candidates(&mut runner, &[(true, None), (true, None)]);
9447
9448 let err = runner
9449 .after_implement()
9450 .expect_err("no candidate declared anything; this is an ordinary failure");
9451
9452 assert!(
9453 err.to_string().contains("no candidate produced a change"),
9454 "{err}"
9455 );
9456 assert_eq!(runner.state.status, RunStatus::Failed);
9457 }
9458}