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 modify(&self, id: &str, f: impl FnOnce(&mut Task) -> bool) -> Result<bool> {
1571 let id = self.resolve_id(id)?;
1572 let _lock = self.lock_task(&id)?;
1573 let mut task = self.get(&id)?;
1574 if !f(&mut task) {
1575 return Ok(false);
1576 }
1577 self.put_unlocked(&mut task)?;
1578 Ok(true)
1579 }
1580
1581 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1586 let id = self.resolve_id(id)?;
1587 let _lock = self.lock_task(&id)?;
1590 let mut task = self.get(&id)?;
1591 if task.link_run(run) {
1592 self.put_unlocked(&mut task)?;
1593 }
1594 Ok(task)
1595 }
1596
1597 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1605 std::fs::create_dir_all(&self.root)
1606 .with_context(|| format!("create {}", self.root.display()))?;
1607 let path = self.root.join(format!("{id}.write-lock"));
1608 let started = std::time::Instant::now();
1609 loop {
1610 match std::fs::OpenOptions::new()
1611 .write(true)
1612 .create_new(true)
1613 .open(&path)
1614 {
1615 Ok(mut file) => {
1616 use std::io::Write;
1617 let token = owner_token();
1618 if let Err(e) = file.write_all(token.as_bytes()) {
1619 drop(file);
1620 let _ = std::fs::remove_file(&path);
1621 return Err(e).with_context(|| format!("lock {}", path.display()));
1622 }
1623 return Ok(TaskLock { path, token });
1624 }
1625 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1626 let judged = std::fs::read_to_string(&path).ok();
1629 let broken = judged
1630 .filter(|_| older_than_stale(&path))
1631 .is_some_and(|judged| break_stale(&path, &judged));
1632 if !broken {
1636 if started.elapsed() > TASK_LOCK_STALE {
1637 bail!("could not lock task {id}");
1638 }
1639 std::thread::sleep(std::time::Duration::from_millis(15));
1640 }
1641 }
1642 Err(e)
1646 if e.kind() == std::io::ErrorKind::PermissionDenied
1647 && started.elapsed() <= TASK_LOCK_STALE =>
1648 {
1649 std::thread::sleep(std::time::Duration::from_millis(15));
1650 }
1651 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1652 }
1653 }
1654 }
1655
1656 pub fn get(&self, id: &str) -> Result<Task> {
1658 let resolved = self.resolve_id(id)?;
1659 read_path(&self.path_of(&resolved))
1660 }
1661
1662 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1694 let _ = questions;
1695 let resolved = self.resolve_id(id)?;
1696 if in_flight {
1697 bail!("task {resolved} is being run by a live daemon right now");
1698 }
1699 self.write_tombstone(&resolved)?;
1700 if let Err(e) = self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))
1701 {
1702 if self.path_of(&resolved).exists() {
1706 let _ = std::fs::remove_file(self.tombstone_path(&resolved));
1707 }
1708 return Err(e);
1709 }
1710 let (released, still_blocked) = self.release_dependents_of(&resolved);
1711 Ok(Removal {
1712 id: resolved,
1713 released,
1714 still_blocked,
1715 })
1716 }
1717
1718 fn tombstone_path(&self, id: &str) -> PathBuf {
1721 self.root.join(format!("{id}.removed"))
1722 }
1723
1724 fn write_tombstone(&self, id: &str) -> Result<()> {
1725 let path = self.tombstone_path(id);
1726 let tmp = self.root.join(format!("{id}.removed.tmp"));
1727 std::fs::write(&tmp, b"")
1728 .and_then(|()| std::fs::rename(&tmp, &path))
1729 .with_context(|| format!("write {}", path.display()))
1730 }
1731
1732 fn was_deleted(&self, id: &str) -> bool {
1736 self.tombstone_path(id).is_file() && !self.path_of(id).exists()
1737 }
1738
1739 pub fn apply_deleted_blockers(&self, task: &mut Task) -> Vec<String> {
1743 let deleted = deleted_blockers(self, &task.blocked_by);
1744 deleted
1745 .into_iter()
1746 .filter(|id| task.dependency_deleted(id))
1747 .collect()
1748 }
1749
1750 pub fn note_dependency_deleted(&self, dependent: &Task, deleted: &str) {
1754 let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty()) else {
1755 return;
1756 };
1757 let ja = crate::lang::is_japanese(&crate::lang::of_repo(&dependent.repo));
1758 let waiting = dependent.status == TaskStatus::Blocked;
1759 let message = match (ja, waiting) {
1760 (true, false) => format!(
1761 "タスク {} は、待っていた {} が削除されたため待機を解除し、元の状態に戻しました",
1762 dependent.short(),
1763 short(deleted)
1764 ),
1765 (true, true) => format!(
1766 "タスク {} は、待っていた {} が削除されたため、残りの依存を待っています",
1767 dependent.short(),
1768 short(deleted)
1769 ),
1770 (false, false) => format!(
1771 "Task {} stopped waiting on {} because it was deleted, and returned to its previous state",
1772 dependent.short(),
1773 short(deleted)
1774 ),
1775 (false, true) => format!(
1776 "Task {} stopped waiting on {} because it was deleted, and is still waiting on its other dependencies",
1777 dependent.short(),
1778 short(deleted)
1779 ),
1780 };
1781 crate::notices::raise_in(
1782 home,
1783 crate::notices::Notice::info(
1784 &format!("unblocked:{}:{}", dependent.id, deleted),
1785 message,
1786 )
1787 .link(crate::notices::Link::Task {
1788 id: dependent.id.clone(),
1789 }),
1790 );
1791 }
1792
1793 fn release_dependents_of(&self, dependency: &str) -> (Vec<String>, Vec<String>) {
1798 let (mut released, mut still_blocked) = (Vec::new(), Vec::new());
1799 for listed in self.list() {
1800 if listed.status != TaskStatus::Blocked
1801 || !listed.blocked_by.iter().any(|b| b == dependency)
1802 {
1803 continue;
1804 }
1805 let Ok(_claim) = self.claim(&listed.id) else {
1806 continue;
1807 };
1808 let Ok(mut task) = self.get(&listed.id) else {
1809 continue;
1810 };
1811 if !task.dependency_deleted(dependency) {
1812 continue;
1813 }
1814 if self.put(&mut task).is_err() {
1815 continue;
1816 }
1817 self.note_dependency_deleted(&task, dependency);
1818 if task.status == TaskStatus::Blocked {
1819 still_blocked.push(task.id.clone());
1820 } else {
1821 released.push(task.id.clone());
1822 }
1823 }
1824 (released, still_blocked)
1825 }
1826
1827 fn remove_record_with_attachments(
1831 &self,
1832 resolved: &str,
1833 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1834 ) -> Result<()> {
1835 self.sweep_removed_attachments();
1842 let attachments = self.attachments_dir(resolved);
1843 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1844 let moved = match std::fs::rename(&attachments, &aside) {
1845 Ok(()) => true,
1846 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1847 Err(e) => {
1848 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1849 }
1850 };
1851 let path = self.path_of(resolved);
1852 if let Err(e) = remove_record(&path) {
1853 if moved {
1854 let _ = std::fs::rename(&aside, &attachments);
1855 }
1856 return Err(e).with_context(|| format!("remove {}", path.display()));
1857 }
1858 if moved {
1859 if let Err(e) = std::fs::remove_dir_all(&aside) {
1860 tracing::warn!("leftover attachments {}: {e}", aside.display());
1861 }
1862 }
1863 let lock = self.lock_path(resolved);
1864 if let Err(e) = std::fs::remove_file(&lock) {
1865 if e.kind() != std::io::ErrorKind::NotFound {
1866 return Err(e).with_context(|| format!("remove {}", lock.display()));
1867 }
1868 }
1869 Ok(())
1870 }
1871
1872 fn sweep_removed_attachments(&self) {
1878 let Ok(entries) = std::fs::read_dir(&self.root) else {
1879 return;
1880 };
1881 for entry in entries.flatten() {
1882 let name = entry.file_name();
1883 let name = name.to_string_lossy();
1884 let Some(id) = name.strip_suffix(".attachments.removing") else {
1885 continue;
1886 };
1887 if !self.path_of(id).exists() {
1890 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1891 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1892 }
1893 }
1894 }
1895 self.sweep_tombstones();
1896 }
1897
1898 fn sweep_tombstones(&self) {
1907 let Ok(entries) = std::fs::read_dir(&self.root) else {
1908 return;
1909 };
1910 let tasks = self.list();
1911 for entry in entries.flatten() {
1912 let name = entry.file_name();
1913 let name = name.to_string_lossy();
1914 let Some(id) = name.strip_suffix(".removed") else {
1915 continue;
1916 };
1917 if self.path_of(id).exists()
1918 || tasks.iter().any(|t| t.blocked_by.iter().any(|b| b == id))
1919 {
1920 continue;
1921 }
1922 let old = entry
1923 .metadata()
1924 .and_then(|m| m.modified())
1925 .ok()
1926 .and_then(|m| m.elapsed().ok())
1927 .is_some_and(|age| age >= TOMBSTONE_GRACE);
1928 if old {
1929 let _ = std::fs::remove_file(entry.path());
1930 }
1931 }
1932 }
1933
1934 fn lock_path(&self, id: &str) -> PathBuf {
1937 self.root.join(format!("{id}.lock"))
1938 }
1939
1940 pub fn list(&self) -> Vec<Task> {
1953 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1954 .into_iter()
1955 .flatten()
1956 .flatten()
1957 .map(|e| e.path())
1958 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1959 .filter_map(|p| read_path(&p).ok())
1960 .collect();
1961 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1962 tasks
1963 }
1964
1965 pub fn superseded(&self) -> HashMap<String, String> {
1978 let mut by = HashMap::new();
1979 for task in self.list() {
1980 for earlier in &task.runs {
1981 if let Some(later) = task.successor_of(earlier) {
1982 by.insert(earlier.clone(), later.clone());
1983 }
1984 }
1985 }
1986 by
1987 }
1988
1989 pub fn superseded_by(&self, run: &str) -> Option<String> {
1997 for task in self.list() {
1998 if task.runs.iter().any(|r| r == run) {
1999 return task.successor_of(run).cloned();
2000 }
2001 }
2002 None
2003 }
2004
2005 pub fn latest_attempt(&self, run: &str) -> Option<String> {
2016 for task in self.list() {
2017 if task.runs.iter().any(|r| r == run) {
2018 return task.runs.last().filter(|last| **last != run).cloned();
2019 }
2020 }
2021 None
2022 }
2023
2024 pub fn next_runnable(&self) -> Option<Task> {
2029 let mut runnable: Vec<Task> = self
2030 .list()
2031 .into_iter()
2032 .filter(|t| t.status.runnable())
2033 .collect();
2034 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
2035 runnable.into_iter().next()
2036 }
2037
2038 pub fn claim(&self, id: &str) -> Result<Claim> {
2045 std::fs::create_dir_all(&self.root)
2046 .with_context(|| format!("create {}", self.root.display()))?;
2047 let path = self.lock_path(id);
2048 match std::fs::OpenOptions::new()
2049 .write(true)
2050 .create_new(true)
2051 .open(&path)
2052 {
2053 Ok(mut f) => {
2054 use std::io::Write as _;
2055 let _ = writeln!(f, "{}", std::process::id());
2057 Ok(Claim { path })
2058 }
2059 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
2060 bail!("task {id} is already claimed ({} exists)", path.display())
2061 }
2062 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
2063 }
2064 }
2065
2066 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
2068 if self.path_of(prefix).is_file() {
2069 return Ok(prefix.to_owned());
2070 }
2071 let hits: Vec<String> = self
2072 .list()
2073 .into_iter()
2074 .map(|t| t.id)
2075 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
2076 .collect();
2077 match hits.len() {
2078 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
2079 0 => bail!("no task matches `{prefix}`"),
2080 _ => bail!(
2081 "`{prefix}` matches {} tasks: {}",
2082 hits.len(),
2083 hits.join(", ")
2084 ),
2085 }
2086 }
2087
2088 pub fn revision(&self) -> u64 {
2095 self.revision_excluding(&std::collections::BTreeSet::new())
2096 }
2097
2098 pub fn revision_excluding(&self, skip: &std::collections::BTreeSet<String>) -> u64 {
2101 use std::hash::{Hash as _, Hasher as _};
2102
2103 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
2104 .into_iter()
2105 .flatten()
2106 .flatten()
2107 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
2108 .filter(|e| {
2109 let path = e.path();
2110 !path
2111 .file_stem()
2112 .is_some_and(|stem| skip.contains(stem.to_string_lossy().as_ref()))
2113 })
2114 .filter_map(|e| {
2115 let name = e.file_name().to_string_lossy().into_owned();
2116 let mtime = e
2117 .metadata()
2118 .ok()?
2119 .modified()
2120 .ok()?
2121 .duration_since(std::time::UNIX_EPOCH)
2122 .ok()?
2123 .as_millis() as u64;
2124 Some((name, mtime))
2125 })
2126 .collect();
2127
2128 if entries.is_empty() {
2129 return 0;
2130 }
2131
2132 entries.sort_unstable();
2133 let mut hasher = std::hash::DefaultHasher::new();
2134 for (name, mtime) in &entries {
2135 name.hash(&mut hasher);
2136 mtime.hash(&mut hasher);
2137 }
2138 let h = hasher.finish();
2139 if h == 0 { 1 } else { h }
2140 }
2141}
2142
2143const TOMBSTONE_GRACE: std::time::Duration = std::time::Duration::from_secs(3600);
2146
2147#[derive(Debug, Clone)]
2149pub struct Removal {
2150 pub id: String,
2152 pub released: Vec<String>,
2156 pub still_blocked: Vec<String>,
2159}
2160
2161#[derive(Debug)]
2163pub struct Claim {
2164 path: PathBuf,
2165}
2166
2167impl Drop for Claim {
2168 fn drop(&mut self) {
2169 let _ = std::fs::remove_file(&self.path);
2170 }
2171}
2172
2173pub fn title_from(instruction: &str, max: usize) -> String {
2176 let Some(line) = first_line(instruction) else {
2177 return "(empty task)".to_owned();
2178 };
2179 if line.chars().count() <= max {
2180 return line.to_owned();
2181 }
2182 let head: String = line.chars().take(max.saturating_sub(1)).collect();
2183 format!("{head}…")
2184}
2185
2186pub(crate) fn first_line(instruction: &str) -> Option<&str> {
2192 let line = instruction
2193 .lines()
2194 .map(str::trim)
2195 .find(|l| !l.is_empty())?
2196 .trim_start_matches(['#', '-', '*', '>', ' '])
2197 .trim();
2198 (!line.is_empty()).then_some(line)
2199}
2200
2201fn read_path(path: &Path) -> Result<Task> {
2202 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2203 let task: Task =
2204 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
2205 if task.schema > SCHEMA {
2211 bail!(
2212 "task {} was written by a different magi (schema {}, this build \
2213 speaks {SCHEMA})",
2214 task.id,
2215 task.schema
2216 );
2217 }
2218 Ok(task)
2219}
2220
2221pub fn missing_blockers(
2235 queue: &Queue,
2236 questions: &Questions,
2237 blocked_by: &[String],
2238) -> Vec<String> {
2239 blocked_by
2240 .iter()
2241 .filter(|id| {
2242 !queue.path_of(id).is_file()
2243 && !questions.path_of(id).is_file()
2244 && !queue.was_deleted(id)
2245 })
2246 .cloned()
2247 .collect()
2248}
2249
2250pub fn deleted_blockers(queue: &Queue, blocked_by: &[String]) -> Vec<String> {
2255 blocked_by
2256 .iter()
2257 .filter(|id| queue.was_deleted(id))
2258 .cloned()
2259 .collect()
2260}
2261
2262pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
2274 missing_blocker_hold_reason_in(blocked_by, missing, "en")
2275}
2276
2277pub fn missing_blocker_hold_reason_in(
2279 blocked_by: &[String],
2280 missing: &[String],
2281 language: &str,
2282) -> String {
2283 if crate::lang::is_japanese(language) {
2284 format!(
2285 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
2286 blocked_by.join(", "),
2287 missing.join(", "),
2288 )
2289 } else {
2290 format!(
2291 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
2292 blocked_by.join(", "),
2293 missing.join(", "),
2294 )
2295 }
2296}
2297
2298pub fn short(id: &str) -> &str {
2300 id.split('-').next_back().unwrap_or(id)
2301}
2302
2303fn new_id() -> String {
2304 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2305 let seed = crate::rng::entropy();
2306 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2307}
2308
2309#[cfg(test)]
2310mod tests {
2311 #[test]
2312 fn merge_over_lets_the_later_choice_win_and_keeps_the_rest() {
2313 let mut base = RunOverrides {
2314 merge: Some("pr".to_owned()),
2315 candidates: Some(3),
2316 ..RunOverrides::default()
2317 };
2318 base.merge_over(&RunOverrides {
2319 merge: Some("none".to_owned()),
2320 seed: Some(7),
2321 ..RunOverrides::default()
2322 });
2323 assert_eq!(base.merge.as_deref(), Some("none"));
2324 assert_eq!(base.candidates, Some(3));
2325 assert_eq!(base.seed, Some(7));
2326 }
2327
2328 #[test]
2329 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
2330 let mut t = task("relanded");
2331 t.attempts = 1;
2332 t.status = TaskStatus::Running;
2333 t.already_landed("already in main as 0e368de");
2334 assert_eq!(t.status, TaskStatus::Done);
2335 assert_eq!(t.attempts, 0);
2336 assert!(t.hold_reason.is_none());
2337 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
2338 }
2339
2340 #[test]
2341 fn missing_blocker_reason_follows_the_language() {
2342 let b = vec!["a".to_owned()];
2343 let en = missing_blocker_hold_reason_in(&b, &b, "en");
2344 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
2345 assert!(en.starts_with("blocked on a"));
2346 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
2347 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
2348 }
2349
2350 use super::*;
2351
2352 #[test]
2353 fn triage_applied_survives_release_and_old_records_read_as_empty() {
2354 let mut t = Task::new(
2355 "t".to_owned(),
2356 "i".to_owned(),
2357 PathBuf::from("r"),
2358 Source::Human,
2359 );
2360 t.mark_triage_applied("q1");
2361 t.mark_triage_applied("q1");
2362 t.hold_machine(Some("x".to_owned()));
2363 t.release();
2364 assert_eq!(t.triage_applied, ["q1"]);
2365 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
2366
2367 let mut v = serde_json::to_value(&t).unwrap();
2368 v.as_object_mut().unwrap().remove("triage_applied");
2369 let old: Task = serde_json::from_value(v).unwrap();
2370 assert!(old.triage_applied.is_empty());
2371 }
2372
2373 #[test]
2374 fn task_counts_of_empty_is_all_zero() {
2375 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
2376 }
2377
2378 #[test]
2379 fn task_counts_of_tallies_every_status() {
2380 let mut queued = Task::new(
2381 "q".to_owned(),
2382 "i".to_owned(),
2383 PathBuf::from("."),
2384 Source::Human,
2385 );
2386 queued.status = TaskStatus::Queued;
2387 let mut running = queued.clone();
2388 running.status = TaskStatus::Running;
2389 let mut done = queued.clone();
2390 done.status = TaskStatus::Done;
2391 let mut failed = queued.clone();
2392 failed.status = TaskStatus::Failed;
2393 let mut held = queued.clone();
2394 held.status = TaskStatus::Held;
2395 let mut blocked = queued.clone();
2396 blocked.status = TaskStatus::Blocked;
2397
2398 let mut parked = queued.clone();
2399 parked.status = TaskStatus::Parked;
2400
2401 let counts = TaskCounts::of(&[
2402 queued,
2403 running,
2404 done.clone(),
2405 done,
2406 failed,
2407 held,
2408 blocked,
2409 parked,
2410 ]);
2411 assert_eq!(
2412 counts,
2413 TaskCounts {
2414 queued: 1,
2415 running: 1,
2416 done: 2,
2417 failed: 1,
2418 held: 1,
2419 blocked: 1,
2420 parked: 1,
2421 }
2422 );
2423 }
2424
2425 fn queue() -> (tempfile::TempDir, Queue) {
2428 let dir = tempfile::tempdir().unwrap();
2429 let q = Queue::at(dir.path().join("queue"));
2430 (dir, q)
2431 }
2432
2433 #[test]
2434 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
2435 let dir = tempfile::tempdir().unwrap();
2436 let q = Queue::at(dir.path().join("queue"));
2437 let mut t = task("held");
2438 q.put(&mut t).unwrap();
2439 assert_eq!(
2440 crate::notices::Notices::at(dir.path().join("notifications"))
2441 .list()
2442 .len(),
2443 0
2444 );
2445 t.hold_machine(Some("out of attempts".to_owned()));
2446 q.put(&mut t).unwrap();
2447 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
2448 assert_eq!(listed.len(), 1);
2449 assert!(listed[0].message.contains("out of attempts"));
2450 }
2451
2452 fn task(title: &str) -> Task {
2453 Task::new(
2454 title.to_owned(),
2455 format!("do {title}"),
2456 PathBuf::from("."),
2457 Source::Human,
2458 )
2459 }
2460
2461 #[test]
2462 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
2463 let mut t = task("retried");
2464 assert!(
2465 t.earlier_attempts().is_empty(),
2466 "a first attempt takes nothing over"
2467 );
2468 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2469 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
2470 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
2472 assert_eq!(t.successor_of("bbbb"), None);
2473 assert_eq!(t.successor_of("zzzz"), None);
2474 }
2475
2476 #[test]
2477 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
2478 let (_dir, q) = queue();
2479 let mut t = task("retried");
2480 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2481 q.put(&mut t).unwrap();
2482
2483 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
2484 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
2485 assert_eq!(
2486 q.superseded_by("cccc"),
2487 None,
2488 "the latest attempt replaces nothing"
2489 );
2490 assert_eq!(
2491 q.superseded_by("never-heard-of-it"),
2492 None,
2493 "a run belonging to no task on this queue is not superseded"
2494 );
2495
2496 let mut by = HashMap::new();
2497 by.insert("aaaa".to_owned(), "bbbb".to_owned());
2498 by.insert("bbbb".to_owned(), "cccc".to_owned());
2499 assert_eq!(
2500 q.superseded(),
2501 by,
2502 "the whole-map and single-run forms must agree"
2503 );
2504 }
2505
2506 #[test]
2507 fn superseded_attempts_is_empty_until_the_task_is_done() {
2508 let mut t = task("retried");
2509 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2510 t.status = TaskStatus::Failed;
2511 assert_eq!(
2512 t.superseded_attempts(true),
2513 &[] as &[String],
2514 "a task still retrying has no attempt yet that a later one made moot"
2515 );
2516
2517 t.status = TaskStatus::Running;
2518 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2519 }
2520
2521 #[test]
2522 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
2523 let mut t = task("retried");
2524 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2525 t.status = TaskStatus::Done;
2526 assert_eq!(
2527 t.superseded_attempts(true),
2528 &["aaaa".to_owned(), "bbbb".to_owned()],
2529 "cccc is the attempt whose success made the task done, and stays out"
2530 );
2531 }
2532
2533 #[test]
2534 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
2535 let mut t = task("first try landed");
2536 t.runs = vec!["aaaa".to_owned()];
2537 t.status = TaskStatus::Done;
2538 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2539 }
2540
2541 #[test]
2542 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
2543 let mut t = task("closed by hand after a manual merge");
2550 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2551 t.status = TaskStatus::Done;
2552 assert_eq!(
2553 t.superseded_attempts(false),
2554 &[] as &[String],
2555 "nothing here is provably why the task is done, so nothing is superseded"
2556 );
2557 }
2558
2559 #[test]
2560 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2561 let (_dir, q) = queue();
2562 let mut t = task("retried twice");
2563 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2564 q.put(&mut t).unwrap();
2565
2566 assert_eq!(
2567 q.latest_attempt("aaaa"),
2568 Some("cccc".to_owned()),
2569 "an old attempt points straight at the chain's current head, not the \
2570 next attempt in the middle of it"
2571 );
2572 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2573 assert_eq!(
2574 q.latest_attempt("cccc"),
2575 None,
2576 "the latest attempt is not superseded by anything"
2577 );
2578 assert_eq!(
2579 q.latest_attempt("never-heard-of-it"),
2580 None,
2581 "a run belonging to no task on this queue is not superseded"
2582 );
2583 }
2584
2585 #[test]
2586 fn a_markdown_heading_is_the_title_not_decoration() {
2587 assert_eq!(
2592 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2593 "Rework the config loader"
2594 );
2595 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2596 assert_eq!(title_from("> quoted task", 40), "quoted task");
2597 assert_eq!(title_from(" \n\n", 40), "(empty task)");
2599 assert_eq!(title_from("###\n", 40), "(empty task)");
2600 }
2601
2602 #[test]
2603 fn a_long_title_is_elided_by_characters_not_bytes() {
2604 let long = "課題".repeat(30);
2606 let title = title_from(&long, 10);
2607 assert_eq!(title.chars().count(), 10);
2608 assert!(title.ends_with('…'));
2609 }
2610
2611 #[test]
2612 fn priority_wins_and_ties_break_oldest_first() {
2613 let (_dir, q) = queue();
2614 let mut a = task("first");
2615 let mut b = task("second");
2616 let mut c = task("urgent");
2617 a.id = "20260101-000001-aaaa".to_owned();
2619 b.id = "20260101-000002-bbbb".to_owned();
2620 c.id = "20260101-000003-cccc".to_owned();
2621 c.priority = 5;
2622 for t in [&mut a, &mut b, &mut c] {
2623 q.put(t).unwrap();
2624 }
2625
2626 assert_eq!(q.next_runnable().unwrap().id, c.id);
2628 c.hold_machine(None);
2629 q.put(&mut c).unwrap();
2630 assert_eq!(q.next_runnable().unwrap().id, a.id);
2632 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2633 }
2634
2635 #[test]
2636 fn a_blocked_task_never_starves_another_runnable_one() {
2637 let (_dir, q) = queue();
2638 let mut blocked = task("blocked");
2639 blocked.block(vec!["something".to_owned()], None);
2640 q.put(&mut blocked).unwrap();
2641
2642 let mut runnable = task("free to go");
2643 q.put(&mut runnable).unwrap();
2644
2645 let next = q.next_runnable().expect("a runnable task is still offered");
2646 assert_eq!(next.id, runnable.id);
2647 }
2648
2649 #[test]
2650 fn a_held_task_is_never_offered_to_the_loop() {
2651 let (_dir, q) = queue();
2652 let mut t = task("held");
2653 q.put(&mut t).unwrap();
2654 assert!(q.next_runnable().is_some());
2655
2656 t.hold_machine(None);
2657 q.put(&mut t).unwrap();
2658 assert!(
2659 q.next_runnable().is_none(),
2660 "a held task must wait for a human"
2661 );
2662
2663 t.status = TaskStatus::Failed;
2665 q.put(&mut t).unwrap();
2666 assert!(q.next_runnable().is_some());
2667 }
2668
2669 #[test]
2670 fn attempts_are_capped_and_then_the_task_is_held() {
2671 let mut t = task("doomed");
2672
2673 t.start("run-1".to_owned());
2674 t.fail("gate red", 2);
2675 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2676
2677 t.start("run-2".to_owned());
2678 t.fail("gate red", 2);
2679 assert_eq!(
2680 t.status,
2681 TaskStatus::Held,
2682 "out of attempts: stop spending money on it"
2683 );
2684 assert_eq!(t.runs, ["run-1", "run-2"]);
2685 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2686 assert_eq!(
2687 t.hold_reason.as_deref(),
2688 Some("gate red"),
2689 "the hold must say why, not leave hold_reason null next to a \
2690 populated last_error"
2691 );
2692 }
2693
2694 #[test]
2695 fn handing_off_a_task_records_a_hold_reason_too() {
2696 let mut t = task("left a pull request");
2697 t.start("run-1".to_owned());
2698 t.handed_off("run ended with a pull request open [run run-1]");
2699 assert_eq!(t.status, TaskStatus::Held);
2700 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2701 assert_eq!(
2702 t.hold_reason.as_deref(),
2703 Some("run ended with a pull request open [run run-1]")
2704 );
2705 assert_eq!(t.hold_reason, t.last_error);
2706 }
2707
2708 #[test]
2709 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2710 let mut t = task("stalled by quota");
2711
2712 t.start("run-1".to_owned());
2713 assert_eq!(t.attempts, 1);
2714 t.stall("judge-1, judge-2 out of quota");
2715 assert_eq!(
2716 t.attempts, 0,
2717 "a closed quota window must not spend the task's retry budget"
2718 );
2719 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2720 assert_eq!(
2721 t.last_error.as_deref(),
2722 Some("judge-1, judge-2 out of quota")
2723 );
2724
2725 for _ in 0..20 {
2728 t.start("run-n".to_owned());
2729 t.stall("still out of quota");
2730 }
2731 t.start("run-real".to_owned());
2732 t.fail("gate red", 2);
2733 assert_eq!(
2734 t.status,
2735 TaskStatus::Failed,
2736 "the first attempt that was really judged is attempt one"
2737 );
2738 }
2739
2740 #[test]
2741 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2742 let mut t = task("retry me");
2743 t.start("run-1".to_owned());
2744 t.fail("gate red", 1);
2745 assert_eq!(t.status, TaskStatus::Held);
2746
2747 t.release();
2748 assert_eq!(t.status, TaskStatus::Queued);
2749 assert_eq!(t.attempts, 0);
2752 assert!(t.last_error.is_none());
2753 assert_eq!(
2754 t.runs.len(),
2755 1,
2756 "history is kept: attempts reset, evidence does not"
2757 );
2758 }
2759
2760 #[test]
2761 fn a_hold_reason_survives_and_a_release_clears_it() {
2762 let mut t = task("waiting on something else");
2763 t.hold_manual(Some(
2764 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2765 ));
2766 assert_eq!(t.status, TaskStatus::Held);
2767 assert_eq!(
2768 t.hold_reason.as_deref(),
2769 Some("waiting for 20260101-000000-aaaa to land first")
2770 );
2771
2772 t.hold_manual(None);
2774 assert_eq!(
2775 t.hold_reason.as_deref(),
2776 Some("waiting for 20260101-000000-aaaa to land first"),
2777 "a bare re-hold keeps whatever a human already wrote down"
2778 );
2779
2780 let mut plain = task("no reason given");
2782 plain.hold_manual(None);
2783 assert_eq!(plain.status, TaskStatus::Held);
2784 assert!(plain.hold_reason.is_none());
2785
2786 t.release();
2787 assert_eq!(t.status, TaskStatus::Queued);
2788 assert!(
2789 t.hold_reason.is_none(),
2790 "a stale reason must not greet the next person who holds this task"
2791 );
2792 }
2793
2794 #[test]
2795 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2796 let mut t = task("landed by hand while held");
2801 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2802 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2803
2804 t.succeed();
2805 assert_eq!(t.status, TaskStatus::Done);
2806 assert!(
2807 t.hold_reason.is_none(),
2808 "a done task cannot still be waiting on something"
2809 );
2810 }
2811
2812 #[test]
2813 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2814 let mut held = task("held straight out of blocked");
2821 held.block(
2822 vec!["20260101-000000-dead".to_owned()],
2823 Some("waiting on the migration script".to_owned()),
2824 );
2825 assert_eq!(held.status, TaskStatus::Blocked);
2826
2827 held.hold_manual(None);
2828 assert_eq!(held.status, TaskStatus::Held);
2829 assert!(
2830 held.blocked_by.is_empty(),
2831 "hold overrides the wait, same as release"
2832 );
2833 assert!(held.block_reason.is_none());
2834
2835 let mut done = task("closed straight out of blocked");
2836 done.block(
2837 vec!["20260101-000000-dead".to_owned()],
2838 Some("waiting on the migration script".to_owned()),
2839 );
2840 done.succeed();
2841 assert_eq!(done.status, TaskStatus::Done);
2842 assert!(
2843 done.blocked_by.is_empty(),
2844 "a done task cannot still be waiting on a dependency"
2845 );
2846 assert!(done.block_reason.is_none());
2847 }
2848
2849 #[test]
2850 fn a_blocked_task_is_never_offered_to_the_loop() {
2851 let mut t = task("blocked");
2852 assert!(t.status.runnable());
2853 t.block(
2854 vec!["dep-id".to_owned()],
2855 Some("waits on dep-id".to_owned()),
2856 );
2857 assert_eq!(t.status, TaskStatus::Blocked);
2858 assert!(!t.status.runnable());
2859 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2860 }
2861
2862 #[test]
2863 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2864 let mut t = task("blocked on two");
2865 t.block(
2866 vec!["a".to_owned(), "b".to_owned()],
2867 Some("waits on a and b".to_owned()),
2868 );
2869
2870 t.unblock("a");
2871 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2872 assert_eq!(t.blocked_by, ["b"]);
2873
2874 t.unblock("b");
2875 assert_eq!(t.status, TaskStatus::Queued);
2876 assert!(t.blocked_by.is_empty());
2877 assert!(t.block_reason.is_none());
2878 }
2879
2880 #[test]
2881 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2882 let mut t = task("never blocked");
2883 t.unblock("whatever");
2884 assert_eq!(t.status, TaskStatus::Queued);
2885 }
2886
2887 #[test]
2888 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2889 let mut t = task("held, then asked about");
2894 t.hold_machine(Some("out of attempts".to_owned()));
2895 assert_eq!(t.status, TaskStatus::Held);
2896
2897 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2898 assert_eq!(t.status, TaskStatus::Blocked);
2899
2900 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2901 t.unblock("q1");
2902 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2903 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2904 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2905 assert!(t.blocked_from.is_none(), "consumed once restored");
2906 }
2907
2908 #[test]
2909 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2910 let mut t = task("manually held, then asked about");
2911 t.hold_manual(Some("waiting on a dependency".to_owned()));
2912
2913 t.block(vec!["q1".to_owned()], None);
2914 t.unblock("q1");
2915
2916 assert_eq!(t.status, TaskStatus::Held);
2917 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2918 }
2919
2920 #[test]
2921 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2922 let mut t = task("held, blocked twice");
2926 t.hold_machine(None);
2927 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2928 t.block(
2929 vec!["q1".to_owned(), "q2".to_owned()],
2930 Some("second".to_owned()),
2931 );
2932
2933 t.unblock("q1");
2934 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2935 t.unblock("q2");
2936 assert_eq!(t.status, TaskStatus::Held);
2937 }
2938
2939 #[test]
2940 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2941 let mut t = task("blocked mid-run");
2945 t.start("run-1".to_owned());
2946 assert_eq!(t.status, TaskStatus::Running);
2947
2948 t.block(vec!["q1".to_owned()], None);
2949 t.unblock("q1");
2950 assert_eq!(t.status, TaskStatus::Queued);
2951 }
2952
2953 #[test]
2954 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2955 let mut t = task("legacy record, held before it was blocked");
2961 t.hold_source = Some(HoldSource::Machine);
2962 t.hold_reason = Some("legacy hold reason".to_owned());
2963 t.status = TaskStatus::Blocked;
2964 t.blocked_by = vec!["q1".to_owned()];
2965 t.blocked_from = None;
2966
2967 t.unblock("q1");
2968 assert_eq!(t.status, TaskStatus::Held);
2969 }
2970
2971 #[test]
2972 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2973 let mut t = task("legacy record, ordinary dependency block");
2974 t.status = TaskStatus::Blocked;
2975 t.blocked_by = vec!["dep".to_owned()];
2976 t.blocked_from = None;
2977
2978 t.unblock("dep");
2979 assert_eq!(t.status, TaskStatus::Queued);
2980 }
2981
2982 #[test]
2983 fn answering_a_question_is_recorded_and_survives_a_release() {
2984 let mut t = task("asked something");
2985 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2986 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2987 t.unblock("q1");
2988 assert_eq!(t.status, TaskStatus::Queued);
2989 assert_eq!(t.answers.len(), 1);
2990 assert_eq!(t.answers[0].answer, "SQLite");
2991
2992 t.release();
2996 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2997 }
2998
2999 #[test]
3000 fn a_refused_handover_keeps_the_review_branch_across_release() {
3001 let mut t = task("refused takeover");
3002 t.start("run-1".to_owned());
3003 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
3004 assert_eq!(t.status, TaskStatus::Held);
3005 assert_eq!(t.attempts, 1);
3006 t.release();
3007 assert_eq!(t.status, TaskStatus::Queued);
3008 assert_eq!(t.attempts, 0);
3009 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
3010
3011 t.hold_for_handover(None, "again".to_owned());
3013 t.requeue();
3014 assert!(t.review_branch.is_none());
3015
3016 let mut m = task("manual");
3018 m.review_branch = Some("magi/x/A".to_owned());
3019 m.hold_manual(None);
3020 m.release();
3021 assert!(m.review_branch.is_none());
3022 }
3023
3024 #[test]
3025 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
3026 let mut t = task("blocked run with a surviving branch");
3027 t.start("run-1".to_owned());
3028 t.fail("blocked with major findings", 5);
3029 assert_eq!(t.status, TaskStatus::Failed);
3030
3031 t.request_review("magi/eba2/A".to_owned());
3032 assert_eq!(t.status, TaskStatus::Queued);
3033 assert_eq!(t.attempts, 0);
3034 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
3035
3036 t.release();
3038 assert!(t.review_branch.is_none());
3039 }
3040
3041 #[test]
3042 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
3043 let mut t = task("retry");
3044 t.start("run-1".to_owned());
3045 t.requeue();
3046 assert!(t.fresh_start);
3047
3048 t.release();
3049 assert!(!t.fresh_start);
3050 }
3051
3052 #[test]
3053 fn priority_can_be_changed_while_queued_but_not_while_running() {
3054 let mut t = task("reprioritise me");
3055 t.set_priority(5).unwrap();
3056 assert_eq!(t.priority, 5);
3057
3058 t.start("run-1".to_owned());
3059 let err = t.set_priority(9).unwrap_err().to_string();
3060 assert!(err.contains("running"), "{err}");
3061 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
3062 }
3063
3064 #[test]
3065 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
3066 let mut t = task("interrupt me");
3067 assert!(!t.interrupt, "off unless asked, same as any other task");
3068
3069 t.set_interrupt(true).unwrap();
3070 assert!(t.interrupt);
3071
3072 t.start("run-1".to_owned());
3073 assert!(
3074 !t.interrupt,
3075 "the mark is one-shot: dispatching the task fulfils it, \
3076 whatever the run that follows ends up doing"
3077 );
3078 let err = t.set_interrupt(true).unwrap_err().to_string();
3079 assert!(err.contains("running"), "{err}");
3080 t.set_interrupt(false).unwrap();
3083 assert!(!t.interrupt);
3084 }
3085
3086 #[test]
3090 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
3091 let mut t = task("interrupt me");
3092 t.set_interrupt(true).unwrap();
3093 t.start("run-1".to_owned());
3094 t.fail("mock failure", 5);
3095 assert_eq!(t.status, TaskStatus::Failed);
3096 assert!(
3097 !t.interrupt,
3098 "one attempt already spent the mark; a retry is an ordinary \
3099 requeue, not a fresh interrupt request"
3100 );
3101 }
3102
3103 #[test]
3104 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
3105 let (_dir, q) = queue();
3106 let mut a = task("first filed");
3107 let mut b = task("second filed");
3108 a.id = "20260101-000001-aaaa".to_owned();
3109 b.id = "20260101-000002-bbbb".to_owned();
3110 q.put(&mut a).unwrap();
3111 q.put(&mut b).unwrap();
3112
3113 assert_eq!(
3114 q.next_runnable().unwrap().id,
3115 a.id,
3116 "with equal priority the older task goes first, so a burst of \
3117 new work cannot starve it"
3118 );
3119 assert_eq!(
3120 q.list()[0].id,
3121 b.id,
3122 "but the list an operator reads is newest first, the same as \
3123 before priority existed - a's turn to run does not make it the \
3124 newest task"
3125 );
3126
3127 let mut a = q.get(&a.id).unwrap();
3128 a.set_priority(10).unwrap();
3129 q.put(&mut a).unwrap();
3130
3131 assert_eq!(
3132 q.next_runnable().unwrap().id,
3133 a.id,
3134 "a raised priority must be reflected the moment it is saved"
3135 );
3136 assert_eq!(
3140 q.list()[0].id,
3141 a.id,
3142 "the raised task must sort first in the list an operator reads, \
3143 not only in next_runnable's own ordering"
3144 );
3145 }
3146
3147 #[test]
3148 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
3149 let mut t = Task::new(
3150 "old title".to_owned(),
3151 "old instruction".to_owned(),
3152 PathBuf::from("/repo"),
3153 Source::Agent {
3154 run: "20260101-000000-beef".to_owned(),
3155 node: "implement".to_owned(),
3156 },
3157 );
3158 let id = t.id.clone();
3159 let created_at = t.created_at;
3160 t.runs.push("20260101-000000-beef".to_owned());
3161
3162 t.edit("new title".to_owned(), "new instruction".to_owned())
3163 .unwrap();
3164
3165 assert_eq!(t.title, "new title");
3166 assert_eq!(t.instruction, "new instruction");
3167 assert_eq!(t.id, id, "editing must not mint a new id");
3168 assert_eq!(t.created_at, created_at);
3169 assert_eq!(
3170 t.source,
3171 Source::Agent {
3172 run: "20260101-000000-beef".to_owned(),
3173 node: "implement".to_owned(),
3174 },
3175 "editing must not turn agent attribution into human"
3176 );
3177 assert_eq!(t.runs, ["20260101-000000-beef"]);
3178 }
3179
3180 #[test]
3181 fn editing_is_refused_once_a_task_is_running_or_finished() {
3182 let mut running = task("in flight");
3183 running.start("run-1".to_owned());
3184 let err = running
3185 .edit("x".to_owned(), "y".to_owned())
3186 .unwrap_err()
3187 .to_string();
3188 assert!(err.contains("running"), "{err}");
3189
3190 let mut done = task("finished");
3191 done.succeed();
3192 let err = done
3193 .edit("x".to_owned(), "y".to_owned())
3194 .unwrap_err()
3195 .to_string();
3196 assert!(err.contains("done"), "{err}");
3197
3198 let mut queued = task("waiting");
3200 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
3201 let mut held = task("parked");
3202 held.hold_machine(None);
3203 held.edit("x".to_owned(), "y".to_owned()).unwrap();
3204 }
3205
3206 #[test]
3207 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
3208 let (_dir, q) = queue();
3209 let path = q.path_of("20260101-000000-aaaa");
3210 std::fs::create_dir_all(q.root()).unwrap();
3211 std::fs::write(
3212 &path,
3213 serde_json::json!({
3214 "schema": SCHEMA,
3215 "id": "20260101-000000-aaaa",
3216 "title": "from before hold reasons existed",
3217 "instruction": "from before hold reasons existed",
3218 "repo": ".",
3219 "source": { "kind": "human" },
3220 "status": "held",
3221 "created_at": Timestamp::now().to_string(),
3222 "updated_at": Timestamp::now().to_string(),
3223 })
3224 .to_string(),
3225 )
3226 .unwrap();
3227
3228 let task = q.get("20260101-000000-aaaa").expect("must still read");
3229 assert!(task.hold_reason.is_none());
3230 assert!(task.operator_held());
3231 }
3232
3233 #[test]
3234 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
3235 let (_dir, q) = queue();
3236 let path = q.path_of("20260101-000000-bbbb");
3237 std::fs::create_dir_all(q.root()).unwrap();
3238 std::fs::write(
3239 &path,
3240 serde_json::json!({
3241 "schema": 2,
3242 "id": "20260101-000000-bbbb",
3243 "title": "old manual recovery",
3244 "instruction": "old manual recovery",
3245 "repo": ".",
3246 "source": { "kind": "human" },
3247 "status": "held",
3248 "hold_reason": "active manual recovery run20260912-224242-daf5",
3249 "created_at": Timestamp::now().to_string(),
3250 "updated_at": Timestamp::now().to_string(),
3251 })
3252 .to_string(),
3253 )
3254 .unwrap();
3255
3256 let task = q.get("20260101-000000-bbbb").expect("must still read");
3257 assert_eq!(task.hold_source, None);
3258 assert!(task.operator_held());
3259 }
3260
3261 #[test]
3262 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
3263 let (_dir, q) = queue();
3264 let path = q.path_of("20260101-000000-aaaa");
3265 std::fs::create_dir_all(q.root()).unwrap();
3266 std::fs::write(
3267 &path,
3268 serde_json::json!({
3269 "schema": SCHEMA,
3270 "id": "20260101-000000-aaaa",
3271 "title": "from before diagnostics existed",
3272 "instruction": "from before diagnostics existed",
3273 "repo": ".",
3274 "source": { "kind": "human" },
3275 "status": "held",
3276 "created_at": Timestamp::now().to_string(),
3277 "updated_at": Timestamp::now().to_string(),
3278 })
3279 .to_string(),
3280 )
3281 .unwrap();
3282
3283 let task = q.get("20260101-000000-aaaa").expect("must still read");
3284 assert!(task.diagnostic.is_none());
3285 }
3286
3287 #[test]
3288 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
3289 let (_dir, q) = queue();
3293 let path = q.path_of("20260101-000000-aaaa");
3294 std::fs::create_dir_all(q.root()).unwrap();
3295 std::fs::write(
3296 &path,
3297 serde_json::json!({
3298 "schema": 1,
3299 "id": "20260101-000000-aaaa",
3300 "title": "from before blocking existed",
3301 "instruction": "from before blocking existed",
3302 "repo": ".",
3303 "source": { "kind": "human" },
3304 "status": "queued",
3305 "created_at": Timestamp::now().to_string(),
3306 "updated_at": Timestamp::now().to_string(),
3307 })
3308 .to_string(),
3309 )
3310 .unwrap();
3311
3312 let task = q.get("20260101-000000-aaaa").expect("must still read");
3313 assert!(task.blocked_by.is_empty());
3314 assert!(task.block_reason.is_none());
3315 assert!(task.answers.is_empty());
3316 assert!(task.review_branch.is_none());
3317 }
3318
3319 #[test]
3320 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
3321 let mut held = task("diagnosed");
3326 held.start("run-1".to_owned());
3327 held.fail("gate red", 1);
3328 held.diagnostic = Some("cargo test failed: ...".to_owned());
3329 assert_eq!(held.status, TaskStatus::Held);
3330
3331 held.release();
3332 assert!(held.diagnostic.is_none());
3333
3334 held.diagnostic = Some("cargo test failed: ...".to_owned());
3335 held.succeed();
3336 assert!(held.diagnostic.is_none());
3337 }
3338
3339 #[test]
3340 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
3341 let mut t = task("retried");
3342 t.start("run-1".to_owned());
3343 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
3344 t.fail("unrelated config error", 5);
3345 assert_eq!(t.status, TaskStatus::Failed);
3346 assert!(
3347 t.diagnostic.is_none(),
3348 "fail() must not let an old diagnostic outlive the run that produced it"
3349 );
3350 }
3351
3352 #[test]
3353 fn a_claim_is_exclusive_and_releases_on_drop() {
3354 let (_dir, q) = queue();
3355 let mut t = task("contended");
3356 q.put(&mut t).unwrap();
3357
3358 let held = q.claim(&t.id).unwrap();
3359 assert!(
3360 q.claim(&t.id).is_err(),
3361 "two daemons must not drive one task into two runs"
3362 );
3363 drop(held);
3364 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
3365 }
3366
3367 #[test]
3368 fn a_round_trip_survives_disk() {
3369 let (_dir, q) = queue();
3370 let mut t = Task::new(
3371 "titled".to_owned(),
3372 "body".to_owned(),
3373 PathBuf::from("/repo"),
3374 Source::Agent {
3375 run: "20260101-000000-beef".to_owned(),
3376 node: "implement".to_owned(),
3377 },
3378 );
3379 t.priority = 3;
3380 q.put(&mut t).unwrap();
3381
3382 let back = q.get(&t.id).unwrap();
3383 assert_eq!(back.id, t.id);
3384 assert_eq!(back.priority, 3);
3385 assert_eq!(back.source.label(), "implement@beef");
3386 assert_eq!(q.get(t.short()).unwrap().id, t.id);
3388 }
3389
3390 #[test]
3391 fn an_unreadable_task_does_not_take_the_queue_down() {
3392 let (_dir, q) = queue();
3393 let mut t = task("fine");
3394 q.put(&mut t).unwrap();
3395 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
3396
3397 let listed = q.list();
3398 assert_eq!(listed.len(), 1, "the readable task still lists");
3399 assert_eq!(listed[0].id, t.id);
3400 }
3401
3402 #[test]
3403 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
3404 let (_dir, q) = queue();
3405 let path = q.path_of("20260101-000000-aaaa");
3406 std::fs::create_dir_all(q.root()).unwrap();
3407 std::fs::write(
3408 &path,
3409 serde_json::json!({
3410 "schema": SCHEMA,
3411 "id": "20260101-000000-aaaa",
3412 "title": "from before solo existed",
3413 "instruction": "from before solo existed",
3414 "repo": ".",
3415 "source": { "kind": "human" },
3416 "status": "queued",
3417 "created_at": Timestamp::now().to_string(),
3418 "updated_at": Timestamp::now().to_string(),
3419 })
3420 .to_string(),
3421 )
3422 .unwrap();
3423
3424 let task = q.get("20260101-000000-aaaa").expect("must still read");
3425 assert!(!task.solo, "a queue file with no `solo` field means false");
3426 }
3427
3428 #[test]
3429 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
3430 let (_dir, q) = queue();
3431 let path = q.path_of("20260101-000000-bbbb");
3432 std::fs::create_dir_all(q.root()).unwrap();
3433 std::fs::write(
3434 &path,
3435 serde_json::json!({
3436 "schema": SCHEMA,
3437 "id": "20260101-000000-bbbb",
3438 "title": "from before urgent existed",
3439 "instruction": "from before urgent existed",
3440 "repo": ".",
3441 "source": { "kind": "human" },
3442 "status": "queued",
3443 "created_at": Timestamp::now().to_string(),
3444 "updated_at": Timestamp::now().to_string(),
3445 })
3446 .to_string(),
3447 )
3448 .unwrap();
3449
3450 let task = q.get("20260101-000000-bbbb").expect("must still read");
3451 assert!(
3452 !task.urgent,
3453 "a queue file with no `urgent` field means false, same as `solo`"
3454 );
3455 }
3456
3457 #[test]
3458 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
3459 let (_dir, q) = queue();
3460 let mut t = task("from the future");
3461 q.put(&mut t).unwrap();
3462 let path = q.path_of(&t.id);
3463 let body = std::fs::read_to_string(&path)
3464 .unwrap()
3465 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
3466 std::fs::write(&path, body).unwrap();
3467
3468 let err = q.get(&t.id).unwrap_err().to_string();
3469 assert!(err.contains("schema 99"), "{err}");
3470 }
3471
3472 #[test]
3473 fn revision_moves_when_the_queue_changes() {
3474 let (_dir, q) = queue();
3475 assert_eq!(q.revision(), 0, "an empty queue has no revision");
3476 let mut t = task("first");
3477 q.put(&mut t).unwrap();
3478 assert!(q.revision() > 0, "a written task moves the revision");
3479 }
3480
3481 #[test]
3482 fn revision_moves_when_deleting_an_older_task() {
3483 let (dir, q) = queue();
3484 let questions = Questions::at(dir.path().join("questions"));
3485 let mut t1 = task("older");
3486 q.put(&mut t1).unwrap();
3487 std::thread::sleep(std::time::Duration::from_millis(10));
3489 let mut t2 = task("newer");
3490 q.put(&mut t2).unwrap();
3491
3492 let rev_before = q.revision();
3493 q.remove(&t1.id, false, &questions).unwrap();
3494 let rev_after = q.revision();
3495
3496 assert_ne!(
3497 rev_before, rev_after,
3498 "deleting an older task must change the revision so other clients see the deletion"
3499 );
3500 }
3501
3502 #[test]
3503 fn removing_a_task_takes_it_out_of_the_listing() {
3504 let (dir, q) = queue();
3505 let questions = Questions::at(dir.path().join("questions"));
3506 let mut t = task("delete me");
3507 q.put(&mut t).unwrap();
3508 let removed = q.remove(t.short(), false, &questions).unwrap();
3509 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
3510 assert!(removed.released.is_empty() && removed.still_blocked.is_empty());
3511 assert!(q.list().is_empty());
3512 assert!(
3513 q.remove(&t.id, false, &questions).is_err(),
3514 "removing twice is an error"
3515 );
3516 }
3517
3518 #[test]
3519 fn removing_a_task_takes_its_stale_lock_with_it() {
3520 let (dir, q) = queue();
3521 let questions = Questions::at(dir.path().join("questions"));
3522 let mut t = task("interrupted");
3523 q.put(&mut t).unwrap();
3524
3525 let claim = q.claim(&t.id).unwrap();
3528 std::mem::forget(claim);
3529 assert!(
3530 q.claim(&t.id).is_err(),
3531 "the orphaned lock is what makes the task look claimed"
3532 );
3533
3534 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
3536 assert!(err.contains("live daemon"), "{err}");
3537 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
3538
3539 q.remove(&t.id, false, &questions).unwrap();
3541 assert!(q.list().is_empty());
3542 let mut again = task("interrupted");
3543 again.id = t.id.clone();
3544 q.put(&mut again).unwrap();
3545 assert!(
3546 q.claim(&t.id).is_ok(),
3547 "a task that comes back must be claimable, which a left-behind lock would prevent"
3548 );
3549 }
3550
3551 fn notices_of(dir: &Path) -> Vec<crate::notices::Notice> {
3552 crate::notices::Notices::at(dir.join("notifications")).list()
3553 }
3554
3555 #[test]
3556 fn removing_a_sole_dependency_releases_the_dependent_without_a_hold() {
3557 let (dir, q) = queue();
3558 let questions = Questions::at(dir.path().join("questions"));
3559 let mut dep = task("dependency");
3560 q.put(&mut dep).unwrap();
3561 let mut blocked = task("waiting");
3562 blocked.block(vec![dep.id.clone()], Some("waits".to_owned()));
3563 q.put(&mut blocked).unwrap();
3564
3565 let removed = q.remove(&dep.id, false, &questions).unwrap();
3566 assert_eq!(removed.released, [blocked.id.clone()]);
3567 assert!(removed.still_blocked.is_empty());
3568
3569 let after = q.get(&blocked.id).unwrap();
3570 assert_eq!(after.status, TaskStatus::Queued);
3571 assert!(after.blocked_by.is_empty());
3572 assert!(after.block_reason.is_none());
3573 let notes = notices_of(dir.path());
3574 assert_eq!(notes.len(), 1, "{notes:?}");
3575 assert_eq!(notes[0].severity, crate::notices::Severity::Info);
3576 assert_eq!(
3577 notes[0].link,
3578 Some(crate::notices::Link::Task {
3579 id: blocked.id.clone(),
3580 })
3581 );
3582 }
3583
3584 #[test]
3585 fn removing_one_of_two_dependencies_keeps_the_other() {
3586 let (dir, q) = queue();
3587 let questions = Questions::at(dir.path().join("questions"));
3588 let mut dep = task("dependency");
3589 q.put(&mut dep).unwrap();
3590 let mut other = task("other");
3591 q.put(&mut other).unwrap();
3592 let mut blocked = task("waiting");
3593 blocked.block(
3594 vec![dep.id.clone(), other.id.clone()],
3595 Some("waits on both".to_owned()),
3596 );
3597 q.put(&mut blocked).unwrap();
3598
3599 let removed = q.remove(&dep.id, false, &questions).unwrap();
3600 assert!(removed.released.is_empty());
3601 assert_eq!(removed.still_blocked, [blocked.id.clone()]);
3602
3603 let after = q.get(&blocked.id).unwrap();
3604 assert_eq!(after.status, TaskStatus::Blocked);
3605 assert_eq!(after.blocked_by, [other.id.clone()]);
3606 assert!(missing_blockers(&q, &questions, &after.blocked_by).is_empty());
3609 }
3610
3611 #[test]
3612 fn a_dependency_deleted_before_its_dependents_were_rewritten_is_released_later() {
3613 let (dir, q) = queue();
3614 let questions = Questions::at(dir.path().join("questions"));
3615 let mut dep = task("dependency");
3616 q.put(&mut dep).unwrap();
3617 let mut blocked = task("waiting");
3618 blocked.block(vec![dep.id.clone()], None);
3619 q.put(&mut blocked).unwrap();
3620
3621 let claim = q.claim(&blocked.id).unwrap();
3623 let removed = q.remove(&dep.id, false, &questions).unwrap();
3624 assert!(removed.released.is_empty());
3625 drop(claim);
3626
3627 let mut task = q.get(&blocked.id).unwrap();
3628 assert_eq!(task.status, TaskStatus::Blocked);
3629 assert!(missing_blockers(&q, &questions, &task.blocked_by).is_empty());
3630 assert_eq!(q.apply_deleted_blockers(&mut task), [dep.id.clone()]);
3631 assert_eq!(task.status, TaskStatus::Queued);
3632 }
3633
3634 #[test]
3635 fn a_dependency_without_a_tombstone_is_still_missing() {
3636 let (dir, q) = queue();
3637 let questions = Questions::at(dir.path().join("questions"));
3638 let mut dep = task("dependency");
3639 q.put(&mut dep).unwrap();
3640 std::fs::remove_file(q.path_of(&dep.id)).unwrap();
3641 assert_eq!(
3642 missing_blockers(&q, &questions, std::slice::from_ref(&dep.id)),
3643 [dep.id.clone()]
3644 );
3645 assert!(deleted_blockers(&q, std::slice::from_ref(&dep.id)).is_empty());
3646 }
3647
3648 #[test]
3649 fn a_held_dependent_is_left_alone_by_a_dependency_deletion() {
3650 let mut t = task("held");
3651 t.hold_machine(Some("because".to_owned()));
3652 t.blocked_by = vec!["gone".to_owned()];
3653 assert!(!t.dependency_deleted("gone"));
3654 assert_eq!(t.status, TaskStatus::Held);
3655 }
3656
3657 #[test]
3658 fn a_failed_record_removal_takes_the_tombstone_back() {
3659 let (_dir, q) = queue();
3660 let mut t = task("stays");
3661 q.put(&mut t).unwrap();
3662 q.write_tombstone(&t.id).unwrap();
3663 let err = q.remove_record_with_attachments(&t.id, |_| Err(std::io::Error::other("nope")));
3664 assert!(err.is_err());
3665 assert!(!q.was_deleted(&t.id));
3668 }
3669
3670 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3671 let p = dir.join(name);
3672 std::fs::write(&p, body).unwrap();
3673 p
3674 }
3675
3676 #[test]
3677 fn an_attachment_copy_survives_deleting_its_source() {
3678 let (dir, q) = queue();
3679 let src = source_file(dir.path(), "shot.png", "pixels");
3680 let mut t = task("with a picture");
3681 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3682 q.put(&mut t).unwrap();
3683 std::fs::remove_file(&src).unwrap();
3684 assert_eq!(names, ["shot.png"]);
3685 let loaded = q.get(&t.id).unwrap();
3686 let paths = q.attachment_paths(&loaded);
3687 assert_eq!(paths.len(), 1);
3688 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3689 }
3690
3691 #[test]
3692 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3693 let (dir, q) = queue();
3694 let mut t = task("bad names");
3695 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3696 let src = source_file(dir.path(), name, "x");
3697 assert!(
3698 q.attach(&mut t, &[src]).is_err(),
3699 "`{name}` must be refused"
3700 );
3701 }
3702 assert!(!crate::ask::valid_asset_name("C:foo.png"));
3706 #[cfg(not(windows))]
3707 {
3708 let src = source_file(dir.path(), "C:foo.png", "x");
3709 assert!(q.attach(&mut t, &[src]).is_err());
3710 }
3711 let long = format!("{}.png", "a".repeat(70));
3712 let src = source_file(dir.path(), &long, "x");
3713 assert!(q.attach(&mut t, &[src]).is_err());
3714 assert!(t.attachments.is_empty());
3715 assert!(!q.attachments_dir(&t.id).exists());
3716 }
3717
3718 #[test]
3719 fn a_taken_attachment_name_is_numbered_not_overwritten() {
3720 let (dir, q) = queue();
3721 let a = source_file(dir.path(), "shot.png", "one");
3722 let sub = dir.path().join("other");
3723 std::fs::create_dir_all(&sub).unwrap();
3724 let b = source_file(&sub, "shot.png", "two");
3725 let mut t = task("collision");
3726 q.attach(&mut t, &[a]).unwrap();
3727 q.attach(&mut t, &[b]).unwrap();
3728 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3729 let paths = q.attachment_paths(&t);
3730 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3731 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3732 }
3733
3734 #[test]
3735 fn a_renumbered_name_stays_inside_the_length_bound() {
3736 let (dir, q) = queue();
3737 let name = format!("{}.png", "a".repeat(60));
3738 assert_eq!(name.len(), 64);
3739 let a = source_file(dir.path(), &name, "one");
3740 let sub = dir.path().join("other");
3741 std::fs::create_dir_all(&sub).unwrap();
3742 let b = source_file(&sub, &name, "two");
3743 let mut t = task("long");
3744 q.attach(&mut t, &[a, b]).unwrap();
3745 assert_eq!(t.attachments.len(), 2);
3746 assert!(
3747 t.attachments
3748 .iter()
3749 .all(|n| crate::ask::valid_asset_name(n))
3750 );
3751 assert!(t.attachments[1].ends_with("-2.png"));
3752 }
3753
3754 #[test]
3755 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3756 let (dir, q) = queue();
3757 let good = source_file(dir.path(), "good.png", "ok");
3758 let mut t = task("partial");
3759 q.attach(&mut t, &[good]).unwrap();
3760 let more = source_file(dir.path(), "more.png", "ok");
3761 let missing = dir.path().join("missing.png");
3762 assert!(q.attach(&mut t, &[more, missing]).is_err());
3763 assert_eq!(t.attachments, ["good.png"]);
3764 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3765 .unwrap()
3766 .flatten()
3767 .collect();
3768 assert_eq!(on_disk.len(), 1);
3769 }
3770
3771 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3773 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3774 std::fs::create_dir_all(&tmp).unwrap();
3775 tmp
3776 }
3777
3778 #[test]
3779 fn a_failed_put_leaves_no_new_attachment_directory() {
3780 let (dir, q) = queue();
3781 let mut t = task("fresh");
3782 let tmp = block_put(&q, &t);
3783 let src = source_file(dir.path(), "shot.png", "x");
3784 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3785 assert!(t.attachments.is_empty());
3786 assert!(!q.attachments_dir(&t.id).exists());
3787 assert!(!q.path_of(&t.id).exists());
3788 std::fs::remove_dir(tmp).unwrap();
3789 }
3790
3791 #[test]
3792 fn a_failed_put_removes_only_the_copy_it_just_made() {
3793 let (dir, q) = queue();
3794 let mut t = task("edited");
3795 let first = source_file(dir.path(), "first.png", "1");
3796 q.attach_and_put(&mut t, &[first]).unwrap();
3797 block_put(&q, &t);
3798 let second = source_file(dir.path(), "second.png", "2");
3799 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3800 assert_eq!(t.attachments, ["first.png"]);
3801 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3802 .unwrap()
3803 .flatten()
3804 .map(|e| e.file_name().to_string_lossy().into_owned())
3805 .collect();
3806 assert_eq!(on_disk, ["first.png"]);
3807 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3808 }
3809
3810 #[test]
3811 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3812 let (dir, q) = queue();
3813 let questions = Questions::at(dir.path().join("questions"));
3814 let gone = task("gone");
3815 let mut other = task("other");
3816 let mut live = task("live");
3817 q.put(&mut other).unwrap();
3818 q.put(&mut live).unwrap();
3819 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3822 std::fs::create_dir_all(&orphan).unwrap();
3823 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3824 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3827 std::fs::create_dir_all(&busy).unwrap();
3828
3829 q.remove(&other.id, false, &questions).unwrap();
3830 assert!(!orphan.exists(), "an orphan is swept");
3831 assert!(busy.exists(), "a removal in progress is left alone");
3832 }
3833
3834 #[test]
3835 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3836 let (dir, q) = queue();
3837 let questions = Questions::at(dir.path().join("questions"));
3838 let mut t = task("stuck");
3839 let src = source_file(dir.path(), "shot.png", "x");
3840 q.attach_and_put(&mut t, &[src]).unwrap();
3841 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3844 std::fs::create_dir_all(&aside).unwrap();
3845 std::fs::write(aside.join("old.png"), "o").unwrap();
3846 assert!(q.remove(&t.id, false, &questions).is_err());
3847 assert!(q.path_of(&t.id).exists());
3848 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3849 }
3850
3851 #[test]
3852 fn a_failed_record_removal_puts_the_attachments_back() {
3853 let (dir, q) = queue();
3854 let mut t = task("rollback");
3855 let src = source_file(dir.path(), "shot.png", "x");
3856 q.attach_and_put(&mut t, &[src]).unwrap();
3857 let err = q
3858 .remove_record_with_attachments(&t.id, |_| {
3859 Err(std::io::Error::other("injected failure"))
3860 })
3861 .unwrap_err();
3862 assert!(format!("{err:#}").contains("injected failure"));
3863 assert!(q.path_of(&t.id).exists());
3864 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3865 assert!(
3866 !q.root
3867 .join(format!("{}.attachments.removing", t.id))
3868 .exists()
3869 );
3870 }
3871
3872 #[test]
3873 fn editing_a_task_keeps_its_attachments() {
3874 let (dir, q) = queue();
3875 let src = source_file(dir.path(), "shot.png", "x");
3876 let mut t = task("editable");
3877 q.attach(&mut t, &[src]).unwrap();
3878 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3879 q.put(&mut t).unwrap();
3880 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3881 }
3882
3883 #[test]
3884 fn removing_a_task_deletes_its_attachments() {
3885 let (dir, q) = queue();
3886 let questions = Questions::at(dir.path().join("questions"));
3887 let src = source_file(dir.path(), "shot.png", "x");
3888 let mut t = task("doomed");
3889 q.attach(&mut t, &[src]).unwrap();
3890 q.put(&mut t).unwrap();
3891 assert!(q.attachments_dir(&t.id).is_dir());
3892 q.remove(&t.id, false, &questions).unwrap();
3893 assert!(!q.attachments_dir(&t.id).exists());
3894 assert!(q.list().is_empty());
3895 }
3896
3897 #[test]
3898 fn a_task_written_before_attachments_still_reads() {
3899 let (_dir, q) = queue();
3900 let mut t = task("old");
3901 q.put(&mut t).unwrap();
3902 let path = q.path_of(&t.id);
3903 let mut v: serde_json::Value =
3904 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3905 v.as_object_mut().unwrap().remove("attachments");
3906 std::fs::write(&path, v.to_string()).unwrap();
3907 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3908 }
3909
3910 #[test]
3911 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3912 let q = Queue::at(PathBuf::from("relative-queue"));
3913 let mut t = task("rel");
3914 t.attachments.push("shot.png".to_owned());
3915 let paths = q.attachment_paths(&t);
3916 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3917 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3918 }
3919
3920 #[test]
3921 fn link_run_adds_a_run_once_and_touches_nothing_else() {
3922 let dir = tempfile::tempdir().unwrap();
3923 let queue = Queue::at(dir.path().join("queue"));
3924 let mut t = Task::new(
3925 "t".to_owned(),
3926 "do it".to_owned(),
3927 PathBuf::from("."),
3928 Source::Human,
3929 );
3930 queue.put(&mut t).unwrap();
3931 let before = queue.get(&t.id).unwrap();
3932
3933 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3934 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3935 assert_eq!(linked.status, before.status);
3936 assert_eq!(linked.attempts, before.attempts);
3937 assert_eq!(linked.interrupt, before.interrupt);
3938
3939 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3940 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3941 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3942 assert!(queue.link_run("no-such-task", "r").is_err());
3943 }
3944
3945 #[test]
3946 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
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 mut snapshot = queue.get(&t.id).unwrap();
3958 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3959
3960 snapshot.start("20260930-000000-aaaa".to_owned());
3961 queue.put(&mut snapshot).unwrap();
3962
3963 let stored = queue.get(&t.id).unwrap();
3964 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3965 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3966 assert_eq!(stored.attempts, 1);
3967 }
3968
3969 #[test]
3970 fn concurrent_links_and_daemon_saves_lose_nothing() {
3971 let dir = tempfile::tempdir().unwrap();
3972 let queue = Queue::at(dir.path().join("queue"));
3973 let mut t = Task::new(
3974 "t".to_owned(),
3975 "do it".to_owned(),
3976 PathBuf::from("."),
3977 Source::Human,
3978 );
3979 queue.put(&mut t).unwrap();
3980 let id = t.id.clone();
3981
3982 let linkers: Vec<_> = (0..4)
3983 .map(|n| {
3984 let (queue, id) = (queue.clone(), id.clone());
3985 std::thread::spawn(move || {
3986 for k in 0..10 {
3987 queue
3988 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3989 .unwrap();
3990 }
3991 })
3992 })
3993 .collect();
3994 let mut mine = queue.get(&id).unwrap();
3997 for k in 0..10 {
3998 mine.start(format!("20260930-000009-d{k:03}"));
3999 queue.put(&mut mine).unwrap();
4000 }
4001 for l in linkers {
4002 l.join().unwrap();
4003 }
4004
4005 let stored = queue.get(&id).unwrap();
4006 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
4007 assert_eq!(
4008 stored.attempts, 10,
4009 "linking never rewinds the daemon's work"
4010 );
4011 }
4012
4013 fn age_lock(path: &Path) {
4014 let f = std::fs::OpenOptions::new().write(true).open(path).unwrap();
4015 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
4016 .unwrap();
4017 }
4018
4019 #[test]
4020 fn concurrent_stale_takeover_yields_one_holder() {
4021 use std::sync::atomic::{AtomicUsize, Ordering};
4022 use std::sync::{Arc, Barrier};
4023 for _ in 0..5 {
4024 let dir = tempfile::tempdir().unwrap();
4025 let q = Queue::at(dir.path().to_path_buf());
4026 let lock = dir.path().join("t.write-lock");
4027 std::fs::write(&lock, "dead-0000").unwrap();
4028 age_lock(&lock);
4029 let n = 6;
4030 let barrier = Arc::new(Barrier::new(n));
4031 let (now, max) = (Arc::new(AtomicUsize::new(0)), Arc::new(AtomicUsize::new(0)));
4032 let handles: Vec<_> = (0..n)
4033 .map(|_| {
4034 let (q, b, now, max) = (q.clone(), barrier.clone(), now.clone(), max.clone());
4035 std::thread::spawn(move || {
4036 b.wait();
4037 let g = q.lock_task("t").unwrap();
4038 let held = now.fetch_add(1, Ordering::SeqCst) + 1;
4039 max.fetch_max(held, Ordering::SeqCst);
4040 std::thread::sleep(std::time::Duration::from_millis(20));
4041 now.fetch_sub(1, Ordering::SeqCst);
4042 drop(g);
4043 })
4044 })
4045 .collect();
4046 for h in handles {
4047 h.join().unwrap();
4048 }
4049 assert_eq!(max.load(Ordering::SeqCst), 1);
4050 assert!(!lock.exists());
4051 }
4052 }
4053
4054 #[test]
4055 fn dropping_a_stolen_lock_leaves_the_new_holders_lock() {
4056 let dir = tempfile::tempdir().unwrap();
4057 let q = Queue::at(dir.path().to_path_buf());
4058 let lock = dir.path().join("t.write-lock");
4059 let a = q.lock_task("t").unwrap();
4060 age_lock(&lock);
4061 let b = q.lock_task("t").unwrap();
4062 assert_ne!(a.token, b.token);
4063 drop(a);
4064 assert_eq!(std::fs::read_to_string(&lock).unwrap(), b.token);
4065 drop(b);
4066 assert!(!lock.exists());
4067 }
4068
4069 #[test]
4070 fn break_stale_leaves_a_lock_that_replaced_the_one_judged() {
4071 let dir = tempfile::tempdir().unwrap();
4072 let lock = dir.path().join("t.write-lock");
4073 std::fs::write(&lock, "new-token").unwrap();
4074 assert!(!break_stale(&lock, "old-token"));
4075 assert_eq!(std::fs::read_to_string(&lock).unwrap(), "new-token");
4076 assert!(!break_stale(&lock, "new-token"));
4078 assert!(lock.exists());
4079 age_lock(&lock);
4080 assert!(break_stale(&lock, "new-token"));
4081 assert!(!lock.exists());
4082 }
4083
4084 #[test]
4085 fn a_stale_break_marker_is_recovered_and_a_fresh_one_is_respected() {
4086 let dir = tempfile::tempdir().unwrap();
4087 let lock = dir.path().join("t.write-lock");
4088 std::fs::write(&lock, "dead-1").unwrap();
4089 age_lock(&lock);
4090 let marker = break_marker(&lock, "dead-1");
4091 std::fs::write(&marker, "crashed-remover").unwrap();
4092 assert!(!break_stale(&lock, "dead-1"));
4094 assert!(lock.exists());
4095 age_lock(&marker);
4097 assert!(break_stale(&lock, "dead-1"));
4098 assert!(!lock.exists());
4099 assert!(!marker.exists());
4100 }
4101}