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
1595fn short(id: &str) -> &str {
1596 id.split('-').next_back().unwrap_or(id)
1597}
1598
1599fn new_id() -> String {
1600 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1601 let seed = crate::rng::entropy();
1602 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1603}
1604
1605#[cfg(test)]
1606mod tests {
1607 #[test]
1608 fn missing_blocker_reason_follows_the_language() {
1609 let b = vec!["a".to_owned()];
1610 let en = missing_blocker_hold_reason_in(&b, &b, "en");
1611 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
1612 assert!(en.starts_with("blocked on a"));
1613 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
1614 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
1615 }
1616
1617 use super::*;
1618
1619 #[test]
1620 fn triage_applied_survives_release_and_old_records_read_as_empty() {
1621 let mut t = Task::new(
1622 "t".to_owned(),
1623 "i".to_owned(),
1624 PathBuf::from("r"),
1625 Source::Human,
1626 );
1627 t.mark_triage_applied("q1");
1628 t.mark_triage_applied("q1");
1629 t.hold_machine(Some("x".to_owned()));
1630 t.release();
1631 assert_eq!(t.triage_applied, ["q1"]);
1632 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
1633
1634 let mut v = serde_json::to_value(&t).unwrap();
1635 v.as_object_mut().unwrap().remove("triage_applied");
1636 let old: Task = serde_json::from_value(v).unwrap();
1637 assert!(old.triage_applied.is_empty());
1638 }
1639
1640 #[test]
1641 fn task_counts_of_empty_is_all_zero() {
1642 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
1643 }
1644
1645 #[test]
1646 fn task_counts_of_tallies_every_status() {
1647 let mut queued = Task::new(
1648 "q".to_owned(),
1649 "i".to_owned(),
1650 PathBuf::from("."),
1651 Source::Human,
1652 );
1653 queued.status = TaskStatus::Queued;
1654 let mut running = queued.clone();
1655 running.status = TaskStatus::Running;
1656 let mut done = queued.clone();
1657 done.status = TaskStatus::Done;
1658 let mut failed = queued.clone();
1659 failed.status = TaskStatus::Failed;
1660 let mut held = queued.clone();
1661 held.status = TaskStatus::Held;
1662 let mut blocked = queued.clone();
1663 blocked.status = TaskStatus::Blocked;
1664
1665 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
1666 assert_eq!(
1667 counts,
1668 TaskCounts {
1669 queued: 1,
1670 running: 1,
1671 done: 2,
1672 failed: 1,
1673 held: 1,
1674 blocked: 1,
1675 }
1676 );
1677 }
1678
1679 fn queue() -> (tempfile::TempDir, Queue) {
1682 let dir = tempfile::tempdir().unwrap();
1683 let q = Queue::at(dir.path().join("queue"));
1684 (dir, q)
1685 }
1686
1687 #[test]
1688 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
1689 let dir = tempfile::tempdir().unwrap();
1690 let q = Queue::at(dir.path().join("queue"));
1691 let mut t = task("held");
1692 q.put(&mut t).unwrap();
1693 assert_eq!(
1694 crate::notices::Notices::at(dir.path().join("notifications"))
1695 .list()
1696 .len(),
1697 0
1698 );
1699 t.hold_machine(Some("out of attempts".to_owned()));
1700 q.put(&mut t).unwrap();
1701 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
1702 assert_eq!(listed.len(), 1);
1703 assert!(listed[0].message.contains("out of attempts"));
1704 }
1705
1706 fn task(title: &str) -> Task {
1707 Task::new(
1708 title.to_owned(),
1709 format!("do {title}"),
1710 PathBuf::from("."),
1711 Source::Human,
1712 )
1713 }
1714
1715 #[test]
1716 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
1717 let mut t = task("retried");
1718 assert!(
1719 t.earlier_attempts().is_empty(),
1720 "a first attempt takes nothing over"
1721 );
1722 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1723 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
1724 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
1726 assert_eq!(t.successor_of("bbbb"), None);
1727 assert_eq!(t.successor_of("zzzz"), None);
1728 }
1729
1730 #[test]
1731 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
1732 let (_dir, q) = queue();
1733 let mut t = task("retried");
1734 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1735 q.put(&mut t).unwrap();
1736
1737 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
1738 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
1739 assert_eq!(
1740 q.superseded_by("cccc"),
1741 None,
1742 "the latest attempt replaces nothing"
1743 );
1744 assert_eq!(
1745 q.superseded_by("never-heard-of-it"),
1746 None,
1747 "a run belonging to no task on this queue is not superseded"
1748 );
1749
1750 let mut by = HashMap::new();
1751 by.insert("aaaa".to_owned(), "bbbb".to_owned());
1752 by.insert("bbbb".to_owned(), "cccc".to_owned());
1753 assert_eq!(
1754 q.superseded(),
1755 by,
1756 "the whole-map and single-run forms must agree"
1757 );
1758 }
1759
1760 #[test]
1761 fn superseded_attempts_is_empty_until_the_task_is_done() {
1762 let mut t = task("retried");
1763 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1764 t.status = TaskStatus::Failed;
1765 assert_eq!(
1766 t.superseded_attempts(true),
1767 &[] as &[String],
1768 "a task still retrying has no attempt yet that a later one made moot"
1769 );
1770
1771 t.status = TaskStatus::Running;
1772 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1773 }
1774
1775 #[test]
1776 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
1777 let mut t = task("retried");
1778 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1779 t.status = TaskStatus::Done;
1780 assert_eq!(
1781 t.superseded_attempts(true),
1782 &["aaaa".to_owned(), "bbbb".to_owned()],
1783 "cccc is the attempt whose success made the task done, and stays out"
1784 );
1785 }
1786
1787 #[test]
1788 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
1789 let mut t = task("first try landed");
1790 t.runs = vec!["aaaa".to_owned()];
1791 t.status = TaskStatus::Done;
1792 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1793 }
1794
1795 #[test]
1796 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
1797 let mut t = task("closed by hand after a manual merge");
1804 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1805 t.status = TaskStatus::Done;
1806 assert_eq!(
1807 t.superseded_attempts(false),
1808 &[] as &[String],
1809 "nothing here is provably why the task is done, so nothing is superseded"
1810 );
1811 }
1812
1813 #[test]
1814 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
1815 let (_dir, q) = queue();
1816 let mut t = task("retried twice");
1817 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1818 q.put(&mut t).unwrap();
1819
1820 assert_eq!(
1821 q.latest_attempt("aaaa"),
1822 Some("cccc".to_owned()),
1823 "an old attempt points straight at the chain's current head, not the \
1824 next attempt in the middle of it"
1825 );
1826 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
1827 assert_eq!(
1828 q.latest_attempt("cccc"),
1829 None,
1830 "the latest attempt is not superseded by anything"
1831 );
1832 assert_eq!(
1833 q.latest_attempt("never-heard-of-it"),
1834 None,
1835 "a run belonging to no task on this queue is not superseded"
1836 );
1837 }
1838
1839 #[test]
1840 fn a_markdown_heading_is_the_title_not_decoration() {
1841 assert_eq!(
1846 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
1847 "Rework the config loader"
1848 );
1849 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
1850 assert_eq!(title_from("> quoted task", 40), "quoted task");
1851 assert_eq!(title_from(" \n\n", 40), "(empty task)");
1853 assert_eq!(title_from("###\n", 40), "(empty task)");
1854 }
1855
1856 #[test]
1857 fn a_long_title_is_elided_by_characters_not_bytes() {
1858 let long = "課題".repeat(30);
1860 let title = title_from(&long, 10);
1861 assert_eq!(title.chars().count(), 10);
1862 assert!(title.ends_with('…'));
1863 }
1864
1865 #[test]
1866 fn priority_wins_and_ties_break_oldest_first() {
1867 let (_dir, q) = queue();
1868 let mut a = task("first");
1869 let mut b = task("second");
1870 let mut c = task("urgent");
1871 a.id = "20260101-000001-aaaa".to_owned();
1873 b.id = "20260101-000002-bbbb".to_owned();
1874 c.id = "20260101-000003-cccc".to_owned();
1875 c.priority = 5;
1876 for t in [&mut a, &mut b, &mut c] {
1877 q.put(t).unwrap();
1878 }
1879
1880 assert_eq!(q.next_runnable().unwrap().id, c.id);
1882 c.hold_machine(None);
1883 q.put(&mut c).unwrap();
1884 assert_eq!(q.next_runnable().unwrap().id, a.id);
1886 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
1887 }
1888
1889 #[test]
1890 fn a_blocked_task_never_starves_another_runnable_one() {
1891 let (_dir, q) = queue();
1892 let mut blocked = task("blocked");
1893 blocked.block(vec!["something".to_owned()], None);
1894 q.put(&mut blocked).unwrap();
1895
1896 let mut runnable = task("free to go");
1897 q.put(&mut runnable).unwrap();
1898
1899 let next = q.next_runnable().expect("a runnable task is still offered");
1900 assert_eq!(next.id, runnable.id);
1901 }
1902
1903 #[test]
1904 fn a_held_task_is_never_offered_to_the_loop() {
1905 let (_dir, q) = queue();
1906 let mut t = task("held");
1907 q.put(&mut t).unwrap();
1908 assert!(q.next_runnable().is_some());
1909
1910 t.hold_machine(None);
1911 q.put(&mut t).unwrap();
1912 assert!(
1913 q.next_runnable().is_none(),
1914 "a held task must wait for a human"
1915 );
1916
1917 t.status = TaskStatus::Failed;
1919 q.put(&mut t).unwrap();
1920 assert!(q.next_runnable().is_some());
1921 }
1922
1923 #[test]
1924 fn attempts_are_capped_and_then_the_task_is_held() {
1925 let mut t = task("doomed");
1926
1927 t.start("run-1".to_owned());
1928 t.fail("gate red", 2);
1929 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
1930
1931 t.start("run-2".to_owned());
1932 t.fail("gate red", 2);
1933 assert_eq!(
1934 t.status,
1935 TaskStatus::Held,
1936 "out of attempts: stop spending money on it"
1937 );
1938 assert_eq!(t.runs, ["run-1", "run-2"]);
1939 assert_eq!(t.last_error.as_deref(), Some("gate red"));
1940 assert_eq!(
1941 t.hold_reason.as_deref(),
1942 Some("gate red"),
1943 "the hold must say why, not leave hold_reason null next to a \
1944 populated last_error"
1945 );
1946 }
1947
1948 #[test]
1949 fn handing_off_a_task_records_a_hold_reason_too() {
1950 let mut t = task("left a pull request");
1951 t.start("run-1".to_owned());
1952 t.handed_off("run ended with a pull request open [run run-1]");
1953 assert_eq!(t.status, TaskStatus::Held);
1954 assert_eq!(t.hold_source, Some(HoldSource::Machine));
1955 assert_eq!(
1956 t.hold_reason.as_deref(),
1957 Some("run ended with a pull request open [run run-1]")
1958 );
1959 assert_eq!(t.hold_reason, t.last_error);
1960 }
1961
1962 #[test]
1963 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
1964 let mut t = task("stalled by quota");
1965
1966 t.start("run-1".to_owned());
1967 assert_eq!(t.attempts, 1);
1968 t.stall("judge-1, judge-2 out of quota");
1969 assert_eq!(
1970 t.attempts, 0,
1971 "a closed quota window must not spend the task's retry budget"
1972 );
1973 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
1974 assert_eq!(
1975 t.last_error.as_deref(),
1976 Some("judge-1, judge-2 out of quota")
1977 );
1978
1979 for _ in 0..20 {
1982 t.start("run-n".to_owned());
1983 t.stall("still out of quota");
1984 }
1985 t.start("run-real".to_owned());
1986 t.fail("gate red", 2);
1987 assert_eq!(
1988 t.status,
1989 TaskStatus::Failed,
1990 "the first attempt that was really judged is attempt one"
1991 );
1992 }
1993
1994 #[test]
1995 fn releasing_a_held_task_gives_it_a_real_second_chance() {
1996 let mut t = task("retry me");
1997 t.start("run-1".to_owned());
1998 t.fail("gate red", 1);
1999 assert_eq!(t.status, TaskStatus::Held);
2000
2001 t.release();
2002 assert_eq!(t.status, TaskStatus::Queued);
2003 assert_eq!(t.attempts, 0);
2006 assert!(t.last_error.is_none());
2007 assert_eq!(
2008 t.runs.len(),
2009 1,
2010 "history is kept: attempts reset, evidence does not"
2011 );
2012 }
2013
2014 #[test]
2015 fn a_hold_reason_survives_and_a_release_clears_it() {
2016 let mut t = task("waiting on something else");
2017 t.hold_manual(Some(
2018 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2019 ));
2020 assert_eq!(t.status, TaskStatus::Held);
2021 assert_eq!(
2022 t.hold_reason.as_deref(),
2023 Some("waiting for 20260101-000000-aaaa to land first")
2024 );
2025
2026 t.hold_manual(None);
2028 assert_eq!(
2029 t.hold_reason.as_deref(),
2030 Some("waiting for 20260101-000000-aaaa to land first"),
2031 "a bare re-hold keeps whatever a human already wrote down"
2032 );
2033
2034 let mut plain = task("no reason given");
2036 plain.hold_manual(None);
2037 assert_eq!(plain.status, TaskStatus::Held);
2038 assert!(plain.hold_reason.is_none());
2039
2040 t.release();
2041 assert_eq!(t.status, TaskStatus::Queued);
2042 assert!(
2043 t.hold_reason.is_none(),
2044 "a stale reason must not greet the next person who holds this task"
2045 );
2046 }
2047
2048 #[test]
2049 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2050 let mut t = task("landed by hand while held");
2055 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2056 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2057
2058 t.succeed();
2059 assert_eq!(t.status, TaskStatus::Done);
2060 assert!(
2061 t.hold_reason.is_none(),
2062 "a done task cannot still be waiting on something"
2063 );
2064 }
2065
2066 #[test]
2067 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2068 let mut held = task("held straight out of blocked");
2075 held.block(
2076 vec!["20260101-000000-dead".to_owned()],
2077 Some("waiting on the migration script".to_owned()),
2078 );
2079 assert_eq!(held.status, TaskStatus::Blocked);
2080
2081 held.hold_manual(None);
2082 assert_eq!(held.status, TaskStatus::Held);
2083 assert!(
2084 held.blocked_by.is_empty(),
2085 "hold overrides the wait, same as release"
2086 );
2087 assert!(held.block_reason.is_none());
2088
2089 let mut done = task("closed straight out of blocked");
2090 done.block(
2091 vec!["20260101-000000-dead".to_owned()],
2092 Some("waiting on the migration script".to_owned()),
2093 );
2094 done.succeed();
2095 assert_eq!(done.status, TaskStatus::Done);
2096 assert!(
2097 done.blocked_by.is_empty(),
2098 "a done task cannot still be waiting on a dependency"
2099 );
2100 assert!(done.block_reason.is_none());
2101 }
2102
2103 #[test]
2104 fn a_blocked_task_is_never_offered_to_the_loop() {
2105 let mut t = task("blocked");
2106 assert!(t.status.runnable());
2107 t.block(
2108 vec!["dep-id".to_owned()],
2109 Some("waits on dep-id".to_owned()),
2110 );
2111 assert_eq!(t.status, TaskStatus::Blocked);
2112 assert!(!t.status.runnable());
2113 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2114 }
2115
2116 #[test]
2117 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2118 let mut t = task("blocked on two");
2119 t.block(
2120 vec!["a".to_owned(), "b".to_owned()],
2121 Some("waits on a and b".to_owned()),
2122 );
2123
2124 t.unblock("a");
2125 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2126 assert_eq!(t.blocked_by, ["b"]);
2127
2128 t.unblock("b");
2129 assert_eq!(t.status, TaskStatus::Queued);
2130 assert!(t.blocked_by.is_empty());
2131 assert!(t.block_reason.is_none());
2132 }
2133
2134 #[test]
2135 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2136 let mut t = task("never blocked");
2137 t.unblock("whatever");
2138 assert_eq!(t.status, TaskStatus::Queued);
2139 }
2140
2141 #[test]
2142 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2143 let mut t = task("held, then asked about");
2148 t.hold_machine(Some("out of attempts".to_owned()));
2149 assert_eq!(t.status, TaskStatus::Held);
2150
2151 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2152 assert_eq!(t.status, TaskStatus::Blocked);
2153
2154 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2155 t.unblock("q1");
2156 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2157 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2158 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2159 assert!(t.blocked_from.is_none(), "consumed once restored");
2160 }
2161
2162 #[test]
2163 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2164 let mut t = task("manually held, then asked about");
2165 t.hold_manual(Some("waiting on a dependency".to_owned()));
2166
2167 t.block(vec!["q1".to_owned()], None);
2168 t.unblock("q1");
2169
2170 assert_eq!(t.status, TaskStatus::Held);
2171 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2172 }
2173
2174 #[test]
2175 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2176 let mut t = task("held, blocked twice");
2180 t.hold_machine(None);
2181 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2182 t.block(
2183 vec!["q1".to_owned(), "q2".to_owned()],
2184 Some("second".to_owned()),
2185 );
2186
2187 t.unblock("q1");
2188 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2189 t.unblock("q2");
2190 assert_eq!(t.status, TaskStatus::Held);
2191 }
2192
2193 #[test]
2194 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2195 let mut t = task("blocked mid-run");
2199 t.start("run-1".to_owned());
2200 assert_eq!(t.status, TaskStatus::Running);
2201
2202 t.block(vec!["q1".to_owned()], None);
2203 t.unblock("q1");
2204 assert_eq!(t.status, TaskStatus::Queued);
2205 }
2206
2207 #[test]
2208 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2209 let mut t = task("legacy record, held before it was blocked");
2215 t.hold_source = Some(HoldSource::Machine);
2216 t.hold_reason = Some("legacy hold reason".to_owned());
2217 t.status = TaskStatus::Blocked;
2218 t.blocked_by = vec!["q1".to_owned()];
2219 t.blocked_from = None;
2220
2221 t.unblock("q1");
2222 assert_eq!(t.status, TaskStatus::Held);
2223 }
2224
2225 #[test]
2226 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2227 let mut t = task("legacy record, ordinary dependency block");
2228 t.status = TaskStatus::Blocked;
2229 t.blocked_by = vec!["dep".to_owned()];
2230 t.blocked_from = None;
2231
2232 t.unblock("dep");
2233 assert_eq!(t.status, TaskStatus::Queued);
2234 }
2235
2236 #[test]
2237 fn answering_a_question_is_recorded_and_survives_a_release() {
2238 let mut t = task("asked something");
2239 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2240 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2241 t.unblock("q1");
2242 assert_eq!(t.status, TaskStatus::Queued);
2243 assert_eq!(t.answers.len(), 1);
2244 assert_eq!(t.answers[0].answer, "SQLite");
2245
2246 t.release();
2250 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2251 }
2252
2253 #[test]
2254 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2255 let mut t = task("blocked run with a surviving branch");
2256 t.start("run-1".to_owned());
2257 t.fail("blocked with major findings", 5);
2258 assert_eq!(t.status, TaskStatus::Failed);
2259
2260 t.request_review("magi/eba2/A".to_owned());
2261 assert_eq!(t.status, TaskStatus::Queued);
2262 assert_eq!(t.attempts, 0);
2263 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2264
2265 t.release();
2267 assert!(t.review_branch.is_none());
2268 }
2269
2270 #[test]
2271 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2272 let mut t = task("retry");
2273 t.start("run-1".to_owned());
2274 t.requeue();
2275 assert!(t.fresh_start);
2276
2277 t.release();
2278 assert!(!t.fresh_start);
2279 }
2280
2281 #[test]
2282 fn priority_can_be_changed_while_queued_but_not_while_running() {
2283 let mut t = task("reprioritise me");
2284 t.set_priority(5).unwrap();
2285 assert_eq!(t.priority, 5);
2286
2287 t.start("run-1".to_owned());
2288 let err = t.set_priority(9).unwrap_err().to_string();
2289 assert!(err.contains("running"), "{err}");
2290 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2291 }
2292
2293 #[test]
2294 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2295 let mut t = task("interrupt me");
2296 assert!(!t.interrupt, "off unless asked, same as any other task");
2297
2298 t.set_interrupt(true).unwrap();
2299 assert!(t.interrupt);
2300
2301 t.start("run-1".to_owned());
2302 assert!(
2303 !t.interrupt,
2304 "the mark is one-shot: dispatching the task fulfils it, \
2305 whatever the run that follows ends up doing"
2306 );
2307 let err = t.set_interrupt(true).unwrap_err().to_string();
2308 assert!(err.contains("running"), "{err}");
2309 t.set_interrupt(false).unwrap();
2312 assert!(!t.interrupt);
2313 }
2314
2315 #[test]
2319 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2320 let mut t = task("interrupt me");
2321 t.set_interrupt(true).unwrap();
2322 t.start("run-1".to_owned());
2323 t.fail("mock failure", 5);
2324 assert_eq!(t.status, TaskStatus::Failed);
2325 assert!(
2326 !t.interrupt,
2327 "one attempt already spent the mark; a retry is an ordinary \
2328 requeue, not a fresh interrupt request"
2329 );
2330 }
2331
2332 #[test]
2333 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2334 let (_dir, q) = queue();
2335 let mut a = task("first filed");
2336 let mut b = task("second filed");
2337 a.id = "20260101-000001-aaaa".to_owned();
2338 b.id = "20260101-000002-bbbb".to_owned();
2339 q.put(&mut a).unwrap();
2340 q.put(&mut b).unwrap();
2341
2342 assert_eq!(
2343 q.next_runnable().unwrap().id,
2344 a.id,
2345 "with equal priority the older task goes first, so a burst of \
2346 new work cannot starve it"
2347 );
2348 assert_eq!(
2349 q.list()[0].id,
2350 b.id,
2351 "but the list an operator reads is newest first, the same as \
2352 before priority existed - a's turn to run does not make it the \
2353 newest task"
2354 );
2355
2356 let mut a = q.get(&a.id).unwrap();
2357 a.set_priority(10).unwrap();
2358 q.put(&mut a).unwrap();
2359
2360 assert_eq!(
2361 q.next_runnable().unwrap().id,
2362 a.id,
2363 "a raised priority must be reflected the moment it is saved"
2364 );
2365 assert_eq!(
2369 q.list()[0].id,
2370 a.id,
2371 "the raised task must sort first in the list an operator reads, \
2372 not only in next_runnable's own ordering"
2373 );
2374 }
2375
2376 #[test]
2377 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2378 let mut t = Task::new(
2379 "old title".to_owned(),
2380 "old instruction".to_owned(),
2381 PathBuf::from("/repo"),
2382 Source::Agent {
2383 run: "20260101-000000-beef".to_owned(),
2384 node: "implement".to_owned(),
2385 },
2386 );
2387 let id = t.id.clone();
2388 let created_at = t.created_at;
2389 t.runs.push("20260101-000000-beef".to_owned());
2390
2391 t.edit("new title".to_owned(), "new instruction".to_owned())
2392 .unwrap();
2393
2394 assert_eq!(t.title, "new title");
2395 assert_eq!(t.instruction, "new instruction");
2396 assert_eq!(t.id, id, "editing must not mint a new id");
2397 assert_eq!(t.created_at, created_at);
2398 assert_eq!(
2399 t.source,
2400 Source::Agent {
2401 run: "20260101-000000-beef".to_owned(),
2402 node: "implement".to_owned(),
2403 },
2404 "editing must not turn agent attribution into human"
2405 );
2406 assert_eq!(t.runs, ["20260101-000000-beef"]);
2407 }
2408
2409 #[test]
2410 fn editing_is_refused_once_a_task_is_running_or_finished() {
2411 let mut running = task("in flight");
2412 running.start("run-1".to_owned());
2413 let err = running
2414 .edit("x".to_owned(), "y".to_owned())
2415 .unwrap_err()
2416 .to_string();
2417 assert!(err.contains("running"), "{err}");
2418
2419 let mut done = task("finished");
2420 done.succeed();
2421 let err = done
2422 .edit("x".to_owned(), "y".to_owned())
2423 .unwrap_err()
2424 .to_string();
2425 assert!(err.contains("done"), "{err}");
2426
2427 let mut queued = task("waiting");
2429 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2430 let mut held = task("parked");
2431 held.hold_machine(None);
2432 held.edit("x".to_owned(), "y".to_owned()).unwrap();
2433 }
2434
2435 #[test]
2436 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2437 let (_dir, q) = queue();
2438 let path = q.path_of("20260101-000000-aaaa");
2439 std::fs::create_dir_all(q.root()).unwrap();
2440 std::fs::write(
2441 &path,
2442 serde_json::json!({
2443 "schema": SCHEMA,
2444 "id": "20260101-000000-aaaa",
2445 "title": "from before hold reasons existed",
2446 "instruction": "from before hold reasons existed",
2447 "repo": ".",
2448 "source": { "kind": "human" },
2449 "status": "held",
2450 "created_at": Timestamp::now().to_string(),
2451 "updated_at": Timestamp::now().to_string(),
2452 })
2453 .to_string(),
2454 )
2455 .unwrap();
2456
2457 let task = q.get("20260101-000000-aaaa").expect("must still read");
2458 assert!(task.hold_reason.is_none());
2459 assert!(task.operator_held());
2460 }
2461
2462 #[test]
2463 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
2464 let (_dir, q) = queue();
2465 let path = q.path_of("20260101-000000-bbbb");
2466 std::fs::create_dir_all(q.root()).unwrap();
2467 std::fs::write(
2468 &path,
2469 serde_json::json!({
2470 "schema": 2,
2471 "id": "20260101-000000-bbbb",
2472 "title": "old manual recovery",
2473 "instruction": "old manual recovery",
2474 "repo": ".",
2475 "source": { "kind": "human" },
2476 "status": "held",
2477 "hold_reason": "active manual recovery run20260912-224242-daf5",
2478 "created_at": Timestamp::now().to_string(),
2479 "updated_at": Timestamp::now().to_string(),
2480 })
2481 .to_string(),
2482 )
2483 .unwrap();
2484
2485 let task = q.get("20260101-000000-bbbb").expect("must still read");
2486 assert_eq!(task.hold_source, None);
2487 assert!(task.operator_held());
2488 }
2489
2490 #[test]
2491 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
2492 let (_dir, q) = queue();
2493 let path = q.path_of("20260101-000000-aaaa");
2494 std::fs::create_dir_all(q.root()).unwrap();
2495 std::fs::write(
2496 &path,
2497 serde_json::json!({
2498 "schema": SCHEMA,
2499 "id": "20260101-000000-aaaa",
2500 "title": "from before diagnostics existed",
2501 "instruction": "from before diagnostics existed",
2502 "repo": ".",
2503 "source": { "kind": "human" },
2504 "status": "held",
2505 "created_at": Timestamp::now().to_string(),
2506 "updated_at": Timestamp::now().to_string(),
2507 })
2508 .to_string(),
2509 )
2510 .unwrap();
2511
2512 let task = q.get("20260101-000000-aaaa").expect("must still read");
2513 assert!(task.diagnostic.is_none());
2514 }
2515
2516 #[test]
2517 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
2518 let (_dir, q) = queue();
2522 let path = q.path_of("20260101-000000-aaaa");
2523 std::fs::create_dir_all(q.root()).unwrap();
2524 std::fs::write(
2525 &path,
2526 serde_json::json!({
2527 "schema": 1,
2528 "id": "20260101-000000-aaaa",
2529 "title": "from before blocking existed",
2530 "instruction": "from before blocking existed",
2531 "repo": ".",
2532 "source": { "kind": "human" },
2533 "status": "queued",
2534 "created_at": Timestamp::now().to_string(),
2535 "updated_at": Timestamp::now().to_string(),
2536 })
2537 .to_string(),
2538 )
2539 .unwrap();
2540
2541 let task = q.get("20260101-000000-aaaa").expect("must still read");
2542 assert!(task.blocked_by.is_empty());
2543 assert!(task.block_reason.is_none());
2544 assert!(task.answers.is_empty());
2545 assert!(task.review_branch.is_none());
2546 }
2547
2548 #[test]
2549 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
2550 let mut held = task("diagnosed");
2555 held.start("run-1".to_owned());
2556 held.fail("gate red", 1);
2557 held.diagnostic = Some("cargo test failed: ...".to_owned());
2558 assert_eq!(held.status, TaskStatus::Held);
2559
2560 held.release();
2561 assert!(held.diagnostic.is_none());
2562
2563 held.diagnostic = Some("cargo test failed: ...".to_owned());
2564 held.succeed();
2565 assert!(held.diagnostic.is_none());
2566 }
2567
2568 #[test]
2569 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
2570 let mut t = task("retried");
2571 t.start("run-1".to_owned());
2572 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
2573 t.fail("unrelated config error", 5);
2574 assert_eq!(t.status, TaskStatus::Failed);
2575 assert!(
2576 t.diagnostic.is_none(),
2577 "fail() must not let an old diagnostic outlive the run that produced it"
2578 );
2579 }
2580
2581 #[test]
2582 fn a_claim_is_exclusive_and_releases_on_drop() {
2583 let (_dir, q) = queue();
2584 let mut t = task("contended");
2585 q.put(&mut t).unwrap();
2586
2587 let held = q.claim(&t.id).unwrap();
2588 assert!(
2589 q.claim(&t.id).is_err(),
2590 "two daemons must not drive one task into two runs"
2591 );
2592 drop(held);
2593 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
2594 }
2595
2596 #[test]
2597 fn a_round_trip_survives_disk() {
2598 let (_dir, q) = queue();
2599 let mut t = Task::new(
2600 "titled".to_owned(),
2601 "body".to_owned(),
2602 PathBuf::from("/repo"),
2603 Source::Agent {
2604 run: "20260101-000000-beef".to_owned(),
2605 node: "implement".to_owned(),
2606 },
2607 );
2608 t.priority = 3;
2609 q.put(&mut t).unwrap();
2610
2611 let back = q.get(&t.id).unwrap();
2612 assert_eq!(back.id, t.id);
2613 assert_eq!(back.priority, 3);
2614 assert_eq!(back.source.label(), "implement@beef");
2615 assert_eq!(q.get(t.short()).unwrap().id, t.id);
2617 }
2618
2619 #[test]
2620 fn an_unreadable_task_does_not_take_the_queue_down() {
2621 let (_dir, q) = queue();
2622 let mut t = task("fine");
2623 q.put(&mut t).unwrap();
2624 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
2625
2626 let listed = q.list();
2627 assert_eq!(listed.len(), 1, "the readable task still lists");
2628 assert_eq!(listed[0].id, t.id);
2629 }
2630
2631 #[test]
2632 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
2633 let (_dir, q) = queue();
2634 let path = q.path_of("20260101-000000-aaaa");
2635 std::fs::create_dir_all(q.root()).unwrap();
2636 std::fs::write(
2637 &path,
2638 serde_json::json!({
2639 "schema": SCHEMA,
2640 "id": "20260101-000000-aaaa",
2641 "title": "from before solo existed",
2642 "instruction": "from before solo existed",
2643 "repo": ".",
2644 "source": { "kind": "human" },
2645 "status": "queued",
2646 "created_at": Timestamp::now().to_string(),
2647 "updated_at": Timestamp::now().to_string(),
2648 })
2649 .to_string(),
2650 )
2651 .unwrap();
2652
2653 let task = q.get("20260101-000000-aaaa").expect("must still read");
2654 assert!(!task.solo, "a queue file with no `solo` field means false");
2655 }
2656
2657 #[test]
2658 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
2659 let (_dir, q) = queue();
2660 let path = q.path_of("20260101-000000-bbbb");
2661 std::fs::create_dir_all(q.root()).unwrap();
2662 std::fs::write(
2663 &path,
2664 serde_json::json!({
2665 "schema": SCHEMA,
2666 "id": "20260101-000000-bbbb",
2667 "title": "from before urgent existed",
2668 "instruction": "from before urgent existed",
2669 "repo": ".",
2670 "source": { "kind": "human" },
2671 "status": "queued",
2672 "created_at": Timestamp::now().to_string(),
2673 "updated_at": Timestamp::now().to_string(),
2674 })
2675 .to_string(),
2676 )
2677 .unwrap();
2678
2679 let task = q.get("20260101-000000-bbbb").expect("must still read");
2680 assert!(
2681 !task.urgent,
2682 "a queue file with no `urgent` field means false, same as `solo`"
2683 );
2684 }
2685
2686 #[test]
2687 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
2688 let (_dir, q) = queue();
2689 let mut t = task("from the future");
2690 q.put(&mut t).unwrap();
2691 let path = q.path_of(&t.id);
2692 let body = std::fs::read_to_string(&path)
2693 .unwrap()
2694 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
2695 std::fs::write(&path, body).unwrap();
2696
2697 let err = q.get(&t.id).unwrap_err().to_string();
2698 assert!(err.contains("schema 99"), "{err}");
2699 }
2700
2701 #[test]
2702 fn revision_moves_when_the_queue_changes() {
2703 let (_dir, q) = queue();
2704 assert_eq!(q.revision(), 0, "an empty queue has no revision");
2705 let mut t = task("first");
2706 q.put(&mut t).unwrap();
2707 assert!(q.revision() > 0, "a written task moves the revision");
2708 }
2709
2710 #[test]
2711 fn revision_moves_when_deleting_an_older_task() {
2712 let (dir, q) = queue();
2713 let questions = Questions::at(dir.path().join("questions"));
2714 let mut t1 = task("older");
2715 q.put(&mut t1).unwrap();
2716 std::thread::sleep(std::time::Duration::from_millis(10));
2718 let mut t2 = task("newer");
2719 q.put(&mut t2).unwrap();
2720
2721 let rev_before = q.revision();
2722 q.remove(&t1.id, false, &questions).unwrap();
2723 let rev_after = q.revision();
2724
2725 assert_ne!(
2726 rev_before, rev_after,
2727 "deleting an older task must change the revision so other clients see the deletion"
2728 );
2729 }
2730
2731 #[test]
2732 fn removing_a_task_takes_it_out_of_the_listing() {
2733 let (dir, q) = queue();
2734 let questions = Questions::at(dir.path().join("questions"));
2735 let mut t = task("delete me");
2736 q.put(&mut t).unwrap();
2737 let removed = q.remove(t.short(), false, &questions).unwrap();
2738 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
2739 assert!(removed.quarantined.is_empty(), "nothing was blocked on it");
2740 assert!(q.list().is_empty());
2741 assert!(
2742 q.remove(&t.id, false, &questions).is_err(),
2743 "removing twice is an error"
2744 );
2745 }
2746
2747 #[test]
2748 fn removing_a_task_takes_its_stale_lock_with_it() {
2749 let (dir, q) = queue();
2750 let questions = Questions::at(dir.path().join("questions"));
2751 let mut t = task("interrupted");
2752 q.put(&mut t).unwrap();
2753
2754 let claim = q.claim(&t.id).unwrap();
2757 std::mem::forget(claim);
2758 assert!(
2759 q.claim(&t.id).is_err(),
2760 "the orphaned lock is what makes the task look claimed"
2761 );
2762
2763 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
2765 assert!(err.contains("live daemon"), "{err}");
2766 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
2767
2768 q.remove(&t.id, false, &questions).unwrap();
2770 assert!(q.list().is_empty());
2771 let mut again = task("interrupted");
2772 again.id = t.id.clone();
2773 q.put(&mut again).unwrap();
2774 assert!(
2775 q.claim(&t.id).is_ok(),
2776 "a task that comes back must be claimable, which a left-behind lock would prevent"
2777 );
2778 }
2779
2780 #[test]
2781 fn removing_a_task_quarantines_what_was_blocked_on_it() {
2782 let (dir, q) = queue();
2783 let questions = Questions::at(dir.path().join("questions"));
2784
2785 let mut dep = task("dependency");
2786 q.put(&mut dep).unwrap();
2787
2788 let mut still_valid = task("still valid");
2789 q.put(&mut still_valid).unwrap();
2790
2791 let mut blocked = task("waiting");
2792 blocked.block(
2793 vec![dep.id.clone(), still_valid.id.clone()],
2794 Some("waits on both".to_owned()),
2795 );
2796 q.put(&mut blocked).unwrap();
2797
2798 let removed = q.remove(&dep.id, false, &questions).unwrap();
2799 assert_eq!(removed.quarantined, [blocked.id.clone()]);
2800
2801 let after = q.get(&blocked.id).unwrap();
2802 assert_eq!(after.status, TaskStatus::Held);
2803 assert_eq!(after.hold_source, Some(HoldSource::Machine));
2804 assert!(after.blocked_by.is_empty());
2805 let reason = after.hold_reason.as_deref().unwrap_or_default();
2806 assert!(reason.contains(&dep.id), "{reason}");
2807 assert!(
2808 reason.contains(&still_valid.id),
2809 "the still-valid dependency must survive in the reason text: {reason}"
2810 );
2811 }
2812
2813 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
2814 let p = dir.join(name);
2815 std::fs::write(&p, body).unwrap();
2816 p
2817 }
2818
2819 #[test]
2820 fn an_attachment_copy_survives_deleting_its_source() {
2821 let (dir, q) = queue();
2822 let src = source_file(dir.path(), "shot.png", "pixels");
2823 let mut t = task("with a picture");
2824 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
2825 q.put(&mut t).unwrap();
2826 std::fs::remove_file(&src).unwrap();
2827 assert_eq!(names, ["shot.png"]);
2828 let loaded = q.get(&t.id).unwrap();
2829 let paths = q.attachment_paths(&loaded);
2830 assert_eq!(paths.len(), 1);
2831 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
2832 }
2833
2834 #[test]
2835 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
2836 let (dir, q) = queue();
2837 let mut t = task("bad names");
2838 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
2839 let src = source_file(dir.path(), name, "x");
2840 assert!(
2841 q.attach(&mut t, &[src]).is_err(),
2842 "`{name}` must be refused"
2843 );
2844 }
2845 assert!(!crate::ask::valid_asset_name("C:foo.png"));
2849 #[cfg(not(windows))]
2850 {
2851 let src = source_file(dir.path(), "C:foo.png", "x");
2852 assert!(q.attach(&mut t, &[src]).is_err());
2853 }
2854 let long = format!("{}.png", "a".repeat(70));
2855 let src = source_file(dir.path(), &long, "x");
2856 assert!(q.attach(&mut t, &[src]).is_err());
2857 assert!(t.attachments.is_empty());
2858 assert!(!q.attachments_dir(&t.id).exists());
2859 }
2860
2861 #[test]
2862 fn a_taken_attachment_name_is_numbered_not_overwritten() {
2863 let (dir, q) = queue();
2864 let a = source_file(dir.path(), "shot.png", "one");
2865 let sub = dir.path().join("other");
2866 std::fs::create_dir_all(&sub).unwrap();
2867 let b = source_file(&sub, "shot.png", "two");
2868 let mut t = task("collision");
2869 q.attach(&mut t, &[a]).unwrap();
2870 q.attach(&mut t, &[b]).unwrap();
2871 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
2872 let paths = q.attachment_paths(&t);
2873 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
2874 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
2875 }
2876
2877 #[test]
2878 fn a_renumbered_name_stays_inside_the_length_bound() {
2879 let (dir, q) = queue();
2880 let name = format!("{}.png", "a".repeat(60));
2881 assert_eq!(name.len(), 64);
2882 let a = source_file(dir.path(), &name, "one");
2883 let sub = dir.path().join("other");
2884 std::fs::create_dir_all(&sub).unwrap();
2885 let b = source_file(&sub, &name, "two");
2886 let mut t = task("long");
2887 q.attach(&mut t, &[a, b]).unwrap();
2888 assert_eq!(t.attachments.len(), 2);
2889 assert!(
2890 t.attachments
2891 .iter()
2892 .all(|n| crate::ask::valid_asset_name(n))
2893 );
2894 assert!(t.attachments[1].ends_with("-2.png"));
2895 }
2896
2897 #[test]
2898 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
2899 let (dir, q) = queue();
2900 let good = source_file(dir.path(), "good.png", "ok");
2901 let mut t = task("partial");
2902 q.attach(&mut t, &[good]).unwrap();
2903 let more = source_file(dir.path(), "more.png", "ok");
2904 let missing = dir.path().join("missing.png");
2905 assert!(q.attach(&mut t, &[more, missing]).is_err());
2906 assert_eq!(t.attachments, ["good.png"]);
2907 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
2908 .unwrap()
2909 .flatten()
2910 .collect();
2911 assert_eq!(on_disk.len(), 1);
2912 }
2913
2914 fn block_put(q: &Queue, t: &Task) -> PathBuf {
2916 let tmp = q.path_of(&t.id).with_extension("json.tmp");
2917 std::fs::create_dir_all(&tmp).unwrap();
2918 tmp
2919 }
2920
2921 #[test]
2922 fn a_failed_put_leaves_no_new_attachment_directory() {
2923 let (dir, q) = queue();
2924 let mut t = task("fresh");
2925 let tmp = block_put(&q, &t);
2926 let src = source_file(dir.path(), "shot.png", "x");
2927 assert!(q.attach_and_put(&mut t, &[src]).is_err());
2928 assert!(t.attachments.is_empty());
2929 assert!(!q.attachments_dir(&t.id).exists());
2930 assert!(!q.path_of(&t.id).exists());
2931 std::fs::remove_dir(tmp).unwrap();
2932 }
2933
2934 #[test]
2935 fn a_failed_put_removes_only_the_copy_it_just_made() {
2936 let (dir, q) = queue();
2937 let mut t = task("edited");
2938 let first = source_file(dir.path(), "first.png", "1");
2939 q.attach_and_put(&mut t, &[first]).unwrap();
2940 block_put(&q, &t);
2941 let second = source_file(dir.path(), "second.png", "2");
2942 assert!(q.attach_and_put(&mut t, &[second]).is_err());
2943 assert_eq!(t.attachments, ["first.png"]);
2944 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
2945 .unwrap()
2946 .flatten()
2947 .map(|e| e.file_name().to_string_lossy().into_owned())
2948 .collect();
2949 assert_eq!(on_disk, ["first.png"]);
2950 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
2951 }
2952
2953 #[test]
2954 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
2955 let (dir, q) = queue();
2956 let questions = Questions::at(dir.path().join("questions"));
2957 let gone = task("gone");
2958 let mut other = task("other");
2959 let mut live = task("live");
2960 q.put(&mut other).unwrap();
2961 q.put(&mut live).unwrap();
2962 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
2965 std::fs::create_dir_all(&orphan).unwrap();
2966 std::fs::write(orphan.join("shot.png"), "x").unwrap();
2967 let busy = q.root.join(format!("{}.attachments.removing", live.id));
2970 std::fs::create_dir_all(&busy).unwrap();
2971
2972 q.remove(&other.id, false, &questions).unwrap();
2973 assert!(!orphan.exists(), "an orphan is swept");
2974 assert!(busy.exists(), "a removal in progress is left alone");
2975 }
2976
2977 #[test]
2978 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
2979 let (dir, q) = queue();
2980 let questions = Questions::at(dir.path().join("questions"));
2981 let mut t = task("stuck");
2982 let src = source_file(dir.path(), "shot.png", "x");
2983 q.attach_and_put(&mut t, &[src]).unwrap();
2984 let aside = q.root.join(format!("{}.attachments.removing", t.id));
2987 std::fs::create_dir_all(&aside).unwrap();
2988 std::fs::write(aside.join("old.png"), "o").unwrap();
2989 assert!(q.remove(&t.id, false, &questions).is_err());
2990 assert!(q.path_of(&t.id).exists());
2991 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
2992 }
2993
2994 #[test]
2995 fn a_failed_record_removal_puts_the_attachments_back() {
2996 let (dir, q) = queue();
2997 let mut t = task("rollback");
2998 let src = source_file(dir.path(), "shot.png", "x");
2999 q.attach_and_put(&mut t, &[src]).unwrap();
3000 let err = q
3001 .remove_record_with_attachments(&t.id, |_| {
3002 Err(std::io::Error::other("injected failure"))
3003 })
3004 .unwrap_err();
3005 assert!(format!("{err:#}").contains("injected failure"));
3006 assert!(q.path_of(&t.id).exists());
3007 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3008 assert!(
3009 !q.root
3010 .join(format!("{}.attachments.removing", t.id))
3011 .exists()
3012 );
3013 }
3014
3015 #[test]
3016 fn editing_a_task_keeps_its_attachments() {
3017 let (dir, q) = queue();
3018 let src = source_file(dir.path(), "shot.png", "x");
3019 let mut t = task("editable");
3020 q.attach(&mut t, &[src]).unwrap();
3021 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3022 q.put(&mut t).unwrap();
3023 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3024 }
3025
3026 #[test]
3027 fn removing_a_task_deletes_its_attachments() {
3028 let (dir, q) = queue();
3029 let questions = Questions::at(dir.path().join("questions"));
3030 let src = source_file(dir.path(), "shot.png", "x");
3031 let mut t = task("doomed");
3032 q.attach(&mut t, &[src]).unwrap();
3033 q.put(&mut t).unwrap();
3034 assert!(q.attachments_dir(&t.id).is_dir());
3035 q.remove(&t.id, false, &questions).unwrap();
3036 assert!(!q.attachments_dir(&t.id).exists());
3037 assert!(q.list().is_empty());
3038 }
3039
3040 #[test]
3041 fn a_task_written_before_attachments_still_reads() {
3042 let (_dir, q) = queue();
3043 let mut t = task("old");
3044 q.put(&mut t).unwrap();
3045 let path = q.path_of(&t.id);
3046 let mut v: serde_json::Value =
3047 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3048 v.as_object_mut().unwrap().remove("attachments");
3049 std::fs::write(&path, v.to_string()).unwrap();
3050 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3051 }
3052
3053 #[test]
3054 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3055 let q = Queue::at(PathBuf::from("relative-queue"));
3056 let mut t = task("rel");
3057 t.attachments.push("shot.png".to_owned());
3058 let paths = q.attachment_paths(&t);
3059 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3060 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3061 }
3062}