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, base: &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 match crate::reconcile::reconcile(repo, remote, branch, &local_sha, &remote_sha, base)
323 .await?
324 {
325 crate::reconcile::Reconciliation::Pushed => {
326 tracing::warn!(
327 "local `{branch}` ({}) is {tracking} ({}) rebased; pushed it over",
328 short(&local_sha),
329 short(&remote_sha)
330 );
331 return Ok(());
332 }
333 crate::reconcile::Reconciliation::Placeholder => {}
334 crate::reconcile::Reconciliation::Genuine(d) => return Err((*d).into()),
335 }
336 }
337 let out = git::git_raw(repo, &["branch", "-f", branch, &tracking]).await?;
338 if !out.ok() {
339 bail!(
340 "local `{branch}` ({}) is stale against {tracking} ({}) but git will not move it: {}",
341 short(&local_sha),
342 short(&remote_sha),
343 out.stderr
344 );
345 }
346 tracing::warn!(
347 "local `{branch}` was stale: fast-forwarded {} -> {}",
348 short(&local_sha),
349 short(&remote_sha)
350 );
351 Ok(())
352}
353
354async fn resolve_base(repo: &Path, base_branch: &str, remote: &str) -> Result<String> {
355 let tracking = format!("{remote}/{base_branch}");
356 let fetched = git::fetch(repo, remote, base_branch).await;
357 if let Ok(out) = &fetched
358 && out.ok()
359 && git::rev_exists(repo, &tracking).await
360 {
361 return git::rev_parse(repo, &tracking).await;
362 }
363 let why = match &fetched {
364 Ok(out) if !out.ok() => out.stderr.lines().next().unwrap_or("").to_owned(),
365 Ok(_) => format!("{remote} has no {base_branch}"),
366 Err(e) => e.to_string(),
367 };
368 tracing::warn!(
369 "could not read {tracking} ({why}); branching off the local \
370 {base_branch} instead, which may be behind"
371 );
372 git::rev_parse(repo, base_branch).await.with_context(|| {
373 format!(
374 "cannot resolve `{base_branch}`; set [merge] base in magi.toml to a \
375 branch that exists"
376 )
377 })
378}
379
380struct FixClaim {
397 path: PathBuf,
398}
399
400impl FixClaim {
401 fn acquire(dir: &Path) -> Result<Self> {
402 std::fs::create_dir_all(dir).with_context(|| format!("create {}", dir.display()))?;
403 let path = dir.join("fix.lock");
404 match Self::create(&path) {
405 Ok(claim) => Ok(claim),
406 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
407 if Self::reclaim_if_dead(&path) {
408 Self::create(&path).with_context(|| format!("lock {}", path.display()))
409 } else {
410 bail!(
411 "another `magi fix` is already running for this run ({} exists)",
412 path.display()
413 )
414 }
415 }
416 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
417 }
418 }
419
420 fn create(path: &Path) -> std::io::Result<Self> {
421 let mut f = std::fs::OpenOptions::new()
422 .write(true)
423 .create_new(true)
424 .open(path)?;
425 use std::io::Write as _;
426 writeln!(f, "{}", std::process::id())?;
428 Ok(Self {
429 path: path.to_owned(),
430 })
431 }
432
433 fn reclaim_if_dead(path: &Path) -> bool {
437 let dead = std::fs::read_to_string(path)
438 .ok()
439 .and_then(|body| body.trim().parse::<u32>().ok())
440 .is_some_and(|pid| !crate::proc::pid_alive(pid));
441 dead && std::fs::remove_file(path).is_ok()
442 }
443}
444
445impl Drop for FixClaim {
446 fn drop(&mut self) {
447 let _ = std::fs::remove_file(&self.path);
448 }
449}
450
451impl Runner {
452 pub async fn start(repo: &Path, instruction: String, config: Config) -> Result<Self> {
454 let repo = git::toplevel(repo).await?;
455 let missing = agent::missing_programs(&config.agents);
456 if !missing.is_empty() {
457 bail!(
458 "these agent programs are not on PATH: {}. Fix the roster in \
459 magi.toml or install them.",
460 missing.join(", ")
461 );
462 }
463 let base_branch = match config.merge.base.clone() {
464 Some(b) => b,
465 None => git::current_branch(&repo)
466 .await?
467 .context("HEAD is detached; set [merge] base in magi.toml")?,
468 };
469 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
470 if !git::is_clean(&repo).await? {
474 tracing::warn!(
475 "{} has uncommitted changes; they are not part of this run, \
476 which branches off {base_branch} ({})",
477 repo.display(),
478 &base_commit[..base_commit.len().min(8)]
479 );
480 }
481 let roles = config.resolve_roles()?;
482 let max_parallel = config.graph.max_parallel.max(1);
483 let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
484 state.event("start", format!("run {} created", state.id));
485 state.save()?;
486 Ok(Self {
487 state,
488 roles,
489 sem: Arc::new(Semaphore::new(max_parallel)),
490 pause: Pause::new(),
491 interrupt: Pause::new(),
492 })
493 }
494
495 pub async fn review(repo: &Path, branch: &str, config: Config) -> Result<Self> {
509 Self::review_taking_over(repo, branch, config, None).await
510 }
511
512 pub async fn review_taking_over(
518 repo: &Path,
519 branch: &str,
520 config: Config,
521 takeover: Option<crate::handover::Takeover>,
522 ) -> Result<Self> {
523 let repo = git::toplevel(repo).await?;
524 let missing = agent::missing_programs(&config.agents);
525 if !missing.is_empty() {
526 bail!(
527 "these agent programs are not on PATH: {}. Fix the roster in \
528 magi.toml or install them.",
529 missing.join(", ")
530 );
531 }
532 let base_branch = match config.merge.base.clone() {
533 Some(b) => b,
534 None => git::current_branch(&repo)
535 .await?
536 .context("HEAD is detached; set [merge] base in magi.toml")?,
537 };
538 if base_branch == branch {
539 bail!("`{branch}` is the base branch; there is nothing to review against");
540 }
541 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
542
543 let roles = config.resolve_roles()?;
544 let max_parallel = config.graph.max_parallel.max(1);
545 let mut state = RunState::new(
546 repo.clone(),
547 base_branch,
548 base_commit.clone(),
549 String::new(),
550 config,
551 );
552
553 let released = match &takeover {
558 Some(takeover) => crate::handover::release(&repo, branch, &state.id, takeover).await?,
559 None => None,
560 };
561 if let Some(released) = &released {
562 state.event(
563 "release",
564 format!(
565 "took `{branch}` over from run {}: its worktree was released",
566 crate::run::short_of(&released.old_id)
567 ),
568 );
569 }
570 if let Some(choice) = takeover.as_ref().and_then(|t| t.choice.as_ref())
574 && let Err(e) =
575 crate::reconcile::apply_choice(&repo, &state.config.merge.remote, branch, choice)
576 .await
577 {
578 if let Some(released) = &released {
579 released.restore(&repo, branch).await;
580 }
581 return Err(e.context("applying the owner's answer about the diverged branch"));
582 }
583 let opened =
584 Self::open_review(&repo, branch, state, roles, max_parallel, base_commit).await;
585 if opened.is_err()
586 && let Some(released) = &released
587 {
588 released.restore(&repo, branch).await;
589 }
590 opened
591 }
592
593 async fn open_review(
596 repo: &Path,
597 branch: &str,
598 mut state: RunState,
599 roles: ResolvedRoles,
600 max_parallel: usize,
601 base_commit: String,
602 ) -> Result<Self> {
603 sync_review_branch(repo, branch, &state.config.merge.remote, &base_commit).await?;
604 let log = git::log_oneline(repo, &base_commit, branch)
607 .await
608 .unwrap_or_default();
609 let instruction = format!(
610 "Review the work already on branch `{branch}`. There is no task \
611 statement: what the change claims to do is whatever its commits \
612 say.\n\n{}",
613 if log.trim().is_empty() {
614 "(no commit messages)"
615 } else {
616 log.trim()
617 }
618 );
619 state.instruction = instruction;
620
621 let worktree = state.worktree_root().join("under-review");
624 if let Some(parent) = worktree.parent() {
625 tokio::fs::create_dir_all(parent).await.ok();
626 }
627 let path = worktree.to_string_lossy().to_string();
628 git::git(repo, &["worktree", "add", &path, branch])
629 .await
630 .with_context(|| {
631 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
632 })?;
633
634 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
635 .await
636 .unwrap_or(0);
637 if commits == 0 {
638 git::worktree_remove(repo, &worktree).await.ok();
639 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
640 }
641 let files = git::changed_files(&worktree, &base_commit, "HEAD")
642 .await
643 .map(|f| f.len())
644 .unwrap_or(0);
645 if files == 0
646 && let (Ok(head_tree), Ok(base_tree)) = (
647 git::tree_of(&worktree, "HEAD").await,
648 git::tree_of(&worktree, &base_commit).await,
649 )
650 && head_tree == base_tree
651 {
652 let head = git::rev_parse(&worktree, "HEAD").await.unwrap_or_default();
653 git::worktree_remove(repo, &worktree).await.ok();
654 bail!(
655 "`{branch}` at {} has a tree identical to base {}; this usually means \
656 the branch ref is stale (check `git rev-parse refs/heads/{branch}` \
657 against `{}/{branch}`) rather than an empty change",
658 short(&head),
659 short(&base_commit),
660 state.config.merge.remote
661 );
662 }
663 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
664 .await
665 .unwrap_or_default();
666
667 state.candidates.push(Candidate {
668 index: 0,
669 label: 'A',
670 agent: "(existing branch)".to_owned(),
673 branch: branch.to_owned(),
674 worktree,
675 summary: String::new(),
676 stat,
677 files,
678 commits,
679 empty: false,
680 failed: None,
681 verified_noop: None,
682 duration_ms: 0,
683 folded: false,
684 });
685 state.tally = Some(Tally {
686 first_choice: BTreeMap::from([('A', 0)]),
687 borda: BTreeMap::new(),
688 winner: 'A',
689 rankings: 0,
690 unanimous_initial: false,
691 deliberated: false,
692 changed_votes: 0,
693 unanimous_final: false,
694 tie_break: None,
695 judges: 0,
699 present: 0,
700 quorum: 0,
701 met_quorum: true,
702 uncontested: Some("review-only run: nothing competed".to_owned()),
703 });
704 state.status = RunStatus::Reviewing;
705 state.event(
706 "start",
707 format!(
708 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
709 state.id
710 ),
711 );
712 state.save()?;
713 Ok(Self {
714 state,
715 roles,
716 sem: Arc::new(Semaphore::new(max_parallel)),
717 pause: Pause::new(),
718 interrupt: Pause::new(),
719 })
720 }
721
722 pub fn resume(id: &str) -> Result<Self> {
724 let state = RunState::load(id)?;
725 if let Some(to) = &state.released_to {
726 bail!(
727 "run {} cannot be resumed: its worktree was released to run {}",
728 state.short(),
729 crate::run::short_of(to)
730 );
731 }
732 let roles = state.config.resolve_roles()?;
733 let max_parallel = state.config.graph.max_parallel.max(1);
734 Ok(Self {
735 state,
736 roles,
737 sem: Arc::new(Semaphore::new(max_parallel)),
738 pause: Pause::new(),
739 interrupt: Pause::new(),
740 })
741 }
742
743 pub async fn execute(&mut self) -> Result<()> {
750 let result = self.execute_graph().await;
751 self.mark_driver_exited();
752 let ended = if result.is_err() {
753 Some(crate::notices::run_stopped(&self.state.id, &self.state))
754 } else {
755 crate::notices::run_ended(&self.state)
756 };
757 if let Some(notice) = ended {
758 crate::notices::raise(notice);
759 }
760 result
761 }
762
763 fn mark_driver_exited(&mut self) {
772 self.state.driver_exited = true;
773 let pid = std::process::id();
774 let Ok(mut disk) = RunState::load(&self.state.id) else {
775 return;
776 };
777 if disk.released_to.is_some() || disk.driver_pid != Some(pid) || disk.driver_exited {
778 return;
779 }
780 disk.driver_exited = true;
781 if let Err(e) = disk.save() {
782 tracing::warn!("could not record that run {} stopped: {e:#}", self.state.id);
783 }
784 }
785
786 async fn execute_graph(&mut self) -> Result<()> {
787 self.state.parked = false;
792 self.state.clear_active();
799 if let Ok(disk) = RunState::load(&self.state.id)
817 && let Some(to) = &disk.released_to
818 {
819 bail!(
820 "run {} cannot continue: its worktree was released to run {}",
821 self.state.short(),
822 crate::run::short_of(to)
823 );
824 }
825 let pid = std::process::id();
826 self.state.driver_pid = Some(pid);
827 self.state.driver_started_at = crate::proc::process_started_at(pid);
828 self.state.driver_exited = false;
829 self.state.save()?;
830 if self.state.status == RunStatus::Stalled {
843 if self.recover_stall().await? {
844 self.finish_after_tally().await?;
845 } else {
846 self.state.save()?;
848 }
849 return Ok(());
850 }
851 if self.state.status == RunStatus::Landing {
861 self.run_land().await?;
862 self.settle_questions();
867 return Ok(());
868 }
869 self.prep().await?;
870 if self.park_here()? {
871 return Ok(());
872 }
873 self.advise().await?;
874 if self.park_here()? {
875 return Ok(());
876 }
877 self.implement().await?;
878 if self.park_here()? {
879 return Ok(());
880 }
881 if self.state.status == RunStatus::VerifiedNoop {
885 return Ok(());
886 }
887 self.judge().await?;
888 if self.park_here()? {
889 return Ok(());
890 }
891 self.deliberate().await?;
892 if self.park_here()? {
893 return Ok(());
894 }
895 self.vote().await?;
896 if self.park_here()? {
897 return Ok(());
898 }
899 self.tally()?;
900 if self.state.status == RunStatus::Stalled {
905 self.state.save()?;
909 return Ok(());
910 }
911 self.finish_after_tally().await?;
912 Ok(())
913 }
914
915 fn park_here(&mut self) -> Result<bool> {
922 if !self.pause.parked() && !self.interrupt.parked() {
927 return Ok(false);
928 }
929 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
930 Some(reason) => format!(
931 "parked after `{}` ({reason}) — resume to carry on from here",
932 self.state.status.as_str()
933 ),
934 None => format!(
935 "parked after `{}` — resume to carry on from here",
936 self.state.status.as_str()
937 ),
938 };
939 self.state.event("park", why);
940 self.state.parked = true;
941 self.state.save()?;
942 Ok(true)
943 }
944
945 pub fn on_pause(&mut self, pause: Pause) {
947 self.pause = pause;
948 }
949
950 pub fn watch_interrupt(&mut self, pause: Pause) {
956 self.interrupt = pause;
957 }
958
959 fn settle_questions(&mut self) {
979 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
980 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
981 }
982 }
983
984 async fn finish_after_tally(&mut self) -> Result<()> {
987 self.fold_losers().await?;
988 self.sync_to_base().await?;
993 self.review_loop().await?;
994 self.sync_to_base().await?;
995 self.gate().await?;
996 self.merge().await?;
997 self.state.save()?;
998 Ok(())
999 }
1000
1001 async fn prep(&mut self) -> Result<()> {
1004 if !self.state.candidates.is_empty() {
1005 return Ok(());
1006 }
1007 self.state.status = RunStatus::Prep;
1008 let repo = self.state.repo.clone();
1009 let base = self.state.base_commit.clone();
1010 let root = self.state.worktree_root();
1011 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
1012
1013 let hooks_dir = self.state.dir().join("hooks");
1016 if self.state.config.blind.commit_msg_hook {
1017 std::fs::create_dir_all(&hooks_dir)
1018 .with_context(|| format!("create {}", hooks_dir.display()))?;
1019 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
1020 let path = hooks_dir.join("commit-msg");
1021 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
1022 make_executable(&path)?;
1023 git::acquire_worktree_config(&repo).await?;
1031 self.state.enabled_worktree_config = true;
1032 }
1033
1034 for (index, (spec, label)) in self
1035 .roles
1036 .implementers
1037 .clone()
1038 .into_iter()
1039 .zip(labels)
1040 .enumerate()
1041 {
1042 let branch = self.state.branch_for(label);
1043 let worktree = root.join(format!("cand-{label}"));
1044 git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
1045 if self.state.config.blind.commit_msg_hook {
1046 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
1047 }
1048 git::local_exclude(&worktree, "/.magi/").await?;
1049 self.state.candidates.push(Candidate {
1050 index,
1051 label,
1052 agent: spec.id.clone(),
1053 branch,
1054 worktree,
1055 summary: String::new(),
1056 stat: String::new(),
1057 files: 0,
1058 commits: 0,
1059 empty: false,
1060 failed: None,
1061 verified_noop: None,
1062 duration_ms: 0,
1063 folded: false,
1064 });
1065 }
1066
1067 for j in 1..=self.roles.judges.len() {
1068 let wt = root.join(format!("judge-{j}"));
1069 if !wt.exists() {
1070 git::worktree_add_detached(&repo, &wt, &base).await?;
1071 }
1072 }
1073
1074 if self.state.config.graph.advise {
1082 for k in 1..=self.state.config.graph.advisors {
1083 let wt = root.join(format!("advisor-{k}"));
1084 if !wt.exists() {
1085 git::worktree_add_detached(&repo, &wt, &base).await?;
1086 }
1087 }
1088 }
1089
1090 let authors: Vec<&str> = self
1095 .roles
1096 .implementers
1097 .iter()
1098 .map(|a| a.id.as_str())
1099 .collect();
1100 let overlap: Vec<String> = self
1101 .roles
1102 .judges
1103 .iter()
1104 .enumerate()
1105 .filter(|(_, j)| authors.contains(&j.id.as_str()))
1106 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
1107 .collect();
1108 if !overlap.is_empty() {
1109 let note = format!(
1110 "{} also authored a candidate; blind, but the panel is less \
1111 independent than {} distinct agents would be",
1112 overlap.join(", "),
1113 self.roles.judges.len()
1114 );
1115 self.state.event("prep", note);
1116 }
1117
1118 self.state.event(
1119 "prep",
1120 format!(
1121 "{} candidates, {} judges, base {} ({})",
1122 self.state.candidates.len(),
1123 self.roles.judges.len(),
1124 &self.state.base_commit[..7.min(self.state.base_commit.len())],
1125 self.state.base_branch
1126 ),
1127 );
1128 self.state.status = RunStatus::Implementing;
1129 self.state.save()?;
1130 Ok(())
1131 }
1132
1133 async fn advise(&mut self) -> Result<()> {
1166 let implement_untouched = self
1167 .state
1168 .candidates
1169 .iter()
1170 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
1171 if !self.state.config.graph.advise || self.state.advise_attempted {
1172 return Ok(());
1173 }
1174 if !implement_untouched {
1175 self.state.event(
1176 "advise",
1177 "skipping the design-deliberation stage: at least one \
1178 candidate already shows implementation progress, so this \
1179 run is past the point the stage exists to run before"
1180 .to_owned(),
1181 );
1182 self.state.advise_attempted = true;
1183 self.state.save()?;
1184 return Ok(());
1185 }
1186 let run_id = self.state.id.clone();
1187 let prompts = self.state.config.prompts.clone();
1188 let instruction = self.state.instruction.clone();
1189 let language = self.state.config.graph.language.clone();
1190 let root = self.state.worktree_root();
1191 let n = self.state.config.graph.advisors;
1192 let where_recorded = self.state.dir().join("run.json");
1193
1194 let seats = match self.state.config.advisors() {
1195 Ok(seats) if !seats.is_empty() => seats,
1196 Ok(_) => {
1197 self.state.event(
1198 "advise",
1199 format!(
1200 "[graph] advisors is 0; skipping the design-deliberation \
1201 stage and continuing without a synthesis brief (see {})",
1202 where_recorded.display()
1203 ),
1204 );
1205 self.state.advise_attempted = true;
1206 self.state.save()?;
1207 return Ok(());
1208 }
1209 Err(e) => {
1210 self.state.event(
1211 "advise",
1212 format!(
1213 "could not resolve advisor seats ({e:#}); continuing \
1214 without a design-deliberation brief (see {})",
1215 where_recorded.display()
1216 ),
1217 );
1218 self.state.advise_attempted = true;
1219 self.state.save()?;
1220 return Ok(());
1221 }
1222 };
1223
1224 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1225 let artifacts = agent::artifacts_dir(&self.state.dir());
1226 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
1227
1228 let mut jobs = Vec::new();
1229 for (i, spec) in seats.iter().cloned().enumerate() {
1230 let seat_key = format!("advisor-{}", i + 1);
1231 let seat = self.seat(&seat_key, &spec.id);
1232 jobs.push(SeatJob {
1233 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
1234 spec,
1235 seat,
1236 cwd: worktrees[i % worktrees.len()].clone(),
1237 timeout,
1238 allow_write: false,
1239 sessions: false,
1240 artifacts: artifacts.clone(),
1241 stem: seat_key,
1242 });
1243 }
1244
1245 self.state.event(
1246 "advise",
1247 format!(
1248 "{} advisor seat(s) sketching a design in parallel",
1249 jobs.len()
1250 ),
1251 );
1252 let mut quota_losses = Vec::new();
1253 let cache = self.state.config.cache_dir();
1254 let ctx = WaveCtx {
1255 run: &run_id,
1256 node: "advise",
1257 prompts: &prompts,
1258 cache: cache.as_deref(),
1259 round: None,
1260 };
1261 let results = ask_json_wave::<Proposal>(
1262 jobs,
1263 Arc::clone(&self.sem),
1264 self.state.config.graph.retries,
1265 &ctx,
1266 &mut quota_losses,
1267 &mut self.state,
1268 &|p: &Proposal| p.validate(),
1269 )
1270 .await;
1271 self.state.quota.extend(quota_losses);
1272
1273 let mut records = Vec::with_capacity(results.len());
1274 for (i, (seat, res, _attempts)) in results.into_iter().enumerate() {
1275 let agent_id = seat.agent.clone();
1276 self.state.seats.insert(seat.key.clone(), seat);
1277 match res {
1278 Ok((proposal, out)) => {
1279 self.state
1280 .event("advise", format!("advisor-{} proposed a design", i + 1));
1281 records.push(advise::AdvisorRecord::proposed(
1282 i + 1,
1283 agent_id,
1284 proposal,
1285 out.duration_ms,
1286 ));
1287 }
1288 Err(e) => {
1289 self.state.event(
1290 "advise",
1291 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
1292 );
1293 records.push(advise::AdvisorRecord::failed(
1294 i + 1,
1295 agent_id,
1296 e.to_string(),
1297 ));
1298 }
1299 }
1300 }
1301
1302 let mut advice = advise::Advice {
1303 records,
1304 synthesis: None,
1305 };
1306 if advice.proposals().is_empty() {
1307 self.state.event(
1308 "advise",
1309 "no advisor produced a usable proposal; continuing without a \
1310 synthesis brief"
1311 .to_owned(),
1312 );
1313 } else {
1314 match self
1315 .synthesize_brief(
1316 &advice,
1317 &instruction,
1318 &language,
1319 &worktrees[0],
1320 &artifacts,
1321 &run_id,
1322 &prompts,
1323 cache.as_deref(),
1324 )
1325 .await
1326 {
1327 Ok(Some(text)) => {
1328 self.state.event(
1329 "advise",
1330 "synthesized a design brief for the implementer".to_owned(),
1331 );
1332 advice.synthesis = Some(text);
1333 }
1334 Ok(None) => {
1335 self.state.event(
1336 "advise",
1337 "the synthesis seat produced nothing usable; continuing \
1338 without a design brief"
1339 .to_owned(),
1340 );
1341 }
1342 Err(e) => {
1343 self.state.event(
1344 "advise",
1345 format!("could not synthesize a design brief: {e:#}"),
1346 );
1347 }
1348 }
1349 }
1350 advise::apply_reflection(&mut advice);
1351
1352 self.state.advice = Some(advice);
1353 self.state.advise_attempted = true;
1354 self.state.save()?;
1355 Ok(())
1356 }
1357
1358 #[allow(clippy::too_many_arguments)]
1370 async fn synthesize_brief(
1371 &mut self,
1372 advice: &advise::Advice,
1373 instruction: &str,
1374 language: &str,
1375 cwd: &Path,
1376 artifacts: &Path,
1377 run_id: &str,
1378 prompts: &Prompts,
1379 cache: Option<&Path>,
1380 ) -> Result<Option<String>> {
1381 let want = self.state.config.roles.synthesizer.as_deref();
1382 let spec = agent::pick(&self.state.config.agents, want, &agent::installed)?;
1383 let mut seat = self.seat("advise-synthesis", &spec.id);
1384 let proposals = advice.proposals();
1385 let mut prompt = prompt::with_overlay(
1386 prompt::synthesize_brief(instruction, &proposals, language),
1387 prompts.overlay("advise"),
1388 );
1389 if cache.is_some() {
1390 prompt.push('\n');
1395 prompt.push_str(&prompt::build_cache_note("advise", false));
1396 }
1397 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1398 let out = agent::invoke(
1399 &spec,
1400 &mut seat,
1401 &Invocation {
1402 cwd,
1403 prompt: &prompt,
1404 timeout,
1405 allow_write: false,
1406 sessions: false,
1407 artifacts,
1408 stem: "advise-synthesis",
1409 run: run_id,
1410 node: "advise",
1411 cache_dir: None,
1412 attachments: &[],
1413 },
1414 )
1415 .await?;
1416 self.state.seats.insert(seat.key.clone(), seat);
1417 if !out.usable() {
1418 return Ok(None);
1419 }
1420 let text =
1421 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1422 Ok((!text.trim().is_empty()).then_some(text))
1423 }
1424
1425 async fn implement(&mut self) -> Result<()> {
1428 let run_id = self.state.id.clone();
1433 let prompts = self.state.config.prompts.clone();
1434 let todo: Vec<usize> = self
1435 .state
1436 .candidates
1437 .iter()
1438 .enumerate()
1439 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1440 .map(|(i, _)| i)
1441 .collect();
1442 if todo.is_empty() {
1443 return self.after_implement();
1444 }
1445 self.state.status = RunStatus::Implementing;
1446
1447 let language = self.state.config.graph.language.clone();
1448 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1449 let sessions = self.state.config.graph.sessions;
1450 let artifacts = agent::artifacts_dir(&self.state.dir());
1451 let brief = self
1455 .state
1456 .advice
1457 .as_ref()
1458 .and_then(|a| a.synthesis.as_deref())
1459 .map(str::to_owned);
1460
1461 let mut jobs = Vec::new();
1462 for &i in &todo {
1463 let (index, label, worktree) = {
1464 let c = &self.state.candidates[i];
1465 (c.index, c.label, c.worktree.clone())
1466 };
1467 let spec = self.roles.implementers[index].clone();
1468 let seat_key = format!("impl-{label}");
1469 let seat = self.seat(&seat_key, &spec.id);
1470 let instruction = self.state.instruction.clone();
1471 jobs.push(SeatJob {
1472 spec,
1473 seat,
1474 prompt: prompt::implement(
1475 &instruction,
1476 &worktree.to_string_lossy(),
1477 &language,
1478 brief.as_deref(),
1479 ),
1480 cwd: worktree,
1481 timeout,
1482 allow_write: true,
1483 sessions,
1484 artifacts: artifacts.clone(),
1485 stem: format!("impl-{label}"),
1486 });
1487 }
1488
1489 self.state.event(
1490 "implement",
1491 format!("{} candidates in parallel", jobs.len()),
1492 );
1493 let mut sent = jobs.clone();
1499 let cache = self.state.config.cache_dir();
1500 let ctx = WaveCtx {
1501 run: &run_id,
1502 node: "implement",
1503 prompts: &prompts,
1504 cache: cache.as_deref(),
1505 round: None,
1506 };
1507 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1508 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1509 .await;
1510 self.resume_quota_losses(&mut results, &mut sent, &prompts, &run_id)
1511 .await;
1512 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1513 .await;
1514
1515 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1516 let seat_key = seat.key.clone();
1517 let agent = seat.agent.clone();
1526 let exhausted_the_fallback_chain = matches!(&out, AgentOutcome::Quota(_));
1527 self.state.seats.insert(seat.key.clone(), seat);
1528 let label = self.state.candidates[i].label;
1529 let worktree = self.state.candidates[i].worktree.clone();
1530 let base = self.state.base_commit.clone();
1531
1532 let (summary, duration, failed, verified_claim) = match out {
1533 AgentOutcome::Ok(o) => {
1534 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1535 let failed = (!o.usable()).then(|| {
1536 if o.timed_out {
1537 "agent timed out".to_owned()
1538 } else {
1539 format!("agent exited with {:?}", o.exit_code)
1540 }
1541 });
1542 let verified_claim = verified_noop_claim(failed.is_none(), &o.commands, &text);
1543 (text, o.duration_ms, failed, verified_claim)
1544 }
1545 AgentOutcome::Dropped(o) => {
1551 let why = o
1552 .dropped
1553 .as_ref()
1554 .map(|d| d.why.as_str())
1555 .unwrap_or("the CLI ended the stream without delivering its answer");
1556 (
1557 String::new(),
1558 o.duration_ms,
1559 Some(format!("the CLI dropped the stream ({why})")),
1560 None,
1561 )
1562 }
1563 AgentOutcome::Quota(o) => {
1564 self.state.quota.push(QuotaLoss {
1565 seat: seat_key,
1566 node: "implement".to_owned(),
1567 at: Timestamp::now(),
1568 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1569 });
1570 (
1571 String::new(),
1572 o.duration_ms,
1573 Some("rate limited (quota); produced no change".to_owned()),
1574 None,
1575 )
1576 }
1577 AgentOutcome::Failed(e) => (String::new(), 0, Some(e), None),
1578 };
1579
1580 let rescued = match git::rescue_commit(
1583 &worktree,
1584 &format!("magi: candidate {label} (uncommitted work)"),
1585 )
1586 .await
1587 {
1588 Ok(r) => {
1589 self.state.note_withheld("implement", &r.withheld);
1590 r.committed
1591 }
1592 Err(_) => false,
1593 };
1594 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1595 .await
1596 .unwrap_or(0);
1597 let patch = git::diff(&worktree, &base, "HEAD")
1598 .await
1599 .unwrap_or_default();
1600 let stat = git::diff_stat(&worktree, &base, "HEAD")
1601 .await
1602 .unwrap_or_default();
1603 let files = git::changed_files(&worktree, &base, "HEAD")
1604 .await
1605 .map(|f| f.len())
1606 .unwrap_or(0);
1607 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1608
1609 let c = &mut self.state.candidates[i];
1610 if !exhausted_the_fallback_chain {
1611 c.agent = agent;
1612 }
1613 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1614 c.stat = stat;
1615 c.files = files;
1616 c.commits = commits;
1617 c.duration_ms = duration;
1618 c.empty = commits == 0 || patch.trim().is_empty();
1619 c.failed = match failed {
1622 Some(_) if c.empty => failed,
1623 _ => None,
1624 };
1625 c.verified_noop = if c.empty { verified_claim } else { None };
1630 let note = match (&c.failed, c.empty, &c.verified_noop, rescued) {
1631 (Some(e), _, _, _) => format!("candidate {label}: {e}"),
1632 (None, true, Some(_), _) => {
1633 format!("candidate {label}: no change produced (agent-verified no-op)")
1634 }
1635 (None, true, None, _) => format!("candidate {label}: no change produced"),
1636 (None, false, _, true) => {
1637 format!(
1638 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1639 )
1640 }
1641 (None, false, _, false) => {
1642 format!("candidate {label}: {files} files, {commits} commits")
1643 }
1644 };
1645 self.state.event("implement", note);
1646 self.state.save()?;
1647 }
1648
1649 self.after_implement()
1650 }
1651
1652 async fn resume_undelivered(
1680 &mut self,
1681 results: &mut [(usize, SeatState, AgentOutcome)],
1682 sent: &[SeatJob],
1683 prompts: &Prompts,
1684 run_id: &str,
1685 ) {
1686 for (wi, seat, out) in results.iter_mut() {
1687 let Some(dropped) = (match &*out {
1688 AgentOutcome::Dropped(o) => o.dropped.clone(),
1689 _ => None,
1690 }) else {
1691 continue;
1692 };
1693 let Some(job) = sent.get(*wi) else { continue };
1694 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1696 self.state.event(
1697 "implement",
1698 format!(
1699 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1700 work is in the tree",
1701 seat.key, dropped.output_tokens, dropped.why
1702 ),
1703 );
1704 continue;
1705 }
1706 if !has_context(&job.spec, seat, job.sessions) {
1714 self.state.event(
1715 "implement",
1716 format!(
1717 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1718 is no session left to resume",
1719 seat.key, dropped.output_tokens, dropped.why
1720 ),
1721 );
1722 continue;
1723 }
1724 self.state.event(
1725 "implement",
1726 format!(
1727 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1728 conversation",
1729 seat.key, dropped.output_tokens, dropped.why
1730 ),
1731 );
1732 let mut retry = job.clone();
1733 retry.seat = seat.clone();
1734 retry.prompt = prompt::resume_after_drop(&dropped.why);
1735 retry.timeout = retry_budget(job.timeout, true);
1736 retry.stem = format!("{}-resume", job.stem);
1737 let cache = self.state.config.cache_dir();
1738 let ctx = WaveCtx {
1739 run: run_id,
1740 node: "implement",
1741 prompts,
1742 cache: cache.as_deref(),
1743 round: None,
1744 };
1745 let (resumed_seat, resumed) =
1746 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1747 *seat = resumed_seat;
1748 *out = resumed;
1749 }
1750 }
1751
1752 async fn resume_quota_losses(
1814 &mut self,
1815 results: &mut [(usize, SeatState, AgentOutcome)],
1816 sent: &mut [SeatJob],
1817 prompts: &Prompts,
1818 run_id: &str,
1819 ) {
1820 let instruction = self.state.instruction.clone();
1821 let language = self.state.config.graph.language.clone();
1822 let brief = self
1823 .state
1824 .advice
1825 .as_ref()
1826 .and_then(|a| a.synthesis.as_deref())
1827 .map(str::to_owned);
1828 for (wi, seat, out) in results.iter_mut() {
1829 let Some(job) = sent.get_mut(*wi) else {
1830 continue;
1831 };
1832 let start = self
1837 .roles
1838 .implementer_roster
1839 .iter()
1840 .position(|s| s.id == job.spec.id)
1841 .unwrap_or(0);
1842 let mut tried: BTreeSet<String> = BTreeSet::from([job.spec.id.clone()]);
1843 let mut fallback_attempt = 0usize;
1844 while matches!(&*out, AgentOutcome::Quota(_)) {
1845 let Some(next) =
1846 next_untried_implementer(&self.roles.implementer_roster, start, &tried)
1847 .cloned()
1848 else {
1849 break;
1850 };
1851 tried.insert(next.id.clone());
1852 fallback_attempt += 1;
1853
1854 if let Ok(r) = git::rescue_commit(
1855 &job.cwd,
1856 &format!(
1857 "magi: candidate {} (uncommitted work before quota fallback)",
1858 seat.key
1859 ),
1860 )
1861 .await
1862 {
1863 self.state.note_withheld("implement", &r.withheld);
1864 }
1865
1866 self.state.event(
1867 "implement",
1868 format!(
1869 "{}: rate limited (quota) on {}; retrying with {}",
1870 seat.key, seat.agent, next.id
1871 ),
1872 );
1873
1874 let new_seat = self.seat(&seat.key, &next.id);
1875 job.spec = next.clone();
1883 let mut retry = job.clone();
1884 retry.seat = new_seat;
1885 retry.prompt = prompt::implement(
1886 &instruction,
1887 &job.cwd.to_string_lossy(),
1888 &language,
1889 brief.as_deref(),
1890 );
1891 retry.stem = format!("{}-quota-{}", job.stem, next.id);
1892 let cache = self.state.config.cache_dir();
1893 let ctx = WaveCtx {
1894 run: run_id,
1895 node: "implement",
1896 prompts,
1897 cache: cache.as_deref(),
1898 round: None,
1899 };
1900 let (fallback_seat, fallback_out) = run_one(
1901 retry,
1902 Arc::clone(&self.sem),
1903 &ctx,
1904 &mut self.state,
1905 fallback_attempt,
1906 )
1907 .await;
1908 *seat = fallback_seat;
1909 *out = fallback_out;
1910 }
1911 }
1912 }
1913
1914 async fn resume_unconfirmed_commands(
1938 &mut self,
1939 results: &mut [(usize, SeatState, AgentOutcome)],
1940 sent: &[SeatJob],
1941 prompts: &Prompts,
1942 run_id: &str,
1943 ) {
1944 for (wi, seat, out) in results.iter_mut() {
1945 let AgentOutcome::Ok(o) = &*out else {
1946 continue;
1947 };
1948 if !has_unconfirmed_command(&o.commands) {
1949 continue;
1950 }
1951 let Some(job) = sent.get(*wi) else { continue };
1952 if !has_context(&job.spec, seat, job.sessions) {
1953 self.state.event(
1954 "implement",
1955 format!(
1956 "{}: the reply named a command whose own CLI never confirmed the exit \
1957 status of, but there is no session left to resume",
1958 seat.key
1959 ),
1960 );
1961 continue;
1962 }
1963 self.state.event(
1964 "implement",
1965 format!(
1966 "{}: the reply named a command whose own CLI never confirmed the exit \
1967 status of; resuming the conversation",
1968 seat.key
1969 ),
1970 );
1971 let mut retry = job.clone();
1972 retry.seat = seat.clone();
1973 retry.prompt = prompt::resume_incomplete(
1974 "a command in your last reply had no confirmed exit status",
1975 );
1976 retry.timeout = retry_budget(job.timeout, true);
1977 retry.stem = format!("{}-confirm", job.stem);
1978 let cache = self.state.config.cache_dir();
1979 let ctx = WaveCtx {
1980 run: run_id,
1981 node: "implement",
1982 prompts,
1983 cache: cache.as_deref(),
1984 round: None,
1985 };
1986 let (resumed_seat, resumed) =
1987 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1988 *seat = resumed_seat;
1989 *out = resumed;
1990 }
1991 }
1992
1993 async fn continue_fix_report(
2014 &mut self,
2015 mut seat: SeatState,
2016 parse_err: String,
2017 job: &SeatJob,
2018 prompts: &Prompts,
2019 run_id: &str,
2020 round: usize,
2021 ) -> (
2022 SeatState,
2023 Option<FixReport>,
2024 Option<String>,
2025 ContinuationRecord,
2026 ) {
2027 let mut last_err = parse_err;
2028 let mut cumulative_wait_ms = 0u64;
2029 let mut attempts = 0usize;
2030 loop {
2031 if !has_context(&job.spec, &seat, job.sessions) {
2032 self.state.event(
2033 "fix",
2034 format!(
2035 "round {round}: fixer's reply had no adoption report ({last_err}); no \
2036 session left to resume into"
2037 ),
2038 );
2039 let outcome = if attempts == 0 {
2040 ContinuationOutcome::NoSession
2041 } else {
2042 ContinuationOutcome::Exhausted
2043 };
2044 return (
2045 seat,
2046 None,
2047 Some(format!("unparsable fix report: {last_err}")),
2048 ContinuationRecord {
2049 attempts,
2050 cumulative_wait_ms,
2051 outcome,
2052 },
2053 );
2054 }
2055 if attempts >= MAX_FIX_CONTINUATIONS {
2056 self.state.event(
2057 "fix",
2058 format!(
2059 "round {round}: fixer's reply still had no adoption report after \
2060 {attempts} continuation(s) ({last_err}); giving up"
2061 ),
2062 );
2063 return (
2064 seat,
2065 None,
2066 Some(format!(
2067 "unparsable fix report after {attempts} continuation(s): {last_err}"
2068 )),
2069 ContinuationRecord {
2070 attempts,
2071 cumulative_wait_ms,
2072 outcome: ContinuationOutcome::Exhausted,
2073 },
2074 );
2075 }
2076 attempts += 1;
2077 self.state.event(
2078 "fix",
2079 format!(
2080 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
2081 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
2082 ),
2083 );
2084 let mut retry = job.clone();
2085 retry.seat = seat.clone();
2086 retry.prompt = prompt::resume_incomplete(&last_err);
2087 retry.timeout = retry_budget(job.timeout, true);
2088 retry.stem = format!("{}-continue{attempts}", job.stem);
2089 let cache = self.state.config.cache_dir();
2090 let ctx = WaveCtx {
2091 run: run_id,
2092 node: "fix",
2093 prompts,
2094 cache: cache.as_deref(),
2095 round: Some(round),
2096 };
2097 let (resumed_seat, resumed_out) = run_one(
2098 retry,
2099 Arc::clone(&self.sem),
2100 &ctx,
2101 &mut self.state,
2102 attempts,
2103 )
2104 .await;
2105 seat = resumed_seat;
2106 match resumed_out {
2107 AgentOutcome::Ok(o) => {
2108 cumulative_wait_ms += o.duration_ms;
2109 match verdict::extract_json::<FixReport>(&o.text) {
2110 Ok(report) if !has_unconfirmed_command(&o.commands) => {
2111 self.state.event(
2112 "fix",
2113 format!(
2114 "round {round}: fixer's adoption report recovered after \
2115 {attempts} continuation(s)"
2116 ),
2117 );
2118 return (
2119 seat,
2120 Some(report),
2121 None,
2122 ContinuationRecord {
2123 attempts,
2124 cumulative_wait_ms,
2125 outcome: ContinuationOutcome::Resumed,
2126 },
2127 );
2128 }
2129 Ok(_) => {
2137 last_err = "the reply parsed, but it reported a command whose own CLI \
2138 never confirmed an exit status"
2139 .to_owned();
2140 }
2141 Err(e) => last_err = e.to_string(),
2142 }
2143 }
2144 AgentOutcome::Quota(o) => {
2145 cumulative_wait_ms += o.duration_ms;
2146 self.state.quota.push(QuotaLoss {
2147 seat: seat.key.clone(),
2148 node: "fix".to_owned(),
2149 at: Timestamp::now(),
2150 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2151 });
2152 self.state.event(
2153 "fix",
2154 format!(
2155 "round {round}: continuation rate limited (quota); not retrying now"
2156 ),
2157 );
2158 return (
2159 seat,
2160 None,
2161 Some("rate limited (quota) while recovering the fix report".to_owned()),
2162 ContinuationRecord {
2163 attempts,
2164 cumulative_wait_ms,
2165 outcome: ContinuationOutcome::QuotaLost,
2166 },
2167 );
2168 }
2169 AgentOutcome::Dropped(o) => {
2170 cumulative_wait_ms += o.duration_ms;
2171 let why = o
2172 .dropped
2173 .as_ref()
2174 .map(|d| d.why.as_str())
2175 .unwrap_or("the CLI ended the stream without delivering its answer");
2176 last_err = format!("the CLI dropped the stream ({why})");
2177 }
2178 AgentOutcome::Failed(e) => last_err = e,
2179 }
2180 }
2181 }
2182
2183 fn after_implement(&mut self) -> Result<()> {
2184 if self.state.leaks.is_empty() {
2186 let cfg = self.state.config.blind.clone();
2187 let mut leaks = Vec::new();
2188 for c in &self.state.candidates {
2189 let Some(patch) =
2190 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
2191 else {
2192 continue;
2193 };
2194 leaks.extend(blind::scan(
2195 &format!("candidate {} patch", c.label),
2196 &patch,
2197 &cfg.vendor_tokens,
2198 ));
2199 }
2200 if !leaks.is_empty() {
2201 let summary = leaks
2202 .iter()
2203 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
2204 .collect::<Vec<_>>()
2205 .join(", ");
2206 match cfg.on_leak {
2207 LeakPolicy::Fail => {
2208 self.state.status = RunStatus::Failed;
2209 self.state
2210 .event("blind", format!("vendor text in a patch: {summary}"));
2211 self.state.leaks = leaks;
2212 self.state.save()?;
2213 self.settle_questions();
2214 bail!(
2215 "blind.on_leak = \"fail\" and vendor text reached a \
2216 judged patch: {summary}"
2217 );
2218 }
2219 LeakPolicy::Redact => self.state.event(
2220 "blind",
2221 format!("redacting vendor text for judging: {summary}"),
2222 ),
2223 LeakPolicy::Warn => self.state.event(
2224 "blind",
2225 format!("vendor text present in a judged patch (shown as-is): {summary}"),
2226 ),
2227 }
2228 self.state.leaks = leaks;
2229 }
2230 }
2231
2232 if self.state.viable().is_empty() {
2233 if self.state.all_candidates_verified_noop() {
2234 self.state.status = RunStatus::VerifiedNoop;
2245 self.state.save()?;
2246 self.settle_questions();
2247 return Ok(());
2248 }
2249 self.state.status = RunStatus::Failed;
2250 self.state.save()?;
2251 self.settle_questions();
2252 bail!("no candidate produced a change; nothing to judge");
2253 }
2254 self.state.status = RunStatus::Judging;
2255 self.state.save()?;
2256 Ok(())
2257 }
2258
2259 async fn judge(&mut self) -> Result<()> {
2262 let run_id = self.state.id.clone();
2267 let prompts = self.state.config.prompts.clone();
2268 if !self.state.judgements.is_empty() || self.state.judge_skipped {
2269 return Ok(());
2270 }
2271 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2272 if viable.len() == 1 {
2273 self.state.judge_skipped = true;
2280 self.state.event(
2281 "judge",
2282 format!(
2283 "only candidate {} produced a change; judging skipped",
2284 viable[0].label
2285 ),
2286 );
2287 self.state.save()?;
2288 return Ok(());
2289 }
2290 self.state.status = RunStatus::Judging;
2291
2292 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2293 let language = self.state.config.graph.language.clone();
2294 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2295 let sessions = self.state.config.graph.sessions;
2296 let artifacts = agent::artifacts_dir(&self.state.dir());
2297 let root = self.state.worktree_root();
2298 let base_short = short(&self.state.base_commit);
2299
2300 let mut jobs = Vec::new();
2301 let mut orders = Vec::new();
2302 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2303 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2304 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2305 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
2306 let seat_key = format!("judge-{}", j + 1);
2307 let seat = self.seat(&seat_key, &spec.id);
2308 jobs.push(SeatJob {
2309 prompt: prompt::judge(
2310 &self.state.instruction,
2311 &views,
2312 self.roles.judges.len(),
2313 &base_short,
2314 &language,
2315 ),
2316 spec,
2317 seat,
2318 cwd: root.join(format!("judge-{}", j + 1)),
2319 timeout,
2320 allow_write: false,
2321 sessions,
2322 artifacts: artifacts.clone(),
2323 stem: format!("judge-{}", j + 1),
2324 });
2325 }
2326
2327 self.state.event(
2328 "judge",
2329 format!(
2330 "{} judges ranking {} candidates blind",
2331 jobs.len(),
2332 viable.len()
2333 ),
2334 );
2335 let labels_for_check = labels.clone();
2336 let mut quota_losses = Vec::new();
2337 let cache = self.state.config.cache_dir();
2338 let ctx = WaveCtx {
2339 run: &run_id,
2340 node: "judge",
2341 prompts: &prompts,
2342 cache: cache.as_deref(),
2343 round: None,
2344 };
2345 let results = ask_json_wave::<Ranking>(
2346 jobs,
2347 Arc::clone(&self.sem),
2348 self.state.config.graph.retries,
2349 &ctx,
2350 &mut quota_losses,
2351 &mut self.state,
2352 &move |r: &Ranking| r.validate(&labels_for_check),
2353 )
2354 .await;
2355 self.state.quota.extend(quota_losses);
2356
2357 for (j, (seat, res, _attempts)) in results.into_iter().enumerate() {
2358 let agent_id = seat.agent.clone();
2359 self.state.seats.insert(seat.key.clone(), seat);
2360 let mut record = Judgement {
2361 judge: j + 1,
2362 seat: format!("judge-{}", j + 1),
2363 agent: agent_id,
2364 ranking: Vec::new(),
2365 reasons: BTreeMap::new(),
2366 confidence: None,
2367 order: orders[j].clone(),
2368 failed: None,
2369 duration_ms: 0,
2370 };
2371 match res {
2372 Ok((ranking, out)) => {
2373 record.ranking = ranking.normalized();
2374 record.reasons = ranking.reasons;
2375 record.confidence = ranking.confidence;
2376 record.duration_ms = out.duration_ms;
2377 self.state.event(
2378 "judge",
2379 format!(
2380 "judge {} ranked {}",
2381 j + 1,
2382 record.ranking.iter().collect::<String>()
2383 ),
2384 );
2385 }
2386 Err(e) => {
2387 record.failed = Some(e.to_string());
2388 self.state
2389 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
2390 }
2391 }
2392 self.state.judgements.push(record);
2393 self.state.save()?;
2394 }
2395 Ok(())
2396 }
2397
2398 async fn deliberate(&mut self) -> Result<()> {
2401 let run_id = self.state.id.clone();
2406 let prompts = self.state.config.prompts.clone();
2407 if !self.state.deliberation.is_empty() {
2408 return Ok(());
2409 }
2410 let tops: Vec<char> = self
2411 .state
2412 .judgements
2413 .iter()
2414 .filter_map(|j| j.ranking.first().copied())
2415 .collect();
2416 let rounds = self.state.config.graph.deliberate_rounds;
2417 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
2418 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
2419 self.state.event(
2420 "deliberate",
2421 format!("judges agreed on {} outright; no deliberation", tops[0]),
2422 );
2423 }
2424 self.state.status = RunStatus::Voting;
2425 self.state.save()?;
2426 return Ok(());
2427 }
2428
2429 self.state.status = RunStatus::Deliberating;
2430 self.state.event(
2431 "deliberate",
2432 format!(
2433 "split: first choices were {} — opening {rounds} round(s)",
2434 tops.iter().collect::<String>()
2435 ),
2436 );
2437
2438 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2439 let language = self.state.config.graph.language.clone();
2440 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2441 let sessions = self.state.config.graph.sessions;
2442 let artifacts = agent::artifacts_dir(&self.state.dir());
2443 let root = self.state.worktree_root();
2444 let base_short = short(&self.state.base_commit);
2445
2446 for round in 1..=rounds {
2450 let mut turns: Vec<DeliberationTurn> = Vec::new();
2451 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2452 if self.state.judgements[j].failed.is_some() {
2453 continue;
2454 }
2455 let seat_key = format!("judge-{}", j + 1);
2456 let mut seat = self.seat(&seat_key, &spec.id);
2457 let transcript = self.transcript(&turns, j);
2458 let context = if has_context(&spec, &seat, sessions) {
2459 None
2460 } else {
2461 Some(self.candidate_block(&viable, &base_short))
2462 };
2463 let text = prompt::deliberate(
2464 &self.state.instruction,
2465 context.as_deref(),
2466 &transcript,
2467 round,
2468 rounds,
2469 &language,
2470 );
2471 let job = SeatJob {
2472 spec,
2473 seat: seat.clone(),
2474 prompt: text,
2475 cwd: root.join(format!("judge-{}", j + 1)),
2476 timeout,
2477 allow_write: false,
2478 sessions,
2479 artifacts: artifacts.clone(),
2480 stem: format!("delib-{round}-judge-{}", j + 1),
2481 };
2482 let cache = self.state.config.cache_dir();
2483 let ctx = WaveCtx {
2484 run: &run_id,
2485 node: "deliberate",
2486 prompts: &prompts,
2487 cache: cache.as_deref(),
2488 round: None,
2489 };
2490 let (updated, out) =
2491 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2492 seat = updated;
2493 let agent_id = seat.agent.clone();
2494 let seat_key = seat.key.clone();
2495 self.state.seats.insert(seat.key.clone(), seat);
2496 let body = match out {
2497 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
2498 AgentOutcome::Dropped(o) => {
2502 let why =
2503 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
2504 "the CLI ended the stream without delivering its answer",
2505 );
2506 self.state.event(
2507 "deliberate",
2508 format!(
2509 "judge {} skipped: the CLI dropped the stream ({why})",
2510 j + 1
2511 ),
2512 );
2513 continue;
2514 }
2515 AgentOutcome::Quota(o) => {
2516 self.state.quota.push(QuotaLoss {
2517 seat: seat_key,
2518 node: "deliberate".to_owned(),
2519 at: Timestamp::now(),
2520 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2521 });
2522 self.state.event(
2523 "deliberate",
2524 format!("judge {} skipped: rate limited (quota)", j + 1),
2525 );
2526 continue;
2527 }
2528 AgentOutcome::Failed(e) => {
2529 self.state
2530 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
2531 continue;
2532 }
2533 };
2534 let tentative = verdict::extract_json::<Position>(&body)
2535 .ok()
2536 .and_then(|p| p.tentative)
2537 .and_then(|s| s.trim().chars().next())
2538 .map(|c| c.to_ascii_uppercase());
2539 self.state.event(
2540 "deliberate",
2541 format!(
2542 "round {round}: judge {} now favours {}",
2543 j + 1,
2544 tentative.map_or("—".to_owned(), |c| c.to_string())
2545 ),
2546 );
2547 turns.push(DeliberationTurn {
2548 judge: j + 1,
2549 agent: agent_id,
2550 body: blind::sanitize_prose(&body, &self.state.config.blind),
2551 tentative,
2552 });
2553 }
2554 self.state
2555 .deliberation
2556 .push(DeliberationRound { round, turns });
2557 self.state.save()?;
2558 }
2559
2560 self.state.status = RunStatus::Voting;
2561 self.state.save()?;
2562 Ok(())
2563 }
2564
2565 async fn vote(&mut self) -> Result<()> {
2568 let run_id = self.state.id.clone();
2573 let prompts = self.state.config.prompts.clone();
2574 if !self.state.votes.is_empty() {
2575 return Ok(());
2576 }
2577 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2578 if viable.len() == 1 {
2579 return Ok(());
2580 }
2581 self.state.status = RunStatus::Voting;
2582
2583 let language = self.state.config.graph.language.clone();
2584 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2585 let sessions = self.state.config.graph.sessions;
2586 let artifacts = agent::artifacts_dir(&self.state.dir());
2587 let root = self.state.worktree_root();
2588 let base_short = short(&self.state.base_commit);
2589 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2590
2591 let mut jobs = Vec::new();
2592 let mut seats_at = Vec::new();
2593 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2594 if self
2595 .state
2596 .judgements
2597 .get(j)
2598 .is_some_and(|r| r.failed.is_some())
2599 {
2600 continue;
2601 }
2602 let seat_key = format!("judge-{}", j + 1);
2603 let seat = self.seat(&seat_key, &spec.id);
2604 let mut text = prompt::final_vote(&viable, &language);
2605 if !has_context(&spec, &seat, sessions) {
2606 text = format!(
2607 "{}\n\n# Candidates\n\n{}",
2608 text,
2609 self.candidate_block(&candidates, &base_short)
2610 );
2611 }
2612 jobs.push(SeatJob {
2613 spec,
2614 seat,
2615 prompt: text,
2616 cwd: root.join(format!("judge-{}", j + 1)),
2617 timeout,
2618 allow_write: false,
2619 sessions,
2620 artifacts: artifacts.clone(),
2621 stem: format!("vote-judge-{}", j + 1),
2622 });
2623 seats_at.push(j);
2624 }
2625
2626 self.state.event(
2627 "vote",
2628 format!(
2629 "collecting {} final votes one by one, privately",
2630 jobs.len()
2631 ),
2632 );
2633 let allowed = viable.clone();
2634 let mut quota_losses = Vec::new();
2635 let cache = self.state.config.cache_dir();
2636 let ctx = WaveCtx {
2637 run: &run_id,
2638 node: "vote",
2639 prompts: &prompts,
2640 cache: cache.as_deref(),
2641 round: None,
2642 };
2643 let results = ask_json_wave::<FinalVote>(
2644 jobs,
2645 Arc::clone(&self.sem),
2646 self.state.config.graph.retries,
2647 &ctx,
2648 &mut quota_losses,
2649 &mut self.state,
2650 &move |v: &FinalVote| match v.label() {
2651 Some(c) if allowed.contains(&c) => Ok(()),
2652 other => bail!("vote {other:?} is not one of {allowed:?}"),
2653 },
2654 )
2655 .await;
2656 self.state.quota.extend(quota_losses);
2657
2658 for (&j, (seat, res, _attempts)) in seats_at.iter().zip(results) {
2659 let agent_id = seat.agent.clone();
2660 self.state.seats.insert(seat.key.clone(), seat);
2661 let initial = self
2662 .state
2663 .judgements
2664 .get(j)
2665 .and_then(|r| r.ranking.first().copied());
2666 let mut record = VoteRecord {
2667 judge: j + 1,
2668 agent: agent_id,
2669 vote: None,
2670 reason: String::new(),
2671 changed: false,
2672 };
2673 match res {
2674 Ok((v, _)) => {
2675 record.vote = v.label();
2676 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2677 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2678 self.state.event(
2679 "vote",
2680 format!(
2681 "judge {} voted {}{}",
2682 j + 1,
2683 record.vote.unwrap_or('?'),
2684 if record.changed { " (changed)" } else { "" }
2685 ),
2686 );
2687 }
2688 Err(e) => {
2689 self.state
2690 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2691 }
2692 }
2693 self.state.votes.push(record);
2694 self.state.save()?;
2695 }
2696 Ok(())
2697 }
2698
2699 fn tally(&mut self) -> Result<()> {
2702 if self.state.tally.is_some() {
2703 return Ok(());
2704 }
2705 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2706 let tops: Vec<char> = self
2707 .state
2708 .judgements
2709 .iter()
2710 .filter_map(|j| j.ranking.first().copied())
2711 .collect();
2712 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2713
2714 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2717 let mut cast: Vec<char> = Vec::new();
2718 for (i, j) in self.state.judgements.iter().enumerate() {
2719 let vote = self
2720 .state
2721 .votes
2722 .iter()
2723 .find(|v| v.judge == i + 1)
2724 .and_then(|v| v.vote)
2725 .or_else(|| j.ranking.first().copied());
2726 if let Some(v) = vote {
2727 *first_choice.entry(v).or_insert(0) += 1;
2728 cast.push(v);
2729 }
2730 }
2731
2732 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2733 for j in &self.state.judgements {
2734 let n = j.ranking.len();
2735 for (pos, label) in j.ranking.iter().enumerate() {
2736 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2737 }
2738 }
2739
2740 let best = first_choice.values().copied().max().unwrap_or(0);
2741 let mut leaders: Vec<char> = first_choice
2742 .iter()
2743 .filter(|(_, v)| **v == best)
2744 .map(|(k, _)| *k)
2745 .collect();
2746 let mut tie_break = None;
2747 if leaders.len() > 1 {
2748 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2749 let borda_leaders: Vec<char> = leaders
2750 .iter()
2751 .copied()
2752 .filter(|l| borda[l] == top_borda)
2753 .collect();
2754 tie_break = Some(if borda_leaders.len() == 1 {
2755 format!(
2756 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2757 leaders.len()
2758 )
2759 } else {
2760 format!(
2761 "{} way tie on both first-choice votes and Borda points, broken by label order",
2762 leaders.len()
2763 )
2764 });
2765 leaders = borda_leaders;
2766 leaders.sort_unstable();
2767 }
2768 let winner = *leaders
2769 .first()
2770 .or(viable.first())
2771 .context("no candidate to declare a winner from")?;
2772
2773 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2774 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2775 let deliberated = !self.state.deliberation.is_empty();
2776
2777 let quota_seats: std::collections::BTreeSet<&str> =
2781 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2782 let mut present = 0usize;
2783 for (i, j) in self.state.judgements.iter().enumerate() {
2784 if quota_seats.contains(j.seat.as_str()) {
2785 continue;
2786 }
2787 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2788 let voted = self
2789 .state
2790 .votes
2791 .iter()
2792 .any(|v| v.judge == i + 1 && v.vote.is_some());
2793 if ranked || voted {
2794 present += 1;
2795 }
2796 }
2797 let needs_quorum = viable.len() > 1;
2803 let judges_total = if needs_quorum {
2804 self.roles.judges.len()
2805 } else {
2806 0
2807 };
2808 let quorum = if needs_quorum {
2809 judges_total / 2 + 1
2810 } else {
2811 0
2812 };
2813 let met_quorum = !needs_quorum || present >= quorum;
2814 let uncontested = (!needs_quorum).then(|| {
2815 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2816 });
2817
2818 self.state.event(
2819 "tally",
2820 match &uncontested {
2821 Some(reason) => format!("winner {winner} — {reason}"),
2822 None => format!(
2823 "winner {winner} — votes {} | initial {} | {} changed | \
2824 {present}/{judges_total} judges{}",
2825 first_choice
2826 .iter()
2827 .map(|(k, v)| format!("{k}:{v}"))
2828 .collect::<Vec<_>>()
2829 .join(" "),
2830 if unanimous_initial {
2831 "unanimous"
2832 } else {
2833 "split"
2834 },
2835 changed_votes,
2836 if met_quorum {
2837 String::new()
2838 } else {
2839 format!(" — below quorum ({quorum} required)")
2840 },
2841 ),
2842 },
2843 );
2844 if !met_quorum {
2845 self.state.event(
2846 "stall",
2847 format!(
2848 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2849 the run stops here, resumable"
2850 ),
2851 );
2852 }
2853 self.state.tally = Some(Tally {
2854 first_choice,
2855 borda,
2856 winner,
2857 rankings: tops.len(),
2858 unanimous_initial,
2859 deliberated,
2860 changed_votes,
2861 unanimous_final,
2862 tie_break,
2863 judges: judges_total,
2864 present,
2865 quorum,
2866 met_quorum,
2867 uncontested,
2868 });
2869 self.state.status = if met_quorum {
2870 RunStatus::Reviewing
2871 } else {
2872 RunStatus::Stalled
2873 };
2874 self.state.save()?;
2875 Ok(())
2876 }
2877
2878 #[allow(clippy::too_many_lines)]
2899 async fn recover_stall(&mut self) -> Result<bool> {
2900 let run_id = self.state.id.clone();
2905 let prompts = self.state.config.prompts.clone();
2906 let quota_seats: BTreeSet<&str> =
2911 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2912 let absent: Vec<String> = self
2913 .state
2914 .judgements
2915 .iter()
2916 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2917 .map(|j| j.seat.clone())
2918 .collect();
2919 if absent.is_empty() {
2920 return Ok(false);
2921 }
2922 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2923 if viable.len() <= 1 {
2924 return Ok(false);
2925 }
2926 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2927 let language = self.state.config.graph.language.clone();
2928 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2929 let sessions = self.state.config.graph.sessions;
2930 let artifacts = agent::artifacts_dir(&self.state.dir());
2931 let root = self.state.worktree_root();
2932 let base_short = short(&self.state.base_commit);
2933 let candidates: Vec<Candidate> = viable.clone();
2934
2935 let mut positions: Vec<usize> = absent
2937 .iter()
2938 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2939 .collect();
2940 if positions.is_empty() {
2941 return Ok(false);
2942 }
2943 positions.sort_unstable();
2944 positions.dedup();
2945
2946 let mut judge_jobs = Vec::new();
2948 for &j in &positions {
2949 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2950 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2951 let seat_key = format!("judge-{}", j + 1);
2952 let spec = self.roles.judges[j].clone();
2953 let seat = self.seat(&seat_key, &spec.id);
2954 judge_jobs.push(SeatJob {
2955 spec,
2956 seat,
2957 prompt: prompt::judge(
2958 &self.state.instruction,
2959 &views,
2960 self.roles.judges.len(),
2961 &base_short,
2962 &language,
2963 ),
2964 cwd: root.join(seat_key),
2965 timeout,
2966 allow_write: false,
2967 sessions,
2968 artifacts: artifacts.clone(),
2969 stem: format!("judge-{}-recover", j + 1),
2970 });
2971 }
2972
2973 let labels_for_check = labels.clone();
2974 let mut judge_losses = Vec::new();
2975 let retries = self.state.config.graph.retries;
2976 let cache = self.state.config.cache_dir();
2977 let ctx = WaveCtx {
2978 run: &run_id,
2979 node: "judge",
2980 prompts: &prompts,
2981 cache: cache.as_deref(),
2982 round: None,
2983 };
2984 let results = ask_json_wave::<Ranking>(
2985 judge_jobs,
2986 Arc::clone(&self.sem),
2987 retries,
2988 &ctx,
2989 &mut judge_losses,
2990 &mut self.state,
2991 &move |r: &Ranking| r.validate(&labels_for_check),
2992 )
2993 .await;
2994
2995 let mut recovered: BTreeSet<usize> = BTreeSet::new();
2997 for (&j, (seat, res, _attempts)) in positions.iter().zip(results) {
2998 self.state.seats.insert(seat.key.clone(), seat);
2999 let record = &mut self.state.judgements[j];
3000 match res {
3001 Ok((ranking, out)) => {
3002 record.ranking = ranking.normalized();
3003 record.reasons = ranking.reasons;
3004 record.confidence = ranking.confidence;
3005 record.failed = None;
3006 record.duration_ms = out.duration_ms;
3007 recovered.insert(j);
3008 self.state.event(
3009 "recover",
3010 format!("judge {} ranked again after the limit", j + 1),
3011 );
3012 }
3013 Err(e) => {
3014 self.state
3015 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
3016 }
3017 }
3018 }
3019
3020 let mut vote_jobs = Vec::new();
3022 let mut vote_pos: Vec<usize> = Vec::new();
3023 for &j in &recovered {
3024 let seat_key = format!("judge-{}", j + 1);
3025 let spec = self.roles.judges[j].clone();
3026 let seat = self.seat(&seat_key, &spec.id);
3027 let mut text = prompt::final_vote(&labels, &language);
3028 if !has_context(&spec, &seat, sessions) {
3029 text = format!(
3030 "{}\n\n# Candidates\n\n{}",
3031 text,
3032 self.candidate_block(&candidates, &base_short)
3033 );
3034 }
3035 vote_jobs.push(SeatJob {
3036 spec,
3037 seat,
3038 prompt: text,
3039 cwd: root.join(seat_key),
3040 timeout,
3041 allow_write: false,
3042 sessions,
3043 artifacts: artifacts.clone(),
3044 stem: format!("vote-judge-{}-recover", j + 1),
3045 });
3046 vote_pos.push(j);
3047 }
3048 let allowed = labels.clone();
3049 let mut vote_losses = Vec::new();
3050 let vote_retries = self.state.config.graph.retries;
3051 let vote_cache = self.state.config.cache_dir();
3052 let ctx = WaveCtx {
3053 run: &run_id,
3054 node: "vote",
3055 prompts: &prompts,
3056 cache: vote_cache.as_deref(),
3057 round: None,
3058 };
3059 let votes = ask_json_wave::<FinalVote>(
3060 vote_jobs,
3061 Arc::clone(&self.sem),
3062 vote_retries,
3063 &ctx,
3064 &mut vote_losses,
3065 &mut self.state,
3066 &move |v: &FinalVote| match v.label() {
3067 Some(c) if allowed.contains(&c) => Ok(()),
3068 other => bail!("vote {other:?} is not one of {allowed:?}"),
3069 },
3070 )
3071 .await;
3072 for (&j, (seat, res, _attempts)) in vote_pos.iter().zip(votes) {
3073 let agent_id = seat.agent.clone();
3074 self.state.seats.insert(seat.key.clone(), seat);
3075 match res {
3076 Ok((v, _)) => {
3077 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
3078 rec.vote = v.label();
3079 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
3080 } else {
3081 self.state.votes.push(VoteRecord {
3082 judge: j + 1,
3083 agent: agent_id,
3084 vote: v.label(),
3085 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
3086 changed: false,
3087 });
3088 }
3089 self.state.event(
3090 "recover",
3091 format!("judge {} voted again after the limit", j + 1),
3092 );
3093 }
3094 Err(e) => {
3095 self.state
3096 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
3097 }
3098 }
3099 }
3100
3101 let recovered_keys: BTreeSet<String> = recovered
3105 .iter()
3106 .map(|&j| format!("judge-{}", j + 1))
3107 .collect();
3108 self.state
3109 .quota
3110 .retain(|q| !recovered_keys.contains(&q.seat));
3111 for loss in judge_losses.into_iter().chain(vote_losses) {
3115 if recovered_keys.contains(&loss.seat) {
3116 continue;
3117 }
3118 self.state.quota.retain(|q| q.seat != loss.seat);
3119 self.state.quota.push(loss);
3120 }
3121
3122 self.state.tally = None;
3124 self.tally()?;
3125 Ok(self
3126 .state
3127 .tally
3128 .as_ref()
3129 .map(|t| t.met_quorum)
3130 .unwrap_or(false))
3131 }
3132
3133 async fn fold_losers(&mut self) -> Result<()> {
3136 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
3137 return Ok(());
3138 };
3139 let repo = self.state.repo.clone();
3140 let mut folded = Vec::new();
3141 for i in 0..self.state.candidates.len() {
3142 let c = &self.state.candidates[i];
3143 if c.label == winner || c.folded {
3144 continue;
3145 }
3146 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
3147 git::worktree_remove(&repo, &wt).await.ok();
3148 git::branch_delete(&repo, &branch).await.ok();
3149 self.state.candidates[i].folded = true;
3150 folded.push(label.to_string());
3151 }
3152 let root = self.state.worktree_root();
3154 for j in 1..=self.roles.judges.len() {
3155 let wt = root.join(format!("judge-{j}"));
3156 if wt.exists() {
3157 git::worktree_remove(&repo, &wt).await.ok();
3158 }
3159 }
3160 if self.state.config.graph.advise {
3163 for k in 1..=self.state.config.graph.advisors {
3164 let wt = root.join(format!("advisor-{k}"));
3165 if wt.exists() {
3166 git::worktree_remove(&repo, &wt).await.ok();
3167 }
3168 }
3169 }
3170 if !folded.is_empty() {
3171 self.state
3172 .event("fold", format!("folded candidates {}", folded.join(", ")));
3173 self.state.save()?;
3174 }
3175 Ok(())
3176 }
3177
3178 async fn sync_to_base(&mut self) -> Result<()> {
3208 if self
3209 .state
3210 .base_sync
3211 .as_ref()
3212 .is_some_and(|s| s.conflict.is_some())
3213 {
3214 return Ok(());
3215 }
3216 let Some(winner) = self.state.winner().cloned() else {
3217 return Ok(());
3218 };
3219
3220 let repo = self.state.repo.clone();
3221 let remote = self.state.config.merge.remote.clone();
3222 let base_branch = self.state.base_branch.clone();
3223 let tracking = format!("{remote}/{base_branch}");
3224
3225 git::fetch(&repo, &remote, &base_branch).await.ok();
3226 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
3230 return Ok(());
3231 };
3232
3233 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3234 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
3235 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
3236
3237 if behind == 0 {
3238 self.state.base_sync = Some(BaseSync {
3239 tip,
3240 behind: 0,
3241 attempts,
3242 conflict: None,
3243 });
3244 self.state.save()?;
3245 return Ok(());
3246 }
3247
3248 if attempts >= BASE_SYNC_ROUNDS {
3249 let why = format!(
3250 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
3251 rebase(s); rebasing again would only race it",
3252 winner.branch
3253 );
3254 self.state.status = RunStatus::Blocked;
3255 self.state.base_sync = Some(BaseSync {
3256 tip,
3257 behind,
3258 attempts,
3259 conflict: Some(why.clone()),
3260 });
3261 self.state.event("land", why);
3262 self.state.save()?;
3263 return Ok(());
3264 }
3265
3266 self.state.event(
3267 "land",
3268 format!(
3269 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
3270 winner.branch
3271 ),
3272 );
3273 self.state.save()?;
3274
3275 let branch_tracking = format!("{remote}/{}", winner.branch);
3281 let fetched_branch = git::fetch(&repo, &remote, &winner.branch).await;
3282 let remote_tip = if matches!(&fetched_branch, Ok(o) if o.ok()) {
3283 git::rev_parse(&repo, &branch_tracking).await.ok()
3284 } else {
3285 None
3286 };
3287 if let Some(theirs) = &remote_tip
3291 && !git::is_ancestor(&repo, theirs, &head).await
3292 && !crate::reconcile::origin_missing(&repo, &head, theirs)
3293 .await
3294 .is_ok_and(|missing| missing.is_empty())
3295 {
3296 let why = format!(
3297 "{branch_tracking} ({}) has commits {} does not contain; not rebasing over \
3298 them",
3299 short(theirs),
3300 winner.branch
3301 );
3302 self.state.status = RunStatus::Blocked;
3303 self.state.base_sync = Some(BaseSync {
3304 tip,
3305 behind,
3306 attempts,
3307 conflict: Some(why.clone()),
3308 });
3309 self.state.event("land", why);
3310 self.state.save()?;
3311 return Ok(());
3312 }
3313
3314 let scratch = self.state.dir().join("base-sync");
3315 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
3316 let attempts = attempts + 1;
3317 match rebased {
3318 Ok(None) => {
3319 git::sync_to_head(&winner.worktree).await?;
3323 let mut conflict = None;
3324 if let Some(pinned) = &remote_tip {
3325 let pushed = git::push_pinned(&repo, &remote, &winner.branch, pinned).await;
3326 match pushed {
3327 Ok(o) if o.ok() => self.state.event(
3328 "land",
3329 format!("pushed rebased {} to {remote}", winner.branch),
3330 ),
3331 Ok(o) => {
3332 conflict = Some(format!(
3333 "rebased {} locally but {remote} refused the push (it moved since {}; someone may have pushed): {}",
3334 winner.branch,
3335 short(pinned),
3336 o.stderr.chars().take(600).collect::<String>()
3337 ));
3338 }
3339 Err(e) => {
3340 conflict = Some(format!(
3341 "rebased {} locally but could not push it: {e:#}",
3342 winner.branch
3343 ));
3344 }
3345 }
3346 }
3347 if let Some(why) = &conflict {
3348 self.state.status = RunStatus::Blocked;
3349 self.state.event("land", why.clone());
3350 }
3351 self.state.base_sync = Some(BaseSync {
3352 tip: tip.clone(),
3353 behind: 0,
3354 attempts,
3355 conflict,
3356 });
3357 self.state
3358 .event("land", format!("rebased {} onto {tracking}", winner.branch));
3359 }
3360 Ok(Some(conflict)) => {
3361 let why = format!(
3362 "{} conflicts with {tracking} and did not rebase: {}",
3363 winner.branch,
3364 conflict.chars().take(600).collect::<String>()
3365 );
3366 self.state.status = RunStatus::Blocked;
3367 self.state.base_sync = Some(BaseSync {
3368 tip,
3369 behind,
3370 attempts,
3371 conflict: Some(why.clone()),
3372 });
3373 self.state.event("land", why);
3374 }
3375 Err(e) => {
3376 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
3377 self.state.status = RunStatus::Blocked;
3378 self.state.base_sync = Some(BaseSync {
3379 tip,
3380 behind,
3381 attempts,
3382 conflict: Some(why.clone()),
3383 });
3384 self.state.event("land", why);
3385 }
3386 }
3387 self.state.save()?;
3388 Ok(())
3389 }
3390
3391 fn landing_base(&self) -> String {
3401 self.state
3402 .base_sync
3403 .as_ref()
3404 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
3405 }
3406
3407 pub async fn fix_selected(
3440 &mut self,
3441 ids: &[String],
3442 reason: &str,
3443 allow_stale: bool,
3444 ) -> Result<()> {
3445 let reason = reason.trim();
3446 if reason.is_empty() {
3447 bail!("a fix request needs a reason — that is the operator's own record of why");
3448 }
3449 if ids.is_empty() {
3450 bail!("no finding id given");
3451 }
3452 if !matches!(self.state.status, RunStatus::Ready | RunStatus::Blocked) {
3453 bail!(
3454 "run {} is `{}`; only a `ready` or `blocked` run — one whose review \
3455 has already concluded — can be given a targeted fix. A run still \
3456 in progress should simply be resumed; a `merged` run's branch has \
3457 already landed, so its answer is a fresh `magi review <branch>`, \
3458 not reopening this run's own record",
3459 self.state.id,
3460 self.state.status.as_str()
3461 );
3462 }
3463 let Some(winner) = self.state.winner().cloned() else {
3464 bail!("run {} has no winning candidate to fix", self.state.id);
3465 };
3466 if !git::branch_exists(&self.state.repo, &winner.branch).await? {
3467 bail!(
3468 "branch `{}` no longer exists; this run cannot be extended",
3469 winner.branch
3470 );
3471 }
3472 let home = crate::run::home();
3473 if crate::daemon::is_working_on(&home, &self.state.id, Timestamp::now()) {
3474 bail!(
3475 "run {} is currently being worked on by another magi process",
3476 self.state.id
3477 );
3478 }
3479 let _claim = FixClaim::acquire(&self.state.dir())?;
3485
3486 let mut seen = BTreeSet::new();
3490 let mut findings = Vec::new();
3491 let mut missing = Vec::new();
3492 for id in ids {
3493 if !seen.insert(id.clone()) {
3494 continue;
3495 }
3496 match self.state.finding(id) {
3497 Some((round, rec, f)) => findings.push(OperatorFixFinding {
3498 id: f.id.clone(),
3499 severity: f.severity,
3500 reviewer_vote: rec.vote,
3501 round: round.round,
3502 round_head: round.head.clone(),
3503 reviewer: rec.reviewer,
3504 agent: rec.agent.clone(),
3505 file: f.file.clone(),
3506 line: f.line,
3507 title: f.title.clone(),
3508 detail: f.detail.clone(),
3509 outcome: OperatorFixOutcome::Pending,
3510 }),
3511 None => missing.push(id.clone()),
3512 }
3513 }
3514 if !missing.is_empty() {
3515 bail!(
3516 "unknown finding id(s): {}; nothing was changed",
3517 missing.join(", ")
3518 );
3519 }
3520
3521 let head_at_request = git::rev_parse(&self.state.repo, &winner.branch).await?;
3522 let stale_details: Vec<(String, String)> = findings
3523 .iter()
3524 .filter(|f| f.round_head != head_at_request)
3525 .map(|f| (f.id.clone(), f.round_head.clone()))
3526 .collect();
3527 let stale = !stale_details.is_empty();
3528 if stale && !allow_stale {
3529 bail!(
3530 "the branch has moved since some finding(s) were raised — {} — now \
3531 at {}; pass --allow-stale to fix anyway, or re-run review first",
3532 stale_details
3533 .iter()
3534 .map(|(id, head)| format!("{id} (raised against {})", short(head)))
3535 .collect::<Vec<_>>()
3536 .join(", "),
3537 short(&head_at_request)
3538 );
3539 }
3540
3541 let request = OperatorFixRequest {
3542 requested_at: Timestamp::now(),
3543 reason: reason.to_owned(),
3544 findings,
3545 head_at_request: head_at_request.clone(),
3546 allow_stale,
3547 stale,
3548 fix: None,
3549 result_head: None,
3550 follow_up_review_run: None,
3551 };
3552 self.state.event(
3553 "fix",
3554 format!(
3555 "operator requested a targeted fix on {} finding(s) ({}): {reason}",
3556 request.findings.len(),
3557 request
3558 .findings
3559 .iter()
3560 .map(|f| f.id.as_str())
3561 .collect::<Vec<_>>()
3562 .join(", "),
3563 ),
3564 );
3565 self.state.operator_fixes.push(request);
3572 self.state.save()?;
3573 let request_index = self.state.operator_fixes.len() - 1;
3574
3575 if winner.worktree.exists() {
3584 let dirty = git::git(
3587 &winner.worktree,
3588 &["status", "--porcelain", "--untracked-files=all"],
3589 )
3590 .await?;
3591 let only_withheld = dirty.lines().all(|l| {
3592 l.strip_prefix("?? ")
3593 .is_some_and(|p| self.state.withheld.iter().any(|w| w.path == p))
3594 });
3595 if !only_withheld {
3596 bail!(
3597 "`{}` has uncommitted changes; refusing to touch it — commit or \
3598 discard them first",
3599 winner.worktree.display()
3600 );
3601 }
3602 git::worktree_remove(&self.state.repo, &winner.worktree)
3603 .await
3604 .ok();
3605 }
3606 let fix_worktree = self.state.worktree_root().join("operator-fix");
3607 let fix_worktree_s = fix_worktree.to_string_lossy().to_string();
3608 git::git(
3609 &self.state.repo,
3610 &["worktree", "add", &fix_worktree_s, winner.branch.as_str()],
3611 )
3612 .await
3613 .with_context(|| format!("checking out `{}` for the fix", winner.branch))?;
3614 if !git::is_clean(&fix_worktree).await? {
3615 git::worktree_remove(&self.state.repo, &fix_worktree)
3616 .await
3617 .ok();
3618 bail!(
3619 "`{}` has uncommitted changes; refusing to start a fix on a dirty tree",
3620 winner.branch
3621 );
3622 }
3623
3624 let run_id = self.state.id.clone();
3625 let prompts = self.state.config.prompts.clone();
3626 let language = self.state.config.graph.language.clone();
3627 let sessions = self.state.config.graph.sessions;
3628 let artifacts = agent::artifacts_dir(&self.state.dir());
3629 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3630 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3631 _ => (
3632 self.state
3633 .config
3634 .agent(&winner.agent)
3635 .cloned()
3636 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3637 format!("impl-{}", winner.label),
3638 ),
3639 };
3640 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3641 let finding_list: Vec<Finding> = self.state.operator_fixes[request_index]
3642 .findings
3643 .iter()
3644 .map(|f| Finding {
3645 id: f.id.clone(),
3646 severity: f.severity,
3647 file: f.file.clone(),
3648 line: f.line,
3649 title: f.title.clone(),
3650 detail: f.detail.clone(),
3651 })
3652 .collect();
3653 let job = SeatJob {
3654 prompt: prompt::operator_fix(
3655 &self.state.instruction,
3656 &finding_list,
3657 reason,
3658 &stale_details,
3659 &head_at_request,
3660 &language,
3661 ),
3662 spec: fix_spec.clone(),
3663 seat,
3664 cwd: fix_worktree.clone(),
3665 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3666 allow_write: true,
3667 sessions,
3668 artifacts: artifacts.clone(),
3669 stem: "operator-fix".to_owned(),
3670 };
3671 let cache = self.state.config.cache_dir();
3672 let ctx = WaveCtx {
3673 run: &run_id,
3674 node: "fix",
3675 prompts: &prompts,
3676 cache: cache.as_deref(),
3677 round: None,
3678 };
3679 let (seat, out) =
3680 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3681 let agent_id = seat.agent.clone();
3682
3683 let mut fix = FixRecord {
3684 agent: agent_id,
3685 addressed: Vec::new(),
3686 rejected: Vec::new(),
3687 notes: String::new(),
3688 committed: false,
3689 failed: None,
3690 duration_ms: 0,
3691 continuation: None,
3692 };
3693 let mut final_seat = seat.clone();
3694 match out {
3695 AgentOutcome::Ok(o) => {
3696 fix.duration_ms = o.duration_ms;
3697 let parsed = verdict::extract_json::<FixReport>(&o.text);
3698 let incomplete_reason = match &parsed {
3699 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3700 "the reply parsed, but it reported a command whose own CLI \
3701 never confirmed an exit status"
3702 .to_owned(),
3703 ),
3704 Ok(_) => None,
3705 Err(e) => Some(e.to_string()),
3706 };
3707 match incomplete_reason {
3708 None => {
3709 let report = parsed.expect("checked Ok above");
3710 fix.addressed = report.addressed;
3711 fix.rejected = report.rejected;
3712 fix.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
3713 }
3714 Some(reason) => {
3715 let (resumed_seat, resolved, failure, cont) = self
3716 .continue_fix_report(seat, reason, &job, &prompts, &run_id, 0)
3717 .await;
3718 fix.duration_ms += cont.cumulative_wait_ms;
3719 fix.continuation = Some(cont);
3720 final_seat = resumed_seat;
3721 match resolved {
3722 Some(report) => {
3723 fix.addressed = report.addressed;
3724 fix.rejected = report.rejected;
3725 fix.notes =
3726 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3727 }
3728 None => fix.failed = failure,
3729 }
3730 }
3731 }
3732 }
3733 AgentOutcome::Dropped(o) => {
3734 fix.duration_ms = o.duration_ms;
3735 let why = o
3736 .dropped
3737 .as_ref()
3738 .map(|d| d.why.as_str())
3739 .unwrap_or("the CLI ended the stream without delivering its answer");
3740 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3741 }
3742 AgentOutcome::Quota(o) => {
3743 self.state.quota.push(QuotaLoss {
3744 seat: final_seat.key.clone(),
3745 node: "fix".to_owned(),
3746 at: Timestamp::now(),
3747 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3748 });
3749 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3750 }
3751 AgentOutcome::Failed(e) => fix.failed = Some(e),
3752 }
3753 if fix.continuation.is_none() {
3754 fix.continuation = Some(ContinuationRecord::not_needed());
3755 }
3756 self.state.seats.insert(final_seat.key.clone(), final_seat);
3757
3758 let rescue_message = format!(
3759 "magi: operator-selected fix ({}) (uncommitted work)",
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 if let Ok(r) = git::rescue_commit(&fix_worktree, &rescue_message).await {
3768 self.state.note_withheld("fix", &r.withheld);
3769 }
3770 let after = git::rev_parse(&fix_worktree, "HEAD").await?;
3771 fix.committed = after != head_at_request;
3772 git::worktree_remove(&self.state.repo, &fix_worktree)
3773 .await
3774 .ok();
3775
3776 self.state.event(
3777 "fix",
3778 match &fix.failed {
3779 Some(reason) => format!(
3780 "operator fix: adoption report was lost ({reason}); {}",
3781 if fix.committed {
3782 "committed"
3783 } else {
3784 "NO new commit"
3785 }
3786 ),
3787 None => format!(
3788 "operator fix: {} addressed, {} rejected, {}",
3789 fix.addressed.len(),
3790 fix.rejected.len(),
3791 if fix.committed {
3792 "committed"
3793 } else {
3794 "NO new commit"
3795 }
3796 ),
3797 },
3798 );
3799
3800 for f in &mut self.state.operator_fixes[request_index].findings {
3807 f.outcome = if fix.failed.is_some() {
3808 OperatorFixOutcome::Unreported
3809 } else if fix.addressed.contains(&f.id) {
3810 OperatorFixOutcome::Addressed
3811 } else if let Some(r) = fix.rejected.iter().find(|r| r.id == f.id) {
3812 OperatorFixOutcome::Rejected { why: r.why.clone() }
3813 } else {
3814 OperatorFixOutcome::Unreported
3815 };
3816 }
3817
3818 let committed = fix.committed;
3819 if committed {
3820 self.state.operator_fixes[request_index].result_head = Some(after.clone());
3821 }
3822 self.state.operator_fixes[request_index].fix = Some(fix);
3823 self.state.save()?;
3826
3827 if committed {
3828 self.state.event(
3829 "fix",
3830 format!(
3831 "operator fix committed {}; opening a follow-up review-only run",
3832 short(&after)
3833 ),
3834 );
3835 match Self::review(&self.state.repo, &winner.branch, self.state.config.clone()).await {
3836 Ok(mut follow_up) => {
3837 follow_up.state.event(
3838 "start",
3839 format!(
3840 "requested by an operator fix on run {} for finding(s) {}",
3841 self.state.id,
3842 self.state.operator_fixes[request_index]
3843 .findings
3844 .iter()
3845 .map(|f| f.id.as_str())
3846 .collect::<Vec<_>>()
3847 .join(", "),
3848 ),
3849 );
3850 follow_up.state.save()?;
3851 let follow_up_id = follow_up.state.id.clone();
3852 if let Err(e) = follow_up.execute().await {
3853 self.state.event(
3854 "fix",
3855 format!(
3856 "follow-up review {follow_up_id} did not complete cleanly: {e:#}"
3857 ),
3858 );
3859 }
3860 self.state.operator_fixes[request_index].follow_up_review_run =
3861 Some(follow_up_id);
3862 }
3863 Err(e) => {
3864 self.state.event(
3865 "fix",
3866 format!("committed the fix but could not open a follow-up review: {e:#}"),
3867 );
3868 }
3869 }
3870 self.state.save()?;
3871 }
3872
3873 Ok(())
3874 }
3875
3876 fn fixer_spec(&self, winner: &Candidate) -> (AgentSpec, String) {
3883 match &self.roles.fixer {
3884 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3885 _ => (
3886 self.state
3887 .config
3888 .agent(&winner.agent)
3889 .cloned()
3890 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3891 format!("impl-{}", winner.label),
3892 ),
3893 }
3894 }
3895
3896 async fn review_loop(&mut self) -> Result<()> {
3897 if self
3902 .state
3903 .base_sync
3904 .as_ref()
3905 .is_some_and(|s| s.conflict.is_some())
3906 {
3907 return Ok(());
3908 }
3909 let run_id = self.state.id.clone();
3914 let prompts = self.state.config.prompts.clone();
3915 let Some(winner) = self.state.winner().cloned() else {
3916 return Ok(());
3917 };
3918 let max_rounds = self.state.config.graph.review_rounds;
3919 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
3929 self.state.status = status;
3930 self.state.save()?;
3931 return Ok(());
3932 }
3933 self.state.status = RunStatus::Reviewing;
3934 if self
3944 .state
3945 .reviews
3946 .last()
3947 .is_some_and(|r| r.e2e_status() == E2eStatus::ResourceBlocked)
3948 {
3949 let shell = self.state.config.shell();
3950 return self
3951 .stop_reviewing(
3952 "the last round's own verification never resolved",
3953 &shell,
3954 &winner.worktree,
3955 )
3956 .await;
3957 }
3958
3959 let repo = self.state.repo.clone();
3960 let root = self.state.worktree_root();
3961 let language = self.state.config.graph.language.clone();
3962 let sessions = self.state.config.graph.sessions;
3963 let artifacts = agent::artifacts_dir(&self.state.dir());
3964 let base = self.landing_base();
3965 let base_short = short(&base);
3966 let reviewers = self.roles.reviewers.clone();
3967 let shell = self.state.config.shell();
3968
3969 for round in (self.state.reviews.len() + 1)..=max_rounds {
3970 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3971 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
3972 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
3973 let prev_verification = self
3982 .state
3983 .reviews
3984 .last()
3985 .and_then(|r| r.verification_summary(&head));
3986
3987 let mut jobs = Vec::new();
3991 for (r, spec) in reviewers.iter().cloned().enumerate() {
3992 let wt = root.join(format!("review-{}", r + 1));
3993 if wt.exists() {
3994 git::reset_detached(&wt, &head).await?;
3995 } else {
3996 git::worktree_add_detached(&repo, &wt, &head).await?;
3997 }
3998 let seat_key = format!("review-{}", r + 1);
3999 let seat = self.seat(&seat_key, &spec.id);
4000 jobs.push(SeatJob {
4001 prompt: prompt::review(&prompt::ReviewCtx {
4002 instruction: &self.state.instruction,
4003 branch: &winner.branch,
4004 base_short: &base_short,
4005 stat: &stat,
4006 patch: &patch,
4007 verification: prev_verification.as_ref(),
4008 reviewers: reviewers.len(),
4009 round,
4010 rounds: max_rounds,
4011 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
4014 lens: Lens::for_seat(r),
4015 language: &language,
4016 }),
4017 spec,
4018 seat,
4019 cwd: wt,
4020 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4021 allow_write: false,
4022 sessions,
4023 artifacts: artifacts.clone(),
4024 stem: format!("review-{round}-{}", r + 1),
4025 });
4026 }
4027
4028 self.state.event(
4029 "review",
4030 format!(
4031 "round {round}: {} reviewers on {}",
4032 jobs.len(),
4033 short(&head)
4034 ),
4035 );
4036 let mut quota_losses = Vec::new();
4037 let review_retries = self.state.config.graph.retries;
4038 let review_cache = self.state.config.cache_dir();
4039 let ctx = WaveCtx {
4040 run: &run_id,
4041 node: "review",
4042 prompts: &prompts,
4043 cache: review_cache.as_deref(),
4044 round: Some(round),
4045 };
4046 let results = ask_json_wave::<Review>(
4047 jobs,
4048 Arc::clone(&self.sem),
4049 review_retries,
4050 &ctx,
4051 &mut quota_losses,
4052 &mut self.state,
4053 &|_: &Review| Ok(()),
4054 )
4055 .await;
4056 let round_quota_missing = quota_losses.len();
4060 self.state.quota.extend(quota_losses);
4061
4062 let mut records = Vec::new();
4063 let mut all_findings = Vec::new();
4064 for (r, (seat, res, attempts)) in results.into_iter().enumerate() {
4065 let agent_id = seat.agent.clone();
4066 self.state.seats.insert(seat.key.clone(), seat);
4067 let mut record = ReviewRecord {
4068 reviewer: r + 1,
4069 agent: agent_id,
4070 summary: String::new(),
4071 findings: Vec::new(),
4072 vote: None,
4073 failed: None,
4074 duration_ms: 0,
4075 attempts,
4081 };
4082 match res {
4083 Ok((review, out)) => {
4084 record.summary =
4092 blind::sanitize_prose(&review.summary, &self.state.config.blind);
4093 record.vote = Some(review.vote);
4094 record.duration_ms = out.duration_ms;
4095 for (n, mut f) in review.findings.into_iter().enumerate() {
4096 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
4099 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
4100 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
4101 f.file = f
4107 .file
4108 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
4109 all_findings.push(f.clone());
4110 record.findings.push(f);
4111 }
4112 self.state.event(
4113 "review",
4114 format!(
4115 "round {round}: reviewer {} voted {} with {} finding(s)",
4116 r + 1,
4117 review.vote.label(),
4118 record.findings.len()
4119 ),
4120 );
4121 }
4122 Err(e) => {
4123 record.failed = Some(e.to_string());
4124 self.state.event(
4125 "review",
4126 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
4127 );
4128 }
4129 }
4130 records.push(record);
4131 }
4132
4133 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
4140 let vote_split =
4141 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
4142 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
4143 if vote_split {
4144 self.state.event(
4145 "review",
4146 format!(
4147 "round {round}: votes split ({}) — one round of reconsideration",
4148 initial_votes
4149 .iter()
4150 .map(|v| v.label())
4151 .collect::<Vec<_>>()
4152 .join(", ")
4153 ),
4154 );
4155 let panel: Vec<ReviewSeatReport<'_>> = records
4158 .iter()
4159 .filter_map(|r| {
4160 r.vote.map(|vote| ReviewSeatReport {
4161 reviewer: r.reviewer,
4162 vote,
4163 summary: &r.summary,
4164 findings: &r.findings,
4165 })
4166 })
4167 .collect();
4168
4169 let mut jobs = Vec::new();
4170 let mut seats_at = Vec::new();
4171 for (r, spec) in reviewers.iter().cloned().enumerate() {
4172 if records[r].vote.is_none() {
4176 continue;
4177 }
4178 let wt = root.join(format!("review-{}", r + 1));
4179 let seat_key = format!("review-{}", r + 1);
4180 let seat = self.seat(&seat_key, &spec.id);
4181 let patch_ctx = if has_context(&spec, &seat, sessions) {
4186 None
4187 } else {
4188 Some(ReviewPatch {
4189 branch: &winner.branch,
4190 base_short: &base_short,
4191 stat: &stat,
4192 patch: &patch,
4193 })
4194 };
4195 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
4196 instruction: &self.state.instruction,
4197 reviewer: r + 1,
4198 lens: Lens::for_seat(r),
4199 panel: &panel,
4200 patch: patch_ctx,
4201 round,
4202 rounds: max_rounds,
4203 language: &language,
4204 });
4205 jobs.push(SeatJob {
4206 prompt,
4207 spec,
4208 seat,
4209 cwd: wt,
4210 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
4211 allow_write: false,
4212 sessions,
4213 artifacts: artifacts.clone(),
4214 stem: format!("review-{round}-reconsider-{}", r + 1),
4215 });
4216 seats_at.push(r);
4217 }
4218
4219 let mut recon_quota_losses = Vec::new();
4220 let recon_cache = self.state.config.cache_dir();
4221 let recon_ctx = WaveCtx {
4222 run: &run_id,
4223 node: "review",
4224 prompts: &prompts,
4225 cache: recon_cache.as_deref(),
4226 round: Some(round),
4227 };
4228 let recon_results = ask_json_wave::<ReviewRevote>(
4229 jobs,
4230 Arc::clone(&self.sem),
4231 review_retries,
4232 &recon_ctx,
4233 &mut recon_quota_losses,
4234 &mut self.state,
4235 &|_: &ReviewRevote| Ok(()),
4236 )
4237 .await;
4238 self.state.quota.extend(recon_quota_losses);
4239
4240 for (&r, (seat, res, _attempts)) in seats_at.iter().zip(recon_results) {
4241 let agent_id = seat.agent.clone();
4242 self.state.seats.insert(seat.key.clone(), seat);
4243 let mut rec = ReviewRevoteRecord {
4244 reviewer: r + 1,
4245 agent: agent_id,
4246 vote: None,
4247 reason: String::new(),
4248 failed: None,
4249 };
4250 match res {
4251 Ok((rv, _)) => {
4252 rec.vote = Some(rv.vote);
4253 rec.reason =
4254 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
4255 self.state.event(
4256 "review",
4257 format!(
4258 "round {round}: reviewer {} revoted {}",
4259 r + 1,
4260 rv.vote.label()
4261 ),
4262 );
4263 }
4264 Err(e) => {
4265 rec.failed = Some(e.to_string());
4266 self.state.event(
4267 "review",
4268 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
4269 );
4270 }
4271 }
4272 reconsideration.push(rec);
4273 }
4274 } else if initial_votes.len() > 1 {
4275 self.state.event(
4276 "review",
4277 format!(
4278 "round {round}: votes agreed ({}) — no reconsideration",
4279 initial_votes[0].label()
4280 ),
4281 );
4282 }
4283
4284 let final_votes: Vec<ReviewVote> = records
4288 .iter()
4289 .filter_map(|r| {
4290 reconsideration
4291 .iter()
4292 .find(|rv| rv.reviewer == r.reviewer)
4293 .and_then(|rv| rv.vote)
4294 .or(r.vote)
4295 })
4296 .collect();
4297 let round_verdict = ReviewVote::worst(final_votes);
4298
4299 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
4300 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4301 let defer_e2e =
4312 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
4313 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
4314 let reason =
4315 format!("{blocking} blocking finding(s) already required a fix this round");
4316 self.state.event(
4317 "verify",
4318 format!(
4319 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
4320 {}); it will run once a round has none left",
4321 short(&head)
4322 ),
4323 );
4324 (Vec::new(), false, true, Some(reason))
4325 } else {
4326 let e2e_commands = self.state.config.verify.e2e.clone();
4327 let cache_dir = self.state.config.cache_dir();
4328 let context = format!("round {round}");
4329 let (e2e, verify_retried) = with_cache_lease(
4330 &mut self.state,
4331 cache_dir.as_deref(),
4332 "e2e",
4333 "e2e",
4334 &winner.worktree,
4335 &head,
4336 verify_timeout,
4337 &context,
4338 |state, budget| {
4339 let shell = shell.clone();
4340 let e2e_commands = e2e_commands.clone();
4341 let worktree = winner.worktree.clone();
4342 let context = context.clone();
4343 async move {
4344 run_e2e_with_retry(
4345 state,
4346 &shell,
4347 &e2e_commands,
4348 &worktree,
4349 budget,
4350 &context,
4351 )
4352 .await
4353 }
4354 },
4355 )
4356 .await;
4357 (e2e, verify_retried, false, None)
4358 };
4359
4360 let expected = records.len();
4361 let answered = records.iter().filter(|r| r.failed.is_none()).count();
4362 let incomplete = answered < expected;
4363 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
4364 let policy = self.state.config.graph.incomplete_review;
4365 let clean = round_is_clean(
4366 blocking,
4367 e2e_ok,
4368 answered,
4369 expected,
4370 round_quota_missing,
4371 policy,
4372 );
4373
4374 let mut round_record = ReviewRound {
4375 round,
4376 head: head.clone(),
4377 verified_head: None,
4378 verified_at: None,
4379 reviews: records,
4380 e2e,
4381 verify_retried,
4382 e2e_deferred,
4383 e2e_defer_reason,
4384 fix: None,
4385 blocking,
4386 answered,
4387 expected,
4388 clean,
4389 progressed: false,
4390 vote_split,
4391 reconsideration,
4392 verdict: round_verdict,
4393 };
4394 if !matches!(
4405 round_record.e2e_status(),
4406 E2eStatus::Deferred | E2eStatus::NotConfigured
4407 ) {
4408 round_record.verified_head = Some(head.clone());
4409 round_record.verified_at = Some(Timestamp::now());
4410 }
4411 let this_round_verification = round_record.verification_summary(&head);
4412
4413 if incomplete {
4414 let missing: Vec<String> = round_record
4415 .reviews
4416 .iter()
4417 .filter(|r| r.failed.is_some())
4418 .map(|r| format!("review-{}", r.reviewer))
4419 .collect();
4420 self.state.event(
4421 "review",
4422 format!(
4423 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
4424 missing.join(", ")
4425 ),
4426 );
4427 }
4428
4429 if clean {
4430 self.state.event(
4431 "review",
4432 if incomplete && policy == IncompleteReviewPolicy::Warn {
4433 format!(
4434 "round {round}: clean (warn policy, incomplete panel) — no \
4435 blocking findings from the seats that answered, verification green"
4436 )
4437 } else if incomplete {
4438 format!(
4439 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
4440 quorum) — no blocking findings from the seats that answered, \
4441 verification green",
4442 expected - answered
4443 )
4444 } else {
4445 format!("round {round}: clean — no blocking findings, verification green")
4446 },
4447 );
4448 self.state.reviews.push(round_record);
4449 self.state.status = RunStatus::Gating;
4450 self.state.save()?;
4451 return Ok(());
4452 }
4453
4454 if incomplete && blocking == 0 && e2e_ok {
4462 self.state.reviews.push(round_record);
4463 self.state.save()?;
4464 if round == max_rounds {
4465 self.state.status = RunStatus::Blocked;
4466 self.state.event(
4467 "review",
4468 format!(
4469 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
4470 refusing to call it clean",
4471 expected - answered
4472 ),
4473 );
4474 return Ok(());
4475 }
4476 continue;
4477 }
4478
4479 if blocking == 0 && round_record.e2e_status() == E2eStatus::ResourceBlocked {
4491 self.state.reviews.push(round_record);
4492 return self
4493 .stop_reviewing(
4494 "the round's own verification could not run",
4495 &shell,
4496 &winner.worktree,
4497 )
4498 .await;
4499 }
4500
4501 if round == max_rounds {
4502 self.state.reviews.push(round_record);
4503 return self
4504 .stop_reviewing(
4505 &format!(
4506 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
4507 ),
4508 &shell,
4509 &winner.worktree,
4510 )
4511 .await;
4512 }
4513
4514 let (fix_spec, fix_seat_key) = self.fixer_spec(&winner);
4517 let seat = self.seat(&fix_seat_key, &fix_spec.id);
4518 let blocking_findings: Vec<_> = all_findings
4519 .iter()
4520 .filter(|f| f.severity.blocks())
4521 .cloned()
4522 .collect();
4523 let job = SeatJob {
4524 prompt: prompt::fix(
4525 &self.state.instruction,
4526 &blocking_findings,
4527 this_round_verification.as_ref(),
4528 round,
4529 max_rounds,
4530 &language,
4531 ),
4532 spec: fix_spec.clone(),
4533 seat,
4534 cwd: winner.worktree.clone(),
4535 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
4536 allow_write: true,
4537 sessions,
4538 artifacts: artifacts.clone(),
4539 stem: format!("fix-{round}"),
4540 };
4541 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
4542 let cache = self.state.config.cache_dir();
4543 let ctx = WaveCtx {
4544 run: &run_id,
4545 node: "fix",
4546 prompts: &prompts,
4547 cache: cache.as_deref(),
4548 round: Some(round),
4549 };
4550 let (seat, out) =
4551 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
4552 let agent_id = seat.agent.clone();
4553
4554 let mut fix = FixRecord {
4555 agent: agent_id,
4556 addressed: Vec::new(),
4557 rejected: Vec::new(),
4558 notes: String::new(),
4559 committed: false,
4560 failed: None,
4561 duration_ms: 0,
4562 continuation: None,
4563 };
4564 let mut continuation = ContinuationRecord::not_needed();
4565 let mut final_seat = seat.clone();
4566 match out {
4567 AgentOutcome::Ok(o) => {
4568 fix.duration_ms = o.duration_ms;
4569 let parsed = verdict::extract_json::<FixReport>(&o.text);
4570 let incomplete_reason = match &parsed {
4577 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
4578 "the reply parsed, but it reported a command whose own CLI never \
4579 confirmed an exit status"
4580 .to_owned(),
4581 ),
4582 Ok(_) => None,
4583 Err(e) => Some(e.to_string()),
4584 };
4585 match incomplete_reason {
4586 None => {
4587 let report = parsed.expect("checked Ok above");
4588 fix.addressed = report.addressed;
4589 fix.rejected = report.rejected;
4590 fix.notes =
4591 blind::sanitize_prose(&report.notes, &self.state.config.blind);
4592 }
4593 Some(reason) => {
4594 let (resumed_seat, resolved, failure, cont) = self
4595 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
4596 .await;
4597 fix.duration_ms += cont.cumulative_wait_ms;
4598 continuation = cont;
4599 final_seat = resumed_seat;
4600 match resolved {
4601 Some(report) => {
4602 fix.addressed = report.addressed;
4603 fix.rejected = report.rejected;
4604 fix.notes = blind::sanitize_prose(
4605 &report.notes,
4606 &self.state.config.blind,
4607 );
4608 }
4609 None => fix.failed = failure,
4610 }
4611 }
4612 }
4613 }
4614 AgentOutcome::Dropped(o) => {
4616 fix.duration_ms = o.duration_ms;
4617 let why = o
4618 .dropped
4619 .as_ref()
4620 .map(|d| d.why.as_str())
4621 .unwrap_or("the CLI ended the stream without delivering its answer");
4622 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
4623 }
4624 AgentOutcome::Quota(o) => {
4625 self.state.quota.push(QuotaLoss {
4626 seat: final_seat.key.clone(),
4627 node: "fix".to_owned(),
4628 at: Timestamp::now(),
4629 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4630 });
4631 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
4632 }
4633 AgentOutcome::Failed(e) => fix.failed = Some(e),
4634 }
4635 fix.continuation = Some(continuation);
4636 self.state.seats.insert(final_seat.key.clone(), final_seat);
4637 if let Ok(r) = git::rescue_commit(
4638 &winner.worktree,
4639 &format!("magi: review round {round} fixes (uncommitted work)"),
4640 )
4641 .await
4642 {
4643 self.state.note_withheld("fix", &r.withheld);
4644 }
4645 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
4646 fix.committed = after != before;
4647 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
4655 let progressed = diff_after != patch;
4656 let commit_note = if fix.committed {
4657 "committed"
4658 } else {
4659 "NO new commit"
4660 };
4661 let tree_note = if progressed {
4662 "changed vs base"
4663 } else {
4664 "unchanged vs base"
4665 };
4666 self.state.event(
4667 "fix",
4668 match &fix.failed {
4669 Some(reason) => {
4675 format!(
4676 "round {round}: fixer's adoption report was lost ({reason}); \
4677 {commit_note}, tree {tree_note}"
4678 )
4679 }
4680 None => format!(
4681 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
4682 {tree_note}{}",
4683 fix.addressed.len(),
4684 fix.rejected.len(),
4685 if continuation.outcome == ContinuationOutcome::Resumed {
4686 format!(
4687 " (adoption report recovered after {} continuation(s))",
4688 continuation.attempts
4689 )
4690 } else {
4691 String::new()
4692 },
4693 ),
4694 },
4695 );
4696 round_record.fix = Some(fix);
4697 round_record.progressed = progressed;
4698 self.state.reviews.push(round_record);
4699 self.state.save()?;
4700
4701 if matches!(
4714 continuation.outcome,
4715 ContinuationOutcome::Exhausted
4716 | ContinuationOutcome::QuotaLost
4717 | ContinuationOutcome::NoSession
4718 ) {
4719 return self
4720 .stop_reviewing(
4721 "the fixer's adoption report never came back, even after resuming its \
4722 own seat; refusing to start another round against the same worktree \
4723 while that is unresolved",
4724 &shell,
4725 &winner.worktree,
4726 )
4727 .await;
4728 }
4729
4730 let streak = self
4731 .state
4732 .reviews
4733 .iter()
4734 .rev()
4735 .take_while(|r| !r.progressed)
4736 .count();
4737 if streak >= STAGNANT_LIMIT {
4738 return self
4739 .stop_reviewing(
4740 &format!(
4741 "the tree has not moved against base for {streak} round(s) in a row"
4742 ),
4743 &shell,
4744 &winner.worktree,
4745 )
4746 .await;
4747 }
4748 }
4749 Ok(())
4750 }
4751
4752 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
4782 let round_idx = self.state.reviews.len() - 1;
4783 let needs_catchup_run = matches!(
4791 self.state.reviews[round_idx].e2e_status(),
4792 E2eStatus::Deferred | E2eStatus::ResourceBlocked
4793 );
4794 if needs_catchup_run {
4795 let round = self.state.reviews[round_idx].round;
4796 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
4797 let commands = self.state.config.verify.e2e.clone();
4798 let attempted_head = git::rev_parse(worktree, "HEAD").await?;
4799 let cache_dir = self.state.config.cache_dir();
4800 let context = format!(
4801 "round {round}: verification unresolved, catching up before the final decision"
4802 );
4803 let (outcomes, verify_retried) = with_cache_lease(
4804 &mut self.state,
4805 cache_dir.as_deref(),
4806 "e2e",
4807 "e2e",
4808 worktree,
4809 &attempted_head,
4810 timeout,
4811 &context,
4812 |state, budget| {
4813 let shell = shell.to_vec();
4814 let commands = commands.clone();
4815 let context = context.clone();
4816 async move {
4817 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
4818 .await
4819 }
4820 },
4821 )
4822 .await;
4823 let last = &mut self.state.reviews[round_idx];
4824 last.e2e = outcomes;
4825 last.verify_retried = verify_retried;
4826 last.verified_head = Some(attempted_head);
4833 last.verified_at = Some(Timestamp::now());
4834 if verify_inconclusive(&last.e2e) {
4835 self.state.save()?;
4842 return Ok(());
4843 }
4844 last.e2e_deferred = false;
4845 }
4846 let last = &self.state.reviews[round_idx];
4847 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
4848
4849 match last.e2e_status() {
4850 E2eStatus::Failed => {
4851 let red: Vec<String> = last
4852 .e2e
4853 .iter()
4854 .filter(|o| !o.ok())
4855 .map(|o| {
4856 format!(
4857 "`{}` -> {:?}\n{}",
4858 o.command,
4859 o.code,
4860 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
4861 )
4862 })
4863 .collect();
4864 self.state
4865 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
4866 self.state.status = RunStatus::Blocked;
4867 }
4868 E2eStatus::ResourceBlocked => {
4873 self.state.event(
4874 "review",
4875 format!(
4876 "{why}; e2e could not run (shared build cache unavailable); not \
4877 deciding yet"
4878 ),
4879 );
4880 }
4881 E2eStatus::Passed | E2eStatus::Deferred | E2eStatus::NotConfigured => {
4882 self.state.event(
4883 "review",
4884 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
4885 );
4886 self.state.status = RunStatus::Gating;
4887 }
4888 }
4889 self.state.save()?;
4890 Ok(())
4891 }
4892
4893 async fn gate(&mut self) -> Result<()> {
4896 if self.state.status == RunStatus::Failed
4908 || self
4909 .state
4910 .base_sync
4911 .as_ref()
4912 .is_some_and(|s| s.conflict.is_some())
4913 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
4914 != Some(RunStatus::Gating)
4915 {
4916 return Ok(());
4917 }
4918 if self.state.gate_ran {
4919 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
4930 self.state.status = RunStatus::Blocked;
4931 self.state.save()?;
4932 }
4933 return Ok(());
4934 }
4935 let Some(winner) = self.state.winner().cloned() else {
4936 return Ok(());
4937 };
4938 self.state.status = RunStatus::Gating;
4939 let mut outcomes = self.run_gate(&winner).await?;
4940 loop {
4941 if verify_inconclusive(&outcomes) {
4952 self.state.save()?;
4953 return Ok(());
4954 }
4955 if outcomes.iter().all(CommandOutcome::ok) {
4956 break;
4957 }
4958 match self.gate_fix_round(&winner, &outcomes).await? {
4959 GateFix::Retry => outcomes = self.run_gate(&winner).await?,
4960 GateFix::Stop => break,
4961 GateFix::Defer => {
4962 self.state.save()?;
4963 return Ok(());
4964 }
4965 }
4966 }
4967 let passed = outcomes.iter().all(CommandOutcome::ok);
4968 self.state.gate = outcomes;
4969 self.state.gate_ran = true;
4970 if !passed {
4971 self.state.status = RunStatus::Blocked;
4972 let spent = self.state.gate_fixes.len();
4973 self.state.event(
4974 "gate",
4975 if spent == 0 {
4976 "gate failed; not merging".to_owned()
4977 } else {
4978 format!("gate failed after {spent} gate-fix round(s); not merging")
4979 },
4980 );
4981 }
4982 self.state.save()?;
4983 Ok(())
4984 }
4985
4986 async fn run_pre_gate(&mut self, winner: &Candidate) {
4996 let commands = self.state.config.verify.pre_gate.clone();
4997 if commands.is_empty() {
4998 return;
4999 }
5000 let shell = self.state.config.shell();
5001 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5002 let (outcomes, _) = run_commands(
5003 &mut self.state,
5004 "pre_gate",
5005 "pre_gate",
5006 0,
5007 &shell,
5008 &commands,
5009 &winner.worktree,
5010 timeout,
5011 )
5012 .await;
5013 for o in &outcomes {
5014 if !o.ok() {
5015 tracing::warn!(
5016 "pre_gate `{}` failed ({:?}); the gate decides",
5017 o.command,
5018 o.code
5019 );
5020 }
5021 self.state.event(
5022 "pre_gate",
5023 format!(
5024 "`{}` -> {}",
5025 o.command,
5026 if o.ok() {
5027 "pass".to_owned()
5028 } else {
5029 format!(
5030 "FAIL ({:?})\n{}",
5031 o.code,
5032 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5033 )
5034 }
5035 ),
5036 );
5037 }
5038 self.state.pre_gate = outcomes;
5039 match git::commit_all(&winner.worktree, "magi: pre_gate (mechanical fixes)").await {
5040 Ok(true) => match git::rev_parse(&winner.worktree, "HEAD").await {
5041 Ok(head) => {
5042 self.state
5043 .event("pre_gate", format!("committed mechanical fixes ({head})"));
5044 self.state.pre_gate_commit = Some(head);
5045 }
5046 Err(e) => tracing::warn!("pre_gate committed but HEAD unreadable: {e:#}"),
5047 },
5048 Ok(false) => {}
5049 Err(e) => tracing::warn!("pre_gate could not commit its changes: {e:#}"),
5050 }
5051 if let Err(e) = self.state.save() {
5052 tracing::warn!("could not persist the pre_gate record: {e:#}");
5053 }
5054 }
5055
5056 async fn run_gate(&mut self, winner: &Candidate) -> Result<Vec<CommandOutcome>> {
5059 self.run_pre_gate(winner).await;
5060 let shell = self.state.config.shell();
5061 let gate_commands = self.state.config.verify.gate.clone();
5062 let outcomes = if gate_commands.is_empty() {
5071 Vec::new()
5072 } else {
5073 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5074 let cache_dir = self.state.config.cache_dir();
5075 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
5076 let (outcomes, _) = with_cache_lease(
5077 &mut self.state,
5078 cache_dir.as_deref(),
5079 "gate",
5080 "gate",
5081 &winner.worktree,
5082 &head,
5083 timeout,
5084 "final gate",
5085 |state, budget| {
5086 let shell = shell.clone();
5087 let gate_commands = gate_commands.clone();
5088 let worktree = winner.worktree.clone();
5089 async move {
5090 let (outcomes, timed_out_pids) = run_commands(
5091 state,
5092 "gate",
5093 "gate",
5094 0,
5095 &shell,
5096 &gate_commands,
5097 &worktree,
5098 budget,
5099 )
5100 .await;
5101 (outcomes, false, timed_out_pids)
5102 }
5103 },
5104 )
5105 .await;
5106 outcomes
5107 };
5108 if outcomes.is_empty() {
5109 self.state.event(
5114 "gate",
5115 "no gate commands configured; nothing to check, passing",
5116 );
5117 }
5118 for o in &outcomes {
5119 self.state.event(
5120 "gate",
5121 format!(
5122 "`{}` -> {}",
5123 o.command,
5124 if o.ok() {
5125 "pass".to_owned()
5126 } else {
5127 format!(
5128 "FAIL ({:?})\n{}",
5129 o.code,
5130 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
5131 )
5132 }
5133 ),
5134 );
5135 }
5136 Ok(outcomes)
5137 }
5138
5139 async fn gate_fix_round(
5152 &mut self,
5153 winner: &Candidate,
5154 outcomes: &[CommandOutcome],
5155 ) -> Result<GateFix> {
5156 let cap = self.state.config.graph.gate_fix_rounds;
5157 let spent = self.state.gate_fixes.len();
5158 if spent >= cap {
5159 if cap > 0 {
5160 self.state.event(
5161 "gate",
5162 format!("{spent} gate-fix round(s) spent and the gate still fails"),
5163 );
5164 }
5165 return Ok(GateFix::Stop);
5166 }
5167 if !gate_fixable(outcomes) {
5168 self.state.event(
5169 "gate",
5170 "gate failure is not an ordinary non-zero exit with output (timeout, missing \
5171 command or similar); not spending a fix round on it",
5172 );
5173 return Ok(GateFix::Stop);
5174 }
5175 let min_free = self.state.config.disk.min_free_bytes;
5176 if min_free > 0 {
5177 match crate::disk::free_bytes(&winner.worktree) {
5178 Ok(free) if crate::disk::enough_space(free, min_free) => {}
5179 Ok(free) => {
5180 self.state.event(
5181 "gate",
5182 format!(
5183 "only {free} bytes free ({min_free} required by `[disk] \
5184 min_free_bytes`); not spending a fix round on a failure the disk \
5185 may explain"
5186 ),
5187 );
5188 return Ok(GateFix::Stop);
5189 }
5190 Err(e) => {
5191 self.state.event(
5192 "gate",
5193 format!("free disk space could not be measured ({e:#}); no fix round"),
5194 );
5195 return Ok(GateFix::Stop);
5196 }
5197 }
5198 }
5199
5200 let attempt = spent + 1;
5201 let run_id = self.state.id.clone();
5202 let prompts = self.state.config.prompts.clone();
5203 let failed: Vec<CommandOutcome> = outcomes.iter().filter(|o| !o.ok()).cloned().collect();
5204 let base = self.landing_base();
5205 let (fix_spec, fix_seat_key) = self.fixer_spec(winner);
5206 let seat = self.seat(&fix_seat_key, &fix_spec.id);
5207 let job = SeatJob {
5208 prompt: prompt::gate_fix(
5209 &self.state.instruction,
5210 &failed,
5211 attempt,
5212 cap,
5213 &self.state.config.graph.language,
5214 ),
5215 spec: fix_spec,
5216 seat,
5217 cwd: winner.worktree.clone(),
5218 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
5219 allow_write: true,
5220 sessions: self.state.config.graph.sessions,
5221 artifacts: agent::artifacts_dir(&self.state.dir()),
5222 stem: format!("gate-fix-{attempt}"),
5223 };
5224 self.state.event(
5225 "gate",
5226 format!("gate failed; gate-fix round {attempt} of {cap}"),
5227 );
5228 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
5229 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
5230 let cache = self.state.config.cache_dir();
5231 let ctx = WaveCtx {
5232 run: &run_id,
5233 node: "gate-fix",
5234 prompts: &prompts,
5235 cache: cache.as_deref(),
5236 round: None,
5237 };
5238 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
5239 let mut record = GateFixRecord {
5240 agent: seat.agent.clone(),
5241 failed,
5242 notes: String::new(),
5243 committed: false,
5244 error: None,
5245 };
5246 match out {
5247 AgentOutcome::Ok(o) => {
5248 if let Ok(report) = verdict::extract_json::<FixReport>(&o.text) {
5251 record.notes = blind::sanitize_prose(&report.notes, &self.state.config.blind);
5252 }
5253 }
5254 AgentOutcome::Dropped(_) => {
5255 record.error = Some("the CLI dropped the stream".to_owned());
5256 }
5257 AgentOutcome::Quota(o) => {
5258 self.state.quota.push(QuotaLoss {
5259 seat: seat.key.clone(),
5260 node: "gate-fix".to_owned(),
5261 at: Timestamp::now(),
5262 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
5263 });
5264 record.error = Some("rate limited (quota); fixer could not run".to_owned());
5265 }
5266 AgentOutcome::Failed(e) => record.error = Some(e),
5267 }
5268 self.state.seats.insert(seat.key.clone(), seat);
5269 if let Ok(r) = git::rescue_commit(
5270 &winner.worktree,
5271 &format!("magi: gate fix {attempt} (uncommitted work)"),
5272 )
5273 .await
5274 {
5275 self.state.note_withheld("gate-fix", &r.withheld);
5276 }
5277 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
5278 record.committed = after != before;
5279 let changed = git::diff(&winner.worktree, &base, "HEAD").await? != patch;
5280 let note = record.error.clone();
5281 self.state.gate_fixes.push(record);
5282 self.state.save()?;
5283 if !changed {
5284 self.state.event(
5285 "gate",
5286 match note {
5287 Some(why) => format!("gate-fix round {attempt}: fixer failed ({why})"),
5288 None => format!("gate-fix round {attempt}: the tree did not change"),
5289 },
5290 );
5291 return Ok(GateFix::Stop);
5292 }
5293 self.state.event(
5294 "gate",
5295 format!("gate-fix round {attempt}: tree changed vs base; re-running verify.e2e"),
5296 );
5297
5298 let commands = self.state.config.verify.e2e.clone();
5299 if !commands.is_empty() {
5300 let shell = self.state.config.shell();
5301 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
5302 let cache_dir = self.state.config.cache_dir();
5303 let context = format!("gate-fix round {attempt}");
5304 let (e2e, _) = with_cache_lease(
5305 &mut self.state,
5306 cache_dir.as_deref(),
5307 "e2e",
5308 "e2e",
5309 &winner.worktree,
5310 &after,
5311 timeout,
5312 &context,
5313 |state, budget| {
5314 let shell = shell.clone();
5315 let commands = commands.clone();
5316 let context = context.clone();
5317 let worktree = winner.worktree.clone();
5318 async move {
5319 run_e2e_with_retry(state, &shell, &commands, &worktree, budget, &context)
5320 .await
5321 }
5322 },
5323 )
5324 .await;
5325 if verify_inconclusive(&e2e) {
5326 return Ok(GateFix::Defer);
5327 }
5328 if e2e.iter().any(|o| !o.ok()) {
5329 self.state.event(
5330 "gate",
5331 format!("gate-fix round {attempt}: verify.e2e failed after the fix"),
5332 );
5333 return Ok(GateFix::Stop);
5334 }
5335 }
5336 Ok(GateFix::Retry)
5337 }
5338
5339 async fn merge(&mut self) -> Result<()> {
5342 if self
5357 .state
5358 .base_sync
5359 .as_ref()
5360 .is_some_and(|s| s.conflict.is_some())
5361 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
5362 != Some(RunStatus::Gating)
5363 || !self.state.gate_status().ok()
5372 {
5373 return Ok(());
5374 }
5375 if self.state.merge.is_some() {
5384 return Ok(());
5385 }
5386 let Some(winner) = self.state.winner().cloned() else {
5387 return Ok(());
5388 };
5389 let repo = self.state.repo.clone();
5390 let base = self.state.base_branch.clone();
5391 let mode = self.state.config.merge.mode;
5392 let style = self.state.config.merge.style;
5393 let pr = pr_message(&self.state, winner.label);
5394 let message = pr.commit_message();
5395
5396 let outcome = match mode {
5397 MergeMode::None => MergeOutcome {
5398 mode,
5399 ok: true,
5400 detail: manual_merge_command(style, &repo, &winner.branch, &message),
5401 },
5402 MergeMode::Local => {
5403 let on = git::current_branch(&repo).await?;
5404 if on.as_deref() != Some(base.as_str()) {
5405 MergeOutcome {
5406 mode,
5407 ok: false,
5408 detail: format!(
5409 "{} has {} checked out, not the base branch {base}",
5410 repo.display(),
5411 on.unwrap_or_else(|| "a detached HEAD".to_owned())
5412 ),
5413 }
5414 } else if !git::is_clean(&repo).await? {
5415 MergeOutcome {
5416 mode,
5417 ok: false,
5418 detail: format!("{} is dirty; refusing to merge", repo.display()),
5419 }
5420 } else {
5421 let out = match style {
5422 MergeStyle::Merge => {
5423 git::merge_no_ff(&repo, &winner.branch, &message).await?
5424 }
5425 MergeStyle::Squash => {
5426 git::merge_squash(&repo, &winner.branch, &message).await?
5427 }
5428 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
5429 };
5430 MergeOutcome {
5431 mode,
5432 ok: out.ok(),
5433 detail: if out.ok() { out.stdout } else { out.stderr },
5434 }
5435 }
5436 }
5437 MergeMode::Pr => {
5438 let remote = self.state.config.merge.remote.clone();
5439 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
5440 if !pushed.ok() {
5441 MergeOutcome {
5442 mode,
5443 ok: false,
5444 detail: pushed.stderr,
5445 }
5446 } else {
5447 let out =
5448 gh_pr_create(&winner.worktree, &base, &winner.branch, &pr.title, &pr.body)
5449 .await;
5450 match out {
5451 Ok(url) => MergeOutcome {
5452 mode,
5453 ok: true,
5454 detail: url,
5455 },
5456 Err(e) => MergeOutcome {
5457 mode,
5458 ok: false,
5459 detail: e.to_string(),
5460 },
5461 }
5462 }
5463 }
5464 };
5465
5466 self.state.status = match (mode, outcome.ok) {
5467 (MergeMode::None, _) => RunStatus::Ready,
5468 (_, true) => RunStatus::Merged,
5469 (_, false) => RunStatus::Blocked,
5470 };
5471 self.state.event(
5472 "merge",
5473 format!(
5474 "{:?}: {}",
5475 mode,
5476 outcome.detail.lines().next().unwrap_or("")
5477 ),
5478 );
5479 self.state.merge = Some(outcome);
5480 self.state.save()?;
5481
5482 if self.state.config.graph.land
5488 && mode == MergeMode::Pr
5489 && self.state.status == RunStatus::Merged
5490 {
5491 self.run_land().await?;
5492 }
5493 self.settle_questions();
5498 Ok(())
5499 }
5500
5501 async fn run_land(&mut self) -> Result<()> {
5512 let url = self
5513 .state
5514 .merge
5515 .as_ref()
5516 .map(|m| m.detail.clone())
5517 .unwrap_or_default();
5518 let url = url.lines().next().unwrap_or("").trim().to_owned();
5519 if !url.starts_with("http") {
5520 return Ok(());
5521 }
5522 match land::land(&mut self.state, &url).await {
5525 Ok(pr) if self.state.parked => {
5526 let _ = pr;
5530 }
5531 Ok(pr) => {
5532 self.state.status = match pr.state {
5533 land::PrLifecycle::Merged => RunStatus::Merged,
5534 _ => RunStatus::Blocked,
5535 };
5536 if bump::should_release_bump(self.state.status)
5543 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
5544 {
5545 self.state
5551 .event("bump", format!("release bump skipped: {e:#}"));
5552 }
5553 self.state.save()?;
5554 }
5555 Err(e) => {
5556 self.state.status = RunStatus::Blocked;
5557 self.state.event("land", format!("gave up: {e}"));
5558 self.state.save()?;
5559 }
5560 }
5561 Ok(())
5562 }
5563
5564 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
5568 if let Some(existing) = self.state.seats.get(key)
5569 && existing.agent == agent
5570 {
5571 return existing.clone();
5572 }
5573 let fresh = SeatState::new(key, agent, self.state.seed);
5574 self.state.seats.insert(key.to_owned(), fresh.clone());
5575 fresh
5576 }
5577
5578 fn view(&self, c: &Candidate) -> CandidateView {
5580 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
5581 .unwrap_or_default();
5582 let (patch, _) = blind::sanitize_patch(
5583 &format!("candidate {} patch", c.label),
5584 &raw,
5585 &self.state.config.blind,
5586 );
5587 CandidateView {
5588 label: c.label,
5589 branch: c.branch.clone(),
5590 summary: c.summary.clone(),
5591 stat: c.stat.clone(),
5592 patch,
5593 }
5594 }
5595
5596 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
5598 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
5599 prompt::judge(
5600 "(see above)",
5601 &views,
5602 self.roles.judges.len(),
5603 base_short,
5604 "en",
5605 )
5606 }
5607
5608 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
5615 let mut turns = Vec::new();
5616 for j in &self.state.judgements {
5617 if j.ranking.is_empty() {
5618 continue;
5619 }
5620 let reasons = j
5621 .reasons
5622 .iter()
5623 .map(|(k, v)| format!("- {k}: {v}"))
5624 .collect::<Vec<_>>()
5625 .join("\n");
5626 turns.push(Turn {
5627 who: format!("Judge {} (opening ranking)", j.judge),
5628 is_self: j.judge == self_idx + 1,
5629 body: format!(
5630 "Ranked {}{}{reasons}",
5631 j.ranking.iter().collect::<String>(),
5632 if reasons.is_empty() {
5633 ""
5634 } else {
5635 ", because:\n"
5636 }
5637 ),
5638 });
5639 }
5640 for t in self
5641 .state
5642 .deliberation
5643 .iter()
5644 .flat_map(|r| r.turns.iter())
5645 .chain(current)
5646 {
5647 turns.push(Turn {
5648 who: format!("Judge {}", t.judge),
5649 is_self: t.judge == self_idx + 1,
5650 body: t.body.clone(),
5651 });
5652 }
5653 turns
5654 }
5655}
5656
5657fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
5659 agent::has_session(spec.kind, seat, sessions)
5660}
5661
5662fn next_untried_implementer<'a>(
5683 roster: &'a [AgentSpec],
5684 start: usize,
5685 tried: &BTreeSet<String>,
5686) -> Option<&'a AgentSpec> {
5687 roster
5688 .get(start + 1..)?
5689 .iter()
5690 .find(|s| !tried.contains(&s.id))
5691}
5692
5693fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
5707 commands.iter().any(|c| c.exit_code.is_none())
5708}
5709
5710fn verified_noop_claim(
5723 usable: bool,
5724 commands: &[agent::CommandEvidence],
5725 text: &str,
5726) -> Option<String> {
5727 (usable && !has_unconfirmed_command(commands))
5728 .then(|| verdict::verified_noop(text))
5729 .flatten()
5730}
5731
5732fn short(commit: &str) -> String {
5733 commit.chars().take(7).collect()
5734}
5735
5736fn make_executable(path: &Path) -> Result<()> {
5737 #[cfg(unix)]
5738 {
5739 use std::os::unix::fs::PermissionsExt as _;
5740 let mut perms = std::fs::metadata(path)?.permissions();
5741 perms.set_mode(0o755);
5742 std::fs::set_permissions(path, perms)?;
5743 }
5744 #[cfg(not(unix))]
5745 {
5746 let _ = path;
5747 }
5748 Ok(())
5749}
5750
5751struct WaveCtx<'a> {
5758 run: &'a str,
5761 node: &'a str,
5763 prompts: &'a Prompts,
5764 cache: Option<&'a Path>,
5766 round: Option<usize>,
5769}
5770
5771async fn run_one(
5773 job: SeatJob,
5774 sem: Arc<Semaphore>,
5775 ctx: &WaveCtx<'_>,
5776 state: &mut RunState,
5777 attempt: usize,
5778) -> (SeatState, AgentOutcome) {
5779 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
5780 .await
5781 .pop()
5782 .expect("one job in, one result out");
5783 (seat, out)
5784}
5785
5786async fn wave(
5792 jobs: Vec<SeatJob>,
5793 sem: Arc<Semaphore>,
5794 ctx: &WaveCtx<'_>,
5795 state: &mut RunState,
5796 attempt: usize,
5797) -> Vec<(usize, SeatState, AgentOutcome)> {
5798 let WaveCtx {
5799 run,
5800 node,
5801 prompts,
5802 cache,
5803 round,
5804 } = *ctx;
5805 for job in &jobs {
5806 state.seat_started(node, &job.seat.key, job.timeout, attempt);
5807 }
5808 if let Err(e) = state.save() {
5809 tracing::warn!("could not persist in-progress seats: {e:#}");
5814 }
5815 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
5831 let wait_started = Instant::now();
5832 let cache_guard = if let Some(cache_dir) = cache {
5833 if jobs_had_a_writer {
5834 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
5835 let budget = jobs
5836 .iter()
5837 .map(|j| j.timeout)
5838 .max()
5839 .unwrap_or(Duration::from_secs(60));
5840 acquire_cache_lease(state, cache_dir, &owner, budget, node)
5841 .await
5842 .ok()
5843 } else {
5844 None
5845 }
5846 } else {
5847 None
5848 };
5849 let waited_for_lease = wait_started.elapsed();
5856 let mut set = tokio::task::JoinSet::new();
5857 let overlay = prompts.overlay(node);
5858 for (i, mut job) in jobs.into_iter().enumerate() {
5859 job.timeout = job.timeout.saturating_sub(waited_for_lease);
5860 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
5861 if cache.is_some() {
5862 job.prompt.push('\n');
5863 job.prompt
5864 .push_str(&prompt::build_cache_note(node, job.allow_write));
5865 }
5866 let sem = Arc::clone(&sem);
5867 let run = run.to_owned();
5868 let node = node.to_owned();
5869 let cache = cache
5880 .filter(|_| job.allow_write && cache_guard.is_some())
5881 .map(Path::to_path_buf);
5882 set.spawn(async move {
5883 let _permit = sem.acquire().await;
5884 let mut seat = job.seat;
5885 let out = agent::invoke(
5886 &job.spec,
5887 &mut seat,
5888 &Invocation {
5889 cwd: &job.cwd,
5890 prompt: &job.prompt,
5891 timeout: job.timeout,
5892 allow_write: job.allow_write,
5893 sessions: job.sessions,
5894 artifacts: &job.artifacts,
5895 stem: &job.stem,
5896 run: &run,
5897 node: &node,
5898 cache_dir: cache.as_deref(),
5899 attachments: &[],
5900 },
5901 )
5902 .await;
5903 let out = match out {
5904 Ok(o) if o.usable() => AgentOutcome::Ok(o),
5905 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
5906 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
5914 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
5915 Ok(o) => AgentOutcome::Failed(format!(
5916 "exited with {:?} and no usable output",
5917 o.exit_code
5918 )),
5919 Err(e) => AgentOutcome::Failed(e.to_string()),
5920 };
5921 (i, seat, out)
5922 });
5923 }
5924 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
5925 while let Some(joined) = set.join_next().await {
5926 let (i, seat, out) = match joined {
5927 Ok(v) => v,
5928 Err(e) => {
5932 tracing::error!("agent task panicked: {e}");
5933 continue;
5934 }
5935 };
5936 state.seat_finished(&seat.key);
5937 record_jobs(state, node, round, &seat.key, &out);
5938 if let Err(e) = state.save() {
5939 tracing::warn!("could not persist a seat's completion: {e:#}");
5940 }
5941 if collected.len() <= i {
5942 collected.resize_with(i + 1, || None);
5943 }
5944 collected[i] = Some((i, seat, out));
5945 }
5946 if state
5952 .active
5953 .values()
5954 .any(|a| a.node == node && a.attempt == attempt)
5955 {
5956 state
5957 .active
5958 .retain(|_, a| !(a.node == node && a.attempt == attempt));
5959 if let Err(e) = state.save() {
5960 tracing::warn!("could not persist the end of a wave: {e:#}");
5961 }
5962 }
5963 if let Some(cache_dir) = cache
5970 && jobs_had_a_writer
5971 {
5972 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
5973 }
5974 if let Some(guard) = cache_guard {
5975 guard.release();
5976 }
5977 collected.into_iter().flatten().collect()
5978}
5979
5980fn record_jobs(
5991 state: &mut RunState,
5992 node: &str,
5993 round: Option<usize>,
5994 seat: &str,
5995 out: &AgentOutcome,
5996) {
5997 let commands: &[agent::CommandEvidence] = match out {
5998 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
5999 AgentOutcome::Failed(_) => &[],
6000 };
6001 let checked_at = Timestamp::now();
6002 for c in commands {
6003 state.jobs.push(JobRecord {
6004 node: node.to_owned(),
6005 round,
6006 seat: seat.to_owned(),
6007 id: c.id.clone(),
6008 description: c.description.clone(),
6009 checked_at,
6010 status: match c.exit_code {
6011 Some(0) => JobStatus::Completed,
6012 Some(_) => JobStatus::Failed,
6013 None => JobStatus::Unknown,
6014 },
6015 exit_code: c.exit_code,
6016 result_summary: c.result_summary.clone(),
6017 source: c.source.clone(),
6018 });
6019 }
6020}
6021
6022fn round_is_clean(
6043 blocking: usize,
6044 e2e_ok: bool,
6045 answered: usize,
6046 expected: usize,
6047 quota_missing: usize,
6048 policy: IncompleteReviewPolicy,
6049) -> bool {
6050 if blocking != 0 || !e2e_ok {
6051 return false;
6052 }
6053 if answered == expected || policy == IncompleteReviewPolicy::Warn {
6054 return true;
6055 }
6056 answered > 0 && expected - answered <= quota_missing
6057}
6058
6059fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
6083 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
6084 return Some(RunStatus::Gating);
6085 }
6086 let last = reviews.last()?;
6087 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
6088 if reviews.len() < max_rounds && !stagnant {
6089 return None;
6090 }
6091 if last.incomplete() && last.blocking == 0 {
6092 return Some(RunStatus::Blocked);
6093 }
6094 if last.e2e_status() == E2eStatus::ResourceBlocked {
6095 return None;
6096 }
6097 Some(if last.e2e.iter().all(CommandOutcome::ok) {
6098 RunStatus::Gating
6099 } else {
6100 RunStatus::Blocked
6101 })
6102}
6103
6104fn retry_budget(full: Duration, nudged: bool) -> Duration {
6119 if nudged {
6120 (full / 4).max(Duration::from_secs(120)).min(full)
6121 } else {
6122 full
6123 }
6124}
6125
6126#[allow(clippy::too_many_arguments)]
6139async fn ask_json_wave<T>(
6140 jobs: Vec<SeatJob>,
6141 sem: Arc<Semaphore>,
6142 retries: usize,
6143 ctx: &WaveCtx<'_>,
6144 losses: &mut Vec<QuotaLoss>,
6145 state: &mut RunState,
6146 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
6147) -> Vec<(SeatState, Result<(T, AgentOutput)>, usize)>
6148where
6149 T: serde::de::DeserializeOwned + Send + 'static,
6150{
6151 let n = jobs.len();
6152 let originals: Vec<SeatJob> = jobs;
6153 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
6154 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
6155 let mut attempts_used: Vec<usize> = vec![0; n];
6162 let mut pending: Vec<usize> = (0..n).collect();
6163
6164 for attempt in 0..=retries {
6165 if pending.is_empty() {
6166 break;
6167 }
6168 let mut batch = Vec::with_capacity(pending.len());
6169 for &i in &pending {
6170 let src = &originals[i];
6171 let (prompt, timeout) = if attempt == 0 {
6174 (src.prompt.clone(), src.timeout)
6175 } else {
6176 let why = done[i]
6177 .as_ref()
6178 .and_then(|r| r.as_ref().err().map(ToString::to_string))
6179 .unwrap_or_else(|| "no parsable answer".to_owned());
6180 let nudge = prompt::nudge(&why);
6181 let nudged = has_context(&src.spec, &seats[i], src.sessions);
6182 let prompt = if nudged {
6183 nudge
6184 } else {
6185 format!("{}\n\n---\n\n{}", src.prompt, nudge)
6186 };
6187 (prompt, retry_budget(src.timeout, nudged))
6188 };
6189 batch.push(SeatJob {
6190 spec: src.spec.clone(),
6191 seat: seats[i].clone(),
6192 cwd: src.cwd.clone(),
6193 prompt,
6194 timeout,
6195 allow_write: src.allow_write,
6196 sessions: src.sessions,
6197 artifacts: src.artifacts.clone(),
6198 stem: if attempt == 0 {
6199 src.stem.clone()
6200 } else {
6201 format!("{}-retry{attempt}", src.stem)
6202 },
6203 });
6204 }
6205
6206 if attempt > 0 {
6207 let seats_out: Vec<&str> = pending
6208 .iter()
6209 .map(|&i| originals[i].seat.key.as_str())
6210 .collect();
6211 state.event(
6212 ctx.node,
6213 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
6214 );
6215 }
6216 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
6217 let mut still = Vec::new();
6218 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
6219 seats[i] = seat;
6220 let (parsed, quota) = match out {
6221 AgentOutcome::Ok(o) => (
6222 match verdict::extract_json::<T>(&o.text) {
6223 Ok(v) => match validate(&v) {
6224 Ok(()) => Ok((v, o)),
6225 Err(e) => Err(e),
6226 },
6227 Err(e) => Err(e),
6228 },
6229 false,
6230 ),
6231 AgentOutcome::Quota(o) => {
6232 losses.push(QuotaLoss {
6233 seat: originals[i].seat.key.clone(),
6234 node: ctx.node.to_owned(),
6235 at: Timestamp::now(),
6236 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
6237 });
6238 (
6239 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
6240 true,
6241 )
6242 }
6243 AgentOutcome::Dropped(o) => {
6248 let why = o
6249 .dropped
6250 .as_ref()
6251 .map(|d| d.why.as_str())
6252 .unwrap_or("the CLI ended the stream without delivering its answer");
6253 (
6254 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
6255 false,
6256 )
6257 }
6258 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
6259 };
6260 let failed = parsed.is_err();
6261 done[i] = Some(parsed);
6262 attempts_used[i] = attempt;
6263 if failed && !quota {
6266 still.push(i);
6267 }
6268 }
6269 pending = still;
6270 }
6271
6272 seats
6273 .into_iter()
6274 .zip(done)
6275 .zip(attempts_used)
6276 .map(|((seat, res), attempts)| {
6277 (
6278 seat,
6279 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
6280 attempts,
6281 )
6282 })
6283 .collect()
6284}
6285
6286async fn acquire_cache_lease(
6299 state: &mut RunState,
6300 cache_dir: &Path,
6301 owner: &crate::cache::Owner,
6302 budget: Duration,
6303 context: &str,
6304) -> Result<crate::cache::Guard> {
6305 let home = crate::run::home();
6306 let started = Instant::now();
6307 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
6308 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
6309 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
6310 Err(e) => {
6311 state.event(
6312 "verify",
6313 format!("{context}: could not check the shared build cache: {e:#}"),
6314 );
6315 if let Err(e2) = state.save() {
6316 tracing::warn!("could not persist a cache-check failure: {e2:#}");
6317 }
6318 return Err(e);
6319 }
6320 };
6321 state.event(
6322 "verify",
6323 format!(
6324 "{context}: waiting for the shared build cache at {} ({})",
6325 cache_dir.display(),
6326 busy.describe()
6327 ),
6328 );
6329 if let Err(e) = state.save() {
6330 tracing::warn!("could not persist a cache wait: {e:#}");
6331 }
6332 let remaining = budget.saturating_sub(started.elapsed());
6333 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
6334 Ok(g) => Ok(g),
6335 Err(e) => {
6336 state.event("verify", format!("{context}: {e:#}"));
6337 if let Err(e2) = state.save() {
6338 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
6339 }
6340 Err(e)
6341 }
6342 }
6343}
6344
6345#[allow(clippy::too_many_arguments)]
6366async fn with_cache_lease<'s, F, Fut>(
6367 state: &'s mut RunState,
6368 cache_dir: Option<&Path>,
6369 node: &str,
6370 seat: &str,
6371 worktree: &Path,
6372 head: &str,
6373 budget: Duration,
6374 context: &str,
6375 body: F,
6376) -> (Vec<CommandOutcome>, bool)
6377where
6378 F: FnOnce(&'s mut RunState, Duration) -> Fut,
6379 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
6380{
6381 let Some(cache_dir) = cache_dir else {
6382 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
6383 return (outcomes, retried);
6384 };
6385 let home = crate::run::home();
6386 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
6387 let started = Instant::now();
6388 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
6389 Ok(g) => g,
6390 Err(e) => {
6391 return (
6392 vec![CommandOutcome {
6393 command: "(waiting for the shared build cache)".to_owned(),
6394 code: None,
6395 output_tail: e.to_string(),
6396 duration_ms: started.elapsed().as_millis() as u64,
6397 resource_blocked: true,
6398 }],
6399 false,
6400 );
6401 }
6402 };
6403 let identity = crate::cache::Identity::new(worktree, head);
6404 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
6405 state.event(
6413 "verify",
6414 format!(
6415 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
6416 worktree.display(),
6417 short(head)
6418 ),
6419 );
6420 guard.release();
6421 return (
6422 vec![CommandOutcome {
6423 command: "(confirming the shared build cache is fresh)".to_owned(),
6424 code: None,
6425 output_tail: e.to_string(),
6426 duration_ms: started.elapsed().as_millis() as u64,
6427 resource_blocked: true,
6428 }],
6429 false,
6430 );
6431 }
6432 let remaining = budget.saturating_sub(started.elapsed());
6433 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
6434 if !timed_out_pids.is_empty() {
6439 wait_for_timed_out_children_to_die(&timed_out_pids).await;
6440 }
6441 guard.release();
6442 (outcomes, retried)
6443}
6444
6445async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
6457 wait_for_pids_with(
6458 pids,
6459 crate::proc::pid_alive,
6460 LEASE_RELEASE_POLL,
6461 LEASE_RELEASE_MAX_WAIT,
6462 )
6463 .await;
6464}
6465
6466async fn wait_for_pids_with<F: Fn(u32) -> bool>(
6472 pids: &[u32],
6473 alive: F,
6474 poll: Duration,
6475 max_wait: Duration,
6476) {
6477 let deadline = Instant::now() + max_wait;
6478 loop {
6479 if pids.iter().all(|&pid| !alive(pid)) {
6480 return;
6481 }
6482 if Instant::now() >= deadline {
6483 return;
6484 }
6485 tokio::time::sleep(poll).await;
6486 }
6487}
6488
6489fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
6496 outcomes.iter().any(|o| o.resource_blocked)
6497}
6498
6499enum GateFix {
6501 Retry,
6503 Stop,
6506 Defer,
6509}
6510
6511fn gate_fixable(outcomes: &[CommandOutcome]) -> bool {
6519 let mut red = outcomes.iter().filter(|o| !o.ok()).peekable();
6520 red.peek().is_some()
6521 && red.all(|o| {
6522 !o.resource_blocked
6523 && matches!(o.code, Some(c) if c != 0 && c != 126 && c != 127)
6524 && !o.output_tail.trim().is_empty()
6525 })
6526}
6527
6528fn e2e_outcome_label(o: &CommandOutcome) -> String {
6532 if o.ok() {
6533 return "pass".to_owned();
6534 }
6535 let reason = if o.build_failed() {
6536 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
6537 } else {
6538 format!("FAIL ({:?})", o.code)
6539 };
6540 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
6541}
6542
6543async fn run_e2e_with_retry(
6551 state: &mut RunState,
6552 shell: &[String],
6553 commands: &[String],
6554 worktree: &Path,
6555 timeout: Duration,
6556 context: &str,
6557) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
6558 let (mut e2e, mut timed_out_pids) = run_commands(
6559 state, "verify", "e2e", 0, shell, commands, worktree, timeout,
6560 )
6561 .await;
6562 for o in &e2e {
6563 state.event(
6564 "verify",
6565 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
6566 );
6567 }
6568 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
6572 if verify_retried {
6573 state.event(
6574 "verify",
6575 format!(
6576 "{context}: verify could not build/link, not a test result — retrying once \
6577 before concluding"
6578 ),
6579 );
6580 let retried = run_commands(
6581 state, "verify", "e2e", 1, shell, commands, worktree, timeout,
6582 )
6583 .await;
6584 e2e = retried.0;
6585 timed_out_pids.extend(retried.1);
6588 for o in &e2e {
6589 state.event(
6590 "verify",
6591 format!(
6592 "{context}: retry `{}` -> {}",
6593 o.command,
6594 e2e_outcome_label(o)
6595 ),
6596 );
6597 }
6598 }
6599 (e2e, verify_retried, timed_out_pids)
6600}
6601
6602#[allow(clippy::too_many_arguments)]
6618async fn run_commands(
6619 state: &mut RunState,
6620 node: &str,
6621 task: &str,
6622 attempt: usize,
6623 shell: &[String],
6624 commands: &[String],
6625 cwd: &Path,
6626 timeout: Duration,
6627) -> (Vec<CommandOutcome>, Vec<u32>) {
6628 if commands.is_empty() {
6629 return (Vec::new(), Vec::new());
6634 }
6635 let mut out = Vec::new();
6636 let mut timed_out_pids = Vec::new();
6637 let total = commands.len();
6638 for (idx, command) in commands.iter().enumerate() {
6639 state.task_command(task, node, attempt, command, idx + 1, total, timeout);
6640 if let Err(e) = state.save() {
6641 tracing::warn!("could not persist an in-progress {task} command: {e:#}");
6642 }
6643 let started = Instant::now();
6644 let mut cmd = tokio::process::Command::new(&shell[0]);
6645 cmd.quiet();
6646 cmd.args(&shell[1..])
6647 .arg(command)
6648 .current_dir(cwd)
6649 .stdin(std::process::Stdio::null())
6650 .stdout(std::process::Stdio::piped())
6651 .stderr(std::process::Stdio::piped())
6652 .kill_on_drop(true);
6653 let spawned = cmd.spawn();
6654 let (code, body) = match spawned {
6655 Ok(child) => {
6656 let pid = child.id();
6661 match tokio::time::timeout(timeout, child.wait_with_output()).await {
6662 Ok(Ok(o)) => {
6663 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
6664 body.push_str(&String::from_utf8_lossy(&o.stderr));
6665 (o.status.code(), body)
6666 }
6667 Ok(Err(e)) => (None, format!("failed to run: {e}")),
6668 Err(_) => {
6669 if let Some(pid) = pid {
6670 timed_out_pids.push(pid);
6671 }
6672 (None, format!("timed out after {}s", timeout.as_secs()))
6673 }
6674 }
6675 }
6676 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
6677 };
6678 out.push(CommandOutcome {
6679 command: command.clone(),
6680 code,
6681 output_tail: tail(&body, OUTPUT_TAIL),
6682 duration_ms: started.elapsed().as_millis() as u64,
6683 resource_blocked: false,
6684 });
6685 }
6686 state.task_finished(task);
6687 if let Err(e) = state.save() {
6688 tracing::warn!("could not persist the end of {task}: {e:#}");
6689 }
6690 (out, timed_out_pids)
6691}
6692
6693fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
6705 let repo = repo.display();
6706 match style {
6707 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
6708 MergeStyle::Squash => {
6709 let subject = message
6712 .lines()
6713 .next()
6714 .unwrap_or(branch)
6715 .replace(['\\', '"', '$', '`'], "");
6716 format!(
6717 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
6718 )
6719 }
6720 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
6721 }
6722}
6723
6724const PR_TITLE_MAX: usize = 240;
6735
6736struct PrMessage {
6741 title: String,
6742 body: String,
6743}
6744
6745impl PrMessage {
6746 fn commit_message(&self) -> String {
6750 format!("{}\n\n{}", self.title, self.body)
6751 }
6752}
6753
6754fn title_marker(line: &str) -> Option<&str> {
6756 let line = line.trim();
6757 let head = line.get(..6)?;
6758 head.eq_ignore_ascii_case("title:")
6759 .then(|| line[6..].trim())
6760}
6761
6762fn summary_title(summary: &str) -> Option<String> {
6767 let first = summary.lines().find(|l| !l.trim().is_empty())?;
6768 let raw = title_marker(first)?;
6769 if raw.is_empty() {
6770 return None;
6771 }
6772 let title = queue::title_from(raw, PR_TITLE_MAX);
6773 let lower = title.to_ascii_lowercase();
6774 if lower.starts_with("magi:") || lower.contains("(uncommitted work)") {
6775 return None;
6776 }
6777 Some(title)
6778}
6779
6780fn summary_without_title(summary: &str) -> String {
6783 let mut lines = summary.trim().lines().peekable();
6784 if lines.peek().is_some_and(|l| title_marker(l).is_some()) {
6785 lines.next();
6786 }
6787 lines.collect::<Vec<_>>().join("\n").trim().to_owned()
6788}
6789
6790fn pr_message(state: &RunState, winner: char) -> PrMessage {
6804 let summary = state
6805 .candidates
6806 .iter()
6807 .find(|c| c.label == winner)
6808 .map(|c| c.summary.as_str())
6809 .unwrap_or_default();
6810 let title = summary_title(summary).unwrap_or_else(|| {
6813 let t = queue::title_from(&state.instruction, PR_TITLE_MAX);
6814 if t.is_ascii() && t.chars().any(|c| c.is_ascii_alphabetic()) {
6815 t
6816 } else {
6817 format!(
6818 "chore: land candidate {} of run {}",
6819 winner.to_ascii_uppercase(),
6820 state.id
6821 )
6822 }
6823 });
6824
6825 let mut body = String::new();
6826 let what = summary_without_title(summary);
6827 if !what.is_empty() {
6828 body.push_str("## Summary\n\n");
6829 body.push_str(&what);
6830 body.push_str("\n\n");
6831 }
6832
6833 let fix = state.reviews.last().and_then(|r| r.fix.as_ref());
6834 if let Some(fix) = fix
6835 && !fix.notes.trim().is_empty()
6836 {
6837 body.push_str("## Review fixes\n\n");
6838 body.push_str(fix.notes.trim());
6839 body.push_str("\n\n");
6840 }
6841
6842 let open = state.open_findings();
6843 if !open.is_empty() {
6844 body.push_str("## Open review findings\n\n");
6845 for f in &open {
6846 body.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
6847 }
6848 body.push('\n');
6849 }
6850
6851 if let Some(fix) = fix
6852 && !fix.rejected.is_empty()
6853 {
6854 body.push_str("## Declined by the fixer\n\n");
6855 for r in &fix.rejected {
6856 body.push_str(&format!("- `{}`: {}\n", r.id, r.why));
6857 }
6858 body.push('\n');
6859 }
6860
6861 let task = state.instruction.trim();
6862 let task = if task.is_empty() {
6863 "(empty task)"
6864 } else {
6865 task
6866 };
6867 body.push_str(&format!(
6868 "<details>\n<summary>Original task</summary>\n\n{}\n\n</details>\n",
6869 task.replace("</details>", "</details>")
6870 ));
6871
6872 body.push_str(&format!(
6873 "\n---\nmagi:run/{} magi:candidate-{}\n",
6874 state.id,
6875 winner.to_ascii_lowercase()
6876 ));
6877
6878 let id = crate::scrub::Identity::current();
6881 PrMessage {
6882 title: crate::scrub::scrub(&title, &id),
6883 body: crate::scrub::scrub(&body, &id),
6884 }
6885}
6886
6887async fn gh_pr_create(
6889 cwd: &Path,
6890 base: &str,
6891 head: &str,
6892 title: &str,
6893 body: &str,
6894) -> Result<String> {
6895 let out = tokio::process::Command::new("gh")
6896 .args([
6897 "pr", "create", "--base", base, "--head", head, "--title", title, "--body", body,
6898 ])
6899 .current_dir(cwd)
6900 .quiet()
6901 .stdin(std::process::Stdio::null())
6902 .output()
6903 .await
6904 .context("spawn gh")?;
6905 if out.status.success() {
6906 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
6907 } else {
6908 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
6909 }
6910}
6911
6912pub async fn fold_run(state: &mut RunState, drop_winner: bool, home: &Path) -> Result<Vec<String>> {
6921 let repo = state.repo.clone();
6922 let root = state.worktree_root();
6923 let winner = state.tally.as_ref().map(|t| t.winner);
6924 let mut removed = Vec::new();
6925
6926 for i in 0..state.candidates.len() {
6927 let c = state.candidates[i].clone();
6928 let is_winner = Some(c.label) == winner;
6929 if is_winner && !drop_winner {
6930 continue;
6931 }
6932 if c.worktree.exists() {
6933 git::worktree_remove(&repo, &c.worktree).await.ok();
6934 removed.push(c.worktree.to_string_lossy().into_owned());
6935 }
6936 let handed_over = state.released_branches.contains(&c.branch);
6939 if !handed_over && git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
6940 git::branch_delete(&repo, &c.branch).await.ok();
6941 removed.push(c.branch.clone());
6942 }
6943 state.candidates[i].folded = true;
6944 }
6945
6946 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
6947 let path = name.path();
6948 let keep = !drop_winner
6949 && winner.is_some_and(|w| {
6950 path.file_name()
6951 .is_some_and(|n| n == format!("cand-{w}").as_str())
6952 });
6953 if keep {
6954 continue;
6955 }
6956 git::worktree_remove(&repo, &path).await.ok();
6957 removed.push(path.to_string_lossy().into_owned());
6958 }
6959
6960 remove_if_empty(&root);
6969
6970 if state.enabled_worktree_config && drop_winner {
6971 git::release_worktree_config(&repo).await.ok();
6975 state.enabled_worktree_config = false;
6976 }
6977 state.save_under(home)?;
6978 Ok(removed)
6979}
6980
6981fn remove_if_empty(dir: &Path) {
6992 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
6993 std::fs::remove_dir(dir).ok();
6994 }
6995}
6996
6997pub fn worst_open(state: &RunState) -> Option<Severity> {
6999 state
7000 .reviews
7001 .last()?
7002 .reviews
7003 .iter()
7004 .flat_map(|r| r.findings.iter())
7005 .map(|f| f.severity)
7006 .max()
7007}
7008
7009#[cfg(test)]
7010mod tests {
7011 use super::*;
7012 use crate::run::GateStatus;
7013 use std::collections::BTreeMap;
7014 use std::time::Duration;
7015
7016 fn conductor() -> AgentSpec {
7017 AgentSpec {
7018 id: "conductor".to_owned(),
7019 kind: crate::config::AgentKind::Command,
7020 model: None,
7021 command: vec!["true".to_owned()],
7022 extra_args: Vec::new(),
7023 env: BTreeMap::new(),
7024 prompt_delivery: None,
7025 }
7026 }
7027
7028 fn spec(id: &str) -> AgentSpec {
7029 AgentSpec {
7030 id: id.to_owned(),
7031 kind: crate::config::AgentKind::Command,
7032 model: None,
7033 command: vec!["true".to_owned()],
7034 extra_args: Vec::new(),
7035 env: BTreeMap::new(),
7036 prompt_delivery: None,
7037 }
7038 }
7039
7040 #[test]
7047 fn next_untried_implementer_walks_forward_from_the_seats_own_position() {
7048 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7049 let tried = BTreeSet::from(["beta".to_owned()]);
7050 let next = next_untried_implementer(&roster, 1, &tried);
7053 assert_eq!(next.map(|s| s.id.as_str()), Some("gamma"));
7054 }
7055
7056 #[test]
7057 fn next_untried_implementer_does_not_wrap_back_past_its_own_start() {
7058 let roster = vec![spec("alpha"), spec("beta")];
7059 let tried = BTreeSet::from(["beta".to_owned()]);
7060 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7064 }
7065
7066 #[test]
7067 fn next_untried_implementer_stops_once_the_tail_is_exhausted_even_if_earlier_ids_are_untried() {
7068 let roster = vec![spec("alpha"), spec("beta"), spec("gamma")];
7069 let tried = BTreeSet::from(["beta".to_owned(), "gamma".to_owned()]);
7070 assert!(next_untried_implementer(&roster, 1, &tried).is_none());
7074 }
7075
7076 #[test]
7077 fn next_untried_implementer_skips_ids_already_tried_even_when_duplicated() {
7078 let roster = vec![spec("a"), spec("a"), spec("b")];
7079 let tried = BTreeSet::from(["a".to_owned()]);
7080 let next = next_untried_implementer(&roster, 0, &tried);
7081 assert_eq!(next.map(|s| s.id.as_str()), Some("b"));
7082 }
7083
7084 #[test]
7085 fn next_untried_implementer_returns_none_once_every_id_is_tried() {
7086 let roster = vec![spec("a"), spec("b")];
7087 let tried = BTreeSet::from(["a".to_owned(), "b".to_owned()]);
7088 assert!(next_untried_implementer(&roster, 0, &tried).is_none());
7089 }
7090
7091 #[test]
7092 fn remove_if_empty_only_ever_takes_a_bare_directory() {
7093 let dir = tempfile::tempdir().unwrap();
7094 let bay = dir.path().join("ffff");
7095
7096 remove_if_empty(&bay);
7098 assert!(!bay.exists());
7099
7100 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
7103 remove_if_empty(&bay);
7104 assert!(bay.exists(), "non-empty directory must survive");
7105
7106 std::fs::remove_dir(bay.join("cand-A")).unwrap();
7108 remove_if_empty(&bay);
7109 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
7110 }
7111
7112 #[test]
7121 fn a_full_panel_that_found_nothing_is_clean() {
7122 assert!(round_is_clean(
7123 0,
7124 true,
7125 2,
7126 2,
7127 0,
7128 IncompleteReviewPolicy::Block
7129 ));
7130 }
7131
7132 #[test]
7133 fn a_missing_seat_is_never_clean_under_the_default_policy() {
7134 assert!(!round_is_clean(
7135 0,
7136 true,
7137 1,
7138 2,
7139 0,
7140 IncompleteReviewPolicy::Block
7141 ));
7142 }
7143
7144 #[test]
7145 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
7146 assert!(!round_is_clean(
7147 1,
7148 true,
7149 1,
7150 2,
7151 0,
7152 IncompleteReviewPolicy::Warn
7153 ));
7154 }
7155
7156 #[test]
7157 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
7158 assert!(round_is_clean(
7159 0,
7160 true,
7161 1,
7162 2,
7163 0,
7164 IncompleteReviewPolicy::Warn
7165 ));
7166 }
7167
7168 #[test]
7169 fn a_full_panel_with_an_open_finding_is_not_clean() {
7170 assert!(!round_is_clean(
7171 1,
7172 true,
7173 2,
7174 2,
7175 0,
7176 IncompleteReviewPolicy::Block
7177 ));
7178 }
7179
7180 #[test]
7181 fn a_full_panel_with_a_red_e2e_is_not_clean() {
7182 assert!(!round_is_clean(
7183 0,
7184 false,
7185 2,
7186 2,
7187 0,
7188 IncompleteReviewPolicy::Block
7189 ));
7190 }
7191
7192 #[test]
7199 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
7200 assert!(round_is_clean(
7203 0,
7204 true,
7205 1,
7206 2,
7207 1,
7208 IncompleteReviewPolicy::Block
7209 ));
7210 }
7211
7212 #[test]
7213 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
7214 assert!(!round_is_clean(
7217 0,
7218 true,
7219 1,
7220 2,
7221 0,
7222 IncompleteReviewPolicy::Block
7223 ));
7224 }
7225
7226 #[test]
7227 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
7228 assert!(!round_is_clean(
7229 1,
7230 true,
7231 1,
7232 2,
7233 1,
7234 IncompleteReviewPolicy::Block
7235 ));
7236 assert!(!round_is_clean(
7237 0,
7238 false,
7239 1,
7240 2,
7241 1,
7242 IncompleteReviewPolicy::Block
7243 ));
7244 }
7245
7246 #[test]
7247 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
7248 assert!(!round_is_clean(
7252 0,
7253 true,
7254 0,
7255 2,
7256 2,
7257 IncompleteReviewPolicy::Block
7258 ));
7259 }
7260
7261 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
7262 CommandOutcome {
7263 command: "test".to_owned(),
7264 code,
7265 output_tail: String::new(),
7266 duration_ms: 0,
7267 resource_blocked,
7268 }
7269 }
7270
7271 #[test]
7272 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
7273 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
7274 assert!(
7275 !verify_inconclusive(&[outcome(Some(1), false)]),
7276 "an ordinary failure is still evidence about the patch"
7277 );
7278 assert!(verify_inconclusive(&[outcome(None, true)]));
7279 assert!(
7280 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
7281 "one inconclusive outcome taints the whole batch"
7282 );
7283 assert!(!verify_inconclusive(&[]));
7284 }
7285
7286 #[tokio::test]
7287 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
7288 let calls = std::sync::atomic::AtomicUsize::new(0);
7292 let started = Instant::now();
7293 wait_for_pids_with(
7294 &[123],
7295 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
7296 Duration::from_millis(5),
7297 Duration::from_secs(5),
7298 )
7299 .await;
7300 assert!(
7301 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
7302 "must keep checking rather than deciding on the first answer"
7303 );
7304 assert!(
7305 started.elapsed() < Duration::from_secs(1),
7306 "must return the moment it is confirmed dead, not wait out the ceiling"
7307 );
7308 }
7309
7310 #[tokio::test]
7311 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
7312 let started = Instant::now();
7313 wait_for_pids_with(
7314 &[123],
7315 |_| true, Duration::from_millis(5),
7317 Duration::from_millis(30),
7318 )
7319 .await;
7320 let elapsed = started.elapsed();
7321 assert!(
7322 elapsed >= Duration::from_millis(30),
7323 "must not give up before its own ceiling: {elapsed:?}"
7324 );
7325 assert!(
7326 elapsed < Duration::from_secs(1),
7327 "must not wait past its own ceiling either: {elapsed:?}"
7328 );
7329 }
7330
7331 #[tokio::test]
7332 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
7333 let started = Instant::now();
7334 wait_for_pids_with(
7335 &[],
7336 |_| true,
7337 Duration::from_secs(5),
7338 Duration::from_secs(5),
7339 )
7340 .await;
7341 assert!(
7342 started.elapsed() < Duration::from_millis(200),
7343 "an empty pid list has nothing to confirm"
7344 );
7345 }
7346
7347 fn review_round(
7353 clean: bool,
7354 blocking: usize,
7355 answered: usize,
7356 expected: usize,
7357 progressed: bool,
7358 e2e_ok: bool,
7359 ) -> ReviewRound {
7360 ReviewRound {
7361 round: 1,
7362 head: "h".to_owned(),
7363 verified_head: None,
7364 verified_at: None,
7365 reviews: Vec::new(),
7366 e2e: vec![CommandOutcome {
7367 command: "test".to_owned(),
7368 code: Some(if e2e_ok { 0 } else { 1 }),
7369 output_tail: String::new(),
7370 duration_ms: 0,
7371 resource_blocked: false,
7372 }],
7373 verify_retried: false,
7374 e2e_deferred: false,
7375 e2e_defer_reason: None,
7376 fix: None,
7377 blocking,
7378 answered,
7379 expected,
7380 clean,
7381 progressed,
7382 vote_split: false,
7383 reconsideration: Vec::new(),
7384 verdict: None,
7385 }
7386 }
7387
7388 #[test]
7389 fn review_conclusion_is_none_when_nothing_has_run() {
7390 assert_eq!(review_conclusion(&[], 3), None);
7391 }
7392
7393 #[test]
7394 fn review_conclusion_is_none_while_rounds_remain() {
7395 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
7396 assert_eq!(review_conclusion(&rounds, 3), None);
7397 }
7398
7399 #[test]
7400 fn review_conclusion_is_gating_once_a_round_is_clean() {
7401 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
7402 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
7403 }
7404
7405 #[test]
7406 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
7407 let rounds = vec![
7408 review_round(false, 1, 2, 2, true, true),
7409 review_round(false, 1, 2, 2, true, true),
7410 ];
7411 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
7412 }
7413
7414 #[test]
7415 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
7416 let rounds = vec![
7417 review_round(false, 1, 2, 2, true, true),
7418 review_round(false, 1, 2, 2, true, false),
7419 ];
7420 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
7421 }
7422
7423 #[test]
7424 fn review_conclusion_stays_none_when_the_budget_is_spent_but_the_last_round_could_not_run() {
7425 let mut blocked = review_round(false, 1, 2, 2, true, false);
7432 blocked.e2e[0].resource_blocked = true;
7433 let rounds = vec![review_round(false, 1, 2, 2, true, true), blocked];
7434 assert_eq!(review_conclusion(&rounds, 2), None);
7435 }
7436
7437 #[test]
7438 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
7439 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
7441 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
7442 }
7443
7444 #[test]
7445 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
7446 let rounds = vec![
7447 review_round(false, 1, 2, 2, false, true),
7448 review_round(false, 1, 2, 2, false, true),
7449 ];
7450 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
7451 }
7452
7453 fn secs(n: u64) -> Duration {
7454 Duration::from_secs(n)
7455 }
7456
7457 fn init_repo(dir: &Path) {
7460 let run = |args: &[&str]| {
7461 let out = std::process::Command::new("git")
7462 .args(args)
7463 .current_dir(dir)
7464 .quiet()
7465 .output()
7466 .expect("spawn git");
7467 assert!(
7468 out.status.success(),
7469 "git {args:?} failed: {}",
7470 String::from_utf8_lossy(&out.stderr)
7471 );
7472 };
7473 run(&["init", "-b", "main"]);
7474 run(&["config", "user.name", "magi test"]);
7475 run(&["config", "user.email", "magi@example.com"]);
7476 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
7477 run(&["add", "-A"]);
7478 run(&["commit", "-m", "init"]);
7479 }
7480
7481 fn ask_test_home() {
7489 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
7490 }
7491
7492 fn runner_at(status: RunStatus) -> Runner {
7495 let mut state = RunState::new(
7496 PathBuf::from("/nonexistent/repo"),
7497 "main".to_owned(),
7498 "deadbeef".to_owned(),
7499 "task".to_owned(),
7500 Config::default(),
7501 );
7502 state.status = status;
7503 Runner {
7504 state,
7505 roles: ResolvedRoles {
7506 implementers: Vec::new(),
7507 judges: Vec::new(),
7508 reviewers: Vec::new(),
7509 fixer: None,
7510 conductor: conductor(),
7511 implementer_roster: Vec::new(),
7512 },
7513 sem: Arc::new(Semaphore::new(1)),
7514 pause: Pause::new(),
7515 interrupt: Pause::new(),
7516 }
7517 }
7518
7519 #[test]
7523 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
7524 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7525 let mut runner = runner_at(RunStatus::Implementing);
7526 let interrupt = Pause::new();
7527 runner.watch_interrupt(interrupt.clone());
7528
7529 interrupt.park_because("task a1b2 asked to run first");
7530
7531 assert!(runner.park_here().expect("park_here"));
7532 assert!(runner.state.parked);
7533 let last = runner.state.events.last().expect("a park event");
7534 assert_eq!(last.node, "park");
7535 assert!(
7536 last.message.contains("task a1b2 asked to run first"),
7537 "expected the interrupt reason in {:?}",
7538 last.message
7539 );
7540 }
7541
7542 #[test]
7550 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
7551 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7552 let mut runner = runner_at(RunStatus::Implementing);
7553 let shutdown = Pause::new();
7554 runner.on_pause(shutdown.clone());
7555 let interrupt = Pause::new();
7556 runner.watch_interrupt(interrupt.clone());
7557
7558 assert!(!runner.park_here().expect("park_here"));
7560 assert!(!runner.state.parked);
7561
7562 interrupt.park_because("test");
7564 assert!(!shutdown.parked());
7565 assert!(runner.park_here().expect("park_here"));
7566 }
7567
7568 #[tokio::test]
7582 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
7583 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7584 let mut runner = runner_at(RunStatus::Implementing);
7585 let interrupt = Pause::new();
7586 runner.watch_interrupt(interrupt.clone());
7587
7588 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
7589 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
7590
7591 let node = async move {
7595 started_tx.send(()).expect("send started");
7596 finish_rx.await.expect("recv finish");
7597 "node finished"
7598 };
7599
7600 let interrupter = async move {
7601 started_rx.await.expect("recv started");
7602 interrupt.park_because("higher-priority task waiting");
7604 tokio::task::yield_now().await;
7608 finish_tx.send(()).expect("send finish");
7609 };
7610
7611 let (node_result, ()) = tokio::join!(node, interrupter);
7612 assert_eq!(
7613 node_result, "node finished",
7614 "the in-flight call ran to completion"
7615 );
7616
7617 assert!(runner.park_here().expect("park_here"));
7620 assert!(runner.state.parked);
7621 }
7622
7623 #[test]
7629 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
7630 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
7631 let mut runner = runner_at(RunStatus::Judging);
7632 runner.state.config.agents = vec![conductor()];
7636 runner.state.candidates = vec![Candidate {
7637 index: 0,
7638 label: 'A',
7639 agent: "alpha".to_owned(),
7640 branch: "magi/x/A".to_owned(),
7641 worktree: PathBuf::from("/nonexistent/worktree"),
7642 summary: "did the thing".to_owned(),
7643 stat: "1 file changed".to_owned(),
7644 files: 1,
7645 commits: 1,
7646 empty: false,
7647 failed: None,
7648 verified_noop: None,
7649 duration_ms: 1234,
7650 folded: false,
7651 }];
7652 let run_id = runner.state.id.clone();
7653
7654 let interrupt = Pause::new();
7655 runner.watch_interrupt(interrupt.clone());
7656 interrupt.park_because("task c3d4 asked to run first");
7657 assert!(runner.park_here().expect("park_here"));
7658
7659 let resumed = Runner::resume(&run_id).expect("resume");
7660 assert_eq!(resumed.state.candidates.len(), 1);
7661 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
7662 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
7663 assert_eq!(resumed.state.status, runner.state.status);
7664 assert!(
7665 resumed.state.parked,
7666 "still parked until `execute` actually walks the graph again"
7667 );
7668 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
7669 }
7670
7671 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
7673 let mut q = ask::Question::new(
7674 run.to_owned(),
7675 "implement".to_owned(),
7676 "impl-A".to_owned(),
7677 "Which storage backend should the cache use?".to_owned(),
7678 String::new(),
7679 vec!["SQLite".to_owned(), "Redis".to_owned()],
7680 );
7681 store.put(&mut q).unwrap();
7682 q
7683 }
7684
7685 #[test]
7686 fn a_failed_runs_open_question_is_abandoned() {
7687 ask_test_home();
7688 let store = ask::Questions::open();
7689 let mut runner = runner_at(RunStatus::Failed);
7690 let run = runner.state.id.clone();
7691 let q = ask_open_question(&store, &run);
7692
7693 runner.settle_questions();
7694
7695 let back = store.get(&q.id).unwrap();
7696 assert!(
7697 !back.status.open(),
7698 "the seat that asked died with the run; nobody is left to read an answer"
7699 );
7700 assert!(
7701 back.detail.contains(&run) && back.detail.contains("failed"),
7702 "the reason names what the run became, not just that it is gone: {}",
7703 back.detail
7704 );
7705 }
7706
7707 #[test]
7708 fn a_merged_runs_open_question_is_abandoned_too() {
7709 ask_test_home();
7710 let store = ask::Questions::open();
7711 for status in [RunStatus::Merged, RunStatus::Ready] {
7714 let mut runner = runner_at(status);
7715 let run = runner.state.id.clone();
7716 let q = ask_open_question(&store, &run);
7717
7718 runner.settle_questions();
7719
7720 let back = store.get(&q.id).unwrap();
7721 assert!(
7722 !back.status.open(),
7723 "{status:?} run's question must not outlive the run"
7724 );
7725 }
7726 }
7727
7728 #[test]
7729 fn a_still_resumable_runs_open_question_is_left_alone() {
7730 ask_test_home();
7731 let store = ask::Questions::open();
7732 for status in [RunStatus::Blocked, RunStatus::Stalled] {
7738 let mut runner = runner_at(status);
7739 let run = runner.state.id.clone();
7740 let q = ask_open_question(&store, &run);
7741
7742 runner.settle_questions();
7743
7744 let back = store.get(&q.id).unwrap();
7745 assert!(
7746 back.status.open(),
7747 "{status:?} is still alive; the question must still be waiting"
7748 );
7749 }
7750 }
7751
7752 #[test]
7753 fn settle_questions_never_touches_an_already_answered_question() {
7754 ask_test_home();
7755 let store = ask::Questions::open();
7756 let mut runner = runner_at(RunStatus::Failed);
7757 let run = runner.state.id.clone();
7758 let mut q = ask_open_question(&store, &run);
7759 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
7760 .unwrap();
7761 store.put(&mut q).unwrap();
7762
7763 runner.settle_questions();
7768 runner.settle_questions();
7769
7770 let back = store.get(&q.id).unwrap();
7771 assert_eq!(
7772 back.status,
7773 ask::QuestionStatus::Answered,
7774 "a real answer is a decision on record, never overwritten by a sweep"
7775 );
7776 }
7777
7778 #[tokio::test]
7789 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
7790 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
7791 let tmp = tempfile::tempdir().expect("tempdir");
7792 let repo = tmp.path().join("repo");
7793 std::fs::create_dir_all(&repo).unwrap();
7794 init_repo(&repo);
7795
7796 let mut config = Config::default();
7797 config.graph.worktree_root = Some(tmp.path().join("wt"));
7798
7799 let mut state = RunState::new(
7800 repo.clone(),
7801 "main".to_owned(),
7802 "deadbeef".to_owned(),
7803 "task".to_owned(),
7804 config,
7805 );
7806 let root = state.worktree_root();
7807 let wt_a = root.join("cand-A");
7808 let wt_b = root.join("cand-B");
7809 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
7810 .await
7811 .expect("worktree A");
7812 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
7813 .await
7814 .expect("worktree B");
7815
7816 state.candidates = vec![
7817 Candidate {
7818 index: 0,
7819 label: 'A',
7820 agent: "alpha".to_owned(),
7821 branch: "magi/x/A".to_owned(),
7822 worktree: wt_a.clone(),
7823 summary: String::new(),
7824 stat: String::new(),
7825 files: 0,
7826 commits: 0,
7827 empty: false,
7828 failed: None,
7829 verified_noop: None,
7830 duration_ms: 0,
7831 folded: false,
7832 },
7833 Candidate {
7834 index: 1,
7835 label: 'B',
7836 agent: "beta".to_owned(),
7837 branch: "magi/x/B".to_owned(),
7838 worktree: wt_b.clone(),
7839 summary: String::new(),
7840 stat: String::new(),
7841 files: 0,
7842 commits: 0,
7843 empty: false,
7844 failed: None,
7845 verified_noop: None,
7846 duration_ms: 0,
7847 folded: false,
7848 },
7849 ];
7850 state.tally = Some(Tally {
7851 first_choice: BTreeMap::from([('A', 1)]),
7852 borda: BTreeMap::new(),
7853 winner: 'A',
7854 rankings: 1,
7855 unanimous_initial: true,
7856 deliberated: false,
7857 changed_votes: 0,
7858 unanimous_final: true,
7859 tie_break: None,
7860 judges: 1,
7861 present: 1,
7862 quorum: 1,
7863 met_quorum: true,
7864 uncontested: None,
7865 });
7866 state.status = RunStatus::Ready;
7867
7868 fold_run(&mut state, false, &crate::run::home())
7869 .await
7870 .expect("fold_run");
7871
7872 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
7873 assert!(
7874 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7875 "the unmerged winner's branch survives"
7876 );
7877 assert!(
7878 !state.candidates[0].folded,
7879 "the winner is not marked folded"
7880 );
7881
7882 assert!(!wt_b.exists(), "the loser's worktree is removed");
7883 assert!(
7884 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
7885 "the loser's branch is removed"
7886 );
7887 assert!(state.candidates[1].folded, "the loser is marked folded");
7888 }
7889
7890 #[tokio::test]
7893 async fn fold_run_keeps_a_branch_that_was_handed_to_a_later_run() {
7894 let tmp = tempfile::tempdir().expect("tempdir");
7895 let repo = tmp.path().join("repo");
7896 std::fs::create_dir_all(&repo).unwrap();
7897 init_repo(&repo);
7898 let home = tmp.path().join("home");
7899
7900 let mut config = Config::default();
7901 config.graph.worktree_root = Some(tmp.path().join("wt"));
7902 let mut state = RunState::new(
7903 repo.clone(),
7904 "main".to_owned(),
7905 "deadbeef".to_owned(),
7906 "task".to_owned(),
7907 config,
7908 );
7909 git::git(&repo, &["branch", "magi/x/A", "main"])
7911 .await
7912 .expect("branch");
7913 state.candidates = vec![Candidate {
7914 index: 0,
7915 label: 'A',
7916 agent: "alpha".to_owned(),
7917 branch: "magi/x/A".to_owned(),
7918 worktree: state.worktree_root().join("cand-A"),
7919 summary: String::new(),
7920 stat: String::new(),
7921 files: 0,
7922 commits: 0,
7923 empty: false,
7924 failed: None,
7925 verified_noop: None,
7926 duration_ms: 0,
7927 folded: true,
7928 }];
7929 state.released_to = Some("20260901-000000-new1".to_owned());
7930 state.released_branches = vec!["magi/x/A".to_owned()];
7931
7932 fold_run(&mut state, true, &home).await.expect("fold_run");
7933
7934 assert!(
7935 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
7936 "the handed-over branch survives a fold"
7937 );
7938 }
7939
7940 #[tokio::test]
7949 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
7950 let tmp = tempfile::tempdir().expect("tempdir");
7951 let repo = tmp.path().join("repo");
7952 std::fs::create_dir_all(&repo).unwrap();
7953 init_repo(&repo);
7954
7955 let mut config = Config::default();
7956 config.merge.mode = MergeMode::Local;
7957
7958 let mut state = RunState::new(
7959 repo.clone(),
7960 "main".to_owned(),
7961 "deadbeef".to_owned(),
7962 "task".to_owned(),
7963 config,
7964 );
7965 state.candidates = vec![Candidate {
7966 index: 0,
7967 label: 'A',
7968 agent: "alpha".to_owned(),
7969 branch: "does-not-exist".to_owned(),
7970 worktree: repo.clone(),
7971 summary: String::new(),
7972 stat: String::new(),
7973 files: 0,
7974 commits: 0,
7975 empty: false,
7976 failed: None,
7977 verified_noop: None,
7978 duration_ms: 0,
7979 folded: false,
7980 }];
7981 state.tally = Some(Tally {
7982 first_choice: BTreeMap::from([('A', 1)]),
7983 borda: BTreeMap::new(),
7984 winner: 'A',
7985 rankings: 1,
7986 unanimous_initial: true,
7987 deliberated: false,
7988 changed_votes: 0,
7989 unanimous_final: true,
7990 tie_break: None,
7991 judges: 0,
7992 present: 0,
7993 quorum: 0,
7994 met_quorum: true,
7995 uncontested: Some("only candidate A produced a change".to_owned()),
7996 });
7997 state.reviews = vec![ReviewRound {
7998 round: 1,
7999 head: "deadbeef".to_owned(),
8000 verified_head: None,
8001 verified_at: None,
8002 reviews: Vec::new(),
8003 e2e: Vec::new(),
8004 fix: None,
8005 blocking: 0,
8006 answered: 0,
8007 expected: 0,
8008 clean: true,
8009 verify_retried: false,
8010 e2e_deferred: false,
8011 e2e_defer_reason: None,
8012 progressed: false,
8013 vote_split: false,
8014 reconsideration: Vec::new(),
8015 verdict: None,
8016 }];
8017 state.gate = vec![CommandOutcome {
8018 command: "test".to_owned(),
8019 code: Some(0),
8020 output_tail: String::new(),
8021 duration_ms: 0,
8022 resource_blocked: false,
8023 }];
8024 state.gate_ran = true;
8025 state.status = RunStatus::Ready;
8030 state.merge = Some(MergeOutcome {
8031 mode: MergeMode::Local,
8032 ok: false,
8033 detail: "already concluded".to_owned(),
8034 });
8035
8036 let mut runner = Runner {
8037 state,
8038 roles: ResolvedRoles {
8039 implementers: Vec::new(),
8040 judges: Vec::new(),
8041 reviewers: Vec::new(),
8042 fixer: None,
8043 conductor: conductor(),
8044 implementer_roster: Vec::new(),
8045 },
8046 sem: Arc::new(Semaphore::new(1)),
8047 pause: Pause::new(),
8048 interrupt: Pause::new(),
8049 };
8050
8051 runner.merge().await.expect("merge");
8052
8053 assert_eq!(
8054 runner.state.status,
8055 RunStatus::Ready,
8056 "a concluded run's status must not change on reentry"
8057 );
8058 assert_eq!(
8059 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
8060 Some("already concluded"),
8061 "merge must not run again once the node already recorded an outcome"
8062 );
8063 }
8064
8065 #[tokio::test]
8074 async fn merge_refuses_a_gate_that_has_not_actually_run() {
8075 let tmp = tempfile::tempdir().expect("tempdir");
8076 let repo = tmp.path().join("repo");
8077 std::fs::create_dir_all(&repo).unwrap();
8078 init_repo(&repo);
8079
8080 let mut config = Config::default();
8081 config.merge.mode = MergeMode::Local;
8082
8083 let mut state = RunState::new(
8084 repo.clone(),
8085 "main".to_owned(),
8086 "deadbeef".to_owned(),
8087 "task".to_owned(),
8088 config,
8089 );
8090 state.candidates = vec![Candidate {
8091 index: 0,
8092 label: 'A',
8093 agent: "alpha".to_owned(),
8094 branch: "does-not-exist".to_owned(),
8095 worktree: repo.clone(),
8096 summary: String::new(),
8097 stat: String::new(),
8098 files: 0,
8099 commits: 0,
8100 empty: false,
8101 failed: None,
8102 verified_noop: None,
8103 duration_ms: 0,
8104 folded: false,
8105 }];
8106 state.tally = Some(Tally {
8107 first_choice: BTreeMap::from([('A', 1)]),
8108 borda: BTreeMap::new(),
8109 winner: 'A',
8110 rankings: 1,
8111 unanimous_initial: true,
8112 deliberated: false,
8113 changed_votes: 0,
8114 unanimous_final: true,
8115 tie_break: None,
8116 judges: 0,
8117 present: 0,
8118 quorum: 0,
8119 met_quorum: true,
8120 uncontested: Some("only candidate A produced a change".to_owned()),
8121 });
8122 state.reviews = vec![ReviewRound {
8123 round: 1,
8124 head: "deadbeef".to_owned(),
8125 verified_head: None,
8126 verified_at: None,
8127 reviews: Vec::new(),
8128 e2e: Vec::new(),
8129 fix: None,
8130 blocking: 0,
8131 answered: 0,
8132 expected: 0,
8133 clean: true,
8134 verify_retried: false,
8135 e2e_deferred: false,
8136 e2e_defer_reason: None,
8137 progressed: false,
8138 vote_split: false,
8139 reconsideration: Vec::new(),
8140 verdict: None,
8141 }];
8142 state.gate = Vec::new();
8144 state.gate_ran = false;
8145 state.status = RunStatus::Gating;
8146
8147 let mut runner = Runner {
8148 state,
8149 roles: ResolvedRoles {
8150 implementers: Vec::new(),
8151 judges: Vec::new(),
8152 reviewers: Vec::new(),
8153 fixer: None,
8154 conductor: conductor(),
8155 implementer_roster: Vec::new(),
8156 },
8157 sem: Arc::new(Semaphore::new(1)),
8158 pause: Pause::new(),
8159 interrupt: Pause::new(),
8160 };
8161
8162 runner.merge().await.expect("merge");
8163
8164 assert!(
8165 runner.state.merge.is_none(),
8166 "an empty gate must never be read as a passing one: {:?}",
8167 runner.state.merge
8168 );
8169 }
8170
8171 #[tokio::test]
8178 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
8179 let tmp = tempfile::tempdir().expect("tempdir");
8180 let repo = tmp.path().join("repo");
8181 std::fs::create_dir_all(&repo).unwrap();
8182 init_repo(&repo);
8183
8184 let config = Config::default();
8186
8187 let mut state = RunState::new(
8188 repo.clone(),
8189 "main".to_owned(),
8190 "deadbeef".to_owned(),
8191 "task".to_owned(),
8192 config,
8193 );
8194 state.candidates = vec![Candidate {
8195 index: 0,
8196 label: 'A',
8197 agent: "alpha".to_owned(),
8198 branch: "does-not-exist".to_owned(),
8199 worktree: repo.clone(),
8200 summary: String::new(),
8201 stat: String::new(),
8202 files: 0,
8203 commits: 0,
8204 empty: false,
8205 failed: None,
8206 verified_noop: None,
8207 duration_ms: 0,
8208 folded: false,
8209 }];
8210 state.tally = Some(Tally {
8211 first_choice: BTreeMap::from([('A', 1)]),
8212 borda: BTreeMap::new(),
8213 winner: 'A',
8214 rankings: 1,
8215 unanimous_initial: true,
8216 deliberated: false,
8217 changed_votes: 0,
8218 unanimous_final: true,
8219 tie_break: None,
8220 judges: 0,
8221 present: 0,
8222 quorum: 0,
8223 met_quorum: true,
8224 uncontested: Some("only candidate A produced a change".to_owned()),
8225 });
8226 state.reviews = vec![ReviewRound {
8227 round: 1,
8228 head: "deadbeef".to_owned(),
8229 verified_head: None,
8230 verified_at: None,
8231 reviews: Vec::new(),
8232 e2e: Vec::new(),
8233 fix: None,
8234 blocking: 0,
8235 answered: 0,
8236 expected: 0,
8237 clean: true,
8238 verify_retried: false,
8239 e2e_deferred: false,
8240 e2e_defer_reason: None,
8241 progressed: false,
8242 vote_split: false,
8243 reconsideration: Vec::new(),
8244 verdict: None,
8245 }];
8246
8247 let mut runner = Runner {
8248 state,
8249 roles: ResolvedRoles {
8250 implementers: Vec::new(),
8251 judges: Vec::new(),
8252 reviewers: Vec::new(),
8253 fixer: None,
8254 conductor: conductor(),
8255 implementer_roster: Vec::new(),
8256 },
8257 sem: Arc::new(Semaphore::new(1)),
8258 pause: Pause::new(),
8259 interrupt: Pause::new(),
8260 };
8261
8262 runner.gate().await.expect("gate");
8263 assert!(
8264 runner.state.gate_ran,
8265 "zero configured commands is still a real attempt, not an unrun gate"
8266 );
8267 assert!(runner.state.gate.is_empty());
8268 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
8269 assert_ne!(
8270 runner.state.status,
8271 RunStatus::Blocked,
8272 "a gate with nothing to check must not read as failed"
8273 );
8274
8275 runner.merge().await.expect("merge");
8276 assert_eq!(
8277 runner.state.status,
8278 RunStatus::Ready,
8279 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
8280 );
8281 }
8282
8283 #[tokio::test]
8294 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
8295 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8296 let home = crate::run::home();
8297
8298 let tmp = tempfile::tempdir().expect("tempdir");
8299 let repo = tmp.path().join("repo");
8300 std::fs::create_dir_all(&repo).unwrap();
8301 init_repo(&repo);
8302 let cache_dir = tmp.path().join("target");
8305
8306 let mut config = Config::default();
8307 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
8308 config.graph.timeout_verify = Some(2);
8311
8312 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8313 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8314 .expect("no io error acquiring directly")
8315 {
8316 crate::cache::AcquireOutcome::Acquired(g) => g,
8317 crate::cache::AcquireOutcome::Busy(b) => {
8318 panic!("expected the direct acquire to win the lease first: {b:?}")
8319 }
8320 };
8321
8322 let mut state = RunState::new(
8323 repo.clone(),
8324 "main".to_owned(),
8325 "deadbeef".to_owned(),
8326 "task".to_owned(),
8327 config,
8328 );
8329 state.candidates = vec![Candidate {
8330 index: 0,
8331 label: 'A',
8332 agent: "alpha".to_owned(),
8333 branch: "does-not-exist".to_owned(),
8334 worktree: repo.clone(),
8335 summary: String::new(),
8336 stat: String::new(),
8337 files: 0,
8338 commits: 0,
8339 empty: false,
8340 failed: None,
8341 verified_noop: None,
8342 duration_ms: 0,
8343 folded: false,
8344 }];
8345 state.tally = Some(Tally {
8346 first_choice: BTreeMap::from([('A', 1)]),
8347 borda: BTreeMap::new(),
8348 winner: 'A',
8349 rankings: 1,
8350 unanimous_initial: true,
8351 deliberated: false,
8352 changed_votes: 0,
8353 unanimous_final: true,
8354 tie_break: None,
8355 judges: 0,
8356 present: 0,
8357 quorum: 0,
8358 met_quorum: true,
8359 uncontested: Some("only candidate A produced a change".to_owned()),
8360 });
8361 state.reviews = vec![ReviewRound {
8362 round: 1,
8363 head: "deadbeef".to_owned(),
8364 verified_head: None,
8365 verified_at: None,
8366 reviews: Vec::new(),
8367 e2e: Vec::new(),
8368 fix: None,
8369 blocking: 0,
8370 answered: 0,
8371 expected: 0,
8372 clean: true,
8373 verify_retried: false,
8374 e2e_deferred: false,
8375 e2e_defer_reason: None,
8376 progressed: false,
8377 vote_split: false,
8378 reconsideration: Vec::new(),
8379 verdict: None,
8380 }];
8381
8382 let mut runner = Runner {
8383 state,
8384 roles: ResolvedRoles {
8385 implementers: Vec::new(),
8386 judges: Vec::new(),
8387 reviewers: Vec::new(),
8388 fixer: None,
8389 conductor: conductor(),
8390 implementer_roster: Vec::new(),
8391 },
8392 sem: Arc::new(Semaphore::new(1)),
8393 pause: Pause::new(),
8394 interrupt: Pause::new(),
8395 };
8396
8397 let started = std::time::Instant::now();
8398 runner.gate().await.expect("gate");
8399 assert!(
8400 started.elapsed() < Duration::from_secs(1),
8401 "a gate with nothing to run must never wait on a lease it never needed"
8402 );
8403 assert!(
8404 runner.state.gate_ran,
8405 "zero commands is still a real, immediate attempt"
8406 );
8407 assert!(runner.state.gate.is_empty());
8408 assert_ne!(
8409 runner.state.status,
8410 RunStatus::Blocked,
8411 "must not read as resource-blocked on a lease it never asked for"
8412 );
8413 }
8414
8415 #[tokio::test]
8425 async fn gate_records_a_running_task_entry_while_its_command_is_still_in_flight() {
8426 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8427
8428 let tmp = tempfile::tempdir().expect("tempdir");
8429 let repo = tmp.path().join("repo");
8430 std::fs::create_dir_all(&repo).unwrap();
8431 init_repo(&repo);
8432
8433 let mut config = Config::default();
8434 config.verify.gate = vec![
8435 "printf started > started.marker; i=0; while [ ! -f release.marker ] && \
8436 [ \"$i\" -lt 100 ]; do i=$((i+1)); sleep 0.05; done"
8437 .to_owned(),
8438 ];
8439
8440 let mut state = RunState::new(
8441 repo.clone(),
8442 "main".to_owned(),
8443 "deadbeef".to_owned(),
8444 "task".to_owned(),
8445 config,
8446 );
8447 let run_id = state.id.clone();
8448 state.candidates = vec![Candidate {
8449 index: 0,
8450 label: 'A',
8451 agent: "alpha".to_owned(),
8452 branch: "does-not-exist".to_owned(),
8453 worktree: repo.clone(),
8454 summary: String::new(),
8455 stat: String::new(),
8456 files: 0,
8457 commits: 0,
8458 empty: false,
8459 failed: None,
8460 verified_noop: None,
8461 duration_ms: 0,
8462 folded: false,
8463 }];
8464 state.tally = Some(Tally {
8465 first_choice: BTreeMap::from([('A', 1)]),
8466 borda: BTreeMap::new(),
8467 winner: 'A',
8468 rankings: 1,
8469 unanimous_initial: true,
8470 deliberated: false,
8471 changed_votes: 0,
8472 unanimous_final: true,
8473 tie_break: None,
8474 judges: 0,
8475 present: 0,
8476 quorum: 0,
8477 met_quorum: true,
8478 uncontested: Some("only candidate A produced a change".to_owned()),
8479 });
8480 state.reviews = vec![ReviewRound {
8481 round: 1,
8482 head: "deadbeef".to_owned(),
8483 verified_head: None,
8484 verified_at: None,
8485 reviews: Vec::new(),
8486 e2e: Vec::new(),
8487 fix: None,
8488 blocking: 0,
8489 answered: 0,
8490 expected: 0,
8491 clean: true,
8492 verify_retried: false,
8493 e2e_deferred: false,
8494 e2e_defer_reason: None,
8495 progressed: false,
8496 vote_split: false,
8497 reconsideration: Vec::new(),
8498 verdict: None,
8499 }];
8500
8501 let mut runner = Runner {
8502 state,
8503 roles: ResolvedRoles {
8504 implementers: Vec::new(),
8505 judges: Vec::new(),
8506 reviewers: Vec::new(),
8507 fixer: None,
8508 conductor: conductor(),
8509 implementer_roster: Vec::new(),
8510 },
8511 sem: Arc::new(Semaphore::new(1)),
8512 pause: Pause::new(),
8513 interrupt: Pause::new(),
8514 };
8515
8516 let started_marker = repo.join("started.marker");
8517 let release_marker = repo.join("release.marker");
8518 let poller = tokio::spawn(async move {
8519 for _ in 0..100 {
8524 if started_marker.exists()
8525 && let Ok(s) = crate::run::RunState::load(&run_id)
8526 && let Some(a) = s.active.get("gate")
8527 {
8528 std::fs::write(&release_marker, b"go").expect("release marker");
8529 return Some(a.clone());
8530 }
8531 tokio::time::sleep(Duration::from_millis(50)).await;
8532 }
8533 None
8534 });
8535
8536 runner.gate().await.expect("gate");
8537 let captured = poller.await.expect("poller task");
8538 let captured = captured.expect(
8539 "the poller never saw a `gate` task entry in run.json while the command was \
8540 still blocked on its own release marker",
8541 );
8542
8543 assert_eq!(captured.task.as_deref(), Some("gate"));
8544 assert_eq!(captured.node, "gate");
8545 assert_eq!(captured.index, Some(1));
8546 assert_eq!(captured.total, Some(1));
8547 assert!(
8548 captured
8549 .command
8550 .as_deref()
8551 .is_some_and(|c| c.contains("started.marker")),
8552 "{captured:?}"
8553 );
8554
8555 assert!(
8556 runner.state.active.is_empty(),
8557 "the entry must be cleared once the command actually finished: {:?}",
8558 runner.state.active
8559 );
8560 assert!(runner.state.gate_ran);
8561 assert!(runner.state.gate.iter().all(CommandOutcome::ok));
8562 }
8563
8564 #[tokio::test]
8577 async fn stop_reviewing_retries_a_resource_blocked_e2e_instead_of_reading_it_as_red() {
8578 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8579 let home = crate::run::home();
8580
8581 let tmp = tempfile::tempdir().expect("tempdir");
8582 let repo = tmp.path().join("repo");
8583 std::fs::create_dir_all(&repo).unwrap();
8584 init_repo(&repo);
8585 let head = crate::git::rev_parse(&repo, "HEAD")
8586 .await
8587 .expect("rev-parse");
8588 let cache_dir = tmp.path().join("target");
8591
8592 let mut config = Config::default();
8593 config.verify.e2e = vec![format!(
8594 "CARGO_TARGET_DIR='{}' test -f README.md",
8595 cache_dir.display()
8596 )];
8597 config.graph.review_rounds = 1;
8598 config.graph.timeout_verify = Some(2);
8601
8602 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8603 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8604 .expect("no io error acquiring directly")
8605 {
8606 crate::cache::AcquireOutcome::Acquired(g) => g,
8607 crate::cache::AcquireOutcome::Busy(b) => {
8608 panic!("expected the direct acquire to win the lease first: {b:?}")
8609 }
8610 };
8611
8612 let mut state = RunState::new(
8613 repo.clone(),
8614 "main".to_owned(),
8615 head.clone(),
8616 "task".to_owned(),
8617 config,
8618 );
8619 state.candidates = vec![Candidate {
8620 index: 0,
8621 label: 'A',
8622 agent: "alpha".to_owned(),
8623 branch: "does-not-exist".to_owned(),
8624 worktree: repo.clone(),
8625 summary: String::new(),
8626 stat: String::new(),
8627 files: 0,
8628 commits: 0,
8629 empty: false,
8630 failed: None,
8631 verified_noop: None,
8632 duration_ms: 0,
8633 folded: false,
8634 }];
8635 state.tally = Some(Tally {
8636 first_choice: BTreeMap::from([('A', 1)]),
8637 borda: BTreeMap::new(),
8638 winner: 'A',
8639 rankings: 1,
8640 unanimous_initial: true,
8641 deliberated: false,
8642 changed_votes: 0,
8643 unanimous_final: true,
8644 tie_break: None,
8645 judges: 0,
8646 present: 0,
8647 quorum: 0,
8648 met_quorum: true,
8649 uncontested: Some("only candidate A produced a change".to_owned()),
8650 });
8651 state.reviews = vec![ReviewRound {
8655 round: 1,
8656 head: head.clone(),
8657 verified_head: None,
8658 verified_at: None,
8659 reviews: Vec::new(),
8660 e2e: Vec::new(),
8661 fix: None,
8662 blocking: 1,
8663 answered: 1,
8664 expected: 1,
8665 clean: false,
8666 verify_retried: false,
8667 e2e_deferred: true,
8668 e2e_defer_reason: Some("1 blocking finding(s) already required a fix".to_owned()),
8669 progressed: false,
8670 vote_split: false,
8671 reconsideration: Vec::new(),
8672 verdict: None,
8673 }];
8674
8675 let mut runner = Runner {
8676 state,
8677 roles: ResolvedRoles {
8678 implementers: Vec::new(),
8679 judges: Vec::new(),
8680 reviewers: Vec::new(),
8681 fixer: None,
8682 conductor: conductor(),
8683 implementer_roster: Vec::new(),
8684 },
8685 sem: Arc::new(Semaphore::new(1)),
8686 pause: Pause::new(),
8687 interrupt: Pause::new(),
8688 };
8689
8690 let shell = runner.state.config.shell();
8691 runner
8692 .stop_reviewing("round budget spent", &shell, &repo)
8693 .await
8694 .expect("stop_reviewing");
8695
8696 let last = runner.state.reviews.last().expect("round record");
8697 assert_eq!(
8698 last.e2e_status(),
8699 E2eStatus::ResourceBlocked,
8700 "the shared cache is still held; the attempt must read as blocked, not deferred or \
8701 failed: {last:?}"
8702 );
8703 assert_eq!(
8704 last.verified_head.as_deref(),
8705 Some(head.as_str()),
8706 "which commit this attempt targeted is known even though nothing finished checking \
8707 it"
8708 );
8709 let first_attempt_at = last
8710 .verified_at
8711 .expect("when this attempt ran is known too");
8712 assert_ne!(
8713 runner.state.status,
8714 RunStatus::Blocked,
8715 "contention is evidence about the machine, not the patch — it must not settle the \
8716 run as blocked: {:?}",
8717 runner.state.status
8718 );
8719 assert!(
8720 !runner
8721 .state
8722 .events
8723 .iter()
8724 .any(|e| e.node == "review" && e.message.contains("e2e failed")),
8725 "a resource-blocked attempt must never be logged as a failed e2e: {:?}",
8726 runner.state.events
8727 );
8728
8729 runner
8734 .stop_reviewing("round budget spent", &shell, &repo)
8735 .await
8736 .expect("stop_reviewing retry");
8737 assert_eq!(
8738 runner.state.reviews.len(),
8739 1,
8740 "no new round was started: {:?}",
8741 runner.state.reviews
8742 );
8743 let last = runner.state.reviews.last().expect("round record");
8744 assert_eq!(last.e2e_status(), E2eStatus::ResourceBlocked, "{last:?}");
8745 assert!(
8746 last.verified_at.expect("still known") > first_attempt_at,
8747 "a second reentry must be a fresh attempt, not a stale copy of the first"
8748 );
8749 assert_ne!(runner.state.status, RunStatus::Blocked);
8750
8751 held.release();
8752 }
8753
8754 #[tokio::test]
8766 async fn a_resumed_review_loop_retries_a_last_round_left_resource_blocked() {
8767 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8768 let home = crate::run::home();
8769
8770 let tmp = tempfile::tempdir().expect("tempdir");
8771 let repo = tmp.path().join("repo");
8772 std::fs::create_dir_all(&repo).unwrap();
8773 init_repo(&repo);
8774 let head = crate::git::rev_parse(&repo, "HEAD")
8775 .await
8776 .expect("rev-parse");
8777 let cache_dir = tmp.path().join("target");
8778
8779 let mut config = Config::default();
8780 config.verify.e2e = vec![format!(
8781 "CARGO_TARGET_DIR='{}' test -f README.md",
8782 cache_dir.display()
8783 )];
8784 config.graph.review_rounds = 1;
8785 config.graph.timeout_verify = Some(2);
8786
8787 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
8788 let held = match crate::cache::try_acquire(&home, &cache_dir, &other)
8789 .expect("no io error acquiring directly")
8790 {
8791 crate::cache::AcquireOutcome::Acquired(g) => g,
8792 crate::cache::AcquireOutcome::Busy(b) => {
8793 panic!("expected the direct acquire to win the lease first: {b:?}")
8794 }
8795 };
8796
8797 let mut state = RunState::new(
8798 repo.clone(),
8799 "main".to_owned(),
8800 head.clone(),
8801 "task".to_owned(),
8802 config,
8803 );
8804 state.candidates = vec![Candidate {
8805 index: 0,
8806 label: 'A',
8807 agent: "alpha".to_owned(),
8808 branch: "does-not-exist".to_owned(),
8809 worktree: repo.clone(),
8810 summary: String::new(),
8811 stat: String::new(),
8812 files: 0,
8813 commits: 0,
8814 empty: false,
8815 failed: None,
8816 verified_noop: None,
8817 duration_ms: 0,
8818 folded: false,
8819 }];
8820 state.tally = Some(Tally {
8821 first_choice: BTreeMap::from([('A', 1)]),
8822 borda: BTreeMap::new(),
8823 winner: 'A',
8824 rankings: 1,
8825 unanimous_initial: true,
8826 deliberated: false,
8827 changed_votes: 0,
8828 unanimous_final: true,
8829 tie_break: None,
8830 judges: 0,
8831 present: 0,
8832 quorum: 0,
8833 met_quorum: true,
8834 uncontested: Some("only candidate A produced a change".to_owned()),
8835 });
8836 state.reviews = vec![ReviewRound {
8840 round: 1,
8841 head: head.clone(),
8842 verified_head: Some(head.clone()),
8843 verified_at: Some(jiff::Timestamp::now()),
8844 reviews: Vec::new(),
8845 e2e: vec![CommandOutcome {
8846 command: format!(
8847 "CARGO_TARGET_DIR='{}' test -f README.md",
8848 cache_dir.display()
8849 ),
8850 code: None,
8851 output_tail: "waiting for the shared build cache".to_owned(),
8852 duration_ms: 0,
8853 resource_blocked: true,
8854 }],
8855 fix: None,
8856 blocking: 1,
8857 answered: 1,
8858 expected: 1,
8859 clean: false,
8860 verify_retried: false,
8861 e2e_deferred: false,
8862 e2e_defer_reason: None,
8863 progressed: false,
8864 vote_split: false,
8865 reconsideration: Vec::new(),
8866 verdict: None,
8867 }];
8868
8869 let first_attempt_at = state.reviews[0].verified_at.expect("set above");
8870 let mut runner = Runner {
8871 state,
8872 roles: ResolvedRoles {
8873 implementers: Vec::new(),
8874 judges: Vec::new(),
8875 reviewers: Vec::new(),
8876 fixer: None,
8877 conductor: conductor(),
8878 implementer_roster: Vec::new(),
8879 },
8880 sem: Arc::new(Semaphore::new(1)),
8881 pause: Pause::new(),
8882 interrupt: Pause::new(),
8883 };
8884
8885 runner.review_loop().await.expect("review_loop");
8890
8891 assert_eq!(
8892 runner.state.reviews.len(),
8893 1,
8894 "no new round was started on top of the unresolved one: {:?}",
8895 runner.state.reviews
8896 );
8897 let last = &runner.state.reviews[0];
8898 assert_eq!(
8899 last.e2e_status(),
8900 E2eStatus::ResourceBlocked,
8901 "still contended: {last:?}"
8902 );
8903 assert!(
8904 last.verified_at.expect("still known") > first_attempt_at,
8905 "review_loop must have actually retried the check, not left it exactly as found"
8906 );
8907 assert_ne!(
8908 runner.state.status,
8909 RunStatus::Blocked,
8910 "a resumed run must not read leftover contention as a verdict on the patch: {:?}",
8911 runner.state.status
8912 );
8913
8914 held.release();
8915 }
8916
8917 #[tokio::test]
8918 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
8919 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
8920 let tmp = tempfile::tempdir().expect("tempdir");
8921 let repo = tmp.path().join("repo");
8922 std::fs::create_dir_all(&repo).unwrap();
8923 init_repo(&repo);
8924
8925 let mut config = Config::default();
8926 config.merge.mode = MergeMode::Pr;
8927 config.graph.land = true;
8928 config.graph.land_approval = false;
8929
8930 let mut state = RunState::new(
8931 repo.clone(),
8932 "main".to_owned(),
8933 "deadbeef".to_owned(),
8934 "task".to_owned(),
8935 config,
8936 );
8937 state.candidates = vec![Candidate {
8938 index: 0,
8939 label: 'A',
8940 agent: "alpha".to_owned(),
8941 branch: "does-not-exist".to_owned(),
8942 worktree: repo.clone(),
8943 summary: String::new(),
8944 stat: String::new(),
8945 files: 0,
8946 commits: 0,
8947 empty: false,
8948 failed: None,
8949 verified_noop: None,
8950 duration_ms: 0,
8951 folded: false,
8952 }];
8953 state.tally = Some(Tally {
8954 first_choice: BTreeMap::from([('A', 1)]),
8955 borda: BTreeMap::new(),
8956 winner: 'A',
8957 rankings: 1,
8958 unanimous_initial: true,
8959 deliberated: false,
8960 changed_votes: 0,
8961 unanimous_final: true,
8962 tie_break: None,
8963 judges: 0,
8964 present: 0,
8965 quorum: 0,
8966 met_quorum: true,
8967 uncontested: Some("only candidate A produced a change".to_owned()),
8968 });
8969 state.reviews = vec![ReviewRound {
8970 round: 1,
8971 head: "deadbeef".to_owned(),
8972 verified_head: None,
8973 verified_at: None,
8974 reviews: Vec::new(),
8975 e2e: Vec::new(),
8976 fix: None,
8977 blocking: 0,
8978 answered: 0,
8979 expected: 0,
8980 clean: true,
8981 verify_retried: false,
8982 e2e_deferred: false,
8983 e2e_defer_reason: None,
8984 progressed: false,
8985 vote_split: false,
8986 reconsideration: Vec::new(),
8987 verdict: None,
8988 }];
8989 state.gate = vec![CommandOutcome {
8990 command: "test".to_owned(),
8991 code: Some(0),
8992 output_tail: String::new(),
8993 duration_ms: 0,
8994 resource_blocked: false,
8995 }];
8996 state.gate_ran = true;
8997 state.status = RunStatus::Landing;
9001 state.merge = Some(MergeOutcome {
9002 mode: MergeMode::Pr,
9003 ok: true,
9004 detail: "https://example.invalid/x/y/pull/1".to_owned(),
9005 });
9006
9007 ask_test_home();
9011 let store = ask::Questions::open();
9012 let q = ask_open_question(&store, &state.id);
9013
9014 let mut runner = Runner {
9015 state,
9016 roles: ResolvedRoles {
9017 implementers: Vec::new(),
9018 judges: Vec::new(),
9019 reviewers: Vec::new(),
9020 fixer: None,
9021 conductor: conductor(),
9022 implementer_roster: Vec::new(),
9023 },
9024 sem: Arc::new(Semaphore::new(1)),
9025 pause: Pause::new(),
9026 interrupt: Pause::new(),
9027 };
9028
9029 runner.execute().await.expect("execute");
9034
9035 assert_eq!(
9036 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
9037 Some("https://example.invalid/x/y/pull/1"),
9038 "reentry must not push again or open a second pull request over the \
9039 one `land` is already watching"
9040 );
9041 assert_ne!(
9042 runner.state.status,
9043 RunStatus::Landing,
9044 "land could not actually reach the fake pull request, so it must \
9045 have given up rather than left the run silently parked forever"
9046 );
9047 assert_eq!(runner.state.status, RunStatus::Blocked);
9051 assert!(
9052 store.get(&q.id).unwrap().status.open(),
9053 "Blocked is still alive; settle_questions must have been a no-op here"
9054 );
9055 }
9056
9057 fn state_with_round(round: ReviewRound) -> RunState {
9058 let mut s = RunState::new(
9059 PathBuf::from("/repo"),
9060 "main".to_owned(),
9061 "abc1234".to_owned(),
9062 "add retries".to_owned(),
9063 Config::default(),
9064 );
9065 s.reviews = vec![round];
9066 s
9067 }
9068
9069 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
9070 crate::verdict::Finding {
9071 id: id.to_owned(),
9072 severity,
9073 file: None,
9074 line: None,
9075 title: title.to_owned(),
9076 detail: String::new(),
9077 }
9078 }
9079
9080 #[test]
9081 fn pr_body_names_open_findings_and_declined_ones() {
9082 let round = ReviewRound {
9083 round: 2,
9084 head: "deadbee".to_owned(),
9085 verified_head: None,
9086 verified_at: None,
9087 reviews: vec![ReviewRecord {
9088 attempts: 0,
9089 reviewer: 1,
9090 agent: "alpha".to_owned(),
9091 summary: String::new(),
9092 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
9093 vote: None,
9094 failed: None,
9095 duration_ms: 0,
9096 }],
9097 e2e: vec![CommandOutcome {
9098 command: "cargo test".to_owned(),
9099 code: Some(0),
9100 output_tail: String::new(),
9101 duration_ms: 0,
9102 resource_blocked: false,
9103 }],
9104 verify_retried: false,
9105 e2e_deferred: false,
9106 e2e_defer_reason: None,
9107 fix: Some(FixRecord {
9108 agent: "alpha".to_owned(),
9109 addressed: Vec::new(),
9110 rejected: vec![crate::verdict::Rejection {
9111 id: "R1-1-1".to_owned(),
9112 why: "not reachable from any caller".to_owned(),
9113 }],
9114 notes: String::new(),
9115 committed: true,
9116 failed: None,
9117 duration_ms: 0,
9118 continuation: None,
9119 }),
9120 blocking: 0,
9121 answered: 1,
9122 expected: 1,
9123 clean: false,
9124 progressed: true,
9125 vote_split: false,
9126 reconsideration: Vec::new(),
9127 verdict: None,
9128 };
9129 let state = state_with_round(round);
9130 let body = pr_message(&state, 'A').body;
9131
9132 assert!(body.contains("add retries"), "the task must still be there");
9133 assert!(body.contains("R2-1-1"), "{body}");
9134 assert!(body.contains("unused import"), "{body}");
9135 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
9136 assert!(
9137 body.contains("not reachable from any caller"),
9138 "the reason it was declined: {body}"
9139 );
9140 }
9141
9142 #[test]
9143 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
9144 let round = ReviewRound {
9145 round: 1,
9146 head: "deadbee".to_owned(),
9147 verified_head: None,
9148 verified_at: None,
9149 reviews: vec![ReviewRecord {
9150 attempts: 0,
9151 reviewer: 1,
9152 agent: "alpha".to_owned(),
9153 summary: String::new(),
9154 findings: Vec::new(),
9155 vote: None,
9156 failed: None,
9157 duration_ms: 0,
9158 }],
9159 e2e: Vec::new(),
9160 verify_retried: false,
9161 e2e_deferred: false,
9162 e2e_defer_reason: None,
9163 fix: None,
9164 blocking: 0,
9165 answered: 1,
9166 expected: 1,
9167 clean: true,
9168 progressed: false,
9169 vote_split: false,
9170 reconsideration: Vec::new(),
9171 verdict: None,
9172 };
9173 let state = state_with_round(round);
9174 let body = pr_message(&state, 'A').body;
9175 assert!(!body.contains("Open review findings"), "{body}");
9176 assert!(!body.contains("Declined"), "{body}");
9177 }
9178
9179 fn state_with_summary(instruction: &str, summary: &str) -> RunState {
9180 let mut state = RunState::new(
9181 PathBuf::from("/repo"),
9182 "main".to_owned(),
9183 "abc1234".to_owned(),
9184 instruction.to_owned(),
9185 Config::default(),
9186 );
9187 state.candidates.push(Candidate {
9188 index: 0,
9189 label: 'A',
9190 agent: "alpha".to_owned(),
9191 branch: "magi/x/A".to_owned(),
9192 worktree: PathBuf::from("/wt"),
9193 summary: summary.to_owned(),
9194 stat: String::new(),
9195 files: 1,
9196 commits: 1,
9197 empty: false,
9198 failed: None,
9199 verified_noop: None,
9200 folded: false,
9201 duration_ms: 0,
9202 });
9203 state
9204 }
9205
9206 #[test]
9207 fn pr_message_describes_the_change_not_the_task() {
9208 let state = state_with_summary(
9209 "今回やってほしいこと: results projector を直す",
9210 "TITLE: fix(web): batch the runs list reads\n- reads run.json once\n- risk: none",
9211 );
9212 let m = pr_message(&state, 'A');
9213 assert_eq!(m.title, "fix(web): batch the runs list reads");
9214 assert!(
9215 m.body.starts_with("## Summary\n\n- reads run.json once"),
9216 "{}",
9217 m.body
9218 );
9219 assert!(!m.body.contains("TITLE:"), "{}", m.body);
9220 let task_at = m.body.find("今回やってほしいこと").unwrap();
9221 let details_at = m.body.find("<details>").unwrap();
9222 assert!(
9223 details_at < task_at,
9224 "the task lives inside <details>: {}",
9225 m.body
9226 );
9227 assert!(m.body.contains(&format!("magi:run/{}", state.id)));
9228 assert!(m.body.contains("magi:candidate-a"));
9229 }
9230
9231 #[test]
9232 fn pr_message_falls_back_to_the_task_without_a_title_line() {
9233 let state = state_with_summary("\n\nadd retries\n\ndetails", "- did some things");
9234 let m = pr_message(&state, 'A');
9235 assert_eq!(m.title, "add retries");
9236 assert!(
9237 m.body.contains("## Summary\n\n- did some things"),
9238 "{}",
9239 m.body
9240 );
9241
9242 let none = RunState::new(
9243 PathBuf::from("/repo"),
9244 "main".to_owned(),
9245 "abc1234".to_owned(),
9246 "add retries".to_owned(),
9247 Config::default(),
9248 );
9249 let m = pr_message(&none, 'A');
9250 assert_eq!(m.title, "add retries");
9251 assert!(!m.body.contains("## Summary"), "{}", m.body);
9252 }
9253
9254 #[test]
9255 fn pr_message_refuses_the_candidate_commit_subject() {
9256 for bad in [
9257 "TITLE: magi: candidate A (uncommitted work)",
9258 "TITLE: chore: stuff (uncommitted work)",
9259 "TITLE: ",
9260 ] {
9261 let state = state_with_summary("add retries", bad);
9262 assert_eq!(pr_message(&state, 'A').title, "add retries", "{bad}");
9263 }
9264 }
9265
9266 #[test]
9267 fn pr_message_bounds_a_very_long_task_and_title() {
9268 let long = format!("fix the thing 🎉 {}", "x".repeat(5000));
9269 let state = state_with_summary(&long, "- nothing");
9270 let m = pr_message(&state, 'A');
9271 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9272 assert!(!m.title.contains('\n'));
9273
9274 let state = state_with_summary("task", &format!("TITLE: feat: {}", "y".repeat(5000)));
9275 let m = pr_message(&state, 'A');
9276 assert!(m.title.starts_with("feat: "));
9277 assert!(m.title.chars().count() <= PR_TITLE_MAX, "{}", m.title);
9278 assert_eq!(m.commit_message().lines().next(), Some(m.title.as_str()));
9279 }
9280
9281 #[test]
9282 fn pr_message_magi_text_is_english_and_the_task_is_verbatim() {
9283 let mut state = state_with_summary(
9287 "add retries",
9288 "TITLE: fix(web): batch reads\n- reads run.json once",
9289 );
9290 state.config.graph.language = "ja".to_owned();
9291 let m = pr_message(&state, 'A');
9292 assert!(m.title.is_ascii() && m.body.is_ascii(), "{}", m.body);
9293
9294 let task = "今回やってほしいこと: results projector を直す";
9297 let mut state = state_with_summary(task, "- no title line");
9298 state.config.graph.language = "ja".to_owned();
9299 let m = pr_message(&state, 'A');
9300 assert_eq!(
9301 m.title,
9302 format!("chore: land candidate A of run {}", state.id)
9303 );
9304 assert!(
9305 m.body.contains(&format!(
9306 "<summary>Original task</summary>\n\n{task}\n\n</details>"
9307 )),
9308 "{}",
9309 m.body
9310 );
9311 }
9312
9313 #[test]
9314 fn pr_message_scrubs_home_paths_and_addresses() {
9315 let state = state_with_summary(
9316 "fix it in /Users/someone/src/x",
9317 "TITLE: fix(x): y\n- edited /home/someone/repo/src/a.rs on 10.1.2.3",
9318 );
9319 let m = pr_message(&state, 'A');
9320 for leak in ["/Users/someone", "/home/someone", "10.1.2.3"] {
9321 assert!(!m.body.contains(leak), "{}", m.body);
9322 }
9323 assert!(m.body.contains("~/repo/src/a.rs"), "{}", m.body);
9324 }
9325
9326 #[test]
9327 fn pr_message_survives_a_task_that_closes_details() {
9328 let state = state_with_summary("a </details> b", "TITLE: fix: x");
9329 let m = pr_message(&state, 'A');
9330 assert_eq!(m.body.matches("</details>").count(), 1, "{}", m.body);
9331 }
9332
9333 #[test]
9334 fn manual_squash_subject_cannot_break_out_of_its_quotes() {
9335 let cmd = manual_merge_command(
9336 MergeStyle::Squash,
9337 Path::new("/repo"),
9338 "b",
9339 "fix: \"quoted\" $(x) `y`\n\nbody",
9340 );
9341 assert!(cmd.ends_with("commit -m \"fix: quoted (x) y\""), "{cmd}");
9342 }
9343
9344 #[test]
9345 fn manual_merge_command_matches_the_configured_style() {
9346 let repo = Path::new("/repo");
9347 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
9348
9349 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
9350 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
9351
9352 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
9353 assert_eq!(
9354 squash,
9355 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
9356 \"Merge magi run 0832 (candidate A)\""
9357 );
9358
9359 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
9360 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
9361 }
9362
9363 #[test]
9364 fn a_nudge_gets_a_quarter_of_the_budget() {
9365 assert_eq!(retry_budget(secs(1200), true), secs(300));
9367 assert_eq!(retry_budget(secs(3600), true), secs(900));
9368 }
9369
9370 #[test]
9371 fn a_resent_prompt_keeps_the_whole_budget() {
9372 assert_eq!(retry_budget(secs(1200), false), secs(1200));
9375 assert_eq!(retry_budget(secs(60), false), secs(60));
9376 }
9377
9378 #[test]
9379 fn the_floor_never_exceeds_the_original_budget() {
9380 assert_eq!(retry_budget(secs(60), true), secs(60));
9384 assert_eq!(retry_budget(secs(480), true), secs(120));
9385 assert_eq!(retry_budget(secs(0), true), secs(0));
9386 }
9387
9388 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
9389 agent::CommandEvidence {
9390 id: "item1".to_owned(),
9391 description: "cargo test".to_owned(),
9392 exit_code,
9393 result_summary: String::new(),
9394 source: "codex".to_owned(),
9395 }
9396 }
9397
9398 #[test]
9399 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
9400 assert!(!has_unconfirmed_command(&[]));
9404 }
9405
9406 #[test]
9407 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
9408 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
9412 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
9413 assert!(!has_unconfirmed_command(&[
9414 evidence(Some(0)),
9415 evidence(Some(101))
9416 ]));
9417 }
9418
9419 #[test]
9420 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
9421 assert!(has_unconfirmed_command(&[
9422 evidence(Some(0)),
9423 evidence(None)
9424 ]));
9425 }
9426
9427 #[test]
9428 fn a_clean_usable_reply_with_the_marker_is_a_verified_claim() {
9429 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9430 assert_eq!(
9431 verified_noop_claim(true, &[], text).as_deref(),
9432 Some("already fixed by b32cfc4, on main.")
9433 );
9434 }
9435
9436 #[test]
9437 fn an_unusable_reply_never_earns_the_benefit_of_the_doubt() {
9438 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9441 assert!(verified_noop_claim(false, &[], text).is_none());
9442 }
9443
9444 #[test]
9445 fn an_unconfirmed_command_disqualifies_the_claim_even_on_a_usable_reply() {
9446 let text = "NO CHANGE NEEDED: already fixed by b32cfc4, on main.";
9447 assert!(verified_noop_claim(true, &[evidence(None)], text).is_none());
9448 assert!(verified_noop_claim(true, &[evidence(Some(0))], text).is_some());
9450 }
9451
9452 #[test]
9453 fn an_ordinary_reply_with_no_marker_is_never_a_claim() {
9454 assert!(verified_noop_claim(true, &[], "- did the thing\n- tested it").is_none());
9455 }
9456
9457 fn set_candidates(runner: &mut Runner, shape: &[(bool, Option<&str>)]) {
9460 runner.state.candidates = shape
9461 .iter()
9462 .enumerate()
9463 .map(|(i, &(empty, verified))| Candidate {
9464 index: i,
9465 label: (b'A' + i as u8) as char,
9466 agent: "sonnet".to_owned(),
9467 branch: format!("magi/x/{}", (b'A' + i as u8) as char),
9468 worktree: PathBuf::from(format!("/wt/{i}")),
9469 summary: String::new(),
9470 stat: String::new(),
9471 files: 0,
9472 commits: 0,
9473 empty,
9474 failed: None,
9475 verified_noop: verified.map(str::to_owned),
9476 duration_ms: 0,
9477 folded: false,
9478 })
9479 .collect();
9480 }
9481
9482 #[test]
9483 fn after_implement_reads_all_candidates_verified_as_a_noop_not_a_failure() {
9484 ask_test_home();
9485 let mut runner = runner_at(RunStatus::Implementing);
9486 set_candidates(
9487 &mut runner,
9488 &[
9489 (true, Some("already on main at b32cfc4")),
9490 (true, Some("same fix, see the existing test")),
9491 ],
9492 );
9493
9494 runner
9495 .after_implement()
9496 .expect("a verified no-op is not an error");
9497
9498 assert_eq!(runner.state.status, RunStatus::VerifiedNoop);
9499 }
9500
9501 #[test]
9502 fn after_implement_does_not_accept_one_candidates_claim_next_to_an_ordinary_loss() {
9503 ask_test_home();
9504 let mut runner = runner_at(RunStatus::Implementing);
9505 set_candidates(
9509 &mut runner,
9510 &[(true, Some("already on main at b32cfc4")), (true, None)],
9511 );
9512
9513 let err = runner
9514 .after_implement()
9515 .expect_err("an unverified empty candidate must still fail the run");
9516
9517 assert!(
9518 err.to_string().contains("no candidate produced a change"),
9519 "{err}"
9520 );
9521 assert_eq!(runner.state.status, RunStatus::Failed);
9522 }
9523
9524 #[test]
9525 fn after_implement_still_fails_an_ordinary_all_empty_run() {
9526 ask_test_home();
9527 let mut runner = runner_at(RunStatus::Implementing);
9528 set_candidates(&mut runner, &[(true, None), (true, None)]);
9529
9530 let err = runner
9531 .after_implement()
9532 .expect_err("no candidate declared anything; this is an ordinary failure");
9533
9534 assert!(
9535 err.to_string().contains("no candidate produced a change"),
9536 "{err}"
9537 );
9538 assert_eq!(runner.state.status, RunStatus::Failed);
9539 }
9540}