1use std::path::{Path, PathBuf};
50use std::sync::Arc;
51use std::sync::atomic::{AtomicBool, Ordering};
52use std::sync::{Mutex, MutexGuard};
53use std::time::Duration;
54
55use anyhow::{Context, Result, bail};
56use jiff::Timestamp;
57use serde::{Deserialize, Serialize};
58use tokio::sync::Notify;
59
60use crate::ask::{self, Questions};
61use crate::clean;
62use crate::conduct::Conductor;
63use crate::config::{Config, MergeMode};
64use crate::graph::Runner;
65use crate::land;
66use crate::notices::{self, Link, Notice};
67use crate::queue::{Queue, Task, TaskStatus};
68use crate::run::{Liveness, QuotaLoss, RunState, RunStatus};
69use crate::triage;
70
71pub const SCHEMA: u32 = 1;
73
74pub const HEARTBEAT: Duration = Duration::from_secs(5);
78
79pub const STALE_SECS: i64 = 30;
88
89pub const POLL: Duration = Duration::from_secs(5);
91
92pub const STALE_CLAIM: Duration = Duration::from_secs(6 * 60 * 60);
96
97pub const STALLED_RUNNING: Duration = Duration::from_secs(30 * 60);
116
117#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
119#[serde(default)]
120pub struct Current {
121 pub task: String,
123 pub run: String,
125}
126
127#[derive(Debug, Clone, Serialize, Deserialize)]
134pub struct Status {
135 pub schema: u32,
137 pub pid: u32,
139 pub started_at: Timestamp,
141 pub updated_at: Timestamp,
143 pub idle: bool,
145 pub current: Vec<Current>,
151 pub completed: usize,
153 pub polls: u64,
155}
156
157impl Status {
158 #[must_use]
160 pub fn new() -> Self {
161 let now = Timestamp::now();
162 Self {
163 schema: SCHEMA,
164 pid: std::process::id(),
165 started_at: now,
166 updated_at: now,
167 idle: true,
168 current: Vec::new(),
169 completed: 0,
170 polls: 0,
171 }
172 }
173}
174
175impl Default for Status {
176 fn default() -> Self {
177 Self::new()
178 }
179}
180
181#[derive(Debug, Clone)]
183pub struct Opts {
184 pub repo: PathBuf,
186 pub config: Option<PathBuf>,
188 pub poll: Duration,
190 pub max_attempts: usize,
192 pub once: bool,
194 pub merge: Option<String>,
196 pub worktrees_root: Option<PathBuf>,
205}
206
207impl Default for Opts {
208 fn default() -> Self {
209 Self {
210 repo: PathBuf::from("."),
211 config: None,
212 poll: POLL,
213 max_attempts: 2,
214 once: false,
215 merge: None,
216 worktrees_root: None,
217 }
218 }
219}
220
221fn max_concurrent(n: usize) -> usize {
226 n.max(1)
227}
228
229#[must_use]
231pub fn status_path() -> PathBuf {
232 crate::run::home().join("daemon.json")
233}
234
235pub fn write_status(status: &Status) -> Result<()> {
237 write_status_to(&status_path(), status)
238}
239
240pub fn write_status_to(path: &Path, status: &Status) -> Result<()> {
245 if let Some(parent) = path.parent() {
246 std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
247 }
248 let body = serde_json::to_string_pretty(status).context("serialize daemon status")?;
249 let tmp = path.with_extension("json.tmp");
250 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
251 std::fs::rename(&tmp, path).with_context(|| format!("replace {}", path.display()))?;
252 Ok(())
253}
254
255pub fn clear_status() {
258 clear_status_at(&status_path());
259}
260
261fn clear_status_at(path: &Path) {
265 let _ = std::fs::remove_file(path);
266}
267
268#[derive(Debug, Clone, Default)]
281pub struct Stop {
282 stopped: Arc<AtomicBool>,
286 busy: Arc<std::sync::atomic::AtomicUsize>,
292 wake: Arc<Notify>,
296 pause: crate::graph::Pause,
299}
300
301impl Stop {
302 #[must_use]
304 pub fn new() -> Self {
305 Self::default()
306 }
307
308 pub fn stop(&self) {
311 self.stopped.store(true, Ordering::SeqCst);
312 self.wake.notify_one();
316 }
317
318 #[must_use]
320 pub fn stopped(&self) -> bool {
321 self.stopped.load(Ordering::SeqCst)
322 }
323
324 #[must_use]
332 pub fn finishing(&self) -> bool {
333 self.stopped() && self.busy_now()
334 }
335
336 pub fn park(&self) {
347 self.pause.park();
348 self.stop();
349 }
350
351 #[must_use]
353 pub fn parking(&self) -> bool {
354 self.pause.parked()
355 }
356
357 #[must_use]
359 pub fn pause(&self) -> crate::graph::Pause {
360 self.pause.clone()
361 }
362
363 #[must_use]
369 pub fn busy_now(&self) -> bool {
370 self.busy.load(Ordering::SeqCst) > 0
371 }
372
373 fn enter(&self) {
375 self.busy.fetch_add(1, Ordering::SeqCst);
376 }
377
378 fn exit(&self) {
381 self.busy.fetch_sub(1, Ordering::SeqCst);
382 }
383
384 async fn idle(&self, poll: Duration) {
386 tokio::select! {
387 () = tokio::time::sleep(poll) => {}
388 () = self.wake.notified() => {}
389 }
390 }
391}
392
393#[derive(Debug, Clone, Default, Deserialize)]
400#[serde(default)]
401pub struct Reading {
402 pub schema: u32,
404 pub pid: Option<u32>,
406 pub started_at: Option<Timestamp>,
408 pub updated_at: Option<Timestamp>,
410 pub idle: bool,
412 #[serde(deserialize_with = "de_current")]
424 pub current: Vec<Current>,
425 pub completed: u64,
427 pub polls: u64,
429}
430
431fn de_current<'de, D>(deserializer: D) -> std::result::Result<Vec<Current>, D::Error>
434where
435 D: serde::Deserializer<'de>,
436{
437 #[derive(Deserialize)]
438 #[serde(untagged)]
439 enum Shape {
440 Many(Vec<Current>),
441 One(Current),
442 }
443 Ok(
444 Option::<Shape>::deserialize(deserializer)?.map_or_else(Vec::new, |shape| match shape {
445 Shape::Many(v) => v,
446 Shape::One(c) => vec![c],
447 }),
448 )
449}
450
451impl Reading {
452 #[must_use]
455 pub fn age_secs(&self, now: Timestamp) -> Option<i64> {
456 self.updated_at
457 .map(|at| (now.as_second() - at.as_second()).max(0))
458 }
459
460 #[must_use]
464 pub fn running(&self, now: Timestamp) -> bool {
465 self.age_secs(now).is_some_and(|secs| secs <= STALE_SECS)
466 }
467}
468
469#[must_use]
476pub fn read_status(home: &Path) -> Option<Reading> {
477 let body = std::fs::read_to_string(home.join("daemon.json")).ok()?;
478 serde_json::from_str(&body).ok()
479}
480
481#[must_use]
493pub fn current_work(home: &Path, now: Timestamp) -> Vec<Current> {
494 read_status(home)
495 .filter(|reading| reading.running(now))
496 .map(|reading| reading.current)
497 .unwrap_or_default()
498}
499
500#[must_use]
502pub fn is_working_on(home: &Path, run: &str, now: Timestamp) -> bool {
503 current_work(home, now).iter().any(|c| c.run == run)
504}
505
506#[must_use]
516pub fn is_working_on_short(home: &Path, short: &str, now: Timestamp) -> bool {
517 current_work(home, now)
518 .iter()
519 .any(|c| crate::run::short_of(&c.run) == short)
520}
521
522#[must_use]
524pub fn is_working_on_task(home: &Path, task: &str, now: Timestamp) -> bool {
525 current_work(home, now).iter().any(|c| c.task == task)
526}
527
528pub fn sweep_stale_claims(queue: &Queue, older_than: Duration) -> Vec<String> {
572 sweep_stale_claims_with(queue, older_than, crate::proc::pid_alive)
573}
574
575fn sweep_stale_claims_with<F>(queue: &Queue, older_than: Duration, pid_alive: F) -> Vec<String>
579where
580 F: Fn(u32) -> bool,
581{
582 let this_process = std::process::id();
583 let mut swept: Vec<String> = std::fs::read_dir(queue.root())
584 .into_iter()
585 .flatten()
586 .flatten()
587 .map(|e| e.path())
588 .filter(|p| p.extension().is_some_and(|x| x == "lock"))
589 .filter(|p| {
590 match std::fs::read_to_string(p)
591 .ok()
592 .and_then(|body| body.trim().parse::<u32>().ok())
593 {
594 Some(pid) if pid == this_process => false,
598 Some(pid) => !pid_alive(pid),
599 None => p
600 .metadata()
601 .and_then(|m| m.modified())
602 .and_then(|t| t.elapsed().map_err(std::io::Error::other))
603 .is_ok_and(|age| age >= older_than),
604 }
605 })
606 .filter(|p| std::fs::remove_file(p).is_ok())
607 .filter_map(|p| {
608 p.file_stem()
609 .and_then(|s| s.to_str())
610 .map(std::borrow::ToOwned::to_owned)
611 })
612 .collect();
613 swept.sort_unstable();
614 swept
615}
616
617fn is_stalled(task: &Task, home: &Path, now: Timestamp) -> bool {
622 task.status == TaskStatus::Running
623 && (now.as_second() - task.updated_at.as_second()) >= STALLED_RUNNING.as_secs() as i64
624 && !is_working_on_task(home, &task.id, now)
625}
626
627fn stalled_tasks(queue: &Queue, home: &Path, now: Timestamp) -> Vec<Task> {
630 queue
631 .list()
632 .into_iter()
633 .filter(|t| is_stalled(t, home, now))
634 .collect()
635}
636
637fn queued_tasks(queue: &Queue) -> Vec<Task> {
643 queue
644 .list()
645 .into_iter()
646 .filter(|t| t.status == TaskStatus::Queued)
647 .collect()
648}
649
650fn finished_tasks(queue: &Queue) -> Vec<Task> {
653 queue
654 .list()
655 .into_iter()
656 .filter(|t| matches!(t.status, TaskStatus::Failed | TaskStatus::Held))
657 .collect()
658}
659
660fn resolve_blockers(queue: &Queue, questions: &Questions) {
677 for listed in queue.list() {
678 if listed.status != TaskStatus::Blocked || listed.blocked_by.is_empty() {
679 continue;
680 }
681 let Ok(_claim) = queue.claim(&listed.id) else {
682 continue;
683 };
684 let Ok(mut task) = queue.get(&listed.id) else {
685 continue;
686 };
687 if task.status != TaskStatus::Blocked {
688 continue;
689 }
690 let deleted = queue.apply_deleted_blockers(&mut task);
694 if !deleted.is_empty() {
695 record(queue, &mut task);
696 for id in &deleted {
697 queue.note_dependency_deleted(&task, id);
698 }
699 if task.status != TaskStatus::Blocked {
700 continue;
701 }
702 }
703 let missing = crate::queue::missing_blockers(queue, questions, &task.blocked_by);
704 if !missing.is_empty() {
705 let language = language_of(&task, Path::new("."));
706 task.hold_machine(Some(crate::queue::missing_blocker_hold_reason_in(
707 &task.blocked_by,
708 &missing,
709 &language,
710 )));
711 record(queue, &mut task);
712 continue;
713 }
714 if let Some(q) = task.blocked_by.iter().find_map(|id| {
717 questions.get(id).ok().filter(|q| {
718 q.node == crate::conduct::NODE && q.status == ask::QuestionStatus::Abandoned
719 })
720 }) {
721 let language = language_of(&task, Path::new("."));
722 task.hold_machine(Some(unanswered_question_hold_reason(&q, &language)));
723 record(queue, &mut task);
724 continue;
725 }
726 let mut changed = false;
727 for id in task.blocked_by.clone() {
728 if let Ok(dep) = queue.get(&id) {
729 if dep.status == TaskStatus::Done {
730 task.unblock(&id);
731 changed = true;
732 }
733 continue;
734 }
735 if let Ok(q) = questions.get(&id)
736 && q.status == ask::QuestionStatus::Answered
737 {
738 let answer = match &q.answer {
739 Some(ask::Answer::Choice(c) | ask::Answer::Text(c)) => c.clone(),
740 None => String::new(),
741 };
742 task.record_answer(q.summary.clone(), answer);
743 task.unblock(&id);
744 changed = true;
745 }
746 }
747 if changed {
748 record(queue, &mut task);
749 }
750 }
751}
752
753fn unanswered_question_hold_reason(q: &ask::Question, language: &str) -> String {
755 if crate::lang::is_japanese(language) {
756 format!(
757 "質問 {} 「{}」 に期限内の回答がなく、取り下げられました - `magi task triage` を参照",
758 q.short(),
759 q.summary
760 )
761 } else {
762 format!(
763 "question {} \"{}\" went unanswered and was abandoned - see `magi task triage`",
764 q.short(),
765 q.summary
766 )
767 }
768}
769
770#[derive(Debug, Clone, Copy, PartialEq, Eq)]
773pub(crate) enum ActionStanding {
774 NoAction,
776 Pending,
778 Applied,
781 Stale,
783 Busy,
785}
786
787impl ActionStanding {
788 pub(crate) fn daemon_owns(self) -> bool {
790 matches!(self, Self::Pending | Self::Applied | Self::Stale)
791 }
792}
793
794pub(crate) fn action_standing(task: &Task, q: &ask::Question) -> ActionStanding {
798 if q.chosen_action().is_none() {
799 ActionStanding::NoAction
800 } else if task.action_applied(&q.id) {
801 ActionStanding::Applied
802 } else if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
803 ActionStanding::Stale
804 } else if matches!(
805 task.status,
806 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
807 ) {
808 ActionStanding::Busy
809 } else {
810 ActionStanding::Pending
811 }
812}
813
814#[derive(Debug, Clone, PartialEq, Eq)]
816enum ActionDecision {
817 Skip,
819 Resume(String),
821 Requeue,
823 Done,
825 Stale,
828 Refuse(String),
831}
832
833fn decide_action<F>(task: &Task, q: &ask::Question, p: &Phrases, load: F) -> ActionDecision
841where
842 F: FnOnce(&str) -> Result<RunState>,
843{
844 let Some(action) = q.chosen_action() else {
845 return ActionDecision::Skip;
846 };
847 if task.action_applied(&q.id)
848 || matches!(
849 task.status,
850 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
851 )
852 {
853 return ActionDecision::Skip;
854 }
855 if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
858 return ActionDecision::Stale;
859 }
860 match action {
861 ask::ChoiceAction::Requeue => ActionDecision::Requeue,
862 ask::ChoiceAction::Done => ActionDecision::Done,
863 ask::ChoiceAction::Resume { run } => {
864 if task.runs.last() != Some(run) {
865 return ActionDecision::Refuse((p.resume_not_latest)(
866 q.short(),
867 ask::short_id(run),
868 ));
869 }
870 match load(run) {
871 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
872 ActionDecision::Resume(run.clone())
873 }
874 Ok(_) => ActionDecision::Refuse((p.resume_cannot_progress)(
875 q.short(),
876 ask::short_id(run),
877 )),
878 Err(e) => ActionDecision::Refuse((p.resume_unreadable)(
879 q.short(),
880 ask::short_id(run),
881 &format!("{e:#}"),
882 )),
883 }
884 }
885 }
886}
887
888struct Phrases {
902 graph_stopped: fn(&str, &str) -> String,
904 quorum_lost: &'static str,
905 quota_took_out: &'static str,
907 run_ended: &'static str,
909 waiting_for_answer: &'static str,
911 recovered_running: &'static str,
912 no_run_to_recover: &'static str,
913 could_not_start: &'static str,
914 resume_not_latest: fn(&str, &str) -> String,
915 resume_cannot_progress: fn(&str, &str) -> String,
916 resume_unreadable: fn(&str, &str, &str) -> String,
917 handover_refused: fn(&crate::handover::Refused) -> String,
920 handover_hint: &'static str,
922 reviewers_never_answered: fn(usize, usize) -> String,
925}
926
927const PHRASES_EN: Phrases = Phrases {
928 graph_stopped: |status, detail| {
929 format!("the graph stopped at `{status}` without reaching a terminal status: {detail}")
930 },
931 quorum_lost: "the judging panel lost its quorum",
932 quota_took_out: "; quota took out ",
933 run_ended: "run ended ",
934 waiting_for_answer: " - waiting for operator answer to question ",
935 recovered_running: "recovered a `running` task whose daemon never recorded the outcome: ",
936 no_run_to_recover: "task was `running` with no live daemon and no readable \
937 run to recover; held for a human to check what happened",
938 could_not_start: "could not start the run: ",
939 resume_not_latest: |q, run| {
940 format!("question {q} asked to resume run {run}, which is not this task's latest run")
941 },
942 resume_cannot_progress: |q, run| {
943 format!("question {q} asked to resume run {run}, which cannot make progress")
944 },
945 resume_unreadable: |q, run, e| {
946 format!("question {q} asked to resume run {run}, which could not be read: {e}")
947 },
948 handover_refused: |r| r.to_string(),
949 handover_hint: " (clean up the other worktree, then release the task from the queue)",
950 reviewers_never_answered: |missing, rounds| {
951 format!(
952 "{missing} reviewer seat(s) never answered after {rounds} rounds; \
953 refusing to call it clean"
954 )
955 },
956};
957
958const PHRASES_JA: Phrases = Phrases {
959 graph_stopped: |status, detail| {
960 format!("グラフが終端状態に達しないまま `{status}` で停止しました: {detail}")
961 },
962 quorum_lost: "審査パネルが定足数を失いました",
963 quota_took_out: "。クォータで脱落: ",
964 run_ended: "run 終了: ",
965 waiting_for_answer: " - オペレーターの回答待ち: 質問 ",
966 recovered_running: "daemon が結果を記録しないまま `running` だったタスクを回収しました: ",
967 no_run_to_recover: "タスクは `running` でしたが、生きた daemon も回収できる run も見つかりません。\
968 何が起きたか人が確認するため保留にしました",
969 could_not_start: "run を開始できませんでした: ",
970 resume_not_latest: |q, run| {
971 format!(
972 "質問 {q} は run {run} の再開を求めましたが、これはタスクの最新の run ではありません"
973 )
974 },
975 resume_cannot_progress: |q, run| {
976 format!("質問 {q} は run {run} の再開を求めましたが、これは進行できません")
977 },
978 resume_unreadable: |q, run, e| {
979 format!("質問 {q} は run {run} の再開を求めましたが、読み込めませんでした: {e}")
980 },
981 handover_refused: |r| {
982 use crate::handover::Refused;
983 match r {
984 Refused::Foreign { branch, path, why } => format!(
985 "ブランチ `{branch}` は {path} にチェックアウトされており、magi は自動では\
986 削除しません。不要なら `git worktree remove` でその worktree を削除して\
987 から、やり直してください(詳細: {why})"
988 ),
989 Refused::Unsafe { branch, path, why } => format!(
990 "ブランチ `{branch}` は {path} にチェックアウトされています。そこの作業を\
991 コミットか破棄したうえで `git worktree remove` で worktree を削除する\
992 か、run を破棄してよいと伝えてから、やり直してください(詳細: {why})"
993 ),
994 Refused::ReleaseFailed { branch, path, run } => format!(
995 "ブランチ `{branch}` は {path} で run {run} がチェックアウトしており、その\
996 worktree の解放に失敗したか、変更が見つかりました(未コミットの変更が\
997 ある worktree は git が削除を拒否します)。worktree はそのまま残しました"
998 ),
999 }
1000 },
1001 handover_hint: "(他の worktree を片付けてから、タスクをキューから解放してください)",
1002 reviewers_never_answered: |missing, rounds| {
1003 format!(
1004 "{missing} 席のレビュアーが {rounds} ラウンドの間に一度も回答しなかったため、\
1005 クリーンとは認めません"
1006 )
1007 },
1008};
1009
1010fn phrases(language: &str) -> &'static Phrases {
1011 if crate::lang::is_japanese(language) {
1012 &PHRASES_JA
1013 } else {
1014 &PHRASES_EN
1015 }
1016}
1017
1018fn language_of(task: &Task, fallback: &Path) -> String {
1022 crate::lang::of_repo(&repo_for(task, fallback))
1023}
1024
1025pub(crate) fn task_of_question<'a>(tasks: &'a [Task], q: &ask::Question) -> Option<&'a Task> {
1029 if q.node == crate::conduct::NODE {
1030 return tasks.iter().find(|t| t.id == q.run);
1031 }
1032 tasks.iter().find(|t| t.runs.contains(&q.run))
1033}
1034
1035fn apply_choice_actions(queue: &Queue, questions: &Questions, home: &Path) {
1044 let tasks = queue.list();
1045 for q in questions.list() {
1046 if q.chosen_action().is_none() {
1047 continue;
1048 }
1049 let Some(listed) = task_of_question(&tasks, &q) else {
1050 continue;
1051 };
1052 if listed.action_applied(&q.id) {
1053 continue;
1054 }
1055 let Ok(_claim) = queue.claim(&listed.id) else {
1056 continue;
1057 };
1058 let Ok(mut task) = queue.get(&listed.id) else {
1059 continue;
1060 };
1061 let language = language_of(&task, Path::new("."));
1062 let decision = decide_action(&task, &q, phrases(&language), |id| {
1063 RunState::load_under(id, home)
1064 });
1065 let ran = matches!(
1066 decision,
1067 ActionDecision::Resume(_) | ActionDecision::Requeue | ActionDecision::Done
1068 );
1069 if ran && asker_may_still_read(questions, &q, home) {
1070 continue;
1079 }
1080 let done = matches!(decision, ActionDecision::Done);
1081 match decision {
1082 ActionDecision::Skip => continue,
1083 ActionDecision::Resume(run) => {
1084 task.release();
1085 task.resume_override = Some(crate::queue::OperatorResume {
1086 question_id: q.id.clone(),
1087 at: Timestamp::now(),
1088 conductor_rehold: None,
1089 forced: true,
1090 pinned_run: Some(run),
1091 });
1092 }
1093 ActionDecision::Stale => {}
1094 ActionDecision::Requeue => task.requeue(),
1095 ActionDecision::Done => task.succeed(),
1096 ActionDecision::Refuse(why) => {
1097 task.hold_machine(Some(why));
1098 notices::raise(
1099 Notice::warn(
1100 &format!("action:{}", q.id),
1101 "An answer asked the daemon to resume a run that cannot be resumed; the task stays held.",
1102 )
1103 .link(Link::Task {
1104 id: task.id.clone(),
1105 }),
1106 );
1107 }
1108 }
1109 task.mark_action_applied(&q.id);
1110 record(queue, &mut task);
1113 if done && queue.get(&task.id).is_ok_and(|t| t.action_applied(&q.id)) {
1114 supersede_prior_runs(&task, home);
1115 }
1116 let _ = questions.update(&q.id, |r| {
1119 r.answer_delivered = true;
1120 Ok(())
1121 });
1122 }
1123}
1124
1125fn asker_may_still_read(questions: &Questions, q: &ask::Question, home: &Path) -> bool {
1129 let now = Timestamp::now();
1130 if questions.read_lease(&q.id).is_some_and(|l| l.fresh(now)) {
1131 return true;
1132 }
1133 q.node != crate::conduct::NODE
1134 && RunState::load_under(&q.run, home).is_ok_and(|s| {
1135 s.seats_active()
1136 .any(|(k, a)| *k == q.seat && a.remaining_secs(now) > 0)
1137 })
1138}
1139
1140fn reconcile_task_questions(queue: &Queue, questions: &Questions) {
1151 let tasks = queue.list();
1152 let by_id: std::collections::BTreeMap<&str, &Task> =
1153 tasks.iter().map(|t| (t.id.as_str(), t)).collect();
1154 let referenced: std::collections::BTreeSet<&str> = tasks
1155 .iter()
1156 .flat_map(|task| task.blocked_by.iter().map(String::as_str))
1157 .collect();
1158
1159 for mut question in questions.list() {
1160 if !question.status.open() || question.node != crate::conduct::NODE {
1161 continue;
1162 }
1163 if referenced.contains(question.id.as_str()) {
1166 continue;
1167 }
1168 let Some(task) = by_id.get(question.run.as_str()) else {
1169 continue;
1170 };
1171 question.abandon(format!(
1172 "task {} no longer waits for this answer",
1173 task.short()
1174 ));
1175 if let Err(e) = questions.put(&mut question) {
1176 tracing::warn!(
1177 "could not retire question {} for task {}: {e:#}",
1178 question.short(),
1179 task.short()
1180 );
1181 }
1182 }
1183}
1184
1185#[derive(Debug, Clone, Copy)]
1192pub struct Verdict {
1193 pub status: RunStatus,
1195 pub left_pr: bool,
1197 pub quota_hit: bool,
1199 pub parked: bool,
1201 pub no_viable_candidates: bool,
1208}
1209
1210pub fn settle(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize) {
1267 settle_in(task, verdict, detail, max_attempts, &PHRASES_EN)
1268}
1269
1270fn settle_in(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize, p: &Phrases) {
1272 if verdict.parked {
1278 task.park(detail);
1279 return;
1280 }
1281 match verdict.status {
1282 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1283 RunStatus::AlreadyInBase => task.already_landed(detail),
1284 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1285 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1286 task.stall(detail)
1287 }
1288 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1289 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1290 RunStatus::Blocked => task.fail(detail, max_attempts),
1291 RunStatus::VerifiedNoop => task.handed_off(detail),
1292 other => task.fail((p.graph_stopped)(label(other), detail), max_attempts),
1293 }
1294}
1295
1296pub fn supersede_prior_runs(task: &Task, home: &Path) {
1345 let now = Timestamp::now();
1346 let last_run_succeeded = task
1353 .runs
1354 .last()
1355 .and_then(|id| RunState::load_under(id, home).ok())
1356 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1357 for id in task.superseded_attempts(last_run_succeeded) {
1358 let mut state = match RunState::load_under(id, home) {
1359 Ok(s) => s,
1360 Err(e) => {
1361 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1362 continue;
1363 }
1364 };
1365 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1366 continue;
1367 }
1368 let daemon_claims = is_working_on(home, id, now);
1369 if state.liveness(daemon_claims) == Liveness::Live {
1370 continue;
1371 }
1372 state.status = RunStatus::Superseded;
1373 if let Err(e) = state.save_under(home) {
1374 tracing::warn!("could not mark run {id} superseded: {e:#}");
1375 }
1376 }
1377}
1378
1379fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1401 for task in queue.list() {
1402 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1403 continue;
1404 }
1405 supersede_prior_runs(&task, home);
1406 }
1407}
1408
1409fn settle_and_diagnose(
1416 task: &mut Task,
1417 verdict: Verdict,
1418 detail: &str,
1419 max_attempts: usize,
1420 state: &RunState,
1421) {
1422 let p = phrases(&state.config.graph.language);
1423 settle_in(task, verdict, detail, max_attempts, p);
1424 if task.status == TaskStatus::Held {
1425 task.diagnostic = diagnostic(state);
1426 note_open_question(task, &state.id, p);
1427 }
1428}
1429
1430fn note_open_question(task: &mut Task, run: &str, p: &Phrases) {
1445 let Some(home) = crate::run::try_home() else {
1446 return;
1447 };
1448 let open = Questions::at(home.join("questions")).open_for(run);
1449 let Some(q) = open.first() else {
1450 return;
1451 };
1452 let base = task.hold_reason.clone().unwrap_or_default();
1453 task.hold_reason = Some(format!("{base}{}{}", p.waiting_for_answer, q.short()));
1454}
1455
1456fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize, language: &str) {
1467 match last_run {
1468 Some(state) => {
1469 let verdict = Verdict {
1470 status: state.status,
1471 left_pr: state.pr.is_some(),
1472 quota_hit: !state.quota.is_empty(),
1473 parked: state.parked,
1474 no_viable_candidates: state.viable().is_empty(),
1475 };
1476 let detail = format!(
1477 "{}{}",
1478 phrases(&state.config.graph.language).recovered_running,
1479 describe(&state)
1480 );
1481 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1482 }
1483 None => {
1484 let why = phrases(language).no_run_to_recover;
1487 task.last_error = Some(why.to_owned());
1488 task.hold_machine(Some(why.to_owned()));
1491 }
1492 }
1493}
1494
1495fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1517 let mut reclaimed = Vec::new();
1518 for listed in queue.list() {
1519 if listed.status != TaskStatus::Running {
1520 continue;
1521 }
1522 let Ok(_claim) = queue.claim(&listed.id) else {
1523 continue;
1524 };
1525 let Ok(mut task) = queue.get(&listed.id) else {
1529 continue;
1530 };
1531 if task.status != TaskStatus::Running {
1532 continue;
1533 }
1534 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1535 if let Some(state) = &last_run
1545 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1546 {
1547 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1548 }
1549 let language = if last_run.is_none() {
1550 language_of(&task, Path::new("."))
1551 } else {
1552 String::new()
1553 };
1554 reclaim(&mut task, last_run, max_attempts, &language);
1555 if task.status == TaskStatus::Done {
1556 supersede_prior_runs(&task, &crate::run::home());
1557 }
1558 record(queue, &mut task);
1559 reclaimed.push(task.id.clone());
1560 }
1561 reclaimed
1562}
1563
1564fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1596 reclaim_abandoned_runs_with(
1597 home,
1598 now,
1599 crate::proc::pid_status,
1600 crate::proc::process_started_at,
1601 )
1602}
1603
1604fn reclaim_abandoned_runs_with<F, G>(
1610 home: &Path,
1611 now: Timestamp,
1612 query: F,
1613 identity: G,
1614) -> Vec<String>
1615where
1616 F: Fn(u32) -> Option<bool>,
1617 G: Fn(u32) -> Option<String>,
1618{
1619 let mut abandoned = Vec::new();
1620 for entry in std::fs::read_dir(home.join("runs"))
1621 .into_iter()
1622 .flatten()
1623 .flatten()
1624 {
1625 let id = entry.file_name().to_string_lossy().into_owned();
1626 if !crate::run::is_run_id(&id) {
1627 continue;
1628 }
1629 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1637 continue;
1638 };
1639 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1640 continue;
1641 };
1642 if state.status.done() || !state.active_all_overrun(now) {
1643 continue;
1644 }
1645 let daemon_claims = is_working_on(home, &id, now);
1657 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1658 continue;
1659 }
1660 state.abandon("daemon");
1661 if let Err(e) = state.save_under(home) {
1662 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1663 continue;
1664 }
1665 if let Some(notice) = notices::run_ended(&state) {
1668 notices::raise_in(home, notice);
1669 }
1670 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1678 tracing::warn!("abandon questions for {id}: {e:#}");
1679 }
1680 abandoned.push(id);
1681 }
1682 abandoned
1683}
1684
1685pub async fn serve(opts: Opts) -> Result<()> {
1691 serve_until(opts, Stop::new()).await
1692}
1693
1694pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1711 let signal = {
1712 let stop = stop.clone();
1713 tokio::spawn(async move {
1714 if tokio::signal::ctrl_c().await.is_ok() {
1715 stop.stop();
1716 tracing::info!("shutdown requested; a run in flight will be finished first");
1717 }
1718 })
1719 };
1720
1721 let worktrees_root = opts
1722 .worktrees_root
1723 .clone()
1724 .unwrap_or_else(crate::run::default_worktree_root);
1725 let outcome = drive(
1726 &opts,
1727 &Queue::open(),
1728 &status_path(),
1729 &crate::run::home(),
1730 &worktrees_root,
1731 &stop,
1732 )
1733 .await;
1734
1735 signal.abort();
1736 outcome
1737}
1738
1739async fn drive(
1752 opts: &Opts,
1753 queue: &Queue,
1754 status_file: &Path,
1755 home: &Path,
1756 worktrees_root: &Path,
1757 stop: &Stop,
1758) -> Result<()> {
1759 let status = Arc::new(Mutex::new(Status::new()));
1767 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1768 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1769
1770 let daemon_cfg = prepare(&opts.repo, opts)
1776 .map(|c| c.daemon)
1777 .unwrap_or_default();
1778 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1779
1780 let waiter = tokio::spawn(crate::waiter::run(
1785 crate::waiter::Waiter::new(
1786 crate::ask::Questions::at(home.join("questions")),
1787 home.to_path_buf(),
1788 prepare(&opts.repo, opts).ok(),
1789 ),
1790 stop.clone(),
1791 ));
1792
1793 let deputies = tokio::spawn(crate::deputy::run(
1797 crate::deputy::Deputies::new(
1798 crate::ask::Questions::at(home.join("questions")),
1799 home.to_path_buf(),
1800 prepare(&opts.repo, opts).ok(),
1801 opts.repo.clone(),
1802 daemon_cfg.max_deputies,
1803 {
1804 let stop = stop.clone();
1805 Arc::new(move || stop.parking())
1806 },
1807 ),
1808 stop.clone(),
1809 ));
1810
1811 let fetcher = tokio::spawn(fetch_loop(opts.repo.clone(), opts.clone(), stop.clone()));
1814
1815 let release_watch = tokio::spawn(crate::release_watch::run(
1818 crate::release_watch::Watcher::new(
1819 Box::new(crate::release_watch::GhForge),
1820 home.to_path_buf(),
1821 ),
1822 {
1823 let (repo, opts) = (opts.repo.clone(), opts.clone());
1824 move || release_watch_settings(&repo, &opts)
1825 },
1826 stop.clone(),
1827 ));
1828
1829 tracing::info!(
1830 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1831 queue.root().display(),
1832 opts.poll.as_secs(),
1833 opts.max_attempts,
1834 concurrency,
1835 if daemon_cfg.pause_for_interrupts {
1836 ", interrupts enabled"
1837 } else {
1838 ""
1839 }
1840 );
1841
1842 janitor(&opts.repo, opts, home, worktrees_root).await;
1845 resweep_superseded_attempts(queue, home);
1846
1847 let outcome = poll(
1848 opts,
1849 queue,
1850 &status,
1851 home,
1852 worktrees_root,
1853 stop,
1854 DispatchLimits {
1855 max_concurrent: concurrency,
1856 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1857 },
1858 )
1859 .await;
1860
1861 beat.abort();
1862 waiter.abort();
1863 deputies.abort();
1864 fetcher.abort();
1865 release_watch.abort();
1866 clear_status_at(status_file);
1867 outcome
1868}
1869
1870fn release_watch_settings(repo: &Path, opts: &Opts) -> (Vec<PathBuf>, u64) {
1873 let Ok(cfg) = prepare(repo, opts) else {
1874 return (Vec::new(), 0);
1875 };
1876 let mut paths: Vec<PathBuf> = crate::repos::scan(&cfg.repos.roots)
1877 .into_iter()
1878 .map(|r| r.path)
1879 .collect();
1880 let mut known: Vec<PathBuf> = Queue::open().list().into_iter().map(|t| t.repo).collect();
1883 let runs = crate::run::home().join("runs");
1884 for id in crate::run::list_ids_in(&runs) {
1885 let repo = std::fs::read_to_string(runs.join(&id).join("run.json"))
1886 .ok()
1887 .and_then(|s| serde_json::from_str::<serde_json::Value>(&s).ok())
1888 .and_then(|v| v.get("repo")?.as_str().map(PathBuf::from));
1889 known.extend(repo);
1890 }
1891 for k in known {
1892 let k = k.canonicalize().unwrap_or(k);
1893 if k.join(".git").exists() && !paths.contains(&k) {
1894 paths.push(k);
1895 }
1896 }
1897 let own = repo.canonicalize().unwrap_or_else(|_| repo.to_path_buf());
1898 if !paths.contains(&own) {
1899 paths.push(own);
1900 }
1901 (paths, cfg.daemon.release_stall_minutes)
1902}
1903
1904const FETCH_TIMEOUT: Duration = Duration::from_secs(30);
1906
1907async fn fetch_loop(repo: PathBuf, opts: Opts, stop: Stop) {
1911 while !stop.stopped() {
1912 let (interval, roots) = match prepare(&repo, &opts) {
1913 Ok(c) => (c.repos.fetch_interval, c.repos.roots),
1914 Err(_) => (0, Vec::new()),
1915 };
1916 if interval > 0 && !roots.is_empty() {
1917 let r = crate::clean::fetch_origins(&roots, FETCH_TIMEOUT, || stop.stopped()).await;
1918 tracing::debug!("fetch origins: {r:?}");
1919 }
1920 let wait = if interval > 0 { interval } else { 60 };
1922 let mut slept = 0;
1923 while slept < wait && !stop.stopped() {
1924 tokio::time::sleep(Duration::from_secs(1)).await;
1925 slept += 1;
1926 }
1927 }
1928}
1929
1930async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1936 loop {
1937 tokio::time::sleep(HEARTBEAT).await;
1938 let snapshot = {
1939 let mut guard = lock(&status);
1940 guard.updated_at = Timestamp::now();
1941 guard.clone()
1942 };
1943 if let Err(e) = write_status_to(&path, &snapshot) {
1944 tracing::warn!("could not refresh the daemon status file: {e:#}");
1947 }
1948 }
1949}
1950
1951#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1954enum LandResume {
1955 NotLanding,
1958 StillWaiting,
1963 Ready,
1967}
1968
1969fn land_resume_state(task: &Task) -> LandResume {
1973 let Some(run_id) = task.runs.last() else {
1974 return LandResume::NotLanding;
1975 };
1976 let Ok(state) = RunState::load(run_id) else {
1977 return LandResume::NotLanding;
1978 };
1979 if !state.parked {
1980 return LandResume::NotLanding;
1981 }
1982 let node = if state.status == RunStatus::Landing {
1986 land::APPROVAL_NODE
1987 } else if state.github_text.as_ref().is_some_and(|g| !g.resolved) {
1988 crate::github_text::ASK_NODE
1989 } else {
1990 return LandResume::NotLanding;
1991 };
1992 let store = ask::Questions::open();
1993 let waiting = store
1994 .list()
1995 .into_iter()
1996 .filter(|q| &q.run == run_id && q.node == node)
1997 .max_by(|a, b| a.id.cmp(&b.id));
1998 let Some(mut q) = waiting else {
1999 return LandResume::Ready;
2000 };
2001 if !q.status.open() {
2002 return LandResume::Ready;
2003 }
2004 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
2011 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
2012 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
2013 q.abandon(format!(
2014 "no answer within {}s of asking",
2015 timeout.as_secs().max(1)
2016 ));
2017 if store.put(&mut q).is_ok() {
2020 return LandResume::Ready;
2021 }
2022 }
2023 LandResume::StillWaiting
2024}
2025
2026const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
2034
2035const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
2047
2048struct InFlightGuard<'a> {
2061 status: &'a Arc<Mutex<Status>>,
2062 stop: &'a Stop,
2063 task_id: &'a str,
2064}
2065
2066impl Drop for InFlightGuard<'_> {
2067 fn drop(&mut self) {
2068 lock(self.status).current.retain(|c| c.task != self.task_id);
2069 self.stop.exit();
2070 }
2071}
2072
2073#[derive(Debug, Clone, PartialEq, Eq)]
2095enum Interrupt {
2096 Idle,
2098 Parking {
2110 parked: Vec<String>,
2111 interrupt_task: String,
2112 },
2113 Running {
2121 parked: Vec<String>,
2122 interrupt_task: String,
2123 },
2124 Resuming { parked: Vec<String> },
2131}
2132
2133fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
2154 match state {
2155 Interrupt::Idle => {
2156 if in_flight.len() != 1 {
2168 return Interrupt::Idle;
2169 }
2170 match runnable.iter().find(|t| t.interrupt) {
2171 Some(t) => Interrupt::Parking {
2172 parked: in_flight.to_vec(),
2173 interrupt_task: t.id.clone(),
2174 },
2175 None => Interrupt::Idle,
2176 }
2177 }
2178 Interrupt::Parking {
2179 parked,
2180 interrupt_task,
2181 } => {
2182 if in_flight.iter().any(|id| parked.contains(id)) {
2183 Interrupt::Parking {
2185 parked,
2186 interrupt_task,
2187 }
2188 } else if in_flight.contains(&interrupt_task) {
2189 Interrupt::Running {
2190 parked,
2191 interrupt_task,
2192 }
2193 } else if runnable.iter().any(|t| t.id == interrupt_task) {
2194 Interrupt::Parking {
2198 parked,
2199 interrupt_task,
2200 }
2201 } else {
2202 Interrupt::Resuming { parked }
2207 }
2208 }
2209 Interrupt::Running {
2210 parked,
2211 interrupt_task,
2212 } => {
2213 if in_flight.contains(&interrupt_task) {
2214 Interrupt::Running {
2215 parked,
2216 interrupt_task,
2217 }
2218 } else {
2219 Interrupt::Resuming { parked }
2225 }
2226 }
2227 Interrupt::Resuming { parked } => {
2228 if in_flight.iter().any(|id| parked.contains(id)) {
2229 Interrupt::Idle
2235 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
2236 Interrupt::Resuming { parked }
2237 } else {
2238 Interrupt::Idle
2241 }
2242 }
2243 }
2244}
2245
2246fn advance_interrupt_tick(
2252 enabled: bool,
2253 state: Interrupt,
2254 in_flight: &[String],
2255 runnable: &[Task],
2256) -> Interrupt {
2257 if !enabled {
2258 return Interrupt::Idle;
2259 }
2260 advance_interrupt(state, in_flight, runnable)
2261}
2262
2263fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
2268 match state {
2269 Interrupt::Idle => candidates,
2270 Interrupt::Parking {
2271 parked,
2272 interrupt_task,
2273 } => {
2274 if in_flight.iter().any(|id| parked.contains(id)) {
2275 Vec::new()
2276 } else {
2277 candidates
2278 .into_iter()
2279 .filter(|t| &t.id == interrupt_task)
2280 .collect()
2281 }
2282 }
2283 Interrupt::Running { .. } => Vec::new(),
2284 Interrupt::Resuming { parked } => candidates
2292 .into_iter()
2293 .find(|t| parked.contains(&t.id))
2294 .into_iter()
2295 .collect(),
2296 }
2297}
2298
2299struct DispatchLimits {
2303 max_concurrent: usize,
2306 pause_for_interrupts: bool,
2308}
2309
2310#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2319enum PermitKind {
2320 None,
2325 Urgent,
2332 Ordinary,
2335}
2336
2337fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
2340 if priority {
2341 PermitKind::None
2342 } else if urgent {
2343 PermitKind::Urgent
2344 } else {
2345 PermitKind::Ordinary
2346 }
2347}
2348
2349async fn poll(
2368 opts: &Opts,
2369 queue: &Queue,
2370 status: &Arc<Mutex<Status>>,
2371 home: &Path,
2372 worktrees_root: &Path,
2373 stop: &Stop,
2374 limits: DispatchLimits,
2375) -> Result<()> {
2376 let DispatchLimits {
2377 max_concurrent,
2378 pause_for_interrupts,
2379 } = limits;
2380 let mut attempted: Vec<String> = Vec::new();
2385 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2386 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2394 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2400 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2401 let mut conductor = Conductor::new();
2402 let mut cache_last_checked: Option<Timestamp> = None;
2405 let mut interrupt = Interrupt::Idle;
2407 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2413 std::collections::HashMap::new();
2414
2415 while !stop.stopped() {
2416 lock(status).polls += 1;
2417
2418 while let Some(result) = inflight.try_join_next() {
2423 if let Err(e) = result {
2424 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2425 notices::raise(Notice::error(
2426 "loop:attempt",
2427 "A queued attempt ended abnormally; check the task it was running.",
2428 ));
2429 }
2430 }
2431
2432 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2433 if !swept.is_empty() {
2434 tracing::warn!(
2435 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2436 swept.len(),
2437 swept.join(", ")
2438 );
2439 }
2440 let now = Timestamp::now();
2445
2446 if !stop.busy_now() {
2451 maybe_prune_cache_between_runs(
2452 &opts.repo,
2453 opts,
2454 home,
2455 stop,
2456 &mut cache_last_checked,
2457 now,
2458 )
2459 .await;
2460 }
2461
2462 let stalled = stalled_tasks(queue, home, now);
2463 let stalled_ids: std::collections::BTreeSet<_> =
2464 stalled.iter().map(|task| task.id.clone()).collect();
2465 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2466 if !reclaimed.is_empty() {
2467 tracing::warn!(
2468 "reclaimed {} task(s) left `running` by a daemon that never \
2469 recorded the outcome: {}",
2470 reclaimed.len(),
2471 reclaimed.join(", ")
2472 );
2473 }
2474 let abandoned_runs = reclaim_abandoned_runs(home, now);
2475 if !abandoned_runs.is_empty() {
2476 tracing::warn!(
2477 "failed {} run(s) left behind by a killed process, past every \
2478 active seat's own timeout: {}",
2479 abandoned_runs.len(),
2480 abandoned_runs.join(", ")
2481 );
2482 }
2483
2484 let questions = Questions::at(home.join("questions"));
2489
2490 resolve_blockers(queue, &questions);
2493 apply_choice_actions(queue, &questions, home);
2494 reconcile_task_questions(queue, &questions);
2495 {
2499 let talks = crate::talk::Talks::at(home.join("talks"));
2500 crate::chat_notify::sweep(queue, &talks, home);
2501 crate::chat_notify::start_turns(&talks, stop.parking(), &mut inflight);
2502 }
2503
2504 let finished: Vec<Task> = finished_tasks(queue)
2510 .into_iter()
2511 .filter(|task| !stalled_ids.contains(&task.id))
2512 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2516 .collect();
2517 let queued = queued_tasks(queue);
2518 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2522 && conductor.worth_a_look(queue, &stalled, &finished)
2523 {
2524 match prepare(&opts.repo, opts) {
2525 Ok(cfg) => {
2526 conductor
2527 .maybe_run(
2528 &cfg,
2529 &opts.repo,
2530 queue,
2531 &questions,
2532 home,
2533 &queued,
2534 &stalled,
2535 &finished,
2536 opts.max_attempts,
2537 )
2538 .await;
2539 }
2540 Err(e) => {
2541 tracing::warn!("conductor: no config: {e:#}");
2542 notices::raise(Notice::warn(
2543 "loop:no-config",
2544 "The loop could not read this repository's config, so held tasks are not being triaged.",
2545 ));
2546 }
2547 }
2548 }
2549
2550 let candidates: Vec<Task> = runnable(queue)
2551 .into_iter()
2552 .filter(|t| !opts.once || !attempted.contains(&t.id))
2553 .collect();
2554
2555 let in_flight: Vec<String> = lock(status)
2560 .current
2561 .iter()
2562 .map(|c| c.task.clone())
2563 .collect();
2564 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2565
2566 interrupt =
2567 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2568 if let Interrupt::Parking {
2569 parked,
2570 interrupt_task,
2571 } = &interrupt
2572 {
2573 let reason = format!(
2574 "task {} asked to run first",
2575 crate::run::short_of(interrupt_task)
2576 );
2577 for id in parked {
2578 if let Some(pause) = interrupt_pauses.get(id) {
2579 pause.park_because(reason.clone());
2580 }
2581 }
2582 }
2583 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2584
2585 let cooling_down =
2586 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2587
2588 let mut started_any = false;
2589 for candidate in candidates {
2590 if stop.stopped() {
2591 break;
2592 }
2593
2594 let resume = land_resume_state(&candidate);
2595 if resume == LandResume::StillWaiting {
2596 continue;
2597 }
2598 let priority = resume == LandResume::Ready;
2599
2600 if !priority && cooling_down {
2601 continue;
2602 }
2603 let permit = match permit_kind(priority, candidate.urgent) {
2604 PermitKind::None => None,
2605 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2606 Ok(p) => Some(p),
2607 Err(_) => continue,
2613 },
2614 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2615 Ok(p) => Some(p),
2616 Err(_) => continue,
2620 },
2621 };
2622
2623 let Ok(claim) = queue.claim(&candidate.id) else {
2628 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2629 continue;
2630 };
2631 let mut task = match queue.get(&candidate.id) {
2634 Ok(t) if t.status.runnable() => t,
2635 Ok(_) => continue,
2636 Err(e) => {
2637 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2638 continue;
2639 }
2640 };
2641 let task_id = task.id.clone();
2642 attempted.push(task_id.clone());
2643 lock(status).idle = false;
2644 stop.enter();
2647 started_any = true;
2648
2649 let run_pause = crate::graph::Pause::new();
2653 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2654
2655 let opts = opts.clone();
2656 let queue = queue.clone();
2657 let status = Arc::clone(status);
2658 let stop = stop.clone();
2659 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2660 inflight.spawn(async move {
2661 let _claim = claim;
2665 let _permit = permit;
2666 let _inflight = InFlightGuard {
2668 status: &status,
2669 stop: &stop,
2670 task_id: &task_id,
2671 };
2672 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2673 lock(&status).completed += 1;
2674 let now = Timestamp::now();
2680 if let Some(until) = cooldown_until("a, now) {
2681 let wait = until.as_second() - now.as_second();
2682 *lock("a_cooldown_until) = Some(until);
2683 let hint = quota
2684 .iter()
2685 .find(|q| q.reset.is_some())
2686 .and_then(|q| q.reset.as_deref());
2687 match hint {
2688 Some(h) => tracing::warn!(
2689 "quota hit; waiting {wait}s before taking another ordinary task \
2690 (CLI reported reset: {h})"
2691 ),
2692 None => tracing::warn!(
2693 "quota hit; waiting {wait}s before taking another ordinary task \
2694 (no reset hint reported)"
2695 ),
2696 }
2697 }
2698 });
2699 }
2700
2701 if started_any {
2702 continue;
2703 }
2704
2705 if stop.busy_now() {
2706 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2711 continue;
2712 }
2713
2714 lock(status).idle = true;
2716 if opts.once {
2717 janitor(&opts.repo, opts, home, worktrees_root).await;
2721 resweep_superseded_attempts(queue, home);
2722 triage_held(queue, home, opts).await;
2723 break;
2724 }
2725 stop.idle(opts.poll).await;
2726 if stop.stopped() {
2727 continue;
2728 }
2729 janitor(&opts.repo, opts, home, worktrees_root).await;
2735 resweep_superseded_attempts(queue, home);
2736 triage_held(queue, home, opts).await;
2737 }
2738
2739 while let Some(result) = inflight.join_next().await {
2744 if let Err(e) = result {
2745 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2746 notices::raise(Notice::error(
2747 "loop:attempt",
2748 "A queued attempt ended abnormally; check the task it was running.",
2749 ));
2750 }
2751 }
2752 Ok(())
2753}
2754
2755async fn attempt(
2761 opts: &Opts,
2762 queue: &Queue,
2763 status: &Arc<Mutex<Status>>,
2764 stop: &Stop,
2765 interrupt_pause: crate::graph::Pause,
2766 task: &mut Task,
2767) -> Vec<QuotaLoss> {
2768 let repo = repo_for(task, &opts.repo);
2769 tracing::info!(
2770 "task {} — {} (repo {})",
2771 task.short(),
2772 task.title,
2773 repo.display()
2774 );
2775
2776 let explicit = match &task.overrides {
2780 Some(o) => o.config.as_deref(),
2781 None => opts.config.as_deref(),
2782 };
2783 if let Err(e) = Config::discover_fetched(&repo, explicit).await {
2784 let reason = format!("[config] {e:#}");
2785 task.last_error = Some(reason.clone());
2786 task.hold_machine(Some(reason.clone()));
2787 record(queue, task);
2788 tracing::warn!("holding {}: {reason}", task.short());
2789 return Vec::new();
2790 }
2791 let mut config = match prepare_for(&repo, opts, task) {
2792 Ok(c) => c,
2793 Err(e) => {
2794 task.attempts += 1;
2798 task.fail(format!("config: {e:#}"), opts.max_attempts);
2799 record(queue, task);
2800 return Vec::new();
2801 }
2802 };
2803 apply_solo(&mut config, task);
2804 let p = phrases(&config.graph.language);
2805 let start_failed = p.could_not_start;
2806
2807 if let Some(reason) = disk_gate(&repo, &config) {
2815 task.last_error = Some(reason.clone());
2816 task.hold_machine(Some(reason.clone()));
2817 record(queue, task);
2818 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2819 notices::raise(
2822 Notice::warn(
2823 &format!("disk:{}", repo.display()),
2824 "A task was held for want of free disk space; free some, then release it from the queue.",
2825 )
2826 .link(Link::Task {
2827 id: task.id.clone(),
2828 }),
2829 );
2830 return Vec::new();
2831 }
2832
2833 let unfinished = (!task.fresh_start)
2853 .then(|| unfinished_run(&task.runs, task.short()))
2854 .flatten();
2855 let review_branch = task.review_branch.take();
2861 let branch_exists = match &review_branch {
2862 Some(branch) => crate::git::branch_exists(&repo, branch)
2863 .await
2864 .unwrap_or(false),
2865 None => false,
2866 };
2867 let starter = match (&review_branch, &unfinished, &task.review_of) {
2871 (None, None, Some(branch)) => Starter::Review(branch.clone()),
2872 _ => choose_starter(
2873 review_branch.as_deref(),
2874 branch_exists,
2875 unfinished.as_deref(),
2876 ),
2877 };
2878 let attachments = match task_attachments(queue, task) {
2883 Ok(a) => a,
2884 Err(e) => {
2885 task.attempts += 1;
2886 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2887 record(queue, task);
2888 return Vec::new();
2889 }
2890 };
2891 let started = match &starter {
2892 Starter::Review(branch) => {
2893 tracing::info!(
2894 "task {} reopens `{branch}` as a review-only pass",
2895 task.short()
2896 );
2897 let takeover = crate::handover::Takeover {
2900 earlier: task.earlier_attempts().to_vec(),
2901 home: crate::run::home(),
2902 choice: take_divergence_answer(branch, &config.merge.remote, task),
2903 };
2904 Runner::review_taking_over(
2905 &repo,
2906 branch,
2907 config,
2908 Some(takeover),
2909 crate::run::Origin::queue(&task.id),
2910 )
2911 .await
2912 }
2913 Starter::Resume(id) => {
2914 tracing::info!("resuming run {id} rather than competing again");
2915 Runner::resume(id).map(|mut r| {
2916 if let Some(instruction) =
2917 prepare_instruction(&starter, Some(&r.state.instruction), task)
2918 {
2919 r.state.instruction = instruction;
2920 }
2921 r.state.attachments = attachments.clone();
2922 if let Some(mode) = task
2925 .overrides
2926 .as_ref()
2927 .and_then(|o| o.merge.as_deref())
2928 .and_then(|m| merge_mode(m).ok())
2929 {
2930 r.state.config.merge.mode = mode;
2931 }
2932 r
2933 })
2934 }
2935 Starter::Start => {
2936 if let Some(branch) = &review_branch {
2937 tracing::warn!(
2938 "conductor chose review for task {} but branch `{branch}` no longer \
2939 exists; requeuing as a fresh competition instead",
2940 task.short()
2941 );
2942 }
2943 let instruction = prepare_instruction(&starter, None, task)
2944 .unwrap_or_else(|| task.instruction.clone());
2945 Runner::start_naming(
2946 &repo,
2947 instruction,
2948 &task.title,
2949 config,
2950 crate::run::Origin::queue(&task.id),
2951 )
2952 .await
2953 .map(|mut r| {
2954 r.state.attachments = attachments.clone();
2955 r
2956 })
2957 }
2958 };
2959 let mut runner = match started {
2960 Ok(r) => r,
2961 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2965 let detail = match e.downcast_ref::<crate::handover::Refused>() {
2966 Some(r) => (p.handover_refused)(r),
2967 None => format!("{e:#}"),
2968 };
2969 let reason = format!("{start_failed}{detail}");
2970 task.last_error = Some(reason.clone());
2971 let branch = match &starter {
2974 Starter::Review(branch) => Some(branch.clone()),
2975 _ => None,
2976 };
2977 task.hold_for_handover(branch, format!("{reason}{}", p.handover_hint));
2981 record(queue, task);
2982 tracing::warn!(
2983 "holding {} for a branch it cannot take over: {e:#}",
2984 task.short()
2985 );
2986 return Vec::new();
2987 }
2988 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2993 let d = e
2994 .downcast_ref::<crate::reconcile::Diverged>()
2995 .expect("checked by the guard");
2996 let mut q = ask::Question::new(
2997 task.id.clone(),
2998 crate::reconcile::NODE.to_owned(),
2999 crate::reconcile::SEAT.to_owned(),
3000 d.summary(),
3001 d.detail(),
3002 d.choices(),
3003 );
3004 match Questions::open().put(&mut q) {
3005 Ok(()) => {
3006 task.last_error = Some(format!("{e:#}"));
3007 task.review_branch = Some(d.branch.clone());
3010 task.block(vec![q.id.clone()], Some(d.summary()));
3011 }
3012 Err(put) => {
3013 tracing::warn!("could not file the divergence question: {put:#}");
3014 task.attempts += 1;
3015 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
3016 }
3017 }
3018 record(queue, task);
3019 return Vec::new();
3020 }
3021 Err(e) => {
3022 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
3028 && let Starter::Review(branch) = &starter
3029 {
3030 task.review_branch = Some(branch.clone());
3031 task.last_error = Some(format!("{start_failed}{e:#}"));
3032 task.status = crate::queue::TaskStatus::Failed;
3033 record(queue, task);
3034 return Vec::new();
3035 }
3036 task.attempts += 1;
3037 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
3038 record(queue, task);
3039 return Vec::new();
3040 }
3041 };
3042 runner.state.followup_generation = Some(task.followup.as_ref().map_or(0, |f| f.generation));
3045 if runner.state.origin_chat.is_none() {
3046 runner.state.origin_chat = task.chat_talk().map(str::to_owned);
3047 }
3048 runner.on_pause(stop.pause());
3050 runner.watch_interrupt(interrupt_pause);
3054
3055 let run = runner.state.id.clone();
3058 task.start(run.clone());
3059 record(queue, task);
3060 lock(status).current.push(Current {
3061 task: task.id.clone(),
3062 run,
3063 });
3064
3065 let quota_before = runner.state.quota.clone();
3068 let result = runner.execute().await;
3069 finish_attempt(
3070 opts.max_attempts,
3071 queue,
3072 task,
3073 &runner.state,
3074 "a_before,
3075 result,
3076 )
3077}
3078
3079pub fn finish_attempt(
3084 max_attempts: usize,
3085 queue: &Queue,
3086 task: &mut Task,
3087 state: &RunState,
3088 quota_before: &[QuotaLoss],
3089 result: Result<()>,
3090) -> Vec<QuotaLoss> {
3091 let detail = match result {
3092 Ok(()) => describe(state),
3093 Err(e) => format!("{e:#}"),
3094 };
3095 let fresh = losses_this_attempt(quota_before, &state.quota);
3096 let verdict = Verdict {
3097 status: state.status,
3098 left_pr: state.pr.is_some(),
3101 quota_hit: !fresh.is_empty(),
3107 parked: state.parked,
3111 no_viable_candidates: state.viable().is_empty(),
3114 };
3115 settle_and_diagnose(task, verdict, &detail, max_attempts, state);
3116 if task.status == TaskStatus::Done {
3117 supersede_prior_runs(task, &crate::run::home());
3118 }
3119 record(queue, task);
3120 tracing::info!(
3121 "task {} is {} after run {} ({})",
3122 task.short(),
3123 task.status.as_str(),
3124 state.short(),
3125 label(state.status)
3126 );
3127 fresh
3128}
3129
3130pub fn hold_if_runnable(queue: &Queue, task: &mut Task) {
3134 let waits_on_owner = withheld_text_wait(task).is_some();
3139 if task.status.runnable() && !waits_on_owner {
3140 let why = task.last_error.clone().map_or_else(
3141 || "the run did not finish".to_owned(),
3142 |e| format!("the run did not finish: {e}"),
3143 );
3144 task.hold_manual(Some(format!(
3145 "{why}. It was started by hand, so it is not retried \
3146 automatically; `magi task release` retries it."
3147 )));
3148 record(queue, task);
3149 }
3150}
3151
3152#[must_use]
3156pub fn loop_would_resume(run: &str) -> bool {
3157 unfinished_run(&[run.to_owned()], crate::run::short_of(run)).is_some()
3158}
3159
3160#[must_use]
3170pub fn foreign_loop(
3171 reading: Option<&Reading>,
3172 now: Timestamp,
3173 own_pid: u32,
3174) -> Option<Option<u32>> {
3175 let reading = reading.filter(|r| r.running(now))?;
3176 match reading.pid {
3177 Some(pid) if pid == own_pid => None,
3178 pid => Some(pid),
3179 }
3180}
3181
3182pub async fn run_claimed(opts: &Opts, queue: &Queue, task: &mut Task) {
3197 let status = Arc::new(Mutex::new(Status::new()));
3198 let stop = Stop::new();
3199 attempt(
3200 opts,
3201 queue,
3202 &status,
3203 &stop,
3204 crate::graph::Pause::new(),
3205 task,
3206 )
3207 .await;
3208 let mut resumes = 0;
3217 while resumes < MAX_WITHHELD_RESUMES {
3218 let Some(timeout) = withheld_text_wait(task) else {
3219 break;
3220 };
3221 eprintln!(
3222 "waiting for the owner's answer on a withheld pull request text \
3223 (up to {timeout}s; Ctrl-C leaves the run parked for a later `magi serve`)"
3224 );
3225 let limit = std::time::Instant::now() + Duration::from_secs(timeout + 60);
3228 let ready = loop {
3229 match land_resume_state(task) {
3230 LandResume::StillWaiting if std::time::Instant::now() < limit => {
3231 tokio::time::sleep(WITHHELD_POLL).await;
3232 }
3233 LandResume::Ready => break true,
3234 _ => break false,
3235 }
3236 };
3237 if !ready {
3238 break;
3239 }
3240 resumes += 1;
3241 attempt(
3242 opts,
3243 queue,
3244 &status,
3245 &stop,
3246 crate::graph::Pause::new(),
3247 task,
3248 )
3249 .await;
3250 }
3251 hold_if_runnable(queue, task);
3252}
3253
3254const MAX_WITHHELD_RESUMES: u32 = 3;
3257
3258const WITHHELD_POLL: Duration = Duration::from_secs(15);
3260
3261fn withheld_text_wait(task: &Task) -> Option<u64> {
3264 let id = task.runs.last()?;
3265 let s = RunState::load(id).ok()?;
3266 (s.parked && s.github_text.as_ref().is_some_and(|g| !g.resolved))
3267 .then_some(s.config.graph.answer_timeout)
3268}
3269
3270fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
3280 after
3281 .iter()
3282 .filter(|q| !before.contains(q))
3283 .cloned()
3284 .collect()
3285}
3286
3287fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
3290 if quota.is_empty() {
3291 return None;
3292 }
3293 let with_hint = quota.iter().find(|q| q.reset.is_some());
3294 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
3295 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
3296 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
3297 Some(
3298 now.checked_add(jiff::SignedDuration::from_secs(secs))
3299 .unwrap_or(Timestamp::MAX),
3300 )
3301}
3302
3303fn apply_solo(config: &mut Config, task: &Task) {
3313 if task.solo {
3314 config.graph.implementers = 1;
3315 }
3316}
3317
3318fn prepare_for(repo: &Path, opts: &Opts, task: &Task) -> Result<Config> {
3323 let Some(o) = &task.overrides else {
3324 return prepare(repo, opts);
3325 };
3326 let own = Opts {
3330 config: o.config.clone(),
3331 merge: None,
3332 ..opts.clone()
3333 };
3334 let mut config = prepare(repo, &own)?;
3335 o.apply(&mut config);
3336 Ok(config)
3337}
3338
3339fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
3341 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
3342 if let Some(mode) = &opts.merge {
3343 config.merge.mode = merge_mode(mode)?;
3344 }
3345 Ok(config)
3346}
3347
3348async fn maybe_prune_cache_between_runs(
3379 repo: &Path,
3380 opts: &Opts,
3381 home: &Path,
3382 stop: &Stop,
3383 last_checked: &mut Option<Timestamp>,
3384 now: Timestamp,
3385) {
3386 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
3387 return;
3388 }
3389 *last_checked = Some(now);
3390 let cfg = match prepare(repo, opts) {
3391 Ok(cfg) => cfg,
3392 Err(e) => {
3393 tracing::warn!("cache check: no config: {e:#}");
3394 return;
3395 }
3396 };
3397 match clean::prune_cache_if_over_limit(&cfg, home) {
3398 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
3399 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
3400 pruned.files,
3401 pruned.freed
3402 ),
3403 Ok(_) => {}
3404 Err(e) => {
3405 tracing::warn!("housekeep: prune cache: {e:#}");
3406 notices::raise_in(
3407 home,
3408 Notice::warn(
3409 "housekeep:cache",
3410 "Pruning the shared build cache failed; disk usage may keep growing.",
3411 ),
3412 );
3413 }
3414 }
3415}
3416
3417fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
3421 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
3422}
3423
3424async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
3443 let cfg = match prepare(repo, opts) {
3444 Ok(cfg) => cfg,
3445 Err(e) => {
3446 tracing::warn!("housekeep: no config: {e:#}");
3447 return;
3448 }
3449 };
3450 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
3459 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
3460 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
3465 let mut extra = Vec::new();
3466 if out.unreadable > 0 {
3467 extra.push(format!("{} unreadable", out.unreadable));
3468 }
3469 if out.orphaned_worktrees > 0 {
3470 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
3471 }
3472 let detail = if extra.is_empty() {
3473 String::new()
3474 } else {
3475 format!(" ({})", extra.join(", "))
3476 };
3477 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
3478 }
3479 if out.external_merges_recorded > 0 {
3480 tracing::info!(
3481 "housekeep: recorded {} run(s) as merged externally",
3482 out.external_merges_recorded
3483 );
3484 }
3485 if out.stale_pr_states_repaired > 0 {
3486 tracing::info!(
3487 "housekeep: rewrote {} run record(s) whose pull request had already settled",
3488 out.stale_pr_states_repaired
3489 );
3490 }
3491 if out.cache_files > 0 {
3492 tracing::info!(
3493 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
3494 out.cache_files,
3495 out.cache_freed
3496 );
3497 }
3498 if out.questions_abandoned > 0 {
3499 tracing::info!(
3500 "housekeep: abandoned {} question(s) left open by a finished run",
3501 out.questions_abandoned
3502 );
3503 }
3504}
3505
3506async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3515 let questions = Questions::at(home.join("questions"));
3516 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3517 if report.is_empty() {
3518 return;
3519 }
3520 if !report.quarantined.is_empty() {
3521 tracing::info!(
3522 "triage: held {} blocked task(s) whose blocked-on task or \
3523 question no longer exists: {}",
3524 report.quarantined.len(),
3525 report.quarantined.join(", ")
3526 );
3527 }
3528 if !report.resumed.is_empty() {
3529 tracing::info!(
3530 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3531 report.resumed.len(),
3532 report.resumed.join(", ")
3533 );
3534 }
3535 if !report.asked.is_empty() {
3536 tracing::info!(
3537 "triage: asked about {} held task(s): {}",
3538 report.asked.len(),
3539 report.asked.join(", ")
3540 );
3541 }
3542 if !report.answered.is_empty() {
3543 tracing::info!(
3544 "triage: applied {} operator answer(s): {}",
3545 report.answered.len(),
3546 report.answered.join(", ")
3547 );
3548 }
3549}
3550
3551fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3558 disk_gate_with(repo, config, crate::disk::free_bytes)
3559}
3560
3561fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3565 repo: &Path,
3566 config: &Config,
3567 free_bytes: F,
3568) -> Option<String> {
3569 let min = config.disk.min_free_bytes;
3570 if min == 0 {
3571 return None;
3572 }
3573 match free_bytes(repo) {
3574 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3575 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3576 }
3577}
3578
3579const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3586
3587const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3591
3592fn quota_wait(
3601 reset_at: Option<Timestamp>,
3602 now: Timestamp,
3603 fallback: Duration,
3604 cap: Duration,
3605) -> Duration {
3606 match reset_at {
3607 Some(at) if at > now => {
3608 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3609 Duration::from_secs(secs).min(cap)
3610 }
3611 _ => fallback,
3612 }
3613}
3614
3615fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3629 parse_reset_hint_zoned(text, now)
3630 .or_else(|| parse_reset_hint_dated(text))
3631 .or_else(|| parse_reset_hint_relative(text, recorded))
3632}
3633
3634fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3638 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3639 let mut rest = rest.trim();
3640 if rest.is_empty() {
3641 return None;
3642 }
3643 let mut total: i64 = 0;
3644 let mut matched = false;
3645 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3646 if let Some((digits, tail)) = rest.split_once(unit)
3647 && !digits.is_empty()
3648 && digits.bytes().all(|b| b.is_ascii_digit())
3649 {
3650 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3651 rest = tail;
3652 matched = true;
3653 }
3654 }
3655 if !rest.is_empty() || !matched {
3656 return None;
3657 }
3658 recorded
3659 .checked_add(jiff::SignedDuration::from_secs(total))
3660 .ok()
3661}
3662
3663fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3667 let clock = clock.trim().to_lowercase();
3668 let (digits, pm) = clock
3669 .strip_suffix("am")
3670 .map(|d| (d, false))
3671 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3672 let (h, m) = digits.trim().split_once(':')?;
3673 let mut hour: i8 = h.trim().parse().ok()?;
3674 let minute: i8 = m.trim().parse().ok()?;
3675 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3676 return None;
3677 }
3678 if pm && hour != 12 {
3679 hour += 12;
3680 } else if !pm && hour == 12 {
3681 hour = 0;
3682 }
3683 Some((hour, minute))
3684}
3685
3686fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3691 let open = text.find('(')?;
3692 let close = text.rfind(')')?;
3693 if close <= open {
3694 return None;
3695 }
3696 let zone = text[open + 1..close].trim();
3697 let (hour, minute) = parse_12h_clock(&text[..open])?;
3698 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3699 let candidate = now
3700 .to_zoned(tz)
3701 .with()
3702 .hour(hour)
3703 .minute(minute)
3704 .second(0)
3705 .millisecond(0)
3706 .microsecond(0)
3707 .nanosecond(0)
3708 .build()
3709 .ok()?;
3710 let mut at = candidate.timestamp();
3711 if at <= now {
3712 at += jiff::SignedDuration::from_hours(24);
3713 }
3714 Some(at)
3715}
3716
3717fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3726 let words: Vec<&str> = text.split_whitespace().collect();
3727 if words.len() < 5 {
3728 return None;
3729 }
3730 (0..=words.len() - 5)
3731 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3732}
3733
3734fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3740 if trailing.is_some_and(|next| next.starts_with('(')) {
3741 return None;
3742 }
3743 let month = month_number(window[0])?;
3744 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3745 let day_digits = ["st", "nd", "rd", "th"]
3746 .iter()
3747 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3748 let day: i8 = day_digits.parse().ok()?;
3749 let year_token = window[2];
3750 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3751 return None;
3752 }
3753 let year: i16 = year_token.parse().ok()?;
3754 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3758 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3759 let date = jiff::civil::Date::new(year, month, day).ok()?;
3760 let candidate = date
3761 .at(hour, minute, 0, 0)
3762 .to_zoned(jiff::tz::TimeZone::UTC)
3763 .ok()?;
3764 Some(candidate.timestamp())
3765}
3766
3767fn month_number(name: &str) -> Option<i8> {
3770 const NAMES: [&str; 12] = [
3771 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3772 ];
3773 let lower = name.to_lowercase();
3774 NAMES
3775 .iter()
3776 .position(|n| *n == lower.as_str())
3777 .map(|i| i as i8 + 1)
3778}
3779
3780fn exhausted_review_budget(state: &RunState) -> bool {
3792 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3793}
3794
3795fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3832 unfinished_run_with(runs, short, RunState::load)
3833}
3834
3835fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3838where
3839 F: FnOnce(&str) -> Result<RunState>,
3840{
3841 let id = runs.last()?;
3842 match load(id) {
3843 Ok(s)
3848 if s.status.resumable()
3849 && !s.released()
3850 && !exhausted_review_budget(&s)
3851 && s.liveness(false) != crate::run::Liveness::Live =>
3852 {
3853 Some(id.clone())
3854 }
3855 Ok(_) => None,
3856 Err(e) => {
3857 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3858 None
3859 }
3860 }
3861}
3862
3863fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3871where
3872 F: FnOnce(&str) -> Result<RunState>,
3873{
3874 if !matches!(task.status, TaskStatus::Parked | TaskStatus::Failed)
3876 || task.fresh_start
3877 || task.review_branch.is_some()
3878 {
3879 return false;
3880 }
3881 let Some(id) = task.runs.last() else {
3882 return false;
3883 };
3884 load(id).is_ok_and(|s| {
3885 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3886 })
3887}
3888
3889#[derive(Debug, Clone, PartialEq, Eq)]
3892enum Starter {
3893 Review(String),
3896 Resume(String),
3898 Start,
3900}
3901
3902fn take_divergence_answer(
3906 branch: &str,
3907 remote: &str,
3908 task: &mut Task,
3909) -> Option<crate::reconcile::Choice> {
3910 let summary = crate::reconcile::summary_for(branch, remote);
3911 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3912 (a.question == summary)
3913 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3914 .flatten()
3915 .map(|c| (i, c))
3916 })?;
3917 task.answers.remove(idx);
3918 Some(choice)
3919}
3920
3921fn choose_starter(
3933 review_branch: Option<&str>,
3934 branch_exists: bool,
3935 unfinished: Option<&str>,
3936) -> Starter {
3937 match review_branch {
3938 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3939 Some(_) => Starter::Start,
3940 None => match unfinished {
3941 Some(id) => Starter::Resume(id.to_owned()),
3942 None => Starter::Start,
3943 },
3944 }
3945}
3946
3947fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3950 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3951 return fallback.to_path_buf();
3952 }
3953 task.repo.clone()
3954}
3955
3956const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3960
3961fn answers_block(task: &Task, count: usize) -> String {
3963 let mut s = ANSWERS_HEADER.to_owned();
3964 for a in &task.answers[..count] {
3965 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3966 }
3967 s
3968}
3969
3970fn append_answers(base: &str, task: &Task) -> String {
3973 if task.answers.is_empty() {
3974 return base.to_owned();
3975 }
3976 let mut s = base.to_owned();
3977 s.push_str(&answers_block(task, task.answers.len()));
3978 s
3979}
3980
3981fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3985 for count in (1..=task.answers.len()).rev() {
3986 let block = answers_block(task, count);
3987 if let Some(base) = instruction.strip_suffix(&block) {
3988 return base;
3989 }
3990 }
3991 instruction
3992}
3993
3994fn instruction_for(task: &Task) -> String {
4002 append_answers(&task.instruction, task)
4003}
4004
4005fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
4017 append_answers(strip_answers_block(old_instruction, task), task)
4018}
4019
4020fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
4022 let paths = queue.attachment_paths(task);
4023 for (name, path) in task.attachments.iter().zip(&paths) {
4024 if !path.is_file() {
4025 bail!(
4026 "attachment `{name}` is recorded on the task but {} is missing",
4027 path.display()
4028 );
4029 }
4030 }
4031 Ok(paths)
4032}
4033
4034fn prepare_instruction(
4045 starter: &Starter,
4046 old_instruction: Option<&str>,
4047 task: &Task,
4048) -> Option<String> {
4049 match starter {
4050 Starter::Start => Some(instruction_for(task)),
4051 Starter::Resume(_) => Some(resumed_instruction(
4052 old_instruction.expect("a resumed run always has a prior instruction"),
4053 task,
4054 )),
4055 Starter::Review(_) => None,
4056 }
4057}
4058
4059fn record(queue: &Queue, task: &mut Task) {
4063 if let Err(e) = queue.put(task) {
4064 tracing::error!("could not record task {}: {e:#}", task.short());
4065 notices::raise(
4066 Notice::error(
4067 "loop:record",
4068 "The loop could not save a task's state; check the disk.",
4069 )
4070 .link(Link::Task {
4071 id: task.id.clone(),
4072 }),
4073 );
4074 }
4075}
4076
4077fn runnable(queue: &Queue) -> Vec<Task> {
4083 let mut tasks: Vec<Task> = queue
4084 .list()
4085 .into_iter()
4086 .filter(|t| t.status.runnable())
4087 .collect();
4088 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
4089 tasks
4090}
4091
4092fn describe(state: &RunState) -> String {
4106 let p = phrases(&state.config.graph.language);
4107 let mut detail = if state.status == RunStatus::Stalled {
4108 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
4109 seats.sort_unstable();
4110 seats.dedup();
4111 if seats.is_empty() {
4112 p.quorum_lost.to_owned()
4113 } else {
4114 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
4115 }
4116 } else {
4117 format!("{}{}", p.run_ended, state.status.display_label())
4118 };
4119 if let Some(last) = state.events.last() {
4120 let unanswered = state
4124 .reviews
4125 .last()
4126 .filter(|r| {
4127 state.status == RunStatus::Blocked
4128 && last.node == "review"
4129 && r.incomplete()
4130 && r.blocking == 0
4131 && r.round == state.config.graph.review_rounds
4132 && r.e2e.iter().all(crate::run::CommandOutcome::ok)
4133 })
4134 .map(|r| (r.expected - r.answered, r.round));
4135 match unanswered {
4136 Some((missing, rounds)) if crate::lang::is_japanese(&state.config.graph.language) => {
4137 detail.push_str(&format!(
4138 " ({}: {})",
4139 last.node,
4140 (p.reviewers_never_answered)(missing, rounds)
4141 ));
4142 }
4143 _ => detail.push_str(&format!(" ({}: {})", last.node, last.message)),
4144 }
4145 }
4146 detail.push_str(&format!(" [run {}]", state.id));
4147 detail
4148}
4149
4150const DIAGNOSTIC_MAX: usize = 4_000;
4156
4157const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
4162
4163fn diagnostic(state: &RunState) -> Option<String> {
4177 let mut parts: Vec<String> = Vec::new();
4178
4179 for o in state.gate.iter().filter(|o| !o.ok()) {
4181 parts.push(format!(
4182 "gate `{}` failed ({:?}):\n{}",
4183 o.command,
4184 o.code,
4185 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
4186 ));
4187 }
4188
4189 if let Some(last) = state
4192 .events
4193 .iter()
4194 .rev()
4195 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
4196 {
4197 parts.push(last.message.clone());
4198 }
4199
4200 if state.viable().is_empty() {
4207 for c in &state.candidates {
4208 if let Some(evidence) = &c.verified_noop {
4209 parts.push(format!(
4210 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
4211 c.label
4212 ));
4213 } else if !c.summary.trim().is_empty() {
4214 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
4215 } else if let Some(why) = &c.failed {
4216 parts.push(format!("candidate {}: {why}", c.label));
4217 }
4218 }
4219 }
4220
4221 if parts.is_empty() {
4222 return None;
4223 }
4224 Some(crate::run::tail(
4229 &parts.join("\n\n"),
4230 DIAGNOSTIC_MAX.saturating_sub(100),
4231 ))
4232}
4233
4234fn label(status: RunStatus) -> &'static str {
4242 status.as_str()
4243}
4244
4245pub(crate) fn merge_mode(mode: &str) -> Result<MergeMode> {
4247 match mode {
4248 "none" => Ok(MergeMode::None),
4249 "local" => Ok(MergeMode::Local),
4250 "pr" => Ok(MergeMode::Pr),
4251 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
4252 }
4253}
4254
4255fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
4260 mutex
4261 .lock()
4262 .unwrap_or_else(std::sync::PoisonError::into_inner)
4263}
4264
4265#[cfg(test)]
4266mod tests {
4267 use super::*;
4268 use crate::queue::{Source, TaskStatus};
4269 use crate::run::{Candidate, CommandOutcome};
4270 use pretty_assertions::assert_eq;
4271
4272 fn task() -> Task {
4273 Task::new(
4274 "add retries".to_owned(),
4275 "add retries".to_owned(),
4276 PathBuf::from("/repo"),
4277 Source::Human,
4278 )
4279 }
4280
4281 fn interrupt_task(id: &str) -> Task {
4284 let mut t = task();
4285 t.id = id.to_owned();
4286 t.interrupt = true;
4287 t
4288 }
4289
4290 fn task_with_id(id: &str) -> Task {
4292 let mut t = task();
4293 t.id = id.to_owned();
4294 t
4295 }
4296
4297 fn urgent_task(id: &str) -> Task {
4299 let mut t = task();
4300 t.id = id.to_owned();
4301 t.urgent = true;
4302 t
4303 }
4304
4305 #[test]
4310 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
4311 assert_eq!(permit_kind(true, false), PermitKind::None);
4312 assert_eq!(permit_kind(true, true), PermitKind::None);
4313 }
4314
4315 #[test]
4320 fn permit_kind_separates_urgent_from_ordinary() {
4321 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
4322 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
4323 }
4324
4325 #[test]
4332 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
4333 let cfg = Config::default();
4334 let repo = Path::new("/any/repo/path");
4335
4336 let reason =
4337 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
4338 assert!(reason.contains("1024"), "{reason}");
4339 assert!(
4340 reason.contains(&cfg.disk.min_free_bytes.to_string()),
4341 "{reason}"
4342 );
4343
4344 assert_eq!(
4345 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
4346 None,
4347 "exactly at the floor is open"
4348 );
4349 assert_eq!(
4350 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
4351 None,
4352 "comfortably above the floor is open"
4353 );
4354 }
4355
4356 #[test]
4357 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
4358 let mut cfg = Config::default();
4359 cfg.disk.min_free_bytes = 0;
4360 let repo = Path::new("/any/repo/path");
4361 assert_eq!(
4362 disk_gate_with(repo, &cfg, |_| Ok(0)),
4363 None,
4364 "a zero floor never measures at all"
4365 );
4366 }
4367
4368 #[test]
4369 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
4370 let cfg = Config::default();
4371 let repo = Path::new("/any/repo/path");
4372 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
4373 .expect("a measurement failure must close the gate, not open it");
4374 assert!(reason.contains("could not measure"), "{reason}");
4375 }
4376
4377 #[test]
4378 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
4379 let ordinary = task();
4380 let next = advance_interrupt(
4381 Interrupt::Idle,
4382 std::slice::from_ref(&ordinary.id),
4383 std::slice::from_ref(&ordinary),
4384 );
4385 assert_eq!(next, Interrupt::Idle);
4386 }
4387
4388 #[test]
4389 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
4390 let marked = interrupt_task("marked");
4393 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4394 assert_eq!(next, Interrupt::Idle);
4395 }
4396
4397 #[test]
4398 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
4399 let marked = interrupt_task("marked");
4400 let next = advance_interrupt(
4401 Interrupt::Idle,
4402 &["running".to_owned()],
4403 std::slice::from_ref(&marked),
4404 );
4405 assert_eq!(
4406 next,
4407 Interrupt::Parking {
4408 parked: vec!["running".to_owned()],
4409 interrupt_task: "marked".to_owned(),
4410 }
4411 );
4412 }
4413
4414 #[test]
4423 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
4424 let marked = interrupt_task("marked");
4425
4426 let two = advance_interrupt(
4427 Interrupt::Idle,
4428 &["a".to_owned(), "b".to_owned()],
4429 std::slice::from_ref(&marked),
4430 );
4431 assert_eq!(two, Interrupt::Idle);
4432
4433 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4434 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
4435 }
4436
4437 #[test]
4438 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
4439 let state = Interrupt::Parking {
4440 parked: vec!["running".to_owned()],
4441 interrupt_task: "marked".to_owned(),
4442 };
4443 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
4445 assert_eq!(still_going, state);
4446
4447 let stopped_but_not_yet_dispatched =
4451 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
4452 assert_eq!(stopped_but_not_yet_dispatched, state);
4453
4454 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
4456 assert_eq!(
4457 dispatched,
4458 Interrupt::Running {
4459 parked: vec!["running".to_owned()],
4460 interrupt_task: "marked".to_owned(),
4461 }
4462 );
4463 }
4464
4465 #[test]
4466 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
4467 let state = Interrupt::Running {
4468 parked: vec!["running".to_owned()],
4469 interrupt_task: "marked".to_owned(),
4470 };
4471 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
4472 assert_eq!(still_running, state);
4473
4474 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
4481 assert_eq!(
4482 ended,
4483 Interrupt::Resuming {
4484 parked: vec!["running".to_owned()]
4485 }
4486 );
4487 }
4488
4489 #[test]
4490 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
4491 let state = Interrupt::Resuming {
4492 parked: vec!["running".to_owned()],
4493 };
4494 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
4495 assert_eq!(still_waiting, state);
4496
4497 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
4498 assert_eq!(dispatched, Interrupt::Idle);
4499 }
4500
4501 #[test]
4507 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
4508 {
4509 let state = Interrupt::Parking {
4510 parked: vec!["running".to_owned()],
4511 interrupt_task: "marked".to_owned(),
4512 };
4513 let next = advance_interrupt(state, &[], &[]);
4516 assert_eq!(
4517 next,
4518 Interrupt::Resuming {
4519 parked: vec!["running".to_owned()]
4520 },
4521 "abandoning the interrupt must not abandon the resume it owes"
4522 );
4523 }
4524
4525 #[test]
4528 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
4529 let state = Interrupt::Resuming {
4530 parked: vec!["running".to_owned()],
4531 };
4532 let next = advance_interrupt(state, &[], &[]);
4533 assert_eq!(
4534 next,
4535 Interrupt::Idle,
4536 "nothing is left to wait for; the loop must not stay wedged"
4537 );
4538 }
4539
4540 #[test]
4541 fn disabled_by_config_the_sequence_can_never_leave_idle() {
4542 let marked = interrupt_task("marked");
4543 let next = advance_interrupt_tick(
4544 false,
4545 Interrupt::Idle,
4546 &["running".to_owned()],
4547 std::slice::from_ref(&marked),
4548 );
4549 assert_eq!(
4550 next,
4551 Interrupt::Idle,
4552 "an unmarked, unconfigured daemon must behave exactly as before"
4553 );
4554 }
4555
4556 #[test]
4557 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4558 let state = Interrupt::Parking {
4559 parked: vec!["running".to_owned()],
4560 interrupt_task: "marked".to_owned(),
4561 };
4562 let candidates = vec![interrupt_task("marked"), task()];
4563 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4564 assert!(
4565 allowed.is_empty(),
4566 "nothing may dispatch - not even the interrupt task itself - \
4567 until the parked run has actually stopped"
4568 );
4569 }
4570
4571 #[test]
4585 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4586 for state in [
4587 Interrupt::Parking {
4588 parked: vec!["running".to_owned()],
4589 interrupt_task: "marked".to_owned(),
4590 },
4591 Interrupt::Running {
4592 parked: vec!["running".to_owned()],
4593 interrupt_task: "marked".to_owned(),
4594 },
4595 Interrupt::Resuming {
4596 parked: vec!["running".to_owned()],
4597 },
4598 ] {
4599 let candidates = vec![
4600 interrupt_task("marked"),
4601 urgent_task("hot"),
4602 task_with_id("ordinary"),
4603 ];
4604 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4605 assert!(
4606 !allowed.iter().any(|t| t.id == "hot"),
4607 "an urgent candidate must wait out the same gate as anything \
4608 else while the run it would run alongside has not actually \
4609 left flight, for state {state:?}: {allowed:?}"
4610 );
4611 }
4612 }
4613
4614 #[test]
4620 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4621 let state = Interrupt::Resuming {
4622 parked: vec!["hot".to_owned()],
4623 };
4624 let candidates = vec![urgent_task("hot"), task()];
4625 let allowed = interrupt_gate(&state, &[], candidates);
4626 assert_eq!(
4627 allowed.iter().filter(|t| t.id == "hot").count(),
4628 1,
4629 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4630 );
4631 }
4632
4633 #[test]
4634 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4635 let state = Interrupt::Parking {
4636 parked: vec!["running".to_owned()],
4637 interrupt_task: "marked".to_owned(),
4638 };
4639 let other = task();
4640 let candidates = vec![interrupt_task("marked"), other.clone()];
4641 let allowed = interrupt_gate(&state, &[], candidates);
4642 assert_eq!(allowed.len(), 1);
4643 assert_eq!(allowed[0].id, "marked");
4644 }
4645
4646 #[test]
4647 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4648 let state = Interrupt::Running {
4649 parked: vec!["running".to_owned()],
4650 interrupt_task: "marked".to_owned(),
4651 };
4652 let candidates = vec![task(), task()];
4653 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4654 assert!(allowed.is_empty());
4655 }
4656
4657 #[test]
4664 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4665 let state = Interrupt::Resuming {
4666 parked: vec!["a".to_owned(), "c".to_owned()],
4667 };
4668 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4669 let allowed = interrupt_gate(&state, &[], candidates);
4670 assert_eq!(
4671 allowed.len(),
4672 1,
4673 "at most one candidate may be offered while resuming: {allowed:?}"
4674 );
4675 assert_eq!(allowed[0].id, "a");
4676 }
4677
4678 #[test]
4679 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4680 let state = Interrupt::Resuming {
4681 parked: vec!["a".to_owned()],
4682 };
4683 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4684 assert!(allowed.is_empty());
4685 }
4686
4687 #[test]
4693 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4694 let running = task(); let marked = interrupt_task("marked");
4696
4697 let mut state = Interrupt::Idle;
4698 let in_flight = vec![running.id.clone()];
4700 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4701 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4702 assert!(gated.is_empty(), "still waiting on `running` to park");
4703
4704 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4706 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4707 assert_eq!(
4708 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4709 vec!["marked"],
4710 "only the interrupt task may be offered to the dispatcher now"
4711 );
4712
4713 state = advance_interrupt_tick(
4715 true,
4716 state,
4717 &["marked".to_owned()],
4718 std::slice::from_ref(&running),
4719 );
4720 let gated = interrupt_gate(
4721 &state,
4722 &["marked".to_owned()],
4723 vec![marked.clone(), running.clone()],
4724 );
4725 assert!(
4726 gated.is_empty(),
4727 "the parked run must not be offered back while the interrupt \
4728 task is still running"
4729 );
4730
4731 let other = task_with_id("other");
4735 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4736 assert_eq!(
4737 state,
4738 Interrupt::Resuming {
4739 parked: vec![running.id.clone()]
4740 }
4741 );
4742 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4743 assert_eq!(
4744 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4745 vec![running.id.as_str()],
4746 "exactly the parked run resumes - not the unrelated task, even \
4747 though it was offered first"
4748 );
4749
4750 state = advance_interrupt_tick(
4754 true,
4755 state,
4756 std::slice::from_ref(&running.id),
4757 std::slice::from_ref(&other),
4758 );
4759 assert_eq!(state, Interrupt::Idle);
4760 let gated = interrupt_gate(
4761 &state,
4762 std::slice::from_ref(&running.id),
4763 vec![other.clone()],
4764 );
4765 assert_eq!(
4766 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4767 vec![other.id.as_str()],
4768 "ordinary dispatch is unrestricted again"
4769 );
4770 }
4771
4772 #[test]
4773 fn every_run_status_settles_the_task_it_came_from() {
4774 let table = [
4776 (RunStatus::Merged, TaskStatus::Done, 1),
4777 (RunStatus::Ready, TaskStatus::Done, 1),
4778 (RunStatus::Stalled, TaskStatus::Failed, 0),
4779 (RunStatus::Blocked, TaskStatus::Failed, 1),
4780 (RunStatus::Failed, TaskStatus::Failed, 1),
4781 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4782 (RunStatus::Prep, TaskStatus::Failed, 1),
4783 (RunStatus::Implementing, TaskStatus::Failed, 1),
4784 (RunStatus::Judging, TaskStatus::Failed, 1),
4785 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4786 (RunStatus::Voting, TaskStatus::Failed, 1),
4787 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4788 (RunStatus::Gating, TaskStatus::Failed, 1),
4789 ];
4790 for (run, want, attempts) in table {
4791 let mut t = task();
4792 t.start("20260902-000000-aaaa".to_owned());
4793 settle(
4794 &mut t,
4795 Verdict {
4796 status: run,
4797 left_pr: false,
4798 parked: false,
4799 quota_hit: matches!(run, RunStatus::Stalled),
4800 no_viable_candidates: false,
4801 },
4802 "why",
4803 2,
4804 );
4805 assert_eq!(t.status, want, "task status after {}", label(run));
4806 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4807 }
4808 }
4809
4810 #[test]
4811 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4812 let mut stalled = task();
4813 stalled.start("20260902-000000-aaaa".to_owned());
4814 settle(
4815 &mut stalled,
4816 Verdict {
4817 status: RunStatus::Stalled,
4818 left_pr: false,
4819 parked: false,
4820 quota_hit: true,
4821 no_viable_candidates: false,
4822 },
4823 "quota",
4824 1,
4825 );
4826 assert_eq!(stalled.attempts, 0);
4827 assert!(
4828 stalled.status.runnable(),
4829 "a machine problem must leave the task in line"
4830 );
4831
4832 let mut blocked = task();
4833 blocked.start("20260902-000000-aaaa".to_owned());
4834 settle(
4835 &mut blocked,
4836 Verdict {
4837 status: RunStatus::Blocked,
4838 left_pr: false,
4839 parked: false,
4840 quota_hit: false,
4841 no_viable_candidates: false,
4842 },
4843 "findings open",
4844 1,
4845 );
4846 assert_eq!(blocked.attempts, 1);
4847 assert_eq!(
4848 blocked.status,
4849 TaskStatus::Held,
4850 "the last attempt hands the task to a human"
4851 );
4852 }
4853
4854 #[test]
4855 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4856 let mut delivered = task();
4859 delivered.start("20260903-080619-01c2".to_owned());
4860 settle(
4861 &mut delivered,
4862 Verdict {
4863 status: RunStatus::Blocked,
4864 left_pr: true,
4865 parked: false,
4866 quota_hit: false,
4867 no_viable_candidates: false,
4868 },
4869 "no check status",
4870 4,
4871 );
4872 assert_eq!(
4873 delivered.status,
4874 TaskStatus::Held,
4875 "a pull request waiting on CI or a person is not a retryable failure"
4876 );
4877 assert!(
4878 !delivered.status.runnable(),
4879 "the loop must not pick this task up again"
4880 );
4881 assert_eq!(
4882 delivered.last_error.as_deref(),
4883 Some("no check status"),
4884 "the operator needs to be told what the gate was waiting for"
4885 );
4886
4887 let mut empty_handed = task();
4890 empty_handed.start("20260903-080619-01c2".to_owned());
4891 settle(
4892 &mut empty_handed,
4893 Verdict {
4894 status: RunStatus::Blocked,
4895 left_pr: false,
4896 parked: false,
4897 quota_hit: false,
4898 no_viable_candidates: false,
4899 },
4900 "findings open",
4901 4,
4902 );
4903 assert_eq!(empty_handed.status, TaskStatus::Failed);
4904 assert!(empty_handed.status.runnable());
4905 }
4906
4907 #[test]
4908 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4909 let mut noop = task();
4915 noop.start("20260912-131304-391f".to_owned());
4916 settle(
4917 &mut noop,
4918 Verdict {
4919 status: RunStatus::VerifiedNoop,
4920 left_pr: false,
4921 parked: false,
4922 quota_hit: false,
4923 no_viable_candidates: true,
4924 },
4925 "candidate A: already fixed by b32cfc4, on main",
4926 4,
4927 );
4928 assert_eq!(
4929 noop.status,
4930 TaskStatus::Held,
4931 "an unverified claim is a request for a human, not a failure"
4932 );
4933 assert!(
4934 !noop.status.runnable(),
4935 "the loop must not requeue this on the same unverified claim"
4936 );
4937 assert_eq!(noop.attempts, 1);
4942 }
4943
4944 #[test]
4945 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4946 let mut parked = task();
4951 parked.start("20260903-183634-2d98".to_owned());
4952 settle(
4953 &mut parked,
4954 Verdict {
4955 status: RunStatus::Implementing,
4956 left_pr: false,
4957 quota_hit: false,
4958 parked: true,
4959 no_viable_candidates: false,
4960 },
4961 "parked after `implementing`",
4962 2,
4963 );
4964 assert_eq!(parked.attempts, 0, "a park is refunded");
4965 assert!(
4966 parked.status.runnable(),
4967 "and the task stays in line so the next loop resumes its run"
4968 );
4969 assert_eq!(parked.status, TaskStatus::Parked);
4970 assert_eq!(
4971 parked.park_reason.as_deref(),
4972 Some("parked after `implementing`"),
4973 "the card says where it stopped"
4974 );
4975 assert_eq!(parked.last_error, None, "nothing failed");
4976
4977 let mut broken = task();
4981 broken.start("20260903-183634-2d98".to_owned());
4982 settle(
4983 &mut broken,
4984 Verdict {
4985 status: RunStatus::Implementing,
4986 left_pr: false,
4987 quota_hit: false,
4988 parked: false,
4989 no_viable_candidates: false,
4990 },
4991 "returned mid-flight",
4992 2,
4993 );
4994 assert_eq!(broken.attempts, 1);
4995 }
4996
4997 #[test]
4998 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4999 let mut flaky = task();
5004 flaky.start("20260903-123023-e633".to_owned());
5005 settle(
5006 &mut flaky,
5007 Verdict {
5008 status: RunStatus::Stalled,
5009 left_pr: false,
5010 parked: false,
5011 quota_hit: false,
5012 no_viable_candidates: false,
5013 },
5014 "verdict rests on 1 of 3 judges",
5015 2,
5016 );
5017 assert_eq!(
5018 flaky.attempts, 1,
5019 "flakiness spends an attempt, so `max_attempts` still bounds it"
5020 );
5021 assert!(flaky.status.runnable(), "and it is still worth retrying");
5022
5023 let mut limited = task();
5025 limited.start("20260903-123023-e633".to_owned());
5026 settle(
5027 &mut limited,
5028 Verdict {
5029 status: RunStatus::Stalled,
5030 left_pr: false,
5031 parked: false,
5032 quota_hit: true,
5033 no_viable_candidates: false,
5034 },
5035 "judge-2, judge-3 out of quota",
5036 2,
5037 );
5038 assert_eq!(limited.attempts, 0, "a quota window is refunded");
5039 assert!(limited.status.runnable());
5040
5041 let mut worn = task();
5044 for _ in 0..2 {
5045 worn.release();
5046 }
5047 worn.start("20260903-123023-e633".to_owned());
5048 worn.attempts = 2;
5049 settle(
5050 &mut worn,
5051 Verdict {
5052 status: RunStatus::Stalled,
5053 left_pr: false,
5054 parked: false,
5055 quota_hit: false,
5056 no_viable_candidates: false,
5057 },
5058 "no quorum again",
5059 2,
5060 );
5061 assert_eq!(worn.status, TaskStatus::Held);
5062 assert!(!worn.status.runnable());
5063 }
5064
5065 #[test]
5066 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
5067 let mut wiped_out = task();
5074 wiped_out.start("20260907-025000-a1b2".to_owned());
5075 settle(
5076 &mut wiped_out,
5077 Verdict {
5078 status: RunStatus::Failed,
5079 left_pr: false,
5080 parked: false,
5081 quota_hit: true,
5082 no_viable_candidates: true,
5083 },
5084 "no candidate produced a change; nothing to judge",
5085 2,
5086 );
5087 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
5088 assert!(
5089 wiped_out.status.runnable(),
5090 "a machine problem must leave the task in line"
5091 );
5092
5093 let mut partial_progress = task();
5099 partial_progress.start("20260907-025500-c3d4".to_owned());
5100 settle(
5101 &mut partial_progress,
5102 Verdict {
5103 status: RunStatus::Failed,
5104 left_pr: false,
5105 parked: false,
5106 quota_hit: true,
5107 no_viable_candidates: false,
5108 },
5109 "gate failed on the winning candidate",
5110 2,
5111 );
5112 assert_eq!(
5113 partial_progress.attempts, 1,
5114 "a candidate that actually produced a change spends the attempt \
5115 even though some other seat hit its quota"
5116 );
5117 assert!(partial_progress.status.runnable());
5118 }
5119
5120 #[test]
5121 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
5122 let mut t = task();
5128 t.start("20260907-025000-a1b2".to_owned());
5129 let mut state = run_state(RunStatus::Failed);
5130 state.quota.push(QuotaLoss {
5131 seat: "cand-a".to_owned(),
5132 node: "implement".to_owned(),
5133 at: Timestamp::now(),
5134 reset: None,
5135 });
5136 assert!(
5137 state.viable().is_empty(),
5138 "no candidate was added, so nothing is viable"
5139 );
5140 reclaim(&mut t, Some(state), 2, "en");
5141 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
5142 assert!(t.status.runnable());
5143 }
5144
5145 #[test]
5146 fn a_held_task_is_never_offered_to_the_loop() {
5147 let dir = tempfile::tempdir().unwrap();
5148 let queue = Queue::at(dir.path().to_path_buf());
5149 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
5150 let mut t = task();
5151 t.id = format!("2026090{n}-000000-000{n}");
5152 t.priority = priority;
5153 queue.put(&mut t).unwrap();
5154 }
5155 let mut held = task();
5156 held.id = "20260909-000000-9999".to_owned();
5157 held.priority = 99;
5158 held.hold_machine(None);
5159 queue.put(&mut held).unwrap();
5160
5161 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
5162 assert_eq!(order.len(), 3);
5163 assert!(!order.contains(&held.id));
5164 assert_eq!(
5165 order.first().cloned(),
5166 queue.next_runnable().map(|t| t.id),
5167 "the loop's first candidate is exactly what the queue offers"
5168 );
5169 assert_eq!(
5170 order,
5171 vec![
5172 "20260902-000000-0002".to_owned(),
5173 "20260903-000000-0003".to_owned(),
5174 "20260901-000000-0001".to_owned(),
5175 ],
5176 "priority first, then oldest, so nothing starves"
5177 );
5178 }
5179
5180 #[test]
5181 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
5182 let dir = tempfile::tempdir().unwrap();
5183 let queue = Queue::at(dir.path().to_path_buf());
5184 let mut old = task();
5185 old.id = "20260101-000000-old0".to_owned();
5186 queue.put(&mut old).unwrap();
5187 let mut fresh = task();
5188 fresh.id = "20260101-000000-new0".to_owned();
5189 queue.put(&mut fresh).unwrap();
5190
5191 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
5195 std::thread::sleep(Duration::from_millis(60));
5196 let live = queue.claim(&fresh.id).unwrap();
5197
5198 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
5199 assert_eq!(swept, vec![old.id.clone()]);
5200 assert!(
5201 queue.claim(&old.id).is_ok(),
5202 "an unparseable lock older than the threshold is swept"
5203 );
5204 assert!(
5205 queue.claim(&fresh.id).is_err(),
5206 "a live pid protects its lock regardless of age"
5207 );
5208 drop(live);
5209 }
5210
5211 #[test]
5212 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
5213 let dir = tempfile::tempdir().unwrap();
5223 let queue = Queue::at(dir.path().to_path_buf());
5224 let mut t = task();
5225 t.id = "20260101-000000-live".to_owned();
5226 queue.put(&mut t).unwrap();
5227
5228 let claim = queue.claim(&t.id).unwrap();
5229 std::thread::sleep(Duration::from_millis(60));
5230
5231 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
5232 assert!(
5233 swept.is_empty(),
5234 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
5235 );
5236 assert!(
5237 queue.claim(&t.id).is_err(),
5238 "the lock still protects its task"
5239 );
5240 drop(claim);
5241 }
5242
5243 fn injected_dead_pid() -> u32 {
5246 std::process::id().checked_add(1).unwrap_or(1)
5247 }
5248
5249 #[test]
5250 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
5251 let dir = tempfile::tempdir().unwrap();
5252 let queue = Queue::at(dir.path().to_path_buf());
5253 let mut t = task();
5254 t.id = "20260101-000000-dead".to_owned();
5255 queue.put(&mut t).unwrap();
5256 let dead_pid = injected_dead_pid();
5257
5258 std::fs::write(
5263 dir.path().join(format!("{}.lock", t.id)),
5264 dead_pid.to_string(),
5265 )
5266 .unwrap();
5267
5268 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5269 pid != dead_pid
5270 });
5271 assert_eq!(
5272 swept,
5273 vec![t.id.clone()],
5274 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
5275 );
5276 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
5277 }
5278
5279 #[test]
5280 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
5281 let dir = tempfile::tempdir().unwrap();
5282 let queue = Queue::at(dir.path().to_path_buf());
5283 let mut t = task();
5284 t.id = "20260101-000000-late".to_owned();
5285 queue.put(&mut t).unwrap();
5286 let dead_pid = injected_dead_pid();
5287
5288 assert!(
5291 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
5292 "nothing has claimed the task yet"
5293 );
5294
5295 std::fs::write(
5298 dir.path().join(format!("{}.lock", t.id)),
5299 dead_pid.to_string(),
5300 )
5301 .unwrap();
5302
5303 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5307 pid != dead_pid
5308 });
5309 assert_eq!(swept, vec![t.id.clone()]);
5310 }
5311
5312 #[test]
5313 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
5314 crate::run::pin_test_home();
5319 let dir = tempfile::tempdir().unwrap();
5320 let queue = Queue::at(dir.path().to_path_buf());
5321 let mut t = task();
5322 t.id = "20260101-000000-crsh".to_owned();
5323 t.status = TaskStatus::Running;
5324 t.attempts = 1;
5325 t.runs.push("20260904-000000-4043".to_owned());
5329 queue.put(&mut t).unwrap();
5330 let dead_pid = injected_dead_pid();
5331
5332 std::fs::write(
5335 dir.path().join(format!("{}.lock", t.id)),
5336 dead_pid.to_string(),
5337 )
5338 .unwrap();
5339
5340 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
5346 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
5347
5348 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5349 pid != dead_pid
5350 });
5351 assert_eq!(swept, vec![t.id.clone()]);
5352
5353 let reclaimed = reclaim_orphaned_running(&queue, 2);
5354 assert_eq!(reclaimed, vec![t.id.clone()]);
5355 let after = queue.get(&t.id).unwrap();
5356 assert_eq!(
5357 after.status,
5358 TaskStatus::Held,
5359 "no run.json to recover from, so a human is asked"
5360 );
5361 assert_eq!(
5362 after.runs,
5363 vec!["20260904-000000-4043".to_owned()],
5364 "the crashed run's id is kept as evidence, not discarded"
5365 );
5366 }
5367
5368 #[test]
5369 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
5370 let dir = tempfile::tempdir().unwrap();
5371 let queue = Queue::at(dir.path().to_path_buf());
5372 let mut t = task();
5373 t.id = "20260101-000000-unknown".to_owned();
5374 queue.put(&mut t).unwrap();
5375 let dead_pid = injected_dead_pid();
5376 std::fs::write(
5377 dir.path().join(format!("{}.lock", t.id)),
5378 dead_pid.to_string(),
5379 )
5380 .unwrap();
5381
5382 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
5383 assert!(swept.is_empty(), "an unknown pid must keep its lock");
5384 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
5385 }
5386
5387 fn run_state_in(status: RunStatus, language: &str) -> RunState {
5388 let mut s = run_state(status);
5389 s.config.graph.language = language.to_owned();
5390 s
5391 }
5392
5393 fn unstarted_verdict(status: RunStatus) -> Verdict {
5394 Verdict {
5395 status,
5396 left_pr: false,
5397 quota_hit: false,
5398 parked: false,
5399 no_viable_candidates: false,
5400 }
5401 }
5402
5403 #[test]
5404 fn an_already_in_base_run_finishes_the_task_without_spending_an_attempt() {
5405 let mut t = task();
5406 t.attempts = 1;
5407 settle_in(
5408 &mut t,
5409 unstarted_verdict(RunStatus::AlreadyInBase),
5410 "already in main",
5411 1,
5412 phrases("en"),
5413 );
5414 assert_eq!(t.status, TaskStatus::Done);
5415 assert_eq!(t.attempts, 0);
5416 }
5417
5418 #[test]
5419 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
5420 let reason = |language: &str| {
5421 let mut t = task();
5422 settle_in(
5423 &mut t,
5424 unstarted_verdict(RunStatus::Judging),
5425 "boom",
5426 1,
5427 phrases(language),
5428 );
5429 t.last_error.or(t.hold_reason).unwrap_or_default()
5430 };
5431 assert!(
5432 reason("en").starts_with("the graph stopped at `"),
5433 "{}",
5434 reason("en")
5435 );
5436 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
5437 assert!(reason("日本語").contains("boom"));
5438 assert_eq!(reason("fr"), reason("en"));
5439 }
5440
5441 #[test]
5442 fn describe_follows_the_run_language_and_keeps_the_run_id() {
5443 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
5444 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
5445 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
5446 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
5447 assert!(ja.contains("[run "), "{ja}");
5448 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
5449 assert!(ended.starts_with("run 終了: "), "{ended}");
5450 let mut de = run_state_in(RunStatus::Failed, "de");
5451 let mut en = run_state_in(RunStatus::Failed, "en");
5452 de.id = "same".to_owned();
5453 en.id = "same".to_owned();
5454 assert_eq!(describe(&de), describe(&en));
5455 }
5456
5457 #[test]
5458 fn handover_refusals_follow_the_language_and_keep_the_detail_apart() {
5459 use crate::handover::Refused;
5460 let cases = [
5461 Refused::Foreign {
5462 branch: "b".into(),
5463 path: "/w/x".into(),
5464 why: "made by hand".into(),
5465 },
5466 Refused::Unsafe {
5467 branch: "b".into(),
5468 path: "/w/x".into(),
5469 why: "its worktree has uncommitted changes (a.rs)".into(),
5470 },
5471 Refused::ReleaseFailed {
5472 branch: "b".into(),
5473 path: "/w/x".into(),
5474 run: "ab12".into(),
5475 },
5476 ];
5477 for r in &cases {
5478 let en = (phrases("en").handover_refused)(r);
5479 assert_eq!(en, r.to_string());
5480 let ja = (phrases("ja").handover_refused)(r);
5481 assert!(
5482 !ja.contains("is checked out") && !ja.contains("try again"),
5483 "{ja}"
5484 );
5485 assert!(ja.contains("`b`") && ja.contains("/w/x"), "{ja}");
5486 if let Refused::Foreign { why, .. } | Refused::Unsafe { why, .. } = r {
5487 assert!(ja.contains(&format!("(詳細: {why})")), "{ja}");
5488 }
5489 }
5490 assert!(phrases("en").handover_hint.contains("release the task"));
5491 assert!(phrases("ja").handover_hint.contains("解放"));
5492 }
5493
5494 #[test]
5495 fn describe_translates_the_unanswered_reviewer_stop_only_in_ja() {
5496 let build = |lang: &str| {
5497 let mut s = run_state_in(RunStatus::Blocked, lang);
5498 s.id = "same".to_owned();
5499 s.config.graph.review_rounds = 3;
5500 let mut r = review_round(3);
5501 r.expected = 3;
5502 r.answered = 1;
5503 s.reviews.push(r);
5504 s.event(
5505 "review",
5506 "2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean",
5507 );
5508 s
5509 };
5510 let en = describe(&build("en"));
5511 assert!(en.contains("(review: 2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean)"), "{en}");
5512 let ja = describe(&build("ja"));
5513 assert!(ja.contains("2 席のレビュアーが 3 ラウンド"), "{ja}");
5514 assert!(!ja.contains("never answered"), "{ja}");
5515 assert!(ja.contains("[run "), "{ja}");
5516
5517 let mut failed = build("ja");
5520 failed.reviews[0].e2e.push(crate::run::CommandOutcome {
5521 command: "cargo test".to_owned(),
5522 code: Some(1),
5523 output_tail: String::new(),
5524 duration_ms: 0,
5525 resource_blocked: false,
5526 });
5527 failed.event("review", "stopped; e2e failed: cargo test");
5528 let ja = describe(&failed);
5529 assert!(ja.contains("stopped; e2e failed: cargo test"), "{ja}");
5530 assert!(!ja.contains("席のレビュアー"), "{ja}");
5531 }
5532
5533 #[test]
5534 fn refusals_and_recovery_prose_follow_the_language() {
5535 let t = held_task_with("r1");
5536 let q = action_question(
5537 "r1",
5538 ask::ChoiceAction::Resume {
5539 run: "r1".to_owned(),
5540 },
5541 );
5542 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
5543 ActionDecision::Refuse(s) => s,
5544 other => panic!("{other:?}"),
5545 };
5546 assert!(refuse(phrases("en")).contains("could not be read: gone"));
5547 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
5548 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
5549
5550 let mut held = task();
5551 reclaim(&mut held, None, 2, "ja");
5552 assert!(held.hold_reason.unwrap().contains("保留にしました"));
5553 let mut held = task();
5554 reclaim(&mut held, None, 2, "xx");
5555 assert!(held.hold_reason.unwrap().contains("held for a human"));
5556 }
5557
5558 fn run_state(status: RunStatus) -> RunState {
5559 let mut state = RunState::new(
5560 PathBuf::from("/repo"),
5561 "main".to_owned(),
5562 "abc1234def".to_owned(),
5563 "add retries".to_owned(),
5564 Config::default(),
5565 );
5566 state.status = status;
5567 state
5568 }
5569
5570 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
5571 Candidate {
5572 index: 0,
5573 label,
5574 agent: "claude".to_owned(),
5575 branch: format!("magi/x/{label}"),
5576 worktree: PathBuf::from("/repo"),
5577 summary: summary.to_owned(),
5578 stat: String::new(),
5579 files: 0,
5580 commits: usize::from(!empty),
5581 empty,
5582 failed: failed.map(str::to_owned),
5583 verified_noop: None,
5584 duration_ms: 0,
5585 folded: false,
5586 }
5587 }
5588
5589 #[test]
5590 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
5591 let mut state = run_state(RunStatus::Blocked);
5592 state.gate = vec![
5593 CommandOutcome {
5594 command: "cargo make check".to_owned(),
5595 code: Some(0),
5596 output_tail: "ok".to_owned(),
5597 duration_ms: 0,
5598 resource_blocked: false,
5599 },
5600 CommandOutcome {
5601 command: "cargo test".to_owned(),
5602 code: Some(101),
5603 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
5604 duration_ms: 0,
5605 resource_blocked: false,
5606 },
5607 ];
5608 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
5609 assert!(d.contains("cargo test"), "{d}");
5610 assert!(
5611 !d.contains("cargo make check"),
5612 "a passing check is not a diagnostic: {d}"
5613 );
5614 assert!(d.contains("assertion failed"), "{d}");
5615 }
5616
5617 #[test]
5618 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
5619 let mut state = run_state(RunStatus::Blocked);
5620 state.event(
5621 "land",
5622 "stopped: the fixer produced no commit while 2 check(s) were failing \
5623 (build, lint); stopping instead of looping on an unchanged tree",
5624 );
5625 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
5626 assert!(d.contains("build"), "{d}");
5627 assert!(d.contains("lint"), "{d}");
5628 assert!(d.contains("fixer produced no commit"), "{d}");
5629 }
5630
5631 #[test]
5632 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
5633 let state = run_state(RunStatus::VerifiedNoop);
5638 let d = describe(&state);
5639 assert!(
5640 d.contains("agent-verified no-op"),
5641 "expected the display label, not the wire spelling: {d}"
5642 );
5643 assert!(!d.contains("verified_noop"), "{d}");
5644 }
5645
5646 #[test]
5647 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5648 let mut state = run_state(RunStatus::Failed);
5654 state.candidates = vec![candidate(
5655 'A',
5656 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5657 true,
5658 None,
5659 )];
5660 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5661 assert!(d.contains("candidate A"), "{d}");
5662 assert!(d.contains("tagged v1.2.3"), "{d}");
5663 }
5664
5665 #[test]
5666 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5667 let mut state = run_state(RunStatus::Failed);
5668 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5669 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5670 assert!(d.contains("candidate A"), "{d}");
5671 assert!(d.contains("agent timed out"), "{d}");
5672 }
5673
5674 #[test]
5675 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5676 let mut state = run_state(RunStatus::Failed);
5679 state.candidates = vec![candidate('A', "did the work", false, None)];
5680 assert!(diagnostic(&state).is_none());
5681 }
5682
5683 #[test]
5684 fn diagnostic_is_bounded_however_much_a_run_printed() {
5685 let mut state = run_state(RunStatus::Blocked);
5686 state.gate = vec![
5687 CommandOutcome {
5688 command: "cargo test".to_owned(),
5689 code: Some(101),
5690 output_tail: "x".repeat(50_000),
5691 duration_ms: 0,
5692 resource_blocked: false,
5693 },
5694 CommandOutcome {
5695 command: "cargo clippy".to_owned(),
5696 code: Some(1),
5697 output_tail: "y".repeat(50_000),
5698 duration_ms: 0,
5699 resource_blocked: false,
5700 },
5701 ];
5702 state.candidates = vec![
5703 candidate('A', &"z".repeat(50_000), true, None),
5704 candidate('B', &"w".repeat(50_000), true, None),
5705 ];
5706 let d = diagnostic(&state).expect("plenty here to diagnose");
5707 assert!(
5708 d.len() <= DIAGNOSTIC_MAX,
5709 "diagnostic grew to {} bytes, unbounded",
5710 d.len()
5711 );
5712 }
5713
5714 #[test]
5715 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5716 let mut state = run_state(RunStatus::Blocked);
5717 state.gate = vec![CommandOutcome {
5718 command: "cargo test".to_owned(),
5719 code: Some(101),
5720 output_tail: "assertion failed".to_owned(),
5721 duration_ms: 0,
5722 resource_blocked: false,
5723 }];
5724 let verdict = Verdict {
5725 status: RunStatus::Blocked,
5726 left_pr: false,
5727 quota_hit: false,
5728 parked: false,
5729 no_viable_candidates: false,
5730 };
5731
5732 let mut t = task();
5735 t.start("run-1".to_owned());
5736 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5737 assert_eq!(t.status, TaskStatus::Failed);
5738 assert!(t.diagnostic.is_none());
5739
5740 t.start("run-2".to_owned());
5743 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5744 assert_eq!(t.status, TaskStatus::Held);
5745 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5746 assert!(d.contains("cargo test"), "{d}");
5747 }
5748
5749 #[test]
5750 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5751 crate::run::pin_test_home();
5756 let home = crate::run::home();
5757 let state = run_state(RunStatus::VerifiedNoop);
5758 let mut q = ask::Question::new(
5759 state.id.clone(),
5760 "implement".to_owned(),
5761 "impl-A".to_owned(),
5762 "is this really a no-op?".to_owned(),
5763 String::new(),
5764 Vec::new(),
5765 );
5766 Questions::at(home.join("questions")).put(&mut q).unwrap();
5767
5768 let verdict = Verdict {
5769 status: RunStatus::VerifiedNoop,
5770 left_pr: false,
5771 quota_hit: false,
5772 parked: false,
5773 no_viable_candidates: false,
5774 };
5775 let mut t = task();
5776 t.start(state.id.clone());
5777 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5778
5779 assert_eq!(t.status, TaskStatus::Held);
5780 let reason = t.hold_reason.expect("a held task must record why");
5781 assert!(
5782 reason.starts_with("run ended agent-verified no-op"),
5783 "the original settle reason must survive unchanged: {reason}"
5784 );
5785 assert!(
5786 reason.contains(q.short()),
5787 "the open question's id must be named so the notice is actionable: {reason}"
5788 );
5789 }
5790
5791 #[test]
5792 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5793 crate::run::pin_test_home();
5794 let state = run_state(RunStatus::VerifiedNoop);
5795
5796 let verdict = Verdict {
5797 status: RunStatus::VerifiedNoop,
5798 left_pr: false,
5799 quota_hit: false,
5800 parked: false,
5801 no_viable_candidates: false,
5802 };
5803 let mut t = task();
5804 t.start(state.id.clone());
5805 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5806
5807 assert_eq!(t.status, TaskStatus::Held);
5808 assert_eq!(
5809 t.hold_reason.as_deref(),
5810 Some("run ended agent-verified no-op"),
5811 "nothing to append when the question was already answered or never asked"
5812 );
5813 }
5814
5815 #[test]
5816 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5817 crate::run::pin_test_home();
5818 let mut first = run_state(RunStatus::Blocked);
5819 first.id = "20260101-000000-sup1".to_owned();
5820 first.save().unwrap();
5821 let mut second = run_state(RunStatus::Merged);
5822 second.id = "20260101-000000-sup2".to_owned();
5823 second.save().unwrap();
5824
5825 let mut t = task();
5826 t.runs = vec![first.id.clone(), second.id.clone()];
5827 t.status = TaskStatus::Done;
5828
5829 supersede_prior_runs(&t, &crate::run::home());
5830
5831 assert_eq!(
5832 RunState::load(&first.id).unwrap().status,
5833 RunStatus::Superseded,
5834 "the first attempt's Blocked no longer needs anyone's attention"
5835 );
5836 assert_eq!(
5837 RunState::load(&second.id).unwrap().status,
5838 RunStatus::Merged,
5839 "the run that actually succeeded is left exactly as it was"
5840 );
5841 }
5842
5843 #[test]
5844 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5845 crate::run::pin_test_home();
5851 let mut first = run_state(RunStatus::Blocked);
5852 first.id = "20260101-000000-sup9".to_owned();
5853 first.driver_pid = Some(std::process::id());
5856 first.driver_started_at = Some(
5857 crate::proc::process_started_at(std::process::id())
5858 .expect("this test process's own start time must be queryable"),
5859 );
5860 first.save().unwrap();
5861 let mut second = run_state(RunStatus::Merged);
5862 second.id = "20260101-000000-supa".to_owned();
5863 second.save().unwrap();
5864
5865 let mut t = task();
5866 t.runs = vec![first.id.clone(), second.id.clone()];
5867 t.status = TaskStatus::Done;
5868
5869 supersede_prior_runs(&t, &crate::run::home());
5870
5871 assert_eq!(
5872 RunState::load(&first.id).unwrap().status,
5873 RunStatus::Blocked,
5874 "a live driver_pid means something is still actually working this run, \
5875 even though no daemon claims it - rewriting under it would just be \
5876 undone the next time that process saves"
5877 );
5878 }
5879
5880 #[test]
5881 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5882 let dir = tempfile::tempdir().unwrap();
5888 let home = dir.path().to_path_buf();
5889 let queue = Queue::at(dir.path().join("queue"));
5890
5891 let mut first = run_state(RunStatus::Blocked);
5892 first.id = "20260101-000000-supd".to_owned();
5893 first.driver_pid = Some(std::process::id());
5894 first.driver_started_at = Some(
5895 crate::proc::process_started_at(std::process::id())
5896 .expect("this test process's own start time must be queryable"),
5897 );
5898 first.save_under(&home).unwrap();
5899 let mut second = run_state(RunStatus::Merged);
5900 second.id = "20260101-000000-supe".to_owned();
5901 second.save_under(&home).unwrap();
5902
5903 let mut t = task();
5904 t.runs = vec![first.id.clone(), second.id.clone()];
5905 t.status = TaskStatus::Done;
5906 queue.put(&mut t).unwrap();
5907
5908 resweep_superseded_attempts(&queue, &home);
5909 assert_eq!(
5910 RunState::load_under(&first.id, &home).unwrap().status,
5911 RunStatus::Blocked,
5912 "still live on the first pass, so still untouched"
5913 );
5914
5915 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5921 stale.driver_started_at = Some("1".to_owned());
5922 stale.save_under(&home).unwrap();
5923
5924 resweep_superseded_attempts(&queue, &home);
5925 assert_eq!(
5926 RunState::load_under(&first.id, &home).unwrap().status,
5927 RunStatus::Superseded,
5928 "the second pass catches up what the first one correctly skipped"
5929 );
5930 }
5931
5932 #[test]
5933 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5934 crate::run::pin_test_home();
5935 let mut first = run_state(RunStatus::Blocked);
5936 first.id = "20260101-000000-sup3".to_owned();
5937 first.save().unwrap();
5938 let mut second = run_state(RunStatus::Blocked);
5939 second.id = "20260101-000000-sup4".to_owned();
5940 second.save().unwrap();
5941
5942 let mut t = task();
5943 t.runs = vec![first.id.clone(), second.id.clone()];
5944 t.status = TaskStatus::Failed;
5948
5949 supersede_prior_runs(&t, &crate::run::home());
5950
5951 assert_eq!(
5952 RunState::load(&first.id).unwrap().status,
5953 RunStatus::Blocked
5954 );
5955 assert_eq!(
5956 RunState::load(&second.id).unwrap().status,
5957 RunStatus::Blocked
5958 );
5959 }
5960
5961 #[test]
5962 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5963 crate::run::pin_test_home();
5967 let mut first = run_state(RunStatus::Blocked);
5968 first.id = "20260101-000000-sup5".to_owned();
5969 first.save().unwrap();
5970
5971 let mut t = task();
5972 t.runs = vec![first.id.clone()];
5973 t.status = TaskStatus::Done;
5974
5975 supersede_prior_runs(&t, &crate::run::home());
5976
5977 assert_eq!(
5978 RunState::load(&first.id).unwrap().status,
5979 RunStatus::Blocked,
5980 "a single-attempt task has no earlier run to supersede"
5981 );
5982 }
5983
5984 #[test]
5985 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5986 crate::run::pin_test_home();
5993 let mut first = run_state(RunStatus::Blocked);
5994 first.id = "20260101-000000-supb".to_owned();
5995 first.save().unwrap();
5996 let mut second = run_state(RunStatus::Failed);
5997 second.id = "20260101-000000-supc".to_owned();
5998 second.save().unwrap();
5999
6000 let mut t = task();
6001 t.runs = vec![first.id.clone(), second.id.clone()];
6002 t.status = TaskStatus::Done;
6003
6004 supersede_prior_runs(&t, &crate::run::home());
6005
6006 assert_eq!(
6007 RunState::load(&first.id).unwrap().status,
6008 RunStatus::Blocked,
6009 "the task's last attempt never landed, so there is nothing here \
6010 actually superseding it"
6011 );
6012 }
6013
6014 #[test]
6015 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
6016 crate::run::pin_test_home();
6020 let mut failed = run_state(RunStatus::Failed);
6021 failed.id = "20260101-000000-sup6".to_owned();
6022 failed.save().unwrap();
6023 let mut noop = run_state(RunStatus::VerifiedNoop);
6024 noop.id = "20260101-000000-sup7".to_owned();
6025 noop.save().unwrap();
6026 let mut winner = run_state(RunStatus::Ready);
6027 winner.id = "20260101-000000-sup8".to_owned();
6028 winner.save().unwrap();
6029
6030 let mut t = task();
6031 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
6032 t.status = TaskStatus::Done;
6033
6034 supersede_prior_runs(&t, &crate::run::home());
6035
6036 assert_eq!(
6037 RunState::load(&failed.id).unwrap().status,
6038 RunStatus::Failed
6039 );
6040 assert_eq!(
6041 RunState::load(&noop.id).unwrap().status,
6042 RunStatus::VerifiedNoop
6043 );
6044 }
6045
6046 fn approval_question(run: &str) -> ask::Question {
6047 ask::Question::new(
6048 run.to_owned(),
6049 land::APPROVAL_NODE.to_owned(),
6050 "land".to_owned(),
6051 "merge?".to_owned(),
6052 String::new(),
6053 vec!["merge".to_owned(), "hold".to_owned()],
6054 )
6055 }
6056
6057 #[test]
6058 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
6059 crate::run::pin_test_home();
6060 let mut state = run_state(RunStatus::Landing);
6061 state.id = "20260101-000000-fre1".to_owned();
6062 state.parked = true;
6063 state.save().unwrap();
6064 ask::Questions::open()
6065 .put(&mut approval_question(&state.id))
6066 .unwrap();
6067
6068 let mut t = task();
6069 t.runs.push(state.id.clone());
6070 assert_eq!(
6071 land_resume_state(&t),
6072 LandResume::StillWaiting,
6073 "nobody has answered and the timeout has not passed"
6074 );
6075 }
6076
6077 #[test]
6078 fn land_resume_state_waits_on_a_withheld_text_question_then_resumes() {
6079 crate::run::pin_test_home();
6080 let mut state = run_state(RunStatus::Gating);
6081 state.id = "20260101-000000-gt01".to_owned();
6082 state.parked = true;
6083 state.config.graph.answer_timeout = 60;
6084 state.github_text = Some(crate::run::GithubTextAsk {
6085 fingerprint: "f".into(),
6086 title: true,
6087 body: false,
6088 categories: vec!["title-language".into()],
6089 question: None,
6090 asks: 1,
6091 resolved: false,
6092 chosen_title: None,
6093 });
6094 state.save().unwrap();
6095 let store = ask::Questions::open();
6096 let mut q = ask::Question::new(
6097 state.id.clone(),
6098 crate::github_text::ASK_NODE.to_owned(),
6099 crate::github_text::ASK_SEAT.to_owned(),
6100 "s".into(),
6101 String::new(),
6102 vec!["use fallback".into()],
6103 );
6104 store.put(&mut q).unwrap();
6105 let mut t = task();
6106 t.runs.push(state.id.clone());
6107 assert_eq!(land_resume_state(&t), LandResume::StillWaiting);
6108 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
6109 store.put(&mut q).unwrap();
6110 assert_eq!(land_resume_state(&t), LandResume::Ready);
6111 assert!(!store.get(&q.id).unwrap().status.open());
6112 state.github_text.as_mut().unwrap().resolved = true;
6114 state.save().unwrap();
6115 assert_eq!(land_resume_state(&t), LandResume::NotLanding);
6116 }
6117
6118 #[test]
6119 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
6120 crate::run::pin_test_home();
6125 let mut state = run_state(RunStatus::Landing);
6126 state.id = "20260101-000000-exp1".to_owned();
6127 state.parked = true;
6128 state.config.graph.answer_timeout = 60;
6129 state.save().unwrap();
6130
6131 let store = ask::Questions::open();
6132 let mut q = approval_question(&state.id);
6133 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
6134 store.put(&mut q).unwrap();
6135
6136 let mut t = task();
6137 t.runs.push(state.id.clone());
6138 assert_eq!(
6139 land_resume_state(&t),
6140 LandResume::Ready,
6141 "an expired question must not be waited on forever"
6142 );
6143
6144 let after = store.get(&q.id).unwrap();
6145 assert!(
6146 !after.status.open(),
6147 "the question is abandoned, not silently ignored"
6148 );
6149 assert!(
6150 after.resolution().is_none(),
6151 "an abandoned question is not read as a decision"
6152 );
6153 }
6154
6155 #[test]
6156 fn reclaim_settles_a_running_task_against_its_last_run() {
6157 let mut t = task();
6158 t.start("20260904-000000-4043".to_owned());
6159 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
6160 assert_eq!(
6161 t.status,
6162 TaskStatus::Done,
6163 "a run that actually finished must not stay `running` forever"
6164 );
6165 }
6166
6167 #[test]
6168 fn reclaim_maps_a_parked_run_to_parked_without_touching_last_error() {
6169 let mut t = task();
6170 t.start("20260904-000000-4043".to_owned());
6171 t.last_error = Some("earlier trouble".to_owned());
6172 let mut state = run_state(RunStatus::Judging);
6173 state.parked = true;
6174 reclaim(&mut t, Some(state), 1, "en");
6175 assert_eq!(t.status, TaskStatus::Parked);
6176 assert_eq!(t.attempts, 0, "the park is refunded, even at max_attempts");
6177 assert_eq!(t.last_error.as_deref(), Some("earlier trouble"));
6178 assert!(t.park_reason.is_some());
6179 assert_eq!(t.runs, ["20260904-000000-4043"], "the same run is kept");
6180 assert!(!t.fresh_start);
6181 assert!(t.status.runnable());
6182 }
6183
6184 #[test]
6185 fn settle_parks_a_task_without_a_failure_and_a_new_start_clears_the_reason() {
6186 let mut t = task();
6187 t.start("20260904-000000-4043".to_owned());
6188 t.last_error = Some("earlier trouble".to_owned());
6189 settle(
6190 &mut t,
6191 Verdict {
6192 status: RunStatus::Judging,
6193 left_pr: false,
6194 quota_hit: false,
6195 parked: true,
6196 no_viable_candidates: false,
6197 },
6198 "parked after `judging`",
6199 1,
6200 );
6201 assert_eq!(t.status, TaskStatus::Parked);
6202 assert_eq!(t.attempts, 0);
6203 assert_eq!(t.last_error.as_deref(), Some("earlier trouble"));
6204 assert_eq!(t.park_reason.as_deref(), Some("parked after `judging`"));
6205 assert_eq!(t.schema, crate::queue::SCHEMA);
6206 assert!(
6207 crate::notices::task_held(&t).is_none(),
6208 "a park is not news"
6209 );
6210 t.start("20260904-000000-4043".to_owned());
6211 assert_eq!(t.park_reason, None);
6212 }
6213
6214 #[test]
6215 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
6216 let mut t = task();
6220 t.start("20260904-000000-4043".to_owned());
6221 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
6222 assert_eq!(t.status, TaskStatus::Failed);
6223 assert!(t.status.runnable());
6224 }
6225
6226 #[test]
6227 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
6228 let mut t = task();
6229 t.start("20260904-000000-4043".to_owned());
6230 reclaim(&mut t, None, 2, "en");
6231 assert_eq!(t.status, TaskStatus::Held);
6232 assert!(
6233 t.last_error
6234 .as_deref()
6235 .is_some_and(|e| e.contains("running")),
6236 "the operator needs to know why this task was held"
6237 );
6238 }
6239
6240 #[test]
6241 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
6242 let dir = tempfile::tempdir().unwrap();
6243 let queue = Queue::at(dir.path().to_path_buf());
6244
6245 let mut orphaned = task();
6247 orphaned.id = "20260904-000000-orph".to_owned();
6248 orphaned.status = TaskStatus::Running;
6249 orphaned.attempts = 1;
6250 queue.put(&mut orphaned).unwrap();
6251
6252 let mut alive = task();
6253 alive.id = "20260904-000000-live".to_owned();
6254 alive.status = TaskStatus::Running;
6255 alive.attempts = 1;
6256 queue.put(&mut alive).unwrap();
6257 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
6258
6259 let mut queued = task();
6260 queued.id = "20260904-000000-wait".to_owned();
6261 queue.put(&mut queued).unwrap();
6262
6263 let reclaimed = reclaim_orphaned_running(&queue, 2);
6264 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
6265
6266 assert_eq!(
6267 queue.get(&orphaned.id).unwrap().status,
6268 TaskStatus::Held,
6269 "nothing was driving it and there was no run to recover"
6270 );
6271 assert_eq!(
6272 queue.get(&alive.id).unwrap().status,
6273 TaskStatus::Running,
6274 "a live claim must protect the task it belongs to"
6275 );
6276 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
6277 }
6278
6279 fn read_run_under(home: &Path, id: &str) -> RunState {
6285 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
6286 serde_json::from_str(&body).unwrap()
6287 }
6288
6289 #[test]
6290 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
6291 let dir = tempfile::tempdir().unwrap();
6292 let home = dir.path().to_path_buf();
6293 let now = Timestamp::now();
6294 let overrun_seat = || crate::run::ActiveSeat {
6295 node: "implement".to_owned(),
6296 started_at: now - jiff::SignedDuration::new(21_000, 0),
6297 timeout_secs: 3_600,
6298 attempt: 0,
6299 task: None,
6300 command: None,
6301 index: None,
6302 total: None,
6303 };
6304
6305 let mut dead = run_state(RunStatus::Implementing);
6306 dead.id = "20260101-000000-dead".to_owned();
6307 dead.active.insert("impl-A".to_owned(), overrun_seat());
6308 dead.driver_pid = Some(4242);
6311 dead.save_under(&home).unwrap();
6312
6313 let mut alive = run_state(RunStatus::Implementing);
6316 alive.id = "20260101-000000-aliv".to_owned();
6317 alive.active.insert("impl-A".to_owned(), overrun_seat());
6318 alive.save_under(&home).unwrap();
6319 let mut status = Status::new();
6320 status.current = vec![Current {
6321 task: "20260101-000000-task".to_owned(),
6322 run: alive.id.clone(),
6323 }];
6324 write_status_to(&home.join("daemon.json"), &status).unwrap();
6325
6326 let questions = Questions::at(home.join("questions"));
6330 let mut q = ask::Question::new(
6331 dead.id.clone(),
6332 "implement".to_owned(),
6333 "impl-A".to_owned(),
6334 "Which storage backend?".to_owned(),
6335 String::new(),
6336 vec!["SQLite".to_owned(), "Redis".to_owned()],
6337 );
6338 questions.put(&mut q).unwrap();
6339
6340 let abandoned = reclaim_abandoned_runs_with(
6341 &home,
6342 now,
6343 |pid| if pid == 4242 { Some(false) } else { None },
6344 |_| panic!("a query answering Dead outright needs no identity corroboration"),
6345 );
6346 assert_eq!(abandoned, vec![dead.id.clone()]);
6347
6348 let reloaded = read_run_under(&home, &dead.id);
6349 assert_eq!(reloaded.status, RunStatus::Failed);
6350 assert!(reloaded.active.is_empty());
6351 assert!(
6352 !questions.get(&q.id).unwrap().status.open(),
6353 "the failed run's own open question must be settled in the same pass"
6354 );
6355
6356 let still_alive = read_run_under(&home, &alive.id);
6357 assert_eq!(
6358 still_alive.status,
6359 RunStatus::Implementing,
6360 "a live daemon's claim protects it"
6361 );
6362 assert!(!still_alive.active.is_empty());
6363 }
6364
6365 #[test]
6375 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
6376 let dir = tempfile::tempdir().unwrap();
6377 let home = dir.path().to_path_buf();
6378 let now = Timestamp::now();
6379
6380 let mut manual = run_state(RunStatus::Reviewing);
6381 manual.id = "20260101-000000-manl".to_owned();
6382 manual.active.insert(
6383 "review-1".to_owned(),
6384 crate::run::ActiveSeat {
6385 node: "review".to_owned(),
6386 started_at: now - jiff::SignedDuration::new(21_000, 0),
6387 timeout_secs: 3_600,
6388 attempt: 0,
6389 task: None,
6390 command: None,
6391 index: None,
6392 total: None,
6393 },
6394 );
6395 manual.driver_pid = Some(4242);
6399 manual.driver_started_at = Some("1790000000".to_owned());
6400 manual.save_under(&home).unwrap();
6401
6402 let abandoned = reclaim_abandoned_runs_with(
6403 &home,
6404 now,
6405 |pid| if pid == 4242 { Some(true) } else { None },
6406 |pid| {
6407 if pid == 4242 {
6408 Some("1790000000".to_owned())
6409 } else {
6410 None
6411 }
6412 },
6413 );
6414 assert!(
6415 abandoned.is_empty(),
6416 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
6417 );
6418
6419 let reloaded = read_run_under(&home, &manual.id);
6420 assert_eq!(reloaded.status, RunStatus::Reviewing);
6421 assert!(!reloaded.active.is_empty());
6422 }
6423
6424 #[test]
6425 fn an_already_claimed_task_is_skipped_rather_than_failed() {
6426 let dir = tempfile::tempdir().unwrap();
6427 let queue = Queue::at(dir.path().to_path_buf());
6428 let mut only = task();
6429 queue.put(&mut only).unwrap();
6430
6431 let _elsewhere = queue.claim(&only.id).unwrap();
6432 let candidates = runnable(&queue);
6433 assert_eq!(candidates.len(), 1, "the task is still runnable");
6434 assert!(
6435 queue.claim(&candidates[0].id).is_err(),
6436 "the loop cannot take a claim somebody else holds"
6437 );
6438
6439 let after = queue.get(&only.id).unwrap();
6440 assert_eq!(after.status, TaskStatus::Queued);
6441 assert_eq!(
6442 after.attempts, 0,
6443 "losing the race is not an attempt at the task"
6444 );
6445 assert_eq!(after.last_error, None);
6446 }
6447
6448 #[test]
6449 fn the_status_file_round_trips_and_its_heartbeat_advances() {
6450 let dir = tempfile::tempdir().unwrap();
6451 let path = dir.path().join("daemon.json");
6452
6453 let mut status = Status::new();
6454 status.idle = false;
6455 status.completed = 7;
6456 status.current = vec![Current {
6457 task: "20260902-000000-t111".to_owned(),
6458 run: "20260902-000001-r111".to_owned(),
6459 }];
6460 write_status_to(&path, &status).unwrap();
6461 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6462 assert_eq!(first.schema, SCHEMA);
6463 assert_eq!(first.pid, std::process::id());
6464 assert!(!first.idle);
6465 assert_eq!(first.completed, 7);
6466 assert_eq!(first.current, status.current);
6467 assert!(
6468 !path.with_extension("json.tmp").exists(),
6469 "the temp file is renamed, not left behind"
6470 );
6471
6472 std::thread::sleep(Duration::from_millis(5));
6473 status.updated_at = Timestamp::now();
6474 status.polls = 3;
6475 write_status_to(&path, &status).unwrap();
6476 let second: Status =
6477 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6478 assert!(
6479 second.updated_at > first.updated_at,
6480 "a reader can only detect staleness if the heartbeat moves"
6481 );
6482 assert_eq!(
6483 second.started_at, first.started_at,
6484 "the start time is not a heartbeat"
6485 );
6486 assert_eq!(second.polls, 3);
6487 }
6488
6489 #[test]
6490 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
6491 let dir = tempfile::tempdir().unwrap();
6492
6493 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
6494
6495 let mut status = Status::new();
6496 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
6497 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6498 let stale = read_status(dir.path()).unwrap();
6499 assert!(
6500 !stale.running(Timestamp::now()),
6501 "a minute without a heartbeat is a dead daemon, not a busy one"
6502 );
6503 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
6504
6505 status.updated_at = Timestamp::now();
6506 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6507 let fresh = read_status(dir.path()).unwrap();
6508 assert!(fresh.running(Timestamp::now()));
6509 }
6510
6511 #[test]
6512 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
6513 let dir = tempfile::tempdir().unwrap();
6514 let now = Timestamp::now();
6515 let mine = "20260903-080619-01c2";
6516
6517 assert!(
6518 !is_working_on(dir.path(), mine, now),
6519 "no status file means nobody is working on anything"
6520 );
6521
6522 let mut status = Status::new();
6523 status.current = vec![Current {
6524 task: "20260903-080340-0167".to_owned(),
6525 run: mine.to_owned(),
6526 }];
6527 status.updated_at = now;
6528 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6529 assert!(is_working_on(dir.path(), mine, now));
6530 assert!(
6531 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
6532 "a daemon busy with one run is not working on another"
6533 );
6534
6535 status.updated_at = now - jiff::SignedDuration::from_secs(600);
6538 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6539 assert!(
6540 !is_working_on(dir.path(), mine, now),
6541 "a stale heartbeat is a dead daemon, so its run is a leftover"
6542 );
6543 }
6544
6545 #[test]
6546 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
6547 let dir = tempfile::tempdir().unwrap();
6548 let now = Timestamp::now();
6549
6550 assert!(
6551 !is_working_on_short(dir.path(), "01c2", now),
6552 "no status file means nobody is working on anything"
6553 );
6554
6555 let mut status = Status::new();
6556 status.current = vec![Current {
6557 task: "20260903-080340-0167".to_owned(),
6558 run: "20260903-080619-01c2".to_owned(),
6559 }];
6560 status.updated_at = now;
6561 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6562 assert!(
6563 is_working_on_short(dir.path(), "01c2", now),
6564 "the run's short id is the last block of its full id"
6565 );
6566 assert!(
6567 !is_working_on_short(dir.path(), "3cbf", now),
6568 "a daemon busy with one worktree bay is not working on another"
6569 );
6570 }
6571
6572 #[test]
6573 fn a_newer_status_file_still_yields_a_reading() {
6574 let dir = tempfile::tempdir().unwrap();
6575 std::fs::write(
6578 dir.path().join("daemon.json"),
6579 serde_json::json!({
6580 "schema": 2,
6581 "updated_at": Timestamp::now().to_string(),
6582 "idle": true,
6583 "surprise": { "nested": [1, 2, 3] },
6584 })
6585 .to_string(),
6586 )
6587 .unwrap();
6588
6589 let reading = read_status(dir.path()).expect("a forward-compatible read");
6590 assert!(reading.running(Timestamp::now()));
6591 assert!(reading.idle);
6592 assert!(reading.current.is_empty());
6593 }
6594
6595 #[test]
6596 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
6597 let dir = tempfile::tempdir().unwrap();
6603 std::fs::write(
6604 dir.path().join("daemon.json"),
6605 serde_json::json!({
6606 "schema": 1,
6607 "pid": 4242,
6608 "updated_at": Timestamp::now().to_string(),
6609 "idle": false,
6610 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
6611 "completed": 3,
6612 "polls": 9,
6613 })
6614 .to_string(),
6615 )
6616 .unwrap();
6617
6618 let reading = read_status(dir.path()).expect("an older shape must still parse");
6619 assert!(reading.running(Timestamp::now()));
6620 assert_eq!(
6621 reading.current,
6622 vec![Current {
6623 task: "20260902-140501-aaaa".to_owned(),
6624 run: "20260902-140502-bbbb".to_owned(),
6625 }]
6626 );
6627 }
6628
6629 #[test]
6630 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
6631 let dir = tempfile::tempdir().unwrap();
6632 std::fs::write(
6633 dir.path().join("daemon.json"),
6634 serde_json::json!({
6635 "schema": 1,
6636 "updated_at": Timestamp::now().to_string(),
6637 "idle": true,
6638 "current": null,
6639 })
6640 .to_string(),
6641 )
6642 .unwrap();
6643 let with_null = read_status(dir.path()).expect("null must still parse");
6644 assert!(with_null.current.is_empty());
6645
6646 std::fs::write(
6647 dir.path().join("daemon.json"),
6648 serde_json::json!({
6649 "schema": 1,
6650 "updated_at": Timestamp::now().to_string(),
6651 "idle": true,
6652 })
6653 .to_string(),
6654 )
6655 .unwrap();
6656 let absent = read_status(dir.path()).expect("a missing field must still parse");
6657 assert!(absent.current.is_empty());
6658 }
6659
6660 #[test]
6661 fn a_task_without_a_repository_runs_in_the_daemons_default() {
6662 let fallback = Path::new("/default");
6663 let mut blank = task();
6664 blank.repo = PathBuf::new();
6665 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
6666 let mut dot = task();
6667 dot.repo = PathBuf::from(".");
6668 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
6669 assert_eq!(
6670 repo_for(&task(), fallback),
6671 PathBuf::from("/repo"),
6672 "a task that names a repository keeps it"
6673 );
6674 }
6675
6676 #[test]
6677 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
6678 let mut solo_cfg = Config::default();
6684 solo_cfg.graph.implementers = 3;
6685 let mut solo_task = task();
6686 solo_task.solo = true;
6687 apply_solo(&mut solo_cfg, &solo_task);
6688 assert_eq!(solo_cfg.graph.implementers, 1);
6689
6690 let mut plain_cfg = Config::default();
6691 plain_cfg.graph.implementers = 3;
6692 let plain_task = task();
6693 assert!(!plain_task.solo);
6694 apply_solo(&mut plain_cfg, &plain_task);
6695 assert_eq!(
6696 plain_cfg.graph.implementers, 3,
6697 "a task that did not ask to run alone keeps the config's candidates"
6698 );
6699 }
6700
6701 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
6702 QuotaLoss {
6703 seat: seat.into(),
6704 node: "judge".into(),
6705 at: at.parse().unwrap(),
6706 reset: reset.map(str::to_string),
6707 }
6708 }
6709
6710 #[test]
6711 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
6712 let old: Vec<QuotaLoss> = (1..=4)
6713 .map(|i| {
6714 loss(
6715 &format!("judge-{i}"),
6716 "2026-09-23T05:23:00Z",
6717 Some("2:40pm (Asia/Tokyo)"),
6718 )
6719 })
6720 .collect();
6721 let fresh = losses_this_attempt(&old, &old);
6722 assert!(fresh.is_empty());
6723 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
6724 }
6726
6727 #[test]
6728 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6729 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6730 let now = Timestamp::now();
6731 let mut after = old.clone();
6732 after.push(loss("judge-2", &now.to_string(), None));
6733 let fresh = losses_this_attempt(&old, &after);
6734 assert_eq!(fresh, vec![after[1].clone()]);
6735 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6736 assert_eq!(
6737 until,
6738 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6739 );
6740 }
6741
6742 #[test]
6743 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6744 let before = vec![
6747 loss("judge-1", "2026-09-23T05:23:00Z", None),
6748 loss("judge-2", "2026-09-23T05:24:00Z", None),
6749 ];
6750 let after = vec![
6751 loss("judge-2", "2026-09-23T05:24:00Z", None),
6752 loss("judge-1", "2026-09-24T01:00:00Z", None),
6753 ];
6754 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6755 }
6756
6757 #[test]
6758 fn merge_overrides_are_parsed_or_refused() {
6759 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6760 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6761 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6762 assert!(merge_mode("squash").is_err());
6763 }
6764
6765 #[test]
6766 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6767 let now = Timestamp::now();
6768 let fallback = Duration::from_secs(300);
6769 let cap = Duration::from_secs(1800);
6770
6771 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6773
6774 let soon = now + jiff::SignedDuration::from_secs(600);
6776 assert_eq!(
6777 quota_wait(Some(soon), now, fallback, cap),
6778 Duration::from_secs(600)
6779 );
6780
6781 let past = now - jiff::SignedDuration::from_secs(60);
6784 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6785
6786 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6789 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6790 }
6791
6792 #[test]
6793 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6794 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6795
6796 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6797 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6798
6799 let already_past =
6803 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6804 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6805
6806 assert!(
6807 parse_reset_hint("session limit reached", now, now).is_none(),
6808 "free text with no recognised shape is not guessed at"
6809 );
6810 assert!(
6811 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6812 "an unresolvable zone name is not guessed at either"
6813 );
6814 }
6815
6816 #[test]
6817 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6818 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6819
6820 let at = parse_reset_hint(
6821 "You've hit your usage limit. Visit \
6822 https://chatgpt.com/codex/settings/usage to purchase more \
6823 credits or try again at Sep 19th, 2026 5:10 PM.",
6824 now,
6825 now,
6826 )
6827 .expect("the codex reset wording is a recognised shape");
6828 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6829
6830 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6835 .expect("an explicit year needs no rollover");
6836 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6837
6838 assert!(
6839 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6840 "a two-digit year is not the documented shape and is not guessed at"
6841 );
6842 assert!(
6843 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6844 "a four-letter month name is not the documented three-letter abbreviation"
6845 );
6846 assert!(
6847 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6848 "an explicit zone on the dated shape is a format nobody has \
6849 documented, and is refused rather than guessed at as UTC"
6850 );
6851 }
6852
6853 #[test]
6854 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6855 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6856 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6857
6858 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6859 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6860
6861 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6862 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6863
6864 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6865 assert!(
6866 parse_reset_hint(bad, now, recorded).is_none(),
6867 "{bad:?} must not be guessed at"
6868 );
6869 }
6870 }
6871
6872 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6876 let config = dir.join("magi.toml");
6877 std::fs::write(
6878 &config,
6879 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6880 )
6881 .unwrap();
6882 let opts = Opts {
6883 poll: Duration::from_secs(30),
6884 config: Some(config),
6885 repo: dir.join("repo"),
6889 ..Opts::default()
6890 };
6891 let home = dir.join("home");
6900 let worktrees = dir.join("wt");
6901 (
6902 opts,
6903 Queue::at(dir.join("queue")),
6904 home.join("daemon.json"),
6905 home,
6906 worktrees,
6907 )
6908 }
6909
6910 #[test]
6911 fn a_stop_is_idempotent_and_once_set_stays_set() {
6912 let stop = Stop::new();
6913 assert!(!stop.stopped());
6914
6915 stop.stop();
6916 assert!(stop.stopped());
6917 stop.stop();
6918 assert!(stop.stopped(), "a second stop is not a toggle");
6919
6920 let shared = stop.clone();
6921 assert!(
6922 shared.stopped(),
6923 "a clone is the same stop; that is how the loop and its caller share one"
6924 );
6925 }
6926
6927 #[test]
6928 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6929 let stop = Stop::new();
6930 stop.enter();
6931 assert!(
6932 !stop.finishing(),
6933 "a busy loop nobody has asked to stop is just running"
6934 );
6935
6936 stop.stop();
6937 assert!(
6938 stop.finishing(),
6939 "a stop asked for mid-run has not landed until the run is settled"
6940 );
6941
6942 stop.exit();
6943 assert!(
6944 !stop.finishing(),
6945 "once the run is settled the stop has landed and there is nothing to finish"
6946 );
6947 }
6948
6949 #[test]
6950 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6951 let stop = Stop::new();
6952 stop.enter();
6953 stop.enter();
6954 stop.stop();
6955 assert!(stop.finishing(), "two runs still in flight");
6956
6957 stop.exit();
6958 assert!(
6959 stop.finishing(),
6960 "one run finished, but a sibling is still working"
6961 );
6962
6963 stop.exit();
6964 assert!(
6965 !stop.finishing(),
6966 "the last run out is what actually lands the stop"
6967 );
6968 }
6969
6970 #[tokio::test]
6971 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6972 let dir = tempfile::tempdir().unwrap();
6973 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6974 let stop = Stop::new();
6975 stop.stop();
6976
6977 let began = std::time::Instant::now();
6978 tokio::time::timeout(
6979 Duration::from_secs(2),
6980 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6981 )
6982 .await
6983 .expect("a stopped loop must return, not sit out its poll interval")
6984 .expect("the loop's own setup and teardown must not fail");
6985 assert!(
6986 began.elapsed() < opts.poll,
6987 "returned only after {:?}, which is a poll interval, not a stop",
6988 began.elapsed()
6989 );
6990 }
6991
6992 #[tokio::test]
6993 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6994 let dir = tempfile::tempdir().unwrap();
6995 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6996 let stop = Stop::new();
6997
6998 let asker = {
7001 let stop = stop.clone();
7002 tokio::spawn(async move {
7003 tokio::time::sleep(Duration::from_millis(20)).await;
7004 stop.stop();
7005 })
7006 };
7007
7008 let began = std::time::Instant::now();
7009 tokio::time::timeout(
7010 Duration::from_secs(2),
7011 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
7012 )
7013 .await
7014 .expect("a stop asked for while idle must wake the wait")
7015 .expect("the loop's own setup and teardown must not fail");
7016 asker.await.unwrap();
7017 assert!(
7018 began.elapsed() < opts.poll,
7019 "returned only after {:?}, so the stop waited on the sleep",
7020 began.elapsed()
7021 );
7022 }
7023
7024 #[tokio::test]
7025 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
7026 let dir = tempfile::tempdir().unwrap();
7027 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
7028 let stop = Stop::new();
7029 stop.stop();
7030
7031 tokio::time::timeout(
7032 Duration::from_secs(2),
7033 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
7034 )
7035 .await
7036 .expect("a stopped loop must return")
7037 .expect("the loop's own setup and teardown must not fail");
7038
7039 assert!(
7040 home.is_dir(),
7041 "the loop did publish a status file, so its removal is the teardown and not an absence"
7042 );
7043 assert!(
7044 !status_file.exists(),
7045 "a stopped loop clears its status file"
7046 );
7047 assert!(
7048 read_status(&home).is_none(),
7049 "a reader must see no daemon at all, not a heartbeat that merely stopped"
7050 );
7051 }
7052
7053 #[tokio::test]
7054 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
7055 let dir = tempfile::tempdir().unwrap();
7056 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
7057 opts.once = true;
7058
7059 let mut settled = RunState::new(
7060 dir.path().join("repo"),
7061 "main".to_owned(),
7062 "abc1234".to_owned(),
7063 "fixture".to_owned(),
7064 Config::default(),
7065 );
7066 settled.status = RunStatus::Ready;
7067 let run_dir = home.join("runs").join(&settled.id);
7068 std::fs::create_dir_all(&run_dir).unwrap();
7069 std::fs::write(
7070 run_dir.join("run.json"),
7071 serde_json::to_string_pretty(&settled).unwrap(),
7072 )
7073 .unwrap();
7074 let questions = Questions::at(home.join("questions"));
7075 let mut question = ask::Question::new(
7076 settled.id.clone(),
7077 "review".to_owned(),
7078 "reviewer-1".to_owned(),
7079 "Continue?".to_owned(),
7080 String::new(),
7081 Vec::new(),
7082 );
7083 questions.put(&mut question).unwrap();
7084
7085 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
7086 .await
7087 .unwrap();
7088
7089 assert_eq!(
7090 questions.get(&question.id).unwrap().status,
7091 ask::QuestionStatus::Abandoned,
7092 "an empty --once drain still performs startup question cleanup"
7093 );
7094 }
7095
7096 #[test]
7097 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
7098 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
7099
7100 assert!(
7101 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
7102 "never checked before: due at once"
7103 );
7104
7105 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
7106 assert!(
7107 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
7108 "well inside the interval: not due yet"
7109 );
7110
7111 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
7112 assert!(
7113 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
7114 "exactly at the edge: not yet due, same convention as `clean::due`"
7115 );
7116
7117 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
7118 assert!(
7119 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
7120 "past the interval: due again"
7121 );
7122 }
7123
7124 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
7129 let config = dir.join("magi.toml");
7130 std::fs::write(
7136 &config,
7137 format!(
7138 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
7139 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
7140 cache_dir.display()
7141 ),
7142 )
7143 .unwrap();
7144 Opts {
7145 config: Some(config),
7146 repo: dir.join("repo"),
7147 ..Opts::default()
7148 }
7149 }
7150
7151 #[tokio::test]
7152 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
7153 let dir = tempfile::tempdir().unwrap();
7154 let home = dir.path().join("home");
7155 let cache_dir = dir.path().join("cache");
7156 std::fs::create_dir_all(&cache_dir).unwrap();
7157 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
7158 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
7159
7160 let running = Stop::new();
7163 let mut last_checked = None;
7164 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
7165 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
7166 .await;
7167 assert_eq!(
7168 crate::disk::dir_size(&cache_dir),
7169 0,
7170 "over the cap on the first check ever: pruned at once, no idle queue required"
7171 );
7172 assert_eq!(last_checked, Some(t0));
7173
7174 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
7176 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
7177 maybe_prune_cache_between_runs(
7178 &opts.repo,
7179 &opts,
7180 &home,
7181 &running,
7182 &mut last_checked,
7183 too_soon,
7184 )
7185 .await;
7186 assert_eq!(
7187 crate::disk::dir_size(&cache_dir),
7188 10,
7189 "too soon since the last check: left alone rather than rescanned every call"
7190 );
7191 assert_eq!(
7192 last_checked,
7193 Some(t0),
7194 "an idle check does not reset the clock"
7195 );
7196
7197 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
7199 maybe_prune_cache_between_runs(
7200 &opts.repo,
7201 &opts,
7202 &home,
7203 &running,
7204 &mut last_checked,
7205 due_again,
7206 )
7207 .await;
7208 assert_eq!(
7209 crate::disk::dir_size(&cache_dir),
7210 0,
7211 "due again: pruned back under the cap"
7212 );
7213 }
7214
7215 #[tokio::test]
7223 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
7224 let dir = tempfile::tempdir().unwrap();
7225 let home = dir.path().join("home");
7226 let cache_dir = dir.path().join("cache");
7227 std::fs::create_dir_all(&cache_dir).unwrap();
7228 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
7229 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
7230
7231 let stop = Stop::new();
7232 stop.stop();
7233 assert!(
7234 !stop.finishing(),
7235 "no run is in flight at a between-runs boundary, so nothing else \
7236 would tell the operator this stop had not taken effect yet"
7237 );
7238
7239 let mut last_checked = None;
7240 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
7241 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
7242 .await;
7243 assert_eq!(
7244 crate::disk::dir_size(&cache_dir),
7245 10,
7246 "over its cap, and due for the first check ever, but a stop outranks \
7247 it: the cap is a standing policy the next start measures again"
7248 );
7249 assert_eq!(
7250 last_checked, None,
7251 "a check that never happened must not claim the interval"
7252 );
7253 }
7254
7255 #[tokio::test]
7269 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
7270 let dir = tempfile::tempdir().unwrap();
7271 let cache_dir = dir.path().join("cache");
7272 std::fs::create_dir_all(&cache_dir).unwrap();
7273 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
7274
7275 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
7276 opts.poll = Duration::from_millis(20);
7277 opts.max_attempts = 1_000;
7278
7279 let queue = Queue::at(dir.path().join("queue"));
7280 let mut t = Task::new(
7281 "x".to_owned(),
7282 "x".to_owned(),
7283 opts.repo.clone(),
7284 Source::Human,
7285 );
7286 queue.put(&mut t).unwrap();
7287
7288 let home = dir.path().join("home");
7289 let worktrees = dir.path().join("wt");
7290 let status_file = home.join("daemon.json");
7291 let stop = Stop::new();
7292 let stopper = {
7293 let stop = stop.clone();
7294 let queue = Queue::at(dir.path().join("queue"));
7295 let id = t.id.clone();
7296 tokio::spawn(async move {
7297 tokio::time::sleep(Duration::from_millis(400)).await;
7301 for _ in 0..200 {
7302 if queue.get(&id).is_ok_and(|t| t.attempts >= 2) {
7303 break;
7304 }
7305 tokio::time::sleep(Duration::from_millis(50)).await;
7306 }
7307 stop.stop();
7308 })
7309 };
7310
7311 tokio::time::timeout(
7312 Duration::from_secs(20),
7313 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
7314 )
7315 .await
7316 .expect("the loop must not hang on a queue that keeps producing failing work")
7317 .expect("the loop's own setup and teardown must not fail");
7318 stopper.await.unwrap();
7319
7320 let after = queue.get(&t.id).unwrap();
7321 assert!(
7322 after.attempts >= 2,
7323 "the harness must actually have retried more than once, or this is not \
7324 exercising a busy queue at all (got {} attempt(s))",
7325 after.attempts
7326 );
7327 assert!(
7328 after.status.runnable(),
7329 "still under its attempt budget: the queue never reached a natural idle \
7330 on its own, only the external stop ended the test"
7331 );
7332
7333 assert_eq!(
7334 crate::disk::dir_size(&cache_dir),
7335 0,
7336 "an oversized cache must not be left to grow unboundedly just because the \
7337 queue kept the loop busy the whole time"
7338 );
7339 }
7340
7341 #[test]
7342 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
7343 let dir = tempfile::tempdir().unwrap();
7344 let queue = Queue::at(dir.path().join("queue"));
7345 let questions = Questions::at(dir.path().join("questions"));
7346 let mut task = task();
7347 queue.put(&mut task).unwrap();
7348
7349 let mut task_question = ask::Question::new(
7350 task.id.clone(),
7351 crate::conduct::NODE.to_owned(),
7352 "conduct".to_owned(),
7353 "Which backend?".to_owned(),
7354 String::new(),
7355 Vec::new(),
7356 );
7357 questions.put(&mut task_question).unwrap();
7358 task.block(vec![task_question.id.clone()], None);
7359 queue.put(&mut task).unwrap();
7360
7361 let mut run_question = ask::Question::new(
7362 "20260101-000000-run1".to_owned(),
7363 "review".to_owned(),
7364 "reviewer-1".to_owned(),
7365 "Run question".to_owned(),
7366 String::new(),
7367 Vec::new(),
7368 );
7369 questions.put(&mut run_question).unwrap();
7370
7371 let mut coincidental = ask::Question::new(
7376 task.id.clone(),
7377 "review".to_owned(),
7378 "reviewer-1".to_owned(),
7379 "Unrelated review question".to_owned(),
7380 String::new(),
7381 Vec::new(),
7382 );
7383 questions.put(&mut coincidental).unwrap();
7384
7385 reconcile_task_questions(&queue, &questions);
7386 assert!(questions.get(&task_question.id).unwrap().status.open());
7387 assert!(questions.get(&run_question.id).unwrap().status.open());
7388 assert!(questions.get(&coincidental.id).unwrap().status.open());
7389
7390 task.release();
7391 queue.put(&mut task).unwrap();
7392 reconcile_task_questions(&queue, &questions);
7393 assert_eq!(
7394 questions.get(&task_question.id).unwrap().status,
7395 ask::QuestionStatus::Abandoned
7396 );
7397 assert!(
7398 questions.get(&run_question.id).unwrap().status.open(),
7399 "run questions remain the run janitor's responsibility"
7400 );
7401 assert!(
7402 questions.get(&coincidental.id).unwrap().status.open(),
7403 "a non-conductor question must not be abandoned just because its \
7404 run id coincides with a task id"
7405 );
7406 }
7407
7408 #[test]
7409 fn a_freshly_started_running_task_is_never_stalled() {
7410 let dir = tempfile::tempdir().unwrap();
7411 let mut t = task();
7412 t.start("run-1".to_owned());
7413 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
7416 }
7417
7418 #[test]
7419 fn a_long_running_task_with_no_live_daemon_is_stalled() {
7420 let dir = tempfile::tempdir().unwrap();
7421 let mut t = task();
7422 t.start("run-1".to_owned());
7423 t.updated_at = Timestamp::now()
7424 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
7425 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
7426 assert_eq!(
7427 stalled_tasks(
7428 &Queue::at(dir.path().join("q")),
7429 dir.path(),
7430 Timestamp::now()
7431 )
7432 .len(),
7433 0,
7434 "the task was never written to this queue"
7435 );
7436 }
7437
7438 #[test]
7439 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
7440 let dir = tempfile::tempdir().unwrap();
7441 let mut t = task();
7442 t.id = "20260903-080340-0167".to_owned();
7443 t.start("20260903-080619-01c2".to_owned());
7444 t.updated_at = Timestamp::now()
7445 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
7446
7447 let mut status = Status::new();
7448 status.current = vec![Current {
7449 task: t.id.clone(),
7450 run: "20260903-080619-01c2".to_owned(),
7451 }];
7452 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
7453
7454 assert!(
7455 !is_stalled(&t, dir.path(), Timestamp::now()),
7456 "a live daemon's own heartbeat rules out stalled, however long the task has run"
7457 );
7458 }
7459
7460 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
7464 let path = queue.path_of(id);
7465 let body = std::fs::read_to_string(&path).unwrap();
7466 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
7467 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
7468 v["updated_at"] = serde_json::Value::String(old.to_string());
7469 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
7470 }
7471
7472 #[test]
7473 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
7474 let dir = tempfile::tempdir().unwrap();
7487 let queue = Queue::at(dir.path().join("queue"));
7488 let home = dir.path().join("home");
7489
7490 let mut t = task();
7491 t.id = "20260101-000001-lock".to_owned();
7492 t.start("run-1".to_owned());
7493 queue.put(&mut t).unwrap();
7494 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7495 std::fs::write(
7496 dir.path().join("queue").join(format!("{}.lock", t.id)),
7497 "not a pid",
7498 )
7499 .unwrap();
7500
7501 let now = Timestamp::now();
7502 assert!(
7503 reclaim_orphaned_running(&queue, 2).is_empty(),
7504 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
7505 and reclaim must leave the task alone"
7506 );
7507 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
7508
7509 let stalled = stalled_tasks(&queue, &home, now);
7510 assert_eq!(
7511 stalled.len(),
7512 1,
7513 "reclaim's inability to claim it yet must not hide it from the conductor"
7514 );
7515 assert_eq!(stalled[0].id, t.id);
7516 }
7517
7518 #[test]
7519 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
7520 let dir = tempfile::tempdir().unwrap();
7521 crate::run::set_home(dir.path().join("run-home"));
7522 let queue = Queue::at(dir.path().join("queue"));
7523 let home = dir.path().join("home");
7524 let questions = Questions::at(dir.path().join("questions"));
7525
7526 let mut t = task();
7527 t.id = "20260101-000003-dead".to_owned();
7528 t.start("missing-run".to_owned());
7529 queue.put(&mut t).unwrap();
7530 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7531
7532 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7535 assert_eq!(
7536 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
7537 [&t.id]
7538 );
7539 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
7540 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7541
7542 crate::conduct::apply(
7545 &queue,
7546 &questions,
7547 &crate::conduct::Verdict {
7548 decisions: vec![crate::conduct::Decision {
7549 id: t.id.clone(),
7550 recovery: Some(crate::conduct::Recovery::Requeue),
7551 ..crate::conduct::Decision::default()
7552 }],
7553 },
7554 )
7555 .unwrap();
7556 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7557 }
7558
7559 #[test]
7560 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
7561 let dir = tempfile::tempdir().unwrap();
7562 let queue = Queue::at(dir.path().join("queue"));
7563 let home = dir.path().join("home");
7564
7565 let mut fresh = task();
7566 fresh.id = "20260101-000001-aaaa".to_owned();
7567 fresh.start("run-1".to_owned());
7568 queue.put(&mut fresh).unwrap();
7569
7570 let mut old = task();
7571 old.id = "20260101-000002-bbbb".to_owned();
7572 old.start("run-2".to_owned());
7573 queue.put(&mut old).unwrap();
7574 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
7575
7576 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7577 assert_eq!(stalled.len(), 1);
7578 assert_eq!(stalled[0].id, old.id);
7579 }
7580
7581 #[test]
7582 fn queued_and_finished_task_views_partition_by_status() {
7583 let dir = tempfile::tempdir().unwrap();
7584 let queue = Queue::at(dir.path().join("queue"));
7585
7586 let mut queued = task();
7587 queued.id = "20260101-000001-aaaa".to_owned();
7588 queue.put(&mut queued).unwrap();
7589
7590 let mut failed = task();
7591 failed.id = "20260101-000002-bbbb".to_owned();
7592 failed.start("run-1".to_owned());
7593 failed.fail("gate red", 5);
7594 queue.put(&mut failed).unwrap();
7595
7596 let mut held = task();
7597 held.id = "20260101-000003-cccc".to_owned();
7598 held.hold_machine(None);
7599 queue.put(&mut held).unwrap();
7600
7601 let mut running = task();
7602 running.id = "20260101-000004-dddd".to_owned();
7603 running.start("run-2".to_owned());
7604 queue.put(&mut running).unwrap();
7605
7606 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
7607 assert_eq!(queued_ids, [queued.id.clone()]);
7608
7609 let mut finished_ids: Vec<String> =
7610 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
7611 finished_ids.sort_unstable();
7612 let mut want = vec![failed.id.clone(), held.id.clone()];
7613 want.sort_unstable();
7614 assert_eq!(finished_ids, want);
7615 }
7616
7617 #[test]
7618 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
7619 let dir = tempfile::tempdir().unwrap();
7620 let queue = Queue::at(dir.path().join("queue"));
7621 let questions = ask::Questions::at(dir.path().join("questions"));
7622
7623 let mut dep = task();
7624 dep.id = "20260101-000001-dep0".to_owned();
7625 dep.succeed();
7626 queue.put(&mut dep).unwrap();
7627
7628 let mut still_going = task();
7629 still_going.id = "20260101-000002-dep1".to_owned();
7630 queue.put(&mut still_going).unwrap();
7631
7632 let mut blocked = task();
7633 blocked.id = "20260101-000003-main".to_owned();
7634 blocked.block(
7635 vec![dep.id.clone(), still_going.id.clone()],
7636 Some("waits on both".to_owned()),
7637 );
7638 queue.put(&mut blocked).unwrap();
7639
7640 resolve_blockers(&queue, &questions);
7641
7642 let after = queue.get(&blocked.id).unwrap();
7643 assert_eq!(
7644 after.status,
7645 TaskStatus::Blocked,
7646 "one dependency is still outstanding"
7647 );
7648 assert_eq!(after.blocked_by, [still_going.id.clone()]);
7649 }
7650
7651 #[test]
7652 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
7653 let dir = tempfile::tempdir().unwrap();
7654 let queue = Queue::at(dir.path().join("queue"));
7655 let questions = ask::Questions::at(dir.path().join("questions"));
7656
7657 let mut q = crate::ask::Question::new(
7658 "20260101-000001-main".to_owned(),
7659 crate::conduct::NODE.to_owned(),
7660 "conduct".to_owned(),
7661 "Which backend?".to_owned(),
7662 String::new(),
7663 Vec::new(),
7664 );
7665 questions.put(&mut q).unwrap();
7666 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
7667 .unwrap();
7668 questions.put(&mut q).unwrap();
7669
7670 let mut blocked = task();
7671 blocked.id = "20260101-000001-main".to_owned();
7672 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
7673 queue.put(&mut blocked).unwrap();
7674
7675 resolve_blockers(&queue, &questions);
7676
7677 let after = queue.get(&blocked.id).unwrap();
7678 assert_eq!(
7679 after.status,
7680 TaskStatus::Queued,
7681 "the only blocker resolved"
7682 );
7683 assert_eq!(after.answers.len(), 1);
7684 assert_eq!(after.answers[0].question, "Which backend?");
7685 assert_eq!(after.answers[0].answer, "SQLite");
7686
7687 let instruction = instruction_for(&after);
7689 assert!(instruction.contains("Which backend?"));
7690 assert!(instruction.contains("SQLite"));
7691 }
7692
7693 #[test]
7694 fn resolve_blockers_holds_a_task_whose_conductor_question_was_abandoned() {
7695 let dir = tempfile::tempdir().unwrap();
7696 let queue = Queue::at(dir.path().join("queue"));
7697 let questions = ask::Questions::at(dir.path().join("questions"));
7698
7699 let mut q = crate::ask::Question::new(
7700 "20260101-000001-main".to_owned(),
7701 crate::conduct::NODE.to_owned(),
7702 "conduct".to_owned(),
7703 "Is the setup done?".to_owned(),
7704 String::new(),
7705 Vec::new(),
7706 );
7707 q.abandon("no answer within 60s of asking");
7708 questions.put(&mut q).unwrap();
7709
7710 let mut blocked = task();
7711 blocked.id = "20260101-000001-main".to_owned();
7712 blocked.block(vec![q.id.clone()], Some("setup?".to_owned()));
7713 queue.put(&mut blocked).unwrap();
7714
7715 resolve_blockers(&queue, &questions);
7716
7717 let after = queue.get(&blocked.id).unwrap();
7718 assert_eq!(
7719 after.status,
7720 TaskStatus::Held,
7721 "never left blocked on nothing"
7722 );
7723 assert!(!after.operator_held(), "a machine hold, for triage");
7724 assert!(
7725 after
7726 .hold_reason
7727 .as_deref()
7728 .unwrap_or_default()
7729 .contains("went unanswered")
7730 );
7731 }
7732
7733 #[test]
7734 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
7735 let dir = tempfile::tempdir().unwrap();
7741 let queue = Queue::at(dir.path().join("queue"));
7742 let questions = ask::Questions::at(dir.path().join("questions"));
7743
7744 let mut q = crate::ask::Question::new(
7745 "20260101-000001-main".to_owned(),
7746 crate::conduct::NODE.to_owned(),
7747 "conduct".to_owned(),
7748 "How should this be handled?".to_owned(),
7749 String::new(),
7750 Vec::new(),
7751 );
7752 questions.put(&mut q).unwrap();
7753 q.answer(crate::ask::Answer::Text(
7754 "leave it held, a human will look at it later".to_owned(),
7755 ))
7756 .unwrap();
7757 questions.put(&mut q).unwrap();
7758
7759 let mut held = task();
7760 held.id = "20260101-000001-main".to_owned();
7761 held.hold_machine(Some("out of attempts".to_owned()));
7762 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
7763 queue.put(&mut held).unwrap();
7764
7765 resolve_blockers(&queue, &questions);
7766
7767 let after = queue.get(&held.id).unwrap();
7768 assert_eq!(after.status, TaskStatus::Held);
7769 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
7770 assert_eq!(
7771 after.answers[0].answer,
7772 "leave it held, a human will look at it later"
7773 );
7774 }
7775
7776 #[test]
7777 fn resolve_blockers_releases_a_task_whose_dependency_was_deleted_on_purpose() {
7778 let dir = tempfile::tempdir().unwrap();
7779 let queue = Queue::at(dir.path().join("queue"));
7780 let questions = ask::Questions::at(dir.path().join("questions"));
7781 let mut dep = task();
7782 dep.id = "20260101-000001-gone".to_owned();
7783 queue.put(&mut dep).unwrap();
7784 let mut blocked = task();
7785 blocked.id = "20260101-000003-main".to_owned();
7786 blocked.block(vec![dep.id.clone()], None);
7787 queue.put(&mut blocked).unwrap();
7788
7789 let claim = queue.claim(&blocked.id).unwrap();
7791 queue.remove(&dep.id, false, &questions).unwrap();
7792 drop(claim);
7793 resolve_blockers(&queue, &questions);
7794
7795 let after = queue.get(&blocked.id).unwrap();
7796 assert_eq!(after.status, TaskStatus::Queued);
7797 assert!(after.blocked_by.is_empty());
7798 }
7799
7800 #[test]
7801 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
7802 let dir = tempfile::tempdir().unwrap();
7808 let queue = Queue::at(dir.path().join("queue"));
7809 let questions = ask::Questions::at(dir.path().join("questions"));
7810
7811 let mut still_going = task();
7812 still_going.id = "20260101-000002-dep1".to_owned();
7813 queue.put(&mut still_going).unwrap();
7814
7815 let mut blocked = task();
7816 blocked.id = "20260101-000003-main".to_owned();
7817 blocked.block(
7818 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7819 Some("waits on both".to_owned()),
7820 );
7821 queue.put(&mut blocked).unwrap();
7822
7823 resolve_blockers(&queue, &questions);
7824
7825 let after = queue.get(&blocked.id).unwrap();
7826 assert_eq!(
7827 after.status,
7828 TaskStatus::Held,
7829 "a missing dependency must not leave the task blocked forever"
7830 );
7831 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7832 assert!(after.blocked_by.is_empty());
7833 let reason = after.hold_reason.as_deref().unwrap_or_default();
7834 assert!(
7835 reason.contains("20260101-000001-gone"),
7836 "the missing id must be named so an operator can tell what happened: {reason}"
7837 );
7838 assert!(
7839 reason.contains(&still_going.id),
7840 "the still-valid dependency must not silently vanish from the record: {reason}"
7841 );
7842 }
7843
7844 #[test]
7845 fn instruction_for_is_unchanged_without_any_answers() {
7846 let t = task();
7847 assert_eq!(instruction_for(&t), t.instruction);
7848 }
7849
7850 #[test]
7851 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7852 let dir = tempfile::tempdir().unwrap();
7853 let q = Queue::at(dir.path().join("queue"));
7854 let src = dir.path().join("shot.png");
7855 std::fs::write(&src, "x").unwrap();
7856 let mut t = task();
7857 q.attach(&mut t, &[src]).unwrap();
7858 let paths = task_attachments(&q, &t).unwrap();
7859 assert_eq!(paths.len(), 1);
7860 assert!(paths[0].is_absolute() && paths[0].is_file());
7861 std::fs::remove_file(&paths[0]).unwrap();
7862 let err = task_attachments(&q, &t).unwrap_err().to_string();
7863 assert!(err.contains("shot.png"), "{err}");
7864 }
7865
7866 #[test]
7867 fn resumed_instruction_is_unchanged_without_any_answers() {
7868 let t = task();
7869 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7870 }
7871
7872 #[test]
7873 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7874 let mut t = task();
7875 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7876 let old = t.instruction.clone();
7880
7881 let refreshed = resumed_instruction(&old, &t);
7882 assert!(refreshed.starts_with(&old), "the original text is kept");
7883 assert!(refreshed.contains("Which backend?"));
7884 assert!(refreshed.contains("SQLite"));
7885 }
7886
7887 #[test]
7888 fn resumed_instruction_keeps_an_original_answers_heading() {
7889 let mut t = task();
7890 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7891 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7892
7893 let refreshed = resumed_instruction(&t.instruction, &t);
7894
7895 assert!(
7896 refreshed.starts_with(&t.instruction),
7897 "an answers heading in the original instruction is not the appended block"
7898 );
7899 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7900 assert!(refreshed.contains("Which backend?"));
7901 assert!(refreshed.contains("SQLite"));
7902
7903 let repeated = resumed_instruction(&refreshed, &t);
7904 assert_eq!(
7905 repeated, refreshed,
7906 "only the final appended block is refreshed"
7907 );
7908 }
7909
7910 #[test]
7911 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7912 let mut t = task();
7913 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7914
7915 let once = resumed_instruction(&t.instruction, &t);
7919 let twice = resumed_instruction(&once, &t);
7920 assert_eq!(once, twice);
7921 assert_eq!(once.matches("Which backend?").count(), 1);
7922
7923 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7925 let refreshed = resumed_instruction(&once, &t);
7926 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7927 assert!(refreshed.contains("Which backend?"));
7928 assert!(refreshed.contains("Which cache?"));
7929 }
7930
7931 #[test]
7932 fn prepare_instruction_covers_all_three_starters() {
7933 let mut t = task();
7934 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7935
7936 assert_eq!(
7939 prepare_instruction(&Starter::Start, None, &t),
7940 Some(instruction_for(&t))
7941 );
7942
7943 let old = t.instruction.clone();
7946 assert_eq!(
7947 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7948 Some(resumed_instruction(&old, &t))
7949 );
7950
7951 assert_eq!(
7955 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7956 None
7957 );
7958 }
7959
7960 #[test]
7961 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7962 assert_eq!(
7963 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7964 Starter::Review("magi/eba2/A".to_owned())
7965 );
7966 }
7967
7968 #[test]
7969 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7970 assert_eq!(
7971 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7972 Starter::Start,
7973 "a vanished review branch must not fall back to resuming the old run either"
7974 );
7975 }
7976
7977 #[test]
7978 fn a_refused_handover_retries_as_a_review_of_the_same_branch() {
7979 let mut t = task();
7980 t.start("old-run".to_owned());
7981 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
7982 t.release();
7983 let branch = t.review_branch.take();
7984 assert_eq!(
7985 choose_starter(branch.as_deref(), true, Some("old-run")),
7986 Starter::Review("magi/eba2/A".to_owned()),
7987 "a review wins over resuming the old run"
7988 );
7989 }
7990
7991 #[test]
7992 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7993 assert_eq!(
7994 choose_starter(None, false, Some("some-run")),
7995 Starter::Resume("some-run".to_owned())
7996 );
7997 assert_eq!(choose_starter(None, false, None), Starter::Start);
7998 }
7999
8000 #[test]
8001 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
8002 let mut released = task();
8003 released.start("stalled-run".to_owned());
8004 released.requeue();
8005 let unfinished = (!released.fresh_start)
8006 .then(|| Some("stalled-run".to_owned()))
8007 .flatten();
8008 assert_eq!(
8009 choose_starter(None, false, unfinished.as_deref()),
8010 Starter::Start,
8011 "release keeps run history but must not resume it"
8012 );
8013 assert_eq!(released.runs, ["stalled-run"]);
8014 }
8015
8016 #[test]
8017 fn an_ordinary_release_keeps_a_resumable_run_available() {
8018 let mut released = task();
8019 released.start("stalled-run".to_owned());
8020 released.release();
8021 let unfinished = (!released.fresh_start)
8022 .then(|| Some("stalled-run".to_owned()))
8023 .flatten();
8024 assert_eq!(
8025 choose_starter(None, false, unfinished.as_deref()),
8026 Starter::Resume("stalled-run".to_owned()),
8027 "manual release must preserve the normal resume path"
8028 );
8029 }
8030
8031 #[test]
8032 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
8033 let mut state = run_state(RunStatus::Blocked);
8034 state.config.graph.review_rounds = 3;
8035 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
8036 assert!(exhausted_review_budget(&state));
8037
8038 state.reviews.pop();
8040 assert!(!exhausted_review_budget(&state));
8041
8042 let mut stalled = run_state(RunStatus::Stalled);
8045 stalled.config.graph.review_rounds = 1;
8046 stalled.reviews = vec![review_round(1)];
8047 assert!(!exhausted_review_budget(&stalled));
8048 }
8049
8050 fn review_round(round: usize) -> crate::run::ReviewRound {
8051 crate::run::ReviewRound {
8052 round,
8053 head: "deadbeef".to_owned(),
8054 verified_head: None,
8055 verified_at: None,
8056 reviews: Vec::new(),
8057 e2e: Vec::new(),
8058 verify_retried: false,
8059 e2e_deferred: false,
8060 e2e_defer_reason: None,
8061 fix: None,
8062 blocking: 0,
8063 answered: 1,
8064 expected: 1,
8065 clean: false,
8066 progressed: true,
8067 vote_split: false,
8068 reconsideration: Vec::new(),
8069 verdict: None,
8070 }
8071 }
8072
8073 #[test]
8074 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
8075 let mut run = RunState::new(
8076 PathBuf::from("/repo"),
8077 "main".to_owned(),
8078 "abc1234def".to_owned(),
8079 "add retries".to_owned(),
8080 Config::default(),
8081 );
8082 run.status = RunStatus::Judging;
8083 run.parked = true;
8084 let mut task = Task::new(
8085 "add retries".to_owned(),
8086 "add retries".to_owned(),
8087 PathBuf::from("/repo"),
8088 crate::queue::Source::Human,
8089 );
8090 task.status = TaskStatus::Failed;
8091 task.runs = vec![run.id.clone()];
8092 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
8093 assert!(with(&task, &run), "parked after judging is the case");
8094 let mut parked_task = task.clone();
8095 parked_task.status = TaskStatus::Parked;
8096 assert!(with(&parked_task, &run), "the Parked status is the case");
8097
8098 let mut not_parked = run.clone();
8099 not_parked.parked = false;
8100 not_parked.status = RunStatus::Stalled;
8101 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
8102
8103 let mut fresh = task.clone();
8104 fresh.fresh_start = true;
8105 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
8106
8107 let mut review = task.clone();
8108 review.review_branch = Some("magi/x/A".to_owned());
8109 assert!(!with(&review, &run), "review is ranked before resume");
8110
8111 let mut held = task.clone();
8112 held.status = TaskStatus::Held;
8113 assert!(!with(&held, &run), "a hold stays visible to the conductor");
8114
8115 let mut released = run.clone();
8116 released.released_to = Some("20260901-000000-new1".to_owned());
8117 assert!(!with(&task, &released), "nothing left to resume into");
8118
8119 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
8120 "unreadable"
8121 )));
8122 }
8123
8124 #[test]
8125 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
8126 let mut released = RunState::new(
8127 PathBuf::from("/repo"),
8128 "main".to_owned(),
8129 "abc1234def".to_owned(),
8130 "add retries".to_owned(),
8131 Config::default(),
8132 );
8133 released.status = RunStatus::Blocked;
8134 assert_eq!(
8135 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
8136 Some(released.id.clone())
8137 );
8138 released.released_to = Some("20260901-000000-new1".to_owned());
8139 assert_eq!(
8140 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
8141 None,
8142 "there is nothing left to resume it into"
8143 );
8144 }
8145
8146 #[test]
8147 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
8148 let mut exhausted = RunState::new(
8157 PathBuf::from("/repo"),
8158 "main".to_owned(),
8159 "abc1234def".to_owned(),
8160 "add retries".to_owned(),
8161 Config::default(),
8162 );
8163 exhausted.status = RunStatus::Blocked;
8164 exhausted.config.graph.review_rounds = 1;
8165 exhausted.reviews = vec![review_round(1)];
8166
8167 assert_eq!(
8168 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
8169 None,
8170 "an exhausted `Blocked` run must not be offered as resumable"
8171 );
8172
8173 let mut has_budget_left = RunState::new(
8176 PathBuf::from("/repo"),
8177 "main".to_owned(),
8178 "abc1234def".to_owned(),
8179 "add retries".to_owned(),
8180 Config::default(),
8181 );
8182 has_budget_left.status = RunStatus::Blocked;
8183 has_budget_left.config.graph.review_rounds = 3;
8184 has_budget_left.reviews = vec![review_round(1)];
8185
8186 assert_eq!(
8187 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
8188 Ok(has_budget_left.clone())
8189 }),
8190 Some(has_budget_left.id.clone())
8191 );
8192 }
8193
8194 #[test]
8195 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
8196 let mut older_stalled = RunState::new(
8204 PathBuf::from("/repo"),
8205 "main".to_owned(),
8206 "abc1234def".to_owned(),
8207 "add retries".to_owned(),
8208 Config::default(),
8209 );
8210 older_stalled.status = RunStatus::Stalled;
8211
8212 let mut newest_exhausted = RunState::new(
8213 PathBuf::from("/repo"),
8214 "main".to_owned(),
8215 "abc1234def".to_owned(),
8216 "add retries".to_owned(),
8217 Config::default(),
8218 );
8219 newest_exhausted.status = RunStatus::Blocked;
8220 newest_exhausted.config.graph.review_rounds = 1;
8221 newest_exhausted.reviews = vec![review_round(1)];
8222
8223 assert_eq!(
8224 unfinished_run_with(
8225 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
8226 "t",
8227 |_| Ok(newest_exhausted.clone())
8228 ),
8229 None,
8230 "the newest run is exhausted, so nothing here is worth resuming - \
8231 least of all the older, already-superseded run"
8232 );
8233 }
8234
8235 #[test]
8236 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
8237 assert_eq!(
8238 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
8239 Err(anyhow::anyhow!("fixture is absent"))
8240 }),
8241 None
8242 );
8243 }
8244
8245 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
8246 let mut q = ask::Question::new(
8247 run.to_owned(),
8248 "implement".to_owned(),
8249 "impl-A".to_owned(),
8250 "continue?".to_owned(),
8251 String::new(),
8252 vec!["resume で続行する".to_owned(), "other".to_owned()],
8253 );
8254 q.actions.insert("resume で続行する".to_owned(), action);
8255 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
8256 .unwrap();
8257 q
8258 }
8259
8260 fn held_task_with(run: &str) -> Task {
8261 let mut t = task();
8262 t.runs = vec![run.to_owned()];
8263 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
8264 t
8265 }
8266
8267 fn resume_action(run: &str) -> ask::ChoiceAction {
8268 ask::ChoiceAction::Resume { run: run.into() }
8269 }
8270
8271 #[test]
8272 fn decide_action_resumes_only_the_latest_resumable_run() {
8273 let t = held_task_with("r1");
8274 let q = action_question("r1", resume_action("r1"));
8275 let load = |s: RunState| move |_: &str| Ok(s);
8276 assert_eq!(
8277 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
8278 ActionDecision::Resume("r1".into())
8279 );
8280 let q_other = action_question("r1", resume_action("r0"));
8282 assert!(matches!(
8283 decide_action(
8284 &t,
8285 &q_other,
8286 &PHRASES_EN,
8287 load(run_state(RunStatus::Blocked))
8288 ),
8289 ActionDecision::Refuse(_)
8290 ));
8291 assert!(matches!(
8293 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
8294 ActionDecision::Refuse(_)
8295 ));
8296 let mut released = run_state(RunStatus::Blocked);
8298 released.released_to = Some("elsewhere".into());
8299 assert!(matches!(
8300 decide_action(&t, &q, &PHRASES_EN, load(released)),
8301 ActionDecision::Refuse(_)
8302 ));
8303 assert!(matches!(
8305 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
8306 ActionDecision::Refuse(_)
8307 ));
8308 }
8309
8310 #[test]
8311 fn decide_action_ignores_a_question_about_an_earlier_run() {
8312 let mut t = held_task_with("r1");
8313 t.runs.push("r2".to_owned());
8314 let never = |_: &str| -> Result<RunState> { bail!("not read") };
8315 assert_eq!(
8316 decide_action(
8317 &t,
8318 &action_question("r1", ask::ChoiceAction::Done),
8319 &PHRASES_EN,
8320 never
8321 ),
8322 ActionDecision::Stale
8323 );
8324 }
8325
8326 #[test]
8327 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
8328 let mut t = held_task_with("r1");
8329 let never = |_: &str| -> Result<RunState> { bail!("not read") };
8330 assert_eq!(
8331 decide_action(
8332 &t,
8333 &action_question("r1", ask::ChoiceAction::Requeue),
8334 &PHRASES_EN,
8335 never
8336 ),
8337 ActionDecision::Requeue
8338 );
8339 let done_q = action_question("r1", ask::ChoiceAction::Done);
8340 assert_eq!(
8341 decide_action(&t, &done_q, &PHRASES_EN, never),
8342 ActionDecision::Done
8343 );
8344 t.mark_action_applied(&done_q.id);
8345 assert_eq!(
8346 decide_action(&t, &done_q, &PHRASES_EN, never),
8347 ActionDecision::Skip
8348 );
8349
8350 let mut plain = action_question("r1", ask::ChoiceAction::Done);
8352 plain.actions.clear();
8353 assert_eq!(
8354 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
8355 ActionDecision::Skip
8356 );
8357 let mut running = held_task_with("r1");
8359 running.status = TaskStatus::Running;
8360 assert_eq!(
8361 decide_action(
8362 &running,
8363 &action_question("r1", ask::ChoiceAction::Done),
8364 &PHRASES_EN,
8365 never
8366 ),
8367 ActionDecision::Skip
8368 );
8369 }
8370
8371 #[test]
8372 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
8373 let dir = tempfile::tempdir().unwrap();
8374 let queue = Queue::at(dir.path().join("queue"));
8375 let questions = Questions::at(dir.path().join("questions"));
8376 let home = dir.path().join("home");
8377 let mut state = run_state(RunStatus::Blocked);
8378 state.id = "20260101-000000-act1".to_owned();
8379 state.save_under(&home).unwrap();
8380
8381 let mut t = held_task_with(&state.id);
8382 queue.put(&mut t).unwrap();
8383 let mut q = action_question(&state.id, resume_action(&state.id));
8384 questions.put(&mut q).unwrap();
8385
8386 apply_choice_actions(&queue, &questions, &home);
8387 let after = queue.get(&t.id).unwrap();
8388 assert_eq!(after.status, TaskStatus::Queued);
8389 assert!(!after.fresh_start);
8390 assert!(after.action_applied(&q.id));
8391 let pin = after.resume_override.clone().unwrap();
8392 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
8393 assert!(pin.forced);
8394
8395 let mut again = queue.get(&t.id).unwrap();
8397 again.hold_machine(Some("later".into()));
8398 queue.put(&mut again).unwrap();
8399 apply_choice_actions(&queue, &questions, &home);
8400 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8401 }
8402
8403 #[test]
8404 fn a_delivered_answer_is_still_acted_on_exactly_once() {
8405 let dir = tempfile::tempdir().unwrap();
8406 let queue = Queue::at(dir.path().join("queue"));
8407 let questions = Questions::at(dir.path().join("questions"));
8408 let home = dir.path().join("home");
8409 let mut state = run_state(RunStatus::Blocked);
8410 state.id = "20260101-000000-act2".to_owned();
8411 state.save_under(&home).unwrap();
8412
8413 let mut t = held_task_with(&state.id);
8414 queue.put(&mut t).unwrap();
8415 let mut q = action_question(&state.id, resume_action(&state.id));
8416 q.answer_delivered = true;
8417 questions.put(&mut q).unwrap();
8418
8419 apply_choice_actions(&queue, &questions, &home);
8422 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
8423 assert!(queue.get(&t.id).unwrap().action_applied(&q.id));
8424
8425 let mut again = queue.get(&t.id).unwrap();
8427 again.hold_machine(Some("later".into()));
8428 queue.put(&mut again).unwrap();
8429 apply_choice_actions(&queue, &questions, &home);
8430 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8431 }
8432
8433 #[test]
8434 fn a_fresh_asker_defers_the_action_until_it_goes_quiet() {
8435 let dir = tempfile::tempdir().unwrap();
8436 let queue = Queue::at(dir.path().join("queue"));
8437 let questions = Questions::at(dir.path().join("questions"));
8438 let home = dir.path().join("home");
8439 let mut state = run_state(RunStatus::Blocked);
8440 state.id = "20260101-000000-act3".to_owned();
8441 state.save_under(&home).unwrap();
8442 let mut t = held_task_with(&state.id);
8443 queue.put(&mut t).unwrap();
8444 let mut q = action_question(&state.id, ask::ChoiceAction::Requeue);
8445 questions.put(&mut q).unwrap();
8446
8447 questions.beat(&q.id, crate::ask::WaiterKind::Asker);
8448 apply_choice_actions(&queue, &questions, &home);
8449 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8450
8451 std::fs::remove_file(questions.root().join(format!("{}.lease", q.id))).unwrap();
8452 apply_choice_actions(&queue, &questions, &home);
8453 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
8454 }
8455
8456 #[test]
8457 fn action_standing_tells_the_reasons_apart() {
8458 let mut t = held_task_with("r1");
8459 let q = action_question("r1", ask::ChoiceAction::Requeue);
8460 assert_eq!(action_standing(&t, &q), ActionStanding::Pending);
8461 let mut plain = q.clone();
8462 plain.actions.clear();
8463 assert_eq!(action_standing(&t, &plain), ActionStanding::NoAction);
8464
8465 t.mark_action_applied(&q.id);
8467 assert_eq!(action_standing(&t, &q), ActionStanding::Applied);
8468
8469 let mut t = held_task_with("r1");
8471 t.start("r2".to_owned());
8472 assert_eq!(action_standing(&t, &q), ActionStanding::Stale);
8473 let q2 = action_question("r2", ask::ChoiceAction::Requeue);
8475 assert_eq!(action_standing(&t, &q2), ActionStanding::Busy);
8476 t.status = TaskStatus::Blocked;
8477 assert_eq!(action_standing(&t, &q2), ActionStanding::Busy);
8478
8479 let mut c = action_question("r0", ask::ChoiceAction::Requeue);
8481 c.node = crate::conduct::NODE.to_owned();
8482 assert_eq!(action_standing(&t, &c), ActionStanding::Busy);
8483 assert!(!ActionStanding::Busy.daemon_owns());
8484 }
8485
8486 #[test]
8487 fn a_running_task_waits_and_a_dead_asker_with_a_cwd_is_still_actioned() {
8488 let dir = tempfile::tempdir().unwrap();
8489 let queue = Queue::at(dir.path().join("queue"));
8490 let questions = Questions::at(dir.path().join("questions"));
8491 let home = dir.path().join("home");
8492 let mut state = run_state(RunStatus::Blocked);
8493 state.id = "20260101-000000-act4".to_owned();
8494 state.save_under(&home).unwrap();
8495 let mut t = held_task_with(&state.id);
8496 t.status = TaskStatus::Running;
8497 queue.put(&mut t).unwrap();
8498 let mut q = action_question(&state.id, ask::ChoiceAction::Done);
8499 q.cwd = Some(dir.path().display().to_string());
8500 questions.put(&mut q).unwrap();
8501
8502 let running = queue.get(&t.id).unwrap();
8505 assert_eq!(action_standing(&running, &q), ActionStanding::Busy);
8506 assert!(matches!(
8507 crate::waiter::decide_owned(&q, None, false, 86_400, Timestamp::now(), false),
8508 crate::waiter::Action::Deliver(_)
8509 ));
8510 apply_choice_actions(&queue, &questions, &home);
8512 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
8513 let mut back = queue.get(&t.id).unwrap();
8514 back.hold_machine(Some("later".into()));
8515 queue.put(&mut back).unwrap();
8516 apply_choice_actions(&queue, &questions, &home);
8517 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Done);
8518 let mut again = queue.get(&t.id).unwrap();
8519 again.hold_machine(Some("again".into()));
8520 queue.put(&mut again).unwrap();
8521 apply_choice_actions(&queue, &questions, &home);
8522 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8523 }
8524}