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 = 8;
93
94#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
96#[serde(rename_all = "lowercase")]
97pub enum HoldSource {
98 Manual,
100 Machine,
102}
103
104impl HoldSource {
105 pub fn label(self) -> &'static str {
107 match self {
108 Self::Manual => "manual",
109 Self::Machine => "machine",
110 }
111 }
112}
113
114#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
117#[serde(tag = "kind", rename_all = "lowercase")]
118pub enum Source {
119 Human,
121 Agent {
124 run: String,
126 node: String,
128 },
129 Issue {
131 number: u64,
133 repo: String,
135 },
136}
137
138impl Source {
139 pub fn label(&self) -> String {
141 match self {
142 Self::Human => "human".to_owned(),
143 Self::Agent { run, node } => format!("{node}@{}", short(run)),
144 Self::Issue { number, .. } => format!("issue #{number}"),
145 }
146 }
147}
148
149#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
151#[serde(rename_all = "lowercase")]
152pub enum TaskStatus {
153 Queued,
155 Running,
157 Done,
159 Failed,
161 Held,
163 Blocked,
167}
168
169impl TaskStatus {
170 pub fn runnable(self) -> bool {
172 matches!(self, Self::Queued | Self::Failed)
173 }
174
175 pub fn as_str(self) -> &'static str {
177 match self {
178 Self::Queued => "queued",
179 Self::Running => "running",
180 Self::Done => "done",
181 Self::Failed => "failed",
182 Self::Held => "held",
183 Self::Blocked => "blocked",
184 }
185 }
186}
187
188#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
191pub struct TaskCounts {
192 pub queued: usize,
194 pub running: usize,
196 pub done: usize,
198 pub failed: usize,
200 pub held: usize,
202 pub blocked: usize,
204}
205
206impl TaskCounts {
207 pub fn of(tasks: &[Task]) -> Self {
211 let mut counts = Self::default();
212 for t in tasks {
213 match t.status {
214 TaskStatus::Queued => counts.queued += 1,
215 TaskStatus::Running => counts.running += 1,
216 TaskStatus::Done => counts.done += 1,
217 TaskStatus::Failed => counts.failed += 1,
218 TaskStatus::Held => counts.held += 1,
219 TaskStatus::Blocked => counts.blocked += 1,
220 }
221 }
222 counts
223 }
224}
225
226#[derive(Debug, Clone, Serialize, Deserialize)]
228#[serde(deny_unknown_fields)]
229pub struct Task {
230 pub schema: u32,
232 pub id: String,
234 pub title: String,
236 pub instruction: String,
238 pub repo: PathBuf,
240 pub source: Source,
242 #[serde(default)]
244 pub priority: i32,
245 #[serde(default)]
256 pub solo: bool,
257 pub status: TaskStatus,
259 #[serde(default)]
261 pub attempts: usize,
262 #[serde(default)]
264 pub runs: Vec<String>,
265 #[serde(default)]
267 pub last_error: Option<String>,
268 #[serde(default)]
281 pub hold_reason: Option<String>,
282 #[serde(default)]
285 pub hold_source: Option<HoldSource>,
286 #[serde(default)]
298 pub diagnostic: Option<String>,
299 #[serde(default)]
309 pub blocked_by: Vec<String>,
310 #[serde(default)]
313 pub block_reason: Option<String>,
314 #[serde(default)]
327 pub blocked_from: Option<TaskStatus>,
328 #[serde(default)]
339 pub answers: Vec<AnsweredQuestion>,
340 #[serde(default)]
346 pub triage_applied: Vec<String>,
347 #[serde(default)]
353 pub actions_applied: Vec<String>,
354 #[serde(default)]
359 pub resume_override: Option<OperatorResume>,
360 #[serde(default)]
372 pub review_branch: Option<String>,
373 #[serde(default)]
376 pub fresh_start: bool,
377 #[serde(default)]
388 pub interrupt: bool,
389 #[serde(default)]
412 pub urgent: bool,
413 #[serde(default)]
419 pub attachments: Vec<String>,
420 #[serde(default)]
424 pub followup: Option<FollowUp>,
425 pub created_at: Timestamp,
427 pub updated_at: Timestamp,
429}
430
431#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
433pub struct FollowUp {
434 pub run: String,
436 #[serde(default)]
438 pub origin_task: Option<String>,
439 pub pr: String,
441 pub findings: Vec<String>,
443 pub generation: u32,
446}
447
448#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
451pub struct OperatorResume {
452 pub question_id: String,
454 pub at: Timestamp,
456 #[serde(default)]
459 pub conductor_rehold: Option<String>,
460 #[serde(default)]
463 pub forced: bool,
464 #[serde(default)]
471 pub pinned_run: Option<String>,
472}
473
474#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
477pub struct AnsweredQuestion {
478 pub question: String,
480 pub answer: String,
482}
483
484impl Task {
485 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
487 let now = Timestamp::now();
488 Self {
489 schema: SCHEMA,
490 id: new_id(),
491 title,
492 instruction,
493 repo,
494 source,
495 priority: 0,
496 solo: false,
497 status: TaskStatus::Queued,
498 attempts: 0,
499 runs: Vec::new(),
500 last_error: None,
501 hold_reason: None,
502 hold_source: None,
503 diagnostic: None,
504 blocked_by: Vec::new(),
505 block_reason: None,
506 blocked_from: None,
507 answers: Vec::new(),
508 triage_applied: Vec::new(),
509 actions_applied: Vec::new(),
510 resume_override: None,
511 review_branch: None,
512 fresh_start: false,
513 interrupt: false,
514 urgent: false,
515 attachments: Vec::new(),
516 followup: None,
517 created_at: now,
518 updated_at: now,
519 }
520 }
521
522 pub fn short(&self) -> &str {
524 short(&self.id)
525 }
526
527 pub fn mark_triage_applied(&mut self, question_id: &str) {
529 if !self.triage_applied(question_id) {
530 self.triage_applied.push(question_id.to_owned());
531 }
532 }
533
534 pub fn action_applied(&self, question_id: &str) -> bool {
537 self.actions_applied.iter().any(|id| id == question_id)
538 }
539
540 pub fn mark_action_applied(&mut self, question_id: &str) {
542 if !self.action_applied(question_id) {
543 self.actions_applied.push(question_id.to_owned());
544 }
545 }
546
547 pub fn triage_applied(&self, question_id: &str) -> bool {
549 self.triage_applied.iter().any(|id| id == question_id)
550 }
551
552 pub fn start(&mut self, run: String) {
563 self.status = TaskStatus::Running;
564 self.attempts += 1;
565 self.runs.push(run);
566 self.last_error = None;
567 self.fresh_start = false;
568 self.interrupt = false;
569 self.resume_override = None;
571 }
572
573 pub fn link_run(&mut self, run: &str) -> bool {
578 if self.runs.iter().any(|r| r == run) {
579 return false;
580 }
581 self.runs.push(run.to_owned());
582 true
583 }
584
585 pub fn succeed(&mut self) {
595 self.status = TaskStatus::Done;
596 self.resume_override = None;
597 self.last_error = None;
598 self.hold_reason = None;
599 self.hold_source = None;
600 self.diagnostic = None;
601 self.blocked_by.clear();
602 self.block_reason = None;
603 self.blocked_from = None;
604 }
605
606 pub fn already_landed(&mut self, note: impl Into<String>) {
614 self.succeed();
615 self.attempts = self.attempts.saturating_sub(1);
616 self.last_error = Some(note.into());
617 }
618
619 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
649 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
650 return &[];
651 }
652 &self.runs[..self.runs.len() - 1]
653 }
654
655 pub fn successor_of(&self, run: &str) -> Option<&String> {
662 let pos = self.runs.iter().position(|r| r == run)?;
663 self.runs.get(pos + 1)
664 }
665
666 pub fn earlier_attempts(&self) -> &[String] {
673 &self.runs
674 }
675
676 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
692 let why = why.into();
693 self.diagnostic = None;
694 self.status = if self.attempts >= max_attempts {
695 self.hold_source = Some(HoldSource::Machine);
696 self.hold_reason = Some(why.clone());
697 TaskStatus::Held
698 } else {
699 TaskStatus::Failed
700 };
701 self.last_error = Some(why);
702 }
703
704 pub fn stall(&mut self, why: impl Into<String>) {
714 self.last_error = Some(why.into());
715 self.diagnostic = None;
716 self.attempts = self.attempts.saturating_sub(1);
717 self.status = TaskStatus::Failed;
718 }
719
720 pub fn operator_held(&self) -> bool {
726 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
727 }
728
729 pub fn hold_manual(&mut self, reason: Option<String>) {
739 self.status = TaskStatus::Held;
740 if reason.is_some() {
741 self.hold_reason = reason;
742 }
743 self.hold_source = Some(HoldSource::Manual);
744 self.blocked_by.clear();
745 self.block_reason = None;
746 self.blocked_from = None;
747 }
748
749 pub fn hold_machine(&mut self, reason: Option<String>) {
754 self.status = TaskStatus::Held;
755 if reason.is_some() {
756 self.hold_reason = reason;
757 }
758 self.hold_source = Some(HoldSource::Machine);
759 self.blocked_by.clear();
760 self.block_reason = None;
761 self.blocked_from = None;
762 }
763
764 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
774 if self.status != TaskStatus::Blocked {
775 self.blocked_from = Some(self.status);
776 }
777 self.status = TaskStatus::Blocked;
778 self.blocked_by = blocked_by;
779 self.block_reason = reason;
780 }
781
782 pub fn unblock(&mut self, resolved_id: &str) {
805 if self.status != TaskStatus::Blocked {
806 return;
807 }
808 self.blocked_by.retain(|id| id != resolved_id);
809 if self.blocked_by.is_empty() {
810 self.status = match self.blocked_from {
811 Some(TaskStatus::Running) => TaskStatus::Queued,
812 Some(other) => other,
813 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
814 TaskStatus::Held
815 }
816 None => TaskStatus::Queued,
817 };
818 self.block_reason = None;
819 self.blocked_from = None;
820 }
821 }
822
823 pub fn record_answer(&mut self, question: String, answer: String) {
828 self.answers.push(AnsweredQuestion { question, answer });
829 }
830
831 pub fn request_review(&mut self, branch: String) {
835 self.release();
836 self.review_branch = Some(branch);
837 }
838
839 pub fn requeue(&mut self) {
842 self.release();
843 self.review_branch = None;
845 self.fresh_start = true;
846 }
847
848 pub fn hold_for_handover(&mut self, branch: Option<String>, reason: String) {
857 if branch.is_some() {
858 self.review_branch = branch;
859 }
860 self.hold_machine(Some(reason));
861 }
862
863 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
872 if self.status == TaskStatus::Running {
873 bail!(
874 "task {} is running; its priority cannot be changed until \
875 this attempt finishes",
876 self.short()
877 );
878 }
879 self.priority = priority;
880 Ok(())
881 }
882
883 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
898 if interrupt && !self.status.runnable() {
899 bail!(
900 "task {} is {}; only a queued or failed task can be marked \
901 to interrupt",
902 self.short(),
903 self.status.as_str()
904 );
905 }
906 self.interrupt = interrupt;
907 Ok(())
908 }
909
910 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
922 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
923 bail!(
924 "task {} is {}; only a queued or held task's instruction can \
925 be edited",
926 self.short(),
927 self.status.as_str()
928 );
929 }
930 self.title = title;
931 self.instruction = instruction;
932 Ok(())
933 }
934
935 pub fn handed_off(&mut self, why: impl Into<String>) {
955 let why = why.into();
956 self.diagnostic = None;
957 self.status = TaskStatus::Held;
958 self.hold_source = Some(HoldSource::Machine);
959 self.hold_reason = Some(why.clone());
960 self.last_error = Some(why);
961 }
962
963 pub fn release(&mut self) {
967 let refused_handover = self.status == TaskStatus::Held
968 && self.hold_source == Some(HoldSource::Machine)
969 && self.review_branch.is_some();
970 self.status = TaskStatus::Queued;
971 self.attempts = 0;
972 self.last_error = None;
973 self.hold_reason = None;
976 self.hold_source = None;
977 self.diagnostic = None;
978 self.blocked_by.clear();
983 self.block_reason = None;
984 self.blocked_from = None;
985 if !refused_handover {
989 self.review_branch = None;
990 }
991 self.fresh_start = false;
992 }
993}
994
995fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
999 use std::io::ErrorKind;
1000 let (stem, ext) = match name.rfind('.') {
1001 Some(i) if i > 0 => (&name[..i], &name[i..]),
1002 _ => (name, ""),
1003 };
1004 for n in 1u32.. {
1005 let candidate = if n == 1 {
1006 name.to_owned()
1007 } else {
1008 let suffix = format!("-{n}");
1009 let room = 64usize.saturating_sub(suffix.len() + ext.len());
1010 let stem: String = stem.chars().take(room).collect();
1011 format!("{stem}{suffix}{ext}")
1012 };
1013 if !crate::ask::valid_asset_name(&candidate) {
1014 bail!("no valid attachment name is left for `{name}`");
1015 }
1016 let path = dir.join(&candidate);
1017 match std::fs::OpenOptions::new()
1018 .write(true)
1019 .create_new(true)
1020 .open(&path)
1021 {
1022 Ok(mut out) => {
1023 let copied = std::fs::File::open(src)
1024 .and_then(|mut input| std::io::copy(&mut input, &mut out));
1025 if let Err(e) = copied {
1026 drop(out);
1027 let _ = std::fs::remove_file(&path);
1028 return Err(e).with_context(|| format!("copy {}", src.display()));
1029 }
1030 return Ok((candidate, path));
1031 }
1032 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
1033 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
1034 }
1035 }
1036 unreachable!("the counter never runs out")
1037}
1038
1039const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
1042
1043struct TaskLock(PathBuf);
1045
1046impl Drop for TaskLock {
1047 fn drop(&mut self) {
1048 let _ = std::fs::remove_file(&self.0);
1049 }
1050}
1051
1052#[derive(Debug, Clone)]
1054pub struct Queue {
1055 root: PathBuf,
1056}
1057
1058impl Queue {
1059 pub fn open() -> Self {
1061 Self::at(crate::run::home().join("queue"))
1062 }
1063
1064 pub fn at(root: PathBuf) -> Self {
1067 Self { root }
1068 }
1069
1070 pub fn root(&self) -> &Path {
1072 &self.root
1073 }
1074
1075 pub fn path_of(&self, id: &str) -> PathBuf {
1077 self.root.join(format!("{id}.json"))
1078 }
1079
1080 pub fn attachments_dir(&self, id: &str) -> PathBuf {
1083 self.root.join(format!("{id}.attachments"))
1084 }
1085
1086 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1097 let mut wanted = Vec::new();
1098 for src in sources {
1099 let name = src
1100 .file_name()
1101 .and_then(|n| n.to_str())
1102 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1103 if !crate::ask::valid_asset_name(name) {
1104 bail!(
1105 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1106 with no `..`; rename the file and try again"
1107 );
1108 }
1109 if !src.is_file() {
1110 bail!("attachment `{}` is not a file", src.display());
1111 }
1112 wanted.push((src, name));
1113 }
1114 let dir = self.attachments_dir(&task.id);
1115 let existed = dir.is_dir();
1116 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1117 let mut created: Vec<PathBuf> = Vec::new();
1118 let mut names = Vec::new();
1119 let mut copy_all = || -> Result<()> {
1120 for (src, name) in &wanted {
1121 let (stored, path) = copy_new(&dir, src, name)?;
1122 created.push(path);
1123 names.push(stored);
1124 }
1125 Ok(())
1126 };
1127 if let Err(e) = copy_all() {
1128 for path in &created {
1129 let _ = std::fs::remove_file(path);
1130 }
1131 if !existed {
1132 let _ = std::fs::remove_dir(&dir);
1133 }
1134 return Err(e);
1135 }
1136 task.attachments.extend(names.iter().cloned());
1137 Ok(names)
1138 }
1139
1140 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1151 let dir = self.attachments_dir(&task.id);
1152 let existed = dir.is_dir();
1153 let before = task.attachments.len();
1154 let names = self.attach(task, sources)?;
1155 if let Err(e) = self.put(task) {
1156 for name in &names {
1157 let _ = std::fs::remove_file(dir.join(name));
1158 }
1159 if !existed {
1160 let _ = std::fs::remove_dir(&dir);
1161 }
1162 task.attachments.truncate(before);
1163 return Err(e);
1164 }
1165 Ok(names)
1166 }
1167
1168 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1172 let dir = self.attachments_dir(&task.id);
1173 task.attachments
1174 .iter()
1175 .map(|n| {
1176 let p = dir.join(n);
1177 std::path::absolute(&p).unwrap_or(p)
1178 })
1179 .collect()
1180 }
1181
1182 pub fn put(&self, task: &mut Task) -> Result<()> {
1191 let _lock = self.lock_task(&task.id)?;
1192 self.put_unlocked(task)
1193 }
1194
1195 pub fn create_new(&self, task: &mut Task) -> Result<bool> {
1200 let _lock = self.lock_task(&task.id)?;
1201 if self.path_of(&task.id).exists() {
1202 return Ok(false);
1203 }
1204 self.put_unlocked(task)?;
1205 Ok(true)
1206 }
1207
1208 fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1210 if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1211 for run in stored.runs {
1212 if !task.runs.contains(&run) {
1213 task.runs.push(run);
1214 }
1215 }
1216 }
1217 task.updated_at = Timestamp::now();
1218 std::fs::create_dir_all(&self.root)
1219 .with_context(|| format!("create {}", self.root.display()))?;
1220 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1221 let path = self.path_of(&task.id);
1222 let tmp = path.with_extension("json.tmp");
1223 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1224 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1225 if let (Some(notice), Some(home)) = (
1230 crate::notices::task_held(task),
1231 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1232 ) {
1233 crate::notices::raise_in(home, notice);
1234 }
1235 Ok(())
1236 }
1237
1238 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1243 let id = self.resolve_id(id)?;
1244 let _lock = self.lock_task(&id)?;
1247 let mut task = self.get(&id)?;
1248 if task.link_run(run) {
1249 self.put_unlocked(&mut task)?;
1250 }
1251 Ok(task)
1252 }
1253
1254 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1262 std::fs::create_dir_all(&self.root)
1263 .with_context(|| format!("create {}", self.root.display()))?;
1264 let path = self.root.join(format!("{id}.write-lock"));
1265 let started = std::time::Instant::now();
1266 loop {
1267 match std::fs::OpenOptions::new()
1268 .write(true)
1269 .create_new(true)
1270 .open(&path)
1271 {
1272 Ok(_) => return Ok(TaskLock(path)),
1273 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1274 let stale = std::fs::metadata(&path)
1275 .and_then(|m| m.modified())
1276 .ok()
1277 .and_then(|t| t.elapsed().ok())
1278 .is_some_and(|age| age > TASK_LOCK_STALE);
1279 if stale {
1280 let _ = std::fs::remove_file(&path);
1281 } else if started.elapsed() > TASK_LOCK_STALE {
1282 bail!("could not lock task {id}");
1283 } else {
1284 std::thread::sleep(std::time::Duration::from_millis(15));
1285 }
1286 }
1287 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1288 }
1289 }
1290 }
1291
1292 pub fn get(&self, id: &str) -> Result<Task> {
1294 let resolved = self.resolve_id(id)?;
1295 read_path(&self.path_of(&resolved))
1296 }
1297
1298 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1326 let resolved = self.resolve_id(id)?;
1327 if in_flight {
1328 bail!("task {resolved} is being run by a live daemon right now");
1329 }
1330 self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))?;
1331 let quarantined = self.quarantine_dependents_of(&resolved, questions);
1332 Ok(Removal {
1333 id: resolved,
1334 quarantined,
1335 })
1336 }
1337
1338 fn remove_record_with_attachments(
1342 &self,
1343 resolved: &str,
1344 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1345 ) -> Result<()> {
1346 self.sweep_removed_attachments();
1353 let attachments = self.attachments_dir(resolved);
1354 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1355 let moved = match std::fs::rename(&attachments, &aside) {
1356 Ok(()) => true,
1357 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1358 Err(e) => {
1359 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1360 }
1361 };
1362 let path = self.path_of(resolved);
1363 if let Err(e) = remove_record(&path) {
1364 if moved {
1365 let _ = std::fs::rename(&aside, &attachments);
1366 }
1367 return Err(e).with_context(|| format!("remove {}", path.display()));
1368 }
1369 if moved {
1370 if let Err(e) = std::fs::remove_dir_all(&aside) {
1371 tracing::warn!("leftover attachments {}: {e}", aside.display());
1372 }
1373 }
1374 let lock = self.lock_path(resolved);
1375 if let Err(e) = std::fs::remove_file(&lock) {
1376 if e.kind() != std::io::ErrorKind::NotFound {
1377 return Err(e).with_context(|| format!("remove {}", lock.display()));
1378 }
1379 }
1380 Ok(())
1381 }
1382
1383 fn sweep_removed_attachments(&self) {
1389 let Ok(entries) = std::fs::read_dir(&self.root) else {
1390 return;
1391 };
1392 for entry in entries.flatten() {
1393 let name = entry.file_name();
1394 let name = name.to_string_lossy();
1395 let Some(id) = name.strip_suffix(".attachments.removing") else {
1396 continue;
1397 };
1398 if !self.path_of(id).exists() {
1401 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1402 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1403 }
1404 }
1405 }
1406 }
1407
1408 fn quarantine_dependents_of(&self, dependency: &str, questions: &Questions) -> Vec<String> {
1412 let mut quarantined = Vec::new();
1413 for listed in self.list() {
1414 if listed.status != TaskStatus::Blocked
1415 || !listed.blocked_by.iter().any(|b| b == dependency)
1416 {
1417 continue;
1418 }
1419 let Ok(_claim) = self.claim(&listed.id) else {
1420 continue;
1421 };
1422 let Ok(mut task) = self.get(&listed.id) else {
1423 continue;
1424 };
1425 if task.status != TaskStatus::Blocked
1426 || !task.blocked_by.iter().any(|b| b == dependency)
1427 {
1428 continue;
1429 }
1430 let missing = missing_blockers(self, questions, &task.blocked_by);
1431 let language = crate::lang::of_repo(&task.repo);
1432 task.hold_machine(Some(missing_blocker_hold_reason_in(
1433 &task.blocked_by,
1434 &missing,
1435 &language,
1436 )));
1437 if self.put(&mut task).is_ok() {
1438 quarantined.push(task.id.clone());
1439 }
1440 }
1441 quarantined
1442 }
1443
1444 fn lock_path(&self, id: &str) -> PathBuf {
1447 self.root.join(format!("{id}.lock"))
1448 }
1449
1450 pub fn list(&self) -> Vec<Task> {
1463 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1464 .into_iter()
1465 .flatten()
1466 .flatten()
1467 .map(|e| e.path())
1468 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1469 .filter_map(|p| read_path(&p).ok())
1470 .collect();
1471 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1472 tasks
1473 }
1474
1475 pub fn superseded(&self) -> HashMap<String, String> {
1488 let mut by = HashMap::new();
1489 for task in self.list() {
1490 for earlier in &task.runs {
1491 if let Some(later) = task.successor_of(earlier) {
1492 by.insert(earlier.clone(), later.clone());
1493 }
1494 }
1495 }
1496 by
1497 }
1498
1499 pub fn superseded_by(&self, run: &str) -> Option<String> {
1507 for task in self.list() {
1508 if task.runs.iter().any(|r| r == run) {
1509 return task.successor_of(run).cloned();
1510 }
1511 }
1512 None
1513 }
1514
1515 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1526 for task in self.list() {
1527 if task.runs.iter().any(|r| r == run) {
1528 return task.runs.last().filter(|last| **last != run).cloned();
1529 }
1530 }
1531 None
1532 }
1533
1534 pub fn next_runnable(&self) -> Option<Task> {
1539 let mut runnable: Vec<Task> = self
1540 .list()
1541 .into_iter()
1542 .filter(|t| t.status.runnable())
1543 .collect();
1544 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1545 runnable.into_iter().next()
1546 }
1547
1548 pub fn claim(&self, id: &str) -> Result<Claim> {
1555 std::fs::create_dir_all(&self.root)
1556 .with_context(|| format!("create {}", self.root.display()))?;
1557 let path = self.lock_path(id);
1558 match std::fs::OpenOptions::new()
1559 .write(true)
1560 .create_new(true)
1561 .open(&path)
1562 {
1563 Ok(mut f) => {
1564 use std::io::Write as _;
1565 let _ = writeln!(f, "{}", std::process::id());
1567 Ok(Claim { path })
1568 }
1569 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1570 bail!("task {id} is already claimed ({} exists)", path.display())
1571 }
1572 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1573 }
1574 }
1575
1576 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1578 if self.path_of(prefix).is_file() {
1579 return Ok(prefix.to_owned());
1580 }
1581 let hits: Vec<String> = self
1582 .list()
1583 .into_iter()
1584 .map(|t| t.id)
1585 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1586 .collect();
1587 match hits.len() {
1588 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1589 0 => bail!("no task matches `{prefix}`"),
1590 _ => bail!(
1591 "`{prefix}` matches {} tasks: {}",
1592 hits.len(),
1593 hits.join(", ")
1594 ),
1595 }
1596 }
1597
1598 pub fn revision(&self) -> u64 {
1605 use std::hash::{Hash as _, Hasher as _};
1606
1607 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
1608 .into_iter()
1609 .flatten()
1610 .flatten()
1611 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
1612 .filter_map(|e| {
1613 let name = e.file_name().to_string_lossy().into_owned();
1614 let mtime = e
1615 .metadata()
1616 .ok()?
1617 .modified()
1618 .ok()?
1619 .duration_since(std::time::UNIX_EPOCH)
1620 .ok()?
1621 .as_millis() as u64;
1622 Some((name, mtime))
1623 })
1624 .collect();
1625
1626 if entries.is_empty() {
1627 return 0;
1628 }
1629
1630 entries.sort_unstable();
1631 let mut hasher = std::hash::DefaultHasher::new();
1632 for (name, mtime) in &entries {
1633 name.hash(&mut hasher);
1634 mtime.hash(&mut hasher);
1635 }
1636 let h = hasher.finish();
1637 if h == 0 { 1 } else { h }
1638 }
1639}
1640
1641#[derive(Debug, Clone)]
1643pub struct Removal {
1644 pub id: String,
1646 pub quarantined: Vec<String>,
1650}
1651
1652#[derive(Debug)]
1654pub struct Claim {
1655 path: PathBuf,
1656}
1657
1658impl Drop for Claim {
1659 fn drop(&mut self) {
1660 let _ = std::fs::remove_file(&self.path);
1661 }
1662}
1663
1664pub fn title_from(instruction: &str, max: usize) -> String {
1667 let line = instruction
1673 .lines()
1674 .map(str::trim)
1675 .find(|l| !l.is_empty())
1676 .unwrap_or("(empty task)")
1677 .trim_start_matches(['#', '-', '*', '>', ' '])
1678 .trim();
1679 if line.is_empty() {
1680 return "(empty task)".to_owned();
1681 }
1682 if line.chars().count() <= max {
1683 return line.to_owned();
1684 }
1685 let head: String = line.chars().take(max.saturating_sub(1)).collect();
1686 format!("{head}…")
1687}
1688
1689fn read_path(path: &Path) -> Result<Task> {
1690 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1691 let task: Task =
1692 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1693 if task.schema > SCHEMA {
1699 bail!(
1700 "task {} was written by a different magi (schema {}, this build \
1701 speaks {SCHEMA})",
1702 task.id,
1703 task.schema
1704 );
1705 }
1706 Ok(task)
1707}
1708
1709pub fn missing_blockers(
1723 queue: &Queue,
1724 questions: &Questions,
1725 blocked_by: &[String],
1726) -> Vec<String> {
1727 blocked_by
1728 .iter()
1729 .filter(|id| !queue.path_of(id).is_file() && !questions.path_of(id).is_file())
1730 .cloned()
1731 .collect()
1732}
1733
1734pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
1746 missing_blocker_hold_reason_in(blocked_by, missing, "en")
1747}
1748
1749pub fn missing_blocker_hold_reason_in(
1751 blocked_by: &[String],
1752 missing: &[String],
1753 language: &str,
1754) -> String {
1755 if crate::lang::is_japanese(language) {
1756 format!(
1757 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
1758 blocked_by.join(", "),
1759 missing.join(", "),
1760 )
1761 } else {
1762 format!(
1763 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
1764 blocked_by.join(", "),
1765 missing.join(", "),
1766 )
1767 }
1768}
1769
1770pub fn short(id: &str) -> &str {
1772 id.split('-').next_back().unwrap_or(id)
1773}
1774
1775fn new_id() -> String {
1776 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1777 let seed = crate::rng::entropy();
1778 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1779}
1780
1781#[cfg(test)]
1782mod tests {
1783 #[test]
1784 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
1785 let mut t = task("relanded");
1786 t.attempts = 1;
1787 t.status = TaskStatus::Running;
1788 t.already_landed("already in main as 0e368de");
1789 assert_eq!(t.status, TaskStatus::Done);
1790 assert_eq!(t.attempts, 0);
1791 assert!(t.hold_reason.is_none());
1792 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
1793 }
1794
1795 #[test]
1796 fn missing_blocker_reason_follows_the_language() {
1797 let b = vec!["a".to_owned()];
1798 let en = missing_blocker_hold_reason_in(&b, &b, "en");
1799 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
1800 assert!(en.starts_with("blocked on a"));
1801 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
1802 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
1803 }
1804
1805 use super::*;
1806
1807 #[test]
1808 fn triage_applied_survives_release_and_old_records_read_as_empty() {
1809 let mut t = Task::new(
1810 "t".to_owned(),
1811 "i".to_owned(),
1812 PathBuf::from("r"),
1813 Source::Human,
1814 );
1815 t.mark_triage_applied("q1");
1816 t.mark_triage_applied("q1");
1817 t.hold_machine(Some("x".to_owned()));
1818 t.release();
1819 assert_eq!(t.triage_applied, ["q1"]);
1820 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
1821
1822 let mut v = serde_json::to_value(&t).unwrap();
1823 v.as_object_mut().unwrap().remove("triage_applied");
1824 let old: Task = serde_json::from_value(v).unwrap();
1825 assert!(old.triage_applied.is_empty());
1826 }
1827
1828 #[test]
1829 fn task_counts_of_empty_is_all_zero() {
1830 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
1831 }
1832
1833 #[test]
1834 fn task_counts_of_tallies_every_status() {
1835 let mut queued = Task::new(
1836 "q".to_owned(),
1837 "i".to_owned(),
1838 PathBuf::from("."),
1839 Source::Human,
1840 );
1841 queued.status = TaskStatus::Queued;
1842 let mut running = queued.clone();
1843 running.status = TaskStatus::Running;
1844 let mut done = queued.clone();
1845 done.status = TaskStatus::Done;
1846 let mut failed = queued.clone();
1847 failed.status = TaskStatus::Failed;
1848 let mut held = queued.clone();
1849 held.status = TaskStatus::Held;
1850 let mut blocked = queued.clone();
1851 blocked.status = TaskStatus::Blocked;
1852
1853 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
1854 assert_eq!(
1855 counts,
1856 TaskCounts {
1857 queued: 1,
1858 running: 1,
1859 done: 2,
1860 failed: 1,
1861 held: 1,
1862 blocked: 1,
1863 }
1864 );
1865 }
1866
1867 fn queue() -> (tempfile::TempDir, Queue) {
1870 let dir = tempfile::tempdir().unwrap();
1871 let q = Queue::at(dir.path().join("queue"));
1872 (dir, q)
1873 }
1874
1875 #[test]
1876 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
1877 let dir = tempfile::tempdir().unwrap();
1878 let q = Queue::at(dir.path().join("queue"));
1879 let mut t = task("held");
1880 q.put(&mut t).unwrap();
1881 assert_eq!(
1882 crate::notices::Notices::at(dir.path().join("notifications"))
1883 .list()
1884 .len(),
1885 0
1886 );
1887 t.hold_machine(Some("out of attempts".to_owned()));
1888 q.put(&mut t).unwrap();
1889 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
1890 assert_eq!(listed.len(), 1);
1891 assert!(listed[0].message.contains("out of attempts"));
1892 }
1893
1894 fn task(title: &str) -> Task {
1895 Task::new(
1896 title.to_owned(),
1897 format!("do {title}"),
1898 PathBuf::from("."),
1899 Source::Human,
1900 )
1901 }
1902
1903 #[test]
1904 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
1905 let mut t = task("retried");
1906 assert!(
1907 t.earlier_attempts().is_empty(),
1908 "a first attempt takes nothing over"
1909 );
1910 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1911 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
1912 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
1914 assert_eq!(t.successor_of("bbbb"), None);
1915 assert_eq!(t.successor_of("zzzz"), None);
1916 }
1917
1918 #[test]
1919 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
1920 let (_dir, q) = queue();
1921 let mut t = task("retried");
1922 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1923 q.put(&mut t).unwrap();
1924
1925 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
1926 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
1927 assert_eq!(
1928 q.superseded_by("cccc"),
1929 None,
1930 "the latest attempt replaces nothing"
1931 );
1932 assert_eq!(
1933 q.superseded_by("never-heard-of-it"),
1934 None,
1935 "a run belonging to no task on this queue is not superseded"
1936 );
1937
1938 let mut by = HashMap::new();
1939 by.insert("aaaa".to_owned(), "bbbb".to_owned());
1940 by.insert("bbbb".to_owned(), "cccc".to_owned());
1941 assert_eq!(
1942 q.superseded(),
1943 by,
1944 "the whole-map and single-run forms must agree"
1945 );
1946 }
1947
1948 #[test]
1949 fn superseded_attempts_is_empty_until_the_task_is_done() {
1950 let mut t = task("retried");
1951 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1952 t.status = TaskStatus::Failed;
1953 assert_eq!(
1954 t.superseded_attempts(true),
1955 &[] as &[String],
1956 "a task still retrying has no attempt yet that a later one made moot"
1957 );
1958
1959 t.status = TaskStatus::Running;
1960 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1961 }
1962
1963 #[test]
1964 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
1965 let mut t = task("retried");
1966 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1967 t.status = TaskStatus::Done;
1968 assert_eq!(
1969 t.superseded_attempts(true),
1970 &["aaaa".to_owned(), "bbbb".to_owned()],
1971 "cccc is the attempt whose success made the task done, and stays out"
1972 );
1973 }
1974
1975 #[test]
1976 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
1977 let mut t = task("first try landed");
1978 t.runs = vec!["aaaa".to_owned()];
1979 t.status = TaskStatus::Done;
1980 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1981 }
1982
1983 #[test]
1984 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
1985 let mut t = task("closed by hand after a manual merge");
1992 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1993 t.status = TaskStatus::Done;
1994 assert_eq!(
1995 t.superseded_attempts(false),
1996 &[] as &[String],
1997 "nothing here is provably why the task is done, so nothing is superseded"
1998 );
1999 }
2000
2001 #[test]
2002 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2003 let (_dir, q) = queue();
2004 let mut t = task("retried twice");
2005 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2006 q.put(&mut t).unwrap();
2007
2008 assert_eq!(
2009 q.latest_attempt("aaaa"),
2010 Some("cccc".to_owned()),
2011 "an old attempt points straight at the chain's current head, not the \
2012 next attempt in the middle of it"
2013 );
2014 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2015 assert_eq!(
2016 q.latest_attempt("cccc"),
2017 None,
2018 "the latest attempt is not superseded by anything"
2019 );
2020 assert_eq!(
2021 q.latest_attempt("never-heard-of-it"),
2022 None,
2023 "a run belonging to no task on this queue is not superseded"
2024 );
2025 }
2026
2027 #[test]
2028 fn a_markdown_heading_is_the_title_not_decoration() {
2029 assert_eq!(
2034 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2035 "Rework the config loader"
2036 );
2037 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2038 assert_eq!(title_from("> quoted task", 40), "quoted task");
2039 assert_eq!(title_from(" \n\n", 40), "(empty task)");
2041 assert_eq!(title_from("###\n", 40), "(empty task)");
2042 }
2043
2044 #[test]
2045 fn a_long_title_is_elided_by_characters_not_bytes() {
2046 let long = "課題".repeat(30);
2048 let title = title_from(&long, 10);
2049 assert_eq!(title.chars().count(), 10);
2050 assert!(title.ends_with('…'));
2051 }
2052
2053 #[test]
2054 fn priority_wins_and_ties_break_oldest_first() {
2055 let (_dir, q) = queue();
2056 let mut a = task("first");
2057 let mut b = task("second");
2058 let mut c = task("urgent");
2059 a.id = "20260101-000001-aaaa".to_owned();
2061 b.id = "20260101-000002-bbbb".to_owned();
2062 c.id = "20260101-000003-cccc".to_owned();
2063 c.priority = 5;
2064 for t in [&mut a, &mut b, &mut c] {
2065 q.put(t).unwrap();
2066 }
2067
2068 assert_eq!(q.next_runnable().unwrap().id, c.id);
2070 c.hold_machine(None);
2071 q.put(&mut c).unwrap();
2072 assert_eq!(q.next_runnable().unwrap().id, a.id);
2074 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2075 }
2076
2077 #[test]
2078 fn a_blocked_task_never_starves_another_runnable_one() {
2079 let (_dir, q) = queue();
2080 let mut blocked = task("blocked");
2081 blocked.block(vec!["something".to_owned()], None);
2082 q.put(&mut blocked).unwrap();
2083
2084 let mut runnable = task("free to go");
2085 q.put(&mut runnable).unwrap();
2086
2087 let next = q.next_runnable().expect("a runnable task is still offered");
2088 assert_eq!(next.id, runnable.id);
2089 }
2090
2091 #[test]
2092 fn a_held_task_is_never_offered_to_the_loop() {
2093 let (_dir, q) = queue();
2094 let mut t = task("held");
2095 q.put(&mut t).unwrap();
2096 assert!(q.next_runnable().is_some());
2097
2098 t.hold_machine(None);
2099 q.put(&mut t).unwrap();
2100 assert!(
2101 q.next_runnable().is_none(),
2102 "a held task must wait for a human"
2103 );
2104
2105 t.status = TaskStatus::Failed;
2107 q.put(&mut t).unwrap();
2108 assert!(q.next_runnable().is_some());
2109 }
2110
2111 #[test]
2112 fn attempts_are_capped_and_then_the_task_is_held() {
2113 let mut t = task("doomed");
2114
2115 t.start("run-1".to_owned());
2116 t.fail("gate red", 2);
2117 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2118
2119 t.start("run-2".to_owned());
2120 t.fail("gate red", 2);
2121 assert_eq!(
2122 t.status,
2123 TaskStatus::Held,
2124 "out of attempts: stop spending money on it"
2125 );
2126 assert_eq!(t.runs, ["run-1", "run-2"]);
2127 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2128 assert_eq!(
2129 t.hold_reason.as_deref(),
2130 Some("gate red"),
2131 "the hold must say why, not leave hold_reason null next to a \
2132 populated last_error"
2133 );
2134 }
2135
2136 #[test]
2137 fn handing_off_a_task_records_a_hold_reason_too() {
2138 let mut t = task("left a pull request");
2139 t.start("run-1".to_owned());
2140 t.handed_off("run ended with a pull request open [run run-1]");
2141 assert_eq!(t.status, TaskStatus::Held);
2142 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2143 assert_eq!(
2144 t.hold_reason.as_deref(),
2145 Some("run ended with a pull request open [run run-1]")
2146 );
2147 assert_eq!(t.hold_reason, t.last_error);
2148 }
2149
2150 #[test]
2151 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2152 let mut t = task("stalled by quota");
2153
2154 t.start("run-1".to_owned());
2155 assert_eq!(t.attempts, 1);
2156 t.stall("judge-1, judge-2 out of quota");
2157 assert_eq!(
2158 t.attempts, 0,
2159 "a closed quota window must not spend the task's retry budget"
2160 );
2161 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2162 assert_eq!(
2163 t.last_error.as_deref(),
2164 Some("judge-1, judge-2 out of quota")
2165 );
2166
2167 for _ in 0..20 {
2170 t.start("run-n".to_owned());
2171 t.stall("still out of quota");
2172 }
2173 t.start("run-real".to_owned());
2174 t.fail("gate red", 2);
2175 assert_eq!(
2176 t.status,
2177 TaskStatus::Failed,
2178 "the first attempt that was really judged is attempt one"
2179 );
2180 }
2181
2182 #[test]
2183 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2184 let mut t = task("retry me");
2185 t.start("run-1".to_owned());
2186 t.fail("gate red", 1);
2187 assert_eq!(t.status, TaskStatus::Held);
2188
2189 t.release();
2190 assert_eq!(t.status, TaskStatus::Queued);
2191 assert_eq!(t.attempts, 0);
2194 assert!(t.last_error.is_none());
2195 assert_eq!(
2196 t.runs.len(),
2197 1,
2198 "history is kept: attempts reset, evidence does not"
2199 );
2200 }
2201
2202 #[test]
2203 fn a_hold_reason_survives_and_a_release_clears_it() {
2204 let mut t = task("waiting on something else");
2205 t.hold_manual(Some(
2206 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2207 ));
2208 assert_eq!(t.status, TaskStatus::Held);
2209 assert_eq!(
2210 t.hold_reason.as_deref(),
2211 Some("waiting for 20260101-000000-aaaa to land first")
2212 );
2213
2214 t.hold_manual(None);
2216 assert_eq!(
2217 t.hold_reason.as_deref(),
2218 Some("waiting for 20260101-000000-aaaa to land first"),
2219 "a bare re-hold keeps whatever a human already wrote down"
2220 );
2221
2222 let mut plain = task("no reason given");
2224 plain.hold_manual(None);
2225 assert_eq!(plain.status, TaskStatus::Held);
2226 assert!(plain.hold_reason.is_none());
2227
2228 t.release();
2229 assert_eq!(t.status, TaskStatus::Queued);
2230 assert!(
2231 t.hold_reason.is_none(),
2232 "a stale reason must not greet the next person who holds this task"
2233 );
2234 }
2235
2236 #[test]
2237 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2238 let mut t = task("landed by hand while held");
2243 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2244 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2245
2246 t.succeed();
2247 assert_eq!(t.status, TaskStatus::Done);
2248 assert!(
2249 t.hold_reason.is_none(),
2250 "a done task cannot still be waiting on something"
2251 );
2252 }
2253
2254 #[test]
2255 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2256 let mut held = task("held straight out of blocked");
2263 held.block(
2264 vec!["20260101-000000-dead".to_owned()],
2265 Some("waiting on the migration script".to_owned()),
2266 );
2267 assert_eq!(held.status, TaskStatus::Blocked);
2268
2269 held.hold_manual(None);
2270 assert_eq!(held.status, TaskStatus::Held);
2271 assert!(
2272 held.blocked_by.is_empty(),
2273 "hold overrides the wait, same as release"
2274 );
2275 assert!(held.block_reason.is_none());
2276
2277 let mut done = task("closed straight out of blocked");
2278 done.block(
2279 vec!["20260101-000000-dead".to_owned()],
2280 Some("waiting on the migration script".to_owned()),
2281 );
2282 done.succeed();
2283 assert_eq!(done.status, TaskStatus::Done);
2284 assert!(
2285 done.blocked_by.is_empty(),
2286 "a done task cannot still be waiting on a dependency"
2287 );
2288 assert!(done.block_reason.is_none());
2289 }
2290
2291 #[test]
2292 fn a_blocked_task_is_never_offered_to_the_loop() {
2293 let mut t = task("blocked");
2294 assert!(t.status.runnable());
2295 t.block(
2296 vec!["dep-id".to_owned()],
2297 Some("waits on dep-id".to_owned()),
2298 );
2299 assert_eq!(t.status, TaskStatus::Blocked);
2300 assert!(!t.status.runnable());
2301 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2302 }
2303
2304 #[test]
2305 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2306 let mut t = task("blocked on two");
2307 t.block(
2308 vec!["a".to_owned(), "b".to_owned()],
2309 Some("waits on a and b".to_owned()),
2310 );
2311
2312 t.unblock("a");
2313 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2314 assert_eq!(t.blocked_by, ["b"]);
2315
2316 t.unblock("b");
2317 assert_eq!(t.status, TaskStatus::Queued);
2318 assert!(t.blocked_by.is_empty());
2319 assert!(t.block_reason.is_none());
2320 }
2321
2322 #[test]
2323 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2324 let mut t = task("never blocked");
2325 t.unblock("whatever");
2326 assert_eq!(t.status, TaskStatus::Queued);
2327 }
2328
2329 #[test]
2330 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2331 let mut t = task("held, then asked about");
2336 t.hold_machine(Some("out of attempts".to_owned()));
2337 assert_eq!(t.status, TaskStatus::Held);
2338
2339 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2340 assert_eq!(t.status, TaskStatus::Blocked);
2341
2342 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2343 t.unblock("q1");
2344 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2345 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2346 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2347 assert!(t.blocked_from.is_none(), "consumed once restored");
2348 }
2349
2350 #[test]
2351 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2352 let mut t = task("manually held, then asked about");
2353 t.hold_manual(Some("waiting on a dependency".to_owned()));
2354
2355 t.block(vec!["q1".to_owned()], None);
2356 t.unblock("q1");
2357
2358 assert_eq!(t.status, TaskStatus::Held);
2359 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2360 }
2361
2362 #[test]
2363 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2364 let mut t = task("held, blocked twice");
2368 t.hold_machine(None);
2369 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2370 t.block(
2371 vec!["q1".to_owned(), "q2".to_owned()],
2372 Some("second".to_owned()),
2373 );
2374
2375 t.unblock("q1");
2376 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2377 t.unblock("q2");
2378 assert_eq!(t.status, TaskStatus::Held);
2379 }
2380
2381 #[test]
2382 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2383 let mut t = task("blocked mid-run");
2387 t.start("run-1".to_owned());
2388 assert_eq!(t.status, TaskStatus::Running);
2389
2390 t.block(vec!["q1".to_owned()], None);
2391 t.unblock("q1");
2392 assert_eq!(t.status, TaskStatus::Queued);
2393 }
2394
2395 #[test]
2396 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2397 let mut t = task("legacy record, held before it was blocked");
2403 t.hold_source = Some(HoldSource::Machine);
2404 t.hold_reason = Some("legacy hold reason".to_owned());
2405 t.status = TaskStatus::Blocked;
2406 t.blocked_by = vec!["q1".to_owned()];
2407 t.blocked_from = None;
2408
2409 t.unblock("q1");
2410 assert_eq!(t.status, TaskStatus::Held);
2411 }
2412
2413 #[test]
2414 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2415 let mut t = task("legacy record, ordinary dependency block");
2416 t.status = TaskStatus::Blocked;
2417 t.blocked_by = vec!["dep".to_owned()];
2418 t.blocked_from = None;
2419
2420 t.unblock("dep");
2421 assert_eq!(t.status, TaskStatus::Queued);
2422 }
2423
2424 #[test]
2425 fn answering_a_question_is_recorded_and_survives_a_release() {
2426 let mut t = task("asked something");
2427 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2428 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2429 t.unblock("q1");
2430 assert_eq!(t.status, TaskStatus::Queued);
2431 assert_eq!(t.answers.len(), 1);
2432 assert_eq!(t.answers[0].answer, "SQLite");
2433
2434 t.release();
2438 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2439 }
2440
2441 #[test]
2442 fn a_refused_handover_keeps_the_review_branch_across_release() {
2443 let mut t = task("refused takeover");
2444 t.start("run-1".to_owned());
2445 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
2446 assert_eq!(t.status, TaskStatus::Held);
2447 assert_eq!(t.attempts, 1);
2448 t.release();
2449 assert_eq!(t.status, TaskStatus::Queued);
2450 assert_eq!(t.attempts, 0);
2451 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2452
2453 t.hold_for_handover(None, "again".to_owned());
2455 t.requeue();
2456 assert!(t.review_branch.is_none());
2457
2458 let mut m = task("manual");
2460 m.review_branch = Some("magi/x/A".to_owned());
2461 m.hold_manual(None);
2462 m.release();
2463 assert!(m.review_branch.is_none());
2464 }
2465
2466 #[test]
2467 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2468 let mut t = task("blocked run with a surviving branch");
2469 t.start("run-1".to_owned());
2470 t.fail("blocked with major findings", 5);
2471 assert_eq!(t.status, TaskStatus::Failed);
2472
2473 t.request_review("magi/eba2/A".to_owned());
2474 assert_eq!(t.status, TaskStatus::Queued);
2475 assert_eq!(t.attempts, 0);
2476 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2477
2478 t.release();
2480 assert!(t.review_branch.is_none());
2481 }
2482
2483 #[test]
2484 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2485 let mut t = task("retry");
2486 t.start("run-1".to_owned());
2487 t.requeue();
2488 assert!(t.fresh_start);
2489
2490 t.release();
2491 assert!(!t.fresh_start);
2492 }
2493
2494 #[test]
2495 fn priority_can_be_changed_while_queued_but_not_while_running() {
2496 let mut t = task("reprioritise me");
2497 t.set_priority(5).unwrap();
2498 assert_eq!(t.priority, 5);
2499
2500 t.start("run-1".to_owned());
2501 let err = t.set_priority(9).unwrap_err().to_string();
2502 assert!(err.contains("running"), "{err}");
2503 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2504 }
2505
2506 #[test]
2507 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2508 let mut t = task("interrupt me");
2509 assert!(!t.interrupt, "off unless asked, same as any other task");
2510
2511 t.set_interrupt(true).unwrap();
2512 assert!(t.interrupt);
2513
2514 t.start("run-1".to_owned());
2515 assert!(
2516 !t.interrupt,
2517 "the mark is one-shot: dispatching the task fulfils it, \
2518 whatever the run that follows ends up doing"
2519 );
2520 let err = t.set_interrupt(true).unwrap_err().to_string();
2521 assert!(err.contains("running"), "{err}");
2522 t.set_interrupt(false).unwrap();
2525 assert!(!t.interrupt);
2526 }
2527
2528 #[test]
2532 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2533 let mut t = task("interrupt me");
2534 t.set_interrupt(true).unwrap();
2535 t.start("run-1".to_owned());
2536 t.fail("mock failure", 5);
2537 assert_eq!(t.status, TaskStatus::Failed);
2538 assert!(
2539 !t.interrupt,
2540 "one attempt already spent the mark; a retry is an ordinary \
2541 requeue, not a fresh interrupt request"
2542 );
2543 }
2544
2545 #[test]
2546 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2547 let (_dir, q) = queue();
2548 let mut a = task("first filed");
2549 let mut b = task("second filed");
2550 a.id = "20260101-000001-aaaa".to_owned();
2551 b.id = "20260101-000002-bbbb".to_owned();
2552 q.put(&mut a).unwrap();
2553 q.put(&mut b).unwrap();
2554
2555 assert_eq!(
2556 q.next_runnable().unwrap().id,
2557 a.id,
2558 "with equal priority the older task goes first, so a burst of \
2559 new work cannot starve it"
2560 );
2561 assert_eq!(
2562 q.list()[0].id,
2563 b.id,
2564 "but the list an operator reads is newest first, the same as \
2565 before priority existed - a's turn to run does not make it the \
2566 newest task"
2567 );
2568
2569 let mut a = q.get(&a.id).unwrap();
2570 a.set_priority(10).unwrap();
2571 q.put(&mut a).unwrap();
2572
2573 assert_eq!(
2574 q.next_runnable().unwrap().id,
2575 a.id,
2576 "a raised priority must be reflected the moment it is saved"
2577 );
2578 assert_eq!(
2582 q.list()[0].id,
2583 a.id,
2584 "the raised task must sort first in the list an operator reads, \
2585 not only in next_runnable's own ordering"
2586 );
2587 }
2588
2589 #[test]
2590 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2591 let mut t = Task::new(
2592 "old title".to_owned(),
2593 "old instruction".to_owned(),
2594 PathBuf::from("/repo"),
2595 Source::Agent {
2596 run: "20260101-000000-beef".to_owned(),
2597 node: "implement".to_owned(),
2598 },
2599 );
2600 let id = t.id.clone();
2601 let created_at = t.created_at;
2602 t.runs.push("20260101-000000-beef".to_owned());
2603
2604 t.edit("new title".to_owned(), "new instruction".to_owned())
2605 .unwrap();
2606
2607 assert_eq!(t.title, "new title");
2608 assert_eq!(t.instruction, "new instruction");
2609 assert_eq!(t.id, id, "editing must not mint a new id");
2610 assert_eq!(t.created_at, created_at);
2611 assert_eq!(
2612 t.source,
2613 Source::Agent {
2614 run: "20260101-000000-beef".to_owned(),
2615 node: "implement".to_owned(),
2616 },
2617 "editing must not turn agent attribution into human"
2618 );
2619 assert_eq!(t.runs, ["20260101-000000-beef"]);
2620 }
2621
2622 #[test]
2623 fn editing_is_refused_once_a_task_is_running_or_finished() {
2624 let mut running = task("in flight");
2625 running.start("run-1".to_owned());
2626 let err = running
2627 .edit("x".to_owned(), "y".to_owned())
2628 .unwrap_err()
2629 .to_string();
2630 assert!(err.contains("running"), "{err}");
2631
2632 let mut done = task("finished");
2633 done.succeed();
2634 let err = done
2635 .edit("x".to_owned(), "y".to_owned())
2636 .unwrap_err()
2637 .to_string();
2638 assert!(err.contains("done"), "{err}");
2639
2640 let mut queued = task("waiting");
2642 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2643 let mut held = task("parked");
2644 held.hold_machine(None);
2645 held.edit("x".to_owned(), "y".to_owned()).unwrap();
2646 }
2647
2648 #[test]
2649 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2650 let (_dir, q) = queue();
2651 let path = q.path_of("20260101-000000-aaaa");
2652 std::fs::create_dir_all(q.root()).unwrap();
2653 std::fs::write(
2654 &path,
2655 serde_json::json!({
2656 "schema": SCHEMA,
2657 "id": "20260101-000000-aaaa",
2658 "title": "from before hold reasons existed",
2659 "instruction": "from before hold reasons existed",
2660 "repo": ".",
2661 "source": { "kind": "human" },
2662 "status": "held",
2663 "created_at": Timestamp::now().to_string(),
2664 "updated_at": Timestamp::now().to_string(),
2665 })
2666 .to_string(),
2667 )
2668 .unwrap();
2669
2670 let task = q.get("20260101-000000-aaaa").expect("must still read");
2671 assert!(task.hold_reason.is_none());
2672 assert!(task.operator_held());
2673 }
2674
2675 #[test]
2676 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
2677 let (_dir, q) = queue();
2678 let path = q.path_of("20260101-000000-bbbb");
2679 std::fs::create_dir_all(q.root()).unwrap();
2680 std::fs::write(
2681 &path,
2682 serde_json::json!({
2683 "schema": 2,
2684 "id": "20260101-000000-bbbb",
2685 "title": "old manual recovery",
2686 "instruction": "old manual recovery",
2687 "repo": ".",
2688 "source": { "kind": "human" },
2689 "status": "held",
2690 "hold_reason": "active manual recovery run20260912-224242-daf5",
2691 "created_at": Timestamp::now().to_string(),
2692 "updated_at": Timestamp::now().to_string(),
2693 })
2694 .to_string(),
2695 )
2696 .unwrap();
2697
2698 let task = q.get("20260101-000000-bbbb").expect("must still read");
2699 assert_eq!(task.hold_source, None);
2700 assert!(task.operator_held());
2701 }
2702
2703 #[test]
2704 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
2705 let (_dir, q) = queue();
2706 let path = q.path_of("20260101-000000-aaaa");
2707 std::fs::create_dir_all(q.root()).unwrap();
2708 std::fs::write(
2709 &path,
2710 serde_json::json!({
2711 "schema": SCHEMA,
2712 "id": "20260101-000000-aaaa",
2713 "title": "from before diagnostics existed",
2714 "instruction": "from before diagnostics existed",
2715 "repo": ".",
2716 "source": { "kind": "human" },
2717 "status": "held",
2718 "created_at": Timestamp::now().to_string(),
2719 "updated_at": Timestamp::now().to_string(),
2720 })
2721 .to_string(),
2722 )
2723 .unwrap();
2724
2725 let task = q.get("20260101-000000-aaaa").expect("must still read");
2726 assert!(task.diagnostic.is_none());
2727 }
2728
2729 #[test]
2730 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
2731 let (_dir, q) = queue();
2735 let path = q.path_of("20260101-000000-aaaa");
2736 std::fs::create_dir_all(q.root()).unwrap();
2737 std::fs::write(
2738 &path,
2739 serde_json::json!({
2740 "schema": 1,
2741 "id": "20260101-000000-aaaa",
2742 "title": "from before blocking existed",
2743 "instruction": "from before blocking existed",
2744 "repo": ".",
2745 "source": { "kind": "human" },
2746 "status": "queued",
2747 "created_at": Timestamp::now().to_string(),
2748 "updated_at": Timestamp::now().to_string(),
2749 })
2750 .to_string(),
2751 )
2752 .unwrap();
2753
2754 let task = q.get("20260101-000000-aaaa").expect("must still read");
2755 assert!(task.blocked_by.is_empty());
2756 assert!(task.block_reason.is_none());
2757 assert!(task.answers.is_empty());
2758 assert!(task.review_branch.is_none());
2759 }
2760
2761 #[test]
2762 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
2763 let mut held = task("diagnosed");
2768 held.start("run-1".to_owned());
2769 held.fail("gate red", 1);
2770 held.diagnostic = Some("cargo test failed: ...".to_owned());
2771 assert_eq!(held.status, TaskStatus::Held);
2772
2773 held.release();
2774 assert!(held.diagnostic.is_none());
2775
2776 held.diagnostic = Some("cargo test failed: ...".to_owned());
2777 held.succeed();
2778 assert!(held.diagnostic.is_none());
2779 }
2780
2781 #[test]
2782 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
2783 let mut t = task("retried");
2784 t.start("run-1".to_owned());
2785 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
2786 t.fail("unrelated config error", 5);
2787 assert_eq!(t.status, TaskStatus::Failed);
2788 assert!(
2789 t.diagnostic.is_none(),
2790 "fail() must not let an old diagnostic outlive the run that produced it"
2791 );
2792 }
2793
2794 #[test]
2795 fn a_claim_is_exclusive_and_releases_on_drop() {
2796 let (_dir, q) = queue();
2797 let mut t = task("contended");
2798 q.put(&mut t).unwrap();
2799
2800 let held = q.claim(&t.id).unwrap();
2801 assert!(
2802 q.claim(&t.id).is_err(),
2803 "two daemons must not drive one task into two runs"
2804 );
2805 drop(held);
2806 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
2807 }
2808
2809 #[test]
2810 fn a_round_trip_survives_disk() {
2811 let (_dir, q) = queue();
2812 let mut t = Task::new(
2813 "titled".to_owned(),
2814 "body".to_owned(),
2815 PathBuf::from("/repo"),
2816 Source::Agent {
2817 run: "20260101-000000-beef".to_owned(),
2818 node: "implement".to_owned(),
2819 },
2820 );
2821 t.priority = 3;
2822 q.put(&mut t).unwrap();
2823
2824 let back = q.get(&t.id).unwrap();
2825 assert_eq!(back.id, t.id);
2826 assert_eq!(back.priority, 3);
2827 assert_eq!(back.source.label(), "implement@beef");
2828 assert_eq!(q.get(t.short()).unwrap().id, t.id);
2830 }
2831
2832 #[test]
2833 fn an_unreadable_task_does_not_take_the_queue_down() {
2834 let (_dir, q) = queue();
2835 let mut t = task("fine");
2836 q.put(&mut t).unwrap();
2837 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
2838
2839 let listed = q.list();
2840 assert_eq!(listed.len(), 1, "the readable task still lists");
2841 assert_eq!(listed[0].id, t.id);
2842 }
2843
2844 #[test]
2845 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
2846 let (_dir, q) = queue();
2847 let path = q.path_of("20260101-000000-aaaa");
2848 std::fs::create_dir_all(q.root()).unwrap();
2849 std::fs::write(
2850 &path,
2851 serde_json::json!({
2852 "schema": SCHEMA,
2853 "id": "20260101-000000-aaaa",
2854 "title": "from before solo existed",
2855 "instruction": "from before solo existed",
2856 "repo": ".",
2857 "source": { "kind": "human" },
2858 "status": "queued",
2859 "created_at": Timestamp::now().to_string(),
2860 "updated_at": Timestamp::now().to_string(),
2861 })
2862 .to_string(),
2863 )
2864 .unwrap();
2865
2866 let task = q.get("20260101-000000-aaaa").expect("must still read");
2867 assert!(!task.solo, "a queue file with no `solo` field means false");
2868 }
2869
2870 #[test]
2871 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
2872 let (_dir, q) = queue();
2873 let path = q.path_of("20260101-000000-bbbb");
2874 std::fs::create_dir_all(q.root()).unwrap();
2875 std::fs::write(
2876 &path,
2877 serde_json::json!({
2878 "schema": SCHEMA,
2879 "id": "20260101-000000-bbbb",
2880 "title": "from before urgent existed",
2881 "instruction": "from before urgent existed",
2882 "repo": ".",
2883 "source": { "kind": "human" },
2884 "status": "queued",
2885 "created_at": Timestamp::now().to_string(),
2886 "updated_at": Timestamp::now().to_string(),
2887 })
2888 .to_string(),
2889 )
2890 .unwrap();
2891
2892 let task = q.get("20260101-000000-bbbb").expect("must still read");
2893 assert!(
2894 !task.urgent,
2895 "a queue file with no `urgent` field means false, same as `solo`"
2896 );
2897 }
2898
2899 #[test]
2900 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
2901 let (_dir, q) = queue();
2902 let mut t = task("from the future");
2903 q.put(&mut t).unwrap();
2904 let path = q.path_of(&t.id);
2905 let body = std::fs::read_to_string(&path)
2906 .unwrap()
2907 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
2908 std::fs::write(&path, body).unwrap();
2909
2910 let err = q.get(&t.id).unwrap_err().to_string();
2911 assert!(err.contains("schema 99"), "{err}");
2912 }
2913
2914 #[test]
2915 fn revision_moves_when_the_queue_changes() {
2916 let (_dir, q) = queue();
2917 assert_eq!(q.revision(), 0, "an empty queue has no revision");
2918 let mut t = task("first");
2919 q.put(&mut t).unwrap();
2920 assert!(q.revision() > 0, "a written task moves the revision");
2921 }
2922
2923 #[test]
2924 fn revision_moves_when_deleting_an_older_task() {
2925 let (dir, q) = queue();
2926 let questions = Questions::at(dir.path().join("questions"));
2927 let mut t1 = task("older");
2928 q.put(&mut t1).unwrap();
2929 std::thread::sleep(std::time::Duration::from_millis(10));
2931 let mut t2 = task("newer");
2932 q.put(&mut t2).unwrap();
2933
2934 let rev_before = q.revision();
2935 q.remove(&t1.id, false, &questions).unwrap();
2936 let rev_after = q.revision();
2937
2938 assert_ne!(
2939 rev_before, rev_after,
2940 "deleting an older task must change the revision so other clients see the deletion"
2941 );
2942 }
2943
2944 #[test]
2945 fn removing_a_task_takes_it_out_of_the_listing() {
2946 let (dir, q) = queue();
2947 let questions = Questions::at(dir.path().join("questions"));
2948 let mut t = task("delete me");
2949 q.put(&mut t).unwrap();
2950 let removed = q.remove(t.short(), false, &questions).unwrap();
2951 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
2952 assert!(removed.quarantined.is_empty(), "nothing was blocked on it");
2953 assert!(q.list().is_empty());
2954 assert!(
2955 q.remove(&t.id, false, &questions).is_err(),
2956 "removing twice is an error"
2957 );
2958 }
2959
2960 #[test]
2961 fn removing_a_task_takes_its_stale_lock_with_it() {
2962 let (dir, q) = queue();
2963 let questions = Questions::at(dir.path().join("questions"));
2964 let mut t = task("interrupted");
2965 q.put(&mut t).unwrap();
2966
2967 let claim = q.claim(&t.id).unwrap();
2970 std::mem::forget(claim);
2971 assert!(
2972 q.claim(&t.id).is_err(),
2973 "the orphaned lock is what makes the task look claimed"
2974 );
2975
2976 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
2978 assert!(err.contains("live daemon"), "{err}");
2979 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
2980
2981 q.remove(&t.id, false, &questions).unwrap();
2983 assert!(q.list().is_empty());
2984 let mut again = task("interrupted");
2985 again.id = t.id.clone();
2986 q.put(&mut again).unwrap();
2987 assert!(
2988 q.claim(&t.id).is_ok(),
2989 "a task that comes back must be claimable, which a left-behind lock would prevent"
2990 );
2991 }
2992
2993 #[test]
2994 fn removing_a_task_quarantines_what_was_blocked_on_it() {
2995 let (dir, q) = queue();
2996 let questions = Questions::at(dir.path().join("questions"));
2997
2998 let mut dep = task("dependency");
2999 q.put(&mut dep).unwrap();
3000
3001 let mut still_valid = task("still valid");
3002 q.put(&mut still_valid).unwrap();
3003
3004 let mut blocked = task("waiting");
3005 blocked.block(
3006 vec![dep.id.clone(), still_valid.id.clone()],
3007 Some("waits on both".to_owned()),
3008 );
3009 q.put(&mut blocked).unwrap();
3010
3011 let removed = q.remove(&dep.id, false, &questions).unwrap();
3012 assert_eq!(removed.quarantined, [blocked.id.clone()]);
3013
3014 let after = q.get(&blocked.id).unwrap();
3015 assert_eq!(after.status, TaskStatus::Held);
3016 assert_eq!(after.hold_source, Some(HoldSource::Machine));
3017 assert!(after.blocked_by.is_empty());
3018 let reason = after.hold_reason.as_deref().unwrap_or_default();
3019 assert!(reason.contains(&dep.id), "{reason}");
3020 assert!(
3021 reason.contains(&still_valid.id),
3022 "the still-valid dependency must survive in the reason text: {reason}"
3023 );
3024 }
3025
3026 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3027 let p = dir.join(name);
3028 std::fs::write(&p, body).unwrap();
3029 p
3030 }
3031
3032 #[test]
3033 fn an_attachment_copy_survives_deleting_its_source() {
3034 let (dir, q) = queue();
3035 let src = source_file(dir.path(), "shot.png", "pixels");
3036 let mut t = task("with a picture");
3037 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3038 q.put(&mut t).unwrap();
3039 std::fs::remove_file(&src).unwrap();
3040 assert_eq!(names, ["shot.png"]);
3041 let loaded = q.get(&t.id).unwrap();
3042 let paths = q.attachment_paths(&loaded);
3043 assert_eq!(paths.len(), 1);
3044 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3045 }
3046
3047 #[test]
3048 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3049 let (dir, q) = queue();
3050 let mut t = task("bad names");
3051 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3052 let src = source_file(dir.path(), name, "x");
3053 assert!(
3054 q.attach(&mut t, &[src]).is_err(),
3055 "`{name}` must be refused"
3056 );
3057 }
3058 assert!(!crate::ask::valid_asset_name("C:foo.png"));
3062 #[cfg(not(windows))]
3063 {
3064 let src = source_file(dir.path(), "C:foo.png", "x");
3065 assert!(q.attach(&mut t, &[src]).is_err());
3066 }
3067 let long = format!("{}.png", "a".repeat(70));
3068 let src = source_file(dir.path(), &long, "x");
3069 assert!(q.attach(&mut t, &[src]).is_err());
3070 assert!(t.attachments.is_empty());
3071 assert!(!q.attachments_dir(&t.id).exists());
3072 }
3073
3074 #[test]
3075 fn a_taken_attachment_name_is_numbered_not_overwritten() {
3076 let (dir, q) = queue();
3077 let a = source_file(dir.path(), "shot.png", "one");
3078 let sub = dir.path().join("other");
3079 std::fs::create_dir_all(&sub).unwrap();
3080 let b = source_file(&sub, "shot.png", "two");
3081 let mut t = task("collision");
3082 q.attach(&mut t, &[a]).unwrap();
3083 q.attach(&mut t, &[b]).unwrap();
3084 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3085 let paths = q.attachment_paths(&t);
3086 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3087 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3088 }
3089
3090 #[test]
3091 fn a_renumbered_name_stays_inside_the_length_bound() {
3092 let (dir, q) = queue();
3093 let name = format!("{}.png", "a".repeat(60));
3094 assert_eq!(name.len(), 64);
3095 let a = source_file(dir.path(), &name, "one");
3096 let sub = dir.path().join("other");
3097 std::fs::create_dir_all(&sub).unwrap();
3098 let b = source_file(&sub, &name, "two");
3099 let mut t = task("long");
3100 q.attach(&mut t, &[a, b]).unwrap();
3101 assert_eq!(t.attachments.len(), 2);
3102 assert!(
3103 t.attachments
3104 .iter()
3105 .all(|n| crate::ask::valid_asset_name(n))
3106 );
3107 assert!(t.attachments[1].ends_with("-2.png"));
3108 }
3109
3110 #[test]
3111 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3112 let (dir, q) = queue();
3113 let good = source_file(dir.path(), "good.png", "ok");
3114 let mut t = task("partial");
3115 q.attach(&mut t, &[good]).unwrap();
3116 let more = source_file(dir.path(), "more.png", "ok");
3117 let missing = dir.path().join("missing.png");
3118 assert!(q.attach(&mut t, &[more, missing]).is_err());
3119 assert_eq!(t.attachments, ["good.png"]);
3120 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3121 .unwrap()
3122 .flatten()
3123 .collect();
3124 assert_eq!(on_disk.len(), 1);
3125 }
3126
3127 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3129 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3130 std::fs::create_dir_all(&tmp).unwrap();
3131 tmp
3132 }
3133
3134 #[test]
3135 fn a_failed_put_leaves_no_new_attachment_directory() {
3136 let (dir, q) = queue();
3137 let mut t = task("fresh");
3138 let tmp = block_put(&q, &t);
3139 let src = source_file(dir.path(), "shot.png", "x");
3140 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3141 assert!(t.attachments.is_empty());
3142 assert!(!q.attachments_dir(&t.id).exists());
3143 assert!(!q.path_of(&t.id).exists());
3144 std::fs::remove_dir(tmp).unwrap();
3145 }
3146
3147 #[test]
3148 fn a_failed_put_removes_only_the_copy_it_just_made() {
3149 let (dir, q) = queue();
3150 let mut t = task("edited");
3151 let first = source_file(dir.path(), "first.png", "1");
3152 q.attach_and_put(&mut t, &[first]).unwrap();
3153 block_put(&q, &t);
3154 let second = source_file(dir.path(), "second.png", "2");
3155 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3156 assert_eq!(t.attachments, ["first.png"]);
3157 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3158 .unwrap()
3159 .flatten()
3160 .map(|e| e.file_name().to_string_lossy().into_owned())
3161 .collect();
3162 assert_eq!(on_disk, ["first.png"]);
3163 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3164 }
3165
3166 #[test]
3167 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3168 let (dir, q) = queue();
3169 let questions = Questions::at(dir.path().join("questions"));
3170 let gone = task("gone");
3171 let mut other = task("other");
3172 let mut live = task("live");
3173 q.put(&mut other).unwrap();
3174 q.put(&mut live).unwrap();
3175 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3178 std::fs::create_dir_all(&orphan).unwrap();
3179 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3180 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3183 std::fs::create_dir_all(&busy).unwrap();
3184
3185 q.remove(&other.id, false, &questions).unwrap();
3186 assert!(!orphan.exists(), "an orphan is swept");
3187 assert!(busy.exists(), "a removal in progress is left alone");
3188 }
3189
3190 #[test]
3191 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3192 let (dir, q) = queue();
3193 let questions = Questions::at(dir.path().join("questions"));
3194 let mut t = task("stuck");
3195 let src = source_file(dir.path(), "shot.png", "x");
3196 q.attach_and_put(&mut t, &[src]).unwrap();
3197 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3200 std::fs::create_dir_all(&aside).unwrap();
3201 std::fs::write(aside.join("old.png"), "o").unwrap();
3202 assert!(q.remove(&t.id, false, &questions).is_err());
3203 assert!(q.path_of(&t.id).exists());
3204 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3205 }
3206
3207 #[test]
3208 fn a_failed_record_removal_puts_the_attachments_back() {
3209 let (dir, q) = queue();
3210 let mut t = task("rollback");
3211 let src = source_file(dir.path(), "shot.png", "x");
3212 q.attach_and_put(&mut t, &[src]).unwrap();
3213 let err = q
3214 .remove_record_with_attachments(&t.id, |_| {
3215 Err(std::io::Error::other("injected failure"))
3216 })
3217 .unwrap_err();
3218 assert!(format!("{err:#}").contains("injected failure"));
3219 assert!(q.path_of(&t.id).exists());
3220 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3221 assert!(
3222 !q.root
3223 .join(format!("{}.attachments.removing", t.id))
3224 .exists()
3225 );
3226 }
3227
3228 #[test]
3229 fn editing_a_task_keeps_its_attachments() {
3230 let (dir, q) = queue();
3231 let src = source_file(dir.path(), "shot.png", "x");
3232 let mut t = task("editable");
3233 q.attach(&mut t, &[src]).unwrap();
3234 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3235 q.put(&mut t).unwrap();
3236 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3237 }
3238
3239 #[test]
3240 fn removing_a_task_deletes_its_attachments() {
3241 let (dir, q) = queue();
3242 let questions = Questions::at(dir.path().join("questions"));
3243 let src = source_file(dir.path(), "shot.png", "x");
3244 let mut t = task("doomed");
3245 q.attach(&mut t, &[src]).unwrap();
3246 q.put(&mut t).unwrap();
3247 assert!(q.attachments_dir(&t.id).is_dir());
3248 q.remove(&t.id, false, &questions).unwrap();
3249 assert!(!q.attachments_dir(&t.id).exists());
3250 assert!(q.list().is_empty());
3251 }
3252
3253 #[test]
3254 fn a_task_written_before_attachments_still_reads() {
3255 let (_dir, q) = queue();
3256 let mut t = task("old");
3257 q.put(&mut t).unwrap();
3258 let path = q.path_of(&t.id);
3259 let mut v: serde_json::Value =
3260 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3261 v.as_object_mut().unwrap().remove("attachments");
3262 std::fs::write(&path, v.to_string()).unwrap();
3263 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3264 }
3265
3266 #[test]
3267 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3268 let q = Queue::at(PathBuf::from("relative-queue"));
3269 let mut t = task("rel");
3270 t.attachments.push("shot.png".to_owned());
3271 let paths = q.attachment_paths(&t);
3272 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3273 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3274 }
3275
3276 #[test]
3277 fn link_run_adds_a_run_once_and_touches_nothing_else() {
3278 let dir = tempfile::tempdir().unwrap();
3279 let queue = Queue::at(dir.path().join("queue"));
3280 let mut t = Task::new(
3281 "t".to_owned(),
3282 "do it".to_owned(),
3283 PathBuf::from("."),
3284 Source::Human,
3285 );
3286 queue.put(&mut t).unwrap();
3287 let before = queue.get(&t.id).unwrap();
3288
3289 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3290 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3291 assert_eq!(linked.status, before.status);
3292 assert_eq!(linked.attempts, before.attempts);
3293 assert_eq!(linked.interrupt, before.interrupt);
3294
3295 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3296 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3297 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3298 assert!(queue.link_run("no-such-task", "r").is_err());
3299 }
3300
3301 #[test]
3302 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
3303 let dir = tempfile::tempdir().unwrap();
3304 let queue = Queue::at(dir.path().join("queue"));
3305 let mut t = Task::new(
3306 "t".to_owned(),
3307 "do it".to_owned(),
3308 PathBuf::from("."),
3309 Source::Human,
3310 );
3311 queue.put(&mut t).unwrap();
3312 let mut snapshot = queue.get(&t.id).unwrap();
3314 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3315
3316 snapshot.start("20260930-000000-aaaa".to_owned());
3317 queue.put(&mut snapshot).unwrap();
3318
3319 let stored = queue.get(&t.id).unwrap();
3320 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3321 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3322 assert_eq!(stored.attempts, 1);
3323 }
3324
3325 #[test]
3326 fn concurrent_links_and_daemon_saves_lose_nothing() {
3327 let dir = tempfile::tempdir().unwrap();
3328 let queue = Queue::at(dir.path().join("queue"));
3329 let mut t = Task::new(
3330 "t".to_owned(),
3331 "do it".to_owned(),
3332 PathBuf::from("."),
3333 Source::Human,
3334 );
3335 queue.put(&mut t).unwrap();
3336 let id = t.id.clone();
3337
3338 let linkers: Vec<_> = (0..4)
3339 .map(|n| {
3340 let (queue, id) = (queue.clone(), id.clone());
3341 std::thread::spawn(move || {
3342 for k in 0..10 {
3343 queue
3344 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3345 .unwrap();
3346 }
3347 })
3348 })
3349 .collect();
3350 let mut mine = queue.get(&id).unwrap();
3353 for k in 0..10 {
3354 mine.start(format!("20260930-000009-d{k:03}"));
3355 queue.put(&mut mine).unwrap();
3356 }
3357 for l in linkers {
3358 l.join().unwrap();
3359 }
3360
3361 let stored = queue.get(&id).unwrap();
3362 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
3363 assert_eq!(
3364 stored.attempts, 10,
3365 "linking never rewinds the daemon's work"
3366 );
3367 }
3368}