1use std::collections::{BTreeMap, BTreeSet};
44use std::path::{Path, PathBuf};
45use std::time::Duration;
46
47use anyhow::{Context as _, Result};
48use serde::Deserialize;
49
50use crate::agent::{self, Invocation, SeatState};
51use crate::ask::{Deputy, Question, Questions};
52use crate::config::Config;
53use crate::prompt;
54use crate::queue::{Queue, Task, TaskStatus};
55use crate::run::RunState;
56use crate::verdict;
57
58const SEAT: &str = "conduct";
61
62pub const NODE: &str = "conduct";
66
67const TURN_TIMEOUT: Duration = Duration::from_secs(300);
72
73#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
76#[serde(rename_all = "lowercase")]
77pub enum Recovery {
78 Requeue,
81 Hold,
84 Review,
89 Done,
97}
98
99#[derive(Debug, Clone, Default, Deserialize)]
102pub struct Decision {
103 pub id: String,
105 #[serde(default)]
108 pub blocked_by: Vec<String>,
109 #[serde(default)]
111 pub reason: Option<String>,
112 #[serde(default)]
120 pub recovery: Option<Recovery>,
121 #[serde(default)]
125 pub question: Option<String>,
126 #[serde(default)]
128 pub choices: Vec<String>,
129}
130
131#[derive(Debug, Clone, Default, Deserialize)]
144pub struct Verdict {
145 pub decisions: Vec<Decision>,
148}
149
150fn view(t: &Task, max_attempts: usize) -> prompt::ConductTask {
153 prompt::ConductTask {
154 id: t.id.clone(),
155 title: t.title.clone(),
156 instruction: t.instruction.clone(),
157 repo: t.repo.display().to_string(),
158 priority: t.priority,
159 status: t.status.as_str().to_owned(),
160 attempts: t.attempts,
161 max_attempts,
162 last_error: t.last_error.clone(),
163 hold_reason: t.hold_reason.clone(),
164 hold_source: t.hold_source.map(|source| source.label().to_owned()),
165 blocked_by: t.blocked_by.clone(),
166 answers: t
167 .answers
168 .iter()
169 .map(|a| prompt::ConductAnswer {
170 question: a.question.clone(),
171 answer: a.answer.clone(),
172 })
173 .collect(),
174 operator_resume: t.resume_override.as_ref().map(|o| {
175 format!(
176 "the operator explicitly answered \"resume\" at {}; do not hold this \
177 task again for the same reason unless there is new information",
178 o.at
179 )
180 }),
181 }
182}
183
184fn pinned_resume(task: &Task) -> bool {
188 task.resume_override
189 .as_ref()
190 .is_some_and(|o| o.pinned_run.is_some())
191}
192
193fn may_hold(task: &mut Task, reason: &str) -> bool {
202 let Some(o) = task.resume_override.as_mut() else {
203 return true;
204 };
205 if o.forced || o.pinned_run.is_some() {
206 tracing::warn!(
207 "conductor tried to hold task {} after the operator forced a resume: {reason}",
208 task.id
209 );
210 return false;
211 }
212 if o.conductor_rehold.is_some() {
213 return false;
214 }
215 o.conductor_rehold = Some(reason.to_owned());
216 true
217}
218
219fn severity_str(s: crate::verdict::Severity) -> &'static str {
222 match s {
223 crate::verdict::Severity::Nit => "nit",
224 crate::verdict::Severity::Minor => "minor",
225 crate::verdict::Severity::Major => "major",
226 crate::verdict::Severity::Blocker => "blocker",
227 }
228}
229
230fn surviving_branch(task: &Task) -> Option<String> {
235 let last = task.runs.last()?;
236 let state = RunState::load(last).ok()?;
237 state.winner().map(|c| c.branch.clone())
238}
239
240fn reaffirmed_hold_reason(task: &Task, d: &Decision) -> Option<String> {
267 let note = match &d.reason {
268 Some(reason) => reason.clone(),
269 None => match task.answers.last() {
270 Some(a) => format!(
271 "conduct held this again with no new reason given; last operator \
272 answer on record: {}",
273 a.answer
274 ),
275 None => "conduct held this again with no reason given".to_owned(),
276 },
277 };
278 match task.hold_reason.as_deref() {
279 Some(prior) if !prior.is_empty() => {
280 let outer = outermost_hold_note(prior);
281 if note.trim() == prior.trim() || note.trim() == outer {
282 None
283 } else {
284 Some(format!("{note}\n\n(previously: {outer})"))
285 }
286 }
287 _ => Some(note),
288 }
289}
290
291const PREVIOUSLY: &str = "\n\n(previously: ";
294
295fn outermost_hold_note(reason: &str) -> &str {
298 reason.split(PREVIOUSLY).next().unwrap_or(reason).trim()
299}
300
301fn hold_note(d: &Decision) -> String {
303 d.reason
304 .clone()
305 .unwrap_or_else(|| "(no reason given)".to_owned())
306}
307
308async fn outcome_for(task: &Task, repo: &Path) -> prompt::ConductOutcome {
310 let Some(run_id) = task.runs.last().cloned() else {
311 return prompt::ConductOutcome {
312 run_id: "(none)".to_owned(),
313 unreadable: Some("this task has not produced a run yet".to_owned()),
314 run_status: None,
315 open_findings: Vec::new(),
316 rounds_used: 0,
317 rounds_max: 0,
318 rounds: Vec::new(),
319 branch: None,
320 branch_head: None,
321 references: None,
322 empty_candidate: false,
323 };
324 };
325 let state = match RunState::load(&run_id) {
326 Ok(s) => s,
327 Err(e) => {
328 tracing::warn!(
333 "conductor: could not read run {run_id} for task {}: {e:#}",
334 task.short()
335 );
336 return prompt::ConductOutcome {
337 run_id,
338 unreadable: Some(format!("{e:#}")),
339 run_status: None,
340 open_findings: Vec::new(),
341 rounds_used: 0,
342 rounds_max: 0,
343 rounds: Vec::new(),
344 branch: None,
345 branch_head: None,
346 references: None,
347 empty_candidate: false,
348 };
349 }
350 };
351
352 let finding_view = |f: &crate::verdict::Finding| prompt::ConductFinding {
353 id: f.id.clone(),
354 title: f.title.clone(),
355 severity: severity_str(f.severity).to_owned(),
356 };
357 let open_findings = state
358 .open_findings()
359 .into_iter()
360 .map(finding_view)
361 .collect();
362 let rounds = state
363 .reviews
364 .iter()
365 .map(|r| prompt::ConductRound {
366 round: r.round,
367 findings: r
368 .reviews
369 .iter()
370 .flat_map(|rec| rec.findings.iter())
371 .map(finding_view)
372 .collect(),
373 addressed: r
374 .fix
375 .as_ref()
376 .map(|fx| fx.addressed.clone())
377 .unwrap_or_default(),
378 rejected: r
379 .fix
380 .as_ref()
381 .map(|fx| {
382 fx.rejected
383 .iter()
384 .map(|rej| prompt::ConductRejection {
385 id: rej.id.clone(),
386 why: rej.why.clone(),
387 })
388 .collect()
389 })
390 .unwrap_or_default(),
391 })
392 .collect();
393 let branch = state.winner().map(|c| c.branch.clone());
394 let branch_head = match &branch {
395 Some(b) => crate::git::rev_parse(repo, b)
396 .await
397 .ok()
398 .map(|h| h.chars().take(8).collect()),
399 None => None,
400 };
401
402 prompt::ConductOutcome {
403 run_id,
404 unreadable: None,
405 run_status: Some(state.status.as_str().to_owned()),
406 open_findings,
407 rounds_used: state.reviews.len(),
408 rounds_max: state.config.graph.review_rounds,
409 rounds,
410 branch,
411 branch_head,
412 references: crate::refs::describe(&state.seeds),
413 empty_candidate: state.merge.as_ref().is_some_and(|m| m.empty),
414 }
415}
416
417async fn open_pr_fact(task: &Task, repo: &Path) -> Option<String> {
421 let id = task.runs.last()?;
422 let state = match RunState::load(id) {
423 Ok(s) => s,
424 Err(e) => return Some(format!("could not check open pull requests: {e:#}")),
425 };
426 let branch = state.winner()?.branch.clone();
427 let base = state.base_branch.clone();
428 match crate::land::find_open_pr(repo, &branch, &base).await {
429 Ok(crate::land::OpenPr::None) => None,
430 Ok(crate::land::OpenPr::One { url, .. }) => Some(format!(
431 "pull request {url} is already open for {branch} into {base}"
432 )),
433 Ok(crate::land::OpenPr::Many(urls)) => Some(format!(
434 "several pull requests are already open for {branch} into {base}: {}",
435 urls.join(" ")
436 )),
437 Err(e) => Some(format!("could not check open pull requests: {e:#}")),
438 }
439}
440
441async fn attach_facts(cfg: &Config, repo: &Path, queue: &Queue, verdict: &mut Verdict) {
448 for d in verdict
449 .decisions
450 .iter_mut()
451 .filter(|d| d.question.is_some())
452 {
453 let Ok(task) = queue.get(&d.id) else {
454 continue;
455 };
456 let text = format!("{}\n{}", task.title, task.instruction);
457 let pr_fact = open_pr_fact(&task, repo_for(&task, repo).as_path()).await;
458 if crate::refs::scan(&text).is_empty() {
459 if let (Some(fact), Some(q)) = (pr_fact, d.question.as_mut()) {
460 q.push_str(&format!(
461 "\n\nChecked against the repository (magi did this, not the model):\n{fact}"
462 ));
463 }
464 continue;
465 }
466 let repo = repo_for(&task, repo);
467 let remote = &cfg.merge.remote;
468 let base = match cfg.merge.base.clone() {
469 Some(b) => Some(b),
470 None => crate::git::current_branch(&repo).await.ok().flatten(),
471 };
472 let facts = match base {
473 Some(base) => {
474 let base_name = base.clone();
475 let tracking = format!("{remote}/{base}");
476 let refreshed = crate::git::fetch(&repo, remote, &base)
479 .await
480 .is_ok_and(|o| o.ok());
481 let against = if crate::git::rev_exists(&repo, &tracking).await {
482 tracking
483 } else {
484 base
485 };
486 match crate::git::rev_parse(&repo, &against).await {
487 Ok(tip) => crate::refs::describe(
488 &crate::refs::resolve(&repo, &tip, remote, &text).await,
489 )
490 .map(|facts| {
491 if refreshed {
492 facts
493 } else {
494 format!(
495 "{facts}\n(could not fetch {remote}/{base_name}: this is \
496 against the local `{against}`, which may be behind the \
497 remote)"
498 )
499 }
500 }),
501 Err(e) => Some(format!("could not check the repository: {e:#}")),
502 }
503 }
504 None => Some("could not check the repository: no base branch known".to_owned()),
505 };
506 let facts = match (facts, pr_fact) {
507 (Some(f), Some(p)) => Some(format!("{f}\n{p}")),
508 (f, p) => f.or(p),
509 };
510 if let (Some(facts), Some(q)) = (facts, d.question.as_mut()) {
511 q.push_str(&format!(
512 "\n\nChecked against the repository (magi did this, not the model):\n{facts}"
513 ));
514 }
515 }
516}
517
518async fn finished_view(t: &Task, repo: &Path, max_attempts: usize) -> prompt::ConductFinished {
520 prompt::ConductFinished {
521 task: view(t, max_attempts),
522 outcome: outcome_for(t, &repo_for(t, repo)).await,
523 }
524}
525
526fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
529 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
530 fallback.to_path_buf()
531 } else {
532 task.repo.clone()
533 }
534}
535
536fn apply_one(queue: &Queue, questions: &Questions, d: &Decision) -> Result<()> {
553 let _claim = queue
554 .claim(&d.id)
555 .with_context(|| format!("task {} is claimed elsewhere right now", d.id))?;
556 let mut task = queue.get(&d.id).context("no such task")?;
557
558 if task.operator_held() {
562 return Ok(());
563 }
564
565 if task.status == TaskStatus::Held && crate::triage::pending_for(questions, &task) {
575 return Ok(());
576 }
577
578 let pinned = pinned_resume(&task);
581 if pinned && (d.question.is_some() || !d.blocked_by.is_empty()) {
582 return Ok(());
583 }
584
585 if let Some(text) = &d.question {
586 if task.status == TaskStatus::Done {
587 return Ok(());
588 }
589 let question_id = match questions
593 .list()
594 .into_iter()
595 .find(|q| q.status.open() && q.node == NODE && q.run == task.id)
596 {
597 Some(existing) => existing.id,
598 None => {
606 let mut q = Question::new(
607 task.id.clone(),
608 NODE.to_owned(),
609 SEAT.to_owned(),
610 text.clone(),
611 d.reason.clone().unwrap_or_default(),
612 d.choices.clone(),
613 );
614 q.deputy = Some(Deputy::new(crate::deputy::brief(
619 &task.id,
620 d.reason.as_deref().unwrap_or_default(),
621 &q.choices,
622 &q.actions,
623 )));
624 if task.repo.is_dir() && task.repo != Path::new(".") {
625 q.cwd = Some(task.repo.to_string_lossy().into_owned());
626 }
627 questions.put(&mut q)?;
628 q.id
629 }
630 };
631 task.block(vec![question_id], d.reason.clone());
632 return queue.put(&mut task);
633 }
634
635 match task.status {
636 TaskStatus::Queued if !d.blocked_by.is_empty() => {
637 task.block(d.blocked_by.clone(), d.reason.clone());
638 queue.put(&mut task)?;
639 }
640 TaskStatus::Queued if d.recovery == Some(Recovery::Hold) => {
644 if may_hold(&mut task, &hold_note(d)) {
645 task.hold_machine(d.reason.clone());
646 queue.put(&mut task)?;
647 }
648 }
649 TaskStatus::Running => match d.recovery {
650 Some(Recovery::Requeue) if pinned_resume(&task) => {}
651 Some(Recovery::Requeue) => {
652 task.requeue();
653 queue.put(&mut task)?;
654 }
655 Some(Recovery::Hold) if may_hold(&mut task, &hold_note(d)) => {
656 task.hold_machine(d.reason.clone());
657 queue.put(&mut task)?;
658 }
659 _ => {}
663 },
664 TaskStatus::Failed | TaskStatus::Held => match d.recovery {
665 Some(Recovery::Requeue) if pinned_resume(&task) => {}
668 Some(Recovery::Requeue) => {
669 task.requeue();
670 queue.put(&mut task)?;
671 }
672 Some(Recovery::Hold) => {
673 if let Some(reason) = reaffirmed_hold_reason(&task, d) {
677 if may_hold(&mut task, &hold_note(d)) {
678 task.hold_machine(Some(reason));
679 queue.put(&mut task)?;
680 }
681 }
682 }
683 Some(Recovery::Review) => {
684 if let Some(branch) = surviving_branch(&task) {
685 task.request_review(branch);
686 queue.put(&mut task)?;
687 }
688 }
694 Some(Recovery::Done) => {
695 task.succeed();
696 crate::daemon::supersede_prior_runs(&task, &crate::run::home());
700 queue.put(&mut task)?;
701 }
702 None => {}
703 },
704 _ => {}
707 }
708 Ok(())
709}
710
711pub fn apply(queue: &Queue, questions: &Questions, verdict: &Verdict) -> Result<()> {
715 for d in &verdict.decisions {
716 if let Err(e) = apply_one(queue, questions, d) {
717 tracing::warn!("conductor decision for task {}: {e:#}", d.id);
718 }
719 }
720 Ok(())
721}
722
723pub fn seat_path(home: &Path) -> PathBuf {
729 home.join("conduct").join("seat.json")
730}
731
732pub fn load_seat(home: &Path) -> Option<SeatState> {
734 serde_json::from_str(&std::fs::read_to_string(seat_path(home)).ok()?).ok()
735}
736
737fn busy_path(home: &Path) -> PathBuf {
738 home.join("conduct").join("busy")
739}
740
741pub fn busy(home: &Path) -> bool {
747 std::fs::metadata(busy_path(home))
748 .and_then(|m| m.modified())
749 .ok()
750 .and_then(|t| t.elapsed().ok())
751 .is_some_and(|age| age < TURN_TIMEOUT + Duration::from_secs(30))
752}
753
754struct Busy(PathBuf);
756
757impl Busy {
758 fn mark(home: &Path) -> Self {
759 let path = busy_path(home);
760 if let Some(dir) = path.parent() {
761 let _ = std::fs::create_dir_all(dir);
762 }
763 let _ = std::fs::write(&path, std::process::id().to_string());
764 Self(path)
765 }
766}
767
768impl Drop for Busy {
769 fn drop(&mut self) {
770 let _ = std::fs::remove_file(&self.0);
771 }
772}
773
774#[derive(Debug, Default)]
777pub struct Conductor {
778 seat: Option<SeatState>,
779 last_seen: Option<(u64, BTreeSet<String>)>,
780 considered: BTreeMap<String, String>,
784}
785
786fn input_fingerprint(task: &Task) -> String {
791 let mut v = serde_json::to_value(task).unwrap_or_default();
792 if let Some(o) = v.as_object_mut() {
793 o.remove("updated_at");
794 o.remove("hold_reason");
795 if let Some(ro) = o.get_mut("resume_override").and_then(|r| r.as_object_mut()) {
796 ro.remove("conductor_rehold");
797 }
798 }
799 v.to_string()
800}
801
802impl Conductor {
803 #[must_use]
805 pub fn new() -> Self {
806 Self::default()
807 }
808
809 fn settled(
814 considered: &BTreeMap<String, String>,
815 stalled: &[Task],
816 finished: &[Task],
817 ) -> BTreeSet<String> {
818 stalled
819 .iter()
820 .chain(finished)
821 .filter(|t| t.status == TaskStatus::Held)
822 .filter(|t| considered.get(&t.id) == Some(&input_fingerprint(t)))
823 .map(|t| t.id.clone())
824 .collect()
825 }
826
827 fn snapshot_with(
828 considered: &BTreeMap<String, String>,
829 queue: &Queue,
830 stalled: &[Task],
831 finished: &[Task],
832 ) -> (u64, BTreeSet<String>) {
833 let ids = stalled
834 .iter()
835 .chain(finished)
836 .map(|t| t.id.clone())
837 .collect();
838 let skip = Self::settled(considered, stalled, finished);
839 (queue.revision_excluding(&skip), ids)
840 }
841
842 fn snapshot(
843 &self,
844 queue: &Queue,
845 stalled: &[Task],
846 finished: &[Task],
847 ) -> (u64, BTreeSet<String>) {
848 Self::snapshot_with(&self.considered, queue, stalled, finished)
849 }
850
851 #[must_use]
868 pub fn worth_a_look(&self, queue: &Queue, stalled: &[Task], finished: &[Task]) -> bool {
869 self.last_seen.as_ref() != Some(&self.snapshot(queue, stalled, finished))
870 }
871
872 #[allow(clippy::too_many_arguments)]
876 pub async fn maybe_run(
877 &mut self,
878 cfg: &Config,
879 repo: &Path,
880 queue: &Queue,
881 questions: &Questions,
882 home: &Path,
883 queued: &[Task],
884 stalled: &[Task],
885 finished: &[Task],
886 max_attempts: usize,
887 ) {
888 if self.last_seen.as_ref() == Some(&self.snapshot(queue, stalled, finished)) {
889 return;
890 }
891 self.considered = stalled
894 .iter()
895 .chain(finished)
896 .filter(|t| t.status == TaskStatus::Held)
897 .map(|t| (t.id.clone(), input_fingerprint(t)))
898 .collect();
899 self.last_seen = Some(self.snapshot(queue, stalled, finished));
900 if let Err(e) = self
901 .run_once(
902 cfg,
903 repo,
904 queue,
905 questions,
906 home,
907 queued,
908 stalled,
909 finished,
910 max_attempts,
911 )
912 .await
913 {
914 tracing::warn!("conductor: {e:#}");
920 }
921 }
922
923 #[allow(clippy::too_many_arguments)]
924 async fn run_once(
925 &mut self,
926 cfg: &Config,
927 repo: &Path,
928 queue: &Queue,
929 questions: &Questions,
930 home: &Path,
931 queued: &[Task],
932 stalled: &[Task],
933 finished: &[Task],
934 max_attempts: usize,
935 ) -> Result<()> {
936 if queued.is_empty() && stalled.is_empty() && finished.is_empty() {
937 return Ok(());
938 }
939
940 let primary = cfg
941 .resolve_roles()
942 .context("resolving the conductor seat")?
943 .conductor;
944 let chain = match cfg.roles.conductor.as_ref() {
948 Some(c) if c.ids().len() > 1 => {
949 agent::pick_chain(&cfg.agents, Some(c), &agent::installed, "conductor")?
950 }
951 _ => vec![primary],
952 };
953
954 let runnable_views: Vec<prompt::ConductTask> =
955 queued.iter().map(|t| view(t, max_attempts)).collect();
956 let stalled_views: Vec<prompt::ConductTask> =
957 stalled.iter().map(|t| view(t, max_attempts)).collect();
958 let mut finished_views = Vec::with_capacity(finished.len());
959 for t in finished {
960 finished_views.push(finished_view(t, repo, max_attempts).await);
961 }
962
963 let body = prompt::with_overlay(
964 prompt::conduct(
965 &runnable_views,
966 &stalled_views,
967 &finished_views,
968 &cfg.graph.language,
969 ),
970 cfg.prompts.overlay(NODE),
971 );
972
973 let artifacts = home.join("conduct").join("artifacts");
974 let cache_dir = cfg.cache_dir();
977
978 let busy = Busy::mark(home);
980 let mut result: Result<Verdict> = Err(anyhow::anyhow!("no conductor agent ran"));
981 for spec in &chain {
982 let needs_new_seat = !matches!(&self.seat, Some(s) if s.agent == spec.id);
985 if needs_new_seat {
986 self.seat = Some(SeatState::new(SEAT, &spec.id, crate::rng::entropy()));
987 }
988 let seat = self.seat.as_mut().expect("just ensured a seat exists");
989 let stem = if std::ptr::eq(spec, &chain[0]) {
992 format!("turn-{}", seat.turns + 1)
993 } else {
994 format!("turn-{}-{}", seat.turns + 1, spec.id)
995 };
996 let inv = Invocation {
997 cwd: repo,
998 prompt: &body,
999 timeout: TURN_TIMEOUT,
1000 allow_write: false,
1003 sessions: cfg.graph.sessions,
1004 artifacts: &artifacts,
1005 stem: &stem,
1006 run: NODE,
1007 node: NODE,
1008 cache_dir: cache_dir.as_deref(),
1009 attachments: &[],
1010 writable: &[],
1011 };
1012 let out = agent::invoke(spec, seat, &inv).await;
1013 if let Ok(body) = serde_json::to_string(&*seat) {
1016 let path = seat_path(home);
1017 if let Some(dir) = path.parent() {
1018 let _ = std::fs::create_dir_all(dir);
1019 }
1020 let _ = std::fs::write(path, body);
1021 }
1022 let advance = agent::chain_advances(&out);
1023 result = match out {
1024 Err(e) => Err(e.context("invoking the conductor")),
1025 Ok(out) if advance => Err(anyhow::anyhow!(
1026 "no usable reply (exit {:?}, timed out {})",
1027 out.exit_code,
1028 out.timed_out
1029 )),
1030 Ok(out) => verdict::extract_json(&out.text)
1031 .context("the conductor's reply could not be parsed"),
1032 };
1033 if result.is_ok() {
1034 break;
1035 }
1036 if chain.len() > 1 {
1037 tracing::warn!("conductor: `{}` failed, trying the next agent", spec.id);
1038 }
1039 }
1040 drop(busy);
1041 let mut verdict = result?;
1042 attach_facts(cfg, repo, queue, &mut verdict).await;
1043 apply(queue, questions, &verdict)
1044 }
1045}
1046
1047#[cfg(test)]
1048mod tests {
1049 use std::collections::BTreeMap;
1050
1051 use tempfile::tempdir;
1052
1053 use super::*;
1054 use crate::ask::{Answer, QuestionStatus};
1055 use crate::config::{AgentKind, AgentSpec, Graph};
1056 use crate::queue::Source;
1057
1058 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
1059 let path = dir.join("mock-conduct-agent.sh");
1060 std::fs::write(&path, script).expect("write mock");
1061 AgentSpec {
1062 id: "mock".to_owned(),
1063 kind: AgentKind::Command,
1064 model: None,
1065 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1066 extra_args: Vec::new(),
1067 env,
1068 prompt_delivery: None,
1069 }
1070 }
1071
1072 fn config(spec: AgentSpec) -> Config {
1073 Config {
1074 agents: vec![spec],
1075 graph: Graph {
1076 language: "en".to_owned(),
1077 ..Graph::default()
1078 },
1079 ..Config::default()
1080 }
1081 }
1082
1083 fn task(title: &str) -> Task {
1084 Task::new(
1085 title.to_owned(),
1086 format!("do {title}"),
1087 std::path::PathBuf::from("."),
1088 Source::Human,
1089 )
1090 }
1091
1092 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
1093 const GARBAGE: &str = "#!/bin/sh\ncat >/dev/null\nprintf 'not json at all\\n'\n";
1094
1095 fn env(reply: &str) -> BTreeMap<String, String> {
1096 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
1097 }
1098
1099 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
1100
1101 fn init_repo_with_branch(dir: &Path, branch: &str) {
1105 use crate::proc::Quiet as _;
1106 let run = |args: &[&str]| {
1107 let out = std::process::Command::new("git")
1108 .args(args)
1109 .current_dir(dir)
1110 .quiet()
1111 .output()
1112 .expect("spawn git");
1113 assert!(
1114 out.status.success(),
1115 "git {args:?} failed: {}",
1116 String::from_utf8_lossy(&out.stderr)
1117 );
1118 };
1119 run(&["init", "-b", "main"]);
1120 run(&["config", "user.name", "magi test"]);
1121 run(&["config", "user.email", "magi@example.com"]);
1122 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
1123 run(&["add", "-A"]);
1124 run(&["commit", "-m", "init"]);
1125 run(&["checkout", "-b", branch]);
1126 std::fs::write(dir.join("change.txt"), "x\n").unwrap();
1127 run(&["add", "-A"]);
1128 run(&["commit", "-m", "candidate work"]);
1129 }
1130
1131 fn review_round_with_finding(
1132 round: usize,
1133 finding_id: &str,
1134 title: &str,
1135 addressed: &[&str],
1136 rejected: &[(&str, &str)],
1137 ) -> crate::run::ReviewRound {
1138 crate::run::ReviewRound {
1139 round,
1140 head: "deadbeef".to_owned(),
1141 verified_head: None,
1142 verified_at: None,
1143 reviews: vec![crate::run::ReviewRecord {
1144 attempts: 0,
1145 reviewer: 1,
1146 agent: "mock".to_owned(),
1147 summary: String::new(),
1148 findings: vec![crate::verdict::Finding {
1149 id: finding_id.to_owned(),
1150 severity: crate::verdict::Severity::Major,
1151 file: None,
1152 line: None,
1153 title: title.to_owned(),
1154 detail: String::new(),
1155 }],
1156 vote: None,
1157 failed: None,
1158 duration_ms: 0,
1159 }],
1160 e2e: Vec::new(),
1161 verify_retried: false,
1162 e2e_deferred: false,
1163 e2e_defer_reason: None,
1164 fix: Some(crate::run::FixRecord {
1165 agent: "mock".to_owned(),
1166 addressed: addressed.iter().map(|s| (*s).to_owned()).collect(),
1167 rejected: rejected
1168 .iter()
1169 .map(|(id, why)| crate::verdict::Rejection {
1170 id: (*id).to_owned(),
1171 why: (*why).to_owned(),
1172 })
1173 .collect(),
1174 notes: String::new(),
1175 committed: false,
1176 failed: None,
1177 duration_ms: 0,
1178 continuation: None,
1179 }),
1180 blocking: 1,
1181 answered: 1,
1182 expected: 1,
1183 clean: false,
1184 progressed: true,
1185 vote_split: false,
1186 reconsideration: Vec::new(),
1187 verdict: None,
1188 }
1189 }
1190
1191 #[test]
1192 fn outcome_for_carries_every_rounds_findings_and_the_branch_head() {
1193 crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
1194 let dir = tempdir().unwrap();
1195 let default_repo = dir.path().join("default");
1196 let task_repo = dir.path().join("task");
1197 std::fs::create_dir_all(&default_repo).unwrap();
1198 std::fs::create_dir_all(&task_repo).unwrap();
1199 init_repo_with_branch(&default_repo, "other-branch");
1200 init_repo_with_branch(&task_repo, "magi/f00d/A");
1201
1202 let mut config = Config::default();
1203 config.graph.review_rounds = 6;
1204 let mut state = crate::run::RunState::new(
1205 task_repo.clone(),
1206 "main".to_owned(),
1207 "deadbeef".to_owned(),
1208 "task".to_owned(),
1209 config,
1210 );
1211 state.status = crate::run::RunStatus::Blocked;
1212 state.candidates.push(crate::run::Candidate {
1213 index: 0,
1214 label: 'A',
1215 agent: "mock".to_owned(),
1216 branch: "magi/f00d/A".to_owned(),
1217 worktree: task_repo.clone(),
1218 summary: String::new(),
1219 stat: String::new(),
1220 files: 1,
1221 commits: 1,
1222 empty: false,
1223 failed: None,
1224 verified_noop: None,
1225 duration_ms: 0,
1226 folded: false,
1227 });
1228 state.tally = Some(crate::run::Tally {
1229 first_choice: std::collections::BTreeMap::new(),
1230 borda: std::collections::BTreeMap::new(),
1231 winner: 'A',
1232 rankings: 0,
1233 unanimous_initial: false,
1234 deliberated: false,
1235 changed_votes: 0,
1236 unanimous_final: false,
1237 tie_break: None,
1238 judges: 0,
1239 present: 0,
1240 quorum: 0,
1241 met_quorum: true,
1242 uncontested: Some("solo".to_owned()),
1243 });
1244 state.reviews = vec![
1245 review_round_with_finding(
1246 1,
1247 "R1-1-2",
1248 "answer content is dropped",
1249 &[],
1250 &[("R1-1-2", "the id leaving blocked_by is enough")],
1251 ),
1252 review_round_with_finding(2, "R2-1-3", "answer content is still dropped", &[], &[]),
1253 ];
1254 state.save().unwrap();
1255
1256 let mut t = task("outcome test");
1257 t.repo = task_repo;
1258 t.runs.push(state.id.clone());
1259
1260 let finished = tokio_test_block_on(finished_view(&t, &default_repo, 2));
1261 let outcome = finished.outcome;
1262
1263 assert!(outcome.unreadable.is_none());
1264 assert_eq!(outcome.run_status.as_deref(), Some("blocked"));
1265 assert_eq!(outcome.rounds_used, 2);
1266 assert_eq!(outcome.rounds_max, 6);
1267 assert_eq!(outcome.rounds.len(), 2);
1268 assert_eq!(outcome.rounds[0].findings[0].id, "R1-1-2");
1269 assert_eq!(outcome.rounds[0].rejected[0].id, "R1-1-2");
1270 assert!(outcome.rounds[1].addressed.is_empty());
1271 assert!(outcome.rounds[1].rejected.is_empty());
1272 assert_eq!(outcome.branch.as_deref(), Some("magi/f00d/A"));
1273 assert!(
1274 outcome.branch_head.is_some(),
1275 "a real branch must resolve a head commit: {outcome:?}"
1276 );
1277 }
1278
1279 fn tokio_test_block_on<F: std::future::Future>(f: F) -> F::Output {
1283 tokio::runtime::Builder::new_current_thread()
1284 .enable_all()
1285 .build()
1286 .unwrap()
1287 .block_on(f)
1288 }
1289
1290 #[test]
1291 fn view_carries_a_tasks_recorded_answers_into_the_conductor_prompt_input() {
1292 let mut t = task("answered");
1293 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
1294 let v = view(&t, 2);
1295 assert_eq!(v.answers.len(), 1);
1296 assert_eq!(v.answers[0].question, "Which backend?");
1297 assert_eq!(v.answers[0].answer, "SQLite");
1298 }
1299
1300 #[test]
1301 fn a_dependency_decision_blocks_the_task_and_leaves_priority_alone() {
1302 let dir = tempdir().unwrap();
1303 let queue = Queue::at(dir.path().join("queue"));
1304 let questions = Questions::at(dir.path().join("questions"));
1305 let mut a = task("a");
1306 a.priority = 9;
1307 queue.put(&mut a).unwrap();
1308
1309 let verdict = Verdict {
1310 decisions: vec![Decision {
1311 id: a.id.clone(),
1312 blocked_by: vec!["20260101-000000-dead".to_owned()],
1313 reason: Some("waits on the other task".to_owned()),
1314 recovery: None,
1315 question: None,
1316 choices: Vec::new(),
1317 }],
1318 };
1319 apply(&queue, &questions, &verdict).unwrap();
1320
1321 let back = queue.get(&a.id).unwrap();
1322 assert_eq!(back.status, TaskStatus::Blocked);
1323 assert_eq!(back.blocked_by, ["20260101-000000-dead"]);
1324 assert_eq!(
1325 back.priority, 9,
1326 "the conductor's reply cannot carry priority"
1327 );
1328 }
1329
1330 #[test]
1331 fn a_question_decision_files_one_and_blocks_on_its_id() {
1332 let dir = tempdir().unwrap();
1333 let queue = Queue::at(dir.path().join("queue"));
1334 let questions = Questions::at(dir.path().join("questions"));
1335 let mut t = task("ambiguous");
1336 queue.put(&mut t).unwrap();
1337
1338 let verdict = Verdict {
1339 decisions: vec![Decision {
1340 id: t.id.clone(),
1341 blocked_by: Vec::new(),
1342 reason: Some("which backend?".to_owned()),
1343 recovery: None,
1344 question: Some("Which storage backend?".to_owned()),
1345 choices: vec!["SQLite".to_owned(), "Redis".to_owned()],
1346 }],
1347 };
1348 apply(&queue, &questions, &verdict).unwrap();
1349
1350 let back = queue.get(&t.id).unwrap();
1351 assert_eq!(back.status, TaskStatus::Blocked);
1352 assert_eq!(back.blocked_by.len(), 1);
1353 let q = questions.get(&back.blocked_by[0]).unwrap();
1354 assert_eq!(q.summary, "Which storage backend?");
1355 assert_eq!(q.node, NODE);
1356 assert!(q.status.open());
1357 }
1358
1359 #[test]
1360 fn a_task_with_an_open_question_already_reuses_it_rather_than_filing_a_second_one() {
1361 let dir = tempdir().unwrap();
1362 let queue = Queue::at(dir.path().join("queue"));
1363 let questions = Questions::at(dir.path().join("questions"));
1364 let mut t = task("asked once");
1365 queue.put(&mut t).unwrap();
1366
1367 let decision = Decision {
1368 id: t.id.clone(),
1369 reason: Some("still deciding".to_owned()),
1370 question: Some("Which backend?".to_owned()),
1371 ..Decision::default()
1372 };
1373 apply(
1374 &queue,
1375 &questions,
1376 &Verdict {
1377 decisions: vec![decision.clone()],
1378 },
1379 )
1380 .unwrap();
1381 assert_eq!(questions.list().len(), 1);
1382 let first_question_id = queue.get(&t.id).unwrap().blocked_by[0].clone();
1383
1384 let mut released = queue.get(&t.id).unwrap();
1389 released.release();
1390 queue.put(&mut released).unwrap();
1391
1392 apply(
1393 &queue,
1394 &questions,
1395 &Verdict {
1396 decisions: vec![decision],
1397 },
1398 )
1399 .unwrap();
1400 assert_eq!(questions.list().len(), 1, "no duplicate question was filed");
1401 let after = queue.get(&t.id).unwrap();
1402 assert_eq!(
1403 after.blocked_by,
1404 [first_question_id],
1405 "the existing open question is reused, not replaced"
1406 );
1407 }
1408
1409 #[test]
1410 fn a_same_id_question_from_another_node_is_not_reused() {
1411 let dir = tempdir().unwrap();
1412 let queue = Queue::at(dir.path().join("queue"));
1413 let questions = Questions::at(dir.path().join("questions"));
1414 let mut t = task("must ask the conductor");
1415 queue.put(&mut t).unwrap();
1416
1417 let mut unrelated = Question::new(
1418 t.id.clone(),
1419 "review".to_owned(),
1420 "reviewer-1".to_owned(),
1421 "An unrelated review question".to_owned(),
1422 String::new(),
1423 Vec::new(),
1424 );
1425 questions.put(&mut unrelated).unwrap();
1426
1427 apply(
1428 &queue,
1429 &questions,
1430 &Verdict {
1431 decisions: vec![Decision {
1432 id: t.id.clone(),
1433 question: Some("Which backend?".to_owned()),
1434 ..Decision::default()
1435 }],
1436 },
1437 )
1438 .unwrap();
1439
1440 let blocked_by = &queue.get(&t.id).unwrap().blocked_by;
1441 assert_eq!(blocked_by.len(), 1);
1442 assert_ne!(blocked_by[0], unrelated.id);
1443 assert!(questions.get(&unrelated.id).unwrap().status.open());
1444 assert_eq!(questions.get(&blocked_by[0]).unwrap().node, NODE);
1445 }
1446
1447 #[test]
1448 fn answering_the_question_lets_the_resolver_clear_the_block_with_the_answer_kept() {
1449 let dir = tempdir().unwrap();
1450 let queue = Queue::at(dir.path().join("queue"));
1451 let questions = Questions::at(dir.path().join("questions"));
1452 let mut t = task("waits on an answer");
1453 queue.put(&mut t).unwrap();
1454
1455 apply(
1456 &queue,
1457 &questions,
1458 &Verdict {
1459 decisions: vec![Decision {
1460 id: t.id.clone(),
1461 blocked_by: Vec::new(),
1462 reason: None,
1463 recovery: None,
1464 question: Some("Which backend?".to_owned()),
1465 choices: Vec::new(),
1466 }],
1467 },
1468 )
1469 .unwrap();
1470 let blocked = queue.get(&t.id).unwrap();
1471 let question_id = blocked.blocked_by[0].clone();
1472
1473 let mut q = questions.get(&question_id).unwrap();
1474 q.answer(Answer::Text("SQLite".to_owned())).unwrap();
1475 questions.put(&mut q).unwrap();
1476 assert_eq!(q.status, QuestionStatus::Answered);
1477
1478 let mut task_after = queue.get(&t.id).unwrap();
1482 task_after.record_answer(q.summary.clone(), "SQLite".to_owned());
1483 task_after.unblock(&question_id);
1484 assert_eq!(task_after.status, TaskStatus::Queued);
1485 assert_eq!(task_after.answers[0].answer, "SQLite");
1486 }
1487
1488 #[test]
1489 fn a_stalled_task_can_be_requeued_or_held() {
1490 let dir = tempdir().unwrap();
1491 let queue = Queue::at(dir.path().join("queue"));
1492 let questions = Questions::at(dir.path().join("questions"));
1493
1494 let mut requeue_me = task("stuck a");
1495 requeue_me.start("run-1".to_owned());
1496 queue.put(&mut requeue_me).unwrap();
1497
1498 let mut hold_me = task("stuck b");
1499 hold_me.start("run-2".to_owned());
1500 queue.put(&mut hold_me).unwrap();
1501
1502 apply(
1503 &queue,
1504 &questions,
1505 &Verdict {
1506 decisions: vec![
1507 Decision {
1508 id: requeue_me.id.clone(),
1509 recovery: Some(Recovery::Requeue),
1510 ..Decision::default()
1511 },
1512 Decision {
1513 id: hold_me.id.clone(),
1514 recovery: Some(Recovery::Hold),
1515 reason: Some("looks broken".to_owned()),
1516 ..Decision::default()
1517 },
1518 ],
1519 },
1520 )
1521 .unwrap();
1522
1523 let requeued = queue.get(&requeue_me.id).unwrap();
1524 assert_eq!(requeued.status, TaskStatus::Queued);
1525 assert_eq!(requeued.attempts, 0);
1526
1527 let held = queue.get(&hold_me.id).unwrap();
1528 assert_eq!(held.status, TaskStatus::Held);
1529 assert_eq!(held.hold_reason.as_deref(), Some("looks broken"));
1530 }
1531
1532 #[test]
1533 fn a_machine_held_task_asked_about_restores_to_held_once_answered() {
1534 let dir = tempdir().unwrap();
1539 let queue = Queue::at(dir.path().join("queue"));
1540 let questions = Questions::at(dir.path().join("questions"));
1541 let mut t = task("held out of attempts");
1542 t.hold_machine(Some("out of attempts".to_owned()));
1543 queue.put(&mut t).unwrap();
1544
1545 apply(
1546 &queue,
1547 &questions,
1548 &Verdict {
1549 decisions: vec![Decision {
1550 id: t.id.clone(),
1551 reason: Some("what should happen to this one?".to_owned()),
1552 question: Some("Hold it, or try again?".to_owned()),
1553 ..Decision::default()
1554 }],
1555 },
1556 )
1557 .unwrap();
1558 let blocked = queue.get(&t.id).unwrap();
1559 assert_eq!(blocked.status, TaskStatus::Blocked);
1560 let question_id = blocked.blocked_by[0].clone();
1561
1562 let mut q = questions.get(&question_id).unwrap();
1563 q.answer(Answer::Text("leave it held".to_owned())).unwrap();
1564 questions.put(&mut q).unwrap();
1565
1566 let mut after = queue.get(&t.id).unwrap();
1568 after.record_answer(q.summary.clone(), "leave it held".to_owned());
1569 after.unblock(&question_id);
1570 assert_eq!(
1571 after.status,
1572 TaskStatus::Held,
1573 "must not fall back to queued"
1574 );
1575 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
1576 }
1577
1578 #[test]
1579 fn a_reaffirmed_hold_with_no_new_reason_is_not_silently_auto_released_by_triage() {
1580 let dir = tempdir().unwrap();
1590 let queue = Queue::at(dir.path().join("queue"));
1591 let questions = Questions::at(dir.path().join("questions"));
1592 let mut t = task("disk pressure, then reconsidered");
1593 t.hold_machine(Some(
1594 "not enough free space to start a run: 10 bytes free, 100 required by \
1595 `[disk] min_free_bytes`"
1596 .to_owned(),
1597 ));
1598 t.record_answer(
1599 "How should this be handled?".to_owned(),
1600 "keep it held, a human will look at it later".to_owned(),
1601 );
1602 queue.put(&mut t).unwrap();
1603
1604 apply(
1605 &queue,
1606 &questions,
1607 &Verdict {
1608 decisions: vec![Decision {
1609 id: t.id.clone(),
1610 recovery: Some(Recovery::Hold),
1611 ..Decision::default()
1612 }],
1613 },
1614 )
1615 .unwrap();
1616
1617 let after = queue.get(&t.id).unwrap();
1618 assert_eq!(after.status, TaskStatus::Held);
1619 assert!(
1620 !after
1621 .hold_reason
1622 .as_deref()
1623 .unwrap_or_default()
1624 .starts_with("not enough free space"),
1625 "the stale disk-pressure text must not survive a reconfirmed hold: {:?}",
1626 after.hold_reason
1627 );
1628
1629 let cfg_dir = tempdir().unwrap();
1632 let config = cfg_dir.path().join("magi.toml");
1633 std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1634 let report =
1635 crate::triage::run_once(&queue, &questions, Some(&config), jiff::Timestamp::now());
1636 assert!(report.resumed.is_empty(), "must not be auto-released");
1637 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1638 }
1639
1640 #[test]
1641 fn a_conductor_rehold_after_a_resume_answer_is_not_asked_again_identically() {
1642 let dir = tempdir().unwrap();
1643 let queue = Queue::at(dir.path().join("queue"));
1644 let questions = Questions::at(dir.path().join("questions"));
1645 let config = dir.path().join("magi.toml");
1646 std::fs::write(&config, "[disk]\nmin_free_bytes = 0\n").unwrap();
1647 let now = jiff::Timestamp::now();
1648 let triage = || crate::triage::run_once(&queue, &questions, Some(&config), now);
1649 let open = || {
1650 questions
1651 .list()
1652 .into_iter()
1653 .filter(|q| q.node == "triage" && q.status.open())
1654 .collect::<Vec<_>>()
1655 };
1656 let hold = Decision {
1657 recovery: Some(Recovery::Hold),
1658 reason: Some("waiting on manual worktree cleanup".to_owned()),
1659 ..Decision::default()
1660 };
1661
1662 let mut t = task("looping hold");
1663 t.hold_machine(Some("waiting on manual worktree cleanup".to_owned()));
1664 queue.put(&mut t).unwrap();
1665 let hold = Decision {
1666 id: t.id.clone(),
1667 ..hold
1668 };
1669
1670 assert_eq!(triage().asked.len(), 1);
1672 let first = open().remove(0);
1673 let mut q = questions.get(&first.id).unwrap();
1674 let resume = q.choices[0].clone();
1675 q.answer(Answer::Choice(resume)).unwrap();
1676 questions.put(&mut q).unwrap();
1677 assert_eq!(triage().answered.len(), 1);
1678 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1679
1680 apply(
1682 &queue,
1683 &questions,
1684 &Verdict {
1685 decisions: vec![hold.clone()],
1686 },
1687 )
1688 .unwrap();
1689 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
1690
1691 assert_eq!(triage().asked.len(), 1);
1693 let second = open().remove(0);
1694 assert_ne!(second.summary, first.summary);
1695 assert_ne!(second.choices, first.choices);
1696 assert!(second.detail.contains("waiting on manual worktree cleanup"));
1697 assert!(triage().asked.is_empty(), "no duplicate question");
1698 assert_eq!(open().len(), 1);
1699
1700 let mut q = questions.get(&second.id).unwrap();
1702 let force = q.choices[0].clone();
1703 q.answer(Answer::Choice(force)).unwrap();
1704 questions.put(&mut q).unwrap();
1705 assert_eq!(triage().answered.len(), 1);
1706 apply(
1707 &queue,
1708 &questions,
1709 &Verdict {
1710 decisions: vec![hold],
1711 },
1712 )
1713 .unwrap();
1714 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
1715 }
1716
1717 #[test]
1718 fn a_runnable_task_can_be_held_directly_without_a_question() {
1719 let dir = tempdir().unwrap();
1720 let queue = Queue::at(dir.path().join("queue"));
1721 let questions = Questions::at(dir.path().join("questions"));
1722 let mut t = task("already answered, should stay put");
1723 queue.put(&mut t).unwrap();
1724
1725 apply(
1726 &queue,
1727 &questions,
1728 &Verdict {
1729 decisions: vec![Decision {
1730 id: t.id.clone(),
1731 recovery: Some(Recovery::Hold),
1732 reason: Some("operator already said keep this held".to_owned()),
1733 ..Decision::default()
1734 }],
1735 },
1736 )
1737 .unwrap();
1738
1739 let after = queue.get(&t.id).unwrap();
1740 assert_eq!(after.status, TaskStatus::Held);
1741 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
1742 }
1743
1744 #[test]
1745 fn done_recovery_closes_a_held_task_whose_goal_is_already_met() {
1746 let dir = tempdir().unwrap();
1752 let queue = Queue::at(dir.path().join("queue"));
1753 let questions = Questions::at(dir.path().join("questions"));
1754 let mut t = task("already merged by hand");
1755 t.hold_machine(Some("branch survived, awaiting a decision".to_owned()));
1756 t.record_answer(
1757 "Handle this one?".to_owned(),
1758 "already merged and cleaned up, close it".to_owned(),
1759 );
1760 queue.put(&mut t).unwrap();
1761
1762 apply(
1763 &queue,
1764 &questions,
1765 &Verdict {
1766 decisions: vec![Decision {
1767 id: t.id.clone(),
1768 recovery: Some(Recovery::Done),
1769 reason: Some("operator confirmed this already landed".to_owned()),
1770 ..Decision::default()
1771 }],
1772 },
1773 )
1774 .unwrap();
1775
1776 let after = queue.get(&t.id).unwrap();
1777 assert_eq!(after.status, TaskStatus::Done);
1778 assert!(after.hold_reason.is_none());
1779 assert_eq!(after.answers.len(), 1, "the record of why is kept");
1780 }
1781
1782 #[test]
1783 fn done_recovery_is_ignored_for_a_runnable_or_running_task() {
1784 let dir = tempdir().unwrap();
1785 let queue = Queue::at(dir.path().join("queue"));
1786 let questions = Questions::at(dir.path().join("questions"));
1787
1788 let mut queued = task("never ran yet");
1789 queue.put(&mut queued).unwrap();
1790
1791 let mut running = task("mid-run");
1792 running.start("run-1".to_owned());
1793 queue.put(&mut running).unwrap();
1794
1795 for id in [queued.id.clone(), running.id.clone()] {
1796 apply(
1797 &queue,
1798 &questions,
1799 &Verdict {
1800 decisions: vec![Decision {
1801 id,
1802 recovery: Some(Recovery::Done),
1803 ..Decision::default()
1804 }],
1805 },
1806 )
1807 .unwrap();
1808 }
1809
1810 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
1811 assert_eq!(queue.get(&running.id).unwrap().status, TaskStatus::Running);
1812 }
1813
1814 #[test]
1815 fn a_question_after_two_settled_answers_is_still_filed_and_blocks() {
1816 let dir = tempdir().unwrap();
1819 let queue = Queue::at(dir.path().join("queue"));
1820 let questions = Questions::at(dir.path().join("questions"));
1821 let mut t = task("asked about repeatedly");
1822 t.hold_machine(Some("out of attempts".to_owned()));
1823 t.record_answer("Handle this one? (1)".to_owned(), "not yet".to_owned());
1824 t.record_answer(
1825 "Handle this one? (2)".to_owned(),
1826 "still not yet".to_owned(),
1827 );
1828 queue.put(&mut t).unwrap();
1829 assert_eq!(questions.list().len(), 0);
1830
1831 let verdict = Verdict {
1832 decisions: vec![Decision {
1833 id: t.id.clone(),
1834 question: Some("Branch conflicts with origin/main, how do we proceed?".to_owned()),
1835 ..Decision::default()
1836 }],
1837 };
1838 apply(&queue, &questions, &verdict).unwrap();
1839
1840 let filed = questions.list();
1841 assert_eq!(filed.len(), 1, "the question was filed");
1842 assert_eq!(
1843 filed[0].summary,
1844 "Branch conflicts with origin/main, how do we proceed?"
1845 );
1846 let deputy = filed[0]
1847 .deputy
1848 .as_ref()
1849 .expect("the wait is handed to a deputy");
1850 assert!(deputy.brief.contains(&t.id), "{}", deputy.brief);
1851 assert_eq!(
1852 deputy.starts, 0,
1853 "filing never starts anything: the loop does not wait"
1854 );
1855 assert!(filed[0].status.open());
1856 assert_eq!(filed[0].node, NODE);
1857 let after = queue.get(&t.id).unwrap();
1858 assert_eq!(after.status, TaskStatus::Blocked);
1859 assert_eq!(after.blocked_by, vec![filed[0].id.clone()]);
1860 assert_eq!(after.answers.len(), 2, "the prior answers are untouched");
1861 assert!(
1862 !after
1863 .hold_reason
1864 .clone()
1865 .unwrap_or_default()
1866 .contains("conduct tried to ask"),
1867 "no hold was applied"
1868 );
1869
1870 apply(&queue, &questions, &verdict).unwrap();
1871 let again = questions.list();
1872 assert_eq!(again.len(), 1, "the open question is reused");
1873 assert_eq!(again[0].id, filed[0].id);
1874 }
1875
1876 #[test]
1877 fn a_second_conductor_question_is_still_allowed_after_one_settled_answer() {
1878 let dir = tempdir().unwrap();
1879 let queue = Queue::at(dir.path().join("queue"));
1880 let questions = Questions::at(dir.path().join("questions"));
1881 let mut t = task("asked about once already");
1882 t.hold_machine(Some("out of attempts".to_owned()));
1883 t.record_answer("Handle this one?".to_owned(), "not yet".to_owned());
1884 queue.put(&mut t).unwrap();
1885
1886 apply(
1887 &queue,
1888 &questions,
1889 &Verdict {
1890 decisions: vec![Decision {
1891 id: t.id.clone(),
1892 question: Some("Still not sure - now what?".to_owned()),
1893 ..Decision::default()
1894 }],
1895 },
1896 )
1897 .unwrap();
1898
1899 assert_eq!(questions.list().len(), 1, "the second question was filed");
1900 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
1901 }
1902
1903 #[test]
1904 fn a_held_task_with_an_open_triage_question_is_left_to_triage() {
1905 let dir = tempdir().unwrap();
1911 let queue = Queue::at(dir.path().join("queue"));
1912 let questions = Questions::at(dir.path().join("questions"));
1913 let mut t = task("held, triage already asking about it");
1914 t.hold_machine(Some("cause unclear".to_owned()));
1915 queue.put(&mut t).unwrap();
1916
1917 let mut triage_q = Question::new(
1918 t.id.clone(),
1919 crate::triage::NODE.to_owned(),
1920 "triage".to_owned(),
1921 "Still needed?".to_owned(),
1922 String::new(),
1923 vec![
1924 "resume".to_owned(),
1925 "not yet".to_owned(),
1926 "discard".to_owned(),
1927 ],
1928 );
1929 questions.put(&mut triage_q).unwrap();
1930
1931 for decision in [
1932 Decision {
1933 id: t.id.clone(),
1934 question: Some("what now?".to_owned()),
1935 ..Decision::default()
1936 },
1937 Decision {
1938 id: t.id.clone(),
1939 recovery: Some(Recovery::Requeue),
1940 ..Decision::default()
1941 },
1942 ] {
1943 apply(
1944 &queue,
1945 &questions,
1946 &Verdict {
1947 decisions: vec![decision],
1948 },
1949 )
1950 .unwrap();
1951 }
1952
1953 let after = queue.get(&t.id).unwrap();
1954 assert_eq!(
1955 after.status,
1956 TaskStatus::Held,
1957 "triage still owns this hold"
1958 );
1959 assert!(after.blocked_by.is_empty());
1960 assert_eq!(
1961 questions.list().len(),
1962 1,
1963 "no second, conductor-owned question was filed"
1964 );
1965 }
1966
1967 #[test]
1968 fn a_held_task_with_an_answered_but_unapplied_triage_question_is_still_left_alone() {
1969 let dir = tempdir().unwrap();
1977 let queue = Queue::at(dir.path().join("queue"));
1978 let questions = Questions::at(dir.path().join("questions"));
1979 let mut t = task("held, triage question answered but not yet applied");
1980 t.hold_machine(Some("cause unclear".to_owned()));
1981 queue.put(&mut t).unwrap();
1982
1983 let mut triage_q = Question::new(
1984 t.id.clone(),
1985 crate::triage::NODE.to_owned(),
1986 "triage".to_owned(),
1987 "Still needed?".to_owned(),
1988 String::new(),
1989 vec![
1990 "resume".to_owned(),
1991 "not yet".to_owned(),
1992 "discard".to_owned(),
1993 ],
1994 );
1995 questions.put(&mut triage_q).unwrap();
1996 triage_q
1997 .answer(Answer::Choice("not yet".to_owned()))
1998 .unwrap();
1999 questions.put(&mut triage_q).unwrap();
2000 assert!(!triage_q.status.open());
2001
2002 apply(
2003 &queue,
2004 &questions,
2005 &Verdict {
2006 decisions: vec![Decision {
2007 id: t.id.clone(),
2008 question: Some("what now?".to_owned()),
2009 ..Decision::default()
2010 }],
2011 },
2012 )
2013 .unwrap();
2014
2015 let after = queue.get(&t.id).unwrap();
2016 assert_eq!(
2017 after.status,
2018 TaskStatus::Held,
2019 "triage's own answer is not yet applied - conduct must wait"
2020 );
2021 assert_eq!(
2022 questions.list().len(),
2023 1,
2024 "no conductor question was filed over the pending triage answer"
2025 );
2026 }
2027
2028 #[test]
2029 fn manual_hold_rejects_hostile_or_stale_conductor_recovery() {
2030 let dir = tempdir().unwrap();
2031 let queue = Queue::at(dir.path().join("queue"));
2032 let questions = Questions::at(dir.path().join("questions"));
2033 let mut held = task("manual recovery");
2034 held.priority = 300;
2035 held.runs.push("run20260912-224242-daf5".to_owned());
2036 held.hold_manual(Some(
2037 "active manual recovery run20260912-224242-daf5".to_owned(),
2038 ));
2039 queue.put(&mut held).unwrap();
2040
2041 for decision in [
2045 Decision {
2046 id: held.id.clone(),
2047 recovery: Some(Recovery::Requeue),
2048 ..Decision::default()
2049 },
2050 Decision {
2051 id: held.id.clone(),
2052 recovery: Some(Recovery::Hold),
2053 reason: Some("stale replacement reason".to_owned()),
2054 ..Decision::default()
2055 },
2056 Decision {
2057 id: held.id.clone(),
2058 recovery: Some(Recovery::Review),
2059 ..Decision::default()
2060 },
2061 Decision {
2062 id: held.id.clone(),
2063 blocked_by: vec!["other-task".to_owned()],
2064 question: Some("retry now?".to_owned()),
2065 ..Decision::default()
2066 },
2067 ] {
2068 apply(
2069 &queue,
2070 &questions,
2071 &Verdict {
2072 decisions: vec![decision],
2073 },
2074 )
2075 .unwrap();
2076 }
2077
2078 let after = queue.get(&held.id).unwrap();
2079 assert_eq!(after.status, TaskStatus::Held);
2080 assert!(after.operator_held());
2081 assert_eq!(after.priority, 300);
2082 assert_eq!(after.runs, ["run20260912-224242-daf5"]);
2083 assert_eq!(
2084 after.hold_reason.as_deref(),
2085 Some("active manual recovery run20260912-224242-daf5")
2086 );
2087 assert!(after.blocked_by.is_empty());
2088 assert!(questions.list().is_empty());
2089 assert!(
2090 queue.next_runnable().is_none(),
2091 "must not dispatch a duplicate"
2092 );
2093 }
2094
2095 #[test]
2096 fn machine_holds_remain_recoverable_and_manual_release_is_authorization() {
2097 let dir = tempdir().unwrap();
2098 let queue = Queue::at(dir.path().join("queue"));
2099 let questions = Questions::at(dir.path().join("questions"));
2100
2101 let mut automatic = task("disk gate");
2102 automatic.hold_machine(Some("disk full".to_owned()));
2103 queue.put(&mut automatic).unwrap();
2104 let requeue = || Verdict {
2105 decisions: vec![Decision {
2106 id: automatic.id.clone(),
2107 recovery: Some(Recovery::Requeue),
2108 ..Decision::default()
2109 }],
2110 };
2111 apply(&queue, &questions, &requeue()).unwrap();
2112 assert_eq!(queue.get(&automatic.id).unwrap().status, TaskStatus::Queued);
2113
2114 let mut manual = task("operator gate");
2115 manual.hold_manual(Some("wait for operator".to_owned()));
2116 queue.put(&mut manual).unwrap();
2117 apply(
2118 &queue,
2119 &questions,
2120 &Verdict {
2121 decisions: vec![Decision {
2122 id: manual.id.clone(),
2123 recovery: Some(Recovery::Requeue),
2124 ..Decision::default()
2125 }],
2126 },
2127 )
2128 .unwrap();
2129 assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Held);
2130
2131 let mut released = queue.get(&manual.id).unwrap();
2134 released.release();
2135 queue.put(&mut released).unwrap();
2136 assert_eq!(queue.get(&manual.id).unwrap().status, TaskStatus::Queued);
2137 }
2138
2139 #[test]
2140 fn legacy_reasoned_hold_is_protected_without_losing_its_metadata() {
2141 let dir = tempdir().unwrap();
2142 let queue = Queue::at(dir.path().join("queue"));
2143 let questions = Questions::at(dir.path().join("questions"));
2144 let mut legacy = task("old explicit hold");
2145 legacy.status = TaskStatus::Held;
2146 legacy.hold_reason = Some("manual recovery already active".to_owned());
2147 legacy.hold_source = None;
2148 legacy.blocked_by = vec!["dependency".to_owned()];
2149 queue.put(&mut legacy).unwrap();
2150
2151 apply(
2152 &queue,
2153 &questions,
2154 &Verdict {
2155 decisions: vec![Decision {
2156 id: legacy.id.clone(),
2157 recovery: Some(Recovery::Requeue),
2158 ..Decision::default()
2159 }],
2160 },
2161 )
2162 .unwrap();
2163
2164 let after = queue.get(&legacy.id).unwrap();
2165 assert_eq!(after.status, TaskStatus::Held);
2166 assert_eq!(after.hold_source, None);
2167 assert_eq!(after.hold_reason, legacy.hold_reason);
2168 assert_eq!(after.blocked_by, legacy.blocked_by);
2169 }
2170
2171 #[test]
2172 fn review_recovery_is_a_no_op_without_a_survivable_branch() {
2173 crate::run::set_home(std::env::temp_dir().join("magi-conduct-tests-home"));
2179 let dir = tempdir().unwrap();
2180 let queue = Queue::at(dir.path().join("queue"));
2181 let questions = Questions::at(dir.path().join("questions"));
2182 let mut t = task("blocked with no readable run");
2183 t.start("20260101-000000-dead".to_owned()); t.fail("blocked", 5);
2185 queue.put(&mut t).unwrap();
2186
2187 apply(
2188 &queue,
2189 &questions,
2190 &Verdict {
2191 decisions: vec![Decision {
2192 id: t.id.clone(),
2193 recovery: Some(Recovery::Review),
2194 ..Decision::default()
2195 }],
2196 },
2197 )
2198 .unwrap();
2199
2200 let after = queue.get(&t.id).unwrap();
2201 assert_eq!(
2202 after.status,
2203 TaskStatus::Failed,
2204 "with nothing to reopen, the decision is dropped rather than guessed at"
2205 );
2206 assert!(after.review_branch.is_none());
2207 }
2208
2209 #[test]
2210 fn requeue_and_review_recovery_are_ignored_for_a_runnable_task() {
2211 let dir = tempdir().unwrap();
2215 let queue = Queue::at(dir.path().join("queue"));
2216 let questions = Questions::at(dir.path().join("questions"));
2217
2218 for recovery in [Recovery::Requeue, Recovery::Review] {
2219 let mut t = task("ordinary");
2220 queue.put(&mut t).unwrap();
2221
2222 apply(
2223 &queue,
2224 &questions,
2225 &Verdict {
2226 decisions: vec![Decision {
2227 id: t.id.clone(),
2228 recovery: Some(recovery),
2229 ..Decision::default()
2230 }],
2231 },
2232 )
2233 .unwrap();
2234
2235 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2236 }
2237 }
2238
2239 #[tokio::test]
2240 async fn a_broken_agent_leaves_the_queue_untouched_and_does_not_error() {
2241 let dir = tempdir().unwrap();
2242 let cfg = config(mock_agent(dir.path(), BROKEN, BTreeMap::new()));
2243 let queue = Queue::at(dir.path().join("queue"));
2244 let questions = Questions::at(dir.path().join("questions"));
2245 let mut t = task("normal");
2246 queue.put(&mut t).unwrap();
2247
2248 let mut conductor = Conductor::new();
2249 conductor
2250 .maybe_run(
2251 &cfg,
2252 dir.path(),
2253 &queue,
2254 &questions,
2255 dir.path(),
2256 &[t.clone()],
2257 &[],
2258 &[],
2259 2,
2260 )
2261 .await;
2262
2263 assert_eq!(
2264 queue.get(&t.id).unwrap().status,
2265 TaskStatus::Queued,
2266 "a failed invocation must change nothing"
2267 );
2268 assert!(
2269 queue.next_runnable().is_some(),
2270 "the loop must still be able to take the next task"
2271 );
2272 }
2273
2274 #[tokio::test]
2275 async fn a_reply_with_no_json_leaves_the_queue_untouched() {
2276 let dir = tempdir().unwrap();
2277 let cfg = config(mock_agent(dir.path(), GARBAGE, BTreeMap::new()));
2278 let queue = Queue::at(dir.path().join("queue"));
2279 let questions = Questions::at(dir.path().join("questions"));
2280 let mut t = task("normal");
2281 queue.put(&mut t).unwrap();
2282
2283 let mut conductor = Conductor::new();
2284 conductor
2285 .maybe_run(
2286 &cfg,
2287 dir.path(),
2288 &queue,
2289 &questions,
2290 dir.path(),
2291 &[t.clone()],
2292 &[],
2293 &[],
2294 2,
2295 )
2296 .await;
2297
2298 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2299 }
2300
2301 #[tokio::test]
2302 async fn a_conductor_chain_falls_to_the_next_agent_once_each() {
2303 let dir = tempdir().unwrap();
2304 let mut t = task("chained");
2305 let reply = format!(
2306 "{{\"decisions\":[{{\"id\":\"{}\",\"blocked_by\":[\"x\"],\"reason\":\"why\"}}]}}",
2307 t.id
2308 );
2309 let counted = |id: &str, tail: &str| {
2310 let calls = dir.path().join(format!("{id}.calls"));
2311 let path = dir.path().join(format!("{id}.sh"));
2312 std::fs::write(
2313 &path,
2314 format!(
2315 "#!/bin/sh\ncat >/dev/null\necho x >> '{}'\n{tail}\n",
2316 calls.display()
2317 ),
2318 )
2319 .unwrap();
2320 AgentSpec {
2321 id: id.to_owned(),
2322 kind: AgentKind::Command,
2323 model: None,
2324 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2325 extra_args: Vec::new(),
2326 env: BTreeMap::new(),
2327 prompt_delivery: None,
2328 }
2329 };
2330 let n_calls = |id: &str| {
2331 std::fs::read_to_string(dir.path().join(format!("{id}.calls")))
2332 .map_or(0, |s| s.lines().count())
2333 };
2334 let mut cfg = config(counted("a", "exit 3"));
2335 cfg.agents = vec![
2336 counted("a", "exit 3"),
2337 counted("b", &format!("printf '%s' '{reply}'")),
2338 ];
2339 cfg.roles.conductor = Some(crate::config::AgentChoice::Chain(vec![
2340 "a".into(),
2341 "b".into(),
2342 "a".into(),
2343 ]));
2344 let queue = Queue::at(dir.path().join("queue"));
2345 let questions = Questions::at(dir.path().join("questions"));
2346 queue.put(&mut t).unwrap();
2347
2348 let mut conductor = Conductor::new();
2349 conductor
2350 .maybe_run(
2351 &cfg,
2352 dir.path(),
2353 &queue,
2354 &questions,
2355 dir.path(),
2356 &[t.clone()],
2357 &[],
2358 &[],
2359 2,
2360 )
2361 .await;
2362
2363 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
2364 assert_eq!(n_calls("a"), 1, "each id is tried once");
2365 assert_eq!(n_calls("b"), 1);
2366 let saved = std::fs::read_to_string(seat_path(dir.path())).unwrap();
2367 assert!(
2368 saved.contains("\"b\""),
2369 "the seat on disk is the one that ran: {saved}"
2370 );
2371 }
2372
2373 #[tokio::test]
2374 async fn json_survives_code_fences_and_a_preamble() {
2375 let dir = tempdir().unwrap();
2376 let mut t = task("fenced");
2377 let reply = format!(
2378 "Sure, here is my decision.\n\n```json\n{{\"decisions\":[{{\"id\":\"{}\",\
2379 \"blocked_by\":[\"x\"],\"reason\":\"why\"}}]}}\n```\n",
2380 t.id
2381 );
2382 let cfg = config(mock_agent(dir.path(), REPLY, env(&reply)));
2383 let queue = Queue::at(dir.path().join("queue"));
2384 let questions = Questions::at(dir.path().join("questions"));
2385 queue.put(&mut t).unwrap();
2386
2387 let mut conductor = Conductor::new();
2388 conductor
2389 .maybe_run(
2390 &cfg,
2391 dir.path(),
2392 &queue,
2393 &questions,
2394 dir.path(),
2395 &[t.clone()],
2396 &[],
2397 &[],
2398 2,
2399 )
2400 .await;
2401
2402 let back = queue.get(&t.id).unwrap();
2403 assert_eq!(back.status, TaskStatus::Blocked);
2404 assert_eq!(back.blocked_by, ["x"]);
2405 }
2406
2407 #[tokio::test]
2408 async fn the_conductor_is_not_called_again_when_nothing_worth_looking_at_has_changed() {
2409 let dir = tempdir().unwrap();
2412 let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2413 let queue = Queue::at(dir.path().join("queue"));
2414 let questions = Questions::at(dir.path().join("questions"));
2415 let mut t = task("stable");
2416 queue.put(&mut t).unwrap();
2417 let artifacts = dir.path().join("conduct").join("artifacts");
2418 let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2419
2420 let mut conductor = Conductor::new();
2421 conductor
2422 .maybe_run(
2423 &cfg,
2424 dir.path(),
2425 &queue,
2426 &questions,
2427 dir.path(),
2428 &[t.clone()],
2429 &[],
2430 &[],
2431 2,
2432 )
2433 .await;
2434 assert!(turn(1).is_file(), "the first cycle must call the conductor");
2435
2436 conductor
2437 .maybe_run(
2438 &cfg,
2439 dir.path(),
2440 &queue,
2441 &questions,
2442 dir.path(),
2443 &[t.clone()],
2444 &[],
2445 &[],
2446 2,
2447 )
2448 .await;
2449 assert!(
2450 !turn(2).is_file(),
2451 "an unchanged revision and an unchanged stalled/finished set must not call the \
2452 conductor twice"
2453 );
2454
2455 t.priority = 1;
2457 queue.put(&mut t).unwrap();
2458 conductor
2459 .maybe_run(
2460 &cfg,
2461 dir.path(),
2462 &queue,
2463 &questions,
2464 dir.path(),
2465 &[t.clone()],
2466 &[],
2467 &[],
2468 2,
2469 )
2470 .await;
2471 assert!(turn(2).is_file(), "a moved revision calls it again");
2472 }
2473
2474 #[tokio::test]
2475 async fn a_task_turning_stalled_calls_the_conductor_again_despite_an_unchanged_revision() {
2476 let dir = tempdir().unwrap();
2482 let cfg = config(mock_agent(dir.path(), REPLY, env("{\"decisions\":[]}")));
2483 let queue = Queue::at(dir.path().join("queue"));
2484 let questions = Questions::at(dir.path().join("questions"));
2485 let mut t = task("quiet");
2486 queue.put(&mut t).unwrap();
2487 let artifacts = dir.path().join("conduct").join("artifacts");
2488 let turn = |n: usize| artifacts.join(format!("turn-{n}.out"));
2489
2490 let mut conductor = Conductor::new();
2491 conductor
2492 .maybe_run(
2493 &cfg,
2494 dir.path(),
2495 &queue,
2496 &questions,
2497 dir.path(),
2498 &[t.clone()],
2499 &[],
2500 &[],
2501 2,
2502 )
2503 .await;
2504 assert!(turn(1).is_file());
2505
2506 conductor
2507 .maybe_run(
2508 &cfg,
2509 dir.path(),
2510 &queue,
2511 &questions,
2512 dir.path(),
2513 &[],
2514 &[t.clone()],
2515 &[],
2516 2,
2517 )
2518 .await;
2519 assert!(
2520 turn(2).is_file(),
2521 "a task turning stalled must call the conductor again"
2522 );
2523
2524 conductor
2527 .maybe_run(
2528 &cfg,
2529 dir.path(),
2530 &queue,
2531 &questions,
2532 dir.path(),
2533 &[],
2534 &[t.clone()],
2535 &[],
2536 2,
2537 )
2538 .await;
2539 assert!(
2540 !turn(3).is_file(),
2541 "the same stalled task lingering must not call the conductor every cycle"
2542 );
2543 }
2544
2545 fn rehold(queue: &Queue, questions: &Questions, id: &str, note: &str) {
2546 apply(
2547 queue,
2548 questions,
2549 &Verdict {
2550 decisions: vec![Decision {
2551 id: id.to_owned(),
2552 reason: Some(note.to_owned()),
2553 recovery: Some(Recovery::Hold),
2554 ..Decision::default()
2555 }],
2556 },
2557 )
2558 .unwrap();
2559 }
2560
2561 #[test]
2562 fn reholding_with_an_identical_note_changes_nothing() {
2563 let dir = tempdir().unwrap();
2564 let queue = Queue::at(dir.path().join("queue"));
2565 let questions = Questions::at(dir.path().join("questions"));
2566 let mut t = task("held");
2567 t.hold_machine(Some("waiting on a human".to_owned()));
2568 queue.put(&mut t).unwrap();
2569
2570 rehold(&queue, &questions, &t.id, "owner said keep it held");
2571 let first = queue.get(&t.id).unwrap();
2572 let rev = queue.revision();
2573 for _ in 0..5 {
2574 rehold(&queue, &questions, &t.id, "owner said keep it held");
2575 }
2576 let after = queue.get(&t.id).unwrap();
2577 assert_eq!(after.hold_reason, first.hold_reason);
2578 assert_eq!(after.updated_at, first.updated_at);
2579 assert_eq!(queue.revision(), rev, "nothing was written");
2580 }
2581
2582 #[test]
2583 fn reholding_with_different_notes_keeps_one_previously_level() {
2584 let dir = tempdir().unwrap();
2585 let queue = Queue::at(dir.path().join("queue"));
2586 let questions = Questions::at(dir.path().join("questions"));
2587 let mut t = task("held");
2588 t.hold_machine(Some("original".to_owned()));
2589 queue.put(&mut t).unwrap();
2590
2591 for i in 0..20 {
2592 rehold(&queue, &questions, &t.id, &format!("note {}", i % 2));
2593 }
2594 let reason = queue.get(&t.id).unwrap().hold_reason.unwrap();
2595 assert_eq!(reason.matches("(previously:").count(), 1, "{reason}");
2596 assert!(reason.starts_with("note "), "the new note leads: {reason}");
2597 }
2598
2599 #[test]
2600 fn an_already_nested_reason_is_matched_by_its_outermost_note() {
2601 let mut t = task("deep");
2602 let nested = "same\n\n(previously: same\n\n(previously: same))";
2603 t.hold_machine(Some(nested.to_owned()));
2604 let d = Decision {
2605 id: t.id.clone(),
2606 reason: Some("same".to_owned()),
2607 ..Decision::default()
2608 };
2609 assert_eq!(reaffirmed_hold_reason(&t, &d), None);
2610 let d2 = Decision {
2611 reason: Some("other".to_owned()),
2612 ..d
2613 };
2614 assert_eq!(
2615 reaffirmed_hold_reason(&t, &d2).as_deref(),
2616 Some("other\n\n(previously: same)")
2617 );
2618 }
2619
2620 #[tokio::test]
2621 async fn a_repeated_hold_decision_settles_instead_of_looping() {
2622 let dir = tempdir().unwrap();
2623 let queue = Queue::at(dir.path().join("queue"));
2624 let questions = Questions::at(dir.path().join("questions"));
2625 let reply = r#"{"decisions":[{"id":"ID","reason":"keep held","recovery":"hold"}]}"#;
2626 let mut t = task("held");
2627 t.hold_machine(Some("first".to_owned()));
2628 queue.put(&mut t).unwrap();
2629 let cfg = config(mock_agent(
2630 dir.path(),
2631 REPLY,
2632 env(&reply.replace("ID", &t.id)),
2633 ));
2634 let mut conductor = Conductor::new();
2635 for _ in 0..2 {
2637 let held = queue.get(&t.id).unwrap();
2638 conductor
2639 .maybe_run(
2640 &cfg,
2641 dir.path(),
2642 &queue,
2643 &questions,
2644 dir.path(),
2645 &[],
2646 &[],
2647 &[held],
2648 2,
2649 )
2650 .await;
2651 }
2652 let held = queue.get(&t.id).unwrap();
2653 assert!(
2654 held.hold_reason
2655 .as_deref()
2656 .unwrap()
2657 .starts_with("keep held")
2658 );
2659 assert!(
2660 !conductor.worth_a_look(&queue, &[], &[held]),
2661 "an unchanged held task must stop being reconsidered"
2662 );
2663 }
2664
2665 #[tokio::test]
2666 async fn alternating_wording_does_not_keep_a_held_task_in_the_loop() {
2667 let dir = tempdir().unwrap();
2668 let queue = Queue::at(dir.path().join("queue"));
2669 let questions = Questions::at(dir.path().join("questions"));
2670 let mut t = task("held");
2671 t.hold_machine(Some("first".to_owned()));
2672 queue.put(&mut t).unwrap();
2673 let reply = |note: &str| {
2674 format!(
2675 r#"{{"decisions":[{{"id":"{}","reason":"{note}","recovery":"hold"}}]}}"#,
2676 t.id
2677 )
2678 };
2679 let cfg_a = config(mock_agent(dir.path(), REPLY, env(&reply("keep held"))));
2680 let cfg_b = config(mock_agent(dir.path(), REPLY, env(&reply("leave held"))));
2681 let mut conductor = Conductor::new();
2682 let mut writes = 0;
2683 for i in 0..6 {
2684 let held = queue.get(&t.id).unwrap();
2685 let before = queue.revision();
2686 conductor
2687 .maybe_run(
2688 if i % 2 == 0 { &cfg_a } else { &cfg_b },
2689 dir.path(),
2690 &queue,
2691 &questions,
2692 dir.path(),
2693 &[],
2694 &[],
2695 &[held],
2696 2,
2697 )
2698 .await;
2699 if queue.revision() != before {
2700 writes += 1;
2701 }
2702 }
2703 assert!(writes <= 2, "the loop must settle, saw {writes} writes");
2704 let mut held = queue.get(&t.id).unwrap();
2705 held.instruction.push_str(" (edited)");
2706 queue.put(&mut held).unwrap();
2707 let held = queue.get(&t.id).unwrap();
2708 assert!(conductor.worth_a_look(&queue, &[], &[held]));
2709 }
2710
2711 #[test]
2712 fn worth_a_look_is_config_free_and_matches_maybe_runs_own_gate() {
2713 let dir = tempdir().unwrap();
2714 let queue = Queue::at(dir.path().join("queue"));
2715 let mut t = task("t");
2716 queue.put(&mut t).unwrap();
2717
2718 let mut conductor = Conductor::new();
2719 assert!(
2720 conductor.worth_a_look(&queue, &[], &[]),
2721 "a conductor that has never run has something to look at"
2722 );
2723
2724 conductor.last_seen = Some(conductor.snapshot(&queue, &[], &[]));
2725 assert!(
2726 !conductor.worth_a_look(&queue, &[], &[]),
2727 "nothing changed and nothing is stalled or finished"
2728 );
2729 assert!(
2730 conductor.worth_a_look(&queue, &[t.clone()], &[]),
2731 "a stalled task is worth a look even at the same revision"
2732 );
2733 assert!(
2734 conductor.worth_a_look(&queue, &[], &[t.clone()]),
2735 "a finished task is worth a look even at the same revision"
2736 );
2737 }
2738
2739 #[tokio::test]
2740 async fn the_conduct_path_never_calls_ask_and_wait() {
2741 let dir = tempdir().unwrap();
2747 let queue = Queue::at(dir.path().join("queue"));
2748 let questions = Questions::at(dir.path().join("questions"));
2749 let mut t = task("asks without blocking");
2750 queue.put(&mut t).unwrap();
2751
2752 apply(
2753 &queue,
2754 &questions,
2755 &Verdict {
2756 decisions: vec![Decision {
2757 id: t.id.clone(),
2758 question: Some("ok?".to_owned()),
2759 ..Decision::default()
2760 }],
2761 },
2762 )
2763 .unwrap();
2764 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Blocked);
2766 }
2767
2768 #[test]
2769 fn a_pinned_resume_is_not_held_or_requeued_by_the_conductor() {
2770 let mut t = Task::new("t".into(), "t".into(), PathBuf::new(), Source::Human);
2771 t.resume_override = Some(crate::queue::OperatorResume {
2772 question_id: "q".into(),
2773 at: jiff::Timestamp::now(),
2774 conductor_rehold: None,
2775 forced: false,
2776 pinned_run: Some("run-1".into()),
2777 });
2778 assert!(pinned_resume(&t));
2779 assert!(!may_hold(&mut t, "waiting for magi resume"));
2780 assert!(
2781 t.resume_override
2782 .as_ref()
2783 .unwrap()
2784 .conductor_rehold
2785 .is_none(),
2786 "a refused hold is not recorded as an override"
2787 );
2788 t.resume_override = None;
2789 assert!(!pinned_resume(&t));
2790 assert!(may_hold(&mut t, "no override, so a hold is allowed"));
2791 }
2792
2793 #[test]
2794 fn a_pinned_resume_is_not_blocked_by_a_conductor_question_or_dependency() {
2795 let dir = tempdir().unwrap();
2796 let queue = Queue::at(dir.path().join("queue"));
2797 let questions = Questions::at(dir.path().join("questions"));
2798 let mut t = task("resume me");
2799 t.resume_override = Some(crate::queue::OperatorResume {
2800 question_id: "q".into(),
2801 at: jiff::Timestamp::now(),
2802 conductor_rehold: None,
2803 forced: true,
2804 pinned_run: Some("run-1".into()),
2805 });
2806 queue.put(&mut t).unwrap();
2807
2808 apply(
2809 &queue,
2810 &questions,
2811 &Verdict {
2812 decisions: vec![
2813 Decision {
2814 id: t.id.clone(),
2815 question: Some("really?".to_owned()),
2816 ..Decision::default()
2817 },
2818 Decision {
2819 id: t.id.clone(),
2820 blocked_by: vec!["other".to_owned()],
2821 ..Decision::default()
2822 },
2823 ],
2824 },
2825 )
2826 .unwrap();
2827
2828 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
2829 assert!(questions.list().is_empty());
2830 }
2831}