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 = 7;
90
91#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
93#[serde(rename_all = "lowercase")]
94pub enum HoldSource {
95 Manual,
97 Machine,
99}
100
101impl HoldSource {
102 pub fn label(self) -> &'static str {
104 match self {
105 Self::Manual => "manual",
106 Self::Machine => "machine",
107 }
108 }
109}
110
111#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
114#[serde(tag = "kind", rename_all = "lowercase")]
115pub enum Source {
116 Human,
118 Agent {
121 run: String,
123 node: String,
125 },
126 Issue {
128 number: u64,
130 repo: String,
132 },
133}
134
135impl Source {
136 pub fn label(&self) -> String {
138 match self {
139 Self::Human => "human".to_owned(),
140 Self::Agent { run, node } => format!("{node}@{}", short(run)),
141 Self::Issue { number, .. } => format!("issue #{number}"),
142 }
143 }
144}
145
146#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
148#[serde(rename_all = "lowercase")]
149pub enum TaskStatus {
150 Queued,
152 Running,
154 Done,
156 Failed,
158 Held,
160 Blocked,
164}
165
166impl TaskStatus {
167 pub fn runnable(self) -> bool {
169 matches!(self, Self::Queued | Self::Failed)
170 }
171
172 pub fn as_str(self) -> &'static str {
174 match self {
175 Self::Queued => "queued",
176 Self::Running => "running",
177 Self::Done => "done",
178 Self::Failed => "failed",
179 Self::Held => "held",
180 Self::Blocked => "blocked",
181 }
182 }
183}
184
185#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
188pub struct TaskCounts {
189 pub queued: usize,
191 pub running: usize,
193 pub done: usize,
195 pub failed: usize,
197 pub held: usize,
199 pub blocked: usize,
201}
202
203impl TaskCounts {
204 pub fn of(tasks: &[Task]) -> Self {
208 let mut counts = Self::default();
209 for t in tasks {
210 match t.status {
211 TaskStatus::Queued => counts.queued += 1,
212 TaskStatus::Running => counts.running += 1,
213 TaskStatus::Done => counts.done += 1,
214 TaskStatus::Failed => counts.failed += 1,
215 TaskStatus::Held => counts.held += 1,
216 TaskStatus::Blocked => counts.blocked += 1,
217 }
218 }
219 counts
220 }
221}
222
223#[derive(Debug, Clone, Serialize, Deserialize)]
225#[serde(deny_unknown_fields)]
226pub struct Task {
227 pub schema: u32,
229 pub id: String,
231 pub title: String,
233 pub instruction: String,
235 pub repo: PathBuf,
237 pub source: Source,
239 #[serde(default)]
241 pub priority: i32,
242 #[serde(default)]
253 pub solo: bool,
254 pub status: TaskStatus,
256 #[serde(default)]
258 pub attempts: usize,
259 #[serde(default)]
261 pub runs: Vec<String>,
262 #[serde(default)]
264 pub last_error: Option<String>,
265 #[serde(default)]
278 pub hold_reason: Option<String>,
279 #[serde(default)]
282 pub hold_source: Option<HoldSource>,
283 #[serde(default)]
295 pub diagnostic: Option<String>,
296 #[serde(default)]
306 pub blocked_by: Vec<String>,
307 #[serde(default)]
310 pub block_reason: Option<String>,
311 #[serde(default)]
324 pub blocked_from: Option<TaskStatus>,
325 #[serde(default)]
336 pub answers: Vec<AnsweredQuestion>,
337 #[serde(default)]
343 pub triage_applied: Vec<String>,
344 #[serde(default)]
350 pub actions_applied: Vec<String>,
351 #[serde(default)]
356 pub resume_override: Option<OperatorResume>,
357 #[serde(default)]
369 pub review_branch: Option<String>,
370 #[serde(default)]
373 pub fresh_start: bool,
374 #[serde(default)]
385 pub interrupt: bool,
386 #[serde(default)]
409 pub urgent: bool,
410 #[serde(default)]
416 pub attachments: Vec<String>,
417 pub created_at: Timestamp,
419 pub updated_at: Timestamp,
421}
422
423#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
426pub struct OperatorResume {
427 pub question_id: String,
429 pub at: Timestamp,
431 #[serde(default)]
434 pub conductor_rehold: Option<String>,
435 #[serde(default)]
438 pub forced: bool,
439 #[serde(default)]
446 pub pinned_run: Option<String>,
447}
448
449#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
452pub struct AnsweredQuestion {
453 pub question: String,
455 pub answer: String,
457}
458
459impl Task {
460 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
462 let now = Timestamp::now();
463 Self {
464 schema: SCHEMA,
465 id: new_id(),
466 title,
467 instruction,
468 repo,
469 source,
470 priority: 0,
471 solo: false,
472 status: TaskStatus::Queued,
473 attempts: 0,
474 runs: Vec::new(),
475 last_error: None,
476 hold_reason: None,
477 hold_source: None,
478 diagnostic: None,
479 blocked_by: Vec::new(),
480 block_reason: None,
481 blocked_from: None,
482 answers: Vec::new(),
483 triage_applied: Vec::new(),
484 actions_applied: Vec::new(),
485 resume_override: None,
486 review_branch: None,
487 fresh_start: false,
488 interrupt: false,
489 urgent: false,
490 attachments: Vec::new(),
491 created_at: now,
492 updated_at: now,
493 }
494 }
495
496 pub fn short(&self) -> &str {
498 short(&self.id)
499 }
500
501 pub fn mark_triage_applied(&mut self, question_id: &str) {
503 if !self.triage_applied(question_id) {
504 self.triage_applied.push(question_id.to_owned());
505 }
506 }
507
508 pub fn action_applied(&self, question_id: &str) -> bool {
511 self.actions_applied.iter().any(|id| id == question_id)
512 }
513
514 pub fn mark_action_applied(&mut self, question_id: &str) {
516 if !self.action_applied(question_id) {
517 self.actions_applied.push(question_id.to_owned());
518 }
519 }
520
521 pub fn triage_applied(&self, question_id: &str) -> bool {
523 self.triage_applied.iter().any(|id| id == question_id)
524 }
525
526 pub fn start(&mut self, run: String) {
537 self.status = TaskStatus::Running;
538 self.attempts += 1;
539 self.runs.push(run);
540 self.last_error = None;
541 self.fresh_start = false;
542 self.interrupt = false;
543 self.resume_override = None;
545 }
546
547 pub fn succeed(&mut self) {
557 self.status = TaskStatus::Done;
558 self.resume_override = None;
559 self.last_error = None;
560 self.hold_reason = None;
561 self.hold_source = None;
562 self.diagnostic = None;
563 self.blocked_by.clear();
564 self.block_reason = None;
565 self.blocked_from = None;
566 }
567
568 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
598 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
599 return &[];
600 }
601 &self.runs[..self.runs.len() - 1]
602 }
603
604 pub fn successor_of(&self, run: &str) -> Option<&String> {
611 let pos = self.runs.iter().position(|r| r == run)?;
612 self.runs.get(pos + 1)
613 }
614
615 pub fn earlier_attempts(&self) -> &[String] {
622 &self.runs
623 }
624
625 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
641 let why = why.into();
642 self.diagnostic = None;
643 self.status = if self.attempts >= max_attempts {
644 self.hold_source = Some(HoldSource::Machine);
645 self.hold_reason = Some(why.clone());
646 TaskStatus::Held
647 } else {
648 TaskStatus::Failed
649 };
650 self.last_error = Some(why);
651 }
652
653 pub fn stall(&mut self, why: impl Into<String>) {
663 self.last_error = Some(why.into());
664 self.diagnostic = None;
665 self.attempts = self.attempts.saturating_sub(1);
666 self.status = TaskStatus::Failed;
667 }
668
669 pub fn operator_held(&self) -> bool {
675 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
676 }
677
678 pub fn hold_manual(&mut self, reason: Option<String>) {
688 self.status = TaskStatus::Held;
689 if reason.is_some() {
690 self.hold_reason = reason;
691 }
692 self.hold_source = Some(HoldSource::Manual);
693 self.blocked_by.clear();
694 self.block_reason = None;
695 self.blocked_from = None;
696 }
697
698 pub fn hold_machine(&mut self, reason: Option<String>) {
703 self.status = TaskStatus::Held;
704 if reason.is_some() {
705 self.hold_reason = reason;
706 }
707 self.hold_source = Some(HoldSource::Machine);
708 self.blocked_by.clear();
709 self.block_reason = None;
710 self.blocked_from = None;
711 }
712
713 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
723 if self.status != TaskStatus::Blocked {
724 self.blocked_from = Some(self.status);
725 }
726 self.status = TaskStatus::Blocked;
727 self.blocked_by = blocked_by;
728 self.block_reason = reason;
729 }
730
731 pub fn unblock(&mut self, resolved_id: &str) {
754 if self.status != TaskStatus::Blocked {
755 return;
756 }
757 self.blocked_by.retain(|id| id != resolved_id);
758 if self.blocked_by.is_empty() {
759 self.status = match self.blocked_from {
760 Some(TaskStatus::Running) => TaskStatus::Queued,
761 Some(other) => other,
762 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
763 TaskStatus::Held
764 }
765 None => TaskStatus::Queued,
766 };
767 self.block_reason = None;
768 self.blocked_from = None;
769 }
770 }
771
772 pub fn record_answer(&mut self, question: String, answer: String) {
777 self.answers.push(AnsweredQuestion { question, answer });
778 }
779
780 pub fn request_review(&mut self, branch: String) {
784 self.release();
785 self.review_branch = Some(branch);
786 }
787
788 pub fn requeue(&mut self) {
791 self.release();
792 self.fresh_start = true;
793 }
794
795 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
804 if self.status == TaskStatus::Running {
805 bail!(
806 "task {} is running; its priority cannot be changed until \
807 this attempt finishes",
808 self.short()
809 );
810 }
811 self.priority = priority;
812 Ok(())
813 }
814
815 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
830 if interrupt && !self.status.runnable() {
831 bail!(
832 "task {} is {}; only a queued or failed task can be marked \
833 to interrupt",
834 self.short(),
835 self.status.as_str()
836 );
837 }
838 self.interrupt = interrupt;
839 Ok(())
840 }
841
842 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
854 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
855 bail!(
856 "task {} is {}; only a queued or held task's instruction can \
857 be edited",
858 self.short(),
859 self.status.as_str()
860 );
861 }
862 self.title = title;
863 self.instruction = instruction;
864 Ok(())
865 }
866
867 pub fn handed_off(&mut self, why: impl Into<String>) {
887 let why = why.into();
888 self.diagnostic = None;
889 self.status = TaskStatus::Held;
890 self.hold_source = Some(HoldSource::Machine);
891 self.hold_reason = Some(why.clone());
892 self.last_error = Some(why);
893 }
894
895 pub fn release(&mut self) {
899 self.status = TaskStatus::Queued;
900 self.attempts = 0;
901 self.last_error = None;
902 self.hold_reason = None;
905 self.hold_source = None;
906 self.diagnostic = None;
907 self.blocked_by.clear();
912 self.block_reason = None;
913 self.blocked_from = None;
914 self.review_branch = None;
915 self.fresh_start = false;
916 }
917}
918
919fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
923 use std::io::ErrorKind;
924 let (stem, ext) = match name.rfind('.') {
925 Some(i) if i > 0 => (&name[..i], &name[i..]),
926 _ => (name, ""),
927 };
928 for n in 1u32.. {
929 let candidate = if n == 1 {
930 name.to_owned()
931 } else {
932 let suffix = format!("-{n}");
933 let room = 64usize.saturating_sub(suffix.len() + ext.len());
934 let stem: String = stem.chars().take(room).collect();
935 format!("{stem}{suffix}{ext}")
936 };
937 if !crate::ask::valid_asset_name(&candidate) {
938 bail!("no valid attachment name is left for `{name}`");
939 }
940 let path = dir.join(&candidate);
941 match std::fs::OpenOptions::new()
942 .write(true)
943 .create_new(true)
944 .open(&path)
945 {
946 Ok(mut out) => {
947 let copied = std::fs::File::open(src)
948 .and_then(|mut input| std::io::copy(&mut input, &mut out));
949 if let Err(e) = copied {
950 drop(out);
951 let _ = std::fs::remove_file(&path);
952 return Err(e).with_context(|| format!("copy {}", src.display()));
953 }
954 return Ok((candidate, path));
955 }
956 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
957 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
958 }
959 }
960 unreachable!("the counter never runs out")
961}
962
963#[derive(Debug, Clone)]
965pub struct Queue {
966 root: PathBuf,
967}
968
969impl Queue {
970 pub fn open() -> Self {
972 Self::at(crate::run::home().join("queue"))
973 }
974
975 pub fn at(root: PathBuf) -> Self {
978 Self { root }
979 }
980
981 pub fn root(&self) -> &Path {
983 &self.root
984 }
985
986 pub fn path_of(&self, id: &str) -> PathBuf {
988 self.root.join(format!("{id}.json"))
989 }
990
991 pub fn attachments_dir(&self, id: &str) -> PathBuf {
994 self.root.join(format!("{id}.attachments"))
995 }
996
997 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1008 let mut wanted = Vec::new();
1009 for src in sources {
1010 let name = src
1011 .file_name()
1012 .and_then(|n| n.to_str())
1013 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1014 if !crate::ask::valid_asset_name(name) {
1015 bail!(
1016 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1017 with no `..`; rename the file and try again"
1018 );
1019 }
1020 if !src.is_file() {
1021 bail!("attachment `{}` is not a file", src.display());
1022 }
1023 wanted.push((src, name));
1024 }
1025 let dir = self.attachments_dir(&task.id);
1026 let existed = dir.is_dir();
1027 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1028 let mut created: Vec<PathBuf> = Vec::new();
1029 let mut names = Vec::new();
1030 let mut copy_all = || -> Result<()> {
1031 for (src, name) in &wanted {
1032 let (stored, path) = copy_new(&dir, src, name)?;
1033 created.push(path);
1034 names.push(stored);
1035 }
1036 Ok(())
1037 };
1038 if let Err(e) = copy_all() {
1039 for path in &created {
1040 let _ = std::fs::remove_file(path);
1041 }
1042 if !existed {
1043 let _ = std::fs::remove_dir(&dir);
1044 }
1045 return Err(e);
1046 }
1047 task.attachments.extend(names.iter().cloned());
1048 Ok(names)
1049 }
1050
1051 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1062 let dir = self.attachments_dir(&task.id);
1063 let existed = dir.is_dir();
1064 let before = task.attachments.len();
1065 let names = self.attach(task, sources)?;
1066 if let Err(e) = self.put(task) {
1067 for name in &names {
1068 let _ = std::fs::remove_file(dir.join(name));
1069 }
1070 if !existed {
1071 let _ = std::fs::remove_dir(&dir);
1072 }
1073 task.attachments.truncate(before);
1074 return Err(e);
1075 }
1076 Ok(names)
1077 }
1078
1079 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1083 let dir = self.attachments_dir(&task.id);
1084 task.attachments
1085 .iter()
1086 .map(|n| {
1087 let p = dir.join(n);
1088 std::path::absolute(&p).unwrap_or(p)
1089 })
1090 .collect()
1091 }
1092
1093 pub fn put(&self, task: &mut Task) -> Result<()> {
1096 task.updated_at = Timestamp::now();
1097 std::fs::create_dir_all(&self.root)
1098 .with_context(|| format!("create {}", self.root.display()))?;
1099 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1100 let path = self.path_of(&task.id);
1101 let tmp = path.with_extension("json.tmp");
1102 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1103 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1104 if let (Some(notice), Some(home)) = (
1109 crate::notices::task_held(task),
1110 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1111 ) {
1112 crate::notices::raise_in(home, notice);
1113 }
1114 Ok(())
1115 }
1116
1117 pub fn get(&self, id: &str) -> Result<Task> {
1119 let resolved = self.resolve_id(id)?;
1120 read_path(&self.path_of(&resolved))
1121 }
1122
1123 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1151 let resolved = self.resolve_id(id)?;
1152 if in_flight {
1153 bail!("task {resolved} is being run by a live daemon right now");
1154 }
1155 self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))?;
1156 let quarantined = self.quarantine_dependents_of(&resolved, questions);
1157 Ok(Removal {
1158 id: resolved,
1159 quarantined,
1160 })
1161 }
1162
1163 fn remove_record_with_attachments(
1167 &self,
1168 resolved: &str,
1169 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1170 ) -> Result<()> {
1171 self.sweep_removed_attachments();
1178 let attachments = self.attachments_dir(resolved);
1179 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1180 let moved = match std::fs::rename(&attachments, &aside) {
1181 Ok(()) => true,
1182 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1183 Err(e) => {
1184 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1185 }
1186 };
1187 let path = self.path_of(resolved);
1188 if let Err(e) = remove_record(&path) {
1189 if moved {
1190 let _ = std::fs::rename(&aside, &attachments);
1191 }
1192 return Err(e).with_context(|| format!("remove {}", path.display()));
1193 }
1194 if moved {
1195 if let Err(e) = std::fs::remove_dir_all(&aside) {
1196 tracing::warn!("leftover attachments {}: {e}", aside.display());
1197 }
1198 }
1199 let lock = self.lock_path(resolved);
1200 if let Err(e) = std::fs::remove_file(&lock) {
1201 if e.kind() != std::io::ErrorKind::NotFound {
1202 return Err(e).with_context(|| format!("remove {}", lock.display()));
1203 }
1204 }
1205 Ok(())
1206 }
1207
1208 fn sweep_removed_attachments(&self) {
1214 let Ok(entries) = std::fs::read_dir(&self.root) else {
1215 return;
1216 };
1217 for entry in entries.flatten() {
1218 let name = entry.file_name();
1219 let name = name.to_string_lossy();
1220 let Some(id) = name.strip_suffix(".attachments.removing") else {
1221 continue;
1222 };
1223 if !self.path_of(id).exists() {
1226 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1227 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1228 }
1229 }
1230 }
1231 }
1232
1233 fn quarantine_dependents_of(&self, dependency: &str, questions: &Questions) -> Vec<String> {
1237 let mut quarantined = Vec::new();
1238 for listed in self.list() {
1239 if listed.status != TaskStatus::Blocked
1240 || !listed.blocked_by.iter().any(|b| b == dependency)
1241 {
1242 continue;
1243 }
1244 let Ok(_claim) = self.claim(&listed.id) else {
1245 continue;
1246 };
1247 let Ok(mut task) = self.get(&listed.id) else {
1248 continue;
1249 };
1250 if task.status != TaskStatus::Blocked
1251 || !task.blocked_by.iter().any(|b| b == dependency)
1252 {
1253 continue;
1254 }
1255 let missing = missing_blockers(self, questions, &task.blocked_by);
1256 let language = crate::lang::of_repo(&task.repo);
1257 task.hold_machine(Some(missing_blocker_hold_reason_in(
1258 &task.blocked_by,
1259 &missing,
1260 &language,
1261 )));
1262 if self.put(&mut task).is_ok() {
1263 quarantined.push(task.id.clone());
1264 }
1265 }
1266 quarantined
1267 }
1268
1269 fn lock_path(&self, id: &str) -> PathBuf {
1272 self.root.join(format!("{id}.lock"))
1273 }
1274
1275 pub fn list(&self) -> Vec<Task> {
1288 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1289 .into_iter()
1290 .flatten()
1291 .flatten()
1292 .map(|e| e.path())
1293 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1294 .filter_map(|p| read_path(&p).ok())
1295 .collect();
1296 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1297 tasks
1298 }
1299
1300 pub fn superseded(&self) -> HashMap<String, String> {
1313 let mut by = HashMap::new();
1314 for task in self.list() {
1315 for earlier in &task.runs {
1316 if let Some(later) = task.successor_of(earlier) {
1317 by.insert(earlier.clone(), later.clone());
1318 }
1319 }
1320 }
1321 by
1322 }
1323
1324 pub fn superseded_by(&self, run: &str) -> Option<String> {
1332 for task in self.list() {
1333 if task.runs.iter().any(|r| r == run) {
1334 return task.successor_of(run).cloned();
1335 }
1336 }
1337 None
1338 }
1339
1340 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1351 for task in self.list() {
1352 if task.runs.iter().any(|r| r == run) {
1353 return task.runs.last().filter(|last| **last != run).cloned();
1354 }
1355 }
1356 None
1357 }
1358
1359 pub fn next_runnable(&self) -> Option<Task> {
1364 let mut runnable: Vec<Task> = self
1365 .list()
1366 .into_iter()
1367 .filter(|t| t.status.runnable())
1368 .collect();
1369 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1370 runnable.into_iter().next()
1371 }
1372
1373 pub fn claim(&self, id: &str) -> Result<Claim> {
1380 std::fs::create_dir_all(&self.root)
1381 .with_context(|| format!("create {}", self.root.display()))?;
1382 let path = self.lock_path(id);
1383 match std::fs::OpenOptions::new()
1384 .write(true)
1385 .create_new(true)
1386 .open(&path)
1387 {
1388 Ok(mut f) => {
1389 use std::io::Write as _;
1390 let _ = writeln!(f, "{}", std::process::id());
1392 Ok(Claim { path })
1393 }
1394 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1395 bail!("task {id} is already claimed ({} exists)", path.display())
1396 }
1397 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1398 }
1399 }
1400
1401 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1403 if self.path_of(prefix).is_file() {
1404 return Ok(prefix.to_owned());
1405 }
1406 let hits: Vec<String> = self
1407 .list()
1408 .into_iter()
1409 .map(|t| t.id)
1410 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1411 .collect();
1412 match hits.len() {
1413 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1414 0 => bail!("no task matches `{prefix}`"),
1415 _ => bail!(
1416 "`{prefix}` matches {} tasks: {}",
1417 hits.len(),
1418 hits.join(", ")
1419 ),
1420 }
1421 }
1422
1423 pub fn revision(&self) -> u64 {
1430 use std::hash::{Hash as _, Hasher as _};
1431
1432 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
1433 .into_iter()
1434 .flatten()
1435 .flatten()
1436 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
1437 .filter_map(|e| {
1438 let name = e.file_name().to_string_lossy().into_owned();
1439 let mtime = e
1440 .metadata()
1441 .ok()?
1442 .modified()
1443 .ok()?
1444 .duration_since(std::time::UNIX_EPOCH)
1445 .ok()?
1446 .as_millis() as u64;
1447 Some((name, mtime))
1448 })
1449 .collect();
1450
1451 if entries.is_empty() {
1452 return 0;
1453 }
1454
1455 entries.sort_unstable();
1456 let mut hasher = std::hash::DefaultHasher::new();
1457 for (name, mtime) in &entries {
1458 name.hash(&mut hasher);
1459 mtime.hash(&mut hasher);
1460 }
1461 let h = hasher.finish();
1462 if h == 0 { 1 } else { h }
1463 }
1464}
1465
1466#[derive(Debug, Clone)]
1468pub struct Removal {
1469 pub id: String,
1471 pub quarantined: Vec<String>,
1475}
1476
1477#[derive(Debug)]
1479pub struct Claim {
1480 path: PathBuf,
1481}
1482
1483impl Drop for Claim {
1484 fn drop(&mut self) {
1485 let _ = std::fs::remove_file(&self.path);
1486 }
1487}
1488
1489pub fn title_from(instruction: &str, max: usize) -> String {
1492 let line = instruction
1498 .lines()
1499 .map(str::trim)
1500 .find(|l| !l.is_empty())
1501 .unwrap_or("(empty task)")
1502 .trim_start_matches(['#', '-', '*', '>', ' '])
1503 .trim();
1504 if line.is_empty() {
1505 return "(empty task)".to_owned();
1506 }
1507 if line.chars().count() <= max {
1508 return line.to_owned();
1509 }
1510 let head: String = line.chars().take(max.saturating_sub(1)).collect();
1511 format!("{head}…")
1512}
1513
1514fn read_path(path: &Path) -> Result<Task> {
1515 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1516 let task: Task =
1517 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1518 if task.schema > SCHEMA {
1524 bail!(
1525 "task {} was written by a different magi (schema {}, this build \
1526 speaks {SCHEMA})",
1527 task.id,
1528 task.schema
1529 );
1530 }
1531 Ok(task)
1532}
1533
1534pub fn missing_blockers(
1548 queue: &Queue,
1549 questions: &Questions,
1550 blocked_by: &[String],
1551) -> Vec<String> {
1552 blocked_by
1553 .iter()
1554 .filter(|id| !queue.path_of(id).is_file() && !questions.path_of(id).is_file())
1555 .cloned()
1556 .collect()
1557}
1558
1559pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
1571 missing_blocker_hold_reason_in(blocked_by, missing, "en")
1572}
1573
1574pub fn missing_blocker_hold_reason_in(
1576 blocked_by: &[String],
1577 missing: &[String],
1578 language: &str,
1579) -> String {
1580 if crate::lang::is_japanese(language) {
1581 format!(
1582 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
1583 blocked_by.join(", "),
1584 missing.join(", "),
1585 )
1586 } else {
1587 format!(
1588 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
1589 blocked_by.join(", "),
1590 missing.join(", "),
1591 )
1592 }
1593}
1594
1595pub fn short(id: &str) -> &str {
1597 id.split('-').next_back().unwrap_or(id)
1598}
1599
1600fn new_id() -> String {
1601 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1602 let seed = crate::rng::entropy();
1603 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1604}
1605
1606#[cfg(test)]
1607mod tests {
1608 #[test]
1609 fn missing_blocker_reason_follows_the_language() {
1610 let b = vec!["a".to_owned()];
1611 let en = missing_blocker_hold_reason_in(&b, &b, "en");
1612 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
1613 assert!(en.starts_with("blocked on a"));
1614 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
1615 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
1616 }
1617
1618 use super::*;
1619
1620 #[test]
1621 fn triage_applied_survives_release_and_old_records_read_as_empty() {
1622 let mut t = Task::new(
1623 "t".to_owned(),
1624 "i".to_owned(),
1625 PathBuf::from("r"),
1626 Source::Human,
1627 );
1628 t.mark_triage_applied("q1");
1629 t.mark_triage_applied("q1");
1630 t.hold_machine(Some("x".to_owned()));
1631 t.release();
1632 assert_eq!(t.triage_applied, ["q1"]);
1633 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
1634
1635 let mut v = serde_json::to_value(&t).unwrap();
1636 v.as_object_mut().unwrap().remove("triage_applied");
1637 let old: Task = serde_json::from_value(v).unwrap();
1638 assert!(old.triage_applied.is_empty());
1639 }
1640
1641 #[test]
1642 fn task_counts_of_empty_is_all_zero() {
1643 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
1644 }
1645
1646 #[test]
1647 fn task_counts_of_tallies_every_status() {
1648 let mut queued = Task::new(
1649 "q".to_owned(),
1650 "i".to_owned(),
1651 PathBuf::from("."),
1652 Source::Human,
1653 );
1654 queued.status = TaskStatus::Queued;
1655 let mut running = queued.clone();
1656 running.status = TaskStatus::Running;
1657 let mut done = queued.clone();
1658 done.status = TaskStatus::Done;
1659 let mut failed = queued.clone();
1660 failed.status = TaskStatus::Failed;
1661 let mut held = queued.clone();
1662 held.status = TaskStatus::Held;
1663 let mut blocked = queued.clone();
1664 blocked.status = TaskStatus::Blocked;
1665
1666 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
1667 assert_eq!(
1668 counts,
1669 TaskCounts {
1670 queued: 1,
1671 running: 1,
1672 done: 2,
1673 failed: 1,
1674 held: 1,
1675 blocked: 1,
1676 }
1677 );
1678 }
1679
1680 fn queue() -> (tempfile::TempDir, Queue) {
1683 let dir = tempfile::tempdir().unwrap();
1684 let q = Queue::at(dir.path().join("queue"));
1685 (dir, q)
1686 }
1687
1688 #[test]
1689 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
1690 let dir = tempfile::tempdir().unwrap();
1691 let q = Queue::at(dir.path().join("queue"));
1692 let mut t = task("held");
1693 q.put(&mut t).unwrap();
1694 assert_eq!(
1695 crate::notices::Notices::at(dir.path().join("notifications"))
1696 .list()
1697 .len(),
1698 0
1699 );
1700 t.hold_machine(Some("out of attempts".to_owned()));
1701 q.put(&mut t).unwrap();
1702 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
1703 assert_eq!(listed.len(), 1);
1704 assert!(listed[0].message.contains("out of attempts"));
1705 }
1706
1707 fn task(title: &str) -> Task {
1708 Task::new(
1709 title.to_owned(),
1710 format!("do {title}"),
1711 PathBuf::from("."),
1712 Source::Human,
1713 )
1714 }
1715
1716 #[test]
1717 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
1718 let mut t = task("retried");
1719 assert!(
1720 t.earlier_attempts().is_empty(),
1721 "a first attempt takes nothing over"
1722 );
1723 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1724 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
1725 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
1727 assert_eq!(t.successor_of("bbbb"), None);
1728 assert_eq!(t.successor_of("zzzz"), None);
1729 }
1730
1731 #[test]
1732 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
1733 let (_dir, q) = queue();
1734 let mut t = task("retried");
1735 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1736 q.put(&mut t).unwrap();
1737
1738 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
1739 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
1740 assert_eq!(
1741 q.superseded_by("cccc"),
1742 None,
1743 "the latest attempt replaces nothing"
1744 );
1745 assert_eq!(
1746 q.superseded_by("never-heard-of-it"),
1747 None,
1748 "a run belonging to no task on this queue is not superseded"
1749 );
1750
1751 let mut by = HashMap::new();
1752 by.insert("aaaa".to_owned(), "bbbb".to_owned());
1753 by.insert("bbbb".to_owned(), "cccc".to_owned());
1754 assert_eq!(
1755 q.superseded(),
1756 by,
1757 "the whole-map and single-run forms must agree"
1758 );
1759 }
1760
1761 #[test]
1762 fn superseded_attempts_is_empty_until_the_task_is_done() {
1763 let mut t = task("retried");
1764 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1765 t.status = TaskStatus::Failed;
1766 assert_eq!(
1767 t.superseded_attempts(true),
1768 &[] as &[String],
1769 "a task still retrying has no attempt yet that a later one made moot"
1770 );
1771
1772 t.status = TaskStatus::Running;
1773 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1774 }
1775
1776 #[test]
1777 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
1778 let mut t = task("retried");
1779 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1780 t.status = TaskStatus::Done;
1781 assert_eq!(
1782 t.superseded_attempts(true),
1783 &["aaaa".to_owned(), "bbbb".to_owned()],
1784 "cccc is the attempt whose success made the task done, and stays out"
1785 );
1786 }
1787
1788 #[test]
1789 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
1790 let mut t = task("first try landed");
1791 t.runs = vec!["aaaa".to_owned()];
1792 t.status = TaskStatus::Done;
1793 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1794 }
1795
1796 #[test]
1797 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
1798 let mut t = task("closed by hand after a manual merge");
1805 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1806 t.status = TaskStatus::Done;
1807 assert_eq!(
1808 t.superseded_attempts(false),
1809 &[] as &[String],
1810 "nothing here is provably why the task is done, so nothing is superseded"
1811 );
1812 }
1813
1814 #[test]
1815 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
1816 let (_dir, q) = queue();
1817 let mut t = task("retried twice");
1818 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1819 q.put(&mut t).unwrap();
1820
1821 assert_eq!(
1822 q.latest_attempt("aaaa"),
1823 Some("cccc".to_owned()),
1824 "an old attempt points straight at the chain's current head, not the \
1825 next attempt in the middle of it"
1826 );
1827 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
1828 assert_eq!(
1829 q.latest_attempt("cccc"),
1830 None,
1831 "the latest attempt is not superseded by anything"
1832 );
1833 assert_eq!(
1834 q.latest_attempt("never-heard-of-it"),
1835 None,
1836 "a run belonging to no task on this queue is not superseded"
1837 );
1838 }
1839
1840 #[test]
1841 fn a_markdown_heading_is_the_title_not_decoration() {
1842 assert_eq!(
1847 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
1848 "Rework the config loader"
1849 );
1850 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
1851 assert_eq!(title_from("> quoted task", 40), "quoted task");
1852 assert_eq!(title_from(" \n\n", 40), "(empty task)");
1854 assert_eq!(title_from("###\n", 40), "(empty task)");
1855 }
1856
1857 #[test]
1858 fn a_long_title_is_elided_by_characters_not_bytes() {
1859 let long = "課題".repeat(30);
1861 let title = title_from(&long, 10);
1862 assert_eq!(title.chars().count(), 10);
1863 assert!(title.ends_with('…'));
1864 }
1865
1866 #[test]
1867 fn priority_wins_and_ties_break_oldest_first() {
1868 let (_dir, q) = queue();
1869 let mut a = task("first");
1870 let mut b = task("second");
1871 let mut c = task("urgent");
1872 a.id = "20260101-000001-aaaa".to_owned();
1874 b.id = "20260101-000002-bbbb".to_owned();
1875 c.id = "20260101-000003-cccc".to_owned();
1876 c.priority = 5;
1877 for t in [&mut a, &mut b, &mut c] {
1878 q.put(t).unwrap();
1879 }
1880
1881 assert_eq!(q.next_runnable().unwrap().id, c.id);
1883 c.hold_machine(None);
1884 q.put(&mut c).unwrap();
1885 assert_eq!(q.next_runnable().unwrap().id, a.id);
1887 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
1888 }
1889
1890 #[test]
1891 fn a_blocked_task_never_starves_another_runnable_one() {
1892 let (_dir, q) = queue();
1893 let mut blocked = task("blocked");
1894 blocked.block(vec!["something".to_owned()], None);
1895 q.put(&mut blocked).unwrap();
1896
1897 let mut runnable = task("free to go");
1898 q.put(&mut runnable).unwrap();
1899
1900 let next = q.next_runnable().expect("a runnable task is still offered");
1901 assert_eq!(next.id, runnable.id);
1902 }
1903
1904 #[test]
1905 fn a_held_task_is_never_offered_to_the_loop() {
1906 let (_dir, q) = queue();
1907 let mut t = task("held");
1908 q.put(&mut t).unwrap();
1909 assert!(q.next_runnable().is_some());
1910
1911 t.hold_machine(None);
1912 q.put(&mut t).unwrap();
1913 assert!(
1914 q.next_runnable().is_none(),
1915 "a held task must wait for a human"
1916 );
1917
1918 t.status = TaskStatus::Failed;
1920 q.put(&mut t).unwrap();
1921 assert!(q.next_runnable().is_some());
1922 }
1923
1924 #[test]
1925 fn attempts_are_capped_and_then_the_task_is_held() {
1926 let mut t = task("doomed");
1927
1928 t.start("run-1".to_owned());
1929 t.fail("gate red", 2);
1930 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
1931
1932 t.start("run-2".to_owned());
1933 t.fail("gate red", 2);
1934 assert_eq!(
1935 t.status,
1936 TaskStatus::Held,
1937 "out of attempts: stop spending money on it"
1938 );
1939 assert_eq!(t.runs, ["run-1", "run-2"]);
1940 assert_eq!(t.last_error.as_deref(), Some("gate red"));
1941 assert_eq!(
1942 t.hold_reason.as_deref(),
1943 Some("gate red"),
1944 "the hold must say why, not leave hold_reason null next to a \
1945 populated last_error"
1946 );
1947 }
1948
1949 #[test]
1950 fn handing_off_a_task_records_a_hold_reason_too() {
1951 let mut t = task("left a pull request");
1952 t.start("run-1".to_owned());
1953 t.handed_off("run ended with a pull request open [run run-1]");
1954 assert_eq!(t.status, TaskStatus::Held);
1955 assert_eq!(t.hold_source, Some(HoldSource::Machine));
1956 assert_eq!(
1957 t.hold_reason.as_deref(),
1958 Some("run ended with a pull request open [run run-1]")
1959 );
1960 assert_eq!(t.hold_reason, t.last_error);
1961 }
1962
1963 #[test]
1964 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
1965 let mut t = task("stalled by quota");
1966
1967 t.start("run-1".to_owned());
1968 assert_eq!(t.attempts, 1);
1969 t.stall("judge-1, judge-2 out of quota");
1970 assert_eq!(
1971 t.attempts, 0,
1972 "a closed quota window must not spend the task's retry budget"
1973 );
1974 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
1975 assert_eq!(
1976 t.last_error.as_deref(),
1977 Some("judge-1, judge-2 out of quota")
1978 );
1979
1980 for _ in 0..20 {
1983 t.start("run-n".to_owned());
1984 t.stall("still out of quota");
1985 }
1986 t.start("run-real".to_owned());
1987 t.fail("gate red", 2);
1988 assert_eq!(
1989 t.status,
1990 TaskStatus::Failed,
1991 "the first attempt that was really judged is attempt one"
1992 );
1993 }
1994
1995 #[test]
1996 fn releasing_a_held_task_gives_it_a_real_second_chance() {
1997 let mut t = task("retry me");
1998 t.start("run-1".to_owned());
1999 t.fail("gate red", 1);
2000 assert_eq!(t.status, TaskStatus::Held);
2001
2002 t.release();
2003 assert_eq!(t.status, TaskStatus::Queued);
2004 assert_eq!(t.attempts, 0);
2007 assert!(t.last_error.is_none());
2008 assert_eq!(
2009 t.runs.len(),
2010 1,
2011 "history is kept: attempts reset, evidence does not"
2012 );
2013 }
2014
2015 #[test]
2016 fn a_hold_reason_survives_and_a_release_clears_it() {
2017 let mut t = task("waiting on something else");
2018 t.hold_manual(Some(
2019 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2020 ));
2021 assert_eq!(t.status, TaskStatus::Held);
2022 assert_eq!(
2023 t.hold_reason.as_deref(),
2024 Some("waiting for 20260101-000000-aaaa to land first")
2025 );
2026
2027 t.hold_manual(None);
2029 assert_eq!(
2030 t.hold_reason.as_deref(),
2031 Some("waiting for 20260101-000000-aaaa to land first"),
2032 "a bare re-hold keeps whatever a human already wrote down"
2033 );
2034
2035 let mut plain = task("no reason given");
2037 plain.hold_manual(None);
2038 assert_eq!(plain.status, TaskStatus::Held);
2039 assert!(plain.hold_reason.is_none());
2040
2041 t.release();
2042 assert_eq!(t.status, TaskStatus::Queued);
2043 assert!(
2044 t.hold_reason.is_none(),
2045 "a stale reason must not greet the next person who holds this task"
2046 );
2047 }
2048
2049 #[test]
2050 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2051 let mut t = task("landed by hand while held");
2056 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2057 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2058
2059 t.succeed();
2060 assert_eq!(t.status, TaskStatus::Done);
2061 assert!(
2062 t.hold_reason.is_none(),
2063 "a done task cannot still be waiting on something"
2064 );
2065 }
2066
2067 #[test]
2068 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2069 let mut held = task("held straight out of blocked");
2076 held.block(
2077 vec!["20260101-000000-dead".to_owned()],
2078 Some("waiting on the migration script".to_owned()),
2079 );
2080 assert_eq!(held.status, TaskStatus::Blocked);
2081
2082 held.hold_manual(None);
2083 assert_eq!(held.status, TaskStatus::Held);
2084 assert!(
2085 held.blocked_by.is_empty(),
2086 "hold overrides the wait, same as release"
2087 );
2088 assert!(held.block_reason.is_none());
2089
2090 let mut done = task("closed straight out of blocked");
2091 done.block(
2092 vec!["20260101-000000-dead".to_owned()],
2093 Some("waiting on the migration script".to_owned()),
2094 );
2095 done.succeed();
2096 assert_eq!(done.status, TaskStatus::Done);
2097 assert!(
2098 done.blocked_by.is_empty(),
2099 "a done task cannot still be waiting on a dependency"
2100 );
2101 assert!(done.block_reason.is_none());
2102 }
2103
2104 #[test]
2105 fn a_blocked_task_is_never_offered_to_the_loop() {
2106 let mut t = task("blocked");
2107 assert!(t.status.runnable());
2108 t.block(
2109 vec!["dep-id".to_owned()],
2110 Some("waits on dep-id".to_owned()),
2111 );
2112 assert_eq!(t.status, TaskStatus::Blocked);
2113 assert!(!t.status.runnable());
2114 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2115 }
2116
2117 #[test]
2118 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2119 let mut t = task("blocked on two");
2120 t.block(
2121 vec!["a".to_owned(), "b".to_owned()],
2122 Some("waits on a and b".to_owned()),
2123 );
2124
2125 t.unblock("a");
2126 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2127 assert_eq!(t.blocked_by, ["b"]);
2128
2129 t.unblock("b");
2130 assert_eq!(t.status, TaskStatus::Queued);
2131 assert!(t.blocked_by.is_empty());
2132 assert!(t.block_reason.is_none());
2133 }
2134
2135 #[test]
2136 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2137 let mut t = task("never blocked");
2138 t.unblock("whatever");
2139 assert_eq!(t.status, TaskStatus::Queued);
2140 }
2141
2142 #[test]
2143 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2144 let mut t = task("held, then asked about");
2149 t.hold_machine(Some("out of attempts".to_owned()));
2150 assert_eq!(t.status, TaskStatus::Held);
2151
2152 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2153 assert_eq!(t.status, TaskStatus::Blocked);
2154
2155 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2156 t.unblock("q1");
2157 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2158 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2159 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2160 assert!(t.blocked_from.is_none(), "consumed once restored");
2161 }
2162
2163 #[test]
2164 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2165 let mut t = task("manually held, then asked about");
2166 t.hold_manual(Some("waiting on a dependency".to_owned()));
2167
2168 t.block(vec!["q1".to_owned()], None);
2169 t.unblock("q1");
2170
2171 assert_eq!(t.status, TaskStatus::Held);
2172 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2173 }
2174
2175 #[test]
2176 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2177 let mut t = task("held, blocked twice");
2181 t.hold_machine(None);
2182 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2183 t.block(
2184 vec!["q1".to_owned(), "q2".to_owned()],
2185 Some("second".to_owned()),
2186 );
2187
2188 t.unblock("q1");
2189 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2190 t.unblock("q2");
2191 assert_eq!(t.status, TaskStatus::Held);
2192 }
2193
2194 #[test]
2195 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2196 let mut t = task("blocked mid-run");
2200 t.start("run-1".to_owned());
2201 assert_eq!(t.status, TaskStatus::Running);
2202
2203 t.block(vec!["q1".to_owned()], None);
2204 t.unblock("q1");
2205 assert_eq!(t.status, TaskStatus::Queued);
2206 }
2207
2208 #[test]
2209 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2210 let mut t = task("legacy record, held before it was blocked");
2216 t.hold_source = Some(HoldSource::Machine);
2217 t.hold_reason = Some("legacy hold reason".to_owned());
2218 t.status = TaskStatus::Blocked;
2219 t.blocked_by = vec!["q1".to_owned()];
2220 t.blocked_from = None;
2221
2222 t.unblock("q1");
2223 assert_eq!(t.status, TaskStatus::Held);
2224 }
2225
2226 #[test]
2227 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2228 let mut t = task("legacy record, ordinary dependency block");
2229 t.status = TaskStatus::Blocked;
2230 t.blocked_by = vec!["dep".to_owned()];
2231 t.blocked_from = None;
2232
2233 t.unblock("dep");
2234 assert_eq!(t.status, TaskStatus::Queued);
2235 }
2236
2237 #[test]
2238 fn answering_a_question_is_recorded_and_survives_a_release() {
2239 let mut t = task("asked something");
2240 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2241 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2242 t.unblock("q1");
2243 assert_eq!(t.status, TaskStatus::Queued);
2244 assert_eq!(t.answers.len(), 1);
2245 assert_eq!(t.answers[0].answer, "SQLite");
2246
2247 t.release();
2251 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2252 }
2253
2254 #[test]
2255 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2256 let mut t = task("blocked run with a surviving branch");
2257 t.start("run-1".to_owned());
2258 t.fail("blocked with major findings", 5);
2259 assert_eq!(t.status, TaskStatus::Failed);
2260
2261 t.request_review("magi/eba2/A".to_owned());
2262 assert_eq!(t.status, TaskStatus::Queued);
2263 assert_eq!(t.attempts, 0);
2264 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2265
2266 t.release();
2268 assert!(t.review_branch.is_none());
2269 }
2270
2271 #[test]
2272 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2273 let mut t = task("retry");
2274 t.start("run-1".to_owned());
2275 t.requeue();
2276 assert!(t.fresh_start);
2277
2278 t.release();
2279 assert!(!t.fresh_start);
2280 }
2281
2282 #[test]
2283 fn priority_can_be_changed_while_queued_but_not_while_running() {
2284 let mut t = task("reprioritise me");
2285 t.set_priority(5).unwrap();
2286 assert_eq!(t.priority, 5);
2287
2288 t.start("run-1".to_owned());
2289 let err = t.set_priority(9).unwrap_err().to_string();
2290 assert!(err.contains("running"), "{err}");
2291 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2292 }
2293
2294 #[test]
2295 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2296 let mut t = task("interrupt me");
2297 assert!(!t.interrupt, "off unless asked, same as any other task");
2298
2299 t.set_interrupt(true).unwrap();
2300 assert!(t.interrupt);
2301
2302 t.start("run-1".to_owned());
2303 assert!(
2304 !t.interrupt,
2305 "the mark is one-shot: dispatching the task fulfils it, \
2306 whatever the run that follows ends up doing"
2307 );
2308 let err = t.set_interrupt(true).unwrap_err().to_string();
2309 assert!(err.contains("running"), "{err}");
2310 t.set_interrupt(false).unwrap();
2313 assert!(!t.interrupt);
2314 }
2315
2316 #[test]
2320 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2321 let mut t = task("interrupt me");
2322 t.set_interrupt(true).unwrap();
2323 t.start("run-1".to_owned());
2324 t.fail("mock failure", 5);
2325 assert_eq!(t.status, TaskStatus::Failed);
2326 assert!(
2327 !t.interrupt,
2328 "one attempt already spent the mark; a retry is an ordinary \
2329 requeue, not a fresh interrupt request"
2330 );
2331 }
2332
2333 #[test]
2334 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2335 let (_dir, q) = queue();
2336 let mut a = task("first filed");
2337 let mut b = task("second filed");
2338 a.id = "20260101-000001-aaaa".to_owned();
2339 b.id = "20260101-000002-bbbb".to_owned();
2340 q.put(&mut a).unwrap();
2341 q.put(&mut b).unwrap();
2342
2343 assert_eq!(
2344 q.next_runnable().unwrap().id,
2345 a.id,
2346 "with equal priority the older task goes first, so a burst of \
2347 new work cannot starve it"
2348 );
2349 assert_eq!(
2350 q.list()[0].id,
2351 b.id,
2352 "but the list an operator reads is newest first, the same as \
2353 before priority existed - a's turn to run does not make it the \
2354 newest task"
2355 );
2356
2357 let mut a = q.get(&a.id).unwrap();
2358 a.set_priority(10).unwrap();
2359 q.put(&mut a).unwrap();
2360
2361 assert_eq!(
2362 q.next_runnable().unwrap().id,
2363 a.id,
2364 "a raised priority must be reflected the moment it is saved"
2365 );
2366 assert_eq!(
2370 q.list()[0].id,
2371 a.id,
2372 "the raised task must sort first in the list an operator reads, \
2373 not only in next_runnable's own ordering"
2374 );
2375 }
2376
2377 #[test]
2378 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2379 let mut t = Task::new(
2380 "old title".to_owned(),
2381 "old instruction".to_owned(),
2382 PathBuf::from("/repo"),
2383 Source::Agent {
2384 run: "20260101-000000-beef".to_owned(),
2385 node: "implement".to_owned(),
2386 },
2387 );
2388 let id = t.id.clone();
2389 let created_at = t.created_at;
2390 t.runs.push("20260101-000000-beef".to_owned());
2391
2392 t.edit("new title".to_owned(), "new instruction".to_owned())
2393 .unwrap();
2394
2395 assert_eq!(t.title, "new title");
2396 assert_eq!(t.instruction, "new instruction");
2397 assert_eq!(t.id, id, "editing must not mint a new id");
2398 assert_eq!(t.created_at, created_at);
2399 assert_eq!(
2400 t.source,
2401 Source::Agent {
2402 run: "20260101-000000-beef".to_owned(),
2403 node: "implement".to_owned(),
2404 },
2405 "editing must not turn agent attribution into human"
2406 );
2407 assert_eq!(t.runs, ["20260101-000000-beef"]);
2408 }
2409
2410 #[test]
2411 fn editing_is_refused_once_a_task_is_running_or_finished() {
2412 let mut running = task("in flight");
2413 running.start("run-1".to_owned());
2414 let err = running
2415 .edit("x".to_owned(), "y".to_owned())
2416 .unwrap_err()
2417 .to_string();
2418 assert!(err.contains("running"), "{err}");
2419
2420 let mut done = task("finished");
2421 done.succeed();
2422 let err = done
2423 .edit("x".to_owned(), "y".to_owned())
2424 .unwrap_err()
2425 .to_string();
2426 assert!(err.contains("done"), "{err}");
2427
2428 let mut queued = task("waiting");
2430 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2431 let mut held = task("parked");
2432 held.hold_machine(None);
2433 held.edit("x".to_owned(), "y".to_owned()).unwrap();
2434 }
2435
2436 #[test]
2437 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2438 let (_dir, q) = queue();
2439 let path = q.path_of("20260101-000000-aaaa");
2440 std::fs::create_dir_all(q.root()).unwrap();
2441 std::fs::write(
2442 &path,
2443 serde_json::json!({
2444 "schema": SCHEMA,
2445 "id": "20260101-000000-aaaa",
2446 "title": "from before hold reasons existed",
2447 "instruction": "from before hold reasons existed",
2448 "repo": ".",
2449 "source": { "kind": "human" },
2450 "status": "held",
2451 "created_at": Timestamp::now().to_string(),
2452 "updated_at": Timestamp::now().to_string(),
2453 })
2454 .to_string(),
2455 )
2456 .unwrap();
2457
2458 let task = q.get("20260101-000000-aaaa").expect("must still read");
2459 assert!(task.hold_reason.is_none());
2460 assert!(task.operator_held());
2461 }
2462
2463 #[test]
2464 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
2465 let (_dir, q) = queue();
2466 let path = q.path_of("20260101-000000-bbbb");
2467 std::fs::create_dir_all(q.root()).unwrap();
2468 std::fs::write(
2469 &path,
2470 serde_json::json!({
2471 "schema": 2,
2472 "id": "20260101-000000-bbbb",
2473 "title": "old manual recovery",
2474 "instruction": "old manual recovery",
2475 "repo": ".",
2476 "source": { "kind": "human" },
2477 "status": "held",
2478 "hold_reason": "active manual recovery run20260912-224242-daf5",
2479 "created_at": Timestamp::now().to_string(),
2480 "updated_at": Timestamp::now().to_string(),
2481 })
2482 .to_string(),
2483 )
2484 .unwrap();
2485
2486 let task = q.get("20260101-000000-bbbb").expect("must still read");
2487 assert_eq!(task.hold_source, None);
2488 assert!(task.operator_held());
2489 }
2490
2491 #[test]
2492 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
2493 let (_dir, q) = queue();
2494 let path = q.path_of("20260101-000000-aaaa");
2495 std::fs::create_dir_all(q.root()).unwrap();
2496 std::fs::write(
2497 &path,
2498 serde_json::json!({
2499 "schema": SCHEMA,
2500 "id": "20260101-000000-aaaa",
2501 "title": "from before diagnostics existed",
2502 "instruction": "from before diagnostics existed",
2503 "repo": ".",
2504 "source": { "kind": "human" },
2505 "status": "held",
2506 "created_at": Timestamp::now().to_string(),
2507 "updated_at": Timestamp::now().to_string(),
2508 })
2509 .to_string(),
2510 )
2511 .unwrap();
2512
2513 let task = q.get("20260101-000000-aaaa").expect("must still read");
2514 assert!(task.diagnostic.is_none());
2515 }
2516
2517 #[test]
2518 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
2519 let (_dir, q) = queue();
2523 let path = q.path_of("20260101-000000-aaaa");
2524 std::fs::create_dir_all(q.root()).unwrap();
2525 std::fs::write(
2526 &path,
2527 serde_json::json!({
2528 "schema": 1,
2529 "id": "20260101-000000-aaaa",
2530 "title": "from before blocking existed",
2531 "instruction": "from before blocking existed",
2532 "repo": ".",
2533 "source": { "kind": "human" },
2534 "status": "queued",
2535 "created_at": Timestamp::now().to_string(),
2536 "updated_at": Timestamp::now().to_string(),
2537 })
2538 .to_string(),
2539 )
2540 .unwrap();
2541
2542 let task = q.get("20260101-000000-aaaa").expect("must still read");
2543 assert!(task.blocked_by.is_empty());
2544 assert!(task.block_reason.is_none());
2545 assert!(task.answers.is_empty());
2546 assert!(task.review_branch.is_none());
2547 }
2548
2549 #[test]
2550 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
2551 let mut held = task("diagnosed");
2556 held.start("run-1".to_owned());
2557 held.fail("gate red", 1);
2558 held.diagnostic = Some("cargo test failed: ...".to_owned());
2559 assert_eq!(held.status, TaskStatus::Held);
2560
2561 held.release();
2562 assert!(held.diagnostic.is_none());
2563
2564 held.diagnostic = Some("cargo test failed: ...".to_owned());
2565 held.succeed();
2566 assert!(held.diagnostic.is_none());
2567 }
2568
2569 #[test]
2570 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
2571 let mut t = task("retried");
2572 t.start("run-1".to_owned());
2573 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
2574 t.fail("unrelated config error", 5);
2575 assert_eq!(t.status, TaskStatus::Failed);
2576 assert!(
2577 t.diagnostic.is_none(),
2578 "fail() must not let an old diagnostic outlive the run that produced it"
2579 );
2580 }
2581
2582 #[test]
2583 fn a_claim_is_exclusive_and_releases_on_drop() {
2584 let (_dir, q) = queue();
2585 let mut t = task("contended");
2586 q.put(&mut t).unwrap();
2587
2588 let held = q.claim(&t.id).unwrap();
2589 assert!(
2590 q.claim(&t.id).is_err(),
2591 "two daemons must not drive one task into two runs"
2592 );
2593 drop(held);
2594 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
2595 }
2596
2597 #[test]
2598 fn a_round_trip_survives_disk() {
2599 let (_dir, q) = queue();
2600 let mut t = Task::new(
2601 "titled".to_owned(),
2602 "body".to_owned(),
2603 PathBuf::from("/repo"),
2604 Source::Agent {
2605 run: "20260101-000000-beef".to_owned(),
2606 node: "implement".to_owned(),
2607 },
2608 );
2609 t.priority = 3;
2610 q.put(&mut t).unwrap();
2611
2612 let back = q.get(&t.id).unwrap();
2613 assert_eq!(back.id, t.id);
2614 assert_eq!(back.priority, 3);
2615 assert_eq!(back.source.label(), "implement@beef");
2616 assert_eq!(q.get(t.short()).unwrap().id, t.id);
2618 }
2619
2620 #[test]
2621 fn an_unreadable_task_does_not_take_the_queue_down() {
2622 let (_dir, q) = queue();
2623 let mut t = task("fine");
2624 q.put(&mut t).unwrap();
2625 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
2626
2627 let listed = q.list();
2628 assert_eq!(listed.len(), 1, "the readable task still lists");
2629 assert_eq!(listed[0].id, t.id);
2630 }
2631
2632 #[test]
2633 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
2634 let (_dir, q) = queue();
2635 let path = q.path_of("20260101-000000-aaaa");
2636 std::fs::create_dir_all(q.root()).unwrap();
2637 std::fs::write(
2638 &path,
2639 serde_json::json!({
2640 "schema": SCHEMA,
2641 "id": "20260101-000000-aaaa",
2642 "title": "from before solo existed",
2643 "instruction": "from before solo existed",
2644 "repo": ".",
2645 "source": { "kind": "human" },
2646 "status": "queued",
2647 "created_at": Timestamp::now().to_string(),
2648 "updated_at": Timestamp::now().to_string(),
2649 })
2650 .to_string(),
2651 )
2652 .unwrap();
2653
2654 let task = q.get("20260101-000000-aaaa").expect("must still read");
2655 assert!(!task.solo, "a queue file with no `solo` field means false");
2656 }
2657
2658 #[test]
2659 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
2660 let (_dir, q) = queue();
2661 let path = q.path_of("20260101-000000-bbbb");
2662 std::fs::create_dir_all(q.root()).unwrap();
2663 std::fs::write(
2664 &path,
2665 serde_json::json!({
2666 "schema": SCHEMA,
2667 "id": "20260101-000000-bbbb",
2668 "title": "from before urgent existed",
2669 "instruction": "from before urgent existed",
2670 "repo": ".",
2671 "source": { "kind": "human" },
2672 "status": "queued",
2673 "created_at": Timestamp::now().to_string(),
2674 "updated_at": Timestamp::now().to_string(),
2675 })
2676 .to_string(),
2677 )
2678 .unwrap();
2679
2680 let task = q.get("20260101-000000-bbbb").expect("must still read");
2681 assert!(
2682 !task.urgent,
2683 "a queue file with no `urgent` field means false, same as `solo`"
2684 );
2685 }
2686
2687 #[test]
2688 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
2689 let (_dir, q) = queue();
2690 let mut t = task("from the future");
2691 q.put(&mut t).unwrap();
2692 let path = q.path_of(&t.id);
2693 let body = std::fs::read_to_string(&path)
2694 .unwrap()
2695 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
2696 std::fs::write(&path, body).unwrap();
2697
2698 let err = q.get(&t.id).unwrap_err().to_string();
2699 assert!(err.contains("schema 99"), "{err}");
2700 }
2701
2702 #[test]
2703 fn revision_moves_when_the_queue_changes() {
2704 let (_dir, q) = queue();
2705 assert_eq!(q.revision(), 0, "an empty queue has no revision");
2706 let mut t = task("first");
2707 q.put(&mut t).unwrap();
2708 assert!(q.revision() > 0, "a written task moves the revision");
2709 }
2710
2711 #[test]
2712 fn revision_moves_when_deleting_an_older_task() {
2713 let (dir, q) = queue();
2714 let questions = Questions::at(dir.path().join("questions"));
2715 let mut t1 = task("older");
2716 q.put(&mut t1).unwrap();
2717 std::thread::sleep(std::time::Duration::from_millis(10));
2719 let mut t2 = task("newer");
2720 q.put(&mut t2).unwrap();
2721
2722 let rev_before = q.revision();
2723 q.remove(&t1.id, false, &questions).unwrap();
2724 let rev_after = q.revision();
2725
2726 assert_ne!(
2727 rev_before, rev_after,
2728 "deleting an older task must change the revision so other clients see the deletion"
2729 );
2730 }
2731
2732 #[test]
2733 fn removing_a_task_takes_it_out_of_the_listing() {
2734 let (dir, q) = queue();
2735 let questions = Questions::at(dir.path().join("questions"));
2736 let mut t = task("delete me");
2737 q.put(&mut t).unwrap();
2738 let removed = q.remove(t.short(), false, &questions).unwrap();
2739 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
2740 assert!(removed.quarantined.is_empty(), "nothing was blocked on it");
2741 assert!(q.list().is_empty());
2742 assert!(
2743 q.remove(&t.id, false, &questions).is_err(),
2744 "removing twice is an error"
2745 );
2746 }
2747
2748 #[test]
2749 fn removing_a_task_takes_its_stale_lock_with_it() {
2750 let (dir, q) = queue();
2751 let questions = Questions::at(dir.path().join("questions"));
2752 let mut t = task("interrupted");
2753 q.put(&mut t).unwrap();
2754
2755 let claim = q.claim(&t.id).unwrap();
2758 std::mem::forget(claim);
2759 assert!(
2760 q.claim(&t.id).is_err(),
2761 "the orphaned lock is what makes the task look claimed"
2762 );
2763
2764 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
2766 assert!(err.contains("live daemon"), "{err}");
2767 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
2768
2769 q.remove(&t.id, false, &questions).unwrap();
2771 assert!(q.list().is_empty());
2772 let mut again = task("interrupted");
2773 again.id = t.id.clone();
2774 q.put(&mut again).unwrap();
2775 assert!(
2776 q.claim(&t.id).is_ok(),
2777 "a task that comes back must be claimable, which a left-behind lock would prevent"
2778 );
2779 }
2780
2781 #[test]
2782 fn removing_a_task_quarantines_what_was_blocked_on_it() {
2783 let (dir, q) = queue();
2784 let questions = Questions::at(dir.path().join("questions"));
2785
2786 let mut dep = task("dependency");
2787 q.put(&mut dep).unwrap();
2788
2789 let mut still_valid = task("still valid");
2790 q.put(&mut still_valid).unwrap();
2791
2792 let mut blocked = task("waiting");
2793 blocked.block(
2794 vec![dep.id.clone(), still_valid.id.clone()],
2795 Some("waits on both".to_owned()),
2796 );
2797 q.put(&mut blocked).unwrap();
2798
2799 let removed = q.remove(&dep.id, false, &questions).unwrap();
2800 assert_eq!(removed.quarantined, [blocked.id.clone()]);
2801
2802 let after = q.get(&blocked.id).unwrap();
2803 assert_eq!(after.status, TaskStatus::Held);
2804 assert_eq!(after.hold_source, Some(HoldSource::Machine));
2805 assert!(after.blocked_by.is_empty());
2806 let reason = after.hold_reason.as_deref().unwrap_or_default();
2807 assert!(reason.contains(&dep.id), "{reason}");
2808 assert!(
2809 reason.contains(&still_valid.id),
2810 "the still-valid dependency must survive in the reason text: {reason}"
2811 );
2812 }
2813
2814 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
2815 let p = dir.join(name);
2816 std::fs::write(&p, body).unwrap();
2817 p
2818 }
2819
2820 #[test]
2821 fn an_attachment_copy_survives_deleting_its_source() {
2822 let (dir, q) = queue();
2823 let src = source_file(dir.path(), "shot.png", "pixels");
2824 let mut t = task("with a picture");
2825 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
2826 q.put(&mut t).unwrap();
2827 std::fs::remove_file(&src).unwrap();
2828 assert_eq!(names, ["shot.png"]);
2829 let loaded = q.get(&t.id).unwrap();
2830 let paths = q.attachment_paths(&loaded);
2831 assert_eq!(paths.len(), 1);
2832 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
2833 }
2834
2835 #[test]
2836 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
2837 let (dir, q) = queue();
2838 let mut t = task("bad names");
2839 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
2840 let src = source_file(dir.path(), name, "x");
2841 assert!(
2842 q.attach(&mut t, &[src]).is_err(),
2843 "`{name}` must be refused"
2844 );
2845 }
2846 assert!(!crate::ask::valid_asset_name("C:foo.png"));
2850 #[cfg(not(windows))]
2851 {
2852 let src = source_file(dir.path(), "C:foo.png", "x");
2853 assert!(q.attach(&mut t, &[src]).is_err());
2854 }
2855 let long = format!("{}.png", "a".repeat(70));
2856 let src = source_file(dir.path(), &long, "x");
2857 assert!(q.attach(&mut t, &[src]).is_err());
2858 assert!(t.attachments.is_empty());
2859 assert!(!q.attachments_dir(&t.id).exists());
2860 }
2861
2862 #[test]
2863 fn a_taken_attachment_name_is_numbered_not_overwritten() {
2864 let (dir, q) = queue();
2865 let a = source_file(dir.path(), "shot.png", "one");
2866 let sub = dir.path().join("other");
2867 std::fs::create_dir_all(&sub).unwrap();
2868 let b = source_file(&sub, "shot.png", "two");
2869 let mut t = task("collision");
2870 q.attach(&mut t, &[a]).unwrap();
2871 q.attach(&mut t, &[b]).unwrap();
2872 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
2873 let paths = q.attachment_paths(&t);
2874 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
2875 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
2876 }
2877
2878 #[test]
2879 fn a_renumbered_name_stays_inside_the_length_bound() {
2880 let (dir, q) = queue();
2881 let name = format!("{}.png", "a".repeat(60));
2882 assert_eq!(name.len(), 64);
2883 let a = source_file(dir.path(), &name, "one");
2884 let sub = dir.path().join("other");
2885 std::fs::create_dir_all(&sub).unwrap();
2886 let b = source_file(&sub, &name, "two");
2887 let mut t = task("long");
2888 q.attach(&mut t, &[a, b]).unwrap();
2889 assert_eq!(t.attachments.len(), 2);
2890 assert!(
2891 t.attachments
2892 .iter()
2893 .all(|n| crate::ask::valid_asset_name(n))
2894 );
2895 assert!(t.attachments[1].ends_with("-2.png"));
2896 }
2897
2898 #[test]
2899 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
2900 let (dir, q) = queue();
2901 let good = source_file(dir.path(), "good.png", "ok");
2902 let mut t = task("partial");
2903 q.attach(&mut t, &[good]).unwrap();
2904 let more = source_file(dir.path(), "more.png", "ok");
2905 let missing = dir.path().join("missing.png");
2906 assert!(q.attach(&mut t, &[more, missing]).is_err());
2907 assert_eq!(t.attachments, ["good.png"]);
2908 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
2909 .unwrap()
2910 .flatten()
2911 .collect();
2912 assert_eq!(on_disk.len(), 1);
2913 }
2914
2915 fn block_put(q: &Queue, t: &Task) -> PathBuf {
2917 let tmp = q.path_of(&t.id).with_extension("json.tmp");
2918 std::fs::create_dir_all(&tmp).unwrap();
2919 tmp
2920 }
2921
2922 #[test]
2923 fn a_failed_put_leaves_no_new_attachment_directory() {
2924 let (dir, q) = queue();
2925 let mut t = task("fresh");
2926 let tmp = block_put(&q, &t);
2927 let src = source_file(dir.path(), "shot.png", "x");
2928 assert!(q.attach_and_put(&mut t, &[src]).is_err());
2929 assert!(t.attachments.is_empty());
2930 assert!(!q.attachments_dir(&t.id).exists());
2931 assert!(!q.path_of(&t.id).exists());
2932 std::fs::remove_dir(tmp).unwrap();
2933 }
2934
2935 #[test]
2936 fn a_failed_put_removes_only_the_copy_it_just_made() {
2937 let (dir, q) = queue();
2938 let mut t = task("edited");
2939 let first = source_file(dir.path(), "first.png", "1");
2940 q.attach_and_put(&mut t, &[first]).unwrap();
2941 block_put(&q, &t);
2942 let second = source_file(dir.path(), "second.png", "2");
2943 assert!(q.attach_and_put(&mut t, &[second]).is_err());
2944 assert_eq!(t.attachments, ["first.png"]);
2945 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
2946 .unwrap()
2947 .flatten()
2948 .map(|e| e.file_name().to_string_lossy().into_owned())
2949 .collect();
2950 assert_eq!(on_disk, ["first.png"]);
2951 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
2952 }
2953
2954 #[test]
2955 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
2956 let (dir, q) = queue();
2957 let questions = Questions::at(dir.path().join("questions"));
2958 let gone = task("gone");
2959 let mut other = task("other");
2960 let mut live = task("live");
2961 q.put(&mut other).unwrap();
2962 q.put(&mut live).unwrap();
2963 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
2966 std::fs::create_dir_all(&orphan).unwrap();
2967 std::fs::write(orphan.join("shot.png"), "x").unwrap();
2968 let busy = q.root.join(format!("{}.attachments.removing", live.id));
2971 std::fs::create_dir_all(&busy).unwrap();
2972
2973 q.remove(&other.id, false, &questions).unwrap();
2974 assert!(!orphan.exists(), "an orphan is swept");
2975 assert!(busy.exists(), "a removal in progress is left alone");
2976 }
2977
2978 #[test]
2979 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
2980 let (dir, q) = queue();
2981 let questions = Questions::at(dir.path().join("questions"));
2982 let mut t = task("stuck");
2983 let src = source_file(dir.path(), "shot.png", "x");
2984 q.attach_and_put(&mut t, &[src]).unwrap();
2985 let aside = q.root.join(format!("{}.attachments.removing", t.id));
2988 std::fs::create_dir_all(&aside).unwrap();
2989 std::fs::write(aside.join("old.png"), "o").unwrap();
2990 assert!(q.remove(&t.id, false, &questions).is_err());
2991 assert!(q.path_of(&t.id).exists());
2992 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
2993 }
2994
2995 #[test]
2996 fn a_failed_record_removal_puts_the_attachments_back() {
2997 let (dir, q) = queue();
2998 let mut t = task("rollback");
2999 let src = source_file(dir.path(), "shot.png", "x");
3000 q.attach_and_put(&mut t, &[src]).unwrap();
3001 let err = q
3002 .remove_record_with_attachments(&t.id, |_| {
3003 Err(std::io::Error::other("injected failure"))
3004 })
3005 .unwrap_err();
3006 assert!(format!("{err:#}").contains("injected failure"));
3007 assert!(q.path_of(&t.id).exists());
3008 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3009 assert!(
3010 !q.root
3011 .join(format!("{}.attachments.removing", t.id))
3012 .exists()
3013 );
3014 }
3015
3016 #[test]
3017 fn editing_a_task_keeps_its_attachments() {
3018 let (dir, q) = queue();
3019 let src = source_file(dir.path(), "shot.png", "x");
3020 let mut t = task("editable");
3021 q.attach(&mut t, &[src]).unwrap();
3022 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3023 q.put(&mut t).unwrap();
3024 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3025 }
3026
3027 #[test]
3028 fn removing_a_task_deletes_its_attachments() {
3029 let (dir, q) = queue();
3030 let questions = Questions::at(dir.path().join("questions"));
3031 let src = source_file(dir.path(), "shot.png", "x");
3032 let mut t = task("doomed");
3033 q.attach(&mut t, &[src]).unwrap();
3034 q.put(&mut t).unwrap();
3035 assert!(q.attachments_dir(&t.id).is_dir());
3036 q.remove(&t.id, false, &questions).unwrap();
3037 assert!(!q.attachments_dir(&t.id).exists());
3038 assert!(q.list().is_empty());
3039 }
3040
3041 #[test]
3042 fn a_task_written_before_attachments_still_reads() {
3043 let (_dir, q) = queue();
3044 let mut t = task("old");
3045 q.put(&mut t).unwrap();
3046 let path = q.path_of(&t.id);
3047 let mut v: serde_json::Value =
3048 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3049 v.as_object_mut().unwrap().remove("attachments");
3050 std::fs::write(&path, v.to_string()).unwrap();
3051 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3052 }
3053
3054 #[test]
3055 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3056 let q = Queue::at(PathBuf::from("relative-queue"));
3057 let mut t = task("rel");
3058 t.attachments.push("shot.png".to_owned());
3059 let paths = q.attachment_paths(&t);
3060 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3061 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3062 }
3063}