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 task.hold_machine(Some(crate::queue::missing_blocker_hold_reason(
693 &task.blocked_by,
694 &missing,
695 )));
696 record(queue, &mut task);
697 continue;
698 }
699 let mut changed = false;
700 for id in task.blocked_by.clone() {
701 if let Ok(dep) = queue.get(&id) {
702 if dep.status == TaskStatus::Done {
703 task.unblock(&id);
704 changed = true;
705 }
706 continue;
707 }
708 if let Ok(q) = questions.get(&id)
709 && q.status == ask::QuestionStatus::Answered
710 {
711 let answer = match &q.answer {
712 Some(ask::Answer::Choice(c) | ask::Answer::Text(c)) => c.clone(),
713 None => String::new(),
714 };
715 task.record_answer(q.summary.clone(), answer);
716 task.unblock(&id);
717 changed = true;
718 }
719 }
720 if changed {
721 record(queue, &mut task);
722 }
723 }
724}
725
726#[derive(Debug, Clone, PartialEq, Eq)]
728enum ActionDecision {
729 Skip,
731 Resume(String),
733 Requeue,
735 Done,
737 Stale,
740 Refuse(String),
743}
744
745fn decide_action<F>(task: &Task, q: &ask::Question, load: F) -> ActionDecision
753where
754 F: FnOnce(&str) -> Result<RunState>,
755{
756 let Some(action) = q.chosen_action() else {
757 return ActionDecision::Skip;
758 };
759 if task.action_applied(&q.id)
760 || matches!(
761 task.status,
762 TaskStatus::Running | TaskStatus::Blocked | TaskStatus::Done
763 )
764 {
765 return ActionDecision::Skip;
766 }
767 if q.node != crate::conduct::NODE && task.runs.last() != Some(&q.run) {
770 return ActionDecision::Stale;
771 }
772 match action {
773 ask::ChoiceAction::Requeue => ActionDecision::Requeue,
774 ask::ChoiceAction::Done => ActionDecision::Done,
775 ask::ChoiceAction::Resume { run } => {
776 if task.runs.last() != Some(run) {
777 return ActionDecision::Refuse(format!(
778 "question {} asked to resume run {}, which is not this task's latest run",
779 q.short(),
780 ask::short_id(run)
781 ));
782 }
783 match load(run) {
784 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
785 ActionDecision::Resume(run.clone())
786 }
787 Ok(_) => ActionDecision::Refuse(format!(
788 "question {} asked to resume run {}, which cannot make progress",
789 q.short(),
790 ask::short_id(run)
791 )),
792 Err(e) => ActionDecision::Refuse(format!(
793 "question {} asked to resume run {}, which could not be read: {e:#}",
794 q.short(),
795 ask::short_id(run)
796 )),
797 }
798 }
799 }
800}
801
802fn task_of_question<'a>(tasks: &'a [Task], q: &ask::Question) -> Option<&'a Task> {
806 if q.node == crate::conduct::NODE {
807 return tasks.iter().find(|t| t.id == q.run);
808 }
809 tasks.iter().find(|t| t.runs.contains(&q.run))
810}
811
812fn apply_choice_actions(queue: &Queue, questions: &Questions, home: &Path) {
821 let tasks = queue.list();
822 for q in questions.list() {
823 if q.chosen_action().is_none() {
824 continue;
825 }
826 let Some(listed) = task_of_question(&tasks, &q) else {
827 continue;
828 };
829 if listed.action_applied(&q.id) {
830 continue;
831 }
832 let Ok(_claim) = queue.claim(&listed.id) else {
833 continue;
834 };
835 let Ok(mut task) = queue.get(&listed.id) else {
836 continue;
837 };
838 let decision = decide_action(&task, &q, |id| RunState::load_under(id, home));
839 let ran = matches!(
840 decision,
841 ActionDecision::Resume(_) | ActionDecision::Requeue | ActionDecision::Done
842 );
843 if ran {
844 if questions
849 .read_lease(&q.id)
850 .is_some_and(|l| l.fresh(Timestamp::now()))
851 {
852 continue;
853 }
854 let taken = questions.update(&q.id, |r| {
855 let free = !r.answer_delivered;
856 r.answer_delivered = true;
857 Ok(free)
858 });
859 if !matches!(taken, Ok((_, true))) {
860 continue;
861 }
862 }
863 match decision {
864 ActionDecision::Skip => continue,
865 ActionDecision::Resume(run) => {
866 task.release();
867 task.resume_override = Some(crate::queue::OperatorResume {
868 question_id: q.id.clone(),
869 at: Timestamp::now(),
870 conductor_rehold: None,
871 forced: true,
872 pinned_run: Some(run),
873 });
874 }
875 ActionDecision::Stale => {}
876 ActionDecision::Requeue => task.requeue(),
877 ActionDecision::Done => {
878 task.succeed();
879 supersede_prior_runs(&task, home);
880 }
881 ActionDecision::Refuse(why) => {
882 task.hold_machine(Some(why));
883 notices::raise(
884 Notice::warn(
885 &format!("action:{}", q.id),
886 "An answer asked the daemon to resume a run that cannot be resumed; the task stays held.",
887 )
888 .link(Link::Task {
889 id: task.id.clone(),
890 }),
891 );
892 }
893 }
894 task.mark_action_applied(&q.id);
895 record(queue, &mut task);
896 }
897}
898
899fn reconcile_task_questions(queue: &Queue, questions: &Questions) {
910 let tasks = queue.list();
911 let by_id: std::collections::BTreeMap<&str, &Task> =
912 tasks.iter().map(|t| (t.id.as_str(), t)).collect();
913 let referenced: std::collections::BTreeSet<&str> = tasks
914 .iter()
915 .flat_map(|task| task.blocked_by.iter().map(String::as_str))
916 .collect();
917
918 for mut question in questions.list() {
919 if !question.status.open() || question.node != crate::conduct::NODE {
920 continue;
921 }
922 if referenced.contains(question.id.as_str()) {
925 continue;
926 }
927 let Some(task) = by_id.get(question.run.as_str()) else {
928 continue;
929 };
930 question.abandon(format!(
931 "task {} no longer waits for this answer",
932 task.short()
933 ));
934 if let Err(e) = questions.put(&mut question) {
935 tracing::warn!(
936 "could not retire question {} for task {}: {e:#}",
937 question.short(),
938 task.short()
939 );
940 }
941 }
942}
943
944#[derive(Debug, Clone, Copy)]
951pub struct Verdict {
952 pub status: RunStatus,
954 pub left_pr: bool,
956 pub quota_hit: bool,
958 pub parked: bool,
960 pub no_viable_candidates: bool,
967}
968
969pub fn settle(task: &mut Task, verdict: Verdict, detail: &str, max_attempts: usize) {
1025 if verdict.parked {
1031 task.stall(detail);
1032 return;
1033 }
1034 match verdict.status {
1035 RunStatus::Merged | RunStatus::Ready => task.succeed(),
1036 RunStatus::Stalled if verdict.quota_hit => task.stall(detail),
1037 RunStatus::Failed if verdict.quota_hit && verdict.no_viable_candidates => {
1038 task.stall(detail)
1039 }
1040 RunStatus::Stalled | RunStatus::Failed => task.fail(detail, max_attempts),
1041 RunStatus::Blocked if verdict.left_pr => task.handed_off(detail),
1042 RunStatus::Blocked => task.fail(detail, max_attempts),
1043 RunStatus::VerifiedNoop => task.handed_off(detail),
1044 other => task.fail(
1045 format!(
1046 "the graph stopped at `{}` without reaching a terminal status: {detail}",
1047 label(other)
1048 ),
1049 max_attempts,
1050 ),
1051 }
1052}
1053
1054pub fn supersede_prior_runs(task: &Task, home: &Path) {
1103 let now = Timestamp::now();
1104 let last_run_succeeded = task
1111 .runs
1112 .last()
1113 .and_then(|id| RunState::load_under(id, home).ok())
1114 .is_some_and(|s| matches!(s.status, RunStatus::Merged | RunStatus::Ready));
1115 for id in task.superseded_attempts(last_run_succeeded) {
1116 let mut state = match RunState::load_under(id, home) {
1117 Ok(s) => s,
1118 Err(e) => {
1119 tracing::warn!("could not load run {id} to mark it superseded: {e:#}");
1120 continue;
1121 }
1122 };
1123 if !matches!(state.status, RunStatus::Blocked | RunStatus::Stalled) {
1124 continue;
1125 }
1126 let daemon_claims = is_working_on(home, id, now);
1127 if state.liveness(daemon_claims) == Liveness::Live {
1128 continue;
1129 }
1130 state.status = RunStatus::Superseded;
1131 if let Err(e) = state.save_under(home) {
1132 tracing::warn!("could not mark run {id} superseded: {e:#}");
1133 }
1134 }
1135}
1136
1137fn resweep_superseded_attempts(queue: &Queue, home: &Path) {
1159 for task in queue.list() {
1160 if task.status != TaskStatus::Done || task.runs.len() < 2 {
1161 continue;
1162 }
1163 supersede_prior_runs(&task, home);
1164 }
1165}
1166
1167fn settle_and_diagnose(
1174 task: &mut Task,
1175 verdict: Verdict,
1176 detail: &str,
1177 max_attempts: usize,
1178 state: &RunState,
1179) {
1180 settle(task, verdict, detail, max_attempts);
1181 if task.status == TaskStatus::Held {
1182 task.diagnostic = diagnostic(state);
1183 note_open_question(task, &state.id);
1184 }
1185}
1186
1187fn note_open_question(task: &mut Task, run: &str) {
1202 let Some(home) = crate::run::try_home() else {
1203 return;
1204 };
1205 let open = Questions::at(home.join("questions")).open_for(run);
1206 let Some(q) = open.first() else {
1207 return;
1208 };
1209 let base = task.hold_reason.clone().unwrap_or_default();
1210 task.hold_reason = Some(format!(
1211 "{base} - waiting for operator answer to question {}",
1212 q.short()
1213 ));
1214}
1215
1216fn reclaim(task: &mut Task, last_run: Option<RunState>, max_attempts: usize) {
1227 match last_run {
1228 Some(state) => {
1229 let verdict = Verdict {
1230 status: state.status,
1231 left_pr: state.pr.is_some(),
1232 quota_hit: !state.quota.is_empty(),
1233 parked: state.parked,
1234 no_viable_candidates: state.viable().is_empty(),
1235 };
1236 let detail = format!(
1237 "recovered a `running` task whose daemon never recorded the outcome: {}",
1238 describe(&state)
1239 );
1240 settle_and_diagnose(task, verdict, &detail, max_attempts, &state);
1241 }
1242 None => {
1243 let why = "task was `running` with no live daemon and no readable \
1244 run to recover; held for a human to check what happened";
1245 task.last_error = Some(why.to_owned());
1246 task.hold_machine(Some(why.to_owned()));
1249 }
1250 }
1251}
1252
1253fn reclaim_orphaned_running(queue: &Queue, max_attempts: usize) -> Vec<String> {
1275 let mut reclaimed = Vec::new();
1276 for listed in queue.list() {
1277 if listed.status != TaskStatus::Running {
1278 continue;
1279 }
1280 let Ok(_claim) = queue.claim(&listed.id) else {
1281 continue;
1282 };
1283 let Ok(mut task) = queue.get(&listed.id) else {
1287 continue;
1288 };
1289 if task.status != TaskStatus::Running {
1290 continue;
1291 }
1292 let last_run = task.runs.last().and_then(|id| RunState::load(id).ok());
1293 if let Some(state) = &last_run
1303 && let Err(e) = ask::Questions::open().settle_run(&state.id, state.status)
1304 {
1305 tracing::warn!("abandon questions for {}: {e:#}", state.id);
1306 }
1307 reclaim(&mut task, last_run, max_attempts);
1308 if task.status == TaskStatus::Done {
1309 supersede_prior_runs(&task, &crate::run::home());
1310 }
1311 record(queue, &mut task);
1312 reclaimed.push(task.id.clone());
1313 }
1314 reclaimed
1315}
1316
1317fn reclaim_abandoned_runs(home: &Path, now: Timestamp) -> Vec<String> {
1349 reclaim_abandoned_runs_with(
1350 home,
1351 now,
1352 crate::proc::pid_status,
1353 crate::proc::process_started_at,
1354 )
1355}
1356
1357fn reclaim_abandoned_runs_with<F, G>(
1363 home: &Path,
1364 now: Timestamp,
1365 query: F,
1366 identity: G,
1367) -> Vec<String>
1368where
1369 F: Fn(u32) -> Option<bool>,
1370 G: Fn(u32) -> Option<String>,
1371{
1372 let mut abandoned = Vec::new();
1373 for entry in std::fs::read_dir(home.join("runs"))
1374 .into_iter()
1375 .flatten()
1376 .flatten()
1377 {
1378 let id = entry.file_name().to_string_lossy().into_owned();
1379 if !crate::run::is_run_id(&id) {
1380 continue;
1381 }
1382 let Ok(body) = std::fs::read_to_string(entry.path().join("run.json")) else {
1390 continue;
1391 };
1392 let Ok(mut state) = serde_json::from_str::<RunState>(&body) else {
1393 continue;
1394 };
1395 if state.status.done() || !state.active_all_overrun(now) {
1396 continue;
1397 }
1398 let daemon_claims = is_working_on(home, &id, now);
1410 if state.liveness_with(daemon_claims, &query, &identity) != crate::run::Liveness::Dead {
1411 continue;
1412 }
1413 state.abandon("daemon");
1414 if let Err(e) = state.save_under(home) {
1415 tracing::warn!("could not persist abandoned run {id}: {e:#}");
1416 continue;
1417 }
1418 if let Some(notice) = notices::run_ended(&state) {
1421 notices::raise_in(home, notice);
1422 }
1423 if let Err(e) = Questions::at(home.join("questions")).settle_run(&id, state.status) {
1431 tracing::warn!("abandon questions for {id}: {e:#}");
1432 }
1433 abandoned.push(id);
1434 }
1435 abandoned
1436}
1437
1438pub async fn serve(opts: Opts) -> Result<()> {
1444 serve_until(opts, Stop::new()).await
1445}
1446
1447pub async fn serve_until(opts: Opts, stop: Stop) -> Result<()> {
1464 let signal = {
1465 let stop = stop.clone();
1466 tokio::spawn(async move {
1467 if tokio::signal::ctrl_c().await.is_ok() {
1468 stop.stop();
1469 tracing::info!("shutdown requested; a run in flight will be finished first");
1470 }
1471 })
1472 };
1473
1474 let worktrees_root = opts
1475 .worktrees_root
1476 .clone()
1477 .unwrap_or_else(crate::run::default_worktree_root);
1478 let outcome = drive(
1479 &opts,
1480 &Queue::open(),
1481 &status_path(),
1482 &crate::run::home(),
1483 &worktrees_root,
1484 &stop,
1485 )
1486 .await;
1487
1488 signal.abort();
1489 outcome
1490}
1491
1492async fn drive(
1505 opts: &Opts,
1506 queue: &Queue,
1507 status_file: &Path,
1508 home: &Path,
1509 worktrees_root: &Path,
1510 stop: &Stop,
1511) -> Result<()> {
1512 let status = Arc::new(Mutex::new(Status::new()));
1520 write_status_to(status_file, &lock(&status)).context("publish the daemon status file")?;
1521 let beat = tokio::spawn(heartbeat(Arc::clone(&status), status_file.to_path_buf()));
1522
1523 let daemon_cfg = prepare(&opts.repo, opts)
1529 .map(|c| c.daemon)
1530 .unwrap_or_default();
1531 let concurrency = max_concurrent(daemon_cfg.max_concurrent_runs);
1532
1533 let waiter = tokio::spawn(crate::waiter::run(
1538 crate::waiter::Waiter::new(
1539 crate::ask::Questions::at(home.join("questions")),
1540 home.to_path_buf(),
1541 prepare(&opts.repo, opts).ok(),
1542 ),
1543 stop.clone(),
1544 ));
1545
1546 tracing::info!(
1547 "magi serve: queue {} (poll {}s, {} attempts per task, {} run(s) at once{})",
1548 queue.root().display(),
1549 opts.poll.as_secs(),
1550 opts.max_attempts,
1551 concurrency,
1552 if daemon_cfg.pause_for_interrupts {
1553 ", interrupts enabled"
1554 } else {
1555 ""
1556 }
1557 );
1558
1559 janitor(&opts.repo, opts, home, worktrees_root).await;
1562 resweep_superseded_attempts(queue, home);
1563
1564 let outcome = poll(
1565 opts,
1566 queue,
1567 &status,
1568 home,
1569 worktrees_root,
1570 stop,
1571 DispatchLimits {
1572 max_concurrent: concurrency,
1573 pause_for_interrupts: daemon_cfg.pause_for_interrupts,
1574 },
1575 )
1576 .await;
1577
1578 beat.abort();
1579 waiter.abort();
1580 clear_status_at(status_file);
1581 outcome
1582}
1583
1584async fn heartbeat(status: Arc<Mutex<Status>>, path: PathBuf) {
1590 loop {
1591 tokio::time::sleep(HEARTBEAT).await;
1592 let snapshot = {
1593 let mut guard = lock(&status);
1594 guard.updated_at = Timestamp::now();
1595 guard.clone()
1596 };
1597 if let Err(e) = write_status_to(&path, &snapshot) {
1598 tracing::warn!("could not refresh the daemon status file: {e:#}");
1601 }
1602 }
1603}
1604
1605#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1608enum LandResume {
1609 NotLanding,
1612 StillWaiting,
1617 Ready,
1621}
1622
1623fn land_resume_state(task: &Task) -> LandResume {
1627 let Some(run_id) = task.runs.last() else {
1628 return LandResume::NotLanding;
1629 };
1630 let Ok(state) = RunState::load(run_id) else {
1631 return LandResume::NotLanding;
1632 };
1633 if state.status != RunStatus::Landing || !state.parked {
1634 return LandResume::NotLanding;
1635 }
1636 let store = ask::Questions::open();
1637 let waiting = store
1638 .list()
1639 .into_iter()
1640 .filter(|q| &q.run == run_id && q.node == land::APPROVAL_NODE)
1641 .max_by(|a, b| a.id.cmp(&b.id));
1642 let Some(mut q) = waiting else {
1643 return LandResume::Ready;
1644 };
1645 if !q.status.open() {
1646 return LandResume::Ready;
1647 }
1648 let timeout = Duration::from_secs(state.config.graph.answer_timeout);
1655 let elapsed = Timestamp::now().as_second() - q.asked_at.as_second();
1656 if elapsed >= 0 && elapsed as u64 >= timeout.as_secs() {
1657 q.abandon(format!(
1658 "no answer within {}s of asking",
1659 timeout.as_secs().max(1)
1660 ));
1661 if store.put(&mut q).is_ok() {
1664 return LandResume::Ready;
1665 }
1666 }
1667 LandResume::StillWaiting
1668}
1669
1670const RECHECK_WHILE_BUSY: Duration = Duration::from_millis(200);
1678
1679const CACHE_CHECK_INTERVAL_SECS: u64 = 5 * 60;
1691
1692struct InFlightGuard<'a> {
1705 status: &'a Arc<Mutex<Status>>,
1706 stop: &'a Stop,
1707 task_id: &'a str,
1708}
1709
1710impl Drop for InFlightGuard<'_> {
1711 fn drop(&mut self) {
1712 lock(self.status).current.retain(|c| c.task != self.task_id);
1713 self.stop.exit();
1714 }
1715}
1716
1717#[derive(Debug, Clone, PartialEq, Eq)]
1739enum Interrupt {
1740 Idle,
1742 Parking {
1754 parked: Vec<String>,
1755 interrupt_task: String,
1756 },
1757 Running {
1765 parked: Vec<String>,
1766 interrupt_task: String,
1767 },
1768 Resuming { parked: Vec<String> },
1775}
1776
1777fn advance_interrupt(state: Interrupt, in_flight: &[String], runnable: &[Task]) -> Interrupt {
1798 match state {
1799 Interrupt::Idle => {
1800 if in_flight.len() != 1 {
1812 return Interrupt::Idle;
1813 }
1814 match runnable.iter().find(|t| t.interrupt) {
1815 Some(t) => Interrupt::Parking {
1816 parked: in_flight.to_vec(),
1817 interrupt_task: t.id.clone(),
1818 },
1819 None => Interrupt::Idle,
1820 }
1821 }
1822 Interrupt::Parking {
1823 parked,
1824 interrupt_task,
1825 } => {
1826 if in_flight.iter().any(|id| parked.contains(id)) {
1827 Interrupt::Parking {
1829 parked,
1830 interrupt_task,
1831 }
1832 } else if in_flight.contains(&interrupt_task) {
1833 Interrupt::Running {
1834 parked,
1835 interrupt_task,
1836 }
1837 } else if runnable.iter().any(|t| t.id == interrupt_task) {
1838 Interrupt::Parking {
1842 parked,
1843 interrupt_task,
1844 }
1845 } else {
1846 Interrupt::Resuming { parked }
1851 }
1852 }
1853 Interrupt::Running {
1854 parked,
1855 interrupt_task,
1856 } => {
1857 if in_flight.contains(&interrupt_task) {
1858 Interrupt::Running {
1859 parked,
1860 interrupt_task,
1861 }
1862 } else {
1863 Interrupt::Resuming { parked }
1869 }
1870 }
1871 Interrupt::Resuming { parked } => {
1872 if in_flight.iter().any(|id| parked.contains(id)) {
1873 Interrupt::Idle
1879 } else if runnable.iter().any(|t| parked.contains(&t.id)) {
1880 Interrupt::Resuming { parked }
1881 } else {
1882 Interrupt::Idle
1885 }
1886 }
1887 }
1888}
1889
1890fn advance_interrupt_tick(
1896 enabled: bool,
1897 state: Interrupt,
1898 in_flight: &[String],
1899 runnable: &[Task],
1900) -> Interrupt {
1901 if !enabled {
1902 return Interrupt::Idle;
1903 }
1904 advance_interrupt(state, in_flight, runnable)
1905}
1906
1907fn interrupt_gate(state: &Interrupt, in_flight: &[String], candidates: Vec<Task>) -> Vec<Task> {
1912 match state {
1913 Interrupt::Idle => candidates,
1914 Interrupt::Parking {
1915 parked,
1916 interrupt_task,
1917 } => {
1918 if in_flight.iter().any(|id| parked.contains(id)) {
1919 Vec::new()
1920 } else {
1921 candidates
1922 .into_iter()
1923 .filter(|t| &t.id == interrupt_task)
1924 .collect()
1925 }
1926 }
1927 Interrupt::Running { .. } => Vec::new(),
1928 Interrupt::Resuming { parked } => candidates
1936 .into_iter()
1937 .find(|t| parked.contains(&t.id))
1938 .into_iter()
1939 .collect(),
1940 }
1941}
1942
1943struct DispatchLimits {
1947 max_concurrent: usize,
1950 pause_for_interrupts: bool,
1952}
1953
1954#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1963enum PermitKind {
1964 None,
1969 Urgent,
1976 Ordinary,
1979}
1980
1981fn permit_kind(priority: bool, urgent: bool) -> PermitKind {
1984 if priority {
1985 PermitKind::None
1986 } else if urgent {
1987 PermitKind::Urgent
1988 } else {
1989 PermitKind::Ordinary
1990 }
1991}
1992
1993async fn poll(
2012 opts: &Opts,
2013 queue: &Queue,
2014 status: &Arc<Mutex<Status>>,
2015 home: &Path,
2016 worktrees_root: &Path,
2017 stop: &Stop,
2018 limits: DispatchLimits,
2019) -> Result<()> {
2020 let DispatchLimits {
2021 max_concurrent,
2022 pause_for_interrupts,
2023 } = limits;
2024 let mut attempted: Vec<String> = Vec::new();
2029 let sem = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
2030 let urgent_sem = Arc::new(tokio::sync::Semaphore::new(1));
2038 let quota_cooldown_until: Arc<Mutex<Option<Timestamp>>> = Arc::new(Mutex::new(None));
2044 let mut inflight: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
2045 let mut conductor = Conductor::new();
2046 let mut cache_last_checked: Option<Timestamp> = None;
2049 let mut interrupt = Interrupt::Idle;
2051 let mut interrupt_pauses: std::collections::HashMap<String, crate::graph::Pause> =
2057 std::collections::HashMap::new();
2058
2059 while !stop.stopped() {
2060 lock(status).polls += 1;
2061
2062 while let Some(result) = inflight.try_join_next() {
2067 if let Err(e) = result {
2068 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2069 notices::raise(Notice::error(
2070 "loop:attempt",
2071 "A queued attempt ended abnormally; check the task it was running.",
2072 ));
2073 }
2074 }
2075
2076 let swept = sweep_stale_claims(queue, STALE_CLAIM);
2077 if !swept.is_empty() {
2078 tracing::warn!(
2079 "swept {} stale claim(s) left behind by an earlier daemon: {}",
2080 swept.len(),
2081 swept.join(", ")
2082 );
2083 }
2084 let now = Timestamp::now();
2089
2090 if !stop.busy_now() {
2095 maybe_prune_cache_between_runs(
2096 &opts.repo,
2097 opts,
2098 home,
2099 stop,
2100 &mut cache_last_checked,
2101 now,
2102 )
2103 .await;
2104 }
2105
2106 let stalled = stalled_tasks(queue, home, now);
2107 let stalled_ids: std::collections::BTreeSet<_> =
2108 stalled.iter().map(|task| task.id.clone()).collect();
2109 let reclaimed = reclaim_orphaned_running(queue, opts.max_attempts);
2110 if !reclaimed.is_empty() {
2111 tracing::warn!(
2112 "reclaimed {} task(s) left `running` by a daemon that never \
2113 recorded the outcome: {}",
2114 reclaimed.len(),
2115 reclaimed.join(", ")
2116 );
2117 }
2118 let abandoned_runs = reclaim_abandoned_runs(home, now);
2119 if !abandoned_runs.is_empty() {
2120 tracing::warn!(
2121 "failed {} run(s) left behind by a killed process, past every \
2122 active seat's own timeout: {}",
2123 abandoned_runs.len(),
2124 abandoned_runs.join(", ")
2125 );
2126 }
2127
2128 let questions = Questions::at(home.join("questions"));
2133
2134 resolve_blockers(queue, &questions);
2137 apply_choice_actions(queue, &questions, home);
2138 reconcile_task_questions(queue, &questions);
2139
2140 let finished: Vec<Task> = finished_tasks(queue)
2146 .into_iter()
2147 .filter(|task| !stalled_ids.contains(&task.id))
2148 .filter(|task| !awaiting_resume_with(task, |id| RunState::load_under(id, home)))
2152 .collect();
2153 let queued = queued_tasks(queue);
2154 if !(queued.is_empty() && stalled.is_empty() && finished.is_empty())
2158 && conductor.worth_a_look(queue, &stalled, &finished)
2159 {
2160 match prepare(&opts.repo, opts) {
2161 Ok(cfg) => {
2162 conductor
2163 .maybe_run(
2164 &cfg,
2165 &opts.repo,
2166 queue,
2167 &questions,
2168 home,
2169 &queued,
2170 &stalled,
2171 &finished,
2172 opts.max_attempts,
2173 )
2174 .await;
2175 }
2176 Err(e) => {
2177 tracing::warn!("conductor: no config: {e:#}");
2178 notices::raise(Notice::warn(
2179 "loop:no-config",
2180 "The loop could not read this repository's config, so held tasks are not being triaged.",
2181 ));
2182 }
2183 }
2184 }
2185
2186 let candidates: Vec<Task> = runnable(queue)
2187 .into_iter()
2188 .filter(|t| !opts.once || !attempted.contains(&t.id))
2189 .collect();
2190
2191 let in_flight: Vec<String> = lock(status)
2196 .current
2197 .iter()
2198 .map(|c| c.task.clone())
2199 .collect();
2200 interrupt_pauses.retain(|id, _| in_flight.contains(id));
2201
2202 interrupt =
2203 advance_interrupt_tick(pause_for_interrupts, interrupt, &in_flight, &candidates);
2204 if let Interrupt::Parking {
2205 parked,
2206 interrupt_task,
2207 } = &interrupt
2208 {
2209 let reason = format!(
2210 "task {} asked to run first",
2211 crate::run::short_of(interrupt_task)
2212 );
2213 for id in parked {
2214 if let Some(pause) = interrupt_pauses.get(id) {
2215 pause.park_because(reason.clone());
2216 }
2217 }
2218 }
2219 let candidates = interrupt_gate(&interrupt, &in_flight, candidates);
2220
2221 let cooling_down =
2222 lock("a_cooldown_until).is_some_and(|until| Timestamp::now() < until);
2223
2224 let mut started_any = false;
2225 for candidate in candidates {
2226 if stop.stopped() {
2227 break;
2228 }
2229
2230 let resume = land_resume_state(&candidate);
2231 if resume == LandResume::StillWaiting {
2232 continue;
2233 }
2234 let priority = resume == LandResume::Ready;
2235
2236 if !priority && cooling_down {
2237 continue;
2238 }
2239 let permit = match permit_kind(priority, candidate.urgent) {
2240 PermitKind::None => None,
2241 PermitKind::Urgent => match Arc::clone(&urgent_sem).try_acquire_owned() {
2242 Ok(p) => Some(p),
2243 Err(_) => continue,
2249 },
2250 PermitKind::Ordinary => match Arc::clone(&sem).try_acquire_owned() {
2251 Ok(p) => Some(p),
2252 Err(_) => continue,
2256 },
2257 };
2258
2259 let Ok(claim) = queue.claim(&candidate.id) else {
2264 tracing::info!("task {} is claimed elsewhere; skipping", candidate.short());
2265 continue;
2266 };
2267 let mut task = match queue.get(&candidate.id) {
2270 Ok(t) if t.status.runnable() => t,
2271 Ok(_) => continue,
2272 Err(e) => {
2273 tracing::warn!("could not re-read task {}: {e:#}", candidate.short());
2274 continue;
2275 }
2276 };
2277 let task_id = task.id.clone();
2278 attempted.push(task_id.clone());
2279 lock(status).idle = false;
2280 stop.enter();
2283 started_any = true;
2284
2285 let run_pause = crate::graph::Pause::new();
2289 interrupt_pauses.insert(task_id.clone(), run_pause.clone());
2290
2291 let opts = opts.clone();
2292 let queue = queue.clone();
2293 let status = Arc::clone(status);
2294 let stop = stop.clone();
2295 let quota_cooldown_until = Arc::clone("a_cooldown_until);
2296 inflight.spawn(async move {
2297 let _claim = claim;
2301 let _permit = permit;
2302 let _inflight = InFlightGuard {
2304 status: &status,
2305 stop: &stop,
2306 task_id: &task_id,
2307 };
2308 let quota = attempt(&opts, &queue, &status, &stop, run_pause, &mut task).await;
2309 lock(&status).completed += 1;
2310 let now = Timestamp::now();
2316 if let Some(until) = cooldown_until("a, now) {
2317 let wait = until.as_second() - now.as_second();
2318 *lock("a_cooldown_until) = Some(until);
2319 let hint = quota
2320 .iter()
2321 .find(|q| q.reset.is_some())
2322 .and_then(|q| q.reset.as_deref());
2323 match hint {
2324 Some(h) => tracing::warn!(
2325 "quota hit; waiting {wait}s before taking another ordinary task \
2326 (CLI reported reset: {h})"
2327 ),
2328 None => tracing::warn!(
2329 "quota hit; waiting {wait}s before taking another ordinary task \
2330 (no reset hint reported)"
2331 ),
2332 }
2333 }
2334 });
2335 }
2336
2337 if started_any {
2338 continue;
2339 }
2340
2341 if stop.busy_now() {
2342 stop.idle(RECHECK_WHILE_BUSY.min(opts.poll)).await;
2347 continue;
2348 }
2349
2350 lock(status).idle = true;
2352 if opts.once {
2353 janitor(&opts.repo, opts, home, worktrees_root).await;
2357 resweep_superseded_attempts(queue, home);
2358 triage_held(queue, home, opts).await;
2359 break;
2360 }
2361 stop.idle(opts.poll).await;
2362 if stop.stopped() {
2363 continue;
2364 }
2365 janitor(&opts.repo, opts, home, worktrees_root).await;
2371 resweep_superseded_attempts(queue, home);
2372 triage_held(queue, home, opts).await;
2373 }
2374
2375 while let Some(result) = inflight.join_next().await {
2380 if let Err(e) = result {
2381 tracing::error!("a spawned attempt did not finish cleanly: {e}");
2382 notices::raise(Notice::error(
2383 "loop:attempt",
2384 "A queued attempt ended abnormally; check the task it was running.",
2385 ));
2386 }
2387 }
2388 Ok(())
2389}
2390
2391async fn attempt(
2397 opts: &Opts,
2398 queue: &Queue,
2399 status: &Arc<Mutex<Status>>,
2400 stop: &Stop,
2401 interrupt_pause: crate::graph::Pause,
2402 task: &mut Task,
2403) -> Vec<QuotaLoss> {
2404 let repo = repo_for(task, &opts.repo);
2405 tracing::info!(
2406 "task {} — {} (repo {})",
2407 task.short(),
2408 task.title,
2409 repo.display()
2410 );
2411
2412 let mut config = match prepare(&repo, opts) {
2413 Ok(c) => c,
2414 Err(e) => {
2415 task.attempts += 1;
2419 task.fail(format!("config: {e:#}"), opts.max_attempts);
2420 record(queue, task);
2421 return Vec::new();
2422 }
2423 };
2424 apply_solo(&mut config, task);
2425
2426 if let Some(reason) = disk_gate(&repo, &config) {
2434 task.last_error = Some(reason.clone());
2435 task.hold_machine(Some(reason.clone()));
2436 record(queue, task);
2437 tracing::warn!("holding {} for want of disk space: {reason}", task.short());
2438 notices::raise(
2441 Notice::warn(
2442 &format!("disk:{}", repo.display()),
2443 "A task was held for want of free disk space; free some, then release it from the queue.",
2444 )
2445 .link(Link::Task {
2446 id: task.id.clone(),
2447 }),
2448 );
2449 return Vec::new();
2450 }
2451
2452 let unfinished = (!task.fresh_start)
2472 .then(|| unfinished_run(&task.runs, task.short()))
2473 .flatten();
2474 let review_branch = task.review_branch.take();
2480 let branch_exists = match &review_branch {
2481 Some(branch) => crate::git::branch_exists(&repo, branch)
2482 .await
2483 .unwrap_or(false),
2484 None => false,
2485 };
2486 let starter = choose_starter(
2487 review_branch.as_deref(),
2488 branch_exists,
2489 unfinished.as_deref(),
2490 );
2491 let attachments = match task_attachments(queue, task) {
2496 Ok(a) => a,
2497 Err(e) => {
2498 task.attempts += 1;
2499 task.fail(format!("could not start the run: {e:#}"), opts.max_attempts);
2500 record(queue, task);
2501 return Vec::new();
2502 }
2503 };
2504 let started = match &starter {
2505 Starter::Review(branch) => {
2506 tracing::info!(
2507 "task {} reopens `{branch}` as a review-only pass",
2508 task.short()
2509 );
2510 let takeover = crate::handover::Takeover {
2513 earlier: task.earlier_attempts().to_vec(),
2514 home: crate::run::home(),
2515 choice: take_divergence_answer(branch, &config.merge.remote, task),
2516 };
2517 Runner::review_taking_over(&repo, branch, config, Some(takeover)).await
2518 }
2519 Starter::Resume(id) => {
2520 tracing::info!("resuming run {id} rather than competing again");
2521 Runner::resume(id).map(|mut r| {
2522 if let Some(instruction) =
2523 prepare_instruction(&starter, Some(&r.state.instruction), task)
2524 {
2525 r.state.instruction = instruction;
2526 }
2527 r.state.attachments = attachments.clone();
2528 r
2529 })
2530 }
2531 Starter::Start => {
2532 if let Some(branch) = &review_branch {
2533 tracing::warn!(
2534 "conductor chose review for task {} but branch `{branch}` no longer \
2535 exists; requeuing as a fresh competition instead",
2536 task.short()
2537 );
2538 }
2539 let instruction = prepare_instruction(&starter, None, task)
2540 .unwrap_or_else(|| task.instruction.clone());
2541 Runner::start_naming(&repo, instruction, &task.title, config)
2542 .await
2543 .map(|mut r| {
2544 r.state.attachments = attachments.clone();
2545 r
2546 })
2547 }
2548 };
2549 let mut runner = match started {
2550 Ok(r) => r,
2551 Err(e) if e.downcast_ref::<crate::handover::Refused>().is_some() => {
2555 let reason = format!("could not start the run: {e:#}");
2556 task.last_error = Some(reason.clone());
2557 task.hold_machine(Some(reason));
2558 record(queue, task);
2559 tracing::warn!(
2560 "holding {} for a branch it cannot take over: {e:#}",
2561 task.short()
2562 );
2563 notices::raise(
2565 Notice::warn(
2566 &format!("handover:{}", task.id),
2567 "A task was held because an earlier attempt still has its branch checked out; see the task's hold reason, then release it from the queue.",
2568 )
2569 .link(Link::Task {
2570 id: task.id.clone(),
2571 })
2572 .about([task.id.clone()]),
2573 );
2574 return Vec::new();
2575 }
2576 Err(e) if e.downcast_ref::<crate::reconcile::Diverged>().is_some() => {
2581 let d = e
2582 .downcast_ref::<crate::reconcile::Diverged>()
2583 .expect("checked by the guard");
2584 let mut q = ask::Question::new(
2585 task.id.clone(),
2586 "review".to_owned(),
2587 "sync".to_owned(),
2588 d.summary(),
2589 d.detail(),
2590 d.choices(),
2591 );
2592 match Questions::open().put(&mut q) {
2593 Ok(()) => {
2594 task.last_error = Some(format!("{e:#}"));
2595 task.review_branch = Some(d.branch.clone());
2598 task.block(vec![q.id.clone()], Some(d.summary()));
2599 }
2600 Err(put) => {
2601 tracing::warn!("could not file the divergence question: {put:#}");
2602 task.attempts += 1;
2603 task.fail(format!("could not start the run: {e:#}"), opts.max_attempts);
2604 }
2605 }
2606 record(queue, task);
2607 return Vec::new();
2608 }
2609 Err(e) => {
2610 if e.downcast_ref::<crate::reconcile::Stale>().is_some()
2616 && let Starter::Review(branch) = &starter
2617 {
2618 task.review_branch = Some(branch.clone());
2619 task.last_error = Some(format!("could not start the run: {e:#}"));
2620 task.status = crate::queue::TaskStatus::Failed;
2621 record(queue, task);
2622 return Vec::new();
2623 }
2624 task.attempts += 1;
2625 task.fail(format!("could not start the run: {e:#}"), opts.max_attempts);
2626 record(queue, task);
2627 return Vec::new();
2628 }
2629 };
2630 runner.on_pause(stop.pause());
2632 runner.watch_interrupt(interrupt_pause);
2636
2637 let run = runner.state.id.clone();
2640 task.start(run.clone());
2641 record(queue, task);
2642 lock(status).current.push(Current {
2643 task: task.id.clone(),
2644 run,
2645 });
2646
2647 let quota_before = runner.state.quota.clone();
2650 let detail = match runner.execute().await {
2651 Ok(()) => describe(&runner.state),
2652 Err(e) => format!("{e:#}"),
2653 };
2654 let fresh = losses_this_attempt("a_before, &runner.state.quota);
2655 let verdict = Verdict {
2656 status: runner.state.status,
2657 left_pr: runner.state.pr.is_some(),
2660 quota_hit: !fresh.is_empty(),
2666 parked: runner.state.parked,
2670 no_viable_candidates: runner.state.viable().is_empty(),
2673 };
2674 settle_and_diagnose(task, verdict, &detail, opts.max_attempts, &runner.state);
2675 if task.status == TaskStatus::Done {
2676 supersede_prior_runs(task, &crate::run::home());
2677 }
2678 record(queue, task);
2679 tracing::info!(
2680 "task {} is {} after run {} ({})",
2681 task.short(),
2682 task.status.as_str(),
2683 runner.state.short(),
2684 label(runner.state.status)
2685 );
2686 fresh
2687}
2688
2689fn losses_this_attempt(before: &[QuotaLoss], after: &[QuotaLoss]) -> Vec<QuotaLoss> {
2699 after
2700 .iter()
2701 .filter(|q| !before.contains(q))
2702 .cloned()
2703 .collect()
2704}
2705
2706fn cooldown_until(quota: &[QuotaLoss], now: Timestamp) -> Option<Timestamp> {
2709 if quota.is_empty() {
2710 return None;
2711 }
2712 let with_hint = quota.iter().find(|q| q.reset.is_some());
2713 let reset_at = with_hint.and_then(|q| parse_reset_hint(q.reset.as_deref()?, now, q.at));
2714 let wait = quota_wait(reset_at, now, QUOTA_WAIT_FALLBACK, QUOTA_WAIT_CAP);
2715 let secs = i64::try_from(wait.as_secs()).unwrap_or(i64::MAX);
2716 Some(
2717 now.checked_add(jiff::SignedDuration::from_secs(secs))
2718 .unwrap_or(Timestamp::MAX),
2719 )
2720}
2721
2722fn apply_solo(config: &mut Config, task: &Task) {
2732 if task.solo {
2733 config.graph.candidates = 1;
2734 }
2735}
2736
2737fn prepare(repo: &Path, opts: &Opts) -> Result<Config> {
2739 let (mut config, _layers) = Config::discover(repo, opts.config.as_deref())?;
2740 if let Some(mode) = &opts.merge {
2741 config.merge.mode = merge_mode(mode)?;
2742 }
2743 Ok(config)
2744}
2745
2746async fn maybe_prune_cache_between_runs(
2777 repo: &Path,
2778 opts: &Opts,
2779 home: &Path,
2780 stop: &Stop,
2781 last_checked: &mut Option<Timestamp>,
2782 now: Timestamp,
2783) {
2784 if stop.stopped() || !cache_check_due(*last_checked, now, CACHE_CHECK_INTERVAL_SECS) {
2785 return;
2786 }
2787 *last_checked = Some(now);
2788 let cfg = match prepare(repo, opts) {
2789 Ok(cfg) => cfg,
2790 Err(e) => {
2791 tracing::warn!("cache check: no config: {e:#}");
2792 return;
2793 }
2794 };
2795 match clean::prune_cache_if_over_limit(&cfg, home) {
2796 Ok(Some(pruned)) if pruned.files > 0 => tracing::info!(
2797 "housekeep: pruned {} file(s) ({} bytes) from the shared cache between runs",
2798 pruned.files,
2799 pruned.freed
2800 ),
2801 Ok(_) => {}
2802 Err(e) => {
2803 tracing::warn!("housekeep: prune cache: {e:#}");
2804 notices::raise_in(
2805 home,
2806 Notice::warn(
2807 "housekeep:cache",
2808 "Pruning the shared build cache failed; disk usage may keep growing.",
2809 ),
2810 );
2811 }
2812 }
2813}
2814
2815fn cache_check_due(last_checked: Option<Timestamp>, now: Timestamp, interval_secs: u64) -> bool {
2819 last_checked.is_none_or(|last| clean::due(now, last, interval_secs))
2820}
2821
2822async fn janitor(repo: &Path, opts: &Opts, home: &Path, worktrees_root: &Path) {
2841 let cfg = match prepare(repo, opts) {
2842 Ok(cfg) => cfg,
2843 Err(e) => {
2844 tracing::warn!("housekeep: no config: {e:#}");
2845 return;
2846 }
2847 };
2848 let worktrees_root = cfg.graph.worktree_root.as_deref().unwrap_or(worktrees_root);
2857 let out = clean::housekeep(&cfg, home, worktrees_root, repo, Timestamp::now()).await;
2858 if out.folded > 0 || out.unreadable > 0 || out.orphaned_worktrees > 0 {
2863 let mut extra = Vec::new();
2864 if out.unreadable > 0 {
2865 extra.push(format!("{} unreadable", out.unreadable));
2866 }
2867 if out.orphaned_worktrees > 0 {
2868 extra.push(format!("{} orphaned worktree(s)", out.orphaned_worktrees));
2869 }
2870 let detail = if extra.is_empty() {
2871 String::new()
2872 } else {
2873 format!(" ({})", extra.join(", "))
2874 };
2875 tracing::info!("housekeep: folded {} run(s){detail}", out.folded);
2876 }
2877 if out.external_merges_recorded > 0 {
2878 tracing::info!(
2879 "housekeep: recorded {} run(s) as merged externally",
2880 out.external_merges_recorded
2881 );
2882 }
2883 if out.cache_files > 0 {
2884 tracing::info!(
2885 "housekeep: pruned {} file(s) ({} bytes) from the shared cache",
2886 out.cache_files,
2887 out.cache_freed
2888 );
2889 }
2890 if out.questions_abandoned > 0 {
2891 tracing::info!(
2892 "housekeep: abandoned {} question(s) left open by a finished run",
2893 out.questions_abandoned
2894 );
2895 }
2896}
2897
2898async fn triage_held(queue: &Queue, home: &Path, opts: &Opts) {
2907 let questions = Questions::at(home.join("questions"));
2908 let report = triage::run_once(queue, &questions, opts.config.as_deref(), Timestamp::now());
2909 if report.is_empty() {
2910 return;
2911 }
2912 if !report.quarantined.is_empty() {
2913 tracing::info!(
2914 "triage: held {} blocked task(s) whose blocked-on task or \
2915 question no longer exists: {}",
2916 report.quarantined.len(),
2917 report.quarantined.join(", ")
2918 );
2919 }
2920 if !report.resumed.is_empty() {
2921 tracing::info!(
2922 "triage: resumed {} held task(s) whose machine hold had resolved: {}",
2923 report.resumed.len(),
2924 report.resumed.join(", ")
2925 );
2926 }
2927 if !report.asked.is_empty() {
2928 tracing::info!(
2929 "triage: asked about {} held task(s): {}",
2930 report.asked.len(),
2931 report.asked.join(", ")
2932 );
2933 }
2934 if !report.answered.is_empty() {
2935 tracing::info!(
2936 "triage: applied {} operator answer(s): {}",
2937 report.answered.len(),
2938 report.answered.join(", ")
2939 );
2940 }
2941}
2942
2943fn disk_gate(repo: &Path, config: &Config) -> Option<String> {
2950 disk_gate_with(repo, config, crate::disk::free_bytes)
2951}
2952
2953fn disk_gate_with<F: Fn(&Path) -> Result<u64>>(
2957 repo: &Path,
2958 config: &Config,
2959 free_bytes: F,
2960) -> Option<String> {
2961 let min = config.disk.min_free_bytes;
2962 if min == 0 {
2963 return None;
2964 }
2965 match free_bytes(repo) {
2966 Ok(free) => crate::disk::gate(free, min),
2967 Err(e) => Some(format!(
2968 "could not measure free space on {} ({e}); the disk gate refuses \
2969 to let a run start blind",
2970 repo.display()
2971 )),
2972 }
2973}
2974
2975const QUOTA_WAIT_FALLBACK: Duration = Duration::from_secs(5 * 60);
2982
2983const QUOTA_WAIT_CAP: Duration = Duration::from_secs(30 * 60);
2987
2988fn quota_wait(
2997 reset_at: Option<Timestamp>,
2998 now: Timestamp,
2999 fallback: Duration,
3000 cap: Duration,
3001) -> Duration {
3002 match reset_at {
3003 Some(at) if at > now => {
3004 let secs = u64::try_from(at.as_second() - now.as_second()).unwrap_or(0);
3005 Duration::from_secs(secs).min(cap)
3006 }
3007 _ => fallback,
3008 }
3009}
3010
3011fn parse_reset_hint(text: &str, now: Timestamp, recorded: Timestamp) -> Option<Timestamp> {
3025 parse_reset_hint_zoned(text, now)
3026 .or_else(|| parse_reset_hint_dated(text))
3027 .or_else(|| parse_reset_hint_relative(text, recorded))
3028}
3029
3030fn parse_reset_hint_relative(text: &str, recorded: Timestamp) -> Option<Timestamp> {
3034 let rest = text.trim().trim_end_matches('.').strip_prefix("in ")?;
3035 let mut rest = rest.trim();
3036 if rest.is_empty() {
3037 return None;
3038 }
3039 let mut total: i64 = 0;
3040 let mut matched = false;
3041 for (unit, secs) in [('h', 3600), ('m', 60), ('s', 1)] {
3042 if let Some((digits, tail)) = rest.split_once(unit)
3043 && !digits.is_empty()
3044 && digits.bytes().all(|b| b.is_ascii_digit())
3045 {
3046 total += digits.parse::<i64>().ok()?.checked_mul(secs)?;
3047 rest = tail;
3048 matched = true;
3049 }
3050 }
3051 if !rest.is_empty() || !matched {
3052 return None;
3053 }
3054 recorded
3055 .checked_add(jiff::SignedDuration::from_secs(total))
3056 .ok()
3057}
3058
3059fn parse_12h_clock(clock: &str) -> Option<(i8, i8)> {
3063 let clock = clock.trim().to_lowercase();
3064 let (digits, pm) = clock
3065 .strip_suffix("am")
3066 .map(|d| (d, false))
3067 .or_else(|| clock.strip_suffix("pm").map(|d| (d, true)))?;
3068 let (h, m) = digits.trim().split_once(':')?;
3069 let mut hour: i8 = h.trim().parse().ok()?;
3070 let minute: i8 = m.trim().parse().ok()?;
3071 if !(1..=12).contains(&hour) || !(0..=59).contains(&minute) {
3072 return None;
3073 }
3074 if pm && hour != 12 {
3075 hour += 12;
3076 } else if !pm && hour == 12 {
3077 hour = 0;
3078 }
3079 Some((hour, minute))
3080}
3081
3082fn parse_reset_hint_zoned(text: &str, now: Timestamp) -> Option<Timestamp> {
3087 let open = text.find('(')?;
3088 let close = text.rfind(')')?;
3089 if close <= open {
3090 return None;
3091 }
3092 let zone = text[open + 1..close].trim();
3093 let (hour, minute) = parse_12h_clock(&text[..open])?;
3094 let tz = jiff::tz::TimeZone::get(zone).ok()?;
3095 let candidate = now
3096 .to_zoned(tz)
3097 .with()
3098 .hour(hour)
3099 .minute(minute)
3100 .second(0)
3101 .millisecond(0)
3102 .microsecond(0)
3103 .nanosecond(0)
3104 .build()
3105 .ok()?;
3106 let mut at = candidate.timestamp();
3107 if at <= now {
3108 at += jiff::SignedDuration::from_hours(24);
3109 }
3110 Some(at)
3111}
3112
3113fn parse_reset_hint_dated(text: &str) -> Option<Timestamp> {
3122 let words: Vec<&str> = text.split_whitespace().collect();
3123 if words.len() < 5 {
3124 return None;
3125 }
3126 (0..=words.len() - 5)
3127 .find_map(|start| parse_dated_window(&words[start..start + 5], words.get(start + 5)))
3128}
3129
3130fn parse_dated_window(window: &[&str], trailing: Option<&&str>) -> Option<Timestamp> {
3136 if trailing.is_some_and(|next| next.starts_with('(')) {
3137 return None;
3138 }
3139 let month = month_number(window[0])?;
3140 let day_token = window[1].strip_suffix(',')?.to_lowercase();
3141 let day_digits = ["st", "nd", "rd", "th"]
3142 .iter()
3143 .find_map(|suffix| day_token.strip_suffix(*suffix))?;
3144 let day: i8 = day_digits.parse().ok()?;
3145 let year_token = window[2];
3146 if year_token.len() != 4 || !year_token.bytes().all(|b| b.is_ascii_digit()) {
3147 return None;
3148 }
3149 let year: i16 = year_token.parse().ok()?;
3150 let ampm = window[4].trim_matches(|c: char| !c.is_ascii_alphabetic());
3154 let (hour, minute) = parse_12h_clock(&format!("{}{}", window[3], ampm))?;
3155 let date = jiff::civil::Date::new(year, month, day).ok()?;
3156 let candidate = date
3157 .at(hour, minute, 0, 0)
3158 .to_zoned(jiff::tz::TimeZone::UTC)
3159 .ok()?;
3160 Some(candidate.timestamp())
3161}
3162
3163fn month_number(name: &str) -> Option<i8> {
3166 const NAMES: [&str; 12] = [
3167 "jan", "feb", "mar", "apr", "may", "jun", "jul", "aug", "sep", "oct", "nov", "dec",
3168 ];
3169 let lower = name.to_lowercase();
3170 NAMES
3171 .iter()
3172 .position(|n| *n == lower.as_str())
3173 .map(|i| i as i8 + 1)
3174}
3175
3176fn exhausted_review_budget(state: &RunState) -> bool {
3188 state.status == RunStatus::Blocked && state.reviews.len() >= state.config.graph.review_rounds
3189}
3190
3191fn unfinished_run(runs: &[String], short: &str) -> Option<String> {
3228 unfinished_run_with(runs, short, RunState::load)
3229}
3230
3231fn unfinished_run_with<F>(runs: &[String], short: &str, load: F) -> Option<String>
3234where
3235 F: FnOnce(&str) -> Result<RunState>,
3236{
3237 let id = runs.last()?;
3238 match load(id) {
3239 Ok(s) if s.status.resumable() && !s.released() && !exhausted_review_budget(&s) => {
3240 Some(id.clone())
3241 }
3242 Ok(_) => None,
3243 Err(e) => {
3244 tracing::warn!("could not read run {id} for task {short}: {e:#}");
3245 None
3246 }
3247 }
3248}
3249
3250fn awaiting_resume_with<F>(task: &Task, load: F) -> bool
3258where
3259 F: FnOnce(&str) -> Result<RunState>,
3260{
3261 if task.status != TaskStatus::Failed || task.fresh_start || task.review_branch.is_some() {
3262 return false;
3263 }
3264 let Some(id) = task.runs.last() else {
3265 return false;
3266 };
3267 load(id).is_ok_and(|s| {
3268 s.parked && s.status.resumable() && !s.released() && !exhausted_review_budget(&s)
3269 })
3270}
3271
3272#[derive(Debug, Clone, PartialEq, Eq)]
3275enum Starter {
3276 Review(String),
3279 Resume(String),
3281 Start,
3283}
3284
3285fn take_divergence_answer(
3289 branch: &str,
3290 remote: &str,
3291 task: &mut Task,
3292) -> Option<crate::reconcile::Choice> {
3293 let summary = crate::reconcile::summary_for(branch, remote);
3294 let (idx, choice) = task.answers.iter().enumerate().rev().find_map(|(i, a)| {
3295 (a.question == summary)
3296 .then(|| crate::reconcile::Choice::from_answer(&a.answer))
3297 .flatten()
3298 .map(|c| (i, c))
3299 })?;
3300 task.answers.remove(idx);
3301 Some(choice)
3302}
3303
3304fn choose_starter(
3316 review_branch: Option<&str>,
3317 branch_exists: bool,
3318 unfinished: Option<&str>,
3319) -> Starter {
3320 match review_branch {
3321 Some(branch) if branch_exists => Starter::Review(branch.to_owned()),
3322 Some(_) => Starter::Start,
3323 None => match unfinished {
3324 Some(id) => Starter::Resume(id.to_owned()),
3325 None => Starter::Start,
3326 },
3327 }
3328}
3329
3330fn repo_for(task: &Task, fallback: &Path) -> PathBuf {
3333 if task.repo.as_os_str().is_empty() || task.repo == Path::new(".") {
3334 return fallback.to_path_buf();
3335 }
3336 task.repo.clone()
3337}
3338
3339const ANSWERS_HEADER: &str = "\n\n# Operator answers\n\n";
3343
3344fn answers_block(task: &Task, count: usize) -> String {
3346 let mut s = ANSWERS_HEADER.to_owned();
3347 for a in &task.answers[..count] {
3348 s.push_str(&format!("- {}: {}\n", a.question, a.answer));
3349 }
3350 s
3351}
3352
3353fn append_answers(base: &str, task: &Task) -> String {
3356 if task.answers.is_empty() {
3357 return base.to_owned();
3358 }
3359 let mut s = base.to_owned();
3360 s.push_str(&answers_block(task, task.answers.len()));
3361 s
3362}
3363
3364fn strip_answers_block<'a>(instruction: &'a str, task: &Task) -> &'a str {
3368 for count in (1..=task.answers.len()).rev() {
3369 let block = answers_block(task, count);
3370 if let Some(base) = instruction.strip_suffix(&block) {
3371 return base;
3372 }
3373 }
3374 instruction
3375}
3376
3377fn instruction_for(task: &Task) -> String {
3385 append_answers(&task.instruction, task)
3386}
3387
3388fn resumed_instruction(old_instruction: &str, task: &Task) -> String {
3400 append_answers(strip_answers_block(old_instruction, task), task)
3401}
3402
3403fn task_attachments(queue: &Queue, task: &Task) -> Result<Vec<PathBuf>> {
3405 let paths = queue.attachment_paths(task);
3406 for (name, path) in task.attachments.iter().zip(&paths) {
3407 if !path.is_file() {
3408 bail!(
3409 "attachment `{name}` is recorded on the task but {} is missing",
3410 path.display()
3411 );
3412 }
3413 }
3414 Ok(paths)
3415}
3416
3417fn prepare_instruction(
3428 starter: &Starter,
3429 old_instruction: Option<&str>,
3430 task: &Task,
3431) -> Option<String> {
3432 match starter {
3433 Starter::Start => Some(instruction_for(task)),
3434 Starter::Resume(_) => Some(resumed_instruction(
3435 old_instruction.expect("a resumed run always has a prior instruction"),
3436 task,
3437 )),
3438 Starter::Review(_) => None,
3439 }
3440}
3441
3442fn record(queue: &Queue, task: &mut Task) {
3446 if let Err(e) = queue.put(task) {
3447 tracing::error!("could not record task {}: {e:#}", task.short());
3448 notices::raise(Notice::error(
3449 "loop:record",
3450 "The loop could not save a task's state; check the disk.",
3451 ));
3452 }
3453}
3454
3455fn runnable(queue: &Queue) -> Vec<Task> {
3461 let mut tasks: Vec<Task> = queue
3462 .list()
3463 .into_iter()
3464 .filter(|t| t.status.runnable())
3465 .collect();
3466 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
3467 tasks
3468}
3469
3470fn describe(state: &RunState) -> String {
3484 let mut detail = if state.status == RunStatus::Stalled {
3485 let mut seats: Vec<&str> = state.quota.iter().map(|q| q.seat.as_str()).collect();
3486 seats.sort_unstable();
3487 seats.dedup();
3488 if seats.is_empty() {
3489 "the judging panel lost its quorum".to_owned()
3490 } else {
3491 format!(
3492 "the judging panel lost its quorum; quota took out {}",
3493 seats.join(", ")
3494 )
3495 }
3496 } else {
3497 format!("run ended {}", state.status.display_label())
3498 };
3499 if let Some(last) = state.events.last() {
3500 detail.push_str(&format!(" ({}: {})", last.node, last.message));
3501 }
3502 detail.push_str(&format!(" [run {}]", state.id));
3503 detail
3504}
3505
3506const DIAGNOSTIC_MAX: usize = 4_000;
3512
3513const DIAGNOSTIC_OUTPUT_TAIL: usize = 800;
3518
3519fn diagnostic(state: &RunState) -> Option<String> {
3533 let mut parts: Vec<String> = Vec::new();
3534
3535 for o in state.gate.iter().filter(|o| !o.ok()) {
3537 parts.push(format!(
3538 "gate `{}` failed ({:?}):\n{}",
3539 o.command,
3540 o.code,
3541 crate::run::tail(&o.output_tail, DIAGNOSTIC_OUTPUT_TAIL)
3542 ));
3543 }
3544
3545 if let Some(last) = state
3548 .events
3549 .iter()
3550 .rev()
3551 .find(|e| e.node == "land" && e.message.contains("fixer produced no commit"))
3552 {
3553 parts.push(last.message.clone());
3554 }
3555
3556 if state.viable().is_empty() {
3563 for c in &state.candidates {
3564 if let Some(evidence) = &c.verified_noop {
3565 parts.push(format!(
3566 "candidate {} (agent-verified no-op, unconfirmed by magi): {evidence}",
3567 c.label
3568 ));
3569 } else if !c.summary.trim().is_empty() {
3570 parts.push(format!("candidate {}: {}", c.label, c.summary.trim()));
3571 } else if let Some(why) = &c.failed {
3572 parts.push(format!("candidate {}: {why}", c.label));
3573 }
3574 }
3575 }
3576
3577 if parts.is_empty() {
3578 return None;
3579 }
3580 Some(crate::run::tail(
3585 &parts.join("\n\n"),
3586 DIAGNOSTIC_MAX.saturating_sub(100),
3587 ))
3588}
3589
3590fn label(status: RunStatus) -> &'static str {
3598 status.as_str()
3599}
3600
3601fn merge_mode(mode: &str) -> Result<MergeMode> {
3603 match mode {
3604 "none" => Ok(MergeMode::None),
3605 "local" => Ok(MergeMode::Local),
3606 "pr" => Ok(MergeMode::Pr),
3607 other => bail!("unknown merge mode `{other}`; expected none, local or pr"),
3608 }
3609}
3610
3611fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
3616 mutex
3617 .lock()
3618 .unwrap_or_else(std::sync::PoisonError::into_inner)
3619}
3620
3621#[cfg(test)]
3622mod tests {
3623 use super::*;
3624 use crate::queue::{Source, TaskStatus};
3625 use crate::run::{Candidate, CommandOutcome};
3626 use pretty_assertions::assert_eq;
3627
3628 fn task() -> Task {
3629 Task::new(
3630 "add retries".to_owned(),
3631 "add retries".to_owned(),
3632 PathBuf::from("/repo"),
3633 Source::Human,
3634 )
3635 }
3636
3637 fn interrupt_task(id: &str) -> Task {
3640 let mut t = task();
3641 t.id = id.to_owned();
3642 t.interrupt = true;
3643 t
3644 }
3645
3646 fn task_with_id(id: &str) -> Task {
3648 let mut t = task();
3649 t.id = id.to_owned();
3650 t
3651 }
3652
3653 fn urgent_task(id: &str) -> Task {
3655 let mut t = task();
3656 t.id = id.to_owned();
3657 t.urgent = true;
3658 t
3659 }
3660
3661 #[test]
3666 fn permit_kind_prefers_a_land_resume_over_the_urgent_slot() {
3667 assert_eq!(permit_kind(true, false), PermitKind::None);
3668 assert_eq!(permit_kind(true, true), PermitKind::None);
3669 }
3670
3671 #[test]
3676 fn permit_kind_separates_urgent_from_ordinary() {
3677 assert_eq!(permit_kind(false, true), PermitKind::Urgent);
3678 assert_eq!(permit_kind(false, false), PermitKind::Ordinary);
3679 }
3680
3681 #[test]
3688 fn disk_gate_with_holds_a_task_below_the_threshold_and_names_both_numbers() {
3689 let cfg = Config::default();
3690 let repo = Path::new("/any/repo/path");
3691
3692 let reason =
3693 disk_gate_with(repo, &cfg, |_| Ok(1024)).expect("must hold below the threshold");
3694 assert!(reason.contains("1024"), "{reason}");
3695 assert!(
3696 reason.contains(&cfg.disk.min_free_bytes.to_string()),
3697 "{reason}"
3698 );
3699
3700 assert_eq!(
3701 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes)),
3702 None,
3703 "exactly at the floor is open"
3704 );
3705 assert_eq!(
3706 disk_gate_with(repo, &cfg, |_| Ok(cfg.disk.min_free_bytes + 1)),
3707 None,
3708 "comfortably above the floor is open"
3709 );
3710 }
3711
3712 #[test]
3713 fn disk_gate_with_opens_unconditionally_when_the_operator_opted_out() {
3714 let mut cfg = Config::default();
3715 cfg.disk.min_free_bytes = 0;
3716 let repo = Path::new("/any/repo/path");
3717 assert_eq!(
3718 disk_gate_with(repo, &cfg, |_| Ok(0)),
3719 None,
3720 "a zero floor never measures at all"
3721 );
3722 }
3723
3724 #[test]
3725 fn disk_gate_with_closes_rather_than_starts_blind_when_it_cannot_measure() {
3726 let cfg = Config::default();
3727 let repo = Path::new("/any/repo/path");
3728 let reason = disk_gate_with(repo, &cfg, |_| Err(anyhow::anyhow!("no df on this box")))
3729 .expect("a measurement failure must close the gate, not open it");
3730 assert!(reason.contains("could not measure"), "{reason}");
3731 }
3732
3733 #[test]
3734 fn no_interrupt_task_leaves_the_sequence_idle_even_with_something_in_flight() {
3735 let ordinary = task();
3736 let next = advance_interrupt(
3737 Interrupt::Idle,
3738 std::slice::from_ref(&ordinary.id),
3739 std::slice::from_ref(&ordinary),
3740 );
3741 assert_eq!(next, Interrupt::Idle);
3742 }
3743
3744 #[test]
3745 fn an_interrupt_task_with_nothing_in_flight_never_starts_a_sequence() {
3746 let marked = interrupt_task("marked");
3749 let next = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
3750 assert_eq!(next, Interrupt::Idle);
3751 }
3752
3753 #[test]
3754 fn an_interrupt_task_with_something_in_flight_starts_parking_it() {
3755 let marked = interrupt_task("marked");
3756 let next = advance_interrupt(
3757 Interrupt::Idle,
3758 &["running".to_owned()],
3759 std::slice::from_ref(&marked),
3760 );
3761 assert_eq!(
3762 next,
3763 Interrupt::Parking {
3764 parked: vec!["running".to_owned()],
3765 interrupt_task: "marked".to_owned(),
3766 }
3767 );
3768 }
3769
3770 #[test]
3779 fn more_than_one_run_in_flight_never_starts_an_interrupt_sequence() {
3780 let marked = interrupt_task("marked");
3781
3782 let two = advance_interrupt(
3783 Interrupt::Idle,
3784 &["a".to_owned(), "b".to_owned()],
3785 std::slice::from_ref(&marked),
3786 );
3787 assert_eq!(two, Interrupt::Idle);
3788
3789 let none = advance_interrupt(Interrupt::Idle, &[], std::slice::from_ref(&marked));
3790 assert_eq!(none, Interrupt::Idle, "nothing to interrupt either");
3791 }
3792
3793 #[test]
3794 fn parking_holds_until_every_parked_id_has_actually_left_flight() {
3795 let state = Interrupt::Parking {
3796 parked: vec!["running".to_owned()],
3797 interrupt_task: "marked".to_owned(),
3798 };
3799 let still_going = advance_interrupt(state.clone(), &["running".to_owned()], &[]);
3801 assert_eq!(still_going, state);
3802
3803 let stopped_but_not_yet_dispatched =
3807 advance_interrupt(state.clone(), &[], &[interrupt_task("marked")]);
3808 assert_eq!(stopped_but_not_yet_dispatched, state);
3809
3810 let dispatched = advance_interrupt(state, &["marked".to_owned()], &[]);
3812 assert_eq!(
3813 dispatched,
3814 Interrupt::Running {
3815 parked: vec!["running".to_owned()],
3816 interrupt_task: "marked".to_owned(),
3817 }
3818 );
3819 }
3820
3821 #[test]
3822 fn the_sequence_moves_to_resuming_the_instant_the_interrupt_tasks_own_run_leaves_flight() {
3823 let state = Interrupt::Running {
3824 parked: vec!["running".to_owned()],
3825 interrupt_task: "marked".to_owned(),
3826 };
3827 let still_running = advance_interrupt(state.clone(), &["marked".to_owned()], &[]);
3828 assert_eq!(still_running, state);
3829
3830 let ended = advance_interrupt(state, &[], &[task_with_id("running")]);
3837 assert_eq!(
3838 ended,
3839 Interrupt::Resuming {
3840 parked: vec!["running".to_owned()]
3841 }
3842 );
3843 }
3844
3845 #[test]
3846 fn resuming_ends_the_instant_a_parked_task_is_seen_in_flight() {
3847 let state = Interrupt::Resuming {
3848 parked: vec!["running".to_owned()],
3849 };
3850 let still_waiting = advance_interrupt(state.clone(), &[], &[task_with_id("running")]);
3851 assert_eq!(still_waiting, state);
3852
3853 let dispatched = advance_interrupt(state, &["running".to_owned()], &[]);
3854 assert_eq!(dispatched, Interrupt::Idle);
3855 }
3856
3857 #[test]
3863 fn an_interrupt_task_that_stops_being_runnable_abandons_the_wait_without_losing_the_parked_run()
3864 {
3865 let state = Interrupt::Parking {
3866 parked: vec!["running".to_owned()],
3867 interrupt_task: "marked".to_owned(),
3868 };
3869 let next = advance_interrupt(state, &[], &[]);
3872 assert_eq!(
3873 next,
3874 Interrupt::Resuming {
3875 parked: vec!["running".to_owned()]
3876 },
3877 "abandoning the interrupt must not abandon the resume it owes"
3878 );
3879 }
3880
3881 #[test]
3884 fn resuming_abandons_a_parked_task_that_stops_being_runnable() {
3885 let state = Interrupt::Resuming {
3886 parked: vec!["running".to_owned()],
3887 };
3888 let next = advance_interrupt(state, &[], &[]);
3889 assert_eq!(
3890 next,
3891 Interrupt::Idle,
3892 "nothing is left to wait for; the loop must not stay wedged"
3893 );
3894 }
3895
3896 #[test]
3897 fn disabled_by_config_the_sequence_can_never_leave_idle() {
3898 let marked = interrupt_task("marked");
3899 let next = advance_interrupt_tick(
3900 false,
3901 Interrupt::Idle,
3902 &["running".to_owned()],
3903 std::slice::from_ref(&marked),
3904 );
3905 assert_eq!(
3906 next,
3907 Interrupt::Idle,
3908 "an unmarked, unconfigured daemon must behave exactly as before"
3909 );
3910 }
3911
3912 #[test]
3913 fn the_gate_blocks_everyone_while_something_parked_is_still_in_flight() {
3914 let state = Interrupt::Parking {
3915 parked: vec!["running".to_owned()],
3916 interrupt_task: "marked".to_owned(),
3917 };
3918 let candidates = vec![interrupt_task("marked"), task()];
3919 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
3920 assert!(
3921 allowed.is_empty(),
3922 "nothing may dispatch - not even the interrupt task itself - \
3923 until the parked run has actually stopped"
3924 );
3925 }
3926
3927 #[test]
3941 fn urgent_gains_no_exemption_from_an_active_interrupt_sequence() {
3942 for state in [
3943 Interrupt::Parking {
3944 parked: vec!["running".to_owned()],
3945 interrupt_task: "marked".to_owned(),
3946 },
3947 Interrupt::Running {
3948 parked: vec!["running".to_owned()],
3949 interrupt_task: "marked".to_owned(),
3950 },
3951 Interrupt::Resuming {
3952 parked: vec!["running".to_owned()],
3953 },
3954 ] {
3955 let candidates = vec![
3956 interrupt_task("marked"),
3957 urgent_task("hot"),
3958 task_with_id("ordinary"),
3959 ];
3960 let allowed = interrupt_gate(&state, &["running".to_owned()], candidates);
3961 assert!(
3962 !allowed.iter().any(|t| t.id == "hot"),
3963 "an urgent candidate must wait out the same gate as anything \
3964 else while the run it would run alongside has not actually \
3965 left flight, for state {state:?}: {allowed:?}"
3966 );
3967 }
3968 }
3969
3970 #[test]
3976 fn a_task_marked_both_urgent_and_interrupt_is_admitted_once_the_gate_itself_says_so() {
3977 let state = Interrupt::Resuming {
3978 parked: vec!["hot".to_owned()],
3979 };
3980 let candidates = vec![urgent_task("hot"), task()];
3981 let allowed = interrupt_gate(&state, &[], candidates);
3982 assert_eq!(
3983 allowed.iter().filter(|t| t.id == "hot").count(),
3984 1,
3985 "the gate's own decision is unaffected by the urgent flag: {allowed:?}"
3986 );
3987 }
3988
3989 #[test]
3990 fn the_gate_lets_only_the_interrupt_task_through_once_parked_work_has_stopped() {
3991 let state = Interrupt::Parking {
3992 parked: vec!["running".to_owned()],
3993 interrupt_task: "marked".to_owned(),
3994 };
3995 let other = task();
3996 let candidates = vec![interrupt_task("marked"), other.clone()];
3997 let allowed = interrupt_gate(&state, &[], candidates);
3998 assert_eq!(allowed.len(), 1);
3999 assert_eq!(allowed[0].id, "marked");
4000 }
4001
4002 #[test]
4003 fn the_gate_blocks_everyone_while_the_interrupt_task_itself_is_in_flight() {
4004 let state = Interrupt::Running {
4005 parked: vec!["running".to_owned()],
4006 interrupt_task: "marked".to_owned(),
4007 };
4008 let candidates = vec![task(), task()];
4009 let allowed = interrupt_gate(&state, &["marked".to_owned()], candidates);
4010 assert!(allowed.is_empty());
4011 }
4012
4013 #[test]
4020 fn the_gate_offers_at_most_one_candidate_while_resuming_even_with_two_parked() {
4021 let state = Interrupt::Resuming {
4022 parked: vec!["a".to_owned(), "c".to_owned()],
4023 };
4024 let candidates = vec![task_with_id("a"), task_with_id("c"), task_with_id("other")];
4025 let allowed = interrupt_gate(&state, &[], candidates);
4026 assert_eq!(
4027 allowed.len(),
4028 1,
4029 "at most one candidate may be offered while resuming: {allowed:?}"
4030 );
4031 assert_eq!(allowed[0].id, "a");
4032 }
4033
4034 #[test]
4035 fn the_gate_offers_nothing_while_resuming_if_no_parked_task_is_runnable() {
4036 let state = Interrupt::Resuming {
4037 parked: vec!["a".to_owned()],
4038 };
4039 let allowed = interrupt_gate(&state, &[], vec![task_with_id("other")]);
4040 assert!(allowed.is_empty());
4041 }
4042
4043 #[test]
4049 fn a_full_sequence_never_gates_two_runs_through_at_once_and_resumes_exactly_one() {
4050 let running = task(); let marked = interrupt_task("marked");
4052
4053 let mut state = Interrupt::Idle;
4054 let in_flight = vec![running.id.clone()];
4056 state = advance_interrupt_tick(true, state, &in_flight, std::slice::from_ref(&marked));
4057 let gated = interrupt_gate(&state, &in_flight, vec![marked.clone(), running.clone()]);
4058 assert!(gated.is_empty(), "still waiting on `running` to park");
4059
4060 state = advance_interrupt_tick(true, state, &[], &[marked.clone(), running.clone()]);
4062 let gated = interrupt_gate(&state, &[], vec![marked.clone(), running.clone()]);
4063 assert_eq!(
4064 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4065 vec!["marked"],
4066 "only the interrupt task may be offered to the dispatcher now"
4067 );
4068
4069 state = advance_interrupt_tick(
4071 true,
4072 state,
4073 &["marked".to_owned()],
4074 std::slice::from_ref(&running),
4075 );
4076 let gated = interrupt_gate(
4077 &state,
4078 &["marked".to_owned()],
4079 vec![marked.clone(), running.clone()],
4080 );
4081 assert!(
4082 gated.is_empty(),
4083 "the parked run must not be offered back while the interrupt \
4084 task is still running"
4085 );
4086
4087 let other = task_with_id("other");
4091 state = advance_interrupt_tick(true, state, &[], &[running.clone(), other.clone()]);
4092 assert_eq!(
4093 state,
4094 Interrupt::Resuming {
4095 parked: vec![running.id.clone()]
4096 }
4097 );
4098 let gated = interrupt_gate(&state, &[], vec![other.clone(), running.clone()]);
4099 assert_eq!(
4100 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4101 vec![running.id.as_str()],
4102 "exactly the parked run resumes - not the unrelated task, even \
4103 though it was offered first"
4104 );
4105
4106 state = advance_interrupt_tick(
4110 true,
4111 state,
4112 std::slice::from_ref(&running.id),
4113 std::slice::from_ref(&other),
4114 );
4115 assert_eq!(state, Interrupt::Idle);
4116 let gated = interrupt_gate(
4117 &state,
4118 std::slice::from_ref(&running.id),
4119 vec![other.clone()],
4120 );
4121 assert_eq!(
4122 gated.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
4123 vec![other.id.as_str()],
4124 "ordinary dispatch is unrestricted again"
4125 );
4126 }
4127
4128 #[test]
4129 fn every_run_status_settles_the_task_it_came_from() {
4130 let table = [
4132 (RunStatus::Merged, TaskStatus::Done, 1),
4133 (RunStatus::Ready, TaskStatus::Done, 1),
4134 (RunStatus::Stalled, TaskStatus::Failed, 0),
4135 (RunStatus::Blocked, TaskStatus::Failed, 1),
4136 (RunStatus::Failed, TaskStatus::Failed, 1),
4137 (RunStatus::VerifiedNoop, TaskStatus::Held, 1),
4138 (RunStatus::Prep, TaskStatus::Failed, 1),
4139 (RunStatus::Implementing, TaskStatus::Failed, 1),
4140 (RunStatus::Judging, TaskStatus::Failed, 1),
4141 (RunStatus::Deliberating, TaskStatus::Failed, 1),
4142 (RunStatus::Voting, TaskStatus::Failed, 1),
4143 (RunStatus::Reviewing, TaskStatus::Failed, 1),
4144 (RunStatus::Gating, TaskStatus::Failed, 1),
4145 ];
4146 for (run, want, attempts) in table {
4147 let mut t = task();
4148 t.start("20260902-000000-aaaa".to_owned());
4149 settle(
4150 &mut t,
4151 Verdict {
4152 status: run,
4153 left_pr: false,
4154 parked: false,
4155 quota_hit: matches!(run, RunStatus::Stalled),
4156 no_viable_candidates: false,
4157 },
4158 "why",
4159 2,
4160 );
4161 assert_eq!(t.status, want, "task status after {}", label(run));
4162 assert_eq!(t.attempts, attempts, "attempts after {}", label(run));
4163 }
4164 }
4165
4166 #[test]
4167 fn a_quota_stall_costs_the_task_no_attempt_but_a_block_does() {
4168 let mut stalled = task();
4169 stalled.start("20260902-000000-aaaa".to_owned());
4170 settle(
4171 &mut stalled,
4172 Verdict {
4173 status: RunStatus::Stalled,
4174 left_pr: false,
4175 parked: false,
4176 quota_hit: true,
4177 no_viable_candidates: false,
4178 },
4179 "quota",
4180 1,
4181 );
4182 assert_eq!(stalled.attempts, 0);
4183 assert!(
4184 stalled.status.runnable(),
4185 "a machine problem must leave the task in line"
4186 );
4187
4188 let mut blocked = task();
4189 blocked.start("20260902-000000-aaaa".to_owned());
4190 settle(
4191 &mut blocked,
4192 Verdict {
4193 status: RunStatus::Blocked,
4194 left_pr: false,
4195 parked: false,
4196 quota_hit: false,
4197 no_viable_candidates: false,
4198 },
4199 "findings open",
4200 1,
4201 );
4202 assert_eq!(blocked.attempts, 1);
4203 assert_eq!(
4204 blocked.status,
4205 TaskStatus::Held,
4206 "the last attempt hands the task to a human"
4207 );
4208 }
4209
4210 #[test]
4211 fn a_run_that_opened_a_pull_request_is_never_re_competed() {
4212 let mut delivered = task();
4215 delivered.start("20260903-080619-01c2".to_owned());
4216 settle(
4217 &mut delivered,
4218 Verdict {
4219 status: RunStatus::Blocked,
4220 left_pr: true,
4221 parked: false,
4222 quota_hit: false,
4223 no_viable_candidates: false,
4224 },
4225 "no check status",
4226 4,
4227 );
4228 assert_eq!(
4229 delivered.status,
4230 TaskStatus::Held,
4231 "a pull request waiting on CI or a person is not a retryable failure"
4232 );
4233 assert!(
4234 !delivered.status.runnable(),
4235 "the loop must not pick this task up again"
4236 );
4237 assert_eq!(
4238 delivered.last_error.as_deref(),
4239 Some("no check status"),
4240 "the operator needs to be told what the gate was waiting for"
4241 );
4242
4243 let mut empty_handed = task();
4246 empty_handed.start("20260903-080619-01c2".to_owned());
4247 settle(
4248 &mut empty_handed,
4249 Verdict {
4250 status: RunStatus::Blocked,
4251 left_pr: false,
4252 parked: false,
4253 quota_hit: false,
4254 no_viable_candidates: false,
4255 },
4256 "findings open",
4257 4,
4258 );
4259 assert_eq!(empty_handed.status, TaskStatus::Failed);
4260 assert!(empty_handed.status.runnable());
4261 }
4262
4263 #[test]
4264 fn a_verified_noop_run_hands_off_rather_than_closing_or_auto_retrying() {
4265 let mut noop = task();
4271 noop.start("20260912-131304-391f".to_owned());
4272 settle(
4273 &mut noop,
4274 Verdict {
4275 status: RunStatus::VerifiedNoop,
4276 left_pr: false,
4277 parked: false,
4278 quota_hit: false,
4279 no_viable_candidates: true,
4280 },
4281 "candidate A: already fixed by b32cfc4, on main",
4282 4,
4283 );
4284 assert_eq!(
4285 noop.status,
4286 TaskStatus::Held,
4287 "an unverified claim is a request for a human, not a failure"
4288 );
4289 assert!(
4290 !noop.status.runnable(),
4291 "the loop must not requeue this on the same unverified claim"
4292 );
4293 assert_eq!(noop.attempts, 1);
4298 }
4299
4300 #[test]
4301 fn parking_costs_the_task_no_attempt_and_leaves_it_in_line() {
4302 let mut parked = task();
4307 parked.start("20260903-183634-2d98".to_owned());
4308 settle(
4309 &mut parked,
4310 Verdict {
4311 status: RunStatus::Implementing,
4312 left_pr: false,
4313 quota_hit: false,
4314 parked: true,
4315 no_viable_candidates: false,
4316 },
4317 "parked after `implementing`",
4318 2,
4319 );
4320 assert_eq!(parked.attempts, 0, "a park is refunded");
4321 assert!(
4322 parked.status.runnable(),
4323 "and the task stays in line so the next loop resumes its run"
4324 );
4325 assert_eq!(
4326 parked.last_error.as_deref(),
4327 Some("parked after `implementing`"),
4328 "the card says where it stopped"
4329 );
4330
4331 let mut broken = task();
4335 broken.start("20260903-183634-2d98".to_owned());
4336 settle(
4337 &mut broken,
4338 Verdict {
4339 status: RunStatus::Implementing,
4340 left_pr: false,
4341 quota_hit: false,
4342 parked: false,
4343 no_viable_candidates: false,
4344 },
4345 "returned mid-flight",
4346 2,
4347 );
4348 assert_eq!(broken.attempts, 1);
4349 }
4350
4351 #[test]
4352 fn only_a_rate_limit_buys_the_task_its_attempt_back() {
4353 let mut flaky = task();
4358 flaky.start("20260903-123023-e633".to_owned());
4359 settle(
4360 &mut flaky,
4361 Verdict {
4362 status: RunStatus::Stalled,
4363 left_pr: false,
4364 parked: false,
4365 quota_hit: false,
4366 no_viable_candidates: false,
4367 },
4368 "verdict rests on 1 of 3 judges",
4369 2,
4370 );
4371 assert_eq!(
4372 flaky.attempts, 1,
4373 "flakiness spends an attempt, so `max_attempts` still bounds it"
4374 );
4375 assert!(flaky.status.runnable(), "and it is still worth retrying");
4376
4377 let mut limited = task();
4379 limited.start("20260903-123023-e633".to_owned());
4380 settle(
4381 &mut limited,
4382 Verdict {
4383 status: RunStatus::Stalled,
4384 left_pr: false,
4385 parked: false,
4386 quota_hit: true,
4387 no_viable_candidates: false,
4388 },
4389 "judge-2, judge-3 out of quota",
4390 2,
4391 );
4392 assert_eq!(limited.attempts, 0, "a quota window is refunded");
4393 assert!(limited.status.runnable());
4394
4395 let mut worn = task();
4398 for _ in 0..2 {
4399 worn.release();
4400 }
4401 worn.start("20260903-123023-e633".to_owned());
4402 worn.attempts = 2;
4403 settle(
4404 &mut worn,
4405 Verdict {
4406 status: RunStatus::Stalled,
4407 left_pr: false,
4408 parked: false,
4409 quota_hit: false,
4410 no_viable_candidates: false,
4411 },
4412 "no quorum again",
4413 2,
4414 );
4415 assert_eq!(worn.status, TaskStatus::Held);
4416 assert!(!worn.status.runnable());
4417 }
4418
4419 #[test]
4420 fn a_quota_wipeout_that_leaves_nothing_to_judge_also_costs_no_attempt() {
4421 let mut wiped_out = task();
4428 wiped_out.start("20260907-025000-a1b2".to_owned());
4429 settle(
4430 &mut wiped_out,
4431 Verdict {
4432 status: RunStatus::Failed,
4433 left_pr: false,
4434 parked: false,
4435 quota_hit: true,
4436 no_viable_candidates: true,
4437 },
4438 "no candidate produced a change; nothing to judge",
4439 2,
4440 );
4441 assert_eq!(wiped_out.attempts, 0, "a total quota wipeout is refunded");
4442 assert!(
4443 wiped_out.status.runnable(),
4444 "a machine problem must leave the task in line"
4445 );
4446
4447 let mut partial_progress = task();
4453 partial_progress.start("20260907-025500-c3d4".to_owned());
4454 settle(
4455 &mut partial_progress,
4456 Verdict {
4457 status: RunStatus::Failed,
4458 left_pr: false,
4459 parked: false,
4460 quota_hit: true,
4461 no_viable_candidates: false,
4462 },
4463 "gate failed on the winning candidate",
4464 2,
4465 );
4466 assert_eq!(
4467 partial_progress.attempts, 1,
4468 "a candidate that actually produced a change spends the attempt \
4469 even though some other seat hit its quota"
4470 );
4471 assert!(partial_progress.status.runnable());
4472 }
4473
4474 #[test]
4475 fn reclaim_refunds_a_recovered_quota_wipeout_the_same_way_a_live_settle_does() {
4476 let mut t = task();
4482 t.start("20260907-025000-a1b2".to_owned());
4483 let mut state = run_state(RunStatus::Failed);
4484 state.quota.push(QuotaLoss {
4485 seat: "cand-a".to_owned(),
4486 node: "implement".to_owned(),
4487 at: Timestamp::now(),
4488 reset: None,
4489 });
4490 assert!(
4491 state.viable().is_empty(),
4492 "no candidate was added, so nothing is viable"
4493 );
4494 reclaim(&mut t, Some(state), 2);
4495 assert_eq!(t.attempts, 0, "a recovered quota wipeout is refunded");
4496 assert!(t.status.runnable());
4497 }
4498
4499 #[test]
4500 fn a_held_task_is_never_offered_to_the_loop() {
4501 let dir = tempfile::tempdir().unwrap();
4502 let queue = Queue::at(dir.path().to_path_buf());
4503 for (n, priority) in [(1, 0), (2, 5), (3, 5)] {
4504 let mut t = task();
4505 t.id = format!("2026090{n}-000000-000{n}");
4506 t.priority = priority;
4507 queue.put(&mut t).unwrap();
4508 }
4509 let mut held = task();
4510 held.id = "20260909-000000-9999".to_owned();
4511 held.priority = 99;
4512 held.hold_machine(None);
4513 queue.put(&mut held).unwrap();
4514
4515 let order: Vec<String> = runnable(&queue).into_iter().map(|t| t.id).collect();
4516 assert_eq!(order.len(), 3);
4517 assert!(!order.contains(&held.id));
4518 assert_eq!(
4519 order.first().cloned(),
4520 queue.next_runnable().map(|t| t.id),
4521 "the loop's first candidate is exactly what the queue offers"
4522 );
4523 assert_eq!(
4524 order,
4525 vec![
4526 "20260902-000000-0002".to_owned(),
4527 "20260903-000000-0003".to_owned(),
4528 "20260901-000000-0001".to_owned(),
4529 ],
4530 "priority first, then oldest, so nothing starves"
4531 );
4532 }
4533
4534 #[test]
4535 fn sweep_removes_an_old_unparseable_lock_and_keeps_a_live_one() {
4536 let dir = tempfile::tempdir().unwrap();
4537 let queue = Queue::at(dir.path().to_path_buf());
4538 let mut old = task();
4539 old.id = "20260101-000000-old0".to_owned();
4540 queue.put(&mut old).unwrap();
4541 let mut fresh = task();
4542 fresh.id = "20260101-000000-new0".to_owned();
4543 queue.put(&mut fresh).unwrap();
4544
4545 std::fs::write(dir.path().join(format!("{}.lock", old.id)), "not a pid").unwrap();
4549 std::thread::sleep(Duration::from_millis(60));
4550 let live = queue.claim(&fresh.id).unwrap();
4551
4552 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4553 assert_eq!(swept, vec![old.id.clone()]);
4554 assert!(
4555 queue.claim(&old.id).is_ok(),
4556 "an unparseable lock older than the threshold is swept"
4557 );
4558 assert!(
4559 queue.claim(&fresh.id).is_err(),
4560 "a live pid protects its lock regardless of age"
4561 );
4562 drop(live);
4563 }
4564
4565 #[test]
4566 fn an_old_lock_whose_pid_is_still_alive_is_never_swept_by_age_alone() {
4567 let dir = tempfile::tempdir().unwrap();
4577 let queue = Queue::at(dir.path().to_path_buf());
4578 let mut t = task();
4579 t.id = "20260101-000000-live".to_owned();
4580 queue.put(&mut t).unwrap();
4581
4582 let claim = queue.claim(&t.id).unwrap();
4583 std::thread::sleep(Duration::from_millis(60));
4584
4585 let swept = sweep_stale_claims(&queue, Duration::from_millis(50));
4586 assert!(
4587 swept.is_empty(),
4588 "a lock naming a live pid must never be swept by age, no matter how old: {swept:?}"
4589 );
4590 assert!(
4591 queue.claim(&t.id).is_err(),
4592 "the lock still protects its task"
4593 );
4594 drop(claim);
4595 }
4596
4597 fn injected_dead_pid() -> u32 {
4600 std::process::id().checked_add(1).unwrap_or(1)
4601 }
4602
4603 #[test]
4604 fn a_lock_naming_a_dead_pid_is_swept_at_once_regardless_of_age() {
4605 let dir = tempfile::tempdir().unwrap();
4606 let queue = Queue::at(dir.path().to_path_buf());
4607 let mut t = task();
4608 t.id = "20260101-000000-dead".to_owned();
4609 queue.put(&mut t).unwrap();
4610 let dead_pid = injected_dead_pid();
4611
4612 std::fs::write(
4617 dir.path().join(format!("{}.lock", t.id)),
4618 dead_pid.to_string(),
4619 )
4620 .unwrap();
4621
4622 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4623 pid != dead_pid
4624 });
4625 assert_eq!(
4626 swept,
4627 vec![t.id.clone()],
4628 "a dead owner is reclaimed immediately, not after STALE_CLAIM"
4629 );
4630 assert!(queue.claim(&t.id).is_ok(), "the task is claimable again");
4631 }
4632
4633 #[test]
4634 fn sweeping_on_every_poll_catches_a_lock_that_appears_after_the_first_sweep() {
4635 let dir = tempfile::tempdir().unwrap();
4636 let queue = Queue::at(dir.path().to_path_buf());
4637 let mut t = task();
4638 t.id = "20260101-000000-late".to_owned();
4639 queue.put(&mut t).unwrap();
4640 let dead_pid = injected_dead_pid();
4641
4642 assert!(
4645 sweep_stale_claims(&queue, Duration::from_secs(6 * 60 * 60)).is_empty(),
4646 "nothing has claimed the task yet"
4647 );
4648
4649 std::fs::write(
4652 dir.path().join(format!("{}.lock", t.id)),
4653 dead_pid.to_string(),
4654 )
4655 .unwrap();
4656
4657 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4661 pid != dead_pid
4662 });
4663 assert_eq!(swept, vec![t.id.clone()]);
4664 }
4665
4666 #[test]
4667 fn a_running_task_behind_a_dead_daemons_lock_recovers_once_swept_and_keeps_its_history() {
4668 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
4673 let dir = tempfile::tempdir().unwrap();
4674 let queue = Queue::at(dir.path().to_path_buf());
4675 let mut t = task();
4676 t.id = "20260101-000000-crsh".to_owned();
4677 t.status = TaskStatus::Running;
4678 t.attempts = 1;
4679 t.runs.push("20260904-000000-4043".to_owned());
4683 queue.put(&mut t).unwrap();
4684 let dead_pid = injected_dead_pid();
4685
4686 std::fs::write(
4689 dir.path().join(format!("{}.lock", t.id)),
4690 dead_pid.to_string(),
4691 )
4692 .unwrap();
4693
4694 assert!(reclaim_orphaned_running(&queue, 2).is_empty());
4700 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
4701
4702 let swept = sweep_stale_claims_with(&queue, Duration::from_secs(6 * 60 * 60), |pid| {
4703 pid != dead_pid
4704 });
4705 assert_eq!(swept, vec![t.id.clone()]);
4706
4707 let reclaimed = reclaim_orphaned_running(&queue, 2);
4708 assert_eq!(reclaimed, vec![t.id.clone()]);
4709 let after = queue.get(&t.id).unwrap();
4710 assert_eq!(
4711 after.status,
4712 TaskStatus::Held,
4713 "no run.json to recover from, so a human is asked"
4714 );
4715 assert_eq!(
4716 after.runs,
4717 vec!["20260904-000000-4043".to_owned()],
4718 "the crashed run's id is kept as evidence, not discarded"
4719 );
4720 }
4721
4722 #[test]
4723 fn a_lock_is_kept_when_the_process_query_is_unavailable() {
4724 let dir = tempfile::tempdir().unwrap();
4725 let queue = Queue::at(dir.path().to_path_buf());
4726 let mut t = task();
4727 t.id = "20260101-000000-unknown".to_owned();
4728 queue.put(&mut t).unwrap();
4729 let dead_pid = injected_dead_pid();
4730 std::fs::write(
4731 dir.path().join(format!("{}.lock", t.id)),
4732 dead_pid.to_string(),
4733 )
4734 .unwrap();
4735
4736 let swept = sweep_stale_claims_with(&queue, Duration::ZERO, |_| true);
4737 assert!(swept.is_empty(), "an unknown pid must keep its lock");
4738 assert!(queue.claim(&t.id).is_err(), "the lock remains protective");
4739 }
4740
4741 fn run_state(status: RunStatus) -> RunState {
4742 let mut state = RunState::new(
4743 PathBuf::from("/repo"),
4744 "main".to_owned(),
4745 "abc1234def".to_owned(),
4746 "add retries".to_owned(),
4747 Config::default(),
4748 );
4749 state.status = status;
4750 state
4751 }
4752
4753 fn candidate(label: char, summary: &str, empty: bool, failed: Option<&str>) -> Candidate {
4754 Candidate {
4755 index: 0,
4756 label,
4757 agent: "claude".to_owned(),
4758 branch: format!("magi/x/{label}"),
4759 worktree: PathBuf::from("/repo"),
4760 summary: summary.to_owned(),
4761 stat: String::new(),
4762 files: 0,
4763 commits: usize::from(!empty),
4764 empty,
4765 failed: failed.map(str::to_owned),
4766 verified_noop: None,
4767 duration_ms: 0,
4768 folded: false,
4769 }
4770 }
4771
4772 #[test]
4773 fn diagnostic_names_the_failing_gate_checks_and_their_output() {
4774 let mut state = run_state(RunStatus::Blocked);
4775 state.gate = vec![
4776 CommandOutcome {
4777 command: "cargo make check".to_owned(),
4778 code: Some(0),
4779 output_tail: "ok".to_owned(),
4780 duration_ms: 0,
4781 resource_blocked: false,
4782 },
4783 CommandOutcome {
4784 command: "cargo test".to_owned(),
4785 code: Some(101),
4786 output_tail: "thread 'x' panicked: assertion failed".to_owned(),
4787 duration_ms: 0,
4788 resource_blocked: false,
4789 },
4790 ];
4791 let d = diagnostic(&state).expect("a failing gate must produce a diagnostic");
4792 assert!(d.contains("cargo test"), "{d}");
4793 assert!(
4794 !d.contains("cargo make check"),
4795 "a passing check is not a diagnostic: {d}"
4796 );
4797 assert!(d.contains("assertion failed"), "{d}");
4798 }
4799
4800 #[test]
4801 fn diagnostic_names_the_checks_the_fixer_gave_up_in_front_of() {
4802 let mut state = run_state(RunStatus::Blocked);
4803 state.event(
4804 "land",
4805 "stopped: the fixer produced no commit while 2 check(s) were failing \
4806 (build, lint); stopping instead of looping on an unchanged tree",
4807 );
4808 let d = diagnostic(&state).expect("a stalled land loop must produce a diagnostic");
4809 assert!(d.contains("build"), "{d}");
4810 assert!(d.contains("lint"), "{d}");
4811 assert!(d.contains("fixer produced no commit"), "{d}");
4812 }
4813
4814 #[test]
4815 fn describe_never_leaves_a_verified_noop_reading_as_a_bare_status_code() {
4816 let state = run_state(RunStatus::VerifiedNoop);
4821 let d = describe(&state);
4822 assert!(
4823 d.contains("agent-verified no-op"),
4824 "expected the display label, not the wire spelling: {d}"
4825 );
4826 assert!(!d.contains("verified_noop"), "{d}");
4827 }
4828
4829 #[test]
4830 fn diagnostic_carries_a_candidates_own_final_word_when_none_was_viable() {
4831 let mut state = run_state(RunStatus::Failed);
4837 state.candidates = vec![candidate(
4838 'A',
4839 "opened pull request #42, merged it, tagged v1.2.3 and published the release",
4840 true,
4841 None,
4842 )];
4843 let d = diagnostic(&state).expect("an empty candidate with a summary must be surfaced");
4844 assert!(d.contains("candidate A"), "{d}");
4845 assert!(d.contains("tagged v1.2.3"), "{d}");
4846 }
4847
4848 #[test]
4849 fn diagnostic_falls_back_to_a_candidates_failure_reason_when_it_has_no_summary() {
4850 let mut state = run_state(RunStatus::Failed);
4851 state.candidates = vec![candidate('A', "", true, Some("agent timed out"))];
4852 let d = diagnostic(&state).expect("a candidate's own failure reason must be surfaced");
4853 assert!(d.contains("candidate A"), "{d}");
4854 assert!(d.contains("agent timed out"), "{d}");
4855 }
4856
4857 #[test]
4858 fn diagnostic_is_none_when_nothing_recognisable_explains_the_hold() {
4859 let mut state = run_state(RunStatus::Failed);
4862 state.candidates = vec![candidate('A', "did the work", false, None)];
4863 assert!(diagnostic(&state).is_none());
4864 }
4865
4866 #[test]
4867 fn diagnostic_is_bounded_however_much_a_run_printed() {
4868 let mut state = run_state(RunStatus::Blocked);
4869 state.gate = vec![
4870 CommandOutcome {
4871 command: "cargo test".to_owned(),
4872 code: Some(101),
4873 output_tail: "x".repeat(50_000),
4874 duration_ms: 0,
4875 resource_blocked: false,
4876 },
4877 CommandOutcome {
4878 command: "cargo clippy".to_owned(),
4879 code: Some(1),
4880 output_tail: "y".repeat(50_000),
4881 duration_ms: 0,
4882 resource_blocked: false,
4883 },
4884 ];
4885 state.candidates = vec![
4886 candidate('A', &"z".repeat(50_000), true, None),
4887 candidate('B', &"w".repeat(50_000), true, None),
4888 ];
4889 let d = diagnostic(&state).expect("plenty here to diagnose");
4890 assert!(
4891 d.len() <= DIAGNOSTIC_MAX,
4892 "diagnostic grew to {} bytes, unbounded",
4893 d.len()
4894 );
4895 }
4896
4897 #[test]
4898 fn settle_and_diagnose_attaches_a_diagnostic_only_once_the_task_is_held() {
4899 let mut state = run_state(RunStatus::Blocked);
4900 state.gate = vec![CommandOutcome {
4901 command: "cargo test".to_owned(),
4902 code: Some(101),
4903 output_tail: "assertion failed".to_owned(),
4904 duration_ms: 0,
4905 resource_blocked: false,
4906 }];
4907 let verdict = Verdict {
4908 status: RunStatus::Blocked,
4909 left_pr: false,
4910 quota_hit: false,
4911 parked: false,
4912 no_viable_candidates: false,
4913 };
4914
4915 let mut t = task();
4918 t.start("run-1".to_owned());
4919 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
4920 assert_eq!(t.status, TaskStatus::Failed);
4921 assert!(t.diagnostic.is_none());
4922
4923 t.start("run-2".to_owned());
4926 settle_and_diagnose(&mut t, verdict, "gate failed", 2, &state);
4927 assert_eq!(t.status, TaskStatus::Held);
4928 let d = t.diagnostic.expect("a held task must carry its diagnostic");
4929 assert!(d.contains("cargo test"), "{d}");
4930 }
4931
4932 #[test]
4933 fn a_held_task_names_the_open_question_it_is_waiting_on() {
4934 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
4939 let home = crate::run::home();
4940 let state = run_state(RunStatus::VerifiedNoop);
4941 let mut q = ask::Question::new(
4942 state.id.clone(),
4943 "implement".to_owned(),
4944 "impl-A".to_owned(),
4945 "is this really a no-op?".to_owned(),
4946 String::new(),
4947 Vec::new(),
4948 );
4949 Questions::at(home.join("questions")).put(&mut q).unwrap();
4950
4951 let verdict = Verdict {
4952 status: RunStatus::VerifiedNoop,
4953 left_pr: false,
4954 quota_hit: false,
4955 parked: false,
4956 no_viable_candidates: false,
4957 };
4958 let mut t = task();
4959 t.start(state.id.clone());
4960 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
4961
4962 assert_eq!(t.status, TaskStatus::Held);
4963 let reason = t.hold_reason.expect("a held task must record why");
4964 assert!(
4965 reason.starts_with("run ended agent-verified no-op"),
4966 "the original settle reason must survive unchanged: {reason}"
4967 );
4968 assert!(
4969 reason.contains(q.short()),
4970 "the open question's id must be named so the notice is actionable: {reason}"
4971 );
4972 }
4973
4974 #[test]
4975 fn a_held_task_with_no_open_question_keeps_its_plain_reason() {
4976 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
4977 let state = run_state(RunStatus::VerifiedNoop);
4978
4979 let verdict = Verdict {
4980 status: RunStatus::VerifiedNoop,
4981 left_pr: false,
4982 quota_hit: false,
4983 parked: false,
4984 no_viable_candidates: false,
4985 };
4986 let mut t = task();
4987 t.start(state.id.clone());
4988 settle_and_diagnose(&mut t, verdict, "run ended agent-verified no-op", 5, &state);
4989
4990 assert_eq!(t.status, TaskStatus::Held);
4991 assert_eq!(
4992 t.hold_reason.as_deref(),
4993 Some("run ended agent-verified no-op"),
4994 "nothing to append when the question was already answered or never asked"
4995 );
4996 }
4997
4998 #[test]
4999 fn supersede_prior_runs_rewrites_an_earlier_blocked_attempt_once_a_later_one_lands() {
5000 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5001 let mut first = run_state(RunStatus::Blocked);
5002 first.id = "20260101-000000-sup1".to_owned();
5003 first.save().unwrap();
5004 let mut second = run_state(RunStatus::Merged);
5005 second.id = "20260101-000000-sup2".to_owned();
5006 second.save().unwrap();
5007
5008 let mut t = task();
5009 t.runs = vec![first.id.clone(), second.id.clone()];
5010 t.status = TaskStatus::Done;
5011
5012 supersede_prior_runs(&t, &crate::run::home());
5013
5014 assert_eq!(
5015 RunState::load(&first.id).unwrap().status,
5016 RunStatus::Superseded,
5017 "the first attempt's Blocked no longer needs anyone's attention"
5018 );
5019 assert_eq!(
5020 RunState::load(&second.id).unwrap().status,
5021 RunStatus::Merged,
5022 "the run that actually succeeded is left exactly as it was"
5023 );
5024 }
5025
5026 #[test]
5027 fn supersede_prior_runs_leaves_a_manually_resumed_attempt_alone() {
5028 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5034 let mut first = run_state(RunStatus::Blocked);
5035 first.id = "20260101-000000-sup9".to_owned();
5036 first.driver_pid = Some(std::process::id());
5039 first.driver_started_at = Some(
5040 crate::proc::process_started_at(std::process::id())
5041 .expect("this test process's own start time must be queryable"),
5042 );
5043 first.save().unwrap();
5044 let mut second = run_state(RunStatus::Merged);
5045 second.id = "20260101-000000-supa".to_owned();
5046 second.save().unwrap();
5047
5048 let mut t = task();
5049 t.runs = vec![first.id.clone(), second.id.clone()];
5050 t.status = TaskStatus::Done;
5051
5052 supersede_prior_runs(&t, &crate::run::home());
5053
5054 assert_eq!(
5055 RunState::load(&first.id).unwrap().status,
5056 RunStatus::Blocked,
5057 "a live driver_pid means something is still actually working this run, \
5058 even though no daemon claims it - rewriting under it would just be \
5059 undone the next time that process saves"
5060 );
5061 }
5062
5063 #[test]
5064 fn resweep_catches_up_a_run_left_live_once_its_manual_process_is_no_longer_driving_it() {
5065 let dir = tempfile::tempdir().unwrap();
5071 let home = dir.path().to_path_buf();
5072 let queue = Queue::at(dir.path().join("queue"));
5073
5074 let mut first = run_state(RunStatus::Blocked);
5075 first.id = "20260101-000000-supd".to_owned();
5076 first.driver_pid = Some(std::process::id());
5077 first.driver_started_at = Some(
5078 crate::proc::process_started_at(std::process::id())
5079 .expect("this test process's own start time must be queryable"),
5080 );
5081 first.save_under(&home).unwrap();
5082 let mut second = run_state(RunStatus::Merged);
5083 second.id = "20260101-000000-supe".to_owned();
5084 second.save_under(&home).unwrap();
5085
5086 let mut t = task();
5087 t.runs = vec![first.id.clone(), second.id.clone()];
5088 t.status = TaskStatus::Done;
5089 queue.put(&mut t).unwrap();
5090
5091 resweep_superseded_attempts(&queue, &home);
5092 assert_eq!(
5093 RunState::load_under(&first.id, &home).unwrap().status,
5094 RunStatus::Blocked,
5095 "still live on the first pass, so still untouched"
5096 );
5097
5098 let mut stale = RunState::load_under(&first.id, &home).unwrap();
5104 stale.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
5105 stale.save_under(&home).unwrap();
5106
5107 resweep_superseded_attempts(&queue, &home);
5108 assert_eq!(
5109 RunState::load_under(&first.id, &home).unwrap().status,
5110 RunStatus::Superseded,
5111 "the second pass catches up what the first one correctly skipped"
5112 );
5113 }
5114
5115 #[test]
5116 fn supersede_prior_runs_leaves_concurrent_blocked_attempts_alone_while_the_task_is_not_done() {
5117 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5118 let mut first = run_state(RunStatus::Blocked);
5119 first.id = "20260101-000000-sup3".to_owned();
5120 first.save().unwrap();
5121 let mut second = run_state(RunStatus::Blocked);
5122 second.id = "20260101-000000-sup4".to_owned();
5123 second.save().unwrap();
5124
5125 let mut t = task();
5126 t.runs = vec![first.id.clone(), second.id.clone()];
5127 t.status = TaskStatus::Failed;
5131
5132 supersede_prior_runs(&t, &crate::run::home());
5133
5134 assert_eq!(
5135 RunState::load(&first.id).unwrap().status,
5136 RunStatus::Blocked
5137 );
5138 assert_eq!(
5139 RunState::load(&second.id).unwrap().status,
5140 RunStatus::Blocked
5141 );
5142 }
5143
5144 #[test]
5145 fn supersede_prior_runs_does_nothing_when_the_task_was_closed_by_hand() {
5146 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5150 let mut first = run_state(RunStatus::Blocked);
5151 first.id = "20260101-000000-sup5".to_owned();
5152 first.save().unwrap();
5153
5154 let mut t = task();
5155 t.runs = vec![first.id.clone()];
5156 t.status = TaskStatus::Done;
5157
5158 supersede_prior_runs(&t, &crate::run::home());
5159
5160 assert_eq!(
5161 RunState::load(&first.id).unwrap().status,
5162 RunStatus::Blocked,
5163 "a single-attempt task has no earlier run to supersede"
5164 );
5165 }
5166
5167 #[test]
5168 fn supersede_prior_runs_does_nothing_when_the_last_recorded_attempt_never_landed() {
5169 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5176 let mut first = run_state(RunStatus::Blocked);
5177 first.id = "20260101-000000-supb".to_owned();
5178 first.save().unwrap();
5179 let mut second = run_state(RunStatus::Failed);
5180 second.id = "20260101-000000-supc".to_owned();
5181 second.save().unwrap();
5182
5183 let mut t = task();
5184 t.runs = vec![first.id.clone(), second.id.clone()];
5185 t.status = TaskStatus::Done;
5186
5187 supersede_prior_runs(&t, &crate::run::home());
5188
5189 assert_eq!(
5190 RunState::load(&first.id).unwrap().status,
5191 RunStatus::Blocked,
5192 "the task's last attempt never landed, so there is nothing here \
5193 actually superseding it"
5194 );
5195 }
5196
5197 #[test]
5198 fn supersede_prior_runs_leaves_a_failed_or_verified_noop_attempt_as_is() {
5199 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5203 let mut failed = run_state(RunStatus::Failed);
5204 failed.id = "20260101-000000-sup6".to_owned();
5205 failed.save().unwrap();
5206 let mut noop = run_state(RunStatus::VerifiedNoop);
5207 noop.id = "20260101-000000-sup7".to_owned();
5208 noop.save().unwrap();
5209 let mut winner = run_state(RunStatus::Ready);
5210 winner.id = "20260101-000000-sup8".to_owned();
5211 winner.save().unwrap();
5212
5213 let mut t = task();
5214 t.runs = vec![failed.id.clone(), noop.id.clone(), winner.id.clone()];
5215 t.status = TaskStatus::Done;
5216
5217 supersede_prior_runs(&t, &crate::run::home());
5218
5219 assert_eq!(
5220 RunState::load(&failed.id).unwrap().status,
5221 RunStatus::Failed
5222 );
5223 assert_eq!(
5224 RunState::load(&noop.id).unwrap().status,
5225 RunStatus::VerifiedNoop
5226 );
5227 }
5228
5229 fn approval_question(run: &str) -> ask::Question {
5230 ask::Question::new(
5231 run.to_owned(),
5232 land::APPROVAL_NODE.to_owned(),
5233 "land".to_owned(),
5234 "merge?".to_owned(),
5235 String::new(),
5236 vec!["merge".to_owned(), "hold".to_owned()],
5237 )
5238 }
5239
5240 #[test]
5241 fn land_resume_state_leaves_a_fresh_open_question_waiting() {
5242 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5243 let mut state = run_state(RunStatus::Landing);
5244 state.id = "20260101-000000-fre1".to_owned();
5245 state.parked = true;
5246 state.save().unwrap();
5247 ask::Questions::open()
5248 .put(&mut approval_question(&state.id))
5249 .unwrap();
5250
5251 let mut t = task();
5252 t.runs.push(state.id.clone());
5253 assert_eq!(
5254 land_resume_state(&t),
5255 LandResume::StillWaiting,
5256 "nobody has answered and the timeout has not passed"
5257 );
5258 }
5259
5260 #[test]
5261 fn land_resume_state_abandons_a_question_that_outlived_answer_timeout() {
5262 crate::run::set_home(std::env::temp_dir().join("magi-daemon-test-home"));
5267 let mut state = run_state(RunStatus::Landing);
5268 state.id = "20260101-000000-exp1".to_owned();
5269 state.parked = true;
5270 state.config.graph.answer_timeout = 60;
5271 state.save().unwrap();
5272
5273 let store = ask::Questions::open();
5274 let mut q = approval_question(&state.id);
5275 q.asked_at = Timestamp::now() - jiff::SignedDuration::from_secs(120);
5276 store.put(&mut q).unwrap();
5277
5278 let mut t = task();
5279 t.runs.push(state.id.clone());
5280 assert_eq!(
5281 land_resume_state(&t),
5282 LandResume::Ready,
5283 "an expired question must not be waited on forever"
5284 );
5285
5286 let after = store.get(&q.id).unwrap();
5287 assert!(
5288 !after.status.open(),
5289 "the question is abandoned, not silently ignored"
5290 );
5291 assert!(
5292 after.resolution().is_none(),
5293 "an abandoned question is not read as a decision"
5294 );
5295 }
5296
5297 #[test]
5298 fn reclaim_settles_a_running_task_against_its_last_run() {
5299 let mut t = task();
5300 t.start("20260904-000000-4043".to_owned());
5301 reclaim(&mut t, Some(run_state(RunStatus::Ready)), 2);
5302 assert_eq!(
5303 t.status,
5304 TaskStatus::Done,
5305 "a run that actually finished must not stay `running` forever"
5306 );
5307 }
5308
5309 #[test]
5310 fn reclaim_reuses_the_same_retry_policy_as_a_live_settle() {
5311 let mut t = task();
5315 t.start("20260904-000000-4043".to_owned());
5316 reclaim(&mut t, Some(run_state(RunStatus::Blocked)), 2);
5317 assert_eq!(t.status, TaskStatus::Failed);
5318 assert!(t.status.runnable());
5319 }
5320
5321 #[test]
5322 fn reclaim_holds_a_running_task_whose_run_cannot_be_found() {
5323 let mut t = task();
5324 t.start("20260904-000000-4043".to_owned());
5325 reclaim(&mut t, None, 2);
5326 assert_eq!(t.status, TaskStatus::Held);
5327 assert!(
5328 t.last_error
5329 .as_deref()
5330 .is_some_and(|e| e.contains("running")),
5331 "the operator needs to know why this task was held"
5332 );
5333 }
5334
5335 #[test]
5336 fn orphaned_running_tasks_are_reclaimed_but_live_ones_are_left_alone() {
5337 let dir = tempfile::tempdir().unwrap();
5338 let queue = Queue::at(dir.path().to_path_buf());
5339
5340 let mut orphaned = task();
5342 orphaned.id = "20260904-000000-orph".to_owned();
5343 orphaned.status = TaskStatus::Running;
5344 orphaned.attempts = 1;
5345 queue.put(&mut orphaned).unwrap();
5346
5347 let mut alive = task();
5348 alive.id = "20260904-000000-live".to_owned();
5349 alive.status = TaskStatus::Running;
5350 alive.attempts = 1;
5351 queue.put(&mut alive).unwrap();
5352 let _held_by_a_live_daemon = queue.claim(&alive.id).unwrap();
5353
5354 let mut queued = task();
5355 queued.id = "20260904-000000-wait".to_owned();
5356 queue.put(&mut queued).unwrap();
5357
5358 let reclaimed = reclaim_orphaned_running(&queue, 2);
5359 assert_eq!(reclaimed, vec![orphaned.id.clone()]);
5360
5361 assert_eq!(
5362 queue.get(&orphaned.id).unwrap().status,
5363 TaskStatus::Held,
5364 "nothing was driving it and there was no run to recover"
5365 );
5366 assert_eq!(
5367 queue.get(&alive.id).unwrap().status,
5368 TaskStatus::Running,
5369 "a live claim must protect the task it belongs to"
5370 );
5371 assert_eq!(queue.get(&queued.id).unwrap().status, TaskStatus::Queued);
5372 }
5373
5374 fn read_run_under(home: &Path, id: &str) -> RunState {
5380 let body = std::fs::read_to_string(home.join("runs").join(id).join("run.json")).unwrap();
5381 serde_json::from_str(&body).unwrap()
5382 }
5383
5384 #[test]
5385 fn reclaim_abandoned_runs_fails_a_run_whose_active_seats_are_all_provably_dead() {
5386 let dir = tempfile::tempdir().unwrap();
5387 let home = dir.path().to_path_buf();
5388 let now = Timestamp::now();
5389 let overrun_seat = || crate::run::ActiveSeat {
5390 node: "implement".to_owned(),
5391 started_at: now - jiff::SignedDuration::new(21_000, 0),
5392 timeout_secs: 3_600,
5393 attempt: 0,
5394 task: None,
5395 command: None,
5396 index: None,
5397 total: None,
5398 };
5399
5400 let mut dead = run_state(RunStatus::Implementing);
5401 dead.id = "20260101-000000-dead".to_owned();
5402 dead.active.insert("impl-A".to_owned(), overrun_seat());
5403 dead.driver_pid = Some(4242);
5406 dead.save_under(&home).unwrap();
5407
5408 let mut alive = run_state(RunStatus::Implementing);
5411 alive.id = "20260101-000000-aliv".to_owned();
5412 alive.active.insert("impl-A".to_owned(), overrun_seat());
5413 alive.save_under(&home).unwrap();
5414 let mut status = Status::new();
5415 status.current = vec![Current {
5416 task: "20260101-000000-task".to_owned(),
5417 run: alive.id.clone(),
5418 }];
5419 write_status_to(&home.join("daemon.json"), &status).unwrap();
5420
5421 let questions = Questions::at(home.join("questions"));
5425 let mut q = ask::Question::new(
5426 dead.id.clone(),
5427 "implement".to_owned(),
5428 "impl-A".to_owned(),
5429 "Which storage backend?".to_owned(),
5430 String::new(),
5431 vec!["SQLite".to_owned(), "Redis".to_owned()],
5432 );
5433 questions.put(&mut q).unwrap();
5434
5435 let abandoned = reclaim_abandoned_runs_with(
5436 &home,
5437 now,
5438 |pid| if pid == 4242 { Some(false) } else { None },
5439 |_| panic!("a query answering Dead outright needs no identity corroboration"),
5440 );
5441 assert_eq!(abandoned, vec![dead.id.clone()]);
5442
5443 let reloaded = read_run_under(&home, &dead.id);
5444 assert_eq!(reloaded.status, RunStatus::Failed);
5445 assert!(reloaded.active.is_empty());
5446 assert!(
5447 !questions.get(&q.id).unwrap().status.open(),
5448 "the failed run's own open question must be settled in the same pass"
5449 );
5450
5451 let still_alive = read_run_under(&home, &alive.id);
5452 assert_eq!(
5453 still_alive.status,
5454 RunStatus::Implementing,
5455 "a live daemon's claim protects it"
5456 );
5457 assert!(!still_alive.active.is_empty());
5458 }
5459
5460 #[test]
5470 fn reclaim_abandoned_runs_leaves_a_live_manual_run_alone_even_though_no_daemon_claims_it() {
5471 let dir = tempfile::tempdir().unwrap();
5472 let home = dir.path().to_path_buf();
5473 let now = Timestamp::now();
5474
5475 let mut manual = run_state(RunStatus::Reviewing);
5476 manual.id = "20260101-000000-manl".to_owned();
5477 manual.active.insert(
5478 "review-1".to_owned(),
5479 crate::run::ActiveSeat {
5480 node: "review".to_owned(),
5481 started_at: now - jiff::SignedDuration::new(21_000, 0),
5482 timeout_secs: 3_600,
5483 attempt: 0,
5484 task: None,
5485 command: None,
5486 index: None,
5487 total: None,
5488 },
5489 );
5490 manual.driver_pid = Some(4242);
5494 manual.driver_started_at = Some("2026-09-22T10:00:00Z".to_owned());
5495 manual.save_under(&home).unwrap();
5496
5497 let abandoned = reclaim_abandoned_runs_with(
5498 &home,
5499 now,
5500 |pid| if pid == 4242 { Some(true) } else { None },
5501 |pid| {
5502 if pid == 4242 {
5503 Some("2026-09-22T10:00:00Z".to_owned())
5504 } else {
5505 None
5506 }
5507 },
5508 );
5509 assert!(
5510 abandoned.is_empty(),
5511 "a manual run a real process is still driving must never be reclaimed: {abandoned:?}"
5512 );
5513
5514 let reloaded = read_run_under(&home, &manual.id);
5515 assert_eq!(reloaded.status, RunStatus::Reviewing);
5516 assert!(!reloaded.active.is_empty());
5517 }
5518
5519 #[test]
5520 fn an_already_claimed_task_is_skipped_rather_than_failed() {
5521 let dir = tempfile::tempdir().unwrap();
5522 let queue = Queue::at(dir.path().to_path_buf());
5523 let mut only = task();
5524 queue.put(&mut only).unwrap();
5525
5526 let _elsewhere = queue.claim(&only.id).unwrap();
5527 let candidates = runnable(&queue);
5528 assert_eq!(candidates.len(), 1, "the task is still runnable");
5529 assert!(
5530 queue.claim(&candidates[0].id).is_err(),
5531 "the loop cannot take a claim somebody else holds"
5532 );
5533
5534 let after = queue.get(&only.id).unwrap();
5535 assert_eq!(after.status, TaskStatus::Queued);
5536 assert_eq!(
5537 after.attempts, 0,
5538 "losing the race is not an attempt at the task"
5539 );
5540 assert_eq!(after.last_error, None);
5541 }
5542
5543 #[test]
5544 fn the_status_file_round_trips_and_its_heartbeat_advances() {
5545 let dir = tempfile::tempdir().unwrap();
5546 let path = dir.path().join("daemon.json");
5547
5548 let mut status = Status::new();
5549 status.idle = false;
5550 status.completed = 7;
5551 status.current = vec![Current {
5552 task: "20260902-000000-t111".to_owned(),
5553 run: "20260902-000001-r111".to_owned(),
5554 }];
5555 write_status_to(&path, &status).unwrap();
5556 let first: Status = serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
5557 assert_eq!(first.schema, SCHEMA);
5558 assert_eq!(first.pid, std::process::id());
5559 assert!(!first.idle);
5560 assert_eq!(first.completed, 7);
5561 assert_eq!(first.current, status.current);
5562 assert!(
5563 !path.with_extension("json.tmp").exists(),
5564 "the temp file is renamed, not left behind"
5565 );
5566
5567 std::thread::sleep(Duration::from_millis(5));
5568 status.updated_at = Timestamp::now();
5569 status.polls = 3;
5570 write_status_to(&path, &status).unwrap();
5571 let second: Status =
5572 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
5573 assert!(
5574 second.updated_at > first.updated_at,
5575 "a reader can only detect staleness if the heartbeat moves"
5576 );
5577 assert_eq!(
5578 second.started_at, first.started_at,
5579 "the start time is not a heartbeat"
5580 );
5581 assert_eq!(second.polls, 3);
5582 }
5583
5584 #[test]
5585 fn reading_counts_as_running_only_while_its_heartbeat_is_fresh() {
5586 let dir = tempfile::tempdir().unwrap();
5587
5588 assert!(read_status(dir.path()).is_none(), "no file, no daemon");
5589
5590 let mut status = Status::new();
5591 status.updated_at = Timestamp::now() - jiff::SignedDuration::from_secs(60);
5592 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5593 let stale = read_status(dir.path()).unwrap();
5594 assert!(
5595 !stale.running(Timestamp::now()),
5596 "a minute without a heartbeat is a dead daemon, not a busy one"
5597 );
5598 assert!(stale.age_secs(Timestamp::now()).is_some_and(|s| s >= 55));
5599
5600 status.updated_at = Timestamp::now();
5601 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5602 let fresh = read_status(dir.path()).unwrap();
5603 assert!(fresh.running(Timestamp::now()));
5604 }
5605
5606 #[test]
5607 fn only_a_live_daemon_on_this_very_run_counts_as_working_on_it() {
5608 let dir = tempfile::tempdir().unwrap();
5609 let now = Timestamp::now();
5610 let mine = "20260903-080619-01c2";
5611
5612 assert!(
5613 !is_working_on(dir.path(), mine, now),
5614 "no status file means nobody is working on anything"
5615 );
5616
5617 let mut status = Status::new();
5618 status.current = vec![Current {
5619 task: "20260903-080340-0167".to_owned(),
5620 run: mine.to_owned(),
5621 }];
5622 status.updated_at = now;
5623 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5624 assert!(is_working_on(dir.path(), mine, now));
5625 assert!(
5626 !is_working_on(dir.path(), "20260903-105039-3cbf", now),
5627 "a daemon busy with one run is not working on another"
5628 );
5629
5630 status.updated_at = now - jiff::SignedDuration::from_secs(600);
5633 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5634 assert!(
5635 !is_working_on(dir.path(), mine, now),
5636 "a stale heartbeat is a dead daemon, so its run is a leftover"
5637 );
5638 }
5639
5640 #[test]
5641 fn is_working_on_short_matches_by_the_worktree_bays_own_name() {
5642 let dir = tempfile::tempdir().unwrap();
5643 let now = Timestamp::now();
5644
5645 assert!(
5646 !is_working_on_short(dir.path(), "01c2", now),
5647 "no status file means nobody is working on anything"
5648 );
5649
5650 let mut status = Status::new();
5651 status.current = vec![Current {
5652 task: "20260903-080340-0167".to_owned(),
5653 run: "20260903-080619-01c2".to_owned(),
5654 }];
5655 status.updated_at = now;
5656 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
5657 assert!(
5658 is_working_on_short(dir.path(), "01c2", now),
5659 "the run's short id is the last block of its full id"
5660 );
5661 assert!(
5662 !is_working_on_short(dir.path(), "3cbf", now),
5663 "a daemon busy with one worktree bay is not working on another"
5664 );
5665 }
5666
5667 #[test]
5668 fn a_newer_status_file_still_yields_a_reading() {
5669 let dir = tempfile::tempdir().unwrap();
5670 std::fs::write(
5673 dir.path().join("daemon.json"),
5674 serde_json::json!({
5675 "schema": 2,
5676 "updated_at": Timestamp::now().to_string(),
5677 "idle": true,
5678 "surprise": { "nested": [1, 2, 3] },
5679 })
5680 .to_string(),
5681 )
5682 .unwrap();
5683
5684 let reading = read_status(dir.path()).expect("a forward-compatible read");
5685 assert!(reading.running(Timestamp::now()));
5686 assert!(reading.idle);
5687 assert!(reading.current.is_empty());
5688 }
5689
5690 #[test]
5691 fn an_older_daemons_single_object_current_still_reads_as_a_one_item_list() {
5692 let dir = tempfile::tempdir().unwrap();
5698 std::fs::write(
5699 dir.path().join("daemon.json"),
5700 serde_json::json!({
5701 "schema": 1,
5702 "pid": 4242,
5703 "updated_at": Timestamp::now().to_string(),
5704 "idle": false,
5705 "current": {"task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb"},
5706 "completed": 3,
5707 "polls": 9,
5708 })
5709 .to_string(),
5710 )
5711 .unwrap();
5712
5713 let reading = read_status(dir.path()).expect("an older shape must still parse");
5714 assert!(reading.running(Timestamp::now()));
5715 assert_eq!(
5716 reading.current,
5717 vec![Current {
5718 task: "20260902-140501-aaaa".to_owned(),
5719 run: "20260902-140502-bbbb".to_owned(),
5720 }]
5721 );
5722 }
5723
5724 #[test]
5725 fn an_absent_or_null_current_reads_as_idle_not_a_parse_failure() {
5726 let dir = tempfile::tempdir().unwrap();
5727 std::fs::write(
5728 dir.path().join("daemon.json"),
5729 serde_json::json!({
5730 "schema": 1,
5731 "updated_at": Timestamp::now().to_string(),
5732 "idle": true,
5733 "current": null,
5734 })
5735 .to_string(),
5736 )
5737 .unwrap();
5738 let with_null = read_status(dir.path()).expect("null must still parse");
5739 assert!(with_null.current.is_empty());
5740
5741 std::fs::write(
5742 dir.path().join("daemon.json"),
5743 serde_json::json!({
5744 "schema": 1,
5745 "updated_at": Timestamp::now().to_string(),
5746 "idle": true,
5747 })
5748 .to_string(),
5749 )
5750 .unwrap();
5751 let absent = read_status(dir.path()).expect("a missing field must still parse");
5752 assert!(absent.current.is_empty());
5753 }
5754
5755 #[test]
5756 fn a_task_without_a_repository_runs_in_the_daemons_default() {
5757 let fallback = Path::new("/default");
5758 let mut blank = task();
5759 blank.repo = PathBuf::new();
5760 assert_eq!(repo_for(&blank, fallback), PathBuf::from("/default"));
5761 let mut dot = task();
5762 dot.repo = PathBuf::from(".");
5763 assert_eq!(repo_for(&dot, fallback), PathBuf::from("/default"));
5764 assert_eq!(
5765 repo_for(&task(), fallback),
5766 PathBuf::from("/repo"),
5767 "a task that names a repository keeps it"
5768 );
5769 }
5770
5771 #[test]
5772 fn a_solo_task_runs_with_one_candidate_and_a_plain_task_keeps_the_configs() {
5773 let mut solo_cfg = Config::default();
5779 solo_cfg.graph.candidates = 3;
5780 let mut solo_task = task();
5781 solo_task.solo = true;
5782 apply_solo(&mut solo_cfg, &solo_task);
5783 assert_eq!(solo_cfg.graph.candidates, 1);
5784
5785 let mut plain_cfg = Config::default();
5786 plain_cfg.graph.candidates = 3;
5787 let plain_task = task();
5788 assert!(!plain_task.solo);
5789 apply_solo(&mut plain_cfg, &plain_task);
5790 assert_eq!(
5791 plain_cfg.graph.candidates, 3,
5792 "a task that did not ask to run alone keeps the config's candidates"
5793 );
5794 }
5795
5796 fn loss(seat: &str, at: &str, reset: Option<&str>) -> QuotaLoss {
5797 QuotaLoss {
5798 seat: seat.into(),
5799 node: "judge".into(),
5800 at: at.parse().unwrap(),
5801 reset: reset.map(str::to_string),
5802 }
5803 }
5804
5805 #[test]
5806 fn a_resumed_run_with_only_old_quota_losses_arms_no_cooldown() {
5807 let old: Vec<QuotaLoss> = (1..=4)
5808 .map(|i| {
5809 loss(
5810 &format!("judge-{i}"),
5811 "2026-09-23T05:23:00Z",
5812 Some("2:40pm (Asia/Tokyo)"),
5813 )
5814 })
5815 .collect();
5816 let fresh = losses_this_attempt(&old, &old);
5817 assert!(fresh.is_empty());
5818 assert_eq!(cooldown_until(&fresh, Timestamp::now()), None);
5819 }
5821
5822 #[test]
5823 fn a_new_quota_loss_during_the_attempt_still_arms_the_cooldown() {
5824 let old = vec![loss("judge-1", "2026-09-23T05:23:00Z", None)];
5825 let now = Timestamp::now();
5826 let mut after = old.clone();
5827 after.push(loss("judge-2", &now.to_string(), None));
5828 let fresh = losses_this_attempt(&old, &after);
5829 assert_eq!(fresh, vec![after[1].clone()]);
5830 let until = cooldown_until(&fresh, now).expect("a fresh loss arms the cooldown");
5831 assert_eq!(
5832 until,
5833 now + jiff::SignedDuration::from_secs(QUOTA_WAIT_FALLBACK.as_secs() as i64)
5834 );
5835 }
5836
5837 #[test]
5838 fn a_recovered_seat_dropping_out_of_the_history_does_not_hide_a_new_loss() {
5839 let before = vec![
5842 loss("judge-1", "2026-09-23T05:23:00Z", None),
5843 loss("judge-2", "2026-09-23T05:24:00Z", None),
5844 ];
5845 let after = vec![
5846 loss("judge-2", "2026-09-23T05:24:00Z", None),
5847 loss("judge-1", "2026-09-24T01:00:00Z", None),
5848 ];
5849 assert_eq!(losses_this_attempt(&before, &after), vec![after[1].clone()]);
5850 }
5851
5852 #[test]
5853 fn merge_overrides_are_parsed_or_refused() {
5854 assert_eq!(merge_mode("none").unwrap(), MergeMode::None);
5855 assert_eq!(merge_mode("local").unwrap(), MergeMode::Local);
5856 assert_eq!(merge_mode("pr").unwrap(), MergeMode::Pr);
5857 assert!(merge_mode("squash").is_err());
5858 }
5859
5860 #[test]
5861 fn quota_wait_uses_a_future_reset_time_capped_and_falls_back_otherwise() {
5862 let now = Timestamp::now();
5863 let fallback = Duration::from_secs(300);
5864 let cap = Duration::from_secs(1800);
5865
5866 assert_eq!(quota_wait(None, now, fallback, cap), fallback);
5868
5869 let soon = now + jiff::SignedDuration::from_secs(600);
5871 assert_eq!(
5872 quota_wait(Some(soon), now, fallback, cap),
5873 Duration::from_secs(600)
5874 );
5875
5876 let past = now - jiff::SignedDuration::from_secs(60);
5879 assert_eq!(quota_wait(Some(past), now, fallback, cap), fallback);
5880
5881 let far = now + jiff::SignedDuration::from_secs(3 * 3600);
5884 assert_eq!(quota_wait(Some(far), now, fallback, cap), cap);
5885 }
5886
5887 #[test]
5888 fn parse_reset_hint_reads_the_claude_cli_shape_and_rolls_a_past_clock_to_tomorrow() {
5889 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
5890
5891 let at = parse_reset_hint("4:50am (UTC)", now, now).expect("a recognised shape parses");
5892 assert_eq!(at.to_string(), "2026-09-07T04:50:00Z");
5893
5894 let already_past =
5898 parse_reset_hint("1:00am (UTC)", now, now).expect("a recognised shape parses");
5899 assert_eq!(already_past.to_string(), "2026-09-08T01:00:00Z");
5900
5901 assert!(
5902 parse_reset_hint("session limit reached", now, now).is_none(),
5903 "free text with no recognised shape is not guessed at"
5904 );
5905 assert!(
5906 parse_reset_hint("4:50am (Nowhere/Fake)", now, now).is_none(),
5907 "an unresolvable zone name is not guessed at either"
5908 );
5909 }
5910
5911 #[test]
5912 fn parse_reset_hint_reads_the_codex_cli_shape_with_no_year_rollover_needed() {
5913 let now = "2026-09-07T02:50:00Z".parse::<Timestamp>().unwrap();
5914
5915 let at = parse_reset_hint(
5916 "You've hit your usage limit. Visit \
5917 https://chatgpt.com/codex/settings/usage to purchase more \
5918 credits or try again at Sep 19th, 2026 5:10 PM.",
5919 now,
5920 now,
5921 )
5922 .expect("the codex reset wording is a recognised shape");
5923 assert_eq!(at.to_string(), "2026-09-19T17:10:00Z");
5924
5925 let earlier = parse_reset_hint("try again at Jan 2nd, 2026 1:00 AM.", now, now)
5930 .expect("an explicit year needs no rollover");
5931 assert_eq!(earlier.to_string(), "2026-01-02T01:00:00Z");
5932
5933 assert!(
5934 parse_reset_hint("try again at Sep 19th, 26 5:10 PM.", now, now).is_none(),
5935 "a two-digit year is not the documented shape and is not guessed at"
5936 );
5937 assert!(
5938 parse_reset_hint("try again at Sept 19th, 2026 5:10 PM.", now, now).is_none(),
5939 "a four-letter month name is not the documented three-letter abbreviation"
5940 );
5941 assert!(
5942 parse_reset_hint("try again at Sep 19th, 2026 5:10 PM (UTC).", now, now).is_none(),
5943 "an explicit zone on the dated shape is a format nobody has \
5944 documented, and is refused rather than guessed at as UTC"
5945 );
5946 }
5947
5948 #[test]
5949 fn parse_reset_hint_reads_agys_relative_shape_from_when_the_loss_was_recorded() {
5950 let now = "2026-09-24T12:00:00Z".parse::<Timestamp>().unwrap();
5951 let recorded = "2026-09-24T08:00:00Z".parse::<Timestamp>().unwrap();
5952
5953 let at = parse_reset_hint("in 1h2m49s", now, recorded).expect("agy's shape parses");
5954 assert_eq!(at.as_second() - recorded.as_second(), 3769);
5955
5956 let partial = parse_reset_hint("in 45m", now, recorded).expect("units are optional");
5957 assert_eq!(partial.as_second() - recorded.as_second(), 45 * 60);
5958
5959 for bad in ["in ", "in 45", "in 3x", "in m", "in 1h junk", "1h2m"] {
5960 assert!(
5961 parse_reset_hint(bad, now, recorded).is_none(),
5962 "{bad:?} must not be guessed at"
5963 );
5964 }
5965 }
5966
5967 fn idle_loop(dir: &Path) -> (Opts, Queue, PathBuf, PathBuf, PathBuf) {
5971 let config = dir.join("magi.toml");
5972 std::fs::write(
5973 &config,
5974 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = 0\n",
5975 )
5976 .unwrap();
5977 let opts = Opts {
5978 poll: Duration::from_secs(30),
5979 config: Some(config),
5980 repo: dir.join("repo"),
5984 ..Opts::default()
5985 };
5986 let home = dir.join("home");
5995 let worktrees = dir.join("wt");
5996 (
5997 opts,
5998 Queue::at(dir.join("queue")),
5999 home.join("daemon.json"),
6000 home,
6001 worktrees,
6002 )
6003 }
6004
6005 #[test]
6006 fn a_stop_is_idempotent_and_once_set_stays_set() {
6007 let stop = Stop::new();
6008 assert!(!stop.stopped());
6009
6010 stop.stop();
6011 assert!(stop.stopped());
6012 stop.stop();
6013 assert!(stop.stopped(), "a second stop is not a toggle");
6014
6015 let shared = stop.clone();
6016 assert!(
6017 shared.stopped(),
6018 "a clone is the same stop; that is how the loop and its caller share one"
6019 );
6020 }
6021
6022 #[test]
6023 fn only_a_stop_with_a_run_in_flight_reads_as_finishing() {
6024 let stop = Stop::new();
6025 stop.enter();
6026 assert!(
6027 !stop.finishing(),
6028 "a busy loop nobody has asked to stop is just running"
6029 );
6030
6031 stop.stop();
6032 assert!(
6033 stop.finishing(),
6034 "a stop asked for mid-run has not landed until the run is settled"
6035 );
6036
6037 stop.exit();
6038 assert!(
6039 !stop.finishing(),
6040 "once the run is settled the stop has landed and there is nothing to finish"
6041 );
6042 }
6043
6044 #[test]
6045 fn finishing_stays_true_until_the_last_of_several_runs_exits() {
6046 let stop = Stop::new();
6047 stop.enter();
6048 stop.enter();
6049 stop.stop();
6050 assert!(stop.finishing(), "two runs still in flight");
6051
6052 stop.exit();
6053 assert!(
6054 stop.finishing(),
6055 "one run finished, but a sibling is still working"
6056 );
6057
6058 stop.exit();
6059 assert!(
6060 !stop.finishing(),
6061 "the last run out is what actually lands the stop"
6062 );
6063 }
6064
6065 #[tokio::test]
6066 async fn a_loop_already_asked_to_stop_returns_without_waiting_out_a_poll() {
6067 let dir = tempfile::tempdir().unwrap();
6068 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6069 let stop = Stop::new();
6070 stop.stop();
6071
6072 let began = std::time::Instant::now();
6073 tokio::time::timeout(
6074 Duration::from_secs(2),
6075 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6076 )
6077 .await
6078 .expect("a stopped loop must return, not sit out its poll interval")
6079 .expect("the loop's own setup and teardown must not fail");
6080 assert!(
6081 began.elapsed() < opts.poll,
6082 "returned only after {:?}, which is a poll interval, not a stop",
6083 began.elapsed()
6084 );
6085 }
6086
6087 #[tokio::test]
6088 async fn a_stop_while_idle_wakes_the_wait_instead_of_sleeping_it_out() {
6089 let dir = tempfile::tempdir().unwrap();
6090 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6091 let stop = Stop::new();
6092
6093 let asker = {
6096 let stop = stop.clone();
6097 tokio::spawn(async move {
6098 tokio::time::sleep(Duration::from_millis(20)).await;
6099 stop.stop();
6100 })
6101 };
6102
6103 let began = std::time::Instant::now();
6104 tokio::time::timeout(
6105 Duration::from_secs(2),
6106 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6107 )
6108 .await
6109 .expect("a stop asked for while idle must wake the wait")
6110 .expect("the loop's own setup and teardown must not fail");
6111 asker.await.unwrap();
6112 assert!(
6113 began.elapsed() < opts.poll,
6114 "returned only after {:?}, so the stop waited on the sleep",
6115 began.elapsed()
6116 );
6117 }
6118
6119 #[tokio::test]
6120 async fn a_stopped_loop_leaves_no_status_file_claiming_it_is_running() {
6121 let dir = tempfile::tempdir().unwrap();
6122 let (opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6123 let stop = Stop::new();
6124 stop.stop();
6125
6126 tokio::time::timeout(
6127 Duration::from_secs(2),
6128 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6129 )
6130 .await
6131 .expect("a stopped loop must return")
6132 .expect("the loop's own setup and teardown must not fail");
6133
6134 assert!(
6135 home.is_dir(),
6136 "the loop did publish a status file, so its removal is the teardown and not an absence"
6137 );
6138 assert!(
6139 !status_file.exists(),
6140 "a stopped loop clears its status file"
6141 );
6142 assert!(
6143 read_status(&home).is_none(),
6144 "a reader must see no daemon at all, not a heartbeat that merely stopped"
6145 );
6146 }
6147
6148 #[tokio::test]
6149 async fn once_runs_startup_housekeeping_before_an_empty_queue_exits() {
6150 let dir = tempfile::tempdir().unwrap();
6151 let (mut opts, queue, status_file, home, worktrees) = idle_loop(dir.path());
6152 opts.once = true;
6153
6154 let mut settled = RunState::new(
6155 dir.path().join("repo"),
6156 "main".to_owned(),
6157 "abc1234".to_owned(),
6158 "fixture".to_owned(),
6159 Config::default(),
6160 );
6161 settled.status = RunStatus::Ready;
6162 let run_dir = home.join("runs").join(&settled.id);
6163 std::fs::create_dir_all(&run_dir).unwrap();
6164 std::fs::write(
6165 run_dir.join("run.json"),
6166 serde_json::to_string_pretty(&settled).unwrap(),
6167 )
6168 .unwrap();
6169 let questions = Questions::at(home.join("questions"));
6170 let mut question = ask::Question::new(
6171 settled.id.clone(),
6172 "review".to_owned(),
6173 "reviewer-1".to_owned(),
6174 "Continue?".to_owned(),
6175 String::new(),
6176 Vec::new(),
6177 );
6178 questions.put(&mut question).unwrap();
6179
6180 drive(&opts, &queue, &status_file, &home, &worktrees, &Stop::new())
6181 .await
6182 .unwrap();
6183
6184 assert_eq!(
6185 questions.get(&question.id).unwrap().status,
6186 ask::QuestionStatus::Abandoned,
6187 "an empty --once drain still performs startup question cleanup"
6188 );
6189 }
6190
6191 #[test]
6192 fn cache_check_due_fires_immediately_then_waits_out_its_own_interval() {
6193 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6194
6195 assert!(
6196 cache_check_due(None, t0, CACHE_CHECK_INTERVAL_SECS),
6197 "never checked before: due at once"
6198 );
6199
6200 let one_sec_later = t0 + jiff::SignedDuration::from_secs(1);
6201 assert!(
6202 !cache_check_due(Some(t0), one_sec_later, CACHE_CHECK_INTERVAL_SECS),
6203 "well inside the interval: not due yet"
6204 );
6205
6206 let at_the_edge = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64);
6207 assert!(
6208 !cache_check_due(Some(t0), at_the_edge, CACHE_CHECK_INTERVAL_SECS),
6209 "exactly at the edge: not yet due, same convention as `clean::due`"
6210 );
6211
6212 let past_it = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6213 assert!(
6214 cache_check_due(Some(t0), past_it, CACHE_CHECK_INTERVAL_SECS),
6215 "past the interval: due again"
6216 );
6217 }
6218
6219 fn cache_check_opts(dir: &Path, cache_dir: &Path, limit_bytes: u64) -> Opts {
6224 let config = dir.join("magi.toml");
6225 std::fs::write(
6231 &config,
6232 format!(
6233 "[disk]\nmin_free_bytes = 0\nauto_fold = false\ncache_limit_bytes = {limit_bytes}\n\n\
6234 [verify]\ngate = ['CARGO_TARGET_DIR={} cargo make check']\n",
6235 cache_dir.display()
6236 ),
6237 )
6238 .unwrap();
6239 Opts {
6240 config: Some(config),
6241 repo: dir.join("repo"),
6242 ..Opts::default()
6243 }
6244 }
6245
6246 #[tokio::test]
6247 async fn maybe_prune_cache_between_runs_reprunes_only_once_its_own_interval_elapses() {
6248 let dir = tempfile::tempdir().unwrap();
6249 let home = dir.path().join("home");
6250 let cache_dir = dir.path().join("cache");
6251 std::fs::create_dir_all(&cache_dir).unwrap();
6252 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6253 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6254
6255 let running = Stop::new();
6258 let mut last_checked = None;
6259 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6260 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &running, &mut last_checked, t0)
6261 .await;
6262 assert_eq!(
6263 crate::disk::dir_size(&cache_dir),
6264 0,
6265 "over the cap on the first check ever: pruned at once, no idle queue required"
6266 );
6267 assert_eq!(last_checked, Some(t0));
6268
6269 std::fs::write(cache_dir.join("b"), vec![0u8; 10]).unwrap();
6271 let too_soon = t0 + jiff::SignedDuration::from_secs(1);
6272 maybe_prune_cache_between_runs(
6273 &opts.repo,
6274 &opts,
6275 &home,
6276 &running,
6277 &mut last_checked,
6278 too_soon,
6279 )
6280 .await;
6281 assert_eq!(
6282 crate::disk::dir_size(&cache_dir),
6283 10,
6284 "too soon since the last check: left alone rather than rescanned every call"
6285 );
6286 assert_eq!(
6287 last_checked,
6288 Some(t0),
6289 "an idle check does not reset the clock"
6290 );
6291
6292 let due_again = t0 + jiff::SignedDuration::from_secs(CACHE_CHECK_INTERVAL_SECS as i64 + 1);
6294 maybe_prune_cache_between_runs(
6295 &opts.repo,
6296 &opts,
6297 &home,
6298 &running,
6299 &mut last_checked,
6300 due_again,
6301 )
6302 .await;
6303 assert_eq!(
6304 crate::disk::dir_size(&cache_dir),
6305 0,
6306 "due again: pruned back under the cap"
6307 );
6308 }
6309
6310 #[tokio::test]
6318 async fn a_stop_already_asked_for_skips_the_between_runs_cache_walk() {
6319 let dir = tempfile::tempdir().unwrap();
6320 let home = dir.path().join("home");
6321 let cache_dir = dir.path().join("cache");
6322 std::fs::create_dir_all(&cache_dir).unwrap();
6323 std::fs::write(cache_dir.join("a"), vec![0u8; 10]).unwrap();
6324 let opts = cache_check_opts(dir.path(), &cache_dir, 1);
6325
6326 let stop = Stop::new();
6327 stop.stop();
6328 assert!(
6329 !stop.finishing(),
6330 "no run is in flight at a between-runs boundary, so nothing else \
6331 would tell the operator this stop had not taken effect yet"
6332 );
6333
6334 let mut last_checked = None;
6335 let t0 = "2026-09-15T00:00:00Z".parse::<Timestamp>().unwrap();
6336 maybe_prune_cache_between_runs(&opts.repo, &opts, &home, &stop, &mut last_checked, t0)
6337 .await;
6338 assert_eq!(
6339 crate::disk::dir_size(&cache_dir),
6340 10,
6341 "over its cap, and due for the first check ever, but a stop outranks \
6342 it: the cap is a standing policy the next start measures again"
6343 );
6344 assert_eq!(
6345 last_checked, None,
6346 "a check that never happened must not claim the interval"
6347 );
6348 }
6349
6350 #[tokio::test]
6364 async fn cache_prune_reaches_a_queue_that_never_goes_idle() {
6365 let dir = tempfile::tempdir().unwrap();
6366 let cache_dir = dir.path().join("cache");
6367 std::fs::create_dir_all(&cache_dir).unwrap();
6368 std::fs::write(cache_dir.join("stale"), vec![0u8; 4096]).unwrap();
6369
6370 let mut opts = cache_check_opts(dir.path(), &cache_dir, 1);
6371 opts.poll = Duration::from_millis(20);
6372 opts.max_attempts = 1_000;
6373
6374 let queue = Queue::at(dir.path().join("queue"));
6375 let mut t = Task::new(
6376 "x".to_owned(),
6377 "x".to_owned(),
6378 opts.repo.clone(),
6379 Source::Human,
6380 );
6381 queue.put(&mut t).unwrap();
6382
6383 let home = dir.path().join("home");
6384 let worktrees = dir.path().join("wt");
6385 let status_file = home.join("daemon.json");
6386 let stop = Stop::new();
6387 let stopper = {
6388 let stop = stop.clone();
6389 tokio::spawn(async move {
6390 tokio::time::sleep(Duration::from_millis(400)).await;
6391 stop.stop();
6392 })
6393 };
6394
6395 tokio::time::timeout(
6396 Duration::from_secs(10),
6397 drive(&opts, &queue, &status_file, &home, &worktrees, &stop),
6398 )
6399 .await
6400 .expect("the loop must not hang on a queue that keeps producing failing work")
6401 .expect("the loop's own setup and teardown must not fail");
6402 stopper.await.unwrap();
6403
6404 let after = queue.get(&t.id).unwrap();
6405 assert!(
6406 after.attempts >= 2,
6407 "the harness must actually have retried more than once, or this is not \
6408 exercising a busy queue at all (got {} attempt(s))",
6409 after.attempts
6410 );
6411 assert!(
6412 after.status.runnable(),
6413 "still under its attempt budget: the queue never reached a natural idle \
6414 on its own, only the external stop ended the test"
6415 );
6416
6417 assert_eq!(
6418 crate::disk::dir_size(&cache_dir),
6419 0,
6420 "an oversized cache must not be left to grow unboundedly just because the \
6421 queue kept the loop busy the whole time"
6422 );
6423 }
6424
6425 #[test]
6426 fn task_question_reconciliation_keeps_references_and_retires_manual_releases() {
6427 let dir = tempfile::tempdir().unwrap();
6428 let queue = Queue::at(dir.path().join("queue"));
6429 let questions = Questions::at(dir.path().join("questions"));
6430 let mut task = task();
6431 queue.put(&mut task).unwrap();
6432
6433 let mut task_question = ask::Question::new(
6434 task.id.clone(),
6435 crate::conduct::NODE.to_owned(),
6436 "conduct".to_owned(),
6437 "Which backend?".to_owned(),
6438 String::new(),
6439 Vec::new(),
6440 );
6441 questions.put(&mut task_question).unwrap();
6442 task.block(vec![task_question.id.clone()], None);
6443 queue.put(&mut task).unwrap();
6444
6445 let mut run_question = ask::Question::new(
6446 "20260101-000000-run1".to_owned(),
6447 "review".to_owned(),
6448 "reviewer-1".to_owned(),
6449 "Run question".to_owned(),
6450 String::new(),
6451 Vec::new(),
6452 );
6453 questions.put(&mut run_question).unwrap();
6454
6455 let mut coincidental = ask::Question::new(
6460 task.id.clone(),
6461 "review".to_owned(),
6462 "reviewer-1".to_owned(),
6463 "Unrelated review question".to_owned(),
6464 String::new(),
6465 Vec::new(),
6466 );
6467 questions.put(&mut coincidental).unwrap();
6468
6469 reconcile_task_questions(&queue, &questions);
6470 assert!(questions.get(&task_question.id).unwrap().status.open());
6471 assert!(questions.get(&run_question.id).unwrap().status.open());
6472 assert!(questions.get(&coincidental.id).unwrap().status.open());
6473
6474 task.release();
6475 queue.put(&mut task).unwrap();
6476 reconcile_task_questions(&queue, &questions);
6477 assert_eq!(
6478 questions.get(&task_question.id).unwrap().status,
6479 ask::QuestionStatus::Abandoned
6480 );
6481 assert!(
6482 questions.get(&run_question.id).unwrap().status.open(),
6483 "run questions remain the run janitor's responsibility"
6484 );
6485 assert!(
6486 questions.get(&coincidental.id).unwrap().status.open(),
6487 "a non-conductor question must not be abandoned just because its \
6488 run id coincides with a task id"
6489 );
6490 }
6491
6492 #[test]
6493 fn a_freshly_started_running_task_is_never_stalled() {
6494 let dir = tempfile::tempdir().unwrap();
6495 let mut t = task();
6496 t.start("run-1".to_owned());
6497 assert!(!is_stalled(&t, dir.path(), Timestamp::now()));
6500 }
6501
6502 #[test]
6503 fn a_long_running_task_with_no_live_daemon_is_stalled() {
6504 let dir = tempfile::tempdir().unwrap();
6505 let mut t = task();
6506 t.start("run-1".to_owned());
6507 t.updated_at = Timestamp::now()
6508 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6509 assert!(is_stalled(&t, dir.path(), Timestamp::now()));
6510 assert_eq!(
6511 stalled_tasks(
6512 &Queue::at(dir.path().join("q")),
6513 dir.path(),
6514 Timestamp::now()
6515 )
6516 .len(),
6517 0,
6518 "the task was never written to this queue"
6519 );
6520 }
6521
6522 #[test]
6523 fn a_long_running_task_a_live_daemon_still_names_is_not_stalled() {
6524 let dir = tempfile::tempdir().unwrap();
6525 let mut t = task();
6526 t.id = "20260903-080340-0167".to_owned();
6527 t.start("20260903-080619-01c2".to_owned());
6528 t.updated_at = Timestamp::now()
6529 - jiff::SignedDuration::from_secs(STALLED_RUNNING.as_secs() as i64 + 60);
6530
6531 let mut status = Status::new();
6532 status.current = vec![Current {
6533 task: t.id.clone(),
6534 run: "20260903-080619-01c2".to_owned(),
6535 }];
6536 write_status_to(&dir.path().join("daemon.json"), &status).unwrap();
6537
6538 assert!(
6539 !is_stalled(&t, dir.path(), Timestamp::now()),
6540 "a live daemon's own heartbeat rules out stalled, however long the task has run"
6541 );
6542 }
6543
6544 fn backdate_task(queue: &Queue, id: &str, seconds_ago: i64) {
6548 let path = queue.path_of(id);
6549 let body = std::fs::read_to_string(&path).unwrap();
6550 let mut v: serde_json::Value = serde_json::from_str(&body).unwrap();
6551 let old = Timestamp::now() - jiff::SignedDuration::from_secs(seconds_ago);
6552 v["updated_at"] = serde_json::Value::String(old.to_string());
6553 std::fs::write(&path, serde_json::to_string_pretty(&v).unwrap()).unwrap();
6554 }
6555
6556 #[test]
6557 fn stalled_tasks_still_reaches_a_task_reclaim_could_not_claim_yet() {
6558 let dir = tempfile::tempdir().unwrap();
6571 let queue = Queue::at(dir.path().join("queue"));
6572 let home = dir.path().join("home");
6573
6574 let mut t = task();
6575 t.id = "20260101-000001-lock".to_owned();
6576 t.start("run-1".to_owned());
6577 queue.put(&mut t).unwrap();
6578 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
6579 std::fs::write(
6580 dir.path().join("queue").join(format!("{}.lock", t.id)),
6581 "not a pid",
6582 )
6583 .unwrap();
6584
6585 let now = Timestamp::now();
6586 assert!(
6587 reclaim_orphaned_running(&queue, 2).is_empty(),
6588 "the unparseable lock is still well within STALE_CLAIM, so the claim fails \
6589 and reclaim must leave the task alone"
6590 );
6591 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Running);
6592
6593 let stalled = stalled_tasks(&queue, &home, now);
6594 assert_eq!(
6595 stalled.len(),
6596 1,
6597 "reclaim's inability to claim it yet must not hide it from the conductor"
6598 );
6599 assert_eq!(stalled[0].id, t.id);
6600 }
6601
6602 #[test]
6603 fn ordinary_dead_daemon_task_is_shown_stalled_before_reclaim_and_can_be_requeued() {
6604 let dir = tempfile::tempdir().unwrap();
6605 crate::run::set_home(dir.path().join("run-home"));
6606 let queue = Queue::at(dir.path().join("queue"));
6607 let home = dir.path().join("home");
6608 let questions = Questions::at(dir.path().join("questions"));
6609
6610 let mut t = task();
6611 t.id = "20260101-000003-dead".to_owned();
6612 t.start("missing-run".to_owned());
6613 queue.put(&mut t).unwrap();
6614 backdate_task(&queue, &t.id, STALLED_RUNNING.as_secs() as i64 + 60);
6615
6616 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
6619 assert_eq!(
6620 stalled.iter().map(|task| &task.id).collect::<Vec<_>>(),
6621 [&t.id]
6622 );
6623 assert_eq!(reclaim_orphaned_running(&queue, 2), [t.id.clone()]);
6624 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
6625
6626 crate::conduct::apply(
6629 &queue,
6630 &questions,
6631 &crate::conduct::Verdict {
6632 decisions: vec![crate::conduct::Decision {
6633 id: t.id.clone(),
6634 recovery: Some(crate::conduct::Recovery::Requeue),
6635 ..crate::conduct::Decision::default()
6636 }],
6637 },
6638 )
6639 .unwrap();
6640 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
6641 }
6642
6643 #[test]
6644 fn stalled_tasks_reports_exactly_the_tasks_is_stalled_agrees_on() {
6645 let dir = tempfile::tempdir().unwrap();
6646 let queue = Queue::at(dir.path().join("queue"));
6647 let home = dir.path().join("home");
6648
6649 let mut fresh = task();
6650 fresh.id = "20260101-000001-aaaa".to_owned();
6651 fresh.start("run-1".to_owned());
6652 queue.put(&mut fresh).unwrap();
6653
6654 let mut old = task();
6655 old.id = "20260101-000002-bbbb".to_owned();
6656 old.start("run-2".to_owned());
6657 queue.put(&mut old).unwrap();
6658 backdate_task(&queue, &old.id, STALLED_RUNNING.as_secs() as i64 + 60);
6659
6660 let stalled = stalled_tasks(&queue, &home, Timestamp::now());
6661 assert_eq!(stalled.len(), 1);
6662 assert_eq!(stalled[0].id, old.id);
6663 }
6664
6665 #[test]
6666 fn queued_and_finished_task_views_partition_by_status() {
6667 let dir = tempfile::tempdir().unwrap();
6668 let queue = Queue::at(dir.path().join("queue"));
6669
6670 let mut queued = task();
6671 queued.id = "20260101-000001-aaaa".to_owned();
6672 queue.put(&mut queued).unwrap();
6673
6674 let mut failed = task();
6675 failed.id = "20260101-000002-bbbb".to_owned();
6676 failed.start("run-1".to_owned());
6677 failed.fail("gate red", 5);
6678 queue.put(&mut failed).unwrap();
6679
6680 let mut held = task();
6681 held.id = "20260101-000003-cccc".to_owned();
6682 held.hold_machine(None);
6683 queue.put(&mut held).unwrap();
6684
6685 let mut running = task();
6686 running.id = "20260101-000004-dddd".to_owned();
6687 running.start("run-2".to_owned());
6688 queue.put(&mut running).unwrap();
6689
6690 let queued_ids: Vec<String> = queued_tasks(&queue).into_iter().map(|t| t.id).collect();
6691 assert_eq!(queued_ids, [queued.id.clone()]);
6692
6693 let mut finished_ids: Vec<String> =
6694 finished_tasks(&queue).into_iter().map(|t| t.id).collect();
6695 finished_ids.sort_unstable();
6696 let mut want = vec![failed.id.clone(), held.id.clone()];
6697 want.sort_unstable();
6698 assert_eq!(finished_ids, want);
6699 }
6700
6701 #[test]
6702 fn resolve_blockers_clears_a_done_dependency_and_keeps_an_unresolved_one() {
6703 let dir = tempfile::tempdir().unwrap();
6704 let queue = Queue::at(dir.path().join("queue"));
6705 let questions = ask::Questions::at(dir.path().join("questions"));
6706
6707 let mut dep = task();
6708 dep.id = "20260101-000001-dep0".to_owned();
6709 dep.succeed();
6710 queue.put(&mut dep).unwrap();
6711
6712 let mut still_going = task();
6713 still_going.id = "20260101-000002-dep1".to_owned();
6714 queue.put(&mut still_going).unwrap();
6715
6716 let mut blocked = task();
6717 blocked.id = "20260101-000003-main".to_owned();
6718 blocked.block(
6719 vec![dep.id.clone(), still_going.id.clone()],
6720 Some("waits on both".to_owned()),
6721 );
6722 queue.put(&mut blocked).unwrap();
6723
6724 resolve_blockers(&queue, &questions);
6725
6726 let after = queue.get(&blocked.id).unwrap();
6727 assert_eq!(
6728 after.status,
6729 TaskStatus::Blocked,
6730 "one dependency is still outstanding"
6731 );
6732 assert_eq!(after.blocked_by, [still_going.id.clone()]);
6733 }
6734
6735 #[test]
6736 fn resolve_blockers_carries_an_answers_content_onto_the_task_and_unblocks_it() {
6737 let dir = tempfile::tempdir().unwrap();
6738 let queue = Queue::at(dir.path().join("queue"));
6739 let questions = ask::Questions::at(dir.path().join("questions"));
6740
6741 let mut q = crate::ask::Question::new(
6742 "20260101-000001-main".to_owned(),
6743 crate::conduct::NODE.to_owned(),
6744 "conduct".to_owned(),
6745 "Which backend?".to_owned(),
6746 String::new(),
6747 Vec::new(),
6748 );
6749 questions.put(&mut q).unwrap();
6750 q.answer(crate::ask::Answer::Text("SQLite".to_owned()))
6751 .unwrap();
6752 questions.put(&mut q).unwrap();
6753
6754 let mut blocked = task();
6755 blocked.id = "20260101-000001-main".to_owned();
6756 blocked.block(vec![q.id.clone()], Some("which backend?".to_owned()));
6757 queue.put(&mut blocked).unwrap();
6758
6759 resolve_blockers(&queue, &questions);
6760
6761 let after = queue.get(&blocked.id).unwrap();
6762 assert_eq!(
6763 after.status,
6764 TaskStatus::Queued,
6765 "the only blocker resolved"
6766 );
6767 assert_eq!(after.answers.len(), 1);
6768 assert_eq!(after.answers[0].question, "Which backend?");
6769 assert_eq!(after.answers[0].answer, "SQLite");
6770
6771 let instruction = instruction_for(&after);
6773 assert!(instruction.contains("Which backend?"));
6774 assert!(instruction.contains("SQLite"));
6775 }
6776
6777 #[test]
6778 fn resolve_blockers_restores_a_held_task_to_held_instead_of_queuing_it() {
6779 let dir = tempfile::tempdir().unwrap();
6785 let queue = Queue::at(dir.path().join("queue"));
6786 let questions = ask::Questions::at(dir.path().join("questions"));
6787
6788 let mut q = crate::ask::Question::new(
6789 "20260101-000001-main".to_owned(),
6790 crate::conduct::NODE.to_owned(),
6791 "conduct".to_owned(),
6792 "How should this be handled?".to_owned(),
6793 String::new(),
6794 Vec::new(),
6795 );
6796 questions.put(&mut q).unwrap();
6797 q.answer(crate::ask::Answer::Text(
6798 "leave it held, a human will look at it later".to_owned(),
6799 ))
6800 .unwrap();
6801 questions.put(&mut q).unwrap();
6802
6803 let mut held = task();
6804 held.id = "20260101-000001-main".to_owned();
6805 held.hold_machine(Some("out of attempts".to_owned()));
6806 held.block(vec![q.id.clone()], Some("what now?".to_owned()));
6807 queue.put(&mut held).unwrap();
6808
6809 resolve_blockers(&queue, &questions);
6810
6811 let after = queue.get(&held.id).unwrap();
6812 assert_eq!(after.status, TaskStatus::Held);
6813 assert_eq!(after.hold_reason.as_deref(), Some("out of attempts"));
6814 assert_eq!(
6815 after.answers[0].answer,
6816 "leave it held, a human will look at it later"
6817 );
6818 }
6819
6820 #[test]
6821 fn resolve_blockers_holds_a_task_whose_dependency_was_deleted() {
6822 let dir = tempfile::tempdir().unwrap();
6828 let queue = Queue::at(dir.path().join("queue"));
6829 let questions = ask::Questions::at(dir.path().join("questions"));
6830
6831 let mut still_going = task();
6832 still_going.id = "20260101-000002-dep1".to_owned();
6833 queue.put(&mut still_going).unwrap();
6834
6835 let mut blocked = task();
6836 blocked.id = "20260101-000003-main".to_owned();
6837 blocked.block(
6838 vec!["20260101-000001-gone".to_owned(), still_going.id.clone()],
6839 Some("waits on both".to_owned()),
6840 );
6841 queue.put(&mut blocked).unwrap();
6842
6843 resolve_blockers(&queue, &questions);
6844
6845 let after = queue.get(&blocked.id).unwrap();
6846 assert_eq!(
6847 after.status,
6848 TaskStatus::Held,
6849 "a missing dependency must not leave the task blocked forever"
6850 );
6851 assert_eq!(after.hold_source, Some(crate::queue::HoldSource::Machine));
6852 assert!(after.blocked_by.is_empty());
6853 let reason = after.hold_reason.as_deref().unwrap_or_default();
6854 assert!(
6855 reason.contains("20260101-000001-gone"),
6856 "the missing id must be named so an operator can tell what happened: {reason}"
6857 );
6858 assert!(
6859 reason.contains(&still_going.id),
6860 "the still-valid dependency must not silently vanish from the record: {reason}"
6861 );
6862 }
6863
6864 #[test]
6865 fn instruction_for_is_unchanged_without_any_answers() {
6866 let t = task();
6867 assert_eq!(instruction_for(&t), t.instruction);
6868 }
6869
6870 #[test]
6871 fn task_attachments_are_absolute_and_a_missing_file_is_an_error() {
6872 let dir = tempfile::tempdir().unwrap();
6873 let q = Queue::at(dir.path().join("queue"));
6874 let src = dir.path().join("shot.png");
6875 std::fs::write(&src, "x").unwrap();
6876 let mut t = task();
6877 q.attach(&mut t, &[src]).unwrap();
6878 let paths = task_attachments(&q, &t).unwrap();
6879 assert_eq!(paths.len(), 1);
6880 assert!(paths[0].is_absolute() && paths[0].is_file());
6881 std::fs::remove_file(&paths[0]).unwrap();
6882 let err = task_attachments(&q, &t).unwrap_err().to_string();
6883 assert!(err.contains("shot.png"), "{err}");
6884 }
6885
6886 #[test]
6887 fn resumed_instruction_is_unchanged_without_any_answers() {
6888 let t = task();
6889 assert_eq!(resumed_instruction(&t.instruction, &t), t.instruction);
6890 }
6891
6892 #[test]
6893 fn resumed_instruction_carries_a_new_answer_onto_the_old_run() {
6894 let mut t = task();
6895 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
6896 let old = t.instruction.clone();
6900
6901 let refreshed = resumed_instruction(&old, &t);
6902 assert!(refreshed.starts_with(&old), "the original text is kept");
6903 assert!(refreshed.contains("Which backend?"));
6904 assert!(refreshed.contains("SQLite"));
6905 }
6906
6907 #[test]
6908 fn resumed_instruction_keeps_an_original_answers_heading() {
6909 let mut t = task();
6910 t.instruction = "Context\n\n# Operator answers\n\nThis is part of the task.".to_owned();
6911 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
6912
6913 let refreshed = resumed_instruction(&t.instruction, &t);
6914
6915 assert!(
6916 refreshed.starts_with(&t.instruction),
6917 "an answers heading in the original instruction is not the appended block"
6918 );
6919 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 2);
6920 assert!(refreshed.contains("Which backend?"));
6921 assert!(refreshed.contains("SQLite"));
6922
6923 let repeated = resumed_instruction(&refreshed, &t);
6924 assert_eq!(
6925 repeated, refreshed,
6926 "only the final appended block is refreshed"
6927 );
6928 }
6929
6930 #[test]
6931 fn resumed_instruction_does_not_duplicate_across_repeated_resumes() {
6932 let mut t = task();
6933 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
6934
6935 let once = resumed_instruction(&t.instruction, &t);
6939 let twice = resumed_instruction(&once, &t);
6940 assert_eq!(once, twice);
6941 assert_eq!(once.matches("Which backend?").count(), 1);
6942
6943 t.record_answer("Which cache?".to_owned(), "Redis".to_owned());
6945 let refreshed = resumed_instruction(&once, &t);
6946 assert_eq!(refreshed.matches(ANSWERS_HEADER).count(), 1);
6947 assert!(refreshed.contains("Which backend?"));
6948 assert!(refreshed.contains("Which cache?"));
6949 }
6950
6951 #[test]
6952 fn prepare_instruction_covers_all_three_starters() {
6953 let mut t = task();
6954 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
6955
6956 assert_eq!(
6959 prepare_instruction(&Starter::Start, None, &t),
6960 Some(instruction_for(&t))
6961 );
6962
6963 let old = t.instruction.clone();
6966 assert_eq!(
6967 prepare_instruction(&Starter::Resume("some-run".to_owned()), Some(&old), &t),
6968 Some(resumed_instruction(&old, &t))
6969 );
6970
6971 assert_eq!(
6975 prepare_instruction(&Starter::Review("magi/eba2/A".to_owned()), Some(&old), &t),
6976 None
6977 );
6978 }
6979
6980 #[test]
6981 fn choose_starter_prefers_review_over_resume_when_the_branch_survived() {
6982 assert_eq!(
6983 choose_starter(Some("magi/eba2/A"), true, Some("some-run")),
6984 Starter::Review("magi/eba2/A".to_owned())
6985 );
6986 }
6987
6988 #[test]
6989 fn choose_starter_falls_back_to_start_when_the_review_branch_is_gone() {
6990 assert_eq!(
6991 choose_starter(Some("magi/eba2/A"), false, Some("some-run")),
6992 Starter::Start,
6993 "a vanished review branch must not fall back to resuming the old run either"
6994 );
6995 }
6996
6997 #[test]
6998 fn choose_starter_resumes_or_starts_when_there_is_no_review_choice_at_all() {
6999 assert_eq!(
7000 choose_starter(None, false, Some("some-run")),
7001 Starter::Resume("some-run".to_owned())
7002 );
7003 assert_eq!(choose_starter(None, false, None), Starter::Start);
7004 }
7005
7006 #[test]
7007 fn an_explicit_release_forces_a_fresh_competition_even_with_a_resumable_run() {
7008 let mut released = task();
7009 released.start("stalled-run".to_owned());
7010 released.requeue();
7011 let unfinished = (!released.fresh_start)
7012 .then(|| Some("stalled-run".to_owned()))
7013 .flatten();
7014 assert_eq!(
7015 choose_starter(None, false, unfinished.as_deref()),
7016 Starter::Start,
7017 "release keeps run history but must not resume it"
7018 );
7019 assert_eq!(released.runs, ["stalled-run"]);
7020 }
7021
7022 #[test]
7023 fn an_ordinary_release_keeps_a_resumable_run_available() {
7024 let mut released = task();
7025 released.start("stalled-run".to_owned());
7026 released.release();
7027 let unfinished = (!released.fresh_start)
7028 .then(|| Some("stalled-run".to_owned()))
7029 .flatten();
7030 assert_eq!(
7031 choose_starter(None, false, unfinished.as_deref()),
7032 Starter::Resume("stalled-run".to_owned()),
7033 "manual release must preserve the normal resume path"
7034 );
7035 }
7036
7037 #[test]
7038 fn a_blocked_run_that_spent_every_review_round_has_exhausted_its_budget() {
7039 let mut state = run_state(RunStatus::Blocked);
7040 state.config.graph.review_rounds = 3;
7041 state.reviews = vec![review_round(1), review_round(2), review_round(3)];
7042 assert!(exhausted_review_budget(&state));
7043
7044 state.reviews.pop();
7046 assert!(!exhausted_review_budget(&state));
7047
7048 let mut stalled = run_state(RunStatus::Stalled);
7051 stalled.config.graph.review_rounds = 1;
7052 stalled.reviews = vec![review_round(1)];
7053 assert!(!exhausted_review_budget(&stalled));
7054 }
7055
7056 fn review_round(round: usize) -> crate::run::ReviewRound {
7057 crate::run::ReviewRound {
7058 round,
7059 head: "deadbeef".to_owned(),
7060 verified_head: None,
7061 verified_at: None,
7062 reviews: Vec::new(),
7063 e2e: Vec::new(),
7064 verify_retried: false,
7065 e2e_deferred: false,
7066 e2e_defer_reason: None,
7067 fix: None,
7068 blocking: 0,
7069 answered: 1,
7070 expected: 1,
7071 clean: false,
7072 progressed: true,
7073 vote_split: false,
7074 reconsideration: Vec::new(),
7075 verdict: None,
7076 }
7077 }
7078
7079 #[test]
7080 fn awaiting_resume_is_a_failed_task_whose_last_run_parked_and_can_be_resumed() {
7081 let mut run = RunState::new(
7082 PathBuf::from("/repo"),
7083 "main".to_owned(),
7084 "abc1234def".to_owned(),
7085 "add retries".to_owned(),
7086 Config::default(),
7087 );
7088 run.status = RunStatus::Judging;
7089 run.parked = true;
7090 let mut task = Task::new(
7091 "add retries".to_owned(),
7092 "add retries".to_owned(),
7093 PathBuf::from("/repo"),
7094 crate::queue::Source::Human,
7095 );
7096 task.status = TaskStatus::Failed;
7097 task.runs = vec![run.id.clone()];
7098 let with = |t: &Task, r: &RunState| awaiting_resume_with(t, |_| Ok(r.clone()));
7099 assert!(with(&task, &run), "parked after judging is the case");
7100
7101 let mut not_parked = run.clone();
7102 not_parked.parked = false;
7103 not_parked.status = RunStatus::Stalled;
7104 assert!(!with(&task, ¬_parked), "a stall is the conductor's");
7105
7106 let mut fresh = task.clone();
7107 fresh.fresh_start = true;
7108 assert!(!with(&fresh, &run), "a requeue asked for a new competition");
7109
7110 let mut review = task.clone();
7111 review.review_branch = Some("magi/x/A".to_owned());
7112 assert!(!with(&review, &run), "review is ranked before resume");
7113
7114 let mut held = task.clone();
7115 held.status = TaskStatus::Held;
7116 assert!(!with(&held, &run), "a hold stays visible to the conductor");
7117
7118 let mut released = run.clone();
7119 released.released_to = Some("20260901-000000-new1".to_owned());
7120 assert!(!with(&task, &released), "nothing left to resume into");
7121
7122 assert!(!awaiting_resume_with(&task, |_| anyhow::bail!(
7123 "unreadable"
7124 )));
7125 }
7126
7127 #[test]
7128 fn unfinished_run_never_offers_a_run_whose_worktree_was_released() {
7129 let mut released = RunState::new(
7130 PathBuf::from("/repo"),
7131 "main".to_owned(),
7132 "abc1234def".to_owned(),
7133 "add retries".to_owned(),
7134 Config::default(),
7135 );
7136 released.status = RunStatus::Blocked;
7137 assert_eq!(
7138 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7139 Some(released.id.clone())
7140 );
7141 released.released_to = Some("20260901-000000-new1".to_owned());
7142 assert_eq!(
7143 unfinished_run_with(&[released.id.clone()], "t", |_| Ok(released.clone())),
7144 None,
7145 "there is nothing left to resume it into"
7146 );
7147 }
7148
7149 #[test]
7150 fn unfinished_run_skips_a_round_exhausted_blocked_run_so_requeue_means_a_fresh_competition() {
7151 let mut exhausted = RunState::new(
7160 PathBuf::from("/repo"),
7161 "main".to_owned(),
7162 "abc1234def".to_owned(),
7163 "add retries".to_owned(),
7164 Config::default(),
7165 );
7166 exhausted.status = RunStatus::Blocked;
7167 exhausted.config.graph.review_rounds = 1;
7168 exhausted.reviews = vec![review_round(1)];
7169
7170 assert_eq!(
7171 unfinished_run_with(&[exhausted.id.clone()], "t", |_| Ok(exhausted.clone())),
7172 None,
7173 "an exhausted `Blocked` run must not be offered as resumable"
7174 );
7175
7176 let mut has_budget_left = RunState::new(
7179 PathBuf::from("/repo"),
7180 "main".to_owned(),
7181 "abc1234def".to_owned(),
7182 "add retries".to_owned(),
7183 Config::default(),
7184 );
7185 has_budget_left.status = RunStatus::Blocked;
7186 has_budget_left.config.graph.review_rounds = 3;
7187 has_budget_left.reviews = vec![review_round(1)];
7188
7189 assert_eq!(
7190 unfinished_run_with(&[has_budget_left.id.clone()], "t", |_| {
7191 Ok(has_budget_left.clone())
7192 }),
7193 Some(has_budget_left.id.clone())
7194 );
7195 }
7196
7197 #[test]
7198 fn unfinished_run_never_falls_back_to_an_older_resumable_run() {
7199 let mut older_stalled = RunState::new(
7207 PathBuf::from("/repo"),
7208 "main".to_owned(),
7209 "abc1234def".to_owned(),
7210 "add retries".to_owned(),
7211 Config::default(),
7212 );
7213 older_stalled.status = RunStatus::Stalled;
7214
7215 let mut newest_exhausted = RunState::new(
7216 PathBuf::from("/repo"),
7217 "main".to_owned(),
7218 "abc1234def".to_owned(),
7219 "add retries".to_owned(),
7220 Config::default(),
7221 );
7222 newest_exhausted.status = RunStatus::Blocked;
7223 newest_exhausted.config.graph.review_rounds = 1;
7224 newest_exhausted.reviews = vec![review_round(1)];
7225
7226 assert_eq!(
7227 unfinished_run_with(
7228 &[older_stalled.id.clone(), newest_exhausted.id.clone()],
7229 "t",
7230 |_| Ok(newest_exhausted.clone())
7231 ),
7232 None,
7233 "the newest run is exhausted, so nothing here is worth resuming - \
7234 least of all the older, already-superseded run"
7235 );
7236 }
7237
7238 #[test]
7239 fn unfinished_run_warns_and_skips_a_run_it_cannot_read() {
7240 assert_eq!(
7241 unfinished_run_with(&["20260101-000000-gone".to_owned()], "t", |_| {
7242 Err(anyhow::anyhow!("fixture is absent"))
7243 }),
7244 None
7245 );
7246 }
7247
7248 fn action_question(run: &str, action: ask::ChoiceAction) -> ask::Question {
7249 let mut q = ask::Question::new(
7250 run.to_owned(),
7251 "implement".to_owned(),
7252 "impl-A".to_owned(),
7253 "continue?".to_owned(),
7254 String::new(),
7255 vec!["resume で続行する".to_owned(), "other".to_owned()],
7256 );
7257 q.actions.insert("resume で続行する".to_owned(), action);
7258 q.answer(ask::Answer::Choice("resume で続行する".to_owned()))
7259 .unwrap();
7260 q
7261 }
7262
7263 fn held_task_with(run: &str) -> Task {
7264 let mut t = task();
7265 t.runs = vec![run.to_owned()];
7266 t.hold_machine(Some("waiting for magi resume to be executed".to_owned()));
7267 t
7268 }
7269
7270 fn resume_action(run: &str) -> ask::ChoiceAction {
7271 ask::ChoiceAction::Resume { run: run.into() }
7272 }
7273
7274 #[test]
7275 fn decide_action_resumes_only_the_latest_resumable_run() {
7276 let t = held_task_with("r1");
7277 let q = action_question("r1", resume_action("r1"));
7278 let load = |s: RunState| move |_: &str| Ok(s);
7279 assert_eq!(
7280 decide_action(&t, &q, load(run_state(RunStatus::Blocked))),
7281 ActionDecision::Resume("r1".into())
7282 );
7283 let q_other = action_question("r1", resume_action("r0"));
7285 assert!(matches!(
7286 decide_action(&t, &q_other, load(run_state(RunStatus::Blocked))),
7287 ActionDecision::Refuse(_)
7288 ));
7289 assert!(matches!(
7291 decide_action(&t, &q, load(run_state(RunStatus::Ready))),
7292 ActionDecision::Refuse(_)
7293 ));
7294 let mut released = run_state(RunStatus::Blocked);
7296 released.released_to = Some("elsewhere".into());
7297 assert!(matches!(
7298 decide_action(&t, &q, load(released)),
7299 ActionDecision::Refuse(_)
7300 ));
7301 assert!(matches!(
7303 decide_action(&t, &q, |_: &str| bail!("gone")),
7304 ActionDecision::Refuse(_)
7305 ));
7306 }
7307
7308 #[test]
7309 fn decide_action_ignores_a_question_about_an_earlier_run() {
7310 let mut t = held_task_with("r1");
7311 t.runs.push("r2".to_owned());
7312 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7313 assert_eq!(
7314 decide_action(&t, &action_question("r1", ask::ChoiceAction::Done), never),
7315 ActionDecision::Stale
7316 );
7317 }
7318
7319 #[test]
7320 fn decide_action_maps_requeue_and_done_and_never_acts_twice() {
7321 let mut t = held_task_with("r1");
7322 let never = |_: &str| -> Result<RunState> { bail!("not read") };
7323 assert_eq!(
7324 decide_action(
7325 &t,
7326 &action_question("r1", ask::ChoiceAction::Requeue),
7327 never
7328 ),
7329 ActionDecision::Requeue
7330 );
7331 let done_q = action_question("r1", ask::ChoiceAction::Done);
7332 assert_eq!(decide_action(&t, &done_q, never), ActionDecision::Done);
7333 t.mark_action_applied(&done_q.id);
7334 assert_eq!(decide_action(&t, &done_q, never), ActionDecision::Skip);
7335
7336 let mut plain = action_question("r1", ask::ChoiceAction::Done);
7338 plain.actions.clear();
7339 assert_eq!(
7340 decide_action(&held_task_with("r1"), &plain, never),
7341 ActionDecision::Skip
7342 );
7343 let mut running = held_task_with("r1");
7345 running.status = TaskStatus::Running;
7346 assert_eq!(
7347 decide_action(
7348 &running,
7349 &action_question("r1", ask::ChoiceAction::Done),
7350 never
7351 ),
7352 ActionDecision::Skip
7353 );
7354 }
7355
7356 #[test]
7357 fn apply_choice_actions_releases_the_task_pinned_to_its_run_and_only_once() {
7358 let dir = tempfile::tempdir().unwrap();
7359 let queue = Queue::at(dir.path().join("queue"));
7360 let questions = Questions::at(dir.path().join("questions"));
7361 let home = dir.path().join("home");
7362 let mut state = run_state(RunStatus::Blocked);
7363 state.id = "20260101-000000-act1".to_owned();
7364 state.save_under(&home).unwrap();
7365
7366 let mut t = held_task_with(&state.id);
7367 queue.put(&mut t).unwrap();
7368 let mut q = action_question(&state.id, resume_action(&state.id));
7369 questions.put(&mut q).unwrap();
7370
7371 apply_choice_actions(&queue, &questions, &home);
7372 let after = queue.get(&t.id).unwrap();
7373 assert_eq!(after.status, TaskStatus::Queued);
7374 assert!(!after.fresh_start);
7375 assert!(after.action_applied(&q.id));
7376 let pin = after.resume_override.clone().unwrap();
7377 assert_eq!(pin.pinned_run.as_deref(), Some(state.id.as_str()));
7378 assert!(pin.forced);
7379
7380 let mut again = queue.get(&t.id).unwrap();
7382 again.hold_machine(Some("later".into()));
7383 queue.put(&mut again).unwrap();
7384 apply_choice_actions(&queue, &questions, &home);
7385 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7386 }
7387
7388 #[test]
7389 fn an_answer_the_waiter_already_delivered_is_not_acted_on_again() {
7390 let dir = tempfile::tempdir().unwrap();
7391 let queue = Queue::at(dir.path().join("queue"));
7392 let questions = Questions::at(dir.path().join("questions"));
7393 let home = dir.path().join("home");
7394 let mut state = run_state(RunStatus::Blocked);
7395 state.id = "20260101-000000-act2".to_owned();
7396 state.save_under(&home).unwrap();
7397
7398 let mut t = held_task_with(&state.id);
7399 queue.put(&mut t).unwrap();
7400 let mut q = action_question(&state.id, resume_action(&state.id));
7401 q.answer_delivered = true;
7402 questions.put(&mut q).unwrap();
7403
7404 apply_choice_actions(&queue, &questions, &home);
7405 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Held);
7406
7407 let mut q2 = action_question(&state.id, resume_action(&state.id));
7409 questions.put(&mut q2).unwrap();
7410 apply_choice_actions(&queue, &questions, &home);
7411 assert_eq!(queue.get(&t.id).unwrap().status, TaskStatus::Queued);
7412 assert!(questions.get(&q2.id).unwrap().answer_delivered);
7413 }
7414}