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 = 12;
110
111#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
113#[serde(rename_all = "lowercase")]
114pub enum HoldSource {
115 Manual,
117 Machine,
119}
120
121impl HoldSource {
122 pub fn label(self) -> &'static str {
124 match self {
125 Self::Manual => "manual",
126 Self::Machine => "machine",
127 }
128 }
129}
130
131#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
134#[serde(tag = "kind", rename_all = "lowercase")]
135pub enum Source {
136 Human,
138 Agent {
141 run: String,
143 node: String,
145 },
146 Issue {
148 number: u64,
150 repo: String,
152 },
153}
154
155pub const CHAT_NODE: &str = "chat";
158
159impl Source {
160 pub fn label(&self) -> String {
162 match self {
163 Self::Human => "human".to_owned(),
164 Self::Agent { run, node } => format!("{node}@{}", short(run)),
165 Self::Issue { number, .. } => format!("issue #{number}"),
166 }
167 }
168}
169
170#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
172#[serde(rename_all = "lowercase")]
173pub enum TaskStatus {
174 Queued,
176 Running,
178 Done,
180 Failed,
182 Held,
184 Blocked,
188 Parked,
193}
194
195impl TaskStatus {
196 pub fn runnable(self) -> bool {
198 matches!(self, Self::Queued | Self::Failed | Self::Parked)
199 }
200
201 pub fn as_str(self) -> &'static str {
203 match self {
204 Self::Queued => "queued",
205 Self::Running => "running",
206 Self::Done => "done",
207 Self::Failed => "failed",
208 Self::Held => "held",
209 Self::Blocked => "blocked",
210 Self::Parked => "parked",
211 }
212 }
213}
214
215#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
218pub struct TaskCounts {
219 pub queued: usize,
221 pub running: usize,
223 pub done: usize,
225 pub failed: usize,
227 pub held: usize,
229 pub blocked: usize,
231 pub parked: usize,
233}
234
235impl TaskCounts {
236 pub fn of(tasks: &[Task]) -> Self {
240 let mut counts = Self::default();
241 for t in tasks {
242 match t.status {
243 TaskStatus::Queued => counts.queued += 1,
244 TaskStatus::Running => counts.running += 1,
245 TaskStatus::Done => counts.done += 1,
246 TaskStatus::Failed => counts.failed += 1,
247 TaskStatus::Held => counts.held += 1,
248 TaskStatus::Blocked => counts.blocked += 1,
249 TaskStatus::Parked => counts.parked += 1,
250 }
251 }
252 counts
253 }
254}
255
256#[derive(Debug, Clone, Serialize, Deserialize)]
258#[serde(deny_unknown_fields)]
259pub struct Task {
260 pub schema: u32,
262 pub id: String,
264 pub title: String,
266 pub instruction: String,
268 pub repo: PathBuf,
270 pub source: Source,
272 #[serde(default)]
274 pub priority: i32,
275 #[serde(default)]
286 pub solo: bool,
287 pub status: TaskStatus,
289 #[serde(default)]
291 pub attempts: usize,
292 #[serde(default)]
294 pub runs: Vec<String>,
295 #[serde(default)]
297 pub last_error: Option<String>,
298 #[serde(default)]
311 pub hold_reason: Option<String>,
312 #[serde(default)]
315 pub hold_source: Option<HoldSource>,
316 #[serde(default)]
328 pub diagnostic: Option<String>,
329 #[serde(default)]
339 pub blocked_by: Vec<String>,
340 #[serde(default)]
343 pub block_reason: Option<String>,
344 #[serde(default)]
357 pub blocked_from: Option<TaskStatus>,
358 #[serde(default)]
369 pub answers: Vec<AnsweredQuestion>,
370 #[serde(default)]
376 pub triage_applied: Vec<String>,
377 #[serde(default)]
383 pub actions_applied: Vec<String>,
384 #[serde(default)]
389 pub resume_override: Option<OperatorResume>,
390 #[serde(default)]
402 pub review_branch: Option<String>,
403 #[serde(default)]
406 pub fresh_start: bool,
407 #[serde(default)]
418 pub interrupt: bool,
419 #[serde(default)]
442 pub urgent: bool,
443 #[serde(default)]
449 pub attachments: Vec<String>,
450 #[serde(default)]
454 pub followup: Option<FollowUp>,
455 #[serde(default)]
461 pub origin_chat: Option<String>,
462 #[serde(default)]
468 pub overrides: Option<RunOverrides>,
469 #[serde(default)]
476 pub review_of: Option<String>,
477 #[serde(default)]
482 pub held_at: Option<Timestamp>,
483 #[serde(default)]
487 pub park_reason: Option<String>,
488 pub created_at: Timestamp,
490 pub updated_at: Timestamp,
492}
493
494#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
496pub struct RunOverrides {
497 #[serde(default)]
499 pub merge: Option<String>,
500 #[serde(default)]
502 pub candidates: Option<usize>,
503 #[serde(default)]
505 pub judges: Option<usize>,
506 #[serde(default)]
508 pub reviewers: Option<usize>,
509 #[serde(default)]
511 pub review_rounds: Option<usize>,
512 #[serde(default)]
514 pub seed: Option<u64>,
515 #[serde(default)]
517 pub config: Option<PathBuf>,
518}
519
520impl RunOverrides {
521 pub fn merge_over(&mut self, other: &RunOverrides) {
523 macro_rules! take {
524 ($($f:ident),*) => {$(
525 if other.$f.is_some() {
526 self.$f.clone_from(&other.$f);
527 }
528 )*};
529 }
530 take!(
531 merge,
532 candidates,
533 judges,
534 reviewers,
535 review_rounds,
536 seed,
537 config
538 );
539 }
540
541 pub fn apply(&self, config: &mut crate::config::Config) {
543 if let Some(n) = self.candidates {
544 config.graph.implementers = n;
545 }
546 if let Some(n) = self.judges {
547 config.graph.judges = n;
548 }
549 if let Some(n) = self.reviewers {
550 config.graph.reviewers = n;
551 }
552 if let Some(n) = self.review_rounds {
553 config.graph.review_rounds = n;
554 }
555 if let Some(m) = &self.merge
556 && let Ok(mode) = crate::daemon::merge_mode(m)
557 {
558 config.merge.mode = mode;
559 }
560 if let Some(s) = self.seed {
561 config.blind.seed = Some(s);
562 }
563 }
564}
565
566#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
568pub struct FollowUp {
569 pub run: String,
571 #[serde(default)]
573 pub origin_task: Option<String>,
574 pub pr: String,
576 pub findings: Vec<String>,
578 pub generation: u32,
581}
582
583#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
586pub struct OperatorResume {
587 pub question_id: String,
589 pub at: Timestamp,
591 #[serde(default)]
594 pub conductor_rehold: Option<String>,
595 #[serde(default)]
598 pub forced: bool,
599 #[serde(default)]
606 pub pinned_run: Option<String>,
607}
608
609#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
612pub struct AnsweredQuestion {
613 pub question: String,
615 pub answer: String,
617}
618
619impl Task {
620 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
622 let now = Timestamp::now();
623 let origin_chat = match &source {
624 Source::Agent { run, node } if node == CHAT_NODE => Some(run.clone()),
625 _ => None,
626 };
627 Self {
628 schema: SCHEMA,
629 id: new_id(),
630 title,
631 instruction,
632 repo,
633 source,
634 priority: 0,
635 solo: false,
636 status: TaskStatus::Queued,
637 attempts: 0,
638 runs: Vec::new(),
639 last_error: None,
640 hold_reason: None,
641 hold_source: None,
642 diagnostic: None,
643 blocked_by: Vec::new(),
644 block_reason: None,
645 blocked_from: None,
646 answers: Vec::new(),
647 triage_applied: Vec::new(),
648 actions_applied: Vec::new(),
649 resume_override: None,
650 review_branch: None,
651 fresh_start: false,
652 interrupt: false,
653 urgent: false,
654 attachments: Vec::new(),
655 followup: None,
656 origin_chat,
657 overrides: None,
658 review_of: None,
659 held_at: None,
660 park_reason: None,
661 created_at: now,
662 updated_at: now,
663 }
664 }
665
666 pub fn chat_talk(&self) -> Option<&str> {
669 if let Some(id) = &self.origin_chat {
670 return Some(id);
671 }
672 match &self.source {
673 Source::Agent { run, node } if node == CHAT_NODE => Some(run),
674 _ => None,
675 }
676 }
677
678 pub fn short(&self) -> &str {
680 short(&self.id)
681 }
682
683 pub fn mark_triage_applied(&mut self, question_id: &str) {
685 if !self.triage_applied(question_id) {
686 self.triage_applied.push(question_id.to_owned());
687 }
688 }
689
690 pub fn action_applied(&self, question_id: &str) -> bool {
693 self.actions_applied.iter().any(|id| id == question_id)
694 }
695
696 pub fn mark_action_applied(&mut self, question_id: &str) {
698 if !self.action_applied(question_id) {
699 self.actions_applied.push(question_id.to_owned());
700 }
701 }
702
703 pub fn triage_applied(&self, question_id: &str) -> bool {
705 self.triage_applied.iter().any(|id| id == question_id)
706 }
707
708 pub fn start(&mut self, run: String) {
719 self.status = TaskStatus::Running;
720 self.attempts += 1;
721 self.runs.push(run);
722 self.last_error = None;
723 self.park_reason = None;
724 self.fresh_start = false;
725 self.interrupt = false;
726 self.resume_override = None;
728 }
729
730 pub fn link_run(&mut self, run: &str) -> bool {
735 if self.runs.iter().any(|r| r == run) {
736 return false;
737 }
738 self.runs.push(run.to_owned());
739 true
740 }
741
742 pub fn succeed(&mut self) {
752 self.status = TaskStatus::Done;
753 self.resume_override = None;
754 self.last_error = None;
755 self.park_reason = None;
756 self.hold_reason = None;
757 self.hold_source = None;
758 self.diagnostic = None;
759 self.blocked_by.clear();
760 self.block_reason = None;
761 self.blocked_from = None;
762 }
763
764 pub fn already_landed(&mut self, note: impl Into<String>) {
772 self.succeed();
773 self.attempts = self.attempts.saturating_sub(1);
774 self.last_error = Some(note.into());
775 }
776
777 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
807 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
808 return &[];
809 }
810 &self.runs[..self.runs.len() - 1]
811 }
812
813 pub fn successor_of(&self, run: &str) -> Option<&String> {
820 let pos = self.runs.iter().position(|r| r == run)?;
821 self.runs.get(pos + 1)
822 }
823
824 pub fn earlier_attempts(&self) -> &[String] {
831 &self.runs
832 }
833
834 fn note_held(&mut self) {
837 if self.status != TaskStatus::Held || self.held_at.is_none() {
838 self.held_at = Some(Timestamp::now());
839 }
840 }
841
842 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
858 let why = why.into();
859 self.diagnostic = None;
860 self.park_reason = None;
861 self.status = if self.attempts >= max_attempts {
862 self.note_held();
863 self.hold_source = Some(HoldSource::Machine);
864 self.hold_reason = Some(why.clone());
865 TaskStatus::Held
866 } else {
867 TaskStatus::Failed
868 };
869 self.last_error = Some(why);
870 }
871
872 pub fn stall(&mut self, why: impl Into<String>) {
882 self.last_error = Some(why.into());
883 self.diagnostic = None;
884 self.attempts = self.attempts.saturating_sub(1);
885 self.status = TaskStatus::Failed;
886 }
887
888 pub fn park(&mut self, why: impl Into<String>) {
894 self.diagnostic = None;
895 self.attempts = self.attempts.saturating_sub(1);
896 self.status = TaskStatus::Parked;
897 self.park_reason = Some(why.into());
898 self.schema = self.schema.max(SCHEMA);
901 }
902
903 pub fn operator_held(&self) -> bool {
909 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
910 }
911
912 pub fn hold_manual(&mut self, reason: Option<String>) {
922 self.park_reason = None;
923 self.note_held();
924 self.status = TaskStatus::Held;
925 if reason.is_some() {
926 self.hold_reason = reason;
927 }
928 self.hold_source = Some(HoldSource::Manual);
929 self.blocked_by.clear();
930 self.block_reason = None;
931 self.blocked_from = None;
932 }
933
934 pub fn hold_machine(&mut self, reason: Option<String>) {
939 self.park_reason = None;
940 self.note_held();
941 self.status = TaskStatus::Held;
942 if reason.is_some() {
943 self.hold_reason = reason;
944 }
945 self.hold_source = Some(HoldSource::Machine);
946 self.blocked_by.clear();
947 self.block_reason = None;
948 self.blocked_from = None;
949 }
950
951 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
961 if self.status != TaskStatus::Blocked {
962 self.blocked_from = Some(self.status);
963 }
964 self.status = TaskStatus::Blocked;
965 self.blocked_by = blocked_by;
966 self.block_reason = reason;
967 }
968
969 pub fn unblock(&mut self, resolved_id: &str) {
992 if self.status != TaskStatus::Blocked {
993 return;
994 }
995 self.blocked_by.retain(|id| id != resolved_id);
996 self.restore_if_unblocked();
997 }
998
999 pub fn dependency_deleted(&mut self, deleted_id: &str) -> bool {
1010 if self.status != TaskStatus::Blocked || !self.blocked_by.iter().any(|b| b == deleted_id) {
1011 return false;
1012 }
1013 self.blocked_by.retain(|id| id != deleted_id);
1014 self.restore_if_unblocked();
1015 true
1016 }
1017
1018 fn restore_if_unblocked(&mut self) {
1022 if self.blocked_by.is_empty() {
1023 self.status = match self.blocked_from {
1024 Some(TaskStatus::Running) => TaskStatus::Queued,
1025 Some(other) => other,
1026 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
1027 TaskStatus::Held
1028 }
1029 None => TaskStatus::Queued,
1030 };
1031 self.block_reason = None;
1032 self.blocked_from = None;
1033 }
1034 }
1035
1036 pub fn record_answer(&mut self, question: String, answer: String) {
1041 self.answers.push(AnsweredQuestion { question, answer });
1042 }
1043
1044 pub fn request_review(&mut self, branch: String) {
1048 self.release();
1049 self.review_branch = Some(branch);
1050 }
1051
1052 pub fn requeue(&mut self) {
1055 self.release();
1056 self.review_branch = None;
1058 self.fresh_start = true;
1059 }
1060
1061 pub fn hold_for_handover(&mut self, branch: Option<String>, reason: String) {
1070 if branch.is_some() {
1071 self.review_branch = branch;
1072 }
1073 self.hold_machine(Some(reason));
1074 }
1075
1076 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
1085 if self.status == TaskStatus::Running {
1086 bail!(
1087 "task {} is running; its priority cannot be changed until \
1088 this attempt finishes",
1089 self.short()
1090 );
1091 }
1092 self.priority = priority;
1093 Ok(())
1094 }
1095
1096 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
1111 if interrupt && !self.status.runnable() {
1112 bail!(
1113 "task {} is {}; only a queued or failed task can be marked \
1114 to interrupt",
1115 self.short(),
1116 self.status.as_str()
1117 );
1118 }
1119 self.interrupt = interrupt;
1120 Ok(())
1121 }
1122
1123 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
1135 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
1136 bail!(
1137 "task {} is {}; only a queued or held task's instruction can \
1138 be edited",
1139 self.short(),
1140 self.status.as_str()
1141 );
1142 }
1143 self.title = title;
1144 self.instruction = instruction;
1145 Ok(())
1146 }
1147
1148 pub fn handed_off(&mut self, why: impl Into<String>) {
1168 let why = why.into();
1169 self.diagnostic = None;
1170 self.note_held();
1171 self.status = TaskStatus::Held;
1172 self.hold_source = Some(HoldSource::Machine);
1173 self.hold_reason = Some(why.clone());
1174 self.last_error = Some(why);
1175 }
1176
1177 pub fn release(&mut self) {
1181 let refused_handover = self.status == TaskStatus::Held
1182 && self.hold_source == Some(HoldSource::Machine)
1183 && self.review_branch.is_some();
1184 self.status = TaskStatus::Queued;
1185 self.held_at = None;
1186 self.attempts = 0;
1187 self.last_error = None;
1188 self.park_reason = None;
1189 self.hold_reason = None;
1192 self.hold_source = None;
1193 self.diagnostic = None;
1194 self.blocked_by.clear();
1199 self.block_reason = None;
1200 self.blocked_from = None;
1201 if !refused_handover {
1205 self.review_branch = None;
1206 }
1207 self.fresh_start = false;
1208 }
1209}
1210
1211fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
1215 use std::io::ErrorKind;
1216 let (stem, ext) = match name.rfind('.') {
1217 Some(i) if i > 0 => (&name[..i], &name[i..]),
1218 _ => (name, ""),
1219 };
1220 for n in 1u32.. {
1221 let candidate = if n == 1 {
1222 name.to_owned()
1223 } else {
1224 let suffix = format!("-{n}");
1225 let room = 64usize.saturating_sub(suffix.len() + ext.len());
1226 let stem: String = stem.chars().take(room).collect();
1227 format!("{stem}{suffix}{ext}")
1228 };
1229 if !crate::ask::valid_asset_name(&candidate) {
1230 bail!("no valid attachment name is left for `{name}`");
1231 }
1232 let path = dir.join(&candidate);
1233 match std::fs::OpenOptions::new()
1234 .write(true)
1235 .create_new(true)
1236 .open(&path)
1237 {
1238 Ok(mut out) => {
1239 let copied = std::fs::File::open(src)
1240 .and_then(|mut input| std::io::copy(&mut input, &mut out));
1241 if let Err(e) = copied {
1242 drop(out);
1243 let _ = std::fs::remove_file(&path);
1244 return Err(e).with_context(|| format!("copy {}", src.display()));
1245 }
1246 return Ok((candidate, path));
1247 }
1248 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
1249 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
1250 }
1251 }
1252 unreachable!("the counter never runs out")
1253}
1254
1255const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
1258
1259struct TaskLock {
1262 path: PathBuf,
1263 token: String,
1264}
1265
1266fn owner_token() -> String {
1269 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1270 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1271 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy() ^ n.rotate_left(32));
1272 format!("{}-{:016x}", std::process::id(), r.next_u64())
1273}
1274
1275fn break_marker(path: &Path, token: &str) -> PathBuf {
1278 let mut h: u64 = 0xcbf29ce484222325;
1279 for b in token.bytes() {
1280 h ^= u64::from(b);
1281 h = h.wrapping_mul(0x100000001b3);
1282 }
1283 let mut name = path.as_os_str().to_owned();
1284 name.push(format!(".break-{h:016x}"));
1285 PathBuf::from(name)
1286}
1287
1288fn older_than_stale(path: &Path) -> bool {
1289 std::fs::metadata(path)
1290 .and_then(|m| m.modified())
1291 .ok()
1292 .and_then(|t| t.elapsed().ok())
1293 .is_some_and(|age| age > TASK_LOCK_STALE)
1294}
1295
1296const MAX_MARKER_DEPTH: u8 = 3;
1300
1301fn take_marker(path: &Path, token: &str, depth: u8) -> Option<(PathBuf, String)> {
1307 use std::io::Write;
1308 let marker = break_marker(path, token);
1309 for _ in 0..2 {
1310 match std::fs::OpenOptions::new()
1311 .write(true)
1312 .create_new(true)
1313 .open(&marker)
1314 {
1315 Ok(mut f) => {
1316 let mine = owner_token();
1317 if f.write_all(mine.as_bytes()).is_err() {
1318 drop(f);
1319 let _ = std::fs::remove_file(&marker);
1320 return None;
1321 }
1322 return Some((marker, mine));
1323 }
1324 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1325 let judged = std::fs::read_to_string(&marker).ok();
1326 match judged.filter(|_| older_than_stale(&marker)) {
1327 Some(judged) if depth < MAX_MARKER_DEPTH => {
1328 remove_lock_if(&marker, &judged, true, depth + 1)?;
1329 }
1330 _ => return None,
1331 }
1332 }
1333 Err(_) => return None,
1334 }
1335 }
1336 None
1337}
1338
1339fn remove_lock_if(path: &Path, token: &str, require_stale: bool, depth: u8) -> Option<bool> {
1344 let (marker, mine) = take_marker(path, token, depth)?;
1345 let still = std::fs::read_to_string(path).is_ok_and(|c| c == token)
1346 && (!require_stale || older_than_stale(path));
1347 if still {
1348 let _ = std::fs::remove_file(path);
1349 }
1350 if depth >= MAX_MARKER_DEPTH {
1352 if std::fs::read_to_string(&marker).is_ok_and(|c| c == mine) {
1353 let _ = std::fs::remove_file(&marker);
1354 }
1355 } else {
1356 let _ = remove_lock_if(&marker, &mine, false, depth + 1);
1357 }
1358 Some(still)
1359}
1360
1361fn break_stale(path: &Path, judged: &str) -> bool {
1364 remove_lock_if(path, judged, true, 0).unwrap_or(false)
1365}
1366
1367impl Drop for TaskLock {
1368 fn drop(&mut self) {
1369 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
1370 loop {
1371 match remove_lock_if(&self.path, &self.token, false, 0) {
1372 Some(_) => return,
1373 None if std::time::Instant::now() > deadline => return,
1374 None => std::thread::sleep(std::time::Duration::from_millis(5)),
1375 }
1376 }
1377 }
1378}
1379
1380#[derive(Debug, Clone)]
1382pub struct Queue {
1383 root: PathBuf,
1384}
1385
1386impl Queue {
1387 pub fn open() -> Self {
1389 Self::at(crate::run::home().join("queue"))
1390 }
1391
1392 pub fn at(root: PathBuf) -> Self {
1395 Self { root }
1396 }
1397
1398 pub fn root(&self) -> &Path {
1400 &self.root
1401 }
1402
1403 pub fn path_of(&self, id: &str) -> PathBuf {
1405 self.root.join(format!("{id}.json"))
1406 }
1407
1408 pub fn attachments_dir(&self, id: &str) -> PathBuf {
1411 self.root.join(format!("{id}.attachments"))
1412 }
1413
1414 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1425 let mut wanted = Vec::new();
1426 for src in sources {
1427 let name = src
1428 .file_name()
1429 .and_then(|n| n.to_str())
1430 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1431 if !crate::ask::valid_asset_name(name) {
1432 bail!(
1433 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1434 with no `..`; rename the file and try again"
1435 );
1436 }
1437 if !src.is_file() {
1438 bail!("attachment `{}` is not a file", src.display());
1439 }
1440 wanted.push((src, name));
1441 }
1442 let dir = self.attachments_dir(&task.id);
1443 let existed = dir.is_dir();
1444 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1445 let mut created: Vec<PathBuf> = Vec::new();
1446 let mut names = Vec::new();
1447 let mut copy_all = || -> Result<()> {
1448 for (src, name) in &wanted {
1449 let (stored, path) = copy_new(&dir, src, name)?;
1450 created.push(path);
1451 names.push(stored);
1452 }
1453 Ok(())
1454 };
1455 if let Err(e) = copy_all() {
1456 for path in &created {
1457 let _ = std::fs::remove_file(path);
1458 }
1459 if !existed {
1460 let _ = std::fs::remove_dir(&dir);
1461 }
1462 return Err(e);
1463 }
1464 task.attachments.extend(names.iter().cloned());
1465 Ok(names)
1466 }
1467
1468 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1479 let dir = self.attachments_dir(&task.id);
1480 let existed = dir.is_dir();
1481 let before = task.attachments.len();
1482 let names = self.attach(task, sources)?;
1483 if let Err(e) = self.put(task) {
1484 for name in &names {
1485 let _ = std::fs::remove_file(dir.join(name));
1486 }
1487 if !existed {
1488 let _ = std::fs::remove_dir(&dir);
1489 }
1490 task.attachments.truncate(before);
1491 return Err(e);
1492 }
1493 Ok(names)
1494 }
1495
1496 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1500 let dir = self.attachments_dir(&task.id);
1501 task.attachments
1502 .iter()
1503 .map(|n| {
1504 let p = dir.join(n);
1505 std::path::absolute(&p).unwrap_or(p)
1506 })
1507 .collect()
1508 }
1509
1510 pub fn put(&self, task: &mut Task) -> Result<()> {
1519 let _lock = self.lock_task(&task.id)?;
1520 self.put_unlocked(task)
1521 }
1522
1523 pub fn create_new(&self, task: &mut Task) -> Result<bool> {
1528 let _lock = self.lock_task(&task.id)?;
1529 if self.path_of(&task.id).exists() {
1530 return Ok(false);
1531 }
1532 self.put_unlocked(task)?;
1533 Ok(true)
1534 }
1535
1536 fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1538 if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1539 for run in stored.runs {
1540 if !task.runs.contains(&run) {
1541 task.runs.push(run);
1542 }
1543 }
1544 }
1545 task.updated_at = Timestamp::now();
1546 std::fs::create_dir_all(&self.root)
1547 .with_context(|| format!("create {}", self.root.display()))?;
1548 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1549 let path = self.path_of(&task.id);
1550 let tmp = path.with_extension("json.tmp");
1551 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1552 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1553 if let (Some(notice), Some(home)) = (
1558 crate::notices::task_held(task),
1559 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1560 ) {
1561 crate::notices::raise_in(home, notice);
1562 }
1563 Ok(())
1564 }
1565
1566 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1571 let id = self.resolve_id(id)?;
1572 let _lock = self.lock_task(&id)?;
1575 let mut task = self.get(&id)?;
1576 if task.link_run(run) {
1577 self.put_unlocked(&mut task)?;
1578 }
1579 Ok(task)
1580 }
1581
1582 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1590 std::fs::create_dir_all(&self.root)
1591 .with_context(|| format!("create {}", self.root.display()))?;
1592 let path = self.root.join(format!("{id}.write-lock"));
1593 let started = std::time::Instant::now();
1594 loop {
1595 match std::fs::OpenOptions::new()
1596 .write(true)
1597 .create_new(true)
1598 .open(&path)
1599 {
1600 Ok(mut file) => {
1601 use std::io::Write;
1602 let token = owner_token();
1603 if let Err(e) = file.write_all(token.as_bytes()) {
1604 drop(file);
1605 let _ = std::fs::remove_file(&path);
1606 return Err(e).with_context(|| format!("lock {}", path.display()));
1607 }
1608 return Ok(TaskLock { path, token });
1609 }
1610 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1611 let judged = std::fs::read_to_string(&path).ok();
1614 let broken = judged
1615 .filter(|_| older_than_stale(&path))
1616 .is_some_and(|judged| break_stale(&path, &judged));
1617 if !broken {
1621 if started.elapsed() > TASK_LOCK_STALE {
1622 bail!("could not lock task {id}");
1623 }
1624 std::thread::sleep(std::time::Duration::from_millis(15));
1625 }
1626 }
1627 Err(e)
1631 if e.kind() == std::io::ErrorKind::PermissionDenied
1632 && started.elapsed() <= TASK_LOCK_STALE =>
1633 {
1634 std::thread::sleep(std::time::Duration::from_millis(15));
1635 }
1636 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1637 }
1638 }
1639 }
1640
1641 pub fn get(&self, id: &str) -> Result<Task> {
1643 let resolved = self.resolve_id(id)?;
1644 read_path(&self.path_of(&resolved))
1645 }
1646
1647 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1679 let _ = questions;
1680 let resolved = self.resolve_id(id)?;
1681 if in_flight {
1682 bail!("task {resolved} is being run by a live daemon right now");
1683 }
1684 self.write_tombstone(&resolved)?;
1685 if let Err(e) = self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))
1686 {
1687 if self.path_of(&resolved).exists() {
1691 let _ = std::fs::remove_file(self.tombstone_path(&resolved));
1692 }
1693 return Err(e);
1694 }
1695 let (released, still_blocked) = self.release_dependents_of(&resolved);
1696 Ok(Removal {
1697 id: resolved,
1698 released,
1699 still_blocked,
1700 })
1701 }
1702
1703 fn tombstone_path(&self, id: &str) -> PathBuf {
1706 self.root.join(format!("{id}.removed"))
1707 }
1708
1709 fn write_tombstone(&self, id: &str) -> Result<()> {
1710 let path = self.tombstone_path(id);
1711 let tmp = self.root.join(format!("{id}.removed.tmp"));
1712 std::fs::write(&tmp, b"")
1713 .and_then(|()| std::fs::rename(&tmp, &path))
1714 .with_context(|| format!("write {}", path.display()))
1715 }
1716
1717 fn was_deleted(&self, id: &str) -> bool {
1721 self.tombstone_path(id).is_file() && !self.path_of(id).exists()
1722 }
1723
1724 pub fn apply_deleted_blockers(&self, task: &mut Task) -> Vec<String> {
1728 let deleted = deleted_blockers(self, &task.blocked_by);
1729 deleted
1730 .into_iter()
1731 .filter(|id| task.dependency_deleted(id))
1732 .collect()
1733 }
1734
1735 pub fn note_dependency_deleted(&self, dependent: &Task, deleted: &str) {
1739 let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty()) else {
1740 return;
1741 };
1742 let ja = crate::lang::is_japanese(&crate::lang::of_repo(&dependent.repo));
1743 let waiting = dependent.status == TaskStatus::Blocked;
1744 let message = match (ja, waiting) {
1745 (true, false) => format!(
1746 "タスク {} は、待っていた {} が削除されたため待機を解除し、元の状態に戻しました",
1747 dependent.short(),
1748 short(deleted)
1749 ),
1750 (true, true) => format!(
1751 "タスク {} は、待っていた {} が削除されたため、残りの依存を待っています",
1752 dependent.short(),
1753 short(deleted)
1754 ),
1755 (false, false) => format!(
1756 "Task {} stopped waiting on {} because it was deleted, and returned to its previous state",
1757 dependent.short(),
1758 short(deleted)
1759 ),
1760 (false, true) => format!(
1761 "Task {} stopped waiting on {} because it was deleted, and is still waiting on its other dependencies",
1762 dependent.short(),
1763 short(deleted)
1764 ),
1765 };
1766 crate::notices::raise_in(
1767 home,
1768 crate::notices::Notice::info(
1769 &format!("unblocked:{}:{}", dependent.id, deleted),
1770 message,
1771 ),
1772 );
1773 }
1774
1775 fn release_dependents_of(&self, dependency: &str) -> (Vec<String>, Vec<String>) {
1780 let (mut released, mut still_blocked) = (Vec::new(), Vec::new());
1781 for listed in self.list() {
1782 if listed.status != TaskStatus::Blocked
1783 || !listed.blocked_by.iter().any(|b| b == dependency)
1784 {
1785 continue;
1786 }
1787 let Ok(_claim) = self.claim(&listed.id) else {
1788 continue;
1789 };
1790 let Ok(mut task) = self.get(&listed.id) else {
1791 continue;
1792 };
1793 if !task.dependency_deleted(dependency) {
1794 continue;
1795 }
1796 if self.put(&mut task).is_err() {
1797 continue;
1798 }
1799 self.note_dependency_deleted(&task, dependency);
1800 if task.status == TaskStatus::Blocked {
1801 still_blocked.push(task.id.clone());
1802 } else {
1803 released.push(task.id.clone());
1804 }
1805 }
1806 (released, still_blocked)
1807 }
1808
1809 fn remove_record_with_attachments(
1813 &self,
1814 resolved: &str,
1815 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1816 ) -> Result<()> {
1817 self.sweep_removed_attachments();
1824 let attachments = self.attachments_dir(resolved);
1825 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1826 let moved = match std::fs::rename(&attachments, &aside) {
1827 Ok(()) => true,
1828 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1829 Err(e) => {
1830 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1831 }
1832 };
1833 let path = self.path_of(resolved);
1834 if let Err(e) = remove_record(&path) {
1835 if moved {
1836 let _ = std::fs::rename(&aside, &attachments);
1837 }
1838 return Err(e).with_context(|| format!("remove {}", path.display()));
1839 }
1840 if moved {
1841 if let Err(e) = std::fs::remove_dir_all(&aside) {
1842 tracing::warn!("leftover attachments {}: {e}", aside.display());
1843 }
1844 }
1845 let lock = self.lock_path(resolved);
1846 if let Err(e) = std::fs::remove_file(&lock) {
1847 if e.kind() != std::io::ErrorKind::NotFound {
1848 return Err(e).with_context(|| format!("remove {}", lock.display()));
1849 }
1850 }
1851 Ok(())
1852 }
1853
1854 fn sweep_removed_attachments(&self) {
1860 let Ok(entries) = std::fs::read_dir(&self.root) else {
1861 return;
1862 };
1863 for entry in entries.flatten() {
1864 let name = entry.file_name();
1865 let name = name.to_string_lossy();
1866 let Some(id) = name.strip_suffix(".attachments.removing") else {
1867 continue;
1868 };
1869 if !self.path_of(id).exists() {
1872 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1873 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1874 }
1875 }
1876 }
1877 self.sweep_tombstones();
1878 }
1879
1880 fn sweep_tombstones(&self) {
1889 let Ok(entries) = std::fs::read_dir(&self.root) else {
1890 return;
1891 };
1892 let tasks = self.list();
1893 for entry in entries.flatten() {
1894 let name = entry.file_name();
1895 let name = name.to_string_lossy();
1896 let Some(id) = name.strip_suffix(".removed") else {
1897 continue;
1898 };
1899 if self.path_of(id).exists()
1900 || tasks.iter().any(|t| t.blocked_by.iter().any(|b| b == id))
1901 {
1902 continue;
1903 }
1904 let old = entry
1905 .metadata()
1906 .and_then(|m| m.modified())
1907 .ok()
1908 .and_then(|m| m.elapsed().ok())
1909 .is_some_and(|age| age >= TOMBSTONE_GRACE);
1910 if old {
1911 let _ = std::fs::remove_file(entry.path());
1912 }
1913 }
1914 }
1915
1916 fn lock_path(&self, id: &str) -> PathBuf {
1919 self.root.join(format!("{id}.lock"))
1920 }
1921
1922 pub fn list(&self) -> Vec<Task> {
1935 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1936 .into_iter()
1937 .flatten()
1938 .flatten()
1939 .map(|e| e.path())
1940 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1941 .filter_map(|p| read_path(&p).ok())
1942 .collect();
1943 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1944 tasks
1945 }
1946
1947 pub fn superseded(&self) -> HashMap<String, String> {
1960 let mut by = HashMap::new();
1961 for task in self.list() {
1962 for earlier in &task.runs {
1963 if let Some(later) = task.successor_of(earlier) {
1964 by.insert(earlier.clone(), later.clone());
1965 }
1966 }
1967 }
1968 by
1969 }
1970
1971 pub fn superseded_by(&self, run: &str) -> Option<String> {
1979 for task in self.list() {
1980 if task.runs.iter().any(|r| r == run) {
1981 return task.successor_of(run).cloned();
1982 }
1983 }
1984 None
1985 }
1986
1987 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1998 for task in self.list() {
1999 if task.runs.iter().any(|r| r == run) {
2000 return task.runs.last().filter(|last| **last != run).cloned();
2001 }
2002 }
2003 None
2004 }
2005
2006 pub fn next_runnable(&self) -> Option<Task> {
2011 let mut runnable: Vec<Task> = self
2012 .list()
2013 .into_iter()
2014 .filter(|t| t.status.runnable())
2015 .collect();
2016 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
2017 runnable.into_iter().next()
2018 }
2019
2020 pub fn claim(&self, id: &str) -> Result<Claim> {
2027 std::fs::create_dir_all(&self.root)
2028 .with_context(|| format!("create {}", self.root.display()))?;
2029 let path = self.lock_path(id);
2030 match std::fs::OpenOptions::new()
2031 .write(true)
2032 .create_new(true)
2033 .open(&path)
2034 {
2035 Ok(mut f) => {
2036 use std::io::Write as _;
2037 let _ = writeln!(f, "{}", std::process::id());
2039 Ok(Claim { path })
2040 }
2041 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
2042 bail!("task {id} is already claimed ({} exists)", path.display())
2043 }
2044 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
2045 }
2046 }
2047
2048 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
2050 if self.path_of(prefix).is_file() {
2051 return Ok(prefix.to_owned());
2052 }
2053 let hits: Vec<String> = self
2054 .list()
2055 .into_iter()
2056 .map(|t| t.id)
2057 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
2058 .collect();
2059 match hits.len() {
2060 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
2061 0 => bail!("no task matches `{prefix}`"),
2062 _ => bail!(
2063 "`{prefix}` matches {} tasks: {}",
2064 hits.len(),
2065 hits.join(", ")
2066 ),
2067 }
2068 }
2069
2070 pub fn revision(&self) -> u64 {
2077 self.revision_excluding(&std::collections::BTreeSet::new())
2078 }
2079
2080 pub fn revision_excluding(&self, skip: &std::collections::BTreeSet<String>) -> u64 {
2083 use std::hash::{Hash as _, Hasher as _};
2084
2085 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
2086 .into_iter()
2087 .flatten()
2088 .flatten()
2089 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
2090 .filter(|e| {
2091 let path = e.path();
2092 !path
2093 .file_stem()
2094 .is_some_and(|stem| skip.contains(stem.to_string_lossy().as_ref()))
2095 })
2096 .filter_map(|e| {
2097 let name = e.file_name().to_string_lossy().into_owned();
2098 let mtime = e
2099 .metadata()
2100 .ok()?
2101 .modified()
2102 .ok()?
2103 .duration_since(std::time::UNIX_EPOCH)
2104 .ok()?
2105 .as_millis() as u64;
2106 Some((name, mtime))
2107 })
2108 .collect();
2109
2110 if entries.is_empty() {
2111 return 0;
2112 }
2113
2114 entries.sort_unstable();
2115 let mut hasher = std::hash::DefaultHasher::new();
2116 for (name, mtime) in &entries {
2117 name.hash(&mut hasher);
2118 mtime.hash(&mut hasher);
2119 }
2120 let h = hasher.finish();
2121 if h == 0 { 1 } else { h }
2122 }
2123}
2124
2125const TOMBSTONE_GRACE: std::time::Duration = std::time::Duration::from_secs(3600);
2128
2129#[derive(Debug, Clone)]
2131pub struct Removal {
2132 pub id: String,
2134 pub released: Vec<String>,
2138 pub still_blocked: Vec<String>,
2141}
2142
2143#[derive(Debug)]
2145pub struct Claim {
2146 path: PathBuf,
2147}
2148
2149impl Drop for Claim {
2150 fn drop(&mut self) {
2151 let _ = std::fs::remove_file(&self.path);
2152 }
2153}
2154
2155pub fn title_from(instruction: &str, max: usize) -> String {
2158 let Some(line) = first_line(instruction) else {
2159 return "(empty task)".to_owned();
2160 };
2161 if line.chars().count() <= max {
2162 return line.to_owned();
2163 }
2164 let head: String = line.chars().take(max.saturating_sub(1)).collect();
2165 format!("{head}…")
2166}
2167
2168pub(crate) fn first_line(instruction: &str) -> Option<&str> {
2174 let line = instruction
2175 .lines()
2176 .map(str::trim)
2177 .find(|l| !l.is_empty())?
2178 .trim_start_matches(['#', '-', '*', '>', ' '])
2179 .trim();
2180 (!line.is_empty()).then_some(line)
2181}
2182
2183fn read_path(path: &Path) -> Result<Task> {
2184 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2185 let task: Task =
2186 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
2187 if task.schema > SCHEMA {
2193 bail!(
2194 "task {} was written by a different magi (schema {}, this build \
2195 speaks {SCHEMA})",
2196 task.id,
2197 task.schema
2198 );
2199 }
2200 Ok(task)
2201}
2202
2203pub fn missing_blockers(
2217 queue: &Queue,
2218 questions: &Questions,
2219 blocked_by: &[String],
2220) -> Vec<String> {
2221 blocked_by
2222 .iter()
2223 .filter(|id| {
2224 !queue.path_of(id).is_file()
2225 && !questions.path_of(id).is_file()
2226 && !queue.was_deleted(id)
2227 })
2228 .cloned()
2229 .collect()
2230}
2231
2232pub fn deleted_blockers(queue: &Queue, blocked_by: &[String]) -> Vec<String> {
2237 blocked_by
2238 .iter()
2239 .filter(|id| queue.was_deleted(id))
2240 .cloned()
2241 .collect()
2242}
2243
2244pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
2256 missing_blocker_hold_reason_in(blocked_by, missing, "en")
2257}
2258
2259pub fn missing_blocker_hold_reason_in(
2261 blocked_by: &[String],
2262 missing: &[String],
2263 language: &str,
2264) -> String {
2265 if crate::lang::is_japanese(language) {
2266 format!(
2267 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
2268 blocked_by.join(", "),
2269 missing.join(", "),
2270 )
2271 } else {
2272 format!(
2273 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
2274 blocked_by.join(", "),
2275 missing.join(", "),
2276 )
2277 }
2278}
2279
2280pub fn short(id: &str) -> &str {
2282 id.split('-').next_back().unwrap_or(id)
2283}
2284
2285fn new_id() -> String {
2286 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2287 let seed = crate::rng::entropy();
2288 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2289}
2290
2291#[cfg(test)]
2292mod tests {
2293 #[test]
2294 fn merge_over_lets_the_later_choice_win_and_keeps_the_rest() {
2295 let mut base = RunOverrides {
2296 merge: Some("pr".to_owned()),
2297 candidates: Some(3),
2298 ..RunOverrides::default()
2299 };
2300 base.merge_over(&RunOverrides {
2301 merge: Some("none".to_owned()),
2302 seed: Some(7),
2303 ..RunOverrides::default()
2304 });
2305 assert_eq!(base.merge.as_deref(), Some("none"));
2306 assert_eq!(base.candidates, Some(3));
2307 assert_eq!(base.seed, Some(7));
2308 }
2309
2310 #[test]
2311 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
2312 let mut t = task("relanded");
2313 t.attempts = 1;
2314 t.status = TaskStatus::Running;
2315 t.already_landed("already in main as 0e368de");
2316 assert_eq!(t.status, TaskStatus::Done);
2317 assert_eq!(t.attempts, 0);
2318 assert!(t.hold_reason.is_none());
2319 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
2320 }
2321
2322 #[test]
2323 fn missing_blocker_reason_follows_the_language() {
2324 let b = vec!["a".to_owned()];
2325 let en = missing_blocker_hold_reason_in(&b, &b, "en");
2326 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
2327 assert!(en.starts_with("blocked on a"));
2328 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
2329 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
2330 }
2331
2332 use super::*;
2333
2334 #[test]
2335 fn triage_applied_survives_release_and_old_records_read_as_empty() {
2336 let mut t = Task::new(
2337 "t".to_owned(),
2338 "i".to_owned(),
2339 PathBuf::from("r"),
2340 Source::Human,
2341 );
2342 t.mark_triage_applied("q1");
2343 t.mark_triage_applied("q1");
2344 t.hold_machine(Some("x".to_owned()));
2345 t.release();
2346 assert_eq!(t.triage_applied, ["q1"]);
2347 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
2348
2349 let mut v = serde_json::to_value(&t).unwrap();
2350 v.as_object_mut().unwrap().remove("triage_applied");
2351 let old: Task = serde_json::from_value(v).unwrap();
2352 assert!(old.triage_applied.is_empty());
2353 }
2354
2355 #[test]
2356 fn task_counts_of_empty_is_all_zero() {
2357 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
2358 }
2359
2360 #[test]
2361 fn task_counts_of_tallies_every_status() {
2362 let mut queued = Task::new(
2363 "q".to_owned(),
2364 "i".to_owned(),
2365 PathBuf::from("."),
2366 Source::Human,
2367 );
2368 queued.status = TaskStatus::Queued;
2369 let mut running = queued.clone();
2370 running.status = TaskStatus::Running;
2371 let mut done = queued.clone();
2372 done.status = TaskStatus::Done;
2373 let mut failed = queued.clone();
2374 failed.status = TaskStatus::Failed;
2375 let mut held = queued.clone();
2376 held.status = TaskStatus::Held;
2377 let mut blocked = queued.clone();
2378 blocked.status = TaskStatus::Blocked;
2379
2380 let mut parked = queued.clone();
2381 parked.status = TaskStatus::Parked;
2382
2383 let counts = TaskCounts::of(&[
2384 queued,
2385 running,
2386 done.clone(),
2387 done,
2388 failed,
2389 held,
2390 blocked,
2391 parked,
2392 ]);
2393 assert_eq!(
2394 counts,
2395 TaskCounts {
2396 queued: 1,
2397 running: 1,
2398 done: 2,
2399 failed: 1,
2400 held: 1,
2401 blocked: 1,
2402 parked: 1,
2403 }
2404 );
2405 }
2406
2407 fn queue() -> (tempfile::TempDir, Queue) {
2410 let dir = tempfile::tempdir().unwrap();
2411 let q = Queue::at(dir.path().join("queue"));
2412 (dir, q)
2413 }
2414
2415 #[test]
2416 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
2417 let dir = tempfile::tempdir().unwrap();
2418 let q = Queue::at(dir.path().join("queue"));
2419 let mut t = task("held");
2420 q.put(&mut t).unwrap();
2421 assert_eq!(
2422 crate::notices::Notices::at(dir.path().join("notifications"))
2423 .list()
2424 .len(),
2425 0
2426 );
2427 t.hold_machine(Some("out of attempts".to_owned()));
2428 q.put(&mut t).unwrap();
2429 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
2430 assert_eq!(listed.len(), 1);
2431 assert!(listed[0].message.contains("out of attempts"));
2432 }
2433
2434 fn task(title: &str) -> Task {
2435 Task::new(
2436 title.to_owned(),
2437 format!("do {title}"),
2438 PathBuf::from("."),
2439 Source::Human,
2440 )
2441 }
2442
2443 #[test]
2444 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
2445 let mut t = task("retried");
2446 assert!(
2447 t.earlier_attempts().is_empty(),
2448 "a first attempt takes nothing over"
2449 );
2450 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2451 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
2452 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
2454 assert_eq!(t.successor_of("bbbb"), None);
2455 assert_eq!(t.successor_of("zzzz"), None);
2456 }
2457
2458 #[test]
2459 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
2460 let (_dir, q) = queue();
2461 let mut t = task("retried");
2462 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2463 q.put(&mut t).unwrap();
2464
2465 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
2466 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
2467 assert_eq!(
2468 q.superseded_by("cccc"),
2469 None,
2470 "the latest attempt replaces nothing"
2471 );
2472 assert_eq!(
2473 q.superseded_by("never-heard-of-it"),
2474 None,
2475 "a run belonging to no task on this queue is not superseded"
2476 );
2477
2478 let mut by = HashMap::new();
2479 by.insert("aaaa".to_owned(), "bbbb".to_owned());
2480 by.insert("bbbb".to_owned(), "cccc".to_owned());
2481 assert_eq!(
2482 q.superseded(),
2483 by,
2484 "the whole-map and single-run forms must agree"
2485 );
2486 }
2487
2488 #[test]
2489 fn superseded_attempts_is_empty_until_the_task_is_done() {
2490 let mut t = task("retried");
2491 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2492 t.status = TaskStatus::Failed;
2493 assert_eq!(
2494 t.superseded_attempts(true),
2495 &[] as &[String],
2496 "a task still retrying has no attempt yet that a later one made moot"
2497 );
2498
2499 t.status = TaskStatus::Running;
2500 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2501 }
2502
2503 #[test]
2504 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
2505 let mut t = task("retried");
2506 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2507 t.status = TaskStatus::Done;
2508 assert_eq!(
2509 t.superseded_attempts(true),
2510 &["aaaa".to_owned(), "bbbb".to_owned()],
2511 "cccc is the attempt whose success made the task done, and stays out"
2512 );
2513 }
2514
2515 #[test]
2516 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
2517 let mut t = task("first try landed");
2518 t.runs = vec!["aaaa".to_owned()];
2519 t.status = TaskStatus::Done;
2520 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2521 }
2522
2523 #[test]
2524 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
2525 let mut t = task("closed by hand after a manual merge");
2532 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2533 t.status = TaskStatus::Done;
2534 assert_eq!(
2535 t.superseded_attempts(false),
2536 &[] as &[String],
2537 "nothing here is provably why the task is done, so nothing is superseded"
2538 );
2539 }
2540
2541 #[test]
2542 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2543 let (_dir, q) = queue();
2544 let mut t = task("retried twice");
2545 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2546 q.put(&mut t).unwrap();
2547
2548 assert_eq!(
2549 q.latest_attempt("aaaa"),
2550 Some("cccc".to_owned()),
2551 "an old attempt points straight at the chain's current head, not the \
2552 next attempt in the middle of it"
2553 );
2554 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2555 assert_eq!(
2556 q.latest_attempt("cccc"),
2557 None,
2558 "the latest attempt is not superseded by anything"
2559 );
2560 assert_eq!(
2561 q.latest_attempt("never-heard-of-it"),
2562 None,
2563 "a run belonging to no task on this queue is not superseded"
2564 );
2565 }
2566
2567 #[test]
2568 fn a_markdown_heading_is_the_title_not_decoration() {
2569 assert_eq!(
2574 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2575 "Rework the config loader"
2576 );
2577 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2578 assert_eq!(title_from("> quoted task", 40), "quoted task");
2579 assert_eq!(title_from(" \n\n", 40), "(empty task)");
2581 assert_eq!(title_from("###\n", 40), "(empty task)");
2582 }
2583
2584 #[test]
2585 fn a_long_title_is_elided_by_characters_not_bytes() {
2586 let long = "課題".repeat(30);
2588 let title = title_from(&long, 10);
2589 assert_eq!(title.chars().count(), 10);
2590 assert!(title.ends_with('…'));
2591 }
2592
2593 #[test]
2594 fn priority_wins_and_ties_break_oldest_first() {
2595 let (_dir, q) = queue();
2596 let mut a = task("first");
2597 let mut b = task("second");
2598 let mut c = task("urgent");
2599 a.id = "20260101-000001-aaaa".to_owned();
2601 b.id = "20260101-000002-bbbb".to_owned();
2602 c.id = "20260101-000003-cccc".to_owned();
2603 c.priority = 5;
2604 for t in [&mut a, &mut b, &mut c] {
2605 q.put(t).unwrap();
2606 }
2607
2608 assert_eq!(q.next_runnable().unwrap().id, c.id);
2610 c.hold_machine(None);
2611 q.put(&mut c).unwrap();
2612 assert_eq!(q.next_runnable().unwrap().id, a.id);
2614 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2615 }
2616
2617 #[test]
2618 fn a_blocked_task_never_starves_another_runnable_one() {
2619 let (_dir, q) = queue();
2620 let mut blocked = task("blocked");
2621 blocked.block(vec!["something".to_owned()], None);
2622 q.put(&mut blocked).unwrap();
2623
2624 let mut runnable = task("free to go");
2625 q.put(&mut runnable).unwrap();
2626
2627 let next = q.next_runnable().expect("a runnable task is still offered");
2628 assert_eq!(next.id, runnable.id);
2629 }
2630
2631 #[test]
2632 fn a_held_task_is_never_offered_to_the_loop() {
2633 let (_dir, q) = queue();
2634 let mut t = task("held");
2635 q.put(&mut t).unwrap();
2636 assert!(q.next_runnable().is_some());
2637
2638 t.hold_machine(None);
2639 q.put(&mut t).unwrap();
2640 assert!(
2641 q.next_runnable().is_none(),
2642 "a held task must wait for a human"
2643 );
2644
2645 t.status = TaskStatus::Failed;
2647 q.put(&mut t).unwrap();
2648 assert!(q.next_runnable().is_some());
2649 }
2650
2651 #[test]
2652 fn attempts_are_capped_and_then_the_task_is_held() {
2653 let mut t = task("doomed");
2654
2655 t.start("run-1".to_owned());
2656 t.fail("gate red", 2);
2657 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2658
2659 t.start("run-2".to_owned());
2660 t.fail("gate red", 2);
2661 assert_eq!(
2662 t.status,
2663 TaskStatus::Held,
2664 "out of attempts: stop spending money on it"
2665 );
2666 assert_eq!(t.runs, ["run-1", "run-2"]);
2667 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2668 assert_eq!(
2669 t.hold_reason.as_deref(),
2670 Some("gate red"),
2671 "the hold must say why, not leave hold_reason null next to a \
2672 populated last_error"
2673 );
2674 }
2675
2676 #[test]
2677 fn handing_off_a_task_records_a_hold_reason_too() {
2678 let mut t = task("left a pull request");
2679 t.start("run-1".to_owned());
2680 t.handed_off("run ended with a pull request open [run run-1]");
2681 assert_eq!(t.status, TaskStatus::Held);
2682 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2683 assert_eq!(
2684 t.hold_reason.as_deref(),
2685 Some("run ended with a pull request open [run run-1]")
2686 );
2687 assert_eq!(t.hold_reason, t.last_error);
2688 }
2689
2690 #[test]
2691 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2692 let mut t = task("stalled by quota");
2693
2694 t.start("run-1".to_owned());
2695 assert_eq!(t.attempts, 1);
2696 t.stall("judge-1, judge-2 out of quota");
2697 assert_eq!(
2698 t.attempts, 0,
2699 "a closed quota window must not spend the task's retry budget"
2700 );
2701 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2702 assert_eq!(
2703 t.last_error.as_deref(),
2704 Some("judge-1, judge-2 out of quota")
2705 );
2706
2707 for _ in 0..20 {
2710 t.start("run-n".to_owned());
2711 t.stall("still out of quota");
2712 }
2713 t.start("run-real".to_owned());
2714 t.fail("gate red", 2);
2715 assert_eq!(
2716 t.status,
2717 TaskStatus::Failed,
2718 "the first attempt that was really judged is attempt one"
2719 );
2720 }
2721
2722 #[test]
2723 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2724 let mut t = task("retry me");
2725 t.start("run-1".to_owned());
2726 t.fail("gate red", 1);
2727 assert_eq!(t.status, TaskStatus::Held);
2728
2729 t.release();
2730 assert_eq!(t.status, TaskStatus::Queued);
2731 assert_eq!(t.attempts, 0);
2734 assert!(t.last_error.is_none());
2735 assert_eq!(
2736 t.runs.len(),
2737 1,
2738 "history is kept: attempts reset, evidence does not"
2739 );
2740 }
2741
2742 #[test]
2743 fn a_hold_reason_survives_and_a_release_clears_it() {
2744 let mut t = task("waiting on something else");
2745 t.hold_manual(Some(
2746 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2747 ));
2748 assert_eq!(t.status, TaskStatus::Held);
2749 assert_eq!(
2750 t.hold_reason.as_deref(),
2751 Some("waiting for 20260101-000000-aaaa to land first")
2752 );
2753
2754 t.hold_manual(None);
2756 assert_eq!(
2757 t.hold_reason.as_deref(),
2758 Some("waiting for 20260101-000000-aaaa to land first"),
2759 "a bare re-hold keeps whatever a human already wrote down"
2760 );
2761
2762 let mut plain = task("no reason given");
2764 plain.hold_manual(None);
2765 assert_eq!(plain.status, TaskStatus::Held);
2766 assert!(plain.hold_reason.is_none());
2767
2768 t.release();
2769 assert_eq!(t.status, TaskStatus::Queued);
2770 assert!(
2771 t.hold_reason.is_none(),
2772 "a stale reason must not greet the next person who holds this task"
2773 );
2774 }
2775
2776 #[test]
2777 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2778 let mut t = task("landed by hand while held");
2783 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2784 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2785
2786 t.succeed();
2787 assert_eq!(t.status, TaskStatus::Done);
2788 assert!(
2789 t.hold_reason.is_none(),
2790 "a done task cannot still be waiting on something"
2791 );
2792 }
2793
2794 #[test]
2795 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2796 let mut held = task("held straight out of blocked");
2803 held.block(
2804 vec!["20260101-000000-dead".to_owned()],
2805 Some("waiting on the migration script".to_owned()),
2806 );
2807 assert_eq!(held.status, TaskStatus::Blocked);
2808
2809 held.hold_manual(None);
2810 assert_eq!(held.status, TaskStatus::Held);
2811 assert!(
2812 held.blocked_by.is_empty(),
2813 "hold overrides the wait, same as release"
2814 );
2815 assert!(held.block_reason.is_none());
2816
2817 let mut done = task("closed straight out of blocked");
2818 done.block(
2819 vec!["20260101-000000-dead".to_owned()],
2820 Some("waiting on the migration script".to_owned()),
2821 );
2822 done.succeed();
2823 assert_eq!(done.status, TaskStatus::Done);
2824 assert!(
2825 done.blocked_by.is_empty(),
2826 "a done task cannot still be waiting on a dependency"
2827 );
2828 assert!(done.block_reason.is_none());
2829 }
2830
2831 #[test]
2832 fn a_blocked_task_is_never_offered_to_the_loop() {
2833 let mut t = task("blocked");
2834 assert!(t.status.runnable());
2835 t.block(
2836 vec!["dep-id".to_owned()],
2837 Some("waits on dep-id".to_owned()),
2838 );
2839 assert_eq!(t.status, TaskStatus::Blocked);
2840 assert!(!t.status.runnable());
2841 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2842 }
2843
2844 #[test]
2845 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2846 let mut t = task("blocked on two");
2847 t.block(
2848 vec!["a".to_owned(), "b".to_owned()],
2849 Some("waits on a and b".to_owned()),
2850 );
2851
2852 t.unblock("a");
2853 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2854 assert_eq!(t.blocked_by, ["b"]);
2855
2856 t.unblock("b");
2857 assert_eq!(t.status, TaskStatus::Queued);
2858 assert!(t.blocked_by.is_empty());
2859 assert!(t.block_reason.is_none());
2860 }
2861
2862 #[test]
2863 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2864 let mut t = task("never blocked");
2865 t.unblock("whatever");
2866 assert_eq!(t.status, TaskStatus::Queued);
2867 }
2868
2869 #[test]
2870 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2871 let mut t = task("held, then asked about");
2876 t.hold_machine(Some("out of attempts".to_owned()));
2877 assert_eq!(t.status, TaskStatus::Held);
2878
2879 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2880 assert_eq!(t.status, TaskStatus::Blocked);
2881
2882 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2883 t.unblock("q1");
2884 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2885 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2886 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2887 assert!(t.blocked_from.is_none(), "consumed once restored");
2888 }
2889
2890 #[test]
2891 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2892 let mut t = task("manually held, then asked about");
2893 t.hold_manual(Some("waiting on a dependency".to_owned()));
2894
2895 t.block(vec!["q1".to_owned()], None);
2896 t.unblock("q1");
2897
2898 assert_eq!(t.status, TaskStatus::Held);
2899 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2900 }
2901
2902 #[test]
2903 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2904 let mut t = task("held, blocked twice");
2908 t.hold_machine(None);
2909 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2910 t.block(
2911 vec!["q1".to_owned(), "q2".to_owned()],
2912 Some("second".to_owned()),
2913 );
2914
2915 t.unblock("q1");
2916 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2917 t.unblock("q2");
2918 assert_eq!(t.status, TaskStatus::Held);
2919 }
2920
2921 #[test]
2922 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2923 let mut t = task("blocked mid-run");
2927 t.start("run-1".to_owned());
2928 assert_eq!(t.status, TaskStatus::Running);
2929
2930 t.block(vec!["q1".to_owned()], None);
2931 t.unblock("q1");
2932 assert_eq!(t.status, TaskStatus::Queued);
2933 }
2934
2935 #[test]
2936 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2937 let mut t = task("legacy record, held before it was blocked");
2943 t.hold_source = Some(HoldSource::Machine);
2944 t.hold_reason = Some("legacy hold reason".to_owned());
2945 t.status = TaskStatus::Blocked;
2946 t.blocked_by = vec!["q1".to_owned()];
2947 t.blocked_from = None;
2948
2949 t.unblock("q1");
2950 assert_eq!(t.status, TaskStatus::Held);
2951 }
2952
2953 #[test]
2954 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2955 let mut t = task("legacy record, ordinary dependency block");
2956 t.status = TaskStatus::Blocked;
2957 t.blocked_by = vec!["dep".to_owned()];
2958 t.blocked_from = None;
2959
2960 t.unblock("dep");
2961 assert_eq!(t.status, TaskStatus::Queued);
2962 }
2963
2964 #[test]
2965 fn answering_a_question_is_recorded_and_survives_a_release() {
2966 let mut t = task("asked something");
2967 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2968 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2969 t.unblock("q1");
2970 assert_eq!(t.status, TaskStatus::Queued);
2971 assert_eq!(t.answers.len(), 1);
2972 assert_eq!(t.answers[0].answer, "SQLite");
2973
2974 t.release();
2978 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2979 }
2980
2981 #[test]
2982 fn a_refused_handover_keeps_the_review_branch_across_release() {
2983 let mut t = task("refused takeover");
2984 t.start("run-1".to_owned());
2985 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
2986 assert_eq!(t.status, TaskStatus::Held);
2987 assert_eq!(t.attempts, 1);
2988 t.release();
2989 assert_eq!(t.status, TaskStatus::Queued);
2990 assert_eq!(t.attempts, 0);
2991 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2992
2993 t.hold_for_handover(None, "again".to_owned());
2995 t.requeue();
2996 assert!(t.review_branch.is_none());
2997
2998 let mut m = task("manual");
3000 m.review_branch = Some("magi/x/A".to_owned());
3001 m.hold_manual(None);
3002 m.release();
3003 assert!(m.review_branch.is_none());
3004 }
3005
3006 #[test]
3007 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
3008 let mut t = task("blocked run with a surviving branch");
3009 t.start("run-1".to_owned());
3010 t.fail("blocked with major findings", 5);
3011 assert_eq!(t.status, TaskStatus::Failed);
3012
3013 t.request_review("magi/eba2/A".to_owned());
3014 assert_eq!(t.status, TaskStatus::Queued);
3015 assert_eq!(t.attempts, 0);
3016 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
3017
3018 t.release();
3020 assert!(t.review_branch.is_none());
3021 }
3022
3023 #[test]
3024 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
3025 let mut t = task("retry");
3026 t.start("run-1".to_owned());
3027 t.requeue();
3028 assert!(t.fresh_start);
3029
3030 t.release();
3031 assert!(!t.fresh_start);
3032 }
3033
3034 #[test]
3035 fn priority_can_be_changed_while_queued_but_not_while_running() {
3036 let mut t = task("reprioritise me");
3037 t.set_priority(5).unwrap();
3038 assert_eq!(t.priority, 5);
3039
3040 t.start("run-1".to_owned());
3041 let err = t.set_priority(9).unwrap_err().to_string();
3042 assert!(err.contains("running"), "{err}");
3043 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
3044 }
3045
3046 #[test]
3047 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
3048 let mut t = task("interrupt me");
3049 assert!(!t.interrupt, "off unless asked, same as any other task");
3050
3051 t.set_interrupt(true).unwrap();
3052 assert!(t.interrupt);
3053
3054 t.start("run-1".to_owned());
3055 assert!(
3056 !t.interrupt,
3057 "the mark is one-shot: dispatching the task fulfils it, \
3058 whatever the run that follows ends up doing"
3059 );
3060 let err = t.set_interrupt(true).unwrap_err().to_string();
3061 assert!(err.contains("running"), "{err}");
3062 t.set_interrupt(false).unwrap();
3065 assert!(!t.interrupt);
3066 }
3067
3068 #[test]
3072 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
3073 let mut t = task("interrupt me");
3074 t.set_interrupt(true).unwrap();
3075 t.start("run-1".to_owned());
3076 t.fail("mock failure", 5);
3077 assert_eq!(t.status, TaskStatus::Failed);
3078 assert!(
3079 !t.interrupt,
3080 "one attempt already spent the mark; a retry is an ordinary \
3081 requeue, not a fresh interrupt request"
3082 );
3083 }
3084
3085 #[test]
3086 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
3087 let (_dir, q) = queue();
3088 let mut a = task("first filed");
3089 let mut b = task("second filed");
3090 a.id = "20260101-000001-aaaa".to_owned();
3091 b.id = "20260101-000002-bbbb".to_owned();
3092 q.put(&mut a).unwrap();
3093 q.put(&mut b).unwrap();
3094
3095 assert_eq!(
3096 q.next_runnable().unwrap().id,
3097 a.id,
3098 "with equal priority the older task goes first, so a burst of \
3099 new work cannot starve it"
3100 );
3101 assert_eq!(
3102 q.list()[0].id,
3103 b.id,
3104 "but the list an operator reads is newest first, the same as \
3105 before priority existed - a's turn to run does not make it the \
3106 newest task"
3107 );
3108
3109 let mut a = q.get(&a.id).unwrap();
3110 a.set_priority(10).unwrap();
3111 q.put(&mut a).unwrap();
3112
3113 assert_eq!(
3114 q.next_runnable().unwrap().id,
3115 a.id,
3116 "a raised priority must be reflected the moment it is saved"
3117 );
3118 assert_eq!(
3122 q.list()[0].id,
3123 a.id,
3124 "the raised task must sort first in the list an operator reads, \
3125 not only in next_runnable's own ordering"
3126 );
3127 }
3128
3129 #[test]
3130 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
3131 let mut t = Task::new(
3132 "old title".to_owned(),
3133 "old instruction".to_owned(),
3134 PathBuf::from("/repo"),
3135 Source::Agent {
3136 run: "20260101-000000-beef".to_owned(),
3137 node: "implement".to_owned(),
3138 },
3139 );
3140 let id = t.id.clone();
3141 let created_at = t.created_at;
3142 t.runs.push("20260101-000000-beef".to_owned());
3143
3144 t.edit("new title".to_owned(), "new instruction".to_owned())
3145 .unwrap();
3146
3147 assert_eq!(t.title, "new title");
3148 assert_eq!(t.instruction, "new instruction");
3149 assert_eq!(t.id, id, "editing must not mint a new id");
3150 assert_eq!(t.created_at, created_at);
3151 assert_eq!(
3152 t.source,
3153 Source::Agent {
3154 run: "20260101-000000-beef".to_owned(),
3155 node: "implement".to_owned(),
3156 },
3157 "editing must not turn agent attribution into human"
3158 );
3159 assert_eq!(t.runs, ["20260101-000000-beef"]);
3160 }
3161
3162 #[test]
3163 fn editing_is_refused_once_a_task_is_running_or_finished() {
3164 let mut running = task("in flight");
3165 running.start("run-1".to_owned());
3166 let err = running
3167 .edit("x".to_owned(), "y".to_owned())
3168 .unwrap_err()
3169 .to_string();
3170 assert!(err.contains("running"), "{err}");
3171
3172 let mut done = task("finished");
3173 done.succeed();
3174 let err = done
3175 .edit("x".to_owned(), "y".to_owned())
3176 .unwrap_err()
3177 .to_string();
3178 assert!(err.contains("done"), "{err}");
3179
3180 let mut queued = task("waiting");
3182 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
3183 let mut held = task("parked");
3184 held.hold_machine(None);
3185 held.edit("x".to_owned(), "y".to_owned()).unwrap();
3186 }
3187
3188 #[test]
3189 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
3190 let (_dir, q) = queue();
3191 let path = q.path_of("20260101-000000-aaaa");
3192 std::fs::create_dir_all(q.root()).unwrap();
3193 std::fs::write(
3194 &path,
3195 serde_json::json!({
3196 "schema": SCHEMA,
3197 "id": "20260101-000000-aaaa",
3198 "title": "from before hold reasons existed",
3199 "instruction": "from before hold reasons existed",
3200 "repo": ".",
3201 "source": { "kind": "human" },
3202 "status": "held",
3203 "created_at": Timestamp::now().to_string(),
3204 "updated_at": Timestamp::now().to_string(),
3205 })
3206 .to_string(),
3207 )
3208 .unwrap();
3209
3210 let task = q.get("20260101-000000-aaaa").expect("must still read");
3211 assert!(task.hold_reason.is_none());
3212 assert!(task.operator_held());
3213 }
3214
3215 #[test]
3216 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
3217 let (_dir, q) = queue();
3218 let path = q.path_of("20260101-000000-bbbb");
3219 std::fs::create_dir_all(q.root()).unwrap();
3220 std::fs::write(
3221 &path,
3222 serde_json::json!({
3223 "schema": 2,
3224 "id": "20260101-000000-bbbb",
3225 "title": "old manual recovery",
3226 "instruction": "old manual recovery",
3227 "repo": ".",
3228 "source": { "kind": "human" },
3229 "status": "held",
3230 "hold_reason": "active manual recovery run20260912-224242-daf5",
3231 "created_at": Timestamp::now().to_string(),
3232 "updated_at": Timestamp::now().to_string(),
3233 })
3234 .to_string(),
3235 )
3236 .unwrap();
3237
3238 let task = q.get("20260101-000000-bbbb").expect("must still read");
3239 assert_eq!(task.hold_source, None);
3240 assert!(task.operator_held());
3241 }
3242
3243 #[test]
3244 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
3245 let (_dir, q) = queue();
3246 let path = q.path_of("20260101-000000-aaaa");
3247 std::fs::create_dir_all(q.root()).unwrap();
3248 std::fs::write(
3249 &path,
3250 serde_json::json!({
3251 "schema": SCHEMA,
3252 "id": "20260101-000000-aaaa",
3253 "title": "from before diagnostics existed",
3254 "instruction": "from before diagnostics existed",
3255 "repo": ".",
3256 "source": { "kind": "human" },
3257 "status": "held",
3258 "created_at": Timestamp::now().to_string(),
3259 "updated_at": Timestamp::now().to_string(),
3260 })
3261 .to_string(),
3262 )
3263 .unwrap();
3264
3265 let task = q.get("20260101-000000-aaaa").expect("must still read");
3266 assert!(task.diagnostic.is_none());
3267 }
3268
3269 #[test]
3270 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
3271 let (_dir, q) = queue();
3275 let path = q.path_of("20260101-000000-aaaa");
3276 std::fs::create_dir_all(q.root()).unwrap();
3277 std::fs::write(
3278 &path,
3279 serde_json::json!({
3280 "schema": 1,
3281 "id": "20260101-000000-aaaa",
3282 "title": "from before blocking existed",
3283 "instruction": "from before blocking existed",
3284 "repo": ".",
3285 "source": { "kind": "human" },
3286 "status": "queued",
3287 "created_at": Timestamp::now().to_string(),
3288 "updated_at": Timestamp::now().to_string(),
3289 })
3290 .to_string(),
3291 )
3292 .unwrap();
3293
3294 let task = q.get("20260101-000000-aaaa").expect("must still read");
3295 assert!(task.blocked_by.is_empty());
3296 assert!(task.block_reason.is_none());
3297 assert!(task.answers.is_empty());
3298 assert!(task.review_branch.is_none());
3299 }
3300
3301 #[test]
3302 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
3303 let mut held = task("diagnosed");
3308 held.start("run-1".to_owned());
3309 held.fail("gate red", 1);
3310 held.diagnostic = Some("cargo test failed: ...".to_owned());
3311 assert_eq!(held.status, TaskStatus::Held);
3312
3313 held.release();
3314 assert!(held.diagnostic.is_none());
3315
3316 held.diagnostic = Some("cargo test failed: ...".to_owned());
3317 held.succeed();
3318 assert!(held.diagnostic.is_none());
3319 }
3320
3321 #[test]
3322 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
3323 let mut t = task("retried");
3324 t.start("run-1".to_owned());
3325 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
3326 t.fail("unrelated config error", 5);
3327 assert_eq!(t.status, TaskStatus::Failed);
3328 assert!(
3329 t.diagnostic.is_none(),
3330 "fail() must not let an old diagnostic outlive the run that produced it"
3331 );
3332 }
3333
3334 #[test]
3335 fn a_claim_is_exclusive_and_releases_on_drop() {
3336 let (_dir, q) = queue();
3337 let mut t = task("contended");
3338 q.put(&mut t).unwrap();
3339
3340 let held = q.claim(&t.id).unwrap();
3341 assert!(
3342 q.claim(&t.id).is_err(),
3343 "two daemons must not drive one task into two runs"
3344 );
3345 drop(held);
3346 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
3347 }
3348
3349 #[test]
3350 fn a_round_trip_survives_disk() {
3351 let (_dir, q) = queue();
3352 let mut t = Task::new(
3353 "titled".to_owned(),
3354 "body".to_owned(),
3355 PathBuf::from("/repo"),
3356 Source::Agent {
3357 run: "20260101-000000-beef".to_owned(),
3358 node: "implement".to_owned(),
3359 },
3360 );
3361 t.priority = 3;
3362 q.put(&mut t).unwrap();
3363
3364 let back = q.get(&t.id).unwrap();
3365 assert_eq!(back.id, t.id);
3366 assert_eq!(back.priority, 3);
3367 assert_eq!(back.source.label(), "implement@beef");
3368 assert_eq!(q.get(t.short()).unwrap().id, t.id);
3370 }
3371
3372 #[test]
3373 fn an_unreadable_task_does_not_take_the_queue_down() {
3374 let (_dir, q) = queue();
3375 let mut t = task("fine");
3376 q.put(&mut t).unwrap();
3377 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
3378
3379 let listed = q.list();
3380 assert_eq!(listed.len(), 1, "the readable task still lists");
3381 assert_eq!(listed[0].id, t.id);
3382 }
3383
3384 #[test]
3385 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
3386 let (_dir, q) = queue();
3387 let path = q.path_of("20260101-000000-aaaa");
3388 std::fs::create_dir_all(q.root()).unwrap();
3389 std::fs::write(
3390 &path,
3391 serde_json::json!({
3392 "schema": SCHEMA,
3393 "id": "20260101-000000-aaaa",
3394 "title": "from before solo existed",
3395 "instruction": "from before solo existed",
3396 "repo": ".",
3397 "source": { "kind": "human" },
3398 "status": "queued",
3399 "created_at": Timestamp::now().to_string(),
3400 "updated_at": Timestamp::now().to_string(),
3401 })
3402 .to_string(),
3403 )
3404 .unwrap();
3405
3406 let task = q.get("20260101-000000-aaaa").expect("must still read");
3407 assert!(!task.solo, "a queue file with no `solo` field means false");
3408 }
3409
3410 #[test]
3411 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
3412 let (_dir, q) = queue();
3413 let path = q.path_of("20260101-000000-bbbb");
3414 std::fs::create_dir_all(q.root()).unwrap();
3415 std::fs::write(
3416 &path,
3417 serde_json::json!({
3418 "schema": SCHEMA,
3419 "id": "20260101-000000-bbbb",
3420 "title": "from before urgent existed",
3421 "instruction": "from before urgent existed",
3422 "repo": ".",
3423 "source": { "kind": "human" },
3424 "status": "queued",
3425 "created_at": Timestamp::now().to_string(),
3426 "updated_at": Timestamp::now().to_string(),
3427 })
3428 .to_string(),
3429 )
3430 .unwrap();
3431
3432 let task = q.get("20260101-000000-bbbb").expect("must still read");
3433 assert!(
3434 !task.urgent,
3435 "a queue file with no `urgent` field means false, same as `solo`"
3436 );
3437 }
3438
3439 #[test]
3440 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
3441 let (_dir, q) = queue();
3442 let mut t = task("from the future");
3443 q.put(&mut t).unwrap();
3444 let path = q.path_of(&t.id);
3445 let body = std::fs::read_to_string(&path)
3446 .unwrap()
3447 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
3448 std::fs::write(&path, body).unwrap();
3449
3450 let err = q.get(&t.id).unwrap_err().to_string();
3451 assert!(err.contains("schema 99"), "{err}");
3452 }
3453
3454 #[test]
3455 fn revision_moves_when_the_queue_changes() {
3456 let (_dir, q) = queue();
3457 assert_eq!(q.revision(), 0, "an empty queue has no revision");
3458 let mut t = task("first");
3459 q.put(&mut t).unwrap();
3460 assert!(q.revision() > 0, "a written task moves the revision");
3461 }
3462
3463 #[test]
3464 fn revision_moves_when_deleting_an_older_task() {
3465 let (dir, q) = queue();
3466 let questions = Questions::at(dir.path().join("questions"));
3467 let mut t1 = task("older");
3468 q.put(&mut t1).unwrap();
3469 std::thread::sleep(std::time::Duration::from_millis(10));
3471 let mut t2 = task("newer");
3472 q.put(&mut t2).unwrap();
3473
3474 let rev_before = q.revision();
3475 q.remove(&t1.id, false, &questions).unwrap();
3476 let rev_after = q.revision();
3477
3478 assert_ne!(
3479 rev_before, rev_after,
3480 "deleting an older task must change the revision so other clients see the deletion"
3481 );
3482 }
3483
3484 #[test]
3485 fn removing_a_task_takes_it_out_of_the_listing() {
3486 let (dir, q) = queue();
3487 let questions = Questions::at(dir.path().join("questions"));
3488 let mut t = task("delete me");
3489 q.put(&mut t).unwrap();
3490 let removed = q.remove(t.short(), false, &questions).unwrap();
3491 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
3492 assert!(removed.released.is_empty() && removed.still_blocked.is_empty());
3493 assert!(q.list().is_empty());
3494 assert!(
3495 q.remove(&t.id, false, &questions).is_err(),
3496 "removing twice is an error"
3497 );
3498 }
3499
3500 #[test]
3501 fn removing_a_task_takes_its_stale_lock_with_it() {
3502 let (dir, q) = queue();
3503 let questions = Questions::at(dir.path().join("questions"));
3504 let mut t = task("interrupted");
3505 q.put(&mut t).unwrap();
3506
3507 let claim = q.claim(&t.id).unwrap();
3510 std::mem::forget(claim);
3511 assert!(
3512 q.claim(&t.id).is_err(),
3513 "the orphaned lock is what makes the task look claimed"
3514 );
3515
3516 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
3518 assert!(err.contains("live daemon"), "{err}");
3519 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
3520
3521 q.remove(&t.id, false, &questions).unwrap();
3523 assert!(q.list().is_empty());
3524 let mut again = task("interrupted");
3525 again.id = t.id.clone();
3526 q.put(&mut again).unwrap();
3527 assert!(
3528 q.claim(&t.id).is_ok(),
3529 "a task that comes back must be claimable, which a left-behind lock would prevent"
3530 );
3531 }
3532
3533 fn notices_of(dir: &Path) -> Vec<crate::notices::Notice> {
3534 crate::notices::Notices::at(dir.join("notifications")).list()
3535 }
3536
3537 #[test]
3538 fn removing_a_sole_dependency_releases_the_dependent_without_a_hold() {
3539 let (dir, q) = queue();
3540 let questions = Questions::at(dir.path().join("questions"));
3541 let mut dep = task("dependency");
3542 q.put(&mut dep).unwrap();
3543 let mut blocked = task("waiting");
3544 blocked.block(vec![dep.id.clone()], Some("waits".to_owned()));
3545 q.put(&mut blocked).unwrap();
3546
3547 let removed = q.remove(&dep.id, false, &questions).unwrap();
3548 assert_eq!(removed.released, [blocked.id.clone()]);
3549 assert!(removed.still_blocked.is_empty());
3550
3551 let after = q.get(&blocked.id).unwrap();
3552 assert_eq!(after.status, TaskStatus::Queued);
3553 assert!(after.blocked_by.is_empty());
3554 assert!(after.block_reason.is_none());
3555 let notes = notices_of(dir.path());
3556 assert_eq!(notes.len(), 1, "{notes:?}");
3557 assert_eq!(notes[0].severity, crate::notices::Severity::Info);
3558 }
3559
3560 #[test]
3561 fn removing_one_of_two_dependencies_keeps_the_other() {
3562 let (dir, q) = queue();
3563 let questions = Questions::at(dir.path().join("questions"));
3564 let mut dep = task("dependency");
3565 q.put(&mut dep).unwrap();
3566 let mut other = task("other");
3567 q.put(&mut other).unwrap();
3568 let mut blocked = task("waiting");
3569 blocked.block(
3570 vec![dep.id.clone(), other.id.clone()],
3571 Some("waits on both".to_owned()),
3572 );
3573 q.put(&mut blocked).unwrap();
3574
3575 let removed = q.remove(&dep.id, false, &questions).unwrap();
3576 assert!(removed.released.is_empty());
3577 assert_eq!(removed.still_blocked, [blocked.id.clone()]);
3578
3579 let after = q.get(&blocked.id).unwrap();
3580 assert_eq!(after.status, TaskStatus::Blocked);
3581 assert_eq!(after.blocked_by, [other.id.clone()]);
3582 assert!(missing_blockers(&q, &questions, &after.blocked_by).is_empty());
3585 }
3586
3587 #[test]
3588 fn a_dependency_deleted_before_its_dependents_were_rewritten_is_released_later() {
3589 let (dir, q) = queue();
3590 let questions = Questions::at(dir.path().join("questions"));
3591 let mut dep = task("dependency");
3592 q.put(&mut dep).unwrap();
3593 let mut blocked = task("waiting");
3594 blocked.block(vec![dep.id.clone()], None);
3595 q.put(&mut blocked).unwrap();
3596
3597 let claim = q.claim(&blocked.id).unwrap();
3599 let removed = q.remove(&dep.id, false, &questions).unwrap();
3600 assert!(removed.released.is_empty());
3601 drop(claim);
3602
3603 let mut task = q.get(&blocked.id).unwrap();
3604 assert_eq!(task.status, TaskStatus::Blocked);
3605 assert!(missing_blockers(&q, &questions, &task.blocked_by).is_empty());
3606 assert_eq!(q.apply_deleted_blockers(&mut task), [dep.id.clone()]);
3607 assert_eq!(task.status, TaskStatus::Queued);
3608 }
3609
3610 #[test]
3611 fn a_dependency_without_a_tombstone_is_still_missing() {
3612 let (dir, q) = queue();
3613 let questions = Questions::at(dir.path().join("questions"));
3614 let mut dep = task("dependency");
3615 q.put(&mut dep).unwrap();
3616 std::fs::remove_file(q.path_of(&dep.id)).unwrap();
3617 assert_eq!(
3618 missing_blockers(&q, &questions, std::slice::from_ref(&dep.id)),
3619 [dep.id.clone()]
3620 );
3621 assert!(deleted_blockers(&q, std::slice::from_ref(&dep.id)).is_empty());
3622 }
3623
3624 #[test]
3625 fn a_held_dependent_is_left_alone_by_a_dependency_deletion() {
3626 let mut t = task("held");
3627 t.hold_machine(Some("because".to_owned()));
3628 t.blocked_by = vec!["gone".to_owned()];
3629 assert!(!t.dependency_deleted("gone"));
3630 assert_eq!(t.status, TaskStatus::Held);
3631 }
3632
3633 #[test]
3634 fn a_failed_record_removal_takes_the_tombstone_back() {
3635 let (_dir, q) = queue();
3636 let mut t = task("stays");
3637 q.put(&mut t).unwrap();
3638 q.write_tombstone(&t.id).unwrap();
3639 let err = q.remove_record_with_attachments(&t.id, |_| Err(std::io::Error::other("nope")));
3640 assert!(err.is_err());
3641 assert!(!q.was_deleted(&t.id));
3644 }
3645
3646 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3647 let p = dir.join(name);
3648 std::fs::write(&p, body).unwrap();
3649 p
3650 }
3651
3652 #[test]
3653 fn an_attachment_copy_survives_deleting_its_source() {
3654 let (dir, q) = queue();
3655 let src = source_file(dir.path(), "shot.png", "pixels");
3656 let mut t = task("with a picture");
3657 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3658 q.put(&mut t).unwrap();
3659 std::fs::remove_file(&src).unwrap();
3660 assert_eq!(names, ["shot.png"]);
3661 let loaded = q.get(&t.id).unwrap();
3662 let paths = q.attachment_paths(&loaded);
3663 assert_eq!(paths.len(), 1);
3664 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3665 }
3666
3667 #[test]
3668 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3669 let (dir, q) = queue();
3670 let mut t = task("bad names");
3671 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3672 let src = source_file(dir.path(), name, "x");
3673 assert!(
3674 q.attach(&mut t, &[src]).is_err(),
3675 "`{name}` must be refused"
3676 );
3677 }
3678 assert!(!crate::ask::valid_asset_name("C:foo.png"));
3682 #[cfg(not(windows))]
3683 {
3684 let src = source_file(dir.path(), "C:foo.png", "x");
3685 assert!(q.attach(&mut t, &[src]).is_err());
3686 }
3687 let long = format!("{}.png", "a".repeat(70));
3688 let src = source_file(dir.path(), &long, "x");
3689 assert!(q.attach(&mut t, &[src]).is_err());
3690 assert!(t.attachments.is_empty());
3691 assert!(!q.attachments_dir(&t.id).exists());
3692 }
3693
3694 #[test]
3695 fn a_taken_attachment_name_is_numbered_not_overwritten() {
3696 let (dir, q) = queue();
3697 let a = source_file(dir.path(), "shot.png", "one");
3698 let sub = dir.path().join("other");
3699 std::fs::create_dir_all(&sub).unwrap();
3700 let b = source_file(&sub, "shot.png", "two");
3701 let mut t = task("collision");
3702 q.attach(&mut t, &[a]).unwrap();
3703 q.attach(&mut t, &[b]).unwrap();
3704 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3705 let paths = q.attachment_paths(&t);
3706 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3707 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3708 }
3709
3710 #[test]
3711 fn a_renumbered_name_stays_inside_the_length_bound() {
3712 let (dir, q) = queue();
3713 let name = format!("{}.png", "a".repeat(60));
3714 assert_eq!(name.len(), 64);
3715 let a = source_file(dir.path(), &name, "one");
3716 let sub = dir.path().join("other");
3717 std::fs::create_dir_all(&sub).unwrap();
3718 let b = source_file(&sub, &name, "two");
3719 let mut t = task("long");
3720 q.attach(&mut t, &[a, b]).unwrap();
3721 assert_eq!(t.attachments.len(), 2);
3722 assert!(
3723 t.attachments
3724 .iter()
3725 .all(|n| crate::ask::valid_asset_name(n))
3726 );
3727 assert!(t.attachments[1].ends_with("-2.png"));
3728 }
3729
3730 #[test]
3731 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3732 let (dir, q) = queue();
3733 let good = source_file(dir.path(), "good.png", "ok");
3734 let mut t = task("partial");
3735 q.attach(&mut t, &[good]).unwrap();
3736 let more = source_file(dir.path(), "more.png", "ok");
3737 let missing = dir.path().join("missing.png");
3738 assert!(q.attach(&mut t, &[more, missing]).is_err());
3739 assert_eq!(t.attachments, ["good.png"]);
3740 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3741 .unwrap()
3742 .flatten()
3743 .collect();
3744 assert_eq!(on_disk.len(), 1);
3745 }
3746
3747 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3749 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3750 std::fs::create_dir_all(&tmp).unwrap();
3751 tmp
3752 }
3753
3754 #[test]
3755 fn a_failed_put_leaves_no_new_attachment_directory() {
3756 let (dir, q) = queue();
3757 let mut t = task("fresh");
3758 let tmp = block_put(&q, &t);
3759 let src = source_file(dir.path(), "shot.png", "x");
3760 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3761 assert!(t.attachments.is_empty());
3762 assert!(!q.attachments_dir(&t.id).exists());
3763 assert!(!q.path_of(&t.id).exists());
3764 std::fs::remove_dir(tmp).unwrap();
3765 }
3766
3767 #[test]
3768 fn a_failed_put_removes_only_the_copy_it_just_made() {
3769 let (dir, q) = queue();
3770 let mut t = task("edited");
3771 let first = source_file(dir.path(), "first.png", "1");
3772 q.attach_and_put(&mut t, &[first]).unwrap();
3773 block_put(&q, &t);
3774 let second = source_file(dir.path(), "second.png", "2");
3775 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3776 assert_eq!(t.attachments, ["first.png"]);
3777 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3778 .unwrap()
3779 .flatten()
3780 .map(|e| e.file_name().to_string_lossy().into_owned())
3781 .collect();
3782 assert_eq!(on_disk, ["first.png"]);
3783 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3784 }
3785
3786 #[test]
3787 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3788 let (dir, q) = queue();
3789 let questions = Questions::at(dir.path().join("questions"));
3790 let gone = task("gone");
3791 let mut other = task("other");
3792 let mut live = task("live");
3793 q.put(&mut other).unwrap();
3794 q.put(&mut live).unwrap();
3795 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3798 std::fs::create_dir_all(&orphan).unwrap();
3799 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3800 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3803 std::fs::create_dir_all(&busy).unwrap();
3804
3805 q.remove(&other.id, false, &questions).unwrap();
3806 assert!(!orphan.exists(), "an orphan is swept");
3807 assert!(busy.exists(), "a removal in progress is left alone");
3808 }
3809
3810 #[test]
3811 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3812 let (dir, q) = queue();
3813 let questions = Questions::at(dir.path().join("questions"));
3814 let mut t = task("stuck");
3815 let src = source_file(dir.path(), "shot.png", "x");
3816 q.attach_and_put(&mut t, &[src]).unwrap();
3817 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3820 std::fs::create_dir_all(&aside).unwrap();
3821 std::fs::write(aside.join("old.png"), "o").unwrap();
3822 assert!(q.remove(&t.id, false, &questions).is_err());
3823 assert!(q.path_of(&t.id).exists());
3824 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3825 }
3826
3827 #[test]
3828 fn a_failed_record_removal_puts_the_attachments_back() {
3829 let (dir, q) = queue();
3830 let mut t = task("rollback");
3831 let src = source_file(dir.path(), "shot.png", "x");
3832 q.attach_and_put(&mut t, &[src]).unwrap();
3833 let err = q
3834 .remove_record_with_attachments(&t.id, |_| {
3835 Err(std::io::Error::other("injected failure"))
3836 })
3837 .unwrap_err();
3838 assert!(format!("{err:#}").contains("injected failure"));
3839 assert!(q.path_of(&t.id).exists());
3840 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3841 assert!(
3842 !q.root
3843 .join(format!("{}.attachments.removing", t.id))
3844 .exists()
3845 );
3846 }
3847
3848 #[test]
3849 fn editing_a_task_keeps_its_attachments() {
3850 let (dir, q) = queue();
3851 let src = source_file(dir.path(), "shot.png", "x");
3852 let mut t = task("editable");
3853 q.attach(&mut t, &[src]).unwrap();
3854 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3855 q.put(&mut t).unwrap();
3856 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3857 }
3858
3859 #[test]
3860 fn removing_a_task_deletes_its_attachments() {
3861 let (dir, q) = queue();
3862 let questions = Questions::at(dir.path().join("questions"));
3863 let src = source_file(dir.path(), "shot.png", "x");
3864 let mut t = task("doomed");
3865 q.attach(&mut t, &[src]).unwrap();
3866 q.put(&mut t).unwrap();
3867 assert!(q.attachments_dir(&t.id).is_dir());
3868 q.remove(&t.id, false, &questions).unwrap();
3869 assert!(!q.attachments_dir(&t.id).exists());
3870 assert!(q.list().is_empty());
3871 }
3872
3873 #[test]
3874 fn a_task_written_before_attachments_still_reads() {
3875 let (_dir, q) = queue();
3876 let mut t = task("old");
3877 q.put(&mut t).unwrap();
3878 let path = q.path_of(&t.id);
3879 let mut v: serde_json::Value =
3880 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3881 v.as_object_mut().unwrap().remove("attachments");
3882 std::fs::write(&path, v.to_string()).unwrap();
3883 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3884 }
3885
3886 #[test]
3887 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3888 let q = Queue::at(PathBuf::from("relative-queue"));
3889 let mut t = task("rel");
3890 t.attachments.push("shot.png".to_owned());
3891 let paths = q.attachment_paths(&t);
3892 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3893 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3894 }
3895
3896 #[test]
3897 fn link_run_adds_a_run_once_and_touches_nothing_else() {
3898 let dir = tempfile::tempdir().unwrap();
3899 let queue = Queue::at(dir.path().join("queue"));
3900 let mut t = Task::new(
3901 "t".to_owned(),
3902 "do it".to_owned(),
3903 PathBuf::from("."),
3904 Source::Human,
3905 );
3906 queue.put(&mut t).unwrap();
3907 let before = queue.get(&t.id).unwrap();
3908
3909 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3910 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3911 assert_eq!(linked.status, before.status);
3912 assert_eq!(linked.attempts, before.attempts);
3913 assert_eq!(linked.interrupt, before.interrupt);
3914
3915 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3916 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3917 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3918 assert!(queue.link_run("no-such-task", "r").is_err());
3919 }
3920
3921 #[test]
3922 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
3923 let dir = tempfile::tempdir().unwrap();
3924 let queue = Queue::at(dir.path().join("queue"));
3925 let mut t = Task::new(
3926 "t".to_owned(),
3927 "do it".to_owned(),
3928 PathBuf::from("."),
3929 Source::Human,
3930 );
3931 queue.put(&mut t).unwrap();
3932 let mut snapshot = queue.get(&t.id).unwrap();
3934 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3935
3936 snapshot.start("20260930-000000-aaaa".to_owned());
3937 queue.put(&mut snapshot).unwrap();
3938
3939 let stored = queue.get(&t.id).unwrap();
3940 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3941 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3942 assert_eq!(stored.attempts, 1);
3943 }
3944
3945 #[test]
3946 fn concurrent_links_and_daemon_saves_lose_nothing() {
3947 let dir = tempfile::tempdir().unwrap();
3948 let queue = Queue::at(dir.path().join("queue"));
3949 let mut t = Task::new(
3950 "t".to_owned(),
3951 "do it".to_owned(),
3952 PathBuf::from("."),
3953 Source::Human,
3954 );
3955 queue.put(&mut t).unwrap();
3956 let id = t.id.clone();
3957
3958 let linkers: Vec<_> = (0..4)
3959 .map(|n| {
3960 let (queue, id) = (queue.clone(), id.clone());
3961 std::thread::spawn(move || {
3962 for k in 0..10 {
3963 queue
3964 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3965 .unwrap();
3966 }
3967 })
3968 })
3969 .collect();
3970 let mut mine = queue.get(&id).unwrap();
3973 for k in 0..10 {
3974 mine.start(format!("20260930-000009-d{k:03}"));
3975 queue.put(&mut mine).unwrap();
3976 }
3977 for l in linkers {
3978 l.join().unwrap();
3979 }
3980
3981 let stored = queue.get(&id).unwrap();
3982 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
3983 assert_eq!(
3984 stored.attempts, 10,
3985 "linking never rewinds the daemon's work"
3986 );
3987 }
3988
3989 fn age_lock(path: &Path) {
3990 let f = std::fs::OpenOptions::new().write(true).open(path).unwrap();
3991 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
3992 .unwrap();
3993 }
3994
3995 #[test]
3996 fn concurrent_stale_takeover_yields_one_holder() {
3997 use std::sync::atomic::{AtomicUsize, Ordering};
3998 use std::sync::{Arc, Barrier};
3999 for _ in 0..5 {
4000 let dir = tempfile::tempdir().unwrap();
4001 let q = Queue::at(dir.path().to_path_buf());
4002 let lock = dir.path().join("t.write-lock");
4003 std::fs::write(&lock, "dead-0000").unwrap();
4004 age_lock(&lock);
4005 let n = 6;
4006 let barrier = Arc::new(Barrier::new(n));
4007 let (now, max) = (Arc::new(AtomicUsize::new(0)), Arc::new(AtomicUsize::new(0)));
4008 let handles: Vec<_> = (0..n)
4009 .map(|_| {
4010 let (q, b, now, max) = (q.clone(), barrier.clone(), now.clone(), max.clone());
4011 std::thread::spawn(move || {
4012 b.wait();
4013 let g = q.lock_task("t").unwrap();
4014 let held = now.fetch_add(1, Ordering::SeqCst) + 1;
4015 max.fetch_max(held, Ordering::SeqCst);
4016 std::thread::sleep(std::time::Duration::from_millis(20));
4017 now.fetch_sub(1, Ordering::SeqCst);
4018 drop(g);
4019 })
4020 })
4021 .collect();
4022 for h in handles {
4023 h.join().unwrap();
4024 }
4025 assert_eq!(max.load(Ordering::SeqCst), 1);
4026 assert!(!lock.exists());
4027 }
4028 }
4029
4030 #[test]
4031 fn dropping_a_stolen_lock_leaves_the_new_holders_lock() {
4032 let dir = tempfile::tempdir().unwrap();
4033 let q = Queue::at(dir.path().to_path_buf());
4034 let lock = dir.path().join("t.write-lock");
4035 let a = q.lock_task("t").unwrap();
4036 age_lock(&lock);
4037 let b = q.lock_task("t").unwrap();
4038 assert_ne!(a.token, b.token);
4039 drop(a);
4040 assert_eq!(std::fs::read_to_string(&lock).unwrap(), b.token);
4041 drop(b);
4042 assert!(!lock.exists());
4043 }
4044
4045 #[test]
4046 fn break_stale_leaves_a_lock_that_replaced_the_one_judged() {
4047 let dir = tempfile::tempdir().unwrap();
4048 let lock = dir.path().join("t.write-lock");
4049 std::fs::write(&lock, "new-token").unwrap();
4050 assert!(!break_stale(&lock, "old-token"));
4051 assert_eq!(std::fs::read_to_string(&lock).unwrap(), "new-token");
4052 assert!(!break_stale(&lock, "new-token"));
4054 assert!(lock.exists());
4055 age_lock(&lock);
4056 assert!(break_stale(&lock, "new-token"));
4057 assert!(!lock.exists());
4058 }
4059
4060 #[test]
4061 fn a_stale_break_marker_is_recovered_and_a_fresh_one_is_respected() {
4062 let dir = tempfile::tempdir().unwrap();
4063 let lock = dir.path().join("t.write-lock");
4064 std::fs::write(&lock, "dead-1").unwrap();
4065 age_lock(&lock);
4066 let marker = break_marker(&lock, "dead-1");
4067 std::fs::write(&marker, "crashed-remover").unwrap();
4068 assert!(!break_stale(&lock, "dead-1"));
4070 assert!(lock.exists());
4071 age_lock(&marker);
4073 assert!(break_stale(&lock, "dead-1"));
4074 assert!(!lock.exists());
4075 assert!(!marker.exists());
4076 }
4077}