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 = 10;
101
102#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
104#[serde(rename_all = "lowercase")]
105pub enum HoldSource {
106 Manual,
108 Machine,
110}
111
112impl HoldSource {
113 pub fn label(self) -> &'static str {
115 match self {
116 Self::Manual => "manual",
117 Self::Machine => "machine",
118 }
119 }
120}
121
122#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
125#[serde(tag = "kind", rename_all = "lowercase")]
126pub enum Source {
127 Human,
129 Agent {
132 run: String,
134 node: String,
136 },
137 Issue {
139 number: u64,
141 repo: String,
143 },
144}
145
146pub const CHAT_NODE: &str = "chat";
149
150impl Source {
151 pub fn label(&self) -> String {
153 match self {
154 Self::Human => "human".to_owned(),
155 Self::Agent { run, node } => format!("{node}@{}", short(run)),
156 Self::Issue { number, .. } => format!("issue #{number}"),
157 }
158 }
159}
160
161#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
163#[serde(rename_all = "lowercase")]
164pub enum TaskStatus {
165 Queued,
167 Running,
169 Done,
171 Failed,
173 Held,
175 Blocked,
179}
180
181impl TaskStatus {
182 pub fn runnable(self) -> bool {
184 matches!(self, Self::Queued | Self::Failed)
185 }
186
187 pub fn as_str(self) -> &'static str {
189 match self {
190 Self::Queued => "queued",
191 Self::Running => "running",
192 Self::Done => "done",
193 Self::Failed => "failed",
194 Self::Held => "held",
195 Self::Blocked => "blocked",
196 }
197 }
198}
199
200#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
203pub struct TaskCounts {
204 pub queued: usize,
206 pub running: usize,
208 pub done: usize,
210 pub failed: usize,
212 pub held: usize,
214 pub blocked: usize,
216}
217
218impl TaskCounts {
219 pub fn of(tasks: &[Task]) -> Self {
223 let mut counts = Self::default();
224 for t in tasks {
225 match t.status {
226 TaskStatus::Queued => counts.queued += 1,
227 TaskStatus::Running => counts.running += 1,
228 TaskStatus::Done => counts.done += 1,
229 TaskStatus::Failed => counts.failed += 1,
230 TaskStatus::Held => counts.held += 1,
231 TaskStatus::Blocked => counts.blocked += 1,
232 }
233 }
234 counts
235 }
236}
237
238#[derive(Debug, Clone, Serialize, Deserialize)]
240#[serde(deny_unknown_fields)]
241pub struct Task {
242 pub schema: u32,
244 pub id: String,
246 pub title: String,
248 pub instruction: String,
250 pub repo: PathBuf,
252 pub source: Source,
254 #[serde(default)]
256 pub priority: i32,
257 #[serde(default)]
268 pub solo: bool,
269 pub status: TaskStatus,
271 #[serde(default)]
273 pub attempts: usize,
274 #[serde(default)]
276 pub runs: Vec<String>,
277 #[serde(default)]
279 pub last_error: Option<String>,
280 #[serde(default)]
293 pub hold_reason: Option<String>,
294 #[serde(default)]
297 pub hold_source: Option<HoldSource>,
298 #[serde(default)]
310 pub diagnostic: Option<String>,
311 #[serde(default)]
321 pub blocked_by: Vec<String>,
322 #[serde(default)]
325 pub block_reason: Option<String>,
326 #[serde(default)]
339 pub blocked_from: Option<TaskStatus>,
340 #[serde(default)]
351 pub answers: Vec<AnsweredQuestion>,
352 #[serde(default)]
358 pub triage_applied: Vec<String>,
359 #[serde(default)]
365 pub actions_applied: Vec<String>,
366 #[serde(default)]
371 pub resume_override: Option<OperatorResume>,
372 #[serde(default)]
384 pub review_branch: Option<String>,
385 #[serde(default)]
388 pub fresh_start: bool,
389 #[serde(default)]
400 pub interrupt: bool,
401 #[serde(default)]
424 pub urgent: bool,
425 #[serde(default)]
431 pub attachments: Vec<String>,
432 #[serde(default)]
436 pub followup: Option<FollowUp>,
437 #[serde(default)]
443 pub overrides: Option<RunOverrides>,
444 #[serde(default)]
451 pub review_of: Option<String>,
452 #[serde(default)]
457 pub held_at: Option<Timestamp>,
458 pub created_at: Timestamp,
460 pub updated_at: Timestamp,
462}
463
464#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
466pub struct RunOverrides {
467 #[serde(default)]
469 pub merge: Option<String>,
470 #[serde(default)]
472 pub candidates: Option<usize>,
473 #[serde(default)]
475 pub judges: Option<usize>,
476 #[serde(default)]
478 pub reviewers: Option<usize>,
479 #[serde(default)]
481 pub review_rounds: Option<usize>,
482 #[serde(default)]
484 pub seed: Option<u64>,
485 #[serde(default)]
487 pub config: Option<PathBuf>,
488}
489
490impl RunOverrides {
491 pub fn merge_over(&mut self, other: &RunOverrides) {
493 macro_rules! take {
494 ($($f:ident),*) => {$(
495 if other.$f.is_some() {
496 self.$f.clone_from(&other.$f);
497 }
498 )*};
499 }
500 take!(
501 merge,
502 candidates,
503 judges,
504 reviewers,
505 review_rounds,
506 seed,
507 config
508 );
509 }
510
511 pub fn apply(&self, config: &mut crate::config::Config) {
513 if let Some(n) = self.candidates {
514 config.graph.implementers = n;
515 }
516 if let Some(n) = self.judges {
517 config.graph.judges = n;
518 }
519 if let Some(n) = self.reviewers {
520 config.graph.reviewers = n;
521 }
522 if let Some(n) = self.review_rounds {
523 config.graph.review_rounds = n;
524 }
525 if let Some(m) = &self.merge
526 && let Ok(mode) = crate::daemon::merge_mode(m)
527 {
528 config.merge.mode = mode;
529 }
530 if let Some(s) = self.seed {
531 config.blind.seed = Some(s);
532 }
533 }
534}
535
536#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
538pub struct FollowUp {
539 pub run: String,
541 #[serde(default)]
543 pub origin_task: Option<String>,
544 pub pr: String,
546 pub findings: Vec<String>,
548 pub generation: u32,
551}
552
553#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
556pub struct OperatorResume {
557 pub question_id: String,
559 pub at: Timestamp,
561 #[serde(default)]
564 pub conductor_rehold: Option<String>,
565 #[serde(default)]
568 pub forced: bool,
569 #[serde(default)]
576 pub pinned_run: Option<String>,
577}
578
579#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
582pub struct AnsweredQuestion {
583 pub question: String,
585 pub answer: String,
587}
588
589impl Task {
590 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
592 let now = Timestamp::now();
593 Self {
594 schema: SCHEMA,
595 id: new_id(),
596 title,
597 instruction,
598 repo,
599 source,
600 priority: 0,
601 solo: false,
602 status: TaskStatus::Queued,
603 attempts: 0,
604 runs: Vec::new(),
605 last_error: None,
606 hold_reason: None,
607 hold_source: None,
608 diagnostic: None,
609 blocked_by: Vec::new(),
610 block_reason: None,
611 blocked_from: None,
612 answers: Vec::new(),
613 triage_applied: Vec::new(),
614 actions_applied: Vec::new(),
615 resume_override: None,
616 review_branch: None,
617 fresh_start: false,
618 interrupt: false,
619 urgent: false,
620 attachments: Vec::new(),
621 followup: None,
622 overrides: None,
623 review_of: None,
624 held_at: None,
625 created_at: now,
626 updated_at: now,
627 }
628 }
629
630 pub fn short(&self) -> &str {
632 short(&self.id)
633 }
634
635 pub fn mark_triage_applied(&mut self, question_id: &str) {
637 if !self.triage_applied(question_id) {
638 self.triage_applied.push(question_id.to_owned());
639 }
640 }
641
642 pub fn action_applied(&self, question_id: &str) -> bool {
645 self.actions_applied.iter().any(|id| id == question_id)
646 }
647
648 pub fn mark_action_applied(&mut self, question_id: &str) {
650 if !self.action_applied(question_id) {
651 self.actions_applied.push(question_id.to_owned());
652 }
653 }
654
655 pub fn triage_applied(&self, question_id: &str) -> bool {
657 self.triage_applied.iter().any(|id| id == question_id)
658 }
659
660 pub fn start(&mut self, run: String) {
671 self.status = TaskStatus::Running;
672 self.attempts += 1;
673 self.runs.push(run);
674 self.last_error = None;
675 self.fresh_start = false;
676 self.interrupt = false;
677 self.resume_override = None;
679 }
680
681 pub fn link_run(&mut self, run: &str) -> bool {
686 if self.runs.iter().any(|r| r == run) {
687 return false;
688 }
689 self.runs.push(run.to_owned());
690 true
691 }
692
693 pub fn succeed(&mut self) {
703 self.status = TaskStatus::Done;
704 self.resume_override = None;
705 self.last_error = None;
706 self.hold_reason = None;
707 self.hold_source = None;
708 self.diagnostic = None;
709 self.blocked_by.clear();
710 self.block_reason = None;
711 self.blocked_from = None;
712 }
713
714 pub fn already_landed(&mut self, note: impl Into<String>) {
722 self.succeed();
723 self.attempts = self.attempts.saturating_sub(1);
724 self.last_error = Some(note.into());
725 }
726
727 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
757 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
758 return &[];
759 }
760 &self.runs[..self.runs.len() - 1]
761 }
762
763 pub fn successor_of(&self, run: &str) -> Option<&String> {
770 let pos = self.runs.iter().position(|r| r == run)?;
771 self.runs.get(pos + 1)
772 }
773
774 pub fn earlier_attempts(&self) -> &[String] {
781 &self.runs
782 }
783
784 fn note_held(&mut self) {
787 if self.status != TaskStatus::Held || self.held_at.is_none() {
788 self.held_at = Some(Timestamp::now());
789 }
790 }
791
792 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
808 let why = why.into();
809 self.diagnostic = None;
810 self.status = if self.attempts >= max_attempts {
811 self.note_held();
812 self.hold_source = Some(HoldSource::Machine);
813 self.hold_reason = Some(why.clone());
814 TaskStatus::Held
815 } else {
816 TaskStatus::Failed
817 };
818 self.last_error = Some(why);
819 }
820
821 pub fn stall(&mut self, why: impl Into<String>) {
831 self.last_error = Some(why.into());
832 self.diagnostic = None;
833 self.attempts = self.attempts.saturating_sub(1);
834 self.status = TaskStatus::Failed;
835 }
836
837 pub fn operator_held(&self) -> bool {
843 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
844 }
845
846 pub fn hold_manual(&mut self, reason: Option<String>) {
856 self.note_held();
857 self.status = TaskStatus::Held;
858 if reason.is_some() {
859 self.hold_reason = reason;
860 }
861 self.hold_source = Some(HoldSource::Manual);
862 self.blocked_by.clear();
863 self.block_reason = None;
864 self.blocked_from = None;
865 }
866
867 pub fn hold_machine(&mut self, reason: Option<String>) {
872 self.note_held();
873 self.status = TaskStatus::Held;
874 if reason.is_some() {
875 self.hold_reason = reason;
876 }
877 self.hold_source = Some(HoldSource::Machine);
878 self.blocked_by.clear();
879 self.block_reason = None;
880 self.blocked_from = None;
881 }
882
883 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
893 if self.status != TaskStatus::Blocked {
894 self.blocked_from = Some(self.status);
895 }
896 self.status = TaskStatus::Blocked;
897 self.blocked_by = blocked_by;
898 self.block_reason = reason;
899 }
900
901 pub fn unblock(&mut self, resolved_id: &str) {
924 if self.status != TaskStatus::Blocked {
925 return;
926 }
927 self.blocked_by.retain(|id| id != resolved_id);
928 self.restore_if_unblocked();
929 }
930
931 pub fn dependency_deleted(&mut self, deleted_id: &str) -> bool {
942 if self.status != TaskStatus::Blocked || !self.blocked_by.iter().any(|b| b == deleted_id) {
943 return false;
944 }
945 self.blocked_by.retain(|id| id != deleted_id);
946 self.restore_if_unblocked();
947 true
948 }
949
950 fn restore_if_unblocked(&mut self) {
954 if self.blocked_by.is_empty() {
955 self.status = match self.blocked_from {
956 Some(TaskStatus::Running) => TaskStatus::Queued,
957 Some(other) => other,
958 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
959 TaskStatus::Held
960 }
961 None => TaskStatus::Queued,
962 };
963 self.block_reason = None;
964 self.blocked_from = None;
965 }
966 }
967
968 pub fn record_answer(&mut self, question: String, answer: String) {
973 self.answers.push(AnsweredQuestion { question, answer });
974 }
975
976 pub fn request_review(&mut self, branch: String) {
980 self.release();
981 self.review_branch = Some(branch);
982 }
983
984 pub fn requeue(&mut self) {
987 self.release();
988 self.review_branch = None;
990 self.fresh_start = true;
991 }
992
993 pub fn hold_for_handover(&mut self, branch: Option<String>, reason: String) {
1002 if branch.is_some() {
1003 self.review_branch = branch;
1004 }
1005 self.hold_machine(Some(reason));
1006 }
1007
1008 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
1017 if self.status == TaskStatus::Running {
1018 bail!(
1019 "task {} is running; its priority cannot be changed until \
1020 this attempt finishes",
1021 self.short()
1022 );
1023 }
1024 self.priority = priority;
1025 Ok(())
1026 }
1027
1028 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
1043 if interrupt && !self.status.runnable() {
1044 bail!(
1045 "task {} is {}; only a queued or failed task can be marked \
1046 to interrupt",
1047 self.short(),
1048 self.status.as_str()
1049 );
1050 }
1051 self.interrupt = interrupt;
1052 Ok(())
1053 }
1054
1055 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
1067 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
1068 bail!(
1069 "task {} is {}; only a queued or held task's instruction can \
1070 be edited",
1071 self.short(),
1072 self.status.as_str()
1073 );
1074 }
1075 self.title = title;
1076 self.instruction = instruction;
1077 Ok(())
1078 }
1079
1080 pub fn handed_off(&mut self, why: impl Into<String>) {
1100 let why = why.into();
1101 self.diagnostic = None;
1102 self.note_held();
1103 self.status = TaskStatus::Held;
1104 self.hold_source = Some(HoldSource::Machine);
1105 self.hold_reason = Some(why.clone());
1106 self.last_error = Some(why);
1107 }
1108
1109 pub fn release(&mut self) {
1113 let refused_handover = self.status == TaskStatus::Held
1114 && self.hold_source == Some(HoldSource::Machine)
1115 && self.review_branch.is_some();
1116 self.status = TaskStatus::Queued;
1117 self.held_at = None;
1118 self.attempts = 0;
1119 self.last_error = None;
1120 self.hold_reason = None;
1123 self.hold_source = None;
1124 self.diagnostic = None;
1125 self.blocked_by.clear();
1130 self.block_reason = None;
1131 self.blocked_from = None;
1132 if !refused_handover {
1136 self.review_branch = None;
1137 }
1138 self.fresh_start = false;
1139 }
1140}
1141
1142fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
1146 use std::io::ErrorKind;
1147 let (stem, ext) = match name.rfind('.') {
1148 Some(i) if i > 0 => (&name[..i], &name[i..]),
1149 _ => (name, ""),
1150 };
1151 for n in 1u32.. {
1152 let candidate = if n == 1 {
1153 name.to_owned()
1154 } else {
1155 let suffix = format!("-{n}");
1156 let room = 64usize.saturating_sub(suffix.len() + ext.len());
1157 let stem: String = stem.chars().take(room).collect();
1158 format!("{stem}{suffix}{ext}")
1159 };
1160 if !crate::ask::valid_asset_name(&candidate) {
1161 bail!("no valid attachment name is left for `{name}`");
1162 }
1163 let path = dir.join(&candidate);
1164 match std::fs::OpenOptions::new()
1165 .write(true)
1166 .create_new(true)
1167 .open(&path)
1168 {
1169 Ok(mut out) => {
1170 let copied = std::fs::File::open(src)
1171 .and_then(|mut input| std::io::copy(&mut input, &mut out));
1172 if let Err(e) = copied {
1173 drop(out);
1174 let _ = std::fs::remove_file(&path);
1175 return Err(e).with_context(|| format!("copy {}", src.display()));
1176 }
1177 return Ok((candidate, path));
1178 }
1179 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
1180 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
1181 }
1182 }
1183 unreachable!("the counter never runs out")
1184}
1185
1186const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
1189
1190struct TaskLock {
1193 path: PathBuf,
1194 token: String,
1195}
1196
1197fn owner_token() -> String {
1200 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1201 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1202 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy() ^ n.rotate_left(32));
1203 format!("{}-{:016x}", std::process::id(), r.next_u64())
1204}
1205
1206fn break_marker(path: &Path, token: &str) -> PathBuf {
1209 let mut h: u64 = 0xcbf29ce484222325;
1210 for b in token.bytes() {
1211 h ^= u64::from(b);
1212 h = h.wrapping_mul(0x100000001b3);
1213 }
1214 let mut name = path.as_os_str().to_owned();
1215 name.push(format!(".break-{h:016x}"));
1216 PathBuf::from(name)
1217}
1218
1219fn older_than_stale(path: &Path) -> bool {
1220 std::fs::metadata(path)
1221 .and_then(|m| m.modified())
1222 .ok()
1223 .and_then(|t| t.elapsed().ok())
1224 .is_some_and(|age| age > TASK_LOCK_STALE)
1225}
1226
1227const MAX_MARKER_DEPTH: u8 = 3;
1231
1232fn take_marker(path: &Path, token: &str, depth: u8) -> Option<(PathBuf, String)> {
1238 use std::io::Write;
1239 let marker = break_marker(path, token);
1240 for _ in 0..2 {
1241 match std::fs::OpenOptions::new()
1242 .write(true)
1243 .create_new(true)
1244 .open(&marker)
1245 {
1246 Ok(mut f) => {
1247 let mine = owner_token();
1248 if f.write_all(mine.as_bytes()).is_err() {
1249 drop(f);
1250 let _ = std::fs::remove_file(&marker);
1251 return None;
1252 }
1253 return Some((marker, mine));
1254 }
1255 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1256 let judged = std::fs::read_to_string(&marker).ok();
1257 match judged.filter(|_| older_than_stale(&marker)) {
1258 Some(judged) if depth < MAX_MARKER_DEPTH => {
1259 remove_lock_if(&marker, &judged, true, depth + 1)?;
1260 }
1261 _ => return None,
1262 }
1263 }
1264 Err(_) => return None,
1265 }
1266 }
1267 None
1268}
1269
1270fn remove_lock_if(path: &Path, token: &str, require_stale: bool, depth: u8) -> Option<bool> {
1275 let (marker, mine) = take_marker(path, token, depth)?;
1276 let still = std::fs::read_to_string(path).is_ok_and(|c| c == token)
1277 && (!require_stale || older_than_stale(path));
1278 if still {
1279 let _ = std::fs::remove_file(path);
1280 }
1281 if depth >= MAX_MARKER_DEPTH {
1283 if std::fs::read_to_string(&marker).is_ok_and(|c| c == mine) {
1284 let _ = std::fs::remove_file(&marker);
1285 }
1286 } else {
1287 let _ = remove_lock_if(&marker, &mine, false, depth + 1);
1288 }
1289 Some(still)
1290}
1291
1292fn break_stale(path: &Path, judged: &str) -> bool {
1295 remove_lock_if(path, judged, true, 0).unwrap_or(false)
1296}
1297
1298impl Drop for TaskLock {
1299 fn drop(&mut self) {
1300 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
1301 loop {
1302 match remove_lock_if(&self.path, &self.token, false, 0) {
1303 Some(_) => return,
1304 None if std::time::Instant::now() > deadline => return,
1305 None => std::thread::sleep(std::time::Duration::from_millis(5)),
1306 }
1307 }
1308 }
1309}
1310
1311#[derive(Debug, Clone)]
1313pub struct Queue {
1314 root: PathBuf,
1315}
1316
1317impl Queue {
1318 pub fn open() -> Self {
1320 Self::at(crate::run::home().join("queue"))
1321 }
1322
1323 pub fn at(root: PathBuf) -> Self {
1326 Self { root }
1327 }
1328
1329 pub fn root(&self) -> &Path {
1331 &self.root
1332 }
1333
1334 pub fn path_of(&self, id: &str) -> PathBuf {
1336 self.root.join(format!("{id}.json"))
1337 }
1338
1339 pub fn attachments_dir(&self, id: &str) -> PathBuf {
1342 self.root.join(format!("{id}.attachments"))
1343 }
1344
1345 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1356 let mut wanted = Vec::new();
1357 for src in sources {
1358 let name = src
1359 .file_name()
1360 .and_then(|n| n.to_str())
1361 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1362 if !crate::ask::valid_asset_name(name) {
1363 bail!(
1364 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1365 with no `..`; rename the file and try again"
1366 );
1367 }
1368 if !src.is_file() {
1369 bail!("attachment `{}` is not a file", src.display());
1370 }
1371 wanted.push((src, name));
1372 }
1373 let dir = self.attachments_dir(&task.id);
1374 let existed = dir.is_dir();
1375 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1376 let mut created: Vec<PathBuf> = Vec::new();
1377 let mut names = Vec::new();
1378 let mut copy_all = || -> Result<()> {
1379 for (src, name) in &wanted {
1380 let (stored, path) = copy_new(&dir, src, name)?;
1381 created.push(path);
1382 names.push(stored);
1383 }
1384 Ok(())
1385 };
1386 if let Err(e) = copy_all() {
1387 for path in &created {
1388 let _ = std::fs::remove_file(path);
1389 }
1390 if !existed {
1391 let _ = std::fs::remove_dir(&dir);
1392 }
1393 return Err(e);
1394 }
1395 task.attachments.extend(names.iter().cloned());
1396 Ok(names)
1397 }
1398
1399 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1410 let dir = self.attachments_dir(&task.id);
1411 let existed = dir.is_dir();
1412 let before = task.attachments.len();
1413 let names = self.attach(task, sources)?;
1414 if let Err(e) = self.put(task) {
1415 for name in &names {
1416 let _ = std::fs::remove_file(dir.join(name));
1417 }
1418 if !existed {
1419 let _ = std::fs::remove_dir(&dir);
1420 }
1421 task.attachments.truncate(before);
1422 return Err(e);
1423 }
1424 Ok(names)
1425 }
1426
1427 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1431 let dir = self.attachments_dir(&task.id);
1432 task.attachments
1433 .iter()
1434 .map(|n| {
1435 let p = dir.join(n);
1436 std::path::absolute(&p).unwrap_or(p)
1437 })
1438 .collect()
1439 }
1440
1441 pub fn put(&self, task: &mut Task) -> Result<()> {
1450 let _lock = self.lock_task(&task.id)?;
1451 self.put_unlocked(task)
1452 }
1453
1454 pub fn create_new(&self, task: &mut Task) -> Result<bool> {
1459 let _lock = self.lock_task(&task.id)?;
1460 if self.path_of(&task.id).exists() {
1461 return Ok(false);
1462 }
1463 self.put_unlocked(task)?;
1464 Ok(true)
1465 }
1466
1467 fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1469 if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1470 for run in stored.runs {
1471 if !task.runs.contains(&run) {
1472 task.runs.push(run);
1473 }
1474 }
1475 }
1476 task.updated_at = Timestamp::now();
1477 std::fs::create_dir_all(&self.root)
1478 .with_context(|| format!("create {}", self.root.display()))?;
1479 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1480 let path = self.path_of(&task.id);
1481 let tmp = path.with_extension("json.tmp");
1482 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1483 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1484 if let (Some(notice), Some(home)) = (
1489 crate::notices::task_held(task),
1490 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1491 ) {
1492 crate::notices::raise_in(home, notice);
1493 }
1494 Ok(())
1495 }
1496
1497 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1502 let id = self.resolve_id(id)?;
1503 let _lock = self.lock_task(&id)?;
1506 let mut task = self.get(&id)?;
1507 if task.link_run(run) {
1508 self.put_unlocked(&mut task)?;
1509 }
1510 Ok(task)
1511 }
1512
1513 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1521 std::fs::create_dir_all(&self.root)
1522 .with_context(|| format!("create {}", self.root.display()))?;
1523 let path = self.root.join(format!("{id}.write-lock"));
1524 let started = std::time::Instant::now();
1525 loop {
1526 match std::fs::OpenOptions::new()
1527 .write(true)
1528 .create_new(true)
1529 .open(&path)
1530 {
1531 Ok(mut file) => {
1532 use std::io::Write;
1533 let token = owner_token();
1534 if let Err(e) = file.write_all(token.as_bytes()) {
1535 drop(file);
1536 let _ = std::fs::remove_file(&path);
1537 return Err(e).with_context(|| format!("lock {}", path.display()));
1538 }
1539 return Ok(TaskLock { path, token });
1540 }
1541 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1542 let judged = std::fs::read_to_string(&path).ok();
1545 let broken = judged
1546 .filter(|_| older_than_stale(&path))
1547 .is_some_and(|judged| break_stale(&path, &judged));
1548 if !broken {
1552 if started.elapsed() > TASK_LOCK_STALE {
1553 bail!("could not lock task {id}");
1554 }
1555 std::thread::sleep(std::time::Duration::from_millis(15));
1556 }
1557 }
1558 Err(e)
1562 if e.kind() == std::io::ErrorKind::PermissionDenied
1563 && started.elapsed() <= TASK_LOCK_STALE =>
1564 {
1565 std::thread::sleep(std::time::Duration::from_millis(15));
1566 }
1567 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1568 }
1569 }
1570 }
1571
1572 pub fn get(&self, id: &str) -> Result<Task> {
1574 let resolved = self.resolve_id(id)?;
1575 read_path(&self.path_of(&resolved))
1576 }
1577
1578 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1610 let _ = questions;
1611 let resolved = self.resolve_id(id)?;
1612 if in_flight {
1613 bail!("task {resolved} is being run by a live daemon right now");
1614 }
1615 self.write_tombstone(&resolved)?;
1616 if let Err(e) = self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))
1617 {
1618 if self.path_of(&resolved).exists() {
1622 let _ = std::fs::remove_file(self.tombstone_path(&resolved));
1623 }
1624 return Err(e);
1625 }
1626 let (released, still_blocked) = self.release_dependents_of(&resolved);
1627 Ok(Removal {
1628 id: resolved,
1629 released,
1630 still_blocked,
1631 })
1632 }
1633
1634 fn tombstone_path(&self, id: &str) -> PathBuf {
1637 self.root.join(format!("{id}.removed"))
1638 }
1639
1640 fn write_tombstone(&self, id: &str) -> Result<()> {
1641 let path = self.tombstone_path(id);
1642 let tmp = self.root.join(format!("{id}.removed.tmp"));
1643 std::fs::write(&tmp, b"")
1644 .and_then(|()| std::fs::rename(&tmp, &path))
1645 .with_context(|| format!("write {}", path.display()))
1646 }
1647
1648 fn was_deleted(&self, id: &str) -> bool {
1652 self.tombstone_path(id).is_file() && !self.path_of(id).exists()
1653 }
1654
1655 pub fn apply_deleted_blockers(&self, task: &mut Task) -> Vec<String> {
1659 let deleted = deleted_blockers(self, &task.blocked_by);
1660 deleted
1661 .into_iter()
1662 .filter(|id| task.dependency_deleted(id))
1663 .collect()
1664 }
1665
1666 pub fn note_dependency_deleted(&self, dependent: &Task, deleted: &str) {
1670 let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty()) else {
1671 return;
1672 };
1673 let ja = crate::lang::is_japanese(&crate::lang::of_repo(&dependent.repo));
1674 let waiting = dependent.status == TaskStatus::Blocked;
1675 let message = match (ja, waiting) {
1676 (true, false) => format!(
1677 "タスク {} は、待っていた {} が削除されたため待機を解除し、元の状態に戻しました",
1678 dependent.short(),
1679 short(deleted)
1680 ),
1681 (true, true) => format!(
1682 "タスク {} は、待っていた {} が削除されたため、残りの依存を待っています",
1683 dependent.short(),
1684 short(deleted)
1685 ),
1686 (false, false) => format!(
1687 "Task {} stopped waiting on {} because it was deleted, and returned to its previous state",
1688 dependent.short(),
1689 short(deleted)
1690 ),
1691 (false, true) => format!(
1692 "Task {} stopped waiting on {} because it was deleted, and is still waiting on its other dependencies",
1693 dependent.short(),
1694 short(deleted)
1695 ),
1696 };
1697 crate::notices::raise_in(
1698 home,
1699 crate::notices::Notice::info(
1700 &format!("unblocked:{}:{}", dependent.id, deleted),
1701 message,
1702 ),
1703 );
1704 }
1705
1706 fn release_dependents_of(&self, dependency: &str) -> (Vec<String>, Vec<String>) {
1711 let (mut released, mut still_blocked) = (Vec::new(), Vec::new());
1712 for listed in self.list() {
1713 if listed.status != TaskStatus::Blocked
1714 || !listed.blocked_by.iter().any(|b| b == dependency)
1715 {
1716 continue;
1717 }
1718 let Ok(_claim) = self.claim(&listed.id) else {
1719 continue;
1720 };
1721 let Ok(mut task) = self.get(&listed.id) else {
1722 continue;
1723 };
1724 if !task.dependency_deleted(dependency) {
1725 continue;
1726 }
1727 if self.put(&mut task).is_err() {
1728 continue;
1729 }
1730 self.note_dependency_deleted(&task, dependency);
1731 if task.status == TaskStatus::Blocked {
1732 still_blocked.push(task.id.clone());
1733 } else {
1734 released.push(task.id.clone());
1735 }
1736 }
1737 (released, still_blocked)
1738 }
1739
1740 fn remove_record_with_attachments(
1744 &self,
1745 resolved: &str,
1746 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1747 ) -> Result<()> {
1748 self.sweep_removed_attachments();
1755 let attachments = self.attachments_dir(resolved);
1756 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1757 let moved = match std::fs::rename(&attachments, &aside) {
1758 Ok(()) => true,
1759 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1760 Err(e) => {
1761 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1762 }
1763 };
1764 let path = self.path_of(resolved);
1765 if let Err(e) = remove_record(&path) {
1766 if moved {
1767 let _ = std::fs::rename(&aside, &attachments);
1768 }
1769 return Err(e).with_context(|| format!("remove {}", path.display()));
1770 }
1771 if moved {
1772 if let Err(e) = std::fs::remove_dir_all(&aside) {
1773 tracing::warn!("leftover attachments {}: {e}", aside.display());
1774 }
1775 }
1776 let lock = self.lock_path(resolved);
1777 if let Err(e) = std::fs::remove_file(&lock) {
1778 if e.kind() != std::io::ErrorKind::NotFound {
1779 return Err(e).with_context(|| format!("remove {}", lock.display()));
1780 }
1781 }
1782 Ok(())
1783 }
1784
1785 fn sweep_removed_attachments(&self) {
1791 let Ok(entries) = std::fs::read_dir(&self.root) else {
1792 return;
1793 };
1794 for entry in entries.flatten() {
1795 let name = entry.file_name();
1796 let name = name.to_string_lossy();
1797 let Some(id) = name.strip_suffix(".attachments.removing") else {
1798 continue;
1799 };
1800 if !self.path_of(id).exists() {
1803 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1804 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1805 }
1806 }
1807 }
1808 self.sweep_tombstones();
1809 }
1810
1811 fn sweep_tombstones(&self) {
1820 let Ok(entries) = std::fs::read_dir(&self.root) else {
1821 return;
1822 };
1823 let tasks = self.list();
1824 for entry in entries.flatten() {
1825 let name = entry.file_name();
1826 let name = name.to_string_lossy();
1827 let Some(id) = name.strip_suffix(".removed") else {
1828 continue;
1829 };
1830 if self.path_of(id).exists()
1831 || tasks.iter().any(|t| t.blocked_by.iter().any(|b| b == id))
1832 {
1833 continue;
1834 }
1835 let old = entry
1836 .metadata()
1837 .and_then(|m| m.modified())
1838 .ok()
1839 .and_then(|m| m.elapsed().ok())
1840 .is_some_and(|age| age >= TOMBSTONE_GRACE);
1841 if old {
1842 let _ = std::fs::remove_file(entry.path());
1843 }
1844 }
1845 }
1846
1847 fn lock_path(&self, id: &str) -> PathBuf {
1850 self.root.join(format!("{id}.lock"))
1851 }
1852
1853 pub fn list(&self) -> Vec<Task> {
1866 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1867 .into_iter()
1868 .flatten()
1869 .flatten()
1870 .map(|e| e.path())
1871 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1872 .filter_map(|p| read_path(&p).ok())
1873 .collect();
1874 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1875 tasks
1876 }
1877
1878 pub fn superseded(&self) -> HashMap<String, String> {
1891 let mut by = HashMap::new();
1892 for task in self.list() {
1893 for earlier in &task.runs {
1894 if let Some(later) = task.successor_of(earlier) {
1895 by.insert(earlier.clone(), later.clone());
1896 }
1897 }
1898 }
1899 by
1900 }
1901
1902 pub fn superseded_by(&self, run: &str) -> Option<String> {
1910 for task in self.list() {
1911 if task.runs.iter().any(|r| r == run) {
1912 return task.successor_of(run).cloned();
1913 }
1914 }
1915 None
1916 }
1917
1918 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1929 for task in self.list() {
1930 if task.runs.iter().any(|r| r == run) {
1931 return task.runs.last().filter(|last| **last != run).cloned();
1932 }
1933 }
1934 None
1935 }
1936
1937 pub fn next_runnable(&self) -> Option<Task> {
1942 let mut runnable: Vec<Task> = self
1943 .list()
1944 .into_iter()
1945 .filter(|t| t.status.runnable())
1946 .collect();
1947 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1948 runnable.into_iter().next()
1949 }
1950
1951 pub fn claim(&self, id: &str) -> Result<Claim> {
1958 std::fs::create_dir_all(&self.root)
1959 .with_context(|| format!("create {}", self.root.display()))?;
1960 let path = self.lock_path(id);
1961 match std::fs::OpenOptions::new()
1962 .write(true)
1963 .create_new(true)
1964 .open(&path)
1965 {
1966 Ok(mut f) => {
1967 use std::io::Write as _;
1968 let _ = writeln!(f, "{}", std::process::id());
1970 Ok(Claim { path })
1971 }
1972 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1973 bail!("task {id} is already claimed ({} exists)", path.display())
1974 }
1975 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1976 }
1977 }
1978
1979 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1981 if self.path_of(prefix).is_file() {
1982 return Ok(prefix.to_owned());
1983 }
1984 let hits: Vec<String> = self
1985 .list()
1986 .into_iter()
1987 .map(|t| t.id)
1988 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1989 .collect();
1990 match hits.len() {
1991 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1992 0 => bail!("no task matches `{prefix}`"),
1993 _ => bail!(
1994 "`{prefix}` matches {} tasks: {}",
1995 hits.len(),
1996 hits.join(", ")
1997 ),
1998 }
1999 }
2000
2001 pub fn revision(&self) -> u64 {
2008 self.revision_excluding(&std::collections::BTreeSet::new())
2009 }
2010
2011 pub fn revision_excluding(&self, skip: &std::collections::BTreeSet<String>) -> u64 {
2014 use std::hash::{Hash as _, Hasher as _};
2015
2016 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
2017 .into_iter()
2018 .flatten()
2019 .flatten()
2020 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
2021 .filter(|e| {
2022 let path = e.path();
2023 !path
2024 .file_stem()
2025 .is_some_and(|stem| skip.contains(stem.to_string_lossy().as_ref()))
2026 })
2027 .filter_map(|e| {
2028 let name = e.file_name().to_string_lossy().into_owned();
2029 let mtime = e
2030 .metadata()
2031 .ok()?
2032 .modified()
2033 .ok()?
2034 .duration_since(std::time::UNIX_EPOCH)
2035 .ok()?
2036 .as_millis() as u64;
2037 Some((name, mtime))
2038 })
2039 .collect();
2040
2041 if entries.is_empty() {
2042 return 0;
2043 }
2044
2045 entries.sort_unstable();
2046 let mut hasher = std::hash::DefaultHasher::new();
2047 for (name, mtime) in &entries {
2048 name.hash(&mut hasher);
2049 mtime.hash(&mut hasher);
2050 }
2051 let h = hasher.finish();
2052 if h == 0 { 1 } else { h }
2053 }
2054}
2055
2056const TOMBSTONE_GRACE: std::time::Duration = std::time::Duration::from_secs(3600);
2059
2060#[derive(Debug, Clone)]
2062pub struct Removal {
2063 pub id: String,
2065 pub released: Vec<String>,
2069 pub still_blocked: Vec<String>,
2072}
2073
2074#[derive(Debug)]
2076pub struct Claim {
2077 path: PathBuf,
2078}
2079
2080impl Drop for Claim {
2081 fn drop(&mut self) {
2082 let _ = std::fs::remove_file(&self.path);
2083 }
2084}
2085
2086pub fn title_from(instruction: &str, max: usize) -> String {
2089 let Some(line) = first_line(instruction) else {
2090 return "(empty task)".to_owned();
2091 };
2092 if line.chars().count() <= max {
2093 return line.to_owned();
2094 }
2095 let head: String = line.chars().take(max.saturating_sub(1)).collect();
2096 format!("{head}…")
2097}
2098
2099pub(crate) fn first_line(instruction: &str) -> Option<&str> {
2105 let line = instruction
2106 .lines()
2107 .map(str::trim)
2108 .find(|l| !l.is_empty())?
2109 .trim_start_matches(['#', '-', '*', '>', ' '])
2110 .trim();
2111 (!line.is_empty()).then_some(line)
2112}
2113
2114fn read_path(path: &Path) -> Result<Task> {
2115 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2116 let task: Task =
2117 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
2118 if task.schema > SCHEMA {
2124 bail!(
2125 "task {} was written by a different magi (schema {}, this build \
2126 speaks {SCHEMA})",
2127 task.id,
2128 task.schema
2129 );
2130 }
2131 Ok(task)
2132}
2133
2134pub fn missing_blockers(
2148 queue: &Queue,
2149 questions: &Questions,
2150 blocked_by: &[String],
2151) -> Vec<String> {
2152 blocked_by
2153 .iter()
2154 .filter(|id| {
2155 !queue.path_of(id).is_file()
2156 && !questions.path_of(id).is_file()
2157 && !queue.was_deleted(id)
2158 })
2159 .cloned()
2160 .collect()
2161}
2162
2163pub fn deleted_blockers(queue: &Queue, blocked_by: &[String]) -> Vec<String> {
2168 blocked_by
2169 .iter()
2170 .filter(|id| queue.was_deleted(id))
2171 .cloned()
2172 .collect()
2173}
2174
2175pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
2187 missing_blocker_hold_reason_in(blocked_by, missing, "en")
2188}
2189
2190pub fn missing_blocker_hold_reason_in(
2192 blocked_by: &[String],
2193 missing: &[String],
2194 language: &str,
2195) -> String {
2196 if crate::lang::is_japanese(language) {
2197 format!(
2198 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
2199 blocked_by.join(", "),
2200 missing.join(", "),
2201 )
2202 } else {
2203 format!(
2204 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
2205 blocked_by.join(", "),
2206 missing.join(", "),
2207 )
2208 }
2209}
2210
2211pub fn short(id: &str) -> &str {
2213 id.split('-').next_back().unwrap_or(id)
2214}
2215
2216fn new_id() -> String {
2217 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2218 let seed = crate::rng::entropy();
2219 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2220}
2221
2222#[cfg(test)]
2223mod tests {
2224 #[test]
2225 fn merge_over_lets_the_later_choice_win_and_keeps_the_rest() {
2226 let mut base = RunOverrides {
2227 merge: Some("pr".to_owned()),
2228 candidates: Some(3),
2229 ..RunOverrides::default()
2230 };
2231 base.merge_over(&RunOverrides {
2232 merge: Some("none".to_owned()),
2233 seed: Some(7),
2234 ..RunOverrides::default()
2235 });
2236 assert_eq!(base.merge.as_deref(), Some("none"));
2237 assert_eq!(base.candidates, Some(3));
2238 assert_eq!(base.seed, Some(7));
2239 }
2240
2241 #[test]
2242 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
2243 let mut t = task("relanded");
2244 t.attempts = 1;
2245 t.status = TaskStatus::Running;
2246 t.already_landed("already in main as 0e368de");
2247 assert_eq!(t.status, TaskStatus::Done);
2248 assert_eq!(t.attempts, 0);
2249 assert!(t.hold_reason.is_none());
2250 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
2251 }
2252
2253 #[test]
2254 fn missing_blocker_reason_follows_the_language() {
2255 let b = vec!["a".to_owned()];
2256 let en = missing_blocker_hold_reason_in(&b, &b, "en");
2257 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
2258 assert!(en.starts_with("blocked on a"));
2259 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
2260 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
2261 }
2262
2263 use super::*;
2264
2265 #[test]
2266 fn triage_applied_survives_release_and_old_records_read_as_empty() {
2267 let mut t = Task::new(
2268 "t".to_owned(),
2269 "i".to_owned(),
2270 PathBuf::from("r"),
2271 Source::Human,
2272 );
2273 t.mark_triage_applied("q1");
2274 t.mark_triage_applied("q1");
2275 t.hold_machine(Some("x".to_owned()));
2276 t.release();
2277 assert_eq!(t.triage_applied, ["q1"]);
2278 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
2279
2280 let mut v = serde_json::to_value(&t).unwrap();
2281 v.as_object_mut().unwrap().remove("triage_applied");
2282 let old: Task = serde_json::from_value(v).unwrap();
2283 assert!(old.triage_applied.is_empty());
2284 }
2285
2286 #[test]
2287 fn task_counts_of_empty_is_all_zero() {
2288 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
2289 }
2290
2291 #[test]
2292 fn task_counts_of_tallies_every_status() {
2293 let mut queued = Task::new(
2294 "q".to_owned(),
2295 "i".to_owned(),
2296 PathBuf::from("."),
2297 Source::Human,
2298 );
2299 queued.status = TaskStatus::Queued;
2300 let mut running = queued.clone();
2301 running.status = TaskStatus::Running;
2302 let mut done = queued.clone();
2303 done.status = TaskStatus::Done;
2304 let mut failed = queued.clone();
2305 failed.status = TaskStatus::Failed;
2306 let mut held = queued.clone();
2307 held.status = TaskStatus::Held;
2308 let mut blocked = queued.clone();
2309 blocked.status = TaskStatus::Blocked;
2310
2311 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
2312 assert_eq!(
2313 counts,
2314 TaskCounts {
2315 queued: 1,
2316 running: 1,
2317 done: 2,
2318 failed: 1,
2319 held: 1,
2320 blocked: 1,
2321 }
2322 );
2323 }
2324
2325 fn queue() -> (tempfile::TempDir, Queue) {
2328 let dir = tempfile::tempdir().unwrap();
2329 let q = Queue::at(dir.path().join("queue"));
2330 (dir, q)
2331 }
2332
2333 #[test]
2334 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
2335 let dir = tempfile::tempdir().unwrap();
2336 let q = Queue::at(dir.path().join("queue"));
2337 let mut t = task("held");
2338 q.put(&mut t).unwrap();
2339 assert_eq!(
2340 crate::notices::Notices::at(dir.path().join("notifications"))
2341 .list()
2342 .len(),
2343 0
2344 );
2345 t.hold_machine(Some("out of attempts".to_owned()));
2346 q.put(&mut t).unwrap();
2347 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
2348 assert_eq!(listed.len(), 1);
2349 assert!(listed[0].message.contains("out of attempts"));
2350 }
2351
2352 fn task(title: &str) -> Task {
2353 Task::new(
2354 title.to_owned(),
2355 format!("do {title}"),
2356 PathBuf::from("."),
2357 Source::Human,
2358 )
2359 }
2360
2361 #[test]
2362 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
2363 let mut t = task("retried");
2364 assert!(
2365 t.earlier_attempts().is_empty(),
2366 "a first attempt takes nothing over"
2367 );
2368 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2369 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
2370 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
2372 assert_eq!(t.successor_of("bbbb"), None);
2373 assert_eq!(t.successor_of("zzzz"), None);
2374 }
2375
2376 #[test]
2377 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
2378 let (_dir, q) = queue();
2379 let mut t = task("retried");
2380 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2381 q.put(&mut t).unwrap();
2382
2383 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
2384 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
2385 assert_eq!(
2386 q.superseded_by("cccc"),
2387 None,
2388 "the latest attempt replaces nothing"
2389 );
2390 assert_eq!(
2391 q.superseded_by("never-heard-of-it"),
2392 None,
2393 "a run belonging to no task on this queue is not superseded"
2394 );
2395
2396 let mut by = HashMap::new();
2397 by.insert("aaaa".to_owned(), "bbbb".to_owned());
2398 by.insert("bbbb".to_owned(), "cccc".to_owned());
2399 assert_eq!(
2400 q.superseded(),
2401 by,
2402 "the whole-map and single-run forms must agree"
2403 );
2404 }
2405
2406 #[test]
2407 fn superseded_attempts_is_empty_until_the_task_is_done() {
2408 let mut t = task("retried");
2409 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2410 t.status = TaskStatus::Failed;
2411 assert_eq!(
2412 t.superseded_attempts(true),
2413 &[] as &[String],
2414 "a task still retrying has no attempt yet that a later one made moot"
2415 );
2416
2417 t.status = TaskStatus::Running;
2418 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2419 }
2420
2421 #[test]
2422 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
2423 let mut t = task("retried");
2424 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2425 t.status = TaskStatus::Done;
2426 assert_eq!(
2427 t.superseded_attempts(true),
2428 &["aaaa".to_owned(), "bbbb".to_owned()],
2429 "cccc is the attempt whose success made the task done, and stays out"
2430 );
2431 }
2432
2433 #[test]
2434 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
2435 let mut t = task("first try landed");
2436 t.runs = vec!["aaaa".to_owned()];
2437 t.status = TaskStatus::Done;
2438 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2439 }
2440
2441 #[test]
2442 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
2443 let mut t = task("closed by hand after a manual merge");
2450 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2451 t.status = TaskStatus::Done;
2452 assert_eq!(
2453 t.superseded_attempts(false),
2454 &[] as &[String],
2455 "nothing here is provably why the task is done, so nothing is superseded"
2456 );
2457 }
2458
2459 #[test]
2460 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2461 let (_dir, q) = queue();
2462 let mut t = task("retried twice");
2463 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2464 q.put(&mut t).unwrap();
2465
2466 assert_eq!(
2467 q.latest_attempt("aaaa"),
2468 Some("cccc".to_owned()),
2469 "an old attempt points straight at the chain's current head, not the \
2470 next attempt in the middle of it"
2471 );
2472 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2473 assert_eq!(
2474 q.latest_attempt("cccc"),
2475 None,
2476 "the latest attempt is not superseded by anything"
2477 );
2478 assert_eq!(
2479 q.latest_attempt("never-heard-of-it"),
2480 None,
2481 "a run belonging to no task on this queue is not superseded"
2482 );
2483 }
2484
2485 #[test]
2486 fn a_markdown_heading_is_the_title_not_decoration() {
2487 assert_eq!(
2492 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2493 "Rework the config loader"
2494 );
2495 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2496 assert_eq!(title_from("> quoted task", 40), "quoted task");
2497 assert_eq!(title_from(" \n\n", 40), "(empty task)");
2499 assert_eq!(title_from("###\n", 40), "(empty task)");
2500 }
2501
2502 #[test]
2503 fn a_long_title_is_elided_by_characters_not_bytes() {
2504 let long = "課題".repeat(30);
2506 let title = title_from(&long, 10);
2507 assert_eq!(title.chars().count(), 10);
2508 assert!(title.ends_with('…'));
2509 }
2510
2511 #[test]
2512 fn priority_wins_and_ties_break_oldest_first() {
2513 let (_dir, q) = queue();
2514 let mut a = task("first");
2515 let mut b = task("second");
2516 let mut c = task("urgent");
2517 a.id = "20260101-000001-aaaa".to_owned();
2519 b.id = "20260101-000002-bbbb".to_owned();
2520 c.id = "20260101-000003-cccc".to_owned();
2521 c.priority = 5;
2522 for t in [&mut a, &mut b, &mut c] {
2523 q.put(t).unwrap();
2524 }
2525
2526 assert_eq!(q.next_runnable().unwrap().id, c.id);
2528 c.hold_machine(None);
2529 q.put(&mut c).unwrap();
2530 assert_eq!(q.next_runnable().unwrap().id, a.id);
2532 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2533 }
2534
2535 #[test]
2536 fn a_blocked_task_never_starves_another_runnable_one() {
2537 let (_dir, q) = queue();
2538 let mut blocked = task("blocked");
2539 blocked.block(vec!["something".to_owned()], None);
2540 q.put(&mut blocked).unwrap();
2541
2542 let mut runnable = task("free to go");
2543 q.put(&mut runnable).unwrap();
2544
2545 let next = q.next_runnable().expect("a runnable task is still offered");
2546 assert_eq!(next.id, runnable.id);
2547 }
2548
2549 #[test]
2550 fn a_held_task_is_never_offered_to_the_loop() {
2551 let (_dir, q) = queue();
2552 let mut t = task("held");
2553 q.put(&mut t).unwrap();
2554 assert!(q.next_runnable().is_some());
2555
2556 t.hold_machine(None);
2557 q.put(&mut t).unwrap();
2558 assert!(
2559 q.next_runnable().is_none(),
2560 "a held task must wait for a human"
2561 );
2562
2563 t.status = TaskStatus::Failed;
2565 q.put(&mut t).unwrap();
2566 assert!(q.next_runnable().is_some());
2567 }
2568
2569 #[test]
2570 fn attempts_are_capped_and_then_the_task_is_held() {
2571 let mut t = task("doomed");
2572
2573 t.start("run-1".to_owned());
2574 t.fail("gate red", 2);
2575 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2576
2577 t.start("run-2".to_owned());
2578 t.fail("gate red", 2);
2579 assert_eq!(
2580 t.status,
2581 TaskStatus::Held,
2582 "out of attempts: stop spending money on it"
2583 );
2584 assert_eq!(t.runs, ["run-1", "run-2"]);
2585 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2586 assert_eq!(
2587 t.hold_reason.as_deref(),
2588 Some("gate red"),
2589 "the hold must say why, not leave hold_reason null next to a \
2590 populated last_error"
2591 );
2592 }
2593
2594 #[test]
2595 fn handing_off_a_task_records_a_hold_reason_too() {
2596 let mut t = task("left a pull request");
2597 t.start("run-1".to_owned());
2598 t.handed_off("run ended with a pull request open [run run-1]");
2599 assert_eq!(t.status, TaskStatus::Held);
2600 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2601 assert_eq!(
2602 t.hold_reason.as_deref(),
2603 Some("run ended with a pull request open [run run-1]")
2604 );
2605 assert_eq!(t.hold_reason, t.last_error);
2606 }
2607
2608 #[test]
2609 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2610 let mut t = task("stalled by quota");
2611
2612 t.start("run-1".to_owned());
2613 assert_eq!(t.attempts, 1);
2614 t.stall("judge-1, judge-2 out of quota");
2615 assert_eq!(
2616 t.attempts, 0,
2617 "a closed quota window must not spend the task's retry budget"
2618 );
2619 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2620 assert_eq!(
2621 t.last_error.as_deref(),
2622 Some("judge-1, judge-2 out of quota")
2623 );
2624
2625 for _ in 0..20 {
2628 t.start("run-n".to_owned());
2629 t.stall("still out of quota");
2630 }
2631 t.start("run-real".to_owned());
2632 t.fail("gate red", 2);
2633 assert_eq!(
2634 t.status,
2635 TaskStatus::Failed,
2636 "the first attempt that was really judged is attempt one"
2637 );
2638 }
2639
2640 #[test]
2641 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2642 let mut t = task("retry me");
2643 t.start("run-1".to_owned());
2644 t.fail("gate red", 1);
2645 assert_eq!(t.status, TaskStatus::Held);
2646
2647 t.release();
2648 assert_eq!(t.status, TaskStatus::Queued);
2649 assert_eq!(t.attempts, 0);
2652 assert!(t.last_error.is_none());
2653 assert_eq!(
2654 t.runs.len(),
2655 1,
2656 "history is kept: attempts reset, evidence does not"
2657 );
2658 }
2659
2660 #[test]
2661 fn a_hold_reason_survives_and_a_release_clears_it() {
2662 let mut t = task("waiting on something else");
2663 t.hold_manual(Some(
2664 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2665 ));
2666 assert_eq!(t.status, TaskStatus::Held);
2667 assert_eq!(
2668 t.hold_reason.as_deref(),
2669 Some("waiting for 20260101-000000-aaaa to land first")
2670 );
2671
2672 t.hold_manual(None);
2674 assert_eq!(
2675 t.hold_reason.as_deref(),
2676 Some("waiting for 20260101-000000-aaaa to land first"),
2677 "a bare re-hold keeps whatever a human already wrote down"
2678 );
2679
2680 let mut plain = task("no reason given");
2682 plain.hold_manual(None);
2683 assert_eq!(plain.status, TaskStatus::Held);
2684 assert!(plain.hold_reason.is_none());
2685
2686 t.release();
2687 assert_eq!(t.status, TaskStatus::Queued);
2688 assert!(
2689 t.hold_reason.is_none(),
2690 "a stale reason must not greet the next person who holds this task"
2691 );
2692 }
2693
2694 #[test]
2695 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2696 let mut t = task("landed by hand while held");
2701 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2702 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2703
2704 t.succeed();
2705 assert_eq!(t.status, TaskStatus::Done);
2706 assert!(
2707 t.hold_reason.is_none(),
2708 "a done task cannot still be waiting on something"
2709 );
2710 }
2711
2712 #[test]
2713 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2714 let mut held = task("held straight out of blocked");
2721 held.block(
2722 vec!["20260101-000000-dead".to_owned()],
2723 Some("waiting on the migration script".to_owned()),
2724 );
2725 assert_eq!(held.status, TaskStatus::Blocked);
2726
2727 held.hold_manual(None);
2728 assert_eq!(held.status, TaskStatus::Held);
2729 assert!(
2730 held.blocked_by.is_empty(),
2731 "hold overrides the wait, same as release"
2732 );
2733 assert!(held.block_reason.is_none());
2734
2735 let mut done = task("closed straight out of blocked");
2736 done.block(
2737 vec!["20260101-000000-dead".to_owned()],
2738 Some("waiting on the migration script".to_owned()),
2739 );
2740 done.succeed();
2741 assert_eq!(done.status, TaskStatus::Done);
2742 assert!(
2743 done.blocked_by.is_empty(),
2744 "a done task cannot still be waiting on a dependency"
2745 );
2746 assert!(done.block_reason.is_none());
2747 }
2748
2749 #[test]
2750 fn a_blocked_task_is_never_offered_to_the_loop() {
2751 let mut t = task("blocked");
2752 assert!(t.status.runnable());
2753 t.block(
2754 vec!["dep-id".to_owned()],
2755 Some("waits on dep-id".to_owned()),
2756 );
2757 assert_eq!(t.status, TaskStatus::Blocked);
2758 assert!(!t.status.runnable());
2759 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2760 }
2761
2762 #[test]
2763 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2764 let mut t = task("blocked on two");
2765 t.block(
2766 vec!["a".to_owned(), "b".to_owned()],
2767 Some("waits on a and b".to_owned()),
2768 );
2769
2770 t.unblock("a");
2771 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2772 assert_eq!(t.blocked_by, ["b"]);
2773
2774 t.unblock("b");
2775 assert_eq!(t.status, TaskStatus::Queued);
2776 assert!(t.blocked_by.is_empty());
2777 assert!(t.block_reason.is_none());
2778 }
2779
2780 #[test]
2781 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2782 let mut t = task("never blocked");
2783 t.unblock("whatever");
2784 assert_eq!(t.status, TaskStatus::Queued);
2785 }
2786
2787 #[test]
2788 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2789 let mut t = task("held, then asked about");
2794 t.hold_machine(Some("out of attempts".to_owned()));
2795 assert_eq!(t.status, TaskStatus::Held);
2796
2797 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2798 assert_eq!(t.status, TaskStatus::Blocked);
2799
2800 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2801 t.unblock("q1");
2802 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2803 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2804 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2805 assert!(t.blocked_from.is_none(), "consumed once restored");
2806 }
2807
2808 #[test]
2809 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2810 let mut t = task("manually held, then asked about");
2811 t.hold_manual(Some("waiting on a dependency".to_owned()));
2812
2813 t.block(vec!["q1".to_owned()], None);
2814 t.unblock("q1");
2815
2816 assert_eq!(t.status, TaskStatus::Held);
2817 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2818 }
2819
2820 #[test]
2821 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2822 let mut t = task("held, blocked twice");
2826 t.hold_machine(None);
2827 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2828 t.block(
2829 vec!["q1".to_owned(), "q2".to_owned()],
2830 Some("second".to_owned()),
2831 );
2832
2833 t.unblock("q1");
2834 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2835 t.unblock("q2");
2836 assert_eq!(t.status, TaskStatus::Held);
2837 }
2838
2839 #[test]
2840 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2841 let mut t = task("blocked mid-run");
2845 t.start("run-1".to_owned());
2846 assert_eq!(t.status, TaskStatus::Running);
2847
2848 t.block(vec!["q1".to_owned()], None);
2849 t.unblock("q1");
2850 assert_eq!(t.status, TaskStatus::Queued);
2851 }
2852
2853 #[test]
2854 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2855 let mut t = task("legacy record, held before it was blocked");
2861 t.hold_source = Some(HoldSource::Machine);
2862 t.hold_reason = Some("legacy hold reason".to_owned());
2863 t.status = TaskStatus::Blocked;
2864 t.blocked_by = vec!["q1".to_owned()];
2865 t.blocked_from = None;
2866
2867 t.unblock("q1");
2868 assert_eq!(t.status, TaskStatus::Held);
2869 }
2870
2871 #[test]
2872 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2873 let mut t = task("legacy record, ordinary dependency block");
2874 t.status = TaskStatus::Blocked;
2875 t.blocked_by = vec!["dep".to_owned()];
2876 t.blocked_from = None;
2877
2878 t.unblock("dep");
2879 assert_eq!(t.status, TaskStatus::Queued);
2880 }
2881
2882 #[test]
2883 fn answering_a_question_is_recorded_and_survives_a_release() {
2884 let mut t = task("asked something");
2885 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2886 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2887 t.unblock("q1");
2888 assert_eq!(t.status, TaskStatus::Queued);
2889 assert_eq!(t.answers.len(), 1);
2890 assert_eq!(t.answers[0].answer, "SQLite");
2891
2892 t.release();
2896 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2897 }
2898
2899 #[test]
2900 fn a_refused_handover_keeps_the_review_branch_across_release() {
2901 let mut t = task("refused takeover");
2902 t.start("run-1".to_owned());
2903 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
2904 assert_eq!(t.status, TaskStatus::Held);
2905 assert_eq!(t.attempts, 1);
2906 t.release();
2907 assert_eq!(t.status, TaskStatus::Queued);
2908 assert_eq!(t.attempts, 0);
2909 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2910
2911 t.hold_for_handover(None, "again".to_owned());
2913 t.requeue();
2914 assert!(t.review_branch.is_none());
2915
2916 let mut m = task("manual");
2918 m.review_branch = Some("magi/x/A".to_owned());
2919 m.hold_manual(None);
2920 m.release();
2921 assert!(m.review_branch.is_none());
2922 }
2923
2924 #[test]
2925 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2926 let mut t = task("blocked run with a surviving branch");
2927 t.start("run-1".to_owned());
2928 t.fail("blocked with major findings", 5);
2929 assert_eq!(t.status, TaskStatus::Failed);
2930
2931 t.request_review("magi/eba2/A".to_owned());
2932 assert_eq!(t.status, TaskStatus::Queued);
2933 assert_eq!(t.attempts, 0);
2934 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2935
2936 t.release();
2938 assert!(t.review_branch.is_none());
2939 }
2940
2941 #[test]
2942 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2943 let mut t = task("retry");
2944 t.start("run-1".to_owned());
2945 t.requeue();
2946 assert!(t.fresh_start);
2947
2948 t.release();
2949 assert!(!t.fresh_start);
2950 }
2951
2952 #[test]
2953 fn priority_can_be_changed_while_queued_but_not_while_running() {
2954 let mut t = task("reprioritise me");
2955 t.set_priority(5).unwrap();
2956 assert_eq!(t.priority, 5);
2957
2958 t.start("run-1".to_owned());
2959 let err = t.set_priority(9).unwrap_err().to_string();
2960 assert!(err.contains("running"), "{err}");
2961 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2962 }
2963
2964 #[test]
2965 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2966 let mut t = task("interrupt me");
2967 assert!(!t.interrupt, "off unless asked, same as any other task");
2968
2969 t.set_interrupt(true).unwrap();
2970 assert!(t.interrupt);
2971
2972 t.start("run-1".to_owned());
2973 assert!(
2974 !t.interrupt,
2975 "the mark is one-shot: dispatching the task fulfils it, \
2976 whatever the run that follows ends up doing"
2977 );
2978 let err = t.set_interrupt(true).unwrap_err().to_string();
2979 assert!(err.contains("running"), "{err}");
2980 t.set_interrupt(false).unwrap();
2983 assert!(!t.interrupt);
2984 }
2985
2986 #[test]
2990 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2991 let mut t = task("interrupt me");
2992 t.set_interrupt(true).unwrap();
2993 t.start("run-1".to_owned());
2994 t.fail("mock failure", 5);
2995 assert_eq!(t.status, TaskStatus::Failed);
2996 assert!(
2997 !t.interrupt,
2998 "one attempt already spent the mark; a retry is an ordinary \
2999 requeue, not a fresh interrupt request"
3000 );
3001 }
3002
3003 #[test]
3004 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
3005 let (_dir, q) = queue();
3006 let mut a = task("first filed");
3007 let mut b = task("second filed");
3008 a.id = "20260101-000001-aaaa".to_owned();
3009 b.id = "20260101-000002-bbbb".to_owned();
3010 q.put(&mut a).unwrap();
3011 q.put(&mut b).unwrap();
3012
3013 assert_eq!(
3014 q.next_runnable().unwrap().id,
3015 a.id,
3016 "with equal priority the older task goes first, so a burst of \
3017 new work cannot starve it"
3018 );
3019 assert_eq!(
3020 q.list()[0].id,
3021 b.id,
3022 "but the list an operator reads is newest first, the same as \
3023 before priority existed - a's turn to run does not make it the \
3024 newest task"
3025 );
3026
3027 let mut a = q.get(&a.id).unwrap();
3028 a.set_priority(10).unwrap();
3029 q.put(&mut a).unwrap();
3030
3031 assert_eq!(
3032 q.next_runnable().unwrap().id,
3033 a.id,
3034 "a raised priority must be reflected the moment it is saved"
3035 );
3036 assert_eq!(
3040 q.list()[0].id,
3041 a.id,
3042 "the raised task must sort first in the list an operator reads, \
3043 not only in next_runnable's own ordering"
3044 );
3045 }
3046
3047 #[test]
3048 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
3049 let mut t = Task::new(
3050 "old title".to_owned(),
3051 "old instruction".to_owned(),
3052 PathBuf::from("/repo"),
3053 Source::Agent {
3054 run: "20260101-000000-beef".to_owned(),
3055 node: "implement".to_owned(),
3056 },
3057 );
3058 let id = t.id.clone();
3059 let created_at = t.created_at;
3060 t.runs.push("20260101-000000-beef".to_owned());
3061
3062 t.edit("new title".to_owned(), "new instruction".to_owned())
3063 .unwrap();
3064
3065 assert_eq!(t.title, "new title");
3066 assert_eq!(t.instruction, "new instruction");
3067 assert_eq!(t.id, id, "editing must not mint a new id");
3068 assert_eq!(t.created_at, created_at);
3069 assert_eq!(
3070 t.source,
3071 Source::Agent {
3072 run: "20260101-000000-beef".to_owned(),
3073 node: "implement".to_owned(),
3074 },
3075 "editing must not turn agent attribution into human"
3076 );
3077 assert_eq!(t.runs, ["20260101-000000-beef"]);
3078 }
3079
3080 #[test]
3081 fn editing_is_refused_once_a_task_is_running_or_finished() {
3082 let mut running = task("in flight");
3083 running.start("run-1".to_owned());
3084 let err = running
3085 .edit("x".to_owned(), "y".to_owned())
3086 .unwrap_err()
3087 .to_string();
3088 assert!(err.contains("running"), "{err}");
3089
3090 let mut done = task("finished");
3091 done.succeed();
3092 let err = done
3093 .edit("x".to_owned(), "y".to_owned())
3094 .unwrap_err()
3095 .to_string();
3096 assert!(err.contains("done"), "{err}");
3097
3098 let mut queued = task("waiting");
3100 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
3101 let mut held = task("parked");
3102 held.hold_machine(None);
3103 held.edit("x".to_owned(), "y".to_owned()).unwrap();
3104 }
3105
3106 #[test]
3107 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
3108 let (_dir, q) = queue();
3109 let path = q.path_of("20260101-000000-aaaa");
3110 std::fs::create_dir_all(q.root()).unwrap();
3111 std::fs::write(
3112 &path,
3113 serde_json::json!({
3114 "schema": SCHEMA,
3115 "id": "20260101-000000-aaaa",
3116 "title": "from before hold reasons existed",
3117 "instruction": "from before hold reasons existed",
3118 "repo": ".",
3119 "source": { "kind": "human" },
3120 "status": "held",
3121 "created_at": Timestamp::now().to_string(),
3122 "updated_at": Timestamp::now().to_string(),
3123 })
3124 .to_string(),
3125 )
3126 .unwrap();
3127
3128 let task = q.get("20260101-000000-aaaa").expect("must still read");
3129 assert!(task.hold_reason.is_none());
3130 assert!(task.operator_held());
3131 }
3132
3133 #[test]
3134 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
3135 let (_dir, q) = queue();
3136 let path = q.path_of("20260101-000000-bbbb");
3137 std::fs::create_dir_all(q.root()).unwrap();
3138 std::fs::write(
3139 &path,
3140 serde_json::json!({
3141 "schema": 2,
3142 "id": "20260101-000000-bbbb",
3143 "title": "old manual recovery",
3144 "instruction": "old manual recovery",
3145 "repo": ".",
3146 "source": { "kind": "human" },
3147 "status": "held",
3148 "hold_reason": "active manual recovery run20260912-224242-daf5",
3149 "created_at": Timestamp::now().to_string(),
3150 "updated_at": Timestamp::now().to_string(),
3151 })
3152 .to_string(),
3153 )
3154 .unwrap();
3155
3156 let task = q.get("20260101-000000-bbbb").expect("must still read");
3157 assert_eq!(task.hold_source, None);
3158 assert!(task.operator_held());
3159 }
3160
3161 #[test]
3162 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
3163 let (_dir, q) = queue();
3164 let path = q.path_of("20260101-000000-aaaa");
3165 std::fs::create_dir_all(q.root()).unwrap();
3166 std::fs::write(
3167 &path,
3168 serde_json::json!({
3169 "schema": SCHEMA,
3170 "id": "20260101-000000-aaaa",
3171 "title": "from before diagnostics existed",
3172 "instruction": "from before diagnostics existed",
3173 "repo": ".",
3174 "source": { "kind": "human" },
3175 "status": "held",
3176 "created_at": Timestamp::now().to_string(),
3177 "updated_at": Timestamp::now().to_string(),
3178 })
3179 .to_string(),
3180 )
3181 .unwrap();
3182
3183 let task = q.get("20260101-000000-aaaa").expect("must still read");
3184 assert!(task.diagnostic.is_none());
3185 }
3186
3187 #[test]
3188 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
3189 let (_dir, q) = queue();
3193 let path = q.path_of("20260101-000000-aaaa");
3194 std::fs::create_dir_all(q.root()).unwrap();
3195 std::fs::write(
3196 &path,
3197 serde_json::json!({
3198 "schema": 1,
3199 "id": "20260101-000000-aaaa",
3200 "title": "from before blocking existed",
3201 "instruction": "from before blocking existed",
3202 "repo": ".",
3203 "source": { "kind": "human" },
3204 "status": "queued",
3205 "created_at": Timestamp::now().to_string(),
3206 "updated_at": Timestamp::now().to_string(),
3207 })
3208 .to_string(),
3209 )
3210 .unwrap();
3211
3212 let task = q.get("20260101-000000-aaaa").expect("must still read");
3213 assert!(task.blocked_by.is_empty());
3214 assert!(task.block_reason.is_none());
3215 assert!(task.answers.is_empty());
3216 assert!(task.review_branch.is_none());
3217 }
3218
3219 #[test]
3220 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
3221 let mut held = task("diagnosed");
3226 held.start("run-1".to_owned());
3227 held.fail("gate red", 1);
3228 held.diagnostic = Some("cargo test failed: ...".to_owned());
3229 assert_eq!(held.status, TaskStatus::Held);
3230
3231 held.release();
3232 assert!(held.diagnostic.is_none());
3233
3234 held.diagnostic = Some("cargo test failed: ...".to_owned());
3235 held.succeed();
3236 assert!(held.diagnostic.is_none());
3237 }
3238
3239 #[test]
3240 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
3241 let mut t = task("retried");
3242 t.start("run-1".to_owned());
3243 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
3244 t.fail("unrelated config error", 5);
3245 assert_eq!(t.status, TaskStatus::Failed);
3246 assert!(
3247 t.diagnostic.is_none(),
3248 "fail() must not let an old diagnostic outlive the run that produced it"
3249 );
3250 }
3251
3252 #[test]
3253 fn a_claim_is_exclusive_and_releases_on_drop() {
3254 let (_dir, q) = queue();
3255 let mut t = task("contended");
3256 q.put(&mut t).unwrap();
3257
3258 let held = q.claim(&t.id).unwrap();
3259 assert!(
3260 q.claim(&t.id).is_err(),
3261 "two daemons must not drive one task into two runs"
3262 );
3263 drop(held);
3264 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
3265 }
3266
3267 #[test]
3268 fn a_round_trip_survives_disk() {
3269 let (_dir, q) = queue();
3270 let mut t = Task::new(
3271 "titled".to_owned(),
3272 "body".to_owned(),
3273 PathBuf::from("/repo"),
3274 Source::Agent {
3275 run: "20260101-000000-beef".to_owned(),
3276 node: "implement".to_owned(),
3277 },
3278 );
3279 t.priority = 3;
3280 q.put(&mut t).unwrap();
3281
3282 let back = q.get(&t.id).unwrap();
3283 assert_eq!(back.id, t.id);
3284 assert_eq!(back.priority, 3);
3285 assert_eq!(back.source.label(), "implement@beef");
3286 assert_eq!(q.get(t.short()).unwrap().id, t.id);
3288 }
3289
3290 #[test]
3291 fn an_unreadable_task_does_not_take_the_queue_down() {
3292 let (_dir, q) = queue();
3293 let mut t = task("fine");
3294 q.put(&mut t).unwrap();
3295 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
3296
3297 let listed = q.list();
3298 assert_eq!(listed.len(), 1, "the readable task still lists");
3299 assert_eq!(listed[0].id, t.id);
3300 }
3301
3302 #[test]
3303 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
3304 let (_dir, q) = queue();
3305 let path = q.path_of("20260101-000000-aaaa");
3306 std::fs::create_dir_all(q.root()).unwrap();
3307 std::fs::write(
3308 &path,
3309 serde_json::json!({
3310 "schema": SCHEMA,
3311 "id": "20260101-000000-aaaa",
3312 "title": "from before solo existed",
3313 "instruction": "from before solo existed",
3314 "repo": ".",
3315 "source": { "kind": "human" },
3316 "status": "queued",
3317 "created_at": Timestamp::now().to_string(),
3318 "updated_at": Timestamp::now().to_string(),
3319 })
3320 .to_string(),
3321 )
3322 .unwrap();
3323
3324 let task = q.get("20260101-000000-aaaa").expect("must still read");
3325 assert!(!task.solo, "a queue file with no `solo` field means false");
3326 }
3327
3328 #[test]
3329 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
3330 let (_dir, q) = queue();
3331 let path = q.path_of("20260101-000000-bbbb");
3332 std::fs::create_dir_all(q.root()).unwrap();
3333 std::fs::write(
3334 &path,
3335 serde_json::json!({
3336 "schema": SCHEMA,
3337 "id": "20260101-000000-bbbb",
3338 "title": "from before urgent existed",
3339 "instruction": "from before urgent existed",
3340 "repo": ".",
3341 "source": { "kind": "human" },
3342 "status": "queued",
3343 "created_at": Timestamp::now().to_string(),
3344 "updated_at": Timestamp::now().to_string(),
3345 })
3346 .to_string(),
3347 )
3348 .unwrap();
3349
3350 let task = q.get("20260101-000000-bbbb").expect("must still read");
3351 assert!(
3352 !task.urgent,
3353 "a queue file with no `urgent` field means false, same as `solo`"
3354 );
3355 }
3356
3357 #[test]
3358 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
3359 let (_dir, q) = queue();
3360 let mut t = task("from the future");
3361 q.put(&mut t).unwrap();
3362 let path = q.path_of(&t.id);
3363 let body = std::fs::read_to_string(&path)
3364 .unwrap()
3365 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
3366 std::fs::write(&path, body).unwrap();
3367
3368 let err = q.get(&t.id).unwrap_err().to_string();
3369 assert!(err.contains("schema 99"), "{err}");
3370 }
3371
3372 #[test]
3373 fn revision_moves_when_the_queue_changes() {
3374 let (_dir, q) = queue();
3375 assert_eq!(q.revision(), 0, "an empty queue has no revision");
3376 let mut t = task("first");
3377 q.put(&mut t).unwrap();
3378 assert!(q.revision() > 0, "a written task moves the revision");
3379 }
3380
3381 #[test]
3382 fn revision_moves_when_deleting_an_older_task() {
3383 let (dir, q) = queue();
3384 let questions = Questions::at(dir.path().join("questions"));
3385 let mut t1 = task("older");
3386 q.put(&mut t1).unwrap();
3387 std::thread::sleep(std::time::Duration::from_millis(10));
3389 let mut t2 = task("newer");
3390 q.put(&mut t2).unwrap();
3391
3392 let rev_before = q.revision();
3393 q.remove(&t1.id, false, &questions).unwrap();
3394 let rev_after = q.revision();
3395
3396 assert_ne!(
3397 rev_before, rev_after,
3398 "deleting an older task must change the revision so other clients see the deletion"
3399 );
3400 }
3401
3402 #[test]
3403 fn removing_a_task_takes_it_out_of_the_listing() {
3404 let (dir, q) = queue();
3405 let questions = Questions::at(dir.path().join("questions"));
3406 let mut t = task("delete me");
3407 q.put(&mut t).unwrap();
3408 let removed = q.remove(t.short(), false, &questions).unwrap();
3409 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
3410 assert!(removed.released.is_empty() && removed.still_blocked.is_empty());
3411 assert!(q.list().is_empty());
3412 assert!(
3413 q.remove(&t.id, false, &questions).is_err(),
3414 "removing twice is an error"
3415 );
3416 }
3417
3418 #[test]
3419 fn removing_a_task_takes_its_stale_lock_with_it() {
3420 let (dir, q) = queue();
3421 let questions = Questions::at(dir.path().join("questions"));
3422 let mut t = task("interrupted");
3423 q.put(&mut t).unwrap();
3424
3425 let claim = q.claim(&t.id).unwrap();
3428 std::mem::forget(claim);
3429 assert!(
3430 q.claim(&t.id).is_err(),
3431 "the orphaned lock is what makes the task look claimed"
3432 );
3433
3434 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
3436 assert!(err.contains("live daemon"), "{err}");
3437 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
3438
3439 q.remove(&t.id, false, &questions).unwrap();
3441 assert!(q.list().is_empty());
3442 let mut again = task("interrupted");
3443 again.id = t.id.clone();
3444 q.put(&mut again).unwrap();
3445 assert!(
3446 q.claim(&t.id).is_ok(),
3447 "a task that comes back must be claimable, which a left-behind lock would prevent"
3448 );
3449 }
3450
3451 fn notices_of(dir: &Path) -> Vec<crate::notices::Notice> {
3452 crate::notices::Notices::at(dir.join("notifications")).list()
3453 }
3454
3455 #[test]
3456 fn removing_a_sole_dependency_releases_the_dependent_without_a_hold() {
3457 let (dir, q) = queue();
3458 let questions = Questions::at(dir.path().join("questions"));
3459 let mut dep = task("dependency");
3460 q.put(&mut dep).unwrap();
3461 let mut blocked = task("waiting");
3462 blocked.block(vec![dep.id.clone()], Some("waits".to_owned()));
3463 q.put(&mut blocked).unwrap();
3464
3465 let removed = q.remove(&dep.id, false, &questions).unwrap();
3466 assert_eq!(removed.released, [blocked.id.clone()]);
3467 assert!(removed.still_blocked.is_empty());
3468
3469 let after = q.get(&blocked.id).unwrap();
3470 assert_eq!(after.status, TaskStatus::Queued);
3471 assert!(after.blocked_by.is_empty());
3472 assert!(after.block_reason.is_none());
3473 let notes = notices_of(dir.path());
3474 assert_eq!(notes.len(), 1, "{notes:?}");
3475 assert_eq!(notes[0].severity, crate::notices::Severity::Info);
3476 }
3477
3478 #[test]
3479 fn removing_one_of_two_dependencies_keeps_the_other() {
3480 let (dir, q) = queue();
3481 let questions = Questions::at(dir.path().join("questions"));
3482 let mut dep = task("dependency");
3483 q.put(&mut dep).unwrap();
3484 let mut other = task("other");
3485 q.put(&mut other).unwrap();
3486 let mut blocked = task("waiting");
3487 blocked.block(
3488 vec![dep.id.clone(), other.id.clone()],
3489 Some("waits on both".to_owned()),
3490 );
3491 q.put(&mut blocked).unwrap();
3492
3493 let removed = q.remove(&dep.id, false, &questions).unwrap();
3494 assert!(removed.released.is_empty());
3495 assert_eq!(removed.still_blocked, [blocked.id.clone()]);
3496
3497 let after = q.get(&blocked.id).unwrap();
3498 assert_eq!(after.status, TaskStatus::Blocked);
3499 assert_eq!(after.blocked_by, [other.id.clone()]);
3500 assert!(missing_blockers(&q, &questions, &after.blocked_by).is_empty());
3503 }
3504
3505 #[test]
3506 fn a_dependency_deleted_before_its_dependents_were_rewritten_is_released_later() {
3507 let (dir, q) = queue();
3508 let questions = Questions::at(dir.path().join("questions"));
3509 let mut dep = task("dependency");
3510 q.put(&mut dep).unwrap();
3511 let mut blocked = task("waiting");
3512 blocked.block(vec![dep.id.clone()], None);
3513 q.put(&mut blocked).unwrap();
3514
3515 let claim = q.claim(&blocked.id).unwrap();
3517 let removed = q.remove(&dep.id, false, &questions).unwrap();
3518 assert!(removed.released.is_empty());
3519 drop(claim);
3520
3521 let mut task = q.get(&blocked.id).unwrap();
3522 assert_eq!(task.status, TaskStatus::Blocked);
3523 assert!(missing_blockers(&q, &questions, &task.blocked_by).is_empty());
3524 assert_eq!(q.apply_deleted_blockers(&mut task), [dep.id.clone()]);
3525 assert_eq!(task.status, TaskStatus::Queued);
3526 }
3527
3528 #[test]
3529 fn a_dependency_without_a_tombstone_is_still_missing() {
3530 let (dir, q) = queue();
3531 let questions = Questions::at(dir.path().join("questions"));
3532 let mut dep = task("dependency");
3533 q.put(&mut dep).unwrap();
3534 std::fs::remove_file(q.path_of(&dep.id)).unwrap();
3535 assert_eq!(
3536 missing_blockers(&q, &questions, std::slice::from_ref(&dep.id)),
3537 [dep.id.clone()]
3538 );
3539 assert!(deleted_blockers(&q, std::slice::from_ref(&dep.id)).is_empty());
3540 }
3541
3542 #[test]
3543 fn a_held_dependent_is_left_alone_by_a_dependency_deletion() {
3544 let mut t = task("held");
3545 t.hold_machine(Some("because".to_owned()));
3546 t.blocked_by = vec!["gone".to_owned()];
3547 assert!(!t.dependency_deleted("gone"));
3548 assert_eq!(t.status, TaskStatus::Held);
3549 }
3550
3551 #[test]
3552 fn a_failed_record_removal_takes_the_tombstone_back() {
3553 let (_dir, q) = queue();
3554 let mut t = task("stays");
3555 q.put(&mut t).unwrap();
3556 q.write_tombstone(&t.id).unwrap();
3557 let err = q.remove_record_with_attachments(&t.id, |_| Err(std::io::Error::other("nope")));
3558 assert!(err.is_err());
3559 assert!(!q.was_deleted(&t.id));
3562 }
3563
3564 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3565 let p = dir.join(name);
3566 std::fs::write(&p, body).unwrap();
3567 p
3568 }
3569
3570 #[test]
3571 fn an_attachment_copy_survives_deleting_its_source() {
3572 let (dir, q) = queue();
3573 let src = source_file(dir.path(), "shot.png", "pixels");
3574 let mut t = task("with a picture");
3575 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3576 q.put(&mut t).unwrap();
3577 std::fs::remove_file(&src).unwrap();
3578 assert_eq!(names, ["shot.png"]);
3579 let loaded = q.get(&t.id).unwrap();
3580 let paths = q.attachment_paths(&loaded);
3581 assert_eq!(paths.len(), 1);
3582 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3583 }
3584
3585 #[test]
3586 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3587 let (dir, q) = queue();
3588 let mut t = task("bad names");
3589 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3590 let src = source_file(dir.path(), name, "x");
3591 assert!(
3592 q.attach(&mut t, &[src]).is_err(),
3593 "`{name}` must be refused"
3594 );
3595 }
3596 assert!(!crate::ask::valid_asset_name("C:foo.png"));
3600 #[cfg(not(windows))]
3601 {
3602 let src = source_file(dir.path(), "C:foo.png", "x");
3603 assert!(q.attach(&mut t, &[src]).is_err());
3604 }
3605 let long = format!("{}.png", "a".repeat(70));
3606 let src = source_file(dir.path(), &long, "x");
3607 assert!(q.attach(&mut t, &[src]).is_err());
3608 assert!(t.attachments.is_empty());
3609 assert!(!q.attachments_dir(&t.id).exists());
3610 }
3611
3612 #[test]
3613 fn a_taken_attachment_name_is_numbered_not_overwritten() {
3614 let (dir, q) = queue();
3615 let a = source_file(dir.path(), "shot.png", "one");
3616 let sub = dir.path().join("other");
3617 std::fs::create_dir_all(&sub).unwrap();
3618 let b = source_file(&sub, "shot.png", "two");
3619 let mut t = task("collision");
3620 q.attach(&mut t, &[a]).unwrap();
3621 q.attach(&mut t, &[b]).unwrap();
3622 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3623 let paths = q.attachment_paths(&t);
3624 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3625 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3626 }
3627
3628 #[test]
3629 fn a_renumbered_name_stays_inside_the_length_bound() {
3630 let (dir, q) = queue();
3631 let name = format!("{}.png", "a".repeat(60));
3632 assert_eq!(name.len(), 64);
3633 let a = source_file(dir.path(), &name, "one");
3634 let sub = dir.path().join("other");
3635 std::fs::create_dir_all(&sub).unwrap();
3636 let b = source_file(&sub, &name, "two");
3637 let mut t = task("long");
3638 q.attach(&mut t, &[a, b]).unwrap();
3639 assert_eq!(t.attachments.len(), 2);
3640 assert!(
3641 t.attachments
3642 .iter()
3643 .all(|n| crate::ask::valid_asset_name(n))
3644 );
3645 assert!(t.attachments[1].ends_with("-2.png"));
3646 }
3647
3648 #[test]
3649 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3650 let (dir, q) = queue();
3651 let good = source_file(dir.path(), "good.png", "ok");
3652 let mut t = task("partial");
3653 q.attach(&mut t, &[good]).unwrap();
3654 let more = source_file(dir.path(), "more.png", "ok");
3655 let missing = dir.path().join("missing.png");
3656 assert!(q.attach(&mut t, &[more, missing]).is_err());
3657 assert_eq!(t.attachments, ["good.png"]);
3658 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3659 .unwrap()
3660 .flatten()
3661 .collect();
3662 assert_eq!(on_disk.len(), 1);
3663 }
3664
3665 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3667 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3668 std::fs::create_dir_all(&tmp).unwrap();
3669 tmp
3670 }
3671
3672 #[test]
3673 fn a_failed_put_leaves_no_new_attachment_directory() {
3674 let (dir, q) = queue();
3675 let mut t = task("fresh");
3676 let tmp = block_put(&q, &t);
3677 let src = source_file(dir.path(), "shot.png", "x");
3678 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3679 assert!(t.attachments.is_empty());
3680 assert!(!q.attachments_dir(&t.id).exists());
3681 assert!(!q.path_of(&t.id).exists());
3682 std::fs::remove_dir(tmp).unwrap();
3683 }
3684
3685 #[test]
3686 fn a_failed_put_removes_only_the_copy_it_just_made() {
3687 let (dir, q) = queue();
3688 let mut t = task("edited");
3689 let first = source_file(dir.path(), "first.png", "1");
3690 q.attach_and_put(&mut t, &[first]).unwrap();
3691 block_put(&q, &t);
3692 let second = source_file(dir.path(), "second.png", "2");
3693 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3694 assert_eq!(t.attachments, ["first.png"]);
3695 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3696 .unwrap()
3697 .flatten()
3698 .map(|e| e.file_name().to_string_lossy().into_owned())
3699 .collect();
3700 assert_eq!(on_disk, ["first.png"]);
3701 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3702 }
3703
3704 #[test]
3705 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3706 let (dir, q) = queue();
3707 let questions = Questions::at(dir.path().join("questions"));
3708 let gone = task("gone");
3709 let mut other = task("other");
3710 let mut live = task("live");
3711 q.put(&mut other).unwrap();
3712 q.put(&mut live).unwrap();
3713 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3716 std::fs::create_dir_all(&orphan).unwrap();
3717 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3718 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3721 std::fs::create_dir_all(&busy).unwrap();
3722
3723 q.remove(&other.id, false, &questions).unwrap();
3724 assert!(!orphan.exists(), "an orphan is swept");
3725 assert!(busy.exists(), "a removal in progress is left alone");
3726 }
3727
3728 #[test]
3729 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3730 let (dir, q) = queue();
3731 let questions = Questions::at(dir.path().join("questions"));
3732 let mut t = task("stuck");
3733 let src = source_file(dir.path(), "shot.png", "x");
3734 q.attach_and_put(&mut t, &[src]).unwrap();
3735 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3738 std::fs::create_dir_all(&aside).unwrap();
3739 std::fs::write(aside.join("old.png"), "o").unwrap();
3740 assert!(q.remove(&t.id, false, &questions).is_err());
3741 assert!(q.path_of(&t.id).exists());
3742 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3743 }
3744
3745 #[test]
3746 fn a_failed_record_removal_puts_the_attachments_back() {
3747 let (dir, q) = queue();
3748 let mut t = task("rollback");
3749 let src = source_file(dir.path(), "shot.png", "x");
3750 q.attach_and_put(&mut t, &[src]).unwrap();
3751 let err = q
3752 .remove_record_with_attachments(&t.id, |_| {
3753 Err(std::io::Error::other("injected failure"))
3754 })
3755 .unwrap_err();
3756 assert!(format!("{err:#}").contains("injected failure"));
3757 assert!(q.path_of(&t.id).exists());
3758 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3759 assert!(
3760 !q.root
3761 .join(format!("{}.attachments.removing", t.id))
3762 .exists()
3763 );
3764 }
3765
3766 #[test]
3767 fn editing_a_task_keeps_its_attachments() {
3768 let (dir, q) = queue();
3769 let src = source_file(dir.path(), "shot.png", "x");
3770 let mut t = task("editable");
3771 q.attach(&mut t, &[src]).unwrap();
3772 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3773 q.put(&mut t).unwrap();
3774 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3775 }
3776
3777 #[test]
3778 fn removing_a_task_deletes_its_attachments() {
3779 let (dir, q) = queue();
3780 let questions = Questions::at(dir.path().join("questions"));
3781 let src = source_file(dir.path(), "shot.png", "x");
3782 let mut t = task("doomed");
3783 q.attach(&mut t, &[src]).unwrap();
3784 q.put(&mut t).unwrap();
3785 assert!(q.attachments_dir(&t.id).is_dir());
3786 q.remove(&t.id, false, &questions).unwrap();
3787 assert!(!q.attachments_dir(&t.id).exists());
3788 assert!(q.list().is_empty());
3789 }
3790
3791 #[test]
3792 fn a_task_written_before_attachments_still_reads() {
3793 let (_dir, q) = queue();
3794 let mut t = task("old");
3795 q.put(&mut t).unwrap();
3796 let path = q.path_of(&t.id);
3797 let mut v: serde_json::Value =
3798 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3799 v.as_object_mut().unwrap().remove("attachments");
3800 std::fs::write(&path, v.to_string()).unwrap();
3801 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3802 }
3803
3804 #[test]
3805 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3806 let q = Queue::at(PathBuf::from("relative-queue"));
3807 let mut t = task("rel");
3808 t.attachments.push("shot.png".to_owned());
3809 let paths = q.attachment_paths(&t);
3810 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3811 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3812 }
3813
3814 #[test]
3815 fn link_run_adds_a_run_once_and_touches_nothing_else() {
3816 let dir = tempfile::tempdir().unwrap();
3817 let queue = Queue::at(dir.path().join("queue"));
3818 let mut t = Task::new(
3819 "t".to_owned(),
3820 "do it".to_owned(),
3821 PathBuf::from("."),
3822 Source::Human,
3823 );
3824 queue.put(&mut t).unwrap();
3825 let before = queue.get(&t.id).unwrap();
3826
3827 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3828 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3829 assert_eq!(linked.status, before.status);
3830 assert_eq!(linked.attempts, before.attempts);
3831 assert_eq!(linked.interrupt, before.interrupt);
3832
3833 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3834 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3835 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3836 assert!(queue.link_run("no-such-task", "r").is_err());
3837 }
3838
3839 #[test]
3840 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
3841 let dir = tempfile::tempdir().unwrap();
3842 let queue = Queue::at(dir.path().join("queue"));
3843 let mut t = Task::new(
3844 "t".to_owned(),
3845 "do it".to_owned(),
3846 PathBuf::from("."),
3847 Source::Human,
3848 );
3849 queue.put(&mut t).unwrap();
3850 let mut snapshot = queue.get(&t.id).unwrap();
3852 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3853
3854 snapshot.start("20260930-000000-aaaa".to_owned());
3855 queue.put(&mut snapshot).unwrap();
3856
3857 let stored = queue.get(&t.id).unwrap();
3858 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3859 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3860 assert_eq!(stored.attempts, 1);
3861 }
3862
3863 #[test]
3864 fn concurrent_links_and_daemon_saves_lose_nothing() {
3865 let dir = tempfile::tempdir().unwrap();
3866 let queue = Queue::at(dir.path().join("queue"));
3867 let mut t = Task::new(
3868 "t".to_owned(),
3869 "do it".to_owned(),
3870 PathBuf::from("."),
3871 Source::Human,
3872 );
3873 queue.put(&mut t).unwrap();
3874 let id = t.id.clone();
3875
3876 let linkers: Vec<_> = (0..4)
3877 .map(|n| {
3878 let (queue, id) = (queue.clone(), id.clone());
3879 std::thread::spawn(move || {
3880 for k in 0..10 {
3881 queue
3882 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3883 .unwrap();
3884 }
3885 })
3886 })
3887 .collect();
3888 let mut mine = queue.get(&id).unwrap();
3891 for k in 0..10 {
3892 mine.start(format!("20260930-000009-d{k:03}"));
3893 queue.put(&mut mine).unwrap();
3894 }
3895 for l in linkers {
3896 l.join().unwrap();
3897 }
3898
3899 let stored = queue.get(&id).unwrap();
3900 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
3901 assert_eq!(
3902 stored.attempts, 10,
3903 "linking never rewinds the daemon's work"
3904 );
3905 }
3906
3907 fn age_lock(path: &Path) {
3908 let f = std::fs::OpenOptions::new().write(true).open(path).unwrap();
3909 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
3910 .unwrap();
3911 }
3912
3913 #[test]
3914 fn concurrent_stale_takeover_yields_one_holder() {
3915 use std::sync::atomic::{AtomicUsize, Ordering};
3916 use std::sync::{Arc, Barrier};
3917 for _ in 0..5 {
3918 let dir = tempfile::tempdir().unwrap();
3919 let q = Queue::at(dir.path().to_path_buf());
3920 let lock = dir.path().join("t.write-lock");
3921 std::fs::write(&lock, "dead-0000").unwrap();
3922 age_lock(&lock);
3923 let n = 6;
3924 let barrier = Arc::new(Barrier::new(n));
3925 let (now, max) = (Arc::new(AtomicUsize::new(0)), Arc::new(AtomicUsize::new(0)));
3926 let handles: Vec<_> = (0..n)
3927 .map(|_| {
3928 let (q, b, now, max) = (q.clone(), barrier.clone(), now.clone(), max.clone());
3929 std::thread::spawn(move || {
3930 b.wait();
3931 let g = q.lock_task("t").unwrap();
3932 let held = now.fetch_add(1, Ordering::SeqCst) + 1;
3933 max.fetch_max(held, Ordering::SeqCst);
3934 std::thread::sleep(std::time::Duration::from_millis(20));
3935 now.fetch_sub(1, Ordering::SeqCst);
3936 drop(g);
3937 })
3938 })
3939 .collect();
3940 for h in handles {
3941 h.join().unwrap();
3942 }
3943 assert_eq!(max.load(Ordering::SeqCst), 1);
3944 assert!(!lock.exists());
3945 }
3946 }
3947
3948 #[test]
3949 fn dropping_a_stolen_lock_leaves_the_new_holders_lock() {
3950 let dir = tempfile::tempdir().unwrap();
3951 let q = Queue::at(dir.path().to_path_buf());
3952 let lock = dir.path().join("t.write-lock");
3953 let a = q.lock_task("t").unwrap();
3954 age_lock(&lock);
3955 let b = q.lock_task("t").unwrap();
3956 assert_ne!(a.token, b.token);
3957 drop(a);
3958 assert_eq!(std::fs::read_to_string(&lock).unwrap(), b.token);
3959 drop(b);
3960 assert!(!lock.exists());
3961 }
3962
3963 #[test]
3964 fn break_stale_leaves_a_lock_that_replaced_the_one_judged() {
3965 let dir = tempfile::tempdir().unwrap();
3966 let lock = dir.path().join("t.write-lock");
3967 std::fs::write(&lock, "new-token").unwrap();
3968 assert!(!break_stale(&lock, "old-token"));
3969 assert_eq!(std::fs::read_to_string(&lock).unwrap(), "new-token");
3970 assert!(!break_stale(&lock, "new-token"));
3972 assert!(lock.exists());
3973 age_lock(&lock);
3974 assert!(break_stale(&lock, "new-token"));
3975 assert!(!lock.exists());
3976 }
3977
3978 #[test]
3979 fn a_stale_break_marker_is_recovered_and_a_fresh_one_is_respected() {
3980 let dir = tempfile::tempdir().unwrap();
3981 let lock = dir.path().join("t.write-lock");
3982 std::fs::write(&lock, "dead-1").unwrap();
3983 age_lock(&lock);
3984 let marker = break_marker(&lock, "dead-1");
3985 std::fs::write(&marker, "crashed-remover").unwrap();
3986 assert!(!break_stale(&lock, "dead-1"));
3988 assert!(lock.exists());
3989 age_lock(&marker);
3991 assert!(break_stale(&lock, "dead-1"));
3992 assert!(!lock.exists());
3993 assert!(!marker.exists());
3994 }
3995}