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