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, Copy, PartialEq, Eq)]
773pub(crate) enum ActionStanding {
774 NoAction,
776 Pending,
778 Applied,
781 Stale,
783 Busy,
785}
786
787impl ActionStanding {
788 pub(crate) fn daemon_owns(self) -> bool {
790 matches!(self, Self::Pending | Self::Applied | Self::Stale)
791 }
792}
793
794pub(crate) fn action_standing(task: &Task, q: &ask::Question) -> ActionStanding {
798 if q.chosen_action().is_none() {
799 ActionStanding::NoAction
800 } else if task.action_applied(&q.id) {
801 ActionStanding::Applied
802 } else if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
803 ActionStanding::Stale
804 } else if matches!(
805 task.status,
806 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
807 ) {
808 ActionStanding::Busy
809 } else {
810 ActionStanding::Pending
811 }
812}
813
814#[derive(Debug, Clone, PartialEq, Eq)]
816enum ActionDecision {
817 Skip,
819 Resume(String),
821 Requeue,
823 Done,
825 Stale,
828 Refuse(String),
831}
832
833fn decide_action<F>(task: &Task, q: &ask::Question, p: &Phrases, load: F) -> ActionDecision
841where
842 F: FnOnce(&str) -> Result<RunState>,
843{
844 let Some(action) = q.chosen_action() else {
845 return ActionDecision::Skip;
846 };
847 if task.action_applied(&q.id)
848 || matches!(
849 task.status,
850 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
851 )
852 {
853 return ActionDecision::Skip;
854 }
855 if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
858 return ActionDecision::Stale;
859 }
860 match action {
861 ask::ChoiceAction::Requeue => ActionDecision::Requeue,
862 ask::ChoiceAction::Done => ActionDecision::Done,
863 ask::ChoiceAction::Resume { run } => {
864 if task.runs.last() != Some(run) {
865 return ActionDecision::Refuse((p.resume_not_latest)(
866 q.short(),
867 ask::short_id(run),
868 ));
869 }
870 match load(run) {
871 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
872 ActionDecision::Resume(run.clone())
873 }
874 Ok(_) => ActionDecision::Refuse((p.resume_cannot_progress)(
875 q.short(),
876 ask::short_id(run),
877 )),
878 Err(e) => ActionDecision::Refuse((p.resume_unreadable)(
879 q.short(),
880 ask::short_id(run),
881 &format!("{e:#}"),
882 )),
883 }
884 }
885 }
886}
887
888struct Phrases {
902 graph_stopped: fn(&str, &str) -> String,
904 quorum_lost: &'static str,
905 quota_took_out: &'static str,
907 run_ended: &'static str,
909 waiting_for_answer: &'static str,
911 recovered_running: &'static str,
912 no_run_to_recover: &'static str,
913 could_not_start: &'static str,
914 resume_not_latest: fn(&str, &str) -> String,
915 resume_cannot_progress: fn(&str, &str) -> String,
916 resume_unreadable: fn(&str, &str, &str) -> String,
917 handover_refused: fn(&crate::handover::Refused) -> String,
920 handover_hint: &'static str,
922 reviewers_never_answered: fn(usize, usize) -> String,
925}
926
927const PHRASES_EN: Phrases = Phrases {
928 graph_stopped: |status, detail| {
929 format!("the graph stopped at `{status}` without reaching a terminal status: {detail}")
930 },
931 quorum_lost: "the judging panel lost its quorum",
932 quota_took_out: "; quota took out ",
933 run_ended: "run ended ",
934 waiting_for_answer: " - waiting for operator answer to question ",
935 recovered_running: "recovered a `running` task whose daemon never recorded the outcome: ",
936 no_run_to_recover: "task was `running` with no live daemon and no readable \
937 run to recover; held for a human to check what happened",
938 could_not_start: "could not start the run: ",
939 resume_not_latest: |q, run| {
940 format!("question {q} asked to resume run {run}, which is not this task's latest run")
941 },
942 resume_cannot_progress: |q, run| {
943 format!("question {q} asked to resume run {run}, which cannot make progress")
944 },
945 resume_unreadable: |q, run, e| {
946 format!("question {q} asked to resume run {run}, which could not be read: {e}")
947 },
948 handover_refused: |r| r.to_string(),
949 handover_hint: " (clean up the other worktree, then release the task from the queue)",
950 reviewers_never_answered: |missing, rounds| {
951 format!(
952 "{missing} reviewer seat(s) never answered after {rounds} rounds; \
953 refusing to call it clean"
954 )
955 },
956};
957
958const PHRASES_JA: Phrases = Phrases {
959 graph_stopped: |status, detail| {
960 format!("グラフが終端状態に達しないまま `{status}` で停止しました: {detail}")
961 },
962 quorum_lost: "審査パネルが定足数を失いました",
963 quota_took_out: "。クォータで脱落: ",
964 run_ended: "run 終了: ",
965 waiting_for_answer: " - オペレーターの回答待ち: 質問 ",
966 recovered_running: "daemon が結果を記録しないまま `running` だったタスクを回収しました: ",
967 no_run_to_recover: "タスクは `running` でしたが、生きた daemon も回収できる run も見つかりません。\
968 何が起きたか人が確認するため保留にしました",
969 could_not_start: "run を開始できませんでした: ",
970 resume_not_latest: |q, run| {
971 format!(
972 "質問 {q} は run {run} の再開を求めましたが、これはタスクの最新の run ではありません"
973 )
974 },
975 resume_cannot_progress: |q, run| {
976 format!("質問 {q} は run {run} の再開を求めましたが、これは進行できません")
977 },
978 resume_unreadable: |q, run, e| {
979 format!("質問 {q} は run {run} の再開を求めましたが、読み込めませんでした: {e}")
980 },
981 handover_refused: |r| {
982 use crate::handover::Refused;
983 match r {
984 Refused::Foreign { branch, path, why } => format!(
985 "ブランチ `{branch}` は {path} にチェックアウトされており、magi は自動では\
986 削除しません。不要なら `git worktree remove` でその worktree を削除して\
987 から、やり直してください(詳細: {why})"
988 ),
989 Refused::Unsafe { branch, path, why } => format!(
990 "ブランチ `{branch}` は {path} にチェックアウトされています。そこの作業を\
991 コミットか破棄したうえで `git worktree remove` で worktree を削除する\
992 か、run を破棄してよいと伝えてから、やり直してください(詳細: {why})"
993 ),
994 Refused::ReleaseFailed { branch, path, run } => format!(
995 "ブランチ `{branch}` は {path} で run {run} がチェックアウトしており、その\
996 worktree の解放に失敗したか、変更が見つかりました(未コミットの変更が\
997 ある worktree は git が削除を拒否します)。worktree はそのまま残しました"
998 ),
999 }
1000 },
1001 handover_hint: "(他の worktree を片付けてから、タスクをキューから解放してください)",
1002 reviewers_never_answered: |missing, rounds| {
1003 format!(
1004 "{missing} 席のレビュアーが {rounds} ラウンドの間に一度も回答しなかったため、\
1005 クリーンとは認めません"
1006 )
1007 },
1008};
1009
1010fn phrases(language: &str) -> &'static Phrases {
1011 if crate::lang::is_japanese(language) {
1012 &PHRASES_JA
1013 } else {
1014 &PHRASES_EN
1015 }
1016}
1017
1018fn language_of(task: &Task, fallback: &Path) -> String {
1022 crate::lang::of_repo(&repo_for(task, fallback))
1023}
1024
1025pub(crate) fn task_of_question<'a>(tasks: &'a [Task], q: &ask::Question) -> Option<&'a Task> {
1029 if q.node == crate::conduct::NODE {
1030 return tasks.iter().find(|t| t.id == q.run);
1031 }
1032 tasks.iter().find(|t| t.runs.contains(&q.run))
1033}
1034
1035fn apply_choice_actions(queue: &Queue, questions: &Questions, home: &Path) {
1044 let tasks = queue.list();
1045 for q in questions.list() {
1046 if q.chosen_action().is_none() {
1047 continue;
1048 }
1049 let Some(listed) = task_of_question(&tasks, &q) else {
1050 continue;
1051 };
1052 if listed.action_applied(&q.id) {
1053 continue;
1054 }
1055 let Ok(_claim) = queue.claim(&listed.id) else {
1056 continue;
1057 };
1058 let Ok(mut task) = queue.get(&listed.id) else {
1059 continue;
1060 };
1061 let language = language_of(&task, Path::new("."));
1062 let decision = decide_action(&task, &q, phrases(&language), |id| {
1063 RunState::load_under(id, home)
1064 });
1065 let ran = matches!(
1066 decision,
1067 ActionDecision::Resume(_) | ActionDecision::Requeue | ActionDecision::Done
1068 );
1069 if ran && asker_may_still_read(questions, &q, home) {
1070 continue;
1079 }
1080 let done = matches!(decision, ActionDecision::Done);
1081 match decision {
1082 ActionDecision::Skip => continue,
1083 ActionDecision::Resume(run) => {
1084 task.release();
1085 task.resume_override = Some(crate::queue::OperatorResume {
1086 question_id: q.id.clone(),
1087 at: Timestamp::now(),
1088 conductor_rehold: None,
1089 forced: true,
1090 pinned_run: Some(run),
1091 });
1092 }
1093 ActionDecision::Stale => {}
1094 ActionDecision::Requeue => task.requeue(),
1095 ActionDecision::Done => task.succeed(),
1096 ActionDecision::Refuse(why) => {
1097 task.hold_machine(Some(why));
1098 notices::raise(
1099 Notice::warn(
1100 &format!("action:{}", q.id),
1101 "An answer asked the daemon to resume a run that cannot be resumed; the task stays held.",
1102 )
1103 .link(Link::Task {
1104 id: task.id.clone(),
1105 }),
1106 );
1107 }
1108 }
1109 task.mark_action_applied(&q.id);
1110 record(queue, &mut task);
1113 if done && queue.get(&task.id).is_ok_and(|t| t.action_applied(&q.id)) {
1114 supersede_prior_runs(&task, home);
1115 }
1116 let _ = questions.update(&q.id, |r| {
1119 r.answer_delivered = true;
1120 Ok(())
1121 });
1122 }
1123}
1124
1125fn asker_may_still_read(questions: &Questions, q: &ask::Question, home: &Path) -> bool {
1129 let now = Timestamp::now();
1130 if questions.read_lease(&q.id).is_some_and(|l| l.fresh(now)) {
1131 return true;
1132 }
1133 q.node != crate::conduct::NODE
1134 && RunState::load_under(&q.run, home).is_ok_and(|s| {
1135 s.seats_active()
1136 .any(|(k, a)| *k == q.seat && a.remaining_secs(now) > 0)
1137 })
1138}
1139
1140fn reconcile_task_questions(queue: &Queue, questions: &Questions) {
1151 let tasks = queue.list();
1152 let by_id: std::collections::BTreeMap<&str, &Task> =
1153 tasks.iter().map(|t| (t.id.as_str(), t)).collect();
1154 let referenced: std::collections::BTreeSet<&str> = tasks
1155 .iter()
1156 .flat_map(|task| task.blocked_by.iter().map(String::as_str))
1157 .collect();
1158
1159 for mut question in questions.list() {
1160 if !question.status.open() || question.node != crate::conduct::NODE {
1161 continue;
1162 }
1163 if referenced.contains(question.id.as_str()) {
1166 continue;
1167 }
1168 let Some(task) = by_id.get(question.run.as_str()) else {
1169 continue;
1170 };
1171 question.abandon(format!(
1172 "task {} no longer waits for this answer",
1173 task.short()
1174 ));
1175 if let Err(e) = questions.put(&mut question) {
1176 tracing::warn!(
1177 "could not retire question {} for task {}: {e:#}",
1178 question.short(),
1179 task.short()
1180 );
1181 }
1182 }
1183}
1184
1185#[derive(Debug, Clone, Copy)]
1192pub struct Verdict {
1193 pub status: RunStatus,
1195 pub left_pr: bool,
1197 pub quota_hit: bool,
1199 pub parked: bool,
1201 pub no_viable_candidates: bool,
1208}
1209
1210pub fn settle(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize) {
1267 settle_in(task, verdict, detail, max_attempts, &PHRASES_EN)
1268}
1269
1270fn settle_in(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize, p: &Phrases) {
1272 if verdict.parked {
1278 task.park(detail);
1279 return;
1280 }
1281 match verdict.status {
1282 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1283 RunStatus::AlreadyInBase => task.already_landed(detail),
1284 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1285 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1286 task.stall(detail)
1287 }
1288 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1289 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1290 RunStatus::Blocked => task.fail(detail, max_attempts),
1291 RunStatus::VerifiedNoop => task.handed_off(detail),
1292 other => task.fail((p.graph_stopped)(label(other), detail), max_attempts),
1293 }
1294}
1295
1296pub fn supersede_prior_runs(task: &Task, home: &Path) {
1345 let now = Timestamp::now();
1346 let last_run_succeeded = task
1353 .runs
1354 .last()
1355 .and_then(|id| RunState::load_under(id, home).ok())
1356 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1357 for id in task.superseded_attempts(last_run_succeeded) {
1358 let mut state = match RunState::load_under(id, home) {
1359 Ok(s) => s,
1360 Err(e) => {
1361 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1362 continue;
1363 }
1364 };
1365 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1366 continue;
1367 }
1368 let daemon_claims = is_working_on(home, id, now);
1369 if state.liveness(daemon_claims) == Liveness::Live {
1370 continue;
1371 }
1372 state.status = RunStatus::Superseded;
1373 if let Err(e) = state.save_under(home) {
1374 tracing::warn!("could not mark run {id} superseded: {e:#}");
1375 }
1376 }
1377}
1378
1379fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1401 for task in queue.list() {
1402 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1403 continue;
1404 }
1405 supersede_prior_runs(&task, home);
1406 }
1407}
1408
1409fn settle_and_diagnose(
1416 task: &mut Task,
1417 verdict: Verdict,
1418 detail: &str,
1419 max_attempts: usize,
1420 state: &RunState,
1421) {
1422 let p = phrases(&state.config.graph.language);
1423 settle_in(task, verdict, detail, max_attempts, p);
1424 if task.status == TaskStatus::Held {
1425 task.diagnostic = diagnostic(state);
1426 note_open_question(task, &state.id, p);
1427 }
1428}
1429
1430fn note_open_question(task: &mut Task, run: &str, p: &Phrases) {
1445 let Some(home) = crate::run::try_home() else {
1446 return;
1447 };
1448 let open = Questions::at(home.join("questions")).open_for(run);
1449 let Some(q) = open.first() else {
1450 return;
1451 };
1452 let base = task.hold_reason.clone().unwrap_or_default();
1453 task.hold_reason = Some(format!("{base}{}{}", p.waiting_for_answer, q.short()));
1454}
1455
1456fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize, language: &str) {
1467 match last_run {
1468 Some(state) => {
1469 let verdict = Verdict {
1470 status: state.status,
1471 left_pr: state.pr.is_some(),
1472 quota_hit: !state.quota.is_empty(),
1473 parked: state.parked,
1474 no_viable_candidates: state.viable().is_empty(),
1475 };
1476 let detail = format!(
1477 "{}{}",
1478 phrases(&state.config.graph.language).recovered_running,
1479 describe(&state)
1480 );
1481 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1482 }
1483 None => {
1484 let why = phrases(language).no_run_to_recover;
1487 task.last_error = Some(why.to_owned());
1488 task.hold_machine(Some(why.to_owned()));
1491 }
1492 }
1493}
1494
1495fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1517 let mut reclaimed = Vec::new();
1518 for listed in queue.list() {
1519 if listed.status != TaskStatus::Running {
1520 continue;
1521 }
1522 let Ok(_claim) = queue.claim(&listed.id) else {
1523 continue;
1524 };
1525 let Ok(mut task) = queue.get(&listed.id) else {
1529 continue;
1530 };
1531 if task.status != TaskStatus::Running {
1532 continue;
1533 }
1534 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1535 if let Some(state) = &last_run
1545 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1546 {
1547 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1548 }
1549 let language = if last_run.is_none() {
1550 language_of(&task, Path::new("."))
1551 } else {
1552 String::new()
1553 };
1554 reclaim(&mut task, last_run, max_attempts, &language);
1555 if task.status == TaskStatus::Done {
1556 supersede_prior_runs(&task, &crate::run::home());
1557 }
1558 record(queue, &mut task);
1559 reclaimed.push(task.id.clone());
1560 }
1561 reclaimed
1562}
1563
1564fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1596 reclaim_abandoned_runs_with(
1597 home,
1598 now,
1599 crate::proc::pid_status,
1600 crate::proc::process_started_at,
1601 )
1602}
1603
1604fn reclaim_abandoned_runs_with<F, G>(
1610 home: &Path,
1611 now: Timestamp,
1612 query: F,
1613 identity: G,
1614) -> Vec<String>
1615where
1616 F: Fn(u32) -> Option<bool>,
1617 G: Fn(u32) -> Option<String>,
1618{
1619 let mut abandoned = Vec::new();
1620 for entry in std::fs::read_dir(home.join("runs"))
1621 .into_iter()
1622 .flatten()
1623 .flatten()
1624 {
1625 let id = entry.file_name().to_string_lossy().into_owned();
1626 if !crate::run::is_run_id(&id) {
1627 continue;
1628 }
1629 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1637 continue;
1638 };
1639 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1640 continue;
1641 };
1642 if state.status.done() || !state.active_all_overrun(now) {
1643 continue;
1644 }
1645 let daemon_claims = is_working_on(home, &id, now);
1657 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1658 continue;
1659 }
1660 state.abandon("daemon");
1661 if let Err(e) = state.save_under(home) {
1662 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1663 continue;
1664 }
1665 if let Some(notice) = notices::run_ended(&state) {
1668 notices::raise_in(home, notice);
1669 }
1670 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1678 tracing::warn!("abandon questions for {id}: {e:#}");
1679 }
1680 abandoned.push(id);
1681 }
1682 abandoned
1683}
1684
1685pub async fn serve(opts: Opts) -> Result<()> {
1691 serve_until(opts, Stop::new()).await
1692}
1693
1694pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1711 let signal = {
1712 let stop = stop.clone();
1713 tokio::spawn(async move {
1714 if tokio::signal::ctrl_c().await.is_ok() {
1715 stop.stop();
1716 tracing::info!("shutdown requested; a run in flight will be finished first");
1717 }
1718 })
1719 };
1720
1721 let worktrees_root = opts
1722 .worktrees_root
1723 .clone()
1724 .unwrap_or_else(crate::run::default_worktree_root);
1725 let outcome = drive(
1726 &opts,
1727 &Queue::open(),
1728 &status_path(),
1729 &crate::run::home(),
1730 &worktrees_root,
1731 &stop,
1732 )
1733 .await;
1734
1735 signal.abort();
1736 outcome
1737}
1738
1739async fn drive(
1752 opts: &Opts,
1753 queue: &Queue,
1754 status_file: &Path,
1755 home: &Path,
1756 worktrees_root: &Path,
1757 stop: &Stop,
1758) -> Result<()> {
1759 let status = Arc::new(Mutex::new(Status::new()));
1767 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1768 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1769
1770 let daemon_cfg = prepare(&opts.repo, opts)
1776 .map(|c| c.daemon)
1777 .unwrap_or_default();
1778 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1779
1780 let waiter = tokio::spawn(crate::waiter::run(
1785 crate::waiter::Waiter::new(
1786 crate::ask::Questions::at(home.join("questions")),
1787 home.to_path_buf(),
1788 prepare(&opts.repo, opts).ok(),
1789 ),
1790 stop.clone(),
1791 ));
1792
1793 let deputies = tokio::spawn(crate::deputy::run(
1797 crate::deputy::Deputies::new(
1798 crate::ask::Questions::at(home.join("questions")),
1799 home.to_path_buf(),
1800 prepare(&opts.repo, opts).ok(),
1801 opts.repo.clone(),
1802 daemon_cfg.max_deputies,
1803 {
1804 let stop = stop.clone();
1805 Arc::new(move || stop.parking())
1806 },
1807 ),
1808 stop.clone(),
1809 ));
1810
1811 let fetcher = tokio::spawn(fetch_loop(opts.repo.clone(), opts.clone(), stop.clone()));
1814
1815 let release_watch = tokio::spawn(crate::release_watch::run(
1818 crate::release_watch::Watcher::new(
1819 Box::new(crate::release_watch::GhForge),
1820 home.to_path_buf(),
1821 ),
1822 {
1823 let (repo, opts) = (opts.repo.clone(), opts.clone());
1824 move || release_watch_settings(&repo, &opts)
1825 },
1826 stop.clone(),
1827 ));
1828
1829 tracing::info!(
1830 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1831 queue.root().display(),
1832 opts.poll.as_secs(),
1833 opts.max_attempts,
1834 concurrency,
1835 if daemon_cfg.pause_for_interrupts {
1836 ", interrupts enabled"
1837 } else {
1838 ""
1839 }
1840 );
1841
1842 janitor(&opts.repo, opts, home, worktrees_root).await;
1845 resweep_superseded_attempts(queue, home);
1846
1847 let outcome = poll(
1848 opts,
1849 queue,
1850 &status,
1851 home,
1852 worktrees_root,
1853 stop,
1854 DispatchLimits {
1855 max_concurrent: concurrency,
1856 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1857 },
1858 )
1859 .await;
1860
1861 beat.abort();
1862 waiter.abort();
1863 deputies.abort();
1864 fetcher.abort();
1865 release_watch.abort();
1866 clear_status_at(status_file);
1867 outcome
1868}
1869
1870fn release_watch_settings(repo: &Path, opts: &Opts) -> (Vec<PathBuf>, u64) {
1873 let Ok(cfg) = prepare(repo, opts) else {
1874 return (Vec::new(), 0);
1875 };
1876 let mut paths: Vec<PathBuf> = crate::repos::scan(&cfg.repos.roots)
1877 .into_iter()
1878 .map(|r| r.path)
1879 .collect();
1880 let mut known: Vec<PathBuf> = Queue::open().list().into_iter().map(|t| t.repo).collect();
1883 let runs = crate::run::home().join("runs");
1884 for id in crate::run::list_ids_in(&runs) {
1885 let repo = std::fs::read_to_string(runs.join(&id).join("run.json"))
1886 .ok()
1887 .and_then(|s| serde_json::from_str::<serde_json::Value>(&s).ok())
1888 .and_then(|v| v.get("repo")?.as_str().map(PathBuf::from));
1889 known.extend(repo);
1890 }
1891 for k in known {
1892 let k = k.canonicalize().unwrap_or(k);
1893 if k.join(".git").exists() && !paths.contains(&k) {
1894 paths.push(k);
1895 }
1896 }
1897 let own = repo.canonicalize().unwrap_or_else(|_| repo.to_path_buf());
1898 if !paths.contains(&own) {
1899 paths.push(own);
1900 }
1901 (paths, cfg.daemon.release_stall_minutes)
1902}
1903
1904const FETCH_TIMEOUT: Duration = Duration::from_secs(30);
1906
1907async fn fetch_loop(repo: PathBuf, opts: Opts, stop: Stop) {
1911 while !stop.stopped() {
1912 let (interval, roots) = match prepare(&repo, &opts) {
1913 Ok(c) => (c.repos.fetch_interval, c.repos.roots),
1914 Err(_) => (0, Vec::new()),
1915 };
1916 if interval > 0 && !roots.is_empty() {
1917 let r = crate::clean::fetch_origins(&roots, FETCH_TIMEOUT, || stop.stopped()).await;
1918 tracing::debug!("fetch origins: {r:?}");
1919 }
1920 let wait = if interval > 0 { interval } else { 60 };
1922 let mut slept = 0;
1923 while slept < wait && !stop.stopped() {
1924 tokio::time::sleep(Duration::from_secs(1)).await;
1925 slept += 1;
1926 }
1927 }
1928}
1929
1930async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1936 loop {
1937 tokio::time::sleep(HEARTBEAT).await;
1938 let snapshot = {
1939 let mut guard = lock(&status);
1940 guard.updated_at = Timestamp::now();
1941 guard.clone()
1942 };
1943 if let Err(e) = write_status_to(&path, &snapshot) {
1944 tracing::warn!("could not refresh the daemon status file: {e:#}");
1947 }
1948 }
1949}
1950
1951#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1954enum LandResume {
1955 NotLanding,
1958 StillWaiting,
1963 Ready,
1967}
1968
1969fn land_resume_state(task: &Task) -> LandResume {
1973 let Some(run_id) = task.runs.last() else {
1974 return LandResume::NotLanding;
1975 };
1976 let Ok(state) = RunState::load(run_id) else {
1977 return LandResume::NotLanding;
1978 };
1979 if state.status != RunStatus::Landing || !state.parked {
1980 return LandResume::NotLanding;
1981 }
1982 let store = ask::Questions::open();
1983 let waiting = store
1984 .list()
1985 .into_iter()
1986 .filter(|q| &q.run == run_id && q.node == land::APPROVAL_NODE)
1987 .max_by(|a, b| a.id.cmp(&b.id));
1988 let Some(mut q) = waiting else {
1989 return LandResume::Ready;
1990 };
1991 if !q.status.open() {
1992 return LandResume::Ready;
1993 }
1994 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
2001 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
2002 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
2003 q.abandon(format!(
2004 "no answer within {}s of asking",
2005 timeout.as_secs().max(1)
2006 ));
2007 if store.put(&mut q).is_ok() {
2010 return LandResume::Ready;
2011 }
2012 }
2013 LandResume::StillWaiting
2014}
2015
2016const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
2024
2025const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
2037
2038struct InFlightGuard<'a> {
2051 status: &'a Arc<Mutex<Status>>,
2052 stop: &'a Stop,
2053 task_id: &'a str,
2054}
2055
2056impl Drop for InFlightGuard<'_> {
2057 fn drop(&mut self) {
2058 lock(self.status).current.retain(|c| c.task != self.task_id);
2059 self.stop.exit();
2060 }
2061}
2062
2063#[derive(Debug, Clone, PartialEq, Eq)]
2085enum Interrupt {
2086 Idle,
2088 Parking {
2100 parked: Vec<String>,
2101 interrupt_task: String,
2102 },
2103 Running {
2111 parked: Vec<String>,
2112 interrupt_task: String,
2113 },
2114 Resuming { parked: Vec<String> },
2121}
2122
2123fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
2144 match state {
2145 Interrupt::Idle => {
2146 if in_flight.len() != 1 {
2158 return Interrupt::Idle;
2159 }
2160 match runnable.iter().find(|t| t.interrupt) {
2161 Some(t) => Interrupt::Parking {
2162 parked: in_flight.to_vec(),
2163 interrupt_task: t.id.clone(),
2164 },
2165 None => Interrupt::Idle,
2166 }
2167 }
2168 Interrupt::Parking {
2169 parked,
2170 interrupt_task,
2171 } => {
2172 if in_flight.iter().any(|id| parked.contains(id)) {
2173 Interrupt::Parking {
2175 parked,
2176 interrupt_task,
2177 }
2178 } else if in_flight.contains(&interrupt_task) {
2179 Interrupt::Running {
2180 parked,
2181 interrupt_task,
2182 }
2183 } else if runnable.iter().any(|t| t.id == interrupt_task) {
2184 Interrupt::Parking {
2188 parked,
2189 interrupt_task,
2190 }
2191 } else {
2192 Interrupt::Resuming { parked }
2197 }
2198 }
2199 Interrupt::Running {
2200 parked,
2201 interrupt_task,
2202 } => {
2203 if in_flight.contains(&interrupt_task) {
2204 Interrupt::Running {
2205 parked,
2206 interrupt_task,
2207 }
2208 } else {
2209 Interrupt::Resuming { parked }
2215 }
2216 }
2217 Interrupt::Resuming { parked } => {
2218 if in_flight.iter().any(|id| parked.contains(id)) {
2219 Interrupt::Idle
2225 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
2226 Interrupt::Resuming { parked }
2227 } else {
2228 Interrupt::Idle
2231 }
2232 }
2233 }
2234}
2235
2236fn advance_interrupt_tick(
2242 enabled: bool,
2243 state: Interrupt,
2244 in_flight: &[String],
2245 runnable: &[Task],
2246) -> Interrupt {
2247 if !enabled {
2248 return Interrupt::Idle;
2249 }
2250 advance_interrupt(state, in_flight, runnable)
2251}
2252
2253fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
2258 match state {
2259 Interrupt::Idle => candidates,
2260 Interrupt::Parking {
2261 parked,
2262 interrupt_task,
2263 } => {
2264 if in_flight.iter().any(|id| parked.contains(id)) {
2265 Vec::new()
2266 } else {
2267 candidates
2268 .into_iter()
2269 .filter(|t| &t.id == interrupt_task)
2270 .collect()
2271 }
2272 }
2273 Interrupt::Running { .. } => Vec::new(),
2274 Interrupt::Resuming { parked } => candidates
2282 .into_iter()
2283 .find(|t| parked.contains(&t.id))
2284 .into_iter()
2285 .collect(),
2286 }
2287}
2288
2289struct DispatchLimits {
2293 max_concurrent: usize,
2296 pause_for_interrupts: bool,
2298}
2299
2300#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2309enum PermitKind {
2310 None,
2315 Urgent,
2322 Ordinary,
2325}
2326
2327fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
2330 if priority {
2331 PermitKind::None
2332 } else if urgent {
2333 PermitKind::Urgent
2334 } else {
2335 PermitKind::Ordinary
2336 }
2337}
2338
2339async fn poll(
2358 opts: &Opts,
2359 queue: &Queue,
2360 status: &Arc<Mutex<Status>>,
2361 home: &Path,
2362 worktrees_root: &Path,
2363 stop: &Stop,
2364 limits: DispatchLimits,
2365) -> Result<()> {
2366 let DispatchLimits {
2367 max_concurrent,
2368 pause_for_interrupts,
2369 } = limits;
2370 let mut attempted: Vec<String> = Vec::new();
2375 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2376 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2384 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2390 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2391 let mut conductor = Conductor::new();
2392 let mut cache_last_checked: Option<Timestamp> = None;
2395 let mut interrupt = Interrupt::Idle;
2397 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2403 std::collections::HashMap::new();
2404
2405 while !stop.stopped() {
2406 lock(status).polls += 1;
2407
2408 while let Some(result) = inflight.try_join_next() {
2413 if let Err(e) = result {
2414 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2415 notices::raise(Notice::error(
2416 "loop:attempt",
2417 "A queued attempt ended abnormally; check the task it was running.",
2418 ));
2419 }
2420 }
2421
2422 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2423 if !swept.is_empty() {
2424 tracing::warn!(
2425 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2426 swept.len(),
2427 swept.join(", ")
2428 );
2429 }
2430 let now = Timestamp::now();
2435
2436 if !stop.busy_now() {
2441 maybe_prune_cache_between_runs(
2442 &opts.repo,
2443 opts,
2444 home,
2445 stop,
2446 &mut cache_last_checked,
2447 now,
2448 )
2449 .await;
2450 }
2451
2452 let stalled = stalled_tasks(queue, home, now);
2453 let stalled_ids: std::collections::BTreeSet<_> =
2454 stalled.iter().map(|task| task.id.clone()).collect();
2455 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2456 if !reclaimed.is_empty() {
2457 tracing::warn!(
2458 "reclaimed {} task(s) left `running` by a daemon that never \
2459 recorded the outcome: {}",
2460 reclaimed.len(),
2461 reclaimed.join(", ")
2462 );
2463 }
2464 let abandoned_runs = reclaim_abandoned_runs(home, now);
2465 if !abandoned_runs.is_empty() {
2466 tracing::warn!(
2467 "failed {} run(s) left behind by a killed process, past every \
2468 active seat's own timeout: {}",
2469 abandoned_runs.len(),
2470 abandoned_runs.join(", ")
2471 );
2472 }
2473
2474 let questions = Questions::at(home.join("questions"));
2479
2480 resolve_blockers(queue, &questions);
2483 apply_choice_actions(queue, &questions, home);
2484 reconcile_task_questions(queue, &questions);
2485
2486 let finished: Vec<Task> = finished_tasks(queue)
2492 .into_iter()
2493 .filter(|task| !stalled_ids.contains(&task.id))
2494 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2498 .collect();
2499 let queued = queued_tasks(queue);
2500 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2504 && conductor.worth_a_look(queue, &stalled, &finished)
2505 {
2506 match prepare(&opts.repo, opts) {
2507 Ok(cfg) => {
2508 conductor
2509 .maybe_run(
2510 &cfg,
2511 &opts.repo,
2512 queue,
2513 &questions,
2514 home,
2515 &queued,
2516 &stalled,
2517 &finished,
2518 opts.max_attempts,
2519 )
2520 .await;
2521 }
2522 Err(e) => {
2523 tracing::warn!("conductor: no config: {e:#}");
2524 notices::raise(Notice::warn(
2525 "loop:no-config",
2526 "The loop could not read this repository's config, so held tasks are not being triaged.",
2527 ));
2528 }
2529 }
2530 }
2531
2532 let candidates: Vec<Task> = runnable(queue)
2533 .into_iter()
2534 .filter(|t| !opts.once || !attempted.contains(&t.id))
2535 .collect();
2536
2537 let in_flight: Vec<String> = lock(status)
2542 .current
2543 .iter()
2544 .map(|c| c.task.clone())
2545 .collect();
2546 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2547
2548 interrupt =
2549 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2550 if let Interrupt::Parking {
2551 parked,
2552 interrupt_task,
2553 } = &interrupt
2554 {
2555 let reason = format!(
2556 "task {} asked to run first",
2557 crate::run::short_of(interrupt_task)
2558 );
2559 for id in parked {
2560 if let Some(pause) = interrupt_pauses.get(id) {
2561 pause.park_because(reason.clone());
2562 }
2563 }
2564 }
2565 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2566
2567 let cooling_down =
2568 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2569
2570 let mut started_any = false;
2571 for candidate in candidates {
2572 if stop.stopped() {
2573 break;
2574 }
2575
2576 let resume = land_resume_state(&candidate);
2577 if resume == LandResume::StillWaiting {
2578 continue;
2579 }
2580 let priority = resume == LandResume::Ready;
2581
2582 if !priority && cooling_down {
2583 continue;
2584 }
2585 let permit = match permit_kind(priority, candidate.urgent) {
2586 PermitKind::None => None,
2587 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2588 Ok(p) => Some(p),
2589 Err(_) => continue,
2595 },
2596 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2597 Ok(p) => Some(p),
2598 Err(_) => continue,
2602 },
2603 };
2604
2605 let Ok(claim) = queue.claim(&candidate.id) else {
2610 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2611 continue;
2612 };
2613 let mut task = match queue.get(&candidate.id) {
2616 Ok(t) if t.status.runnable() => t,
2617 Ok(_) => continue,
2618 Err(e) => {
2619 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2620 continue;
2621 }
2622 };
2623 let task_id = task.id.clone();
2624 attempted.push(task_id.clone());
2625 lock(status).idle = false;
2626 stop.enter();
2629 started_any = true;
2630
2631 let run_pause = crate::graph::Pause::new();
2635 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2636
2637 let opts = opts.clone();
2638 let queue = queue.clone();
2639 let status = Arc::clone(status);
2640 let stop = stop.clone();
2641 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2642 inflight.spawn(async move {
2643 let _claim = claim;
2647 let _permit = permit;
2648 let _inflight = InFlightGuard {
2650 status: &status,
2651 stop: &stop,
2652 task_id: &task_id,
2653 };
2654 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2655 lock(&status).completed += 1;
2656 let now = Timestamp::now();
2662 if let Some(until) = cooldown_until("a, now) {
2663 let wait = until.as_second() - now.as_second();
2664 *lock("a_cooldown_until) = Some(until);
2665 let hint = quota
2666 .iter()
2667 .find(|q| q.reset.is_some())
2668 .and_then(|q| q.reset.as_deref());
2669 match hint {
2670 Some(h) => tracing::warn!(
2671 "quota hit; waiting {wait}s before taking another ordinary task \
2672 (CLI reported reset: {h})"
2673 ),
2674 None => tracing::warn!(
2675 "quota hit; waiting {wait}s before taking another ordinary task \
2676 (no reset hint reported)"
2677 ),
2678 }
2679 }
2680 });
2681 }
2682
2683 if started_any {
2684 continue;
2685 }
2686
2687 if stop.busy_now() {
2688 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2693 continue;
2694 }
2695
2696 lock(status).idle = true;
2698 if opts.once {
2699 janitor(&opts.repo, opts, home, worktrees_root).await;
2703 resweep_superseded_attempts(queue, home);
2704 triage_held(queue, home, opts).await;
2705 break;
2706 }
2707 stop.idle(opts.poll).await;
2708 if stop.stopped() {
2709 continue;
2710 }
2711 janitor(&opts.repo, opts, home, worktrees_root).await;
2717 resweep_superseded_attempts(queue, home);
2718 triage_held(queue, home, opts).await;
2719 }
2720
2721 while let Some(result) = inflight.join_next().await {
2726 if let Err(e) = result {
2727 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2728 notices::raise(Notice::error(
2729 "loop:attempt",
2730 "A queued attempt ended abnormally; check the task it was running.",
2731 ));
2732 }
2733 }
2734 Ok(())
2735}
2736
2737async fn attempt(
2743 opts: &Opts,
2744 queue: &Queue,
2745 status: &Arc<Mutex<Status>>,
2746 stop: &Stop,
2747 interrupt_pause: crate::graph::Pause,
2748 task: &mut Task,
2749) -> Vec<QuotaLoss> {
2750 let repo = repo_for(task, &opts.repo);
2751 tracing::info!(
2752 "task {} — {} (repo {})",
2753 task.short(),
2754 task.title,
2755 repo.display()
2756 );
2757
2758 let explicit = match &task.overrides {
2762 Some(o) => o.config.as_deref(),
2763 None => opts.config.as_deref(),
2764 };
2765 if let Err(e) = Config::discover_fetched(&repo, explicit).await {
2766 let reason = format!("[config] {e:#}");
2767 task.last_error = Some(reason.clone());
2768 task.hold_machine(Some(reason.clone()));
2769 record(queue, task);
2770 tracing::warn!("holding {}: {reason}", task.short());
2771 return Vec::new();
2772 }
2773 let mut config = match prepare_for(&repo, opts, task) {
2774 Ok(c) => c,
2775 Err(e) => {
2776 task.attempts += 1;
2780 task.fail(format!("config: {e:#}"), opts.max_attempts);
2781 record(queue, task);
2782 return Vec::new();
2783 }
2784 };
2785 apply_solo(&mut config, task);
2786 let p = phrases(&config.graph.language);
2787 let start_failed = p.could_not_start;
2788
2789 if let Some(reason) = disk_gate(&repo, &config) {
2797 task.last_error = Some(reason.clone());
2798 task.hold_machine(Some(reason.clone()));
2799 record(queue, task);
2800 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2801 notices::raise(
2804 Notice::warn(
2805 &format!("disk:{}", repo.display()),
2806 "A task was held for want of free disk space; free some, then release it from the queue.",
2807 )
2808 .link(Link::Task {
2809 id: task.id.clone(),
2810 }),
2811 );
2812 return Vec::new();
2813 }
2814
2815 let unfinished = (!task.fresh_start)
2835 .then(|| unfinished_run(&task.runs, task.short()))
2836 .flatten();
2837 let review_branch = task.review_branch.take();
2843 let branch_exists = match &review_branch {
2844 Some(branch) => crate::git::branch_exists(&repo, branch)
2845 .await
2846 .unwrap_or(false),
2847 None => false,
2848 };
2849 let starter = match (&review_branch, &unfinished, &task.review_of) {
2853 (None, None, Some(branch)) => Starter::Review(branch.clone()),
2854 _ => choose_starter(
2855 review_branch.as_deref(),
2856 branch_exists,
2857 unfinished.as_deref(),
2858 ),
2859 };
2860 let attachments = match task_attachments(queue, task) {
2865 Ok(a) => a,
2866 Err(e) => {
2867 task.attempts += 1;
2868 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2869 record(queue, task);
2870 return Vec::new();
2871 }
2872 };
2873 let started = match &starter {
2874 Starter::Review(branch) => {
2875 tracing::info!(
2876 "task {} reopens `{branch}` as a review-only pass",
2877 task.short()
2878 );
2879 let takeover = crate::handover::Takeover {
2882 earlier: task.earlier_attempts().to_vec(),
2883 home: crate::run::home(),
2884 choice: take_divergence_answer(branch, &config.merge.remote, task),
2885 };
2886 Runner::review_taking_over(
2887 &repo,
2888 branch,
2889 config,
2890 Some(takeover),
2891 crate::run::Origin::queue(&task.id),
2892 )
2893 .await
2894 }
2895 Starter::Resume(id) => {
2896 tracing::info!("resuming run {id} rather than competing again");
2897 Runner::resume(id).map(|mut r| {
2898 if let Some(instruction) =
2899 prepare_instruction(&starter, Some(&r.state.instruction), task)
2900 {
2901 r.state.instruction = instruction;
2902 }
2903 r.state.attachments = attachments.clone();
2904 if let Some(mode) = task
2907 .overrides
2908 .as_ref()
2909 .and_then(|o| o.merge.as_deref())
2910 .and_then(|m| merge_mode(m).ok())
2911 {
2912 r.state.config.merge.mode = mode;
2913 }
2914 r
2915 })
2916 }
2917 Starter::Start => {
2918 if let Some(branch) = &review_branch {
2919 tracing::warn!(
2920 "conductor chose review for task {} but branch `{branch}` no longer \
2921 exists; requeuing as a fresh competition instead",
2922 task.short()
2923 );
2924 }
2925 let instruction = prepare_instruction(&starter, None, task)
2926 .unwrap_or_else(|| task.instruction.clone());
2927 Runner::start_naming(
2928 &repo,
2929 instruction,
2930 &task.title,
2931 config,
2932 crate::run::Origin::queue(&task.id),
2933 )
2934 .await
2935 .map(|mut r| {
2936 r.state.attachments = attachments.clone();
2937 r
2938 })
2939 }
2940 };
2941 let mut runner = match started {
2942 Ok(r) => r,
2943 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2947 let detail = match e.downcast_ref::<crate::handover::Refused>() {
2948 Some(r) => (p.handover_refused)(r),
2949 None => format!("{e:#}"),
2950 };
2951 let reason = format!("{start_failed}{detail}");
2952 task.last_error = Some(reason.clone());
2953 let branch = match &starter {
2956 Starter::Review(branch) => Some(branch.clone()),
2957 _ => None,
2958 };
2959 task.hold_for_handover(branch, format!("{reason}{}", p.handover_hint));
2963 record(queue, task);
2964 tracing::warn!(
2965 "holding {} for a branch it cannot take over: {e:#}",
2966 task.short()
2967 );
2968 return Vec::new();
2969 }
2970 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2975 let d = e
2976 .downcast_ref::<crate::reconcile::Diverged>()
2977 .expect("checked by the guard");
2978 let mut q = ask::Question::new(
2979 task.id.clone(),
2980 crate::reconcile::NODE.to_owned(),
2981 crate::reconcile::SEAT.to_owned(),
2982 d.summary(),
2983 d.detail(),
2984 d.choices(),
2985 );
2986 match Questions::open().put(&mut q) {
2987 Ok(()) => {
2988 task.last_error = Some(format!("{e:#}"));
2989 task.review_branch = Some(d.branch.clone());
2992 task.block(vec![q.id.clone()], Some(d.summary()));
2993 }
2994 Err(put) => {
2995 tracing::warn!("could not file the divergence question: {put:#}");
2996 task.attempts += 1;
2997 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2998 }
2999 }
3000 record(queue, task);
3001 return Vec::new();
3002 }
3003 Err(e) => {
3004 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
3010 && let Starter::Review(branch) = &starter
3011 {
3012 task.review_branch = Some(branch.clone());
3013 task.last_error = Some(format!("{start_failed}{e:#}"));
3014 task.status = crate::queue::TaskStatus::Failed;
3015 record(queue, task);
3016 return Vec::new();
3017 }
3018 task.attempts += 1;
3019 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
3020 record(queue, task);
3021 return Vec::new();
3022 }
3023 };
3024 runner.state.followup_generation = Some(task.followup.as_ref().map_or(0, |f| f.generation));
3027 if runner.state.origin_chat.is_none() {
3028 runner.state.origin_chat = task.chat_talk().map(str::to_owned);
3029 }
3030 runner.on_pause(stop.pause());
3032 runner.watch_interrupt(interrupt_pause);
3036
3037 let run = runner.state.id.clone();
3040 task.start(run.clone());
3041 record(queue, task);
3042 lock(status).current.push(Current {
3043 task: task.id.clone(),
3044 run,
3045 });
3046
3047 let quota_before = runner.state.quota.clone();
3050 let result = runner.execute().await;
3051 finish_attempt(
3052 opts.max_attempts,
3053 queue,
3054 task,
3055 &runner.state,
3056 "a_before,
3057 result,
3058 )
3059}
3060
3061pub fn finish_attempt(
3066 max_attempts: usize,
3067 queue: &Queue,
3068 task: &mut Task,
3069 state: &RunState,
3070 quota_before: &[QuotaLoss],
3071 result: Result<()>,
3072) -> Vec<QuotaLoss> {
3073 let detail = match result {
3074 Ok(()) => describe(state),
3075 Err(e) => format!("{e:#}"),
3076 };
3077 let fresh = losses_this_attempt(quota_before, &state.quota);
3078 let verdict = Verdict {
3079 status: state.status,
3080 left_pr: state.pr.is_some(),
3083 quota_hit: !fresh.is_empty(),
3089 parked: state.parked,
3093 no_viable_candidates: state.viable().is_empty(),
3096 };
3097 settle_and_diagnose(task, verdict, &detail, max_attempts, state);
3098 if task.status == TaskStatus::Done {
3099 supersede_prior_runs(task, &crate::run::home());
3100 }
3101 record(queue, task);
3102 tracing::info!(
3103 "task {} is {} after run {} ({})",
3104 task.short(),
3105 task.status.as_str(),
3106 state.short(),
3107 label(state.status)
3108 );
3109 fresh
3110}
3111
3112pub fn hold_if_runnable(queue: &Queue, task: &mut Task) {
3116 if task.status.runnable() {
3117 let why = task.last_error.clone().map_or_else(
3118 || "the run did not finish".to_owned(),
3119 |e| format!("the run did not finish: {e}"),
3120 );
3121 task.hold_manual(Some(format!(
3122 "{why}. It was started by hand, so it is not retried \
3123 automatically; `magi task release` retries it."
3124 )));
3125 record(queue, task);
3126 }
3127}
3128
3129#[must_use]
3133pub fn loop_would_resume(run: &str) -> bool {
3134 unfinished_run(&[run.to_owned()], crate::run::short_of(run)).is_some()
3135}
3136
3137#[must_use]
3147pub fn foreign_loop(
3148 reading: Option<&Reading>,
3149 now: Timestamp,
3150 own_pid: u32,
3151) -> Option<Option<u32>> {
3152 let reading = reading.filter(|r| r.running(now))?;
3153 match reading.pid {
3154 Some(pid) if pid == own_pid => None,
3155 pid => Some(pid),
3156 }
3157}
3158
3159pub async fn run_claimed(opts: &Opts, queue: &Queue, task: &mut Task) {
3174 let status = Arc::new(Mutex::new(Status::new()));
3175 let stop = Stop::new();
3176 attempt(
3177 opts,
3178 queue,
3179 &status,
3180 &stop,
3181 crate::graph::Pause::new(),
3182 task,
3183 )
3184 .await;
3185 hold_if_runnable(queue, task);
3186}
3187
3188fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
3198 after
3199 .iter()
3200 .filter(|q| !before.contains(q))
3201 .cloned()
3202 .collect()
3203}
3204
3205fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
3208 if quota.is_empty() {
3209 return None;
3210 }
3211 let with_hint = quota.iter().find(|q| q.reset.is_some());
3212 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
3213 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
3214 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
3215 Some(
3216 now.checked_add(jiff::SignedDuration::from_secs(secs))
3217 .unwrap_or(Timestamp::MAX),
3218 )
3219}
3220
3221fn apply_solo(config: &mut Config, task: &Task) {
3231 if task.solo {
3232 config.graph.implementers = 1;
3233 }
3234}
3235
3236fn prepare_for(repo: &Path, opts: &Opts, task: &Task) -> Result<Config> {
3241 let Some(o) = &task.overrides else {
3242 return prepare(repo, opts);
3243 };
3244 let own = Opts {
3248 config: o.config.clone(),
3249 merge: None,
3250 ..opts.clone()
3251 };
3252 let mut config = prepare(repo, &own)?;
3253 o.apply(&mut config);
3254 Ok(config)
3255}
3256
3257fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
3259 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
3260 if let Some(mode) = &opts.merge {
3261 config.merge.mode = merge_mode(mode)?;
3262 }
3263 Ok(config)
3264}
3265
3266async fn maybe_prune_cache_between_runs(
3297 repo: &Path,
3298 opts: &Opts,
3299 home: &Path,
3300 stop: &Stop,
3301 last_checked: &mut Option<Timestamp>,
3302 now: Timestamp,
3303) {
3304 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
3305 return;
3306 }
3307 *last_checked = Some(now);
3308 let cfg = match prepare(repo, opts) {
3309 Ok(cfg) => cfg,
3310 Err(e) => {
3311 tracing::warn!("cache check: no config: {e:#}");
3312 return;
3313 }
3314 };
3315 match clean::prune_cache_if_over_limit(&cfg, home) {
3316 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
3317 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
3318 pruned.files,
3319 pruned.freed
3320 ),
3321 Ok(_) => {}
3322 Err(e) => {
3323 tracing::warn!("housekeep: prune cache: {e:#}");
3324 notices::raise_in(
3325 home,
3326 Notice::warn(
3327 "housekeep:cache",
3328 "Pruning the shared build cache failed; disk usage may keep growing.",
3329 ),
3330 );
3331 }
3332 }
3333}
3334
3335fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
3339 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
3340}
3341
3342async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
3361 let cfg = match prepare(repo, opts) {
3362 Ok(cfg) => cfg,
3363 Err(e) => {
3364 tracing::warn!("housekeep: no config: {e:#}");
3365 return;
3366 }
3367 };
3368 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
3377 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
3378 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
3383 let mut extra = Vec::new();
3384 if out.unreadable > 0 {
3385 extra.push(format!("{} unreadable", out.unreadable));
3386 }
3387 if out.orphaned_worktrees > 0 {
3388 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
3389 }
3390 let detail = if extra.is_empty() {
3391 String::new()
3392 } else {
3393 format!(" ({})", extra.join(", "))
3394 };
3395 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
3396 }
3397 if out.external_merges_recorded > 0 {
3398 tracing::info!(
3399 "housekeep: recorded {} run(s) as merged externally",
3400 out.external_merges_recorded
3401 );
3402 }
3403 if out.stale_pr_states_repaired > 0 {
3404 tracing::info!(
3405 "housekeep: rewrote {} run record(s) whose pull request had already settled",
3406 out.stale_pr_states_repaired
3407 );
3408 }
3409 if out.cache_files > 0 {
3410 tracing::info!(
3411 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
3412 out.cache_files,
3413 out.cache_freed
3414 );
3415 }
3416 if out.questions_abandoned > 0 {
3417 tracing::info!(
3418 "housekeep: abandoned {} question(s) left open by a finished run",
3419 out.questions_abandoned
3420 );
3421 }
3422}
3423
3424async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3433 let questions = Questions::at(home.join("questions"));
3434 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3435 if report.is_empty() {
3436 return;
3437 }
3438 if !report.quarantined.is_empty() {
3439 tracing::info!(
3440 "triage: held {} blocked task(s) whose blocked-on task or \
3441 question no longer exists: {}",
3442 report.quarantined.len(),
3443 report.quarantined.join(", ")
3444 );
3445 }
3446 if !report.resumed.is_empty() {
3447 tracing::info!(
3448 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3449 report.resumed.len(),
3450 report.resumed.join(", ")
3451 );
3452 }
3453 if !report.asked.is_empty() {
3454 tracing::info!(
3455 "triage: asked about {} held task(s): {}",
3456 report.asked.len(),
3457 report.asked.join(", ")
3458 );
3459 }
3460 if !report.answered.is_empty() {
3461 tracing::info!(
3462 "triage: applied {} operator answer(s): {}",
3463 report.answered.len(),
3464 report.answered.join(", ")
3465 );
3466 }
3467}
3468
3469fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3476 disk_gate_with(repo, config, crate::disk::free_bytes)
3477}
3478
3479fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3483 repo: &Path,
3484 config: &Config,
3485 free_bytes: F,
3486) -> Option<String> {
3487 let min = config.disk.min_free_bytes;
3488 if min == 0 {
3489 return None;
3490 }
3491 match free_bytes(repo) {
3492 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3493 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3494 }
3495}
3496
3497const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3504
3505const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3509
3510fn quota_wait(
3519 reset_at: Option<Timestamp>,
3520 now: Timestamp,
3521 fallback: Duration,
3522 cap: Duration,
3523) -> Duration {
3524 match reset_at {
3525 Some(at) if at > now => {
3526 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3527 Duration::from_secs(secs).min(cap)
3528 }
3529 _ => fallback,
3530 }
3531}
3532
3533fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3547 parse_reset_hint_zoned(text, now)
3548 .or_else(|| parse_reset_hint_dated(text))
3549 .or_else(|| parse_reset_hint_relative(text, recorded))
3550}
3551
3552fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3556 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3557 let mut rest = rest.trim();
3558 if rest.is_empty() {
3559 return None;
3560 }
3561 let mut total: i64 = 0;
3562 let mut matched = false;
3563 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3564 if let Some((digits, tail)) = rest.split_once(unit)
3565 && !digits.is_empty()
3566 && digits.bytes().all(|b| b.is_ascii_digit())
3567 {
3568 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3569 rest = tail;
3570 matched = true;
3571 }
3572 }
3573 if !rest.is_empty() || !matched {
3574 return None;
3575 }
3576 recorded
3577 .checked_add(jiff::SignedDuration::from_secs(total))
3578 .ok()
3579}
3580
3581fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3585 let clock = clock.trim().to_lowercase();
3586 let (digits, pm) = clock
3587 .strip_suffix("am")
3588 .map(|d| (d, false))
3589 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3590 let (h, m) = digits.trim().split_once(':')?;
3591 let mut hour: i8 = h.trim().parse().ok()?;
3592 let minute: i8 = m.trim().parse().ok()?;
3593 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3594 return None;
3595 }
3596 if pm && hour != 12 {
3597 hour += 12;
3598 } else if !pm && hour == 12 {
3599 hour = 0;
3600 }
3601 Some((hour, minute))
3602}
3603
3604fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3609 let open = text.find('(')?;
3610 let close = text.rfind(')')?;
3611 if close <= open {
3612 return None;
3613 }
3614 let zone = text[open + 1..close].trim();
3615 let (hour, minute) = parse_12h_clock(&text[..open])?;
3616 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3617 let candidate = now
3618 .to_zoned(tz)
3619 .with()
3620 .hour(hour)
3621 .minute(minute)
3622 .second(0)
3623 .millisecond(0)
3624 .microsecond(0)
3625 .nanosecond(0)
3626 .build()
3627 .ok()?;
3628 let mut at = candidate.timestamp();
3629 if at <= now {
3630 at += jiff::SignedDuration::from_hours(24);
3631 }
3632 Some(at)
3633}
3634
3635fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3644 let words: Vec<&str> = text.split_whitespace().collect();
3645 if words.len() < 5 {
3646 return None;
3647 }
3648 (0..=words.len() - 5)
3649 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3650}
3651
3652fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3658 if trailing.is_some_and(|next| next.starts_with('(')) {
3659 return None;
3660 }
3661 let month = month_number(window[0])?;
3662 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3663 let day_digits = ["st", "nd", "rd", "th"]
3664 .iter()
3665 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3666 let day: i8 = day_digits.parse().ok()?;
3667 let year_token = window[2];
3668 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3669 return None;
3670 }
3671 let year: i16 = year_token.parse().ok()?;
3672 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3676 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3677 let date = jiff::civil::Date::new(year, month, day).ok()?;
3678 let candidate = date
3679 .at(hour, minute, 0, 0)
3680 .to_zoned(jiff::tz::TimeZone::UTC)
3681 .ok()?;
3682 Some(candidate.timestamp())
3683}
3684
3685fn month_number(name: &str) -> Option<i8> {
3688 const NAMES: [&str; 12] = [
3689 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3690 ];
3691 let lower = name.to_lowercase();
3692 NAMES
3693 .iter()
3694 .position(|n| *n == lower.as_str())
3695 .map(|i| i as i8 + 1)
3696}
3697
3698fn exhausted_review_budget(state: &RunState) -> bool {
3710 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3711}
3712
3713fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3750 unfinished_run_with(runs, short, RunState::load)
3751}
3752
3753fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3756where
3757 F: FnOnce(&str) -> Result<RunState>,
3758{
3759 let id = runs.last()?;
3760 match load(id) {
3761 Ok(s)
3766 if s.status.resumable()
3767 && !s.released()
3768 && !exhausted_review_budget(&s)
3769 && s.liveness(false) != crate::run::Liveness::Live =>
3770 {
3771 Some(id.clone())
3772 }
3773 Ok(_) => None,
3774 Err(e) => {
3775 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3776 None
3777 }
3778 }
3779}
3780
3781fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3789where
3790 F: FnOnce(&str) -> Result<RunState>,
3791{
3792 if !matches!(task.status, TaskStatus::Parked | TaskStatus::Failed)
3794 || task.fresh_start
3795 || task.review_branch.is_some()
3796 {
3797 return false;
3798 }
3799 let Some(id) = task.runs.last() else {
3800 return false;
3801 };
3802 load(id).is_ok_and(|s| {
3803 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3804 })
3805}
3806
3807#[derive(Debug, Clone, PartialEq, Eq)]
3810enum Starter {
3811 Review(String),
3814 Resume(String),
3816 Start,
3818}
3819
3820fn take_divergence_answer(
3824 branch: &str,
3825 remote: &str,
3826 task: &mut Task,
3827) -> Option<crate::reconcile::Choice> {
3828 let summary = crate::reconcile::summary_for(branch, remote);
3829 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3830 (a.question == summary)
3831 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3832 .flatten()
3833 .map(|c| (i, c))
3834 })?;
3835 task.answers.remove(idx);
3836 Some(choice)
3837}
3838
3839fn choose_starter(
3851 review_branch: Option<&str>,
3852 branch_exists: bool,
3853 unfinished: Option<&str>,
3854) -> Starter {
3855 match review_branch {
3856 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3857 Some(_) => Starter::Start,
3858 None => match unfinished {
3859 Some(id) => Starter::Resume(id.to_owned()),
3860 None => Starter::Start,
3861 },
3862 }
3863}
3864
3865fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3868 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3869 return fallback.to_path_buf();
3870 }
3871 task.repo.clone()
3872}
3873
3874const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3878
3879fn answers_block(task: &Task, count: usize) -> String {
3881 let mut s = ANSWERS_HEADER.to_owned();
3882 for a in &task.answers[..count] {
3883 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3884 }
3885 s
3886}
3887
3888fn append_answers(base: &str, task: &Task) -> String {
3891 if task.answers.is_empty() {
3892 return base.to_owned();
3893 }
3894 let mut s = base.to_owned();
3895 s.push_str(&answers_block(task, task.answers.len()));
3896 s
3897}
3898
3899fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3903 for count in (1..=task.answers.len()).rev() {
3904 let block = answers_block(task, count);
3905 if let Some(base) = instruction.strip_suffix(&block) {
3906 return base;
3907 }
3908 }
3909 instruction
3910}
3911
3912fn instruction_for(task: &Task) -> String {
3920 append_answers(&task.instruction, task)
3921}
3922
3923fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3935 append_answers(strip_answers_block(old_instruction, task), task)
3936}
3937
3938fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3940 let paths = queue.attachment_paths(task);
3941 for (name, path) in task.attachments.iter().zip(&paths) {
3942 if !path.is_file() {
3943 bail!(
3944 "attachment `{name}` is recorded on the task but {} is missing",
3945 path.display()
3946 );
3947 }
3948 }
3949 Ok(paths)
3950}
3951
3952fn prepare_instruction(
3963 starter: &Starter,
3964 old_instruction: Option<&str>,
3965 task: &Task,
3966) -> Option<String> {
3967 match starter {
3968 Starter::Start => Some(instruction_for(task)),
3969 Starter::Resume(_) => Some(resumed_instruction(
3970 old_instruction.expect("a resumed run always has a prior instruction"),
3971 task,
3972 )),
3973 Starter::Review(_) => None,
3974 }
3975}
3976
3977fn record(queue: &Queue, task: &mut Task) {
3981 if let Err(e) = queue.put(task) {
3982 tracing::error!("could not record task {}: {e:#}", task.short());
3983 notices::raise(Notice::error(
3984 "loop:record",
3985 "The loop could not save a task's state; check the disk.",
3986 ));
3987 }
3988}
3989
3990fn runnable(queue: &Queue) -> Vec<Task> {
3996 let mut tasks: Vec<Task> = queue
3997 .list()
3998 .into_iter()
3999 .filter(|t| t.status.runnable())
4000 .collect();
4001 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
4002 tasks
4003}
4004
4005fn describe(state: &RunState) -> String {
4019 let p = phrases(&state.config.graph.language);
4020 let mut detail = if state.status == RunStatus::Stalled {
4021 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
4022 seats.sort_unstable();
4023 seats.dedup();
4024 if seats.is_empty() {
4025 p.quorum_lost.to_owned()
4026 } else {
4027 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
4028 }
4029 } else {
4030 format!("{}{}", p.run_ended, state.status.display_label())
4031 };
4032 if let Some(last) = state.events.last() {
4033 let unanswered = state
4037 .reviews
4038 .last()
4039 .filter(|r| {
4040 state.status == RunStatus::Blocked
4041 && last.node == "review"
4042 && r.incomplete()
4043 && r.blocking == 0
4044 && r.round == state.config.graph.review_rounds
4045 && r.e2e.iter().all(crate::run::CommandOutcome::ok)
4046 })
4047 .map(|r| (r.expected - r.answered, r.round));
4048 match unanswered {
4049 Some((missing, rounds)) if crate::lang::is_japanese(&state.config.graph.language) => {
4050 detail.push_str(&format!(
4051 " ({}: {})",
4052 last.node,
4053 (p.reviewers_never_answered)(missing, rounds)
4054 ));
4055 }
4056 _ => detail.push_str(&format!(" ({}: {})", last.node, last.message)),
4057 }
4058 }
4059 detail.push_str(&format!(" [run {}]", state.id));
4060 detail
4061}
4062
4063const DIAGNOSTIC_MAX: usize = 4_000;
4069
4070const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
4075
4076fn diagnostic(state: &RunState) -> Option<String> {
4090 let mut parts: Vec<String> = Vec::new();
4091
4092 for o in state.gate.iter().filter(|o| !o.ok()) {
4094 parts.push(format!(
4095 "gate `{}` failed ({:?}):\n{}",
4096 o.command,
4097 o.code,
4098 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
4099 ));
4100 }
4101
4102 if let Some(last) = state
4105 .events
4106 .iter()
4107 .rev()
4108 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
4109 {
4110 parts.push(last.message.clone());
4111 }
4112
4113 if state.viable().is_empty() {
4120 for c in &state.candidates {
4121 if let Some(evidence) = &c.verified_noop {
4122 parts.push(format!(
4123 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
4124 c.label
4125 ));
4126 } else if !c.summary.trim().is_empty() {
4127 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
4128 } else if let Some(why) = &c.failed {
4129 parts.push(format!("candidate {}: {why}", c.label));
4130 }
4131 }
4132 }
4133
4134 if parts.is_empty() {
4135 return None;
4136 }
4137 Some(crate::run::tail(
4142 &parts.join("\n\n"),
4143 DIAGNOSTIC_MAX.saturating_sub(100),
4144 ))
4145}
4146
4147fn label(status: RunStatus) -> &'static str {
4155 status.as_str()
4156}
4157
4158pub(crate) fn merge_mode(mode: &str) -> Result<MergeMode> {
4160 match mode {
4161 "none" => Ok(MergeMode::None),
4162 "local" => Ok(MergeMode::Local),
4163 "pr" => Ok(MergeMode::Pr),
4164 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
4165 }
4166}
4167
4168fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
4173 mutex
4174 .lock()
4175 .unwrap_or_else(std::sync::PoisonError::into_inner)
4176}
4177
4178#[cfg(test)]
4179mod tests {
4180 use super::*;
4181 use crate::queue::{Source, TaskStatus};
4182 use crate::run::{Candidate, CommandOutcome};
4183 use pretty_assertions::assert_eq;
4184
4185 fn task() -> Task {
4186 Task::new(
4187 "add retries".to_owned(),
4188 "add retries".to_owned(),
4189 PathBuf::from("/repo"),
4190 Source::Human,
4191 )
4192 }
4193
4194 fn interrupt_task(id: &str) -> Task {
4197 let mut t = task();
4198 t.id = id.to_owned();
4199 t.interrupt = true;
4200 t
4201 }
4202
4203 fn task_with_id(id: &str) -> Task {
4205 let mut t = task();
4206 t.id = id.to_owned();
4207 t
4208 }
4209
4210 fn urgent_task(id: &str) -> Task {
4212 let mut t = task();
4213 t.id = id.to_owned();
4214 t.urgent = true;
4215 t
4216 }
4217
4218 #[test]
4223 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
4224 assert_eq!(permit_kind(true, false), PermitKind::None);
4225 assert_eq!(permit_kind(true, true), PermitKind::None);
4226 }
4227
4228 #[test]
4233 fn permit_kind_separates_urgent_from_ordinary() {
4234 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
4235 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
4236 }
4237
4238 #[test]
4245 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
4246 let cfg = Config::default();
4247 let repo = Path::new("/any/repo/path");
4248
4249 let reason =
4250 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
4251 assert!(reason.contains("1024"), "{reason}");
4252 assert!(
4253 reason.contains(&cfg.disk.min_free_bytes.to_string()),
4254 "{reason}"
4255 );
4256
4257 assert_eq!(
4258 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
4259 None,
4260 "exactly at the floor is open"
4261 );
4262 assert_eq!(
4263 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
4264 None,
4265 "comfortably above the floor is open"
4266 );
4267 }
4268
4269 #[test]
4270 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
4271 let mut cfg = Config::default();
4272 cfg.disk.min_free_bytes = 0;
4273 let repo = Path::new("/any/repo/path");
4274 assert_eq!(
4275 disk_gate_with(repo, &cfg, |_| Ok(0)),
4276 None,
4277 "a zero floor never measures at all"
4278 );
4279 }
4280
4281 #[test]
4282 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
4283 let cfg = Config::default();
4284 let repo = Path::new("/any/repo/path");
4285 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
4286 .expect("a measurement failure must close the gate, not open it");
4287 assert!(reason.contains("could not measure"), "{reason}");
4288 }
4289
4290 #[test]
4291 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
4292 let ordinary = task();
4293 let next = advance_interrupt(
4294 Interrupt::Idle,
4295 std::slice::from_ref(&ordinary.id),
4296 std::slice::from_ref(&ordinary),
4297 );
4298 assert_eq!(next, Interrupt::Idle);
4299 }
4300
4301 #[test]
4302 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
4303 let marked = interrupt_task("marked");
4306 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4307 assert_eq!(next, Interrupt::Idle);
4308 }
4309
4310 #[test]
4311 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
4312 let marked = interrupt_task("marked");
4313 let next = advance_interrupt(
4314 Interrupt::Idle,
4315 &["running".to_owned()],
4316 std::slice::from_ref(&marked),
4317 );
4318 assert_eq!(
4319 next,
4320 Interrupt::Parking {
4321 parked: vec!["running".to_owned()],
4322 interrupt_task: "marked".to_owned(),
4323 }
4324 );
4325 }
4326
4327 #[test]
4336 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
4337 let marked = interrupt_task("marked");
4338
4339 let two = advance_interrupt(
4340 Interrupt::Idle,
4341 &["a".to_owned(), "b".to_owned()],
4342 std::slice::from_ref(&marked),
4343 );
4344 assert_eq!(two, Interrupt::Idle);
4345
4346 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4347 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
4348 }
4349
4350 #[test]
4351 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
4352 let state = Interrupt::Parking {
4353 parked: vec!["running".to_owned()],
4354 interrupt_task: "marked".to_owned(),
4355 };
4356 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
4358 assert_eq!(still_going, state);
4359
4360 let stopped_but_not_yet_dispatched =
4364 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
4365 assert_eq!(stopped_but_not_yet_dispatched, state);
4366
4367 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
4369 assert_eq!(
4370 dispatched,
4371 Interrupt::Running {
4372 parked: vec!["running".to_owned()],
4373 interrupt_task: "marked".to_owned(),
4374 }
4375 );
4376 }
4377
4378 #[test]
4379 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
4380 let state = Interrupt::Running {
4381 parked: vec!["running".to_owned()],
4382 interrupt_task: "marked".to_owned(),
4383 };
4384 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
4385 assert_eq!(still_running, state);
4386
4387 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
4394 assert_eq!(
4395 ended,
4396 Interrupt::Resuming {
4397 parked: vec!["running".to_owned()]
4398 }
4399 );
4400 }
4401
4402 #[test]
4403 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
4404 let state = Interrupt::Resuming {
4405 parked: vec!["running".to_owned()],
4406 };
4407 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
4408 assert_eq!(still_waiting, state);
4409
4410 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
4411 assert_eq!(dispatched, Interrupt::Idle);
4412 }
4413
4414 #[test]
4420 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
4421 {
4422 let state = Interrupt::Parking {
4423 parked: vec!["running".to_owned()],
4424 interrupt_task: "marked".to_owned(),
4425 };
4426 let next = advance_interrupt(state, &[], &[]);
4429 assert_eq!(
4430 next,
4431 Interrupt::Resuming {
4432 parked: vec!["running".to_owned()]
4433 },
4434 "abandoning the interrupt must not abandon the resume it owes"
4435 );
4436 }
4437
4438 #[test]
4441 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
4442 let state = Interrupt::Resuming {
4443 parked: vec!["running".to_owned()],
4444 };
4445 let next = advance_interrupt(state, &[], &[]);
4446 assert_eq!(
4447 next,
4448 Interrupt::Idle,
4449 "nothing is left to wait for; the loop must not stay wedged"
4450 );
4451 }
4452
4453 #[test]
4454 fn disabled_by_config_the_sequence_can_never_leave_idle() {
4455 let marked = interrupt_task("marked");
4456 let next = advance_interrupt_tick(
4457 false,
4458 Interrupt::Idle,
4459 &["running".to_owned()],
4460 std::slice::from_ref(&marked),
4461 );
4462 assert_eq!(
4463 next,
4464 Interrupt::Idle,
4465 "an unmarked, unconfigured daemon must behave exactly as before"
4466 );
4467 }
4468
4469 #[test]
4470 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4471 let state = Interrupt::Parking {
4472 parked: vec!["running".to_owned()],
4473 interrupt_task: "marked".to_owned(),
4474 };
4475 let candidates = vec![interrupt_task("marked"), task()];
4476 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4477 assert!(
4478 allowed.is_empty(),
4479 "nothing may dispatch - not even the interrupt task itself - \
4480 until the parked run has actually stopped"
4481 );
4482 }
4483
4484 #[test]
4498 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4499 for state in [
4500 Interrupt::Parking {
4501 parked: vec!["running".to_owned()],
4502 interrupt_task: "marked".to_owned(),
4503 },
4504 Interrupt::Running {
4505 parked: vec!["running".to_owned()],
4506 interrupt_task: "marked".to_owned(),
4507 },
4508 Interrupt::Resuming {
4509 parked: vec!["running".to_owned()],
4510 },
4511 ] {
4512 let candidates = vec![
4513 interrupt_task("marked"),
4514 urgent_task("hot"),
4515 task_with_id("ordinary"),
4516 ];
4517 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4518 assert!(
4519 !allowed.iter().any(|t| t.id == "hot"),
4520 "an urgent candidate must wait out the same gate as anything \
4521 else while the run it would run alongside has not actually \
4522 left flight, for state {state:?}: {allowed:?}"
4523 );
4524 }
4525 }
4526
4527 #[test]
4533 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4534 let state = Interrupt::Resuming {
4535 parked: vec!["hot".to_owned()],
4536 };
4537 let candidates = vec![urgent_task("hot"), task()];
4538 let allowed = interrupt_gate(&state, &[], candidates);
4539 assert_eq!(
4540 allowed.iter().filter(|t| t.id == "hot").count(),
4541 1,
4542 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4543 );
4544 }
4545
4546 #[test]
4547 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4548 let state = Interrupt::Parking {
4549 parked: vec!["running".to_owned()],
4550 interrupt_task: "marked".to_owned(),
4551 };
4552 let other = task();
4553 let candidates = vec![interrupt_task("marked"), other.clone()];
4554 let allowed = interrupt_gate(&state, &[], candidates);
4555 assert_eq!(allowed.len(), 1);
4556 assert_eq!(allowed[0].id, "marked");
4557 }
4558
4559 #[test]
4560 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4561 let state = Interrupt::Running {
4562 parked: vec!["running".to_owned()],
4563 interrupt_task: "marked".to_owned(),
4564 };
4565 let candidates = vec![task(), task()];
4566 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4567 assert!(allowed.is_empty());
4568 }
4569
4570 #[test]
4577 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4578 let state = Interrupt::Resuming {
4579 parked: vec!["a".to_owned(), "c".to_owned()],
4580 };
4581 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4582 let allowed = interrupt_gate(&state, &[], candidates);
4583 assert_eq!(
4584 allowed.len(),
4585 1,
4586 "at most one candidate may be offered while resuming: {allowed:?}"
4587 );
4588 assert_eq!(allowed[0].id, "a");
4589 }
4590
4591 #[test]
4592 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4593 let state = Interrupt::Resuming {
4594 parked: vec!["a".to_owned()],
4595 };
4596 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4597 assert!(allowed.is_empty());
4598 }
4599
4600 #[test]
4606 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4607 let running = task(); let marked = interrupt_task("marked");
4609
4610 let mut state = Interrupt::Idle;
4611 let in_flight = vec![running.id.clone()];
4613 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4614 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4615 assert!(gated.is_empty(), "still waiting on `running` to park");
4616
4617 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4619 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4620 assert_eq!(
4621 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4622 vec!["marked"],
4623 "only the interrupt task may be offered to the dispatcher now"
4624 );
4625
4626 state = advance_interrupt_tick(
4628 true,
4629 state,
4630 &["marked".to_owned()],
4631 std::slice::from_ref(&running),
4632 );
4633 let gated = interrupt_gate(
4634 &state,
4635 &["marked".to_owned()],
4636 vec![marked.clone(), running.clone()],
4637 );
4638 assert!(
4639 gated.is_empty(),
4640 "the parked run must not be offered back while the interrupt \
4641 task is still running"
4642 );
4643
4644 let other = task_with_id("other");
4648 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4649 assert_eq!(
4650 state,
4651 Interrupt::Resuming {
4652 parked: vec![running.id.clone()]
4653 }
4654 );
4655 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4656 assert_eq!(
4657 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4658 vec![running.id.as_str()],
4659 "exactly the parked run resumes - not the unrelated task, even \
4660 though it was offered first"
4661 );
4662
4663 state = advance_interrupt_tick(
4667 true,
4668 state,
4669 std::slice::from_ref(&running.id),
4670 std::slice::from_ref(&other),
4671 );
4672 assert_eq!(state, Interrupt::Idle);
4673 let gated = interrupt_gate(
4674 &state,
4675 std::slice::from_ref(&running.id),
4676 vec![other.clone()],
4677 );
4678 assert_eq!(
4679 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4680 vec![other.id.as_str()],
4681 "ordinary dispatch is unrestricted again"
4682 );
4683 }
4684
4685 #[test]
4686 fn every_run_status_settles_the_task_it_came_from() {
4687 let table = [
4689 (RunStatus::Merged, TaskStatus::Done, 1),
4690 (RunStatus::Ready, TaskStatus::Done, 1),
4691 (RunStatus::Stalled, TaskStatus::Failed, 0),
4692 (RunStatus::Blocked, TaskStatus::Failed, 1),
4693 (RunStatus::Failed, TaskStatus::Failed, 1),
4694 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4695 (RunStatus::Prep, TaskStatus::Failed, 1),
4696 (RunStatus::Implementing, TaskStatus::Failed, 1),
4697 (RunStatus::Judging, TaskStatus::Failed, 1),
4698 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4699 (RunStatus::Voting, TaskStatus::Failed, 1),
4700 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4701 (RunStatus::Gating, TaskStatus::Failed, 1),
4702 ];
4703 for (run, want, attempts) in table {
4704 let mut t = task();
4705 t.start("20260902-000000-aaaa".to_owned());
4706 settle(
4707 &mut t,
4708 Verdict {
4709 status: run,
4710 left_pr: false,
4711 parked: false,
4712 quota_hit: matches!(run, RunStatus::Stalled),
4713 no_viable_candidates: false,
4714 },
4715 "why",
4716 2,
4717 );
4718 assert_eq!(t.status, want, "task status after {}", label(run));
4719 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4720 }
4721 }
4722
4723 #[test]
4724 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4725 let mut stalled = task();
4726 stalled.start("20260902-000000-aaaa".to_owned());
4727 settle(
4728 &mut stalled,
4729 Verdict {
4730 status: RunStatus::Stalled,
4731 left_pr: false,
4732 parked: false,
4733 quota_hit: true,
4734 no_viable_candidates: false,
4735 },
4736 "quota",
4737 1,
4738 );
4739 assert_eq!(stalled.attempts, 0);
4740 assert!(
4741 stalled.status.runnable(),
4742 "a machine problem must leave the task in line"
4743 );
4744
4745 let mut blocked = task();
4746 blocked.start("20260902-000000-aaaa".to_owned());
4747 settle(
4748 &mut blocked,
4749 Verdict {
4750 status: RunStatus::Blocked,
4751 left_pr: false,
4752 parked: false,
4753 quota_hit: false,
4754 no_viable_candidates: false,
4755 },
4756 "findings open",
4757 1,
4758 );
4759 assert_eq!(blocked.attempts, 1);
4760 assert_eq!(
4761 blocked.status,
4762 TaskStatus::Held,
4763 "the last attempt hands the task to a human"
4764 );
4765 }
4766
4767 #[test]
4768 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4769 let mut delivered = task();
4772 delivered.start("20260903-080619-01c2".to_owned());
4773 settle(
4774 &mut delivered,
4775 Verdict {
4776 status: RunStatus::Blocked,
4777 left_pr: true,
4778 parked: false,
4779 quota_hit: false,
4780 no_viable_candidates: false,
4781 },
4782 "no check status",
4783 4,
4784 );
4785 assert_eq!(
4786 delivered.status,
4787 TaskStatus::Held,
4788 "a pull request waiting on CI or a person is not a retryable failure"
4789 );
4790 assert!(
4791 !delivered.status.runnable(),
4792 "the loop must not pick this task up again"
4793 );
4794 assert_eq!(
4795 delivered.last_error.as_deref(),
4796 Some("no check status"),
4797 "the operator needs to be told what the gate was waiting for"
4798 );
4799
4800 let mut empty_handed = task();
4803 empty_handed.start("20260903-080619-01c2".to_owned());
4804 settle(
4805 &mut empty_handed,
4806 Verdict {
4807 status: RunStatus::Blocked,
4808 left_pr: false,
4809 parked: false,
4810 quota_hit: false,
4811 no_viable_candidates: false,
4812 },
4813 "findings open",
4814 4,
4815 );
4816 assert_eq!(empty_handed.status, TaskStatus::Failed);
4817 assert!(empty_handed.status.runnable());
4818 }
4819
4820 #[test]
4821 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4822 let mut noop = task();
4828 noop.start("20260912-131304-391f".to_owned());
4829 settle(
4830 &mut noop,
4831 Verdict {
4832 status: RunStatus::VerifiedNoop,
4833 left_pr: false,
4834 parked: false,
4835 quota_hit: false,
4836 no_viable_candidates: true,
4837 },
4838 "candidate A: already fixed by b32cfc4, on main",
4839 4,
4840 );
4841 assert_eq!(
4842 noop.status,
4843 TaskStatus::Held,
4844 "an unverified claim is a request for a human, not a failure"
4845 );
4846 assert!(
4847 !noop.status.runnable(),
4848 "the loop must not requeue this on the same unverified claim"
4849 );
4850 assert_eq!(noop.attempts, 1);
4855 }
4856
4857 #[test]
4858 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4859 let mut parked = task();
4864 parked.start("20260903-183634-2d98".to_owned());
4865 settle(
4866 &mut parked,
4867 Verdict {
4868 status: RunStatus::Implementing,
4869 left_pr: false,
4870 quota_hit: false,
4871 parked: true,
4872 no_viable_candidates: false,
4873 },
4874 "parked after `implementing`",
4875 2,
4876 );
4877 assert_eq!(parked.attempts, 0, "a park is refunded");
4878 assert!(
4879 parked.status.runnable(),
4880 "and the task stays in line so the next loop resumes its run"
4881 );
4882 assert_eq!(parked.status, TaskStatus::Parked);
4883 assert_eq!(
4884 parked.park_reason.as_deref(),
4885 Some("parked after `implementing`"),
4886 "the card says where it stopped"
4887 );
4888 assert_eq!(parked.last_error, None, "nothing failed");
4889
4890 let mut broken = task();
4894 broken.start("20260903-183634-2d98".to_owned());
4895 settle(
4896 &mut broken,
4897 Verdict {
4898 status: RunStatus::Implementing,
4899 left_pr: false,
4900 quota_hit: false,
4901 parked: false,
4902 no_viable_candidates: false,
4903 },
4904 "returned mid-flight",
4905 2,
4906 );
4907 assert_eq!(broken.attempts, 1);
4908 }
4909
4910 #[test]
4911 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4912 let mut flaky = task();
4917 flaky.start("20260903-123023-e633".to_owned());
4918 settle(
4919 &mut flaky,
4920 Verdict {
4921 status: RunStatus::Stalled,
4922 left_pr: false,
4923 parked: false,
4924 quota_hit: false,
4925 no_viable_candidates: false,
4926 },
4927 "verdict rests on 1 of 3 judges",
4928 2,
4929 );
4930 assert_eq!(
4931 flaky.attempts, 1,
4932 "flakiness spends an attempt, so `max_attempts` still bounds it"
4933 );
4934 assert!(flaky.status.runnable(), "and it is still worth retrying");
4935
4936 let mut limited = task();
4938 limited.start("20260903-123023-e633".to_owned());
4939 settle(
4940 &mut limited,
4941 Verdict {
4942 status: RunStatus::Stalled,
4943 left_pr: false,
4944 parked: false,
4945 quota_hit: true,
4946 no_viable_candidates: false,
4947 },
4948 "judge-2, judge-3 out of quota",
4949 2,
4950 );
4951 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4952 assert!(limited.status.runnable());
4953
4954 let mut worn = task();
4957 for _ in 0..2 {
4958 worn.release();
4959 }
4960 worn.start("20260903-123023-e633".to_owned());
4961 worn.attempts = 2;
4962 settle(
4963 &mut worn,
4964 Verdict {
4965 status: RunStatus::Stalled,
4966 left_pr: false,
4967 parked: false,
4968 quota_hit: false,
4969 no_viable_candidates: false,
4970 },
4971 "no quorum again",
4972 2,
4973 );
4974 assert_eq!(worn.status, TaskStatus::Held);
4975 assert!(!worn.status.runnable());
4976 }
4977
4978 #[test]
4979 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4980 let mut wiped_out = task();
4987 wiped_out.start("20260907-025000-a1b2".to_owned());
4988 settle(
4989 &mut wiped_out,
4990 Verdict {
4991 status: RunStatus::Failed,
4992 left_pr: false,
4993 parked: false,
4994 quota_hit: true,
4995 no_viable_candidates: true,
4996 },
4997 "no candidate produced a change; nothing to judge",
4998 2,
4999 );
5000 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
5001 assert!(
5002 wiped_out.status.runnable(),
5003 "a machine problem must leave the task in line"
5004 );
5005
5006 let mut partial_progress = task();
5012 partial_progress.start("20260907-025500-c3d4".to_owned());
5013 settle(
5014 &mut partial_progress,
5015 Verdict {
5016 status: RunStatus::Failed,
5017 left_pr: false,
5018 parked: false,
5019 quota_hit: true,
5020 no_viable_candidates: false,
5021 },
5022 "gate failed on the winning candidate",
5023 2,
5024 );
5025 assert_eq!(
5026 partial_progress.attempts, 1,
5027 "a candidate that actually produced a change spends the attempt \
5028 even though some other seat hit its quota"
5029 );
5030 assert!(partial_progress.status.runnable());
5031 }
5032
5033 #[test]
5034 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
5035 let mut t = task();
5041 t.start("20260907-025000-a1b2".to_owned());
5042 let mut state = run_state(RunStatus::Failed);
5043 state.quota.push(QuotaLoss {
5044 seat: "cand-a".to_owned(),
5045 node: "implement".to_owned(),
5046 at: Timestamp::now(),
5047 reset: None,
5048 });
5049 assert!(
5050 state.viable().is_empty(),
5051 "no candidate was added, so nothing is viable"
5052 );
5053 reclaim(&mut t, Some(state), 2, "en");
5054 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
5055 assert!(t.status.runnable());
5056 }
5057
5058 #[test]
5059 fn a_held_task_is_never_offered_to_the_loop() {
5060 let dir = tempfile::tempdir().unwrap();
5061 let queue = Queue::at(dir.path().to_path_buf());
5062 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
5063 let mut t = task();
5064 t.id = format!("2026090{n}-000000-000{n}");
5065 t.priority = priority;
5066 queue.put(&mut t).unwrap();
5067 }
5068 let mut held = task();
5069 held.id = "20260909-000000-9999".to_owned();
5070 held.priority = 99;
5071 held.hold_machine(None);
5072 queue.put(&mut held).unwrap();
5073
5074 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
5075 assert_eq!(order.len(), 3);
5076 assert!(!order.contains(&held.id));
5077 assert_eq!(
5078 order.first().cloned(),
5079 queue.next_runnable().map(|t| t.id),
5080 "the loop's first candidate is exactly what the queue offers"
5081 );
5082 assert_eq!(
5083 order,
5084 vec![
5085 "20260902-000000-0002".to_owned(),
5086 "20260903-000000-0003".to_owned(),
5087 "20260901-000000-0001".to_owned(),
5088 ],
5089 "priority first, then oldest, so nothing starves"
5090 );
5091 }
5092
5093 #[test]
5094 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
5095 let dir = tempfile::tempdir().unwrap();
5096 let queue = Queue::at(dir.path().to_path_buf());
5097 let mut old = task();
5098 old.id = "20260101-000000-old0".to_owned();
5099 queue.put(&mut old).unwrap();
5100 let mut fresh = task();
5101 fresh.id = "20260101-000000-new0".to_owned();
5102 queue.put(&mut fresh).unwrap();
5103
5104 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
5108 std::thread::sleep(Duration::from_millis(60));
5109 let live = queue.claim(&fresh.id).unwrap();
5110
5111 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
5112 assert_eq!(swept, vec![old.id.clone()]);
5113 assert!(
5114 queue.claim(&old.id).is_ok(),
5115 "an unparseable lock older than the threshold is swept"
5116 );
5117 assert!(
5118 queue.claim(&fresh.id).is_err(),
5119 "a live pid protects its lock regardless of age"
5120 );
5121 drop(live);
5122 }
5123
5124 #[test]
5125 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
5126 let dir = tempfile::tempdir().unwrap();
5136 let queue = Queue::at(dir.path().to_path_buf());
5137 let mut t = task();
5138 t.id = "20260101-000000-live".to_owned();
5139 queue.put(&mut t).unwrap();
5140
5141 let claim = queue.claim(&t.id).unwrap();
5142 std::thread::sleep(Duration::from_millis(60));
5143
5144 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
5145 assert!(
5146 swept.is_empty(),
5147 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
5148 );
5149 assert!(
5150 queue.claim(&t.id).is_err(),
5151 "the lock still protects its task"
5152 );
5153 drop(claim);
5154 }
5155
5156 fn injected_dead_pid() -> u32 {
5159 std::process::id().checked_add(1).unwrap_or(1)
5160 }
5161
5162 #[test]
5163 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
5164 let dir = tempfile::tempdir().unwrap();
5165 let queue = Queue::at(dir.path().to_path_buf());
5166 let mut t = task();
5167 t.id = "20260101-000000-dead".to_owned();
5168 queue.put(&mut t).unwrap();
5169 let dead_pid = injected_dead_pid();
5170
5171 std::fs::write(
5176 dir.path().join(format!("{}.lock", t.id)),
5177 dead_pid.to_string(),
5178 )
5179 .unwrap();
5180
5181 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5182 pid != dead_pid
5183 });
5184 assert_eq!(
5185 swept,
5186 vec![t.id.clone()],
5187 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
5188 );
5189 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
5190 }
5191
5192 #[test]
5193 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
5194 let dir = tempfile::tempdir().unwrap();
5195 let queue = Queue::at(dir.path().to_path_buf());
5196 let mut t = task();
5197 t.id = "20260101-000000-late".to_owned();
5198 queue.put(&mut t).unwrap();
5199 let dead_pid = injected_dead_pid();
5200
5201 assert!(
5204 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
5205 "nothing has claimed the task yet"
5206 );
5207
5208 std::fs::write(
5211 dir.path().join(format!("{}.lock", t.id)),
5212 dead_pid.to_string(),
5213 )
5214 .unwrap();
5215
5216 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5220 pid != dead_pid
5221 });
5222 assert_eq!(swept, vec![t.id.clone()]);
5223 }
5224
5225 #[test]
5226 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
5227 crate::run::pin_test_home();
5232 let dir = tempfile::tempdir().unwrap();
5233 let queue = Queue::at(dir.path().to_path_buf());
5234 let mut t = task();
5235 t.id = "20260101-000000-crsh".to_owned();
5236 t.status = TaskStatus::Running;
5237 t.attempts = 1;
5238 t.runs.push("20260904-000000-4043".to_owned());
5242 queue.put(&mut t).unwrap();
5243 let dead_pid = injected_dead_pid();
5244
5245 std::fs::write(
5248 dir.path().join(format!("{}.lock", t.id)),
5249 dead_pid.to_string(),
5250 )
5251 .unwrap();
5252
5253 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
5259 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
5260
5261 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5262 pid != dead_pid
5263 });
5264 assert_eq!(swept, vec![t.id.clone()]);
5265
5266 let reclaimed = reclaim_orphaned_running(&queue, 2);
5267 assert_eq!(reclaimed, vec![t.id.clone()]);
5268 let after = queue.get(&t.id).unwrap();
5269 assert_eq!(
5270 after.status,
5271 TaskStatus::Held,
5272 "no run.json to recover from, so a human is asked"
5273 );
5274 assert_eq!(
5275 after.runs,
5276 vec!["20260904-000000-4043".to_owned()],
5277 "the crashed run's id is kept as evidence, not discarded"
5278 );
5279 }
5280
5281 #[test]
5282 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
5283 let dir = tempfile::tempdir().unwrap();
5284 let queue = Queue::at(dir.path().to_path_buf());
5285 let mut t = task();
5286 t.id = "20260101-000000-unknown".to_owned();
5287 queue.put(&mut t).unwrap();
5288 let dead_pid = injected_dead_pid();
5289 std::fs::write(
5290 dir.path().join(format!("{}.lock", t.id)),
5291 dead_pid.to_string(),
5292 )
5293 .unwrap();
5294
5295 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
5296 assert!(swept.is_empty(), "an unknown pid must keep its lock");
5297 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
5298 }
5299
5300 fn run_state_in(status: RunStatus, language: &str) -> RunState {
5301 let mut s = run_state(status);
5302 s.config.graph.language = language.to_owned();
5303 s
5304 }
5305
5306 fn unstarted_verdict(status: RunStatus) -> Verdict {
5307 Verdict {
5308 status,
5309 left_pr: false,
5310 quota_hit: false,
5311 parked: false,
5312 no_viable_candidates: false,
5313 }
5314 }
5315
5316 #[test]
5317 fn an_already_in_base_run_finishes_the_task_without_spending_an_attempt() {
5318 let mut t = task();
5319 t.attempts = 1;
5320 settle_in(
5321 &mut t,
5322 unstarted_verdict(RunStatus::AlreadyInBase),
5323 "already in main",
5324 1,
5325 phrases("en"),
5326 );
5327 assert_eq!(t.status, TaskStatus::Done);
5328 assert_eq!(t.attempts, 0);
5329 }
5330
5331 #[test]
5332 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
5333 let reason = |language: &str| {
5334 let mut t = task();
5335 settle_in(
5336 &mut t,
5337 unstarted_verdict(RunStatus::Judging),
5338 "boom",
5339 1,
5340 phrases(language),
5341 );
5342 t.last_error.or(t.hold_reason).unwrap_or_default()
5343 };
5344 assert!(
5345 reason("en").starts_with("the graph stopped at `"),
5346 "{}",
5347 reason("en")
5348 );
5349 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
5350 assert!(reason("日本語").contains("boom"));
5351 assert_eq!(reason("fr"), reason("en"));
5352 }
5353
5354 #[test]
5355 fn describe_follows_the_run_language_and_keeps_the_run_id() {
5356 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
5357 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
5358 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
5359 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
5360 assert!(ja.contains("[run "), "{ja}");
5361 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
5362 assert!(ended.starts_with("run 終了: "), "{ended}");
5363 let mut de = run_state_in(RunStatus::Failed, "de");
5364 let mut en = run_state_in(RunStatus::Failed, "en");
5365 de.id = "same".to_owned();
5366 en.id = "same".to_owned();
5367 assert_eq!(describe(&de), describe(&en));
5368 }
5369
5370 #[test]
5371 fn handover_refusals_follow_the_language_and_keep_the_detail_apart() {
5372 use crate::handover::Refused;
5373 let cases = [
5374 Refused::Foreign {
5375 branch: "b".into(),
5376 path: "/w/x".into(),
5377 why: "made by hand".into(),
5378 },
5379 Refused::Unsafe {
5380 branch: "b".into(),
5381 path: "/w/x".into(),
5382 why: "its worktree has uncommitted changes (a.rs)".into(),
5383 },
5384 Refused::ReleaseFailed {
5385 branch: "b".into(),
5386 path: "/w/x".into(),
5387 run: "ab12".into(),
5388 },
5389 ];
5390 for r in &cases {
5391 let en = (phrases("en").handover_refused)(r);
5392 assert_eq!(en, r.to_string());
5393 let ja = (phrases("ja").handover_refused)(r);
5394 assert!(
5395 !ja.contains("is checked out") && !ja.contains("try again"),
5396 "{ja}"
5397 );
5398 assert!(ja.contains("`b`") && ja.contains("/w/x"), "{ja}");
5399 if let Refused::Foreign { why, .. } | Refused::Unsafe { why, .. } = r {
5400 assert!(ja.contains(&format!("(詳細: {why})")), "{ja}");
5401 }
5402 }
5403 assert!(phrases("en").handover_hint.contains("release the task"));
5404 assert!(phrases("ja").handover_hint.contains("解放"));
5405 }
5406
5407 #[test]
5408 fn describe_translates_the_unanswered_reviewer_stop_only_in_ja() {
5409 let build = |lang: &str| {
5410 let mut s = run_state_in(RunStatus::Blocked, lang);
5411 s.id = "same".to_owned();
5412 s.config.graph.review_rounds = 3;
5413 let mut r = review_round(3);
5414 r.expected = 3;
5415 r.answered = 1;
5416 s.reviews.push(r);
5417 s.event(
5418 "review",
5419 "2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean",
5420 );
5421 s
5422 };
5423 let en = describe(&build("en"));
5424 assert!(en.contains("(review: 2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean)"), "{en}");
5425 let ja = describe(&build("ja"));
5426 assert!(ja.contains("2 席のレビュアーが 3 ラウンド"), "{ja}");
5427 assert!(!ja.contains("never answered"), "{ja}");
5428 assert!(ja.contains("[run "), "{ja}");
5429
5430 let mut failed = build("ja");
5433 failed.reviews[0].e2e.push(crate::run::CommandOutcome {
5434 command: "cargo test".to_owned(),
5435 code: Some(1),
5436 output_tail: String::new(),
5437 duration_ms: 0,
5438 resource_blocked: false,
5439 });
5440 failed.event("review", "stopped; e2e failed: cargo test");
5441 let ja = describe(&failed);
5442 assert!(ja.contains("stopped; e2e failed: cargo test"), "{ja}");
5443 assert!(!ja.contains("席のレビュアー"), "{ja}");
5444 }
5445
5446 #[test]
5447 fn refusals_and_recovery_prose_follow_the_language() {
5448 let t = held_task_with("r1");
5449 let q = action_question(
5450 "r1",
5451 ask::ChoiceAction::Resume {
5452 run: "r1".to_owned(),
5453 },
5454 );
5455 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
5456 ActionDecision::Refuse(s) => s,
5457 other => panic!("{other:?}"),
5458 };
5459 assert!(refuse(phrases("en")).contains("could not be read: gone"));
5460 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
5461 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
5462
5463 let mut held = task();
5464 reclaim(&mut held, None, 2, "ja");
5465 assert!(held.hold_reason.unwrap().contains("保留にしました"));
5466 let mut held = task();
5467 reclaim(&mut held, None, 2, "xx");
5468 assert!(held.hold_reason.unwrap().contains("held for a human"));
5469 }
5470
5471 fn run_state(status: RunStatus) -> RunState {
5472 let mut state = RunState::new(
5473 PathBuf::from("/repo"),
5474 "main".to_owned(),
5475 "abc1234def".to_owned(),
5476 "add retries".to_owned(),
5477 Config::default(),
5478 );
5479 state.status = status;
5480 state
5481 }
5482
5483 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
5484 Candidate {
5485 index: 0,
5486 label,
5487 agent: "claude".to_owned(),
5488 branch: format!("magi/x/{label}"),
5489 worktree: PathBuf::from("/repo"),
5490 summary: summary.to_owned(),
5491 stat: String::new(),
5492 files: 0,
5493 commits: usize::from(!empty),
5494 empty,
5495 failed: failed.map(str::to_owned),
5496 verified_noop: None,
5497 duration_ms: 0,
5498 folded: false,
5499 }
5500 }
5501
5502 #[test]
5503 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
5504 let mut state = run_state(RunStatus::Blocked);
5505 state.gate = vec![
5506 CommandOutcome {
5507 command: "cargo make check".to_owned(),
5508 code: Some(0),
5509 output_tail: "ok".to_owned(),
5510 duration_ms: 0,
5511 resource_blocked: false,
5512 },
5513 CommandOutcome {
5514 command: "cargo test".to_owned(),
5515 code: Some(101),
5516 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
5517 duration_ms: 0,
5518 resource_blocked: false,
5519 },
5520 ];
5521 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
5522 assert!(d.contains("cargo test"), "{d}");
5523 assert!(
5524 !d.contains("cargo make check"),
5525 "a passing check is not a diagnostic: {d}"
5526 );
5527 assert!(d.contains("assertion failed"), "{d}");
5528 }
5529
5530 #[test]
5531 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
5532 let mut state = run_state(RunStatus::Blocked);
5533 state.event(
5534 "land",
5535 "stopped: the fixer produced no commit while 2 check(s) were failing \
5536 (build, lint); stopping instead of looping on an unchanged tree",
5537 );
5538 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
5539 assert!(d.contains("build"), "{d}");
5540 assert!(d.contains("lint"), "{d}");
5541 assert!(d.contains("fixer produced no commit"), "{d}");
5542 }
5543
5544 #[test]
5545 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
5546 let state = run_state(RunStatus::VerifiedNoop);
5551 let d = describe(&state);
5552 assert!(
5553 d.contains("agent-verified no-op"),
5554 "expected the display label, not the wire spelling: {d}"
5555 );
5556 assert!(!d.contains("verified_noop"), "{d}");
5557 }
5558
5559 #[test]
5560 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5561 let mut state = run_state(RunStatus::Failed);
5567 state.candidates = vec![candidate(
5568 'A',
5569 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5570 true,
5571 None,
5572 )];
5573 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5574 assert!(d.contains("candidate A"), "{d}");
5575 assert!(d.contains("tagged v1.2.3"), "{d}");
5576 }
5577
5578 #[test]
5579 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5580 let mut state = run_state(RunStatus::Failed);
5581 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5582 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5583 assert!(d.contains("candidate A"), "{d}");
5584 assert!(d.contains("agent timed out"), "{d}");
5585 }
5586
5587 #[test]
5588 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5589 let mut state = run_state(RunStatus::Failed);
5592 state.candidates = vec![candidate('A', "did the work", false, None)];
5593 assert!(diagnostic(&state).is_none());
5594 }
5595
5596 #[test]
5597 fn diagnostic_is_bounded_however_much_a_run_printed() {
5598 let mut state = run_state(RunStatus::Blocked);
5599 state.gate = vec![
5600 CommandOutcome {
5601 command: "cargo test".to_owned(),
5602 code: Some(101),
5603 output_tail: "x".repeat(50_000),
5604 duration_ms: 0,
5605 resource_blocked: false,
5606 },
5607 CommandOutcome {
5608 command: "cargo clippy".to_owned(),
5609 code: Some(1),
5610 output_tail: "y".repeat(50_000),
5611 duration_ms: 0,
5612 resource_blocked: false,
5613 },
5614 ];
5615 state.candidates = vec![
5616 candidate('A', &"z".repeat(50_000), true, None),
5617 candidate('B', &"w".repeat(50_000), true, None),
5618 ];
5619 let d = diagnostic(&state).expect("plenty here to diagnose");
5620 assert!(
5621 d.len() <= DIAGNOSTIC_MAX,
5622 "diagnostic grew to {} bytes, unbounded",
5623 d.len()
5624 );
5625 }
5626
5627 #[test]
5628 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5629 let mut state = run_state(RunStatus::Blocked);
5630 state.gate = vec![CommandOutcome {
5631 command: "cargo test".to_owned(),
5632 code: Some(101),
5633 output_tail: "assertion failed".to_owned(),
5634 duration_ms: 0,
5635 resource_blocked: false,
5636 }];
5637 let verdict = Verdict {
5638 status: RunStatus::Blocked,
5639 left_pr: false,
5640 quota_hit: false,
5641 parked: false,
5642 no_viable_candidates: false,
5643 };
5644
5645 let mut t = task();
5648 t.start("run-1".to_owned());
5649 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5650 assert_eq!(t.status, TaskStatus::Failed);
5651 assert!(t.diagnostic.is_none());
5652
5653 t.start("run-2".to_owned());
5656 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5657 assert_eq!(t.status, TaskStatus::Held);
5658 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5659 assert!(d.contains("cargo test"), "{d}");
5660 }
5661
5662 #[test]
5663 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5664 crate::run::pin_test_home();
5669 let home = crate::run::home();
5670 let state = run_state(RunStatus::VerifiedNoop);
5671 let mut q = ask::Question::new(
5672 state.id.clone(),
5673 "implement".to_owned(),
5674 "impl-A".to_owned(),
5675 "is this really a no-op?".to_owned(),
5676 String::new(),
5677 Vec::new(),
5678 );
5679 Questions::at(home.join("questions")).put(&mut q).unwrap();
5680
5681 let verdict = Verdict {
5682 status: RunStatus::VerifiedNoop,
5683 left_pr: false,
5684 quota_hit: false,
5685 parked: false,
5686 no_viable_candidates: false,
5687 };
5688 let mut t = task();
5689 t.start(state.id.clone());
5690 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5691
5692 assert_eq!(t.status, TaskStatus::Held);
5693 let reason = t.hold_reason.expect("a held task must record why");
5694 assert!(
5695 reason.starts_with("run ended agent-verified no-op"),
5696 "the original settle reason must survive unchanged: {reason}"
5697 );
5698 assert!(
5699 reason.contains(q.short()),
5700 "the open question's id must be named so the notice is actionable: {reason}"
5701 );
5702 }
5703
5704 #[test]
5705 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5706 crate::run::pin_test_home();
5707 let state = run_state(RunStatus::VerifiedNoop);
5708
5709 let verdict = Verdict {
5710 status: RunStatus::VerifiedNoop,
5711 left_pr: false,
5712 quota_hit: false,
5713 parked: false,
5714 no_viable_candidates: false,
5715 };
5716 let mut t = task();
5717 t.start(state.id.clone());
5718 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5719
5720 assert_eq!(t.status, TaskStatus::Held);
5721 assert_eq!(
5722 t.hold_reason.as_deref(),
5723 Some("run ended agent-verified no-op"),
5724 "nothing to append when the question was already answered or never asked"
5725 );
5726 }
5727
5728 #[test]
5729 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5730 crate::run::pin_test_home();
5731 let mut first = run_state(RunStatus::Blocked);
5732 first.id = "20260101-000000-sup1".to_owned();
5733 first.save().unwrap();
5734 let mut second = run_state(RunStatus::Merged);
5735 second.id = "20260101-000000-sup2".to_owned();
5736 second.save().unwrap();
5737
5738 let mut t = task();
5739 t.runs = vec![first.id.clone(), second.id.clone()];
5740 t.status = TaskStatus::Done;
5741
5742 supersede_prior_runs(&t, &crate::run::home());
5743
5744 assert_eq!(
5745 RunState::load(&first.id).unwrap().status,
5746 RunStatus::Superseded,
5747 "the first attempt's Blocked no longer needs anyone's attention"
5748 );
5749 assert_eq!(
5750 RunState::load(&second.id).unwrap().status,
5751 RunStatus::Merged,
5752 "the run that actually succeeded is left exactly as it was"
5753 );
5754 }
5755
5756 #[test]
5757 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5758 crate::run::pin_test_home();
5764 let mut first = run_state(RunStatus::Blocked);
5765 first.id = "20260101-000000-sup9".to_owned();
5766 first.driver_pid = Some(std::process::id());
5769 first.driver_started_at = Some(
5770 crate::proc::process_started_at(std::process::id())
5771 .expect("this test process's own start time must be queryable"),
5772 );
5773 first.save().unwrap();
5774 let mut second = run_state(RunStatus::Merged);
5775 second.id = "20260101-000000-supa".to_owned();
5776 second.save().unwrap();
5777
5778 let mut t = task();
5779 t.runs = vec![first.id.clone(), second.id.clone()];
5780 t.status = TaskStatus::Done;
5781
5782 supersede_prior_runs(&t, &crate::run::home());
5783
5784 assert_eq!(
5785 RunState::load(&first.id).unwrap().status,
5786 RunStatus::Blocked,
5787 "a live driver_pid means something is still actually working this run, \
5788 even though no daemon claims it - rewriting under it would just be \
5789 undone the next time that process saves"
5790 );
5791 }
5792
5793 #[test]
5794 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5795 let dir = tempfile::tempdir().unwrap();
5801 let home = dir.path().to_path_buf();
5802 let queue = Queue::at(dir.path().join("queue"));
5803
5804 let mut first = run_state(RunStatus::Blocked);
5805 first.id = "20260101-000000-supd".to_owned();
5806 first.driver_pid = Some(std::process::id());
5807 first.driver_started_at = Some(
5808 crate::proc::process_started_at(std::process::id())
5809 .expect("this test process's own start time must be queryable"),
5810 );
5811 first.save_under(&home).unwrap();
5812 let mut second = run_state(RunStatus::Merged);
5813 second.id = "20260101-000000-supe".to_owned();
5814 second.save_under(&home).unwrap();
5815
5816 let mut t = task();
5817 t.runs = vec![first.id.clone(), second.id.clone()];
5818 t.status = TaskStatus::Done;
5819 queue.put(&mut t).unwrap();
5820
5821 resweep_superseded_attempts(&queue, &home);
5822 assert_eq!(
5823 RunState::load_under(&first.id, &home).unwrap().status,
5824 RunStatus::Blocked,
5825 "still live on the first pass, so still untouched"
5826 );
5827
5828 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5834 stale.driver_started_at = Some("1".to_owned());
5835 stale.save_under(&home).unwrap();
5836
5837 resweep_superseded_attempts(&queue, &home);
5838 assert_eq!(
5839 RunState::load_under(&first.id, &home).unwrap().status,
5840 RunStatus::Superseded,
5841 "the second pass catches up what the first one correctly skipped"
5842 );
5843 }
5844
5845 #[test]
5846 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5847 crate::run::pin_test_home();
5848 let mut first = run_state(RunStatus::Blocked);
5849 first.id = "20260101-000000-sup3".to_owned();
5850 first.save().unwrap();
5851 let mut second = run_state(RunStatus::Blocked);
5852 second.id = "20260101-000000-sup4".to_owned();
5853 second.save().unwrap();
5854
5855 let mut t = task();
5856 t.runs = vec![first.id.clone(), second.id.clone()];
5857 t.status = TaskStatus::Failed;
5861
5862 supersede_prior_runs(&t, &crate::run::home());
5863
5864 assert_eq!(
5865 RunState::load(&first.id).unwrap().status,
5866 RunStatus::Blocked
5867 );
5868 assert_eq!(
5869 RunState::load(&second.id).unwrap().status,
5870 RunStatus::Blocked
5871 );
5872 }
5873
5874 #[test]
5875 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5876 crate::run::pin_test_home();
5880 let mut first = run_state(RunStatus::Blocked);
5881 first.id = "20260101-000000-sup5".to_owned();
5882 first.save().unwrap();
5883
5884 let mut t = task();
5885 t.runs = vec![first.id.clone()];
5886 t.status = TaskStatus::Done;
5887
5888 supersede_prior_runs(&t, &crate::run::home());
5889
5890 assert_eq!(
5891 RunState::load(&first.id).unwrap().status,
5892 RunStatus::Blocked,
5893 "a single-attempt task has no earlier run to supersede"
5894 );
5895 }
5896
5897 #[test]
5898 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5899 crate::run::pin_test_home();
5906 let mut first = run_state(RunStatus::Blocked);
5907 first.id = "20260101-000000-supb".to_owned();
5908 first.save().unwrap();
5909 let mut second = run_state(RunStatus::Failed);
5910 second.id = "20260101-000000-supc".to_owned();
5911 second.save().unwrap();
5912
5913 let mut t = task();
5914 t.runs = vec![first.id.clone(), second.id.clone()];
5915 t.status = TaskStatus::Done;
5916
5917 supersede_prior_runs(&t, &crate::run::home());
5918
5919 assert_eq!(
5920 RunState::load(&first.id).unwrap().status,
5921 RunStatus::Blocked,
5922 "the task's last attempt never landed, so there is nothing here \
5923 actually superseding it"
5924 );
5925 }
5926
5927 #[test]
5928 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5929 crate::run::pin_test_home();
5933 let mut failed = run_state(RunStatus::Failed);
5934 failed.id = "20260101-000000-sup6".to_owned();
5935 failed.save().unwrap();
5936 let mut noop = run_state(RunStatus::VerifiedNoop);
5937 noop.id = "20260101-000000-sup7".to_owned();
5938 noop.save().unwrap();
5939 let mut winner = run_state(RunStatus::Ready);
5940 winner.id = "20260101-000000-sup8".to_owned();
5941 winner.save().unwrap();
5942
5943 let mut t = task();
5944 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5945 t.status = TaskStatus::Done;
5946
5947 supersede_prior_runs(&t, &crate::run::home());
5948
5949 assert_eq!(
5950 RunState::load(&failed.id).unwrap().status,
5951 RunStatus::Failed
5952 );
5953 assert_eq!(
5954 RunState::load(&noop.id).unwrap().status,
5955 RunStatus::VerifiedNoop
5956 );
5957 }
5958
5959 fn approval_question(run: &str) -> ask::Question {
5960 ask::Question::new(
5961 run.to_owned(),
5962 land::APPROVAL_NODE.to_owned(),
5963 "land".to_owned(),
5964 "merge?".to_owned(),
5965 String::new(),
5966 vec!["merge".to_owned(), "hold".to_owned()],
5967 )
5968 }
5969
5970 #[test]
5971 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5972 crate::run::pin_test_home();
5973 let mut state = run_state(RunStatus::Landing);
5974 state.id = "20260101-000000-fre1".to_owned();
5975 state.parked = true;
5976 state.save().unwrap();
5977 ask::Questions::open()
5978 .put(&mut approval_question(&state.id))
5979 .unwrap();
5980
5981 let mut t = task();
5982 t.runs.push(state.id.clone());
5983 assert_eq!(
5984 land_resume_state(&t),
5985 LandResume::StillWaiting,
5986 "nobody has answered and the timeout has not passed"
5987 );
5988 }
5989
5990 #[test]
5991 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5992 crate::run::pin_test_home();
5997 let mut state = run_state(RunStatus::Landing);
5998 state.id = "20260101-000000-exp1".to_owned();
5999 state.parked = true;
6000 state.config.graph.answer_timeout = 60;
6001 state.save().unwrap();
6002
6003 let store = ask::Questions::open();
6004 let mut q = approval_question(&state.id);
6005 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
6006 store.put(&mut q).unwrap();
6007
6008 let mut t = task();
6009 t.runs.push(state.id.clone());
6010 assert_eq!(
6011 land_resume_state(&t),
6012 LandResume::Ready,
6013 "an expired question must not be waited on forever"
6014 );
6015
6016 let after = store.get(&q.id).unwrap();
6017 assert!(
6018 !after.status.open(),
6019 "the question is abandoned, not silently ignored"
6020 );
6021 assert!(
6022 after.resolution().is_none(),
6023 "an abandoned question is not read as a decision"
6024 );
6025 }
6026
6027 #[test]
6028 fn reclaim_settles_a_running_task_against_its_last_run() {
6029 let mut t = task();
6030 t.start("20260904-000000-4043".to_owned());
6031 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
6032 assert_eq!(
6033 t.status,
6034 TaskStatus::Done,
6035 "a run that actually finished must not stay `running` forever"
6036 );
6037 }
6038
6039 #[test]
6040 fn reclaim_maps_a_parked_run_to_parked_without_touching_last_error() {
6041 let mut t = task();
6042 t.start("20260904-000000-4043".to_owned());
6043 t.last_error = Some("earlier trouble".to_owned());
6044 let mut state = run_state(RunStatus::Judging);
6045 state.parked = true;
6046 reclaim(&mut t, Some(state), 1, "en");
6047 assert_eq!(t.status, TaskStatus::Parked);
6048 assert_eq!(t.attempts, 0, "the park is refunded, even at max_attempts");
6049 assert_eq!(t.last_error.as_deref(), Some("earlier trouble"));
6050 assert!(t.park_reason.is_some());
6051 assert_eq!(t.runs, ["20260904-000000-4043"], "the same run is kept");
6052 assert!(!t.fresh_start);
6053 assert!(t.status.runnable());
6054 }
6055
6056 #[test]
6057 fn settle_parks_a_task_without_a_failure_and_a_new_start_clears_the_reason() {
6058 let mut t = task();
6059 t.start("20260904-000000-4043".to_owned());
6060 t.last_error = Some("earlier trouble".to_owned());
6061 settle(
6062 &mut t,
6063 Verdict {
6064 status: RunStatus::Judging,
6065 left_pr: false,
6066 quota_hit: false,
6067 parked: true,
6068 no_viable_candidates: false,
6069 },
6070 "parked after `judging`",
6071 1,
6072 );
6073 assert_eq!(t.status, TaskStatus::Parked);
6074 assert_eq!(t.attempts, 0);
6075 assert_eq!(t.last_error.as_deref(), Some("earlier trouble"));
6076 assert_eq!(t.park_reason.as_deref(), Some("parked after `judging`"));
6077 assert_eq!(t.schema, crate::queue::SCHEMA);
6078 assert!(
6079 crate::notices::task_held(&t).is_none(),
6080 "a park is not news"
6081 );
6082 t.start("20260904-000000-4043".to_owned());
6083 assert_eq!(t.park_reason, None);
6084 }
6085
6086 #[test]
6087 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
6088 let mut t = task();
6092 t.start("20260904-000000-4043".to_owned());
6093 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
6094 assert_eq!(t.status, TaskStatus::Failed);
6095 assert!(t.status.runnable());
6096 }
6097
6098 #[test]
6099 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
6100 let mut t = task();
6101 t.start("20260904-000000-4043".to_owned());
6102 reclaim(&mut t, None, 2, "en");
6103 assert_eq!(t.status, TaskStatus::Held);
6104 assert!(
6105 t.last_error
6106 .as_deref()
6107 .is_some_and(|e| e.contains("running")),
6108 "the operator needs to know why this task was held"
6109 );
6110 }
6111
6112 #[test]
6113 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
6114 let dir = tempfile::tempdir().unwrap();
6115 let queue = Queue::at(dir.path().to_path_buf());
6116
6117 let mut orphaned = task();
6119 orphaned.id = "20260904-000000-orph".to_owned();
6120 orphaned.status = TaskStatus::Running;
6121 orphaned.attempts = 1;
6122 queue.put(&mut orphaned).unwrap();
6123
6124 let mut alive = task();
6125 alive.id = "20260904-000000-live".to_owned();
6126 alive.status = TaskStatus::Running;
6127 alive.attempts = 1;
6128 queue.put(&mut alive).unwrap();
6129 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
6130
6131 let mut queued = task();
6132 queued.id = "20260904-000000-wait".to_owned();
6133 queue.put(&mut queued).unwrap();
6134
6135 let reclaimed = reclaim_orphaned_running(&queue, 2);
6136 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
6137
6138 assert_eq!(
6139 queue.get(&orphaned.id).unwrap().status,
6140 TaskStatus::Held,
6141 "nothing was driving it and there was no run to recover"
6142 );
6143 assert_eq!(
6144 queue.get(&alive.id).unwrap().status,
6145 TaskStatus::Running,
6146 "a live claim must protect the task it belongs to"
6147 );
6148 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
6149 }
6150
6151 fn read_run_under(home: &Path, id: &str) -> RunState {
6157 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
6158 serde_json::from_str(&body).unwrap()
6159 }
6160
6161 #[test]
6162 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
6163 let dir = tempfile::tempdir().unwrap();
6164 let home = dir.path().to_path_buf();
6165 let now = Timestamp::now();
6166 let overrun_seat = || crate::run::ActiveSeat {
6167 node: "implement".to_owned(),
6168 started_at: now - jiff::SignedDuration::new(21_000, 0),
6169 timeout_secs: 3_600,
6170 attempt: 0,
6171 task: None,
6172 command: None,
6173 index: None,
6174 total: None,
6175 };
6176
6177 let mut dead = run_state(RunStatus::Implementing);
6178 dead.id = "20260101-000000-dead".to_owned();
6179 dead.active.insert("impl-A".to_owned(), overrun_seat());
6180 dead.driver_pid = Some(4242);
6183 dead.save_under(&home).unwrap();
6184
6185 let mut alive = run_state(RunStatus::Implementing);
6188 alive.id = "20260101-000000-aliv".to_owned();
6189 alive.active.insert("impl-A".to_owned(), overrun_seat());
6190 alive.save_under(&home).unwrap();
6191 let mut status = Status::new();
6192 status.current = vec![Current {
6193 task: "20260101-000000-task".to_owned(),
6194 run: alive.id.clone(),
6195 }];
6196 write_status_to(&home.join("daemon.json"), &status).unwrap();
6197
6198 let questions = Questions::at(home.join("questions"));
6202 let mut q = ask::Question::new(
6203 dead.id.clone(),
6204 "implement".to_owned(),
6205 "impl-A".to_owned(),
6206 "Which storage backend?".to_owned(),
6207 String::new(),
6208 vec!["SQLite".to_owned(), "Redis".to_owned()],
6209 );
6210 questions.put(&mut q).unwrap();
6211
6212 let abandoned = reclaim_abandoned_runs_with(
6213 &home,
6214 now,
6215 |pid| if pid == 4242 { Some(false) } else { None },
6216 |_| panic!("a query answering Dead outright needs no identity corroboration"),
6217 );
6218 assert_eq!(abandoned, vec![dead.id.clone()]);
6219
6220 let reloaded = read_run_under(&home, &dead.id);
6221 assert_eq!(reloaded.status, RunStatus::Failed);
6222 assert!(reloaded.active.is_empty());
6223 assert!(
6224 !questions.get(&q.id).unwrap().status.open(),
6225 "the failed run's own open question must be settled in the same pass"
6226 );
6227
6228 let still_alive = read_run_under(&home, &alive.id);
6229 assert_eq!(
6230 still_alive.status,
6231 RunStatus::Implementing,
6232 "a live daemon's claim protects it"
6233 );
6234 assert!(!still_alive.active.is_empty());
6235 }
6236
6237 #[test]
6247 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
6248 let dir = tempfile::tempdir().unwrap();
6249 let home = dir.path().to_path_buf();
6250 let now = Timestamp::now();
6251
6252 let mut manual = run_state(RunStatus::Reviewing);
6253 manual.id = "20260101-000000-manl".to_owned();
6254 manual.active.insert(
6255 "review-1".to_owned(),
6256 crate::run::ActiveSeat {
6257 node: "review".to_owned(),
6258 started_at: now - jiff::SignedDuration::new(21_000, 0),
6259 timeout_secs: 3_600,
6260 attempt: 0,
6261 task: None,
6262 command: None,
6263 index: None,
6264 total: None,
6265 },
6266 );
6267 manual.driver_pid = Some(4242);
6271 manual.driver_started_at = Some("1790000000".to_owned());
6272 manual.save_under(&home).unwrap();
6273
6274 let abandoned = reclaim_abandoned_runs_with(
6275 &home,
6276 now,
6277 |pid| if pid == 4242 { Some(true) } else { None },
6278 |pid| {
6279 if pid == 4242 {
6280 Some("1790000000".to_owned())
6281 } else {
6282 None
6283 }
6284 },
6285 );
6286 assert!(
6287 abandoned.is_empty(),
6288 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
6289 );
6290
6291 let reloaded = read_run_under(&home, &manual.id);
6292 assert_eq!(reloaded.status, RunStatus::Reviewing);
6293 assert!(!reloaded.active.is_empty());
6294 }
6295
6296 #[test]
6297 fn an_already_claimed_task_is_skipped_rather_than_failed() {
6298 let dir = tempfile::tempdir().unwrap();
6299 let queue = Queue::at(dir.path().to_path_buf());
6300 let mut only = task();
6301 queue.put(&mut only).unwrap();
6302
6303 let _elsewhere = queue.claim(&only.id).unwrap();
6304 let candidates = runnable(&queue);
6305 assert_eq!(candidates.len(), 1, "the task is still runnable");
6306 assert!(
6307 queue.claim(&candidates[0].id).is_err(),
6308 "the loop cannot take a claim somebody else holds"
6309 );
6310
6311 let after = queue.get(&only.id).unwrap();
6312 assert_eq!(after.status, TaskStatus::Queued);
6313 assert_eq!(
6314 after.attempts, 0,
6315 "losing the race is not an attempt at the task"
6316 );
6317 assert_eq!(after.last_error, None);
6318 }
6319
6320 #[test]
6321 fn the_status_file_round_trips_and_its_heartbeat_advances() {
6322 let dir = tempfile::tempdir().unwrap();
6323 let path = dir.path().join("daemon.json");
6324
6325 let mut status = Status::new();
6326 status.idle = false;
6327 status.completed = 7;
6328 status.current = vec![Current {
6329 task: "20260902-000000-t111".to_owned(),
6330 run: "20260902-000001-r111".to_owned(),
6331 }];
6332 write_status_to(&path, &status).unwrap();
6333 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6334 assert_eq!(first.schema, SCHEMA);
6335 assert_eq!(first.pid, std::process::id());
6336 assert!(!first.idle);
6337 assert_eq!(first.completed, 7);
6338 assert_eq!(first.current, status.current);
6339 assert!(
6340 !path.with_extension("json.tmp").exists(),
6341 "the temp file is renamed, not left behind"
6342 );
6343
6344 std::thread::sleep(Duration::from_millis(5));
6345 status.updated_at = Timestamp::now();
6346 status.polls = 3;
6347 write_status_to(&path, &status).unwrap();
6348 let second: Status =
6349 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6350 assert!(
6351 second.updated_at > first.updated_at,
6352 "a reader can only detect staleness if the heartbeat moves"
6353 );
6354 assert_eq!(
6355 second.started_at, first.started_at,
6356 "the start time is not a heartbeat"
6357 );
6358 assert_eq!(second.polls, 3);
6359 }
6360
6361 #[test]
6362 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
6363 let dir = tempfile::tempdir().unwrap();
6364
6365 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
6366
6367 let mut status = Status::new();
6368 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
6369 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6370 let stale = read_status(dir.path()).unwrap();
6371 assert!(
6372 !stale.running(Timestamp::now()),
6373 "a minute without a heartbeat is a dead daemon, not a busy one"
6374 );
6375 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
6376
6377 status.updated_at = Timestamp::now();
6378 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6379 let fresh = read_status(dir.path()).unwrap();
6380 assert!(fresh.running(Timestamp::now()));
6381 }
6382
6383 #[test]
6384 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
6385 let dir = tempfile::tempdir().unwrap();
6386 let now = Timestamp::now();
6387 let mine = "20260903-080619-01c2";
6388
6389 assert!(
6390 !is_working_on(dir.path(), mine, now),
6391 "no status file means nobody is working on anything"
6392 );
6393
6394 let mut status = Status::new();
6395 status.current = vec![Current {
6396 task: "20260903-080340-0167".to_owned(),
6397 run: mine.to_owned(),
6398 }];
6399 status.updated_at = now;
6400 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6401 assert!(is_working_on(dir.path(), mine, now));
6402 assert!(
6403 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
6404 "a daemon busy with one run is not working on another"
6405 );
6406
6407 status.updated_at = now - jiff::SignedDuration::from_secs(600);
6410 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6411 assert!(
6412 !is_working_on(dir.path(), mine, now),
6413 "a stale heartbeat is a dead daemon, so its run is a leftover"
6414 );
6415 }
6416
6417 #[test]
6418 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
6419 let dir = tempfile::tempdir().unwrap();
6420 let now = Timestamp::now();
6421
6422 assert!(
6423 !is_working_on_short(dir.path(), "01c2", now),
6424 "no status file means nobody is working on anything"
6425 );
6426
6427 let mut status = Status::new();
6428 status.current = vec![Current {
6429 task: "20260903-080340-0167".to_owned(),
6430 run: "20260903-080619-01c2".to_owned(),
6431 }];
6432 status.updated_at = now;
6433 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6434 assert!(
6435 is_working_on_short(dir.path(), "01c2", now),
6436 "the run's short id is the last block of its full id"
6437 );
6438 assert!(
6439 !is_working_on_short(dir.path(), "3cbf", now),
6440 "a daemon busy with one worktree bay is not working on another"
6441 );
6442 }
6443
6444 #[test]
6445 fn a_newer_status_file_still_yields_a_reading() {
6446 let dir = tempfile::tempdir().unwrap();
6447 std::fs::write(
6450 dir.path().join("daemon.json"),
6451 serde_json::json!({
6452 "schema": 2,
6453 "updated_at": Timestamp::now().to_string(),
6454 "idle": true,
6455 "surprise": { "nested": [1, 2, 3] },
6456 })
6457 .to_string(),
6458 )
6459 .unwrap();
6460
6461 let reading = read_status(dir.path()).expect("a forward-compatible read");
6462 assert!(reading.running(Timestamp::now()));
6463 assert!(reading.idle);
6464 assert!(reading.current.is_empty());
6465 }
6466
6467 #[test]
6468 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
6469 let dir = tempfile::tempdir().unwrap();
6475 std::fs::write(
6476 dir.path().join("daemon.json"),
6477 serde_json::json!({
6478 "schema": 1,
6479 "pid": 4242,
6480 "updated_at": Timestamp::now().to_string(),
6481 "idle": false,
6482 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
6483 "completed": 3,
6484 "polls": 9,
6485 })
6486 .to_string(),
6487 )
6488 .unwrap();
6489
6490 let reading = read_status(dir.path()).expect("an older shape must still parse");
6491 assert!(reading.running(Timestamp::now()));
6492 assert_eq!(
6493 reading.current,
6494 vec![Current {
6495 task: "20260902-140501-aaaa".to_owned(),
6496 run: "20260902-140502-bbbb".to_owned(),
6497 }]
6498 );
6499 }
6500
6501 #[test]
6502 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
6503 let dir = tempfile::tempdir().unwrap();
6504 std::fs::write(
6505 dir.path().join("daemon.json"),
6506 serde_json::json!({
6507 "schema": 1,
6508 "updated_at": Timestamp::now().to_string(),
6509 "idle": true,
6510 "current": null,
6511 })
6512 .to_string(),
6513 )
6514 .unwrap();
6515 let with_null = read_status(dir.path()).expect("null must still parse");
6516 assert!(with_null.current.is_empty());
6517
6518 std::fs::write(
6519 dir.path().join("daemon.json"),
6520 serde_json::json!({
6521 "schema": 1,
6522 "updated_at": Timestamp::now().to_string(),
6523 "idle": true,
6524 })
6525 .to_string(),
6526 )
6527 .unwrap();
6528 let absent = read_status(dir.path()).expect("a missing field must still parse");
6529 assert!(absent.current.is_empty());
6530 }
6531
6532 #[test]
6533 fn a_task_without_a_repository_runs_in_the_daemons_default() {
6534 let fallback = Path::new("/default");
6535 let mut blank = task();
6536 blank.repo = PathBuf::new();
6537 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
6538 let mut dot = task();
6539 dot.repo = PathBuf::from(".");
6540 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
6541 assert_eq!(
6542 repo_for(&task(), fallback),
6543 PathBuf::from("/repo"),
6544 "a task that names a repository keeps it"
6545 );
6546 }
6547
6548 #[test]
6549 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
6550 let mut solo_cfg = Config::default();
6556 solo_cfg.graph.implementers = 3;
6557 let mut solo_task = task();
6558 solo_task.solo = true;
6559 apply_solo(&mut solo_cfg, &solo_task);
6560 assert_eq!(solo_cfg.graph.implementers, 1);
6561
6562 let mut plain_cfg = Config::default();
6563 plain_cfg.graph.implementers = 3;
6564 let plain_task = task();
6565 assert!(!plain_task.solo);
6566 apply_solo(&mut plain_cfg, &plain_task);
6567 assert_eq!(
6568 plain_cfg.graph.implementers, 3,
6569 "a task that did not ask to run alone keeps the config's candidates"
6570 );
6571 }
6572
6573 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
6574 QuotaLoss {
6575 seat: seat.into(),
6576 node: "judge".into(),
6577 at: at.parse().unwrap(),
6578 reset: reset.map(str::to_string),
6579 }
6580 }
6581
6582 #[test]
6583 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
6584 let old: Vec<QuotaLoss> = (1..=4)
6585 .map(|i| {
6586 loss(
6587 &format!("judge-{i}"),
6588 "2026-09-23T05:23:00Z",
6589 Some("2:40pm (Asia/Tokyo)"),
6590 )
6591 })
6592 .collect();
6593 let fresh = losses_this_attempt(&old, &old);
6594 assert!(fresh.is_empty());
6595 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
6596 }
6598
6599 #[test]
6600 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6601 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6602 let now = Timestamp::now();
6603 let mut after = old.clone();
6604 after.push(loss("judge-2", &now.to_string(), None));
6605 let fresh = losses_this_attempt(&old, &after);
6606 assert_eq!(fresh, vec![after[1].clone()]);
6607 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6608 assert_eq!(
6609 until,
6610 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6611 );
6612 }
6613
6614 #[test]
6615 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6616 let before = vec![
6619 loss("judge-1", "2026-09-23T05:23:00Z", None),
6620 loss("judge-2", "2026-09-23T05:24:00Z", None),
6621 ];
6622 let after = vec![
6623 loss("judge-2", "2026-09-23T05:24:00Z", None),
6624 loss("judge-1", "2026-09-24T01:00:00Z", None),
6625 ];
6626 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6627 }
6628
6629 #[test]
6630 fn merge_overrides_are_parsed_or_refused() {
6631 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6632 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6633 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6634 assert!(merge_mode("squash").is_err());
6635 }
6636
6637 #[test]
6638 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6639 let now = Timestamp::now();
6640 let fallback = Duration::from_secs(300);
6641 let cap = Duration::from_secs(1800);
6642
6643 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6645
6646 let soon = now + jiff::SignedDuration::from_secs(600);
6648 assert_eq!(
6649 quota_wait(Some(soon), now, fallback, cap),
6650 Duration::from_secs(600)
6651 );
6652
6653 let past = now - jiff::SignedDuration::from_secs(60);
6656 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6657
6658 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6661 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6662 }
6663
6664 #[test]
6665 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6666 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6667
6668 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6669 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6670
6671 let already_past =
6675 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6676 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6677
6678 assert!(
6679 parse_reset_hint("session limit reached", now, now).is_none(),
6680 "free text with no recognised shape is not guessed at"
6681 );
6682 assert!(
6683 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6684 "an unresolvable zone name is not guessed at either"
6685 );
6686 }
6687
6688 #[test]
6689 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6690 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6691
6692 let at = parse_reset_hint(
6693 "You've hit your usage limit. Visit \
6694 https://chatgpt.com/codex/settings/usage to purchase more \
6695 credits or try again at Sep 19th, 2026 5:10 PM.",
6696 now,
6697 now,
6698 )
6699 .expect("the codex reset wording is a recognised shape");
6700 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6701
6702 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6707 .expect("an explicit year needs no rollover");
6708 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6709
6710 assert!(
6711 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6712 "a two-digit year is not the documented shape and is not guessed at"
6713 );
6714 assert!(
6715 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6716 "a four-letter month name is not the documented three-letter abbreviation"
6717 );
6718 assert!(
6719 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6720 "an explicit zone on the dated shape is a format nobody has \
6721 documented, and is refused rather than guessed at as UTC"
6722 );
6723 }
6724
6725 #[test]
6726 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6727 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6728 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6729
6730 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6731 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6732
6733 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6734 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6735
6736 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6737 assert!(
6738 parse_reset_hint(bad, now, recorded).is_none(),
6739 "{bad:?} must not be guessed at"
6740 );
6741 }
6742 }
6743
6744 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6748 let config = dir.join("magi.toml");
6749 std::fs::write(
6750 &config,
6751 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6752 )
6753 .unwrap();
6754 let opts = Opts {
6755 poll: Duration::from_secs(30),
6756 config: Some(config),
6757 repo: dir.join("repo"),
6761 ..Opts::default()
6762 };
6763 let home = dir.join("home");
6772 let worktrees = dir.join("wt");
6773 (
6774 opts,
6775 Queue::at(dir.join("queue")),
6776 home.join("daemon.json"),
6777 home,
6778 worktrees,
6779 )
6780 }
6781
6782 #[test]
6783 fn a_stop_is_idempotent_and_once_set_stays_set() {
6784 let stop = Stop::new();
6785 assert!(!stop.stopped());
6786
6787 stop.stop();
6788 assert!(stop.stopped());
6789 stop.stop();
6790 assert!(stop.stopped(), "a second stop is not a toggle");
6791
6792 let shared = stop.clone();
6793 assert!(
6794 shared.stopped(),
6795 "a clone is the same stop; that is how the loop and its caller share one"
6796 );
6797 }
6798
6799 #[test]
6800 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6801 let stop = Stop::new();
6802 stop.enter();
6803 assert!(
6804 !stop.finishing(),
6805 "a busy loop nobody has asked to stop is just running"
6806 );
6807
6808 stop.stop();
6809 assert!(
6810 stop.finishing(),
6811 "a stop asked for mid-run has not landed until the run is settled"
6812 );
6813
6814 stop.exit();
6815 assert!(
6816 !stop.finishing(),
6817 "once the run is settled the stop has landed and there is nothing to finish"
6818 );
6819 }
6820
6821 #[test]
6822 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6823 let stop = Stop::new();
6824 stop.enter();
6825 stop.enter();
6826 stop.stop();
6827 assert!(stop.finishing(), "two runs still in flight");
6828
6829 stop.exit();
6830 assert!(
6831 stop.finishing(),
6832 "one run finished, but a sibling is still working"
6833 );
6834
6835 stop.exit();
6836 assert!(
6837 !stop.finishing(),
6838 "the last run out is what actually lands the stop"
6839 );
6840 }
6841
6842 #[tokio::test]
6843 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6844 let dir = tempfile::tempdir().unwrap();
6845 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6846 let stop = Stop::new();
6847 stop.stop();
6848
6849 let began = std::time::Instant::now();
6850 tokio::time::timeout(
6851 Duration::from_secs(2),
6852 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6853 )
6854 .await
6855 .expect("a stopped loop must return, not sit out its poll interval")
6856 .expect("the loop's own setup and teardown must not fail");
6857 assert!(
6858 began.elapsed() < opts.poll,
6859 "returned only after {:?}, which is a poll interval, not a stop",
6860 began.elapsed()
6861 );
6862 }
6863
6864 #[tokio::test]
6865 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6866 let dir = tempfile::tempdir().unwrap();
6867 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6868 let stop = Stop::new();
6869
6870 let asker = {
6873 let stop = stop.clone();
6874 tokio::spawn(async move {
6875 tokio::time::sleep(Duration::from_millis(20)).await;
6876 stop.stop();
6877 })
6878 };
6879
6880 let began = std::time::Instant::now();
6881 tokio::time::timeout(
6882 Duration::from_secs(2),
6883 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6884 )
6885 .await
6886 .expect("a stop asked for while idle must wake the wait")
6887 .expect("the loop's own setup and teardown must not fail");
6888 asker.await.unwrap();
6889 assert!(
6890 began.elapsed() < opts.poll,
6891 "returned only after {:?}, so the stop waited on the sleep",
6892 began.elapsed()
6893 );
6894 }
6895
6896 #[tokio::test]
6897 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6898 let dir = tempfile::tempdir().unwrap();
6899 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6900 let stop = Stop::new();
6901 stop.stop();
6902
6903 tokio::time::timeout(
6904 Duration::from_secs(2),
6905 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6906 )
6907 .await
6908 .expect("a stopped loop must return")
6909 .expect("the loop's own setup and teardown must not fail");
6910
6911 assert!(
6912 home.is_dir(),
6913 "the loop did publish a status file, so its removal is the teardown and not an absence"
6914 );
6915 assert!(
6916 !status_file.exists(),
6917 "a stopped loop clears its status file"
6918 );
6919 assert!(
6920 read_status(&home).is_none(),
6921 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6922 );
6923 }
6924
6925 #[tokio::test]
6926 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6927 let dir = tempfile::tempdir().unwrap();
6928 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6929 opts.once = true;
6930
6931 let mut settled = RunState::new(
6932 dir.path().join("repo"),
6933 "main".to_owned(),
6934 "abc1234".to_owned(),
6935 "fixture".to_owned(),
6936 Config::default(),
6937 );
6938 settled.status = RunStatus::Ready;
6939 let run_dir = home.join("runs").join(&settled.id);
6940 std::fs::create_dir_all(&run_dir).unwrap();
6941 std::fs::write(
6942 run_dir.join("run.json"),
6943 serde_json::to_string_pretty(&settled).unwrap(),
6944 )
6945 .unwrap();
6946 let questions = Questions::at(home.join("questions"));
6947 let mut question = ask::Question::new(
6948 settled.id.clone(),
6949 "review".to_owned(),
6950 "reviewer-1".to_owned(),
6951 "Continue?".to_owned(),
6952 String::new(),
6953 Vec::new(),
6954 );
6955 questions.put(&mut question).unwrap();
6956
6957 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6958 .await
6959 .unwrap();
6960
6961 assert_eq!(
6962 questions.get(&question.id).unwrap().status,
6963 ask::QuestionStatus::Abandoned,
6964 "an empty --once drain still performs startup question cleanup"
6965 );
6966 }
6967
6968 #[test]
6969 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6970 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6971
6972 assert!(
6973 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6974 "never checked before: due at once"
6975 );
6976
6977 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6978 assert!(
6979 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6980 "well inside the interval: not due yet"
6981 );
6982
6983 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6984 assert!(
6985 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6986 "exactly at the edge: not yet due, same convention as `clean::due`"
6987 );
6988
6989 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6990 assert!(
6991 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6992 "past the interval: due again"
6993 );
6994 }
6995
6996 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
7001 let config = dir.join("magi.toml");
7002 std::fs::write(
7008 &config,
7009 format!(
7010 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
7011 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
7012 cache_dir.display()
7013 ),
7014 )
7015 .unwrap();
7016 Opts {
7017 config: Some(config),
7018 repo: dir.join("repo"),
7019 ..Opts::default()
7020 }
7021 }
7022
7023 #[tokio::test]
7024 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
7025 let dir = tempfile::tempdir().unwrap();
7026 let home = dir.path().join("home");
7027 let cache_dir = dir.path().join("cache");
7028 std::fs::create_dir_all(&cache_dir).unwrap();
7029 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
7030 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
7031
7032 let running = Stop::new();
7035 let mut last_checked = None;
7036 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
7037 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
7038 .await;
7039 assert_eq!(
7040 crate::disk::dir_size(&cache_dir),
7041 0,
7042 "over the cap on the first check ever: pruned at once, no idle queue required"
7043 );
7044 assert_eq!(last_checked, Some(t0));
7045
7046 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
7048 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
7049 maybe_prune_cache_between_runs(
7050 &opts.repo,
7051 &opts,
7052 &home,
7053 &running,
7054 &mut last_checked,
7055 too_soon,
7056 )
7057 .await;
7058 assert_eq!(
7059 crate::disk::dir_size(&cache_dir),
7060 10,
7061 "too soon since the last check: left alone rather than rescanned every call"
7062 );
7063 assert_eq!(
7064 last_checked,
7065 Some(t0),
7066 "an idle check does not reset the clock"
7067 );
7068
7069 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
7071 maybe_prune_cache_between_runs(
7072 &opts.repo,
7073 &opts,
7074 &home,
7075 &running,
7076 &mut last_checked,
7077 due_again,
7078 )
7079 .await;
7080 assert_eq!(
7081 crate::disk::dir_size(&cache_dir),
7082 0,
7083 "due again: pruned back under the cap"
7084 );
7085 }
7086
7087 #[tokio::test]
7095 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
7096 let dir = tempfile::tempdir().unwrap();
7097 let home = dir.path().join("home");
7098 let cache_dir = dir.path().join("cache");
7099 std::fs::create_dir_all(&cache_dir).unwrap();
7100 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
7101 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
7102
7103 let stop = Stop::new();
7104 stop.stop();
7105 assert!(
7106 !stop.finishing(),
7107 "no run is in flight at a between-runs boundary, so nothing else \
7108 would tell the operator this stop had not taken effect yet"
7109 );
7110
7111 let mut last_checked = None;
7112 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
7113 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
7114 .await;
7115 assert_eq!(
7116 crate::disk::dir_size(&cache_dir),
7117 10,
7118 "over its cap, and due for the first check ever, but a stop outranks \
7119 it: the cap is a standing policy the next start measures again"
7120 );
7121 assert_eq!(
7122 last_checked, None,
7123 "a check that never happened must not claim the interval"
7124 );
7125 }
7126
7127 #[tokio::test]
7141 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
7142 let dir = tempfile::tempdir().unwrap();
7143 let cache_dir = dir.path().join("cache");
7144 std::fs::create_dir_all(&cache_dir).unwrap();
7145 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
7146
7147 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
7148 opts.poll = Duration::from_millis(20);
7149 opts.max_attempts = 1_000;
7150
7151 let queue = Queue::at(dir.path().join("queue"));
7152 let mut t = Task::new(
7153 "x".to_owned(),
7154 "x".to_owned(),
7155 opts.repo.clone(),
7156 Source::Human,
7157 );
7158 queue.put(&mut t).unwrap();
7159
7160 let home = dir.path().join("home");
7161 let worktrees = dir.path().join("wt");
7162 let status_file = home.join("daemon.json");
7163 let stop = Stop::new();
7164 let stopper = {
7165 let stop = stop.clone();
7166 let queue = Queue::at(dir.path().join("queue"));
7167 let id = t.id.clone();
7168 tokio::spawn(async move {
7169 tokio::time::sleep(Duration::from_millis(400)).await;
7173 for _ in 0..200 {
7174 if queue.get(&id).is_ok_and(|t| t.attempts >= 2) {
7175 break;
7176 }
7177 tokio::time::sleep(Duration::from_millis(50)).await;
7178 }
7179 stop.stop();
7180 })
7181 };
7182
7183 tokio::time::timeout(
7184 Duration::from_secs(20),
7185 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
7186 )
7187 .await
7188 .expect("the loop must not hang on a queue that keeps producing failing work")
7189 .expect("the loop's own setup and teardown must not fail");
7190 stopper.await.unwrap();
7191
7192 let after = queue.get(&t.id).unwrap();
7193 assert!(
7194 after.attempts >= 2,
7195 "the harness must actually have retried more than once, or this is not \
7196 exercising a busy queue at all (got {} attempt(s))",
7197 after.attempts
7198 );
7199 assert!(
7200 after.status.runnable(),
7201 "still under its attempt budget: the queue never reached a natural idle \
7202 on its own, only the external stop ended the test"
7203 );
7204
7205 assert_eq!(
7206 crate::disk::dir_size(&cache_dir),
7207 0,
7208 "an oversized cache must not be left to grow unboundedly just because the \
7209 queue kept the loop busy the whole time"
7210 );
7211 }
7212
7213 #[test]
7214 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
7215 let dir = tempfile::tempdir().unwrap();
7216 let queue = Queue::at(dir.path().join("queue"));
7217 let questions = Questions::at(dir.path().join("questions"));
7218 let mut task = task();
7219 queue.put(&mut task).unwrap();
7220
7221 let mut task_question = ask::Question::new(
7222 task.id.clone(),
7223 crate::conduct::NODE.to_owned(),
7224 "conduct".to_owned(),
7225 "Which backend?".to_owned(),
7226 String::new(),
7227 Vec::new(),
7228 );
7229 questions.put(&mut task_question).unwrap();
7230 task.block(vec![task_question.id.clone()], None);
7231 queue.put(&mut task).unwrap();
7232
7233 let mut run_question = ask::Question::new(
7234 "20260101-000000-run1".to_owned(),
7235 "review".to_owned(),
7236 "reviewer-1".to_owned(),
7237 "Run question".to_owned(),
7238 String::new(),
7239 Vec::new(),
7240 );
7241 questions.put(&mut run_question).unwrap();
7242
7243 let mut coincidental = ask::Question::new(
7248 task.id.clone(),
7249 "review".to_owned(),
7250 "reviewer-1".to_owned(),
7251 "Unrelated review question".to_owned(),
7252 String::new(),
7253 Vec::new(),
7254 );
7255 questions.put(&mut coincidental).unwrap();
7256
7257 reconcile_task_questions(&queue, &questions);
7258 assert!(questions.get(&task_question.id).unwrap().status.open());
7259 assert!(questions.get(&run_question.id).unwrap().status.open());
7260 assert!(questions.get(&coincidental.id).unwrap().status.open());
7261
7262 task.release();
7263 queue.put(&mut task).unwrap();
7264 reconcile_task_questions(&queue, &questions);
7265 assert_eq!(
7266 questions.get(&task_question.id).unwrap().status,
7267 ask::QuestionStatus::Abandoned
7268 );
7269 assert!(
7270 questions.get(&run_question.id).unwrap().status.open(),
7271 "run questions remain the run janitor's responsibility"
7272 );
7273 assert!(
7274 questions.get(&coincidental.id).unwrap().status.open(),
7275 "a non-conductor question must not be abandoned just because its \
7276 run id coincides with a task id"
7277 );
7278 }
7279
7280 #[test]
7281 fn a_freshly_started_running_task_is_never_stalled() {
7282 let dir = tempfile::tempdir().unwrap();
7283 let mut t = task();
7284 t.start("run-1".to_owned());
7285 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
7288 }
7289
7290 #[test]
7291 fn a_long_running_task_with_no_live_daemon_is_stalled() {
7292 let dir = tempfile::tempdir().unwrap();
7293 let mut t = task();
7294 t.start("run-1".to_owned());
7295 t.updated_at = Timestamp::now()
7296 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
7297 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
7298 assert_eq!(
7299 stalled_tasks(
7300 &Queue::at(dir.path().join("q")),
7301 dir.path(),
7302 Timestamp::now()
7303 )
7304 .len(),
7305 0,
7306 "the task was never written to this queue"
7307 );
7308 }
7309
7310 #[test]
7311 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
7312 let dir = tempfile::tempdir().unwrap();
7313 let mut t = task();
7314 t.id = "20260903-080340-0167".to_owned();
7315 t.start("20260903-080619-01c2".to_owned());
7316 t.updated_at = Timestamp::now()
7317 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
7318
7319 let mut status = Status::new();
7320 status.current = vec![Current {
7321 task: t.id.clone(),
7322 run: "20260903-080619-01c2".to_owned(),
7323 }];
7324 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
7325
7326 assert!(
7327 !is_stalled(&t, dir.path(), Timestamp::now()),
7328 "a live daemon's own heartbeat rules out stalled, however long the task has run"
7329 );
7330 }
7331
7332 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
7336 let path = queue.path_of(id);
7337 let body = std::fs::read_to_string(&path).unwrap();
7338 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
7339 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
7340 v["updated_at"] = serde_json::Value::String(old.to_string());
7341 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
7342 }
7343
7344 #[test]
7345 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
7346 let dir = tempfile::tempdir().unwrap();
7359 let queue = Queue::at(dir.path().join("queue"));
7360 let home = dir.path().join("home");
7361
7362 let mut t = task();
7363 t.id = "20260101-000001-lock".to_owned();
7364 t.start("run-1".to_owned());
7365 queue.put(&mut t).unwrap();
7366 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7367 std::fs::write(
7368 dir.path().join("queue").join(format!("{}.lock", t.id)),
7369 "not a pid",
7370 )
7371 .unwrap();
7372
7373 let now = Timestamp::now();
7374 assert!(
7375 reclaim_orphaned_running(&queue, 2).is_empty(),
7376 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
7377 and reclaim must leave the task alone"
7378 );
7379 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
7380
7381 let stalled = stalled_tasks(&queue, &home, now);
7382 assert_eq!(
7383 stalled.len(),
7384 1,
7385 "reclaim's inability to claim it yet must not hide it from the conductor"
7386 );
7387 assert_eq!(stalled[0].id, t.id);
7388 }
7389
7390 #[test]
7391 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
7392 let dir = tempfile::tempdir().unwrap();
7393 crate::run::set_home(dir.path().join("run-home"));
7394 let queue = Queue::at(dir.path().join("queue"));
7395 let home = dir.path().join("home");
7396 let questions = Questions::at(dir.path().join("questions"));
7397
7398 let mut t = task();
7399 t.id = "20260101-000003-dead".to_owned();
7400 t.start("missing-run".to_owned());
7401 queue.put(&mut t).unwrap();
7402 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7403
7404 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7407 assert_eq!(
7408 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
7409 [&t.id]
7410 );
7411 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
7412 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7413
7414 crate::conduct::apply(
7417 &queue,
7418 &questions,
7419 &crate::conduct::Verdict {
7420 decisions: vec![crate::conduct::Decision {
7421 id: t.id.clone(),
7422 recovery: Some(crate::conduct::Recovery::Requeue),
7423 ..crate::conduct::Decision::default()
7424 }],
7425 },
7426 )
7427 .unwrap();
7428 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7429 }
7430
7431 #[test]
7432 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
7433 let dir = tempfile::tempdir().unwrap();
7434 let queue = Queue::at(dir.path().join("queue"));
7435 let home = dir.path().join("home");
7436
7437 let mut fresh = task();
7438 fresh.id = "20260101-000001-aaaa".to_owned();
7439 fresh.start("run-1".to_owned());
7440 queue.put(&mut fresh).unwrap();
7441
7442 let mut old = task();
7443 old.id = "20260101-000002-bbbb".to_owned();
7444 old.start("run-2".to_owned());
7445 queue.put(&mut old).unwrap();
7446 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
7447
7448 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7449 assert_eq!(stalled.len(), 1);
7450 assert_eq!(stalled[0].id, old.id);
7451 }
7452
7453 #[test]
7454 fn queued_and_finished_task_views_partition_by_status() {
7455 let dir = tempfile::tempdir().unwrap();
7456 let queue = Queue::at(dir.path().join("queue"));
7457
7458 let mut queued = task();
7459 queued.id = "20260101-000001-aaaa".to_owned();
7460 queue.put(&mut queued).unwrap();
7461
7462 let mut failed = task();
7463 failed.id = "20260101-000002-bbbb".to_owned();
7464 failed.start("run-1".to_owned());
7465 failed.fail("gate red", 5);
7466 queue.put(&mut failed).unwrap();
7467
7468 let mut held = task();
7469 held.id = "20260101-000003-cccc".to_owned();
7470 held.hold_machine(None);
7471 queue.put(&mut held).unwrap();
7472
7473 let mut running = task();
7474 running.id = "20260101-000004-dddd".to_owned();
7475 running.start("run-2".to_owned());
7476 queue.put(&mut running).unwrap();
7477
7478 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
7479 assert_eq!(queued_ids, [queued.id.clone()]);
7480
7481 let mut finished_ids: Vec<String> =
7482 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
7483 finished_ids.sort_unstable();
7484 let mut want = vec![failed.id.clone(), held.id.clone()];
7485 want.sort_unstable();
7486 assert_eq!(finished_ids, want);
7487 }
7488
7489 #[test]
7490 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
7491 let dir = tempfile::tempdir().unwrap();
7492 let queue = Queue::at(dir.path().join("queue"));
7493 let questions = ask::Questions::at(dir.path().join("questions"));
7494
7495 let mut dep = task();
7496 dep.id = "20260101-000001-dep0".to_owned();
7497 dep.succeed();
7498 queue.put(&mut dep).unwrap();
7499
7500 let mut still_going = task();
7501 still_going.id = "20260101-000002-dep1".to_owned();
7502 queue.put(&mut still_going).unwrap();
7503
7504 let mut blocked = task();
7505 blocked.id = "20260101-000003-main".to_owned();
7506 blocked.block(
7507 vec![dep.id.clone(), still_going.id.clone()],
7508 Some("waits on both".to_owned()),
7509 );
7510 queue.put(&mut blocked).unwrap();
7511
7512 resolve_blockers(&queue, &questions);
7513
7514 let after = queue.get(&blocked.id).unwrap();
7515 assert_eq!(
7516 after.status,
7517 TaskStatus::Blocked,
7518 "one dependency is still outstanding"
7519 );
7520 assert_eq!(after.blocked_by, [still_going.id.clone()]);
7521 }
7522
7523 #[test]
7524 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
7525 let dir = tempfile::tempdir().unwrap();
7526 let queue = Queue::at(dir.path().join("queue"));
7527 let questions = ask::Questions::at(dir.path().join("questions"));
7528
7529 let mut q = crate::ask::Question::new(
7530 "20260101-000001-main".to_owned(),
7531 crate::conduct::NODE.to_owned(),
7532 "conduct".to_owned(),
7533 "Which backend?".to_owned(),
7534 String::new(),
7535 Vec::new(),
7536 );
7537 questions.put(&mut q).unwrap();
7538 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
7539 .unwrap();
7540 questions.put(&mut q).unwrap();
7541
7542 let mut blocked = task();
7543 blocked.id = "20260101-000001-main".to_owned();
7544 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
7545 queue.put(&mut blocked).unwrap();
7546
7547 resolve_blockers(&queue, &questions);
7548
7549 let after = queue.get(&blocked.id).unwrap();
7550 assert_eq!(
7551 after.status,
7552 TaskStatus::Queued,
7553 "the only blocker resolved"
7554 );
7555 assert_eq!(after.answers.len(), 1);
7556 assert_eq!(after.answers[0].question, "Which backend?");
7557 assert_eq!(after.answers[0].answer, "SQLite");
7558
7559 let instruction = instruction_for(&after);
7561 assert!(instruction.contains("Which backend?"));
7562 assert!(instruction.contains("SQLite"));
7563 }
7564
7565 #[test]
7566 fn resolve_blockers_holds_a_task_whose_conductor_question_was_abandoned() {
7567 let dir = tempfile::tempdir().unwrap();
7568 let queue = Queue::at(dir.path().join("queue"));
7569 let questions = ask::Questions::at(dir.path().join("questions"));
7570
7571 let mut q = crate::ask::Question::new(
7572 "20260101-000001-main".to_owned(),
7573 crate::conduct::NODE.to_owned(),
7574 "conduct".to_owned(),
7575 "Is the setup done?".to_owned(),
7576 String::new(),
7577 Vec::new(),
7578 );
7579 q.abandon("no answer within 60s of asking");
7580 questions.put(&mut q).unwrap();
7581
7582 let mut blocked = task();
7583 blocked.id = "20260101-000001-main".to_owned();
7584 blocked.block(vec![q.id.clone()], Some("setup?".to_owned()));
7585 queue.put(&mut blocked).unwrap();
7586
7587 resolve_blockers(&queue, &questions);
7588
7589 let after = queue.get(&blocked.id).unwrap();
7590 assert_eq!(
7591 after.status,
7592 TaskStatus::Held,
7593 "never left blocked on nothing"
7594 );
7595 assert!(!after.operator_held(), "a machine hold, for triage");
7596 assert!(
7597 after
7598 .hold_reason
7599 .as_deref()
7600 .unwrap_or_default()
7601 .contains("went unanswered")
7602 );
7603 }
7604
7605 #[test]
7606 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
7607 let dir = tempfile::tempdir().unwrap();
7613 let queue = Queue::at(dir.path().join("queue"));
7614 let questions = ask::Questions::at(dir.path().join("questions"));
7615
7616 let mut q = crate::ask::Question::new(
7617 "20260101-000001-main".to_owned(),
7618 crate::conduct::NODE.to_owned(),
7619 "conduct".to_owned(),
7620 "How should this be handled?".to_owned(),
7621 String::new(),
7622 Vec::new(),
7623 );
7624 questions.put(&mut q).unwrap();
7625 q.answer(crate::ask::Answer::Text(
7626 "leave it held, a human will look at it later".to_owned(),
7627 ))
7628 .unwrap();
7629 questions.put(&mut q).unwrap();
7630
7631 let mut held = task();
7632 held.id = "20260101-000001-main".to_owned();
7633 held.hold_machine(Some("out of attempts".to_owned()));
7634 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
7635 queue.put(&mut held).unwrap();
7636
7637 resolve_blockers(&queue, &questions);
7638
7639 let after = queue.get(&held.id).unwrap();
7640 assert_eq!(after.status, TaskStatus::Held);
7641 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
7642 assert_eq!(
7643 after.answers[0].answer,
7644 "leave it held, a human will look at it later"
7645 );
7646 }
7647
7648 #[test]
7649 fn resolve_blockers_releases_a_task_whose_dependency_was_deleted_on_purpose() {
7650 let dir = tempfile::tempdir().unwrap();
7651 let queue = Queue::at(dir.path().join("queue"));
7652 let questions = ask::Questions::at(dir.path().join("questions"));
7653 let mut dep = task();
7654 dep.id = "20260101-000001-gone".to_owned();
7655 queue.put(&mut dep).unwrap();
7656 let mut blocked = task();
7657 blocked.id = "20260101-000003-main".to_owned();
7658 blocked.block(vec![dep.id.clone()], None);
7659 queue.put(&mut blocked).unwrap();
7660
7661 let claim = queue.claim(&blocked.id).unwrap();
7663 queue.remove(&dep.id, false, &questions).unwrap();
7664 drop(claim);
7665 resolve_blockers(&queue, &questions);
7666
7667 let after = queue.get(&blocked.id).unwrap();
7668 assert_eq!(after.status, TaskStatus::Queued);
7669 assert!(after.blocked_by.is_empty());
7670 }
7671
7672 #[test]
7673 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
7674 let dir = tempfile::tempdir().unwrap();
7680 let queue = Queue::at(dir.path().join("queue"));
7681 let questions = ask::Questions::at(dir.path().join("questions"));
7682
7683 let mut still_going = task();
7684 still_going.id = "20260101-000002-dep1".to_owned();
7685 queue.put(&mut still_going).unwrap();
7686
7687 let mut blocked = task();
7688 blocked.id = "20260101-000003-main".to_owned();
7689 blocked.block(
7690 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7691 Some("waits on both".to_owned()),
7692 );
7693 queue.put(&mut blocked).unwrap();
7694
7695 resolve_blockers(&queue, &questions);
7696
7697 let after = queue.get(&blocked.id).unwrap();
7698 assert_eq!(
7699 after.status,
7700 TaskStatus::Held,
7701 "a missing dependency must not leave the task blocked forever"
7702 );
7703 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7704 assert!(after.blocked_by.is_empty());
7705 let reason = after.hold_reason.as_deref().unwrap_or_default();
7706 assert!(
7707 reason.contains("20260101-000001-gone"),
7708 "the missing id must be named so an operator can tell what happened: {reason}"
7709 );
7710 assert!(
7711 reason.contains(&still_going.id),
7712 "the still-valid dependency must not silently vanish from the record: {reason}"
7713 );
7714 }
7715
7716 #[test]
7717 fn instruction_for_is_unchanged_without_any_answers() {
7718 let t = task();
7719 assert_eq!(instruction_for(&t), t.instruction);
7720 }
7721
7722 #[test]
7723 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7724 let dir = tempfile::tempdir().unwrap();
7725 let q = Queue::at(dir.path().join("queue"));
7726 let src = dir.path().join("shot.png");
7727 std::fs::write(&src, "x").unwrap();
7728 let mut t = task();
7729 q.attach(&mut t, &[src]).unwrap();
7730 let paths = task_attachments(&q, &t).unwrap();
7731 assert_eq!(paths.len(), 1);
7732 assert!(paths[0].is_absolute() && paths[0].is_file());
7733 std::fs::remove_file(&paths[0]).unwrap();
7734 let err = task_attachments(&q, &t).unwrap_err().to_string();
7735 assert!(err.contains("shot.png"), "{err}");
7736 }
7737
7738 #[test]
7739 fn resumed_instruction_is_unchanged_without_any_answers() {
7740 let t = task();
7741 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7742 }
7743
7744 #[test]
7745 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7746 let mut t = task();
7747 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7748 let old = t.instruction.clone();
7752
7753 let refreshed = resumed_instruction(&old, &t);
7754 assert!(refreshed.starts_with(&old), "the original text is kept");
7755 assert!(refreshed.contains("Which backend?"));
7756 assert!(refreshed.contains("SQLite"));
7757 }
7758
7759 #[test]
7760 fn resumed_instruction_keeps_an_original_answers_heading() {
7761 let mut t = task();
7762 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7763 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7764
7765 let refreshed = resumed_instruction(&t.instruction, &t);
7766
7767 assert!(
7768 refreshed.starts_with(&t.instruction),
7769 "an answers heading in the original instruction is not the appended block"
7770 );
7771 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7772 assert!(refreshed.contains("Which backend?"));
7773 assert!(refreshed.contains("SQLite"));
7774
7775 let repeated = resumed_instruction(&refreshed, &t);
7776 assert_eq!(
7777 repeated, refreshed,
7778 "only the final appended block is refreshed"
7779 );
7780 }
7781
7782 #[test]
7783 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7784 let mut t = task();
7785 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7786
7787 let once = resumed_instruction(&t.instruction, &t);
7791 let twice = resumed_instruction(&once, &t);
7792 assert_eq!(once, twice);
7793 assert_eq!(once.matches("Which backend?").count(), 1);
7794
7795 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7797 let refreshed = resumed_instruction(&once, &t);
7798 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7799 assert!(refreshed.contains("Which backend?"));
7800 assert!(refreshed.contains("Which cache?"));
7801 }
7802
7803 #[test]
7804 fn prepare_instruction_covers_all_three_starters() {
7805 let mut t = task();
7806 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7807
7808 assert_eq!(
7811 prepare_instruction(&Starter::Start, None, &t),
7812 Some(instruction_for(&t))
7813 );
7814
7815 let old = t.instruction.clone();
7818 assert_eq!(
7819 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7820 Some(resumed_instruction(&old, &t))
7821 );
7822
7823 assert_eq!(
7827 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7828 None
7829 );
7830 }
7831
7832 #[test]
7833 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7834 assert_eq!(
7835 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7836 Starter::Review("magi/eba2/A".to_owned())
7837 );
7838 }
7839
7840 #[test]
7841 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7842 assert_eq!(
7843 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7844 Starter::Start,
7845 "a vanished review branch must not fall back to resuming the old run either"
7846 );
7847 }
7848
7849 #[test]
7850 fn a_refused_handover_retries_as_a_review_of_the_same_branch() {
7851 let mut t = task();
7852 t.start("old-run".to_owned());
7853 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
7854 t.release();
7855 let branch = t.review_branch.take();
7856 assert_eq!(
7857 choose_starter(branch.as_deref(), true, Some("old-run")),
7858 Starter::Review("magi/eba2/A".to_owned()),
7859 "a review wins over resuming the old run"
7860 );
7861 }
7862
7863 #[test]
7864 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7865 assert_eq!(
7866 choose_starter(None, false, Some("some-run")),
7867 Starter::Resume("some-run".to_owned())
7868 );
7869 assert_eq!(choose_starter(None, false, None), Starter::Start);
7870 }
7871
7872 #[test]
7873 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7874 let mut released = task();
7875 released.start("stalled-run".to_owned());
7876 released.requeue();
7877 let unfinished = (!released.fresh_start)
7878 .then(|| Some("stalled-run".to_owned()))
7879 .flatten();
7880 assert_eq!(
7881 choose_starter(None, false, unfinished.as_deref()),
7882 Starter::Start,
7883 "release keeps run history but must not resume it"
7884 );
7885 assert_eq!(released.runs, ["stalled-run"]);
7886 }
7887
7888 #[test]
7889 fn an_ordinary_release_keeps_a_resumable_run_available() {
7890 let mut released = task();
7891 released.start("stalled-run".to_owned());
7892 released.release();
7893 let unfinished = (!released.fresh_start)
7894 .then(|| Some("stalled-run".to_owned()))
7895 .flatten();
7896 assert_eq!(
7897 choose_starter(None, false, unfinished.as_deref()),
7898 Starter::Resume("stalled-run".to_owned()),
7899 "manual release must preserve the normal resume path"
7900 );
7901 }
7902
7903 #[test]
7904 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7905 let mut state = run_state(RunStatus::Blocked);
7906 state.config.graph.review_rounds = 3;
7907 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7908 assert!(exhausted_review_budget(&state));
7909
7910 state.reviews.pop();
7912 assert!(!exhausted_review_budget(&state));
7913
7914 let mut stalled = run_state(RunStatus::Stalled);
7917 stalled.config.graph.review_rounds = 1;
7918 stalled.reviews = vec![review_round(1)];
7919 assert!(!exhausted_review_budget(&stalled));
7920 }
7921
7922 fn review_round(round: usize) -> crate::run::ReviewRound {
7923 crate::run::ReviewRound {
7924 round,
7925 head: "deadbeef".to_owned(),
7926 verified_head: None,
7927 verified_at: None,
7928 reviews: Vec::new(),
7929 e2e: Vec::new(),
7930 verify_retried: false,
7931 e2e_deferred: false,
7932 e2e_defer_reason: None,
7933 fix: None,
7934 blocking: 0,
7935 answered: 1,
7936 expected: 1,
7937 clean: false,
7938 progressed: true,
7939 vote_split: false,
7940 reconsideration: Vec::new(),
7941 verdict: None,
7942 }
7943 }
7944
7945 #[test]
7946 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7947 let mut run = RunState::new(
7948 PathBuf::from("/repo"),
7949 "main".to_owned(),
7950 "abc1234def".to_owned(),
7951 "add retries".to_owned(),
7952 Config::default(),
7953 );
7954 run.status = RunStatus::Judging;
7955 run.parked = true;
7956 let mut task = Task::new(
7957 "add retries".to_owned(),
7958 "add retries".to_owned(),
7959 PathBuf::from("/repo"),
7960 crate::queue::Source::Human,
7961 );
7962 task.status = TaskStatus::Failed;
7963 task.runs = vec![run.id.clone()];
7964 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7965 assert!(with(&task, &run), "parked after judging is the case");
7966 let mut parked_task = task.clone();
7967 parked_task.status = TaskStatus::Parked;
7968 assert!(with(&parked_task, &run), "the Parked status is the case");
7969
7970 let mut not_parked = run.clone();
7971 not_parked.parked = false;
7972 not_parked.status = RunStatus::Stalled;
7973 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7974
7975 let mut fresh = task.clone();
7976 fresh.fresh_start = true;
7977 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7978
7979 let mut review = task.clone();
7980 review.review_branch = Some("magi/x/A".to_owned());
7981 assert!(!with(&review, &run), "review is ranked before resume");
7982
7983 let mut held = task.clone();
7984 held.status = TaskStatus::Held;
7985 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7986
7987 let mut released = run.clone();
7988 released.released_to = Some("20260901-000000-new1".to_owned());
7989 assert!(!with(&task, &released), "nothing left to resume into");
7990
7991 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7992 "unreadable"
7993 )));
7994 }
7995
7996 #[test]
7997 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7998 let mut released = RunState::new(
7999 PathBuf::from("/repo"),
8000 "main".to_owned(),
8001 "abc1234def".to_owned(),
8002 "add retries".to_owned(),
8003 Config::default(),
8004 );
8005 released.status = RunStatus::Blocked;
8006 assert_eq!(
8007 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
8008 Some(released.id.clone())
8009 );
8010 released.released_to = Some("20260901-000000-new1".to_owned());
8011 assert_eq!(
8012 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
8013 None,
8014 "there is nothing left to resume it into"
8015 );
8016 }
8017
8018 #[test]
8019 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
8020 let mut exhausted = RunState::new(
8029 PathBuf::from("/repo"),
8030 "main".to_owned(),
8031 "abc1234def".to_owned(),
8032 "add retries".to_owned(),
8033 Config::default(),
8034 );
8035 exhausted.status = RunStatus::Blocked;
8036 exhausted.config.graph.review_rounds = 1;
8037 exhausted.reviews = vec![review_round(1)];
8038
8039 assert_eq!(
8040 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
8041 None,
8042 "an exhausted `Blocked` run must not be offered as resumable"
8043 );
8044
8045 let mut has_budget_left = RunState::new(
8048 PathBuf::from("/repo"),
8049 "main".to_owned(),
8050 "abc1234def".to_owned(),
8051 "add retries".to_owned(),
8052 Config::default(),
8053 );
8054 has_budget_left.status = RunStatus::Blocked;
8055 has_budget_left.config.graph.review_rounds = 3;
8056 has_budget_left.reviews = vec![review_round(1)];
8057
8058 assert_eq!(
8059 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
8060 Ok(has_budget_left.clone())
8061 }),
8062 Some(has_budget_left.id.clone())
8063 );
8064 }
8065
8066 #[test]
8067 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
8068 let mut older_stalled = RunState::new(
8076 PathBuf::from("/repo"),
8077 "main".to_owned(),
8078 "abc1234def".to_owned(),
8079 "add retries".to_owned(),
8080 Config::default(),
8081 );
8082 older_stalled.status = RunStatus::Stalled;
8083
8084 let mut newest_exhausted = RunState::new(
8085 PathBuf::from("/repo"),
8086 "main".to_owned(),
8087 "abc1234def".to_owned(),
8088 "add retries".to_owned(),
8089 Config::default(),
8090 );
8091 newest_exhausted.status = RunStatus::Blocked;
8092 newest_exhausted.config.graph.review_rounds = 1;
8093 newest_exhausted.reviews = vec![review_round(1)];
8094
8095 assert_eq!(
8096 unfinished_run_with(
8097 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
8098 "t",
8099 |_| Ok(newest_exhausted.clone())
8100 ),
8101 None,
8102 "the newest run is exhausted, so nothing here is worth resuming - \
8103 least of all the older, already-superseded run"
8104 );
8105 }
8106
8107 #[test]
8108 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
8109 assert_eq!(
8110 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
8111 Err(anyhow::anyhow!("fixture is absent"))
8112 }),
8113 None
8114 );
8115 }
8116
8117 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
8118 let mut q = ask::Question::new(
8119 run.to_owned(),
8120 "implement".to_owned(),
8121 "impl-A".to_owned(),
8122 "continue?".to_owned(),
8123 String::new(),
8124 vec!["resume で続行する".to_owned(), "other".to_owned()],
8125 );
8126 q.actions.insert("resume で続行する".to_owned(), action);
8127 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
8128 .unwrap();
8129 q
8130 }
8131
8132 fn held_task_with(run: &str) -> Task {
8133 let mut t = task();
8134 t.runs = vec![run.to_owned()];
8135 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
8136 t
8137 }
8138
8139 fn resume_action(run: &str) -> ask::ChoiceAction {
8140 ask::ChoiceAction::Resume { run: run.into() }
8141 }
8142
8143 #[test]
8144 fn decide_action_resumes_only_the_latest_resumable_run() {
8145 let t = held_task_with("r1");
8146 let q = action_question("r1", resume_action("r1"));
8147 let load = |s: RunState| move |_: &str| Ok(s);
8148 assert_eq!(
8149 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
8150 ActionDecision::Resume("r1".into())
8151 );
8152 let q_other = action_question("r1", resume_action("r0"));
8154 assert!(matches!(
8155 decide_action(
8156 &t,
8157 &q_other,
8158 &PHRASES_EN,
8159 load(run_state(RunStatus::Blocked))
8160 ),
8161 ActionDecision::Refuse(_)
8162 ));
8163 assert!(matches!(
8165 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
8166 ActionDecision::Refuse(_)
8167 ));
8168 let mut released = run_state(RunStatus::Blocked);
8170 released.released_to = Some("elsewhere".into());
8171 assert!(matches!(
8172 decide_action(&t, &q, &PHRASES_EN, load(released)),
8173 ActionDecision::Refuse(_)
8174 ));
8175 assert!(matches!(
8177 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
8178 ActionDecision::Refuse(_)
8179 ));
8180 }
8181
8182 #[test]
8183 fn decide_action_ignores_a_question_about_an_earlier_run() {
8184 let mut t = held_task_with("r1");
8185 t.runs.push("r2".to_owned());
8186 let never = |_: &str| -> Result<RunState> { bail!("not read") };
8187 assert_eq!(
8188 decide_action(
8189 &t,
8190 &action_question("r1", ask::ChoiceAction::Done),
8191 &PHRASES_EN,
8192 never
8193 ),
8194 ActionDecision::Stale
8195 );
8196 }
8197
8198 #[test]
8199 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
8200 let mut t = held_task_with("r1");
8201 let never = |_: &str| -> Result<RunState> { bail!("not read") };
8202 assert_eq!(
8203 decide_action(
8204 &t,
8205 &action_question("r1", ask::ChoiceAction::Requeue),
8206 &PHRASES_EN,
8207 never
8208 ),
8209 ActionDecision::Requeue
8210 );
8211 let done_q = action_question("r1", ask::ChoiceAction::Done);
8212 assert_eq!(
8213 decide_action(&t, &done_q, &PHRASES_EN, never),
8214 ActionDecision::Done
8215 );
8216 t.mark_action_applied(&done_q.id);
8217 assert_eq!(
8218 decide_action(&t, &done_q, &PHRASES_EN, never),
8219 ActionDecision::Skip
8220 );
8221
8222 let mut plain = action_question("r1", ask::ChoiceAction::Done);
8224 plain.actions.clear();
8225 assert_eq!(
8226 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
8227 ActionDecision::Skip
8228 );
8229 let mut running = held_task_with("r1");
8231 running.status = TaskStatus::Running;
8232 assert_eq!(
8233 decide_action(
8234 &running,
8235 &action_question("r1", ask::ChoiceAction::Done),
8236 &PHRASES_EN,
8237 never
8238 ),
8239 ActionDecision::Skip
8240 );
8241 }
8242
8243 #[test]
8244 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
8245 let dir = tempfile::tempdir().unwrap();
8246 let queue = Queue::at(dir.path().join("queue"));
8247 let questions = Questions::at(dir.path().join("questions"));
8248 let home = dir.path().join("home");
8249 let mut state = run_state(RunStatus::Blocked);
8250 state.id = "20260101-000000-act1".to_owned();
8251 state.save_under(&home).unwrap();
8252
8253 let mut t = held_task_with(&state.id);
8254 queue.put(&mut t).unwrap();
8255 let mut q = action_question(&state.id, resume_action(&state.id));
8256 questions.put(&mut q).unwrap();
8257
8258 apply_choice_actions(&queue, &questions, &home);
8259 let after = queue.get(&t.id).unwrap();
8260 assert_eq!(after.status, TaskStatus::Queued);
8261 assert!(!after.fresh_start);
8262 assert!(after.action_applied(&q.id));
8263 let pin = after.resume_override.clone().unwrap();
8264 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
8265 assert!(pin.forced);
8266
8267 let mut again = queue.get(&t.id).unwrap();
8269 again.hold_machine(Some("later".into()));
8270 queue.put(&mut again).unwrap();
8271 apply_choice_actions(&queue, &questions, &home);
8272 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8273 }
8274
8275 #[test]
8276 fn a_delivered_answer_is_still_acted_on_exactly_once() {
8277 let dir = tempfile::tempdir().unwrap();
8278 let queue = Queue::at(dir.path().join("queue"));
8279 let questions = Questions::at(dir.path().join("questions"));
8280 let home = dir.path().join("home");
8281 let mut state = run_state(RunStatus::Blocked);
8282 state.id = "20260101-000000-act2".to_owned();
8283 state.save_under(&home).unwrap();
8284
8285 let mut t = held_task_with(&state.id);
8286 queue.put(&mut t).unwrap();
8287 let mut q = action_question(&state.id, resume_action(&state.id));
8288 q.answer_delivered = true;
8289 questions.put(&mut q).unwrap();
8290
8291 apply_choice_actions(&queue, &questions, &home);
8294 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
8295 assert!(queue.get(&t.id).unwrap().action_applied(&q.id));
8296
8297 let mut again = queue.get(&t.id).unwrap();
8299 again.hold_machine(Some("later".into()));
8300 queue.put(&mut again).unwrap();
8301 apply_choice_actions(&queue, &questions, &home);
8302 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8303 }
8304
8305 #[test]
8306 fn a_fresh_asker_defers_the_action_until_it_goes_quiet() {
8307 let dir = tempfile::tempdir().unwrap();
8308 let queue = Queue::at(dir.path().join("queue"));
8309 let questions = Questions::at(dir.path().join("questions"));
8310 let home = dir.path().join("home");
8311 let mut state = run_state(RunStatus::Blocked);
8312 state.id = "20260101-000000-act3".to_owned();
8313 state.save_under(&home).unwrap();
8314 let mut t = held_task_with(&state.id);
8315 queue.put(&mut t).unwrap();
8316 let mut q = action_question(&state.id, ask::ChoiceAction::Requeue);
8317 questions.put(&mut q).unwrap();
8318
8319 questions.beat(&q.id, crate::ask::WaiterKind::Asker);
8320 apply_choice_actions(&queue, &questions, &home);
8321 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8322
8323 std::fs::remove_file(questions.root().join(format!("{}.lease", q.id))).unwrap();
8324 apply_choice_actions(&queue, &questions, &home);
8325 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
8326 }
8327
8328 #[test]
8329 fn action_standing_tells_the_reasons_apart() {
8330 let mut t = held_task_with("r1");
8331 let q = action_question("r1", ask::ChoiceAction::Requeue);
8332 assert_eq!(action_standing(&t, &q), ActionStanding::Pending);
8333 let mut plain = q.clone();
8334 plain.actions.clear();
8335 assert_eq!(action_standing(&t, &plain), ActionStanding::NoAction);
8336
8337 t.mark_action_applied(&q.id);
8339 assert_eq!(action_standing(&t, &q), ActionStanding::Applied);
8340
8341 let mut t = held_task_with("r1");
8343 t.start("r2".to_owned());
8344 assert_eq!(action_standing(&t, &q), ActionStanding::Stale);
8345 let q2 = action_question("r2", ask::ChoiceAction::Requeue);
8347 assert_eq!(action_standing(&t, &q2), ActionStanding::Busy);
8348 t.status = TaskStatus::Blocked;
8349 assert_eq!(action_standing(&t, &q2), ActionStanding::Busy);
8350
8351 let mut c = action_question("r0", ask::ChoiceAction::Requeue);
8353 c.node = crate::conduct::NODE.to_owned();
8354 assert_eq!(action_standing(&t, &c), ActionStanding::Busy);
8355 assert!(!ActionStanding::Busy.daemon_owns());
8356 }
8357
8358 #[test]
8359 fn a_running_task_waits_and_a_dead_asker_with_a_cwd_is_still_actioned() {
8360 let dir = tempfile::tempdir().unwrap();
8361 let queue = Queue::at(dir.path().join("queue"));
8362 let questions = Questions::at(dir.path().join("questions"));
8363 let home = dir.path().join("home");
8364 let mut state = run_state(RunStatus::Blocked);
8365 state.id = "20260101-000000-act4".to_owned();
8366 state.save_under(&home).unwrap();
8367 let mut t = held_task_with(&state.id);
8368 t.status = TaskStatus::Running;
8369 queue.put(&mut t).unwrap();
8370 let mut q = action_question(&state.id, ask::ChoiceAction::Done);
8371 q.cwd = Some(dir.path().display().to_string());
8372 questions.put(&mut q).unwrap();
8373
8374 let running = queue.get(&t.id).unwrap();
8377 assert_eq!(action_standing(&running, &q), ActionStanding::Busy);
8378 assert!(matches!(
8379 crate::waiter::decide_owned(&q, None, false, 86_400, Timestamp::now(), false),
8380 crate::waiter::Action::Deliver(_)
8381 ));
8382 apply_choice_actions(&queue, &questions, &home);
8384 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
8385 let mut back = queue.get(&t.id).unwrap();
8386 back.hold_machine(Some("later".into()));
8387 queue.put(&mut back).unwrap();
8388 apply_choice_actions(&queue, &questions, &home);
8389 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Done);
8390 let mut again = queue.get(&t.id).unwrap();
8391 again.hold_machine(Some("again".into()));
8392 queue.put(&mut again).unwrap();
8393 apply_choice_actions(&queue, &questions, &home);
8394 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8395 }
8396}