1use std::collections::{BTreeMap, BTreeSet};
20use std::path::{Path, PathBuf};
21use std::sync::Arc;
22use std::sync::atomic::{AtomicBool, Ordering};
23use std::time::{Duration, Instant};
24
25use anyhow::{Context as _, Result, bail};
26use jiff::Timestamp;
27use tokio::sync::Semaphore;
28
29use crate::agent::{self, AgentOutput, Invocation, SeatState};
30use crate::ask;
31use crate::blind;
32use crate::bump;
33use crate::config::{
34 AgentSpec, Config, IncompleteReviewPolicy, LeakPolicy, MergeMode, MergeStyle, Prompts,
35 ResolvedRoles,
36};
37use crate::git;
38use crate::land;
39use crate::proc::Quiet as _;
40use crate::prompt::{
41 self, CandidateView, Lens, ReviewPatch, ReviewReconsiderCtx, ReviewSeatReport, Turn,
42};
43use crate::run::{
44 BaseSync, Candidate, CommandOutcome, DeliberationRound, DeliberationTurn, FixRecord, Judgement,
45 MergeOutcome, QuotaLoss, ReviewRecord, ReviewRevoteRecord, ReviewRound, RunState, RunStatus,
46 Tally, VoteRecord, tail, write_artifact,
47};
48use crate::verdict::{
49 self, FinalVote, FixReport, Position, Ranking, Review, ReviewRevote, ReviewVote, Severity,
50};
51
52const OUTPUT_TAIL: usize = 8_000;
54
55const EVENT_OUTPUT_TAIL: usize = 2_000;
58
59pub(crate) const STAGNANT_LIMIT: usize = 2;
73
74const BASE_SYNC_ROUNDS: usize = 4;
87
88#[derive(Clone)]
94struct SeatJob {
95 spec: AgentSpec,
96 seat: SeatState,
97 cwd: PathBuf,
98 prompt: String,
99 timeout: Duration,
100 allow_write: bool,
101 sessions: bool,
102 artifacts: PathBuf,
103 stem: String,
104}
105
106enum AgentOutcome {
118 Ok(AgentOutput),
120 Quota(AgentOutput),
122 Dropped(AgentOutput),
125 Failed(String),
127}
128
129#[derive(Debug, Clone, Default)]
146pub struct Pause(Arc<AtomicBool>);
147
148impl Pause {
149 #[must_use]
151 pub fn new() -> Self {
152 Self::default()
153 }
154
155 pub fn park(&self) {
157 self.0.store(true, Ordering::SeqCst);
158 }
159
160 #[must_use]
162 pub fn parked(&self) -> bool {
163 self.0.load(Ordering::SeqCst)
164 }
165}
166
167pub struct Runner {
169 pub state: RunState,
171 roles: ResolvedRoles,
172 sem: Arc<Semaphore>,
173 pause: Pause,
175}
176
177async fn resolve_base(repo: &Path, base_branch: &str, remote: &str) -> Result<String> {
197 let tracking = format!("{remote}/{base_branch}");
198 let fetched = git::fetch(repo, remote, base_branch).await;
199 if let Ok(out) = &fetched
200 && out.ok()
201 && git::rev_exists(repo, &tracking).await
202 {
203 return git::rev_parse(repo, &tracking).await;
204 }
205 let why = match &fetched {
206 Ok(out) if !out.ok() => out.stderr.lines().next().unwrap_or("").to_owned(),
207 Ok(_) => format!("{remote} has no {base_branch}"),
208 Err(e) => e.to_string(),
209 };
210 tracing::warn!(
211 "could not read {tracking} ({why}); branching off the local \
212 {base_branch} instead, which may be behind"
213 );
214 git::rev_parse(repo, base_branch).await.with_context(|| {
215 format!(
216 "cannot resolve `{base_branch}`; set [merge] base in magi.toml to a \
217 branch that exists"
218 )
219 })
220}
221
222impl Runner {
223 pub async fn start(repo: &Path, instruction: String, config: Config) -> Result<Self> {
225 let repo = git::toplevel(repo).await?;
226 let missing = agent::missing_programs(&config.agents);
227 if !missing.is_empty() {
228 bail!(
229 "these agent programs are not on PATH: {}. Fix the roster in \
230 magi.toml or install them.",
231 missing.join(", ")
232 );
233 }
234 let base_branch = match config.merge.base.clone() {
235 Some(b) => b,
236 None => git::current_branch(&repo)
237 .await?
238 .context("HEAD is detached; set [merge] base in magi.toml")?,
239 };
240 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
241 if !git::is_clean(&repo).await? {
245 tracing::warn!(
246 "{} has uncommitted changes; they are not part of this run, \
247 which branches off {base_branch} ({})",
248 repo.display(),
249 &base_commit[..base_commit.len().min(8)]
250 );
251 }
252 let roles = config.resolve_roles()?;
253 let max_parallel = config.graph.max_parallel.max(1);
254 let mut state = RunState::new(repo, base_branch, base_commit, instruction, config);
255 state.event("start", format!("run {} created", state.id));
256 state.save()?;
257 Ok(Self {
258 state,
259 roles,
260 sem: Arc::new(Semaphore::new(max_parallel)),
261 pause: Pause::new(),
262 })
263 }
264
265 pub async fn review(repo: &Path, branch: &str, config: Config) -> Result<Self> {
279 let repo = git::toplevel(repo).await?;
280 let missing = agent::missing_programs(&config.agents);
281 if !missing.is_empty() {
282 bail!(
283 "these agent programs are not on PATH: {}. Fix the roster in \
284 magi.toml or install them.",
285 missing.join(", ")
286 );
287 }
288 if !git::branch_exists(&repo, branch).await? {
289 bail!("no branch `{branch}` in {}", repo.display());
290 }
291 let base_branch = match config.merge.base.clone() {
292 Some(b) => b,
293 None => git::current_branch(&repo)
294 .await?
295 .context("HEAD is detached; set [merge] base in magi.toml")?,
296 };
297 if base_branch == branch {
298 bail!("`{branch}` is the base branch; there is nothing to review against");
299 }
300 let base_commit = resolve_base(&repo, &base_branch, &config.merge.remote).await?;
301
302 let roles = config.resolve_roles()?;
303 let max_parallel = config.graph.max_parallel.max(1);
304 let log = git::log_oneline(&repo, &base_commit, branch)
307 .await
308 .unwrap_or_default();
309 let instruction = format!(
310 "Review the work already on branch `{branch}`. There is no task \
311 statement: what the change claims to do is whatever its commits \
312 say.\n\n{}",
313 if log.trim().is_empty() {
314 "(no commit messages)"
315 } else {
316 log.trim()
317 }
318 );
319 let mut state = RunState::new(
320 repo.clone(),
321 base_branch,
322 base_commit.clone(),
323 instruction,
324 config,
325 );
326
327 let worktree = state.worktree_root().join("under-review");
330 if let Some(parent) = worktree.parent() {
331 tokio::fs::create_dir_all(parent).await.ok();
332 }
333 let path = worktree.to_string_lossy().to_string();
334 git::git(&repo, &["worktree", "add", &path, branch])
335 .await
336 .with_context(|| {
337 format!("checking out `{branch}` at {path} (is it checked out elsewhere?)")
338 })?;
339
340 let commits = git::commits_ahead(&worktree, &base_commit, "HEAD")
341 .await
342 .unwrap_or(0);
343 if commits == 0 {
344 bail!("`{branch}` has no commits beyond {}", short(&base_commit));
345 }
346 let files = git::changed_files(&worktree, &base_commit, "HEAD")
347 .await
348 .map(|f| f.len())
349 .unwrap_or(0);
350 let stat = git::diff_stat(&worktree, &base_commit, "HEAD")
351 .await
352 .unwrap_or_default();
353
354 state.candidates.push(Candidate {
355 index: 0,
356 label: 'A',
357 agent: "(existing branch)".to_owned(),
360 branch: branch.to_owned(),
361 worktree,
362 summary: String::new(),
363 stat,
364 files,
365 commits,
366 empty: false,
367 failed: None,
368 duration_ms: 0,
369 folded: false,
370 });
371 state.tally = Some(Tally {
372 first_choice: BTreeMap::from([('A', 0)]),
373 borda: BTreeMap::new(),
374 winner: 'A',
375 rankings: 0,
376 unanimous_initial: false,
377 deliberated: false,
378 changed_votes: 0,
379 unanimous_final: false,
380 tie_break: None,
381 judges: 0,
385 present: 0,
386 quorum: 0,
387 met_quorum: true,
388 uncontested: Some("review-only run: nothing competed".to_owned()),
389 });
390 state.status = RunStatus::Reviewing;
391 state.event(
392 "start",
393 format!(
394 "review-only run {} on `{branch}` ({files} files, {commits} commits)",
395 state.id
396 ),
397 );
398 state.save()?;
399 Ok(Self {
400 state,
401 roles,
402 sem: Arc::new(Semaphore::new(max_parallel)),
403 pause: Pause::new(),
404 })
405 }
406
407 pub fn resume(id: &str) -> Result<Self> {
409 let state = RunState::load(id)?;
410 let roles = state.config.resolve_roles()?;
411 let max_parallel = state.config.graph.max_parallel.max(1);
412 Ok(Self {
413 state,
414 roles,
415 sem: Arc::new(Semaphore::new(max_parallel)),
416 pause: Pause::new(),
417 })
418 }
419
420 pub async fn execute(&mut self) -> Result<()> {
422 self.state.parked = false;
427 if self.state.clear_active() {
434 self.state.save()?;
435 }
436 if self.state.status == RunStatus::Stalled {
449 if self.recover_stall().await? {
450 self.finish_after_tally().await?;
451 } else {
452 self.state.save()?;
454 }
455 return Ok(());
456 }
457 if self.state.status == RunStatus::Landing {
467 self.run_land().await?;
468 self.settle_questions();
473 return Ok(());
474 }
475 self.prep().await?;
476 if self.park_here()? {
477 return Ok(());
478 }
479 self.implement().await?;
480 if self.park_here()? {
481 return Ok(());
482 }
483 self.judge().await?;
484 if self.park_here()? {
485 return Ok(());
486 }
487 self.deliberate().await?;
488 if self.park_here()? {
489 return Ok(());
490 }
491 self.vote().await?;
492 if self.park_here()? {
493 return Ok(());
494 }
495 self.tally()?;
496 if self.state.status == RunStatus::Stalled {
501 self.state.save()?;
505 return Ok(());
506 }
507 self.finish_after_tally().await?;
508 Ok(())
509 }
510
511 fn park_here(&mut self) -> Result<bool> {
518 if !self.pause.parked() {
519 return Ok(false);
520 }
521 self.state.event(
522 "park",
523 format!(
524 "parked after `{}` — resume to carry on from here",
525 self.state.status.as_str()
526 ),
527 );
528 self.state.parked = true;
529 self.state.save()?;
530 Ok(true)
531 }
532
533 pub fn on_pause(&mut self, pause: Pause) {
535 self.pause = pause;
536 }
537
538 fn settle_questions(&mut self) {
556 if let Err(e) = ask::Questions::open().settle_run(&self.state.id, self.state.status) {
557 tracing::warn!("abandon questions for {}: {e:#}", self.state.id);
558 }
559 }
560
561 async fn finish_after_tally(&mut self) -> Result<()> {
564 self.fold_losers().await?;
565 self.sync_to_base().await?;
570 self.review_loop().await?;
571 self.sync_to_base().await?;
572 self.gate().await?;
573 self.merge().await?;
574 self.state.save()?;
575 Ok(())
576 }
577
578 async fn prep(&mut self) -> Result<()> {
581 if !self.state.candidates.is_empty() {
582 return Ok(());
583 }
584 self.state.status = RunStatus::Prep;
585 let repo = self.state.repo.clone();
586 let base = self.state.base_commit.clone();
587 let root = self.state.worktree_root();
588 let labels = blind::assign_labels(self.roles.implementers.len(), self.state.seed);
589
590 let hooks_dir = self.state.dir().join("hooks");
593 if self.state.config.blind.commit_msg_hook {
594 std::fs::create_dir_all(&hooks_dir)
595 .with_context(|| format!("create {}", hooks_dir.display()))?;
596 let script = blind::commit_msg_hook(&self.state.config.blind.strip_lines);
597 let path = hooks_dir.join("commit-msg");
598 std::fs::write(&path, script).with_context(|| format!("write {}", path.display()))?;
599 make_executable(&path)?;
600 git::acquire_worktree_config(&repo).await?;
608 self.state.enabled_worktree_config = true;
609 }
610
611 for (index, (spec, label)) in self
612 .roles
613 .implementers
614 .clone()
615 .into_iter()
616 .zip(labels)
617 .enumerate()
618 {
619 let branch = self.state.branch_for(label);
620 let worktree = root.join(format!("cand-{label}"));
621 git::worktree_add_branch(&repo, &worktree, &branch, &base).await?;
622 if self.state.config.blind.commit_msg_hook {
623 git::set_worktree_hooks_path(&worktree, &hooks_dir).await?;
624 }
625 git::local_exclude(&worktree, "/.magi/").await?;
626 self.state.candidates.push(Candidate {
627 index,
628 label,
629 agent: spec.id.clone(),
630 branch,
631 worktree,
632 summary: String::new(),
633 stat: String::new(),
634 files: 0,
635 commits: 0,
636 empty: false,
637 failed: None,
638 duration_ms: 0,
639 folded: false,
640 });
641 }
642
643 for j in 1..=self.roles.judges.len() {
644 let wt = root.join(format!("judge-{j}"));
645 if !wt.exists() {
646 git::worktree_add_detached(&repo, &wt, &base).await?;
647 }
648 }
649
650 let authors: Vec<&str> = self
655 .roles
656 .implementers
657 .iter()
658 .map(|a| a.id.as_str())
659 .collect();
660 let overlap: Vec<String> = self
661 .roles
662 .judges
663 .iter()
664 .enumerate()
665 .filter(|(_, j)| authors.contains(&j.id.as_str()))
666 .map(|(i, j)| format!("judge {} = {}", i + 1, j.id))
667 .collect();
668 if !overlap.is_empty() {
669 let note = format!(
670 "{} also authored a candidate; blind, but the panel is less \
671 independent than {} distinct agents would be",
672 overlap.join(", "),
673 self.roles.judges.len()
674 );
675 self.state.event("prep", note);
676 }
677
678 self.state.event(
679 "prep",
680 format!(
681 "{} candidates, {} judges, base {} ({})",
682 self.state.candidates.len(),
683 self.roles.judges.len(),
684 &self.state.base_commit[..7.min(self.state.base_commit.len())],
685 self.state.base_branch
686 ),
687 );
688 self.state.status = RunStatus::Implementing;
689 self.state.save()?;
690 Ok(())
691 }
692
693 async fn implement(&mut self) -> Result<()> {
696 let run_id = self.state.id.clone();
701 let prompts = self.state.config.prompts.clone();
702 let todo: Vec<usize> = self
703 .state
704 .candidates
705 .iter()
706 .enumerate()
707 .filter(|(_, c)| c.commits == 0 && c.failed.is_none() && !c.empty)
708 .map(|(i, _)| i)
709 .collect();
710 if todo.is_empty() {
711 return self.after_implement();
712 }
713 self.state.status = RunStatus::Implementing;
714
715 let language = self.state.config.graph.language.clone();
716 let timeout = Duration::from_secs(self.state.config.graph.timeout_implement);
717 let sessions = self.state.config.graph.sessions;
718 let artifacts = agent::artifacts_dir(&self.state.dir());
719
720 let mut jobs = Vec::new();
721 for &i in &todo {
722 let (index, label, worktree) = {
723 let c = &self.state.candidates[i];
724 (c.index, c.label, c.worktree.clone())
725 };
726 let spec = self.roles.implementers[index].clone();
727 let seat_key = format!("impl-{label}");
728 let seat = self.seat(&seat_key, &spec.id);
729 let instruction = self.state.instruction.clone();
730 jobs.push(SeatJob {
731 spec,
732 seat,
733 prompt: prompt::implement(&instruction, &worktree.to_string_lossy(), &language),
734 cwd: worktree,
735 timeout,
736 allow_write: true,
737 sessions,
738 artifacts: artifacts.clone(),
739 stem: format!("impl-{label}"),
740 });
741 }
742
743 self.state.event(
744 "implement",
745 format!("{} candidates in parallel", jobs.len()),
746 );
747 let sent = jobs.clone();
750 let cache = self.state.config.cache_dir();
751 let ctx = WaveCtx {
752 run: &run_id,
753 node: "implement",
754 prompts: &prompts,
755 cache: cache.as_deref(),
756 };
757 let mut results = wave(jobs, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
758 self.resume_undelivered(&mut results, &sent, &prompts, &run_id)
759 .await;
760
761 for (&i, (_wi, seat, out)) in todo.iter().zip(results) {
762 let seat_key = seat.key.clone();
763 self.state.seats.insert(seat.key.clone(), seat);
764 let label = self.state.candidates[i].label;
765 let worktree = self.state.candidates[i].worktree.clone();
766 let base = self.state.base_commit.clone();
767
768 let (summary, duration, failed) = match out {
769 AgentOutcome::Ok(o) => {
770 let text = verdict::section(&o.text, "summary").unwrap_or(o.text.clone());
771 let failed = (!o.usable()).then(|| {
772 if o.timed_out {
773 "agent timed out".to_owned()
774 } else {
775 format!("agent exited with {:?}", o.exit_code)
776 }
777 });
778 (text, o.duration_ms, failed)
779 }
780 AgentOutcome::Dropped(o) => {
786 let why = o
787 .dropped
788 .as_ref()
789 .map(|d| d.why.as_str())
790 .unwrap_or("the CLI ended the stream without delivering its answer");
791 (
792 String::new(),
793 o.duration_ms,
794 Some(format!("the CLI dropped the stream ({why})")),
795 )
796 }
797 AgentOutcome::Quota(o) => {
798 self.state.quota.push(QuotaLoss {
799 seat: seat_key,
800 node: "implement".to_owned(),
801 at: Timestamp::now(),
802 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
803 });
804 (
805 String::new(),
806 o.duration_ms,
807 Some("rate limited (quota); produced no change".to_owned()),
808 )
809 }
810 AgentOutcome::Failed(e) => (String::new(), 0, Some(e)),
811 };
812
813 let rescued = git::commit_all(
816 &worktree,
817 &format!("magi: candidate {label} (uncommitted work)"),
818 )
819 .await
820 .unwrap_or(false);
821 let commits = git::commits_ahead(&worktree, &base, "HEAD")
822 .await
823 .unwrap_or(0);
824 let patch = git::diff(&worktree, &base, "HEAD")
825 .await
826 .unwrap_or_default();
827 let stat = git::diff_stat(&worktree, &base, "HEAD")
828 .await
829 .unwrap_or_default();
830 let files = git::changed_files(&worktree, &base, "HEAD")
831 .await
832 .map(|f| f.len())
833 .unwrap_or(0);
834 write_artifact(&self.state, &format!("cand-{label}.patch"), &patch)?;
835
836 let c = &mut self.state.candidates[i];
837 c.summary = blind::sanitize_prose(&summary, &self.state.config.blind);
838 c.stat = stat;
839 c.files = files;
840 c.commits = commits;
841 c.duration_ms = duration;
842 c.empty = commits == 0 || patch.trim().is_empty();
843 c.failed = match failed {
846 Some(_) if c.empty => failed,
847 _ => None,
848 };
849 let note = match (&c.failed, c.empty, rescued) {
850 (Some(e), _, _) => format!("candidate {label}: {e}"),
851 (None, true, _) => format!("candidate {label}: no change produced"),
852 (None, false, true) => {
853 format!(
854 "candidate {label}: {files} files, {commits} commits (rescued an uncommitted tree)"
855 )
856 }
857 (None, false, false) => {
858 format!("candidate {label}: {files} files, {commits} commits")
859 }
860 };
861 self.state.event("implement", note);
862 self.state.save()?;
863 }
864
865 self.after_implement()
866 }
867
868 async fn resume_undelivered(
896 &mut self,
897 results: &mut [(usize, SeatState, AgentOutcome)],
898 sent: &[SeatJob],
899 prompts: &Prompts,
900 run_id: &str,
901 ) {
902 for (wi, seat, out) in results.iter_mut() {
903 let Some(dropped) = (match &*out {
904 AgentOutcome::Dropped(o) => o.dropped.clone(),
905 _ => None,
906 }) else {
907 continue;
908 };
909 let Some(job) = sent.get(*wi) else { continue };
910 if !git::is_clean(&job.cwd).await.unwrap_or(true) {
912 self.state.event(
913 "implement",
914 format!(
915 "{}: the CLI dropped the stream after {} output tokens ({}), but the \
916 work is in the tree",
917 seat.key, dropped.output_tokens, dropped.why
918 ),
919 );
920 continue;
921 }
922 if !has_context(&job.spec, seat, job.sessions) {
930 self.state.event(
931 "implement",
932 format!(
933 "{}: the CLI dropped the stream after {} output tokens ({}), but there \
934 is no session left to resume",
935 seat.key, dropped.output_tokens, dropped.why
936 ),
937 );
938 continue;
939 }
940 self.state.event(
941 "implement",
942 format!(
943 "{}: the CLI dropped the stream after {} output tokens ({}); resuming the \
944 conversation",
945 seat.key, dropped.output_tokens, dropped.why
946 ),
947 );
948 let mut retry = job.clone();
949 retry.seat = seat.clone();
950 retry.prompt = prompt::resume_after_drop(&dropped.why);
951 retry.timeout = retry_budget(job.timeout, true);
952 retry.stem = format!("{}-resume", job.stem);
953 let cache = self.state.config.cache_dir();
954 let ctx = WaveCtx {
955 run: run_id,
956 node: "implement",
957 prompts,
958 cache: cache.as_deref(),
959 };
960 let (resumed_seat, resumed) =
961 run_one(retry, Arc::clone(&self.sem), &ctx, &mut self.state, 1).await;
962 *seat = resumed_seat;
963 *out = resumed;
964 }
965 }
966
967 fn after_implement(&mut self) -> Result<()> {
968 if self.state.leaks.is_empty() {
970 let cfg = self.state.config.blind.clone();
971 let mut leaks = Vec::new();
972 for c in &self.state.candidates {
973 let Some(patch) =
974 crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
975 else {
976 continue;
977 };
978 leaks.extend(blind::scan(
979 &format!("candidate {} patch", c.label),
980 &patch,
981 &cfg.vendor_tokens,
982 ));
983 }
984 if !leaks.is_empty() {
985 let summary = leaks
986 .iter()
987 .map(|l| format!("{}×{} in {}", l.token, l.count, l.site))
988 .collect::<Vec<_>>()
989 .join(", ");
990 match cfg.on_leak {
991 LeakPolicy::Fail => {
992 self.state.status = RunStatus::Failed;
993 self.state
994 .event("blind", format!("vendor text in a patch: {summary}"));
995 self.state.leaks = leaks;
996 self.state.save()?;
997 self.settle_questions();
998 bail!(
999 "blind.on_leak = \"fail\" and vendor text reached a \
1000 judged patch: {summary}"
1001 );
1002 }
1003 LeakPolicy::Redact => self.state.event(
1004 "blind",
1005 format!("redacting vendor text for judging: {summary}"),
1006 ),
1007 LeakPolicy::Warn => self.state.event(
1008 "blind",
1009 format!("vendor text present in a judged patch (shown as-is): {summary}"),
1010 ),
1011 }
1012 self.state.leaks = leaks;
1013 }
1014 }
1015
1016 if self.state.viable().is_empty() {
1017 self.state.status = RunStatus::Failed;
1018 self.state.save()?;
1019 self.settle_questions();
1020 bail!("no candidate produced a change; nothing to judge");
1021 }
1022 self.state.status = RunStatus::Judging;
1023 self.state.save()?;
1024 Ok(())
1025 }
1026
1027 async fn judge(&mut self) -> Result<()> {
1030 let run_id = self.state.id.clone();
1035 let prompts = self.state.config.prompts.clone();
1036 if !self.state.judgements.is_empty() || self.state.judge_skipped {
1037 return Ok(());
1038 }
1039 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
1040 if viable.len() == 1 {
1041 self.state.judge_skipped = true;
1048 self.state.event(
1049 "judge",
1050 format!(
1051 "only candidate {} produced a change; judging skipped",
1052 viable[0].label
1053 ),
1054 );
1055 self.state.save()?;
1056 return Ok(());
1057 }
1058 self.state.status = RunStatus::Judging;
1059
1060 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
1061 let language = self.state.config.graph.language.clone();
1062 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
1063 let sessions = self.state.config.graph.sessions;
1064 let artifacts = agent::artifacts_dir(&self.state.dir());
1065 let root = self.state.worktree_root();
1066 let base_short = short(&self.state.base_commit);
1067
1068 let mut jobs = Vec::new();
1069 let mut orders = Vec::new();
1070 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
1071 let order = blind::presentation_order(viable.len(), j, self.state.seed);
1072 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
1073 orders.push(order.iter().map(|&k| viable[k].index).collect::<Vec<_>>());
1074 let seat_key = format!("judge-{}", j + 1);
1075 let seat = self.seat(&seat_key, &spec.id);
1076 jobs.push(SeatJob {
1077 prompt: prompt::judge(
1078 &self.state.instruction,
1079 &views,
1080 self.roles.judges.len(),
1081 &base_short,
1082 &language,
1083 ),
1084 spec,
1085 seat,
1086 cwd: root.join(format!("judge-{}", j + 1)),
1087 timeout,
1088 allow_write: false,
1089 sessions,
1090 artifacts: artifacts.clone(),
1091 stem: format!("judge-{}", j + 1),
1092 });
1093 }
1094
1095 self.state.event(
1096 "judge",
1097 format!(
1098 "{} judges ranking {} candidates blind",
1099 jobs.len(),
1100 viable.len()
1101 ),
1102 );
1103 let labels_for_check = labels.clone();
1104 let mut quota_losses = Vec::new();
1105 let cache = self.state.config.cache_dir();
1106 let ctx = WaveCtx {
1107 run: &run_id,
1108 node: "judge",
1109 prompts: &prompts,
1110 cache: cache.as_deref(),
1111 };
1112 let results = ask_json_wave::<Ranking>(
1113 jobs,
1114 Arc::clone(&self.sem),
1115 self.state.config.graph.retries,
1116 &ctx,
1117 &mut quota_losses,
1118 &mut self.state,
1119 &move |r: &Ranking| r.validate(&labels_for_check),
1120 )
1121 .await;
1122 self.state.quota.extend(quota_losses);
1123
1124 for (j, (seat, res)) in results.into_iter().enumerate() {
1125 let agent_id = seat.agent.clone();
1126 self.state.seats.insert(seat.key.clone(), seat);
1127 let mut record = Judgement {
1128 judge: j + 1,
1129 seat: format!("judge-{}", j + 1),
1130 agent: agent_id,
1131 ranking: Vec::new(),
1132 reasons: BTreeMap::new(),
1133 confidence: None,
1134 order: orders[j].clone(),
1135 failed: None,
1136 duration_ms: 0,
1137 };
1138 match res {
1139 Ok((ranking, out)) => {
1140 record.ranking = ranking.normalized();
1141 record.reasons = ranking.reasons;
1142 record.confidence = ranking.confidence;
1143 record.duration_ms = out.duration_ms;
1144 self.state.event(
1145 "judge",
1146 format!(
1147 "judge {} ranked {}",
1148 j + 1,
1149 record.ranking.iter().collect::<String>()
1150 ),
1151 );
1152 }
1153 Err(e) => {
1154 record.failed = Some(e.to_string());
1155 self.state
1156 .event("judge", format!("judge {} produced no ranking: {e}", j + 1));
1157 }
1158 }
1159 self.state.judgements.push(record);
1160 self.state.save()?;
1161 }
1162 Ok(())
1163 }
1164
1165 async fn deliberate(&mut self) -> Result<()> {
1168 let run_id = self.state.id.clone();
1173 let prompts = self.state.config.prompts.clone();
1174 if !self.state.deliberation.is_empty() {
1175 return Ok(());
1176 }
1177 let tops: Vec<char> = self
1178 .state
1179 .judgements
1180 .iter()
1181 .filter_map(|j| j.ranking.first().copied())
1182 .collect();
1183 let rounds = self.state.config.graph.deliberate_rounds;
1184 if tops.len() < 2 || tops.iter().all(|t| *t == tops[0]) || rounds == 0 {
1185 if tops.len() >= 2 && tops.iter().all(|t| *t == tops[0]) {
1186 self.state.event(
1187 "deliberate",
1188 format!("judges agreed on {} outright; no deliberation", tops[0]),
1189 );
1190 }
1191 self.state.status = RunStatus::Voting;
1192 self.state.save()?;
1193 return Ok(());
1194 }
1195
1196 self.state.status = RunStatus::Deliberating;
1197 self.state.event(
1198 "deliberate",
1199 format!(
1200 "split: first choices were {} — opening {rounds} round(s)",
1201 tops.iter().collect::<String>()
1202 ),
1203 );
1204
1205 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
1206 let language = self.state.config.graph.language.clone();
1207 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
1208 let sessions = self.state.config.graph.sessions;
1209 let artifacts = agent::artifacts_dir(&self.state.dir());
1210 let root = self.state.worktree_root();
1211 let base_short = short(&self.state.base_commit);
1212
1213 for round in 1..=rounds {
1217 let mut turns: Vec<DeliberationTurn> = Vec::new();
1218 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
1219 if self.state.judgements[j].failed.is_some() {
1220 continue;
1221 }
1222 let seat_key = format!("judge-{}", j + 1);
1223 let mut seat = self.seat(&seat_key, &spec.id);
1224 let transcript = self.transcript(&turns, j);
1225 let context = if has_context(&spec, &seat, sessions) {
1226 None
1227 } else {
1228 Some(self.candidate_block(&viable, &base_short))
1229 };
1230 let text = prompt::deliberate(
1231 &self.state.instruction,
1232 context.as_deref(),
1233 &transcript,
1234 round,
1235 rounds,
1236 &language,
1237 );
1238 let job = SeatJob {
1239 spec,
1240 seat: seat.clone(),
1241 prompt: text,
1242 cwd: root.join(format!("judge-{}", j + 1)),
1243 timeout,
1244 allow_write: false,
1245 sessions,
1246 artifacts: artifacts.clone(),
1247 stem: format!("delib-{round}-judge-{}", j + 1),
1248 };
1249 let cache = self.state.config.cache_dir();
1250 let ctx = WaveCtx {
1251 run: &run_id,
1252 node: "deliberate",
1253 prompts: &prompts,
1254 cache: cache.as_deref(),
1255 };
1256 let (updated, out) =
1257 run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
1258 seat = updated;
1259 let agent_id = seat.agent.clone();
1260 let seat_key = seat.key.clone();
1261 self.state.seats.insert(seat.key.clone(), seat);
1262 let body = match out {
1263 AgentOutcome::Ok(o) => verdict::section(&o.text, "position").unwrap_or(o.text),
1264 AgentOutcome::Dropped(o) => {
1268 let why =
1269 o.dropped.as_ref().map(|d| d.why.as_str()).unwrap_or(
1270 "the CLI ended the stream without delivering its answer",
1271 );
1272 self.state.event(
1273 "deliberate",
1274 format!(
1275 "judge {} skipped: the CLI dropped the stream ({why})",
1276 j + 1
1277 ),
1278 );
1279 continue;
1280 }
1281 AgentOutcome::Quota(o) => {
1282 self.state.quota.push(QuotaLoss {
1283 seat: seat_key,
1284 node: "deliberate".to_owned(),
1285 at: Timestamp::now(),
1286 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
1287 });
1288 self.state.event(
1289 "deliberate",
1290 format!("judge {} skipped: rate limited (quota)", j + 1),
1291 );
1292 continue;
1293 }
1294 AgentOutcome::Failed(e) => {
1295 self.state
1296 .event("deliberate", format!("judge {} skipped: {e}", j + 1));
1297 continue;
1298 }
1299 };
1300 let tentative = verdict::extract_json::<Position>(&body)
1301 .ok()
1302 .and_then(|p| p.tentative)
1303 .and_then(|s| s.trim().chars().next())
1304 .map(|c| c.to_ascii_uppercase());
1305 self.state.event(
1306 "deliberate",
1307 format!(
1308 "round {round}: judge {} now favours {}",
1309 j + 1,
1310 tentative.map_or("—".to_owned(), |c| c.to_string())
1311 ),
1312 );
1313 turns.push(DeliberationTurn {
1314 judge: j + 1,
1315 agent: agent_id,
1316 body: blind::sanitize_prose(&body, &self.state.config.blind),
1317 tentative,
1318 });
1319 }
1320 self.state
1321 .deliberation
1322 .push(DeliberationRound { round, turns });
1323 self.state.save()?;
1324 }
1325
1326 self.state.status = RunStatus::Voting;
1327 self.state.save()?;
1328 Ok(())
1329 }
1330
1331 async fn vote(&mut self) -> Result<()> {
1334 let run_id = self.state.id.clone();
1339 let prompts = self.state.config.prompts.clone();
1340 if !self.state.votes.is_empty() {
1341 return Ok(());
1342 }
1343 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
1344 if viable.len() == 1 {
1345 return Ok(());
1346 }
1347 self.state.status = RunStatus::Voting;
1348
1349 let language = self.state.config.graph.language.clone();
1350 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
1351 let sessions = self.state.config.graph.sessions;
1352 let artifacts = agent::artifacts_dir(&self.state.dir());
1353 let root = self.state.worktree_root();
1354 let base_short = short(&self.state.base_commit);
1355 let candidates: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
1356
1357 let mut jobs = Vec::new();
1358 let mut seats_at = Vec::new();
1359 for (j, spec) in self.roles.judges.clone().into_iter().enumerate() {
1360 if self
1361 .state
1362 .judgements
1363 .get(j)
1364 .is_some_and(|r| r.failed.is_some())
1365 {
1366 continue;
1367 }
1368 let seat_key = format!("judge-{}", j + 1);
1369 let seat = self.seat(&seat_key, &spec.id);
1370 let mut text = prompt::final_vote(&viable, &language);
1371 if !has_context(&spec, &seat, sessions) {
1372 text = format!(
1373 "{}\n\n# Candidates\n\n{}",
1374 text,
1375 self.candidate_block(&candidates, &base_short)
1376 );
1377 }
1378 jobs.push(SeatJob {
1379 spec,
1380 seat,
1381 prompt: text,
1382 cwd: root.join(format!("judge-{}", j + 1)),
1383 timeout,
1384 allow_write: false,
1385 sessions,
1386 artifacts: artifacts.clone(),
1387 stem: format!("vote-judge-{}", j + 1),
1388 });
1389 seats_at.push(j);
1390 }
1391
1392 self.state.event(
1393 "vote",
1394 format!(
1395 "collecting {} final votes one by one, privately",
1396 jobs.len()
1397 ),
1398 );
1399 let allowed = viable.clone();
1400 let mut quota_losses = Vec::new();
1401 let cache = self.state.config.cache_dir();
1402 let ctx = WaveCtx {
1403 run: &run_id,
1404 node: "vote",
1405 prompts: &prompts,
1406 cache: cache.as_deref(),
1407 };
1408 let results = ask_json_wave::<FinalVote>(
1409 jobs,
1410 Arc::clone(&self.sem),
1411 self.state.config.graph.retries,
1412 &ctx,
1413 &mut quota_losses,
1414 &mut self.state,
1415 &move |v: &FinalVote| match v.label() {
1416 Some(c) if allowed.contains(&c) => Ok(()),
1417 other => bail!("vote {other:?} is not one of {allowed:?}"),
1418 },
1419 )
1420 .await;
1421 self.state.quota.extend(quota_losses);
1422
1423 for (&j, (seat, res)) in seats_at.iter().zip(results) {
1424 let agent_id = seat.agent.clone();
1425 self.state.seats.insert(seat.key.clone(), seat);
1426 let initial = self
1427 .state
1428 .judgements
1429 .get(j)
1430 .and_then(|r| r.ranking.first().copied());
1431 let mut record = VoteRecord {
1432 judge: j + 1,
1433 agent: agent_id,
1434 vote: None,
1435 reason: String::new(),
1436 changed: false,
1437 };
1438 match res {
1439 Ok((v, _)) => {
1440 record.vote = v.label();
1441 record.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
1442 record.changed = matches!((record.vote, initial), (Some(a), Some(b)) if a != b);
1443 self.state.event(
1444 "vote",
1445 format!(
1446 "judge {} voted {}{}",
1447 j + 1,
1448 record.vote.unwrap_or('?'),
1449 if record.changed { " (changed)" } else { "" }
1450 ),
1451 );
1452 }
1453 Err(e) => {
1454 self.state
1455 .event("vote", format!("judge {} cast no vote: {e}", j + 1));
1456 }
1457 }
1458 self.state.votes.push(record);
1459 self.state.save()?;
1460 }
1461 Ok(())
1462 }
1463
1464 fn tally(&mut self) -> Result<()> {
1467 if self.state.tally.is_some() {
1468 return Ok(());
1469 }
1470 let viable: Vec<char> = self.state.viable().into_iter().map(|c| c.label).collect();
1471 let tops: Vec<char> = self
1472 .state
1473 .judgements
1474 .iter()
1475 .filter_map(|j| j.ranking.first().copied())
1476 .collect();
1477 let unanimous_initial = tops.len() > 1 && tops.iter().all(|t| *t == tops[0]);
1478
1479 let mut first_choice: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
1482 let mut cast: Vec<char> = Vec::new();
1483 for (i, j) in self.state.judgements.iter().enumerate() {
1484 let vote = self
1485 .state
1486 .votes
1487 .iter()
1488 .find(|v| v.judge == i + 1)
1489 .and_then(|v| v.vote)
1490 .or_else(|| j.ranking.first().copied());
1491 if let Some(v) = vote {
1492 *first_choice.entry(v).or_insert(0) += 1;
1493 cast.push(v);
1494 }
1495 }
1496
1497 let mut borda: BTreeMap<char, usize> = viable.iter().map(|l| (*l, 0)).collect();
1498 for j in &self.state.judgements {
1499 let n = j.ranking.len();
1500 for (pos, label) in j.ranking.iter().enumerate() {
1501 *borda.entry(*label).or_insert(0) += n.saturating_sub(pos + 1);
1502 }
1503 }
1504
1505 let best = first_choice.values().copied().max().unwrap_or(0);
1506 let mut leaders: Vec<char> = first_choice
1507 .iter()
1508 .filter(|(_, v)| **v == best)
1509 .map(|(k, _)| *k)
1510 .collect();
1511 let mut tie_break = None;
1512 if leaders.len() > 1 {
1513 let top_borda = leaders.iter().map(|l| borda[l]).max().unwrap_or(0);
1514 let borda_leaders: Vec<char> = leaders
1515 .iter()
1516 .copied()
1517 .filter(|l| borda[l] == top_borda)
1518 .collect();
1519 tie_break = Some(if borda_leaders.len() == 1 {
1520 format!(
1521 "{} way tie on first-choice votes, broken by Borda points from the initial rankings",
1522 leaders.len()
1523 )
1524 } else {
1525 format!(
1526 "{} way tie on both first-choice votes and Borda points, broken by label order",
1527 leaders.len()
1528 )
1529 });
1530 leaders = borda_leaders;
1531 leaders.sort_unstable();
1532 }
1533 let winner = *leaders
1534 .first()
1535 .or(viable.first())
1536 .context("no candidate to declare a winner from")?;
1537
1538 let changed_votes = self.state.votes.iter().filter(|v| v.changed).count();
1539 let unanimous_final = !cast.is_empty() && cast.iter().all(|c| *c == cast[0]);
1540 let deliberated = !self.state.deliberation.is_empty();
1541
1542 let quota_seats: std::collections::BTreeSet<&str> =
1546 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
1547 let mut present = 0usize;
1548 for (i, j) in self.state.judgements.iter().enumerate() {
1549 if quota_seats.contains(j.seat.as_str()) {
1550 continue;
1551 }
1552 let ranked = !j.ranking.is_empty() && j.failed.is_none();
1553 let voted = self
1554 .state
1555 .votes
1556 .iter()
1557 .any(|v| v.judge == i + 1 && v.vote.is_some());
1558 if ranked || voted {
1559 present += 1;
1560 }
1561 }
1562 let needs_quorum = viable.len() > 1;
1568 let judges_total = if needs_quorum {
1569 self.roles.judges.len()
1570 } else {
1571 0
1572 };
1573 let quorum = if needs_quorum {
1574 judges_total / 2 + 1
1575 } else {
1576 0
1577 };
1578 let met_quorum = !needs_quorum || present >= quorum;
1579 let uncontested = (!needs_quorum).then(|| {
1580 format!("only one candidate ({winner}) produced a usable change; no panel was asked")
1581 });
1582
1583 self.state.event(
1584 "tally",
1585 match &uncontested {
1586 Some(reason) => format!("winner {winner} — {reason}"),
1587 None => format!(
1588 "winner {winner} — votes {} | initial {} | {} changed | \
1589 {present}/{judges_total} judges{}",
1590 first_choice
1591 .iter()
1592 .map(|(k, v)| format!("{k}:{v}"))
1593 .collect::<Vec<_>>()
1594 .join(" "),
1595 if unanimous_initial {
1596 "unanimous"
1597 } else {
1598 "split"
1599 },
1600 changed_votes,
1601 if met_quorum {
1602 String::new()
1603 } else {
1604 format!(" — below quorum ({quorum} required)")
1605 },
1606 ),
1607 },
1608 );
1609 if !met_quorum {
1610 self.state.event(
1611 "stall",
1612 format!(
1613 "verdict rests on {present} of {judges_total} judges (quorum {quorum}); \
1614 the run stops here, resumable"
1615 ),
1616 );
1617 }
1618 self.state.tally = Some(Tally {
1619 first_choice,
1620 borda,
1621 winner,
1622 rankings: tops.len(),
1623 unanimous_initial,
1624 deliberated,
1625 changed_votes,
1626 unanimous_final,
1627 tie_break,
1628 judges: judges_total,
1629 present,
1630 quorum,
1631 met_quorum,
1632 uncontested,
1633 });
1634 self.state.status = if met_quorum {
1635 RunStatus::Reviewing
1636 } else {
1637 RunStatus::Stalled
1638 };
1639 self.state.save()?;
1640 Ok(())
1641 }
1642
1643 #[allow(clippy::too_many_lines)]
1664 async fn recover_stall(&mut self) -> Result<bool> {
1665 let run_id = self.state.id.clone();
1670 let prompts = self.state.config.prompts.clone();
1671 let quota_seats: BTreeSet<&str> =
1676 self.state.quota.iter().map(|q| q.seat.as_str()).collect();
1677 let absent: Vec<String> = self
1678 .state
1679 .judgements
1680 .iter()
1681 .filter(|j| quota_seats.contains(j.seat.as_str()) || j.failed.is_some())
1682 .map(|j| j.seat.clone())
1683 .collect();
1684 if absent.is_empty() {
1685 return Ok(false);
1686 }
1687 let viable: Vec<Candidate> = self.state.viable().into_iter().cloned().collect();
1688 if viable.len() <= 1 {
1689 return Ok(false);
1690 }
1691 let labels: Vec<char> = viable.iter().map(|c| c.label).collect();
1692 let language = self.state.config.graph.language.clone();
1693 let timeout = Duration::from_secs(self.state.config.graph.timeout_judge);
1694 let sessions = self.state.config.graph.sessions;
1695 let artifacts = agent::artifacts_dir(&self.state.dir());
1696 let root = self.state.worktree_root();
1697 let base_short = short(&self.state.base_commit);
1698 let candidates: Vec<Candidate> = viable.clone();
1699
1700 let mut positions: Vec<usize> = absent
1702 .iter()
1703 .filter_map(|k| self.state.judgements.iter().position(|r| &r.seat == k))
1704 .collect();
1705 if positions.is_empty() {
1706 return Ok(false);
1707 }
1708 positions.sort_unstable();
1709 positions.dedup();
1710
1711 let mut judge_jobs = Vec::new();
1713 for &j in &positions {
1714 let order = blind::presentation_order(viable.len(), j, self.state.seed);
1715 let views: Vec<CandidateView> = order.iter().map(|&k| self.view(&viable[k])).collect();
1716 let seat_key = format!("judge-{}", j + 1);
1717 let spec = self.roles.judges[j].clone();
1718 let seat = self.seat(&seat_key, &spec.id);
1719 judge_jobs.push(SeatJob {
1720 spec,
1721 seat,
1722 prompt: prompt::judge(
1723 &self.state.instruction,
1724 &views,
1725 self.roles.judges.len(),
1726 &base_short,
1727 &language,
1728 ),
1729 cwd: root.join(seat_key),
1730 timeout,
1731 allow_write: false,
1732 sessions,
1733 artifacts: artifacts.clone(),
1734 stem: format!("judge-{}-recover", j + 1),
1735 });
1736 }
1737
1738 let labels_for_check = labels.clone();
1739 let mut judge_losses = Vec::new();
1740 let retries = self.state.config.graph.retries;
1741 let cache = self.state.config.cache_dir();
1742 let ctx = WaveCtx {
1743 run: &run_id,
1744 node: "judge",
1745 prompts: &prompts,
1746 cache: cache.as_deref(),
1747 };
1748 let results = ask_json_wave::<Ranking>(
1749 judge_jobs,
1750 Arc::clone(&self.sem),
1751 retries,
1752 &ctx,
1753 &mut judge_losses,
1754 &mut self.state,
1755 &move |r: &Ranking| r.validate(&labels_for_check),
1756 )
1757 .await;
1758
1759 let mut recovered: BTreeSet<usize> = BTreeSet::new();
1761 for (&j, (seat, res)) in positions.iter().zip(results) {
1762 self.state.seats.insert(seat.key.clone(), seat);
1763 let record = &mut self.state.judgements[j];
1764 match res {
1765 Ok((ranking, out)) => {
1766 record.ranking = ranking.normalized();
1767 record.reasons = ranking.reasons;
1768 record.confidence = ranking.confidence;
1769 record.failed = None;
1770 record.duration_ms = out.duration_ms;
1771 recovered.insert(j);
1772 self.state.event(
1773 "recover",
1774 format!("judge {} ranked again after the limit", j + 1),
1775 );
1776 }
1777 Err(e) => {
1778 self.state
1779 .event("recover", format!("judge {} still cannot rank: {e}", j + 1));
1780 }
1781 }
1782 }
1783
1784 let mut vote_jobs = Vec::new();
1786 let mut vote_pos: Vec<usize> = Vec::new();
1787 for &j in &recovered {
1788 let seat_key = format!("judge-{}", j + 1);
1789 let spec = self.roles.judges[j].clone();
1790 let seat = self.seat(&seat_key, &spec.id);
1791 let mut text = prompt::final_vote(&labels, &language);
1792 if !has_context(&spec, &seat, sessions) {
1793 text = format!(
1794 "{}\n\n# Candidates\n\n{}",
1795 text,
1796 self.candidate_block(&candidates, &base_short)
1797 );
1798 }
1799 vote_jobs.push(SeatJob {
1800 spec,
1801 seat,
1802 prompt: text,
1803 cwd: root.join(seat_key),
1804 timeout,
1805 allow_write: false,
1806 sessions,
1807 artifacts: artifacts.clone(),
1808 stem: format!("vote-judge-{}-recover", j + 1),
1809 });
1810 vote_pos.push(j);
1811 }
1812 let allowed = labels.clone();
1813 let mut vote_losses = Vec::new();
1814 let vote_retries = self.state.config.graph.retries;
1815 let vote_cache = self.state.config.cache_dir();
1816 let ctx = WaveCtx {
1817 run: &run_id,
1818 node: "vote",
1819 prompts: &prompts,
1820 cache: vote_cache.as_deref(),
1821 };
1822 let votes = ask_json_wave::<FinalVote>(
1823 vote_jobs,
1824 Arc::clone(&self.sem),
1825 vote_retries,
1826 &ctx,
1827 &mut vote_losses,
1828 &mut self.state,
1829 &move |v: &FinalVote| match v.label() {
1830 Some(c) if allowed.contains(&c) => Ok(()),
1831 other => bail!("vote {other:?} is not one of {allowed:?}"),
1832 },
1833 )
1834 .await;
1835 for (&j, (seat, res)) in vote_pos.iter().zip(votes) {
1836 let agent_id = seat.agent.clone();
1837 self.state.seats.insert(seat.key.clone(), seat);
1838 match res {
1839 Ok((v, _)) => {
1840 if let Some(rec) = self.state.votes.iter_mut().find(|r| r.judge == j + 1) {
1841 rec.vote = v.label();
1842 rec.reason = blind::sanitize_prose(&v.reason, &self.state.config.blind);
1843 } else {
1844 self.state.votes.push(VoteRecord {
1845 judge: j + 1,
1846 agent: agent_id,
1847 vote: v.label(),
1848 reason: blind::sanitize_prose(&v.reason, &self.state.config.blind),
1849 changed: false,
1850 });
1851 }
1852 self.state.event(
1853 "recover",
1854 format!("judge {} voted again after the limit", j + 1),
1855 );
1856 }
1857 Err(e) => {
1858 self.state
1859 .event("recover", format!("judge {} still cannot vote: {e}", j + 1));
1860 }
1861 }
1862 }
1863
1864 if !recovered.is_empty() {
1868 let recovered_keys: BTreeSet<String> = recovered
1869 .iter()
1870 .map(|&j| format!("judge-{}", j + 1))
1871 .collect();
1872 self.state
1873 .quota
1874 .retain(|q| !recovered_keys.contains(&q.seat));
1875 }
1876
1877 self.state.tally = None;
1879 self.tally()?;
1880 Ok(self
1881 .state
1882 .tally
1883 .as_ref()
1884 .map(|t| t.met_quorum)
1885 .unwrap_or(false))
1886 }
1887
1888 async fn fold_losers(&mut self) -> Result<()> {
1891 let Some(winner) = self.state.tally.as_ref().map(|t| t.winner) else {
1892 return Ok(());
1893 };
1894 let repo = self.state.repo.clone();
1895 let mut folded = Vec::new();
1896 for i in 0..self.state.candidates.len() {
1897 let c = &self.state.candidates[i];
1898 if c.label == winner || c.folded {
1899 continue;
1900 }
1901 let (wt, branch, label) = (c.worktree.clone(), c.branch.clone(), c.label);
1902 git::worktree_remove(&repo, &wt).await.ok();
1903 git::branch_delete(&repo, &branch).await.ok();
1904 self.state.candidates[i].folded = true;
1905 folded.push(label.to_string());
1906 }
1907 let root = self.state.worktree_root();
1909 for j in 1..=self.roles.judges.len() {
1910 let wt = root.join(format!("judge-{j}"));
1911 if wt.exists() {
1912 git::worktree_remove(&repo, &wt).await.ok();
1913 }
1914 }
1915 if !folded.is_empty() {
1916 self.state
1917 .event("fold", format!("folded candidates {}", folded.join(", ")));
1918 self.state.save()?;
1919 }
1920 Ok(())
1921 }
1922
1923 async fn sync_to_base(&mut self) -> Result<()> {
1953 if self
1954 .state
1955 .base_sync
1956 .as_ref()
1957 .is_some_and(|s| s.conflict.is_some())
1958 {
1959 return Ok(());
1960 }
1961 let Some(winner) = self.state.winner().cloned() else {
1962 return Ok(());
1963 };
1964
1965 let repo = self.state.repo.clone();
1966 let remote = self.state.config.merge.remote.clone();
1967 let base_branch = self.state.base_branch.clone();
1968 let tracking = format!("{remote}/{base_branch}");
1969
1970 git::fetch(&repo, &remote, &base_branch).await.ok();
1971 let Ok(tip) = git::rev_parse(&repo, &tracking).await else {
1975 return Ok(());
1976 };
1977
1978 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
1979 let behind = git::commits_ahead(&repo, &head, &tip).await.unwrap_or(0);
1980 let attempts = self.state.base_sync.as_ref().map_or(0, |s| s.attempts);
1981
1982 if behind == 0 {
1983 self.state.base_sync = Some(BaseSync {
1984 tip,
1985 behind: 0,
1986 attempts,
1987 conflict: None,
1988 });
1989 self.state.save()?;
1990 return Ok(());
1991 }
1992
1993 if attempts >= BASE_SYNC_ROUNDS {
1994 let why = format!(
1995 "{base_branch} moved {behind} commit(s) ahead of {} after {BASE_SYNC_ROUNDS} \
1996 rebase(s); rebasing again would only race it",
1997 winner.branch
1998 );
1999 self.state.status = RunStatus::Blocked;
2000 self.state.base_sync = Some(BaseSync {
2001 tip,
2002 behind,
2003 attempts,
2004 conflict: Some(why.clone()),
2005 });
2006 self.state.event("land", why);
2007 self.state.save()?;
2008 return Ok(());
2009 }
2010
2011 self.state.event(
2012 "land",
2013 format!(
2014 "{base_branch} moved {behind} commit(s) ahead of {}; rebasing before verifying",
2015 winner.branch
2016 ),
2017 );
2018 self.state.save()?;
2019
2020 let scratch = self.state.dir().join("base-sync");
2021 let rebased = git::rebase_branch_in_temp(&repo, &scratch, &winner.branch, &tracking).await;
2022 let attempts = attempts + 1;
2023 match rebased {
2024 Ok(None) => {
2025 git::sync_to_head(&winner.worktree).await?;
2029 self.state.base_sync = Some(BaseSync {
2030 tip: tip.clone(),
2031 behind: 0,
2032 attempts,
2033 conflict: None,
2034 });
2035 self.state
2036 .event("land", format!("rebased {} onto {tracking}", winner.branch));
2037 }
2038 Ok(Some(conflict)) => {
2039 let why = format!(
2040 "{} conflicts with {tracking} and did not rebase: {}",
2041 winner.branch,
2042 conflict.chars().take(600).collect::<String>()
2043 );
2044 self.state.status = RunStatus::Blocked;
2045 self.state.base_sync = Some(BaseSync {
2046 tip,
2047 behind,
2048 attempts,
2049 conflict: Some(why.clone()),
2050 });
2051 self.state.event("land", why);
2052 }
2053 Err(e) => {
2054 let why = format!("could not rebase {} onto {tracking}: {e:#}", winner.branch);
2055 self.state.status = RunStatus::Blocked;
2056 self.state.base_sync = Some(BaseSync {
2057 tip,
2058 behind,
2059 attempts,
2060 conflict: Some(why.clone()),
2061 });
2062 self.state.event("land", why);
2063 }
2064 }
2065 self.state.save()?;
2066 Ok(())
2067 }
2068
2069 fn landing_base(&self) -> String {
2079 self.state
2080 .base_sync
2081 .as_ref()
2082 .map_or_else(|| self.state.base_commit.clone(), |s| s.tip.clone())
2083 }
2084
2085 async fn review_loop(&mut self) -> Result<()> {
2088 if self
2093 .state
2094 .base_sync
2095 .as_ref()
2096 .is_some_and(|s| s.conflict.is_some())
2097 {
2098 return Ok(());
2099 }
2100 let run_id = self.state.id.clone();
2105 let prompts = self.state.config.prompts.clone();
2106 let Some(winner) = self.state.winner().cloned() else {
2107 return Ok(());
2108 };
2109 let max_rounds = self.state.config.graph.review_rounds;
2110 if let Some(status) = review_conclusion(&self.state.reviews, max_rounds) {
2120 self.state.status = status;
2121 self.state.save()?;
2122 return Ok(());
2123 }
2124 self.state.status = RunStatus::Reviewing;
2125
2126 let repo = self.state.repo.clone();
2127 let root = self.state.worktree_root();
2128 let language = self.state.config.graph.language.clone();
2129 let sessions = self.state.config.graph.sessions;
2130 let artifacts = agent::artifacts_dir(&self.state.dir());
2131 let base = self.landing_base();
2132 let base_short = short(&base);
2133 let reviewers = self.roles.reviewers.clone();
2134 let shell = self.state.config.shell();
2135
2136 let mut prev_e2e: Option<String> = None;
2137 for round in (self.state.reviews.len() + 1)..=max_rounds {
2138 let head = git::rev_parse(&winner.worktree, "HEAD").await?;
2139 let patch = git::diff(&winner.worktree, &base, "HEAD").await?;
2140 let stat = git::diff_stat(&winner.worktree, &base, "HEAD").await?;
2141
2142 let mut jobs = Vec::new();
2146 for (r, spec) in reviewers.iter().cloned().enumerate() {
2147 let wt = root.join(format!("review-{}", r + 1));
2148 if wt.exists() {
2149 git::reset_detached(&wt, &head).await?;
2150 } else {
2151 git::worktree_add_detached(&repo, &wt, &head).await?;
2152 }
2153 let seat_key = format!("review-{}", r + 1);
2154 let seat = self.seat(&seat_key, &spec.id);
2155 jobs.push(SeatJob {
2156 prompt: prompt::review(&prompt::ReviewCtx {
2157 instruction: &self.state.instruction,
2158 branch: &winner.branch,
2159 base_short: &base_short,
2160 stat: &stat,
2161 patch: &patch,
2162 e2e: prev_e2e.as_deref(),
2163 reviewers: reviewers.len(),
2164 round,
2165 rounds: max_rounds,
2166 competed: self.state.tally.as_ref().is_some_and(|t| t.rankings > 0),
2169 lens: Lens::for_seat(r),
2170 language: &language,
2171 }),
2172 spec,
2173 seat,
2174 cwd: wt,
2175 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
2176 allow_write: false,
2177 sessions,
2178 artifacts: artifacts.clone(),
2179 stem: format!("review-{round}-{}", r + 1),
2180 });
2181 }
2182
2183 self.state.event(
2184 "review",
2185 format!(
2186 "round {round}: {} reviewers on {}",
2187 jobs.len(),
2188 short(&head)
2189 ),
2190 );
2191 let mut quota_losses = Vec::new();
2192 let review_retries = self.state.config.graph.retries;
2193 let review_cache = self.state.config.cache_dir();
2194 let ctx = WaveCtx {
2195 run: &run_id,
2196 node: "review",
2197 prompts: &prompts,
2198 cache: review_cache.as_deref(),
2199 };
2200 let results = ask_json_wave::<Review>(
2201 jobs,
2202 Arc::clone(&self.sem),
2203 review_retries,
2204 &ctx,
2205 &mut quota_losses,
2206 &mut self.state,
2207 &|_: &Review| Ok(()),
2208 )
2209 .await;
2210 self.state.quota.extend(quota_losses);
2211
2212 let mut records = Vec::new();
2213 let mut all_findings = Vec::new();
2214 for (r, (seat, res)) in results.into_iter().enumerate() {
2215 let agent_id = seat.agent.clone();
2216 self.state.seats.insert(seat.key.clone(), seat);
2217 let mut record = ReviewRecord {
2218 reviewer: r + 1,
2219 agent: agent_id,
2220 summary: String::new(),
2221 findings: Vec::new(),
2222 vote: None,
2223 failed: None,
2224 duration_ms: 0,
2225 };
2226 match res {
2227 Ok((review, out)) => {
2228 record.summary =
2236 blind::sanitize_prose(&review.summary, &self.state.config.blind);
2237 record.vote = Some(review.vote);
2238 record.duration_ms = out.duration_ms;
2239 for (n, mut f) in review.findings.into_iter().enumerate() {
2240 f.id = format!("R{round}-{}-{}", r + 1, n + 1);
2243 f.title = blind::sanitize_prose(&f.title, &self.state.config.blind);
2244 f.detail = blind::sanitize_prose(&f.detail, &self.state.config.blind);
2245 f.file = f
2251 .file
2252 .map(|file| blind::sanitize_prose(&file, &self.state.config.blind));
2253 all_findings.push(f.clone());
2254 record.findings.push(f);
2255 }
2256 self.state.event(
2257 "review",
2258 format!(
2259 "round {round}: reviewer {} voted {} with {} finding(s)",
2260 r + 1,
2261 review.vote.label(),
2262 record.findings.len()
2263 ),
2264 );
2265 }
2266 Err(e) => {
2267 record.failed = Some(e.to_string());
2268 self.state.event(
2269 "review",
2270 format!("round {round}: reviewer {} produced nothing: {e}", r + 1),
2271 );
2272 }
2273 }
2274 records.push(record);
2275 }
2276
2277 let initial_votes: Vec<ReviewVote> = records.iter().filter_map(|r| r.vote).collect();
2284 let vote_split =
2285 initial_votes.len() > 1 && !initial_votes.iter().all(|v| *v == initial_votes[0]);
2286 let mut reconsideration: Vec<ReviewRevoteRecord> = Vec::new();
2287 if vote_split {
2288 self.state.event(
2289 "review",
2290 format!(
2291 "round {round}: votes split ({}) — one round of reconsideration",
2292 initial_votes
2293 .iter()
2294 .map(|v| v.label())
2295 .collect::<Vec<_>>()
2296 .join(", ")
2297 ),
2298 );
2299 let panel: Vec<ReviewSeatReport<'_>> = records
2302 .iter()
2303 .filter_map(|r| {
2304 r.vote.map(|vote| ReviewSeatReport {
2305 reviewer: r.reviewer,
2306 vote,
2307 summary: &r.summary,
2308 findings: &r.findings,
2309 })
2310 })
2311 .collect();
2312
2313 let mut jobs = Vec::new();
2314 let mut seats_at = Vec::new();
2315 for (r, spec) in reviewers.iter().cloned().enumerate() {
2316 if records[r].vote.is_none() {
2320 continue;
2321 }
2322 let wt = root.join(format!("review-{}", r + 1));
2323 let seat_key = format!("review-{}", r + 1);
2324 let seat = self.seat(&seat_key, &spec.id);
2325 let patch_ctx = if has_context(&spec, &seat, sessions) {
2330 None
2331 } else {
2332 Some(ReviewPatch {
2333 branch: &winner.branch,
2334 base_short: &base_short,
2335 stat: &stat,
2336 patch: &patch,
2337 })
2338 };
2339 let prompt = prompt::review_reconsider(&ReviewReconsiderCtx {
2340 instruction: &self.state.instruction,
2341 reviewer: r + 1,
2342 lens: Lens::for_seat(r),
2343 panel: &panel,
2344 patch: patch_ctx,
2345 round,
2346 rounds: max_rounds,
2347 language: &language,
2348 });
2349 jobs.push(SeatJob {
2350 prompt,
2351 spec,
2352 seat,
2353 cwd: wt,
2354 timeout: Duration::from_secs(self.state.config.graph.timeout_review),
2355 allow_write: false,
2356 sessions,
2357 artifacts: artifacts.clone(),
2358 stem: format!("review-{round}-reconsider-{}", r + 1),
2359 });
2360 seats_at.push(r);
2361 }
2362
2363 let mut recon_quota_losses = Vec::new();
2364 let recon_cache = self.state.config.cache_dir();
2365 let recon_ctx = WaveCtx {
2366 run: &run_id,
2367 node: "review",
2368 prompts: &prompts,
2369 cache: recon_cache.as_deref(),
2370 };
2371 let recon_results = ask_json_wave::<ReviewRevote>(
2372 jobs,
2373 Arc::clone(&self.sem),
2374 review_retries,
2375 &recon_ctx,
2376 &mut recon_quota_losses,
2377 &mut self.state,
2378 &|_: &ReviewRevote| Ok(()),
2379 )
2380 .await;
2381 self.state.quota.extend(recon_quota_losses);
2382
2383 for (&r, (seat, res)) in seats_at.iter().zip(recon_results) {
2384 let agent_id = seat.agent.clone();
2385 self.state.seats.insert(seat.key.clone(), seat);
2386 let mut rec = ReviewRevoteRecord {
2387 reviewer: r + 1,
2388 agent: agent_id,
2389 vote: None,
2390 reason: String::new(),
2391 failed: None,
2392 };
2393 match res {
2394 Ok((rv, _)) => {
2395 rec.vote = Some(rv.vote);
2396 rec.reason =
2397 blind::sanitize_prose(&rv.reason, &self.state.config.blind);
2398 self.state.event(
2399 "review",
2400 format!(
2401 "round {round}: reviewer {} revoted {}",
2402 r + 1,
2403 rv.vote.label()
2404 ),
2405 );
2406 }
2407 Err(e) => {
2408 rec.failed = Some(e.to_string());
2409 self.state.event(
2410 "review",
2411 format!("round {round}: reviewer {} did not revote: {e}", r + 1),
2412 );
2413 }
2414 }
2415 reconsideration.push(rec);
2416 }
2417 } else if initial_votes.len() > 1 {
2418 self.state.event(
2419 "review",
2420 format!(
2421 "round {round}: votes agreed ({}) — no reconsideration",
2422 initial_votes[0].label()
2423 ),
2424 );
2425 }
2426
2427 let final_votes: Vec<ReviewVote> = records
2431 .iter()
2432 .filter_map(|r| {
2433 reconsideration
2434 .iter()
2435 .find(|rv| rv.reviewer == r.reviewer)
2436 .and_then(|rv| rv.vote)
2437 .or(r.vote)
2438 })
2439 .collect();
2440 let round_verdict = ReviewVote::worst(final_votes);
2441
2442 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
2443 let verify_timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
2444 let defer_e2e =
2455 blocking > 0 && round < max_rounds && !self.state.config.graph.e2e_every_round;
2456 let (e2e, verify_retried, e2e_deferred, e2e_defer_reason) = if defer_e2e {
2457 let reason =
2458 format!("{blocking} blocking finding(s) already required a fix this round");
2459 self.state.event(
2460 "verify",
2461 format!(
2462 "round {round}: {reason} — e2e deferred to the fixer (reviewed head \
2463 {}); it will run once a round has none left",
2464 short(&head)
2465 ),
2466 );
2467 (Vec::new(), false, true, Some(reason))
2468 } else {
2469 let e2e_commands = self.state.config.verify.e2e.clone();
2470 let (e2e, verify_retried) = run_e2e_with_retry(
2471 &mut self.state,
2472 &shell,
2473 &e2e_commands,
2474 &winner.worktree,
2475 verify_timeout,
2476 &format!("round {round}"),
2477 )
2478 .await;
2479 (e2e, verify_retried, false, None)
2480 };
2481
2482 let e2e_failures: String = e2e
2483 .iter()
2484 .filter(|o| !o.ok())
2485 .map(|o| format!("$ {}\n{}\n", o.command, o.output_tail))
2486 .collect();
2487
2488 let expected = records.len();
2489 let answered = records.iter().filter(|r| r.failed.is_none()).count();
2490 let incomplete = answered < expected;
2491 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
2492 let policy = self.state.config.graph.incomplete_review;
2493 let clean = round_is_clean(blocking, e2e_ok, answered, expected, policy);
2494
2495 let mut round_record = ReviewRound {
2496 round,
2497 head: head.clone(),
2498 verified_head: None,
2499 reviews: records,
2500 e2e,
2501 verify_retried,
2502 e2e_deferred,
2503 e2e_defer_reason,
2504 fix: None,
2505 blocking,
2506 answered,
2507 expected,
2508 clean,
2509 progressed: false,
2510 vote_split,
2511 reconsideration,
2512 verdict: round_verdict,
2513 };
2514
2515 if incomplete {
2516 let missing: Vec<String> = round_record
2517 .reviews
2518 .iter()
2519 .filter(|r| r.failed.is_some())
2520 .map(|r| format!("review-{}", r.reviewer))
2521 .collect();
2522 self.state.event(
2523 "review",
2524 format!(
2525 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
2526 missing.join(", ")
2527 ),
2528 );
2529 }
2530
2531 if clean {
2532 self.state.event(
2533 "review",
2534 if incomplete {
2535 format!(
2536 "round {round}: clean (warn policy, incomplete panel) — no \
2537 blocking findings from the seats that answered, verification green"
2538 )
2539 } else {
2540 format!("round {round}: clean — no blocking findings, verification green")
2541 },
2542 );
2543 self.state.reviews.push(round_record);
2544 self.state.status = RunStatus::Gating;
2545 self.state.save()?;
2546 return Ok(());
2547 }
2548
2549 if incomplete && blocking == 0 && e2e_ok {
2553 self.state.reviews.push(round_record);
2554 self.state.save()?;
2555 if round == max_rounds {
2556 self.state.status = RunStatus::Blocked;
2557 self.state.event(
2558 "review",
2559 format!(
2560 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
2561 refusing to call it clean",
2562 expected - answered
2563 ),
2564 );
2565 return Ok(());
2566 }
2567 prev_e2e = None;
2568 continue;
2569 }
2570
2571 if round == max_rounds {
2572 self.state.reviews.push(round_record);
2573 return self
2574 .stop_reviewing(
2575 &format!(
2576 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
2577 ),
2578 &shell,
2579 &winner.worktree,
2580 )
2581 .await;
2582 }
2583
2584 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
2587 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
2588 _ => (
2589 self.state
2590 .config
2591 .agent(&winner.agent)
2592 .cloned()
2593 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
2594 format!("impl-{}", winner.label),
2595 ),
2596 };
2597 let seat = self.seat(&fix_seat_key, &fix_spec.id);
2598 let blocking_findings: Vec<_> = all_findings
2599 .iter()
2600 .filter(|f| f.severity.blocks())
2601 .cloned()
2602 .collect();
2603 let job = SeatJob {
2604 prompt: prompt::fix(
2605 &self.state.instruction,
2606 &blocking_findings,
2607 (!e2e_failures.is_empty()).then_some(e2e_failures.as_str()),
2608 e2e_deferred,
2609 round,
2610 max_rounds,
2611 &language,
2612 ),
2613 spec: fix_spec.clone(),
2614 seat,
2615 cwd: winner.worktree.clone(),
2616 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
2617 allow_write: true,
2618 sessions,
2619 artifacts: artifacts.clone(),
2620 stem: format!("fix-{round}"),
2621 };
2622 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
2623 let cache = self.state.config.cache_dir();
2624 let ctx = WaveCtx {
2625 run: &run_id,
2626 node: "fix",
2627 prompts: &prompts,
2628 cache: cache.as_deref(),
2629 };
2630 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2631 let agent_id = seat.agent.clone();
2632 let seat_key = seat.key.clone();
2633 self.state.seats.insert(seat.key.clone(), seat);
2634
2635 let mut fix = FixRecord {
2636 agent: agent_id,
2637 addressed: Vec::new(),
2638 rejected: Vec::new(),
2639 notes: String::new(),
2640 committed: false,
2641 failed: None,
2642 duration_ms: 0,
2643 };
2644 match out {
2645 AgentOutcome::Ok(o) => {
2646 fix.duration_ms = o.duration_ms;
2647 match verdict::extract_json::<FixReport>(&o.text) {
2648 Ok(report) => {
2649 fix.addressed = report.addressed;
2650 fix.rejected = report.rejected;
2651 fix.notes =
2652 blind::sanitize_prose(&report.notes, &self.state.config.blind);
2653 }
2654 Err(e) => fix.failed = Some(format!("unparsable fix report: {e}")),
2655 }
2656 }
2657 AgentOutcome::Dropped(o) => {
2659 fix.duration_ms = o.duration_ms;
2660 let why = o
2661 .dropped
2662 .as_ref()
2663 .map(|d| d.why.as_str())
2664 .unwrap_or("the CLI ended the stream without delivering its answer");
2665 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
2666 }
2667 AgentOutcome::Quota(o) => {
2668 self.state.quota.push(QuotaLoss {
2669 seat: seat_key,
2670 node: "fix".to_owned(),
2671 at: Timestamp::now(),
2672 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2673 });
2674 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
2675 }
2676 AgentOutcome::Failed(e) => fix.failed = Some(e),
2677 }
2678 git::commit_all(
2679 &winner.worktree,
2680 &format!("magi: review round {round} fixes (uncommitted work)"),
2681 )
2682 .await
2683 .ok();
2684 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
2685 fix.committed = after != before;
2686 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
2694 let progressed = diff_after != patch;
2695 let commit_note = if fix.committed {
2696 "committed"
2697 } else {
2698 "NO new commit"
2699 };
2700 let tree_note = if progressed {
2701 "changed vs base"
2702 } else {
2703 "unchanged vs base"
2704 };
2705 self.state.event(
2706 "fix",
2707 match &fix.failed {
2708 Some(reason) => {
2714 format!(
2715 "round {round}: fixer's adoption report was lost ({reason}); \
2716 {commit_note}, tree {tree_note}"
2717 )
2718 }
2719 None => format!(
2720 "round {round}: {} addressed, {} rejected, {commit_note}, tree {tree_note}",
2721 fix.addressed.len(),
2722 fix.rejected.len(),
2723 ),
2724 },
2725 );
2726 round_record.fix = Some(fix);
2727 round_record.progressed = progressed;
2728 self.state.reviews.push(round_record);
2729 self.state.save()?;
2730
2731 prev_e2e = (!e2e_failures.is_empty()).then_some(e2e_failures);
2732
2733 let streak = self
2734 .state
2735 .reviews
2736 .iter()
2737 .rev()
2738 .take_while(|r| !r.progressed)
2739 .count();
2740 if streak >= STAGNANT_LIMIT {
2741 return self
2742 .stop_reviewing(
2743 &format!(
2744 "the tree has not moved against base for {streak} round(s) in a row"
2745 ),
2746 &shell,
2747 &winner.worktree,
2748 )
2749 .await;
2750 }
2751 }
2752 Ok(())
2753 }
2754
2755 async fn stop_reviewing(&mut self, why: &str, shell: &[String], worktree: &Path) -> Result<()> {
2779 let round_idx = self.state.reviews.len() - 1;
2780 let needs_catchup_run = {
2781 let last = &self.state.reviews[round_idx];
2782 last.e2e.is_empty() && last.e2e_deferred
2783 };
2784 if needs_catchup_run {
2785 let round = self.state.reviews[round_idx].round;
2786 let timeout = Duration::from_secs(self.state.config.graph.verify_timeout());
2787 let commands = self.state.config.verify.e2e.clone();
2788 let verified_head = git::rev_parse(worktree, "HEAD").await?;
2789 let (outcomes, verify_retried) = run_e2e_with_retry(
2790 &mut self.state,
2791 shell,
2792 &commands,
2793 worktree,
2794 timeout,
2795 &format!("round {round}: deferred e2e, now catching up before the final decision"),
2796 )
2797 .await;
2798 let last = &mut self.state.reviews[round_idx];
2799 last.e2e = outcomes;
2800 last.verify_retried = verify_retried;
2801 last.e2e_deferred = false;
2802 if verified_head != last.head {
2803 last.verified_head = Some(verified_head);
2804 }
2805 }
2806 let last = &self.state.reviews[round_idx];
2807 let red: Vec<String> = last
2808 .e2e
2809 .iter()
2810 .filter(|o| !o.ok())
2811 .map(|o| {
2812 format!(
2813 "`{}` -> {:?}\n{}",
2814 o.command,
2815 o.code,
2816 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
2817 )
2818 })
2819 .collect();
2820 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
2821
2822 if red.is_empty() {
2823 self.state.event(
2824 "review",
2825 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
2826 );
2827 self.state.status = RunStatus::Gating;
2828 } else {
2829 self.state
2830 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
2831 self.state.status = RunStatus::Blocked;
2832 }
2833 self.state.save()?;
2834 Ok(())
2835 }
2836
2837 async fn gate(&mut self) -> Result<()> {
2840 if self.state.status == RunStatus::Failed
2852 || self
2853 .state
2854 .base_sync
2855 .as_ref()
2856 .is_some_and(|s| s.conflict.is_some())
2857 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
2858 != Some(RunStatus::Gating)
2859 {
2860 return Ok(());
2861 }
2862 if !self.state.gate.is_empty() {
2863 return Ok(());
2864 }
2865 let Some(winner) = self.state.winner().cloned() else {
2866 return Ok(());
2867 };
2868 self.state.status = RunStatus::Gating;
2869 let shell = self.state.config.shell();
2870 let outcomes = run_commands(
2871 &shell,
2872 &self.state.config.verify.gate,
2873 &winner.worktree,
2874 Duration::from_secs(self.state.config.graph.verify_timeout()),
2875 )
2876 .await;
2877 for o in &outcomes {
2878 self.state.event(
2879 "gate",
2880 format!(
2881 "`{}` -> {}",
2882 o.command,
2883 if o.ok() {
2884 "pass".to_owned()
2885 } else {
2886 format!(
2887 "FAIL ({:?})\n{}",
2888 o.code,
2889 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
2890 )
2891 }
2892 ),
2893 );
2894 }
2895 let passed = outcomes.iter().all(CommandOutcome::ok);
2896 self.state.gate = outcomes;
2897 if !passed {
2898 self.state.status = RunStatus::Blocked;
2899 self.state.event("gate", "gate failed; not merging");
2900 }
2901 self.state.save()?;
2902 Ok(())
2903 }
2904
2905 async fn merge(&mut self) -> Result<()> {
2908 if self
2923 .state
2924 .base_sync
2925 .as_ref()
2926 .is_some_and(|s| s.conflict.is_some())
2927 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
2928 != Some(RunStatus::Gating)
2929 || self.state.gate.iter().any(|o| !o.ok())
2930 {
2931 return Ok(());
2932 }
2933 if self.state.merge.is_some() {
2942 return Ok(());
2943 }
2944 let Some(winner) = self.state.winner().cloned() else {
2945 return Ok(());
2946 };
2947 let repo = self.state.repo.clone();
2948 let base = self.state.base_branch.clone();
2949 let mode = self.state.config.merge.mode;
2950 let style = self.state.config.merge.style;
2951 let message = pr_body(&self.state, winner.label);
2952
2953 let outcome = match mode {
2954 MergeMode::None => MergeOutcome {
2955 mode,
2956 ok: true,
2957 detail: manual_merge_command(style, &repo, &winner.branch, &message),
2958 },
2959 MergeMode::Local => {
2960 let on = git::current_branch(&repo).await?;
2961 if on.as_deref() != Some(base.as_str()) {
2962 MergeOutcome {
2963 mode,
2964 ok: false,
2965 detail: format!(
2966 "{} has {} checked out, not the base branch {base}",
2967 repo.display(),
2968 on.unwrap_or_else(|| "a detached HEAD".to_owned())
2969 ),
2970 }
2971 } else if !git::is_clean(&repo).await? {
2972 MergeOutcome {
2973 mode,
2974 ok: false,
2975 detail: format!("{} is dirty; refusing to merge", repo.display()),
2976 }
2977 } else {
2978 let out = match style {
2979 MergeStyle::Merge => {
2980 git::merge_no_ff(&repo, &winner.branch, &message).await?
2981 }
2982 MergeStyle::Squash => {
2983 git::merge_squash(&repo, &winner.branch, &message).await?
2984 }
2985 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
2986 };
2987 MergeOutcome {
2988 mode,
2989 ok: out.ok(),
2990 detail: if out.ok() { out.stdout } else { out.stderr },
2991 }
2992 }
2993 }
2994 MergeMode::Pr => {
2995 let remote = self.state.config.merge.remote.clone();
2996 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
2997 if !pushed.ok() {
2998 MergeOutcome {
2999 mode,
3000 ok: false,
3001 detail: pushed.stderr,
3002 }
3003 } else {
3004 let out = gh_pr_create(&winner.worktree, &base, &winner.branch, &message).await;
3005 match out {
3006 Ok(url) => MergeOutcome {
3007 mode,
3008 ok: true,
3009 detail: url,
3010 },
3011 Err(e) => MergeOutcome {
3012 mode,
3013 ok: false,
3014 detail: e.to_string(),
3015 },
3016 }
3017 }
3018 }
3019 };
3020
3021 self.state.status = match (mode, outcome.ok) {
3022 (MergeMode::None, _) => RunStatus::Ready,
3023 (_, true) => RunStatus::Merged,
3024 (_, false) => RunStatus::Blocked,
3025 };
3026 self.state.event(
3027 "merge",
3028 format!(
3029 "{:?}: {}",
3030 mode,
3031 outcome.detail.lines().next().unwrap_or("")
3032 ),
3033 );
3034 self.state.merge = Some(outcome);
3035 self.state.save()?;
3036
3037 if self.state.config.graph.land
3043 && mode == MergeMode::Pr
3044 && self.state.status == RunStatus::Merged
3045 {
3046 self.run_land().await?;
3047 }
3048 self.settle_questions();
3053 Ok(())
3054 }
3055
3056 async fn run_land(&mut self) -> Result<()> {
3067 let url = self
3068 .state
3069 .merge
3070 .as_ref()
3071 .map(|m| m.detail.clone())
3072 .unwrap_or_default();
3073 let url = url.lines().next().unwrap_or("").trim().to_owned();
3074 if !url.starts_with("http") {
3075 return Ok(());
3076 }
3077 match land::land(&mut self.state, &url).await {
3080 Ok(pr) if self.state.parked => {
3081 let _ = pr;
3085 }
3086 Ok(pr) => {
3087 self.state.status = match pr.state {
3088 land::PrLifecycle::Merged => RunStatus::Merged,
3089 _ => RunStatus::Blocked,
3090 };
3091 if bump::should_release_bump(self.state.status)
3098 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
3099 {
3100 self.state
3101 .event("bump", format!("release bump skipped: {e:#}"));
3102 }
3103 self.state.save()?;
3104 }
3105 Err(e) => {
3106 self.state.status = RunStatus::Blocked;
3107 self.state.event("land", format!("gave up: {e}"));
3108 self.state.save()?;
3109 }
3110 }
3111 Ok(())
3112 }
3113
3114 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
3118 if let Some(existing) = self.state.seats.get(key)
3119 && existing.agent == agent
3120 {
3121 return existing.clone();
3122 }
3123 let fresh = SeatState::new(key, agent, self.state.seed);
3124 self.state.seats.insert(key.to_owned(), fresh.clone());
3125 fresh
3126 }
3127
3128 fn view(&self, c: &Candidate) -> CandidateView {
3130 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
3131 .unwrap_or_default();
3132 let (patch, _) = blind::sanitize_patch(
3133 &format!("candidate {} patch", c.label),
3134 &raw,
3135 &self.state.config.blind,
3136 );
3137 CandidateView {
3138 label: c.label,
3139 branch: c.branch.clone(),
3140 summary: c.summary.clone(),
3141 stat: c.stat.clone(),
3142 patch,
3143 }
3144 }
3145
3146 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
3148 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
3149 prompt::judge(
3150 "(see above)",
3151 &views,
3152 self.roles.judges.len(),
3153 base_short,
3154 "en",
3155 )
3156 }
3157
3158 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
3165 let mut turns = Vec::new();
3166 for j in &self.state.judgements {
3167 if j.ranking.is_empty() {
3168 continue;
3169 }
3170 let reasons = j
3171 .reasons
3172 .iter()
3173 .map(|(k, v)| format!("- {k}: {v}"))
3174 .collect::<Vec<_>>()
3175 .join("\n");
3176 turns.push(Turn {
3177 who: format!("Judge {} (opening ranking)", j.judge),
3178 is_self: j.judge == self_idx + 1,
3179 body: format!(
3180 "Ranked {}{}{reasons}",
3181 j.ranking.iter().collect::<String>(),
3182 if reasons.is_empty() {
3183 ""
3184 } else {
3185 ", because:\n"
3186 }
3187 ),
3188 });
3189 }
3190 for t in self
3191 .state
3192 .deliberation
3193 .iter()
3194 .flat_map(|r| r.turns.iter())
3195 .chain(current)
3196 {
3197 turns.push(Turn {
3198 who: format!("Judge {}", t.judge),
3199 is_self: t.judge == self_idx + 1,
3200 body: t.body.clone(),
3201 });
3202 }
3203 turns
3204 }
3205}
3206
3207fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
3209 agent::has_session(spec.kind, seat, sessions)
3210}
3211
3212fn short(commit: &str) -> String {
3213 commit.chars().take(7).collect()
3214}
3215
3216fn make_executable(path: &Path) -> Result<()> {
3217 #[cfg(unix)]
3218 {
3219 use std::os::unix::fs::PermissionsExt as _;
3220 let mut perms = std::fs::metadata(path)?.permissions();
3221 perms.set_mode(0o755);
3222 std::fs::set_permissions(path, perms)?;
3223 }
3224 #[cfg(not(unix))]
3225 {
3226 let _ = path;
3227 }
3228 Ok(())
3229}
3230
3231struct WaveCtx<'a> {
3238 run: &'a str,
3241 node: &'a str,
3243 prompts: &'a Prompts,
3244 cache: Option<&'a Path>,
3246}
3247
3248async fn run_one(
3250 job: SeatJob,
3251 sem: Arc<Semaphore>,
3252 ctx: &WaveCtx<'_>,
3253 state: &mut RunState,
3254 attempt: usize,
3255) -> (SeatState, AgentOutcome) {
3256 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
3257 .await
3258 .pop()
3259 .expect("one job in, one result out");
3260 (seat, out)
3261}
3262
3263async fn wave(
3269 jobs: Vec<SeatJob>,
3270 sem: Arc<Semaphore>,
3271 ctx: &WaveCtx<'_>,
3272 state: &mut RunState,
3273 attempt: usize,
3274) -> Vec<(usize, SeatState, AgentOutcome)> {
3275 let WaveCtx {
3276 run,
3277 node,
3278 prompts,
3279 cache,
3280 } = *ctx;
3281 for job in &jobs {
3282 state.seat_started(node, &job.seat.key, job.timeout, attempt);
3283 }
3284 if let Err(e) = state.save() {
3285 tracing::warn!("could not persist in-progress seats: {e:#}");
3290 }
3291 let mut set = tokio::task::JoinSet::new();
3292 let overlay = prompts.overlay(node);
3293 for (i, mut job) in jobs.into_iter().enumerate() {
3294 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
3295 if cache.is_some() {
3296 job.prompt.push('\n');
3297 job.prompt.push_str(&prompt::build_cache_note(node));
3298 }
3299 let sem = Arc::clone(&sem);
3300 let run = run.to_owned();
3301 let node = node.to_owned();
3302 let cache = cache.map(Path::to_path_buf);
3303 set.spawn(async move {
3304 let _permit = sem.acquire().await;
3305 let mut seat = job.seat;
3306 let out = agent::invoke(
3307 &job.spec,
3308 &mut seat,
3309 &Invocation {
3310 cwd: &job.cwd,
3311 prompt: &job.prompt,
3312 timeout: job.timeout,
3313 allow_write: job.allow_write,
3314 sessions: job.sessions,
3315 artifacts: &job.artifacts,
3316 stem: &job.stem,
3317 run: &run,
3318 node: &node,
3319 cache_dir: cache.as_deref(),
3320 attachments: &[],
3321 },
3322 )
3323 .await;
3324 let out = match out {
3325 Ok(o) if o.usable() => AgentOutcome::Ok(o),
3326 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
3327 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
3335 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
3336 Ok(o) => AgentOutcome::Failed(format!(
3337 "exited with {:?} and no usable output",
3338 o.exit_code
3339 )),
3340 Err(e) => AgentOutcome::Failed(e.to_string()),
3341 };
3342 (i, seat, out)
3343 });
3344 }
3345 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
3346 while let Some(joined) = set.join_next().await {
3347 let (i, seat, out) = match joined {
3348 Ok(v) => v,
3349 Err(e) => {
3353 tracing::error!("agent task panicked: {e}");
3354 continue;
3355 }
3356 };
3357 state.seat_finished(&seat.key);
3358 if let Err(e) = state.save() {
3359 tracing::warn!("could not persist a seat's completion: {e:#}");
3360 }
3361 if collected.len() <= i {
3362 collected.resize_with(i + 1, || None);
3363 }
3364 collected[i] = Some((i, seat, out));
3365 }
3366 if state
3372 .active
3373 .values()
3374 .any(|a| a.node == node && a.attempt == attempt)
3375 {
3376 state
3377 .active
3378 .retain(|_, a| !(a.node == node && a.attempt == attempt));
3379 if let Err(e) = state.save() {
3380 tracing::warn!("could not persist the end of a wave: {e:#}");
3381 }
3382 }
3383 collected.into_iter().flatten().collect()
3384}
3385
3386fn round_is_clean(
3396 blocking: usize,
3397 e2e_ok: bool,
3398 answered: usize,
3399 expected: usize,
3400 policy: IncompleteReviewPolicy,
3401) -> bool {
3402 blocking == 0 && e2e_ok && (answered == expected || policy == IncompleteReviewPolicy::Warn)
3403}
3404
3405fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
3422 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
3423 return Some(RunStatus::Gating);
3424 }
3425 let last = reviews.last()?;
3426 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
3427 if reviews.len() < max_rounds && !stagnant {
3428 return None;
3429 }
3430 Some(if last.incomplete() && last.blocking == 0 {
3431 RunStatus::Blocked
3432 } else if last.e2e.iter().all(CommandOutcome::ok) {
3433 RunStatus::Gating
3434 } else {
3435 RunStatus::Blocked
3436 })
3437}
3438
3439fn retry_budget(full: Duration, nudged: bool) -> Duration {
3454 if nudged {
3455 (full / 4).max(Duration::from_secs(120)).min(full)
3456 } else {
3457 full
3458 }
3459}
3460
3461#[allow(clippy::too_many_arguments)]
3474async fn ask_json_wave<T>(
3475 jobs: Vec<SeatJob>,
3476 sem: Arc<Semaphore>,
3477 retries: usize,
3478 ctx: &WaveCtx<'_>,
3479 losses: &mut Vec<QuotaLoss>,
3480 state: &mut RunState,
3481 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
3482) -> Vec<(SeatState, Result<(T, AgentOutput)>)>
3483where
3484 T: serde::de::DeserializeOwned + Send + 'static,
3485{
3486 let n = jobs.len();
3487 let originals: Vec<SeatJob> = jobs;
3488 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
3489 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
3490 let mut pending: Vec<usize> = (0..n).collect();
3491
3492 for attempt in 0..=retries {
3493 if pending.is_empty() {
3494 break;
3495 }
3496 let mut batch = Vec::with_capacity(pending.len());
3497 for &i in &pending {
3498 let src = &originals[i];
3499 let (prompt, timeout) = if attempt == 0 {
3502 (src.prompt.clone(), src.timeout)
3503 } else {
3504 let why = done[i]
3505 .as_ref()
3506 .and_then(|r| r.as_ref().err().map(ToString::to_string))
3507 .unwrap_or_else(|| "no parsable answer".to_owned());
3508 let nudge = prompt::nudge(&why);
3509 let nudged = has_context(&src.spec, &seats[i], src.sessions);
3510 let prompt = if nudged {
3511 nudge
3512 } else {
3513 format!("{}\n\n---\n\n{}", src.prompt, nudge)
3514 };
3515 (prompt, retry_budget(src.timeout, nudged))
3516 };
3517 batch.push(SeatJob {
3518 spec: src.spec.clone(),
3519 seat: seats[i].clone(),
3520 cwd: src.cwd.clone(),
3521 prompt,
3522 timeout,
3523 allow_write: src.allow_write,
3524 sessions: src.sessions,
3525 artifacts: src.artifacts.clone(),
3526 stem: if attempt == 0 {
3527 src.stem.clone()
3528 } else {
3529 format!("{}-retry{attempt}", src.stem)
3530 },
3531 });
3532 }
3533
3534 if attempt > 0 {
3535 let seats_out: Vec<&str> = pending
3536 .iter()
3537 .map(|&i| originals[i].seat.key.as_str())
3538 .collect();
3539 state.event(
3540 ctx.node,
3541 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
3542 );
3543 }
3544 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
3545 let mut still = Vec::new();
3546 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
3547 seats[i] = seat;
3548 let (parsed, quota) = match out {
3549 AgentOutcome::Ok(o) => (
3550 match verdict::extract_json::<T>(&o.text) {
3551 Ok(v) => match validate(&v) {
3552 Ok(()) => Ok((v, o)),
3553 Err(e) => Err(e),
3554 },
3555 Err(e) => Err(e),
3556 },
3557 false,
3558 ),
3559 AgentOutcome::Quota(o) => {
3560 losses.push(QuotaLoss {
3561 seat: originals[i].seat.key.clone(),
3562 node: ctx.node.to_owned(),
3563 at: Timestamp::now(),
3564 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3565 });
3566 (
3567 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
3568 true,
3569 )
3570 }
3571 AgentOutcome::Dropped(o) => {
3576 let why = o
3577 .dropped
3578 .as_ref()
3579 .map(|d| d.why.as_str())
3580 .unwrap_or("the CLI ended the stream without delivering its answer");
3581 (
3582 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
3583 false,
3584 )
3585 }
3586 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
3587 };
3588 let failed = parsed.is_err();
3589 done[i] = Some(parsed);
3590 if failed && !quota {
3593 still.push(i);
3594 }
3595 }
3596 pending = still;
3597 }
3598
3599 seats
3600 .into_iter()
3601 .zip(done)
3602 .map(|(seat, res)| {
3603 (
3604 seat,
3605 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
3606 )
3607 })
3608 .collect()
3609}
3610
3611fn e2e_outcome_label(o: &CommandOutcome) -> String {
3615 if o.ok() {
3616 return "pass".to_owned();
3617 }
3618 let reason = if o.build_failed() {
3619 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
3620 } else {
3621 format!("FAIL ({:?})", o.code)
3622 };
3623 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
3624}
3625
3626async fn run_e2e_with_retry(
3634 state: &mut RunState,
3635 shell: &[String],
3636 commands: &[String],
3637 worktree: &Path,
3638 timeout: Duration,
3639 context: &str,
3640) -> (Vec<CommandOutcome>, bool) {
3641 let mut e2e = run_commands(shell, commands, worktree, timeout).await;
3642 for o in &e2e {
3643 state.event(
3644 "verify",
3645 format!("{context}: `{}` -> {}", o.command, e2e_outcome_label(o)),
3646 );
3647 }
3648 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
3652 if verify_retried {
3653 state.event(
3654 "verify",
3655 format!(
3656 "{context}: verify could not build/link, not a test result — retrying once \
3657 before concluding"
3658 ),
3659 );
3660 e2e = run_commands(shell, commands, worktree, timeout).await;
3661 for o in &e2e {
3662 state.event(
3663 "verify",
3664 format!(
3665 "{context}: retry `{}` -> {}",
3666 o.command,
3667 e2e_outcome_label(o)
3668 ),
3669 );
3670 }
3671 }
3672 (e2e, verify_retried)
3673}
3674
3675async fn run_commands(
3677 shell: &[String],
3678 commands: &[String],
3679 cwd: &Path,
3680 timeout: Duration,
3681) -> Vec<CommandOutcome> {
3682 let mut out = Vec::new();
3683 for command in commands {
3684 let started = Instant::now();
3685 let mut cmd = tokio::process::Command::new(&shell[0]);
3686 cmd.quiet();
3687 cmd.args(&shell[1..])
3688 .arg(command)
3689 .current_dir(cwd)
3690 .stdin(std::process::Stdio::null())
3691 .stdout(std::process::Stdio::piped())
3692 .stderr(std::process::Stdio::piped())
3693 .kill_on_drop(true);
3694 let spawned = cmd.spawn();
3695 let (code, body) = match spawned {
3696 Ok(child) => match tokio::time::timeout(timeout, child.wait_with_output()).await {
3697 Ok(Ok(o)) => {
3698 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
3699 body.push_str(&String::from_utf8_lossy(&o.stderr));
3700 (o.status.code(), body)
3701 }
3702 Ok(Err(e)) => (None, format!("failed to run: {e}")),
3703 Err(_) => (None, format!("timed out after {}s", timeout.as_secs())),
3704 },
3705 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
3706 };
3707 out.push(CommandOutcome {
3708 command: command.clone(),
3709 code,
3710 output_tail: tail(&body, OUTPUT_TAIL),
3711 duration_ms: started.elapsed().as_millis() as u64,
3712 });
3713 }
3714 out
3715}
3716
3717fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
3729 let repo = repo.display();
3730 match style {
3731 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
3732 MergeStyle::Squash => {
3733 let subject = message.lines().next().unwrap_or(branch);
3734 format!(
3735 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
3736 )
3737 }
3738 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
3739 }
3740}
3741
3742fn pr_body(state: &RunState, winner: char) -> String {
3748 let mut message = format!(
3749 "Merge magi run {} (candidate {winner})\n\n{}",
3750 state.id, state.instruction
3751 );
3752
3753 let open = state.open_findings();
3754 if !open.is_empty() {
3755 message.push_str("\n\n## Open review findings\n\n");
3756 for f in &open {
3757 message.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
3758 }
3759 }
3760
3761 if let Some(fix) = state.reviews.last().and_then(|r| r.fix.as_ref())
3762 && !fix.rejected.is_empty()
3763 {
3764 message.push_str("\n## Declined by the fixer\n\n");
3765 for r in &fix.rejected {
3766 message.push_str(&format!("- `{}`: {}\n", r.id, r.why));
3767 }
3768 }
3769
3770 message
3771}
3772
3773async fn gh_pr_create(cwd: &Path, base: &str, head: &str, body: &str) -> Result<String> {
3775 let title = body.lines().next().unwrap_or("magi run").to_owned();
3776 let out = tokio::process::Command::new("gh")
3777 .args([
3778 "pr", "create", "--base", base, "--head", head, "--title", &title, "--body", body,
3779 ])
3780 .current_dir(cwd)
3781 .quiet()
3782 .stdin(std::process::Stdio::null())
3783 .output()
3784 .await
3785 .context("spawn gh")?;
3786 if out.status.success() {
3787 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
3788 } else {
3789 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
3790 }
3791}
3792
3793pub async fn fold_run(state: &mut RunState, drop_winner: bool) -> Result<Vec<String>> {
3795 let repo = state.repo.clone();
3796 let root = state.worktree_root();
3797 let winner = state.tally.as_ref().map(|t| t.winner);
3798 let mut removed = Vec::new();
3799
3800 for i in 0..state.candidates.len() {
3801 let c = state.candidates[i].clone();
3802 let is_winner = Some(c.label) == winner;
3803 if is_winner && !drop_winner {
3804 continue;
3805 }
3806 if c.worktree.exists() {
3807 git::worktree_remove(&repo, &c.worktree).await.ok();
3808 removed.push(c.worktree.to_string_lossy().into_owned());
3809 }
3810 if git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
3811 git::branch_delete(&repo, &c.branch).await.ok();
3812 removed.push(c.branch.clone());
3813 }
3814 state.candidates[i].folded = true;
3815 }
3816
3817 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
3818 let path = name.path();
3819 let keep = !drop_winner
3820 && winner.is_some_and(|w| {
3821 path.file_name()
3822 .is_some_and(|n| n == format!("cand-{w}").as_str())
3823 });
3824 if keep {
3825 continue;
3826 }
3827 git::worktree_remove(&repo, &path).await.ok();
3828 removed.push(path.to_string_lossy().into_owned());
3829 }
3830
3831 remove_if_empty(&root);
3840
3841 if state.enabled_worktree_config && drop_winner {
3842 git::release_worktree_config(&repo).await.ok();
3846 state.enabled_worktree_config = false;
3847 }
3848 state.save()?;
3849 Ok(removed)
3850}
3851
3852fn remove_if_empty(dir: &Path) {
3863 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
3864 std::fs::remove_dir(dir).ok();
3865 }
3866}
3867
3868pub fn worst_open(state: &RunState) -> Option<Severity> {
3870 state
3871 .reviews
3872 .last()?
3873 .reviews
3874 .iter()
3875 .flat_map(|r| r.findings.iter())
3876 .map(|f| f.severity)
3877 .max()
3878}
3879
3880#[cfg(test)]
3881mod tests {
3882 use super::*;
3883 use std::collections::BTreeMap;
3884 use std::time::Duration;
3885
3886 fn conductor() -> AgentSpec {
3887 AgentSpec {
3888 id: "conductor".to_owned(),
3889 kind: crate::config::AgentKind::Command,
3890 model: None,
3891 command: vec!["true".to_owned()],
3892 extra_args: Vec::new(),
3893 env: BTreeMap::new(),
3894 prompt_delivery: None,
3895 }
3896 }
3897
3898 #[test]
3899 fn remove_if_empty_only_ever_takes_a_bare_directory() {
3900 let dir = tempfile::tempdir().unwrap();
3901 let bay = dir.path().join("ffff");
3902
3903 remove_if_empty(&bay);
3905 assert!(!bay.exists());
3906
3907 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
3910 remove_if_empty(&bay);
3911 assert!(bay.exists(), "non-empty directory must survive");
3912
3913 std::fs::remove_dir(bay.join("cand-A")).unwrap();
3915 remove_if_empty(&bay);
3916 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
3917 }
3918
3919 #[test]
3928 fn a_full_panel_that_found_nothing_is_clean() {
3929 assert!(round_is_clean(0, true, 2, 2, IncompleteReviewPolicy::Block));
3930 }
3931
3932 #[test]
3933 fn a_missing_seat_is_never_clean_under_the_default_policy() {
3934 assert!(!round_is_clean(
3935 0,
3936 true,
3937 1,
3938 2,
3939 IncompleteReviewPolicy::Block
3940 ));
3941 }
3942
3943 #[test]
3944 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
3945 assert!(!round_is_clean(1, true, 1, 2, IncompleteReviewPolicy::Warn));
3946 }
3947
3948 #[test]
3949 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
3950 assert!(round_is_clean(0, true, 1, 2, IncompleteReviewPolicy::Warn));
3951 }
3952
3953 #[test]
3954 fn a_full_panel_with_an_open_finding_is_not_clean() {
3955 assert!(!round_is_clean(
3956 1,
3957 true,
3958 2,
3959 2,
3960 IncompleteReviewPolicy::Block
3961 ));
3962 }
3963
3964 #[test]
3965 fn a_full_panel_with_a_red_e2e_is_not_clean() {
3966 assert!(!round_is_clean(
3967 0,
3968 false,
3969 2,
3970 2,
3971 IncompleteReviewPolicy::Block
3972 ));
3973 }
3974
3975 fn review_round(
3981 clean: bool,
3982 blocking: usize,
3983 answered: usize,
3984 expected: usize,
3985 progressed: bool,
3986 e2e_ok: bool,
3987 ) -> ReviewRound {
3988 ReviewRound {
3989 round: 1,
3990 head: "h".to_owned(),
3991 verified_head: None,
3992 reviews: Vec::new(),
3993 e2e: vec![CommandOutcome {
3994 command: "test".to_owned(),
3995 code: Some(if e2e_ok { 0 } else { 1 }),
3996 output_tail: String::new(),
3997 duration_ms: 0,
3998 }],
3999 verify_retried: false,
4000 e2e_deferred: false,
4001 e2e_defer_reason: None,
4002 fix: None,
4003 blocking,
4004 answered,
4005 expected,
4006 clean,
4007 progressed,
4008 vote_split: false,
4009 reconsideration: Vec::new(),
4010 verdict: None,
4011 }
4012 }
4013
4014 #[test]
4015 fn review_conclusion_is_none_when_nothing_has_run() {
4016 assert_eq!(review_conclusion(&[], 3), None);
4017 }
4018
4019 #[test]
4020 fn review_conclusion_is_none_while_rounds_remain() {
4021 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
4022 assert_eq!(review_conclusion(&rounds, 3), None);
4023 }
4024
4025 #[test]
4026 fn review_conclusion_is_gating_once_a_round_is_clean() {
4027 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
4028 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
4029 }
4030
4031 #[test]
4032 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
4033 let rounds = vec![
4034 review_round(false, 1, 2, 2, true, true),
4035 review_round(false, 1, 2, 2, true, true),
4036 ];
4037 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
4038 }
4039
4040 #[test]
4041 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
4042 let rounds = vec![
4043 review_round(false, 1, 2, 2, true, true),
4044 review_round(false, 1, 2, 2, true, false),
4045 ];
4046 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
4047 }
4048
4049 #[test]
4050 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
4051 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
4053 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
4054 }
4055
4056 #[test]
4057 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
4058 let rounds = vec![
4059 review_round(false, 1, 2, 2, false, true),
4060 review_round(false, 1, 2, 2, false, true),
4061 ];
4062 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
4063 }
4064
4065 fn secs(n: u64) -> Duration {
4066 Duration::from_secs(n)
4067 }
4068
4069 fn init_repo(dir: &Path) {
4072 let run = |args: &[&str]| {
4073 let out = std::process::Command::new("git")
4074 .args(args)
4075 .current_dir(dir)
4076 .quiet()
4077 .output()
4078 .expect("spawn git");
4079 assert!(
4080 out.status.success(),
4081 "git {args:?} failed: {}",
4082 String::from_utf8_lossy(&out.stderr)
4083 );
4084 };
4085 run(&["init", "-b", "main"]);
4086 run(&["config", "user.name", "magi test"]);
4087 run(&["config", "user.email", "magi@example.com"]);
4088 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
4089 run(&["add", "-A"]);
4090 run(&["commit", "-m", "init"]);
4091 }
4092
4093 fn ask_test_home() {
4101 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
4102 }
4103
4104 fn runner_at(status: RunStatus) -> Runner {
4107 let mut state = RunState::new(
4108 PathBuf::from("/nonexistent/repo"),
4109 "main".to_owned(),
4110 "deadbeef".to_owned(),
4111 "task".to_owned(),
4112 Config::default(),
4113 );
4114 state.status = status;
4115 Runner {
4116 state,
4117 roles: ResolvedRoles {
4118 implementers: Vec::new(),
4119 judges: Vec::new(),
4120 reviewers: Vec::new(),
4121 fixer: None,
4122 conductor: conductor(),
4123 },
4124 sem: Arc::new(Semaphore::new(1)),
4125 pause: Pause::new(),
4126 }
4127 }
4128
4129 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
4131 let mut q = ask::Question::new(
4132 run.to_owned(),
4133 "implement".to_owned(),
4134 "impl-A".to_owned(),
4135 "Which storage backend should the cache use?".to_owned(),
4136 String::new(),
4137 vec!["SQLite".to_owned(), "Redis".to_owned()],
4138 );
4139 store.put(&mut q).unwrap();
4140 q
4141 }
4142
4143 #[test]
4144 fn a_failed_runs_open_question_is_abandoned() {
4145 ask_test_home();
4146 let store = ask::Questions::open();
4147 let mut runner = runner_at(RunStatus::Failed);
4148 let run = runner.state.id.clone();
4149 let q = ask_open_question(&store, &run);
4150
4151 runner.settle_questions();
4152
4153 let back = store.get(&q.id).unwrap();
4154 assert!(
4155 !back.status.open(),
4156 "the seat that asked died with the run; nobody is left to read an answer"
4157 );
4158 assert!(
4159 back.detail.contains(&run) && back.detail.contains("failed"),
4160 "the reason names what the run became, not just that it is gone: {}",
4161 back.detail
4162 );
4163 }
4164
4165 #[test]
4166 fn a_merged_runs_open_question_is_abandoned_too() {
4167 ask_test_home();
4168 let store = ask::Questions::open();
4169 for status in [RunStatus::Merged, RunStatus::Ready] {
4172 let mut runner = runner_at(status);
4173 let run = runner.state.id.clone();
4174 let q = ask_open_question(&store, &run);
4175
4176 runner.settle_questions();
4177
4178 let back = store.get(&q.id).unwrap();
4179 assert!(
4180 !back.status.open(),
4181 "{status:?} run's question must not outlive the run"
4182 );
4183 }
4184 }
4185
4186 #[test]
4187 fn a_still_resumable_runs_open_question_is_left_alone() {
4188 ask_test_home();
4189 let store = ask::Questions::open();
4190 for status in [RunStatus::Blocked, RunStatus::Stalled] {
4196 let mut runner = runner_at(status);
4197 let run = runner.state.id.clone();
4198 let q = ask_open_question(&store, &run);
4199
4200 runner.settle_questions();
4201
4202 let back = store.get(&q.id).unwrap();
4203 assert!(
4204 back.status.open(),
4205 "{status:?} is still alive; the question must still be waiting"
4206 );
4207 }
4208 }
4209
4210 #[test]
4211 fn settle_questions_never_touches_an_already_answered_question() {
4212 ask_test_home();
4213 let store = ask::Questions::open();
4214 let mut runner = runner_at(RunStatus::Failed);
4215 let run = runner.state.id.clone();
4216 let mut q = ask_open_question(&store, &run);
4217 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
4218 .unwrap();
4219 store.put(&mut q).unwrap();
4220
4221 runner.settle_questions();
4226 runner.settle_questions();
4227
4228 let back = store.get(&q.id).unwrap();
4229 assert_eq!(
4230 back.status,
4231 ask::QuestionStatus::Answered,
4232 "a real answer is a decision on record, never overwritten by a sweep"
4233 );
4234 }
4235
4236 #[tokio::test]
4245 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
4246 let tmp = tempfile::tempdir().expect("tempdir");
4247 let repo = tmp.path().join("repo");
4248 std::fs::create_dir_all(&repo).unwrap();
4249 init_repo(&repo);
4250
4251 let mut config = Config::default();
4252 config.merge.mode = MergeMode::Local;
4253
4254 let mut state = RunState::new(
4255 repo.clone(),
4256 "main".to_owned(),
4257 "deadbeef".to_owned(),
4258 "task".to_owned(),
4259 config,
4260 );
4261 state.candidates = vec![Candidate {
4262 index: 0,
4263 label: 'A',
4264 agent: "alpha".to_owned(),
4265 branch: "does-not-exist".to_owned(),
4266 worktree: repo.clone(),
4267 summary: String::new(),
4268 stat: String::new(),
4269 files: 0,
4270 commits: 0,
4271 empty: false,
4272 failed: None,
4273 duration_ms: 0,
4274 folded: false,
4275 }];
4276 state.tally = Some(Tally {
4277 first_choice: BTreeMap::from([('A', 1)]),
4278 borda: BTreeMap::new(),
4279 winner: 'A',
4280 rankings: 1,
4281 unanimous_initial: true,
4282 deliberated: false,
4283 changed_votes: 0,
4284 unanimous_final: true,
4285 tie_break: None,
4286 judges: 0,
4287 present: 0,
4288 quorum: 0,
4289 met_quorum: true,
4290 uncontested: Some("only candidate A produced a change".to_owned()),
4291 });
4292 state.reviews = vec![ReviewRound {
4293 round: 1,
4294 head: "deadbeef".to_owned(),
4295 verified_head: None,
4296 reviews: Vec::new(),
4297 e2e: Vec::new(),
4298 fix: None,
4299 blocking: 0,
4300 answered: 0,
4301 expected: 0,
4302 clean: true,
4303 verify_retried: false,
4304 e2e_deferred: false,
4305 e2e_defer_reason: None,
4306 progressed: false,
4307 vote_split: false,
4308 reconsideration: Vec::new(),
4309 verdict: None,
4310 }];
4311 state.gate = vec![CommandOutcome {
4312 command: "test".to_owned(),
4313 code: Some(0),
4314 output_tail: String::new(),
4315 duration_ms: 0,
4316 }];
4317 state.status = RunStatus::Ready;
4322 state.merge = Some(MergeOutcome {
4323 mode: MergeMode::Local,
4324 ok: false,
4325 detail: "already concluded".to_owned(),
4326 });
4327
4328 let mut runner = Runner {
4329 state,
4330 roles: ResolvedRoles {
4331 implementers: Vec::new(),
4332 judges: Vec::new(),
4333 reviewers: Vec::new(),
4334 fixer: None,
4335 conductor: conductor(),
4336 },
4337 sem: Arc::new(Semaphore::new(1)),
4338 pause: Pause::new(),
4339 };
4340
4341 runner.merge().await.expect("merge");
4342
4343 assert_eq!(
4344 runner.state.status,
4345 RunStatus::Ready,
4346 "a concluded run's status must not change on reentry"
4347 );
4348 assert_eq!(
4349 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
4350 Some("already concluded"),
4351 "merge must not run again once the node already recorded an outcome"
4352 );
4353 }
4354
4355 #[tokio::test]
4356 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
4357 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
4358 let tmp = tempfile::tempdir().expect("tempdir");
4359 let repo = tmp.path().join("repo");
4360 std::fs::create_dir_all(&repo).unwrap();
4361 init_repo(&repo);
4362
4363 let mut config = Config::default();
4364 config.merge.mode = MergeMode::Pr;
4365 config.graph.land = true;
4366 config.graph.land_approval = false;
4367
4368 let mut state = RunState::new(
4369 repo.clone(),
4370 "main".to_owned(),
4371 "deadbeef".to_owned(),
4372 "task".to_owned(),
4373 config,
4374 );
4375 state.candidates = vec![Candidate {
4376 index: 0,
4377 label: 'A',
4378 agent: "alpha".to_owned(),
4379 branch: "does-not-exist".to_owned(),
4380 worktree: repo.clone(),
4381 summary: String::new(),
4382 stat: String::new(),
4383 files: 0,
4384 commits: 0,
4385 empty: false,
4386 failed: None,
4387 duration_ms: 0,
4388 folded: false,
4389 }];
4390 state.tally = Some(Tally {
4391 first_choice: BTreeMap::from([('A', 1)]),
4392 borda: BTreeMap::new(),
4393 winner: 'A',
4394 rankings: 1,
4395 unanimous_initial: true,
4396 deliberated: false,
4397 changed_votes: 0,
4398 unanimous_final: true,
4399 tie_break: None,
4400 judges: 0,
4401 present: 0,
4402 quorum: 0,
4403 met_quorum: true,
4404 uncontested: Some("only candidate A produced a change".to_owned()),
4405 });
4406 state.reviews = vec![ReviewRound {
4407 round: 1,
4408 head: "deadbeef".to_owned(),
4409 verified_head: None,
4410 reviews: Vec::new(),
4411 e2e: Vec::new(),
4412 fix: None,
4413 blocking: 0,
4414 answered: 0,
4415 expected: 0,
4416 clean: true,
4417 verify_retried: false,
4418 e2e_deferred: false,
4419 e2e_defer_reason: None,
4420 progressed: false,
4421 vote_split: false,
4422 reconsideration: Vec::new(),
4423 verdict: None,
4424 }];
4425 state.gate = vec![CommandOutcome {
4426 command: "test".to_owned(),
4427 code: Some(0),
4428 output_tail: String::new(),
4429 duration_ms: 0,
4430 }];
4431 state.status = RunStatus::Landing;
4435 state.merge = Some(MergeOutcome {
4436 mode: MergeMode::Pr,
4437 ok: true,
4438 detail: "https://example.invalid/x/y/pull/1".to_owned(),
4439 });
4440
4441 ask_test_home();
4445 let store = ask::Questions::open();
4446 let q = ask_open_question(&store, &state.id);
4447
4448 let mut runner = Runner {
4449 state,
4450 roles: ResolvedRoles {
4451 implementers: Vec::new(),
4452 judges: Vec::new(),
4453 reviewers: Vec::new(),
4454 fixer: None,
4455 conductor: conductor(),
4456 },
4457 sem: Arc::new(Semaphore::new(1)),
4458 pause: Pause::new(),
4459 };
4460
4461 runner.execute().await.expect("execute");
4466
4467 assert_eq!(
4468 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
4469 Some("https://example.invalid/x/y/pull/1"),
4470 "reentry must not push again or open a second pull request over the \
4471 one `land` is already watching"
4472 );
4473 assert_ne!(
4474 runner.state.status,
4475 RunStatus::Landing,
4476 "land could not actually reach the fake pull request, so it must \
4477 have given up rather than left the run silently parked forever"
4478 );
4479 assert_eq!(runner.state.status, RunStatus::Blocked);
4483 assert!(
4484 store.get(&q.id).unwrap().status.open(),
4485 "Blocked is still alive; settle_questions must have been a no-op here"
4486 );
4487 }
4488
4489 fn state_with_round(round: ReviewRound) -> RunState {
4490 let mut s = RunState::new(
4491 PathBuf::from("/repo"),
4492 "main".to_owned(),
4493 "abc1234".to_owned(),
4494 "add retries".to_owned(),
4495 Config::default(),
4496 );
4497 s.reviews = vec![round];
4498 s
4499 }
4500
4501 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
4502 crate::verdict::Finding {
4503 id: id.to_owned(),
4504 severity,
4505 file: None,
4506 line: None,
4507 title: title.to_owned(),
4508 detail: String::new(),
4509 }
4510 }
4511
4512 #[test]
4513 fn pr_body_names_open_findings_and_declined_ones() {
4514 let round = ReviewRound {
4515 round: 2,
4516 head: "deadbee".to_owned(),
4517 verified_head: None,
4518 reviews: vec![ReviewRecord {
4519 reviewer: 1,
4520 agent: "alpha".to_owned(),
4521 summary: String::new(),
4522 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
4523 vote: None,
4524 failed: None,
4525 duration_ms: 0,
4526 }],
4527 e2e: vec![CommandOutcome {
4528 command: "cargo test".to_owned(),
4529 code: Some(0),
4530 output_tail: String::new(),
4531 duration_ms: 0,
4532 }],
4533 verify_retried: false,
4534 e2e_deferred: false,
4535 e2e_defer_reason: None,
4536 fix: Some(FixRecord {
4537 agent: "alpha".to_owned(),
4538 addressed: Vec::new(),
4539 rejected: vec![crate::verdict::Rejection {
4540 id: "R1-1-1".to_owned(),
4541 why: "not reachable from any caller".to_owned(),
4542 }],
4543 notes: String::new(),
4544 committed: true,
4545 failed: None,
4546 duration_ms: 0,
4547 }),
4548 blocking: 0,
4549 answered: 1,
4550 expected: 1,
4551 clean: false,
4552 progressed: true,
4553 vote_split: false,
4554 reconsideration: Vec::new(),
4555 verdict: None,
4556 };
4557 let state = state_with_round(round);
4558 let body = pr_body(&state, 'A');
4559
4560 assert!(body.contains("add retries"), "the task must still be there");
4561 assert!(body.contains("R2-1-1"), "{body}");
4562 assert!(body.contains("unused import"), "{body}");
4563 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
4564 assert!(
4565 body.contains("not reachable from any caller"),
4566 "the reason it was declined: {body}"
4567 );
4568 }
4569
4570 #[test]
4571 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
4572 let round = ReviewRound {
4573 round: 1,
4574 head: "deadbee".to_owned(),
4575 verified_head: None,
4576 reviews: vec![ReviewRecord {
4577 reviewer: 1,
4578 agent: "alpha".to_owned(),
4579 summary: String::new(),
4580 findings: Vec::new(),
4581 vote: None,
4582 failed: None,
4583 duration_ms: 0,
4584 }],
4585 e2e: Vec::new(),
4586 verify_retried: false,
4587 e2e_deferred: false,
4588 e2e_defer_reason: None,
4589 fix: None,
4590 blocking: 0,
4591 answered: 1,
4592 expected: 1,
4593 clean: true,
4594 progressed: false,
4595 vote_split: false,
4596 reconsideration: Vec::new(),
4597 verdict: None,
4598 };
4599 let state = state_with_round(round);
4600 let body = pr_body(&state, 'A');
4601 assert!(!body.contains("Open review findings"), "{body}");
4602 assert!(!body.contains("Declined"), "{body}");
4603 }
4604
4605 #[test]
4606 fn manual_merge_command_matches_the_configured_style() {
4607 let repo = Path::new("/repo");
4608 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
4609
4610 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
4611 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
4612
4613 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
4614 assert_eq!(
4615 squash,
4616 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
4617 \"Merge magi run 0832 (candidate A)\""
4618 );
4619
4620 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
4621 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
4622 }
4623
4624 #[test]
4625 fn a_nudge_gets_a_quarter_of_the_budget() {
4626 assert_eq!(retry_budget(secs(1200), true), secs(300));
4628 assert_eq!(retry_budget(secs(3600), true), secs(900));
4629 }
4630
4631 #[test]
4632 fn a_resent_prompt_keeps_the_whole_budget() {
4633 assert_eq!(retry_budget(secs(1200), false), secs(1200));
4636 assert_eq!(retry_budget(secs(60), false), secs(60));
4637 }
4638
4639 #[test]
4640 fn the_floor_never_exceeds_the_original_budget() {
4641 assert_eq!(retry_budget(secs(60), true), secs(60));
4645 assert_eq!(retry_budget(secs(480), true), secs(120));
4646 assert_eq!(retry_budget(secs(0), true), secs(0));
4647 }
4648}