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