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::run::{
45 BaseSync, Candidate, CommandOutcome, ContinuationOutcome, ContinuationRecord,
46 DeliberationRound, DeliberationTurn, FixRecord, JobRecord, JobStatus, Judgement, MergeOutcome,
47 QuotaLoss, ReviewRecord, ReviewRevoteRecord, ReviewRound, RunState, RunStatus, Tally,
48 VoteRecord, tail, write_artifact,
49};
50use crate::verdict::{
51 self, FinalVote, FixReport, Position, Proposal, Ranking, Review, ReviewRevote, ReviewVote,
52 Severity,
53};
54
55const OUTPUT_TAIL: usize = 8_000;
57
58const EVENT_OUTPUT_TAIL: usize = 2_000;
61
62const LEASE_RELEASE_POLL: Duration = Duration::from_secs(1);
65
66const LEASE_RELEASE_MAX_WAIT: Duration = Duration::from_secs(30);
83
84pub(crate) const STAGNANT_LIMIT: usize = 2;
98
99const BASE_SYNC_ROUNDS: usize = 4;
112
113const MAX_FIX_CONTINUATIONS: usize = 2;
129
130#[derive(Clone)]
136struct SeatJob {
137 spec: AgentSpec,
138 seat: SeatState,
139 cwd: PathBuf,
140 prompt: String,
141 timeout: Duration,
142 allow_write: bool,
143 sessions: bool,
144 artifacts: PathBuf,
145 stem: String,
146}
147
148enum AgentOutcome {
160 Ok(AgentOutput),
162 Quota(AgentOutput),
164 Dropped(AgentOutput),
167 Failed(String),
169}
170
171#[derive(Debug, Clone, Default)]
198pub struct Pause(Arc<AtomicBool>, Arc<Mutex<Option<String>>>);
199
200impl Pause {
201 #[must_use]
203 pub fn new() -> Self {
204 Self::default()
205 }
206
207 pub fn park(&self) {
209 self.0.store(true, Ordering::SeqCst);
210 }
211
212 pub fn park_because(&self, reason: impl Into<String>) {
218 let mut reason_guard = self
219 .1
220 .lock()
221 .unwrap_or_else(std::sync::PoisonError::into_inner);
222 if reason_guard.is_none() {
223 *reason_guard = Some(reason.into());
224 }
225 drop(reason_guard);
226 self.park();
227 }
228
229 #[must_use]
231 pub fn parked(&self) -> bool {
232 self.0.load(Ordering::SeqCst)
233 }
234
235 #[must_use]
237 pub fn reason(&self) -> Option<String> {
238 self.1
239 .lock()
240 .unwrap_or_else(std::sync::PoisonError::into_inner)
241 .clone()
242 }
243}
244
245pub struct Runner {
247 pub state: RunState,
249 roles: ResolvedRoles,
250 sem: Arc<Semaphore>,
251 pause: Pause,
255 interrupt: Pause,
261}
262
263async fn resolve_base(repo: &Path, base_branch: &str, remote: &str) -> Result<String> {
283 let tracking = format!("{remote}/{base_branch}");
284 let fetched = git::fetch(repo, remote, base_branch).await;
285 if let Ok(out) = &fetched
286 && out.ok()
287 && git::rev_exists(repo, &tracking).await
288 {
289 return git::rev_parse(repo, &tracking).await;
290 }
291 let why = match &fetched {
292 Ok(out) if !out.ok() => out.stderr.lines().next().unwrap_or("").to_owned(),
293 Ok(_) => format!("{remote} has no {base_branch}"),
294 Err(e) => e.to_string(),
295 };
296 tracing::warn!(
297 "could not read {tracking} ({why}); branching off the local \
298 {base_branch} instead, which may be behind"
299 );
300 git::rev_parse(repo, base_branch).await.with_context(|| {
301 format!(
302 "cannot resolve `{base_branch}`; set [merge] base in magi.toml to a \
303 branch that exists"
304 )
305 })
306}
307
308impl Runner {
309 pub async fn start(repo: &Path, instruction: String, config: Config) -> Result<Self> {
311 let repo = git::toplevel(repo).await?;
312 let missing = agent::missing_programs(&config.agents);
313 if !missing.is_empty() {
314 bail!(
315 "these agent programs are not on PATH: {}. Fix the roster in \
316 magi.toml or install them.",
317 missing.join(", ")
318 );
319 }
320 let base_branch = match config.merge.base.clone() {
321 Some(b) => b,
322 None => git::current_branch(&repo)
323 .await?
324 .context("HEAD is detached; set [merge] base in magi.toml")?,
325 };
326 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
327 if !git::is_clean(&repo).await? {
331 tracing::warn!(
332 "{} has uncommitted changes; they are not part of this run, \
333 which branches off {base_branch} ({})",
334 repo.display(),
335 &base_commit[..base_commit.len().min(8)]
336 );
337 }
338 let roles = config.resolve_roles()?;
339 let max_parallel = config.graph.max_parallel.max(1);
340 let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
341 state.event("start", format!("run {} created", state.id));
342 state.save()?;
343 Ok(Self {
344 state,
345 roles,
346 sem: Arc::new(Semaphore::new(max_parallel)),
347 pause: Pause::new(),
348 interrupt: Pause::new(),
349 })
350 }
351
352 pub async fn review(repo: &Path, branch: &str, config: Config) -> Result<Self> {
366 let repo = git::toplevel(repo).await?;
367 let missing = agent::missing_programs(&config.agents);
368 if !missing.is_empty() {
369 bail!(
370 "these agent programs are not on PATH: {}. Fix the roster in \
371 magi.toml or install them.",
372 missing.join(", ")
373 );
374 }
375 if !git::branch_exists(&repo, branch).await? {
376 bail!("no branch `{branch}` in {}", repo.display());
377 }
378 let base_branch = match config.merge.base.clone() {
379 Some(b) => b,
380 None => git::current_branch(&repo)
381 .await?
382 .context("HEAD is detached; set [merge] base in magi.toml")?,
383 };
384 if base_branch == branch {
385 bail!("`{branch}` is the base branch; there is nothing to review against");
386 }
387 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
388
389 let roles = config.resolve_roles()?;
390 let max_parallel = config.graph.max_parallel.max(1);
391 let log = git::log_oneline(&repo, &base_commit, branch)
394 .await
395 .unwrap_or_default();
396 let instruction = format!(
397 "Review the work already on branch `{branch}`. There is no task \
398 statement: what the change claims to do is whatever its commits \
399 say.\n\n{}",
400 if log.trim().is_empty() {
401 "(no commit messages)"
402 } else {
403 log.trim()
404 }
405 );
406 let mut state = RunState::new(
407 repo.clone(),
408 base_branch,
409 base_commit.clone(),
410 instruction,
411 config,
412 );
413
414 let worktree = state.worktree_root().join("under-review");
417 if let Some(parent) = worktree.parent() {
418 tokio::fs::create_dir_all(parent).await.ok();
419 }
420 let path = worktree.to_string_lossy().to_string();
421 git::git(&repo, &["worktree", "add", &path, branch])
422 .await
423 .with_context(|| {
424 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
425 })?;
426
427 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
428 .await
429 .unwrap_or(0);
430 if commits == 0 {
431 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
432 }
433 let files = git::changed_files(&worktree, &base_commit, "HEAD")
434 .await
435 .map(|f| f.len())
436 .unwrap_or(0);
437 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
438 .await
439 .unwrap_or_default();
440
441 state.candidates.push(Candidate {
442 index: 0,
443 label: 'A',
444 agent: "(existing branch)".to_owned(),
447 branch: branch.to_owned(),
448 worktree,
449 summary: String::new(),
450 stat,
451 files,
452 commits,
453 empty: false,
454 failed: None,
455 duration_ms: 0,
456 folded: false,
457 });
458 state.tally = Some(Tally {
459 first_choice: BTreeMap::from([('A', 0)]),
460 borda: BTreeMap::new(),
461 winner: 'A',
462 rankings: 0,
463 unanimous_initial: false,
464 deliberated: false,
465 changed_votes: 0,
466 unanimous_final: false,
467 tie_break: None,
468 judges: 0,
472 present: 0,
473 quorum: 0,
474 met_quorum: true,
475 uncontested: Some("review-only run: nothing competed".to_owned()),
476 });
477 state.status = RunStatus::Reviewing;
478 state.event(
479 "start",
480 format!(
481 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
482 state.id
483 ),
484 );
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 fn resume(id: &str) -> Result<Self> {
497 let state = RunState::load(id)?;
498 let roles = state.config.resolve_roles()?;
499 let max_parallel = state.config.graph.max_parallel.max(1);
500 Ok(Self {
501 state,
502 roles,
503 sem: Arc::new(Semaphore::new(max_parallel)),
504 pause: Pause::new(),
505 interrupt: Pause::new(),
506 })
507 }
508
509 pub async fn execute(&mut self) -> Result<()> {
511 self.state.parked = false;
516 if self.state.clear_active() {
523 self.state.save()?;
524 }
525 if self.state.status == RunStatus::Stalled {
538 if self.recover_stall().await? {
539 self.finish_after_tally().await?;
540 } else {
541 self.state.save()?;
543 }
544 return Ok(());
545 }
546 if self.state.status == RunStatus::Landing {
556 self.run_land().await?;
557 self.settle_questions();
562 return Ok(());
563 }
564 self.prep().await?;
565 if self.park_here()? {
566 return Ok(());
567 }
568 self.advise().await?;
569 if self.park_here()? {
570 return Ok(());
571 }
572 self.implement().await?;
573 if self.park_here()? {
574 return Ok(());
575 }
576 self.judge().await?;
577 if self.park_here()? {
578 return Ok(());
579 }
580 self.deliberate().await?;
581 if self.park_here()? {
582 return Ok(());
583 }
584 self.vote().await?;
585 if self.park_here()? {
586 return Ok(());
587 }
588 self.tally()?;
589 if self.state.status == RunStatus::Stalled {
594 self.state.save()?;
598 return Ok(());
599 }
600 self.finish_after_tally().await?;
601 Ok(())
602 }
603
604 fn park_here(&mut self) -> Result<bool> {
611 if !self.pause.parked() && !self.interrupt.parked() {
616 return Ok(false);
617 }
618 let why = match self.interrupt.reason().or_else(|| self.pause.reason()) {
619 Some(reason) => format!(
620 "parked after `{}` ({reason}) — resume to carry on from here",
621 self.state.status.as_str()
622 ),
623 None => format!(
624 "parked after `{}` — resume to carry on from here",
625 self.state.status.as_str()
626 ),
627 };
628 self.state.event("park", why);
629 self.state.parked = true;
630 self.state.save()?;
631 Ok(true)
632 }
633
634 pub fn on_pause(&mut self, pause: Pause) {
636 self.pause = pause;
637 }
638
639 pub fn watch_interrupt(&mut self, pause: Pause) {
645 self.interrupt = pause;
646 }
647
648 fn settle_questions(&mut self) {
666 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
667 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
668 }
669 }
670
671 async fn finish_after_tally(&mut self) -> Result<()> {
674 self.fold_losers().await?;
675 self.sync_to_base().await?;
680 self.review_loop().await?;
681 self.sync_to_base().await?;
682 self.gate().await?;
683 self.merge().await?;
684 self.state.save()?;
685 Ok(())
686 }
687
688 async fn prep(&mut self) -> Result<()> {
691 if !self.state.candidates.is_empty() {
692 return Ok(());
693 }
694 self.state.status = RunStatus::Prep;
695 let repo = self.state.repo.clone();
696 let base = self.state.base_commit.clone();
697 let root = self.state.worktree_root();
698 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
699
700 let hooks_dir = self.state.dir().join("hooks");
703 if self.state.config.blind.commit_msg_hook {
704 std::fs::create_dir_all(&hooks_dir)
705 .with_context(|| format!("create {}", hooks_dir.display()))?;
706 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
707 let path = hooks_dir.join("commit-msg");
708 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
709 make_executable(&path)?;
710 git::acquire_worktree_config(&repo).await?;
718 self.state.enabled_worktree_config = true;
719 }
720
721 for (index, (spec, label)) in self
722 .roles
723 .implementers
724 .clone()
725 .into_iter()
726 .zip(labels)
727 .enumerate()
728 {
729 let branch = self.state.branch_for(label);
730 let worktree = root.join(format!("cand-{label}"));
731 git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
732 if self.state.config.blind.commit_msg_hook {
733 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
734 }
735 git::local_exclude(&worktree, "/.magi/").await?;
736 self.state.candidates.push(Candidate {
737 index,
738 label,
739 agent: spec.id.clone(),
740 branch,
741 worktree,
742 summary: String::new(),
743 stat: String::new(),
744 files: 0,
745 commits: 0,
746 empty: false,
747 failed: None,
748 duration_ms: 0,
749 folded: false,
750 });
751 }
752
753 for j in 1..=self.roles.judges.len() {
754 let wt = root.join(format!("judge-{j}"));
755 if !wt.exists() {
756 git::worktree_add_detached(&repo, &wt, &base).await?;
757 }
758 }
759
760 if self.state.config.graph.advise {
768 for k in 1..=self.state.config.graph.advisors {
769 let wt = root.join(format!("advisor-{k}"));
770 if !wt.exists() {
771 git::worktree_add_detached(&repo, &wt, &base).await?;
772 }
773 }
774 }
775
776 let authors: Vec<&str> = self
781 .roles
782 .implementers
783 .iter()
784 .map(|a| a.id.as_str())
785 .collect();
786 let overlap: Vec<String> = self
787 .roles
788 .judges
789 .iter()
790 .enumerate()
791 .filter(|(_, j)| authors.contains(&j.id.as_str()))
792 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
793 .collect();
794 if !overlap.is_empty() {
795 let note = format!(
796 "{} also authored a candidate; blind, but the panel is less \
797 independent than {} distinct agents would be",
798 overlap.join(", "),
799 self.roles.judges.len()
800 );
801 self.state.event("prep", note);
802 }
803
804 self.state.event(
805 "prep",
806 format!(
807 "{} candidates, {} judges, base {} ({})",
808 self.state.candidates.len(),
809 self.roles.judges.len(),
810 &self.state.base_commit[..7.min(self.state.base_commit.len())],
811 self.state.base_branch
812 ),
813 );
814 self.state.status = RunStatus::Implementing;
815 self.state.save()?;
816 Ok(())
817 }
818
819 async fn advise(&mut self) -> Result<()> {
852 let implement_untouched = self
853 .state
854 .candidates
855 .iter()
856 .all(|c| c.commits == 0 && c.failed.is_none() && !c.empty);
857 if !self.state.config.graph.advise || self.state.advise_attempted {
858 return Ok(());
859 }
860 if !implement_untouched {
861 self.state.event(
862 "advise",
863 "skipping the design-deliberation stage: at least one \
864 candidate already shows implementation progress, so this \
865 run is past the point the stage exists to run before"
866 .to_owned(),
867 );
868 self.state.advise_attempted = true;
869 self.state.save()?;
870 return Ok(());
871 }
872 let run_id = self.state.id.clone();
873 let prompts = self.state.config.prompts.clone();
874 let instruction = self.state.instruction.clone();
875 let language = self.state.config.graph.language.clone();
876 let root = self.state.worktree_root();
877 let n = self.state.config.graph.advisors;
878 let where_recorded = self.state.dir().join("run.json");
879
880 let seats = match self.state.config.advisors() {
881 Ok(seats) if !seats.is_empty() => seats,
882 Ok(_) => {
883 self.state.event(
884 "advise",
885 format!(
886 "[graph] advisors is 0; skipping the design-deliberation \
887 stage and continuing without a synthesis brief (see {})",
888 where_recorded.display()
889 ),
890 );
891 self.state.advise_attempted = true;
892 self.state.save()?;
893 return Ok(());
894 }
895 Err(e) => {
896 self.state.event(
897 "advise",
898 format!(
899 "could not resolve advisor seats ({e:#}); continuing \
900 without a design-deliberation brief (see {})",
901 where_recorded.display()
902 ),
903 );
904 self.state.advise_attempted = true;
905 self.state.save()?;
906 return Ok(());
907 }
908 };
909
910 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
911 let artifacts = agent::artifacts_dir(&self.state.dir());
912 let worktrees: Vec<PathBuf> = (1..=n).map(|k| root.join(format!("advisor-{k}"))).collect();
913
914 let mut jobs = Vec::new();
915 for (i, spec) in seats.iter().cloned().enumerate() {
916 let seat_key = format!("advisor-{}", i + 1);
917 let seat = self.seat(&seat_key, &spec.id);
918 jobs.push(SeatJob {
919 prompt: prompt::advisor(&instruction, i + 1, seats.len(), &language),
920 spec,
921 seat,
922 cwd: worktrees[i % worktrees.len()].clone(),
923 timeout,
924 allow_write: false,
925 sessions: false,
926 artifacts: artifacts.clone(),
927 stem: seat_key,
928 });
929 }
930
931 self.state.event(
932 "advise",
933 format!(
934 "{} advisor seat(s) sketching a design in parallel",
935 jobs.len()
936 ),
937 );
938 let mut quota_losses = Vec::new();
939 let cache = self.state.config.cache_dir();
940 let ctx = WaveCtx {
941 run: &run_id,
942 node: "advise",
943 prompts: &prompts,
944 cache: cache.as_deref(),
945 };
946 let results = ask_json_wave::<Proposal>(
947 jobs,
948 Arc::clone(&self.sem),
949 self.state.config.graph.retries,
950 &ctx,
951 &mut quota_losses,
952 &mut self.state,
953 &|p: &Proposal| p.validate(),
954 )
955 .await;
956 self.state.quota.extend(quota_losses);
957
958 let mut records = Vec::with_capacity(results.len());
959 for (i, (seat, res)) in results.into_iter().enumerate() {
960 let agent_id = seat.agent.clone();
961 self.state.seats.insert(seat.key.clone(), seat);
962 match res {
963 Ok((proposal, out)) => {
964 self.state
965 .event("advise", format!("advisor-{} proposed a design", i + 1));
966 records.push(advise::AdvisorRecord::proposed(
967 i + 1,
968 agent_id,
969 proposal,
970 out.duration_ms,
971 ));
972 }
973 Err(e) => {
974 self.state.event(
975 "advise",
976 format!("advisor-{} produced no usable proposal: {e:#}", i + 1),
977 );
978 records.push(advise::AdvisorRecord::failed(
979 i + 1,
980 agent_id,
981 e.to_string(),
982 ));
983 }
984 }
985 }
986
987 let mut advice = advise::Advice {
988 records,
989 synthesis: None,
990 };
991 if advice.proposals().is_empty() {
992 self.state.event(
993 "advise",
994 "no advisor produced a usable proposal; continuing without a \
995 synthesis brief"
996 .to_owned(),
997 );
998 } else {
999 match self
1000 .synthesize_brief(
1001 &advice,
1002 &instruction,
1003 &language,
1004 &worktrees[0],
1005 &artifacts,
1006 &run_id,
1007 &prompts,
1008 cache.as_deref(),
1009 )
1010 .await
1011 {
1012 Ok(Some(text)) => {
1013 self.state.event(
1014 "advise",
1015 "synthesized a design brief for the implementer".to_owned(),
1016 );
1017 advice.synthesis = Some(text);
1018 }
1019 Ok(None) => {
1020 self.state.event(
1021 "advise",
1022 "the synthesis seat produced nothing usable; continuing \
1023 without a design brief"
1024 .to_owned(),
1025 );
1026 }
1027 Err(e) => {
1028 self.state.event(
1029 "advise",
1030 format!("could not synthesize a design brief: {e:#}"),
1031 );
1032 }
1033 }
1034 }
1035 advise::apply_reflection(&mut advice);
1036
1037 self.state.advice = Some(advice);
1038 self.state.advise_attempted = true;
1039 self.state.save()?;
1040 Ok(())
1041 }
1042
1043 #[allow(clippy::too_many_arguments)]
1054 async fn synthesize_brief(
1055 &mut self,
1056 advice: &advise::Advice,
1057 instruction: &str,
1058 language: &str,
1059 cwd: &Path,
1060 artifacts: &Path,
1061 run_id: &str,
1062 prompts: &Prompts,
1063 cache: Option<&Path>,
1064 ) -> Result<Option<String>> {
1065 let spec = agent::pick(&self.state.config.agents, None, &agent::installed)?;
1066 let mut seat = self.seat("advise-synthesis", &spec.id);
1067 let proposals = advice.proposals();
1068 let mut prompt = prompt::with_overlay(
1069 prompt::synthesize_brief(instruction, &proposals, language),
1070 prompts.overlay("advise"),
1071 );
1072 if cache.is_some() {
1073 prompt.push('\n');
1078 prompt.push_str(&prompt::build_cache_note("advise", false));
1079 }
1080 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge.max(1));
1081 let out = agent::invoke(
1082 &spec,
1083 &mut seat,
1084 &Invocation {
1085 cwd,
1086 prompt: &prompt,
1087 timeout,
1088 allow_write: false,
1089 sessions: false,
1090 artifacts,
1091 stem: "advise-synthesis",
1092 run: run_id,
1093 node: "advise",
1094 cache_dir: None,
1095 attachments: &[],
1096 },
1097 )
1098 .await?;
1099 self.state.seats.insert(seat.key.clone(), seat);
1100 if !out.usable() {
1101 return Ok(None);
1102 }
1103 let text =
1104 verdict::section(&out.text, "synthesis").unwrap_or_else(|| out.text.trim().to_owned());
1105 Ok((!text.trim().is_empty()).then_some(text))
1106 }
1107
1108 async fn implement(&mut self) -> Result<()> {
1111 let run_id = self.state.id.clone();
1116 let prompts = self.state.config.prompts.clone();
1117 let todo: Vec<usize> = self
1118 .state
1119 .candidates
1120 .iter()
1121 .enumerate()
1122 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
1123 .map(|(i, _)| i)
1124 .collect();
1125 if todo.is_empty() {
1126 return self.after_implement();
1127 }
1128 self.state.status = RunStatus::Implementing;
1129
1130 let language = self.state.config.graph.language.clone();
1131 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
1132 let sessions = self.state.config.graph.sessions;
1133 let artifacts = agent::artifacts_dir(&self.state.dir());
1134 let brief = self
1138 .state
1139 .advice
1140 .as_ref()
1141 .and_then(|a| a.synthesis.as_deref())
1142 .map(str::to_owned);
1143
1144 let mut jobs = Vec::new();
1145 for &i in &todo {
1146 let (index, label, worktree) = {
1147 let c = &self.state.candidates[i];
1148 (c.index, c.label, c.worktree.clone())
1149 };
1150 let spec = self.roles.implementers[index].clone();
1151 let seat_key = format!("impl-{label}");
1152 let seat = self.seat(&seat_key, &spec.id);
1153 let instruction = self.state.instruction.clone();
1154 jobs.push(SeatJob {
1155 spec,
1156 seat,
1157 prompt: prompt::implement(
1158 &instruction,
1159 &worktree.to_string_lossy(),
1160 &language,
1161 brief.as_deref(),
1162 ),
1163 cwd: worktree,
1164 timeout,
1165 allow_write: true,
1166 sessions,
1167 artifacts: artifacts.clone(),
1168 stem: format!("impl-{label}"),
1169 });
1170 }
1171
1172 self.state.event(
1173 "implement",
1174 format!("{} candidates in parallel", jobs.len()),
1175 );
1176 let sent = jobs.clone();
1179 let cache = self.state.config.cache_dir();
1180 let ctx = WaveCtx {
1181 run: &run_id,
1182 node: "implement",
1183 prompts: &prompts,
1184 cache: cache.as_deref(),
1185 };
1186 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1187 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
1188 .await;
1189 self.resume_unconfirmed_commands(&mut results, &sent, &prompts, &run_id)
1190 .await;
1191
1192 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
1193 let seat_key = seat.key.clone();
1194 self.state.seats.insert(seat.key.clone(), seat);
1195 let label = self.state.candidates[i].label;
1196 let worktree = self.state.candidates[i].worktree.clone();
1197 let base = self.state.base_commit.clone();
1198
1199 let (summary, duration, failed) = match out {
1200 AgentOutcome::Ok(o) => {
1201 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
1202 let failed = (!o.usable()).then(|| {
1203 if o.timed_out {
1204 "agent timed out".to_owned()
1205 } else {
1206 format!("agent exited with {:?}", o.exit_code)
1207 }
1208 });
1209 (text, o.duration_ms, failed)
1210 }
1211 AgentOutcome::Dropped(o) => {
1217 let why = o
1218 .dropped
1219 .as_ref()
1220 .map(|d| d.why.as_str())
1221 .unwrap_or("the CLI ended the stream without delivering its answer");
1222 (
1223 String::new(),
1224 o.duration_ms,
1225 Some(format!("the CLI dropped the stream ({why})")),
1226 )
1227 }
1228 AgentOutcome::Quota(o) => {
1229 self.state.quota.push(QuotaLoss {
1230 seat: seat_key,
1231 node: "implement".to_owned(),
1232 at: Timestamp::now(),
1233 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1234 });
1235 (
1236 String::new(),
1237 o.duration_ms,
1238 Some("rate limited (quota); produced no change".to_owned()),
1239 )
1240 }
1241 AgentOutcome::Failed(e) => (String::new(), 0, Some(e)),
1242 };
1243
1244 let rescued = git::commit_all(
1247 &worktree,
1248 &format!("magi: candidate {label} (uncommitted work)"),
1249 )
1250 .await
1251 .unwrap_or(false);
1252 let commits = git::commits_ahead(&worktree, &base, "HEAD")
1253 .await
1254 .unwrap_or(0);
1255 let patch = git::diff(&worktree, &base, "HEAD")
1256 .await
1257 .unwrap_or_default();
1258 let stat = git::diff_stat(&worktree, &base, "HEAD")
1259 .await
1260 .unwrap_or_default();
1261 let files = git::changed_files(&worktree, &base, "HEAD")
1262 .await
1263 .map(|f| f.len())
1264 .unwrap_or(0);
1265 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
1266
1267 let c = &mut self.state.candidates[i];
1268 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
1269 c.stat = stat;
1270 c.files = files;
1271 c.commits = commits;
1272 c.duration_ms = duration;
1273 c.empty = commits == 0 || patch.trim().is_empty();
1274 c.failed = match failed {
1277 Some(_) if c.empty => failed,
1278 _ => None,
1279 };
1280 let note = match (&c.failed, c.empty, rescued) {
1281 (Some(e), _, _) => format!("candidate {label}: {e}"),
1282 (None, true, _) => format!("candidate {label}: no change produced"),
1283 (None, false, true) => {
1284 format!(
1285 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
1286 )
1287 }
1288 (None, false, false) => {
1289 format!("candidate {label}: {files} files, {commits} commits")
1290 }
1291 };
1292 self.state.event("implement", note);
1293 self.state.save()?;
1294 }
1295
1296 self.after_implement()
1297 }
1298
1299 async fn resume_undelivered(
1327 &mut self,
1328 results: &mut [(usize, SeatState, AgentOutcome)],
1329 sent: &[SeatJob],
1330 prompts: &Prompts,
1331 run_id: &str,
1332 ) {
1333 for (wi, seat, out) in results.iter_mut() {
1334 let Some(dropped) = (match &*out {
1335 AgentOutcome::Dropped(o) => o.dropped.clone(),
1336 _ => None,
1337 }) else {
1338 continue;
1339 };
1340 let Some(job) = sent.get(*wi) else { continue };
1341 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
1343 self.state.event(
1344 "implement",
1345 format!(
1346 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
1347 work is in the tree",
1348 seat.key, dropped.output_tokens, dropped.why
1349 ),
1350 );
1351 continue;
1352 }
1353 if !has_context(&job.spec, seat, job.sessions) {
1361 self.state.event(
1362 "implement",
1363 format!(
1364 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
1365 is no session left to resume",
1366 seat.key, dropped.output_tokens, dropped.why
1367 ),
1368 );
1369 continue;
1370 }
1371 self.state.event(
1372 "implement",
1373 format!(
1374 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
1375 conversation",
1376 seat.key, dropped.output_tokens, dropped.why
1377 ),
1378 );
1379 let mut retry = job.clone();
1380 retry.seat = seat.clone();
1381 retry.prompt = prompt::resume_after_drop(&dropped.why);
1382 retry.timeout = retry_budget(job.timeout, true);
1383 retry.stem = format!("{}-resume", job.stem);
1384 let cache = self.state.config.cache_dir();
1385 let ctx = WaveCtx {
1386 run: run_id,
1387 node: "implement",
1388 prompts,
1389 cache: cache.as_deref(),
1390 };
1391 let (resumed_seat, resumed) =
1392 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1393 *seat = resumed_seat;
1394 *out = resumed;
1395 }
1396 }
1397
1398 async fn resume_unconfirmed_commands(
1422 &mut self,
1423 results: &mut [(usize, SeatState, AgentOutcome)],
1424 sent: &[SeatJob],
1425 prompts: &Prompts,
1426 run_id: &str,
1427 ) {
1428 for (wi, seat, out) in results.iter_mut() {
1429 let AgentOutcome::Ok(o) = &*out else {
1430 continue;
1431 };
1432 if !has_unconfirmed_command(&o.commands) {
1433 continue;
1434 }
1435 let Some(job) = sent.get(*wi) else { continue };
1436 if !has_context(&job.spec, seat, job.sessions) {
1437 self.state.event(
1438 "implement",
1439 format!(
1440 "{}: the reply named a command whose own CLI never confirmed the exit \
1441 status of, but there is no session left to resume",
1442 seat.key
1443 ),
1444 );
1445 continue;
1446 }
1447 self.state.event(
1448 "implement",
1449 format!(
1450 "{}: the reply named a command whose own CLI never confirmed the exit \
1451 status of; resuming the conversation",
1452 seat.key
1453 ),
1454 );
1455 let mut retry = job.clone();
1456 retry.seat = seat.clone();
1457 retry.prompt = prompt::resume_incomplete(
1458 "a command in your last reply had no confirmed exit status",
1459 );
1460 retry.timeout = retry_budget(job.timeout, true);
1461 retry.stem = format!("{}-confirm", job.stem);
1462 let cache = self.state.config.cache_dir();
1463 let ctx = WaveCtx {
1464 run: run_id,
1465 node: "implement",
1466 prompts,
1467 cache: cache.as_deref(),
1468 };
1469 let (resumed_seat, resumed) =
1470 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
1471 *seat = resumed_seat;
1472 *out = resumed;
1473 }
1474 }
1475
1476 async fn continue_fix_report(
1497 &mut self,
1498 mut seat: SeatState,
1499 parse_err: String,
1500 job: &SeatJob,
1501 prompts: &Prompts,
1502 run_id: &str,
1503 round: usize,
1504 ) -> (
1505 SeatState,
1506 Option<FixReport>,
1507 Option<String>,
1508 ContinuationRecord,
1509 ) {
1510 let mut last_err = parse_err;
1511 let mut cumulative_wait_ms = 0u64;
1512 let mut attempts = 0usize;
1513 loop {
1514 if !has_context(&job.spec, &seat, job.sessions) {
1515 self.state.event(
1516 "fix",
1517 format!(
1518 "round {round}: fixer's reply had no adoption report ({last_err}); no \
1519 session left to resume into"
1520 ),
1521 );
1522 let outcome = if attempts == 0 {
1523 ContinuationOutcome::NoSession
1524 } else {
1525 ContinuationOutcome::Exhausted
1526 };
1527 return (
1528 seat,
1529 None,
1530 Some(format!("unparsable fix report: {last_err}")),
1531 ContinuationRecord {
1532 attempts,
1533 cumulative_wait_ms,
1534 outcome,
1535 },
1536 );
1537 }
1538 if attempts >= MAX_FIX_CONTINUATIONS {
1539 self.state.event(
1540 "fix",
1541 format!(
1542 "round {round}: fixer's reply still had no adoption report after \
1543 {attempts} continuation(s) ({last_err}); giving up"
1544 ),
1545 );
1546 return (
1547 seat,
1548 None,
1549 Some(format!(
1550 "unparsable fix report after {attempts} continuation(s): {last_err}"
1551 )),
1552 ContinuationRecord {
1553 attempts,
1554 cumulative_wait_ms,
1555 outcome: ContinuationOutcome::Exhausted,
1556 },
1557 );
1558 }
1559 attempts += 1;
1560 self.state.event(
1561 "fix",
1562 format!(
1563 "round {round}: fixer's reply had no adoption report ({last_err}); resuming \
1564 the conversation (attempt {attempts}/{MAX_FIX_CONTINUATIONS})"
1565 ),
1566 );
1567 let mut retry = job.clone();
1568 retry.seat = seat.clone();
1569 retry.prompt = prompt::resume_incomplete(&last_err);
1570 retry.timeout = retry_budget(job.timeout, true);
1571 retry.stem = format!("{}-continue{attempts}", job.stem);
1572 let cache = self.state.config.cache_dir();
1573 let ctx = WaveCtx {
1574 run: run_id,
1575 node: "fix",
1576 prompts,
1577 cache: cache.as_deref(),
1578 };
1579 let (resumed_seat, resumed_out) = run_one(
1580 retry,
1581 Arc::clone(&self.sem),
1582 &ctx,
1583 &mut self.state,
1584 attempts,
1585 )
1586 .await;
1587 seat = resumed_seat;
1588 match resumed_out {
1589 AgentOutcome::Ok(o) => {
1590 cumulative_wait_ms += o.duration_ms;
1591 match verdict::extract_json::<FixReport>(&o.text) {
1592 Ok(report) if !has_unconfirmed_command(&o.commands) => {
1593 self.state.event(
1594 "fix",
1595 format!(
1596 "round {round}: fixer's adoption report recovered after \
1597 {attempts} continuation(s)"
1598 ),
1599 );
1600 return (
1601 seat,
1602 Some(report),
1603 None,
1604 ContinuationRecord {
1605 attempts,
1606 cumulative_wait_ms,
1607 outcome: ContinuationOutcome::Resumed,
1608 },
1609 );
1610 }
1611 Ok(_) => {
1619 last_err = "the reply parsed, but it reported a command whose own CLI \
1620 never confirmed an exit status"
1621 .to_owned();
1622 }
1623 Err(e) => last_err = e.to_string(),
1624 }
1625 }
1626 AgentOutcome::Quota(o) => {
1627 cumulative_wait_ms += o.duration_ms;
1628 self.state.quota.push(QuotaLoss {
1629 seat: seat.key.clone(),
1630 node: "fix".to_owned(),
1631 at: Timestamp::now(),
1632 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1633 });
1634 self.state.event(
1635 "fix",
1636 format!(
1637 "round {round}: continuation rate limited (quota); not retrying now"
1638 ),
1639 );
1640 return (
1641 seat,
1642 None,
1643 Some("rate limited (quota) while recovering the fix report".to_owned()),
1644 ContinuationRecord {
1645 attempts,
1646 cumulative_wait_ms,
1647 outcome: ContinuationOutcome::QuotaLost,
1648 },
1649 );
1650 }
1651 AgentOutcome::Dropped(o) => {
1652 cumulative_wait_ms += o.duration_ms;
1653 let why = o
1654 .dropped
1655 .as_ref()
1656 .map(|d| d.why.as_str())
1657 .unwrap_or("the CLI ended the stream without delivering its answer");
1658 last_err = format!("the CLI dropped the stream ({why})");
1659 }
1660 AgentOutcome::Failed(e) => last_err = e,
1661 }
1662 }
1663 }
1664
1665 fn after_implement(&mut self) -> Result<()> {
1666 if self.state.leaks.is_empty() {
1668 let cfg = self.state.config.blind.clone();
1669 let mut leaks = Vec::new();
1670 for c in &self.state.candidates {
1671 let Some(patch) =
1672 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
1673 else {
1674 continue;
1675 };
1676 leaks.extend(blind::scan(
1677 &format!("candidate {} patch", c.label),
1678 &patch,
1679 &cfg.vendor_tokens,
1680 ));
1681 }
1682 if !leaks.is_empty() {
1683 let summary = leaks
1684 .iter()
1685 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
1686 .collect::<Vec<_>>()
1687 .join(", ");
1688 match cfg.on_leak {
1689 LeakPolicy::Fail => {
1690 self.state.status = RunStatus::Failed;
1691 self.state
1692 .event("blind", format!("vendor text in a patch: {summary}"));
1693 self.state.leaks = leaks;
1694 self.state.save()?;
1695 self.settle_questions();
1696 bail!(
1697 "blind.on_leak = \"fail\" and vendor text reached a \
1698 judged patch: {summary}"
1699 );
1700 }
1701 LeakPolicy::Redact => self.state.event(
1702 "blind",
1703 format!("redacting vendor text for judging: {summary}"),
1704 ),
1705 LeakPolicy::Warn => self.state.event(
1706 "blind",
1707 format!("vendor text present in a judged patch (shown as-is): {summary}"),
1708 ),
1709 }
1710 self.state.leaks = leaks;
1711 }
1712 }
1713
1714 if self.state.viable().is_empty() {
1715 self.state.status = RunStatus::Failed;
1716 self.state.save()?;
1717 self.settle_questions();
1718 bail!("no candidate produced a change; nothing to judge");
1719 }
1720 self.state.status = RunStatus::Judging;
1721 self.state.save()?;
1722 Ok(())
1723 }
1724
1725 async fn judge(&mut self) -> Result<()> {
1728 let run_id = self.state.id.clone();
1733 let prompts = self.state.config.prompts.clone();
1734 if !self.state.judgements.is_empty() || self.state.judge_skipped {
1735 return Ok(());
1736 }
1737 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
1738 if viable.len() == 1 {
1739 self.state.judge_skipped = true;
1746 self.state.event(
1747 "judge",
1748 format!(
1749 "only candidate {} produced a change; judging skipped",
1750 viable[0].label
1751 ),
1752 );
1753 self.state.save()?;
1754 return Ok(());
1755 }
1756 self.state.status = RunStatus::Judging;
1757
1758 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
1759 let language = self.state.config.graph.language.clone();
1760 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
1761 let sessions = self.state.config.graph.sessions;
1762 let artifacts = agent::artifacts_dir(&self.state.dir());
1763 let root = self.state.worktree_root();
1764 let base_short = short(&self.state.base_commit);
1765
1766 let mut jobs = Vec::new();
1767 let mut orders = Vec::new();
1768 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
1769 let order = blind::presentation_order(viable.len(), j, self.state.seed);
1770 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
1771 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
1772 let seat_key = format!("judge-{}", j + 1);
1773 let seat = self.seat(&seat_key, &spec.id);
1774 jobs.push(SeatJob {
1775 prompt: prompt::judge(
1776 &self.state.instruction,
1777 &views,
1778 self.roles.judges.len(),
1779 &base_short,
1780 &language,
1781 ),
1782 spec,
1783 seat,
1784 cwd: root.join(format!("judge-{}", j + 1)),
1785 timeout,
1786 allow_write: false,
1787 sessions,
1788 artifacts: artifacts.clone(),
1789 stem: format!("judge-{}", j + 1),
1790 });
1791 }
1792
1793 self.state.event(
1794 "judge",
1795 format!(
1796 "{} judges ranking {} candidates blind",
1797 jobs.len(),
1798 viable.len()
1799 ),
1800 );
1801 let labels_for_check = labels.clone();
1802 let mut quota_losses = Vec::new();
1803 let cache = self.state.config.cache_dir();
1804 let ctx = WaveCtx {
1805 run: &run_id,
1806 node: "judge",
1807 prompts: &prompts,
1808 cache: cache.as_deref(),
1809 };
1810 let results = ask_json_wave::<Ranking>(
1811 jobs,
1812 Arc::clone(&self.sem),
1813 self.state.config.graph.retries,
1814 &ctx,
1815 &mut quota_losses,
1816 &mut self.state,
1817 &move |r: &Ranking| r.validate(&labels_for_check),
1818 )
1819 .await;
1820 self.state.quota.extend(quota_losses);
1821
1822 for (j, (seat, res)) in results.into_iter().enumerate() {
1823 let agent_id = seat.agent.clone();
1824 self.state.seats.insert(seat.key.clone(), seat);
1825 let mut record = Judgement {
1826 judge: j + 1,
1827 seat: format!("judge-{}", j + 1),
1828 agent: agent_id,
1829 ranking: Vec::new(),
1830 reasons: BTreeMap::new(),
1831 confidence: None,
1832 order: orders[j].clone(),
1833 failed: None,
1834 duration_ms: 0,
1835 };
1836 match res {
1837 Ok((ranking, out)) => {
1838 record.ranking = ranking.normalized();
1839 record.reasons = ranking.reasons;
1840 record.confidence = ranking.confidence;
1841 record.duration_ms = out.duration_ms;
1842 self.state.event(
1843 "judge",
1844 format!(
1845 "judge {} ranked {}",
1846 j + 1,
1847 record.ranking.iter().collect::<String>()
1848 ),
1849 );
1850 }
1851 Err(e) => {
1852 record.failed = Some(e.to_string());
1853 self.state
1854 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
1855 }
1856 }
1857 self.state.judgements.push(record);
1858 self.state.save()?;
1859 }
1860 Ok(())
1861 }
1862
1863 async fn deliberate(&mut self) -> Result<()> {
1866 let run_id = self.state.id.clone();
1871 let prompts = self.state.config.prompts.clone();
1872 if !self.state.deliberation.is_empty() {
1873 return Ok(());
1874 }
1875 let tops: Vec<char> = self
1876 .state
1877 .judgements
1878 .iter()
1879 .filter_map(|j| j.ranking.first().copied())
1880 .collect();
1881 let rounds = self.state.config.graph.deliberate_rounds;
1882 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
1883 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
1884 self.state.event(
1885 "deliberate",
1886 format!("judges agreed on {} outright; no deliberation", tops[0]),
1887 );
1888 }
1889 self.state.status = RunStatus::Voting;
1890 self.state.save()?;
1891 return Ok(());
1892 }
1893
1894 self.state.status = RunStatus::Deliberating;
1895 self.state.event(
1896 "deliberate",
1897 format!(
1898 "split: first choices were {} — opening {rounds} round(s)",
1899 tops.iter().collect::<String>()
1900 ),
1901 );
1902
1903 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
1904 let language = self.state.config.graph.language.clone();
1905 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
1906 let sessions = self.state.config.graph.sessions;
1907 let artifacts = agent::artifacts_dir(&self.state.dir());
1908 let root = self.state.worktree_root();
1909 let base_short = short(&self.state.base_commit);
1910
1911 for round in 1..=rounds {
1915 let mut turns: Vec<DeliberationTurn> = Vec::new();
1916 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
1917 if self.state.judgements[j].failed.is_some() {
1918 continue;
1919 }
1920 let seat_key = format!("judge-{}", j + 1);
1921 let mut seat = self.seat(&seat_key, &spec.id);
1922 let transcript = self.transcript(&turns, j);
1923 let context = if has_context(&spec, &seat, sessions) {
1924 None
1925 } else {
1926 Some(self.candidate_block(&viable, &base_short))
1927 };
1928 let text = prompt::deliberate(
1929 &self.state.instruction,
1930 context.as_deref(),
1931 &transcript,
1932 round,
1933 rounds,
1934 &language,
1935 );
1936 let job = SeatJob {
1937 spec,
1938 seat: seat.clone(),
1939 prompt: text,
1940 cwd: root.join(format!("judge-{}", j + 1)),
1941 timeout,
1942 allow_write: false,
1943 sessions,
1944 artifacts: artifacts.clone(),
1945 stem: format!("delib-{round}-judge-{}", j + 1),
1946 };
1947 let cache = self.state.config.cache_dir();
1948 let ctx = WaveCtx {
1949 run: &run_id,
1950 node: "deliberate",
1951 prompts: &prompts,
1952 cache: cache.as_deref(),
1953 };
1954 let (updated, out) =
1955 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1956 seat = updated;
1957 let agent_id = seat.agent.clone();
1958 let seat_key = seat.key.clone();
1959 self.state.seats.insert(seat.key.clone(), seat);
1960 let body = match out {
1961 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
1962 AgentOutcome::Dropped(o) => {
1966 let why =
1967 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
1968 "the CLI ended the stream without delivering its answer",
1969 );
1970 self.state.event(
1971 "deliberate",
1972 format!(
1973 "judge {} skipped: the CLI dropped the stream ({why})",
1974 j + 1
1975 ),
1976 );
1977 continue;
1978 }
1979 AgentOutcome::Quota(o) => {
1980 self.state.quota.push(QuotaLoss {
1981 seat: seat_key,
1982 node: "deliberate".to_owned(),
1983 at: Timestamp::now(),
1984 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1985 });
1986 self.state.event(
1987 "deliberate",
1988 format!("judge {} skipped: rate limited (quota)", j + 1),
1989 );
1990 continue;
1991 }
1992 AgentOutcome::Failed(e) => {
1993 self.state
1994 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
1995 continue;
1996 }
1997 };
1998 let tentative = verdict::extract_json::<Position>(&body)
1999 .ok()
2000 .and_then(|p| p.tentative)
2001 .and_then(|s| s.trim().chars().next())
2002 .map(|c| c.to_ascii_uppercase());
2003 self.state.event(
2004 "deliberate",
2005 format!(
2006 "round {round}: judge {} now favours {}",
2007 j + 1,
2008 tentative.map_or("—".to_owned(), |c| c.to_string())
2009 ),
2010 );
2011 turns.push(DeliberationTurn {
2012 judge: j + 1,
2013 agent: agent_id,
2014 body: blind::sanitize_prose(&body, &self.state.config.blind),
2015 tentative,
2016 });
2017 }
2018 self.state
2019 .deliberation
2020 .push(DeliberationRound { round, turns });
2021 self.state.save()?;
2022 }
2023
2024 self.state.status = RunStatus::Voting;
2025 self.state.save()?;
2026 Ok(())
2027 }
2028
2029 async fn vote(&mut self) -> Result<()> {
2032 let run_id = self.state.id.clone();
2037 let prompts = self.state.config.prompts.clone();
2038 if !self.state.votes.is_empty() {
2039 return Ok(());
2040 }
2041 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2042 if viable.len() == 1 {
2043 return Ok(());
2044 }
2045 self.state.status = RunStatus::Voting;
2046
2047 let language = self.state.config.graph.language.clone();
2048 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2049 let sessions = self.state.config.graph.sessions;
2050 let artifacts = agent::artifacts_dir(&self.state.dir());
2051 let root = self.state.worktree_root();
2052 let base_short = short(&self.state.base_commit);
2053 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2054
2055 let mut jobs = Vec::new();
2056 let mut seats_at = Vec::new();
2057 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
2058 if self
2059 .state
2060 .judgements
2061 .get(j)
2062 .is_some_and(|r| r.failed.is_some())
2063 {
2064 continue;
2065 }
2066 let seat_key = format!("judge-{}", j + 1);
2067 let seat = self.seat(&seat_key, &spec.id);
2068 let mut text = prompt::final_vote(&viable, &language);
2069 if !has_context(&spec, &seat, sessions) {
2070 text = format!(
2071 "{}\n\n# Candidates\n\n{}",
2072 text,
2073 self.candidate_block(&candidates, &base_short)
2074 );
2075 }
2076 jobs.push(SeatJob {
2077 spec,
2078 seat,
2079 prompt: text,
2080 cwd: root.join(format!("judge-{}", j + 1)),
2081 timeout,
2082 allow_write: false,
2083 sessions,
2084 artifacts: artifacts.clone(),
2085 stem: format!("vote-judge-{}", j + 1),
2086 });
2087 seats_at.push(j);
2088 }
2089
2090 self.state.event(
2091 "vote",
2092 format!(
2093 "collecting {} final votes one by one, privately",
2094 jobs.len()
2095 ),
2096 );
2097 let allowed = viable.clone();
2098 let mut quota_losses = Vec::new();
2099 let cache = self.state.config.cache_dir();
2100 let ctx = WaveCtx {
2101 run: &run_id,
2102 node: "vote",
2103 prompts: &prompts,
2104 cache: cache.as_deref(),
2105 };
2106 let results = ask_json_wave::<FinalVote>(
2107 jobs,
2108 Arc::clone(&self.sem),
2109 self.state.config.graph.retries,
2110 &ctx,
2111 &mut quota_losses,
2112 &mut self.state,
2113 &move |v: &FinalVote| match v.label() {
2114 Some(c) if allowed.contains(&c) => Ok(()),
2115 other => bail!("vote {other:?} is not one of {allowed:?}"),
2116 },
2117 )
2118 .await;
2119 self.state.quota.extend(quota_losses);
2120
2121 for (&j, (seat, res)) in seats_at.iter().zip(results) {
2122 let agent_id = seat.agent.clone();
2123 self.state.seats.insert(seat.key.clone(), seat);
2124 let initial = self
2125 .state
2126 .judgements
2127 .get(j)
2128 .and_then(|r| r.ranking.first().copied());
2129 let mut record = VoteRecord {
2130 judge: j + 1,
2131 agent: agent_id,
2132 vote: None,
2133 reason: String::new(),
2134 changed: false,
2135 };
2136 match res {
2137 Ok((v, _)) => {
2138 record.vote = v.label();
2139 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2140 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
2141 self.state.event(
2142 "vote",
2143 format!(
2144 "judge {} voted {}{}",
2145 j + 1,
2146 record.vote.unwrap_or('?'),
2147 if record.changed { " (changed)" } else { "" }
2148 ),
2149 );
2150 }
2151 Err(e) => {
2152 self.state
2153 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
2154 }
2155 }
2156 self.state.votes.push(record);
2157 self.state.save()?;
2158 }
2159 Ok(())
2160 }
2161
2162 fn tally(&mut self) -> Result<()> {
2165 if self.state.tally.is_some() {
2166 return Ok(());
2167 }
2168 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
2169 let tops: Vec<char> = self
2170 .state
2171 .judgements
2172 .iter()
2173 .filter_map(|j| j.ranking.first().copied())
2174 .collect();
2175 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
2176
2177 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2180 let mut cast: Vec<char> = Vec::new();
2181 for (i, j) in self.state.judgements.iter().enumerate() {
2182 let vote = self
2183 .state
2184 .votes
2185 .iter()
2186 .find(|v| v.judge == i + 1)
2187 .and_then(|v| v.vote)
2188 .or_else(|| j.ranking.first().copied());
2189 if let Some(v) = vote {
2190 *first_choice.entry(v).or_insert(0) += 1;
2191 cast.push(v);
2192 }
2193 }
2194
2195 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
2196 for j in &self.state.judgements {
2197 let n = j.ranking.len();
2198 for (pos, label) in j.ranking.iter().enumerate() {
2199 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
2200 }
2201 }
2202
2203 let best = first_choice.values().copied().max().unwrap_or(0);
2204 let mut leaders: Vec<char> = first_choice
2205 .iter()
2206 .filter(|(_, v)| **v == best)
2207 .map(|(k, _)| *k)
2208 .collect();
2209 let mut tie_break = None;
2210 if leaders.len() > 1 {
2211 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
2212 let borda_leaders: Vec<char> = leaders
2213 .iter()
2214 .copied()
2215 .filter(|l| borda[l] == top_borda)
2216 .collect();
2217 tie_break = Some(if borda_leaders.len() == 1 {
2218 format!(
2219 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
2220 leaders.len()
2221 )
2222 } else {
2223 format!(
2224 "{} way tie on both first-choice votes and Borda points, broken by label order",
2225 leaders.len()
2226 )
2227 });
2228 leaders = borda_leaders;
2229 leaders.sort_unstable();
2230 }
2231 let winner = *leaders
2232 .first()
2233 .or(viable.first())
2234 .context("no candidate to declare a winner from")?;
2235
2236 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
2237 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
2238 let deliberated = !self.state.deliberation.is_empty();
2239
2240 let quota_seats: std::collections::BTreeSet<&str> =
2244 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2245 let mut present = 0usize;
2246 for (i, j) in self.state.judgements.iter().enumerate() {
2247 if quota_seats.contains(j.seat.as_str()) {
2248 continue;
2249 }
2250 let ranked = !j.ranking.is_empty() && j.failed.is_none();
2251 let voted = self
2252 .state
2253 .votes
2254 .iter()
2255 .any(|v| v.judge == i + 1 && v.vote.is_some());
2256 if ranked || voted {
2257 present += 1;
2258 }
2259 }
2260 let needs_quorum = viable.len() > 1;
2266 let judges_total = if needs_quorum {
2267 self.roles.judges.len()
2268 } else {
2269 0
2270 };
2271 let quorum = if needs_quorum {
2272 judges_total / 2 + 1
2273 } else {
2274 0
2275 };
2276 let met_quorum = !needs_quorum || present >= quorum;
2277 let uncontested = (!needs_quorum).then(|| {
2278 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
2279 });
2280
2281 self.state.event(
2282 "tally",
2283 match &uncontested {
2284 Some(reason) => format!("winner {winner} — {reason}"),
2285 None => format!(
2286 "winner {winner} — votes {} | initial {} | {} changed | \
2287 {present}/{judges_total} judges{}",
2288 first_choice
2289 .iter()
2290 .map(|(k, v)| format!("{k}:{v}"))
2291 .collect::<Vec<_>>()
2292 .join(" "),
2293 if unanimous_initial {
2294 "unanimous"
2295 } else {
2296 "split"
2297 },
2298 changed_votes,
2299 if met_quorum {
2300 String::new()
2301 } else {
2302 format!(" — below quorum ({quorum} required)")
2303 },
2304 ),
2305 },
2306 );
2307 if !met_quorum {
2308 self.state.event(
2309 "stall",
2310 format!(
2311 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
2312 the run stops here, resumable"
2313 ),
2314 );
2315 }
2316 self.state.tally = Some(Tally {
2317 first_choice,
2318 borda,
2319 winner,
2320 rankings: tops.len(),
2321 unanimous_initial,
2322 deliberated,
2323 changed_votes,
2324 unanimous_final,
2325 tie_break,
2326 judges: judges_total,
2327 present,
2328 quorum,
2329 met_quorum,
2330 uncontested,
2331 });
2332 self.state.status = if met_quorum {
2333 RunStatus::Reviewing
2334 } else {
2335 RunStatus::Stalled
2336 };
2337 self.state.save()?;
2338 Ok(())
2339 }
2340
2341 #[allow(clippy::too_many_lines)]
2362 async fn recover_stall(&mut self) -> Result<bool> {
2363 let run_id = self.state.id.clone();
2368 let prompts = self.state.config.prompts.clone();
2369 let quota_seats: BTreeSet<&str> =
2374 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
2375 let absent: Vec<String> = self
2376 .state
2377 .judgements
2378 .iter()
2379 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
2380 .map(|j| j.seat.clone())
2381 .collect();
2382 if absent.is_empty() {
2383 return Ok(false);
2384 }
2385 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
2386 if viable.len() <= 1 {
2387 return Ok(false);
2388 }
2389 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
2390 let language = self.state.config.graph.language.clone();
2391 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
2392 let sessions = self.state.config.graph.sessions;
2393 let artifacts = agent::artifacts_dir(&self.state.dir());
2394 let root = self.state.worktree_root();
2395 let base_short = short(&self.state.base_commit);
2396 let candidates: Vec<Candidate> = viable.clone();
2397
2398 let mut positions: Vec<usize> = absent
2400 .iter()
2401 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
2402 .collect();
2403 if positions.is_empty() {
2404 return Ok(false);
2405 }
2406 positions.sort_unstable();
2407 positions.dedup();
2408
2409 let mut judge_jobs = Vec::new();
2411 for &j in &positions {
2412 let order = blind::presentation_order(viable.len(), j, self.state.seed);
2413 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
2414 let seat_key = format!("judge-{}", j + 1);
2415 let spec = self.roles.judges[j].clone();
2416 let seat = self.seat(&seat_key, &spec.id);
2417 judge_jobs.push(SeatJob {
2418 spec,
2419 seat,
2420 prompt: prompt::judge(
2421 &self.state.instruction,
2422 &views,
2423 self.roles.judges.len(),
2424 &base_short,
2425 &language,
2426 ),
2427 cwd: root.join(seat_key),
2428 timeout,
2429 allow_write: false,
2430 sessions,
2431 artifacts: artifacts.clone(),
2432 stem: format!("judge-{}-recover", j + 1),
2433 });
2434 }
2435
2436 let labels_for_check = labels.clone();
2437 let mut judge_losses = Vec::new();
2438 let retries = self.state.config.graph.retries;
2439 let cache = self.state.config.cache_dir();
2440 let ctx = WaveCtx {
2441 run: &run_id,
2442 node: "judge",
2443 prompts: &prompts,
2444 cache: cache.as_deref(),
2445 };
2446 let results = ask_json_wave::<Ranking>(
2447 judge_jobs,
2448 Arc::clone(&self.sem),
2449 retries,
2450 &ctx,
2451 &mut judge_losses,
2452 &mut self.state,
2453 &move |r: &Ranking| r.validate(&labels_for_check),
2454 )
2455 .await;
2456
2457 let mut recovered: BTreeSet<usize> = BTreeSet::new();
2459 for (&j, (seat, res)) in positions.iter().zip(results) {
2460 self.state.seats.insert(seat.key.clone(), seat);
2461 let record = &mut self.state.judgements[j];
2462 match res {
2463 Ok((ranking, out)) => {
2464 record.ranking = ranking.normalized();
2465 record.reasons = ranking.reasons;
2466 record.confidence = ranking.confidence;
2467 record.failed = None;
2468 record.duration_ms = out.duration_ms;
2469 recovered.insert(j);
2470 self.state.event(
2471 "recover",
2472 format!("judge {} ranked again after the limit", j + 1),
2473 );
2474 }
2475 Err(e) => {
2476 self.state
2477 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
2478 }
2479 }
2480 }
2481
2482 let mut vote_jobs = Vec::new();
2484 let mut vote_pos: Vec<usize> = Vec::new();
2485 for &j in &recovered {
2486 let seat_key = format!("judge-{}", j + 1);
2487 let spec = self.roles.judges[j].clone();
2488 let seat = self.seat(&seat_key, &spec.id);
2489 let mut text = prompt::final_vote(&labels, &language);
2490 if !has_context(&spec, &seat, sessions) {
2491 text = format!(
2492 "{}\n\n# Candidates\n\n{}",
2493 text,
2494 self.candidate_block(&candidates, &base_short)
2495 );
2496 }
2497 vote_jobs.push(SeatJob {
2498 spec,
2499 seat,
2500 prompt: text,
2501 cwd: root.join(seat_key),
2502 timeout,
2503 allow_write: false,
2504 sessions,
2505 artifacts: artifacts.clone(),
2506 stem: format!("vote-judge-{}-recover", j + 1),
2507 });
2508 vote_pos.push(j);
2509 }
2510 let allowed = labels.clone();
2511 let mut vote_losses = Vec::new();
2512 let vote_retries = self.state.config.graph.retries;
2513 let vote_cache = self.state.config.cache_dir();
2514 let ctx = WaveCtx {
2515 run: &run_id,
2516 node: "vote",
2517 prompts: &prompts,
2518 cache: vote_cache.as_deref(),
2519 };
2520 let votes = ask_json_wave::<FinalVote>(
2521 vote_jobs,
2522 Arc::clone(&self.sem),
2523 vote_retries,
2524 &ctx,
2525 &mut vote_losses,
2526 &mut self.state,
2527 &move |v: &FinalVote| match v.label() {
2528 Some(c) if allowed.contains(&c) => Ok(()),
2529 other => bail!("vote {other:?} is not one of {allowed:?}"),
2530 },
2531 )
2532 .await;
2533 for (&j, (seat, res)) in vote_pos.iter().zip(votes) {
2534 let agent_id = seat.agent.clone();
2535 self.state.seats.insert(seat.key.clone(), seat);
2536 match res {
2537 Ok((v, _)) => {
2538 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
2539 rec.vote = v.label();
2540 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
2541 } else {
2542 self.state.votes.push(VoteRecord {
2543 judge: j + 1,
2544 agent: agent_id,
2545 vote: v.label(),
2546 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
2547 changed: false,
2548 });
2549 }
2550 self.state.event(
2551 "recover",
2552 format!("judge {} voted again after the limit", j + 1),
2553 );
2554 }
2555 Err(e) => {
2556 self.state
2557 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
2558 }
2559 }
2560 }
2561
2562 if !recovered.is_empty() {
2566 let recovered_keys: BTreeSet<String> = recovered
2567 .iter()
2568 .map(|&j| format!("judge-{}", j + 1))
2569 .collect();
2570 self.state
2571 .quota
2572 .retain(|q| !recovered_keys.contains(&q.seat));
2573 }
2574
2575 self.state.tally = None;
2577 self.tally()?;
2578 Ok(self
2579 .state
2580 .tally
2581 .as_ref()
2582 .map(|t| t.met_quorum)
2583 .unwrap_or(false))
2584 }
2585
2586 async fn fold_losers(&mut self) -> Result<()> {
2589 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
2590 return Ok(());
2591 };
2592 let repo = self.state.repo.clone();
2593 let mut folded = Vec::new();
2594 for i in 0..self.state.candidates.len() {
2595 let c = &self.state.candidates[i];
2596 if c.label == winner || c.folded {
2597 continue;
2598 }
2599 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
2600 git::worktree_remove(&repo, &wt).await.ok();
2601 git::branch_delete(&repo, &branch).await.ok();
2602 self.state.candidates[i].folded = true;
2603 folded.push(label.to_string());
2604 }
2605 let root = self.state.worktree_root();
2607 for j in 1..=self.roles.judges.len() {
2608 let wt = root.join(format!("judge-{j}"));
2609 if wt.exists() {
2610 git::worktree_remove(&repo, &wt).await.ok();
2611 }
2612 }
2613 if self.state.config.graph.advise {
2616 for k in 1..=self.state.config.graph.advisors {
2617 let wt = root.join(format!("advisor-{k}"));
2618 if wt.exists() {
2619 git::worktree_remove(&repo, &wt).await.ok();
2620 }
2621 }
2622 }
2623 if !folded.is_empty() {
2624 self.state
2625 .event("fold", format!("folded candidates {}", folded.join(", ")));
2626 self.state.save()?;
2627 }
2628 Ok(())
2629 }
2630
2631 async fn sync_to_base(&mut self) -> Result<()> {
2661 if self
2662 .state
2663 .base_sync
2664 .as_ref()
2665 .is_some_and(|s| s.conflict.is_some())
2666 {
2667 return Ok(());
2668 }
2669 let Some(winner) = self.state.winner().cloned() else {
2670 return Ok(());
2671 };
2672
2673 let repo = self.state.repo.clone();
2674 let remote = self.state.config.merge.remote.clone();
2675 let base_branch = self.state.base_branch.clone();
2676 let tracking = format!("{remote}/{base_branch}");
2677
2678 git::fetch(&repo, &remote, &base_branch).await.ok();
2679 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
2683 return Ok(());
2684 };
2685
2686 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
2687 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
2688 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
2689
2690 if behind == 0 {
2691 self.state.base_sync = Some(BaseSync {
2692 tip,
2693 behind: 0,
2694 attempts,
2695 conflict: None,
2696 });
2697 self.state.save()?;
2698 return Ok(());
2699 }
2700
2701 if attempts >= BASE_SYNC_ROUNDS {
2702 let why = format!(
2703 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
2704 rebase(s); rebasing again would only race it",
2705 winner.branch
2706 );
2707 self.state.status = RunStatus::Blocked;
2708 self.state.base_sync = Some(BaseSync {
2709 tip,
2710 behind,
2711 attempts,
2712 conflict: Some(why.clone()),
2713 });
2714 self.state.event("land", why);
2715 self.state.save()?;
2716 return Ok(());
2717 }
2718
2719 self.state.event(
2720 "land",
2721 format!(
2722 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
2723 winner.branch
2724 ),
2725 );
2726 self.state.save()?;
2727
2728 let scratch = self.state.dir().join("base-sync");
2729 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
2730 let attempts = attempts + 1;
2731 match rebased {
2732 Ok(None) => {
2733 git::sync_to_head(&winner.worktree).await?;
2737 self.state.base_sync = Some(BaseSync {
2738 tip: tip.clone(),
2739 behind: 0,
2740 attempts,
2741 conflict: None,
2742 });
2743 self.state
2744 .event("land", format!("rebased {} onto {tracking}", winner.branch));
2745 }
2746 Ok(Some(conflict)) => {
2747 let why = format!(
2748 "{} conflicts with {tracking} and did not rebase: {}",
2749 winner.branch,
2750 conflict.chars().take(600).collect::<String>()
2751 );
2752 self.state.status = RunStatus::Blocked;
2753 self.state.base_sync = Some(BaseSync {
2754 tip,
2755 behind,
2756 attempts,
2757 conflict: Some(why.clone()),
2758 });
2759 self.state.event("land", why);
2760 }
2761 Err(e) => {
2762 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
2763 self.state.status = RunStatus::Blocked;
2764 self.state.base_sync = Some(BaseSync {
2765 tip,
2766 behind,
2767 attempts,
2768 conflict: Some(why.clone()),
2769 });
2770 self.state.event("land", why);
2771 }
2772 }
2773 self.state.save()?;
2774 Ok(())
2775 }
2776
2777 fn landing_base(&self) -> String {
2787 self.state
2788 .base_sync
2789 .as_ref()
2790 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
2791 }
2792
2793 async fn review_loop(&mut self) -> Result<()> {
2796 if self
2801 .state
2802 .base_sync
2803 .as_ref()
2804 .is_some_and(|s| s.conflict.is_some())
2805 {
2806 return Ok(());
2807 }
2808 let run_id = self.state.id.clone();
2813 let prompts = self.state.config.prompts.clone();
2814 let Some(winner) = self.state.winner().cloned() else {
2815 return Ok(());
2816 };
2817 let max_rounds = self.state.config.graph.review_rounds;
2818 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
2828 self.state.status = status;
2829 self.state.save()?;
2830 return Ok(());
2831 }
2832 self.state.status = RunStatus::Reviewing;
2833
2834 let repo = self.state.repo.clone();
2835 let root = self.state.worktree_root();
2836 let language = self.state.config.graph.language.clone();
2837 let sessions = self.state.config.graph.sessions;
2838 let artifacts = agent::artifacts_dir(&self.state.dir());
2839 let base = self.landing_base();
2840 let base_short = short(&base);
2841 let reviewers = self.roles.reviewers.clone();
2842 let shell = self.state.config.shell();
2843
2844 let mut prev_e2e: Option<String> = None;
2845 for round in (self.state.reviews.len() + 1)..=max_rounds {
2846 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
2847 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
2848 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
2849
2850 let mut jobs = Vec::new();
2854 for (r, spec) in reviewers.iter().cloned().enumerate() {
2855 let wt = root.join(format!("review-{}", r + 1));
2856 if wt.exists() {
2857 git::reset_detached(&wt, &head).await?;
2858 } else {
2859 git::worktree_add_detached(&repo, &wt, &head).await?;
2860 }
2861 let seat_key = format!("review-{}", r + 1);
2862 let seat = self.seat(&seat_key, &spec.id);
2863 jobs.push(SeatJob {
2864 prompt: prompt::review(&prompt::ReviewCtx {
2865 instruction: &self.state.instruction,
2866 branch: &winner.branch,
2867 base_short: &base_short,
2868 stat: &stat,
2869 patch: &patch,
2870 e2e: prev_e2e.as_deref(),
2871 reviewers: reviewers.len(),
2872 round,
2873 rounds: max_rounds,
2874 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
2877 lens: Lens::for_seat(r),
2878 language: &language,
2879 }),
2880 spec,
2881 seat,
2882 cwd: wt,
2883 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
2884 allow_write: false,
2885 sessions,
2886 artifacts: artifacts.clone(),
2887 stem: format!("review-{round}-{}", r + 1),
2888 });
2889 }
2890
2891 self.state.event(
2892 "review",
2893 format!(
2894 "round {round}: {} reviewers on {}",
2895 jobs.len(),
2896 short(&head)
2897 ),
2898 );
2899 let mut quota_losses = Vec::new();
2900 let review_retries = self.state.config.graph.retries;
2901 let review_cache = self.state.config.cache_dir();
2902 let ctx = WaveCtx {
2903 run: &run_id,
2904 node: "review",
2905 prompts: &prompts,
2906 cache: review_cache.as_deref(),
2907 };
2908 let results = ask_json_wave::<Review>(
2909 jobs,
2910 Arc::clone(&self.sem),
2911 review_retries,
2912 &ctx,
2913 &mut quota_losses,
2914 &mut self.state,
2915 &|_: &Review| Ok(()),
2916 )
2917 .await;
2918 let round_quota_missing = quota_losses.len();
2922 self.state.quota.extend(quota_losses);
2923
2924 let mut records = Vec::new();
2925 let mut all_findings = Vec::new();
2926 for (r, (seat, res)) in results.into_iter().enumerate() {
2927 let agent_id = seat.agent.clone();
2928 self.state.seats.insert(seat.key.clone(), seat);
2929 let mut record = ReviewRecord {
2930 reviewer: r + 1,
2931 agent: agent_id,
2932 summary: String::new(),
2933 findings: Vec::new(),
2934 vote: None,
2935 failed: None,
2936 duration_ms: 0,
2937 };
2938 match res {
2939 Ok((review, out)) => {
2940 record.summary =
2948 blind::sanitize_prose(&review.summary, &self.state.config.blind);
2949 record.vote = Some(review.vote);
2950 record.duration_ms = out.duration_ms;
2951 for (n, mut f) in review.findings.into_iter().enumerate() {
2952 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
2955 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
2956 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
2957 f.file = f
2963 .file
2964 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
2965 all_findings.push(f.clone());
2966 record.findings.push(f);
2967 }
2968 self.state.event(
2969 "review",
2970 format!(
2971 "round {round}: reviewer {} voted {} with {} finding(s)",
2972 r + 1,
2973 review.vote.label(),
2974 record.findings.len()
2975 ),
2976 );
2977 }
2978 Err(e) => {
2979 record.failed = Some(e.to_string());
2980 self.state.event(
2981 "review",
2982 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
2983 );
2984 }
2985 }
2986 records.push(record);
2987 }
2988
2989 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
2996 let vote_split =
2997 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
2998 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
2999 if vote_split {
3000 self.state.event(
3001 "review",
3002 format!(
3003 "round {round}: votes split ({}) — one round of reconsideration",
3004 initial_votes
3005 .iter()
3006 .map(|v| v.label())
3007 .collect::<Vec<_>>()
3008 .join(", ")
3009 ),
3010 );
3011 let panel: Vec<ReviewSeatReport<'_>> = records
3014 .iter()
3015 .filter_map(|r| {
3016 r.vote.map(|vote| ReviewSeatReport {
3017 reviewer: r.reviewer,
3018 vote,
3019 summary: &r.summary,
3020 findings: &r.findings,
3021 })
3022 })
3023 .collect();
3024
3025 let mut jobs = Vec::new();
3026 let mut seats_at = Vec::new();
3027 for (r, spec) in reviewers.iter().cloned().enumerate() {
3028 if records[r].vote.is_none() {
3032 continue;
3033 }
3034 let wt = root.join(format!("review-{}", r + 1));
3035 let seat_key = format!("review-{}", r + 1);
3036 let seat = self.seat(&seat_key, &spec.id);
3037 let patch_ctx = if has_context(&spec, &seat, sessions) {
3042 None
3043 } else {
3044 Some(ReviewPatch {
3045 branch: &winner.branch,
3046 base_short: &base_short,
3047 stat: &stat,
3048 patch: &patch,
3049 })
3050 };
3051 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
3052 instruction: &self.state.instruction,
3053 reviewer: r + 1,
3054 lens: Lens::for_seat(r),
3055 panel: &panel,
3056 patch: patch_ctx,
3057 round,
3058 rounds: max_rounds,
3059 language: &language,
3060 });
3061 jobs.push(SeatJob {
3062 prompt,
3063 spec,
3064 seat,
3065 cwd: wt,
3066 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
3067 allow_write: false,
3068 sessions,
3069 artifacts: artifacts.clone(),
3070 stem: format!("review-{round}-reconsider-{}", r + 1),
3071 });
3072 seats_at.push(r);
3073 }
3074
3075 let mut recon_quota_losses = Vec::new();
3076 let recon_cache = self.state.config.cache_dir();
3077 let recon_ctx = WaveCtx {
3078 run: &run_id,
3079 node: "review",
3080 prompts: &prompts,
3081 cache: recon_cache.as_deref(),
3082 };
3083 let recon_results = ask_json_wave::<ReviewRevote>(
3084 jobs,
3085 Arc::clone(&self.sem),
3086 review_retries,
3087 &recon_ctx,
3088 &mut recon_quota_losses,
3089 &mut self.state,
3090 &|_: &ReviewRevote| Ok(()),
3091 )
3092 .await;
3093 self.state.quota.extend(recon_quota_losses);
3094
3095 for (&r, (seat, res)) in seats_at.iter().zip(recon_results) {
3096 let agent_id = seat.agent.clone();
3097 self.state.seats.insert(seat.key.clone(), seat);
3098 let mut rec = ReviewRevoteRecord {
3099 reviewer: r + 1,
3100 agent: agent_id,
3101 vote: None,
3102 reason: String::new(),
3103 failed: None,
3104 };
3105 match res {
3106 Ok((rv, _)) => {
3107 rec.vote = Some(rv.vote);
3108 rec.reason =
3109 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
3110 self.state.event(
3111 "review",
3112 format!(
3113 "round {round}: reviewer {} revoted {}",
3114 r + 1,
3115 rv.vote.label()
3116 ),
3117 );
3118 }
3119 Err(e) => {
3120 rec.failed = Some(e.to_string());
3121 self.state.event(
3122 "review",
3123 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
3124 );
3125 }
3126 }
3127 reconsideration.push(rec);
3128 }
3129 } else if initial_votes.len() > 1 {
3130 self.state.event(
3131 "review",
3132 format!(
3133 "round {round}: votes agreed ({}) — no reconsideration",
3134 initial_votes[0].label()
3135 ),
3136 );
3137 }
3138
3139 let final_votes: Vec<ReviewVote> = records
3143 .iter()
3144 .filter_map(|r| {
3145 reconsideration
3146 .iter()
3147 .find(|rv| rv.reviewer == r.reviewer)
3148 .and_then(|rv| rv.vote)
3149 .or(r.vote)
3150 })
3151 .collect();
3152 let round_verdict = ReviewVote::worst(final_votes);
3153
3154 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
3155 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
3156 let defer_e2e =
3167 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
3168 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
3169 let reason =
3170 format!("{blocking} blocking finding(s) already required a fix this round");
3171 self.state.event(
3172 "verify",
3173 format!(
3174 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
3175 {}); it will run once a round has none left",
3176 short(&head)
3177 ),
3178 );
3179 (Vec::new(), false, true, Some(reason))
3180 } else {
3181 let e2e_commands = self.state.config.verify.e2e.clone();
3182 let cache_dir = self.state.config.cache_dir();
3183 let context = format!("round {round}");
3184 let (e2e, verify_retried) = with_cache_lease(
3185 &mut self.state,
3186 cache_dir.as_deref(),
3187 "e2e",
3188 "e2e",
3189 &winner.worktree,
3190 &head,
3191 verify_timeout,
3192 &context,
3193 |state, budget| {
3194 let shell = shell.clone();
3195 let e2e_commands = e2e_commands.clone();
3196 let worktree = winner.worktree.clone();
3197 let context = context.clone();
3198 async move {
3199 run_e2e_with_retry(
3200 state,
3201 &shell,
3202 &e2e_commands,
3203 &worktree,
3204 budget,
3205 &context,
3206 )
3207 .await
3208 }
3209 },
3210 )
3211 .await;
3212 (e2e, verify_retried, false, None)
3213 };
3214
3215 let e2e_failures: String = e2e
3216 .iter()
3217 .filter(|o| !o.ok())
3218 .map(|o| format!("$ {}\n{}\n", o.command, o.output_tail))
3219 .collect();
3220
3221 let expected = records.len();
3222 let answered = records.iter().filter(|r| r.failed.is_none()).count();
3223 let incomplete = answered < expected;
3224 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
3225 let policy = self.state.config.graph.incomplete_review;
3226 let clean = round_is_clean(
3227 blocking,
3228 e2e_ok,
3229 answered,
3230 expected,
3231 round_quota_missing,
3232 policy,
3233 );
3234
3235 let mut round_record = ReviewRound {
3236 round,
3237 head: head.clone(),
3238 verified_head: None,
3239 reviews: records,
3240 e2e,
3241 verify_retried,
3242 e2e_deferred,
3243 e2e_defer_reason,
3244 fix: None,
3245 blocking,
3246 answered,
3247 expected,
3248 clean,
3249 progressed: false,
3250 vote_split,
3251 reconsideration,
3252 verdict: round_verdict,
3253 };
3254
3255 if incomplete {
3256 let missing: Vec<String> = round_record
3257 .reviews
3258 .iter()
3259 .filter(|r| r.failed.is_some())
3260 .map(|r| format!("review-{}", r.reviewer))
3261 .collect();
3262 self.state.event(
3263 "review",
3264 format!(
3265 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
3266 missing.join(", ")
3267 ),
3268 );
3269 }
3270
3271 if clean {
3272 self.state.event(
3273 "review",
3274 if incomplete && policy == IncompleteReviewPolicy::Warn {
3275 format!(
3276 "round {round}: clean (warn policy, incomplete panel) — no \
3277 blocking findings from the seats that answered, verification green"
3278 )
3279 } else if incomplete {
3280 format!(
3281 "round {round}: clean ({} rate-limited reviewer(s) excluded from \
3282 quorum) — no blocking findings from the seats that answered, \
3283 verification green",
3284 expected - answered
3285 )
3286 } else {
3287 format!("round {round}: clean — no blocking findings, verification green")
3288 },
3289 );
3290 self.state.reviews.push(round_record);
3291 self.state.status = RunStatus::Gating;
3292 self.state.save()?;
3293 return Ok(());
3294 }
3295
3296 if incomplete && blocking == 0 && e2e_ok {
3304 self.state.reviews.push(round_record);
3305 self.state.save()?;
3306 if round == max_rounds {
3307 self.state.status = RunStatus::Blocked;
3308 self.state.event(
3309 "review",
3310 format!(
3311 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
3312 refusing to call it clean",
3313 expected - answered
3314 ),
3315 );
3316 return Ok(());
3317 }
3318 prev_e2e = None;
3319 continue;
3320 }
3321
3322 if round == max_rounds {
3323 self.state.reviews.push(round_record);
3324 return self
3325 .stop_reviewing(
3326 &format!(
3327 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
3328 ),
3329 &shell,
3330 &winner.worktree,
3331 )
3332 .await;
3333 }
3334
3335 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
3338 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
3339 _ => (
3340 self.state
3341 .config
3342 .agent(&winner.agent)
3343 .cloned()
3344 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
3345 format!("impl-{}", winner.label),
3346 ),
3347 };
3348 let seat = self.seat(&fix_seat_key, &fix_spec.id);
3349 let blocking_findings: Vec<_> = all_findings
3350 .iter()
3351 .filter(|f| f.severity.blocks())
3352 .cloned()
3353 .collect();
3354 let job = SeatJob {
3355 prompt: prompt::fix(
3356 &self.state.instruction,
3357 &blocking_findings,
3358 (!e2e_failures.is_empty()).then_some(e2e_failures.as_str()),
3359 e2e_deferred,
3360 round,
3361 max_rounds,
3362 &language,
3363 ),
3364 spec: fix_spec.clone(),
3365 seat,
3366 cwd: winner.worktree.clone(),
3367 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
3368 allow_write: true,
3369 sessions,
3370 artifacts: artifacts.clone(),
3371 stem: format!("fix-{round}"),
3372 };
3373 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
3374 let cache = self.state.config.cache_dir();
3375 let ctx = WaveCtx {
3376 run: &run_id,
3377 node: "fix",
3378 prompts: &prompts,
3379 cache: cache.as_deref(),
3380 };
3381 let (seat, out) =
3382 run_one(job.clone(), Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
3383 let agent_id = seat.agent.clone();
3384
3385 let mut fix = FixRecord {
3386 agent: agent_id,
3387 addressed: Vec::new(),
3388 rejected: Vec::new(),
3389 notes: String::new(),
3390 committed: false,
3391 failed: None,
3392 duration_ms: 0,
3393 continuation: None,
3394 };
3395 let mut continuation = ContinuationRecord::not_needed();
3396 let mut final_seat = seat.clone();
3397 match out {
3398 AgentOutcome::Ok(o) => {
3399 fix.duration_ms = o.duration_ms;
3400 let parsed = verdict::extract_json::<FixReport>(&o.text);
3401 let incomplete_reason = match &parsed {
3408 Ok(_) if has_unconfirmed_command(&o.commands) => Some(
3409 "the reply parsed, but it reported a command whose own CLI never \
3410 confirmed an exit status"
3411 .to_owned(),
3412 ),
3413 Ok(_) => None,
3414 Err(e) => Some(e.to_string()),
3415 };
3416 match incomplete_reason {
3417 None => {
3418 let report = parsed.expect("checked Ok above");
3419 fix.addressed = report.addressed;
3420 fix.rejected = report.rejected;
3421 fix.notes =
3422 blind::sanitize_prose(&report.notes, &self.state.config.blind);
3423 }
3424 Some(reason) => {
3425 let (resumed_seat, resolved, failure, cont) = self
3426 .continue_fix_report(seat, reason, &job, &prompts, &run_id, round)
3427 .await;
3428 fix.duration_ms += cont.cumulative_wait_ms;
3429 continuation = cont;
3430 final_seat = resumed_seat;
3431 match resolved {
3432 Some(report) => {
3433 fix.addressed = report.addressed;
3434 fix.rejected = report.rejected;
3435 fix.notes = blind::sanitize_prose(
3436 &report.notes,
3437 &self.state.config.blind,
3438 );
3439 }
3440 None => fix.failed = failure,
3441 }
3442 }
3443 }
3444 }
3445 AgentOutcome::Dropped(o) => {
3447 fix.duration_ms = o.duration_ms;
3448 let why = o
3449 .dropped
3450 .as_ref()
3451 .map(|d| d.why.as_str())
3452 .unwrap_or("the CLI ended the stream without delivering its answer");
3453 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
3454 }
3455 AgentOutcome::Quota(o) => {
3456 self.state.quota.push(QuotaLoss {
3457 seat: final_seat.key.clone(),
3458 node: "fix".to_owned(),
3459 at: Timestamp::now(),
3460 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3461 });
3462 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
3463 }
3464 AgentOutcome::Failed(e) => fix.failed = Some(e),
3465 }
3466 fix.continuation = Some(continuation);
3467 self.state.seats.insert(final_seat.key.clone(), final_seat);
3468 git::commit_all(
3469 &winner.worktree,
3470 &format!("magi: review round {round} fixes (uncommitted work)"),
3471 )
3472 .await
3473 .ok();
3474 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
3475 fix.committed = after != before;
3476 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
3484 let progressed = diff_after != patch;
3485 let commit_note = if fix.committed {
3486 "committed"
3487 } else {
3488 "NO new commit"
3489 };
3490 let tree_note = if progressed {
3491 "changed vs base"
3492 } else {
3493 "unchanged vs base"
3494 };
3495 self.state.event(
3496 "fix",
3497 match &fix.failed {
3498 Some(reason) => {
3504 format!(
3505 "round {round}: fixer's adoption report was lost ({reason}); \
3506 {commit_note}, tree {tree_note}"
3507 )
3508 }
3509 None => format!(
3510 "round {round}: {} addressed, {} rejected, {commit_note}, tree \
3511 {tree_note}{}",
3512 fix.addressed.len(),
3513 fix.rejected.len(),
3514 if continuation.outcome == ContinuationOutcome::Resumed {
3515 format!(
3516 " (adoption report recovered after {} continuation(s))",
3517 continuation.attempts
3518 )
3519 } else {
3520 String::new()
3521 },
3522 ),
3523 },
3524 );
3525 round_record.fix = Some(fix);
3526 round_record.progressed = progressed;
3527 self.state.reviews.push(round_record);
3528 self.state.save()?;
3529
3530 if matches!(
3543 continuation.outcome,
3544 ContinuationOutcome::Exhausted
3545 | ContinuationOutcome::QuotaLost
3546 | ContinuationOutcome::NoSession
3547 ) {
3548 return self
3549 .stop_reviewing(
3550 "the fixer's adoption report never came back, even after resuming its \
3551 own seat; refusing to start another round against the same worktree \
3552 while that is unresolved",
3553 &shell,
3554 &winner.worktree,
3555 )
3556 .await;
3557 }
3558
3559 prev_e2e = (!e2e_failures.is_empty()).then_some(e2e_failures);
3560
3561 let streak = self
3562 .state
3563 .reviews
3564 .iter()
3565 .rev()
3566 .take_while(|r| !r.progressed)
3567 .count();
3568 if streak >= STAGNANT_LIMIT {
3569 return self
3570 .stop_reviewing(
3571 &format!(
3572 "the tree has not moved against base for {streak} round(s) in a row"
3573 ),
3574 &shell,
3575 &winner.worktree,
3576 )
3577 .await;
3578 }
3579 }
3580 Ok(())
3581 }
3582
3583 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
3607 let round_idx = self.state.reviews.len() - 1;
3608 let needs_catchup_run = {
3609 let last = &self.state.reviews[round_idx];
3610 last.e2e.is_empty() && last.e2e_deferred
3611 };
3612 if needs_catchup_run {
3613 let round = self.state.reviews[round_idx].round;
3614 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
3615 let commands = self.state.config.verify.e2e.clone();
3616 let verified_head = git::rev_parse(worktree, "HEAD").await?;
3617 let cache_dir = self.state.config.cache_dir();
3618 let context =
3619 format!("round {round}: deferred e2e, now catching up before the final decision");
3620 let (outcomes, verify_retried) = with_cache_lease(
3621 &mut self.state,
3622 cache_dir.as_deref(),
3623 "e2e",
3624 "e2e",
3625 worktree,
3626 &verified_head,
3627 timeout,
3628 &context,
3629 |state, budget| {
3630 let shell = shell.to_vec();
3631 let commands = commands.clone();
3632 let context = context.clone();
3633 async move {
3634 run_e2e_with_retry(state, &shell, &commands, worktree, budget, &context)
3635 .await
3636 }
3637 },
3638 )
3639 .await;
3640 if verify_inconclusive(&outcomes) {
3648 self.state.save()?;
3649 return Ok(());
3650 }
3651 let last = &mut self.state.reviews[round_idx];
3652 last.e2e = outcomes;
3653 last.verify_retried = verify_retried;
3654 last.e2e_deferred = false;
3655 if verified_head != last.head {
3656 last.verified_head = Some(verified_head);
3657 }
3658 }
3659 let last = &self.state.reviews[round_idx];
3660 let red: Vec<String> = last
3661 .e2e
3662 .iter()
3663 .filter(|o| !o.ok())
3664 .map(|o| {
3665 format!(
3666 "`{}` -> {:?}\n{}",
3667 o.command,
3668 o.code,
3669 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
3670 )
3671 })
3672 .collect();
3673 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
3674
3675 if red.is_empty() {
3676 self.state.event(
3677 "review",
3678 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
3679 );
3680 self.state.status = RunStatus::Gating;
3681 } else {
3682 self.state
3683 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
3684 self.state.status = RunStatus::Blocked;
3685 }
3686 self.state.save()?;
3687 Ok(())
3688 }
3689
3690 async fn gate(&mut self) -> Result<()> {
3693 if self.state.status == RunStatus::Failed
3705 || self
3706 .state
3707 .base_sync
3708 .as_ref()
3709 .is_some_and(|s| s.conflict.is_some())
3710 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
3711 != Some(RunStatus::Gating)
3712 {
3713 return Ok(());
3714 }
3715 if self.state.gate_ran {
3716 if self.state.gate.iter().any(|outcome| !outcome.ok()) {
3727 self.state.status = RunStatus::Blocked;
3728 self.state.save()?;
3729 }
3730 return Ok(());
3731 }
3732 let Some(winner) = self.state.winner().cloned() else {
3733 return Ok(());
3734 };
3735 self.state.status = RunStatus::Gating;
3736 let shell = self.state.config.shell();
3737 let gate_commands = self.state.config.verify.gate.clone();
3738 let outcomes = if gate_commands.is_empty() {
3747 Vec::new()
3748 } else {
3749 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
3750 let cache_dir = self.state.config.cache_dir();
3751 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
3752 let (outcomes, _) = with_cache_lease(
3753 &mut self.state,
3754 cache_dir.as_deref(),
3755 "gate",
3756 "gate",
3757 &winner.worktree,
3758 &head,
3759 timeout,
3760 "final gate",
3761 |_state, budget| {
3762 let shell = shell.clone();
3763 let gate_commands = gate_commands.clone();
3764 let worktree = winner.worktree.clone();
3765 async move {
3766 let (outcomes, timed_out_pids) =
3767 run_commands(&shell, &gate_commands, &worktree, budget).await;
3768 (outcomes, false, timed_out_pids)
3769 }
3770 },
3771 )
3772 .await;
3773 outcomes
3774 };
3775 if outcomes.is_empty() {
3776 self.state.event(
3781 "gate",
3782 "no gate commands configured; nothing to check, passing",
3783 );
3784 }
3785 for o in &outcomes {
3786 self.state.event(
3787 "gate",
3788 format!(
3789 "`{}` -> {}",
3790 o.command,
3791 if o.ok() {
3792 "pass".to_owned()
3793 } else {
3794 format!(
3795 "FAIL ({:?})\n{}",
3796 o.code,
3797 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
3798 )
3799 }
3800 ),
3801 );
3802 }
3803 if verify_inconclusive(&outcomes) {
3814 self.state.save()?;
3815 return Ok(());
3816 }
3817 let passed = outcomes.iter().all(CommandOutcome::ok);
3818 self.state.gate = outcomes;
3819 self.state.gate_ran = true;
3820 if !passed {
3821 self.state.status = RunStatus::Blocked;
3822 self.state.event("gate", "gate failed; not merging");
3823 }
3824 self.state.save()?;
3825 Ok(())
3826 }
3827
3828 async fn merge(&mut self) -> Result<()> {
3831 if self
3846 .state
3847 .base_sync
3848 .as_ref()
3849 .is_some_and(|s| s.conflict.is_some())
3850 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
3851 != Some(RunStatus::Gating)
3852 || !self.state.gate_status().ok()
3861 {
3862 return Ok(());
3863 }
3864 if self.state.merge.is_some() {
3873 return Ok(());
3874 }
3875 let Some(winner) = self.state.winner().cloned() else {
3876 return Ok(());
3877 };
3878 let repo = self.state.repo.clone();
3879 let base = self.state.base_branch.clone();
3880 let mode = self.state.config.merge.mode;
3881 let style = self.state.config.merge.style;
3882 let message = pr_body(&self.state, winner.label);
3883
3884 let outcome = match mode {
3885 MergeMode::None => MergeOutcome {
3886 mode,
3887 ok: true,
3888 detail: manual_merge_command(style, &repo, &winner.branch, &message),
3889 },
3890 MergeMode::Local => {
3891 let on = git::current_branch(&repo).await?;
3892 if on.as_deref() != Some(base.as_str()) {
3893 MergeOutcome {
3894 mode,
3895 ok: false,
3896 detail: format!(
3897 "{} has {} checked out, not the base branch {base}",
3898 repo.display(),
3899 on.unwrap_or_else(|| "a detached HEAD".to_owned())
3900 ),
3901 }
3902 } else if !git::is_clean(&repo).await? {
3903 MergeOutcome {
3904 mode,
3905 ok: false,
3906 detail: format!("{} is dirty; refusing to merge", repo.display()),
3907 }
3908 } else {
3909 let out = match style {
3910 MergeStyle::Merge => {
3911 git::merge_no_ff(&repo, &winner.branch, &message).await?
3912 }
3913 MergeStyle::Squash => {
3914 git::merge_squash(&repo, &winner.branch, &message).await?
3915 }
3916 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
3917 };
3918 MergeOutcome {
3919 mode,
3920 ok: out.ok(),
3921 detail: if out.ok() { out.stdout } else { out.stderr },
3922 }
3923 }
3924 }
3925 MergeMode::Pr => {
3926 let remote = self.state.config.merge.remote.clone();
3927 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
3928 if !pushed.ok() {
3929 MergeOutcome {
3930 mode,
3931 ok: false,
3932 detail: pushed.stderr,
3933 }
3934 } else {
3935 let out = gh_pr_create(&winner.worktree, &base, &winner.branch, &message).await;
3936 match out {
3937 Ok(url) => MergeOutcome {
3938 mode,
3939 ok: true,
3940 detail: url,
3941 },
3942 Err(e) => MergeOutcome {
3943 mode,
3944 ok: false,
3945 detail: e.to_string(),
3946 },
3947 }
3948 }
3949 }
3950 };
3951
3952 self.state.status = match (mode, outcome.ok) {
3953 (MergeMode::None, _) => RunStatus::Ready,
3954 (_, true) => RunStatus::Merged,
3955 (_, false) => RunStatus::Blocked,
3956 };
3957 self.state.event(
3958 "merge",
3959 format!(
3960 "{:?}: {}",
3961 mode,
3962 outcome.detail.lines().next().unwrap_or("")
3963 ),
3964 );
3965 self.state.merge = Some(outcome);
3966 self.state.save()?;
3967
3968 if self.state.config.graph.land
3974 && mode == MergeMode::Pr
3975 && self.state.status == RunStatus::Merged
3976 {
3977 self.run_land().await?;
3978 }
3979 self.settle_questions();
3984 Ok(())
3985 }
3986
3987 async fn run_land(&mut self) -> Result<()> {
3998 let url = self
3999 .state
4000 .merge
4001 .as_ref()
4002 .map(|m| m.detail.clone())
4003 .unwrap_or_default();
4004 let url = url.lines().next().unwrap_or("").trim().to_owned();
4005 if !url.starts_with("http") {
4006 return Ok(());
4007 }
4008 match land::land(&mut self.state, &url).await {
4011 Ok(pr) if self.state.parked => {
4012 let _ = pr;
4016 }
4017 Ok(pr) => {
4018 self.state.status = match pr.state {
4019 land::PrLifecycle::Merged => RunStatus::Merged,
4020 _ => RunStatus::Blocked,
4021 };
4022 if bump::should_release_bump(self.state.status)
4029 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
4030 {
4031 self.state
4032 .event("bump", format!("release bump skipped: {e:#}"));
4033 }
4034 self.state.save()?;
4035 }
4036 Err(e) => {
4037 self.state.status = RunStatus::Blocked;
4038 self.state.event("land", format!("gave up: {e}"));
4039 self.state.save()?;
4040 }
4041 }
4042 Ok(())
4043 }
4044
4045 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
4049 if let Some(existing) = self.state.seats.get(key)
4050 && existing.agent == agent
4051 {
4052 return existing.clone();
4053 }
4054 let fresh = SeatState::new(key, agent, self.state.seed);
4055 self.state.seats.insert(key.to_owned(), fresh.clone());
4056 fresh
4057 }
4058
4059 fn view(&self, c: &Candidate) -> CandidateView {
4061 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
4062 .unwrap_or_default();
4063 let (patch, _) = blind::sanitize_patch(
4064 &format!("candidate {} patch", c.label),
4065 &raw,
4066 &self.state.config.blind,
4067 );
4068 CandidateView {
4069 label: c.label,
4070 branch: c.branch.clone(),
4071 summary: c.summary.clone(),
4072 stat: c.stat.clone(),
4073 patch,
4074 }
4075 }
4076
4077 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
4079 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
4080 prompt::judge(
4081 "(see above)",
4082 &views,
4083 self.roles.judges.len(),
4084 base_short,
4085 "en",
4086 )
4087 }
4088
4089 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
4096 let mut turns = Vec::new();
4097 for j in &self.state.judgements {
4098 if j.ranking.is_empty() {
4099 continue;
4100 }
4101 let reasons = j
4102 .reasons
4103 .iter()
4104 .map(|(k, v)| format!("- {k}: {v}"))
4105 .collect::<Vec<_>>()
4106 .join("\n");
4107 turns.push(Turn {
4108 who: format!("Judge {} (opening ranking)", j.judge),
4109 is_self: j.judge == self_idx + 1,
4110 body: format!(
4111 "Ranked {}{}{reasons}",
4112 j.ranking.iter().collect::<String>(),
4113 if reasons.is_empty() {
4114 ""
4115 } else {
4116 ", because:\n"
4117 }
4118 ),
4119 });
4120 }
4121 for t in self
4122 .state
4123 .deliberation
4124 .iter()
4125 .flat_map(|r| r.turns.iter())
4126 .chain(current)
4127 {
4128 turns.push(Turn {
4129 who: format!("Judge {}", t.judge),
4130 is_self: t.judge == self_idx + 1,
4131 body: t.body.clone(),
4132 });
4133 }
4134 turns
4135 }
4136}
4137
4138fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
4140 agent::has_session(spec.kind, seat, sessions)
4141}
4142
4143fn has_unconfirmed_command(commands: &[agent::CommandEvidence]) -> bool {
4157 commands.iter().any(|c| c.exit_code.is_none())
4158}
4159
4160fn short(commit: &str) -> String {
4161 commit.chars().take(7).collect()
4162}
4163
4164fn make_executable(path: &Path) -> Result<()> {
4165 #[cfg(unix)]
4166 {
4167 use std::os::unix::fs::PermissionsExt as _;
4168 let mut perms = std::fs::metadata(path)?.permissions();
4169 perms.set_mode(0o755);
4170 std::fs::set_permissions(path, perms)?;
4171 }
4172 #[cfg(not(unix))]
4173 {
4174 let _ = path;
4175 }
4176 Ok(())
4177}
4178
4179struct WaveCtx<'a> {
4186 run: &'a str,
4189 node: &'a str,
4191 prompts: &'a Prompts,
4192 cache: Option<&'a Path>,
4194}
4195
4196async fn run_one(
4198 job: SeatJob,
4199 sem: Arc<Semaphore>,
4200 ctx: &WaveCtx<'_>,
4201 state: &mut RunState,
4202 attempt: usize,
4203) -> (SeatState, AgentOutcome) {
4204 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
4205 .await
4206 .pop()
4207 .expect("one job in, one result out");
4208 (seat, out)
4209}
4210
4211async fn wave(
4217 jobs: Vec<SeatJob>,
4218 sem: Arc<Semaphore>,
4219 ctx: &WaveCtx<'_>,
4220 state: &mut RunState,
4221 attempt: usize,
4222) -> Vec<(usize, SeatState, AgentOutcome)> {
4223 let WaveCtx {
4224 run,
4225 node,
4226 prompts,
4227 cache,
4228 } = *ctx;
4229 for job in &jobs {
4230 state.seat_started(node, &job.seat.key, job.timeout, attempt);
4231 }
4232 if let Err(e) = state.save() {
4233 tracing::warn!("could not persist in-progress seats: {e:#}");
4238 }
4239 let jobs_had_a_writer = jobs.iter().any(|j| j.allow_write);
4255 let wait_started = Instant::now();
4256 let cache_guard = if let Some(cache_dir) = cache {
4257 if jobs_had_a_writer {
4258 let owner = crate::cache::Owner::here(run, node, "*", Path::new("(wave)"), "");
4259 let budget = jobs
4260 .iter()
4261 .map(|j| j.timeout)
4262 .max()
4263 .unwrap_or(Duration::from_secs(60));
4264 acquire_cache_lease(state, cache_dir, &owner, budget, node)
4265 .await
4266 .ok()
4267 } else {
4268 None
4269 }
4270 } else {
4271 None
4272 };
4273 let waited_for_lease = wait_started.elapsed();
4280 let mut set = tokio::task::JoinSet::new();
4281 let overlay = prompts.overlay(node);
4282 for (i, mut job) in jobs.into_iter().enumerate() {
4283 job.timeout = job.timeout.saturating_sub(waited_for_lease);
4284 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
4285 if cache.is_some() {
4286 job.prompt.push('\n');
4287 job.prompt
4288 .push_str(&prompt::build_cache_note(node, job.allow_write));
4289 }
4290 let sem = Arc::clone(&sem);
4291 let run = run.to_owned();
4292 let node = node.to_owned();
4293 let cache = cache
4304 .filter(|_| job.allow_write && cache_guard.is_some())
4305 .map(Path::to_path_buf);
4306 set.spawn(async move {
4307 let _permit = sem.acquire().await;
4308 let mut seat = job.seat;
4309 let out = agent::invoke(
4310 &job.spec,
4311 &mut seat,
4312 &Invocation {
4313 cwd: &job.cwd,
4314 prompt: &job.prompt,
4315 timeout: job.timeout,
4316 allow_write: job.allow_write,
4317 sessions: job.sessions,
4318 artifacts: &job.artifacts,
4319 stem: &job.stem,
4320 run: &run,
4321 node: &node,
4322 cache_dir: cache.as_deref(),
4323 attachments: &[],
4324 },
4325 )
4326 .await;
4327 let out = match out {
4328 Ok(o) if o.usable() => AgentOutcome::Ok(o),
4329 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
4330 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
4338 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
4339 Ok(o) => AgentOutcome::Failed(format!(
4340 "exited with {:?} and no usable output",
4341 o.exit_code
4342 )),
4343 Err(e) => AgentOutcome::Failed(e.to_string()),
4344 };
4345 (i, seat, out)
4346 });
4347 }
4348 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
4349 while let Some(joined) = set.join_next().await {
4350 let (i, seat, out) = match joined {
4351 Ok(v) => v,
4352 Err(e) => {
4356 tracing::error!("agent task panicked: {e}");
4357 continue;
4358 }
4359 };
4360 state.seat_finished(&seat.key);
4361 record_jobs(state, node, &seat.key, &out);
4362 if let Err(e) = state.save() {
4363 tracing::warn!("could not persist a seat's completion: {e:#}");
4364 }
4365 if collected.len() <= i {
4366 collected.resize_with(i + 1, || None);
4367 }
4368 collected[i] = Some((i, seat, out));
4369 }
4370 if state
4376 .active
4377 .values()
4378 .any(|a| a.node == node && a.attempt == attempt)
4379 {
4380 state
4381 .active
4382 .retain(|_, a| !(a.node == node && a.attempt == attempt));
4383 if let Err(e) = state.save() {
4384 tracing::warn!("could not persist the end of a wave: {e:#}");
4385 }
4386 }
4387 if let Some(cache_dir) = cache
4394 && jobs_had_a_writer
4395 {
4396 crate::cache::invalidate_identity(&crate::run::home(), cache_dir);
4397 }
4398 if let Some(guard) = cache_guard {
4399 guard.release();
4400 }
4401 collected.into_iter().flatten().collect()
4402}
4403
4404fn record_jobs(state: &mut RunState, node: &str, seat: &str, out: &AgentOutcome) {
4415 let commands: &[agent::CommandEvidence] = match out {
4416 AgentOutcome::Ok(o) | AgentOutcome::Quota(o) | AgentOutcome::Dropped(o) => &o.commands,
4417 AgentOutcome::Failed(_) => &[],
4418 };
4419 let checked_at = Timestamp::now();
4420 for c in commands {
4421 state.jobs.push(JobRecord {
4422 node: node.to_owned(),
4423 seat: seat.to_owned(),
4424 id: c.id.clone(),
4425 description: c.description.clone(),
4426 checked_at,
4427 status: match c.exit_code {
4428 Some(0) => JobStatus::Completed,
4429 Some(_) => JobStatus::Failed,
4430 None => JobStatus::Unknown,
4431 },
4432 exit_code: c.exit_code,
4433 result_summary: c.result_summary.clone(),
4434 source: c.source.clone(),
4435 });
4436 }
4437}
4438
4439fn round_is_clean(
4460 blocking: usize,
4461 e2e_ok: bool,
4462 answered: usize,
4463 expected: usize,
4464 quota_missing: usize,
4465 policy: IncompleteReviewPolicy,
4466) -> bool {
4467 if blocking != 0 || !e2e_ok {
4468 return false;
4469 }
4470 if answered == expected || policy == IncompleteReviewPolicy::Warn {
4471 return true;
4472 }
4473 answered > 0 && expected - answered <= quota_missing
4474}
4475
4476fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
4493 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
4494 return Some(RunStatus::Gating);
4495 }
4496 let last = reviews.last()?;
4497 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
4498 if reviews.len() < max_rounds && !stagnant {
4499 return None;
4500 }
4501 Some(if last.incomplete() && last.blocking == 0 {
4502 RunStatus::Blocked
4503 } else if last.e2e.iter().all(CommandOutcome::ok) {
4504 RunStatus::Gating
4505 } else {
4506 RunStatus::Blocked
4507 })
4508}
4509
4510fn retry_budget(full: Duration, nudged: bool) -> Duration {
4525 if nudged {
4526 (full / 4).max(Duration::from_secs(120)).min(full)
4527 } else {
4528 full
4529 }
4530}
4531
4532#[allow(clippy::too_many_arguments)]
4545async fn ask_json_wave<T>(
4546 jobs: Vec<SeatJob>,
4547 sem: Arc<Semaphore>,
4548 retries: usize,
4549 ctx: &WaveCtx<'_>,
4550 losses: &mut Vec<QuotaLoss>,
4551 state: &mut RunState,
4552 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
4553) -> Vec<(SeatState, Result<(T, AgentOutput)>)>
4554where
4555 T: serde::de::DeserializeOwned + Send + 'static,
4556{
4557 let n = jobs.len();
4558 let originals: Vec<SeatJob> = jobs;
4559 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
4560 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
4561 let mut pending: Vec<usize> = (0..n).collect();
4562
4563 for attempt in 0..=retries {
4564 if pending.is_empty() {
4565 break;
4566 }
4567 let mut batch = Vec::with_capacity(pending.len());
4568 for &i in &pending {
4569 let src = &originals[i];
4570 let (prompt, timeout) = if attempt == 0 {
4573 (src.prompt.clone(), src.timeout)
4574 } else {
4575 let why = done[i]
4576 .as_ref()
4577 .and_then(|r| r.as_ref().err().map(ToString::to_string))
4578 .unwrap_or_else(|| "no parsable answer".to_owned());
4579 let nudge = prompt::nudge(&why);
4580 let nudged = has_context(&src.spec, &seats[i], src.sessions);
4581 let prompt = if nudged {
4582 nudge
4583 } else {
4584 format!("{}\n\n---\n\n{}", src.prompt, nudge)
4585 };
4586 (prompt, retry_budget(src.timeout, nudged))
4587 };
4588 batch.push(SeatJob {
4589 spec: src.spec.clone(),
4590 seat: seats[i].clone(),
4591 cwd: src.cwd.clone(),
4592 prompt,
4593 timeout,
4594 allow_write: src.allow_write,
4595 sessions: src.sessions,
4596 artifacts: src.artifacts.clone(),
4597 stem: if attempt == 0 {
4598 src.stem.clone()
4599 } else {
4600 format!("{}-retry{attempt}", src.stem)
4601 },
4602 });
4603 }
4604
4605 if attempt > 0 {
4606 let seats_out: Vec<&str> = pending
4607 .iter()
4608 .map(|&i| originals[i].seat.key.as_str())
4609 .collect();
4610 state.event(
4611 ctx.node,
4612 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
4613 );
4614 }
4615 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
4616 let mut still = Vec::new();
4617 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
4618 seats[i] = seat;
4619 let (parsed, quota) = match out {
4620 AgentOutcome::Ok(o) => (
4621 match verdict::extract_json::<T>(&o.text) {
4622 Ok(v) => match validate(&v) {
4623 Ok(()) => Ok((v, o)),
4624 Err(e) => Err(e),
4625 },
4626 Err(e) => Err(e),
4627 },
4628 false,
4629 ),
4630 AgentOutcome::Quota(o) => {
4631 losses.push(QuotaLoss {
4632 seat: originals[i].seat.key.clone(),
4633 node: ctx.node.to_owned(),
4634 at: Timestamp::now(),
4635 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
4636 });
4637 (
4638 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
4639 true,
4640 )
4641 }
4642 AgentOutcome::Dropped(o) => {
4647 let why = o
4648 .dropped
4649 .as_ref()
4650 .map(|d| d.why.as_str())
4651 .unwrap_or("the CLI ended the stream without delivering its answer");
4652 (
4653 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
4654 false,
4655 )
4656 }
4657 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
4658 };
4659 let failed = parsed.is_err();
4660 done[i] = Some(parsed);
4661 if failed && !quota {
4664 still.push(i);
4665 }
4666 }
4667 pending = still;
4668 }
4669
4670 seats
4671 .into_iter()
4672 .zip(done)
4673 .map(|(seat, res)| {
4674 (
4675 seat,
4676 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
4677 )
4678 })
4679 .collect()
4680}
4681
4682async fn acquire_cache_lease(
4695 state: &mut RunState,
4696 cache_dir: &Path,
4697 owner: &crate::cache::Owner,
4698 budget: Duration,
4699 context: &str,
4700) -> Result<crate::cache::Guard> {
4701 let home = crate::run::home();
4702 let started = Instant::now();
4703 let busy = match crate::cache::try_acquire(&home, cache_dir, owner) {
4704 Ok(crate::cache::AcquireOutcome::Acquired(g)) => return Ok(g),
4705 Ok(crate::cache::AcquireOutcome::Busy(busy)) => busy,
4706 Err(e) => {
4707 state.event(
4708 "verify",
4709 format!("{context}: could not check the shared build cache: {e:#}"),
4710 );
4711 if let Err(e2) = state.save() {
4712 tracing::warn!("could not persist a cache-check failure: {e2:#}");
4713 }
4714 return Err(e);
4715 }
4716 };
4717 state.event(
4718 "verify",
4719 format!(
4720 "{context}: waiting for the shared build cache at {} ({})",
4721 cache_dir.display(),
4722 busy.describe()
4723 ),
4724 );
4725 if let Err(e) = state.save() {
4726 tracing::warn!("could not persist a cache wait: {e:#}");
4727 }
4728 let remaining = budget.saturating_sub(started.elapsed());
4729 match crate::cache::wait_for(&home, cache_dir, owner, remaining, Duration::from_secs(5)).await {
4730 Ok(g) => Ok(g),
4731 Err(e) => {
4732 state.event("verify", format!("{context}: {e:#}"));
4733 if let Err(e2) = state.save() {
4734 tracing::warn!("could not persist a cache wait timeout: {e2:#}");
4735 }
4736 Err(e)
4737 }
4738 }
4739}
4740
4741#[allow(clippy::too_many_arguments)]
4762async fn with_cache_lease<'s, F, Fut>(
4763 state: &'s mut RunState,
4764 cache_dir: Option<&Path>,
4765 node: &str,
4766 seat: &str,
4767 worktree: &Path,
4768 head: &str,
4769 budget: Duration,
4770 context: &str,
4771 body: F,
4772) -> (Vec<CommandOutcome>, bool)
4773where
4774 F: FnOnce(&'s mut RunState, Duration) -> Fut,
4775 Fut: std::future::Future<Output = (Vec<CommandOutcome>, bool, Vec<u32>)>,
4776{
4777 let Some(cache_dir) = cache_dir else {
4778 let (outcomes, retried, _timed_out_pids) = body(state, budget).await;
4779 return (outcomes, retried);
4780 };
4781 let home = crate::run::home();
4782 let owner = crate::cache::Owner::here(&state.id, node, seat, worktree, head);
4783 let started = Instant::now();
4784 let guard = match acquire_cache_lease(state, cache_dir, &owner, budget, context).await {
4785 Ok(g) => g,
4786 Err(e) => {
4787 return (
4788 vec![CommandOutcome {
4789 command: "(waiting for the shared build cache)".to_owned(),
4790 code: None,
4791 output_tail: e.to_string(),
4792 duration_ms: started.elapsed().as_millis() as u64,
4793 resource_blocked: true,
4794 }],
4795 false,
4796 );
4797 }
4798 };
4799 let identity = crate::cache::Identity::new(worktree, head);
4800 if let Err(e) = crate::cache::ensure_fresh(&home, cache_dir, &identity) {
4801 state.event(
4809 "verify",
4810 format!(
4811 "{context}: could not confirm the shared build cache matches {} at {}: {e:#}",
4812 worktree.display(),
4813 short(head)
4814 ),
4815 );
4816 guard.release();
4817 return (
4818 vec![CommandOutcome {
4819 command: "(confirming the shared build cache is fresh)".to_owned(),
4820 code: None,
4821 output_tail: e.to_string(),
4822 duration_ms: started.elapsed().as_millis() as u64,
4823 resource_blocked: true,
4824 }],
4825 false,
4826 );
4827 }
4828 let remaining = budget.saturating_sub(started.elapsed());
4829 let (outcomes, retried, timed_out_pids) = body(state, remaining).await;
4830 if !timed_out_pids.is_empty() {
4835 wait_for_timed_out_children_to_die(&timed_out_pids).await;
4836 }
4837 guard.release();
4838 (outcomes, retried)
4839}
4840
4841async fn wait_for_timed_out_children_to_die(pids: &[u32]) {
4853 wait_for_pids_with(
4854 pids,
4855 crate::proc::pid_alive,
4856 LEASE_RELEASE_POLL,
4857 LEASE_RELEASE_MAX_WAIT,
4858 )
4859 .await;
4860}
4861
4862async fn wait_for_pids_with<F: Fn(u32) -> bool>(
4868 pids: &[u32],
4869 alive: F,
4870 poll: Duration,
4871 max_wait: Duration,
4872) {
4873 let deadline = Instant::now() + max_wait;
4874 loop {
4875 if pids.iter().all(|&pid| !alive(pid)) {
4876 return;
4877 }
4878 if Instant::now() >= deadline {
4879 return;
4880 }
4881 tokio::time::sleep(poll).await;
4882 }
4883}
4884
4885fn verify_inconclusive(outcomes: &[CommandOutcome]) -> bool {
4892 outcomes.iter().any(|o| o.resource_blocked)
4893}
4894
4895fn e2e_outcome_label(o: &CommandOutcome) -> String {
4899 if o.ok() {
4900 return "pass".to_owned();
4901 }
4902 let reason = if o.build_failed() {
4903 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
4904 } else {
4905 format!("FAIL ({:?})", o.code)
4906 };
4907 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
4908}
4909
4910async fn run_e2e_with_retry(
4918 state: &mut RunState,
4919 shell: &[String],
4920 commands: &[String],
4921 worktree: &Path,
4922 timeout: Duration,
4923 context: &str,
4924) -> (Vec<CommandOutcome>, bool, Vec<u32>) {
4925 let (mut e2e, mut timed_out_pids) = run_commands(shell, commands, worktree, timeout).await;
4926 for o in &e2e {
4927 state.event(
4928 "verify",
4929 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
4930 );
4931 }
4932 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
4936 if verify_retried {
4937 state.event(
4938 "verify",
4939 format!(
4940 "{context}: verify could not build/link, not a test result — retrying once \
4941 before concluding"
4942 ),
4943 );
4944 let retried = run_commands(shell, commands, worktree, timeout).await;
4945 e2e = retried.0;
4946 timed_out_pids.extend(retried.1);
4949 for o in &e2e {
4950 state.event(
4951 "verify",
4952 format!(
4953 "{context}: retry `{}` -> {}",
4954 o.command,
4955 e2e_outcome_label(o)
4956 ),
4957 );
4958 }
4959 }
4960 (e2e, verify_retried, timed_out_pids)
4961}
4962
4963async fn run_commands(
4969 shell: &[String],
4970 commands: &[String],
4971 cwd: &Path,
4972 timeout: Duration,
4973) -> (Vec<CommandOutcome>, Vec<u32>) {
4974 let mut out = Vec::new();
4975 let mut timed_out_pids = Vec::new();
4976 for command in commands {
4977 let started = Instant::now();
4978 let mut cmd = tokio::process::Command::new(&shell[0]);
4979 cmd.quiet();
4980 cmd.args(&shell[1..])
4981 .arg(command)
4982 .current_dir(cwd)
4983 .stdin(std::process::Stdio::null())
4984 .stdout(std::process::Stdio::piped())
4985 .stderr(std::process::Stdio::piped())
4986 .kill_on_drop(true);
4987 let spawned = cmd.spawn();
4988 let (code, body) = match spawned {
4989 Ok(child) => {
4990 let pid = child.id();
4995 match tokio::time::timeout(timeout, child.wait_with_output()).await {
4996 Ok(Ok(o)) => {
4997 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
4998 body.push_str(&String::from_utf8_lossy(&o.stderr));
4999 (o.status.code(), body)
5000 }
5001 Ok(Err(e)) => (None, format!("failed to run: {e}")),
5002 Err(_) => {
5003 if let Some(pid) = pid {
5004 timed_out_pids.push(pid);
5005 }
5006 (None, format!("timed out after {}s", timeout.as_secs()))
5007 }
5008 }
5009 }
5010 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
5011 };
5012 out.push(CommandOutcome {
5013 command: command.clone(),
5014 code,
5015 output_tail: tail(&body, OUTPUT_TAIL),
5016 duration_ms: started.elapsed().as_millis() as u64,
5017 resource_blocked: false,
5018 });
5019 }
5020 (out, timed_out_pids)
5021}
5022
5023fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
5035 let repo = repo.display();
5036 match style {
5037 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
5038 MergeStyle::Squash => {
5039 let subject = message.lines().next().unwrap_or(branch);
5040 format!(
5041 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
5042 )
5043 }
5044 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
5045 }
5046}
5047
5048fn pr_body(state: &RunState, winner: char) -> String {
5074 let instruction = state.instruction.trim_start();
5075 let mut message = if instruction.is_empty() {
5076 "(empty task)".to_owned()
5077 } else {
5078 instruction.to_owned()
5079 };
5080
5081 let open = state.open_findings();
5082 if !open.is_empty() {
5083 message.push_str("\n\n## Open review findings\n\n");
5084 for f in &open {
5085 message.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
5086 }
5087 }
5088
5089 if let Some(fix) = state.reviews.last().and_then(|r| r.fix.as_ref())
5090 && !fix.rejected.is_empty()
5091 {
5092 message.push_str("\n## Declined by the fixer\n\n");
5093 for r in &fix.rejected {
5094 message.push_str(&format!("- `{}`: {}\n", r.id, r.why));
5095 }
5096 }
5097
5098 message.push_str(&format!(
5099 "\n\n---\nmagi:run/{} magi:candidate-{}\n",
5100 state.id,
5101 winner.to_ascii_lowercase()
5102 ));
5103
5104 message
5105}
5106
5107async fn gh_pr_create(cwd: &Path, base: &str, head: &str, body: &str) -> Result<String> {
5109 let title = body.lines().next().unwrap_or("magi run").to_owned();
5110 let out = tokio::process::Command::new("gh")
5111 .args([
5112 "pr", "create", "--base", base, "--head", head, "--title", &title, "--body", body,
5113 ])
5114 .current_dir(cwd)
5115 .quiet()
5116 .stdin(std::process::Stdio::null())
5117 .output()
5118 .await
5119 .context("spawn gh")?;
5120 if out.status.success() {
5121 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
5122 } else {
5123 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
5124 }
5125}
5126
5127pub async fn fold_run(state: &mut RunState, drop_winner: bool) -> Result<Vec<String>> {
5129 let repo = state.repo.clone();
5130 let root = state.worktree_root();
5131 let winner = state.tally.as_ref().map(|t| t.winner);
5132 let mut removed = Vec::new();
5133
5134 for i in 0..state.candidates.len() {
5135 let c = state.candidates[i].clone();
5136 let is_winner = Some(c.label) == winner;
5137 if is_winner && !drop_winner {
5138 continue;
5139 }
5140 if c.worktree.exists() {
5141 git::worktree_remove(&repo, &c.worktree).await.ok();
5142 removed.push(c.worktree.to_string_lossy().into_owned());
5143 }
5144 if git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
5145 git::branch_delete(&repo, &c.branch).await.ok();
5146 removed.push(c.branch.clone());
5147 }
5148 state.candidates[i].folded = true;
5149 }
5150
5151 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
5152 let path = name.path();
5153 let keep = !drop_winner
5154 && winner.is_some_and(|w| {
5155 path.file_name()
5156 .is_some_and(|n| n == format!("cand-{w}").as_str())
5157 });
5158 if keep {
5159 continue;
5160 }
5161 git::worktree_remove(&repo, &path).await.ok();
5162 removed.push(path.to_string_lossy().into_owned());
5163 }
5164
5165 remove_if_empty(&root);
5174
5175 if state.enabled_worktree_config && drop_winner {
5176 git::release_worktree_config(&repo).await.ok();
5180 state.enabled_worktree_config = false;
5181 }
5182 state.save()?;
5183 Ok(removed)
5184}
5185
5186fn remove_if_empty(dir: &Path) {
5197 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
5198 std::fs::remove_dir(dir).ok();
5199 }
5200}
5201
5202pub fn worst_open(state: &RunState) -> Option<Severity> {
5204 state
5205 .reviews
5206 .last()?
5207 .reviews
5208 .iter()
5209 .flat_map(|r| r.findings.iter())
5210 .map(|f| f.severity)
5211 .max()
5212}
5213
5214#[cfg(test)]
5215mod tests {
5216 use super::*;
5217 use crate::run::GateStatus;
5218 use std::collections::BTreeMap;
5219 use std::time::Duration;
5220
5221 fn conductor() -> AgentSpec {
5222 AgentSpec {
5223 id: "conductor".to_owned(),
5224 kind: crate::config::AgentKind::Command,
5225 model: None,
5226 command: vec!["true".to_owned()],
5227 extra_args: Vec::new(),
5228 env: BTreeMap::new(),
5229 prompt_delivery: None,
5230 }
5231 }
5232
5233 #[test]
5234 fn remove_if_empty_only_ever_takes_a_bare_directory() {
5235 let dir = tempfile::tempdir().unwrap();
5236 let bay = dir.path().join("ffff");
5237
5238 remove_if_empty(&bay);
5240 assert!(!bay.exists());
5241
5242 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
5245 remove_if_empty(&bay);
5246 assert!(bay.exists(), "non-empty directory must survive");
5247
5248 std::fs::remove_dir(bay.join("cand-A")).unwrap();
5250 remove_if_empty(&bay);
5251 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
5252 }
5253
5254 #[test]
5263 fn a_full_panel_that_found_nothing_is_clean() {
5264 assert!(round_is_clean(
5265 0,
5266 true,
5267 2,
5268 2,
5269 0,
5270 IncompleteReviewPolicy::Block
5271 ));
5272 }
5273
5274 #[test]
5275 fn a_missing_seat_is_never_clean_under_the_default_policy() {
5276 assert!(!round_is_clean(
5277 0,
5278 true,
5279 1,
5280 2,
5281 0,
5282 IncompleteReviewPolicy::Block
5283 ));
5284 }
5285
5286 #[test]
5287 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
5288 assert!(!round_is_clean(
5289 1,
5290 true,
5291 1,
5292 2,
5293 0,
5294 IncompleteReviewPolicy::Warn
5295 ));
5296 }
5297
5298 #[test]
5299 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
5300 assert!(round_is_clean(
5301 0,
5302 true,
5303 1,
5304 2,
5305 0,
5306 IncompleteReviewPolicy::Warn
5307 ));
5308 }
5309
5310 #[test]
5311 fn a_full_panel_with_an_open_finding_is_not_clean() {
5312 assert!(!round_is_clean(
5313 1,
5314 true,
5315 2,
5316 2,
5317 0,
5318 IncompleteReviewPolicy::Block
5319 ));
5320 }
5321
5322 #[test]
5323 fn a_full_panel_with_a_red_e2e_is_not_clean() {
5324 assert!(!round_is_clean(
5325 0,
5326 false,
5327 2,
5328 2,
5329 0,
5330 IncompleteReviewPolicy::Block
5331 ));
5332 }
5333
5334 #[test]
5341 fn a_seat_missing_only_to_its_own_quota_is_clean_under_the_default_policy() {
5342 assert!(round_is_clean(
5345 0,
5346 true,
5347 1,
5348 2,
5349 1,
5350 IncompleteReviewPolicy::Block
5351 ));
5352 }
5353
5354 #[test]
5355 fn a_seat_missing_for_a_reason_other_than_quota_still_waits() {
5356 assert!(!round_is_clean(
5359 0,
5360 true,
5361 1,
5362 2,
5363 0,
5364 IncompleteReviewPolicy::Block
5365 ));
5366 }
5367
5368 #[test]
5369 fn a_quota_loss_does_not_excuse_an_open_finding_or_a_red_e2e() {
5370 assert!(!round_is_clean(
5371 1,
5372 true,
5373 1,
5374 2,
5375 1,
5376 IncompleteReviewPolicy::Block
5377 ));
5378 assert!(!round_is_clean(
5379 0,
5380 false,
5381 1,
5382 2,
5383 1,
5384 IncompleteReviewPolicy::Block
5385 ));
5386 }
5387
5388 #[test]
5389 fn a_panel_lost_entirely_to_quota_still_waits_rather_than_deciding_on_nobody() {
5390 assert!(!round_is_clean(
5394 0,
5395 true,
5396 0,
5397 2,
5398 2,
5399 IncompleteReviewPolicy::Block
5400 ));
5401 }
5402
5403 fn outcome(code: Option<i32>, resource_blocked: bool) -> CommandOutcome {
5404 CommandOutcome {
5405 command: "test".to_owned(),
5406 code,
5407 output_tail: String::new(),
5408 duration_ms: 0,
5409 resource_blocked,
5410 }
5411 }
5412
5413 #[test]
5414 fn verify_is_inconclusive_only_when_a_resource_blocked_outcome_is_present() {
5415 assert!(!verify_inconclusive(&[outcome(Some(0), false)]));
5416 assert!(
5417 !verify_inconclusive(&[outcome(Some(1), false)]),
5418 "an ordinary failure is still evidence about the patch"
5419 );
5420 assert!(verify_inconclusive(&[outcome(None, true)]));
5421 assert!(
5422 verify_inconclusive(&[outcome(Some(0), false), outcome(None, true)]),
5423 "one inconclusive outcome taints the whole batch"
5424 );
5425 assert!(!verify_inconclusive(&[]));
5426 }
5427
5428 #[tokio::test]
5429 async fn timed_out_pid_waiting_returns_as_soon_as_every_pid_is_confirmed_dead() {
5430 let calls = std::sync::atomic::AtomicUsize::new(0);
5434 let started = Instant::now();
5435 wait_for_pids_with(
5436 &[123],
5437 |_| calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2,
5438 Duration::from_millis(5),
5439 Duration::from_secs(5),
5440 )
5441 .await;
5442 assert!(
5443 calls.load(std::sync::atomic::Ordering::SeqCst) >= 3,
5444 "must keep checking rather than deciding on the first answer"
5445 );
5446 assert!(
5447 started.elapsed() < Duration::from_secs(1),
5448 "must return the moment it is confirmed dead, not wait out the ceiling"
5449 );
5450 }
5451
5452 #[tokio::test]
5453 async fn timed_out_pid_waiting_gives_up_at_its_ceiling_if_never_confirmed_dead() {
5454 let started = Instant::now();
5455 wait_for_pids_with(
5456 &[123],
5457 |_| true, Duration::from_millis(5),
5459 Duration::from_millis(30),
5460 )
5461 .await;
5462 let elapsed = started.elapsed();
5463 assert!(
5464 elapsed >= Duration::from_millis(30),
5465 "must not give up before its own ceiling: {elapsed:?}"
5466 );
5467 assert!(
5468 elapsed < Duration::from_secs(1),
5469 "must not wait past its own ceiling either: {elapsed:?}"
5470 );
5471 }
5472
5473 #[tokio::test]
5474 async fn timed_out_pid_waiting_is_a_no_op_when_nothing_was_still_running() {
5475 let started = Instant::now();
5476 wait_for_pids_with(
5477 &[],
5478 |_| true,
5479 Duration::from_secs(5),
5480 Duration::from_secs(5),
5481 )
5482 .await;
5483 assert!(
5484 started.elapsed() < Duration::from_millis(200),
5485 "an empty pid list has nothing to confirm"
5486 );
5487 }
5488
5489 fn review_round(
5495 clean: bool,
5496 blocking: usize,
5497 answered: usize,
5498 expected: usize,
5499 progressed: bool,
5500 e2e_ok: bool,
5501 ) -> ReviewRound {
5502 ReviewRound {
5503 round: 1,
5504 head: "h".to_owned(),
5505 verified_head: None,
5506 reviews: Vec::new(),
5507 e2e: vec![CommandOutcome {
5508 command: "test".to_owned(),
5509 code: Some(if e2e_ok { 0 } else { 1 }),
5510 output_tail: String::new(),
5511 duration_ms: 0,
5512 resource_blocked: false,
5513 }],
5514 verify_retried: false,
5515 e2e_deferred: false,
5516 e2e_defer_reason: None,
5517 fix: None,
5518 blocking,
5519 answered,
5520 expected,
5521 clean,
5522 progressed,
5523 vote_split: false,
5524 reconsideration: Vec::new(),
5525 verdict: None,
5526 }
5527 }
5528
5529 #[test]
5530 fn review_conclusion_is_none_when_nothing_has_run() {
5531 assert_eq!(review_conclusion(&[], 3), None);
5532 }
5533
5534 #[test]
5535 fn review_conclusion_is_none_while_rounds_remain() {
5536 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
5537 assert_eq!(review_conclusion(&rounds, 3), None);
5538 }
5539
5540 #[test]
5541 fn review_conclusion_is_gating_once_a_round_is_clean() {
5542 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
5543 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
5544 }
5545
5546 #[test]
5547 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
5548 let rounds = vec![
5549 review_round(false, 1, 2, 2, true, true),
5550 review_round(false, 1, 2, 2, true, true),
5551 ];
5552 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
5553 }
5554
5555 #[test]
5556 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
5557 let rounds = vec![
5558 review_round(false, 1, 2, 2, true, true),
5559 review_round(false, 1, 2, 2, true, false),
5560 ];
5561 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
5562 }
5563
5564 #[test]
5565 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
5566 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
5568 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
5569 }
5570
5571 #[test]
5572 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
5573 let rounds = vec![
5574 review_round(false, 1, 2, 2, false, true),
5575 review_round(false, 1, 2, 2, false, true),
5576 ];
5577 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
5578 }
5579
5580 fn secs(n: u64) -> Duration {
5581 Duration::from_secs(n)
5582 }
5583
5584 fn init_repo(dir: &Path) {
5587 let run = |args: &[&str]| {
5588 let out = std::process::Command::new("git")
5589 .args(args)
5590 .current_dir(dir)
5591 .quiet()
5592 .output()
5593 .expect("spawn git");
5594 assert!(
5595 out.status.success(),
5596 "git {args:?} failed: {}",
5597 String::from_utf8_lossy(&out.stderr)
5598 );
5599 };
5600 run(&["init", "-b", "main"]);
5601 run(&["config", "user.name", "magi test"]);
5602 run(&["config", "user.email", "magi@example.com"]);
5603 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
5604 run(&["add", "-A"]);
5605 run(&["commit", "-m", "init"]);
5606 }
5607
5608 fn ask_test_home() {
5616 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
5617 }
5618
5619 fn runner_at(status: RunStatus) -> Runner {
5622 let mut state = RunState::new(
5623 PathBuf::from("/nonexistent/repo"),
5624 "main".to_owned(),
5625 "deadbeef".to_owned(),
5626 "task".to_owned(),
5627 Config::default(),
5628 );
5629 state.status = status;
5630 Runner {
5631 state,
5632 roles: ResolvedRoles {
5633 implementers: Vec::new(),
5634 judges: Vec::new(),
5635 reviewers: Vec::new(),
5636 fixer: None,
5637 conductor: conductor(),
5638 },
5639 sem: Arc::new(Semaphore::new(1)),
5640 pause: Pause::new(),
5641 interrupt: Pause::new(),
5642 }
5643 }
5644
5645 #[test]
5649 fn park_here_folds_the_interrupt_reason_into_the_park_event() {
5650 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
5651 let mut runner = runner_at(RunStatus::Implementing);
5652 let interrupt = Pause::new();
5653 runner.watch_interrupt(interrupt.clone());
5654
5655 interrupt.park_because("task a1b2 asked to run first");
5656
5657 assert!(runner.park_here().expect("park_here"));
5658 assert!(runner.state.parked);
5659 let last = runner.state.events.last().expect("a park event");
5660 assert_eq!(last.node, "park");
5661 assert!(
5662 last.message.contains("task a1b2 asked to run first"),
5663 "expected the interrupt reason in {:?}",
5664 last.message
5665 );
5666 }
5667
5668 #[test]
5676 fn the_stop_level_pause_and_a_runs_interrupt_pause_do_not_leak_into_each_other() {
5677 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
5678 let mut runner = runner_at(RunStatus::Implementing);
5679 let shutdown = Pause::new();
5680 runner.on_pause(shutdown.clone());
5681 let interrupt = Pause::new();
5682 runner.watch_interrupt(interrupt.clone());
5683
5684 assert!(!runner.park_here().expect("park_here"));
5686 assert!(!runner.state.parked);
5687
5688 interrupt.park_because("test");
5690 assert!(!shutdown.parked());
5691 assert!(runner.park_here().expect("park_here"));
5692 }
5693
5694 #[tokio::test]
5708 async fn a_park_request_made_mid_node_only_takes_effect_at_the_next_boundary() {
5709 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
5710 let mut runner = runner_at(RunStatus::Implementing);
5711 let interrupt = Pause::new();
5712 runner.watch_interrupt(interrupt.clone());
5713
5714 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
5715 let (finish_tx, finish_rx) = tokio::sync::oneshot::channel::<()>();
5716
5717 let node = async move {
5721 started_tx.send(()).expect("send started");
5722 finish_rx.await.expect("recv finish");
5723 "node finished"
5724 };
5725
5726 let interrupter = async move {
5727 started_rx.await.expect("recv started");
5728 interrupt.park_because("higher-priority task waiting");
5730 tokio::task::yield_now().await;
5734 finish_tx.send(()).expect("send finish");
5735 };
5736
5737 let (node_result, ()) = tokio::join!(node, interrupter);
5738 assert_eq!(
5739 node_result, "node finished",
5740 "the in-flight call ran to completion"
5741 );
5742
5743 assert!(runner.park_here().expect("park_here"));
5746 assert!(runner.state.parked);
5747 }
5748
5749 #[test]
5755 fn a_run_parked_for_an_interrupt_resumes_with_nothing_lost() {
5756 crate::run::set_home(std::env::temp_dir().join("magi-graph-interrupt-tests-home"));
5757 let mut runner = runner_at(RunStatus::Judging);
5758 runner.state.config.agents = vec![conductor()];
5762 runner.state.candidates = vec![Candidate {
5763 index: 0,
5764 label: 'A',
5765 agent: "alpha".to_owned(),
5766 branch: "magi/x/A".to_owned(),
5767 worktree: PathBuf::from("/nonexistent/worktree"),
5768 summary: "did the thing".to_owned(),
5769 stat: "1 file changed".to_owned(),
5770 files: 1,
5771 commits: 1,
5772 empty: false,
5773 failed: None,
5774 duration_ms: 1234,
5775 folded: false,
5776 }];
5777 let run_id = runner.state.id.clone();
5778
5779 let interrupt = Pause::new();
5780 runner.watch_interrupt(interrupt.clone());
5781 interrupt.park_because("task c3d4 asked to run first");
5782 assert!(runner.park_here().expect("park_here"));
5783
5784 let resumed = Runner::resume(&run_id).expect("resume");
5785 assert_eq!(resumed.state.candidates.len(), 1);
5786 assert_eq!(resumed.state.candidates[0].summary, "did the thing");
5787 assert_eq!(resumed.state.candidates[0].branch, "magi/x/A");
5788 assert_eq!(resumed.state.status, runner.state.status);
5789 assert!(
5790 resumed.state.parked,
5791 "still parked until `execute` actually walks the graph again"
5792 );
5793 assert!(resumed.state.events.iter().any(|e| e.node == "park"));
5794 }
5795
5796 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
5798 let mut q = ask::Question::new(
5799 run.to_owned(),
5800 "implement".to_owned(),
5801 "impl-A".to_owned(),
5802 "Which storage backend should the cache use?".to_owned(),
5803 String::new(),
5804 vec!["SQLite".to_owned(), "Redis".to_owned()],
5805 );
5806 store.put(&mut q).unwrap();
5807 q
5808 }
5809
5810 #[test]
5811 fn a_failed_runs_open_question_is_abandoned() {
5812 ask_test_home();
5813 let store = ask::Questions::open();
5814 let mut runner = runner_at(RunStatus::Failed);
5815 let run = runner.state.id.clone();
5816 let q = ask_open_question(&store, &run);
5817
5818 runner.settle_questions();
5819
5820 let back = store.get(&q.id).unwrap();
5821 assert!(
5822 !back.status.open(),
5823 "the seat that asked died with the run; nobody is left to read an answer"
5824 );
5825 assert!(
5826 back.detail.contains(&run) && back.detail.contains("failed"),
5827 "the reason names what the run became, not just that it is gone: {}",
5828 back.detail
5829 );
5830 }
5831
5832 #[test]
5833 fn a_merged_runs_open_question_is_abandoned_too() {
5834 ask_test_home();
5835 let store = ask::Questions::open();
5836 for status in [RunStatus::Merged, RunStatus::Ready] {
5839 let mut runner = runner_at(status);
5840 let run = runner.state.id.clone();
5841 let q = ask_open_question(&store, &run);
5842
5843 runner.settle_questions();
5844
5845 let back = store.get(&q.id).unwrap();
5846 assert!(
5847 !back.status.open(),
5848 "{status:?} run's question must not outlive the run"
5849 );
5850 }
5851 }
5852
5853 #[test]
5854 fn a_still_resumable_runs_open_question_is_left_alone() {
5855 ask_test_home();
5856 let store = ask::Questions::open();
5857 for status in [RunStatus::Blocked, RunStatus::Stalled] {
5863 let mut runner = runner_at(status);
5864 let run = runner.state.id.clone();
5865 let q = ask_open_question(&store, &run);
5866
5867 runner.settle_questions();
5868
5869 let back = store.get(&q.id).unwrap();
5870 assert!(
5871 back.status.open(),
5872 "{status:?} is still alive; the question must still be waiting"
5873 );
5874 }
5875 }
5876
5877 #[test]
5878 fn settle_questions_never_touches_an_already_answered_question() {
5879 ask_test_home();
5880 let store = ask::Questions::open();
5881 let mut runner = runner_at(RunStatus::Failed);
5882 let run = runner.state.id.clone();
5883 let mut q = ask_open_question(&store, &run);
5884 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
5885 .unwrap();
5886 store.put(&mut q).unwrap();
5887
5888 runner.settle_questions();
5893 runner.settle_questions();
5894
5895 let back = store.get(&q.id).unwrap();
5896 assert_eq!(
5897 back.status,
5898 ask::QuestionStatus::Answered,
5899 "a real answer is a decision on record, never overwritten by a sweep"
5900 );
5901 }
5902
5903 #[tokio::test]
5914 async fn fold_run_keeps_only_the_winner_when_the_winner_is_not_dropped() {
5915 crate::run::set_home(std::env::temp_dir().join("magi-graph-fold-run-tests-home"));
5916 let tmp = tempfile::tempdir().expect("tempdir");
5917 let repo = tmp.path().join("repo");
5918 std::fs::create_dir_all(&repo).unwrap();
5919 init_repo(&repo);
5920
5921 let mut config = Config::default();
5922 config.graph.worktree_root = Some(tmp.path().join("wt"));
5923
5924 let mut state = RunState::new(
5925 repo.clone(),
5926 "main".to_owned(),
5927 "deadbeef".to_owned(),
5928 "task".to_owned(),
5929 config,
5930 );
5931 let root = state.worktree_root();
5932 let wt_a = root.join("cand-A");
5933 let wt_b = root.join("cand-B");
5934 git::worktree_add_branch(&repo, &wt_a, "magi/x/A", "main")
5935 .await
5936 .expect("worktree A");
5937 git::worktree_add_branch(&repo, &wt_b, "magi/x/B", "main")
5938 .await
5939 .expect("worktree B");
5940
5941 state.candidates = vec![
5942 Candidate {
5943 index: 0,
5944 label: 'A',
5945 agent: "alpha".to_owned(),
5946 branch: "magi/x/A".to_owned(),
5947 worktree: wt_a.clone(),
5948 summary: String::new(),
5949 stat: String::new(),
5950 files: 0,
5951 commits: 0,
5952 empty: false,
5953 failed: None,
5954 duration_ms: 0,
5955 folded: false,
5956 },
5957 Candidate {
5958 index: 1,
5959 label: 'B',
5960 agent: "beta".to_owned(),
5961 branch: "magi/x/B".to_owned(),
5962 worktree: wt_b.clone(),
5963 summary: String::new(),
5964 stat: String::new(),
5965 files: 0,
5966 commits: 0,
5967 empty: false,
5968 failed: None,
5969 duration_ms: 0,
5970 folded: false,
5971 },
5972 ];
5973 state.tally = Some(Tally {
5974 first_choice: BTreeMap::from([('A', 1)]),
5975 borda: BTreeMap::new(),
5976 winner: 'A',
5977 rankings: 1,
5978 unanimous_initial: true,
5979 deliberated: false,
5980 changed_votes: 0,
5981 unanimous_final: true,
5982 tie_break: None,
5983 judges: 1,
5984 present: 1,
5985 quorum: 1,
5986 met_quorum: true,
5987 uncontested: None,
5988 });
5989 state.status = RunStatus::Ready;
5990
5991 fold_run(&mut state, false).await.expect("fold_run");
5992
5993 assert!(wt_a.exists(), "the unmerged winner's worktree survives");
5994 assert!(
5995 git::branch_exists(&repo, "magi/x/A").await.unwrap(),
5996 "the unmerged winner's branch survives"
5997 );
5998 assert!(
5999 !state.candidates[0].folded,
6000 "the winner is not marked folded"
6001 );
6002
6003 assert!(!wt_b.exists(), "the loser's worktree is removed");
6004 assert!(
6005 !git::branch_exists(&repo, "magi/x/B").await.unwrap(),
6006 "the loser's branch is removed"
6007 );
6008 assert!(state.candidates[1].folded, "the loser is marked folded");
6009 }
6010
6011 #[tokio::test]
6020 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
6021 let tmp = tempfile::tempdir().expect("tempdir");
6022 let repo = tmp.path().join("repo");
6023 std::fs::create_dir_all(&repo).unwrap();
6024 init_repo(&repo);
6025
6026 let mut config = Config::default();
6027 config.merge.mode = MergeMode::Local;
6028
6029 let mut state = RunState::new(
6030 repo.clone(),
6031 "main".to_owned(),
6032 "deadbeef".to_owned(),
6033 "task".to_owned(),
6034 config,
6035 );
6036 state.candidates = vec![Candidate {
6037 index: 0,
6038 label: 'A',
6039 agent: "alpha".to_owned(),
6040 branch: "does-not-exist".to_owned(),
6041 worktree: repo.clone(),
6042 summary: String::new(),
6043 stat: String::new(),
6044 files: 0,
6045 commits: 0,
6046 empty: false,
6047 failed: None,
6048 duration_ms: 0,
6049 folded: false,
6050 }];
6051 state.tally = Some(Tally {
6052 first_choice: BTreeMap::from([('A', 1)]),
6053 borda: BTreeMap::new(),
6054 winner: 'A',
6055 rankings: 1,
6056 unanimous_initial: true,
6057 deliberated: false,
6058 changed_votes: 0,
6059 unanimous_final: true,
6060 tie_break: None,
6061 judges: 0,
6062 present: 0,
6063 quorum: 0,
6064 met_quorum: true,
6065 uncontested: Some("only candidate A produced a change".to_owned()),
6066 });
6067 state.reviews = vec![ReviewRound {
6068 round: 1,
6069 head: "deadbeef".to_owned(),
6070 verified_head: None,
6071 reviews: Vec::new(),
6072 e2e: Vec::new(),
6073 fix: None,
6074 blocking: 0,
6075 answered: 0,
6076 expected: 0,
6077 clean: true,
6078 verify_retried: false,
6079 e2e_deferred: false,
6080 e2e_defer_reason: None,
6081 progressed: false,
6082 vote_split: false,
6083 reconsideration: Vec::new(),
6084 verdict: None,
6085 }];
6086 state.gate = vec![CommandOutcome {
6087 command: "test".to_owned(),
6088 code: Some(0),
6089 output_tail: String::new(),
6090 duration_ms: 0,
6091 resource_blocked: false,
6092 }];
6093 state.gate_ran = true;
6094 state.status = RunStatus::Ready;
6099 state.merge = Some(MergeOutcome {
6100 mode: MergeMode::Local,
6101 ok: false,
6102 detail: "already concluded".to_owned(),
6103 });
6104
6105 let mut runner = Runner {
6106 state,
6107 roles: ResolvedRoles {
6108 implementers: Vec::new(),
6109 judges: Vec::new(),
6110 reviewers: Vec::new(),
6111 fixer: None,
6112 conductor: conductor(),
6113 },
6114 sem: Arc::new(Semaphore::new(1)),
6115 pause: Pause::new(),
6116 interrupt: Pause::new(),
6117 };
6118
6119 runner.merge().await.expect("merge");
6120
6121 assert_eq!(
6122 runner.state.status,
6123 RunStatus::Ready,
6124 "a concluded run's status must not change on reentry"
6125 );
6126 assert_eq!(
6127 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
6128 Some("already concluded"),
6129 "merge must not run again once the node already recorded an outcome"
6130 );
6131 }
6132
6133 #[tokio::test]
6142 async fn merge_refuses_a_gate_that_has_not_actually_run() {
6143 let tmp = tempfile::tempdir().expect("tempdir");
6144 let repo = tmp.path().join("repo");
6145 std::fs::create_dir_all(&repo).unwrap();
6146 init_repo(&repo);
6147
6148 let mut config = Config::default();
6149 config.merge.mode = MergeMode::Local;
6150
6151 let mut state = RunState::new(
6152 repo.clone(),
6153 "main".to_owned(),
6154 "deadbeef".to_owned(),
6155 "task".to_owned(),
6156 config,
6157 );
6158 state.candidates = vec![Candidate {
6159 index: 0,
6160 label: 'A',
6161 agent: "alpha".to_owned(),
6162 branch: "does-not-exist".to_owned(),
6163 worktree: repo.clone(),
6164 summary: String::new(),
6165 stat: String::new(),
6166 files: 0,
6167 commits: 0,
6168 empty: false,
6169 failed: None,
6170 duration_ms: 0,
6171 folded: false,
6172 }];
6173 state.tally = Some(Tally {
6174 first_choice: BTreeMap::from([('A', 1)]),
6175 borda: BTreeMap::new(),
6176 winner: 'A',
6177 rankings: 1,
6178 unanimous_initial: true,
6179 deliberated: false,
6180 changed_votes: 0,
6181 unanimous_final: true,
6182 tie_break: None,
6183 judges: 0,
6184 present: 0,
6185 quorum: 0,
6186 met_quorum: true,
6187 uncontested: Some("only candidate A produced a change".to_owned()),
6188 });
6189 state.reviews = vec![ReviewRound {
6190 round: 1,
6191 head: "deadbeef".to_owned(),
6192 verified_head: None,
6193 reviews: Vec::new(),
6194 e2e: Vec::new(),
6195 fix: None,
6196 blocking: 0,
6197 answered: 0,
6198 expected: 0,
6199 clean: true,
6200 verify_retried: false,
6201 e2e_deferred: false,
6202 e2e_defer_reason: None,
6203 progressed: false,
6204 vote_split: false,
6205 reconsideration: Vec::new(),
6206 verdict: None,
6207 }];
6208 state.gate = Vec::new();
6210 state.gate_ran = false;
6211 state.status = RunStatus::Gating;
6212
6213 let mut runner = Runner {
6214 state,
6215 roles: ResolvedRoles {
6216 implementers: Vec::new(),
6217 judges: Vec::new(),
6218 reviewers: Vec::new(),
6219 fixer: None,
6220 conductor: conductor(),
6221 },
6222 sem: Arc::new(Semaphore::new(1)),
6223 pause: Pause::new(),
6224 interrupt: Pause::new(),
6225 };
6226
6227 runner.merge().await.expect("merge");
6228
6229 assert!(
6230 runner.state.merge.is_none(),
6231 "an empty gate must never be read as a passing one: {:?}",
6232 runner.state.merge
6233 );
6234 }
6235
6236 #[tokio::test]
6243 async fn gate_and_merge_reach_ready_when_no_gate_commands_are_configured() {
6244 let tmp = tempfile::tempdir().expect("tempdir");
6245 let repo = tmp.path().join("repo");
6246 std::fs::create_dir_all(&repo).unwrap();
6247 init_repo(&repo);
6248
6249 let config = Config::default();
6251
6252 let mut state = RunState::new(
6253 repo.clone(),
6254 "main".to_owned(),
6255 "deadbeef".to_owned(),
6256 "task".to_owned(),
6257 config,
6258 );
6259 state.candidates = vec![Candidate {
6260 index: 0,
6261 label: 'A',
6262 agent: "alpha".to_owned(),
6263 branch: "does-not-exist".to_owned(),
6264 worktree: repo.clone(),
6265 summary: String::new(),
6266 stat: String::new(),
6267 files: 0,
6268 commits: 0,
6269 empty: false,
6270 failed: None,
6271 duration_ms: 0,
6272 folded: false,
6273 }];
6274 state.tally = Some(Tally {
6275 first_choice: BTreeMap::from([('A', 1)]),
6276 borda: BTreeMap::new(),
6277 winner: 'A',
6278 rankings: 1,
6279 unanimous_initial: true,
6280 deliberated: false,
6281 changed_votes: 0,
6282 unanimous_final: true,
6283 tie_break: None,
6284 judges: 0,
6285 present: 0,
6286 quorum: 0,
6287 met_quorum: true,
6288 uncontested: Some("only candidate A produced a change".to_owned()),
6289 });
6290 state.reviews = vec![ReviewRound {
6291 round: 1,
6292 head: "deadbeef".to_owned(),
6293 verified_head: None,
6294 reviews: Vec::new(),
6295 e2e: Vec::new(),
6296 fix: None,
6297 blocking: 0,
6298 answered: 0,
6299 expected: 0,
6300 clean: true,
6301 verify_retried: false,
6302 e2e_deferred: false,
6303 e2e_defer_reason: None,
6304 progressed: false,
6305 vote_split: false,
6306 reconsideration: Vec::new(),
6307 verdict: None,
6308 }];
6309
6310 let mut runner = Runner {
6311 state,
6312 roles: ResolvedRoles {
6313 implementers: Vec::new(),
6314 judges: Vec::new(),
6315 reviewers: Vec::new(),
6316 fixer: None,
6317 conductor: conductor(),
6318 },
6319 sem: Arc::new(Semaphore::new(1)),
6320 pause: Pause::new(),
6321 interrupt: Pause::new(),
6322 };
6323
6324 runner.gate().await.expect("gate");
6325 assert!(
6326 runner.state.gate_ran,
6327 "zero configured commands is still a real attempt, not an unrun gate"
6328 );
6329 assert!(runner.state.gate.is_empty());
6330 assert_eq!(runner.state.gate_status(), GateStatus::PassedWithNoCommands);
6331 assert_ne!(
6332 runner.state.status,
6333 RunStatus::Blocked,
6334 "a gate with nothing to check must not read as failed"
6335 );
6336
6337 runner.merge().await.expect("merge");
6338 assert_eq!(
6339 runner.state.status,
6340 RunStatus::Ready,
6341 "a clean review-only run with no gate commands must reach Ready, not stay stuck in Gating"
6342 );
6343 }
6344
6345 #[tokio::test]
6356 async fn gate_never_asks_for_the_cache_lease_when_it_has_no_commands_to_run() {
6357 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
6358 let home = crate::run::home();
6359
6360 let tmp = tempfile::tempdir().expect("tempdir");
6361 let repo = tmp.path().join("repo");
6362 std::fs::create_dir_all(&repo).unwrap();
6363 init_repo(&repo);
6364 let cache_dir = tmp.path().join("target");
6367
6368 let mut config = Config::default();
6369 config.verify.e2e = vec![format!("CARGO_TARGET_DIR='{}' true", cache_dir.display())];
6370 config.graph.timeout_verify = Some(2);
6373
6374 let other = crate::cache::Owner::here("other-run", "e2e", "e2e", &repo, "deadbeef");
6375 let _held = match crate::cache::try_acquire(&home, &cache_dir, &other)
6376 .expect("no io error acquiring directly")
6377 {
6378 crate::cache::AcquireOutcome::Acquired(g) => g,
6379 crate::cache::AcquireOutcome::Busy(b) => {
6380 panic!("expected the direct acquire to win the lease first: {b:?}")
6381 }
6382 };
6383
6384 let mut state = RunState::new(
6385 repo.clone(),
6386 "main".to_owned(),
6387 "deadbeef".to_owned(),
6388 "task".to_owned(),
6389 config,
6390 );
6391 state.candidates = vec![Candidate {
6392 index: 0,
6393 label: 'A',
6394 agent: "alpha".to_owned(),
6395 branch: "does-not-exist".to_owned(),
6396 worktree: repo.clone(),
6397 summary: String::new(),
6398 stat: String::new(),
6399 files: 0,
6400 commits: 0,
6401 empty: false,
6402 failed: None,
6403 duration_ms: 0,
6404 folded: false,
6405 }];
6406 state.tally = Some(Tally {
6407 first_choice: BTreeMap::from([('A', 1)]),
6408 borda: BTreeMap::new(),
6409 winner: 'A',
6410 rankings: 1,
6411 unanimous_initial: true,
6412 deliberated: false,
6413 changed_votes: 0,
6414 unanimous_final: true,
6415 tie_break: None,
6416 judges: 0,
6417 present: 0,
6418 quorum: 0,
6419 met_quorum: true,
6420 uncontested: Some("only candidate A produced a change".to_owned()),
6421 });
6422 state.reviews = vec![ReviewRound {
6423 round: 1,
6424 head: "deadbeef".to_owned(),
6425 verified_head: None,
6426 reviews: Vec::new(),
6427 e2e: Vec::new(),
6428 fix: None,
6429 blocking: 0,
6430 answered: 0,
6431 expected: 0,
6432 clean: true,
6433 verify_retried: false,
6434 e2e_deferred: false,
6435 e2e_defer_reason: None,
6436 progressed: false,
6437 vote_split: false,
6438 reconsideration: Vec::new(),
6439 verdict: None,
6440 }];
6441
6442 let mut runner = Runner {
6443 state,
6444 roles: ResolvedRoles {
6445 implementers: Vec::new(),
6446 judges: Vec::new(),
6447 reviewers: Vec::new(),
6448 fixer: None,
6449 conductor: conductor(),
6450 },
6451 sem: Arc::new(Semaphore::new(1)),
6452 pause: Pause::new(),
6453 interrupt: Pause::new(),
6454 };
6455
6456 let started = std::time::Instant::now();
6457 runner.gate().await.expect("gate");
6458 assert!(
6459 started.elapsed() < Duration::from_secs(1),
6460 "a gate with nothing to run must never wait on a lease it never needed"
6461 );
6462 assert!(
6463 runner.state.gate_ran,
6464 "zero commands is still a real, immediate attempt"
6465 );
6466 assert!(runner.state.gate.is_empty());
6467 assert_ne!(
6468 runner.state.status,
6469 RunStatus::Blocked,
6470 "must not read as resource-blocked on a lease it never asked for"
6471 );
6472 }
6473
6474 #[tokio::test]
6475 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
6476 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
6477 let tmp = tempfile::tempdir().expect("tempdir");
6478 let repo = tmp.path().join("repo");
6479 std::fs::create_dir_all(&repo).unwrap();
6480 init_repo(&repo);
6481
6482 let mut config = Config::default();
6483 config.merge.mode = MergeMode::Pr;
6484 config.graph.land = true;
6485 config.graph.land_approval = false;
6486
6487 let mut state = RunState::new(
6488 repo.clone(),
6489 "main".to_owned(),
6490 "deadbeef".to_owned(),
6491 "task".to_owned(),
6492 config,
6493 );
6494 state.candidates = vec![Candidate {
6495 index: 0,
6496 label: 'A',
6497 agent: "alpha".to_owned(),
6498 branch: "does-not-exist".to_owned(),
6499 worktree: repo.clone(),
6500 summary: String::new(),
6501 stat: String::new(),
6502 files: 0,
6503 commits: 0,
6504 empty: false,
6505 failed: None,
6506 duration_ms: 0,
6507 folded: false,
6508 }];
6509 state.tally = Some(Tally {
6510 first_choice: BTreeMap::from([('A', 1)]),
6511 borda: BTreeMap::new(),
6512 winner: 'A',
6513 rankings: 1,
6514 unanimous_initial: true,
6515 deliberated: false,
6516 changed_votes: 0,
6517 unanimous_final: true,
6518 tie_break: None,
6519 judges: 0,
6520 present: 0,
6521 quorum: 0,
6522 met_quorum: true,
6523 uncontested: Some("only candidate A produced a change".to_owned()),
6524 });
6525 state.reviews = vec![ReviewRound {
6526 round: 1,
6527 head: "deadbeef".to_owned(),
6528 verified_head: None,
6529 reviews: Vec::new(),
6530 e2e: Vec::new(),
6531 fix: None,
6532 blocking: 0,
6533 answered: 0,
6534 expected: 0,
6535 clean: true,
6536 verify_retried: false,
6537 e2e_deferred: false,
6538 e2e_defer_reason: None,
6539 progressed: false,
6540 vote_split: false,
6541 reconsideration: Vec::new(),
6542 verdict: None,
6543 }];
6544 state.gate = vec![CommandOutcome {
6545 command: "test".to_owned(),
6546 code: Some(0),
6547 output_tail: String::new(),
6548 duration_ms: 0,
6549 resource_blocked: false,
6550 }];
6551 state.gate_ran = true;
6552 state.status = RunStatus::Landing;
6556 state.merge = Some(MergeOutcome {
6557 mode: MergeMode::Pr,
6558 ok: true,
6559 detail: "https://example.invalid/x/y/pull/1".to_owned(),
6560 });
6561
6562 ask_test_home();
6566 let store = ask::Questions::open();
6567 let q = ask_open_question(&store, &state.id);
6568
6569 let mut runner = Runner {
6570 state,
6571 roles: ResolvedRoles {
6572 implementers: Vec::new(),
6573 judges: Vec::new(),
6574 reviewers: Vec::new(),
6575 fixer: None,
6576 conductor: conductor(),
6577 },
6578 sem: Arc::new(Semaphore::new(1)),
6579 pause: Pause::new(),
6580 interrupt: Pause::new(),
6581 };
6582
6583 runner.execute().await.expect("execute");
6588
6589 assert_eq!(
6590 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
6591 Some("https://example.invalid/x/y/pull/1"),
6592 "reentry must not push again or open a second pull request over the \
6593 one `land` is already watching"
6594 );
6595 assert_ne!(
6596 runner.state.status,
6597 RunStatus::Landing,
6598 "land could not actually reach the fake pull request, so it must \
6599 have given up rather than left the run silently parked forever"
6600 );
6601 assert_eq!(runner.state.status, RunStatus::Blocked);
6605 assert!(
6606 store.get(&q.id).unwrap().status.open(),
6607 "Blocked is still alive; settle_questions must have been a no-op here"
6608 );
6609 }
6610
6611 fn state_with_round(round: ReviewRound) -> RunState {
6612 let mut s = RunState::new(
6613 PathBuf::from("/repo"),
6614 "main".to_owned(),
6615 "abc1234".to_owned(),
6616 "add retries".to_owned(),
6617 Config::default(),
6618 );
6619 s.reviews = vec![round];
6620 s
6621 }
6622
6623 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
6624 crate::verdict::Finding {
6625 id: id.to_owned(),
6626 severity,
6627 file: None,
6628 line: None,
6629 title: title.to_owned(),
6630 detail: String::new(),
6631 }
6632 }
6633
6634 #[test]
6635 fn pr_body_names_open_findings_and_declined_ones() {
6636 let round = ReviewRound {
6637 round: 2,
6638 head: "deadbee".to_owned(),
6639 verified_head: None,
6640 reviews: vec![ReviewRecord {
6641 reviewer: 1,
6642 agent: "alpha".to_owned(),
6643 summary: String::new(),
6644 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
6645 vote: None,
6646 failed: None,
6647 duration_ms: 0,
6648 }],
6649 e2e: vec![CommandOutcome {
6650 command: "cargo test".to_owned(),
6651 code: Some(0),
6652 output_tail: String::new(),
6653 duration_ms: 0,
6654 resource_blocked: false,
6655 }],
6656 verify_retried: false,
6657 e2e_deferred: false,
6658 e2e_defer_reason: None,
6659 fix: Some(FixRecord {
6660 agent: "alpha".to_owned(),
6661 addressed: Vec::new(),
6662 rejected: vec![crate::verdict::Rejection {
6663 id: "R1-1-1".to_owned(),
6664 why: "not reachable from any caller".to_owned(),
6665 }],
6666 notes: String::new(),
6667 committed: true,
6668 failed: None,
6669 duration_ms: 0,
6670 continuation: None,
6671 }),
6672 blocking: 0,
6673 answered: 1,
6674 expected: 1,
6675 clean: false,
6676 progressed: true,
6677 vote_split: false,
6678 reconsideration: Vec::new(),
6679 verdict: None,
6680 };
6681 let state = state_with_round(round);
6682 let body = pr_body(&state, 'A');
6683
6684 assert!(body.contains("add retries"), "the task must still be there");
6685 assert!(body.contains("R2-1-1"), "{body}");
6686 assert!(body.contains("unused import"), "{body}");
6687 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
6688 assert!(
6689 body.contains("not reachable from any caller"),
6690 "the reason it was declined: {body}"
6691 );
6692 }
6693
6694 #[test]
6695 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
6696 let round = ReviewRound {
6697 round: 1,
6698 head: "deadbee".to_owned(),
6699 verified_head: None,
6700 reviews: vec![ReviewRecord {
6701 reviewer: 1,
6702 agent: "alpha".to_owned(),
6703 summary: String::new(),
6704 findings: Vec::new(),
6705 vote: None,
6706 failed: None,
6707 duration_ms: 0,
6708 }],
6709 e2e: Vec::new(),
6710 verify_retried: false,
6711 e2e_deferred: false,
6712 e2e_defer_reason: None,
6713 fix: None,
6714 blocking: 0,
6715 answered: 1,
6716 expected: 1,
6717 clean: true,
6718 progressed: false,
6719 vote_split: false,
6720 reconsideration: Vec::new(),
6721 verdict: None,
6722 };
6723 let state = state_with_round(round);
6724 let body = pr_body(&state, 'A');
6725 assert!(!body.contains("Open review findings"), "{body}");
6726 assert!(!body.contains("Declined"), "{body}");
6727 }
6728
6729 #[test]
6730 fn pr_body_titles_itself_from_the_task_not_run_or_candidate() {
6731 let state = RunState::new(
6732 PathBuf::from("/repo"),
6733 "main".to_owned(),
6734 "abc1234".to_owned(),
6735 "add retries".to_owned(),
6736 Config::default(),
6737 );
6738 let body = pr_body(&state, 'A');
6739 let title = body.lines().next().unwrap();
6740
6741 assert_eq!(
6742 title, "add retries",
6743 "the title must be the task, not run/candidate bookkeeping: {body}"
6744 );
6745 assert!(
6746 body.contains(&format!("magi:run/{}", state.id)),
6747 "the run id must still be recoverable from the footer: {body}"
6748 );
6749 assert!(
6750 body.contains("magi:candidate-a"),
6751 "the candidate must still be recoverable from the footer: {body}"
6752 );
6753 }
6754
6755 #[test]
6756 fn pr_body_never_titles_itself_off_a_blank_first_line() {
6757 let leading_blank = RunState::new(
6758 PathBuf::from("/repo"),
6759 "main".to_owned(),
6760 "abc1234".to_owned(),
6761 "\n\n \nadd retries\n\ndetails".to_owned(),
6762 Config::default(),
6763 );
6764 let body = pr_body(&leading_blank, 'A');
6765 assert_eq!(
6766 body.lines().next(),
6767 Some("add retries"),
6768 "a leading blank line must not become an empty title: {body}"
6769 );
6770
6771 let whitespace_only = RunState::new(
6772 PathBuf::from("/repo"),
6773 "main".to_owned(),
6774 "abc1234".to_owned(),
6775 " \n \n".to_owned(),
6776 Config::default(),
6777 );
6778 let body = pr_body(&whitespace_only, 'A');
6779 let title = body.lines().next().unwrap_or_default();
6780 assert!(
6781 !title.is_empty(),
6782 "a whitespace-only instruction must still fall back to a non-empty title: {body}"
6783 );
6784 }
6785
6786 #[test]
6787 fn manual_merge_command_matches_the_configured_style() {
6788 let repo = Path::new("/repo");
6789 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
6790
6791 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
6792 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
6793
6794 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
6795 assert_eq!(
6796 squash,
6797 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
6798 \"Merge magi run 0832 (candidate A)\""
6799 );
6800
6801 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
6802 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
6803 }
6804
6805 #[test]
6806 fn a_nudge_gets_a_quarter_of_the_budget() {
6807 assert_eq!(retry_budget(secs(1200), true), secs(300));
6809 assert_eq!(retry_budget(secs(3600), true), secs(900));
6810 }
6811
6812 #[test]
6813 fn a_resent_prompt_keeps_the_whole_budget() {
6814 assert_eq!(retry_budget(secs(1200), false), secs(1200));
6817 assert_eq!(retry_budget(secs(60), false), secs(60));
6818 }
6819
6820 #[test]
6821 fn the_floor_never_exceeds_the_original_budget() {
6822 assert_eq!(retry_budget(secs(60), true), secs(60));
6826 assert_eq!(retry_budget(secs(480), true), secs(120));
6827 assert_eq!(retry_budget(secs(0), true), secs(0));
6828 }
6829
6830 fn evidence(exit_code: Option<i32>) -> agent::CommandEvidence {
6831 agent::CommandEvidence {
6832 id: "item1".to_owned(),
6833 description: "cargo test".to_owned(),
6834 exit_code,
6835 result_summary: String::new(),
6836 source: "codex".to_owned(),
6837 }
6838 }
6839
6840 #[test]
6841 fn a_reply_with_no_commands_at_all_is_not_unconfirmed() {
6842 assert!(!has_unconfirmed_command(&[]));
6846 }
6847
6848 #[test]
6849 fn a_command_with_a_real_exit_code_is_confirmed_whatever_its_value() {
6850 assert!(!has_unconfirmed_command(&[evidence(Some(0))]));
6854 assert!(!has_unconfirmed_command(&[evidence(Some(1))]));
6855 assert!(!has_unconfirmed_command(&[
6856 evidence(Some(0)),
6857 evidence(Some(101))
6858 ]));
6859 }
6860
6861 #[test]
6862 fn one_command_with_no_readable_exit_code_is_enough_to_flag_the_reply() {
6863 assert!(has_unconfirmed_command(&[
6864 evidence(Some(0)),
6865 evidence(None)
6866 ]));
6867 }
6868}