1use std::collections::BTreeSet;
42use std::path::{Path, PathBuf};
43use std::time::Duration;
44
45use anyhow::{Context as _, Result, bail};
46use serde::Deserialize;
47
48use crate::agent::{self, Invocation, SeatState};
49use crate::ask::{Question, Questions};
50use crate::config::Config;
51use crate::prompt;
52use crate::queue::{Queue, Task, TaskStatus};
53use crate::run::RunState;
54use crate::verdict;
55
56const SEAT: &str = "conduct";
59
60pub const NODE: &str = "conduct";
64
65const TURN_TIMEOUT: Duration = Duration::from_secs(300);
70
71#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
74#[serde(rename_all = "lowercase")]
75pub enum Recovery {
76 Requeue,
79 Hold,
82 Review,
87 Done,
95}
96
97#[derive(Debug, Clone, Default, Deserialize)]
100pub struct Decision {
101 pub id: String,
103 #[serde(default)]
106 pub blocked_by: Vec<String>,
107 #[serde(default)]
109 pub reason: Option<String>,
110 #[serde(default)]
118 pub recovery: Option<Recovery>,
119 #[serde(default)]
123 pub question: Option<String>,
124 #[serde(default)]
126 pub choices: Vec<String>,
127}
128
129#[derive(Debug, Clone, Default, Deserialize)]
142pub struct Verdict {
143 pub decisions: Vec<Decision>,
146}
147
148fn view(t: &Task, max_attempts: usize) -> prompt::ConductTask {
151 prompt::ConductTask {
152 id: t.id.clone(),
153 title: t.title.clone(),
154 instruction: t.instruction.clone(),
155 repo: t.repo.display().to_string(),
156 priority: t.priority,
157 status: t.status.as_str().to_owned(),
158 attempts: t.attempts,
159 max_attempts,
160 last_error: t.last_error.clone(),
161 hold_reason: t.hold_reason.clone(),
162 hold_source: t.hold_source.map(|source| source.label().to_owned()),
163 blocked_by: t.blocked_by.clone(),
164 answers: t
165 .answers
166 .iter()
167 .map(|a| prompt::ConductAnswer {
168 question: a.question.clone(),
169 answer: a.answer.clone(),
170 })
171 .collect(),
172 operator_resume: t.resume_override.as_ref().map(|o| {
173 format!(
174 "the operator explicitly answered \"resume\" at {}; do not hold this \
175 task again for the same reason unless there is new information",
176 o.at
177 )
178 }),
179 }
180}
181
182fn pinned_resume(task: &Task) -> bool {
186 task.resume_override
187 .as_ref()
188 .is_some_and(|o| o.pinned_run.is_some())
189}
190
191fn may_hold(task: &mut Task, reason: &str) -> bool {
200 let Some(o) = task.resume_override.as_mut() else {
201 return true;
202 };
203 if o.forced || o.pinned_run.is_some() {
204 tracing::warn!(
205 "conductor tried to hold task {} after the operator forced a resume: {reason}",
206 task.id
207 );
208 return false;
209 }
210 if o.conductor_rehold.is_some() {
211 return false;
212 }
213 o.conductor_rehold = Some(reason.to_owned());
214 true
215}
216
217fn severity_str(s: crate::verdict::Severity) -> &'static str {
220 match s {
221 crate::verdict::Severity::Nit => "nit",
222 crate::verdict::Severity::Minor => "minor",
223 crate::verdict::Severity::Major => "major",
224 crate::verdict::Severity::Blocker => "blocker",
225 }
226}
227
228fn surviving_branch(task: &Task) -> Option<String> {
233 let last = task.runs.last()?;
234 let state = RunState::load(last).ok()?;
235 state.winner().map(|c| c.branch.clone())
236}
237
238fn reaffirmed_hold_reason(task: &Task, d: &Decision) -> String {
258 let note = match &d.reason {
259 Some(reason) => reason.clone(),
260 None => match task.answers.last() {
261 Some(a) => format!(
262 "conduct held this again with no new reason given; last operator \
263 answer on record: {}",
264 a.answer
265 ),
266 None => "conduct held this again with no reason given".to_owned(),
267 },
268 };
269 match task.hold_reason.as_deref() {
270 Some(prior) if !prior.is_empty() => format!("{note}\n\n(previously: {prior})"),
271 _ => note,
272 }
273}
274
275fn hold_note(d: &Decision) -> String {
277 d.reason
278 .clone()
279 .unwrap_or_else(|| "(no reason given)".to_owned())
280}
281
282async fn outcome_for(task: &Task, repo: &Path) -> prompt::ConductOutcome {
284 let Some(run_id) = task.runs.last().cloned() else {
285 return prompt::ConductOutcome {
286 run_id: "(none)".to_owned(),
287 unreadable: Some("this task has not produced a run yet".to_owned()),
288 run_status: None,
289 open_findings: Vec::new(),
290 rounds_used: 0,
291 rounds_max: 0,
292 rounds: Vec::new(),
293 branch: None,
294 branch_head: None,
295 references: None,
296 empty_candidate: false,
297 };
298 };
299 let state = match RunState::load(&run_id) {
300 Ok(s) => s,
301 Err(e) => {
302 tracing::warn!(
307 "conductor: could not read run {run_id} for task {}: {e:#}",
308 task.short()
309 );
310 return prompt::ConductOutcome {
311 run_id,
312 unreadable: Some(format!("{e:#}")),
313 run_status: None,
314 open_findings: Vec::new(),
315 rounds_used: 0,
316 rounds_max: 0,
317 rounds: Vec::new(),
318 branch: None,
319 branch_head: None,
320 references: None,
321 empty_candidate: false,
322 };
323 }
324 };
325
326 let finding_view = |f: &crate::verdict::Finding| prompt::ConductFinding {
327 id: f.id.clone(),
328 title: f.title.clone(),
329 severity: severity_str(f.severity).to_owned(),
330 };
331 let open_findings = state
332 .open_findings()
333 .into_iter()
334 .map(finding_view)
335 .collect();
336 let rounds = state
337 .reviews
338 .iter()
339 .map(|r| prompt::ConductRound {
340 round: r.round,
341 findings: r
342 .reviews
343 .iter()
344 .flat_map(|rec| rec.findings.iter())
345 .map(finding_view)
346 .collect(),
347 addressed: r
348 .fix
349 .as_ref()
350 .map(|fx| fx.addressed.clone())
351 .unwrap_or_default(),
352 rejected: r
353 .fix
354 .as_ref()
355 .map(|fx| {
356 fx.rejected
357 .iter()
358 .map(|rej| prompt::ConductRejection {
359 id: rej.id.clone(),
360 why: rej.why.clone(),
361 })
362 .collect()
363 })
364 .unwrap_or_default(),
365 })
366 .collect();
367 let branch = state.winner().map(|c| c.branch.clone());
368 let branch_head = match &branch {
369 Some(b) => crate::git::rev_parse(repo, b)
370 .await
371 .ok()
372 .map(|h| h.chars().take(8).collect()),
373 None => None,
374 };
375
376 prompt::ConductOutcome {
377 run_id,
378 unreadable: None,
379 run_status: Some(state.status.as_str().to_owned()),
380 open_findings,
381 rounds_used: state.reviews.len(),
382 rounds_max: state.config.graph.review_rounds,
383 rounds,
384 branch,
385 branch_head,
386 references: crate::refs::describe(&state.seeds),
387 empty_candidate: state.merge.as_ref().is_some_and(|m| m.empty),
388 }
389}
390
391async fn attach_facts(cfg: &Config, repo: &Path, queue: &Queue, verdict: &mut Verdict) {
398 for d in verdict
399 .decisions
400 .iter_mut()
401 .filter(|d| d.question.is_some())
402 {
403 let Ok(task) = queue.get(&d.id) else {
404 continue;
405 };
406 let text = format!("{}\n{}", task.title, task.instruction);
407 if crate::refs::scan(&text).is_empty() {
408 continue;
409 }
410 let repo = repo_for(&task, repo);
411 let remote = &cfg.merge.remote;
412 let base = match cfg.merge.base.clone() {
413 Some(b) => Some(b),
414 None => crate::git::current_branch(&repo).await.ok().flatten(),
415 };
416 let facts = match base {
417 Some(base) => {
418 let base_name = base.clone();
419 let tracking = format!("{remote}/{base}");
420 let refreshed = crate::git::fetch(&repo, remote, &base)
423 .await
424 .is_ok_and(|o| o.ok());
425 let against = if crate::git::rev_exists(&repo, &tracking).await {
426 tracking
427 } else {
428 base
429 };
430 match crate::git::rev_parse(&repo, &against).await {
431 Ok(tip) => crate::refs::describe(
432 &crate::refs::resolve(&repo, &tip, remote, &text).await,
433 )
434 .map(|facts| {
435 if refreshed {
436 facts
437 } else {
438 format!(
439 "{facts}\n(could not fetch {remote}/{base_name}: this is \
440 against the local `{against}`, which may be behind the \
441 remote)"
442 )
443 }
444 }),
445 Err(e) => Some(format!("could not check the repository: {e:#}")),
446 }
447 }
448 None => Some("could not check the repository: no base branch known".to_owned()),
449 };
450 if let (Some(facts), Some(q)) = (facts, d.question.as_mut()) {
451 q.push_str(&format!(
452 "\n\nChecked against the repository (magi did this, not the model):\n{facts}"
453 ));
454 }
455 }
456}
457
458async fn finished_view(t: &Task, repo: &Path, max_attempts: usize) -> prompt::ConductFinished {
460 prompt::ConductFinished {
461 task: view(t, max_attempts),
462 outcome: outcome_for(t, &repo_for(t, repo)).await,
463 }
464}
465
466fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
469 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
470 fallback.to_path_buf()
471 } else {
472 task.repo.clone()
473 }
474}
475
476fn apply_one(queue: &Queue, questions: &Questions, d: &Decision) -> Result<()> {
493 let _claim = queue
494 .claim(&d.id)
495 .with_context(|| format!("task {} is claimed elsewhere right now", d.id))?;
496 let mut task = queue.get(&d.id).context("no such task")?;
497
498 if task.operator_held() {
502 return Ok(());
503 }
504
505 if task.status == TaskStatus::Held && crate::triage::pending_for(questions, &task) {
515 return Ok(());
516 }
517
518 let pinned = pinned_resume(&task);
521 if pinned && (d.question.is_some() || !d.blocked_by.is_empty()) {
522 return Ok(());
523 }
524
525 if let Some(text) = &d.question {
526 if task.status == TaskStatus::Done {
527 return Ok(());
528 }
529 let question_id = match questions
533 .list()
534 .into_iter()
535 .find(|q| q.status.open() && q.node == NODE && q.run == task.id)
536 {
537 Some(existing) => existing.id,
538 None => {
546 let mut q = Question::new(
547 task.id.clone(),
548 NODE.to_owned(),
549 SEAT.to_owned(),
550 text.clone(),
551 d.reason.clone().unwrap_or_default(),
552 d.choices.clone(),
553 );
554 questions.put(&mut q)?;
555 q.id
556 }
557 };
558 task.block(vec![question_id], d.reason.clone());
559 return queue.put(&mut task);
560 }
561
562 match task.status {
563 TaskStatus::Queued if !d.blocked_by.is_empty() => {
564 task.block(d.blocked_by.clone(), d.reason.clone());
565 queue.put(&mut task)?;
566 }
567 TaskStatus::Queued if d.recovery == Some(Recovery::Hold) => {
571 if may_hold(&mut task, &hold_note(d)) {
572 task.hold_machine(d.reason.clone());
573 queue.put(&mut task)?;
574 }
575 }
576 TaskStatus::Running => match d.recovery {
577 Some(Recovery::Requeue) if pinned_resume(&task) => {}
578 Some(Recovery::Requeue) => {
579 task.requeue();
580 queue.put(&mut task)?;
581 }
582 Some(Recovery::Hold) if may_hold(&mut task, &hold_note(d)) => {
583 task.hold_machine(d.reason.clone());
584 queue.put(&mut task)?;
585 }
586 _ => {}
590 },
591 TaskStatus::Failed | TaskStatus::Held => match d.recovery {
592 Some(Recovery::Requeue) if pinned_resume(&task) => {}
595 Some(Recovery::Requeue) => {
596 task.requeue();
597 queue.put(&mut task)?;
598 }
599 Some(Recovery::Hold) => {
600 if may_hold(&mut task, &hold_note(d)) {
601 task.hold_machine(Some(reaffirmed_hold_reason(&task, d)));
602 queue.put(&mut task)?;
603 }
604 }
605 Some(Recovery::Review) => {
606 if let Some(branch) = surviving_branch(&task) {
607 task.request_review(branch);
608 queue.put(&mut task)?;
609 }
610 }
616 Some(Recovery::Done) => {
617 task.succeed();
618 crate::daemon::supersede_prior_runs(&task, &crate::run::home());
622 queue.put(&mut task)?;
623 }
624 None => {}
625 },
626 _ => {}
629 }
630 Ok(())
631}
632
633pub fn apply(queue: &Queue, questions: &Questions, verdict: &Verdict) -> Result<()> {
637 for d in &verdict.decisions {
638 if let Err(e) = apply_one(queue, questions, d) {
639 tracing::warn!("conductor decision for task {}: {e:#}", d.id);
640 }
641 }
642 Ok(())
643}
644
645pub fn seat_path(home: &Path) -> PathBuf {
651 home.join("conduct").join("seat.json")
652}
653
654pub fn load_seat(home: &Path) -> Option<SeatState> {
656 serde_json::from_str(&std::fs::read_to_string(seat_path(home)).ok()?).ok()
657}
658
659fn busy_path(home: &Path) -> PathBuf {
660 home.join("conduct").join("busy")
661}
662
663pub fn busy(home: &Path) -> bool {
669 std::fs::metadata(busy_path(home))
670 .and_then(|m| m.modified())
671 .ok()
672 .and_then(|t| t.elapsed().ok())
673 .is_some_and(|age| age < TURN_TIMEOUT + Duration::from_secs(30))
674}
675
676struct Busy(PathBuf);
678
679impl Busy {
680 fn mark(home: &Path) -> Self {
681 let path = busy_path(home);
682 if let Some(dir) = path.parent() {
683 let _ = std::fs::create_dir_all(dir);
684 }
685 let _ = std::fs::write(&path, std::process::id().to_string());
686 Self(path)
687 }
688}
689
690impl Drop for Busy {
691 fn drop(&mut self) {
692 let _ = std::fs::remove_file(&self.0);
693 }
694}
695
696#[derive(Debug, Default)]
699pub struct Conductor {
700 seat: Option<SeatState>,
701 last_seen: Option<(u64, BTreeSet<String>)>,
702}
703
704impl Conductor {
705 #[must_use]
707 pub fn new() -> Self {
708 Self::default()
709 }
710
711 fn snapshot(queue: &Queue, stalled: &[Task], finished: &[Task]) -> (u64, BTreeSet<String>) {
712 let ids = stalled
713 .iter()
714 .chain(finished)
715 .map(|t| t.id.clone())
716 .collect();
717 (queue.revision(), ids)
718 }
719
720 #[must_use]
737 pub fn worth_a_look(&self, queue: &Queue, stalled: &[Task], finished: &[Task]) -> bool {
738 self.last_seen.as_ref() != Some(&Self::snapshot(queue, stalled, finished))
739 }
740
741 #[allow(clippy::too_many_arguments)]
745 pub async fn maybe_run(
746 &mut self,
747 cfg: &Config,
748 repo: &Path,
749 queue: &Queue,
750 questions: &Questions,
751 home: &Path,
752 queued: &[Task],
753 stalled: &[Task],
754 finished: &[Task],
755 max_attempts: usize,
756 ) {
757 let snapshot = Self::snapshot(queue, stalled, finished);
758 if self.last_seen.as_ref() == Some(&snapshot) {
759 return;
760 }
761 self.last_seen = Some(snapshot);
762 if let Err(e) = self
763 .run_once(
764 cfg,
765 repo,
766 queue,
767 questions,
768 home,
769 queued,
770 stalled,
771 finished,
772 max_attempts,
773 )
774 .await
775 {
776 tracing::warn!("conductor: {e:#}");
777 }
778 }
779
780 #[allow(clippy::too_many_arguments)]
781 async fn run_once(
782 &mut self,
783 cfg: &Config,
784 repo: &Path,
785 queue: &Queue,
786 questions: &Questions,
787 home: &Path,
788 queued: &[Task],
789 stalled: &[Task],
790 finished: &[Task],
791 max_attempts: usize,
792 ) -> Result<()> {
793 if queued.is_empty() && stalled.is_empty() && finished.is_empty() {
794 return Ok(());
795 }
796
797 let spec = cfg
798 .resolve_roles()
799 .context("resolving the conductor seat")?
800 .conductor;
801 let needs_new_seat = !matches!(&self.seat, Some(s) if s.agent == spec.id);
802 if needs_new_seat {
803 self.seat = Some(SeatState::new(SEAT, &spec.id, crate::rng::entropy()));
804 }
805 let seat = self.seat.as_mut().expect("just ensured a seat exists");
806
807 let runnable_views: Vec<prompt::ConductTask> =
808 queued.iter().map(|t| view(t, max_attempts)).collect();
809 let stalled_views: Vec<prompt::ConductTask> =
810 stalled.iter().map(|t| view(t, max_attempts)).collect();
811 let mut finished_views = Vec::with_capacity(finished.len());
812 for t in finished {
813 finished_views.push(finished_view(t, repo, max_attempts).await);
814 }
815
816 let body = prompt::with_overlay(
817 prompt::conduct(
818 &runnable_views,
819 &stalled_views,
820 &finished_views,
821 &cfg.graph.language,
822 ),
823 cfg.prompts.overlay(NODE),
824 );
825
826 let artifacts = home.join("conduct").join("artifacts");
827 let stem = format!("turn-{}", seat.turns + 1);
828 let cache_dir = cfg.cache_dir();
831 let inv = Invocation {
832 cwd: repo,
833 prompt: &body,
834 timeout: TURN_TIMEOUT,
835 allow_write: false,
838 sessions: cfg.graph.sessions,
839 artifacts: &artifacts,
840 stem: &stem,
841 run: NODE,
842 node: NODE,
843 cache_dir: cache_dir.as_deref(),
844 attachments: &[],
845 };
846
847 let busy = Busy::mark(home);
848 let out = agent::invoke(&spec, seat, &inv).await;
849 drop(busy);
850 if let Ok(body) = serde_json::to_string(&*seat) {
853 let path = seat_path(home);
854 if let Some(dir) = path.parent() {
855 let _ = std::fs::create_dir_all(dir);
856 }
857 let _ = std::fs::write(path, body);
858 }
859 let out = out.context("invoking the conductor")?;
860 if !out.usable() {
861 bail!(
862 "no usable reply (exit {:?}, timed out {})",
863 out.exit_code,
864 out.timed_out
865 );
866 }
867 let mut verdict: Verdict = verdict::extract_json(&out.text)
868 .context("the conductor's reply could not be parsed")?;
869 attach_facts(cfg, repo, queue, &mut verdict).await;
870 apply(queue, questions, &verdict)
871 }
872}
873
874#[cfg(test)]
875mod tests {
876 use std::collections::BTreeMap;
877
878 use tempfile::tempdir;
879
880 use super::*;
881 use crate::ask::{Answer, QuestionStatus};
882 use crate::config::{AgentKind, AgentSpec, Graph};
883 use crate::queue::Source;
884
885 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
886 let path = dir.join("mock-conduct-agent.sh");
887 std::fs::write(&path, script).expect("write mock");
888 AgentSpec {
889 id: "mock".to_owned(),
890 kind: AgentKind::Command,
891 model: None,
892 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
893 extra_args: Vec::new(),
894 env,
895 prompt_delivery: None,
896 }
897 }
898
899 fn config(spec: AgentSpec) -> Config {
900 Config {
901 agents: vec![spec],
902 graph: Graph {
903 language: "en".to_owned(),
904 ..Graph::default()
905 },
906 ..Config::default()
907 }
908 }
909
910 fn task(title: &str) -> Task {
911 Task::new(
912 title.to_owned(),
913 format!("do {title}"),
914 std::path::PathBuf::from("."),
915 Source::Human,
916 )
917 }
918
919 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
920 const GARBAGE: &str = "#!/bin/sh\ncat >/dev/null\nprintf 'not json at all\\n'\n";
921
922 fn env(reply: &str) -> BTreeMap<String, String> {
923 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
924 }
925
926 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
927
928 fn init_repo_with_branch(dir: &Path, branch: &str) {
932 use crate::proc::Quiet as _;
933 let run = |args: &[&str]| {
934 let out = std::process::Command::new("git")
935 .args(args)
936 .current_dir(dir)
937 .quiet()
938 .output()
939 .expect("spawn git");
940 assert!(
941 out.status.success(),
942 "git {args:?} failed: {}",
943 String::from_utf8_lossy(&out.stderr)
944 );
945 };
946 run(&["init", "-b", "main"]);
947 run(&["config", "user.name", "magi test"]);
948 run(&["config", "user.email", "magi@example.com"]);
949 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
950 run(&["add", "-A"]);
951 run(&["commit", "-m", "init"]);
952 run(&["checkout", "-b", branch]);
953 std::fs::write(dir.join("change.txt"), "x\n").unwrap();
954 run(&["add", "-A"]);
955 run(&["commit", "-m", "candidate work"]);
956 }
957
958 fn review_round_with_finding(
959 round: usize,
960 finding_id: &str,
961 title: &str,
962 addressed: &[&str],
963 rejected: &[(&str, &str)],
964 ) -> crate::run::ReviewRound {
965 crate::run::ReviewRound {
966 round,
967 head: "deadbeef".to_owned(),
968 verified_head: None,
969 verified_at: None,
970 reviews: vec![crate::run::ReviewRecord {
971 attempts: 0,
972 reviewer: 1,
973 agent: "mock".to_owned(),
974 summary: String::new(),
975 findings: vec![crate::verdict::Finding {
976 id: finding_id.to_owned(),
977 severity: crate::verdict::Severity::Major,
978 file: None,
979 line: None,
980 title: title.to_owned(),
981 detail: String::new(),
982 }],
983 vote: None,
984 failed: None,
985 duration_ms: 0,
986 }],
987 e2e: Vec::new(),
988 verify_retried: false,
989 e2e_deferred: false,
990 e2e_defer_reason: None,
991 fix: Some(crate::run::FixRecord {
992 agent: "mock".to_owned(),
993 addressed: addressed.iter().map(|s| (*s).to_owned()).collect(),
994 rejected: rejected
995 .iter()
996 .map(|(id, why)| crate::verdict::Rejection {
997 id: (*id).to_owned(),
998 why: (*why).to_owned(),
999 })
1000 .collect(),
1001 notes: String::new(),
1002 committed: false,
1003 failed: None,
1004 duration_ms: 0,
1005 continuation: None,
1006 }),
1007 blocking: 1,
1008 answered: 1,
1009 expected: 1,
1010 clean: false,
1011 progressed: true,
1012 vote_split: false,
1013 reconsideration: Vec::new(),
1014 verdict: None,
1015 }
1016 }
1017
1018 #[test]
1019 fn outcome_for_carries_every_rounds_findings_and_the_branch_head() {
1020 crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
1021 let dir = tempdir().unwrap();
1022 let default_repo = dir.path().join("default");
1023 let task_repo = dir.path().join("task");
1024 std::fs::create_dir_all(&default_repo).unwrap();
1025 std::fs::create_dir_all(&task_repo).unwrap();
1026 init_repo_with_branch(&default_repo, "other-branch");
1027 init_repo_with_branch(&task_repo, "magi/f00d/A");
1028
1029 let mut config = Config::default();
1030 config.graph.review_rounds = 6;
1031 let mut state = crate::run::RunState::new(
1032 task_repo.clone(),
1033 "main".to_owned(),
1034 "deadbeef".to_owned(),
1035 "task".to_owned(),
1036 config,
1037 );
1038 state.status = crate::run::RunStatus::Blocked;
1039 state.candidates.push(crate::run::Candidate {
1040 index: 0,
1041 label: 'A',
1042 agent: "mock".to_owned(),
1043 branch: "magi/f00d/A".to_owned(),
1044 worktree: task_repo.clone(),
1045 summary: String::new(),
1046 stat: String::new(),
1047 files: 1,
1048 commits: 1,
1049 empty: false,
1050 failed: None,
1051 verified_noop: None,
1052 duration_ms: 0,
1053 folded: false,
1054 });
1055 state.tally = Some(crate::run::Tally {
1056 first_choice: std::collections::BTreeMap::new(),
1057 borda: std::collections::BTreeMap::new(),
1058 winner: 'A',
1059 rankings: 0,
1060 unanimous_initial: false,
1061 deliberated: false,
1062 changed_votes: 0,
1063 unanimous_final: false,
1064 tie_break: None,
1065 judges: 0,
1066 present: 0,
1067 quorum: 0,
1068 met_quorum: true,
1069 uncontested: Some("solo".to_owned()),
1070 });
1071 state.reviews = vec![
1072 review_round_with_finding(
1073 1,
1074 "R1-1-2",
1075 "answer content is dropped",
1076 &[],
1077 &[("R1-1-2", "the id leaving blocked_by is enough")],
1078 ),
1079 review_round_with_finding(2, "R2-1-3", "answer content is still dropped", &[], &[]),
1080 ];
1081 state.save().unwrap();
1082
1083 let mut t = task("outcome test");
1084 t.repo = task_repo;
1085 t.runs.push(state.id.clone());
1086
1087 let finished = tokio_test_block_on(finished_view(&t, &default_repo, 2));
1088 let outcome = finished.outcome;
1089
1090 assert!(outcome.unreadable.is_none());
1091 assert_eq!(outcome.run_status.as_deref(), Some("blocked"));
1092 assert_eq!(outcome.rounds_used, 2);
1093 assert_eq!(outcome.rounds_max, 6);
1094 assert_eq!(outcome.rounds.len(), 2);
1095 assert_eq!(outcome.rounds[0].findings[0].id, "R1-1-2");
1096 assert_eq!(outcome.rounds[0].rejected[0].id, "R1-1-2");
1097 assert!(outcome.rounds[1].addressed.is_empty());
1098 assert!(outcome.rounds[1].rejected.is_empty());
1099 assert_eq!(outcome.branch.as_deref(), Some("magi/f00d/A"));
1100 assert!(
1101 outcome.branch_head.is_some(),
1102 "a real branch must resolve a head commit: {outcome:?}"
1103 );
1104 }
1105
1106 fn tokio_test_block_on<F: std::future::Future>(f: F) -> F::Output {
1110 tokio::runtime::Builder::new_current_thread()
1111 .enable_all()
1112 .build()
1113 .unwrap()
1114 .block_on(f)
1115 }
1116
1117 #[test]
1118 fn view_carries_a_tasks_recorded_answers_into_the_conductor_prompt_input() {
1119 let mut t = task("answered");
1120 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
1121 let v = view(&t, 2);
1122 assert_eq!(v.answers.len(), 1);
1123 assert_eq!(v.answers[0].question, "Which backend?");
1124 assert_eq!(v.answers[0].answer, "SQLite");
1125 }
1126
1127 #[test]
1128 fn a_dependency_decision_blocks_the_task_and_leaves_priority_alone() {
1129 let dir = tempdir().unwrap();
1130 let queue = Queue::at(dir.path().join("queue"));
1131 let questions = Questions::at(dir.path().join("questions"));
1132 let mut a = task("a");
1133 a.priority = 9;
1134 queue.put(&mut a).unwrap();
1135
1136 let verdict = Verdict {
1137 decisions: vec![Decision {
1138 id: a.id.clone(),
1139 blocked_by: vec!["20260101-000000-dead".to_owned()],
1140 reason: Some("waits on the other task".to_owned()),
1141 recovery: None,
1142 question: None,
1143 choices: Vec::new(),
1144 }],
1145 };
1146 apply(&queue, &questions, &verdict).unwrap();
1147
1148 let back = queue.get(&a.id).unwrap();
1149 assert_eq!(back.status, TaskStatus::Blocked);
1150 assert_eq!(back.blocked_by, ["20260101-000000-dead"]);
1151 assert_eq!(
1152 back.priority, 9,
1153 "the conductor's reply cannot carry priority"
1154 );
1155 }
1156
1157 #[test]
1158 fn a_question_decision_files_one_and_blocks_on_its_id() {
1159 let dir = tempdir().unwrap();
1160 let queue = Queue::at(dir.path().join("queue"));
1161 let questions = Questions::at(dir.path().join("questions"));
1162 let mut t = task("ambiguous");
1163 queue.put(&mut t).unwrap();
1164
1165 let verdict = Verdict {
1166 decisions: vec![Decision {
1167 id: t.id.clone(),
1168 blocked_by: Vec::new(),
1169 reason: Some("which backend?".to_owned()),
1170 recovery: None,
1171 question: Some("Which storage backend?".to_owned()),
1172 choices: vec!["SQLite".to_owned(), "Redis".to_owned()],
1173 }],
1174 };
1175 apply(&queue, &questions, &verdict).unwrap();
1176
1177 let back = queue.get(&t.id).unwrap();
1178 assert_eq!(back.status, TaskStatus::Blocked);
1179 assert_eq!(back.blocked_by.len(), 1);
1180 let q = questions.get(&back.blocked_by[0]).unwrap();
1181 assert_eq!(q.summary, "Which storage backend?");
1182 assert_eq!(q.node, NODE);
1183 assert!(q.status.open());
1184 }
1185
1186 #[test]
1187 fn a_task_with_an_open_question_already_reuses_it_rather_than_filing_a_second_one() {
1188 let dir = tempdir().unwrap();
1189 let queue = Queue::at(dir.path().join("queue"));
1190 let questions = Questions::at(dir.path().join("questions"));
1191 let mut t = task("asked once");
1192 queue.put(&mut t).unwrap();
1193
1194 let decision = Decision {
1195 id: t.id.clone(),
1196 reason: Some("still deciding".to_owned()),
1197 question: Some("Which backend?".to_owned()),
1198 ..Decision::default()
1199 };
1200 apply(
1201 &queue,
1202 &questions,
1203 &Verdict {
1204 decisions: vec![decision.clone()],
1205 },
1206 )
1207 .unwrap();
1208 assert_eq!(questions.list().len(), 1);
1209 let first_question_id = queue.get(&t.id).unwrap().blocked_by[0].clone();
1210
1211 let mut released = queue.get(&t.id).unwrap();
1216 released.release();
1217 queue.put(&mut released).unwrap();
1218
1219 apply(
1220 &queue,
1221 &questions,
1222 &Verdict {
1223 decisions: vec![decision],
1224 },
1225 )
1226 .unwrap();
1227 assert_eq!(questions.list().len(), 1, "no duplicate question was filed");
1228 let after = queue.get(&t.id).unwrap();
1229 assert_eq!(
1230 after.blocked_by,
1231 [first_question_id],
1232 "the existing open question is reused, not replaced"
1233 );
1234 }
1235
1236 #[test]
1237 fn a_same_id_question_from_another_node_is_not_reused() {
1238 let dir = tempdir().unwrap();
1239 let queue = Queue::at(dir.path().join("queue"));
1240 let questions = Questions::at(dir.path().join("questions"));
1241 let mut t = task("must ask the conductor");
1242 queue.put(&mut t).unwrap();
1243
1244 let mut unrelated = Question::new(
1245 t.id.clone(),
1246 "review".to_owned(),
1247 "reviewer-1".to_owned(),
1248 "An unrelated review question".to_owned(),
1249 String::new(),
1250 Vec::new(),
1251 );
1252 questions.put(&mut unrelated).unwrap();
1253
1254 apply(
1255 &queue,
1256 &questions,
1257 &Verdict {
1258 decisions: vec![Decision {
1259 id: t.id.clone(),
1260 question: Some("Which backend?".to_owned()),
1261 ..Decision::default()
1262 }],
1263 },
1264 )
1265 .unwrap();
1266
1267 let blocked_by = &queue.get(&t.id).unwrap().blocked_by;
1268 assert_eq!(blocked_by.len(), 1);
1269 assert_ne!(blocked_by[0], unrelated.id);
1270 assert!(questions.get(&unrelated.id).unwrap().status.open());
1271 assert_eq!(questions.get(&blocked_by[0]).unwrap().node, NODE);
1272 }
1273
1274 #[test]
1275 fn answering_the_question_lets_the_resolver_clear_the_block_with_the_answer_kept() {
1276 let dir = tempdir().unwrap();
1277 let queue = Queue::at(dir.path().join("queue"));
1278 let questions = Questions::at(dir.path().join("questions"));
1279 let mut t = task("waits on an answer");
1280 queue.put(&mut t).unwrap();
1281
1282 apply(
1283 &queue,
1284 &questions,
1285 &Verdict {
1286 decisions: vec![Decision {
1287 id: t.id.clone(),
1288 blocked_by: Vec::new(),
1289 reason: None,
1290 recovery: None,
1291 question: Some("Which backend?".to_owned()),
1292 choices: Vec::new(),
1293 }],
1294 },
1295 )
1296 .unwrap();
1297 let blocked = queue.get(&t.id).unwrap();
1298 let question_id = blocked.blocked_by[0].clone();
1299
1300 let mut q = questions.get(&question_id).unwrap();
1301 q.answer(Answer::Text("SQLite".to_owned())).unwrap();
1302 questions.put(&mut q).unwrap();
1303 assert_eq!(q.status, QuestionStatus::Answered);
1304
1305 let mut task_after = queue.get(&t.id).unwrap();
1309 task_after.record_answer(q.summary.clone(), "SQLite".to_owned());
1310 task_after.unblock(&question_id);
1311 assert_eq!(task_after.status, TaskStatus::Queued);
1312 assert_eq!(task_after.answers[0].answer, "SQLite");
1313 }
1314
1315 #[test]
1316 fn a_stalled_task_can_be_requeued_or_held() {
1317 let dir = tempdir().unwrap();
1318 let queue = Queue::at(dir.path().join("queue"));
1319 let questions = Questions::at(dir.path().join("questions"));
1320
1321 let mut requeue_me = task("stuck a");
1322 requeue_me.start("run-1".to_owned());
1323 queue.put(&mut requeue_me).unwrap();
1324
1325 let mut hold_me = task("stuck b");
1326 hold_me.start("run-2".to_owned());
1327 queue.put(&mut hold_me).unwrap();
1328
1329 apply(
1330 &queue,
1331 &questions,
1332 &Verdict {
1333 decisions: vec![
1334 Decision {
1335 id: requeue_me.id.clone(),
1336 recovery: Some(Recovery::Requeue),
1337 ..Decision::default()
1338 },
1339 Decision {
1340 id: hold_me.id.clone(),
1341 recovery: Some(Recovery::Hold),
1342 reason: Some("looks broken".to_owned()),
1343 ..Decision::default()
1344 },
1345 ],
1346 },
1347 )
1348 .unwrap();
1349
1350 let requeued = queue.get(&requeue_me.id).unwrap();
1351 assert_eq!(requeued.status, TaskStatus::Queued);
1352 assert_eq!(requeued.attempts, 0);
1353
1354 let held = queue.get(&hold_me.id).unwrap();
1355 assert_eq!(held.status, TaskStatus::Held);
1356 assert_eq!(held.hold_reason.as_deref(), Some("looks broken"));
1357 }
1358
1359 #[test]
1360 fn a_machine_held_task_asked_about_restores_to_held_once_answered() {
1361 let dir = tempdir().unwrap();
1366 let queue = Queue::at(dir.path().join("queue"));
1367 let questions = Questions::at(dir.path().join("questions"));
1368 let mut t = task("held out of attempts");
1369 t.hold_machine(Some("out of attempts".to_owned()));
1370 queue.put(&mut t).unwrap();
1371
1372 apply(
1373 &queue,
1374 &questions,
1375 &Verdict {
1376 decisions: vec![Decision {
1377 id: t.id.clone(),
1378 reason: Some("what should happen to this one?".to_owned()),
1379 question: Some("Hold it, or try again?".to_owned()),
1380 ..Decision::default()
1381 }],
1382 },
1383 )
1384 .unwrap();
1385 let blocked = queue.get(&t.id).unwrap();
1386 assert_eq!(blocked.status, TaskStatus::Blocked);
1387 let question_id = blocked.blocked_by[0].clone();
1388
1389 let mut q = questions.get(&question_id).unwrap();
1390 q.answer(Answer::Text("leave it held".to_owned())).unwrap();
1391 questions.put(&mut q).unwrap();
1392
1393 let mut after = queue.get(&t.id).unwrap();
1395 after.record_answer(q.summary.clone(), "leave it held".to_owned());
1396 after.unblock(&question_id);
1397 assert_eq!(
1398 after.status,
1399 TaskStatus::Held,
1400 "must not fall back to queued"
1401 );
1402 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
1403 }
1404
1405 #[test]
1406 fn a_reaffirmed_hold_with_no_new_reason_is_not_silently_auto_released_by_triage() {
1407 let dir = tempdir().unwrap();
1417 let queue = Queue::at(dir.path().join("queue"));
1418 let questions = Questions::at(dir.path().join("questions"));
1419 let mut t = task("disk pressure, then reconsidered");
1420 t.hold_machine(Some(
1421 "not enough free space to start a run: 10 bytes free, 100 required by \
1422 `[disk] min_free_bytes`"
1423 .to_owned(),
1424 ));
1425 t.record_answer(
1426 "How should this be handled?".to_owned(),
1427 "keep it held, a human will look at it later".to_owned(),
1428 );
1429 queue.put(&mut t).unwrap();
1430
1431 apply(
1432 &queue,
1433 &questions,
1434 &Verdict {
1435 decisions: vec![Decision {
1436 id: t.id.clone(),
1437 recovery: Some(Recovery::Hold),
1438 ..Decision::default()
1439 }],
1440 },
1441 )
1442 .unwrap();
1443
1444 let after = queue.get(&t.id).unwrap();
1445 assert_eq!(after.status, TaskStatus::Held);
1446 assert!(
1447 !after
1448 .hold_reason
1449 .as_deref()
1450 .unwrap_or_default()
1451 .starts_with("not enough free space"),
1452 "the stale disk-pressure text must not survive a reconfirmed hold: {:?}",
1453 after.hold_reason
1454 );
1455
1456 let cfg_dir = tempdir().unwrap();
1459 let config = cfg_dir.path().join("magi.toml");
1460 std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1461 let report =
1462 crate::triage::run_once(&queue, &questions, Some(&config), jiff::Timestamp::now());
1463 assert!(report.resumed.is_empty(), "must not be auto-released");
1464 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1465 }
1466
1467 #[test]
1468 fn a_conductor_rehold_after_a_resume_answer_is_not_asked_again_identically() {
1469 let dir = tempdir().unwrap();
1470 let queue = Queue::at(dir.path().join("queue"));
1471 let questions = Questions::at(dir.path().join("questions"));
1472 let config = dir.path().join("magi.toml");
1473 std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1474 let now = jiff::Timestamp::now();
1475 let triage = || crate::triage::run_once(&queue, &questions, Some(&config), now);
1476 let open = || {
1477 questions
1478 .list()
1479 .into_iter()
1480 .filter(|q| q.node == "triage" && q.status.open())
1481 .collect::<Vec<_>>()
1482 };
1483 let hold = Decision {
1484 recovery: Some(Recovery::Hold),
1485 reason: Some("waiting on manual worktree cleanup".to_owned()),
1486 ..Decision::default()
1487 };
1488
1489 let mut t = task("looping hold");
1490 t.hold_machine(Some("waiting on manual worktree cleanup".to_owned()));
1491 queue.put(&mut t).unwrap();
1492 let hold = Decision {
1493 id: t.id.clone(),
1494 ..hold
1495 };
1496
1497 assert_eq!(triage().asked.len(), 1);
1499 let first = open().remove(0);
1500 let mut q = questions.get(&first.id).unwrap();
1501 let resume = q.choices[0].clone();
1502 q.answer(Answer::Choice(resume)).unwrap();
1503 questions.put(&mut q).unwrap();
1504 assert_eq!(triage().answered.len(), 1);
1505 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1506
1507 apply(
1509 &queue,
1510 &questions,
1511 &Verdict {
1512 decisions: vec![hold.clone()],
1513 },
1514 )
1515 .unwrap();
1516 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1517
1518 assert_eq!(triage().asked.len(), 1);
1520 let second = open().remove(0);
1521 assert_ne!(second.summary, first.summary);
1522 assert_ne!(second.choices, first.choices);
1523 assert!(second.detail.contains("waiting on manual worktree cleanup"));
1524 assert!(triage().asked.is_empty(), "no duplicate question");
1525 assert_eq!(open().len(), 1);
1526
1527 let mut q = questions.get(&second.id).unwrap();
1529 let force = q.choices[0].clone();
1530 q.answer(Answer::Choice(force)).unwrap();
1531 questions.put(&mut q).unwrap();
1532 assert_eq!(triage().answered.len(), 1);
1533 apply(
1534 &queue,
1535 &questions,
1536 &Verdict {
1537 decisions: vec![hold],
1538 },
1539 )
1540 .unwrap();
1541 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1542 }
1543
1544 #[test]
1545 fn a_runnable_task_can_be_held_directly_without_a_question() {
1546 let dir = tempdir().unwrap();
1547 let queue = Queue::at(dir.path().join("queue"));
1548 let questions = Questions::at(dir.path().join("questions"));
1549 let mut t = task("already answered, should stay put");
1550 queue.put(&mut t).unwrap();
1551
1552 apply(
1553 &queue,
1554 &questions,
1555 &Verdict {
1556 decisions: vec![Decision {
1557 id: t.id.clone(),
1558 recovery: Some(Recovery::Hold),
1559 reason: Some("operator already said keep this held".to_owned()),
1560 ..Decision::default()
1561 }],
1562 },
1563 )
1564 .unwrap();
1565
1566 let after = queue.get(&t.id).unwrap();
1567 assert_eq!(after.status, TaskStatus::Held);
1568 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
1569 }
1570
1571 #[test]
1572 fn done_recovery_closes_a_held_task_whose_goal_is_already_met() {
1573 let dir = tempdir().unwrap();
1579 let queue = Queue::at(dir.path().join("queue"));
1580 let questions = Questions::at(dir.path().join("questions"));
1581 let mut t = task("already merged by hand");
1582 t.hold_machine(Some("branch survived, awaiting a decision".to_owned()));
1583 t.record_answer(
1584 "Handle this one?".to_owned(),
1585 "already merged and cleaned up, close it".to_owned(),
1586 );
1587 queue.put(&mut t).unwrap();
1588
1589 apply(
1590 &queue,
1591 &questions,
1592 &Verdict {
1593 decisions: vec![Decision {
1594 id: t.id.clone(),
1595 recovery: Some(Recovery::Done),
1596 reason: Some("operator confirmed this already landed".to_owned()),
1597 ..Decision::default()
1598 }],
1599 },
1600 )
1601 .unwrap();
1602
1603 let after = queue.get(&t.id).unwrap();
1604 assert_eq!(after.status, TaskStatus::Done);
1605 assert!(after.hold_reason.is_none());
1606 assert_eq!(after.answers.len(), 1, "the record of why is kept");
1607 }
1608
1609 #[test]
1610 fn done_recovery_is_ignored_for_a_runnable_or_running_task() {
1611 let dir = tempdir().unwrap();
1612 let queue = Queue::at(dir.path().join("queue"));
1613 let questions = Questions::at(dir.path().join("questions"));
1614
1615 let mut queued = task("never ran yet");
1616 queue.put(&mut queued).unwrap();
1617
1618 let mut running = task("mid-run");
1619 running.start("run-1".to_owned());
1620 queue.put(&mut running).unwrap();
1621
1622 for id in [queued.id.clone(), running.id.clone()] {
1623 apply(
1624 &queue,
1625 &questions,
1626 &Verdict {
1627 decisions: vec![Decision {
1628 id,
1629 recovery: Some(Recovery::Done),
1630 ..Decision::default()
1631 }],
1632 },
1633 )
1634 .unwrap();
1635 }
1636
1637 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
1638 assert_eq!(queue.get(&running.id).unwrap().status, TaskStatus::Running);
1639 }
1640
1641 #[test]
1642 fn a_question_after_two_settled_answers_is_still_filed_and_blocks() {
1643 let dir = tempdir().unwrap();
1646 let queue = Queue::at(dir.path().join("queue"));
1647 let questions = Questions::at(dir.path().join("questions"));
1648 let mut t = task("asked about repeatedly");
1649 t.hold_machine(Some("out of attempts".to_owned()));
1650 t.record_answer("Handle this one? (1)".to_owned(), "not yet".to_owned());
1651 t.record_answer(
1652 "Handle this one? (2)".to_owned(),
1653 "still not yet".to_owned(),
1654 );
1655 queue.put(&mut t).unwrap();
1656 assert_eq!(questions.list().len(), 0);
1657
1658 let verdict = Verdict {
1659 decisions: vec![Decision {
1660 id: t.id.clone(),
1661 question: Some("Branch conflicts with origin/main, how do we proceed?".to_owned()),
1662 ..Decision::default()
1663 }],
1664 };
1665 apply(&queue, &questions, &verdict).unwrap();
1666
1667 let filed = questions.list();
1668 assert_eq!(filed.len(), 1, "the question was filed");
1669 assert_eq!(
1670 filed[0].summary,
1671 "Branch conflicts with origin/main, how do we proceed?"
1672 );
1673 assert!(filed[0].status.open());
1674 assert_eq!(filed[0].node, NODE);
1675 let after = queue.get(&t.id).unwrap();
1676 assert_eq!(after.status, TaskStatus::Blocked);
1677 assert_eq!(after.blocked_by, vec![filed[0].id.clone()]);
1678 assert_eq!(after.answers.len(), 2, "the prior answers are untouched");
1679 assert!(
1680 !after
1681 .hold_reason
1682 .clone()
1683 .unwrap_or_default()
1684 .contains("conduct tried to ask"),
1685 "no hold was applied"
1686 );
1687
1688 apply(&queue, &questions, &verdict).unwrap();
1689 let again = questions.list();
1690 assert_eq!(again.len(), 1, "the open question is reused");
1691 assert_eq!(again[0].id, filed[0].id);
1692 }
1693
1694 #[test]
1695 fn a_second_conductor_question_is_still_allowed_after_one_settled_answer() {
1696 let dir = tempdir().unwrap();
1697 let queue = Queue::at(dir.path().join("queue"));
1698 let questions = Questions::at(dir.path().join("questions"));
1699 let mut t = task("asked about once already");
1700 t.hold_machine(Some("out of attempts".to_owned()));
1701 t.record_answer("Handle this one?".to_owned(), "not yet".to_owned());
1702 queue.put(&mut t).unwrap();
1703
1704 apply(
1705 &queue,
1706 &questions,
1707 &Verdict {
1708 decisions: vec![Decision {
1709 id: t.id.clone(),
1710 question: Some("Still not sure - now what?".to_owned()),
1711 ..Decision::default()
1712 }],
1713 },
1714 )
1715 .unwrap();
1716
1717 assert_eq!(questions.list().len(), 1, "the second question was filed");
1718 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
1719 }
1720
1721 #[test]
1722 fn a_held_task_with_an_open_triage_question_is_left_to_triage() {
1723 let dir = tempdir().unwrap();
1729 let queue = Queue::at(dir.path().join("queue"));
1730 let questions = Questions::at(dir.path().join("questions"));
1731 let mut t = task("held, triage already asking about it");
1732 t.hold_machine(Some("cause unclear".to_owned()));
1733 queue.put(&mut t).unwrap();
1734
1735 let mut triage_q = Question::new(
1736 t.id.clone(),
1737 crate::triage::NODE.to_owned(),
1738 "triage".to_owned(),
1739 "Still needed?".to_owned(),
1740 String::new(),
1741 vec![
1742 "resume".to_owned(),
1743 "not yet".to_owned(),
1744 "discard".to_owned(),
1745 ],
1746 );
1747 questions.put(&mut triage_q).unwrap();
1748
1749 for decision in [
1750 Decision {
1751 id: t.id.clone(),
1752 question: Some("what now?".to_owned()),
1753 ..Decision::default()
1754 },
1755 Decision {
1756 id: t.id.clone(),
1757 recovery: Some(Recovery::Requeue),
1758 ..Decision::default()
1759 },
1760 ] {
1761 apply(
1762 &queue,
1763 &questions,
1764 &Verdict {
1765 decisions: vec![decision],
1766 },
1767 )
1768 .unwrap();
1769 }
1770
1771 let after = queue.get(&t.id).unwrap();
1772 assert_eq!(
1773 after.status,
1774 TaskStatus::Held,
1775 "triage still owns this hold"
1776 );
1777 assert!(after.blocked_by.is_empty());
1778 assert_eq!(
1779 questions.list().len(),
1780 1,
1781 "no second, conductor-owned question was filed"
1782 );
1783 }
1784
1785 #[test]
1786 fn a_held_task_with_an_answered_but_unapplied_triage_question_is_still_left_alone() {
1787 let dir = tempdir().unwrap();
1795 let queue = Queue::at(dir.path().join("queue"));
1796 let questions = Questions::at(dir.path().join("questions"));
1797 let mut t = task("held, triage question answered but not yet applied");
1798 t.hold_machine(Some("cause unclear".to_owned()));
1799 queue.put(&mut t).unwrap();
1800
1801 let mut triage_q = Question::new(
1802 t.id.clone(),
1803 crate::triage::NODE.to_owned(),
1804 "triage".to_owned(),
1805 "Still needed?".to_owned(),
1806 String::new(),
1807 vec![
1808 "resume".to_owned(),
1809 "not yet".to_owned(),
1810 "discard".to_owned(),
1811 ],
1812 );
1813 questions.put(&mut triage_q).unwrap();
1814 triage_q
1815 .answer(Answer::Choice("not yet".to_owned()))
1816 .unwrap();
1817 questions.put(&mut triage_q).unwrap();
1818 assert!(!triage_q.status.open());
1819
1820 apply(
1821 &queue,
1822 &questions,
1823 &Verdict {
1824 decisions: vec![Decision {
1825 id: t.id.clone(),
1826 question: Some("what now?".to_owned()),
1827 ..Decision::default()
1828 }],
1829 },
1830 )
1831 .unwrap();
1832
1833 let after = queue.get(&t.id).unwrap();
1834 assert_eq!(
1835 after.status,
1836 TaskStatus::Held,
1837 "triage's own answer is not yet applied - conduct must wait"
1838 );
1839 assert_eq!(
1840 questions.list().len(),
1841 1,
1842 "no conductor question was filed over the pending triage answer"
1843 );
1844 }
1845
1846 #[test]
1847 fn manual_hold_rejects_hostile_or_stale_conductor_recovery() {
1848 let dir = tempdir().unwrap();
1849 let queue = Queue::at(dir.path().join("queue"));
1850 let questions = Questions::at(dir.path().join("questions"));
1851 let mut held = task("manual recovery");
1852 held.priority = 300;
1853 held.runs.push("run20260912-224242-daf5".to_owned());
1854 held.hold_manual(Some(
1855 "active manual recovery run20260912-224242-daf5".to_owned(),
1856 ));
1857 queue.put(&mut held).unwrap();
1858
1859 for decision in [
1863 Decision {
1864 id: held.id.clone(),
1865 recovery: Some(Recovery::Requeue),
1866 ..Decision::default()
1867 },
1868 Decision {
1869 id: held.id.clone(),
1870 recovery: Some(Recovery::Hold),
1871 reason: Some("stale replacement reason".to_owned()),
1872 ..Decision::default()
1873 },
1874 Decision {
1875 id: held.id.clone(),
1876 recovery: Some(Recovery::Review),
1877 ..Decision::default()
1878 },
1879 Decision {
1880 id: held.id.clone(),
1881 blocked_by: vec!["other-task".to_owned()],
1882 question: Some("retry now?".to_owned()),
1883 ..Decision::default()
1884 },
1885 ] {
1886 apply(
1887 &queue,
1888 &questions,
1889 &Verdict {
1890 decisions: vec![decision],
1891 },
1892 )
1893 .unwrap();
1894 }
1895
1896 let after = queue.get(&held.id).unwrap();
1897 assert_eq!(after.status, TaskStatus::Held);
1898 assert!(after.operator_held());
1899 assert_eq!(after.priority, 300);
1900 assert_eq!(after.runs, ["run20260912-224242-daf5"]);
1901 assert_eq!(
1902 after.hold_reason.as_deref(),
1903 Some("active manual recovery run20260912-224242-daf5")
1904 );
1905 assert!(after.blocked_by.is_empty());
1906 assert!(questions.list().is_empty());
1907 assert!(
1908 queue.next_runnable().is_none(),
1909 "must not dispatch a duplicate"
1910 );
1911 }
1912
1913 #[test]
1914 fn machine_holds_remain_recoverable_and_manual_release_is_authorization() {
1915 let dir = tempdir().unwrap();
1916 let queue = Queue::at(dir.path().join("queue"));
1917 let questions = Questions::at(dir.path().join("questions"));
1918
1919 let mut automatic = task("disk gate");
1920 automatic.hold_machine(Some("disk full".to_owned()));
1921 queue.put(&mut automatic).unwrap();
1922 let requeue = || Verdict {
1923 decisions: vec![Decision {
1924 id: automatic.id.clone(),
1925 recovery: Some(Recovery::Requeue),
1926 ..Decision::default()
1927 }],
1928 };
1929 apply(&queue, &questions, &requeue()).unwrap();
1930 assert_eq!(queue.get(&automatic.id).unwrap().status, TaskStatus::Queued);
1931
1932 let mut manual = task("operator gate");
1933 manual.hold_manual(Some("wait for operator".to_owned()));
1934 queue.put(&mut manual).unwrap();
1935 apply(
1936 &queue,
1937 &questions,
1938 &Verdict {
1939 decisions: vec![Decision {
1940 id: manual.id.clone(),
1941 recovery: Some(Recovery::Requeue),
1942 ..Decision::default()
1943 }],
1944 },
1945 )
1946 .unwrap();
1947 assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Held);
1948
1949 let mut released = queue.get(&manual.id).unwrap();
1952 released.release();
1953 queue.put(&mut released).unwrap();
1954 assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Queued);
1955 }
1956
1957 #[test]
1958 fn legacy_reasoned_hold_is_protected_without_losing_its_metadata() {
1959 let dir = tempdir().unwrap();
1960 let queue = Queue::at(dir.path().join("queue"));
1961 let questions = Questions::at(dir.path().join("questions"));
1962 let mut legacy = task("old explicit hold");
1963 legacy.status = TaskStatus::Held;
1964 legacy.hold_reason = Some("manual recovery already active".to_owned());
1965 legacy.hold_source = None;
1966 legacy.blocked_by = vec!["dependency".to_owned()];
1967 queue.put(&mut legacy).unwrap();
1968
1969 apply(
1970 &queue,
1971 &questions,
1972 &Verdict {
1973 decisions: vec![Decision {
1974 id: legacy.id.clone(),
1975 recovery: Some(Recovery::Requeue),
1976 ..Decision::default()
1977 }],
1978 },
1979 )
1980 .unwrap();
1981
1982 let after = queue.get(&legacy.id).unwrap();
1983 assert_eq!(after.status, TaskStatus::Held);
1984 assert_eq!(after.hold_source, None);
1985 assert_eq!(after.hold_reason, legacy.hold_reason);
1986 assert_eq!(after.blocked_by, legacy.blocked_by);
1987 }
1988
1989 #[test]
1990 fn review_recovery_is_a_no_op_without_a_survivable_branch() {
1991 crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
1997 let dir = tempdir().unwrap();
1998 let queue = Queue::at(dir.path().join("queue"));
1999 let questions = Questions::at(dir.path().join("questions"));
2000 let mut t = task("blocked with no readable run");
2001 t.start("20260101-000000-dead".to_owned()); t.fail("blocked", 5);
2003 queue.put(&mut t).unwrap();
2004
2005 apply(
2006 &queue,
2007 &questions,
2008 &Verdict {
2009 decisions: vec![Decision {
2010 id: t.id.clone(),
2011 recovery: Some(Recovery::Review),
2012 ..Decision::default()
2013 }],
2014 },
2015 )
2016 .unwrap();
2017
2018 let after = queue.get(&t.id).unwrap();
2019 assert_eq!(
2020 after.status,
2021 TaskStatus::Failed,
2022 "with nothing to reopen, the decision is dropped rather than guessed at"
2023 );
2024 assert!(after.review_branch.is_none());
2025 }
2026
2027 #[test]
2028 fn requeue_and_review_recovery_are_ignored_for_a_runnable_task() {
2029 let dir = tempdir().unwrap();
2033 let queue = Queue::at(dir.path().join("queue"));
2034 let questions = Questions::at(dir.path().join("questions"));
2035
2036 for recovery in [Recovery::Requeue, Recovery::Review] {
2037 let mut t = task("ordinary");
2038 queue.put(&mut t).unwrap();
2039
2040 apply(
2041 &queue,
2042 &questions,
2043 &Verdict {
2044 decisions: vec![Decision {
2045 id: t.id.clone(),
2046 recovery: Some(recovery),
2047 ..Decision::default()
2048 }],
2049 },
2050 )
2051 .unwrap();
2052
2053 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2054 }
2055 }
2056
2057 #[tokio::test]
2058 async fn a_broken_agent_leaves_the_queue_untouched_and_does_not_error() {
2059 let dir = tempdir().unwrap();
2060 let cfg = config(mock_agent(dir.path(), BROKEN, BTreeMap::new()));
2061 let queue = Queue::at(dir.path().join("queue"));
2062 let questions = Questions::at(dir.path().join("questions"));
2063 let mut t = task("normal");
2064 queue.put(&mut t).unwrap();
2065
2066 let mut conductor = Conductor::new();
2067 conductor
2068 .maybe_run(
2069 &cfg,
2070 dir.path(),
2071 &queue,
2072 &questions,
2073 dir.path(),
2074 &[t.clone()],
2075 &[],
2076 &[],
2077 2,
2078 )
2079 .await;
2080
2081 assert_eq!(
2082 queue.get(&t.id).unwrap().status,
2083 TaskStatus::Queued,
2084 "a failed invocation must change nothing"
2085 );
2086 assert!(
2087 queue.next_runnable().is_some(),
2088 "the loop must still be able to take the next task"
2089 );
2090 }
2091
2092 #[tokio::test]
2093 async fn a_reply_with_no_json_leaves_the_queue_untouched() {
2094 let dir = tempdir().unwrap();
2095 let cfg = config(mock_agent(dir.path(), GARBAGE, BTreeMap::new()));
2096 let queue = Queue::at(dir.path().join("queue"));
2097 let questions = Questions::at(dir.path().join("questions"));
2098 let mut t = task("normal");
2099 queue.put(&mut t).unwrap();
2100
2101 let mut conductor = Conductor::new();
2102 conductor
2103 .maybe_run(
2104 &cfg,
2105 dir.path(),
2106 &queue,
2107 &questions,
2108 dir.path(),
2109 &[t.clone()],
2110 &[],
2111 &[],
2112 2,
2113 )
2114 .await;
2115
2116 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2117 }
2118
2119 #[tokio::test]
2120 async fn json_survives_code_fences_and_a_preamble() {
2121 let dir = tempdir().unwrap();
2122 let mut t = task("fenced");
2123 let reply = format!(
2124 "Sure, here is my decision.\n\n```json\n{{\"decisions\":[{{\"id\":\"{}\",\
2125 \"blocked_by\":[\"x\"],\"reason\":\"why\"}}]}}\n```\n",
2126 t.id
2127 );
2128 let cfg = config(mock_agent(dir.path(), REPLY, env(&reply)));
2129 let queue = Queue::at(dir.path().join("queue"));
2130 let questions = Questions::at(dir.path().join("questions"));
2131 queue.put(&mut t).unwrap();
2132
2133 let mut conductor = Conductor::new();
2134 conductor
2135 .maybe_run(
2136 &cfg,
2137 dir.path(),
2138 &queue,
2139 &questions,
2140 dir.path(),
2141 &[t.clone()],
2142 &[],
2143 &[],
2144 2,
2145 )
2146 .await;
2147
2148 let back = queue.get(&t.id).unwrap();
2149 assert_eq!(back.status, TaskStatus::Blocked);
2150 assert_eq!(back.blocked_by, ["x"]);
2151 }
2152
2153 #[tokio::test]
2154 async fn the_conductor_is_not_called_again_when_nothing_worth_looking_at_has_changed() {
2155 let dir = tempdir().unwrap();
2158 let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2159 let queue = Queue::at(dir.path().join("queue"));
2160 let questions = Questions::at(dir.path().join("questions"));
2161 let mut t = task("stable");
2162 queue.put(&mut t).unwrap();
2163 let artifacts = dir.path().join("conduct").join("artifacts");
2164 let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2165
2166 let mut conductor = Conductor::new();
2167 conductor
2168 .maybe_run(
2169 &cfg,
2170 dir.path(),
2171 &queue,
2172 &questions,
2173 dir.path(),
2174 &[t.clone()],
2175 &[],
2176 &[],
2177 2,
2178 )
2179 .await;
2180 assert!(turn(1).is_file(), "the first cycle must call the conductor");
2181
2182 conductor
2183 .maybe_run(
2184 &cfg,
2185 dir.path(),
2186 &queue,
2187 &questions,
2188 dir.path(),
2189 &[t.clone()],
2190 &[],
2191 &[],
2192 2,
2193 )
2194 .await;
2195 assert!(
2196 !turn(2).is_file(),
2197 "an unchanged revision and an unchanged stalled/finished set must not call the \
2198 conductor twice"
2199 );
2200
2201 t.priority = 1;
2203 queue.put(&mut t).unwrap();
2204 conductor
2205 .maybe_run(
2206 &cfg,
2207 dir.path(),
2208 &queue,
2209 &questions,
2210 dir.path(),
2211 &[t.clone()],
2212 &[],
2213 &[],
2214 2,
2215 )
2216 .await;
2217 assert!(turn(2).is_file(), "a moved revision calls it again");
2218 }
2219
2220 #[tokio::test]
2221 async fn a_task_turning_stalled_calls_the_conductor_again_despite_an_unchanged_revision() {
2222 let dir = tempdir().unwrap();
2228 let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2229 let queue = Queue::at(dir.path().join("queue"));
2230 let questions = Questions::at(dir.path().join("questions"));
2231 let mut t = task("quiet");
2232 queue.put(&mut t).unwrap();
2233 let artifacts = dir.path().join("conduct").join("artifacts");
2234 let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2235
2236 let mut conductor = Conductor::new();
2237 conductor
2238 .maybe_run(
2239 &cfg,
2240 dir.path(),
2241 &queue,
2242 &questions,
2243 dir.path(),
2244 &[t.clone()],
2245 &[],
2246 &[],
2247 2,
2248 )
2249 .await;
2250 assert!(turn(1).is_file());
2251
2252 conductor
2253 .maybe_run(
2254 &cfg,
2255 dir.path(),
2256 &queue,
2257 &questions,
2258 dir.path(),
2259 &[],
2260 &[t.clone()],
2261 &[],
2262 2,
2263 )
2264 .await;
2265 assert!(
2266 turn(2).is_file(),
2267 "a task turning stalled must call the conductor again"
2268 );
2269
2270 conductor
2273 .maybe_run(
2274 &cfg,
2275 dir.path(),
2276 &queue,
2277 &questions,
2278 dir.path(),
2279 &[],
2280 &[t.clone()],
2281 &[],
2282 2,
2283 )
2284 .await;
2285 assert!(
2286 !turn(3).is_file(),
2287 "the same stalled task lingering must not call the conductor every cycle"
2288 );
2289 }
2290
2291 #[test]
2292 fn worth_a_look_is_config_free_and_matches_maybe_runs_own_gate() {
2293 let dir = tempdir().unwrap();
2294 let queue = Queue::at(dir.path().join("queue"));
2295 let mut t = task("t");
2296 queue.put(&mut t).unwrap();
2297
2298 let mut conductor = Conductor::new();
2299 assert!(
2300 conductor.worth_a_look(&queue, &[], &[]),
2301 "a conductor that has never run has something to look at"
2302 );
2303
2304 conductor.last_seen = Some(Conductor::snapshot(&queue, &[], &[]));
2305 assert!(
2306 !conductor.worth_a_look(&queue, &[], &[]),
2307 "nothing changed and nothing is stalled or finished"
2308 );
2309 assert!(
2310 conductor.worth_a_look(&queue, &[t.clone()], &[]),
2311 "a stalled task is worth a look even at the same revision"
2312 );
2313 assert!(
2314 conductor.worth_a_look(&queue, &[], &[t.clone()]),
2315 "a finished task is worth a look even at the same revision"
2316 );
2317 }
2318
2319 #[tokio::test]
2320 async fn the_conduct_path_never_calls_ask_and_wait() {
2321 let dir = tempdir().unwrap();
2327 let queue = Queue::at(dir.path().join("queue"));
2328 let questions = Questions::at(dir.path().join("questions"));
2329 let mut t = task("asks without blocking");
2330 queue.put(&mut t).unwrap();
2331
2332 apply(
2333 &queue,
2334 &questions,
2335 &Verdict {
2336 decisions: vec![Decision {
2337 id: t.id.clone(),
2338 question: Some("ok?".to_owned()),
2339 ..Decision::default()
2340 }],
2341 },
2342 )
2343 .unwrap();
2344 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
2346 }
2347
2348 #[test]
2349 fn a_pinned_resume_is_not_held_or_requeued_by_the_conductor() {
2350 let mut t = Task::new("t".into(), "t".into(), PathBuf::new(), Source::Human);
2351 t.resume_override = Some(crate::queue::OperatorResume {
2352 question_id: "q".into(),
2353 at: jiff::Timestamp::now(),
2354 conductor_rehold: None,
2355 forced: false,
2356 pinned_run: Some("run-1".into()),
2357 });
2358 assert!(pinned_resume(&t));
2359 assert!(!may_hold(&mut t, "waiting for magi resume"));
2360 assert!(
2361 t.resume_override
2362 .as_ref()
2363 .unwrap()
2364 .conductor_rehold
2365 .is_none(),
2366 "a refused hold is not recorded as an override"
2367 );
2368 t.resume_override = None;
2369 assert!(!pinned_resume(&t));
2370 assert!(may_hold(&mut t, "no override, so a hold is allowed"));
2371 }
2372
2373 #[test]
2374 fn a_pinned_resume_is_not_blocked_by_a_conductor_question_or_dependency() {
2375 let dir = tempdir().unwrap();
2376 let queue = Queue::at(dir.path().join("queue"));
2377 let questions = Questions::at(dir.path().join("questions"));
2378 let mut t = task("resume me");
2379 t.resume_override = Some(crate::queue::OperatorResume {
2380 question_id: "q".into(),
2381 at: jiff::Timestamp::now(),
2382 conductor_rehold: None,
2383 forced: true,
2384 pinned_run: Some("run-1".into()),
2385 });
2386 queue.put(&mut t).unwrap();
2387
2388 apply(
2389 &queue,
2390 &questions,
2391 &Verdict {
2392 decisions: vec![
2393 Decision {
2394 id: t.id.clone(),
2395 question: Some("really?".to_owned()),
2396 ..Decision::default()
2397 },
2398 Decision {
2399 id: t.id.clone(),
2400 blocked_by: vec!["other".to_owned()],
2401 ..Decision::default()
2402 },
2403 ],
2404 },
2405 )
2406 .unwrap();
2407
2408 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2409 assert!(questions.list().is_empty());
2410 }
2411}