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) {
1122 settle_in(task, verdict, detail, max_attempts, &PHRASES_EN)
1123}
1124
1125fn settle_in(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize, p: &Phrases) {
1127 if verdict.parked {
1133 task.stall(detail);
1134 return;
1135 }
1136 match verdict.status {
1137 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1138 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1139 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1140 task.stall(detail)
1141 }
1142 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1143 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1144 RunStatus::Blocked => task.fail(detail, max_attempts),
1145 RunStatus::VerifiedNoop => task.handed_off(detail),
1146 other => task.fail((p.graph_stopped)(label(other), detail), max_attempts),
1147 }
1148}
1149
1150pub fn supersede_prior_runs(task: &Task, home: &Path) {
1199 let now = Timestamp::now();
1200 let last_run_succeeded = task
1207 .runs
1208 .last()
1209 .and_then(|id| RunState::load_under(id, home).ok())
1210 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1211 for id in task.superseded_attempts(last_run_succeeded) {
1212 let mut state = match RunState::load_under(id, home) {
1213 Ok(s) => s,
1214 Err(e) => {
1215 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1216 continue;
1217 }
1218 };
1219 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1220 continue;
1221 }
1222 let daemon_claims = is_working_on(home, id, now);
1223 if state.liveness(daemon_claims) == Liveness::Live {
1224 continue;
1225 }
1226 state.status = RunStatus::Superseded;
1227 if let Err(e) = state.save_under(home) {
1228 tracing::warn!("could not mark run {id} superseded: {e:#}");
1229 }
1230 }
1231}
1232
1233fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1255 for task in queue.list() {
1256 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1257 continue;
1258 }
1259 supersede_prior_runs(&task, home);
1260 }
1261}
1262
1263fn settle_and_diagnose(
1270 task: &mut Task,
1271 verdict: Verdict,
1272 detail: &str,
1273 max_attempts: usize,
1274 state: &RunState,
1275) {
1276 let p = phrases(&state.config.graph.language);
1277 settle_in(task, verdict, detail, max_attempts, p);
1278 if task.status == TaskStatus::Held {
1279 task.diagnostic = diagnostic(state);
1280 note_open_question(task, &state.id, p);
1281 }
1282}
1283
1284fn note_open_question(task: &mut Task, run: &str, p: &Phrases) {
1299 let Some(home) = crate::run::try_home() else {
1300 return;
1301 };
1302 let open = Questions::at(home.join("questions")).open_for(run);
1303 let Some(q) = open.first() else {
1304 return;
1305 };
1306 let base = task.hold_reason.clone().unwrap_or_default();
1307 task.hold_reason = Some(format!("{base}{}{}", p.waiting_for_answer, q.short()));
1308}
1309
1310fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize, language: &str) {
1321 match last_run {
1322 Some(state) => {
1323 let verdict = Verdict {
1324 status: state.status,
1325 left_pr: state.pr.is_some(),
1326 quota_hit: !state.quota.is_empty(),
1327 parked: state.parked,
1328 no_viable_candidates: state.viable().is_empty(),
1329 };
1330 let detail = format!(
1331 "{}{}",
1332 phrases(&state.config.graph.language).recovered_running,
1333 describe(&state)
1334 );
1335 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1336 }
1337 None => {
1338 let why = phrases(language).no_run_to_recover;
1341 task.last_error = Some(why.to_owned());
1342 task.hold_machine(Some(why.to_owned()));
1345 }
1346 }
1347}
1348
1349fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1371 let mut reclaimed = Vec::new();
1372 for listed in queue.list() {
1373 if listed.status != TaskStatus::Running {
1374 continue;
1375 }
1376 let Ok(_claim) = queue.claim(&listed.id) else {
1377 continue;
1378 };
1379 let Ok(mut task) = queue.get(&listed.id) else {
1383 continue;
1384 };
1385 if task.status != TaskStatus::Running {
1386 continue;
1387 }
1388 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1389 if let Some(state) = &last_run
1399 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1400 {
1401 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1402 }
1403 let language = if last_run.is_none() {
1404 language_of(&task, Path::new("."))
1405 } else {
1406 String::new()
1407 };
1408 reclaim(&mut task, last_run, max_attempts, &language);
1409 if task.status == TaskStatus::Done {
1410 supersede_prior_runs(&task, &crate::run::home());
1411 }
1412 record(queue, &mut task);
1413 reclaimed.push(task.id.clone());
1414 }
1415 reclaimed
1416}
1417
1418fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1450 reclaim_abandoned_runs_with(
1451 home,
1452 now,
1453 crate::proc::pid_status,
1454 crate::proc::process_started_at,
1455 )
1456}
1457
1458fn reclaim_abandoned_runs_with<F, G>(
1464 home: &Path,
1465 now: Timestamp,
1466 query: F,
1467 identity: G,
1468) -> Vec<String>
1469where
1470 F: Fn(u32) -> Option<bool>,
1471 G: Fn(u32) -> Option<String>,
1472{
1473 let mut abandoned = Vec::new();
1474 for entry in std::fs::read_dir(home.join("runs"))
1475 .into_iter()
1476 .flatten()
1477 .flatten()
1478 {
1479 let id = entry.file_name().to_string_lossy().into_owned();
1480 if !crate::run::is_run_id(&id) {
1481 continue;
1482 }
1483 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1491 continue;
1492 };
1493 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1494 continue;
1495 };
1496 if state.status.done() || !state.active_all_overrun(now) {
1497 continue;
1498 }
1499 let daemon_claims = is_working_on(home, &id, now);
1511 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1512 continue;
1513 }
1514 state.abandon("daemon");
1515 if let Err(e) = state.save_under(home) {
1516 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1517 continue;
1518 }
1519 if let Some(notice) = notices::run_ended(&state) {
1522 notices::raise_in(home, notice);
1523 }
1524 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1532 tracing::warn!("abandon questions for {id}: {e:#}");
1533 }
1534 abandoned.push(id);
1535 }
1536 abandoned
1537}
1538
1539pub async fn serve(opts: Opts) -> Result<()> {
1545 serve_until(opts, Stop::new()).await
1546}
1547
1548pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1565 let signal = {
1566 let stop = stop.clone();
1567 tokio::spawn(async move {
1568 if tokio::signal::ctrl_c().await.is_ok() {
1569 stop.stop();
1570 tracing::info!("shutdown requested; a run in flight will be finished first");
1571 }
1572 })
1573 };
1574
1575 let worktrees_root = opts
1576 .worktrees_root
1577 .clone()
1578 .unwrap_or_else(crate::run::default_worktree_root);
1579 let outcome = drive(
1580 &opts,
1581 &Queue::open(),
1582 &status_path(),
1583 &crate::run::home(),
1584 &worktrees_root,
1585 &stop,
1586 )
1587 .await;
1588
1589 signal.abort();
1590 outcome
1591}
1592
1593async fn drive(
1606 opts: &Opts,
1607 queue: &Queue,
1608 status_file: &Path,
1609 home: &Path,
1610 worktrees_root: &Path,
1611 stop: &Stop,
1612) -> Result<()> {
1613 let status = Arc::new(Mutex::new(Status::new()));
1621 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1622 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1623
1624 let daemon_cfg = prepare(&opts.repo, opts)
1630 .map(|c| c.daemon)
1631 .unwrap_or_default();
1632 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1633
1634 let waiter = tokio::spawn(crate::waiter::run(
1639 crate::waiter::Waiter::new(
1640 crate::ask::Questions::at(home.join("questions")),
1641 home.to_path_buf(),
1642 prepare(&opts.repo, opts).ok(),
1643 ),
1644 stop.clone(),
1645 ));
1646
1647 tracing::info!(
1648 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1649 queue.root().display(),
1650 opts.poll.as_secs(),
1651 opts.max_attempts,
1652 concurrency,
1653 if daemon_cfg.pause_for_interrupts {
1654 ", interrupts enabled"
1655 } else {
1656 ""
1657 }
1658 );
1659
1660 janitor(&opts.repo, opts, home, worktrees_root).await;
1663 resweep_superseded_attempts(queue, home);
1664
1665 let outcome = poll(
1666 opts,
1667 queue,
1668 &status,
1669 home,
1670 worktrees_root,
1671 stop,
1672 DispatchLimits {
1673 max_concurrent: concurrency,
1674 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1675 },
1676 )
1677 .await;
1678
1679 beat.abort();
1680 waiter.abort();
1681 clear_status_at(status_file);
1682 outcome
1683}
1684
1685async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1691 loop {
1692 tokio::time::sleep(HEARTBEAT).await;
1693 let snapshot = {
1694 let mut guard = lock(&status);
1695 guard.updated_at = Timestamp::now();
1696 guard.clone()
1697 };
1698 if let Err(e) = write_status_to(&path, &snapshot) {
1699 tracing::warn!("could not refresh the daemon status file: {e:#}");
1702 }
1703 }
1704}
1705
1706#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1709enum LandResume {
1710 NotLanding,
1713 StillWaiting,
1718 Ready,
1722}
1723
1724fn land_resume_state(task: &Task) -> LandResume {
1728 let Some(run_id) = task.runs.last() else {
1729 return LandResume::NotLanding;
1730 };
1731 let Ok(state) = RunState::load(run_id) else {
1732 return LandResume::NotLanding;
1733 };
1734 if state.status != RunStatus::Landing || !state.parked {
1735 return LandResume::NotLanding;
1736 }
1737 let store = ask::Questions::open();
1738 let waiting = store
1739 .list()
1740 .into_iter()
1741 .filter(|q| &q.run == run_id && q.node == land::APPROVAL_NODE)
1742 .max_by(|a, b| a.id.cmp(&b.id));
1743 let Some(mut q) = waiting else {
1744 return LandResume::Ready;
1745 };
1746 if !q.status.open() {
1747 return LandResume::Ready;
1748 }
1749 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
1756 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
1757 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
1758 q.abandon(format!(
1759 "no answer within {}s of asking",
1760 timeout.as_secs().max(1)
1761 ));
1762 if store.put(&mut q).is_ok() {
1765 return LandResume::Ready;
1766 }
1767 }
1768 LandResume::StillWaiting
1769}
1770
1771const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
1779
1780const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
1792
1793struct InFlightGuard<'a> {
1806 status: &'a Arc<Mutex<Status>>,
1807 stop: &'a Stop,
1808 task_id: &'a str,
1809}
1810
1811impl Drop for InFlightGuard<'_> {
1812 fn drop(&mut self) {
1813 lock(self.status).current.retain(|c| c.task != self.task_id);
1814 self.stop.exit();
1815 }
1816}
1817
1818#[derive(Debug, Clone, PartialEq, Eq)]
1840enum Interrupt {
1841 Idle,
1843 Parking {
1855 parked: Vec<String>,
1856 interrupt_task: String,
1857 },
1858 Running {
1866 parked: Vec<String>,
1867 interrupt_task: String,
1868 },
1869 Resuming { parked: Vec<String> },
1876}
1877
1878fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
1899 match state {
1900 Interrupt::Idle => {
1901 if in_flight.len() != 1 {
1913 return Interrupt::Idle;
1914 }
1915 match runnable.iter().find(|t| t.interrupt) {
1916 Some(t) => Interrupt::Parking {
1917 parked: in_flight.to_vec(),
1918 interrupt_task: t.id.clone(),
1919 },
1920 None => Interrupt::Idle,
1921 }
1922 }
1923 Interrupt::Parking {
1924 parked,
1925 interrupt_task,
1926 } => {
1927 if in_flight.iter().any(|id| parked.contains(id)) {
1928 Interrupt::Parking {
1930 parked,
1931 interrupt_task,
1932 }
1933 } else if in_flight.contains(&interrupt_task) {
1934 Interrupt::Running {
1935 parked,
1936 interrupt_task,
1937 }
1938 } else if runnable.iter().any(|t| t.id == interrupt_task) {
1939 Interrupt::Parking {
1943 parked,
1944 interrupt_task,
1945 }
1946 } else {
1947 Interrupt::Resuming { parked }
1952 }
1953 }
1954 Interrupt::Running {
1955 parked,
1956 interrupt_task,
1957 } => {
1958 if in_flight.contains(&interrupt_task) {
1959 Interrupt::Running {
1960 parked,
1961 interrupt_task,
1962 }
1963 } else {
1964 Interrupt::Resuming { parked }
1970 }
1971 }
1972 Interrupt::Resuming { parked } => {
1973 if in_flight.iter().any(|id| parked.contains(id)) {
1974 Interrupt::Idle
1980 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
1981 Interrupt::Resuming { parked }
1982 } else {
1983 Interrupt::Idle
1986 }
1987 }
1988 }
1989}
1990
1991fn advance_interrupt_tick(
1997 enabled: bool,
1998 state: Interrupt,
1999 in_flight: &[String],
2000 runnable: &[Task],
2001) -> Interrupt {
2002 if !enabled {
2003 return Interrupt::Idle;
2004 }
2005 advance_interrupt(state, in_flight, runnable)
2006}
2007
2008fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
2013 match state {
2014 Interrupt::Idle => candidates,
2015 Interrupt::Parking {
2016 parked,
2017 interrupt_task,
2018 } => {
2019 if in_flight.iter().any(|id| parked.contains(id)) {
2020 Vec::new()
2021 } else {
2022 candidates
2023 .into_iter()
2024 .filter(|t| &t.id == interrupt_task)
2025 .collect()
2026 }
2027 }
2028 Interrupt::Running { .. } => Vec::new(),
2029 Interrupt::Resuming { parked } => candidates
2037 .into_iter()
2038 .find(|t| parked.contains(&t.id))
2039 .into_iter()
2040 .collect(),
2041 }
2042}
2043
2044struct DispatchLimits {
2048 max_concurrent: usize,
2051 pause_for_interrupts: bool,
2053}
2054
2055#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2064enum PermitKind {
2065 None,
2070 Urgent,
2077 Ordinary,
2080}
2081
2082fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
2085 if priority {
2086 PermitKind::None
2087 } else if urgent {
2088 PermitKind::Urgent
2089 } else {
2090 PermitKind::Ordinary
2091 }
2092}
2093
2094async fn poll(
2113 opts: &Opts,
2114 queue: &Queue,
2115 status: &Arc<Mutex<Status>>,
2116 home: &Path,
2117 worktrees_root: &Path,
2118 stop: &Stop,
2119 limits: DispatchLimits,
2120) -> Result<()> {
2121 let DispatchLimits {
2122 max_concurrent,
2123 pause_for_interrupts,
2124 } = limits;
2125 let mut attempted: Vec<String> = Vec::new();
2130 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2131 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2139 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2145 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2146 let mut conductor = Conductor::new();
2147 let mut cache_last_checked: Option<Timestamp> = None;
2150 let mut interrupt = Interrupt::Idle;
2152 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2158 std::collections::HashMap::new();
2159
2160 while !stop.stopped() {
2161 lock(status).polls += 1;
2162
2163 while let Some(result) = inflight.try_join_next() {
2168 if let Err(e) = result {
2169 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2170 notices::raise(Notice::error(
2171 "loop:attempt",
2172 "A queued attempt ended abnormally; check the task it was running.",
2173 ));
2174 }
2175 }
2176
2177 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2178 if !swept.is_empty() {
2179 tracing::warn!(
2180 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2181 swept.len(),
2182 swept.join(", ")
2183 );
2184 }
2185 let now = Timestamp::now();
2190
2191 if !stop.busy_now() {
2196 maybe_prune_cache_between_runs(
2197 &opts.repo,
2198 opts,
2199 home,
2200 stop,
2201 &mut cache_last_checked,
2202 now,
2203 )
2204 .await;
2205 }
2206
2207 let stalled = stalled_tasks(queue, home, now);
2208 let stalled_ids: std::collections::BTreeSet<_> =
2209 stalled.iter().map(|task| task.id.clone()).collect();
2210 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2211 if !reclaimed.is_empty() {
2212 tracing::warn!(
2213 "reclaimed {} task(s) left `running` by a daemon that never \
2214 recorded the outcome: {}",
2215 reclaimed.len(),
2216 reclaimed.join(", ")
2217 );
2218 }
2219 let abandoned_runs = reclaim_abandoned_runs(home, now);
2220 if !abandoned_runs.is_empty() {
2221 tracing::warn!(
2222 "failed {} run(s) left behind by a killed process, past every \
2223 active seat's own timeout: {}",
2224 abandoned_runs.len(),
2225 abandoned_runs.join(", ")
2226 );
2227 }
2228
2229 let questions = Questions::at(home.join("questions"));
2234
2235 resolve_blockers(queue, &questions);
2238 apply_choice_actions(queue, &questions, home);
2239 reconcile_task_questions(queue, &questions);
2240
2241 let finished: Vec<Task> = finished_tasks(queue)
2247 .into_iter()
2248 .filter(|task| !stalled_ids.contains(&task.id))
2249 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2253 .collect();
2254 let queued = queued_tasks(queue);
2255 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2259 && conductor.worth_a_look(queue, &stalled, &finished)
2260 {
2261 match prepare(&opts.repo, opts) {
2262 Ok(cfg) => {
2263 conductor
2264 .maybe_run(
2265 &cfg,
2266 &opts.repo,
2267 queue,
2268 &questions,
2269 home,
2270 &queued,
2271 &stalled,
2272 &finished,
2273 opts.max_attempts,
2274 )
2275 .await;
2276 }
2277 Err(e) => {
2278 tracing::warn!("conductor: no config: {e:#}");
2279 notices::raise(Notice::warn(
2280 "loop:no-config",
2281 "The loop could not read this repository's config, so held tasks are not being triaged.",
2282 ));
2283 }
2284 }
2285 }
2286
2287 let candidates: Vec<Task> = runnable(queue)
2288 .into_iter()
2289 .filter(|t| !opts.once || !attempted.contains(&t.id))
2290 .collect();
2291
2292 let in_flight: Vec<String> = lock(status)
2297 .current
2298 .iter()
2299 .map(|c| c.task.clone())
2300 .collect();
2301 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2302
2303 interrupt =
2304 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2305 if let Interrupt::Parking {
2306 parked,
2307 interrupt_task,
2308 } = &interrupt
2309 {
2310 let reason = format!(
2311 "task {} asked to run first",
2312 crate::run::short_of(interrupt_task)
2313 );
2314 for id in parked {
2315 if let Some(pause) = interrupt_pauses.get(id) {
2316 pause.park_because(reason.clone());
2317 }
2318 }
2319 }
2320 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2321
2322 let cooling_down =
2323 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2324
2325 let mut started_any = false;
2326 for candidate in candidates {
2327 if stop.stopped() {
2328 break;
2329 }
2330
2331 let resume = land_resume_state(&candidate);
2332 if resume == LandResume::StillWaiting {
2333 continue;
2334 }
2335 let priority = resume == LandResume::Ready;
2336
2337 if !priority && cooling_down {
2338 continue;
2339 }
2340 let permit = match permit_kind(priority, candidate.urgent) {
2341 PermitKind::None => None,
2342 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2343 Ok(p) => Some(p),
2344 Err(_) => continue,
2350 },
2351 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2352 Ok(p) => Some(p),
2353 Err(_) => continue,
2357 },
2358 };
2359
2360 let Ok(claim) = queue.claim(&candidate.id) else {
2365 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2366 continue;
2367 };
2368 let mut task = match queue.get(&candidate.id) {
2371 Ok(t) if t.status.runnable() => t,
2372 Ok(_) => continue,
2373 Err(e) => {
2374 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2375 continue;
2376 }
2377 };
2378 let task_id = task.id.clone();
2379 attempted.push(task_id.clone());
2380 lock(status).idle = false;
2381 stop.enter();
2384 started_any = true;
2385
2386 let run_pause = crate::graph::Pause::new();
2390 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2391
2392 let opts = opts.clone();
2393 let queue = queue.clone();
2394 let status = Arc::clone(status);
2395 let stop = stop.clone();
2396 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2397 inflight.spawn(async move {
2398 let _claim = claim;
2402 let _permit = permit;
2403 let _inflight = InFlightGuard {
2405 status: &status,
2406 stop: &stop,
2407 task_id: &task_id,
2408 };
2409 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2410 lock(&status).completed += 1;
2411 let now = Timestamp::now();
2417 if let Some(until) = cooldown_until("a, now) {
2418 let wait = until.as_second() - now.as_second();
2419 *lock("a_cooldown_until) = Some(until);
2420 let hint = quota
2421 .iter()
2422 .find(|q| q.reset.is_some())
2423 .and_then(|q| q.reset.as_deref());
2424 match hint {
2425 Some(h) => tracing::warn!(
2426 "quota hit; waiting {wait}s before taking another ordinary task \
2427 (CLI reported reset: {h})"
2428 ),
2429 None => tracing::warn!(
2430 "quota hit; waiting {wait}s before taking another ordinary task \
2431 (no reset hint reported)"
2432 ),
2433 }
2434 }
2435 });
2436 }
2437
2438 if started_any {
2439 continue;
2440 }
2441
2442 if stop.busy_now() {
2443 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2448 continue;
2449 }
2450
2451 lock(status).idle = true;
2453 if opts.once {
2454 janitor(&opts.repo, opts, home, worktrees_root).await;
2458 resweep_superseded_attempts(queue, home);
2459 triage_held(queue, home, opts).await;
2460 break;
2461 }
2462 stop.idle(opts.poll).await;
2463 if stop.stopped() {
2464 continue;
2465 }
2466 janitor(&opts.repo, opts, home, worktrees_root).await;
2472 resweep_superseded_attempts(queue, home);
2473 triage_held(queue, home, opts).await;
2474 }
2475
2476 while let Some(result) = inflight.join_next().await {
2481 if let Err(e) = result {
2482 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2483 notices::raise(Notice::error(
2484 "loop:attempt",
2485 "A queued attempt ended abnormally; check the task it was running.",
2486 ));
2487 }
2488 }
2489 Ok(())
2490}
2491
2492async fn attempt(
2498 opts: &Opts,
2499 queue: &Queue,
2500 status: &Arc<Mutex<Status>>,
2501 stop: &Stop,
2502 interrupt_pause: crate::graph::Pause,
2503 task: &mut Task,
2504) -> Vec<QuotaLoss> {
2505 let repo = repo_for(task, &opts.repo);
2506 tracing::info!(
2507 "task {} — {} (repo {})",
2508 task.short(),
2509 task.title,
2510 repo.display()
2511 );
2512
2513 let mut config = match prepare(&repo, opts) {
2514 Ok(c) => c,
2515 Err(e) => {
2516 task.attempts += 1;
2520 task.fail(format!("config: {e:#}"), opts.max_attempts);
2521 record(queue, task);
2522 return Vec::new();
2523 }
2524 };
2525 apply_solo(&mut config, task);
2526 let start_failed = phrases(&config.graph.language).could_not_start;
2527
2528 if let Some(reason) = disk_gate(&repo, &config) {
2536 task.last_error = Some(reason.clone());
2537 task.hold_machine(Some(reason.clone()));
2538 record(queue, task);
2539 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2540 notices::raise(
2543 Notice::warn(
2544 &format!("disk:{}", repo.display()),
2545 "A task was held for want of free disk space; free some, then release it from the queue.",
2546 )
2547 .link(Link::Task {
2548 id: task.id.clone(),
2549 }),
2550 );
2551 return Vec::new();
2552 }
2553
2554 let unfinished = (!task.fresh_start)
2574 .then(|| unfinished_run(&task.runs, task.short()))
2575 .flatten();
2576 let review_branch = task.review_branch.take();
2582 let branch_exists = match &review_branch {
2583 Some(branch) => crate::git::branch_exists(&repo, branch)
2584 .await
2585 .unwrap_or(false),
2586 None => false,
2587 };
2588 let starter = choose_starter(
2589 review_branch.as_deref(),
2590 branch_exists,
2591 unfinished.as_deref(),
2592 );
2593 let attachments = match task_attachments(queue, task) {
2598 Ok(a) => a,
2599 Err(e) => {
2600 task.attempts += 1;
2601 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2602 record(queue, task);
2603 return Vec::new();
2604 }
2605 };
2606 let started = match &starter {
2607 Starter::Review(branch) => {
2608 tracing::info!(
2609 "task {} reopens `{branch}` as a review-only pass",
2610 task.short()
2611 );
2612 let takeover = crate::handover::Takeover {
2615 earlier: task.earlier_attempts().to_vec(),
2616 home: crate::run::home(),
2617 choice: take_divergence_answer(branch, &config.merge.remote, task),
2618 };
2619 Runner::review_taking_over(&repo, branch, config, Some(takeover)).await
2620 }
2621 Starter::Resume(id) => {
2622 tracing::info!("resuming run {id} rather than competing again");
2623 Runner::resume(id).map(|mut r| {
2624 if let Some(instruction) =
2625 prepare_instruction(&starter, Some(&r.state.instruction), task)
2626 {
2627 r.state.instruction = instruction;
2628 }
2629 r.state.attachments = attachments.clone();
2630 r
2631 })
2632 }
2633 Starter::Start => {
2634 if let Some(branch) = &review_branch {
2635 tracing::warn!(
2636 "conductor chose review for task {} but branch `{branch}` no longer \
2637 exists; requeuing as a fresh competition instead",
2638 task.short()
2639 );
2640 }
2641 let instruction = prepare_instruction(&starter, None, task)
2642 .unwrap_or_else(|| task.instruction.clone());
2643 Runner::start_naming(&repo, instruction, &task.title, config)
2644 .await
2645 .map(|mut r| {
2646 r.state.attachments = attachments.clone();
2647 r
2648 })
2649 }
2650 };
2651 let mut runner = match started {
2652 Ok(r) => r,
2653 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2657 let reason = format!("{start_failed}{e:#}");
2658 task.last_error = Some(reason.clone());
2659 task.hold_machine(Some(reason));
2660 record(queue, task);
2661 tracing::warn!(
2662 "holding {} for a branch it cannot take over: {e:#}",
2663 task.short()
2664 );
2665 notices::raise(
2667 Notice::warn(
2668 &format!("handover:{}", task.id),
2669 "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.",
2670 )
2671 .link(Link::Task {
2672 id: task.id.clone(),
2673 })
2674 .about([task.id.clone()]),
2675 );
2676 return Vec::new();
2677 }
2678 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2683 let d = e
2684 .downcast_ref::<crate::reconcile::Diverged>()
2685 .expect("checked by the guard");
2686 let mut q = ask::Question::new(
2687 task.id.clone(),
2688 "review".to_owned(),
2689 "sync".to_owned(),
2690 d.summary(),
2691 d.detail(),
2692 d.choices(),
2693 );
2694 match Questions::open().put(&mut q) {
2695 Ok(()) => {
2696 task.last_error = Some(format!("{e:#}"));
2697 task.review_branch = Some(d.branch.clone());
2700 task.block(vec![q.id.clone()], Some(d.summary()));
2701 }
2702 Err(put) => {
2703 tracing::warn!("could not file the divergence question: {put:#}");
2704 task.attempts += 1;
2705 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2706 }
2707 }
2708 record(queue, task);
2709 return Vec::new();
2710 }
2711 Err(e) => {
2712 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
2718 && let Starter::Review(branch) = &starter
2719 {
2720 task.review_branch = Some(branch.clone());
2721 task.last_error = Some(format!("{start_failed}{e:#}"));
2722 task.status = crate::queue::TaskStatus::Failed;
2723 record(queue, task);
2724 return Vec::new();
2725 }
2726 task.attempts += 1;
2727 task.fail(format!("{start_failed}{e:#}"), opts.max_attempts);
2728 record(queue, task);
2729 return Vec::new();
2730 }
2731 };
2732 runner.on_pause(stop.pause());
2734 runner.watch_interrupt(interrupt_pause);
2738
2739 let run = runner.state.id.clone();
2742 task.start(run.clone());
2743 record(queue, task);
2744 lock(status).current.push(Current {
2745 task: task.id.clone(),
2746 run,
2747 });
2748
2749 let quota_before = runner.state.quota.clone();
2752 let detail = match runner.execute().await {
2753 Ok(()) => describe(&runner.state),
2754 Err(e) => format!("{e:#}"),
2755 };
2756 let fresh = losses_this_attempt("a_before, &runner.state.quota);
2757 let verdict = Verdict {
2758 status: runner.state.status,
2759 left_pr: runner.state.pr.is_some(),
2762 quota_hit: !fresh.is_empty(),
2768 parked: runner.state.parked,
2772 no_viable_candidates: runner.state.viable().is_empty(),
2775 };
2776 settle_and_diagnose(task, verdict, &detail, opts.max_attempts, &runner.state);
2777 if task.status == TaskStatus::Done {
2778 supersede_prior_runs(task, &crate::run::home());
2779 }
2780 record(queue, task);
2781 tracing::info!(
2782 "task {} is {} after run {} ({})",
2783 task.short(),
2784 task.status.as_str(),
2785 runner.state.short(),
2786 label(runner.state.status)
2787 );
2788 fresh
2789}
2790
2791fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
2801 after
2802 .iter()
2803 .filter(|q| !before.contains(q))
2804 .cloned()
2805 .collect()
2806}
2807
2808fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
2811 if quota.is_empty() {
2812 return None;
2813 }
2814 let with_hint = quota.iter().find(|q| q.reset.is_some());
2815 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
2816 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
2817 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
2818 Some(
2819 now.checked_add(jiff::SignedDuration::from_secs(secs))
2820 .unwrap_or(Timestamp::MAX),
2821 )
2822}
2823
2824fn apply_solo(config: &mut Config, task: &Task) {
2834 if task.solo {
2835 config.graph.candidates = 1;
2836 }
2837}
2838
2839fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
2841 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
2842 if let Some(mode) = &opts.merge {
2843 config.merge.mode = merge_mode(mode)?;
2844 }
2845 Ok(config)
2846}
2847
2848async fn maybe_prune_cache_between_runs(
2879 repo: &Path,
2880 opts: &Opts,
2881 home: &Path,
2882 stop: &Stop,
2883 last_checked: &mut Option<Timestamp>,
2884 now: Timestamp,
2885) {
2886 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
2887 return;
2888 }
2889 *last_checked = Some(now);
2890 let cfg = match prepare(repo, opts) {
2891 Ok(cfg) => cfg,
2892 Err(e) => {
2893 tracing::warn!("cache check: no config: {e:#}");
2894 return;
2895 }
2896 };
2897 match clean::prune_cache_if_over_limit(&cfg, home) {
2898 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
2899 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
2900 pruned.files,
2901 pruned.freed
2902 ),
2903 Ok(_) => {}
2904 Err(e) => {
2905 tracing::warn!("housekeep: prune cache: {e:#}");
2906 notices::raise_in(
2907 home,
2908 Notice::warn(
2909 "housekeep:cache",
2910 "Pruning the shared build cache failed; disk usage may keep growing.",
2911 ),
2912 );
2913 }
2914 }
2915}
2916
2917fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
2921 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
2922}
2923
2924async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
2943 let cfg = match prepare(repo, opts) {
2944 Ok(cfg) => cfg,
2945 Err(e) => {
2946 tracing::warn!("housekeep: no config: {e:#}");
2947 return;
2948 }
2949 };
2950 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
2959 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
2960 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
2965 let mut extra = Vec::new();
2966 if out.unreadable > 0 {
2967 extra.push(format!("{} unreadable", out.unreadable));
2968 }
2969 if out.orphaned_worktrees > 0 {
2970 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
2971 }
2972 let detail = if extra.is_empty() {
2973 String::new()
2974 } else {
2975 format!(" ({})", extra.join(", "))
2976 };
2977 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
2978 }
2979 if out.external_merges_recorded > 0 {
2980 tracing::info!(
2981 "housekeep: recorded {} run(s) as merged externally",
2982 out.external_merges_recorded
2983 );
2984 }
2985 if out.cache_files > 0 {
2986 tracing::info!(
2987 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
2988 out.cache_files,
2989 out.cache_freed
2990 );
2991 }
2992 if out.questions_abandoned > 0 {
2993 tracing::info!(
2994 "housekeep: abandoned {} question(s) left open by a finished run",
2995 out.questions_abandoned
2996 );
2997 }
2998}
2999
3000async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
3009 let questions = Questions::at(home.join("questions"));
3010 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
3011 if report.is_empty() {
3012 return;
3013 }
3014 if !report.quarantined.is_empty() {
3015 tracing::info!(
3016 "triage: held {} blocked task(s) whose blocked-on task or \
3017 question no longer exists: {}",
3018 report.quarantined.len(),
3019 report.quarantined.join(", ")
3020 );
3021 }
3022 if !report.resumed.is_empty() {
3023 tracing::info!(
3024 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
3025 report.resumed.len(),
3026 report.resumed.join(", ")
3027 );
3028 }
3029 if !report.asked.is_empty() {
3030 tracing::info!(
3031 "triage: asked about {} held task(s): {}",
3032 report.asked.len(),
3033 report.asked.join(", ")
3034 );
3035 }
3036 if !report.answered.is_empty() {
3037 tracing::info!(
3038 "triage: applied {} operator answer(s): {}",
3039 report.answered.len(),
3040 report.answered.join(", ")
3041 );
3042 }
3043}
3044
3045fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
3052 disk_gate_with(repo, config, crate::disk::free_bytes)
3053}
3054
3055fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
3059 repo: &Path,
3060 config: &Config,
3061 free_bytes: F,
3062) -> Option<String> {
3063 let min = config.disk.min_free_bytes;
3064 if min == 0 {
3065 return None;
3066 }
3067 match free_bytes(repo) {
3068 Ok(free) => crate::disk::gate_in(free, min, &config.graph.language),
3069 Err(e) => Some(crate::disk::unmeasured_in(repo, &e, &config.graph.language)),
3070 }
3071}
3072
3073const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
3080
3081const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
3085
3086fn quota_wait(
3095 reset_at: Option<Timestamp>,
3096 now: Timestamp,
3097 fallback: Duration,
3098 cap: Duration,
3099) -> Duration {
3100 match reset_at {
3101 Some(at) if at > now => {
3102 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3103 Duration::from_secs(secs).min(cap)
3104 }
3105 _ => fallback,
3106 }
3107}
3108
3109fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3123 parse_reset_hint_zoned(text, now)
3124 .or_else(|| parse_reset_hint_dated(text))
3125 .or_else(|| parse_reset_hint_relative(text, recorded))
3126}
3127
3128fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3132 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3133 let mut rest = rest.trim();
3134 if rest.is_empty() {
3135 return None;
3136 }
3137 let mut total: i64 = 0;
3138 let mut matched = false;
3139 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3140 if let Some((digits, tail)) = rest.split_once(unit)
3141 && !digits.is_empty()
3142 && digits.bytes().all(|b| b.is_ascii_digit())
3143 {
3144 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3145 rest = tail;
3146 matched = true;
3147 }
3148 }
3149 if !rest.is_empty() || !matched {
3150 return None;
3151 }
3152 recorded
3153 .checked_add(jiff::SignedDuration::from_secs(total))
3154 .ok()
3155}
3156
3157fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3161 let clock = clock.trim().to_lowercase();
3162 let (digits, pm) = clock
3163 .strip_suffix("am")
3164 .map(|d| (d, false))
3165 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3166 let (h, m) = digits.trim().split_once(':')?;
3167 let mut hour: i8 = h.trim().parse().ok()?;
3168 let minute: i8 = m.trim().parse().ok()?;
3169 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3170 return None;
3171 }
3172 if pm && hour != 12 {
3173 hour += 12;
3174 } else if !pm && hour == 12 {
3175 hour = 0;
3176 }
3177 Some((hour, minute))
3178}
3179
3180fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3185 let open = text.find('(')?;
3186 let close = text.rfind(')')?;
3187 if close <= open {
3188 return None;
3189 }
3190 let zone = text[open + 1..close].trim();
3191 let (hour, minute) = parse_12h_clock(&text[..open])?;
3192 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3193 let candidate = now
3194 .to_zoned(tz)
3195 .with()
3196 .hour(hour)
3197 .minute(minute)
3198 .second(0)
3199 .millisecond(0)
3200 .microsecond(0)
3201 .nanosecond(0)
3202 .build()
3203 .ok()?;
3204 let mut at = candidate.timestamp();
3205 if at <= now {
3206 at += jiff::SignedDuration::from_hours(24);
3207 }
3208 Some(at)
3209}
3210
3211fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3220 let words: Vec<&str> = text.split_whitespace().collect();
3221 if words.len() < 5 {
3222 return None;
3223 }
3224 (0..=words.len() - 5)
3225 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3226}
3227
3228fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3234 if trailing.is_some_and(|next| next.starts_with('(')) {
3235 return None;
3236 }
3237 let month = month_number(window[0])?;
3238 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3239 let day_digits = ["st", "nd", "rd", "th"]
3240 .iter()
3241 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3242 let day: i8 = day_digits.parse().ok()?;
3243 let year_token = window[2];
3244 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3245 return None;
3246 }
3247 let year: i16 = year_token.parse().ok()?;
3248 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3252 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3253 let date = jiff::civil::Date::new(year, month, day).ok()?;
3254 let candidate = date
3255 .at(hour, minute, 0, 0)
3256 .to_zoned(jiff::tz::TimeZone::UTC)
3257 .ok()?;
3258 Some(candidate.timestamp())
3259}
3260
3261fn month_number(name: &str) -> Option<i8> {
3264 const NAMES: [&str; 12] = [
3265 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3266 ];
3267 let lower = name.to_lowercase();
3268 NAMES
3269 .iter()
3270 .position(|n| *n == lower.as_str())
3271 .map(|i| i as i8 + 1)
3272}
3273
3274fn exhausted_review_budget(state: &RunState) -> bool {
3286 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3287}
3288
3289fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3326 unfinished_run_with(runs, short, RunState::load)
3327}
3328
3329fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3332where
3333 F: FnOnce(&str) -> Result<RunState>,
3334{
3335 let id = runs.last()?;
3336 match load(id) {
3337 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
3338 Some(id.clone())
3339 }
3340 Ok(_) => None,
3341 Err(e) => {
3342 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3343 None
3344 }
3345 }
3346}
3347
3348fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3356where
3357 F: FnOnce(&str) -> Result<RunState>,
3358{
3359 if task.status != TaskStatus::Failed || task.fresh_start || task.review_branch.is_some() {
3360 return false;
3361 }
3362 let Some(id) = task.runs.last() else {
3363 return false;
3364 };
3365 load(id).is_ok_and(|s| {
3366 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3367 })
3368}
3369
3370#[derive(Debug, Clone, PartialEq, Eq)]
3373enum Starter {
3374 Review(String),
3377 Resume(String),
3379 Start,
3381}
3382
3383fn take_divergence_answer(
3387 branch: &str,
3388 remote: &str,
3389 task: &mut Task,
3390) -> Option<crate::reconcile::Choice> {
3391 let summary = crate::reconcile::summary_for(branch, remote);
3392 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3393 (a.question == summary)
3394 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3395 .flatten()
3396 .map(|c| (i, c))
3397 })?;
3398 task.answers.remove(idx);
3399 Some(choice)
3400}
3401
3402fn choose_starter(
3414 review_branch: Option<&str>,
3415 branch_exists: bool,
3416 unfinished: Option<&str>,
3417) -> Starter {
3418 match review_branch {
3419 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3420 Some(_) => Starter::Start,
3421 None => match unfinished {
3422 Some(id) => Starter::Resume(id.to_owned()),
3423 None => Starter::Start,
3424 },
3425 }
3426}
3427
3428fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3431 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3432 return fallback.to_path_buf();
3433 }
3434 task.repo.clone()
3435}
3436
3437const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3441
3442fn answers_block(task: &Task, count: usize) -> String {
3444 let mut s = ANSWERS_HEADER.to_owned();
3445 for a in &task.answers[..count] {
3446 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3447 }
3448 s
3449}
3450
3451fn append_answers(base: &str, task: &Task) -> String {
3454 if task.answers.is_empty() {
3455 return base.to_owned();
3456 }
3457 let mut s = base.to_owned();
3458 s.push_str(&answers_block(task, task.answers.len()));
3459 s
3460}
3461
3462fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3466 for count in (1..=task.answers.len()).rev() {
3467 let block = answers_block(task, count);
3468 if let Some(base) = instruction.strip_suffix(&block) {
3469 return base;
3470 }
3471 }
3472 instruction
3473}
3474
3475fn instruction_for(task: &Task) -> String {
3483 append_answers(&task.instruction, task)
3484}
3485
3486fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3498 append_answers(strip_answers_block(old_instruction, task), task)
3499}
3500
3501fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3503 let paths = queue.attachment_paths(task);
3504 for (name, path) in task.attachments.iter().zip(&paths) {
3505 if !path.is_file() {
3506 bail!(
3507 "attachment `{name}` is recorded on the task but {} is missing",
3508 path.display()
3509 );
3510 }
3511 }
3512 Ok(paths)
3513}
3514
3515fn prepare_instruction(
3526 starter: &Starter,
3527 old_instruction: Option<&str>,
3528 task: &Task,
3529) -> Option<String> {
3530 match starter {
3531 Starter::Start => Some(instruction_for(task)),
3532 Starter::Resume(_) => Some(resumed_instruction(
3533 old_instruction.expect("a resumed run always has a prior instruction"),
3534 task,
3535 )),
3536 Starter::Review(_) => None,
3537 }
3538}
3539
3540fn record(queue: &Queue, task: &mut Task) {
3544 if let Err(e) = queue.put(task) {
3545 tracing::error!("could not record task {}: {e:#}", task.short());
3546 notices::raise(Notice::error(
3547 "loop:record",
3548 "The loop could not save a task's state; check the disk.",
3549 ));
3550 }
3551}
3552
3553fn runnable(queue: &Queue) -> Vec<Task> {
3559 let mut tasks: Vec<Task> = queue
3560 .list()
3561 .into_iter()
3562 .filter(|t| t.status.runnable())
3563 .collect();
3564 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
3565 tasks
3566}
3567
3568fn describe(state: &RunState) -> String {
3582 let p = phrases(&state.config.graph.language);
3583 let mut detail = if state.status == RunStatus::Stalled {
3584 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
3585 seats.sort_unstable();
3586 seats.dedup();
3587 if seats.is_empty() {
3588 p.quorum_lost.to_owned()
3589 } else {
3590 format!("{}{}{}", p.quorum_lost, p.quota_took_out, seats.join(", "))
3591 }
3592 } else {
3593 format!("{}{}", p.run_ended, state.status.display_label())
3594 };
3595 if let Some(last) = state.events.last() {
3596 detail.push_str(&format!(" ({}: {})", last.node, last.message));
3597 }
3598 detail.push_str(&format!(" [run {}]", state.id));
3599 detail
3600}
3601
3602const DIAGNOSTIC_MAX: usize = 4_000;
3608
3609const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
3614
3615fn diagnostic(state: &RunState) -> Option<String> {
3629 let mut parts: Vec<String> = Vec::new();
3630
3631 for o in state.gate.iter().filter(|o| !o.ok()) {
3633 parts.push(format!(
3634 "gate `{}` failed ({:?}):\n{}",
3635 o.command,
3636 o.code,
3637 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
3638 ));
3639 }
3640
3641 if let Some(last) = state
3644 .events
3645 .iter()
3646 .rev()
3647 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
3648 {
3649 parts.push(last.message.clone());
3650 }
3651
3652 if state.viable().is_empty() {
3659 for c in &state.candidates {
3660 if let Some(evidence) = &c.verified_noop {
3661 parts.push(format!(
3662 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
3663 c.label
3664 ));
3665 } else if !c.summary.trim().is_empty() {
3666 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
3667 } else if let Some(why) = &c.failed {
3668 parts.push(format!("candidate {}: {why}", c.label));
3669 }
3670 }
3671 }
3672
3673 if parts.is_empty() {
3674 return None;
3675 }
3676 Some(crate::run::tail(
3681 &parts.join("\n\n"),
3682 DIAGNOSTIC_MAX.saturating_sub(100),
3683 ))
3684}
3685
3686fn label(status: RunStatus) -> &'static str {
3694 status.as_str()
3695}
3696
3697fn merge_mode(mode: &str) -> Result<MergeMode> {
3699 match mode {
3700 "none" => Ok(MergeMode::None),
3701 "local" => Ok(MergeMode::Local),
3702 "pr" => Ok(MergeMode::Pr),
3703 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
3704 }
3705}
3706
3707fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
3712 mutex
3713 .lock()
3714 .unwrap_or_else(std::sync::PoisonError::into_inner)
3715}
3716
3717#[cfg(test)]
3718mod tests {
3719 use super::*;
3720 use crate::queue::{Source, TaskStatus};
3721 use crate::run::{Candidate, CommandOutcome};
3722 use pretty_assertions::assert_eq;
3723
3724 fn task() -> Task {
3725 Task::new(
3726 "add retries".to_owned(),
3727 "add retries".to_owned(),
3728 PathBuf::from("/repo"),
3729 Source::Human,
3730 )
3731 }
3732
3733 fn interrupt_task(id: &str) -> Task {
3736 let mut t = task();
3737 t.id = id.to_owned();
3738 t.interrupt = true;
3739 t
3740 }
3741
3742 fn task_with_id(id: &str) -> Task {
3744 let mut t = task();
3745 t.id = id.to_owned();
3746 t
3747 }
3748
3749 fn urgent_task(id: &str) -> Task {
3751 let mut t = task();
3752 t.id = id.to_owned();
3753 t.urgent = true;
3754 t
3755 }
3756
3757 #[test]
3762 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
3763 assert_eq!(permit_kind(true, false), PermitKind::None);
3764 assert_eq!(permit_kind(true, true), PermitKind::None);
3765 }
3766
3767 #[test]
3772 fn permit_kind_separates_urgent_from_ordinary() {
3773 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
3774 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
3775 }
3776
3777 #[test]
3784 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
3785 let cfg = Config::default();
3786 let repo = Path::new("/any/repo/path");
3787
3788 let reason =
3789 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
3790 assert!(reason.contains("1024"), "{reason}");
3791 assert!(
3792 reason.contains(&cfg.disk.min_free_bytes.to_string()),
3793 "{reason}"
3794 );
3795
3796 assert_eq!(
3797 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
3798 None,
3799 "exactly at the floor is open"
3800 );
3801 assert_eq!(
3802 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
3803 None,
3804 "comfortably above the floor is open"
3805 );
3806 }
3807
3808 #[test]
3809 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
3810 let mut cfg = Config::default();
3811 cfg.disk.min_free_bytes = 0;
3812 let repo = Path::new("/any/repo/path");
3813 assert_eq!(
3814 disk_gate_with(repo, &cfg, |_| Ok(0)),
3815 None,
3816 "a zero floor never measures at all"
3817 );
3818 }
3819
3820 #[test]
3821 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
3822 let cfg = Config::default();
3823 let repo = Path::new("/any/repo/path");
3824 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
3825 .expect("a measurement failure must close the gate, not open it");
3826 assert!(reason.contains("could not measure"), "{reason}");
3827 }
3828
3829 #[test]
3830 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
3831 let ordinary = task();
3832 let next = advance_interrupt(
3833 Interrupt::Idle,
3834 std::slice::from_ref(&ordinary.id),
3835 std::slice::from_ref(&ordinary),
3836 );
3837 assert_eq!(next, Interrupt::Idle);
3838 }
3839
3840 #[test]
3841 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
3842 let marked = interrupt_task("marked");
3845 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
3846 assert_eq!(next, Interrupt::Idle);
3847 }
3848
3849 #[test]
3850 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
3851 let marked = interrupt_task("marked");
3852 let next = advance_interrupt(
3853 Interrupt::Idle,
3854 &["running".to_owned()],
3855 std::slice::from_ref(&marked),
3856 );
3857 assert_eq!(
3858 next,
3859 Interrupt::Parking {
3860 parked: vec!["running".to_owned()],
3861 interrupt_task: "marked".to_owned(),
3862 }
3863 );
3864 }
3865
3866 #[test]
3875 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
3876 let marked = interrupt_task("marked");
3877
3878 let two = advance_interrupt(
3879 Interrupt::Idle,
3880 &["a".to_owned(), "b".to_owned()],
3881 std::slice::from_ref(&marked),
3882 );
3883 assert_eq!(two, Interrupt::Idle);
3884
3885 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
3886 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
3887 }
3888
3889 #[test]
3890 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
3891 let state = Interrupt::Parking {
3892 parked: vec!["running".to_owned()],
3893 interrupt_task: "marked".to_owned(),
3894 };
3895 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
3897 assert_eq!(still_going, state);
3898
3899 let stopped_but_not_yet_dispatched =
3903 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
3904 assert_eq!(stopped_but_not_yet_dispatched, state);
3905
3906 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
3908 assert_eq!(
3909 dispatched,
3910 Interrupt::Running {
3911 parked: vec!["running".to_owned()],
3912 interrupt_task: "marked".to_owned(),
3913 }
3914 );
3915 }
3916
3917 #[test]
3918 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
3919 let state = Interrupt::Running {
3920 parked: vec!["running".to_owned()],
3921 interrupt_task: "marked".to_owned(),
3922 };
3923 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
3924 assert_eq!(still_running, state);
3925
3926 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
3933 assert_eq!(
3934 ended,
3935 Interrupt::Resuming {
3936 parked: vec!["running".to_owned()]
3937 }
3938 );
3939 }
3940
3941 #[test]
3942 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
3943 let state = Interrupt::Resuming {
3944 parked: vec!["running".to_owned()],
3945 };
3946 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
3947 assert_eq!(still_waiting, state);
3948
3949 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
3950 assert_eq!(dispatched, Interrupt::Idle);
3951 }
3952
3953 #[test]
3959 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
3960 {
3961 let state = Interrupt::Parking {
3962 parked: vec!["running".to_owned()],
3963 interrupt_task: "marked".to_owned(),
3964 };
3965 let next = advance_interrupt(state, &[], &[]);
3968 assert_eq!(
3969 next,
3970 Interrupt::Resuming {
3971 parked: vec!["running".to_owned()]
3972 },
3973 "abandoning the interrupt must not abandon the resume it owes"
3974 );
3975 }
3976
3977 #[test]
3980 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
3981 let state = Interrupt::Resuming {
3982 parked: vec!["running".to_owned()],
3983 };
3984 let next = advance_interrupt(state, &[], &[]);
3985 assert_eq!(
3986 next,
3987 Interrupt::Idle,
3988 "nothing is left to wait for; the loop must not stay wedged"
3989 );
3990 }
3991
3992 #[test]
3993 fn disabled_by_config_the_sequence_can_never_leave_idle() {
3994 let marked = interrupt_task("marked");
3995 let next = advance_interrupt_tick(
3996 false,
3997 Interrupt::Idle,
3998 &["running".to_owned()],
3999 std::slice::from_ref(&marked),
4000 );
4001 assert_eq!(
4002 next,
4003 Interrupt::Idle,
4004 "an unmarked, unconfigured daemon must behave exactly as before"
4005 );
4006 }
4007
4008 #[test]
4009 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
4010 let state = Interrupt::Parking {
4011 parked: vec!["running".to_owned()],
4012 interrupt_task: "marked".to_owned(),
4013 };
4014 let candidates = vec![interrupt_task("marked"), task()];
4015 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4016 assert!(
4017 allowed.is_empty(),
4018 "nothing may dispatch - not even the interrupt task itself - \
4019 until the parked run has actually stopped"
4020 );
4021 }
4022
4023 #[test]
4037 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
4038 for state in [
4039 Interrupt::Parking {
4040 parked: vec!["running".to_owned()],
4041 interrupt_task: "marked".to_owned(),
4042 },
4043 Interrupt::Running {
4044 parked: vec!["running".to_owned()],
4045 interrupt_task: "marked".to_owned(),
4046 },
4047 Interrupt::Resuming {
4048 parked: vec!["running".to_owned()],
4049 },
4050 ] {
4051 let candidates = vec![
4052 interrupt_task("marked"),
4053 urgent_task("hot"),
4054 task_with_id("ordinary"),
4055 ];
4056 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
4057 assert!(
4058 !allowed.iter().any(|t| t.id == "hot"),
4059 "an urgent candidate must wait out the same gate as anything \
4060 else while the run it would run alongside has not actually \
4061 left flight, for state {state:?}: {allowed:?}"
4062 );
4063 }
4064 }
4065
4066 #[test]
4072 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
4073 let state = Interrupt::Resuming {
4074 parked: vec!["hot".to_owned()],
4075 };
4076 let candidates = vec![urgent_task("hot"), task()];
4077 let allowed = interrupt_gate(&state, &[], candidates);
4078 assert_eq!(
4079 allowed.iter().filter(|t| t.id == "hot").count(),
4080 1,
4081 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
4082 );
4083 }
4084
4085 #[test]
4086 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
4087 let state = Interrupt::Parking {
4088 parked: vec!["running".to_owned()],
4089 interrupt_task: "marked".to_owned(),
4090 };
4091 let other = task();
4092 let candidates = vec![interrupt_task("marked"), other.clone()];
4093 let allowed = interrupt_gate(&state, &[], candidates);
4094 assert_eq!(allowed.len(), 1);
4095 assert_eq!(allowed[0].id, "marked");
4096 }
4097
4098 #[test]
4099 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4100 let state = Interrupt::Running {
4101 parked: vec!["running".to_owned()],
4102 interrupt_task: "marked".to_owned(),
4103 };
4104 let candidates = vec![task(), task()];
4105 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4106 assert!(allowed.is_empty());
4107 }
4108
4109 #[test]
4116 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4117 let state = Interrupt::Resuming {
4118 parked: vec!["a".to_owned(), "c".to_owned()],
4119 };
4120 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4121 let allowed = interrupt_gate(&state, &[], candidates);
4122 assert_eq!(
4123 allowed.len(),
4124 1,
4125 "at most one candidate may be offered while resuming: {allowed:?}"
4126 );
4127 assert_eq!(allowed[0].id, "a");
4128 }
4129
4130 #[test]
4131 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4132 let state = Interrupt::Resuming {
4133 parked: vec!["a".to_owned()],
4134 };
4135 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4136 assert!(allowed.is_empty());
4137 }
4138
4139 #[test]
4145 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4146 let running = task(); let marked = interrupt_task("marked");
4148
4149 let mut state = Interrupt::Idle;
4150 let in_flight = vec![running.id.clone()];
4152 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4153 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4154 assert!(gated.is_empty(), "still waiting on `running` to park");
4155
4156 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4158 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4159 assert_eq!(
4160 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4161 vec!["marked"],
4162 "only the interrupt task may be offered to the dispatcher now"
4163 );
4164
4165 state = advance_interrupt_tick(
4167 true,
4168 state,
4169 &["marked".to_owned()],
4170 std::slice::from_ref(&running),
4171 );
4172 let gated = interrupt_gate(
4173 &state,
4174 &["marked".to_owned()],
4175 vec![marked.clone(), running.clone()],
4176 );
4177 assert!(
4178 gated.is_empty(),
4179 "the parked run must not be offered back while the interrupt \
4180 task is still running"
4181 );
4182
4183 let other = task_with_id("other");
4187 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4188 assert_eq!(
4189 state,
4190 Interrupt::Resuming {
4191 parked: vec![running.id.clone()]
4192 }
4193 );
4194 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4195 assert_eq!(
4196 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4197 vec![running.id.as_str()],
4198 "exactly the parked run resumes - not the unrelated task, even \
4199 though it was offered first"
4200 );
4201
4202 state = advance_interrupt_tick(
4206 true,
4207 state,
4208 std::slice::from_ref(&running.id),
4209 std::slice::from_ref(&other),
4210 );
4211 assert_eq!(state, Interrupt::Idle);
4212 let gated = interrupt_gate(
4213 &state,
4214 std::slice::from_ref(&running.id),
4215 vec![other.clone()],
4216 );
4217 assert_eq!(
4218 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4219 vec![other.id.as_str()],
4220 "ordinary dispatch is unrestricted again"
4221 );
4222 }
4223
4224 #[test]
4225 fn every_run_status_settles_the_task_it_came_from() {
4226 let table = [
4228 (RunStatus::Merged, TaskStatus::Done, 1),
4229 (RunStatus::Ready, TaskStatus::Done, 1),
4230 (RunStatus::Stalled, TaskStatus::Failed, 0),
4231 (RunStatus::Blocked, TaskStatus::Failed, 1),
4232 (RunStatus::Failed, TaskStatus::Failed, 1),
4233 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4234 (RunStatus::Prep, TaskStatus::Failed, 1),
4235 (RunStatus::Implementing, TaskStatus::Failed, 1),
4236 (RunStatus::Judging, TaskStatus::Failed, 1),
4237 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4238 (RunStatus::Voting, TaskStatus::Failed, 1),
4239 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4240 (RunStatus::Gating, TaskStatus::Failed, 1),
4241 ];
4242 for (run, want, attempts) in table {
4243 let mut t = task();
4244 t.start("20260902-000000-aaaa".to_owned());
4245 settle(
4246 &mut t,
4247 Verdict {
4248 status: run,
4249 left_pr: false,
4250 parked: false,
4251 quota_hit: matches!(run, RunStatus::Stalled),
4252 no_viable_candidates: false,
4253 },
4254 "why",
4255 2,
4256 );
4257 assert_eq!(t.status, want, "task status after {}", label(run));
4258 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4259 }
4260 }
4261
4262 #[test]
4263 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4264 let mut stalled = task();
4265 stalled.start("20260902-000000-aaaa".to_owned());
4266 settle(
4267 &mut stalled,
4268 Verdict {
4269 status: RunStatus::Stalled,
4270 left_pr: false,
4271 parked: false,
4272 quota_hit: true,
4273 no_viable_candidates: false,
4274 },
4275 "quota",
4276 1,
4277 );
4278 assert_eq!(stalled.attempts, 0);
4279 assert!(
4280 stalled.status.runnable(),
4281 "a machine problem must leave the task in line"
4282 );
4283
4284 let mut blocked = task();
4285 blocked.start("20260902-000000-aaaa".to_owned());
4286 settle(
4287 &mut blocked,
4288 Verdict {
4289 status: RunStatus::Blocked,
4290 left_pr: false,
4291 parked: false,
4292 quota_hit: false,
4293 no_viable_candidates: false,
4294 },
4295 "findings open",
4296 1,
4297 );
4298 assert_eq!(blocked.attempts, 1);
4299 assert_eq!(
4300 blocked.status,
4301 TaskStatus::Held,
4302 "the last attempt hands the task to a human"
4303 );
4304 }
4305
4306 #[test]
4307 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4308 let mut delivered = task();
4311 delivered.start("20260903-080619-01c2".to_owned());
4312 settle(
4313 &mut delivered,
4314 Verdict {
4315 status: RunStatus::Blocked,
4316 left_pr: true,
4317 parked: false,
4318 quota_hit: false,
4319 no_viable_candidates: false,
4320 },
4321 "no check status",
4322 4,
4323 );
4324 assert_eq!(
4325 delivered.status,
4326 TaskStatus::Held,
4327 "a pull request waiting on CI or a person is not a retryable failure"
4328 );
4329 assert!(
4330 !delivered.status.runnable(),
4331 "the loop must not pick this task up again"
4332 );
4333 assert_eq!(
4334 delivered.last_error.as_deref(),
4335 Some("no check status"),
4336 "the operator needs to be told what the gate was waiting for"
4337 );
4338
4339 let mut empty_handed = task();
4342 empty_handed.start("20260903-080619-01c2".to_owned());
4343 settle(
4344 &mut empty_handed,
4345 Verdict {
4346 status: RunStatus::Blocked,
4347 left_pr: false,
4348 parked: false,
4349 quota_hit: false,
4350 no_viable_candidates: false,
4351 },
4352 "findings open",
4353 4,
4354 );
4355 assert_eq!(empty_handed.status, TaskStatus::Failed);
4356 assert!(empty_handed.status.runnable());
4357 }
4358
4359 #[test]
4360 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4361 let mut noop = task();
4367 noop.start("20260912-131304-391f".to_owned());
4368 settle(
4369 &mut noop,
4370 Verdict {
4371 status: RunStatus::VerifiedNoop,
4372 left_pr: false,
4373 parked: false,
4374 quota_hit: false,
4375 no_viable_candidates: true,
4376 },
4377 "candidate A: already fixed by b32cfc4, on main",
4378 4,
4379 );
4380 assert_eq!(
4381 noop.status,
4382 TaskStatus::Held,
4383 "an unverified claim is a request for a human, not a failure"
4384 );
4385 assert!(
4386 !noop.status.runnable(),
4387 "the loop must not requeue this on the same unverified claim"
4388 );
4389 assert_eq!(noop.attempts, 1);
4394 }
4395
4396 #[test]
4397 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4398 let mut parked = task();
4403 parked.start("20260903-183634-2d98".to_owned());
4404 settle(
4405 &mut parked,
4406 Verdict {
4407 status: RunStatus::Implementing,
4408 left_pr: false,
4409 quota_hit: false,
4410 parked: true,
4411 no_viable_candidates: false,
4412 },
4413 "parked after `implementing`",
4414 2,
4415 );
4416 assert_eq!(parked.attempts, 0, "a park is refunded");
4417 assert!(
4418 parked.status.runnable(),
4419 "and the task stays in line so the next loop resumes its run"
4420 );
4421 assert_eq!(
4422 parked.last_error.as_deref(),
4423 Some("parked after `implementing`"),
4424 "the card says where it stopped"
4425 );
4426
4427 let mut broken = task();
4431 broken.start("20260903-183634-2d98".to_owned());
4432 settle(
4433 &mut broken,
4434 Verdict {
4435 status: RunStatus::Implementing,
4436 left_pr: false,
4437 quota_hit: false,
4438 parked: false,
4439 no_viable_candidates: false,
4440 },
4441 "returned mid-flight",
4442 2,
4443 );
4444 assert_eq!(broken.attempts, 1);
4445 }
4446
4447 #[test]
4448 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4449 let mut flaky = task();
4454 flaky.start("20260903-123023-e633".to_owned());
4455 settle(
4456 &mut flaky,
4457 Verdict {
4458 status: RunStatus::Stalled,
4459 left_pr: false,
4460 parked: false,
4461 quota_hit: false,
4462 no_viable_candidates: false,
4463 },
4464 "verdict rests on 1 of 3 judges",
4465 2,
4466 );
4467 assert_eq!(
4468 flaky.attempts, 1,
4469 "flakiness spends an attempt, so `max_attempts` still bounds it"
4470 );
4471 assert!(flaky.status.runnable(), "and it is still worth retrying");
4472
4473 let mut limited = task();
4475 limited.start("20260903-123023-e633".to_owned());
4476 settle(
4477 &mut limited,
4478 Verdict {
4479 status: RunStatus::Stalled,
4480 left_pr: false,
4481 parked: false,
4482 quota_hit: true,
4483 no_viable_candidates: false,
4484 },
4485 "judge-2, judge-3 out of quota",
4486 2,
4487 );
4488 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4489 assert!(limited.status.runnable());
4490
4491 let mut worn = task();
4494 for _ in 0..2 {
4495 worn.release();
4496 }
4497 worn.start("20260903-123023-e633".to_owned());
4498 worn.attempts = 2;
4499 settle(
4500 &mut worn,
4501 Verdict {
4502 status: RunStatus::Stalled,
4503 left_pr: false,
4504 parked: false,
4505 quota_hit: false,
4506 no_viable_candidates: false,
4507 },
4508 "no quorum again",
4509 2,
4510 );
4511 assert_eq!(worn.status, TaskStatus::Held);
4512 assert!(!worn.status.runnable());
4513 }
4514
4515 #[test]
4516 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4517 let mut wiped_out = task();
4524 wiped_out.start("20260907-025000-a1b2".to_owned());
4525 settle(
4526 &mut wiped_out,
4527 Verdict {
4528 status: RunStatus::Failed,
4529 left_pr: false,
4530 parked: false,
4531 quota_hit: true,
4532 no_viable_candidates: true,
4533 },
4534 "no candidate produced a change; nothing to judge",
4535 2,
4536 );
4537 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
4538 assert!(
4539 wiped_out.status.runnable(),
4540 "a machine problem must leave the task in line"
4541 );
4542
4543 let mut partial_progress = task();
4549 partial_progress.start("20260907-025500-c3d4".to_owned());
4550 settle(
4551 &mut partial_progress,
4552 Verdict {
4553 status: RunStatus::Failed,
4554 left_pr: false,
4555 parked: false,
4556 quota_hit: true,
4557 no_viable_candidates: false,
4558 },
4559 "gate failed on the winning candidate",
4560 2,
4561 );
4562 assert_eq!(
4563 partial_progress.attempts, 1,
4564 "a candidate that actually produced a change spends the attempt \
4565 even though some other seat hit its quota"
4566 );
4567 assert!(partial_progress.status.runnable());
4568 }
4569
4570 #[test]
4571 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
4572 let mut t = task();
4578 t.start("20260907-025000-a1b2".to_owned());
4579 let mut state = run_state(RunStatus::Failed);
4580 state.quota.push(QuotaLoss {
4581 seat: "cand-a".to_owned(),
4582 node: "implement".to_owned(),
4583 at: Timestamp::now(),
4584 reset: None,
4585 });
4586 assert!(
4587 state.viable().is_empty(),
4588 "no candidate was added, so nothing is viable"
4589 );
4590 reclaim(&mut t, Some(state), 2, "en");
4591 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
4592 assert!(t.status.runnable());
4593 }
4594
4595 #[test]
4596 fn a_held_task_is_never_offered_to_the_loop() {
4597 let dir = tempfile::tempdir().unwrap();
4598 let queue = Queue::at(dir.path().to_path_buf());
4599 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
4600 let mut t = task();
4601 t.id = format!("2026090{n}-000000-000{n}");
4602 t.priority = priority;
4603 queue.put(&mut t).unwrap();
4604 }
4605 let mut held = task();
4606 held.id = "20260909-000000-9999".to_owned();
4607 held.priority = 99;
4608 held.hold_machine(None);
4609 queue.put(&mut held).unwrap();
4610
4611 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
4612 assert_eq!(order.len(), 3);
4613 assert!(!order.contains(&held.id));
4614 assert_eq!(
4615 order.first().cloned(),
4616 queue.next_runnable().map(|t| t.id),
4617 "the loop's first candidate is exactly what the queue offers"
4618 );
4619 assert_eq!(
4620 order,
4621 vec![
4622 "20260902-000000-0002".to_owned(),
4623 "20260903-000000-0003".to_owned(),
4624 "20260901-000000-0001".to_owned(),
4625 ],
4626 "priority first, then oldest, so nothing starves"
4627 );
4628 }
4629
4630 #[test]
4631 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
4632 let dir = tempfile::tempdir().unwrap();
4633 let queue = Queue::at(dir.path().to_path_buf());
4634 let mut old = task();
4635 old.id = "20260101-000000-old0".to_owned();
4636 queue.put(&mut old).unwrap();
4637 let mut fresh = task();
4638 fresh.id = "20260101-000000-new0".to_owned();
4639 queue.put(&mut fresh).unwrap();
4640
4641 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
4645 std::thread::sleep(Duration::from_millis(60));
4646 let live = queue.claim(&fresh.id).unwrap();
4647
4648 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4649 assert_eq!(swept, vec![old.id.clone()]);
4650 assert!(
4651 queue.claim(&old.id).is_ok(),
4652 "an unparseable lock older than the threshold is swept"
4653 );
4654 assert!(
4655 queue.claim(&fresh.id).is_err(),
4656 "a live pid protects its lock regardless of age"
4657 );
4658 drop(live);
4659 }
4660
4661 #[test]
4662 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
4663 let dir = tempfile::tempdir().unwrap();
4673 let queue = Queue::at(dir.path().to_path_buf());
4674 let mut t = task();
4675 t.id = "20260101-000000-live".to_owned();
4676 queue.put(&mut t).unwrap();
4677
4678 let claim = queue.claim(&t.id).unwrap();
4679 std::thread::sleep(Duration::from_millis(60));
4680
4681 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4682 assert!(
4683 swept.is_empty(),
4684 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
4685 );
4686 assert!(
4687 queue.claim(&t.id).is_err(),
4688 "the lock still protects its task"
4689 );
4690 drop(claim);
4691 }
4692
4693 fn injected_dead_pid() -> u32 {
4696 std::process::id().checked_add(1).unwrap_or(1)
4697 }
4698
4699 #[test]
4700 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
4701 let dir = tempfile::tempdir().unwrap();
4702 let queue = Queue::at(dir.path().to_path_buf());
4703 let mut t = task();
4704 t.id = "20260101-000000-dead".to_owned();
4705 queue.put(&mut t).unwrap();
4706 let dead_pid = injected_dead_pid();
4707
4708 std::fs::write(
4713 dir.path().join(format!("{}.lock", t.id)),
4714 dead_pid.to_string(),
4715 )
4716 .unwrap();
4717
4718 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4719 pid != dead_pid
4720 });
4721 assert_eq!(
4722 swept,
4723 vec![t.id.clone()],
4724 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
4725 );
4726 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
4727 }
4728
4729 #[test]
4730 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
4731 let dir = tempfile::tempdir().unwrap();
4732 let queue = Queue::at(dir.path().to_path_buf());
4733 let mut t = task();
4734 t.id = "20260101-000000-late".to_owned();
4735 queue.put(&mut t).unwrap();
4736 let dead_pid = injected_dead_pid();
4737
4738 assert!(
4741 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
4742 "nothing has claimed the task yet"
4743 );
4744
4745 std::fs::write(
4748 dir.path().join(format!("{}.lock", t.id)),
4749 dead_pid.to_string(),
4750 )
4751 .unwrap();
4752
4753 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4757 pid != dead_pid
4758 });
4759 assert_eq!(swept, vec![t.id.clone()]);
4760 }
4761
4762 #[test]
4763 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
4764 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
4769 let dir = tempfile::tempdir().unwrap();
4770 let queue = Queue::at(dir.path().to_path_buf());
4771 let mut t = task();
4772 t.id = "20260101-000000-crsh".to_owned();
4773 t.status = TaskStatus::Running;
4774 t.attempts = 1;
4775 t.runs.push("20260904-000000-4043".to_owned());
4779 queue.put(&mut t).unwrap();
4780 let dead_pid = injected_dead_pid();
4781
4782 std::fs::write(
4785 dir.path().join(format!("{}.lock", t.id)),
4786 dead_pid.to_string(),
4787 )
4788 .unwrap();
4789
4790 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
4796 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
4797
4798 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4799 pid != dead_pid
4800 });
4801 assert_eq!(swept, vec![t.id.clone()]);
4802
4803 let reclaimed = reclaim_orphaned_running(&queue, 2);
4804 assert_eq!(reclaimed, vec![t.id.clone()]);
4805 let after = queue.get(&t.id).unwrap();
4806 assert_eq!(
4807 after.status,
4808 TaskStatus::Held,
4809 "no run.json to recover from, so a human is asked"
4810 );
4811 assert_eq!(
4812 after.runs,
4813 vec!["20260904-000000-4043".to_owned()],
4814 "the crashed run's id is kept as evidence, not discarded"
4815 );
4816 }
4817
4818 #[test]
4819 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
4820 let dir = tempfile::tempdir().unwrap();
4821 let queue = Queue::at(dir.path().to_path_buf());
4822 let mut t = task();
4823 t.id = "20260101-000000-unknown".to_owned();
4824 queue.put(&mut t).unwrap();
4825 let dead_pid = injected_dead_pid();
4826 std::fs::write(
4827 dir.path().join(format!("{}.lock", t.id)),
4828 dead_pid.to_string(),
4829 )
4830 .unwrap();
4831
4832 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
4833 assert!(swept.is_empty(), "an unknown pid must keep its lock");
4834 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
4835 }
4836
4837 fn run_state_in(status: RunStatus, language: &str) -> RunState {
4838 let mut s = run_state(status);
4839 s.config.graph.language = language.to_owned();
4840 s
4841 }
4842
4843 fn unstarted_verdict(status: RunStatus) -> Verdict {
4844 Verdict {
4845 status,
4846 left_pr: false,
4847 quota_hit: false,
4848 parked: false,
4849 no_viable_candidates: false,
4850 }
4851 }
4852
4853 #[test]
4854 fn settle_renders_the_non_terminal_reason_in_the_configured_language() {
4855 let reason = |language: &str| {
4856 let mut t = task();
4857 settle_in(
4858 &mut t,
4859 unstarted_verdict(RunStatus::Judging),
4860 "boom",
4861 1,
4862 phrases(language),
4863 );
4864 t.last_error.or(t.hold_reason).unwrap_or_default()
4865 };
4866 assert!(
4867 reason("en").starts_with("the graph stopped at `"),
4868 "{}",
4869 reason("en")
4870 );
4871 assert!(reason("ja").starts_with("グラフが終端状態に達しないまま"));
4872 assert!(reason("日本語").contains("boom"));
4873 assert_eq!(reason("fr"), reason("en"));
4874 }
4875
4876 #[test]
4877 fn describe_follows_the_run_language_and_keeps_the_run_id() {
4878 let en = describe(&run_state_in(RunStatus::Stalled, "en"));
4879 assert!(en.starts_with("the judging panel lost its quorum"), "{en}");
4880 let ja = describe(&run_state_in(RunStatus::Stalled, "ja"));
4881 assert!(ja.starts_with("審査パネルが定足数を失いました"), "{ja}");
4882 assert!(ja.contains("[run "), "{ja}");
4883 let ended = describe(&run_state_in(RunStatus::Failed, "jp"));
4884 assert!(ended.starts_with("run 終了: "), "{ended}");
4885 let mut de = run_state_in(RunStatus::Failed, "de");
4886 let mut en = run_state_in(RunStatus::Failed, "en");
4887 de.id = "same".to_owned();
4888 en.id = "same".to_owned();
4889 assert_eq!(describe(&de), describe(&en));
4890 }
4891
4892 #[test]
4893 fn refusals_and_recovery_prose_follow_the_language() {
4894 let t = held_task_with("r1");
4895 let q = action_question(
4896 "r1",
4897 ask::ChoiceAction::Resume {
4898 run: "r1".to_owned(),
4899 },
4900 );
4901 let refuse = |p: &Phrases| match decide_action(&t, &q, p, |_: &str| bail!("gone")) {
4902 ActionDecision::Refuse(s) => s,
4903 other => panic!("{other:?}"),
4904 };
4905 assert!(refuse(phrases("en")).contains("could not be read: gone"));
4906 assert!(refuse(phrases("ja")).contains("読み込めませんでした: gone"));
4907 assert_eq!(refuse(phrases("fr")), refuse(phrases("en")));
4908
4909 let mut held = task();
4910 reclaim(&mut held, None, 2, "ja");
4911 assert!(held.hold_reason.unwrap().contains("保留にしました"));
4912 let mut held = task();
4913 reclaim(&mut held, None, 2, "xx");
4914 assert!(held.hold_reason.unwrap().contains("held for a human"));
4915 }
4916
4917 fn run_state(status: RunStatus) -> RunState {
4918 let mut state = RunState::new(
4919 PathBuf::from("/repo"),
4920 "main".to_owned(),
4921 "abc1234def".to_owned(),
4922 "add retries".to_owned(),
4923 Config::default(),
4924 );
4925 state.status = status;
4926 state
4927 }
4928
4929 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
4930 Candidate {
4931 index: 0,
4932 label,
4933 agent: "claude".to_owned(),
4934 branch: format!("magi/x/{label}"),
4935 worktree: PathBuf::from("/repo"),
4936 summary: summary.to_owned(),
4937 stat: String::new(),
4938 files: 0,
4939 commits: usize::from(!empty),
4940 empty,
4941 failed: failed.map(str::to_owned),
4942 verified_noop: None,
4943 duration_ms: 0,
4944 folded: false,
4945 }
4946 }
4947
4948 #[test]
4949 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
4950 let mut state = run_state(RunStatus::Blocked);
4951 state.gate = vec![
4952 CommandOutcome {
4953 command: "cargo make check".to_owned(),
4954 code: Some(0),
4955 output_tail: "ok".to_owned(),
4956 duration_ms: 0,
4957 resource_blocked: false,
4958 },
4959 CommandOutcome {
4960 command: "cargo test".to_owned(),
4961 code: Some(101),
4962 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
4963 duration_ms: 0,
4964 resource_blocked: false,
4965 },
4966 ];
4967 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
4968 assert!(d.contains("cargo test"), "{d}");
4969 assert!(
4970 !d.contains("cargo make check"),
4971 "a passing check is not a diagnostic: {d}"
4972 );
4973 assert!(d.contains("assertion failed"), "{d}");
4974 }
4975
4976 #[test]
4977 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
4978 let mut state = run_state(RunStatus::Blocked);
4979 state.event(
4980 "land",
4981 "stopped: the fixer produced no commit while 2 check(s) were failing \
4982 (build, lint); stopping instead of looping on an unchanged tree",
4983 );
4984 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
4985 assert!(d.contains("build"), "{d}");
4986 assert!(d.contains("lint"), "{d}");
4987 assert!(d.contains("fixer produced no commit"), "{d}");
4988 }
4989
4990 #[test]
4991 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
4992 let state = run_state(RunStatus::VerifiedNoop);
4997 let d = describe(&state);
4998 assert!(
4999 d.contains("agent-verified no-op"),
5000 "expected the display label, not the wire spelling: {d}"
5001 );
5002 assert!(!d.contains("verified_noop"), "{d}");
5003 }
5004
5005 #[test]
5006 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
5007 let mut state = run_state(RunStatus::Failed);
5013 state.candidates = vec![candidate(
5014 'A',
5015 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
5016 true,
5017 None,
5018 )];
5019 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
5020 assert!(d.contains("candidate A"), "{d}");
5021 assert!(d.contains("tagged v1.2.3"), "{d}");
5022 }
5023
5024 #[test]
5025 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
5026 let mut state = run_state(RunStatus::Failed);
5027 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
5028 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
5029 assert!(d.contains("candidate A"), "{d}");
5030 assert!(d.contains("agent timed out"), "{d}");
5031 }
5032
5033 #[test]
5034 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
5035 let mut state = run_state(RunStatus::Failed);
5038 state.candidates = vec![candidate('A', "did the work", false, None)];
5039 assert!(diagnostic(&state).is_none());
5040 }
5041
5042 #[test]
5043 fn diagnostic_is_bounded_however_much_a_run_printed() {
5044 let mut state = run_state(RunStatus::Blocked);
5045 state.gate = vec![
5046 CommandOutcome {
5047 command: "cargo test".to_owned(),
5048 code: Some(101),
5049 output_tail: "x".repeat(50_000),
5050 duration_ms: 0,
5051 resource_blocked: false,
5052 },
5053 CommandOutcome {
5054 command: "cargo clippy".to_owned(),
5055 code: Some(1),
5056 output_tail: "y".repeat(50_000),
5057 duration_ms: 0,
5058 resource_blocked: false,
5059 },
5060 ];
5061 state.candidates = vec![
5062 candidate('A', &"z".repeat(50_000), true, None),
5063 candidate('B', &"w".repeat(50_000), true, None),
5064 ];
5065 let d = diagnostic(&state).expect("plenty here to diagnose");
5066 assert!(
5067 d.len() <= DIAGNOSTIC_MAX,
5068 "diagnostic grew to {} bytes, unbounded",
5069 d.len()
5070 );
5071 }
5072
5073 #[test]
5074 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
5075 let mut state = run_state(RunStatus::Blocked);
5076 state.gate = vec![CommandOutcome {
5077 command: "cargo test".to_owned(),
5078 code: Some(101),
5079 output_tail: "assertion failed".to_owned(),
5080 duration_ms: 0,
5081 resource_blocked: false,
5082 }];
5083 let verdict = Verdict {
5084 status: RunStatus::Blocked,
5085 left_pr: false,
5086 quota_hit: false,
5087 parked: false,
5088 no_viable_candidates: false,
5089 };
5090
5091 let mut t = task();
5094 t.start("run-1".to_owned());
5095 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5096 assert_eq!(t.status, TaskStatus::Failed);
5097 assert!(t.diagnostic.is_none());
5098
5099 t.start("run-2".to_owned());
5102 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
5103 assert_eq!(t.status, TaskStatus::Held);
5104 let d = t.diagnostic.expect("a held task must carry its diagnostic");
5105 assert!(d.contains("cargo test"), "{d}");
5106 }
5107
5108 #[test]
5109 fn a_held_task_names_the_open_question_it_is_waiting_on() {
5110 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5115 let home = crate::run::home();
5116 let state = run_state(RunStatus::VerifiedNoop);
5117 let mut q = ask::Question::new(
5118 state.id.clone(),
5119 "implement".to_owned(),
5120 "impl-A".to_owned(),
5121 "is this really a no-op?".to_owned(),
5122 String::new(),
5123 Vec::new(),
5124 );
5125 Questions::at(home.join("questions")).put(&mut q).unwrap();
5126
5127 let verdict = Verdict {
5128 status: RunStatus::VerifiedNoop,
5129 left_pr: false,
5130 quota_hit: false,
5131 parked: false,
5132 no_viable_candidates: false,
5133 };
5134 let mut t = task();
5135 t.start(state.id.clone());
5136 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5137
5138 assert_eq!(t.status, TaskStatus::Held);
5139 let reason = t.hold_reason.expect("a held task must record why");
5140 assert!(
5141 reason.starts_with("run ended agent-verified no-op"),
5142 "the original settle reason must survive unchanged: {reason}"
5143 );
5144 assert!(
5145 reason.contains(q.short()),
5146 "the open question's id must be named so the notice is actionable: {reason}"
5147 );
5148 }
5149
5150 #[test]
5151 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
5152 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5153 let state = run_state(RunStatus::VerifiedNoop);
5154
5155 let verdict = Verdict {
5156 status: RunStatus::VerifiedNoop,
5157 left_pr: false,
5158 quota_hit: false,
5159 parked: false,
5160 no_viable_candidates: false,
5161 };
5162 let mut t = task();
5163 t.start(state.id.clone());
5164 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
5165
5166 assert_eq!(t.status, TaskStatus::Held);
5167 assert_eq!(
5168 t.hold_reason.as_deref(),
5169 Some("run ended agent-verified no-op"),
5170 "nothing to append when the question was already answered or never asked"
5171 );
5172 }
5173
5174 #[test]
5175 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5176 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5177 let mut first = run_state(RunStatus::Blocked);
5178 first.id = "20260101-000000-sup1".to_owned();
5179 first.save().unwrap();
5180 let mut second = run_state(RunStatus::Merged);
5181 second.id = "20260101-000000-sup2".to_owned();
5182 second.save().unwrap();
5183
5184 let mut t = task();
5185 t.runs = vec![first.id.clone(), second.id.clone()];
5186 t.status = TaskStatus::Done;
5187
5188 supersede_prior_runs(&t, &crate::run::home());
5189
5190 assert_eq!(
5191 RunState::load(&first.id).unwrap().status,
5192 RunStatus::Superseded,
5193 "the first attempt's Blocked no longer needs anyone's attention"
5194 );
5195 assert_eq!(
5196 RunState::load(&second.id).unwrap().status,
5197 RunStatus::Merged,
5198 "the run that actually succeeded is left exactly as it was"
5199 );
5200 }
5201
5202 #[test]
5203 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5204 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5210 let mut first = run_state(RunStatus::Blocked);
5211 first.id = "20260101-000000-sup9".to_owned();
5212 first.driver_pid = Some(std::process::id());
5215 first.driver_started_at = Some(
5216 crate::proc::process_started_at(std::process::id())
5217 .expect("this test process's own start time must be queryable"),
5218 );
5219 first.save().unwrap();
5220 let mut second = run_state(RunStatus::Merged);
5221 second.id = "20260101-000000-supa".to_owned();
5222 second.save().unwrap();
5223
5224 let mut t = task();
5225 t.runs = vec![first.id.clone(), second.id.clone()];
5226 t.status = TaskStatus::Done;
5227
5228 supersede_prior_runs(&t, &crate::run::home());
5229
5230 assert_eq!(
5231 RunState::load(&first.id).unwrap().status,
5232 RunStatus::Blocked,
5233 "a live driver_pid means something is still actually working this run, \
5234 even though no daemon claims it - rewriting under it would just be \
5235 undone the next time that process saves"
5236 );
5237 }
5238
5239 #[test]
5240 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5241 let dir = tempfile::tempdir().unwrap();
5247 let home = dir.path().to_path_buf();
5248 let queue = Queue::at(dir.path().join("queue"));
5249
5250 let mut first = run_state(RunStatus::Blocked);
5251 first.id = "20260101-000000-supd".to_owned();
5252 first.driver_pid = Some(std::process::id());
5253 first.driver_started_at = Some(
5254 crate::proc::process_started_at(std::process::id())
5255 .expect("this test process's own start time must be queryable"),
5256 );
5257 first.save_under(&home).unwrap();
5258 let mut second = run_state(RunStatus::Merged);
5259 second.id = "20260101-000000-supe".to_owned();
5260 second.save_under(&home).unwrap();
5261
5262 let mut t = task();
5263 t.runs = vec![first.id.clone(), second.id.clone()];
5264 t.status = TaskStatus::Done;
5265 queue.put(&mut t).unwrap();
5266
5267 resweep_superseded_attempts(&queue, &home);
5268 assert_eq!(
5269 RunState::load_under(&first.id, &home).unwrap().status,
5270 RunStatus::Blocked,
5271 "still live on the first pass, so still untouched"
5272 );
5273
5274 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5280 stale.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
5281 stale.save_under(&home).unwrap();
5282
5283 resweep_superseded_attempts(&queue, &home);
5284 assert_eq!(
5285 RunState::load_under(&first.id, &home).unwrap().status,
5286 RunStatus::Superseded,
5287 "the second pass catches up what the first one correctly skipped"
5288 );
5289 }
5290
5291 #[test]
5292 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5293 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5294 let mut first = run_state(RunStatus::Blocked);
5295 first.id = "20260101-000000-sup3".to_owned();
5296 first.save().unwrap();
5297 let mut second = run_state(RunStatus::Blocked);
5298 second.id = "20260101-000000-sup4".to_owned();
5299 second.save().unwrap();
5300
5301 let mut t = task();
5302 t.runs = vec![first.id.clone(), second.id.clone()];
5303 t.status = TaskStatus::Failed;
5307
5308 supersede_prior_runs(&t, &crate::run::home());
5309
5310 assert_eq!(
5311 RunState::load(&first.id).unwrap().status,
5312 RunStatus::Blocked
5313 );
5314 assert_eq!(
5315 RunState::load(&second.id).unwrap().status,
5316 RunStatus::Blocked
5317 );
5318 }
5319
5320 #[test]
5321 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5322 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5326 let mut first = run_state(RunStatus::Blocked);
5327 first.id = "20260101-000000-sup5".to_owned();
5328 first.save().unwrap();
5329
5330 let mut t = task();
5331 t.runs = vec![first.id.clone()];
5332 t.status = TaskStatus::Done;
5333
5334 supersede_prior_runs(&t, &crate::run::home());
5335
5336 assert_eq!(
5337 RunState::load(&first.id).unwrap().status,
5338 RunStatus::Blocked,
5339 "a single-attempt task has no earlier run to supersede"
5340 );
5341 }
5342
5343 #[test]
5344 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5345 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5352 let mut first = run_state(RunStatus::Blocked);
5353 first.id = "20260101-000000-supb".to_owned();
5354 first.save().unwrap();
5355 let mut second = run_state(RunStatus::Failed);
5356 second.id = "20260101-000000-supc".to_owned();
5357 second.save().unwrap();
5358
5359 let mut t = task();
5360 t.runs = vec![first.id.clone(), second.id.clone()];
5361 t.status = TaskStatus::Done;
5362
5363 supersede_prior_runs(&t, &crate::run::home());
5364
5365 assert_eq!(
5366 RunState::load(&first.id).unwrap().status,
5367 RunStatus::Blocked,
5368 "the task's last attempt never landed, so there is nothing here \
5369 actually superseding it"
5370 );
5371 }
5372
5373 #[test]
5374 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5375 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5379 let mut failed = run_state(RunStatus::Failed);
5380 failed.id = "20260101-000000-sup6".to_owned();
5381 failed.save().unwrap();
5382 let mut noop = run_state(RunStatus::VerifiedNoop);
5383 noop.id = "20260101-000000-sup7".to_owned();
5384 noop.save().unwrap();
5385 let mut winner = run_state(RunStatus::Ready);
5386 winner.id = "20260101-000000-sup8".to_owned();
5387 winner.save().unwrap();
5388
5389 let mut t = task();
5390 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5391 t.status = TaskStatus::Done;
5392
5393 supersede_prior_runs(&t, &crate::run::home());
5394
5395 assert_eq!(
5396 RunState::load(&failed.id).unwrap().status,
5397 RunStatus::Failed
5398 );
5399 assert_eq!(
5400 RunState::load(&noop.id).unwrap().status,
5401 RunStatus::VerifiedNoop
5402 );
5403 }
5404
5405 fn approval_question(run: &str) -> ask::Question {
5406 ask::Question::new(
5407 run.to_owned(),
5408 land::APPROVAL_NODE.to_owned(),
5409 "land".to_owned(),
5410 "merge?".to_owned(),
5411 String::new(),
5412 vec!["merge".to_owned(), "hold".to_owned()],
5413 )
5414 }
5415
5416 #[test]
5417 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5418 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5419 let mut state = run_state(RunStatus::Landing);
5420 state.id = "20260101-000000-fre1".to_owned();
5421 state.parked = true;
5422 state.save().unwrap();
5423 ask::Questions::open()
5424 .put(&mut approval_question(&state.id))
5425 .unwrap();
5426
5427 let mut t = task();
5428 t.runs.push(state.id.clone());
5429 assert_eq!(
5430 land_resume_state(&t),
5431 LandResume::StillWaiting,
5432 "nobody has answered and the timeout has not passed"
5433 );
5434 }
5435
5436 #[test]
5437 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5438 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5443 let mut state = run_state(RunStatus::Landing);
5444 state.id = "20260101-000000-exp1".to_owned();
5445 state.parked = true;
5446 state.config.graph.answer_timeout = 60;
5447 state.save().unwrap();
5448
5449 let store = ask::Questions::open();
5450 let mut q = approval_question(&state.id);
5451 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
5452 store.put(&mut q).unwrap();
5453
5454 let mut t = task();
5455 t.runs.push(state.id.clone());
5456 assert_eq!(
5457 land_resume_state(&t),
5458 LandResume::Ready,
5459 "an expired question must not be waited on forever"
5460 );
5461
5462 let after = store.get(&q.id).unwrap();
5463 assert!(
5464 !after.status.open(),
5465 "the question is abandoned, not silently ignored"
5466 );
5467 assert!(
5468 after.resolution().is_none(),
5469 "an abandoned question is not read as a decision"
5470 );
5471 }
5472
5473 #[test]
5474 fn reclaim_settles_a_running_task_against_its_last_run() {
5475 let mut t = task();
5476 t.start("20260904-000000-4043".to_owned());
5477 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2, "en");
5478 assert_eq!(
5479 t.status,
5480 TaskStatus::Done,
5481 "a run that actually finished must not stay `running` forever"
5482 );
5483 }
5484
5485 #[test]
5486 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
5487 let mut t = task();
5491 t.start("20260904-000000-4043".to_owned());
5492 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2, "en");
5493 assert_eq!(t.status, TaskStatus::Failed);
5494 assert!(t.status.runnable());
5495 }
5496
5497 #[test]
5498 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
5499 let mut t = task();
5500 t.start("20260904-000000-4043".to_owned());
5501 reclaim(&mut t, None, 2, "en");
5502 assert_eq!(t.status, TaskStatus::Held);
5503 assert!(
5504 t.last_error
5505 .as_deref()
5506 .is_some_and(|e| e.contains("running")),
5507 "the operator needs to know why this task was held"
5508 );
5509 }
5510
5511 #[test]
5512 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
5513 let dir = tempfile::tempdir().unwrap();
5514 let queue = Queue::at(dir.path().to_path_buf());
5515
5516 let mut orphaned = task();
5518 orphaned.id = "20260904-000000-orph".to_owned();
5519 orphaned.status = TaskStatus::Running;
5520 orphaned.attempts = 1;
5521 queue.put(&mut orphaned).unwrap();
5522
5523 let mut alive = task();
5524 alive.id = "20260904-000000-live".to_owned();
5525 alive.status = TaskStatus::Running;
5526 alive.attempts = 1;
5527 queue.put(&mut alive).unwrap();
5528 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
5529
5530 let mut queued = task();
5531 queued.id = "20260904-000000-wait".to_owned();
5532 queue.put(&mut queued).unwrap();
5533
5534 let reclaimed = reclaim_orphaned_running(&queue, 2);
5535 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
5536
5537 assert_eq!(
5538 queue.get(&orphaned.id).unwrap().status,
5539 TaskStatus::Held,
5540 "nothing was driving it and there was no run to recover"
5541 );
5542 assert_eq!(
5543 queue.get(&alive.id).unwrap().status,
5544 TaskStatus::Running,
5545 "a live claim must protect the task it belongs to"
5546 );
5547 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
5548 }
5549
5550 fn read_run_under(home: &Path, id: &str) -> RunState {
5556 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
5557 serde_json::from_str(&body).unwrap()
5558 }
5559
5560 #[test]
5561 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
5562 let dir = tempfile::tempdir().unwrap();
5563 let home = dir.path().to_path_buf();
5564 let now = Timestamp::now();
5565 let overrun_seat = || crate::run::ActiveSeat {
5566 node: "implement".to_owned(),
5567 started_at: now - jiff::SignedDuration::new(21_000, 0),
5568 timeout_secs: 3_600,
5569 attempt: 0,
5570 task: None,
5571 command: None,
5572 index: None,
5573 total: None,
5574 };
5575
5576 let mut dead = run_state(RunStatus::Implementing);
5577 dead.id = "20260101-000000-dead".to_owned();
5578 dead.active.insert("impl-A".to_owned(), overrun_seat());
5579 dead.driver_pid = Some(4242);
5582 dead.save_under(&home).unwrap();
5583
5584 let mut alive = run_state(RunStatus::Implementing);
5587 alive.id = "20260101-000000-aliv".to_owned();
5588 alive.active.insert("impl-A".to_owned(), overrun_seat());
5589 alive.save_under(&home).unwrap();
5590 let mut status = Status::new();
5591 status.current = vec![Current {
5592 task: "20260101-000000-task".to_owned(),
5593 run: alive.id.clone(),
5594 }];
5595 write_status_to(&home.join("daemon.json"), &status).unwrap();
5596
5597 let questions = Questions::at(home.join("questions"));
5601 let mut q = ask::Question::new(
5602 dead.id.clone(),
5603 "implement".to_owned(),
5604 "impl-A".to_owned(),
5605 "Which storage backend?".to_owned(),
5606 String::new(),
5607 vec!["SQLite".to_owned(), "Redis".to_owned()],
5608 );
5609 questions.put(&mut q).unwrap();
5610
5611 let abandoned = reclaim_abandoned_runs_with(
5612 &home,
5613 now,
5614 |pid| if pid == 4242 { Some(false) } else { None },
5615 |_| panic!("a query answering Dead outright needs no identity corroboration"),
5616 );
5617 assert_eq!(abandoned, vec![dead.id.clone()]);
5618
5619 let reloaded = read_run_under(&home, &dead.id);
5620 assert_eq!(reloaded.status, RunStatus::Failed);
5621 assert!(reloaded.active.is_empty());
5622 assert!(
5623 !questions.get(&q.id).unwrap().status.open(),
5624 "the failed run's own open question must be settled in the same pass"
5625 );
5626
5627 let still_alive = read_run_under(&home, &alive.id);
5628 assert_eq!(
5629 still_alive.status,
5630 RunStatus::Implementing,
5631 "a live daemon's claim protects it"
5632 );
5633 assert!(!still_alive.active.is_empty());
5634 }
5635
5636 #[test]
5646 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
5647 let dir = tempfile::tempdir().unwrap();
5648 let home = dir.path().to_path_buf();
5649 let now = Timestamp::now();
5650
5651 let mut manual = run_state(RunStatus::Reviewing);
5652 manual.id = "20260101-000000-manl".to_owned();
5653 manual.active.insert(
5654 "review-1".to_owned(),
5655 crate::run::ActiveSeat {
5656 node: "review".to_owned(),
5657 started_at: now - jiff::SignedDuration::new(21_000, 0),
5658 timeout_secs: 3_600,
5659 attempt: 0,
5660 task: None,
5661 command: None,
5662 index: None,
5663 total: None,
5664 },
5665 );
5666 manual.driver_pid = Some(4242);
5670 manual.driver_started_at = Some("2026-09-22T10:00:00Z".to_owned());
5671 manual.save_under(&home).unwrap();
5672
5673 let abandoned = reclaim_abandoned_runs_with(
5674 &home,
5675 now,
5676 |pid| if pid == 4242 { Some(true) } else { None },
5677 |pid| {
5678 if pid == 4242 {
5679 Some("2026-09-22T10:00:00Z".to_owned())
5680 } else {
5681 None
5682 }
5683 },
5684 );
5685 assert!(
5686 abandoned.is_empty(),
5687 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
5688 );
5689
5690 let reloaded = read_run_under(&home, &manual.id);
5691 assert_eq!(reloaded.status, RunStatus::Reviewing);
5692 assert!(!reloaded.active.is_empty());
5693 }
5694
5695 #[test]
5696 fn an_already_claimed_task_is_skipped_rather_than_failed() {
5697 let dir = tempfile::tempdir().unwrap();
5698 let queue = Queue::at(dir.path().to_path_buf());
5699 let mut only = task();
5700 queue.put(&mut only).unwrap();
5701
5702 let _elsewhere = queue.claim(&only.id).unwrap();
5703 let candidates = runnable(&queue);
5704 assert_eq!(candidates.len(), 1, "the task is still runnable");
5705 assert!(
5706 queue.claim(&candidates[0].id).is_err(),
5707 "the loop cannot take a claim somebody else holds"
5708 );
5709
5710 let after = queue.get(&only.id).unwrap();
5711 assert_eq!(after.status, TaskStatus::Queued);
5712 assert_eq!(
5713 after.attempts, 0,
5714 "losing the race is not an attempt at the task"
5715 );
5716 assert_eq!(after.last_error, None);
5717 }
5718
5719 #[test]
5720 fn the_status_file_round_trips_and_its_heartbeat_advances() {
5721 let dir = tempfile::tempdir().unwrap();
5722 let path = dir.path().join("daemon.json");
5723
5724 let mut status = Status::new();
5725 status.idle = false;
5726 status.completed = 7;
5727 status.current = vec![Current {
5728 task: "20260902-000000-t111".to_owned(),
5729 run: "20260902-000001-r111".to_owned(),
5730 }];
5731 write_status_to(&path, &status).unwrap();
5732 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
5733 assert_eq!(first.schema, SCHEMA);
5734 assert_eq!(first.pid, std::process::id());
5735 assert!(!first.idle);
5736 assert_eq!(first.completed, 7);
5737 assert_eq!(first.current, status.current);
5738 assert!(
5739 !path.with_extension("json.tmp").exists(),
5740 "the temp file is renamed, not left behind"
5741 );
5742
5743 std::thread::sleep(Duration::from_millis(5));
5744 status.updated_at = Timestamp::now();
5745 status.polls = 3;
5746 write_status_to(&path, &status).unwrap();
5747 let second: Status =
5748 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
5749 assert!(
5750 second.updated_at > first.updated_at,
5751 "a reader can only detect staleness if the heartbeat moves"
5752 );
5753 assert_eq!(
5754 second.started_at, first.started_at,
5755 "the start time is not a heartbeat"
5756 );
5757 assert_eq!(second.polls, 3);
5758 }
5759
5760 #[test]
5761 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
5762 let dir = tempfile::tempdir().unwrap();
5763
5764 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
5765
5766 let mut status = Status::new();
5767 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
5768 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5769 let stale = read_status(dir.path()).unwrap();
5770 assert!(
5771 !stale.running(Timestamp::now()),
5772 "a minute without a heartbeat is a dead daemon, not a busy one"
5773 );
5774 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
5775
5776 status.updated_at = Timestamp::now();
5777 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5778 let fresh = read_status(dir.path()).unwrap();
5779 assert!(fresh.running(Timestamp::now()));
5780 }
5781
5782 #[test]
5783 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
5784 let dir = tempfile::tempdir().unwrap();
5785 let now = Timestamp::now();
5786 let mine = "20260903-080619-01c2";
5787
5788 assert!(
5789 !is_working_on(dir.path(), mine, now),
5790 "no status file means nobody is working on anything"
5791 );
5792
5793 let mut status = Status::new();
5794 status.current = vec![Current {
5795 task: "20260903-080340-0167".to_owned(),
5796 run: mine.to_owned(),
5797 }];
5798 status.updated_at = now;
5799 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5800 assert!(is_working_on(dir.path(), mine, now));
5801 assert!(
5802 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
5803 "a daemon busy with one run is not working on another"
5804 );
5805
5806 status.updated_at = now - jiff::SignedDuration::from_secs(600);
5809 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5810 assert!(
5811 !is_working_on(dir.path(), mine, now),
5812 "a stale heartbeat is a dead daemon, so its run is a leftover"
5813 );
5814 }
5815
5816 #[test]
5817 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
5818 let dir = tempfile::tempdir().unwrap();
5819 let now = Timestamp::now();
5820
5821 assert!(
5822 !is_working_on_short(dir.path(), "01c2", now),
5823 "no status file means nobody is working on anything"
5824 );
5825
5826 let mut status = Status::new();
5827 status.current = vec![Current {
5828 task: "20260903-080340-0167".to_owned(),
5829 run: "20260903-080619-01c2".to_owned(),
5830 }];
5831 status.updated_at = now;
5832 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5833 assert!(
5834 is_working_on_short(dir.path(), "01c2", now),
5835 "the run's short id is the last block of its full id"
5836 );
5837 assert!(
5838 !is_working_on_short(dir.path(), "3cbf", now),
5839 "a daemon busy with one worktree bay is not working on another"
5840 );
5841 }
5842
5843 #[test]
5844 fn a_newer_status_file_still_yields_a_reading() {
5845 let dir = tempfile::tempdir().unwrap();
5846 std::fs::write(
5849 dir.path().join("daemon.json"),
5850 serde_json::json!({
5851 "schema": 2,
5852 "updated_at": Timestamp::now().to_string(),
5853 "idle": true,
5854 "surprise": { "nested": [1, 2, 3] },
5855 })
5856 .to_string(),
5857 )
5858 .unwrap();
5859
5860 let reading = read_status(dir.path()).expect("a forward-compatible read");
5861 assert!(reading.running(Timestamp::now()));
5862 assert!(reading.idle);
5863 assert!(reading.current.is_empty());
5864 }
5865
5866 #[test]
5867 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
5868 let dir = tempfile::tempdir().unwrap();
5874 std::fs::write(
5875 dir.path().join("daemon.json"),
5876 serde_json::json!({
5877 "schema": 1,
5878 "pid": 4242,
5879 "updated_at": Timestamp::now().to_string(),
5880 "idle": false,
5881 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
5882 "completed": 3,
5883 "polls": 9,
5884 })
5885 .to_string(),
5886 )
5887 .unwrap();
5888
5889 let reading = read_status(dir.path()).expect("an older shape must still parse");
5890 assert!(reading.running(Timestamp::now()));
5891 assert_eq!(
5892 reading.current,
5893 vec![Current {
5894 task: "20260902-140501-aaaa".to_owned(),
5895 run: "20260902-140502-bbbb".to_owned(),
5896 }]
5897 );
5898 }
5899
5900 #[test]
5901 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
5902 let dir = tempfile::tempdir().unwrap();
5903 std::fs::write(
5904 dir.path().join("daemon.json"),
5905 serde_json::json!({
5906 "schema": 1,
5907 "updated_at": Timestamp::now().to_string(),
5908 "idle": true,
5909 "current": null,
5910 })
5911 .to_string(),
5912 )
5913 .unwrap();
5914 let with_null = read_status(dir.path()).expect("null must still parse");
5915 assert!(with_null.current.is_empty());
5916
5917 std::fs::write(
5918 dir.path().join("daemon.json"),
5919 serde_json::json!({
5920 "schema": 1,
5921 "updated_at": Timestamp::now().to_string(),
5922 "idle": true,
5923 })
5924 .to_string(),
5925 )
5926 .unwrap();
5927 let absent = read_status(dir.path()).expect("a missing field must still parse");
5928 assert!(absent.current.is_empty());
5929 }
5930
5931 #[test]
5932 fn a_task_without_a_repository_runs_in_the_daemons_default() {
5933 let fallback = Path::new("/default");
5934 let mut blank = task();
5935 blank.repo = PathBuf::new();
5936 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
5937 let mut dot = task();
5938 dot.repo = PathBuf::from(".");
5939 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
5940 assert_eq!(
5941 repo_for(&task(), fallback),
5942 PathBuf::from("/repo"),
5943 "a task that names a repository keeps it"
5944 );
5945 }
5946
5947 #[test]
5948 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
5949 let mut solo_cfg = Config::default();
5955 solo_cfg.graph.candidates = 3;
5956 let mut solo_task = task();
5957 solo_task.solo = true;
5958 apply_solo(&mut solo_cfg, &solo_task);
5959 assert_eq!(solo_cfg.graph.candidates, 1);
5960
5961 let mut plain_cfg = Config::default();
5962 plain_cfg.graph.candidates = 3;
5963 let plain_task = task();
5964 assert!(!plain_task.solo);
5965 apply_solo(&mut plain_cfg, &plain_task);
5966 assert_eq!(
5967 plain_cfg.graph.candidates, 3,
5968 "a task that did not ask to run alone keeps the config's candidates"
5969 );
5970 }
5971
5972 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
5973 QuotaLoss {
5974 seat: seat.into(),
5975 node: "judge".into(),
5976 at: at.parse().unwrap(),
5977 reset: reset.map(str::to_string),
5978 }
5979 }
5980
5981 #[test]
5982 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
5983 let old: Vec<QuotaLoss> = (1..=4)
5984 .map(|i| {
5985 loss(
5986 &format!("judge-{i}"),
5987 "2026-09-23T05:23:00Z",
5988 Some("2:40pm (Asia/Tokyo)"),
5989 )
5990 })
5991 .collect();
5992 let fresh = losses_this_attempt(&old, &old);
5993 assert!(fresh.is_empty());
5994 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
5995 }
5997
5998 #[test]
5999 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
6000 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
6001 let now = Timestamp::now();
6002 let mut after = old.clone();
6003 after.push(loss("judge-2", &now.to_string(), None));
6004 let fresh = losses_this_attempt(&old, &after);
6005 assert_eq!(fresh, vec![after[1].clone()]);
6006 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
6007 assert_eq!(
6008 until,
6009 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
6010 );
6011 }
6012
6013 #[test]
6014 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
6015 let before = vec![
6018 loss("judge-1", "2026-09-23T05:23:00Z", None),
6019 loss("judge-2", "2026-09-23T05:24:00Z", None),
6020 ];
6021 let after = vec![
6022 loss("judge-2", "2026-09-23T05:24:00Z", None),
6023 loss("judge-1", "2026-09-24T01:00:00Z", None),
6024 ];
6025 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
6026 }
6027
6028 #[test]
6029 fn merge_overrides_are_parsed_or_refused() {
6030 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
6031 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
6032 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
6033 assert!(merge_mode("squash").is_err());
6034 }
6035
6036 #[test]
6037 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
6038 let now = Timestamp::now();
6039 let fallback = Duration::from_secs(300);
6040 let cap = Duration::from_secs(1800);
6041
6042 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
6044
6045 let soon = now + jiff::SignedDuration::from_secs(600);
6047 assert_eq!(
6048 quota_wait(Some(soon), now, fallback, cap),
6049 Duration::from_secs(600)
6050 );
6051
6052 let past = now - jiff::SignedDuration::from_secs(60);
6055 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
6056
6057 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
6060 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
6061 }
6062
6063 #[test]
6064 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
6065 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6066
6067 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
6068 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
6069
6070 let already_past =
6074 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
6075 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
6076
6077 assert!(
6078 parse_reset_hint("session limit reached", now, now).is_none(),
6079 "free text with no recognised shape is not guessed at"
6080 );
6081 assert!(
6082 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
6083 "an unresolvable zone name is not guessed at either"
6084 );
6085 }
6086
6087 #[test]
6088 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
6089 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
6090
6091 let at = parse_reset_hint(
6092 "You've hit your usage limit. Visit \
6093 https://chatgpt.com/codex/settings/usage to purchase more \
6094 credits or try again at Sep 19th, 2026 5:10 PM.",
6095 now,
6096 now,
6097 )
6098 .expect("the codex reset wording is a recognised shape");
6099 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
6100
6101 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
6106 .expect("an explicit year needs no rollover");
6107 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
6108
6109 assert!(
6110 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
6111 "a two-digit year is not the documented shape and is not guessed at"
6112 );
6113 assert!(
6114 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
6115 "a four-letter month name is not the documented three-letter abbreviation"
6116 );
6117 assert!(
6118 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
6119 "an explicit zone on the dated shape is a format nobody has \
6120 documented, and is refused rather than guessed at as UTC"
6121 );
6122 }
6123
6124 #[test]
6125 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
6126 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
6127 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
6128
6129 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
6130 assert_eq!(at.as_second() - recorded.as_second(), 3769);
6131
6132 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
6133 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
6134
6135 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
6136 assert!(
6137 parse_reset_hint(bad, now, recorded).is_none(),
6138 "{bad:?} must not be guessed at"
6139 );
6140 }
6141 }
6142
6143 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
6147 let config = dir.join("magi.toml");
6148 std::fs::write(
6149 &config,
6150 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
6151 )
6152 .unwrap();
6153 let opts = Opts {
6154 poll: Duration::from_secs(30),
6155 config: Some(config),
6156 repo: dir.join("repo"),
6160 ..Opts::default()
6161 };
6162 let home = dir.join("home");
6171 let worktrees = dir.join("wt");
6172 (
6173 opts,
6174 Queue::at(dir.join("queue")),
6175 home.join("daemon.json"),
6176 home,
6177 worktrees,
6178 )
6179 }
6180
6181 #[test]
6182 fn a_stop_is_idempotent_and_once_set_stays_set() {
6183 let stop = Stop::new();
6184 assert!(!stop.stopped());
6185
6186 stop.stop();
6187 assert!(stop.stopped());
6188 stop.stop();
6189 assert!(stop.stopped(), "a second stop is not a toggle");
6190
6191 let shared = stop.clone();
6192 assert!(
6193 shared.stopped(),
6194 "a clone is the same stop; that is how the loop and its caller share one"
6195 );
6196 }
6197
6198 #[test]
6199 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6200 let stop = Stop::new();
6201 stop.enter();
6202 assert!(
6203 !stop.finishing(),
6204 "a busy loop nobody has asked to stop is just running"
6205 );
6206
6207 stop.stop();
6208 assert!(
6209 stop.finishing(),
6210 "a stop asked for mid-run has not landed until the run is settled"
6211 );
6212
6213 stop.exit();
6214 assert!(
6215 !stop.finishing(),
6216 "once the run is settled the stop has landed and there is nothing to finish"
6217 );
6218 }
6219
6220 #[test]
6221 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6222 let stop = Stop::new();
6223 stop.enter();
6224 stop.enter();
6225 stop.stop();
6226 assert!(stop.finishing(), "two runs still in flight");
6227
6228 stop.exit();
6229 assert!(
6230 stop.finishing(),
6231 "one run finished, but a sibling is still working"
6232 );
6233
6234 stop.exit();
6235 assert!(
6236 !stop.finishing(),
6237 "the last run out is what actually lands the stop"
6238 );
6239 }
6240
6241 #[tokio::test]
6242 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6243 let dir = tempfile::tempdir().unwrap();
6244 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6245 let stop = Stop::new();
6246 stop.stop();
6247
6248 let began = std::time::Instant::now();
6249 tokio::time::timeout(
6250 Duration::from_secs(2),
6251 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6252 )
6253 .await
6254 .expect("a stopped loop must return, not sit out its poll interval")
6255 .expect("the loop's own setup and teardown must not fail");
6256 assert!(
6257 began.elapsed() < opts.poll,
6258 "returned only after {:?}, which is a poll interval, not a stop",
6259 began.elapsed()
6260 );
6261 }
6262
6263 #[tokio::test]
6264 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6265 let dir = tempfile::tempdir().unwrap();
6266 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6267 let stop = Stop::new();
6268
6269 let asker = {
6272 let stop = stop.clone();
6273 tokio::spawn(async move {
6274 tokio::time::sleep(Duration::from_millis(20)).await;
6275 stop.stop();
6276 })
6277 };
6278
6279 let began = std::time::Instant::now();
6280 tokio::time::timeout(
6281 Duration::from_secs(2),
6282 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6283 )
6284 .await
6285 .expect("a stop asked for while idle must wake the wait")
6286 .expect("the loop's own setup and teardown must not fail");
6287 asker.await.unwrap();
6288 assert!(
6289 began.elapsed() < opts.poll,
6290 "returned only after {:?}, so the stop waited on the sleep",
6291 began.elapsed()
6292 );
6293 }
6294
6295 #[tokio::test]
6296 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6297 let dir = tempfile::tempdir().unwrap();
6298 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6299 let stop = Stop::new();
6300 stop.stop();
6301
6302 tokio::time::timeout(
6303 Duration::from_secs(2),
6304 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6305 )
6306 .await
6307 .expect("a stopped loop must return")
6308 .expect("the loop's own setup and teardown must not fail");
6309
6310 assert!(
6311 home.is_dir(),
6312 "the loop did publish a status file, so its removal is the teardown and not an absence"
6313 );
6314 assert!(
6315 !status_file.exists(),
6316 "a stopped loop clears its status file"
6317 );
6318 assert!(
6319 read_status(&home).is_none(),
6320 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6321 );
6322 }
6323
6324 #[tokio::test]
6325 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6326 let dir = tempfile::tempdir().unwrap();
6327 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6328 opts.once = true;
6329
6330 let mut settled = RunState::new(
6331 dir.path().join("repo"),
6332 "main".to_owned(),
6333 "abc1234".to_owned(),
6334 "fixture".to_owned(),
6335 Config::default(),
6336 );
6337 settled.status = RunStatus::Ready;
6338 let run_dir = home.join("runs").join(&settled.id);
6339 std::fs::create_dir_all(&run_dir).unwrap();
6340 std::fs::write(
6341 run_dir.join("run.json"),
6342 serde_json::to_string_pretty(&settled).unwrap(),
6343 )
6344 .unwrap();
6345 let questions = Questions::at(home.join("questions"));
6346 let mut question = ask::Question::new(
6347 settled.id.clone(),
6348 "review".to_owned(),
6349 "reviewer-1".to_owned(),
6350 "Continue?".to_owned(),
6351 String::new(),
6352 Vec::new(),
6353 );
6354 questions.put(&mut question).unwrap();
6355
6356 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6357 .await
6358 .unwrap();
6359
6360 assert_eq!(
6361 questions.get(&question.id).unwrap().status,
6362 ask::QuestionStatus::Abandoned,
6363 "an empty --once drain still performs startup question cleanup"
6364 );
6365 }
6366
6367 #[test]
6368 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6369 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6370
6371 assert!(
6372 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6373 "never checked before: due at once"
6374 );
6375
6376 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6377 assert!(
6378 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6379 "well inside the interval: not due yet"
6380 );
6381
6382 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6383 assert!(
6384 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6385 "exactly at the edge: not yet due, same convention as `clean::due`"
6386 );
6387
6388 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6389 assert!(
6390 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6391 "past the interval: due again"
6392 );
6393 }
6394
6395 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
6400 let config = dir.join("magi.toml");
6401 std::fs::write(
6407 &config,
6408 format!(
6409 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
6410 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
6411 cache_dir.display()
6412 ),
6413 )
6414 .unwrap();
6415 Opts {
6416 config: Some(config),
6417 repo: dir.join("repo"),
6418 ..Opts::default()
6419 }
6420 }
6421
6422 #[tokio::test]
6423 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
6424 let dir = tempfile::tempdir().unwrap();
6425 let home = dir.path().join("home");
6426 let cache_dir = dir.path().join("cache");
6427 std::fs::create_dir_all(&cache_dir).unwrap();
6428 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6429 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6430
6431 let running = Stop::new();
6434 let mut last_checked = None;
6435 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6436 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
6437 .await;
6438 assert_eq!(
6439 crate::disk::dir_size(&cache_dir),
6440 0,
6441 "over the cap on the first check ever: pruned at once, no idle queue required"
6442 );
6443 assert_eq!(last_checked, Some(t0));
6444
6445 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
6447 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
6448 maybe_prune_cache_between_runs(
6449 &opts.repo,
6450 &opts,
6451 &home,
6452 &running,
6453 &mut last_checked,
6454 too_soon,
6455 )
6456 .await;
6457 assert_eq!(
6458 crate::disk::dir_size(&cache_dir),
6459 10,
6460 "too soon since the last check: left alone rather than rescanned every call"
6461 );
6462 assert_eq!(
6463 last_checked,
6464 Some(t0),
6465 "an idle check does not reset the clock"
6466 );
6467
6468 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6470 maybe_prune_cache_between_runs(
6471 &opts.repo,
6472 &opts,
6473 &home,
6474 &running,
6475 &mut last_checked,
6476 due_again,
6477 )
6478 .await;
6479 assert_eq!(
6480 crate::disk::dir_size(&cache_dir),
6481 0,
6482 "due again: pruned back under the cap"
6483 );
6484 }
6485
6486 #[tokio::test]
6494 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
6495 let dir = tempfile::tempdir().unwrap();
6496 let home = dir.path().join("home");
6497 let cache_dir = dir.path().join("cache");
6498 std::fs::create_dir_all(&cache_dir).unwrap();
6499 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6500 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6501
6502 let stop = Stop::new();
6503 stop.stop();
6504 assert!(
6505 !stop.finishing(),
6506 "no run is in flight at a between-runs boundary, so nothing else \
6507 would tell the operator this stop had not taken effect yet"
6508 );
6509
6510 let mut last_checked = None;
6511 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6512 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
6513 .await;
6514 assert_eq!(
6515 crate::disk::dir_size(&cache_dir),
6516 10,
6517 "over its cap, and due for the first check ever, but a stop outranks \
6518 it: the cap is a standing policy the next start measures again"
6519 );
6520 assert_eq!(
6521 last_checked, None,
6522 "a check that never happened must not claim the interval"
6523 );
6524 }
6525
6526 #[tokio::test]
6540 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
6541 let dir = tempfile::tempdir().unwrap();
6542 let cache_dir = dir.path().join("cache");
6543 std::fs::create_dir_all(&cache_dir).unwrap();
6544 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
6545
6546 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
6547 opts.poll = Duration::from_millis(20);
6548 opts.max_attempts = 1_000;
6549
6550 let queue = Queue::at(dir.path().join("queue"));
6551 let mut t = Task::new(
6552 "x".to_owned(),
6553 "x".to_owned(),
6554 opts.repo.clone(),
6555 Source::Human,
6556 );
6557 queue.put(&mut t).unwrap();
6558
6559 let home = dir.path().join("home");
6560 let worktrees = dir.path().join("wt");
6561 let status_file = home.join("daemon.json");
6562 let stop = Stop::new();
6563 let stopper = {
6564 let stop = stop.clone();
6565 tokio::spawn(async move {
6566 tokio::time::sleep(Duration::from_millis(400)).await;
6567 stop.stop();
6568 })
6569 };
6570
6571 tokio::time::timeout(
6572 Duration::from_secs(10),
6573 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6574 )
6575 .await
6576 .expect("the loop must not hang on a queue that keeps producing failing work")
6577 .expect("the loop's own setup and teardown must not fail");
6578 stopper.await.unwrap();
6579
6580 let after = queue.get(&t.id).unwrap();
6581 assert!(
6582 after.attempts >= 2,
6583 "the harness must actually have retried more than once, or this is not \
6584 exercising a busy queue at all (got {} attempt(s))",
6585 after.attempts
6586 );
6587 assert!(
6588 after.status.runnable(),
6589 "still under its attempt budget: the queue never reached a natural idle \
6590 on its own, only the external stop ended the test"
6591 );
6592
6593 assert_eq!(
6594 crate::disk::dir_size(&cache_dir),
6595 0,
6596 "an oversized cache must not be left to grow unboundedly just because the \
6597 queue kept the loop busy the whole time"
6598 );
6599 }
6600
6601 #[test]
6602 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
6603 let dir = tempfile::tempdir().unwrap();
6604 let queue = Queue::at(dir.path().join("queue"));
6605 let questions = Questions::at(dir.path().join("questions"));
6606 let mut task = task();
6607 queue.put(&mut task).unwrap();
6608
6609 let mut task_question = ask::Question::new(
6610 task.id.clone(),
6611 crate::conduct::NODE.to_owned(),
6612 "conduct".to_owned(),
6613 "Which backend?".to_owned(),
6614 String::new(),
6615 Vec::new(),
6616 );
6617 questions.put(&mut task_question).unwrap();
6618 task.block(vec![task_question.id.clone()], None);
6619 queue.put(&mut task).unwrap();
6620
6621 let mut run_question = ask::Question::new(
6622 "20260101-000000-run1".to_owned(),
6623 "review".to_owned(),
6624 "reviewer-1".to_owned(),
6625 "Run question".to_owned(),
6626 String::new(),
6627 Vec::new(),
6628 );
6629 questions.put(&mut run_question).unwrap();
6630
6631 let mut coincidental = ask::Question::new(
6636 task.id.clone(),
6637 "review".to_owned(),
6638 "reviewer-1".to_owned(),
6639 "Unrelated review question".to_owned(),
6640 String::new(),
6641 Vec::new(),
6642 );
6643 questions.put(&mut coincidental).unwrap();
6644
6645 reconcile_task_questions(&queue, &questions);
6646 assert!(questions.get(&task_question.id).unwrap().status.open());
6647 assert!(questions.get(&run_question.id).unwrap().status.open());
6648 assert!(questions.get(&coincidental.id).unwrap().status.open());
6649
6650 task.release();
6651 queue.put(&mut task).unwrap();
6652 reconcile_task_questions(&queue, &questions);
6653 assert_eq!(
6654 questions.get(&task_question.id).unwrap().status,
6655 ask::QuestionStatus::Abandoned
6656 );
6657 assert!(
6658 questions.get(&run_question.id).unwrap().status.open(),
6659 "run questions remain the run janitor's responsibility"
6660 );
6661 assert!(
6662 questions.get(&coincidental.id).unwrap().status.open(),
6663 "a non-conductor question must not be abandoned just because its \
6664 run id coincides with a task id"
6665 );
6666 }
6667
6668 #[test]
6669 fn a_freshly_started_running_task_is_never_stalled() {
6670 let dir = tempfile::tempdir().unwrap();
6671 let mut t = task();
6672 t.start("run-1".to_owned());
6673 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
6676 }
6677
6678 #[test]
6679 fn a_long_running_task_with_no_live_daemon_is_stalled() {
6680 let dir = tempfile::tempdir().unwrap();
6681 let mut t = task();
6682 t.start("run-1".to_owned());
6683 t.updated_at = Timestamp::now()
6684 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6685 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
6686 assert_eq!(
6687 stalled_tasks(
6688 &Queue::at(dir.path().join("q")),
6689 dir.path(),
6690 Timestamp::now()
6691 )
6692 .len(),
6693 0,
6694 "the task was never written to this queue"
6695 );
6696 }
6697
6698 #[test]
6699 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
6700 let dir = tempfile::tempdir().unwrap();
6701 let mut t = task();
6702 t.id = "20260903-080340-0167".to_owned();
6703 t.start("20260903-080619-01c2".to_owned());
6704 t.updated_at = Timestamp::now()
6705 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6706
6707 let mut status = Status::new();
6708 status.current = vec![Current {
6709 task: t.id.clone(),
6710 run: "20260903-080619-01c2".to_owned(),
6711 }];
6712 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6713
6714 assert!(
6715 !is_stalled(&t, dir.path(), Timestamp::now()),
6716 "a live daemon's own heartbeat rules out stalled, however long the task has run"
6717 );
6718 }
6719
6720 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
6724 let path = queue.path_of(id);
6725 let body = std::fs::read_to_string(&path).unwrap();
6726 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
6727 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
6728 v["updated_at"] = serde_json::Value::String(old.to_string());
6729 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
6730 }
6731
6732 #[test]
6733 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
6734 let dir = tempfile::tempdir().unwrap();
6747 let queue = Queue::at(dir.path().join("queue"));
6748 let home = dir.path().join("home");
6749
6750 let mut t = task();
6751 t.id = "20260101-000001-lock".to_owned();
6752 t.start("run-1".to_owned());
6753 queue.put(&mut t).unwrap();
6754 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
6755 std::fs::write(
6756 dir.path().join("queue").join(format!("{}.lock", t.id)),
6757 "not a pid",
6758 )
6759 .unwrap();
6760
6761 let now = Timestamp::now();
6762 assert!(
6763 reclaim_orphaned_running(&queue, 2).is_empty(),
6764 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
6765 and reclaim must leave the task alone"
6766 );
6767 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
6768
6769 let stalled = stalled_tasks(&queue, &home, now);
6770 assert_eq!(
6771 stalled.len(),
6772 1,
6773 "reclaim's inability to claim it yet must not hide it from the conductor"
6774 );
6775 assert_eq!(stalled[0].id, t.id);
6776 }
6777
6778 #[test]
6779 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
6780 let dir = tempfile::tempdir().unwrap();
6781 crate::run::set_home(dir.path().join("run-home"));
6782 let queue = Queue::at(dir.path().join("queue"));
6783 let home = dir.path().join("home");
6784 let questions = Questions::at(dir.path().join("questions"));
6785
6786 let mut t = task();
6787 t.id = "20260101-000003-dead".to_owned();
6788 t.start("missing-run".to_owned());
6789 queue.put(&mut t).unwrap();
6790 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
6791
6792 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
6795 assert_eq!(
6796 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
6797 [&t.id]
6798 );
6799 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
6800 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
6801
6802 crate::conduct::apply(
6805 &queue,
6806 &questions,
6807 &crate::conduct::Verdict {
6808 decisions: vec![crate::conduct::Decision {
6809 id: t.id.clone(),
6810 recovery: Some(crate::conduct::Recovery::Requeue),
6811 ..crate::conduct::Decision::default()
6812 }],
6813 },
6814 )
6815 .unwrap();
6816 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
6817 }
6818
6819 #[test]
6820 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
6821 let dir = tempfile::tempdir().unwrap();
6822 let queue = Queue::at(dir.path().join("queue"));
6823 let home = dir.path().join("home");
6824
6825 let mut fresh = task();
6826 fresh.id = "20260101-000001-aaaa".to_owned();
6827 fresh.start("run-1".to_owned());
6828 queue.put(&mut fresh).unwrap();
6829
6830 let mut old = task();
6831 old.id = "20260101-000002-bbbb".to_owned();
6832 old.start("run-2".to_owned());
6833 queue.put(&mut old).unwrap();
6834 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
6835
6836 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
6837 assert_eq!(stalled.len(), 1);
6838 assert_eq!(stalled[0].id, old.id);
6839 }
6840
6841 #[test]
6842 fn queued_and_finished_task_views_partition_by_status() {
6843 let dir = tempfile::tempdir().unwrap();
6844 let queue = Queue::at(dir.path().join("queue"));
6845
6846 let mut queued = task();
6847 queued.id = "20260101-000001-aaaa".to_owned();
6848 queue.put(&mut queued).unwrap();
6849
6850 let mut failed = task();
6851 failed.id = "20260101-000002-bbbb".to_owned();
6852 failed.start("run-1".to_owned());
6853 failed.fail("gate red", 5);
6854 queue.put(&mut failed).unwrap();
6855
6856 let mut held = task();
6857 held.id = "20260101-000003-cccc".to_owned();
6858 held.hold_machine(None);
6859 queue.put(&mut held).unwrap();
6860
6861 let mut running = task();
6862 running.id = "20260101-000004-dddd".to_owned();
6863 running.start("run-2".to_owned());
6864 queue.put(&mut running).unwrap();
6865
6866 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
6867 assert_eq!(queued_ids, [queued.id.clone()]);
6868
6869 let mut finished_ids: Vec<String> =
6870 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
6871 finished_ids.sort_unstable();
6872 let mut want = vec![failed.id.clone(), held.id.clone()];
6873 want.sort_unstable();
6874 assert_eq!(finished_ids, want);
6875 }
6876
6877 #[test]
6878 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
6879 let dir = tempfile::tempdir().unwrap();
6880 let queue = Queue::at(dir.path().join("queue"));
6881 let questions = ask::Questions::at(dir.path().join("questions"));
6882
6883 let mut dep = task();
6884 dep.id = "20260101-000001-dep0".to_owned();
6885 dep.succeed();
6886 queue.put(&mut dep).unwrap();
6887
6888 let mut still_going = task();
6889 still_going.id = "20260101-000002-dep1".to_owned();
6890 queue.put(&mut still_going).unwrap();
6891
6892 let mut blocked = task();
6893 blocked.id = "20260101-000003-main".to_owned();
6894 blocked.block(
6895 vec![dep.id.clone(), still_going.id.clone()],
6896 Some("waits on both".to_owned()),
6897 );
6898 queue.put(&mut blocked).unwrap();
6899
6900 resolve_blockers(&queue, &questions);
6901
6902 let after = queue.get(&blocked.id).unwrap();
6903 assert_eq!(
6904 after.status,
6905 TaskStatus::Blocked,
6906 "one dependency is still outstanding"
6907 );
6908 assert_eq!(after.blocked_by, [still_going.id.clone()]);
6909 }
6910
6911 #[test]
6912 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
6913 let dir = tempfile::tempdir().unwrap();
6914 let queue = Queue::at(dir.path().join("queue"));
6915 let questions = ask::Questions::at(dir.path().join("questions"));
6916
6917 let mut q = crate::ask::Question::new(
6918 "20260101-000001-main".to_owned(),
6919 crate::conduct::NODE.to_owned(),
6920 "conduct".to_owned(),
6921 "Which backend?".to_owned(),
6922 String::new(),
6923 Vec::new(),
6924 );
6925 questions.put(&mut q).unwrap();
6926 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
6927 .unwrap();
6928 questions.put(&mut q).unwrap();
6929
6930 let mut blocked = task();
6931 blocked.id = "20260101-000001-main".to_owned();
6932 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
6933 queue.put(&mut blocked).unwrap();
6934
6935 resolve_blockers(&queue, &questions);
6936
6937 let after = queue.get(&blocked.id).unwrap();
6938 assert_eq!(
6939 after.status,
6940 TaskStatus::Queued,
6941 "the only blocker resolved"
6942 );
6943 assert_eq!(after.answers.len(), 1);
6944 assert_eq!(after.answers[0].question, "Which backend?");
6945 assert_eq!(after.answers[0].answer, "SQLite");
6946
6947 let instruction = instruction_for(&after);
6949 assert!(instruction.contains("Which backend?"));
6950 assert!(instruction.contains("SQLite"));
6951 }
6952
6953 #[test]
6954 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
6955 let dir = tempfile::tempdir().unwrap();
6961 let queue = Queue::at(dir.path().join("queue"));
6962 let questions = ask::Questions::at(dir.path().join("questions"));
6963
6964 let mut q = crate::ask::Question::new(
6965 "20260101-000001-main".to_owned(),
6966 crate::conduct::NODE.to_owned(),
6967 "conduct".to_owned(),
6968 "How should this be handled?".to_owned(),
6969 String::new(),
6970 Vec::new(),
6971 );
6972 questions.put(&mut q).unwrap();
6973 q.answer(crate::ask::Answer::Text(
6974 "leave it held, a human will look at it later".to_owned(),
6975 ))
6976 .unwrap();
6977 questions.put(&mut q).unwrap();
6978
6979 let mut held = task();
6980 held.id = "20260101-000001-main".to_owned();
6981 held.hold_machine(Some("out of attempts".to_owned()));
6982 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
6983 queue.put(&mut held).unwrap();
6984
6985 resolve_blockers(&queue, &questions);
6986
6987 let after = queue.get(&held.id).unwrap();
6988 assert_eq!(after.status, TaskStatus::Held);
6989 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
6990 assert_eq!(
6991 after.answers[0].answer,
6992 "leave it held, a human will look at it later"
6993 );
6994 }
6995
6996 #[test]
6997 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
6998 let dir = tempfile::tempdir().unwrap();
7004 let queue = Queue::at(dir.path().join("queue"));
7005 let questions = ask::Questions::at(dir.path().join("questions"));
7006
7007 let mut still_going = task();
7008 still_going.id = "20260101-000002-dep1".to_owned();
7009 queue.put(&mut still_going).unwrap();
7010
7011 let mut blocked = task();
7012 blocked.id = "20260101-000003-main".to_owned();
7013 blocked.block(
7014 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
7015 Some("waits on both".to_owned()),
7016 );
7017 queue.put(&mut blocked).unwrap();
7018
7019 resolve_blockers(&queue, &questions);
7020
7021 let after = queue.get(&blocked.id).unwrap();
7022 assert_eq!(
7023 after.status,
7024 TaskStatus::Held,
7025 "a missing dependency must not leave the task blocked forever"
7026 );
7027 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
7028 assert!(after.blocked_by.is_empty());
7029 let reason = after.hold_reason.as_deref().unwrap_or_default();
7030 assert!(
7031 reason.contains("20260101-000001-gone"),
7032 "the missing id must be named so an operator can tell what happened: {reason}"
7033 );
7034 assert!(
7035 reason.contains(&still_going.id),
7036 "the still-valid dependency must not silently vanish from the record: {reason}"
7037 );
7038 }
7039
7040 #[test]
7041 fn instruction_for_is_unchanged_without_any_answers() {
7042 let t = task();
7043 assert_eq!(instruction_for(&t), t.instruction);
7044 }
7045
7046 #[test]
7047 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
7048 let dir = tempfile::tempdir().unwrap();
7049 let q = Queue::at(dir.path().join("queue"));
7050 let src = dir.path().join("shot.png");
7051 std::fs::write(&src, "x").unwrap();
7052 let mut t = task();
7053 q.attach(&mut t, &[src]).unwrap();
7054 let paths = task_attachments(&q, &t).unwrap();
7055 assert_eq!(paths.len(), 1);
7056 assert!(paths[0].is_absolute() && paths[0].is_file());
7057 std::fs::remove_file(&paths[0]).unwrap();
7058 let err = task_attachments(&q, &t).unwrap_err().to_string();
7059 assert!(err.contains("shot.png"), "{err}");
7060 }
7061
7062 #[test]
7063 fn resumed_instruction_is_unchanged_without_any_answers() {
7064 let t = task();
7065 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
7066 }
7067
7068 #[test]
7069 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
7070 let mut t = task();
7071 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7072 let old = t.instruction.clone();
7076
7077 let refreshed = resumed_instruction(&old, &t);
7078 assert!(refreshed.starts_with(&old), "the original text is kept");
7079 assert!(refreshed.contains("Which backend?"));
7080 assert!(refreshed.contains("SQLite"));
7081 }
7082
7083 #[test]
7084 fn resumed_instruction_keeps_an_original_answers_heading() {
7085 let mut t = task();
7086 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
7087 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7088
7089 let refreshed = resumed_instruction(&t.instruction, &t);
7090
7091 assert!(
7092 refreshed.starts_with(&t.instruction),
7093 "an answers heading in the original instruction is not the appended block"
7094 );
7095 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
7096 assert!(refreshed.contains("Which backend?"));
7097 assert!(refreshed.contains("SQLite"));
7098
7099 let repeated = resumed_instruction(&refreshed, &t);
7100 assert_eq!(
7101 repeated, refreshed,
7102 "only the final appended block is refreshed"
7103 );
7104 }
7105
7106 #[test]
7107 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
7108 let mut t = task();
7109 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7110
7111 let once = resumed_instruction(&t.instruction, &t);
7115 let twice = resumed_instruction(&once, &t);
7116 assert_eq!(once, twice);
7117 assert_eq!(once.matches("Which backend?").count(), 1);
7118
7119 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
7121 let refreshed = resumed_instruction(&once, &t);
7122 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
7123 assert!(refreshed.contains("Which backend?"));
7124 assert!(refreshed.contains("Which cache?"));
7125 }
7126
7127 #[test]
7128 fn prepare_instruction_covers_all_three_starters() {
7129 let mut t = task();
7130 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
7131
7132 assert_eq!(
7135 prepare_instruction(&Starter::Start, None, &t),
7136 Some(instruction_for(&t))
7137 );
7138
7139 let old = t.instruction.clone();
7142 assert_eq!(
7143 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
7144 Some(resumed_instruction(&old, &t))
7145 );
7146
7147 assert_eq!(
7151 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
7152 None
7153 );
7154 }
7155
7156 #[test]
7157 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
7158 assert_eq!(
7159 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
7160 Starter::Review("magi/eba2/A".to_owned())
7161 );
7162 }
7163
7164 #[test]
7165 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
7166 assert_eq!(
7167 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
7168 Starter::Start,
7169 "a vanished review branch must not fall back to resuming the old run either"
7170 );
7171 }
7172
7173 #[test]
7174 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
7175 assert_eq!(
7176 choose_starter(None, false, Some("some-run")),
7177 Starter::Resume("some-run".to_owned())
7178 );
7179 assert_eq!(choose_starter(None, false, None), Starter::Start);
7180 }
7181
7182 #[test]
7183 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7184 let mut released = task();
7185 released.start("stalled-run".to_owned());
7186 released.requeue();
7187 let unfinished = (!released.fresh_start)
7188 .then(|| Some("stalled-run".to_owned()))
7189 .flatten();
7190 assert_eq!(
7191 choose_starter(None, false, unfinished.as_deref()),
7192 Starter::Start,
7193 "release keeps run history but must not resume it"
7194 );
7195 assert_eq!(released.runs, ["stalled-run"]);
7196 }
7197
7198 #[test]
7199 fn an_ordinary_release_keeps_a_resumable_run_available() {
7200 let mut released = task();
7201 released.start("stalled-run".to_owned());
7202 released.release();
7203 let unfinished = (!released.fresh_start)
7204 .then(|| Some("stalled-run".to_owned()))
7205 .flatten();
7206 assert_eq!(
7207 choose_starter(None, false, unfinished.as_deref()),
7208 Starter::Resume("stalled-run".to_owned()),
7209 "manual release must preserve the normal resume path"
7210 );
7211 }
7212
7213 #[test]
7214 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7215 let mut state = run_state(RunStatus::Blocked);
7216 state.config.graph.review_rounds = 3;
7217 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7218 assert!(exhausted_review_budget(&state));
7219
7220 state.reviews.pop();
7222 assert!(!exhausted_review_budget(&state));
7223
7224 let mut stalled = run_state(RunStatus::Stalled);
7227 stalled.config.graph.review_rounds = 1;
7228 stalled.reviews = vec![review_round(1)];
7229 assert!(!exhausted_review_budget(&stalled));
7230 }
7231
7232 fn review_round(round: usize) -> crate::run::ReviewRound {
7233 crate::run::ReviewRound {
7234 round,
7235 head: "deadbeef".to_owned(),
7236 verified_head: None,
7237 verified_at: None,
7238 reviews: Vec::new(),
7239 e2e: Vec::new(),
7240 verify_retried: false,
7241 e2e_deferred: false,
7242 e2e_defer_reason: None,
7243 fix: None,
7244 blocking: 0,
7245 answered: 1,
7246 expected: 1,
7247 clean: false,
7248 progressed: true,
7249 vote_split: false,
7250 reconsideration: Vec::new(),
7251 verdict: None,
7252 }
7253 }
7254
7255 #[test]
7256 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7257 let mut run = RunState::new(
7258 PathBuf::from("/repo"),
7259 "main".to_owned(),
7260 "abc1234def".to_owned(),
7261 "add retries".to_owned(),
7262 Config::default(),
7263 );
7264 run.status = RunStatus::Judging;
7265 run.parked = true;
7266 let mut task = Task::new(
7267 "add retries".to_owned(),
7268 "add retries".to_owned(),
7269 PathBuf::from("/repo"),
7270 crate::queue::Source::Human,
7271 );
7272 task.status = TaskStatus::Failed;
7273 task.runs = vec![run.id.clone()];
7274 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7275 assert!(with(&task, &run), "parked after judging is the case");
7276
7277 let mut not_parked = run.clone();
7278 not_parked.parked = false;
7279 not_parked.status = RunStatus::Stalled;
7280 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7281
7282 let mut fresh = task.clone();
7283 fresh.fresh_start = true;
7284 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7285
7286 let mut review = task.clone();
7287 review.review_branch = Some("magi/x/A".to_owned());
7288 assert!(!with(&review, &run), "review is ranked before resume");
7289
7290 let mut held = task.clone();
7291 held.status = TaskStatus::Held;
7292 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7293
7294 let mut released = run.clone();
7295 released.released_to = Some("20260901-000000-new1".to_owned());
7296 assert!(!with(&task, &released), "nothing left to resume into");
7297
7298 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7299 "unreadable"
7300 )));
7301 }
7302
7303 #[test]
7304 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7305 let mut released = RunState::new(
7306 PathBuf::from("/repo"),
7307 "main".to_owned(),
7308 "abc1234def".to_owned(),
7309 "add retries".to_owned(),
7310 Config::default(),
7311 );
7312 released.status = RunStatus::Blocked;
7313 assert_eq!(
7314 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7315 Some(released.id.clone())
7316 );
7317 released.released_to = Some("20260901-000000-new1".to_owned());
7318 assert_eq!(
7319 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7320 None,
7321 "there is nothing left to resume it into"
7322 );
7323 }
7324
7325 #[test]
7326 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
7327 let mut exhausted = RunState::new(
7336 PathBuf::from("/repo"),
7337 "main".to_owned(),
7338 "abc1234def".to_owned(),
7339 "add retries".to_owned(),
7340 Config::default(),
7341 );
7342 exhausted.status = RunStatus::Blocked;
7343 exhausted.config.graph.review_rounds = 1;
7344 exhausted.reviews = vec![review_round(1)];
7345
7346 assert_eq!(
7347 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
7348 None,
7349 "an exhausted `Blocked` run must not be offered as resumable"
7350 );
7351
7352 let mut has_budget_left = RunState::new(
7355 PathBuf::from("/repo"),
7356 "main".to_owned(),
7357 "abc1234def".to_owned(),
7358 "add retries".to_owned(),
7359 Config::default(),
7360 );
7361 has_budget_left.status = RunStatus::Blocked;
7362 has_budget_left.config.graph.review_rounds = 3;
7363 has_budget_left.reviews = vec![review_round(1)];
7364
7365 assert_eq!(
7366 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
7367 Ok(has_budget_left.clone())
7368 }),
7369 Some(has_budget_left.id.clone())
7370 );
7371 }
7372
7373 #[test]
7374 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
7375 let mut older_stalled = RunState::new(
7383 PathBuf::from("/repo"),
7384 "main".to_owned(),
7385 "abc1234def".to_owned(),
7386 "add retries".to_owned(),
7387 Config::default(),
7388 );
7389 older_stalled.status = RunStatus::Stalled;
7390
7391 let mut newest_exhausted = RunState::new(
7392 PathBuf::from("/repo"),
7393 "main".to_owned(),
7394 "abc1234def".to_owned(),
7395 "add retries".to_owned(),
7396 Config::default(),
7397 );
7398 newest_exhausted.status = RunStatus::Blocked;
7399 newest_exhausted.config.graph.review_rounds = 1;
7400 newest_exhausted.reviews = vec![review_round(1)];
7401
7402 assert_eq!(
7403 unfinished_run_with(
7404 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
7405 "t",
7406 |_| Ok(newest_exhausted.clone())
7407 ),
7408 None,
7409 "the newest run is exhausted, so nothing here is worth resuming - \
7410 least of all the older, already-superseded run"
7411 );
7412 }
7413
7414 #[test]
7415 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
7416 assert_eq!(
7417 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
7418 Err(anyhow::anyhow!("fixture is absent"))
7419 }),
7420 None
7421 );
7422 }
7423
7424 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
7425 let mut q = ask::Question::new(
7426 run.to_owned(),
7427 "implement".to_owned(),
7428 "impl-A".to_owned(),
7429 "continue?".to_owned(),
7430 String::new(),
7431 vec!["resume で続行する".to_owned(), "other".to_owned()],
7432 );
7433 q.actions.insert("resume で続行する".to_owned(), action);
7434 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
7435 .unwrap();
7436 q
7437 }
7438
7439 fn held_task_with(run: &str) -> Task {
7440 let mut t = task();
7441 t.runs = vec![run.to_owned()];
7442 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
7443 t
7444 }
7445
7446 fn resume_action(run: &str) -> ask::ChoiceAction {
7447 ask::ChoiceAction::Resume { run: run.into() }
7448 }
7449
7450 #[test]
7451 fn decide_action_resumes_only_the_latest_resumable_run() {
7452 let t = held_task_with("r1");
7453 let q = action_question("r1", resume_action("r1"));
7454 let load = |s: RunState| move |_: &str| Ok(s);
7455 assert_eq!(
7456 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Blocked))),
7457 ActionDecision::Resume("r1".into())
7458 );
7459 let q_other = action_question("r1", resume_action("r0"));
7461 assert!(matches!(
7462 decide_action(
7463 &t,
7464 &q_other,
7465 &PHRASES_EN,
7466 load(run_state(RunStatus::Blocked))
7467 ),
7468 ActionDecision::Refuse(_)
7469 ));
7470 assert!(matches!(
7472 decide_action(&t, &q, &PHRASES_EN, load(run_state(RunStatus::Ready))),
7473 ActionDecision::Refuse(_)
7474 ));
7475 let mut released = run_state(RunStatus::Blocked);
7477 released.released_to = Some("elsewhere".into());
7478 assert!(matches!(
7479 decide_action(&t, &q, &PHRASES_EN, load(released)),
7480 ActionDecision::Refuse(_)
7481 ));
7482 assert!(matches!(
7484 decide_action(&t, &q, &PHRASES_EN, |_: &str| bail!("gone")),
7485 ActionDecision::Refuse(_)
7486 ));
7487 }
7488
7489 #[test]
7490 fn decide_action_ignores_a_question_about_an_earlier_run() {
7491 let mut t = held_task_with("r1");
7492 t.runs.push("r2".to_owned());
7493 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7494 assert_eq!(
7495 decide_action(
7496 &t,
7497 &action_question("r1", ask::ChoiceAction::Done),
7498 &PHRASES_EN,
7499 never
7500 ),
7501 ActionDecision::Stale
7502 );
7503 }
7504
7505 #[test]
7506 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
7507 let mut t = held_task_with("r1");
7508 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7509 assert_eq!(
7510 decide_action(
7511 &t,
7512 &action_question("r1", ask::ChoiceAction::Requeue),
7513 &PHRASES_EN,
7514 never
7515 ),
7516 ActionDecision::Requeue
7517 );
7518 let done_q = action_question("r1", ask::ChoiceAction::Done);
7519 assert_eq!(
7520 decide_action(&t, &done_q, &PHRASES_EN, never),
7521 ActionDecision::Done
7522 );
7523 t.mark_action_applied(&done_q.id);
7524 assert_eq!(
7525 decide_action(&t, &done_q, &PHRASES_EN, never),
7526 ActionDecision::Skip
7527 );
7528
7529 let mut plain = action_question("r1", ask::ChoiceAction::Done);
7531 plain.actions.clear();
7532 assert_eq!(
7533 decide_action(&held_task_with("r1"), &plain, &PHRASES_EN, never),
7534 ActionDecision::Skip
7535 );
7536 let mut running = held_task_with("r1");
7538 running.status = TaskStatus::Running;
7539 assert_eq!(
7540 decide_action(
7541 &running,
7542 &action_question("r1", ask::ChoiceAction::Done),
7543 &PHRASES_EN,
7544 never
7545 ),
7546 ActionDecision::Skip
7547 );
7548 }
7549
7550 #[test]
7551 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
7552 let dir = tempfile::tempdir().unwrap();
7553 let queue = Queue::at(dir.path().join("queue"));
7554 let questions = Questions::at(dir.path().join("questions"));
7555 let home = dir.path().join("home");
7556 let mut state = run_state(RunStatus::Blocked);
7557 state.id = "20260101-000000-act1".to_owned();
7558 state.save_under(&home).unwrap();
7559
7560 let mut t = held_task_with(&state.id);
7561 queue.put(&mut t).unwrap();
7562 let mut q = action_question(&state.id, resume_action(&state.id));
7563 questions.put(&mut q).unwrap();
7564
7565 apply_choice_actions(&queue, &questions, &home);
7566 let after = queue.get(&t.id).unwrap();
7567 assert_eq!(after.status, TaskStatus::Queued);
7568 assert!(!after.fresh_start);
7569 assert!(after.action_applied(&q.id));
7570 let pin = after.resume_override.clone().unwrap();
7571 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
7572 assert!(pin.forced);
7573
7574 let mut again = queue.get(&t.id).unwrap();
7576 again.hold_machine(Some("later".into()));
7577 queue.put(&mut again).unwrap();
7578 apply_choice_actions(&queue, &questions, &home);
7579 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7580 }
7581
7582 #[test]
7583 fn an_answer_the_waiter_already_delivered_is_not_acted_on_again() {
7584 let dir = tempfile::tempdir().unwrap();
7585 let queue = Queue::at(dir.path().join("queue"));
7586 let questions = Questions::at(dir.path().join("questions"));
7587 let home = dir.path().join("home");
7588 let mut state = run_state(RunStatus::Blocked);
7589 state.id = "20260101-000000-act2".to_owned();
7590 state.save_under(&home).unwrap();
7591
7592 let mut t = held_task_with(&state.id);
7593 queue.put(&mut t).unwrap();
7594 let mut q = action_question(&state.id, resume_action(&state.id));
7595 q.answer_delivered = true;
7596 questions.put(&mut q).unwrap();
7597
7598 apply_choice_actions(&queue, &questions, &home);
7599 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7600
7601 let mut q2 = action_question(&state.id, resume_action(&state.id));
7603 questions.put(&mut q2).unwrap();
7604 apply_choice_actions(&queue, &questions, &home);
7605 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7606 assert!(questions.get(&q2.id).unwrap().answer_delivered);
7607 }
7608}