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 missing = crate::queue::missing_blockers(queue, questions, &task.blocked_by);
691 if !missing.is_empty() {
692 let language = language_of(&task, Path::new("."));
693 task.hold_machine(Some(crate::queue::missing_blocker_hold_reason_in(
694 &task.blocked_by,
695 &missing,
696 &language,
697 )));
698 record(queue, &mut task);
699 continue;
700 }
701 if let Some(q) = task.blocked_by.iter().find_map(|id| {
704 questions.get(id).ok().filter(|q| {
705 q.node == crate::conduct::NODE && q.status == ask::QuestionStatus::Abandoned
706 })
707 }) {
708 let language = language_of(&task, Path::new("."));
709 task.hold_machine(Some(unanswered_question_hold_reason(&q, &language)));
710 record(queue, &mut task);
711 continue;
712 }
713 let mut changed = false;
714 for id in task.blocked_by.clone() {
715 if let Ok(dep) = queue.get(&id) {
716 if dep.status == TaskStatus::Done {
717 task.unblock(&id);
718 changed = true;
719 }
720 continue;
721 }
722 if let Ok(q) = questions.get(&id)
723 && q.status == ask::QuestionStatus::Answered
724 {
725 let answer = match &q.answer {
726 Some(ask::Answer::Choice(c) | ask::Answer::Text(c)) => c.clone(),
727 None => String::new(),
728 };
729 task.record_answer(q.summary.clone(), answer);
730 task.unblock(&id);
731 changed = true;
732 }
733 }
734 if changed {
735 record(queue, &mut task);
736 }
737 }
738}
739
740fn unanswered_question_hold_reason(q: &ask::Question, language: &str) -> String {
742 if crate::lang::is_japanese(language) {
743 format!(
744 "質問 {} 「{}」 に期限内の回答がなく、取り下げられました - `magi task triage` を参照",
745 q.short(),
746 q.summary
747 )
748 } else {
749 format!(
750 "question {} \"{}\" went unanswered and was abandoned - see `magi task triage`",
751 q.short(),
752 q.summary
753 )
754 }
755}
756
757#[derive(Debug, Clone, PartialEq, Eq)]
759enum ActionDecision {
760 Skip,
762 Resume(String),
764 Requeue,
766 Done,
768 Stale,
771 Refuse(String),
774}
775
776fn decide_action<F>(task: &Task, q: &ask::Question, p: &Phrases, load: F) -> ActionDecision
784where
785 F: FnOnce(&str) -> Result<RunState>,
786{
787 let Some(action) = q.chosen_action() else {
788 return ActionDecision::Skip;
789 };
790 if task.action_applied(&q.id)
791 || matches!(
792 task.status,
793 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
794 )
795 {
796 return ActionDecision::Skip;
797 }
798 if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
801 return ActionDecision::Stale;
802 }
803 match action {
804 ask::ChoiceAction::Requeue => ActionDecision::Requeue,
805 ask::ChoiceAction::Done => ActionDecision::Done,
806 ask::ChoiceAction::Resume { run } => {
807 if task.runs.last() != Some(run) {
808 return ActionDecision::Refuse((p.resume_not_latest)(
809 q.short(),
810 ask::short_id(run),
811 ));
812 }
813 match load(run) {
814 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
815 ActionDecision::Resume(run.clone())
816 }
817 Ok(_) => ActionDecision::Refuse((p.resume_cannot_progress)(
818 q.short(),
819 ask::short_id(run),
820 )),
821 Err(e) => ActionDecision::Refuse((p.resume_unreadable)(
822 q.short(),
823 ask::short_id(run),
824 &format!("{e:#}"),
825 )),
826 }
827 }
828 }
829}
830
831struct Phrases {
845 graph_stopped: fn(&str, &str) -> String,
847 quorum_lost: &'static str,
848 quota_took_out: &'static str,
850 run_ended: &'static str,
852 waiting_for_answer: &'static str,
854 recovered_running: &'static str,
855 no_run_to_recover: &'static str,
856 could_not_start: &'static str,
857 resume_not_latest: fn(&str, &str) -> String,
858 resume_cannot_progress: fn(&str, &str) -> String,
859 resume_unreadable: fn(&str, &str, &str) -> String,
860}
861
862const PHRASES_EN: Phrases = Phrases {
863 graph_stopped: |status, detail| {
864 format!("the graph stopped at `{status}` without reaching a terminal status: {detail}")
865 },
866 quorum_lost: "the judging panel lost its quorum",
867 quota_took_out: "; quota took out ",
868 run_ended: "run ended ",
869 waiting_for_answer: " - waiting for operator answer to question ",
870 recovered_running: "recovered a `running` task whose daemon never recorded the outcome: ",
871 no_run_to_recover: "task was `running` with no live daemon and no readable \
872 run to recover; held for a human to check what happened",
873 could_not_start: "could not start the run: ",
874 resume_not_latest: |q, run| {
875 format!("question {q} asked to resume run {run}, which is not this task's latest run")
876 },
877 resume_cannot_progress: |q, run| {
878 format!("question {q} asked to resume run {run}, which cannot make progress")
879 },
880 resume_unreadable: |q, run, e| {
881 format!("question {q} asked to resume run {run}, which could not be read: {e}")
882 },
883};
884
885const PHRASES_JA: Phrases = Phrases {
886 graph_stopped: |status, detail| {
887 format!("グラフが終端状態に達しないまま `{status}` で停止しました: {detail}")
888 },
889 quorum_lost: "審査パネルが定足数を失いました",
890 quota_took_out: "。クォータで脱落: ",
891 run_ended: "run 終了: ",
892 waiting_for_answer: " - オペレーターの回答待ち: 質問 ",
893 recovered_running: "daemon が結果を記録しないまま `running` だったタスクを回収しました: ",
894 no_run_to_recover: "タスクは `running` でしたが、生きた daemon も回収できる run も見つかりません。\
895 何が起きたか人が確認するため保留にしました",
896 could_not_start: "run を開始できませんでした: ",
897 resume_not_latest: |q, run| {
898 format!(
899 "質問 {q} は run {run} の再開を求めましたが、これはタスクの最新の run ではありません"
900 )
901 },
902 resume_cannot_progress: |q, run| {
903 format!("質問 {q} は run {run} の再開を求めましたが、これは進行できません")
904 },
905 resume_unreadable: |q, run, e| {
906 format!("質問 {q} は run {run} の再開を求めましたが、読み込めませんでした: {e}")
907 },
908};
909
910fn phrases(language: &str) -> &'static Phrases {
911 if crate::lang::is_japanese(language) {
912 &PHRASES_JA
913 } else {
914 &PHRASES_EN
915 }
916}
917
918fn language_of(task: &Task, fallback: &Path) -> String {
922 crate::lang::of_repo(&repo_for(task, fallback))
923}
924
925fn task_of_question<'a>(tasks: &'a [Task], q: &ask::Question) -> Option<&'a Task> {
929 if q.node == crate::conduct::NODE {
930 return tasks.iter().find(|t| t.id == q.run);
931 }
932 tasks.iter().find(|t| t.runs.contains(&q.run))
933}
934
935fn apply_choice_actions(queue: &Queue, questions: &Questions, home: &Path) {
944 let tasks = queue.list();
945 for q in questions.list() {
946 if q.chosen_action().is_none() {
947 continue;
948 }
949 let Some(listed) = task_of_question(&tasks, &q) else {
950 continue;
951 };
952 if listed.action_applied(&q.id) {
953 continue;
954 }
955 let Ok(_claim) = queue.claim(&listed.id) else {
956 continue;
957 };
958 let Ok(mut task) = queue.get(&listed.id) else {
959 continue;
960 };
961 let language = language_of(&task, Path::new("."));
962 let decision = decide_action(&task, &q, phrases(&language), |id| {
963 RunState::load_under(id, home)
964 });
965 let ran = matches!(
966 decision,
967 ActionDecision::Resume(_) | ActionDecision::Requeue | ActionDecision::Done
968 );
969 if ran {
970 if questions
975 .read_lease(&q.id)
976 .is_some_and(|l| l.fresh(Timestamp::now()))
977 {
978 continue;
979 }
980 let taken = questions.update(&q.id, |r| {
981 let free = !r.answer_delivered;
982 r.answer_delivered = true;
983 Ok(free)
984 });
985 if !matches!(taken, Ok((_, true))) {
986 continue;
987 }
988 }
989 match decision {
990 ActionDecision::Skip => continue,
991 ActionDecision::Resume(run) => {
992 task.release();
993 task.resume_override = Some(crate::queue::OperatorResume {
994 question_id: q.id.clone(),
995 at: Timestamp::now(),
996 conductor_rehold: None,
997 forced: true,
998 pinned_run: Some(run),
999 });
1000 }
1001 ActionDecision::Stale => {}
1002 ActionDecision::Requeue => task.requeue(),
1003 ActionDecision::Done => {
1004 task.succeed();
1005 supersede_prior_runs(&task, home);
1006 }
1007 ActionDecision::Refuse(why) => {
1008 task.hold_machine(Some(why));
1009 notices::raise(
1010 Notice::warn(
1011 &format!("action:{}", q.id),
1012 "An answer asked the daemon to resume a run that cannot be resumed; the task stays held.",
1013 )
1014 .link(Link::Task {
1015 id: task.id.clone(),
1016 }),
1017 );
1018 }
1019 }
1020 task.mark_action_applied(&q.id);
1021 record(queue, &mut task);
1022 }
1023}
1024
1025fn reconcile_task_questions(queue: &Queue, questions: &Questions) {
1036 let tasks = queue.list();
1037 let by_id: std::collections::BTreeMap<&str, &Task> =
1038 tasks.iter().map(|t| (t.id.as_str(), t)).collect();
1039 let referenced: std::collections::BTreeSet<&str> = tasks
1040 .iter()
1041 .flat_map(|task| task.blocked_by.iter().map(String::as_str))
1042 .collect();
1043
1044 for mut question in questions.list() {
1045 if !question.status.open() || question.node != crate::conduct::NODE {
1046 continue;
1047 }
1048 if referenced.contains(question.id.as_str()) {
1051 continue;
1052 }
1053 let Some(task) = by_id.get(question.run.as_str()) else {
1054 continue;
1055 };
1056 question.abandon(format!(
1057 "task {} no longer waits for this answer",
1058 task.short()
1059 ));
1060 if let Err(e) = questions.put(&mut question) {
1061 tracing::warn!(
1062 "could not retire question {} for task {}: {e:#}",
1063 question.short(),
1064 task.short()
1065 );
1066 }
1067 }
1068}
1069
1070#[derive(Debug, Clone, Copy)]
1077pub struct Verdict {
1078 pub status: RunStatus,
1080 pub left_pr: bool,
1082 pub quota_hit: bool,
1084 pub parked: bool,
1086 pub no_viable_candidates: bool,
1093}
1094
1095pub fn settle(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize) {
1152 settle_in(task, verdict, detail, max_attempts, &PHRASES_EN)
1153}
1154
1155fn settle_in(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize, p: &Phrases) {
1157 if verdict.parked {
1163 task.stall(detail);
1164 return;
1165 }
1166 match verdict.status {
1167 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1168 RunStatus::AlreadyInBase => task.already_landed(detail),
1169 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1170 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1171 task.stall(detail)
1172 }
1173 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1174 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1175 RunStatus::Blocked => task.fail(detail, max_attempts),
1176 RunStatus::VerifiedNoop => task.handed_off(detail),
1177 other => task.fail((p.graph_stopped)(label(other), detail), max_attempts),
1178 }
1179}
1180
1181pub fn supersede_prior_runs(task: &Task, home: &Path) {
1230 let now = Timestamp::now();
1231 let last_run_succeeded = task
1238 .runs
1239 .last()
1240 .and_then(|id| RunState::load_under(id, home).ok())
1241 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1242 for id in task.superseded_attempts(last_run_succeeded) {
1243 let mut state = match RunState::load_under(id, home) {
1244 Ok(s) => s,
1245 Err(e) => {
1246 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1247 continue;
1248 }
1249 };
1250 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1251 continue;
1252 }
1253 let daemon_claims = is_working_on(home, id, now);
1254 if state.liveness(daemon_claims) == Liveness::Live {
1255 continue;
1256 }
1257 state.status = RunStatus::Superseded;
1258 if let Err(e) = state.save_under(home) {
1259 tracing::warn!("could not mark run {id} superseded: {e:#}");
1260 }
1261 }
1262}
1263
1264fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1286 for task in queue.list() {
1287 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1288 continue;
1289 }
1290 supersede_prior_runs(&task, home);
1291 }
1292}
1293
1294fn settle_and_diagnose(
1301 task: &mut Task,
1302 verdict: Verdict,
1303 detail: &str,
1304 max_attempts: usize,
1305 state: &RunState,
1306) {
1307 let p = phrases(&state.config.graph.language);
1308 settle_in(task, verdict, detail, max_attempts, p);
1309 if task.status == TaskStatus::Held {
1310 task.diagnostic = diagnostic(state);
1311 note_open_question(task, &state.id, p);
1312 }
1313}
1314
1315fn note_open_question(task: &mut Task, run: &str, p: &Phrases) {
1330 let Some(home) = crate::run::try_home() else {
1331 return;
1332 };
1333 let open = Questions::at(home.join("questions")).open_for(run);
1334 let Some(q) = open.first() else {
1335 return;
1336 };
1337 let base = task.hold_reason.clone().unwrap_or_default();
1338 task.hold_reason = Some(format!("{base}{}{}", p.waiting_for_answer, q.short()));
1339}
1340
1341fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize, language: &str) {
1352 match last_run {
1353 Some(state) => {
1354 let verdict = Verdict {
1355 status: state.status,
1356 left_pr: state.pr.is_some(),
1357 quota_hit: !state.quota.is_empty(),
1358 parked: state.parked,
1359 no_viable_candidates: state.viable().is_empty(),
1360 };
1361 let detail = format!(
1362 "{}{}",
1363 phrases(&state.config.graph.language).recovered_running,
1364 describe(&state)
1365 );
1366 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1367 }
1368 None => {
1369 let why = phrases(language).no_run_to_recover;
1372 task.last_error = Some(why.to_owned());
1373 task.hold_machine(Some(why.to_owned()));
1376 }
1377 }
1378}
1379
1380fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1402 let mut reclaimed = Vec::new();
1403 for listed in queue.list() {
1404 if listed.status != TaskStatus::Running {
1405 continue;
1406 }
1407 let Ok(_claim) = queue.claim(&listed.id) else {
1408 continue;
1409 };
1410 let Ok(mut task) = queue.get(&listed.id) else {
1414 continue;
1415 };
1416 if task.status != TaskStatus::Running {
1417 continue;
1418 }
1419 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1420 if let Some(state) = &last_run
1430 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1431 {
1432 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1433 }
1434 let language = if last_run.is_none() {
1435 language_of(&task, Path::new("."))
1436 } else {
1437 String::new()
1438 };
1439 reclaim(&mut task, last_run, max_attempts, &language);
1440 if task.status == TaskStatus::Done {
1441 supersede_prior_runs(&task, &crate::run::home());
1442 }
1443 record(queue, &mut task);
1444 reclaimed.push(task.id.clone());
1445 }
1446 reclaimed
1447}
1448
1449fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1481 reclaim_abandoned_runs_with(
1482 home,
1483 now,
1484 crate::proc::pid_status,
1485 crate::proc::process_started_at,
1486 )
1487}
1488
1489fn reclaim_abandoned_runs_with<F, G>(
1495 home: &Path,
1496 now: Timestamp,
1497 query: F,
1498 identity: G,
1499) -> Vec<String>
1500where
1501 F: Fn(u32) -> Option<bool>,
1502 G: Fn(u32) -> Option<String>,
1503{
1504 let mut abandoned = Vec::new();
1505 for entry in std::fs::read_dir(home.join("runs"))
1506 .into_iter()
1507 .flatten()
1508 .flatten()
1509 {
1510 let id = entry.file_name().to_string_lossy().into_owned();
1511 if !crate::run::is_run_id(&id) {
1512 continue;
1513 }
1514 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1522 continue;
1523 };
1524 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1525 continue;
1526 };
1527 if state.status.done() || !state.active_all_overrun(now) {
1528 continue;
1529 }
1530 let daemon_claims = is_working_on(home, &id, now);
1542 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1543 continue;
1544 }
1545 state.abandon("daemon");
1546 if let Err(e) = state.save_under(home) {
1547 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1548 continue;
1549 }
1550 if let Some(notice) = notices::run_ended(&state) {
1553 notices::raise_in(home, notice);
1554 }
1555 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1563 tracing::warn!("abandon questions for {id}: {e:#}");
1564 }
1565 abandoned.push(id);
1566 }
1567 abandoned
1568}
1569
1570pub async fn serve(opts: Opts) -> Result<()> {
1576 serve_until(opts, Stop::new()).await
1577}
1578
1579pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1596 let signal = {
1597 let stop = stop.clone();
1598 tokio::spawn(async move {
1599 if tokio::signal::ctrl_c().await.is_ok() {
1600 stop.stop();
1601 tracing::info!("shutdown requested; a run in flight will be finished first");
1602 }
1603 })
1604 };
1605
1606 let worktrees_root = opts
1607 .worktrees_root
1608 .clone()
1609 .unwrap_or_else(crate::run::default_worktree_root);
1610 let outcome = drive(
1611 &opts,
1612 &Queue::open(),
1613 &status_path(),
1614 &crate::run::home(),
1615 &worktrees_root,
1616 &stop,
1617 )
1618 .await;
1619
1620 signal.abort();
1621 outcome
1622}
1623
1624async fn drive(
1637 opts: &Opts,
1638 queue: &Queue,
1639 status_file: &Path,
1640 home: &Path,
1641 worktrees_root: &Path,
1642 stop: &Stop,
1643) -> Result<()> {
1644 let status = Arc::new(Mutex::new(Status::new()));
1652 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1653 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1654
1655 let daemon_cfg = prepare(&opts.repo, opts)
1661 .map(|c| c.daemon)
1662 .unwrap_or_default();
1663 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1664
1665 let waiter = tokio::spawn(crate::waiter::run(
1670 crate::waiter::Waiter::new(
1671 crate::ask::Questions::at(home.join("questions")),
1672 home.to_path_buf(),
1673 prepare(&opts.repo, opts).ok(),
1674 ),
1675 stop.clone(),
1676 ));
1677
1678 let deputies = tokio::spawn(crate::deputy::run(
1682 crate::deputy::Deputies::new(
1683 crate::ask::Questions::at(home.join("questions")),
1684 home.to_path_buf(),
1685 prepare(&opts.repo, opts).ok(),
1686 opts.repo.clone(),
1687 daemon_cfg.max_deputies,
1688 {
1689 let stop = stop.clone();
1690 Arc::new(move || stop.parking())
1691 },
1692 ),
1693 stop.clone(),
1694 ));
1695
1696 tracing::info!(
1697 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1698 queue.root().display(),
1699 opts.poll.as_secs(),
1700 opts.max_attempts,
1701 concurrency,
1702 if daemon_cfg.pause_for_interrupts {
1703 ", interrupts enabled"
1704 } else {
1705 ""
1706 }
1707 );
1708
1709 janitor(&opts.repo, opts, home, worktrees_root).await;
1712 resweep_superseded_attempts(queue, home);
1713
1714 let outcome = poll(
1715 opts,
1716 queue,
1717 &status,
1718 home,
1719 worktrees_root,
1720 stop,
1721 DispatchLimits {
1722 max_concurrent: concurrency,
1723 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1724 },
1725 )
1726 .await;
1727
1728 beat.abort();
1729 waiter.abort();
1730 deputies.abort();
1731 clear_status_at(status_file);
1732 outcome
1733}
1734
1735async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1741 loop {
1742 tokio::time::sleep(HEARTBEAT).await;
1743 let snapshot = {
1744 let mut guard = lock(&status);
1745 guard.updated_at = Timestamp::now();
1746 guard.clone()
1747 };
1748 if let Err(e) = write_status_to(&path, &snapshot) {
1749 tracing::warn!("could not refresh the daemon status file: {e:#}");
1752 }
1753 }
1754}
1755
1756#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1759enum LandResume {
1760 NotLanding,
1763 StillWaiting,
1768 Ready,
1772}
1773
1774fn land_resume_state(task: &Task) -> LandResume {
1778 let Some(run_id) = task.runs.last() else {
1779 return LandResume::NotLanding;
1780 };
1781 let Ok(state) = RunState::load(run_id) else {
1782 return LandResume::NotLanding;
1783 };
1784 if state.status != RunStatus::Landing || !state.parked {
1785 return LandResume::NotLanding;
1786 }
1787 let store = ask::Questions::open();
1788 let waiting = store
1789 .list()
1790 .into_iter()
1791 .filter(|q| &q.run == run_id && q.node == land::APPROVAL_NODE)
1792 .max_by(|a, b| a.id.cmp(&b.id));
1793 let Some(mut q) = waiting else {
1794 return LandResume::Ready;
1795 };
1796 if !q.status.open() {
1797 return LandResume::Ready;
1798 }
1799 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
1806 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
1807 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
1808 q.abandon(format!(
1809 "no answer within {}s of asking",
1810 timeout.as_secs().max(1)
1811 ));
1812 if store.put(&mut q).is_ok() {
1815 return LandResume::Ready;
1816 }
1817 }
1818 LandResume::StillWaiting
1819}
1820
1821const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
1829
1830const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
1842
1843struct InFlightGuard<'a> {
1856 status: &'a Arc<Mutex<Status>>,
1857 stop: &'a Stop,
1858 task_id: &'a str,
1859}
1860
1861impl Drop for InFlightGuard<'_> {
1862 fn drop(&mut self) {
1863 lock(self.status).current.retain(|c| c.task != self.task_id);
1864 self.stop.exit();
1865 }
1866}
1867
1868#[derive(Debug, Clone, PartialEq, Eq)]
1890enum Interrupt {
1891 Idle,
1893 Parking {
1905 parked: Vec<String>,
1906 interrupt_task: String,
1907 },
1908 Running {
1916 parked: Vec<String>,
1917 interrupt_task: String,
1918 },
1919 Resuming { parked: Vec<String> },
1926}
1927
1928fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
1949 match state {
1950 Interrupt::Idle => {
1951 if in_flight.len() != 1 {
1963 return Interrupt::Idle;
1964 }
1965 match runnable.iter().find(|t| t.interrupt) {
1966 Some(t) => Interrupt::Parking {
1967 parked: in_flight.to_vec(),
1968 interrupt_task: t.id.clone(),
1969 },
1970 None => Interrupt::Idle,
1971 }
1972 }
1973 Interrupt::Parking {
1974 parked,
1975 interrupt_task,
1976 } => {
1977 if in_flight.iter().any(|id| parked.contains(id)) {
1978 Interrupt::Parking {
1980 parked,
1981 interrupt_task,
1982 }
1983 } else if in_flight.contains(&interrupt_task) {
1984 Interrupt::Running {
1985 parked,
1986 interrupt_task,
1987 }
1988 } else if runnable.iter().any(|t| t.id == interrupt_task) {
1989 Interrupt::Parking {
1993 parked,
1994 interrupt_task,
1995 }
1996 } else {
1997 Interrupt::Resuming { parked }
2002 }
2003 }
2004 Interrupt::Running {
2005 parked,
2006 interrupt_task,
2007 } => {
2008 if in_flight.contains(&interrupt_task) {
2009 Interrupt::Running {
2010 parked,
2011 interrupt_task,
2012 }
2013 } else {
2014 Interrupt::Resuming { parked }
2020 }
2021 }
2022 Interrupt::Resuming { parked } => {
2023 if in_flight.iter().any(|id| parked.contains(id)) {
2024 Interrupt::Idle
2030 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
2031 Interrupt::Resuming { parked }
2032 } else {
2033 Interrupt::Idle
2036 }
2037 }
2038 }
2039}
2040
2041fn advance_interrupt_tick(
2047 enabled: bool,
2048 state: Interrupt,
2049 in_flight: &[String],
2050 runnable: &[Task],
2051) -> Interrupt {
2052 if !enabled {
2053 return Interrupt::Idle;
2054 }
2055 advance_interrupt(state, in_flight, runnable)
2056}
2057
2058fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
2063 match state {
2064 Interrupt::Idle => candidates,
2065 Interrupt::Parking {
2066 parked,
2067 interrupt_task,
2068 } => {
2069 if in_flight.iter().any(|id| parked.contains(id)) {
2070 Vec::new()
2071 } else {
2072 candidates
2073 .into_iter()
2074 .filter(|t| &t.id == interrupt_task)
2075 .collect()
2076 }
2077 }
2078 Interrupt::Running { .. } => Vec::new(),
2079 Interrupt::Resuming { parked } => candidates
2087 .into_iter()
2088 .find(|t| parked.contains(&t.id))
2089 .into_iter()
2090 .collect(),
2091 }
2092}
2093
2094struct DispatchLimits {
2098 max_concurrent: usize,
2101 pause_for_interrupts: bool,
2103}
2104
2105#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2114enum PermitKind {
2115 None,
2120 Urgent,
2127 Ordinary,
2130}
2131
2132fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
2135 if priority {
2136 PermitKind::None
2137 } else if urgent {
2138 PermitKind::Urgent
2139 } else {
2140 PermitKind::Ordinary
2141 }
2142}
2143
2144async fn poll(
2163 opts: &Opts,
2164 queue: &Queue,
2165 status: &Arc<Mutex<Status>>,
2166 home: &Path,
2167 worktrees_root: &Path,
2168 stop: &Stop,
2169 limits: DispatchLimits,
2170) -> Result<()> {
2171 let DispatchLimits {
2172 max_concurrent,
2173 pause_for_interrupts,
2174 } = limits;
2175 let mut attempted: Vec<String> = Vec::new();
2180 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2181 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2189 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2195 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2196 let mut conductor = Conductor::new();
2197 let mut cache_last_checked: Option<Timestamp> = None;
2200 let mut interrupt = Interrupt::Idle;
2202 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2208 std::collections::HashMap::new();
2209
2210 while !stop.stopped() {
2211 lock(status).polls += 1;
2212
2213 while let Some(result) = inflight.try_join_next() {
2218 if let Err(e) = result {
2219 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2220 notices::raise(Notice::error(
2221 "loop:attempt",
2222 "A queued attempt ended abnormally; check the task it was running.",
2223 ));
2224 }
2225 }
2226
2227 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2228 if !swept.is_empty() {
2229 tracing::warn!(
2230 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2231 swept.len(),
2232 swept.join(", ")
2233 );
2234 }
2235 let now = Timestamp::now();
2240
2241 if !stop.busy_now() {
2246 maybe_prune_cache_between_runs(
2247 &opts.repo,
2248 opts,
2249 home,
2250 stop,
2251 &mut cache_last_checked,
2252 now,
2253 )
2254 .await;
2255 }
2256
2257 let stalled = stalled_tasks(queue, home, now);
2258 let stalled_ids: std::collections::BTreeSet<_> =
2259 stalled.iter().map(|task| task.id.clone()).collect();
2260 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2261 if !reclaimed.is_empty() {
2262 tracing::warn!(
2263 "reclaimed {} task(s) left `running` by a daemon that never \
2264 recorded the outcome: {}",
2265 reclaimed.len(),
2266 reclaimed.join(", ")
2267 );
2268 }
2269 let abandoned_runs = reclaim_abandoned_runs(home, now);
2270 if !abandoned_runs.is_empty() {
2271 tracing::warn!(
2272 "failed {} run(s) left behind by a killed process, past every \
2273 active seat's own timeout: {}",
2274 abandoned_runs.len(),
2275 abandoned_runs.join(", ")
2276 );
2277 }
2278
2279 let questions = Questions::at(home.join("questions"));
2284
2285 resolve_blockers(queue, &questions);
2288 apply_choice_actions(queue, &questions, home);
2289 reconcile_task_questions(queue, &questions);
2290
2291 let finished: Vec<Task> = finished_tasks(queue)
2297 .into_iter()
2298 .filter(|task| !stalled_ids.contains(&task.id))
2299 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2303 .collect();
2304 let queued = queued_tasks(queue);
2305 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2309 && conductor.worth_a_look(queue, &stalled, &finished)
2310 {
2311 match prepare(&opts.repo, opts) {
2312 Ok(cfg) => {
2313 conductor
2314 .maybe_run(
2315 &cfg,
2316 &opts.repo,
2317 queue,
2318 &questions,
2319 home,
2320 &queued,
2321 &stalled,
2322 &finished,
2323 opts.max_attempts,
2324 )
2325 .await;
2326 }
2327 Err(e) => {
2328 tracing::warn!("conductor: no config: {e:#}");
2329 notices::raise(Notice::warn(
2330 "loop:no-config",
2331 "The loop could not read this repository's config, so held tasks are not being triaged.",
2332 ));
2333 }
2334 }
2335 }
2336
2337 let candidates: Vec<Task> = runnable(queue)
2338 .into_iter()
2339 .filter(|t| !opts.once || !attempted.contains(&t.id))
2340 .collect();
2341
2342 let in_flight: Vec<String> = lock(status)
2347 .current
2348 .iter()
2349 .map(|c| c.task.clone())
2350 .collect();
2351 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2352
2353 interrupt =
2354 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2355 if let Interrupt::Parking {
2356 parked,
2357 interrupt_task,
2358 } = &interrupt
2359 {
2360 let reason = format!(
2361 "task {} asked to run first",
2362 crate::run::short_of(interrupt_task)
2363 );
2364 for id in parked {
2365 if let Some(pause) = interrupt_pauses.get(id) {
2366 pause.park_because(reason.clone());
2367 }
2368 }
2369 }
2370 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2371
2372 let cooling_down =
2373 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2374
2375 let mut started_any = false;
2376 for candidate in candidates {
2377 if stop.stopped() {
2378 break;
2379 }
2380
2381 let resume = land_resume_state(&candidate);
2382 if resume == LandResume::StillWaiting {
2383 continue;
2384 }
2385 let priority = resume == LandResume::Ready;
2386
2387 if !priority && cooling_down {
2388 continue;
2389 }
2390 let permit = match permit_kind(priority, candidate.urgent) {
2391 PermitKind::None => None,
2392 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2393 Ok(p) => Some(p),
2394 Err(_) => continue,
2400 },
2401 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2402 Ok(p) => Some(p),
2403 Err(_) => continue,
2407 },
2408 };
2409
2410 let Ok(claim) = queue.claim(&candidate.id) else {
2415 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2416 continue;
2417 };
2418 let mut task = match queue.get(&candidate.id) {
2421 Ok(t) if t.status.runnable() => t,
2422 Ok(_) => continue,
2423 Err(e) => {
2424 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2425 continue;
2426 }
2427 };
2428 let task_id = task.id.clone();
2429 attempted.push(task_id.clone());
2430 lock(status).idle = false;
2431 stop.enter();
2434 started_any = true;
2435
2436 let run_pause = crate::graph::Pause::new();
2440 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2441
2442 let opts = opts.clone();
2443 let queue = queue.clone();
2444 let status = Arc::clone(status);
2445 let stop = stop.clone();
2446 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2447 inflight.spawn(async move {
2448 let _claim = claim;
2452 let _permit = permit;
2453 let _inflight = InFlightGuard {
2455 status: &status,
2456 stop: &stop,
2457 task_id: &task_id,
2458 };
2459 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2460 lock(&status).completed += 1;
2461 let now = Timestamp::now();
2467 if let Some(until) = cooldown_until("a, now) {
2468 let wait = until.as_second() - now.as_second();
2469 *lock("a_cooldown_until) = Some(until);
2470 let hint = quota
2471 .iter()
2472 .find(|q| q.reset.is_some())
2473 .and_then(|q| q.reset.as_deref());
2474 match hint {
2475 Some(h) => tracing::warn!(
2476 "quota hit; waiting {wait}s before taking another ordinary task \
2477 (CLI reported reset: {h})"
2478 ),
2479 None => tracing::warn!(
2480 "quota hit; waiting {wait}s before taking another ordinary task \
2481 (no reset hint reported)"
2482 ),
2483 }
2484 }
2485 });
2486 }
2487
2488 if started_any {
2489 continue;
2490 }
2491
2492 if stop.busy_now() {
2493 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2498 continue;
2499 }
2500
2501 lock(status).idle = true;
2503 if opts.once {
2504 janitor(&opts.repo, opts, home, worktrees_root).await;
2508 resweep_superseded_attempts(queue, home);
2509 triage_held(queue, home, opts).await;
2510 break;
2511 }
2512 stop.idle(opts.poll).await;
2513 if stop.stopped() {
2514 continue;
2515 }
2516 janitor(&opts.repo, opts, home, worktrees_root).await;
2522 resweep_superseded_attempts(queue, home);
2523 triage_held(queue, home, opts).await;
2524 }
2525
2526 while let Some(result) = inflight.join_next().await {
2531 if let Err(e) = result {
2532 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2533 notices::raise(Notice::error(
2534 "loop:attempt",
2535 "A queued attempt ended abnormally; check the task it was running.",
2536 ));
2537 }
2538 }
2539 Ok(())
2540}
2541
2542async fn attempt(
2548 opts: &Opts,
2549 queue: &Queue,
2550 status: &Arc<Mutex<Status>>,
2551 stop: &Stop,
2552 interrupt_pause: crate::graph::Pause,
2553 task: &mut Task,
2554) -> Vec<QuotaLoss> {
2555 let repo = repo_for(task, &opts.repo);
2556 tracing::info!(
2557 "task {} — {} (repo {})",
2558 task.short(),
2559 task.title,
2560 repo.display()
2561 );
2562
2563 let mut config = match prepare(&repo, opts) {
2564 Ok(c) => c,
2565 Err(e) => {
2566 task.attempts += 1;
2570 task.fail(format!("config: {e:#}"), opts.max_attempts);
2571 record(queue, task);
2572 return Vec::new();
2573 }
2574 };
2575 apply_solo(&mut config, task);
2576 let start_failed = phrases(&config.graph.language).could_not_start;
2577
2578 if let Some(reason) = disk_gate(&repo, &config) {
2586 task.last_error = Some(reason.clone());
2587 task.hold_machine(Some(reason.clone()));
2588 record(queue, task);
2589 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2590 notices::raise(
2593 Notice::warn(
2594 &format!("disk:{}", repo.display()),
2595 "A task was held for want of free disk space; free some, then release it from the queue.",
2596 )
2597 .link(Link::Task {
2598 id: task.id.clone(),
2599 }),
2600 );
2601 return Vec::new();
2602 }
2603
2604 let unfinished = (!task.fresh_start)
2624 .then(|| unfinished_run(&task.runs, task.short()))
2625 .flatten();
2626 let review_branch = task.review_branch.take();
2632 let branch_exists = match &review_branch {
2633 Some(branch) => crate::git::branch_exists(&repo, branch)
2634 .await
2635 .unwrap_or(false),
2636 None => false,
2637 };
2638 let starter = choose_starter(
2639 review_branch.as_deref(),
2640 branch_exists,
2641 unfinished.as_deref(),
2642 );
2643 let attachments = match task_attachments(queue, task) {
2648 Ok(a) => a,
2649 Err(e) => {
2650 task.attempts += 1;
2651 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2652 record(queue, task);
2653 return Vec::new();
2654 }
2655 };
2656 let started = match &starter {
2657 Starter::Review(branch) => {
2658 tracing::info!(
2659 "task {} reopens `{branch}` as a review-only pass",
2660 task.short()
2661 );
2662 let takeover = crate::handover::Takeover {
2665 earlier: task.earlier_attempts().to_vec(),
2666 home: crate::run::home(),
2667 choice: take_divergence_answer(branch, &config.merge.remote, task),
2668 };
2669 Runner::review_taking_over(
2670 &repo,
2671 branch,
2672 config,
2673 Some(takeover),
2674 crate::run::Origin::queue(&task.id),
2675 )
2676 .await
2677 }
2678 Starter::Resume(id) => {
2679 tracing::info!("resuming run {id} rather than competing again");
2680 Runner::resume(id).map(|mut r| {
2681 if let Some(instruction) =
2682 prepare_instruction(&starter, Some(&r.state.instruction), task)
2683 {
2684 r.state.instruction = instruction;
2685 }
2686 r.state.attachments = attachments.clone();
2687 r
2688 })
2689 }
2690 Starter::Start => {
2691 if let Some(branch) = &review_branch {
2692 tracing::warn!(
2693 "conductor chose review for task {} but branch `{branch}` no longer \
2694 exists; requeuing as a fresh competition instead",
2695 task.short()
2696 );
2697 }
2698 let instruction = prepare_instruction(&starter, None, task)
2699 .unwrap_or_else(|| task.instruction.clone());
2700 Runner::start_naming(
2701 &repo,
2702 instruction,
2703 &task.title,
2704 config,
2705 crate::run::Origin::queue(&task.id),
2706 )
2707 .await
2708 .map(|mut r| {
2709 r.state.attachments = attachments.clone();
2710 r
2711 })
2712 }
2713 };
2714 let mut runner = match started {
2715 Ok(r) => r,
2716 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2720 let reason = format!("{start_failed}{e:#}");
2721 task.last_error = Some(reason.clone());
2722 let branch = match &starter {
2725 Starter::Review(branch) => Some(branch.clone()),
2726 _ => None,
2727 };
2728 task.hold_for_handover(branch, reason);
2729 record(queue, task);
2730 tracing::warn!(
2731 "holding {} for a branch it cannot take over: {e:#}",
2732 task.short()
2733 );
2734 notices::raise(
2736 Notice::warn(
2737 &format!("handover:{}", task.id),
2738 "A task was held because its branch is still checked out in another worktree that magi would not remove by itself; see the task's hold reason, then release it from the queue.",
2739 )
2740 .link(Link::Task {
2741 id: task.id.clone(),
2742 })
2743 .about([task.id.clone()]),
2744 );
2745 return Vec::new();
2746 }
2747 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2752 let d = e
2753 .downcast_ref::<crate::reconcile::Diverged>()
2754 .expect("checked by the guard");
2755 let mut q = ask::Question::new(
2756 task.id.clone(),
2757 "review".to_owned(),
2758 "sync".to_owned(),
2759 d.summary(),
2760 d.detail(),
2761 d.choices(),
2762 );
2763 match Questions::open().put(&mut q) {
2764 Ok(()) => {
2765 task.last_error = Some(format!("{e:#}"));
2766 task.review_branch = Some(d.branch.clone());
2769 task.block(vec![q.id.clone()], Some(d.summary()));
2770 }
2771 Err(put) => {
2772 tracing::warn!("could not file the divergence question: {put:#}");
2773 task.attempts += 1;
2774 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2775 }
2776 }
2777 record(queue, task);
2778 return Vec::new();
2779 }
2780 Err(e) => {
2781 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
2787 && let Starter::Review(branch) = &starter
2788 {
2789 task.review_branch = Some(branch.clone());
2790 task.last_error = Some(format!("{start_failed}{e:#}"));
2791 task.status = crate::queue::TaskStatus::Failed;
2792 record(queue, task);
2793 return Vec::new();
2794 }
2795 task.attempts += 1;
2796 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2797 record(queue, task);
2798 return Vec::new();
2799 }
2800 };
2801 runner.state.followup_generation = Some(task.followup.as_ref().map_or(0, |f| f.generation));
2804 runner.on_pause(stop.pause());
2806 runner.watch_interrupt(interrupt_pause);
2810
2811 let run = runner.state.id.clone();
2814 task.start(run.clone());
2815 record(queue, task);
2816 lock(status).current.push(Current {
2817 task: task.id.clone(),
2818 run,
2819 });
2820
2821 let quota_before = runner.state.quota.clone();
2824 let detail = match runner.execute().await {
2825 Ok(()) => describe(&runner.state),
2826 Err(e) => format!("{e:#}"),
2827 };
2828 let fresh = losses_this_attempt("a_before, &runner.state.quota);
2829 let verdict = Verdict {
2830 status: runner.state.status,
2831 left_pr: runner.state.pr.is_some(),
2834 quota_hit: !fresh.is_empty(),
2840 parked: runner.state.parked,
2844 no_viable_candidates: runner.state.viable().is_empty(),
2847 };
2848 settle_and_diagnose(task, verdict, &detail, opts.max_attempts, &runner.state);
2849 if task.status == TaskStatus::Done {
2850 supersede_prior_runs(task, &crate::run::home());
2851 }
2852 record(queue, task);
2853 tracing::info!(
2854 "task {} is {} after run {} ({})",
2855 task.short(),
2856 task.status.as_str(),
2857 runner.state.short(),
2858 label(runner.state.status)
2859 );
2860 fresh
2861}
2862
2863fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
2873 after
2874 .iter()
2875 .filter(|q| !before.contains(q))
2876 .cloned()
2877 .collect()
2878}
2879
2880fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
2883 if quota.is_empty() {
2884 return None;
2885 }
2886 let with_hint = quota.iter().find(|q| q.reset.is_some());
2887 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
2888 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
2889 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
2890 Some(
2891 now.checked_add(jiff::SignedDuration::from_secs(secs))
2892 .unwrap_or(Timestamp::MAX),
2893 )
2894}
2895
2896fn apply_solo(config: &mut Config, task: &Task) {
2906 if task.solo {
2907 config.graph.candidates = 1;
2908 }
2909}
2910
2911fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
2913 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
2914 if let Some(mode) = &opts.merge {
2915 config.merge.mode = merge_mode(mode)?;
2916 }
2917 Ok(config)
2918}
2919
2920async fn maybe_prune_cache_between_runs(
2951 repo: &Path,
2952 opts: &Opts,
2953 home: &Path,
2954 stop: &Stop,
2955 last_checked: &mut Option<Timestamp>,
2956 now: Timestamp,
2957) {
2958 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
2959 return;
2960 }
2961 *last_checked = Some(now);
2962 let cfg = match prepare(repo, opts) {
2963 Ok(cfg) => cfg,
2964 Err(e) => {
2965 tracing::warn!("cache check: no config: {e:#}");
2966 return;
2967 }
2968 };
2969 match clean::prune_cache_if_over_limit(&cfg, home) {
2970 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
2971 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
2972 pruned.files,
2973 pruned.freed
2974 ),
2975 Ok(_) => {}
2976 Err(e) => {
2977 tracing::warn!("housekeep: prune cache: {e:#}");
2978 notices::raise_in(
2979 home,
2980 Notice::warn(
2981 "housekeep:cache",
2982 "Pruning the shared build cache failed; disk usage may keep growing.",
2983 ),
2984 );
2985 }
2986 }
2987}
2988
2989fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
2993 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
2994}
2995
2996async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
3015 let cfg = match prepare(repo, opts) {
3016 Ok(cfg) => cfg,
3017 Err(e) => {
3018 tracing::warn!("housekeep: no config: {e:#}");
3019 return;
3020 }
3021 };
3022 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
3031 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
3032 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
3037 let mut extra = Vec::new();
3038 if out.unreadable > 0 {
3039 extra.push(format!("{} unreadable", out.unreadable));
3040 }
3041 if out.orphaned_worktrees > 0 {
3042 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
3043 }
3044 let detail = if extra.is_empty() {
3045 String::new()
3046 } else {
3047 format!(" ({})", extra.join(", "))
3048 };
3049 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
3050 }
3051 if out.external_merges_recorded > 0 {
3052 tracing::info!(
3053 "housekeep: recorded {} run(s) as merged externally",
3054 out.external_merges_recorded
3055 );
3056 }
3057 if out.stale_pr_states_repaired > 0 {
3058 tracing::info!(
3059 "housekeep: rewrote {} run record(s) whose pull request had already settled",
3060 out.stale_pr_states_repaired
3061 );
3062 }
3063 if out.cache_files > 0 {
3064 tracing::info!(
3065 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
3066 out.cache_files,
3067 out.cache_freed
3068 );
3069 }
3070 if out.questions_abandoned > 0 {
3071 tracing::info!(
3072 "housekeep: abandoned {} question(s) left open by a finished run",
3073 out.questions_abandoned
3074 );
3075 }
3076}
3077
3078async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3087 let questions = Questions::at(home.join("questions"));
3088 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3089 if report.is_empty() {
3090 return;
3091 }
3092 if !report.quarantined.is_empty() {
3093 tracing::info!(
3094 "triage: held {} blocked task(s) whose blocked-on task or \
3095 question no longer exists: {}",
3096 report.quarantined.len(),
3097 report.quarantined.join(", ")
3098 );
3099 }
3100 if !report.resumed.is_empty() {
3101 tracing::info!(
3102 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3103 report.resumed.len(),
3104 report.resumed.join(", ")
3105 );
3106 }
3107 if !report.asked.is_empty() {
3108 tracing::info!(
3109 "triage: asked about {} held task(s): {}",
3110 report.asked.len(),
3111 report.asked.join(", ")
3112 );
3113 }
3114 if !report.answered.is_empty() {
3115 tracing::info!(
3116 "triage: applied {} operator answer(s): {}",
3117 report.answered.len(),
3118 report.answered.join(", ")
3119 );
3120 }
3121}
3122
3123fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3130 disk_gate_with(repo, config, crate::disk::free_bytes)
3131}
3132
3133fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3137 repo: &Path,
3138 config: &Config,
3139 free_bytes: F,
3140) -> Option<String> {
3141 let min = config.disk.min_free_bytes;
3142 if min == 0 {
3143 return None;
3144 }
3145 match free_bytes(repo) {
3146 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3147 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3148 }
3149}
3150
3151const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3158
3159const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3163
3164fn quota_wait(
3173 reset_at: Option<Timestamp>,
3174 now: Timestamp,
3175 fallback: Duration,
3176 cap: Duration,
3177) -> Duration {
3178 match reset_at {
3179 Some(at) if at > now => {
3180 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3181 Duration::from_secs(secs).min(cap)
3182 }
3183 _ => fallback,
3184 }
3185}
3186
3187fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3201 parse_reset_hint_zoned(text, now)
3202 .or_else(|| parse_reset_hint_dated(text))
3203 .or_else(|| parse_reset_hint_relative(text, recorded))
3204}
3205
3206fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3210 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3211 let mut rest = rest.trim();
3212 if rest.is_empty() {
3213 return None;
3214 }
3215 let mut total: i64 = 0;
3216 let mut matched = false;
3217 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3218 if let Some((digits, tail)) = rest.split_once(unit)
3219 && !digits.is_empty()
3220 && digits.bytes().all(|b| b.is_ascii_digit())
3221 {
3222 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3223 rest = tail;
3224 matched = true;
3225 }
3226 }
3227 if !rest.is_empty() || !matched {
3228 return None;
3229 }
3230 recorded
3231 .checked_add(jiff::SignedDuration::from_secs(total))
3232 .ok()
3233}
3234
3235fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3239 let clock = clock.trim().to_lowercase();
3240 let (digits, pm) = clock
3241 .strip_suffix("am")
3242 .map(|d| (d, false))
3243 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3244 let (h, m) = digits.trim().split_once(':')?;
3245 let mut hour: i8 = h.trim().parse().ok()?;
3246 let minute: i8 = m.trim().parse().ok()?;
3247 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3248 return None;
3249 }
3250 if pm && hour != 12 {
3251 hour += 12;
3252 } else if !pm && hour == 12 {
3253 hour = 0;
3254 }
3255 Some((hour, minute))
3256}
3257
3258fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3263 let open = text.find('(')?;
3264 let close = text.rfind(')')?;
3265 if close <= open {
3266 return None;
3267 }
3268 let zone = text[open + 1..close].trim();
3269 let (hour, minute) = parse_12h_clock(&text[..open])?;
3270 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3271 let candidate = now
3272 .to_zoned(tz)
3273 .with()
3274 .hour(hour)
3275 .minute(minute)
3276 .second(0)
3277 .millisecond(0)
3278 .microsecond(0)
3279 .nanosecond(0)
3280 .build()
3281 .ok()?;
3282 let mut at = candidate.timestamp();
3283 if at <= now {
3284 at += jiff::SignedDuration::from_hours(24);
3285 }
3286 Some(at)
3287}
3288
3289fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3298 let words: Vec<&str> = text.split_whitespace().collect();
3299 if words.len() < 5 {
3300 return None;
3301 }
3302 (0..=words.len() - 5)
3303 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3304}
3305
3306fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3312 if trailing.is_some_and(|next| next.starts_with('(')) {
3313 return None;
3314 }
3315 let month = month_number(window[0])?;
3316 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3317 let day_digits = ["st", "nd", "rd", "th"]
3318 .iter()
3319 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3320 let day: i8 = day_digits.parse().ok()?;
3321 let year_token = window[2];
3322 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3323 return None;
3324 }
3325 let year: i16 = year_token.parse().ok()?;
3326 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3330 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3331 let date = jiff::civil::Date::new(year, month, day).ok()?;
3332 let candidate = date
3333 .at(hour, minute, 0, 0)
3334 .to_zoned(jiff::tz::TimeZone::UTC)
3335 .ok()?;
3336 Some(candidate.timestamp())
3337}
3338
3339fn month_number(name: &str) -> Option<i8> {
3342 const NAMES: [&str; 12] = [
3343 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3344 ];
3345 let lower = name.to_lowercase();
3346 NAMES
3347 .iter()
3348 .position(|n| *n == lower.as_str())
3349 .map(|i| i as i8 + 1)
3350}
3351
3352fn exhausted_review_budget(state: &RunState) -> bool {
3364 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3365}
3366
3367fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3404 unfinished_run_with(runs, short, RunState::load)
3405}
3406
3407fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3410where
3411 F: FnOnce(&str) -> Result<RunState>,
3412{
3413 let id = runs.last()?;
3414 match load(id) {
3415 Ok(s)
3420 if s.status.resumable()
3421 && !s.released()
3422 && !exhausted_review_budget(&s)
3423 && s.liveness(false) != crate::run::Liveness::Live =>
3424 {
3425 Some(id.clone())
3426 }
3427 Ok(_) => None,
3428 Err(e) => {
3429 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3430 None
3431 }
3432 }
3433}
3434
3435fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3443where
3444 F: FnOnce(&str) -> Result<RunState>,
3445{
3446 if task.status != TaskStatus::Failed || task.fresh_start || task.review_branch.is_some() {
3447 return false;
3448 }
3449 let Some(id) = task.runs.last() else {
3450 return false;
3451 };
3452 load(id).is_ok_and(|s| {
3453 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3454 })
3455}
3456
3457#[derive(Debug, Clone, PartialEq, Eq)]
3460enum Starter {
3461 Review(String),
3464 Resume(String),
3466 Start,
3468}
3469
3470fn take_divergence_answer(
3474 branch: &str,
3475 remote: &str,
3476 task: &mut Task,
3477) -> Option<crate::reconcile::Choice> {
3478 let summary = crate::reconcile::summary_for(branch, remote);
3479 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3480 (a.question == summary)
3481 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3482 .flatten()
3483 .map(|c| (i, c))
3484 })?;
3485 task.answers.remove(idx);
3486 Some(choice)
3487}
3488
3489fn choose_starter(
3501 review_branch: Option<&str>,
3502 branch_exists: bool,
3503 unfinished: Option<&str>,
3504) -> Starter {
3505 match review_branch {
3506 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3507 Some(_) => Starter::Start,
3508 None => match unfinished {
3509 Some(id) => Starter::Resume(id.to_owned()),
3510 None => Starter::Start,
3511 },
3512 }
3513}
3514
3515fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3518 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3519 return fallback.to_path_buf();
3520 }
3521 task.repo.clone()
3522}
3523
3524const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3528
3529fn answers_block(task: &Task, count: usize) -> String {
3531 let mut s = ANSWERS_HEADER.to_owned();
3532 for a in &task.answers[..count] {
3533 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3534 }
3535 s
3536}
3537
3538fn append_answers(base: &str, task: &Task) -> String {
3541 if task.answers.is_empty() {
3542 return base.to_owned();
3543 }
3544 let mut s = base.to_owned();
3545 s.push_str(&answers_block(task, task.answers.len()));
3546 s
3547}
3548
3549fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3553 for count in (1..=task.answers.len()).rev() {
3554 let block = answers_block(task, count);
3555 if let Some(base) = instruction.strip_suffix(&block) {
3556 return base;
3557 }
3558 }
3559 instruction
3560}
3561
3562fn instruction_for(task: &Task) -> String {
3570 append_answers(&task.instruction, task)
3571}
3572
3573fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3585 append_answers(strip_answers_block(old_instruction, task), task)
3586}
3587
3588fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3590 let paths = queue.attachment_paths(task);
3591 for (name, path) in task.attachments.iter().zip(&paths) {
3592 if !path.is_file() {
3593 bail!(
3594 "attachment `{name}` is recorded on the task but {} is missing",
3595 path.display()
3596 );
3597 }
3598 }
3599 Ok(paths)
3600}
3601
3602fn prepare_instruction(
3613 starter: &Starter,
3614 old_instruction: Option<&str>,
3615 task: &Task,
3616) -> Option<String> {
3617 match starter {
3618 Starter::Start => Some(instruction_for(task)),
3619 Starter::Resume(_) => Some(resumed_instruction(
3620 old_instruction.expect("a resumed run always has a prior instruction"),
3621 task,
3622 )),
3623 Starter::Review(_) => None,
3624 }
3625}
3626
3627fn record(queue: &Queue, task: &mut Task) {
3631 if let Err(e) = queue.put(task) {
3632 tracing::error!("could not record task {}: {e:#}", task.short());
3633 notices::raise(Notice::error(
3634 "loop:record",
3635 "The loop could not save a task's state; check the disk.",
3636 ));
3637 }
3638}
3639
3640fn runnable(queue: &Queue) -> Vec<Task> {
3646 let mut tasks: Vec<Task> = queue
3647 .list()
3648 .into_iter()
3649 .filter(|t| t.status.runnable())
3650 .collect();
3651 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
3652 tasks
3653}
3654
3655fn describe(state: &RunState) -> String {
3669 let p = phrases(&state.config.graph.language);
3670 let mut detail = if state.status == RunStatus::Stalled {
3671 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
3672 seats.sort_unstable();
3673 seats.dedup();
3674 if seats.is_empty() {
3675 p.quorum_lost.to_owned()
3676 } else {
3677 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
3678 }
3679 } else {
3680 format!("{}{}", p.run_ended, state.status.display_label())
3681 };
3682 if let Some(last) = state.events.last() {
3683 detail.push_str(&format!(" ({}: {})", last.node, last.message));
3684 }
3685 detail.push_str(&format!(" [run {}]", state.id));
3686 detail
3687}
3688
3689const DIAGNOSTIC_MAX: usize = 4_000;
3695
3696const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
3701
3702fn diagnostic(state: &RunState) -> Option<String> {
3716 let mut parts: Vec<String> = Vec::new();
3717
3718 for o in state.gate.iter().filter(|o| !o.ok()) {
3720 parts.push(format!(
3721 "gate `{}` failed ({:?}):\n{}",
3722 o.command,
3723 o.code,
3724 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
3725 ));
3726 }
3727
3728 if let Some(last) = state
3731 .events
3732 .iter()
3733 .rev()
3734 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
3735 {
3736 parts.push(last.message.clone());
3737 }
3738
3739 if state.viable().is_empty() {
3746 for c in &state.candidates {
3747 if let Some(evidence) = &c.verified_noop {
3748 parts.push(format!(
3749 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
3750 c.label
3751 ));
3752 } else if !c.summary.trim().is_empty() {
3753 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
3754 } else if let Some(why) = &c.failed {
3755 parts.push(format!("candidate {}: {why}", c.label));
3756 }
3757 }
3758 }
3759
3760 if parts.is_empty() {
3761 return None;
3762 }
3763 Some(crate::run::tail(
3768 &parts.join("\n\n"),
3769 DIAGNOSTIC_MAX.saturating_sub(100),
3770 ))
3771}
3772
3773fn label(status: RunStatus) -> &'static str {
3781 status.as_str()
3782}
3783
3784fn merge_mode(mode: &str) -> Result<MergeMode> {
3786 match mode {
3787 "none" => Ok(MergeMode::None),
3788 "local" => Ok(MergeMode::Local),
3789 "pr" => Ok(MergeMode::Pr),
3790 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
3791 }
3792}
3793
3794fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
3799 mutex
3800 .lock()
3801 .unwrap_or_else(std::sync::PoisonError::into_inner)
3802}
3803
3804#[cfg(test)]
3805mod tests {
3806 use super::*;
3807 use crate::queue::{Source, TaskStatus};
3808 use crate::run::{Candidate, CommandOutcome};
3809 use pretty_assertions::assert_eq;
3810
3811 fn task() -> Task {
3812 Task::new(
3813 "add retries".to_owned(),
3814 "add retries".to_owned(),
3815 PathBuf::from("/repo"),
3816 Source::Human,
3817 )
3818 }
3819
3820 fn interrupt_task(id: &str) -> Task {
3823 let mut t = task();
3824 t.id = id.to_owned();
3825 t.interrupt = true;
3826 t
3827 }
3828
3829 fn task_with_id(id: &str) -> Task {
3831 let mut t = task();
3832 t.id = id.to_owned();
3833 t
3834 }
3835
3836 fn urgent_task(id: &str) -> Task {
3838 let mut t = task();
3839 t.id = id.to_owned();
3840 t.urgent = true;
3841 t
3842 }
3843
3844 #[test]
3849 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
3850 assert_eq!(permit_kind(true, false), PermitKind::None);
3851 assert_eq!(permit_kind(true, true), PermitKind::None);
3852 }
3853
3854 #[test]
3859 fn permit_kind_separates_urgent_from_ordinary() {
3860 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
3861 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
3862 }
3863
3864 #[test]
3871 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
3872 let cfg = Config::default();
3873 let repo = Path::new("/any/repo/path");
3874
3875 let reason =
3876 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
3877 assert!(reason.contains("1024"), "{reason}");
3878 assert!(
3879 reason.contains(&cfg.disk.min_free_bytes.to_string()),
3880 "{reason}"
3881 );
3882
3883 assert_eq!(
3884 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
3885 None,
3886 "exactly at the floor is open"
3887 );
3888 assert_eq!(
3889 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
3890 None,
3891 "comfortably above the floor is open"
3892 );
3893 }
3894
3895 #[test]
3896 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
3897 let mut cfg = Config::default();
3898 cfg.disk.min_free_bytes = 0;
3899 let repo = Path::new("/any/repo/path");
3900 assert_eq!(
3901 disk_gate_with(repo, &cfg, |_| Ok(0)),
3902 None,
3903 "a zero floor never measures at all"
3904 );
3905 }
3906
3907 #[test]
3908 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
3909 let cfg = Config::default();
3910 let repo = Path::new("/any/repo/path");
3911 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
3912 .expect("a measurement failure must close the gate, not open it");
3913 assert!(reason.contains("could not measure"), "{reason}");
3914 }
3915
3916 #[test]
3917 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
3918 let ordinary = task();
3919 let next = advance_interrupt(
3920 Interrupt::Idle,
3921 std::slice::from_ref(&ordinary.id),
3922 std::slice::from_ref(&ordinary),
3923 );
3924 assert_eq!(next, Interrupt::Idle);
3925 }
3926
3927 #[test]
3928 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
3929 let marked = interrupt_task("marked");
3932 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
3933 assert_eq!(next, Interrupt::Idle);
3934 }
3935
3936 #[test]
3937 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
3938 let marked = interrupt_task("marked");
3939 let next = advance_interrupt(
3940 Interrupt::Idle,
3941 &["running".to_owned()],
3942 std::slice::from_ref(&marked),
3943 );
3944 assert_eq!(
3945 next,
3946 Interrupt::Parking {
3947 parked: vec!["running".to_owned()],
3948 interrupt_task: "marked".to_owned(),
3949 }
3950 );
3951 }
3952
3953 #[test]
3962 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
3963 let marked = interrupt_task("marked");
3964
3965 let two = advance_interrupt(
3966 Interrupt::Idle,
3967 &["a".to_owned(), "b".to_owned()],
3968 std::slice::from_ref(&marked),
3969 );
3970 assert_eq!(two, Interrupt::Idle);
3971
3972 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
3973 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
3974 }
3975
3976 #[test]
3977 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
3978 let state = Interrupt::Parking {
3979 parked: vec!["running".to_owned()],
3980 interrupt_task: "marked".to_owned(),
3981 };
3982 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
3984 assert_eq!(still_going, state);
3985
3986 let stopped_but_not_yet_dispatched =
3990 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
3991 assert_eq!(stopped_but_not_yet_dispatched, state);
3992
3993 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
3995 assert_eq!(
3996 dispatched,
3997 Interrupt::Running {
3998 parked: vec!["running".to_owned()],
3999 interrupt_task: "marked".to_owned(),
4000 }
4001 );
4002 }
4003
4004 #[test]
4005 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
4006 let state = Interrupt::Running {
4007 parked: vec!["running".to_owned()],
4008 interrupt_task: "marked".to_owned(),
4009 };
4010 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
4011 assert_eq!(still_running, state);
4012
4013 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
4020 assert_eq!(
4021 ended,
4022 Interrupt::Resuming {
4023 parked: vec!["running".to_owned()]
4024 }
4025 );
4026 }
4027
4028 #[test]
4029 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
4030 let state = Interrupt::Resuming {
4031 parked: vec!["running".to_owned()],
4032 };
4033 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
4034 assert_eq!(still_waiting, state);
4035
4036 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
4037 assert_eq!(dispatched, Interrupt::Idle);
4038 }
4039
4040 #[test]
4046 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
4047 {
4048 let state = Interrupt::Parking {
4049 parked: vec!["running".to_owned()],
4050 interrupt_task: "marked".to_owned(),
4051 };
4052 let next = advance_interrupt(state, &[], &[]);
4055 assert_eq!(
4056 next,
4057 Interrupt::Resuming {
4058 parked: vec!["running".to_owned()]
4059 },
4060 "abandoning the interrupt must not abandon the resume it owes"
4061 );
4062 }
4063
4064 #[test]
4067 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
4068 let state = Interrupt::Resuming {
4069 parked: vec!["running".to_owned()],
4070 };
4071 let next = advance_interrupt(state, &[], &[]);
4072 assert_eq!(
4073 next,
4074 Interrupt::Idle,
4075 "nothing is left to wait for; the loop must not stay wedged"
4076 );
4077 }
4078
4079 #[test]
4080 fn disabled_by_config_the_sequence_can_never_leave_idle() {
4081 let marked = interrupt_task("marked");
4082 let next = advance_interrupt_tick(
4083 false,
4084 Interrupt::Idle,
4085 &["running".to_owned()],
4086 std::slice::from_ref(&marked),
4087 );
4088 assert_eq!(
4089 next,
4090 Interrupt::Idle,
4091 "an unmarked, unconfigured daemon must behave exactly as before"
4092 );
4093 }
4094
4095 #[test]
4096 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4097 let state = Interrupt::Parking {
4098 parked: vec!["running".to_owned()],
4099 interrupt_task: "marked".to_owned(),
4100 };
4101 let candidates = vec![interrupt_task("marked"), task()];
4102 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4103 assert!(
4104 allowed.is_empty(),
4105 "nothing may dispatch - not even the interrupt task itself - \
4106 until the parked run has actually stopped"
4107 );
4108 }
4109
4110 #[test]
4124 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4125 for state in [
4126 Interrupt::Parking {
4127 parked: vec!["running".to_owned()],
4128 interrupt_task: "marked".to_owned(),
4129 },
4130 Interrupt::Running {
4131 parked: vec!["running".to_owned()],
4132 interrupt_task: "marked".to_owned(),
4133 },
4134 Interrupt::Resuming {
4135 parked: vec!["running".to_owned()],
4136 },
4137 ] {
4138 let candidates = vec![
4139 interrupt_task("marked"),
4140 urgent_task("hot"),
4141 task_with_id("ordinary"),
4142 ];
4143 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4144 assert!(
4145 !allowed.iter().any(|t| t.id == "hot"),
4146 "an urgent candidate must wait out the same gate as anything \
4147 else while the run it would run alongside has not actually \
4148 left flight, for state {state:?}: {allowed:?}"
4149 );
4150 }
4151 }
4152
4153 #[test]
4159 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4160 let state = Interrupt::Resuming {
4161 parked: vec!["hot".to_owned()],
4162 };
4163 let candidates = vec![urgent_task("hot"), task()];
4164 let allowed = interrupt_gate(&state, &[], candidates);
4165 assert_eq!(
4166 allowed.iter().filter(|t| t.id == "hot").count(),
4167 1,
4168 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4169 );
4170 }
4171
4172 #[test]
4173 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4174 let state = Interrupt::Parking {
4175 parked: vec!["running".to_owned()],
4176 interrupt_task: "marked".to_owned(),
4177 };
4178 let other = task();
4179 let candidates = vec![interrupt_task("marked"), other.clone()];
4180 let allowed = interrupt_gate(&state, &[], candidates);
4181 assert_eq!(allowed.len(), 1);
4182 assert_eq!(allowed[0].id, "marked");
4183 }
4184
4185 #[test]
4186 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4187 let state = Interrupt::Running {
4188 parked: vec!["running".to_owned()],
4189 interrupt_task: "marked".to_owned(),
4190 };
4191 let candidates = vec![task(), task()];
4192 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4193 assert!(allowed.is_empty());
4194 }
4195
4196 #[test]
4203 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4204 let state = Interrupt::Resuming {
4205 parked: vec!["a".to_owned(), "c".to_owned()],
4206 };
4207 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4208 let allowed = interrupt_gate(&state, &[], candidates);
4209 assert_eq!(
4210 allowed.len(),
4211 1,
4212 "at most one candidate may be offered while resuming: {allowed:?}"
4213 );
4214 assert_eq!(allowed[0].id, "a");
4215 }
4216
4217 #[test]
4218 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4219 let state = Interrupt::Resuming {
4220 parked: vec!["a".to_owned()],
4221 };
4222 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4223 assert!(allowed.is_empty());
4224 }
4225
4226 #[test]
4232 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4233 let running = task(); let marked = interrupt_task("marked");
4235
4236 let mut state = Interrupt::Idle;
4237 let in_flight = vec![running.id.clone()];
4239 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4240 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4241 assert!(gated.is_empty(), "still waiting on `running` to park");
4242
4243 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4245 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4246 assert_eq!(
4247 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4248 vec!["marked"],
4249 "only the interrupt task may be offered to the dispatcher now"
4250 );
4251
4252 state = advance_interrupt_tick(
4254 true,
4255 state,
4256 &["marked".to_owned()],
4257 std::slice::from_ref(&running),
4258 );
4259 let gated = interrupt_gate(
4260 &state,
4261 &["marked".to_owned()],
4262 vec![marked.clone(), running.clone()],
4263 );
4264 assert!(
4265 gated.is_empty(),
4266 "the parked run must not be offered back while the interrupt \
4267 task is still running"
4268 );
4269
4270 let other = task_with_id("other");
4274 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4275 assert_eq!(
4276 state,
4277 Interrupt::Resuming {
4278 parked: vec![running.id.clone()]
4279 }
4280 );
4281 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4282 assert_eq!(
4283 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4284 vec![running.id.as_str()],
4285 "exactly the parked run resumes - not the unrelated task, even \
4286 though it was offered first"
4287 );
4288
4289 state = advance_interrupt_tick(
4293 true,
4294 state,
4295 std::slice::from_ref(&running.id),
4296 std::slice::from_ref(&other),
4297 );
4298 assert_eq!(state, Interrupt::Idle);
4299 let gated = interrupt_gate(
4300 &state,
4301 std::slice::from_ref(&running.id),
4302 vec![other.clone()],
4303 );
4304 assert_eq!(
4305 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4306 vec![other.id.as_str()],
4307 "ordinary dispatch is unrestricted again"
4308 );
4309 }
4310
4311 #[test]
4312 fn every_run_status_settles_the_task_it_came_from() {
4313 let table = [
4315 (RunStatus::Merged, TaskStatus::Done, 1),
4316 (RunStatus::Ready, TaskStatus::Done, 1),
4317 (RunStatus::Stalled, TaskStatus::Failed, 0),
4318 (RunStatus::Blocked, TaskStatus::Failed, 1),
4319 (RunStatus::Failed, TaskStatus::Failed, 1),
4320 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4321 (RunStatus::Prep, TaskStatus::Failed, 1),
4322 (RunStatus::Implementing, TaskStatus::Failed, 1),
4323 (RunStatus::Judging, TaskStatus::Failed, 1),
4324 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4325 (RunStatus::Voting, TaskStatus::Failed, 1),
4326 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4327 (RunStatus::Gating, TaskStatus::Failed, 1),
4328 ];
4329 for (run, want, attempts) in table {
4330 let mut t = task();
4331 t.start("20260902-000000-aaaa".to_owned());
4332 settle(
4333 &mut t,
4334 Verdict {
4335 status: run,
4336 left_pr: false,
4337 parked: false,
4338 quota_hit: matches!(run, RunStatus::Stalled),
4339 no_viable_candidates: false,
4340 },
4341 "why",
4342 2,
4343 );
4344 assert_eq!(t.status, want, "task status after {}", label(run));
4345 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4346 }
4347 }
4348
4349 #[test]
4350 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4351 let mut stalled = task();
4352 stalled.start("20260902-000000-aaaa".to_owned());
4353 settle(
4354 &mut stalled,
4355 Verdict {
4356 status: RunStatus::Stalled,
4357 left_pr: false,
4358 parked: false,
4359 quota_hit: true,
4360 no_viable_candidates: false,
4361 },
4362 "quota",
4363 1,
4364 );
4365 assert_eq!(stalled.attempts, 0);
4366 assert!(
4367 stalled.status.runnable(),
4368 "a machine problem must leave the task in line"
4369 );
4370
4371 let mut blocked = task();
4372 blocked.start("20260902-000000-aaaa".to_owned());
4373 settle(
4374 &mut blocked,
4375 Verdict {
4376 status: RunStatus::Blocked,
4377 left_pr: false,
4378 parked: false,
4379 quota_hit: false,
4380 no_viable_candidates: false,
4381 },
4382 "findings open",
4383 1,
4384 );
4385 assert_eq!(blocked.attempts, 1);
4386 assert_eq!(
4387 blocked.status,
4388 TaskStatus::Held,
4389 "the last attempt hands the task to a human"
4390 );
4391 }
4392
4393 #[test]
4394 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4395 let mut delivered = task();
4398 delivered.start("20260903-080619-01c2".to_owned());
4399 settle(
4400 &mut delivered,
4401 Verdict {
4402 status: RunStatus::Blocked,
4403 left_pr: true,
4404 parked: false,
4405 quota_hit: false,
4406 no_viable_candidates: false,
4407 },
4408 "no check status",
4409 4,
4410 );
4411 assert_eq!(
4412 delivered.status,
4413 TaskStatus::Held,
4414 "a pull request waiting on CI or a person is not a retryable failure"
4415 );
4416 assert!(
4417 !delivered.status.runnable(),
4418 "the loop must not pick this task up again"
4419 );
4420 assert_eq!(
4421 delivered.last_error.as_deref(),
4422 Some("no check status"),
4423 "the operator needs to be told what the gate was waiting for"
4424 );
4425
4426 let mut empty_handed = task();
4429 empty_handed.start("20260903-080619-01c2".to_owned());
4430 settle(
4431 &mut empty_handed,
4432 Verdict {
4433 status: RunStatus::Blocked,
4434 left_pr: false,
4435 parked: false,
4436 quota_hit: false,
4437 no_viable_candidates: false,
4438 },
4439 "findings open",
4440 4,
4441 );
4442 assert_eq!(empty_handed.status, TaskStatus::Failed);
4443 assert!(empty_handed.status.runnable());
4444 }
4445
4446 #[test]
4447 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4448 let mut noop = task();
4454 noop.start("20260912-131304-391f".to_owned());
4455 settle(
4456 &mut noop,
4457 Verdict {
4458 status: RunStatus::VerifiedNoop,
4459 left_pr: false,
4460 parked: false,
4461 quota_hit: false,
4462 no_viable_candidates: true,
4463 },
4464 "candidate A: already fixed by b32cfc4, on main",
4465 4,
4466 );
4467 assert_eq!(
4468 noop.status,
4469 TaskStatus::Held,
4470 "an unverified claim is a request for a human, not a failure"
4471 );
4472 assert!(
4473 !noop.status.runnable(),
4474 "the loop must not requeue this on the same unverified claim"
4475 );
4476 assert_eq!(noop.attempts, 1);
4481 }
4482
4483 #[test]
4484 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4485 let mut parked = task();
4490 parked.start("20260903-183634-2d98".to_owned());
4491 settle(
4492 &mut parked,
4493 Verdict {
4494 status: RunStatus::Implementing,
4495 left_pr: false,
4496 quota_hit: false,
4497 parked: true,
4498 no_viable_candidates: false,
4499 },
4500 "parked after `implementing`",
4501 2,
4502 );
4503 assert_eq!(parked.attempts, 0, "a park is refunded");
4504 assert!(
4505 parked.status.runnable(),
4506 "and the task stays in line so the next loop resumes its run"
4507 );
4508 assert_eq!(
4509 parked.last_error.as_deref(),
4510 Some("parked after `implementing`"),
4511 "the card says where it stopped"
4512 );
4513
4514 let mut broken = task();
4518 broken.start("20260903-183634-2d98".to_owned());
4519 settle(
4520 &mut broken,
4521 Verdict {
4522 status: RunStatus::Implementing,
4523 left_pr: false,
4524 quota_hit: false,
4525 parked: false,
4526 no_viable_candidates: false,
4527 },
4528 "returned mid-flight",
4529 2,
4530 );
4531 assert_eq!(broken.attempts, 1);
4532 }
4533
4534 #[test]
4535 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4536 let mut flaky = task();
4541 flaky.start("20260903-123023-e633".to_owned());
4542 settle(
4543 &mut flaky,
4544 Verdict {
4545 status: RunStatus::Stalled,
4546 left_pr: false,
4547 parked: false,
4548 quota_hit: false,
4549 no_viable_candidates: false,
4550 },
4551 "verdict rests on 1 of 3 judges",
4552 2,
4553 );
4554 assert_eq!(
4555 flaky.attempts, 1,
4556 "flakiness spends an attempt, so `max_attempts` still bounds it"
4557 );
4558 assert!(flaky.status.runnable(), "and it is still worth retrying");
4559
4560 let mut limited = task();
4562 limited.start("20260903-123023-e633".to_owned());
4563 settle(
4564 &mut limited,
4565 Verdict {
4566 status: RunStatus::Stalled,
4567 left_pr: false,
4568 parked: false,
4569 quota_hit: true,
4570 no_viable_candidates: false,
4571 },
4572 "judge-2, judge-3 out of quota",
4573 2,
4574 );
4575 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4576 assert!(limited.status.runnable());
4577
4578 let mut worn = task();
4581 for _ in 0..2 {
4582 worn.release();
4583 }
4584 worn.start("20260903-123023-e633".to_owned());
4585 worn.attempts = 2;
4586 settle(
4587 &mut worn,
4588 Verdict {
4589 status: RunStatus::Stalled,
4590 left_pr: false,
4591 parked: false,
4592 quota_hit: false,
4593 no_viable_candidates: false,
4594 },
4595 "no quorum again",
4596 2,
4597 );
4598 assert_eq!(worn.status, TaskStatus::Held);
4599 assert!(!worn.status.runnable());
4600 }
4601
4602 #[test]
4603 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4604 let mut wiped_out = task();
4611 wiped_out.start("20260907-025000-a1b2".to_owned());
4612 settle(
4613 &mut wiped_out,
4614 Verdict {
4615 status: RunStatus::Failed,
4616 left_pr: false,
4617 parked: false,
4618 quota_hit: true,
4619 no_viable_candidates: true,
4620 },
4621 "no candidate produced a change; nothing to judge",
4622 2,
4623 );
4624 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
4625 assert!(
4626 wiped_out.status.runnable(),
4627 "a machine problem must leave the task in line"
4628 );
4629
4630 let mut partial_progress = task();
4636 partial_progress.start("20260907-025500-c3d4".to_owned());
4637 settle(
4638 &mut partial_progress,
4639 Verdict {
4640 status: RunStatus::Failed,
4641 left_pr: false,
4642 parked: false,
4643 quota_hit: true,
4644 no_viable_candidates: false,
4645 },
4646 "gate failed on the winning candidate",
4647 2,
4648 );
4649 assert_eq!(
4650 partial_progress.attempts, 1,
4651 "a candidate that actually produced a change spends the attempt \
4652 even though some other seat hit its quota"
4653 );
4654 assert!(partial_progress.status.runnable());
4655 }
4656
4657 #[test]
4658 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
4659 let mut t = task();
4665 t.start("20260907-025000-a1b2".to_owned());
4666 let mut state = run_state(RunStatus::Failed);
4667 state.quota.push(QuotaLoss {
4668 seat: "cand-a".to_owned(),
4669 node: "implement".to_owned(),
4670 at: Timestamp::now(),
4671 reset: None,
4672 });
4673 assert!(
4674 state.viable().is_empty(),
4675 "no candidate was added, so nothing is viable"
4676 );
4677 reclaim(&mut t, Some(state), 2, "en");
4678 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
4679 assert!(t.status.runnable());
4680 }
4681
4682 #[test]
4683 fn a_held_task_is_never_offered_to_the_loop() {
4684 let dir = tempfile::tempdir().unwrap();
4685 let queue = Queue::at(dir.path().to_path_buf());
4686 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
4687 let mut t = task();
4688 t.id = format!("2026090{n}-000000-000{n}");
4689 t.priority = priority;
4690 queue.put(&mut t).unwrap();
4691 }
4692 let mut held = task();
4693 held.id = "20260909-000000-9999".to_owned();
4694 held.priority = 99;
4695 held.hold_machine(None);
4696 queue.put(&mut held).unwrap();
4697
4698 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
4699 assert_eq!(order.len(), 3);
4700 assert!(!order.contains(&held.id));
4701 assert_eq!(
4702 order.first().cloned(),
4703 queue.next_runnable().map(|t| t.id),
4704 "the loop's first candidate is exactly what the queue offers"
4705 );
4706 assert_eq!(
4707 order,
4708 vec![
4709 "20260902-000000-0002".to_owned(),
4710 "20260903-000000-0003".to_owned(),
4711 "20260901-000000-0001".to_owned(),
4712 ],
4713 "priority first, then oldest, so nothing starves"
4714 );
4715 }
4716
4717 #[test]
4718 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
4719 let dir = tempfile::tempdir().unwrap();
4720 let queue = Queue::at(dir.path().to_path_buf());
4721 let mut old = task();
4722 old.id = "20260101-000000-old0".to_owned();
4723 queue.put(&mut old).unwrap();
4724 let mut fresh = task();
4725 fresh.id = "20260101-000000-new0".to_owned();
4726 queue.put(&mut fresh).unwrap();
4727
4728 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
4732 std::thread::sleep(Duration::from_millis(60));
4733 let live = queue.claim(&fresh.id).unwrap();
4734
4735 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4736 assert_eq!(swept, vec![old.id.clone()]);
4737 assert!(
4738 queue.claim(&old.id).is_ok(),
4739 "an unparseable lock older than the threshold is swept"
4740 );
4741 assert!(
4742 queue.claim(&fresh.id).is_err(),
4743 "a live pid protects its lock regardless of age"
4744 );
4745 drop(live);
4746 }
4747
4748 #[test]
4749 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
4750 let dir = tempfile::tempdir().unwrap();
4760 let queue = Queue::at(dir.path().to_path_buf());
4761 let mut t = task();
4762 t.id = "20260101-000000-live".to_owned();
4763 queue.put(&mut t).unwrap();
4764
4765 let claim = queue.claim(&t.id).unwrap();
4766 std::thread::sleep(Duration::from_millis(60));
4767
4768 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4769 assert!(
4770 swept.is_empty(),
4771 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
4772 );
4773 assert!(
4774 queue.claim(&t.id).is_err(),
4775 "the lock still protects its task"
4776 );
4777 drop(claim);
4778 }
4779
4780 fn injected_dead_pid() -> u32 {
4783 std::process::id().checked_add(1).unwrap_or(1)
4784 }
4785
4786 #[test]
4787 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
4788 let dir = tempfile::tempdir().unwrap();
4789 let queue = Queue::at(dir.path().to_path_buf());
4790 let mut t = task();
4791 t.id = "20260101-000000-dead".to_owned();
4792 queue.put(&mut t).unwrap();
4793 let dead_pid = injected_dead_pid();
4794
4795 std::fs::write(
4800 dir.path().join(format!("{}.lock", t.id)),
4801 dead_pid.to_string(),
4802 )
4803 .unwrap();
4804
4805 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4806 pid != dead_pid
4807 });
4808 assert_eq!(
4809 swept,
4810 vec![t.id.clone()],
4811 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
4812 );
4813 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
4814 }
4815
4816 #[test]
4817 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
4818 let dir = tempfile::tempdir().unwrap();
4819 let queue = Queue::at(dir.path().to_path_buf());
4820 let mut t = task();
4821 t.id = "20260101-000000-late".to_owned();
4822 queue.put(&mut t).unwrap();
4823 let dead_pid = injected_dead_pid();
4824
4825 assert!(
4828 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
4829 "nothing has claimed the task yet"
4830 );
4831
4832 std::fs::write(
4835 dir.path().join(format!("{}.lock", t.id)),
4836 dead_pid.to_string(),
4837 )
4838 .unwrap();
4839
4840 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4844 pid != dead_pid
4845 });
4846 assert_eq!(swept, vec![t.id.clone()]);
4847 }
4848
4849 #[test]
4850 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
4851 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
4856 let dir = tempfile::tempdir().unwrap();
4857 let queue = Queue::at(dir.path().to_path_buf());
4858 let mut t = task();
4859 t.id = "20260101-000000-crsh".to_owned();
4860 t.status = TaskStatus::Running;
4861 t.attempts = 1;
4862 t.runs.push("20260904-000000-4043".to_owned());
4866 queue.put(&mut t).unwrap();
4867 let dead_pid = injected_dead_pid();
4868
4869 std::fs::write(
4872 dir.path().join(format!("{}.lock", t.id)),
4873 dead_pid.to_string(),
4874 )
4875 .unwrap();
4876
4877 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
4883 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
4884
4885 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4886 pid != dead_pid
4887 });
4888 assert_eq!(swept, vec![t.id.clone()]);
4889
4890 let reclaimed = reclaim_orphaned_running(&queue, 2);
4891 assert_eq!(reclaimed, vec![t.id.clone()]);
4892 let after = queue.get(&t.id).unwrap();
4893 assert_eq!(
4894 after.status,
4895 TaskStatus::Held,
4896 "no run.json to recover from, so a human is asked"
4897 );
4898 assert_eq!(
4899 after.runs,
4900 vec!["20260904-000000-4043".to_owned()],
4901 "the crashed run's id is kept as evidence, not discarded"
4902 );
4903 }
4904
4905 #[test]
4906 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
4907 let dir = tempfile::tempdir().unwrap();
4908 let queue = Queue::at(dir.path().to_path_buf());
4909 let mut t = task();
4910 t.id = "20260101-000000-unknown".to_owned();
4911 queue.put(&mut t).unwrap();
4912 let dead_pid = injected_dead_pid();
4913 std::fs::write(
4914 dir.path().join(format!("{}.lock", t.id)),
4915 dead_pid.to_string(),
4916 )
4917 .unwrap();
4918
4919 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
4920 assert!(swept.is_empty(), "an unknown pid must keep its lock");
4921 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
4922 }
4923
4924 fn run_state_in(status: RunStatus, language: &str) -> RunState {
4925 let mut s = run_state(status);
4926 s.config.graph.language = language.to_owned();
4927 s
4928 }
4929
4930 fn unstarted_verdict(status: RunStatus) -> Verdict {
4931 Verdict {
4932 status,
4933 left_pr: false,
4934 quota_hit: false,
4935 parked: false,
4936 no_viable_candidates: false,
4937 }
4938 }
4939
4940 #[test]
4941 fn an_already_in_base_run_finishes_the_task_without_spending_an_attempt() {
4942 let mut t = task();
4943 t.attempts = 1;
4944 settle_in(
4945 &mut t,
4946 unstarted_verdict(RunStatus::AlreadyInBase),
4947 "already in main",
4948 1,
4949 phrases("en"),
4950 );
4951 assert_eq!(t.status, TaskStatus::Done);
4952 assert_eq!(t.attempts, 0);
4953 }
4954
4955 #[test]
4956 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
4957 let reason = |language: &str| {
4958 let mut t = task();
4959 settle_in(
4960 &mut t,
4961 unstarted_verdict(RunStatus::Judging),
4962 "boom",
4963 1,
4964 phrases(language),
4965 );
4966 t.last_error.or(t.hold_reason).unwrap_or_default()
4967 };
4968 assert!(
4969 reason("en").starts_with("the graph stopped at `"),
4970 "{}",
4971 reason("en")
4972 );
4973 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
4974 assert!(reason("日本語").contains("boom"));
4975 assert_eq!(reason("fr"), reason("en"));
4976 }
4977
4978 #[test]
4979 fn describe_follows_the_run_language_and_keeps_the_run_id() {
4980 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
4981 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
4982 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
4983 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
4984 assert!(ja.contains("[run "), "{ja}");
4985 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
4986 assert!(ended.starts_with("run 終了: "), "{ended}");
4987 let mut de = run_state_in(RunStatus::Failed, "de");
4988 let mut en = run_state_in(RunStatus::Failed, "en");
4989 de.id = "same".to_owned();
4990 en.id = "same".to_owned();
4991 assert_eq!(describe(&de), describe(&en));
4992 }
4993
4994 #[test]
4995 fn refusals_and_recovery_prose_follow_the_language() {
4996 let t = held_task_with("r1");
4997 let q = action_question(
4998 "r1",
4999 ask::ChoiceAction::Resume {
5000 run: "r1".to_owned(),
5001 },
5002 );
5003 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
5004 ActionDecision::Refuse(s) => s,
5005 other => panic!("{other:?}"),
5006 };
5007 assert!(refuse(phrases("en")).contains("could not be read: gone"));
5008 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
5009 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
5010
5011 let mut held = task();
5012 reclaim(&mut held, None, 2, "ja");
5013 assert!(held.hold_reason.unwrap().contains("保留にしました"));
5014 let mut held = task();
5015 reclaim(&mut held, None, 2, "xx");
5016 assert!(held.hold_reason.unwrap().contains("held for a human"));
5017 }
5018
5019 fn run_state(status: RunStatus) -> RunState {
5020 let mut state = RunState::new(
5021 PathBuf::from("/repo"),
5022 "main".to_owned(),
5023 "abc1234def".to_owned(),
5024 "add retries".to_owned(),
5025 Config::default(),
5026 );
5027 state.status = status;
5028 state
5029 }
5030
5031 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
5032 Candidate {
5033 index: 0,
5034 label,
5035 agent: "claude".to_owned(),
5036 branch: format!("magi/x/{label}"),
5037 worktree: PathBuf::from("/repo"),
5038 summary: summary.to_owned(),
5039 stat: String::new(),
5040 files: 0,
5041 commits: usize::from(!empty),
5042 empty,
5043 failed: failed.map(str::to_owned),
5044 verified_noop: None,
5045 duration_ms: 0,
5046 folded: false,
5047 }
5048 }
5049
5050 #[test]
5051 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
5052 let mut state = run_state(RunStatus::Blocked);
5053 state.gate = vec![
5054 CommandOutcome {
5055 command: "cargo make check".to_owned(),
5056 code: Some(0),
5057 output_tail: "ok".to_owned(),
5058 duration_ms: 0,
5059 resource_blocked: false,
5060 },
5061 CommandOutcome {
5062 command: "cargo test".to_owned(),
5063 code: Some(101),
5064 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
5065 duration_ms: 0,
5066 resource_blocked: false,
5067 },
5068 ];
5069 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
5070 assert!(d.contains("cargo test"), "{d}");
5071 assert!(
5072 !d.contains("cargo make check"),
5073 "a passing check is not a diagnostic: {d}"
5074 );
5075 assert!(d.contains("assertion failed"), "{d}");
5076 }
5077
5078 #[test]
5079 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
5080 let mut state = run_state(RunStatus::Blocked);
5081 state.event(
5082 "land",
5083 "stopped: the fixer produced no commit while 2 check(s) were failing \
5084 (build, lint); stopping instead of looping on an unchanged tree",
5085 );
5086 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
5087 assert!(d.contains("build"), "{d}");
5088 assert!(d.contains("lint"), "{d}");
5089 assert!(d.contains("fixer produced no commit"), "{d}");
5090 }
5091
5092 #[test]
5093 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
5094 let state = run_state(RunStatus::VerifiedNoop);
5099 let d = describe(&state);
5100 assert!(
5101 d.contains("agent-verified no-op"),
5102 "expected the display label, not the wire spelling: {d}"
5103 );
5104 assert!(!d.contains("verified_noop"), "{d}");
5105 }
5106
5107 #[test]
5108 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5109 let mut state = run_state(RunStatus::Failed);
5115 state.candidates = vec![candidate(
5116 'A',
5117 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5118 true,
5119 None,
5120 )];
5121 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5122 assert!(d.contains("candidate A"), "{d}");
5123 assert!(d.contains("tagged v1.2.3"), "{d}");
5124 }
5125
5126 #[test]
5127 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5128 let mut state = run_state(RunStatus::Failed);
5129 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5130 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5131 assert!(d.contains("candidate A"), "{d}");
5132 assert!(d.contains("agent timed out"), "{d}");
5133 }
5134
5135 #[test]
5136 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5137 let mut state = run_state(RunStatus::Failed);
5140 state.candidates = vec![candidate('A', "did the work", false, None)];
5141 assert!(diagnostic(&state).is_none());
5142 }
5143
5144 #[test]
5145 fn diagnostic_is_bounded_however_much_a_run_printed() {
5146 let mut state = run_state(RunStatus::Blocked);
5147 state.gate = vec![
5148 CommandOutcome {
5149 command: "cargo test".to_owned(),
5150 code: Some(101),
5151 output_tail: "x".repeat(50_000),
5152 duration_ms: 0,
5153 resource_blocked: false,
5154 },
5155 CommandOutcome {
5156 command: "cargo clippy".to_owned(),
5157 code: Some(1),
5158 output_tail: "y".repeat(50_000),
5159 duration_ms: 0,
5160 resource_blocked: false,
5161 },
5162 ];
5163 state.candidates = vec![
5164 candidate('A', &"z".repeat(50_000), true, None),
5165 candidate('B', &"w".repeat(50_000), true, None),
5166 ];
5167 let d = diagnostic(&state).expect("plenty here to diagnose");
5168 assert!(
5169 d.len() <= DIAGNOSTIC_MAX,
5170 "diagnostic grew to {} bytes, unbounded",
5171 d.len()
5172 );
5173 }
5174
5175 #[test]
5176 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5177 let mut state = run_state(RunStatus::Blocked);
5178 state.gate = vec![CommandOutcome {
5179 command: "cargo test".to_owned(),
5180 code: Some(101),
5181 output_tail: "assertion failed".to_owned(),
5182 duration_ms: 0,
5183 resource_blocked: false,
5184 }];
5185 let verdict = Verdict {
5186 status: RunStatus::Blocked,
5187 left_pr: false,
5188 quota_hit: false,
5189 parked: false,
5190 no_viable_candidates: false,
5191 };
5192
5193 let mut t = task();
5196 t.start("run-1".to_owned());
5197 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5198 assert_eq!(t.status, TaskStatus::Failed);
5199 assert!(t.diagnostic.is_none());
5200
5201 t.start("run-2".to_owned());
5204 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5205 assert_eq!(t.status, TaskStatus::Held);
5206 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5207 assert!(d.contains("cargo test"), "{d}");
5208 }
5209
5210 #[test]
5211 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5212 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5217 let home = crate::run::home();
5218 let state = run_state(RunStatus::VerifiedNoop);
5219 let mut q = ask::Question::new(
5220 state.id.clone(),
5221 "implement".to_owned(),
5222 "impl-A".to_owned(),
5223 "is this really a no-op?".to_owned(),
5224 String::new(),
5225 Vec::new(),
5226 );
5227 Questions::at(home.join("questions")).put(&mut q).unwrap();
5228
5229 let verdict = Verdict {
5230 status: RunStatus::VerifiedNoop,
5231 left_pr: false,
5232 quota_hit: false,
5233 parked: false,
5234 no_viable_candidates: false,
5235 };
5236 let mut t = task();
5237 t.start(state.id.clone());
5238 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5239
5240 assert_eq!(t.status, TaskStatus::Held);
5241 let reason = t.hold_reason.expect("a held task must record why");
5242 assert!(
5243 reason.starts_with("run ended agent-verified no-op"),
5244 "the original settle reason must survive unchanged: {reason}"
5245 );
5246 assert!(
5247 reason.contains(q.short()),
5248 "the open question's id must be named so the notice is actionable: {reason}"
5249 );
5250 }
5251
5252 #[test]
5253 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5254 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5255 let state = run_state(RunStatus::VerifiedNoop);
5256
5257 let verdict = Verdict {
5258 status: RunStatus::VerifiedNoop,
5259 left_pr: false,
5260 quota_hit: false,
5261 parked: false,
5262 no_viable_candidates: false,
5263 };
5264 let mut t = task();
5265 t.start(state.id.clone());
5266 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5267
5268 assert_eq!(t.status, TaskStatus::Held);
5269 assert_eq!(
5270 t.hold_reason.as_deref(),
5271 Some("run ended agent-verified no-op"),
5272 "nothing to append when the question was already answered or never asked"
5273 );
5274 }
5275
5276 #[test]
5277 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5278 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5279 let mut first = run_state(RunStatus::Blocked);
5280 first.id = "20260101-000000-sup1".to_owned();
5281 first.save().unwrap();
5282 let mut second = run_state(RunStatus::Merged);
5283 second.id = "20260101-000000-sup2".to_owned();
5284 second.save().unwrap();
5285
5286 let mut t = task();
5287 t.runs = vec![first.id.clone(), second.id.clone()];
5288 t.status = TaskStatus::Done;
5289
5290 supersede_prior_runs(&t, &crate::run::home());
5291
5292 assert_eq!(
5293 RunState::load(&first.id).unwrap().status,
5294 RunStatus::Superseded,
5295 "the first attempt's Blocked no longer needs anyone's attention"
5296 );
5297 assert_eq!(
5298 RunState::load(&second.id).unwrap().status,
5299 RunStatus::Merged,
5300 "the run that actually succeeded is left exactly as it was"
5301 );
5302 }
5303
5304 #[test]
5305 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5306 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5312 let mut first = run_state(RunStatus::Blocked);
5313 first.id = "20260101-000000-sup9".to_owned();
5314 first.driver_pid = Some(std::process::id());
5317 first.driver_started_at = Some(
5318 crate::proc::process_started_at(std::process::id())
5319 .expect("this test process's own start time must be queryable"),
5320 );
5321 first.save().unwrap();
5322 let mut second = run_state(RunStatus::Merged);
5323 second.id = "20260101-000000-supa".to_owned();
5324 second.save().unwrap();
5325
5326 let mut t = task();
5327 t.runs = vec![first.id.clone(), second.id.clone()];
5328 t.status = TaskStatus::Done;
5329
5330 supersede_prior_runs(&t, &crate::run::home());
5331
5332 assert_eq!(
5333 RunState::load(&first.id).unwrap().status,
5334 RunStatus::Blocked,
5335 "a live driver_pid means something is still actually working this run, \
5336 even though no daemon claims it - rewriting under it would just be \
5337 undone the next time that process saves"
5338 );
5339 }
5340
5341 #[test]
5342 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5343 let dir = tempfile::tempdir().unwrap();
5349 let home = dir.path().to_path_buf();
5350 let queue = Queue::at(dir.path().join("queue"));
5351
5352 let mut first = run_state(RunStatus::Blocked);
5353 first.id = "20260101-000000-supd".to_owned();
5354 first.driver_pid = Some(std::process::id());
5355 first.driver_started_at = Some(
5356 crate::proc::process_started_at(std::process::id())
5357 .expect("this test process's own start time must be queryable"),
5358 );
5359 first.save_under(&home).unwrap();
5360 let mut second = run_state(RunStatus::Merged);
5361 second.id = "20260101-000000-supe".to_owned();
5362 second.save_under(&home).unwrap();
5363
5364 let mut t = task();
5365 t.runs = vec![first.id.clone(), second.id.clone()];
5366 t.status = TaskStatus::Done;
5367 queue.put(&mut t).unwrap();
5368
5369 resweep_superseded_attempts(&queue, &home);
5370 assert_eq!(
5371 RunState::load_under(&first.id, &home).unwrap().status,
5372 RunStatus::Blocked,
5373 "still live on the first pass, so still untouched"
5374 );
5375
5376 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5382 stale.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
5383 stale.save_under(&home).unwrap();
5384
5385 resweep_superseded_attempts(&queue, &home);
5386 assert_eq!(
5387 RunState::load_under(&first.id, &home).unwrap().status,
5388 RunStatus::Superseded,
5389 "the second pass catches up what the first one correctly skipped"
5390 );
5391 }
5392
5393 #[test]
5394 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5395 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5396 let mut first = run_state(RunStatus::Blocked);
5397 first.id = "20260101-000000-sup3".to_owned();
5398 first.save().unwrap();
5399 let mut second = run_state(RunStatus::Blocked);
5400 second.id = "20260101-000000-sup4".to_owned();
5401 second.save().unwrap();
5402
5403 let mut t = task();
5404 t.runs = vec![first.id.clone(), second.id.clone()];
5405 t.status = TaskStatus::Failed;
5409
5410 supersede_prior_runs(&t, &crate::run::home());
5411
5412 assert_eq!(
5413 RunState::load(&first.id).unwrap().status,
5414 RunStatus::Blocked
5415 );
5416 assert_eq!(
5417 RunState::load(&second.id).unwrap().status,
5418 RunStatus::Blocked
5419 );
5420 }
5421
5422 #[test]
5423 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5424 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5428 let mut first = run_state(RunStatus::Blocked);
5429 first.id = "20260101-000000-sup5".to_owned();
5430 first.save().unwrap();
5431
5432 let mut t = task();
5433 t.runs = vec![first.id.clone()];
5434 t.status = TaskStatus::Done;
5435
5436 supersede_prior_runs(&t, &crate::run::home());
5437
5438 assert_eq!(
5439 RunState::load(&first.id).unwrap().status,
5440 RunStatus::Blocked,
5441 "a single-attempt task has no earlier run to supersede"
5442 );
5443 }
5444
5445 #[test]
5446 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5447 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5454 let mut first = run_state(RunStatus::Blocked);
5455 first.id = "20260101-000000-supb".to_owned();
5456 first.save().unwrap();
5457 let mut second = run_state(RunStatus::Failed);
5458 second.id = "20260101-000000-supc".to_owned();
5459 second.save().unwrap();
5460
5461 let mut t = task();
5462 t.runs = vec![first.id.clone(), second.id.clone()];
5463 t.status = TaskStatus::Done;
5464
5465 supersede_prior_runs(&t, &crate::run::home());
5466
5467 assert_eq!(
5468 RunState::load(&first.id).unwrap().status,
5469 RunStatus::Blocked,
5470 "the task's last attempt never landed, so there is nothing here \
5471 actually superseding it"
5472 );
5473 }
5474
5475 #[test]
5476 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5477 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5481 let mut failed = run_state(RunStatus::Failed);
5482 failed.id = "20260101-000000-sup6".to_owned();
5483 failed.save().unwrap();
5484 let mut noop = run_state(RunStatus::VerifiedNoop);
5485 noop.id = "20260101-000000-sup7".to_owned();
5486 noop.save().unwrap();
5487 let mut winner = run_state(RunStatus::Ready);
5488 winner.id = "20260101-000000-sup8".to_owned();
5489 winner.save().unwrap();
5490
5491 let mut t = task();
5492 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5493 t.status = TaskStatus::Done;
5494
5495 supersede_prior_runs(&t, &crate::run::home());
5496
5497 assert_eq!(
5498 RunState::load(&failed.id).unwrap().status,
5499 RunStatus::Failed
5500 );
5501 assert_eq!(
5502 RunState::load(&noop.id).unwrap().status,
5503 RunStatus::VerifiedNoop
5504 );
5505 }
5506
5507 fn approval_question(run: &str) -> ask::Question {
5508 ask::Question::new(
5509 run.to_owned(),
5510 land::APPROVAL_NODE.to_owned(),
5511 "land".to_owned(),
5512 "merge?".to_owned(),
5513 String::new(),
5514 vec!["merge".to_owned(), "hold".to_owned()],
5515 )
5516 }
5517
5518 #[test]
5519 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5520 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5521 let mut state = run_state(RunStatus::Landing);
5522 state.id = "20260101-000000-fre1".to_owned();
5523 state.parked = true;
5524 state.save().unwrap();
5525 ask::Questions::open()
5526 .put(&mut approval_question(&state.id))
5527 .unwrap();
5528
5529 let mut t = task();
5530 t.runs.push(state.id.clone());
5531 assert_eq!(
5532 land_resume_state(&t),
5533 LandResume::StillWaiting,
5534 "nobody has answered and the timeout has not passed"
5535 );
5536 }
5537
5538 #[test]
5539 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5540 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5545 let mut state = run_state(RunStatus::Landing);
5546 state.id = "20260101-000000-exp1".to_owned();
5547 state.parked = true;
5548 state.config.graph.answer_timeout = 60;
5549 state.save().unwrap();
5550
5551 let store = ask::Questions::open();
5552 let mut q = approval_question(&state.id);
5553 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
5554 store.put(&mut q).unwrap();
5555
5556 let mut t = task();
5557 t.runs.push(state.id.clone());
5558 assert_eq!(
5559 land_resume_state(&t),
5560 LandResume::Ready,
5561 "an expired question must not be waited on forever"
5562 );
5563
5564 let after = store.get(&q.id).unwrap();
5565 assert!(
5566 !after.status.open(),
5567 "the question is abandoned, not silently ignored"
5568 );
5569 assert!(
5570 after.resolution().is_none(),
5571 "an abandoned question is not read as a decision"
5572 );
5573 }
5574
5575 #[test]
5576 fn reclaim_settles_a_running_task_against_its_last_run() {
5577 let mut t = task();
5578 t.start("20260904-000000-4043".to_owned());
5579 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
5580 assert_eq!(
5581 t.status,
5582 TaskStatus::Done,
5583 "a run that actually finished must not stay `running` forever"
5584 );
5585 }
5586
5587 #[test]
5588 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
5589 let mut t = task();
5593 t.start("20260904-000000-4043".to_owned());
5594 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
5595 assert_eq!(t.status, TaskStatus::Failed);
5596 assert!(t.status.runnable());
5597 }
5598
5599 #[test]
5600 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
5601 let mut t = task();
5602 t.start("20260904-000000-4043".to_owned());
5603 reclaim(&mut t, None, 2, "en");
5604 assert_eq!(t.status, TaskStatus::Held);
5605 assert!(
5606 t.last_error
5607 .as_deref()
5608 .is_some_and(|e| e.contains("running")),
5609 "the operator needs to know why this task was held"
5610 );
5611 }
5612
5613 #[test]
5614 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
5615 let dir = tempfile::tempdir().unwrap();
5616 let queue = Queue::at(dir.path().to_path_buf());
5617
5618 let mut orphaned = task();
5620 orphaned.id = "20260904-000000-orph".to_owned();
5621 orphaned.status = TaskStatus::Running;
5622 orphaned.attempts = 1;
5623 queue.put(&mut orphaned).unwrap();
5624
5625 let mut alive = task();
5626 alive.id = "20260904-000000-live".to_owned();
5627 alive.status = TaskStatus::Running;
5628 alive.attempts = 1;
5629 queue.put(&mut alive).unwrap();
5630 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
5631
5632 let mut queued = task();
5633 queued.id = "20260904-000000-wait".to_owned();
5634 queue.put(&mut queued).unwrap();
5635
5636 let reclaimed = reclaim_orphaned_running(&queue, 2);
5637 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
5638
5639 assert_eq!(
5640 queue.get(&orphaned.id).unwrap().status,
5641 TaskStatus::Held,
5642 "nothing was driving it and there was no run to recover"
5643 );
5644 assert_eq!(
5645 queue.get(&alive.id).unwrap().status,
5646 TaskStatus::Running,
5647 "a live claim must protect the task it belongs to"
5648 );
5649 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
5650 }
5651
5652 fn read_run_under(home: &Path, id: &str) -> RunState {
5658 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
5659 serde_json::from_str(&body).unwrap()
5660 }
5661
5662 #[test]
5663 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
5664 let dir = tempfile::tempdir().unwrap();
5665 let home = dir.path().to_path_buf();
5666 let now = Timestamp::now();
5667 let overrun_seat = || crate::run::ActiveSeat {
5668 node: "implement".to_owned(),
5669 started_at: now - jiff::SignedDuration::new(21_000, 0),
5670 timeout_secs: 3_600,
5671 attempt: 0,
5672 task: None,
5673 command: None,
5674 index: None,
5675 total: None,
5676 };
5677
5678 let mut dead = run_state(RunStatus::Implementing);
5679 dead.id = "20260101-000000-dead".to_owned();
5680 dead.active.insert("impl-A".to_owned(), overrun_seat());
5681 dead.driver_pid = Some(4242);
5684 dead.save_under(&home).unwrap();
5685
5686 let mut alive = run_state(RunStatus::Implementing);
5689 alive.id = "20260101-000000-aliv".to_owned();
5690 alive.active.insert("impl-A".to_owned(), overrun_seat());
5691 alive.save_under(&home).unwrap();
5692 let mut status = Status::new();
5693 status.current = vec![Current {
5694 task: "20260101-000000-task".to_owned(),
5695 run: alive.id.clone(),
5696 }];
5697 write_status_to(&home.join("daemon.json"), &status).unwrap();
5698
5699 let questions = Questions::at(home.join("questions"));
5703 let mut q = ask::Question::new(
5704 dead.id.clone(),
5705 "implement".to_owned(),
5706 "impl-A".to_owned(),
5707 "Which storage backend?".to_owned(),
5708 String::new(),
5709 vec!["SQLite".to_owned(), "Redis".to_owned()],
5710 );
5711 questions.put(&mut q).unwrap();
5712
5713 let abandoned = reclaim_abandoned_runs_with(
5714 &home,
5715 now,
5716 |pid| if pid == 4242 { Some(false) } else { None },
5717 |_| panic!("a query answering Dead outright needs no identity corroboration"),
5718 );
5719 assert_eq!(abandoned, vec![dead.id.clone()]);
5720
5721 let reloaded = read_run_under(&home, &dead.id);
5722 assert_eq!(reloaded.status, RunStatus::Failed);
5723 assert!(reloaded.active.is_empty());
5724 assert!(
5725 !questions.get(&q.id).unwrap().status.open(),
5726 "the failed run's own open question must be settled in the same pass"
5727 );
5728
5729 let still_alive = read_run_under(&home, &alive.id);
5730 assert_eq!(
5731 still_alive.status,
5732 RunStatus::Implementing,
5733 "a live daemon's claim protects it"
5734 );
5735 assert!(!still_alive.active.is_empty());
5736 }
5737
5738 #[test]
5748 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
5749 let dir = tempfile::tempdir().unwrap();
5750 let home = dir.path().to_path_buf();
5751 let now = Timestamp::now();
5752
5753 let mut manual = run_state(RunStatus::Reviewing);
5754 manual.id = "20260101-000000-manl".to_owned();
5755 manual.active.insert(
5756 "review-1".to_owned(),
5757 crate::run::ActiveSeat {
5758 node: "review".to_owned(),
5759 started_at: now - jiff::SignedDuration::new(21_000, 0),
5760 timeout_secs: 3_600,
5761 attempt: 0,
5762 task: None,
5763 command: None,
5764 index: None,
5765 total: None,
5766 },
5767 );
5768 manual.driver_pid = Some(4242);
5772 manual.driver_started_at = Some("2026-09-22T10:00:00Z".to_owned());
5773 manual.save_under(&home).unwrap();
5774
5775 let abandoned = reclaim_abandoned_runs_with(
5776 &home,
5777 now,
5778 |pid| if pid == 4242 { Some(true) } else { None },
5779 |pid| {
5780 if pid == 4242 {
5781 Some("2026-09-22T10:00:00Z".to_owned())
5782 } else {
5783 None
5784 }
5785 },
5786 );
5787 assert!(
5788 abandoned.is_empty(),
5789 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
5790 );
5791
5792 let reloaded = read_run_under(&home, &manual.id);
5793 assert_eq!(reloaded.status, RunStatus::Reviewing);
5794 assert!(!reloaded.active.is_empty());
5795 }
5796
5797 #[test]
5798 fn an_already_claimed_task_is_skipped_rather_than_failed() {
5799 let dir = tempfile::tempdir().unwrap();
5800 let queue = Queue::at(dir.path().to_path_buf());
5801 let mut only = task();
5802 queue.put(&mut only).unwrap();
5803
5804 let _elsewhere = queue.claim(&only.id).unwrap();
5805 let candidates = runnable(&queue);
5806 assert_eq!(candidates.len(), 1, "the task is still runnable");
5807 assert!(
5808 queue.claim(&candidates[0].id).is_err(),
5809 "the loop cannot take a claim somebody else holds"
5810 );
5811
5812 let after = queue.get(&only.id).unwrap();
5813 assert_eq!(after.status, TaskStatus::Queued);
5814 assert_eq!(
5815 after.attempts, 0,
5816 "losing the race is not an attempt at the task"
5817 );
5818 assert_eq!(after.last_error, None);
5819 }
5820
5821 #[test]
5822 fn the_status_file_round_trips_and_its_heartbeat_advances() {
5823 let dir = tempfile::tempdir().unwrap();
5824 let path = dir.path().join("daemon.json");
5825
5826 let mut status = Status::new();
5827 status.idle = false;
5828 status.completed = 7;
5829 status.current = vec![Current {
5830 task: "20260902-000000-t111".to_owned(),
5831 run: "20260902-000001-r111".to_owned(),
5832 }];
5833 write_status_to(&path, &status).unwrap();
5834 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
5835 assert_eq!(first.schema, SCHEMA);
5836 assert_eq!(first.pid, std::process::id());
5837 assert!(!first.idle);
5838 assert_eq!(first.completed, 7);
5839 assert_eq!(first.current, status.current);
5840 assert!(
5841 !path.with_extension("json.tmp").exists(),
5842 "the temp file is renamed, not left behind"
5843 );
5844
5845 std::thread::sleep(Duration::from_millis(5));
5846 status.updated_at = Timestamp::now();
5847 status.polls = 3;
5848 write_status_to(&path, &status).unwrap();
5849 let second: Status =
5850 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
5851 assert!(
5852 second.updated_at > first.updated_at,
5853 "a reader can only detect staleness if the heartbeat moves"
5854 );
5855 assert_eq!(
5856 second.started_at, first.started_at,
5857 "the start time is not a heartbeat"
5858 );
5859 assert_eq!(second.polls, 3);
5860 }
5861
5862 #[test]
5863 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
5864 let dir = tempfile::tempdir().unwrap();
5865
5866 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
5867
5868 let mut status = Status::new();
5869 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
5870 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5871 let stale = read_status(dir.path()).unwrap();
5872 assert!(
5873 !stale.running(Timestamp::now()),
5874 "a minute without a heartbeat is a dead daemon, not a busy one"
5875 );
5876 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
5877
5878 status.updated_at = Timestamp::now();
5879 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5880 let fresh = read_status(dir.path()).unwrap();
5881 assert!(fresh.running(Timestamp::now()));
5882 }
5883
5884 #[test]
5885 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
5886 let dir = tempfile::tempdir().unwrap();
5887 let now = Timestamp::now();
5888 let mine = "20260903-080619-01c2";
5889
5890 assert!(
5891 !is_working_on(dir.path(), mine, now),
5892 "no status file means nobody is working on anything"
5893 );
5894
5895 let mut status = Status::new();
5896 status.current = vec![Current {
5897 task: "20260903-080340-0167".to_owned(),
5898 run: mine.to_owned(),
5899 }];
5900 status.updated_at = now;
5901 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5902 assert!(is_working_on(dir.path(), mine, now));
5903 assert!(
5904 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
5905 "a daemon busy with one run is not working on another"
5906 );
5907
5908 status.updated_at = now - jiff::SignedDuration::from_secs(600);
5911 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5912 assert!(
5913 !is_working_on(dir.path(), mine, now),
5914 "a stale heartbeat is a dead daemon, so its run is a leftover"
5915 );
5916 }
5917
5918 #[test]
5919 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
5920 let dir = tempfile::tempdir().unwrap();
5921 let now = Timestamp::now();
5922
5923 assert!(
5924 !is_working_on_short(dir.path(), "01c2", now),
5925 "no status file means nobody is working on anything"
5926 );
5927
5928 let mut status = Status::new();
5929 status.current = vec![Current {
5930 task: "20260903-080340-0167".to_owned(),
5931 run: "20260903-080619-01c2".to_owned(),
5932 }];
5933 status.updated_at = now;
5934 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5935 assert!(
5936 is_working_on_short(dir.path(), "01c2", now),
5937 "the run's short id is the last block of its full id"
5938 );
5939 assert!(
5940 !is_working_on_short(dir.path(), "3cbf", now),
5941 "a daemon busy with one worktree bay is not working on another"
5942 );
5943 }
5944
5945 #[test]
5946 fn a_newer_status_file_still_yields_a_reading() {
5947 let dir = tempfile::tempdir().unwrap();
5948 std::fs::write(
5951 dir.path().join("daemon.json"),
5952 serde_json::json!({
5953 "schema": 2,
5954 "updated_at": Timestamp::now().to_string(),
5955 "idle": true,
5956 "surprise": { "nested": [1, 2, 3] },
5957 })
5958 .to_string(),
5959 )
5960 .unwrap();
5961
5962 let reading = read_status(dir.path()).expect("a forward-compatible read");
5963 assert!(reading.running(Timestamp::now()));
5964 assert!(reading.idle);
5965 assert!(reading.current.is_empty());
5966 }
5967
5968 #[test]
5969 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
5970 let dir = tempfile::tempdir().unwrap();
5976 std::fs::write(
5977 dir.path().join("daemon.json"),
5978 serde_json::json!({
5979 "schema": 1,
5980 "pid": 4242,
5981 "updated_at": Timestamp::now().to_string(),
5982 "idle": false,
5983 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
5984 "completed": 3,
5985 "polls": 9,
5986 })
5987 .to_string(),
5988 )
5989 .unwrap();
5990
5991 let reading = read_status(dir.path()).expect("an older shape must still parse");
5992 assert!(reading.running(Timestamp::now()));
5993 assert_eq!(
5994 reading.current,
5995 vec![Current {
5996 task: "20260902-140501-aaaa".to_owned(),
5997 run: "20260902-140502-bbbb".to_owned(),
5998 }]
5999 );
6000 }
6001
6002 #[test]
6003 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
6004 let dir = tempfile::tempdir().unwrap();
6005 std::fs::write(
6006 dir.path().join("daemon.json"),
6007 serde_json::json!({
6008 "schema": 1,
6009 "updated_at": Timestamp::now().to_string(),
6010 "idle": true,
6011 "current": null,
6012 })
6013 .to_string(),
6014 )
6015 .unwrap();
6016 let with_null = read_status(dir.path()).expect("null must still parse");
6017 assert!(with_null.current.is_empty());
6018
6019 std::fs::write(
6020 dir.path().join("daemon.json"),
6021 serde_json::json!({
6022 "schema": 1,
6023 "updated_at": Timestamp::now().to_string(),
6024 "idle": true,
6025 })
6026 .to_string(),
6027 )
6028 .unwrap();
6029 let absent = read_status(dir.path()).expect("a missing field must still parse");
6030 assert!(absent.current.is_empty());
6031 }
6032
6033 #[test]
6034 fn a_task_without_a_repository_runs_in_the_daemons_default() {
6035 let fallback = Path::new("/default");
6036 let mut blank = task();
6037 blank.repo = PathBuf::new();
6038 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
6039 let mut dot = task();
6040 dot.repo = PathBuf::from(".");
6041 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
6042 assert_eq!(
6043 repo_for(&task(), fallback),
6044 PathBuf::from("/repo"),
6045 "a task that names a repository keeps it"
6046 );
6047 }
6048
6049 #[test]
6050 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
6051 let mut solo_cfg = Config::default();
6057 solo_cfg.graph.candidates = 3;
6058 let mut solo_task = task();
6059 solo_task.solo = true;
6060 apply_solo(&mut solo_cfg, &solo_task);
6061 assert_eq!(solo_cfg.graph.candidates, 1);
6062
6063 let mut plain_cfg = Config::default();
6064 plain_cfg.graph.candidates = 3;
6065 let plain_task = task();
6066 assert!(!plain_task.solo);
6067 apply_solo(&mut plain_cfg, &plain_task);
6068 assert_eq!(
6069 plain_cfg.graph.candidates, 3,
6070 "a task that did not ask to run alone keeps the config's candidates"
6071 );
6072 }
6073
6074 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
6075 QuotaLoss {
6076 seat: seat.into(),
6077 node: "judge".into(),
6078 at: at.parse().unwrap(),
6079 reset: reset.map(str::to_string),
6080 }
6081 }
6082
6083 #[test]
6084 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
6085 let old: Vec<QuotaLoss> = (1..=4)
6086 .map(|i| {
6087 loss(
6088 &format!("judge-{i}"),
6089 "2026-09-23T05:23:00Z",
6090 Some("2:40pm (Asia/Tokyo)"),
6091 )
6092 })
6093 .collect();
6094 let fresh = losses_this_attempt(&old, &old);
6095 assert!(fresh.is_empty());
6096 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
6097 }
6099
6100 #[test]
6101 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6102 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6103 let now = Timestamp::now();
6104 let mut after = old.clone();
6105 after.push(loss("judge-2", &now.to_string(), None));
6106 let fresh = losses_this_attempt(&old, &after);
6107 assert_eq!(fresh, vec![after[1].clone()]);
6108 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6109 assert_eq!(
6110 until,
6111 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6112 );
6113 }
6114
6115 #[test]
6116 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6117 let before = vec![
6120 loss("judge-1", "2026-09-23T05:23:00Z", None),
6121 loss("judge-2", "2026-09-23T05:24:00Z", None),
6122 ];
6123 let after = vec![
6124 loss("judge-2", "2026-09-23T05:24:00Z", None),
6125 loss("judge-1", "2026-09-24T01:00:00Z", None),
6126 ];
6127 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6128 }
6129
6130 #[test]
6131 fn merge_overrides_are_parsed_or_refused() {
6132 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6133 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6134 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6135 assert!(merge_mode("squash").is_err());
6136 }
6137
6138 #[test]
6139 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6140 let now = Timestamp::now();
6141 let fallback = Duration::from_secs(300);
6142 let cap = Duration::from_secs(1800);
6143
6144 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6146
6147 let soon = now + jiff::SignedDuration::from_secs(600);
6149 assert_eq!(
6150 quota_wait(Some(soon), now, fallback, cap),
6151 Duration::from_secs(600)
6152 );
6153
6154 let past = now - jiff::SignedDuration::from_secs(60);
6157 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6158
6159 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6162 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6163 }
6164
6165 #[test]
6166 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6167 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6168
6169 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6170 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6171
6172 let already_past =
6176 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6177 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6178
6179 assert!(
6180 parse_reset_hint("session limit reached", now, now).is_none(),
6181 "free text with no recognised shape is not guessed at"
6182 );
6183 assert!(
6184 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6185 "an unresolvable zone name is not guessed at either"
6186 );
6187 }
6188
6189 #[test]
6190 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6191 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6192
6193 let at = parse_reset_hint(
6194 "You've hit your usage limit. Visit \
6195 https://chatgpt.com/codex/settings/usage to purchase more \
6196 credits or try again at Sep 19th, 2026 5:10 PM.",
6197 now,
6198 now,
6199 )
6200 .expect("the codex reset wording is a recognised shape");
6201 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6202
6203 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6208 .expect("an explicit year needs no rollover");
6209 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6210
6211 assert!(
6212 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6213 "a two-digit year is not the documented shape and is not guessed at"
6214 );
6215 assert!(
6216 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6217 "a four-letter month name is not the documented three-letter abbreviation"
6218 );
6219 assert!(
6220 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6221 "an explicit zone on the dated shape is a format nobody has \
6222 documented, and is refused rather than guessed at as UTC"
6223 );
6224 }
6225
6226 #[test]
6227 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6228 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6229 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6230
6231 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6232 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6233
6234 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6235 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6236
6237 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6238 assert!(
6239 parse_reset_hint(bad, now, recorded).is_none(),
6240 "{bad:?} must not be guessed at"
6241 );
6242 }
6243 }
6244
6245 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6249 let config = dir.join("magi.toml");
6250 std::fs::write(
6251 &config,
6252 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6253 )
6254 .unwrap();
6255 let opts = Opts {
6256 poll: Duration::from_secs(30),
6257 config: Some(config),
6258 repo: dir.join("repo"),
6262 ..Opts::default()
6263 };
6264 let home = dir.join("home");
6273 let worktrees = dir.join("wt");
6274 (
6275 opts,
6276 Queue::at(dir.join("queue")),
6277 home.join("daemon.json"),
6278 home,
6279 worktrees,
6280 )
6281 }
6282
6283 #[test]
6284 fn a_stop_is_idempotent_and_once_set_stays_set() {
6285 let stop = Stop::new();
6286 assert!(!stop.stopped());
6287
6288 stop.stop();
6289 assert!(stop.stopped());
6290 stop.stop();
6291 assert!(stop.stopped(), "a second stop is not a toggle");
6292
6293 let shared = stop.clone();
6294 assert!(
6295 shared.stopped(),
6296 "a clone is the same stop; that is how the loop and its caller share one"
6297 );
6298 }
6299
6300 #[test]
6301 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6302 let stop = Stop::new();
6303 stop.enter();
6304 assert!(
6305 !stop.finishing(),
6306 "a busy loop nobody has asked to stop is just running"
6307 );
6308
6309 stop.stop();
6310 assert!(
6311 stop.finishing(),
6312 "a stop asked for mid-run has not landed until the run is settled"
6313 );
6314
6315 stop.exit();
6316 assert!(
6317 !stop.finishing(),
6318 "once the run is settled the stop has landed and there is nothing to finish"
6319 );
6320 }
6321
6322 #[test]
6323 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6324 let stop = Stop::new();
6325 stop.enter();
6326 stop.enter();
6327 stop.stop();
6328 assert!(stop.finishing(), "two runs still in flight");
6329
6330 stop.exit();
6331 assert!(
6332 stop.finishing(),
6333 "one run finished, but a sibling is still working"
6334 );
6335
6336 stop.exit();
6337 assert!(
6338 !stop.finishing(),
6339 "the last run out is what actually lands the stop"
6340 );
6341 }
6342
6343 #[tokio::test]
6344 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6345 let dir = tempfile::tempdir().unwrap();
6346 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6347 let stop = Stop::new();
6348 stop.stop();
6349
6350 let began = std::time::Instant::now();
6351 tokio::time::timeout(
6352 Duration::from_secs(2),
6353 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6354 )
6355 .await
6356 .expect("a stopped loop must return, not sit out its poll interval")
6357 .expect("the loop's own setup and teardown must not fail");
6358 assert!(
6359 began.elapsed() < opts.poll,
6360 "returned only after {:?}, which is a poll interval, not a stop",
6361 began.elapsed()
6362 );
6363 }
6364
6365 #[tokio::test]
6366 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6367 let dir = tempfile::tempdir().unwrap();
6368 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6369 let stop = Stop::new();
6370
6371 let asker = {
6374 let stop = stop.clone();
6375 tokio::spawn(async move {
6376 tokio::time::sleep(Duration::from_millis(20)).await;
6377 stop.stop();
6378 })
6379 };
6380
6381 let began = std::time::Instant::now();
6382 tokio::time::timeout(
6383 Duration::from_secs(2),
6384 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6385 )
6386 .await
6387 .expect("a stop asked for while idle must wake the wait")
6388 .expect("the loop's own setup and teardown must not fail");
6389 asker.await.unwrap();
6390 assert!(
6391 began.elapsed() < opts.poll,
6392 "returned only after {:?}, so the stop waited on the sleep",
6393 began.elapsed()
6394 );
6395 }
6396
6397 #[tokio::test]
6398 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6399 let dir = tempfile::tempdir().unwrap();
6400 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6401 let stop = Stop::new();
6402 stop.stop();
6403
6404 tokio::time::timeout(
6405 Duration::from_secs(2),
6406 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6407 )
6408 .await
6409 .expect("a stopped loop must return")
6410 .expect("the loop's own setup and teardown must not fail");
6411
6412 assert!(
6413 home.is_dir(),
6414 "the loop did publish a status file, so its removal is the teardown and not an absence"
6415 );
6416 assert!(
6417 !status_file.exists(),
6418 "a stopped loop clears its status file"
6419 );
6420 assert!(
6421 read_status(&home).is_none(),
6422 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6423 );
6424 }
6425
6426 #[tokio::test]
6427 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6428 let dir = tempfile::tempdir().unwrap();
6429 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6430 opts.once = true;
6431
6432 let mut settled = RunState::new(
6433 dir.path().join("repo"),
6434 "main".to_owned(),
6435 "abc1234".to_owned(),
6436 "fixture".to_owned(),
6437 Config::default(),
6438 );
6439 settled.status = RunStatus::Ready;
6440 let run_dir = home.join("runs").join(&settled.id);
6441 std::fs::create_dir_all(&run_dir).unwrap();
6442 std::fs::write(
6443 run_dir.join("run.json"),
6444 serde_json::to_string_pretty(&settled).unwrap(),
6445 )
6446 .unwrap();
6447 let questions = Questions::at(home.join("questions"));
6448 let mut question = ask::Question::new(
6449 settled.id.clone(),
6450 "review".to_owned(),
6451 "reviewer-1".to_owned(),
6452 "Continue?".to_owned(),
6453 String::new(),
6454 Vec::new(),
6455 );
6456 questions.put(&mut question).unwrap();
6457
6458 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6459 .await
6460 .unwrap();
6461
6462 assert_eq!(
6463 questions.get(&question.id).unwrap().status,
6464 ask::QuestionStatus::Abandoned,
6465 "an empty --once drain still performs startup question cleanup"
6466 );
6467 }
6468
6469 #[test]
6470 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6471 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6472
6473 assert!(
6474 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6475 "never checked before: due at once"
6476 );
6477
6478 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6479 assert!(
6480 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6481 "well inside the interval: not due yet"
6482 );
6483
6484 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6485 assert!(
6486 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6487 "exactly at the edge: not yet due, same convention as `clean::due`"
6488 );
6489
6490 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6491 assert!(
6492 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6493 "past the interval: due again"
6494 );
6495 }
6496
6497 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
6502 let config = dir.join("magi.toml");
6503 std::fs::write(
6509 &config,
6510 format!(
6511 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
6512 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
6513 cache_dir.display()
6514 ),
6515 )
6516 .unwrap();
6517 Opts {
6518 config: Some(config),
6519 repo: dir.join("repo"),
6520 ..Opts::default()
6521 }
6522 }
6523
6524 #[tokio::test]
6525 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
6526 let dir = tempfile::tempdir().unwrap();
6527 let home = dir.path().join("home");
6528 let cache_dir = dir.path().join("cache");
6529 std::fs::create_dir_all(&cache_dir).unwrap();
6530 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6531 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6532
6533 let running = Stop::new();
6536 let mut last_checked = None;
6537 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6538 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
6539 .await;
6540 assert_eq!(
6541 crate::disk::dir_size(&cache_dir),
6542 0,
6543 "over the cap on the first check ever: pruned at once, no idle queue required"
6544 );
6545 assert_eq!(last_checked, Some(t0));
6546
6547 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
6549 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
6550 maybe_prune_cache_between_runs(
6551 &opts.repo,
6552 &opts,
6553 &home,
6554 &running,
6555 &mut last_checked,
6556 too_soon,
6557 )
6558 .await;
6559 assert_eq!(
6560 crate::disk::dir_size(&cache_dir),
6561 10,
6562 "too soon since the last check: left alone rather than rescanned every call"
6563 );
6564 assert_eq!(
6565 last_checked,
6566 Some(t0),
6567 "an idle check does not reset the clock"
6568 );
6569
6570 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6572 maybe_prune_cache_between_runs(
6573 &opts.repo,
6574 &opts,
6575 &home,
6576 &running,
6577 &mut last_checked,
6578 due_again,
6579 )
6580 .await;
6581 assert_eq!(
6582 crate::disk::dir_size(&cache_dir),
6583 0,
6584 "due again: pruned back under the cap"
6585 );
6586 }
6587
6588 #[tokio::test]
6596 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
6597 let dir = tempfile::tempdir().unwrap();
6598 let home = dir.path().join("home");
6599 let cache_dir = dir.path().join("cache");
6600 std::fs::create_dir_all(&cache_dir).unwrap();
6601 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6602 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6603
6604 let stop = Stop::new();
6605 stop.stop();
6606 assert!(
6607 !stop.finishing(),
6608 "no run is in flight at a between-runs boundary, so nothing else \
6609 would tell the operator this stop had not taken effect yet"
6610 );
6611
6612 let mut last_checked = None;
6613 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6614 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
6615 .await;
6616 assert_eq!(
6617 crate::disk::dir_size(&cache_dir),
6618 10,
6619 "over its cap, and due for the first check ever, but a stop outranks \
6620 it: the cap is a standing policy the next start measures again"
6621 );
6622 assert_eq!(
6623 last_checked, None,
6624 "a check that never happened must not claim the interval"
6625 );
6626 }
6627
6628 #[tokio::test]
6642 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
6643 let dir = tempfile::tempdir().unwrap();
6644 let cache_dir = dir.path().join("cache");
6645 std::fs::create_dir_all(&cache_dir).unwrap();
6646 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
6647
6648 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
6649 opts.poll = Duration::from_millis(20);
6650 opts.max_attempts = 1_000;
6651
6652 let queue = Queue::at(dir.path().join("queue"));
6653 let mut t = Task::new(
6654 "x".to_owned(),
6655 "x".to_owned(),
6656 opts.repo.clone(),
6657 Source::Human,
6658 );
6659 queue.put(&mut t).unwrap();
6660
6661 let home = dir.path().join("home");
6662 let worktrees = dir.path().join("wt");
6663 let status_file = home.join("daemon.json");
6664 let stop = Stop::new();
6665 let stopper = {
6666 let stop = stop.clone();
6667 tokio::spawn(async move {
6668 tokio::time::sleep(Duration::from_millis(400)).await;
6669 stop.stop();
6670 })
6671 };
6672
6673 tokio::time::timeout(
6674 Duration::from_secs(10),
6675 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6676 )
6677 .await
6678 .expect("the loop must not hang on a queue that keeps producing failing work")
6679 .expect("the loop's own setup and teardown must not fail");
6680 stopper.await.unwrap();
6681
6682 let after = queue.get(&t.id).unwrap();
6683 assert!(
6684 after.attempts >= 2,
6685 "the harness must actually have retried more than once, or this is not \
6686 exercising a busy queue at all (got {} attempt(s))",
6687 after.attempts
6688 );
6689 assert!(
6690 after.status.runnable(),
6691 "still under its attempt budget: the queue never reached a natural idle \
6692 on its own, only the external stop ended the test"
6693 );
6694
6695 assert_eq!(
6696 crate::disk::dir_size(&cache_dir),
6697 0,
6698 "an oversized cache must not be left to grow unboundedly just because the \
6699 queue kept the loop busy the whole time"
6700 );
6701 }
6702
6703 #[test]
6704 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
6705 let dir = tempfile::tempdir().unwrap();
6706 let queue = Queue::at(dir.path().join("queue"));
6707 let questions = Questions::at(dir.path().join("questions"));
6708 let mut task = task();
6709 queue.put(&mut task).unwrap();
6710
6711 let mut task_question = ask::Question::new(
6712 task.id.clone(),
6713 crate::conduct::NODE.to_owned(),
6714 "conduct".to_owned(),
6715 "Which backend?".to_owned(),
6716 String::new(),
6717 Vec::new(),
6718 );
6719 questions.put(&mut task_question).unwrap();
6720 task.block(vec![task_question.id.clone()], None);
6721 queue.put(&mut task).unwrap();
6722
6723 let mut run_question = ask::Question::new(
6724 "20260101-000000-run1".to_owned(),
6725 "review".to_owned(),
6726 "reviewer-1".to_owned(),
6727 "Run question".to_owned(),
6728 String::new(),
6729 Vec::new(),
6730 );
6731 questions.put(&mut run_question).unwrap();
6732
6733 let mut coincidental = ask::Question::new(
6738 task.id.clone(),
6739 "review".to_owned(),
6740 "reviewer-1".to_owned(),
6741 "Unrelated review question".to_owned(),
6742 String::new(),
6743 Vec::new(),
6744 );
6745 questions.put(&mut coincidental).unwrap();
6746
6747 reconcile_task_questions(&queue, &questions);
6748 assert!(questions.get(&task_question.id).unwrap().status.open());
6749 assert!(questions.get(&run_question.id).unwrap().status.open());
6750 assert!(questions.get(&coincidental.id).unwrap().status.open());
6751
6752 task.release();
6753 queue.put(&mut task).unwrap();
6754 reconcile_task_questions(&queue, &questions);
6755 assert_eq!(
6756 questions.get(&task_question.id).unwrap().status,
6757 ask::QuestionStatus::Abandoned
6758 );
6759 assert!(
6760 questions.get(&run_question.id).unwrap().status.open(),
6761 "run questions remain the run janitor's responsibility"
6762 );
6763 assert!(
6764 questions.get(&coincidental.id).unwrap().status.open(),
6765 "a non-conductor question must not be abandoned just because its \
6766 run id coincides with a task id"
6767 );
6768 }
6769
6770 #[test]
6771 fn a_freshly_started_running_task_is_never_stalled() {
6772 let dir = tempfile::tempdir().unwrap();
6773 let mut t = task();
6774 t.start("run-1".to_owned());
6775 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
6778 }
6779
6780 #[test]
6781 fn a_long_running_task_with_no_live_daemon_is_stalled() {
6782 let dir = tempfile::tempdir().unwrap();
6783 let mut t = task();
6784 t.start("run-1".to_owned());
6785 t.updated_at = Timestamp::now()
6786 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6787 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
6788 assert_eq!(
6789 stalled_tasks(
6790 &Queue::at(dir.path().join("q")),
6791 dir.path(),
6792 Timestamp::now()
6793 )
6794 .len(),
6795 0,
6796 "the task was never written to this queue"
6797 );
6798 }
6799
6800 #[test]
6801 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
6802 let dir = tempfile::tempdir().unwrap();
6803 let mut t = task();
6804 t.id = "20260903-080340-0167".to_owned();
6805 t.start("20260903-080619-01c2".to_owned());
6806 t.updated_at = Timestamp::now()
6807 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6808
6809 let mut status = Status::new();
6810 status.current = vec![Current {
6811 task: t.id.clone(),
6812 run: "20260903-080619-01c2".to_owned(),
6813 }];
6814 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6815
6816 assert!(
6817 !is_stalled(&t, dir.path(), Timestamp::now()),
6818 "a live daemon's own heartbeat rules out stalled, however long the task has run"
6819 );
6820 }
6821
6822 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
6826 let path = queue.path_of(id);
6827 let body = std::fs::read_to_string(&path).unwrap();
6828 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
6829 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
6830 v["updated_at"] = serde_json::Value::String(old.to_string());
6831 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
6832 }
6833
6834 #[test]
6835 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
6836 let dir = tempfile::tempdir().unwrap();
6849 let queue = Queue::at(dir.path().join("queue"));
6850 let home = dir.path().join("home");
6851
6852 let mut t = task();
6853 t.id = "20260101-000001-lock".to_owned();
6854 t.start("run-1".to_owned());
6855 queue.put(&mut t).unwrap();
6856 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
6857 std::fs::write(
6858 dir.path().join("queue").join(format!("{}.lock", t.id)),
6859 "not a pid",
6860 )
6861 .unwrap();
6862
6863 let now = Timestamp::now();
6864 assert!(
6865 reclaim_orphaned_running(&queue, 2).is_empty(),
6866 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
6867 and reclaim must leave the task alone"
6868 );
6869 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
6870
6871 let stalled = stalled_tasks(&queue, &home, now);
6872 assert_eq!(
6873 stalled.len(),
6874 1,
6875 "reclaim's inability to claim it yet must not hide it from the conductor"
6876 );
6877 assert_eq!(stalled[0].id, t.id);
6878 }
6879
6880 #[test]
6881 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
6882 let dir = tempfile::tempdir().unwrap();
6883 crate::run::set_home(dir.path().join("run-home"));
6884 let queue = Queue::at(dir.path().join("queue"));
6885 let home = dir.path().join("home");
6886 let questions = Questions::at(dir.path().join("questions"));
6887
6888 let mut t = task();
6889 t.id = "20260101-000003-dead".to_owned();
6890 t.start("missing-run".to_owned());
6891 queue.put(&mut t).unwrap();
6892 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
6893
6894 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
6897 assert_eq!(
6898 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
6899 [&t.id]
6900 );
6901 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
6902 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
6903
6904 crate::conduct::apply(
6907 &queue,
6908 &questions,
6909 &crate::conduct::Verdict {
6910 decisions: vec![crate::conduct::Decision {
6911 id: t.id.clone(),
6912 recovery: Some(crate::conduct::Recovery::Requeue),
6913 ..crate::conduct::Decision::default()
6914 }],
6915 },
6916 )
6917 .unwrap();
6918 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
6919 }
6920
6921 #[test]
6922 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
6923 let dir = tempfile::tempdir().unwrap();
6924 let queue = Queue::at(dir.path().join("queue"));
6925 let home = dir.path().join("home");
6926
6927 let mut fresh = task();
6928 fresh.id = "20260101-000001-aaaa".to_owned();
6929 fresh.start("run-1".to_owned());
6930 queue.put(&mut fresh).unwrap();
6931
6932 let mut old = task();
6933 old.id = "20260101-000002-bbbb".to_owned();
6934 old.start("run-2".to_owned());
6935 queue.put(&mut old).unwrap();
6936 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
6937
6938 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
6939 assert_eq!(stalled.len(), 1);
6940 assert_eq!(stalled[0].id, old.id);
6941 }
6942
6943 #[test]
6944 fn queued_and_finished_task_views_partition_by_status() {
6945 let dir = tempfile::tempdir().unwrap();
6946 let queue = Queue::at(dir.path().join("queue"));
6947
6948 let mut queued = task();
6949 queued.id = "20260101-000001-aaaa".to_owned();
6950 queue.put(&mut queued).unwrap();
6951
6952 let mut failed = task();
6953 failed.id = "20260101-000002-bbbb".to_owned();
6954 failed.start("run-1".to_owned());
6955 failed.fail("gate red", 5);
6956 queue.put(&mut failed).unwrap();
6957
6958 let mut held = task();
6959 held.id = "20260101-000003-cccc".to_owned();
6960 held.hold_machine(None);
6961 queue.put(&mut held).unwrap();
6962
6963 let mut running = task();
6964 running.id = "20260101-000004-dddd".to_owned();
6965 running.start("run-2".to_owned());
6966 queue.put(&mut running).unwrap();
6967
6968 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
6969 assert_eq!(queued_ids, [queued.id.clone()]);
6970
6971 let mut finished_ids: Vec<String> =
6972 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
6973 finished_ids.sort_unstable();
6974 let mut want = vec![failed.id.clone(), held.id.clone()];
6975 want.sort_unstable();
6976 assert_eq!(finished_ids, want);
6977 }
6978
6979 #[test]
6980 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
6981 let dir = tempfile::tempdir().unwrap();
6982 let queue = Queue::at(dir.path().join("queue"));
6983 let questions = ask::Questions::at(dir.path().join("questions"));
6984
6985 let mut dep = task();
6986 dep.id = "20260101-000001-dep0".to_owned();
6987 dep.succeed();
6988 queue.put(&mut dep).unwrap();
6989
6990 let mut still_going = task();
6991 still_going.id = "20260101-000002-dep1".to_owned();
6992 queue.put(&mut still_going).unwrap();
6993
6994 let mut blocked = task();
6995 blocked.id = "20260101-000003-main".to_owned();
6996 blocked.block(
6997 vec![dep.id.clone(), still_going.id.clone()],
6998 Some("waits on both".to_owned()),
6999 );
7000 queue.put(&mut blocked).unwrap();
7001
7002 resolve_blockers(&queue, &questions);
7003
7004 let after = queue.get(&blocked.id).unwrap();
7005 assert_eq!(
7006 after.status,
7007 TaskStatus::Blocked,
7008 "one dependency is still outstanding"
7009 );
7010 assert_eq!(after.blocked_by, [still_going.id.clone()]);
7011 }
7012
7013 #[test]
7014 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
7015 let dir = tempfile::tempdir().unwrap();
7016 let queue = Queue::at(dir.path().join("queue"));
7017 let questions = ask::Questions::at(dir.path().join("questions"));
7018
7019 let mut q = crate::ask::Question::new(
7020 "20260101-000001-main".to_owned(),
7021 crate::conduct::NODE.to_owned(),
7022 "conduct".to_owned(),
7023 "Which backend?".to_owned(),
7024 String::new(),
7025 Vec::new(),
7026 );
7027 questions.put(&mut q).unwrap();
7028 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
7029 .unwrap();
7030 questions.put(&mut q).unwrap();
7031
7032 let mut blocked = task();
7033 blocked.id = "20260101-000001-main".to_owned();
7034 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
7035 queue.put(&mut blocked).unwrap();
7036
7037 resolve_blockers(&queue, &questions);
7038
7039 let after = queue.get(&blocked.id).unwrap();
7040 assert_eq!(
7041 after.status,
7042 TaskStatus::Queued,
7043 "the only blocker resolved"
7044 );
7045 assert_eq!(after.answers.len(), 1);
7046 assert_eq!(after.answers[0].question, "Which backend?");
7047 assert_eq!(after.answers[0].answer, "SQLite");
7048
7049 let instruction = instruction_for(&after);
7051 assert!(instruction.contains("Which backend?"));
7052 assert!(instruction.contains("SQLite"));
7053 }
7054
7055 #[test]
7056 fn resolve_blockers_holds_a_task_whose_conductor_question_was_abandoned() {
7057 let dir = tempfile::tempdir().unwrap();
7058 let queue = Queue::at(dir.path().join("queue"));
7059 let questions = ask::Questions::at(dir.path().join("questions"));
7060
7061 let mut q = crate::ask::Question::new(
7062 "20260101-000001-main".to_owned(),
7063 crate::conduct::NODE.to_owned(),
7064 "conduct".to_owned(),
7065 "Is the setup done?".to_owned(),
7066 String::new(),
7067 Vec::new(),
7068 );
7069 q.abandon("no answer within 60s of asking");
7070 questions.put(&mut q).unwrap();
7071
7072 let mut blocked = task();
7073 blocked.id = "20260101-000001-main".to_owned();
7074 blocked.block(vec![q.id.clone()], Some("setup?".to_owned()));
7075 queue.put(&mut blocked).unwrap();
7076
7077 resolve_blockers(&queue, &questions);
7078
7079 let after = queue.get(&blocked.id).unwrap();
7080 assert_eq!(
7081 after.status,
7082 TaskStatus::Held,
7083 "never left blocked on nothing"
7084 );
7085 assert!(!after.operator_held(), "a machine hold, for triage");
7086 assert!(
7087 after
7088 .hold_reason
7089 .as_deref()
7090 .unwrap_or_default()
7091 .contains("went unanswered")
7092 );
7093 }
7094
7095 #[test]
7096 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
7097 let dir = tempfile::tempdir().unwrap();
7103 let queue = Queue::at(dir.path().join("queue"));
7104 let questions = ask::Questions::at(dir.path().join("questions"));
7105
7106 let mut q = crate::ask::Question::new(
7107 "20260101-000001-main".to_owned(),
7108 crate::conduct::NODE.to_owned(),
7109 "conduct".to_owned(),
7110 "How should this be handled?".to_owned(),
7111 String::new(),
7112 Vec::new(),
7113 );
7114 questions.put(&mut q).unwrap();
7115 q.answer(crate::ask::Answer::Text(
7116 "leave it held, a human will look at it later".to_owned(),
7117 ))
7118 .unwrap();
7119 questions.put(&mut q).unwrap();
7120
7121 let mut held = task();
7122 held.id = "20260101-000001-main".to_owned();
7123 held.hold_machine(Some("out of attempts".to_owned()));
7124 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
7125 queue.put(&mut held).unwrap();
7126
7127 resolve_blockers(&queue, &questions);
7128
7129 let after = queue.get(&held.id).unwrap();
7130 assert_eq!(after.status, TaskStatus::Held);
7131 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
7132 assert_eq!(
7133 after.answers[0].answer,
7134 "leave it held, a human will look at it later"
7135 );
7136 }
7137
7138 #[test]
7139 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
7140 let dir = tempfile::tempdir().unwrap();
7146 let queue = Queue::at(dir.path().join("queue"));
7147 let questions = ask::Questions::at(dir.path().join("questions"));
7148
7149 let mut still_going = task();
7150 still_going.id = "20260101-000002-dep1".to_owned();
7151 queue.put(&mut still_going).unwrap();
7152
7153 let mut blocked = task();
7154 blocked.id = "20260101-000003-main".to_owned();
7155 blocked.block(
7156 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7157 Some("waits on both".to_owned()),
7158 );
7159 queue.put(&mut blocked).unwrap();
7160
7161 resolve_blockers(&queue, &questions);
7162
7163 let after = queue.get(&blocked.id).unwrap();
7164 assert_eq!(
7165 after.status,
7166 TaskStatus::Held,
7167 "a missing dependency must not leave the task blocked forever"
7168 );
7169 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7170 assert!(after.blocked_by.is_empty());
7171 let reason = after.hold_reason.as_deref().unwrap_or_default();
7172 assert!(
7173 reason.contains("20260101-000001-gone"),
7174 "the missing id must be named so an operator can tell what happened: {reason}"
7175 );
7176 assert!(
7177 reason.contains(&still_going.id),
7178 "the still-valid dependency must not silently vanish from the record: {reason}"
7179 );
7180 }
7181
7182 #[test]
7183 fn instruction_for_is_unchanged_without_any_answers() {
7184 let t = task();
7185 assert_eq!(instruction_for(&t), t.instruction);
7186 }
7187
7188 #[test]
7189 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7190 let dir = tempfile::tempdir().unwrap();
7191 let q = Queue::at(dir.path().join("queue"));
7192 let src = dir.path().join("shot.png");
7193 std::fs::write(&src, "x").unwrap();
7194 let mut t = task();
7195 q.attach(&mut t, &[src]).unwrap();
7196 let paths = task_attachments(&q, &t).unwrap();
7197 assert_eq!(paths.len(), 1);
7198 assert!(paths[0].is_absolute() && paths[0].is_file());
7199 std::fs::remove_file(&paths[0]).unwrap();
7200 let err = task_attachments(&q, &t).unwrap_err().to_string();
7201 assert!(err.contains("shot.png"), "{err}");
7202 }
7203
7204 #[test]
7205 fn resumed_instruction_is_unchanged_without_any_answers() {
7206 let t = task();
7207 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7208 }
7209
7210 #[test]
7211 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7212 let mut t = task();
7213 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7214 let old = t.instruction.clone();
7218
7219 let refreshed = resumed_instruction(&old, &t);
7220 assert!(refreshed.starts_with(&old), "the original text is kept");
7221 assert!(refreshed.contains("Which backend?"));
7222 assert!(refreshed.contains("SQLite"));
7223 }
7224
7225 #[test]
7226 fn resumed_instruction_keeps_an_original_answers_heading() {
7227 let mut t = task();
7228 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7229 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7230
7231 let refreshed = resumed_instruction(&t.instruction, &t);
7232
7233 assert!(
7234 refreshed.starts_with(&t.instruction),
7235 "an answers heading in the original instruction is not the appended block"
7236 );
7237 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7238 assert!(refreshed.contains("Which backend?"));
7239 assert!(refreshed.contains("SQLite"));
7240
7241 let repeated = resumed_instruction(&refreshed, &t);
7242 assert_eq!(
7243 repeated, refreshed,
7244 "only the final appended block is refreshed"
7245 );
7246 }
7247
7248 #[test]
7249 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7250 let mut t = task();
7251 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7252
7253 let once = resumed_instruction(&t.instruction, &t);
7257 let twice = resumed_instruction(&once, &t);
7258 assert_eq!(once, twice);
7259 assert_eq!(once.matches("Which backend?").count(), 1);
7260
7261 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7263 let refreshed = resumed_instruction(&once, &t);
7264 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7265 assert!(refreshed.contains("Which backend?"));
7266 assert!(refreshed.contains("Which cache?"));
7267 }
7268
7269 #[test]
7270 fn prepare_instruction_covers_all_three_starters() {
7271 let mut t = task();
7272 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7273
7274 assert_eq!(
7277 prepare_instruction(&Starter::Start, None, &t),
7278 Some(instruction_for(&t))
7279 );
7280
7281 let old = t.instruction.clone();
7284 assert_eq!(
7285 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7286 Some(resumed_instruction(&old, &t))
7287 );
7288
7289 assert_eq!(
7293 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7294 None
7295 );
7296 }
7297
7298 #[test]
7299 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7300 assert_eq!(
7301 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7302 Starter::Review("magi/eba2/A".to_owned())
7303 );
7304 }
7305
7306 #[test]
7307 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7308 assert_eq!(
7309 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7310 Starter::Start,
7311 "a vanished review branch must not fall back to resuming the old run either"
7312 );
7313 }
7314
7315 #[test]
7316 fn a_refused_handover_retries_as_a_review_of_the_same_branch() {
7317 let mut t = task();
7318 t.start("old-run".to_owned());
7319 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
7320 t.release();
7321 let branch = t.review_branch.take();
7322 assert_eq!(
7323 choose_starter(branch.as_deref(), true, Some("old-run")),
7324 Starter::Review("magi/eba2/A".to_owned()),
7325 "a review wins over resuming the old run"
7326 );
7327 }
7328
7329 #[test]
7330 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7331 assert_eq!(
7332 choose_starter(None, false, Some("some-run")),
7333 Starter::Resume("some-run".to_owned())
7334 );
7335 assert_eq!(choose_starter(None, false, None), Starter::Start);
7336 }
7337
7338 #[test]
7339 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7340 let mut released = task();
7341 released.start("stalled-run".to_owned());
7342 released.requeue();
7343 let unfinished = (!released.fresh_start)
7344 .then(|| Some("stalled-run".to_owned()))
7345 .flatten();
7346 assert_eq!(
7347 choose_starter(None, false, unfinished.as_deref()),
7348 Starter::Start,
7349 "release keeps run history but must not resume it"
7350 );
7351 assert_eq!(released.runs, ["stalled-run"]);
7352 }
7353
7354 #[test]
7355 fn an_ordinary_release_keeps_a_resumable_run_available() {
7356 let mut released = task();
7357 released.start("stalled-run".to_owned());
7358 released.release();
7359 let unfinished = (!released.fresh_start)
7360 .then(|| Some("stalled-run".to_owned()))
7361 .flatten();
7362 assert_eq!(
7363 choose_starter(None, false, unfinished.as_deref()),
7364 Starter::Resume("stalled-run".to_owned()),
7365 "manual release must preserve the normal resume path"
7366 );
7367 }
7368
7369 #[test]
7370 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7371 let mut state = run_state(RunStatus::Blocked);
7372 state.config.graph.review_rounds = 3;
7373 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7374 assert!(exhausted_review_budget(&state));
7375
7376 state.reviews.pop();
7378 assert!(!exhausted_review_budget(&state));
7379
7380 let mut stalled = run_state(RunStatus::Stalled);
7383 stalled.config.graph.review_rounds = 1;
7384 stalled.reviews = vec![review_round(1)];
7385 assert!(!exhausted_review_budget(&stalled));
7386 }
7387
7388 fn review_round(round: usize) -> crate::run::ReviewRound {
7389 crate::run::ReviewRound {
7390 round,
7391 head: "deadbeef".to_owned(),
7392 verified_head: None,
7393 verified_at: None,
7394 reviews: Vec::new(),
7395 e2e: Vec::new(),
7396 verify_retried: false,
7397 e2e_deferred: false,
7398 e2e_defer_reason: None,
7399 fix: None,
7400 blocking: 0,
7401 answered: 1,
7402 expected: 1,
7403 clean: false,
7404 progressed: true,
7405 vote_split: false,
7406 reconsideration: Vec::new(),
7407 verdict: None,
7408 }
7409 }
7410
7411 #[test]
7412 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7413 let mut run = RunState::new(
7414 PathBuf::from("/repo"),
7415 "main".to_owned(),
7416 "abc1234def".to_owned(),
7417 "add retries".to_owned(),
7418 Config::default(),
7419 );
7420 run.status = RunStatus::Judging;
7421 run.parked = true;
7422 let mut task = Task::new(
7423 "add retries".to_owned(),
7424 "add retries".to_owned(),
7425 PathBuf::from("/repo"),
7426 crate::queue::Source::Human,
7427 );
7428 task.status = TaskStatus::Failed;
7429 task.runs = vec![run.id.clone()];
7430 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7431 assert!(with(&task, &run), "parked after judging is the case");
7432
7433 let mut not_parked = run.clone();
7434 not_parked.parked = false;
7435 not_parked.status = RunStatus::Stalled;
7436 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7437
7438 let mut fresh = task.clone();
7439 fresh.fresh_start = true;
7440 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7441
7442 let mut review = task.clone();
7443 review.review_branch = Some("magi/x/A".to_owned());
7444 assert!(!with(&review, &run), "review is ranked before resume");
7445
7446 let mut held = task.clone();
7447 held.status = TaskStatus::Held;
7448 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7449
7450 let mut released = run.clone();
7451 released.released_to = Some("20260901-000000-new1".to_owned());
7452 assert!(!with(&task, &released), "nothing left to resume into");
7453
7454 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7455 "unreadable"
7456 )));
7457 }
7458
7459 #[test]
7460 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7461 let mut released = RunState::new(
7462 PathBuf::from("/repo"),
7463 "main".to_owned(),
7464 "abc1234def".to_owned(),
7465 "add retries".to_owned(),
7466 Config::default(),
7467 );
7468 released.status = RunStatus::Blocked;
7469 assert_eq!(
7470 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7471 Some(released.id.clone())
7472 );
7473 released.released_to = Some("20260901-000000-new1".to_owned());
7474 assert_eq!(
7475 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7476 None,
7477 "there is nothing left to resume it into"
7478 );
7479 }
7480
7481 #[test]
7482 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
7483 let mut exhausted = RunState::new(
7492 PathBuf::from("/repo"),
7493 "main".to_owned(),
7494 "abc1234def".to_owned(),
7495 "add retries".to_owned(),
7496 Config::default(),
7497 );
7498 exhausted.status = RunStatus::Blocked;
7499 exhausted.config.graph.review_rounds = 1;
7500 exhausted.reviews = vec![review_round(1)];
7501
7502 assert_eq!(
7503 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
7504 None,
7505 "an exhausted `Blocked` run must not be offered as resumable"
7506 );
7507
7508 let mut has_budget_left = RunState::new(
7511 PathBuf::from("/repo"),
7512 "main".to_owned(),
7513 "abc1234def".to_owned(),
7514 "add retries".to_owned(),
7515 Config::default(),
7516 );
7517 has_budget_left.status = RunStatus::Blocked;
7518 has_budget_left.config.graph.review_rounds = 3;
7519 has_budget_left.reviews = vec![review_round(1)];
7520
7521 assert_eq!(
7522 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
7523 Ok(has_budget_left.clone())
7524 }),
7525 Some(has_budget_left.id.clone())
7526 );
7527 }
7528
7529 #[test]
7530 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
7531 let mut older_stalled = RunState::new(
7539 PathBuf::from("/repo"),
7540 "main".to_owned(),
7541 "abc1234def".to_owned(),
7542 "add retries".to_owned(),
7543 Config::default(),
7544 );
7545 older_stalled.status = RunStatus::Stalled;
7546
7547 let mut newest_exhausted = RunState::new(
7548 PathBuf::from("/repo"),
7549 "main".to_owned(),
7550 "abc1234def".to_owned(),
7551 "add retries".to_owned(),
7552 Config::default(),
7553 );
7554 newest_exhausted.status = RunStatus::Blocked;
7555 newest_exhausted.config.graph.review_rounds = 1;
7556 newest_exhausted.reviews = vec![review_round(1)];
7557
7558 assert_eq!(
7559 unfinished_run_with(
7560 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
7561 "t",
7562 |_| Ok(newest_exhausted.clone())
7563 ),
7564 None,
7565 "the newest run is exhausted, so nothing here is worth resuming - \
7566 least of all the older, already-superseded run"
7567 );
7568 }
7569
7570 #[test]
7571 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
7572 assert_eq!(
7573 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
7574 Err(anyhow::anyhow!("fixture is absent"))
7575 }),
7576 None
7577 );
7578 }
7579
7580 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
7581 let mut q = ask::Question::new(
7582 run.to_owned(),
7583 "implement".to_owned(),
7584 "impl-A".to_owned(),
7585 "continue?".to_owned(),
7586 String::new(),
7587 vec!["resume で続行する".to_owned(), "other".to_owned()],
7588 );
7589 q.actions.insert("resume で続行する".to_owned(), action);
7590 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
7591 .unwrap();
7592 q
7593 }
7594
7595 fn held_task_with(run: &str) -> Task {
7596 let mut t = task();
7597 t.runs = vec![run.to_owned()];
7598 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
7599 t
7600 }
7601
7602 fn resume_action(run: &str) -> ask::ChoiceAction {
7603 ask::ChoiceAction::Resume { run: run.into() }
7604 }
7605
7606 #[test]
7607 fn decide_action_resumes_only_the_latest_resumable_run() {
7608 let t = held_task_with("r1");
7609 let q = action_question("r1", resume_action("r1"));
7610 let load = |s: RunState| move |_: &str| Ok(s);
7611 assert_eq!(
7612 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
7613 ActionDecision::Resume("r1".into())
7614 );
7615 let q_other = action_question("r1", resume_action("r0"));
7617 assert!(matches!(
7618 decide_action(
7619 &t,
7620 &q_other,
7621 &PHRASES_EN,
7622 load(run_state(RunStatus::Blocked))
7623 ),
7624 ActionDecision::Refuse(_)
7625 ));
7626 assert!(matches!(
7628 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
7629 ActionDecision::Refuse(_)
7630 ));
7631 let mut released = run_state(RunStatus::Blocked);
7633 released.released_to = Some("elsewhere".into());
7634 assert!(matches!(
7635 decide_action(&t, &q, &PHRASES_EN, load(released)),
7636 ActionDecision::Refuse(_)
7637 ));
7638 assert!(matches!(
7640 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
7641 ActionDecision::Refuse(_)
7642 ));
7643 }
7644
7645 #[test]
7646 fn decide_action_ignores_a_question_about_an_earlier_run() {
7647 let mut t = held_task_with("r1");
7648 t.runs.push("r2".to_owned());
7649 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7650 assert_eq!(
7651 decide_action(
7652 &t,
7653 &action_question("r1", ask::ChoiceAction::Done),
7654 &PHRASES_EN,
7655 never
7656 ),
7657 ActionDecision::Stale
7658 );
7659 }
7660
7661 #[test]
7662 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
7663 let mut t = held_task_with("r1");
7664 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7665 assert_eq!(
7666 decide_action(
7667 &t,
7668 &action_question("r1", ask::ChoiceAction::Requeue),
7669 &PHRASES_EN,
7670 never
7671 ),
7672 ActionDecision::Requeue
7673 );
7674 let done_q = action_question("r1", ask::ChoiceAction::Done);
7675 assert_eq!(
7676 decide_action(&t, &done_q, &PHRASES_EN, never),
7677 ActionDecision::Done
7678 );
7679 t.mark_action_applied(&done_q.id);
7680 assert_eq!(
7681 decide_action(&t, &done_q, &PHRASES_EN, never),
7682 ActionDecision::Skip
7683 );
7684
7685 let mut plain = action_question("r1", ask::ChoiceAction::Done);
7687 plain.actions.clear();
7688 assert_eq!(
7689 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
7690 ActionDecision::Skip
7691 );
7692 let mut running = held_task_with("r1");
7694 running.status = TaskStatus::Running;
7695 assert_eq!(
7696 decide_action(
7697 &running,
7698 &action_question("r1", ask::ChoiceAction::Done),
7699 &PHRASES_EN,
7700 never
7701 ),
7702 ActionDecision::Skip
7703 );
7704 }
7705
7706 #[test]
7707 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
7708 let dir = tempfile::tempdir().unwrap();
7709 let queue = Queue::at(dir.path().join("queue"));
7710 let questions = Questions::at(dir.path().join("questions"));
7711 let home = dir.path().join("home");
7712 let mut state = run_state(RunStatus::Blocked);
7713 state.id = "20260101-000000-act1".to_owned();
7714 state.save_under(&home).unwrap();
7715
7716 let mut t = held_task_with(&state.id);
7717 queue.put(&mut t).unwrap();
7718 let mut q = action_question(&state.id, resume_action(&state.id));
7719 questions.put(&mut q).unwrap();
7720
7721 apply_choice_actions(&queue, &questions, &home);
7722 let after = queue.get(&t.id).unwrap();
7723 assert_eq!(after.status, TaskStatus::Queued);
7724 assert!(!after.fresh_start);
7725 assert!(after.action_applied(&q.id));
7726 let pin = after.resume_override.clone().unwrap();
7727 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
7728 assert!(pin.forced);
7729
7730 let mut again = queue.get(&t.id).unwrap();
7732 again.hold_machine(Some("later".into()));
7733 queue.put(&mut again).unwrap();
7734 apply_choice_actions(&queue, &questions, &home);
7735 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7736 }
7737
7738 #[test]
7739 fn an_answer_the_waiter_already_delivered_is_not_acted_on_again() {
7740 let dir = tempfile::tempdir().unwrap();
7741 let queue = Queue::at(dir.path().join("queue"));
7742 let questions = Questions::at(dir.path().join("questions"));
7743 let home = dir.path().join("home");
7744 let mut state = run_state(RunStatus::Blocked);
7745 state.id = "20260101-000000-act2".to_owned();
7746 state.save_under(&home).unwrap();
7747
7748 let mut t = held_task_with(&state.id);
7749 queue.put(&mut t).unwrap();
7750 let mut q = action_question(&state.id, resume_action(&state.id));
7751 q.answer_delivered = true;
7752 questions.put(&mut q).unwrap();
7753
7754 apply_choice_actions(&queue, &questions, &home);
7755 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7756
7757 let mut q2 = action_question(&state.id, resume_action(&state.id));
7759 questions.put(&mut q2).unwrap();
7760 apply_choice_actions(&queue, &questions, &home);
7761 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7762 assert!(questions.get(&q2.id).unwrap().answer_delivered);
7763 }
7764}