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