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 may_hold(task: &mut Task, reason: &str) -> bool {
210 let Some(o) = task.resume_override.as_mut() else {
211 return true;
212 };
213 if o.forced {
214 tracing::warn!(
215 "conductor tried to hold task {} after the operator forced a resume: {reason}",
216 task.id
217 );
218 return false;
219 }
220 if o.conductor_rehold.is_some() {
221 return false;
222 }
223 o.conductor_rehold = Some(reason.to_owned());
224 true
225}
226
227fn severity_str(s: crate::verdict::Severity) -> &'static str {
230 match s {
231 crate::verdict::Severity::Nit => "nit",
232 crate::verdict::Severity::Minor => "minor",
233 crate::verdict::Severity::Major => "major",
234 crate::verdict::Severity::Blocker => "blocker",
235 }
236}
237
238fn surviving_branch(task: &Task) -> Option<String> {
243 let last = task.runs.last()?;
244 let state = RunState::load(last).ok()?;
245 state.winner().map(|c| c.branch.clone())
246}
247
248fn reaffirmed_hold_reason(task: &Task, d: &Decision) -> String {
268 let note = match &d.reason {
269 Some(reason) => reason.clone(),
270 None => match task.answers.last() {
271 Some(a) => format!(
272 "conduct held this again with no new reason given; last operator \
273 answer on record: {}",
274 a.answer
275 ),
276 None => "conduct held this again with no reason given".to_owned(),
277 },
278 };
279 match task.hold_reason.as_deref() {
280 Some(prior) if !prior.is_empty() => format!("{note}\n\n(previously: {prior})"),
281 _ => note,
282 }
283}
284
285fn hold_note(d: &Decision) -> String {
287 d.reason
288 .clone()
289 .unwrap_or_else(|| "(no reason given)".to_owned())
290}
291
292async fn outcome_for(task: &Task, repo: &Path) -> prompt::ConductOutcome {
294 let Some(run_id) = task.runs.last().cloned() else {
295 return prompt::ConductOutcome {
296 run_id: "(none)".to_owned(),
297 unreadable: Some("this task has not produced a run yet".to_owned()),
298 run_status: None,
299 open_findings: Vec::new(),
300 rounds_used: 0,
301 rounds_max: 0,
302 rounds: Vec::new(),
303 branch: None,
304 branch_head: None,
305 };
306 };
307 let state = match RunState::load(&run_id) {
308 Ok(s) => s,
309 Err(e) => {
310 tracing::warn!(
315 "conductor: could not read run {run_id} for task {}: {e:#}",
316 task.short()
317 );
318 return prompt::ConductOutcome {
319 run_id,
320 unreadable: Some(format!("{e:#}")),
321 run_status: None,
322 open_findings: Vec::new(),
323 rounds_used: 0,
324 rounds_max: 0,
325 rounds: Vec::new(),
326 branch: None,
327 branch_head: None,
328 };
329 }
330 };
331
332 let finding_view = |f: &crate::verdict::Finding| prompt::ConductFinding {
333 id: f.id.clone(),
334 title: f.title.clone(),
335 severity: severity_str(f.severity).to_owned(),
336 };
337 let open_findings = state
338 .open_findings()
339 .into_iter()
340 .map(finding_view)
341 .collect();
342 let rounds = state
343 .reviews
344 .iter()
345 .map(|r| prompt::ConductRound {
346 round: r.round,
347 findings: r
348 .reviews
349 .iter()
350 .flat_map(|rec| rec.findings.iter())
351 .map(finding_view)
352 .collect(),
353 addressed: r
354 .fix
355 .as_ref()
356 .map(|fx| fx.addressed.clone())
357 .unwrap_or_default(),
358 rejected: r
359 .fix
360 .as_ref()
361 .map(|fx| {
362 fx.rejected
363 .iter()
364 .map(|rej| prompt::ConductRejection {
365 id: rej.id.clone(),
366 why: rej.why.clone(),
367 })
368 .collect()
369 })
370 .unwrap_or_default(),
371 })
372 .collect();
373 let branch = state.winner().map(|c| c.branch.clone());
374 let branch_head = match &branch {
375 Some(b) => crate::git::rev_parse(repo, b)
376 .await
377 .ok()
378 .map(|h| h.chars().take(8).collect()),
379 None => None,
380 };
381
382 prompt::ConductOutcome {
383 run_id,
384 unreadable: None,
385 run_status: Some(state.status.as_str().to_owned()),
386 open_findings,
387 rounds_used: state.reviews.len(),
388 rounds_max: state.config.graph.review_rounds,
389 rounds,
390 branch,
391 branch_head,
392 }
393}
394
395async fn finished_view(t: &Task, repo: &Path, max_attempts: usize) -> prompt::ConductFinished {
397 prompt::ConductFinished {
398 task: view(t, max_attempts),
399 outcome: outcome_for(t, &repo_for(t, repo)).await,
400 }
401}
402
403fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
406 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
407 fallback.to_path_buf()
408 } else {
409 task.repo.clone()
410 }
411}
412
413fn apply_one(queue: &Queue, questions: &Questions, d: &Decision) -> Result<()> {
430 let _claim = queue
431 .claim(&d.id)
432 .with_context(|| format!("task {} is claimed elsewhere right now", d.id))?;
433 let mut task = queue.get(&d.id).context("no such task")?;
434
435 if task.operator_held() {
439 return Ok(());
440 }
441
442 if task.status == TaskStatus::Held && crate::triage::pending_for(questions, &task) {
452 return Ok(());
453 }
454
455 if let Some(text) = &d.question {
456 if task.status == TaskStatus::Done {
457 return Ok(());
458 }
459 let question_id = match questions
463 .list()
464 .into_iter()
465 .find(|q| q.status.open() && q.node == NODE && q.run == task.id)
466 {
467 Some(existing) => existing.id,
468 None if task.answers.len() >= MAX_SETTLED_CONDUCT_ANSWERS => {
473 task.hold_machine(Some(format!(
474 "conduct tried to ask another question after {} were \
475 already answered about this task: {text}",
476 task.answers.len()
477 )));
478 return queue.put(&mut task);
479 }
480 None => {
481 let mut q = Question::new(
482 task.id.clone(),
483 NODE.to_owned(),
484 SEAT.to_owned(),
485 text.clone(),
486 d.reason.clone().unwrap_or_default(),
487 d.choices.clone(),
488 );
489 questions.put(&mut q)?;
490 q.id
491 }
492 };
493 task.block(vec![question_id], d.reason.clone());
494 return queue.put(&mut task);
495 }
496
497 match task.status {
498 TaskStatus::Queued if !d.blocked_by.is_empty() => {
499 task.block(d.blocked_by.clone(), d.reason.clone());
500 queue.put(&mut task)?;
501 }
502 TaskStatus::Queued if d.recovery == Some(Recovery::Hold) => {
506 if may_hold(&mut task, &hold_note(d)) {
507 task.hold_machine(d.reason.clone());
508 queue.put(&mut task)?;
509 }
510 }
511 TaskStatus::Running => match d.recovery {
512 Some(Recovery::Requeue) => {
513 task.requeue();
514 queue.put(&mut task)?;
515 }
516 Some(Recovery::Hold) if may_hold(&mut task, &hold_note(d)) => {
517 task.hold_machine(d.reason.clone());
518 queue.put(&mut task)?;
519 }
520 _ => {}
524 },
525 TaskStatus::Failed | TaskStatus::Held => match d.recovery {
526 Some(Recovery::Requeue) => {
527 task.requeue();
528 queue.put(&mut task)?;
529 }
530 Some(Recovery::Hold) => {
531 if may_hold(&mut task, &hold_note(d)) {
532 task.hold_machine(Some(reaffirmed_hold_reason(&task, d)));
533 queue.put(&mut task)?;
534 }
535 }
536 Some(Recovery::Review) => {
537 if let Some(branch) = surviving_branch(&task) {
538 task.request_review(branch);
539 queue.put(&mut task)?;
540 }
541 }
547 Some(Recovery::Done) => {
548 task.succeed();
549 crate::daemon::supersede_prior_runs(&task, &crate::run::home());
553 queue.put(&mut task)?;
554 }
555 None => {}
556 },
557 _ => {}
560 }
561 Ok(())
562}
563
564pub fn apply(queue: &Queue, questions: &Questions, verdict: &Verdict) -> Result<()> {
568 for d in &verdict.decisions {
569 if let Err(e) = apply_one(queue, questions, d) {
570 tracing::warn!("conductor decision for task {}: {e:#}", d.id);
571 }
572 }
573 Ok(())
574}
575
576#[derive(Debug, Default)]
579pub struct Conductor {
580 seat: Option<SeatState>,
581 last_seen: Option<(u64, BTreeSet<String>)>,
582}
583
584impl Conductor {
585 #[must_use]
587 pub fn new() -> Self {
588 Self::default()
589 }
590
591 fn snapshot(queue: &Queue, stalled: &[Task], finished: &[Task]) -> (u64, BTreeSet<String>) {
592 let ids = stalled
593 .iter()
594 .chain(finished)
595 .map(|t| t.id.clone())
596 .collect();
597 (queue.revision(), ids)
598 }
599
600 #[must_use]
617 pub fn worth_a_look(&self, queue: &Queue, stalled: &[Task], finished: &[Task]) -> bool {
618 self.last_seen.as_ref() != Some(&Self::snapshot(queue, stalled, finished))
619 }
620
621 #[allow(clippy::too_many_arguments)]
625 pub async fn maybe_run(
626 &mut self,
627 cfg: &Config,
628 repo: &Path,
629 queue: &Queue,
630 questions: &Questions,
631 home: &Path,
632 queued: &[Task],
633 stalled: &[Task],
634 finished: &[Task],
635 max_attempts: usize,
636 ) {
637 let snapshot = Self::snapshot(queue, stalled, finished);
638 if self.last_seen.as_ref() == Some(&snapshot) {
639 return;
640 }
641 self.last_seen = Some(snapshot);
642 if let Err(e) = self
643 .run_once(
644 cfg,
645 repo,
646 queue,
647 questions,
648 home,
649 queued,
650 stalled,
651 finished,
652 max_attempts,
653 )
654 .await
655 {
656 tracing::warn!("conductor: {e:#}");
657 }
658 }
659
660 #[allow(clippy::too_many_arguments)]
661 async fn run_once(
662 &mut self,
663 cfg: &Config,
664 repo: &Path,
665 queue: &Queue,
666 questions: &Questions,
667 home: &Path,
668 queued: &[Task],
669 stalled: &[Task],
670 finished: &[Task],
671 max_attempts: usize,
672 ) -> Result<()> {
673 if queued.is_empty() && stalled.is_empty() && finished.is_empty() {
674 return Ok(());
675 }
676
677 let spec = cfg
678 .resolve_roles()
679 .context("resolving the conductor seat")?
680 .conductor;
681 let needs_new_seat = !matches!(&self.seat, Some(s) if s.agent == spec.id);
682 if needs_new_seat {
683 self.seat = Some(SeatState::new(SEAT, &spec.id, crate::rng::entropy()));
684 }
685 let seat = self.seat.as_mut().expect("just ensured a seat exists");
686
687 let runnable_views: Vec<prompt::ConductTask> =
688 queued.iter().map(|t| view(t, max_attempts)).collect();
689 let stalled_views: Vec<prompt::ConductTask> =
690 stalled.iter().map(|t| view(t, max_attempts)).collect();
691 let mut finished_views = Vec::with_capacity(finished.len());
692 for t in finished {
693 finished_views.push(finished_view(t, repo, max_attempts).await);
694 }
695
696 let body = prompt::with_overlay(
697 prompt::conduct(
698 &runnable_views,
699 &stalled_views,
700 &finished_views,
701 &cfg.graph.language,
702 ),
703 cfg.prompts.overlay(NODE),
704 );
705
706 let artifacts = home.join("conduct").join("artifacts");
707 let stem = format!("turn-{}", seat.turns + 1);
708 let cache_dir = cfg.cache_dir();
711 let inv = Invocation {
712 cwd: repo,
713 prompt: &body,
714 timeout: TURN_TIMEOUT,
715 allow_write: false,
718 sessions: cfg.graph.sessions,
719 artifacts: &artifacts,
720 stem: &stem,
721 run: NODE,
722 node: NODE,
723 cache_dir: cache_dir.as_deref(),
724 attachments: &[],
725 };
726
727 let out = agent::invoke(&spec, seat, &inv)
728 .await
729 .context("invoking the conductor")?;
730 if !out.usable() {
731 bail!(
732 "no usable reply (exit {:?}, timed out {})",
733 out.exit_code,
734 out.timed_out
735 );
736 }
737 let verdict: Verdict = verdict::extract_json(&out.text)
738 .context("the conductor's reply could not be parsed")?;
739 apply(queue, questions, &verdict)
740 }
741}
742
743#[cfg(test)]
744mod tests {
745 use std::collections::BTreeMap;
746
747 use tempfile::tempdir;
748
749 use super::*;
750 use crate::ask::{Answer, QuestionStatus};
751 use crate::config::{AgentKind, AgentSpec, Graph};
752 use crate::queue::Source;
753
754 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
755 let path = dir.join("mock-conduct-agent.sh");
756 std::fs::write(&path, script).expect("write mock");
757 AgentSpec {
758 id: "mock".to_owned(),
759 kind: AgentKind::Command,
760 model: None,
761 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
762 extra_args: Vec::new(),
763 env,
764 prompt_delivery: None,
765 }
766 }
767
768 fn config(spec: AgentSpec) -> Config {
769 Config {
770 agents: vec![spec],
771 graph: Graph {
772 language: "en".to_owned(),
773 ..Graph::default()
774 },
775 ..Config::default()
776 }
777 }
778
779 fn task(title: &str) -> Task {
780 Task::new(
781 title.to_owned(),
782 format!("do {title}"),
783 std::path::PathBuf::from("."),
784 Source::Human,
785 )
786 }
787
788 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
789 const GARBAGE: &str = "#!/bin/sh\ncat >/dev/null\nprintf 'not json at all\\n'\n";
790
791 fn env(reply: &str) -> BTreeMap<String, String> {
792 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
793 }
794
795 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
796
797 fn init_repo_with_branch(dir: &Path, branch: &str) {
801 use crate::proc::Quiet as _;
802 let run = |args: &[&str]| {
803 let out = std::process::Command::new("git")
804 .args(args)
805 .current_dir(dir)
806 .quiet()
807 .output()
808 .expect("spawn git");
809 assert!(
810 out.status.success(),
811 "git {args:?} failed: {}",
812 String::from_utf8_lossy(&out.stderr)
813 );
814 };
815 run(&["init", "-b", "main"]);
816 run(&["config", "user.name", "magi test"]);
817 run(&["config", "user.email", "magi@example.com"]);
818 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
819 run(&["add", "-A"]);
820 run(&["commit", "-m", "init"]);
821 run(&["checkout", "-b", branch]);
822 std::fs::write(dir.join("change.txt"), "x\n").unwrap();
823 run(&["add", "-A"]);
824 run(&["commit", "-m", "candidate work"]);
825 }
826
827 fn review_round_with_finding(
828 round: usize,
829 finding_id: &str,
830 title: &str,
831 addressed: &[&str],
832 rejected: &[(&str, &str)],
833 ) -> crate::run::ReviewRound {
834 crate::run::ReviewRound {
835 round,
836 head: "deadbeef".to_owned(),
837 verified_head: None,
838 verified_at: None,
839 reviews: vec![crate::run::ReviewRecord {
840 attempts: 0,
841 reviewer: 1,
842 agent: "mock".to_owned(),
843 summary: String::new(),
844 findings: vec![crate::verdict::Finding {
845 id: finding_id.to_owned(),
846 severity: crate::verdict::Severity::Major,
847 file: None,
848 line: None,
849 title: title.to_owned(),
850 detail: String::new(),
851 }],
852 vote: None,
853 failed: None,
854 duration_ms: 0,
855 }],
856 e2e: Vec::new(),
857 verify_retried: false,
858 e2e_deferred: false,
859 e2e_defer_reason: None,
860 fix: Some(crate::run::FixRecord {
861 agent: "mock".to_owned(),
862 addressed: addressed.iter().map(|s| (*s).to_owned()).collect(),
863 rejected: rejected
864 .iter()
865 .map(|(id, why)| crate::verdict::Rejection {
866 id: (*id).to_owned(),
867 why: (*why).to_owned(),
868 })
869 .collect(),
870 notes: String::new(),
871 committed: false,
872 failed: None,
873 duration_ms: 0,
874 continuation: None,
875 }),
876 blocking: 1,
877 answered: 1,
878 expected: 1,
879 clean: false,
880 progressed: true,
881 vote_split: false,
882 reconsideration: Vec::new(),
883 verdict: None,
884 }
885 }
886
887 #[test]
888 fn outcome_for_carries_every_rounds_findings_and_the_branch_head() {
889 crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
890 let dir = tempdir().unwrap();
891 let default_repo = dir.path().join("default");
892 let task_repo = dir.path().join("task");
893 std::fs::create_dir_all(&default_repo).unwrap();
894 std::fs::create_dir_all(&task_repo).unwrap();
895 init_repo_with_branch(&default_repo, "other-branch");
896 init_repo_with_branch(&task_repo, "magi/f00d/A");
897
898 let mut config = Config::default();
899 config.graph.review_rounds = 6;
900 let mut state = crate::run::RunState::new(
901 task_repo.clone(),
902 "main".to_owned(),
903 "deadbeef".to_owned(),
904 "task".to_owned(),
905 config,
906 );
907 state.status = crate::run::RunStatus::Blocked;
908 state.candidates.push(crate::run::Candidate {
909 index: 0,
910 label: 'A',
911 agent: "mock".to_owned(),
912 branch: "magi/f00d/A".to_owned(),
913 worktree: task_repo.clone(),
914 summary: String::new(),
915 stat: String::new(),
916 files: 1,
917 commits: 1,
918 empty: false,
919 failed: None,
920 verified_noop: None,
921 duration_ms: 0,
922 folded: false,
923 });
924 state.tally = Some(crate::run::Tally {
925 first_choice: std::collections::BTreeMap::new(),
926 borda: std::collections::BTreeMap::new(),
927 winner: 'A',
928 rankings: 0,
929 unanimous_initial: false,
930 deliberated: false,
931 changed_votes: 0,
932 unanimous_final: false,
933 tie_break: None,
934 judges: 0,
935 present: 0,
936 quorum: 0,
937 met_quorum: true,
938 uncontested: Some("solo".to_owned()),
939 });
940 state.reviews = vec![
941 review_round_with_finding(
942 1,
943 "R1-1-2",
944 "answer content is dropped",
945 &[],
946 &[("R1-1-2", "the id leaving blocked_by is enough")],
947 ),
948 review_round_with_finding(2, "R2-1-3", "answer content is still dropped", &[], &[]),
949 ];
950 state.save().unwrap();
951
952 let mut t = task("outcome test");
953 t.repo = task_repo;
954 t.runs.push(state.id.clone());
955
956 let finished = tokio_test_block_on(finished_view(&t, &default_repo, 2));
957 let outcome = finished.outcome;
958
959 assert!(outcome.unreadable.is_none());
960 assert_eq!(outcome.run_status.as_deref(), Some("blocked"));
961 assert_eq!(outcome.rounds_used, 2);
962 assert_eq!(outcome.rounds_max, 6);
963 assert_eq!(outcome.rounds.len(), 2);
964 assert_eq!(outcome.rounds[0].findings[0].id, "R1-1-2");
965 assert_eq!(outcome.rounds[0].rejected[0].id, "R1-1-2");
966 assert!(outcome.rounds[1].addressed.is_empty());
967 assert!(outcome.rounds[1].rejected.is_empty());
968 assert_eq!(outcome.branch.as_deref(), Some("magi/f00d/A"));
969 assert!(
970 outcome.branch_head.is_some(),
971 "a real branch must resolve a head commit: {outcome:?}"
972 );
973 }
974
975 fn tokio_test_block_on<F: std::future::Future>(f: F) -> F::Output {
979 tokio::runtime::Builder::new_current_thread()
980 .enable_all()
981 .build()
982 .unwrap()
983 .block_on(f)
984 }
985
986 #[test]
987 fn view_carries_a_tasks_recorded_answers_into_the_conductor_prompt_input() {
988 let mut t = task("answered");
989 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
990 let v = view(&t, 2);
991 assert_eq!(v.answers.len(), 1);
992 assert_eq!(v.answers[0].question, "Which backend?");
993 assert_eq!(v.answers[0].answer, "SQLite");
994 }
995
996 #[test]
997 fn a_dependency_decision_blocks_the_task_and_leaves_priority_alone() {
998 let dir = tempdir().unwrap();
999 let queue = Queue::at(dir.path().join("queue"));
1000 let questions = Questions::at(dir.path().join("questions"));
1001 let mut a = task("a");
1002 a.priority = 9;
1003 queue.put(&mut a).unwrap();
1004
1005 let verdict = Verdict {
1006 decisions: vec![Decision {
1007 id: a.id.clone(),
1008 blocked_by: vec!["20260101-000000-dead".to_owned()],
1009 reason: Some("waits on the other task".to_owned()),
1010 recovery: None,
1011 question: None,
1012 choices: Vec::new(),
1013 }],
1014 };
1015 apply(&queue, &questions, &verdict).unwrap();
1016
1017 let back = queue.get(&a.id).unwrap();
1018 assert_eq!(back.status, TaskStatus::Blocked);
1019 assert_eq!(back.blocked_by, ["20260101-000000-dead"]);
1020 assert_eq!(
1021 back.priority, 9,
1022 "the conductor's reply cannot carry priority"
1023 );
1024 }
1025
1026 #[test]
1027 fn a_question_decision_files_one_and_blocks_on_its_id() {
1028 let dir = tempdir().unwrap();
1029 let queue = Queue::at(dir.path().join("queue"));
1030 let questions = Questions::at(dir.path().join("questions"));
1031 let mut t = task("ambiguous");
1032 queue.put(&mut t).unwrap();
1033
1034 let verdict = Verdict {
1035 decisions: vec![Decision {
1036 id: t.id.clone(),
1037 blocked_by: Vec::new(),
1038 reason: Some("which backend?".to_owned()),
1039 recovery: None,
1040 question: Some("Which storage backend?".to_owned()),
1041 choices: vec!["SQLite".to_owned(), "Redis".to_owned()],
1042 }],
1043 };
1044 apply(&queue, &questions, &verdict).unwrap();
1045
1046 let back = queue.get(&t.id).unwrap();
1047 assert_eq!(back.status, TaskStatus::Blocked);
1048 assert_eq!(back.blocked_by.len(), 1);
1049 let q = questions.get(&back.blocked_by[0]).unwrap();
1050 assert_eq!(q.summary, "Which storage backend?");
1051 assert_eq!(q.node, NODE);
1052 assert!(q.status.open());
1053 }
1054
1055 #[test]
1056 fn a_task_with_an_open_question_already_reuses_it_rather_than_filing_a_second_one() {
1057 let dir = tempdir().unwrap();
1058 let queue = Queue::at(dir.path().join("queue"));
1059 let questions = Questions::at(dir.path().join("questions"));
1060 let mut t = task("asked once");
1061 queue.put(&mut t).unwrap();
1062
1063 let decision = Decision {
1064 id: t.id.clone(),
1065 reason: Some("still deciding".to_owned()),
1066 question: Some("Which backend?".to_owned()),
1067 ..Decision::default()
1068 };
1069 apply(
1070 &queue,
1071 &questions,
1072 &Verdict {
1073 decisions: vec![decision.clone()],
1074 },
1075 )
1076 .unwrap();
1077 assert_eq!(questions.list().len(), 1);
1078 let first_question_id = queue.get(&t.id).unwrap().blocked_by[0].clone();
1079
1080 let mut released = queue.get(&t.id).unwrap();
1085 released.release();
1086 queue.put(&mut released).unwrap();
1087
1088 apply(
1089 &queue,
1090 &questions,
1091 &Verdict {
1092 decisions: vec![decision],
1093 },
1094 )
1095 .unwrap();
1096 assert_eq!(questions.list().len(), 1, "no duplicate question was filed");
1097 let after = queue.get(&t.id).unwrap();
1098 assert_eq!(
1099 after.blocked_by,
1100 [first_question_id],
1101 "the existing open question is reused, not replaced"
1102 );
1103 }
1104
1105 #[test]
1106 fn a_same_id_question_from_another_node_is_not_reused() {
1107 let dir = tempdir().unwrap();
1108 let queue = Queue::at(dir.path().join("queue"));
1109 let questions = Questions::at(dir.path().join("questions"));
1110 let mut t = task("must ask the conductor");
1111 queue.put(&mut t).unwrap();
1112
1113 let mut unrelated = Question::new(
1114 t.id.clone(),
1115 "review".to_owned(),
1116 "reviewer-1".to_owned(),
1117 "An unrelated review question".to_owned(),
1118 String::new(),
1119 Vec::new(),
1120 );
1121 questions.put(&mut unrelated).unwrap();
1122
1123 apply(
1124 &queue,
1125 &questions,
1126 &Verdict {
1127 decisions: vec![Decision {
1128 id: t.id.clone(),
1129 question: Some("Which backend?".to_owned()),
1130 ..Decision::default()
1131 }],
1132 },
1133 )
1134 .unwrap();
1135
1136 let blocked_by = &queue.get(&t.id).unwrap().blocked_by;
1137 assert_eq!(blocked_by.len(), 1);
1138 assert_ne!(blocked_by[0], unrelated.id);
1139 assert!(questions.get(&unrelated.id).unwrap().status.open());
1140 assert_eq!(questions.get(&blocked_by[0]).unwrap().node, NODE);
1141 }
1142
1143 #[test]
1144 fn answering_the_question_lets_the_resolver_clear_the_block_with_the_answer_kept() {
1145 let dir = tempdir().unwrap();
1146 let queue = Queue::at(dir.path().join("queue"));
1147 let questions = Questions::at(dir.path().join("questions"));
1148 let mut t = task("waits on an answer");
1149 queue.put(&mut t).unwrap();
1150
1151 apply(
1152 &queue,
1153 &questions,
1154 &Verdict {
1155 decisions: vec![Decision {
1156 id: t.id.clone(),
1157 blocked_by: Vec::new(),
1158 reason: None,
1159 recovery: None,
1160 question: Some("Which backend?".to_owned()),
1161 choices: Vec::new(),
1162 }],
1163 },
1164 )
1165 .unwrap();
1166 let blocked = queue.get(&t.id).unwrap();
1167 let question_id = blocked.blocked_by[0].clone();
1168
1169 let mut q = questions.get(&question_id).unwrap();
1170 q.answer(Answer::Text("SQLite".to_owned())).unwrap();
1171 questions.put(&mut q).unwrap();
1172 assert_eq!(q.status, QuestionStatus::Answered);
1173
1174 let mut task_after = queue.get(&t.id).unwrap();
1178 task_after.record_answer(q.summary.clone(), "SQLite".to_owned());
1179 task_after.unblock(&question_id);
1180 assert_eq!(task_after.status, TaskStatus::Queued);
1181 assert_eq!(task_after.answers[0].answer, "SQLite");
1182 }
1183
1184 #[test]
1185 fn a_stalled_task_can_be_requeued_or_held() {
1186 let dir = tempdir().unwrap();
1187 let queue = Queue::at(dir.path().join("queue"));
1188 let questions = Questions::at(dir.path().join("questions"));
1189
1190 let mut requeue_me = task("stuck a");
1191 requeue_me.start("run-1".to_owned());
1192 queue.put(&mut requeue_me).unwrap();
1193
1194 let mut hold_me = task("stuck b");
1195 hold_me.start("run-2".to_owned());
1196 queue.put(&mut hold_me).unwrap();
1197
1198 apply(
1199 &queue,
1200 &questions,
1201 &Verdict {
1202 decisions: vec![
1203 Decision {
1204 id: requeue_me.id.clone(),
1205 recovery: Some(Recovery::Requeue),
1206 ..Decision::default()
1207 },
1208 Decision {
1209 id: hold_me.id.clone(),
1210 recovery: Some(Recovery::Hold),
1211 reason: Some("looks broken".to_owned()),
1212 ..Decision::default()
1213 },
1214 ],
1215 },
1216 )
1217 .unwrap();
1218
1219 let requeued = queue.get(&requeue_me.id).unwrap();
1220 assert_eq!(requeued.status, TaskStatus::Queued);
1221 assert_eq!(requeued.attempts, 0);
1222
1223 let held = queue.get(&hold_me.id).unwrap();
1224 assert_eq!(held.status, TaskStatus::Held);
1225 assert_eq!(held.hold_reason.as_deref(), Some("looks broken"));
1226 }
1227
1228 #[test]
1229 fn a_machine_held_task_asked_about_restores_to_held_once_answered() {
1230 let dir = tempdir().unwrap();
1235 let queue = Queue::at(dir.path().join("queue"));
1236 let questions = Questions::at(dir.path().join("questions"));
1237 let mut t = task("held out of attempts");
1238 t.hold_machine(Some("out of attempts".to_owned()));
1239 queue.put(&mut t).unwrap();
1240
1241 apply(
1242 &queue,
1243 &questions,
1244 &Verdict {
1245 decisions: vec![Decision {
1246 id: t.id.clone(),
1247 reason: Some("what should happen to this one?".to_owned()),
1248 question: Some("Hold it, or try again?".to_owned()),
1249 ..Decision::default()
1250 }],
1251 },
1252 )
1253 .unwrap();
1254 let blocked = queue.get(&t.id).unwrap();
1255 assert_eq!(blocked.status, TaskStatus::Blocked);
1256 let question_id = blocked.blocked_by[0].clone();
1257
1258 let mut q = questions.get(&question_id).unwrap();
1259 q.answer(Answer::Text("leave it held".to_owned())).unwrap();
1260 questions.put(&mut q).unwrap();
1261
1262 let mut after = queue.get(&t.id).unwrap();
1264 after.record_answer(q.summary.clone(), "leave it held".to_owned());
1265 after.unblock(&question_id);
1266 assert_eq!(
1267 after.status,
1268 TaskStatus::Held,
1269 "must not fall back to queued"
1270 );
1271 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
1272 }
1273
1274 #[test]
1275 fn a_reaffirmed_hold_with_no_new_reason_is_not_silently_auto_released_by_triage() {
1276 let dir = tempdir().unwrap();
1286 let queue = Queue::at(dir.path().join("queue"));
1287 let questions = Questions::at(dir.path().join("questions"));
1288 let mut t = task("disk pressure, then reconsidered");
1289 t.hold_machine(Some(
1290 "not enough free space to start a run: 10 bytes free, 100 required by \
1291 `[disk] min_free_bytes`"
1292 .to_owned(),
1293 ));
1294 t.record_answer(
1295 "How should this be handled?".to_owned(),
1296 "keep it held, a human will look at it later".to_owned(),
1297 );
1298 queue.put(&mut t).unwrap();
1299
1300 apply(
1301 &queue,
1302 &questions,
1303 &Verdict {
1304 decisions: vec![Decision {
1305 id: t.id.clone(),
1306 recovery: Some(Recovery::Hold),
1307 ..Decision::default()
1308 }],
1309 },
1310 )
1311 .unwrap();
1312
1313 let after = queue.get(&t.id).unwrap();
1314 assert_eq!(after.status, TaskStatus::Held);
1315 assert!(
1316 !after
1317 .hold_reason
1318 .as_deref()
1319 .unwrap_or_default()
1320 .starts_with("not enough free space"),
1321 "the stale disk-pressure text must not survive a reconfirmed hold: {:?}",
1322 after.hold_reason
1323 );
1324
1325 let cfg_dir = tempdir().unwrap();
1328 let config = cfg_dir.path().join("magi.toml");
1329 std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1330 let report =
1331 crate::triage::run_once(&queue, &questions, Some(&config), jiff::Timestamp::now());
1332 assert!(report.resumed.is_empty(), "must not be auto-released");
1333 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1334 }
1335
1336 #[test]
1337 fn a_conductor_rehold_after_a_resume_answer_is_not_asked_again_identically() {
1338 let dir = tempdir().unwrap();
1339 let queue = Queue::at(dir.path().join("queue"));
1340 let questions = Questions::at(dir.path().join("questions"));
1341 let config = dir.path().join("magi.toml");
1342 std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1343 let now = jiff::Timestamp::now();
1344 let triage = || crate::triage::run_once(&queue, &questions, Some(&config), now);
1345 let open = || {
1346 questions
1347 .list()
1348 .into_iter()
1349 .filter(|q| q.node == "triage" && q.status.open())
1350 .collect::<Vec<_>>()
1351 };
1352 let hold = Decision {
1353 recovery: Some(Recovery::Hold),
1354 reason: Some("waiting on manual worktree cleanup".to_owned()),
1355 ..Decision::default()
1356 };
1357
1358 let mut t = task("looping hold");
1359 t.hold_machine(Some("waiting on manual worktree cleanup".to_owned()));
1360 queue.put(&mut t).unwrap();
1361 let hold = Decision {
1362 id: t.id.clone(),
1363 ..hold
1364 };
1365
1366 assert_eq!(triage().asked.len(), 1);
1368 let first = open().remove(0);
1369 let mut q = questions.get(&first.id).unwrap();
1370 let resume = q.choices[0].clone();
1371 q.answer(Answer::Choice(resume)).unwrap();
1372 questions.put(&mut q).unwrap();
1373 assert_eq!(triage().answered.len(), 1);
1374 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1375
1376 apply(
1378 &queue,
1379 &questions,
1380 &Verdict {
1381 decisions: vec![hold.clone()],
1382 },
1383 )
1384 .unwrap();
1385 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1386
1387 assert_eq!(triage().asked.len(), 1);
1389 let second = open().remove(0);
1390 assert_ne!(second.summary, first.summary);
1391 assert_ne!(second.choices, first.choices);
1392 assert!(second.detail.contains("waiting on manual worktree cleanup"));
1393 assert!(triage().asked.is_empty(), "no duplicate question");
1394 assert_eq!(open().len(), 1);
1395
1396 let mut q = questions.get(&second.id).unwrap();
1398 let force = q.choices[0].clone();
1399 q.answer(Answer::Choice(force)).unwrap();
1400 questions.put(&mut q).unwrap();
1401 assert_eq!(triage().answered.len(), 1);
1402 apply(
1403 &queue,
1404 &questions,
1405 &Verdict {
1406 decisions: vec![hold],
1407 },
1408 )
1409 .unwrap();
1410 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1411 }
1412
1413 #[test]
1414 fn a_runnable_task_can_be_held_directly_without_a_question() {
1415 let dir = tempdir().unwrap();
1416 let queue = Queue::at(dir.path().join("queue"));
1417 let questions = Questions::at(dir.path().join("questions"));
1418 let mut t = task("already answered, should stay put");
1419 queue.put(&mut t).unwrap();
1420
1421 apply(
1422 &queue,
1423 &questions,
1424 &Verdict {
1425 decisions: vec![Decision {
1426 id: t.id.clone(),
1427 recovery: Some(Recovery::Hold),
1428 reason: Some("operator already said keep this held".to_owned()),
1429 ..Decision::default()
1430 }],
1431 },
1432 )
1433 .unwrap();
1434
1435 let after = queue.get(&t.id).unwrap();
1436 assert_eq!(after.status, TaskStatus::Held);
1437 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
1438 }
1439
1440 #[test]
1441 fn done_recovery_closes_a_held_task_whose_goal_is_already_met() {
1442 let dir = tempdir().unwrap();
1448 let queue = Queue::at(dir.path().join("queue"));
1449 let questions = Questions::at(dir.path().join("questions"));
1450 let mut t = task("already merged by hand");
1451 t.hold_machine(Some("branch survived, awaiting a decision".to_owned()));
1452 t.record_answer(
1453 "Handle this one?".to_owned(),
1454 "already merged and cleaned up, close it".to_owned(),
1455 );
1456 queue.put(&mut t).unwrap();
1457
1458 apply(
1459 &queue,
1460 &questions,
1461 &Verdict {
1462 decisions: vec![Decision {
1463 id: t.id.clone(),
1464 recovery: Some(Recovery::Done),
1465 reason: Some("operator confirmed this already landed".to_owned()),
1466 ..Decision::default()
1467 }],
1468 },
1469 )
1470 .unwrap();
1471
1472 let after = queue.get(&t.id).unwrap();
1473 assert_eq!(after.status, TaskStatus::Done);
1474 assert!(after.hold_reason.is_none());
1475 assert_eq!(after.answers.len(), 1, "the record of why is kept");
1476 }
1477
1478 #[test]
1479 fn done_recovery_is_ignored_for_a_runnable_or_running_task() {
1480 let dir = tempdir().unwrap();
1481 let queue = Queue::at(dir.path().join("queue"));
1482 let questions = Questions::at(dir.path().join("questions"));
1483
1484 let mut queued = task("never ran yet");
1485 queue.put(&mut queued).unwrap();
1486
1487 let mut running = task("mid-run");
1488 running.start("run-1".to_owned());
1489 queue.put(&mut running).unwrap();
1490
1491 for id in [queued.id.clone(), running.id.clone()] {
1492 apply(
1493 &queue,
1494 &questions,
1495 &Verdict {
1496 decisions: vec![Decision {
1497 id,
1498 recovery: Some(Recovery::Done),
1499 ..Decision::default()
1500 }],
1501 },
1502 )
1503 .unwrap();
1504 }
1505
1506 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
1507 assert_eq!(queue.get(&running.id).unwrap().status, TaskStatus::Running);
1508 }
1509
1510 #[test]
1511 fn a_third_conductor_question_after_two_settled_answers_holds_instead_of_asking_again() {
1512 let dir = tempdir().unwrap();
1518 let queue = Queue::at(dir.path().join("queue"));
1519 let questions = Questions::at(dir.path().join("questions"));
1520 let mut t = task("asked about repeatedly");
1521 t.hold_machine(Some("out of attempts".to_owned()));
1522 t.record_answer("Handle this one? (1)".to_owned(), "not yet".to_owned());
1523 t.record_answer(
1524 "Handle this one? (2)".to_owned(),
1525 "still not yet".to_owned(),
1526 );
1527 queue.put(&mut t).unwrap();
1528 assert_eq!(questions.list().len(), 0);
1529
1530 apply(
1531 &queue,
1532 &questions,
1533 &Verdict {
1534 decisions: vec![Decision {
1535 id: t.id.clone(),
1536 question: Some("Handle this one? (3)".to_owned()),
1537 ..Decision::default()
1538 }],
1539 },
1540 )
1541 .unwrap();
1542
1543 assert_eq!(questions.list().len(), 0, "no third question was filed");
1544 let after = queue.get(&t.id).unwrap();
1545 assert_eq!(after.status, TaskStatus::Held);
1546 assert!(after.blocked_by.is_empty());
1547 assert_eq!(after.answers.len(), 2, "the prior answers are untouched");
1548 }
1549
1550 #[test]
1551 fn a_second_conductor_question_is_still_allowed_after_one_settled_answer() {
1552 let dir = tempdir().unwrap();
1553 let queue = Queue::at(dir.path().join("queue"));
1554 let questions = Questions::at(dir.path().join("questions"));
1555 let mut t = task("asked about once already");
1556 t.hold_machine(Some("out of attempts".to_owned()));
1557 t.record_answer("Handle this one?".to_owned(), "not yet".to_owned());
1558 queue.put(&mut t).unwrap();
1559
1560 apply(
1561 &queue,
1562 &questions,
1563 &Verdict {
1564 decisions: vec![Decision {
1565 id: t.id.clone(),
1566 question: Some("Still not sure - now what?".to_owned()),
1567 ..Decision::default()
1568 }],
1569 },
1570 )
1571 .unwrap();
1572
1573 assert_eq!(questions.list().len(), 1, "the second question was filed");
1574 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
1575 }
1576
1577 #[test]
1578 fn a_held_task_with_an_open_triage_question_is_left_to_triage() {
1579 let dir = tempdir().unwrap();
1585 let queue = Queue::at(dir.path().join("queue"));
1586 let questions = Questions::at(dir.path().join("questions"));
1587 let mut t = task("held, triage already asking about it");
1588 t.hold_machine(Some("cause unclear".to_owned()));
1589 queue.put(&mut t).unwrap();
1590
1591 let mut triage_q = Question::new(
1592 t.id.clone(),
1593 crate::triage::NODE.to_owned(),
1594 "triage".to_owned(),
1595 "Still needed?".to_owned(),
1596 String::new(),
1597 vec![
1598 "resume".to_owned(),
1599 "not yet".to_owned(),
1600 "discard".to_owned(),
1601 ],
1602 );
1603 questions.put(&mut triage_q).unwrap();
1604
1605 for decision in [
1606 Decision {
1607 id: t.id.clone(),
1608 question: Some("what now?".to_owned()),
1609 ..Decision::default()
1610 },
1611 Decision {
1612 id: t.id.clone(),
1613 recovery: Some(Recovery::Requeue),
1614 ..Decision::default()
1615 },
1616 ] {
1617 apply(
1618 &queue,
1619 &questions,
1620 &Verdict {
1621 decisions: vec![decision],
1622 },
1623 )
1624 .unwrap();
1625 }
1626
1627 let after = queue.get(&t.id).unwrap();
1628 assert_eq!(
1629 after.status,
1630 TaskStatus::Held,
1631 "triage still owns this hold"
1632 );
1633 assert!(after.blocked_by.is_empty());
1634 assert_eq!(
1635 questions.list().len(),
1636 1,
1637 "no second, conductor-owned question was filed"
1638 );
1639 }
1640
1641 #[test]
1642 fn a_held_task_with_an_answered_but_unapplied_triage_question_is_still_left_alone() {
1643 let dir = tempdir().unwrap();
1651 let queue = Queue::at(dir.path().join("queue"));
1652 let questions = Questions::at(dir.path().join("questions"));
1653 let mut t = task("held, triage question answered but not yet applied");
1654 t.hold_machine(Some("cause unclear".to_owned()));
1655 queue.put(&mut t).unwrap();
1656
1657 let mut triage_q = Question::new(
1658 t.id.clone(),
1659 crate::triage::NODE.to_owned(),
1660 "triage".to_owned(),
1661 "Still needed?".to_owned(),
1662 String::new(),
1663 vec![
1664 "resume".to_owned(),
1665 "not yet".to_owned(),
1666 "discard".to_owned(),
1667 ],
1668 );
1669 questions.put(&mut triage_q).unwrap();
1670 triage_q
1671 .answer(Answer::Choice("not yet".to_owned()))
1672 .unwrap();
1673 questions.put(&mut triage_q).unwrap();
1674 assert!(!triage_q.status.open());
1675
1676 apply(
1677 &queue,
1678 &questions,
1679 &Verdict {
1680 decisions: vec![Decision {
1681 id: t.id.clone(),
1682 question: Some("what now?".to_owned()),
1683 ..Decision::default()
1684 }],
1685 },
1686 )
1687 .unwrap();
1688
1689 let after = queue.get(&t.id).unwrap();
1690 assert_eq!(
1691 after.status,
1692 TaskStatus::Held,
1693 "triage's own answer is not yet applied - conduct must wait"
1694 );
1695 assert_eq!(
1696 questions.list().len(),
1697 1,
1698 "no conductor question was filed over the pending triage answer"
1699 );
1700 }
1701
1702 #[test]
1703 fn manual_hold_rejects_hostile_or_stale_conductor_recovery() {
1704 let dir = tempdir().unwrap();
1705 let queue = Queue::at(dir.path().join("queue"));
1706 let questions = Questions::at(dir.path().join("questions"));
1707 let mut held = task("manual recovery");
1708 held.priority = 300;
1709 held.runs.push("run20260912-224242-daf5".to_owned());
1710 held.hold_manual(Some(
1711 "active manual recovery run20260912-224242-daf5".to_owned(),
1712 ));
1713 queue.put(&mut held).unwrap();
1714
1715 for decision in [
1719 Decision {
1720 id: held.id.clone(),
1721 recovery: Some(Recovery::Requeue),
1722 ..Decision::default()
1723 },
1724 Decision {
1725 id: held.id.clone(),
1726 recovery: Some(Recovery::Hold),
1727 reason: Some("stale replacement reason".to_owned()),
1728 ..Decision::default()
1729 },
1730 Decision {
1731 id: held.id.clone(),
1732 recovery: Some(Recovery::Review),
1733 ..Decision::default()
1734 },
1735 Decision {
1736 id: held.id.clone(),
1737 blocked_by: vec!["other-task".to_owned()],
1738 question: Some("retry now?".to_owned()),
1739 ..Decision::default()
1740 },
1741 ] {
1742 apply(
1743 &queue,
1744 &questions,
1745 &Verdict {
1746 decisions: vec![decision],
1747 },
1748 )
1749 .unwrap();
1750 }
1751
1752 let after = queue.get(&held.id).unwrap();
1753 assert_eq!(after.status, TaskStatus::Held);
1754 assert!(after.operator_held());
1755 assert_eq!(after.priority, 300);
1756 assert_eq!(after.runs, ["run20260912-224242-daf5"]);
1757 assert_eq!(
1758 after.hold_reason.as_deref(),
1759 Some("active manual recovery run20260912-224242-daf5")
1760 );
1761 assert!(after.blocked_by.is_empty());
1762 assert!(questions.list().is_empty());
1763 assert!(
1764 queue.next_runnable().is_none(),
1765 "must not dispatch a duplicate"
1766 );
1767 }
1768
1769 #[test]
1770 fn machine_holds_remain_recoverable_and_manual_release_is_authorization() {
1771 let dir = tempdir().unwrap();
1772 let queue = Queue::at(dir.path().join("queue"));
1773 let questions = Questions::at(dir.path().join("questions"));
1774
1775 let mut automatic = task("disk gate");
1776 automatic.hold_machine(Some("disk full".to_owned()));
1777 queue.put(&mut automatic).unwrap();
1778 let requeue = || Verdict {
1779 decisions: vec![Decision {
1780 id: automatic.id.clone(),
1781 recovery: Some(Recovery::Requeue),
1782 ..Decision::default()
1783 }],
1784 };
1785 apply(&queue, &questions, &requeue()).unwrap();
1786 assert_eq!(queue.get(&automatic.id).unwrap().status, TaskStatus::Queued);
1787
1788 let mut manual = task("operator gate");
1789 manual.hold_manual(Some("wait for operator".to_owned()));
1790 queue.put(&mut manual).unwrap();
1791 apply(
1792 &queue,
1793 &questions,
1794 &Verdict {
1795 decisions: vec![Decision {
1796 id: manual.id.clone(),
1797 recovery: Some(Recovery::Requeue),
1798 ..Decision::default()
1799 }],
1800 },
1801 )
1802 .unwrap();
1803 assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Held);
1804
1805 let mut released = queue.get(&manual.id).unwrap();
1808 released.release();
1809 queue.put(&mut released).unwrap();
1810 assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Queued);
1811 }
1812
1813 #[test]
1814 fn legacy_reasoned_hold_is_protected_without_losing_its_metadata() {
1815 let dir = tempdir().unwrap();
1816 let queue = Queue::at(dir.path().join("queue"));
1817 let questions = Questions::at(dir.path().join("questions"));
1818 let mut legacy = task("old explicit hold");
1819 legacy.status = TaskStatus::Held;
1820 legacy.hold_reason = Some("manual recovery already active".to_owned());
1821 legacy.hold_source = None;
1822 legacy.blocked_by = vec!["dependency".to_owned()];
1823 queue.put(&mut legacy).unwrap();
1824
1825 apply(
1826 &queue,
1827 &questions,
1828 &Verdict {
1829 decisions: vec![Decision {
1830 id: legacy.id.clone(),
1831 recovery: Some(Recovery::Requeue),
1832 ..Decision::default()
1833 }],
1834 },
1835 )
1836 .unwrap();
1837
1838 let after = queue.get(&legacy.id).unwrap();
1839 assert_eq!(after.status, TaskStatus::Held);
1840 assert_eq!(after.hold_source, None);
1841 assert_eq!(after.hold_reason, legacy.hold_reason);
1842 assert_eq!(after.blocked_by, legacy.blocked_by);
1843 }
1844
1845 #[test]
1846 fn review_recovery_is_a_no_op_without_a_survivable_branch() {
1847 crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
1853 let dir = tempdir().unwrap();
1854 let queue = Queue::at(dir.path().join("queue"));
1855 let questions = Questions::at(dir.path().join("questions"));
1856 let mut t = task("blocked with no readable run");
1857 t.start("20260101-000000-dead".to_owned()); t.fail("blocked", 5);
1859 queue.put(&mut t).unwrap();
1860
1861 apply(
1862 &queue,
1863 &questions,
1864 &Verdict {
1865 decisions: vec![Decision {
1866 id: t.id.clone(),
1867 recovery: Some(Recovery::Review),
1868 ..Decision::default()
1869 }],
1870 },
1871 )
1872 .unwrap();
1873
1874 let after = queue.get(&t.id).unwrap();
1875 assert_eq!(
1876 after.status,
1877 TaskStatus::Failed,
1878 "with nothing to reopen, the decision is dropped rather than guessed at"
1879 );
1880 assert!(after.review_branch.is_none());
1881 }
1882
1883 #[test]
1884 fn requeue_and_review_recovery_are_ignored_for_a_runnable_task() {
1885 let dir = tempdir().unwrap();
1889 let queue = Queue::at(dir.path().join("queue"));
1890 let questions = Questions::at(dir.path().join("questions"));
1891
1892 for recovery in [Recovery::Requeue, Recovery::Review] {
1893 let mut t = task("ordinary");
1894 queue.put(&mut t).unwrap();
1895
1896 apply(
1897 &queue,
1898 &questions,
1899 &Verdict {
1900 decisions: vec![Decision {
1901 id: t.id.clone(),
1902 recovery: Some(recovery),
1903 ..Decision::default()
1904 }],
1905 },
1906 )
1907 .unwrap();
1908
1909 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1910 }
1911 }
1912
1913 #[tokio::test]
1914 async fn a_broken_agent_leaves_the_queue_untouched_and_does_not_error() {
1915 let dir = tempdir().unwrap();
1916 let cfg = config(mock_agent(dir.path(), BROKEN, BTreeMap::new()));
1917 let queue = Queue::at(dir.path().join("queue"));
1918 let questions = Questions::at(dir.path().join("questions"));
1919 let mut t = task("normal");
1920 queue.put(&mut t).unwrap();
1921
1922 let mut conductor = Conductor::new();
1923 conductor
1924 .maybe_run(
1925 &cfg,
1926 dir.path(),
1927 &queue,
1928 &questions,
1929 dir.path(),
1930 &[t.clone()],
1931 &[],
1932 &[],
1933 2,
1934 )
1935 .await;
1936
1937 assert_eq!(
1938 queue.get(&t.id).unwrap().status,
1939 TaskStatus::Queued,
1940 "a failed invocation must change nothing"
1941 );
1942 assert!(
1943 queue.next_runnable().is_some(),
1944 "the loop must still be able to take the next task"
1945 );
1946 }
1947
1948 #[tokio::test]
1949 async fn a_reply_with_no_json_leaves_the_queue_untouched() {
1950 let dir = tempdir().unwrap();
1951 let cfg = config(mock_agent(dir.path(), GARBAGE, BTreeMap::new()));
1952 let queue = Queue::at(dir.path().join("queue"));
1953 let questions = Questions::at(dir.path().join("questions"));
1954 let mut t = task("normal");
1955 queue.put(&mut t).unwrap();
1956
1957 let mut conductor = Conductor::new();
1958 conductor
1959 .maybe_run(
1960 &cfg,
1961 dir.path(),
1962 &queue,
1963 &questions,
1964 dir.path(),
1965 &[t.clone()],
1966 &[],
1967 &[],
1968 2,
1969 )
1970 .await;
1971
1972 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1973 }
1974
1975 #[tokio::test]
1976 async fn json_survives_code_fences_and_a_preamble() {
1977 let dir = tempdir().unwrap();
1978 let mut t = task("fenced");
1979 let reply = format!(
1980 "Sure, here is my decision.\n\n```json\n{{\"decisions\":[{{\"id\":\"{}\",\
1981 \"blocked_by\":[\"x\"],\"reason\":\"why\"}}]}}\n```\n",
1982 t.id
1983 );
1984 let cfg = config(mock_agent(dir.path(), REPLY, env(&reply)));
1985 let queue = Queue::at(dir.path().join("queue"));
1986 let questions = Questions::at(dir.path().join("questions"));
1987 queue.put(&mut t).unwrap();
1988
1989 let mut conductor = Conductor::new();
1990 conductor
1991 .maybe_run(
1992 &cfg,
1993 dir.path(),
1994 &queue,
1995 &questions,
1996 dir.path(),
1997 &[t.clone()],
1998 &[],
1999 &[],
2000 2,
2001 )
2002 .await;
2003
2004 let back = queue.get(&t.id).unwrap();
2005 assert_eq!(back.status, TaskStatus::Blocked);
2006 assert_eq!(back.blocked_by, ["x"]);
2007 }
2008
2009 #[tokio::test]
2010 async fn the_conductor_is_not_called_again_when_nothing_worth_looking_at_has_changed() {
2011 let dir = tempdir().unwrap();
2014 let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2015 let queue = Queue::at(dir.path().join("queue"));
2016 let questions = Questions::at(dir.path().join("questions"));
2017 let mut t = task("stable");
2018 queue.put(&mut t).unwrap();
2019 let artifacts = dir.path().join("conduct").join("artifacts");
2020 let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2021
2022 let mut conductor = Conductor::new();
2023 conductor
2024 .maybe_run(
2025 &cfg,
2026 dir.path(),
2027 &queue,
2028 &questions,
2029 dir.path(),
2030 &[t.clone()],
2031 &[],
2032 &[],
2033 2,
2034 )
2035 .await;
2036 assert!(turn(1).is_file(), "the first cycle must call the conductor");
2037
2038 conductor
2039 .maybe_run(
2040 &cfg,
2041 dir.path(),
2042 &queue,
2043 &questions,
2044 dir.path(),
2045 &[t.clone()],
2046 &[],
2047 &[],
2048 2,
2049 )
2050 .await;
2051 assert!(
2052 !turn(2).is_file(),
2053 "an unchanged revision and an unchanged stalled/finished set must not call the \
2054 conductor twice"
2055 );
2056
2057 t.priority = 1;
2059 queue.put(&mut t).unwrap();
2060 conductor
2061 .maybe_run(
2062 &cfg,
2063 dir.path(),
2064 &queue,
2065 &questions,
2066 dir.path(),
2067 &[t.clone()],
2068 &[],
2069 &[],
2070 2,
2071 )
2072 .await;
2073 assert!(turn(2).is_file(), "a moved revision calls it again");
2074 }
2075
2076 #[tokio::test]
2077 async fn a_task_turning_stalled_calls_the_conductor_again_despite_an_unchanged_revision() {
2078 let dir = tempdir().unwrap();
2084 let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2085 let queue = Queue::at(dir.path().join("queue"));
2086 let questions = Questions::at(dir.path().join("questions"));
2087 let mut t = task("quiet");
2088 queue.put(&mut t).unwrap();
2089 let artifacts = dir.path().join("conduct").join("artifacts");
2090 let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2091
2092 let mut conductor = Conductor::new();
2093 conductor
2094 .maybe_run(
2095 &cfg,
2096 dir.path(),
2097 &queue,
2098 &questions,
2099 dir.path(),
2100 &[t.clone()],
2101 &[],
2102 &[],
2103 2,
2104 )
2105 .await;
2106 assert!(turn(1).is_file());
2107
2108 conductor
2109 .maybe_run(
2110 &cfg,
2111 dir.path(),
2112 &queue,
2113 &questions,
2114 dir.path(),
2115 &[],
2116 &[t.clone()],
2117 &[],
2118 2,
2119 )
2120 .await;
2121 assert!(
2122 turn(2).is_file(),
2123 "a task turning stalled must call the conductor again"
2124 );
2125
2126 conductor
2129 .maybe_run(
2130 &cfg,
2131 dir.path(),
2132 &queue,
2133 &questions,
2134 dir.path(),
2135 &[],
2136 &[t.clone()],
2137 &[],
2138 2,
2139 )
2140 .await;
2141 assert!(
2142 !turn(3).is_file(),
2143 "the same stalled task lingering must not call the conductor every cycle"
2144 );
2145 }
2146
2147 #[test]
2148 fn worth_a_look_is_config_free_and_matches_maybe_runs_own_gate() {
2149 let dir = tempdir().unwrap();
2150 let queue = Queue::at(dir.path().join("queue"));
2151 let mut t = task("t");
2152 queue.put(&mut t).unwrap();
2153
2154 let mut conductor = Conductor::new();
2155 assert!(
2156 conductor.worth_a_look(&queue, &[], &[]),
2157 "a conductor that has never run has something to look at"
2158 );
2159
2160 conductor.last_seen = Some(Conductor::snapshot(&queue, &[], &[]));
2161 assert!(
2162 !conductor.worth_a_look(&queue, &[], &[]),
2163 "nothing changed and nothing is stalled or finished"
2164 );
2165 assert!(
2166 conductor.worth_a_look(&queue, &[t.clone()], &[]),
2167 "a stalled task is worth a look even at the same revision"
2168 );
2169 assert!(
2170 conductor.worth_a_look(&queue, &[], &[t.clone()]),
2171 "a finished task is worth a look even at the same revision"
2172 );
2173 }
2174
2175 #[tokio::test]
2176 async fn the_conduct_path_never_calls_ask_and_wait() {
2177 let dir = tempdir().unwrap();
2183 let queue = Queue::at(dir.path().join("queue"));
2184 let questions = Questions::at(dir.path().join("questions"));
2185 let mut t = task("asks without blocking");
2186 queue.put(&mut t).unwrap();
2187
2188 apply(
2189 &queue,
2190 &questions,
2191 &Verdict {
2192 decisions: vec![Decision {
2193 id: t.id.clone(),
2194 question: Some("ok?".to_owned()),
2195 ..Decision::default()
2196 }],
2197 },
2198 )
2199 .unwrap();
2200 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
2202 }
2203}