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 let mut changed = false;
702 for id in task.blocked_by.clone() {
703 if let Ok(dep) = queue.get(&id) {
704 if dep.status == TaskStatus::Done {
705 task.unblock(&id);
706 changed = true;
707 }
708 continue;
709 }
710 if let Ok(q) = questions.get(&id)
711 && q.status == ask::QuestionStatus::Answered
712 {
713 let answer = match &q.answer {
714 Some(ask::Answer::Choice(c) | ask::Answer::Text(c)) => c.clone(),
715 None => String::new(),
716 };
717 task.record_answer(q.summary.clone(), answer);
718 task.unblock(&id);
719 changed = true;
720 }
721 }
722 if changed {
723 record(queue, &mut task);
724 }
725 }
726}
727
728#[derive(Debug, Clone, PartialEq, Eq)]
730enum ActionDecision {
731 Skip,
733 Resume(String),
735 Requeue,
737 Done,
739 Stale,
742 Refuse(String),
745}
746
747fn decide_action<F>(task: &Task, q: &ask::Question, p: &Phrases, load: F) -> ActionDecision
755where
756 F: FnOnce(&str) -> Result<RunState>,
757{
758 let Some(action) = q.chosen_action() else {
759 return ActionDecision::Skip;
760 };
761 if task.action_applied(&q.id)
762 || matches!(
763 task.status,
764 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
765 )
766 {
767 return ActionDecision::Skip;
768 }
769 if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
772 return ActionDecision::Stale;
773 }
774 match action {
775 ask::ChoiceAction::Requeue => ActionDecision::Requeue,
776 ask::ChoiceAction::Done => ActionDecision::Done,
777 ask::ChoiceAction::Resume { run } => {
778 if task.runs.last() != Some(run) {
779 return ActionDecision::Refuse((p.resume_not_latest)(
780 q.short(),
781 ask::short_id(run),
782 ));
783 }
784 match load(run) {
785 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
786 ActionDecision::Resume(run.clone())
787 }
788 Ok(_) => ActionDecision::Refuse((p.resume_cannot_progress)(
789 q.short(),
790 ask::short_id(run),
791 )),
792 Err(e) => ActionDecision::Refuse((p.resume_unreadable)(
793 q.short(),
794 ask::short_id(run),
795 &format!("{e:#}"),
796 )),
797 }
798 }
799 }
800}
801
802struct Phrases {
816 graph_stopped: fn(&str, &str) -> String,
818 quorum_lost: &'static str,
819 quota_took_out: &'static str,
821 run_ended: &'static str,
823 waiting_for_answer: &'static str,
825 recovered_running: &'static str,
826 no_run_to_recover: &'static str,
827 could_not_start: &'static str,
828 resume_not_latest: fn(&str, &str) -> String,
829 resume_cannot_progress: fn(&str, &str) -> String,
830 resume_unreadable: fn(&str, &str, &str) -> String,
831}
832
833const PHRASES_EN: Phrases = Phrases {
834 graph_stopped: |status, detail| {
835 format!("the graph stopped at `{status}` without reaching a terminal status: {detail}")
836 },
837 quorum_lost: "the judging panel lost its quorum",
838 quota_took_out: "; quota took out ",
839 run_ended: "run ended ",
840 waiting_for_answer: " - waiting for operator answer to question ",
841 recovered_running: "recovered a `running` task whose daemon never recorded the outcome: ",
842 no_run_to_recover: "task was `running` with no live daemon and no readable \
843 run to recover; held for a human to check what happened",
844 could_not_start: "could not start the run: ",
845 resume_not_latest: |q, run| {
846 format!("question {q} asked to resume run {run}, which is not this task's latest run")
847 },
848 resume_cannot_progress: |q, run| {
849 format!("question {q} asked to resume run {run}, which cannot make progress")
850 },
851 resume_unreadable: |q, run, e| {
852 format!("question {q} asked to resume run {run}, which could not be read: {e}")
853 },
854};
855
856const PHRASES_JA: Phrases = Phrases {
857 graph_stopped: |status, detail| {
858 format!("グラフが終端状態に達しないまま `{status}` で停止しました: {detail}")
859 },
860 quorum_lost: "審査パネルが定足数を失いました",
861 quota_took_out: "。クォータで脱落: ",
862 run_ended: "run 終了: ",
863 waiting_for_answer: " - オペレーターの回答待ち: 質問 ",
864 recovered_running: "daemon が結果を記録しないまま `running` だったタスクを回収しました: ",
865 no_run_to_recover: "タスクは `running` でしたが、生きた daemon も回収できる run も見つかりません。\
866 何が起きたか人が確認するため保留にしました",
867 could_not_start: "run を開始できませんでした: ",
868 resume_not_latest: |q, run| {
869 format!(
870 "質問 {q} は run {run} の再開を求めましたが、これはタスクの最新の run ではありません"
871 )
872 },
873 resume_cannot_progress: |q, run| {
874 format!("質問 {q} は run {run} の再開を求めましたが、これは進行できません")
875 },
876 resume_unreadable: |q, run, e| {
877 format!("質問 {q} は run {run} の再開を求めましたが、読み込めませんでした: {e}")
878 },
879};
880
881fn phrases(language: &str) -> &'static Phrases {
882 if crate::lang::is_japanese(language) {
883 &PHRASES_JA
884 } else {
885 &PHRASES_EN
886 }
887}
888
889fn language_of(task: &Task, fallback: &Path) -> String {
893 crate::lang::of_repo(&repo_for(task, fallback))
894}
895
896fn task_of_question<'a>(tasks: &'a [Task], q: &ask::Question) -> Option<&'a Task> {
900 if q.node == crate::conduct::NODE {
901 return tasks.iter().find(|t| t.id == q.run);
902 }
903 tasks.iter().find(|t| t.runs.contains(&q.run))
904}
905
906fn apply_choice_actions(queue: &Queue, questions: &Questions, home: &Path) {
915 let tasks = queue.list();
916 for q in questions.list() {
917 if q.chosen_action().is_none() {
918 continue;
919 }
920 let Some(listed) = task_of_question(&tasks, &q) else {
921 continue;
922 };
923 if listed.action_applied(&q.id) {
924 continue;
925 }
926 let Ok(_claim) = queue.claim(&listed.id) else {
927 continue;
928 };
929 let Ok(mut task) = queue.get(&listed.id) else {
930 continue;
931 };
932 let language = language_of(&task, Path::new("."));
933 let decision = decide_action(&task, &q, phrases(&language), |id| {
934 RunState::load_under(id, home)
935 });
936 let ran = matches!(
937 decision,
938 ActionDecision::Resume(_) | ActionDecision::Requeue | ActionDecision::Done
939 );
940 if ran {
941 if questions
946 .read_lease(&q.id)
947 .is_some_and(|l| l.fresh(Timestamp::now()))
948 {
949 continue;
950 }
951 let taken = questions.update(&q.id, |r| {
952 let free = !r.answer_delivered;
953 r.answer_delivered = true;
954 Ok(free)
955 });
956 if !matches!(taken, Ok((_, true))) {
957 continue;
958 }
959 }
960 match decision {
961 ActionDecision::Skip => continue,
962 ActionDecision::Resume(run) => {
963 task.release();
964 task.resume_override = Some(crate::queue::OperatorResume {
965 question_id: q.id.clone(),
966 at: Timestamp::now(),
967 conductor_rehold: None,
968 forced: true,
969 pinned_run: Some(run),
970 });
971 }
972 ActionDecision::Stale => {}
973 ActionDecision::Requeue => task.requeue(),
974 ActionDecision::Done => {
975 task.succeed();
976 supersede_prior_runs(&task, home);
977 }
978 ActionDecision::Refuse(why) => {
979 task.hold_machine(Some(why));
980 notices::raise(
981 Notice::warn(
982 &format!("action:{}", q.id),
983 "An answer asked the daemon to resume a run that cannot be resumed; the task stays held.",
984 )
985 .link(Link::Task {
986 id: task.id.clone(),
987 }),
988 );
989 }
990 }
991 task.mark_action_applied(&q.id);
992 record(queue, &mut task);
993 }
994}
995
996fn reconcile_task_questions(queue: &Queue, questions: &Questions) {
1007 let tasks = queue.list();
1008 let by_id: std::collections::BTreeMap<&str, &Task> =
1009 tasks.iter().map(|t| (t.id.as_str(), t)).collect();
1010 let referenced: std::collections::BTreeSet<&str> = tasks
1011 .iter()
1012 .flat_map(|task| task.blocked_by.iter().map(String::as_str))
1013 .collect();
1014
1015 for mut question in questions.list() {
1016 if !question.status.open() || question.node != crate::conduct::NODE {
1017 continue;
1018 }
1019 if referenced.contains(question.id.as_str()) {
1022 continue;
1023 }
1024 let Some(task) = by_id.get(question.run.as_str()) else {
1025 continue;
1026 };
1027 question.abandon(format!(
1028 "task {} no longer waits for this answer",
1029 task.short()
1030 ));
1031 if let Err(e) = questions.put(&mut question) {
1032 tracing::warn!(
1033 "could not retire question {} for task {}: {e:#}",
1034 question.short(),
1035 task.short()
1036 );
1037 }
1038 }
1039}
1040
1041#[derive(Debug, Clone, Copy)]
1048pub struct Verdict {
1049 pub status: RunStatus,
1051 pub left_pr: bool,
1053 pub quota_hit: bool,
1055 pub parked: bool,
1057 pub no_viable_candidates: bool,
1064}
1065
1066pub fn settle(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize) {
1123 settle_in(task, verdict, detail, max_attempts, &PHRASES_EN)
1124}
1125
1126fn settle_in(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize, p: &Phrases) {
1128 if verdict.parked {
1134 task.stall(detail);
1135 return;
1136 }
1137 match verdict.status {
1138 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1139 RunStatus::AlreadyInBase => task.already_landed(detail),
1140 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1141 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1142 task.stall(detail)
1143 }
1144 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1145 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1146 RunStatus::Blocked => task.fail(detail, max_attempts),
1147 RunStatus::VerifiedNoop => task.handed_off(detail),
1148 other => task.fail((p.graph_stopped)(label(other), detail), max_attempts),
1149 }
1150}
1151
1152pub fn supersede_prior_runs(task: &Task, home: &Path) {
1201 let now = Timestamp::now();
1202 let last_run_succeeded = task
1209 .runs
1210 .last()
1211 .and_then(|id| RunState::load_under(id, home).ok())
1212 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1213 for id in task.superseded_attempts(last_run_succeeded) {
1214 let mut state = match RunState::load_under(id, home) {
1215 Ok(s) => s,
1216 Err(e) => {
1217 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1218 continue;
1219 }
1220 };
1221 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1222 continue;
1223 }
1224 let daemon_claims = is_working_on(home, id, now);
1225 if state.liveness(daemon_claims) == Liveness::Live {
1226 continue;
1227 }
1228 state.status = RunStatus::Superseded;
1229 if let Err(e) = state.save_under(home) {
1230 tracing::warn!("could not mark run {id} superseded: {e:#}");
1231 }
1232 }
1233}
1234
1235fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1257 for task in queue.list() {
1258 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1259 continue;
1260 }
1261 supersede_prior_runs(&task, home);
1262 }
1263}
1264
1265fn settle_and_diagnose(
1272 task: &mut Task,
1273 verdict: Verdict,
1274 detail: &str,
1275 max_attempts: usize,
1276 state: &RunState,
1277) {
1278 let p = phrases(&state.config.graph.language);
1279 settle_in(task, verdict, detail, max_attempts, p);
1280 if task.status == TaskStatus::Held {
1281 task.diagnostic = diagnostic(state);
1282 note_open_question(task, &state.id, p);
1283 }
1284}
1285
1286fn note_open_question(task: &mut Task, run: &str, p: &Phrases) {
1301 let Some(home) = crate::run::try_home() else {
1302 return;
1303 };
1304 let open = Questions::at(home.join("questions")).open_for(run);
1305 let Some(q) = open.first() else {
1306 return;
1307 };
1308 let base = task.hold_reason.clone().unwrap_or_default();
1309 task.hold_reason = Some(format!("{base}{}{}", p.waiting_for_answer, q.short()));
1310}
1311
1312fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize, language: &str) {
1323 match last_run {
1324 Some(state) => {
1325 let verdict = Verdict {
1326 status: state.status,
1327 left_pr: state.pr.is_some(),
1328 quota_hit: !state.quota.is_empty(),
1329 parked: state.parked,
1330 no_viable_candidates: state.viable().is_empty(),
1331 };
1332 let detail = format!(
1333 "{}{}",
1334 phrases(&state.config.graph.language).recovered_running,
1335 describe(&state)
1336 );
1337 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1338 }
1339 None => {
1340 let why = phrases(language).no_run_to_recover;
1343 task.last_error = Some(why.to_owned());
1344 task.hold_machine(Some(why.to_owned()));
1347 }
1348 }
1349}
1350
1351fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1373 let mut reclaimed = Vec::new();
1374 for listed in queue.list() {
1375 if listed.status != TaskStatus::Running {
1376 continue;
1377 }
1378 let Ok(_claim) = queue.claim(&listed.id) else {
1379 continue;
1380 };
1381 let Ok(mut task) = queue.get(&listed.id) else {
1385 continue;
1386 };
1387 if task.status != TaskStatus::Running {
1388 continue;
1389 }
1390 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1391 if let Some(state) = &last_run
1401 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1402 {
1403 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1404 }
1405 let language = if last_run.is_none() {
1406 language_of(&task, Path::new("."))
1407 } else {
1408 String::new()
1409 };
1410 reclaim(&mut task, last_run, max_attempts, &language);
1411 if task.status == TaskStatus::Done {
1412 supersede_prior_runs(&task, &crate::run::home());
1413 }
1414 record(queue, &mut task);
1415 reclaimed.push(task.id.clone());
1416 }
1417 reclaimed
1418}
1419
1420fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1452 reclaim_abandoned_runs_with(
1453 home,
1454 now,
1455 crate::proc::pid_status,
1456 crate::proc::process_started_at,
1457 )
1458}
1459
1460fn reclaim_abandoned_runs_with<F, G>(
1466 home: &Path,
1467 now: Timestamp,
1468 query: F,
1469 identity: G,
1470) -> Vec<String>
1471where
1472 F: Fn(u32) -> Option<bool>,
1473 G: Fn(u32) -> Option<String>,
1474{
1475 let mut abandoned = Vec::new();
1476 for entry in std::fs::read_dir(home.join("runs"))
1477 .into_iter()
1478 .flatten()
1479 .flatten()
1480 {
1481 let id = entry.file_name().to_string_lossy().into_owned();
1482 if !crate::run::is_run_id(&id) {
1483 continue;
1484 }
1485 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1493 continue;
1494 };
1495 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1496 continue;
1497 };
1498 if state.status.done() || !state.active_all_overrun(now) {
1499 continue;
1500 }
1501 let daemon_claims = is_working_on(home, &id, now);
1513 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1514 continue;
1515 }
1516 state.abandon("daemon");
1517 if let Err(e) = state.save_under(home) {
1518 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1519 continue;
1520 }
1521 if let Some(notice) = notices::run_ended(&state) {
1524 notices::raise_in(home, notice);
1525 }
1526 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1534 tracing::warn!("abandon questions for {id}: {e:#}");
1535 }
1536 abandoned.push(id);
1537 }
1538 abandoned
1539}
1540
1541pub async fn serve(opts: Opts) -> Result<()> {
1547 serve_until(opts, Stop::new()).await
1548}
1549
1550pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1567 let signal = {
1568 let stop = stop.clone();
1569 tokio::spawn(async move {
1570 if tokio::signal::ctrl_c().await.is_ok() {
1571 stop.stop();
1572 tracing::info!("shutdown requested; a run in flight will be finished first");
1573 }
1574 })
1575 };
1576
1577 let worktrees_root = opts
1578 .worktrees_root
1579 .clone()
1580 .unwrap_or_else(crate::run::default_worktree_root);
1581 let outcome = drive(
1582 &opts,
1583 &Queue::open(),
1584 &status_path(),
1585 &crate::run::home(),
1586 &worktrees_root,
1587 &stop,
1588 )
1589 .await;
1590
1591 signal.abort();
1592 outcome
1593}
1594
1595async fn drive(
1608 opts: &Opts,
1609 queue: &Queue,
1610 status_file: &Path,
1611 home: &Path,
1612 worktrees_root: &Path,
1613 stop: &Stop,
1614) -> Result<()> {
1615 let status = Arc::new(Mutex::new(Status::new()));
1623 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1624 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1625
1626 let daemon_cfg = prepare(&opts.repo, opts)
1632 .map(|c| c.daemon)
1633 .unwrap_or_default();
1634 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1635
1636 let waiter = tokio::spawn(crate::waiter::run(
1641 crate::waiter::Waiter::new(
1642 crate::ask::Questions::at(home.join("questions")),
1643 home.to_path_buf(),
1644 prepare(&opts.repo, opts).ok(),
1645 ),
1646 stop.clone(),
1647 ));
1648
1649 tracing::info!(
1650 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1651 queue.root().display(),
1652 opts.poll.as_secs(),
1653 opts.max_attempts,
1654 concurrency,
1655 if daemon_cfg.pause_for_interrupts {
1656 ", interrupts enabled"
1657 } else {
1658 ""
1659 }
1660 );
1661
1662 janitor(&opts.repo, opts, home, worktrees_root).await;
1665 resweep_superseded_attempts(queue, home);
1666
1667 let outcome = poll(
1668 opts,
1669 queue,
1670 &status,
1671 home,
1672 worktrees_root,
1673 stop,
1674 DispatchLimits {
1675 max_concurrent: concurrency,
1676 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1677 },
1678 )
1679 .await;
1680
1681 beat.abort();
1682 waiter.abort();
1683 clear_status_at(status_file);
1684 outcome
1685}
1686
1687async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1693 loop {
1694 tokio::time::sleep(HEARTBEAT).await;
1695 let snapshot = {
1696 let mut guard = lock(&status);
1697 guard.updated_at = Timestamp::now();
1698 guard.clone()
1699 };
1700 if let Err(e) = write_status_to(&path, &snapshot) {
1701 tracing::warn!("could not refresh the daemon status file: {e:#}");
1704 }
1705 }
1706}
1707
1708#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1711enum LandResume {
1712 NotLanding,
1715 StillWaiting,
1720 Ready,
1724}
1725
1726fn land_resume_state(task: &Task) -> LandResume {
1730 let Some(run_id) = task.runs.last() else {
1731 return LandResume::NotLanding;
1732 };
1733 let Ok(state) = RunState::load(run_id) else {
1734 return LandResume::NotLanding;
1735 };
1736 if state.status != RunStatus::Landing || !state.parked {
1737 return LandResume::NotLanding;
1738 }
1739 let store = ask::Questions::open();
1740 let waiting = store
1741 .list()
1742 .into_iter()
1743 .filter(|q| &q.run == run_id && q.node == land::APPROVAL_NODE)
1744 .max_by(|a, b| a.id.cmp(&b.id));
1745 let Some(mut q) = waiting else {
1746 return LandResume::Ready;
1747 };
1748 if !q.status.open() {
1749 return LandResume::Ready;
1750 }
1751 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
1758 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
1759 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
1760 q.abandon(format!(
1761 "no answer within {}s of asking",
1762 timeout.as_secs().max(1)
1763 ));
1764 if store.put(&mut q).is_ok() {
1767 return LandResume::Ready;
1768 }
1769 }
1770 LandResume::StillWaiting
1771}
1772
1773const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
1781
1782const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
1794
1795struct InFlightGuard<'a> {
1808 status: &'a Arc<Mutex<Status>>,
1809 stop: &'a Stop,
1810 task_id: &'a str,
1811}
1812
1813impl Drop for InFlightGuard<'_> {
1814 fn drop(&mut self) {
1815 lock(self.status).current.retain(|c| c.task != self.task_id);
1816 self.stop.exit();
1817 }
1818}
1819
1820#[derive(Debug, Clone, PartialEq, Eq)]
1842enum Interrupt {
1843 Idle,
1845 Parking {
1857 parked: Vec<String>,
1858 interrupt_task: String,
1859 },
1860 Running {
1868 parked: Vec<String>,
1869 interrupt_task: String,
1870 },
1871 Resuming { parked: Vec<String> },
1878}
1879
1880fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
1901 match state {
1902 Interrupt::Idle => {
1903 if in_flight.len() != 1 {
1915 return Interrupt::Idle;
1916 }
1917 match runnable.iter().find(|t| t.interrupt) {
1918 Some(t) => Interrupt::Parking {
1919 parked: in_flight.to_vec(),
1920 interrupt_task: t.id.clone(),
1921 },
1922 None => Interrupt::Idle,
1923 }
1924 }
1925 Interrupt::Parking {
1926 parked,
1927 interrupt_task,
1928 } => {
1929 if in_flight.iter().any(|id| parked.contains(id)) {
1930 Interrupt::Parking {
1932 parked,
1933 interrupt_task,
1934 }
1935 } else if in_flight.contains(&interrupt_task) {
1936 Interrupt::Running {
1937 parked,
1938 interrupt_task,
1939 }
1940 } else if runnable.iter().any(|t| t.id == interrupt_task) {
1941 Interrupt::Parking {
1945 parked,
1946 interrupt_task,
1947 }
1948 } else {
1949 Interrupt::Resuming { parked }
1954 }
1955 }
1956 Interrupt::Running {
1957 parked,
1958 interrupt_task,
1959 } => {
1960 if in_flight.contains(&interrupt_task) {
1961 Interrupt::Running {
1962 parked,
1963 interrupt_task,
1964 }
1965 } else {
1966 Interrupt::Resuming { parked }
1972 }
1973 }
1974 Interrupt::Resuming { parked } => {
1975 if in_flight.iter().any(|id| parked.contains(id)) {
1976 Interrupt::Idle
1982 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
1983 Interrupt::Resuming { parked }
1984 } else {
1985 Interrupt::Idle
1988 }
1989 }
1990 }
1991}
1992
1993fn advance_interrupt_tick(
1999 enabled: bool,
2000 state: Interrupt,
2001 in_flight: &[String],
2002 runnable: &[Task],
2003) -> Interrupt {
2004 if !enabled {
2005 return Interrupt::Idle;
2006 }
2007 advance_interrupt(state, in_flight, runnable)
2008}
2009
2010fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
2015 match state {
2016 Interrupt::Idle => candidates,
2017 Interrupt::Parking {
2018 parked,
2019 interrupt_task,
2020 } => {
2021 if in_flight.iter().any(|id| parked.contains(id)) {
2022 Vec::new()
2023 } else {
2024 candidates
2025 .into_iter()
2026 .filter(|t| &t.id == interrupt_task)
2027 .collect()
2028 }
2029 }
2030 Interrupt::Running { .. } => Vec::new(),
2031 Interrupt::Resuming { parked } => candidates
2039 .into_iter()
2040 .find(|t| parked.contains(&t.id))
2041 .into_iter()
2042 .collect(),
2043 }
2044}
2045
2046struct DispatchLimits {
2050 max_concurrent: usize,
2053 pause_for_interrupts: bool,
2055}
2056
2057#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2066enum PermitKind {
2067 None,
2072 Urgent,
2079 Ordinary,
2082}
2083
2084fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
2087 if priority {
2088 PermitKind::None
2089 } else if urgent {
2090 PermitKind::Urgent
2091 } else {
2092 PermitKind::Ordinary
2093 }
2094}
2095
2096async fn poll(
2115 opts: &Opts,
2116 queue: &Queue,
2117 status: &Arc<Mutex<Status>>,
2118 home: &Path,
2119 worktrees_root: &Path,
2120 stop: &Stop,
2121 limits: DispatchLimits,
2122) -> Result<()> {
2123 let DispatchLimits {
2124 max_concurrent,
2125 pause_for_interrupts,
2126 } = limits;
2127 let mut attempted: Vec<String> = Vec::new();
2132 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2133 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2141 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2147 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2148 let mut conductor = Conductor::new();
2149 let mut cache_last_checked: Option<Timestamp> = None;
2152 let mut interrupt = Interrupt::Idle;
2154 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2160 std::collections::HashMap::new();
2161
2162 while !stop.stopped() {
2163 lock(status).polls += 1;
2164
2165 while let Some(result) = inflight.try_join_next() {
2170 if let Err(e) = result {
2171 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2172 notices::raise(Notice::error(
2173 "loop:attempt",
2174 "A queued attempt ended abnormally; check the task it was running.",
2175 ));
2176 }
2177 }
2178
2179 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2180 if !swept.is_empty() {
2181 tracing::warn!(
2182 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2183 swept.len(),
2184 swept.join(", ")
2185 );
2186 }
2187 let now = Timestamp::now();
2192
2193 if !stop.busy_now() {
2198 maybe_prune_cache_between_runs(
2199 &opts.repo,
2200 opts,
2201 home,
2202 stop,
2203 &mut cache_last_checked,
2204 now,
2205 )
2206 .await;
2207 }
2208
2209 let stalled = stalled_tasks(queue, home, now);
2210 let stalled_ids: std::collections::BTreeSet<_> =
2211 stalled.iter().map(|task| task.id.clone()).collect();
2212 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2213 if !reclaimed.is_empty() {
2214 tracing::warn!(
2215 "reclaimed {} task(s) left `running` by a daemon that never \
2216 recorded the outcome: {}",
2217 reclaimed.len(),
2218 reclaimed.join(", ")
2219 );
2220 }
2221 let abandoned_runs = reclaim_abandoned_runs(home, now);
2222 if !abandoned_runs.is_empty() {
2223 tracing::warn!(
2224 "failed {} run(s) left behind by a killed process, past every \
2225 active seat's own timeout: {}",
2226 abandoned_runs.len(),
2227 abandoned_runs.join(", ")
2228 );
2229 }
2230
2231 let questions = Questions::at(home.join("questions"));
2236
2237 resolve_blockers(queue, &questions);
2240 apply_choice_actions(queue, &questions, home);
2241 reconcile_task_questions(queue, &questions);
2242
2243 let finished: Vec<Task> = finished_tasks(queue)
2249 .into_iter()
2250 .filter(|task| !stalled_ids.contains(&task.id))
2251 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2255 .collect();
2256 let queued = queued_tasks(queue);
2257 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2261 && conductor.worth_a_look(queue, &stalled, &finished)
2262 {
2263 match prepare(&opts.repo, opts) {
2264 Ok(cfg) => {
2265 conductor
2266 .maybe_run(
2267 &cfg,
2268 &opts.repo,
2269 queue,
2270 &questions,
2271 home,
2272 &queued,
2273 &stalled,
2274 &finished,
2275 opts.max_attempts,
2276 )
2277 .await;
2278 }
2279 Err(e) => {
2280 tracing::warn!("conductor: no config: {e:#}");
2281 notices::raise(Notice::warn(
2282 "loop:no-config",
2283 "The loop could not read this repository's config, so held tasks are not being triaged.",
2284 ));
2285 }
2286 }
2287 }
2288
2289 let candidates: Vec<Task> = runnable(queue)
2290 .into_iter()
2291 .filter(|t| !opts.once || !attempted.contains(&t.id))
2292 .collect();
2293
2294 let in_flight: Vec<String> = lock(status)
2299 .current
2300 .iter()
2301 .map(|c| c.task.clone())
2302 .collect();
2303 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2304
2305 interrupt =
2306 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2307 if let Interrupt::Parking {
2308 parked,
2309 interrupt_task,
2310 } = &interrupt
2311 {
2312 let reason = format!(
2313 "task {} asked to run first",
2314 crate::run::short_of(interrupt_task)
2315 );
2316 for id in parked {
2317 if let Some(pause) = interrupt_pauses.get(id) {
2318 pause.park_because(reason.clone());
2319 }
2320 }
2321 }
2322 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2323
2324 let cooling_down =
2325 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2326
2327 let mut started_any = false;
2328 for candidate in candidates {
2329 if stop.stopped() {
2330 break;
2331 }
2332
2333 let resume = land_resume_state(&candidate);
2334 if resume == LandResume::StillWaiting {
2335 continue;
2336 }
2337 let priority = resume == LandResume::Ready;
2338
2339 if !priority && cooling_down {
2340 continue;
2341 }
2342 let permit = match permit_kind(priority, candidate.urgent) {
2343 PermitKind::None => None,
2344 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2345 Ok(p) => Some(p),
2346 Err(_) => continue,
2352 },
2353 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2354 Ok(p) => Some(p),
2355 Err(_) => continue,
2359 },
2360 };
2361
2362 let Ok(claim) = queue.claim(&candidate.id) else {
2367 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2368 continue;
2369 };
2370 let mut task = match queue.get(&candidate.id) {
2373 Ok(t) if t.status.runnable() => t,
2374 Ok(_) => continue,
2375 Err(e) => {
2376 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2377 continue;
2378 }
2379 };
2380 let task_id = task.id.clone();
2381 attempted.push(task_id.clone());
2382 lock(status).idle = false;
2383 stop.enter();
2386 started_any = true;
2387
2388 let run_pause = crate::graph::Pause::new();
2392 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2393
2394 let opts = opts.clone();
2395 let queue = queue.clone();
2396 let status = Arc::clone(status);
2397 let stop = stop.clone();
2398 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2399 inflight.spawn(async move {
2400 let _claim = claim;
2404 let _permit = permit;
2405 let _inflight = InFlightGuard {
2407 status: &status,
2408 stop: &stop,
2409 task_id: &task_id,
2410 };
2411 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2412 lock(&status).completed += 1;
2413 let now = Timestamp::now();
2419 if let Some(until) = cooldown_until("a, now) {
2420 let wait = until.as_second() - now.as_second();
2421 *lock("a_cooldown_until) = Some(until);
2422 let hint = quota
2423 .iter()
2424 .find(|q| q.reset.is_some())
2425 .and_then(|q| q.reset.as_deref());
2426 match hint {
2427 Some(h) => tracing::warn!(
2428 "quota hit; waiting {wait}s before taking another ordinary task \
2429 (CLI reported reset: {h})"
2430 ),
2431 None => tracing::warn!(
2432 "quota hit; waiting {wait}s before taking another ordinary task \
2433 (no reset hint reported)"
2434 ),
2435 }
2436 }
2437 });
2438 }
2439
2440 if started_any {
2441 continue;
2442 }
2443
2444 if stop.busy_now() {
2445 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2450 continue;
2451 }
2452
2453 lock(status).idle = true;
2455 if opts.once {
2456 janitor(&opts.repo, opts, home, worktrees_root).await;
2460 resweep_superseded_attempts(queue, home);
2461 triage_held(queue, home, opts).await;
2462 break;
2463 }
2464 stop.idle(opts.poll).await;
2465 if stop.stopped() {
2466 continue;
2467 }
2468 janitor(&opts.repo, opts, home, worktrees_root).await;
2474 resweep_superseded_attempts(queue, home);
2475 triage_held(queue, home, opts).await;
2476 }
2477
2478 while let Some(result) = inflight.join_next().await {
2483 if let Err(e) = result {
2484 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2485 notices::raise(Notice::error(
2486 "loop:attempt",
2487 "A queued attempt ended abnormally; check the task it was running.",
2488 ));
2489 }
2490 }
2491 Ok(())
2492}
2493
2494async fn attempt(
2500 opts: &Opts,
2501 queue: &Queue,
2502 status: &Arc<Mutex<Status>>,
2503 stop: &Stop,
2504 interrupt_pause: crate::graph::Pause,
2505 task: &mut Task,
2506) -> Vec<QuotaLoss> {
2507 let repo = repo_for(task, &opts.repo);
2508 tracing::info!(
2509 "task {} — {} (repo {})",
2510 task.short(),
2511 task.title,
2512 repo.display()
2513 );
2514
2515 let mut config = match prepare(&repo, opts) {
2516 Ok(c) => c,
2517 Err(e) => {
2518 task.attempts += 1;
2522 task.fail(format!("config: {e:#}"), opts.max_attempts);
2523 record(queue, task);
2524 return Vec::new();
2525 }
2526 };
2527 apply_solo(&mut config, task);
2528 let start_failed = phrases(&config.graph.language).could_not_start;
2529
2530 if let Some(reason) = disk_gate(&repo, &config) {
2538 task.last_error = Some(reason.clone());
2539 task.hold_machine(Some(reason.clone()));
2540 record(queue, task);
2541 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2542 notices::raise(
2545 Notice::warn(
2546 &format!("disk:{}", repo.display()),
2547 "A task was held for want of free disk space; free some, then release it from the queue.",
2548 )
2549 .link(Link::Task {
2550 id: task.id.clone(),
2551 }),
2552 );
2553 return Vec::new();
2554 }
2555
2556 let unfinished = (!task.fresh_start)
2576 .then(|| unfinished_run(&task.runs, task.short()))
2577 .flatten();
2578 let review_branch = task.review_branch.take();
2584 let branch_exists = match &review_branch {
2585 Some(branch) => crate::git::branch_exists(&repo, branch)
2586 .await
2587 .unwrap_or(false),
2588 None => false,
2589 };
2590 let starter = choose_starter(
2591 review_branch.as_deref(),
2592 branch_exists,
2593 unfinished.as_deref(),
2594 );
2595 let attachments = match task_attachments(queue, task) {
2600 Ok(a) => a,
2601 Err(e) => {
2602 task.attempts += 1;
2603 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2604 record(queue, task);
2605 return Vec::new();
2606 }
2607 };
2608 let started = match &starter {
2609 Starter::Review(branch) => {
2610 tracing::info!(
2611 "task {} reopens `{branch}` as a review-only pass",
2612 task.short()
2613 );
2614 let takeover = crate::handover::Takeover {
2617 earlier: task.earlier_attempts().to_vec(),
2618 home: crate::run::home(),
2619 choice: take_divergence_answer(branch, &config.merge.remote, task),
2620 };
2621 Runner::review_taking_over(
2622 &repo,
2623 branch,
2624 config,
2625 Some(takeover),
2626 crate::run::Origin::queue(&task.id),
2627 )
2628 .await
2629 }
2630 Starter::Resume(id) => {
2631 tracing::info!("resuming run {id} rather than competing again");
2632 Runner::resume(id).map(|mut r| {
2633 if let Some(instruction) =
2634 prepare_instruction(&starter, Some(&r.state.instruction), task)
2635 {
2636 r.state.instruction = instruction;
2637 }
2638 r.state.attachments = attachments.clone();
2639 r
2640 })
2641 }
2642 Starter::Start => {
2643 if let Some(branch) = &review_branch {
2644 tracing::warn!(
2645 "conductor chose review for task {} but branch `{branch}` no longer \
2646 exists; requeuing as a fresh competition instead",
2647 task.short()
2648 );
2649 }
2650 let instruction = prepare_instruction(&starter, None, task)
2651 .unwrap_or_else(|| task.instruction.clone());
2652 Runner::start_naming(
2653 &repo,
2654 instruction,
2655 &task.title,
2656 config,
2657 crate::run::Origin::queue(&task.id),
2658 )
2659 .await
2660 .map(|mut r| {
2661 r.state.attachments = attachments.clone();
2662 r
2663 })
2664 }
2665 };
2666 let mut runner = match started {
2667 Ok(r) => r,
2668 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2672 let reason = format!("{start_failed}{e:#}");
2673 task.last_error = Some(reason.clone());
2674 task.hold_machine(Some(reason));
2675 record(queue, task);
2676 tracing::warn!(
2677 "holding {} for a branch it cannot take over: {e:#}",
2678 task.short()
2679 );
2680 notices::raise(
2682 Notice::warn(
2683 &format!("handover:{}", task.id),
2684 "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.",
2685 )
2686 .link(Link::Task {
2687 id: task.id.clone(),
2688 })
2689 .about([task.id.clone()]),
2690 );
2691 return Vec::new();
2692 }
2693 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2698 let d = e
2699 .downcast_ref::<crate::reconcile::Diverged>()
2700 .expect("checked by the guard");
2701 let mut q = ask::Question::new(
2702 task.id.clone(),
2703 "review".to_owned(),
2704 "sync".to_owned(),
2705 d.summary(),
2706 d.detail(),
2707 d.choices(),
2708 );
2709 match Questions::open().put(&mut q) {
2710 Ok(()) => {
2711 task.last_error = Some(format!("{e:#}"));
2712 task.review_branch = Some(d.branch.clone());
2715 task.block(vec![q.id.clone()], Some(d.summary()));
2716 }
2717 Err(put) => {
2718 tracing::warn!("could not file the divergence question: {put:#}");
2719 task.attempts += 1;
2720 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2721 }
2722 }
2723 record(queue, task);
2724 return Vec::new();
2725 }
2726 Err(e) => {
2727 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
2733 && let Starter::Review(branch) = &starter
2734 {
2735 task.review_branch = Some(branch.clone());
2736 task.last_error = Some(format!("{start_failed}{e:#}"));
2737 task.status = crate::queue::TaskStatus::Failed;
2738 record(queue, task);
2739 return Vec::new();
2740 }
2741 task.attempts += 1;
2742 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2743 record(queue, task);
2744 return Vec::new();
2745 }
2746 };
2747 runner.on_pause(stop.pause());
2749 runner.watch_interrupt(interrupt_pause);
2753
2754 let run = runner.state.id.clone();
2757 task.start(run.clone());
2758 record(queue, task);
2759 lock(status).current.push(Current {
2760 task: task.id.clone(),
2761 run,
2762 });
2763
2764 let quota_before = runner.state.quota.clone();
2767 let detail = match runner.execute().await {
2768 Ok(()) => describe(&runner.state),
2769 Err(e) => format!("{e:#}"),
2770 };
2771 let fresh = losses_this_attempt("a_before, &runner.state.quota);
2772 let verdict = Verdict {
2773 status: runner.state.status,
2774 left_pr: runner.state.pr.is_some(),
2777 quota_hit: !fresh.is_empty(),
2783 parked: runner.state.parked,
2787 no_viable_candidates: runner.state.viable().is_empty(),
2790 };
2791 settle_and_diagnose(task, verdict, &detail, opts.max_attempts, &runner.state);
2792 if task.status == TaskStatus::Done {
2793 supersede_prior_runs(task, &crate::run::home());
2794 }
2795 record(queue, task);
2796 tracing::info!(
2797 "task {} is {} after run {} ({})",
2798 task.short(),
2799 task.status.as_str(),
2800 runner.state.short(),
2801 label(runner.state.status)
2802 );
2803 fresh
2804}
2805
2806fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
2816 after
2817 .iter()
2818 .filter(|q| !before.contains(q))
2819 .cloned()
2820 .collect()
2821}
2822
2823fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
2826 if quota.is_empty() {
2827 return None;
2828 }
2829 let with_hint = quota.iter().find(|q| q.reset.is_some());
2830 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
2831 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
2832 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
2833 Some(
2834 now.checked_add(jiff::SignedDuration::from_secs(secs))
2835 .unwrap_or(Timestamp::MAX),
2836 )
2837}
2838
2839fn apply_solo(config: &mut Config, task: &Task) {
2849 if task.solo {
2850 config.graph.candidates = 1;
2851 }
2852}
2853
2854fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
2856 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
2857 if let Some(mode) = &opts.merge {
2858 config.merge.mode = merge_mode(mode)?;
2859 }
2860 Ok(config)
2861}
2862
2863async fn maybe_prune_cache_between_runs(
2894 repo: &Path,
2895 opts: &Opts,
2896 home: &Path,
2897 stop: &Stop,
2898 last_checked: &mut Option<Timestamp>,
2899 now: Timestamp,
2900) {
2901 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
2902 return;
2903 }
2904 *last_checked = Some(now);
2905 let cfg = match prepare(repo, opts) {
2906 Ok(cfg) => cfg,
2907 Err(e) => {
2908 tracing::warn!("cache check: no config: {e:#}");
2909 return;
2910 }
2911 };
2912 match clean::prune_cache_if_over_limit(&cfg, home) {
2913 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
2914 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
2915 pruned.files,
2916 pruned.freed
2917 ),
2918 Ok(_) => {}
2919 Err(e) => {
2920 tracing::warn!("housekeep: prune cache: {e:#}");
2921 notices::raise_in(
2922 home,
2923 Notice::warn(
2924 "housekeep:cache",
2925 "Pruning the shared build cache failed; disk usage may keep growing.",
2926 ),
2927 );
2928 }
2929 }
2930}
2931
2932fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
2936 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
2937}
2938
2939async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
2958 let cfg = match prepare(repo, opts) {
2959 Ok(cfg) => cfg,
2960 Err(e) => {
2961 tracing::warn!("housekeep: no config: {e:#}");
2962 return;
2963 }
2964 };
2965 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
2974 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
2975 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
2980 let mut extra = Vec::new();
2981 if out.unreadable > 0 {
2982 extra.push(format!("{} unreadable", out.unreadable));
2983 }
2984 if out.orphaned_worktrees > 0 {
2985 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
2986 }
2987 let detail = if extra.is_empty() {
2988 String::new()
2989 } else {
2990 format!(" ({})", extra.join(", "))
2991 };
2992 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
2993 }
2994 if out.external_merges_recorded > 0 {
2995 tracing::info!(
2996 "housekeep: recorded {} run(s) as merged externally",
2997 out.external_merges_recorded
2998 );
2999 }
3000 if out.cache_files > 0 {
3001 tracing::info!(
3002 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
3003 out.cache_files,
3004 out.cache_freed
3005 );
3006 }
3007 if out.questions_abandoned > 0 {
3008 tracing::info!(
3009 "housekeep: abandoned {} question(s) left open by a finished run",
3010 out.questions_abandoned
3011 );
3012 }
3013}
3014
3015async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3024 let questions = Questions::at(home.join("questions"));
3025 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3026 if report.is_empty() {
3027 return;
3028 }
3029 if !report.quarantined.is_empty() {
3030 tracing::info!(
3031 "triage: held {} blocked task(s) whose blocked-on task or \
3032 question no longer exists: {}",
3033 report.quarantined.len(),
3034 report.quarantined.join(", ")
3035 );
3036 }
3037 if !report.resumed.is_empty() {
3038 tracing::info!(
3039 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3040 report.resumed.len(),
3041 report.resumed.join(", ")
3042 );
3043 }
3044 if !report.asked.is_empty() {
3045 tracing::info!(
3046 "triage: asked about {} held task(s): {}",
3047 report.asked.len(),
3048 report.asked.join(", ")
3049 );
3050 }
3051 if !report.answered.is_empty() {
3052 tracing::info!(
3053 "triage: applied {} operator answer(s): {}",
3054 report.answered.len(),
3055 report.answered.join(", ")
3056 );
3057 }
3058}
3059
3060fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3067 disk_gate_with(repo, config, crate::disk::free_bytes)
3068}
3069
3070fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3074 repo: &Path,
3075 config: &Config,
3076 free_bytes: F,
3077) -> Option<String> {
3078 let min = config.disk.min_free_bytes;
3079 if min == 0 {
3080 return None;
3081 }
3082 match free_bytes(repo) {
3083 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3084 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3085 }
3086}
3087
3088const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3095
3096const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3100
3101fn quota_wait(
3110 reset_at: Option<Timestamp>,
3111 now: Timestamp,
3112 fallback: Duration,
3113 cap: Duration,
3114) -> Duration {
3115 match reset_at {
3116 Some(at) if at > now => {
3117 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3118 Duration::from_secs(secs).min(cap)
3119 }
3120 _ => fallback,
3121 }
3122}
3123
3124fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3138 parse_reset_hint_zoned(text, now)
3139 .or_else(|| parse_reset_hint_dated(text))
3140 .or_else(|| parse_reset_hint_relative(text, recorded))
3141}
3142
3143fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3147 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3148 let mut rest = rest.trim();
3149 if rest.is_empty() {
3150 return None;
3151 }
3152 let mut total: i64 = 0;
3153 let mut matched = false;
3154 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3155 if let Some((digits, tail)) = rest.split_once(unit)
3156 && !digits.is_empty()
3157 && digits.bytes().all(|b| b.is_ascii_digit())
3158 {
3159 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3160 rest = tail;
3161 matched = true;
3162 }
3163 }
3164 if !rest.is_empty() || !matched {
3165 return None;
3166 }
3167 recorded
3168 .checked_add(jiff::SignedDuration::from_secs(total))
3169 .ok()
3170}
3171
3172fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3176 let clock = clock.trim().to_lowercase();
3177 let (digits, pm) = clock
3178 .strip_suffix("am")
3179 .map(|d| (d, false))
3180 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3181 let (h, m) = digits.trim().split_once(':')?;
3182 let mut hour: i8 = h.trim().parse().ok()?;
3183 let minute: i8 = m.trim().parse().ok()?;
3184 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3185 return None;
3186 }
3187 if pm && hour != 12 {
3188 hour += 12;
3189 } else if !pm && hour == 12 {
3190 hour = 0;
3191 }
3192 Some((hour, minute))
3193}
3194
3195fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3200 let open = text.find('(')?;
3201 let close = text.rfind(')')?;
3202 if close <= open {
3203 return None;
3204 }
3205 let zone = text[open + 1..close].trim();
3206 let (hour, minute) = parse_12h_clock(&text[..open])?;
3207 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3208 let candidate = now
3209 .to_zoned(tz)
3210 .with()
3211 .hour(hour)
3212 .minute(minute)
3213 .second(0)
3214 .millisecond(0)
3215 .microsecond(0)
3216 .nanosecond(0)
3217 .build()
3218 .ok()?;
3219 let mut at = candidate.timestamp();
3220 if at <= now {
3221 at += jiff::SignedDuration::from_hours(24);
3222 }
3223 Some(at)
3224}
3225
3226fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3235 let words: Vec<&str> = text.split_whitespace().collect();
3236 if words.len() < 5 {
3237 return None;
3238 }
3239 (0..=words.len() - 5)
3240 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3241}
3242
3243fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3249 if trailing.is_some_and(|next| next.starts_with('(')) {
3250 return None;
3251 }
3252 let month = month_number(window[0])?;
3253 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3254 let day_digits = ["st", "nd", "rd", "th"]
3255 .iter()
3256 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3257 let day: i8 = day_digits.parse().ok()?;
3258 let year_token = window[2];
3259 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3260 return None;
3261 }
3262 let year: i16 = year_token.parse().ok()?;
3263 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3267 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3268 let date = jiff::civil::Date::new(year, month, day).ok()?;
3269 let candidate = date
3270 .at(hour, minute, 0, 0)
3271 .to_zoned(jiff::tz::TimeZone::UTC)
3272 .ok()?;
3273 Some(candidate.timestamp())
3274}
3275
3276fn month_number(name: &str) -> Option<i8> {
3279 const NAMES: [&str; 12] = [
3280 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3281 ];
3282 let lower = name.to_lowercase();
3283 NAMES
3284 .iter()
3285 .position(|n| *n == lower.as_str())
3286 .map(|i| i as i8 + 1)
3287}
3288
3289fn exhausted_review_budget(state: &RunState) -> bool {
3301 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3302}
3303
3304fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3341 unfinished_run_with(runs, short, RunState::load)
3342}
3343
3344fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3347where
3348 F: FnOnce(&str) -> Result<RunState>,
3349{
3350 let id = runs.last()?;
3351 match load(id) {
3352 Ok(s)
3357 if s.status.resumable()
3358 && !s.released()
3359 && !exhausted_review_budget(&s)
3360 && s.liveness(false) != crate::run::Liveness::Live =>
3361 {
3362 Some(id.clone())
3363 }
3364 Ok(_) => None,
3365 Err(e) => {
3366 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3367 None
3368 }
3369 }
3370}
3371
3372fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3380where
3381 F: FnOnce(&str) -> Result<RunState>,
3382{
3383 if task.status != TaskStatus::Failed || task.fresh_start || task.review_branch.is_some() {
3384 return false;
3385 }
3386 let Some(id) = task.runs.last() else {
3387 return false;
3388 };
3389 load(id).is_ok_and(|s| {
3390 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3391 })
3392}
3393
3394#[derive(Debug, Clone, PartialEq, Eq)]
3397enum Starter {
3398 Review(String),
3401 Resume(String),
3403 Start,
3405}
3406
3407fn take_divergence_answer(
3411 branch: &str,
3412 remote: &str,
3413 task: &mut Task,
3414) -> Option<crate::reconcile::Choice> {
3415 let summary = crate::reconcile::summary_for(branch, remote);
3416 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3417 (a.question == summary)
3418 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3419 .flatten()
3420 .map(|c| (i, c))
3421 })?;
3422 task.answers.remove(idx);
3423 Some(choice)
3424}
3425
3426fn choose_starter(
3438 review_branch: Option<&str>,
3439 branch_exists: bool,
3440 unfinished: Option<&str>,
3441) -> Starter {
3442 match review_branch {
3443 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3444 Some(_) => Starter::Start,
3445 None => match unfinished {
3446 Some(id) => Starter::Resume(id.to_owned()),
3447 None => Starter::Start,
3448 },
3449 }
3450}
3451
3452fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3455 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3456 return fallback.to_path_buf();
3457 }
3458 task.repo.clone()
3459}
3460
3461const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3465
3466fn answers_block(task: &Task, count: usize) -> String {
3468 let mut s = ANSWERS_HEADER.to_owned();
3469 for a in &task.answers[..count] {
3470 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3471 }
3472 s
3473}
3474
3475fn append_answers(base: &str, task: &Task) -> String {
3478 if task.answers.is_empty() {
3479 return base.to_owned();
3480 }
3481 let mut s = base.to_owned();
3482 s.push_str(&answers_block(task, task.answers.len()));
3483 s
3484}
3485
3486fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3490 for count in (1..=task.answers.len()).rev() {
3491 let block = answers_block(task, count);
3492 if let Some(base) = instruction.strip_suffix(&block) {
3493 return base;
3494 }
3495 }
3496 instruction
3497}
3498
3499fn instruction_for(task: &Task) -> String {
3507 append_answers(&task.instruction, task)
3508}
3509
3510fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3522 append_answers(strip_answers_block(old_instruction, task), task)
3523}
3524
3525fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3527 let paths = queue.attachment_paths(task);
3528 for (name, path) in task.attachments.iter().zip(&paths) {
3529 if !path.is_file() {
3530 bail!(
3531 "attachment `{name}` is recorded on the task but {} is missing",
3532 path.display()
3533 );
3534 }
3535 }
3536 Ok(paths)
3537}
3538
3539fn prepare_instruction(
3550 starter: &Starter,
3551 old_instruction: Option<&str>,
3552 task: &Task,
3553) -> Option<String> {
3554 match starter {
3555 Starter::Start => Some(instruction_for(task)),
3556 Starter::Resume(_) => Some(resumed_instruction(
3557 old_instruction.expect("a resumed run always has a prior instruction"),
3558 task,
3559 )),
3560 Starter::Review(_) => None,
3561 }
3562}
3563
3564fn record(queue: &Queue, task: &mut Task) {
3568 if let Err(e) = queue.put(task) {
3569 tracing::error!("could not record task {}: {e:#}", task.short());
3570 notices::raise(Notice::error(
3571 "loop:record",
3572 "The loop could not save a task's state; check the disk.",
3573 ));
3574 }
3575}
3576
3577fn runnable(queue: &Queue) -> Vec<Task> {
3583 let mut tasks: Vec<Task> = queue
3584 .list()
3585 .into_iter()
3586 .filter(|t| t.status.runnable())
3587 .collect();
3588 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
3589 tasks
3590}
3591
3592fn describe(state: &RunState) -> String {
3606 let p = phrases(&state.config.graph.language);
3607 let mut detail = if state.status == RunStatus::Stalled {
3608 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
3609 seats.sort_unstable();
3610 seats.dedup();
3611 if seats.is_empty() {
3612 p.quorum_lost.to_owned()
3613 } else {
3614 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
3615 }
3616 } else {
3617 format!("{}{}", p.run_ended, state.status.display_label())
3618 };
3619 if let Some(last) = state.events.last() {
3620 detail.push_str(&format!(" ({}: {})", last.node, last.message));
3621 }
3622 detail.push_str(&format!(" [run {}]", state.id));
3623 detail
3624}
3625
3626const DIAGNOSTIC_MAX: usize = 4_000;
3632
3633const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
3638
3639fn diagnostic(state: &RunState) -> Option<String> {
3653 let mut parts: Vec<String> = Vec::new();
3654
3655 for o in state.gate.iter().filter(|o| !o.ok()) {
3657 parts.push(format!(
3658 "gate `{}` failed ({:?}):\n{}",
3659 o.command,
3660 o.code,
3661 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
3662 ));
3663 }
3664
3665 if let Some(last) = state
3668 .events
3669 .iter()
3670 .rev()
3671 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
3672 {
3673 parts.push(last.message.clone());
3674 }
3675
3676 if state.viable().is_empty() {
3683 for c in &state.candidates {
3684 if let Some(evidence) = &c.verified_noop {
3685 parts.push(format!(
3686 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
3687 c.label
3688 ));
3689 } else if !c.summary.trim().is_empty() {
3690 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
3691 } else if let Some(why) = &c.failed {
3692 parts.push(format!("candidate {}: {why}", c.label));
3693 }
3694 }
3695 }
3696
3697 if parts.is_empty() {
3698 return None;
3699 }
3700 Some(crate::run::tail(
3705 &parts.join("\n\n"),
3706 DIAGNOSTIC_MAX.saturating_sub(100),
3707 ))
3708}
3709
3710fn label(status: RunStatus) -> &'static str {
3718 status.as_str()
3719}
3720
3721fn merge_mode(mode: &str) -> Result<MergeMode> {
3723 match mode {
3724 "none" => Ok(MergeMode::None),
3725 "local" => Ok(MergeMode::Local),
3726 "pr" => Ok(MergeMode::Pr),
3727 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
3728 }
3729}
3730
3731fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
3736 mutex
3737 .lock()
3738 .unwrap_or_else(std::sync::PoisonError::into_inner)
3739}
3740
3741#[cfg(test)]
3742mod tests {
3743 use super::*;
3744 use crate::queue::{Source, TaskStatus};
3745 use crate::run::{Candidate, CommandOutcome};
3746 use pretty_assertions::assert_eq;
3747
3748 fn task() -> Task {
3749 Task::new(
3750 "add retries".to_owned(),
3751 "add retries".to_owned(),
3752 PathBuf::from("/repo"),
3753 Source::Human,
3754 )
3755 }
3756
3757 fn interrupt_task(id: &str) -> Task {
3760 let mut t = task();
3761 t.id = id.to_owned();
3762 t.interrupt = true;
3763 t
3764 }
3765
3766 fn task_with_id(id: &str) -> Task {
3768 let mut t = task();
3769 t.id = id.to_owned();
3770 t
3771 }
3772
3773 fn urgent_task(id: &str) -> Task {
3775 let mut t = task();
3776 t.id = id.to_owned();
3777 t.urgent = true;
3778 t
3779 }
3780
3781 #[test]
3786 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
3787 assert_eq!(permit_kind(true, false), PermitKind::None);
3788 assert_eq!(permit_kind(true, true), PermitKind::None);
3789 }
3790
3791 #[test]
3796 fn permit_kind_separates_urgent_from_ordinary() {
3797 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
3798 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
3799 }
3800
3801 #[test]
3808 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
3809 let cfg = Config::default();
3810 let repo = Path::new("/any/repo/path");
3811
3812 let reason =
3813 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
3814 assert!(reason.contains("1024"), "{reason}");
3815 assert!(
3816 reason.contains(&cfg.disk.min_free_bytes.to_string()),
3817 "{reason}"
3818 );
3819
3820 assert_eq!(
3821 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
3822 None,
3823 "exactly at the floor is open"
3824 );
3825 assert_eq!(
3826 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
3827 None,
3828 "comfortably above the floor is open"
3829 );
3830 }
3831
3832 #[test]
3833 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
3834 let mut cfg = Config::default();
3835 cfg.disk.min_free_bytes = 0;
3836 let repo = Path::new("/any/repo/path");
3837 assert_eq!(
3838 disk_gate_with(repo, &cfg, |_| Ok(0)),
3839 None,
3840 "a zero floor never measures at all"
3841 );
3842 }
3843
3844 #[test]
3845 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
3846 let cfg = Config::default();
3847 let repo = Path::new("/any/repo/path");
3848 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
3849 .expect("a measurement failure must close the gate, not open it");
3850 assert!(reason.contains("could not measure"), "{reason}");
3851 }
3852
3853 #[test]
3854 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
3855 let ordinary = task();
3856 let next = advance_interrupt(
3857 Interrupt::Idle,
3858 std::slice::from_ref(&ordinary.id),
3859 std::slice::from_ref(&ordinary),
3860 );
3861 assert_eq!(next, Interrupt::Idle);
3862 }
3863
3864 #[test]
3865 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
3866 let marked = interrupt_task("marked");
3869 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
3870 assert_eq!(next, Interrupt::Idle);
3871 }
3872
3873 #[test]
3874 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
3875 let marked = interrupt_task("marked");
3876 let next = advance_interrupt(
3877 Interrupt::Idle,
3878 &["running".to_owned()],
3879 std::slice::from_ref(&marked),
3880 );
3881 assert_eq!(
3882 next,
3883 Interrupt::Parking {
3884 parked: vec!["running".to_owned()],
3885 interrupt_task: "marked".to_owned(),
3886 }
3887 );
3888 }
3889
3890 #[test]
3899 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
3900 let marked = interrupt_task("marked");
3901
3902 let two = advance_interrupt(
3903 Interrupt::Idle,
3904 &["a".to_owned(), "b".to_owned()],
3905 std::slice::from_ref(&marked),
3906 );
3907 assert_eq!(two, Interrupt::Idle);
3908
3909 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
3910 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
3911 }
3912
3913 #[test]
3914 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
3915 let state = Interrupt::Parking {
3916 parked: vec!["running".to_owned()],
3917 interrupt_task: "marked".to_owned(),
3918 };
3919 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
3921 assert_eq!(still_going, state);
3922
3923 let stopped_but_not_yet_dispatched =
3927 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
3928 assert_eq!(stopped_but_not_yet_dispatched, state);
3929
3930 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
3932 assert_eq!(
3933 dispatched,
3934 Interrupt::Running {
3935 parked: vec!["running".to_owned()],
3936 interrupt_task: "marked".to_owned(),
3937 }
3938 );
3939 }
3940
3941 #[test]
3942 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
3943 let state = Interrupt::Running {
3944 parked: vec!["running".to_owned()],
3945 interrupt_task: "marked".to_owned(),
3946 };
3947 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
3948 assert_eq!(still_running, state);
3949
3950 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
3957 assert_eq!(
3958 ended,
3959 Interrupt::Resuming {
3960 parked: vec!["running".to_owned()]
3961 }
3962 );
3963 }
3964
3965 #[test]
3966 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
3967 let state = Interrupt::Resuming {
3968 parked: vec!["running".to_owned()],
3969 };
3970 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
3971 assert_eq!(still_waiting, state);
3972
3973 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
3974 assert_eq!(dispatched, Interrupt::Idle);
3975 }
3976
3977 #[test]
3983 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
3984 {
3985 let state = Interrupt::Parking {
3986 parked: vec!["running".to_owned()],
3987 interrupt_task: "marked".to_owned(),
3988 };
3989 let next = advance_interrupt(state, &[], &[]);
3992 assert_eq!(
3993 next,
3994 Interrupt::Resuming {
3995 parked: vec!["running".to_owned()]
3996 },
3997 "abandoning the interrupt must not abandon the resume it owes"
3998 );
3999 }
4000
4001 #[test]
4004 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
4005 let state = Interrupt::Resuming {
4006 parked: vec!["running".to_owned()],
4007 };
4008 let next = advance_interrupt(state, &[], &[]);
4009 assert_eq!(
4010 next,
4011 Interrupt::Idle,
4012 "nothing is left to wait for; the loop must not stay wedged"
4013 );
4014 }
4015
4016 #[test]
4017 fn disabled_by_config_the_sequence_can_never_leave_idle() {
4018 let marked = interrupt_task("marked");
4019 let next = advance_interrupt_tick(
4020 false,
4021 Interrupt::Idle,
4022 &["running".to_owned()],
4023 std::slice::from_ref(&marked),
4024 );
4025 assert_eq!(
4026 next,
4027 Interrupt::Idle,
4028 "an unmarked, unconfigured daemon must behave exactly as before"
4029 );
4030 }
4031
4032 #[test]
4033 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4034 let state = Interrupt::Parking {
4035 parked: vec!["running".to_owned()],
4036 interrupt_task: "marked".to_owned(),
4037 };
4038 let candidates = vec![interrupt_task("marked"), task()];
4039 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4040 assert!(
4041 allowed.is_empty(),
4042 "nothing may dispatch - not even the interrupt task itself - \
4043 until the parked run has actually stopped"
4044 );
4045 }
4046
4047 #[test]
4061 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4062 for state in [
4063 Interrupt::Parking {
4064 parked: vec!["running".to_owned()],
4065 interrupt_task: "marked".to_owned(),
4066 },
4067 Interrupt::Running {
4068 parked: vec!["running".to_owned()],
4069 interrupt_task: "marked".to_owned(),
4070 },
4071 Interrupt::Resuming {
4072 parked: vec!["running".to_owned()],
4073 },
4074 ] {
4075 let candidates = vec![
4076 interrupt_task("marked"),
4077 urgent_task("hot"),
4078 task_with_id("ordinary"),
4079 ];
4080 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4081 assert!(
4082 !allowed.iter().any(|t| t.id == "hot"),
4083 "an urgent candidate must wait out the same gate as anything \
4084 else while the run it would run alongside has not actually \
4085 left flight, for state {state:?}: {allowed:?}"
4086 );
4087 }
4088 }
4089
4090 #[test]
4096 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4097 let state = Interrupt::Resuming {
4098 parked: vec!["hot".to_owned()],
4099 };
4100 let candidates = vec![urgent_task("hot"), task()];
4101 let allowed = interrupt_gate(&state, &[], candidates);
4102 assert_eq!(
4103 allowed.iter().filter(|t| t.id == "hot").count(),
4104 1,
4105 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4106 );
4107 }
4108
4109 #[test]
4110 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4111 let state = Interrupt::Parking {
4112 parked: vec!["running".to_owned()],
4113 interrupt_task: "marked".to_owned(),
4114 };
4115 let other = task();
4116 let candidates = vec![interrupt_task("marked"), other.clone()];
4117 let allowed = interrupt_gate(&state, &[], candidates);
4118 assert_eq!(allowed.len(), 1);
4119 assert_eq!(allowed[0].id, "marked");
4120 }
4121
4122 #[test]
4123 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4124 let state = Interrupt::Running {
4125 parked: vec!["running".to_owned()],
4126 interrupt_task: "marked".to_owned(),
4127 };
4128 let candidates = vec![task(), task()];
4129 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4130 assert!(allowed.is_empty());
4131 }
4132
4133 #[test]
4140 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4141 let state = Interrupt::Resuming {
4142 parked: vec!["a".to_owned(), "c".to_owned()],
4143 };
4144 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4145 let allowed = interrupt_gate(&state, &[], candidates);
4146 assert_eq!(
4147 allowed.len(),
4148 1,
4149 "at most one candidate may be offered while resuming: {allowed:?}"
4150 );
4151 assert_eq!(allowed[0].id, "a");
4152 }
4153
4154 #[test]
4155 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4156 let state = Interrupt::Resuming {
4157 parked: vec!["a".to_owned()],
4158 };
4159 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4160 assert!(allowed.is_empty());
4161 }
4162
4163 #[test]
4169 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4170 let running = task(); let marked = interrupt_task("marked");
4172
4173 let mut state = Interrupt::Idle;
4174 let in_flight = vec![running.id.clone()];
4176 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4177 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4178 assert!(gated.is_empty(), "still waiting on `running` to park");
4179
4180 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4182 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4183 assert_eq!(
4184 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4185 vec!["marked"],
4186 "only the interrupt task may be offered to the dispatcher now"
4187 );
4188
4189 state = advance_interrupt_tick(
4191 true,
4192 state,
4193 &["marked".to_owned()],
4194 std::slice::from_ref(&running),
4195 );
4196 let gated = interrupt_gate(
4197 &state,
4198 &["marked".to_owned()],
4199 vec![marked.clone(), running.clone()],
4200 );
4201 assert!(
4202 gated.is_empty(),
4203 "the parked run must not be offered back while the interrupt \
4204 task is still running"
4205 );
4206
4207 let other = task_with_id("other");
4211 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4212 assert_eq!(
4213 state,
4214 Interrupt::Resuming {
4215 parked: vec![running.id.clone()]
4216 }
4217 );
4218 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4219 assert_eq!(
4220 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4221 vec![running.id.as_str()],
4222 "exactly the parked run resumes - not the unrelated task, even \
4223 though it was offered first"
4224 );
4225
4226 state = advance_interrupt_tick(
4230 true,
4231 state,
4232 std::slice::from_ref(&running.id),
4233 std::slice::from_ref(&other),
4234 );
4235 assert_eq!(state, Interrupt::Idle);
4236 let gated = interrupt_gate(
4237 &state,
4238 std::slice::from_ref(&running.id),
4239 vec![other.clone()],
4240 );
4241 assert_eq!(
4242 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4243 vec![other.id.as_str()],
4244 "ordinary dispatch is unrestricted again"
4245 );
4246 }
4247
4248 #[test]
4249 fn every_run_status_settles_the_task_it_came_from() {
4250 let table = [
4252 (RunStatus::Merged, TaskStatus::Done, 1),
4253 (RunStatus::Ready, TaskStatus::Done, 1),
4254 (RunStatus::Stalled, TaskStatus::Failed, 0),
4255 (RunStatus::Blocked, TaskStatus::Failed, 1),
4256 (RunStatus::Failed, TaskStatus::Failed, 1),
4257 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4258 (RunStatus::Prep, TaskStatus::Failed, 1),
4259 (RunStatus::Implementing, TaskStatus::Failed, 1),
4260 (RunStatus::Judging, TaskStatus::Failed, 1),
4261 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4262 (RunStatus::Voting, TaskStatus::Failed, 1),
4263 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4264 (RunStatus::Gating, TaskStatus::Failed, 1),
4265 ];
4266 for (run, want, attempts) in table {
4267 let mut t = task();
4268 t.start("20260902-000000-aaaa".to_owned());
4269 settle(
4270 &mut t,
4271 Verdict {
4272 status: run,
4273 left_pr: false,
4274 parked: false,
4275 quota_hit: matches!(run, RunStatus::Stalled),
4276 no_viable_candidates: false,
4277 },
4278 "why",
4279 2,
4280 );
4281 assert_eq!(t.status, want, "task status after {}", label(run));
4282 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4283 }
4284 }
4285
4286 #[test]
4287 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4288 let mut stalled = task();
4289 stalled.start("20260902-000000-aaaa".to_owned());
4290 settle(
4291 &mut stalled,
4292 Verdict {
4293 status: RunStatus::Stalled,
4294 left_pr: false,
4295 parked: false,
4296 quota_hit: true,
4297 no_viable_candidates: false,
4298 },
4299 "quota",
4300 1,
4301 );
4302 assert_eq!(stalled.attempts, 0);
4303 assert!(
4304 stalled.status.runnable(),
4305 "a machine problem must leave the task in line"
4306 );
4307
4308 let mut blocked = task();
4309 blocked.start("20260902-000000-aaaa".to_owned());
4310 settle(
4311 &mut blocked,
4312 Verdict {
4313 status: RunStatus::Blocked,
4314 left_pr: false,
4315 parked: false,
4316 quota_hit: false,
4317 no_viable_candidates: false,
4318 },
4319 "findings open",
4320 1,
4321 );
4322 assert_eq!(blocked.attempts, 1);
4323 assert_eq!(
4324 blocked.status,
4325 TaskStatus::Held,
4326 "the last attempt hands the task to a human"
4327 );
4328 }
4329
4330 #[test]
4331 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4332 let mut delivered = task();
4335 delivered.start("20260903-080619-01c2".to_owned());
4336 settle(
4337 &mut delivered,
4338 Verdict {
4339 status: RunStatus::Blocked,
4340 left_pr: true,
4341 parked: false,
4342 quota_hit: false,
4343 no_viable_candidates: false,
4344 },
4345 "no check status",
4346 4,
4347 );
4348 assert_eq!(
4349 delivered.status,
4350 TaskStatus::Held,
4351 "a pull request waiting on CI or a person is not a retryable failure"
4352 );
4353 assert!(
4354 !delivered.status.runnable(),
4355 "the loop must not pick this task up again"
4356 );
4357 assert_eq!(
4358 delivered.last_error.as_deref(),
4359 Some("no check status"),
4360 "the operator needs to be told what the gate was waiting for"
4361 );
4362
4363 let mut empty_handed = task();
4366 empty_handed.start("20260903-080619-01c2".to_owned());
4367 settle(
4368 &mut empty_handed,
4369 Verdict {
4370 status: RunStatus::Blocked,
4371 left_pr: false,
4372 parked: false,
4373 quota_hit: false,
4374 no_viable_candidates: false,
4375 },
4376 "findings open",
4377 4,
4378 );
4379 assert_eq!(empty_handed.status, TaskStatus::Failed);
4380 assert!(empty_handed.status.runnable());
4381 }
4382
4383 #[test]
4384 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4385 let mut noop = task();
4391 noop.start("20260912-131304-391f".to_owned());
4392 settle(
4393 &mut noop,
4394 Verdict {
4395 status: RunStatus::VerifiedNoop,
4396 left_pr: false,
4397 parked: false,
4398 quota_hit: false,
4399 no_viable_candidates: true,
4400 },
4401 "candidate A: already fixed by b32cfc4, on main",
4402 4,
4403 );
4404 assert_eq!(
4405 noop.status,
4406 TaskStatus::Held,
4407 "an unverified claim is a request for a human, not a failure"
4408 );
4409 assert!(
4410 !noop.status.runnable(),
4411 "the loop must not requeue this on the same unverified claim"
4412 );
4413 assert_eq!(noop.attempts, 1);
4418 }
4419
4420 #[test]
4421 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4422 let mut parked = task();
4427 parked.start("20260903-183634-2d98".to_owned());
4428 settle(
4429 &mut parked,
4430 Verdict {
4431 status: RunStatus::Implementing,
4432 left_pr: false,
4433 quota_hit: false,
4434 parked: true,
4435 no_viable_candidates: false,
4436 },
4437 "parked after `implementing`",
4438 2,
4439 );
4440 assert_eq!(parked.attempts, 0, "a park is refunded");
4441 assert!(
4442 parked.status.runnable(),
4443 "and the task stays in line so the next loop resumes its run"
4444 );
4445 assert_eq!(
4446 parked.last_error.as_deref(),
4447 Some("parked after `implementing`"),
4448 "the card says where it stopped"
4449 );
4450
4451 let mut broken = task();
4455 broken.start("20260903-183634-2d98".to_owned());
4456 settle(
4457 &mut broken,
4458 Verdict {
4459 status: RunStatus::Implementing,
4460 left_pr: false,
4461 quota_hit: false,
4462 parked: false,
4463 no_viable_candidates: false,
4464 },
4465 "returned mid-flight",
4466 2,
4467 );
4468 assert_eq!(broken.attempts, 1);
4469 }
4470
4471 #[test]
4472 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4473 let mut flaky = task();
4478 flaky.start("20260903-123023-e633".to_owned());
4479 settle(
4480 &mut flaky,
4481 Verdict {
4482 status: RunStatus::Stalled,
4483 left_pr: false,
4484 parked: false,
4485 quota_hit: false,
4486 no_viable_candidates: false,
4487 },
4488 "verdict rests on 1 of 3 judges",
4489 2,
4490 );
4491 assert_eq!(
4492 flaky.attempts, 1,
4493 "flakiness spends an attempt, so `max_attempts` still bounds it"
4494 );
4495 assert!(flaky.status.runnable(), "and it is still worth retrying");
4496
4497 let mut limited = task();
4499 limited.start("20260903-123023-e633".to_owned());
4500 settle(
4501 &mut limited,
4502 Verdict {
4503 status: RunStatus::Stalled,
4504 left_pr: false,
4505 parked: false,
4506 quota_hit: true,
4507 no_viable_candidates: false,
4508 },
4509 "judge-2, judge-3 out of quota",
4510 2,
4511 );
4512 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4513 assert!(limited.status.runnable());
4514
4515 let mut worn = task();
4518 for _ in 0..2 {
4519 worn.release();
4520 }
4521 worn.start("20260903-123023-e633".to_owned());
4522 worn.attempts = 2;
4523 settle(
4524 &mut worn,
4525 Verdict {
4526 status: RunStatus::Stalled,
4527 left_pr: false,
4528 parked: false,
4529 quota_hit: false,
4530 no_viable_candidates: false,
4531 },
4532 "no quorum again",
4533 2,
4534 );
4535 assert_eq!(worn.status, TaskStatus::Held);
4536 assert!(!worn.status.runnable());
4537 }
4538
4539 #[test]
4540 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4541 let mut wiped_out = task();
4548 wiped_out.start("20260907-025000-a1b2".to_owned());
4549 settle(
4550 &mut wiped_out,
4551 Verdict {
4552 status: RunStatus::Failed,
4553 left_pr: false,
4554 parked: false,
4555 quota_hit: true,
4556 no_viable_candidates: true,
4557 },
4558 "no candidate produced a change; nothing to judge",
4559 2,
4560 );
4561 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
4562 assert!(
4563 wiped_out.status.runnable(),
4564 "a machine problem must leave the task in line"
4565 );
4566
4567 let mut partial_progress = task();
4573 partial_progress.start("20260907-025500-c3d4".to_owned());
4574 settle(
4575 &mut partial_progress,
4576 Verdict {
4577 status: RunStatus::Failed,
4578 left_pr: false,
4579 parked: false,
4580 quota_hit: true,
4581 no_viable_candidates: false,
4582 },
4583 "gate failed on the winning candidate",
4584 2,
4585 );
4586 assert_eq!(
4587 partial_progress.attempts, 1,
4588 "a candidate that actually produced a change spends the attempt \
4589 even though some other seat hit its quota"
4590 );
4591 assert!(partial_progress.status.runnable());
4592 }
4593
4594 #[test]
4595 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
4596 let mut t = task();
4602 t.start("20260907-025000-a1b2".to_owned());
4603 let mut state = run_state(RunStatus::Failed);
4604 state.quota.push(QuotaLoss {
4605 seat: "cand-a".to_owned(),
4606 node: "implement".to_owned(),
4607 at: Timestamp::now(),
4608 reset: None,
4609 });
4610 assert!(
4611 state.viable().is_empty(),
4612 "no candidate was added, so nothing is viable"
4613 );
4614 reclaim(&mut t, Some(state), 2, "en");
4615 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
4616 assert!(t.status.runnable());
4617 }
4618
4619 #[test]
4620 fn a_held_task_is_never_offered_to_the_loop() {
4621 let dir = tempfile::tempdir().unwrap();
4622 let queue = Queue::at(dir.path().to_path_buf());
4623 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
4624 let mut t = task();
4625 t.id = format!("2026090{n}-000000-000{n}");
4626 t.priority = priority;
4627 queue.put(&mut t).unwrap();
4628 }
4629 let mut held = task();
4630 held.id = "20260909-000000-9999".to_owned();
4631 held.priority = 99;
4632 held.hold_machine(None);
4633 queue.put(&mut held).unwrap();
4634
4635 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
4636 assert_eq!(order.len(), 3);
4637 assert!(!order.contains(&held.id));
4638 assert_eq!(
4639 order.first().cloned(),
4640 queue.next_runnable().map(|t| t.id),
4641 "the loop's first candidate is exactly what the queue offers"
4642 );
4643 assert_eq!(
4644 order,
4645 vec![
4646 "20260902-000000-0002".to_owned(),
4647 "20260903-000000-0003".to_owned(),
4648 "20260901-000000-0001".to_owned(),
4649 ],
4650 "priority first, then oldest, so nothing starves"
4651 );
4652 }
4653
4654 #[test]
4655 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
4656 let dir = tempfile::tempdir().unwrap();
4657 let queue = Queue::at(dir.path().to_path_buf());
4658 let mut old = task();
4659 old.id = "20260101-000000-old0".to_owned();
4660 queue.put(&mut old).unwrap();
4661 let mut fresh = task();
4662 fresh.id = "20260101-000000-new0".to_owned();
4663 queue.put(&mut fresh).unwrap();
4664
4665 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
4669 std::thread::sleep(Duration::from_millis(60));
4670 let live = queue.claim(&fresh.id).unwrap();
4671
4672 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4673 assert_eq!(swept, vec![old.id.clone()]);
4674 assert!(
4675 queue.claim(&old.id).is_ok(),
4676 "an unparseable lock older than the threshold is swept"
4677 );
4678 assert!(
4679 queue.claim(&fresh.id).is_err(),
4680 "a live pid protects its lock regardless of age"
4681 );
4682 drop(live);
4683 }
4684
4685 #[test]
4686 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
4687 let dir = tempfile::tempdir().unwrap();
4697 let queue = Queue::at(dir.path().to_path_buf());
4698 let mut t = task();
4699 t.id = "20260101-000000-live".to_owned();
4700 queue.put(&mut t).unwrap();
4701
4702 let claim = queue.claim(&t.id).unwrap();
4703 std::thread::sleep(Duration::from_millis(60));
4704
4705 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4706 assert!(
4707 swept.is_empty(),
4708 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
4709 );
4710 assert!(
4711 queue.claim(&t.id).is_err(),
4712 "the lock still protects its task"
4713 );
4714 drop(claim);
4715 }
4716
4717 fn injected_dead_pid() -> u32 {
4720 std::process::id().checked_add(1).unwrap_or(1)
4721 }
4722
4723 #[test]
4724 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
4725 let dir = tempfile::tempdir().unwrap();
4726 let queue = Queue::at(dir.path().to_path_buf());
4727 let mut t = task();
4728 t.id = "20260101-000000-dead".to_owned();
4729 queue.put(&mut t).unwrap();
4730 let dead_pid = injected_dead_pid();
4731
4732 std::fs::write(
4737 dir.path().join(format!("{}.lock", t.id)),
4738 dead_pid.to_string(),
4739 )
4740 .unwrap();
4741
4742 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4743 pid != dead_pid
4744 });
4745 assert_eq!(
4746 swept,
4747 vec![t.id.clone()],
4748 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
4749 );
4750 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
4751 }
4752
4753 #[test]
4754 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
4755 let dir = tempfile::tempdir().unwrap();
4756 let queue = Queue::at(dir.path().to_path_buf());
4757 let mut t = task();
4758 t.id = "20260101-000000-late".to_owned();
4759 queue.put(&mut t).unwrap();
4760 let dead_pid = injected_dead_pid();
4761
4762 assert!(
4765 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
4766 "nothing has claimed the task yet"
4767 );
4768
4769 std::fs::write(
4772 dir.path().join(format!("{}.lock", t.id)),
4773 dead_pid.to_string(),
4774 )
4775 .unwrap();
4776
4777 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4781 pid != dead_pid
4782 });
4783 assert_eq!(swept, vec![t.id.clone()]);
4784 }
4785
4786 #[test]
4787 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
4788 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
4793 let dir = tempfile::tempdir().unwrap();
4794 let queue = Queue::at(dir.path().to_path_buf());
4795 let mut t = task();
4796 t.id = "20260101-000000-crsh".to_owned();
4797 t.status = TaskStatus::Running;
4798 t.attempts = 1;
4799 t.runs.push("20260904-000000-4043".to_owned());
4803 queue.put(&mut t).unwrap();
4804 let dead_pid = injected_dead_pid();
4805
4806 std::fs::write(
4809 dir.path().join(format!("{}.lock", t.id)),
4810 dead_pid.to_string(),
4811 )
4812 .unwrap();
4813
4814 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
4820 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
4821
4822 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4823 pid != dead_pid
4824 });
4825 assert_eq!(swept, vec![t.id.clone()]);
4826
4827 let reclaimed = reclaim_orphaned_running(&queue, 2);
4828 assert_eq!(reclaimed, vec![t.id.clone()]);
4829 let after = queue.get(&t.id).unwrap();
4830 assert_eq!(
4831 after.status,
4832 TaskStatus::Held,
4833 "no run.json to recover from, so a human is asked"
4834 );
4835 assert_eq!(
4836 after.runs,
4837 vec!["20260904-000000-4043".to_owned()],
4838 "the crashed run's id is kept as evidence, not discarded"
4839 );
4840 }
4841
4842 #[test]
4843 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
4844 let dir = tempfile::tempdir().unwrap();
4845 let queue = Queue::at(dir.path().to_path_buf());
4846 let mut t = task();
4847 t.id = "20260101-000000-unknown".to_owned();
4848 queue.put(&mut t).unwrap();
4849 let dead_pid = injected_dead_pid();
4850 std::fs::write(
4851 dir.path().join(format!("{}.lock", t.id)),
4852 dead_pid.to_string(),
4853 )
4854 .unwrap();
4855
4856 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
4857 assert!(swept.is_empty(), "an unknown pid must keep its lock");
4858 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
4859 }
4860
4861 fn run_state_in(status: RunStatus, language: &str) -> RunState {
4862 let mut s = run_state(status);
4863 s.config.graph.language = language.to_owned();
4864 s
4865 }
4866
4867 fn unstarted_verdict(status: RunStatus) -> Verdict {
4868 Verdict {
4869 status,
4870 left_pr: false,
4871 quota_hit: false,
4872 parked: false,
4873 no_viable_candidates: false,
4874 }
4875 }
4876
4877 #[test]
4878 fn an_already_in_base_run_finishes_the_task_without_spending_an_attempt() {
4879 let mut t = task();
4880 t.attempts = 1;
4881 settle_in(
4882 &mut t,
4883 unstarted_verdict(RunStatus::AlreadyInBase),
4884 "already in main",
4885 1,
4886 phrases("en"),
4887 );
4888 assert_eq!(t.status, TaskStatus::Done);
4889 assert_eq!(t.attempts, 0);
4890 }
4891
4892 #[test]
4893 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
4894 let reason = |language: &str| {
4895 let mut t = task();
4896 settle_in(
4897 &mut t,
4898 unstarted_verdict(RunStatus::Judging),
4899 "boom",
4900 1,
4901 phrases(language),
4902 );
4903 t.last_error.or(t.hold_reason).unwrap_or_default()
4904 };
4905 assert!(
4906 reason("en").starts_with("the graph stopped at `"),
4907 "{}",
4908 reason("en")
4909 );
4910 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
4911 assert!(reason("日本語").contains("boom"));
4912 assert_eq!(reason("fr"), reason("en"));
4913 }
4914
4915 #[test]
4916 fn describe_follows_the_run_language_and_keeps_the_run_id() {
4917 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
4918 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
4919 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
4920 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
4921 assert!(ja.contains("[run "), "{ja}");
4922 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
4923 assert!(ended.starts_with("run 終了: "), "{ended}");
4924 let mut de = run_state_in(RunStatus::Failed, "de");
4925 let mut en = run_state_in(RunStatus::Failed, "en");
4926 de.id = "same".to_owned();
4927 en.id = "same".to_owned();
4928 assert_eq!(describe(&de), describe(&en));
4929 }
4930
4931 #[test]
4932 fn refusals_and_recovery_prose_follow_the_language() {
4933 let t = held_task_with("r1");
4934 let q = action_question(
4935 "r1",
4936 ask::ChoiceAction::Resume {
4937 run: "r1".to_owned(),
4938 },
4939 );
4940 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
4941 ActionDecision::Refuse(s) => s,
4942 other => panic!("{other:?}"),
4943 };
4944 assert!(refuse(phrases("en")).contains("could not be read: gone"));
4945 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
4946 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
4947
4948 let mut held = task();
4949 reclaim(&mut held, None, 2, "ja");
4950 assert!(held.hold_reason.unwrap().contains("保留にしました"));
4951 let mut held = task();
4952 reclaim(&mut held, None, 2, "xx");
4953 assert!(held.hold_reason.unwrap().contains("held for a human"));
4954 }
4955
4956 fn run_state(status: RunStatus) -> RunState {
4957 let mut state = RunState::new(
4958 PathBuf::from("/repo"),
4959 "main".to_owned(),
4960 "abc1234def".to_owned(),
4961 "add retries".to_owned(),
4962 Config::default(),
4963 );
4964 state.status = status;
4965 state
4966 }
4967
4968 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
4969 Candidate {
4970 index: 0,
4971 label,
4972 agent: "claude".to_owned(),
4973 branch: format!("magi/x/{label}"),
4974 worktree: PathBuf::from("/repo"),
4975 summary: summary.to_owned(),
4976 stat: String::new(),
4977 files: 0,
4978 commits: usize::from(!empty),
4979 empty,
4980 failed: failed.map(str::to_owned),
4981 verified_noop: None,
4982 duration_ms: 0,
4983 folded: false,
4984 }
4985 }
4986
4987 #[test]
4988 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
4989 let mut state = run_state(RunStatus::Blocked);
4990 state.gate = vec![
4991 CommandOutcome {
4992 command: "cargo make check".to_owned(),
4993 code: Some(0),
4994 output_tail: "ok".to_owned(),
4995 duration_ms: 0,
4996 resource_blocked: false,
4997 },
4998 CommandOutcome {
4999 command: "cargo test".to_owned(),
5000 code: Some(101),
5001 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
5002 duration_ms: 0,
5003 resource_blocked: false,
5004 },
5005 ];
5006 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
5007 assert!(d.contains("cargo test"), "{d}");
5008 assert!(
5009 !d.contains("cargo make check"),
5010 "a passing check is not a diagnostic: {d}"
5011 );
5012 assert!(d.contains("assertion failed"), "{d}");
5013 }
5014
5015 #[test]
5016 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
5017 let mut state = run_state(RunStatus::Blocked);
5018 state.event(
5019 "land",
5020 "stopped: the fixer produced no commit while 2 check(s) were failing \
5021 (build, lint); stopping instead of looping on an unchanged tree",
5022 );
5023 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
5024 assert!(d.contains("build"), "{d}");
5025 assert!(d.contains("lint"), "{d}");
5026 assert!(d.contains("fixer produced no commit"), "{d}");
5027 }
5028
5029 #[test]
5030 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
5031 let state = run_state(RunStatus::VerifiedNoop);
5036 let d = describe(&state);
5037 assert!(
5038 d.contains("agent-verified no-op"),
5039 "expected the display label, not the wire spelling: {d}"
5040 );
5041 assert!(!d.contains("verified_noop"), "{d}");
5042 }
5043
5044 #[test]
5045 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5046 let mut state = run_state(RunStatus::Failed);
5052 state.candidates = vec![candidate(
5053 'A',
5054 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5055 true,
5056 None,
5057 )];
5058 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5059 assert!(d.contains("candidate A"), "{d}");
5060 assert!(d.contains("tagged v1.2.3"), "{d}");
5061 }
5062
5063 #[test]
5064 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5065 let mut state = run_state(RunStatus::Failed);
5066 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5067 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5068 assert!(d.contains("candidate A"), "{d}");
5069 assert!(d.contains("agent timed out"), "{d}");
5070 }
5071
5072 #[test]
5073 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5074 let mut state = run_state(RunStatus::Failed);
5077 state.candidates = vec![candidate('A', "did the work", false, None)];
5078 assert!(diagnostic(&state).is_none());
5079 }
5080
5081 #[test]
5082 fn diagnostic_is_bounded_however_much_a_run_printed() {
5083 let mut state = run_state(RunStatus::Blocked);
5084 state.gate = vec![
5085 CommandOutcome {
5086 command: "cargo test".to_owned(),
5087 code: Some(101),
5088 output_tail: "x".repeat(50_000),
5089 duration_ms: 0,
5090 resource_blocked: false,
5091 },
5092 CommandOutcome {
5093 command: "cargo clippy".to_owned(),
5094 code: Some(1),
5095 output_tail: "y".repeat(50_000),
5096 duration_ms: 0,
5097 resource_blocked: false,
5098 },
5099 ];
5100 state.candidates = vec![
5101 candidate('A', &"z".repeat(50_000), true, None),
5102 candidate('B', &"w".repeat(50_000), true, None),
5103 ];
5104 let d = diagnostic(&state).expect("plenty here to diagnose");
5105 assert!(
5106 d.len() <= DIAGNOSTIC_MAX,
5107 "diagnostic grew to {} bytes, unbounded",
5108 d.len()
5109 );
5110 }
5111
5112 #[test]
5113 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5114 let mut state = run_state(RunStatus::Blocked);
5115 state.gate = vec![CommandOutcome {
5116 command: "cargo test".to_owned(),
5117 code: Some(101),
5118 output_tail: "assertion failed".to_owned(),
5119 duration_ms: 0,
5120 resource_blocked: false,
5121 }];
5122 let verdict = Verdict {
5123 status: RunStatus::Blocked,
5124 left_pr: false,
5125 quota_hit: false,
5126 parked: false,
5127 no_viable_candidates: false,
5128 };
5129
5130 let mut t = task();
5133 t.start("run-1".to_owned());
5134 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5135 assert_eq!(t.status, TaskStatus::Failed);
5136 assert!(t.diagnostic.is_none());
5137
5138 t.start("run-2".to_owned());
5141 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5142 assert_eq!(t.status, TaskStatus::Held);
5143 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5144 assert!(d.contains("cargo test"), "{d}");
5145 }
5146
5147 #[test]
5148 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5149 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5154 let home = crate::run::home();
5155 let state = run_state(RunStatus::VerifiedNoop);
5156 let mut q = ask::Question::new(
5157 state.id.clone(),
5158 "implement".to_owned(),
5159 "impl-A".to_owned(),
5160 "is this really a no-op?".to_owned(),
5161 String::new(),
5162 Vec::new(),
5163 );
5164 Questions::at(home.join("questions")).put(&mut q).unwrap();
5165
5166 let verdict = Verdict {
5167 status: RunStatus::VerifiedNoop,
5168 left_pr: false,
5169 quota_hit: false,
5170 parked: false,
5171 no_viable_candidates: false,
5172 };
5173 let mut t = task();
5174 t.start(state.id.clone());
5175 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5176
5177 assert_eq!(t.status, TaskStatus::Held);
5178 let reason = t.hold_reason.expect("a held task must record why");
5179 assert!(
5180 reason.starts_with("run ended agent-verified no-op"),
5181 "the original settle reason must survive unchanged: {reason}"
5182 );
5183 assert!(
5184 reason.contains(q.short()),
5185 "the open question's id must be named so the notice is actionable: {reason}"
5186 );
5187 }
5188
5189 #[test]
5190 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5191 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5192 let state = run_state(RunStatus::VerifiedNoop);
5193
5194 let verdict = Verdict {
5195 status: RunStatus::VerifiedNoop,
5196 left_pr: false,
5197 quota_hit: false,
5198 parked: false,
5199 no_viable_candidates: false,
5200 };
5201 let mut t = task();
5202 t.start(state.id.clone());
5203 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5204
5205 assert_eq!(t.status, TaskStatus::Held);
5206 assert_eq!(
5207 t.hold_reason.as_deref(),
5208 Some("run ended agent-verified no-op"),
5209 "nothing to append when the question was already answered or never asked"
5210 );
5211 }
5212
5213 #[test]
5214 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5215 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5216 let mut first = run_state(RunStatus::Blocked);
5217 first.id = "20260101-000000-sup1".to_owned();
5218 first.save().unwrap();
5219 let mut second = run_state(RunStatus::Merged);
5220 second.id = "20260101-000000-sup2".to_owned();
5221 second.save().unwrap();
5222
5223 let mut t = task();
5224 t.runs = vec![first.id.clone(), second.id.clone()];
5225 t.status = TaskStatus::Done;
5226
5227 supersede_prior_runs(&t, &crate::run::home());
5228
5229 assert_eq!(
5230 RunState::load(&first.id).unwrap().status,
5231 RunStatus::Superseded,
5232 "the first attempt's Blocked no longer needs anyone's attention"
5233 );
5234 assert_eq!(
5235 RunState::load(&second.id).unwrap().status,
5236 RunStatus::Merged,
5237 "the run that actually succeeded is left exactly as it was"
5238 );
5239 }
5240
5241 #[test]
5242 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5243 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5249 let mut first = run_state(RunStatus::Blocked);
5250 first.id = "20260101-000000-sup9".to_owned();
5251 first.driver_pid = Some(std::process::id());
5254 first.driver_started_at = Some(
5255 crate::proc::process_started_at(std::process::id())
5256 .expect("this test process's own start time must be queryable"),
5257 );
5258 first.save().unwrap();
5259 let mut second = run_state(RunStatus::Merged);
5260 second.id = "20260101-000000-supa".to_owned();
5261 second.save().unwrap();
5262
5263 let mut t = task();
5264 t.runs = vec![first.id.clone(), second.id.clone()];
5265 t.status = TaskStatus::Done;
5266
5267 supersede_prior_runs(&t, &crate::run::home());
5268
5269 assert_eq!(
5270 RunState::load(&first.id).unwrap().status,
5271 RunStatus::Blocked,
5272 "a live driver_pid means something is still actually working this run, \
5273 even though no daemon claims it - rewriting under it would just be \
5274 undone the next time that process saves"
5275 );
5276 }
5277
5278 #[test]
5279 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5280 let dir = tempfile::tempdir().unwrap();
5286 let home = dir.path().to_path_buf();
5287 let queue = Queue::at(dir.path().join("queue"));
5288
5289 let mut first = run_state(RunStatus::Blocked);
5290 first.id = "20260101-000000-supd".to_owned();
5291 first.driver_pid = Some(std::process::id());
5292 first.driver_started_at = Some(
5293 crate::proc::process_started_at(std::process::id())
5294 .expect("this test process's own start time must be queryable"),
5295 );
5296 first.save_under(&home).unwrap();
5297 let mut second = run_state(RunStatus::Merged);
5298 second.id = "20260101-000000-supe".to_owned();
5299 second.save_under(&home).unwrap();
5300
5301 let mut t = task();
5302 t.runs = vec![first.id.clone(), second.id.clone()];
5303 t.status = TaskStatus::Done;
5304 queue.put(&mut t).unwrap();
5305
5306 resweep_superseded_attempts(&queue, &home);
5307 assert_eq!(
5308 RunState::load_under(&first.id, &home).unwrap().status,
5309 RunStatus::Blocked,
5310 "still live on the first pass, so still untouched"
5311 );
5312
5313 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5319 stale.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
5320 stale.save_under(&home).unwrap();
5321
5322 resweep_superseded_attempts(&queue, &home);
5323 assert_eq!(
5324 RunState::load_under(&first.id, &home).unwrap().status,
5325 RunStatus::Superseded,
5326 "the second pass catches up what the first one correctly skipped"
5327 );
5328 }
5329
5330 #[test]
5331 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5332 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5333 let mut first = run_state(RunStatus::Blocked);
5334 first.id = "20260101-000000-sup3".to_owned();
5335 first.save().unwrap();
5336 let mut second = run_state(RunStatus::Blocked);
5337 second.id = "20260101-000000-sup4".to_owned();
5338 second.save().unwrap();
5339
5340 let mut t = task();
5341 t.runs = vec![first.id.clone(), second.id.clone()];
5342 t.status = TaskStatus::Failed;
5346
5347 supersede_prior_runs(&t, &crate::run::home());
5348
5349 assert_eq!(
5350 RunState::load(&first.id).unwrap().status,
5351 RunStatus::Blocked
5352 );
5353 assert_eq!(
5354 RunState::load(&second.id).unwrap().status,
5355 RunStatus::Blocked
5356 );
5357 }
5358
5359 #[test]
5360 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5361 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5365 let mut first = run_state(RunStatus::Blocked);
5366 first.id = "20260101-000000-sup5".to_owned();
5367 first.save().unwrap();
5368
5369 let mut t = task();
5370 t.runs = vec![first.id.clone()];
5371 t.status = TaskStatus::Done;
5372
5373 supersede_prior_runs(&t, &crate::run::home());
5374
5375 assert_eq!(
5376 RunState::load(&first.id).unwrap().status,
5377 RunStatus::Blocked,
5378 "a single-attempt task has no earlier run to supersede"
5379 );
5380 }
5381
5382 #[test]
5383 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5384 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5391 let mut first = run_state(RunStatus::Blocked);
5392 first.id = "20260101-000000-supb".to_owned();
5393 first.save().unwrap();
5394 let mut second = run_state(RunStatus::Failed);
5395 second.id = "20260101-000000-supc".to_owned();
5396 second.save().unwrap();
5397
5398 let mut t = task();
5399 t.runs = vec![first.id.clone(), second.id.clone()];
5400 t.status = TaskStatus::Done;
5401
5402 supersede_prior_runs(&t, &crate::run::home());
5403
5404 assert_eq!(
5405 RunState::load(&first.id).unwrap().status,
5406 RunStatus::Blocked,
5407 "the task's last attempt never landed, so there is nothing here \
5408 actually superseding it"
5409 );
5410 }
5411
5412 #[test]
5413 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5414 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5418 let mut failed = run_state(RunStatus::Failed);
5419 failed.id = "20260101-000000-sup6".to_owned();
5420 failed.save().unwrap();
5421 let mut noop = run_state(RunStatus::VerifiedNoop);
5422 noop.id = "20260101-000000-sup7".to_owned();
5423 noop.save().unwrap();
5424 let mut winner = run_state(RunStatus::Ready);
5425 winner.id = "20260101-000000-sup8".to_owned();
5426 winner.save().unwrap();
5427
5428 let mut t = task();
5429 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5430 t.status = TaskStatus::Done;
5431
5432 supersede_prior_runs(&t, &crate::run::home());
5433
5434 assert_eq!(
5435 RunState::load(&failed.id).unwrap().status,
5436 RunStatus::Failed
5437 );
5438 assert_eq!(
5439 RunState::load(&noop.id).unwrap().status,
5440 RunStatus::VerifiedNoop
5441 );
5442 }
5443
5444 fn approval_question(run: &str) -> ask::Question {
5445 ask::Question::new(
5446 run.to_owned(),
5447 land::APPROVAL_NODE.to_owned(),
5448 "land".to_owned(),
5449 "merge?".to_owned(),
5450 String::new(),
5451 vec!["merge".to_owned(), "hold".to_owned()],
5452 )
5453 }
5454
5455 #[test]
5456 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5457 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5458 let mut state = run_state(RunStatus::Landing);
5459 state.id = "20260101-000000-fre1".to_owned();
5460 state.parked = true;
5461 state.save().unwrap();
5462 ask::Questions::open()
5463 .put(&mut approval_question(&state.id))
5464 .unwrap();
5465
5466 let mut t = task();
5467 t.runs.push(state.id.clone());
5468 assert_eq!(
5469 land_resume_state(&t),
5470 LandResume::StillWaiting,
5471 "nobody has answered and the timeout has not passed"
5472 );
5473 }
5474
5475 #[test]
5476 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5477 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5482 let mut state = run_state(RunStatus::Landing);
5483 state.id = "20260101-000000-exp1".to_owned();
5484 state.parked = true;
5485 state.config.graph.answer_timeout = 60;
5486 state.save().unwrap();
5487
5488 let store = ask::Questions::open();
5489 let mut q = approval_question(&state.id);
5490 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
5491 store.put(&mut q).unwrap();
5492
5493 let mut t = task();
5494 t.runs.push(state.id.clone());
5495 assert_eq!(
5496 land_resume_state(&t),
5497 LandResume::Ready,
5498 "an expired question must not be waited on forever"
5499 );
5500
5501 let after = store.get(&q.id).unwrap();
5502 assert!(
5503 !after.status.open(),
5504 "the question is abandoned, not silently ignored"
5505 );
5506 assert!(
5507 after.resolution().is_none(),
5508 "an abandoned question is not read as a decision"
5509 );
5510 }
5511
5512 #[test]
5513 fn reclaim_settles_a_running_task_against_its_last_run() {
5514 let mut t = task();
5515 t.start("20260904-000000-4043".to_owned());
5516 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
5517 assert_eq!(
5518 t.status,
5519 TaskStatus::Done,
5520 "a run that actually finished must not stay `running` forever"
5521 );
5522 }
5523
5524 #[test]
5525 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
5526 let mut t = task();
5530 t.start("20260904-000000-4043".to_owned());
5531 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
5532 assert_eq!(t.status, TaskStatus::Failed);
5533 assert!(t.status.runnable());
5534 }
5535
5536 #[test]
5537 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
5538 let mut t = task();
5539 t.start("20260904-000000-4043".to_owned());
5540 reclaim(&mut t, None, 2, "en");
5541 assert_eq!(t.status, TaskStatus::Held);
5542 assert!(
5543 t.last_error
5544 .as_deref()
5545 .is_some_and(|e| e.contains("running")),
5546 "the operator needs to know why this task was held"
5547 );
5548 }
5549
5550 #[test]
5551 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
5552 let dir = tempfile::tempdir().unwrap();
5553 let queue = Queue::at(dir.path().to_path_buf());
5554
5555 let mut orphaned = task();
5557 orphaned.id = "20260904-000000-orph".to_owned();
5558 orphaned.status = TaskStatus::Running;
5559 orphaned.attempts = 1;
5560 queue.put(&mut orphaned).unwrap();
5561
5562 let mut alive = task();
5563 alive.id = "20260904-000000-live".to_owned();
5564 alive.status = TaskStatus::Running;
5565 alive.attempts = 1;
5566 queue.put(&mut alive).unwrap();
5567 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
5568
5569 let mut queued = task();
5570 queued.id = "20260904-000000-wait".to_owned();
5571 queue.put(&mut queued).unwrap();
5572
5573 let reclaimed = reclaim_orphaned_running(&queue, 2);
5574 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
5575
5576 assert_eq!(
5577 queue.get(&orphaned.id).unwrap().status,
5578 TaskStatus::Held,
5579 "nothing was driving it and there was no run to recover"
5580 );
5581 assert_eq!(
5582 queue.get(&alive.id).unwrap().status,
5583 TaskStatus::Running,
5584 "a live claim must protect the task it belongs to"
5585 );
5586 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
5587 }
5588
5589 fn read_run_under(home: &Path, id: &str) -> RunState {
5595 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
5596 serde_json::from_str(&body).unwrap()
5597 }
5598
5599 #[test]
5600 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
5601 let dir = tempfile::tempdir().unwrap();
5602 let home = dir.path().to_path_buf();
5603 let now = Timestamp::now();
5604 let overrun_seat = || crate::run::ActiveSeat {
5605 node: "implement".to_owned(),
5606 started_at: now - jiff::SignedDuration::new(21_000, 0),
5607 timeout_secs: 3_600,
5608 attempt: 0,
5609 task: None,
5610 command: None,
5611 index: None,
5612 total: None,
5613 };
5614
5615 let mut dead = run_state(RunStatus::Implementing);
5616 dead.id = "20260101-000000-dead".to_owned();
5617 dead.active.insert("impl-A".to_owned(), overrun_seat());
5618 dead.driver_pid = Some(4242);
5621 dead.save_under(&home).unwrap();
5622
5623 let mut alive = run_state(RunStatus::Implementing);
5626 alive.id = "20260101-000000-aliv".to_owned();
5627 alive.active.insert("impl-A".to_owned(), overrun_seat());
5628 alive.save_under(&home).unwrap();
5629 let mut status = Status::new();
5630 status.current = vec![Current {
5631 task: "20260101-000000-task".to_owned(),
5632 run: alive.id.clone(),
5633 }];
5634 write_status_to(&home.join("daemon.json"), &status).unwrap();
5635
5636 let questions = Questions::at(home.join("questions"));
5640 let mut q = ask::Question::new(
5641 dead.id.clone(),
5642 "implement".to_owned(),
5643 "impl-A".to_owned(),
5644 "Which storage backend?".to_owned(),
5645 String::new(),
5646 vec!["SQLite".to_owned(), "Redis".to_owned()],
5647 );
5648 questions.put(&mut q).unwrap();
5649
5650 let abandoned = reclaim_abandoned_runs_with(
5651 &home,
5652 now,
5653 |pid| if pid == 4242 { Some(false) } else { None },
5654 |_| panic!("a query answering Dead outright needs no identity corroboration"),
5655 );
5656 assert_eq!(abandoned, vec![dead.id.clone()]);
5657
5658 let reloaded = read_run_under(&home, &dead.id);
5659 assert_eq!(reloaded.status, RunStatus::Failed);
5660 assert!(reloaded.active.is_empty());
5661 assert!(
5662 !questions.get(&q.id).unwrap().status.open(),
5663 "the failed run's own open question must be settled in the same pass"
5664 );
5665
5666 let still_alive = read_run_under(&home, &alive.id);
5667 assert_eq!(
5668 still_alive.status,
5669 RunStatus::Implementing,
5670 "a live daemon's claim protects it"
5671 );
5672 assert!(!still_alive.active.is_empty());
5673 }
5674
5675 #[test]
5685 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
5686 let dir = tempfile::tempdir().unwrap();
5687 let home = dir.path().to_path_buf();
5688 let now = Timestamp::now();
5689
5690 let mut manual = run_state(RunStatus::Reviewing);
5691 manual.id = "20260101-000000-manl".to_owned();
5692 manual.active.insert(
5693 "review-1".to_owned(),
5694 crate::run::ActiveSeat {
5695 node: "review".to_owned(),
5696 started_at: now - jiff::SignedDuration::new(21_000, 0),
5697 timeout_secs: 3_600,
5698 attempt: 0,
5699 task: None,
5700 command: None,
5701 index: None,
5702 total: None,
5703 },
5704 );
5705 manual.driver_pid = Some(4242);
5709 manual.driver_started_at = Some("2026-09-22T10:00:00Z".to_owned());
5710 manual.save_under(&home).unwrap();
5711
5712 let abandoned = reclaim_abandoned_runs_with(
5713 &home,
5714 now,
5715 |pid| if pid == 4242 { Some(true) } else { None },
5716 |pid| {
5717 if pid == 4242 {
5718 Some("2026-09-22T10:00:00Z".to_owned())
5719 } else {
5720 None
5721 }
5722 },
5723 );
5724 assert!(
5725 abandoned.is_empty(),
5726 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
5727 );
5728
5729 let reloaded = read_run_under(&home, &manual.id);
5730 assert_eq!(reloaded.status, RunStatus::Reviewing);
5731 assert!(!reloaded.active.is_empty());
5732 }
5733
5734 #[test]
5735 fn an_already_claimed_task_is_skipped_rather_than_failed() {
5736 let dir = tempfile::tempdir().unwrap();
5737 let queue = Queue::at(dir.path().to_path_buf());
5738 let mut only = task();
5739 queue.put(&mut only).unwrap();
5740
5741 let _elsewhere = queue.claim(&only.id).unwrap();
5742 let candidates = runnable(&queue);
5743 assert_eq!(candidates.len(), 1, "the task is still runnable");
5744 assert!(
5745 queue.claim(&candidates[0].id).is_err(),
5746 "the loop cannot take a claim somebody else holds"
5747 );
5748
5749 let after = queue.get(&only.id).unwrap();
5750 assert_eq!(after.status, TaskStatus::Queued);
5751 assert_eq!(
5752 after.attempts, 0,
5753 "losing the race is not an attempt at the task"
5754 );
5755 assert_eq!(after.last_error, None);
5756 }
5757
5758 #[test]
5759 fn the_status_file_round_trips_and_its_heartbeat_advances() {
5760 let dir = tempfile::tempdir().unwrap();
5761 let path = dir.path().join("daemon.json");
5762
5763 let mut status = Status::new();
5764 status.idle = false;
5765 status.completed = 7;
5766 status.current = vec![Current {
5767 task: "20260902-000000-t111".to_owned(),
5768 run: "20260902-000001-r111".to_owned(),
5769 }];
5770 write_status_to(&path, &status).unwrap();
5771 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
5772 assert_eq!(first.schema, SCHEMA);
5773 assert_eq!(first.pid, std::process::id());
5774 assert!(!first.idle);
5775 assert_eq!(first.completed, 7);
5776 assert_eq!(first.current, status.current);
5777 assert!(
5778 !path.with_extension("json.tmp").exists(),
5779 "the temp file is renamed, not left behind"
5780 );
5781
5782 std::thread::sleep(Duration::from_millis(5));
5783 status.updated_at = Timestamp::now();
5784 status.polls = 3;
5785 write_status_to(&path, &status).unwrap();
5786 let second: Status =
5787 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
5788 assert!(
5789 second.updated_at > first.updated_at,
5790 "a reader can only detect staleness if the heartbeat moves"
5791 );
5792 assert_eq!(
5793 second.started_at, first.started_at,
5794 "the start time is not a heartbeat"
5795 );
5796 assert_eq!(second.polls, 3);
5797 }
5798
5799 #[test]
5800 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
5801 let dir = tempfile::tempdir().unwrap();
5802
5803 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
5804
5805 let mut status = Status::new();
5806 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
5807 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5808 let stale = read_status(dir.path()).unwrap();
5809 assert!(
5810 !stale.running(Timestamp::now()),
5811 "a minute without a heartbeat is a dead daemon, not a busy one"
5812 );
5813 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
5814
5815 status.updated_at = Timestamp::now();
5816 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5817 let fresh = read_status(dir.path()).unwrap();
5818 assert!(fresh.running(Timestamp::now()));
5819 }
5820
5821 #[test]
5822 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
5823 let dir = tempfile::tempdir().unwrap();
5824 let now = Timestamp::now();
5825 let mine = "20260903-080619-01c2";
5826
5827 assert!(
5828 !is_working_on(dir.path(), mine, now),
5829 "no status file means nobody is working on anything"
5830 );
5831
5832 let mut status = Status::new();
5833 status.current = vec![Current {
5834 task: "20260903-080340-0167".to_owned(),
5835 run: mine.to_owned(),
5836 }];
5837 status.updated_at = now;
5838 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5839 assert!(is_working_on(dir.path(), mine, now));
5840 assert!(
5841 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
5842 "a daemon busy with one run is not working on another"
5843 );
5844
5845 status.updated_at = now - jiff::SignedDuration::from_secs(600);
5848 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5849 assert!(
5850 !is_working_on(dir.path(), mine, now),
5851 "a stale heartbeat is a dead daemon, so its run is a leftover"
5852 );
5853 }
5854
5855 #[test]
5856 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
5857 let dir = tempfile::tempdir().unwrap();
5858 let now = Timestamp::now();
5859
5860 assert!(
5861 !is_working_on_short(dir.path(), "01c2", now),
5862 "no status file means nobody is working on anything"
5863 );
5864
5865 let mut status = Status::new();
5866 status.current = vec![Current {
5867 task: "20260903-080340-0167".to_owned(),
5868 run: "20260903-080619-01c2".to_owned(),
5869 }];
5870 status.updated_at = now;
5871 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5872 assert!(
5873 is_working_on_short(dir.path(), "01c2", now),
5874 "the run's short id is the last block of its full id"
5875 );
5876 assert!(
5877 !is_working_on_short(dir.path(), "3cbf", now),
5878 "a daemon busy with one worktree bay is not working on another"
5879 );
5880 }
5881
5882 #[test]
5883 fn a_newer_status_file_still_yields_a_reading() {
5884 let dir = tempfile::tempdir().unwrap();
5885 std::fs::write(
5888 dir.path().join("daemon.json"),
5889 serde_json::json!({
5890 "schema": 2,
5891 "updated_at": Timestamp::now().to_string(),
5892 "idle": true,
5893 "surprise": { "nested": [1, 2, 3] },
5894 })
5895 .to_string(),
5896 )
5897 .unwrap();
5898
5899 let reading = read_status(dir.path()).expect("a forward-compatible read");
5900 assert!(reading.running(Timestamp::now()));
5901 assert!(reading.idle);
5902 assert!(reading.current.is_empty());
5903 }
5904
5905 #[test]
5906 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
5907 let dir = tempfile::tempdir().unwrap();
5913 std::fs::write(
5914 dir.path().join("daemon.json"),
5915 serde_json::json!({
5916 "schema": 1,
5917 "pid": 4242,
5918 "updated_at": Timestamp::now().to_string(),
5919 "idle": false,
5920 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
5921 "completed": 3,
5922 "polls": 9,
5923 })
5924 .to_string(),
5925 )
5926 .unwrap();
5927
5928 let reading = read_status(dir.path()).expect("an older shape must still parse");
5929 assert!(reading.running(Timestamp::now()));
5930 assert_eq!(
5931 reading.current,
5932 vec![Current {
5933 task: "20260902-140501-aaaa".to_owned(),
5934 run: "20260902-140502-bbbb".to_owned(),
5935 }]
5936 );
5937 }
5938
5939 #[test]
5940 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
5941 let dir = tempfile::tempdir().unwrap();
5942 std::fs::write(
5943 dir.path().join("daemon.json"),
5944 serde_json::json!({
5945 "schema": 1,
5946 "updated_at": Timestamp::now().to_string(),
5947 "idle": true,
5948 "current": null,
5949 })
5950 .to_string(),
5951 )
5952 .unwrap();
5953 let with_null = read_status(dir.path()).expect("null must still parse");
5954 assert!(with_null.current.is_empty());
5955
5956 std::fs::write(
5957 dir.path().join("daemon.json"),
5958 serde_json::json!({
5959 "schema": 1,
5960 "updated_at": Timestamp::now().to_string(),
5961 "idle": true,
5962 })
5963 .to_string(),
5964 )
5965 .unwrap();
5966 let absent = read_status(dir.path()).expect("a missing field must still parse");
5967 assert!(absent.current.is_empty());
5968 }
5969
5970 #[test]
5971 fn a_task_without_a_repository_runs_in_the_daemons_default() {
5972 let fallback = Path::new("/default");
5973 let mut blank = task();
5974 blank.repo = PathBuf::new();
5975 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
5976 let mut dot = task();
5977 dot.repo = PathBuf::from(".");
5978 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
5979 assert_eq!(
5980 repo_for(&task(), fallback),
5981 PathBuf::from("/repo"),
5982 "a task that names a repository keeps it"
5983 );
5984 }
5985
5986 #[test]
5987 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
5988 let mut solo_cfg = Config::default();
5994 solo_cfg.graph.candidates = 3;
5995 let mut solo_task = task();
5996 solo_task.solo = true;
5997 apply_solo(&mut solo_cfg, &solo_task);
5998 assert_eq!(solo_cfg.graph.candidates, 1);
5999
6000 let mut plain_cfg = Config::default();
6001 plain_cfg.graph.candidates = 3;
6002 let plain_task = task();
6003 assert!(!plain_task.solo);
6004 apply_solo(&mut plain_cfg, &plain_task);
6005 assert_eq!(
6006 plain_cfg.graph.candidates, 3,
6007 "a task that did not ask to run alone keeps the config's candidates"
6008 );
6009 }
6010
6011 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
6012 QuotaLoss {
6013 seat: seat.into(),
6014 node: "judge".into(),
6015 at: at.parse().unwrap(),
6016 reset: reset.map(str::to_string),
6017 }
6018 }
6019
6020 #[test]
6021 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
6022 let old: Vec<QuotaLoss> = (1..=4)
6023 .map(|i| {
6024 loss(
6025 &format!("judge-{i}"),
6026 "2026-09-23T05:23:00Z",
6027 Some("2:40pm (Asia/Tokyo)"),
6028 )
6029 })
6030 .collect();
6031 let fresh = losses_this_attempt(&old, &old);
6032 assert!(fresh.is_empty());
6033 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
6034 }
6036
6037 #[test]
6038 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6039 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6040 let now = Timestamp::now();
6041 let mut after = old.clone();
6042 after.push(loss("judge-2", &now.to_string(), None));
6043 let fresh = losses_this_attempt(&old, &after);
6044 assert_eq!(fresh, vec![after[1].clone()]);
6045 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6046 assert_eq!(
6047 until,
6048 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6049 );
6050 }
6051
6052 #[test]
6053 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6054 let before = vec![
6057 loss("judge-1", "2026-09-23T05:23:00Z", None),
6058 loss("judge-2", "2026-09-23T05:24:00Z", None),
6059 ];
6060 let after = vec![
6061 loss("judge-2", "2026-09-23T05:24:00Z", None),
6062 loss("judge-1", "2026-09-24T01:00:00Z", None),
6063 ];
6064 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6065 }
6066
6067 #[test]
6068 fn merge_overrides_are_parsed_or_refused() {
6069 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6070 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6071 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6072 assert!(merge_mode("squash").is_err());
6073 }
6074
6075 #[test]
6076 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6077 let now = Timestamp::now();
6078 let fallback = Duration::from_secs(300);
6079 let cap = Duration::from_secs(1800);
6080
6081 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6083
6084 let soon = now + jiff::SignedDuration::from_secs(600);
6086 assert_eq!(
6087 quota_wait(Some(soon), now, fallback, cap),
6088 Duration::from_secs(600)
6089 );
6090
6091 let past = now - jiff::SignedDuration::from_secs(60);
6094 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6095
6096 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6099 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6100 }
6101
6102 #[test]
6103 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6104 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6105
6106 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6107 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6108
6109 let already_past =
6113 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6114 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6115
6116 assert!(
6117 parse_reset_hint("session limit reached", now, now).is_none(),
6118 "free text with no recognised shape is not guessed at"
6119 );
6120 assert!(
6121 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6122 "an unresolvable zone name is not guessed at either"
6123 );
6124 }
6125
6126 #[test]
6127 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6128 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6129
6130 let at = parse_reset_hint(
6131 "You've hit your usage limit. Visit \
6132 https://chatgpt.com/codex/settings/usage to purchase more \
6133 credits or try again at Sep 19th, 2026 5:10 PM.",
6134 now,
6135 now,
6136 )
6137 .expect("the codex reset wording is a recognised shape");
6138 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6139
6140 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6145 .expect("an explicit year needs no rollover");
6146 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6147
6148 assert!(
6149 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6150 "a two-digit year is not the documented shape and is not guessed at"
6151 );
6152 assert!(
6153 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6154 "a four-letter month name is not the documented three-letter abbreviation"
6155 );
6156 assert!(
6157 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6158 "an explicit zone on the dated shape is a format nobody has \
6159 documented, and is refused rather than guessed at as UTC"
6160 );
6161 }
6162
6163 #[test]
6164 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6165 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6166 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6167
6168 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6169 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6170
6171 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6172 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6173
6174 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6175 assert!(
6176 parse_reset_hint(bad, now, recorded).is_none(),
6177 "{bad:?} must not be guessed at"
6178 );
6179 }
6180 }
6181
6182 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6186 let config = dir.join("magi.toml");
6187 std::fs::write(
6188 &config,
6189 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6190 )
6191 .unwrap();
6192 let opts = Opts {
6193 poll: Duration::from_secs(30),
6194 config: Some(config),
6195 repo: dir.join("repo"),
6199 ..Opts::default()
6200 };
6201 let home = dir.join("home");
6210 let worktrees = dir.join("wt");
6211 (
6212 opts,
6213 Queue::at(dir.join("queue")),
6214 home.join("daemon.json"),
6215 home,
6216 worktrees,
6217 )
6218 }
6219
6220 #[test]
6221 fn a_stop_is_idempotent_and_once_set_stays_set() {
6222 let stop = Stop::new();
6223 assert!(!stop.stopped());
6224
6225 stop.stop();
6226 assert!(stop.stopped());
6227 stop.stop();
6228 assert!(stop.stopped(), "a second stop is not a toggle");
6229
6230 let shared = stop.clone();
6231 assert!(
6232 shared.stopped(),
6233 "a clone is the same stop; that is how the loop and its caller share one"
6234 );
6235 }
6236
6237 #[test]
6238 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6239 let stop = Stop::new();
6240 stop.enter();
6241 assert!(
6242 !stop.finishing(),
6243 "a busy loop nobody has asked to stop is just running"
6244 );
6245
6246 stop.stop();
6247 assert!(
6248 stop.finishing(),
6249 "a stop asked for mid-run has not landed until the run is settled"
6250 );
6251
6252 stop.exit();
6253 assert!(
6254 !stop.finishing(),
6255 "once the run is settled the stop has landed and there is nothing to finish"
6256 );
6257 }
6258
6259 #[test]
6260 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6261 let stop = Stop::new();
6262 stop.enter();
6263 stop.enter();
6264 stop.stop();
6265 assert!(stop.finishing(), "two runs still in flight");
6266
6267 stop.exit();
6268 assert!(
6269 stop.finishing(),
6270 "one run finished, but a sibling is still working"
6271 );
6272
6273 stop.exit();
6274 assert!(
6275 !stop.finishing(),
6276 "the last run out is what actually lands the stop"
6277 );
6278 }
6279
6280 #[tokio::test]
6281 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6282 let dir = tempfile::tempdir().unwrap();
6283 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6284 let stop = Stop::new();
6285 stop.stop();
6286
6287 let began = std::time::Instant::now();
6288 tokio::time::timeout(
6289 Duration::from_secs(2),
6290 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6291 )
6292 .await
6293 .expect("a stopped loop must return, not sit out its poll interval")
6294 .expect("the loop's own setup and teardown must not fail");
6295 assert!(
6296 began.elapsed() < opts.poll,
6297 "returned only after {:?}, which is a poll interval, not a stop",
6298 began.elapsed()
6299 );
6300 }
6301
6302 #[tokio::test]
6303 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6304 let dir = tempfile::tempdir().unwrap();
6305 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6306 let stop = Stop::new();
6307
6308 let asker = {
6311 let stop = stop.clone();
6312 tokio::spawn(async move {
6313 tokio::time::sleep(Duration::from_millis(20)).await;
6314 stop.stop();
6315 })
6316 };
6317
6318 let began = std::time::Instant::now();
6319 tokio::time::timeout(
6320 Duration::from_secs(2),
6321 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6322 )
6323 .await
6324 .expect("a stop asked for while idle must wake the wait")
6325 .expect("the loop's own setup and teardown must not fail");
6326 asker.await.unwrap();
6327 assert!(
6328 began.elapsed() < opts.poll,
6329 "returned only after {:?}, so the stop waited on the sleep",
6330 began.elapsed()
6331 );
6332 }
6333
6334 #[tokio::test]
6335 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6336 let dir = tempfile::tempdir().unwrap();
6337 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6338 let stop = Stop::new();
6339 stop.stop();
6340
6341 tokio::time::timeout(
6342 Duration::from_secs(2),
6343 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6344 )
6345 .await
6346 .expect("a stopped loop must return")
6347 .expect("the loop's own setup and teardown must not fail");
6348
6349 assert!(
6350 home.is_dir(),
6351 "the loop did publish a status file, so its removal is the teardown and not an absence"
6352 );
6353 assert!(
6354 !status_file.exists(),
6355 "a stopped loop clears its status file"
6356 );
6357 assert!(
6358 read_status(&home).is_none(),
6359 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6360 );
6361 }
6362
6363 #[tokio::test]
6364 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6365 let dir = tempfile::tempdir().unwrap();
6366 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6367 opts.once = true;
6368
6369 let mut settled = RunState::new(
6370 dir.path().join("repo"),
6371 "main".to_owned(),
6372 "abc1234".to_owned(),
6373 "fixture".to_owned(),
6374 Config::default(),
6375 );
6376 settled.status = RunStatus::Ready;
6377 let run_dir = home.join("runs").join(&settled.id);
6378 std::fs::create_dir_all(&run_dir).unwrap();
6379 std::fs::write(
6380 run_dir.join("run.json"),
6381 serde_json::to_string_pretty(&settled).unwrap(),
6382 )
6383 .unwrap();
6384 let questions = Questions::at(home.join("questions"));
6385 let mut question = ask::Question::new(
6386 settled.id.clone(),
6387 "review".to_owned(),
6388 "reviewer-1".to_owned(),
6389 "Continue?".to_owned(),
6390 String::new(),
6391 Vec::new(),
6392 );
6393 questions.put(&mut question).unwrap();
6394
6395 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6396 .await
6397 .unwrap();
6398
6399 assert_eq!(
6400 questions.get(&question.id).unwrap().status,
6401 ask::QuestionStatus::Abandoned,
6402 "an empty --once drain still performs startup question cleanup"
6403 );
6404 }
6405
6406 #[test]
6407 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6408 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6409
6410 assert!(
6411 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6412 "never checked before: due at once"
6413 );
6414
6415 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6416 assert!(
6417 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6418 "well inside the interval: not due yet"
6419 );
6420
6421 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6422 assert!(
6423 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6424 "exactly at the edge: not yet due, same convention as `clean::due`"
6425 );
6426
6427 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6428 assert!(
6429 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6430 "past the interval: due again"
6431 );
6432 }
6433
6434 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
6439 let config = dir.join("magi.toml");
6440 std::fs::write(
6446 &config,
6447 format!(
6448 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
6449 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
6450 cache_dir.display()
6451 ),
6452 )
6453 .unwrap();
6454 Opts {
6455 config: Some(config),
6456 repo: dir.join("repo"),
6457 ..Opts::default()
6458 }
6459 }
6460
6461 #[tokio::test]
6462 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
6463 let dir = tempfile::tempdir().unwrap();
6464 let home = dir.path().join("home");
6465 let cache_dir = dir.path().join("cache");
6466 std::fs::create_dir_all(&cache_dir).unwrap();
6467 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6468 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6469
6470 let running = Stop::new();
6473 let mut last_checked = None;
6474 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6475 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
6476 .await;
6477 assert_eq!(
6478 crate::disk::dir_size(&cache_dir),
6479 0,
6480 "over the cap on the first check ever: pruned at once, no idle queue required"
6481 );
6482 assert_eq!(last_checked, Some(t0));
6483
6484 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
6486 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
6487 maybe_prune_cache_between_runs(
6488 &opts.repo,
6489 &opts,
6490 &home,
6491 &running,
6492 &mut last_checked,
6493 too_soon,
6494 )
6495 .await;
6496 assert_eq!(
6497 crate::disk::dir_size(&cache_dir),
6498 10,
6499 "too soon since the last check: left alone rather than rescanned every call"
6500 );
6501 assert_eq!(
6502 last_checked,
6503 Some(t0),
6504 "an idle check does not reset the clock"
6505 );
6506
6507 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6509 maybe_prune_cache_between_runs(
6510 &opts.repo,
6511 &opts,
6512 &home,
6513 &running,
6514 &mut last_checked,
6515 due_again,
6516 )
6517 .await;
6518 assert_eq!(
6519 crate::disk::dir_size(&cache_dir),
6520 0,
6521 "due again: pruned back under the cap"
6522 );
6523 }
6524
6525 #[tokio::test]
6533 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
6534 let dir = tempfile::tempdir().unwrap();
6535 let home = dir.path().join("home");
6536 let cache_dir = dir.path().join("cache");
6537 std::fs::create_dir_all(&cache_dir).unwrap();
6538 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6539 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6540
6541 let stop = Stop::new();
6542 stop.stop();
6543 assert!(
6544 !stop.finishing(),
6545 "no run is in flight at a between-runs boundary, so nothing else \
6546 would tell the operator this stop had not taken effect yet"
6547 );
6548
6549 let mut last_checked = None;
6550 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6551 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
6552 .await;
6553 assert_eq!(
6554 crate::disk::dir_size(&cache_dir),
6555 10,
6556 "over its cap, and due for the first check ever, but a stop outranks \
6557 it: the cap is a standing policy the next start measures again"
6558 );
6559 assert_eq!(
6560 last_checked, None,
6561 "a check that never happened must not claim the interval"
6562 );
6563 }
6564
6565 #[tokio::test]
6579 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
6580 let dir = tempfile::tempdir().unwrap();
6581 let cache_dir = dir.path().join("cache");
6582 std::fs::create_dir_all(&cache_dir).unwrap();
6583 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
6584
6585 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
6586 opts.poll = Duration::from_millis(20);
6587 opts.max_attempts = 1_000;
6588
6589 let queue = Queue::at(dir.path().join("queue"));
6590 let mut t = Task::new(
6591 "x".to_owned(),
6592 "x".to_owned(),
6593 opts.repo.clone(),
6594 Source::Human,
6595 );
6596 queue.put(&mut t).unwrap();
6597
6598 let home = dir.path().join("home");
6599 let worktrees = dir.path().join("wt");
6600 let status_file = home.join("daemon.json");
6601 let stop = Stop::new();
6602 let stopper = {
6603 let stop = stop.clone();
6604 tokio::spawn(async move {
6605 tokio::time::sleep(Duration::from_millis(400)).await;
6606 stop.stop();
6607 })
6608 };
6609
6610 tokio::time::timeout(
6611 Duration::from_secs(10),
6612 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6613 )
6614 .await
6615 .expect("the loop must not hang on a queue that keeps producing failing work")
6616 .expect("the loop's own setup and teardown must not fail");
6617 stopper.await.unwrap();
6618
6619 let after = queue.get(&t.id).unwrap();
6620 assert!(
6621 after.attempts >= 2,
6622 "the harness must actually have retried more than once, or this is not \
6623 exercising a busy queue at all (got {} attempt(s))",
6624 after.attempts
6625 );
6626 assert!(
6627 after.status.runnable(),
6628 "still under its attempt budget: the queue never reached a natural idle \
6629 on its own, only the external stop ended the test"
6630 );
6631
6632 assert_eq!(
6633 crate::disk::dir_size(&cache_dir),
6634 0,
6635 "an oversized cache must not be left to grow unboundedly just because the \
6636 queue kept the loop busy the whole time"
6637 );
6638 }
6639
6640 #[test]
6641 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
6642 let dir = tempfile::tempdir().unwrap();
6643 let queue = Queue::at(dir.path().join("queue"));
6644 let questions = Questions::at(dir.path().join("questions"));
6645 let mut task = task();
6646 queue.put(&mut task).unwrap();
6647
6648 let mut task_question = ask::Question::new(
6649 task.id.clone(),
6650 crate::conduct::NODE.to_owned(),
6651 "conduct".to_owned(),
6652 "Which backend?".to_owned(),
6653 String::new(),
6654 Vec::new(),
6655 );
6656 questions.put(&mut task_question).unwrap();
6657 task.block(vec![task_question.id.clone()], None);
6658 queue.put(&mut task).unwrap();
6659
6660 let mut run_question = ask::Question::new(
6661 "20260101-000000-run1".to_owned(),
6662 "review".to_owned(),
6663 "reviewer-1".to_owned(),
6664 "Run question".to_owned(),
6665 String::new(),
6666 Vec::new(),
6667 );
6668 questions.put(&mut run_question).unwrap();
6669
6670 let mut coincidental = ask::Question::new(
6675 task.id.clone(),
6676 "review".to_owned(),
6677 "reviewer-1".to_owned(),
6678 "Unrelated review question".to_owned(),
6679 String::new(),
6680 Vec::new(),
6681 );
6682 questions.put(&mut coincidental).unwrap();
6683
6684 reconcile_task_questions(&queue, &questions);
6685 assert!(questions.get(&task_question.id).unwrap().status.open());
6686 assert!(questions.get(&run_question.id).unwrap().status.open());
6687 assert!(questions.get(&coincidental.id).unwrap().status.open());
6688
6689 task.release();
6690 queue.put(&mut task).unwrap();
6691 reconcile_task_questions(&queue, &questions);
6692 assert_eq!(
6693 questions.get(&task_question.id).unwrap().status,
6694 ask::QuestionStatus::Abandoned
6695 );
6696 assert!(
6697 questions.get(&run_question.id).unwrap().status.open(),
6698 "run questions remain the run janitor's responsibility"
6699 );
6700 assert!(
6701 questions.get(&coincidental.id).unwrap().status.open(),
6702 "a non-conductor question must not be abandoned just because its \
6703 run id coincides with a task id"
6704 );
6705 }
6706
6707 #[test]
6708 fn a_freshly_started_running_task_is_never_stalled() {
6709 let dir = tempfile::tempdir().unwrap();
6710 let mut t = task();
6711 t.start("run-1".to_owned());
6712 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
6715 }
6716
6717 #[test]
6718 fn a_long_running_task_with_no_live_daemon_is_stalled() {
6719 let dir = tempfile::tempdir().unwrap();
6720 let mut t = task();
6721 t.start("run-1".to_owned());
6722 t.updated_at = Timestamp::now()
6723 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6724 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
6725 assert_eq!(
6726 stalled_tasks(
6727 &Queue::at(dir.path().join("q")),
6728 dir.path(),
6729 Timestamp::now()
6730 )
6731 .len(),
6732 0,
6733 "the task was never written to this queue"
6734 );
6735 }
6736
6737 #[test]
6738 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
6739 let dir = tempfile::tempdir().unwrap();
6740 let mut t = task();
6741 t.id = "20260903-080340-0167".to_owned();
6742 t.start("20260903-080619-01c2".to_owned());
6743 t.updated_at = Timestamp::now()
6744 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6745
6746 let mut status = Status::new();
6747 status.current = vec![Current {
6748 task: t.id.clone(),
6749 run: "20260903-080619-01c2".to_owned(),
6750 }];
6751 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6752
6753 assert!(
6754 !is_stalled(&t, dir.path(), Timestamp::now()),
6755 "a live daemon's own heartbeat rules out stalled, however long the task has run"
6756 );
6757 }
6758
6759 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
6763 let path = queue.path_of(id);
6764 let body = std::fs::read_to_string(&path).unwrap();
6765 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
6766 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
6767 v["updated_at"] = serde_json::Value::String(old.to_string());
6768 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
6769 }
6770
6771 #[test]
6772 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
6773 let dir = tempfile::tempdir().unwrap();
6786 let queue = Queue::at(dir.path().join("queue"));
6787 let home = dir.path().join("home");
6788
6789 let mut t = task();
6790 t.id = "20260101-000001-lock".to_owned();
6791 t.start("run-1".to_owned());
6792 queue.put(&mut t).unwrap();
6793 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
6794 std::fs::write(
6795 dir.path().join("queue").join(format!("{}.lock", t.id)),
6796 "not a pid",
6797 )
6798 .unwrap();
6799
6800 let now = Timestamp::now();
6801 assert!(
6802 reclaim_orphaned_running(&queue, 2).is_empty(),
6803 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
6804 and reclaim must leave the task alone"
6805 );
6806 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
6807
6808 let stalled = stalled_tasks(&queue, &home, now);
6809 assert_eq!(
6810 stalled.len(),
6811 1,
6812 "reclaim's inability to claim it yet must not hide it from the conductor"
6813 );
6814 assert_eq!(stalled[0].id, t.id);
6815 }
6816
6817 #[test]
6818 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
6819 let dir = tempfile::tempdir().unwrap();
6820 crate::run::set_home(dir.path().join("run-home"));
6821 let queue = Queue::at(dir.path().join("queue"));
6822 let home = dir.path().join("home");
6823 let questions = Questions::at(dir.path().join("questions"));
6824
6825 let mut t = task();
6826 t.id = "20260101-000003-dead".to_owned();
6827 t.start("missing-run".to_owned());
6828 queue.put(&mut t).unwrap();
6829 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
6830
6831 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
6834 assert_eq!(
6835 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
6836 [&t.id]
6837 );
6838 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
6839 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
6840
6841 crate::conduct::apply(
6844 &queue,
6845 &questions,
6846 &crate::conduct::Verdict {
6847 decisions: vec![crate::conduct::Decision {
6848 id: t.id.clone(),
6849 recovery: Some(crate::conduct::Recovery::Requeue),
6850 ..crate::conduct::Decision::default()
6851 }],
6852 },
6853 )
6854 .unwrap();
6855 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
6856 }
6857
6858 #[test]
6859 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
6860 let dir = tempfile::tempdir().unwrap();
6861 let queue = Queue::at(dir.path().join("queue"));
6862 let home = dir.path().join("home");
6863
6864 let mut fresh = task();
6865 fresh.id = "20260101-000001-aaaa".to_owned();
6866 fresh.start("run-1".to_owned());
6867 queue.put(&mut fresh).unwrap();
6868
6869 let mut old = task();
6870 old.id = "20260101-000002-bbbb".to_owned();
6871 old.start("run-2".to_owned());
6872 queue.put(&mut old).unwrap();
6873 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
6874
6875 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
6876 assert_eq!(stalled.len(), 1);
6877 assert_eq!(stalled[0].id, old.id);
6878 }
6879
6880 #[test]
6881 fn queued_and_finished_task_views_partition_by_status() {
6882 let dir = tempfile::tempdir().unwrap();
6883 let queue = Queue::at(dir.path().join("queue"));
6884
6885 let mut queued = task();
6886 queued.id = "20260101-000001-aaaa".to_owned();
6887 queue.put(&mut queued).unwrap();
6888
6889 let mut failed = task();
6890 failed.id = "20260101-000002-bbbb".to_owned();
6891 failed.start("run-1".to_owned());
6892 failed.fail("gate red", 5);
6893 queue.put(&mut failed).unwrap();
6894
6895 let mut held = task();
6896 held.id = "20260101-000003-cccc".to_owned();
6897 held.hold_machine(None);
6898 queue.put(&mut held).unwrap();
6899
6900 let mut running = task();
6901 running.id = "20260101-000004-dddd".to_owned();
6902 running.start("run-2".to_owned());
6903 queue.put(&mut running).unwrap();
6904
6905 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
6906 assert_eq!(queued_ids, [queued.id.clone()]);
6907
6908 let mut finished_ids: Vec<String> =
6909 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
6910 finished_ids.sort_unstable();
6911 let mut want = vec![failed.id.clone(), held.id.clone()];
6912 want.sort_unstable();
6913 assert_eq!(finished_ids, want);
6914 }
6915
6916 #[test]
6917 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
6918 let dir = tempfile::tempdir().unwrap();
6919 let queue = Queue::at(dir.path().join("queue"));
6920 let questions = ask::Questions::at(dir.path().join("questions"));
6921
6922 let mut dep = task();
6923 dep.id = "20260101-000001-dep0".to_owned();
6924 dep.succeed();
6925 queue.put(&mut dep).unwrap();
6926
6927 let mut still_going = task();
6928 still_going.id = "20260101-000002-dep1".to_owned();
6929 queue.put(&mut still_going).unwrap();
6930
6931 let mut blocked = task();
6932 blocked.id = "20260101-000003-main".to_owned();
6933 blocked.block(
6934 vec![dep.id.clone(), still_going.id.clone()],
6935 Some("waits on both".to_owned()),
6936 );
6937 queue.put(&mut blocked).unwrap();
6938
6939 resolve_blockers(&queue, &questions);
6940
6941 let after = queue.get(&blocked.id).unwrap();
6942 assert_eq!(
6943 after.status,
6944 TaskStatus::Blocked,
6945 "one dependency is still outstanding"
6946 );
6947 assert_eq!(after.blocked_by, [still_going.id.clone()]);
6948 }
6949
6950 #[test]
6951 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
6952 let dir = tempfile::tempdir().unwrap();
6953 let queue = Queue::at(dir.path().join("queue"));
6954 let questions = ask::Questions::at(dir.path().join("questions"));
6955
6956 let mut q = crate::ask::Question::new(
6957 "20260101-000001-main".to_owned(),
6958 crate::conduct::NODE.to_owned(),
6959 "conduct".to_owned(),
6960 "Which backend?".to_owned(),
6961 String::new(),
6962 Vec::new(),
6963 );
6964 questions.put(&mut q).unwrap();
6965 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
6966 .unwrap();
6967 questions.put(&mut q).unwrap();
6968
6969 let mut blocked = task();
6970 blocked.id = "20260101-000001-main".to_owned();
6971 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
6972 queue.put(&mut blocked).unwrap();
6973
6974 resolve_blockers(&queue, &questions);
6975
6976 let after = queue.get(&blocked.id).unwrap();
6977 assert_eq!(
6978 after.status,
6979 TaskStatus::Queued,
6980 "the only blocker resolved"
6981 );
6982 assert_eq!(after.answers.len(), 1);
6983 assert_eq!(after.answers[0].question, "Which backend?");
6984 assert_eq!(after.answers[0].answer, "SQLite");
6985
6986 let instruction = instruction_for(&after);
6988 assert!(instruction.contains("Which backend?"));
6989 assert!(instruction.contains("SQLite"));
6990 }
6991
6992 #[test]
6993 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
6994 let dir = tempfile::tempdir().unwrap();
7000 let queue = Queue::at(dir.path().join("queue"));
7001 let questions = ask::Questions::at(dir.path().join("questions"));
7002
7003 let mut q = crate::ask::Question::new(
7004 "20260101-000001-main".to_owned(),
7005 crate::conduct::NODE.to_owned(),
7006 "conduct".to_owned(),
7007 "How should this be handled?".to_owned(),
7008 String::new(),
7009 Vec::new(),
7010 );
7011 questions.put(&mut q).unwrap();
7012 q.answer(crate::ask::Answer::Text(
7013 "leave it held, a human will look at it later".to_owned(),
7014 ))
7015 .unwrap();
7016 questions.put(&mut q).unwrap();
7017
7018 let mut held = task();
7019 held.id = "20260101-000001-main".to_owned();
7020 held.hold_machine(Some("out of attempts".to_owned()));
7021 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
7022 queue.put(&mut held).unwrap();
7023
7024 resolve_blockers(&queue, &questions);
7025
7026 let after = queue.get(&held.id).unwrap();
7027 assert_eq!(after.status, TaskStatus::Held);
7028 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
7029 assert_eq!(
7030 after.answers[0].answer,
7031 "leave it held, a human will look at it later"
7032 );
7033 }
7034
7035 #[test]
7036 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
7037 let dir = tempfile::tempdir().unwrap();
7043 let queue = Queue::at(dir.path().join("queue"));
7044 let questions = ask::Questions::at(dir.path().join("questions"));
7045
7046 let mut still_going = task();
7047 still_going.id = "20260101-000002-dep1".to_owned();
7048 queue.put(&mut still_going).unwrap();
7049
7050 let mut blocked = task();
7051 blocked.id = "20260101-000003-main".to_owned();
7052 blocked.block(
7053 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7054 Some("waits on both".to_owned()),
7055 );
7056 queue.put(&mut blocked).unwrap();
7057
7058 resolve_blockers(&queue, &questions);
7059
7060 let after = queue.get(&blocked.id).unwrap();
7061 assert_eq!(
7062 after.status,
7063 TaskStatus::Held,
7064 "a missing dependency must not leave the task blocked forever"
7065 );
7066 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7067 assert!(after.blocked_by.is_empty());
7068 let reason = after.hold_reason.as_deref().unwrap_or_default();
7069 assert!(
7070 reason.contains("20260101-000001-gone"),
7071 "the missing id must be named so an operator can tell what happened: {reason}"
7072 );
7073 assert!(
7074 reason.contains(&still_going.id),
7075 "the still-valid dependency must not silently vanish from the record: {reason}"
7076 );
7077 }
7078
7079 #[test]
7080 fn instruction_for_is_unchanged_without_any_answers() {
7081 let t = task();
7082 assert_eq!(instruction_for(&t), t.instruction);
7083 }
7084
7085 #[test]
7086 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7087 let dir = tempfile::tempdir().unwrap();
7088 let q = Queue::at(dir.path().join("queue"));
7089 let src = dir.path().join("shot.png");
7090 std::fs::write(&src, "x").unwrap();
7091 let mut t = task();
7092 q.attach(&mut t, &[src]).unwrap();
7093 let paths = task_attachments(&q, &t).unwrap();
7094 assert_eq!(paths.len(), 1);
7095 assert!(paths[0].is_absolute() && paths[0].is_file());
7096 std::fs::remove_file(&paths[0]).unwrap();
7097 let err = task_attachments(&q, &t).unwrap_err().to_string();
7098 assert!(err.contains("shot.png"), "{err}");
7099 }
7100
7101 #[test]
7102 fn resumed_instruction_is_unchanged_without_any_answers() {
7103 let t = task();
7104 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7105 }
7106
7107 #[test]
7108 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7109 let mut t = task();
7110 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7111 let old = t.instruction.clone();
7115
7116 let refreshed = resumed_instruction(&old, &t);
7117 assert!(refreshed.starts_with(&old), "the original text is kept");
7118 assert!(refreshed.contains("Which backend?"));
7119 assert!(refreshed.contains("SQLite"));
7120 }
7121
7122 #[test]
7123 fn resumed_instruction_keeps_an_original_answers_heading() {
7124 let mut t = task();
7125 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7126 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7127
7128 let refreshed = resumed_instruction(&t.instruction, &t);
7129
7130 assert!(
7131 refreshed.starts_with(&t.instruction),
7132 "an answers heading in the original instruction is not the appended block"
7133 );
7134 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7135 assert!(refreshed.contains("Which backend?"));
7136 assert!(refreshed.contains("SQLite"));
7137
7138 let repeated = resumed_instruction(&refreshed, &t);
7139 assert_eq!(
7140 repeated, refreshed,
7141 "only the final appended block is refreshed"
7142 );
7143 }
7144
7145 #[test]
7146 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7147 let mut t = task();
7148 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7149
7150 let once = resumed_instruction(&t.instruction, &t);
7154 let twice = resumed_instruction(&once, &t);
7155 assert_eq!(once, twice);
7156 assert_eq!(once.matches("Which backend?").count(), 1);
7157
7158 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7160 let refreshed = resumed_instruction(&once, &t);
7161 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7162 assert!(refreshed.contains("Which backend?"));
7163 assert!(refreshed.contains("Which cache?"));
7164 }
7165
7166 #[test]
7167 fn prepare_instruction_covers_all_three_starters() {
7168 let mut t = task();
7169 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7170
7171 assert_eq!(
7174 prepare_instruction(&Starter::Start, None, &t),
7175 Some(instruction_for(&t))
7176 );
7177
7178 let old = t.instruction.clone();
7181 assert_eq!(
7182 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7183 Some(resumed_instruction(&old, &t))
7184 );
7185
7186 assert_eq!(
7190 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7191 None
7192 );
7193 }
7194
7195 #[test]
7196 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7197 assert_eq!(
7198 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7199 Starter::Review("magi/eba2/A".to_owned())
7200 );
7201 }
7202
7203 #[test]
7204 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7205 assert_eq!(
7206 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7207 Starter::Start,
7208 "a vanished review branch must not fall back to resuming the old run either"
7209 );
7210 }
7211
7212 #[test]
7213 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7214 assert_eq!(
7215 choose_starter(None, false, Some("some-run")),
7216 Starter::Resume("some-run".to_owned())
7217 );
7218 assert_eq!(choose_starter(None, false, None), Starter::Start);
7219 }
7220
7221 #[test]
7222 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7223 let mut released = task();
7224 released.start("stalled-run".to_owned());
7225 released.requeue();
7226 let unfinished = (!released.fresh_start)
7227 .then(|| Some("stalled-run".to_owned()))
7228 .flatten();
7229 assert_eq!(
7230 choose_starter(None, false, unfinished.as_deref()),
7231 Starter::Start,
7232 "release keeps run history but must not resume it"
7233 );
7234 assert_eq!(released.runs, ["stalled-run"]);
7235 }
7236
7237 #[test]
7238 fn an_ordinary_release_keeps_a_resumable_run_available() {
7239 let mut released = task();
7240 released.start("stalled-run".to_owned());
7241 released.release();
7242 let unfinished = (!released.fresh_start)
7243 .then(|| Some("stalled-run".to_owned()))
7244 .flatten();
7245 assert_eq!(
7246 choose_starter(None, false, unfinished.as_deref()),
7247 Starter::Resume("stalled-run".to_owned()),
7248 "manual release must preserve the normal resume path"
7249 );
7250 }
7251
7252 #[test]
7253 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7254 let mut state = run_state(RunStatus::Blocked);
7255 state.config.graph.review_rounds = 3;
7256 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7257 assert!(exhausted_review_budget(&state));
7258
7259 state.reviews.pop();
7261 assert!(!exhausted_review_budget(&state));
7262
7263 let mut stalled = run_state(RunStatus::Stalled);
7266 stalled.config.graph.review_rounds = 1;
7267 stalled.reviews = vec![review_round(1)];
7268 assert!(!exhausted_review_budget(&stalled));
7269 }
7270
7271 fn review_round(round: usize) -> crate::run::ReviewRound {
7272 crate::run::ReviewRound {
7273 round,
7274 head: "deadbeef".to_owned(),
7275 verified_head: None,
7276 verified_at: None,
7277 reviews: Vec::new(),
7278 e2e: Vec::new(),
7279 verify_retried: false,
7280 e2e_deferred: false,
7281 e2e_defer_reason: None,
7282 fix: None,
7283 blocking: 0,
7284 answered: 1,
7285 expected: 1,
7286 clean: false,
7287 progressed: true,
7288 vote_split: false,
7289 reconsideration: Vec::new(),
7290 verdict: None,
7291 }
7292 }
7293
7294 #[test]
7295 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7296 let mut run = RunState::new(
7297 PathBuf::from("/repo"),
7298 "main".to_owned(),
7299 "abc1234def".to_owned(),
7300 "add retries".to_owned(),
7301 Config::default(),
7302 );
7303 run.status = RunStatus::Judging;
7304 run.parked = true;
7305 let mut task = Task::new(
7306 "add retries".to_owned(),
7307 "add retries".to_owned(),
7308 PathBuf::from("/repo"),
7309 crate::queue::Source::Human,
7310 );
7311 task.status = TaskStatus::Failed;
7312 task.runs = vec![run.id.clone()];
7313 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7314 assert!(with(&task, &run), "parked after judging is the case");
7315
7316 let mut not_parked = run.clone();
7317 not_parked.parked = false;
7318 not_parked.status = RunStatus::Stalled;
7319 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7320
7321 let mut fresh = task.clone();
7322 fresh.fresh_start = true;
7323 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7324
7325 let mut review = task.clone();
7326 review.review_branch = Some("magi/x/A".to_owned());
7327 assert!(!with(&review, &run), "review is ranked before resume");
7328
7329 let mut held = task.clone();
7330 held.status = TaskStatus::Held;
7331 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7332
7333 let mut released = run.clone();
7334 released.released_to = Some("20260901-000000-new1".to_owned());
7335 assert!(!with(&task, &released), "nothing left to resume into");
7336
7337 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7338 "unreadable"
7339 )));
7340 }
7341
7342 #[test]
7343 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7344 let mut released = RunState::new(
7345 PathBuf::from("/repo"),
7346 "main".to_owned(),
7347 "abc1234def".to_owned(),
7348 "add retries".to_owned(),
7349 Config::default(),
7350 );
7351 released.status = RunStatus::Blocked;
7352 assert_eq!(
7353 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7354 Some(released.id.clone())
7355 );
7356 released.released_to = Some("20260901-000000-new1".to_owned());
7357 assert_eq!(
7358 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7359 None,
7360 "there is nothing left to resume it into"
7361 );
7362 }
7363
7364 #[test]
7365 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
7366 let mut exhausted = RunState::new(
7375 PathBuf::from("/repo"),
7376 "main".to_owned(),
7377 "abc1234def".to_owned(),
7378 "add retries".to_owned(),
7379 Config::default(),
7380 );
7381 exhausted.status = RunStatus::Blocked;
7382 exhausted.config.graph.review_rounds = 1;
7383 exhausted.reviews = vec![review_round(1)];
7384
7385 assert_eq!(
7386 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
7387 None,
7388 "an exhausted `Blocked` run must not be offered as resumable"
7389 );
7390
7391 let mut has_budget_left = RunState::new(
7394 PathBuf::from("/repo"),
7395 "main".to_owned(),
7396 "abc1234def".to_owned(),
7397 "add retries".to_owned(),
7398 Config::default(),
7399 );
7400 has_budget_left.status = RunStatus::Blocked;
7401 has_budget_left.config.graph.review_rounds = 3;
7402 has_budget_left.reviews = vec![review_round(1)];
7403
7404 assert_eq!(
7405 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
7406 Ok(has_budget_left.clone())
7407 }),
7408 Some(has_budget_left.id.clone())
7409 );
7410 }
7411
7412 #[test]
7413 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
7414 let mut older_stalled = RunState::new(
7422 PathBuf::from("/repo"),
7423 "main".to_owned(),
7424 "abc1234def".to_owned(),
7425 "add retries".to_owned(),
7426 Config::default(),
7427 );
7428 older_stalled.status = RunStatus::Stalled;
7429
7430 let mut newest_exhausted = RunState::new(
7431 PathBuf::from("/repo"),
7432 "main".to_owned(),
7433 "abc1234def".to_owned(),
7434 "add retries".to_owned(),
7435 Config::default(),
7436 );
7437 newest_exhausted.status = RunStatus::Blocked;
7438 newest_exhausted.config.graph.review_rounds = 1;
7439 newest_exhausted.reviews = vec![review_round(1)];
7440
7441 assert_eq!(
7442 unfinished_run_with(
7443 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
7444 "t",
7445 |_| Ok(newest_exhausted.clone())
7446 ),
7447 None,
7448 "the newest run is exhausted, so nothing here is worth resuming - \
7449 least of all the older, already-superseded run"
7450 );
7451 }
7452
7453 #[test]
7454 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
7455 assert_eq!(
7456 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
7457 Err(anyhow::anyhow!("fixture is absent"))
7458 }),
7459 None
7460 );
7461 }
7462
7463 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
7464 let mut q = ask::Question::new(
7465 run.to_owned(),
7466 "implement".to_owned(),
7467 "impl-A".to_owned(),
7468 "continue?".to_owned(),
7469 String::new(),
7470 vec!["resume で続行する".to_owned(), "other".to_owned()],
7471 );
7472 q.actions.insert("resume で続行する".to_owned(), action);
7473 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
7474 .unwrap();
7475 q
7476 }
7477
7478 fn held_task_with(run: &str) -> Task {
7479 let mut t = task();
7480 t.runs = vec![run.to_owned()];
7481 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
7482 t
7483 }
7484
7485 fn resume_action(run: &str) -> ask::ChoiceAction {
7486 ask::ChoiceAction::Resume { run: run.into() }
7487 }
7488
7489 #[test]
7490 fn decide_action_resumes_only_the_latest_resumable_run() {
7491 let t = held_task_with("r1");
7492 let q = action_question("r1", resume_action("r1"));
7493 let load = |s: RunState| move |_: &str| Ok(s);
7494 assert_eq!(
7495 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
7496 ActionDecision::Resume("r1".into())
7497 );
7498 let q_other = action_question("r1", resume_action("r0"));
7500 assert!(matches!(
7501 decide_action(
7502 &t,
7503 &q_other,
7504 &PHRASES_EN,
7505 load(run_state(RunStatus::Blocked))
7506 ),
7507 ActionDecision::Refuse(_)
7508 ));
7509 assert!(matches!(
7511 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
7512 ActionDecision::Refuse(_)
7513 ));
7514 let mut released = run_state(RunStatus::Blocked);
7516 released.released_to = Some("elsewhere".into());
7517 assert!(matches!(
7518 decide_action(&t, &q, &PHRASES_EN, load(released)),
7519 ActionDecision::Refuse(_)
7520 ));
7521 assert!(matches!(
7523 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
7524 ActionDecision::Refuse(_)
7525 ));
7526 }
7527
7528 #[test]
7529 fn decide_action_ignores_a_question_about_an_earlier_run() {
7530 let mut t = held_task_with("r1");
7531 t.runs.push("r2".to_owned());
7532 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7533 assert_eq!(
7534 decide_action(
7535 &t,
7536 &action_question("r1", ask::ChoiceAction::Done),
7537 &PHRASES_EN,
7538 never
7539 ),
7540 ActionDecision::Stale
7541 );
7542 }
7543
7544 #[test]
7545 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
7546 let mut t = held_task_with("r1");
7547 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7548 assert_eq!(
7549 decide_action(
7550 &t,
7551 &action_question("r1", ask::ChoiceAction::Requeue),
7552 &PHRASES_EN,
7553 never
7554 ),
7555 ActionDecision::Requeue
7556 );
7557 let done_q = action_question("r1", ask::ChoiceAction::Done);
7558 assert_eq!(
7559 decide_action(&t, &done_q, &PHRASES_EN, never),
7560 ActionDecision::Done
7561 );
7562 t.mark_action_applied(&done_q.id);
7563 assert_eq!(
7564 decide_action(&t, &done_q, &PHRASES_EN, never),
7565 ActionDecision::Skip
7566 );
7567
7568 let mut plain = action_question("r1", ask::ChoiceAction::Done);
7570 plain.actions.clear();
7571 assert_eq!(
7572 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
7573 ActionDecision::Skip
7574 );
7575 let mut running = held_task_with("r1");
7577 running.status = TaskStatus::Running;
7578 assert_eq!(
7579 decide_action(
7580 &running,
7581 &action_question("r1", ask::ChoiceAction::Done),
7582 &PHRASES_EN,
7583 never
7584 ),
7585 ActionDecision::Skip
7586 );
7587 }
7588
7589 #[test]
7590 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
7591 let dir = tempfile::tempdir().unwrap();
7592 let queue = Queue::at(dir.path().join("queue"));
7593 let questions = Questions::at(dir.path().join("questions"));
7594 let home = dir.path().join("home");
7595 let mut state = run_state(RunStatus::Blocked);
7596 state.id = "20260101-000000-act1".to_owned();
7597 state.save_under(&home).unwrap();
7598
7599 let mut t = held_task_with(&state.id);
7600 queue.put(&mut t).unwrap();
7601 let mut q = action_question(&state.id, resume_action(&state.id));
7602 questions.put(&mut q).unwrap();
7603
7604 apply_choice_actions(&queue, &questions, &home);
7605 let after = queue.get(&t.id).unwrap();
7606 assert_eq!(after.status, TaskStatus::Queued);
7607 assert!(!after.fresh_start);
7608 assert!(after.action_applied(&q.id));
7609 let pin = after.resume_override.clone().unwrap();
7610 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
7611 assert!(pin.forced);
7612
7613 let mut again = queue.get(&t.id).unwrap();
7615 again.hold_machine(Some("later".into()));
7616 queue.put(&mut again).unwrap();
7617 apply_choice_actions(&queue, &questions, &home);
7618 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7619 }
7620
7621 #[test]
7622 fn an_answer_the_waiter_already_delivered_is_not_acted_on_again() {
7623 let dir = tempfile::tempdir().unwrap();
7624 let queue = Queue::at(dir.path().join("queue"));
7625 let questions = Questions::at(dir.path().join("questions"));
7626 let home = dir.path().join("home");
7627 let mut state = run_state(RunStatus::Blocked);
7628 state.id = "20260101-000000-act2".to_owned();
7629 state.save_under(&home).unwrap();
7630
7631 let mut t = held_task_with(&state.id);
7632 queue.put(&mut t).unwrap();
7633 let mut q = action_question(&state.id, resume_action(&state.id));
7634 q.answer_delivered = true;
7635 questions.put(&mut q).unwrap();
7636
7637 apply_choice_actions(&queue, &questions, &home);
7638 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7639
7640 let mut q2 = action_question(&state.id, resume_action(&state.id));
7642 questions.put(&mut q2).unwrap();
7643 apply_choice_actions(&queue, &questions, &home);
7644 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7645 assert!(questions.get(&q2.id).unwrap().answer_delivered);
7646 }
7647}