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 attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1055 let dir = self.attachments_dir(&task.id);
1056 task.attachments
1057 .iter()
1058 .map(|n| {
1059 let p = dir.join(n);
1060 std::path::absolute(&p).unwrap_or(p)
1061 })
1062 .collect()
1063 }
1064
1065 pub fn put(&self, task: &mut Task) -> Result<()> {
1068 task.updated_at = Timestamp::now();
1069 std::fs::create_dir_all(&self.root)
1070 .with_context(|| format!("create {}", self.root.display()))?;
1071 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1072 let path = self.path_of(&task.id);
1073 let tmp = path.with_extension("json.tmp");
1074 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1075 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1076 if let (Some(notice), Some(home)) = (
1081 crate::notices::task_held(task),
1082 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1083 ) {
1084 crate::notices::raise_in(home, notice);
1085 }
1086 Ok(())
1087 }
1088
1089 pub fn get(&self, id: &str) -> Result<Task> {
1091 let resolved = self.resolve_id(id)?;
1092 read_path(&self.path_of(&resolved))
1093 }
1094
1095 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1123 let resolved = self.resolve_id(id)?;
1124 if in_flight {
1125 bail!("task {resolved} is being run by a live daemon right now");
1126 }
1127 self.sweep_removed_attachments();
1134 let attachments = self.attachments_dir(&resolved);
1135 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1136 let moved = match std::fs::rename(&attachments, &aside) {
1137 Ok(()) => true,
1138 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1139 Err(e) => {
1140 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1141 }
1142 };
1143 let path = self.path_of(&resolved);
1144 if let Err(e) = std::fs::remove_file(&path) {
1145 if moved {
1146 let _ = std::fs::rename(&aside, &attachments);
1147 }
1148 return Err(e).with_context(|| format!("remove {}", path.display()));
1149 }
1150 if moved {
1151 if let Err(e) = std::fs::remove_dir_all(&aside) {
1152 tracing::warn!("leftover attachments {}: {e}", aside.display());
1153 }
1154 }
1155 let lock = self.lock_path(&resolved);
1156 if let Err(e) = std::fs::remove_file(&lock) {
1157 if e.kind() != std::io::ErrorKind::NotFound {
1158 return Err(e).with_context(|| format!("remove {}", lock.display()));
1159 }
1160 }
1161 let quarantined = self.quarantine_dependents_of(&resolved, questions);
1162 Ok(Removal {
1163 id: resolved,
1164 quarantined,
1165 })
1166 }
1167
1168 fn sweep_removed_attachments(&self) {
1174 let Ok(entries) = std::fs::read_dir(&self.root) else {
1175 return;
1176 };
1177 for entry in entries.flatten() {
1178 let name = entry.file_name();
1179 let name = name.to_string_lossy();
1180 let Some(id) = name.strip_suffix(".attachments.removing") else {
1181 continue;
1182 };
1183 if !self.path_of(id).exists() {
1186 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1187 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1188 }
1189 }
1190 }
1191 }
1192
1193 fn quarantine_dependents_of(&self, dependency: &str, questions: &Questions) -> Vec<String> {
1197 let mut quarantined = Vec::new();
1198 for listed in self.list() {
1199 if listed.status != TaskStatus::Blocked
1200 || !listed.blocked_by.iter().any(|b| b == dependency)
1201 {
1202 continue;
1203 }
1204 let Ok(_claim) = self.claim(&listed.id) else {
1205 continue;
1206 };
1207 let Ok(mut task) = self.get(&listed.id) else {
1208 continue;
1209 };
1210 if task.status != TaskStatus::Blocked
1211 || !task.blocked_by.iter().any(|b| b == dependency)
1212 {
1213 continue;
1214 }
1215 let missing = missing_blockers(self, questions, &task.blocked_by);
1216 let language = crate::lang::of_repo(&task.repo);
1217 task.hold_machine(Some(missing_blocker_hold_reason_in(
1218 &task.blocked_by,
1219 &missing,
1220 &language,
1221 )));
1222 if self.put(&mut task).is_ok() {
1223 quarantined.push(task.id.clone());
1224 }
1225 }
1226 quarantined
1227 }
1228
1229 fn lock_path(&self, id: &str) -> PathBuf {
1232 self.root.join(format!("{id}.lock"))
1233 }
1234
1235 pub fn list(&self) -> Vec<Task> {
1248 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1249 .into_iter()
1250 .flatten()
1251 .flatten()
1252 .map(|e| e.path())
1253 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1254 .filter_map(|p| read_path(&p).ok())
1255 .collect();
1256 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1257 tasks
1258 }
1259
1260 pub fn superseded(&self) -> HashMap<String, String> {
1273 let mut by = HashMap::new();
1274 for task in self.list() {
1275 for earlier in &task.runs {
1276 if let Some(later) = task.successor_of(earlier) {
1277 by.insert(earlier.clone(), later.clone());
1278 }
1279 }
1280 }
1281 by
1282 }
1283
1284 pub fn superseded_by(&self, run: &str) -> Option<String> {
1292 for task in self.list() {
1293 if task.runs.iter().any(|r| r == run) {
1294 return task.successor_of(run).cloned();
1295 }
1296 }
1297 None
1298 }
1299
1300 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1311 for task in self.list() {
1312 if task.runs.iter().any(|r| r == run) {
1313 return task.runs.last().filter(|last| **last != run).cloned();
1314 }
1315 }
1316 None
1317 }
1318
1319 pub fn next_runnable(&self) -> Option<Task> {
1324 let mut runnable: Vec<Task> = self
1325 .list()
1326 .into_iter()
1327 .filter(|t| t.status.runnable())
1328 .collect();
1329 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1330 runnable.into_iter().next()
1331 }
1332
1333 pub fn claim(&self, id: &str) -> Result<Claim> {
1340 std::fs::create_dir_all(&self.root)
1341 .with_context(|| format!("create {}", self.root.display()))?;
1342 let path = self.lock_path(id);
1343 match std::fs::OpenOptions::new()
1344 .write(true)
1345 .create_new(true)
1346 .open(&path)
1347 {
1348 Ok(mut f) => {
1349 use std::io::Write as _;
1350 let _ = writeln!(f, "{}", std::process::id());
1352 Ok(Claim { path })
1353 }
1354 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1355 bail!("task {id} is already claimed ({} exists)", path.display())
1356 }
1357 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1358 }
1359 }
1360
1361 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1363 if self.path_of(prefix).is_file() {
1364 return Ok(prefix.to_owned());
1365 }
1366 let hits: Vec<String> = self
1367 .list()
1368 .into_iter()
1369 .map(|t| t.id)
1370 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1371 .collect();
1372 match hits.len() {
1373 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1374 0 => bail!("no task matches `{prefix}`"),
1375 _ => bail!(
1376 "`{prefix}` matches {} tasks: {}",
1377 hits.len(),
1378 hits.join(", ")
1379 ),
1380 }
1381 }
1382
1383 pub fn revision(&self) -> u64 {
1390 use std::hash::{Hash as _, Hasher as _};
1391
1392 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
1393 .into_iter()
1394 .flatten()
1395 .flatten()
1396 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
1397 .filter_map(|e| {
1398 let name = e.file_name().to_string_lossy().into_owned();
1399 let mtime = e
1400 .metadata()
1401 .ok()?
1402 .modified()
1403 .ok()?
1404 .duration_since(std::time::UNIX_EPOCH)
1405 .ok()?
1406 .as_millis() as u64;
1407 Some((name, mtime))
1408 })
1409 .collect();
1410
1411 if entries.is_empty() {
1412 return 0;
1413 }
1414
1415 entries.sort_unstable();
1416 let mut hasher = std::hash::DefaultHasher::new();
1417 for (name, mtime) in &entries {
1418 name.hash(&mut hasher);
1419 mtime.hash(&mut hasher);
1420 }
1421 let h = hasher.finish();
1422 if h == 0 { 1 } else { h }
1423 }
1424}
1425
1426#[derive(Debug, Clone)]
1428pub struct Removal {
1429 pub id: String,
1431 pub quarantined: Vec<String>,
1435}
1436
1437#[derive(Debug)]
1439pub struct Claim {
1440 path: PathBuf,
1441}
1442
1443impl Drop for Claim {
1444 fn drop(&mut self) {
1445 let _ = std::fs::remove_file(&self.path);
1446 }
1447}
1448
1449pub fn title_from(instruction: &str, max: usize) -> String {
1452 let line = instruction
1458 .lines()
1459 .map(str::trim)
1460 .find(|l| !l.is_empty())
1461 .unwrap_or("(empty task)")
1462 .trim_start_matches(['#', '-', '*', '>', ' '])
1463 .trim();
1464 if line.is_empty() {
1465 return "(empty task)".to_owned();
1466 }
1467 if line.chars().count() <= max {
1468 return line.to_owned();
1469 }
1470 let head: String = line.chars().take(max.saturating_sub(1)).collect();
1471 format!("{head}…")
1472}
1473
1474fn read_path(path: &Path) -> Result<Task> {
1475 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1476 let task: Task =
1477 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1478 if task.schema > SCHEMA {
1484 bail!(
1485 "task {} was written by a different magi (schema {}, this build \
1486 speaks {SCHEMA})",
1487 task.id,
1488 task.schema
1489 );
1490 }
1491 Ok(task)
1492}
1493
1494pub fn missing_blockers(
1508 queue: &Queue,
1509 questions: &Questions,
1510 blocked_by: &[String],
1511) -> Vec<String> {
1512 blocked_by
1513 .iter()
1514 .filter(|id| !queue.path_of(id).is_file() && !questions.path_of(id).is_file())
1515 .cloned()
1516 .collect()
1517}
1518
1519pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
1531 missing_blocker_hold_reason_in(blocked_by, missing, "en")
1532}
1533
1534pub fn missing_blocker_hold_reason_in(
1536 blocked_by: &[String],
1537 missing: &[String],
1538 language: &str,
1539) -> String {
1540 if crate::lang::is_japanese(language) {
1541 format!(
1542 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
1543 blocked_by.join(", "),
1544 missing.join(", "),
1545 )
1546 } else {
1547 format!(
1548 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
1549 blocked_by.join(", "),
1550 missing.join(", "),
1551 )
1552 }
1553}
1554
1555fn short(id: &str) -> &str {
1556 id.split('-').next_back().unwrap_or(id)
1557}
1558
1559fn new_id() -> String {
1560 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1561 let seed = crate::rng::entropy();
1562 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1563}
1564
1565#[cfg(test)]
1566mod tests {
1567 #[test]
1568 fn missing_blocker_reason_follows_the_language() {
1569 let b = vec!["a".to_owned()];
1570 let en = missing_blocker_hold_reason_in(&b, &b, "en");
1571 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
1572 assert!(en.starts_with("blocked on a"));
1573 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
1574 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
1575 }
1576
1577 use super::*;
1578
1579 #[test]
1580 fn triage_applied_survives_release_and_old_records_read_as_empty() {
1581 let mut t = Task::new(
1582 "t".to_owned(),
1583 "i".to_owned(),
1584 PathBuf::from("r"),
1585 Source::Human,
1586 );
1587 t.mark_triage_applied("q1");
1588 t.mark_triage_applied("q1");
1589 t.hold_machine(Some("x".to_owned()));
1590 t.release();
1591 assert_eq!(t.triage_applied, ["q1"]);
1592 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
1593
1594 let mut v = serde_json::to_value(&t).unwrap();
1595 v.as_object_mut().unwrap().remove("triage_applied");
1596 let old: Task = serde_json::from_value(v).unwrap();
1597 assert!(old.triage_applied.is_empty());
1598 }
1599
1600 #[test]
1601 fn task_counts_of_empty_is_all_zero() {
1602 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
1603 }
1604
1605 #[test]
1606 fn task_counts_of_tallies_every_status() {
1607 let mut queued = Task::new(
1608 "q".to_owned(),
1609 "i".to_owned(),
1610 PathBuf::from("."),
1611 Source::Human,
1612 );
1613 queued.status = TaskStatus::Queued;
1614 let mut running = queued.clone();
1615 running.status = TaskStatus::Running;
1616 let mut done = queued.clone();
1617 done.status = TaskStatus::Done;
1618 let mut failed = queued.clone();
1619 failed.status = TaskStatus::Failed;
1620 let mut held = queued.clone();
1621 held.status = TaskStatus::Held;
1622 let mut blocked = queued.clone();
1623 blocked.status = TaskStatus::Blocked;
1624
1625 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
1626 assert_eq!(
1627 counts,
1628 TaskCounts {
1629 queued: 1,
1630 running: 1,
1631 done: 2,
1632 failed: 1,
1633 held: 1,
1634 blocked: 1,
1635 }
1636 );
1637 }
1638
1639 fn queue() -> (tempfile::TempDir, Queue) {
1642 let dir = tempfile::tempdir().unwrap();
1643 let q = Queue::at(dir.path().join("queue"));
1644 (dir, q)
1645 }
1646
1647 #[test]
1648 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
1649 let dir = tempfile::tempdir().unwrap();
1650 let q = Queue::at(dir.path().join("queue"));
1651 let mut t = task("held");
1652 q.put(&mut t).unwrap();
1653 assert_eq!(
1654 crate::notices::Notices::at(dir.path().join("notifications"))
1655 .list()
1656 .len(),
1657 0
1658 );
1659 t.hold_machine(Some("out of attempts".to_owned()));
1660 q.put(&mut t).unwrap();
1661 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
1662 assert_eq!(listed.len(), 1);
1663 assert!(listed[0].message.contains("out of attempts"));
1664 }
1665
1666 fn task(title: &str) -> Task {
1667 Task::new(
1668 title.to_owned(),
1669 format!("do {title}"),
1670 PathBuf::from("."),
1671 Source::Human,
1672 )
1673 }
1674
1675 #[test]
1676 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
1677 let mut t = task("retried");
1678 assert!(
1679 t.earlier_attempts().is_empty(),
1680 "a first attempt takes nothing over"
1681 );
1682 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1683 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
1684 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
1686 assert_eq!(t.successor_of("bbbb"), None);
1687 assert_eq!(t.successor_of("zzzz"), None);
1688 }
1689
1690 #[test]
1691 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
1692 let (_dir, q) = queue();
1693 let mut t = task("retried");
1694 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1695 q.put(&mut t).unwrap();
1696
1697 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
1698 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
1699 assert_eq!(
1700 q.superseded_by("cccc"),
1701 None,
1702 "the latest attempt replaces nothing"
1703 );
1704 assert_eq!(
1705 q.superseded_by("never-heard-of-it"),
1706 None,
1707 "a run belonging to no task on this queue is not superseded"
1708 );
1709
1710 let mut by = HashMap::new();
1711 by.insert("aaaa".to_owned(), "bbbb".to_owned());
1712 by.insert("bbbb".to_owned(), "cccc".to_owned());
1713 assert_eq!(
1714 q.superseded(),
1715 by,
1716 "the whole-map and single-run forms must agree"
1717 );
1718 }
1719
1720 #[test]
1721 fn superseded_attempts_is_empty_until_the_task_is_done() {
1722 let mut t = task("retried");
1723 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1724 t.status = TaskStatus::Failed;
1725 assert_eq!(
1726 t.superseded_attempts(true),
1727 &[] as &[String],
1728 "a task still retrying has no attempt yet that a later one made moot"
1729 );
1730
1731 t.status = TaskStatus::Running;
1732 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1733 }
1734
1735 #[test]
1736 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
1737 let mut t = task("retried");
1738 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1739 t.status = TaskStatus::Done;
1740 assert_eq!(
1741 t.superseded_attempts(true),
1742 &["aaaa".to_owned(), "bbbb".to_owned()],
1743 "cccc is the attempt whose success made the task done, and stays out"
1744 );
1745 }
1746
1747 #[test]
1748 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
1749 let mut t = task("first try landed");
1750 t.runs = vec!["aaaa".to_owned()];
1751 t.status = TaskStatus::Done;
1752 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1753 }
1754
1755 #[test]
1756 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
1757 let mut t = task("closed by hand after a manual merge");
1764 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1765 t.status = TaskStatus::Done;
1766 assert_eq!(
1767 t.superseded_attempts(false),
1768 &[] as &[String],
1769 "nothing here is provably why the task is done, so nothing is superseded"
1770 );
1771 }
1772
1773 #[test]
1774 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
1775 let (_dir, q) = queue();
1776 let mut t = task("retried twice");
1777 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1778 q.put(&mut t).unwrap();
1779
1780 assert_eq!(
1781 q.latest_attempt("aaaa"),
1782 Some("cccc".to_owned()),
1783 "an old attempt points straight at the chain's current head, not the \
1784 next attempt in the middle of it"
1785 );
1786 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
1787 assert_eq!(
1788 q.latest_attempt("cccc"),
1789 None,
1790 "the latest attempt is not superseded by anything"
1791 );
1792 assert_eq!(
1793 q.latest_attempt("never-heard-of-it"),
1794 None,
1795 "a run belonging to no task on this queue is not superseded"
1796 );
1797 }
1798
1799 #[test]
1800 fn a_markdown_heading_is_the_title_not_decoration() {
1801 assert_eq!(
1806 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
1807 "Rework the config loader"
1808 );
1809 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
1810 assert_eq!(title_from("> quoted task", 40), "quoted task");
1811 assert_eq!(title_from(" \n\n", 40), "(empty task)");
1813 assert_eq!(title_from("###\n", 40), "(empty task)");
1814 }
1815
1816 #[test]
1817 fn a_long_title_is_elided_by_characters_not_bytes() {
1818 let long = "課題".repeat(30);
1820 let title = title_from(&long, 10);
1821 assert_eq!(title.chars().count(), 10);
1822 assert!(title.ends_with('…'));
1823 }
1824
1825 #[test]
1826 fn priority_wins_and_ties_break_oldest_first() {
1827 let (_dir, q) = queue();
1828 let mut a = task("first");
1829 let mut b = task("second");
1830 let mut c = task("urgent");
1831 a.id = "20260101-000001-aaaa".to_owned();
1833 b.id = "20260101-000002-bbbb".to_owned();
1834 c.id = "20260101-000003-cccc".to_owned();
1835 c.priority = 5;
1836 for t in [&mut a, &mut b, &mut c] {
1837 q.put(t).unwrap();
1838 }
1839
1840 assert_eq!(q.next_runnable().unwrap().id, c.id);
1842 c.hold_machine(None);
1843 q.put(&mut c).unwrap();
1844 assert_eq!(q.next_runnable().unwrap().id, a.id);
1846 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
1847 }
1848
1849 #[test]
1850 fn a_blocked_task_never_starves_another_runnable_one() {
1851 let (_dir, q) = queue();
1852 let mut blocked = task("blocked");
1853 blocked.block(vec!["something".to_owned()], None);
1854 q.put(&mut blocked).unwrap();
1855
1856 let mut runnable = task("free to go");
1857 q.put(&mut runnable).unwrap();
1858
1859 let next = q.next_runnable().expect("a runnable task is still offered");
1860 assert_eq!(next.id, runnable.id);
1861 }
1862
1863 #[test]
1864 fn a_held_task_is_never_offered_to_the_loop() {
1865 let (_dir, q) = queue();
1866 let mut t = task("held");
1867 q.put(&mut t).unwrap();
1868 assert!(q.next_runnable().is_some());
1869
1870 t.hold_machine(None);
1871 q.put(&mut t).unwrap();
1872 assert!(
1873 q.next_runnable().is_none(),
1874 "a held task must wait for a human"
1875 );
1876
1877 t.status = TaskStatus::Failed;
1879 q.put(&mut t).unwrap();
1880 assert!(q.next_runnable().is_some());
1881 }
1882
1883 #[test]
1884 fn attempts_are_capped_and_then_the_task_is_held() {
1885 let mut t = task("doomed");
1886
1887 t.start("run-1".to_owned());
1888 t.fail("gate red", 2);
1889 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
1890
1891 t.start("run-2".to_owned());
1892 t.fail("gate red", 2);
1893 assert_eq!(
1894 t.status,
1895 TaskStatus::Held,
1896 "out of attempts: stop spending money on it"
1897 );
1898 assert_eq!(t.runs, ["run-1", "run-2"]);
1899 assert_eq!(t.last_error.as_deref(), Some("gate red"));
1900 assert_eq!(
1901 t.hold_reason.as_deref(),
1902 Some("gate red"),
1903 "the hold must say why, not leave hold_reason null next to a \
1904 populated last_error"
1905 );
1906 }
1907
1908 #[test]
1909 fn handing_off_a_task_records_a_hold_reason_too() {
1910 let mut t = task("left a pull request");
1911 t.start("run-1".to_owned());
1912 t.handed_off("run ended with a pull request open [run run-1]");
1913 assert_eq!(t.status, TaskStatus::Held);
1914 assert_eq!(t.hold_source, Some(HoldSource::Machine));
1915 assert_eq!(
1916 t.hold_reason.as_deref(),
1917 Some("run ended with a pull request open [run run-1]")
1918 );
1919 assert_eq!(t.hold_reason, t.last_error);
1920 }
1921
1922 #[test]
1923 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
1924 let mut t = task("stalled by quota");
1925
1926 t.start("run-1".to_owned());
1927 assert_eq!(t.attempts, 1);
1928 t.stall("judge-1, judge-2 out of quota");
1929 assert_eq!(
1930 t.attempts, 0,
1931 "a closed quota window must not spend the task's retry budget"
1932 );
1933 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
1934 assert_eq!(
1935 t.last_error.as_deref(),
1936 Some("judge-1, judge-2 out of quota")
1937 );
1938
1939 for _ in 0..20 {
1942 t.start("run-n".to_owned());
1943 t.stall("still out of quota");
1944 }
1945 t.start("run-real".to_owned());
1946 t.fail("gate red", 2);
1947 assert_eq!(
1948 t.status,
1949 TaskStatus::Failed,
1950 "the first attempt that was really judged is attempt one"
1951 );
1952 }
1953
1954 #[test]
1955 fn releasing_a_held_task_gives_it_a_real_second_chance() {
1956 let mut t = task("retry me");
1957 t.start("run-1".to_owned());
1958 t.fail("gate red", 1);
1959 assert_eq!(t.status, TaskStatus::Held);
1960
1961 t.release();
1962 assert_eq!(t.status, TaskStatus::Queued);
1963 assert_eq!(t.attempts, 0);
1966 assert!(t.last_error.is_none());
1967 assert_eq!(
1968 t.runs.len(),
1969 1,
1970 "history is kept: attempts reset, evidence does not"
1971 );
1972 }
1973
1974 #[test]
1975 fn a_hold_reason_survives_and_a_release_clears_it() {
1976 let mut t = task("waiting on something else");
1977 t.hold_manual(Some(
1978 "waiting for 20260101-000000-aaaa to land first".to_owned(),
1979 ));
1980 assert_eq!(t.status, TaskStatus::Held);
1981 assert_eq!(
1982 t.hold_reason.as_deref(),
1983 Some("waiting for 20260101-000000-aaaa to land first")
1984 );
1985
1986 t.hold_manual(None);
1988 assert_eq!(
1989 t.hold_reason.as_deref(),
1990 Some("waiting for 20260101-000000-aaaa to land first"),
1991 "a bare re-hold keeps whatever a human already wrote down"
1992 );
1993
1994 let mut plain = task("no reason given");
1996 plain.hold_manual(None);
1997 assert_eq!(plain.status, TaskStatus::Held);
1998 assert!(plain.hold_reason.is_none());
1999
2000 t.release();
2001 assert_eq!(t.status, TaskStatus::Queued);
2002 assert!(
2003 t.hold_reason.is_none(),
2004 "a stale reason must not greet the next person who holds this task"
2005 );
2006 }
2007
2008 #[test]
2009 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2010 let mut t = task("landed by hand while held");
2015 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2016 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2017
2018 t.succeed();
2019 assert_eq!(t.status, TaskStatus::Done);
2020 assert!(
2021 t.hold_reason.is_none(),
2022 "a done task cannot still be waiting on something"
2023 );
2024 }
2025
2026 #[test]
2027 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2028 let mut held = task("held straight out of blocked");
2035 held.block(
2036 vec!["20260101-000000-dead".to_owned()],
2037 Some("waiting on the migration script".to_owned()),
2038 );
2039 assert_eq!(held.status, TaskStatus::Blocked);
2040
2041 held.hold_manual(None);
2042 assert_eq!(held.status, TaskStatus::Held);
2043 assert!(
2044 held.blocked_by.is_empty(),
2045 "hold overrides the wait, same as release"
2046 );
2047 assert!(held.block_reason.is_none());
2048
2049 let mut done = task("closed straight out of blocked");
2050 done.block(
2051 vec!["20260101-000000-dead".to_owned()],
2052 Some("waiting on the migration script".to_owned()),
2053 );
2054 done.succeed();
2055 assert_eq!(done.status, TaskStatus::Done);
2056 assert!(
2057 done.blocked_by.is_empty(),
2058 "a done task cannot still be waiting on a dependency"
2059 );
2060 assert!(done.block_reason.is_none());
2061 }
2062
2063 #[test]
2064 fn a_blocked_task_is_never_offered_to_the_loop() {
2065 let mut t = task("blocked");
2066 assert!(t.status.runnable());
2067 t.block(
2068 vec!["dep-id".to_owned()],
2069 Some("waits on dep-id".to_owned()),
2070 );
2071 assert_eq!(t.status, TaskStatus::Blocked);
2072 assert!(!t.status.runnable());
2073 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2074 }
2075
2076 #[test]
2077 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2078 let mut t = task("blocked on two");
2079 t.block(
2080 vec!["a".to_owned(), "b".to_owned()],
2081 Some("waits on a and b".to_owned()),
2082 );
2083
2084 t.unblock("a");
2085 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2086 assert_eq!(t.blocked_by, ["b"]);
2087
2088 t.unblock("b");
2089 assert_eq!(t.status, TaskStatus::Queued);
2090 assert!(t.blocked_by.is_empty());
2091 assert!(t.block_reason.is_none());
2092 }
2093
2094 #[test]
2095 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2096 let mut t = task("never blocked");
2097 t.unblock("whatever");
2098 assert_eq!(t.status, TaskStatus::Queued);
2099 }
2100
2101 #[test]
2102 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2103 let mut t = task("held, then asked about");
2108 t.hold_machine(Some("out of attempts".to_owned()));
2109 assert_eq!(t.status, TaskStatus::Held);
2110
2111 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2112 assert_eq!(t.status, TaskStatus::Blocked);
2113
2114 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2115 t.unblock("q1");
2116 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2117 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2118 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2119 assert!(t.blocked_from.is_none(), "consumed once restored");
2120 }
2121
2122 #[test]
2123 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2124 let mut t = task("manually held, then asked about");
2125 t.hold_manual(Some("waiting on a dependency".to_owned()));
2126
2127 t.block(vec!["q1".to_owned()], None);
2128 t.unblock("q1");
2129
2130 assert_eq!(t.status, TaskStatus::Held);
2131 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2132 }
2133
2134 #[test]
2135 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2136 let mut t = task("held, blocked twice");
2140 t.hold_machine(None);
2141 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2142 t.block(
2143 vec!["q1".to_owned(), "q2".to_owned()],
2144 Some("second".to_owned()),
2145 );
2146
2147 t.unblock("q1");
2148 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2149 t.unblock("q2");
2150 assert_eq!(t.status, TaskStatus::Held);
2151 }
2152
2153 #[test]
2154 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2155 let mut t = task("blocked mid-run");
2159 t.start("run-1".to_owned());
2160 assert_eq!(t.status, TaskStatus::Running);
2161
2162 t.block(vec!["q1".to_owned()], None);
2163 t.unblock("q1");
2164 assert_eq!(t.status, TaskStatus::Queued);
2165 }
2166
2167 #[test]
2168 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2169 let mut t = task("legacy record, held before it was blocked");
2175 t.hold_source = Some(HoldSource::Machine);
2176 t.hold_reason = Some("legacy hold reason".to_owned());
2177 t.status = TaskStatus::Blocked;
2178 t.blocked_by = vec!["q1".to_owned()];
2179 t.blocked_from = None;
2180
2181 t.unblock("q1");
2182 assert_eq!(t.status, TaskStatus::Held);
2183 }
2184
2185 #[test]
2186 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2187 let mut t = task("legacy record, ordinary dependency block");
2188 t.status = TaskStatus::Blocked;
2189 t.blocked_by = vec!["dep".to_owned()];
2190 t.blocked_from = None;
2191
2192 t.unblock("dep");
2193 assert_eq!(t.status, TaskStatus::Queued);
2194 }
2195
2196 #[test]
2197 fn answering_a_question_is_recorded_and_survives_a_release() {
2198 let mut t = task("asked something");
2199 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2200 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2201 t.unblock("q1");
2202 assert_eq!(t.status, TaskStatus::Queued);
2203 assert_eq!(t.answers.len(), 1);
2204 assert_eq!(t.answers[0].answer, "SQLite");
2205
2206 t.release();
2210 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2211 }
2212
2213 #[test]
2214 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2215 let mut t = task("blocked run with a surviving branch");
2216 t.start("run-1".to_owned());
2217 t.fail("blocked with major findings", 5);
2218 assert_eq!(t.status, TaskStatus::Failed);
2219
2220 t.request_review("magi/eba2/A".to_owned());
2221 assert_eq!(t.status, TaskStatus::Queued);
2222 assert_eq!(t.attempts, 0);
2223 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2224
2225 t.release();
2227 assert!(t.review_branch.is_none());
2228 }
2229
2230 #[test]
2231 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2232 let mut t = task("retry");
2233 t.start("run-1".to_owned());
2234 t.requeue();
2235 assert!(t.fresh_start);
2236
2237 t.release();
2238 assert!(!t.fresh_start);
2239 }
2240
2241 #[test]
2242 fn priority_can_be_changed_while_queued_but_not_while_running() {
2243 let mut t = task("reprioritise me");
2244 t.set_priority(5).unwrap();
2245 assert_eq!(t.priority, 5);
2246
2247 t.start("run-1".to_owned());
2248 let err = t.set_priority(9).unwrap_err().to_string();
2249 assert!(err.contains("running"), "{err}");
2250 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2251 }
2252
2253 #[test]
2254 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2255 let mut t = task("interrupt me");
2256 assert!(!t.interrupt, "off unless asked, same as any other task");
2257
2258 t.set_interrupt(true).unwrap();
2259 assert!(t.interrupt);
2260
2261 t.start("run-1".to_owned());
2262 assert!(
2263 !t.interrupt,
2264 "the mark is one-shot: dispatching the task fulfils it, \
2265 whatever the run that follows ends up doing"
2266 );
2267 let err = t.set_interrupt(true).unwrap_err().to_string();
2268 assert!(err.contains("running"), "{err}");
2269 t.set_interrupt(false).unwrap();
2272 assert!(!t.interrupt);
2273 }
2274
2275 #[test]
2279 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2280 let mut t = task("interrupt me");
2281 t.set_interrupt(true).unwrap();
2282 t.start("run-1".to_owned());
2283 t.fail("mock failure", 5);
2284 assert_eq!(t.status, TaskStatus::Failed);
2285 assert!(
2286 !t.interrupt,
2287 "one attempt already spent the mark; a retry is an ordinary \
2288 requeue, not a fresh interrupt request"
2289 );
2290 }
2291
2292 #[test]
2293 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2294 let (_dir, q) = queue();
2295 let mut a = task("first filed");
2296 let mut b = task("second filed");
2297 a.id = "20260101-000001-aaaa".to_owned();
2298 b.id = "20260101-000002-bbbb".to_owned();
2299 q.put(&mut a).unwrap();
2300 q.put(&mut b).unwrap();
2301
2302 assert_eq!(
2303 q.next_runnable().unwrap().id,
2304 a.id,
2305 "with equal priority the older task goes first, so a burst of \
2306 new work cannot starve it"
2307 );
2308 assert_eq!(
2309 q.list()[0].id,
2310 b.id,
2311 "but the list an operator reads is newest first, the same as \
2312 before priority existed - a's turn to run does not make it the \
2313 newest task"
2314 );
2315
2316 let mut a = q.get(&a.id).unwrap();
2317 a.set_priority(10).unwrap();
2318 q.put(&mut a).unwrap();
2319
2320 assert_eq!(
2321 q.next_runnable().unwrap().id,
2322 a.id,
2323 "a raised priority must be reflected the moment it is saved"
2324 );
2325 assert_eq!(
2329 q.list()[0].id,
2330 a.id,
2331 "the raised task must sort first in the list an operator reads, \
2332 not only in next_runnable's own ordering"
2333 );
2334 }
2335
2336 #[test]
2337 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2338 let mut t = Task::new(
2339 "old title".to_owned(),
2340 "old instruction".to_owned(),
2341 PathBuf::from("/repo"),
2342 Source::Agent {
2343 run: "20260101-000000-beef".to_owned(),
2344 node: "implement".to_owned(),
2345 },
2346 );
2347 let id = t.id.clone();
2348 let created_at = t.created_at;
2349 t.runs.push("20260101-000000-beef".to_owned());
2350
2351 t.edit("new title".to_owned(), "new instruction".to_owned())
2352 .unwrap();
2353
2354 assert_eq!(t.title, "new title");
2355 assert_eq!(t.instruction, "new instruction");
2356 assert_eq!(t.id, id, "editing must not mint a new id");
2357 assert_eq!(t.created_at, created_at);
2358 assert_eq!(
2359 t.source,
2360 Source::Agent {
2361 run: "20260101-000000-beef".to_owned(),
2362 node: "implement".to_owned(),
2363 },
2364 "editing must not turn agent attribution into human"
2365 );
2366 assert_eq!(t.runs, ["20260101-000000-beef"]);
2367 }
2368
2369 #[test]
2370 fn editing_is_refused_once_a_task_is_running_or_finished() {
2371 let mut running = task("in flight");
2372 running.start("run-1".to_owned());
2373 let err = running
2374 .edit("x".to_owned(), "y".to_owned())
2375 .unwrap_err()
2376 .to_string();
2377 assert!(err.contains("running"), "{err}");
2378
2379 let mut done = task("finished");
2380 done.succeed();
2381 let err = done
2382 .edit("x".to_owned(), "y".to_owned())
2383 .unwrap_err()
2384 .to_string();
2385 assert!(err.contains("done"), "{err}");
2386
2387 let mut queued = task("waiting");
2389 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2390 let mut held = task("parked");
2391 held.hold_machine(None);
2392 held.edit("x".to_owned(), "y".to_owned()).unwrap();
2393 }
2394
2395 #[test]
2396 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2397 let (_dir, q) = queue();
2398 let path = q.path_of("20260101-000000-aaaa");
2399 std::fs::create_dir_all(q.root()).unwrap();
2400 std::fs::write(
2401 &path,
2402 serde_json::json!({
2403 "schema": SCHEMA,
2404 "id": "20260101-000000-aaaa",
2405 "title": "from before hold reasons existed",
2406 "instruction": "from before hold reasons existed",
2407 "repo": ".",
2408 "source": { "kind": "human" },
2409 "status": "held",
2410 "created_at": Timestamp::now().to_string(),
2411 "updated_at": Timestamp::now().to_string(),
2412 })
2413 .to_string(),
2414 )
2415 .unwrap();
2416
2417 let task = q.get("20260101-000000-aaaa").expect("must still read");
2418 assert!(task.hold_reason.is_none());
2419 assert!(task.operator_held());
2420 }
2421
2422 #[test]
2423 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
2424 let (_dir, q) = queue();
2425 let path = q.path_of("20260101-000000-bbbb");
2426 std::fs::create_dir_all(q.root()).unwrap();
2427 std::fs::write(
2428 &path,
2429 serde_json::json!({
2430 "schema": 2,
2431 "id": "20260101-000000-bbbb",
2432 "title": "old manual recovery",
2433 "instruction": "old manual recovery",
2434 "repo": ".",
2435 "source": { "kind": "human" },
2436 "status": "held",
2437 "hold_reason": "active manual recovery run20260912-224242-daf5",
2438 "created_at": Timestamp::now().to_string(),
2439 "updated_at": Timestamp::now().to_string(),
2440 })
2441 .to_string(),
2442 )
2443 .unwrap();
2444
2445 let task = q.get("20260101-000000-bbbb").expect("must still read");
2446 assert_eq!(task.hold_source, None);
2447 assert!(task.operator_held());
2448 }
2449
2450 #[test]
2451 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
2452 let (_dir, q) = queue();
2453 let path = q.path_of("20260101-000000-aaaa");
2454 std::fs::create_dir_all(q.root()).unwrap();
2455 std::fs::write(
2456 &path,
2457 serde_json::json!({
2458 "schema": SCHEMA,
2459 "id": "20260101-000000-aaaa",
2460 "title": "from before diagnostics existed",
2461 "instruction": "from before diagnostics existed",
2462 "repo": ".",
2463 "source": { "kind": "human" },
2464 "status": "held",
2465 "created_at": Timestamp::now().to_string(),
2466 "updated_at": Timestamp::now().to_string(),
2467 })
2468 .to_string(),
2469 )
2470 .unwrap();
2471
2472 let task = q.get("20260101-000000-aaaa").expect("must still read");
2473 assert!(task.diagnostic.is_none());
2474 }
2475
2476 #[test]
2477 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
2478 let (_dir, q) = queue();
2482 let path = q.path_of("20260101-000000-aaaa");
2483 std::fs::create_dir_all(q.root()).unwrap();
2484 std::fs::write(
2485 &path,
2486 serde_json::json!({
2487 "schema": 1,
2488 "id": "20260101-000000-aaaa",
2489 "title": "from before blocking existed",
2490 "instruction": "from before blocking existed",
2491 "repo": ".",
2492 "source": { "kind": "human" },
2493 "status": "queued",
2494 "created_at": Timestamp::now().to_string(),
2495 "updated_at": Timestamp::now().to_string(),
2496 })
2497 .to_string(),
2498 )
2499 .unwrap();
2500
2501 let task = q.get("20260101-000000-aaaa").expect("must still read");
2502 assert!(task.blocked_by.is_empty());
2503 assert!(task.block_reason.is_none());
2504 assert!(task.answers.is_empty());
2505 assert!(task.review_branch.is_none());
2506 }
2507
2508 #[test]
2509 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
2510 let mut held = task("diagnosed");
2515 held.start("run-1".to_owned());
2516 held.fail("gate red", 1);
2517 held.diagnostic = Some("cargo test failed: ...".to_owned());
2518 assert_eq!(held.status, TaskStatus::Held);
2519
2520 held.release();
2521 assert!(held.diagnostic.is_none());
2522
2523 held.diagnostic = Some("cargo test failed: ...".to_owned());
2524 held.succeed();
2525 assert!(held.diagnostic.is_none());
2526 }
2527
2528 #[test]
2529 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
2530 let mut t = task("retried");
2531 t.start("run-1".to_owned());
2532 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
2533 t.fail("unrelated config error", 5);
2534 assert_eq!(t.status, TaskStatus::Failed);
2535 assert!(
2536 t.diagnostic.is_none(),
2537 "fail() must not let an old diagnostic outlive the run that produced it"
2538 );
2539 }
2540
2541 #[test]
2542 fn a_claim_is_exclusive_and_releases_on_drop() {
2543 let (_dir, q) = queue();
2544 let mut t = task("contended");
2545 q.put(&mut t).unwrap();
2546
2547 let held = q.claim(&t.id).unwrap();
2548 assert!(
2549 q.claim(&t.id).is_err(),
2550 "two daemons must not drive one task into two runs"
2551 );
2552 drop(held);
2553 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
2554 }
2555
2556 #[test]
2557 fn a_round_trip_survives_disk() {
2558 let (_dir, q) = queue();
2559 let mut t = Task::new(
2560 "titled".to_owned(),
2561 "body".to_owned(),
2562 PathBuf::from("/repo"),
2563 Source::Agent {
2564 run: "20260101-000000-beef".to_owned(),
2565 node: "implement".to_owned(),
2566 },
2567 );
2568 t.priority = 3;
2569 q.put(&mut t).unwrap();
2570
2571 let back = q.get(&t.id).unwrap();
2572 assert_eq!(back.id, t.id);
2573 assert_eq!(back.priority, 3);
2574 assert_eq!(back.source.label(), "implement@beef");
2575 assert_eq!(q.get(t.short()).unwrap().id, t.id);
2577 }
2578
2579 #[test]
2580 fn an_unreadable_task_does_not_take_the_queue_down() {
2581 let (_dir, q) = queue();
2582 let mut t = task("fine");
2583 q.put(&mut t).unwrap();
2584 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
2585
2586 let listed = q.list();
2587 assert_eq!(listed.len(), 1, "the readable task still lists");
2588 assert_eq!(listed[0].id, t.id);
2589 }
2590
2591 #[test]
2592 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
2593 let (_dir, q) = queue();
2594 let path = q.path_of("20260101-000000-aaaa");
2595 std::fs::create_dir_all(q.root()).unwrap();
2596 std::fs::write(
2597 &path,
2598 serde_json::json!({
2599 "schema": SCHEMA,
2600 "id": "20260101-000000-aaaa",
2601 "title": "from before solo existed",
2602 "instruction": "from before solo existed",
2603 "repo": ".",
2604 "source": { "kind": "human" },
2605 "status": "queued",
2606 "created_at": Timestamp::now().to_string(),
2607 "updated_at": Timestamp::now().to_string(),
2608 })
2609 .to_string(),
2610 )
2611 .unwrap();
2612
2613 let task = q.get("20260101-000000-aaaa").expect("must still read");
2614 assert!(!task.solo, "a queue file with no `solo` field means false");
2615 }
2616
2617 #[test]
2618 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
2619 let (_dir, q) = queue();
2620 let path = q.path_of("20260101-000000-bbbb");
2621 std::fs::create_dir_all(q.root()).unwrap();
2622 std::fs::write(
2623 &path,
2624 serde_json::json!({
2625 "schema": SCHEMA,
2626 "id": "20260101-000000-bbbb",
2627 "title": "from before urgent existed",
2628 "instruction": "from before urgent existed",
2629 "repo": ".",
2630 "source": { "kind": "human" },
2631 "status": "queued",
2632 "created_at": Timestamp::now().to_string(),
2633 "updated_at": Timestamp::now().to_string(),
2634 })
2635 .to_string(),
2636 )
2637 .unwrap();
2638
2639 let task = q.get("20260101-000000-bbbb").expect("must still read");
2640 assert!(
2641 !task.urgent,
2642 "a queue file with no `urgent` field means false, same as `solo`"
2643 );
2644 }
2645
2646 #[test]
2647 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
2648 let (_dir, q) = queue();
2649 let mut t = task("from the future");
2650 q.put(&mut t).unwrap();
2651 let path = q.path_of(&t.id);
2652 let body = std::fs::read_to_string(&path)
2653 .unwrap()
2654 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
2655 std::fs::write(&path, body).unwrap();
2656
2657 let err = q.get(&t.id).unwrap_err().to_string();
2658 assert!(err.contains("schema 99"), "{err}");
2659 }
2660
2661 #[test]
2662 fn revision_moves_when_the_queue_changes() {
2663 let (_dir, q) = queue();
2664 assert_eq!(q.revision(), 0, "an empty queue has no revision");
2665 let mut t = task("first");
2666 q.put(&mut t).unwrap();
2667 assert!(q.revision() > 0, "a written task moves the revision");
2668 }
2669
2670 #[test]
2671 fn revision_moves_when_deleting_an_older_task() {
2672 let (dir, q) = queue();
2673 let questions = Questions::at(dir.path().join("questions"));
2674 let mut t1 = task("older");
2675 q.put(&mut t1).unwrap();
2676 std::thread::sleep(std::time::Duration::from_millis(10));
2678 let mut t2 = task("newer");
2679 q.put(&mut t2).unwrap();
2680
2681 let rev_before = q.revision();
2682 q.remove(&t1.id, false, &questions).unwrap();
2683 let rev_after = q.revision();
2684
2685 assert_ne!(
2686 rev_before, rev_after,
2687 "deleting an older task must change the revision so other clients see the deletion"
2688 );
2689 }
2690
2691 #[test]
2692 fn removing_a_task_takes_it_out_of_the_listing() {
2693 let (dir, q) = queue();
2694 let questions = Questions::at(dir.path().join("questions"));
2695 let mut t = task("delete me");
2696 q.put(&mut t).unwrap();
2697 let removed = q.remove(t.short(), false, &questions).unwrap();
2698 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
2699 assert!(removed.quarantined.is_empty(), "nothing was blocked on it");
2700 assert!(q.list().is_empty());
2701 assert!(
2702 q.remove(&t.id, false, &questions).is_err(),
2703 "removing twice is an error"
2704 );
2705 }
2706
2707 #[test]
2708 fn removing_a_task_takes_its_stale_lock_with_it() {
2709 let (dir, q) = queue();
2710 let questions = Questions::at(dir.path().join("questions"));
2711 let mut t = task("interrupted");
2712 q.put(&mut t).unwrap();
2713
2714 let claim = q.claim(&t.id).unwrap();
2717 std::mem::forget(claim);
2718 assert!(
2719 q.claim(&t.id).is_err(),
2720 "the orphaned lock is what makes the task look claimed"
2721 );
2722
2723 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
2725 assert!(err.contains("live daemon"), "{err}");
2726 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
2727
2728 q.remove(&t.id, false, &questions).unwrap();
2730 assert!(q.list().is_empty());
2731 let mut again = task("interrupted");
2732 again.id = t.id.clone();
2733 q.put(&mut again).unwrap();
2734 assert!(
2735 q.claim(&t.id).is_ok(),
2736 "a task that comes back must be claimable, which a left-behind lock would prevent"
2737 );
2738 }
2739
2740 #[test]
2741 fn removing_a_task_quarantines_what_was_blocked_on_it() {
2742 let (dir, q) = queue();
2743 let questions = Questions::at(dir.path().join("questions"));
2744
2745 let mut dep = task("dependency");
2746 q.put(&mut dep).unwrap();
2747
2748 let mut still_valid = task("still valid");
2749 q.put(&mut still_valid).unwrap();
2750
2751 let mut blocked = task("waiting");
2752 blocked.block(
2753 vec![dep.id.clone(), still_valid.id.clone()],
2754 Some("waits on both".to_owned()),
2755 );
2756 q.put(&mut blocked).unwrap();
2757
2758 let removed = q.remove(&dep.id, false, &questions).unwrap();
2759 assert_eq!(removed.quarantined, [blocked.id.clone()]);
2760
2761 let after = q.get(&blocked.id).unwrap();
2762 assert_eq!(after.status, TaskStatus::Held);
2763 assert_eq!(after.hold_source, Some(HoldSource::Machine));
2764 assert!(after.blocked_by.is_empty());
2765 let reason = after.hold_reason.as_deref().unwrap_or_default();
2766 assert!(reason.contains(&dep.id), "{reason}");
2767 assert!(
2768 reason.contains(&still_valid.id),
2769 "the still-valid dependency must survive in the reason text: {reason}"
2770 );
2771 }
2772
2773 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
2774 let p = dir.join(name);
2775 std::fs::write(&p, body).unwrap();
2776 p
2777 }
2778
2779 #[test]
2780 fn an_attachment_copy_survives_deleting_its_source() {
2781 let (dir, q) = queue();
2782 let src = source_file(dir.path(), "shot.png", "pixels");
2783 let mut t = task("with a picture");
2784 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
2785 q.put(&mut t).unwrap();
2786 std::fs::remove_file(&src).unwrap();
2787 assert_eq!(names, ["shot.png"]);
2788 let loaded = q.get(&t.id).unwrap();
2789 let paths = q.attachment_paths(&loaded);
2790 assert_eq!(paths.len(), 1);
2791 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
2792 }
2793
2794 #[test]
2795 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
2796 let (dir, q) = queue();
2797 let mut t = task("bad names");
2798 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
2799 let src = source_file(dir.path(), name, "x");
2800 assert!(
2801 q.attach(&mut t, &[src]).is_err(),
2802 "`{name}` must be refused"
2803 );
2804 }
2805 assert!(!crate::ask::valid_asset_name("C:foo.png"));
2809 #[cfg(not(windows))]
2810 {
2811 let src = source_file(dir.path(), "C:foo.png", "x");
2812 assert!(q.attach(&mut t, &[src]).is_err());
2813 }
2814 let long = format!("{}.png", "a".repeat(70));
2815 let src = source_file(dir.path(), &long, "x");
2816 assert!(q.attach(&mut t, &[src]).is_err());
2817 assert!(t.attachments.is_empty());
2818 assert!(!q.attachments_dir(&t.id).exists());
2819 }
2820
2821 #[test]
2822 fn a_taken_attachment_name_is_numbered_not_overwritten() {
2823 let (dir, q) = queue();
2824 let a = source_file(dir.path(), "shot.png", "one");
2825 let sub = dir.path().join("other");
2826 std::fs::create_dir_all(&sub).unwrap();
2827 let b = source_file(&sub, "shot.png", "two");
2828 let mut t = task("collision");
2829 q.attach(&mut t, &[a]).unwrap();
2830 q.attach(&mut t, &[b]).unwrap();
2831 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
2832 let paths = q.attachment_paths(&t);
2833 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
2834 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
2835 }
2836
2837 #[test]
2838 fn a_renumbered_name_stays_inside_the_length_bound() {
2839 let (dir, q) = queue();
2840 let name = format!("{}.png", "a".repeat(60));
2841 assert_eq!(name.len(), 64);
2842 let a = source_file(dir.path(), &name, "one");
2843 let sub = dir.path().join("other");
2844 std::fs::create_dir_all(&sub).unwrap();
2845 let b = source_file(&sub, &name, "two");
2846 let mut t = task("long");
2847 q.attach(&mut t, &[a, b]).unwrap();
2848 assert_eq!(t.attachments.len(), 2);
2849 assert!(
2850 t.attachments
2851 .iter()
2852 .all(|n| crate::ask::valid_asset_name(n))
2853 );
2854 assert!(t.attachments[1].ends_with("-2.png"));
2855 }
2856
2857 #[test]
2858 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
2859 let (dir, q) = queue();
2860 let good = source_file(dir.path(), "good.png", "ok");
2861 let mut t = task("partial");
2862 q.attach(&mut t, &[good]).unwrap();
2863 let more = source_file(dir.path(), "more.png", "ok");
2864 let missing = dir.path().join("missing.png");
2865 assert!(q.attach(&mut t, &[more, missing]).is_err());
2866 assert_eq!(t.attachments, ["good.png"]);
2867 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
2868 .unwrap()
2869 .flatten()
2870 .collect();
2871 assert_eq!(on_disk.len(), 1);
2872 }
2873
2874 #[test]
2875 fn editing_a_task_keeps_its_attachments() {
2876 let (dir, q) = queue();
2877 let src = source_file(dir.path(), "shot.png", "x");
2878 let mut t = task("editable");
2879 q.attach(&mut t, &[src]).unwrap();
2880 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
2881 q.put(&mut t).unwrap();
2882 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
2883 }
2884
2885 #[test]
2886 fn removing_a_task_deletes_its_attachments() {
2887 let (dir, q) = queue();
2888 let questions = Questions::at(dir.path().join("questions"));
2889 let src = source_file(dir.path(), "shot.png", "x");
2890 let mut t = task("doomed");
2891 q.attach(&mut t, &[src]).unwrap();
2892 q.put(&mut t).unwrap();
2893 assert!(q.attachments_dir(&t.id).is_dir());
2894 q.remove(&t.id, false, &questions).unwrap();
2895 assert!(!q.attachments_dir(&t.id).exists());
2896 assert!(q.list().is_empty());
2897 }
2898
2899 #[test]
2900 fn a_task_written_before_attachments_still_reads() {
2901 let (_dir, q) = queue();
2902 let mut t = task("old");
2903 q.put(&mut t).unwrap();
2904 let path = q.path_of(&t.id);
2905 let mut v: serde_json::Value =
2906 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
2907 v.as_object_mut().unwrap().remove("attachments");
2908 std::fs::write(&path, v.to_string()).unwrap();
2909 assert!(q.get(&t.id).unwrap().attachments.is_empty());
2910 }
2911
2912 #[test]
2913 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
2914 let q = Queue::at(PathBuf::from("relative-queue"));
2915 let mut t = task("rel");
2916 t.attachments.push("shot.png".to_owned());
2917 let paths = q.attachment_paths(&t);
2918 assert!(paths[0].is_absolute(), "{}", paths[0].display());
2919 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
2920 }
2921}