1use std::collections::HashMap;
34use std::path::{Path, PathBuf};
35
36use anyhow::{Context, Result, bail};
37use jiff::Timestamp;
38use serde::{Deserialize, Serialize};
39
40use crate::ask::Questions;
41
42pub const SCHEMA: u32 = 9;
97
98#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
100#[serde(rename_all = "lowercase")]
101pub enum HoldSource {
102 Manual,
104 Machine,
106}
107
108impl HoldSource {
109 pub fn label(self) -> &'static str {
111 match self {
112 Self::Manual => "manual",
113 Self::Machine => "machine",
114 }
115 }
116}
117
118#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
121#[serde(tag = "kind", rename_all = "lowercase")]
122pub enum Source {
123 Human,
125 Agent {
128 run: String,
130 node: String,
132 },
133 Issue {
135 number: u64,
137 repo: String,
139 },
140}
141
142impl Source {
143 pub fn label(&self) -> String {
145 match self {
146 Self::Human => "human".to_owned(),
147 Self::Agent { run, node } => format!("{node}@{}", short(run)),
148 Self::Issue { number, .. } => format!("issue #{number}"),
149 }
150 }
151}
152
153#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
155#[serde(rename_all = "lowercase")]
156pub enum TaskStatus {
157 Queued,
159 Running,
161 Done,
163 Failed,
165 Held,
167 Blocked,
171}
172
173impl TaskStatus {
174 pub fn runnable(self) -> bool {
176 matches!(self, Self::Queued | Self::Failed)
177 }
178
179 pub fn as_str(self) -> &'static str {
181 match self {
182 Self::Queued => "queued",
183 Self::Running => "running",
184 Self::Done => "done",
185 Self::Failed => "failed",
186 Self::Held => "held",
187 Self::Blocked => "blocked",
188 }
189 }
190}
191
192#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
195pub struct TaskCounts {
196 pub queued: usize,
198 pub running: usize,
200 pub done: usize,
202 pub failed: usize,
204 pub held: usize,
206 pub blocked: usize,
208}
209
210impl TaskCounts {
211 pub fn of(tasks: &[Task]) -> Self {
215 let mut counts = Self::default();
216 for t in tasks {
217 match t.status {
218 TaskStatus::Queued => counts.queued += 1,
219 TaskStatus::Running => counts.running += 1,
220 TaskStatus::Done => counts.done += 1,
221 TaskStatus::Failed => counts.failed += 1,
222 TaskStatus::Held => counts.held += 1,
223 TaskStatus::Blocked => counts.blocked += 1,
224 }
225 }
226 counts
227 }
228}
229
230#[derive(Debug, Clone, Serialize, Deserialize)]
232#[serde(deny_unknown_fields)]
233pub struct Task {
234 pub schema: u32,
236 pub id: String,
238 pub title: String,
240 pub instruction: String,
242 pub repo: PathBuf,
244 pub source: Source,
246 #[serde(default)]
248 pub priority: i32,
249 #[serde(default)]
260 pub solo: bool,
261 pub status: TaskStatus,
263 #[serde(default)]
265 pub attempts: usize,
266 #[serde(default)]
268 pub runs: Vec<String>,
269 #[serde(default)]
271 pub last_error: Option<String>,
272 #[serde(default)]
285 pub hold_reason: Option<String>,
286 #[serde(default)]
289 pub hold_source: Option<HoldSource>,
290 #[serde(default)]
302 pub diagnostic: Option<String>,
303 #[serde(default)]
313 pub blocked_by: Vec<String>,
314 #[serde(default)]
317 pub block_reason: Option<String>,
318 #[serde(default)]
331 pub blocked_from: Option<TaskStatus>,
332 #[serde(default)]
343 pub answers: Vec<AnsweredQuestion>,
344 #[serde(default)]
350 pub triage_applied: Vec<String>,
351 #[serde(default)]
357 pub actions_applied: Vec<String>,
358 #[serde(default)]
363 pub resume_override: Option<OperatorResume>,
364 #[serde(default)]
376 pub review_branch: Option<String>,
377 #[serde(default)]
380 pub fresh_start: bool,
381 #[serde(default)]
392 pub interrupt: bool,
393 #[serde(default)]
416 pub urgent: bool,
417 #[serde(default)]
423 pub attachments: Vec<String>,
424 #[serde(default)]
428 pub followup: Option<FollowUp>,
429 #[serde(default)]
434 pub held_at: Option<Timestamp>,
435 pub created_at: Timestamp,
437 pub updated_at: Timestamp,
439}
440
441#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
443pub struct FollowUp {
444 pub run: String,
446 #[serde(default)]
448 pub origin_task: Option<String>,
449 pub pr: String,
451 pub findings: Vec<String>,
453 pub generation: u32,
456}
457
458#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
461pub struct OperatorResume {
462 pub question_id: String,
464 pub at: Timestamp,
466 #[serde(default)]
469 pub conductor_rehold: Option<String>,
470 #[serde(default)]
473 pub forced: bool,
474 #[serde(default)]
481 pub pinned_run: Option<String>,
482}
483
484#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
487pub struct AnsweredQuestion {
488 pub question: String,
490 pub answer: String,
492}
493
494impl Task {
495 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
497 let now = Timestamp::now();
498 Self {
499 schema: SCHEMA,
500 id: new_id(),
501 title,
502 instruction,
503 repo,
504 source,
505 priority: 0,
506 solo: false,
507 status: TaskStatus::Queued,
508 attempts: 0,
509 runs: Vec::new(),
510 last_error: None,
511 hold_reason: None,
512 hold_source: None,
513 diagnostic: None,
514 blocked_by: Vec::new(),
515 block_reason: None,
516 blocked_from: None,
517 answers: Vec::new(),
518 triage_applied: Vec::new(),
519 actions_applied: Vec::new(),
520 resume_override: None,
521 review_branch: None,
522 fresh_start: false,
523 interrupt: false,
524 urgent: false,
525 attachments: Vec::new(),
526 followup: None,
527 held_at: None,
528 created_at: now,
529 updated_at: now,
530 }
531 }
532
533 pub fn short(&self) -> &str {
535 short(&self.id)
536 }
537
538 pub fn mark_triage_applied(&mut self, question_id: &str) {
540 if !self.triage_applied(question_id) {
541 self.triage_applied.push(question_id.to_owned());
542 }
543 }
544
545 pub fn action_applied(&self, question_id: &str) -> bool {
548 self.actions_applied.iter().any(|id| id == question_id)
549 }
550
551 pub fn mark_action_applied(&mut self, question_id: &str) {
553 if !self.action_applied(question_id) {
554 self.actions_applied.push(question_id.to_owned());
555 }
556 }
557
558 pub fn triage_applied(&self, question_id: &str) -> bool {
560 self.triage_applied.iter().any(|id| id == question_id)
561 }
562
563 pub fn start(&mut self, run: String) {
574 self.status = TaskStatus::Running;
575 self.attempts += 1;
576 self.runs.push(run);
577 self.last_error = None;
578 self.fresh_start = false;
579 self.interrupt = false;
580 self.resume_override = None;
582 }
583
584 pub fn link_run(&mut self, run: &str) -> bool {
589 if self.runs.iter().any(|r| r == run) {
590 return false;
591 }
592 self.runs.push(run.to_owned());
593 true
594 }
595
596 pub fn succeed(&mut self) {
606 self.status = TaskStatus::Done;
607 self.resume_override = None;
608 self.last_error = None;
609 self.hold_reason = None;
610 self.hold_source = None;
611 self.diagnostic = None;
612 self.blocked_by.clear();
613 self.block_reason = None;
614 self.blocked_from = None;
615 }
616
617 pub fn already_landed(&mut self, note: impl Into<String>) {
625 self.succeed();
626 self.attempts = self.attempts.saturating_sub(1);
627 self.last_error = Some(note.into());
628 }
629
630 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
660 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
661 return &[];
662 }
663 &self.runs[..self.runs.len() - 1]
664 }
665
666 pub fn successor_of(&self, run: &str) -> Option<&String> {
673 let pos = self.runs.iter().position(|r| r == run)?;
674 self.runs.get(pos + 1)
675 }
676
677 pub fn earlier_attempts(&self) -> &[String] {
684 &self.runs
685 }
686
687 fn note_held(&mut self) {
690 if self.status != TaskStatus::Held || self.held_at.is_none() {
691 self.held_at = Some(Timestamp::now());
692 }
693 }
694
695 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
711 let why = why.into();
712 self.diagnostic = None;
713 self.status = if self.attempts >= max_attempts {
714 self.note_held();
715 self.hold_source = Some(HoldSource::Machine);
716 self.hold_reason = Some(why.clone());
717 TaskStatus::Held
718 } else {
719 TaskStatus::Failed
720 };
721 self.last_error = Some(why);
722 }
723
724 pub fn stall(&mut self, why: impl Into<String>) {
734 self.last_error = Some(why.into());
735 self.diagnostic = None;
736 self.attempts = self.attempts.saturating_sub(1);
737 self.status = TaskStatus::Failed;
738 }
739
740 pub fn operator_held(&self) -> bool {
746 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
747 }
748
749 pub fn hold_manual(&mut self, reason: Option<String>) {
759 self.note_held();
760 self.status = TaskStatus::Held;
761 if reason.is_some() {
762 self.hold_reason = reason;
763 }
764 self.hold_source = Some(HoldSource::Manual);
765 self.blocked_by.clear();
766 self.block_reason = None;
767 self.blocked_from = None;
768 }
769
770 pub fn hold_machine(&mut self, reason: Option<String>) {
775 self.note_held();
776 self.status = TaskStatus::Held;
777 if reason.is_some() {
778 self.hold_reason = reason;
779 }
780 self.hold_source = Some(HoldSource::Machine);
781 self.blocked_by.clear();
782 self.block_reason = None;
783 self.blocked_from = None;
784 }
785
786 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
796 if self.status != TaskStatus::Blocked {
797 self.blocked_from = Some(self.status);
798 }
799 self.status = TaskStatus::Blocked;
800 self.blocked_by = blocked_by;
801 self.block_reason = reason;
802 }
803
804 pub fn unblock(&mut self, resolved_id: &str) {
827 if self.status != TaskStatus::Blocked {
828 return;
829 }
830 self.blocked_by.retain(|id| id != resolved_id);
831 if self.blocked_by.is_empty() {
832 self.status = match self.blocked_from {
833 Some(TaskStatus::Running) => TaskStatus::Queued,
834 Some(other) => other,
835 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
836 TaskStatus::Held
837 }
838 None => TaskStatus::Queued,
839 };
840 self.block_reason = None;
841 self.blocked_from = None;
842 }
843 }
844
845 pub fn record_answer(&mut self, question: String, answer: String) {
850 self.answers.push(AnsweredQuestion { question, answer });
851 }
852
853 pub fn request_review(&mut self, branch: String) {
857 self.release();
858 self.review_branch = Some(branch);
859 }
860
861 pub fn requeue(&mut self) {
864 self.release();
865 self.review_branch = None;
867 self.fresh_start = true;
868 }
869
870 pub fn hold_for_handover(&mut self, branch: Option<String>, reason: String) {
879 if branch.is_some() {
880 self.review_branch = branch;
881 }
882 self.hold_machine(Some(reason));
883 }
884
885 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
894 if self.status == TaskStatus::Running {
895 bail!(
896 "task {} is running; its priority cannot be changed until \
897 this attempt finishes",
898 self.short()
899 );
900 }
901 self.priority = priority;
902 Ok(())
903 }
904
905 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
920 if interrupt && !self.status.runnable() {
921 bail!(
922 "task {} is {}; only a queued or failed task can be marked \
923 to interrupt",
924 self.short(),
925 self.status.as_str()
926 );
927 }
928 self.interrupt = interrupt;
929 Ok(())
930 }
931
932 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
944 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
945 bail!(
946 "task {} is {}; only a queued or held task's instruction can \
947 be edited",
948 self.short(),
949 self.status.as_str()
950 );
951 }
952 self.title = title;
953 self.instruction = instruction;
954 Ok(())
955 }
956
957 pub fn handed_off(&mut self, why: impl Into<String>) {
977 let why = why.into();
978 self.diagnostic = None;
979 self.note_held();
980 self.status = TaskStatus::Held;
981 self.hold_source = Some(HoldSource::Machine);
982 self.hold_reason = Some(why.clone());
983 self.last_error = Some(why);
984 }
985
986 pub fn release(&mut self) {
990 let refused_handover = self.status == TaskStatus::Held
991 && self.hold_source == Some(HoldSource::Machine)
992 && self.review_branch.is_some();
993 self.status = TaskStatus::Queued;
994 self.held_at = None;
995 self.attempts = 0;
996 self.last_error = None;
997 self.hold_reason = None;
1000 self.hold_source = None;
1001 self.diagnostic = None;
1002 self.blocked_by.clear();
1007 self.block_reason = None;
1008 self.blocked_from = None;
1009 if !refused_handover {
1013 self.review_branch = None;
1014 }
1015 self.fresh_start = false;
1016 }
1017}
1018
1019fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
1023 use std::io::ErrorKind;
1024 let (stem, ext) = match name.rfind('.') {
1025 Some(i) if i > 0 => (&name[..i], &name[i..]),
1026 _ => (name, ""),
1027 };
1028 for n in 1u32.. {
1029 let candidate = if n == 1 {
1030 name.to_owned()
1031 } else {
1032 let suffix = format!("-{n}");
1033 let room = 64usize.saturating_sub(suffix.len() + ext.len());
1034 let stem: String = stem.chars().take(room).collect();
1035 format!("{stem}{suffix}{ext}")
1036 };
1037 if !crate::ask::valid_asset_name(&candidate) {
1038 bail!("no valid attachment name is left for `{name}`");
1039 }
1040 let path = dir.join(&candidate);
1041 match std::fs::OpenOptions::new()
1042 .write(true)
1043 .create_new(true)
1044 .open(&path)
1045 {
1046 Ok(mut out) => {
1047 let copied = std::fs::File::open(src)
1048 .and_then(|mut input| std::io::copy(&mut input, &mut out));
1049 if let Err(e) = copied {
1050 drop(out);
1051 let _ = std::fs::remove_file(&path);
1052 return Err(e).with_context(|| format!("copy {}", src.display()));
1053 }
1054 return Ok((candidate, path));
1055 }
1056 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
1057 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
1058 }
1059 }
1060 unreachable!("the counter never runs out")
1061}
1062
1063const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
1066
1067struct TaskLock {
1070 path: PathBuf,
1071 token: String,
1072}
1073
1074fn owner_token() -> String {
1077 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1078 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1079 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy() ^ n.rotate_left(32));
1080 format!("{}-{:016x}", std::process::id(), r.next_u64())
1081}
1082
1083fn break_marker(path: &Path, token: &str) -> PathBuf {
1086 let mut h: u64 = 0xcbf29ce484222325;
1087 for b in token.bytes() {
1088 h ^= u64::from(b);
1089 h = h.wrapping_mul(0x100000001b3);
1090 }
1091 let mut name = path.as_os_str().to_owned();
1092 name.push(format!(".break-{h:016x}"));
1093 PathBuf::from(name)
1094}
1095
1096fn older_than_stale(path: &Path) -> bool {
1097 std::fs::metadata(path)
1098 .and_then(|m| m.modified())
1099 .ok()
1100 .and_then(|t| t.elapsed().ok())
1101 .is_some_and(|age| age > TASK_LOCK_STALE)
1102}
1103
1104const MAX_MARKER_DEPTH: u8 = 3;
1108
1109fn take_marker(path: &Path, token: &str, depth: u8) -> Option<(PathBuf, String)> {
1115 use std::io::Write;
1116 let marker = break_marker(path, token);
1117 for _ in 0..2 {
1118 match std::fs::OpenOptions::new()
1119 .write(true)
1120 .create_new(true)
1121 .open(&marker)
1122 {
1123 Ok(mut f) => {
1124 let mine = owner_token();
1125 if f.write_all(mine.as_bytes()).is_err() {
1126 drop(f);
1127 let _ = std::fs::remove_file(&marker);
1128 return None;
1129 }
1130 return Some((marker, mine));
1131 }
1132 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1133 let judged = std::fs::read_to_string(&marker).ok();
1134 match judged.filter(|_| older_than_stale(&marker)) {
1135 Some(judged) if depth < MAX_MARKER_DEPTH => {
1136 remove_lock_if(&marker, &judged, true, depth + 1)?;
1137 }
1138 _ => return None,
1139 }
1140 }
1141 Err(_) => return None,
1142 }
1143 }
1144 None
1145}
1146
1147fn remove_lock_if(path: &Path, token: &str, require_stale: bool, depth: u8) -> Option<bool> {
1152 let (marker, mine) = take_marker(path, token, depth)?;
1153 let still = std::fs::read_to_string(path).is_ok_and(|c| c == token)
1154 && (!require_stale || older_than_stale(path));
1155 if still {
1156 let _ = std::fs::remove_file(path);
1157 }
1158 if depth >= MAX_MARKER_DEPTH {
1160 if std::fs::read_to_string(&marker).is_ok_and(|c| c == mine) {
1161 let _ = std::fs::remove_file(&marker);
1162 }
1163 } else {
1164 let _ = remove_lock_if(&marker, &mine, false, depth + 1);
1165 }
1166 Some(still)
1167}
1168
1169fn break_stale(path: &Path, judged: &str) -> bool {
1172 remove_lock_if(path, judged, true, 0).unwrap_or(false)
1173}
1174
1175impl Drop for TaskLock {
1176 fn drop(&mut self) {
1177 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
1178 loop {
1179 match remove_lock_if(&self.path, &self.token, false, 0) {
1180 Some(_) => return,
1181 None if std::time::Instant::now() > deadline => return,
1182 None => std::thread::sleep(std::time::Duration::from_millis(5)),
1183 }
1184 }
1185 }
1186}
1187
1188#[derive(Debug, Clone)]
1190pub struct Queue {
1191 root: PathBuf,
1192}
1193
1194impl Queue {
1195 pub fn open() -> Self {
1197 Self::at(crate::run::home().join("queue"))
1198 }
1199
1200 pub fn at(root: PathBuf) -> Self {
1203 Self { root }
1204 }
1205
1206 pub fn root(&self) -> &Path {
1208 &self.root
1209 }
1210
1211 pub fn path_of(&self, id: &str) -> PathBuf {
1213 self.root.join(format!("{id}.json"))
1214 }
1215
1216 pub fn attachments_dir(&self, id: &str) -> PathBuf {
1219 self.root.join(format!("{id}.attachments"))
1220 }
1221
1222 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1233 let mut wanted = Vec::new();
1234 for src in sources {
1235 let name = src
1236 .file_name()
1237 .and_then(|n| n.to_str())
1238 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1239 if !crate::ask::valid_asset_name(name) {
1240 bail!(
1241 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1242 with no `..`; rename the file and try again"
1243 );
1244 }
1245 if !src.is_file() {
1246 bail!("attachment `{}` is not a file", src.display());
1247 }
1248 wanted.push((src, name));
1249 }
1250 let dir = self.attachments_dir(&task.id);
1251 let existed = dir.is_dir();
1252 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1253 let mut created: Vec<PathBuf> = Vec::new();
1254 let mut names = Vec::new();
1255 let mut copy_all = || -> Result<()> {
1256 for (src, name) in &wanted {
1257 let (stored, path) = copy_new(&dir, src, name)?;
1258 created.push(path);
1259 names.push(stored);
1260 }
1261 Ok(())
1262 };
1263 if let Err(e) = copy_all() {
1264 for path in &created {
1265 let _ = std::fs::remove_file(path);
1266 }
1267 if !existed {
1268 let _ = std::fs::remove_dir(&dir);
1269 }
1270 return Err(e);
1271 }
1272 task.attachments.extend(names.iter().cloned());
1273 Ok(names)
1274 }
1275
1276 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1287 let dir = self.attachments_dir(&task.id);
1288 let existed = dir.is_dir();
1289 let before = task.attachments.len();
1290 let names = self.attach(task, sources)?;
1291 if let Err(e) = self.put(task) {
1292 for name in &names {
1293 let _ = std::fs::remove_file(dir.join(name));
1294 }
1295 if !existed {
1296 let _ = std::fs::remove_dir(&dir);
1297 }
1298 task.attachments.truncate(before);
1299 return Err(e);
1300 }
1301 Ok(names)
1302 }
1303
1304 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1308 let dir = self.attachments_dir(&task.id);
1309 task.attachments
1310 .iter()
1311 .map(|n| {
1312 let p = dir.join(n);
1313 std::path::absolute(&p).unwrap_or(p)
1314 })
1315 .collect()
1316 }
1317
1318 pub fn put(&self, task: &mut Task) -> Result<()> {
1327 let _lock = self.lock_task(&task.id)?;
1328 self.put_unlocked(task)
1329 }
1330
1331 pub fn create_new(&self, task: &mut Task) -> Result<bool> {
1336 let _lock = self.lock_task(&task.id)?;
1337 if self.path_of(&task.id).exists() {
1338 return Ok(false);
1339 }
1340 self.put_unlocked(task)?;
1341 Ok(true)
1342 }
1343
1344 fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1346 if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1347 for run in stored.runs {
1348 if !task.runs.contains(&run) {
1349 task.runs.push(run);
1350 }
1351 }
1352 }
1353 task.updated_at = Timestamp::now();
1354 std::fs::create_dir_all(&self.root)
1355 .with_context(|| format!("create {}", self.root.display()))?;
1356 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1357 let path = self.path_of(&task.id);
1358 let tmp = path.with_extension("json.tmp");
1359 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1360 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1361 if let (Some(notice), Some(home)) = (
1366 crate::notices::task_held(task),
1367 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1368 ) {
1369 crate::notices::raise_in(home, notice);
1370 }
1371 Ok(())
1372 }
1373
1374 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1379 let id = self.resolve_id(id)?;
1380 let _lock = self.lock_task(&id)?;
1383 let mut task = self.get(&id)?;
1384 if task.link_run(run) {
1385 self.put_unlocked(&mut task)?;
1386 }
1387 Ok(task)
1388 }
1389
1390 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1398 std::fs::create_dir_all(&self.root)
1399 .with_context(|| format!("create {}", self.root.display()))?;
1400 let path = self.root.join(format!("{id}.write-lock"));
1401 let started = std::time::Instant::now();
1402 loop {
1403 match std::fs::OpenOptions::new()
1404 .write(true)
1405 .create_new(true)
1406 .open(&path)
1407 {
1408 Ok(mut file) => {
1409 use std::io::Write;
1410 let token = owner_token();
1411 if let Err(e) = file.write_all(token.as_bytes()) {
1412 drop(file);
1413 let _ = std::fs::remove_file(&path);
1414 return Err(e).with_context(|| format!("lock {}", path.display()));
1415 }
1416 return Ok(TaskLock { path, token });
1417 }
1418 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1419 let judged = std::fs::read_to_string(&path).ok();
1422 let broken = judged
1423 .filter(|_| older_than_stale(&path))
1424 .is_some_and(|judged| break_stale(&path, &judged));
1425 if !broken {
1429 if started.elapsed() > TASK_LOCK_STALE {
1430 bail!("could not lock task {id}");
1431 }
1432 std::thread::sleep(std::time::Duration::from_millis(15));
1433 }
1434 }
1435 Err(e)
1439 if e.kind() == std::io::ErrorKind::PermissionDenied
1440 && started.elapsed() <= TASK_LOCK_STALE =>
1441 {
1442 std::thread::sleep(std::time::Duration::from_millis(15));
1443 }
1444 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1445 }
1446 }
1447 }
1448
1449 pub fn get(&self, id: &str) -> Result<Task> {
1451 let resolved = self.resolve_id(id)?;
1452 read_path(&self.path_of(&resolved))
1453 }
1454
1455 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1483 let resolved = self.resolve_id(id)?;
1484 if in_flight {
1485 bail!("task {resolved} is being run by a live daemon right now");
1486 }
1487 self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))?;
1488 let quarantined = self.quarantine_dependents_of(&resolved, questions);
1489 Ok(Removal {
1490 id: resolved,
1491 quarantined,
1492 })
1493 }
1494
1495 fn remove_record_with_attachments(
1499 &self,
1500 resolved: &str,
1501 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1502 ) -> Result<()> {
1503 self.sweep_removed_attachments();
1510 let attachments = self.attachments_dir(resolved);
1511 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1512 let moved = match std::fs::rename(&attachments, &aside) {
1513 Ok(()) => true,
1514 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1515 Err(e) => {
1516 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1517 }
1518 };
1519 let path = self.path_of(resolved);
1520 if let Err(e) = remove_record(&path) {
1521 if moved {
1522 let _ = std::fs::rename(&aside, &attachments);
1523 }
1524 return Err(e).with_context(|| format!("remove {}", path.display()));
1525 }
1526 if moved {
1527 if let Err(e) = std::fs::remove_dir_all(&aside) {
1528 tracing::warn!("leftover attachments {}: {e}", aside.display());
1529 }
1530 }
1531 let lock = self.lock_path(resolved);
1532 if let Err(e) = std::fs::remove_file(&lock) {
1533 if e.kind() != std::io::ErrorKind::NotFound {
1534 return Err(e).with_context(|| format!("remove {}", lock.display()));
1535 }
1536 }
1537 Ok(())
1538 }
1539
1540 fn sweep_removed_attachments(&self) {
1546 let Ok(entries) = std::fs::read_dir(&self.root) else {
1547 return;
1548 };
1549 for entry in entries.flatten() {
1550 let name = entry.file_name();
1551 let name = name.to_string_lossy();
1552 let Some(id) = name.strip_suffix(".attachments.removing") else {
1553 continue;
1554 };
1555 if !self.path_of(id).exists() {
1558 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1559 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1560 }
1561 }
1562 }
1563 }
1564
1565 fn quarantine_dependents_of(&self, dependency: &str, questions: &Questions) -> Vec<String> {
1569 let mut quarantined = Vec::new();
1570 for listed in self.list() {
1571 if listed.status != TaskStatus::Blocked
1572 || !listed.blocked_by.iter().any(|b| b == dependency)
1573 {
1574 continue;
1575 }
1576 let Ok(_claim) = self.claim(&listed.id) else {
1577 continue;
1578 };
1579 let Ok(mut task) = self.get(&listed.id) else {
1580 continue;
1581 };
1582 if task.status != TaskStatus::Blocked
1583 || !task.blocked_by.iter().any(|b| b == dependency)
1584 {
1585 continue;
1586 }
1587 let missing = missing_blockers(self, questions, &task.blocked_by);
1588 let language = crate::lang::of_repo(&task.repo);
1589 task.hold_machine(Some(missing_blocker_hold_reason_in(
1590 &task.blocked_by,
1591 &missing,
1592 &language,
1593 )));
1594 if self.put(&mut task).is_ok() {
1595 quarantined.push(task.id.clone());
1596 }
1597 }
1598 quarantined
1599 }
1600
1601 fn lock_path(&self, id: &str) -> PathBuf {
1604 self.root.join(format!("{id}.lock"))
1605 }
1606
1607 pub fn list(&self) -> Vec<Task> {
1620 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1621 .into_iter()
1622 .flatten()
1623 .flatten()
1624 .map(|e| e.path())
1625 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1626 .filter_map(|p| read_path(&p).ok())
1627 .collect();
1628 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1629 tasks
1630 }
1631
1632 pub fn superseded(&self) -> HashMap<String, String> {
1645 let mut by = HashMap::new();
1646 for task in self.list() {
1647 for earlier in &task.runs {
1648 if let Some(later) = task.successor_of(earlier) {
1649 by.insert(earlier.clone(), later.clone());
1650 }
1651 }
1652 }
1653 by
1654 }
1655
1656 pub fn superseded_by(&self, run: &str) -> Option<String> {
1664 for task in self.list() {
1665 if task.runs.iter().any(|r| r == run) {
1666 return task.successor_of(run).cloned();
1667 }
1668 }
1669 None
1670 }
1671
1672 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1683 for task in self.list() {
1684 if task.runs.iter().any(|r| r == run) {
1685 return task.runs.last().filter(|last| **last != run).cloned();
1686 }
1687 }
1688 None
1689 }
1690
1691 pub fn next_runnable(&self) -> Option<Task> {
1696 let mut runnable: Vec<Task> = self
1697 .list()
1698 .into_iter()
1699 .filter(|t| t.status.runnable())
1700 .collect();
1701 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1702 runnable.into_iter().next()
1703 }
1704
1705 pub fn claim(&self, id: &str) -> Result<Claim> {
1712 std::fs::create_dir_all(&self.root)
1713 .with_context(|| format!("create {}", self.root.display()))?;
1714 let path = self.lock_path(id);
1715 match std::fs::OpenOptions::new()
1716 .write(true)
1717 .create_new(true)
1718 .open(&path)
1719 {
1720 Ok(mut f) => {
1721 use std::io::Write as _;
1722 let _ = writeln!(f, "{}", std::process::id());
1724 Ok(Claim { path })
1725 }
1726 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1727 bail!("task {id} is already claimed ({} exists)", path.display())
1728 }
1729 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1730 }
1731 }
1732
1733 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1735 if self.path_of(prefix).is_file() {
1736 return Ok(prefix.to_owned());
1737 }
1738 let hits: Vec<String> = self
1739 .list()
1740 .into_iter()
1741 .map(|t| t.id)
1742 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1743 .collect();
1744 match hits.len() {
1745 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1746 0 => bail!("no task matches `{prefix}`"),
1747 _ => bail!(
1748 "`{prefix}` matches {} tasks: {}",
1749 hits.len(),
1750 hits.join(", ")
1751 ),
1752 }
1753 }
1754
1755 pub fn revision(&self) -> u64 {
1762 self.revision_excluding(&std::collections::BTreeSet::new())
1763 }
1764
1765 pub fn revision_excluding(&self, skip: &std::collections::BTreeSet<String>) -> u64 {
1768 use std::hash::{Hash as _, Hasher as _};
1769
1770 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
1771 .into_iter()
1772 .flatten()
1773 .flatten()
1774 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
1775 .filter(|e| {
1776 let path = e.path();
1777 !path
1778 .file_stem()
1779 .is_some_and(|stem| skip.contains(stem.to_string_lossy().as_ref()))
1780 })
1781 .filter_map(|e| {
1782 let name = e.file_name().to_string_lossy().into_owned();
1783 let mtime = e
1784 .metadata()
1785 .ok()?
1786 .modified()
1787 .ok()?
1788 .duration_since(std::time::UNIX_EPOCH)
1789 .ok()?
1790 .as_millis() as u64;
1791 Some((name, mtime))
1792 })
1793 .collect();
1794
1795 if entries.is_empty() {
1796 return 0;
1797 }
1798
1799 entries.sort_unstable();
1800 let mut hasher = std::hash::DefaultHasher::new();
1801 for (name, mtime) in &entries {
1802 name.hash(&mut hasher);
1803 mtime.hash(&mut hasher);
1804 }
1805 let h = hasher.finish();
1806 if h == 0 { 1 } else { h }
1807 }
1808}
1809
1810#[derive(Debug, Clone)]
1812pub struct Removal {
1813 pub id: String,
1815 pub quarantined: Vec<String>,
1819}
1820
1821#[derive(Debug)]
1823pub struct Claim {
1824 path: PathBuf,
1825}
1826
1827impl Drop for Claim {
1828 fn drop(&mut self) {
1829 let _ = std::fs::remove_file(&self.path);
1830 }
1831}
1832
1833pub fn title_from(instruction: &str, max: usize) -> String {
1836 let line = instruction
1842 .lines()
1843 .map(str::trim)
1844 .find(|l| !l.is_empty())
1845 .unwrap_or("(empty task)")
1846 .trim_start_matches(['#', '-', '*', '>', ' '])
1847 .trim();
1848 if line.is_empty() {
1849 return "(empty task)".to_owned();
1850 }
1851 if line.chars().count() <= max {
1852 return line.to_owned();
1853 }
1854 let head: String = line.chars().take(max.saturating_sub(1)).collect();
1855 format!("{head}…")
1856}
1857
1858fn read_path(path: &Path) -> Result<Task> {
1859 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1860 let task: Task =
1861 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1862 if task.schema > SCHEMA {
1868 bail!(
1869 "task {} was written by a different magi (schema {}, this build \
1870 speaks {SCHEMA})",
1871 task.id,
1872 task.schema
1873 );
1874 }
1875 Ok(task)
1876}
1877
1878pub fn missing_blockers(
1892 queue: &Queue,
1893 questions: &Questions,
1894 blocked_by: &[String],
1895) -> Vec<String> {
1896 blocked_by
1897 .iter()
1898 .filter(|id| !queue.path_of(id).is_file() && !questions.path_of(id).is_file())
1899 .cloned()
1900 .collect()
1901}
1902
1903pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
1915 missing_blocker_hold_reason_in(blocked_by, missing, "en")
1916}
1917
1918pub fn missing_blocker_hold_reason_in(
1920 blocked_by: &[String],
1921 missing: &[String],
1922 language: &str,
1923) -> String {
1924 if crate::lang::is_japanese(language) {
1925 format!(
1926 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
1927 blocked_by.join(", "),
1928 missing.join(", "),
1929 )
1930 } else {
1931 format!(
1932 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
1933 blocked_by.join(", "),
1934 missing.join(", "),
1935 )
1936 }
1937}
1938
1939pub fn short(id: &str) -> &str {
1941 id.split('-').next_back().unwrap_or(id)
1942}
1943
1944fn new_id() -> String {
1945 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1946 let seed = crate::rng::entropy();
1947 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1948}
1949
1950#[cfg(test)]
1951mod tests {
1952 #[test]
1953 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
1954 let mut t = task("relanded");
1955 t.attempts = 1;
1956 t.status = TaskStatus::Running;
1957 t.already_landed("already in main as 0e368de");
1958 assert_eq!(t.status, TaskStatus::Done);
1959 assert_eq!(t.attempts, 0);
1960 assert!(t.hold_reason.is_none());
1961 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
1962 }
1963
1964 #[test]
1965 fn missing_blocker_reason_follows_the_language() {
1966 let b = vec!["a".to_owned()];
1967 let en = missing_blocker_hold_reason_in(&b, &b, "en");
1968 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
1969 assert!(en.starts_with("blocked on a"));
1970 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
1971 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
1972 }
1973
1974 use super::*;
1975
1976 #[test]
1977 fn triage_applied_survives_release_and_old_records_read_as_empty() {
1978 let mut t = Task::new(
1979 "t".to_owned(),
1980 "i".to_owned(),
1981 PathBuf::from("r"),
1982 Source::Human,
1983 );
1984 t.mark_triage_applied("q1");
1985 t.mark_triage_applied("q1");
1986 t.hold_machine(Some("x".to_owned()));
1987 t.release();
1988 assert_eq!(t.triage_applied, ["q1"]);
1989 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
1990
1991 let mut v = serde_json::to_value(&t).unwrap();
1992 v.as_object_mut().unwrap().remove("triage_applied");
1993 let old: Task = serde_json::from_value(v).unwrap();
1994 assert!(old.triage_applied.is_empty());
1995 }
1996
1997 #[test]
1998 fn task_counts_of_empty_is_all_zero() {
1999 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
2000 }
2001
2002 #[test]
2003 fn task_counts_of_tallies_every_status() {
2004 let mut queued = Task::new(
2005 "q".to_owned(),
2006 "i".to_owned(),
2007 PathBuf::from("."),
2008 Source::Human,
2009 );
2010 queued.status = TaskStatus::Queued;
2011 let mut running = queued.clone();
2012 running.status = TaskStatus::Running;
2013 let mut done = queued.clone();
2014 done.status = TaskStatus::Done;
2015 let mut failed = queued.clone();
2016 failed.status = TaskStatus::Failed;
2017 let mut held = queued.clone();
2018 held.status = TaskStatus::Held;
2019 let mut blocked = queued.clone();
2020 blocked.status = TaskStatus::Blocked;
2021
2022 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
2023 assert_eq!(
2024 counts,
2025 TaskCounts {
2026 queued: 1,
2027 running: 1,
2028 done: 2,
2029 failed: 1,
2030 held: 1,
2031 blocked: 1,
2032 }
2033 );
2034 }
2035
2036 fn queue() -> (tempfile::TempDir, Queue) {
2039 let dir = tempfile::tempdir().unwrap();
2040 let q = Queue::at(dir.path().join("queue"));
2041 (dir, q)
2042 }
2043
2044 #[test]
2045 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
2046 let dir = tempfile::tempdir().unwrap();
2047 let q = Queue::at(dir.path().join("queue"));
2048 let mut t = task("held");
2049 q.put(&mut t).unwrap();
2050 assert_eq!(
2051 crate::notices::Notices::at(dir.path().join("notifications"))
2052 .list()
2053 .len(),
2054 0
2055 );
2056 t.hold_machine(Some("out of attempts".to_owned()));
2057 q.put(&mut t).unwrap();
2058 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
2059 assert_eq!(listed.len(), 1);
2060 assert!(listed[0].message.contains("out of attempts"));
2061 }
2062
2063 fn task(title: &str) -> Task {
2064 Task::new(
2065 title.to_owned(),
2066 format!("do {title}"),
2067 PathBuf::from("."),
2068 Source::Human,
2069 )
2070 }
2071
2072 #[test]
2073 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
2074 let mut t = task("retried");
2075 assert!(
2076 t.earlier_attempts().is_empty(),
2077 "a first attempt takes nothing over"
2078 );
2079 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2080 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
2081 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
2083 assert_eq!(t.successor_of("bbbb"), None);
2084 assert_eq!(t.successor_of("zzzz"), None);
2085 }
2086
2087 #[test]
2088 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
2089 let (_dir, q) = queue();
2090 let mut t = task("retried");
2091 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2092 q.put(&mut t).unwrap();
2093
2094 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
2095 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
2096 assert_eq!(
2097 q.superseded_by("cccc"),
2098 None,
2099 "the latest attempt replaces nothing"
2100 );
2101 assert_eq!(
2102 q.superseded_by("never-heard-of-it"),
2103 None,
2104 "a run belonging to no task on this queue is not superseded"
2105 );
2106
2107 let mut by = HashMap::new();
2108 by.insert("aaaa".to_owned(), "bbbb".to_owned());
2109 by.insert("bbbb".to_owned(), "cccc".to_owned());
2110 assert_eq!(
2111 q.superseded(),
2112 by,
2113 "the whole-map and single-run forms must agree"
2114 );
2115 }
2116
2117 #[test]
2118 fn superseded_attempts_is_empty_until_the_task_is_done() {
2119 let mut t = task("retried");
2120 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2121 t.status = TaskStatus::Failed;
2122 assert_eq!(
2123 t.superseded_attempts(true),
2124 &[] as &[String],
2125 "a task still retrying has no attempt yet that a later one made moot"
2126 );
2127
2128 t.status = TaskStatus::Running;
2129 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2130 }
2131
2132 #[test]
2133 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
2134 let mut t = task("retried");
2135 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2136 t.status = TaskStatus::Done;
2137 assert_eq!(
2138 t.superseded_attempts(true),
2139 &["aaaa".to_owned(), "bbbb".to_owned()],
2140 "cccc is the attempt whose success made the task done, and stays out"
2141 );
2142 }
2143
2144 #[test]
2145 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
2146 let mut t = task("first try landed");
2147 t.runs = vec!["aaaa".to_owned()];
2148 t.status = TaskStatus::Done;
2149 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2150 }
2151
2152 #[test]
2153 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
2154 let mut t = task("closed by hand after a manual merge");
2161 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2162 t.status = TaskStatus::Done;
2163 assert_eq!(
2164 t.superseded_attempts(false),
2165 &[] as &[String],
2166 "nothing here is provably why the task is done, so nothing is superseded"
2167 );
2168 }
2169
2170 #[test]
2171 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2172 let (_dir, q) = queue();
2173 let mut t = task("retried twice");
2174 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2175 q.put(&mut t).unwrap();
2176
2177 assert_eq!(
2178 q.latest_attempt("aaaa"),
2179 Some("cccc".to_owned()),
2180 "an old attempt points straight at the chain's current head, not the \
2181 next attempt in the middle of it"
2182 );
2183 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2184 assert_eq!(
2185 q.latest_attempt("cccc"),
2186 None,
2187 "the latest attempt is not superseded by anything"
2188 );
2189 assert_eq!(
2190 q.latest_attempt("never-heard-of-it"),
2191 None,
2192 "a run belonging to no task on this queue is not superseded"
2193 );
2194 }
2195
2196 #[test]
2197 fn a_markdown_heading_is_the_title_not_decoration() {
2198 assert_eq!(
2203 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2204 "Rework the config loader"
2205 );
2206 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2207 assert_eq!(title_from("> quoted task", 40), "quoted task");
2208 assert_eq!(title_from(" \n\n", 40), "(empty task)");
2210 assert_eq!(title_from("###\n", 40), "(empty task)");
2211 }
2212
2213 #[test]
2214 fn a_long_title_is_elided_by_characters_not_bytes() {
2215 let long = "課題".repeat(30);
2217 let title = title_from(&long, 10);
2218 assert_eq!(title.chars().count(), 10);
2219 assert!(title.ends_with('…'));
2220 }
2221
2222 #[test]
2223 fn priority_wins_and_ties_break_oldest_first() {
2224 let (_dir, q) = queue();
2225 let mut a = task("first");
2226 let mut b = task("second");
2227 let mut c = task("urgent");
2228 a.id = "20260101-000001-aaaa".to_owned();
2230 b.id = "20260101-000002-bbbb".to_owned();
2231 c.id = "20260101-000003-cccc".to_owned();
2232 c.priority = 5;
2233 for t in [&mut a, &mut b, &mut c] {
2234 q.put(t).unwrap();
2235 }
2236
2237 assert_eq!(q.next_runnable().unwrap().id, c.id);
2239 c.hold_machine(None);
2240 q.put(&mut c).unwrap();
2241 assert_eq!(q.next_runnable().unwrap().id, a.id);
2243 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2244 }
2245
2246 #[test]
2247 fn a_blocked_task_never_starves_another_runnable_one() {
2248 let (_dir, q) = queue();
2249 let mut blocked = task("blocked");
2250 blocked.block(vec!["something".to_owned()], None);
2251 q.put(&mut blocked).unwrap();
2252
2253 let mut runnable = task("free to go");
2254 q.put(&mut runnable).unwrap();
2255
2256 let next = q.next_runnable().expect("a runnable task is still offered");
2257 assert_eq!(next.id, runnable.id);
2258 }
2259
2260 #[test]
2261 fn a_held_task_is_never_offered_to_the_loop() {
2262 let (_dir, q) = queue();
2263 let mut t = task("held");
2264 q.put(&mut t).unwrap();
2265 assert!(q.next_runnable().is_some());
2266
2267 t.hold_machine(None);
2268 q.put(&mut t).unwrap();
2269 assert!(
2270 q.next_runnable().is_none(),
2271 "a held task must wait for a human"
2272 );
2273
2274 t.status = TaskStatus::Failed;
2276 q.put(&mut t).unwrap();
2277 assert!(q.next_runnable().is_some());
2278 }
2279
2280 #[test]
2281 fn attempts_are_capped_and_then_the_task_is_held() {
2282 let mut t = task("doomed");
2283
2284 t.start("run-1".to_owned());
2285 t.fail("gate red", 2);
2286 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2287
2288 t.start("run-2".to_owned());
2289 t.fail("gate red", 2);
2290 assert_eq!(
2291 t.status,
2292 TaskStatus::Held,
2293 "out of attempts: stop spending money on it"
2294 );
2295 assert_eq!(t.runs, ["run-1", "run-2"]);
2296 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2297 assert_eq!(
2298 t.hold_reason.as_deref(),
2299 Some("gate red"),
2300 "the hold must say why, not leave hold_reason null next to a \
2301 populated last_error"
2302 );
2303 }
2304
2305 #[test]
2306 fn handing_off_a_task_records_a_hold_reason_too() {
2307 let mut t = task("left a pull request");
2308 t.start("run-1".to_owned());
2309 t.handed_off("run ended with a pull request open [run run-1]");
2310 assert_eq!(t.status, TaskStatus::Held);
2311 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2312 assert_eq!(
2313 t.hold_reason.as_deref(),
2314 Some("run ended with a pull request open [run run-1]")
2315 );
2316 assert_eq!(t.hold_reason, t.last_error);
2317 }
2318
2319 #[test]
2320 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2321 let mut t = task("stalled by quota");
2322
2323 t.start("run-1".to_owned());
2324 assert_eq!(t.attempts, 1);
2325 t.stall("judge-1, judge-2 out of quota");
2326 assert_eq!(
2327 t.attempts, 0,
2328 "a closed quota window must not spend the task's retry budget"
2329 );
2330 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2331 assert_eq!(
2332 t.last_error.as_deref(),
2333 Some("judge-1, judge-2 out of quota")
2334 );
2335
2336 for _ in 0..20 {
2339 t.start("run-n".to_owned());
2340 t.stall("still out of quota");
2341 }
2342 t.start("run-real".to_owned());
2343 t.fail("gate red", 2);
2344 assert_eq!(
2345 t.status,
2346 TaskStatus::Failed,
2347 "the first attempt that was really judged is attempt one"
2348 );
2349 }
2350
2351 #[test]
2352 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2353 let mut t = task("retry me");
2354 t.start("run-1".to_owned());
2355 t.fail("gate red", 1);
2356 assert_eq!(t.status, TaskStatus::Held);
2357
2358 t.release();
2359 assert_eq!(t.status, TaskStatus::Queued);
2360 assert_eq!(t.attempts, 0);
2363 assert!(t.last_error.is_none());
2364 assert_eq!(
2365 t.runs.len(),
2366 1,
2367 "history is kept: attempts reset, evidence does not"
2368 );
2369 }
2370
2371 #[test]
2372 fn a_hold_reason_survives_and_a_release_clears_it() {
2373 let mut t = task("waiting on something else");
2374 t.hold_manual(Some(
2375 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2376 ));
2377 assert_eq!(t.status, TaskStatus::Held);
2378 assert_eq!(
2379 t.hold_reason.as_deref(),
2380 Some("waiting for 20260101-000000-aaaa to land first")
2381 );
2382
2383 t.hold_manual(None);
2385 assert_eq!(
2386 t.hold_reason.as_deref(),
2387 Some("waiting for 20260101-000000-aaaa to land first"),
2388 "a bare re-hold keeps whatever a human already wrote down"
2389 );
2390
2391 let mut plain = task("no reason given");
2393 plain.hold_manual(None);
2394 assert_eq!(plain.status, TaskStatus::Held);
2395 assert!(plain.hold_reason.is_none());
2396
2397 t.release();
2398 assert_eq!(t.status, TaskStatus::Queued);
2399 assert!(
2400 t.hold_reason.is_none(),
2401 "a stale reason must not greet the next person who holds this task"
2402 );
2403 }
2404
2405 #[test]
2406 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2407 let mut t = task("landed by hand while held");
2412 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2413 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2414
2415 t.succeed();
2416 assert_eq!(t.status, TaskStatus::Done);
2417 assert!(
2418 t.hold_reason.is_none(),
2419 "a done task cannot still be waiting on something"
2420 );
2421 }
2422
2423 #[test]
2424 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2425 let mut held = task("held straight out of blocked");
2432 held.block(
2433 vec!["20260101-000000-dead".to_owned()],
2434 Some("waiting on the migration script".to_owned()),
2435 );
2436 assert_eq!(held.status, TaskStatus::Blocked);
2437
2438 held.hold_manual(None);
2439 assert_eq!(held.status, TaskStatus::Held);
2440 assert!(
2441 held.blocked_by.is_empty(),
2442 "hold overrides the wait, same as release"
2443 );
2444 assert!(held.block_reason.is_none());
2445
2446 let mut done = task("closed straight out of blocked");
2447 done.block(
2448 vec!["20260101-000000-dead".to_owned()],
2449 Some("waiting on the migration script".to_owned()),
2450 );
2451 done.succeed();
2452 assert_eq!(done.status, TaskStatus::Done);
2453 assert!(
2454 done.blocked_by.is_empty(),
2455 "a done task cannot still be waiting on a dependency"
2456 );
2457 assert!(done.block_reason.is_none());
2458 }
2459
2460 #[test]
2461 fn a_blocked_task_is_never_offered_to_the_loop() {
2462 let mut t = task("blocked");
2463 assert!(t.status.runnable());
2464 t.block(
2465 vec!["dep-id".to_owned()],
2466 Some("waits on dep-id".to_owned()),
2467 );
2468 assert_eq!(t.status, TaskStatus::Blocked);
2469 assert!(!t.status.runnable());
2470 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2471 }
2472
2473 #[test]
2474 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2475 let mut t = task("blocked on two");
2476 t.block(
2477 vec!["a".to_owned(), "b".to_owned()],
2478 Some("waits on a and b".to_owned()),
2479 );
2480
2481 t.unblock("a");
2482 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2483 assert_eq!(t.blocked_by, ["b"]);
2484
2485 t.unblock("b");
2486 assert_eq!(t.status, TaskStatus::Queued);
2487 assert!(t.blocked_by.is_empty());
2488 assert!(t.block_reason.is_none());
2489 }
2490
2491 #[test]
2492 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2493 let mut t = task("never blocked");
2494 t.unblock("whatever");
2495 assert_eq!(t.status, TaskStatus::Queued);
2496 }
2497
2498 #[test]
2499 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2500 let mut t = task("held, then asked about");
2505 t.hold_machine(Some("out of attempts".to_owned()));
2506 assert_eq!(t.status, TaskStatus::Held);
2507
2508 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2509 assert_eq!(t.status, TaskStatus::Blocked);
2510
2511 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2512 t.unblock("q1");
2513 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2514 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2515 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2516 assert!(t.blocked_from.is_none(), "consumed once restored");
2517 }
2518
2519 #[test]
2520 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2521 let mut t = task("manually held, then asked about");
2522 t.hold_manual(Some("waiting on a dependency".to_owned()));
2523
2524 t.block(vec!["q1".to_owned()], None);
2525 t.unblock("q1");
2526
2527 assert_eq!(t.status, TaskStatus::Held);
2528 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2529 }
2530
2531 #[test]
2532 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2533 let mut t = task("held, blocked twice");
2537 t.hold_machine(None);
2538 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2539 t.block(
2540 vec!["q1".to_owned(), "q2".to_owned()],
2541 Some("second".to_owned()),
2542 );
2543
2544 t.unblock("q1");
2545 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2546 t.unblock("q2");
2547 assert_eq!(t.status, TaskStatus::Held);
2548 }
2549
2550 #[test]
2551 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2552 let mut t = task("blocked mid-run");
2556 t.start("run-1".to_owned());
2557 assert_eq!(t.status, TaskStatus::Running);
2558
2559 t.block(vec!["q1".to_owned()], None);
2560 t.unblock("q1");
2561 assert_eq!(t.status, TaskStatus::Queued);
2562 }
2563
2564 #[test]
2565 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2566 let mut t = task("legacy record, held before it was blocked");
2572 t.hold_source = Some(HoldSource::Machine);
2573 t.hold_reason = Some("legacy hold reason".to_owned());
2574 t.status = TaskStatus::Blocked;
2575 t.blocked_by = vec!["q1".to_owned()];
2576 t.blocked_from = None;
2577
2578 t.unblock("q1");
2579 assert_eq!(t.status, TaskStatus::Held);
2580 }
2581
2582 #[test]
2583 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2584 let mut t = task("legacy record, ordinary dependency block");
2585 t.status = TaskStatus::Blocked;
2586 t.blocked_by = vec!["dep".to_owned()];
2587 t.blocked_from = None;
2588
2589 t.unblock("dep");
2590 assert_eq!(t.status, TaskStatus::Queued);
2591 }
2592
2593 #[test]
2594 fn answering_a_question_is_recorded_and_survives_a_release() {
2595 let mut t = task("asked something");
2596 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2597 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2598 t.unblock("q1");
2599 assert_eq!(t.status, TaskStatus::Queued);
2600 assert_eq!(t.answers.len(), 1);
2601 assert_eq!(t.answers[0].answer, "SQLite");
2602
2603 t.release();
2607 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2608 }
2609
2610 #[test]
2611 fn a_refused_handover_keeps_the_review_branch_across_release() {
2612 let mut t = task("refused takeover");
2613 t.start("run-1".to_owned());
2614 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
2615 assert_eq!(t.status, TaskStatus::Held);
2616 assert_eq!(t.attempts, 1);
2617 t.release();
2618 assert_eq!(t.status, TaskStatus::Queued);
2619 assert_eq!(t.attempts, 0);
2620 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2621
2622 t.hold_for_handover(None, "again".to_owned());
2624 t.requeue();
2625 assert!(t.review_branch.is_none());
2626
2627 let mut m = task("manual");
2629 m.review_branch = Some("magi/x/A".to_owned());
2630 m.hold_manual(None);
2631 m.release();
2632 assert!(m.review_branch.is_none());
2633 }
2634
2635 #[test]
2636 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2637 let mut t = task("blocked run with a surviving branch");
2638 t.start("run-1".to_owned());
2639 t.fail("blocked with major findings", 5);
2640 assert_eq!(t.status, TaskStatus::Failed);
2641
2642 t.request_review("magi/eba2/A".to_owned());
2643 assert_eq!(t.status, TaskStatus::Queued);
2644 assert_eq!(t.attempts, 0);
2645 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2646
2647 t.release();
2649 assert!(t.review_branch.is_none());
2650 }
2651
2652 #[test]
2653 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2654 let mut t = task("retry");
2655 t.start("run-1".to_owned());
2656 t.requeue();
2657 assert!(t.fresh_start);
2658
2659 t.release();
2660 assert!(!t.fresh_start);
2661 }
2662
2663 #[test]
2664 fn priority_can_be_changed_while_queued_but_not_while_running() {
2665 let mut t = task("reprioritise me");
2666 t.set_priority(5).unwrap();
2667 assert_eq!(t.priority, 5);
2668
2669 t.start("run-1".to_owned());
2670 let err = t.set_priority(9).unwrap_err().to_string();
2671 assert!(err.contains("running"), "{err}");
2672 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2673 }
2674
2675 #[test]
2676 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2677 let mut t = task("interrupt me");
2678 assert!(!t.interrupt, "off unless asked, same as any other task");
2679
2680 t.set_interrupt(true).unwrap();
2681 assert!(t.interrupt);
2682
2683 t.start("run-1".to_owned());
2684 assert!(
2685 !t.interrupt,
2686 "the mark is one-shot: dispatching the task fulfils it, \
2687 whatever the run that follows ends up doing"
2688 );
2689 let err = t.set_interrupt(true).unwrap_err().to_string();
2690 assert!(err.contains("running"), "{err}");
2691 t.set_interrupt(false).unwrap();
2694 assert!(!t.interrupt);
2695 }
2696
2697 #[test]
2701 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2702 let mut t = task("interrupt me");
2703 t.set_interrupt(true).unwrap();
2704 t.start("run-1".to_owned());
2705 t.fail("mock failure", 5);
2706 assert_eq!(t.status, TaskStatus::Failed);
2707 assert!(
2708 !t.interrupt,
2709 "one attempt already spent the mark; a retry is an ordinary \
2710 requeue, not a fresh interrupt request"
2711 );
2712 }
2713
2714 #[test]
2715 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2716 let (_dir, q) = queue();
2717 let mut a = task("first filed");
2718 let mut b = task("second filed");
2719 a.id = "20260101-000001-aaaa".to_owned();
2720 b.id = "20260101-000002-bbbb".to_owned();
2721 q.put(&mut a).unwrap();
2722 q.put(&mut b).unwrap();
2723
2724 assert_eq!(
2725 q.next_runnable().unwrap().id,
2726 a.id,
2727 "with equal priority the older task goes first, so a burst of \
2728 new work cannot starve it"
2729 );
2730 assert_eq!(
2731 q.list()[0].id,
2732 b.id,
2733 "but the list an operator reads is newest first, the same as \
2734 before priority existed - a's turn to run does not make it the \
2735 newest task"
2736 );
2737
2738 let mut a = q.get(&a.id).unwrap();
2739 a.set_priority(10).unwrap();
2740 q.put(&mut a).unwrap();
2741
2742 assert_eq!(
2743 q.next_runnable().unwrap().id,
2744 a.id,
2745 "a raised priority must be reflected the moment it is saved"
2746 );
2747 assert_eq!(
2751 q.list()[0].id,
2752 a.id,
2753 "the raised task must sort first in the list an operator reads, \
2754 not only in next_runnable's own ordering"
2755 );
2756 }
2757
2758 #[test]
2759 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2760 let mut t = Task::new(
2761 "old title".to_owned(),
2762 "old instruction".to_owned(),
2763 PathBuf::from("/repo"),
2764 Source::Agent {
2765 run: "20260101-000000-beef".to_owned(),
2766 node: "implement".to_owned(),
2767 },
2768 );
2769 let id = t.id.clone();
2770 let created_at = t.created_at;
2771 t.runs.push("20260101-000000-beef".to_owned());
2772
2773 t.edit("new title".to_owned(), "new instruction".to_owned())
2774 .unwrap();
2775
2776 assert_eq!(t.title, "new title");
2777 assert_eq!(t.instruction, "new instruction");
2778 assert_eq!(t.id, id, "editing must not mint a new id");
2779 assert_eq!(t.created_at, created_at);
2780 assert_eq!(
2781 t.source,
2782 Source::Agent {
2783 run: "20260101-000000-beef".to_owned(),
2784 node: "implement".to_owned(),
2785 },
2786 "editing must not turn agent attribution into human"
2787 );
2788 assert_eq!(t.runs, ["20260101-000000-beef"]);
2789 }
2790
2791 #[test]
2792 fn editing_is_refused_once_a_task_is_running_or_finished() {
2793 let mut running = task("in flight");
2794 running.start("run-1".to_owned());
2795 let err = running
2796 .edit("x".to_owned(), "y".to_owned())
2797 .unwrap_err()
2798 .to_string();
2799 assert!(err.contains("running"), "{err}");
2800
2801 let mut done = task("finished");
2802 done.succeed();
2803 let err = done
2804 .edit("x".to_owned(), "y".to_owned())
2805 .unwrap_err()
2806 .to_string();
2807 assert!(err.contains("done"), "{err}");
2808
2809 let mut queued = task("waiting");
2811 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2812 let mut held = task("parked");
2813 held.hold_machine(None);
2814 held.edit("x".to_owned(), "y".to_owned()).unwrap();
2815 }
2816
2817 #[test]
2818 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2819 let (_dir, q) = queue();
2820 let path = q.path_of("20260101-000000-aaaa");
2821 std::fs::create_dir_all(q.root()).unwrap();
2822 std::fs::write(
2823 &path,
2824 serde_json::json!({
2825 "schema": SCHEMA,
2826 "id": "20260101-000000-aaaa",
2827 "title": "from before hold reasons existed",
2828 "instruction": "from before hold reasons existed",
2829 "repo": ".",
2830 "source": { "kind": "human" },
2831 "status": "held",
2832 "created_at": Timestamp::now().to_string(),
2833 "updated_at": Timestamp::now().to_string(),
2834 })
2835 .to_string(),
2836 )
2837 .unwrap();
2838
2839 let task = q.get("20260101-000000-aaaa").expect("must still read");
2840 assert!(task.hold_reason.is_none());
2841 assert!(task.operator_held());
2842 }
2843
2844 #[test]
2845 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
2846 let (_dir, q) = queue();
2847 let path = q.path_of("20260101-000000-bbbb");
2848 std::fs::create_dir_all(q.root()).unwrap();
2849 std::fs::write(
2850 &path,
2851 serde_json::json!({
2852 "schema": 2,
2853 "id": "20260101-000000-bbbb",
2854 "title": "old manual recovery",
2855 "instruction": "old manual recovery",
2856 "repo": ".",
2857 "source": { "kind": "human" },
2858 "status": "held",
2859 "hold_reason": "active manual recovery run20260912-224242-daf5",
2860 "created_at": Timestamp::now().to_string(),
2861 "updated_at": Timestamp::now().to_string(),
2862 })
2863 .to_string(),
2864 )
2865 .unwrap();
2866
2867 let task = q.get("20260101-000000-bbbb").expect("must still read");
2868 assert_eq!(task.hold_source, None);
2869 assert!(task.operator_held());
2870 }
2871
2872 #[test]
2873 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
2874 let (_dir, q) = queue();
2875 let path = q.path_of("20260101-000000-aaaa");
2876 std::fs::create_dir_all(q.root()).unwrap();
2877 std::fs::write(
2878 &path,
2879 serde_json::json!({
2880 "schema": SCHEMA,
2881 "id": "20260101-000000-aaaa",
2882 "title": "from before diagnostics existed",
2883 "instruction": "from before diagnostics existed",
2884 "repo": ".",
2885 "source": { "kind": "human" },
2886 "status": "held",
2887 "created_at": Timestamp::now().to_string(),
2888 "updated_at": Timestamp::now().to_string(),
2889 })
2890 .to_string(),
2891 )
2892 .unwrap();
2893
2894 let task = q.get("20260101-000000-aaaa").expect("must still read");
2895 assert!(task.diagnostic.is_none());
2896 }
2897
2898 #[test]
2899 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
2900 let (_dir, q) = queue();
2904 let path = q.path_of("20260101-000000-aaaa");
2905 std::fs::create_dir_all(q.root()).unwrap();
2906 std::fs::write(
2907 &path,
2908 serde_json::json!({
2909 "schema": 1,
2910 "id": "20260101-000000-aaaa",
2911 "title": "from before blocking existed",
2912 "instruction": "from before blocking existed",
2913 "repo": ".",
2914 "source": { "kind": "human" },
2915 "status": "queued",
2916 "created_at": Timestamp::now().to_string(),
2917 "updated_at": Timestamp::now().to_string(),
2918 })
2919 .to_string(),
2920 )
2921 .unwrap();
2922
2923 let task = q.get("20260101-000000-aaaa").expect("must still read");
2924 assert!(task.blocked_by.is_empty());
2925 assert!(task.block_reason.is_none());
2926 assert!(task.answers.is_empty());
2927 assert!(task.review_branch.is_none());
2928 }
2929
2930 #[test]
2931 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
2932 let mut held = task("diagnosed");
2937 held.start("run-1".to_owned());
2938 held.fail("gate red", 1);
2939 held.diagnostic = Some("cargo test failed: ...".to_owned());
2940 assert_eq!(held.status, TaskStatus::Held);
2941
2942 held.release();
2943 assert!(held.diagnostic.is_none());
2944
2945 held.diagnostic = Some("cargo test failed: ...".to_owned());
2946 held.succeed();
2947 assert!(held.diagnostic.is_none());
2948 }
2949
2950 #[test]
2951 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
2952 let mut t = task("retried");
2953 t.start("run-1".to_owned());
2954 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
2955 t.fail("unrelated config error", 5);
2956 assert_eq!(t.status, TaskStatus::Failed);
2957 assert!(
2958 t.diagnostic.is_none(),
2959 "fail() must not let an old diagnostic outlive the run that produced it"
2960 );
2961 }
2962
2963 #[test]
2964 fn a_claim_is_exclusive_and_releases_on_drop() {
2965 let (_dir, q) = queue();
2966 let mut t = task("contended");
2967 q.put(&mut t).unwrap();
2968
2969 let held = q.claim(&t.id).unwrap();
2970 assert!(
2971 q.claim(&t.id).is_err(),
2972 "two daemons must not drive one task into two runs"
2973 );
2974 drop(held);
2975 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
2976 }
2977
2978 #[test]
2979 fn a_round_trip_survives_disk() {
2980 let (_dir, q) = queue();
2981 let mut t = Task::new(
2982 "titled".to_owned(),
2983 "body".to_owned(),
2984 PathBuf::from("/repo"),
2985 Source::Agent {
2986 run: "20260101-000000-beef".to_owned(),
2987 node: "implement".to_owned(),
2988 },
2989 );
2990 t.priority = 3;
2991 q.put(&mut t).unwrap();
2992
2993 let back = q.get(&t.id).unwrap();
2994 assert_eq!(back.id, t.id);
2995 assert_eq!(back.priority, 3);
2996 assert_eq!(back.source.label(), "implement@beef");
2997 assert_eq!(q.get(t.short()).unwrap().id, t.id);
2999 }
3000
3001 #[test]
3002 fn an_unreadable_task_does_not_take_the_queue_down() {
3003 let (_dir, q) = queue();
3004 let mut t = task("fine");
3005 q.put(&mut t).unwrap();
3006 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
3007
3008 let listed = q.list();
3009 assert_eq!(listed.len(), 1, "the readable task still lists");
3010 assert_eq!(listed[0].id, t.id);
3011 }
3012
3013 #[test]
3014 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
3015 let (_dir, q) = queue();
3016 let path = q.path_of("20260101-000000-aaaa");
3017 std::fs::create_dir_all(q.root()).unwrap();
3018 std::fs::write(
3019 &path,
3020 serde_json::json!({
3021 "schema": SCHEMA,
3022 "id": "20260101-000000-aaaa",
3023 "title": "from before solo existed",
3024 "instruction": "from before solo existed",
3025 "repo": ".",
3026 "source": { "kind": "human" },
3027 "status": "queued",
3028 "created_at": Timestamp::now().to_string(),
3029 "updated_at": Timestamp::now().to_string(),
3030 })
3031 .to_string(),
3032 )
3033 .unwrap();
3034
3035 let task = q.get("20260101-000000-aaaa").expect("must still read");
3036 assert!(!task.solo, "a queue file with no `solo` field means false");
3037 }
3038
3039 #[test]
3040 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
3041 let (_dir, q) = queue();
3042 let path = q.path_of("20260101-000000-bbbb");
3043 std::fs::create_dir_all(q.root()).unwrap();
3044 std::fs::write(
3045 &path,
3046 serde_json::json!({
3047 "schema": SCHEMA,
3048 "id": "20260101-000000-bbbb",
3049 "title": "from before urgent existed",
3050 "instruction": "from before urgent existed",
3051 "repo": ".",
3052 "source": { "kind": "human" },
3053 "status": "queued",
3054 "created_at": Timestamp::now().to_string(),
3055 "updated_at": Timestamp::now().to_string(),
3056 })
3057 .to_string(),
3058 )
3059 .unwrap();
3060
3061 let task = q.get("20260101-000000-bbbb").expect("must still read");
3062 assert!(
3063 !task.urgent,
3064 "a queue file with no `urgent` field means false, same as `solo`"
3065 );
3066 }
3067
3068 #[test]
3069 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
3070 let (_dir, q) = queue();
3071 let mut t = task("from the future");
3072 q.put(&mut t).unwrap();
3073 let path = q.path_of(&t.id);
3074 let body = std::fs::read_to_string(&path)
3075 .unwrap()
3076 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
3077 std::fs::write(&path, body).unwrap();
3078
3079 let err = q.get(&t.id).unwrap_err().to_string();
3080 assert!(err.contains("schema 99"), "{err}");
3081 }
3082
3083 #[test]
3084 fn revision_moves_when_the_queue_changes() {
3085 let (_dir, q) = queue();
3086 assert_eq!(q.revision(), 0, "an empty queue has no revision");
3087 let mut t = task("first");
3088 q.put(&mut t).unwrap();
3089 assert!(q.revision() > 0, "a written task moves the revision");
3090 }
3091
3092 #[test]
3093 fn revision_moves_when_deleting_an_older_task() {
3094 let (dir, q) = queue();
3095 let questions = Questions::at(dir.path().join("questions"));
3096 let mut t1 = task("older");
3097 q.put(&mut t1).unwrap();
3098 std::thread::sleep(std::time::Duration::from_millis(10));
3100 let mut t2 = task("newer");
3101 q.put(&mut t2).unwrap();
3102
3103 let rev_before = q.revision();
3104 q.remove(&t1.id, false, &questions).unwrap();
3105 let rev_after = q.revision();
3106
3107 assert_ne!(
3108 rev_before, rev_after,
3109 "deleting an older task must change the revision so other clients see the deletion"
3110 );
3111 }
3112
3113 #[test]
3114 fn removing_a_task_takes_it_out_of_the_listing() {
3115 let (dir, q) = queue();
3116 let questions = Questions::at(dir.path().join("questions"));
3117 let mut t = task("delete me");
3118 q.put(&mut t).unwrap();
3119 let removed = q.remove(t.short(), false, &questions).unwrap();
3120 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
3121 assert!(removed.quarantined.is_empty(), "nothing was blocked on it");
3122 assert!(q.list().is_empty());
3123 assert!(
3124 q.remove(&t.id, false, &questions).is_err(),
3125 "removing twice is an error"
3126 );
3127 }
3128
3129 #[test]
3130 fn removing_a_task_takes_its_stale_lock_with_it() {
3131 let (dir, q) = queue();
3132 let questions = Questions::at(dir.path().join("questions"));
3133 let mut t = task("interrupted");
3134 q.put(&mut t).unwrap();
3135
3136 let claim = q.claim(&t.id).unwrap();
3139 std::mem::forget(claim);
3140 assert!(
3141 q.claim(&t.id).is_err(),
3142 "the orphaned lock is what makes the task look claimed"
3143 );
3144
3145 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
3147 assert!(err.contains("live daemon"), "{err}");
3148 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
3149
3150 q.remove(&t.id, false, &questions).unwrap();
3152 assert!(q.list().is_empty());
3153 let mut again = task("interrupted");
3154 again.id = t.id.clone();
3155 q.put(&mut again).unwrap();
3156 assert!(
3157 q.claim(&t.id).is_ok(),
3158 "a task that comes back must be claimable, which a left-behind lock would prevent"
3159 );
3160 }
3161
3162 #[test]
3163 fn removing_a_task_quarantines_what_was_blocked_on_it() {
3164 let (dir, q) = queue();
3165 let questions = Questions::at(dir.path().join("questions"));
3166
3167 let mut dep = task("dependency");
3168 q.put(&mut dep).unwrap();
3169
3170 let mut still_valid = task("still valid");
3171 q.put(&mut still_valid).unwrap();
3172
3173 let mut blocked = task("waiting");
3174 blocked.block(
3175 vec![dep.id.clone(), still_valid.id.clone()],
3176 Some("waits on both".to_owned()),
3177 );
3178 q.put(&mut blocked).unwrap();
3179
3180 let removed = q.remove(&dep.id, false, &questions).unwrap();
3181 assert_eq!(removed.quarantined, [blocked.id.clone()]);
3182
3183 let after = q.get(&blocked.id).unwrap();
3184 assert_eq!(after.status, TaskStatus::Held);
3185 assert_eq!(after.hold_source, Some(HoldSource::Machine));
3186 assert!(after.blocked_by.is_empty());
3187 let reason = after.hold_reason.as_deref().unwrap_or_default();
3188 assert!(reason.contains(&dep.id), "{reason}");
3189 assert!(
3190 reason.contains(&still_valid.id),
3191 "the still-valid dependency must survive in the reason text: {reason}"
3192 );
3193 }
3194
3195 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3196 let p = dir.join(name);
3197 std::fs::write(&p, body).unwrap();
3198 p
3199 }
3200
3201 #[test]
3202 fn an_attachment_copy_survives_deleting_its_source() {
3203 let (dir, q) = queue();
3204 let src = source_file(dir.path(), "shot.png", "pixels");
3205 let mut t = task("with a picture");
3206 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3207 q.put(&mut t).unwrap();
3208 std::fs::remove_file(&src).unwrap();
3209 assert_eq!(names, ["shot.png"]);
3210 let loaded = q.get(&t.id).unwrap();
3211 let paths = q.attachment_paths(&loaded);
3212 assert_eq!(paths.len(), 1);
3213 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3214 }
3215
3216 #[test]
3217 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3218 let (dir, q) = queue();
3219 let mut t = task("bad names");
3220 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3221 let src = source_file(dir.path(), name, "x");
3222 assert!(
3223 q.attach(&mut t, &[src]).is_err(),
3224 "`{name}` must be refused"
3225 );
3226 }
3227 assert!(!crate::ask::valid_asset_name("C:foo.png"));
3231 #[cfg(not(windows))]
3232 {
3233 let src = source_file(dir.path(), "C:foo.png", "x");
3234 assert!(q.attach(&mut t, &[src]).is_err());
3235 }
3236 let long = format!("{}.png", "a".repeat(70));
3237 let src = source_file(dir.path(), &long, "x");
3238 assert!(q.attach(&mut t, &[src]).is_err());
3239 assert!(t.attachments.is_empty());
3240 assert!(!q.attachments_dir(&t.id).exists());
3241 }
3242
3243 #[test]
3244 fn a_taken_attachment_name_is_numbered_not_overwritten() {
3245 let (dir, q) = queue();
3246 let a = source_file(dir.path(), "shot.png", "one");
3247 let sub = dir.path().join("other");
3248 std::fs::create_dir_all(&sub).unwrap();
3249 let b = source_file(&sub, "shot.png", "two");
3250 let mut t = task("collision");
3251 q.attach(&mut t, &[a]).unwrap();
3252 q.attach(&mut t, &[b]).unwrap();
3253 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3254 let paths = q.attachment_paths(&t);
3255 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3256 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3257 }
3258
3259 #[test]
3260 fn a_renumbered_name_stays_inside_the_length_bound() {
3261 let (dir, q) = queue();
3262 let name = format!("{}.png", "a".repeat(60));
3263 assert_eq!(name.len(), 64);
3264 let a = source_file(dir.path(), &name, "one");
3265 let sub = dir.path().join("other");
3266 std::fs::create_dir_all(&sub).unwrap();
3267 let b = source_file(&sub, &name, "two");
3268 let mut t = task("long");
3269 q.attach(&mut t, &[a, b]).unwrap();
3270 assert_eq!(t.attachments.len(), 2);
3271 assert!(
3272 t.attachments
3273 .iter()
3274 .all(|n| crate::ask::valid_asset_name(n))
3275 );
3276 assert!(t.attachments[1].ends_with("-2.png"));
3277 }
3278
3279 #[test]
3280 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3281 let (dir, q) = queue();
3282 let good = source_file(dir.path(), "good.png", "ok");
3283 let mut t = task("partial");
3284 q.attach(&mut t, &[good]).unwrap();
3285 let more = source_file(dir.path(), "more.png", "ok");
3286 let missing = dir.path().join("missing.png");
3287 assert!(q.attach(&mut t, &[more, missing]).is_err());
3288 assert_eq!(t.attachments, ["good.png"]);
3289 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3290 .unwrap()
3291 .flatten()
3292 .collect();
3293 assert_eq!(on_disk.len(), 1);
3294 }
3295
3296 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3298 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3299 std::fs::create_dir_all(&tmp).unwrap();
3300 tmp
3301 }
3302
3303 #[test]
3304 fn a_failed_put_leaves_no_new_attachment_directory() {
3305 let (dir, q) = queue();
3306 let mut t = task("fresh");
3307 let tmp = block_put(&q, &t);
3308 let src = source_file(dir.path(), "shot.png", "x");
3309 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3310 assert!(t.attachments.is_empty());
3311 assert!(!q.attachments_dir(&t.id).exists());
3312 assert!(!q.path_of(&t.id).exists());
3313 std::fs::remove_dir(tmp).unwrap();
3314 }
3315
3316 #[test]
3317 fn a_failed_put_removes_only_the_copy_it_just_made() {
3318 let (dir, q) = queue();
3319 let mut t = task("edited");
3320 let first = source_file(dir.path(), "first.png", "1");
3321 q.attach_and_put(&mut t, &[first]).unwrap();
3322 block_put(&q, &t);
3323 let second = source_file(dir.path(), "second.png", "2");
3324 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3325 assert_eq!(t.attachments, ["first.png"]);
3326 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3327 .unwrap()
3328 .flatten()
3329 .map(|e| e.file_name().to_string_lossy().into_owned())
3330 .collect();
3331 assert_eq!(on_disk, ["first.png"]);
3332 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3333 }
3334
3335 #[test]
3336 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3337 let (dir, q) = queue();
3338 let questions = Questions::at(dir.path().join("questions"));
3339 let gone = task("gone");
3340 let mut other = task("other");
3341 let mut live = task("live");
3342 q.put(&mut other).unwrap();
3343 q.put(&mut live).unwrap();
3344 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3347 std::fs::create_dir_all(&orphan).unwrap();
3348 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3349 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3352 std::fs::create_dir_all(&busy).unwrap();
3353
3354 q.remove(&other.id, false, &questions).unwrap();
3355 assert!(!orphan.exists(), "an orphan is swept");
3356 assert!(busy.exists(), "a removal in progress is left alone");
3357 }
3358
3359 #[test]
3360 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3361 let (dir, q) = queue();
3362 let questions = Questions::at(dir.path().join("questions"));
3363 let mut t = task("stuck");
3364 let src = source_file(dir.path(), "shot.png", "x");
3365 q.attach_and_put(&mut t, &[src]).unwrap();
3366 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3369 std::fs::create_dir_all(&aside).unwrap();
3370 std::fs::write(aside.join("old.png"), "o").unwrap();
3371 assert!(q.remove(&t.id, false, &questions).is_err());
3372 assert!(q.path_of(&t.id).exists());
3373 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3374 }
3375
3376 #[test]
3377 fn a_failed_record_removal_puts_the_attachments_back() {
3378 let (dir, q) = queue();
3379 let mut t = task("rollback");
3380 let src = source_file(dir.path(), "shot.png", "x");
3381 q.attach_and_put(&mut t, &[src]).unwrap();
3382 let err = q
3383 .remove_record_with_attachments(&t.id, |_| {
3384 Err(std::io::Error::other("injected failure"))
3385 })
3386 .unwrap_err();
3387 assert!(format!("{err:#}").contains("injected failure"));
3388 assert!(q.path_of(&t.id).exists());
3389 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3390 assert!(
3391 !q.root
3392 .join(format!("{}.attachments.removing", t.id))
3393 .exists()
3394 );
3395 }
3396
3397 #[test]
3398 fn editing_a_task_keeps_its_attachments() {
3399 let (dir, q) = queue();
3400 let src = source_file(dir.path(), "shot.png", "x");
3401 let mut t = task("editable");
3402 q.attach(&mut t, &[src]).unwrap();
3403 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3404 q.put(&mut t).unwrap();
3405 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3406 }
3407
3408 #[test]
3409 fn removing_a_task_deletes_its_attachments() {
3410 let (dir, q) = queue();
3411 let questions = Questions::at(dir.path().join("questions"));
3412 let src = source_file(dir.path(), "shot.png", "x");
3413 let mut t = task("doomed");
3414 q.attach(&mut t, &[src]).unwrap();
3415 q.put(&mut t).unwrap();
3416 assert!(q.attachments_dir(&t.id).is_dir());
3417 q.remove(&t.id, false, &questions).unwrap();
3418 assert!(!q.attachments_dir(&t.id).exists());
3419 assert!(q.list().is_empty());
3420 }
3421
3422 #[test]
3423 fn a_task_written_before_attachments_still_reads() {
3424 let (_dir, q) = queue();
3425 let mut t = task("old");
3426 q.put(&mut t).unwrap();
3427 let path = q.path_of(&t.id);
3428 let mut v: serde_json::Value =
3429 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3430 v.as_object_mut().unwrap().remove("attachments");
3431 std::fs::write(&path, v.to_string()).unwrap();
3432 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3433 }
3434
3435 #[test]
3436 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3437 let q = Queue::at(PathBuf::from("relative-queue"));
3438 let mut t = task("rel");
3439 t.attachments.push("shot.png".to_owned());
3440 let paths = q.attachment_paths(&t);
3441 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3442 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3443 }
3444
3445 #[test]
3446 fn link_run_adds_a_run_once_and_touches_nothing_else() {
3447 let dir = tempfile::tempdir().unwrap();
3448 let queue = Queue::at(dir.path().join("queue"));
3449 let mut t = Task::new(
3450 "t".to_owned(),
3451 "do it".to_owned(),
3452 PathBuf::from("."),
3453 Source::Human,
3454 );
3455 queue.put(&mut t).unwrap();
3456 let before = queue.get(&t.id).unwrap();
3457
3458 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3459 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3460 assert_eq!(linked.status, before.status);
3461 assert_eq!(linked.attempts, before.attempts);
3462 assert_eq!(linked.interrupt, before.interrupt);
3463
3464 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3465 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3466 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3467 assert!(queue.link_run("no-such-task", "r").is_err());
3468 }
3469
3470 #[test]
3471 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
3472 let dir = tempfile::tempdir().unwrap();
3473 let queue = Queue::at(dir.path().join("queue"));
3474 let mut t = Task::new(
3475 "t".to_owned(),
3476 "do it".to_owned(),
3477 PathBuf::from("."),
3478 Source::Human,
3479 );
3480 queue.put(&mut t).unwrap();
3481 let mut snapshot = queue.get(&t.id).unwrap();
3483 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3484
3485 snapshot.start("20260930-000000-aaaa".to_owned());
3486 queue.put(&mut snapshot).unwrap();
3487
3488 let stored = queue.get(&t.id).unwrap();
3489 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3490 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3491 assert_eq!(stored.attempts, 1);
3492 }
3493
3494 #[test]
3495 fn concurrent_links_and_daemon_saves_lose_nothing() {
3496 let dir = tempfile::tempdir().unwrap();
3497 let queue = Queue::at(dir.path().join("queue"));
3498 let mut t = Task::new(
3499 "t".to_owned(),
3500 "do it".to_owned(),
3501 PathBuf::from("."),
3502 Source::Human,
3503 );
3504 queue.put(&mut t).unwrap();
3505 let id = t.id.clone();
3506
3507 let linkers: Vec<_> = (0..4)
3508 .map(|n| {
3509 let (queue, id) = (queue.clone(), id.clone());
3510 std::thread::spawn(move || {
3511 for k in 0..10 {
3512 queue
3513 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3514 .unwrap();
3515 }
3516 })
3517 })
3518 .collect();
3519 let mut mine = queue.get(&id).unwrap();
3522 for k in 0..10 {
3523 mine.start(format!("20260930-000009-d{k:03}"));
3524 queue.put(&mut mine).unwrap();
3525 }
3526 for l in linkers {
3527 l.join().unwrap();
3528 }
3529
3530 let stored = queue.get(&id).unwrap();
3531 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
3532 assert_eq!(
3533 stored.attempts, 10,
3534 "linking never rewinds the daemon's work"
3535 );
3536 }
3537
3538 fn age_lock(path: &Path) {
3539 let f = std::fs::OpenOptions::new().write(true).open(path).unwrap();
3540 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
3541 .unwrap();
3542 }
3543
3544 #[test]
3545 fn concurrent_stale_takeover_yields_one_holder() {
3546 use std::sync::atomic::{AtomicUsize, Ordering};
3547 use std::sync::{Arc, Barrier};
3548 for _ in 0..5 {
3549 let dir = tempfile::tempdir().unwrap();
3550 let q = Queue::at(dir.path().to_path_buf());
3551 let lock = dir.path().join("t.write-lock");
3552 std::fs::write(&lock, "dead-0000").unwrap();
3553 age_lock(&lock);
3554 let n = 6;
3555 let barrier = Arc::new(Barrier::new(n));
3556 let (now, max) = (Arc::new(AtomicUsize::new(0)), Arc::new(AtomicUsize::new(0)));
3557 let handles: Vec<_> = (0..n)
3558 .map(|_| {
3559 let (q, b, now, max) = (q.clone(), barrier.clone(), now.clone(), max.clone());
3560 std::thread::spawn(move || {
3561 b.wait();
3562 let g = q.lock_task("t").unwrap();
3563 let held = now.fetch_add(1, Ordering::SeqCst) + 1;
3564 max.fetch_max(held, Ordering::SeqCst);
3565 std::thread::sleep(std::time::Duration::from_millis(20));
3566 now.fetch_sub(1, Ordering::SeqCst);
3567 drop(g);
3568 })
3569 })
3570 .collect();
3571 for h in handles {
3572 h.join().unwrap();
3573 }
3574 assert_eq!(max.load(Ordering::SeqCst), 1);
3575 assert!(!lock.exists());
3576 }
3577 }
3578
3579 #[test]
3580 fn dropping_a_stolen_lock_leaves_the_new_holders_lock() {
3581 let dir = tempfile::tempdir().unwrap();
3582 let q = Queue::at(dir.path().to_path_buf());
3583 let lock = dir.path().join("t.write-lock");
3584 let a = q.lock_task("t").unwrap();
3585 age_lock(&lock);
3586 let b = q.lock_task("t").unwrap();
3587 assert_ne!(a.token, b.token);
3588 drop(a);
3589 assert_eq!(std::fs::read_to_string(&lock).unwrap(), b.token);
3590 drop(b);
3591 assert!(!lock.exists());
3592 }
3593
3594 #[test]
3595 fn break_stale_leaves_a_lock_that_replaced_the_one_judged() {
3596 let dir = tempfile::tempdir().unwrap();
3597 let lock = dir.path().join("t.write-lock");
3598 std::fs::write(&lock, "new-token").unwrap();
3599 assert!(!break_stale(&lock, "old-token"));
3600 assert_eq!(std::fs::read_to_string(&lock).unwrap(), "new-token");
3601 assert!(!break_stale(&lock, "new-token"));
3603 assert!(lock.exists());
3604 age_lock(&lock);
3605 assert!(break_stale(&lock, "new-token"));
3606 assert!(!lock.exists());
3607 }
3608
3609 #[test]
3610 fn a_stale_break_marker_is_recovered_and_a_fresh_one_is_respected() {
3611 let dir = tempfile::tempdir().unwrap();
3612 let lock = dir.path().join("t.write-lock");
3613 std::fs::write(&lock, "dead-1").unwrap();
3614 age_lock(&lock);
3615 let marker = break_marker(&lock, "dead-1");
3616 std::fs::write(&marker, "crashed-remover").unwrap();
3617 assert!(!break_stale(&lock, "dead-1"));
3619 assert!(lock.exists());
3620 age_lock(&marker);
3622 assert!(break_stale(&lock, "dead-1"));
3623 assert!(!lock.exists());
3624 assert!(!marker.exists());
3625 }
3626}