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.stall(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 mut config = match prepare_for(&repo, opts, task) {
2759 Ok(c) => c,
2760 Err(e) => {
2761 task.attempts += 1;
2765 task.fail(format!("config: {e:#}"), opts.max_attempts);
2766 record(queue, task);
2767 return Vec::new();
2768 }
2769 };
2770 apply_solo(&mut config, task);
2771 let p = phrases(&config.graph.language);
2772 let start_failed = p.could_not_start;
2773
2774 if let Some(reason) = disk_gate(&repo, &config) {
2782 task.last_error = Some(reason.clone());
2783 task.hold_machine(Some(reason.clone()));
2784 record(queue, task);
2785 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2786 notices::raise(
2789 Notice::warn(
2790 &format!("disk:{}", repo.display()),
2791 "A task was held for want of free disk space; free some, then release it from the queue.",
2792 )
2793 .link(Link::Task {
2794 id: task.id.clone(),
2795 }),
2796 );
2797 return Vec::new();
2798 }
2799
2800 let unfinished = (!task.fresh_start)
2820 .then(|| unfinished_run(&task.runs, task.short()))
2821 .flatten();
2822 let review_branch = task.review_branch.take();
2828 let branch_exists = match &review_branch {
2829 Some(branch) => crate::git::branch_exists(&repo, branch)
2830 .await
2831 .unwrap_or(false),
2832 None => false,
2833 };
2834 let starter = match (&review_branch, &unfinished, &task.review_of) {
2838 (None, None, Some(branch)) => Starter::Review(branch.clone()),
2839 _ => choose_starter(
2840 review_branch.as_deref(),
2841 branch_exists,
2842 unfinished.as_deref(),
2843 ),
2844 };
2845 let attachments = match task_attachments(queue, task) {
2850 Ok(a) => a,
2851 Err(e) => {
2852 task.attempts += 1;
2853 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2854 record(queue, task);
2855 return Vec::new();
2856 }
2857 };
2858 let started = match &starter {
2859 Starter::Review(branch) => {
2860 tracing::info!(
2861 "task {} reopens `{branch}` as a review-only pass",
2862 task.short()
2863 );
2864 let takeover = crate::handover::Takeover {
2867 earlier: task.earlier_attempts().to_vec(),
2868 home: crate::run::home(),
2869 choice: take_divergence_answer(branch, &config.merge.remote, task),
2870 };
2871 Runner::review_taking_over(
2872 &repo,
2873 branch,
2874 config,
2875 Some(takeover),
2876 crate::run::Origin::queue(&task.id),
2877 )
2878 .await
2879 }
2880 Starter::Resume(id) => {
2881 tracing::info!("resuming run {id} rather than competing again");
2882 Runner::resume(id).map(|mut r| {
2883 if let Some(instruction) =
2884 prepare_instruction(&starter, Some(&r.state.instruction), task)
2885 {
2886 r.state.instruction = instruction;
2887 }
2888 r.state.attachments = attachments.clone();
2889 if let Some(mode) = task
2892 .overrides
2893 .as_ref()
2894 .and_then(|o| o.merge.as_deref())
2895 .and_then(|m| merge_mode(m).ok())
2896 {
2897 r.state.config.merge.mode = mode;
2898 }
2899 r
2900 })
2901 }
2902 Starter::Start => {
2903 if let Some(branch) = &review_branch {
2904 tracing::warn!(
2905 "conductor chose review for task {} but branch `{branch}` no longer \
2906 exists; requeuing as a fresh competition instead",
2907 task.short()
2908 );
2909 }
2910 let instruction = prepare_instruction(&starter, None, task)
2911 .unwrap_or_else(|| task.instruction.clone());
2912 Runner::start_naming(
2913 &repo,
2914 instruction,
2915 &task.title,
2916 config,
2917 crate::run::Origin::queue(&task.id),
2918 )
2919 .await
2920 .map(|mut r| {
2921 r.state.attachments = attachments.clone();
2922 r
2923 })
2924 }
2925 };
2926 let mut runner = match started {
2927 Ok(r) => r,
2928 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2932 let detail = match e.downcast_ref::<crate::handover::Refused>() {
2933 Some(r) => (p.handover_refused)(r),
2934 None => format!("{e:#}"),
2935 };
2936 let reason = format!("{start_failed}{detail}");
2937 task.last_error = Some(reason.clone());
2938 let branch = match &starter {
2941 Starter::Review(branch) => Some(branch.clone()),
2942 _ => None,
2943 };
2944 task.hold_for_handover(branch, format!("{reason}{}", p.handover_hint));
2948 record(queue, task);
2949 tracing::warn!(
2950 "holding {} for a branch it cannot take over: {e:#}",
2951 task.short()
2952 );
2953 return Vec::new();
2954 }
2955 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2960 let d = e
2961 .downcast_ref::<crate::reconcile::Diverged>()
2962 .expect("checked by the guard");
2963 let mut q = ask::Question::new(
2964 task.id.clone(),
2965 "review".to_owned(),
2966 "sync".to_owned(),
2967 d.summary(),
2968 d.detail(),
2969 d.choices(),
2970 );
2971 match Questions::open().put(&mut q) {
2972 Ok(()) => {
2973 task.last_error = Some(format!("{e:#}"));
2974 task.review_branch = Some(d.branch.clone());
2977 task.block(vec![q.id.clone()], Some(d.summary()));
2978 }
2979 Err(put) => {
2980 tracing::warn!("could not file the divergence question: {put:#}");
2981 task.attempts += 1;
2982 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2983 }
2984 }
2985 record(queue, task);
2986 return Vec::new();
2987 }
2988 Err(e) => {
2989 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
2995 && let Starter::Review(branch) = &starter
2996 {
2997 task.review_branch = Some(branch.clone());
2998 task.last_error = Some(format!("{start_failed}{e:#}"));
2999 task.status = crate::queue::TaskStatus::Failed;
3000 record(queue, task);
3001 return Vec::new();
3002 }
3003 task.attempts += 1;
3004 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
3005 record(queue, task);
3006 return Vec::new();
3007 }
3008 };
3009 runner.state.followup_generation = Some(task.followup.as_ref().map_or(0, |f| f.generation));
3012 runner.on_pause(stop.pause());
3014 runner.watch_interrupt(interrupt_pause);
3018
3019 let run = runner.state.id.clone();
3022 task.start(run.clone());
3023 record(queue, task);
3024 lock(status).current.push(Current {
3025 task: task.id.clone(),
3026 run,
3027 });
3028
3029 let quota_before = runner.state.quota.clone();
3032 let result = runner.execute().await;
3033 finish_attempt(
3034 opts.max_attempts,
3035 queue,
3036 task,
3037 &runner.state,
3038 "a_before,
3039 result,
3040 )
3041}
3042
3043pub fn finish_attempt(
3048 max_attempts: usize,
3049 queue: &Queue,
3050 task: &mut Task,
3051 state: &RunState,
3052 quota_before: &[QuotaLoss],
3053 result: Result<()>,
3054) -> Vec<QuotaLoss> {
3055 let detail = match result {
3056 Ok(()) => describe(state),
3057 Err(e) => format!("{e:#}"),
3058 };
3059 let fresh = losses_this_attempt(quota_before, &state.quota);
3060 let verdict = Verdict {
3061 status: state.status,
3062 left_pr: state.pr.is_some(),
3065 quota_hit: !fresh.is_empty(),
3071 parked: state.parked,
3075 no_viable_candidates: state.viable().is_empty(),
3078 };
3079 settle_and_diagnose(task, verdict, &detail, max_attempts, state);
3080 if task.status == TaskStatus::Done {
3081 supersede_prior_runs(task, &crate::run::home());
3082 }
3083 record(queue, task);
3084 tracing::info!(
3085 "task {} is {} after run {} ({})",
3086 task.short(),
3087 task.status.as_str(),
3088 state.short(),
3089 label(state.status)
3090 );
3091 fresh
3092}
3093
3094pub fn hold_if_runnable(queue: &Queue, task: &mut Task) {
3098 if task.status.runnable() {
3099 let why = task.last_error.clone().map_or_else(
3100 || "the run did not finish".to_owned(),
3101 |e| format!("the run did not finish: {e}"),
3102 );
3103 task.hold_manual(Some(format!(
3104 "{why}. It was started by hand, so it is not retried \
3105 automatically; `magi task release` retries it."
3106 )));
3107 record(queue, task);
3108 }
3109}
3110
3111#[must_use]
3115pub fn loop_would_resume(run: &str) -> bool {
3116 unfinished_run(&[run.to_owned()], crate::run::short_of(run)).is_some()
3117}
3118
3119#[must_use]
3129pub fn foreign_loop(
3130 reading: Option<&Reading>,
3131 now: Timestamp,
3132 own_pid: u32,
3133) -> Option<Option<u32>> {
3134 let reading = reading.filter(|r| r.running(now))?;
3135 match reading.pid {
3136 Some(pid) if pid == own_pid => None,
3137 pid => Some(pid),
3138 }
3139}
3140
3141pub async fn run_claimed(opts: &Opts, queue: &Queue, task: &mut Task) {
3156 let status = Arc::new(Mutex::new(Status::new()));
3157 let stop = Stop::new();
3158 attempt(
3159 opts,
3160 queue,
3161 &status,
3162 &stop,
3163 crate::graph::Pause::new(),
3164 task,
3165 )
3166 .await;
3167 hold_if_runnable(queue, task);
3168}
3169
3170fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
3180 after
3181 .iter()
3182 .filter(|q| !before.contains(q))
3183 .cloned()
3184 .collect()
3185}
3186
3187fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
3190 if quota.is_empty() {
3191 return None;
3192 }
3193 let with_hint = quota.iter().find(|q| q.reset.is_some());
3194 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
3195 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
3196 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
3197 Some(
3198 now.checked_add(jiff::SignedDuration::from_secs(secs))
3199 .unwrap_or(Timestamp::MAX),
3200 )
3201}
3202
3203fn apply_solo(config: &mut Config, task: &Task) {
3213 if task.solo {
3214 config.graph.candidates = 1;
3215 }
3216}
3217
3218fn prepare_for(repo: &Path, opts: &Opts, task: &Task) -> Result<Config> {
3223 let Some(o) = &task.overrides else {
3224 return prepare(repo, opts);
3225 };
3226 let own = Opts {
3230 config: o.config.clone(),
3231 merge: None,
3232 ..opts.clone()
3233 };
3234 let mut config = prepare(repo, &own)?;
3235 o.apply(&mut config);
3236 Ok(config)
3237}
3238
3239fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
3241 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
3242 if let Some(mode) = &opts.merge {
3243 config.merge.mode = merge_mode(mode)?;
3244 }
3245 Ok(config)
3246}
3247
3248async fn maybe_prune_cache_between_runs(
3279 repo: &Path,
3280 opts: &Opts,
3281 home: &Path,
3282 stop: &Stop,
3283 last_checked: &mut Option<Timestamp>,
3284 now: Timestamp,
3285) {
3286 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
3287 return;
3288 }
3289 *last_checked = Some(now);
3290 let cfg = match prepare(repo, opts) {
3291 Ok(cfg) => cfg,
3292 Err(e) => {
3293 tracing::warn!("cache check: no config: {e:#}");
3294 return;
3295 }
3296 };
3297 match clean::prune_cache_if_over_limit(&cfg, home) {
3298 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
3299 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
3300 pruned.files,
3301 pruned.freed
3302 ),
3303 Ok(_) => {}
3304 Err(e) => {
3305 tracing::warn!("housekeep: prune cache: {e:#}");
3306 notices::raise_in(
3307 home,
3308 Notice::warn(
3309 "housekeep:cache",
3310 "Pruning the shared build cache failed; disk usage may keep growing.",
3311 ),
3312 );
3313 }
3314 }
3315}
3316
3317fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
3321 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
3322}
3323
3324async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
3343 let cfg = match prepare(repo, opts) {
3344 Ok(cfg) => cfg,
3345 Err(e) => {
3346 tracing::warn!("housekeep: no config: {e:#}");
3347 return;
3348 }
3349 };
3350 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
3359 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
3360 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
3365 let mut extra = Vec::new();
3366 if out.unreadable > 0 {
3367 extra.push(format!("{} unreadable", out.unreadable));
3368 }
3369 if out.orphaned_worktrees > 0 {
3370 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
3371 }
3372 let detail = if extra.is_empty() {
3373 String::new()
3374 } else {
3375 format!(" ({})", extra.join(", "))
3376 };
3377 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
3378 }
3379 if out.external_merges_recorded > 0 {
3380 tracing::info!(
3381 "housekeep: recorded {} run(s) as merged externally",
3382 out.external_merges_recorded
3383 );
3384 }
3385 if out.stale_pr_states_repaired > 0 {
3386 tracing::info!(
3387 "housekeep: rewrote {} run record(s) whose pull request had already settled",
3388 out.stale_pr_states_repaired
3389 );
3390 }
3391 if out.cache_files > 0 {
3392 tracing::info!(
3393 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
3394 out.cache_files,
3395 out.cache_freed
3396 );
3397 }
3398 if out.questions_abandoned > 0 {
3399 tracing::info!(
3400 "housekeep: abandoned {} question(s) left open by a finished run",
3401 out.questions_abandoned
3402 );
3403 }
3404}
3405
3406async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3415 let questions = Questions::at(home.join("questions"));
3416 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3417 if report.is_empty() {
3418 return;
3419 }
3420 if !report.quarantined.is_empty() {
3421 tracing::info!(
3422 "triage: held {} blocked task(s) whose blocked-on task or \
3423 question no longer exists: {}",
3424 report.quarantined.len(),
3425 report.quarantined.join(", ")
3426 );
3427 }
3428 if !report.resumed.is_empty() {
3429 tracing::info!(
3430 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3431 report.resumed.len(),
3432 report.resumed.join(", ")
3433 );
3434 }
3435 if !report.asked.is_empty() {
3436 tracing::info!(
3437 "triage: asked about {} held task(s): {}",
3438 report.asked.len(),
3439 report.asked.join(", ")
3440 );
3441 }
3442 if !report.answered.is_empty() {
3443 tracing::info!(
3444 "triage: applied {} operator answer(s): {}",
3445 report.answered.len(),
3446 report.answered.join(", ")
3447 );
3448 }
3449}
3450
3451fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3458 disk_gate_with(repo, config, crate::disk::free_bytes)
3459}
3460
3461fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3465 repo: &Path,
3466 config: &Config,
3467 free_bytes: F,
3468) -> Option<String> {
3469 let min = config.disk.min_free_bytes;
3470 if min == 0 {
3471 return None;
3472 }
3473 match free_bytes(repo) {
3474 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3475 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3476 }
3477}
3478
3479const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3486
3487const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3491
3492fn quota_wait(
3501 reset_at: Option<Timestamp>,
3502 now: Timestamp,
3503 fallback: Duration,
3504 cap: Duration,
3505) -> Duration {
3506 match reset_at {
3507 Some(at) if at > now => {
3508 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3509 Duration::from_secs(secs).min(cap)
3510 }
3511 _ => fallback,
3512 }
3513}
3514
3515fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3529 parse_reset_hint_zoned(text, now)
3530 .or_else(|| parse_reset_hint_dated(text))
3531 .or_else(|| parse_reset_hint_relative(text, recorded))
3532}
3533
3534fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3538 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3539 let mut rest = rest.trim();
3540 if rest.is_empty() {
3541 return None;
3542 }
3543 let mut total: i64 = 0;
3544 let mut matched = false;
3545 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3546 if let Some((digits, tail)) = rest.split_once(unit)
3547 && !digits.is_empty()
3548 && digits.bytes().all(|b| b.is_ascii_digit())
3549 {
3550 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3551 rest = tail;
3552 matched = true;
3553 }
3554 }
3555 if !rest.is_empty() || !matched {
3556 return None;
3557 }
3558 recorded
3559 .checked_add(jiff::SignedDuration::from_secs(total))
3560 .ok()
3561}
3562
3563fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3567 let clock = clock.trim().to_lowercase();
3568 let (digits, pm) = clock
3569 .strip_suffix("am")
3570 .map(|d| (d, false))
3571 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3572 let (h, m) = digits.trim().split_once(':')?;
3573 let mut hour: i8 = h.trim().parse().ok()?;
3574 let minute: i8 = m.trim().parse().ok()?;
3575 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3576 return None;
3577 }
3578 if pm && hour != 12 {
3579 hour += 12;
3580 } else if !pm && hour == 12 {
3581 hour = 0;
3582 }
3583 Some((hour, minute))
3584}
3585
3586fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3591 let open = text.find('(')?;
3592 let close = text.rfind(')')?;
3593 if close <= open {
3594 return None;
3595 }
3596 let zone = text[open + 1..close].trim();
3597 let (hour, minute) = parse_12h_clock(&text[..open])?;
3598 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3599 let candidate = now
3600 .to_zoned(tz)
3601 .with()
3602 .hour(hour)
3603 .minute(minute)
3604 .second(0)
3605 .millisecond(0)
3606 .microsecond(0)
3607 .nanosecond(0)
3608 .build()
3609 .ok()?;
3610 let mut at = candidate.timestamp();
3611 if at <= now {
3612 at += jiff::SignedDuration::from_hours(24);
3613 }
3614 Some(at)
3615}
3616
3617fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3626 let words: Vec<&str> = text.split_whitespace().collect();
3627 if words.len() < 5 {
3628 return None;
3629 }
3630 (0..=words.len() - 5)
3631 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3632}
3633
3634fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3640 if trailing.is_some_and(|next| next.starts_with('(')) {
3641 return None;
3642 }
3643 let month = month_number(window[0])?;
3644 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3645 let day_digits = ["st", "nd", "rd", "th"]
3646 .iter()
3647 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3648 let day: i8 = day_digits.parse().ok()?;
3649 let year_token = window[2];
3650 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3651 return None;
3652 }
3653 let year: i16 = year_token.parse().ok()?;
3654 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3658 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3659 let date = jiff::civil::Date::new(year, month, day).ok()?;
3660 let candidate = date
3661 .at(hour, minute, 0, 0)
3662 .to_zoned(jiff::tz::TimeZone::UTC)
3663 .ok()?;
3664 Some(candidate.timestamp())
3665}
3666
3667fn month_number(name: &str) -> Option<i8> {
3670 const NAMES: [&str; 12] = [
3671 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3672 ];
3673 let lower = name.to_lowercase();
3674 NAMES
3675 .iter()
3676 .position(|n| *n == lower.as_str())
3677 .map(|i| i as i8 + 1)
3678}
3679
3680fn exhausted_review_budget(state: &RunState) -> bool {
3692 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3693}
3694
3695fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3732 unfinished_run_with(runs, short, RunState::load)
3733}
3734
3735fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3738where
3739 F: FnOnce(&str) -> Result<RunState>,
3740{
3741 let id = runs.last()?;
3742 match load(id) {
3743 Ok(s)
3748 if s.status.resumable()
3749 && !s.released()
3750 && !exhausted_review_budget(&s)
3751 && s.liveness(false) != crate::run::Liveness::Live =>
3752 {
3753 Some(id.clone())
3754 }
3755 Ok(_) => None,
3756 Err(e) => {
3757 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3758 None
3759 }
3760 }
3761}
3762
3763fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3771where
3772 F: FnOnce(&str) -> Result<RunState>,
3773{
3774 if task.status != TaskStatus::Failed || task.fresh_start || task.review_branch.is_some() {
3775 return false;
3776 }
3777 let Some(id) = task.runs.last() else {
3778 return false;
3779 };
3780 load(id).is_ok_and(|s| {
3781 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3782 })
3783}
3784
3785#[derive(Debug, Clone, PartialEq, Eq)]
3788enum Starter {
3789 Review(String),
3792 Resume(String),
3794 Start,
3796}
3797
3798fn take_divergence_answer(
3802 branch: &str,
3803 remote: &str,
3804 task: &mut Task,
3805) -> Option<crate::reconcile::Choice> {
3806 let summary = crate::reconcile::summary_for(branch, remote);
3807 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3808 (a.question == summary)
3809 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3810 .flatten()
3811 .map(|c| (i, c))
3812 })?;
3813 task.answers.remove(idx);
3814 Some(choice)
3815}
3816
3817fn choose_starter(
3829 review_branch: Option<&str>,
3830 branch_exists: bool,
3831 unfinished: Option<&str>,
3832) -> Starter {
3833 match review_branch {
3834 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3835 Some(_) => Starter::Start,
3836 None => match unfinished {
3837 Some(id) => Starter::Resume(id.to_owned()),
3838 None => Starter::Start,
3839 },
3840 }
3841}
3842
3843fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3846 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3847 return fallback.to_path_buf();
3848 }
3849 task.repo.clone()
3850}
3851
3852const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3856
3857fn answers_block(task: &Task, count: usize) -> String {
3859 let mut s = ANSWERS_HEADER.to_owned();
3860 for a in &task.answers[..count] {
3861 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3862 }
3863 s
3864}
3865
3866fn append_answers(base: &str, task: &Task) -> String {
3869 if task.answers.is_empty() {
3870 return base.to_owned();
3871 }
3872 let mut s = base.to_owned();
3873 s.push_str(&answers_block(task, task.answers.len()));
3874 s
3875}
3876
3877fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3881 for count in (1..=task.answers.len()).rev() {
3882 let block = answers_block(task, count);
3883 if let Some(base) = instruction.strip_suffix(&block) {
3884 return base;
3885 }
3886 }
3887 instruction
3888}
3889
3890fn instruction_for(task: &Task) -> String {
3898 append_answers(&task.instruction, task)
3899}
3900
3901fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3913 append_answers(strip_answers_block(old_instruction, task), task)
3914}
3915
3916fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3918 let paths = queue.attachment_paths(task);
3919 for (name, path) in task.attachments.iter().zip(&paths) {
3920 if !path.is_file() {
3921 bail!(
3922 "attachment `{name}` is recorded on the task but {} is missing",
3923 path.display()
3924 );
3925 }
3926 }
3927 Ok(paths)
3928}
3929
3930fn prepare_instruction(
3941 starter: &Starter,
3942 old_instruction: Option<&str>,
3943 task: &Task,
3944) -> Option<String> {
3945 match starter {
3946 Starter::Start => Some(instruction_for(task)),
3947 Starter::Resume(_) => Some(resumed_instruction(
3948 old_instruction.expect("a resumed run always has a prior instruction"),
3949 task,
3950 )),
3951 Starter::Review(_) => None,
3952 }
3953}
3954
3955fn record(queue: &Queue, task: &mut Task) {
3959 if let Err(e) = queue.put(task) {
3960 tracing::error!("could not record task {}: {e:#}", task.short());
3961 notices::raise(Notice::error(
3962 "loop:record",
3963 "The loop could not save a task's state; check the disk.",
3964 ));
3965 }
3966}
3967
3968fn runnable(queue: &Queue) -> Vec<Task> {
3974 let mut tasks: Vec<Task> = queue
3975 .list()
3976 .into_iter()
3977 .filter(|t| t.status.runnable())
3978 .collect();
3979 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
3980 tasks
3981}
3982
3983fn describe(state: &RunState) -> String {
3997 let p = phrases(&state.config.graph.language);
3998 let mut detail = if state.status == RunStatus::Stalled {
3999 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
4000 seats.sort_unstable();
4001 seats.dedup();
4002 if seats.is_empty() {
4003 p.quorum_lost.to_owned()
4004 } else {
4005 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
4006 }
4007 } else {
4008 format!("{}{}", p.run_ended, state.status.display_label())
4009 };
4010 if let Some(last) = state.events.last() {
4011 let unanswered = state
4015 .reviews
4016 .last()
4017 .filter(|r| {
4018 state.status == RunStatus::Blocked
4019 && last.node == "review"
4020 && r.incomplete()
4021 && r.blocking == 0
4022 && r.round == state.config.graph.review_rounds
4023 && r.e2e.iter().all(crate::run::CommandOutcome::ok)
4024 })
4025 .map(|r| (r.expected - r.answered, r.round));
4026 match unanswered {
4027 Some((missing, rounds)) if crate::lang::is_japanese(&state.config.graph.language) => {
4028 detail.push_str(&format!(
4029 " ({}: {})",
4030 last.node,
4031 (p.reviewers_never_answered)(missing, rounds)
4032 ));
4033 }
4034 _ => detail.push_str(&format!(" ({}: {})", last.node, last.message)),
4035 }
4036 }
4037 detail.push_str(&format!(" [run {}]", state.id));
4038 detail
4039}
4040
4041const DIAGNOSTIC_MAX: usize = 4_000;
4047
4048const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
4053
4054fn diagnostic(state: &RunState) -> Option<String> {
4068 let mut parts: Vec<String> = Vec::new();
4069
4070 for o in state.gate.iter().filter(|o| !o.ok()) {
4072 parts.push(format!(
4073 "gate `{}` failed ({:?}):\n{}",
4074 o.command,
4075 o.code,
4076 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
4077 ));
4078 }
4079
4080 if let Some(last) = state
4083 .events
4084 .iter()
4085 .rev()
4086 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
4087 {
4088 parts.push(last.message.clone());
4089 }
4090
4091 if state.viable().is_empty() {
4098 for c in &state.candidates {
4099 if let Some(evidence) = &c.verified_noop {
4100 parts.push(format!(
4101 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
4102 c.label
4103 ));
4104 } else if !c.summary.trim().is_empty() {
4105 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
4106 } else if let Some(why) = &c.failed {
4107 parts.push(format!("candidate {}: {why}", c.label));
4108 }
4109 }
4110 }
4111
4112 if parts.is_empty() {
4113 return None;
4114 }
4115 Some(crate::run::tail(
4120 &parts.join("\n\n"),
4121 DIAGNOSTIC_MAX.saturating_sub(100),
4122 ))
4123}
4124
4125fn label(status: RunStatus) -> &'static str {
4133 status.as_str()
4134}
4135
4136pub(crate) fn merge_mode(mode: &str) -> Result<MergeMode> {
4138 match mode {
4139 "none" => Ok(MergeMode::None),
4140 "local" => Ok(MergeMode::Local),
4141 "pr" => Ok(MergeMode::Pr),
4142 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
4143 }
4144}
4145
4146fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
4151 mutex
4152 .lock()
4153 .unwrap_or_else(std::sync::PoisonError::into_inner)
4154}
4155
4156#[cfg(test)]
4157mod tests {
4158 use super::*;
4159 use crate::queue::{Source, TaskStatus};
4160 use crate::run::{Candidate, CommandOutcome};
4161 use pretty_assertions::assert_eq;
4162
4163 fn task() -> Task {
4164 Task::new(
4165 "add retries".to_owned(),
4166 "add retries".to_owned(),
4167 PathBuf::from("/repo"),
4168 Source::Human,
4169 )
4170 }
4171
4172 fn interrupt_task(id: &str) -> Task {
4175 let mut t = task();
4176 t.id = id.to_owned();
4177 t.interrupt = true;
4178 t
4179 }
4180
4181 fn task_with_id(id: &str) -> Task {
4183 let mut t = task();
4184 t.id = id.to_owned();
4185 t
4186 }
4187
4188 fn urgent_task(id: &str) -> Task {
4190 let mut t = task();
4191 t.id = id.to_owned();
4192 t.urgent = true;
4193 t
4194 }
4195
4196 #[test]
4201 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
4202 assert_eq!(permit_kind(true, false), PermitKind::None);
4203 assert_eq!(permit_kind(true, true), PermitKind::None);
4204 }
4205
4206 #[test]
4211 fn permit_kind_separates_urgent_from_ordinary() {
4212 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
4213 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
4214 }
4215
4216 #[test]
4223 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
4224 let cfg = Config::default();
4225 let repo = Path::new("/any/repo/path");
4226
4227 let reason =
4228 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
4229 assert!(reason.contains("1024"), "{reason}");
4230 assert!(
4231 reason.contains(&cfg.disk.min_free_bytes.to_string()),
4232 "{reason}"
4233 );
4234
4235 assert_eq!(
4236 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
4237 None,
4238 "exactly at the floor is open"
4239 );
4240 assert_eq!(
4241 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
4242 None,
4243 "comfortably above the floor is open"
4244 );
4245 }
4246
4247 #[test]
4248 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
4249 let mut cfg = Config::default();
4250 cfg.disk.min_free_bytes = 0;
4251 let repo = Path::new("/any/repo/path");
4252 assert_eq!(
4253 disk_gate_with(repo, &cfg, |_| Ok(0)),
4254 None,
4255 "a zero floor never measures at all"
4256 );
4257 }
4258
4259 #[test]
4260 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
4261 let cfg = Config::default();
4262 let repo = Path::new("/any/repo/path");
4263 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
4264 .expect("a measurement failure must close the gate, not open it");
4265 assert!(reason.contains("could not measure"), "{reason}");
4266 }
4267
4268 #[test]
4269 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
4270 let ordinary = task();
4271 let next = advance_interrupt(
4272 Interrupt::Idle,
4273 std::slice::from_ref(&ordinary.id),
4274 std::slice::from_ref(&ordinary),
4275 );
4276 assert_eq!(next, Interrupt::Idle);
4277 }
4278
4279 #[test]
4280 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
4281 let marked = interrupt_task("marked");
4284 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4285 assert_eq!(next, Interrupt::Idle);
4286 }
4287
4288 #[test]
4289 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
4290 let marked = interrupt_task("marked");
4291 let next = advance_interrupt(
4292 Interrupt::Idle,
4293 &["running".to_owned()],
4294 std::slice::from_ref(&marked),
4295 );
4296 assert_eq!(
4297 next,
4298 Interrupt::Parking {
4299 parked: vec!["running".to_owned()],
4300 interrupt_task: "marked".to_owned(),
4301 }
4302 );
4303 }
4304
4305 #[test]
4314 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
4315 let marked = interrupt_task("marked");
4316
4317 let two = advance_interrupt(
4318 Interrupt::Idle,
4319 &["a".to_owned(), "b".to_owned()],
4320 std::slice::from_ref(&marked),
4321 );
4322 assert_eq!(two, Interrupt::Idle);
4323
4324 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4325 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
4326 }
4327
4328 #[test]
4329 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
4330 let state = Interrupt::Parking {
4331 parked: vec!["running".to_owned()],
4332 interrupt_task: "marked".to_owned(),
4333 };
4334 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
4336 assert_eq!(still_going, state);
4337
4338 let stopped_but_not_yet_dispatched =
4342 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
4343 assert_eq!(stopped_but_not_yet_dispatched, state);
4344
4345 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
4347 assert_eq!(
4348 dispatched,
4349 Interrupt::Running {
4350 parked: vec!["running".to_owned()],
4351 interrupt_task: "marked".to_owned(),
4352 }
4353 );
4354 }
4355
4356 #[test]
4357 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
4358 let state = Interrupt::Running {
4359 parked: vec!["running".to_owned()],
4360 interrupt_task: "marked".to_owned(),
4361 };
4362 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
4363 assert_eq!(still_running, state);
4364
4365 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
4372 assert_eq!(
4373 ended,
4374 Interrupt::Resuming {
4375 parked: vec!["running".to_owned()]
4376 }
4377 );
4378 }
4379
4380 #[test]
4381 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
4382 let state = Interrupt::Resuming {
4383 parked: vec!["running".to_owned()],
4384 };
4385 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
4386 assert_eq!(still_waiting, state);
4387
4388 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
4389 assert_eq!(dispatched, Interrupt::Idle);
4390 }
4391
4392 #[test]
4398 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
4399 {
4400 let state = Interrupt::Parking {
4401 parked: vec!["running".to_owned()],
4402 interrupt_task: "marked".to_owned(),
4403 };
4404 let next = advance_interrupt(state, &[], &[]);
4407 assert_eq!(
4408 next,
4409 Interrupt::Resuming {
4410 parked: vec!["running".to_owned()]
4411 },
4412 "abandoning the interrupt must not abandon the resume it owes"
4413 );
4414 }
4415
4416 #[test]
4419 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
4420 let state = Interrupt::Resuming {
4421 parked: vec!["running".to_owned()],
4422 };
4423 let next = advance_interrupt(state, &[], &[]);
4424 assert_eq!(
4425 next,
4426 Interrupt::Idle,
4427 "nothing is left to wait for; the loop must not stay wedged"
4428 );
4429 }
4430
4431 #[test]
4432 fn disabled_by_config_the_sequence_can_never_leave_idle() {
4433 let marked = interrupt_task("marked");
4434 let next = advance_interrupt_tick(
4435 false,
4436 Interrupt::Idle,
4437 &["running".to_owned()],
4438 std::slice::from_ref(&marked),
4439 );
4440 assert_eq!(
4441 next,
4442 Interrupt::Idle,
4443 "an unmarked, unconfigured daemon must behave exactly as before"
4444 );
4445 }
4446
4447 #[test]
4448 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4449 let state = Interrupt::Parking {
4450 parked: vec!["running".to_owned()],
4451 interrupt_task: "marked".to_owned(),
4452 };
4453 let candidates = vec![interrupt_task("marked"), task()];
4454 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4455 assert!(
4456 allowed.is_empty(),
4457 "nothing may dispatch - not even the interrupt task itself - \
4458 until the parked run has actually stopped"
4459 );
4460 }
4461
4462 #[test]
4476 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4477 for state in [
4478 Interrupt::Parking {
4479 parked: vec!["running".to_owned()],
4480 interrupt_task: "marked".to_owned(),
4481 },
4482 Interrupt::Running {
4483 parked: vec!["running".to_owned()],
4484 interrupt_task: "marked".to_owned(),
4485 },
4486 Interrupt::Resuming {
4487 parked: vec!["running".to_owned()],
4488 },
4489 ] {
4490 let candidates = vec![
4491 interrupt_task("marked"),
4492 urgent_task("hot"),
4493 task_with_id("ordinary"),
4494 ];
4495 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4496 assert!(
4497 !allowed.iter().any(|t| t.id == "hot"),
4498 "an urgent candidate must wait out the same gate as anything \
4499 else while the run it would run alongside has not actually \
4500 left flight, for state {state:?}: {allowed:?}"
4501 );
4502 }
4503 }
4504
4505 #[test]
4511 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4512 let state = Interrupt::Resuming {
4513 parked: vec!["hot".to_owned()],
4514 };
4515 let candidates = vec![urgent_task("hot"), task()];
4516 let allowed = interrupt_gate(&state, &[], candidates);
4517 assert_eq!(
4518 allowed.iter().filter(|t| t.id == "hot").count(),
4519 1,
4520 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4521 );
4522 }
4523
4524 #[test]
4525 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4526 let state = Interrupt::Parking {
4527 parked: vec!["running".to_owned()],
4528 interrupt_task: "marked".to_owned(),
4529 };
4530 let other = task();
4531 let candidates = vec![interrupt_task("marked"), other.clone()];
4532 let allowed = interrupt_gate(&state, &[], candidates);
4533 assert_eq!(allowed.len(), 1);
4534 assert_eq!(allowed[0].id, "marked");
4535 }
4536
4537 #[test]
4538 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4539 let state = Interrupt::Running {
4540 parked: vec!["running".to_owned()],
4541 interrupt_task: "marked".to_owned(),
4542 };
4543 let candidates = vec![task(), task()];
4544 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4545 assert!(allowed.is_empty());
4546 }
4547
4548 #[test]
4555 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4556 let state = Interrupt::Resuming {
4557 parked: vec!["a".to_owned(), "c".to_owned()],
4558 };
4559 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4560 let allowed = interrupt_gate(&state, &[], candidates);
4561 assert_eq!(
4562 allowed.len(),
4563 1,
4564 "at most one candidate may be offered while resuming: {allowed:?}"
4565 );
4566 assert_eq!(allowed[0].id, "a");
4567 }
4568
4569 #[test]
4570 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4571 let state = Interrupt::Resuming {
4572 parked: vec!["a".to_owned()],
4573 };
4574 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4575 assert!(allowed.is_empty());
4576 }
4577
4578 #[test]
4584 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4585 let running = task(); let marked = interrupt_task("marked");
4587
4588 let mut state = Interrupt::Idle;
4589 let in_flight = vec![running.id.clone()];
4591 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4592 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4593 assert!(gated.is_empty(), "still waiting on `running` to park");
4594
4595 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4597 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4598 assert_eq!(
4599 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4600 vec!["marked"],
4601 "only the interrupt task may be offered to the dispatcher now"
4602 );
4603
4604 state = advance_interrupt_tick(
4606 true,
4607 state,
4608 &["marked".to_owned()],
4609 std::slice::from_ref(&running),
4610 );
4611 let gated = interrupt_gate(
4612 &state,
4613 &["marked".to_owned()],
4614 vec![marked.clone(), running.clone()],
4615 );
4616 assert!(
4617 gated.is_empty(),
4618 "the parked run must not be offered back while the interrupt \
4619 task is still running"
4620 );
4621
4622 let other = task_with_id("other");
4626 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4627 assert_eq!(
4628 state,
4629 Interrupt::Resuming {
4630 parked: vec![running.id.clone()]
4631 }
4632 );
4633 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4634 assert_eq!(
4635 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4636 vec![running.id.as_str()],
4637 "exactly the parked run resumes - not the unrelated task, even \
4638 though it was offered first"
4639 );
4640
4641 state = advance_interrupt_tick(
4645 true,
4646 state,
4647 std::slice::from_ref(&running.id),
4648 std::slice::from_ref(&other),
4649 );
4650 assert_eq!(state, Interrupt::Idle);
4651 let gated = interrupt_gate(
4652 &state,
4653 std::slice::from_ref(&running.id),
4654 vec![other.clone()],
4655 );
4656 assert_eq!(
4657 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4658 vec![other.id.as_str()],
4659 "ordinary dispatch is unrestricted again"
4660 );
4661 }
4662
4663 #[test]
4664 fn every_run_status_settles_the_task_it_came_from() {
4665 let table = [
4667 (RunStatus::Merged, TaskStatus::Done, 1),
4668 (RunStatus::Ready, TaskStatus::Done, 1),
4669 (RunStatus::Stalled, TaskStatus::Failed, 0),
4670 (RunStatus::Blocked, TaskStatus::Failed, 1),
4671 (RunStatus::Failed, TaskStatus::Failed, 1),
4672 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4673 (RunStatus::Prep, TaskStatus::Failed, 1),
4674 (RunStatus::Implementing, TaskStatus::Failed, 1),
4675 (RunStatus::Judging, TaskStatus::Failed, 1),
4676 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4677 (RunStatus::Voting, TaskStatus::Failed, 1),
4678 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4679 (RunStatus::Gating, TaskStatus::Failed, 1),
4680 ];
4681 for (run, want, attempts) in table {
4682 let mut t = task();
4683 t.start("20260902-000000-aaaa".to_owned());
4684 settle(
4685 &mut t,
4686 Verdict {
4687 status: run,
4688 left_pr: false,
4689 parked: false,
4690 quota_hit: matches!(run, RunStatus::Stalled),
4691 no_viable_candidates: false,
4692 },
4693 "why",
4694 2,
4695 );
4696 assert_eq!(t.status, want, "task status after {}", label(run));
4697 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4698 }
4699 }
4700
4701 #[test]
4702 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4703 let mut stalled = task();
4704 stalled.start("20260902-000000-aaaa".to_owned());
4705 settle(
4706 &mut stalled,
4707 Verdict {
4708 status: RunStatus::Stalled,
4709 left_pr: false,
4710 parked: false,
4711 quota_hit: true,
4712 no_viable_candidates: false,
4713 },
4714 "quota",
4715 1,
4716 );
4717 assert_eq!(stalled.attempts, 0);
4718 assert!(
4719 stalled.status.runnable(),
4720 "a machine problem must leave the task in line"
4721 );
4722
4723 let mut blocked = task();
4724 blocked.start("20260902-000000-aaaa".to_owned());
4725 settle(
4726 &mut blocked,
4727 Verdict {
4728 status: RunStatus::Blocked,
4729 left_pr: false,
4730 parked: false,
4731 quota_hit: false,
4732 no_viable_candidates: false,
4733 },
4734 "findings open",
4735 1,
4736 );
4737 assert_eq!(blocked.attempts, 1);
4738 assert_eq!(
4739 blocked.status,
4740 TaskStatus::Held,
4741 "the last attempt hands the task to a human"
4742 );
4743 }
4744
4745 #[test]
4746 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4747 let mut delivered = task();
4750 delivered.start("20260903-080619-01c2".to_owned());
4751 settle(
4752 &mut delivered,
4753 Verdict {
4754 status: RunStatus::Blocked,
4755 left_pr: true,
4756 parked: false,
4757 quota_hit: false,
4758 no_viable_candidates: false,
4759 },
4760 "no check status",
4761 4,
4762 );
4763 assert_eq!(
4764 delivered.status,
4765 TaskStatus::Held,
4766 "a pull request waiting on CI or a person is not a retryable failure"
4767 );
4768 assert!(
4769 !delivered.status.runnable(),
4770 "the loop must not pick this task up again"
4771 );
4772 assert_eq!(
4773 delivered.last_error.as_deref(),
4774 Some("no check status"),
4775 "the operator needs to be told what the gate was waiting for"
4776 );
4777
4778 let mut empty_handed = task();
4781 empty_handed.start("20260903-080619-01c2".to_owned());
4782 settle(
4783 &mut empty_handed,
4784 Verdict {
4785 status: RunStatus::Blocked,
4786 left_pr: false,
4787 parked: false,
4788 quota_hit: false,
4789 no_viable_candidates: false,
4790 },
4791 "findings open",
4792 4,
4793 );
4794 assert_eq!(empty_handed.status, TaskStatus::Failed);
4795 assert!(empty_handed.status.runnable());
4796 }
4797
4798 #[test]
4799 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4800 let mut noop = task();
4806 noop.start("20260912-131304-391f".to_owned());
4807 settle(
4808 &mut noop,
4809 Verdict {
4810 status: RunStatus::VerifiedNoop,
4811 left_pr: false,
4812 parked: false,
4813 quota_hit: false,
4814 no_viable_candidates: true,
4815 },
4816 "candidate A: already fixed by b32cfc4, on main",
4817 4,
4818 );
4819 assert_eq!(
4820 noop.status,
4821 TaskStatus::Held,
4822 "an unverified claim is a request for a human, not a failure"
4823 );
4824 assert!(
4825 !noop.status.runnable(),
4826 "the loop must not requeue this on the same unverified claim"
4827 );
4828 assert_eq!(noop.attempts, 1);
4833 }
4834
4835 #[test]
4836 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4837 let mut parked = task();
4842 parked.start("20260903-183634-2d98".to_owned());
4843 settle(
4844 &mut parked,
4845 Verdict {
4846 status: RunStatus::Implementing,
4847 left_pr: false,
4848 quota_hit: false,
4849 parked: true,
4850 no_viable_candidates: false,
4851 },
4852 "parked after `implementing`",
4853 2,
4854 );
4855 assert_eq!(parked.attempts, 0, "a park is refunded");
4856 assert!(
4857 parked.status.runnable(),
4858 "and the task stays in line so the next loop resumes its run"
4859 );
4860 assert_eq!(
4861 parked.last_error.as_deref(),
4862 Some("parked after `implementing`"),
4863 "the card says where it stopped"
4864 );
4865
4866 let mut broken = task();
4870 broken.start("20260903-183634-2d98".to_owned());
4871 settle(
4872 &mut broken,
4873 Verdict {
4874 status: RunStatus::Implementing,
4875 left_pr: false,
4876 quota_hit: false,
4877 parked: false,
4878 no_viable_candidates: false,
4879 },
4880 "returned mid-flight",
4881 2,
4882 );
4883 assert_eq!(broken.attempts, 1);
4884 }
4885
4886 #[test]
4887 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4888 let mut flaky = task();
4893 flaky.start("20260903-123023-e633".to_owned());
4894 settle(
4895 &mut flaky,
4896 Verdict {
4897 status: RunStatus::Stalled,
4898 left_pr: false,
4899 parked: false,
4900 quota_hit: false,
4901 no_viable_candidates: false,
4902 },
4903 "verdict rests on 1 of 3 judges",
4904 2,
4905 );
4906 assert_eq!(
4907 flaky.attempts, 1,
4908 "flakiness spends an attempt, so `max_attempts` still bounds it"
4909 );
4910 assert!(flaky.status.runnable(), "and it is still worth retrying");
4911
4912 let mut limited = task();
4914 limited.start("20260903-123023-e633".to_owned());
4915 settle(
4916 &mut limited,
4917 Verdict {
4918 status: RunStatus::Stalled,
4919 left_pr: false,
4920 parked: false,
4921 quota_hit: true,
4922 no_viable_candidates: false,
4923 },
4924 "judge-2, judge-3 out of quota",
4925 2,
4926 );
4927 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4928 assert!(limited.status.runnable());
4929
4930 let mut worn = task();
4933 for _ in 0..2 {
4934 worn.release();
4935 }
4936 worn.start("20260903-123023-e633".to_owned());
4937 worn.attempts = 2;
4938 settle(
4939 &mut worn,
4940 Verdict {
4941 status: RunStatus::Stalled,
4942 left_pr: false,
4943 parked: false,
4944 quota_hit: false,
4945 no_viable_candidates: false,
4946 },
4947 "no quorum again",
4948 2,
4949 );
4950 assert_eq!(worn.status, TaskStatus::Held);
4951 assert!(!worn.status.runnable());
4952 }
4953
4954 #[test]
4955 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4956 let mut wiped_out = task();
4963 wiped_out.start("20260907-025000-a1b2".to_owned());
4964 settle(
4965 &mut wiped_out,
4966 Verdict {
4967 status: RunStatus::Failed,
4968 left_pr: false,
4969 parked: false,
4970 quota_hit: true,
4971 no_viable_candidates: true,
4972 },
4973 "no candidate produced a change; nothing to judge",
4974 2,
4975 );
4976 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
4977 assert!(
4978 wiped_out.status.runnable(),
4979 "a machine problem must leave the task in line"
4980 );
4981
4982 let mut partial_progress = task();
4988 partial_progress.start("20260907-025500-c3d4".to_owned());
4989 settle(
4990 &mut partial_progress,
4991 Verdict {
4992 status: RunStatus::Failed,
4993 left_pr: false,
4994 parked: false,
4995 quota_hit: true,
4996 no_viable_candidates: false,
4997 },
4998 "gate failed on the winning candidate",
4999 2,
5000 );
5001 assert_eq!(
5002 partial_progress.attempts, 1,
5003 "a candidate that actually produced a change spends the attempt \
5004 even though some other seat hit its quota"
5005 );
5006 assert!(partial_progress.status.runnable());
5007 }
5008
5009 #[test]
5010 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
5011 let mut t = task();
5017 t.start("20260907-025000-a1b2".to_owned());
5018 let mut state = run_state(RunStatus::Failed);
5019 state.quota.push(QuotaLoss {
5020 seat: "cand-a".to_owned(),
5021 node: "implement".to_owned(),
5022 at: Timestamp::now(),
5023 reset: None,
5024 });
5025 assert!(
5026 state.viable().is_empty(),
5027 "no candidate was added, so nothing is viable"
5028 );
5029 reclaim(&mut t, Some(state), 2, "en");
5030 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
5031 assert!(t.status.runnable());
5032 }
5033
5034 #[test]
5035 fn a_held_task_is_never_offered_to_the_loop() {
5036 let dir = tempfile::tempdir().unwrap();
5037 let queue = Queue::at(dir.path().to_path_buf());
5038 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
5039 let mut t = task();
5040 t.id = format!("2026090{n}-000000-000{n}");
5041 t.priority = priority;
5042 queue.put(&mut t).unwrap();
5043 }
5044 let mut held = task();
5045 held.id = "20260909-000000-9999".to_owned();
5046 held.priority = 99;
5047 held.hold_machine(None);
5048 queue.put(&mut held).unwrap();
5049
5050 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
5051 assert_eq!(order.len(), 3);
5052 assert!(!order.contains(&held.id));
5053 assert_eq!(
5054 order.first().cloned(),
5055 queue.next_runnable().map(|t| t.id),
5056 "the loop's first candidate is exactly what the queue offers"
5057 );
5058 assert_eq!(
5059 order,
5060 vec![
5061 "20260902-000000-0002".to_owned(),
5062 "20260903-000000-0003".to_owned(),
5063 "20260901-000000-0001".to_owned(),
5064 ],
5065 "priority first, then oldest, so nothing starves"
5066 );
5067 }
5068
5069 #[test]
5070 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
5071 let dir = tempfile::tempdir().unwrap();
5072 let queue = Queue::at(dir.path().to_path_buf());
5073 let mut old = task();
5074 old.id = "20260101-000000-old0".to_owned();
5075 queue.put(&mut old).unwrap();
5076 let mut fresh = task();
5077 fresh.id = "20260101-000000-new0".to_owned();
5078 queue.put(&mut fresh).unwrap();
5079
5080 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
5084 std::thread::sleep(Duration::from_millis(60));
5085 let live = queue.claim(&fresh.id).unwrap();
5086
5087 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
5088 assert_eq!(swept, vec![old.id.clone()]);
5089 assert!(
5090 queue.claim(&old.id).is_ok(),
5091 "an unparseable lock older than the threshold is swept"
5092 );
5093 assert!(
5094 queue.claim(&fresh.id).is_err(),
5095 "a live pid protects its lock regardless of age"
5096 );
5097 drop(live);
5098 }
5099
5100 #[test]
5101 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
5102 let dir = tempfile::tempdir().unwrap();
5112 let queue = Queue::at(dir.path().to_path_buf());
5113 let mut t = task();
5114 t.id = "20260101-000000-live".to_owned();
5115 queue.put(&mut t).unwrap();
5116
5117 let claim = queue.claim(&t.id).unwrap();
5118 std::thread::sleep(Duration::from_millis(60));
5119
5120 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
5121 assert!(
5122 swept.is_empty(),
5123 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
5124 );
5125 assert!(
5126 queue.claim(&t.id).is_err(),
5127 "the lock still protects its task"
5128 );
5129 drop(claim);
5130 }
5131
5132 fn injected_dead_pid() -> u32 {
5135 std::process::id().checked_add(1).unwrap_or(1)
5136 }
5137
5138 #[test]
5139 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
5140 let dir = tempfile::tempdir().unwrap();
5141 let queue = Queue::at(dir.path().to_path_buf());
5142 let mut t = task();
5143 t.id = "20260101-000000-dead".to_owned();
5144 queue.put(&mut t).unwrap();
5145 let dead_pid = injected_dead_pid();
5146
5147 std::fs::write(
5152 dir.path().join(format!("{}.lock", t.id)),
5153 dead_pid.to_string(),
5154 )
5155 .unwrap();
5156
5157 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5158 pid != dead_pid
5159 });
5160 assert_eq!(
5161 swept,
5162 vec![t.id.clone()],
5163 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
5164 );
5165 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
5166 }
5167
5168 #[test]
5169 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
5170 let dir = tempfile::tempdir().unwrap();
5171 let queue = Queue::at(dir.path().to_path_buf());
5172 let mut t = task();
5173 t.id = "20260101-000000-late".to_owned();
5174 queue.put(&mut t).unwrap();
5175 let dead_pid = injected_dead_pid();
5176
5177 assert!(
5180 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
5181 "nothing has claimed the task yet"
5182 );
5183
5184 std::fs::write(
5187 dir.path().join(format!("{}.lock", t.id)),
5188 dead_pid.to_string(),
5189 )
5190 .unwrap();
5191
5192 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5196 pid != dead_pid
5197 });
5198 assert_eq!(swept, vec![t.id.clone()]);
5199 }
5200
5201 #[test]
5202 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
5203 crate::run::pin_test_home();
5208 let dir = tempfile::tempdir().unwrap();
5209 let queue = Queue::at(dir.path().to_path_buf());
5210 let mut t = task();
5211 t.id = "20260101-000000-crsh".to_owned();
5212 t.status = TaskStatus::Running;
5213 t.attempts = 1;
5214 t.runs.push("20260904-000000-4043".to_owned());
5218 queue.put(&mut t).unwrap();
5219 let dead_pid = injected_dead_pid();
5220
5221 std::fs::write(
5224 dir.path().join(format!("{}.lock", t.id)),
5225 dead_pid.to_string(),
5226 )
5227 .unwrap();
5228
5229 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
5235 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
5236
5237 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5238 pid != dead_pid
5239 });
5240 assert_eq!(swept, vec![t.id.clone()]);
5241
5242 let reclaimed = reclaim_orphaned_running(&queue, 2);
5243 assert_eq!(reclaimed, vec![t.id.clone()]);
5244 let after = queue.get(&t.id).unwrap();
5245 assert_eq!(
5246 after.status,
5247 TaskStatus::Held,
5248 "no run.json to recover from, so a human is asked"
5249 );
5250 assert_eq!(
5251 after.runs,
5252 vec!["20260904-000000-4043".to_owned()],
5253 "the crashed run's id is kept as evidence, not discarded"
5254 );
5255 }
5256
5257 #[test]
5258 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
5259 let dir = tempfile::tempdir().unwrap();
5260 let queue = Queue::at(dir.path().to_path_buf());
5261 let mut t = task();
5262 t.id = "20260101-000000-unknown".to_owned();
5263 queue.put(&mut t).unwrap();
5264 let dead_pid = injected_dead_pid();
5265 std::fs::write(
5266 dir.path().join(format!("{}.lock", t.id)),
5267 dead_pid.to_string(),
5268 )
5269 .unwrap();
5270
5271 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
5272 assert!(swept.is_empty(), "an unknown pid must keep its lock");
5273 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
5274 }
5275
5276 fn run_state_in(status: RunStatus, language: &str) -> RunState {
5277 let mut s = run_state(status);
5278 s.config.graph.language = language.to_owned();
5279 s
5280 }
5281
5282 fn unstarted_verdict(status: RunStatus) -> Verdict {
5283 Verdict {
5284 status,
5285 left_pr: false,
5286 quota_hit: false,
5287 parked: false,
5288 no_viable_candidates: false,
5289 }
5290 }
5291
5292 #[test]
5293 fn an_already_in_base_run_finishes_the_task_without_spending_an_attempt() {
5294 let mut t = task();
5295 t.attempts = 1;
5296 settle_in(
5297 &mut t,
5298 unstarted_verdict(RunStatus::AlreadyInBase),
5299 "already in main",
5300 1,
5301 phrases("en"),
5302 );
5303 assert_eq!(t.status, TaskStatus::Done);
5304 assert_eq!(t.attempts, 0);
5305 }
5306
5307 #[test]
5308 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
5309 let reason = |language: &str| {
5310 let mut t = task();
5311 settle_in(
5312 &mut t,
5313 unstarted_verdict(RunStatus::Judging),
5314 "boom",
5315 1,
5316 phrases(language),
5317 );
5318 t.last_error.or(t.hold_reason).unwrap_or_default()
5319 };
5320 assert!(
5321 reason("en").starts_with("the graph stopped at `"),
5322 "{}",
5323 reason("en")
5324 );
5325 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
5326 assert!(reason("日本語").contains("boom"));
5327 assert_eq!(reason("fr"), reason("en"));
5328 }
5329
5330 #[test]
5331 fn describe_follows_the_run_language_and_keeps_the_run_id() {
5332 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
5333 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
5334 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
5335 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
5336 assert!(ja.contains("[run "), "{ja}");
5337 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
5338 assert!(ended.starts_with("run 終了: "), "{ended}");
5339 let mut de = run_state_in(RunStatus::Failed, "de");
5340 let mut en = run_state_in(RunStatus::Failed, "en");
5341 de.id = "same".to_owned();
5342 en.id = "same".to_owned();
5343 assert_eq!(describe(&de), describe(&en));
5344 }
5345
5346 #[test]
5347 fn handover_refusals_follow_the_language_and_keep_the_detail_apart() {
5348 use crate::handover::Refused;
5349 let cases = [
5350 Refused::Foreign {
5351 branch: "b".into(),
5352 path: "/w/x".into(),
5353 why: "made by hand".into(),
5354 },
5355 Refused::Unsafe {
5356 branch: "b".into(),
5357 path: "/w/x".into(),
5358 why: "its worktree has uncommitted changes (a.rs)".into(),
5359 },
5360 Refused::ReleaseFailed {
5361 branch: "b".into(),
5362 path: "/w/x".into(),
5363 run: "ab12".into(),
5364 },
5365 ];
5366 for r in &cases {
5367 let en = (phrases("en").handover_refused)(r);
5368 assert_eq!(en, r.to_string());
5369 let ja = (phrases("ja").handover_refused)(r);
5370 assert!(
5371 !ja.contains("is checked out") && !ja.contains("try again"),
5372 "{ja}"
5373 );
5374 assert!(ja.contains("`b`") && ja.contains("/w/x"), "{ja}");
5375 if let Refused::Foreign { why, .. } | Refused::Unsafe { why, .. } = r {
5376 assert!(ja.contains(&format!("(詳細: {why})")), "{ja}");
5377 }
5378 }
5379 assert!(phrases("en").handover_hint.contains("release the task"));
5380 assert!(phrases("ja").handover_hint.contains("解放"));
5381 }
5382
5383 #[test]
5384 fn describe_translates_the_unanswered_reviewer_stop_only_in_ja() {
5385 let build = |lang: &str| {
5386 let mut s = run_state_in(RunStatus::Blocked, lang);
5387 s.id = "same".to_owned();
5388 s.config.graph.review_rounds = 3;
5389 let mut r = review_round(3);
5390 r.expected = 3;
5391 r.answered = 1;
5392 s.reviews.push(r);
5393 s.event(
5394 "review",
5395 "2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean",
5396 );
5397 s
5398 };
5399 let en = describe(&build("en"));
5400 assert!(en.contains("(review: 2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean)"), "{en}");
5401 let ja = describe(&build("ja"));
5402 assert!(ja.contains("2 席のレビュアーが 3 ラウンド"), "{ja}");
5403 assert!(!ja.contains("never answered"), "{ja}");
5404 assert!(ja.contains("[run "), "{ja}");
5405
5406 let mut failed = build("ja");
5409 failed.reviews[0].e2e.push(crate::run::CommandOutcome {
5410 command: "cargo test".to_owned(),
5411 code: Some(1),
5412 output_tail: String::new(),
5413 duration_ms: 0,
5414 resource_blocked: false,
5415 });
5416 failed.event("review", "stopped; e2e failed: cargo test");
5417 let ja = describe(&failed);
5418 assert!(ja.contains("stopped; e2e failed: cargo test"), "{ja}");
5419 assert!(!ja.contains("席のレビュアー"), "{ja}");
5420 }
5421
5422 #[test]
5423 fn refusals_and_recovery_prose_follow_the_language() {
5424 let t = held_task_with("r1");
5425 let q = action_question(
5426 "r1",
5427 ask::ChoiceAction::Resume {
5428 run: "r1".to_owned(),
5429 },
5430 );
5431 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
5432 ActionDecision::Refuse(s) => s,
5433 other => panic!("{other:?}"),
5434 };
5435 assert!(refuse(phrases("en")).contains("could not be read: gone"));
5436 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
5437 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
5438
5439 let mut held = task();
5440 reclaim(&mut held, None, 2, "ja");
5441 assert!(held.hold_reason.unwrap().contains("保留にしました"));
5442 let mut held = task();
5443 reclaim(&mut held, None, 2, "xx");
5444 assert!(held.hold_reason.unwrap().contains("held for a human"));
5445 }
5446
5447 fn run_state(status: RunStatus) -> RunState {
5448 let mut state = RunState::new(
5449 PathBuf::from("/repo"),
5450 "main".to_owned(),
5451 "abc1234def".to_owned(),
5452 "add retries".to_owned(),
5453 Config::default(),
5454 );
5455 state.status = status;
5456 state
5457 }
5458
5459 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
5460 Candidate {
5461 index: 0,
5462 label,
5463 agent: "claude".to_owned(),
5464 branch: format!("magi/x/{label}"),
5465 worktree: PathBuf::from("/repo"),
5466 summary: summary.to_owned(),
5467 stat: String::new(),
5468 files: 0,
5469 commits: usize::from(!empty),
5470 empty,
5471 failed: failed.map(str::to_owned),
5472 verified_noop: None,
5473 duration_ms: 0,
5474 folded: false,
5475 }
5476 }
5477
5478 #[test]
5479 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
5480 let mut state = run_state(RunStatus::Blocked);
5481 state.gate = vec![
5482 CommandOutcome {
5483 command: "cargo make check".to_owned(),
5484 code: Some(0),
5485 output_tail: "ok".to_owned(),
5486 duration_ms: 0,
5487 resource_blocked: false,
5488 },
5489 CommandOutcome {
5490 command: "cargo test".to_owned(),
5491 code: Some(101),
5492 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
5493 duration_ms: 0,
5494 resource_blocked: false,
5495 },
5496 ];
5497 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
5498 assert!(d.contains("cargo test"), "{d}");
5499 assert!(
5500 !d.contains("cargo make check"),
5501 "a passing check is not a diagnostic: {d}"
5502 );
5503 assert!(d.contains("assertion failed"), "{d}");
5504 }
5505
5506 #[test]
5507 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
5508 let mut state = run_state(RunStatus::Blocked);
5509 state.event(
5510 "land",
5511 "stopped: the fixer produced no commit while 2 check(s) were failing \
5512 (build, lint); stopping instead of looping on an unchanged tree",
5513 );
5514 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
5515 assert!(d.contains("build"), "{d}");
5516 assert!(d.contains("lint"), "{d}");
5517 assert!(d.contains("fixer produced no commit"), "{d}");
5518 }
5519
5520 #[test]
5521 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
5522 let state = run_state(RunStatus::VerifiedNoop);
5527 let d = describe(&state);
5528 assert!(
5529 d.contains("agent-verified no-op"),
5530 "expected the display label, not the wire spelling: {d}"
5531 );
5532 assert!(!d.contains("verified_noop"), "{d}");
5533 }
5534
5535 #[test]
5536 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5537 let mut state = run_state(RunStatus::Failed);
5543 state.candidates = vec![candidate(
5544 'A',
5545 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5546 true,
5547 None,
5548 )];
5549 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5550 assert!(d.contains("candidate A"), "{d}");
5551 assert!(d.contains("tagged v1.2.3"), "{d}");
5552 }
5553
5554 #[test]
5555 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5556 let mut state = run_state(RunStatus::Failed);
5557 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5558 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5559 assert!(d.contains("candidate A"), "{d}");
5560 assert!(d.contains("agent timed out"), "{d}");
5561 }
5562
5563 #[test]
5564 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5565 let mut state = run_state(RunStatus::Failed);
5568 state.candidates = vec![candidate('A', "did the work", false, None)];
5569 assert!(diagnostic(&state).is_none());
5570 }
5571
5572 #[test]
5573 fn diagnostic_is_bounded_however_much_a_run_printed() {
5574 let mut state = run_state(RunStatus::Blocked);
5575 state.gate = vec![
5576 CommandOutcome {
5577 command: "cargo test".to_owned(),
5578 code: Some(101),
5579 output_tail: "x".repeat(50_000),
5580 duration_ms: 0,
5581 resource_blocked: false,
5582 },
5583 CommandOutcome {
5584 command: "cargo clippy".to_owned(),
5585 code: Some(1),
5586 output_tail: "y".repeat(50_000),
5587 duration_ms: 0,
5588 resource_blocked: false,
5589 },
5590 ];
5591 state.candidates = vec![
5592 candidate('A', &"z".repeat(50_000), true, None),
5593 candidate('B', &"w".repeat(50_000), true, None),
5594 ];
5595 let d = diagnostic(&state).expect("plenty here to diagnose");
5596 assert!(
5597 d.len() <= DIAGNOSTIC_MAX,
5598 "diagnostic grew to {} bytes, unbounded",
5599 d.len()
5600 );
5601 }
5602
5603 #[test]
5604 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5605 let mut state = run_state(RunStatus::Blocked);
5606 state.gate = vec![CommandOutcome {
5607 command: "cargo test".to_owned(),
5608 code: Some(101),
5609 output_tail: "assertion failed".to_owned(),
5610 duration_ms: 0,
5611 resource_blocked: false,
5612 }];
5613 let verdict = Verdict {
5614 status: RunStatus::Blocked,
5615 left_pr: false,
5616 quota_hit: false,
5617 parked: false,
5618 no_viable_candidates: false,
5619 };
5620
5621 let mut t = task();
5624 t.start("run-1".to_owned());
5625 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5626 assert_eq!(t.status, TaskStatus::Failed);
5627 assert!(t.diagnostic.is_none());
5628
5629 t.start("run-2".to_owned());
5632 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5633 assert_eq!(t.status, TaskStatus::Held);
5634 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5635 assert!(d.contains("cargo test"), "{d}");
5636 }
5637
5638 #[test]
5639 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5640 crate::run::pin_test_home();
5645 let home = crate::run::home();
5646 let state = run_state(RunStatus::VerifiedNoop);
5647 let mut q = ask::Question::new(
5648 state.id.clone(),
5649 "implement".to_owned(),
5650 "impl-A".to_owned(),
5651 "is this really a no-op?".to_owned(),
5652 String::new(),
5653 Vec::new(),
5654 );
5655 Questions::at(home.join("questions")).put(&mut q).unwrap();
5656
5657 let verdict = Verdict {
5658 status: RunStatus::VerifiedNoop,
5659 left_pr: false,
5660 quota_hit: false,
5661 parked: false,
5662 no_viable_candidates: false,
5663 };
5664 let mut t = task();
5665 t.start(state.id.clone());
5666 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5667
5668 assert_eq!(t.status, TaskStatus::Held);
5669 let reason = t.hold_reason.expect("a held task must record why");
5670 assert!(
5671 reason.starts_with("run ended agent-verified no-op"),
5672 "the original settle reason must survive unchanged: {reason}"
5673 );
5674 assert!(
5675 reason.contains(q.short()),
5676 "the open question's id must be named so the notice is actionable: {reason}"
5677 );
5678 }
5679
5680 #[test]
5681 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5682 crate::run::pin_test_home();
5683 let state = run_state(RunStatus::VerifiedNoop);
5684
5685 let verdict = Verdict {
5686 status: RunStatus::VerifiedNoop,
5687 left_pr: false,
5688 quota_hit: false,
5689 parked: false,
5690 no_viable_candidates: false,
5691 };
5692 let mut t = task();
5693 t.start(state.id.clone());
5694 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5695
5696 assert_eq!(t.status, TaskStatus::Held);
5697 assert_eq!(
5698 t.hold_reason.as_deref(),
5699 Some("run ended agent-verified no-op"),
5700 "nothing to append when the question was already answered or never asked"
5701 );
5702 }
5703
5704 #[test]
5705 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5706 crate::run::pin_test_home();
5707 let mut first = run_state(RunStatus::Blocked);
5708 first.id = "20260101-000000-sup1".to_owned();
5709 first.save().unwrap();
5710 let mut second = run_state(RunStatus::Merged);
5711 second.id = "20260101-000000-sup2".to_owned();
5712 second.save().unwrap();
5713
5714 let mut t = task();
5715 t.runs = vec![first.id.clone(), second.id.clone()];
5716 t.status = TaskStatus::Done;
5717
5718 supersede_prior_runs(&t, &crate::run::home());
5719
5720 assert_eq!(
5721 RunState::load(&first.id).unwrap().status,
5722 RunStatus::Superseded,
5723 "the first attempt's Blocked no longer needs anyone's attention"
5724 );
5725 assert_eq!(
5726 RunState::load(&second.id).unwrap().status,
5727 RunStatus::Merged,
5728 "the run that actually succeeded is left exactly as it was"
5729 );
5730 }
5731
5732 #[test]
5733 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5734 crate::run::pin_test_home();
5740 let mut first = run_state(RunStatus::Blocked);
5741 first.id = "20260101-000000-sup9".to_owned();
5742 first.driver_pid = Some(std::process::id());
5745 first.driver_started_at = Some(
5746 crate::proc::process_started_at(std::process::id())
5747 .expect("this test process's own start time must be queryable"),
5748 );
5749 first.save().unwrap();
5750 let mut second = run_state(RunStatus::Merged);
5751 second.id = "20260101-000000-supa".to_owned();
5752 second.save().unwrap();
5753
5754 let mut t = task();
5755 t.runs = vec![first.id.clone(), second.id.clone()];
5756 t.status = TaskStatus::Done;
5757
5758 supersede_prior_runs(&t, &crate::run::home());
5759
5760 assert_eq!(
5761 RunState::load(&first.id).unwrap().status,
5762 RunStatus::Blocked,
5763 "a live driver_pid means something is still actually working this run, \
5764 even though no daemon claims it - rewriting under it would just be \
5765 undone the next time that process saves"
5766 );
5767 }
5768
5769 #[test]
5770 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5771 let dir = tempfile::tempdir().unwrap();
5777 let home = dir.path().to_path_buf();
5778 let queue = Queue::at(dir.path().join("queue"));
5779
5780 let mut first = run_state(RunStatus::Blocked);
5781 first.id = "20260101-000000-supd".to_owned();
5782 first.driver_pid = Some(std::process::id());
5783 first.driver_started_at = Some(
5784 crate::proc::process_started_at(std::process::id())
5785 .expect("this test process's own start time must be queryable"),
5786 );
5787 first.save_under(&home).unwrap();
5788 let mut second = run_state(RunStatus::Merged);
5789 second.id = "20260101-000000-supe".to_owned();
5790 second.save_under(&home).unwrap();
5791
5792 let mut t = task();
5793 t.runs = vec![first.id.clone(), second.id.clone()];
5794 t.status = TaskStatus::Done;
5795 queue.put(&mut t).unwrap();
5796
5797 resweep_superseded_attempts(&queue, &home);
5798 assert_eq!(
5799 RunState::load_under(&first.id, &home).unwrap().status,
5800 RunStatus::Blocked,
5801 "still live on the first pass, so still untouched"
5802 );
5803
5804 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5810 stale.driver_started_at = Some("1".to_owned());
5811 stale.save_under(&home).unwrap();
5812
5813 resweep_superseded_attempts(&queue, &home);
5814 assert_eq!(
5815 RunState::load_under(&first.id, &home).unwrap().status,
5816 RunStatus::Superseded,
5817 "the second pass catches up what the first one correctly skipped"
5818 );
5819 }
5820
5821 #[test]
5822 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5823 crate::run::pin_test_home();
5824 let mut first = run_state(RunStatus::Blocked);
5825 first.id = "20260101-000000-sup3".to_owned();
5826 first.save().unwrap();
5827 let mut second = run_state(RunStatus::Blocked);
5828 second.id = "20260101-000000-sup4".to_owned();
5829 second.save().unwrap();
5830
5831 let mut t = task();
5832 t.runs = vec![first.id.clone(), second.id.clone()];
5833 t.status = TaskStatus::Failed;
5837
5838 supersede_prior_runs(&t, &crate::run::home());
5839
5840 assert_eq!(
5841 RunState::load(&first.id).unwrap().status,
5842 RunStatus::Blocked
5843 );
5844 assert_eq!(
5845 RunState::load(&second.id).unwrap().status,
5846 RunStatus::Blocked
5847 );
5848 }
5849
5850 #[test]
5851 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5852 crate::run::pin_test_home();
5856 let mut first = run_state(RunStatus::Blocked);
5857 first.id = "20260101-000000-sup5".to_owned();
5858 first.save().unwrap();
5859
5860 let mut t = task();
5861 t.runs = vec![first.id.clone()];
5862 t.status = TaskStatus::Done;
5863
5864 supersede_prior_runs(&t, &crate::run::home());
5865
5866 assert_eq!(
5867 RunState::load(&first.id).unwrap().status,
5868 RunStatus::Blocked,
5869 "a single-attempt task has no earlier run to supersede"
5870 );
5871 }
5872
5873 #[test]
5874 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5875 crate::run::pin_test_home();
5882 let mut first = run_state(RunStatus::Blocked);
5883 first.id = "20260101-000000-supb".to_owned();
5884 first.save().unwrap();
5885 let mut second = run_state(RunStatus::Failed);
5886 second.id = "20260101-000000-supc".to_owned();
5887 second.save().unwrap();
5888
5889 let mut t = task();
5890 t.runs = vec![first.id.clone(), second.id.clone()];
5891 t.status = TaskStatus::Done;
5892
5893 supersede_prior_runs(&t, &crate::run::home());
5894
5895 assert_eq!(
5896 RunState::load(&first.id).unwrap().status,
5897 RunStatus::Blocked,
5898 "the task's last attempt never landed, so there is nothing here \
5899 actually superseding it"
5900 );
5901 }
5902
5903 #[test]
5904 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5905 crate::run::pin_test_home();
5909 let mut failed = run_state(RunStatus::Failed);
5910 failed.id = "20260101-000000-sup6".to_owned();
5911 failed.save().unwrap();
5912 let mut noop = run_state(RunStatus::VerifiedNoop);
5913 noop.id = "20260101-000000-sup7".to_owned();
5914 noop.save().unwrap();
5915 let mut winner = run_state(RunStatus::Ready);
5916 winner.id = "20260101-000000-sup8".to_owned();
5917 winner.save().unwrap();
5918
5919 let mut t = task();
5920 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5921 t.status = TaskStatus::Done;
5922
5923 supersede_prior_runs(&t, &crate::run::home());
5924
5925 assert_eq!(
5926 RunState::load(&failed.id).unwrap().status,
5927 RunStatus::Failed
5928 );
5929 assert_eq!(
5930 RunState::load(&noop.id).unwrap().status,
5931 RunStatus::VerifiedNoop
5932 );
5933 }
5934
5935 fn approval_question(run: &str) -> ask::Question {
5936 ask::Question::new(
5937 run.to_owned(),
5938 land::APPROVAL_NODE.to_owned(),
5939 "land".to_owned(),
5940 "merge?".to_owned(),
5941 String::new(),
5942 vec!["merge".to_owned(), "hold".to_owned()],
5943 )
5944 }
5945
5946 #[test]
5947 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5948 crate::run::pin_test_home();
5949 let mut state = run_state(RunStatus::Landing);
5950 state.id = "20260101-000000-fre1".to_owned();
5951 state.parked = true;
5952 state.save().unwrap();
5953 ask::Questions::open()
5954 .put(&mut approval_question(&state.id))
5955 .unwrap();
5956
5957 let mut t = task();
5958 t.runs.push(state.id.clone());
5959 assert_eq!(
5960 land_resume_state(&t),
5961 LandResume::StillWaiting,
5962 "nobody has answered and the timeout has not passed"
5963 );
5964 }
5965
5966 #[test]
5967 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5968 crate::run::pin_test_home();
5973 let mut state = run_state(RunStatus::Landing);
5974 state.id = "20260101-000000-exp1".to_owned();
5975 state.parked = true;
5976 state.config.graph.answer_timeout = 60;
5977 state.save().unwrap();
5978
5979 let store = ask::Questions::open();
5980 let mut q = approval_question(&state.id);
5981 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
5982 store.put(&mut q).unwrap();
5983
5984 let mut t = task();
5985 t.runs.push(state.id.clone());
5986 assert_eq!(
5987 land_resume_state(&t),
5988 LandResume::Ready,
5989 "an expired question must not be waited on forever"
5990 );
5991
5992 let after = store.get(&q.id).unwrap();
5993 assert!(
5994 !after.status.open(),
5995 "the question is abandoned, not silently ignored"
5996 );
5997 assert!(
5998 after.resolution().is_none(),
5999 "an abandoned question is not read as a decision"
6000 );
6001 }
6002
6003 #[test]
6004 fn reclaim_settles_a_running_task_against_its_last_run() {
6005 let mut t = task();
6006 t.start("20260904-000000-4043".to_owned());
6007 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
6008 assert_eq!(
6009 t.status,
6010 TaskStatus::Done,
6011 "a run that actually finished must not stay `running` forever"
6012 );
6013 }
6014
6015 #[test]
6016 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
6017 let mut t = task();
6021 t.start("20260904-000000-4043".to_owned());
6022 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
6023 assert_eq!(t.status, TaskStatus::Failed);
6024 assert!(t.status.runnable());
6025 }
6026
6027 #[test]
6028 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
6029 let mut t = task();
6030 t.start("20260904-000000-4043".to_owned());
6031 reclaim(&mut t, None, 2, "en");
6032 assert_eq!(t.status, TaskStatus::Held);
6033 assert!(
6034 t.last_error
6035 .as_deref()
6036 .is_some_and(|e| e.contains("running")),
6037 "the operator needs to know why this task was held"
6038 );
6039 }
6040
6041 #[test]
6042 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
6043 let dir = tempfile::tempdir().unwrap();
6044 let queue = Queue::at(dir.path().to_path_buf());
6045
6046 let mut orphaned = task();
6048 orphaned.id = "20260904-000000-orph".to_owned();
6049 orphaned.status = TaskStatus::Running;
6050 orphaned.attempts = 1;
6051 queue.put(&mut orphaned).unwrap();
6052
6053 let mut alive = task();
6054 alive.id = "20260904-000000-live".to_owned();
6055 alive.status = TaskStatus::Running;
6056 alive.attempts = 1;
6057 queue.put(&mut alive).unwrap();
6058 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
6059
6060 let mut queued = task();
6061 queued.id = "20260904-000000-wait".to_owned();
6062 queue.put(&mut queued).unwrap();
6063
6064 let reclaimed = reclaim_orphaned_running(&queue, 2);
6065 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
6066
6067 assert_eq!(
6068 queue.get(&orphaned.id).unwrap().status,
6069 TaskStatus::Held,
6070 "nothing was driving it and there was no run to recover"
6071 );
6072 assert_eq!(
6073 queue.get(&alive.id).unwrap().status,
6074 TaskStatus::Running,
6075 "a live claim must protect the task it belongs to"
6076 );
6077 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
6078 }
6079
6080 fn read_run_under(home: &Path, id: &str) -> RunState {
6086 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
6087 serde_json::from_str(&body).unwrap()
6088 }
6089
6090 #[test]
6091 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
6092 let dir = tempfile::tempdir().unwrap();
6093 let home = dir.path().to_path_buf();
6094 let now = Timestamp::now();
6095 let overrun_seat = || crate::run::ActiveSeat {
6096 node: "implement".to_owned(),
6097 started_at: now - jiff::SignedDuration::new(21_000, 0),
6098 timeout_secs: 3_600,
6099 attempt: 0,
6100 task: None,
6101 command: None,
6102 index: None,
6103 total: None,
6104 };
6105
6106 let mut dead = run_state(RunStatus::Implementing);
6107 dead.id = "20260101-000000-dead".to_owned();
6108 dead.active.insert("impl-A".to_owned(), overrun_seat());
6109 dead.driver_pid = Some(4242);
6112 dead.save_under(&home).unwrap();
6113
6114 let mut alive = run_state(RunStatus::Implementing);
6117 alive.id = "20260101-000000-aliv".to_owned();
6118 alive.active.insert("impl-A".to_owned(), overrun_seat());
6119 alive.save_under(&home).unwrap();
6120 let mut status = Status::new();
6121 status.current = vec![Current {
6122 task: "20260101-000000-task".to_owned(),
6123 run: alive.id.clone(),
6124 }];
6125 write_status_to(&home.join("daemon.json"), &status).unwrap();
6126
6127 let questions = Questions::at(home.join("questions"));
6131 let mut q = ask::Question::new(
6132 dead.id.clone(),
6133 "implement".to_owned(),
6134 "impl-A".to_owned(),
6135 "Which storage backend?".to_owned(),
6136 String::new(),
6137 vec!["SQLite".to_owned(), "Redis".to_owned()],
6138 );
6139 questions.put(&mut q).unwrap();
6140
6141 let abandoned = reclaim_abandoned_runs_with(
6142 &home,
6143 now,
6144 |pid| if pid == 4242 { Some(false) } else { None },
6145 |_| panic!("a query answering Dead outright needs no identity corroboration"),
6146 );
6147 assert_eq!(abandoned, vec![dead.id.clone()]);
6148
6149 let reloaded = read_run_under(&home, &dead.id);
6150 assert_eq!(reloaded.status, RunStatus::Failed);
6151 assert!(reloaded.active.is_empty());
6152 assert!(
6153 !questions.get(&q.id).unwrap().status.open(),
6154 "the failed run's own open question must be settled in the same pass"
6155 );
6156
6157 let still_alive = read_run_under(&home, &alive.id);
6158 assert_eq!(
6159 still_alive.status,
6160 RunStatus::Implementing,
6161 "a live daemon's claim protects it"
6162 );
6163 assert!(!still_alive.active.is_empty());
6164 }
6165
6166 #[test]
6176 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
6177 let dir = tempfile::tempdir().unwrap();
6178 let home = dir.path().to_path_buf();
6179 let now = Timestamp::now();
6180
6181 let mut manual = run_state(RunStatus::Reviewing);
6182 manual.id = "20260101-000000-manl".to_owned();
6183 manual.active.insert(
6184 "review-1".to_owned(),
6185 crate::run::ActiveSeat {
6186 node: "review".to_owned(),
6187 started_at: now - jiff::SignedDuration::new(21_000, 0),
6188 timeout_secs: 3_600,
6189 attempt: 0,
6190 task: None,
6191 command: None,
6192 index: None,
6193 total: None,
6194 },
6195 );
6196 manual.driver_pid = Some(4242);
6200 manual.driver_started_at = Some("1790000000".to_owned());
6201 manual.save_under(&home).unwrap();
6202
6203 let abandoned = reclaim_abandoned_runs_with(
6204 &home,
6205 now,
6206 |pid| if pid == 4242 { Some(true) } else { None },
6207 |pid| {
6208 if pid == 4242 {
6209 Some("1790000000".to_owned())
6210 } else {
6211 None
6212 }
6213 },
6214 );
6215 assert!(
6216 abandoned.is_empty(),
6217 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
6218 );
6219
6220 let reloaded = read_run_under(&home, &manual.id);
6221 assert_eq!(reloaded.status, RunStatus::Reviewing);
6222 assert!(!reloaded.active.is_empty());
6223 }
6224
6225 #[test]
6226 fn an_already_claimed_task_is_skipped_rather_than_failed() {
6227 let dir = tempfile::tempdir().unwrap();
6228 let queue = Queue::at(dir.path().to_path_buf());
6229 let mut only = task();
6230 queue.put(&mut only).unwrap();
6231
6232 let _elsewhere = queue.claim(&only.id).unwrap();
6233 let candidates = runnable(&queue);
6234 assert_eq!(candidates.len(), 1, "the task is still runnable");
6235 assert!(
6236 queue.claim(&candidates[0].id).is_err(),
6237 "the loop cannot take a claim somebody else holds"
6238 );
6239
6240 let after = queue.get(&only.id).unwrap();
6241 assert_eq!(after.status, TaskStatus::Queued);
6242 assert_eq!(
6243 after.attempts, 0,
6244 "losing the race is not an attempt at the task"
6245 );
6246 assert_eq!(after.last_error, None);
6247 }
6248
6249 #[test]
6250 fn the_status_file_round_trips_and_its_heartbeat_advances() {
6251 let dir = tempfile::tempdir().unwrap();
6252 let path = dir.path().join("daemon.json");
6253
6254 let mut status = Status::new();
6255 status.idle = false;
6256 status.completed = 7;
6257 status.current = vec![Current {
6258 task: "20260902-000000-t111".to_owned(),
6259 run: "20260902-000001-r111".to_owned(),
6260 }];
6261 write_status_to(&path, &status).unwrap();
6262 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6263 assert_eq!(first.schema, SCHEMA);
6264 assert_eq!(first.pid, std::process::id());
6265 assert!(!first.idle);
6266 assert_eq!(first.completed, 7);
6267 assert_eq!(first.current, status.current);
6268 assert!(
6269 !path.with_extension("json.tmp").exists(),
6270 "the temp file is renamed, not left behind"
6271 );
6272
6273 std::thread::sleep(Duration::from_millis(5));
6274 status.updated_at = Timestamp::now();
6275 status.polls = 3;
6276 write_status_to(&path, &status).unwrap();
6277 let second: Status =
6278 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6279 assert!(
6280 second.updated_at > first.updated_at,
6281 "a reader can only detect staleness if the heartbeat moves"
6282 );
6283 assert_eq!(
6284 second.started_at, first.started_at,
6285 "the start time is not a heartbeat"
6286 );
6287 assert_eq!(second.polls, 3);
6288 }
6289
6290 #[test]
6291 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
6292 let dir = tempfile::tempdir().unwrap();
6293
6294 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
6295
6296 let mut status = Status::new();
6297 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
6298 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6299 let stale = read_status(dir.path()).unwrap();
6300 assert!(
6301 !stale.running(Timestamp::now()),
6302 "a minute without a heartbeat is a dead daemon, not a busy one"
6303 );
6304 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
6305
6306 status.updated_at = Timestamp::now();
6307 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6308 let fresh = read_status(dir.path()).unwrap();
6309 assert!(fresh.running(Timestamp::now()));
6310 }
6311
6312 #[test]
6313 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
6314 let dir = tempfile::tempdir().unwrap();
6315 let now = Timestamp::now();
6316 let mine = "20260903-080619-01c2";
6317
6318 assert!(
6319 !is_working_on(dir.path(), mine, now),
6320 "no status file means nobody is working on anything"
6321 );
6322
6323 let mut status = Status::new();
6324 status.current = vec![Current {
6325 task: "20260903-080340-0167".to_owned(),
6326 run: mine.to_owned(),
6327 }];
6328 status.updated_at = now;
6329 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6330 assert!(is_working_on(dir.path(), mine, now));
6331 assert!(
6332 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
6333 "a daemon busy with one run is not working on another"
6334 );
6335
6336 status.updated_at = now - jiff::SignedDuration::from_secs(600);
6339 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6340 assert!(
6341 !is_working_on(dir.path(), mine, now),
6342 "a stale heartbeat is a dead daemon, so its run is a leftover"
6343 );
6344 }
6345
6346 #[test]
6347 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
6348 let dir = tempfile::tempdir().unwrap();
6349 let now = Timestamp::now();
6350
6351 assert!(
6352 !is_working_on_short(dir.path(), "01c2", now),
6353 "no status file means nobody is working on anything"
6354 );
6355
6356 let mut status = Status::new();
6357 status.current = vec![Current {
6358 task: "20260903-080340-0167".to_owned(),
6359 run: "20260903-080619-01c2".to_owned(),
6360 }];
6361 status.updated_at = now;
6362 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6363 assert!(
6364 is_working_on_short(dir.path(), "01c2", now),
6365 "the run's short id is the last block of its full id"
6366 );
6367 assert!(
6368 !is_working_on_short(dir.path(), "3cbf", now),
6369 "a daemon busy with one worktree bay is not working on another"
6370 );
6371 }
6372
6373 #[test]
6374 fn a_newer_status_file_still_yields_a_reading() {
6375 let dir = tempfile::tempdir().unwrap();
6376 std::fs::write(
6379 dir.path().join("daemon.json"),
6380 serde_json::json!({
6381 "schema": 2,
6382 "updated_at": Timestamp::now().to_string(),
6383 "idle": true,
6384 "surprise": { "nested": [1, 2, 3] },
6385 })
6386 .to_string(),
6387 )
6388 .unwrap();
6389
6390 let reading = read_status(dir.path()).expect("a forward-compatible read");
6391 assert!(reading.running(Timestamp::now()));
6392 assert!(reading.idle);
6393 assert!(reading.current.is_empty());
6394 }
6395
6396 #[test]
6397 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
6398 let dir = tempfile::tempdir().unwrap();
6404 std::fs::write(
6405 dir.path().join("daemon.json"),
6406 serde_json::json!({
6407 "schema": 1,
6408 "pid": 4242,
6409 "updated_at": Timestamp::now().to_string(),
6410 "idle": false,
6411 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
6412 "completed": 3,
6413 "polls": 9,
6414 })
6415 .to_string(),
6416 )
6417 .unwrap();
6418
6419 let reading = read_status(dir.path()).expect("an older shape must still parse");
6420 assert!(reading.running(Timestamp::now()));
6421 assert_eq!(
6422 reading.current,
6423 vec![Current {
6424 task: "20260902-140501-aaaa".to_owned(),
6425 run: "20260902-140502-bbbb".to_owned(),
6426 }]
6427 );
6428 }
6429
6430 #[test]
6431 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
6432 let dir = tempfile::tempdir().unwrap();
6433 std::fs::write(
6434 dir.path().join("daemon.json"),
6435 serde_json::json!({
6436 "schema": 1,
6437 "updated_at": Timestamp::now().to_string(),
6438 "idle": true,
6439 "current": null,
6440 })
6441 .to_string(),
6442 )
6443 .unwrap();
6444 let with_null = read_status(dir.path()).expect("null must still parse");
6445 assert!(with_null.current.is_empty());
6446
6447 std::fs::write(
6448 dir.path().join("daemon.json"),
6449 serde_json::json!({
6450 "schema": 1,
6451 "updated_at": Timestamp::now().to_string(),
6452 "idle": true,
6453 })
6454 .to_string(),
6455 )
6456 .unwrap();
6457 let absent = read_status(dir.path()).expect("a missing field must still parse");
6458 assert!(absent.current.is_empty());
6459 }
6460
6461 #[test]
6462 fn a_task_without_a_repository_runs_in_the_daemons_default() {
6463 let fallback = Path::new("/default");
6464 let mut blank = task();
6465 blank.repo = PathBuf::new();
6466 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
6467 let mut dot = task();
6468 dot.repo = PathBuf::from(".");
6469 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
6470 assert_eq!(
6471 repo_for(&task(), fallback),
6472 PathBuf::from("/repo"),
6473 "a task that names a repository keeps it"
6474 );
6475 }
6476
6477 #[test]
6478 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
6479 let mut solo_cfg = Config::default();
6485 solo_cfg.graph.candidates = 3;
6486 let mut solo_task = task();
6487 solo_task.solo = true;
6488 apply_solo(&mut solo_cfg, &solo_task);
6489 assert_eq!(solo_cfg.graph.candidates, 1);
6490
6491 let mut plain_cfg = Config::default();
6492 plain_cfg.graph.candidates = 3;
6493 let plain_task = task();
6494 assert!(!plain_task.solo);
6495 apply_solo(&mut plain_cfg, &plain_task);
6496 assert_eq!(
6497 plain_cfg.graph.candidates, 3,
6498 "a task that did not ask to run alone keeps the config's candidates"
6499 );
6500 }
6501
6502 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
6503 QuotaLoss {
6504 seat: seat.into(),
6505 node: "judge".into(),
6506 at: at.parse().unwrap(),
6507 reset: reset.map(str::to_string),
6508 }
6509 }
6510
6511 #[test]
6512 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
6513 let old: Vec<QuotaLoss> = (1..=4)
6514 .map(|i| {
6515 loss(
6516 &format!("judge-{i}"),
6517 "2026-09-23T05:23:00Z",
6518 Some("2:40pm (Asia/Tokyo)"),
6519 )
6520 })
6521 .collect();
6522 let fresh = losses_this_attempt(&old, &old);
6523 assert!(fresh.is_empty());
6524 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
6525 }
6527
6528 #[test]
6529 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6530 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6531 let now = Timestamp::now();
6532 let mut after = old.clone();
6533 after.push(loss("judge-2", &now.to_string(), None));
6534 let fresh = losses_this_attempt(&old, &after);
6535 assert_eq!(fresh, vec![after[1].clone()]);
6536 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6537 assert_eq!(
6538 until,
6539 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6540 );
6541 }
6542
6543 #[test]
6544 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6545 let before = vec![
6548 loss("judge-1", "2026-09-23T05:23:00Z", None),
6549 loss("judge-2", "2026-09-23T05:24:00Z", None),
6550 ];
6551 let after = vec![
6552 loss("judge-2", "2026-09-23T05:24:00Z", None),
6553 loss("judge-1", "2026-09-24T01:00:00Z", None),
6554 ];
6555 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6556 }
6557
6558 #[test]
6559 fn merge_overrides_are_parsed_or_refused() {
6560 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6561 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6562 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6563 assert!(merge_mode("squash").is_err());
6564 }
6565
6566 #[test]
6567 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6568 let now = Timestamp::now();
6569 let fallback = Duration::from_secs(300);
6570 let cap = Duration::from_secs(1800);
6571
6572 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6574
6575 let soon = now + jiff::SignedDuration::from_secs(600);
6577 assert_eq!(
6578 quota_wait(Some(soon), now, fallback, cap),
6579 Duration::from_secs(600)
6580 );
6581
6582 let past = now - jiff::SignedDuration::from_secs(60);
6585 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6586
6587 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6590 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6591 }
6592
6593 #[test]
6594 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6595 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6596
6597 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6598 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6599
6600 let already_past =
6604 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6605 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6606
6607 assert!(
6608 parse_reset_hint("session limit reached", now, now).is_none(),
6609 "free text with no recognised shape is not guessed at"
6610 );
6611 assert!(
6612 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6613 "an unresolvable zone name is not guessed at either"
6614 );
6615 }
6616
6617 #[test]
6618 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6619 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6620
6621 let at = parse_reset_hint(
6622 "You've hit your usage limit. Visit \
6623 https://chatgpt.com/codex/settings/usage to purchase more \
6624 credits or try again at Sep 19th, 2026 5:10 PM.",
6625 now,
6626 now,
6627 )
6628 .expect("the codex reset wording is a recognised shape");
6629 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6630
6631 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6636 .expect("an explicit year needs no rollover");
6637 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6638
6639 assert!(
6640 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6641 "a two-digit year is not the documented shape and is not guessed at"
6642 );
6643 assert!(
6644 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6645 "a four-letter month name is not the documented three-letter abbreviation"
6646 );
6647 assert!(
6648 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6649 "an explicit zone on the dated shape is a format nobody has \
6650 documented, and is refused rather than guessed at as UTC"
6651 );
6652 }
6653
6654 #[test]
6655 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6656 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6657 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6658
6659 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6660 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6661
6662 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6663 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6664
6665 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6666 assert!(
6667 parse_reset_hint(bad, now, recorded).is_none(),
6668 "{bad:?} must not be guessed at"
6669 );
6670 }
6671 }
6672
6673 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6677 let config = dir.join("magi.toml");
6678 std::fs::write(
6679 &config,
6680 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6681 )
6682 .unwrap();
6683 let opts = Opts {
6684 poll: Duration::from_secs(30),
6685 config: Some(config),
6686 repo: dir.join("repo"),
6690 ..Opts::default()
6691 };
6692 let home = dir.join("home");
6701 let worktrees = dir.join("wt");
6702 (
6703 opts,
6704 Queue::at(dir.join("queue")),
6705 home.join("daemon.json"),
6706 home,
6707 worktrees,
6708 )
6709 }
6710
6711 #[test]
6712 fn a_stop_is_idempotent_and_once_set_stays_set() {
6713 let stop = Stop::new();
6714 assert!(!stop.stopped());
6715
6716 stop.stop();
6717 assert!(stop.stopped());
6718 stop.stop();
6719 assert!(stop.stopped(), "a second stop is not a toggle");
6720
6721 let shared = stop.clone();
6722 assert!(
6723 shared.stopped(),
6724 "a clone is the same stop; that is how the loop and its caller share one"
6725 );
6726 }
6727
6728 #[test]
6729 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6730 let stop = Stop::new();
6731 stop.enter();
6732 assert!(
6733 !stop.finishing(),
6734 "a busy loop nobody has asked to stop is just running"
6735 );
6736
6737 stop.stop();
6738 assert!(
6739 stop.finishing(),
6740 "a stop asked for mid-run has not landed until the run is settled"
6741 );
6742
6743 stop.exit();
6744 assert!(
6745 !stop.finishing(),
6746 "once the run is settled the stop has landed and there is nothing to finish"
6747 );
6748 }
6749
6750 #[test]
6751 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6752 let stop = Stop::new();
6753 stop.enter();
6754 stop.enter();
6755 stop.stop();
6756 assert!(stop.finishing(), "two runs still in flight");
6757
6758 stop.exit();
6759 assert!(
6760 stop.finishing(),
6761 "one run finished, but a sibling is still working"
6762 );
6763
6764 stop.exit();
6765 assert!(
6766 !stop.finishing(),
6767 "the last run out is what actually lands the stop"
6768 );
6769 }
6770
6771 #[tokio::test]
6772 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6773 let dir = tempfile::tempdir().unwrap();
6774 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6775 let stop = Stop::new();
6776 stop.stop();
6777
6778 let began = std::time::Instant::now();
6779 tokio::time::timeout(
6780 Duration::from_secs(2),
6781 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6782 )
6783 .await
6784 .expect("a stopped loop must return, not sit out its poll interval")
6785 .expect("the loop's own setup and teardown must not fail");
6786 assert!(
6787 began.elapsed() < opts.poll,
6788 "returned only after {:?}, which is a poll interval, not a stop",
6789 began.elapsed()
6790 );
6791 }
6792
6793 #[tokio::test]
6794 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6795 let dir = tempfile::tempdir().unwrap();
6796 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6797 let stop = Stop::new();
6798
6799 let asker = {
6802 let stop = stop.clone();
6803 tokio::spawn(async move {
6804 tokio::time::sleep(Duration::from_millis(20)).await;
6805 stop.stop();
6806 })
6807 };
6808
6809 let began = std::time::Instant::now();
6810 tokio::time::timeout(
6811 Duration::from_secs(2),
6812 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6813 )
6814 .await
6815 .expect("a stop asked for while idle must wake the wait")
6816 .expect("the loop's own setup and teardown must not fail");
6817 asker.await.unwrap();
6818 assert!(
6819 began.elapsed() < opts.poll,
6820 "returned only after {:?}, so the stop waited on the sleep",
6821 began.elapsed()
6822 );
6823 }
6824
6825 #[tokio::test]
6826 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6827 let dir = tempfile::tempdir().unwrap();
6828 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6829 let stop = Stop::new();
6830 stop.stop();
6831
6832 tokio::time::timeout(
6833 Duration::from_secs(2),
6834 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6835 )
6836 .await
6837 .expect("a stopped loop must return")
6838 .expect("the loop's own setup and teardown must not fail");
6839
6840 assert!(
6841 home.is_dir(),
6842 "the loop did publish a status file, so its removal is the teardown and not an absence"
6843 );
6844 assert!(
6845 !status_file.exists(),
6846 "a stopped loop clears its status file"
6847 );
6848 assert!(
6849 read_status(&home).is_none(),
6850 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6851 );
6852 }
6853
6854 #[tokio::test]
6855 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6856 let dir = tempfile::tempdir().unwrap();
6857 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6858 opts.once = true;
6859
6860 let mut settled = RunState::new(
6861 dir.path().join("repo"),
6862 "main".to_owned(),
6863 "abc1234".to_owned(),
6864 "fixture".to_owned(),
6865 Config::default(),
6866 );
6867 settled.status = RunStatus::Ready;
6868 let run_dir = home.join("runs").join(&settled.id);
6869 std::fs::create_dir_all(&run_dir).unwrap();
6870 std::fs::write(
6871 run_dir.join("run.json"),
6872 serde_json::to_string_pretty(&settled).unwrap(),
6873 )
6874 .unwrap();
6875 let questions = Questions::at(home.join("questions"));
6876 let mut question = ask::Question::new(
6877 settled.id.clone(),
6878 "review".to_owned(),
6879 "reviewer-1".to_owned(),
6880 "Continue?".to_owned(),
6881 String::new(),
6882 Vec::new(),
6883 );
6884 questions.put(&mut question).unwrap();
6885
6886 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6887 .await
6888 .unwrap();
6889
6890 assert_eq!(
6891 questions.get(&question.id).unwrap().status,
6892 ask::QuestionStatus::Abandoned,
6893 "an empty --once drain still performs startup question cleanup"
6894 );
6895 }
6896
6897 #[test]
6898 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6899 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6900
6901 assert!(
6902 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6903 "never checked before: due at once"
6904 );
6905
6906 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6907 assert!(
6908 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6909 "well inside the interval: not due yet"
6910 );
6911
6912 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6913 assert!(
6914 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6915 "exactly at the edge: not yet due, same convention as `clean::due`"
6916 );
6917
6918 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6919 assert!(
6920 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6921 "past the interval: due again"
6922 );
6923 }
6924
6925 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
6930 let config = dir.join("magi.toml");
6931 std::fs::write(
6937 &config,
6938 format!(
6939 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
6940 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
6941 cache_dir.display()
6942 ),
6943 )
6944 .unwrap();
6945 Opts {
6946 config: Some(config),
6947 repo: dir.join("repo"),
6948 ..Opts::default()
6949 }
6950 }
6951
6952 #[tokio::test]
6953 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
6954 let dir = tempfile::tempdir().unwrap();
6955 let home = dir.path().join("home");
6956 let cache_dir = dir.path().join("cache");
6957 std::fs::create_dir_all(&cache_dir).unwrap();
6958 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6959 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6960
6961 let running = Stop::new();
6964 let mut last_checked = None;
6965 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6966 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
6967 .await;
6968 assert_eq!(
6969 crate::disk::dir_size(&cache_dir),
6970 0,
6971 "over the cap on the first check ever: pruned at once, no idle queue required"
6972 );
6973 assert_eq!(last_checked, Some(t0));
6974
6975 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
6977 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
6978 maybe_prune_cache_between_runs(
6979 &opts.repo,
6980 &opts,
6981 &home,
6982 &running,
6983 &mut last_checked,
6984 too_soon,
6985 )
6986 .await;
6987 assert_eq!(
6988 crate::disk::dir_size(&cache_dir),
6989 10,
6990 "too soon since the last check: left alone rather than rescanned every call"
6991 );
6992 assert_eq!(
6993 last_checked,
6994 Some(t0),
6995 "an idle check does not reset the clock"
6996 );
6997
6998 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
7000 maybe_prune_cache_between_runs(
7001 &opts.repo,
7002 &opts,
7003 &home,
7004 &running,
7005 &mut last_checked,
7006 due_again,
7007 )
7008 .await;
7009 assert_eq!(
7010 crate::disk::dir_size(&cache_dir),
7011 0,
7012 "due again: pruned back under the cap"
7013 );
7014 }
7015
7016 #[tokio::test]
7024 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
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 stop = Stop::new();
7033 stop.stop();
7034 assert!(
7035 !stop.finishing(),
7036 "no run is in flight at a between-runs boundary, so nothing else \
7037 would tell the operator this stop had not taken effect yet"
7038 );
7039
7040 let mut last_checked = None;
7041 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
7042 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
7043 .await;
7044 assert_eq!(
7045 crate::disk::dir_size(&cache_dir),
7046 10,
7047 "over its cap, and due for the first check ever, but a stop outranks \
7048 it: the cap is a standing policy the next start measures again"
7049 );
7050 assert_eq!(
7051 last_checked, None,
7052 "a check that never happened must not claim the interval"
7053 );
7054 }
7055
7056 #[tokio::test]
7070 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
7071 let dir = tempfile::tempdir().unwrap();
7072 let cache_dir = dir.path().join("cache");
7073 std::fs::create_dir_all(&cache_dir).unwrap();
7074 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
7075
7076 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
7077 opts.poll = Duration::from_millis(20);
7078 opts.max_attempts = 1_000;
7079
7080 let queue = Queue::at(dir.path().join("queue"));
7081 let mut t = Task::new(
7082 "x".to_owned(),
7083 "x".to_owned(),
7084 opts.repo.clone(),
7085 Source::Human,
7086 );
7087 queue.put(&mut t).unwrap();
7088
7089 let home = dir.path().join("home");
7090 let worktrees = dir.path().join("wt");
7091 let status_file = home.join("daemon.json");
7092 let stop = Stop::new();
7093 let stopper = {
7094 let stop = stop.clone();
7095 tokio::spawn(async move {
7096 tokio::time::sleep(Duration::from_millis(400)).await;
7097 stop.stop();
7098 })
7099 };
7100
7101 tokio::time::timeout(
7102 Duration::from_secs(10),
7103 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
7104 )
7105 .await
7106 .expect("the loop must not hang on a queue that keeps producing failing work")
7107 .expect("the loop's own setup and teardown must not fail");
7108 stopper.await.unwrap();
7109
7110 let after = queue.get(&t.id).unwrap();
7111 assert!(
7112 after.attempts >= 2,
7113 "the harness must actually have retried more than once, or this is not \
7114 exercising a busy queue at all (got {} attempt(s))",
7115 after.attempts
7116 );
7117 assert!(
7118 after.status.runnable(),
7119 "still under its attempt budget: the queue never reached a natural idle \
7120 on its own, only the external stop ended the test"
7121 );
7122
7123 assert_eq!(
7124 crate::disk::dir_size(&cache_dir),
7125 0,
7126 "an oversized cache must not be left to grow unboundedly just because the \
7127 queue kept the loop busy the whole time"
7128 );
7129 }
7130
7131 #[test]
7132 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
7133 let dir = tempfile::tempdir().unwrap();
7134 let queue = Queue::at(dir.path().join("queue"));
7135 let questions = Questions::at(dir.path().join("questions"));
7136 let mut task = task();
7137 queue.put(&mut task).unwrap();
7138
7139 let mut task_question = ask::Question::new(
7140 task.id.clone(),
7141 crate::conduct::NODE.to_owned(),
7142 "conduct".to_owned(),
7143 "Which backend?".to_owned(),
7144 String::new(),
7145 Vec::new(),
7146 );
7147 questions.put(&mut task_question).unwrap();
7148 task.block(vec![task_question.id.clone()], None);
7149 queue.put(&mut task).unwrap();
7150
7151 let mut run_question = ask::Question::new(
7152 "20260101-000000-run1".to_owned(),
7153 "review".to_owned(),
7154 "reviewer-1".to_owned(),
7155 "Run question".to_owned(),
7156 String::new(),
7157 Vec::new(),
7158 );
7159 questions.put(&mut run_question).unwrap();
7160
7161 let mut coincidental = ask::Question::new(
7166 task.id.clone(),
7167 "review".to_owned(),
7168 "reviewer-1".to_owned(),
7169 "Unrelated review question".to_owned(),
7170 String::new(),
7171 Vec::new(),
7172 );
7173 questions.put(&mut coincidental).unwrap();
7174
7175 reconcile_task_questions(&queue, &questions);
7176 assert!(questions.get(&task_question.id).unwrap().status.open());
7177 assert!(questions.get(&run_question.id).unwrap().status.open());
7178 assert!(questions.get(&coincidental.id).unwrap().status.open());
7179
7180 task.release();
7181 queue.put(&mut task).unwrap();
7182 reconcile_task_questions(&queue, &questions);
7183 assert_eq!(
7184 questions.get(&task_question.id).unwrap().status,
7185 ask::QuestionStatus::Abandoned
7186 );
7187 assert!(
7188 questions.get(&run_question.id).unwrap().status.open(),
7189 "run questions remain the run janitor's responsibility"
7190 );
7191 assert!(
7192 questions.get(&coincidental.id).unwrap().status.open(),
7193 "a non-conductor question must not be abandoned just because its \
7194 run id coincides with a task id"
7195 );
7196 }
7197
7198 #[test]
7199 fn a_freshly_started_running_task_is_never_stalled() {
7200 let dir = tempfile::tempdir().unwrap();
7201 let mut t = task();
7202 t.start("run-1".to_owned());
7203 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
7206 }
7207
7208 #[test]
7209 fn a_long_running_task_with_no_live_daemon_is_stalled() {
7210 let dir = tempfile::tempdir().unwrap();
7211 let mut t = task();
7212 t.start("run-1".to_owned());
7213 t.updated_at = Timestamp::now()
7214 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
7215 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
7216 assert_eq!(
7217 stalled_tasks(
7218 &Queue::at(dir.path().join("q")),
7219 dir.path(),
7220 Timestamp::now()
7221 )
7222 .len(),
7223 0,
7224 "the task was never written to this queue"
7225 );
7226 }
7227
7228 #[test]
7229 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
7230 let dir = tempfile::tempdir().unwrap();
7231 let mut t = task();
7232 t.id = "20260903-080340-0167".to_owned();
7233 t.start("20260903-080619-01c2".to_owned());
7234 t.updated_at = Timestamp::now()
7235 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
7236
7237 let mut status = Status::new();
7238 status.current = vec![Current {
7239 task: t.id.clone(),
7240 run: "20260903-080619-01c2".to_owned(),
7241 }];
7242 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
7243
7244 assert!(
7245 !is_stalled(&t, dir.path(), Timestamp::now()),
7246 "a live daemon's own heartbeat rules out stalled, however long the task has run"
7247 );
7248 }
7249
7250 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
7254 let path = queue.path_of(id);
7255 let body = std::fs::read_to_string(&path).unwrap();
7256 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
7257 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
7258 v["updated_at"] = serde_json::Value::String(old.to_string());
7259 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
7260 }
7261
7262 #[test]
7263 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
7264 let dir = tempfile::tempdir().unwrap();
7277 let queue = Queue::at(dir.path().join("queue"));
7278 let home = dir.path().join("home");
7279
7280 let mut t = task();
7281 t.id = "20260101-000001-lock".to_owned();
7282 t.start("run-1".to_owned());
7283 queue.put(&mut t).unwrap();
7284 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7285 std::fs::write(
7286 dir.path().join("queue").join(format!("{}.lock", t.id)),
7287 "not a pid",
7288 )
7289 .unwrap();
7290
7291 let now = Timestamp::now();
7292 assert!(
7293 reclaim_orphaned_running(&queue, 2).is_empty(),
7294 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
7295 and reclaim must leave the task alone"
7296 );
7297 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
7298
7299 let stalled = stalled_tasks(&queue, &home, now);
7300 assert_eq!(
7301 stalled.len(),
7302 1,
7303 "reclaim's inability to claim it yet must not hide it from the conductor"
7304 );
7305 assert_eq!(stalled[0].id, t.id);
7306 }
7307
7308 #[test]
7309 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
7310 let dir = tempfile::tempdir().unwrap();
7311 crate::run::set_home(dir.path().join("run-home"));
7312 let queue = Queue::at(dir.path().join("queue"));
7313 let home = dir.path().join("home");
7314 let questions = Questions::at(dir.path().join("questions"));
7315
7316 let mut t = task();
7317 t.id = "20260101-000003-dead".to_owned();
7318 t.start("missing-run".to_owned());
7319 queue.put(&mut t).unwrap();
7320 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7321
7322 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7325 assert_eq!(
7326 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
7327 [&t.id]
7328 );
7329 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
7330 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7331
7332 crate::conduct::apply(
7335 &queue,
7336 &questions,
7337 &crate::conduct::Verdict {
7338 decisions: vec![crate::conduct::Decision {
7339 id: t.id.clone(),
7340 recovery: Some(crate::conduct::Recovery::Requeue),
7341 ..crate::conduct::Decision::default()
7342 }],
7343 },
7344 )
7345 .unwrap();
7346 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7347 }
7348
7349 #[test]
7350 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
7351 let dir = tempfile::tempdir().unwrap();
7352 let queue = Queue::at(dir.path().join("queue"));
7353 let home = dir.path().join("home");
7354
7355 let mut fresh = task();
7356 fresh.id = "20260101-000001-aaaa".to_owned();
7357 fresh.start("run-1".to_owned());
7358 queue.put(&mut fresh).unwrap();
7359
7360 let mut old = task();
7361 old.id = "20260101-000002-bbbb".to_owned();
7362 old.start("run-2".to_owned());
7363 queue.put(&mut old).unwrap();
7364 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
7365
7366 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7367 assert_eq!(stalled.len(), 1);
7368 assert_eq!(stalled[0].id, old.id);
7369 }
7370
7371 #[test]
7372 fn queued_and_finished_task_views_partition_by_status() {
7373 let dir = tempfile::tempdir().unwrap();
7374 let queue = Queue::at(dir.path().join("queue"));
7375
7376 let mut queued = task();
7377 queued.id = "20260101-000001-aaaa".to_owned();
7378 queue.put(&mut queued).unwrap();
7379
7380 let mut failed = task();
7381 failed.id = "20260101-000002-bbbb".to_owned();
7382 failed.start("run-1".to_owned());
7383 failed.fail("gate red", 5);
7384 queue.put(&mut failed).unwrap();
7385
7386 let mut held = task();
7387 held.id = "20260101-000003-cccc".to_owned();
7388 held.hold_machine(None);
7389 queue.put(&mut held).unwrap();
7390
7391 let mut running = task();
7392 running.id = "20260101-000004-dddd".to_owned();
7393 running.start("run-2".to_owned());
7394 queue.put(&mut running).unwrap();
7395
7396 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
7397 assert_eq!(queued_ids, [queued.id.clone()]);
7398
7399 let mut finished_ids: Vec<String> =
7400 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
7401 finished_ids.sort_unstable();
7402 let mut want = vec![failed.id.clone(), held.id.clone()];
7403 want.sort_unstable();
7404 assert_eq!(finished_ids, want);
7405 }
7406
7407 #[test]
7408 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
7409 let dir = tempfile::tempdir().unwrap();
7410 let queue = Queue::at(dir.path().join("queue"));
7411 let questions = ask::Questions::at(dir.path().join("questions"));
7412
7413 let mut dep = task();
7414 dep.id = "20260101-000001-dep0".to_owned();
7415 dep.succeed();
7416 queue.put(&mut dep).unwrap();
7417
7418 let mut still_going = task();
7419 still_going.id = "20260101-000002-dep1".to_owned();
7420 queue.put(&mut still_going).unwrap();
7421
7422 let mut blocked = task();
7423 blocked.id = "20260101-000003-main".to_owned();
7424 blocked.block(
7425 vec![dep.id.clone(), still_going.id.clone()],
7426 Some("waits on both".to_owned()),
7427 );
7428 queue.put(&mut blocked).unwrap();
7429
7430 resolve_blockers(&queue, &questions);
7431
7432 let after = queue.get(&blocked.id).unwrap();
7433 assert_eq!(
7434 after.status,
7435 TaskStatus::Blocked,
7436 "one dependency is still outstanding"
7437 );
7438 assert_eq!(after.blocked_by, [still_going.id.clone()]);
7439 }
7440
7441 #[test]
7442 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
7443 let dir = tempfile::tempdir().unwrap();
7444 let queue = Queue::at(dir.path().join("queue"));
7445 let questions = ask::Questions::at(dir.path().join("questions"));
7446
7447 let mut q = crate::ask::Question::new(
7448 "20260101-000001-main".to_owned(),
7449 crate::conduct::NODE.to_owned(),
7450 "conduct".to_owned(),
7451 "Which backend?".to_owned(),
7452 String::new(),
7453 Vec::new(),
7454 );
7455 questions.put(&mut q).unwrap();
7456 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
7457 .unwrap();
7458 questions.put(&mut q).unwrap();
7459
7460 let mut blocked = task();
7461 blocked.id = "20260101-000001-main".to_owned();
7462 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
7463 queue.put(&mut blocked).unwrap();
7464
7465 resolve_blockers(&queue, &questions);
7466
7467 let after = queue.get(&blocked.id).unwrap();
7468 assert_eq!(
7469 after.status,
7470 TaskStatus::Queued,
7471 "the only blocker resolved"
7472 );
7473 assert_eq!(after.answers.len(), 1);
7474 assert_eq!(after.answers[0].question, "Which backend?");
7475 assert_eq!(after.answers[0].answer, "SQLite");
7476
7477 let instruction = instruction_for(&after);
7479 assert!(instruction.contains("Which backend?"));
7480 assert!(instruction.contains("SQLite"));
7481 }
7482
7483 #[test]
7484 fn resolve_blockers_holds_a_task_whose_conductor_question_was_abandoned() {
7485 let dir = tempfile::tempdir().unwrap();
7486 let queue = Queue::at(dir.path().join("queue"));
7487 let questions = ask::Questions::at(dir.path().join("questions"));
7488
7489 let mut q = crate::ask::Question::new(
7490 "20260101-000001-main".to_owned(),
7491 crate::conduct::NODE.to_owned(),
7492 "conduct".to_owned(),
7493 "Is the setup done?".to_owned(),
7494 String::new(),
7495 Vec::new(),
7496 );
7497 q.abandon("no answer within 60s of asking");
7498 questions.put(&mut q).unwrap();
7499
7500 let mut blocked = task();
7501 blocked.id = "20260101-000001-main".to_owned();
7502 blocked.block(vec![q.id.clone()], Some("setup?".to_owned()));
7503 queue.put(&mut blocked).unwrap();
7504
7505 resolve_blockers(&queue, &questions);
7506
7507 let after = queue.get(&blocked.id).unwrap();
7508 assert_eq!(
7509 after.status,
7510 TaskStatus::Held,
7511 "never left blocked on nothing"
7512 );
7513 assert!(!after.operator_held(), "a machine hold, for triage");
7514 assert!(
7515 after
7516 .hold_reason
7517 .as_deref()
7518 .unwrap_or_default()
7519 .contains("went unanswered")
7520 );
7521 }
7522
7523 #[test]
7524 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
7525 let dir = tempfile::tempdir().unwrap();
7531 let queue = Queue::at(dir.path().join("queue"));
7532 let questions = ask::Questions::at(dir.path().join("questions"));
7533
7534 let mut q = crate::ask::Question::new(
7535 "20260101-000001-main".to_owned(),
7536 crate::conduct::NODE.to_owned(),
7537 "conduct".to_owned(),
7538 "How should this be handled?".to_owned(),
7539 String::new(),
7540 Vec::new(),
7541 );
7542 questions.put(&mut q).unwrap();
7543 q.answer(crate::ask::Answer::Text(
7544 "leave it held, a human will look at it later".to_owned(),
7545 ))
7546 .unwrap();
7547 questions.put(&mut q).unwrap();
7548
7549 let mut held = task();
7550 held.id = "20260101-000001-main".to_owned();
7551 held.hold_machine(Some("out of attempts".to_owned()));
7552 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
7553 queue.put(&mut held).unwrap();
7554
7555 resolve_blockers(&queue, &questions);
7556
7557 let after = queue.get(&held.id).unwrap();
7558 assert_eq!(after.status, TaskStatus::Held);
7559 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
7560 assert_eq!(
7561 after.answers[0].answer,
7562 "leave it held, a human will look at it later"
7563 );
7564 }
7565
7566 #[test]
7567 fn resolve_blockers_releases_a_task_whose_dependency_was_deleted_on_purpose() {
7568 let dir = tempfile::tempdir().unwrap();
7569 let queue = Queue::at(dir.path().join("queue"));
7570 let questions = ask::Questions::at(dir.path().join("questions"));
7571 let mut dep = task();
7572 dep.id = "20260101-000001-gone".to_owned();
7573 queue.put(&mut dep).unwrap();
7574 let mut blocked = task();
7575 blocked.id = "20260101-000003-main".to_owned();
7576 blocked.block(vec![dep.id.clone()], None);
7577 queue.put(&mut blocked).unwrap();
7578
7579 let claim = queue.claim(&blocked.id).unwrap();
7581 queue.remove(&dep.id, false, &questions).unwrap();
7582 drop(claim);
7583 resolve_blockers(&queue, &questions);
7584
7585 let after = queue.get(&blocked.id).unwrap();
7586 assert_eq!(after.status, TaskStatus::Queued);
7587 assert!(after.blocked_by.is_empty());
7588 }
7589
7590 #[test]
7591 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
7592 let dir = tempfile::tempdir().unwrap();
7598 let queue = Queue::at(dir.path().join("queue"));
7599 let questions = ask::Questions::at(dir.path().join("questions"));
7600
7601 let mut still_going = task();
7602 still_going.id = "20260101-000002-dep1".to_owned();
7603 queue.put(&mut still_going).unwrap();
7604
7605 let mut blocked = task();
7606 blocked.id = "20260101-000003-main".to_owned();
7607 blocked.block(
7608 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7609 Some("waits on both".to_owned()),
7610 );
7611 queue.put(&mut blocked).unwrap();
7612
7613 resolve_blockers(&queue, &questions);
7614
7615 let after = queue.get(&blocked.id).unwrap();
7616 assert_eq!(
7617 after.status,
7618 TaskStatus::Held,
7619 "a missing dependency must not leave the task blocked forever"
7620 );
7621 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7622 assert!(after.blocked_by.is_empty());
7623 let reason = after.hold_reason.as_deref().unwrap_or_default();
7624 assert!(
7625 reason.contains("20260101-000001-gone"),
7626 "the missing id must be named so an operator can tell what happened: {reason}"
7627 );
7628 assert!(
7629 reason.contains(&still_going.id),
7630 "the still-valid dependency must not silently vanish from the record: {reason}"
7631 );
7632 }
7633
7634 #[test]
7635 fn instruction_for_is_unchanged_without_any_answers() {
7636 let t = task();
7637 assert_eq!(instruction_for(&t), t.instruction);
7638 }
7639
7640 #[test]
7641 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7642 let dir = tempfile::tempdir().unwrap();
7643 let q = Queue::at(dir.path().join("queue"));
7644 let src = dir.path().join("shot.png");
7645 std::fs::write(&src, "x").unwrap();
7646 let mut t = task();
7647 q.attach(&mut t, &[src]).unwrap();
7648 let paths = task_attachments(&q, &t).unwrap();
7649 assert_eq!(paths.len(), 1);
7650 assert!(paths[0].is_absolute() && paths[0].is_file());
7651 std::fs::remove_file(&paths[0]).unwrap();
7652 let err = task_attachments(&q, &t).unwrap_err().to_string();
7653 assert!(err.contains("shot.png"), "{err}");
7654 }
7655
7656 #[test]
7657 fn resumed_instruction_is_unchanged_without_any_answers() {
7658 let t = task();
7659 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7660 }
7661
7662 #[test]
7663 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7664 let mut t = task();
7665 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7666 let old = t.instruction.clone();
7670
7671 let refreshed = resumed_instruction(&old, &t);
7672 assert!(refreshed.starts_with(&old), "the original text is kept");
7673 assert!(refreshed.contains("Which backend?"));
7674 assert!(refreshed.contains("SQLite"));
7675 }
7676
7677 #[test]
7678 fn resumed_instruction_keeps_an_original_answers_heading() {
7679 let mut t = task();
7680 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7681 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7682
7683 let refreshed = resumed_instruction(&t.instruction, &t);
7684
7685 assert!(
7686 refreshed.starts_with(&t.instruction),
7687 "an answers heading in the original instruction is not the appended block"
7688 );
7689 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7690 assert!(refreshed.contains("Which backend?"));
7691 assert!(refreshed.contains("SQLite"));
7692
7693 let repeated = resumed_instruction(&refreshed, &t);
7694 assert_eq!(
7695 repeated, refreshed,
7696 "only the final appended block is refreshed"
7697 );
7698 }
7699
7700 #[test]
7701 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7702 let mut t = task();
7703 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7704
7705 let once = resumed_instruction(&t.instruction, &t);
7709 let twice = resumed_instruction(&once, &t);
7710 assert_eq!(once, twice);
7711 assert_eq!(once.matches("Which backend?").count(), 1);
7712
7713 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7715 let refreshed = resumed_instruction(&once, &t);
7716 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7717 assert!(refreshed.contains("Which backend?"));
7718 assert!(refreshed.contains("Which cache?"));
7719 }
7720
7721 #[test]
7722 fn prepare_instruction_covers_all_three_starters() {
7723 let mut t = task();
7724 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7725
7726 assert_eq!(
7729 prepare_instruction(&Starter::Start, None, &t),
7730 Some(instruction_for(&t))
7731 );
7732
7733 let old = t.instruction.clone();
7736 assert_eq!(
7737 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7738 Some(resumed_instruction(&old, &t))
7739 );
7740
7741 assert_eq!(
7745 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7746 None
7747 );
7748 }
7749
7750 #[test]
7751 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7752 assert_eq!(
7753 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7754 Starter::Review("magi/eba2/A".to_owned())
7755 );
7756 }
7757
7758 #[test]
7759 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7760 assert_eq!(
7761 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7762 Starter::Start,
7763 "a vanished review branch must not fall back to resuming the old run either"
7764 );
7765 }
7766
7767 #[test]
7768 fn a_refused_handover_retries_as_a_review_of_the_same_branch() {
7769 let mut t = task();
7770 t.start("old-run".to_owned());
7771 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
7772 t.release();
7773 let branch = t.review_branch.take();
7774 assert_eq!(
7775 choose_starter(branch.as_deref(), true, Some("old-run")),
7776 Starter::Review("magi/eba2/A".to_owned()),
7777 "a review wins over resuming the old run"
7778 );
7779 }
7780
7781 #[test]
7782 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7783 assert_eq!(
7784 choose_starter(None, false, Some("some-run")),
7785 Starter::Resume("some-run".to_owned())
7786 );
7787 assert_eq!(choose_starter(None, false, None), Starter::Start);
7788 }
7789
7790 #[test]
7791 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7792 let mut released = task();
7793 released.start("stalled-run".to_owned());
7794 released.requeue();
7795 let unfinished = (!released.fresh_start)
7796 .then(|| Some("stalled-run".to_owned()))
7797 .flatten();
7798 assert_eq!(
7799 choose_starter(None, false, unfinished.as_deref()),
7800 Starter::Start,
7801 "release keeps run history but must not resume it"
7802 );
7803 assert_eq!(released.runs, ["stalled-run"]);
7804 }
7805
7806 #[test]
7807 fn an_ordinary_release_keeps_a_resumable_run_available() {
7808 let mut released = task();
7809 released.start("stalled-run".to_owned());
7810 released.release();
7811 let unfinished = (!released.fresh_start)
7812 .then(|| Some("stalled-run".to_owned()))
7813 .flatten();
7814 assert_eq!(
7815 choose_starter(None, false, unfinished.as_deref()),
7816 Starter::Resume("stalled-run".to_owned()),
7817 "manual release must preserve the normal resume path"
7818 );
7819 }
7820
7821 #[test]
7822 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7823 let mut state = run_state(RunStatus::Blocked);
7824 state.config.graph.review_rounds = 3;
7825 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7826 assert!(exhausted_review_budget(&state));
7827
7828 state.reviews.pop();
7830 assert!(!exhausted_review_budget(&state));
7831
7832 let mut stalled = run_state(RunStatus::Stalled);
7835 stalled.config.graph.review_rounds = 1;
7836 stalled.reviews = vec![review_round(1)];
7837 assert!(!exhausted_review_budget(&stalled));
7838 }
7839
7840 fn review_round(round: usize) -> crate::run::ReviewRound {
7841 crate::run::ReviewRound {
7842 round,
7843 head: "deadbeef".to_owned(),
7844 verified_head: None,
7845 verified_at: None,
7846 reviews: Vec::new(),
7847 e2e: Vec::new(),
7848 verify_retried: false,
7849 e2e_deferred: false,
7850 e2e_defer_reason: None,
7851 fix: None,
7852 blocking: 0,
7853 answered: 1,
7854 expected: 1,
7855 clean: false,
7856 progressed: true,
7857 vote_split: false,
7858 reconsideration: Vec::new(),
7859 verdict: None,
7860 }
7861 }
7862
7863 #[test]
7864 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7865 let mut run = RunState::new(
7866 PathBuf::from("/repo"),
7867 "main".to_owned(),
7868 "abc1234def".to_owned(),
7869 "add retries".to_owned(),
7870 Config::default(),
7871 );
7872 run.status = RunStatus::Judging;
7873 run.parked = true;
7874 let mut task = Task::new(
7875 "add retries".to_owned(),
7876 "add retries".to_owned(),
7877 PathBuf::from("/repo"),
7878 crate::queue::Source::Human,
7879 );
7880 task.status = TaskStatus::Failed;
7881 task.runs = vec![run.id.clone()];
7882 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7883 assert!(with(&task, &run), "parked after judging is the case");
7884
7885 let mut not_parked = run.clone();
7886 not_parked.parked = false;
7887 not_parked.status = RunStatus::Stalled;
7888 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7889
7890 let mut fresh = task.clone();
7891 fresh.fresh_start = true;
7892 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7893
7894 let mut review = task.clone();
7895 review.review_branch = Some("magi/x/A".to_owned());
7896 assert!(!with(&review, &run), "review is ranked before resume");
7897
7898 let mut held = task.clone();
7899 held.status = TaskStatus::Held;
7900 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7901
7902 let mut released = run.clone();
7903 released.released_to = Some("20260901-000000-new1".to_owned());
7904 assert!(!with(&task, &released), "nothing left to resume into");
7905
7906 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7907 "unreadable"
7908 )));
7909 }
7910
7911 #[test]
7912 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7913 let mut released = RunState::new(
7914 PathBuf::from("/repo"),
7915 "main".to_owned(),
7916 "abc1234def".to_owned(),
7917 "add retries".to_owned(),
7918 Config::default(),
7919 );
7920 released.status = RunStatus::Blocked;
7921 assert_eq!(
7922 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7923 Some(released.id.clone())
7924 );
7925 released.released_to = Some("20260901-000000-new1".to_owned());
7926 assert_eq!(
7927 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7928 None,
7929 "there is nothing left to resume it into"
7930 );
7931 }
7932
7933 #[test]
7934 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
7935 let mut exhausted = RunState::new(
7944 PathBuf::from("/repo"),
7945 "main".to_owned(),
7946 "abc1234def".to_owned(),
7947 "add retries".to_owned(),
7948 Config::default(),
7949 );
7950 exhausted.status = RunStatus::Blocked;
7951 exhausted.config.graph.review_rounds = 1;
7952 exhausted.reviews = vec![review_round(1)];
7953
7954 assert_eq!(
7955 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
7956 None,
7957 "an exhausted `Blocked` run must not be offered as resumable"
7958 );
7959
7960 let mut has_budget_left = RunState::new(
7963 PathBuf::from("/repo"),
7964 "main".to_owned(),
7965 "abc1234def".to_owned(),
7966 "add retries".to_owned(),
7967 Config::default(),
7968 );
7969 has_budget_left.status = RunStatus::Blocked;
7970 has_budget_left.config.graph.review_rounds = 3;
7971 has_budget_left.reviews = vec![review_round(1)];
7972
7973 assert_eq!(
7974 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
7975 Ok(has_budget_left.clone())
7976 }),
7977 Some(has_budget_left.id.clone())
7978 );
7979 }
7980
7981 #[test]
7982 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
7983 let mut older_stalled = RunState::new(
7991 PathBuf::from("/repo"),
7992 "main".to_owned(),
7993 "abc1234def".to_owned(),
7994 "add retries".to_owned(),
7995 Config::default(),
7996 );
7997 older_stalled.status = RunStatus::Stalled;
7998
7999 let mut newest_exhausted = RunState::new(
8000 PathBuf::from("/repo"),
8001 "main".to_owned(),
8002 "abc1234def".to_owned(),
8003 "add retries".to_owned(),
8004 Config::default(),
8005 );
8006 newest_exhausted.status = RunStatus::Blocked;
8007 newest_exhausted.config.graph.review_rounds = 1;
8008 newest_exhausted.reviews = vec![review_round(1)];
8009
8010 assert_eq!(
8011 unfinished_run_with(
8012 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
8013 "t",
8014 |_| Ok(newest_exhausted.clone())
8015 ),
8016 None,
8017 "the newest run is exhausted, so nothing here is worth resuming - \
8018 least of all the older, already-superseded run"
8019 );
8020 }
8021
8022 #[test]
8023 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
8024 assert_eq!(
8025 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
8026 Err(anyhow::anyhow!("fixture is absent"))
8027 }),
8028 None
8029 );
8030 }
8031
8032 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
8033 let mut q = ask::Question::new(
8034 run.to_owned(),
8035 "implement".to_owned(),
8036 "impl-A".to_owned(),
8037 "continue?".to_owned(),
8038 String::new(),
8039 vec!["resume で続行する".to_owned(), "other".to_owned()],
8040 );
8041 q.actions.insert("resume で続行する".to_owned(), action);
8042 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
8043 .unwrap();
8044 q
8045 }
8046
8047 fn held_task_with(run: &str) -> Task {
8048 let mut t = task();
8049 t.runs = vec![run.to_owned()];
8050 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
8051 t
8052 }
8053
8054 fn resume_action(run: &str) -> ask::ChoiceAction {
8055 ask::ChoiceAction::Resume { run: run.into() }
8056 }
8057
8058 #[test]
8059 fn decide_action_resumes_only_the_latest_resumable_run() {
8060 let t = held_task_with("r1");
8061 let q = action_question("r1", resume_action("r1"));
8062 let load = |s: RunState| move |_: &str| Ok(s);
8063 assert_eq!(
8064 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
8065 ActionDecision::Resume("r1".into())
8066 );
8067 let q_other = action_question("r1", resume_action("r0"));
8069 assert!(matches!(
8070 decide_action(
8071 &t,
8072 &q_other,
8073 &PHRASES_EN,
8074 load(run_state(RunStatus::Blocked))
8075 ),
8076 ActionDecision::Refuse(_)
8077 ));
8078 assert!(matches!(
8080 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
8081 ActionDecision::Refuse(_)
8082 ));
8083 let mut released = run_state(RunStatus::Blocked);
8085 released.released_to = Some("elsewhere".into());
8086 assert!(matches!(
8087 decide_action(&t, &q, &PHRASES_EN, load(released)),
8088 ActionDecision::Refuse(_)
8089 ));
8090 assert!(matches!(
8092 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
8093 ActionDecision::Refuse(_)
8094 ));
8095 }
8096
8097 #[test]
8098 fn decide_action_ignores_a_question_about_an_earlier_run() {
8099 let mut t = held_task_with("r1");
8100 t.runs.push("r2".to_owned());
8101 let never = |_: &str| -> Result<RunState> { bail!("not read") };
8102 assert_eq!(
8103 decide_action(
8104 &t,
8105 &action_question("r1", ask::ChoiceAction::Done),
8106 &PHRASES_EN,
8107 never
8108 ),
8109 ActionDecision::Stale
8110 );
8111 }
8112
8113 #[test]
8114 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
8115 let mut t = held_task_with("r1");
8116 let never = |_: &str| -> Result<RunState> { bail!("not read") };
8117 assert_eq!(
8118 decide_action(
8119 &t,
8120 &action_question("r1", ask::ChoiceAction::Requeue),
8121 &PHRASES_EN,
8122 never
8123 ),
8124 ActionDecision::Requeue
8125 );
8126 let done_q = action_question("r1", ask::ChoiceAction::Done);
8127 assert_eq!(
8128 decide_action(&t, &done_q, &PHRASES_EN, never),
8129 ActionDecision::Done
8130 );
8131 t.mark_action_applied(&done_q.id);
8132 assert_eq!(
8133 decide_action(&t, &done_q, &PHRASES_EN, never),
8134 ActionDecision::Skip
8135 );
8136
8137 let mut plain = action_question("r1", ask::ChoiceAction::Done);
8139 plain.actions.clear();
8140 assert_eq!(
8141 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
8142 ActionDecision::Skip
8143 );
8144 let mut running = held_task_with("r1");
8146 running.status = TaskStatus::Running;
8147 assert_eq!(
8148 decide_action(
8149 &running,
8150 &action_question("r1", ask::ChoiceAction::Done),
8151 &PHRASES_EN,
8152 never
8153 ),
8154 ActionDecision::Skip
8155 );
8156 }
8157
8158 #[test]
8159 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
8160 let dir = tempfile::tempdir().unwrap();
8161 let queue = Queue::at(dir.path().join("queue"));
8162 let questions = Questions::at(dir.path().join("questions"));
8163 let home = dir.path().join("home");
8164 let mut state = run_state(RunStatus::Blocked);
8165 state.id = "20260101-000000-act1".to_owned();
8166 state.save_under(&home).unwrap();
8167
8168 let mut t = held_task_with(&state.id);
8169 queue.put(&mut t).unwrap();
8170 let mut q = action_question(&state.id, resume_action(&state.id));
8171 questions.put(&mut q).unwrap();
8172
8173 apply_choice_actions(&queue, &questions, &home);
8174 let after = queue.get(&t.id).unwrap();
8175 assert_eq!(after.status, TaskStatus::Queued);
8176 assert!(!after.fresh_start);
8177 assert!(after.action_applied(&q.id));
8178 let pin = after.resume_override.clone().unwrap();
8179 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
8180 assert!(pin.forced);
8181
8182 let mut again = queue.get(&t.id).unwrap();
8184 again.hold_machine(Some("later".into()));
8185 queue.put(&mut again).unwrap();
8186 apply_choice_actions(&queue, &questions, &home);
8187 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8188 }
8189
8190 #[test]
8191 fn a_delivered_answer_is_still_acted_on_exactly_once() {
8192 let dir = tempfile::tempdir().unwrap();
8193 let queue = Queue::at(dir.path().join("queue"));
8194 let questions = Questions::at(dir.path().join("questions"));
8195 let home = dir.path().join("home");
8196 let mut state = run_state(RunStatus::Blocked);
8197 state.id = "20260101-000000-act2".to_owned();
8198 state.save_under(&home).unwrap();
8199
8200 let mut t = held_task_with(&state.id);
8201 queue.put(&mut t).unwrap();
8202 let mut q = action_question(&state.id, resume_action(&state.id));
8203 q.answer_delivered = true;
8204 questions.put(&mut q).unwrap();
8205
8206 apply_choice_actions(&queue, &questions, &home);
8209 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
8210 assert!(queue.get(&t.id).unwrap().action_applied(&q.id));
8211
8212 let mut again = queue.get(&t.id).unwrap();
8214 again.hold_machine(Some("later".into()));
8215 queue.put(&mut again).unwrap();
8216 apply_choice_actions(&queue, &questions, &home);
8217 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8218 }
8219
8220 #[test]
8221 fn a_fresh_asker_defers_the_action_until_it_goes_quiet() {
8222 let dir = tempfile::tempdir().unwrap();
8223 let queue = Queue::at(dir.path().join("queue"));
8224 let questions = Questions::at(dir.path().join("questions"));
8225 let home = dir.path().join("home");
8226 let mut state = run_state(RunStatus::Blocked);
8227 state.id = "20260101-000000-act3".to_owned();
8228 state.save_under(&home).unwrap();
8229 let mut t = held_task_with(&state.id);
8230 queue.put(&mut t).unwrap();
8231 let mut q = action_question(&state.id, ask::ChoiceAction::Requeue);
8232 questions.put(&mut q).unwrap();
8233
8234 questions.beat(&q.id, crate::ask::WaiterKind::Asker);
8235 apply_choice_actions(&queue, &questions, &home);
8236 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8237
8238 std::fs::remove_file(questions.root().join(format!("{}.lease", q.id))).unwrap();
8239 apply_choice_actions(&queue, &questions, &home);
8240 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
8241 }
8242
8243 #[test]
8244 fn action_standing_tells_the_reasons_apart() {
8245 let mut t = held_task_with("r1");
8246 let q = action_question("r1", ask::ChoiceAction::Requeue);
8247 assert_eq!(action_standing(&t, &q), ActionStanding::Pending);
8248 let mut plain = q.clone();
8249 plain.actions.clear();
8250 assert_eq!(action_standing(&t, &plain), ActionStanding::NoAction);
8251
8252 t.mark_action_applied(&q.id);
8254 assert_eq!(action_standing(&t, &q), ActionStanding::Applied);
8255
8256 let mut t = held_task_with("r1");
8258 t.start("r2".to_owned());
8259 assert_eq!(action_standing(&t, &q), ActionStanding::Stale);
8260 let q2 = action_question("r2", ask::ChoiceAction::Requeue);
8262 assert_eq!(action_standing(&t, &q2), ActionStanding::Busy);
8263 t.status = TaskStatus::Blocked;
8264 assert_eq!(action_standing(&t, &q2), ActionStanding::Busy);
8265
8266 let mut c = action_question("r0", ask::ChoiceAction::Requeue);
8268 c.node = crate::conduct::NODE.to_owned();
8269 assert_eq!(action_standing(&t, &c), ActionStanding::Busy);
8270 assert!(!ActionStanding::Busy.daemon_owns());
8271 }
8272
8273 #[test]
8274 fn a_running_task_waits_and_a_dead_asker_with_a_cwd_is_still_actioned() {
8275 let dir = tempfile::tempdir().unwrap();
8276 let queue = Queue::at(dir.path().join("queue"));
8277 let questions = Questions::at(dir.path().join("questions"));
8278 let home = dir.path().join("home");
8279 let mut state = run_state(RunStatus::Blocked);
8280 state.id = "20260101-000000-act4".to_owned();
8281 state.save_under(&home).unwrap();
8282 let mut t = held_task_with(&state.id);
8283 t.status = TaskStatus::Running;
8284 queue.put(&mut t).unwrap();
8285 let mut q = action_question(&state.id, ask::ChoiceAction::Done);
8286 q.cwd = Some(dir.path().display().to_string());
8287 questions.put(&mut q).unwrap();
8288
8289 let running = queue.get(&t.id).unwrap();
8292 assert_eq!(action_standing(&running, &q), ActionStanding::Busy);
8293 assert!(matches!(
8294 crate::waiter::decide_owned(&q, None, false, 86_400, Timestamp::now(), false),
8295 crate::waiter::Action::Deliver(_)
8296 ));
8297 apply_choice_actions(&queue, &questions, &home);
8299 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
8300 let mut back = queue.get(&t.id).unwrap();
8301 back.hold_machine(Some("later".into()));
8302 queue.put(&mut back).unwrap();
8303 apply_choice_actions(&queue, &questions, &home);
8304 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Done);
8305 let mut again = queue.get(&t.id).unwrap();
8306 again.hold_machine(Some("again".into()));
8307 queue.put(&mut again).unwrap();
8308 apply_choice_actions(&queue, &questions, &home);
8309 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8310 }
8311}