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