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
142pub const CHAT_NODE: &str = "chat";
145
146impl Source {
147 pub fn label(&self) -> String {
149 match self {
150 Self::Human => "human".to_owned(),
151 Self::Agent { run, node } => format!("{node}@{}", short(run)),
152 Self::Issue { number, .. } => format!("issue #{number}"),
153 }
154 }
155}
156
157#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
159#[serde(rename_all = "lowercase")]
160pub enum TaskStatus {
161 Queued,
163 Running,
165 Done,
167 Failed,
169 Held,
171 Blocked,
175}
176
177impl TaskStatus {
178 pub fn runnable(self) -> bool {
180 matches!(self, Self::Queued | Self::Failed)
181 }
182
183 pub fn as_str(self) -> &'static str {
185 match self {
186 Self::Queued => "queued",
187 Self::Running => "running",
188 Self::Done => "done",
189 Self::Failed => "failed",
190 Self::Held => "held",
191 Self::Blocked => "blocked",
192 }
193 }
194}
195
196#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
199pub struct TaskCounts {
200 pub queued: usize,
202 pub running: usize,
204 pub done: usize,
206 pub failed: usize,
208 pub held: usize,
210 pub blocked: usize,
212}
213
214impl TaskCounts {
215 pub fn of(tasks: &[Task]) -> Self {
219 let mut counts = Self::default();
220 for t in tasks {
221 match t.status {
222 TaskStatus::Queued => counts.queued += 1,
223 TaskStatus::Running => counts.running += 1,
224 TaskStatus::Done => counts.done += 1,
225 TaskStatus::Failed => counts.failed += 1,
226 TaskStatus::Held => counts.held += 1,
227 TaskStatus::Blocked => counts.blocked += 1,
228 }
229 }
230 counts
231 }
232}
233
234#[derive(Debug, Clone, Serialize, Deserialize)]
236#[serde(deny_unknown_fields)]
237pub struct Task {
238 pub schema: u32,
240 pub id: String,
242 pub title: String,
244 pub instruction: String,
246 pub repo: PathBuf,
248 pub source: Source,
250 #[serde(default)]
252 pub priority: i32,
253 #[serde(default)]
264 pub solo: bool,
265 pub status: TaskStatus,
267 #[serde(default)]
269 pub attempts: usize,
270 #[serde(default)]
272 pub runs: Vec<String>,
273 #[serde(default)]
275 pub last_error: Option<String>,
276 #[serde(default)]
289 pub hold_reason: Option<String>,
290 #[serde(default)]
293 pub hold_source: Option<HoldSource>,
294 #[serde(default)]
306 pub diagnostic: Option<String>,
307 #[serde(default)]
317 pub blocked_by: Vec<String>,
318 #[serde(default)]
321 pub block_reason: Option<String>,
322 #[serde(default)]
335 pub blocked_from: Option<TaskStatus>,
336 #[serde(default)]
347 pub answers: Vec<AnsweredQuestion>,
348 #[serde(default)]
354 pub triage_applied: Vec<String>,
355 #[serde(default)]
361 pub actions_applied: Vec<String>,
362 #[serde(default)]
367 pub resume_override: Option<OperatorResume>,
368 #[serde(default)]
380 pub review_branch: Option<String>,
381 #[serde(default)]
384 pub fresh_start: bool,
385 #[serde(default)]
396 pub interrupt: bool,
397 #[serde(default)]
420 pub urgent: bool,
421 #[serde(default)]
427 pub attachments: Vec<String>,
428 #[serde(default)]
432 pub followup: Option<FollowUp>,
433 #[serde(default)]
438 pub held_at: Option<Timestamp>,
439 pub created_at: Timestamp,
441 pub updated_at: Timestamp,
443}
444
445#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
447pub struct FollowUp {
448 pub run: String,
450 #[serde(default)]
452 pub origin_task: Option<String>,
453 pub pr: String,
455 pub findings: Vec<String>,
457 pub generation: u32,
460}
461
462#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
465pub struct OperatorResume {
466 pub question_id: String,
468 pub at: Timestamp,
470 #[serde(default)]
473 pub conductor_rehold: Option<String>,
474 #[serde(default)]
477 pub forced: bool,
478 #[serde(default)]
485 pub pinned_run: Option<String>,
486}
487
488#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
491pub struct AnsweredQuestion {
492 pub question: String,
494 pub answer: String,
496}
497
498impl Task {
499 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
501 let now = Timestamp::now();
502 Self {
503 schema: SCHEMA,
504 id: new_id(),
505 title,
506 instruction,
507 repo,
508 source,
509 priority: 0,
510 solo: false,
511 status: TaskStatus::Queued,
512 attempts: 0,
513 runs: Vec::new(),
514 last_error: None,
515 hold_reason: None,
516 hold_source: None,
517 diagnostic: None,
518 blocked_by: Vec::new(),
519 block_reason: None,
520 blocked_from: None,
521 answers: Vec::new(),
522 triage_applied: Vec::new(),
523 actions_applied: Vec::new(),
524 resume_override: None,
525 review_branch: None,
526 fresh_start: false,
527 interrupt: false,
528 urgent: false,
529 attachments: Vec::new(),
530 followup: None,
531 held_at: None,
532 created_at: now,
533 updated_at: now,
534 }
535 }
536
537 pub fn short(&self) -> &str {
539 short(&self.id)
540 }
541
542 pub fn mark_triage_applied(&mut self, question_id: &str) {
544 if !self.triage_applied(question_id) {
545 self.triage_applied.push(question_id.to_owned());
546 }
547 }
548
549 pub fn action_applied(&self, question_id: &str) -> bool {
552 self.actions_applied.iter().any(|id| id == question_id)
553 }
554
555 pub fn mark_action_applied(&mut self, question_id: &str) {
557 if !self.action_applied(question_id) {
558 self.actions_applied.push(question_id.to_owned());
559 }
560 }
561
562 pub fn triage_applied(&self, question_id: &str) -> bool {
564 self.triage_applied.iter().any(|id| id == question_id)
565 }
566
567 pub fn start(&mut self, run: String) {
578 self.status = TaskStatus::Running;
579 self.attempts += 1;
580 self.runs.push(run);
581 self.last_error = None;
582 self.fresh_start = false;
583 self.interrupt = false;
584 self.resume_override = None;
586 }
587
588 pub fn link_run(&mut self, run: &str) -> bool {
593 if self.runs.iter().any(|r| r == run) {
594 return false;
595 }
596 self.runs.push(run.to_owned());
597 true
598 }
599
600 pub fn succeed(&mut self) {
610 self.status = TaskStatus::Done;
611 self.resume_override = None;
612 self.last_error = None;
613 self.hold_reason = None;
614 self.hold_source = None;
615 self.diagnostic = None;
616 self.blocked_by.clear();
617 self.block_reason = None;
618 self.blocked_from = None;
619 }
620
621 pub fn already_landed(&mut self, note: impl Into<String>) {
629 self.succeed();
630 self.attempts = self.attempts.saturating_sub(1);
631 self.last_error = Some(note.into());
632 }
633
634 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
664 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
665 return &[];
666 }
667 &self.runs[..self.runs.len() - 1]
668 }
669
670 pub fn successor_of(&self, run: &str) -> Option<&String> {
677 let pos = self.runs.iter().position(|r| r == run)?;
678 self.runs.get(pos + 1)
679 }
680
681 pub fn earlier_attempts(&self) -> &[String] {
688 &self.runs
689 }
690
691 fn note_held(&mut self) {
694 if self.status != TaskStatus::Held || self.held_at.is_none() {
695 self.held_at = Some(Timestamp::now());
696 }
697 }
698
699 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
715 let why = why.into();
716 self.diagnostic = None;
717 self.status = if self.attempts >= max_attempts {
718 self.note_held();
719 self.hold_source = Some(HoldSource::Machine);
720 self.hold_reason = Some(why.clone());
721 TaskStatus::Held
722 } else {
723 TaskStatus::Failed
724 };
725 self.last_error = Some(why);
726 }
727
728 pub fn stall(&mut self, why: impl Into<String>) {
738 self.last_error = Some(why.into());
739 self.diagnostic = None;
740 self.attempts = self.attempts.saturating_sub(1);
741 self.status = TaskStatus::Failed;
742 }
743
744 pub fn operator_held(&self) -> bool {
750 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
751 }
752
753 pub fn hold_manual(&mut self, reason: Option<String>) {
763 self.note_held();
764 self.status = TaskStatus::Held;
765 if reason.is_some() {
766 self.hold_reason = reason;
767 }
768 self.hold_source = Some(HoldSource::Manual);
769 self.blocked_by.clear();
770 self.block_reason = None;
771 self.blocked_from = None;
772 }
773
774 pub fn hold_machine(&mut self, reason: Option<String>) {
779 self.note_held();
780 self.status = TaskStatus::Held;
781 if reason.is_some() {
782 self.hold_reason = reason;
783 }
784 self.hold_source = Some(HoldSource::Machine);
785 self.blocked_by.clear();
786 self.block_reason = None;
787 self.blocked_from = None;
788 }
789
790 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
800 if self.status != TaskStatus::Blocked {
801 self.blocked_from = Some(self.status);
802 }
803 self.status = TaskStatus::Blocked;
804 self.blocked_by = blocked_by;
805 self.block_reason = reason;
806 }
807
808 pub fn unblock(&mut self, resolved_id: &str) {
831 if self.status != TaskStatus::Blocked {
832 return;
833 }
834 self.blocked_by.retain(|id| id != resolved_id);
835 self.restore_if_unblocked();
836 }
837
838 pub fn dependency_deleted(&mut self, deleted_id: &str) -> bool {
849 if self.status != TaskStatus::Blocked || !self.blocked_by.iter().any(|b| b == deleted_id) {
850 return false;
851 }
852 self.blocked_by.retain(|id| id != deleted_id);
853 self.restore_if_unblocked();
854 true
855 }
856
857 fn restore_if_unblocked(&mut self) {
861 if self.blocked_by.is_empty() {
862 self.status = match self.blocked_from {
863 Some(TaskStatus::Running) => TaskStatus::Queued,
864 Some(other) => other,
865 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
866 TaskStatus::Held
867 }
868 None => TaskStatus::Queued,
869 };
870 self.block_reason = None;
871 self.blocked_from = None;
872 }
873 }
874
875 pub fn record_answer(&mut self, question: String, answer: String) {
880 self.answers.push(AnsweredQuestion { question, answer });
881 }
882
883 pub fn request_review(&mut self, branch: String) {
887 self.release();
888 self.review_branch = Some(branch);
889 }
890
891 pub fn requeue(&mut self) {
894 self.release();
895 self.review_branch = None;
897 self.fresh_start = true;
898 }
899
900 pub fn hold_for_handover(&mut self, branch: Option<String>, reason: String) {
909 if branch.is_some() {
910 self.review_branch = branch;
911 }
912 self.hold_machine(Some(reason));
913 }
914
915 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
924 if self.status == TaskStatus::Running {
925 bail!(
926 "task {} is running; its priority cannot be changed until \
927 this attempt finishes",
928 self.short()
929 );
930 }
931 self.priority = priority;
932 Ok(())
933 }
934
935 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
950 if interrupt && !self.status.runnable() {
951 bail!(
952 "task {} is {}; only a queued or failed task can be marked \
953 to interrupt",
954 self.short(),
955 self.status.as_str()
956 );
957 }
958 self.interrupt = interrupt;
959 Ok(())
960 }
961
962 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
974 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
975 bail!(
976 "task {} is {}; only a queued or held task's instruction can \
977 be edited",
978 self.short(),
979 self.status.as_str()
980 );
981 }
982 self.title = title;
983 self.instruction = instruction;
984 Ok(())
985 }
986
987 pub fn handed_off(&mut self, why: impl Into<String>) {
1007 let why = why.into();
1008 self.diagnostic = None;
1009 self.note_held();
1010 self.status = TaskStatus::Held;
1011 self.hold_source = Some(HoldSource::Machine);
1012 self.hold_reason = Some(why.clone());
1013 self.last_error = Some(why);
1014 }
1015
1016 pub fn release(&mut self) {
1020 let refused_handover = self.status == TaskStatus::Held
1021 && self.hold_source == Some(HoldSource::Machine)
1022 && self.review_branch.is_some();
1023 self.status = TaskStatus::Queued;
1024 self.held_at = None;
1025 self.attempts = 0;
1026 self.last_error = None;
1027 self.hold_reason = None;
1030 self.hold_source = None;
1031 self.diagnostic = None;
1032 self.blocked_by.clear();
1037 self.block_reason = None;
1038 self.blocked_from = None;
1039 if !refused_handover {
1043 self.review_branch = None;
1044 }
1045 self.fresh_start = false;
1046 }
1047}
1048
1049fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
1053 use std::io::ErrorKind;
1054 let (stem, ext) = match name.rfind('.') {
1055 Some(i) if i > 0 => (&name[..i], &name[i..]),
1056 _ => (name, ""),
1057 };
1058 for n in 1u32.. {
1059 let candidate = if n == 1 {
1060 name.to_owned()
1061 } else {
1062 let suffix = format!("-{n}");
1063 let room = 64usize.saturating_sub(suffix.len() + ext.len());
1064 let stem: String = stem.chars().take(room).collect();
1065 format!("{stem}{suffix}{ext}")
1066 };
1067 if !crate::ask::valid_asset_name(&candidate) {
1068 bail!("no valid attachment name is left for `{name}`");
1069 }
1070 let path = dir.join(&candidate);
1071 match std::fs::OpenOptions::new()
1072 .write(true)
1073 .create_new(true)
1074 .open(&path)
1075 {
1076 Ok(mut out) => {
1077 let copied = std::fs::File::open(src)
1078 .and_then(|mut input| std::io::copy(&mut input, &mut out));
1079 if let Err(e) = copied {
1080 drop(out);
1081 let _ = std::fs::remove_file(&path);
1082 return Err(e).with_context(|| format!("copy {}", src.display()));
1083 }
1084 return Ok((candidate, path));
1085 }
1086 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
1087 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
1088 }
1089 }
1090 unreachable!("the counter never runs out")
1091}
1092
1093const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
1096
1097struct TaskLock {
1100 path: PathBuf,
1101 token: String,
1102}
1103
1104fn owner_token() -> String {
1107 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1108 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1109 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy() ^ n.rotate_left(32));
1110 format!("{}-{:016x}", std::process::id(), r.next_u64())
1111}
1112
1113fn break_marker(path: &Path, token: &str) -> PathBuf {
1116 let mut h: u64 = 0xcbf29ce484222325;
1117 for b in token.bytes() {
1118 h ^= u64::from(b);
1119 h = h.wrapping_mul(0x100000001b3);
1120 }
1121 let mut name = path.as_os_str().to_owned();
1122 name.push(format!(".break-{h:016x}"));
1123 PathBuf::from(name)
1124}
1125
1126fn older_than_stale(path: &Path) -> bool {
1127 std::fs::metadata(path)
1128 .and_then(|m| m.modified())
1129 .ok()
1130 .and_then(|t| t.elapsed().ok())
1131 .is_some_and(|age| age > TASK_LOCK_STALE)
1132}
1133
1134const MAX_MARKER_DEPTH: u8 = 3;
1138
1139fn take_marker(path: &Path, token: &str, depth: u8) -> Option<(PathBuf, String)> {
1145 use std::io::Write;
1146 let marker = break_marker(path, token);
1147 for _ in 0..2 {
1148 match std::fs::OpenOptions::new()
1149 .write(true)
1150 .create_new(true)
1151 .open(&marker)
1152 {
1153 Ok(mut f) => {
1154 let mine = owner_token();
1155 if f.write_all(mine.as_bytes()).is_err() {
1156 drop(f);
1157 let _ = std::fs::remove_file(&marker);
1158 return None;
1159 }
1160 return Some((marker, mine));
1161 }
1162 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1163 let judged = std::fs::read_to_string(&marker).ok();
1164 match judged.filter(|_| older_than_stale(&marker)) {
1165 Some(judged) if depth < MAX_MARKER_DEPTH => {
1166 remove_lock_if(&marker, &judged, true, depth + 1)?;
1167 }
1168 _ => return None,
1169 }
1170 }
1171 Err(_) => return None,
1172 }
1173 }
1174 None
1175}
1176
1177fn remove_lock_if(path: &Path, token: &str, require_stale: bool, depth: u8) -> Option<bool> {
1182 let (marker, mine) = take_marker(path, token, depth)?;
1183 let still = std::fs::read_to_string(path).is_ok_and(|c| c == token)
1184 && (!require_stale || older_than_stale(path));
1185 if still {
1186 let _ = std::fs::remove_file(path);
1187 }
1188 if depth >= MAX_MARKER_DEPTH {
1190 if std::fs::read_to_string(&marker).is_ok_and(|c| c == mine) {
1191 let _ = std::fs::remove_file(&marker);
1192 }
1193 } else {
1194 let _ = remove_lock_if(&marker, &mine, false, depth + 1);
1195 }
1196 Some(still)
1197}
1198
1199fn break_stale(path: &Path, judged: &str) -> bool {
1202 remove_lock_if(path, judged, true, 0).unwrap_or(false)
1203}
1204
1205impl Drop for TaskLock {
1206 fn drop(&mut self) {
1207 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
1208 loop {
1209 match remove_lock_if(&self.path, &self.token, false, 0) {
1210 Some(_) => return,
1211 None if std::time::Instant::now() > deadline => return,
1212 None => std::thread::sleep(std::time::Duration::from_millis(5)),
1213 }
1214 }
1215 }
1216}
1217
1218#[derive(Debug, Clone)]
1220pub struct Queue {
1221 root: PathBuf,
1222}
1223
1224impl Queue {
1225 pub fn open() -> Self {
1227 Self::at(crate::run::home().join("queue"))
1228 }
1229
1230 pub fn at(root: PathBuf) -> Self {
1233 Self { root }
1234 }
1235
1236 pub fn root(&self) -> &Path {
1238 &self.root
1239 }
1240
1241 pub fn path_of(&self, id: &str) -> PathBuf {
1243 self.root.join(format!("{id}.json"))
1244 }
1245
1246 pub fn attachments_dir(&self, id: &str) -> PathBuf {
1249 self.root.join(format!("{id}.attachments"))
1250 }
1251
1252 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1263 let mut wanted = Vec::new();
1264 for src in sources {
1265 let name = src
1266 .file_name()
1267 .and_then(|n| n.to_str())
1268 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1269 if !crate::ask::valid_asset_name(name) {
1270 bail!(
1271 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1272 with no `..`; rename the file and try again"
1273 );
1274 }
1275 if !src.is_file() {
1276 bail!("attachment `{}` is not a file", src.display());
1277 }
1278 wanted.push((src, name));
1279 }
1280 let dir = self.attachments_dir(&task.id);
1281 let existed = dir.is_dir();
1282 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1283 let mut created: Vec<PathBuf> = Vec::new();
1284 let mut names = Vec::new();
1285 let mut copy_all = || -> Result<()> {
1286 for (src, name) in &wanted {
1287 let (stored, path) = copy_new(&dir, src, name)?;
1288 created.push(path);
1289 names.push(stored);
1290 }
1291 Ok(())
1292 };
1293 if let Err(e) = copy_all() {
1294 for path in &created {
1295 let _ = std::fs::remove_file(path);
1296 }
1297 if !existed {
1298 let _ = std::fs::remove_dir(&dir);
1299 }
1300 return Err(e);
1301 }
1302 task.attachments.extend(names.iter().cloned());
1303 Ok(names)
1304 }
1305
1306 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1317 let dir = self.attachments_dir(&task.id);
1318 let existed = dir.is_dir();
1319 let before = task.attachments.len();
1320 let names = self.attach(task, sources)?;
1321 if let Err(e) = self.put(task) {
1322 for name in &names {
1323 let _ = std::fs::remove_file(dir.join(name));
1324 }
1325 if !existed {
1326 let _ = std::fs::remove_dir(&dir);
1327 }
1328 task.attachments.truncate(before);
1329 return Err(e);
1330 }
1331 Ok(names)
1332 }
1333
1334 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1338 let dir = self.attachments_dir(&task.id);
1339 task.attachments
1340 .iter()
1341 .map(|n| {
1342 let p = dir.join(n);
1343 std::path::absolute(&p).unwrap_or(p)
1344 })
1345 .collect()
1346 }
1347
1348 pub fn put(&self, task: &mut Task) -> Result<()> {
1357 let _lock = self.lock_task(&task.id)?;
1358 self.put_unlocked(task)
1359 }
1360
1361 pub fn create_new(&self, task: &mut Task) -> Result<bool> {
1366 let _lock = self.lock_task(&task.id)?;
1367 if self.path_of(&task.id).exists() {
1368 return Ok(false);
1369 }
1370 self.put_unlocked(task)?;
1371 Ok(true)
1372 }
1373
1374 fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1376 if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1377 for run in stored.runs {
1378 if !task.runs.contains(&run) {
1379 task.runs.push(run);
1380 }
1381 }
1382 }
1383 task.updated_at = Timestamp::now();
1384 std::fs::create_dir_all(&self.root)
1385 .with_context(|| format!("create {}", self.root.display()))?;
1386 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1387 let path = self.path_of(&task.id);
1388 let tmp = path.with_extension("json.tmp");
1389 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1390 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1391 if let (Some(notice), Some(home)) = (
1396 crate::notices::task_held(task),
1397 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1398 ) {
1399 crate::notices::raise_in(home, notice);
1400 }
1401 Ok(())
1402 }
1403
1404 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1409 let id = self.resolve_id(id)?;
1410 let _lock = self.lock_task(&id)?;
1413 let mut task = self.get(&id)?;
1414 if task.link_run(run) {
1415 self.put_unlocked(&mut task)?;
1416 }
1417 Ok(task)
1418 }
1419
1420 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1428 std::fs::create_dir_all(&self.root)
1429 .with_context(|| format!("create {}", self.root.display()))?;
1430 let path = self.root.join(format!("{id}.write-lock"));
1431 let started = std::time::Instant::now();
1432 loop {
1433 match std::fs::OpenOptions::new()
1434 .write(true)
1435 .create_new(true)
1436 .open(&path)
1437 {
1438 Ok(mut file) => {
1439 use std::io::Write;
1440 let token = owner_token();
1441 if let Err(e) = file.write_all(token.as_bytes()) {
1442 drop(file);
1443 let _ = std::fs::remove_file(&path);
1444 return Err(e).with_context(|| format!("lock {}", path.display()));
1445 }
1446 return Ok(TaskLock { path, token });
1447 }
1448 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1449 let judged = std::fs::read_to_string(&path).ok();
1452 let broken = judged
1453 .filter(|_| older_than_stale(&path))
1454 .is_some_and(|judged| break_stale(&path, &judged));
1455 if !broken {
1459 if started.elapsed() > TASK_LOCK_STALE {
1460 bail!("could not lock task {id}");
1461 }
1462 std::thread::sleep(std::time::Duration::from_millis(15));
1463 }
1464 }
1465 Err(e)
1469 if e.kind() == std::io::ErrorKind::PermissionDenied
1470 && started.elapsed() <= TASK_LOCK_STALE =>
1471 {
1472 std::thread::sleep(std::time::Duration::from_millis(15));
1473 }
1474 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1475 }
1476 }
1477 }
1478
1479 pub fn get(&self, id: &str) -> Result<Task> {
1481 let resolved = self.resolve_id(id)?;
1482 read_path(&self.path_of(&resolved))
1483 }
1484
1485 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1517 let _ = questions;
1518 let resolved = self.resolve_id(id)?;
1519 if in_flight {
1520 bail!("task {resolved} is being run by a live daemon right now");
1521 }
1522 self.write_tombstone(&resolved)?;
1523 if let Err(e) = self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))
1524 {
1525 if self.path_of(&resolved).exists() {
1529 let _ = std::fs::remove_file(self.tombstone_path(&resolved));
1530 }
1531 return Err(e);
1532 }
1533 let (released, still_blocked) = self.release_dependents_of(&resolved);
1534 Ok(Removal {
1535 id: resolved,
1536 released,
1537 still_blocked,
1538 })
1539 }
1540
1541 fn tombstone_path(&self, id: &str) -> PathBuf {
1544 self.root.join(format!("{id}.removed"))
1545 }
1546
1547 fn write_tombstone(&self, id: &str) -> Result<()> {
1548 let path = self.tombstone_path(id);
1549 let tmp = self.root.join(format!("{id}.removed.tmp"));
1550 std::fs::write(&tmp, b"")
1551 .and_then(|()| std::fs::rename(&tmp, &path))
1552 .with_context(|| format!("write {}", path.display()))
1553 }
1554
1555 fn was_deleted(&self, id: &str) -> bool {
1559 self.tombstone_path(id).is_file() && !self.path_of(id).exists()
1560 }
1561
1562 pub fn apply_deleted_blockers(&self, task: &mut Task) -> Vec<String> {
1566 let deleted = deleted_blockers(self, &task.blocked_by);
1567 deleted
1568 .into_iter()
1569 .filter(|id| task.dependency_deleted(id))
1570 .collect()
1571 }
1572
1573 pub fn note_dependency_deleted(&self, dependent: &Task, deleted: &str) {
1577 let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty()) else {
1578 return;
1579 };
1580 let ja = crate::lang::is_japanese(&crate::lang::of_repo(&dependent.repo));
1581 let waiting = dependent.status == TaskStatus::Blocked;
1582 let message = match (ja, waiting) {
1583 (true, false) => format!(
1584 "タスク {} は、待っていた {} が削除されたため待機を解除し、元の状態に戻しました",
1585 dependent.short(),
1586 short(deleted)
1587 ),
1588 (true, true) => format!(
1589 "タスク {} は、待っていた {} が削除されたため、残りの依存を待っています",
1590 dependent.short(),
1591 short(deleted)
1592 ),
1593 (false, false) => format!(
1594 "Task {} stopped waiting on {} because it was deleted, and returned to its previous state",
1595 dependent.short(),
1596 short(deleted)
1597 ),
1598 (false, true) => format!(
1599 "Task {} stopped waiting on {} because it was deleted, and is still waiting on its other dependencies",
1600 dependent.short(),
1601 short(deleted)
1602 ),
1603 };
1604 crate::notices::raise_in(
1605 home,
1606 crate::notices::Notice::info(
1607 &format!("unblocked:{}:{}", dependent.id, deleted),
1608 message,
1609 ),
1610 );
1611 }
1612
1613 fn release_dependents_of(&self, dependency: &str) -> (Vec<String>, Vec<String>) {
1618 let (mut released, mut still_blocked) = (Vec::new(), Vec::new());
1619 for listed in self.list() {
1620 if listed.status != TaskStatus::Blocked
1621 || !listed.blocked_by.iter().any(|b| b == dependency)
1622 {
1623 continue;
1624 }
1625 let Ok(_claim) = self.claim(&listed.id) else {
1626 continue;
1627 };
1628 let Ok(mut task) = self.get(&listed.id) else {
1629 continue;
1630 };
1631 if !task.dependency_deleted(dependency) {
1632 continue;
1633 }
1634 if self.put(&mut task).is_err() {
1635 continue;
1636 }
1637 self.note_dependency_deleted(&task, dependency);
1638 if task.status == TaskStatus::Blocked {
1639 still_blocked.push(task.id.clone());
1640 } else {
1641 released.push(task.id.clone());
1642 }
1643 }
1644 (released, still_blocked)
1645 }
1646
1647 fn remove_record_with_attachments(
1651 &self,
1652 resolved: &str,
1653 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1654 ) -> Result<()> {
1655 self.sweep_removed_attachments();
1662 let attachments = self.attachments_dir(resolved);
1663 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1664 let moved = match std::fs::rename(&attachments, &aside) {
1665 Ok(()) => true,
1666 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1667 Err(e) => {
1668 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1669 }
1670 };
1671 let path = self.path_of(resolved);
1672 if let Err(e) = remove_record(&path) {
1673 if moved {
1674 let _ = std::fs::rename(&aside, &attachments);
1675 }
1676 return Err(e).with_context(|| format!("remove {}", path.display()));
1677 }
1678 if moved {
1679 if let Err(e) = std::fs::remove_dir_all(&aside) {
1680 tracing::warn!("leftover attachments {}: {e}", aside.display());
1681 }
1682 }
1683 let lock = self.lock_path(resolved);
1684 if let Err(e) = std::fs::remove_file(&lock) {
1685 if e.kind() != std::io::ErrorKind::NotFound {
1686 return Err(e).with_context(|| format!("remove {}", lock.display()));
1687 }
1688 }
1689 Ok(())
1690 }
1691
1692 fn sweep_removed_attachments(&self) {
1698 let Ok(entries) = std::fs::read_dir(&self.root) else {
1699 return;
1700 };
1701 for entry in entries.flatten() {
1702 let name = entry.file_name();
1703 let name = name.to_string_lossy();
1704 let Some(id) = name.strip_suffix(".attachments.removing") else {
1705 continue;
1706 };
1707 if !self.path_of(id).exists() {
1710 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1711 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1712 }
1713 }
1714 }
1715 self.sweep_tombstones();
1716 }
1717
1718 fn sweep_tombstones(&self) {
1727 let Ok(entries) = std::fs::read_dir(&self.root) else {
1728 return;
1729 };
1730 let tasks = self.list();
1731 for entry in entries.flatten() {
1732 let name = entry.file_name();
1733 let name = name.to_string_lossy();
1734 let Some(id) = name.strip_suffix(".removed") else {
1735 continue;
1736 };
1737 if self.path_of(id).exists()
1738 || tasks.iter().any(|t| t.blocked_by.iter().any(|b| b == id))
1739 {
1740 continue;
1741 }
1742 let old = entry
1743 .metadata()
1744 .and_then(|m| m.modified())
1745 .ok()
1746 .and_then(|m| m.elapsed().ok())
1747 .is_some_and(|age| age >= TOMBSTONE_GRACE);
1748 if old {
1749 let _ = std::fs::remove_file(entry.path());
1750 }
1751 }
1752 }
1753
1754 fn lock_path(&self, id: &str) -> PathBuf {
1757 self.root.join(format!("{id}.lock"))
1758 }
1759
1760 pub fn list(&self) -> Vec<Task> {
1773 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1774 .into_iter()
1775 .flatten()
1776 .flatten()
1777 .map(|e| e.path())
1778 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1779 .filter_map(|p| read_path(&p).ok())
1780 .collect();
1781 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1782 tasks
1783 }
1784
1785 pub fn superseded(&self) -> HashMap<String, String> {
1798 let mut by = HashMap::new();
1799 for task in self.list() {
1800 for earlier in &task.runs {
1801 if let Some(later) = task.successor_of(earlier) {
1802 by.insert(earlier.clone(), later.clone());
1803 }
1804 }
1805 }
1806 by
1807 }
1808
1809 pub fn superseded_by(&self, run: &str) -> Option<String> {
1817 for task in self.list() {
1818 if task.runs.iter().any(|r| r == run) {
1819 return task.successor_of(run).cloned();
1820 }
1821 }
1822 None
1823 }
1824
1825 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1836 for task in self.list() {
1837 if task.runs.iter().any(|r| r == run) {
1838 return task.runs.last().filter(|last| **last != run).cloned();
1839 }
1840 }
1841 None
1842 }
1843
1844 pub fn next_runnable(&self) -> Option<Task> {
1849 let mut runnable: Vec<Task> = self
1850 .list()
1851 .into_iter()
1852 .filter(|t| t.status.runnable())
1853 .collect();
1854 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1855 runnable.into_iter().next()
1856 }
1857
1858 pub fn claim(&self, id: &str) -> Result<Claim> {
1865 std::fs::create_dir_all(&self.root)
1866 .with_context(|| format!("create {}", self.root.display()))?;
1867 let path = self.lock_path(id);
1868 match std::fs::OpenOptions::new()
1869 .write(true)
1870 .create_new(true)
1871 .open(&path)
1872 {
1873 Ok(mut f) => {
1874 use std::io::Write as _;
1875 let _ = writeln!(f, "{}", std::process::id());
1877 Ok(Claim { path })
1878 }
1879 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1880 bail!("task {id} is already claimed ({} exists)", path.display())
1881 }
1882 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1883 }
1884 }
1885
1886 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1888 if self.path_of(prefix).is_file() {
1889 return Ok(prefix.to_owned());
1890 }
1891 let hits: Vec<String> = self
1892 .list()
1893 .into_iter()
1894 .map(|t| t.id)
1895 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1896 .collect();
1897 match hits.len() {
1898 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1899 0 => bail!("no task matches `{prefix}`"),
1900 _ => bail!(
1901 "`{prefix}` matches {} tasks: {}",
1902 hits.len(),
1903 hits.join(", ")
1904 ),
1905 }
1906 }
1907
1908 pub fn revision(&self) -> u64 {
1915 self.revision_excluding(&std::collections::BTreeSet::new())
1916 }
1917
1918 pub fn revision_excluding(&self, skip: &std::collections::BTreeSet<String>) -> u64 {
1921 use std::hash::{Hash as _, Hasher as _};
1922
1923 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
1924 .into_iter()
1925 .flatten()
1926 .flatten()
1927 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
1928 .filter(|e| {
1929 let path = e.path();
1930 !path
1931 .file_stem()
1932 .is_some_and(|stem| skip.contains(stem.to_string_lossy().as_ref()))
1933 })
1934 .filter_map(|e| {
1935 let name = e.file_name().to_string_lossy().into_owned();
1936 let mtime = e
1937 .metadata()
1938 .ok()?
1939 .modified()
1940 .ok()?
1941 .duration_since(std::time::UNIX_EPOCH)
1942 .ok()?
1943 .as_millis() as u64;
1944 Some((name, mtime))
1945 })
1946 .collect();
1947
1948 if entries.is_empty() {
1949 return 0;
1950 }
1951
1952 entries.sort_unstable();
1953 let mut hasher = std::hash::DefaultHasher::new();
1954 for (name, mtime) in &entries {
1955 name.hash(&mut hasher);
1956 mtime.hash(&mut hasher);
1957 }
1958 let h = hasher.finish();
1959 if h == 0 { 1 } else { h }
1960 }
1961}
1962
1963const TOMBSTONE_GRACE: std::time::Duration = std::time::Duration::from_secs(3600);
1966
1967#[derive(Debug, Clone)]
1969pub struct Removal {
1970 pub id: String,
1972 pub released: Vec<String>,
1976 pub still_blocked: Vec<String>,
1979}
1980
1981#[derive(Debug)]
1983pub struct Claim {
1984 path: PathBuf,
1985}
1986
1987impl Drop for Claim {
1988 fn drop(&mut self) {
1989 let _ = std::fs::remove_file(&self.path);
1990 }
1991}
1992
1993pub fn title_from(instruction: &str, max: usize) -> String {
1996 let Some(line) = first_line(instruction) else {
1997 return "(empty task)".to_owned();
1998 };
1999 if line.chars().count() <= max {
2000 return line.to_owned();
2001 }
2002 let head: String = line.chars().take(max.saturating_sub(1)).collect();
2003 format!("{head}…")
2004}
2005
2006pub(crate) fn first_line(instruction: &str) -> Option<&str> {
2012 let line = instruction
2013 .lines()
2014 .map(str::trim)
2015 .find(|l| !l.is_empty())?
2016 .trim_start_matches(['#', '-', '*', '>', ' '])
2017 .trim();
2018 (!line.is_empty()).then_some(line)
2019}
2020
2021fn read_path(path: &Path) -> Result<Task> {
2022 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2023 let task: Task =
2024 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
2025 if task.schema > SCHEMA {
2031 bail!(
2032 "task {} was written by a different magi (schema {}, this build \
2033 speaks {SCHEMA})",
2034 task.id,
2035 task.schema
2036 );
2037 }
2038 Ok(task)
2039}
2040
2041pub fn missing_blockers(
2055 queue: &Queue,
2056 questions: &Questions,
2057 blocked_by: &[String],
2058) -> Vec<String> {
2059 blocked_by
2060 .iter()
2061 .filter(|id| {
2062 !queue.path_of(id).is_file()
2063 && !questions.path_of(id).is_file()
2064 && !queue.was_deleted(id)
2065 })
2066 .cloned()
2067 .collect()
2068}
2069
2070pub fn deleted_blockers(queue: &Queue, blocked_by: &[String]) -> Vec<String> {
2075 blocked_by
2076 .iter()
2077 .filter(|id| queue.was_deleted(id))
2078 .cloned()
2079 .collect()
2080}
2081
2082pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
2094 missing_blocker_hold_reason_in(blocked_by, missing, "en")
2095}
2096
2097pub fn missing_blocker_hold_reason_in(
2099 blocked_by: &[String],
2100 missing: &[String],
2101 language: &str,
2102) -> String {
2103 if crate::lang::is_japanese(language) {
2104 format!(
2105 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
2106 blocked_by.join(", "),
2107 missing.join(", "),
2108 )
2109 } else {
2110 format!(
2111 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
2112 blocked_by.join(", "),
2113 missing.join(", "),
2114 )
2115 }
2116}
2117
2118pub fn short(id: &str) -> &str {
2120 id.split('-').next_back().unwrap_or(id)
2121}
2122
2123fn new_id() -> String {
2124 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2125 let seed = crate::rng::entropy();
2126 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2127}
2128
2129#[cfg(test)]
2130mod tests {
2131 #[test]
2132 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
2133 let mut t = task("relanded");
2134 t.attempts = 1;
2135 t.status = TaskStatus::Running;
2136 t.already_landed("already in main as 0e368de");
2137 assert_eq!(t.status, TaskStatus::Done);
2138 assert_eq!(t.attempts, 0);
2139 assert!(t.hold_reason.is_none());
2140 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
2141 }
2142
2143 #[test]
2144 fn missing_blocker_reason_follows_the_language() {
2145 let b = vec!["a".to_owned()];
2146 let en = missing_blocker_hold_reason_in(&b, &b, "en");
2147 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
2148 assert!(en.starts_with("blocked on a"));
2149 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
2150 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
2151 }
2152
2153 use super::*;
2154
2155 #[test]
2156 fn triage_applied_survives_release_and_old_records_read_as_empty() {
2157 let mut t = Task::new(
2158 "t".to_owned(),
2159 "i".to_owned(),
2160 PathBuf::from("r"),
2161 Source::Human,
2162 );
2163 t.mark_triage_applied("q1");
2164 t.mark_triage_applied("q1");
2165 t.hold_machine(Some("x".to_owned()));
2166 t.release();
2167 assert_eq!(t.triage_applied, ["q1"]);
2168 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
2169
2170 let mut v = serde_json::to_value(&t).unwrap();
2171 v.as_object_mut().unwrap().remove("triage_applied");
2172 let old: Task = serde_json::from_value(v).unwrap();
2173 assert!(old.triage_applied.is_empty());
2174 }
2175
2176 #[test]
2177 fn task_counts_of_empty_is_all_zero() {
2178 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
2179 }
2180
2181 #[test]
2182 fn task_counts_of_tallies_every_status() {
2183 let mut queued = Task::new(
2184 "q".to_owned(),
2185 "i".to_owned(),
2186 PathBuf::from("."),
2187 Source::Human,
2188 );
2189 queued.status = TaskStatus::Queued;
2190 let mut running = queued.clone();
2191 running.status = TaskStatus::Running;
2192 let mut done = queued.clone();
2193 done.status = TaskStatus::Done;
2194 let mut failed = queued.clone();
2195 failed.status = TaskStatus::Failed;
2196 let mut held = queued.clone();
2197 held.status = TaskStatus::Held;
2198 let mut blocked = queued.clone();
2199 blocked.status = TaskStatus::Blocked;
2200
2201 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
2202 assert_eq!(
2203 counts,
2204 TaskCounts {
2205 queued: 1,
2206 running: 1,
2207 done: 2,
2208 failed: 1,
2209 held: 1,
2210 blocked: 1,
2211 }
2212 );
2213 }
2214
2215 fn queue() -> (tempfile::TempDir, Queue) {
2218 let dir = tempfile::tempdir().unwrap();
2219 let q = Queue::at(dir.path().join("queue"));
2220 (dir, q)
2221 }
2222
2223 #[test]
2224 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
2225 let dir = tempfile::tempdir().unwrap();
2226 let q = Queue::at(dir.path().join("queue"));
2227 let mut t = task("held");
2228 q.put(&mut t).unwrap();
2229 assert_eq!(
2230 crate::notices::Notices::at(dir.path().join("notifications"))
2231 .list()
2232 .len(),
2233 0
2234 );
2235 t.hold_machine(Some("out of attempts".to_owned()));
2236 q.put(&mut t).unwrap();
2237 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
2238 assert_eq!(listed.len(), 1);
2239 assert!(listed[0].message.contains("out of attempts"));
2240 }
2241
2242 fn task(title: &str) -> Task {
2243 Task::new(
2244 title.to_owned(),
2245 format!("do {title}"),
2246 PathBuf::from("."),
2247 Source::Human,
2248 )
2249 }
2250
2251 #[test]
2252 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
2253 let mut t = task("retried");
2254 assert!(
2255 t.earlier_attempts().is_empty(),
2256 "a first attempt takes nothing over"
2257 );
2258 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2259 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
2260 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
2262 assert_eq!(t.successor_of("bbbb"), None);
2263 assert_eq!(t.successor_of("zzzz"), None);
2264 }
2265
2266 #[test]
2267 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
2268 let (_dir, q) = queue();
2269 let mut t = task("retried");
2270 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2271 q.put(&mut t).unwrap();
2272
2273 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
2274 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
2275 assert_eq!(
2276 q.superseded_by("cccc"),
2277 None,
2278 "the latest attempt replaces nothing"
2279 );
2280 assert_eq!(
2281 q.superseded_by("never-heard-of-it"),
2282 None,
2283 "a run belonging to no task on this queue is not superseded"
2284 );
2285
2286 let mut by = HashMap::new();
2287 by.insert("aaaa".to_owned(), "bbbb".to_owned());
2288 by.insert("bbbb".to_owned(), "cccc".to_owned());
2289 assert_eq!(
2290 q.superseded(),
2291 by,
2292 "the whole-map and single-run forms must agree"
2293 );
2294 }
2295
2296 #[test]
2297 fn superseded_attempts_is_empty_until_the_task_is_done() {
2298 let mut t = task("retried");
2299 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2300 t.status = TaskStatus::Failed;
2301 assert_eq!(
2302 t.superseded_attempts(true),
2303 &[] as &[String],
2304 "a task still retrying has no attempt yet that a later one made moot"
2305 );
2306
2307 t.status = TaskStatus::Running;
2308 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2309 }
2310
2311 #[test]
2312 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
2313 let mut t = task("retried");
2314 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2315 t.status = TaskStatus::Done;
2316 assert_eq!(
2317 t.superseded_attempts(true),
2318 &["aaaa".to_owned(), "bbbb".to_owned()],
2319 "cccc is the attempt whose success made the task done, and stays out"
2320 );
2321 }
2322
2323 #[test]
2324 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
2325 let mut t = task("first try landed");
2326 t.runs = vec!["aaaa".to_owned()];
2327 t.status = TaskStatus::Done;
2328 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2329 }
2330
2331 #[test]
2332 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
2333 let mut t = task("closed by hand after a manual merge");
2340 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2341 t.status = TaskStatus::Done;
2342 assert_eq!(
2343 t.superseded_attempts(false),
2344 &[] as &[String],
2345 "nothing here is provably why the task is done, so nothing is superseded"
2346 );
2347 }
2348
2349 #[test]
2350 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2351 let (_dir, q) = queue();
2352 let mut t = task("retried twice");
2353 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2354 q.put(&mut t).unwrap();
2355
2356 assert_eq!(
2357 q.latest_attempt("aaaa"),
2358 Some("cccc".to_owned()),
2359 "an old attempt points straight at the chain's current head, not the \
2360 next attempt in the middle of it"
2361 );
2362 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2363 assert_eq!(
2364 q.latest_attempt("cccc"),
2365 None,
2366 "the latest attempt is not superseded by anything"
2367 );
2368 assert_eq!(
2369 q.latest_attempt("never-heard-of-it"),
2370 None,
2371 "a run belonging to no task on this queue is not superseded"
2372 );
2373 }
2374
2375 #[test]
2376 fn a_markdown_heading_is_the_title_not_decoration() {
2377 assert_eq!(
2382 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2383 "Rework the config loader"
2384 );
2385 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2386 assert_eq!(title_from("> quoted task", 40), "quoted task");
2387 assert_eq!(title_from(" \n\n", 40), "(empty task)");
2389 assert_eq!(title_from("###\n", 40), "(empty task)");
2390 }
2391
2392 #[test]
2393 fn a_long_title_is_elided_by_characters_not_bytes() {
2394 let long = "課題".repeat(30);
2396 let title = title_from(&long, 10);
2397 assert_eq!(title.chars().count(), 10);
2398 assert!(title.ends_with('…'));
2399 }
2400
2401 #[test]
2402 fn priority_wins_and_ties_break_oldest_first() {
2403 let (_dir, q) = queue();
2404 let mut a = task("first");
2405 let mut b = task("second");
2406 let mut c = task("urgent");
2407 a.id = "20260101-000001-aaaa".to_owned();
2409 b.id = "20260101-000002-bbbb".to_owned();
2410 c.id = "20260101-000003-cccc".to_owned();
2411 c.priority = 5;
2412 for t in [&mut a, &mut b, &mut c] {
2413 q.put(t).unwrap();
2414 }
2415
2416 assert_eq!(q.next_runnable().unwrap().id, c.id);
2418 c.hold_machine(None);
2419 q.put(&mut c).unwrap();
2420 assert_eq!(q.next_runnable().unwrap().id, a.id);
2422 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2423 }
2424
2425 #[test]
2426 fn a_blocked_task_never_starves_another_runnable_one() {
2427 let (_dir, q) = queue();
2428 let mut blocked = task("blocked");
2429 blocked.block(vec!["something".to_owned()], None);
2430 q.put(&mut blocked).unwrap();
2431
2432 let mut runnable = task("free to go");
2433 q.put(&mut runnable).unwrap();
2434
2435 let next = q.next_runnable().expect("a runnable task is still offered");
2436 assert_eq!(next.id, runnable.id);
2437 }
2438
2439 #[test]
2440 fn a_held_task_is_never_offered_to_the_loop() {
2441 let (_dir, q) = queue();
2442 let mut t = task("held");
2443 q.put(&mut t).unwrap();
2444 assert!(q.next_runnable().is_some());
2445
2446 t.hold_machine(None);
2447 q.put(&mut t).unwrap();
2448 assert!(
2449 q.next_runnable().is_none(),
2450 "a held task must wait for a human"
2451 );
2452
2453 t.status = TaskStatus::Failed;
2455 q.put(&mut t).unwrap();
2456 assert!(q.next_runnable().is_some());
2457 }
2458
2459 #[test]
2460 fn attempts_are_capped_and_then_the_task_is_held() {
2461 let mut t = task("doomed");
2462
2463 t.start("run-1".to_owned());
2464 t.fail("gate red", 2);
2465 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2466
2467 t.start("run-2".to_owned());
2468 t.fail("gate red", 2);
2469 assert_eq!(
2470 t.status,
2471 TaskStatus::Held,
2472 "out of attempts: stop spending money on it"
2473 );
2474 assert_eq!(t.runs, ["run-1", "run-2"]);
2475 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2476 assert_eq!(
2477 t.hold_reason.as_deref(),
2478 Some("gate red"),
2479 "the hold must say why, not leave hold_reason null next to a \
2480 populated last_error"
2481 );
2482 }
2483
2484 #[test]
2485 fn handing_off_a_task_records_a_hold_reason_too() {
2486 let mut t = task("left a pull request");
2487 t.start("run-1".to_owned());
2488 t.handed_off("run ended with a pull request open [run run-1]");
2489 assert_eq!(t.status, TaskStatus::Held);
2490 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2491 assert_eq!(
2492 t.hold_reason.as_deref(),
2493 Some("run ended with a pull request open [run run-1]")
2494 );
2495 assert_eq!(t.hold_reason, t.last_error);
2496 }
2497
2498 #[test]
2499 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2500 let mut t = task("stalled by quota");
2501
2502 t.start("run-1".to_owned());
2503 assert_eq!(t.attempts, 1);
2504 t.stall("judge-1, judge-2 out of quota");
2505 assert_eq!(
2506 t.attempts, 0,
2507 "a closed quota window must not spend the task's retry budget"
2508 );
2509 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2510 assert_eq!(
2511 t.last_error.as_deref(),
2512 Some("judge-1, judge-2 out of quota")
2513 );
2514
2515 for _ in 0..20 {
2518 t.start("run-n".to_owned());
2519 t.stall("still out of quota");
2520 }
2521 t.start("run-real".to_owned());
2522 t.fail("gate red", 2);
2523 assert_eq!(
2524 t.status,
2525 TaskStatus::Failed,
2526 "the first attempt that was really judged is attempt one"
2527 );
2528 }
2529
2530 #[test]
2531 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2532 let mut t = task("retry me");
2533 t.start("run-1".to_owned());
2534 t.fail("gate red", 1);
2535 assert_eq!(t.status, TaskStatus::Held);
2536
2537 t.release();
2538 assert_eq!(t.status, TaskStatus::Queued);
2539 assert_eq!(t.attempts, 0);
2542 assert!(t.last_error.is_none());
2543 assert_eq!(
2544 t.runs.len(),
2545 1,
2546 "history is kept: attempts reset, evidence does not"
2547 );
2548 }
2549
2550 #[test]
2551 fn a_hold_reason_survives_and_a_release_clears_it() {
2552 let mut t = task("waiting on something else");
2553 t.hold_manual(Some(
2554 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2555 ));
2556 assert_eq!(t.status, TaskStatus::Held);
2557 assert_eq!(
2558 t.hold_reason.as_deref(),
2559 Some("waiting for 20260101-000000-aaaa to land first")
2560 );
2561
2562 t.hold_manual(None);
2564 assert_eq!(
2565 t.hold_reason.as_deref(),
2566 Some("waiting for 20260101-000000-aaaa to land first"),
2567 "a bare re-hold keeps whatever a human already wrote down"
2568 );
2569
2570 let mut plain = task("no reason given");
2572 plain.hold_manual(None);
2573 assert_eq!(plain.status, TaskStatus::Held);
2574 assert!(plain.hold_reason.is_none());
2575
2576 t.release();
2577 assert_eq!(t.status, TaskStatus::Queued);
2578 assert!(
2579 t.hold_reason.is_none(),
2580 "a stale reason must not greet the next person who holds this task"
2581 );
2582 }
2583
2584 #[test]
2585 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2586 let mut t = task("landed by hand while held");
2591 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2592 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2593
2594 t.succeed();
2595 assert_eq!(t.status, TaskStatus::Done);
2596 assert!(
2597 t.hold_reason.is_none(),
2598 "a done task cannot still be waiting on something"
2599 );
2600 }
2601
2602 #[test]
2603 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2604 let mut held = task("held straight out of blocked");
2611 held.block(
2612 vec!["20260101-000000-dead".to_owned()],
2613 Some("waiting on the migration script".to_owned()),
2614 );
2615 assert_eq!(held.status, TaskStatus::Blocked);
2616
2617 held.hold_manual(None);
2618 assert_eq!(held.status, TaskStatus::Held);
2619 assert!(
2620 held.blocked_by.is_empty(),
2621 "hold overrides the wait, same as release"
2622 );
2623 assert!(held.block_reason.is_none());
2624
2625 let mut done = task("closed straight out of blocked");
2626 done.block(
2627 vec!["20260101-000000-dead".to_owned()],
2628 Some("waiting on the migration script".to_owned()),
2629 );
2630 done.succeed();
2631 assert_eq!(done.status, TaskStatus::Done);
2632 assert!(
2633 done.blocked_by.is_empty(),
2634 "a done task cannot still be waiting on a dependency"
2635 );
2636 assert!(done.block_reason.is_none());
2637 }
2638
2639 #[test]
2640 fn a_blocked_task_is_never_offered_to_the_loop() {
2641 let mut t = task("blocked");
2642 assert!(t.status.runnable());
2643 t.block(
2644 vec!["dep-id".to_owned()],
2645 Some("waits on dep-id".to_owned()),
2646 );
2647 assert_eq!(t.status, TaskStatus::Blocked);
2648 assert!(!t.status.runnable());
2649 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2650 }
2651
2652 #[test]
2653 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2654 let mut t = task("blocked on two");
2655 t.block(
2656 vec!["a".to_owned(), "b".to_owned()],
2657 Some("waits on a and b".to_owned()),
2658 );
2659
2660 t.unblock("a");
2661 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2662 assert_eq!(t.blocked_by, ["b"]);
2663
2664 t.unblock("b");
2665 assert_eq!(t.status, TaskStatus::Queued);
2666 assert!(t.blocked_by.is_empty());
2667 assert!(t.block_reason.is_none());
2668 }
2669
2670 #[test]
2671 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2672 let mut t = task("never blocked");
2673 t.unblock("whatever");
2674 assert_eq!(t.status, TaskStatus::Queued);
2675 }
2676
2677 #[test]
2678 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2679 let mut t = task("held, then asked about");
2684 t.hold_machine(Some("out of attempts".to_owned()));
2685 assert_eq!(t.status, TaskStatus::Held);
2686
2687 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2688 assert_eq!(t.status, TaskStatus::Blocked);
2689
2690 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2691 t.unblock("q1");
2692 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2693 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2694 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2695 assert!(t.blocked_from.is_none(), "consumed once restored");
2696 }
2697
2698 #[test]
2699 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2700 let mut t = task("manually held, then asked about");
2701 t.hold_manual(Some("waiting on a dependency".to_owned()));
2702
2703 t.block(vec!["q1".to_owned()], None);
2704 t.unblock("q1");
2705
2706 assert_eq!(t.status, TaskStatus::Held);
2707 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2708 }
2709
2710 #[test]
2711 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2712 let mut t = task("held, blocked twice");
2716 t.hold_machine(None);
2717 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2718 t.block(
2719 vec!["q1".to_owned(), "q2".to_owned()],
2720 Some("second".to_owned()),
2721 );
2722
2723 t.unblock("q1");
2724 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2725 t.unblock("q2");
2726 assert_eq!(t.status, TaskStatus::Held);
2727 }
2728
2729 #[test]
2730 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2731 let mut t = task("blocked mid-run");
2735 t.start("run-1".to_owned());
2736 assert_eq!(t.status, TaskStatus::Running);
2737
2738 t.block(vec!["q1".to_owned()], None);
2739 t.unblock("q1");
2740 assert_eq!(t.status, TaskStatus::Queued);
2741 }
2742
2743 #[test]
2744 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2745 let mut t = task("legacy record, held before it was blocked");
2751 t.hold_source = Some(HoldSource::Machine);
2752 t.hold_reason = Some("legacy hold reason".to_owned());
2753 t.status = TaskStatus::Blocked;
2754 t.blocked_by = vec!["q1".to_owned()];
2755 t.blocked_from = None;
2756
2757 t.unblock("q1");
2758 assert_eq!(t.status, TaskStatus::Held);
2759 }
2760
2761 #[test]
2762 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2763 let mut t = task("legacy record, ordinary dependency block");
2764 t.status = TaskStatus::Blocked;
2765 t.blocked_by = vec!["dep".to_owned()];
2766 t.blocked_from = None;
2767
2768 t.unblock("dep");
2769 assert_eq!(t.status, TaskStatus::Queued);
2770 }
2771
2772 #[test]
2773 fn answering_a_question_is_recorded_and_survives_a_release() {
2774 let mut t = task("asked something");
2775 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2776 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2777 t.unblock("q1");
2778 assert_eq!(t.status, TaskStatus::Queued);
2779 assert_eq!(t.answers.len(), 1);
2780 assert_eq!(t.answers[0].answer, "SQLite");
2781
2782 t.release();
2786 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2787 }
2788
2789 #[test]
2790 fn a_refused_handover_keeps_the_review_branch_across_release() {
2791 let mut t = task("refused takeover");
2792 t.start("run-1".to_owned());
2793 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
2794 assert_eq!(t.status, TaskStatus::Held);
2795 assert_eq!(t.attempts, 1);
2796 t.release();
2797 assert_eq!(t.status, TaskStatus::Queued);
2798 assert_eq!(t.attempts, 0);
2799 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2800
2801 t.hold_for_handover(None, "again".to_owned());
2803 t.requeue();
2804 assert!(t.review_branch.is_none());
2805
2806 let mut m = task("manual");
2808 m.review_branch = Some("magi/x/A".to_owned());
2809 m.hold_manual(None);
2810 m.release();
2811 assert!(m.review_branch.is_none());
2812 }
2813
2814 #[test]
2815 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2816 let mut t = task("blocked run with a surviving branch");
2817 t.start("run-1".to_owned());
2818 t.fail("blocked with major findings", 5);
2819 assert_eq!(t.status, TaskStatus::Failed);
2820
2821 t.request_review("magi/eba2/A".to_owned());
2822 assert_eq!(t.status, TaskStatus::Queued);
2823 assert_eq!(t.attempts, 0);
2824 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2825
2826 t.release();
2828 assert!(t.review_branch.is_none());
2829 }
2830
2831 #[test]
2832 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2833 let mut t = task("retry");
2834 t.start("run-1".to_owned());
2835 t.requeue();
2836 assert!(t.fresh_start);
2837
2838 t.release();
2839 assert!(!t.fresh_start);
2840 }
2841
2842 #[test]
2843 fn priority_can_be_changed_while_queued_but_not_while_running() {
2844 let mut t = task("reprioritise me");
2845 t.set_priority(5).unwrap();
2846 assert_eq!(t.priority, 5);
2847
2848 t.start("run-1".to_owned());
2849 let err = t.set_priority(9).unwrap_err().to_string();
2850 assert!(err.contains("running"), "{err}");
2851 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2852 }
2853
2854 #[test]
2855 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2856 let mut t = task("interrupt me");
2857 assert!(!t.interrupt, "off unless asked, same as any other task");
2858
2859 t.set_interrupt(true).unwrap();
2860 assert!(t.interrupt);
2861
2862 t.start("run-1".to_owned());
2863 assert!(
2864 !t.interrupt,
2865 "the mark is one-shot: dispatching the task fulfils it, \
2866 whatever the run that follows ends up doing"
2867 );
2868 let err = t.set_interrupt(true).unwrap_err().to_string();
2869 assert!(err.contains("running"), "{err}");
2870 t.set_interrupt(false).unwrap();
2873 assert!(!t.interrupt);
2874 }
2875
2876 #[test]
2880 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2881 let mut t = task("interrupt me");
2882 t.set_interrupt(true).unwrap();
2883 t.start("run-1".to_owned());
2884 t.fail("mock failure", 5);
2885 assert_eq!(t.status, TaskStatus::Failed);
2886 assert!(
2887 !t.interrupt,
2888 "one attempt already spent the mark; a retry is an ordinary \
2889 requeue, not a fresh interrupt request"
2890 );
2891 }
2892
2893 #[test]
2894 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2895 let (_dir, q) = queue();
2896 let mut a = task("first filed");
2897 let mut b = task("second filed");
2898 a.id = "20260101-000001-aaaa".to_owned();
2899 b.id = "20260101-000002-bbbb".to_owned();
2900 q.put(&mut a).unwrap();
2901 q.put(&mut b).unwrap();
2902
2903 assert_eq!(
2904 q.next_runnable().unwrap().id,
2905 a.id,
2906 "with equal priority the older task goes first, so a burst of \
2907 new work cannot starve it"
2908 );
2909 assert_eq!(
2910 q.list()[0].id,
2911 b.id,
2912 "but the list an operator reads is newest first, the same as \
2913 before priority existed - a's turn to run does not make it the \
2914 newest task"
2915 );
2916
2917 let mut a = q.get(&a.id).unwrap();
2918 a.set_priority(10).unwrap();
2919 q.put(&mut a).unwrap();
2920
2921 assert_eq!(
2922 q.next_runnable().unwrap().id,
2923 a.id,
2924 "a raised priority must be reflected the moment it is saved"
2925 );
2926 assert_eq!(
2930 q.list()[0].id,
2931 a.id,
2932 "the raised task must sort first in the list an operator reads, \
2933 not only in next_runnable's own ordering"
2934 );
2935 }
2936
2937 #[test]
2938 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2939 let mut t = Task::new(
2940 "old title".to_owned(),
2941 "old instruction".to_owned(),
2942 PathBuf::from("/repo"),
2943 Source::Agent {
2944 run: "20260101-000000-beef".to_owned(),
2945 node: "implement".to_owned(),
2946 },
2947 );
2948 let id = t.id.clone();
2949 let created_at = t.created_at;
2950 t.runs.push("20260101-000000-beef".to_owned());
2951
2952 t.edit("new title".to_owned(), "new instruction".to_owned())
2953 .unwrap();
2954
2955 assert_eq!(t.title, "new title");
2956 assert_eq!(t.instruction, "new instruction");
2957 assert_eq!(t.id, id, "editing must not mint a new id");
2958 assert_eq!(t.created_at, created_at);
2959 assert_eq!(
2960 t.source,
2961 Source::Agent {
2962 run: "20260101-000000-beef".to_owned(),
2963 node: "implement".to_owned(),
2964 },
2965 "editing must not turn agent attribution into human"
2966 );
2967 assert_eq!(t.runs, ["20260101-000000-beef"]);
2968 }
2969
2970 #[test]
2971 fn editing_is_refused_once_a_task_is_running_or_finished() {
2972 let mut running = task("in flight");
2973 running.start("run-1".to_owned());
2974 let err = running
2975 .edit("x".to_owned(), "y".to_owned())
2976 .unwrap_err()
2977 .to_string();
2978 assert!(err.contains("running"), "{err}");
2979
2980 let mut done = task("finished");
2981 done.succeed();
2982 let err = done
2983 .edit("x".to_owned(), "y".to_owned())
2984 .unwrap_err()
2985 .to_string();
2986 assert!(err.contains("done"), "{err}");
2987
2988 let mut queued = task("waiting");
2990 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2991 let mut held = task("parked");
2992 held.hold_machine(None);
2993 held.edit("x".to_owned(), "y".to_owned()).unwrap();
2994 }
2995
2996 #[test]
2997 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2998 let (_dir, q) = queue();
2999 let path = q.path_of("20260101-000000-aaaa");
3000 std::fs::create_dir_all(q.root()).unwrap();
3001 std::fs::write(
3002 &path,
3003 serde_json::json!({
3004 "schema": SCHEMA,
3005 "id": "20260101-000000-aaaa",
3006 "title": "from before hold reasons existed",
3007 "instruction": "from before hold reasons existed",
3008 "repo": ".",
3009 "source": { "kind": "human" },
3010 "status": "held",
3011 "created_at": Timestamp::now().to_string(),
3012 "updated_at": Timestamp::now().to_string(),
3013 })
3014 .to_string(),
3015 )
3016 .unwrap();
3017
3018 let task = q.get("20260101-000000-aaaa").expect("must still read");
3019 assert!(task.hold_reason.is_none());
3020 assert!(task.operator_held());
3021 }
3022
3023 #[test]
3024 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
3025 let (_dir, q) = queue();
3026 let path = q.path_of("20260101-000000-bbbb");
3027 std::fs::create_dir_all(q.root()).unwrap();
3028 std::fs::write(
3029 &path,
3030 serde_json::json!({
3031 "schema": 2,
3032 "id": "20260101-000000-bbbb",
3033 "title": "old manual recovery",
3034 "instruction": "old manual recovery",
3035 "repo": ".",
3036 "source": { "kind": "human" },
3037 "status": "held",
3038 "hold_reason": "active manual recovery run20260912-224242-daf5",
3039 "created_at": Timestamp::now().to_string(),
3040 "updated_at": Timestamp::now().to_string(),
3041 })
3042 .to_string(),
3043 )
3044 .unwrap();
3045
3046 let task = q.get("20260101-000000-bbbb").expect("must still read");
3047 assert_eq!(task.hold_source, None);
3048 assert!(task.operator_held());
3049 }
3050
3051 #[test]
3052 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
3053 let (_dir, q) = queue();
3054 let path = q.path_of("20260101-000000-aaaa");
3055 std::fs::create_dir_all(q.root()).unwrap();
3056 std::fs::write(
3057 &path,
3058 serde_json::json!({
3059 "schema": SCHEMA,
3060 "id": "20260101-000000-aaaa",
3061 "title": "from before diagnostics existed",
3062 "instruction": "from before diagnostics existed",
3063 "repo": ".",
3064 "source": { "kind": "human" },
3065 "status": "held",
3066 "created_at": Timestamp::now().to_string(),
3067 "updated_at": Timestamp::now().to_string(),
3068 })
3069 .to_string(),
3070 )
3071 .unwrap();
3072
3073 let task = q.get("20260101-000000-aaaa").expect("must still read");
3074 assert!(task.diagnostic.is_none());
3075 }
3076
3077 #[test]
3078 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
3079 let (_dir, q) = queue();
3083 let path = q.path_of("20260101-000000-aaaa");
3084 std::fs::create_dir_all(q.root()).unwrap();
3085 std::fs::write(
3086 &path,
3087 serde_json::json!({
3088 "schema": 1,
3089 "id": "20260101-000000-aaaa",
3090 "title": "from before blocking existed",
3091 "instruction": "from before blocking existed",
3092 "repo": ".",
3093 "source": { "kind": "human" },
3094 "status": "queued",
3095 "created_at": Timestamp::now().to_string(),
3096 "updated_at": Timestamp::now().to_string(),
3097 })
3098 .to_string(),
3099 )
3100 .unwrap();
3101
3102 let task = q.get("20260101-000000-aaaa").expect("must still read");
3103 assert!(task.blocked_by.is_empty());
3104 assert!(task.block_reason.is_none());
3105 assert!(task.answers.is_empty());
3106 assert!(task.review_branch.is_none());
3107 }
3108
3109 #[test]
3110 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
3111 let mut held = task("diagnosed");
3116 held.start("run-1".to_owned());
3117 held.fail("gate red", 1);
3118 held.diagnostic = Some("cargo test failed: ...".to_owned());
3119 assert_eq!(held.status, TaskStatus::Held);
3120
3121 held.release();
3122 assert!(held.diagnostic.is_none());
3123
3124 held.diagnostic = Some("cargo test failed: ...".to_owned());
3125 held.succeed();
3126 assert!(held.diagnostic.is_none());
3127 }
3128
3129 #[test]
3130 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
3131 let mut t = task("retried");
3132 t.start("run-1".to_owned());
3133 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
3134 t.fail("unrelated config error", 5);
3135 assert_eq!(t.status, TaskStatus::Failed);
3136 assert!(
3137 t.diagnostic.is_none(),
3138 "fail() must not let an old diagnostic outlive the run that produced it"
3139 );
3140 }
3141
3142 #[test]
3143 fn a_claim_is_exclusive_and_releases_on_drop() {
3144 let (_dir, q) = queue();
3145 let mut t = task("contended");
3146 q.put(&mut t).unwrap();
3147
3148 let held = q.claim(&t.id).unwrap();
3149 assert!(
3150 q.claim(&t.id).is_err(),
3151 "two daemons must not drive one task into two runs"
3152 );
3153 drop(held);
3154 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
3155 }
3156
3157 #[test]
3158 fn a_round_trip_survives_disk() {
3159 let (_dir, q) = queue();
3160 let mut t = Task::new(
3161 "titled".to_owned(),
3162 "body".to_owned(),
3163 PathBuf::from("/repo"),
3164 Source::Agent {
3165 run: "20260101-000000-beef".to_owned(),
3166 node: "implement".to_owned(),
3167 },
3168 );
3169 t.priority = 3;
3170 q.put(&mut t).unwrap();
3171
3172 let back = q.get(&t.id).unwrap();
3173 assert_eq!(back.id, t.id);
3174 assert_eq!(back.priority, 3);
3175 assert_eq!(back.source.label(), "implement@beef");
3176 assert_eq!(q.get(t.short()).unwrap().id, t.id);
3178 }
3179
3180 #[test]
3181 fn an_unreadable_task_does_not_take_the_queue_down() {
3182 let (_dir, q) = queue();
3183 let mut t = task("fine");
3184 q.put(&mut t).unwrap();
3185 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
3186
3187 let listed = q.list();
3188 assert_eq!(listed.len(), 1, "the readable task still lists");
3189 assert_eq!(listed[0].id, t.id);
3190 }
3191
3192 #[test]
3193 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
3194 let (_dir, q) = queue();
3195 let path = q.path_of("20260101-000000-aaaa");
3196 std::fs::create_dir_all(q.root()).unwrap();
3197 std::fs::write(
3198 &path,
3199 serde_json::json!({
3200 "schema": SCHEMA,
3201 "id": "20260101-000000-aaaa",
3202 "title": "from before solo existed",
3203 "instruction": "from before solo existed",
3204 "repo": ".",
3205 "source": { "kind": "human" },
3206 "status": "queued",
3207 "created_at": Timestamp::now().to_string(),
3208 "updated_at": Timestamp::now().to_string(),
3209 })
3210 .to_string(),
3211 )
3212 .unwrap();
3213
3214 let task = q.get("20260101-000000-aaaa").expect("must still read");
3215 assert!(!task.solo, "a queue file with no `solo` field means false");
3216 }
3217
3218 #[test]
3219 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
3220 let (_dir, q) = queue();
3221 let path = q.path_of("20260101-000000-bbbb");
3222 std::fs::create_dir_all(q.root()).unwrap();
3223 std::fs::write(
3224 &path,
3225 serde_json::json!({
3226 "schema": SCHEMA,
3227 "id": "20260101-000000-bbbb",
3228 "title": "from before urgent existed",
3229 "instruction": "from before urgent existed",
3230 "repo": ".",
3231 "source": { "kind": "human" },
3232 "status": "queued",
3233 "created_at": Timestamp::now().to_string(),
3234 "updated_at": Timestamp::now().to_string(),
3235 })
3236 .to_string(),
3237 )
3238 .unwrap();
3239
3240 let task = q.get("20260101-000000-bbbb").expect("must still read");
3241 assert!(
3242 !task.urgent,
3243 "a queue file with no `urgent` field means false, same as `solo`"
3244 );
3245 }
3246
3247 #[test]
3248 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
3249 let (_dir, q) = queue();
3250 let mut t = task("from the future");
3251 q.put(&mut t).unwrap();
3252 let path = q.path_of(&t.id);
3253 let body = std::fs::read_to_string(&path)
3254 .unwrap()
3255 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
3256 std::fs::write(&path, body).unwrap();
3257
3258 let err = q.get(&t.id).unwrap_err().to_string();
3259 assert!(err.contains("schema 99"), "{err}");
3260 }
3261
3262 #[test]
3263 fn revision_moves_when_the_queue_changes() {
3264 let (_dir, q) = queue();
3265 assert_eq!(q.revision(), 0, "an empty queue has no revision");
3266 let mut t = task("first");
3267 q.put(&mut t).unwrap();
3268 assert!(q.revision() > 0, "a written task moves the revision");
3269 }
3270
3271 #[test]
3272 fn revision_moves_when_deleting_an_older_task() {
3273 let (dir, q) = queue();
3274 let questions = Questions::at(dir.path().join("questions"));
3275 let mut t1 = task("older");
3276 q.put(&mut t1).unwrap();
3277 std::thread::sleep(std::time::Duration::from_millis(10));
3279 let mut t2 = task("newer");
3280 q.put(&mut t2).unwrap();
3281
3282 let rev_before = q.revision();
3283 q.remove(&t1.id, false, &questions).unwrap();
3284 let rev_after = q.revision();
3285
3286 assert_ne!(
3287 rev_before, rev_after,
3288 "deleting an older task must change the revision so other clients see the deletion"
3289 );
3290 }
3291
3292 #[test]
3293 fn removing_a_task_takes_it_out_of_the_listing() {
3294 let (dir, q) = queue();
3295 let questions = Questions::at(dir.path().join("questions"));
3296 let mut t = task("delete me");
3297 q.put(&mut t).unwrap();
3298 let removed = q.remove(t.short(), false, &questions).unwrap();
3299 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
3300 assert!(removed.released.is_empty() && removed.still_blocked.is_empty());
3301 assert!(q.list().is_empty());
3302 assert!(
3303 q.remove(&t.id, false, &questions).is_err(),
3304 "removing twice is an error"
3305 );
3306 }
3307
3308 #[test]
3309 fn removing_a_task_takes_its_stale_lock_with_it() {
3310 let (dir, q) = queue();
3311 let questions = Questions::at(dir.path().join("questions"));
3312 let mut t = task("interrupted");
3313 q.put(&mut t).unwrap();
3314
3315 let claim = q.claim(&t.id).unwrap();
3318 std::mem::forget(claim);
3319 assert!(
3320 q.claim(&t.id).is_err(),
3321 "the orphaned lock is what makes the task look claimed"
3322 );
3323
3324 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
3326 assert!(err.contains("live daemon"), "{err}");
3327 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
3328
3329 q.remove(&t.id, false, &questions).unwrap();
3331 assert!(q.list().is_empty());
3332 let mut again = task("interrupted");
3333 again.id = t.id.clone();
3334 q.put(&mut again).unwrap();
3335 assert!(
3336 q.claim(&t.id).is_ok(),
3337 "a task that comes back must be claimable, which a left-behind lock would prevent"
3338 );
3339 }
3340
3341 fn notices_of(dir: &Path) -> Vec<crate::notices::Notice> {
3342 crate::notices::Notices::at(dir.join("notifications")).list()
3343 }
3344
3345 #[test]
3346 fn removing_a_sole_dependency_releases_the_dependent_without_a_hold() {
3347 let (dir, q) = queue();
3348 let questions = Questions::at(dir.path().join("questions"));
3349 let mut dep = task("dependency");
3350 q.put(&mut dep).unwrap();
3351 let mut blocked = task("waiting");
3352 blocked.block(vec![dep.id.clone()], Some("waits".to_owned()));
3353 q.put(&mut blocked).unwrap();
3354
3355 let removed = q.remove(&dep.id, false, &questions).unwrap();
3356 assert_eq!(removed.released, [blocked.id.clone()]);
3357 assert!(removed.still_blocked.is_empty());
3358
3359 let after = q.get(&blocked.id).unwrap();
3360 assert_eq!(after.status, TaskStatus::Queued);
3361 assert!(after.blocked_by.is_empty());
3362 assert!(after.block_reason.is_none());
3363 let notes = notices_of(dir.path());
3364 assert_eq!(notes.len(), 1, "{notes:?}");
3365 assert_eq!(notes[0].severity, crate::notices::Severity::Info);
3366 }
3367
3368 #[test]
3369 fn removing_one_of_two_dependencies_keeps_the_other() {
3370 let (dir, q) = queue();
3371 let questions = Questions::at(dir.path().join("questions"));
3372 let mut dep = task("dependency");
3373 q.put(&mut dep).unwrap();
3374 let mut other = task("other");
3375 q.put(&mut other).unwrap();
3376 let mut blocked = task("waiting");
3377 blocked.block(
3378 vec![dep.id.clone(), other.id.clone()],
3379 Some("waits on both".to_owned()),
3380 );
3381 q.put(&mut blocked).unwrap();
3382
3383 let removed = q.remove(&dep.id, false, &questions).unwrap();
3384 assert!(removed.released.is_empty());
3385 assert_eq!(removed.still_blocked, [blocked.id.clone()]);
3386
3387 let after = q.get(&blocked.id).unwrap();
3388 assert_eq!(after.status, TaskStatus::Blocked);
3389 assert_eq!(after.blocked_by, [other.id.clone()]);
3390 assert!(missing_blockers(&q, &questions, &after.blocked_by).is_empty());
3393 }
3394
3395 #[test]
3396 fn a_dependency_deleted_before_its_dependents_were_rewritten_is_released_later() {
3397 let (dir, q) = queue();
3398 let questions = Questions::at(dir.path().join("questions"));
3399 let mut dep = task("dependency");
3400 q.put(&mut dep).unwrap();
3401 let mut blocked = task("waiting");
3402 blocked.block(vec![dep.id.clone()], None);
3403 q.put(&mut blocked).unwrap();
3404
3405 let claim = q.claim(&blocked.id).unwrap();
3407 let removed = q.remove(&dep.id, false, &questions).unwrap();
3408 assert!(removed.released.is_empty());
3409 drop(claim);
3410
3411 let mut task = q.get(&blocked.id).unwrap();
3412 assert_eq!(task.status, TaskStatus::Blocked);
3413 assert!(missing_blockers(&q, &questions, &task.blocked_by).is_empty());
3414 assert_eq!(q.apply_deleted_blockers(&mut task), [dep.id.clone()]);
3415 assert_eq!(task.status, TaskStatus::Queued);
3416 }
3417
3418 #[test]
3419 fn a_dependency_without_a_tombstone_is_still_missing() {
3420 let (dir, q) = queue();
3421 let questions = Questions::at(dir.path().join("questions"));
3422 let mut dep = task("dependency");
3423 q.put(&mut dep).unwrap();
3424 std::fs::remove_file(q.path_of(&dep.id)).unwrap();
3425 assert_eq!(
3426 missing_blockers(&q, &questions, std::slice::from_ref(&dep.id)),
3427 [dep.id.clone()]
3428 );
3429 assert!(deleted_blockers(&q, std::slice::from_ref(&dep.id)).is_empty());
3430 }
3431
3432 #[test]
3433 fn a_held_dependent_is_left_alone_by_a_dependency_deletion() {
3434 let mut t = task("held");
3435 t.hold_machine(Some("because".to_owned()));
3436 t.blocked_by = vec!["gone".to_owned()];
3437 assert!(!t.dependency_deleted("gone"));
3438 assert_eq!(t.status, TaskStatus::Held);
3439 }
3440
3441 #[test]
3442 fn a_failed_record_removal_takes_the_tombstone_back() {
3443 let (_dir, q) = queue();
3444 let mut t = task("stays");
3445 q.put(&mut t).unwrap();
3446 q.write_tombstone(&t.id).unwrap();
3447 let err = q.remove_record_with_attachments(&t.id, |_| Err(std::io::Error::other("nope")));
3448 assert!(err.is_err());
3449 assert!(!q.was_deleted(&t.id));
3452 }
3453
3454 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3455 let p = dir.join(name);
3456 std::fs::write(&p, body).unwrap();
3457 p
3458 }
3459
3460 #[test]
3461 fn an_attachment_copy_survives_deleting_its_source() {
3462 let (dir, q) = queue();
3463 let src = source_file(dir.path(), "shot.png", "pixels");
3464 let mut t = task("with a picture");
3465 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3466 q.put(&mut t).unwrap();
3467 std::fs::remove_file(&src).unwrap();
3468 assert_eq!(names, ["shot.png"]);
3469 let loaded = q.get(&t.id).unwrap();
3470 let paths = q.attachment_paths(&loaded);
3471 assert_eq!(paths.len(), 1);
3472 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3473 }
3474
3475 #[test]
3476 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3477 let (dir, q) = queue();
3478 let mut t = task("bad names");
3479 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3480 let src = source_file(dir.path(), name, "x");
3481 assert!(
3482 q.attach(&mut t, &[src]).is_err(),
3483 "`{name}` must be refused"
3484 );
3485 }
3486 assert!(!crate::ask::valid_asset_name("C:foo.png"));
3490 #[cfg(not(windows))]
3491 {
3492 let src = source_file(dir.path(), "C:foo.png", "x");
3493 assert!(q.attach(&mut t, &[src]).is_err());
3494 }
3495 let long = format!("{}.png", "a".repeat(70));
3496 let src = source_file(dir.path(), &long, "x");
3497 assert!(q.attach(&mut t, &[src]).is_err());
3498 assert!(t.attachments.is_empty());
3499 assert!(!q.attachments_dir(&t.id).exists());
3500 }
3501
3502 #[test]
3503 fn a_taken_attachment_name_is_numbered_not_overwritten() {
3504 let (dir, q) = queue();
3505 let a = source_file(dir.path(), "shot.png", "one");
3506 let sub = dir.path().join("other");
3507 std::fs::create_dir_all(&sub).unwrap();
3508 let b = source_file(&sub, "shot.png", "two");
3509 let mut t = task("collision");
3510 q.attach(&mut t, &[a]).unwrap();
3511 q.attach(&mut t, &[b]).unwrap();
3512 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3513 let paths = q.attachment_paths(&t);
3514 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3515 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3516 }
3517
3518 #[test]
3519 fn a_renumbered_name_stays_inside_the_length_bound() {
3520 let (dir, q) = queue();
3521 let name = format!("{}.png", "a".repeat(60));
3522 assert_eq!(name.len(), 64);
3523 let a = source_file(dir.path(), &name, "one");
3524 let sub = dir.path().join("other");
3525 std::fs::create_dir_all(&sub).unwrap();
3526 let b = source_file(&sub, &name, "two");
3527 let mut t = task("long");
3528 q.attach(&mut t, &[a, b]).unwrap();
3529 assert_eq!(t.attachments.len(), 2);
3530 assert!(
3531 t.attachments
3532 .iter()
3533 .all(|n| crate::ask::valid_asset_name(n))
3534 );
3535 assert!(t.attachments[1].ends_with("-2.png"));
3536 }
3537
3538 #[test]
3539 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3540 let (dir, q) = queue();
3541 let good = source_file(dir.path(), "good.png", "ok");
3542 let mut t = task("partial");
3543 q.attach(&mut t, &[good]).unwrap();
3544 let more = source_file(dir.path(), "more.png", "ok");
3545 let missing = dir.path().join("missing.png");
3546 assert!(q.attach(&mut t, &[more, missing]).is_err());
3547 assert_eq!(t.attachments, ["good.png"]);
3548 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3549 .unwrap()
3550 .flatten()
3551 .collect();
3552 assert_eq!(on_disk.len(), 1);
3553 }
3554
3555 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3557 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3558 std::fs::create_dir_all(&tmp).unwrap();
3559 tmp
3560 }
3561
3562 #[test]
3563 fn a_failed_put_leaves_no_new_attachment_directory() {
3564 let (dir, q) = queue();
3565 let mut t = task("fresh");
3566 let tmp = block_put(&q, &t);
3567 let src = source_file(dir.path(), "shot.png", "x");
3568 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3569 assert!(t.attachments.is_empty());
3570 assert!(!q.attachments_dir(&t.id).exists());
3571 assert!(!q.path_of(&t.id).exists());
3572 std::fs::remove_dir(tmp).unwrap();
3573 }
3574
3575 #[test]
3576 fn a_failed_put_removes_only_the_copy_it_just_made() {
3577 let (dir, q) = queue();
3578 let mut t = task("edited");
3579 let first = source_file(dir.path(), "first.png", "1");
3580 q.attach_and_put(&mut t, &[first]).unwrap();
3581 block_put(&q, &t);
3582 let second = source_file(dir.path(), "second.png", "2");
3583 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3584 assert_eq!(t.attachments, ["first.png"]);
3585 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3586 .unwrap()
3587 .flatten()
3588 .map(|e| e.file_name().to_string_lossy().into_owned())
3589 .collect();
3590 assert_eq!(on_disk, ["first.png"]);
3591 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3592 }
3593
3594 #[test]
3595 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3596 let (dir, q) = queue();
3597 let questions = Questions::at(dir.path().join("questions"));
3598 let gone = task("gone");
3599 let mut other = task("other");
3600 let mut live = task("live");
3601 q.put(&mut other).unwrap();
3602 q.put(&mut live).unwrap();
3603 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3606 std::fs::create_dir_all(&orphan).unwrap();
3607 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3608 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3611 std::fs::create_dir_all(&busy).unwrap();
3612
3613 q.remove(&other.id, false, &questions).unwrap();
3614 assert!(!orphan.exists(), "an orphan is swept");
3615 assert!(busy.exists(), "a removal in progress is left alone");
3616 }
3617
3618 #[test]
3619 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3620 let (dir, q) = queue();
3621 let questions = Questions::at(dir.path().join("questions"));
3622 let mut t = task("stuck");
3623 let src = source_file(dir.path(), "shot.png", "x");
3624 q.attach_and_put(&mut t, &[src]).unwrap();
3625 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3628 std::fs::create_dir_all(&aside).unwrap();
3629 std::fs::write(aside.join("old.png"), "o").unwrap();
3630 assert!(q.remove(&t.id, false, &questions).is_err());
3631 assert!(q.path_of(&t.id).exists());
3632 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3633 }
3634
3635 #[test]
3636 fn a_failed_record_removal_puts_the_attachments_back() {
3637 let (dir, q) = queue();
3638 let mut t = task("rollback");
3639 let src = source_file(dir.path(), "shot.png", "x");
3640 q.attach_and_put(&mut t, &[src]).unwrap();
3641 let err = q
3642 .remove_record_with_attachments(&t.id, |_| {
3643 Err(std::io::Error::other("injected failure"))
3644 })
3645 .unwrap_err();
3646 assert!(format!("{err:#}").contains("injected failure"));
3647 assert!(q.path_of(&t.id).exists());
3648 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3649 assert!(
3650 !q.root
3651 .join(format!("{}.attachments.removing", t.id))
3652 .exists()
3653 );
3654 }
3655
3656 #[test]
3657 fn editing_a_task_keeps_its_attachments() {
3658 let (dir, q) = queue();
3659 let src = source_file(dir.path(), "shot.png", "x");
3660 let mut t = task("editable");
3661 q.attach(&mut t, &[src]).unwrap();
3662 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3663 q.put(&mut t).unwrap();
3664 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3665 }
3666
3667 #[test]
3668 fn removing_a_task_deletes_its_attachments() {
3669 let (dir, q) = queue();
3670 let questions = Questions::at(dir.path().join("questions"));
3671 let src = source_file(dir.path(), "shot.png", "x");
3672 let mut t = task("doomed");
3673 q.attach(&mut t, &[src]).unwrap();
3674 q.put(&mut t).unwrap();
3675 assert!(q.attachments_dir(&t.id).is_dir());
3676 q.remove(&t.id, false, &questions).unwrap();
3677 assert!(!q.attachments_dir(&t.id).exists());
3678 assert!(q.list().is_empty());
3679 }
3680
3681 #[test]
3682 fn a_task_written_before_attachments_still_reads() {
3683 let (_dir, q) = queue();
3684 let mut t = task("old");
3685 q.put(&mut t).unwrap();
3686 let path = q.path_of(&t.id);
3687 let mut v: serde_json::Value =
3688 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3689 v.as_object_mut().unwrap().remove("attachments");
3690 std::fs::write(&path, v.to_string()).unwrap();
3691 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3692 }
3693
3694 #[test]
3695 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3696 let q = Queue::at(PathBuf::from("relative-queue"));
3697 let mut t = task("rel");
3698 t.attachments.push("shot.png".to_owned());
3699 let paths = q.attachment_paths(&t);
3700 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3701 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3702 }
3703
3704 #[test]
3705 fn link_run_adds_a_run_once_and_touches_nothing_else() {
3706 let dir = tempfile::tempdir().unwrap();
3707 let queue = Queue::at(dir.path().join("queue"));
3708 let mut t = Task::new(
3709 "t".to_owned(),
3710 "do it".to_owned(),
3711 PathBuf::from("."),
3712 Source::Human,
3713 );
3714 queue.put(&mut t).unwrap();
3715 let before = queue.get(&t.id).unwrap();
3716
3717 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3718 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3719 assert_eq!(linked.status, before.status);
3720 assert_eq!(linked.attempts, before.attempts);
3721 assert_eq!(linked.interrupt, before.interrupt);
3722
3723 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3724 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3725 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3726 assert!(queue.link_run("no-such-task", "r").is_err());
3727 }
3728
3729 #[test]
3730 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
3731 let dir = tempfile::tempdir().unwrap();
3732 let queue = Queue::at(dir.path().join("queue"));
3733 let mut t = Task::new(
3734 "t".to_owned(),
3735 "do it".to_owned(),
3736 PathBuf::from("."),
3737 Source::Human,
3738 );
3739 queue.put(&mut t).unwrap();
3740 let mut snapshot = queue.get(&t.id).unwrap();
3742 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3743
3744 snapshot.start("20260930-000000-aaaa".to_owned());
3745 queue.put(&mut snapshot).unwrap();
3746
3747 let stored = queue.get(&t.id).unwrap();
3748 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3749 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3750 assert_eq!(stored.attempts, 1);
3751 }
3752
3753 #[test]
3754 fn concurrent_links_and_daemon_saves_lose_nothing() {
3755 let dir = tempfile::tempdir().unwrap();
3756 let queue = Queue::at(dir.path().join("queue"));
3757 let mut t = Task::new(
3758 "t".to_owned(),
3759 "do it".to_owned(),
3760 PathBuf::from("."),
3761 Source::Human,
3762 );
3763 queue.put(&mut t).unwrap();
3764 let id = t.id.clone();
3765
3766 let linkers: Vec<_> = (0..4)
3767 .map(|n| {
3768 let (queue, id) = (queue.clone(), id.clone());
3769 std::thread::spawn(move || {
3770 for k in 0..10 {
3771 queue
3772 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3773 .unwrap();
3774 }
3775 })
3776 })
3777 .collect();
3778 let mut mine = queue.get(&id).unwrap();
3781 for k in 0..10 {
3782 mine.start(format!("20260930-000009-d{k:03}"));
3783 queue.put(&mut mine).unwrap();
3784 }
3785 for l in linkers {
3786 l.join().unwrap();
3787 }
3788
3789 let stored = queue.get(&id).unwrap();
3790 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
3791 assert_eq!(
3792 stored.attempts, 10,
3793 "linking never rewinds the daemon's work"
3794 );
3795 }
3796
3797 fn age_lock(path: &Path) {
3798 let f = std::fs::OpenOptions::new().write(true).open(path).unwrap();
3799 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
3800 .unwrap();
3801 }
3802
3803 #[test]
3804 fn concurrent_stale_takeover_yields_one_holder() {
3805 use std::sync::atomic::{AtomicUsize, Ordering};
3806 use std::sync::{Arc, Barrier};
3807 for _ in 0..5 {
3808 let dir = tempfile::tempdir().unwrap();
3809 let q = Queue::at(dir.path().to_path_buf());
3810 let lock = dir.path().join("t.write-lock");
3811 std::fs::write(&lock, "dead-0000").unwrap();
3812 age_lock(&lock);
3813 let n = 6;
3814 let barrier = Arc::new(Barrier::new(n));
3815 let (now, max) = (Arc::new(AtomicUsize::new(0)), Arc::new(AtomicUsize::new(0)));
3816 let handles: Vec<_> = (0..n)
3817 .map(|_| {
3818 let (q, b, now, max) = (q.clone(), barrier.clone(), now.clone(), max.clone());
3819 std::thread::spawn(move || {
3820 b.wait();
3821 let g = q.lock_task("t").unwrap();
3822 let held = now.fetch_add(1, Ordering::SeqCst) + 1;
3823 max.fetch_max(held, Ordering::SeqCst);
3824 std::thread::sleep(std::time::Duration::from_millis(20));
3825 now.fetch_sub(1, Ordering::SeqCst);
3826 drop(g);
3827 })
3828 })
3829 .collect();
3830 for h in handles {
3831 h.join().unwrap();
3832 }
3833 assert_eq!(max.load(Ordering::SeqCst), 1);
3834 assert!(!lock.exists());
3835 }
3836 }
3837
3838 #[test]
3839 fn dropping_a_stolen_lock_leaves_the_new_holders_lock() {
3840 let dir = tempfile::tempdir().unwrap();
3841 let q = Queue::at(dir.path().to_path_buf());
3842 let lock = dir.path().join("t.write-lock");
3843 let a = q.lock_task("t").unwrap();
3844 age_lock(&lock);
3845 let b = q.lock_task("t").unwrap();
3846 assert_ne!(a.token, b.token);
3847 drop(a);
3848 assert_eq!(std::fs::read_to_string(&lock).unwrap(), b.token);
3849 drop(b);
3850 assert!(!lock.exists());
3851 }
3852
3853 #[test]
3854 fn break_stale_leaves_a_lock_that_replaced_the_one_judged() {
3855 let dir = tempfile::tempdir().unwrap();
3856 let lock = dir.path().join("t.write-lock");
3857 std::fs::write(&lock, "new-token").unwrap();
3858 assert!(!break_stale(&lock, "old-token"));
3859 assert_eq!(std::fs::read_to_string(&lock).unwrap(), "new-token");
3860 assert!(!break_stale(&lock, "new-token"));
3862 assert!(lock.exists());
3863 age_lock(&lock);
3864 assert!(break_stale(&lock, "new-token"));
3865 assert!(!lock.exists());
3866 }
3867
3868 #[test]
3869 fn a_stale_break_marker_is_recovered_and_a_fresh_one_is_respected() {
3870 let dir = tempfile::tempdir().unwrap();
3871 let lock = dir.path().join("t.write-lock");
3872 std::fs::write(&lock, "dead-1").unwrap();
3873 age_lock(&lock);
3874 let marker = break_marker(&lock, "dead-1");
3875 std::fs::write(&marker, "crashed-remover").unwrap();
3876 assert!(!break_stale(&lock, "dead-1"));
3878 assert!(lock.exists());
3879 age_lock(&marker);
3881 assert!(break_stale(&lock, "dead-1"));
3882 assert!(!lock.exists());
3883 assert!(!marker.exists());
3884 }
3885}