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 deleted = queue.apply_deleted_blockers(&mut task);
694 if !deleted.is_empty() {
695 record(queue, &mut task);
696 for id in &deleted {
697 queue.note_dependency_deleted(&task, id);
698 }
699 if task.status != TaskStatus::Blocked {
700 continue;
701 }
702 }
703 let missing = crate::queue::missing_blockers(queue, questions, &task.blocked_by);
704 if !missing.is_empty() {
705 let language = language_of(&task, Path::new("."));
706 task.hold_machine(Some(crate::queue::missing_blocker_hold_reason_in(
707 &task.blocked_by,
708 &missing,
709 &language,
710 )));
711 record(queue, &mut task);
712 continue;
713 }
714 if let Some(q) = task.blocked_by.iter().find_map(|id| {
717 questions.get(id).ok().filter(|q| {
718 q.node == crate::conduct::NODE && q.status == ask::QuestionStatus::Abandoned
719 })
720 }) {
721 let language = language_of(&task, Path::new("."));
722 task.hold_machine(Some(unanswered_question_hold_reason(&q, &language)));
723 record(queue, &mut task);
724 continue;
725 }
726 let mut changed = false;
727 for id in task.blocked_by.clone() {
728 if let Ok(dep) = queue.get(&id) {
729 if dep.status == TaskStatus::Done {
730 task.unblock(&id);
731 changed = true;
732 }
733 continue;
734 }
735 if let Ok(q) = questions.get(&id)
736 && q.status == ask::QuestionStatus::Answered
737 {
738 let answer = match &q.answer {
739 Some(ask::Answer::Choice(c) | ask::Answer::Text(c)) => c.clone(),
740 None => String::new(),
741 };
742 task.record_answer(q.summary.clone(), answer);
743 task.unblock(&id);
744 changed = true;
745 }
746 }
747 if changed {
748 record(queue, &mut task);
749 }
750 }
751}
752
753fn unanswered_question_hold_reason(q: &ask::Question, language: &str) -> String {
755 if crate::lang::is_japanese(language) {
756 format!(
757 "質問 {} 「{}」 に期限内の回答がなく、取り下げられました - `magi task triage` を参照",
758 q.short(),
759 q.summary
760 )
761 } else {
762 format!(
763 "question {} \"{}\" went unanswered and was abandoned - see `magi task triage`",
764 q.short(),
765 q.summary
766 )
767 }
768}
769
770#[derive(Debug, Clone, PartialEq, Eq)]
772enum ActionDecision {
773 Skip,
775 Resume(String),
777 Requeue,
779 Done,
781 Stale,
784 Refuse(String),
787}
788
789fn decide_action<F>(task: &Task, q: &ask::Question, p: &Phrases, load: F) -> ActionDecision
797where
798 F: FnOnce(&str) -> Result<RunState>,
799{
800 let Some(action) = q.chosen_action() else {
801 return ActionDecision::Skip;
802 };
803 if task.action_applied(&q.id)
804 || matches!(
805 task.status,
806 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
807 )
808 {
809 return ActionDecision::Skip;
810 }
811 if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
814 return ActionDecision::Stale;
815 }
816 match action {
817 ask::ChoiceAction::Requeue => ActionDecision::Requeue,
818 ask::ChoiceAction::Done => ActionDecision::Done,
819 ask::ChoiceAction::Resume { run } => {
820 if task.runs.last() != Some(run) {
821 return ActionDecision::Refuse((p.resume_not_latest)(
822 q.short(),
823 ask::short_id(run),
824 ));
825 }
826 match load(run) {
827 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
828 ActionDecision::Resume(run.clone())
829 }
830 Ok(_) => ActionDecision::Refuse((p.resume_cannot_progress)(
831 q.short(),
832 ask::short_id(run),
833 )),
834 Err(e) => ActionDecision::Refuse((p.resume_unreadable)(
835 q.short(),
836 ask::short_id(run),
837 &format!("{e:#}"),
838 )),
839 }
840 }
841 }
842}
843
844struct Phrases {
858 graph_stopped: fn(&str, &str) -> String,
860 quorum_lost: &'static str,
861 quota_took_out: &'static str,
863 run_ended: &'static str,
865 waiting_for_answer: &'static str,
867 recovered_running: &'static str,
868 no_run_to_recover: &'static str,
869 could_not_start: &'static str,
870 resume_not_latest: fn(&str, &str) -> String,
871 resume_cannot_progress: fn(&str, &str) -> String,
872 resume_unreadable: fn(&str, &str, &str) -> String,
873 handover_refused: fn(&crate::handover::Refused) -> String,
876 handover_hint: &'static str,
878 reviewers_never_answered: fn(usize, usize) -> String,
881}
882
883const PHRASES_EN: Phrases = Phrases {
884 graph_stopped: |status, detail| {
885 format!("the graph stopped at `{status}` without reaching a terminal status: {detail}")
886 },
887 quorum_lost: "the judging panel lost its quorum",
888 quota_took_out: "; quota took out ",
889 run_ended: "run ended ",
890 waiting_for_answer: " - waiting for operator answer to question ",
891 recovered_running: "recovered a `running` task whose daemon never recorded the outcome: ",
892 no_run_to_recover: "task was `running` with no live daemon and no readable \
893 run to recover; held for a human to check what happened",
894 could_not_start: "could not start the run: ",
895 resume_not_latest: |q, run| {
896 format!("question {q} asked to resume run {run}, which is not this task's latest run")
897 },
898 resume_cannot_progress: |q, run| {
899 format!("question {q} asked to resume run {run}, which cannot make progress")
900 },
901 resume_unreadable: |q, run, e| {
902 format!("question {q} asked to resume run {run}, which could not be read: {e}")
903 },
904 handover_refused: |r| r.to_string(),
905 handover_hint: " (clean up the other worktree, then release the task from the queue)",
906 reviewers_never_answered: |missing, rounds| {
907 format!(
908 "{missing} reviewer seat(s) never answered after {rounds} rounds; \
909 refusing to call it clean"
910 )
911 },
912};
913
914const PHRASES_JA: Phrases = Phrases {
915 graph_stopped: |status, detail| {
916 format!("グラフが終端状態に達しないまま `{status}` で停止しました: {detail}")
917 },
918 quorum_lost: "審査パネルが定足数を失いました",
919 quota_took_out: "。クォータで脱落: ",
920 run_ended: "run 終了: ",
921 waiting_for_answer: " - オペレーターの回答待ち: 質問 ",
922 recovered_running: "daemon が結果を記録しないまま `running` だったタスクを回収しました: ",
923 no_run_to_recover: "タスクは `running` でしたが、生きた daemon も回収できる run も見つかりません。\
924 何が起きたか人が確認するため保留にしました",
925 could_not_start: "run を開始できませんでした: ",
926 resume_not_latest: |q, run| {
927 format!(
928 "質問 {q} は run {run} の再開を求めましたが、これはタスクの最新の run ではありません"
929 )
930 },
931 resume_cannot_progress: |q, run| {
932 format!("質問 {q} は run {run} の再開を求めましたが、これは進行できません")
933 },
934 resume_unreadable: |q, run, e| {
935 format!("質問 {q} は run {run} の再開を求めましたが、読み込めませんでした: {e}")
936 },
937 handover_refused: |r| {
938 use crate::handover::Refused;
939 match r {
940 Refused::Foreign { branch, path, why } => format!(
941 "ブランチ `{branch}` は {path} にチェックアウトされており、magi は自動では\
942 削除しません。不要なら `git worktree remove` でその worktree を削除して\
943 から、やり直してください(詳細: {why})"
944 ),
945 Refused::Unsafe { branch, path, why } => format!(
946 "ブランチ `{branch}` は {path} にチェックアウトされています。そこの作業を\
947 コミットか破棄したうえで `git worktree remove` で worktree を削除する\
948 か、run を破棄してよいと伝えてから、やり直してください(詳細: {why})"
949 ),
950 Refused::ReleaseFailed { branch, path, run } => format!(
951 "ブランチ `{branch}` は {path} で run {run} がチェックアウトしており、その\
952 worktree の解放に失敗したか、変更が見つかりました(未コミットの変更が\
953 ある worktree は git が削除を拒否します)。worktree はそのまま残しました"
954 ),
955 }
956 },
957 handover_hint: "(他の worktree を片付けてから、タスクをキューから解放してください)",
958 reviewers_never_answered: |missing, rounds| {
959 format!(
960 "{missing} 席のレビュアーが {rounds} ラウンドの間に一度も回答しなかったため、\
961 クリーンとは認めません"
962 )
963 },
964};
965
966fn phrases(language: &str) -> &'static Phrases {
967 if crate::lang::is_japanese(language) {
968 &PHRASES_JA
969 } else {
970 &PHRASES_EN
971 }
972}
973
974fn language_of(task: &Task, fallback: &Path) -> String {
978 crate::lang::of_repo(&repo_for(task, fallback))
979}
980
981fn task_of_question<'a>(tasks: &'a [Task], q: &ask::Question) -> Option<&'a Task> {
985 if q.node == crate::conduct::NODE {
986 return tasks.iter().find(|t| t.id == q.run);
987 }
988 tasks.iter().find(|t| t.runs.contains(&q.run))
989}
990
991fn apply_choice_actions(queue: &Queue, questions: &Questions, home: &Path) {
1000 let tasks = queue.list();
1001 for q in questions.list() {
1002 if q.chosen_action().is_none() {
1003 continue;
1004 }
1005 let Some(listed) = task_of_question(&tasks, &q) else {
1006 continue;
1007 };
1008 if listed.action_applied(&q.id) {
1009 continue;
1010 }
1011 let Ok(_claim) = queue.claim(&listed.id) else {
1012 continue;
1013 };
1014 let Ok(mut task) = queue.get(&listed.id) else {
1015 continue;
1016 };
1017 let language = language_of(&task, Path::new("."));
1018 let decision = decide_action(&task, &q, phrases(&language), |id| {
1019 RunState::load_under(id, home)
1020 });
1021 let ran = matches!(
1022 decision,
1023 ActionDecision::Resume(_) | ActionDecision::Requeue | ActionDecision::Done
1024 );
1025 if ran {
1026 if questions
1031 .read_lease(&q.id)
1032 .is_some_and(|l| l.fresh(Timestamp::now()))
1033 {
1034 continue;
1035 }
1036 let taken = questions.update(&q.id, |r| {
1037 let free = !r.answer_delivered;
1038 r.answer_delivered = true;
1039 Ok(free)
1040 });
1041 if !matches!(taken, Ok((_, true))) {
1042 continue;
1043 }
1044 }
1045 match decision {
1046 ActionDecision::Skip => continue,
1047 ActionDecision::Resume(run) => {
1048 task.release();
1049 task.resume_override = Some(crate::queue::OperatorResume {
1050 question_id: q.id.clone(),
1051 at: Timestamp::now(),
1052 conductor_rehold: None,
1053 forced: true,
1054 pinned_run: Some(run),
1055 });
1056 }
1057 ActionDecision::Stale => {}
1058 ActionDecision::Requeue => task.requeue(),
1059 ActionDecision::Done => {
1060 task.succeed();
1061 supersede_prior_runs(&task, home);
1062 }
1063 ActionDecision::Refuse(why) => {
1064 task.hold_machine(Some(why));
1065 notices::raise(
1066 Notice::warn(
1067 &format!("action:{}", q.id),
1068 "An answer asked the daemon to resume a run that cannot be resumed; the task stays held.",
1069 )
1070 .link(Link::Task {
1071 id: task.id.clone(),
1072 }),
1073 );
1074 }
1075 }
1076 task.mark_action_applied(&q.id);
1077 record(queue, &mut task);
1078 }
1079}
1080
1081fn reconcile_task_questions(queue: &Queue, questions: &Questions) {
1092 let tasks = queue.list();
1093 let by_id: std::collections::BTreeMap<&str, &Task> =
1094 tasks.iter().map(|t| (t.id.as_str(), t)).collect();
1095 let referenced: std::collections::BTreeSet<&str> = tasks
1096 .iter()
1097 .flat_map(|task| task.blocked_by.iter().map(String::as_str))
1098 .collect();
1099
1100 for mut question in questions.list() {
1101 if !question.status.open() || question.node != crate::conduct::NODE {
1102 continue;
1103 }
1104 if referenced.contains(question.id.as_str()) {
1107 continue;
1108 }
1109 let Some(task) = by_id.get(question.run.as_str()) else {
1110 continue;
1111 };
1112 question.abandon(format!(
1113 "task {} no longer waits for this answer",
1114 task.short()
1115 ));
1116 if let Err(e) = questions.put(&mut question) {
1117 tracing::warn!(
1118 "could not retire question {} for task {}: {e:#}",
1119 question.short(),
1120 task.short()
1121 );
1122 }
1123 }
1124}
1125
1126#[derive(Debug, Clone, Copy)]
1133pub struct Verdict {
1134 pub status: RunStatus,
1136 pub left_pr: bool,
1138 pub quota_hit: bool,
1140 pub parked: bool,
1142 pub no_viable_candidates: bool,
1149}
1150
1151pub fn settle(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize) {
1208 settle_in(task, verdict, detail, max_attempts, &PHRASES_EN)
1209}
1210
1211fn settle_in(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize, p: &Phrases) {
1213 if verdict.parked {
1219 task.stall(detail);
1220 return;
1221 }
1222 match verdict.status {
1223 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1224 RunStatus::AlreadyInBase => task.already_landed(detail),
1225 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1226 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1227 task.stall(detail)
1228 }
1229 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1230 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1231 RunStatus::Blocked => task.fail(detail, max_attempts),
1232 RunStatus::VerifiedNoop => task.handed_off(detail),
1233 other => task.fail((p.graph_stopped)(label(other), detail), max_attempts),
1234 }
1235}
1236
1237pub fn supersede_prior_runs(task: &Task, home: &Path) {
1286 let now = Timestamp::now();
1287 let last_run_succeeded = task
1294 .runs
1295 .last()
1296 .and_then(|id| RunState::load_under(id, home).ok())
1297 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1298 for id in task.superseded_attempts(last_run_succeeded) {
1299 let mut state = match RunState::load_under(id, home) {
1300 Ok(s) => s,
1301 Err(e) => {
1302 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1303 continue;
1304 }
1305 };
1306 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1307 continue;
1308 }
1309 let daemon_claims = is_working_on(home, id, now);
1310 if state.liveness(daemon_claims) == Liveness::Live {
1311 continue;
1312 }
1313 state.status = RunStatus::Superseded;
1314 if let Err(e) = state.save_under(home) {
1315 tracing::warn!("could not mark run {id} superseded: {e:#}");
1316 }
1317 }
1318}
1319
1320fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1342 for task in queue.list() {
1343 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1344 continue;
1345 }
1346 supersede_prior_runs(&task, home);
1347 }
1348}
1349
1350fn settle_and_diagnose(
1357 task: &mut Task,
1358 verdict: Verdict,
1359 detail: &str,
1360 max_attempts: usize,
1361 state: &RunState,
1362) {
1363 let p = phrases(&state.config.graph.language);
1364 settle_in(task, verdict, detail, max_attempts, p);
1365 if task.status == TaskStatus::Held {
1366 task.diagnostic = diagnostic(state);
1367 note_open_question(task, &state.id, p);
1368 }
1369}
1370
1371fn note_open_question(task: &mut Task, run: &str, p: &Phrases) {
1386 let Some(home) = crate::run::try_home() else {
1387 return;
1388 };
1389 let open = Questions::at(home.join("questions")).open_for(run);
1390 let Some(q) = open.first() else {
1391 return;
1392 };
1393 let base = task.hold_reason.clone().unwrap_or_default();
1394 task.hold_reason = Some(format!("{base}{}{}", p.waiting_for_answer, q.short()));
1395}
1396
1397fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize, language: &str) {
1408 match last_run {
1409 Some(state) => {
1410 let verdict = Verdict {
1411 status: state.status,
1412 left_pr: state.pr.is_some(),
1413 quota_hit: !state.quota.is_empty(),
1414 parked: state.parked,
1415 no_viable_candidates: state.viable().is_empty(),
1416 };
1417 let detail = format!(
1418 "{}{}",
1419 phrases(&state.config.graph.language).recovered_running,
1420 describe(&state)
1421 );
1422 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1423 }
1424 None => {
1425 let why = phrases(language).no_run_to_recover;
1428 task.last_error = Some(why.to_owned());
1429 task.hold_machine(Some(why.to_owned()));
1432 }
1433 }
1434}
1435
1436fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1458 let mut reclaimed = Vec::new();
1459 for listed in queue.list() {
1460 if listed.status != TaskStatus::Running {
1461 continue;
1462 }
1463 let Ok(_claim) = queue.claim(&listed.id) else {
1464 continue;
1465 };
1466 let Ok(mut task) = queue.get(&listed.id) else {
1470 continue;
1471 };
1472 if task.status != TaskStatus::Running {
1473 continue;
1474 }
1475 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1476 if let Some(state) = &last_run
1486 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1487 {
1488 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1489 }
1490 let language = if last_run.is_none() {
1491 language_of(&task, Path::new("."))
1492 } else {
1493 String::new()
1494 };
1495 reclaim(&mut task, last_run, max_attempts, &language);
1496 if task.status == TaskStatus::Done {
1497 supersede_prior_runs(&task, &crate::run::home());
1498 }
1499 record(queue, &mut task);
1500 reclaimed.push(task.id.clone());
1501 }
1502 reclaimed
1503}
1504
1505fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1537 reclaim_abandoned_runs_with(
1538 home,
1539 now,
1540 crate::proc::pid_status,
1541 crate::proc::process_started_at,
1542 )
1543}
1544
1545fn reclaim_abandoned_runs_with<F, G>(
1551 home: &Path,
1552 now: Timestamp,
1553 query: F,
1554 identity: G,
1555) -> Vec<String>
1556where
1557 F: Fn(u32) -> Option<bool>,
1558 G: Fn(u32) -> Option<String>,
1559{
1560 let mut abandoned = Vec::new();
1561 for entry in std::fs::read_dir(home.join("runs"))
1562 .into_iter()
1563 .flatten()
1564 .flatten()
1565 {
1566 let id = entry.file_name().to_string_lossy().into_owned();
1567 if !crate::run::is_run_id(&id) {
1568 continue;
1569 }
1570 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1578 continue;
1579 };
1580 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1581 continue;
1582 };
1583 if state.status.done() || !state.active_all_overrun(now) {
1584 continue;
1585 }
1586 let daemon_claims = is_working_on(home, &id, now);
1598 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1599 continue;
1600 }
1601 state.abandon("daemon");
1602 if let Err(e) = state.save_under(home) {
1603 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1604 continue;
1605 }
1606 if let Some(notice) = notices::run_ended(&state) {
1609 notices::raise_in(home, notice);
1610 }
1611 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1619 tracing::warn!("abandon questions for {id}: {e:#}");
1620 }
1621 abandoned.push(id);
1622 }
1623 abandoned
1624}
1625
1626pub async fn serve(opts: Opts) -> Result<()> {
1632 serve_until(opts, Stop::new()).await
1633}
1634
1635pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1652 let signal = {
1653 let stop = stop.clone();
1654 tokio::spawn(async move {
1655 if tokio::signal::ctrl_c().await.is_ok() {
1656 stop.stop();
1657 tracing::info!("shutdown requested; a run in flight will be finished first");
1658 }
1659 })
1660 };
1661
1662 let worktrees_root = opts
1663 .worktrees_root
1664 .clone()
1665 .unwrap_or_else(crate::run::default_worktree_root);
1666 let outcome = drive(
1667 &opts,
1668 &Queue::open(),
1669 &status_path(),
1670 &crate::run::home(),
1671 &worktrees_root,
1672 &stop,
1673 )
1674 .await;
1675
1676 signal.abort();
1677 outcome
1678}
1679
1680async fn drive(
1693 opts: &Opts,
1694 queue: &Queue,
1695 status_file: &Path,
1696 home: &Path,
1697 worktrees_root: &Path,
1698 stop: &Stop,
1699) -> Result<()> {
1700 let status = Arc::new(Mutex::new(Status::new()));
1708 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1709 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1710
1711 let daemon_cfg = prepare(&opts.repo, opts)
1717 .map(|c| c.daemon)
1718 .unwrap_or_default();
1719 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1720
1721 let waiter = tokio::spawn(crate::waiter::run(
1726 crate::waiter::Waiter::new(
1727 crate::ask::Questions::at(home.join("questions")),
1728 home.to_path_buf(),
1729 prepare(&opts.repo, opts).ok(),
1730 ),
1731 stop.clone(),
1732 ));
1733
1734 let deputies = tokio::spawn(crate::deputy::run(
1738 crate::deputy::Deputies::new(
1739 crate::ask::Questions::at(home.join("questions")),
1740 home.to_path_buf(),
1741 prepare(&opts.repo, opts).ok(),
1742 opts.repo.clone(),
1743 daemon_cfg.max_deputies,
1744 {
1745 let stop = stop.clone();
1746 Arc::new(move || stop.parking())
1747 },
1748 ),
1749 stop.clone(),
1750 ));
1751
1752 let fetcher = tokio::spawn(fetch_loop(opts.repo.clone(), opts.clone(), stop.clone()));
1755
1756 tracing::info!(
1757 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1758 queue.root().display(),
1759 opts.poll.as_secs(),
1760 opts.max_attempts,
1761 concurrency,
1762 if daemon_cfg.pause_for_interrupts {
1763 ", interrupts enabled"
1764 } else {
1765 ""
1766 }
1767 );
1768
1769 janitor(&opts.repo, opts, home, worktrees_root).await;
1772 resweep_superseded_attempts(queue, home);
1773
1774 let outcome = poll(
1775 opts,
1776 queue,
1777 &status,
1778 home,
1779 worktrees_root,
1780 stop,
1781 DispatchLimits {
1782 max_concurrent: concurrency,
1783 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1784 },
1785 )
1786 .await;
1787
1788 beat.abort();
1789 waiter.abort();
1790 deputies.abort();
1791 fetcher.abort();
1792 clear_status_at(status_file);
1793 outcome
1794}
1795
1796const FETCH_TIMEOUT: Duration = Duration::from_secs(30);
1798
1799async fn fetch_loop(repo: PathBuf, opts: Opts, stop: Stop) {
1803 while !stop.stopped() {
1804 let (interval, roots) = match prepare(&repo, &opts) {
1805 Ok(c) => (c.repos.fetch_interval, c.repos.roots),
1806 Err(_) => (0, Vec::new()),
1807 };
1808 if interval > 0 && !roots.is_empty() {
1809 let r = crate::clean::fetch_origins(&roots, FETCH_TIMEOUT, || stop.stopped()).await;
1810 tracing::debug!("fetch origins: {r:?}");
1811 }
1812 let wait = if interval > 0 { interval } else { 60 };
1814 let mut slept = 0;
1815 while slept < wait && !stop.stopped() {
1816 tokio::time::sleep(Duration::from_secs(1)).await;
1817 slept += 1;
1818 }
1819 }
1820}
1821
1822async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1828 loop {
1829 tokio::time::sleep(HEARTBEAT).await;
1830 let snapshot = {
1831 let mut guard = lock(&status);
1832 guard.updated_at = Timestamp::now();
1833 guard.clone()
1834 };
1835 if let Err(e) = write_status_to(&path, &snapshot) {
1836 tracing::warn!("could not refresh the daemon status file: {e:#}");
1839 }
1840 }
1841}
1842
1843#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1846enum LandResume {
1847 NotLanding,
1850 StillWaiting,
1855 Ready,
1859}
1860
1861fn land_resume_state(task: &Task) -> LandResume {
1865 let Some(run_id) = task.runs.last() else {
1866 return LandResume::NotLanding;
1867 };
1868 let Ok(state) = RunState::load(run_id) else {
1869 return LandResume::NotLanding;
1870 };
1871 if state.status != RunStatus::Landing || !state.parked {
1872 return LandResume::NotLanding;
1873 }
1874 let store = ask::Questions::open();
1875 let waiting = store
1876 .list()
1877 .into_iter()
1878 .filter(|q| &q.run == run_id && q.node == land::APPROVAL_NODE)
1879 .max_by(|a, b| a.id.cmp(&b.id));
1880 let Some(mut q) = waiting else {
1881 return LandResume::Ready;
1882 };
1883 if !q.status.open() {
1884 return LandResume::Ready;
1885 }
1886 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
1893 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
1894 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
1895 q.abandon(format!(
1896 "no answer within {}s of asking",
1897 timeout.as_secs().max(1)
1898 ));
1899 if store.put(&mut q).is_ok() {
1902 return LandResume::Ready;
1903 }
1904 }
1905 LandResume::StillWaiting
1906}
1907
1908const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
1916
1917const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
1929
1930struct InFlightGuard<'a> {
1943 status: &'a Arc<Mutex<Status>>,
1944 stop: &'a Stop,
1945 task_id: &'a str,
1946}
1947
1948impl Drop for InFlightGuard<'_> {
1949 fn drop(&mut self) {
1950 lock(self.status).current.retain(|c| c.task != self.task_id);
1951 self.stop.exit();
1952 }
1953}
1954
1955#[derive(Debug, Clone, PartialEq, Eq)]
1977enum Interrupt {
1978 Idle,
1980 Parking {
1992 parked: Vec<String>,
1993 interrupt_task: String,
1994 },
1995 Running {
2003 parked: Vec<String>,
2004 interrupt_task: String,
2005 },
2006 Resuming { parked: Vec<String> },
2013}
2014
2015fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
2036 match state {
2037 Interrupt::Idle => {
2038 if in_flight.len() != 1 {
2050 return Interrupt::Idle;
2051 }
2052 match runnable.iter().find(|t| t.interrupt) {
2053 Some(t) => Interrupt::Parking {
2054 parked: in_flight.to_vec(),
2055 interrupt_task: t.id.clone(),
2056 },
2057 None => Interrupt::Idle,
2058 }
2059 }
2060 Interrupt::Parking {
2061 parked,
2062 interrupt_task,
2063 } => {
2064 if in_flight.iter().any(|id| parked.contains(id)) {
2065 Interrupt::Parking {
2067 parked,
2068 interrupt_task,
2069 }
2070 } else if in_flight.contains(&interrupt_task) {
2071 Interrupt::Running {
2072 parked,
2073 interrupt_task,
2074 }
2075 } else if runnable.iter().any(|t| t.id == interrupt_task) {
2076 Interrupt::Parking {
2080 parked,
2081 interrupt_task,
2082 }
2083 } else {
2084 Interrupt::Resuming { parked }
2089 }
2090 }
2091 Interrupt::Running {
2092 parked,
2093 interrupt_task,
2094 } => {
2095 if in_flight.contains(&interrupt_task) {
2096 Interrupt::Running {
2097 parked,
2098 interrupt_task,
2099 }
2100 } else {
2101 Interrupt::Resuming { parked }
2107 }
2108 }
2109 Interrupt::Resuming { parked } => {
2110 if in_flight.iter().any(|id| parked.contains(id)) {
2111 Interrupt::Idle
2117 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
2118 Interrupt::Resuming { parked }
2119 } else {
2120 Interrupt::Idle
2123 }
2124 }
2125 }
2126}
2127
2128fn advance_interrupt_tick(
2134 enabled: bool,
2135 state: Interrupt,
2136 in_flight: &[String],
2137 runnable: &[Task],
2138) -> Interrupt {
2139 if !enabled {
2140 return Interrupt::Idle;
2141 }
2142 advance_interrupt(state, in_flight, runnable)
2143}
2144
2145fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
2150 match state {
2151 Interrupt::Idle => candidates,
2152 Interrupt::Parking {
2153 parked,
2154 interrupt_task,
2155 } => {
2156 if in_flight.iter().any(|id| parked.contains(id)) {
2157 Vec::new()
2158 } else {
2159 candidates
2160 .into_iter()
2161 .filter(|t| &t.id == interrupt_task)
2162 .collect()
2163 }
2164 }
2165 Interrupt::Running { .. } => Vec::new(),
2166 Interrupt::Resuming { parked } => candidates
2174 .into_iter()
2175 .find(|t| parked.contains(&t.id))
2176 .into_iter()
2177 .collect(),
2178 }
2179}
2180
2181struct DispatchLimits {
2185 max_concurrent: usize,
2188 pause_for_interrupts: bool,
2190}
2191
2192#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2201enum PermitKind {
2202 None,
2207 Urgent,
2214 Ordinary,
2217}
2218
2219fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
2222 if priority {
2223 PermitKind::None
2224 } else if urgent {
2225 PermitKind::Urgent
2226 } else {
2227 PermitKind::Ordinary
2228 }
2229}
2230
2231async fn poll(
2250 opts: &Opts,
2251 queue: &Queue,
2252 status: &Arc<Mutex<Status>>,
2253 home: &Path,
2254 worktrees_root: &Path,
2255 stop: &Stop,
2256 limits: DispatchLimits,
2257) -> Result<()> {
2258 let DispatchLimits {
2259 max_concurrent,
2260 pause_for_interrupts,
2261 } = limits;
2262 let mut attempted: Vec<String> = Vec::new();
2267 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2268 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2276 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2282 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2283 let mut conductor = Conductor::new();
2284 let mut cache_last_checked: Option<Timestamp> = None;
2287 let mut interrupt = Interrupt::Idle;
2289 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2295 std::collections::HashMap::new();
2296
2297 while !stop.stopped() {
2298 lock(status).polls += 1;
2299
2300 while let Some(result) = inflight.try_join_next() {
2305 if let Err(e) = result {
2306 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2307 notices::raise(Notice::error(
2308 "loop:attempt",
2309 "A queued attempt ended abnormally; check the task it was running.",
2310 ));
2311 }
2312 }
2313
2314 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2315 if !swept.is_empty() {
2316 tracing::warn!(
2317 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2318 swept.len(),
2319 swept.join(", ")
2320 );
2321 }
2322 let now = Timestamp::now();
2327
2328 if !stop.busy_now() {
2333 maybe_prune_cache_between_runs(
2334 &opts.repo,
2335 opts,
2336 home,
2337 stop,
2338 &mut cache_last_checked,
2339 now,
2340 )
2341 .await;
2342 }
2343
2344 let stalled = stalled_tasks(queue, home, now);
2345 let stalled_ids: std::collections::BTreeSet<_> =
2346 stalled.iter().map(|task| task.id.clone()).collect();
2347 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2348 if !reclaimed.is_empty() {
2349 tracing::warn!(
2350 "reclaimed {} task(s) left `running` by a daemon that never \
2351 recorded the outcome: {}",
2352 reclaimed.len(),
2353 reclaimed.join(", ")
2354 );
2355 }
2356 let abandoned_runs = reclaim_abandoned_runs(home, now);
2357 if !abandoned_runs.is_empty() {
2358 tracing::warn!(
2359 "failed {} run(s) left behind by a killed process, past every \
2360 active seat's own timeout: {}",
2361 abandoned_runs.len(),
2362 abandoned_runs.join(", ")
2363 );
2364 }
2365
2366 let questions = Questions::at(home.join("questions"));
2371
2372 resolve_blockers(queue, &questions);
2375 apply_choice_actions(queue, &questions, home);
2376 reconcile_task_questions(queue, &questions);
2377
2378 let finished: Vec<Task> = finished_tasks(queue)
2384 .into_iter()
2385 .filter(|task| !stalled_ids.contains(&task.id))
2386 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2390 .collect();
2391 let queued = queued_tasks(queue);
2392 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2396 && conductor.worth_a_look(queue, &stalled, &finished)
2397 {
2398 match prepare(&opts.repo, opts) {
2399 Ok(cfg) => {
2400 conductor
2401 .maybe_run(
2402 &cfg,
2403 &opts.repo,
2404 queue,
2405 &questions,
2406 home,
2407 &queued,
2408 &stalled,
2409 &finished,
2410 opts.max_attempts,
2411 )
2412 .await;
2413 }
2414 Err(e) => {
2415 tracing::warn!("conductor: no config: {e:#}");
2416 notices::raise(Notice::warn(
2417 "loop:no-config",
2418 "The loop could not read this repository's config, so held tasks are not being triaged.",
2419 ));
2420 }
2421 }
2422 }
2423
2424 let candidates: Vec<Task> = runnable(queue)
2425 .into_iter()
2426 .filter(|t| !opts.once || !attempted.contains(&t.id))
2427 .collect();
2428
2429 let in_flight: Vec<String> = lock(status)
2434 .current
2435 .iter()
2436 .map(|c| c.task.clone())
2437 .collect();
2438 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2439
2440 interrupt =
2441 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2442 if let Interrupt::Parking {
2443 parked,
2444 interrupt_task,
2445 } = &interrupt
2446 {
2447 let reason = format!(
2448 "task {} asked to run first",
2449 crate::run::short_of(interrupt_task)
2450 );
2451 for id in parked {
2452 if let Some(pause) = interrupt_pauses.get(id) {
2453 pause.park_because(reason.clone());
2454 }
2455 }
2456 }
2457 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2458
2459 let cooling_down =
2460 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2461
2462 let mut started_any = false;
2463 for candidate in candidates {
2464 if stop.stopped() {
2465 break;
2466 }
2467
2468 let resume = land_resume_state(&candidate);
2469 if resume == LandResume::StillWaiting {
2470 continue;
2471 }
2472 let priority = resume == LandResume::Ready;
2473
2474 if !priority && cooling_down {
2475 continue;
2476 }
2477 let permit = match permit_kind(priority, candidate.urgent) {
2478 PermitKind::None => None,
2479 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2480 Ok(p) => Some(p),
2481 Err(_) => continue,
2487 },
2488 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2489 Ok(p) => Some(p),
2490 Err(_) => continue,
2494 },
2495 };
2496
2497 let Ok(claim) = queue.claim(&candidate.id) else {
2502 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2503 continue;
2504 };
2505 let mut task = match queue.get(&candidate.id) {
2508 Ok(t) if t.status.runnable() => t,
2509 Ok(_) => continue,
2510 Err(e) => {
2511 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2512 continue;
2513 }
2514 };
2515 let task_id = task.id.clone();
2516 attempted.push(task_id.clone());
2517 lock(status).idle = false;
2518 stop.enter();
2521 started_any = true;
2522
2523 let run_pause = crate::graph::Pause::new();
2527 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2528
2529 let opts = opts.clone();
2530 let queue = queue.clone();
2531 let status = Arc::clone(status);
2532 let stop = stop.clone();
2533 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2534 inflight.spawn(async move {
2535 let _claim = claim;
2539 let _permit = permit;
2540 let _inflight = InFlightGuard {
2542 status: &status,
2543 stop: &stop,
2544 task_id: &task_id,
2545 };
2546 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2547 lock(&status).completed += 1;
2548 let now = Timestamp::now();
2554 if let Some(until) = cooldown_until("a, now) {
2555 let wait = until.as_second() - now.as_second();
2556 *lock("a_cooldown_until) = Some(until);
2557 let hint = quota
2558 .iter()
2559 .find(|q| q.reset.is_some())
2560 .and_then(|q| q.reset.as_deref());
2561 match hint {
2562 Some(h) => tracing::warn!(
2563 "quota hit; waiting {wait}s before taking another ordinary task \
2564 (CLI reported reset: {h})"
2565 ),
2566 None => tracing::warn!(
2567 "quota hit; waiting {wait}s before taking another ordinary task \
2568 (no reset hint reported)"
2569 ),
2570 }
2571 }
2572 });
2573 }
2574
2575 if started_any {
2576 continue;
2577 }
2578
2579 if stop.busy_now() {
2580 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2585 continue;
2586 }
2587
2588 lock(status).idle = true;
2590 if opts.once {
2591 janitor(&opts.repo, opts, home, worktrees_root).await;
2595 resweep_superseded_attempts(queue, home);
2596 triage_held(queue, home, opts).await;
2597 break;
2598 }
2599 stop.idle(opts.poll).await;
2600 if stop.stopped() {
2601 continue;
2602 }
2603 janitor(&opts.repo, opts, home, worktrees_root).await;
2609 resweep_superseded_attempts(queue, home);
2610 triage_held(queue, home, opts).await;
2611 }
2612
2613 while let Some(result) = inflight.join_next().await {
2618 if let Err(e) = result {
2619 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2620 notices::raise(Notice::error(
2621 "loop:attempt",
2622 "A queued attempt ended abnormally; check the task it was running.",
2623 ));
2624 }
2625 }
2626 Ok(())
2627}
2628
2629async fn attempt(
2635 opts: &Opts,
2636 queue: &Queue,
2637 status: &Arc<Mutex<Status>>,
2638 stop: &Stop,
2639 interrupt_pause: crate::graph::Pause,
2640 task: &mut Task,
2641) -> Vec<QuotaLoss> {
2642 let repo = repo_for(task, &opts.repo);
2643 tracing::info!(
2644 "task {} — {} (repo {})",
2645 task.short(),
2646 task.title,
2647 repo.display()
2648 );
2649
2650 let mut config = match prepare(&repo, opts) {
2651 Ok(c) => c,
2652 Err(e) => {
2653 task.attempts += 1;
2657 task.fail(format!("config: {e:#}"), opts.max_attempts);
2658 record(queue, task);
2659 return Vec::new();
2660 }
2661 };
2662 apply_solo(&mut config, task);
2663 let p = phrases(&config.graph.language);
2664 let start_failed = p.could_not_start;
2665
2666 if let Some(reason) = disk_gate(&repo, &config) {
2674 task.last_error = Some(reason.clone());
2675 task.hold_machine(Some(reason.clone()));
2676 record(queue, task);
2677 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2678 notices::raise(
2681 Notice::warn(
2682 &format!("disk:{}", repo.display()),
2683 "A task was held for want of free disk space; free some, then release it from the queue.",
2684 )
2685 .link(Link::Task {
2686 id: task.id.clone(),
2687 }),
2688 );
2689 return Vec::new();
2690 }
2691
2692 let unfinished = (!task.fresh_start)
2712 .then(|| unfinished_run(&task.runs, task.short()))
2713 .flatten();
2714 let review_branch = task.review_branch.take();
2720 let branch_exists = match &review_branch {
2721 Some(branch) => crate::git::branch_exists(&repo, branch)
2722 .await
2723 .unwrap_or(false),
2724 None => false,
2725 };
2726 let starter = choose_starter(
2727 review_branch.as_deref(),
2728 branch_exists,
2729 unfinished.as_deref(),
2730 );
2731 let attachments = match task_attachments(queue, task) {
2736 Ok(a) => a,
2737 Err(e) => {
2738 task.attempts += 1;
2739 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2740 record(queue, task);
2741 return Vec::new();
2742 }
2743 };
2744 let started = match &starter {
2745 Starter::Review(branch) => {
2746 tracing::info!(
2747 "task {} reopens `{branch}` as a review-only pass",
2748 task.short()
2749 );
2750 let takeover = crate::handover::Takeover {
2753 earlier: task.earlier_attempts().to_vec(),
2754 home: crate::run::home(),
2755 choice: take_divergence_answer(branch, &config.merge.remote, task),
2756 };
2757 Runner::review_taking_over(
2758 &repo,
2759 branch,
2760 config,
2761 Some(takeover),
2762 crate::run::Origin::queue(&task.id),
2763 )
2764 .await
2765 }
2766 Starter::Resume(id) => {
2767 tracing::info!("resuming run {id} rather than competing again");
2768 Runner::resume(id).map(|mut r| {
2769 if let Some(instruction) =
2770 prepare_instruction(&starter, Some(&r.state.instruction), task)
2771 {
2772 r.state.instruction = instruction;
2773 }
2774 r.state.attachments = attachments.clone();
2775 r
2776 })
2777 }
2778 Starter::Start => {
2779 if let Some(branch) = &review_branch {
2780 tracing::warn!(
2781 "conductor chose review for task {} but branch `{branch}` no longer \
2782 exists; requeuing as a fresh competition instead",
2783 task.short()
2784 );
2785 }
2786 let instruction = prepare_instruction(&starter, None, task)
2787 .unwrap_or_else(|| task.instruction.clone());
2788 Runner::start_naming(
2789 &repo,
2790 instruction,
2791 &task.title,
2792 config,
2793 crate::run::Origin::queue(&task.id),
2794 )
2795 .await
2796 .map(|mut r| {
2797 r.state.attachments = attachments.clone();
2798 r
2799 })
2800 }
2801 };
2802 let mut runner = match started {
2803 Ok(r) => r,
2804 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2808 let detail = match e.downcast_ref::<crate::handover::Refused>() {
2809 Some(r) => (p.handover_refused)(r),
2810 None => format!("{e:#}"),
2811 };
2812 let reason = format!("{start_failed}{detail}");
2813 task.last_error = Some(reason.clone());
2814 let branch = match &starter {
2817 Starter::Review(branch) => Some(branch.clone()),
2818 _ => None,
2819 };
2820 task.hold_for_handover(branch, format!("{reason}{}", p.handover_hint));
2824 record(queue, task);
2825 tracing::warn!(
2826 "holding {} for a branch it cannot take over: {e:#}",
2827 task.short()
2828 );
2829 return Vec::new();
2830 }
2831 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2836 let d = e
2837 .downcast_ref::<crate::reconcile::Diverged>()
2838 .expect("checked by the guard");
2839 let mut q = ask::Question::new(
2840 task.id.clone(),
2841 "review".to_owned(),
2842 "sync".to_owned(),
2843 d.summary(),
2844 d.detail(),
2845 d.choices(),
2846 );
2847 match Questions::open().put(&mut q) {
2848 Ok(()) => {
2849 task.last_error = Some(format!("{e:#}"));
2850 task.review_branch = Some(d.branch.clone());
2853 task.block(vec![q.id.clone()], Some(d.summary()));
2854 }
2855 Err(put) => {
2856 tracing::warn!("could not file the divergence question: {put:#}");
2857 task.attempts += 1;
2858 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2859 }
2860 }
2861 record(queue, task);
2862 return Vec::new();
2863 }
2864 Err(e) => {
2865 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
2871 && let Starter::Review(branch) = &starter
2872 {
2873 task.review_branch = Some(branch.clone());
2874 task.last_error = Some(format!("{start_failed}{e:#}"));
2875 task.status = crate::queue::TaskStatus::Failed;
2876 record(queue, task);
2877 return Vec::new();
2878 }
2879 task.attempts += 1;
2880 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2881 record(queue, task);
2882 return Vec::new();
2883 }
2884 };
2885 runner.state.followup_generation = Some(task.followup.as_ref().map_or(0, |f| f.generation));
2888 runner.on_pause(stop.pause());
2890 runner.watch_interrupt(interrupt_pause);
2894
2895 let run = runner.state.id.clone();
2898 task.start(run.clone());
2899 record(queue, task);
2900 lock(status).current.push(Current {
2901 task: task.id.clone(),
2902 run,
2903 });
2904
2905 let quota_before = runner.state.quota.clone();
2908 let detail = match runner.execute().await {
2909 Ok(()) => describe(&runner.state),
2910 Err(e) => format!("{e:#}"),
2911 };
2912 let fresh = losses_this_attempt("a_before, &runner.state.quota);
2913 let verdict = Verdict {
2914 status: runner.state.status,
2915 left_pr: runner.state.pr.is_some(),
2918 quota_hit: !fresh.is_empty(),
2924 parked: runner.state.parked,
2928 no_viable_candidates: runner.state.viable().is_empty(),
2931 };
2932 settle_and_diagnose(task, verdict, &detail, opts.max_attempts, &runner.state);
2933 if task.status == TaskStatus::Done {
2934 supersede_prior_runs(task, &crate::run::home());
2935 }
2936 record(queue, task);
2937 tracing::info!(
2938 "task {} is {} after run {} ({})",
2939 task.short(),
2940 task.status.as_str(),
2941 runner.state.short(),
2942 label(runner.state.status)
2943 );
2944 fresh
2945}
2946
2947fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
2957 after
2958 .iter()
2959 .filter(|q| !before.contains(q))
2960 .cloned()
2961 .collect()
2962}
2963
2964fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
2967 if quota.is_empty() {
2968 return None;
2969 }
2970 let with_hint = quota.iter().find(|q| q.reset.is_some());
2971 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
2972 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
2973 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
2974 Some(
2975 now.checked_add(jiff::SignedDuration::from_secs(secs))
2976 .unwrap_or(Timestamp::MAX),
2977 )
2978}
2979
2980fn apply_solo(config: &mut Config, task: &Task) {
2990 if task.solo {
2991 config.graph.candidates = 1;
2992 }
2993}
2994
2995fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
2997 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
2998 if let Some(mode) = &opts.merge {
2999 config.merge.mode = merge_mode(mode)?;
3000 }
3001 Ok(config)
3002}
3003
3004async fn maybe_prune_cache_between_runs(
3035 repo: &Path,
3036 opts: &Opts,
3037 home: &Path,
3038 stop: &Stop,
3039 last_checked: &mut Option<Timestamp>,
3040 now: Timestamp,
3041) {
3042 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
3043 return;
3044 }
3045 *last_checked = Some(now);
3046 let cfg = match prepare(repo, opts) {
3047 Ok(cfg) => cfg,
3048 Err(e) => {
3049 tracing::warn!("cache check: no config: {e:#}");
3050 return;
3051 }
3052 };
3053 match clean::prune_cache_if_over_limit(&cfg, home) {
3054 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
3055 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
3056 pruned.files,
3057 pruned.freed
3058 ),
3059 Ok(_) => {}
3060 Err(e) => {
3061 tracing::warn!("housekeep: prune cache: {e:#}");
3062 notices::raise_in(
3063 home,
3064 Notice::warn(
3065 "housekeep:cache",
3066 "Pruning the shared build cache failed; disk usage may keep growing.",
3067 ),
3068 );
3069 }
3070 }
3071}
3072
3073fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
3077 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
3078}
3079
3080async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
3099 let cfg = match prepare(repo, opts) {
3100 Ok(cfg) => cfg,
3101 Err(e) => {
3102 tracing::warn!("housekeep: no config: {e:#}");
3103 return;
3104 }
3105 };
3106 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
3115 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
3116 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
3121 let mut extra = Vec::new();
3122 if out.unreadable > 0 {
3123 extra.push(format!("{} unreadable", out.unreadable));
3124 }
3125 if out.orphaned_worktrees > 0 {
3126 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
3127 }
3128 let detail = if extra.is_empty() {
3129 String::new()
3130 } else {
3131 format!(" ({})", extra.join(", "))
3132 };
3133 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
3134 }
3135 if out.external_merges_recorded > 0 {
3136 tracing::info!(
3137 "housekeep: recorded {} run(s) as merged externally",
3138 out.external_merges_recorded
3139 );
3140 }
3141 if out.stale_pr_states_repaired > 0 {
3142 tracing::info!(
3143 "housekeep: rewrote {} run record(s) whose pull request had already settled",
3144 out.stale_pr_states_repaired
3145 );
3146 }
3147 if out.cache_files > 0 {
3148 tracing::info!(
3149 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
3150 out.cache_files,
3151 out.cache_freed
3152 );
3153 }
3154 if out.questions_abandoned > 0 {
3155 tracing::info!(
3156 "housekeep: abandoned {} question(s) left open by a finished run",
3157 out.questions_abandoned
3158 );
3159 }
3160}
3161
3162async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3171 let questions = Questions::at(home.join("questions"));
3172 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3173 if report.is_empty() {
3174 return;
3175 }
3176 if !report.quarantined.is_empty() {
3177 tracing::info!(
3178 "triage: held {} blocked task(s) whose blocked-on task or \
3179 question no longer exists: {}",
3180 report.quarantined.len(),
3181 report.quarantined.join(", ")
3182 );
3183 }
3184 if !report.resumed.is_empty() {
3185 tracing::info!(
3186 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3187 report.resumed.len(),
3188 report.resumed.join(", ")
3189 );
3190 }
3191 if !report.asked.is_empty() {
3192 tracing::info!(
3193 "triage: asked about {} held task(s): {}",
3194 report.asked.len(),
3195 report.asked.join(", ")
3196 );
3197 }
3198 if !report.answered.is_empty() {
3199 tracing::info!(
3200 "triage: applied {} operator answer(s): {}",
3201 report.answered.len(),
3202 report.answered.join(", ")
3203 );
3204 }
3205}
3206
3207fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3214 disk_gate_with(repo, config, crate::disk::free_bytes)
3215}
3216
3217fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3221 repo: &Path,
3222 config: &Config,
3223 free_bytes: F,
3224) -> Option<String> {
3225 let min = config.disk.min_free_bytes;
3226 if min == 0 {
3227 return None;
3228 }
3229 match free_bytes(repo) {
3230 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3231 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3232 }
3233}
3234
3235const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3242
3243const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3247
3248fn quota_wait(
3257 reset_at: Option<Timestamp>,
3258 now: Timestamp,
3259 fallback: Duration,
3260 cap: Duration,
3261) -> Duration {
3262 match reset_at {
3263 Some(at) if at > now => {
3264 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3265 Duration::from_secs(secs).min(cap)
3266 }
3267 _ => fallback,
3268 }
3269}
3270
3271fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3285 parse_reset_hint_zoned(text, now)
3286 .or_else(|| parse_reset_hint_dated(text))
3287 .or_else(|| parse_reset_hint_relative(text, recorded))
3288}
3289
3290fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3294 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3295 let mut rest = rest.trim();
3296 if rest.is_empty() {
3297 return None;
3298 }
3299 let mut total: i64 = 0;
3300 let mut matched = false;
3301 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3302 if let Some((digits, tail)) = rest.split_once(unit)
3303 && !digits.is_empty()
3304 && digits.bytes().all(|b| b.is_ascii_digit())
3305 {
3306 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3307 rest = tail;
3308 matched = true;
3309 }
3310 }
3311 if !rest.is_empty() || !matched {
3312 return None;
3313 }
3314 recorded
3315 .checked_add(jiff::SignedDuration::from_secs(total))
3316 .ok()
3317}
3318
3319fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3323 let clock = clock.trim().to_lowercase();
3324 let (digits, pm) = clock
3325 .strip_suffix("am")
3326 .map(|d| (d, false))
3327 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3328 let (h, m) = digits.trim().split_once(':')?;
3329 let mut hour: i8 = h.trim().parse().ok()?;
3330 let minute: i8 = m.trim().parse().ok()?;
3331 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3332 return None;
3333 }
3334 if pm && hour != 12 {
3335 hour += 12;
3336 } else if !pm && hour == 12 {
3337 hour = 0;
3338 }
3339 Some((hour, minute))
3340}
3341
3342fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3347 let open = text.find('(')?;
3348 let close = text.rfind(')')?;
3349 if close <= open {
3350 return None;
3351 }
3352 let zone = text[open + 1..close].trim();
3353 let (hour, minute) = parse_12h_clock(&text[..open])?;
3354 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3355 let candidate = now
3356 .to_zoned(tz)
3357 .with()
3358 .hour(hour)
3359 .minute(minute)
3360 .second(0)
3361 .millisecond(0)
3362 .microsecond(0)
3363 .nanosecond(0)
3364 .build()
3365 .ok()?;
3366 let mut at = candidate.timestamp();
3367 if at <= now {
3368 at += jiff::SignedDuration::from_hours(24);
3369 }
3370 Some(at)
3371}
3372
3373fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3382 let words: Vec<&str> = text.split_whitespace().collect();
3383 if words.len() < 5 {
3384 return None;
3385 }
3386 (0..=words.len() - 5)
3387 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3388}
3389
3390fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3396 if trailing.is_some_and(|next| next.starts_with('(')) {
3397 return None;
3398 }
3399 let month = month_number(window[0])?;
3400 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3401 let day_digits = ["st", "nd", "rd", "th"]
3402 .iter()
3403 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3404 let day: i8 = day_digits.parse().ok()?;
3405 let year_token = window[2];
3406 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3407 return None;
3408 }
3409 let year: i16 = year_token.parse().ok()?;
3410 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3414 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3415 let date = jiff::civil::Date::new(year, month, day).ok()?;
3416 let candidate = date
3417 .at(hour, minute, 0, 0)
3418 .to_zoned(jiff::tz::TimeZone::UTC)
3419 .ok()?;
3420 Some(candidate.timestamp())
3421}
3422
3423fn month_number(name: &str) -> Option<i8> {
3426 const NAMES: [&str; 12] = [
3427 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3428 ];
3429 let lower = name.to_lowercase();
3430 NAMES
3431 .iter()
3432 .position(|n| *n == lower.as_str())
3433 .map(|i| i as i8 + 1)
3434}
3435
3436fn exhausted_review_budget(state: &RunState) -> bool {
3448 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3449}
3450
3451fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3488 unfinished_run_with(runs, short, RunState::load)
3489}
3490
3491fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3494where
3495 F: FnOnce(&str) -> Result<RunState>,
3496{
3497 let id = runs.last()?;
3498 match load(id) {
3499 Ok(s)
3504 if s.status.resumable()
3505 && !s.released()
3506 && !exhausted_review_budget(&s)
3507 && s.liveness(false) != crate::run::Liveness::Live =>
3508 {
3509 Some(id.clone())
3510 }
3511 Ok(_) => None,
3512 Err(e) => {
3513 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3514 None
3515 }
3516 }
3517}
3518
3519fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3527where
3528 F: FnOnce(&str) -> Result<RunState>,
3529{
3530 if task.status != TaskStatus::Failed || task.fresh_start || task.review_branch.is_some() {
3531 return false;
3532 }
3533 let Some(id) = task.runs.last() else {
3534 return false;
3535 };
3536 load(id).is_ok_and(|s| {
3537 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3538 })
3539}
3540
3541#[derive(Debug, Clone, PartialEq, Eq)]
3544enum Starter {
3545 Review(String),
3548 Resume(String),
3550 Start,
3552}
3553
3554fn take_divergence_answer(
3558 branch: &str,
3559 remote: &str,
3560 task: &mut Task,
3561) -> Option<crate::reconcile::Choice> {
3562 let summary = crate::reconcile::summary_for(branch, remote);
3563 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3564 (a.question == summary)
3565 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3566 .flatten()
3567 .map(|c| (i, c))
3568 })?;
3569 task.answers.remove(idx);
3570 Some(choice)
3571}
3572
3573fn choose_starter(
3585 review_branch: Option<&str>,
3586 branch_exists: bool,
3587 unfinished: Option<&str>,
3588) -> Starter {
3589 match review_branch {
3590 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3591 Some(_) => Starter::Start,
3592 None => match unfinished {
3593 Some(id) => Starter::Resume(id.to_owned()),
3594 None => Starter::Start,
3595 },
3596 }
3597}
3598
3599fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3602 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3603 return fallback.to_path_buf();
3604 }
3605 task.repo.clone()
3606}
3607
3608const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3612
3613fn answers_block(task: &Task, count: usize) -> String {
3615 let mut s = ANSWERS_HEADER.to_owned();
3616 for a in &task.answers[..count] {
3617 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3618 }
3619 s
3620}
3621
3622fn append_answers(base: &str, task: &Task) -> String {
3625 if task.answers.is_empty() {
3626 return base.to_owned();
3627 }
3628 let mut s = base.to_owned();
3629 s.push_str(&answers_block(task, task.answers.len()));
3630 s
3631}
3632
3633fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3637 for count in (1..=task.answers.len()).rev() {
3638 let block = answers_block(task, count);
3639 if let Some(base) = instruction.strip_suffix(&block) {
3640 return base;
3641 }
3642 }
3643 instruction
3644}
3645
3646fn instruction_for(task: &Task) -> String {
3654 append_answers(&task.instruction, task)
3655}
3656
3657fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3669 append_answers(strip_answers_block(old_instruction, task), task)
3670}
3671
3672fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3674 let paths = queue.attachment_paths(task);
3675 for (name, path) in task.attachments.iter().zip(&paths) {
3676 if !path.is_file() {
3677 bail!(
3678 "attachment `{name}` is recorded on the task but {} is missing",
3679 path.display()
3680 );
3681 }
3682 }
3683 Ok(paths)
3684}
3685
3686fn prepare_instruction(
3697 starter: &Starter,
3698 old_instruction: Option<&str>,
3699 task: &Task,
3700) -> Option<String> {
3701 match starter {
3702 Starter::Start => Some(instruction_for(task)),
3703 Starter::Resume(_) => Some(resumed_instruction(
3704 old_instruction.expect("a resumed run always has a prior instruction"),
3705 task,
3706 )),
3707 Starter::Review(_) => None,
3708 }
3709}
3710
3711fn record(queue: &Queue, task: &mut Task) {
3715 if let Err(e) = queue.put(task) {
3716 tracing::error!("could not record task {}: {e:#}", task.short());
3717 notices::raise(Notice::error(
3718 "loop:record",
3719 "The loop could not save a task's state; check the disk.",
3720 ));
3721 }
3722}
3723
3724fn runnable(queue: &Queue) -> Vec<Task> {
3730 let mut tasks: Vec<Task> = queue
3731 .list()
3732 .into_iter()
3733 .filter(|t| t.status.runnable())
3734 .collect();
3735 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
3736 tasks
3737}
3738
3739fn describe(state: &RunState) -> String {
3753 let p = phrases(&state.config.graph.language);
3754 let mut detail = if state.status == RunStatus::Stalled {
3755 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
3756 seats.sort_unstable();
3757 seats.dedup();
3758 if seats.is_empty() {
3759 p.quorum_lost.to_owned()
3760 } else {
3761 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
3762 }
3763 } else {
3764 format!("{}{}", p.run_ended, state.status.display_label())
3765 };
3766 if let Some(last) = state.events.last() {
3767 let unanswered = state
3771 .reviews
3772 .last()
3773 .filter(|r| {
3774 state.status == RunStatus::Blocked
3775 && last.node == "review"
3776 && r.incomplete()
3777 && r.blocking == 0
3778 && r.round == state.config.graph.review_rounds
3779 && r.e2e.iter().all(crate::run::CommandOutcome::ok)
3780 })
3781 .map(|r| (r.expected - r.answered, r.round));
3782 match unanswered {
3783 Some((missing, rounds)) if crate::lang::is_japanese(&state.config.graph.language) => {
3784 detail.push_str(&format!(
3785 " ({}: {})",
3786 last.node,
3787 (p.reviewers_never_answered)(missing, rounds)
3788 ));
3789 }
3790 _ => detail.push_str(&format!(" ({}: {})", last.node, last.message)),
3791 }
3792 }
3793 detail.push_str(&format!(" [run {}]", state.id));
3794 detail
3795}
3796
3797const DIAGNOSTIC_MAX: usize = 4_000;
3803
3804const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
3809
3810fn diagnostic(state: &RunState) -> Option<String> {
3824 let mut parts: Vec<String> = Vec::new();
3825
3826 for o in state.gate.iter().filter(|o| !o.ok()) {
3828 parts.push(format!(
3829 "gate `{}` failed ({:?}):\n{}",
3830 o.command,
3831 o.code,
3832 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
3833 ));
3834 }
3835
3836 if let Some(last) = state
3839 .events
3840 .iter()
3841 .rev()
3842 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
3843 {
3844 parts.push(last.message.clone());
3845 }
3846
3847 if state.viable().is_empty() {
3854 for c in &state.candidates {
3855 if let Some(evidence) = &c.verified_noop {
3856 parts.push(format!(
3857 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
3858 c.label
3859 ));
3860 } else if !c.summary.trim().is_empty() {
3861 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
3862 } else if let Some(why) = &c.failed {
3863 parts.push(format!("candidate {}: {why}", c.label));
3864 }
3865 }
3866 }
3867
3868 if parts.is_empty() {
3869 return None;
3870 }
3871 Some(crate::run::tail(
3876 &parts.join("\n\n"),
3877 DIAGNOSTIC_MAX.saturating_sub(100),
3878 ))
3879}
3880
3881fn label(status: RunStatus) -> &'static str {
3889 status.as_str()
3890}
3891
3892fn merge_mode(mode: &str) -> Result<MergeMode> {
3894 match mode {
3895 "none" => Ok(MergeMode::None),
3896 "local" => Ok(MergeMode::Local),
3897 "pr" => Ok(MergeMode::Pr),
3898 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
3899 }
3900}
3901
3902fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
3907 mutex
3908 .lock()
3909 .unwrap_or_else(std::sync::PoisonError::into_inner)
3910}
3911
3912#[cfg(test)]
3913mod tests {
3914 use super::*;
3915 use crate::queue::{Source, TaskStatus};
3916 use crate::run::{Candidate, CommandOutcome};
3917 use pretty_assertions::assert_eq;
3918
3919 fn task() -> Task {
3920 Task::new(
3921 "add retries".to_owned(),
3922 "add retries".to_owned(),
3923 PathBuf::from("/repo"),
3924 Source::Human,
3925 )
3926 }
3927
3928 fn interrupt_task(id: &str) -> Task {
3931 let mut t = task();
3932 t.id = id.to_owned();
3933 t.interrupt = true;
3934 t
3935 }
3936
3937 fn task_with_id(id: &str) -> Task {
3939 let mut t = task();
3940 t.id = id.to_owned();
3941 t
3942 }
3943
3944 fn urgent_task(id: &str) -> Task {
3946 let mut t = task();
3947 t.id = id.to_owned();
3948 t.urgent = true;
3949 t
3950 }
3951
3952 #[test]
3957 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
3958 assert_eq!(permit_kind(true, false), PermitKind::None);
3959 assert_eq!(permit_kind(true, true), PermitKind::None);
3960 }
3961
3962 #[test]
3967 fn permit_kind_separates_urgent_from_ordinary() {
3968 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
3969 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
3970 }
3971
3972 #[test]
3979 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
3980 let cfg = Config::default();
3981 let repo = Path::new("/any/repo/path");
3982
3983 let reason =
3984 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
3985 assert!(reason.contains("1024"), "{reason}");
3986 assert!(
3987 reason.contains(&cfg.disk.min_free_bytes.to_string()),
3988 "{reason}"
3989 );
3990
3991 assert_eq!(
3992 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
3993 None,
3994 "exactly at the floor is open"
3995 );
3996 assert_eq!(
3997 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
3998 None,
3999 "comfortably above the floor is open"
4000 );
4001 }
4002
4003 #[test]
4004 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
4005 let mut cfg = Config::default();
4006 cfg.disk.min_free_bytes = 0;
4007 let repo = Path::new("/any/repo/path");
4008 assert_eq!(
4009 disk_gate_with(repo, &cfg, |_| Ok(0)),
4010 None,
4011 "a zero floor never measures at all"
4012 );
4013 }
4014
4015 #[test]
4016 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
4017 let cfg = Config::default();
4018 let repo = Path::new("/any/repo/path");
4019 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
4020 .expect("a measurement failure must close the gate, not open it");
4021 assert!(reason.contains("could not measure"), "{reason}");
4022 }
4023
4024 #[test]
4025 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
4026 let ordinary = task();
4027 let next = advance_interrupt(
4028 Interrupt::Idle,
4029 std::slice::from_ref(&ordinary.id),
4030 std::slice::from_ref(&ordinary),
4031 );
4032 assert_eq!(next, Interrupt::Idle);
4033 }
4034
4035 #[test]
4036 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
4037 let marked = interrupt_task("marked");
4040 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4041 assert_eq!(next, Interrupt::Idle);
4042 }
4043
4044 #[test]
4045 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
4046 let marked = interrupt_task("marked");
4047 let next = advance_interrupt(
4048 Interrupt::Idle,
4049 &["running".to_owned()],
4050 std::slice::from_ref(&marked),
4051 );
4052 assert_eq!(
4053 next,
4054 Interrupt::Parking {
4055 parked: vec!["running".to_owned()],
4056 interrupt_task: "marked".to_owned(),
4057 }
4058 );
4059 }
4060
4061 #[test]
4070 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
4071 let marked = interrupt_task("marked");
4072
4073 let two = advance_interrupt(
4074 Interrupt::Idle,
4075 &["a".to_owned(), "b".to_owned()],
4076 std::slice::from_ref(&marked),
4077 );
4078 assert_eq!(two, Interrupt::Idle);
4079
4080 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4081 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
4082 }
4083
4084 #[test]
4085 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
4086 let state = Interrupt::Parking {
4087 parked: vec!["running".to_owned()],
4088 interrupt_task: "marked".to_owned(),
4089 };
4090 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
4092 assert_eq!(still_going, state);
4093
4094 let stopped_but_not_yet_dispatched =
4098 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
4099 assert_eq!(stopped_but_not_yet_dispatched, state);
4100
4101 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
4103 assert_eq!(
4104 dispatched,
4105 Interrupt::Running {
4106 parked: vec!["running".to_owned()],
4107 interrupt_task: "marked".to_owned(),
4108 }
4109 );
4110 }
4111
4112 #[test]
4113 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
4114 let state = Interrupt::Running {
4115 parked: vec!["running".to_owned()],
4116 interrupt_task: "marked".to_owned(),
4117 };
4118 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
4119 assert_eq!(still_running, state);
4120
4121 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
4128 assert_eq!(
4129 ended,
4130 Interrupt::Resuming {
4131 parked: vec!["running".to_owned()]
4132 }
4133 );
4134 }
4135
4136 #[test]
4137 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
4138 let state = Interrupt::Resuming {
4139 parked: vec!["running".to_owned()],
4140 };
4141 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
4142 assert_eq!(still_waiting, state);
4143
4144 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
4145 assert_eq!(dispatched, Interrupt::Idle);
4146 }
4147
4148 #[test]
4154 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
4155 {
4156 let state = Interrupt::Parking {
4157 parked: vec!["running".to_owned()],
4158 interrupt_task: "marked".to_owned(),
4159 };
4160 let next = advance_interrupt(state, &[], &[]);
4163 assert_eq!(
4164 next,
4165 Interrupt::Resuming {
4166 parked: vec!["running".to_owned()]
4167 },
4168 "abandoning the interrupt must not abandon the resume it owes"
4169 );
4170 }
4171
4172 #[test]
4175 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
4176 let state = Interrupt::Resuming {
4177 parked: vec!["running".to_owned()],
4178 };
4179 let next = advance_interrupt(state, &[], &[]);
4180 assert_eq!(
4181 next,
4182 Interrupt::Idle,
4183 "nothing is left to wait for; the loop must not stay wedged"
4184 );
4185 }
4186
4187 #[test]
4188 fn disabled_by_config_the_sequence_can_never_leave_idle() {
4189 let marked = interrupt_task("marked");
4190 let next = advance_interrupt_tick(
4191 false,
4192 Interrupt::Idle,
4193 &["running".to_owned()],
4194 std::slice::from_ref(&marked),
4195 );
4196 assert_eq!(
4197 next,
4198 Interrupt::Idle,
4199 "an unmarked, unconfigured daemon must behave exactly as before"
4200 );
4201 }
4202
4203 #[test]
4204 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4205 let state = Interrupt::Parking {
4206 parked: vec!["running".to_owned()],
4207 interrupt_task: "marked".to_owned(),
4208 };
4209 let candidates = vec![interrupt_task("marked"), task()];
4210 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4211 assert!(
4212 allowed.is_empty(),
4213 "nothing may dispatch - not even the interrupt task itself - \
4214 until the parked run has actually stopped"
4215 );
4216 }
4217
4218 #[test]
4232 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4233 for state in [
4234 Interrupt::Parking {
4235 parked: vec!["running".to_owned()],
4236 interrupt_task: "marked".to_owned(),
4237 },
4238 Interrupt::Running {
4239 parked: vec!["running".to_owned()],
4240 interrupt_task: "marked".to_owned(),
4241 },
4242 Interrupt::Resuming {
4243 parked: vec!["running".to_owned()],
4244 },
4245 ] {
4246 let candidates = vec![
4247 interrupt_task("marked"),
4248 urgent_task("hot"),
4249 task_with_id("ordinary"),
4250 ];
4251 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4252 assert!(
4253 !allowed.iter().any(|t| t.id == "hot"),
4254 "an urgent candidate must wait out the same gate as anything \
4255 else while the run it would run alongside has not actually \
4256 left flight, for state {state:?}: {allowed:?}"
4257 );
4258 }
4259 }
4260
4261 #[test]
4267 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4268 let state = Interrupt::Resuming {
4269 parked: vec!["hot".to_owned()],
4270 };
4271 let candidates = vec![urgent_task("hot"), task()];
4272 let allowed = interrupt_gate(&state, &[], candidates);
4273 assert_eq!(
4274 allowed.iter().filter(|t| t.id == "hot").count(),
4275 1,
4276 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4277 );
4278 }
4279
4280 #[test]
4281 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4282 let state = Interrupt::Parking {
4283 parked: vec!["running".to_owned()],
4284 interrupt_task: "marked".to_owned(),
4285 };
4286 let other = task();
4287 let candidates = vec![interrupt_task("marked"), other.clone()];
4288 let allowed = interrupt_gate(&state, &[], candidates);
4289 assert_eq!(allowed.len(), 1);
4290 assert_eq!(allowed[0].id, "marked");
4291 }
4292
4293 #[test]
4294 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4295 let state = Interrupt::Running {
4296 parked: vec!["running".to_owned()],
4297 interrupt_task: "marked".to_owned(),
4298 };
4299 let candidates = vec![task(), task()];
4300 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4301 assert!(allowed.is_empty());
4302 }
4303
4304 #[test]
4311 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4312 let state = Interrupt::Resuming {
4313 parked: vec!["a".to_owned(), "c".to_owned()],
4314 };
4315 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4316 let allowed = interrupt_gate(&state, &[], candidates);
4317 assert_eq!(
4318 allowed.len(),
4319 1,
4320 "at most one candidate may be offered while resuming: {allowed:?}"
4321 );
4322 assert_eq!(allowed[0].id, "a");
4323 }
4324
4325 #[test]
4326 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4327 let state = Interrupt::Resuming {
4328 parked: vec!["a".to_owned()],
4329 };
4330 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4331 assert!(allowed.is_empty());
4332 }
4333
4334 #[test]
4340 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4341 let running = task(); let marked = interrupt_task("marked");
4343
4344 let mut state = Interrupt::Idle;
4345 let in_flight = vec![running.id.clone()];
4347 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4348 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4349 assert!(gated.is_empty(), "still waiting on `running` to park");
4350
4351 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4353 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4354 assert_eq!(
4355 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4356 vec!["marked"],
4357 "only the interrupt task may be offered to the dispatcher now"
4358 );
4359
4360 state = advance_interrupt_tick(
4362 true,
4363 state,
4364 &["marked".to_owned()],
4365 std::slice::from_ref(&running),
4366 );
4367 let gated = interrupt_gate(
4368 &state,
4369 &["marked".to_owned()],
4370 vec![marked.clone(), running.clone()],
4371 );
4372 assert!(
4373 gated.is_empty(),
4374 "the parked run must not be offered back while the interrupt \
4375 task is still running"
4376 );
4377
4378 let other = task_with_id("other");
4382 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4383 assert_eq!(
4384 state,
4385 Interrupt::Resuming {
4386 parked: vec![running.id.clone()]
4387 }
4388 );
4389 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4390 assert_eq!(
4391 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4392 vec![running.id.as_str()],
4393 "exactly the parked run resumes - not the unrelated task, even \
4394 though it was offered first"
4395 );
4396
4397 state = advance_interrupt_tick(
4401 true,
4402 state,
4403 std::slice::from_ref(&running.id),
4404 std::slice::from_ref(&other),
4405 );
4406 assert_eq!(state, Interrupt::Idle);
4407 let gated = interrupt_gate(
4408 &state,
4409 std::slice::from_ref(&running.id),
4410 vec![other.clone()],
4411 );
4412 assert_eq!(
4413 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4414 vec![other.id.as_str()],
4415 "ordinary dispatch is unrestricted again"
4416 );
4417 }
4418
4419 #[test]
4420 fn every_run_status_settles_the_task_it_came_from() {
4421 let table = [
4423 (RunStatus::Merged, TaskStatus::Done, 1),
4424 (RunStatus::Ready, TaskStatus::Done, 1),
4425 (RunStatus::Stalled, TaskStatus::Failed, 0),
4426 (RunStatus::Blocked, TaskStatus::Failed, 1),
4427 (RunStatus::Failed, TaskStatus::Failed, 1),
4428 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4429 (RunStatus::Prep, TaskStatus::Failed, 1),
4430 (RunStatus::Implementing, TaskStatus::Failed, 1),
4431 (RunStatus::Judging, TaskStatus::Failed, 1),
4432 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4433 (RunStatus::Voting, TaskStatus::Failed, 1),
4434 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4435 (RunStatus::Gating, TaskStatus::Failed, 1),
4436 ];
4437 for (run, want, attempts) in table {
4438 let mut t = task();
4439 t.start("20260902-000000-aaaa".to_owned());
4440 settle(
4441 &mut t,
4442 Verdict {
4443 status: run,
4444 left_pr: false,
4445 parked: false,
4446 quota_hit: matches!(run, RunStatus::Stalled),
4447 no_viable_candidates: false,
4448 },
4449 "why",
4450 2,
4451 );
4452 assert_eq!(t.status, want, "task status after {}", label(run));
4453 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4454 }
4455 }
4456
4457 #[test]
4458 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4459 let mut stalled = task();
4460 stalled.start("20260902-000000-aaaa".to_owned());
4461 settle(
4462 &mut stalled,
4463 Verdict {
4464 status: RunStatus::Stalled,
4465 left_pr: false,
4466 parked: false,
4467 quota_hit: true,
4468 no_viable_candidates: false,
4469 },
4470 "quota",
4471 1,
4472 );
4473 assert_eq!(stalled.attempts, 0);
4474 assert!(
4475 stalled.status.runnable(),
4476 "a machine problem must leave the task in line"
4477 );
4478
4479 let mut blocked = task();
4480 blocked.start("20260902-000000-aaaa".to_owned());
4481 settle(
4482 &mut blocked,
4483 Verdict {
4484 status: RunStatus::Blocked,
4485 left_pr: false,
4486 parked: false,
4487 quota_hit: false,
4488 no_viable_candidates: false,
4489 },
4490 "findings open",
4491 1,
4492 );
4493 assert_eq!(blocked.attempts, 1);
4494 assert_eq!(
4495 blocked.status,
4496 TaskStatus::Held,
4497 "the last attempt hands the task to a human"
4498 );
4499 }
4500
4501 #[test]
4502 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4503 let mut delivered = task();
4506 delivered.start("20260903-080619-01c2".to_owned());
4507 settle(
4508 &mut delivered,
4509 Verdict {
4510 status: RunStatus::Blocked,
4511 left_pr: true,
4512 parked: false,
4513 quota_hit: false,
4514 no_viable_candidates: false,
4515 },
4516 "no check status",
4517 4,
4518 );
4519 assert_eq!(
4520 delivered.status,
4521 TaskStatus::Held,
4522 "a pull request waiting on CI or a person is not a retryable failure"
4523 );
4524 assert!(
4525 !delivered.status.runnable(),
4526 "the loop must not pick this task up again"
4527 );
4528 assert_eq!(
4529 delivered.last_error.as_deref(),
4530 Some("no check status"),
4531 "the operator needs to be told what the gate was waiting for"
4532 );
4533
4534 let mut empty_handed = task();
4537 empty_handed.start("20260903-080619-01c2".to_owned());
4538 settle(
4539 &mut empty_handed,
4540 Verdict {
4541 status: RunStatus::Blocked,
4542 left_pr: false,
4543 parked: false,
4544 quota_hit: false,
4545 no_viable_candidates: false,
4546 },
4547 "findings open",
4548 4,
4549 );
4550 assert_eq!(empty_handed.status, TaskStatus::Failed);
4551 assert!(empty_handed.status.runnable());
4552 }
4553
4554 #[test]
4555 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4556 let mut noop = task();
4562 noop.start("20260912-131304-391f".to_owned());
4563 settle(
4564 &mut noop,
4565 Verdict {
4566 status: RunStatus::VerifiedNoop,
4567 left_pr: false,
4568 parked: false,
4569 quota_hit: false,
4570 no_viable_candidates: true,
4571 },
4572 "candidate A: already fixed by b32cfc4, on main",
4573 4,
4574 );
4575 assert_eq!(
4576 noop.status,
4577 TaskStatus::Held,
4578 "an unverified claim is a request for a human, not a failure"
4579 );
4580 assert!(
4581 !noop.status.runnable(),
4582 "the loop must not requeue this on the same unverified claim"
4583 );
4584 assert_eq!(noop.attempts, 1);
4589 }
4590
4591 #[test]
4592 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4593 let mut parked = task();
4598 parked.start("20260903-183634-2d98".to_owned());
4599 settle(
4600 &mut parked,
4601 Verdict {
4602 status: RunStatus::Implementing,
4603 left_pr: false,
4604 quota_hit: false,
4605 parked: true,
4606 no_viable_candidates: false,
4607 },
4608 "parked after `implementing`",
4609 2,
4610 );
4611 assert_eq!(parked.attempts, 0, "a park is refunded");
4612 assert!(
4613 parked.status.runnable(),
4614 "and the task stays in line so the next loop resumes its run"
4615 );
4616 assert_eq!(
4617 parked.last_error.as_deref(),
4618 Some("parked after `implementing`"),
4619 "the card says where it stopped"
4620 );
4621
4622 let mut broken = task();
4626 broken.start("20260903-183634-2d98".to_owned());
4627 settle(
4628 &mut broken,
4629 Verdict {
4630 status: RunStatus::Implementing,
4631 left_pr: false,
4632 quota_hit: false,
4633 parked: false,
4634 no_viable_candidates: false,
4635 },
4636 "returned mid-flight",
4637 2,
4638 );
4639 assert_eq!(broken.attempts, 1);
4640 }
4641
4642 #[test]
4643 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4644 let mut flaky = task();
4649 flaky.start("20260903-123023-e633".to_owned());
4650 settle(
4651 &mut flaky,
4652 Verdict {
4653 status: RunStatus::Stalled,
4654 left_pr: false,
4655 parked: false,
4656 quota_hit: false,
4657 no_viable_candidates: false,
4658 },
4659 "verdict rests on 1 of 3 judges",
4660 2,
4661 );
4662 assert_eq!(
4663 flaky.attempts, 1,
4664 "flakiness spends an attempt, so `max_attempts` still bounds it"
4665 );
4666 assert!(flaky.status.runnable(), "and it is still worth retrying");
4667
4668 let mut limited = task();
4670 limited.start("20260903-123023-e633".to_owned());
4671 settle(
4672 &mut limited,
4673 Verdict {
4674 status: RunStatus::Stalled,
4675 left_pr: false,
4676 parked: false,
4677 quota_hit: true,
4678 no_viable_candidates: false,
4679 },
4680 "judge-2, judge-3 out of quota",
4681 2,
4682 );
4683 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4684 assert!(limited.status.runnable());
4685
4686 let mut worn = task();
4689 for _ in 0..2 {
4690 worn.release();
4691 }
4692 worn.start("20260903-123023-e633".to_owned());
4693 worn.attempts = 2;
4694 settle(
4695 &mut worn,
4696 Verdict {
4697 status: RunStatus::Stalled,
4698 left_pr: false,
4699 parked: false,
4700 quota_hit: false,
4701 no_viable_candidates: false,
4702 },
4703 "no quorum again",
4704 2,
4705 );
4706 assert_eq!(worn.status, TaskStatus::Held);
4707 assert!(!worn.status.runnable());
4708 }
4709
4710 #[test]
4711 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4712 let mut wiped_out = task();
4719 wiped_out.start("20260907-025000-a1b2".to_owned());
4720 settle(
4721 &mut wiped_out,
4722 Verdict {
4723 status: RunStatus::Failed,
4724 left_pr: false,
4725 parked: false,
4726 quota_hit: true,
4727 no_viable_candidates: true,
4728 },
4729 "no candidate produced a change; nothing to judge",
4730 2,
4731 );
4732 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
4733 assert!(
4734 wiped_out.status.runnable(),
4735 "a machine problem must leave the task in line"
4736 );
4737
4738 let mut partial_progress = task();
4744 partial_progress.start("20260907-025500-c3d4".to_owned());
4745 settle(
4746 &mut partial_progress,
4747 Verdict {
4748 status: RunStatus::Failed,
4749 left_pr: false,
4750 parked: false,
4751 quota_hit: true,
4752 no_viable_candidates: false,
4753 },
4754 "gate failed on the winning candidate",
4755 2,
4756 );
4757 assert_eq!(
4758 partial_progress.attempts, 1,
4759 "a candidate that actually produced a change spends the attempt \
4760 even though some other seat hit its quota"
4761 );
4762 assert!(partial_progress.status.runnable());
4763 }
4764
4765 #[test]
4766 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
4767 let mut t = task();
4773 t.start("20260907-025000-a1b2".to_owned());
4774 let mut state = run_state(RunStatus::Failed);
4775 state.quota.push(QuotaLoss {
4776 seat: "cand-a".to_owned(),
4777 node: "implement".to_owned(),
4778 at: Timestamp::now(),
4779 reset: None,
4780 });
4781 assert!(
4782 state.viable().is_empty(),
4783 "no candidate was added, so nothing is viable"
4784 );
4785 reclaim(&mut t, Some(state), 2, "en");
4786 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
4787 assert!(t.status.runnable());
4788 }
4789
4790 #[test]
4791 fn a_held_task_is_never_offered_to_the_loop() {
4792 let dir = tempfile::tempdir().unwrap();
4793 let queue = Queue::at(dir.path().to_path_buf());
4794 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
4795 let mut t = task();
4796 t.id = format!("2026090{n}-000000-000{n}");
4797 t.priority = priority;
4798 queue.put(&mut t).unwrap();
4799 }
4800 let mut held = task();
4801 held.id = "20260909-000000-9999".to_owned();
4802 held.priority = 99;
4803 held.hold_machine(None);
4804 queue.put(&mut held).unwrap();
4805
4806 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
4807 assert_eq!(order.len(), 3);
4808 assert!(!order.contains(&held.id));
4809 assert_eq!(
4810 order.first().cloned(),
4811 queue.next_runnable().map(|t| t.id),
4812 "the loop's first candidate is exactly what the queue offers"
4813 );
4814 assert_eq!(
4815 order,
4816 vec![
4817 "20260902-000000-0002".to_owned(),
4818 "20260903-000000-0003".to_owned(),
4819 "20260901-000000-0001".to_owned(),
4820 ],
4821 "priority first, then oldest, so nothing starves"
4822 );
4823 }
4824
4825 #[test]
4826 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
4827 let dir = tempfile::tempdir().unwrap();
4828 let queue = Queue::at(dir.path().to_path_buf());
4829 let mut old = task();
4830 old.id = "20260101-000000-old0".to_owned();
4831 queue.put(&mut old).unwrap();
4832 let mut fresh = task();
4833 fresh.id = "20260101-000000-new0".to_owned();
4834 queue.put(&mut fresh).unwrap();
4835
4836 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
4840 std::thread::sleep(Duration::from_millis(60));
4841 let live = queue.claim(&fresh.id).unwrap();
4842
4843 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4844 assert_eq!(swept, vec![old.id.clone()]);
4845 assert!(
4846 queue.claim(&old.id).is_ok(),
4847 "an unparseable lock older than the threshold is swept"
4848 );
4849 assert!(
4850 queue.claim(&fresh.id).is_err(),
4851 "a live pid protects its lock regardless of age"
4852 );
4853 drop(live);
4854 }
4855
4856 #[test]
4857 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
4858 let dir = tempfile::tempdir().unwrap();
4868 let queue = Queue::at(dir.path().to_path_buf());
4869 let mut t = task();
4870 t.id = "20260101-000000-live".to_owned();
4871 queue.put(&mut t).unwrap();
4872
4873 let claim = queue.claim(&t.id).unwrap();
4874 std::thread::sleep(Duration::from_millis(60));
4875
4876 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4877 assert!(
4878 swept.is_empty(),
4879 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
4880 );
4881 assert!(
4882 queue.claim(&t.id).is_err(),
4883 "the lock still protects its task"
4884 );
4885 drop(claim);
4886 }
4887
4888 fn injected_dead_pid() -> u32 {
4891 std::process::id().checked_add(1).unwrap_or(1)
4892 }
4893
4894 #[test]
4895 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
4896 let dir = tempfile::tempdir().unwrap();
4897 let queue = Queue::at(dir.path().to_path_buf());
4898 let mut t = task();
4899 t.id = "20260101-000000-dead".to_owned();
4900 queue.put(&mut t).unwrap();
4901 let dead_pid = injected_dead_pid();
4902
4903 std::fs::write(
4908 dir.path().join(format!("{}.lock", t.id)),
4909 dead_pid.to_string(),
4910 )
4911 .unwrap();
4912
4913 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4914 pid != dead_pid
4915 });
4916 assert_eq!(
4917 swept,
4918 vec![t.id.clone()],
4919 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
4920 );
4921 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
4922 }
4923
4924 #[test]
4925 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
4926 let dir = tempfile::tempdir().unwrap();
4927 let queue = Queue::at(dir.path().to_path_buf());
4928 let mut t = task();
4929 t.id = "20260101-000000-late".to_owned();
4930 queue.put(&mut t).unwrap();
4931 let dead_pid = injected_dead_pid();
4932
4933 assert!(
4936 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
4937 "nothing has claimed the task yet"
4938 );
4939
4940 std::fs::write(
4943 dir.path().join(format!("{}.lock", t.id)),
4944 dead_pid.to_string(),
4945 )
4946 .unwrap();
4947
4948 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4952 pid != dead_pid
4953 });
4954 assert_eq!(swept, vec![t.id.clone()]);
4955 }
4956
4957 #[test]
4958 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
4959 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
4964 let dir = tempfile::tempdir().unwrap();
4965 let queue = Queue::at(dir.path().to_path_buf());
4966 let mut t = task();
4967 t.id = "20260101-000000-crsh".to_owned();
4968 t.status = TaskStatus::Running;
4969 t.attempts = 1;
4970 t.runs.push("20260904-000000-4043".to_owned());
4974 queue.put(&mut t).unwrap();
4975 let dead_pid = injected_dead_pid();
4976
4977 std::fs::write(
4980 dir.path().join(format!("{}.lock", t.id)),
4981 dead_pid.to_string(),
4982 )
4983 .unwrap();
4984
4985 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
4991 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
4992
4993 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4994 pid != dead_pid
4995 });
4996 assert_eq!(swept, vec![t.id.clone()]);
4997
4998 let reclaimed = reclaim_orphaned_running(&queue, 2);
4999 assert_eq!(reclaimed, vec![t.id.clone()]);
5000 let after = queue.get(&t.id).unwrap();
5001 assert_eq!(
5002 after.status,
5003 TaskStatus::Held,
5004 "no run.json to recover from, so a human is asked"
5005 );
5006 assert_eq!(
5007 after.runs,
5008 vec!["20260904-000000-4043".to_owned()],
5009 "the crashed run's id is kept as evidence, not discarded"
5010 );
5011 }
5012
5013 #[test]
5014 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
5015 let dir = tempfile::tempdir().unwrap();
5016 let queue = Queue::at(dir.path().to_path_buf());
5017 let mut t = task();
5018 t.id = "20260101-000000-unknown".to_owned();
5019 queue.put(&mut t).unwrap();
5020 let dead_pid = injected_dead_pid();
5021 std::fs::write(
5022 dir.path().join(format!("{}.lock", t.id)),
5023 dead_pid.to_string(),
5024 )
5025 .unwrap();
5026
5027 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
5028 assert!(swept.is_empty(), "an unknown pid must keep its lock");
5029 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
5030 }
5031
5032 fn run_state_in(status: RunStatus, language: &str) -> RunState {
5033 let mut s = run_state(status);
5034 s.config.graph.language = language.to_owned();
5035 s
5036 }
5037
5038 fn unstarted_verdict(status: RunStatus) -> Verdict {
5039 Verdict {
5040 status,
5041 left_pr: false,
5042 quota_hit: false,
5043 parked: false,
5044 no_viable_candidates: false,
5045 }
5046 }
5047
5048 #[test]
5049 fn an_already_in_base_run_finishes_the_task_without_spending_an_attempt() {
5050 let mut t = task();
5051 t.attempts = 1;
5052 settle_in(
5053 &mut t,
5054 unstarted_verdict(RunStatus::AlreadyInBase),
5055 "already in main",
5056 1,
5057 phrases("en"),
5058 );
5059 assert_eq!(t.status, TaskStatus::Done);
5060 assert_eq!(t.attempts, 0);
5061 }
5062
5063 #[test]
5064 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
5065 let reason = |language: &str| {
5066 let mut t = task();
5067 settle_in(
5068 &mut t,
5069 unstarted_verdict(RunStatus::Judging),
5070 "boom",
5071 1,
5072 phrases(language),
5073 );
5074 t.last_error.or(t.hold_reason).unwrap_or_default()
5075 };
5076 assert!(
5077 reason("en").starts_with("the graph stopped at `"),
5078 "{}",
5079 reason("en")
5080 );
5081 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
5082 assert!(reason("日本語").contains("boom"));
5083 assert_eq!(reason("fr"), reason("en"));
5084 }
5085
5086 #[test]
5087 fn describe_follows_the_run_language_and_keeps_the_run_id() {
5088 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
5089 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
5090 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
5091 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
5092 assert!(ja.contains("[run "), "{ja}");
5093 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
5094 assert!(ended.starts_with("run 終了: "), "{ended}");
5095 let mut de = run_state_in(RunStatus::Failed, "de");
5096 let mut en = run_state_in(RunStatus::Failed, "en");
5097 de.id = "same".to_owned();
5098 en.id = "same".to_owned();
5099 assert_eq!(describe(&de), describe(&en));
5100 }
5101
5102 #[test]
5103 fn handover_refusals_follow_the_language_and_keep_the_detail_apart() {
5104 use crate::handover::Refused;
5105 let cases = [
5106 Refused::Foreign {
5107 branch: "b".into(),
5108 path: "/w/x".into(),
5109 why: "made by hand".into(),
5110 },
5111 Refused::Unsafe {
5112 branch: "b".into(),
5113 path: "/w/x".into(),
5114 why: "its worktree has uncommitted changes (a.rs)".into(),
5115 },
5116 Refused::ReleaseFailed {
5117 branch: "b".into(),
5118 path: "/w/x".into(),
5119 run: "ab12".into(),
5120 },
5121 ];
5122 for r in &cases {
5123 let en = (phrases("en").handover_refused)(r);
5124 assert_eq!(en, r.to_string());
5125 let ja = (phrases("ja").handover_refused)(r);
5126 assert!(
5127 !ja.contains("is checked out") && !ja.contains("try again"),
5128 "{ja}"
5129 );
5130 assert!(ja.contains("`b`") && ja.contains("/w/x"), "{ja}");
5131 if let Refused::Foreign { why, .. } | Refused::Unsafe { why, .. } = r {
5132 assert!(ja.contains(&format!("(詳細: {why})")), "{ja}");
5133 }
5134 }
5135 assert!(phrases("en").handover_hint.contains("release the task"));
5136 assert!(phrases("ja").handover_hint.contains("解放"));
5137 }
5138
5139 #[test]
5140 fn describe_translates_the_unanswered_reviewer_stop_only_in_ja() {
5141 let build = |lang: &str| {
5142 let mut s = run_state_in(RunStatus::Blocked, lang);
5143 s.id = "same".to_owned();
5144 s.config.graph.review_rounds = 3;
5145 let mut r = review_round(3);
5146 r.expected = 3;
5147 r.answered = 1;
5148 s.reviews.push(r);
5149 s.event(
5150 "review",
5151 "2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean",
5152 );
5153 s
5154 };
5155 let en = describe(&build("en"));
5156 assert!(en.contains("(review: 2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean)"), "{en}");
5157 let ja = describe(&build("ja"));
5158 assert!(ja.contains("2 席のレビュアーが 3 ラウンド"), "{ja}");
5159 assert!(!ja.contains("never answered"), "{ja}");
5160 assert!(ja.contains("[run "), "{ja}");
5161
5162 let mut failed = build("ja");
5165 failed.reviews[0].e2e.push(crate::run::CommandOutcome {
5166 command: "cargo test".to_owned(),
5167 code: Some(1),
5168 output_tail: String::new(),
5169 duration_ms: 0,
5170 resource_blocked: false,
5171 });
5172 failed.event("review", "stopped; e2e failed: cargo test");
5173 let ja = describe(&failed);
5174 assert!(ja.contains("stopped; e2e failed: cargo test"), "{ja}");
5175 assert!(!ja.contains("席のレビュアー"), "{ja}");
5176 }
5177
5178 #[test]
5179 fn refusals_and_recovery_prose_follow_the_language() {
5180 let t = held_task_with("r1");
5181 let q = action_question(
5182 "r1",
5183 ask::ChoiceAction::Resume {
5184 run: "r1".to_owned(),
5185 },
5186 );
5187 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
5188 ActionDecision::Refuse(s) => s,
5189 other => panic!("{other:?}"),
5190 };
5191 assert!(refuse(phrases("en")).contains("could not be read: gone"));
5192 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
5193 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
5194
5195 let mut held = task();
5196 reclaim(&mut held, None, 2, "ja");
5197 assert!(held.hold_reason.unwrap().contains("保留にしました"));
5198 let mut held = task();
5199 reclaim(&mut held, None, 2, "xx");
5200 assert!(held.hold_reason.unwrap().contains("held for a human"));
5201 }
5202
5203 fn run_state(status: RunStatus) -> RunState {
5204 let mut state = RunState::new(
5205 PathBuf::from("/repo"),
5206 "main".to_owned(),
5207 "abc1234def".to_owned(),
5208 "add retries".to_owned(),
5209 Config::default(),
5210 );
5211 state.status = status;
5212 state
5213 }
5214
5215 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
5216 Candidate {
5217 index: 0,
5218 label,
5219 agent: "claude".to_owned(),
5220 branch: format!("magi/x/{label}"),
5221 worktree: PathBuf::from("/repo"),
5222 summary: summary.to_owned(),
5223 stat: String::new(),
5224 files: 0,
5225 commits: usize::from(!empty),
5226 empty,
5227 failed: failed.map(str::to_owned),
5228 verified_noop: None,
5229 duration_ms: 0,
5230 folded: false,
5231 }
5232 }
5233
5234 #[test]
5235 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
5236 let mut state = run_state(RunStatus::Blocked);
5237 state.gate = vec![
5238 CommandOutcome {
5239 command: "cargo make check".to_owned(),
5240 code: Some(0),
5241 output_tail: "ok".to_owned(),
5242 duration_ms: 0,
5243 resource_blocked: false,
5244 },
5245 CommandOutcome {
5246 command: "cargo test".to_owned(),
5247 code: Some(101),
5248 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
5249 duration_ms: 0,
5250 resource_blocked: false,
5251 },
5252 ];
5253 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
5254 assert!(d.contains("cargo test"), "{d}");
5255 assert!(
5256 !d.contains("cargo make check"),
5257 "a passing check is not a diagnostic: {d}"
5258 );
5259 assert!(d.contains("assertion failed"), "{d}");
5260 }
5261
5262 #[test]
5263 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
5264 let mut state = run_state(RunStatus::Blocked);
5265 state.event(
5266 "land",
5267 "stopped: the fixer produced no commit while 2 check(s) were failing \
5268 (build, lint); stopping instead of looping on an unchanged tree",
5269 );
5270 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
5271 assert!(d.contains("build"), "{d}");
5272 assert!(d.contains("lint"), "{d}");
5273 assert!(d.contains("fixer produced no commit"), "{d}");
5274 }
5275
5276 #[test]
5277 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
5278 let state = run_state(RunStatus::VerifiedNoop);
5283 let d = describe(&state);
5284 assert!(
5285 d.contains("agent-verified no-op"),
5286 "expected the display label, not the wire spelling: {d}"
5287 );
5288 assert!(!d.contains("verified_noop"), "{d}");
5289 }
5290
5291 #[test]
5292 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5293 let mut state = run_state(RunStatus::Failed);
5299 state.candidates = vec![candidate(
5300 'A',
5301 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5302 true,
5303 None,
5304 )];
5305 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5306 assert!(d.contains("candidate A"), "{d}");
5307 assert!(d.contains("tagged v1.2.3"), "{d}");
5308 }
5309
5310 #[test]
5311 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5312 let mut state = run_state(RunStatus::Failed);
5313 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5314 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5315 assert!(d.contains("candidate A"), "{d}");
5316 assert!(d.contains("agent timed out"), "{d}");
5317 }
5318
5319 #[test]
5320 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5321 let mut state = run_state(RunStatus::Failed);
5324 state.candidates = vec![candidate('A', "did the work", false, None)];
5325 assert!(diagnostic(&state).is_none());
5326 }
5327
5328 #[test]
5329 fn diagnostic_is_bounded_however_much_a_run_printed() {
5330 let mut state = run_state(RunStatus::Blocked);
5331 state.gate = vec![
5332 CommandOutcome {
5333 command: "cargo test".to_owned(),
5334 code: Some(101),
5335 output_tail: "x".repeat(50_000),
5336 duration_ms: 0,
5337 resource_blocked: false,
5338 },
5339 CommandOutcome {
5340 command: "cargo clippy".to_owned(),
5341 code: Some(1),
5342 output_tail: "y".repeat(50_000),
5343 duration_ms: 0,
5344 resource_blocked: false,
5345 },
5346 ];
5347 state.candidates = vec![
5348 candidate('A', &"z".repeat(50_000), true, None),
5349 candidate('B', &"w".repeat(50_000), true, None),
5350 ];
5351 let d = diagnostic(&state).expect("plenty here to diagnose");
5352 assert!(
5353 d.len() <= DIAGNOSTIC_MAX,
5354 "diagnostic grew to {} bytes, unbounded",
5355 d.len()
5356 );
5357 }
5358
5359 #[test]
5360 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5361 let mut state = run_state(RunStatus::Blocked);
5362 state.gate = vec![CommandOutcome {
5363 command: "cargo test".to_owned(),
5364 code: Some(101),
5365 output_tail: "assertion failed".to_owned(),
5366 duration_ms: 0,
5367 resource_blocked: false,
5368 }];
5369 let verdict = Verdict {
5370 status: RunStatus::Blocked,
5371 left_pr: false,
5372 quota_hit: false,
5373 parked: false,
5374 no_viable_candidates: false,
5375 };
5376
5377 let mut t = task();
5380 t.start("run-1".to_owned());
5381 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5382 assert_eq!(t.status, TaskStatus::Failed);
5383 assert!(t.diagnostic.is_none());
5384
5385 t.start("run-2".to_owned());
5388 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5389 assert_eq!(t.status, TaskStatus::Held);
5390 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5391 assert!(d.contains("cargo test"), "{d}");
5392 }
5393
5394 #[test]
5395 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5396 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5401 let home = crate::run::home();
5402 let state = run_state(RunStatus::VerifiedNoop);
5403 let mut q = ask::Question::new(
5404 state.id.clone(),
5405 "implement".to_owned(),
5406 "impl-A".to_owned(),
5407 "is this really a no-op?".to_owned(),
5408 String::new(),
5409 Vec::new(),
5410 );
5411 Questions::at(home.join("questions")).put(&mut q).unwrap();
5412
5413 let verdict = Verdict {
5414 status: RunStatus::VerifiedNoop,
5415 left_pr: false,
5416 quota_hit: false,
5417 parked: false,
5418 no_viable_candidates: false,
5419 };
5420 let mut t = task();
5421 t.start(state.id.clone());
5422 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5423
5424 assert_eq!(t.status, TaskStatus::Held);
5425 let reason = t.hold_reason.expect("a held task must record why");
5426 assert!(
5427 reason.starts_with("run ended agent-verified no-op"),
5428 "the original settle reason must survive unchanged: {reason}"
5429 );
5430 assert!(
5431 reason.contains(q.short()),
5432 "the open question's id must be named so the notice is actionable: {reason}"
5433 );
5434 }
5435
5436 #[test]
5437 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5438 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5439 let state = run_state(RunStatus::VerifiedNoop);
5440
5441 let verdict = Verdict {
5442 status: RunStatus::VerifiedNoop,
5443 left_pr: false,
5444 quota_hit: false,
5445 parked: false,
5446 no_viable_candidates: false,
5447 };
5448 let mut t = task();
5449 t.start(state.id.clone());
5450 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5451
5452 assert_eq!(t.status, TaskStatus::Held);
5453 assert_eq!(
5454 t.hold_reason.as_deref(),
5455 Some("run ended agent-verified no-op"),
5456 "nothing to append when the question was already answered or never asked"
5457 );
5458 }
5459
5460 #[test]
5461 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5462 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5463 let mut first = run_state(RunStatus::Blocked);
5464 first.id = "20260101-000000-sup1".to_owned();
5465 first.save().unwrap();
5466 let mut second = run_state(RunStatus::Merged);
5467 second.id = "20260101-000000-sup2".to_owned();
5468 second.save().unwrap();
5469
5470 let mut t = task();
5471 t.runs = vec![first.id.clone(), second.id.clone()];
5472 t.status = TaskStatus::Done;
5473
5474 supersede_prior_runs(&t, &crate::run::home());
5475
5476 assert_eq!(
5477 RunState::load(&first.id).unwrap().status,
5478 RunStatus::Superseded,
5479 "the first attempt's Blocked no longer needs anyone's attention"
5480 );
5481 assert_eq!(
5482 RunState::load(&second.id).unwrap().status,
5483 RunStatus::Merged,
5484 "the run that actually succeeded is left exactly as it was"
5485 );
5486 }
5487
5488 #[test]
5489 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5490 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5496 let mut first = run_state(RunStatus::Blocked);
5497 first.id = "20260101-000000-sup9".to_owned();
5498 first.driver_pid = Some(std::process::id());
5501 first.driver_started_at = Some(
5502 crate::proc::process_started_at(std::process::id())
5503 .expect("this test process's own start time must be queryable"),
5504 );
5505 first.save().unwrap();
5506 let mut second = run_state(RunStatus::Merged);
5507 second.id = "20260101-000000-supa".to_owned();
5508 second.save().unwrap();
5509
5510 let mut t = task();
5511 t.runs = vec![first.id.clone(), second.id.clone()];
5512 t.status = TaskStatus::Done;
5513
5514 supersede_prior_runs(&t, &crate::run::home());
5515
5516 assert_eq!(
5517 RunState::load(&first.id).unwrap().status,
5518 RunStatus::Blocked,
5519 "a live driver_pid means something is still actually working this run, \
5520 even though no daemon claims it - rewriting under it would just be \
5521 undone the next time that process saves"
5522 );
5523 }
5524
5525 #[test]
5526 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5527 let dir = tempfile::tempdir().unwrap();
5533 let home = dir.path().to_path_buf();
5534 let queue = Queue::at(dir.path().join("queue"));
5535
5536 let mut first = run_state(RunStatus::Blocked);
5537 first.id = "20260101-000000-supd".to_owned();
5538 first.driver_pid = Some(std::process::id());
5539 first.driver_started_at = Some(
5540 crate::proc::process_started_at(std::process::id())
5541 .expect("this test process's own start time must be queryable"),
5542 );
5543 first.save_under(&home).unwrap();
5544 let mut second = run_state(RunStatus::Merged);
5545 second.id = "20260101-000000-supe".to_owned();
5546 second.save_under(&home).unwrap();
5547
5548 let mut t = task();
5549 t.runs = vec![first.id.clone(), second.id.clone()];
5550 t.status = TaskStatus::Done;
5551 queue.put(&mut t).unwrap();
5552
5553 resweep_superseded_attempts(&queue, &home);
5554 assert_eq!(
5555 RunState::load_under(&first.id, &home).unwrap().status,
5556 RunStatus::Blocked,
5557 "still live on the first pass, so still untouched"
5558 );
5559
5560 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5566 stale.driver_started_at = Some("1".to_owned());
5567 stale.save_under(&home).unwrap();
5568
5569 resweep_superseded_attempts(&queue, &home);
5570 assert_eq!(
5571 RunState::load_under(&first.id, &home).unwrap().status,
5572 RunStatus::Superseded,
5573 "the second pass catches up what the first one correctly skipped"
5574 );
5575 }
5576
5577 #[test]
5578 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5579 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5580 let mut first = run_state(RunStatus::Blocked);
5581 first.id = "20260101-000000-sup3".to_owned();
5582 first.save().unwrap();
5583 let mut second = run_state(RunStatus::Blocked);
5584 second.id = "20260101-000000-sup4".to_owned();
5585 second.save().unwrap();
5586
5587 let mut t = task();
5588 t.runs = vec![first.id.clone(), second.id.clone()];
5589 t.status = TaskStatus::Failed;
5593
5594 supersede_prior_runs(&t, &crate::run::home());
5595
5596 assert_eq!(
5597 RunState::load(&first.id).unwrap().status,
5598 RunStatus::Blocked
5599 );
5600 assert_eq!(
5601 RunState::load(&second.id).unwrap().status,
5602 RunStatus::Blocked
5603 );
5604 }
5605
5606 #[test]
5607 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5608 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5612 let mut first = run_state(RunStatus::Blocked);
5613 first.id = "20260101-000000-sup5".to_owned();
5614 first.save().unwrap();
5615
5616 let mut t = task();
5617 t.runs = vec![first.id.clone()];
5618 t.status = TaskStatus::Done;
5619
5620 supersede_prior_runs(&t, &crate::run::home());
5621
5622 assert_eq!(
5623 RunState::load(&first.id).unwrap().status,
5624 RunStatus::Blocked,
5625 "a single-attempt task has no earlier run to supersede"
5626 );
5627 }
5628
5629 #[test]
5630 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5631 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5638 let mut first = run_state(RunStatus::Blocked);
5639 first.id = "20260101-000000-supb".to_owned();
5640 first.save().unwrap();
5641 let mut second = run_state(RunStatus::Failed);
5642 second.id = "20260101-000000-supc".to_owned();
5643 second.save().unwrap();
5644
5645 let mut t = task();
5646 t.runs = vec![first.id.clone(), second.id.clone()];
5647 t.status = TaskStatus::Done;
5648
5649 supersede_prior_runs(&t, &crate::run::home());
5650
5651 assert_eq!(
5652 RunState::load(&first.id).unwrap().status,
5653 RunStatus::Blocked,
5654 "the task's last attempt never landed, so there is nothing here \
5655 actually superseding it"
5656 );
5657 }
5658
5659 #[test]
5660 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5661 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5665 let mut failed = run_state(RunStatus::Failed);
5666 failed.id = "20260101-000000-sup6".to_owned();
5667 failed.save().unwrap();
5668 let mut noop = run_state(RunStatus::VerifiedNoop);
5669 noop.id = "20260101-000000-sup7".to_owned();
5670 noop.save().unwrap();
5671 let mut winner = run_state(RunStatus::Ready);
5672 winner.id = "20260101-000000-sup8".to_owned();
5673 winner.save().unwrap();
5674
5675 let mut t = task();
5676 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5677 t.status = TaskStatus::Done;
5678
5679 supersede_prior_runs(&t, &crate::run::home());
5680
5681 assert_eq!(
5682 RunState::load(&failed.id).unwrap().status,
5683 RunStatus::Failed
5684 );
5685 assert_eq!(
5686 RunState::load(&noop.id).unwrap().status,
5687 RunStatus::VerifiedNoop
5688 );
5689 }
5690
5691 fn approval_question(run: &str) -> ask::Question {
5692 ask::Question::new(
5693 run.to_owned(),
5694 land::APPROVAL_NODE.to_owned(),
5695 "land".to_owned(),
5696 "merge?".to_owned(),
5697 String::new(),
5698 vec!["merge".to_owned(), "hold".to_owned()],
5699 )
5700 }
5701
5702 #[test]
5703 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5704 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5705 let mut state = run_state(RunStatus::Landing);
5706 state.id = "20260101-000000-fre1".to_owned();
5707 state.parked = true;
5708 state.save().unwrap();
5709 ask::Questions::open()
5710 .put(&mut approval_question(&state.id))
5711 .unwrap();
5712
5713 let mut t = task();
5714 t.runs.push(state.id.clone());
5715 assert_eq!(
5716 land_resume_state(&t),
5717 LandResume::StillWaiting,
5718 "nobody has answered and the timeout has not passed"
5719 );
5720 }
5721
5722 #[test]
5723 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5724 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5729 let mut state = run_state(RunStatus::Landing);
5730 state.id = "20260101-000000-exp1".to_owned();
5731 state.parked = true;
5732 state.config.graph.answer_timeout = 60;
5733 state.save().unwrap();
5734
5735 let store = ask::Questions::open();
5736 let mut q = approval_question(&state.id);
5737 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
5738 store.put(&mut q).unwrap();
5739
5740 let mut t = task();
5741 t.runs.push(state.id.clone());
5742 assert_eq!(
5743 land_resume_state(&t),
5744 LandResume::Ready,
5745 "an expired question must not be waited on forever"
5746 );
5747
5748 let after = store.get(&q.id).unwrap();
5749 assert!(
5750 !after.status.open(),
5751 "the question is abandoned, not silently ignored"
5752 );
5753 assert!(
5754 after.resolution().is_none(),
5755 "an abandoned question is not read as a decision"
5756 );
5757 }
5758
5759 #[test]
5760 fn reclaim_settles_a_running_task_against_its_last_run() {
5761 let mut t = task();
5762 t.start("20260904-000000-4043".to_owned());
5763 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
5764 assert_eq!(
5765 t.status,
5766 TaskStatus::Done,
5767 "a run that actually finished must not stay `running` forever"
5768 );
5769 }
5770
5771 #[test]
5772 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
5773 let mut t = task();
5777 t.start("20260904-000000-4043".to_owned());
5778 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
5779 assert_eq!(t.status, TaskStatus::Failed);
5780 assert!(t.status.runnable());
5781 }
5782
5783 #[test]
5784 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
5785 let mut t = task();
5786 t.start("20260904-000000-4043".to_owned());
5787 reclaim(&mut t, None, 2, "en");
5788 assert_eq!(t.status, TaskStatus::Held);
5789 assert!(
5790 t.last_error
5791 .as_deref()
5792 .is_some_and(|e| e.contains("running")),
5793 "the operator needs to know why this task was held"
5794 );
5795 }
5796
5797 #[test]
5798 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
5799 let dir = tempfile::tempdir().unwrap();
5800 let queue = Queue::at(dir.path().to_path_buf());
5801
5802 let mut orphaned = task();
5804 orphaned.id = "20260904-000000-orph".to_owned();
5805 orphaned.status = TaskStatus::Running;
5806 orphaned.attempts = 1;
5807 queue.put(&mut orphaned).unwrap();
5808
5809 let mut alive = task();
5810 alive.id = "20260904-000000-live".to_owned();
5811 alive.status = TaskStatus::Running;
5812 alive.attempts = 1;
5813 queue.put(&mut alive).unwrap();
5814 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
5815
5816 let mut queued = task();
5817 queued.id = "20260904-000000-wait".to_owned();
5818 queue.put(&mut queued).unwrap();
5819
5820 let reclaimed = reclaim_orphaned_running(&queue, 2);
5821 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
5822
5823 assert_eq!(
5824 queue.get(&orphaned.id).unwrap().status,
5825 TaskStatus::Held,
5826 "nothing was driving it and there was no run to recover"
5827 );
5828 assert_eq!(
5829 queue.get(&alive.id).unwrap().status,
5830 TaskStatus::Running,
5831 "a live claim must protect the task it belongs to"
5832 );
5833 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
5834 }
5835
5836 fn read_run_under(home: &Path, id: &str) -> RunState {
5842 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
5843 serde_json::from_str(&body).unwrap()
5844 }
5845
5846 #[test]
5847 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
5848 let dir = tempfile::tempdir().unwrap();
5849 let home = dir.path().to_path_buf();
5850 let now = Timestamp::now();
5851 let overrun_seat = || crate::run::ActiveSeat {
5852 node: "implement".to_owned(),
5853 started_at: now - jiff::SignedDuration::new(21_000, 0),
5854 timeout_secs: 3_600,
5855 attempt: 0,
5856 task: None,
5857 command: None,
5858 index: None,
5859 total: None,
5860 };
5861
5862 let mut dead = run_state(RunStatus::Implementing);
5863 dead.id = "20260101-000000-dead".to_owned();
5864 dead.active.insert("impl-A".to_owned(), overrun_seat());
5865 dead.driver_pid = Some(4242);
5868 dead.save_under(&home).unwrap();
5869
5870 let mut alive = run_state(RunStatus::Implementing);
5873 alive.id = "20260101-000000-aliv".to_owned();
5874 alive.active.insert("impl-A".to_owned(), overrun_seat());
5875 alive.save_under(&home).unwrap();
5876 let mut status = Status::new();
5877 status.current = vec![Current {
5878 task: "20260101-000000-task".to_owned(),
5879 run: alive.id.clone(),
5880 }];
5881 write_status_to(&home.join("daemon.json"), &status).unwrap();
5882
5883 let questions = Questions::at(home.join("questions"));
5887 let mut q = ask::Question::new(
5888 dead.id.clone(),
5889 "implement".to_owned(),
5890 "impl-A".to_owned(),
5891 "Which storage backend?".to_owned(),
5892 String::new(),
5893 vec!["SQLite".to_owned(), "Redis".to_owned()],
5894 );
5895 questions.put(&mut q).unwrap();
5896
5897 let abandoned = reclaim_abandoned_runs_with(
5898 &home,
5899 now,
5900 |pid| if pid == 4242 { Some(false) } else { None },
5901 |_| panic!("a query answering Dead outright needs no identity corroboration"),
5902 );
5903 assert_eq!(abandoned, vec![dead.id.clone()]);
5904
5905 let reloaded = read_run_under(&home, &dead.id);
5906 assert_eq!(reloaded.status, RunStatus::Failed);
5907 assert!(reloaded.active.is_empty());
5908 assert!(
5909 !questions.get(&q.id).unwrap().status.open(),
5910 "the failed run's own open question must be settled in the same pass"
5911 );
5912
5913 let still_alive = read_run_under(&home, &alive.id);
5914 assert_eq!(
5915 still_alive.status,
5916 RunStatus::Implementing,
5917 "a live daemon's claim protects it"
5918 );
5919 assert!(!still_alive.active.is_empty());
5920 }
5921
5922 #[test]
5932 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
5933 let dir = tempfile::tempdir().unwrap();
5934 let home = dir.path().to_path_buf();
5935 let now = Timestamp::now();
5936
5937 let mut manual = run_state(RunStatus::Reviewing);
5938 manual.id = "20260101-000000-manl".to_owned();
5939 manual.active.insert(
5940 "review-1".to_owned(),
5941 crate::run::ActiveSeat {
5942 node: "review".to_owned(),
5943 started_at: now - jiff::SignedDuration::new(21_000, 0),
5944 timeout_secs: 3_600,
5945 attempt: 0,
5946 task: None,
5947 command: None,
5948 index: None,
5949 total: None,
5950 },
5951 );
5952 manual.driver_pid = Some(4242);
5956 manual.driver_started_at = Some("1790000000".to_owned());
5957 manual.save_under(&home).unwrap();
5958
5959 let abandoned = reclaim_abandoned_runs_with(
5960 &home,
5961 now,
5962 |pid| if pid == 4242 { Some(true) } else { None },
5963 |pid| {
5964 if pid == 4242 {
5965 Some("1790000000".to_owned())
5966 } else {
5967 None
5968 }
5969 },
5970 );
5971 assert!(
5972 abandoned.is_empty(),
5973 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
5974 );
5975
5976 let reloaded = read_run_under(&home, &manual.id);
5977 assert_eq!(reloaded.status, RunStatus::Reviewing);
5978 assert!(!reloaded.active.is_empty());
5979 }
5980
5981 #[test]
5982 fn an_already_claimed_task_is_skipped_rather_than_failed() {
5983 let dir = tempfile::tempdir().unwrap();
5984 let queue = Queue::at(dir.path().to_path_buf());
5985 let mut only = task();
5986 queue.put(&mut only).unwrap();
5987
5988 let _elsewhere = queue.claim(&only.id).unwrap();
5989 let candidates = runnable(&queue);
5990 assert_eq!(candidates.len(), 1, "the task is still runnable");
5991 assert!(
5992 queue.claim(&candidates[0].id).is_err(),
5993 "the loop cannot take a claim somebody else holds"
5994 );
5995
5996 let after = queue.get(&only.id).unwrap();
5997 assert_eq!(after.status, TaskStatus::Queued);
5998 assert_eq!(
5999 after.attempts, 0,
6000 "losing the race is not an attempt at the task"
6001 );
6002 assert_eq!(after.last_error, None);
6003 }
6004
6005 #[test]
6006 fn the_status_file_round_trips_and_its_heartbeat_advances() {
6007 let dir = tempfile::tempdir().unwrap();
6008 let path = dir.path().join("daemon.json");
6009
6010 let mut status = Status::new();
6011 status.idle = false;
6012 status.completed = 7;
6013 status.current = vec![Current {
6014 task: "20260902-000000-t111".to_owned(),
6015 run: "20260902-000001-r111".to_owned(),
6016 }];
6017 write_status_to(&path, &status).unwrap();
6018 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6019 assert_eq!(first.schema, SCHEMA);
6020 assert_eq!(first.pid, std::process::id());
6021 assert!(!first.idle);
6022 assert_eq!(first.completed, 7);
6023 assert_eq!(first.current, status.current);
6024 assert!(
6025 !path.with_extension("json.tmp").exists(),
6026 "the temp file is renamed, not left behind"
6027 );
6028
6029 std::thread::sleep(Duration::from_millis(5));
6030 status.updated_at = Timestamp::now();
6031 status.polls = 3;
6032 write_status_to(&path, &status).unwrap();
6033 let second: Status =
6034 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6035 assert!(
6036 second.updated_at > first.updated_at,
6037 "a reader can only detect staleness if the heartbeat moves"
6038 );
6039 assert_eq!(
6040 second.started_at, first.started_at,
6041 "the start time is not a heartbeat"
6042 );
6043 assert_eq!(second.polls, 3);
6044 }
6045
6046 #[test]
6047 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
6048 let dir = tempfile::tempdir().unwrap();
6049
6050 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
6051
6052 let mut status = Status::new();
6053 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
6054 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6055 let stale = read_status(dir.path()).unwrap();
6056 assert!(
6057 !stale.running(Timestamp::now()),
6058 "a minute without a heartbeat is a dead daemon, not a busy one"
6059 );
6060 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
6061
6062 status.updated_at = Timestamp::now();
6063 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6064 let fresh = read_status(dir.path()).unwrap();
6065 assert!(fresh.running(Timestamp::now()));
6066 }
6067
6068 #[test]
6069 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
6070 let dir = tempfile::tempdir().unwrap();
6071 let now = Timestamp::now();
6072 let mine = "20260903-080619-01c2";
6073
6074 assert!(
6075 !is_working_on(dir.path(), mine, now),
6076 "no status file means nobody is working on anything"
6077 );
6078
6079 let mut status = Status::new();
6080 status.current = vec![Current {
6081 task: "20260903-080340-0167".to_owned(),
6082 run: mine.to_owned(),
6083 }];
6084 status.updated_at = now;
6085 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6086 assert!(is_working_on(dir.path(), mine, now));
6087 assert!(
6088 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
6089 "a daemon busy with one run is not working on another"
6090 );
6091
6092 status.updated_at = now - jiff::SignedDuration::from_secs(600);
6095 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6096 assert!(
6097 !is_working_on(dir.path(), mine, now),
6098 "a stale heartbeat is a dead daemon, so its run is a leftover"
6099 );
6100 }
6101
6102 #[test]
6103 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
6104 let dir = tempfile::tempdir().unwrap();
6105 let now = Timestamp::now();
6106
6107 assert!(
6108 !is_working_on_short(dir.path(), "01c2", now),
6109 "no status file means nobody is working on anything"
6110 );
6111
6112 let mut status = Status::new();
6113 status.current = vec![Current {
6114 task: "20260903-080340-0167".to_owned(),
6115 run: "20260903-080619-01c2".to_owned(),
6116 }];
6117 status.updated_at = now;
6118 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6119 assert!(
6120 is_working_on_short(dir.path(), "01c2", now),
6121 "the run's short id is the last block of its full id"
6122 );
6123 assert!(
6124 !is_working_on_short(dir.path(), "3cbf", now),
6125 "a daemon busy with one worktree bay is not working on another"
6126 );
6127 }
6128
6129 #[test]
6130 fn a_newer_status_file_still_yields_a_reading() {
6131 let dir = tempfile::tempdir().unwrap();
6132 std::fs::write(
6135 dir.path().join("daemon.json"),
6136 serde_json::json!({
6137 "schema": 2,
6138 "updated_at": Timestamp::now().to_string(),
6139 "idle": true,
6140 "surprise": { "nested": [1, 2, 3] },
6141 })
6142 .to_string(),
6143 )
6144 .unwrap();
6145
6146 let reading = read_status(dir.path()).expect("a forward-compatible read");
6147 assert!(reading.running(Timestamp::now()));
6148 assert!(reading.idle);
6149 assert!(reading.current.is_empty());
6150 }
6151
6152 #[test]
6153 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
6154 let dir = tempfile::tempdir().unwrap();
6160 std::fs::write(
6161 dir.path().join("daemon.json"),
6162 serde_json::json!({
6163 "schema": 1,
6164 "pid": 4242,
6165 "updated_at": Timestamp::now().to_string(),
6166 "idle": false,
6167 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
6168 "completed": 3,
6169 "polls": 9,
6170 })
6171 .to_string(),
6172 )
6173 .unwrap();
6174
6175 let reading = read_status(dir.path()).expect("an older shape must still parse");
6176 assert!(reading.running(Timestamp::now()));
6177 assert_eq!(
6178 reading.current,
6179 vec![Current {
6180 task: "20260902-140501-aaaa".to_owned(),
6181 run: "20260902-140502-bbbb".to_owned(),
6182 }]
6183 );
6184 }
6185
6186 #[test]
6187 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
6188 let dir = tempfile::tempdir().unwrap();
6189 std::fs::write(
6190 dir.path().join("daemon.json"),
6191 serde_json::json!({
6192 "schema": 1,
6193 "updated_at": Timestamp::now().to_string(),
6194 "idle": true,
6195 "current": null,
6196 })
6197 .to_string(),
6198 )
6199 .unwrap();
6200 let with_null = read_status(dir.path()).expect("null must still parse");
6201 assert!(with_null.current.is_empty());
6202
6203 std::fs::write(
6204 dir.path().join("daemon.json"),
6205 serde_json::json!({
6206 "schema": 1,
6207 "updated_at": Timestamp::now().to_string(),
6208 "idle": true,
6209 })
6210 .to_string(),
6211 )
6212 .unwrap();
6213 let absent = read_status(dir.path()).expect("a missing field must still parse");
6214 assert!(absent.current.is_empty());
6215 }
6216
6217 #[test]
6218 fn a_task_without_a_repository_runs_in_the_daemons_default() {
6219 let fallback = Path::new("/default");
6220 let mut blank = task();
6221 blank.repo = PathBuf::new();
6222 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
6223 let mut dot = task();
6224 dot.repo = PathBuf::from(".");
6225 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
6226 assert_eq!(
6227 repo_for(&task(), fallback),
6228 PathBuf::from("/repo"),
6229 "a task that names a repository keeps it"
6230 );
6231 }
6232
6233 #[test]
6234 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
6235 let mut solo_cfg = Config::default();
6241 solo_cfg.graph.candidates = 3;
6242 let mut solo_task = task();
6243 solo_task.solo = true;
6244 apply_solo(&mut solo_cfg, &solo_task);
6245 assert_eq!(solo_cfg.graph.candidates, 1);
6246
6247 let mut plain_cfg = Config::default();
6248 plain_cfg.graph.candidates = 3;
6249 let plain_task = task();
6250 assert!(!plain_task.solo);
6251 apply_solo(&mut plain_cfg, &plain_task);
6252 assert_eq!(
6253 plain_cfg.graph.candidates, 3,
6254 "a task that did not ask to run alone keeps the config's candidates"
6255 );
6256 }
6257
6258 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
6259 QuotaLoss {
6260 seat: seat.into(),
6261 node: "judge".into(),
6262 at: at.parse().unwrap(),
6263 reset: reset.map(str::to_string),
6264 }
6265 }
6266
6267 #[test]
6268 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
6269 let old: Vec<QuotaLoss> = (1..=4)
6270 .map(|i| {
6271 loss(
6272 &format!("judge-{i}"),
6273 "2026-09-23T05:23:00Z",
6274 Some("2:40pm (Asia/Tokyo)"),
6275 )
6276 })
6277 .collect();
6278 let fresh = losses_this_attempt(&old, &old);
6279 assert!(fresh.is_empty());
6280 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
6281 }
6283
6284 #[test]
6285 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6286 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6287 let now = Timestamp::now();
6288 let mut after = old.clone();
6289 after.push(loss("judge-2", &now.to_string(), None));
6290 let fresh = losses_this_attempt(&old, &after);
6291 assert_eq!(fresh, vec![after[1].clone()]);
6292 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6293 assert_eq!(
6294 until,
6295 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6296 );
6297 }
6298
6299 #[test]
6300 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6301 let before = vec![
6304 loss("judge-1", "2026-09-23T05:23:00Z", None),
6305 loss("judge-2", "2026-09-23T05:24:00Z", None),
6306 ];
6307 let after = vec![
6308 loss("judge-2", "2026-09-23T05:24:00Z", None),
6309 loss("judge-1", "2026-09-24T01:00:00Z", None),
6310 ];
6311 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6312 }
6313
6314 #[test]
6315 fn merge_overrides_are_parsed_or_refused() {
6316 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6317 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6318 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6319 assert!(merge_mode("squash").is_err());
6320 }
6321
6322 #[test]
6323 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6324 let now = Timestamp::now();
6325 let fallback = Duration::from_secs(300);
6326 let cap = Duration::from_secs(1800);
6327
6328 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6330
6331 let soon = now + jiff::SignedDuration::from_secs(600);
6333 assert_eq!(
6334 quota_wait(Some(soon), now, fallback, cap),
6335 Duration::from_secs(600)
6336 );
6337
6338 let past = now - jiff::SignedDuration::from_secs(60);
6341 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6342
6343 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6346 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6347 }
6348
6349 #[test]
6350 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6351 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6352
6353 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6354 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6355
6356 let already_past =
6360 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6361 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6362
6363 assert!(
6364 parse_reset_hint("session limit reached", now, now).is_none(),
6365 "free text with no recognised shape is not guessed at"
6366 );
6367 assert!(
6368 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6369 "an unresolvable zone name is not guessed at either"
6370 );
6371 }
6372
6373 #[test]
6374 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6375 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6376
6377 let at = parse_reset_hint(
6378 "You've hit your usage limit. Visit \
6379 https://chatgpt.com/codex/settings/usage to purchase more \
6380 credits or try again at Sep 19th, 2026 5:10 PM.",
6381 now,
6382 now,
6383 )
6384 .expect("the codex reset wording is a recognised shape");
6385 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6386
6387 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6392 .expect("an explicit year needs no rollover");
6393 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6394
6395 assert!(
6396 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6397 "a two-digit year is not the documented shape and is not guessed at"
6398 );
6399 assert!(
6400 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6401 "a four-letter month name is not the documented three-letter abbreviation"
6402 );
6403 assert!(
6404 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6405 "an explicit zone on the dated shape is a format nobody has \
6406 documented, and is refused rather than guessed at as UTC"
6407 );
6408 }
6409
6410 #[test]
6411 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6412 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6413 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6414
6415 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6416 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6417
6418 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6419 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6420
6421 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6422 assert!(
6423 parse_reset_hint(bad, now, recorded).is_none(),
6424 "{bad:?} must not be guessed at"
6425 );
6426 }
6427 }
6428
6429 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6433 let config = dir.join("magi.toml");
6434 std::fs::write(
6435 &config,
6436 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6437 )
6438 .unwrap();
6439 let opts = Opts {
6440 poll: Duration::from_secs(30),
6441 config: Some(config),
6442 repo: dir.join("repo"),
6446 ..Opts::default()
6447 };
6448 let home = dir.join("home");
6457 let worktrees = dir.join("wt");
6458 (
6459 opts,
6460 Queue::at(dir.join("queue")),
6461 home.join("daemon.json"),
6462 home,
6463 worktrees,
6464 )
6465 }
6466
6467 #[test]
6468 fn a_stop_is_idempotent_and_once_set_stays_set() {
6469 let stop = Stop::new();
6470 assert!(!stop.stopped());
6471
6472 stop.stop();
6473 assert!(stop.stopped());
6474 stop.stop();
6475 assert!(stop.stopped(), "a second stop is not a toggle");
6476
6477 let shared = stop.clone();
6478 assert!(
6479 shared.stopped(),
6480 "a clone is the same stop; that is how the loop and its caller share one"
6481 );
6482 }
6483
6484 #[test]
6485 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6486 let stop = Stop::new();
6487 stop.enter();
6488 assert!(
6489 !stop.finishing(),
6490 "a busy loop nobody has asked to stop is just running"
6491 );
6492
6493 stop.stop();
6494 assert!(
6495 stop.finishing(),
6496 "a stop asked for mid-run has not landed until the run is settled"
6497 );
6498
6499 stop.exit();
6500 assert!(
6501 !stop.finishing(),
6502 "once the run is settled the stop has landed and there is nothing to finish"
6503 );
6504 }
6505
6506 #[test]
6507 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6508 let stop = Stop::new();
6509 stop.enter();
6510 stop.enter();
6511 stop.stop();
6512 assert!(stop.finishing(), "two runs still in flight");
6513
6514 stop.exit();
6515 assert!(
6516 stop.finishing(),
6517 "one run finished, but a sibling is still working"
6518 );
6519
6520 stop.exit();
6521 assert!(
6522 !stop.finishing(),
6523 "the last run out is what actually lands the stop"
6524 );
6525 }
6526
6527 #[tokio::test]
6528 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6529 let dir = tempfile::tempdir().unwrap();
6530 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6531 let stop = Stop::new();
6532 stop.stop();
6533
6534 let began = std::time::Instant::now();
6535 tokio::time::timeout(
6536 Duration::from_secs(2),
6537 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6538 )
6539 .await
6540 .expect("a stopped loop must return, not sit out its poll interval")
6541 .expect("the loop's own setup and teardown must not fail");
6542 assert!(
6543 began.elapsed() < opts.poll,
6544 "returned only after {:?}, which is a poll interval, not a stop",
6545 began.elapsed()
6546 );
6547 }
6548
6549 #[tokio::test]
6550 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6551 let dir = tempfile::tempdir().unwrap();
6552 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6553 let stop = Stop::new();
6554
6555 let asker = {
6558 let stop = stop.clone();
6559 tokio::spawn(async move {
6560 tokio::time::sleep(Duration::from_millis(20)).await;
6561 stop.stop();
6562 })
6563 };
6564
6565 let began = std::time::Instant::now();
6566 tokio::time::timeout(
6567 Duration::from_secs(2),
6568 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6569 )
6570 .await
6571 .expect("a stop asked for while idle must wake the wait")
6572 .expect("the loop's own setup and teardown must not fail");
6573 asker.await.unwrap();
6574 assert!(
6575 began.elapsed() < opts.poll,
6576 "returned only after {:?}, so the stop waited on the sleep",
6577 began.elapsed()
6578 );
6579 }
6580
6581 #[tokio::test]
6582 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6583 let dir = tempfile::tempdir().unwrap();
6584 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6585 let stop = Stop::new();
6586 stop.stop();
6587
6588 tokio::time::timeout(
6589 Duration::from_secs(2),
6590 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6591 )
6592 .await
6593 .expect("a stopped loop must return")
6594 .expect("the loop's own setup and teardown must not fail");
6595
6596 assert!(
6597 home.is_dir(),
6598 "the loop did publish a status file, so its removal is the teardown and not an absence"
6599 );
6600 assert!(
6601 !status_file.exists(),
6602 "a stopped loop clears its status file"
6603 );
6604 assert!(
6605 read_status(&home).is_none(),
6606 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6607 );
6608 }
6609
6610 #[tokio::test]
6611 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6612 let dir = tempfile::tempdir().unwrap();
6613 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6614 opts.once = true;
6615
6616 let mut settled = RunState::new(
6617 dir.path().join("repo"),
6618 "main".to_owned(),
6619 "abc1234".to_owned(),
6620 "fixture".to_owned(),
6621 Config::default(),
6622 );
6623 settled.status = RunStatus::Ready;
6624 let run_dir = home.join("runs").join(&settled.id);
6625 std::fs::create_dir_all(&run_dir).unwrap();
6626 std::fs::write(
6627 run_dir.join("run.json"),
6628 serde_json::to_string_pretty(&settled).unwrap(),
6629 )
6630 .unwrap();
6631 let questions = Questions::at(home.join("questions"));
6632 let mut question = ask::Question::new(
6633 settled.id.clone(),
6634 "review".to_owned(),
6635 "reviewer-1".to_owned(),
6636 "Continue?".to_owned(),
6637 String::new(),
6638 Vec::new(),
6639 );
6640 questions.put(&mut question).unwrap();
6641
6642 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6643 .await
6644 .unwrap();
6645
6646 assert_eq!(
6647 questions.get(&question.id).unwrap().status,
6648 ask::QuestionStatus::Abandoned,
6649 "an empty --once drain still performs startup question cleanup"
6650 );
6651 }
6652
6653 #[test]
6654 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6655 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6656
6657 assert!(
6658 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6659 "never checked before: due at once"
6660 );
6661
6662 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6663 assert!(
6664 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6665 "well inside the interval: not due yet"
6666 );
6667
6668 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6669 assert!(
6670 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6671 "exactly at the edge: not yet due, same convention as `clean::due`"
6672 );
6673
6674 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6675 assert!(
6676 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6677 "past the interval: due again"
6678 );
6679 }
6680
6681 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
6686 let config = dir.join("magi.toml");
6687 std::fs::write(
6693 &config,
6694 format!(
6695 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
6696 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
6697 cache_dir.display()
6698 ),
6699 )
6700 .unwrap();
6701 Opts {
6702 config: Some(config),
6703 repo: dir.join("repo"),
6704 ..Opts::default()
6705 }
6706 }
6707
6708 #[tokio::test]
6709 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
6710 let dir = tempfile::tempdir().unwrap();
6711 let home = dir.path().join("home");
6712 let cache_dir = dir.path().join("cache");
6713 std::fs::create_dir_all(&cache_dir).unwrap();
6714 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6715 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6716
6717 let running = Stop::new();
6720 let mut last_checked = None;
6721 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6722 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
6723 .await;
6724 assert_eq!(
6725 crate::disk::dir_size(&cache_dir),
6726 0,
6727 "over the cap on the first check ever: pruned at once, no idle queue required"
6728 );
6729 assert_eq!(last_checked, Some(t0));
6730
6731 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
6733 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
6734 maybe_prune_cache_between_runs(
6735 &opts.repo,
6736 &opts,
6737 &home,
6738 &running,
6739 &mut last_checked,
6740 too_soon,
6741 )
6742 .await;
6743 assert_eq!(
6744 crate::disk::dir_size(&cache_dir),
6745 10,
6746 "too soon since the last check: left alone rather than rescanned every call"
6747 );
6748 assert_eq!(
6749 last_checked,
6750 Some(t0),
6751 "an idle check does not reset the clock"
6752 );
6753
6754 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6756 maybe_prune_cache_between_runs(
6757 &opts.repo,
6758 &opts,
6759 &home,
6760 &running,
6761 &mut last_checked,
6762 due_again,
6763 )
6764 .await;
6765 assert_eq!(
6766 crate::disk::dir_size(&cache_dir),
6767 0,
6768 "due again: pruned back under the cap"
6769 );
6770 }
6771
6772 #[tokio::test]
6780 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
6781 let dir = tempfile::tempdir().unwrap();
6782 let home = dir.path().join("home");
6783 let cache_dir = dir.path().join("cache");
6784 std::fs::create_dir_all(&cache_dir).unwrap();
6785 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6786 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6787
6788 let stop = Stop::new();
6789 stop.stop();
6790 assert!(
6791 !stop.finishing(),
6792 "no run is in flight at a between-runs boundary, so nothing else \
6793 would tell the operator this stop had not taken effect yet"
6794 );
6795
6796 let mut last_checked = None;
6797 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6798 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
6799 .await;
6800 assert_eq!(
6801 crate::disk::dir_size(&cache_dir),
6802 10,
6803 "over its cap, and due for the first check ever, but a stop outranks \
6804 it: the cap is a standing policy the next start measures again"
6805 );
6806 assert_eq!(
6807 last_checked, None,
6808 "a check that never happened must not claim the interval"
6809 );
6810 }
6811
6812 #[tokio::test]
6826 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
6827 let dir = tempfile::tempdir().unwrap();
6828 let cache_dir = dir.path().join("cache");
6829 std::fs::create_dir_all(&cache_dir).unwrap();
6830 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
6831
6832 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
6833 opts.poll = Duration::from_millis(20);
6834 opts.max_attempts = 1_000;
6835
6836 let queue = Queue::at(dir.path().join("queue"));
6837 let mut t = Task::new(
6838 "x".to_owned(),
6839 "x".to_owned(),
6840 opts.repo.clone(),
6841 Source::Human,
6842 );
6843 queue.put(&mut t).unwrap();
6844
6845 let home = dir.path().join("home");
6846 let worktrees = dir.path().join("wt");
6847 let status_file = home.join("daemon.json");
6848 let stop = Stop::new();
6849 let stopper = {
6850 let stop = stop.clone();
6851 tokio::spawn(async move {
6852 tokio::time::sleep(Duration::from_millis(400)).await;
6853 stop.stop();
6854 })
6855 };
6856
6857 tokio::time::timeout(
6858 Duration::from_secs(10),
6859 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6860 )
6861 .await
6862 .expect("the loop must not hang on a queue that keeps producing failing work")
6863 .expect("the loop's own setup and teardown must not fail");
6864 stopper.await.unwrap();
6865
6866 let after = queue.get(&t.id).unwrap();
6867 assert!(
6868 after.attempts >= 2,
6869 "the harness must actually have retried more than once, or this is not \
6870 exercising a busy queue at all (got {} attempt(s))",
6871 after.attempts
6872 );
6873 assert!(
6874 after.status.runnable(),
6875 "still under its attempt budget: the queue never reached a natural idle \
6876 on its own, only the external stop ended the test"
6877 );
6878
6879 assert_eq!(
6880 crate::disk::dir_size(&cache_dir),
6881 0,
6882 "an oversized cache must not be left to grow unboundedly just because the \
6883 queue kept the loop busy the whole time"
6884 );
6885 }
6886
6887 #[test]
6888 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
6889 let dir = tempfile::tempdir().unwrap();
6890 let queue = Queue::at(dir.path().join("queue"));
6891 let questions = Questions::at(dir.path().join("questions"));
6892 let mut task = task();
6893 queue.put(&mut task).unwrap();
6894
6895 let mut task_question = ask::Question::new(
6896 task.id.clone(),
6897 crate::conduct::NODE.to_owned(),
6898 "conduct".to_owned(),
6899 "Which backend?".to_owned(),
6900 String::new(),
6901 Vec::new(),
6902 );
6903 questions.put(&mut task_question).unwrap();
6904 task.block(vec![task_question.id.clone()], None);
6905 queue.put(&mut task).unwrap();
6906
6907 let mut run_question = ask::Question::new(
6908 "20260101-000000-run1".to_owned(),
6909 "review".to_owned(),
6910 "reviewer-1".to_owned(),
6911 "Run question".to_owned(),
6912 String::new(),
6913 Vec::new(),
6914 );
6915 questions.put(&mut run_question).unwrap();
6916
6917 let mut coincidental = ask::Question::new(
6922 task.id.clone(),
6923 "review".to_owned(),
6924 "reviewer-1".to_owned(),
6925 "Unrelated review question".to_owned(),
6926 String::new(),
6927 Vec::new(),
6928 );
6929 questions.put(&mut coincidental).unwrap();
6930
6931 reconcile_task_questions(&queue, &questions);
6932 assert!(questions.get(&task_question.id).unwrap().status.open());
6933 assert!(questions.get(&run_question.id).unwrap().status.open());
6934 assert!(questions.get(&coincidental.id).unwrap().status.open());
6935
6936 task.release();
6937 queue.put(&mut task).unwrap();
6938 reconcile_task_questions(&queue, &questions);
6939 assert_eq!(
6940 questions.get(&task_question.id).unwrap().status,
6941 ask::QuestionStatus::Abandoned
6942 );
6943 assert!(
6944 questions.get(&run_question.id).unwrap().status.open(),
6945 "run questions remain the run janitor's responsibility"
6946 );
6947 assert!(
6948 questions.get(&coincidental.id).unwrap().status.open(),
6949 "a non-conductor question must not be abandoned just because its \
6950 run id coincides with a task id"
6951 );
6952 }
6953
6954 #[test]
6955 fn a_freshly_started_running_task_is_never_stalled() {
6956 let dir = tempfile::tempdir().unwrap();
6957 let mut t = task();
6958 t.start("run-1".to_owned());
6959 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
6962 }
6963
6964 #[test]
6965 fn a_long_running_task_with_no_live_daemon_is_stalled() {
6966 let dir = tempfile::tempdir().unwrap();
6967 let mut t = task();
6968 t.start("run-1".to_owned());
6969 t.updated_at = Timestamp::now()
6970 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6971 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
6972 assert_eq!(
6973 stalled_tasks(
6974 &Queue::at(dir.path().join("q")),
6975 dir.path(),
6976 Timestamp::now()
6977 )
6978 .len(),
6979 0,
6980 "the task was never written to this queue"
6981 );
6982 }
6983
6984 #[test]
6985 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
6986 let dir = tempfile::tempdir().unwrap();
6987 let mut t = task();
6988 t.id = "20260903-080340-0167".to_owned();
6989 t.start("20260903-080619-01c2".to_owned());
6990 t.updated_at = Timestamp::now()
6991 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6992
6993 let mut status = Status::new();
6994 status.current = vec![Current {
6995 task: t.id.clone(),
6996 run: "20260903-080619-01c2".to_owned(),
6997 }];
6998 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6999
7000 assert!(
7001 !is_stalled(&t, dir.path(), Timestamp::now()),
7002 "a live daemon's own heartbeat rules out stalled, however long the task has run"
7003 );
7004 }
7005
7006 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
7010 let path = queue.path_of(id);
7011 let body = std::fs::read_to_string(&path).unwrap();
7012 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
7013 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
7014 v["updated_at"] = serde_json::Value::String(old.to_string());
7015 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
7016 }
7017
7018 #[test]
7019 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
7020 let dir = tempfile::tempdir().unwrap();
7033 let queue = Queue::at(dir.path().join("queue"));
7034 let home = dir.path().join("home");
7035
7036 let mut t = task();
7037 t.id = "20260101-000001-lock".to_owned();
7038 t.start("run-1".to_owned());
7039 queue.put(&mut t).unwrap();
7040 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7041 std::fs::write(
7042 dir.path().join("queue").join(format!("{}.lock", t.id)),
7043 "not a pid",
7044 )
7045 .unwrap();
7046
7047 let now = Timestamp::now();
7048 assert!(
7049 reclaim_orphaned_running(&queue, 2).is_empty(),
7050 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
7051 and reclaim must leave the task alone"
7052 );
7053 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
7054
7055 let stalled = stalled_tasks(&queue, &home, now);
7056 assert_eq!(
7057 stalled.len(),
7058 1,
7059 "reclaim's inability to claim it yet must not hide it from the conductor"
7060 );
7061 assert_eq!(stalled[0].id, t.id);
7062 }
7063
7064 #[test]
7065 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
7066 let dir = tempfile::tempdir().unwrap();
7067 crate::run::set_home(dir.path().join("run-home"));
7068 let queue = Queue::at(dir.path().join("queue"));
7069 let home = dir.path().join("home");
7070 let questions = Questions::at(dir.path().join("questions"));
7071
7072 let mut t = task();
7073 t.id = "20260101-000003-dead".to_owned();
7074 t.start("missing-run".to_owned());
7075 queue.put(&mut t).unwrap();
7076 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7077
7078 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7081 assert_eq!(
7082 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
7083 [&t.id]
7084 );
7085 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
7086 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7087
7088 crate::conduct::apply(
7091 &queue,
7092 &questions,
7093 &crate::conduct::Verdict {
7094 decisions: vec![crate::conduct::Decision {
7095 id: t.id.clone(),
7096 recovery: Some(crate::conduct::Recovery::Requeue),
7097 ..crate::conduct::Decision::default()
7098 }],
7099 },
7100 )
7101 .unwrap();
7102 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7103 }
7104
7105 #[test]
7106 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
7107 let dir = tempfile::tempdir().unwrap();
7108 let queue = Queue::at(dir.path().join("queue"));
7109 let home = dir.path().join("home");
7110
7111 let mut fresh = task();
7112 fresh.id = "20260101-000001-aaaa".to_owned();
7113 fresh.start("run-1".to_owned());
7114 queue.put(&mut fresh).unwrap();
7115
7116 let mut old = task();
7117 old.id = "20260101-000002-bbbb".to_owned();
7118 old.start("run-2".to_owned());
7119 queue.put(&mut old).unwrap();
7120 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
7121
7122 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7123 assert_eq!(stalled.len(), 1);
7124 assert_eq!(stalled[0].id, old.id);
7125 }
7126
7127 #[test]
7128 fn queued_and_finished_task_views_partition_by_status() {
7129 let dir = tempfile::tempdir().unwrap();
7130 let queue = Queue::at(dir.path().join("queue"));
7131
7132 let mut queued = task();
7133 queued.id = "20260101-000001-aaaa".to_owned();
7134 queue.put(&mut queued).unwrap();
7135
7136 let mut failed = task();
7137 failed.id = "20260101-000002-bbbb".to_owned();
7138 failed.start("run-1".to_owned());
7139 failed.fail("gate red", 5);
7140 queue.put(&mut failed).unwrap();
7141
7142 let mut held = task();
7143 held.id = "20260101-000003-cccc".to_owned();
7144 held.hold_machine(None);
7145 queue.put(&mut held).unwrap();
7146
7147 let mut running = task();
7148 running.id = "20260101-000004-dddd".to_owned();
7149 running.start("run-2".to_owned());
7150 queue.put(&mut running).unwrap();
7151
7152 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
7153 assert_eq!(queued_ids, [queued.id.clone()]);
7154
7155 let mut finished_ids: Vec<String> =
7156 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
7157 finished_ids.sort_unstable();
7158 let mut want = vec![failed.id.clone(), held.id.clone()];
7159 want.sort_unstable();
7160 assert_eq!(finished_ids, want);
7161 }
7162
7163 #[test]
7164 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
7165 let dir = tempfile::tempdir().unwrap();
7166 let queue = Queue::at(dir.path().join("queue"));
7167 let questions = ask::Questions::at(dir.path().join("questions"));
7168
7169 let mut dep = task();
7170 dep.id = "20260101-000001-dep0".to_owned();
7171 dep.succeed();
7172 queue.put(&mut dep).unwrap();
7173
7174 let mut still_going = task();
7175 still_going.id = "20260101-000002-dep1".to_owned();
7176 queue.put(&mut still_going).unwrap();
7177
7178 let mut blocked = task();
7179 blocked.id = "20260101-000003-main".to_owned();
7180 blocked.block(
7181 vec![dep.id.clone(), still_going.id.clone()],
7182 Some("waits on both".to_owned()),
7183 );
7184 queue.put(&mut blocked).unwrap();
7185
7186 resolve_blockers(&queue, &questions);
7187
7188 let after = queue.get(&blocked.id).unwrap();
7189 assert_eq!(
7190 after.status,
7191 TaskStatus::Blocked,
7192 "one dependency is still outstanding"
7193 );
7194 assert_eq!(after.blocked_by, [still_going.id.clone()]);
7195 }
7196
7197 #[test]
7198 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
7199 let dir = tempfile::tempdir().unwrap();
7200 let queue = Queue::at(dir.path().join("queue"));
7201 let questions = ask::Questions::at(dir.path().join("questions"));
7202
7203 let mut q = crate::ask::Question::new(
7204 "20260101-000001-main".to_owned(),
7205 crate::conduct::NODE.to_owned(),
7206 "conduct".to_owned(),
7207 "Which backend?".to_owned(),
7208 String::new(),
7209 Vec::new(),
7210 );
7211 questions.put(&mut q).unwrap();
7212 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
7213 .unwrap();
7214 questions.put(&mut q).unwrap();
7215
7216 let mut blocked = task();
7217 blocked.id = "20260101-000001-main".to_owned();
7218 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
7219 queue.put(&mut blocked).unwrap();
7220
7221 resolve_blockers(&queue, &questions);
7222
7223 let after = queue.get(&blocked.id).unwrap();
7224 assert_eq!(
7225 after.status,
7226 TaskStatus::Queued,
7227 "the only blocker resolved"
7228 );
7229 assert_eq!(after.answers.len(), 1);
7230 assert_eq!(after.answers[0].question, "Which backend?");
7231 assert_eq!(after.answers[0].answer, "SQLite");
7232
7233 let instruction = instruction_for(&after);
7235 assert!(instruction.contains("Which backend?"));
7236 assert!(instruction.contains("SQLite"));
7237 }
7238
7239 #[test]
7240 fn resolve_blockers_holds_a_task_whose_conductor_question_was_abandoned() {
7241 let dir = tempfile::tempdir().unwrap();
7242 let queue = Queue::at(dir.path().join("queue"));
7243 let questions = ask::Questions::at(dir.path().join("questions"));
7244
7245 let mut q = crate::ask::Question::new(
7246 "20260101-000001-main".to_owned(),
7247 crate::conduct::NODE.to_owned(),
7248 "conduct".to_owned(),
7249 "Is the setup done?".to_owned(),
7250 String::new(),
7251 Vec::new(),
7252 );
7253 q.abandon("no answer within 60s of asking");
7254 questions.put(&mut q).unwrap();
7255
7256 let mut blocked = task();
7257 blocked.id = "20260101-000001-main".to_owned();
7258 blocked.block(vec![q.id.clone()], Some("setup?".to_owned()));
7259 queue.put(&mut blocked).unwrap();
7260
7261 resolve_blockers(&queue, &questions);
7262
7263 let after = queue.get(&blocked.id).unwrap();
7264 assert_eq!(
7265 after.status,
7266 TaskStatus::Held,
7267 "never left blocked on nothing"
7268 );
7269 assert!(!after.operator_held(), "a machine hold, for triage");
7270 assert!(
7271 after
7272 .hold_reason
7273 .as_deref()
7274 .unwrap_or_default()
7275 .contains("went unanswered")
7276 );
7277 }
7278
7279 #[test]
7280 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
7281 let dir = tempfile::tempdir().unwrap();
7287 let queue = Queue::at(dir.path().join("queue"));
7288 let questions = ask::Questions::at(dir.path().join("questions"));
7289
7290 let mut q = crate::ask::Question::new(
7291 "20260101-000001-main".to_owned(),
7292 crate::conduct::NODE.to_owned(),
7293 "conduct".to_owned(),
7294 "How should this be handled?".to_owned(),
7295 String::new(),
7296 Vec::new(),
7297 );
7298 questions.put(&mut q).unwrap();
7299 q.answer(crate::ask::Answer::Text(
7300 "leave it held, a human will look at it later".to_owned(),
7301 ))
7302 .unwrap();
7303 questions.put(&mut q).unwrap();
7304
7305 let mut held = task();
7306 held.id = "20260101-000001-main".to_owned();
7307 held.hold_machine(Some("out of attempts".to_owned()));
7308 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
7309 queue.put(&mut held).unwrap();
7310
7311 resolve_blockers(&queue, &questions);
7312
7313 let after = queue.get(&held.id).unwrap();
7314 assert_eq!(after.status, TaskStatus::Held);
7315 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
7316 assert_eq!(
7317 after.answers[0].answer,
7318 "leave it held, a human will look at it later"
7319 );
7320 }
7321
7322 #[test]
7323 fn resolve_blockers_releases_a_task_whose_dependency_was_deleted_on_purpose() {
7324 let dir = tempfile::tempdir().unwrap();
7325 let queue = Queue::at(dir.path().join("queue"));
7326 let questions = ask::Questions::at(dir.path().join("questions"));
7327 let mut dep = task();
7328 dep.id = "20260101-000001-gone".to_owned();
7329 queue.put(&mut dep).unwrap();
7330 let mut blocked = task();
7331 blocked.id = "20260101-000003-main".to_owned();
7332 blocked.block(vec![dep.id.clone()], None);
7333 queue.put(&mut blocked).unwrap();
7334
7335 let claim = queue.claim(&blocked.id).unwrap();
7337 queue.remove(&dep.id, false, &questions).unwrap();
7338 drop(claim);
7339 resolve_blockers(&queue, &questions);
7340
7341 let after = queue.get(&blocked.id).unwrap();
7342 assert_eq!(after.status, TaskStatus::Queued);
7343 assert!(after.blocked_by.is_empty());
7344 }
7345
7346 #[test]
7347 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
7348 let dir = tempfile::tempdir().unwrap();
7354 let queue = Queue::at(dir.path().join("queue"));
7355 let questions = ask::Questions::at(dir.path().join("questions"));
7356
7357 let mut still_going = task();
7358 still_going.id = "20260101-000002-dep1".to_owned();
7359 queue.put(&mut still_going).unwrap();
7360
7361 let mut blocked = task();
7362 blocked.id = "20260101-000003-main".to_owned();
7363 blocked.block(
7364 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7365 Some("waits on both".to_owned()),
7366 );
7367 queue.put(&mut blocked).unwrap();
7368
7369 resolve_blockers(&queue, &questions);
7370
7371 let after = queue.get(&blocked.id).unwrap();
7372 assert_eq!(
7373 after.status,
7374 TaskStatus::Held,
7375 "a missing dependency must not leave the task blocked forever"
7376 );
7377 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7378 assert!(after.blocked_by.is_empty());
7379 let reason = after.hold_reason.as_deref().unwrap_or_default();
7380 assert!(
7381 reason.contains("20260101-000001-gone"),
7382 "the missing id must be named so an operator can tell what happened: {reason}"
7383 );
7384 assert!(
7385 reason.contains(&still_going.id),
7386 "the still-valid dependency must not silently vanish from the record: {reason}"
7387 );
7388 }
7389
7390 #[test]
7391 fn instruction_for_is_unchanged_without_any_answers() {
7392 let t = task();
7393 assert_eq!(instruction_for(&t), t.instruction);
7394 }
7395
7396 #[test]
7397 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7398 let dir = tempfile::tempdir().unwrap();
7399 let q = Queue::at(dir.path().join("queue"));
7400 let src = dir.path().join("shot.png");
7401 std::fs::write(&src, "x").unwrap();
7402 let mut t = task();
7403 q.attach(&mut t, &[src]).unwrap();
7404 let paths = task_attachments(&q, &t).unwrap();
7405 assert_eq!(paths.len(), 1);
7406 assert!(paths[0].is_absolute() && paths[0].is_file());
7407 std::fs::remove_file(&paths[0]).unwrap();
7408 let err = task_attachments(&q, &t).unwrap_err().to_string();
7409 assert!(err.contains("shot.png"), "{err}");
7410 }
7411
7412 #[test]
7413 fn resumed_instruction_is_unchanged_without_any_answers() {
7414 let t = task();
7415 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7416 }
7417
7418 #[test]
7419 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7420 let mut t = task();
7421 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7422 let old = t.instruction.clone();
7426
7427 let refreshed = resumed_instruction(&old, &t);
7428 assert!(refreshed.starts_with(&old), "the original text is kept");
7429 assert!(refreshed.contains("Which backend?"));
7430 assert!(refreshed.contains("SQLite"));
7431 }
7432
7433 #[test]
7434 fn resumed_instruction_keeps_an_original_answers_heading() {
7435 let mut t = task();
7436 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7437 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7438
7439 let refreshed = resumed_instruction(&t.instruction, &t);
7440
7441 assert!(
7442 refreshed.starts_with(&t.instruction),
7443 "an answers heading in the original instruction is not the appended block"
7444 );
7445 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7446 assert!(refreshed.contains("Which backend?"));
7447 assert!(refreshed.contains("SQLite"));
7448
7449 let repeated = resumed_instruction(&refreshed, &t);
7450 assert_eq!(
7451 repeated, refreshed,
7452 "only the final appended block is refreshed"
7453 );
7454 }
7455
7456 #[test]
7457 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7458 let mut t = task();
7459 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7460
7461 let once = resumed_instruction(&t.instruction, &t);
7465 let twice = resumed_instruction(&once, &t);
7466 assert_eq!(once, twice);
7467 assert_eq!(once.matches("Which backend?").count(), 1);
7468
7469 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7471 let refreshed = resumed_instruction(&once, &t);
7472 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7473 assert!(refreshed.contains("Which backend?"));
7474 assert!(refreshed.contains("Which cache?"));
7475 }
7476
7477 #[test]
7478 fn prepare_instruction_covers_all_three_starters() {
7479 let mut t = task();
7480 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7481
7482 assert_eq!(
7485 prepare_instruction(&Starter::Start, None, &t),
7486 Some(instruction_for(&t))
7487 );
7488
7489 let old = t.instruction.clone();
7492 assert_eq!(
7493 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7494 Some(resumed_instruction(&old, &t))
7495 );
7496
7497 assert_eq!(
7501 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7502 None
7503 );
7504 }
7505
7506 #[test]
7507 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7508 assert_eq!(
7509 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7510 Starter::Review("magi/eba2/A".to_owned())
7511 );
7512 }
7513
7514 #[test]
7515 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7516 assert_eq!(
7517 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7518 Starter::Start,
7519 "a vanished review branch must not fall back to resuming the old run either"
7520 );
7521 }
7522
7523 #[test]
7524 fn a_refused_handover_retries_as_a_review_of_the_same_branch() {
7525 let mut t = task();
7526 t.start("old-run".to_owned());
7527 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
7528 t.release();
7529 let branch = t.review_branch.take();
7530 assert_eq!(
7531 choose_starter(branch.as_deref(), true, Some("old-run")),
7532 Starter::Review("magi/eba2/A".to_owned()),
7533 "a review wins over resuming the old run"
7534 );
7535 }
7536
7537 #[test]
7538 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7539 assert_eq!(
7540 choose_starter(None, false, Some("some-run")),
7541 Starter::Resume("some-run".to_owned())
7542 );
7543 assert_eq!(choose_starter(None, false, None), Starter::Start);
7544 }
7545
7546 #[test]
7547 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7548 let mut released = task();
7549 released.start("stalled-run".to_owned());
7550 released.requeue();
7551 let unfinished = (!released.fresh_start)
7552 .then(|| Some("stalled-run".to_owned()))
7553 .flatten();
7554 assert_eq!(
7555 choose_starter(None, false, unfinished.as_deref()),
7556 Starter::Start,
7557 "release keeps run history but must not resume it"
7558 );
7559 assert_eq!(released.runs, ["stalled-run"]);
7560 }
7561
7562 #[test]
7563 fn an_ordinary_release_keeps_a_resumable_run_available() {
7564 let mut released = task();
7565 released.start("stalled-run".to_owned());
7566 released.release();
7567 let unfinished = (!released.fresh_start)
7568 .then(|| Some("stalled-run".to_owned()))
7569 .flatten();
7570 assert_eq!(
7571 choose_starter(None, false, unfinished.as_deref()),
7572 Starter::Resume("stalled-run".to_owned()),
7573 "manual release must preserve the normal resume path"
7574 );
7575 }
7576
7577 #[test]
7578 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7579 let mut state = run_state(RunStatus::Blocked);
7580 state.config.graph.review_rounds = 3;
7581 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7582 assert!(exhausted_review_budget(&state));
7583
7584 state.reviews.pop();
7586 assert!(!exhausted_review_budget(&state));
7587
7588 let mut stalled = run_state(RunStatus::Stalled);
7591 stalled.config.graph.review_rounds = 1;
7592 stalled.reviews = vec![review_round(1)];
7593 assert!(!exhausted_review_budget(&stalled));
7594 }
7595
7596 fn review_round(round: usize) -> crate::run::ReviewRound {
7597 crate::run::ReviewRound {
7598 round,
7599 head: "deadbeef".to_owned(),
7600 verified_head: None,
7601 verified_at: None,
7602 reviews: Vec::new(),
7603 e2e: Vec::new(),
7604 verify_retried: false,
7605 e2e_deferred: false,
7606 e2e_defer_reason: None,
7607 fix: None,
7608 blocking: 0,
7609 answered: 1,
7610 expected: 1,
7611 clean: false,
7612 progressed: true,
7613 vote_split: false,
7614 reconsideration: Vec::new(),
7615 verdict: None,
7616 }
7617 }
7618
7619 #[test]
7620 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7621 let mut run = RunState::new(
7622 PathBuf::from("/repo"),
7623 "main".to_owned(),
7624 "abc1234def".to_owned(),
7625 "add retries".to_owned(),
7626 Config::default(),
7627 );
7628 run.status = RunStatus::Judging;
7629 run.parked = true;
7630 let mut task = Task::new(
7631 "add retries".to_owned(),
7632 "add retries".to_owned(),
7633 PathBuf::from("/repo"),
7634 crate::queue::Source::Human,
7635 );
7636 task.status = TaskStatus::Failed;
7637 task.runs = vec![run.id.clone()];
7638 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7639 assert!(with(&task, &run), "parked after judging is the case");
7640
7641 let mut not_parked = run.clone();
7642 not_parked.parked = false;
7643 not_parked.status = RunStatus::Stalled;
7644 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7645
7646 let mut fresh = task.clone();
7647 fresh.fresh_start = true;
7648 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7649
7650 let mut review = task.clone();
7651 review.review_branch = Some("magi/x/A".to_owned());
7652 assert!(!with(&review, &run), "review is ranked before resume");
7653
7654 let mut held = task.clone();
7655 held.status = TaskStatus::Held;
7656 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7657
7658 let mut released = run.clone();
7659 released.released_to = Some("20260901-000000-new1".to_owned());
7660 assert!(!with(&task, &released), "nothing left to resume into");
7661
7662 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7663 "unreadable"
7664 )));
7665 }
7666
7667 #[test]
7668 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7669 let mut released = RunState::new(
7670 PathBuf::from("/repo"),
7671 "main".to_owned(),
7672 "abc1234def".to_owned(),
7673 "add retries".to_owned(),
7674 Config::default(),
7675 );
7676 released.status = RunStatus::Blocked;
7677 assert_eq!(
7678 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7679 Some(released.id.clone())
7680 );
7681 released.released_to = Some("20260901-000000-new1".to_owned());
7682 assert_eq!(
7683 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7684 None,
7685 "there is nothing left to resume it into"
7686 );
7687 }
7688
7689 #[test]
7690 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
7691 let mut exhausted = RunState::new(
7700 PathBuf::from("/repo"),
7701 "main".to_owned(),
7702 "abc1234def".to_owned(),
7703 "add retries".to_owned(),
7704 Config::default(),
7705 );
7706 exhausted.status = RunStatus::Blocked;
7707 exhausted.config.graph.review_rounds = 1;
7708 exhausted.reviews = vec![review_round(1)];
7709
7710 assert_eq!(
7711 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
7712 None,
7713 "an exhausted `Blocked` run must not be offered as resumable"
7714 );
7715
7716 let mut has_budget_left = RunState::new(
7719 PathBuf::from("/repo"),
7720 "main".to_owned(),
7721 "abc1234def".to_owned(),
7722 "add retries".to_owned(),
7723 Config::default(),
7724 );
7725 has_budget_left.status = RunStatus::Blocked;
7726 has_budget_left.config.graph.review_rounds = 3;
7727 has_budget_left.reviews = vec![review_round(1)];
7728
7729 assert_eq!(
7730 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
7731 Ok(has_budget_left.clone())
7732 }),
7733 Some(has_budget_left.id.clone())
7734 );
7735 }
7736
7737 #[test]
7738 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
7739 let mut older_stalled = RunState::new(
7747 PathBuf::from("/repo"),
7748 "main".to_owned(),
7749 "abc1234def".to_owned(),
7750 "add retries".to_owned(),
7751 Config::default(),
7752 );
7753 older_stalled.status = RunStatus::Stalled;
7754
7755 let mut newest_exhausted = RunState::new(
7756 PathBuf::from("/repo"),
7757 "main".to_owned(),
7758 "abc1234def".to_owned(),
7759 "add retries".to_owned(),
7760 Config::default(),
7761 );
7762 newest_exhausted.status = RunStatus::Blocked;
7763 newest_exhausted.config.graph.review_rounds = 1;
7764 newest_exhausted.reviews = vec![review_round(1)];
7765
7766 assert_eq!(
7767 unfinished_run_with(
7768 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
7769 "t",
7770 |_| Ok(newest_exhausted.clone())
7771 ),
7772 None,
7773 "the newest run is exhausted, so nothing here is worth resuming - \
7774 least of all the older, already-superseded run"
7775 );
7776 }
7777
7778 #[test]
7779 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
7780 assert_eq!(
7781 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
7782 Err(anyhow::anyhow!("fixture is absent"))
7783 }),
7784 None
7785 );
7786 }
7787
7788 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
7789 let mut q = ask::Question::new(
7790 run.to_owned(),
7791 "implement".to_owned(),
7792 "impl-A".to_owned(),
7793 "continue?".to_owned(),
7794 String::new(),
7795 vec!["resume で続行する".to_owned(), "other".to_owned()],
7796 );
7797 q.actions.insert("resume で続行する".to_owned(), action);
7798 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
7799 .unwrap();
7800 q
7801 }
7802
7803 fn held_task_with(run: &str) -> Task {
7804 let mut t = task();
7805 t.runs = vec![run.to_owned()];
7806 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
7807 t
7808 }
7809
7810 fn resume_action(run: &str) -> ask::ChoiceAction {
7811 ask::ChoiceAction::Resume { run: run.into() }
7812 }
7813
7814 #[test]
7815 fn decide_action_resumes_only_the_latest_resumable_run() {
7816 let t = held_task_with("r1");
7817 let q = action_question("r1", resume_action("r1"));
7818 let load = |s: RunState| move |_: &str| Ok(s);
7819 assert_eq!(
7820 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
7821 ActionDecision::Resume("r1".into())
7822 );
7823 let q_other = action_question("r1", resume_action("r0"));
7825 assert!(matches!(
7826 decide_action(
7827 &t,
7828 &q_other,
7829 &PHRASES_EN,
7830 load(run_state(RunStatus::Blocked))
7831 ),
7832 ActionDecision::Refuse(_)
7833 ));
7834 assert!(matches!(
7836 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
7837 ActionDecision::Refuse(_)
7838 ));
7839 let mut released = run_state(RunStatus::Blocked);
7841 released.released_to = Some("elsewhere".into());
7842 assert!(matches!(
7843 decide_action(&t, &q, &PHRASES_EN, load(released)),
7844 ActionDecision::Refuse(_)
7845 ));
7846 assert!(matches!(
7848 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
7849 ActionDecision::Refuse(_)
7850 ));
7851 }
7852
7853 #[test]
7854 fn decide_action_ignores_a_question_about_an_earlier_run() {
7855 let mut t = held_task_with("r1");
7856 t.runs.push("r2".to_owned());
7857 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7858 assert_eq!(
7859 decide_action(
7860 &t,
7861 &action_question("r1", ask::ChoiceAction::Done),
7862 &PHRASES_EN,
7863 never
7864 ),
7865 ActionDecision::Stale
7866 );
7867 }
7868
7869 #[test]
7870 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
7871 let mut t = held_task_with("r1");
7872 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7873 assert_eq!(
7874 decide_action(
7875 &t,
7876 &action_question("r1", ask::ChoiceAction::Requeue),
7877 &PHRASES_EN,
7878 never
7879 ),
7880 ActionDecision::Requeue
7881 );
7882 let done_q = action_question("r1", ask::ChoiceAction::Done);
7883 assert_eq!(
7884 decide_action(&t, &done_q, &PHRASES_EN, never),
7885 ActionDecision::Done
7886 );
7887 t.mark_action_applied(&done_q.id);
7888 assert_eq!(
7889 decide_action(&t, &done_q, &PHRASES_EN, never),
7890 ActionDecision::Skip
7891 );
7892
7893 let mut plain = action_question("r1", ask::ChoiceAction::Done);
7895 plain.actions.clear();
7896 assert_eq!(
7897 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
7898 ActionDecision::Skip
7899 );
7900 let mut running = held_task_with("r1");
7902 running.status = TaskStatus::Running;
7903 assert_eq!(
7904 decide_action(
7905 &running,
7906 &action_question("r1", ask::ChoiceAction::Done),
7907 &PHRASES_EN,
7908 never
7909 ),
7910 ActionDecision::Skip
7911 );
7912 }
7913
7914 #[test]
7915 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
7916 let dir = tempfile::tempdir().unwrap();
7917 let queue = Queue::at(dir.path().join("queue"));
7918 let questions = Questions::at(dir.path().join("questions"));
7919 let home = dir.path().join("home");
7920 let mut state = run_state(RunStatus::Blocked);
7921 state.id = "20260101-000000-act1".to_owned();
7922 state.save_under(&home).unwrap();
7923
7924 let mut t = held_task_with(&state.id);
7925 queue.put(&mut t).unwrap();
7926 let mut q = action_question(&state.id, resume_action(&state.id));
7927 questions.put(&mut q).unwrap();
7928
7929 apply_choice_actions(&queue, &questions, &home);
7930 let after = queue.get(&t.id).unwrap();
7931 assert_eq!(after.status, TaskStatus::Queued);
7932 assert!(!after.fresh_start);
7933 assert!(after.action_applied(&q.id));
7934 let pin = after.resume_override.clone().unwrap();
7935 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
7936 assert!(pin.forced);
7937
7938 let mut again = queue.get(&t.id).unwrap();
7940 again.hold_machine(Some("later".into()));
7941 queue.put(&mut again).unwrap();
7942 apply_choice_actions(&queue, &questions, &home);
7943 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7944 }
7945
7946 #[test]
7947 fn an_answer_the_waiter_already_delivered_is_not_acted_on_again() {
7948 let dir = tempfile::tempdir().unwrap();
7949 let queue = Queue::at(dir.path().join("queue"));
7950 let questions = Questions::at(dir.path().join("questions"));
7951 let home = dir.path().join("home");
7952 let mut state = run_state(RunStatus::Blocked);
7953 state.id = "20260101-000000-act2".to_owned();
7954 state.save_under(&home).unwrap();
7955
7956 let mut t = held_task_with(&state.id);
7957 queue.put(&mut t).unwrap();
7958 let mut q = action_question(&state.id, resume_action(&state.id));
7959 q.answer_delivered = true;
7960 questions.put(&mut q).unwrap();
7961
7962 apply_choice_actions(&queue, &questions, &home);
7963 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7964
7965 let mut q2 = action_question(&state.id, resume_action(&state.id));
7967 questions.put(&mut q2).unwrap();
7968 apply_choice_actions(&queue, &questions, &home);
7969 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7970 assert!(questions.get(&q2.id).unwrap().answer_delivered);
7971 }
7972}