1use std::path::{Path, PathBuf};
50use std::sync::Arc;
51use std::sync::atomic::{AtomicBool, Ordering};
52use std::sync::{Mutex, MutexGuard};
53use std::time::Duration;
54
55use anyhow::{Context, Result, bail};
56use jiff::Timestamp;
57use serde::{Deserialize, Serialize};
58use tokio::sync::Notify;
59
60use crate::ask::{self, Questions};
61use crate::clean;
62use crate::conduct::Conductor;
63use crate::config::{Config, MergeMode};
64use crate::graph::Runner;
65use crate::land;
66use crate::notices::{self, Link, Notice};
67use crate::queue::{Queue, Task, TaskStatus};
68use crate::run::{Liveness, QuotaLoss, RunState, RunStatus};
69use crate::triage;
70
71pub const SCHEMA: u32 = 1;
73
74pub const HEARTBEAT: Duration = Duration::from_secs(5);
78
79pub const STALE_SECS: i64 = 30;
88
89pub const POLL: Duration = Duration::from_secs(5);
91
92pub const STALE_CLAIM: Duration = Duration::from_secs(6 * 60 * 60);
96
97pub const STALLED_RUNNING: Duration = Duration::from_secs(30 * 60);
116
117#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
119#[serde(default)]
120pub struct Current {
121 pub task: String,
123 pub run: String,
125}
126
127#[derive(Debug, Clone, Serialize, Deserialize)]
134pub struct Status {
135 pub schema: u32,
137 pub pid: u32,
139 pub started_at: Timestamp,
141 pub updated_at: Timestamp,
143 pub idle: bool,
145 pub current: Vec<Current>,
151 pub completed: usize,
153 pub polls: u64,
155}
156
157impl Status {
158 #[must_use]
160 pub fn new() -> Self {
161 let now = Timestamp::now();
162 Self {
163 schema: SCHEMA,
164 pid: std::process::id(),
165 started_at: now,
166 updated_at: now,
167 idle: true,
168 current: Vec::new(),
169 completed: 0,
170 polls: 0,
171 }
172 }
173}
174
175impl Default for Status {
176 fn default() -> Self {
177 Self::new()
178 }
179}
180
181#[derive(Debug, Clone)]
183pub struct Opts {
184 pub repo: PathBuf,
186 pub config: Option<PathBuf>,
188 pub poll: Duration,
190 pub max_attempts: usize,
192 pub once: bool,
194 pub merge: Option<String>,
196 pub worktrees_root: Option<PathBuf>,
205}
206
207impl Default for Opts {
208 fn default() -> Self {
209 Self {
210 repo: PathBuf::from("."),
211 config: None,
212 poll: POLL,
213 max_attempts: 2,
214 once: false,
215 merge: None,
216 worktrees_root: None,
217 }
218 }
219}
220
221fn max_concurrent(n: usize) -> usize {
226 n.max(1)
227}
228
229#[must_use]
231pub fn status_path() -> PathBuf {
232 crate::run::home().join("daemon.json")
233}
234
235pub fn write_status(status: &Status) -> Result<()> {
237 write_status_to(&status_path(), status)
238}
239
240pub fn write_status_to(path: &Path, status: &Status) -> Result<()> {
245 if let Some(parent) = path.parent() {
246 std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
247 }
248 let body = serde_json::to_string_pretty(status).context("serialize daemon status")?;
249 let tmp = path.with_extension("json.tmp");
250 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
251 std::fs::rename(&tmp, path).with_context(|| format!("replace {}", path.display()))?;
252 Ok(())
253}
254
255pub fn clear_status() {
258 clear_status_at(&status_path());
259}
260
261fn clear_status_at(path: &Path) {
265 let _ = std::fs::remove_file(path);
266}
267
268#[derive(Debug, Clone, Default)]
281pub struct Stop {
282 stopped: Arc<AtomicBool>,
286 busy: Arc<std::sync::atomic::AtomicUsize>,
292 wake: Arc<Notify>,
296 pause: crate::graph::Pause,
299}
300
301impl Stop {
302 #[must_use]
304 pub fn new() -> Self {
305 Self::default()
306 }
307
308 pub fn stop(&self) {
311 self.stopped.store(true, Ordering::SeqCst);
312 self.wake.notify_one();
316 }
317
318 #[must_use]
320 pub fn stopped(&self) -> bool {
321 self.stopped.load(Ordering::SeqCst)
322 }
323
324 #[must_use]
332 pub fn finishing(&self) -> bool {
333 self.stopped() && self.busy_now()
334 }
335
336 pub fn park(&self) {
347 self.pause.park();
348 self.stop();
349 }
350
351 #[must_use]
353 pub fn parking(&self) -> bool {
354 self.pause.parked()
355 }
356
357 #[must_use]
359 pub fn pause(&self) -> crate::graph::Pause {
360 self.pause.clone()
361 }
362
363 #[must_use]
369 pub fn busy_now(&self) -> bool {
370 self.busy.load(Ordering::SeqCst) > 0
371 }
372
373 fn enter(&self) {
375 self.busy.fetch_add(1, Ordering::SeqCst);
376 }
377
378 fn exit(&self) {
381 self.busy.fetch_sub(1, Ordering::SeqCst);
382 }
383
384 async fn idle(&self, poll: Duration) {
386 tokio::select! {
387 () = tokio::time::sleep(poll) => {}
388 () = self.wake.notified() => {}
389 }
390 }
391}
392
393#[derive(Debug, Clone, Default, Deserialize)]
400#[serde(default)]
401pub struct Reading {
402 pub schema: u32,
404 pub pid: Option<u32>,
406 pub started_at: Option<Timestamp>,
408 pub updated_at: Option<Timestamp>,
410 pub idle: bool,
412 #[serde(deserialize_with = "de_current")]
424 pub current: Vec<Current>,
425 pub completed: u64,
427 pub polls: u64,
429}
430
431fn de_current<'de, D>(deserializer: D) -> std::result::Result<Vec<Current>, D::Error>
434where
435 D: serde::Deserializer<'de>,
436{
437 #[derive(Deserialize)]
438 #[serde(untagged)]
439 enum Shape {
440 Many(Vec<Current>),
441 One(Current),
442 }
443 Ok(
444 Option::<Shape>::deserialize(deserializer)?.map_or_else(Vec::new, |shape| match shape {
445 Shape::Many(v) => v,
446 Shape::One(c) => vec![c],
447 }),
448 )
449}
450
451impl Reading {
452 #[must_use]
455 pub fn age_secs(&self, now: Timestamp) -> Option<i64> {
456 self.updated_at
457 .map(|at| (now.as_second() - at.as_second()).max(0))
458 }
459
460 #[must_use]
464 pub fn running(&self, now: Timestamp) -> bool {
465 self.age_secs(now).is_some_and(|secs| secs <= STALE_SECS)
466 }
467}
468
469#[must_use]
476pub fn read_status(home: &Path) -> Option<Reading> {
477 let body = std::fs::read_to_string(home.join("daemon.json")).ok()?;
478 serde_json::from_str(&body).ok()
479}
480
481#[must_use]
493pub fn current_work(home: &Path, now: Timestamp) -> Vec<Current> {
494 read_status(home)
495 .filter(|reading| reading.running(now))
496 .map(|reading| reading.current)
497 .unwrap_or_default()
498}
499
500#[must_use]
502pub fn is_working_on(home: &Path, run: &str, now: Timestamp) -> bool {
503 current_work(home, now).iter().any(|c| c.run == run)
504}
505
506#[must_use]
516pub fn is_working_on_short(home: &Path, short: &str, now: Timestamp) -> bool {
517 current_work(home, now)
518 .iter()
519 .any(|c| crate::run::short_of(&c.run) == short)
520}
521
522#[must_use]
524pub fn is_working_on_task(home: &Path, task: &str, now: Timestamp) -> bool {
525 current_work(home, now).iter().any(|c| c.task == task)
526}
527
528pub fn sweep_stale_claims(queue: &Queue, older_than: Duration) -> Vec<String> {
572 sweep_stale_claims_with(queue, older_than, crate::proc::pid_alive)
573}
574
575fn sweep_stale_claims_with<F>(queue: &Queue, older_than: Duration, pid_alive: F) -> Vec<String>
579where
580 F: Fn(u32) -> bool,
581{
582 let this_process = std::process::id();
583 let mut swept: Vec<String> = std::fs::read_dir(queue.root())
584 .into_iter()
585 .flatten()
586 .flatten()
587 .map(|e| e.path())
588 .filter(|p| p.extension().is_some_and(|x| x == "lock"))
589 .filter(|p| {
590 match std::fs::read_to_string(p)
591 .ok()
592 .and_then(|body| body.trim().parse::<u32>().ok())
593 {
594 Some(pid) if pid == this_process => false,
598 Some(pid) => !pid_alive(pid),
599 None => p
600 .metadata()
601 .and_then(|m| m.modified())
602 .and_then(|t| t.elapsed().map_err(std::io::Error::other))
603 .is_ok_and(|age| age >= older_than),
604 }
605 })
606 .filter(|p| std::fs::remove_file(p).is_ok())
607 .filter_map(|p| {
608 p.file_stem()
609 .and_then(|s| s.to_str())
610 .map(std::borrow::ToOwned::to_owned)
611 })
612 .collect();
613 swept.sort_unstable();
614 swept
615}
616
617fn is_stalled(task: &Task, home: &Path, now: Timestamp) -> bool {
622 task.status == TaskStatus::Running
623 && (now.as_second() - task.updated_at.as_second()) >= STALLED_RUNNING.as_secs() as i64
624 && !is_working_on_task(home, &task.id, now)
625}
626
627fn stalled_tasks(queue: &Queue, home: &Path, now: Timestamp) -> Vec<Task> {
630 queue
631 .list()
632 .into_iter()
633 .filter(|t| is_stalled(t, home, now))
634 .collect()
635}
636
637fn queued_tasks(queue: &Queue) -> Vec<Task> {
643 queue
644 .list()
645 .into_iter()
646 .filter(|t| t.status == TaskStatus::Queued)
647 .collect()
648}
649
650fn finished_tasks(queue: &Queue) -> Vec<Task> {
653 queue
654 .list()
655 .into_iter()
656 .filter(|t| matches!(t.status, TaskStatus::Failed | TaskStatus::Held))
657 .collect()
658}
659
660fn resolve_blockers(queue: &Queue, questions: &Questions) {
677 for listed in queue.list() {
678 if listed.status != TaskStatus::Blocked || listed.blocked_by.is_empty() {
679 continue;
680 }
681 let Ok(_claim) = queue.claim(&listed.id) else {
682 continue;
683 };
684 let Ok(mut task) = queue.get(&listed.id) else {
685 continue;
686 };
687 if task.status != TaskStatus::Blocked {
688 continue;
689 }
690 let deleted = queue.apply_deleted_blockers(&mut task);
694 if !deleted.is_empty() {
695 record(queue, &mut task);
696 for id in &deleted {
697 queue.note_dependency_deleted(&task, id);
698 }
699 if task.status != TaskStatus::Blocked {
700 continue;
701 }
702 }
703 let missing = crate::queue::missing_blockers(queue, questions, &task.blocked_by);
704 if !missing.is_empty() {
705 let language = language_of(&task, Path::new("."));
706 task.hold_machine(Some(crate::queue::missing_blocker_hold_reason_in(
707 &task.blocked_by,
708 &missing,
709 &language,
710 )));
711 record(queue, &mut task);
712 continue;
713 }
714 if let Some(q) = task.blocked_by.iter().find_map(|id| {
717 questions.get(id).ok().filter(|q| {
718 q.node == crate::conduct::NODE && q.status == ask::QuestionStatus::Abandoned
719 })
720 }) {
721 let language = language_of(&task, Path::new("."));
722 task.hold_machine(Some(unanswered_question_hold_reason(&q, &language)));
723 record(queue, &mut task);
724 continue;
725 }
726 let mut changed = false;
727 for id in task.blocked_by.clone() {
728 if let Ok(dep) = queue.get(&id) {
729 if dep.status == TaskStatus::Done {
730 task.unblock(&id);
731 changed = true;
732 }
733 continue;
734 }
735 if let Ok(q) = questions.get(&id)
736 && q.status == ask::QuestionStatus::Answered
737 {
738 let answer = match &q.answer {
739 Some(ask::Answer::Choice(c) | ask::Answer::Text(c)) => c.clone(),
740 None => String::new(),
741 };
742 task.record_answer(q.summary.clone(), answer);
743 task.unblock(&id);
744 changed = true;
745 }
746 }
747 if changed {
748 record(queue, &mut task);
749 }
750 }
751}
752
753fn unanswered_question_hold_reason(q: &ask::Question, language: &str) -> String {
755 if crate::lang::is_japanese(language) {
756 format!(
757 "質問 {} 「{}」 に期限内の回答がなく、取り下げられました - `magi task triage` を参照",
758 q.short(),
759 q.summary
760 )
761 } else {
762 format!(
763 "question {} \"{}\" went unanswered and was abandoned - see `magi task triage`",
764 q.short(),
765 q.summary
766 )
767 }
768}
769
770#[derive(Debug, Clone, Copy, PartialEq, Eq)]
773pub(crate) enum ActionStanding {
774 NoAction,
776 Pending,
778 Applied,
781 Stale,
783 Busy,
785}
786
787impl ActionStanding {
788 pub(crate) fn daemon_owns(self) -> bool {
790 matches!(self, Self::Pending | Self::Applied | Self::Stale)
791 }
792}
793
794pub(crate) fn action_standing(task: &Task, q: &ask::Question) -> ActionStanding {
798 if q.chosen_action().is_none() {
799 ActionStanding::NoAction
800 } else if task.action_applied(&q.id) {
801 ActionStanding::Applied
802 } else if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
803 ActionStanding::Stale
804 } else if matches!(
805 task.status,
806 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
807 ) {
808 ActionStanding::Busy
809 } else {
810 ActionStanding::Pending
811 }
812}
813
814#[derive(Debug, Clone, PartialEq, Eq)]
816enum ActionDecision {
817 Skip,
819 Resume(String),
821 Requeue,
823 Done,
825 Stale,
828 Refuse(String),
831}
832
833fn decide_action<F>(task: &Task, q: &ask::Question, p: &Phrases, load: F) -> ActionDecision
841where
842 F: FnOnce(&str) -> Result<RunState>,
843{
844 let Some(action) = q.chosen_action() else {
845 return ActionDecision::Skip;
846 };
847 if task.action_applied(&q.id)
848 || matches!(
849 task.status,
850 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
851 )
852 {
853 return ActionDecision::Skip;
854 }
855 if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
858 return ActionDecision::Stale;
859 }
860 match action {
861 ask::ChoiceAction::Requeue => ActionDecision::Requeue,
862 ask::ChoiceAction::Done => ActionDecision::Done,
863 ask::ChoiceAction::Resume { run } => {
864 if task.runs.last() != Some(run) {
865 return ActionDecision::Refuse((p.resume_not_latest)(
866 q.short(),
867 ask::short_id(run),
868 ));
869 }
870 match load(run) {
871 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
872 ActionDecision::Resume(run.clone())
873 }
874 Ok(_) => ActionDecision::Refuse((p.resume_cannot_progress)(
875 q.short(),
876 ask::short_id(run),
877 )),
878 Err(e) => ActionDecision::Refuse((p.resume_unreadable)(
879 q.short(),
880 ask::short_id(run),
881 &format!("{e:#}"),
882 )),
883 }
884 }
885 }
886}
887
888struct Phrases {
902 graph_stopped: fn(&str, &str) -> String,
904 quorum_lost: &'static str,
905 quota_took_out: &'static str,
907 run_ended: &'static str,
909 waiting_for_answer: &'static str,
911 recovered_running: &'static str,
912 no_run_to_recover: &'static str,
913 could_not_start: &'static str,
914 resume_not_latest: fn(&str, &str) -> String,
915 resume_cannot_progress: fn(&str, &str) -> String,
916 resume_unreadable: fn(&str, &str, &str) -> String,
917 handover_refused: fn(&crate::handover::Refused) -> String,
920 handover_hint: &'static str,
922 reviewers_never_answered: fn(usize, usize) -> String,
925}
926
927const PHRASES_EN: Phrases = Phrases {
928 graph_stopped: |status, detail| {
929 format!("the graph stopped at `{status}` without reaching a terminal status: {detail}")
930 },
931 quorum_lost: "the judging panel lost its quorum",
932 quota_took_out: "; quota took out ",
933 run_ended: "run ended ",
934 waiting_for_answer: " - waiting for operator answer to question ",
935 recovered_running: "recovered a `running` task whose daemon never recorded the outcome: ",
936 no_run_to_recover: "task was `running` with no live daemon and no readable \
937 run to recover; held for a human to check what happened",
938 could_not_start: "could not start the run: ",
939 resume_not_latest: |q, run| {
940 format!("question {q} asked to resume run {run}, which is not this task's latest run")
941 },
942 resume_cannot_progress: |q, run| {
943 format!("question {q} asked to resume run {run}, which cannot make progress")
944 },
945 resume_unreadable: |q, run, e| {
946 format!("question {q} asked to resume run {run}, which could not be read: {e}")
947 },
948 handover_refused: |r| r.to_string(),
949 handover_hint: " (clean up the other worktree, then release the task from the queue)",
950 reviewers_never_answered: |missing, rounds| {
951 format!(
952 "{missing} reviewer seat(s) never answered after {rounds} rounds; \
953 refusing to call it clean"
954 )
955 },
956};
957
958const PHRASES_JA: Phrases = Phrases {
959 graph_stopped: |status, detail| {
960 format!("グラフが終端状態に達しないまま `{status}` で停止しました: {detail}")
961 },
962 quorum_lost: "審査パネルが定足数を失いました",
963 quota_took_out: "。クォータで脱落: ",
964 run_ended: "run 終了: ",
965 waiting_for_answer: " - オペレーターの回答待ち: 質問 ",
966 recovered_running: "daemon が結果を記録しないまま `running` だったタスクを回収しました: ",
967 no_run_to_recover: "タスクは `running` でしたが、生きた daemon も回収できる run も見つかりません。\
968 何が起きたか人が確認するため保留にしました",
969 could_not_start: "run を開始できませんでした: ",
970 resume_not_latest: |q, run| {
971 format!(
972 "質問 {q} は run {run} の再開を求めましたが、これはタスクの最新の run ではありません"
973 )
974 },
975 resume_cannot_progress: |q, run| {
976 format!("質問 {q} は run {run} の再開を求めましたが、これは進行できません")
977 },
978 resume_unreadable: |q, run, e| {
979 format!("質問 {q} は run {run} の再開を求めましたが、読み込めませんでした: {e}")
980 },
981 handover_refused: |r| {
982 use crate::handover::Refused;
983 match r {
984 Refused::Foreign { branch, path, why } => format!(
985 "ブランチ `{branch}` は {path} にチェックアウトされており、magi は自動では\
986 削除しません。不要なら `git worktree remove` でその worktree を削除して\
987 から、やり直してください(詳細: {why})"
988 ),
989 Refused::Unsafe { branch, path, why } => format!(
990 "ブランチ `{branch}` は {path} にチェックアウトされています。そこの作業を\
991 コミットか破棄したうえで `git worktree remove` で worktree を削除する\
992 か、run を破棄してよいと伝えてから、やり直してください(詳細: {why})"
993 ),
994 Refused::ReleaseFailed { branch, path, run } => format!(
995 "ブランチ `{branch}` は {path} で run {run} がチェックアウトしており、その\
996 worktree の解放に失敗したか、変更が見つかりました(未コミットの変更が\
997 ある worktree は git が削除を拒否します)。worktree はそのまま残しました"
998 ),
999 }
1000 },
1001 handover_hint: "(他の worktree を片付けてから、タスクをキューから解放してください)",
1002 reviewers_never_answered: |missing, rounds| {
1003 format!(
1004 "{missing} 席のレビュアーが {rounds} ラウンドの間に一度も回答しなかったため、\
1005 クリーンとは認めません"
1006 )
1007 },
1008};
1009
1010fn phrases(language: &str) -> &'static Phrases {
1011 if crate::lang::is_japanese(language) {
1012 &PHRASES_JA
1013 } else {
1014 &PHRASES_EN
1015 }
1016}
1017
1018fn language_of(task: &Task, fallback: &Path) -> String {
1022 crate::lang::of_repo(&repo_for(task, fallback))
1023}
1024
1025pub(crate) fn task_of_question<'a>(tasks: &'a [Task], q: &ask::Question) -> Option<&'a Task> {
1029 if q.node == crate::conduct::NODE {
1030 return tasks.iter().find(|t| t.id == q.run);
1031 }
1032 tasks.iter().find(|t| t.runs.contains(&q.run))
1033}
1034
1035fn apply_choice_actions(queue: &Queue, questions: &Questions, home: &Path) {
1044 let tasks = queue.list();
1045 for q in questions.list() {
1046 if q.chosen_action().is_none() {
1047 continue;
1048 }
1049 let Some(listed) = task_of_question(&tasks, &q) else {
1050 continue;
1051 };
1052 if listed.action_applied(&q.id) {
1053 continue;
1054 }
1055 let Ok(_claim) = queue.claim(&listed.id) else {
1056 continue;
1057 };
1058 let Ok(mut task) = queue.get(&listed.id) else {
1059 continue;
1060 };
1061 let language = language_of(&task, Path::new("."));
1062 let decision = decide_action(&task, &q, phrases(&language), |id| {
1063 RunState::load_under(id, home)
1064 });
1065 let ran = matches!(
1066 decision,
1067 ActionDecision::Resume(_) | ActionDecision::Requeue | ActionDecision::Done
1068 );
1069 if ran && asker_may_still_read(questions, &q, home) {
1070 continue;
1079 }
1080 let done = matches!(decision, ActionDecision::Done);
1081 match decision {
1082 ActionDecision::Skip => continue,
1083 ActionDecision::Resume(run) => {
1084 task.release();
1085 task.resume_override = Some(crate::queue::OperatorResume {
1086 question_id: q.id.clone(),
1087 at: Timestamp::now(),
1088 conductor_rehold: None,
1089 forced: true,
1090 pinned_run: Some(run),
1091 });
1092 }
1093 ActionDecision::Stale => {}
1094 ActionDecision::Requeue => task.requeue(),
1095 ActionDecision::Done => task.succeed(),
1096 ActionDecision::Refuse(why) => {
1097 task.hold_machine(Some(why));
1098 notices::raise(
1099 Notice::warn(
1100 &format!("action:{}", q.id),
1101 "An answer asked the daemon to resume a run that cannot be resumed; the task stays held.",
1102 )
1103 .link(Link::Task {
1104 id: task.id.clone(),
1105 }),
1106 );
1107 }
1108 }
1109 task.mark_action_applied(&q.id);
1110 record(queue, &mut task);
1113 if done && queue.get(&task.id).is_ok_and(|t| t.action_applied(&q.id)) {
1114 supersede_prior_runs(&task, home);
1115 }
1116 let _ = questions.update(&q.id, |r| {
1119 r.answer_delivered = true;
1120 Ok(())
1121 });
1122 }
1123}
1124
1125fn asker_may_still_read(questions: &Questions, q: &ask::Question, home: &Path) -> bool {
1129 let now = Timestamp::now();
1130 if questions.read_lease(&q.id).is_some_and(|l| l.fresh(now)) {
1131 return true;
1132 }
1133 q.node != crate::conduct::NODE
1134 && RunState::load_under(&q.run, home).is_ok_and(|s| {
1135 s.seats_active()
1136 .any(|(k, a)| *k == q.seat && a.remaining_secs(now) > 0)
1137 })
1138}
1139
1140fn reconcile_task_questions(queue: &Queue, questions: &Questions) {
1151 let tasks = queue.list();
1152 let by_id: std::collections::BTreeMap<&str, &Task> =
1153 tasks.iter().map(|t| (t.id.as_str(), t)).collect();
1154 let referenced: std::collections::BTreeSet<&str> = tasks
1155 .iter()
1156 .flat_map(|task| task.blocked_by.iter().map(String::as_str))
1157 .collect();
1158
1159 for mut question in questions.list() {
1160 if !question.status.open() || question.node != crate::conduct::NODE {
1161 continue;
1162 }
1163 if referenced.contains(question.id.as_str()) {
1166 continue;
1167 }
1168 let Some(task) = by_id.get(question.run.as_str()) else {
1169 continue;
1170 };
1171 question.abandon(format!(
1172 "task {} no longer waits for this answer",
1173 task.short()
1174 ));
1175 if let Err(e) = questions.put(&mut question) {
1176 tracing::warn!(
1177 "could not retire question {} for task {}: {e:#}",
1178 question.short(),
1179 task.short()
1180 );
1181 }
1182 }
1183}
1184
1185#[derive(Debug, Clone, Copy)]
1192pub struct Verdict {
1193 pub status: RunStatus,
1195 pub left_pr: bool,
1197 pub quota_hit: bool,
1199 pub parked: bool,
1201 pub no_viable_candidates: bool,
1208}
1209
1210pub fn settle(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize) {
1267 settle_in(task, verdict, detail, max_attempts, &PHRASES_EN)
1268}
1269
1270fn settle_in(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize, p: &Phrases) {
1272 if verdict.parked {
1278 task.stall(detail);
1279 return;
1280 }
1281 match verdict.status {
1282 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1283 RunStatus::AlreadyInBase => task.already_landed(detail),
1284 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1285 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1286 task.stall(detail)
1287 }
1288 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1289 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1290 RunStatus::Blocked => task.fail(detail, max_attempts),
1291 RunStatus::VerifiedNoop => task.handed_off(detail),
1292 other => task.fail((p.graph_stopped)(label(other), detail), max_attempts),
1293 }
1294}
1295
1296pub fn supersede_prior_runs(task: &Task, home: &Path) {
1345 let now = Timestamp::now();
1346 let last_run_succeeded = task
1353 .runs
1354 .last()
1355 .and_then(|id| RunState::load_under(id, home).ok())
1356 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1357 for id in task.superseded_attempts(last_run_succeeded) {
1358 let mut state = match RunState::load_under(id, home) {
1359 Ok(s) => s,
1360 Err(e) => {
1361 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1362 continue;
1363 }
1364 };
1365 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1366 continue;
1367 }
1368 let daemon_claims = is_working_on(home, id, now);
1369 if state.liveness(daemon_claims) == Liveness::Live {
1370 continue;
1371 }
1372 state.status = RunStatus::Superseded;
1373 if let Err(e) = state.save_under(home) {
1374 tracing::warn!("could not mark run {id} superseded: {e:#}");
1375 }
1376 }
1377}
1378
1379fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1401 for task in queue.list() {
1402 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1403 continue;
1404 }
1405 supersede_prior_runs(&task, home);
1406 }
1407}
1408
1409fn settle_and_diagnose(
1416 task: &mut Task,
1417 verdict: Verdict,
1418 detail: &str,
1419 max_attempts: usize,
1420 state: &RunState,
1421) {
1422 let p = phrases(&state.config.graph.language);
1423 settle_in(task, verdict, detail, max_attempts, p);
1424 if task.status == TaskStatus::Held {
1425 task.diagnostic = diagnostic(state);
1426 note_open_question(task, &state.id, p);
1427 }
1428}
1429
1430fn note_open_question(task: &mut Task, run: &str, p: &Phrases) {
1445 let Some(home) = crate::run::try_home() else {
1446 return;
1447 };
1448 let open = Questions::at(home.join("questions")).open_for(run);
1449 let Some(q) = open.first() else {
1450 return;
1451 };
1452 let base = task.hold_reason.clone().unwrap_or_default();
1453 task.hold_reason = Some(format!("{base}{}{}", p.waiting_for_answer, q.short()));
1454}
1455
1456fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize, language: &str) {
1467 match last_run {
1468 Some(state) => {
1469 let verdict = Verdict {
1470 status: state.status,
1471 left_pr: state.pr.is_some(),
1472 quota_hit: !state.quota.is_empty(),
1473 parked: state.parked,
1474 no_viable_candidates: state.viable().is_empty(),
1475 };
1476 let detail = format!(
1477 "{}{}",
1478 phrases(&state.config.graph.language).recovered_running,
1479 describe(&state)
1480 );
1481 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1482 }
1483 None => {
1484 let why = phrases(language).no_run_to_recover;
1487 task.last_error = Some(why.to_owned());
1488 task.hold_machine(Some(why.to_owned()));
1491 }
1492 }
1493}
1494
1495fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1517 let mut reclaimed = Vec::new();
1518 for listed in queue.list() {
1519 if listed.status != TaskStatus::Running {
1520 continue;
1521 }
1522 let Ok(_claim) = queue.claim(&listed.id) else {
1523 continue;
1524 };
1525 let Ok(mut task) = queue.get(&listed.id) else {
1529 continue;
1530 };
1531 if task.status != TaskStatus::Running {
1532 continue;
1533 }
1534 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1535 if let Some(state) = &last_run
1545 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1546 {
1547 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1548 }
1549 let language = if last_run.is_none() {
1550 language_of(&task, Path::new("."))
1551 } else {
1552 String::new()
1553 };
1554 reclaim(&mut task, last_run, max_attempts, &language);
1555 if task.status == TaskStatus::Done {
1556 supersede_prior_runs(&task, &crate::run::home());
1557 }
1558 record(queue, &mut task);
1559 reclaimed.push(task.id.clone());
1560 }
1561 reclaimed
1562}
1563
1564fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1596 reclaim_abandoned_runs_with(
1597 home,
1598 now,
1599 crate::proc::pid_status,
1600 crate::proc::process_started_at,
1601 )
1602}
1603
1604fn reclaim_abandoned_runs_with<F, G>(
1610 home: &Path,
1611 now: Timestamp,
1612 query: F,
1613 identity: G,
1614) -> Vec<String>
1615where
1616 F: Fn(u32) -> Option<bool>,
1617 G: Fn(u32) -> Option<String>,
1618{
1619 let mut abandoned = Vec::new();
1620 for entry in std::fs::read_dir(home.join("runs"))
1621 .into_iter()
1622 .flatten()
1623 .flatten()
1624 {
1625 let id = entry.file_name().to_string_lossy().into_owned();
1626 if !crate::run::is_run_id(&id) {
1627 continue;
1628 }
1629 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1637 continue;
1638 };
1639 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1640 continue;
1641 };
1642 if state.status.done() || !state.active_all_overrun(now) {
1643 continue;
1644 }
1645 let daemon_claims = is_working_on(home, &id, now);
1657 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1658 continue;
1659 }
1660 state.abandon("daemon");
1661 if let Err(e) = state.save_under(home) {
1662 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1663 continue;
1664 }
1665 if let Some(notice) = notices::run_ended(&state) {
1668 notices::raise_in(home, notice);
1669 }
1670 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1678 tracing::warn!("abandon questions for {id}: {e:#}");
1679 }
1680 abandoned.push(id);
1681 }
1682 abandoned
1683}
1684
1685pub async fn serve(opts: Opts) -> Result<()> {
1691 serve_until(opts, Stop::new()).await
1692}
1693
1694pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1711 let signal = {
1712 let stop = stop.clone();
1713 tokio::spawn(async move {
1714 if tokio::signal::ctrl_c().await.is_ok() {
1715 stop.stop();
1716 tracing::info!("shutdown requested; a run in flight will be finished first");
1717 }
1718 })
1719 };
1720
1721 let worktrees_root = opts
1722 .worktrees_root
1723 .clone()
1724 .unwrap_or_else(crate::run::default_worktree_root);
1725 let outcome = drive(
1726 &opts,
1727 &Queue::open(),
1728 &status_path(),
1729 &crate::run::home(),
1730 &worktrees_root,
1731 &stop,
1732 )
1733 .await;
1734
1735 signal.abort();
1736 outcome
1737}
1738
1739async fn drive(
1752 opts: &Opts,
1753 queue: &Queue,
1754 status_file: &Path,
1755 home: &Path,
1756 worktrees_root: &Path,
1757 stop: &Stop,
1758) -> Result<()> {
1759 let status = Arc::new(Mutex::new(Status::new()));
1767 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1768 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1769
1770 let daemon_cfg = prepare(&opts.repo, opts)
1776 .map(|c| c.daemon)
1777 .unwrap_or_default();
1778 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1779
1780 let waiter = tokio::spawn(crate::waiter::run(
1785 crate::waiter::Waiter::new(
1786 crate::ask::Questions::at(home.join("questions")),
1787 home.to_path_buf(),
1788 prepare(&opts.repo, opts).ok(),
1789 ),
1790 stop.clone(),
1791 ));
1792
1793 let deputies = tokio::spawn(crate::deputy::run(
1797 crate::deputy::Deputies::new(
1798 crate::ask::Questions::at(home.join("questions")),
1799 home.to_path_buf(),
1800 prepare(&opts.repo, opts).ok(),
1801 opts.repo.clone(),
1802 daemon_cfg.max_deputies,
1803 {
1804 let stop = stop.clone();
1805 Arc::new(move || stop.parking())
1806 },
1807 ),
1808 stop.clone(),
1809 ));
1810
1811 let fetcher = tokio::spawn(fetch_loop(opts.repo.clone(), opts.clone(), stop.clone()));
1814
1815 tracing::info!(
1816 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1817 queue.root().display(),
1818 opts.poll.as_secs(),
1819 opts.max_attempts,
1820 concurrency,
1821 if daemon_cfg.pause_for_interrupts {
1822 ", interrupts enabled"
1823 } else {
1824 ""
1825 }
1826 );
1827
1828 janitor(&opts.repo, opts, home, worktrees_root).await;
1831 resweep_superseded_attempts(queue, home);
1832
1833 let outcome = poll(
1834 opts,
1835 queue,
1836 &status,
1837 home,
1838 worktrees_root,
1839 stop,
1840 DispatchLimits {
1841 max_concurrent: concurrency,
1842 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1843 },
1844 )
1845 .await;
1846
1847 beat.abort();
1848 waiter.abort();
1849 deputies.abort();
1850 fetcher.abort();
1851 clear_status_at(status_file);
1852 outcome
1853}
1854
1855const FETCH_TIMEOUT: Duration = Duration::from_secs(30);
1857
1858async fn fetch_loop(repo: PathBuf, opts: Opts, stop: Stop) {
1862 while !stop.stopped() {
1863 let (interval, roots) = match prepare(&repo, &opts) {
1864 Ok(c) => (c.repos.fetch_interval, c.repos.roots),
1865 Err(_) => (0, Vec::new()),
1866 };
1867 if interval > 0 && !roots.is_empty() {
1868 let r = crate::clean::fetch_origins(&roots, FETCH_TIMEOUT, || stop.stopped()).await;
1869 tracing::debug!("fetch origins: {r:?}");
1870 }
1871 let wait = if interval > 0 { interval } else { 60 };
1873 let mut slept = 0;
1874 while slept < wait && !stop.stopped() {
1875 tokio::time::sleep(Duration::from_secs(1)).await;
1876 slept += 1;
1877 }
1878 }
1879}
1880
1881async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1887 loop {
1888 tokio::time::sleep(HEARTBEAT).await;
1889 let snapshot = {
1890 let mut guard = lock(&status);
1891 guard.updated_at = Timestamp::now();
1892 guard.clone()
1893 };
1894 if let Err(e) = write_status_to(&path, &snapshot) {
1895 tracing::warn!("could not refresh the daemon status file: {e:#}");
1898 }
1899 }
1900}
1901
1902#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1905enum LandResume {
1906 NotLanding,
1909 StillWaiting,
1914 Ready,
1918}
1919
1920fn land_resume_state(task: &Task) -> LandResume {
1924 let Some(run_id) = task.runs.last() else {
1925 return LandResume::NotLanding;
1926 };
1927 let Ok(state) = RunState::load(run_id) else {
1928 return LandResume::NotLanding;
1929 };
1930 if state.status != RunStatus::Landing || !state.parked {
1931 return LandResume::NotLanding;
1932 }
1933 let store = ask::Questions::open();
1934 let waiting = store
1935 .list()
1936 .into_iter()
1937 .filter(|q| &q.run == run_id && q.node == land::APPROVAL_NODE)
1938 .max_by(|a, b| a.id.cmp(&b.id));
1939 let Some(mut q) = waiting else {
1940 return LandResume::Ready;
1941 };
1942 if !q.status.open() {
1943 return LandResume::Ready;
1944 }
1945 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
1952 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
1953 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
1954 q.abandon(format!(
1955 "no answer within {}s of asking",
1956 timeout.as_secs().max(1)
1957 ));
1958 if store.put(&mut q).is_ok() {
1961 return LandResume::Ready;
1962 }
1963 }
1964 LandResume::StillWaiting
1965}
1966
1967const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
1975
1976const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
1988
1989struct InFlightGuard<'a> {
2002 status: &'a Arc<Mutex<Status>>,
2003 stop: &'a Stop,
2004 task_id: &'a str,
2005}
2006
2007impl Drop for InFlightGuard<'_> {
2008 fn drop(&mut self) {
2009 lock(self.status).current.retain(|c| c.task != self.task_id);
2010 self.stop.exit();
2011 }
2012}
2013
2014#[derive(Debug, Clone, PartialEq, Eq)]
2036enum Interrupt {
2037 Idle,
2039 Parking {
2051 parked: Vec<String>,
2052 interrupt_task: String,
2053 },
2054 Running {
2062 parked: Vec<String>,
2063 interrupt_task: String,
2064 },
2065 Resuming { parked: Vec<String> },
2072}
2073
2074fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
2095 match state {
2096 Interrupt::Idle => {
2097 if in_flight.len() != 1 {
2109 return Interrupt::Idle;
2110 }
2111 match runnable.iter().find(|t| t.interrupt) {
2112 Some(t) => Interrupt::Parking {
2113 parked: in_flight.to_vec(),
2114 interrupt_task: t.id.clone(),
2115 },
2116 None => Interrupt::Idle,
2117 }
2118 }
2119 Interrupt::Parking {
2120 parked,
2121 interrupt_task,
2122 } => {
2123 if in_flight.iter().any(|id| parked.contains(id)) {
2124 Interrupt::Parking {
2126 parked,
2127 interrupt_task,
2128 }
2129 } else if in_flight.contains(&interrupt_task) {
2130 Interrupt::Running {
2131 parked,
2132 interrupt_task,
2133 }
2134 } else if runnable.iter().any(|t| t.id == interrupt_task) {
2135 Interrupt::Parking {
2139 parked,
2140 interrupt_task,
2141 }
2142 } else {
2143 Interrupt::Resuming { parked }
2148 }
2149 }
2150 Interrupt::Running {
2151 parked,
2152 interrupt_task,
2153 } => {
2154 if in_flight.contains(&interrupt_task) {
2155 Interrupt::Running {
2156 parked,
2157 interrupt_task,
2158 }
2159 } else {
2160 Interrupt::Resuming { parked }
2166 }
2167 }
2168 Interrupt::Resuming { parked } => {
2169 if in_flight.iter().any(|id| parked.contains(id)) {
2170 Interrupt::Idle
2176 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
2177 Interrupt::Resuming { parked }
2178 } else {
2179 Interrupt::Idle
2182 }
2183 }
2184 }
2185}
2186
2187fn advance_interrupt_tick(
2193 enabled: bool,
2194 state: Interrupt,
2195 in_flight: &[String],
2196 runnable: &[Task],
2197) -> Interrupt {
2198 if !enabled {
2199 return Interrupt::Idle;
2200 }
2201 advance_interrupt(state, in_flight, runnable)
2202}
2203
2204fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
2209 match state {
2210 Interrupt::Idle => candidates,
2211 Interrupt::Parking {
2212 parked,
2213 interrupt_task,
2214 } => {
2215 if in_flight.iter().any(|id| parked.contains(id)) {
2216 Vec::new()
2217 } else {
2218 candidates
2219 .into_iter()
2220 .filter(|t| &t.id == interrupt_task)
2221 .collect()
2222 }
2223 }
2224 Interrupt::Running { .. } => Vec::new(),
2225 Interrupt::Resuming { parked } => candidates
2233 .into_iter()
2234 .find(|t| parked.contains(&t.id))
2235 .into_iter()
2236 .collect(),
2237 }
2238}
2239
2240struct DispatchLimits {
2244 max_concurrent: usize,
2247 pause_for_interrupts: bool,
2249}
2250
2251#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2260enum PermitKind {
2261 None,
2266 Urgent,
2273 Ordinary,
2276}
2277
2278fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
2281 if priority {
2282 PermitKind::None
2283 } else if urgent {
2284 PermitKind::Urgent
2285 } else {
2286 PermitKind::Ordinary
2287 }
2288}
2289
2290async fn poll(
2309 opts: &Opts,
2310 queue: &Queue,
2311 status: &Arc<Mutex<Status>>,
2312 home: &Path,
2313 worktrees_root: &Path,
2314 stop: &Stop,
2315 limits: DispatchLimits,
2316) -> Result<()> {
2317 let DispatchLimits {
2318 max_concurrent,
2319 pause_for_interrupts,
2320 } = limits;
2321 let mut attempted: Vec<String> = Vec::new();
2326 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2327 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2335 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2341 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2342 let mut conductor = Conductor::new();
2343 let mut cache_last_checked: Option<Timestamp> = None;
2346 let mut interrupt = Interrupt::Idle;
2348 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2354 std::collections::HashMap::new();
2355
2356 while !stop.stopped() {
2357 lock(status).polls += 1;
2358
2359 while let Some(result) = inflight.try_join_next() {
2364 if let Err(e) = result {
2365 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2366 notices::raise(Notice::error(
2367 "loop:attempt",
2368 "A queued attempt ended abnormally; check the task it was running.",
2369 ));
2370 }
2371 }
2372
2373 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2374 if !swept.is_empty() {
2375 tracing::warn!(
2376 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2377 swept.len(),
2378 swept.join(", ")
2379 );
2380 }
2381 let now = Timestamp::now();
2386
2387 if !stop.busy_now() {
2392 maybe_prune_cache_between_runs(
2393 &opts.repo,
2394 opts,
2395 home,
2396 stop,
2397 &mut cache_last_checked,
2398 now,
2399 )
2400 .await;
2401 }
2402
2403 let stalled = stalled_tasks(queue, home, now);
2404 let stalled_ids: std::collections::BTreeSet<_> =
2405 stalled.iter().map(|task| task.id.clone()).collect();
2406 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2407 if !reclaimed.is_empty() {
2408 tracing::warn!(
2409 "reclaimed {} task(s) left `running` by a daemon that never \
2410 recorded the outcome: {}",
2411 reclaimed.len(),
2412 reclaimed.join(", ")
2413 );
2414 }
2415 let abandoned_runs = reclaim_abandoned_runs(home, now);
2416 if !abandoned_runs.is_empty() {
2417 tracing::warn!(
2418 "failed {} run(s) left behind by a killed process, past every \
2419 active seat's own timeout: {}",
2420 abandoned_runs.len(),
2421 abandoned_runs.join(", ")
2422 );
2423 }
2424
2425 let questions = Questions::at(home.join("questions"));
2430
2431 resolve_blockers(queue, &questions);
2434 apply_choice_actions(queue, &questions, home);
2435 reconcile_task_questions(queue, &questions);
2436
2437 let finished: Vec<Task> = finished_tasks(queue)
2443 .into_iter()
2444 .filter(|task| !stalled_ids.contains(&task.id))
2445 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2449 .collect();
2450 let queued = queued_tasks(queue);
2451 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2455 && conductor.worth_a_look(queue, &stalled, &finished)
2456 {
2457 match prepare(&opts.repo, opts) {
2458 Ok(cfg) => {
2459 conductor
2460 .maybe_run(
2461 &cfg,
2462 &opts.repo,
2463 queue,
2464 &questions,
2465 home,
2466 &queued,
2467 &stalled,
2468 &finished,
2469 opts.max_attempts,
2470 )
2471 .await;
2472 }
2473 Err(e) => {
2474 tracing::warn!("conductor: no config: {e:#}");
2475 notices::raise(Notice::warn(
2476 "loop:no-config",
2477 "The loop could not read this repository's config, so held tasks are not being triaged.",
2478 ));
2479 }
2480 }
2481 }
2482
2483 let candidates: Vec<Task> = runnable(queue)
2484 .into_iter()
2485 .filter(|t| !opts.once || !attempted.contains(&t.id))
2486 .collect();
2487
2488 let in_flight: Vec<String> = lock(status)
2493 .current
2494 .iter()
2495 .map(|c| c.task.clone())
2496 .collect();
2497 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2498
2499 interrupt =
2500 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2501 if let Interrupt::Parking {
2502 parked,
2503 interrupt_task,
2504 } = &interrupt
2505 {
2506 let reason = format!(
2507 "task {} asked to run first",
2508 crate::run::short_of(interrupt_task)
2509 );
2510 for id in parked {
2511 if let Some(pause) = interrupt_pauses.get(id) {
2512 pause.park_because(reason.clone());
2513 }
2514 }
2515 }
2516 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2517
2518 let cooling_down =
2519 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2520
2521 let mut started_any = false;
2522 for candidate in candidates {
2523 if stop.stopped() {
2524 break;
2525 }
2526
2527 let resume = land_resume_state(&candidate);
2528 if resume == LandResume::StillWaiting {
2529 continue;
2530 }
2531 let priority = resume == LandResume::Ready;
2532
2533 if !priority && cooling_down {
2534 continue;
2535 }
2536 let permit = match permit_kind(priority, candidate.urgent) {
2537 PermitKind::None => None,
2538 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2539 Ok(p) => Some(p),
2540 Err(_) => continue,
2546 },
2547 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2548 Ok(p) => Some(p),
2549 Err(_) => continue,
2553 },
2554 };
2555
2556 let Ok(claim) = queue.claim(&candidate.id) else {
2561 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2562 continue;
2563 };
2564 let mut task = match queue.get(&candidate.id) {
2567 Ok(t) if t.status.runnable() => t,
2568 Ok(_) => continue,
2569 Err(e) => {
2570 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2571 continue;
2572 }
2573 };
2574 let task_id = task.id.clone();
2575 attempted.push(task_id.clone());
2576 lock(status).idle = false;
2577 stop.enter();
2580 started_any = true;
2581
2582 let run_pause = crate::graph::Pause::new();
2586 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2587
2588 let opts = opts.clone();
2589 let queue = queue.clone();
2590 let status = Arc::clone(status);
2591 let stop = stop.clone();
2592 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2593 inflight.spawn(async move {
2594 let _claim = claim;
2598 let _permit = permit;
2599 let _inflight = InFlightGuard {
2601 status: &status,
2602 stop: &stop,
2603 task_id: &task_id,
2604 };
2605 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2606 lock(&status).completed += 1;
2607 let now = Timestamp::now();
2613 if let Some(until) = cooldown_until("a, now) {
2614 let wait = until.as_second() - now.as_second();
2615 *lock("a_cooldown_until) = Some(until);
2616 let hint = quota
2617 .iter()
2618 .find(|q| q.reset.is_some())
2619 .and_then(|q| q.reset.as_deref());
2620 match hint {
2621 Some(h) => tracing::warn!(
2622 "quota hit; waiting {wait}s before taking another ordinary task \
2623 (CLI reported reset: {h})"
2624 ),
2625 None => tracing::warn!(
2626 "quota hit; waiting {wait}s before taking another ordinary task \
2627 (no reset hint reported)"
2628 ),
2629 }
2630 }
2631 });
2632 }
2633
2634 if started_any {
2635 continue;
2636 }
2637
2638 if stop.busy_now() {
2639 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2644 continue;
2645 }
2646
2647 lock(status).idle = true;
2649 if opts.once {
2650 janitor(&opts.repo, opts, home, worktrees_root).await;
2654 resweep_superseded_attempts(queue, home);
2655 triage_held(queue, home, opts).await;
2656 break;
2657 }
2658 stop.idle(opts.poll).await;
2659 if stop.stopped() {
2660 continue;
2661 }
2662 janitor(&opts.repo, opts, home, worktrees_root).await;
2668 resweep_superseded_attempts(queue, home);
2669 triage_held(queue, home, opts).await;
2670 }
2671
2672 while let Some(result) = inflight.join_next().await {
2677 if let Err(e) = result {
2678 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2679 notices::raise(Notice::error(
2680 "loop:attempt",
2681 "A queued attempt ended abnormally; check the task it was running.",
2682 ));
2683 }
2684 }
2685 Ok(())
2686}
2687
2688async fn attempt(
2694 opts: &Opts,
2695 queue: &Queue,
2696 status: &Arc<Mutex<Status>>,
2697 stop: &Stop,
2698 interrupt_pause: crate::graph::Pause,
2699 task: &mut Task,
2700) -> Vec<QuotaLoss> {
2701 let repo = repo_for(task, &opts.repo);
2702 tracing::info!(
2703 "task {} — {} (repo {})",
2704 task.short(),
2705 task.title,
2706 repo.display()
2707 );
2708
2709 let mut config = match prepare_for(&repo, opts, task) {
2710 Ok(c) => c,
2711 Err(e) => {
2712 task.attempts += 1;
2716 task.fail(format!("config: {e:#}"), opts.max_attempts);
2717 record(queue, task);
2718 return Vec::new();
2719 }
2720 };
2721 apply_solo(&mut config, task);
2722 let p = phrases(&config.graph.language);
2723 let start_failed = p.could_not_start;
2724
2725 if let Some(reason) = disk_gate(&repo, &config) {
2733 task.last_error = Some(reason.clone());
2734 task.hold_machine(Some(reason.clone()));
2735 record(queue, task);
2736 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2737 notices::raise(
2740 Notice::warn(
2741 &format!("disk:{}", repo.display()),
2742 "A task was held for want of free disk space; free some, then release it from the queue.",
2743 )
2744 .link(Link::Task {
2745 id: task.id.clone(),
2746 }),
2747 );
2748 return Vec::new();
2749 }
2750
2751 let unfinished = (!task.fresh_start)
2771 .then(|| unfinished_run(&task.runs, task.short()))
2772 .flatten();
2773 let review_branch = task.review_branch.take();
2779 let branch_exists = match &review_branch {
2780 Some(branch) => crate::git::branch_exists(&repo, branch)
2781 .await
2782 .unwrap_or(false),
2783 None => false,
2784 };
2785 let starter = match (&review_branch, &unfinished, &task.review_of) {
2789 (None, None, Some(branch)) => Starter::Review(branch.clone()),
2790 _ => choose_starter(
2791 review_branch.as_deref(),
2792 branch_exists,
2793 unfinished.as_deref(),
2794 ),
2795 };
2796 let attachments = match task_attachments(queue, task) {
2801 Ok(a) => a,
2802 Err(e) => {
2803 task.attempts += 1;
2804 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2805 record(queue, task);
2806 return Vec::new();
2807 }
2808 };
2809 let started = match &starter {
2810 Starter::Review(branch) => {
2811 tracing::info!(
2812 "task {} reopens `{branch}` as a review-only pass",
2813 task.short()
2814 );
2815 let takeover = crate::handover::Takeover {
2818 earlier: task.earlier_attempts().to_vec(),
2819 home: crate::run::home(),
2820 choice: take_divergence_answer(branch, &config.merge.remote, task),
2821 };
2822 Runner::review_taking_over(
2823 &repo,
2824 branch,
2825 config,
2826 Some(takeover),
2827 crate::run::Origin::queue(&task.id),
2828 )
2829 .await
2830 }
2831 Starter::Resume(id) => {
2832 tracing::info!("resuming run {id} rather than competing again");
2833 Runner::resume(id).map(|mut r| {
2834 if let Some(instruction) =
2835 prepare_instruction(&starter, Some(&r.state.instruction), task)
2836 {
2837 r.state.instruction = instruction;
2838 }
2839 r.state.attachments = attachments.clone();
2840 if let Some(mode) = task
2843 .overrides
2844 .as_ref()
2845 .and_then(|o| o.merge.as_deref())
2846 .and_then(|m| merge_mode(m).ok())
2847 {
2848 r.state.config.merge.mode = mode;
2849 }
2850 r
2851 })
2852 }
2853 Starter::Start => {
2854 if let Some(branch) = &review_branch {
2855 tracing::warn!(
2856 "conductor chose review for task {} but branch `{branch}` no longer \
2857 exists; requeuing as a fresh competition instead",
2858 task.short()
2859 );
2860 }
2861 let instruction = prepare_instruction(&starter, None, task)
2862 .unwrap_or_else(|| task.instruction.clone());
2863 Runner::start_naming(
2864 &repo,
2865 instruction,
2866 &task.title,
2867 config,
2868 crate::run::Origin::queue(&task.id),
2869 )
2870 .await
2871 .map(|mut r| {
2872 r.state.attachments = attachments.clone();
2873 r
2874 })
2875 }
2876 };
2877 let mut runner = match started {
2878 Ok(r) => r,
2879 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2883 let detail = match e.downcast_ref::<crate::handover::Refused>() {
2884 Some(r) => (p.handover_refused)(r),
2885 None => format!("{e:#}"),
2886 };
2887 let reason = format!("{start_failed}{detail}");
2888 task.last_error = Some(reason.clone());
2889 let branch = match &starter {
2892 Starter::Review(branch) => Some(branch.clone()),
2893 _ => None,
2894 };
2895 task.hold_for_handover(branch, format!("{reason}{}", p.handover_hint));
2899 record(queue, task);
2900 tracing::warn!(
2901 "holding {} for a branch it cannot take over: {e:#}",
2902 task.short()
2903 );
2904 return Vec::new();
2905 }
2906 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2911 let d = e
2912 .downcast_ref::<crate::reconcile::Diverged>()
2913 .expect("checked by the guard");
2914 let mut q = ask::Question::new(
2915 task.id.clone(),
2916 "review".to_owned(),
2917 "sync".to_owned(),
2918 d.summary(),
2919 d.detail(),
2920 d.choices(),
2921 );
2922 match Questions::open().put(&mut q) {
2923 Ok(()) => {
2924 task.last_error = Some(format!("{e:#}"));
2925 task.review_branch = Some(d.branch.clone());
2928 task.block(vec![q.id.clone()], Some(d.summary()));
2929 }
2930 Err(put) => {
2931 tracing::warn!("could not file the divergence question: {put:#}");
2932 task.attempts += 1;
2933 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2934 }
2935 }
2936 record(queue, task);
2937 return Vec::new();
2938 }
2939 Err(e) => {
2940 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
2946 && let Starter::Review(branch) = &starter
2947 {
2948 task.review_branch = Some(branch.clone());
2949 task.last_error = Some(format!("{start_failed}{e:#}"));
2950 task.status = crate::queue::TaskStatus::Failed;
2951 record(queue, task);
2952 return Vec::new();
2953 }
2954 task.attempts += 1;
2955 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2956 record(queue, task);
2957 return Vec::new();
2958 }
2959 };
2960 runner.state.followup_generation = Some(task.followup.as_ref().map_or(0, |f| f.generation));
2963 runner.on_pause(stop.pause());
2965 runner.watch_interrupt(interrupt_pause);
2969
2970 let run = runner.state.id.clone();
2973 task.start(run.clone());
2974 record(queue, task);
2975 lock(status).current.push(Current {
2976 task: task.id.clone(),
2977 run,
2978 });
2979
2980 let quota_before = runner.state.quota.clone();
2983 let result = runner.execute().await;
2984 finish_attempt(
2985 opts.max_attempts,
2986 queue,
2987 task,
2988 &runner.state,
2989 "a_before,
2990 result,
2991 )
2992}
2993
2994pub fn finish_attempt(
2999 max_attempts: usize,
3000 queue: &Queue,
3001 task: &mut Task,
3002 state: &RunState,
3003 quota_before: &[QuotaLoss],
3004 result: Result<()>,
3005) -> Vec<QuotaLoss> {
3006 let detail = match result {
3007 Ok(()) => describe(state),
3008 Err(e) => format!("{e:#}"),
3009 };
3010 let fresh = losses_this_attempt(quota_before, &state.quota);
3011 let verdict = Verdict {
3012 status: state.status,
3013 left_pr: state.pr.is_some(),
3016 quota_hit: !fresh.is_empty(),
3022 parked: state.parked,
3026 no_viable_candidates: state.viable().is_empty(),
3029 };
3030 settle_and_diagnose(task, verdict, &detail, max_attempts, state);
3031 if task.status == TaskStatus::Done {
3032 supersede_prior_runs(task, &crate::run::home());
3033 }
3034 record(queue, task);
3035 tracing::info!(
3036 "task {} is {} after run {} ({})",
3037 task.short(),
3038 task.status.as_str(),
3039 state.short(),
3040 label(state.status)
3041 );
3042 fresh
3043}
3044
3045pub fn hold_if_runnable(queue: &Queue, task: &mut Task) {
3049 if task.status.runnable() {
3050 let why = task.last_error.clone().map_or_else(
3051 || "the run did not finish".to_owned(),
3052 |e| format!("the run did not finish: {e}"),
3053 );
3054 task.hold_manual(Some(format!(
3055 "{why}. It was started by hand, so it is not retried \
3056 automatically; `magi task release` retries it."
3057 )));
3058 record(queue, task);
3059 }
3060}
3061
3062#[must_use]
3066pub fn loop_would_resume(run: &str) -> bool {
3067 unfinished_run(&[run.to_owned()], crate::run::short_of(run)).is_some()
3068}
3069
3070#[must_use]
3080pub fn foreign_loop(
3081 reading: Option<&Reading>,
3082 now: Timestamp,
3083 own_pid: u32,
3084) -> Option<Option<u32>> {
3085 let reading = reading.filter(|r| r.running(now))?;
3086 match reading.pid {
3087 Some(pid) if pid == own_pid => None,
3088 pid => Some(pid),
3089 }
3090}
3091
3092pub async fn run_claimed(opts: &Opts, queue: &Queue, task: &mut Task) {
3107 let status = Arc::new(Mutex::new(Status::new()));
3108 let stop = Stop::new();
3109 attempt(
3110 opts,
3111 queue,
3112 &status,
3113 &stop,
3114 crate::graph::Pause::new(),
3115 task,
3116 )
3117 .await;
3118 hold_if_runnable(queue, task);
3119}
3120
3121fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
3131 after
3132 .iter()
3133 .filter(|q| !before.contains(q))
3134 .cloned()
3135 .collect()
3136}
3137
3138fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
3141 if quota.is_empty() {
3142 return None;
3143 }
3144 let with_hint = quota.iter().find(|q| q.reset.is_some());
3145 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
3146 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
3147 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
3148 Some(
3149 now.checked_add(jiff::SignedDuration::from_secs(secs))
3150 .unwrap_or(Timestamp::MAX),
3151 )
3152}
3153
3154fn apply_solo(config: &mut Config, task: &Task) {
3164 if task.solo {
3165 config.graph.candidates = 1;
3166 }
3167}
3168
3169fn prepare_for(repo: &Path, opts: &Opts, task: &Task) -> Result<Config> {
3174 let Some(o) = &task.overrides else {
3175 return prepare(repo, opts);
3176 };
3177 let own = Opts {
3181 config: o.config.clone(),
3182 merge: None,
3183 ..opts.clone()
3184 };
3185 let mut config = prepare(repo, &own)?;
3186 o.apply(&mut config);
3187 Ok(config)
3188}
3189
3190fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
3192 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
3193 if let Some(mode) = &opts.merge {
3194 config.merge.mode = merge_mode(mode)?;
3195 }
3196 Ok(config)
3197}
3198
3199async fn maybe_prune_cache_between_runs(
3230 repo: &Path,
3231 opts: &Opts,
3232 home: &Path,
3233 stop: &Stop,
3234 last_checked: &mut Option<Timestamp>,
3235 now: Timestamp,
3236) {
3237 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
3238 return;
3239 }
3240 *last_checked = Some(now);
3241 let cfg = match prepare(repo, opts) {
3242 Ok(cfg) => cfg,
3243 Err(e) => {
3244 tracing::warn!("cache check: no config: {e:#}");
3245 return;
3246 }
3247 };
3248 match clean::prune_cache_if_over_limit(&cfg, home) {
3249 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
3250 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
3251 pruned.files,
3252 pruned.freed
3253 ),
3254 Ok(_) => {}
3255 Err(e) => {
3256 tracing::warn!("housekeep: prune cache: {e:#}");
3257 notices::raise_in(
3258 home,
3259 Notice::warn(
3260 "housekeep:cache",
3261 "Pruning the shared build cache failed; disk usage may keep growing.",
3262 ),
3263 );
3264 }
3265 }
3266}
3267
3268fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
3272 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
3273}
3274
3275async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
3294 let cfg = match prepare(repo, opts) {
3295 Ok(cfg) => cfg,
3296 Err(e) => {
3297 tracing::warn!("housekeep: no config: {e:#}");
3298 return;
3299 }
3300 };
3301 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
3310 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
3311 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
3316 let mut extra = Vec::new();
3317 if out.unreadable > 0 {
3318 extra.push(format!("{} unreadable", out.unreadable));
3319 }
3320 if out.orphaned_worktrees > 0 {
3321 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
3322 }
3323 let detail = if extra.is_empty() {
3324 String::new()
3325 } else {
3326 format!(" ({})", extra.join(", "))
3327 };
3328 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
3329 }
3330 if out.external_merges_recorded > 0 {
3331 tracing::info!(
3332 "housekeep: recorded {} run(s) as merged externally",
3333 out.external_merges_recorded
3334 );
3335 }
3336 if out.stale_pr_states_repaired > 0 {
3337 tracing::info!(
3338 "housekeep: rewrote {} run record(s) whose pull request had already settled",
3339 out.stale_pr_states_repaired
3340 );
3341 }
3342 if out.cache_files > 0 {
3343 tracing::info!(
3344 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
3345 out.cache_files,
3346 out.cache_freed
3347 );
3348 }
3349 if out.questions_abandoned > 0 {
3350 tracing::info!(
3351 "housekeep: abandoned {} question(s) left open by a finished run",
3352 out.questions_abandoned
3353 );
3354 }
3355}
3356
3357async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3366 let questions = Questions::at(home.join("questions"));
3367 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3368 if report.is_empty() {
3369 return;
3370 }
3371 if !report.quarantined.is_empty() {
3372 tracing::info!(
3373 "triage: held {} blocked task(s) whose blocked-on task or \
3374 question no longer exists: {}",
3375 report.quarantined.len(),
3376 report.quarantined.join(", ")
3377 );
3378 }
3379 if !report.resumed.is_empty() {
3380 tracing::info!(
3381 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3382 report.resumed.len(),
3383 report.resumed.join(", ")
3384 );
3385 }
3386 if !report.asked.is_empty() {
3387 tracing::info!(
3388 "triage: asked about {} held task(s): {}",
3389 report.asked.len(),
3390 report.asked.join(", ")
3391 );
3392 }
3393 if !report.answered.is_empty() {
3394 tracing::info!(
3395 "triage: applied {} operator answer(s): {}",
3396 report.answered.len(),
3397 report.answered.join(", ")
3398 );
3399 }
3400}
3401
3402fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3409 disk_gate_with(repo, config, crate::disk::free_bytes)
3410}
3411
3412fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3416 repo: &Path,
3417 config: &Config,
3418 free_bytes: F,
3419) -> Option<String> {
3420 let min = config.disk.min_free_bytes;
3421 if min == 0 {
3422 return None;
3423 }
3424 match free_bytes(repo) {
3425 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3426 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3427 }
3428}
3429
3430const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3437
3438const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3442
3443fn quota_wait(
3452 reset_at: Option<Timestamp>,
3453 now: Timestamp,
3454 fallback: Duration,
3455 cap: Duration,
3456) -> Duration {
3457 match reset_at {
3458 Some(at) if at > now => {
3459 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3460 Duration::from_secs(secs).min(cap)
3461 }
3462 _ => fallback,
3463 }
3464}
3465
3466fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3480 parse_reset_hint_zoned(text, now)
3481 .or_else(|| parse_reset_hint_dated(text))
3482 .or_else(|| parse_reset_hint_relative(text, recorded))
3483}
3484
3485fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3489 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3490 let mut rest = rest.trim();
3491 if rest.is_empty() {
3492 return None;
3493 }
3494 let mut total: i64 = 0;
3495 let mut matched = false;
3496 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3497 if let Some((digits, tail)) = rest.split_once(unit)
3498 && !digits.is_empty()
3499 && digits.bytes().all(|b| b.is_ascii_digit())
3500 {
3501 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3502 rest = tail;
3503 matched = true;
3504 }
3505 }
3506 if !rest.is_empty() || !matched {
3507 return None;
3508 }
3509 recorded
3510 .checked_add(jiff::SignedDuration::from_secs(total))
3511 .ok()
3512}
3513
3514fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3518 let clock = clock.trim().to_lowercase();
3519 let (digits, pm) = clock
3520 .strip_suffix("am")
3521 .map(|d| (d, false))
3522 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3523 let (h, m) = digits.trim().split_once(':')?;
3524 let mut hour: i8 = h.trim().parse().ok()?;
3525 let minute: i8 = m.trim().parse().ok()?;
3526 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3527 return None;
3528 }
3529 if pm && hour != 12 {
3530 hour += 12;
3531 } else if !pm && hour == 12 {
3532 hour = 0;
3533 }
3534 Some((hour, minute))
3535}
3536
3537fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3542 let open = text.find('(')?;
3543 let close = text.rfind(')')?;
3544 if close <= open {
3545 return None;
3546 }
3547 let zone = text[open + 1..close].trim();
3548 let (hour, minute) = parse_12h_clock(&text[..open])?;
3549 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3550 let candidate = now
3551 .to_zoned(tz)
3552 .with()
3553 .hour(hour)
3554 .minute(minute)
3555 .second(0)
3556 .millisecond(0)
3557 .microsecond(0)
3558 .nanosecond(0)
3559 .build()
3560 .ok()?;
3561 let mut at = candidate.timestamp();
3562 if at <= now {
3563 at += jiff::SignedDuration::from_hours(24);
3564 }
3565 Some(at)
3566}
3567
3568fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3577 let words: Vec<&str> = text.split_whitespace().collect();
3578 if words.len() < 5 {
3579 return None;
3580 }
3581 (0..=words.len() - 5)
3582 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3583}
3584
3585fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3591 if trailing.is_some_and(|next| next.starts_with('(')) {
3592 return None;
3593 }
3594 let month = month_number(window[0])?;
3595 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3596 let day_digits = ["st", "nd", "rd", "th"]
3597 .iter()
3598 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3599 let day: i8 = day_digits.parse().ok()?;
3600 let year_token = window[2];
3601 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3602 return None;
3603 }
3604 let year: i16 = year_token.parse().ok()?;
3605 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3609 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3610 let date = jiff::civil::Date::new(year, month, day).ok()?;
3611 let candidate = date
3612 .at(hour, minute, 0, 0)
3613 .to_zoned(jiff::tz::TimeZone::UTC)
3614 .ok()?;
3615 Some(candidate.timestamp())
3616}
3617
3618fn month_number(name: &str) -> Option<i8> {
3621 const NAMES: [&str; 12] = [
3622 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3623 ];
3624 let lower = name.to_lowercase();
3625 NAMES
3626 .iter()
3627 .position(|n| *n == lower.as_str())
3628 .map(|i| i as i8 + 1)
3629}
3630
3631fn exhausted_review_budget(state: &RunState) -> bool {
3643 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3644}
3645
3646fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3683 unfinished_run_with(runs, short, RunState::load)
3684}
3685
3686fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3689where
3690 F: FnOnce(&str) -> Result<RunState>,
3691{
3692 let id = runs.last()?;
3693 match load(id) {
3694 Ok(s)
3699 if s.status.resumable()
3700 && !s.released()
3701 && !exhausted_review_budget(&s)
3702 && s.liveness(false) != crate::run::Liveness::Live =>
3703 {
3704 Some(id.clone())
3705 }
3706 Ok(_) => None,
3707 Err(e) => {
3708 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3709 None
3710 }
3711 }
3712}
3713
3714fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3722where
3723 F: FnOnce(&str) -> Result<RunState>,
3724{
3725 if task.status != TaskStatus::Failed || task.fresh_start || task.review_branch.is_some() {
3726 return false;
3727 }
3728 let Some(id) = task.runs.last() else {
3729 return false;
3730 };
3731 load(id).is_ok_and(|s| {
3732 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3733 })
3734}
3735
3736#[derive(Debug, Clone, PartialEq, Eq)]
3739enum Starter {
3740 Review(String),
3743 Resume(String),
3745 Start,
3747}
3748
3749fn take_divergence_answer(
3753 branch: &str,
3754 remote: &str,
3755 task: &mut Task,
3756) -> Option<crate::reconcile::Choice> {
3757 let summary = crate::reconcile::summary_for(branch, remote);
3758 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3759 (a.question == summary)
3760 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3761 .flatten()
3762 .map(|c| (i, c))
3763 })?;
3764 task.answers.remove(idx);
3765 Some(choice)
3766}
3767
3768fn choose_starter(
3780 review_branch: Option<&str>,
3781 branch_exists: bool,
3782 unfinished: Option<&str>,
3783) -> Starter {
3784 match review_branch {
3785 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3786 Some(_) => Starter::Start,
3787 None => match unfinished {
3788 Some(id) => Starter::Resume(id.to_owned()),
3789 None => Starter::Start,
3790 },
3791 }
3792}
3793
3794fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3797 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3798 return fallback.to_path_buf();
3799 }
3800 task.repo.clone()
3801}
3802
3803const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3807
3808fn answers_block(task: &Task, count: usize) -> String {
3810 let mut s = ANSWERS_HEADER.to_owned();
3811 for a in &task.answers[..count] {
3812 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3813 }
3814 s
3815}
3816
3817fn append_answers(base: &str, task: &Task) -> String {
3820 if task.answers.is_empty() {
3821 return base.to_owned();
3822 }
3823 let mut s = base.to_owned();
3824 s.push_str(&answers_block(task, task.answers.len()));
3825 s
3826}
3827
3828fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3832 for count in (1..=task.answers.len()).rev() {
3833 let block = answers_block(task, count);
3834 if let Some(base) = instruction.strip_suffix(&block) {
3835 return base;
3836 }
3837 }
3838 instruction
3839}
3840
3841fn instruction_for(task: &Task) -> String {
3849 append_answers(&task.instruction, task)
3850}
3851
3852fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3864 append_answers(strip_answers_block(old_instruction, task), task)
3865}
3866
3867fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3869 let paths = queue.attachment_paths(task);
3870 for (name, path) in task.attachments.iter().zip(&paths) {
3871 if !path.is_file() {
3872 bail!(
3873 "attachment `{name}` is recorded on the task but {} is missing",
3874 path.display()
3875 );
3876 }
3877 }
3878 Ok(paths)
3879}
3880
3881fn prepare_instruction(
3892 starter: &Starter,
3893 old_instruction: Option<&str>,
3894 task: &Task,
3895) -> Option<String> {
3896 match starter {
3897 Starter::Start => Some(instruction_for(task)),
3898 Starter::Resume(_) => Some(resumed_instruction(
3899 old_instruction.expect("a resumed run always has a prior instruction"),
3900 task,
3901 )),
3902 Starter::Review(_) => None,
3903 }
3904}
3905
3906fn record(queue: &Queue, task: &mut Task) {
3910 if let Err(e) = queue.put(task) {
3911 tracing::error!("could not record task {}: {e:#}", task.short());
3912 notices::raise(Notice::error(
3913 "loop:record",
3914 "The loop could not save a task's state; check the disk.",
3915 ));
3916 }
3917}
3918
3919fn runnable(queue: &Queue) -> Vec<Task> {
3925 let mut tasks: Vec<Task> = queue
3926 .list()
3927 .into_iter()
3928 .filter(|t| t.status.runnable())
3929 .collect();
3930 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
3931 tasks
3932}
3933
3934fn describe(state: &RunState) -> String {
3948 let p = phrases(&state.config.graph.language);
3949 let mut detail = if state.status == RunStatus::Stalled {
3950 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
3951 seats.sort_unstable();
3952 seats.dedup();
3953 if seats.is_empty() {
3954 p.quorum_lost.to_owned()
3955 } else {
3956 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
3957 }
3958 } else {
3959 format!("{}{}", p.run_ended, state.status.display_label())
3960 };
3961 if let Some(last) = state.events.last() {
3962 let unanswered = state
3966 .reviews
3967 .last()
3968 .filter(|r| {
3969 state.status == RunStatus::Blocked
3970 && last.node == "review"
3971 && r.incomplete()
3972 && r.blocking == 0
3973 && r.round == state.config.graph.review_rounds
3974 && r.e2e.iter().all(crate::run::CommandOutcome::ok)
3975 })
3976 .map(|r| (r.expected - r.answered, r.round));
3977 match unanswered {
3978 Some((missing, rounds)) if crate::lang::is_japanese(&state.config.graph.language) => {
3979 detail.push_str(&format!(
3980 " ({}: {})",
3981 last.node,
3982 (p.reviewers_never_answered)(missing, rounds)
3983 ));
3984 }
3985 _ => detail.push_str(&format!(" ({}: {})", last.node, last.message)),
3986 }
3987 }
3988 detail.push_str(&format!(" [run {}]", state.id));
3989 detail
3990}
3991
3992const DIAGNOSTIC_MAX: usize = 4_000;
3998
3999const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
4004
4005fn diagnostic(state: &RunState) -> Option<String> {
4019 let mut parts: Vec<String> = Vec::new();
4020
4021 for o in state.gate.iter().filter(|o| !o.ok()) {
4023 parts.push(format!(
4024 "gate `{}` failed ({:?}):\n{}",
4025 o.command,
4026 o.code,
4027 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
4028 ));
4029 }
4030
4031 if let Some(last) = state
4034 .events
4035 .iter()
4036 .rev()
4037 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
4038 {
4039 parts.push(last.message.clone());
4040 }
4041
4042 if state.viable().is_empty() {
4049 for c in &state.candidates {
4050 if let Some(evidence) = &c.verified_noop {
4051 parts.push(format!(
4052 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
4053 c.label
4054 ));
4055 } else if !c.summary.trim().is_empty() {
4056 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
4057 } else if let Some(why) = &c.failed {
4058 parts.push(format!("candidate {}: {why}", c.label));
4059 }
4060 }
4061 }
4062
4063 if parts.is_empty() {
4064 return None;
4065 }
4066 Some(crate::run::tail(
4071 &parts.join("\n\n"),
4072 DIAGNOSTIC_MAX.saturating_sub(100),
4073 ))
4074}
4075
4076fn label(status: RunStatus) -> &'static str {
4084 status.as_str()
4085}
4086
4087pub(crate) fn merge_mode(mode: &str) -> Result<MergeMode> {
4089 match mode {
4090 "none" => Ok(MergeMode::None),
4091 "local" => Ok(MergeMode::Local),
4092 "pr" => Ok(MergeMode::Pr),
4093 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
4094 }
4095}
4096
4097fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
4102 mutex
4103 .lock()
4104 .unwrap_or_else(std::sync::PoisonError::into_inner)
4105}
4106
4107#[cfg(test)]
4108mod tests {
4109 use super::*;
4110 use crate::queue::{Source, TaskStatus};
4111 use crate::run::{Candidate, CommandOutcome};
4112 use pretty_assertions::assert_eq;
4113
4114 fn task() -> Task {
4115 Task::new(
4116 "add retries".to_owned(),
4117 "add retries".to_owned(),
4118 PathBuf::from("/repo"),
4119 Source::Human,
4120 )
4121 }
4122
4123 fn interrupt_task(id: &str) -> Task {
4126 let mut t = task();
4127 t.id = id.to_owned();
4128 t.interrupt = true;
4129 t
4130 }
4131
4132 fn task_with_id(id: &str) -> Task {
4134 let mut t = task();
4135 t.id = id.to_owned();
4136 t
4137 }
4138
4139 fn urgent_task(id: &str) -> Task {
4141 let mut t = task();
4142 t.id = id.to_owned();
4143 t.urgent = true;
4144 t
4145 }
4146
4147 #[test]
4152 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
4153 assert_eq!(permit_kind(true, false), PermitKind::None);
4154 assert_eq!(permit_kind(true, true), PermitKind::None);
4155 }
4156
4157 #[test]
4162 fn permit_kind_separates_urgent_from_ordinary() {
4163 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
4164 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
4165 }
4166
4167 #[test]
4174 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
4175 let cfg = Config::default();
4176 let repo = Path::new("/any/repo/path");
4177
4178 let reason =
4179 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
4180 assert!(reason.contains("1024"), "{reason}");
4181 assert!(
4182 reason.contains(&cfg.disk.min_free_bytes.to_string()),
4183 "{reason}"
4184 );
4185
4186 assert_eq!(
4187 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
4188 None,
4189 "exactly at the floor is open"
4190 );
4191 assert_eq!(
4192 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
4193 None,
4194 "comfortably above the floor is open"
4195 );
4196 }
4197
4198 #[test]
4199 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
4200 let mut cfg = Config::default();
4201 cfg.disk.min_free_bytes = 0;
4202 let repo = Path::new("/any/repo/path");
4203 assert_eq!(
4204 disk_gate_with(repo, &cfg, |_| Ok(0)),
4205 None,
4206 "a zero floor never measures at all"
4207 );
4208 }
4209
4210 #[test]
4211 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
4212 let cfg = Config::default();
4213 let repo = Path::new("/any/repo/path");
4214 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
4215 .expect("a measurement failure must close the gate, not open it");
4216 assert!(reason.contains("could not measure"), "{reason}");
4217 }
4218
4219 #[test]
4220 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
4221 let ordinary = task();
4222 let next = advance_interrupt(
4223 Interrupt::Idle,
4224 std::slice::from_ref(&ordinary.id),
4225 std::slice::from_ref(&ordinary),
4226 );
4227 assert_eq!(next, Interrupt::Idle);
4228 }
4229
4230 #[test]
4231 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
4232 let marked = interrupt_task("marked");
4235 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4236 assert_eq!(next, Interrupt::Idle);
4237 }
4238
4239 #[test]
4240 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
4241 let marked = interrupt_task("marked");
4242 let next = advance_interrupt(
4243 Interrupt::Idle,
4244 &["running".to_owned()],
4245 std::slice::from_ref(&marked),
4246 );
4247 assert_eq!(
4248 next,
4249 Interrupt::Parking {
4250 parked: vec!["running".to_owned()],
4251 interrupt_task: "marked".to_owned(),
4252 }
4253 );
4254 }
4255
4256 #[test]
4265 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
4266 let marked = interrupt_task("marked");
4267
4268 let two = advance_interrupt(
4269 Interrupt::Idle,
4270 &["a".to_owned(), "b".to_owned()],
4271 std::slice::from_ref(&marked),
4272 );
4273 assert_eq!(two, Interrupt::Idle);
4274
4275 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
4276 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
4277 }
4278
4279 #[test]
4280 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
4281 let state = Interrupt::Parking {
4282 parked: vec!["running".to_owned()],
4283 interrupt_task: "marked".to_owned(),
4284 };
4285 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
4287 assert_eq!(still_going, state);
4288
4289 let stopped_but_not_yet_dispatched =
4293 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
4294 assert_eq!(stopped_but_not_yet_dispatched, state);
4295
4296 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
4298 assert_eq!(
4299 dispatched,
4300 Interrupt::Running {
4301 parked: vec!["running".to_owned()],
4302 interrupt_task: "marked".to_owned(),
4303 }
4304 );
4305 }
4306
4307 #[test]
4308 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
4309 let state = Interrupt::Running {
4310 parked: vec!["running".to_owned()],
4311 interrupt_task: "marked".to_owned(),
4312 };
4313 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
4314 assert_eq!(still_running, state);
4315
4316 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
4323 assert_eq!(
4324 ended,
4325 Interrupt::Resuming {
4326 parked: vec!["running".to_owned()]
4327 }
4328 );
4329 }
4330
4331 #[test]
4332 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
4333 let state = Interrupt::Resuming {
4334 parked: vec!["running".to_owned()],
4335 };
4336 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
4337 assert_eq!(still_waiting, state);
4338
4339 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
4340 assert_eq!(dispatched, Interrupt::Idle);
4341 }
4342
4343 #[test]
4349 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
4350 {
4351 let state = Interrupt::Parking {
4352 parked: vec!["running".to_owned()],
4353 interrupt_task: "marked".to_owned(),
4354 };
4355 let next = advance_interrupt(state, &[], &[]);
4358 assert_eq!(
4359 next,
4360 Interrupt::Resuming {
4361 parked: vec!["running".to_owned()]
4362 },
4363 "abandoning the interrupt must not abandon the resume it owes"
4364 );
4365 }
4366
4367 #[test]
4370 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
4371 let state = Interrupt::Resuming {
4372 parked: vec!["running".to_owned()],
4373 };
4374 let next = advance_interrupt(state, &[], &[]);
4375 assert_eq!(
4376 next,
4377 Interrupt::Idle,
4378 "nothing is left to wait for; the loop must not stay wedged"
4379 );
4380 }
4381
4382 #[test]
4383 fn disabled_by_config_the_sequence_can_never_leave_idle() {
4384 let marked = interrupt_task("marked");
4385 let next = advance_interrupt_tick(
4386 false,
4387 Interrupt::Idle,
4388 &["running".to_owned()],
4389 std::slice::from_ref(&marked),
4390 );
4391 assert_eq!(
4392 next,
4393 Interrupt::Idle,
4394 "an unmarked, unconfigured daemon must behave exactly as before"
4395 );
4396 }
4397
4398 #[test]
4399 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4400 let state = Interrupt::Parking {
4401 parked: vec!["running".to_owned()],
4402 interrupt_task: "marked".to_owned(),
4403 };
4404 let candidates = vec![interrupt_task("marked"), task()];
4405 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4406 assert!(
4407 allowed.is_empty(),
4408 "nothing may dispatch - not even the interrupt task itself - \
4409 until the parked run has actually stopped"
4410 );
4411 }
4412
4413 #[test]
4427 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4428 for state in [
4429 Interrupt::Parking {
4430 parked: vec!["running".to_owned()],
4431 interrupt_task: "marked".to_owned(),
4432 },
4433 Interrupt::Running {
4434 parked: vec!["running".to_owned()],
4435 interrupt_task: "marked".to_owned(),
4436 },
4437 Interrupt::Resuming {
4438 parked: vec!["running".to_owned()],
4439 },
4440 ] {
4441 let candidates = vec![
4442 interrupt_task("marked"),
4443 urgent_task("hot"),
4444 task_with_id("ordinary"),
4445 ];
4446 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4447 assert!(
4448 !allowed.iter().any(|t| t.id == "hot"),
4449 "an urgent candidate must wait out the same gate as anything \
4450 else while the run it would run alongside has not actually \
4451 left flight, for state {state:?}: {allowed:?}"
4452 );
4453 }
4454 }
4455
4456 #[test]
4462 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4463 let state = Interrupt::Resuming {
4464 parked: vec!["hot".to_owned()],
4465 };
4466 let candidates = vec![urgent_task("hot"), task()];
4467 let allowed = interrupt_gate(&state, &[], candidates);
4468 assert_eq!(
4469 allowed.iter().filter(|t| t.id == "hot").count(),
4470 1,
4471 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4472 );
4473 }
4474
4475 #[test]
4476 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4477 let state = Interrupt::Parking {
4478 parked: vec!["running".to_owned()],
4479 interrupt_task: "marked".to_owned(),
4480 };
4481 let other = task();
4482 let candidates = vec![interrupt_task("marked"), other.clone()];
4483 let allowed = interrupt_gate(&state, &[], candidates);
4484 assert_eq!(allowed.len(), 1);
4485 assert_eq!(allowed[0].id, "marked");
4486 }
4487
4488 #[test]
4489 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4490 let state = Interrupt::Running {
4491 parked: vec!["running".to_owned()],
4492 interrupt_task: "marked".to_owned(),
4493 };
4494 let candidates = vec![task(), task()];
4495 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4496 assert!(allowed.is_empty());
4497 }
4498
4499 #[test]
4506 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4507 let state = Interrupt::Resuming {
4508 parked: vec!["a".to_owned(), "c".to_owned()],
4509 };
4510 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4511 let allowed = interrupt_gate(&state, &[], candidates);
4512 assert_eq!(
4513 allowed.len(),
4514 1,
4515 "at most one candidate may be offered while resuming: {allowed:?}"
4516 );
4517 assert_eq!(allowed[0].id, "a");
4518 }
4519
4520 #[test]
4521 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4522 let state = Interrupt::Resuming {
4523 parked: vec!["a".to_owned()],
4524 };
4525 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4526 assert!(allowed.is_empty());
4527 }
4528
4529 #[test]
4535 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4536 let running = task(); let marked = interrupt_task("marked");
4538
4539 let mut state = Interrupt::Idle;
4540 let in_flight = vec![running.id.clone()];
4542 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4543 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4544 assert!(gated.is_empty(), "still waiting on `running` to park");
4545
4546 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4548 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4549 assert_eq!(
4550 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4551 vec!["marked"],
4552 "only the interrupt task may be offered to the dispatcher now"
4553 );
4554
4555 state = advance_interrupt_tick(
4557 true,
4558 state,
4559 &["marked".to_owned()],
4560 std::slice::from_ref(&running),
4561 );
4562 let gated = interrupt_gate(
4563 &state,
4564 &["marked".to_owned()],
4565 vec![marked.clone(), running.clone()],
4566 );
4567 assert!(
4568 gated.is_empty(),
4569 "the parked run must not be offered back while the interrupt \
4570 task is still running"
4571 );
4572
4573 let other = task_with_id("other");
4577 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4578 assert_eq!(
4579 state,
4580 Interrupt::Resuming {
4581 parked: vec![running.id.clone()]
4582 }
4583 );
4584 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4585 assert_eq!(
4586 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4587 vec![running.id.as_str()],
4588 "exactly the parked run resumes - not the unrelated task, even \
4589 though it was offered first"
4590 );
4591
4592 state = advance_interrupt_tick(
4596 true,
4597 state,
4598 std::slice::from_ref(&running.id),
4599 std::slice::from_ref(&other),
4600 );
4601 assert_eq!(state, Interrupt::Idle);
4602 let gated = interrupt_gate(
4603 &state,
4604 std::slice::from_ref(&running.id),
4605 vec![other.clone()],
4606 );
4607 assert_eq!(
4608 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4609 vec![other.id.as_str()],
4610 "ordinary dispatch is unrestricted again"
4611 );
4612 }
4613
4614 #[test]
4615 fn every_run_status_settles_the_task_it_came_from() {
4616 let table = [
4618 (RunStatus::Merged, TaskStatus::Done, 1),
4619 (RunStatus::Ready, TaskStatus::Done, 1),
4620 (RunStatus::Stalled, TaskStatus::Failed, 0),
4621 (RunStatus::Blocked, TaskStatus::Failed, 1),
4622 (RunStatus::Failed, TaskStatus::Failed, 1),
4623 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4624 (RunStatus::Prep, TaskStatus::Failed, 1),
4625 (RunStatus::Implementing, TaskStatus::Failed, 1),
4626 (RunStatus::Judging, TaskStatus::Failed, 1),
4627 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4628 (RunStatus::Voting, TaskStatus::Failed, 1),
4629 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4630 (RunStatus::Gating, TaskStatus::Failed, 1),
4631 ];
4632 for (run, want, attempts) in table {
4633 let mut t = task();
4634 t.start("20260902-000000-aaaa".to_owned());
4635 settle(
4636 &mut t,
4637 Verdict {
4638 status: run,
4639 left_pr: false,
4640 parked: false,
4641 quota_hit: matches!(run, RunStatus::Stalled),
4642 no_viable_candidates: false,
4643 },
4644 "why",
4645 2,
4646 );
4647 assert_eq!(t.status, want, "task status after {}", label(run));
4648 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4649 }
4650 }
4651
4652 #[test]
4653 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4654 let mut stalled = task();
4655 stalled.start("20260902-000000-aaaa".to_owned());
4656 settle(
4657 &mut stalled,
4658 Verdict {
4659 status: RunStatus::Stalled,
4660 left_pr: false,
4661 parked: false,
4662 quota_hit: true,
4663 no_viable_candidates: false,
4664 },
4665 "quota",
4666 1,
4667 );
4668 assert_eq!(stalled.attempts, 0);
4669 assert!(
4670 stalled.status.runnable(),
4671 "a machine problem must leave the task in line"
4672 );
4673
4674 let mut blocked = task();
4675 blocked.start("20260902-000000-aaaa".to_owned());
4676 settle(
4677 &mut blocked,
4678 Verdict {
4679 status: RunStatus::Blocked,
4680 left_pr: false,
4681 parked: false,
4682 quota_hit: false,
4683 no_viable_candidates: false,
4684 },
4685 "findings open",
4686 1,
4687 );
4688 assert_eq!(blocked.attempts, 1);
4689 assert_eq!(
4690 blocked.status,
4691 TaskStatus::Held,
4692 "the last attempt hands the task to a human"
4693 );
4694 }
4695
4696 #[test]
4697 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4698 let mut delivered = task();
4701 delivered.start("20260903-080619-01c2".to_owned());
4702 settle(
4703 &mut delivered,
4704 Verdict {
4705 status: RunStatus::Blocked,
4706 left_pr: true,
4707 parked: false,
4708 quota_hit: false,
4709 no_viable_candidates: false,
4710 },
4711 "no check status",
4712 4,
4713 );
4714 assert_eq!(
4715 delivered.status,
4716 TaskStatus::Held,
4717 "a pull request waiting on CI or a person is not a retryable failure"
4718 );
4719 assert!(
4720 !delivered.status.runnable(),
4721 "the loop must not pick this task up again"
4722 );
4723 assert_eq!(
4724 delivered.last_error.as_deref(),
4725 Some("no check status"),
4726 "the operator needs to be told what the gate was waiting for"
4727 );
4728
4729 let mut empty_handed = task();
4732 empty_handed.start("20260903-080619-01c2".to_owned());
4733 settle(
4734 &mut empty_handed,
4735 Verdict {
4736 status: RunStatus::Blocked,
4737 left_pr: false,
4738 parked: false,
4739 quota_hit: false,
4740 no_viable_candidates: false,
4741 },
4742 "findings open",
4743 4,
4744 );
4745 assert_eq!(empty_handed.status, TaskStatus::Failed);
4746 assert!(empty_handed.status.runnable());
4747 }
4748
4749 #[test]
4750 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4751 let mut noop = task();
4757 noop.start("20260912-131304-391f".to_owned());
4758 settle(
4759 &mut noop,
4760 Verdict {
4761 status: RunStatus::VerifiedNoop,
4762 left_pr: false,
4763 parked: false,
4764 quota_hit: false,
4765 no_viable_candidates: true,
4766 },
4767 "candidate A: already fixed by b32cfc4, on main",
4768 4,
4769 );
4770 assert_eq!(
4771 noop.status,
4772 TaskStatus::Held,
4773 "an unverified claim is a request for a human, not a failure"
4774 );
4775 assert!(
4776 !noop.status.runnable(),
4777 "the loop must not requeue this on the same unverified claim"
4778 );
4779 assert_eq!(noop.attempts, 1);
4784 }
4785
4786 #[test]
4787 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4788 let mut parked = task();
4793 parked.start("20260903-183634-2d98".to_owned());
4794 settle(
4795 &mut parked,
4796 Verdict {
4797 status: RunStatus::Implementing,
4798 left_pr: false,
4799 quota_hit: false,
4800 parked: true,
4801 no_viable_candidates: false,
4802 },
4803 "parked after `implementing`",
4804 2,
4805 );
4806 assert_eq!(parked.attempts, 0, "a park is refunded");
4807 assert!(
4808 parked.status.runnable(),
4809 "and the task stays in line so the next loop resumes its run"
4810 );
4811 assert_eq!(
4812 parked.last_error.as_deref(),
4813 Some("parked after `implementing`"),
4814 "the card says where it stopped"
4815 );
4816
4817 let mut broken = task();
4821 broken.start("20260903-183634-2d98".to_owned());
4822 settle(
4823 &mut broken,
4824 Verdict {
4825 status: RunStatus::Implementing,
4826 left_pr: false,
4827 quota_hit: false,
4828 parked: false,
4829 no_viable_candidates: false,
4830 },
4831 "returned mid-flight",
4832 2,
4833 );
4834 assert_eq!(broken.attempts, 1);
4835 }
4836
4837 #[test]
4838 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4839 let mut flaky = task();
4844 flaky.start("20260903-123023-e633".to_owned());
4845 settle(
4846 &mut flaky,
4847 Verdict {
4848 status: RunStatus::Stalled,
4849 left_pr: false,
4850 parked: false,
4851 quota_hit: false,
4852 no_viable_candidates: false,
4853 },
4854 "verdict rests on 1 of 3 judges",
4855 2,
4856 );
4857 assert_eq!(
4858 flaky.attempts, 1,
4859 "flakiness spends an attempt, so `max_attempts` still bounds it"
4860 );
4861 assert!(flaky.status.runnable(), "and it is still worth retrying");
4862
4863 let mut limited = task();
4865 limited.start("20260903-123023-e633".to_owned());
4866 settle(
4867 &mut limited,
4868 Verdict {
4869 status: RunStatus::Stalled,
4870 left_pr: false,
4871 parked: false,
4872 quota_hit: true,
4873 no_viable_candidates: false,
4874 },
4875 "judge-2, judge-3 out of quota",
4876 2,
4877 );
4878 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4879 assert!(limited.status.runnable());
4880
4881 let mut worn = task();
4884 for _ in 0..2 {
4885 worn.release();
4886 }
4887 worn.start("20260903-123023-e633".to_owned());
4888 worn.attempts = 2;
4889 settle(
4890 &mut worn,
4891 Verdict {
4892 status: RunStatus::Stalled,
4893 left_pr: false,
4894 parked: false,
4895 quota_hit: false,
4896 no_viable_candidates: false,
4897 },
4898 "no quorum again",
4899 2,
4900 );
4901 assert_eq!(worn.status, TaskStatus::Held);
4902 assert!(!worn.status.runnable());
4903 }
4904
4905 #[test]
4906 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4907 let mut wiped_out = task();
4914 wiped_out.start("20260907-025000-a1b2".to_owned());
4915 settle(
4916 &mut wiped_out,
4917 Verdict {
4918 status: RunStatus::Failed,
4919 left_pr: false,
4920 parked: false,
4921 quota_hit: true,
4922 no_viable_candidates: true,
4923 },
4924 "no candidate produced a change; nothing to judge",
4925 2,
4926 );
4927 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
4928 assert!(
4929 wiped_out.status.runnable(),
4930 "a machine problem must leave the task in line"
4931 );
4932
4933 let mut partial_progress = task();
4939 partial_progress.start("20260907-025500-c3d4".to_owned());
4940 settle(
4941 &mut partial_progress,
4942 Verdict {
4943 status: RunStatus::Failed,
4944 left_pr: false,
4945 parked: false,
4946 quota_hit: true,
4947 no_viable_candidates: false,
4948 },
4949 "gate failed on the winning candidate",
4950 2,
4951 );
4952 assert_eq!(
4953 partial_progress.attempts, 1,
4954 "a candidate that actually produced a change spends the attempt \
4955 even though some other seat hit its quota"
4956 );
4957 assert!(partial_progress.status.runnable());
4958 }
4959
4960 #[test]
4961 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
4962 let mut t = task();
4968 t.start("20260907-025000-a1b2".to_owned());
4969 let mut state = run_state(RunStatus::Failed);
4970 state.quota.push(QuotaLoss {
4971 seat: "cand-a".to_owned(),
4972 node: "implement".to_owned(),
4973 at: Timestamp::now(),
4974 reset: None,
4975 });
4976 assert!(
4977 state.viable().is_empty(),
4978 "no candidate was added, so nothing is viable"
4979 );
4980 reclaim(&mut t, Some(state), 2, "en");
4981 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
4982 assert!(t.status.runnable());
4983 }
4984
4985 #[test]
4986 fn a_held_task_is_never_offered_to_the_loop() {
4987 let dir = tempfile::tempdir().unwrap();
4988 let queue = Queue::at(dir.path().to_path_buf());
4989 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
4990 let mut t = task();
4991 t.id = format!("2026090{n}-000000-000{n}");
4992 t.priority = priority;
4993 queue.put(&mut t).unwrap();
4994 }
4995 let mut held = task();
4996 held.id = "20260909-000000-9999".to_owned();
4997 held.priority = 99;
4998 held.hold_machine(None);
4999 queue.put(&mut held).unwrap();
5000
5001 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
5002 assert_eq!(order.len(), 3);
5003 assert!(!order.contains(&held.id));
5004 assert_eq!(
5005 order.first().cloned(),
5006 queue.next_runnable().map(|t| t.id),
5007 "the loop's first candidate is exactly what the queue offers"
5008 );
5009 assert_eq!(
5010 order,
5011 vec![
5012 "20260902-000000-0002".to_owned(),
5013 "20260903-000000-0003".to_owned(),
5014 "20260901-000000-0001".to_owned(),
5015 ],
5016 "priority first, then oldest, so nothing starves"
5017 );
5018 }
5019
5020 #[test]
5021 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
5022 let dir = tempfile::tempdir().unwrap();
5023 let queue = Queue::at(dir.path().to_path_buf());
5024 let mut old = task();
5025 old.id = "20260101-000000-old0".to_owned();
5026 queue.put(&mut old).unwrap();
5027 let mut fresh = task();
5028 fresh.id = "20260101-000000-new0".to_owned();
5029 queue.put(&mut fresh).unwrap();
5030
5031 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
5035 std::thread::sleep(Duration::from_millis(60));
5036 let live = queue.claim(&fresh.id).unwrap();
5037
5038 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
5039 assert_eq!(swept, vec![old.id.clone()]);
5040 assert!(
5041 queue.claim(&old.id).is_ok(),
5042 "an unparseable lock older than the threshold is swept"
5043 );
5044 assert!(
5045 queue.claim(&fresh.id).is_err(),
5046 "a live pid protects its lock regardless of age"
5047 );
5048 drop(live);
5049 }
5050
5051 #[test]
5052 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
5053 let dir = tempfile::tempdir().unwrap();
5063 let queue = Queue::at(dir.path().to_path_buf());
5064 let mut t = task();
5065 t.id = "20260101-000000-live".to_owned();
5066 queue.put(&mut t).unwrap();
5067
5068 let claim = queue.claim(&t.id).unwrap();
5069 std::thread::sleep(Duration::from_millis(60));
5070
5071 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
5072 assert!(
5073 swept.is_empty(),
5074 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
5075 );
5076 assert!(
5077 queue.claim(&t.id).is_err(),
5078 "the lock still protects its task"
5079 );
5080 drop(claim);
5081 }
5082
5083 fn injected_dead_pid() -> u32 {
5086 std::process::id().checked_add(1).unwrap_or(1)
5087 }
5088
5089 #[test]
5090 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
5091 let dir = tempfile::tempdir().unwrap();
5092 let queue = Queue::at(dir.path().to_path_buf());
5093 let mut t = task();
5094 t.id = "20260101-000000-dead".to_owned();
5095 queue.put(&mut t).unwrap();
5096 let dead_pid = injected_dead_pid();
5097
5098 std::fs::write(
5103 dir.path().join(format!("{}.lock", t.id)),
5104 dead_pid.to_string(),
5105 )
5106 .unwrap();
5107
5108 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5109 pid != dead_pid
5110 });
5111 assert_eq!(
5112 swept,
5113 vec![t.id.clone()],
5114 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
5115 );
5116 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
5117 }
5118
5119 #[test]
5120 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
5121 let dir = tempfile::tempdir().unwrap();
5122 let queue = Queue::at(dir.path().to_path_buf());
5123 let mut t = task();
5124 t.id = "20260101-000000-late".to_owned();
5125 queue.put(&mut t).unwrap();
5126 let dead_pid = injected_dead_pid();
5127
5128 assert!(
5131 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
5132 "nothing has claimed the task yet"
5133 );
5134
5135 std::fs::write(
5138 dir.path().join(format!("{}.lock", t.id)),
5139 dead_pid.to_string(),
5140 )
5141 .unwrap();
5142
5143 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5147 pid != dead_pid
5148 });
5149 assert_eq!(swept, vec![t.id.clone()]);
5150 }
5151
5152 #[test]
5153 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
5154 crate::run::pin_test_home();
5159 let dir = tempfile::tempdir().unwrap();
5160 let queue = Queue::at(dir.path().to_path_buf());
5161 let mut t = task();
5162 t.id = "20260101-000000-crsh".to_owned();
5163 t.status = TaskStatus::Running;
5164 t.attempts = 1;
5165 t.runs.push("20260904-000000-4043".to_owned());
5169 queue.put(&mut t).unwrap();
5170 let dead_pid = injected_dead_pid();
5171
5172 std::fs::write(
5175 dir.path().join(format!("{}.lock", t.id)),
5176 dead_pid.to_string(),
5177 )
5178 .unwrap();
5179
5180 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
5186 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
5187
5188 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
5189 pid != dead_pid
5190 });
5191 assert_eq!(swept, vec![t.id.clone()]);
5192
5193 let reclaimed = reclaim_orphaned_running(&queue, 2);
5194 assert_eq!(reclaimed, vec![t.id.clone()]);
5195 let after = queue.get(&t.id).unwrap();
5196 assert_eq!(
5197 after.status,
5198 TaskStatus::Held,
5199 "no run.json to recover from, so a human is asked"
5200 );
5201 assert_eq!(
5202 after.runs,
5203 vec!["20260904-000000-4043".to_owned()],
5204 "the crashed run's id is kept as evidence, not discarded"
5205 );
5206 }
5207
5208 #[test]
5209 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
5210 let dir = tempfile::tempdir().unwrap();
5211 let queue = Queue::at(dir.path().to_path_buf());
5212 let mut t = task();
5213 t.id = "20260101-000000-unknown".to_owned();
5214 queue.put(&mut t).unwrap();
5215 let dead_pid = injected_dead_pid();
5216 std::fs::write(
5217 dir.path().join(format!("{}.lock", t.id)),
5218 dead_pid.to_string(),
5219 )
5220 .unwrap();
5221
5222 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
5223 assert!(swept.is_empty(), "an unknown pid must keep its lock");
5224 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
5225 }
5226
5227 fn run_state_in(status: RunStatus, language: &str) -> RunState {
5228 let mut s = run_state(status);
5229 s.config.graph.language = language.to_owned();
5230 s
5231 }
5232
5233 fn unstarted_verdict(status: RunStatus) -> Verdict {
5234 Verdict {
5235 status,
5236 left_pr: false,
5237 quota_hit: false,
5238 parked: false,
5239 no_viable_candidates: false,
5240 }
5241 }
5242
5243 #[test]
5244 fn an_already_in_base_run_finishes_the_task_without_spending_an_attempt() {
5245 let mut t = task();
5246 t.attempts = 1;
5247 settle_in(
5248 &mut t,
5249 unstarted_verdict(RunStatus::AlreadyInBase),
5250 "already in main",
5251 1,
5252 phrases("en"),
5253 );
5254 assert_eq!(t.status, TaskStatus::Done);
5255 assert_eq!(t.attempts, 0);
5256 }
5257
5258 #[test]
5259 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
5260 let reason = |language: &str| {
5261 let mut t = task();
5262 settle_in(
5263 &mut t,
5264 unstarted_verdict(RunStatus::Judging),
5265 "boom",
5266 1,
5267 phrases(language),
5268 );
5269 t.last_error.or(t.hold_reason).unwrap_or_default()
5270 };
5271 assert!(
5272 reason("en").starts_with("the graph stopped at `"),
5273 "{}",
5274 reason("en")
5275 );
5276 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
5277 assert!(reason("日本語").contains("boom"));
5278 assert_eq!(reason("fr"), reason("en"));
5279 }
5280
5281 #[test]
5282 fn describe_follows_the_run_language_and_keeps_the_run_id() {
5283 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
5284 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
5285 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
5286 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
5287 assert!(ja.contains("[run "), "{ja}");
5288 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
5289 assert!(ended.starts_with("run 終了: "), "{ended}");
5290 let mut de = run_state_in(RunStatus::Failed, "de");
5291 let mut en = run_state_in(RunStatus::Failed, "en");
5292 de.id = "same".to_owned();
5293 en.id = "same".to_owned();
5294 assert_eq!(describe(&de), describe(&en));
5295 }
5296
5297 #[test]
5298 fn handover_refusals_follow_the_language_and_keep_the_detail_apart() {
5299 use crate::handover::Refused;
5300 let cases = [
5301 Refused::Foreign {
5302 branch: "b".into(),
5303 path: "/w/x".into(),
5304 why: "made by hand".into(),
5305 },
5306 Refused::Unsafe {
5307 branch: "b".into(),
5308 path: "/w/x".into(),
5309 why: "its worktree has uncommitted changes (a.rs)".into(),
5310 },
5311 Refused::ReleaseFailed {
5312 branch: "b".into(),
5313 path: "/w/x".into(),
5314 run: "ab12".into(),
5315 },
5316 ];
5317 for r in &cases {
5318 let en = (phrases("en").handover_refused)(r);
5319 assert_eq!(en, r.to_string());
5320 let ja = (phrases("ja").handover_refused)(r);
5321 assert!(
5322 !ja.contains("is checked out") && !ja.contains("try again"),
5323 "{ja}"
5324 );
5325 assert!(ja.contains("`b`") && ja.contains("/w/x"), "{ja}");
5326 if let Refused::Foreign { why, .. } | Refused::Unsafe { why, .. } = r {
5327 assert!(ja.contains(&format!("(詳細: {why})")), "{ja}");
5328 }
5329 }
5330 assert!(phrases("en").handover_hint.contains("release the task"));
5331 assert!(phrases("ja").handover_hint.contains("解放"));
5332 }
5333
5334 #[test]
5335 fn describe_translates_the_unanswered_reviewer_stop_only_in_ja() {
5336 let build = |lang: &str| {
5337 let mut s = run_state_in(RunStatus::Blocked, lang);
5338 s.id = "same".to_owned();
5339 s.config.graph.review_rounds = 3;
5340 let mut r = review_round(3);
5341 r.expected = 3;
5342 r.answered = 1;
5343 s.reviews.push(r);
5344 s.event(
5345 "review",
5346 "2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean",
5347 );
5348 s
5349 };
5350 let en = describe(&build("en"));
5351 assert!(en.contains("(review: 2 reviewer seat(s) never answered after 3 rounds; refusing to call it clean)"), "{en}");
5352 let ja = describe(&build("ja"));
5353 assert!(ja.contains("2 席のレビュアーが 3 ラウンド"), "{ja}");
5354 assert!(!ja.contains("never answered"), "{ja}");
5355 assert!(ja.contains("[run "), "{ja}");
5356
5357 let mut failed = build("ja");
5360 failed.reviews[0].e2e.push(crate::run::CommandOutcome {
5361 command: "cargo test".to_owned(),
5362 code: Some(1),
5363 output_tail: String::new(),
5364 duration_ms: 0,
5365 resource_blocked: false,
5366 });
5367 failed.event("review", "stopped; e2e failed: cargo test");
5368 let ja = describe(&failed);
5369 assert!(ja.contains("stopped; e2e failed: cargo test"), "{ja}");
5370 assert!(!ja.contains("席のレビュアー"), "{ja}");
5371 }
5372
5373 #[test]
5374 fn refusals_and_recovery_prose_follow_the_language() {
5375 let t = held_task_with("r1");
5376 let q = action_question(
5377 "r1",
5378 ask::ChoiceAction::Resume {
5379 run: "r1".to_owned(),
5380 },
5381 );
5382 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
5383 ActionDecision::Refuse(s) => s,
5384 other => panic!("{other:?}"),
5385 };
5386 assert!(refuse(phrases("en")).contains("could not be read: gone"));
5387 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
5388 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
5389
5390 let mut held = task();
5391 reclaim(&mut held, None, 2, "ja");
5392 assert!(held.hold_reason.unwrap().contains("保留にしました"));
5393 let mut held = task();
5394 reclaim(&mut held, None, 2, "xx");
5395 assert!(held.hold_reason.unwrap().contains("held for a human"));
5396 }
5397
5398 fn run_state(status: RunStatus) -> RunState {
5399 let mut state = RunState::new(
5400 PathBuf::from("/repo"),
5401 "main".to_owned(),
5402 "abc1234def".to_owned(),
5403 "add retries".to_owned(),
5404 Config::default(),
5405 );
5406 state.status = status;
5407 state
5408 }
5409
5410 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
5411 Candidate {
5412 index: 0,
5413 label,
5414 agent: "claude".to_owned(),
5415 branch: format!("magi/x/{label}"),
5416 worktree: PathBuf::from("/repo"),
5417 summary: summary.to_owned(),
5418 stat: String::new(),
5419 files: 0,
5420 commits: usize::from(!empty),
5421 empty,
5422 failed: failed.map(str::to_owned),
5423 verified_noop: None,
5424 duration_ms: 0,
5425 folded: false,
5426 }
5427 }
5428
5429 #[test]
5430 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
5431 let mut state = run_state(RunStatus::Blocked);
5432 state.gate = vec![
5433 CommandOutcome {
5434 command: "cargo make check".to_owned(),
5435 code: Some(0),
5436 output_tail: "ok".to_owned(),
5437 duration_ms: 0,
5438 resource_blocked: false,
5439 },
5440 CommandOutcome {
5441 command: "cargo test".to_owned(),
5442 code: Some(101),
5443 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
5444 duration_ms: 0,
5445 resource_blocked: false,
5446 },
5447 ];
5448 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
5449 assert!(d.contains("cargo test"), "{d}");
5450 assert!(
5451 !d.contains("cargo make check"),
5452 "a passing check is not a diagnostic: {d}"
5453 );
5454 assert!(d.contains("assertion failed"), "{d}");
5455 }
5456
5457 #[test]
5458 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
5459 let mut state = run_state(RunStatus::Blocked);
5460 state.event(
5461 "land",
5462 "stopped: the fixer produced no commit while 2 check(s) were failing \
5463 (build, lint); stopping instead of looping on an unchanged tree",
5464 );
5465 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
5466 assert!(d.contains("build"), "{d}");
5467 assert!(d.contains("lint"), "{d}");
5468 assert!(d.contains("fixer produced no commit"), "{d}");
5469 }
5470
5471 #[test]
5472 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
5473 let state = run_state(RunStatus::VerifiedNoop);
5478 let d = describe(&state);
5479 assert!(
5480 d.contains("agent-verified no-op"),
5481 "expected the display label, not the wire spelling: {d}"
5482 );
5483 assert!(!d.contains("verified_noop"), "{d}");
5484 }
5485
5486 #[test]
5487 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5488 let mut state = run_state(RunStatus::Failed);
5494 state.candidates = vec![candidate(
5495 'A',
5496 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5497 true,
5498 None,
5499 )];
5500 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5501 assert!(d.contains("candidate A"), "{d}");
5502 assert!(d.contains("tagged v1.2.3"), "{d}");
5503 }
5504
5505 #[test]
5506 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5507 let mut state = run_state(RunStatus::Failed);
5508 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5509 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5510 assert!(d.contains("candidate A"), "{d}");
5511 assert!(d.contains("agent timed out"), "{d}");
5512 }
5513
5514 #[test]
5515 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5516 let mut state = run_state(RunStatus::Failed);
5519 state.candidates = vec![candidate('A', "did the work", false, None)];
5520 assert!(diagnostic(&state).is_none());
5521 }
5522
5523 #[test]
5524 fn diagnostic_is_bounded_however_much_a_run_printed() {
5525 let mut state = run_state(RunStatus::Blocked);
5526 state.gate = vec![
5527 CommandOutcome {
5528 command: "cargo test".to_owned(),
5529 code: Some(101),
5530 output_tail: "x".repeat(50_000),
5531 duration_ms: 0,
5532 resource_blocked: false,
5533 },
5534 CommandOutcome {
5535 command: "cargo clippy".to_owned(),
5536 code: Some(1),
5537 output_tail: "y".repeat(50_000),
5538 duration_ms: 0,
5539 resource_blocked: false,
5540 },
5541 ];
5542 state.candidates = vec![
5543 candidate('A', &"z".repeat(50_000), true, None),
5544 candidate('B', &"w".repeat(50_000), true, None),
5545 ];
5546 let d = diagnostic(&state).expect("plenty here to diagnose");
5547 assert!(
5548 d.len() <= DIAGNOSTIC_MAX,
5549 "diagnostic grew to {} bytes, unbounded",
5550 d.len()
5551 );
5552 }
5553
5554 #[test]
5555 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5556 let mut state = run_state(RunStatus::Blocked);
5557 state.gate = vec![CommandOutcome {
5558 command: "cargo test".to_owned(),
5559 code: Some(101),
5560 output_tail: "assertion failed".to_owned(),
5561 duration_ms: 0,
5562 resource_blocked: false,
5563 }];
5564 let verdict = Verdict {
5565 status: RunStatus::Blocked,
5566 left_pr: false,
5567 quota_hit: false,
5568 parked: false,
5569 no_viable_candidates: false,
5570 };
5571
5572 let mut t = task();
5575 t.start("run-1".to_owned());
5576 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5577 assert_eq!(t.status, TaskStatus::Failed);
5578 assert!(t.diagnostic.is_none());
5579
5580 t.start("run-2".to_owned());
5583 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5584 assert_eq!(t.status, TaskStatus::Held);
5585 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5586 assert!(d.contains("cargo test"), "{d}");
5587 }
5588
5589 #[test]
5590 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5591 crate::run::pin_test_home();
5596 let home = crate::run::home();
5597 let state = run_state(RunStatus::VerifiedNoop);
5598 let mut q = ask::Question::new(
5599 state.id.clone(),
5600 "implement".to_owned(),
5601 "impl-A".to_owned(),
5602 "is this really a no-op?".to_owned(),
5603 String::new(),
5604 Vec::new(),
5605 );
5606 Questions::at(home.join("questions")).put(&mut q).unwrap();
5607
5608 let verdict = Verdict {
5609 status: RunStatus::VerifiedNoop,
5610 left_pr: false,
5611 quota_hit: false,
5612 parked: false,
5613 no_viable_candidates: false,
5614 };
5615 let mut t = task();
5616 t.start(state.id.clone());
5617 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5618
5619 assert_eq!(t.status, TaskStatus::Held);
5620 let reason = t.hold_reason.expect("a held task must record why");
5621 assert!(
5622 reason.starts_with("run ended agent-verified no-op"),
5623 "the original settle reason must survive unchanged: {reason}"
5624 );
5625 assert!(
5626 reason.contains(q.short()),
5627 "the open question's id must be named so the notice is actionable: {reason}"
5628 );
5629 }
5630
5631 #[test]
5632 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5633 crate::run::pin_test_home();
5634 let state = run_state(RunStatus::VerifiedNoop);
5635
5636 let verdict = Verdict {
5637 status: RunStatus::VerifiedNoop,
5638 left_pr: false,
5639 quota_hit: false,
5640 parked: false,
5641 no_viable_candidates: false,
5642 };
5643 let mut t = task();
5644 t.start(state.id.clone());
5645 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5646
5647 assert_eq!(t.status, TaskStatus::Held);
5648 assert_eq!(
5649 t.hold_reason.as_deref(),
5650 Some("run ended agent-verified no-op"),
5651 "nothing to append when the question was already answered or never asked"
5652 );
5653 }
5654
5655 #[test]
5656 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5657 crate::run::pin_test_home();
5658 let mut first = run_state(RunStatus::Blocked);
5659 first.id = "20260101-000000-sup1".to_owned();
5660 first.save().unwrap();
5661 let mut second = run_state(RunStatus::Merged);
5662 second.id = "20260101-000000-sup2".to_owned();
5663 second.save().unwrap();
5664
5665 let mut t = task();
5666 t.runs = vec![first.id.clone(), second.id.clone()];
5667 t.status = TaskStatus::Done;
5668
5669 supersede_prior_runs(&t, &crate::run::home());
5670
5671 assert_eq!(
5672 RunState::load(&first.id).unwrap().status,
5673 RunStatus::Superseded,
5674 "the first attempt's Blocked no longer needs anyone's attention"
5675 );
5676 assert_eq!(
5677 RunState::load(&second.id).unwrap().status,
5678 RunStatus::Merged,
5679 "the run that actually succeeded is left exactly as it was"
5680 );
5681 }
5682
5683 #[test]
5684 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5685 crate::run::pin_test_home();
5691 let mut first = run_state(RunStatus::Blocked);
5692 first.id = "20260101-000000-sup9".to_owned();
5693 first.driver_pid = Some(std::process::id());
5696 first.driver_started_at = Some(
5697 crate::proc::process_started_at(std::process::id())
5698 .expect("this test process's own start time must be queryable"),
5699 );
5700 first.save().unwrap();
5701 let mut second = run_state(RunStatus::Merged);
5702 second.id = "20260101-000000-supa".to_owned();
5703 second.save().unwrap();
5704
5705 let mut t = task();
5706 t.runs = vec![first.id.clone(), second.id.clone()];
5707 t.status = TaskStatus::Done;
5708
5709 supersede_prior_runs(&t, &crate::run::home());
5710
5711 assert_eq!(
5712 RunState::load(&first.id).unwrap().status,
5713 RunStatus::Blocked,
5714 "a live driver_pid means something is still actually working this run, \
5715 even though no daemon claims it - rewriting under it would just be \
5716 undone the next time that process saves"
5717 );
5718 }
5719
5720 #[test]
5721 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5722 let dir = tempfile::tempdir().unwrap();
5728 let home = dir.path().to_path_buf();
5729 let queue = Queue::at(dir.path().join("queue"));
5730
5731 let mut first = run_state(RunStatus::Blocked);
5732 first.id = "20260101-000000-supd".to_owned();
5733 first.driver_pid = Some(std::process::id());
5734 first.driver_started_at = Some(
5735 crate::proc::process_started_at(std::process::id())
5736 .expect("this test process's own start time must be queryable"),
5737 );
5738 first.save_under(&home).unwrap();
5739 let mut second = run_state(RunStatus::Merged);
5740 second.id = "20260101-000000-supe".to_owned();
5741 second.save_under(&home).unwrap();
5742
5743 let mut t = task();
5744 t.runs = vec![first.id.clone(), second.id.clone()];
5745 t.status = TaskStatus::Done;
5746 queue.put(&mut t).unwrap();
5747
5748 resweep_superseded_attempts(&queue, &home);
5749 assert_eq!(
5750 RunState::load_under(&first.id, &home).unwrap().status,
5751 RunStatus::Blocked,
5752 "still live on the first pass, so still untouched"
5753 );
5754
5755 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5761 stale.driver_started_at = Some("1".to_owned());
5762 stale.save_under(&home).unwrap();
5763
5764 resweep_superseded_attempts(&queue, &home);
5765 assert_eq!(
5766 RunState::load_under(&first.id, &home).unwrap().status,
5767 RunStatus::Superseded,
5768 "the second pass catches up what the first one correctly skipped"
5769 );
5770 }
5771
5772 #[test]
5773 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5774 crate::run::pin_test_home();
5775 let mut first = run_state(RunStatus::Blocked);
5776 first.id = "20260101-000000-sup3".to_owned();
5777 first.save().unwrap();
5778 let mut second = run_state(RunStatus::Blocked);
5779 second.id = "20260101-000000-sup4".to_owned();
5780 second.save().unwrap();
5781
5782 let mut t = task();
5783 t.runs = vec![first.id.clone(), second.id.clone()];
5784 t.status = TaskStatus::Failed;
5788
5789 supersede_prior_runs(&t, &crate::run::home());
5790
5791 assert_eq!(
5792 RunState::load(&first.id).unwrap().status,
5793 RunStatus::Blocked
5794 );
5795 assert_eq!(
5796 RunState::load(&second.id).unwrap().status,
5797 RunStatus::Blocked
5798 );
5799 }
5800
5801 #[test]
5802 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5803 crate::run::pin_test_home();
5807 let mut first = run_state(RunStatus::Blocked);
5808 first.id = "20260101-000000-sup5".to_owned();
5809 first.save().unwrap();
5810
5811 let mut t = task();
5812 t.runs = vec![first.id.clone()];
5813 t.status = TaskStatus::Done;
5814
5815 supersede_prior_runs(&t, &crate::run::home());
5816
5817 assert_eq!(
5818 RunState::load(&first.id).unwrap().status,
5819 RunStatus::Blocked,
5820 "a single-attempt task has no earlier run to supersede"
5821 );
5822 }
5823
5824 #[test]
5825 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5826 crate::run::pin_test_home();
5833 let mut first = run_state(RunStatus::Blocked);
5834 first.id = "20260101-000000-supb".to_owned();
5835 first.save().unwrap();
5836 let mut second = run_state(RunStatus::Failed);
5837 second.id = "20260101-000000-supc".to_owned();
5838 second.save().unwrap();
5839
5840 let mut t = task();
5841 t.runs = vec![first.id.clone(), second.id.clone()];
5842 t.status = TaskStatus::Done;
5843
5844 supersede_prior_runs(&t, &crate::run::home());
5845
5846 assert_eq!(
5847 RunState::load(&first.id).unwrap().status,
5848 RunStatus::Blocked,
5849 "the task's last attempt never landed, so there is nothing here \
5850 actually superseding it"
5851 );
5852 }
5853
5854 #[test]
5855 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5856 crate::run::pin_test_home();
5860 let mut failed = run_state(RunStatus::Failed);
5861 failed.id = "20260101-000000-sup6".to_owned();
5862 failed.save().unwrap();
5863 let mut noop = run_state(RunStatus::VerifiedNoop);
5864 noop.id = "20260101-000000-sup7".to_owned();
5865 noop.save().unwrap();
5866 let mut winner = run_state(RunStatus::Ready);
5867 winner.id = "20260101-000000-sup8".to_owned();
5868 winner.save().unwrap();
5869
5870 let mut t = task();
5871 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5872 t.status = TaskStatus::Done;
5873
5874 supersede_prior_runs(&t, &crate::run::home());
5875
5876 assert_eq!(
5877 RunState::load(&failed.id).unwrap().status,
5878 RunStatus::Failed
5879 );
5880 assert_eq!(
5881 RunState::load(&noop.id).unwrap().status,
5882 RunStatus::VerifiedNoop
5883 );
5884 }
5885
5886 fn approval_question(run: &str) -> ask::Question {
5887 ask::Question::new(
5888 run.to_owned(),
5889 land::APPROVAL_NODE.to_owned(),
5890 "land".to_owned(),
5891 "merge?".to_owned(),
5892 String::new(),
5893 vec!["merge".to_owned(), "hold".to_owned()],
5894 )
5895 }
5896
5897 #[test]
5898 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5899 crate::run::pin_test_home();
5900 let mut state = run_state(RunStatus::Landing);
5901 state.id = "20260101-000000-fre1".to_owned();
5902 state.parked = true;
5903 state.save().unwrap();
5904 ask::Questions::open()
5905 .put(&mut approval_question(&state.id))
5906 .unwrap();
5907
5908 let mut t = task();
5909 t.runs.push(state.id.clone());
5910 assert_eq!(
5911 land_resume_state(&t),
5912 LandResume::StillWaiting,
5913 "nobody has answered and the timeout has not passed"
5914 );
5915 }
5916
5917 #[test]
5918 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5919 crate::run::pin_test_home();
5924 let mut state = run_state(RunStatus::Landing);
5925 state.id = "20260101-000000-exp1".to_owned();
5926 state.parked = true;
5927 state.config.graph.answer_timeout = 60;
5928 state.save().unwrap();
5929
5930 let store = ask::Questions::open();
5931 let mut q = approval_question(&state.id);
5932 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
5933 store.put(&mut q).unwrap();
5934
5935 let mut t = task();
5936 t.runs.push(state.id.clone());
5937 assert_eq!(
5938 land_resume_state(&t),
5939 LandResume::Ready,
5940 "an expired question must not be waited on forever"
5941 );
5942
5943 let after = store.get(&q.id).unwrap();
5944 assert!(
5945 !after.status.open(),
5946 "the question is abandoned, not silently ignored"
5947 );
5948 assert!(
5949 after.resolution().is_none(),
5950 "an abandoned question is not read as a decision"
5951 );
5952 }
5953
5954 #[test]
5955 fn reclaim_settles_a_running_task_against_its_last_run() {
5956 let mut t = task();
5957 t.start("20260904-000000-4043".to_owned());
5958 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
5959 assert_eq!(
5960 t.status,
5961 TaskStatus::Done,
5962 "a run that actually finished must not stay `running` forever"
5963 );
5964 }
5965
5966 #[test]
5967 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
5968 let mut t = task();
5972 t.start("20260904-000000-4043".to_owned());
5973 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
5974 assert_eq!(t.status, TaskStatus::Failed);
5975 assert!(t.status.runnable());
5976 }
5977
5978 #[test]
5979 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
5980 let mut t = task();
5981 t.start("20260904-000000-4043".to_owned());
5982 reclaim(&mut t, None, 2, "en");
5983 assert_eq!(t.status, TaskStatus::Held);
5984 assert!(
5985 t.last_error
5986 .as_deref()
5987 .is_some_and(|e| e.contains("running")),
5988 "the operator needs to know why this task was held"
5989 );
5990 }
5991
5992 #[test]
5993 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
5994 let dir = tempfile::tempdir().unwrap();
5995 let queue = Queue::at(dir.path().to_path_buf());
5996
5997 let mut orphaned = task();
5999 orphaned.id = "20260904-000000-orph".to_owned();
6000 orphaned.status = TaskStatus::Running;
6001 orphaned.attempts = 1;
6002 queue.put(&mut orphaned).unwrap();
6003
6004 let mut alive = task();
6005 alive.id = "20260904-000000-live".to_owned();
6006 alive.status = TaskStatus::Running;
6007 alive.attempts = 1;
6008 queue.put(&mut alive).unwrap();
6009 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
6010
6011 let mut queued = task();
6012 queued.id = "20260904-000000-wait".to_owned();
6013 queue.put(&mut queued).unwrap();
6014
6015 let reclaimed = reclaim_orphaned_running(&queue, 2);
6016 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
6017
6018 assert_eq!(
6019 queue.get(&orphaned.id).unwrap().status,
6020 TaskStatus::Held,
6021 "nothing was driving it and there was no run to recover"
6022 );
6023 assert_eq!(
6024 queue.get(&alive.id).unwrap().status,
6025 TaskStatus::Running,
6026 "a live claim must protect the task it belongs to"
6027 );
6028 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
6029 }
6030
6031 fn read_run_under(home: &Path, id: &str) -> RunState {
6037 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
6038 serde_json::from_str(&body).unwrap()
6039 }
6040
6041 #[test]
6042 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
6043 let dir = tempfile::tempdir().unwrap();
6044 let home = dir.path().to_path_buf();
6045 let now = Timestamp::now();
6046 let overrun_seat = || crate::run::ActiveSeat {
6047 node: "implement".to_owned(),
6048 started_at: now - jiff::SignedDuration::new(21_000, 0),
6049 timeout_secs: 3_600,
6050 attempt: 0,
6051 task: None,
6052 command: None,
6053 index: None,
6054 total: None,
6055 };
6056
6057 let mut dead = run_state(RunStatus::Implementing);
6058 dead.id = "20260101-000000-dead".to_owned();
6059 dead.active.insert("impl-A".to_owned(), overrun_seat());
6060 dead.driver_pid = Some(4242);
6063 dead.save_under(&home).unwrap();
6064
6065 let mut alive = run_state(RunStatus::Implementing);
6068 alive.id = "20260101-000000-aliv".to_owned();
6069 alive.active.insert("impl-A".to_owned(), overrun_seat());
6070 alive.save_under(&home).unwrap();
6071 let mut status = Status::new();
6072 status.current = vec![Current {
6073 task: "20260101-000000-task".to_owned(),
6074 run: alive.id.clone(),
6075 }];
6076 write_status_to(&home.join("daemon.json"), &status).unwrap();
6077
6078 let questions = Questions::at(home.join("questions"));
6082 let mut q = ask::Question::new(
6083 dead.id.clone(),
6084 "implement".to_owned(),
6085 "impl-A".to_owned(),
6086 "Which storage backend?".to_owned(),
6087 String::new(),
6088 vec!["SQLite".to_owned(), "Redis".to_owned()],
6089 );
6090 questions.put(&mut q).unwrap();
6091
6092 let abandoned = reclaim_abandoned_runs_with(
6093 &home,
6094 now,
6095 |pid| if pid == 4242 { Some(false) } else { None },
6096 |_| panic!("a query answering Dead outright needs no identity corroboration"),
6097 );
6098 assert_eq!(abandoned, vec![dead.id.clone()]);
6099
6100 let reloaded = read_run_under(&home, &dead.id);
6101 assert_eq!(reloaded.status, RunStatus::Failed);
6102 assert!(reloaded.active.is_empty());
6103 assert!(
6104 !questions.get(&q.id).unwrap().status.open(),
6105 "the failed run's own open question must be settled in the same pass"
6106 );
6107
6108 let still_alive = read_run_under(&home, &alive.id);
6109 assert_eq!(
6110 still_alive.status,
6111 RunStatus::Implementing,
6112 "a live daemon's claim protects it"
6113 );
6114 assert!(!still_alive.active.is_empty());
6115 }
6116
6117 #[test]
6127 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
6128 let dir = tempfile::tempdir().unwrap();
6129 let home = dir.path().to_path_buf();
6130 let now = Timestamp::now();
6131
6132 let mut manual = run_state(RunStatus::Reviewing);
6133 manual.id = "20260101-000000-manl".to_owned();
6134 manual.active.insert(
6135 "review-1".to_owned(),
6136 crate::run::ActiveSeat {
6137 node: "review".to_owned(),
6138 started_at: now - jiff::SignedDuration::new(21_000, 0),
6139 timeout_secs: 3_600,
6140 attempt: 0,
6141 task: None,
6142 command: None,
6143 index: None,
6144 total: None,
6145 },
6146 );
6147 manual.driver_pid = Some(4242);
6151 manual.driver_started_at = Some("1790000000".to_owned());
6152 manual.save_under(&home).unwrap();
6153
6154 let abandoned = reclaim_abandoned_runs_with(
6155 &home,
6156 now,
6157 |pid| if pid == 4242 { Some(true) } else { None },
6158 |pid| {
6159 if pid == 4242 {
6160 Some("1790000000".to_owned())
6161 } else {
6162 None
6163 }
6164 },
6165 );
6166 assert!(
6167 abandoned.is_empty(),
6168 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
6169 );
6170
6171 let reloaded = read_run_under(&home, &manual.id);
6172 assert_eq!(reloaded.status, RunStatus::Reviewing);
6173 assert!(!reloaded.active.is_empty());
6174 }
6175
6176 #[test]
6177 fn an_already_claimed_task_is_skipped_rather_than_failed() {
6178 let dir = tempfile::tempdir().unwrap();
6179 let queue = Queue::at(dir.path().to_path_buf());
6180 let mut only = task();
6181 queue.put(&mut only).unwrap();
6182
6183 let _elsewhere = queue.claim(&only.id).unwrap();
6184 let candidates = runnable(&queue);
6185 assert_eq!(candidates.len(), 1, "the task is still runnable");
6186 assert!(
6187 queue.claim(&candidates[0].id).is_err(),
6188 "the loop cannot take a claim somebody else holds"
6189 );
6190
6191 let after = queue.get(&only.id).unwrap();
6192 assert_eq!(after.status, TaskStatus::Queued);
6193 assert_eq!(
6194 after.attempts, 0,
6195 "losing the race is not an attempt at the task"
6196 );
6197 assert_eq!(after.last_error, None);
6198 }
6199
6200 #[test]
6201 fn the_status_file_round_trips_and_its_heartbeat_advances() {
6202 let dir = tempfile::tempdir().unwrap();
6203 let path = dir.path().join("daemon.json");
6204
6205 let mut status = Status::new();
6206 status.idle = false;
6207 status.completed = 7;
6208 status.current = vec![Current {
6209 task: "20260902-000000-t111".to_owned(),
6210 run: "20260902-000001-r111".to_owned(),
6211 }];
6212 write_status_to(&path, &status).unwrap();
6213 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6214 assert_eq!(first.schema, SCHEMA);
6215 assert_eq!(first.pid, std::process::id());
6216 assert!(!first.idle);
6217 assert_eq!(first.completed, 7);
6218 assert_eq!(first.current, status.current);
6219 assert!(
6220 !path.with_extension("json.tmp").exists(),
6221 "the temp file is renamed, not left behind"
6222 );
6223
6224 std::thread::sleep(Duration::from_millis(5));
6225 status.updated_at = Timestamp::now();
6226 status.polls = 3;
6227 write_status_to(&path, &status).unwrap();
6228 let second: Status =
6229 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
6230 assert!(
6231 second.updated_at > first.updated_at,
6232 "a reader can only detect staleness if the heartbeat moves"
6233 );
6234 assert_eq!(
6235 second.started_at, first.started_at,
6236 "the start time is not a heartbeat"
6237 );
6238 assert_eq!(second.polls, 3);
6239 }
6240
6241 #[test]
6242 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
6243 let dir = tempfile::tempdir().unwrap();
6244
6245 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
6246
6247 let mut status = Status::new();
6248 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
6249 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6250 let stale = read_status(dir.path()).unwrap();
6251 assert!(
6252 !stale.running(Timestamp::now()),
6253 "a minute without a heartbeat is a dead daemon, not a busy one"
6254 );
6255 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
6256
6257 status.updated_at = Timestamp::now();
6258 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6259 let fresh = read_status(dir.path()).unwrap();
6260 assert!(fresh.running(Timestamp::now()));
6261 }
6262
6263 #[test]
6264 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
6265 let dir = tempfile::tempdir().unwrap();
6266 let now = Timestamp::now();
6267 let mine = "20260903-080619-01c2";
6268
6269 assert!(
6270 !is_working_on(dir.path(), mine, now),
6271 "no status file means nobody is working on anything"
6272 );
6273
6274 let mut status = Status::new();
6275 status.current = vec![Current {
6276 task: "20260903-080340-0167".to_owned(),
6277 run: mine.to_owned(),
6278 }];
6279 status.updated_at = now;
6280 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6281 assert!(is_working_on(dir.path(), mine, now));
6282 assert!(
6283 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
6284 "a daemon busy with one run is not working on another"
6285 );
6286
6287 status.updated_at = now - jiff::SignedDuration::from_secs(600);
6290 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6291 assert!(
6292 !is_working_on(dir.path(), mine, now),
6293 "a stale heartbeat is a dead daemon, so its run is a leftover"
6294 );
6295 }
6296
6297 #[test]
6298 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
6299 let dir = tempfile::tempdir().unwrap();
6300 let now = Timestamp::now();
6301
6302 assert!(
6303 !is_working_on_short(dir.path(), "01c2", now),
6304 "no status file means nobody is working on anything"
6305 );
6306
6307 let mut status = Status::new();
6308 status.current = vec![Current {
6309 task: "20260903-080340-0167".to_owned(),
6310 run: "20260903-080619-01c2".to_owned(),
6311 }];
6312 status.updated_at = now;
6313 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6314 assert!(
6315 is_working_on_short(dir.path(), "01c2", now),
6316 "the run's short id is the last block of its full id"
6317 );
6318 assert!(
6319 !is_working_on_short(dir.path(), "3cbf", now),
6320 "a daemon busy with one worktree bay is not working on another"
6321 );
6322 }
6323
6324 #[test]
6325 fn a_newer_status_file_still_yields_a_reading() {
6326 let dir = tempfile::tempdir().unwrap();
6327 std::fs::write(
6330 dir.path().join("daemon.json"),
6331 serde_json::json!({
6332 "schema": 2,
6333 "updated_at": Timestamp::now().to_string(),
6334 "idle": true,
6335 "surprise": { "nested": [1, 2, 3] },
6336 })
6337 .to_string(),
6338 )
6339 .unwrap();
6340
6341 let reading = read_status(dir.path()).expect("a forward-compatible read");
6342 assert!(reading.running(Timestamp::now()));
6343 assert!(reading.idle);
6344 assert!(reading.current.is_empty());
6345 }
6346
6347 #[test]
6348 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
6349 let dir = tempfile::tempdir().unwrap();
6355 std::fs::write(
6356 dir.path().join("daemon.json"),
6357 serde_json::json!({
6358 "schema": 1,
6359 "pid": 4242,
6360 "updated_at": Timestamp::now().to_string(),
6361 "idle": false,
6362 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
6363 "completed": 3,
6364 "polls": 9,
6365 })
6366 .to_string(),
6367 )
6368 .unwrap();
6369
6370 let reading = read_status(dir.path()).expect("an older shape must still parse");
6371 assert!(reading.running(Timestamp::now()));
6372 assert_eq!(
6373 reading.current,
6374 vec![Current {
6375 task: "20260902-140501-aaaa".to_owned(),
6376 run: "20260902-140502-bbbb".to_owned(),
6377 }]
6378 );
6379 }
6380
6381 #[test]
6382 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
6383 let dir = tempfile::tempdir().unwrap();
6384 std::fs::write(
6385 dir.path().join("daemon.json"),
6386 serde_json::json!({
6387 "schema": 1,
6388 "updated_at": Timestamp::now().to_string(),
6389 "idle": true,
6390 "current": null,
6391 })
6392 .to_string(),
6393 )
6394 .unwrap();
6395 let with_null = read_status(dir.path()).expect("null must still parse");
6396 assert!(with_null.current.is_empty());
6397
6398 std::fs::write(
6399 dir.path().join("daemon.json"),
6400 serde_json::json!({
6401 "schema": 1,
6402 "updated_at": Timestamp::now().to_string(),
6403 "idle": true,
6404 })
6405 .to_string(),
6406 )
6407 .unwrap();
6408 let absent = read_status(dir.path()).expect("a missing field must still parse");
6409 assert!(absent.current.is_empty());
6410 }
6411
6412 #[test]
6413 fn a_task_without_a_repository_runs_in_the_daemons_default() {
6414 let fallback = Path::new("/default");
6415 let mut blank = task();
6416 blank.repo = PathBuf::new();
6417 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
6418 let mut dot = task();
6419 dot.repo = PathBuf::from(".");
6420 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
6421 assert_eq!(
6422 repo_for(&task(), fallback),
6423 PathBuf::from("/repo"),
6424 "a task that names a repository keeps it"
6425 );
6426 }
6427
6428 #[test]
6429 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
6430 let mut solo_cfg = Config::default();
6436 solo_cfg.graph.candidates = 3;
6437 let mut solo_task = task();
6438 solo_task.solo = true;
6439 apply_solo(&mut solo_cfg, &solo_task);
6440 assert_eq!(solo_cfg.graph.candidates, 1);
6441
6442 let mut plain_cfg = Config::default();
6443 plain_cfg.graph.candidates = 3;
6444 let plain_task = task();
6445 assert!(!plain_task.solo);
6446 apply_solo(&mut plain_cfg, &plain_task);
6447 assert_eq!(
6448 plain_cfg.graph.candidates, 3,
6449 "a task that did not ask to run alone keeps the config's candidates"
6450 );
6451 }
6452
6453 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
6454 QuotaLoss {
6455 seat: seat.into(),
6456 node: "judge".into(),
6457 at: at.parse().unwrap(),
6458 reset: reset.map(str::to_string),
6459 }
6460 }
6461
6462 #[test]
6463 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
6464 let old: Vec<QuotaLoss> = (1..=4)
6465 .map(|i| {
6466 loss(
6467 &format!("judge-{i}"),
6468 "2026-09-23T05:23:00Z",
6469 Some("2:40pm (Asia/Tokyo)"),
6470 )
6471 })
6472 .collect();
6473 let fresh = losses_this_attempt(&old, &old);
6474 assert!(fresh.is_empty());
6475 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
6476 }
6478
6479 #[test]
6480 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6481 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6482 let now = Timestamp::now();
6483 let mut after = old.clone();
6484 after.push(loss("judge-2", &now.to_string(), None));
6485 let fresh = losses_this_attempt(&old, &after);
6486 assert_eq!(fresh, vec![after[1].clone()]);
6487 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6488 assert_eq!(
6489 until,
6490 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6491 );
6492 }
6493
6494 #[test]
6495 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6496 let before = vec![
6499 loss("judge-1", "2026-09-23T05:23:00Z", None),
6500 loss("judge-2", "2026-09-23T05:24:00Z", None),
6501 ];
6502 let after = vec![
6503 loss("judge-2", "2026-09-23T05:24:00Z", None),
6504 loss("judge-1", "2026-09-24T01:00:00Z", None),
6505 ];
6506 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6507 }
6508
6509 #[test]
6510 fn merge_overrides_are_parsed_or_refused() {
6511 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6512 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6513 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6514 assert!(merge_mode("squash").is_err());
6515 }
6516
6517 #[test]
6518 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6519 let now = Timestamp::now();
6520 let fallback = Duration::from_secs(300);
6521 let cap = Duration::from_secs(1800);
6522
6523 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6525
6526 let soon = now + jiff::SignedDuration::from_secs(600);
6528 assert_eq!(
6529 quota_wait(Some(soon), now, fallback, cap),
6530 Duration::from_secs(600)
6531 );
6532
6533 let past = now - jiff::SignedDuration::from_secs(60);
6536 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6537
6538 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6541 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6542 }
6543
6544 #[test]
6545 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6546 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6547
6548 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6549 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6550
6551 let already_past =
6555 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6556 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6557
6558 assert!(
6559 parse_reset_hint("session limit reached", now, now).is_none(),
6560 "free text with no recognised shape is not guessed at"
6561 );
6562 assert!(
6563 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6564 "an unresolvable zone name is not guessed at either"
6565 );
6566 }
6567
6568 #[test]
6569 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6570 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6571
6572 let at = parse_reset_hint(
6573 "You've hit your usage limit. Visit \
6574 https://chatgpt.com/codex/settings/usage to purchase more \
6575 credits or try again at Sep 19th, 2026 5:10 PM.",
6576 now,
6577 now,
6578 )
6579 .expect("the codex reset wording is a recognised shape");
6580 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6581
6582 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6587 .expect("an explicit year needs no rollover");
6588 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6589
6590 assert!(
6591 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6592 "a two-digit year is not the documented shape and is not guessed at"
6593 );
6594 assert!(
6595 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6596 "a four-letter month name is not the documented three-letter abbreviation"
6597 );
6598 assert!(
6599 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6600 "an explicit zone on the dated shape is a format nobody has \
6601 documented, and is refused rather than guessed at as UTC"
6602 );
6603 }
6604
6605 #[test]
6606 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6607 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6608 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6609
6610 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6611 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6612
6613 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6614 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6615
6616 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6617 assert!(
6618 parse_reset_hint(bad, now, recorded).is_none(),
6619 "{bad:?} must not be guessed at"
6620 );
6621 }
6622 }
6623
6624 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6628 let config = dir.join("magi.toml");
6629 std::fs::write(
6630 &config,
6631 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6632 )
6633 .unwrap();
6634 let opts = Opts {
6635 poll: Duration::from_secs(30),
6636 config: Some(config),
6637 repo: dir.join("repo"),
6641 ..Opts::default()
6642 };
6643 let home = dir.join("home");
6652 let worktrees = dir.join("wt");
6653 (
6654 opts,
6655 Queue::at(dir.join("queue")),
6656 home.join("daemon.json"),
6657 home,
6658 worktrees,
6659 )
6660 }
6661
6662 #[test]
6663 fn a_stop_is_idempotent_and_once_set_stays_set() {
6664 let stop = Stop::new();
6665 assert!(!stop.stopped());
6666
6667 stop.stop();
6668 assert!(stop.stopped());
6669 stop.stop();
6670 assert!(stop.stopped(), "a second stop is not a toggle");
6671
6672 let shared = stop.clone();
6673 assert!(
6674 shared.stopped(),
6675 "a clone is the same stop; that is how the loop and its caller share one"
6676 );
6677 }
6678
6679 #[test]
6680 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6681 let stop = Stop::new();
6682 stop.enter();
6683 assert!(
6684 !stop.finishing(),
6685 "a busy loop nobody has asked to stop is just running"
6686 );
6687
6688 stop.stop();
6689 assert!(
6690 stop.finishing(),
6691 "a stop asked for mid-run has not landed until the run is settled"
6692 );
6693
6694 stop.exit();
6695 assert!(
6696 !stop.finishing(),
6697 "once the run is settled the stop has landed and there is nothing to finish"
6698 );
6699 }
6700
6701 #[test]
6702 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6703 let stop = Stop::new();
6704 stop.enter();
6705 stop.enter();
6706 stop.stop();
6707 assert!(stop.finishing(), "two runs still in flight");
6708
6709 stop.exit();
6710 assert!(
6711 stop.finishing(),
6712 "one run finished, but a sibling is still working"
6713 );
6714
6715 stop.exit();
6716 assert!(
6717 !stop.finishing(),
6718 "the last run out is what actually lands the stop"
6719 );
6720 }
6721
6722 #[tokio::test]
6723 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6724 let dir = tempfile::tempdir().unwrap();
6725 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6726 let stop = Stop::new();
6727 stop.stop();
6728
6729 let began = std::time::Instant::now();
6730 tokio::time::timeout(
6731 Duration::from_secs(2),
6732 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6733 )
6734 .await
6735 .expect("a stopped loop must return, not sit out its poll interval")
6736 .expect("the loop's own setup and teardown must not fail");
6737 assert!(
6738 began.elapsed() < opts.poll,
6739 "returned only after {:?}, which is a poll interval, not a stop",
6740 began.elapsed()
6741 );
6742 }
6743
6744 #[tokio::test]
6745 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6746 let dir = tempfile::tempdir().unwrap();
6747 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6748 let stop = Stop::new();
6749
6750 let asker = {
6753 let stop = stop.clone();
6754 tokio::spawn(async move {
6755 tokio::time::sleep(Duration::from_millis(20)).await;
6756 stop.stop();
6757 })
6758 };
6759
6760 let began = std::time::Instant::now();
6761 tokio::time::timeout(
6762 Duration::from_secs(2),
6763 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6764 )
6765 .await
6766 .expect("a stop asked for while idle must wake the wait")
6767 .expect("the loop's own setup and teardown must not fail");
6768 asker.await.unwrap();
6769 assert!(
6770 began.elapsed() < opts.poll,
6771 "returned only after {:?}, so the stop waited on the sleep",
6772 began.elapsed()
6773 );
6774 }
6775
6776 #[tokio::test]
6777 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6778 let dir = tempfile::tempdir().unwrap();
6779 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6780 let stop = Stop::new();
6781 stop.stop();
6782
6783 tokio::time::timeout(
6784 Duration::from_secs(2),
6785 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6786 )
6787 .await
6788 .expect("a stopped loop must return")
6789 .expect("the loop's own setup and teardown must not fail");
6790
6791 assert!(
6792 home.is_dir(),
6793 "the loop did publish a status file, so its removal is the teardown and not an absence"
6794 );
6795 assert!(
6796 !status_file.exists(),
6797 "a stopped loop clears its status file"
6798 );
6799 assert!(
6800 read_status(&home).is_none(),
6801 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6802 );
6803 }
6804
6805 #[tokio::test]
6806 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6807 let dir = tempfile::tempdir().unwrap();
6808 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6809 opts.once = true;
6810
6811 let mut settled = RunState::new(
6812 dir.path().join("repo"),
6813 "main".to_owned(),
6814 "abc1234".to_owned(),
6815 "fixture".to_owned(),
6816 Config::default(),
6817 );
6818 settled.status = RunStatus::Ready;
6819 let run_dir = home.join("runs").join(&settled.id);
6820 std::fs::create_dir_all(&run_dir).unwrap();
6821 std::fs::write(
6822 run_dir.join("run.json"),
6823 serde_json::to_string_pretty(&settled).unwrap(),
6824 )
6825 .unwrap();
6826 let questions = Questions::at(home.join("questions"));
6827 let mut question = ask::Question::new(
6828 settled.id.clone(),
6829 "review".to_owned(),
6830 "reviewer-1".to_owned(),
6831 "Continue?".to_owned(),
6832 String::new(),
6833 Vec::new(),
6834 );
6835 questions.put(&mut question).unwrap();
6836
6837 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6838 .await
6839 .unwrap();
6840
6841 assert_eq!(
6842 questions.get(&question.id).unwrap().status,
6843 ask::QuestionStatus::Abandoned,
6844 "an empty --once drain still performs startup question cleanup"
6845 );
6846 }
6847
6848 #[test]
6849 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6850 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6851
6852 assert!(
6853 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6854 "never checked before: due at once"
6855 );
6856
6857 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6858 assert!(
6859 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6860 "well inside the interval: not due yet"
6861 );
6862
6863 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6864 assert!(
6865 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6866 "exactly at the edge: not yet due, same convention as `clean::due`"
6867 );
6868
6869 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6870 assert!(
6871 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6872 "past the interval: due again"
6873 );
6874 }
6875
6876 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
6881 let config = dir.join("magi.toml");
6882 std::fs::write(
6888 &config,
6889 format!(
6890 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
6891 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
6892 cache_dir.display()
6893 ),
6894 )
6895 .unwrap();
6896 Opts {
6897 config: Some(config),
6898 repo: dir.join("repo"),
6899 ..Opts::default()
6900 }
6901 }
6902
6903 #[tokio::test]
6904 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
6905 let dir = tempfile::tempdir().unwrap();
6906 let home = dir.path().join("home");
6907 let cache_dir = dir.path().join("cache");
6908 std::fs::create_dir_all(&cache_dir).unwrap();
6909 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6910 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6911
6912 let running = Stop::new();
6915 let mut last_checked = None;
6916 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6917 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
6918 .await;
6919 assert_eq!(
6920 crate::disk::dir_size(&cache_dir),
6921 0,
6922 "over the cap on the first check ever: pruned at once, no idle queue required"
6923 );
6924 assert_eq!(last_checked, Some(t0));
6925
6926 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
6928 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
6929 maybe_prune_cache_between_runs(
6930 &opts.repo,
6931 &opts,
6932 &home,
6933 &running,
6934 &mut last_checked,
6935 too_soon,
6936 )
6937 .await;
6938 assert_eq!(
6939 crate::disk::dir_size(&cache_dir),
6940 10,
6941 "too soon since the last check: left alone rather than rescanned every call"
6942 );
6943 assert_eq!(
6944 last_checked,
6945 Some(t0),
6946 "an idle check does not reset the clock"
6947 );
6948
6949 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6951 maybe_prune_cache_between_runs(
6952 &opts.repo,
6953 &opts,
6954 &home,
6955 &running,
6956 &mut last_checked,
6957 due_again,
6958 )
6959 .await;
6960 assert_eq!(
6961 crate::disk::dir_size(&cache_dir),
6962 0,
6963 "due again: pruned back under the cap"
6964 );
6965 }
6966
6967 #[tokio::test]
6975 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
6976 let dir = tempfile::tempdir().unwrap();
6977 let home = dir.path().join("home");
6978 let cache_dir = dir.path().join("cache");
6979 std::fs::create_dir_all(&cache_dir).unwrap();
6980 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6981 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6982
6983 let stop = Stop::new();
6984 stop.stop();
6985 assert!(
6986 !stop.finishing(),
6987 "no run is in flight at a between-runs boundary, so nothing else \
6988 would tell the operator this stop had not taken effect yet"
6989 );
6990
6991 let mut last_checked = None;
6992 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6993 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
6994 .await;
6995 assert_eq!(
6996 crate::disk::dir_size(&cache_dir),
6997 10,
6998 "over its cap, and due for the first check ever, but a stop outranks \
6999 it: the cap is a standing policy the next start measures again"
7000 );
7001 assert_eq!(
7002 last_checked, None,
7003 "a check that never happened must not claim the interval"
7004 );
7005 }
7006
7007 #[tokio::test]
7021 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
7022 let dir = tempfile::tempdir().unwrap();
7023 let cache_dir = dir.path().join("cache");
7024 std::fs::create_dir_all(&cache_dir).unwrap();
7025 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
7026
7027 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
7028 opts.poll = Duration::from_millis(20);
7029 opts.max_attempts = 1_000;
7030
7031 let queue = Queue::at(dir.path().join("queue"));
7032 let mut t = Task::new(
7033 "x".to_owned(),
7034 "x".to_owned(),
7035 opts.repo.clone(),
7036 Source::Human,
7037 );
7038 queue.put(&mut t).unwrap();
7039
7040 let home = dir.path().join("home");
7041 let worktrees = dir.path().join("wt");
7042 let status_file = home.join("daemon.json");
7043 let stop = Stop::new();
7044 let stopper = {
7045 let stop = stop.clone();
7046 tokio::spawn(async move {
7047 tokio::time::sleep(Duration::from_millis(400)).await;
7048 stop.stop();
7049 })
7050 };
7051
7052 tokio::time::timeout(
7053 Duration::from_secs(10),
7054 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
7055 )
7056 .await
7057 .expect("the loop must not hang on a queue that keeps producing failing work")
7058 .expect("the loop's own setup and teardown must not fail");
7059 stopper.await.unwrap();
7060
7061 let after = queue.get(&t.id).unwrap();
7062 assert!(
7063 after.attempts >= 2,
7064 "the harness must actually have retried more than once, or this is not \
7065 exercising a busy queue at all (got {} attempt(s))",
7066 after.attempts
7067 );
7068 assert!(
7069 after.status.runnable(),
7070 "still under its attempt budget: the queue never reached a natural idle \
7071 on its own, only the external stop ended the test"
7072 );
7073
7074 assert_eq!(
7075 crate::disk::dir_size(&cache_dir),
7076 0,
7077 "an oversized cache must not be left to grow unboundedly just because the \
7078 queue kept the loop busy the whole time"
7079 );
7080 }
7081
7082 #[test]
7083 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
7084 let dir = tempfile::tempdir().unwrap();
7085 let queue = Queue::at(dir.path().join("queue"));
7086 let questions = Questions::at(dir.path().join("questions"));
7087 let mut task = task();
7088 queue.put(&mut task).unwrap();
7089
7090 let mut task_question = ask::Question::new(
7091 task.id.clone(),
7092 crate::conduct::NODE.to_owned(),
7093 "conduct".to_owned(),
7094 "Which backend?".to_owned(),
7095 String::new(),
7096 Vec::new(),
7097 );
7098 questions.put(&mut task_question).unwrap();
7099 task.block(vec![task_question.id.clone()], None);
7100 queue.put(&mut task).unwrap();
7101
7102 let mut run_question = ask::Question::new(
7103 "20260101-000000-run1".to_owned(),
7104 "review".to_owned(),
7105 "reviewer-1".to_owned(),
7106 "Run question".to_owned(),
7107 String::new(),
7108 Vec::new(),
7109 );
7110 questions.put(&mut run_question).unwrap();
7111
7112 let mut coincidental = ask::Question::new(
7117 task.id.clone(),
7118 "review".to_owned(),
7119 "reviewer-1".to_owned(),
7120 "Unrelated review question".to_owned(),
7121 String::new(),
7122 Vec::new(),
7123 );
7124 questions.put(&mut coincidental).unwrap();
7125
7126 reconcile_task_questions(&queue, &questions);
7127 assert!(questions.get(&task_question.id).unwrap().status.open());
7128 assert!(questions.get(&run_question.id).unwrap().status.open());
7129 assert!(questions.get(&coincidental.id).unwrap().status.open());
7130
7131 task.release();
7132 queue.put(&mut task).unwrap();
7133 reconcile_task_questions(&queue, &questions);
7134 assert_eq!(
7135 questions.get(&task_question.id).unwrap().status,
7136 ask::QuestionStatus::Abandoned
7137 );
7138 assert!(
7139 questions.get(&run_question.id).unwrap().status.open(),
7140 "run questions remain the run janitor's responsibility"
7141 );
7142 assert!(
7143 questions.get(&coincidental.id).unwrap().status.open(),
7144 "a non-conductor question must not be abandoned just because its \
7145 run id coincides with a task id"
7146 );
7147 }
7148
7149 #[test]
7150 fn a_freshly_started_running_task_is_never_stalled() {
7151 let dir = tempfile::tempdir().unwrap();
7152 let mut t = task();
7153 t.start("run-1".to_owned());
7154 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
7157 }
7158
7159 #[test]
7160 fn a_long_running_task_with_no_live_daemon_is_stalled() {
7161 let dir = tempfile::tempdir().unwrap();
7162 let mut t = task();
7163 t.start("run-1".to_owned());
7164 t.updated_at = Timestamp::now()
7165 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
7166 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
7167 assert_eq!(
7168 stalled_tasks(
7169 &Queue::at(dir.path().join("q")),
7170 dir.path(),
7171 Timestamp::now()
7172 )
7173 .len(),
7174 0,
7175 "the task was never written to this queue"
7176 );
7177 }
7178
7179 #[test]
7180 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
7181 let dir = tempfile::tempdir().unwrap();
7182 let mut t = task();
7183 t.id = "20260903-080340-0167".to_owned();
7184 t.start("20260903-080619-01c2".to_owned());
7185 t.updated_at = Timestamp::now()
7186 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
7187
7188 let mut status = Status::new();
7189 status.current = vec![Current {
7190 task: t.id.clone(),
7191 run: "20260903-080619-01c2".to_owned(),
7192 }];
7193 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
7194
7195 assert!(
7196 !is_stalled(&t, dir.path(), Timestamp::now()),
7197 "a live daemon's own heartbeat rules out stalled, however long the task has run"
7198 );
7199 }
7200
7201 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
7205 let path = queue.path_of(id);
7206 let body = std::fs::read_to_string(&path).unwrap();
7207 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
7208 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
7209 v["updated_at"] = serde_json::Value::String(old.to_string());
7210 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
7211 }
7212
7213 #[test]
7214 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
7215 let dir = tempfile::tempdir().unwrap();
7228 let queue = Queue::at(dir.path().join("queue"));
7229 let home = dir.path().join("home");
7230
7231 let mut t = task();
7232 t.id = "20260101-000001-lock".to_owned();
7233 t.start("run-1".to_owned());
7234 queue.put(&mut t).unwrap();
7235 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7236 std::fs::write(
7237 dir.path().join("queue").join(format!("{}.lock", t.id)),
7238 "not a pid",
7239 )
7240 .unwrap();
7241
7242 let now = Timestamp::now();
7243 assert!(
7244 reclaim_orphaned_running(&queue, 2).is_empty(),
7245 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
7246 and reclaim must leave the task alone"
7247 );
7248 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
7249
7250 let stalled = stalled_tasks(&queue, &home, now);
7251 assert_eq!(
7252 stalled.len(),
7253 1,
7254 "reclaim's inability to claim it yet must not hide it from the conductor"
7255 );
7256 assert_eq!(stalled[0].id, t.id);
7257 }
7258
7259 #[test]
7260 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
7261 let dir = tempfile::tempdir().unwrap();
7262 crate::run::set_home(dir.path().join("run-home"));
7263 let queue = Queue::at(dir.path().join("queue"));
7264 let home = dir.path().join("home");
7265 let questions = Questions::at(dir.path().join("questions"));
7266
7267 let mut t = task();
7268 t.id = "20260101-000003-dead".to_owned();
7269 t.start("missing-run".to_owned());
7270 queue.put(&mut t).unwrap();
7271 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
7272
7273 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7276 assert_eq!(
7277 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
7278 [&t.id]
7279 );
7280 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
7281 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7282
7283 crate::conduct::apply(
7286 &queue,
7287 &questions,
7288 &crate::conduct::Verdict {
7289 decisions: vec![crate::conduct::Decision {
7290 id: t.id.clone(),
7291 recovery: Some(crate::conduct::Recovery::Requeue),
7292 ..crate::conduct::Decision::default()
7293 }],
7294 },
7295 )
7296 .unwrap();
7297 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7298 }
7299
7300 #[test]
7301 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
7302 let dir = tempfile::tempdir().unwrap();
7303 let queue = Queue::at(dir.path().join("queue"));
7304 let home = dir.path().join("home");
7305
7306 let mut fresh = task();
7307 fresh.id = "20260101-000001-aaaa".to_owned();
7308 fresh.start("run-1".to_owned());
7309 queue.put(&mut fresh).unwrap();
7310
7311 let mut old = task();
7312 old.id = "20260101-000002-bbbb".to_owned();
7313 old.start("run-2".to_owned());
7314 queue.put(&mut old).unwrap();
7315 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
7316
7317 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
7318 assert_eq!(stalled.len(), 1);
7319 assert_eq!(stalled[0].id, old.id);
7320 }
7321
7322 #[test]
7323 fn queued_and_finished_task_views_partition_by_status() {
7324 let dir = tempfile::tempdir().unwrap();
7325 let queue = Queue::at(dir.path().join("queue"));
7326
7327 let mut queued = task();
7328 queued.id = "20260101-000001-aaaa".to_owned();
7329 queue.put(&mut queued).unwrap();
7330
7331 let mut failed = task();
7332 failed.id = "20260101-000002-bbbb".to_owned();
7333 failed.start("run-1".to_owned());
7334 failed.fail("gate red", 5);
7335 queue.put(&mut failed).unwrap();
7336
7337 let mut held = task();
7338 held.id = "20260101-000003-cccc".to_owned();
7339 held.hold_machine(None);
7340 queue.put(&mut held).unwrap();
7341
7342 let mut running = task();
7343 running.id = "20260101-000004-dddd".to_owned();
7344 running.start("run-2".to_owned());
7345 queue.put(&mut running).unwrap();
7346
7347 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
7348 assert_eq!(queued_ids, [queued.id.clone()]);
7349
7350 let mut finished_ids: Vec<String> =
7351 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
7352 finished_ids.sort_unstable();
7353 let mut want = vec![failed.id.clone(), held.id.clone()];
7354 want.sort_unstable();
7355 assert_eq!(finished_ids, want);
7356 }
7357
7358 #[test]
7359 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
7360 let dir = tempfile::tempdir().unwrap();
7361 let queue = Queue::at(dir.path().join("queue"));
7362 let questions = ask::Questions::at(dir.path().join("questions"));
7363
7364 let mut dep = task();
7365 dep.id = "20260101-000001-dep0".to_owned();
7366 dep.succeed();
7367 queue.put(&mut dep).unwrap();
7368
7369 let mut still_going = task();
7370 still_going.id = "20260101-000002-dep1".to_owned();
7371 queue.put(&mut still_going).unwrap();
7372
7373 let mut blocked = task();
7374 blocked.id = "20260101-000003-main".to_owned();
7375 blocked.block(
7376 vec![dep.id.clone(), still_going.id.clone()],
7377 Some("waits on both".to_owned()),
7378 );
7379 queue.put(&mut blocked).unwrap();
7380
7381 resolve_blockers(&queue, &questions);
7382
7383 let after = queue.get(&blocked.id).unwrap();
7384 assert_eq!(
7385 after.status,
7386 TaskStatus::Blocked,
7387 "one dependency is still outstanding"
7388 );
7389 assert_eq!(after.blocked_by, [still_going.id.clone()]);
7390 }
7391
7392 #[test]
7393 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
7394 let dir = tempfile::tempdir().unwrap();
7395 let queue = Queue::at(dir.path().join("queue"));
7396 let questions = ask::Questions::at(dir.path().join("questions"));
7397
7398 let mut q = crate::ask::Question::new(
7399 "20260101-000001-main".to_owned(),
7400 crate::conduct::NODE.to_owned(),
7401 "conduct".to_owned(),
7402 "Which backend?".to_owned(),
7403 String::new(),
7404 Vec::new(),
7405 );
7406 questions.put(&mut q).unwrap();
7407 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
7408 .unwrap();
7409 questions.put(&mut q).unwrap();
7410
7411 let mut blocked = task();
7412 blocked.id = "20260101-000001-main".to_owned();
7413 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
7414 queue.put(&mut blocked).unwrap();
7415
7416 resolve_blockers(&queue, &questions);
7417
7418 let after = queue.get(&blocked.id).unwrap();
7419 assert_eq!(
7420 after.status,
7421 TaskStatus::Queued,
7422 "the only blocker resolved"
7423 );
7424 assert_eq!(after.answers.len(), 1);
7425 assert_eq!(after.answers[0].question, "Which backend?");
7426 assert_eq!(after.answers[0].answer, "SQLite");
7427
7428 let instruction = instruction_for(&after);
7430 assert!(instruction.contains("Which backend?"));
7431 assert!(instruction.contains("SQLite"));
7432 }
7433
7434 #[test]
7435 fn resolve_blockers_holds_a_task_whose_conductor_question_was_abandoned() {
7436 let dir = tempfile::tempdir().unwrap();
7437 let queue = Queue::at(dir.path().join("queue"));
7438 let questions = ask::Questions::at(dir.path().join("questions"));
7439
7440 let mut q = crate::ask::Question::new(
7441 "20260101-000001-main".to_owned(),
7442 crate::conduct::NODE.to_owned(),
7443 "conduct".to_owned(),
7444 "Is the setup done?".to_owned(),
7445 String::new(),
7446 Vec::new(),
7447 );
7448 q.abandon("no answer within 60s of asking");
7449 questions.put(&mut q).unwrap();
7450
7451 let mut blocked = task();
7452 blocked.id = "20260101-000001-main".to_owned();
7453 blocked.block(vec![q.id.clone()], Some("setup?".to_owned()));
7454 queue.put(&mut blocked).unwrap();
7455
7456 resolve_blockers(&queue, &questions);
7457
7458 let after = queue.get(&blocked.id).unwrap();
7459 assert_eq!(
7460 after.status,
7461 TaskStatus::Held,
7462 "never left blocked on nothing"
7463 );
7464 assert!(!after.operator_held(), "a machine hold, for triage");
7465 assert!(
7466 after
7467 .hold_reason
7468 .as_deref()
7469 .unwrap_or_default()
7470 .contains("went unanswered")
7471 );
7472 }
7473
7474 #[test]
7475 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
7476 let dir = tempfile::tempdir().unwrap();
7482 let queue = Queue::at(dir.path().join("queue"));
7483 let questions = ask::Questions::at(dir.path().join("questions"));
7484
7485 let mut q = crate::ask::Question::new(
7486 "20260101-000001-main".to_owned(),
7487 crate::conduct::NODE.to_owned(),
7488 "conduct".to_owned(),
7489 "How should this be handled?".to_owned(),
7490 String::new(),
7491 Vec::new(),
7492 );
7493 questions.put(&mut q).unwrap();
7494 q.answer(crate::ask::Answer::Text(
7495 "leave it held, a human will look at it later".to_owned(),
7496 ))
7497 .unwrap();
7498 questions.put(&mut q).unwrap();
7499
7500 let mut held = task();
7501 held.id = "20260101-000001-main".to_owned();
7502 held.hold_machine(Some("out of attempts".to_owned()));
7503 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
7504 queue.put(&mut held).unwrap();
7505
7506 resolve_blockers(&queue, &questions);
7507
7508 let after = queue.get(&held.id).unwrap();
7509 assert_eq!(after.status, TaskStatus::Held);
7510 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
7511 assert_eq!(
7512 after.answers[0].answer,
7513 "leave it held, a human will look at it later"
7514 );
7515 }
7516
7517 #[test]
7518 fn resolve_blockers_releases_a_task_whose_dependency_was_deleted_on_purpose() {
7519 let dir = tempfile::tempdir().unwrap();
7520 let queue = Queue::at(dir.path().join("queue"));
7521 let questions = ask::Questions::at(dir.path().join("questions"));
7522 let mut dep = task();
7523 dep.id = "20260101-000001-gone".to_owned();
7524 queue.put(&mut dep).unwrap();
7525 let mut blocked = task();
7526 blocked.id = "20260101-000003-main".to_owned();
7527 blocked.block(vec![dep.id.clone()], None);
7528 queue.put(&mut blocked).unwrap();
7529
7530 let claim = queue.claim(&blocked.id).unwrap();
7532 queue.remove(&dep.id, false, &questions).unwrap();
7533 drop(claim);
7534 resolve_blockers(&queue, &questions);
7535
7536 let after = queue.get(&blocked.id).unwrap();
7537 assert_eq!(after.status, TaskStatus::Queued);
7538 assert!(after.blocked_by.is_empty());
7539 }
7540
7541 #[test]
7542 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
7543 let dir = tempfile::tempdir().unwrap();
7549 let queue = Queue::at(dir.path().join("queue"));
7550 let questions = ask::Questions::at(dir.path().join("questions"));
7551
7552 let mut still_going = task();
7553 still_going.id = "20260101-000002-dep1".to_owned();
7554 queue.put(&mut still_going).unwrap();
7555
7556 let mut blocked = task();
7557 blocked.id = "20260101-000003-main".to_owned();
7558 blocked.block(
7559 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7560 Some("waits on both".to_owned()),
7561 );
7562 queue.put(&mut blocked).unwrap();
7563
7564 resolve_blockers(&queue, &questions);
7565
7566 let after = queue.get(&blocked.id).unwrap();
7567 assert_eq!(
7568 after.status,
7569 TaskStatus::Held,
7570 "a missing dependency must not leave the task blocked forever"
7571 );
7572 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7573 assert!(after.blocked_by.is_empty());
7574 let reason = after.hold_reason.as_deref().unwrap_or_default();
7575 assert!(
7576 reason.contains("20260101-000001-gone"),
7577 "the missing id must be named so an operator can tell what happened: {reason}"
7578 );
7579 assert!(
7580 reason.contains(&still_going.id),
7581 "the still-valid dependency must not silently vanish from the record: {reason}"
7582 );
7583 }
7584
7585 #[test]
7586 fn instruction_for_is_unchanged_without_any_answers() {
7587 let t = task();
7588 assert_eq!(instruction_for(&t), t.instruction);
7589 }
7590
7591 #[test]
7592 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7593 let dir = tempfile::tempdir().unwrap();
7594 let q = Queue::at(dir.path().join("queue"));
7595 let src = dir.path().join("shot.png");
7596 std::fs::write(&src, "x").unwrap();
7597 let mut t = task();
7598 q.attach(&mut t, &[src]).unwrap();
7599 let paths = task_attachments(&q, &t).unwrap();
7600 assert_eq!(paths.len(), 1);
7601 assert!(paths[0].is_absolute() && paths[0].is_file());
7602 std::fs::remove_file(&paths[0]).unwrap();
7603 let err = task_attachments(&q, &t).unwrap_err().to_string();
7604 assert!(err.contains("shot.png"), "{err}");
7605 }
7606
7607 #[test]
7608 fn resumed_instruction_is_unchanged_without_any_answers() {
7609 let t = task();
7610 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7611 }
7612
7613 #[test]
7614 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7615 let mut t = task();
7616 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7617 let old = t.instruction.clone();
7621
7622 let refreshed = resumed_instruction(&old, &t);
7623 assert!(refreshed.starts_with(&old), "the original text is kept");
7624 assert!(refreshed.contains("Which backend?"));
7625 assert!(refreshed.contains("SQLite"));
7626 }
7627
7628 #[test]
7629 fn resumed_instruction_keeps_an_original_answers_heading() {
7630 let mut t = task();
7631 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7632 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7633
7634 let refreshed = resumed_instruction(&t.instruction, &t);
7635
7636 assert!(
7637 refreshed.starts_with(&t.instruction),
7638 "an answers heading in the original instruction is not the appended block"
7639 );
7640 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7641 assert!(refreshed.contains("Which backend?"));
7642 assert!(refreshed.contains("SQLite"));
7643
7644 let repeated = resumed_instruction(&refreshed, &t);
7645 assert_eq!(
7646 repeated, refreshed,
7647 "only the final appended block is refreshed"
7648 );
7649 }
7650
7651 #[test]
7652 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7653 let mut t = task();
7654 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7655
7656 let once = resumed_instruction(&t.instruction, &t);
7660 let twice = resumed_instruction(&once, &t);
7661 assert_eq!(once, twice);
7662 assert_eq!(once.matches("Which backend?").count(), 1);
7663
7664 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7666 let refreshed = resumed_instruction(&once, &t);
7667 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7668 assert!(refreshed.contains("Which backend?"));
7669 assert!(refreshed.contains("Which cache?"));
7670 }
7671
7672 #[test]
7673 fn prepare_instruction_covers_all_three_starters() {
7674 let mut t = task();
7675 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7676
7677 assert_eq!(
7680 prepare_instruction(&Starter::Start, None, &t),
7681 Some(instruction_for(&t))
7682 );
7683
7684 let old = t.instruction.clone();
7687 assert_eq!(
7688 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7689 Some(resumed_instruction(&old, &t))
7690 );
7691
7692 assert_eq!(
7696 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7697 None
7698 );
7699 }
7700
7701 #[test]
7702 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7703 assert_eq!(
7704 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7705 Starter::Review("magi/eba2/A".to_owned())
7706 );
7707 }
7708
7709 #[test]
7710 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7711 assert_eq!(
7712 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7713 Starter::Start,
7714 "a vanished review branch must not fall back to resuming the old run either"
7715 );
7716 }
7717
7718 #[test]
7719 fn a_refused_handover_retries_as_a_review_of_the_same_branch() {
7720 let mut t = task();
7721 t.start("old-run".to_owned());
7722 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
7723 t.release();
7724 let branch = t.review_branch.take();
7725 assert_eq!(
7726 choose_starter(branch.as_deref(), true, Some("old-run")),
7727 Starter::Review("magi/eba2/A".to_owned()),
7728 "a review wins over resuming the old run"
7729 );
7730 }
7731
7732 #[test]
7733 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7734 assert_eq!(
7735 choose_starter(None, false, Some("some-run")),
7736 Starter::Resume("some-run".to_owned())
7737 );
7738 assert_eq!(choose_starter(None, false, None), Starter::Start);
7739 }
7740
7741 #[test]
7742 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7743 let mut released = task();
7744 released.start("stalled-run".to_owned());
7745 released.requeue();
7746 let unfinished = (!released.fresh_start)
7747 .then(|| Some("stalled-run".to_owned()))
7748 .flatten();
7749 assert_eq!(
7750 choose_starter(None, false, unfinished.as_deref()),
7751 Starter::Start,
7752 "release keeps run history but must not resume it"
7753 );
7754 assert_eq!(released.runs, ["stalled-run"]);
7755 }
7756
7757 #[test]
7758 fn an_ordinary_release_keeps_a_resumable_run_available() {
7759 let mut released = task();
7760 released.start("stalled-run".to_owned());
7761 released.release();
7762 let unfinished = (!released.fresh_start)
7763 .then(|| Some("stalled-run".to_owned()))
7764 .flatten();
7765 assert_eq!(
7766 choose_starter(None, false, unfinished.as_deref()),
7767 Starter::Resume("stalled-run".to_owned()),
7768 "manual release must preserve the normal resume path"
7769 );
7770 }
7771
7772 #[test]
7773 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7774 let mut state = run_state(RunStatus::Blocked);
7775 state.config.graph.review_rounds = 3;
7776 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7777 assert!(exhausted_review_budget(&state));
7778
7779 state.reviews.pop();
7781 assert!(!exhausted_review_budget(&state));
7782
7783 let mut stalled = run_state(RunStatus::Stalled);
7786 stalled.config.graph.review_rounds = 1;
7787 stalled.reviews = vec![review_round(1)];
7788 assert!(!exhausted_review_budget(&stalled));
7789 }
7790
7791 fn review_round(round: usize) -> crate::run::ReviewRound {
7792 crate::run::ReviewRound {
7793 round,
7794 head: "deadbeef".to_owned(),
7795 verified_head: None,
7796 verified_at: None,
7797 reviews: Vec::new(),
7798 e2e: Vec::new(),
7799 verify_retried: false,
7800 e2e_deferred: false,
7801 e2e_defer_reason: None,
7802 fix: None,
7803 blocking: 0,
7804 answered: 1,
7805 expected: 1,
7806 clean: false,
7807 progressed: true,
7808 vote_split: false,
7809 reconsideration: Vec::new(),
7810 verdict: None,
7811 }
7812 }
7813
7814 #[test]
7815 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7816 let mut run = RunState::new(
7817 PathBuf::from("/repo"),
7818 "main".to_owned(),
7819 "abc1234def".to_owned(),
7820 "add retries".to_owned(),
7821 Config::default(),
7822 );
7823 run.status = RunStatus::Judging;
7824 run.parked = true;
7825 let mut task = Task::new(
7826 "add retries".to_owned(),
7827 "add retries".to_owned(),
7828 PathBuf::from("/repo"),
7829 crate::queue::Source::Human,
7830 );
7831 task.status = TaskStatus::Failed;
7832 task.runs = vec![run.id.clone()];
7833 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7834 assert!(with(&task, &run), "parked after judging is the case");
7835
7836 let mut not_parked = run.clone();
7837 not_parked.parked = false;
7838 not_parked.status = RunStatus::Stalled;
7839 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7840
7841 let mut fresh = task.clone();
7842 fresh.fresh_start = true;
7843 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7844
7845 let mut review = task.clone();
7846 review.review_branch = Some("magi/x/A".to_owned());
7847 assert!(!with(&review, &run), "review is ranked before resume");
7848
7849 let mut held = task.clone();
7850 held.status = TaskStatus::Held;
7851 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7852
7853 let mut released = run.clone();
7854 released.released_to = Some("20260901-000000-new1".to_owned());
7855 assert!(!with(&task, &released), "nothing left to resume into");
7856
7857 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7858 "unreadable"
7859 )));
7860 }
7861
7862 #[test]
7863 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7864 let mut released = RunState::new(
7865 PathBuf::from("/repo"),
7866 "main".to_owned(),
7867 "abc1234def".to_owned(),
7868 "add retries".to_owned(),
7869 Config::default(),
7870 );
7871 released.status = RunStatus::Blocked;
7872 assert_eq!(
7873 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7874 Some(released.id.clone())
7875 );
7876 released.released_to = Some("20260901-000000-new1".to_owned());
7877 assert_eq!(
7878 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7879 None,
7880 "there is nothing left to resume it into"
7881 );
7882 }
7883
7884 #[test]
7885 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
7886 let mut exhausted = RunState::new(
7895 PathBuf::from("/repo"),
7896 "main".to_owned(),
7897 "abc1234def".to_owned(),
7898 "add retries".to_owned(),
7899 Config::default(),
7900 );
7901 exhausted.status = RunStatus::Blocked;
7902 exhausted.config.graph.review_rounds = 1;
7903 exhausted.reviews = vec![review_round(1)];
7904
7905 assert_eq!(
7906 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
7907 None,
7908 "an exhausted `Blocked` run must not be offered as resumable"
7909 );
7910
7911 let mut has_budget_left = RunState::new(
7914 PathBuf::from("/repo"),
7915 "main".to_owned(),
7916 "abc1234def".to_owned(),
7917 "add retries".to_owned(),
7918 Config::default(),
7919 );
7920 has_budget_left.status = RunStatus::Blocked;
7921 has_budget_left.config.graph.review_rounds = 3;
7922 has_budget_left.reviews = vec![review_round(1)];
7923
7924 assert_eq!(
7925 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
7926 Ok(has_budget_left.clone())
7927 }),
7928 Some(has_budget_left.id.clone())
7929 );
7930 }
7931
7932 #[test]
7933 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
7934 let mut older_stalled = RunState::new(
7942 PathBuf::from("/repo"),
7943 "main".to_owned(),
7944 "abc1234def".to_owned(),
7945 "add retries".to_owned(),
7946 Config::default(),
7947 );
7948 older_stalled.status = RunStatus::Stalled;
7949
7950 let mut newest_exhausted = RunState::new(
7951 PathBuf::from("/repo"),
7952 "main".to_owned(),
7953 "abc1234def".to_owned(),
7954 "add retries".to_owned(),
7955 Config::default(),
7956 );
7957 newest_exhausted.status = RunStatus::Blocked;
7958 newest_exhausted.config.graph.review_rounds = 1;
7959 newest_exhausted.reviews = vec![review_round(1)];
7960
7961 assert_eq!(
7962 unfinished_run_with(
7963 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
7964 "t",
7965 |_| Ok(newest_exhausted.clone())
7966 ),
7967 None,
7968 "the newest run is exhausted, so nothing here is worth resuming - \
7969 least of all the older, already-superseded run"
7970 );
7971 }
7972
7973 #[test]
7974 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
7975 assert_eq!(
7976 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
7977 Err(anyhow::anyhow!("fixture is absent"))
7978 }),
7979 None
7980 );
7981 }
7982
7983 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
7984 let mut q = ask::Question::new(
7985 run.to_owned(),
7986 "implement".to_owned(),
7987 "impl-A".to_owned(),
7988 "continue?".to_owned(),
7989 String::new(),
7990 vec!["resume で続行する".to_owned(), "other".to_owned()],
7991 );
7992 q.actions.insert("resume で続行する".to_owned(), action);
7993 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
7994 .unwrap();
7995 q
7996 }
7997
7998 fn held_task_with(run: &str) -> Task {
7999 let mut t = task();
8000 t.runs = vec![run.to_owned()];
8001 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
8002 t
8003 }
8004
8005 fn resume_action(run: &str) -> ask::ChoiceAction {
8006 ask::ChoiceAction::Resume { run: run.into() }
8007 }
8008
8009 #[test]
8010 fn decide_action_resumes_only_the_latest_resumable_run() {
8011 let t = held_task_with("r1");
8012 let q = action_question("r1", resume_action("r1"));
8013 let load = |s: RunState| move |_: &str| Ok(s);
8014 assert_eq!(
8015 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
8016 ActionDecision::Resume("r1".into())
8017 );
8018 let q_other = action_question("r1", resume_action("r0"));
8020 assert!(matches!(
8021 decide_action(
8022 &t,
8023 &q_other,
8024 &PHRASES_EN,
8025 load(run_state(RunStatus::Blocked))
8026 ),
8027 ActionDecision::Refuse(_)
8028 ));
8029 assert!(matches!(
8031 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
8032 ActionDecision::Refuse(_)
8033 ));
8034 let mut released = run_state(RunStatus::Blocked);
8036 released.released_to = Some("elsewhere".into());
8037 assert!(matches!(
8038 decide_action(&t, &q, &PHRASES_EN, load(released)),
8039 ActionDecision::Refuse(_)
8040 ));
8041 assert!(matches!(
8043 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
8044 ActionDecision::Refuse(_)
8045 ));
8046 }
8047
8048 #[test]
8049 fn decide_action_ignores_a_question_about_an_earlier_run() {
8050 let mut t = held_task_with("r1");
8051 t.runs.push("r2".to_owned());
8052 let never = |_: &str| -> Result<RunState> { bail!("not read") };
8053 assert_eq!(
8054 decide_action(
8055 &t,
8056 &action_question("r1", ask::ChoiceAction::Done),
8057 &PHRASES_EN,
8058 never
8059 ),
8060 ActionDecision::Stale
8061 );
8062 }
8063
8064 #[test]
8065 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
8066 let mut t = held_task_with("r1");
8067 let never = |_: &str| -> Result<RunState> { bail!("not read") };
8068 assert_eq!(
8069 decide_action(
8070 &t,
8071 &action_question("r1", ask::ChoiceAction::Requeue),
8072 &PHRASES_EN,
8073 never
8074 ),
8075 ActionDecision::Requeue
8076 );
8077 let done_q = action_question("r1", ask::ChoiceAction::Done);
8078 assert_eq!(
8079 decide_action(&t, &done_q, &PHRASES_EN, never),
8080 ActionDecision::Done
8081 );
8082 t.mark_action_applied(&done_q.id);
8083 assert_eq!(
8084 decide_action(&t, &done_q, &PHRASES_EN, never),
8085 ActionDecision::Skip
8086 );
8087
8088 let mut plain = action_question("r1", ask::ChoiceAction::Done);
8090 plain.actions.clear();
8091 assert_eq!(
8092 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
8093 ActionDecision::Skip
8094 );
8095 let mut running = held_task_with("r1");
8097 running.status = TaskStatus::Running;
8098 assert_eq!(
8099 decide_action(
8100 &running,
8101 &action_question("r1", ask::ChoiceAction::Done),
8102 &PHRASES_EN,
8103 never
8104 ),
8105 ActionDecision::Skip
8106 );
8107 }
8108
8109 #[test]
8110 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
8111 let dir = tempfile::tempdir().unwrap();
8112 let queue = Queue::at(dir.path().join("queue"));
8113 let questions = Questions::at(dir.path().join("questions"));
8114 let home = dir.path().join("home");
8115 let mut state = run_state(RunStatus::Blocked);
8116 state.id = "20260101-000000-act1".to_owned();
8117 state.save_under(&home).unwrap();
8118
8119 let mut t = held_task_with(&state.id);
8120 queue.put(&mut t).unwrap();
8121 let mut q = action_question(&state.id, resume_action(&state.id));
8122 questions.put(&mut q).unwrap();
8123
8124 apply_choice_actions(&queue, &questions, &home);
8125 let after = queue.get(&t.id).unwrap();
8126 assert_eq!(after.status, TaskStatus::Queued);
8127 assert!(!after.fresh_start);
8128 assert!(after.action_applied(&q.id));
8129 let pin = after.resume_override.clone().unwrap();
8130 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
8131 assert!(pin.forced);
8132
8133 let mut again = queue.get(&t.id).unwrap();
8135 again.hold_machine(Some("later".into()));
8136 queue.put(&mut again).unwrap();
8137 apply_choice_actions(&queue, &questions, &home);
8138 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8139 }
8140
8141 #[test]
8142 fn a_delivered_answer_is_still_acted_on_exactly_once() {
8143 let dir = tempfile::tempdir().unwrap();
8144 let queue = Queue::at(dir.path().join("queue"));
8145 let questions = Questions::at(dir.path().join("questions"));
8146 let home = dir.path().join("home");
8147 let mut state = run_state(RunStatus::Blocked);
8148 state.id = "20260101-000000-act2".to_owned();
8149 state.save_under(&home).unwrap();
8150
8151 let mut t = held_task_with(&state.id);
8152 queue.put(&mut t).unwrap();
8153 let mut q = action_question(&state.id, resume_action(&state.id));
8154 q.answer_delivered = true;
8155 questions.put(&mut q).unwrap();
8156
8157 apply_choice_actions(&queue, &questions, &home);
8160 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
8161 assert!(queue.get(&t.id).unwrap().action_applied(&q.id));
8162
8163 let mut again = queue.get(&t.id).unwrap();
8165 again.hold_machine(Some("later".into()));
8166 queue.put(&mut again).unwrap();
8167 apply_choice_actions(&queue, &questions, &home);
8168 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8169 }
8170
8171 #[test]
8172 fn a_fresh_asker_defers_the_action_until_it_goes_quiet() {
8173 let dir = tempfile::tempdir().unwrap();
8174 let queue = Queue::at(dir.path().join("queue"));
8175 let questions = Questions::at(dir.path().join("questions"));
8176 let home = dir.path().join("home");
8177 let mut state = run_state(RunStatus::Blocked);
8178 state.id = "20260101-000000-act3".to_owned();
8179 state.save_under(&home).unwrap();
8180 let mut t = held_task_with(&state.id);
8181 queue.put(&mut t).unwrap();
8182 let mut q = action_question(&state.id, ask::ChoiceAction::Requeue);
8183 questions.put(&mut q).unwrap();
8184
8185 questions.beat(&q.id, crate::ask::WaiterKind::Asker);
8186 apply_choice_actions(&queue, &questions, &home);
8187 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8188
8189 std::fs::remove_file(questions.root().join(format!("{}.lease", q.id))).unwrap();
8190 apply_choice_actions(&queue, &questions, &home);
8191 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
8192 }
8193
8194 #[test]
8195 fn action_standing_tells_the_reasons_apart() {
8196 let mut t = held_task_with("r1");
8197 let q = action_question("r1", ask::ChoiceAction::Requeue);
8198 assert_eq!(action_standing(&t, &q), ActionStanding::Pending);
8199 let mut plain = q.clone();
8200 plain.actions.clear();
8201 assert_eq!(action_standing(&t, &plain), ActionStanding::NoAction);
8202
8203 t.mark_action_applied(&q.id);
8205 assert_eq!(action_standing(&t, &q), ActionStanding::Applied);
8206
8207 let mut t = held_task_with("r1");
8209 t.start("r2".to_owned());
8210 assert_eq!(action_standing(&t, &q), ActionStanding::Stale);
8211 let q2 = action_question("r2", ask::ChoiceAction::Requeue);
8213 assert_eq!(action_standing(&t, &q2), ActionStanding::Busy);
8214 t.status = TaskStatus::Blocked;
8215 assert_eq!(action_standing(&t, &q2), ActionStanding::Busy);
8216
8217 let mut c = action_question("r0", ask::ChoiceAction::Requeue);
8219 c.node = crate::conduct::NODE.to_owned();
8220 assert_eq!(action_standing(&t, &c), ActionStanding::Busy);
8221 assert!(!ActionStanding::Busy.daemon_owns());
8222 }
8223
8224 #[test]
8225 fn a_running_task_waits_and_a_dead_asker_with_a_cwd_is_still_actioned() {
8226 let dir = tempfile::tempdir().unwrap();
8227 let queue = Queue::at(dir.path().join("queue"));
8228 let questions = Questions::at(dir.path().join("questions"));
8229 let home = dir.path().join("home");
8230 let mut state = run_state(RunStatus::Blocked);
8231 state.id = "20260101-000000-act4".to_owned();
8232 state.save_under(&home).unwrap();
8233 let mut t = held_task_with(&state.id);
8234 t.status = TaskStatus::Running;
8235 queue.put(&mut t).unwrap();
8236 let mut q = action_question(&state.id, ask::ChoiceAction::Done);
8237 q.cwd = Some(dir.path().display().to_string());
8238 questions.put(&mut q).unwrap();
8239
8240 let running = queue.get(&t.id).unwrap();
8243 assert_eq!(action_standing(&running, &q), ActionStanding::Busy);
8244 assert!(matches!(
8245 crate::waiter::decide_owned(&q, None, false, 86_400, Timestamp::now(), false),
8246 crate::waiter::Action::Deliver(_)
8247 ));
8248 apply_choice_actions(&queue, &questions, &home);
8250 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
8251 let mut back = queue.get(&t.id).unwrap();
8252 back.hold_machine(Some("later".into()));
8253 queue.put(&mut back).unwrap();
8254 apply_choice_actions(&queue, &questions, &home);
8255 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Done);
8256 let mut again = queue.get(&t.id).unwrap();
8257 again.hold_machine(Some("again".into()));
8258 queue.put(&mut again).unwrap();
8259 apply_choice_actions(&queue, &questions, &home);
8260 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
8261 }
8262}