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
576pub fn seat_path(home: &Path) -> PathBuf {
582 home.join("conduct").join("seat.json")
583}
584
585pub fn load_seat(home: &Path) -> Option<SeatState> {
587 serde_json::from_str(&std::fs::read_to_string(seat_path(home)).ok()?).ok()
588}
589
590fn busy_path(home: &Path) -> PathBuf {
591 home.join("conduct").join("busy")
592}
593
594pub fn busy(home: &Path) -> bool {
600 std::fs::metadata(busy_path(home))
601 .and_then(|m| m.modified())
602 .ok()
603 .and_then(|t| t.elapsed().ok())
604 .is_some_and(|age| age < TURN_TIMEOUT + Duration::from_secs(30))
605}
606
607struct Busy(PathBuf);
609
610impl Busy {
611 fn mark(home: &Path) -> Self {
612 let path = busy_path(home);
613 if let Some(dir) = path.parent() {
614 let _ = std::fs::create_dir_all(dir);
615 }
616 let _ = std::fs::write(&path, std::process::id().to_string());
617 Self(path)
618 }
619}
620
621impl Drop for Busy {
622 fn drop(&mut self) {
623 let _ = std::fs::remove_file(&self.0);
624 }
625}
626
627#[derive(Debug, Default)]
630pub struct Conductor {
631 seat: Option<SeatState>,
632 last_seen: Option<(u64, BTreeSet<String>)>,
633}
634
635impl Conductor {
636 #[must_use]
638 pub fn new() -> Self {
639 Self::default()
640 }
641
642 fn snapshot(queue: &Queue, stalled: &[Task], finished: &[Task]) -> (u64, BTreeSet<String>) {
643 let ids = stalled
644 .iter()
645 .chain(finished)
646 .map(|t| t.id.clone())
647 .collect();
648 (queue.revision(), ids)
649 }
650
651 #[must_use]
668 pub fn worth_a_look(&self, queue: &Queue, stalled: &[Task], finished: &[Task]) -> bool {
669 self.last_seen.as_ref() != Some(&Self::snapshot(queue, stalled, finished))
670 }
671
672 #[allow(clippy::too_many_arguments)]
676 pub async fn maybe_run(
677 &mut self,
678 cfg: &Config,
679 repo: &Path,
680 queue: &Queue,
681 questions: &Questions,
682 home: &Path,
683 queued: &[Task],
684 stalled: &[Task],
685 finished: &[Task],
686 max_attempts: usize,
687 ) {
688 let snapshot = Self::snapshot(queue, stalled, finished);
689 if self.last_seen.as_ref() == Some(&snapshot) {
690 return;
691 }
692 self.last_seen = Some(snapshot);
693 if let Err(e) = self
694 .run_once(
695 cfg,
696 repo,
697 queue,
698 questions,
699 home,
700 queued,
701 stalled,
702 finished,
703 max_attempts,
704 )
705 .await
706 {
707 tracing::warn!("conductor: {e:#}");
708 }
709 }
710
711 #[allow(clippy::too_many_arguments)]
712 async fn run_once(
713 &mut self,
714 cfg: &Config,
715 repo: &Path,
716 queue: &Queue,
717 questions: &Questions,
718 home: &Path,
719 queued: &[Task],
720 stalled: &[Task],
721 finished: &[Task],
722 max_attempts: usize,
723 ) -> Result<()> {
724 if queued.is_empty() && stalled.is_empty() && finished.is_empty() {
725 return Ok(());
726 }
727
728 let spec = cfg
729 .resolve_roles()
730 .context("resolving the conductor seat")?
731 .conductor;
732 let needs_new_seat = !matches!(&self.seat, Some(s) if s.agent == spec.id);
733 if needs_new_seat {
734 self.seat = Some(SeatState::new(SEAT, &spec.id, crate::rng::entropy()));
735 }
736 let seat = self.seat.as_mut().expect("just ensured a seat exists");
737
738 let runnable_views: Vec<prompt::ConductTask> =
739 queued.iter().map(|t| view(t, max_attempts)).collect();
740 let stalled_views: Vec<prompt::ConductTask> =
741 stalled.iter().map(|t| view(t, max_attempts)).collect();
742 let mut finished_views = Vec::with_capacity(finished.len());
743 for t in finished {
744 finished_views.push(finished_view(t, repo, max_attempts).await);
745 }
746
747 let body = prompt::with_overlay(
748 prompt::conduct(
749 &runnable_views,
750 &stalled_views,
751 &finished_views,
752 &cfg.graph.language,
753 ),
754 cfg.prompts.overlay(NODE),
755 );
756
757 let artifacts = home.join("conduct").join("artifacts");
758 let stem = format!("turn-{}", seat.turns + 1);
759 let cache_dir = cfg.cache_dir();
762 let inv = Invocation {
763 cwd: repo,
764 prompt: &body,
765 timeout: TURN_TIMEOUT,
766 allow_write: false,
769 sessions: cfg.graph.sessions,
770 artifacts: &artifacts,
771 stem: &stem,
772 run: NODE,
773 node: NODE,
774 cache_dir: cache_dir.as_deref(),
775 attachments: &[],
776 };
777
778 let busy = Busy::mark(home);
779 let out = agent::invoke(&spec, seat, &inv).await;
780 drop(busy);
781 if let Ok(body) = serde_json::to_string(&*seat) {
784 let path = seat_path(home);
785 if let Some(dir) = path.parent() {
786 let _ = std::fs::create_dir_all(dir);
787 }
788 let _ = std::fs::write(path, body);
789 }
790 let out = out.context("invoking the conductor")?;
791 if !out.usable() {
792 bail!(
793 "no usable reply (exit {:?}, timed out {})",
794 out.exit_code,
795 out.timed_out
796 );
797 }
798 let verdict: Verdict = verdict::extract_json(&out.text)
799 .context("the conductor's reply could not be parsed")?;
800 apply(queue, questions, &verdict)
801 }
802}
803
804#[cfg(test)]
805mod tests {
806 use std::collections::BTreeMap;
807
808 use tempfile::tempdir;
809
810 use super::*;
811 use crate::ask::{Answer, QuestionStatus};
812 use crate::config::{AgentKind, AgentSpec, Graph};
813 use crate::queue::Source;
814
815 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
816 let path = dir.join("mock-conduct-agent.sh");
817 std::fs::write(&path, script).expect("write mock");
818 AgentSpec {
819 id: "mock".to_owned(),
820 kind: AgentKind::Command,
821 model: None,
822 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
823 extra_args: Vec::new(),
824 env,
825 prompt_delivery: None,
826 }
827 }
828
829 fn config(spec: AgentSpec) -> Config {
830 Config {
831 agents: vec![spec],
832 graph: Graph {
833 language: "en".to_owned(),
834 ..Graph::default()
835 },
836 ..Config::default()
837 }
838 }
839
840 fn task(title: &str) -> Task {
841 Task::new(
842 title.to_owned(),
843 format!("do {title}"),
844 std::path::PathBuf::from("."),
845 Source::Human,
846 )
847 }
848
849 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
850 const GARBAGE: &str = "#!/bin/sh\ncat >/dev/null\nprintf 'not json at all\\n'\n";
851
852 fn env(reply: &str) -> BTreeMap<String, String> {
853 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
854 }
855
856 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
857
858 fn init_repo_with_branch(dir: &Path, branch: &str) {
862 use crate::proc::Quiet as _;
863 let run = |args: &[&str]| {
864 let out = std::process::Command::new("git")
865 .args(args)
866 .current_dir(dir)
867 .quiet()
868 .output()
869 .expect("spawn git");
870 assert!(
871 out.status.success(),
872 "git {args:?} failed: {}",
873 String::from_utf8_lossy(&out.stderr)
874 );
875 };
876 run(&["init", "-b", "main"]);
877 run(&["config", "user.name", "magi test"]);
878 run(&["config", "user.email", "magi@example.com"]);
879 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
880 run(&["add", "-A"]);
881 run(&["commit", "-m", "init"]);
882 run(&["checkout", "-b", branch]);
883 std::fs::write(dir.join("change.txt"), "x\n").unwrap();
884 run(&["add", "-A"]);
885 run(&["commit", "-m", "candidate work"]);
886 }
887
888 fn review_round_with_finding(
889 round: usize,
890 finding_id: &str,
891 title: &str,
892 addressed: &[&str],
893 rejected: &[(&str, &str)],
894 ) -> crate::run::ReviewRound {
895 crate::run::ReviewRound {
896 round,
897 head: "deadbeef".to_owned(),
898 verified_head: None,
899 verified_at: None,
900 reviews: vec![crate::run::ReviewRecord {
901 attempts: 0,
902 reviewer: 1,
903 agent: "mock".to_owned(),
904 summary: String::new(),
905 findings: vec![crate::verdict::Finding {
906 id: finding_id.to_owned(),
907 severity: crate::verdict::Severity::Major,
908 file: None,
909 line: None,
910 title: title.to_owned(),
911 detail: String::new(),
912 }],
913 vote: None,
914 failed: None,
915 duration_ms: 0,
916 }],
917 e2e: Vec::new(),
918 verify_retried: false,
919 e2e_deferred: false,
920 e2e_defer_reason: None,
921 fix: Some(crate::run::FixRecord {
922 agent: "mock".to_owned(),
923 addressed: addressed.iter().map(|s| (*s).to_owned()).collect(),
924 rejected: rejected
925 .iter()
926 .map(|(id, why)| crate::verdict::Rejection {
927 id: (*id).to_owned(),
928 why: (*why).to_owned(),
929 })
930 .collect(),
931 notes: String::new(),
932 committed: false,
933 failed: None,
934 duration_ms: 0,
935 continuation: None,
936 }),
937 blocking: 1,
938 answered: 1,
939 expected: 1,
940 clean: false,
941 progressed: true,
942 vote_split: false,
943 reconsideration: Vec::new(),
944 verdict: None,
945 }
946 }
947
948 #[test]
949 fn outcome_for_carries_every_rounds_findings_and_the_branch_head() {
950 crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
951 let dir = tempdir().unwrap();
952 let default_repo = dir.path().join("default");
953 let task_repo = dir.path().join("task");
954 std::fs::create_dir_all(&default_repo).unwrap();
955 std::fs::create_dir_all(&task_repo).unwrap();
956 init_repo_with_branch(&default_repo, "other-branch");
957 init_repo_with_branch(&task_repo, "magi/f00d/A");
958
959 let mut config = Config::default();
960 config.graph.review_rounds = 6;
961 let mut state = crate::run::RunState::new(
962 task_repo.clone(),
963 "main".to_owned(),
964 "deadbeef".to_owned(),
965 "task".to_owned(),
966 config,
967 );
968 state.status = crate::run::RunStatus::Blocked;
969 state.candidates.push(crate::run::Candidate {
970 index: 0,
971 label: 'A',
972 agent: "mock".to_owned(),
973 branch: "magi/f00d/A".to_owned(),
974 worktree: task_repo.clone(),
975 summary: String::new(),
976 stat: String::new(),
977 files: 1,
978 commits: 1,
979 empty: false,
980 failed: None,
981 verified_noop: None,
982 duration_ms: 0,
983 folded: false,
984 });
985 state.tally = Some(crate::run::Tally {
986 first_choice: std::collections::BTreeMap::new(),
987 borda: std::collections::BTreeMap::new(),
988 winner: 'A',
989 rankings: 0,
990 unanimous_initial: false,
991 deliberated: false,
992 changed_votes: 0,
993 unanimous_final: false,
994 tie_break: None,
995 judges: 0,
996 present: 0,
997 quorum: 0,
998 met_quorum: true,
999 uncontested: Some("solo".to_owned()),
1000 });
1001 state.reviews = vec![
1002 review_round_with_finding(
1003 1,
1004 "R1-1-2",
1005 "answer content is dropped",
1006 &[],
1007 &[("R1-1-2", "the id leaving blocked_by is enough")],
1008 ),
1009 review_round_with_finding(2, "R2-1-3", "answer content is still dropped", &[], &[]),
1010 ];
1011 state.save().unwrap();
1012
1013 let mut t = task("outcome test");
1014 t.repo = task_repo;
1015 t.runs.push(state.id.clone());
1016
1017 let finished = tokio_test_block_on(finished_view(&t, &default_repo, 2));
1018 let outcome = finished.outcome;
1019
1020 assert!(outcome.unreadable.is_none());
1021 assert_eq!(outcome.run_status.as_deref(), Some("blocked"));
1022 assert_eq!(outcome.rounds_used, 2);
1023 assert_eq!(outcome.rounds_max, 6);
1024 assert_eq!(outcome.rounds.len(), 2);
1025 assert_eq!(outcome.rounds[0].findings[0].id, "R1-1-2");
1026 assert_eq!(outcome.rounds[0].rejected[0].id, "R1-1-2");
1027 assert!(outcome.rounds[1].addressed.is_empty());
1028 assert!(outcome.rounds[1].rejected.is_empty());
1029 assert_eq!(outcome.branch.as_deref(), Some("magi/f00d/A"));
1030 assert!(
1031 outcome.branch_head.is_some(),
1032 "a real branch must resolve a head commit: {outcome:?}"
1033 );
1034 }
1035
1036 fn tokio_test_block_on<F: std::future::Future>(f: F) -> F::Output {
1040 tokio::runtime::Builder::new_current_thread()
1041 .enable_all()
1042 .build()
1043 .unwrap()
1044 .block_on(f)
1045 }
1046
1047 #[test]
1048 fn view_carries_a_tasks_recorded_answers_into_the_conductor_prompt_input() {
1049 let mut t = task("answered");
1050 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
1051 let v = view(&t, 2);
1052 assert_eq!(v.answers.len(), 1);
1053 assert_eq!(v.answers[0].question, "Which backend?");
1054 assert_eq!(v.answers[0].answer, "SQLite");
1055 }
1056
1057 #[test]
1058 fn a_dependency_decision_blocks_the_task_and_leaves_priority_alone() {
1059 let dir = tempdir().unwrap();
1060 let queue = Queue::at(dir.path().join("queue"));
1061 let questions = Questions::at(dir.path().join("questions"));
1062 let mut a = task("a");
1063 a.priority = 9;
1064 queue.put(&mut a).unwrap();
1065
1066 let verdict = Verdict {
1067 decisions: vec![Decision {
1068 id: a.id.clone(),
1069 blocked_by: vec!["20260101-000000-dead".to_owned()],
1070 reason: Some("waits on the other task".to_owned()),
1071 recovery: None,
1072 question: None,
1073 choices: Vec::new(),
1074 }],
1075 };
1076 apply(&queue, &questions, &verdict).unwrap();
1077
1078 let back = queue.get(&a.id).unwrap();
1079 assert_eq!(back.status, TaskStatus::Blocked);
1080 assert_eq!(back.blocked_by, ["20260101-000000-dead"]);
1081 assert_eq!(
1082 back.priority, 9,
1083 "the conductor's reply cannot carry priority"
1084 );
1085 }
1086
1087 #[test]
1088 fn a_question_decision_files_one_and_blocks_on_its_id() {
1089 let dir = tempdir().unwrap();
1090 let queue = Queue::at(dir.path().join("queue"));
1091 let questions = Questions::at(dir.path().join("questions"));
1092 let mut t = task("ambiguous");
1093 queue.put(&mut t).unwrap();
1094
1095 let verdict = Verdict {
1096 decisions: vec![Decision {
1097 id: t.id.clone(),
1098 blocked_by: Vec::new(),
1099 reason: Some("which backend?".to_owned()),
1100 recovery: None,
1101 question: Some("Which storage backend?".to_owned()),
1102 choices: vec!["SQLite".to_owned(), "Redis".to_owned()],
1103 }],
1104 };
1105 apply(&queue, &questions, &verdict).unwrap();
1106
1107 let back = queue.get(&t.id).unwrap();
1108 assert_eq!(back.status, TaskStatus::Blocked);
1109 assert_eq!(back.blocked_by.len(), 1);
1110 let q = questions.get(&back.blocked_by[0]).unwrap();
1111 assert_eq!(q.summary, "Which storage backend?");
1112 assert_eq!(q.node, NODE);
1113 assert!(q.status.open());
1114 }
1115
1116 #[test]
1117 fn a_task_with_an_open_question_already_reuses_it_rather_than_filing_a_second_one() {
1118 let dir = tempdir().unwrap();
1119 let queue = Queue::at(dir.path().join("queue"));
1120 let questions = Questions::at(dir.path().join("questions"));
1121 let mut t = task("asked once");
1122 queue.put(&mut t).unwrap();
1123
1124 let decision = Decision {
1125 id: t.id.clone(),
1126 reason: Some("still deciding".to_owned()),
1127 question: Some("Which backend?".to_owned()),
1128 ..Decision::default()
1129 };
1130 apply(
1131 &queue,
1132 &questions,
1133 &Verdict {
1134 decisions: vec![decision.clone()],
1135 },
1136 )
1137 .unwrap();
1138 assert_eq!(questions.list().len(), 1);
1139 let first_question_id = queue.get(&t.id).unwrap().blocked_by[0].clone();
1140
1141 let mut released = queue.get(&t.id).unwrap();
1146 released.release();
1147 queue.put(&mut released).unwrap();
1148
1149 apply(
1150 &queue,
1151 &questions,
1152 &Verdict {
1153 decisions: vec![decision],
1154 },
1155 )
1156 .unwrap();
1157 assert_eq!(questions.list().len(), 1, "no duplicate question was filed");
1158 let after = queue.get(&t.id).unwrap();
1159 assert_eq!(
1160 after.blocked_by,
1161 [first_question_id],
1162 "the existing open question is reused, not replaced"
1163 );
1164 }
1165
1166 #[test]
1167 fn a_same_id_question_from_another_node_is_not_reused() {
1168 let dir = tempdir().unwrap();
1169 let queue = Queue::at(dir.path().join("queue"));
1170 let questions = Questions::at(dir.path().join("questions"));
1171 let mut t = task("must ask the conductor");
1172 queue.put(&mut t).unwrap();
1173
1174 let mut unrelated = Question::new(
1175 t.id.clone(),
1176 "review".to_owned(),
1177 "reviewer-1".to_owned(),
1178 "An unrelated review question".to_owned(),
1179 String::new(),
1180 Vec::new(),
1181 );
1182 questions.put(&mut unrelated).unwrap();
1183
1184 apply(
1185 &queue,
1186 &questions,
1187 &Verdict {
1188 decisions: vec![Decision {
1189 id: t.id.clone(),
1190 question: Some("Which backend?".to_owned()),
1191 ..Decision::default()
1192 }],
1193 },
1194 )
1195 .unwrap();
1196
1197 let blocked_by = &queue.get(&t.id).unwrap().blocked_by;
1198 assert_eq!(blocked_by.len(), 1);
1199 assert_ne!(blocked_by[0], unrelated.id);
1200 assert!(questions.get(&unrelated.id).unwrap().status.open());
1201 assert_eq!(questions.get(&blocked_by[0]).unwrap().node, NODE);
1202 }
1203
1204 #[test]
1205 fn answering_the_question_lets_the_resolver_clear_the_block_with_the_answer_kept() {
1206 let dir = tempdir().unwrap();
1207 let queue = Queue::at(dir.path().join("queue"));
1208 let questions = Questions::at(dir.path().join("questions"));
1209 let mut t = task("waits on an answer");
1210 queue.put(&mut t).unwrap();
1211
1212 apply(
1213 &queue,
1214 &questions,
1215 &Verdict {
1216 decisions: vec![Decision {
1217 id: t.id.clone(),
1218 blocked_by: Vec::new(),
1219 reason: None,
1220 recovery: None,
1221 question: Some("Which backend?".to_owned()),
1222 choices: Vec::new(),
1223 }],
1224 },
1225 )
1226 .unwrap();
1227 let blocked = queue.get(&t.id).unwrap();
1228 let question_id = blocked.blocked_by[0].clone();
1229
1230 let mut q = questions.get(&question_id).unwrap();
1231 q.answer(Answer::Text("SQLite".to_owned())).unwrap();
1232 questions.put(&mut q).unwrap();
1233 assert_eq!(q.status, QuestionStatus::Answered);
1234
1235 let mut task_after = queue.get(&t.id).unwrap();
1239 task_after.record_answer(q.summary.clone(), "SQLite".to_owned());
1240 task_after.unblock(&question_id);
1241 assert_eq!(task_after.status, TaskStatus::Queued);
1242 assert_eq!(task_after.answers[0].answer, "SQLite");
1243 }
1244
1245 #[test]
1246 fn a_stalled_task_can_be_requeued_or_held() {
1247 let dir = tempdir().unwrap();
1248 let queue = Queue::at(dir.path().join("queue"));
1249 let questions = Questions::at(dir.path().join("questions"));
1250
1251 let mut requeue_me = task("stuck a");
1252 requeue_me.start("run-1".to_owned());
1253 queue.put(&mut requeue_me).unwrap();
1254
1255 let mut hold_me = task("stuck b");
1256 hold_me.start("run-2".to_owned());
1257 queue.put(&mut hold_me).unwrap();
1258
1259 apply(
1260 &queue,
1261 &questions,
1262 &Verdict {
1263 decisions: vec![
1264 Decision {
1265 id: requeue_me.id.clone(),
1266 recovery: Some(Recovery::Requeue),
1267 ..Decision::default()
1268 },
1269 Decision {
1270 id: hold_me.id.clone(),
1271 recovery: Some(Recovery::Hold),
1272 reason: Some("looks broken".to_owned()),
1273 ..Decision::default()
1274 },
1275 ],
1276 },
1277 )
1278 .unwrap();
1279
1280 let requeued = queue.get(&requeue_me.id).unwrap();
1281 assert_eq!(requeued.status, TaskStatus::Queued);
1282 assert_eq!(requeued.attempts, 0);
1283
1284 let held = queue.get(&hold_me.id).unwrap();
1285 assert_eq!(held.status, TaskStatus::Held);
1286 assert_eq!(held.hold_reason.as_deref(), Some("looks broken"));
1287 }
1288
1289 #[test]
1290 fn a_machine_held_task_asked_about_restores_to_held_once_answered() {
1291 let dir = tempdir().unwrap();
1296 let queue = Queue::at(dir.path().join("queue"));
1297 let questions = Questions::at(dir.path().join("questions"));
1298 let mut t = task("held out of attempts");
1299 t.hold_machine(Some("out of attempts".to_owned()));
1300 queue.put(&mut t).unwrap();
1301
1302 apply(
1303 &queue,
1304 &questions,
1305 &Verdict {
1306 decisions: vec![Decision {
1307 id: t.id.clone(),
1308 reason: Some("what should happen to this one?".to_owned()),
1309 question: Some("Hold it, or try again?".to_owned()),
1310 ..Decision::default()
1311 }],
1312 },
1313 )
1314 .unwrap();
1315 let blocked = queue.get(&t.id).unwrap();
1316 assert_eq!(blocked.status, TaskStatus::Blocked);
1317 let question_id = blocked.blocked_by[0].clone();
1318
1319 let mut q = questions.get(&question_id).unwrap();
1320 q.answer(Answer::Text("leave it held".to_owned())).unwrap();
1321 questions.put(&mut q).unwrap();
1322
1323 let mut after = queue.get(&t.id).unwrap();
1325 after.record_answer(q.summary.clone(), "leave it held".to_owned());
1326 after.unblock(&question_id);
1327 assert_eq!(
1328 after.status,
1329 TaskStatus::Held,
1330 "must not fall back to queued"
1331 );
1332 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
1333 }
1334
1335 #[test]
1336 fn a_reaffirmed_hold_with_no_new_reason_is_not_silently_auto_released_by_triage() {
1337 let dir = tempdir().unwrap();
1347 let queue = Queue::at(dir.path().join("queue"));
1348 let questions = Questions::at(dir.path().join("questions"));
1349 let mut t = task("disk pressure, then reconsidered");
1350 t.hold_machine(Some(
1351 "not enough free space to start a run: 10 bytes free, 100 required by \
1352 `[disk] min_free_bytes`"
1353 .to_owned(),
1354 ));
1355 t.record_answer(
1356 "How should this be handled?".to_owned(),
1357 "keep it held, a human will look at it later".to_owned(),
1358 );
1359 queue.put(&mut t).unwrap();
1360
1361 apply(
1362 &queue,
1363 &questions,
1364 &Verdict {
1365 decisions: vec![Decision {
1366 id: t.id.clone(),
1367 recovery: Some(Recovery::Hold),
1368 ..Decision::default()
1369 }],
1370 },
1371 )
1372 .unwrap();
1373
1374 let after = queue.get(&t.id).unwrap();
1375 assert_eq!(after.status, TaskStatus::Held);
1376 assert!(
1377 !after
1378 .hold_reason
1379 .as_deref()
1380 .unwrap_or_default()
1381 .starts_with("not enough free space"),
1382 "the stale disk-pressure text must not survive a reconfirmed hold: {:?}",
1383 after.hold_reason
1384 );
1385
1386 let cfg_dir = tempdir().unwrap();
1389 let config = cfg_dir.path().join("magi.toml");
1390 std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1391 let report =
1392 crate::triage::run_once(&queue, &questions, Some(&config), jiff::Timestamp::now());
1393 assert!(report.resumed.is_empty(), "must not be auto-released");
1394 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1395 }
1396
1397 #[test]
1398 fn a_conductor_rehold_after_a_resume_answer_is_not_asked_again_identically() {
1399 let dir = tempdir().unwrap();
1400 let queue = Queue::at(dir.path().join("queue"));
1401 let questions = Questions::at(dir.path().join("questions"));
1402 let config = dir.path().join("magi.toml");
1403 std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1404 let now = jiff::Timestamp::now();
1405 let triage = || crate::triage::run_once(&queue, &questions, Some(&config), now);
1406 let open = || {
1407 questions
1408 .list()
1409 .into_iter()
1410 .filter(|q| q.node == "triage" && q.status.open())
1411 .collect::<Vec<_>>()
1412 };
1413 let hold = Decision {
1414 recovery: Some(Recovery::Hold),
1415 reason: Some("waiting on manual worktree cleanup".to_owned()),
1416 ..Decision::default()
1417 };
1418
1419 let mut t = task("looping hold");
1420 t.hold_machine(Some("waiting on manual worktree cleanup".to_owned()));
1421 queue.put(&mut t).unwrap();
1422 let hold = Decision {
1423 id: t.id.clone(),
1424 ..hold
1425 };
1426
1427 assert_eq!(triage().asked.len(), 1);
1429 let first = open().remove(0);
1430 let mut q = questions.get(&first.id).unwrap();
1431 let resume = q.choices[0].clone();
1432 q.answer(Answer::Choice(resume)).unwrap();
1433 questions.put(&mut q).unwrap();
1434 assert_eq!(triage().answered.len(), 1);
1435 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1436
1437 apply(
1439 &queue,
1440 &questions,
1441 &Verdict {
1442 decisions: vec![hold.clone()],
1443 },
1444 )
1445 .unwrap();
1446 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1447
1448 assert_eq!(triage().asked.len(), 1);
1450 let second = open().remove(0);
1451 assert_ne!(second.summary, first.summary);
1452 assert_ne!(second.choices, first.choices);
1453 assert!(second.detail.contains("waiting on manual worktree cleanup"));
1454 assert!(triage().asked.is_empty(), "no duplicate question");
1455 assert_eq!(open().len(), 1);
1456
1457 let mut q = questions.get(&second.id).unwrap();
1459 let force = q.choices[0].clone();
1460 q.answer(Answer::Choice(force)).unwrap();
1461 questions.put(&mut q).unwrap();
1462 assert_eq!(triage().answered.len(), 1);
1463 apply(
1464 &queue,
1465 &questions,
1466 &Verdict {
1467 decisions: vec![hold],
1468 },
1469 )
1470 .unwrap();
1471 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1472 }
1473
1474 #[test]
1475 fn a_runnable_task_can_be_held_directly_without_a_question() {
1476 let dir = tempdir().unwrap();
1477 let queue = Queue::at(dir.path().join("queue"));
1478 let questions = Questions::at(dir.path().join("questions"));
1479 let mut t = task("already answered, should stay put");
1480 queue.put(&mut t).unwrap();
1481
1482 apply(
1483 &queue,
1484 &questions,
1485 &Verdict {
1486 decisions: vec![Decision {
1487 id: t.id.clone(),
1488 recovery: Some(Recovery::Hold),
1489 reason: Some("operator already said keep this held".to_owned()),
1490 ..Decision::default()
1491 }],
1492 },
1493 )
1494 .unwrap();
1495
1496 let after = queue.get(&t.id).unwrap();
1497 assert_eq!(after.status, TaskStatus::Held);
1498 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
1499 }
1500
1501 #[test]
1502 fn done_recovery_closes_a_held_task_whose_goal_is_already_met() {
1503 let dir = tempdir().unwrap();
1509 let queue = Queue::at(dir.path().join("queue"));
1510 let questions = Questions::at(dir.path().join("questions"));
1511 let mut t = task("already merged by hand");
1512 t.hold_machine(Some("branch survived, awaiting a decision".to_owned()));
1513 t.record_answer(
1514 "Handle this one?".to_owned(),
1515 "already merged and cleaned up, close it".to_owned(),
1516 );
1517 queue.put(&mut t).unwrap();
1518
1519 apply(
1520 &queue,
1521 &questions,
1522 &Verdict {
1523 decisions: vec![Decision {
1524 id: t.id.clone(),
1525 recovery: Some(Recovery::Done),
1526 reason: Some("operator confirmed this already landed".to_owned()),
1527 ..Decision::default()
1528 }],
1529 },
1530 )
1531 .unwrap();
1532
1533 let after = queue.get(&t.id).unwrap();
1534 assert_eq!(after.status, TaskStatus::Done);
1535 assert!(after.hold_reason.is_none());
1536 assert_eq!(after.answers.len(), 1, "the record of why is kept");
1537 }
1538
1539 #[test]
1540 fn done_recovery_is_ignored_for_a_runnable_or_running_task() {
1541 let dir = tempdir().unwrap();
1542 let queue = Queue::at(dir.path().join("queue"));
1543 let questions = Questions::at(dir.path().join("questions"));
1544
1545 let mut queued = task("never ran yet");
1546 queue.put(&mut queued).unwrap();
1547
1548 let mut running = task("mid-run");
1549 running.start("run-1".to_owned());
1550 queue.put(&mut running).unwrap();
1551
1552 for id in [queued.id.clone(), running.id.clone()] {
1553 apply(
1554 &queue,
1555 &questions,
1556 &Verdict {
1557 decisions: vec![Decision {
1558 id,
1559 recovery: Some(Recovery::Done),
1560 ..Decision::default()
1561 }],
1562 },
1563 )
1564 .unwrap();
1565 }
1566
1567 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
1568 assert_eq!(queue.get(&running.id).unwrap().status, TaskStatus::Running);
1569 }
1570
1571 #[test]
1572 fn a_third_conductor_question_after_two_settled_answers_holds_instead_of_asking_again() {
1573 let dir = tempdir().unwrap();
1579 let queue = Queue::at(dir.path().join("queue"));
1580 let questions = Questions::at(dir.path().join("questions"));
1581 let mut t = task("asked about repeatedly");
1582 t.hold_machine(Some("out of attempts".to_owned()));
1583 t.record_answer("Handle this one? (1)".to_owned(), "not yet".to_owned());
1584 t.record_answer(
1585 "Handle this one? (2)".to_owned(),
1586 "still not yet".to_owned(),
1587 );
1588 queue.put(&mut t).unwrap();
1589 assert_eq!(questions.list().len(), 0);
1590
1591 apply(
1592 &queue,
1593 &questions,
1594 &Verdict {
1595 decisions: vec![Decision {
1596 id: t.id.clone(),
1597 question: Some("Handle this one? (3)".to_owned()),
1598 ..Decision::default()
1599 }],
1600 },
1601 )
1602 .unwrap();
1603
1604 assert_eq!(questions.list().len(), 0, "no third question was filed");
1605 let after = queue.get(&t.id).unwrap();
1606 assert_eq!(after.status, TaskStatus::Held);
1607 assert!(after.blocked_by.is_empty());
1608 assert_eq!(after.answers.len(), 2, "the prior answers are untouched");
1609 }
1610
1611 #[test]
1612 fn a_second_conductor_question_is_still_allowed_after_one_settled_answer() {
1613 let dir = tempdir().unwrap();
1614 let queue = Queue::at(dir.path().join("queue"));
1615 let questions = Questions::at(dir.path().join("questions"));
1616 let mut t = task("asked about once already");
1617 t.hold_machine(Some("out of attempts".to_owned()));
1618 t.record_answer("Handle this one?".to_owned(), "not yet".to_owned());
1619 queue.put(&mut t).unwrap();
1620
1621 apply(
1622 &queue,
1623 &questions,
1624 &Verdict {
1625 decisions: vec![Decision {
1626 id: t.id.clone(),
1627 question: Some("Still not sure - now what?".to_owned()),
1628 ..Decision::default()
1629 }],
1630 },
1631 )
1632 .unwrap();
1633
1634 assert_eq!(questions.list().len(), 1, "the second question was filed");
1635 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
1636 }
1637
1638 #[test]
1639 fn a_held_task_with_an_open_triage_question_is_left_to_triage() {
1640 let dir = tempdir().unwrap();
1646 let queue = Queue::at(dir.path().join("queue"));
1647 let questions = Questions::at(dir.path().join("questions"));
1648 let mut t = task("held, triage already asking about it");
1649 t.hold_machine(Some("cause unclear".to_owned()));
1650 queue.put(&mut t).unwrap();
1651
1652 let mut triage_q = Question::new(
1653 t.id.clone(),
1654 crate::triage::NODE.to_owned(),
1655 "triage".to_owned(),
1656 "Still needed?".to_owned(),
1657 String::new(),
1658 vec![
1659 "resume".to_owned(),
1660 "not yet".to_owned(),
1661 "discard".to_owned(),
1662 ],
1663 );
1664 questions.put(&mut triage_q).unwrap();
1665
1666 for decision in [
1667 Decision {
1668 id: t.id.clone(),
1669 question: Some("what now?".to_owned()),
1670 ..Decision::default()
1671 },
1672 Decision {
1673 id: t.id.clone(),
1674 recovery: Some(Recovery::Requeue),
1675 ..Decision::default()
1676 },
1677 ] {
1678 apply(
1679 &queue,
1680 &questions,
1681 &Verdict {
1682 decisions: vec![decision],
1683 },
1684 )
1685 .unwrap();
1686 }
1687
1688 let after = queue.get(&t.id).unwrap();
1689 assert_eq!(
1690 after.status,
1691 TaskStatus::Held,
1692 "triage still owns this hold"
1693 );
1694 assert!(after.blocked_by.is_empty());
1695 assert_eq!(
1696 questions.list().len(),
1697 1,
1698 "no second, conductor-owned question was filed"
1699 );
1700 }
1701
1702 #[test]
1703 fn a_held_task_with_an_answered_but_unapplied_triage_question_is_still_left_alone() {
1704 let dir = tempdir().unwrap();
1712 let queue = Queue::at(dir.path().join("queue"));
1713 let questions = Questions::at(dir.path().join("questions"));
1714 let mut t = task("held, triage question answered but not yet applied");
1715 t.hold_machine(Some("cause unclear".to_owned()));
1716 queue.put(&mut t).unwrap();
1717
1718 let mut triage_q = Question::new(
1719 t.id.clone(),
1720 crate::triage::NODE.to_owned(),
1721 "triage".to_owned(),
1722 "Still needed?".to_owned(),
1723 String::new(),
1724 vec![
1725 "resume".to_owned(),
1726 "not yet".to_owned(),
1727 "discard".to_owned(),
1728 ],
1729 );
1730 questions.put(&mut triage_q).unwrap();
1731 triage_q
1732 .answer(Answer::Choice("not yet".to_owned()))
1733 .unwrap();
1734 questions.put(&mut triage_q).unwrap();
1735 assert!(!triage_q.status.open());
1736
1737 apply(
1738 &queue,
1739 &questions,
1740 &Verdict {
1741 decisions: vec![Decision {
1742 id: t.id.clone(),
1743 question: Some("what now?".to_owned()),
1744 ..Decision::default()
1745 }],
1746 },
1747 )
1748 .unwrap();
1749
1750 let after = queue.get(&t.id).unwrap();
1751 assert_eq!(
1752 after.status,
1753 TaskStatus::Held,
1754 "triage's own answer is not yet applied - conduct must wait"
1755 );
1756 assert_eq!(
1757 questions.list().len(),
1758 1,
1759 "no conductor question was filed over the pending triage answer"
1760 );
1761 }
1762
1763 #[test]
1764 fn manual_hold_rejects_hostile_or_stale_conductor_recovery() {
1765 let dir = tempdir().unwrap();
1766 let queue = Queue::at(dir.path().join("queue"));
1767 let questions = Questions::at(dir.path().join("questions"));
1768 let mut held = task("manual recovery");
1769 held.priority = 300;
1770 held.runs.push("run20260912-224242-daf5".to_owned());
1771 held.hold_manual(Some(
1772 "active manual recovery run20260912-224242-daf5".to_owned(),
1773 ));
1774 queue.put(&mut held).unwrap();
1775
1776 for decision in [
1780 Decision {
1781 id: held.id.clone(),
1782 recovery: Some(Recovery::Requeue),
1783 ..Decision::default()
1784 },
1785 Decision {
1786 id: held.id.clone(),
1787 recovery: Some(Recovery::Hold),
1788 reason: Some("stale replacement reason".to_owned()),
1789 ..Decision::default()
1790 },
1791 Decision {
1792 id: held.id.clone(),
1793 recovery: Some(Recovery::Review),
1794 ..Decision::default()
1795 },
1796 Decision {
1797 id: held.id.clone(),
1798 blocked_by: vec!["other-task".to_owned()],
1799 question: Some("retry now?".to_owned()),
1800 ..Decision::default()
1801 },
1802 ] {
1803 apply(
1804 &queue,
1805 &questions,
1806 &Verdict {
1807 decisions: vec![decision],
1808 },
1809 )
1810 .unwrap();
1811 }
1812
1813 let after = queue.get(&held.id).unwrap();
1814 assert_eq!(after.status, TaskStatus::Held);
1815 assert!(after.operator_held());
1816 assert_eq!(after.priority, 300);
1817 assert_eq!(after.runs, ["run20260912-224242-daf5"]);
1818 assert_eq!(
1819 after.hold_reason.as_deref(),
1820 Some("active manual recovery run20260912-224242-daf5")
1821 );
1822 assert!(after.blocked_by.is_empty());
1823 assert!(questions.list().is_empty());
1824 assert!(
1825 queue.next_runnable().is_none(),
1826 "must not dispatch a duplicate"
1827 );
1828 }
1829
1830 #[test]
1831 fn machine_holds_remain_recoverable_and_manual_release_is_authorization() {
1832 let dir = tempdir().unwrap();
1833 let queue = Queue::at(dir.path().join("queue"));
1834 let questions = Questions::at(dir.path().join("questions"));
1835
1836 let mut automatic = task("disk gate");
1837 automatic.hold_machine(Some("disk full".to_owned()));
1838 queue.put(&mut automatic).unwrap();
1839 let requeue = || Verdict {
1840 decisions: vec![Decision {
1841 id: automatic.id.clone(),
1842 recovery: Some(Recovery::Requeue),
1843 ..Decision::default()
1844 }],
1845 };
1846 apply(&queue, &questions, &requeue()).unwrap();
1847 assert_eq!(queue.get(&automatic.id).unwrap().status, TaskStatus::Queued);
1848
1849 let mut manual = task("operator gate");
1850 manual.hold_manual(Some("wait for operator".to_owned()));
1851 queue.put(&mut manual).unwrap();
1852 apply(
1853 &queue,
1854 &questions,
1855 &Verdict {
1856 decisions: vec![Decision {
1857 id: manual.id.clone(),
1858 recovery: Some(Recovery::Requeue),
1859 ..Decision::default()
1860 }],
1861 },
1862 )
1863 .unwrap();
1864 assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Held);
1865
1866 let mut released = queue.get(&manual.id).unwrap();
1869 released.release();
1870 queue.put(&mut released).unwrap();
1871 assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Queued);
1872 }
1873
1874 #[test]
1875 fn legacy_reasoned_hold_is_protected_without_losing_its_metadata() {
1876 let dir = tempdir().unwrap();
1877 let queue = Queue::at(dir.path().join("queue"));
1878 let questions = Questions::at(dir.path().join("questions"));
1879 let mut legacy = task("old explicit hold");
1880 legacy.status = TaskStatus::Held;
1881 legacy.hold_reason = Some("manual recovery already active".to_owned());
1882 legacy.hold_source = None;
1883 legacy.blocked_by = vec!["dependency".to_owned()];
1884 queue.put(&mut legacy).unwrap();
1885
1886 apply(
1887 &queue,
1888 &questions,
1889 &Verdict {
1890 decisions: vec![Decision {
1891 id: legacy.id.clone(),
1892 recovery: Some(Recovery::Requeue),
1893 ..Decision::default()
1894 }],
1895 },
1896 )
1897 .unwrap();
1898
1899 let after = queue.get(&legacy.id).unwrap();
1900 assert_eq!(after.status, TaskStatus::Held);
1901 assert_eq!(after.hold_source, None);
1902 assert_eq!(after.hold_reason, legacy.hold_reason);
1903 assert_eq!(after.blocked_by, legacy.blocked_by);
1904 }
1905
1906 #[test]
1907 fn review_recovery_is_a_no_op_without_a_survivable_branch() {
1908 crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
1914 let dir = tempdir().unwrap();
1915 let queue = Queue::at(dir.path().join("queue"));
1916 let questions = Questions::at(dir.path().join("questions"));
1917 let mut t = task("blocked with no readable run");
1918 t.start("20260101-000000-dead".to_owned()); t.fail("blocked", 5);
1920 queue.put(&mut t).unwrap();
1921
1922 apply(
1923 &queue,
1924 &questions,
1925 &Verdict {
1926 decisions: vec![Decision {
1927 id: t.id.clone(),
1928 recovery: Some(Recovery::Review),
1929 ..Decision::default()
1930 }],
1931 },
1932 )
1933 .unwrap();
1934
1935 let after = queue.get(&t.id).unwrap();
1936 assert_eq!(
1937 after.status,
1938 TaskStatus::Failed,
1939 "with nothing to reopen, the decision is dropped rather than guessed at"
1940 );
1941 assert!(after.review_branch.is_none());
1942 }
1943
1944 #[test]
1945 fn requeue_and_review_recovery_are_ignored_for_a_runnable_task() {
1946 let dir = tempdir().unwrap();
1950 let queue = Queue::at(dir.path().join("queue"));
1951 let questions = Questions::at(dir.path().join("questions"));
1952
1953 for recovery in [Recovery::Requeue, Recovery::Review] {
1954 let mut t = task("ordinary");
1955 queue.put(&mut t).unwrap();
1956
1957 apply(
1958 &queue,
1959 &questions,
1960 &Verdict {
1961 decisions: vec![Decision {
1962 id: t.id.clone(),
1963 recovery: Some(recovery),
1964 ..Decision::default()
1965 }],
1966 },
1967 )
1968 .unwrap();
1969
1970 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1971 }
1972 }
1973
1974 #[tokio::test]
1975 async fn a_broken_agent_leaves_the_queue_untouched_and_does_not_error() {
1976 let dir = tempdir().unwrap();
1977 let cfg = config(mock_agent(dir.path(), BROKEN, BTreeMap::new()));
1978 let queue = Queue::at(dir.path().join("queue"));
1979 let questions = Questions::at(dir.path().join("questions"));
1980 let mut t = task("normal");
1981 queue.put(&mut t).unwrap();
1982
1983 let mut conductor = Conductor::new();
1984 conductor
1985 .maybe_run(
1986 &cfg,
1987 dir.path(),
1988 &queue,
1989 &questions,
1990 dir.path(),
1991 &[t.clone()],
1992 &[],
1993 &[],
1994 2,
1995 )
1996 .await;
1997
1998 assert_eq!(
1999 queue.get(&t.id).unwrap().status,
2000 TaskStatus::Queued,
2001 "a failed invocation must change nothing"
2002 );
2003 assert!(
2004 queue.next_runnable().is_some(),
2005 "the loop must still be able to take the next task"
2006 );
2007 }
2008
2009 #[tokio::test]
2010 async fn a_reply_with_no_json_leaves_the_queue_untouched() {
2011 let dir = tempdir().unwrap();
2012 let cfg = config(mock_agent(dir.path(), GARBAGE, BTreeMap::new()));
2013 let queue = Queue::at(dir.path().join("queue"));
2014 let questions = Questions::at(dir.path().join("questions"));
2015 let mut t = task("normal");
2016 queue.put(&mut t).unwrap();
2017
2018 let mut conductor = Conductor::new();
2019 conductor
2020 .maybe_run(
2021 &cfg,
2022 dir.path(),
2023 &queue,
2024 &questions,
2025 dir.path(),
2026 &[t.clone()],
2027 &[],
2028 &[],
2029 2,
2030 )
2031 .await;
2032
2033 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2034 }
2035
2036 #[tokio::test]
2037 async fn json_survives_code_fences_and_a_preamble() {
2038 let dir = tempdir().unwrap();
2039 let mut t = task("fenced");
2040 let reply = format!(
2041 "Sure, here is my decision.\n\n```json\n{{\"decisions\":[{{\"id\":\"{}\",\
2042 \"blocked_by\":[\"x\"],\"reason\":\"why\"}}]}}\n```\n",
2043 t.id
2044 );
2045 let cfg = config(mock_agent(dir.path(), REPLY, env(&reply)));
2046 let queue = Queue::at(dir.path().join("queue"));
2047 let questions = Questions::at(dir.path().join("questions"));
2048 queue.put(&mut t).unwrap();
2049
2050 let mut conductor = Conductor::new();
2051 conductor
2052 .maybe_run(
2053 &cfg,
2054 dir.path(),
2055 &queue,
2056 &questions,
2057 dir.path(),
2058 &[t.clone()],
2059 &[],
2060 &[],
2061 2,
2062 )
2063 .await;
2064
2065 let back = queue.get(&t.id).unwrap();
2066 assert_eq!(back.status, TaskStatus::Blocked);
2067 assert_eq!(back.blocked_by, ["x"]);
2068 }
2069
2070 #[tokio::test]
2071 async fn the_conductor_is_not_called_again_when_nothing_worth_looking_at_has_changed() {
2072 let dir = tempdir().unwrap();
2075 let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2076 let queue = Queue::at(dir.path().join("queue"));
2077 let questions = Questions::at(dir.path().join("questions"));
2078 let mut t = task("stable");
2079 queue.put(&mut t).unwrap();
2080 let artifacts = dir.path().join("conduct").join("artifacts");
2081 let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2082
2083 let mut conductor = Conductor::new();
2084 conductor
2085 .maybe_run(
2086 &cfg,
2087 dir.path(),
2088 &queue,
2089 &questions,
2090 dir.path(),
2091 &[t.clone()],
2092 &[],
2093 &[],
2094 2,
2095 )
2096 .await;
2097 assert!(turn(1).is_file(), "the first cycle must call the conductor");
2098
2099 conductor
2100 .maybe_run(
2101 &cfg,
2102 dir.path(),
2103 &queue,
2104 &questions,
2105 dir.path(),
2106 &[t.clone()],
2107 &[],
2108 &[],
2109 2,
2110 )
2111 .await;
2112 assert!(
2113 !turn(2).is_file(),
2114 "an unchanged revision and an unchanged stalled/finished set must not call the \
2115 conductor twice"
2116 );
2117
2118 t.priority = 1;
2120 queue.put(&mut t).unwrap();
2121 conductor
2122 .maybe_run(
2123 &cfg,
2124 dir.path(),
2125 &queue,
2126 &questions,
2127 dir.path(),
2128 &[t.clone()],
2129 &[],
2130 &[],
2131 2,
2132 )
2133 .await;
2134 assert!(turn(2).is_file(), "a moved revision calls it again");
2135 }
2136
2137 #[tokio::test]
2138 async fn a_task_turning_stalled_calls_the_conductor_again_despite_an_unchanged_revision() {
2139 let dir = tempdir().unwrap();
2145 let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2146 let queue = Queue::at(dir.path().join("queue"));
2147 let questions = Questions::at(dir.path().join("questions"));
2148 let mut t = task("quiet");
2149 queue.put(&mut t).unwrap();
2150 let artifacts = dir.path().join("conduct").join("artifacts");
2151 let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2152
2153 let mut conductor = Conductor::new();
2154 conductor
2155 .maybe_run(
2156 &cfg,
2157 dir.path(),
2158 &queue,
2159 &questions,
2160 dir.path(),
2161 &[t.clone()],
2162 &[],
2163 &[],
2164 2,
2165 )
2166 .await;
2167 assert!(turn(1).is_file());
2168
2169 conductor
2170 .maybe_run(
2171 &cfg,
2172 dir.path(),
2173 &queue,
2174 &questions,
2175 dir.path(),
2176 &[],
2177 &[t.clone()],
2178 &[],
2179 2,
2180 )
2181 .await;
2182 assert!(
2183 turn(2).is_file(),
2184 "a task turning stalled must call the conductor again"
2185 );
2186
2187 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(3).is_file(),
2204 "the same stalled task lingering must not call the conductor every cycle"
2205 );
2206 }
2207
2208 #[test]
2209 fn worth_a_look_is_config_free_and_matches_maybe_runs_own_gate() {
2210 let dir = tempdir().unwrap();
2211 let queue = Queue::at(dir.path().join("queue"));
2212 let mut t = task("t");
2213 queue.put(&mut t).unwrap();
2214
2215 let mut conductor = Conductor::new();
2216 assert!(
2217 conductor.worth_a_look(&queue, &[], &[]),
2218 "a conductor that has never run has something to look at"
2219 );
2220
2221 conductor.last_seen = Some(Conductor::snapshot(&queue, &[], &[]));
2222 assert!(
2223 !conductor.worth_a_look(&queue, &[], &[]),
2224 "nothing changed and nothing is stalled or finished"
2225 );
2226 assert!(
2227 conductor.worth_a_look(&queue, &[t.clone()], &[]),
2228 "a stalled task is worth a look even at the same revision"
2229 );
2230 assert!(
2231 conductor.worth_a_look(&queue, &[], &[t.clone()]),
2232 "a finished task is worth a look even at the same revision"
2233 );
2234 }
2235
2236 #[tokio::test]
2237 async fn the_conduct_path_never_calls_ask_and_wait() {
2238 let dir = tempdir().unwrap();
2244 let queue = Queue::at(dir.path().join("queue"));
2245 let questions = Questions::at(dir.path().join("questions"));
2246 let mut t = task("asks without blocking");
2247 queue.put(&mut t).unwrap();
2248
2249 apply(
2250 &queue,
2251 &questions,
2252 &Verdict {
2253 decisions: vec![Decision {
2254 id: t.id.clone(),
2255 question: Some("ok?".to_owned()),
2256 ..Decision::default()
2257 }],
2258 },
2259 )
2260 .unwrap();
2261 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
2263 }
2264}