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