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