1use std::path::{Path, PathBuf};
50use std::sync::Arc;
51use std::sync::atomic::{AtomicBool, Ordering};
52use std::sync::{Mutex, MutexGuard};
53use std::time::Duration;
54
55use anyhow::{Context, Result, bail};
56use jiff::Timestamp;
57use serde::{Deserialize, Serialize};
58use tokio::sync::Notify;
59
60use crate::ask::{self, Questions};
61use crate::clean;
62use crate::conduct::Conductor;
63use crate::config::{Config, MergeMode};
64use crate::graph::Runner;
65use crate::land;
66use crate::notices::{self, Link, Notice};
67use crate::queue::{Queue, Task, TaskStatus};
68use crate::run::{Liveness, QuotaLoss, RunState, RunStatus};
69use crate::triage;
70
71pub const SCHEMA: u32 = 1;
73
74pub const HEARTBEAT: Duration = Duration::from_secs(5);
78
79pub const STALE_SECS: i64 = 30;
88
89pub const POLL: Duration = Duration::from_secs(5);
91
92pub const STALE_CLAIM: Duration = Duration::from_secs(6 * 60 * 60);
96
97pub const STALLED_RUNNING: Duration = Duration::from_secs(30 * 60);
116
117#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
119#[serde(default)]
120pub struct Current {
121 pub task: String,
123 pub run: String,
125}
126
127#[derive(Debug, Clone, Serialize, Deserialize)]
134pub struct Status {
135 pub schema: u32,
137 pub pid: u32,
139 pub started_at: Timestamp,
141 pub updated_at: Timestamp,
143 pub idle: bool,
145 pub current: Vec<Current>,
151 pub completed: usize,
153 pub polls: u64,
155}
156
157impl Status {
158 #[must_use]
160 pub fn new() -> Self {
161 let now = Timestamp::now();
162 Self {
163 schema: SCHEMA,
164 pid: std::process::id(),
165 started_at: now,
166 updated_at: now,
167 idle: true,
168 current: Vec::new(),
169 completed: 0,
170 polls: 0,
171 }
172 }
173}
174
175impl Default for Status {
176 fn default() -> Self {
177 Self::new()
178 }
179}
180
181#[derive(Debug, Clone)]
183pub struct Opts {
184 pub repo: PathBuf,
186 pub config: Option<PathBuf>,
188 pub poll: Duration,
190 pub max_attempts: usize,
192 pub once: bool,
194 pub merge: Option<String>,
196 pub worktrees_root: Option<PathBuf>,
205}
206
207impl Default for Opts {
208 fn default() -> Self {
209 Self {
210 repo: PathBuf::from("."),
211 config: None,
212 poll: POLL,
213 max_attempts: 2,
214 once: false,
215 merge: None,
216 worktrees_root: None,
217 }
218 }
219}
220
221fn max_concurrent(n: usize) -> usize {
226 n.max(1)
227}
228
229#[must_use]
231pub fn status_path() -> PathBuf {
232 crate::run::home().join("daemon.json")
233}
234
235pub fn write_status(status: &Status) -> Result<()> {
237 write_status_to(&status_path(), status)
238}
239
240pub fn write_status_to(path: &Path, status: &Status) -> Result<()> {
245 if let Some(parent) = path.parent() {
246 std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
247 }
248 let body = serde_json::to_string_pretty(status).context("serialize daemon status")?;
249 let tmp = path.with_extension("json.tmp");
250 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
251 std::fs::rename(&tmp, path).with_context(|| format!("replace {}", path.display()))?;
252 Ok(())
253}
254
255pub fn clear_status() {
258 clear_status_at(&status_path());
259}
260
261fn clear_status_at(path: &Path) {
265 let _ = std::fs::remove_file(path);
266}
267
268#[derive(Debug, Clone, Default)]
281pub struct Stop {
282 stopped: Arc<AtomicBool>,
286 busy: Arc<std::sync::atomic::AtomicUsize>,
292 wake: Arc<Notify>,
296 pause: crate::graph::Pause,
299}
300
301impl Stop {
302 #[must_use]
304 pub fn new() -> Self {
305 Self::default()
306 }
307
308 pub fn stop(&self) {
311 self.stopped.store(true, Ordering::SeqCst);
312 self.wake.notify_one();
316 }
317
318 #[must_use]
320 pub fn stopped(&self) -> bool {
321 self.stopped.load(Ordering::SeqCst)
322 }
323
324 #[must_use]
332 pub fn finishing(&self) -> bool {
333 self.stopped() && self.busy_now()
334 }
335
336 pub fn park(&self) {
347 self.pause.park();
348 self.stop();
349 }
350
351 #[must_use]
353 pub fn parking(&self) -> bool {
354 self.pause.parked()
355 }
356
357 #[must_use]
359 pub fn pause(&self) -> crate::graph::Pause {
360 self.pause.clone()
361 }
362
363 #[must_use]
369 pub fn busy_now(&self) -> bool {
370 self.busy.load(Ordering::SeqCst) > 0
371 }
372
373 fn enter(&self) {
375 self.busy.fetch_add(1, Ordering::SeqCst);
376 }
377
378 fn exit(&self) {
381 self.busy.fetch_sub(1, Ordering::SeqCst);
382 }
383
384 async fn idle(&self, poll: Duration) {
386 tokio::select! {
387 () = tokio::time::sleep(poll) => {}
388 () = self.wake.notified() => {}
389 }
390 }
391}
392
393#[derive(Debug, Clone, Default, Deserialize)]
400#[serde(default)]
401pub struct Reading {
402 pub schema: u32,
404 pub pid: Option<u32>,
406 pub started_at: Option<Timestamp>,
408 pub updated_at: Option<Timestamp>,
410 pub idle: bool,
412 #[serde(deserialize_with = "de_current")]
424 pub current: Vec<Current>,
425 pub completed: u64,
427 pub polls: u64,
429}
430
431fn de_current<'de, D>(deserializer: D) -> std::result::Result<Vec<Current>, D::Error>
434where
435 D: serde::Deserializer<'de>,
436{
437 #[derive(Deserialize)]
438 #[serde(untagged)]
439 enum Shape {
440 Many(Vec<Current>),
441 One(Current),
442 }
443 Ok(
444 Option::<Shape>::deserialize(deserializer)?.map_or_else(Vec::new, |shape| match shape {
445 Shape::Many(v) => v,
446 Shape::One(c) => vec![c],
447 }),
448 )
449}
450
451impl Reading {
452 #[must_use]
455 pub fn age_secs(&self, now: Timestamp) -> Option<i64> {
456 self.updated_at
457 .map(|at| (now.as_second() - at.as_second()).max(0))
458 }
459
460 #[must_use]
464 pub fn running(&self, now: Timestamp) -> bool {
465 self.age_secs(now).is_some_and(|secs| secs <= STALE_SECS)
466 }
467}
468
469#[must_use]
476pub fn read_status(home: &Path) -> Option<Reading> {
477 let body = std::fs::read_to_string(home.join("daemon.json")).ok()?;
478 serde_json::from_str(&body).ok()
479}
480
481#[must_use]
493pub fn current_work(home: &Path, now: Timestamp) -> Vec<Current> {
494 read_status(home)
495 .filter(|reading| reading.running(now))
496 .map(|reading| reading.current)
497 .unwrap_or_default()
498}
499
500#[must_use]
502pub fn is_working_on(home: &Path, run: &str, now: Timestamp) -> bool {
503 current_work(home, now).iter().any(|c| c.run == run)
504}
505
506#[must_use]
516pub fn is_working_on_short(home: &Path, short: &str, now: Timestamp) -> bool {
517 current_work(home, now)
518 .iter()
519 .any(|c| crate::run::short_of(&c.run) == short)
520}
521
522#[must_use]
524pub fn is_working_on_task(home: &Path, task: &str, now: Timestamp) -> bool {
525 current_work(home, now).iter().any(|c| c.task == task)
526}
527
528pub fn sweep_stale_claims(queue: &Queue, older_than: Duration) -> Vec<String> {
572 sweep_stale_claims_with(queue, older_than, crate::proc::pid_alive)
573}
574
575fn sweep_stale_claims_with<F>(queue: &Queue, older_than: Duration, pid_alive: F) -> Vec<String>
579where
580 F: Fn(u32) -> bool,
581{
582 let this_process = std::process::id();
583 let mut swept: Vec<String> = std::fs::read_dir(queue.root())
584 .into_iter()
585 .flatten()
586 .flatten()
587 .map(|e| e.path())
588 .filter(|p| p.extension().is_some_and(|x| x == "lock"))
589 .filter(|p| {
590 match std::fs::read_to_string(p)
591 .ok()
592 .and_then(|body| body.trim().parse::<u32>().ok())
593 {
594 Some(pid) if pid == this_process => false,
598 Some(pid) => !pid_alive(pid),
599 None => p
600 .metadata()
601 .and_then(|m| m.modified())
602 .and_then(|t| t.elapsed().map_err(std::io::Error::other))
603 .is_ok_and(|age| age >= older_than),
604 }
605 })
606 .filter(|p| std::fs::remove_file(p).is_ok())
607 .filter_map(|p| {
608 p.file_stem()
609 .and_then(|s| s.to_str())
610 .map(std::borrow::ToOwned::to_owned)
611 })
612 .collect();
613 swept.sort_unstable();
614 swept
615}
616
617fn is_stalled(task: &Task, home: &Path, now: Timestamp) -> bool {
622 task.status == TaskStatus::Running
623 && (now.as_second() - task.updated_at.as_second()) >= STALLED_RUNNING.as_secs() as i64
624 && !is_working_on_task(home, &task.id, now)
625}
626
627fn stalled_tasks(queue: &Queue, home: &Path, now: Timestamp) -> Vec<Task> {
630 queue
631 .list()
632 .into_iter()
633 .filter(|t| is_stalled(t, home, now))
634 .collect()
635}
636
637fn queued_tasks(queue: &Queue) -> Vec<Task> {
643 queue
644 .list()
645 .into_iter()
646 .filter(|t| t.status == TaskStatus::Queued)
647 .collect()
648}
649
650fn finished_tasks(queue: &Queue) -> Vec<Task> {
653 queue
654 .list()
655 .into_iter()
656 .filter(|t| matches!(t.status, TaskStatus::Failed | TaskStatus::Held))
657 .collect()
658}
659
660fn resolve_blockers(queue: &Queue, questions: &Questions) {
677 for listed in queue.list() {
678 if listed.status != TaskStatus::Blocked || listed.blocked_by.is_empty() {
679 continue;
680 }
681 let Ok(_claim) = queue.claim(&listed.id) else {
682 continue;
683 };
684 let Ok(mut task) = queue.get(&listed.id) else {
685 continue;
686 };
687 if task.status != TaskStatus::Blocked {
688 continue;
689 }
690 let missing = crate::queue::missing_blockers(queue, questions, &task.blocked_by);
691 if !missing.is_empty() {
692 let language = language_of(&task, Path::new("."));
693 task.hold_machine(Some(crate::queue::missing_blocker_hold_reason_in(
694 &task.blocked_by,
695 &missing,
696 &language,
697 )));
698 record(queue, &mut task);
699 continue;
700 }
701 if let Some(q) = task.blocked_by.iter().find_map(|id| {
704 questions.get(id).ok().filter(|q| {
705 q.node == crate::conduct::NODE && q.status == ask::QuestionStatus::Abandoned
706 })
707 }) {
708 let language = language_of(&task, Path::new("."));
709 task.hold_machine(Some(unanswered_question_hold_reason(&q, &language)));
710 record(queue, &mut task);
711 continue;
712 }
713 let mut changed = false;
714 for id in task.blocked_by.clone() {
715 if let Ok(dep) = queue.get(&id) {
716 if dep.status == TaskStatus::Done {
717 task.unblock(&id);
718 changed = true;
719 }
720 continue;
721 }
722 if let Ok(q) = questions.get(&id)
723 && q.status == ask::QuestionStatus::Answered
724 {
725 let answer = match &q.answer {
726 Some(ask::Answer::Choice(c) | ask::Answer::Text(c)) => c.clone(),
727 None => String::new(),
728 };
729 task.record_answer(q.summary.clone(), answer);
730 task.unblock(&id);
731 changed = true;
732 }
733 }
734 if changed {
735 record(queue, &mut task);
736 }
737 }
738}
739
740fn unanswered_question_hold_reason(q: &ask::Question, language: &str) -> String {
742 if crate::lang::is_japanese(language) {
743 format!(
744 "質問 {} 「{}」 に期限内の回答がなく、取り下げられました - `magi task triage` を参照",
745 q.short(),
746 q.summary
747 )
748 } else {
749 format!(
750 "question {} \"{}\" went unanswered and was abandoned - see `magi task triage`",
751 q.short(),
752 q.summary
753 )
754 }
755}
756
757#[derive(Debug, Clone, PartialEq, Eq)]
759enum ActionDecision {
760 Skip,
762 Resume(String),
764 Requeue,
766 Done,
768 Stale,
771 Refuse(String),
774}
775
776fn decide_action<F>(task: &Task, q: &ask::Question, p: &Phrases, load: F) -> ActionDecision
784where
785 F: FnOnce(&str) -> Result<RunState>,
786{
787 let Some(action) = q.chosen_action() else {
788 return ActionDecision::Skip;
789 };
790 if task.action_applied(&q.id)
791 || matches!(
792 task.status,
793 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
794 )
795 {
796 return ActionDecision::Skip;
797 }
798 if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
801 return ActionDecision::Stale;
802 }
803 match action {
804 ask::ChoiceAction::Requeue => ActionDecision::Requeue,
805 ask::ChoiceAction::Done => ActionDecision::Done,
806 ask::ChoiceAction::Resume { run } => {
807 if task.runs.last() != Some(run) {
808 return ActionDecision::Refuse((p.resume_not_latest)(
809 q.short(),
810 ask::short_id(run),
811 ));
812 }
813 match load(run) {
814 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
815 ActionDecision::Resume(run.clone())
816 }
817 Ok(_) => ActionDecision::Refuse((p.resume_cannot_progress)(
818 q.short(),
819 ask::short_id(run),
820 )),
821 Err(e) => ActionDecision::Refuse((p.resume_unreadable)(
822 q.short(),
823 ask::short_id(run),
824 &format!("{e:#}"),
825 )),
826 }
827 }
828 }
829}
830
831struct Phrases {
845 graph_stopped: fn(&str, &str) -> String,
847 quorum_lost: &'static str,
848 quota_took_out: &'static str,
850 run_ended: &'static str,
852 waiting_for_answer: &'static str,
854 recovered_running: &'static str,
855 no_run_to_recover: &'static str,
856 could_not_start: &'static str,
857 resume_not_latest: fn(&str, &str) -> String,
858 resume_cannot_progress: fn(&str, &str) -> String,
859 resume_unreadable: fn(&str, &str, &str) -> String,
860 handover_refused: fn(&crate::handover::Refused) -> String,
863 handover_hint: &'static str,
865 reviewers_never_answered: fn(usize, usize) -> String,
868}
869
870const PHRASES_EN: Phrases = Phrases {
871 graph_stopped: |status, detail| {
872 format!("the graph stopped at `{status}` without reaching a terminal status: {detail}")
873 },
874 quorum_lost: "the judging panel lost its quorum",
875 quota_took_out: "; quota took out ",
876 run_ended: "run ended ",
877 waiting_for_answer: " - waiting for operator answer to question ",
878 recovered_running: "recovered a `running` task whose daemon never recorded the outcome: ",
879 no_run_to_recover: "task was `running` with no live daemon and no readable \
880 run to recover; held for a human to check what happened",
881 could_not_start: "could not start the run: ",
882 resume_not_latest: |q, run| {
883 format!("question {q} asked to resume run {run}, which is not this task's latest run")
884 },
885 resume_cannot_progress: |q, run| {
886 format!("question {q} asked to resume run {run}, which cannot make progress")
887 },
888 resume_unreadable: |q, run, e| {
889 format!("question {q} asked to resume run {run}, which could not be read: {e}")
890 },
891 handover_refused: |r| r.to_string(),
892 handover_hint: " (clean up the other worktree, then release the task from the queue)",
893 reviewers_never_answered: |missing, rounds| {
894 format!(
895 "{missing} reviewer seat(s) never answered after {rounds} rounds; \
896 refusing to call it clean"
897 )
898 },
899};
900
901const PHRASES_JA: Phrases = Phrases {
902 graph_stopped: |status, detail| {
903 format!("グラフが終端状態に達しないまま `{status}` で停止しました: {detail}")
904 },
905 quorum_lost: "審査パネルが定足数を失いました",
906 quota_took_out: "。クォータで脱落: ",
907 run_ended: "run 終了: ",
908 waiting_for_answer: " - オペレーターの回答待ち: 質問 ",
909 recovered_running: "daemon が結果を記録しないまま `running` だったタスクを回収しました: ",
910 no_run_to_recover: "タスクは `running` でしたが、生きた daemon も回収できる run も見つかりません。\
911 何が起きたか人が確認するため保留にしました",
912 could_not_start: "run を開始できませんでした: ",
913 resume_not_latest: |q, run| {
914 format!(
915 "質問 {q} は run {run} の再開を求めましたが、これはタスクの最新の run ではありません"
916 )
917 },
918 resume_cannot_progress: |q, run| {
919 format!("質問 {q} は run {run} の再開を求めましたが、これは進行できません")
920 },
921 resume_unreadable: |q, run, e| {
922 format!("質問 {q} は run {run} の再開を求めましたが、読み込めませんでした: {e}")
923 },
924 handover_refused: |r| {
925 use crate::handover::Refused;
926 match r {
927 Refused::Foreign { branch, path, why } => format!(
928 "ブランチ `{branch}` は {path} にチェックアウトされており、magi は自動では\
929 削除しません。不要なら `git worktree remove` でその worktree を削除して\
930 から、やり直してください(詳細: {why})"
931 ),
932 Refused::Unsafe { branch, path, why } => format!(
933 "ブランチ `{branch}` は {path} にチェックアウトされています。そこの作業を\
934 コミットか破棄したうえで `git worktree remove` で worktree を削除する\
935 か、run を破棄してよいと伝えてから、やり直してください(詳細: {why})"
936 ),
937 Refused::ReleaseFailed { branch, path, run } => format!(
938 "ブランチ `{branch}` は {path} で run {run} がチェックアウトしており、その\
939 worktree の解放に失敗したか、変更が見つかりました(未コミットの変更が\
940 ある worktree は git が削除を拒否します)。worktree はそのまま残しました"
941 ),
942 }
943 },
944 handover_hint: "(他の worktree を片付けてから、タスクをキューから解放してください)",
945 reviewers_never_answered: |missing, rounds| {
946 format!(
947 "{missing} 席のレビュアーが {rounds} ラウンドの間に一度も回答しなかったため、\
948 クリーンとは認めません"
949 )
950 },
951};
952
953fn phrases(language: &str) -> &'static Phrases {
954 if crate::lang::is_japanese(language) {
955 &PHRASES_JA
956 } else {
957 &PHRASES_EN
958 }
959}
960
961fn language_of(task: &Task, fallback: &Path) -> String {
965 crate::lang::of_repo(&repo_for(task, fallback))
966}
967
968fn task_of_question<'a>(tasks: &'a [Task], q: &ask::Question) -> Option<&'a Task> {
972 if q.node == crate::conduct::NODE {
973 return tasks.iter().find(|t| t.id == q.run);
974 }
975 tasks.iter().find(|t| t.runs.contains(&q.run))
976}
977
978fn apply_choice_actions(queue: &Queue, questions: &Questions, home: &Path) {
987 let tasks = queue.list();
988 for q in questions.list() {
989 if q.chosen_action().is_none() {
990 continue;
991 }
992 let Some(listed) = task_of_question(&tasks, &q) else {
993 continue;
994 };
995 if listed.action_applied(&q.id) {
996 continue;
997 }
998 let Ok(_claim) = queue.claim(&listed.id) else {
999 continue;
1000 };
1001 let Ok(mut task) = queue.get(&listed.id) else {
1002 continue;
1003 };
1004 let language = language_of(&task, Path::new("."));
1005 let decision = decide_action(&task, &q, phrases(&language), |id| {
1006 RunState::load_under(id, home)
1007 });
1008 let ran = matches!(
1009 decision,
1010 ActionDecision::Resume(_) | ActionDecision::Requeue | ActionDecision::Done
1011 );
1012 if ran {
1013 if questions
1018 .read_lease(&q.id)
1019 .is_some_and(|l| l.fresh(Timestamp::now()))
1020 {
1021 continue;
1022 }
1023 let taken = questions.update(&q.id, |r| {
1024 let free = !r.answer_delivered;
1025 r.answer_delivered = true;
1026 Ok(free)
1027 });
1028 if !matches!(taken, Ok((_, true))) {
1029 continue;
1030 }
1031 }
1032 match decision {
1033 ActionDecision::Skip => continue,
1034 ActionDecision::Resume(run) => {
1035 task.release();
1036 task.resume_override = Some(crate::queue::OperatorResume {
1037 question_id: q.id.clone(),
1038 at: Timestamp::now(),
1039 conductor_rehold: None,
1040 forced: true,
1041 pinned_run: Some(run),
1042 });
1043 }
1044 ActionDecision::Stale => {}
1045 ActionDecision::Requeue => task.requeue(),
1046 ActionDecision::Done => {
1047 task.succeed();
1048 supersede_prior_runs(&task, home);
1049 }
1050 ActionDecision::Refuse(why) => {
1051 task.hold_machine(Some(why));
1052 notices::raise(
1053 Notice::warn(
1054 &format!("action:{}", q.id),
1055 "An answer asked the daemon to resume a run that cannot be resumed; the task stays held.",
1056 )
1057 .link(Link::Task {
1058 id: task.id.clone(),
1059 }),
1060 );
1061 }
1062 }
1063 task.mark_action_applied(&q.id);
1064 record(queue, &mut task);
1065 }
1066}
1067
1068fn reconcile_task_questions(queue: &Queue, questions: &Questions) {
1079 let tasks = queue.list();
1080 let by_id: std::collections::BTreeMap<&str, &Task> =
1081 tasks.iter().map(|t| (t.id.as_str(), t)).collect();
1082 let referenced: std::collections::BTreeSet<&str> = tasks
1083 .iter()
1084 .flat_map(|task| task.blocked_by.iter().map(String::as_str))
1085 .collect();
1086
1087 for mut question in questions.list() {
1088 if !question.status.open() || question.node != crate::conduct::NODE {
1089 continue;
1090 }
1091 if referenced.contains(question.id.as_str()) {
1094 continue;
1095 }
1096 let Some(task) = by_id.get(question.run.as_str()) else {
1097 continue;
1098 };
1099 question.abandon(format!(
1100 "task {} no longer waits for this answer",
1101 task.short()
1102 ));
1103 if let Err(e) = questions.put(&mut question) {
1104 tracing::warn!(
1105 "could not retire question {} for task {}: {e:#}",
1106 question.short(),
1107 task.short()
1108 );
1109 }
1110 }
1111}
1112
1113#[derive(Debug, Clone, Copy)]
1120pub struct Verdict {
1121 pub status: RunStatus,
1123 pub left_pr: bool,
1125 pub quota_hit: bool,
1127 pub parked: bool,
1129 pub no_viable_candidates: bool,
1136}
1137
1138pub fn settle(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize) {
1195 settle_in(task, verdict, detail, max_attempts, &PHRASES_EN)
1196}
1197
1198fn settle_in(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize, p: &Phrases) {
1200 if verdict.parked {
1206 task.stall(detail);
1207 return;
1208 }
1209 match verdict.status {
1210 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1211 RunStatus::AlreadyInBase => task.already_landed(detail),
1212 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1213 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1214 task.stall(detail)
1215 }
1216 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1217 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1218 RunStatus::Blocked => task.fail(detail, max_attempts),
1219 RunStatus::VerifiedNoop => task.handed_off(detail),
1220 other => task.fail((p.graph_stopped)(label(other), detail), max_attempts),
1221 }
1222}
1223
1224pub fn supersede_prior_runs(task: &Task, home: &Path) {
1273 let now = Timestamp::now();
1274 let last_run_succeeded = task
1281 .runs
1282 .last()
1283 .and_then(|id| RunState::load_under(id, home).ok())
1284 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1285 for id in task.superseded_attempts(last_run_succeeded) {
1286 let mut state = match RunState::load_under(id, home) {
1287 Ok(s) => s,
1288 Err(e) => {
1289 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1290 continue;
1291 }
1292 };
1293 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1294 continue;
1295 }
1296 let daemon_claims = is_working_on(home, id, now);
1297 if state.liveness(daemon_claims) == Liveness::Live {
1298 continue;
1299 }
1300 state.status = RunStatus::Superseded;
1301 if let Err(e) = state.save_under(home) {
1302 tracing::warn!("could not mark run {id} superseded: {e:#}");
1303 }
1304 }
1305}
1306
1307fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1329 for task in queue.list() {
1330 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1331 continue;
1332 }
1333 supersede_prior_runs(&task, home);
1334 }
1335}
1336
1337fn settle_and_diagnose(
1344 task: &mut Task,
1345 verdict: Verdict,
1346 detail: &str,
1347 max_attempts: usize,
1348 state: &RunState,
1349) {
1350 let p = phrases(&state.config.graph.language);
1351 settle_in(task, verdict, detail, max_attempts, p);
1352 if task.status == TaskStatus::Held {
1353 task.diagnostic = diagnostic(state);
1354 note_open_question(task, &state.id, p);
1355 }
1356}
1357
1358fn note_open_question(task: &mut Task, run: &str, p: &Phrases) {
1373 let Some(home) = crate::run::try_home() else {
1374 return;
1375 };
1376 let open = Questions::at(home.join("questions")).open_for(run);
1377 let Some(q) = open.first() else {
1378 return;
1379 };
1380 let base = task.hold_reason.clone().unwrap_or_default();
1381 task.hold_reason = Some(format!("{base}{}{}", p.waiting_for_answer, q.short()));
1382}
1383
1384fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize, language: &str) {
1395 match last_run {
1396 Some(state) => {
1397 let verdict = Verdict {
1398 status: state.status,
1399 left_pr: state.pr.is_some(),
1400 quota_hit: !state.quota.is_empty(),
1401 parked: state.parked,
1402 no_viable_candidates: state.viable().is_empty(),
1403 };
1404 let detail = format!(
1405 "{}{}",
1406 phrases(&state.config.graph.language).recovered_running,
1407 describe(&state)
1408 );
1409 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1410 }
1411 None => {
1412 let why = phrases(language).no_run_to_recover;
1415 task.last_error = Some(why.to_owned());
1416 task.hold_machine(Some(why.to_owned()));
1419 }
1420 }
1421}
1422
1423fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1445 let mut reclaimed = Vec::new();
1446 for listed in queue.list() {
1447 if listed.status != TaskStatus::Running {
1448 continue;
1449 }
1450 let Ok(_claim) = queue.claim(&listed.id) else {
1451 continue;
1452 };
1453 let Ok(mut task) = queue.get(&listed.id) else {
1457 continue;
1458 };
1459 if task.status != TaskStatus::Running {
1460 continue;
1461 }
1462 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1463 if let Some(state) = &last_run
1473 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1474 {
1475 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1476 }
1477 let language = if last_run.is_none() {
1478 language_of(&task, Path::new("."))
1479 } else {
1480 String::new()
1481 };
1482 reclaim(&mut task, last_run, max_attempts, &language);
1483 if task.status == TaskStatus::Done {
1484 supersede_prior_runs(&task, &crate::run::home());
1485 }
1486 record(queue, &mut task);
1487 reclaimed.push(task.id.clone());
1488 }
1489 reclaimed
1490}
1491
1492fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1524 reclaim_abandoned_runs_with(
1525 home,
1526 now,
1527 crate::proc::pid_status,
1528 crate::proc::process_started_at,
1529 )
1530}
1531
1532fn reclaim_abandoned_runs_with<F, G>(
1538 home: &Path,
1539 now: Timestamp,
1540 query: F,
1541 identity: G,
1542) -> Vec<String>
1543where
1544 F: Fn(u32) -> Option<bool>,
1545 G: Fn(u32) -> Option<String>,
1546{
1547 let mut abandoned = Vec::new();
1548 for entry in std::fs::read_dir(home.join("runs"))
1549 .into_iter()
1550 .flatten()
1551 .flatten()
1552 {
1553 let id = entry.file_name().to_string_lossy().into_owned();
1554 if !crate::run::is_run_id(&id) {
1555 continue;
1556 }
1557 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1565 continue;
1566 };
1567 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1568 continue;
1569 };
1570 if state.status.done() || !state.active_all_overrun(now) {
1571 continue;
1572 }
1573 let daemon_claims = is_working_on(home, &id, now);
1585 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1586 continue;
1587 }
1588 state.abandon("daemon");
1589 if let Err(e) = state.save_under(home) {
1590 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1591 continue;
1592 }
1593 if let Some(notice) = notices::run_ended(&state) {
1596 notices::raise_in(home, notice);
1597 }
1598 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1606 tracing::warn!("abandon questions for {id}: {e:#}");
1607 }
1608 abandoned.push(id);
1609 }
1610 abandoned
1611}
1612
1613pub async fn serve(opts: Opts) -> Result<()> {
1619 serve_until(opts, Stop::new()).await
1620}
1621
1622pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1639 let signal = {
1640 let stop = stop.clone();
1641 tokio::spawn(async move {
1642 if tokio::signal::ctrl_c().await.is_ok() {
1643 stop.stop();
1644 tracing::info!("shutdown requested; a run in flight will be finished first");
1645 }
1646 })
1647 };
1648
1649 let worktrees_root = opts
1650 .worktrees_root
1651 .clone()
1652 .unwrap_or_else(crate::run::default_worktree_root);
1653 let outcome = drive(
1654 &opts,
1655 &Queue::open(),
1656 &status_path(),
1657 &crate::run::home(),
1658 &worktrees_root,
1659 &stop,
1660 )
1661 .await;
1662
1663 signal.abort();
1664 outcome
1665}
1666
1667async fn drive(
1680 opts: &Opts,
1681 queue: &Queue,
1682 status_file: &Path,
1683 home: &Path,
1684 worktrees_root: &Path,
1685 stop: &Stop,
1686) -> Result<()> {
1687 let status = Arc::new(Mutex::new(Status::new()));
1695 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1696 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1697
1698 let daemon_cfg = prepare(&opts.repo, opts)
1704 .map(|c| c.daemon)
1705 .unwrap_or_default();
1706 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1707
1708 let waiter = tokio::spawn(crate::waiter::run(
1713 crate::waiter::Waiter::new(
1714 crate::ask::Questions::at(home.join("questions")),
1715 home.to_path_buf(),
1716 prepare(&opts.repo, opts).ok(),
1717 ),
1718 stop.clone(),
1719 ));
1720
1721 let deputies = tokio::spawn(crate::deputy::run(
1725 crate::deputy::Deputies::new(
1726 crate::ask::Questions::at(home.join("questions")),
1727 home.to_path_buf(),
1728 prepare(&opts.repo, opts).ok(),
1729 opts.repo.clone(),
1730 daemon_cfg.max_deputies,
1731 {
1732 let stop = stop.clone();
1733 Arc::new(move || stop.parking())
1734 },
1735 ),
1736 stop.clone(),
1737 ));
1738
1739 let fetcher = tokio::spawn(fetch_loop(opts.repo.clone(), opts.clone(), stop.clone()));
1742
1743 tracing::info!(
1744 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1745 queue.root().display(),
1746 opts.poll.as_secs(),
1747 opts.max_attempts,
1748 concurrency,
1749 if daemon_cfg.pause_for_interrupts {
1750 ", interrupts enabled"
1751 } else {
1752 ""
1753 }
1754 );
1755
1756 janitor(&opts.repo, opts, home, worktrees_root).await;
1759 resweep_superseded_attempts(queue, home);
1760
1761 let outcome = poll(
1762 opts,
1763 queue,
1764 &status,
1765 home,
1766 worktrees_root,
1767 stop,
1768 DispatchLimits {
1769 max_concurrent: concurrency,
1770 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1771 },
1772 )
1773 .await;
1774
1775 beat.abort();
1776 waiter.abort();
1777 deputies.abort();
1778 fetcher.abort();
1779 clear_status_at(status_file);
1780 outcome
1781}
1782
1783const FETCH_TIMEOUT: Duration = Duration::from_secs(30);
1785
1786async fn fetch_loop(repo: PathBuf, opts: Opts, stop: Stop) {
1790 while !stop.stopped() {
1791 let (interval, roots) = match prepare(&repo, &opts) {
1792 Ok(c) => (c.repos.fetch_interval, c.repos.roots),
1793 Err(_) => (0, Vec::new()),
1794 };
1795 if interval > 0 && !roots.is_empty() {
1796 let r = crate::clean::fetch_origins(&roots, FETCH_TIMEOUT, || stop.stopped()).await;
1797 tracing::debug!("fetch origins: {r:?}");
1798 }
1799 let wait = if interval > 0 { interval } else { 60 };
1801 let mut slept = 0;
1802 while slept < wait && !stop.stopped() {
1803 tokio::time::sleep(Duration::from_secs(1)).await;
1804 slept += 1;
1805 }
1806 }
1807}
1808
1809async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1815 loop {
1816 tokio::time::sleep(HEARTBEAT).await;
1817 let snapshot = {
1818 let mut guard = lock(&status);
1819 guard.updated_at = Timestamp::now();
1820 guard.clone()
1821 };
1822 if let Err(e) = write_status_to(&path, &snapshot) {
1823 tracing::warn!("could not refresh the daemon status file: {e:#}");
1826 }
1827 }
1828}
1829
1830#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1833enum LandResume {
1834 NotLanding,
1837 StillWaiting,
1842 Ready,
1846}
1847
1848fn land_resume_state(task: &Task) -> LandResume {
1852 let Some(run_id) = task.runs.last() else {
1853 return LandResume::NotLanding;
1854 };
1855 let Ok(state) = RunState::load(run_id) else {
1856 return LandResume::NotLanding;
1857 };
1858 if state.status != RunStatus::Landing || !state.parked {
1859 return LandResume::NotLanding;
1860 }
1861 let store = ask::Questions::open();
1862 let waiting = store
1863 .list()
1864 .into_iter()
1865 .filter(|q| &q.run == run_id && q.node == land::APPROVAL_NODE)
1866 .max_by(|a, b| a.id.cmp(&b.id));
1867 let Some(mut q) = waiting else {
1868 return LandResume::Ready;
1869 };
1870 if !q.status.open() {
1871 return LandResume::Ready;
1872 }
1873 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
1880 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
1881 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
1882 q.abandon(format!(
1883 "no answer within {}s of asking",
1884 timeout.as_secs().max(1)
1885 ));
1886 if store.put(&mut q).is_ok() {
1889 return LandResume::Ready;
1890 }
1891 }
1892 LandResume::StillWaiting
1893}
1894
1895const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
1903
1904const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
1916
1917struct InFlightGuard<'a> {
1930 status: &'a Arc<Mutex<Status>>,
1931 stop: &'a Stop,
1932 task_id: &'a str,
1933}
1934
1935impl Drop for InFlightGuard<'_> {
1936 fn drop(&mut self) {
1937 lock(self.status).current.retain(|c| c.task != self.task_id);
1938 self.stop.exit();
1939 }
1940}
1941
1942#[derive(Debug, Clone, PartialEq, Eq)]
1964enum Interrupt {
1965 Idle,
1967 Parking {
1979 parked: Vec<String>,
1980 interrupt_task: String,
1981 },
1982 Running {
1990 parked: Vec<String>,
1991 interrupt_task: String,
1992 },
1993 Resuming { parked: Vec<String> },
2000}
2001
2002fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
2023 match state {
2024 Interrupt::Idle => {
2025 if in_flight.len() != 1 {
2037 return Interrupt::Idle;
2038 }
2039 match runnable.iter().find(|t| t.interrupt) {
2040 Some(t) => Interrupt::Parking {
2041 parked: in_flight.to_vec(),
2042 interrupt_task: t.id.clone(),
2043 },
2044 None => Interrupt::Idle,
2045 }
2046 }
2047 Interrupt::Parking {
2048 parked,
2049 interrupt_task,
2050 } => {
2051 if in_flight.iter().any(|id| parked.contains(id)) {
2052 Interrupt::Parking {
2054 parked,
2055 interrupt_task,
2056 }
2057 } else if in_flight.contains(&interrupt_task) {
2058 Interrupt::Running {
2059 parked,
2060 interrupt_task,
2061 }
2062 } else if runnable.iter().any(|t| t.id == interrupt_task) {
2063 Interrupt::Parking {
2067 parked,
2068 interrupt_task,
2069 }
2070 } else {
2071 Interrupt::Resuming { parked }
2076 }
2077 }
2078 Interrupt::Running {
2079 parked,
2080 interrupt_task,
2081 } => {
2082 if in_flight.contains(&interrupt_task) {
2083 Interrupt::Running {
2084 parked,
2085 interrupt_task,
2086 }
2087 } else {
2088 Interrupt::Resuming { parked }
2094 }
2095 }
2096 Interrupt::Resuming { parked } => {
2097 if in_flight.iter().any(|id| parked.contains(id)) {
2098 Interrupt::Idle
2104 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
2105 Interrupt::Resuming { parked }
2106 } else {
2107 Interrupt::Idle
2110 }
2111 }
2112 }
2113}
2114
2115fn advance_interrupt_tick(
2121 enabled: bool,
2122 state: Interrupt,
2123 in_flight: &[String],
2124 runnable: &[Task],
2125) -> Interrupt {
2126 if !enabled {
2127 return Interrupt::Idle;
2128 }
2129 advance_interrupt(state, in_flight, runnable)
2130}
2131
2132fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
2137 match state {
2138 Interrupt::Idle => candidates,
2139 Interrupt::Parking {
2140 parked,
2141 interrupt_task,
2142 } => {
2143 if in_flight.iter().any(|id| parked.contains(id)) {
2144 Vec::new()
2145 } else {
2146 candidates
2147 .into_iter()
2148 .filter(|t| &t.id == interrupt_task)
2149 .collect()
2150 }
2151 }
2152 Interrupt::Running { .. } => Vec::new(),
2153 Interrupt::Resuming { parked } => candidates
2161 .into_iter()
2162 .find(|t| parked.contains(&t.id))
2163 .into_iter()
2164 .collect(),
2165 }
2166}
2167
2168struct DispatchLimits {
2172 max_concurrent: usize,
2175 pause_for_interrupts: bool,
2177}
2178
2179#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2188enum PermitKind {
2189 None,
2194 Urgent,
2201 Ordinary,
2204}
2205
2206fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
2209 if priority {
2210 PermitKind::None
2211 } else if urgent {
2212 PermitKind::Urgent
2213 } else {
2214 PermitKind::Ordinary
2215 }
2216}
2217
2218async fn poll(
2237 opts: &Opts,
2238 queue: &Queue,
2239 status: &Arc<Mutex<Status>>,
2240 home: &Path,
2241 worktrees_root: &Path,
2242 stop: &Stop,
2243 limits: DispatchLimits,
2244) -> Result<()> {
2245 let DispatchLimits {
2246 max_concurrent,
2247 pause_for_interrupts,
2248 } = limits;
2249 let mut attempted: Vec<String> = Vec::new();
2254 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2255 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2263 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2269 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2270 let mut conductor = Conductor::new();
2271 let mut cache_last_checked: Option<Timestamp> = None;
2274 let mut interrupt = Interrupt::Idle;
2276 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2282 std::collections::HashMap::new();
2283
2284 while !stop.stopped() {
2285 lock(status).polls += 1;
2286
2287 while let Some(result) = inflight.try_join_next() {
2292 if let Err(e) = result {
2293 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2294 notices::raise(Notice::error(
2295 "loop:attempt",
2296 "A queued attempt ended abnormally; check the task it was running.",
2297 ));
2298 }
2299 }
2300
2301 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2302 if !swept.is_empty() {
2303 tracing::warn!(
2304 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2305 swept.len(),
2306 swept.join(", ")
2307 );
2308 }
2309 let now = Timestamp::now();
2314
2315 if !stop.busy_now() {
2320 maybe_prune_cache_between_runs(
2321 &opts.repo,
2322 opts,
2323 home,
2324 stop,
2325 &mut cache_last_checked,
2326 now,
2327 )
2328 .await;
2329 }
2330
2331 let stalled = stalled_tasks(queue, home, now);
2332 let stalled_ids: std::collections::BTreeSet<_> =
2333 stalled.iter().map(|task| task.id.clone()).collect();
2334 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2335 if !reclaimed.is_empty() {
2336 tracing::warn!(
2337 "reclaimed {} task(s) left `running` by a daemon that never \
2338 recorded the outcome: {}",
2339 reclaimed.len(),
2340 reclaimed.join(", ")
2341 );
2342 }
2343 let abandoned_runs = reclaim_abandoned_runs(home, now);
2344 if !abandoned_runs.is_empty() {
2345 tracing::warn!(
2346 "failed {} run(s) left behind by a killed process, past every \
2347 active seat's own timeout: {}",
2348 abandoned_runs.len(),
2349 abandoned_runs.join(", ")
2350 );
2351 }
2352
2353 let questions = Questions::at(home.join("questions"));
2358
2359 resolve_blockers(queue, &questions);
2362 apply_choice_actions(queue, &questions, home);
2363 reconcile_task_questions(queue, &questions);
2364
2365 let finished: Vec<Task> = finished_tasks(queue)
2371 .into_iter()
2372 .filter(|task| !stalled_ids.contains(&task.id))
2373 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2377 .collect();
2378 let queued = queued_tasks(queue);
2379 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2383 && conductor.worth_a_look(queue, &stalled, &finished)
2384 {
2385 match prepare(&opts.repo, opts) {
2386 Ok(cfg) => {
2387 conductor
2388 .maybe_run(
2389 &cfg,
2390 &opts.repo,
2391 queue,
2392 &questions,
2393 home,
2394 &queued,
2395 &stalled,
2396 &finished,
2397 opts.max_attempts,
2398 )
2399 .await;
2400 }
2401 Err(e) => {
2402 tracing::warn!("conductor: no config: {e:#}");
2403 notices::raise(Notice::warn(
2404 "loop:no-config",
2405 "The loop could not read this repository's config, so held tasks are not being triaged.",
2406 ));
2407 }
2408 }
2409 }
2410
2411 let candidates: Vec<Task> = runnable(queue)
2412 .into_iter()
2413 .filter(|t| !opts.once || !attempted.contains(&t.id))
2414 .collect();
2415
2416 let in_flight: Vec<String> = lock(status)
2421 .current
2422 .iter()
2423 .map(|c| c.task.clone())
2424 .collect();
2425 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2426
2427 interrupt =
2428 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2429 if let Interrupt::Parking {
2430 parked,
2431 interrupt_task,
2432 } = &interrupt
2433 {
2434 let reason = format!(
2435 "task {} asked to run first",
2436 crate::run::short_of(interrupt_task)
2437 );
2438 for id in parked {
2439 if let Some(pause) = interrupt_pauses.get(id) {
2440 pause.park_because(reason.clone());
2441 }
2442 }
2443 }
2444 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2445
2446 let cooling_down =
2447 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2448
2449 let mut started_any = false;
2450 for candidate in candidates {
2451 if stop.stopped() {
2452 break;
2453 }
2454
2455 let resume = land_resume_state(&candidate);
2456 if resume == LandResume::StillWaiting {
2457 continue;
2458 }
2459 let priority = resume == LandResume::Ready;
2460
2461 if !priority && cooling_down {
2462 continue;
2463 }
2464 let permit = match permit_kind(priority, candidate.urgent) {
2465 PermitKind::None => None,
2466 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2467 Ok(p) => Some(p),
2468 Err(_) => continue,
2474 },
2475 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2476 Ok(p) => Some(p),
2477 Err(_) => continue,
2481 },
2482 };
2483
2484 let Ok(claim) = queue.claim(&candidate.id) else {
2489 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2490 continue;
2491 };
2492 let mut task = match queue.get(&candidate.id) {
2495 Ok(t) if t.status.runnable() => t,
2496 Ok(_) => continue,
2497 Err(e) => {
2498 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2499 continue;
2500 }
2501 };
2502 let task_id = task.id.clone();
2503 attempted.push(task_id.clone());
2504 lock(status).idle = false;
2505 stop.enter();
2508 started_any = true;
2509
2510 let run_pause = crate::graph::Pause::new();
2514 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2515
2516 let opts = opts.clone();
2517 let queue = queue.clone();
2518 let status = Arc::clone(status);
2519 let stop = stop.clone();
2520 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2521 inflight.spawn(async move {
2522 let _claim = claim;
2526 let _permit = permit;
2527 let _inflight = InFlightGuard {
2529 status: &status,
2530 stop: &stop,
2531 task_id: &task_id,
2532 };
2533 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2534 lock(&status).completed += 1;
2535 let now = Timestamp::now();
2541 if let Some(until) = cooldown_until("a, now) {
2542 let wait = until.as_second() - now.as_second();
2543 *lock("a_cooldown_until) = Some(until);
2544 let hint = quota
2545 .iter()
2546 .find(|q| q.reset.is_some())
2547 .and_then(|q| q.reset.as_deref());
2548 match hint {
2549 Some(h) => tracing::warn!(
2550 "quota hit; waiting {wait}s before taking another ordinary task \
2551 (CLI reported reset: {h})"
2552 ),
2553 None => tracing::warn!(
2554 "quota hit; waiting {wait}s before taking another ordinary task \
2555 (no reset hint reported)"
2556 ),
2557 }
2558 }
2559 });
2560 }
2561
2562 if started_any {
2563 continue;
2564 }
2565
2566 if stop.busy_now() {
2567 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2572 continue;
2573 }
2574
2575 lock(status).idle = true;
2577 if opts.once {
2578 janitor(&opts.repo, opts, home, worktrees_root).await;
2582 resweep_superseded_attempts(queue, home);
2583 triage_held(queue, home, opts).await;
2584 break;
2585 }
2586 stop.idle(opts.poll).await;
2587 if stop.stopped() {
2588 continue;
2589 }
2590 janitor(&opts.repo, opts, home, worktrees_root).await;
2596 resweep_superseded_attempts(queue, home);
2597 triage_held(queue, home, opts).await;
2598 }
2599
2600 while let Some(result) = inflight.join_next().await {
2605 if let Err(e) = result {
2606 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2607 notices::raise(Notice::error(
2608 "loop:attempt",
2609 "A queued attempt ended abnormally; check the task it was running.",
2610 ));
2611 }
2612 }
2613 Ok(())
2614}
2615
2616async fn attempt(
2622 opts: &Opts,
2623 queue: &Queue,
2624 status: &Arc<Mutex<Status>>,
2625 stop: &Stop,
2626 interrupt_pause: crate::graph::Pause,
2627 task: &mut Task,
2628) -> Vec<QuotaLoss> {
2629 let repo = repo_for(task, &opts.repo);
2630 tracing::info!(
2631 "task {} — {} (repo {})",
2632 task.short(),
2633 task.title,
2634 repo.display()
2635 );
2636
2637 let mut config = match prepare(&repo, opts) {
2638 Ok(c) => c,
2639 Err(e) => {
2640 task.attempts += 1;
2644 task.fail(format!("config: {e:#}"), opts.max_attempts);
2645 record(queue, task);
2646 return Vec::new();
2647 }
2648 };
2649 apply_solo(&mut config, task);
2650 let p = phrases(&config.graph.language);
2651 let start_failed = p.could_not_start;
2652
2653 if let Some(reason) = disk_gate(&repo, &config) {
2661 task.last_error = Some(reason.clone());
2662 task.hold_machine(Some(reason.clone()));
2663 record(queue, task);
2664 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2665 notices::raise(
2668 Notice::warn(
2669 &format!("disk:{}", repo.display()),
2670 "A task was held for want of free disk space; free some, then release it from the queue.",
2671 )
2672 .link(Link::Task {
2673 id: task.id.clone(),
2674 }),
2675 );
2676 return Vec::new();
2677 }
2678
2679 let unfinished = (!task.fresh_start)
2699 .then(|| unfinished_run(&task.runs, task.short()))
2700 .flatten();
2701 let review_branch = task.review_branch.take();
2707 let branch_exists = match &review_branch {
2708 Some(branch) => crate::git::branch_exists(&repo, branch)
2709 .await
2710 .unwrap_or(false),
2711 None => false,
2712 };
2713 let starter = choose_starter(
2714 review_branch.as_deref(),
2715 branch_exists,
2716 unfinished.as_deref(),
2717 );
2718 let attachments = match task_attachments(queue, task) {
2723 Ok(a) => a,
2724 Err(e) => {
2725 task.attempts += 1;
2726 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2727 record(queue, task);
2728 return Vec::new();
2729 }
2730 };
2731 let started = match &starter {
2732 Starter::Review(branch) => {
2733 tracing::info!(
2734 "task {} reopens `{branch}` as a review-only pass",
2735 task.short()
2736 );
2737 let takeover = crate::handover::Takeover {
2740 earlier: task.earlier_attempts().to_vec(),
2741 home: crate::run::home(),
2742 choice: take_divergence_answer(branch, &config.merge.remote, task),
2743 };
2744 Runner::review_taking_over(
2745 &repo,
2746 branch,
2747 config,
2748 Some(takeover),
2749 crate::run::Origin::queue(&task.id),
2750 )
2751 .await
2752 }
2753 Starter::Resume(id) => {
2754 tracing::info!("resuming run {id} rather than competing again");
2755 Runner::resume(id).map(|mut r| {
2756 if let Some(instruction) =
2757 prepare_instruction(&starter, Some(&r.state.instruction), task)
2758 {
2759 r.state.instruction = instruction;
2760 }
2761 r.state.attachments = attachments.clone();
2762 r
2763 })
2764 }
2765 Starter::Start => {
2766 if let Some(branch) = &review_branch {
2767 tracing::warn!(
2768 "conductor chose review for task {} but branch `{branch}` no longer \
2769 exists; requeuing as a fresh competition instead",
2770 task.short()
2771 );
2772 }
2773 let instruction = prepare_instruction(&starter, None, task)
2774 .unwrap_or_else(|| task.instruction.clone());
2775 Runner::start_naming(
2776 &repo,
2777 instruction,
2778 &task.title,
2779 config,
2780 crate::run::Origin::queue(&task.id),
2781 )
2782 .await
2783 .map(|mut r| {
2784 r.state.attachments = attachments.clone();
2785 r
2786 })
2787 }
2788 };
2789 let mut runner = match started {
2790 Ok(r) => r,
2791 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2795 let detail = match e.downcast_ref::<crate::handover::Refused>() {
2796 Some(r) => (p.handover_refused)(r),
2797 None => format!("{e:#}"),
2798 };
2799 let reason = format!("{start_failed}{detail}");
2800 task.last_error = Some(reason.clone());
2801 let branch = match &starter {
2804 Starter::Review(branch) => Some(branch.clone()),
2805 _ => None,
2806 };
2807 task.hold_for_handover(branch, format!("{reason}{}", p.handover_hint));
2811 record(queue, task);
2812 tracing::warn!(
2813 "holding {} for a branch it cannot take over: {e:#}",
2814 task.short()
2815 );
2816 return Vec::new();
2817 }
2818 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2823 let d = e
2824 .downcast_ref::<crate::reconcile::Diverged>()
2825 .expect("checked by the guard");
2826 let mut q = ask::Question::new(
2827 task.id.clone(),
2828 "review".to_owned(),
2829 "sync".to_owned(),
2830 d.summary(),
2831 d.detail(),
2832 d.choices(),
2833 );
2834 match Questions::open().put(&mut q) {
2835 Ok(()) => {
2836 task.last_error = Some(format!("{e:#}"));
2837 task.review_branch = Some(d.branch.clone());
2840 task.block(vec![q.id.clone()], Some(d.summary()));
2841 }
2842 Err(put) => {
2843 tracing::warn!("could not file the divergence question: {put:#}");
2844 task.attempts += 1;
2845 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2846 }
2847 }
2848 record(queue, task);
2849 return Vec::new();
2850 }
2851 Err(e) => {
2852 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
2858 && let Starter::Review(branch) = &starter
2859 {
2860 task.review_branch = Some(branch.clone());
2861 task.last_error = Some(format!("{start_failed}{e:#}"));
2862 task.status = crate::queue::TaskStatus::Failed;
2863 record(queue, task);
2864 return Vec::new();
2865 }
2866 task.attempts += 1;
2867 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2868 record(queue, task);
2869 return Vec::new();
2870 }
2871 };
2872 runner.state.followup_generation = Some(task.followup.as_ref().map_or(0, |f| f.generation));
2875 runner.on_pause(stop.pause());
2877 runner.watch_interrupt(interrupt_pause);
2881
2882 let run = runner.state.id.clone();
2885 task.start(run.clone());
2886 record(queue, task);
2887 lock(status).current.push(Current {
2888 task: task.id.clone(),
2889 run,
2890 });
2891
2892 let quota_before = runner.state.quota.clone();
2895 let detail = match runner.execute().await {
2896 Ok(()) => describe(&runner.state),
2897 Err(e) => format!("{e:#}"),
2898 };
2899 let fresh = losses_this_attempt("a_before, &runner.state.quota);
2900 let verdict = Verdict {
2901 status: runner.state.status,
2902 left_pr: runner.state.pr.is_some(),
2905 quota_hit: !fresh.is_empty(),
2911 parked: runner.state.parked,
2915 no_viable_candidates: runner.state.viable().is_empty(),
2918 };
2919 settle_and_diagnose(task, verdict, &detail, opts.max_attempts, &runner.state);
2920 if task.status == TaskStatus::Done {
2921 supersede_prior_runs(task, &crate::run::home());
2922 }
2923 record(queue, task);
2924 tracing::info!(
2925 "task {} is {} after run {} ({})",
2926 task.short(),
2927 task.status.as_str(),
2928 runner.state.short(),
2929 label(runner.state.status)
2930 );
2931 fresh
2932}
2933
2934fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
2944 after
2945 .iter()
2946 .filter(|q| !before.contains(q))
2947 .cloned()
2948 .collect()
2949}
2950
2951fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
2954 if quota.is_empty() {
2955 return None;
2956 }
2957 let with_hint = quota.iter().find(|q| q.reset.is_some());
2958 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
2959 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
2960 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
2961 Some(
2962 now.checked_add(jiff::SignedDuration::from_secs(secs))
2963 .unwrap_or(Timestamp::MAX),
2964 )
2965}
2966
2967fn apply_solo(config: &mut Config, task: &Task) {
2977 if task.solo {
2978 config.graph.candidates = 1;
2979 }
2980}
2981
2982fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
2984 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
2985 if let Some(mode) = &opts.merge {
2986 config.merge.mode = merge_mode(mode)?;
2987 }
2988 Ok(config)
2989}
2990
2991async fn maybe_prune_cache_between_runs(
3022 repo: &Path,
3023 opts: &Opts,
3024 home: &Path,
3025 stop: &Stop,
3026 last_checked: &mut Option<Timestamp>,
3027 now: Timestamp,
3028) {
3029 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
3030 return;
3031 }
3032 *last_checked = Some(now);
3033 let cfg = match prepare(repo, opts) {
3034 Ok(cfg) => cfg,
3035 Err(e) => {
3036 tracing::warn!("cache check: no config: {e:#}");
3037 return;
3038 }
3039 };
3040 match clean::prune_cache_if_over_limit(&cfg, home) {
3041 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
3042 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
3043 pruned.files,
3044 pruned.freed
3045 ),
3046 Ok(_) => {}
3047 Err(e) => {
3048 tracing::warn!("housekeep: prune cache: {e:#}");
3049 notices::raise_in(
3050 home,
3051 Notice::warn(
3052 "housekeep:cache",
3053 "Pruning the shared build cache failed; disk usage may keep growing.",
3054 ),
3055 );
3056 }
3057 }
3058}
3059
3060fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
3064 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
3065}
3066
3067async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
3086 let cfg = match prepare(repo, opts) {
3087 Ok(cfg) => cfg,
3088 Err(e) => {
3089 tracing::warn!("housekeep: no config: {e:#}");
3090 return;
3091 }
3092 };
3093 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
3102 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
3103 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
3108 let mut extra = Vec::new();
3109 if out.unreadable > 0 {
3110 extra.push(format!("{} unreadable", out.unreadable));
3111 }
3112 if out.orphaned_worktrees > 0 {
3113 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
3114 }
3115 let detail = if extra.is_empty() {
3116 String::new()
3117 } else {
3118 format!(" ({})", extra.join(", "))
3119 };
3120 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
3121 }
3122 if out.external_merges_recorded > 0 {
3123 tracing::info!(
3124 "housekeep: recorded {} run(s) as merged externally",
3125 out.external_merges_recorded
3126 );
3127 }
3128 if out.stale_pr_states_repaired > 0 {
3129 tracing::info!(
3130 "housekeep: rewrote {} run record(s) whose pull request had already settled",
3131 out.stale_pr_states_repaired
3132 );
3133 }
3134 if out.cache_files > 0 {
3135 tracing::info!(
3136 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
3137 out.cache_files,
3138 out.cache_freed
3139 );
3140 }
3141 if out.questions_abandoned > 0 {
3142 tracing::info!(
3143 "housekeep: abandoned {} question(s) left open by a finished run",
3144 out.questions_abandoned
3145 );
3146 }
3147}
3148
3149async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3158 let questions = Questions::at(home.join("questions"));
3159 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3160 if report.is_empty() {
3161 return;
3162 }
3163 if !report.quarantined.is_empty() {
3164 tracing::info!(
3165 "triage: held {} blocked task(s) whose blocked-on task or \
3166 question no longer exists: {}",
3167 report.quarantined.len(),
3168 report.quarantined.join(", ")
3169 );
3170 }
3171 if !report.resumed.is_empty() {
3172 tracing::info!(
3173 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3174 report.resumed.len(),
3175 report.resumed.join(", ")
3176 );
3177 }
3178 if !report.asked.is_empty() {
3179 tracing::info!(
3180 "triage: asked about {} held task(s): {}",
3181 report.asked.len(),
3182 report.asked.join(", ")
3183 );
3184 }
3185 if !report.answered.is_empty() {
3186 tracing::info!(
3187 "triage: applied {} operator answer(s): {}",
3188 report.answered.len(),
3189 report.answered.join(", ")
3190 );
3191 }
3192}
3193
3194fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3201 disk_gate_with(repo, config, crate::disk::free_bytes)
3202}
3203
3204fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3208 repo: &Path,
3209 config: &Config,
3210 free_bytes: F,
3211) -> Option<String> {
3212 let min = config.disk.min_free_bytes;
3213 if min == 0 {
3214 return None;
3215 }
3216 match free_bytes(repo) {
3217 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3218 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3219 }
3220}
3221
3222const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3229
3230const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3234
3235fn quota_wait(
3244 reset_at: Option<Timestamp>,
3245 now: Timestamp,
3246 fallback: Duration,
3247 cap: Duration,
3248) -> Duration {
3249 match reset_at {
3250 Some(at) if at > now => {
3251 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3252 Duration::from_secs(secs).min(cap)
3253 }
3254 _ => fallback,
3255 }
3256}
3257
3258fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3272 parse_reset_hint_zoned(text, now)
3273 .or_else(|| parse_reset_hint_dated(text))
3274 .or_else(|| parse_reset_hint_relative(text, recorded))
3275}
3276
3277fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3281 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3282 let mut rest = rest.trim();
3283 if rest.is_empty() {
3284 return None;
3285 }
3286 let mut total: i64 = 0;
3287 let mut matched = false;
3288 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3289 if let Some((digits, tail)) = rest.split_once(unit)
3290 && !digits.is_empty()
3291 && digits.bytes().all(|b| b.is_ascii_digit())
3292 {
3293 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3294 rest = tail;
3295 matched = true;
3296 }
3297 }
3298 if !rest.is_empty() || !matched {
3299 return None;
3300 }
3301 recorded
3302 .checked_add(jiff::SignedDuration::from_secs(total))
3303 .ok()
3304}
3305
3306fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3310 let clock = clock.trim().to_lowercase();
3311 let (digits, pm) = clock
3312 .strip_suffix("am")
3313 .map(|d| (d, false))
3314 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3315 let (h, m) = digits.trim().split_once(':')?;
3316 let mut hour: i8 = h.trim().parse().ok()?;
3317 let minute: i8 = m.trim().parse().ok()?;
3318 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3319 return None;
3320 }
3321 if pm && hour != 12 {
3322 hour += 12;
3323 } else if !pm && hour == 12 {
3324 hour = 0;
3325 }
3326 Some((hour, minute))
3327}
3328
3329fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3334 let open = text.find('(')?;
3335 let close = text.rfind(')')?;
3336 if close <= open {
3337 return None;
3338 }
3339 let zone = text[open + 1..close].trim();
3340 let (hour, minute) = parse_12h_clock(&text[..open])?;
3341 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3342 let candidate = now
3343 .to_zoned(tz)
3344 .with()
3345 .hour(hour)
3346 .minute(minute)
3347 .second(0)
3348 .millisecond(0)
3349 .microsecond(0)
3350 .nanosecond(0)
3351 .build()
3352 .ok()?;
3353 let mut at = candidate.timestamp();
3354 if at <= now {
3355 at += jiff::SignedDuration::from_hours(24);
3356 }
3357 Some(at)
3358}
3359
3360fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3369 let words: Vec<&str> = text.split_whitespace().collect();
3370 if words.len() < 5 {
3371 return None;
3372 }
3373 (0..=words.len() - 5)
3374 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3375}
3376
3377fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3383 if trailing.is_some_and(|next| next.starts_with('(')) {
3384 return None;
3385 }
3386 let month = month_number(window[0])?;
3387 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3388 let day_digits = ["st", "nd", "rd", "th"]
3389 .iter()
3390 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3391 let day: i8 = day_digits.parse().ok()?;
3392 let year_token = window[2];
3393 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3394 return None;
3395 }
3396 let year: i16 = year_token.parse().ok()?;
3397 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3401 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3402 let date = jiff::civil::Date::new(year, month, day).ok()?;
3403 let candidate = date
3404 .at(hour, minute, 0, 0)
3405 .to_zoned(jiff::tz::TimeZone::UTC)
3406 .ok()?;
3407 Some(candidate.timestamp())
3408}
3409
3410fn month_number(name: &str) -> Option<i8> {
3413 const NAMES: [&str; 12] = [
3414 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3415 ];
3416 let lower = name.to_lowercase();
3417 NAMES
3418 .iter()
3419 .position(|n| *n == lower.as_str())
3420 .map(|i| i as i8 + 1)
3421}
3422
3423fn exhausted_review_budget(state: &RunState) -> bool {
3435 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3436}
3437
3438fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3475 unfinished_run_with(runs, short, RunState::load)
3476}
3477
3478fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3481where
3482 F: FnOnce(&str) -> Result<RunState>,
3483{
3484 let id = runs.last()?;
3485 match load(id) {
3486 Ok(s)
3491 if s.status.resumable()
3492 && !s.released()
3493 && !exhausted_review_budget(&s)
3494 && s.liveness(false) != crate::run::Liveness::Live =>
3495 {
3496 Some(id.clone())
3497 }
3498 Ok(_) => None,
3499 Err(e) => {
3500 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3501 None
3502 }
3503 }
3504}
3505
3506fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3514where
3515 F: FnOnce(&str) -> Result<RunState>,
3516{
3517 if task.status != TaskStatus::Failed || task.fresh_start || task.review_branch.is_some() {
3518 return false;
3519 }
3520 let Some(id) = task.runs.last() else {
3521 return false;
3522 };
3523 load(id).is_ok_and(|s| {
3524 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3525 })
3526}
3527
3528#[derive(Debug, Clone, PartialEq, Eq)]
3531enum Starter {
3532 Review(String),
3535 Resume(String),
3537 Start,
3539}
3540
3541fn take_divergence_answer(
3545 branch: &str,
3546 remote: &str,
3547 task: &mut Task,
3548) -> Option<crate::reconcile::Choice> {
3549 let summary = crate::reconcile::summary_for(branch, remote);
3550 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3551 (a.question == summary)
3552 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3553 .flatten()
3554 .map(|c| (i, c))
3555 })?;
3556 task.answers.remove(idx);
3557 Some(choice)
3558}
3559
3560fn choose_starter(
3572 review_branch: Option<&str>,
3573 branch_exists: bool,
3574 unfinished: Option<&str>,
3575) -> Starter {
3576 match review_branch {
3577 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3578 Some(_) => Starter::Start,
3579 None => match unfinished {
3580 Some(id) => Starter::Resume(id.to_owned()),
3581 None => Starter::Start,
3582 },
3583 }
3584}
3585
3586fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3589 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3590 return fallback.to_path_buf();
3591 }
3592 task.repo.clone()
3593}
3594
3595const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3599
3600fn answers_block(task: &Task, count: usize) -> String {
3602 let mut s = ANSWERS_HEADER.to_owned();
3603 for a in &task.answers[..count] {
3604 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3605 }
3606 s
3607}
3608
3609fn append_answers(base: &str, task: &Task) -> String {
3612 if task.answers.is_empty() {
3613 return base.to_owned();
3614 }
3615 let mut s = base.to_owned();
3616 s.push_str(&answers_block(task, task.answers.len()));
3617 s
3618}
3619
3620fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3624 for count in (1..=task.answers.len()).rev() {
3625 let block = answers_block(task, count);
3626 if let Some(base) = instruction.strip_suffix(&block) {
3627 return base;
3628 }
3629 }
3630 instruction
3631}
3632
3633fn instruction_for(task: &Task) -> String {
3641 append_answers(&task.instruction, task)
3642}
3643
3644fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3656 append_answers(strip_answers_block(old_instruction, task), task)
3657}
3658
3659fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3661 let paths = queue.attachment_paths(task);
3662 for (name, path) in task.attachments.iter().zip(&paths) {
3663 if !path.is_file() {
3664 bail!(
3665 "attachment `{name}` is recorded on the task but {} is missing",
3666 path.display()
3667 );
3668 }
3669 }
3670 Ok(paths)
3671}
3672
3673fn prepare_instruction(
3684 starter: &Starter,
3685 old_instruction: Option<&str>,
3686 task: &Task,
3687) -> Option<String> {
3688 match starter {
3689 Starter::Start => Some(instruction_for(task)),
3690 Starter::Resume(_) => Some(resumed_instruction(
3691 old_instruction.expect("a resumed run always has a prior instruction"),
3692 task,
3693 )),
3694 Starter::Review(_) => None,
3695 }
3696}
3697
3698fn record(queue: &Queue, task: &mut Task) {
3702 if let Err(e) = queue.put(task) {
3703 tracing::error!("could not record task {}: {e:#}", task.short());
3704 notices::raise(Notice::error(
3705 "loop:record",
3706 "The loop could not save a task's state; check the disk.",
3707 ));
3708 }
3709}
3710
3711fn runnable(queue: &Queue) -> Vec<Task> {
3717 let mut tasks: Vec<Task> = queue
3718 .list()
3719 .into_iter()
3720 .filter(|t| t.status.runnable())
3721 .collect();
3722 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
3723 tasks
3724}
3725
3726fn describe(state: &RunState) -> String {
3740 let p = phrases(&state.config.graph.language);
3741 let mut detail = if state.status == RunStatus::Stalled {
3742 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
3743 seats.sort_unstable();
3744 seats.dedup();
3745 if seats.is_empty() {
3746 p.quorum_lost.to_owned()
3747 } else {
3748 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
3749 }
3750 } else {
3751 format!("{}{}", p.run_ended, state.status.display_label())
3752 };
3753 if let Some(last) = state.events.last() {
3754 let unanswered = state
3758 .reviews
3759 .last()
3760 .filter(|r| {
3761 state.status == RunStatus::Blocked
3762 && last.node == "review"
3763 && r.incomplete()
3764 && r.blocking == 0
3765 && r.round == state.config.graph.review_rounds
3766 && r.e2e.iter().all(crate::run::CommandOutcome::ok)
3767 })
3768 .map(|r| (r.expected - r.answered, r.round));
3769 match unanswered {
3770 Some((missing, rounds)) if crate::lang::is_japanese(&state.config.graph.language) => {
3771 detail.push_str(&format!(
3772 " ({}: {})",
3773 last.node,
3774 (p.reviewers_never_answered)(missing, rounds)
3775 ));
3776 }
3777 _ => detail.push_str(&format!(" ({}: {})", last.node, last.message)),
3778 }
3779 }
3780 detail.push_str(&format!(" [run {}]", state.id));
3781 detail
3782}
3783
3784const DIAGNOSTIC_MAX: usize = 4_000;
3790
3791const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
3796
3797fn diagnostic(state: &RunState) -> Option<String> {
3811 let mut parts: Vec<String> = Vec::new();
3812
3813 for o in state.gate.iter().filter(|o| !o.ok()) {
3815 parts.push(format!(
3816 "gate `{}` failed ({:?}):\n{}",
3817 o.command,
3818 o.code,
3819 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
3820 ));
3821 }
3822
3823 if let Some(last) = state
3826 .events
3827 .iter()
3828 .rev()
3829 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
3830 {
3831 parts.push(last.message.clone());
3832 }
3833
3834 if state.viable().is_empty() {
3841 for c in &state.candidates {
3842 if let Some(evidence) = &c.verified_noop {
3843 parts.push(format!(
3844 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
3845 c.label
3846 ));
3847 } else if !c.summary.trim().is_empty() {
3848 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
3849 } else if let Some(why) = &c.failed {
3850 parts.push(format!("candidate {}: {why}", c.label));
3851 }
3852 }
3853 }
3854
3855 if parts.is_empty() {
3856 return None;
3857 }
3858 Some(crate::run::tail(
3863 &parts.join("\n\n"),
3864 DIAGNOSTIC_MAX.saturating_sub(100),
3865 ))
3866}
3867
3868fn label(status: RunStatus) -> &'static str {
3876 status.as_str()
3877}
3878
3879fn merge_mode(mode: &str) -> Result<MergeMode> {
3881 match mode {
3882 "none" => Ok(MergeMode::None),
3883 "local" => Ok(MergeMode::Local),
3884 "pr" => Ok(MergeMode::Pr),
3885 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
3886 }
3887}
3888
3889fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
3894 mutex
3895 .lock()
3896 .unwrap_or_else(std::sync::PoisonError::into_inner)
3897}
3898
3899#[cfg(test)]
3900mod tests {
3901 use super::*;
3902 use crate::queue::{Source, TaskStatus};
3903 use crate::run::{Candidate, CommandOutcome};
3904 use pretty_assertions::assert_eq;
3905
3906 fn task() -> Task {
3907 Task::new(
3908 "add retries".to_owned(),
3909 "add retries".to_owned(),
3910 PathBuf::from("/repo"),
3911 Source::Human,
3912 )
3913 }
3914
3915 fn interrupt_task(id: &str) -> Task {
3918 let mut t = task();
3919 t.id = id.to_owned();
3920 t.interrupt = true;
3921 t
3922 }
3923
3924 fn task_with_id(id: &str) -> Task {
3926 let mut t = task();
3927 t.id = id.to_owned();
3928 t
3929 }
3930
3931 fn urgent_task(id: &str) -> Task {
3933 let mut t = task();
3934 t.id = id.to_owned();
3935 t.urgent = true;
3936 t
3937 }
3938
3939 #[test]
3944 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
3945 assert_eq!(permit_kind(true, false), PermitKind::None);
3946 assert_eq!(permit_kind(true, true), PermitKind::None);
3947 }
3948
3949 #[test]
3954 fn permit_kind_separates_urgent_from_ordinary() {
3955 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
3956 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
3957 }
3958
3959 #[test]
3966 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
3967 let cfg = Config::default();
3968 let repo = Path::new("/any/repo/path");
3969
3970 let reason =
3971 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
3972 assert!(reason.contains("1024"), "{reason}");
3973 assert!(
3974 reason.contains(&cfg.disk.min_free_bytes.to_string()),
3975 "{reason}"
3976 );
3977
3978 assert_eq!(
3979 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
3980 None,
3981 "exactly at the floor is open"
3982 );
3983 assert_eq!(
3984 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
3985 None,
3986 "comfortably above the floor is open"
3987 );
3988 }
3989
3990 #[test]
3991 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
3992 let mut cfg = Config::default();
3993 cfg.disk.min_free_bytes = 0;
3994 let repo = Path::new("/any/repo/path");
3995 assert_eq!(
3996 disk_gate_with(repo, &cfg, |_| Ok(0)),
3997 None,
3998 "a zero floor never measures at all"
3999 );
4000 }
4001
4002 #[test]
4003 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
4004 let cfg = Config::default();
4005 let repo = Path::new("/any/repo/path");
4006 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
4007 .expect("a measurement failure must close the gate, not open it");
4008 assert!(reason.contains("could not measure"), "{reason}");
4009 }
4010
4011 #[test]
4012 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
4013 let ordinary = task();
4014 let next = advance_interrupt(
4015 Interrupt::Idle,
4016 std::slice::from_ref(&ordinary.id),
4017 std::slice::from_ref(&ordinary),
4018 );
4019 assert_eq!(next, Interrupt::Idle);
4020 }
4021
4022 #[test]
4023 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
4024 let marked = interrupt_task("marked");
4027 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4028 assert_eq!(next, Interrupt::Idle);
4029 }
4030
4031 #[test]
4032 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
4033 let marked = interrupt_task("marked");
4034 let next = advance_interrupt(
4035 Interrupt::Idle,
4036 &["running".to_owned()],
4037 std::slice::from_ref(&marked),
4038 );
4039 assert_eq!(
4040 next,
4041 Interrupt::Parking {
4042 parked: vec!["running".to_owned()],
4043 interrupt_task: "marked".to_owned(),
4044 }
4045 );
4046 }
4047
4048 #[test]
4057 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
4058 let marked = interrupt_task("marked");
4059
4060 let two = advance_interrupt(
4061 Interrupt::Idle,
4062 &["a".to_owned(), "b".to_owned()],
4063 std::slice::from_ref(&marked),
4064 );
4065 assert_eq!(two, Interrupt::Idle);
4066
4067 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4068 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
4069 }
4070
4071 #[test]
4072 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
4073 let state = Interrupt::Parking {
4074 parked: vec!["running".to_owned()],
4075 interrupt_task: "marked".to_owned(),
4076 };
4077 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
4079 assert_eq!(still_going, state);
4080
4081 let stopped_but_not_yet_dispatched =
4085 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
4086 assert_eq!(stopped_but_not_yet_dispatched, state);
4087
4088 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
4090 assert_eq!(
4091 dispatched,
4092 Interrupt::Running {
4093 parked: vec!["running".to_owned()],
4094 interrupt_task: "marked".to_owned(),
4095 }
4096 );
4097 }
4098
4099 #[test]
4100 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
4101 let state = Interrupt::Running {
4102 parked: vec!["running".to_owned()],
4103 interrupt_task: "marked".to_owned(),
4104 };
4105 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
4106 assert_eq!(still_running, state);
4107
4108 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
4115 assert_eq!(
4116 ended,
4117 Interrupt::Resuming {
4118 parked: vec!["running".to_owned()]
4119 }
4120 );
4121 }
4122
4123 #[test]
4124 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
4125 let state = Interrupt::Resuming {
4126 parked: vec!["running".to_owned()],
4127 };
4128 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
4129 assert_eq!(still_waiting, state);
4130
4131 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
4132 assert_eq!(dispatched, Interrupt::Idle);
4133 }
4134
4135 #[test]
4141 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
4142 {
4143 let state = Interrupt::Parking {
4144 parked: vec!["running".to_owned()],
4145 interrupt_task: "marked".to_owned(),
4146 };
4147 let next = advance_interrupt(state, &[], &[]);
4150 assert_eq!(
4151 next,
4152 Interrupt::Resuming {
4153 parked: vec!["running".to_owned()]
4154 },
4155 "abandoning the interrupt must not abandon the resume it owes"
4156 );
4157 }
4158
4159 #[test]
4162 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
4163 let state = Interrupt::Resuming {
4164 parked: vec!["running".to_owned()],
4165 };
4166 let next = advance_interrupt(state, &[], &[]);
4167 assert_eq!(
4168 next,
4169 Interrupt::Idle,
4170 "nothing is left to wait for; the loop must not stay wedged"
4171 );
4172 }
4173
4174 #[test]
4175 fn disabled_by_config_the_sequence_can_never_leave_idle() {
4176 let marked = interrupt_task("marked");
4177 let next = advance_interrupt_tick(
4178 false,
4179 Interrupt::Idle,
4180 &["running".to_owned()],
4181 std::slice::from_ref(&marked),
4182 );
4183 assert_eq!(
4184 next,
4185 Interrupt::Idle,
4186 "an unmarked, unconfigured daemon must behave exactly as before"
4187 );
4188 }
4189
4190 #[test]
4191 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4192 let state = Interrupt::Parking {
4193 parked: vec!["running".to_owned()],
4194 interrupt_task: "marked".to_owned(),
4195 };
4196 let candidates = vec![interrupt_task("marked"), task()];
4197 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4198 assert!(
4199 allowed.is_empty(),
4200 "nothing may dispatch - not even the interrupt task itself - \
4201 until the parked run has actually stopped"
4202 );
4203 }
4204
4205 #[test]
4219 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4220 for state in [
4221 Interrupt::Parking {
4222 parked: vec!["running".to_owned()],
4223 interrupt_task: "marked".to_owned(),
4224 },
4225 Interrupt::Running {
4226 parked: vec!["running".to_owned()],
4227 interrupt_task: "marked".to_owned(),
4228 },
4229 Interrupt::Resuming {
4230 parked: vec!["running".to_owned()],
4231 },
4232 ] {
4233 let candidates = vec![
4234 interrupt_task("marked"),
4235 urgent_task("hot"),
4236 task_with_id("ordinary"),
4237 ];
4238 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4239 assert!(
4240 !allowed.iter().any(|t| t.id == "hot"),
4241 "an urgent candidate must wait out the same gate as anything \
4242 else while the run it would run alongside has not actually \
4243 left flight, for state {state:?}: {allowed:?}"
4244 );
4245 }
4246 }
4247
4248 #[test]
4254 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4255 let state = Interrupt::Resuming {
4256 parked: vec!["hot".to_owned()],
4257 };
4258 let candidates = vec![urgent_task("hot"), task()];
4259 let allowed = interrupt_gate(&state, &[], candidates);
4260 assert_eq!(
4261 allowed.iter().filter(|t| t.id == "hot").count(),
4262 1,
4263 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4264 );
4265 }
4266
4267 #[test]
4268 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4269 let state = Interrupt::Parking {
4270 parked: vec!["running".to_owned()],
4271 interrupt_task: "marked".to_owned(),
4272 };
4273 let other = task();
4274 let candidates = vec![interrupt_task("marked"), other.clone()];
4275 let allowed = interrupt_gate(&state, &[], candidates);
4276 assert_eq!(allowed.len(), 1);
4277 assert_eq!(allowed[0].id, "marked");
4278 }
4279
4280 #[test]
4281 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4282 let state = Interrupt::Running {
4283 parked: vec!["running".to_owned()],
4284 interrupt_task: "marked".to_owned(),
4285 };
4286 let candidates = vec![task(), task()];
4287 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4288 assert!(allowed.is_empty());
4289 }
4290
4291 #[test]
4298 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4299 let state = Interrupt::Resuming {
4300 parked: vec!["a".to_owned(), "c".to_owned()],
4301 };
4302 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4303 let allowed = interrupt_gate(&state, &[], candidates);
4304 assert_eq!(
4305 allowed.len(),
4306 1,
4307 "at most one candidate may be offered while resuming: {allowed:?}"
4308 );
4309 assert_eq!(allowed[0].id, "a");
4310 }
4311
4312 #[test]
4313 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4314 let state = Interrupt::Resuming {
4315 parked: vec!["a".to_owned()],
4316 };
4317 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4318 assert!(allowed.is_empty());
4319 }
4320
4321 #[test]
4327 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4328 let running = task(); let marked = interrupt_task("marked");
4330
4331 let mut state = Interrupt::Idle;
4332 let in_flight = vec![running.id.clone()];
4334 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4335 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4336 assert!(gated.is_empty(), "still waiting on `running` to park");
4337
4338 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4340 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4341 assert_eq!(
4342 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4343 vec!["marked"],
4344 "only the interrupt task may be offered to the dispatcher now"
4345 );
4346
4347 state = advance_interrupt_tick(
4349 true,
4350 state,
4351 &["marked".to_owned()],
4352 std::slice::from_ref(&running),
4353 );
4354 let gated = interrupt_gate(
4355 &state,
4356 &["marked".to_owned()],
4357 vec![marked.clone(), running.clone()],
4358 );
4359 assert!(
4360 gated.is_empty(),
4361 "the parked run must not be offered back while the interrupt \
4362 task is still running"
4363 );
4364
4365 let other = task_with_id("other");
4369 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4370 assert_eq!(
4371 state,
4372 Interrupt::Resuming {
4373 parked: vec![running.id.clone()]
4374 }
4375 );
4376 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4377 assert_eq!(
4378 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4379 vec![running.id.as_str()],
4380 "exactly the parked run resumes - not the unrelated task, even \
4381 though it was offered first"
4382 );
4383
4384 state = advance_interrupt_tick(
4388 true,
4389 state,
4390 std::slice::from_ref(&running.id),
4391 std::slice::from_ref(&other),
4392 );
4393 assert_eq!(state, Interrupt::Idle);
4394 let gated = interrupt_gate(
4395 &state,
4396 std::slice::from_ref(&running.id),
4397 vec![other.clone()],
4398 );
4399 assert_eq!(
4400 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4401 vec![other.id.as_str()],
4402 "ordinary dispatch is unrestricted again"
4403 );
4404 }
4405
4406 #[test]
4407 fn every_run_status_settles_the_task_it_came_from() {
4408 let table = [
4410 (RunStatus::Merged, TaskStatus::Done, 1),
4411 (RunStatus::Ready, TaskStatus::Done, 1),
4412 (RunStatus::Stalled, TaskStatus::Failed, 0),
4413 (RunStatus::Blocked, TaskStatus::Failed, 1),
4414 (RunStatus::Failed, TaskStatus::Failed, 1),
4415 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4416 (RunStatus::Prep, TaskStatus::Failed, 1),
4417 (RunStatus::Implementing, TaskStatus::Failed, 1),
4418 (RunStatus::Judging, TaskStatus::Failed, 1),
4419 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4420 (RunStatus::Voting, TaskStatus::Failed, 1),
4421 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4422 (RunStatus::Gating, TaskStatus::Failed, 1),
4423 ];
4424 for (run, want, attempts) in table {
4425 let mut t = task();
4426 t.start("20260902-000000-aaaa".to_owned());
4427 settle(
4428 &mut t,
4429 Verdict {
4430 status: run,
4431 left_pr: false,
4432 parked: false,
4433 quota_hit: matches!(run, RunStatus::Stalled),
4434 no_viable_candidates: false,
4435 },
4436 "why",
4437 2,
4438 );
4439 assert_eq!(t.status, want, "task status after {}", label(run));
4440 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4441 }
4442 }
4443
4444 #[test]
4445 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4446 let mut stalled = task();
4447 stalled.start("20260902-000000-aaaa".to_owned());
4448 settle(
4449 &mut stalled,
4450 Verdict {
4451 status: RunStatus::Stalled,
4452 left_pr: false,
4453 parked: false,
4454 quota_hit: true,
4455 no_viable_candidates: false,
4456 },
4457 "quota",
4458 1,
4459 );
4460 assert_eq!(stalled.attempts, 0);
4461 assert!(
4462 stalled.status.runnable(),
4463 "a machine problem must leave the task in line"
4464 );
4465
4466 let mut blocked = task();
4467 blocked.start("20260902-000000-aaaa".to_owned());
4468 settle(
4469 &mut blocked,
4470 Verdict {
4471 status: RunStatus::Blocked,
4472 left_pr: false,
4473 parked: false,
4474 quota_hit: false,
4475 no_viable_candidates: false,
4476 },
4477 "findings open",
4478 1,
4479 );
4480 assert_eq!(blocked.attempts, 1);
4481 assert_eq!(
4482 blocked.status,
4483 TaskStatus::Held,
4484 "the last attempt hands the task to a human"
4485 );
4486 }
4487
4488 #[test]
4489 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4490 let mut delivered = task();
4493 delivered.start("20260903-080619-01c2".to_owned());
4494 settle(
4495 &mut delivered,
4496 Verdict {
4497 status: RunStatus::Blocked,
4498 left_pr: true,
4499 parked: false,
4500 quota_hit: false,
4501 no_viable_candidates: false,
4502 },
4503 "no check status",
4504 4,
4505 );
4506 assert_eq!(
4507 delivered.status,
4508 TaskStatus::Held,
4509 "a pull request waiting on CI or a person is not a retryable failure"
4510 );
4511 assert!(
4512 !delivered.status.runnable(),
4513 "the loop must not pick this task up again"
4514 );
4515 assert_eq!(
4516 delivered.last_error.as_deref(),
4517 Some("no check status"),
4518 "the operator needs to be told what the gate was waiting for"
4519 );
4520
4521 let mut empty_handed = task();
4524 empty_handed.start("20260903-080619-01c2".to_owned());
4525 settle(
4526 &mut empty_handed,
4527 Verdict {
4528 status: RunStatus::Blocked,
4529 left_pr: false,
4530 parked: false,
4531 quota_hit: false,
4532 no_viable_candidates: false,
4533 },
4534 "findings open",
4535 4,
4536 );
4537 assert_eq!(empty_handed.status, TaskStatus::Failed);
4538 assert!(empty_handed.status.runnable());
4539 }
4540
4541 #[test]
4542 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4543 let mut noop = task();
4549 noop.start("20260912-131304-391f".to_owned());
4550 settle(
4551 &mut noop,
4552 Verdict {
4553 status: RunStatus::VerifiedNoop,
4554 left_pr: false,
4555 parked: false,
4556 quota_hit: false,
4557 no_viable_candidates: true,
4558 },
4559 "candidate A: already fixed by b32cfc4, on main",
4560 4,
4561 );
4562 assert_eq!(
4563 noop.status,
4564 TaskStatus::Held,
4565 "an unverified claim is a request for a human, not a failure"
4566 );
4567 assert!(
4568 !noop.status.runnable(),
4569 "the loop must not requeue this on the same unverified claim"
4570 );
4571 assert_eq!(noop.attempts, 1);
4576 }
4577
4578 #[test]
4579 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4580 let mut parked = task();
4585 parked.start("20260903-183634-2d98".to_owned());
4586 settle(
4587 &mut parked,
4588 Verdict {
4589 status: RunStatus::Implementing,
4590 left_pr: false,
4591 quota_hit: false,
4592 parked: true,
4593 no_viable_candidates: false,
4594 },
4595 "parked after `implementing`",
4596 2,
4597 );
4598 assert_eq!(parked.attempts, 0, "a park is refunded");
4599 assert!(
4600 parked.status.runnable(),
4601 "and the task stays in line so the next loop resumes its run"
4602 );
4603 assert_eq!(
4604 parked.last_error.as_deref(),
4605 Some("parked after `implementing`"),
4606 "the card says where it stopped"
4607 );
4608
4609 let mut broken = task();
4613 broken.start("20260903-183634-2d98".to_owned());
4614 settle(
4615 &mut broken,
4616 Verdict {
4617 status: RunStatus::Implementing,
4618 left_pr: false,
4619 quota_hit: false,
4620 parked: false,
4621 no_viable_candidates: false,
4622 },
4623 "returned mid-flight",
4624 2,
4625 );
4626 assert_eq!(broken.attempts, 1);
4627 }
4628
4629 #[test]
4630 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4631 let mut flaky = task();
4636 flaky.start("20260903-123023-e633".to_owned());
4637 settle(
4638 &mut flaky,
4639 Verdict {
4640 status: RunStatus::Stalled,
4641 left_pr: false,
4642 parked: false,
4643 quota_hit: false,
4644 no_viable_candidates: false,
4645 },
4646 "verdict rests on 1 of 3 judges",
4647 2,
4648 );
4649 assert_eq!(
4650 flaky.attempts, 1,
4651 "flakiness spends an attempt, so `max_attempts` still bounds it"
4652 );
4653 assert!(flaky.status.runnable(), "and it is still worth retrying");
4654
4655 let mut limited = task();
4657 limited.start("20260903-123023-e633".to_owned());
4658 settle(
4659 &mut limited,
4660 Verdict {
4661 status: RunStatus::Stalled,
4662 left_pr: false,
4663 parked: false,
4664 quota_hit: true,
4665 no_viable_candidates: false,
4666 },
4667 "judge-2, judge-3 out of quota",
4668 2,
4669 );
4670 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4671 assert!(limited.status.runnable());
4672
4673 let mut worn = task();
4676 for _ in 0..2 {
4677 worn.release();
4678 }
4679 worn.start("20260903-123023-e633".to_owned());
4680 worn.attempts = 2;
4681 settle(
4682 &mut worn,
4683 Verdict {
4684 status: RunStatus::Stalled,
4685 left_pr: false,
4686 parked: false,
4687 quota_hit: false,
4688 no_viable_candidates: false,
4689 },
4690 "no quorum again",
4691 2,
4692 );
4693 assert_eq!(worn.status, TaskStatus::Held);
4694 assert!(!worn.status.runnable());
4695 }
4696
4697 #[test]
4698 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4699 let mut wiped_out = task();
4706 wiped_out.start("20260907-025000-a1b2".to_owned());
4707 settle(
4708 &mut wiped_out,
4709 Verdict {
4710 status: RunStatus::Failed,
4711 left_pr: false,
4712 parked: false,
4713 quota_hit: true,
4714 no_viable_candidates: true,
4715 },
4716 "no candidate produced a change; nothing to judge",
4717 2,
4718 );
4719 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
4720 assert!(
4721 wiped_out.status.runnable(),
4722 "a machine problem must leave the task in line"
4723 );
4724
4725 let mut partial_progress = task();
4731 partial_progress.start("20260907-025500-c3d4".to_owned());
4732 settle(
4733 &mut partial_progress,
4734 Verdict {
4735 status: RunStatus::Failed,
4736 left_pr: false,
4737 parked: false,
4738 quota_hit: true,
4739 no_viable_candidates: false,
4740 },
4741 "gate failed on the winning candidate",
4742 2,
4743 );
4744 assert_eq!(
4745 partial_progress.attempts, 1,
4746 "a candidate that actually produced a change spends the attempt \
4747 even though some other seat hit its quota"
4748 );
4749 assert!(partial_progress.status.runnable());
4750 }
4751
4752 #[test]
4753 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
4754 let mut t = task();
4760 t.start("20260907-025000-a1b2".to_owned());
4761 let mut state = run_state(RunStatus::Failed);
4762 state.quota.push(QuotaLoss {
4763 seat: "cand-a".to_owned(),
4764 node: "implement".to_owned(),
4765 at: Timestamp::now(),
4766 reset: None,
4767 });
4768 assert!(
4769 state.viable().is_empty(),
4770 "no candidate was added, so nothing is viable"
4771 );
4772 reclaim(&mut t, Some(state), 2, "en");
4773 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
4774 assert!(t.status.runnable());
4775 }
4776
4777 #[test]
4778 fn a_held_task_is_never_offered_to_the_loop() {
4779 let dir = tempfile::tempdir().unwrap();
4780 let queue = Queue::at(dir.path().to_path_buf());
4781 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
4782 let mut t = task();
4783 t.id = format!("2026090{n}-000000-000{n}");
4784 t.priority = priority;
4785 queue.put(&mut t).unwrap();
4786 }
4787 let mut held = task();
4788 held.id = "20260909-000000-9999".to_owned();
4789 held.priority = 99;
4790 held.hold_machine(None);
4791 queue.put(&mut held).unwrap();
4792
4793 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
4794 assert_eq!(order.len(), 3);
4795 assert!(!order.contains(&held.id));
4796 assert_eq!(
4797 order.first().cloned(),
4798 queue.next_runnable().map(|t| t.id),
4799 "the loop's first candidate is exactly what the queue offers"
4800 );
4801 assert_eq!(
4802 order,
4803 vec![
4804 "20260902-000000-0002".to_owned(),
4805 "20260903-000000-0003".to_owned(),
4806 "20260901-000000-0001".to_owned(),
4807 ],
4808 "priority first, then oldest, so nothing starves"
4809 );
4810 }
4811
4812 #[test]
4813 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
4814 let dir = tempfile::tempdir().unwrap();
4815 let queue = Queue::at(dir.path().to_path_buf());
4816 let mut old = task();
4817 old.id = "20260101-000000-old0".to_owned();
4818 queue.put(&mut old).unwrap();
4819 let mut fresh = task();
4820 fresh.id = "20260101-000000-new0".to_owned();
4821 queue.put(&mut fresh).unwrap();
4822
4823 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
4827 std::thread::sleep(Duration::from_millis(60));
4828 let live = queue.claim(&fresh.id).unwrap();
4829
4830 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4831 assert_eq!(swept, vec![old.id.clone()]);
4832 assert!(
4833 queue.claim(&old.id).is_ok(),
4834 "an unparseable lock older than the threshold is swept"
4835 );
4836 assert!(
4837 queue.claim(&fresh.id).is_err(),
4838 "a live pid protects its lock regardless of age"
4839 );
4840 drop(live);
4841 }
4842
4843 #[test]
4844 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
4845 let dir = tempfile::tempdir().unwrap();
4855 let queue = Queue::at(dir.path().to_path_buf());
4856 let mut t = task();
4857 t.id = "20260101-000000-live".to_owned();
4858 queue.put(&mut t).unwrap();
4859
4860 let claim = queue.claim(&t.id).unwrap();
4861 std::thread::sleep(Duration::from_millis(60));
4862
4863 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4864 assert!(
4865 swept.is_empty(),
4866 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
4867 );
4868 assert!(
4869 queue.claim(&t.id).is_err(),
4870 "the lock still protects its task"
4871 );
4872 drop(claim);
4873 }
4874
4875 fn injected_dead_pid() -> u32 {
4878 std::process::id().checked_add(1).unwrap_or(1)
4879 }
4880
4881 #[test]
4882 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
4883 let dir = tempfile::tempdir().unwrap();
4884 let queue = Queue::at(dir.path().to_path_buf());
4885 let mut t = task();
4886 t.id = "20260101-000000-dead".to_owned();
4887 queue.put(&mut t).unwrap();
4888 let dead_pid = injected_dead_pid();
4889
4890 std::fs::write(
4895 dir.path().join(format!("{}.lock", t.id)),
4896 dead_pid.to_string(),
4897 )
4898 .unwrap();
4899
4900 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4901 pid != dead_pid
4902 });
4903 assert_eq!(
4904 swept,
4905 vec![t.id.clone()],
4906 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
4907 );
4908 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
4909 }
4910
4911 #[test]
4912 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
4913 let dir = tempfile::tempdir().unwrap();
4914 let queue = Queue::at(dir.path().to_path_buf());
4915 let mut t = task();
4916 t.id = "20260101-000000-late".to_owned();
4917 queue.put(&mut t).unwrap();
4918 let dead_pid = injected_dead_pid();
4919
4920 assert!(
4923 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
4924 "nothing has claimed the task yet"
4925 );
4926
4927 std::fs::write(
4930 dir.path().join(format!("{}.lock", t.id)),
4931 dead_pid.to_string(),
4932 )
4933 .unwrap();
4934
4935 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4939 pid != dead_pid
4940 });
4941 assert_eq!(swept, vec![t.id.clone()]);
4942 }
4943
4944 #[test]
4945 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
4946 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
4951 let dir = tempfile::tempdir().unwrap();
4952 let queue = Queue::at(dir.path().to_path_buf());
4953 let mut t = task();
4954 t.id = "20260101-000000-crsh".to_owned();
4955 t.status = TaskStatus::Running;
4956 t.attempts = 1;
4957 t.runs.push("20260904-000000-4043".to_owned());
4961 queue.put(&mut t).unwrap();
4962 let dead_pid = injected_dead_pid();
4963
4964 std::fs::write(
4967 dir.path().join(format!("{}.lock", t.id)),
4968 dead_pid.to_string(),
4969 )
4970 .unwrap();
4971
4972 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
4978 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
4979
4980 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4981 pid != dead_pid
4982 });
4983 assert_eq!(swept, vec![t.id.clone()]);
4984
4985 let reclaimed = reclaim_orphaned_running(&queue, 2);
4986 assert_eq!(reclaimed, vec![t.id.clone()]);
4987 let after = queue.get(&t.id).unwrap();
4988 assert_eq!(
4989 after.status,
4990 TaskStatus::Held,
4991 "no run.json to recover from, so a human is asked"
4992 );
4993 assert_eq!(
4994 after.runs,
4995 vec!["20260904-000000-4043".to_owned()],
4996 "the crashed run's id is kept as evidence, not discarded"
4997 );
4998 }
4999
5000 #[test]
5001 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
5002 let dir = tempfile::tempdir().unwrap();
5003 let queue = Queue::at(dir.path().to_path_buf());
5004 let mut t = task();
5005 t.id = "20260101-000000-unknown".to_owned();
5006 queue.put(&mut t).unwrap();
5007 let dead_pid = injected_dead_pid();
5008 std::fs::write(
5009 dir.path().join(format!("{}.lock", t.id)),
5010 dead_pid.to_string(),
5011 )
5012 .unwrap();
5013
5014 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
5015 assert!(swept.is_empty(), "an unknown pid must keep its lock");
5016 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
5017 }
5018
5019 fn run_state_in(status: RunStatus, language: &str) -> RunState {
5020 let mut s = run_state(status);
5021 s.config.graph.language = language.to_owned();
5022 s
5023 }
5024
5025 fn unstarted_verdict(status: RunStatus) -> Verdict {
5026 Verdict {
5027 status,
5028 left_pr: false,
5029 quota_hit: false,
5030 parked: false,
5031 no_viable_candidates: false,
5032 }
5033 }
5034
5035 #[test]
5036 fn an_already_in_base_run_finishes_the_task_without_spending_an_attempt() {
5037 let mut t = task();
5038 t.attempts = 1;
5039 settle_in(
5040 &mut t,
5041 unstarted_verdict(RunStatus::AlreadyInBase),
5042 "already in main",
5043 1,
5044 phrases("en"),
5045 );
5046 assert_eq!(t.status, TaskStatus::Done);
5047 assert_eq!(t.attempts, 0);
5048 }
5049
5050 #[test]
5051 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
5052 let reason = |language: &str| {
5053 let mut t = task();
5054 settle_in(
5055 &mut t,
5056 unstarted_verdict(RunStatus::Judging),
5057 "boom",
5058 1,
5059 phrases(language),
5060 );
5061 t.last_error.or(t.hold_reason).unwrap_or_default()
5062 };
5063 assert!(
5064 reason("en").starts_with("the graph stopped at `"),
5065 "{}",
5066 reason("en")
5067 );
5068 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
5069 assert!(reason("日本語").contains("boom"));
5070 assert_eq!(reason("fr"), reason("en"));
5071 }
5072
5073 #[test]
5074 fn describe_follows_the_run_language_and_keeps_the_run_id() {
5075 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
5076 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
5077 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
5078 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
5079 assert!(ja.contains("[run "), "{ja}");
5080 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
5081 assert!(ended.starts_with("run 終了: "), "{ended}");
5082 let mut de = run_state_in(RunStatus::Failed, "de");
5083 let mut en = run_state_in(RunStatus::Failed, "en");
5084 de.id = "same".to_owned();
5085 en.id = "same".to_owned();
5086 assert_eq!(describe(&de), describe(&en));
5087 }
5088
5089 #[test]
5090 fn handover_refusals_follow_the_language_and_keep_the_detail_apart() {
5091 use crate::handover::Refused;
5092 let cases = [
5093 Refused::Foreign {
5094 branch: "b".into(),
5095 path: "/w/x".into(),
5096 why: "made by hand".into(),
5097 },
5098 Refused::Unsafe {
5099 branch: "b".into(),
5100 path: "/w/x".into(),
5101 why: "its worktree has uncommitted changes (a.rs)".into(),
5102 },
5103 Refused::ReleaseFailed {
5104 branch: "b".into(),
5105 path: "/w/x".into(),
5106 run: "ab12".into(),
5107 },
5108 ];
5109 for r in &cases {
5110 let en = (phrases("en").handover_refused)(r);
5111 assert_eq!(en, r.to_string());
5112 let ja = (phrases("ja").handover_refused)(r);
5113 assert!(
5114 !ja.contains("is checked out") && !ja.contains("try again"),
5115 "{ja}"
5116 );
5117 assert!(ja.contains("`b`") && ja.contains("/w/x"), "{ja}");
5118 if let Refused::Foreign { why, .. } | Refused::Unsafe { why, .. } = r {
5119 assert!(ja.contains(&format!("(詳細: {why})")), "{ja}");
5120 }
5121 }
5122 assert!(phrases("en").handover_hint.contains("release the task"));
5123 assert!(phrases("ja").handover_hint.contains("解放"));
5124 }
5125
5126 #[test]
5127 fn describe_translates_the_unanswered_reviewer_stop_only_in_ja() {
5128 let build = |lang: &str| {
5129 let mut s = run_state_in(RunStatus::Blocked, lang);
5130 s.id = "same".to_owned();
5131 s.config.graph.review_rounds = 3;
5132 let mut r = review_round(3);
5133 r.expected = 3;
5134 r.answered = 1;
5135 s.reviews.push(r);
5136 s.event(
5137 "review",
5138 "2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean",
5139 );
5140 s
5141 };
5142 let en = describe(&build("en"));
5143 assert!(en.contains("(review: 2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean)"), "{en}");
5144 let ja = describe(&build("ja"));
5145 assert!(ja.contains("2 席のレビュアーが 3 ラウンド"), "{ja}");
5146 assert!(!ja.contains("never answered"), "{ja}");
5147 assert!(ja.contains("[run "), "{ja}");
5148
5149 let mut failed = build("ja");
5152 failed.reviews[0].e2e.push(crate::run::CommandOutcome {
5153 command: "cargo test".to_owned(),
5154 code: Some(1),
5155 output_tail: String::new(),
5156 duration_ms: 0,
5157 resource_blocked: false,
5158 });
5159 failed.event("review", "stopped; e2e failed: cargo test");
5160 let ja = describe(&failed);
5161 assert!(ja.contains("stopped; e2e failed: cargo test"), "{ja}");
5162 assert!(!ja.contains("席のレビュアー"), "{ja}");
5163 }
5164
5165 #[test]
5166 fn refusals_and_recovery_prose_follow_the_language() {
5167 let t = held_task_with("r1");
5168 let q = action_question(
5169 "r1",
5170 ask::ChoiceAction::Resume {
5171 run: "r1".to_owned(),
5172 },
5173 );
5174 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
5175 ActionDecision::Refuse(s) => s,
5176 other => panic!("{other:?}"),
5177 };
5178 assert!(refuse(phrases("en")).contains("could not be read: gone"));
5179 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
5180 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
5181
5182 let mut held = task();
5183 reclaim(&mut held, None, 2, "ja");
5184 assert!(held.hold_reason.unwrap().contains("保留にしました"));
5185 let mut held = task();
5186 reclaim(&mut held, None, 2, "xx");
5187 assert!(held.hold_reason.unwrap().contains("held for a human"));
5188 }
5189
5190 fn run_state(status: RunStatus) -> RunState {
5191 let mut state = RunState::new(
5192 PathBuf::from("/repo"),
5193 "main".to_owned(),
5194 "abc1234def".to_owned(),
5195 "add retries".to_owned(),
5196 Config::default(),
5197 );
5198 state.status = status;
5199 state
5200 }
5201
5202 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
5203 Candidate {
5204 index: 0,
5205 label,
5206 agent: "claude".to_owned(),
5207 branch: format!("magi/x/{label}"),
5208 worktree: PathBuf::from("/repo"),
5209 summary: summary.to_owned(),
5210 stat: String::new(),
5211 files: 0,
5212 commits: usize::from(!empty),
5213 empty,
5214 failed: failed.map(str::to_owned),
5215 verified_noop: None,
5216 duration_ms: 0,
5217 folded: false,
5218 }
5219 }
5220
5221 #[test]
5222 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
5223 let mut state = run_state(RunStatus::Blocked);
5224 state.gate = vec![
5225 CommandOutcome {
5226 command: "cargo make check".to_owned(),
5227 code: Some(0),
5228 output_tail: "ok".to_owned(),
5229 duration_ms: 0,
5230 resource_blocked: false,
5231 },
5232 CommandOutcome {
5233 command: "cargo test".to_owned(),
5234 code: Some(101),
5235 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
5236 duration_ms: 0,
5237 resource_blocked: false,
5238 },
5239 ];
5240 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
5241 assert!(d.contains("cargo test"), "{d}");
5242 assert!(
5243 !d.contains("cargo make check"),
5244 "a passing check is not a diagnostic: {d}"
5245 );
5246 assert!(d.contains("assertion failed"), "{d}");
5247 }
5248
5249 #[test]
5250 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
5251 let mut state = run_state(RunStatus::Blocked);
5252 state.event(
5253 "land",
5254 "stopped: the fixer produced no commit while 2 check(s) were failing \
5255 (build, lint); stopping instead of looping on an unchanged tree",
5256 );
5257 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
5258 assert!(d.contains("build"), "{d}");
5259 assert!(d.contains("lint"), "{d}");
5260 assert!(d.contains("fixer produced no commit"), "{d}");
5261 }
5262
5263 #[test]
5264 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
5265 let state = run_state(RunStatus::VerifiedNoop);
5270 let d = describe(&state);
5271 assert!(
5272 d.contains("agent-verified no-op"),
5273 "expected the display label, not the wire spelling: {d}"
5274 );
5275 assert!(!d.contains("verified_noop"), "{d}");
5276 }
5277
5278 #[test]
5279 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5280 let mut state = run_state(RunStatus::Failed);
5286 state.candidates = vec![candidate(
5287 'A',
5288 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5289 true,
5290 None,
5291 )];
5292 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5293 assert!(d.contains("candidate A"), "{d}");
5294 assert!(d.contains("tagged v1.2.3"), "{d}");
5295 }
5296
5297 #[test]
5298 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5299 let mut state = run_state(RunStatus::Failed);
5300 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5301 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5302 assert!(d.contains("candidate A"), "{d}");
5303 assert!(d.contains("agent timed out"), "{d}");
5304 }
5305
5306 #[test]
5307 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5308 let mut state = run_state(RunStatus::Failed);
5311 state.candidates = vec![candidate('A', "did the work", false, None)];
5312 assert!(diagnostic(&state).is_none());
5313 }
5314
5315 #[test]
5316 fn diagnostic_is_bounded_however_much_a_run_printed() {
5317 let mut state = run_state(RunStatus::Blocked);
5318 state.gate = vec![
5319 CommandOutcome {
5320 command: "cargo test".to_owned(),
5321 code: Some(101),
5322 output_tail: "x".repeat(50_000),
5323 duration_ms: 0,
5324 resource_blocked: false,
5325 },
5326 CommandOutcome {
5327 command: "cargo clippy".to_owned(),
5328 code: Some(1),
5329 output_tail: "y".repeat(50_000),
5330 duration_ms: 0,
5331 resource_blocked: false,
5332 },
5333 ];
5334 state.candidates = vec![
5335 candidate('A', &"z".repeat(50_000), true, None),
5336 candidate('B', &"w".repeat(50_000), true, None),
5337 ];
5338 let d = diagnostic(&state).expect("plenty here to diagnose");
5339 assert!(
5340 d.len() <= DIAGNOSTIC_MAX,
5341 "diagnostic grew to {} bytes, unbounded",
5342 d.len()
5343 );
5344 }
5345
5346 #[test]
5347 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5348 let mut state = run_state(RunStatus::Blocked);
5349 state.gate = vec![CommandOutcome {
5350 command: "cargo test".to_owned(),
5351 code: Some(101),
5352 output_tail: "assertion failed".to_owned(),
5353 duration_ms: 0,
5354 resource_blocked: false,
5355 }];
5356 let verdict = Verdict {
5357 status: RunStatus::Blocked,
5358 left_pr: false,
5359 quota_hit: false,
5360 parked: false,
5361 no_viable_candidates: false,
5362 };
5363
5364 let mut t = task();
5367 t.start("run-1".to_owned());
5368 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5369 assert_eq!(t.status, TaskStatus::Failed);
5370 assert!(t.diagnostic.is_none());
5371
5372 t.start("run-2".to_owned());
5375 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5376 assert_eq!(t.status, TaskStatus::Held);
5377 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5378 assert!(d.contains("cargo test"), "{d}");
5379 }
5380
5381 #[test]
5382 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5383 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5388 let home = crate::run::home();
5389 let state = run_state(RunStatus::VerifiedNoop);
5390 let mut q = ask::Question::new(
5391 state.id.clone(),
5392 "implement".to_owned(),
5393 "impl-A".to_owned(),
5394 "is this really a no-op?".to_owned(),
5395 String::new(),
5396 Vec::new(),
5397 );
5398 Questions::at(home.join("questions")).put(&mut q).unwrap();
5399
5400 let verdict = Verdict {
5401 status: RunStatus::VerifiedNoop,
5402 left_pr: false,
5403 quota_hit: false,
5404 parked: false,
5405 no_viable_candidates: false,
5406 };
5407 let mut t = task();
5408 t.start(state.id.clone());
5409 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5410
5411 assert_eq!(t.status, TaskStatus::Held);
5412 let reason = t.hold_reason.expect("a held task must record why");
5413 assert!(
5414 reason.starts_with("run ended agent-verified no-op"),
5415 "the original settle reason must survive unchanged: {reason}"
5416 );
5417 assert!(
5418 reason.contains(q.short()),
5419 "the open question's id must be named so the notice is actionable: {reason}"
5420 );
5421 }
5422
5423 #[test]
5424 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5425 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5426 let state = run_state(RunStatus::VerifiedNoop);
5427
5428 let verdict = Verdict {
5429 status: RunStatus::VerifiedNoop,
5430 left_pr: false,
5431 quota_hit: false,
5432 parked: false,
5433 no_viable_candidates: false,
5434 };
5435 let mut t = task();
5436 t.start(state.id.clone());
5437 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5438
5439 assert_eq!(t.status, TaskStatus::Held);
5440 assert_eq!(
5441 t.hold_reason.as_deref(),
5442 Some("run ended agent-verified no-op"),
5443 "nothing to append when the question was already answered or never asked"
5444 );
5445 }
5446
5447 #[test]
5448 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5449 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5450 let mut first = run_state(RunStatus::Blocked);
5451 first.id = "20260101-000000-sup1".to_owned();
5452 first.save().unwrap();
5453 let mut second = run_state(RunStatus::Merged);
5454 second.id = "20260101-000000-sup2".to_owned();
5455 second.save().unwrap();
5456
5457 let mut t = task();
5458 t.runs = vec![first.id.clone(), second.id.clone()];
5459 t.status = TaskStatus::Done;
5460
5461 supersede_prior_runs(&t, &crate::run::home());
5462
5463 assert_eq!(
5464 RunState::load(&first.id).unwrap().status,
5465 RunStatus::Superseded,
5466 "the first attempt's Blocked no longer needs anyone's attention"
5467 );
5468 assert_eq!(
5469 RunState::load(&second.id).unwrap().status,
5470 RunStatus::Merged,
5471 "the run that actually succeeded is left exactly as it was"
5472 );
5473 }
5474
5475 #[test]
5476 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5477 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5483 let mut first = run_state(RunStatus::Blocked);
5484 first.id = "20260101-000000-sup9".to_owned();
5485 first.driver_pid = Some(std::process::id());
5488 first.driver_started_at = Some(
5489 crate::proc::process_started_at(std::process::id())
5490 .expect("this test process's own start time must be queryable"),
5491 );
5492 first.save().unwrap();
5493 let mut second = run_state(RunStatus::Merged);
5494 second.id = "20260101-000000-supa".to_owned();
5495 second.save().unwrap();
5496
5497 let mut t = task();
5498 t.runs = vec![first.id.clone(), second.id.clone()];
5499 t.status = TaskStatus::Done;
5500
5501 supersede_prior_runs(&t, &crate::run::home());
5502
5503 assert_eq!(
5504 RunState::load(&first.id).unwrap().status,
5505 RunStatus::Blocked,
5506 "a live driver_pid means something is still actually working this run, \
5507 even though no daemon claims it - rewriting under it would just be \
5508 undone the next time that process saves"
5509 );
5510 }
5511
5512 #[test]
5513 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5514 let dir = tempfile::tempdir().unwrap();
5520 let home = dir.path().to_path_buf();
5521 let queue = Queue::at(dir.path().join("queue"));
5522
5523 let mut first = run_state(RunStatus::Blocked);
5524 first.id = "20260101-000000-supd".to_owned();
5525 first.driver_pid = Some(std::process::id());
5526 first.driver_started_at = Some(
5527 crate::proc::process_started_at(std::process::id())
5528 .expect("this test process's own start time must be queryable"),
5529 );
5530 first.save_under(&home).unwrap();
5531 let mut second = run_state(RunStatus::Merged);
5532 second.id = "20260101-000000-supe".to_owned();
5533 second.save_under(&home).unwrap();
5534
5535 let mut t = task();
5536 t.runs = vec![first.id.clone(), second.id.clone()];
5537 t.status = TaskStatus::Done;
5538 queue.put(&mut t).unwrap();
5539
5540 resweep_superseded_attempts(&queue, &home);
5541 assert_eq!(
5542 RunState::load_under(&first.id, &home).unwrap().status,
5543 RunStatus::Blocked,
5544 "still live on the first pass, so still untouched"
5545 );
5546
5547 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5553 stale.driver_started_at = Some("1".to_owned());
5554 stale.save_under(&home).unwrap();
5555
5556 resweep_superseded_attempts(&queue, &home);
5557 assert_eq!(
5558 RunState::load_under(&first.id, &home).unwrap().status,
5559 RunStatus::Superseded,
5560 "the second pass catches up what the first one correctly skipped"
5561 );
5562 }
5563
5564 #[test]
5565 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5566 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5567 let mut first = run_state(RunStatus::Blocked);
5568 first.id = "20260101-000000-sup3".to_owned();
5569 first.save().unwrap();
5570 let mut second = run_state(RunStatus::Blocked);
5571 second.id = "20260101-000000-sup4".to_owned();
5572 second.save().unwrap();
5573
5574 let mut t = task();
5575 t.runs = vec![first.id.clone(), second.id.clone()];
5576 t.status = TaskStatus::Failed;
5580
5581 supersede_prior_runs(&t, &crate::run::home());
5582
5583 assert_eq!(
5584 RunState::load(&first.id).unwrap().status,
5585 RunStatus::Blocked
5586 );
5587 assert_eq!(
5588 RunState::load(&second.id).unwrap().status,
5589 RunStatus::Blocked
5590 );
5591 }
5592
5593 #[test]
5594 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5595 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5599 let mut first = run_state(RunStatus::Blocked);
5600 first.id = "20260101-000000-sup5".to_owned();
5601 first.save().unwrap();
5602
5603 let mut t = task();
5604 t.runs = vec![first.id.clone()];
5605 t.status = TaskStatus::Done;
5606
5607 supersede_prior_runs(&t, &crate::run::home());
5608
5609 assert_eq!(
5610 RunState::load(&first.id).unwrap().status,
5611 RunStatus::Blocked,
5612 "a single-attempt task has no earlier run to supersede"
5613 );
5614 }
5615
5616 #[test]
5617 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5618 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5625 let mut first = run_state(RunStatus::Blocked);
5626 first.id = "20260101-000000-supb".to_owned();
5627 first.save().unwrap();
5628 let mut second = run_state(RunStatus::Failed);
5629 second.id = "20260101-000000-supc".to_owned();
5630 second.save().unwrap();
5631
5632 let mut t = task();
5633 t.runs = vec![first.id.clone(), second.id.clone()];
5634 t.status = TaskStatus::Done;
5635
5636 supersede_prior_runs(&t, &crate::run::home());
5637
5638 assert_eq!(
5639 RunState::load(&first.id).unwrap().status,
5640 RunStatus::Blocked,
5641 "the task's last attempt never landed, so there is nothing here \
5642 actually superseding it"
5643 );
5644 }
5645
5646 #[test]
5647 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5648 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5652 let mut failed = run_state(RunStatus::Failed);
5653 failed.id = "20260101-000000-sup6".to_owned();
5654 failed.save().unwrap();
5655 let mut noop = run_state(RunStatus::VerifiedNoop);
5656 noop.id = "20260101-000000-sup7".to_owned();
5657 noop.save().unwrap();
5658 let mut winner = run_state(RunStatus::Ready);
5659 winner.id = "20260101-000000-sup8".to_owned();
5660 winner.save().unwrap();
5661
5662 let mut t = task();
5663 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5664 t.status = TaskStatus::Done;
5665
5666 supersede_prior_runs(&t, &crate::run::home());
5667
5668 assert_eq!(
5669 RunState::load(&failed.id).unwrap().status,
5670 RunStatus::Failed
5671 );
5672 assert_eq!(
5673 RunState::load(&noop.id).unwrap().status,
5674 RunStatus::VerifiedNoop
5675 );
5676 }
5677
5678 fn approval_question(run: &str) -> ask::Question {
5679 ask::Question::new(
5680 run.to_owned(),
5681 land::APPROVAL_NODE.to_owned(),
5682 "land".to_owned(),
5683 "merge?".to_owned(),
5684 String::new(),
5685 vec!["merge".to_owned(), "hold".to_owned()],
5686 )
5687 }
5688
5689 #[test]
5690 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5691 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5692 let mut state = run_state(RunStatus::Landing);
5693 state.id = "20260101-000000-fre1".to_owned();
5694 state.parked = true;
5695 state.save().unwrap();
5696 ask::Questions::open()
5697 .put(&mut approval_question(&state.id))
5698 .unwrap();
5699
5700 let mut t = task();
5701 t.runs.push(state.id.clone());
5702 assert_eq!(
5703 land_resume_state(&t),
5704 LandResume::StillWaiting,
5705 "nobody has answered and the timeout has not passed"
5706 );
5707 }
5708
5709 #[test]
5710 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5711 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5716 let mut state = run_state(RunStatus::Landing);
5717 state.id = "20260101-000000-exp1".to_owned();
5718 state.parked = true;
5719 state.config.graph.answer_timeout = 60;
5720 state.save().unwrap();
5721
5722 let store = ask::Questions::open();
5723 let mut q = approval_question(&state.id);
5724 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
5725 store.put(&mut q).unwrap();
5726
5727 let mut t = task();
5728 t.runs.push(state.id.clone());
5729 assert_eq!(
5730 land_resume_state(&t),
5731 LandResume::Ready,
5732 "an expired question must not be waited on forever"
5733 );
5734
5735 let after = store.get(&q.id).unwrap();
5736 assert!(
5737 !after.status.open(),
5738 "the question is abandoned, not silently ignored"
5739 );
5740 assert!(
5741 after.resolution().is_none(),
5742 "an abandoned question is not read as a decision"
5743 );
5744 }
5745
5746 #[test]
5747 fn reclaim_settles_a_running_task_against_its_last_run() {
5748 let mut t = task();
5749 t.start("20260904-000000-4043".to_owned());
5750 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
5751 assert_eq!(
5752 t.status,
5753 TaskStatus::Done,
5754 "a run that actually finished must not stay `running` forever"
5755 );
5756 }
5757
5758 #[test]
5759 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
5760 let mut t = task();
5764 t.start("20260904-000000-4043".to_owned());
5765 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
5766 assert_eq!(t.status, TaskStatus::Failed);
5767 assert!(t.status.runnable());
5768 }
5769
5770 #[test]
5771 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
5772 let mut t = task();
5773 t.start("20260904-000000-4043".to_owned());
5774 reclaim(&mut t, None, 2, "en");
5775 assert_eq!(t.status, TaskStatus::Held);
5776 assert!(
5777 t.last_error
5778 .as_deref()
5779 .is_some_and(|e| e.contains("running")),
5780 "the operator needs to know why this task was held"
5781 );
5782 }
5783
5784 #[test]
5785 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
5786 let dir = tempfile::tempdir().unwrap();
5787 let queue = Queue::at(dir.path().to_path_buf());
5788
5789 let mut orphaned = task();
5791 orphaned.id = "20260904-000000-orph".to_owned();
5792 orphaned.status = TaskStatus::Running;
5793 orphaned.attempts = 1;
5794 queue.put(&mut orphaned).unwrap();
5795
5796 let mut alive = task();
5797 alive.id = "20260904-000000-live".to_owned();
5798 alive.status = TaskStatus::Running;
5799 alive.attempts = 1;
5800 queue.put(&mut alive).unwrap();
5801 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
5802
5803 let mut queued = task();
5804 queued.id = "20260904-000000-wait".to_owned();
5805 queue.put(&mut queued).unwrap();
5806
5807 let reclaimed = reclaim_orphaned_running(&queue, 2);
5808 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
5809
5810 assert_eq!(
5811 queue.get(&orphaned.id).unwrap().status,
5812 TaskStatus::Held,
5813 "nothing was driving it and there was no run to recover"
5814 );
5815 assert_eq!(
5816 queue.get(&alive.id).unwrap().status,
5817 TaskStatus::Running,
5818 "a live claim must protect the task it belongs to"
5819 );
5820 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
5821 }
5822
5823 fn read_run_under(home: &Path, id: &str) -> RunState {
5829 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
5830 serde_json::from_str(&body).unwrap()
5831 }
5832
5833 #[test]
5834 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
5835 let dir = tempfile::tempdir().unwrap();
5836 let home = dir.path().to_path_buf();
5837 let now = Timestamp::now();
5838 let overrun_seat = || crate::run::ActiveSeat {
5839 node: "implement".to_owned(),
5840 started_at: now - jiff::SignedDuration::new(21_000, 0),
5841 timeout_secs: 3_600,
5842 attempt: 0,
5843 task: None,
5844 command: None,
5845 index: None,
5846 total: None,
5847 };
5848
5849 let mut dead = run_state(RunStatus::Implementing);
5850 dead.id = "20260101-000000-dead".to_owned();
5851 dead.active.insert("impl-A".to_owned(), overrun_seat());
5852 dead.driver_pid = Some(4242);
5855 dead.save_under(&home).unwrap();
5856
5857 let mut alive = run_state(RunStatus::Implementing);
5860 alive.id = "20260101-000000-aliv".to_owned();
5861 alive.active.insert("impl-A".to_owned(), overrun_seat());
5862 alive.save_under(&home).unwrap();
5863 let mut status = Status::new();
5864 status.current = vec![Current {
5865 task: "20260101-000000-task".to_owned(),
5866 run: alive.id.clone(),
5867 }];
5868 write_status_to(&home.join("daemon.json"), &status).unwrap();
5869
5870 let questions = Questions::at(home.join("questions"));
5874 let mut q = ask::Question::new(
5875 dead.id.clone(),
5876 "implement".to_owned(),
5877 "impl-A".to_owned(),
5878 "Which storage backend?".to_owned(),
5879 String::new(),
5880 vec!["SQLite".to_owned(), "Redis".to_owned()],
5881 );
5882 questions.put(&mut q).unwrap();
5883
5884 let abandoned = reclaim_abandoned_runs_with(
5885 &home,
5886 now,
5887 |pid| if pid == 4242 { Some(false) } else { None },
5888 |_| panic!("a query answering Dead outright needs no identity corroboration"),
5889 );
5890 assert_eq!(abandoned, vec![dead.id.clone()]);
5891
5892 let reloaded = read_run_under(&home, &dead.id);
5893 assert_eq!(reloaded.status, RunStatus::Failed);
5894 assert!(reloaded.active.is_empty());
5895 assert!(
5896 !questions.get(&q.id).unwrap().status.open(),
5897 "the failed run's own open question must be settled in the same pass"
5898 );
5899
5900 let still_alive = read_run_under(&home, &alive.id);
5901 assert_eq!(
5902 still_alive.status,
5903 RunStatus::Implementing,
5904 "a live daemon's claim protects it"
5905 );
5906 assert!(!still_alive.active.is_empty());
5907 }
5908
5909 #[test]
5919 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
5920 let dir = tempfile::tempdir().unwrap();
5921 let home = dir.path().to_path_buf();
5922 let now = Timestamp::now();
5923
5924 let mut manual = run_state(RunStatus::Reviewing);
5925 manual.id = "20260101-000000-manl".to_owned();
5926 manual.active.insert(
5927 "review-1".to_owned(),
5928 crate::run::ActiveSeat {
5929 node: "review".to_owned(),
5930 started_at: now - jiff::SignedDuration::new(21_000, 0),
5931 timeout_secs: 3_600,
5932 attempt: 0,
5933 task: None,
5934 command: None,
5935 index: None,
5936 total: None,
5937 },
5938 );
5939 manual.driver_pid = Some(4242);
5943 manual.driver_started_at = Some("1790000000".to_owned());
5944 manual.save_under(&home).unwrap();
5945
5946 let abandoned = reclaim_abandoned_runs_with(
5947 &home,
5948 now,
5949 |pid| if pid == 4242 { Some(true) } else { None },
5950 |pid| {
5951 if pid == 4242 {
5952 Some("1790000000".to_owned())
5953 } else {
5954 None
5955 }
5956 },
5957 );
5958 assert!(
5959 abandoned.is_empty(),
5960 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
5961 );
5962
5963 let reloaded = read_run_under(&home, &manual.id);
5964 assert_eq!(reloaded.status, RunStatus::Reviewing);
5965 assert!(!reloaded.active.is_empty());
5966 }
5967
5968 #[test]
5969 fn an_already_claimed_task_is_skipped_rather_than_failed() {
5970 let dir = tempfile::tempdir().unwrap();
5971 let queue = Queue::at(dir.path().to_path_buf());
5972 let mut only = task();
5973 queue.put(&mut only).unwrap();
5974
5975 let _elsewhere = queue.claim(&only.id).unwrap();
5976 let candidates = runnable(&queue);
5977 assert_eq!(candidates.len(), 1, "the task is still runnable");
5978 assert!(
5979 queue.claim(&candidates[0].id).is_err(),
5980 "the loop cannot take a claim somebody else holds"
5981 );
5982
5983 let after = queue.get(&only.id).unwrap();
5984 assert_eq!(after.status, TaskStatus::Queued);
5985 assert_eq!(
5986 after.attempts, 0,
5987 "losing the race is not an attempt at the task"
5988 );
5989 assert_eq!(after.last_error, None);
5990 }
5991
5992 #[test]
5993 fn the_status_file_round_trips_and_its_heartbeat_advances() {
5994 let dir = tempfile::tempdir().unwrap();
5995 let path = dir.path().join("daemon.json");
5996
5997 let mut status = Status::new();
5998 status.idle = false;
5999 status.completed = 7;
6000 status.current = vec![Current {
6001 task: "20260902-000000-t111".to_owned(),
6002 run: "20260902-000001-r111".to_owned(),
6003 }];
6004 write_status_to(&path, &status).unwrap();
6005 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6006 assert_eq!(first.schema, SCHEMA);
6007 assert_eq!(first.pid, std::process::id());
6008 assert!(!first.idle);
6009 assert_eq!(first.completed, 7);
6010 assert_eq!(first.current, status.current);
6011 assert!(
6012 !path.with_extension("json.tmp").exists(),
6013 "the temp file is renamed, not left behind"
6014 );
6015
6016 std::thread::sleep(Duration::from_millis(5));
6017 status.updated_at = Timestamp::now();
6018 status.polls = 3;
6019 write_status_to(&path, &status).unwrap();
6020 let second: Status =
6021 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6022 assert!(
6023 second.updated_at > first.updated_at,
6024 "a reader can only detect staleness if the heartbeat moves"
6025 );
6026 assert_eq!(
6027 second.started_at, first.started_at,
6028 "the start time is not a heartbeat"
6029 );
6030 assert_eq!(second.polls, 3);
6031 }
6032
6033 #[test]
6034 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
6035 let dir = tempfile::tempdir().unwrap();
6036
6037 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
6038
6039 let mut status = Status::new();
6040 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
6041 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6042 let stale = read_status(dir.path()).unwrap();
6043 assert!(
6044 !stale.running(Timestamp::now()),
6045 "a minute without a heartbeat is a dead daemon, not a busy one"
6046 );
6047 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
6048
6049 status.updated_at = Timestamp::now();
6050 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6051 let fresh = read_status(dir.path()).unwrap();
6052 assert!(fresh.running(Timestamp::now()));
6053 }
6054
6055 #[test]
6056 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
6057 let dir = tempfile::tempdir().unwrap();
6058 let now = Timestamp::now();
6059 let mine = "20260903-080619-01c2";
6060
6061 assert!(
6062 !is_working_on(dir.path(), mine, now),
6063 "no status file means nobody is working on anything"
6064 );
6065
6066 let mut status = Status::new();
6067 status.current = vec![Current {
6068 task: "20260903-080340-0167".to_owned(),
6069 run: mine.to_owned(),
6070 }];
6071 status.updated_at = now;
6072 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6073 assert!(is_working_on(dir.path(), mine, now));
6074 assert!(
6075 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
6076 "a daemon busy with one run is not working on another"
6077 );
6078
6079 status.updated_at = now - jiff::SignedDuration::from_secs(600);
6082 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6083 assert!(
6084 !is_working_on(dir.path(), mine, now),
6085 "a stale heartbeat is a dead daemon, so its run is a leftover"
6086 );
6087 }
6088
6089 #[test]
6090 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
6091 let dir = tempfile::tempdir().unwrap();
6092 let now = Timestamp::now();
6093
6094 assert!(
6095 !is_working_on_short(dir.path(), "01c2", now),
6096 "no status file means nobody is working on anything"
6097 );
6098
6099 let mut status = Status::new();
6100 status.current = vec![Current {
6101 task: "20260903-080340-0167".to_owned(),
6102 run: "20260903-080619-01c2".to_owned(),
6103 }];
6104 status.updated_at = now;
6105 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6106 assert!(
6107 is_working_on_short(dir.path(), "01c2", now),
6108 "the run's short id is the last block of its full id"
6109 );
6110 assert!(
6111 !is_working_on_short(dir.path(), "3cbf", now),
6112 "a daemon busy with one worktree bay is not working on another"
6113 );
6114 }
6115
6116 #[test]
6117 fn a_newer_status_file_still_yields_a_reading() {
6118 let dir = tempfile::tempdir().unwrap();
6119 std::fs::write(
6122 dir.path().join("daemon.json"),
6123 serde_json::json!({
6124 "schema": 2,
6125 "updated_at": Timestamp::now().to_string(),
6126 "idle": true,
6127 "surprise": { "nested": [1, 2, 3] },
6128 })
6129 .to_string(),
6130 )
6131 .unwrap();
6132
6133 let reading = read_status(dir.path()).expect("a forward-compatible read");
6134 assert!(reading.running(Timestamp::now()));
6135 assert!(reading.idle);
6136 assert!(reading.current.is_empty());
6137 }
6138
6139 #[test]
6140 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
6141 let dir = tempfile::tempdir().unwrap();
6147 std::fs::write(
6148 dir.path().join("daemon.json"),
6149 serde_json::json!({
6150 "schema": 1,
6151 "pid": 4242,
6152 "updated_at": Timestamp::now().to_string(),
6153 "idle": false,
6154 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
6155 "completed": 3,
6156 "polls": 9,
6157 })
6158 .to_string(),
6159 )
6160 .unwrap();
6161
6162 let reading = read_status(dir.path()).expect("an older shape must still parse");
6163 assert!(reading.running(Timestamp::now()));
6164 assert_eq!(
6165 reading.current,
6166 vec![Current {
6167 task: "20260902-140501-aaaa".to_owned(),
6168 run: "20260902-140502-bbbb".to_owned(),
6169 }]
6170 );
6171 }
6172
6173 #[test]
6174 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
6175 let dir = tempfile::tempdir().unwrap();
6176 std::fs::write(
6177 dir.path().join("daemon.json"),
6178 serde_json::json!({
6179 "schema": 1,
6180 "updated_at": Timestamp::now().to_string(),
6181 "idle": true,
6182 "current": null,
6183 })
6184 .to_string(),
6185 )
6186 .unwrap();
6187 let with_null = read_status(dir.path()).expect("null must still parse");
6188 assert!(with_null.current.is_empty());
6189
6190 std::fs::write(
6191 dir.path().join("daemon.json"),
6192 serde_json::json!({
6193 "schema": 1,
6194 "updated_at": Timestamp::now().to_string(),
6195 "idle": true,
6196 })
6197 .to_string(),
6198 )
6199 .unwrap();
6200 let absent = read_status(dir.path()).expect("a missing field must still parse");
6201 assert!(absent.current.is_empty());
6202 }
6203
6204 #[test]
6205 fn a_task_without_a_repository_runs_in_the_daemons_default() {
6206 let fallback = Path::new("/default");
6207 let mut blank = task();
6208 blank.repo = PathBuf::new();
6209 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
6210 let mut dot = task();
6211 dot.repo = PathBuf::from(".");
6212 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
6213 assert_eq!(
6214 repo_for(&task(), fallback),
6215 PathBuf::from("/repo"),
6216 "a task that names a repository keeps it"
6217 );
6218 }
6219
6220 #[test]
6221 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
6222 let mut solo_cfg = Config::default();
6228 solo_cfg.graph.candidates = 3;
6229 let mut solo_task = task();
6230 solo_task.solo = true;
6231 apply_solo(&mut solo_cfg, &solo_task);
6232 assert_eq!(solo_cfg.graph.candidates, 1);
6233
6234 let mut plain_cfg = Config::default();
6235 plain_cfg.graph.candidates = 3;
6236 let plain_task = task();
6237 assert!(!plain_task.solo);
6238 apply_solo(&mut plain_cfg, &plain_task);
6239 assert_eq!(
6240 plain_cfg.graph.candidates, 3,
6241 "a task that did not ask to run alone keeps the config's candidates"
6242 );
6243 }
6244
6245 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
6246 QuotaLoss {
6247 seat: seat.into(),
6248 node: "judge".into(),
6249 at: at.parse().unwrap(),
6250 reset: reset.map(str::to_string),
6251 }
6252 }
6253
6254 #[test]
6255 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
6256 let old: Vec<QuotaLoss> = (1..=4)
6257 .map(|i| {
6258 loss(
6259 &format!("judge-{i}"),
6260 "2026-09-23T05:23:00Z",
6261 Some("2:40pm (Asia/Tokyo)"),
6262 )
6263 })
6264 .collect();
6265 let fresh = losses_this_attempt(&old, &old);
6266 assert!(fresh.is_empty());
6267 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
6268 }
6270
6271 #[test]
6272 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6273 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6274 let now = Timestamp::now();
6275 let mut after = old.clone();
6276 after.push(loss("judge-2", &now.to_string(), None));
6277 let fresh = losses_this_attempt(&old, &after);
6278 assert_eq!(fresh, vec![after[1].clone()]);
6279 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6280 assert_eq!(
6281 until,
6282 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6283 );
6284 }
6285
6286 #[test]
6287 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6288 let before = vec![
6291 loss("judge-1", "2026-09-23T05:23:00Z", None),
6292 loss("judge-2", "2026-09-23T05:24:00Z", None),
6293 ];
6294 let after = vec![
6295 loss("judge-2", "2026-09-23T05:24:00Z", None),
6296 loss("judge-1", "2026-09-24T01:00:00Z", None),
6297 ];
6298 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6299 }
6300
6301 #[test]
6302 fn merge_overrides_are_parsed_or_refused() {
6303 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6304 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6305 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6306 assert!(merge_mode("squash").is_err());
6307 }
6308
6309 #[test]
6310 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6311 let now = Timestamp::now();
6312 let fallback = Duration::from_secs(300);
6313 let cap = Duration::from_secs(1800);
6314
6315 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6317
6318 let soon = now + jiff::SignedDuration::from_secs(600);
6320 assert_eq!(
6321 quota_wait(Some(soon), now, fallback, cap),
6322 Duration::from_secs(600)
6323 );
6324
6325 let past = now - jiff::SignedDuration::from_secs(60);
6328 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6329
6330 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6333 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6334 }
6335
6336 #[test]
6337 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6338 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6339
6340 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6341 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6342
6343 let already_past =
6347 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6348 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6349
6350 assert!(
6351 parse_reset_hint("session limit reached", now, now).is_none(),
6352 "free text with no recognised shape is not guessed at"
6353 );
6354 assert!(
6355 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6356 "an unresolvable zone name is not guessed at either"
6357 );
6358 }
6359
6360 #[test]
6361 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6362 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6363
6364 let at = parse_reset_hint(
6365 "You've hit your usage limit. Visit \
6366 https://chatgpt.com/codex/settings/usage to purchase more \
6367 credits or try again at Sep 19th, 2026 5:10 PM.",
6368 now,
6369 now,
6370 )
6371 .expect("the codex reset wording is a recognised shape");
6372 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6373
6374 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6379 .expect("an explicit year needs no rollover");
6380 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6381
6382 assert!(
6383 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6384 "a two-digit year is not the documented shape and is not guessed at"
6385 );
6386 assert!(
6387 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6388 "a four-letter month name is not the documented three-letter abbreviation"
6389 );
6390 assert!(
6391 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6392 "an explicit zone on the dated shape is a format nobody has \
6393 documented, and is refused rather than guessed at as UTC"
6394 );
6395 }
6396
6397 #[test]
6398 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6399 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6400 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6401
6402 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6403 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6404
6405 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6406 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6407
6408 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6409 assert!(
6410 parse_reset_hint(bad, now, recorded).is_none(),
6411 "{bad:?} must not be guessed at"
6412 );
6413 }
6414 }
6415
6416 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6420 let config = dir.join("magi.toml");
6421 std::fs::write(
6422 &config,
6423 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6424 )
6425 .unwrap();
6426 let opts = Opts {
6427 poll: Duration::from_secs(30),
6428 config: Some(config),
6429 repo: dir.join("repo"),
6433 ..Opts::default()
6434 };
6435 let home = dir.join("home");
6444 let worktrees = dir.join("wt");
6445 (
6446 opts,
6447 Queue::at(dir.join("queue")),
6448 home.join("daemon.json"),
6449 home,
6450 worktrees,
6451 )
6452 }
6453
6454 #[test]
6455 fn a_stop_is_idempotent_and_once_set_stays_set() {
6456 let stop = Stop::new();
6457 assert!(!stop.stopped());
6458
6459 stop.stop();
6460 assert!(stop.stopped());
6461 stop.stop();
6462 assert!(stop.stopped(), "a second stop is not a toggle");
6463
6464 let shared = stop.clone();
6465 assert!(
6466 shared.stopped(),
6467 "a clone is the same stop; that is how the loop and its caller share one"
6468 );
6469 }
6470
6471 #[test]
6472 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6473 let stop = Stop::new();
6474 stop.enter();
6475 assert!(
6476 !stop.finishing(),
6477 "a busy loop nobody has asked to stop is just running"
6478 );
6479
6480 stop.stop();
6481 assert!(
6482 stop.finishing(),
6483 "a stop asked for mid-run has not landed until the run is settled"
6484 );
6485
6486 stop.exit();
6487 assert!(
6488 !stop.finishing(),
6489 "once the run is settled the stop has landed and there is nothing to finish"
6490 );
6491 }
6492
6493 #[test]
6494 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6495 let stop = Stop::new();
6496 stop.enter();
6497 stop.enter();
6498 stop.stop();
6499 assert!(stop.finishing(), "two runs still in flight");
6500
6501 stop.exit();
6502 assert!(
6503 stop.finishing(),
6504 "one run finished, but a sibling is still working"
6505 );
6506
6507 stop.exit();
6508 assert!(
6509 !stop.finishing(),
6510 "the last run out is what actually lands the stop"
6511 );
6512 }
6513
6514 #[tokio::test]
6515 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6516 let dir = tempfile::tempdir().unwrap();
6517 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6518 let stop = Stop::new();
6519 stop.stop();
6520
6521 let began = std::time::Instant::now();
6522 tokio::time::timeout(
6523 Duration::from_secs(2),
6524 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6525 )
6526 .await
6527 .expect("a stopped loop must return, not sit out its poll interval")
6528 .expect("the loop's own setup and teardown must not fail");
6529 assert!(
6530 began.elapsed() < opts.poll,
6531 "returned only after {:?}, which is a poll interval, not a stop",
6532 began.elapsed()
6533 );
6534 }
6535
6536 #[tokio::test]
6537 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6538 let dir = tempfile::tempdir().unwrap();
6539 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6540 let stop = Stop::new();
6541
6542 let asker = {
6545 let stop = stop.clone();
6546 tokio::spawn(async move {
6547 tokio::time::sleep(Duration::from_millis(20)).await;
6548 stop.stop();
6549 })
6550 };
6551
6552 let began = std::time::Instant::now();
6553 tokio::time::timeout(
6554 Duration::from_secs(2),
6555 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6556 )
6557 .await
6558 .expect("a stop asked for while idle must wake the wait")
6559 .expect("the loop's own setup and teardown must not fail");
6560 asker.await.unwrap();
6561 assert!(
6562 began.elapsed() < opts.poll,
6563 "returned only after {:?}, so the stop waited on the sleep",
6564 began.elapsed()
6565 );
6566 }
6567
6568 #[tokio::test]
6569 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6570 let dir = tempfile::tempdir().unwrap();
6571 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6572 let stop = Stop::new();
6573 stop.stop();
6574
6575 tokio::time::timeout(
6576 Duration::from_secs(2),
6577 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6578 )
6579 .await
6580 .expect("a stopped loop must return")
6581 .expect("the loop's own setup and teardown must not fail");
6582
6583 assert!(
6584 home.is_dir(),
6585 "the loop did publish a status file, so its removal is the teardown and not an absence"
6586 );
6587 assert!(
6588 !status_file.exists(),
6589 "a stopped loop clears its status file"
6590 );
6591 assert!(
6592 read_status(&home).is_none(),
6593 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6594 );
6595 }
6596
6597 #[tokio::test]
6598 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6599 let dir = tempfile::tempdir().unwrap();
6600 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6601 opts.once = true;
6602
6603 let mut settled = RunState::new(
6604 dir.path().join("repo"),
6605 "main".to_owned(),
6606 "abc1234".to_owned(),
6607 "fixture".to_owned(),
6608 Config::default(),
6609 );
6610 settled.status = RunStatus::Ready;
6611 let run_dir = home.join("runs").join(&settled.id);
6612 std::fs::create_dir_all(&run_dir).unwrap();
6613 std::fs::write(
6614 run_dir.join("run.json"),
6615 serde_json::to_string_pretty(&settled).unwrap(),
6616 )
6617 .unwrap();
6618 let questions = Questions::at(home.join("questions"));
6619 let mut question = ask::Question::new(
6620 settled.id.clone(),
6621 "review".to_owned(),
6622 "reviewer-1".to_owned(),
6623 "Continue?".to_owned(),
6624 String::new(),
6625 Vec::new(),
6626 );
6627 questions.put(&mut question).unwrap();
6628
6629 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6630 .await
6631 .unwrap();
6632
6633 assert_eq!(
6634 questions.get(&question.id).unwrap().status,
6635 ask::QuestionStatus::Abandoned,
6636 "an empty --once drain still performs startup question cleanup"
6637 );
6638 }
6639
6640 #[test]
6641 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6642 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6643
6644 assert!(
6645 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6646 "never checked before: due at once"
6647 );
6648
6649 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6650 assert!(
6651 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6652 "well inside the interval: not due yet"
6653 );
6654
6655 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6656 assert!(
6657 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6658 "exactly at the edge: not yet due, same convention as `clean::due`"
6659 );
6660
6661 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6662 assert!(
6663 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6664 "past the interval: due again"
6665 );
6666 }
6667
6668 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
6673 let config = dir.join("magi.toml");
6674 std::fs::write(
6680 &config,
6681 format!(
6682 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
6683 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
6684 cache_dir.display()
6685 ),
6686 )
6687 .unwrap();
6688 Opts {
6689 config: Some(config),
6690 repo: dir.join("repo"),
6691 ..Opts::default()
6692 }
6693 }
6694
6695 #[tokio::test]
6696 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
6697 let dir = tempfile::tempdir().unwrap();
6698 let home = dir.path().join("home");
6699 let cache_dir = dir.path().join("cache");
6700 std::fs::create_dir_all(&cache_dir).unwrap();
6701 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6702 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6703
6704 let running = Stop::new();
6707 let mut last_checked = None;
6708 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6709 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
6710 .await;
6711 assert_eq!(
6712 crate::disk::dir_size(&cache_dir),
6713 0,
6714 "over the cap on the first check ever: pruned at once, no idle queue required"
6715 );
6716 assert_eq!(last_checked, Some(t0));
6717
6718 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
6720 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
6721 maybe_prune_cache_between_runs(
6722 &opts.repo,
6723 &opts,
6724 &home,
6725 &running,
6726 &mut last_checked,
6727 too_soon,
6728 )
6729 .await;
6730 assert_eq!(
6731 crate::disk::dir_size(&cache_dir),
6732 10,
6733 "too soon since the last check: left alone rather than rescanned every call"
6734 );
6735 assert_eq!(
6736 last_checked,
6737 Some(t0),
6738 "an idle check does not reset the clock"
6739 );
6740
6741 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6743 maybe_prune_cache_between_runs(
6744 &opts.repo,
6745 &opts,
6746 &home,
6747 &running,
6748 &mut last_checked,
6749 due_again,
6750 )
6751 .await;
6752 assert_eq!(
6753 crate::disk::dir_size(&cache_dir),
6754 0,
6755 "due again: pruned back under the cap"
6756 );
6757 }
6758
6759 #[tokio::test]
6767 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
6768 let dir = tempfile::tempdir().unwrap();
6769 let home = dir.path().join("home");
6770 let cache_dir = dir.path().join("cache");
6771 std::fs::create_dir_all(&cache_dir).unwrap();
6772 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6773 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6774
6775 let stop = Stop::new();
6776 stop.stop();
6777 assert!(
6778 !stop.finishing(),
6779 "no run is in flight at a between-runs boundary, so nothing else \
6780 would tell the operator this stop had not taken effect yet"
6781 );
6782
6783 let mut last_checked = None;
6784 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6785 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
6786 .await;
6787 assert_eq!(
6788 crate::disk::dir_size(&cache_dir),
6789 10,
6790 "over its cap, and due for the first check ever, but a stop outranks \
6791 it: the cap is a standing policy the next start measures again"
6792 );
6793 assert_eq!(
6794 last_checked, None,
6795 "a check that never happened must not claim the interval"
6796 );
6797 }
6798
6799 #[tokio::test]
6813 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
6814 let dir = tempfile::tempdir().unwrap();
6815 let cache_dir = dir.path().join("cache");
6816 std::fs::create_dir_all(&cache_dir).unwrap();
6817 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
6818
6819 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
6820 opts.poll = Duration::from_millis(20);
6821 opts.max_attempts = 1_000;
6822
6823 let queue = Queue::at(dir.path().join("queue"));
6824 let mut t = Task::new(
6825 "x".to_owned(),
6826 "x".to_owned(),
6827 opts.repo.clone(),
6828 Source::Human,
6829 );
6830 queue.put(&mut t).unwrap();
6831
6832 let home = dir.path().join("home");
6833 let worktrees = dir.path().join("wt");
6834 let status_file = home.join("daemon.json");
6835 let stop = Stop::new();
6836 let stopper = {
6837 let stop = stop.clone();
6838 tokio::spawn(async move {
6839 tokio::time::sleep(Duration::from_millis(400)).await;
6840 stop.stop();
6841 })
6842 };
6843
6844 tokio::time::timeout(
6845 Duration::from_secs(10),
6846 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6847 )
6848 .await
6849 .expect("the loop must not hang on a queue that keeps producing failing work")
6850 .expect("the loop's own setup and teardown must not fail");
6851 stopper.await.unwrap();
6852
6853 let after = queue.get(&t.id).unwrap();
6854 assert!(
6855 after.attempts >= 2,
6856 "the harness must actually have retried more than once, or this is not \
6857 exercising a busy queue at all (got {} attempt(s))",
6858 after.attempts
6859 );
6860 assert!(
6861 after.status.runnable(),
6862 "still under its attempt budget: the queue never reached a natural idle \
6863 on its own, only the external stop ended the test"
6864 );
6865
6866 assert_eq!(
6867 crate::disk::dir_size(&cache_dir),
6868 0,
6869 "an oversized cache must not be left to grow unboundedly just because the \
6870 queue kept the loop busy the whole time"
6871 );
6872 }
6873
6874 #[test]
6875 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
6876 let dir = tempfile::tempdir().unwrap();
6877 let queue = Queue::at(dir.path().join("queue"));
6878 let questions = Questions::at(dir.path().join("questions"));
6879 let mut task = task();
6880 queue.put(&mut task).unwrap();
6881
6882 let mut task_question = ask::Question::new(
6883 task.id.clone(),
6884 crate::conduct::NODE.to_owned(),
6885 "conduct".to_owned(),
6886 "Which backend?".to_owned(),
6887 String::new(),
6888 Vec::new(),
6889 );
6890 questions.put(&mut task_question).unwrap();
6891 task.block(vec![task_question.id.clone()], None);
6892 queue.put(&mut task).unwrap();
6893
6894 let mut run_question = ask::Question::new(
6895 "20260101-000000-run1".to_owned(),
6896 "review".to_owned(),
6897 "reviewer-1".to_owned(),
6898 "Run question".to_owned(),
6899 String::new(),
6900 Vec::new(),
6901 );
6902 questions.put(&mut run_question).unwrap();
6903
6904 let mut coincidental = ask::Question::new(
6909 task.id.clone(),
6910 "review".to_owned(),
6911 "reviewer-1".to_owned(),
6912 "Unrelated review question".to_owned(),
6913 String::new(),
6914 Vec::new(),
6915 );
6916 questions.put(&mut coincidental).unwrap();
6917
6918 reconcile_task_questions(&queue, &questions);
6919 assert!(questions.get(&task_question.id).unwrap().status.open());
6920 assert!(questions.get(&run_question.id).unwrap().status.open());
6921 assert!(questions.get(&coincidental.id).unwrap().status.open());
6922
6923 task.release();
6924 queue.put(&mut task).unwrap();
6925 reconcile_task_questions(&queue, &questions);
6926 assert_eq!(
6927 questions.get(&task_question.id).unwrap().status,
6928 ask::QuestionStatus::Abandoned
6929 );
6930 assert!(
6931 questions.get(&run_question.id).unwrap().status.open(),
6932 "run questions remain the run janitor's responsibility"
6933 );
6934 assert!(
6935 questions.get(&coincidental.id).unwrap().status.open(),
6936 "a non-conductor question must not be abandoned just because its \
6937 run id coincides with a task id"
6938 );
6939 }
6940
6941 #[test]
6942 fn a_freshly_started_running_task_is_never_stalled() {
6943 let dir = tempfile::tempdir().unwrap();
6944 let mut t = task();
6945 t.start("run-1".to_owned());
6946 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
6949 }
6950
6951 #[test]
6952 fn a_long_running_task_with_no_live_daemon_is_stalled() {
6953 let dir = tempfile::tempdir().unwrap();
6954 let mut t = task();
6955 t.start("run-1".to_owned());
6956 t.updated_at = Timestamp::now()
6957 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6958 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
6959 assert_eq!(
6960 stalled_tasks(
6961 &Queue::at(dir.path().join("q")),
6962 dir.path(),
6963 Timestamp::now()
6964 )
6965 .len(),
6966 0,
6967 "the task was never written to this queue"
6968 );
6969 }
6970
6971 #[test]
6972 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
6973 let dir = tempfile::tempdir().unwrap();
6974 let mut t = task();
6975 t.id = "20260903-080340-0167".to_owned();
6976 t.start("20260903-080619-01c2".to_owned());
6977 t.updated_at = Timestamp::now()
6978 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6979
6980 let mut status = Status::new();
6981 status.current = vec![Current {
6982 task: t.id.clone(),
6983 run: "20260903-080619-01c2".to_owned(),
6984 }];
6985 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6986
6987 assert!(
6988 !is_stalled(&t, dir.path(), Timestamp::now()),
6989 "a live daemon's own heartbeat rules out stalled, however long the task has run"
6990 );
6991 }
6992
6993 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
6997 let path = queue.path_of(id);
6998 let body = std::fs::read_to_string(&path).unwrap();
6999 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
7000 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
7001 v["updated_at"] = serde_json::Value::String(old.to_string());
7002 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
7003 }
7004
7005 #[test]
7006 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
7007 let dir = tempfile::tempdir().unwrap();
7020 let queue = Queue::at(dir.path().join("queue"));
7021 let home = dir.path().join("home");
7022
7023 let mut t = task();
7024 t.id = "20260101-000001-lock".to_owned();
7025 t.start("run-1".to_owned());
7026 queue.put(&mut t).unwrap();
7027 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7028 std::fs::write(
7029 dir.path().join("queue").join(format!("{}.lock", t.id)),
7030 "not a pid",
7031 )
7032 .unwrap();
7033
7034 let now = Timestamp::now();
7035 assert!(
7036 reclaim_orphaned_running(&queue, 2).is_empty(),
7037 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
7038 and reclaim must leave the task alone"
7039 );
7040 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
7041
7042 let stalled = stalled_tasks(&queue, &home, now);
7043 assert_eq!(
7044 stalled.len(),
7045 1,
7046 "reclaim's inability to claim it yet must not hide it from the conductor"
7047 );
7048 assert_eq!(stalled[0].id, t.id);
7049 }
7050
7051 #[test]
7052 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
7053 let dir = tempfile::tempdir().unwrap();
7054 crate::run::set_home(dir.path().join("run-home"));
7055 let queue = Queue::at(dir.path().join("queue"));
7056 let home = dir.path().join("home");
7057 let questions = Questions::at(dir.path().join("questions"));
7058
7059 let mut t = task();
7060 t.id = "20260101-000003-dead".to_owned();
7061 t.start("missing-run".to_owned());
7062 queue.put(&mut t).unwrap();
7063 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7064
7065 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7068 assert_eq!(
7069 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
7070 [&t.id]
7071 );
7072 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
7073 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7074
7075 crate::conduct::apply(
7078 &queue,
7079 &questions,
7080 &crate::conduct::Verdict {
7081 decisions: vec![crate::conduct::Decision {
7082 id: t.id.clone(),
7083 recovery: Some(crate::conduct::Recovery::Requeue),
7084 ..crate::conduct::Decision::default()
7085 }],
7086 },
7087 )
7088 .unwrap();
7089 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7090 }
7091
7092 #[test]
7093 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
7094 let dir = tempfile::tempdir().unwrap();
7095 let queue = Queue::at(dir.path().join("queue"));
7096 let home = dir.path().join("home");
7097
7098 let mut fresh = task();
7099 fresh.id = "20260101-000001-aaaa".to_owned();
7100 fresh.start("run-1".to_owned());
7101 queue.put(&mut fresh).unwrap();
7102
7103 let mut old = task();
7104 old.id = "20260101-000002-bbbb".to_owned();
7105 old.start("run-2".to_owned());
7106 queue.put(&mut old).unwrap();
7107 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
7108
7109 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7110 assert_eq!(stalled.len(), 1);
7111 assert_eq!(stalled[0].id, old.id);
7112 }
7113
7114 #[test]
7115 fn queued_and_finished_task_views_partition_by_status() {
7116 let dir = tempfile::tempdir().unwrap();
7117 let queue = Queue::at(dir.path().join("queue"));
7118
7119 let mut queued = task();
7120 queued.id = "20260101-000001-aaaa".to_owned();
7121 queue.put(&mut queued).unwrap();
7122
7123 let mut failed = task();
7124 failed.id = "20260101-000002-bbbb".to_owned();
7125 failed.start("run-1".to_owned());
7126 failed.fail("gate red", 5);
7127 queue.put(&mut failed).unwrap();
7128
7129 let mut held = task();
7130 held.id = "20260101-000003-cccc".to_owned();
7131 held.hold_machine(None);
7132 queue.put(&mut held).unwrap();
7133
7134 let mut running = task();
7135 running.id = "20260101-000004-dddd".to_owned();
7136 running.start("run-2".to_owned());
7137 queue.put(&mut running).unwrap();
7138
7139 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
7140 assert_eq!(queued_ids, [queued.id.clone()]);
7141
7142 let mut finished_ids: Vec<String> =
7143 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
7144 finished_ids.sort_unstable();
7145 let mut want = vec![failed.id.clone(), held.id.clone()];
7146 want.sort_unstable();
7147 assert_eq!(finished_ids, want);
7148 }
7149
7150 #[test]
7151 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
7152 let dir = tempfile::tempdir().unwrap();
7153 let queue = Queue::at(dir.path().join("queue"));
7154 let questions = ask::Questions::at(dir.path().join("questions"));
7155
7156 let mut dep = task();
7157 dep.id = "20260101-000001-dep0".to_owned();
7158 dep.succeed();
7159 queue.put(&mut dep).unwrap();
7160
7161 let mut still_going = task();
7162 still_going.id = "20260101-000002-dep1".to_owned();
7163 queue.put(&mut still_going).unwrap();
7164
7165 let mut blocked = task();
7166 blocked.id = "20260101-000003-main".to_owned();
7167 blocked.block(
7168 vec![dep.id.clone(), still_going.id.clone()],
7169 Some("waits on both".to_owned()),
7170 );
7171 queue.put(&mut blocked).unwrap();
7172
7173 resolve_blockers(&queue, &questions);
7174
7175 let after = queue.get(&blocked.id).unwrap();
7176 assert_eq!(
7177 after.status,
7178 TaskStatus::Blocked,
7179 "one dependency is still outstanding"
7180 );
7181 assert_eq!(after.blocked_by, [still_going.id.clone()]);
7182 }
7183
7184 #[test]
7185 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
7186 let dir = tempfile::tempdir().unwrap();
7187 let queue = Queue::at(dir.path().join("queue"));
7188 let questions = ask::Questions::at(dir.path().join("questions"));
7189
7190 let mut q = crate::ask::Question::new(
7191 "20260101-000001-main".to_owned(),
7192 crate::conduct::NODE.to_owned(),
7193 "conduct".to_owned(),
7194 "Which backend?".to_owned(),
7195 String::new(),
7196 Vec::new(),
7197 );
7198 questions.put(&mut q).unwrap();
7199 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
7200 .unwrap();
7201 questions.put(&mut q).unwrap();
7202
7203 let mut blocked = task();
7204 blocked.id = "20260101-000001-main".to_owned();
7205 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
7206 queue.put(&mut blocked).unwrap();
7207
7208 resolve_blockers(&queue, &questions);
7209
7210 let after = queue.get(&blocked.id).unwrap();
7211 assert_eq!(
7212 after.status,
7213 TaskStatus::Queued,
7214 "the only blocker resolved"
7215 );
7216 assert_eq!(after.answers.len(), 1);
7217 assert_eq!(after.answers[0].question, "Which backend?");
7218 assert_eq!(after.answers[0].answer, "SQLite");
7219
7220 let instruction = instruction_for(&after);
7222 assert!(instruction.contains("Which backend?"));
7223 assert!(instruction.contains("SQLite"));
7224 }
7225
7226 #[test]
7227 fn resolve_blockers_holds_a_task_whose_conductor_question_was_abandoned() {
7228 let dir = tempfile::tempdir().unwrap();
7229 let queue = Queue::at(dir.path().join("queue"));
7230 let questions = ask::Questions::at(dir.path().join("questions"));
7231
7232 let mut q = crate::ask::Question::new(
7233 "20260101-000001-main".to_owned(),
7234 crate::conduct::NODE.to_owned(),
7235 "conduct".to_owned(),
7236 "Is the setup done?".to_owned(),
7237 String::new(),
7238 Vec::new(),
7239 );
7240 q.abandon("no answer within 60s of asking");
7241 questions.put(&mut q).unwrap();
7242
7243 let mut blocked = task();
7244 blocked.id = "20260101-000001-main".to_owned();
7245 blocked.block(vec![q.id.clone()], Some("setup?".to_owned()));
7246 queue.put(&mut blocked).unwrap();
7247
7248 resolve_blockers(&queue, &questions);
7249
7250 let after = queue.get(&blocked.id).unwrap();
7251 assert_eq!(
7252 after.status,
7253 TaskStatus::Held,
7254 "never left blocked on nothing"
7255 );
7256 assert!(!after.operator_held(), "a machine hold, for triage");
7257 assert!(
7258 after
7259 .hold_reason
7260 .as_deref()
7261 .unwrap_or_default()
7262 .contains("went unanswered")
7263 );
7264 }
7265
7266 #[test]
7267 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
7268 let dir = tempfile::tempdir().unwrap();
7274 let queue = Queue::at(dir.path().join("queue"));
7275 let questions = ask::Questions::at(dir.path().join("questions"));
7276
7277 let mut q = crate::ask::Question::new(
7278 "20260101-000001-main".to_owned(),
7279 crate::conduct::NODE.to_owned(),
7280 "conduct".to_owned(),
7281 "How should this be handled?".to_owned(),
7282 String::new(),
7283 Vec::new(),
7284 );
7285 questions.put(&mut q).unwrap();
7286 q.answer(crate::ask::Answer::Text(
7287 "leave it held, a human will look at it later".to_owned(),
7288 ))
7289 .unwrap();
7290 questions.put(&mut q).unwrap();
7291
7292 let mut held = task();
7293 held.id = "20260101-000001-main".to_owned();
7294 held.hold_machine(Some("out of attempts".to_owned()));
7295 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
7296 queue.put(&mut held).unwrap();
7297
7298 resolve_blockers(&queue, &questions);
7299
7300 let after = queue.get(&held.id).unwrap();
7301 assert_eq!(after.status, TaskStatus::Held);
7302 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
7303 assert_eq!(
7304 after.answers[0].answer,
7305 "leave it held, a human will look at it later"
7306 );
7307 }
7308
7309 #[test]
7310 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
7311 let dir = tempfile::tempdir().unwrap();
7317 let queue = Queue::at(dir.path().join("queue"));
7318 let questions = ask::Questions::at(dir.path().join("questions"));
7319
7320 let mut still_going = task();
7321 still_going.id = "20260101-000002-dep1".to_owned();
7322 queue.put(&mut still_going).unwrap();
7323
7324 let mut blocked = task();
7325 blocked.id = "20260101-000003-main".to_owned();
7326 blocked.block(
7327 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7328 Some("waits on both".to_owned()),
7329 );
7330 queue.put(&mut blocked).unwrap();
7331
7332 resolve_blockers(&queue, &questions);
7333
7334 let after = queue.get(&blocked.id).unwrap();
7335 assert_eq!(
7336 after.status,
7337 TaskStatus::Held,
7338 "a missing dependency must not leave the task blocked forever"
7339 );
7340 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7341 assert!(after.blocked_by.is_empty());
7342 let reason = after.hold_reason.as_deref().unwrap_or_default();
7343 assert!(
7344 reason.contains("20260101-000001-gone"),
7345 "the missing id must be named so an operator can tell what happened: {reason}"
7346 );
7347 assert!(
7348 reason.contains(&still_going.id),
7349 "the still-valid dependency must not silently vanish from the record: {reason}"
7350 );
7351 }
7352
7353 #[test]
7354 fn instruction_for_is_unchanged_without_any_answers() {
7355 let t = task();
7356 assert_eq!(instruction_for(&t), t.instruction);
7357 }
7358
7359 #[test]
7360 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7361 let dir = tempfile::tempdir().unwrap();
7362 let q = Queue::at(dir.path().join("queue"));
7363 let src = dir.path().join("shot.png");
7364 std::fs::write(&src, "x").unwrap();
7365 let mut t = task();
7366 q.attach(&mut t, &[src]).unwrap();
7367 let paths = task_attachments(&q, &t).unwrap();
7368 assert_eq!(paths.len(), 1);
7369 assert!(paths[0].is_absolute() && paths[0].is_file());
7370 std::fs::remove_file(&paths[0]).unwrap();
7371 let err = task_attachments(&q, &t).unwrap_err().to_string();
7372 assert!(err.contains("shot.png"), "{err}");
7373 }
7374
7375 #[test]
7376 fn resumed_instruction_is_unchanged_without_any_answers() {
7377 let t = task();
7378 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7379 }
7380
7381 #[test]
7382 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7383 let mut t = task();
7384 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7385 let old = t.instruction.clone();
7389
7390 let refreshed = resumed_instruction(&old, &t);
7391 assert!(refreshed.starts_with(&old), "the original text is kept");
7392 assert!(refreshed.contains("Which backend?"));
7393 assert!(refreshed.contains("SQLite"));
7394 }
7395
7396 #[test]
7397 fn resumed_instruction_keeps_an_original_answers_heading() {
7398 let mut t = task();
7399 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7400 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7401
7402 let refreshed = resumed_instruction(&t.instruction, &t);
7403
7404 assert!(
7405 refreshed.starts_with(&t.instruction),
7406 "an answers heading in the original instruction is not the appended block"
7407 );
7408 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7409 assert!(refreshed.contains("Which backend?"));
7410 assert!(refreshed.contains("SQLite"));
7411
7412 let repeated = resumed_instruction(&refreshed, &t);
7413 assert_eq!(
7414 repeated, refreshed,
7415 "only the final appended block is refreshed"
7416 );
7417 }
7418
7419 #[test]
7420 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7421 let mut t = task();
7422 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7423
7424 let once = resumed_instruction(&t.instruction, &t);
7428 let twice = resumed_instruction(&once, &t);
7429 assert_eq!(once, twice);
7430 assert_eq!(once.matches("Which backend?").count(), 1);
7431
7432 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7434 let refreshed = resumed_instruction(&once, &t);
7435 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7436 assert!(refreshed.contains("Which backend?"));
7437 assert!(refreshed.contains("Which cache?"));
7438 }
7439
7440 #[test]
7441 fn prepare_instruction_covers_all_three_starters() {
7442 let mut t = task();
7443 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7444
7445 assert_eq!(
7448 prepare_instruction(&Starter::Start, None, &t),
7449 Some(instruction_for(&t))
7450 );
7451
7452 let old = t.instruction.clone();
7455 assert_eq!(
7456 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7457 Some(resumed_instruction(&old, &t))
7458 );
7459
7460 assert_eq!(
7464 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7465 None
7466 );
7467 }
7468
7469 #[test]
7470 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7471 assert_eq!(
7472 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7473 Starter::Review("magi/eba2/A".to_owned())
7474 );
7475 }
7476
7477 #[test]
7478 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7479 assert_eq!(
7480 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7481 Starter::Start,
7482 "a vanished review branch must not fall back to resuming the old run either"
7483 );
7484 }
7485
7486 #[test]
7487 fn a_refused_handover_retries_as_a_review_of_the_same_branch() {
7488 let mut t = task();
7489 t.start("old-run".to_owned());
7490 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
7491 t.release();
7492 let branch = t.review_branch.take();
7493 assert_eq!(
7494 choose_starter(branch.as_deref(), true, Some("old-run")),
7495 Starter::Review("magi/eba2/A".to_owned()),
7496 "a review wins over resuming the old run"
7497 );
7498 }
7499
7500 #[test]
7501 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7502 assert_eq!(
7503 choose_starter(None, false, Some("some-run")),
7504 Starter::Resume("some-run".to_owned())
7505 );
7506 assert_eq!(choose_starter(None, false, None), Starter::Start);
7507 }
7508
7509 #[test]
7510 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7511 let mut released = task();
7512 released.start("stalled-run".to_owned());
7513 released.requeue();
7514 let unfinished = (!released.fresh_start)
7515 .then(|| Some("stalled-run".to_owned()))
7516 .flatten();
7517 assert_eq!(
7518 choose_starter(None, false, unfinished.as_deref()),
7519 Starter::Start,
7520 "release keeps run history but must not resume it"
7521 );
7522 assert_eq!(released.runs, ["stalled-run"]);
7523 }
7524
7525 #[test]
7526 fn an_ordinary_release_keeps_a_resumable_run_available() {
7527 let mut released = task();
7528 released.start("stalled-run".to_owned());
7529 released.release();
7530 let unfinished = (!released.fresh_start)
7531 .then(|| Some("stalled-run".to_owned()))
7532 .flatten();
7533 assert_eq!(
7534 choose_starter(None, false, unfinished.as_deref()),
7535 Starter::Resume("stalled-run".to_owned()),
7536 "manual release must preserve the normal resume path"
7537 );
7538 }
7539
7540 #[test]
7541 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7542 let mut state = run_state(RunStatus::Blocked);
7543 state.config.graph.review_rounds = 3;
7544 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7545 assert!(exhausted_review_budget(&state));
7546
7547 state.reviews.pop();
7549 assert!(!exhausted_review_budget(&state));
7550
7551 let mut stalled = run_state(RunStatus::Stalled);
7554 stalled.config.graph.review_rounds = 1;
7555 stalled.reviews = vec![review_round(1)];
7556 assert!(!exhausted_review_budget(&stalled));
7557 }
7558
7559 fn review_round(round: usize) -> crate::run::ReviewRound {
7560 crate::run::ReviewRound {
7561 round,
7562 head: "deadbeef".to_owned(),
7563 verified_head: None,
7564 verified_at: None,
7565 reviews: Vec::new(),
7566 e2e: Vec::new(),
7567 verify_retried: false,
7568 e2e_deferred: false,
7569 e2e_defer_reason: None,
7570 fix: None,
7571 blocking: 0,
7572 answered: 1,
7573 expected: 1,
7574 clean: false,
7575 progressed: true,
7576 vote_split: false,
7577 reconsideration: Vec::new(),
7578 verdict: None,
7579 }
7580 }
7581
7582 #[test]
7583 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7584 let mut run = RunState::new(
7585 PathBuf::from("/repo"),
7586 "main".to_owned(),
7587 "abc1234def".to_owned(),
7588 "add retries".to_owned(),
7589 Config::default(),
7590 );
7591 run.status = RunStatus::Judging;
7592 run.parked = true;
7593 let mut task = Task::new(
7594 "add retries".to_owned(),
7595 "add retries".to_owned(),
7596 PathBuf::from("/repo"),
7597 crate::queue::Source::Human,
7598 );
7599 task.status = TaskStatus::Failed;
7600 task.runs = vec![run.id.clone()];
7601 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7602 assert!(with(&task, &run), "parked after judging is the case");
7603
7604 let mut not_parked = run.clone();
7605 not_parked.parked = false;
7606 not_parked.status = RunStatus::Stalled;
7607 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7608
7609 let mut fresh = task.clone();
7610 fresh.fresh_start = true;
7611 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7612
7613 let mut review = task.clone();
7614 review.review_branch = Some("magi/x/A".to_owned());
7615 assert!(!with(&review, &run), "review is ranked before resume");
7616
7617 let mut held = task.clone();
7618 held.status = TaskStatus::Held;
7619 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7620
7621 let mut released = run.clone();
7622 released.released_to = Some("20260901-000000-new1".to_owned());
7623 assert!(!with(&task, &released), "nothing left to resume into");
7624
7625 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7626 "unreadable"
7627 )));
7628 }
7629
7630 #[test]
7631 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7632 let mut released = RunState::new(
7633 PathBuf::from("/repo"),
7634 "main".to_owned(),
7635 "abc1234def".to_owned(),
7636 "add retries".to_owned(),
7637 Config::default(),
7638 );
7639 released.status = RunStatus::Blocked;
7640 assert_eq!(
7641 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7642 Some(released.id.clone())
7643 );
7644 released.released_to = Some("20260901-000000-new1".to_owned());
7645 assert_eq!(
7646 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7647 None,
7648 "there is nothing left to resume it into"
7649 );
7650 }
7651
7652 #[test]
7653 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
7654 let mut exhausted = RunState::new(
7663 PathBuf::from("/repo"),
7664 "main".to_owned(),
7665 "abc1234def".to_owned(),
7666 "add retries".to_owned(),
7667 Config::default(),
7668 );
7669 exhausted.status = RunStatus::Blocked;
7670 exhausted.config.graph.review_rounds = 1;
7671 exhausted.reviews = vec![review_round(1)];
7672
7673 assert_eq!(
7674 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
7675 None,
7676 "an exhausted `Blocked` run must not be offered as resumable"
7677 );
7678
7679 let mut has_budget_left = RunState::new(
7682 PathBuf::from("/repo"),
7683 "main".to_owned(),
7684 "abc1234def".to_owned(),
7685 "add retries".to_owned(),
7686 Config::default(),
7687 );
7688 has_budget_left.status = RunStatus::Blocked;
7689 has_budget_left.config.graph.review_rounds = 3;
7690 has_budget_left.reviews = vec![review_round(1)];
7691
7692 assert_eq!(
7693 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
7694 Ok(has_budget_left.clone())
7695 }),
7696 Some(has_budget_left.id.clone())
7697 );
7698 }
7699
7700 #[test]
7701 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
7702 let mut older_stalled = RunState::new(
7710 PathBuf::from("/repo"),
7711 "main".to_owned(),
7712 "abc1234def".to_owned(),
7713 "add retries".to_owned(),
7714 Config::default(),
7715 );
7716 older_stalled.status = RunStatus::Stalled;
7717
7718 let mut newest_exhausted = RunState::new(
7719 PathBuf::from("/repo"),
7720 "main".to_owned(),
7721 "abc1234def".to_owned(),
7722 "add retries".to_owned(),
7723 Config::default(),
7724 );
7725 newest_exhausted.status = RunStatus::Blocked;
7726 newest_exhausted.config.graph.review_rounds = 1;
7727 newest_exhausted.reviews = vec![review_round(1)];
7728
7729 assert_eq!(
7730 unfinished_run_with(
7731 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
7732 "t",
7733 |_| Ok(newest_exhausted.clone())
7734 ),
7735 None,
7736 "the newest run is exhausted, so nothing here is worth resuming - \
7737 least of all the older, already-superseded run"
7738 );
7739 }
7740
7741 #[test]
7742 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
7743 assert_eq!(
7744 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
7745 Err(anyhow::anyhow!("fixture is absent"))
7746 }),
7747 None
7748 );
7749 }
7750
7751 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
7752 let mut q = ask::Question::new(
7753 run.to_owned(),
7754 "implement".to_owned(),
7755 "impl-A".to_owned(),
7756 "continue?".to_owned(),
7757 String::new(),
7758 vec!["resume で続行する".to_owned(), "other".to_owned()],
7759 );
7760 q.actions.insert("resume で続行する".to_owned(), action);
7761 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
7762 .unwrap();
7763 q
7764 }
7765
7766 fn held_task_with(run: &str) -> Task {
7767 let mut t = task();
7768 t.runs = vec![run.to_owned()];
7769 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
7770 t
7771 }
7772
7773 fn resume_action(run: &str) -> ask::ChoiceAction {
7774 ask::ChoiceAction::Resume { run: run.into() }
7775 }
7776
7777 #[test]
7778 fn decide_action_resumes_only_the_latest_resumable_run() {
7779 let t = held_task_with("r1");
7780 let q = action_question("r1", resume_action("r1"));
7781 let load = |s: RunState| move |_: &str| Ok(s);
7782 assert_eq!(
7783 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
7784 ActionDecision::Resume("r1".into())
7785 );
7786 let q_other = action_question("r1", resume_action("r0"));
7788 assert!(matches!(
7789 decide_action(
7790 &t,
7791 &q_other,
7792 &PHRASES_EN,
7793 load(run_state(RunStatus::Blocked))
7794 ),
7795 ActionDecision::Refuse(_)
7796 ));
7797 assert!(matches!(
7799 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
7800 ActionDecision::Refuse(_)
7801 ));
7802 let mut released = run_state(RunStatus::Blocked);
7804 released.released_to = Some("elsewhere".into());
7805 assert!(matches!(
7806 decide_action(&t, &q, &PHRASES_EN, load(released)),
7807 ActionDecision::Refuse(_)
7808 ));
7809 assert!(matches!(
7811 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
7812 ActionDecision::Refuse(_)
7813 ));
7814 }
7815
7816 #[test]
7817 fn decide_action_ignores_a_question_about_an_earlier_run() {
7818 let mut t = held_task_with("r1");
7819 t.runs.push("r2".to_owned());
7820 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7821 assert_eq!(
7822 decide_action(
7823 &t,
7824 &action_question("r1", ask::ChoiceAction::Done),
7825 &PHRASES_EN,
7826 never
7827 ),
7828 ActionDecision::Stale
7829 );
7830 }
7831
7832 #[test]
7833 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
7834 let mut t = held_task_with("r1");
7835 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7836 assert_eq!(
7837 decide_action(
7838 &t,
7839 &action_question("r1", ask::ChoiceAction::Requeue),
7840 &PHRASES_EN,
7841 never
7842 ),
7843 ActionDecision::Requeue
7844 );
7845 let done_q = action_question("r1", ask::ChoiceAction::Done);
7846 assert_eq!(
7847 decide_action(&t, &done_q, &PHRASES_EN, never),
7848 ActionDecision::Done
7849 );
7850 t.mark_action_applied(&done_q.id);
7851 assert_eq!(
7852 decide_action(&t, &done_q, &PHRASES_EN, never),
7853 ActionDecision::Skip
7854 );
7855
7856 let mut plain = action_question("r1", ask::ChoiceAction::Done);
7858 plain.actions.clear();
7859 assert_eq!(
7860 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
7861 ActionDecision::Skip
7862 );
7863 let mut running = held_task_with("r1");
7865 running.status = TaskStatus::Running;
7866 assert_eq!(
7867 decide_action(
7868 &running,
7869 &action_question("r1", ask::ChoiceAction::Done),
7870 &PHRASES_EN,
7871 never
7872 ),
7873 ActionDecision::Skip
7874 );
7875 }
7876
7877 #[test]
7878 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
7879 let dir = tempfile::tempdir().unwrap();
7880 let queue = Queue::at(dir.path().join("queue"));
7881 let questions = Questions::at(dir.path().join("questions"));
7882 let home = dir.path().join("home");
7883 let mut state = run_state(RunStatus::Blocked);
7884 state.id = "20260101-000000-act1".to_owned();
7885 state.save_under(&home).unwrap();
7886
7887 let mut t = held_task_with(&state.id);
7888 queue.put(&mut t).unwrap();
7889 let mut q = action_question(&state.id, resume_action(&state.id));
7890 questions.put(&mut q).unwrap();
7891
7892 apply_choice_actions(&queue, &questions, &home);
7893 let after = queue.get(&t.id).unwrap();
7894 assert_eq!(after.status, TaskStatus::Queued);
7895 assert!(!after.fresh_start);
7896 assert!(after.action_applied(&q.id));
7897 let pin = after.resume_override.clone().unwrap();
7898 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
7899 assert!(pin.forced);
7900
7901 let mut again = queue.get(&t.id).unwrap();
7903 again.hold_machine(Some("later".into()));
7904 queue.put(&mut again).unwrap();
7905 apply_choice_actions(&queue, &questions, &home);
7906 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7907 }
7908
7909 #[test]
7910 fn an_answer_the_waiter_already_delivered_is_not_acted_on_again() {
7911 let dir = tempfile::tempdir().unwrap();
7912 let queue = Queue::at(dir.path().join("queue"));
7913 let questions = Questions::at(dir.path().join("questions"));
7914 let home = dir.path().join("home");
7915 let mut state = run_state(RunStatus::Blocked);
7916 state.id = "20260101-000000-act2".to_owned();
7917 state.save_under(&home).unwrap();
7918
7919 let mut t = held_task_with(&state.id);
7920 queue.put(&mut t).unwrap();
7921 let mut q = action_question(&state.id, resume_action(&state.id));
7922 q.answer_delivered = true;
7923 questions.put(&mut q).unwrap();
7924
7925 apply_choice_actions(&queue, &questions, &home);
7926 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7927
7928 let mut q2 = action_question(&state.id, resume_action(&state.id));
7930 questions.put(&mut q2).unwrap();
7931 apply_choice_actions(&queue, &questions, &home);
7932 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7933 assert!(questions.get(&q2.id).unwrap().answer_delivered);
7934 }
7935}