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 link_run(&mut self, run: &str) -> bool {
552 if self.runs.iter().any(|r| r == run) {
553 return false;
554 }
555 self.runs.push(run.to_owned());
556 true
557 }
558
559 pub fn succeed(&mut self) {
569 self.status = TaskStatus::Done;
570 self.resume_override = None;
571 self.last_error = None;
572 self.hold_reason = None;
573 self.hold_source = None;
574 self.diagnostic = None;
575 self.blocked_by.clear();
576 self.block_reason = None;
577 self.blocked_from = None;
578 }
579
580 pub fn already_landed(&mut self, note: impl Into<String>) {
588 self.succeed();
589 self.attempts = self.attempts.saturating_sub(1);
590 self.last_error = Some(note.into());
591 }
592
593 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
623 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
624 return &[];
625 }
626 &self.runs[..self.runs.len() - 1]
627 }
628
629 pub fn successor_of(&self, run: &str) -> Option<&String> {
636 let pos = self.runs.iter().position(|r| r == run)?;
637 self.runs.get(pos + 1)
638 }
639
640 pub fn earlier_attempts(&self) -> &[String] {
647 &self.runs
648 }
649
650 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
666 let why = why.into();
667 self.diagnostic = None;
668 self.status = if self.attempts >= max_attempts {
669 self.hold_source = Some(HoldSource::Machine);
670 self.hold_reason = Some(why.clone());
671 TaskStatus::Held
672 } else {
673 TaskStatus::Failed
674 };
675 self.last_error = Some(why);
676 }
677
678 pub fn stall(&mut self, why: impl Into<String>) {
688 self.last_error = Some(why.into());
689 self.diagnostic = None;
690 self.attempts = self.attempts.saturating_sub(1);
691 self.status = TaskStatus::Failed;
692 }
693
694 pub fn operator_held(&self) -> bool {
700 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
701 }
702
703 pub fn hold_manual(&mut self, reason: Option<String>) {
713 self.status = TaskStatus::Held;
714 if reason.is_some() {
715 self.hold_reason = reason;
716 }
717 self.hold_source = Some(HoldSource::Manual);
718 self.blocked_by.clear();
719 self.block_reason = None;
720 self.blocked_from = None;
721 }
722
723 pub fn hold_machine(&mut self, reason: Option<String>) {
728 self.status = TaskStatus::Held;
729 if reason.is_some() {
730 self.hold_reason = reason;
731 }
732 self.hold_source = Some(HoldSource::Machine);
733 self.blocked_by.clear();
734 self.block_reason = None;
735 self.blocked_from = None;
736 }
737
738 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
748 if self.status != TaskStatus::Blocked {
749 self.blocked_from = Some(self.status);
750 }
751 self.status = TaskStatus::Blocked;
752 self.blocked_by = blocked_by;
753 self.block_reason = reason;
754 }
755
756 pub fn unblock(&mut self, resolved_id: &str) {
779 if self.status != TaskStatus::Blocked {
780 return;
781 }
782 self.blocked_by.retain(|id| id != resolved_id);
783 if self.blocked_by.is_empty() {
784 self.status = match self.blocked_from {
785 Some(TaskStatus::Running) => TaskStatus::Queued,
786 Some(other) => other,
787 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
788 TaskStatus::Held
789 }
790 None => TaskStatus::Queued,
791 };
792 self.block_reason = None;
793 self.blocked_from = None;
794 }
795 }
796
797 pub fn record_answer(&mut self, question: String, answer: String) {
802 self.answers.push(AnsweredQuestion { question, answer });
803 }
804
805 pub fn request_review(&mut self, branch: String) {
809 self.release();
810 self.review_branch = Some(branch);
811 }
812
813 pub fn requeue(&mut self) {
816 self.release();
817 self.fresh_start = true;
818 }
819
820 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
829 if self.status == TaskStatus::Running {
830 bail!(
831 "task {} is running; its priority cannot be changed until \
832 this attempt finishes",
833 self.short()
834 );
835 }
836 self.priority = priority;
837 Ok(())
838 }
839
840 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
855 if interrupt && !self.status.runnable() {
856 bail!(
857 "task {} is {}; only a queued or failed task can be marked \
858 to interrupt",
859 self.short(),
860 self.status.as_str()
861 );
862 }
863 self.interrupt = interrupt;
864 Ok(())
865 }
866
867 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
879 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
880 bail!(
881 "task {} is {}; only a queued or held task's instruction can \
882 be edited",
883 self.short(),
884 self.status.as_str()
885 );
886 }
887 self.title = title;
888 self.instruction = instruction;
889 Ok(())
890 }
891
892 pub fn handed_off(&mut self, why: impl Into<String>) {
912 let why = why.into();
913 self.diagnostic = None;
914 self.status = TaskStatus::Held;
915 self.hold_source = Some(HoldSource::Machine);
916 self.hold_reason = Some(why.clone());
917 self.last_error = Some(why);
918 }
919
920 pub fn release(&mut self) {
924 self.status = TaskStatus::Queued;
925 self.attempts = 0;
926 self.last_error = None;
927 self.hold_reason = None;
930 self.hold_source = None;
931 self.diagnostic = None;
932 self.blocked_by.clear();
937 self.block_reason = None;
938 self.blocked_from = None;
939 self.review_branch = None;
940 self.fresh_start = false;
941 }
942}
943
944fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
948 use std::io::ErrorKind;
949 let (stem, ext) = match name.rfind('.') {
950 Some(i) if i > 0 => (&name[..i], &name[i..]),
951 _ => (name, ""),
952 };
953 for n in 1u32.. {
954 let candidate = if n == 1 {
955 name.to_owned()
956 } else {
957 let suffix = format!("-{n}");
958 let room = 64usize.saturating_sub(suffix.len() + ext.len());
959 let stem: String = stem.chars().take(room).collect();
960 format!("{stem}{suffix}{ext}")
961 };
962 if !crate::ask::valid_asset_name(&candidate) {
963 bail!("no valid attachment name is left for `{name}`");
964 }
965 let path = dir.join(&candidate);
966 match std::fs::OpenOptions::new()
967 .write(true)
968 .create_new(true)
969 .open(&path)
970 {
971 Ok(mut out) => {
972 let copied = std::fs::File::open(src)
973 .and_then(|mut input| std::io::copy(&mut input, &mut out));
974 if let Err(e) = copied {
975 drop(out);
976 let _ = std::fs::remove_file(&path);
977 return Err(e).with_context(|| format!("copy {}", src.display()));
978 }
979 return Ok((candidate, path));
980 }
981 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
982 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
983 }
984 }
985 unreachable!("the counter never runs out")
986}
987
988const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
991
992struct TaskLock(PathBuf);
994
995impl Drop for TaskLock {
996 fn drop(&mut self) {
997 let _ = std::fs::remove_file(&self.0);
998 }
999}
1000
1001#[derive(Debug, Clone)]
1003pub struct Queue {
1004 root: PathBuf,
1005}
1006
1007impl Queue {
1008 pub fn open() -> Self {
1010 Self::at(crate::run::home().join("queue"))
1011 }
1012
1013 pub fn at(root: PathBuf) -> Self {
1016 Self { root }
1017 }
1018
1019 pub fn root(&self) -> &Path {
1021 &self.root
1022 }
1023
1024 pub fn path_of(&self, id: &str) -> PathBuf {
1026 self.root.join(format!("{id}.json"))
1027 }
1028
1029 pub fn attachments_dir(&self, id: &str) -> PathBuf {
1032 self.root.join(format!("{id}.attachments"))
1033 }
1034
1035 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1046 let mut wanted = Vec::new();
1047 for src in sources {
1048 let name = src
1049 .file_name()
1050 .and_then(|n| n.to_str())
1051 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1052 if !crate::ask::valid_asset_name(name) {
1053 bail!(
1054 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1055 with no `..`; rename the file and try again"
1056 );
1057 }
1058 if !src.is_file() {
1059 bail!("attachment `{}` is not a file", src.display());
1060 }
1061 wanted.push((src, name));
1062 }
1063 let dir = self.attachments_dir(&task.id);
1064 let existed = dir.is_dir();
1065 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1066 let mut created: Vec<PathBuf> = Vec::new();
1067 let mut names = Vec::new();
1068 let mut copy_all = || -> Result<()> {
1069 for (src, name) in &wanted {
1070 let (stored, path) = copy_new(&dir, src, name)?;
1071 created.push(path);
1072 names.push(stored);
1073 }
1074 Ok(())
1075 };
1076 if let Err(e) = copy_all() {
1077 for path in &created {
1078 let _ = std::fs::remove_file(path);
1079 }
1080 if !existed {
1081 let _ = std::fs::remove_dir(&dir);
1082 }
1083 return Err(e);
1084 }
1085 task.attachments.extend(names.iter().cloned());
1086 Ok(names)
1087 }
1088
1089 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1100 let dir = self.attachments_dir(&task.id);
1101 let existed = dir.is_dir();
1102 let before = task.attachments.len();
1103 let names = self.attach(task, sources)?;
1104 if let Err(e) = self.put(task) {
1105 for name in &names {
1106 let _ = std::fs::remove_file(dir.join(name));
1107 }
1108 if !existed {
1109 let _ = std::fs::remove_dir(&dir);
1110 }
1111 task.attachments.truncate(before);
1112 return Err(e);
1113 }
1114 Ok(names)
1115 }
1116
1117 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1121 let dir = self.attachments_dir(&task.id);
1122 task.attachments
1123 .iter()
1124 .map(|n| {
1125 let p = dir.join(n);
1126 std::path::absolute(&p).unwrap_or(p)
1127 })
1128 .collect()
1129 }
1130
1131 pub fn put(&self, task: &mut Task) -> Result<()> {
1140 let _lock = self.lock_task(&task.id)?;
1141 self.put_unlocked(task)
1142 }
1143
1144 fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1146 if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1147 for run in stored.runs {
1148 if !task.runs.contains(&run) {
1149 task.runs.push(run);
1150 }
1151 }
1152 }
1153 task.updated_at = Timestamp::now();
1154 std::fs::create_dir_all(&self.root)
1155 .with_context(|| format!("create {}", self.root.display()))?;
1156 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1157 let path = self.path_of(&task.id);
1158 let tmp = path.with_extension("json.tmp");
1159 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1160 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1161 if let (Some(notice), Some(home)) = (
1166 crate::notices::task_held(task),
1167 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1168 ) {
1169 crate::notices::raise_in(home, notice);
1170 }
1171 Ok(())
1172 }
1173
1174 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1179 let id = self.resolve_id(id)?;
1180 let _lock = self.lock_task(&id)?;
1183 let mut task = self.get(&id)?;
1184 if task.link_run(run) {
1185 self.put_unlocked(&mut task)?;
1186 }
1187 Ok(task)
1188 }
1189
1190 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1198 std::fs::create_dir_all(&self.root)
1199 .with_context(|| format!("create {}", self.root.display()))?;
1200 let path = self.root.join(format!("{id}.write-lock"));
1201 let started = std::time::Instant::now();
1202 loop {
1203 match std::fs::OpenOptions::new()
1204 .write(true)
1205 .create_new(true)
1206 .open(&path)
1207 {
1208 Ok(_) => return Ok(TaskLock(path)),
1209 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1210 let stale = std::fs::metadata(&path)
1211 .and_then(|m| m.modified())
1212 .ok()
1213 .and_then(|t| t.elapsed().ok())
1214 .is_some_and(|age| age > TASK_LOCK_STALE);
1215 if stale {
1216 let _ = std::fs::remove_file(&path);
1217 } else if started.elapsed() > TASK_LOCK_STALE {
1218 bail!("could not lock task {id}");
1219 } else {
1220 std::thread::sleep(std::time::Duration::from_millis(15));
1221 }
1222 }
1223 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1224 }
1225 }
1226 }
1227
1228 pub fn get(&self, id: &str) -> Result<Task> {
1230 let resolved = self.resolve_id(id)?;
1231 read_path(&self.path_of(&resolved))
1232 }
1233
1234 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1262 let resolved = self.resolve_id(id)?;
1263 if in_flight {
1264 bail!("task {resolved} is being run by a live daemon right now");
1265 }
1266 self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))?;
1267 let quarantined = self.quarantine_dependents_of(&resolved, questions);
1268 Ok(Removal {
1269 id: resolved,
1270 quarantined,
1271 })
1272 }
1273
1274 fn remove_record_with_attachments(
1278 &self,
1279 resolved: &str,
1280 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1281 ) -> Result<()> {
1282 self.sweep_removed_attachments();
1289 let attachments = self.attachments_dir(resolved);
1290 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1291 let moved = match std::fs::rename(&attachments, &aside) {
1292 Ok(()) => true,
1293 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1294 Err(e) => {
1295 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1296 }
1297 };
1298 let path = self.path_of(resolved);
1299 if let Err(e) = remove_record(&path) {
1300 if moved {
1301 let _ = std::fs::rename(&aside, &attachments);
1302 }
1303 return Err(e).with_context(|| format!("remove {}", path.display()));
1304 }
1305 if moved {
1306 if let Err(e) = std::fs::remove_dir_all(&aside) {
1307 tracing::warn!("leftover attachments {}: {e}", aside.display());
1308 }
1309 }
1310 let lock = self.lock_path(resolved);
1311 if let Err(e) = std::fs::remove_file(&lock) {
1312 if e.kind() != std::io::ErrorKind::NotFound {
1313 return Err(e).with_context(|| format!("remove {}", lock.display()));
1314 }
1315 }
1316 Ok(())
1317 }
1318
1319 fn sweep_removed_attachments(&self) {
1325 let Ok(entries) = std::fs::read_dir(&self.root) else {
1326 return;
1327 };
1328 for entry in entries.flatten() {
1329 let name = entry.file_name();
1330 let name = name.to_string_lossy();
1331 let Some(id) = name.strip_suffix(".attachments.removing") else {
1332 continue;
1333 };
1334 if !self.path_of(id).exists() {
1337 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1338 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1339 }
1340 }
1341 }
1342 }
1343
1344 fn quarantine_dependents_of(&self, dependency: &str, questions: &Questions) -> Vec<String> {
1348 let mut quarantined = Vec::new();
1349 for listed in self.list() {
1350 if listed.status != TaskStatus::Blocked
1351 || !listed.blocked_by.iter().any(|b| b == dependency)
1352 {
1353 continue;
1354 }
1355 let Ok(_claim) = self.claim(&listed.id) else {
1356 continue;
1357 };
1358 let Ok(mut task) = self.get(&listed.id) else {
1359 continue;
1360 };
1361 if task.status != TaskStatus::Blocked
1362 || !task.blocked_by.iter().any(|b| b == dependency)
1363 {
1364 continue;
1365 }
1366 let missing = missing_blockers(self, questions, &task.blocked_by);
1367 let language = crate::lang::of_repo(&task.repo);
1368 task.hold_machine(Some(missing_blocker_hold_reason_in(
1369 &task.blocked_by,
1370 &missing,
1371 &language,
1372 )));
1373 if self.put(&mut task).is_ok() {
1374 quarantined.push(task.id.clone());
1375 }
1376 }
1377 quarantined
1378 }
1379
1380 fn lock_path(&self, id: &str) -> PathBuf {
1383 self.root.join(format!("{id}.lock"))
1384 }
1385
1386 pub fn list(&self) -> Vec<Task> {
1399 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1400 .into_iter()
1401 .flatten()
1402 .flatten()
1403 .map(|e| e.path())
1404 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1405 .filter_map(|p| read_path(&p).ok())
1406 .collect();
1407 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1408 tasks
1409 }
1410
1411 pub fn superseded(&self) -> HashMap<String, String> {
1424 let mut by = HashMap::new();
1425 for task in self.list() {
1426 for earlier in &task.runs {
1427 if let Some(later) = task.successor_of(earlier) {
1428 by.insert(earlier.clone(), later.clone());
1429 }
1430 }
1431 }
1432 by
1433 }
1434
1435 pub fn superseded_by(&self, run: &str) -> Option<String> {
1443 for task in self.list() {
1444 if task.runs.iter().any(|r| r == run) {
1445 return task.successor_of(run).cloned();
1446 }
1447 }
1448 None
1449 }
1450
1451 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1462 for task in self.list() {
1463 if task.runs.iter().any(|r| r == run) {
1464 return task.runs.last().filter(|last| **last != run).cloned();
1465 }
1466 }
1467 None
1468 }
1469
1470 pub fn next_runnable(&self) -> Option<Task> {
1475 let mut runnable: Vec<Task> = self
1476 .list()
1477 .into_iter()
1478 .filter(|t| t.status.runnable())
1479 .collect();
1480 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1481 runnable.into_iter().next()
1482 }
1483
1484 pub fn claim(&self, id: &str) -> Result<Claim> {
1491 std::fs::create_dir_all(&self.root)
1492 .with_context(|| format!("create {}", self.root.display()))?;
1493 let path = self.lock_path(id);
1494 match std::fs::OpenOptions::new()
1495 .write(true)
1496 .create_new(true)
1497 .open(&path)
1498 {
1499 Ok(mut f) => {
1500 use std::io::Write as _;
1501 let _ = writeln!(f, "{}", std::process::id());
1503 Ok(Claim { path })
1504 }
1505 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1506 bail!("task {id} is already claimed ({} exists)", path.display())
1507 }
1508 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1509 }
1510 }
1511
1512 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1514 if self.path_of(prefix).is_file() {
1515 return Ok(prefix.to_owned());
1516 }
1517 let hits: Vec<String> = self
1518 .list()
1519 .into_iter()
1520 .map(|t| t.id)
1521 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1522 .collect();
1523 match hits.len() {
1524 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1525 0 => bail!("no task matches `{prefix}`"),
1526 _ => bail!(
1527 "`{prefix}` matches {} tasks: {}",
1528 hits.len(),
1529 hits.join(", ")
1530 ),
1531 }
1532 }
1533
1534 pub fn revision(&self) -> u64 {
1541 use std::hash::{Hash as _, Hasher as _};
1542
1543 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
1544 .into_iter()
1545 .flatten()
1546 .flatten()
1547 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
1548 .filter_map(|e| {
1549 let name = e.file_name().to_string_lossy().into_owned();
1550 let mtime = e
1551 .metadata()
1552 .ok()?
1553 .modified()
1554 .ok()?
1555 .duration_since(std::time::UNIX_EPOCH)
1556 .ok()?
1557 .as_millis() as u64;
1558 Some((name, mtime))
1559 })
1560 .collect();
1561
1562 if entries.is_empty() {
1563 return 0;
1564 }
1565
1566 entries.sort_unstable();
1567 let mut hasher = std::hash::DefaultHasher::new();
1568 for (name, mtime) in &entries {
1569 name.hash(&mut hasher);
1570 mtime.hash(&mut hasher);
1571 }
1572 let h = hasher.finish();
1573 if h == 0 { 1 } else { h }
1574 }
1575}
1576
1577#[derive(Debug, Clone)]
1579pub struct Removal {
1580 pub id: String,
1582 pub quarantined: Vec<String>,
1586}
1587
1588#[derive(Debug)]
1590pub struct Claim {
1591 path: PathBuf,
1592}
1593
1594impl Drop for Claim {
1595 fn drop(&mut self) {
1596 let _ = std::fs::remove_file(&self.path);
1597 }
1598}
1599
1600pub fn title_from(instruction: &str, max: usize) -> String {
1603 let line = instruction
1609 .lines()
1610 .map(str::trim)
1611 .find(|l| !l.is_empty())
1612 .unwrap_or("(empty task)")
1613 .trim_start_matches(['#', '-', '*', '>', ' '])
1614 .trim();
1615 if line.is_empty() {
1616 return "(empty task)".to_owned();
1617 }
1618 if line.chars().count() <= max {
1619 return line.to_owned();
1620 }
1621 let head: String = line.chars().take(max.saturating_sub(1)).collect();
1622 format!("{head}…")
1623}
1624
1625fn read_path(path: &Path) -> Result<Task> {
1626 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1627 let task: Task =
1628 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1629 if task.schema > SCHEMA {
1635 bail!(
1636 "task {} was written by a different magi (schema {}, this build \
1637 speaks {SCHEMA})",
1638 task.id,
1639 task.schema
1640 );
1641 }
1642 Ok(task)
1643}
1644
1645pub fn missing_blockers(
1659 queue: &Queue,
1660 questions: &Questions,
1661 blocked_by: &[String],
1662) -> Vec<String> {
1663 blocked_by
1664 .iter()
1665 .filter(|id| !queue.path_of(id).is_file() && !questions.path_of(id).is_file())
1666 .cloned()
1667 .collect()
1668}
1669
1670pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
1682 missing_blocker_hold_reason_in(blocked_by, missing, "en")
1683}
1684
1685pub fn missing_blocker_hold_reason_in(
1687 blocked_by: &[String],
1688 missing: &[String],
1689 language: &str,
1690) -> String {
1691 if crate::lang::is_japanese(language) {
1692 format!(
1693 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
1694 blocked_by.join(", "),
1695 missing.join(", "),
1696 )
1697 } else {
1698 format!(
1699 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
1700 blocked_by.join(", "),
1701 missing.join(", "),
1702 )
1703 }
1704}
1705
1706pub fn short(id: &str) -> &str {
1708 id.split('-').next_back().unwrap_or(id)
1709}
1710
1711fn new_id() -> String {
1712 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1713 let seed = crate::rng::entropy();
1714 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1715}
1716
1717#[cfg(test)]
1718mod tests {
1719 #[test]
1720 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
1721 let mut t = task("relanded");
1722 t.attempts = 1;
1723 t.status = TaskStatus::Running;
1724 t.already_landed("already in main as 0e368de");
1725 assert_eq!(t.status, TaskStatus::Done);
1726 assert_eq!(t.attempts, 0);
1727 assert!(t.hold_reason.is_none());
1728 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
1729 }
1730
1731 #[test]
1732 fn missing_blocker_reason_follows_the_language() {
1733 let b = vec!["a".to_owned()];
1734 let en = missing_blocker_hold_reason_in(&b, &b, "en");
1735 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
1736 assert!(en.starts_with("blocked on a"));
1737 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
1738 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
1739 }
1740
1741 use super::*;
1742
1743 #[test]
1744 fn triage_applied_survives_release_and_old_records_read_as_empty() {
1745 let mut t = Task::new(
1746 "t".to_owned(),
1747 "i".to_owned(),
1748 PathBuf::from("r"),
1749 Source::Human,
1750 );
1751 t.mark_triage_applied("q1");
1752 t.mark_triage_applied("q1");
1753 t.hold_machine(Some("x".to_owned()));
1754 t.release();
1755 assert_eq!(t.triage_applied, ["q1"]);
1756 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
1757
1758 let mut v = serde_json::to_value(&t).unwrap();
1759 v.as_object_mut().unwrap().remove("triage_applied");
1760 let old: Task = serde_json::from_value(v).unwrap();
1761 assert!(old.triage_applied.is_empty());
1762 }
1763
1764 #[test]
1765 fn task_counts_of_empty_is_all_zero() {
1766 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
1767 }
1768
1769 #[test]
1770 fn task_counts_of_tallies_every_status() {
1771 let mut queued = Task::new(
1772 "q".to_owned(),
1773 "i".to_owned(),
1774 PathBuf::from("."),
1775 Source::Human,
1776 );
1777 queued.status = TaskStatus::Queued;
1778 let mut running = queued.clone();
1779 running.status = TaskStatus::Running;
1780 let mut done = queued.clone();
1781 done.status = TaskStatus::Done;
1782 let mut failed = queued.clone();
1783 failed.status = TaskStatus::Failed;
1784 let mut held = queued.clone();
1785 held.status = TaskStatus::Held;
1786 let mut blocked = queued.clone();
1787 blocked.status = TaskStatus::Blocked;
1788
1789 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
1790 assert_eq!(
1791 counts,
1792 TaskCounts {
1793 queued: 1,
1794 running: 1,
1795 done: 2,
1796 failed: 1,
1797 held: 1,
1798 blocked: 1,
1799 }
1800 );
1801 }
1802
1803 fn queue() -> (tempfile::TempDir, Queue) {
1806 let dir = tempfile::tempdir().unwrap();
1807 let q = Queue::at(dir.path().join("queue"));
1808 (dir, q)
1809 }
1810
1811 #[test]
1812 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
1813 let dir = tempfile::tempdir().unwrap();
1814 let q = Queue::at(dir.path().join("queue"));
1815 let mut t = task("held");
1816 q.put(&mut t).unwrap();
1817 assert_eq!(
1818 crate::notices::Notices::at(dir.path().join("notifications"))
1819 .list()
1820 .len(),
1821 0
1822 );
1823 t.hold_machine(Some("out of attempts".to_owned()));
1824 q.put(&mut t).unwrap();
1825 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
1826 assert_eq!(listed.len(), 1);
1827 assert!(listed[0].message.contains("out of attempts"));
1828 }
1829
1830 fn task(title: &str) -> Task {
1831 Task::new(
1832 title.to_owned(),
1833 format!("do {title}"),
1834 PathBuf::from("."),
1835 Source::Human,
1836 )
1837 }
1838
1839 #[test]
1840 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
1841 let mut t = task("retried");
1842 assert!(
1843 t.earlier_attempts().is_empty(),
1844 "a first attempt takes nothing over"
1845 );
1846 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1847 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
1848 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
1850 assert_eq!(t.successor_of("bbbb"), None);
1851 assert_eq!(t.successor_of("zzzz"), None);
1852 }
1853
1854 #[test]
1855 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
1856 let (_dir, q) = queue();
1857 let mut t = task("retried");
1858 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1859 q.put(&mut t).unwrap();
1860
1861 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
1862 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
1863 assert_eq!(
1864 q.superseded_by("cccc"),
1865 None,
1866 "the latest attempt replaces nothing"
1867 );
1868 assert_eq!(
1869 q.superseded_by("never-heard-of-it"),
1870 None,
1871 "a run belonging to no task on this queue is not superseded"
1872 );
1873
1874 let mut by = HashMap::new();
1875 by.insert("aaaa".to_owned(), "bbbb".to_owned());
1876 by.insert("bbbb".to_owned(), "cccc".to_owned());
1877 assert_eq!(
1878 q.superseded(),
1879 by,
1880 "the whole-map and single-run forms must agree"
1881 );
1882 }
1883
1884 #[test]
1885 fn superseded_attempts_is_empty_until_the_task_is_done() {
1886 let mut t = task("retried");
1887 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1888 t.status = TaskStatus::Failed;
1889 assert_eq!(
1890 t.superseded_attempts(true),
1891 &[] as &[String],
1892 "a task still retrying has no attempt yet that a later one made moot"
1893 );
1894
1895 t.status = TaskStatus::Running;
1896 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1897 }
1898
1899 #[test]
1900 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
1901 let mut t = task("retried");
1902 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1903 t.status = TaskStatus::Done;
1904 assert_eq!(
1905 t.superseded_attempts(true),
1906 &["aaaa".to_owned(), "bbbb".to_owned()],
1907 "cccc is the attempt whose success made the task done, and stays out"
1908 );
1909 }
1910
1911 #[test]
1912 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
1913 let mut t = task("first try landed");
1914 t.runs = vec!["aaaa".to_owned()];
1915 t.status = TaskStatus::Done;
1916 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1917 }
1918
1919 #[test]
1920 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
1921 let mut t = task("closed by hand after a manual merge");
1928 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1929 t.status = TaskStatus::Done;
1930 assert_eq!(
1931 t.superseded_attempts(false),
1932 &[] as &[String],
1933 "nothing here is provably why the task is done, so nothing is superseded"
1934 );
1935 }
1936
1937 #[test]
1938 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
1939 let (_dir, q) = queue();
1940 let mut t = task("retried twice");
1941 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1942 q.put(&mut t).unwrap();
1943
1944 assert_eq!(
1945 q.latest_attempt("aaaa"),
1946 Some("cccc".to_owned()),
1947 "an old attempt points straight at the chain's current head, not the \
1948 next attempt in the middle of it"
1949 );
1950 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
1951 assert_eq!(
1952 q.latest_attempt("cccc"),
1953 None,
1954 "the latest attempt is not superseded by anything"
1955 );
1956 assert_eq!(
1957 q.latest_attempt("never-heard-of-it"),
1958 None,
1959 "a run belonging to no task on this queue is not superseded"
1960 );
1961 }
1962
1963 #[test]
1964 fn a_markdown_heading_is_the_title_not_decoration() {
1965 assert_eq!(
1970 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
1971 "Rework the config loader"
1972 );
1973 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
1974 assert_eq!(title_from("> quoted task", 40), "quoted task");
1975 assert_eq!(title_from(" \n\n", 40), "(empty task)");
1977 assert_eq!(title_from("###\n", 40), "(empty task)");
1978 }
1979
1980 #[test]
1981 fn a_long_title_is_elided_by_characters_not_bytes() {
1982 let long = "課題".repeat(30);
1984 let title = title_from(&long, 10);
1985 assert_eq!(title.chars().count(), 10);
1986 assert!(title.ends_with('…'));
1987 }
1988
1989 #[test]
1990 fn priority_wins_and_ties_break_oldest_first() {
1991 let (_dir, q) = queue();
1992 let mut a = task("first");
1993 let mut b = task("second");
1994 let mut c = task("urgent");
1995 a.id = "20260101-000001-aaaa".to_owned();
1997 b.id = "20260101-000002-bbbb".to_owned();
1998 c.id = "20260101-000003-cccc".to_owned();
1999 c.priority = 5;
2000 for t in [&mut a, &mut b, &mut c] {
2001 q.put(t).unwrap();
2002 }
2003
2004 assert_eq!(q.next_runnable().unwrap().id, c.id);
2006 c.hold_machine(None);
2007 q.put(&mut c).unwrap();
2008 assert_eq!(q.next_runnable().unwrap().id, a.id);
2010 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2011 }
2012
2013 #[test]
2014 fn a_blocked_task_never_starves_another_runnable_one() {
2015 let (_dir, q) = queue();
2016 let mut blocked = task("blocked");
2017 blocked.block(vec!["something".to_owned()], None);
2018 q.put(&mut blocked).unwrap();
2019
2020 let mut runnable = task("free to go");
2021 q.put(&mut runnable).unwrap();
2022
2023 let next = q.next_runnable().expect("a runnable task is still offered");
2024 assert_eq!(next.id, runnable.id);
2025 }
2026
2027 #[test]
2028 fn a_held_task_is_never_offered_to_the_loop() {
2029 let (_dir, q) = queue();
2030 let mut t = task("held");
2031 q.put(&mut t).unwrap();
2032 assert!(q.next_runnable().is_some());
2033
2034 t.hold_machine(None);
2035 q.put(&mut t).unwrap();
2036 assert!(
2037 q.next_runnable().is_none(),
2038 "a held task must wait for a human"
2039 );
2040
2041 t.status = TaskStatus::Failed;
2043 q.put(&mut t).unwrap();
2044 assert!(q.next_runnable().is_some());
2045 }
2046
2047 #[test]
2048 fn attempts_are_capped_and_then_the_task_is_held() {
2049 let mut t = task("doomed");
2050
2051 t.start("run-1".to_owned());
2052 t.fail("gate red", 2);
2053 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2054
2055 t.start("run-2".to_owned());
2056 t.fail("gate red", 2);
2057 assert_eq!(
2058 t.status,
2059 TaskStatus::Held,
2060 "out of attempts: stop spending money on it"
2061 );
2062 assert_eq!(t.runs, ["run-1", "run-2"]);
2063 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2064 assert_eq!(
2065 t.hold_reason.as_deref(),
2066 Some("gate red"),
2067 "the hold must say why, not leave hold_reason null next to a \
2068 populated last_error"
2069 );
2070 }
2071
2072 #[test]
2073 fn handing_off_a_task_records_a_hold_reason_too() {
2074 let mut t = task("left a pull request");
2075 t.start("run-1".to_owned());
2076 t.handed_off("run ended with a pull request open [run run-1]");
2077 assert_eq!(t.status, TaskStatus::Held);
2078 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2079 assert_eq!(
2080 t.hold_reason.as_deref(),
2081 Some("run ended with a pull request open [run run-1]")
2082 );
2083 assert_eq!(t.hold_reason, t.last_error);
2084 }
2085
2086 #[test]
2087 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2088 let mut t = task("stalled by quota");
2089
2090 t.start("run-1".to_owned());
2091 assert_eq!(t.attempts, 1);
2092 t.stall("judge-1, judge-2 out of quota");
2093 assert_eq!(
2094 t.attempts, 0,
2095 "a closed quota window must not spend the task's retry budget"
2096 );
2097 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2098 assert_eq!(
2099 t.last_error.as_deref(),
2100 Some("judge-1, judge-2 out of quota")
2101 );
2102
2103 for _ in 0..20 {
2106 t.start("run-n".to_owned());
2107 t.stall("still out of quota");
2108 }
2109 t.start("run-real".to_owned());
2110 t.fail("gate red", 2);
2111 assert_eq!(
2112 t.status,
2113 TaskStatus::Failed,
2114 "the first attempt that was really judged is attempt one"
2115 );
2116 }
2117
2118 #[test]
2119 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2120 let mut t = task("retry me");
2121 t.start("run-1".to_owned());
2122 t.fail("gate red", 1);
2123 assert_eq!(t.status, TaskStatus::Held);
2124
2125 t.release();
2126 assert_eq!(t.status, TaskStatus::Queued);
2127 assert_eq!(t.attempts, 0);
2130 assert!(t.last_error.is_none());
2131 assert_eq!(
2132 t.runs.len(),
2133 1,
2134 "history is kept: attempts reset, evidence does not"
2135 );
2136 }
2137
2138 #[test]
2139 fn a_hold_reason_survives_and_a_release_clears_it() {
2140 let mut t = task("waiting on something else");
2141 t.hold_manual(Some(
2142 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2143 ));
2144 assert_eq!(t.status, TaskStatus::Held);
2145 assert_eq!(
2146 t.hold_reason.as_deref(),
2147 Some("waiting for 20260101-000000-aaaa to land first")
2148 );
2149
2150 t.hold_manual(None);
2152 assert_eq!(
2153 t.hold_reason.as_deref(),
2154 Some("waiting for 20260101-000000-aaaa to land first"),
2155 "a bare re-hold keeps whatever a human already wrote down"
2156 );
2157
2158 let mut plain = task("no reason given");
2160 plain.hold_manual(None);
2161 assert_eq!(plain.status, TaskStatus::Held);
2162 assert!(plain.hold_reason.is_none());
2163
2164 t.release();
2165 assert_eq!(t.status, TaskStatus::Queued);
2166 assert!(
2167 t.hold_reason.is_none(),
2168 "a stale reason must not greet the next person who holds this task"
2169 );
2170 }
2171
2172 #[test]
2173 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2174 let mut t = task("landed by hand while held");
2179 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2180 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2181
2182 t.succeed();
2183 assert_eq!(t.status, TaskStatus::Done);
2184 assert!(
2185 t.hold_reason.is_none(),
2186 "a done task cannot still be waiting on something"
2187 );
2188 }
2189
2190 #[test]
2191 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2192 let mut held = task("held straight out of blocked");
2199 held.block(
2200 vec!["20260101-000000-dead".to_owned()],
2201 Some("waiting on the migration script".to_owned()),
2202 );
2203 assert_eq!(held.status, TaskStatus::Blocked);
2204
2205 held.hold_manual(None);
2206 assert_eq!(held.status, TaskStatus::Held);
2207 assert!(
2208 held.blocked_by.is_empty(),
2209 "hold overrides the wait, same as release"
2210 );
2211 assert!(held.block_reason.is_none());
2212
2213 let mut done = task("closed straight out of blocked");
2214 done.block(
2215 vec!["20260101-000000-dead".to_owned()],
2216 Some("waiting on the migration script".to_owned()),
2217 );
2218 done.succeed();
2219 assert_eq!(done.status, TaskStatus::Done);
2220 assert!(
2221 done.blocked_by.is_empty(),
2222 "a done task cannot still be waiting on a dependency"
2223 );
2224 assert!(done.block_reason.is_none());
2225 }
2226
2227 #[test]
2228 fn a_blocked_task_is_never_offered_to_the_loop() {
2229 let mut t = task("blocked");
2230 assert!(t.status.runnable());
2231 t.block(
2232 vec!["dep-id".to_owned()],
2233 Some("waits on dep-id".to_owned()),
2234 );
2235 assert_eq!(t.status, TaskStatus::Blocked);
2236 assert!(!t.status.runnable());
2237 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2238 }
2239
2240 #[test]
2241 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2242 let mut t = task("blocked on two");
2243 t.block(
2244 vec!["a".to_owned(), "b".to_owned()],
2245 Some("waits on a and b".to_owned()),
2246 );
2247
2248 t.unblock("a");
2249 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2250 assert_eq!(t.blocked_by, ["b"]);
2251
2252 t.unblock("b");
2253 assert_eq!(t.status, TaskStatus::Queued);
2254 assert!(t.blocked_by.is_empty());
2255 assert!(t.block_reason.is_none());
2256 }
2257
2258 #[test]
2259 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2260 let mut t = task("never blocked");
2261 t.unblock("whatever");
2262 assert_eq!(t.status, TaskStatus::Queued);
2263 }
2264
2265 #[test]
2266 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2267 let mut t = task("held, then asked about");
2272 t.hold_machine(Some("out of attempts".to_owned()));
2273 assert_eq!(t.status, TaskStatus::Held);
2274
2275 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2276 assert_eq!(t.status, TaskStatus::Blocked);
2277
2278 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2279 t.unblock("q1");
2280 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2281 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2282 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2283 assert!(t.blocked_from.is_none(), "consumed once restored");
2284 }
2285
2286 #[test]
2287 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2288 let mut t = task("manually held, then asked about");
2289 t.hold_manual(Some("waiting on a dependency".to_owned()));
2290
2291 t.block(vec!["q1".to_owned()], None);
2292 t.unblock("q1");
2293
2294 assert_eq!(t.status, TaskStatus::Held);
2295 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2296 }
2297
2298 #[test]
2299 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2300 let mut t = task("held, blocked twice");
2304 t.hold_machine(None);
2305 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2306 t.block(
2307 vec!["q1".to_owned(), "q2".to_owned()],
2308 Some("second".to_owned()),
2309 );
2310
2311 t.unblock("q1");
2312 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2313 t.unblock("q2");
2314 assert_eq!(t.status, TaskStatus::Held);
2315 }
2316
2317 #[test]
2318 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2319 let mut t = task("blocked mid-run");
2323 t.start("run-1".to_owned());
2324 assert_eq!(t.status, TaskStatus::Running);
2325
2326 t.block(vec!["q1".to_owned()], None);
2327 t.unblock("q1");
2328 assert_eq!(t.status, TaskStatus::Queued);
2329 }
2330
2331 #[test]
2332 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2333 let mut t = task("legacy record, held before it was blocked");
2339 t.hold_source = Some(HoldSource::Machine);
2340 t.hold_reason = Some("legacy hold reason".to_owned());
2341 t.status = TaskStatus::Blocked;
2342 t.blocked_by = vec!["q1".to_owned()];
2343 t.blocked_from = None;
2344
2345 t.unblock("q1");
2346 assert_eq!(t.status, TaskStatus::Held);
2347 }
2348
2349 #[test]
2350 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2351 let mut t = task("legacy record, ordinary dependency block");
2352 t.status = TaskStatus::Blocked;
2353 t.blocked_by = vec!["dep".to_owned()];
2354 t.blocked_from = None;
2355
2356 t.unblock("dep");
2357 assert_eq!(t.status, TaskStatus::Queued);
2358 }
2359
2360 #[test]
2361 fn answering_a_question_is_recorded_and_survives_a_release() {
2362 let mut t = task("asked something");
2363 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2364 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2365 t.unblock("q1");
2366 assert_eq!(t.status, TaskStatus::Queued);
2367 assert_eq!(t.answers.len(), 1);
2368 assert_eq!(t.answers[0].answer, "SQLite");
2369
2370 t.release();
2374 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2375 }
2376
2377 #[test]
2378 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2379 let mut t = task("blocked run with a surviving branch");
2380 t.start("run-1".to_owned());
2381 t.fail("blocked with major findings", 5);
2382 assert_eq!(t.status, TaskStatus::Failed);
2383
2384 t.request_review("magi/eba2/A".to_owned());
2385 assert_eq!(t.status, TaskStatus::Queued);
2386 assert_eq!(t.attempts, 0);
2387 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2388
2389 t.release();
2391 assert!(t.review_branch.is_none());
2392 }
2393
2394 #[test]
2395 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2396 let mut t = task("retry");
2397 t.start("run-1".to_owned());
2398 t.requeue();
2399 assert!(t.fresh_start);
2400
2401 t.release();
2402 assert!(!t.fresh_start);
2403 }
2404
2405 #[test]
2406 fn priority_can_be_changed_while_queued_but_not_while_running() {
2407 let mut t = task("reprioritise me");
2408 t.set_priority(5).unwrap();
2409 assert_eq!(t.priority, 5);
2410
2411 t.start("run-1".to_owned());
2412 let err = t.set_priority(9).unwrap_err().to_string();
2413 assert!(err.contains("running"), "{err}");
2414 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2415 }
2416
2417 #[test]
2418 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2419 let mut t = task("interrupt me");
2420 assert!(!t.interrupt, "off unless asked, same as any other task");
2421
2422 t.set_interrupt(true).unwrap();
2423 assert!(t.interrupt);
2424
2425 t.start("run-1".to_owned());
2426 assert!(
2427 !t.interrupt,
2428 "the mark is one-shot: dispatching the task fulfils it, \
2429 whatever the run that follows ends up doing"
2430 );
2431 let err = t.set_interrupt(true).unwrap_err().to_string();
2432 assert!(err.contains("running"), "{err}");
2433 t.set_interrupt(false).unwrap();
2436 assert!(!t.interrupt);
2437 }
2438
2439 #[test]
2443 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2444 let mut t = task("interrupt me");
2445 t.set_interrupt(true).unwrap();
2446 t.start("run-1".to_owned());
2447 t.fail("mock failure", 5);
2448 assert_eq!(t.status, TaskStatus::Failed);
2449 assert!(
2450 !t.interrupt,
2451 "one attempt already spent the mark; a retry is an ordinary \
2452 requeue, not a fresh interrupt request"
2453 );
2454 }
2455
2456 #[test]
2457 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2458 let (_dir, q) = queue();
2459 let mut a = task("first filed");
2460 let mut b = task("second filed");
2461 a.id = "20260101-000001-aaaa".to_owned();
2462 b.id = "20260101-000002-bbbb".to_owned();
2463 q.put(&mut a).unwrap();
2464 q.put(&mut b).unwrap();
2465
2466 assert_eq!(
2467 q.next_runnable().unwrap().id,
2468 a.id,
2469 "with equal priority the older task goes first, so a burst of \
2470 new work cannot starve it"
2471 );
2472 assert_eq!(
2473 q.list()[0].id,
2474 b.id,
2475 "but the list an operator reads is newest first, the same as \
2476 before priority existed - a's turn to run does not make it the \
2477 newest task"
2478 );
2479
2480 let mut a = q.get(&a.id).unwrap();
2481 a.set_priority(10).unwrap();
2482 q.put(&mut a).unwrap();
2483
2484 assert_eq!(
2485 q.next_runnable().unwrap().id,
2486 a.id,
2487 "a raised priority must be reflected the moment it is saved"
2488 );
2489 assert_eq!(
2493 q.list()[0].id,
2494 a.id,
2495 "the raised task must sort first in the list an operator reads, \
2496 not only in next_runnable's own ordering"
2497 );
2498 }
2499
2500 #[test]
2501 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2502 let mut t = Task::new(
2503 "old title".to_owned(),
2504 "old instruction".to_owned(),
2505 PathBuf::from("/repo"),
2506 Source::Agent {
2507 run: "20260101-000000-beef".to_owned(),
2508 node: "implement".to_owned(),
2509 },
2510 );
2511 let id = t.id.clone();
2512 let created_at = t.created_at;
2513 t.runs.push("20260101-000000-beef".to_owned());
2514
2515 t.edit("new title".to_owned(), "new instruction".to_owned())
2516 .unwrap();
2517
2518 assert_eq!(t.title, "new title");
2519 assert_eq!(t.instruction, "new instruction");
2520 assert_eq!(t.id, id, "editing must not mint a new id");
2521 assert_eq!(t.created_at, created_at);
2522 assert_eq!(
2523 t.source,
2524 Source::Agent {
2525 run: "20260101-000000-beef".to_owned(),
2526 node: "implement".to_owned(),
2527 },
2528 "editing must not turn agent attribution into human"
2529 );
2530 assert_eq!(t.runs, ["20260101-000000-beef"]);
2531 }
2532
2533 #[test]
2534 fn editing_is_refused_once_a_task_is_running_or_finished() {
2535 let mut running = task("in flight");
2536 running.start("run-1".to_owned());
2537 let err = running
2538 .edit("x".to_owned(), "y".to_owned())
2539 .unwrap_err()
2540 .to_string();
2541 assert!(err.contains("running"), "{err}");
2542
2543 let mut done = task("finished");
2544 done.succeed();
2545 let err = done
2546 .edit("x".to_owned(), "y".to_owned())
2547 .unwrap_err()
2548 .to_string();
2549 assert!(err.contains("done"), "{err}");
2550
2551 let mut queued = task("waiting");
2553 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2554 let mut held = task("parked");
2555 held.hold_machine(None);
2556 held.edit("x".to_owned(), "y".to_owned()).unwrap();
2557 }
2558
2559 #[test]
2560 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2561 let (_dir, q) = queue();
2562 let path = q.path_of("20260101-000000-aaaa");
2563 std::fs::create_dir_all(q.root()).unwrap();
2564 std::fs::write(
2565 &path,
2566 serde_json::json!({
2567 "schema": SCHEMA,
2568 "id": "20260101-000000-aaaa",
2569 "title": "from before hold reasons existed",
2570 "instruction": "from before hold reasons existed",
2571 "repo": ".",
2572 "source": { "kind": "human" },
2573 "status": "held",
2574 "created_at": Timestamp::now().to_string(),
2575 "updated_at": Timestamp::now().to_string(),
2576 })
2577 .to_string(),
2578 )
2579 .unwrap();
2580
2581 let task = q.get("20260101-000000-aaaa").expect("must still read");
2582 assert!(task.hold_reason.is_none());
2583 assert!(task.operator_held());
2584 }
2585
2586 #[test]
2587 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
2588 let (_dir, q) = queue();
2589 let path = q.path_of("20260101-000000-bbbb");
2590 std::fs::create_dir_all(q.root()).unwrap();
2591 std::fs::write(
2592 &path,
2593 serde_json::json!({
2594 "schema": 2,
2595 "id": "20260101-000000-bbbb",
2596 "title": "old manual recovery",
2597 "instruction": "old manual recovery",
2598 "repo": ".",
2599 "source": { "kind": "human" },
2600 "status": "held",
2601 "hold_reason": "active manual recovery run20260912-224242-daf5",
2602 "created_at": Timestamp::now().to_string(),
2603 "updated_at": Timestamp::now().to_string(),
2604 })
2605 .to_string(),
2606 )
2607 .unwrap();
2608
2609 let task = q.get("20260101-000000-bbbb").expect("must still read");
2610 assert_eq!(task.hold_source, None);
2611 assert!(task.operator_held());
2612 }
2613
2614 #[test]
2615 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
2616 let (_dir, q) = queue();
2617 let path = q.path_of("20260101-000000-aaaa");
2618 std::fs::create_dir_all(q.root()).unwrap();
2619 std::fs::write(
2620 &path,
2621 serde_json::json!({
2622 "schema": SCHEMA,
2623 "id": "20260101-000000-aaaa",
2624 "title": "from before diagnostics existed",
2625 "instruction": "from before diagnostics existed",
2626 "repo": ".",
2627 "source": { "kind": "human" },
2628 "status": "held",
2629 "created_at": Timestamp::now().to_string(),
2630 "updated_at": Timestamp::now().to_string(),
2631 })
2632 .to_string(),
2633 )
2634 .unwrap();
2635
2636 let task = q.get("20260101-000000-aaaa").expect("must still read");
2637 assert!(task.diagnostic.is_none());
2638 }
2639
2640 #[test]
2641 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
2642 let (_dir, q) = queue();
2646 let path = q.path_of("20260101-000000-aaaa");
2647 std::fs::create_dir_all(q.root()).unwrap();
2648 std::fs::write(
2649 &path,
2650 serde_json::json!({
2651 "schema": 1,
2652 "id": "20260101-000000-aaaa",
2653 "title": "from before blocking existed",
2654 "instruction": "from before blocking existed",
2655 "repo": ".",
2656 "source": { "kind": "human" },
2657 "status": "queued",
2658 "created_at": Timestamp::now().to_string(),
2659 "updated_at": Timestamp::now().to_string(),
2660 })
2661 .to_string(),
2662 )
2663 .unwrap();
2664
2665 let task = q.get("20260101-000000-aaaa").expect("must still read");
2666 assert!(task.blocked_by.is_empty());
2667 assert!(task.block_reason.is_none());
2668 assert!(task.answers.is_empty());
2669 assert!(task.review_branch.is_none());
2670 }
2671
2672 #[test]
2673 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
2674 let mut held = task("diagnosed");
2679 held.start("run-1".to_owned());
2680 held.fail("gate red", 1);
2681 held.diagnostic = Some("cargo test failed: ...".to_owned());
2682 assert_eq!(held.status, TaskStatus::Held);
2683
2684 held.release();
2685 assert!(held.diagnostic.is_none());
2686
2687 held.diagnostic = Some("cargo test failed: ...".to_owned());
2688 held.succeed();
2689 assert!(held.diagnostic.is_none());
2690 }
2691
2692 #[test]
2693 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
2694 let mut t = task("retried");
2695 t.start("run-1".to_owned());
2696 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
2697 t.fail("unrelated config error", 5);
2698 assert_eq!(t.status, TaskStatus::Failed);
2699 assert!(
2700 t.diagnostic.is_none(),
2701 "fail() must not let an old diagnostic outlive the run that produced it"
2702 );
2703 }
2704
2705 #[test]
2706 fn a_claim_is_exclusive_and_releases_on_drop() {
2707 let (_dir, q) = queue();
2708 let mut t = task("contended");
2709 q.put(&mut t).unwrap();
2710
2711 let held = q.claim(&t.id).unwrap();
2712 assert!(
2713 q.claim(&t.id).is_err(),
2714 "two daemons must not drive one task into two runs"
2715 );
2716 drop(held);
2717 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
2718 }
2719
2720 #[test]
2721 fn a_round_trip_survives_disk() {
2722 let (_dir, q) = queue();
2723 let mut t = Task::new(
2724 "titled".to_owned(),
2725 "body".to_owned(),
2726 PathBuf::from("/repo"),
2727 Source::Agent {
2728 run: "20260101-000000-beef".to_owned(),
2729 node: "implement".to_owned(),
2730 },
2731 );
2732 t.priority = 3;
2733 q.put(&mut t).unwrap();
2734
2735 let back = q.get(&t.id).unwrap();
2736 assert_eq!(back.id, t.id);
2737 assert_eq!(back.priority, 3);
2738 assert_eq!(back.source.label(), "implement@beef");
2739 assert_eq!(q.get(t.short()).unwrap().id, t.id);
2741 }
2742
2743 #[test]
2744 fn an_unreadable_task_does_not_take_the_queue_down() {
2745 let (_dir, q) = queue();
2746 let mut t = task("fine");
2747 q.put(&mut t).unwrap();
2748 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
2749
2750 let listed = q.list();
2751 assert_eq!(listed.len(), 1, "the readable task still lists");
2752 assert_eq!(listed[0].id, t.id);
2753 }
2754
2755 #[test]
2756 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
2757 let (_dir, q) = queue();
2758 let path = q.path_of("20260101-000000-aaaa");
2759 std::fs::create_dir_all(q.root()).unwrap();
2760 std::fs::write(
2761 &path,
2762 serde_json::json!({
2763 "schema": SCHEMA,
2764 "id": "20260101-000000-aaaa",
2765 "title": "from before solo existed",
2766 "instruction": "from before solo existed",
2767 "repo": ".",
2768 "source": { "kind": "human" },
2769 "status": "queued",
2770 "created_at": Timestamp::now().to_string(),
2771 "updated_at": Timestamp::now().to_string(),
2772 })
2773 .to_string(),
2774 )
2775 .unwrap();
2776
2777 let task = q.get("20260101-000000-aaaa").expect("must still read");
2778 assert!(!task.solo, "a queue file with no `solo` field means false");
2779 }
2780
2781 #[test]
2782 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
2783 let (_dir, q) = queue();
2784 let path = q.path_of("20260101-000000-bbbb");
2785 std::fs::create_dir_all(q.root()).unwrap();
2786 std::fs::write(
2787 &path,
2788 serde_json::json!({
2789 "schema": SCHEMA,
2790 "id": "20260101-000000-bbbb",
2791 "title": "from before urgent existed",
2792 "instruction": "from before urgent existed",
2793 "repo": ".",
2794 "source": { "kind": "human" },
2795 "status": "queued",
2796 "created_at": Timestamp::now().to_string(),
2797 "updated_at": Timestamp::now().to_string(),
2798 })
2799 .to_string(),
2800 )
2801 .unwrap();
2802
2803 let task = q.get("20260101-000000-bbbb").expect("must still read");
2804 assert!(
2805 !task.urgent,
2806 "a queue file with no `urgent` field means false, same as `solo`"
2807 );
2808 }
2809
2810 #[test]
2811 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
2812 let (_dir, q) = queue();
2813 let mut t = task("from the future");
2814 q.put(&mut t).unwrap();
2815 let path = q.path_of(&t.id);
2816 let body = std::fs::read_to_string(&path)
2817 .unwrap()
2818 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
2819 std::fs::write(&path, body).unwrap();
2820
2821 let err = q.get(&t.id).unwrap_err().to_string();
2822 assert!(err.contains("schema 99"), "{err}");
2823 }
2824
2825 #[test]
2826 fn revision_moves_when_the_queue_changes() {
2827 let (_dir, q) = queue();
2828 assert_eq!(q.revision(), 0, "an empty queue has no revision");
2829 let mut t = task("first");
2830 q.put(&mut t).unwrap();
2831 assert!(q.revision() > 0, "a written task moves the revision");
2832 }
2833
2834 #[test]
2835 fn revision_moves_when_deleting_an_older_task() {
2836 let (dir, q) = queue();
2837 let questions = Questions::at(dir.path().join("questions"));
2838 let mut t1 = task("older");
2839 q.put(&mut t1).unwrap();
2840 std::thread::sleep(std::time::Duration::from_millis(10));
2842 let mut t2 = task("newer");
2843 q.put(&mut t2).unwrap();
2844
2845 let rev_before = q.revision();
2846 q.remove(&t1.id, false, &questions).unwrap();
2847 let rev_after = q.revision();
2848
2849 assert_ne!(
2850 rev_before, rev_after,
2851 "deleting an older task must change the revision so other clients see the deletion"
2852 );
2853 }
2854
2855 #[test]
2856 fn removing_a_task_takes_it_out_of_the_listing() {
2857 let (dir, q) = queue();
2858 let questions = Questions::at(dir.path().join("questions"));
2859 let mut t = task("delete me");
2860 q.put(&mut t).unwrap();
2861 let removed = q.remove(t.short(), false, &questions).unwrap();
2862 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
2863 assert!(removed.quarantined.is_empty(), "nothing was blocked on it");
2864 assert!(q.list().is_empty());
2865 assert!(
2866 q.remove(&t.id, false, &questions).is_err(),
2867 "removing twice is an error"
2868 );
2869 }
2870
2871 #[test]
2872 fn removing_a_task_takes_its_stale_lock_with_it() {
2873 let (dir, q) = queue();
2874 let questions = Questions::at(dir.path().join("questions"));
2875 let mut t = task("interrupted");
2876 q.put(&mut t).unwrap();
2877
2878 let claim = q.claim(&t.id).unwrap();
2881 std::mem::forget(claim);
2882 assert!(
2883 q.claim(&t.id).is_err(),
2884 "the orphaned lock is what makes the task look claimed"
2885 );
2886
2887 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
2889 assert!(err.contains("live daemon"), "{err}");
2890 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
2891
2892 q.remove(&t.id, false, &questions).unwrap();
2894 assert!(q.list().is_empty());
2895 let mut again = task("interrupted");
2896 again.id = t.id.clone();
2897 q.put(&mut again).unwrap();
2898 assert!(
2899 q.claim(&t.id).is_ok(),
2900 "a task that comes back must be claimable, which a left-behind lock would prevent"
2901 );
2902 }
2903
2904 #[test]
2905 fn removing_a_task_quarantines_what_was_blocked_on_it() {
2906 let (dir, q) = queue();
2907 let questions = Questions::at(dir.path().join("questions"));
2908
2909 let mut dep = task("dependency");
2910 q.put(&mut dep).unwrap();
2911
2912 let mut still_valid = task("still valid");
2913 q.put(&mut still_valid).unwrap();
2914
2915 let mut blocked = task("waiting");
2916 blocked.block(
2917 vec![dep.id.clone(), still_valid.id.clone()],
2918 Some("waits on both".to_owned()),
2919 );
2920 q.put(&mut blocked).unwrap();
2921
2922 let removed = q.remove(&dep.id, false, &questions).unwrap();
2923 assert_eq!(removed.quarantined, [blocked.id.clone()]);
2924
2925 let after = q.get(&blocked.id).unwrap();
2926 assert_eq!(after.status, TaskStatus::Held);
2927 assert_eq!(after.hold_source, Some(HoldSource::Machine));
2928 assert!(after.blocked_by.is_empty());
2929 let reason = after.hold_reason.as_deref().unwrap_or_default();
2930 assert!(reason.contains(&dep.id), "{reason}");
2931 assert!(
2932 reason.contains(&still_valid.id),
2933 "the still-valid dependency must survive in the reason text: {reason}"
2934 );
2935 }
2936
2937 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
2938 let p = dir.join(name);
2939 std::fs::write(&p, body).unwrap();
2940 p
2941 }
2942
2943 #[test]
2944 fn an_attachment_copy_survives_deleting_its_source() {
2945 let (dir, q) = queue();
2946 let src = source_file(dir.path(), "shot.png", "pixels");
2947 let mut t = task("with a picture");
2948 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
2949 q.put(&mut t).unwrap();
2950 std::fs::remove_file(&src).unwrap();
2951 assert_eq!(names, ["shot.png"]);
2952 let loaded = q.get(&t.id).unwrap();
2953 let paths = q.attachment_paths(&loaded);
2954 assert_eq!(paths.len(), 1);
2955 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
2956 }
2957
2958 #[test]
2959 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
2960 let (dir, q) = queue();
2961 let mut t = task("bad names");
2962 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
2963 let src = source_file(dir.path(), name, "x");
2964 assert!(
2965 q.attach(&mut t, &[src]).is_err(),
2966 "`{name}` must be refused"
2967 );
2968 }
2969 assert!(!crate::ask::valid_asset_name("C:foo.png"));
2973 #[cfg(not(windows))]
2974 {
2975 let src = source_file(dir.path(), "C:foo.png", "x");
2976 assert!(q.attach(&mut t, &[src]).is_err());
2977 }
2978 let long = format!("{}.png", "a".repeat(70));
2979 let src = source_file(dir.path(), &long, "x");
2980 assert!(q.attach(&mut t, &[src]).is_err());
2981 assert!(t.attachments.is_empty());
2982 assert!(!q.attachments_dir(&t.id).exists());
2983 }
2984
2985 #[test]
2986 fn a_taken_attachment_name_is_numbered_not_overwritten() {
2987 let (dir, q) = queue();
2988 let a = source_file(dir.path(), "shot.png", "one");
2989 let sub = dir.path().join("other");
2990 std::fs::create_dir_all(&sub).unwrap();
2991 let b = source_file(&sub, "shot.png", "two");
2992 let mut t = task("collision");
2993 q.attach(&mut t, &[a]).unwrap();
2994 q.attach(&mut t, &[b]).unwrap();
2995 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
2996 let paths = q.attachment_paths(&t);
2997 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
2998 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
2999 }
3000
3001 #[test]
3002 fn a_renumbered_name_stays_inside_the_length_bound() {
3003 let (dir, q) = queue();
3004 let name = format!("{}.png", "a".repeat(60));
3005 assert_eq!(name.len(), 64);
3006 let a = source_file(dir.path(), &name, "one");
3007 let sub = dir.path().join("other");
3008 std::fs::create_dir_all(&sub).unwrap();
3009 let b = source_file(&sub, &name, "two");
3010 let mut t = task("long");
3011 q.attach(&mut t, &[a, b]).unwrap();
3012 assert_eq!(t.attachments.len(), 2);
3013 assert!(
3014 t.attachments
3015 .iter()
3016 .all(|n| crate::ask::valid_asset_name(n))
3017 );
3018 assert!(t.attachments[1].ends_with("-2.png"));
3019 }
3020
3021 #[test]
3022 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3023 let (dir, q) = queue();
3024 let good = source_file(dir.path(), "good.png", "ok");
3025 let mut t = task("partial");
3026 q.attach(&mut t, &[good]).unwrap();
3027 let more = source_file(dir.path(), "more.png", "ok");
3028 let missing = dir.path().join("missing.png");
3029 assert!(q.attach(&mut t, &[more, missing]).is_err());
3030 assert_eq!(t.attachments, ["good.png"]);
3031 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3032 .unwrap()
3033 .flatten()
3034 .collect();
3035 assert_eq!(on_disk.len(), 1);
3036 }
3037
3038 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3040 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3041 std::fs::create_dir_all(&tmp).unwrap();
3042 tmp
3043 }
3044
3045 #[test]
3046 fn a_failed_put_leaves_no_new_attachment_directory() {
3047 let (dir, q) = queue();
3048 let mut t = task("fresh");
3049 let tmp = block_put(&q, &t);
3050 let src = source_file(dir.path(), "shot.png", "x");
3051 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3052 assert!(t.attachments.is_empty());
3053 assert!(!q.attachments_dir(&t.id).exists());
3054 assert!(!q.path_of(&t.id).exists());
3055 std::fs::remove_dir(tmp).unwrap();
3056 }
3057
3058 #[test]
3059 fn a_failed_put_removes_only_the_copy_it_just_made() {
3060 let (dir, q) = queue();
3061 let mut t = task("edited");
3062 let first = source_file(dir.path(), "first.png", "1");
3063 q.attach_and_put(&mut t, &[first]).unwrap();
3064 block_put(&q, &t);
3065 let second = source_file(dir.path(), "second.png", "2");
3066 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3067 assert_eq!(t.attachments, ["first.png"]);
3068 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3069 .unwrap()
3070 .flatten()
3071 .map(|e| e.file_name().to_string_lossy().into_owned())
3072 .collect();
3073 assert_eq!(on_disk, ["first.png"]);
3074 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3075 }
3076
3077 #[test]
3078 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3079 let (dir, q) = queue();
3080 let questions = Questions::at(dir.path().join("questions"));
3081 let gone = task("gone");
3082 let mut other = task("other");
3083 let mut live = task("live");
3084 q.put(&mut other).unwrap();
3085 q.put(&mut live).unwrap();
3086 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3089 std::fs::create_dir_all(&orphan).unwrap();
3090 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3091 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3094 std::fs::create_dir_all(&busy).unwrap();
3095
3096 q.remove(&other.id, false, &questions).unwrap();
3097 assert!(!orphan.exists(), "an orphan is swept");
3098 assert!(busy.exists(), "a removal in progress is left alone");
3099 }
3100
3101 #[test]
3102 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3103 let (dir, q) = queue();
3104 let questions = Questions::at(dir.path().join("questions"));
3105 let mut t = task("stuck");
3106 let src = source_file(dir.path(), "shot.png", "x");
3107 q.attach_and_put(&mut t, &[src]).unwrap();
3108 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3111 std::fs::create_dir_all(&aside).unwrap();
3112 std::fs::write(aside.join("old.png"), "o").unwrap();
3113 assert!(q.remove(&t.id, false, &questions).is_err());
3114 assert!(q.path_of(&t.id).exists());
3115 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3116 }
3117
3118 #[test]
3119 fn a_failed_record_removal_puts_the_attachments_back() {
3120 let (dir, q) = queue();
3121 let mut t = task("rollback");
3122 let src = source_file(dir.path(), "shot.png", "x");
3123 q.attach_and_put(&mut t, &[src]).unwrap();
3124 let err = q
3125 .remove_record_with_attachments(&t.id, |_| {
3126 Err(std::io::Error::other("injected failure"))
3127 })
3128 .unwrap_err();
3129 assert!(format!("{err:#}").contains("injected failure"));
3130 assert!(q.path_of(&t.id).exists());
3131 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3132 assert!(
3133 !q.root
3134 .join(format!("{}.attachments.removing", t.id))
3135 .exists()
3136 );
3137 }
3138
3139 #[test]
3140 fn editing_a_task_keeps_its_attachments() {
3141 let (dir, q) = queue();
3142 let src = source_file(dir.path(), "shot.png", "x");
3143 let mut t = task("editable");
3144 q.attach(&mut t, &[src]).unwrap();
3145 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3146 q.put(&mut t).unwrap();
3147 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3148 }
3149
3150 #[test]
3151 fn removing_a_task_deletes_its_attachments() {
3152 let (dir, q) = queue();
3153 let questions = Questions::at(dir.path().join("questions"));
3154 let src = source_file(dir.path(), "shot.png", "x");
3155 let mut t = task("doomed");
3156 q.attach(&mut t, &[src]).unwrap();
3157 q.put(&mut t).unwrap();
3158 assert!(q.attachments_dir(&t.id).is_dir());
3159 q.remove(&t.id, false, &questions).unwrap();
3160 assert!(!q.attachments_dir(&t.id).exists());
3161 assert!(q.list().is_empty());
3162 }
3163
3164 #[test]
3165 fn a_task_written_before_attachments_still_reads() {
3166 let (_dir, q) = queue();
3167 let mut t = task("old");
3168 q.put(&mut t).unwrap();
3169 let path = q.path_of(&t.id);
3170 let mut v: serde_json::Value =
3171 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3172 v.as_object_mut().unwrap().remove("attachments");
3173 std::fs::write(&path, v.to_string()).unwrap();
3174 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3175 }
3176
3177 #[test]
3178 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3179 let q = Queue::at(PathBuf::from("relative-queue"));
3180 let mut t = task("rel");
3181 t.attachments.push("shot.png".to_owned());
3182 let paths = q.attachment_paths(&t);
3183 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3184 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3185 }
3186
3187 #[test]
3188 fn link_run_adds_a_run_once_and_touches_nothing_else() {
3189 let dir = tempfile::tempdir().unwrap();
3190 let queue = Queue::at(dir.path().join("queue"));
3191 let mut t = Task::new(
3192 "t".to_owned(),
3193 "do it".to_owned(),
3194 PathBuf::from("."),
3195 Source::Human,
3196 );
3197 queue.put(&mut t).unwrap();
3198 let before = queue.get(&t.id).unwrap();
3199
3200 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3201 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3202 assert_eq!(linked.status, before.status);
3203 assert_eq!(linked.attempts, before.attempts);
3204 assert_eq!(linked.interrupt, before.interrupt);
3205
3206 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3207 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3208 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3209 assert!(queue.link_run("no-such-task", "r").is_err());
3210 }
3211
3212 #[test]
3213 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
3214 let dir = tempfile::tempdir().unwrap();
3215 let queue = Queue::at(dir.path().join("queue"));
3216 let mut t = Task::new(
3217 "t".to_owned(),
3218 "do it".to_owned(),
3219 PathBuf::from("."),
3220 Source::Human,
3221 );
3222 queue.put(&mut t).unwrap();
3223 let mut snapshot = queue.get(&t.id).unwrap();
3225 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3226
3227 snapshot.start("20260930-000000-aaaa".to_owned());
3228 queue.put(&mut snapshot).unwrap();
3229
3230 let stored = queue.get(&t.id).unwrap();
3231 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3232 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3233 assert_eq!(stored.attempts, 1);
3234 }
3235
3236 #[test]
3237 fn concurrent_links_and_daemon_saves_lose_nothing() {
3238 let dir = tempfile::tempdir().unwrap();
3239 let queue = Queue::at(dir.path().join("queue"));
3240 let mut t = Task::new(
3241 "t".to_owned(),
3242 "do it".to_owned(),
3243 PathBuf::from("."),
3244 Source::Human,
3245 );
3246 queue.put(&mut t).unwrap();
3247 let id = t.id.clone();
3248
3249 let linkers: Vec<_> = (0..4)
3250 .map(|n| {
3251 let (queue, id) = (queue.clone(), id.clone());
3252 std::thread::spawn(move || {
3253 for k in 0..10 {
3254 queue
3255 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3256 .unwrap();
3257 }
3258 })
3259 })
3260 .collect();
3261 let mut mine = queue.get(&id).unwrap();
3264 for k in 0..10 {
3265 mine.start(format!("20260930-000009-d{k:03}"));
3266 queue.put(&mut mine).unwrap();
3267 }
3268 for l in linkers {
3269 l.join().unwrap();
3270 }
3271
3272 let stored = queue.get(&id).unwrap();
3273 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
3274 assert_eq!(
3275 stored.attempts, 10,
3276 "linking never rewinds the daemon's work"
3277 );
3278 }
3279}