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 mut e2e = run_commands(
2443 &shell,
2444 &self.state.config.verify.e2e,
2445 &winner.worktree,
2446 Duration::from_secs(self.state.config.graph.timeout_review),
2447 )
2448 .await;
2449 for o in &e2e {
2450 self.state.event(
2451 "verify",
2452 format!("round {round}: `{}` -> {}", o.command, e2e_outcome_label(o)),
2453 );
2454 }
2455
2456 let verify_retried = e2e.iter().any(CommandOutcome::build_failed);
2461 if verify_retried {
2462 self.state.event(
2463 "verify",
2464 format!(
2465 "round {round}: verify could not build/link, not a test result — \
2466 retrying once before concluding"
2467 ),
2468 );
2469 e2e = run_commands(
2470 &shell,
2471 &self.state.config.verify.e2e,
2472 &winner.worktree,
2473 Duration::from_secs(self.state.config.graph.timeout_review),
2474 )
2475 .await;
2476 for o in &e2e {
2477 self.state.event(
2478 "verify",
2479 format!(
2480 "round {round}: retry `{}` -> {}",
2481 o.command,
2482 e2e_outcome_label(o)
2483 ),
2484 );
2485 }
2486 }
2487
2488 let e2e_failures: String = e2e
2489 .iter()
2490 .filter(|o| !o.ok())
2491 .map(|o| format!("$ {}\n{}\n", o.command, o.output_tail))
2492 .collect();
2493
2494 let expected = records.len();
2495 let answered = records.iter().filter(|r| r.failed.is_none()).count();
2496 let incomplete = answered < expected;
2497 let blocking = all_findings.iter().filter(|f| f.severity.blocks()).count();
2498 let e2e_ok = e2e.iter().all(CommandOutcome::ok);
2499 let policy = self.state.config.graph.incomplete_review;
2500 let clean = round_is_clean(blocking, e2e_ok, answered, expected, policy);
2501
2502 let mut round_record = ReviewRound {
2503 round,
2504 head: head.clone(),
2505 reviews: records,
2506 e2e,
2507 verify_retried,
2508 fix: None,
2509 blocking,
2510 answered,
2511 expected,
2512 clean,
2513 progressed: false,
2514 vote_split,
2515 reconsideration,
2516 verdict: round_verdict,
2517 };
2518
2519 if incomplete {
2520 let missing: Vec<String> = round_record
2521 .reviews
2522 .iter()
2523 .filter(|r| r.failed.is_some())
2524 .map(|r| format!("review-{}", r.reviewer))
2525 .collect();
2526 self.state.event(
2527 "review",
2528 format!(
2529 "round {round}: {answered}/{expected} reviewer(s) answered ({} never answered)",
2530 missing.join(", ")
2531 ),
2532 );
2533 }
2534
2535 if clean {
2536 self.state.event(
2537 "review",
2538 if incomplete {
2539 format!(
2540 "round {round}: clean (warn policy, incomplete panel) — no \
2541 blocking findings from the seats that answered, verification green"
2542 )
2543 } else {
2544 format!("round {round}: clean — no blocking findings, verification green")
2545 },
2546 );
2547 self.state.reviews.push(round_record);
2548 self.state.status = RunStatus::Gating;
2549 self.state.save()?;
2550 return Ok(());
2551 }
2552
2553 if incomplete && blocking == 0 && e2e_ok {
2557 self.state.reviews.push(round_record);
2558 self.state.save()?;
2559 if round == max_rounds {
2560 self.state.status = RunStatus::Blocked;
2561 self.state.event(
2562 "review",
2563 format!(
2564 "{} reviewer seat(s) never answered after {max_rounds} rounds; \
2565 refusing to call it clean",
2566 expected - answered
2567 ),
2568 );
2569 return Ok(());
2570 }
2571 prev_e2e = None;
2572 continue;
2573 }
2574
2575 if round == max_rounds {
2576 self.state.reviews.push(round_record);
2577 return self.stop_reviewing(&format!(
2578 "{blocking} blocking finding(s) still open after {max_rounds} round(s)"
2579 ));
2580 }
2581
2582 let (fix_spec, fix_seat_key) = match &self.roles.fixer {
2585 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
2586 _ => (
2587 self.state
2588 .config
2589 .agent(&winner.agent)
2590 .cloned()
2591 .unwrap_or_else(|_| self.roles.implementers[winner.index].clone()),
2592 format!("impl-{}", winner.label),
2593 ),
2594 };
2595 let seat = self.seat(&fix_seat_key, &fix_spec.id);
2596 let blocking_findings: Vec<_> = all_findings
2597 .iter()
2598 .filter(|f| f.severity.blocks())
2599 .cloned()
2600 .collect();
2601 let job = SeatJob {
2602 prompt: prompt::fix(
2603 &self.state.instruction,
2604 &blocking_findings,
2605 (!e2e_failures.is_empty()).then_some(e2e_failures.as_str()),
2606 round,
2607 max_rounds,
2608 &language,
2609 ),
2610 spec: fix_spec.clone(),
2611 seat,
2612 cwd: winner.worktree.clone(),
2613 timeout: Duration::from_secs(self.state.config.graph.timeout_fix),
2614 allow_write: true,
2615 sessions,
2616 artifacts: artifacts.clone(),
2617 stem: format!("fix-{round}"),
2618 };
2619 let before = git::rev_parse(&winner.worktree, "HEAD").await?;
2620 let cache = self.state.config.cache_dir();
2621 let ctx = WaveCtx {
2622 run: &run_id,
2623 node: "fix",
2624 prompts: &prompts,
2625 cache: cache.as_deref(),
2626 };
2627 let (seat, out) = run_one(job, Arc::clone(&self.sem), &ctx, &mut self.state, 0).await;
2628 let agent_id = seat.agent.clone();
2629 let seat_key = seat.key.clone();
2630 self.state.seats.insert(seat.key.clone(), seat);
2631
2632 let mut fix = FixRecord {
2633 agent: agent_id,
2634 addressed: Vec::new(),
2635 rejected: Vec::new(),
2636 notes: String::new(),
2637 committed: false,
2638 failed: None,
2639 duration_ms: 0,
2640 };
2641 match out {
2642 AgentOutcome::Ok(o) => {
2643 fix.duration_ms = o.duration_ms;
2644 match verdict::extract_json::<FixReport>(&o.text) {
2645 Ok(report) => {
2646 fix.addressed = report.addressed;
2647 fix.rejected = report.rejected;
2648 fix.notes =
2649 blind::sanitize_prose(&report.notes, &self.state.config.blind);
2650 }
2651 Err(e) => fix.failed = Some(format!("unparsable fix report: {e}")),
2652 }
2653 }
2654 AgentOutcome::Dropped(o) => {
2656 fix.duration_ms = o.duration_ms;
2657 let why = o
2658 .dropped
2659 .as_ref()
2660 .map(|d| d.why.as_str())
2661 .unwrap_or("the CLI ended the stream without delivering its answer");
2662 fix.failed = Some(format!("the CLI dropped the stream ({why})"));
2663 }
2664 AgentOutcome::Quota(o) => {
2665 self.state.quota.push(QuotaLoss {
2666 seat: seat_key,
2667 node: "fix".to_owned(),
2668 at: Timestamp::now(),
2669 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
2670 });
2671 fix.failed = Some("rate limited (quota); fixer could not run".to_owned());
2672 }
2673 AgentOutcome::Failed(e) => fix.failed = Some(e),
2674 }
2675 git::commit_all(
2676 &winner.worktree,
2677 &format!("magi: review round {round} fixes (uncommitted work)"),
2678 )
2679 .await
2680 .ok();
2681 let after = git::rev_parse(&winner.worktree, "HEAD").await?;
2682 fix.committed = after != before;
2683 let diff_after = git::diff(&winner.worktree, &base, "HEAD").await?;
2691 let progressed = diff_after != patch;
2692 let commit_note = if fix.committed {
2693 "committed"
2694 } else {
2695 "NO new commit"
2696 };
2697 let tree_note = if progressed {
2698 "changed vs base"
2699 } else {
2700 "unchanged vs base"
2701 };
2702 self.state.event(
2703 "fix",
2704 match &fix.failed {
2705 Some(reason) => {
2711 format!(
2712 "round {round}: fixer's adoption report was lost ({reason}); \
2713 {commit_note}, tree {tree_note}"
2714 )
2715 }
2716 None => format!(
2717 "round {round}: {} addressed, {} rejected, {commit_note}, tree {tree_note}",
2718 fix.addressed.len(),
2719 fix.rejected.len(),
2720 ),
2721 },
2722 );
2723 round_record.fix = Some(fix);
2724 round_record.progressed = progressed;
2725 self.state.reviews.push(round_record);
2726 self.state.save()?;
2727
2728 prev_e2e = (!e2e_failures.is_empty()).then_some(e2e_failures);
2729
2730 let streak = self
2731 .state
2732 .reviews
2733 .iter()
2734 .rev()
2735 .take_while(|r| !r.progressed)
2736 .count();
2737 if streak >= STAGNANT_LIMIT {
2738 return self.stop_reviewing(&format!(
2739 "the tree has not moved against base for {streak} round(s) in a row"
2740 ));
2741 }
2742 }
2743 Ok(())
2744 }
2745
2746 fn stop_reviewing(&mut self, why: &str) -> Result<()> {
2762 let last = self
2763 .state
2764 .reviews
2765 .last()
2766 .expect("a round was just recorded before this is called");
2767 let red: Vec<String> = last
2768 .e2e
2769 .iter()
2770 .filter(|o| !o.ok())
2771 .map(|o| {
2772 format!(
2773 "`{}` -> {:?}\n{}",
2774 o.command,
2775 o.code,
2776 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
2777 )
2778 })
2779 .collect();
2780 let open: usize = last.reviews.iter().map(|r| r.findings.len()).sum();
2781
2782 if red.is_empty() {
2783 self.state.event(
2784 "review",
2785 format!("{why}; e2e is green — handing off with {open} finding(s) still open"),
2786 );
2787 self.state.status = RunStatus::Gating;
2788 } else {
2789 self.state
2790 .event("review", format!("{why}; e2e failed:\n{}", red.join("\n")));
2791 self.state.status = RunStatus::Blocked;
2792 }
2793 self.state.save()?;
2794 Ok(())
2795 }
2796
2797 async fn gate(&mut self) -> Result<()> {
2800 if self.state.status == RunStatus::Failed
2812 || self
2813 .state
2814 .base_sync
2815 .as_ref()
2816 .is_some_and(|s| s.conflict.is_some())
2817 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
2818 != Some(RunStatus::Gating)
2819 {
2820 return Ok(());
2821 }
2822 if !self.state.gate.is_empty() {
2823 return Ok(());
2824 }
2825 let Some(winner) = self.state.winner().cloned() else {
2826 return Ok(());
2827 };
2828 self.state.status = RunStatus::Gating;
2829 let shell = self.state.config.shell();
2830 let outcomes = run_commands(
2831 &shell,
2832 &self.state.config.verify.gate,
2833 &winner.worktree,
2834 Duration::from_secs(self.state.config.graph.timeout_review),
2835 )
2836 .await;
2837 for o in &outcomes {
2838 self.state.event(
2839 "gate",
2840 format!(
2841 "`{}` -> {}",
2842 o.command,
2843 if o.ok() {
2844 "pass".to_owned()
2845 } else {
2846 format!(
2847 "FAIL ({:?})\n{}",
2848 o.code,
2849 tail(&o.output_tail, EVENT_OUTPUT_TAIL)
2850 )
2851 }
2852 ),
2853 );
2854 }
2855 let passed = outcomes.iter().all(CommandOutcome::ok);
2856 self.state.gate = outcomes;
2857 if !passed {
2858 self.state.status = RunStatus::Blocked;
2859 self.state.event("gate", "gate failed; not merging");
2860 }
2861 self.state.save()?;
2862 Ok(())
2863 }
2864
2865 async fn merge(&mut self) -> Result<()> {
2868 if self
2883 .state
2884 .base_sync
2885 .as_ref()
2886 .is_some_and(|s| s.conflict.is_some())
2887 || review_conclusion(&self.state.reviews, self.state.config.graph.review_rounds)
2888 != Some(RunStatus::Gating)
2889 || self.state.gate.iter().any(|o| !o.ok())
2890 {
2891 return Ok(());
2892 }
2893 if self.state.merge.is_some() {
2902 return Ok(());
2903 }
2904 let Some(winner) = self.state.winner().cloned() else {
2905 return Ok(());
2906 };
2907 let repo = self.state.repo.clone();
2908 let base = self.state.base_branch.clone();
2909 let mode = self.state.config.merge.mode;
2910 let style = self.state.config.merge.style;
2911 let message = pr_body(&self.state, winner.label);
2912
2913 let outcome = match mode {
2914 MergeMode::None => MergeOutcome {
2915 mode,
2916 ok: true,
2917 detail: manual_merge_command(style, &repo, &winner.branch, &message),
2918 },
2919 MergeMode::Local => {
2920 let on = git::current_branch(&repo).await?;
2921 if on.as_deref() != Some(base.as_str()) {
2922 MergeOutcome {
2923 mode,
2924 ok: false,
2925 detail: format!(
2926 "{} has {} checked out, not the base branch {base}",
2927 repo.display(),
2928 on.unwrap_or_else(|| "a detached HEAD".to_owned())
2929 ),
2930 }
2931 } else if !git::is_clean(&repo).await? {
2932 MergeOutcome {
2933 mode,
2934 ok: false,
2935 detail: format!("{} is dirty; refusing to merge", repo.display()),
2936 }
2937 } else {
2938 let out = match style {
2939 MergeStyle::Merge => {
2940 git::merge_no_ff(&repo, &winner.branch, &message).await?
2941 }
2942 MergeStyle::Squash => {
2943 git::merge_squash(&repo, &winner.branch, &message).await?
2944 }
2945 MergeStyle::Rebase => git::merge_ff_only(&repo, &winner.branch).await?,
2946 };
2947 MergeOutcome {
2948 mode,
2949 ok: out.ok(),
2950 detail: if out.ok() { out.stdout } else { out.stderr },
2951 }
2952 }
2953 }
2954 MergeMode::Pr => {
2955 let remote = self.state.config.merge.remote.clone();
2956 let pushed = git::push(&winner.worktree, &remote, &winner.branch).await?;
2957 if !pushed.ok() {
2958 MergeOutcome {
2959 mode,
2960 ok: false,
2961 detail: pushed.stderr,
2962 }
2963 } else {
2964 let out = gh_pr_create(&winner.worktree, &base, &winner.branch, &message).await;
2965 match out {
2966 Ok(url) => MergeOutcome {
2967 mode,
2968 ok: true,
2969 detail: url,
2970 },
2971 Err(e) => MergeOutcome {
2972 mode,
2973 ok: false,
2974 detail: e.to_string(),
2975 },
2976 }
2977 }
2978 }
2979 };
2980
2981 self.state.status = match (mode, outcome.ok) {
2982 (MergeMode::None, _) => RunStatus::Ready,
2983 (_, true) => RunStatus::Merged,
2984 (_, false) => RunStatus::Blocked,
2985 };
2986 self.state.event(
2987 "merge",
2988 format!(
2989 "{:?}: {}",
2990 mode,
2991 outcome.detail.lines().next().unwrap_or("")
2992 ),
2993 );
2994 self.state.merge = Some(outcome);
2995 self.state.save()?;
2996
2997 if self.state.config.graph.land
3003 && mode == MergeMode::Pr
3004 && self.state.status == RunStatus::Merged
3005 {
3006 self.run_land().await?;
3007 }
3008 self.settle_questions();
3013 Ok(())
3014 }
3015
3016 async fn run_land(&mut self) -> Result<()> {
3027 let url = self
3028 .state
3029 .merge
3030 .as_ref()
3031 .map(|m| m.detail.clone())
3032 .unwrap_or_default();
3033 let url = url.lines().next().unwrap_or("").trim().to_owned();
3034 if !url.starts_with("http") {
3035 return Ok(());
3036 }
3037 match land::land(&mut self.state, &url).await {
3040 Ok(pr) if self.state.parked => {
3041 let _ = pr;
3045 }
3046 Ok(pr) => {
3047 self.state.status = match pr.state {
3048 land::PrLifecycle::Merged => RunStatus::Merged,
3049 _ => RunStatus::Blocked,
3050 };
3051 if bump::should_release_bump(self.state.status)
3058 && let Err(e) = bump::after_merge(&mut self.state, &pr.url).await
3059 {
3060 self.state
3061 .event("bump", format!("release bump skipped: {e:#}"));
3062 }
3063 self.state.save()?;
3064 }
3065 Err(e) => {
3066 self.state.status = RunStatus::Blocked;
3067 self.state.event("land", format!("gave up: {e}"));
3068 self.state.save()?;
3069 }
3070 }
3071 Ok(())
3072 }
3073
3074 fn seat(&mut self, key: &str, agent: &str) -> SeatState {
3078 if let Some(existing) = self.state.seats.get(key)
3079 && existing.agent == agent
3080 {
3081 return existing.clone();
3082 }
3083 let fresh = SeatState::new(key, agent, self.state.seed);
3084 self.state.seats.insert(key.to_owned(), fresh.clone());
3085 fresh
3086 }
3087
3088 fn view(&self, c: &Candidate) -> CandidateView {
3090 let raw = crate::run::read_artifact(&self.state, &format!("cand-{}.patch", c.label))
3091 .unwrap_or_default();
3092 let (patch, _) = blind::sanitize_patch(
3093 &format!("candidate {} patch", c.label),
3094 &raw,
3095 &self.state.config.blind,
3096 );
3097 CandidateView {
3098 label: c.label,
3099 branch: c.branch.clone(),
3100 summary: c.summary.clone(),
3101 stat: c.stat.clone(),
3102 patch,
3103 }
3104 }
3105
3106 fn candidate_block(&self, candidates: &[Candidate], base_short: &str) -> String {
3108 let views: Vec<CandidateView> = candidates.iter().map(|c| self.view(c)).collect();
3109 prompt::judge(
3110 "(see above)",
3111 &views,
3112 self.roles.judges.len(),
3113 base_short,
3114 "en",
3115 )
3116 }
3117
3118 fn transcript(&self, current: &[DeliberationTurn], self_idx: usize) -> Vec<Turn> {
3125 let mut turns = Vec::new();
3126 for j in &self.state.judgements {
3127 if j.ranking.is_empty() {
3128 continue;
3129 }
3130 let reasons = j
3131 .reasons
3132 .iter()
3133 .map(|(k, v)| format!("- {k}: {v}"))
3134 .collect::<Vec<_>>()
3135 .join("\n");
3136 turns.push(Turn {
3137 who: format!("Judge {} (opening ranking)", j.judge),
3138 is_self: j.judge == self_idx + 1,
3139 body: format!(
3140 "Ranked {}{}{reasons}",
3141 j.ranking.iter().collect::<String>(),
3142 if reasons.is_empty() {
3143 ""
3144 } else {
3145 ", because:\n"
3146 }
3147 ),
3148 });
3149 }
3150 for t in self
3151 .state
3152 .deliberation
3153 .iter()
3154 .flat_map(|r| r.turns.iter())
3155 .chain(current)
3156 {
3157 turns.push(Turn {
3158 who: format!("Judge {}", t.judge),
3159 is_self: t.judge == self_idx + 1,
3160 body: t.body.clone(),
3161 });
3162 }
3163 turns
3164 }
3165}
3166
3167fn has_context(spec: &AgentSpec, seat: &SeatState, sessions: bool) -> bool {
3169 agent::has_session(spec.kind, seat, sessions)
3170}
3171
3172fn short(commit: &str) -> String {
3173 commit.chars().take(7).collect()
3174}
3175
3176fn make_executable(path: &Path) -> Result<()> {
3177 #[cfg(unix)]
3178 {
3179 use std::os::unix::fs::PermissionsExt as _;
3180 let mut perms = std::fs::metadata(path)?.permissions();
3181 perms.set_mode(0o755);
3182 std::fs::set_permissions(path, perms)?;
3183 }
3184 #[cfg(not(unix))]
3185 {
3186 let _ = path;
3187 }
3188 Ok(())
3189}
3190
3191struct WaveCtx<'a> {
3198 run: &'a str,
3201 node: &'a str,
3203 prompts: &'a Prompts,
3204 cache: Option<&'a Path>,
3206}
3207
3208async fn run_one(
3210 job: SeatJob,
3211 sem: Arc<Semaphore>,
3212 ctx: &WaveCtx<'_>,
3213 state: &mut RunState,
3214 attempt: usize,
3215) -> (SeatState, AgentOutcome) {
3216 let (_, seat, out) = wave(vec![job], sem, ctx, state, attempt)
3217 .await
3218 .pop()
3219 .expect("one job in, one result out");
3220 (seat, out)
3221}
3222
3223async fn wave(
3229 jobs: Vec<SeatJob>,
3230 sem: Arc<Semaphore>,
3231 ctx: &WaveCtx<'_>,
3232 state: &mut RunState,
3233 attempt: usize,
3234) -> Vec<(usize, SeatState, AgentOutcome)> {
3235 let WaveCtx {
3236 run,
3237 node,
3238 prompts,
3239 cache,
3240 } = *ctx;
3241 for job in &jobs {
3242 state.seat_started(node, &job.seat.key, job.timeout, attempt);
3243 }
3244 if let Err(e) = state.save() {
3245 tracing::warn!("could not persist in-progress seats: {e:#}");
3250 }
3251 let mut set = tokio::task::JoinSet::new();
3252 let overlay = prompts.overlay(node);
3253 for (i, mut job) in jobs.into_iter().enumerate() {
3254 job.prompt = prompt::with_overlay(job.prompt, overlay.clone());
3255 if cache.is_some() {
3256 job.prompt.push('\n');
3257 job.prompt.push_str(prompt::build_cache_note());
3258 }
3259 let sem = Arc::clone(&sem);
3260 let run = run.to_owned();
3261 let node = node.to_owned();
3262 let cache = cache.map(Path::to_path_buf);
3263 set.spawn(async move {
3264 let _permit = sem.acquire().await;
3265 let mut seat = job.seat;
3266 let out = agent::invoke(
3267 &job.spec,
3268 &mut seat,
3269 &Invocation {
3270 cwd: &job.cwd,
3271 prompt: &job.prompt,
3272 timeout: job.timeout,
3273 allow_write: job.allow_write,
3274 sessions: job.sessions,
3275 artifacts: &job.artifacts,
3276 stem: &job.stem,
3277 run: &run,
3278 node: &node,
3279 cache_dir: cache.as_deref(),
3280 attachments: &[],
3281 },
3282 )
3283 .await;
3284 let out = match out {
3285 Ok(o) if o.usable() => AgentOutcome::Ok(o),
3286 Ok(o) if o.quota_exhausted() => AgentOutcome::Quota(o),
3287 Ok(o) if o.work_undelivered() => AgentOutcome::Dropped(o),
3295 Ok(o) if o.timed_out => AgentOutcome::Failed("timed out".to_owned()),
3296 Ok(o) => AgentOutcome::Failed(format!(
3297 "exited with {:?} and no usable output",
3298 o.exit_code
3299 )),
3300 Err(e) => AgentOutcome::Failed(e.to_string()),
3301 };
3302 (i, seat, out)
3303 });
3304 }
3305 let mut collected: Vec<Option<(usize, SeatState, AgentOutcome)>> = Vec::new();
3306 while let Some(joined) = set.join_next().await {
3307 let (i, seat, out) = match joined {
3308 Ok(v) => v,
3309 Err(e) => {
3313 tracing::error!("agent task panicked: {e}");
3314 continue;
3315 }
3316 };
3317 state.seat_finished(&seat.key);
3318 if let Err(e) = state.save() {
3319 tracing::warn!("could not persist a seat's completion: {e:#}");
3320 }
3321 if collected.len() <= i {
3322 collected.resize_with(i + 1, || None);
3323 }
3324 collected[i] = Some((i, seat, out));
3325 }
3326 if state
3332 .active
3333 .values()
3334 .any(|a| a.node == node && a.attempt == attempt)
3335 {
3336 state
3337 .active
3338 .retain(|_, a| !(a.node == node && a.attempt == attempt));
3339 if let Err(e) = state.save() {
3340 tracing::warn!("could not persist the end of a wave: {e:#}");
3341 }
3342 }
3343 collected.into_iter().flatten().collect()
3344}
3345
3346fn round_is_clean(
3356 blocking: usize,
3357 e2e_ok: bool,
3358 answered: usize,
3359 expected: usize,
3360 policy: IncompleteReviewPolicy,
3361) -> bool {
3362 blocking == 0 && e2e_ok && (answered == expected || policy == IncompleteReviewPolicy::Warn)
3363}
3364
3365fn review_conclusion(reviews: &[ReviewRound], max_rounds: usize) -> Option<RunStatus> {
3382 if max_rounds == 0 || reviews.iter().any(|r| r.clean) {
3383 return Some(RunStatus::Gating);
3384 }
3385 let last = reviews.last()?;
3386 let stagnant = reviews.iter().rev().take_while(|r| !r.progressed).count() >= STAGNANT_LIMIT;
3387 if reviews.len() < max_rounds && !stagnant {
3388 return None;
3389 }
3390 Some(if last.incomplete() && last.blocking == 0 {
3391 RunStatus::Blocked
3392 } else if last.e2e.iter().all(CommandOutcome::ok) {
3393 RunStatus::Gating
3394 } else {
3395 RunStatus::Blocked
3396 })
3397}
3398
3399fn retry_budget(full: Duration, nudged: bool) -> Duration {
3414 if nudged {
3415 (full / 4).max(Duration::from_secs(120)).min(full)
3416 } else {
3417 full
3418 }
3419}
3420
3421#[allow(clippy::too_many_arguments)]
3434async fn ask_json_wave<T>(
3435 jobs: Vec<SeatJob>,
3436 sem: Arc<Semaphore>,
3437 retries: usize,
3438 ctx: &WaveCtx<'_>,
3439 losses: &mut Vec<QuotaLoss>,
3440 state: &mut RunState,
3441 validate: &(dyn Fn(&T) -> Result<()> + Send + Sync),
3442) -> Vec<(SeatState, Result<(T, AgentOutput)>)>
3443where
3444 T: serde::de::DeserializeOwned + Send + 'static,
3445{
3446 let n = jobs.len();
3447 let originals: Vec<SeatJob> = jobs;
3448 let mut seats: Vec<SeatState> = originals.iter().map(|j| j.seat.clone()).collect();
3449 let mut done: Vec<Option<Result<(T, AgentOutput)>>> = (0..n).map(|_| None).collect();
3450 let mut pending: Vec<usize> = (0..n).collect();
3451
3452 for attempt in 0..=retries {
3453 if pending.is_empty() {
3454 break;
3455 }
3456 let mut batch = Vec::with_capacity(pending.len());
3457 for &i in &pending {
3458 let src = &originals[i];
3459 let (prompt, timeout) = if attempt == 0 {
3462 (src.prompt.clone(), src.timeout)
3463 } else {
3464 let why = done[i]
3465 .as_ref()
3466 .and_then(|r| r.as_ref().err().map(ToString::to_string))
3467 .unwrap_or_else(|| "no parsable answer".to_owned());
3468 let nudge = prompt::nudge(&why);
3469 let nudged = has_context(&src.spec, &seats[i], src.sessions);
3470 let prompt = if nudged {
3471 nudge
3472 } else {
3473 format!("{}\n\n---\n\n{}", src.prompt, nudge)
3474 };
3475 (prompt, retry_budget(src.timeout, nudged))
3476 };
3477 batch.push(SeatJob {
3478 spec: src.spec.clone(),
3479 seat: seats[i].clone(),
3480 cwd: src.cwd.clone(),
3481 prompt,
3482 timeout,
3483 allow_write: src.allow_write,
3484 sessions: src.sessions,
3485 artifacts: src.artifacts.clone(),
3486 stem: if attempt == 0 {
3487 src.stem.clone()
3488 } else {
3489 format!("{}-retry{attempt}", src.stem)
3490 },
3491 });
3492 }
3493
3494 if attempt > 0 {
3495 let seats_out: Vec<&str> = pending
3496 .iter()
3497 .map(|&i| originals[i].seat.key.as_str())
3498 .collect();
3499 state.event(
3500 ctx.node,
3501 format!("retry {attempt}: re-asking {}", seats_out.join(", ")),
3502 );
3503 }
3504 let results = wave(batch, Arc::clone(&sem), ctx, state, attempt).await;
3505 let mut still = Vec::new();
3506 for (&i, (_wi, seat, out)) in pending.iter().zip(results) {
3507 seats[i] = seat;
3508 let (parsed, quota) = match out {
3509 AgentOutcome::Ok(o) => (
3510 match verdict::extract_json::<T>(&o.text) {
3511 Ok(v) => match validate(&v) {
3512 Ok(()) => Ok((v, o)),
3513 Err(e) => Err(e),
3514 },
3515 Err(e) => Err(e),
3516 },
3517 false,
3518 ),
3519 AgentOutcome::Quota(o) => {
3520 losses.push(QuotaLoss {
3521 seat: originals[i].seat.key.clone(),
3522 node: ctx.node.to_owned(),
3523 at: Timestamp::now(),
3524 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
3525 });
3526 (
3527 Err(anyhow::anyhow!("rate limited (quota); not retrying now")),
3528 true,
3529 )
3530 }
3531 AgentOutcome::Dropped(o) => {
3536 let why = o
3537 .dropped
3538 .as_ref()
3539 .map(|d| d.why.as_str())
3540 .unwrap_or("the CLI ended the stream without delivering its answer");
3541 (
3542 Err(anyhow::anyhow!("the CLI dropped the stream ({why})")),
3543 false,
3544 )
3545 }
3546 AgentOutcome::Failed(e) => (Err(anyhow::anyhow!(e)), false),
3547 };
3548 let failed = parsed.is_err();
3549 done[i] = Some(parsed);
3550 if failed && !quota {
3553 still.push(i);
3554 }
3555 }
3556 pending = still;
3557 }
3558
3559 seats
3560 .into_iter()
3561 .zip(done)
3562 .map(|(seat, res)| {
3563 (
3564 seat,
3565 res.unwrap_or_else(|| Err(anyhow::anyhow!("no attempt was made"))),
3566 )
3567 })
3568 .collect()
3569}
3570
3571fn e2e_outcome_label(o: &CommandOutcome) -> String {
3575 if o.ok() {
3576 return "pass".to_owned();
3577 }
3578 let reason = if o.build_failed() {
3579 format!("COULD NOT RUN ({:?}, build/link failure)", o.code)
3580 } else {
3581 format!("FAIL ({:?})", o.code)
3582 };
3583 format!("{reason}\n{}", tail(&o.output_tail, EVENT_OUTPUT_TAIL))
3584}
3585
3586async fn run_commands(
3588 shell: &[String],
3589 commands: &[String],
3590 cwd: &Path,
3591 timeout: Duration,
3592) -> Vec<CommandOutcome> {
3593 let mut out = Vec::new();
3594 for command in commands {
3595 let started = Instant::now();
3596 let mut cmd = tokio::process::Command::new(&shell[0]);
3597 cmd.quiet();
3598 cmd.args(&shell[1..])
3599 .arg(command)
3600 .current_dir(cwd)
3601 .stdin(std::process::Stdio::null())
3602 .stdout(std::process::Stdio::piped())
3603 .stderr(std::process::Stdio::piped())
3604 .kill_on_drop(true);
3605 let spawned = cmd.spawn();
3606 let (code, body) = match spawned {
3607 Ok(child) => match tokio::time::timeout(timeout, child.wait_with_output()).await {
3608 Ok(Ok(o)) => {
3609 let mut body = String::from_utf8_lossy(&o.stdout).into_owned();
3610 body.push_str(&String::from_utf8_lossy(&o.stderr));
3611 (o.status.code(), body)
3612 }
3613 Ok(Err(e)) => (None, format!("failed to run: {e}")),
3614 Err(_) => (None, format!("timed out after {}s", timeout.as_secs())),
3615 },
3616 Err(e) => (None, format!("failed to spawn `{}`: {e}", shell[0])),
3617 };
3618 out.push(CommandOutcome {
3619 command: command.clone(),
3620 code,
3621 output_tail: tail(&body, OUTPUT_TAIL),
3622 duration_ms: started.elapsed().as_millis() as u64,
3623 });
3624 }
3625 out
3626}
3627
3628fn manual_merge_command(style: MergeStyle, repo: &Path, branch: &str, message: &str) -> String {
3640 let repo = repo.display();
3641 match style {
3642 MergeStyle::Merge => format!("git -C {repo} merge --no-ff {branch}"),
3643 MergeStyle::Squash => {
3644 let subject = message.lines().next().unwrap_or(branch);
3645 format!(
3646 "git -C {repo} merge --squash {branch} && git -C {repo} commit -m \"{subject}\""
3647 )
3648 }
3649 MergeStyle::Rebase => format!("git -C {repo} merge --ff-only {branch}"),
3650 }
3651}
3652
3653fn pr_body(state: &RunState, winner: char) -> String {
3659 let mut message = format!(
3660 "Merge magi run {} (candidate {winner})\n\n{}",
3661 state.id, state.instruction
3662 );
3663
3664 let open = state.open_findings();
3665 if !open.is_empty() {
3666 message.push_str("\n\n## Open review findings\n\n");
3667 for f in &open {
3668 message.push_str(&format!("- `{}` [{:?}] {}\n", f.id, f.severity, f.title));
3669 }
3670 }
3671
3672 if let Some(fix) = state.reviews.last().and_then(|r| r.fix.as_ref())
3673 && !fix.rejected.is_empty()
3674 {
3675 message.push_str("\n## Declined by the fixer\n\n");
3676 for r in &fix.rejected {
3677 message.push_str(&format!("- `{}`: {}\n", r.id, r.why));
3678 }
3679 }
3680
3681 message
3682}
3683
3684async fn gh_pr_create(cwd: &Path, base: &str, head: &str, body: &str) -> Result<String> {
3686 let title = body.lines().next().unwrap_or("magi run").to_owned();
3687 let out = tokio::process::Command::new("gh")
3688 .args([
3689 "pr", "create", "--base", base, "--head", head, "--title", &title, "--body", body,
3690 ])
3691 .current_dir(cwd)
3692 .quiet()
3693 .stdin(std::process::Stdio::null())
3694 .output()
3695 .await
3696 .context("spawn gh")?;
3697 if out.status.success() {
3698 Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned())
3699 } else {
3700 bail!("{}", String::from_utf8_lossy(&out.stderr).trim().to_owned())
3701 }
3702}
3703
3704pub async fn fold_run(state: &mut RunState, drop_winner: bool) -> Result<Vec<String>> {
3706 let repo = state.repo.clone();
3707 let root = state.worktree_root();
3708 let winner = state.tally.as_ref().map(|t| t.winner);
3709 let mut removed = Vec::new();
3710
3711 for i in 0..state.candidates.len() {
3712 let c = state.candidates[i].clone();
3713 let is_winner = Some(c.label) == winner;
3714 if is_winner && !drop_winner {
3715 continue;
3716 }
3717 if c.worktree.exists() {
3718 git::worktree_remove(&repo, &c.worktree).await.ok();
3719 removed.push(c.worktree.to_string_lossy().into_owned());
3720 }
3721 if git::branch_exists(&repo, &c.branch).await.unwrap_or(false) {
3722 git::branch_delete(&repo, &c.branch).await.ok();
3723 removed.push(c.branch.clone());
3724 }
3725 state.candidates[i].folded = true;
3726 }
3727
3728 for name in std::fs::read_dir(&root).into_iter().flatten().flatten() {
3729 let path = name.path();
3730 let keep = !drop_winner
3731 && winner.is_some_and(|w| {
3732 path.file_name()
3733 .is_some_and(|n| n == format!("cand-{w}").as_str())
3734 });
3735 if keep {
3736 continue;
3737 }
3738 git::worktree_remove(&repo, &path).await.ok();
3739 removed.push(path.to_string_lossy().into_owned());
3740 }
3741
3742 remove_if_empty(&root);
3751
3752 if state.enabled_worktree_config && drop_winner {
3753 git::release_worktree_config(&repo).await.ok();
3757 state.enabled_worktree_config = false;
3758 }
3759 state.save()?;
3760 Ok(removed)
3761}
3762
3763fn remove_if_empty(dir: &Path) {
3774 if dir.is_dir() && std::fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) {
3775 std::fs::remove_dir(dir).ok();
3776 }
3777}
3778
3779pub fn worst_open(state: &RunState) -> Option<Severity> {
3781 state
3782 .reviews
3783 .last()?
3784 .reviews
3785 .iter()
3786 .flat_map(|r| r.findings.iter())
3787 .map(|f| f.severity)
3788 .max()
3789}
3790
3791#[cfg(test)]
3792mod tests {
3793 use super::*;
3794 use std::time::Duration;
3795
3796 #[test]
3797 fn remove_if_empty_only_ever_takes_a_bare_directory() {
3798 let dir = tempfile::tempdir().unwrap();
3799 let bay = dir.path().join("ffff");
3800
3801 remove_if_empty(&bay);
3803 assert!(!bay.exists());
3804
3805 std::fs::create_dir_all(bay.join("cand-A")).unwrap();
3808 remove_if_empty(&bay);
3809 assert!(bay.exists(), "non-empty directory must survive");
3810
3811 std::fs::remove_dir(bay.join("cand-A")).unwrap();
3813 remove_if_empty(&bay);
3814 assert!(!bay.exists(), "an empty bay is a leftover, not a record");
3815 }
3816
3817 #[test]
3826 fn a_full_panel_that_found_nothing_is_clean() {
3827 assert!(round_is_clean(0, true, 2, 2, IncompleteReviewPolicy::Block));
3828 }
3829
3830 #[test]
3831 fn a_missing_seat_is_never_clean_under_the_default_policy() {
3832 assert!(!round_is_clean(
3833 0,
3834 true,
3835 1,
3836 2,
3837 IncompleteReviewPolicy::Block
3838 ));
3839 }
3840
3841 #[test]
3842 fn warn_policy_still_refuses_a_missing_seat_with_open_findings() {
3843 assert!(!round_is_clean(1, true, 1, 2, IncompleteReviewPolicy::Warn));
3844 }
3845
3846 #[test]
3847 fn warn_policy_gates_a_missing_seat_once_what_answered_is_clean() {
3848 assert!(round_is_clean(0, true, 1, 2, IncompleteReviewPolicy::Warn));
3849 }
3850
3851 #[test]
3852 fn a_full_panel_with_an_open_finding_is_not_clean() {
3853 assert!(!round_is_clean(
3854 1,
3855 true,
3856 2,
3857 2,
3858 IncompleteReviewPolicy::Block
3859 ));
3860 }
3861
3862 #[test]
3863 fn a_full_panel_with_a_red_e2e_is_not_clean() {
3864 assert!(!round_is_clean(
3865 0,
3866 false,
3867 2,
3868 2,
3869 IncompleteReviewPolicy::Block
3870 ));
3871 }
3872
3873 fn review_round(
3879 clean: bool,
3880 blocking: usize,
3881 answered: usize,
3882 expected: usize,
3883 progressed: bool,
3884 e2e_ok: bool,
3885 ) -> ReviewRound {
3886 ReviewRound {
3887 round: 1,
3888 head: "h".to_owned(),
3889 reviews: Vec::new(),
3890 e2e: vec![CommandOutcome {
3891 command: "test".to_owned(),
3892 code: Some(if e2e_ok { 0 } else { 1 }),
3893 output_tail: String::new(),
3894 duration_ms: 0,
3895 }],
3896 verify_retried: false,
3897 fix: None,
3898 blocking,
3899 answered,
3900 expected,
3901 clean,
3902 progressed,
3903 vote_split: false,
3904 reconsideration: Vec::new(),
3905 verdict: None,
3906 }
3907 }
3908
3909 #[test]
3910 fn review_conclusion_is_none_when_nothing_has_run() {
3911 assert_eq!(review_conclusion(&[], 3), None);
3912 }
3913
3914 #[test]
3915 fn review_conclusion_is_none_while_rounds_remain() {
3916 let rounds = vec![review_round(false, 1, 2, 2, true, true)];
3917 assert_eq!(review_conclusion(&rounds, 3), None);
3918 }
3919
3920 #[test]
3921 fn review_conclusion_is_gating_once_a_round_is_clean() {
3922 let rounds = vec![review_round(true, 0, 2, 2, false, true)];
3923 assert_eq!(review_conclusion(&rounds, 3), Some(RunStatus::Gating));
3924 }
3925
3926 #[test]
3927 fn review_conclusion_hands_off_when_the_budget_is_spent_and_e2e_is_green() {
3928 let rounds = vec![
3929 review_round(false, 1, 2, 2, true, true),
3930 review_round(false, 1, 2, 2, true, true),
3931 ];
3932 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Gating));
3933 }
3934
3935 #[test]
3936 fn review_conclusion_blocks_when_the_budget_is_spent_and_e2e_is_red() {
3937 let rounds = vec![
3938 review_round(false, 1, 2, 2, true, true),
3939 review_round(false, 1, 2, 2, true, false),
3940 ];
3941 assert_eq!(review_conclusion(&rounds, 2), Some(RunStatus::Blocked));
3942 }
3943
3944 #[test]
3945 fn review_conclusion_blocks_an_incomplete_panel_that_raised_nothing_even_with_green_e2e() {
3946 let rounds = vec![review_round(false, 0, 1, 2, false, true)];
3948 assert_eq!(review_conclusion(&rounds, 1), Some(RunStatus::Blocked));
3949 }
3950
3951 #[test]
3952 fn review_conclusion_hands_off_when_the_tree_stagnates_before_the_budget_is_spent() {
3953 let rounds = vec![
3954 review_round(false, 1, 2, 2, false, true),
3955 review_round(false, 1, 2, 2, false, true),
3956 ];
3957 assert_eq!(review_conclusion(&rounds, 10), Some(RunStatus::Gating));
3958 }
3959
3960 fn secs(n: u64) -> Duration {
3961 Duration::from_secs(n)
3962 }
3963
3964 fn init_repo(dir: &Path) {
3967 let run = |args: &[&str]| {
3968 let out = std::process::Command::new("git")
3969 .args(args)
3970 .current_dir(dir)
3971 .quiet()
3972 .output()
3973 .expect("spawn git");
3974 assert!(
3975 out.status.success(),
3976 "git {args:?} failed: {}",
3977 String::from_utf8_lossy(&out.stderr)
3978 );
3979 };
3980 run(&["init", "-b", "main"]);
3981 run(&["config", "user.name", "magi test"]);
3982 run(&["config", "user.email", "magi@example.com"]);
3983 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
3984 run(&["add", "-A"]);
3985 run(&["commit", "-m", "init"]);
3986 }
3987
3988 fn ask_test_home() {
3996 crate::run::set_home(std::env::temp_dir().join("magi-graph-ask-tests-home"));
3997 }
3998
3999 fn runner_at(status: RunStatus) -> Runner {
4002 let mut state = RunState::new(
4003 PathBuf::from("/nonexistent/repo"),
4004 "main".to_owned(),
4005 "deadbeef".to_owned(),
4006 "task".to_owned(),
4007 Config::default(),
4008 );
4009 state.status = status;
4010 Runner {
4011 state,
4012 roles: ResolvedRoles {
4013 implementers: Vec::new(),
4014 judges: Vec::new(),
4015 reviewers: Vec::new(),
4016 fixer: None,
4017 },
4018 sem: Arc::new(Semaphore::new(1)),
4019 pause: Pause::new(),
4020 }
4021 }
4022
4023 fn ask_open_question(store: &ask::Questions, run: &str) -> ask::Question {
4025 let mut q = ask::Question::new(
4026 run.to_owned(),
4027 "implement".to_owned(),
4028 "impl-A".to_owned(),
4029 "Which storage backend should the cache use?".to_owned(),
4030 String::new(),
4031 vec!["SQLite".to_owned(), "Redis".to_owned()],
4032 );
4033 store.put(&mut q).unwrap();
4034 q
4035 }
4036
4037 #[test]
4038 fn a_failed_runs_open_question_is_abandoned() {
4039 ask_test_home();
4040 let store = ask::Questions::open();
4041 let mut runner = runner_at(RunStatus::Failed);
4042 let run = runner.state.id.clone();
4043 let q = ask_open_question(&store, &run);
4044
4045 runner.settle_questions();
4046
4047 let back = store.get(&q.id).unwrap();
4048 assert!(
4049 !back.status.open(),
4050 "the seat that asked died with the run; nobody is left to read an answer"
4051 );
4052 assert!(
4053 back.detail.contains(&run) && back.detail.contains("failed"),
4054 "the reason names what the run became, not just that it is gone: {}",
4055 back.detail
4056 );
4057 }
4058
4059 #[test]
4060 fn a_merged_runs_open_question_is_abandoned_too() {
4061 ask_test_home();
4062 let store = ask::Questions::open();
4063 for status in [RunStatus::Merged, RunStatus::Ready] {
4066 let mut runner = runner_at(status);
4067 let run = runner.state.id.clone();
4068 let q = ask_open_question(&store, &run);
4069
4070 runner.settle_questions();
4071
4072 let back = store.get(&q.id).unwrap();
4073 assert!(
4074 !back.status.open(),
4075 "{status:?} run's question must not outlive the run"
4076 );
4077 }
4078 }
4079
4080 #[test]
4081 fn a_still_resumable_runs_open_question_is_left_alone() {
4082 ask_test_home();
4083 let store = ask::Questions::open();
4084 for status in [RunStatus::Blocked, RunStatus::Stalled] {
4090 let mut runner = runner_at(status);
4091 let run = runner.state.id.clone();
4092 let q = ask_open_question(&store, &run);
4093
4094 runner.settle_questions();
4095
4096 let back = store.get(&q.id).unwrap();
4097 assert!(
4098 back.status.open(),
4099 "{status:?} is still alive; the question must still be waiting"
4100 );
4101 }
4102 }
4103
4104 #[test]
4105 fn settle_questions_never_touches_an_already_answered_question() {
4106 ask_test_home();
4107 let store = ask::Questions::open();
4108 let mut runner = runner_at(RunStatus::Failed);
4109 let run = runner.state.id.clone();
4110 let mut q = ask_open_question(&store, &run);
4111 q.answer(crate::ask::Answer::Choice("SQLite".to_owned()))
4112 .unwrap();
4113 store.put(&mut q).unwrap();
4114
4115 runner.settle_questions();
4120 runner.settle_questions();
4121
4122 let back = store.get(&q.id).unwrap();
4123 assert_eq!(
4124 back.status,
4125 ask::QuestionStatus::Answered,
4126 "a real answer is a decision on record, never overwritten by a sweep"
4127 );
4128 }
4129
4130 #[tokio::test]
4139 async fn merge_does_not_reattempt_once_a_run_has_concluded() {
4140 let tmp = tempfile::tempdir().expect("tempdir");
4141 let repo = tmp.path().join("repo");
4142 std::fs::create_dir_all(&repo).unwrap();
4143 init_repo(&repo);
4144
4145 let mut config = Config::default();
4146 config.merge.mode = MergeMode::Local;
4147
4148 let mut state = RunState::new(
4149 repo.clone(),
4150 "main".to_owned(),
4151 "deadbeef".to_owned(),
4152 "task".to_owned(),
4153 config,
4154 );
4155 state.candidates = vec![Candidate {
4156 index: 0,
4157 label: 'A',
4158 agent: "alpha".to_owned(),
4159 branch: "does-not-exist".to_owned(),
4160 worktree: repo.clone(),
4161 summary: String::new(),
4162 stat: String::new(),
4163 files: 0,
4164 commits: 0,
4165 empty: false,
4166 failed: None,
4167 duration_ms: 0,
4168 folded: false,
4169 }];
4170 state.tally = Some(Tally {
4171 first_choice: BTreeMap::from([('A', 1)]),
4172 borda: BTreeMap::new(),
4173 winner: 'A',
4174 rankings: 1,
4175 unanimous_initial: true,
4176 deliberated: false,
4177 changed_votes: 0,
4178 unanimous_final: true,
4179 tie_break: None,
4180 judges: 0,
4181 present: 0,
4182 quorum: 0,
4183 met_quorum: true,
4184 uncontested: Some("only candidate A produced a change".to_owned()),
4185 });
4186 state.reviews = vec![ReviewRound {
4187 round: 1,
4188 head: "deadbeef".to_owned(),
4189 reviews: Vec::new(),
4190 e2e: Vec::new(),
4191 fix: None,
4192 blocking: 0,
4193 answered: 0,
4194 expected: 0,
4195 clean: true,
4196 verify_retried: false,
4197 progressed: false,
4198 vote_split: false,
4199 reconsideration: Vec::new(),
4200 verdict: None,
4201 }];
4202 state.gate = vec![CommandOutcome {
4203 command: "test".to_owned(),
4204 code: Some(0),
4205 output_tail: String::new(),
4206 duration_ms: 0,
4207 }];
4208 state.status = RunStatus::Ready;
4213 state.merge = Some(MergeOutcome {
4214 mode: MergeMode::Local,
4215 ok: false,
4216 detail: "already concluded".to_owned(),
4217 });
4218
4219 let mut runner = Runner {
4220 state,
4221 roles: ResolvedRoles {
4222 implementers: Vec::new(),
4223 judges: Vec::new(),
4224 reviewers: Vec::new(),
4225 fixer: None,
4226 },
4227 sem: Arc::new(Semaphore::new(1)),
4228 pause: Pause::new(),
4229 };
4230
4231 runner.merge().await.expect("merge");
4232
4233 assert_eq!(
4234 runner.state.status,
4235 RunStatus::Ready,
4236 "a concluded run's status must not change on reentry"
4237 );
4238 assert_eq!(
4239 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
4240 Some("already concluded"),
4241 "merge must not run again once the node already recorded an outcome"
4242 );
4243 }
4244
4245 #[tokio::test]
4246 async fn a_run_resumed_mid_landing_reenters_land_instead_of_opening_a_second_pull_request() {
4247 crate::run::set_home(std::env::temp_dir().join("magi-graph-test-home"));
4248 let tmp = tempfile::tempdir().expect("tempdir");
4249 let repo = tmp.path().join("repo");
4250 std::fs::create_dir_all(&repo).unwrap();
4251 init_repo(&repo);
4252
4253 let mut config = Config::default();
4254 config.merge.mode = MergeMode::Pr;
4255 config.graph.land = true;
4256 config.graph.land_approval = false;
4257
4258 let mut state = RunState::new(
4259 repo.clone(),
4260 "main".to_owned(),
4261 "deadbeef".to_owned(),
4262 "task".to_owned(),
4263 config,
4264 );
4265 state.candidates = vec![Candidate {
4266 index: 0,
4267 label: 'A',
4268 agent: "alpha".to_owned(),
4269 branch: "does-not-exist".to_owned(),
4270 worktree: repo.clone(),
4271 summary: String::new(),
4272 stat: String::new(),
4273 files: 0,
4274 commits: 0,
4275 empty: false,
4276 failed: None,
4277 duration_ms: 0,
4278 folded: false,
4279 }];
4280 state.tally = Some(Tally {
4281 first_choice: BTreeMap::from([('A', 1)]),
4282 borda: BTreeMap::new(),
4283 winner: 'A',
4284 rankings: 1,
4285 unanimous_initial: true,
4286 deliberated: false,
4287 changed_votes: 0,
4288 unanimous_final: true,
4289 tie_break: None,
4290 judges: 0,
4291 present: 0,
4292 quorum: 0,
4293 met_quorum: true,
4294 uncontested: Some("only candidate A produced a change".to_owned()),
4295 });
4296 state.reviews = vec![ReviewRound {
4297 round: 1,
4298 head: "deadbeef".to_owned(),
4299 reviews: Vec::new(),
4300 e2e: Vec::new(),
4301 fix: None,
4302 blocking: 0,
4303 answered: 0,
4304 expected: 0,
4305 clean: true,
4306 verify_retried: false,
4307 progressed: false,
4308 vote_split: false,
4309 reconsideration: Vec::new(),
4310 verdict: None,
4311 }];
4312 state.gate = vec![CommandOutcome {
4313 command: "test".to_owned(),
4314 code: Some(0),
4315 output_tail: String::new(),
4316 duration_ms: 0,
4317 }];
4318 state.status = RunStatus::Landing;
4322 state.merge = Some(MergeOutcome {
4323 mode: MergeMode::Pr,
4324 ok: true,
4325 detail: "https://example.invalid/x/y/pull/1".to_owned(),
4326 });
4327
4328 ask_test_home();
4332 let store = ask::Questions::open();
4333 let q = ask_open_question(&store, &state.id);
4334
4335 let mut runner = Runner {
4336 state,
4337 roles: ResolvedRoles {
4338 implementers: Vec::new(),
4339 judges: Vec::new(),
4340 reviewers: Vec::new(),
4341 fixer: None,
4342 },
4343 sem: Arc::new(Semaphore::new(1)),
4344 pause: Pause::new(),
4345 };
4346
4347 runner.execute().await.expect("execute");
4352
4353 assert_eq!(
4354 runner.state.merge.as_ref().map(|m| m.detail.as_str()),
4355 Some("https://example.invalid/x/y/pull/1"),
4356 "reentry must not push again or open a second pull request over the \
4357 one `land` is already watching"
4358 );
4359 assert_ne!(
4360 runner.state.status,
4361 RunStatus::Landing,
4362 "land could not actually reach the fake pull request, so it must \
4363 have given up rather than left the run silently parked forever"
4364 );
4365 assert_eq!(runner.state.status, RunStatus::Blocked);
4369 assert!(
4370 store.get(&q.id).unwrap().status.open(),
4371 "Blocked is still alive; settle_questions must have been a no-op here"
4372 );
4373 }
4374
4375 fn state_with_round(round: ReviewRound) -> RunState {
4376 let mut s = RunState::new(
4377 PathBuf::from("/repo"),
4378 "main".to_owned(),
4379 "abc1234".to_owned(),
4380 "add retries".to_owned(),
4381 Config::default(),
4382 );
4383 s.reviews = vec![round];
4384 s
4385 }
4386
4387 fn finding(id: &str, severity: Severity, title: &str) -> crate::verdict::Finding {
4388 crate::verdict::Finding {
4389 id: id.to_owned(),
4390 severity,
4391 file: None,
4392 line: None,
4393 title: title.to_owned(),
4394 detail: String::new(),
4395 }
4396 }
4397
4398 #[test]
4399 fn pr_body_names_open_findings_and_declined_ones() {
4400 let round = ReviewRound {
4401 round: 2,
4402 head: "deadbee".to_owned(),
4403 reviews: vec![ReviewRecord {
4404 reviewer: 1,
4405 agent: "alpha".to_owned(),
4406 summary: String::new(),
4407 findings: vec![finding("R2-1-1", Severity::Minor, "unused import")],
4408 vote: None,
4409 failed: None,
4410 duration_ms: 0,
4411 }],
4412 e2e: vec![CommandOutcome {
4413 command: "cargo test".to_owned(),
4414 code: Some(0),
4415 output_tail: String::new(),
4416 duration_ms: 0,
4417 }],
4418 verify_retried: false,
4419 fix: Some(FixRecord {
4420 agent: "alpha".to_owned(),
4421 addressed: Vec::new(),
4422 rejected: vec![crate::verdict::Rejection {
4423 id: "R1-1-1".to_owned(),
4424 why: "not reachable from any caller".to_owned(),
4425 }],
4426 notes: String::new(),
4427 committed: true,
4428 failed: None,
4429 duration_ms: 0,
4430 }),
4431 blocking: 0,
4432 answered: 1,
4433 expected: 1,
4434 clean: false,
4435 progressed: true,
4436 vote_split: false,
4437 reconsideration: Vec::new(),
4438 verdict: None,
4439 };
4440 let state = state_with_round(round);
4441 let body = pr_body(&state, 'A');
4442
4443 assert!(body.contains("add retries"), "the task must still be there");
4444 assert!(body.contains("R2-1-1"), "{body}");
4445 assert!(body.contains("unused import"), "{body}");
4446 assert!(body.contains("R1-1-1"), "the declined finding: {body}");
4447 assert!(
4448 body.contains("not reachable from any caller"),
4449 "the reason it was declined: {body}"
4450 );
4451 }
4452
4453 #[test]
4454 fn pr_body_says_nothing_extra_when_the_round_was_clean() {
4455 let round = ReviewRound {
4456 round: 1,
4457 head: "deadbee".to_owned(),
4458 reviews: vec![ReviewRecord {
4459 reviewer: 1,
4460 agent: "alpha".to_owned(),
4461 summary: String::new(),
4462 findings: Vec::new(),
4463 vote: None,
4464 failed: None,
4465 duration_ms: 0,
4466 }],
4467 e2e: Vec::new(),
4468 verify_retried: false,
4469 fix: None,
4470 blocking: 0,
4471 answered: 1,
4472 expected: 1,
4473 clean: true,
4474 progressed: false,
4475 vote_split: false,
4476 reconsideration: Vec::new(),
4477 verdict: None,
4478 };
4479 let state = state_with_round(round);
4480 let body = pr_body(&state, 'A');
4481 assert!(!body.contains("Open review findings"), "{body}");
4482 assert!(!body.contains("Declined"), "{body}");
4483 }
4484
4485 #[test]
4486 fn manual_merge_command_matches_the_configured_style() {
4487 let repo = Path::new("/repo");
4488 let message = "Merge magi run 0832 (candidate A)\n\nadd retries";
4489
4490 let merge = manual_merge_command(MergeStyle::Merge, repo, "magi/0832/A", message);
4491 assert_eq!(merge, "git -C /repo merge --no-ff magi/0832/A");
4492
4493 let squash = manual_merge_command(MergeStyle::Squash, repo, "magi/0832/A", message);
4494 assert_eq!(
4495 squash,
4496 "git -C /repo merge --squash magi/0832/A && git -C /repo commit -m \
4497 \"Merge magi run 0832 (candidate A)\""
4498 );
4499
4500 let rebase = manual_merge_command(MergeStyle::Rebase, repo, "magi/0832/A", message);
4501 assert_eq!(rebase, "git -C /repo merge --ff-only magi/0832/A");
4502 }
4503
4504 #[test]
4505 fn a_nudge_gets_a_quarter_of_the_budget() {
4506 assert_eq!(retry_budget(secs(1200), true), secs(300));
4508 assert_eq!(retry_budget(secs(3600), true), secs(900));
4509 }
4510
4511 #[test]
4512 fn a_resent_prompt_keeps_the_whole_budget() {
4513 assert_eq!(retry_budget(secs(1200), false), secs(1200));
4516 assert_eq!(retry_budget(secs(60), false), secs(60));
4517 }
4518
4519 #[test]
4520 fn the_floor_never_exceeds_the_original_budget() {
4521 assert_eq!(retry_budget(secs(60), true), secs(60));
4525 assert_eq!(retry_budget(secs(480), true), secs(120));
4526 assert_eq!(retry_budget(secs(0), true), secs(0));
4527 }
4528}