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 = 11;
106
107#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
109#[serde(rename_all = "lowercase")]
110pub enum HoldSource {
111 Manual,
113 Machine,
115}
116
117impl HoldSource {
118 pub fn label(self) -> &'static str {
120 match self {
121 Self::Manual => "manual",
122 Self::Machine => "machine",
123 }
124 }
125}
126
127#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
130#[serde(tag = "kind", rename_all = "lowercase")]
131pub enum Source {
132 Human,
134 Agent {
137 run: String,
139 node: String,
141 },
142 Issue {
144 number: u64,
146 repo: String,
148 },
149}
150
151pub const CHAT_NODE: &str = "chat";
154
155impl Source {
156 pub fn label(&self) -> String {
158 match self {
159 Self::Human => "human".to_owned(),
160 Self::Agent { run, node } => format!("{node}@{}", short(run)),
161 Self::Issue { number, .. } => format!("issue #{number}"),
162 }
163 }
164}
165
166#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
168#[serde(rename_all = "lowercase")]
169pub enum TaskStatus {
170 Queued,
172 Running,
174 Done,
176 Failed,
178 Held,
180 Blocked,
184 Parked,
189}
190
191impl TaskStatus {
192 pub fn runnable(self) -> bool {
194 matches!(self, Self::Queued | Self::Failed | Self::Parked)
195 }
196
197 pub fn as_str(self) -> &'static str {
199 match self {
200 Self::Queued => "queued",
201 Self::Running => "running",
202 Self::Done => "done",
203 Self::Failed => "failed",
204 Self::Held => "held",
205 Self::Blocked => "blocked",
206 Self::Parked => "parked",
207 }
208 }
209}
210
211#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
214pub struct TaskCounts {
215 pub queued: usize,
217 pub running: usize,
219 pub done: usize,
221 pub failed: usize,
223 pub held: usize,
225 pub blocked: usize,
227 pub parked: usize,
229}
230
231impl TaskCounts {
232 pub fn of(tasks: &[Task]) -> Self {
236 let mut counts = Self::default();
237 for t in tasks {
238 match t.status {
239 TaskStatus::Queued => counts.queued += 1,
240 TaskStatus::Running => counts.running += 1,
241 TaskStatus::Done => counts.done += 1,
242 TaskStatus::Failed => counts.failed += 1,
243 TaskStatus::Held => counts.held += 1,
244 TaskStatus::Blocked => counts.blocked += 1,
245 TaskStatus::Parked => counts.parked += 1,
246 }
247 }
248 counts
249 }
250}
251
252#[derive(Debug, Clone, Serialize, Deserialize)]
254#[serde(deny_unknown_fields)]
255pub struct Task {
256 pub schema: u32,
258 pub id: String,
260 pub title: String,
262 pub instruction: String,
264 pub repo: PathBuf,
266 pub source: Source,
268 #[serde(default)]
270 pub priority: i32,
271 #[serde(default)]
282 pub solo: bool,
283 pub status: TaskStatus,
285 #[serde(default)]
287 pub attempts: usize,
288 #[serde(default)]
290 pub runs: Vec<String>,
291 #[serde(default)]
293 pub last_error: Option<String>,
294 #[serde(default)]
307 pub hold_reason: Option<String>,
308 #[serde(default)]
311 pub hold_source: Option<HoldSource>,
312 #[serde(default)]
324 pub diagnostic: Option<String>,
325 #[serde(default)]
335 pub blocked_by: Vec<String>,
336 #[serde(default)]
339 pub block_reason: Option<String>,
340 #[serde(default)]
353 pub blocked_from: Option<TaskStatus>,
354 #[serde(default)]
365 pub answers: Vec<AnsweredQuestion>,
366 #[serde(default)]
372 pub triage_applied: Vec<String>,
373 #[serde(default)]
379 pub actions_applied: Vec<String>,
380 #[serde(default)]
385 pub resume_override: Option<OperatorResume>,
386 #[serde(default)]
398 pub review_branch: Option<String>,
399 #[serde(default)]
402 pub fresh_start: bool,
403 #[serde(default)]
414 pub interrupt: bool,
415 #[serde(default)]
438 pub urgent: bool,
439 #[serde(default)]
445 pub attachments: Vec<String>,
446 #[serde(default)]
450 pub followup: Option<FollowUp>,
451 #[serde(default)]
457 pub overrides: Option<RunOverrides>,
458 #[serde(default)]
465 pub review_of: Option<String>,
466 #[serde(default)]
471 pub held_at: Option<Timestamp>,
472 #[serde(default)]
476 pub park_reason: Option<String>,
477 pub created_at: Timestamp,
479 pub updated_at: Timestamp,
481}
482
483#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
485pub struct RunOverrides {
486 #[serde(default)]
488 pub merge: Option<String>,
489 #[serde(default)]
491 pub candidates: Option<usize>,
492 #[serde(default)]
494 pub judges: Option<usize>,
495 #[serde(default)]
497 pub reviewers: Option<usize>,
498 #[serde(default)]
500 pub review_rounds: Option<usize>,
501 #[serde(default)]
503 pub seed: Option<u64>,
504 #[serde(default)]
506 pub config: Option<PathBuf>,
507}
508
509impl RunOverrides {
510 pub fn merge_over(&mut self, other: &RunOverrides) {
512 macro_rules! take {
513 ($($f:ident),*) => {$(
514 if other.$f.is_some() {
515 self.$f.clone_from(&other.$f);
516 }
517 )*};
518 }
519 take!(
520 merge,
521 candidates,
522 judges,
523 reviewers,
524 review_rounds,
525 seed,
526 config
527 );
528 }
529
530 pub fn apply(&self, config: &mut crate::config::Config) {
532 if let Some(n) = self.candidates {
533 config.graph.implementers = n;
534 }
535 if let Some(n) = self.judges {
536 config.graph.judges = n;
537 }
538 if let Some(n) = self.reviewers {
539 config.graph.reviewers = n;
540 }
541 if let Some(n) = self.review_rounds {
542 config.graph.review_rounds = n;
543 }
544 if let Some(m) = &self.merge
545 && let Ok(mode) = crate::daemon::merge_mode(m)
546 {
547 config.merge.mode = mode;
548 }
549 if let Some(s) = self.seed {
550 config.blind.seed = Some(s);
551 }
552 }
553}
554
555#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
557pub struct FollowUp {
558 pub run: String,
560 #[serde(default)]
562 pub origin_task: Option<String>,
563 pub pr: String,
565 pub findings: Vec<String>,
567 pub generation: u32,
570}
571
572#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
575pub struct OperatorResume {
576 pub question_id: String,
578 pub at: Timestamp,
580 #[serde(default)]
583 pub conductor_rehold: Option<String>,
584 #[serde(default)]
587 pub forced: bool,
588 #[serde(default)]
595 pub pinned_run: Option<String>,
596}
597
598#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
601pub struct AnsweredQuestion {
602 pub question: String,
604 pub answer: String,
606}
607
608impl Task {
609 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
611 let now = Timestamp::now();
612 Self {
613 schema: SCHEMA,
614 id: new_id(),
615 title,
616 instruction,
617 repo,
618 source,
619 priority: 0,
620 solo: false,
621 status: TaskStatus::Queued,
622 attempts: 0,
623 runs: Vec::new(),
624 last_error: None,
625 hold_reason: None,
626 hold_source: None,
627 diagnostic: None,
628 blocked_by: Vec::new(),
629 block_reason: None,
630 blocked_from: None,
631 answers: Vec::new(),
632 triage_applied: Vec::new(),
633 actions_applied: Vec::new(),
634 resume_override: None,
635 review_branch: None,
636 fresh_start: false,
637 interrupt: false,
638 urgent: false,
639 attachments: Vec::new(),
640 followup: None,
641 overrides: None,
642 review_of: None,
643 held_at: None,
644 park_reason: None,
645 created_at: now,
646 updated_at: now,
647 }
648 }
649
650 pub fn short(&self) -> &str {
652 short(&self.id)
653 }
654
655 pub fn mark_triage_applied(&mut self, question_id: &str) {
657 if !self.triage_applied(question_id) {
658 self.triage_applied.push(question_id.to_owned());
659 }
660 }
661
662 pub fn action_applied(&self, question_id: &str) -> bool {
665 self.actions_applied.iter().any(|id| id == question_id)
666 }
667
668 pub fn mark_action_applied(&mut self, question_id: &str) {
670 if !self.action_applied(question_id) {
671 self.actions_applied.push(question_id.to_owned());
672 }
673 }
674
675 pub fn triage_applied(&self, question_id: &str) -> bool {
677 self.triage_applied.iter().any(|id| id == question_id)
678 }
679
680 pub fn start(&mut self, run: String) {
691 self.status = TaskStatus::Running;
692 self.attempts += 1;
693 self.runs.push(run);
694 self.last_error = None;
695 self.park_reason = None;
696 self.fresh_start = false;
697 self.interrupt = false;
698 self.resume_override = None;
700 }
701
702 pub fn link_run(&mut self, run: &str) -> bool {
707 if self.runs.iter().any(|r| r == run) {
708 return false;
709 }
710 self.runs.push(run.to_owned());
711 true
712 }
713
714 pub fn succeed(&mut self) {
724 self.status = TaskStatus::Done;
725 self.resume_override = None;
726 self.last_error = None;
727 self.park_reason = None;
728 self.hold_reason = None;
729 self.hold_source = None;
730 self.diagnostic = None;
731 self.blocked_by.clear();
732 self.block_reason = None;
733 self.blocked_from = None;
734 }
735
736 pub fn already_landed(&mut self, note: impl Into<String>) {
744 self.succeed();
745 self.attempts = self.attempts.saturating_sub(1);
746 self.last_error = Some(note.into());
747 }
748
749 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
779 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
780 return &[];
781 }
782 &self.runs[..self.runs.len() - 1]
783 }
784
785 pub fn successor_of(&self, run: &str) -> Option<&String> {
792 let pos = self.runs.iter().position(|r| r == run)?;
793 self.runs.get(pos + 1)
794 }
795
796 pub fn earlier_attempts(&self) -> &[String] {
803 &self.runs
804 }
805
806 fn note_held(&mut self) {
809 if self.status != TaskStatus::Held || self.held_at.is_none() {
810 self.held_at = Some(Timestamp::now());
811 }
812 }
813
814 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
830 let why = why.into();
831 self.diagnostic = None;
832 self.park_reason = None;
833 self.status = if self.attempts >= max_attempts {
834 self.note_held();
835 self.hold_source = Some(HoldSource::Machine);
836 self.hold_reason = Some(why.clone());
837 TaskStatus::Held
838 } else {
839 TaskStatus::Failed
840 };
841 self.last_error = Some(why);
842 }
843
844 pub fn stall(&mut self, why: impl Into<String>) {
854 self.last_error = Some(why.into());
855 self.diagnostic = None;
856 self.attempts = self.attempts.saturating_sub(1);
857 self.status = TaskStatus::Failed;
858 }
859
860 pub fn park(&mut self, why: impl Into<String>) {
866 self.diagnostic = None;
867 self.attempts = self.attempts.saturating_sub(1);
868 self.status = TaskStatus::Parked;
869 self.park_reason = Some(why.into());
870 self.schema = self.schema.max(SCHEMA);
873 }
874
875 pub fn operator_held(&self) -> bool {
881 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
882 }
883
884 pub fn hold_manual(&mut self, reason: Option<String>) {
894 self.park_reason = None;
895 self.note_held();
896 self.status = TaskStatus::Held;
897 if reason.is_some() {
898 self.hold_reason = reason;
899 }
900 self.hold_source = Some(HoldSource::Manual);
901 self.blocked_by.clear();
902 self.block_reason = None;
903 self.blocked_from = None;
904 }
905
906 pub fn hold_machine(&mut self, reason: Option<String>) {
911 self.park_reason = None;
912 self.note_held();
913 self.status = TaskStatus::Held;
914 if reason.is_some() {
915 self.hold_reason = reason;
916 }
917 self.hold_source = Some(HoldSource::Machine);
918 self.blocked_by.clear();
919 self.block_reason = None;
920 self.blocked_from = None;
921 }
922
923 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
933 if self.status != TaskStatus::Blocked {
934 self.blocked_from = Some(self.status);
935 }
936 self.status = TaskStatus::Blocked;
937 self.blocked_by = blocked_by;
938 self.block_reason = reason;
939 }
940
941 pub fn unblock(&mut self, resolved_id: &str) {
964 if self.status != TaskStatus::Blocked {
965 return;
966 }
967 self.blocked_by.retain(|id| id != resolved_id);
968 self.restore_if_unblocked();
969 }
970
971 pub fn dependency_deleted(&mut self, deleted_id: &str) -> bool {
982 if self.status != TaskStatus::Blocked || !self.blocked_by.iter().any(|b| b == deleted_id) {
983 return false;
984 }
985 self.blocked_by.retain(|id| id != deleted_id);
986 self.restore_if_unblocked();
987 true
988 }
989
990 fn restore_if_unblocked(&mut self) {
994 if self.blocked_by.is_empty() {
995 self.status = match self.blocked_from {
996 Some(TaskStatus::Running) => TaskStatus::Queued,
997 Some(other) => other,
998 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
999 TaskStatus::Held
1000 }
1001 None => TaskStatus::Queued,
1002 };
1003 self.block_reason = None;
1004 self.blocked_from = None;
1005 }
1006 }
1007
1008 pub fn record_answer(&mut self, question: String, answer: String) {
1013 self.answers.push(AnsweredQuestion { question, answer });
1014 }
1015
1016 pub fn request_review(&mut self, branch: String) {
1020 self.release();
1021 self.review_branch = Some(branch);
1022 }
1023
1024 pub fn requeue(&mut self) {
1027 self.release();
1028 self.review_branch = None;
1030 self.fresh_start = true;
1031 }
1032
1033 pub fn hold_for_handover(&mut self, branch: Option<String>, reason: String) {
1042 if branch.is_some() {
1043 self.review_branch = branch;
1044 }
1045 self.hold_machine(Some(reason));
1046 }
1047
1048 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
1057 if self.status == TaskStatus::Running {
1058 bail!(
1059 "task {} is running; its priority cannot be changed until \
1060 this attempt finishes",
1061 self.short()
1062 );
1063 }
1064 self.priority = priority;
1065 Ok(())
1066 }
1067
1068 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
1083 if interrupt && !self.status.runnable() {
1084 bail!(
1085 "task {} is {}; only a queued or failed task can be marked \
1086 to interrupt",
1087 self.short(),
1088 self.status.as_str()
1089 );
1090 }
1091 self.interrupt = interrupt;
1092 Ok(())
1093 }
1094
1095 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
1107 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
1108 bail!(
1109 "task {} is {}; only a queued or held task's instruction can \
1110 be edited",
1111 self.short(),
1112 self.status.as_str()
1113 );
1114 }
1115 self.title = title;
1116 self.instruction = instruction;
1117 Ok(())
1118 }
1119
1120 pub fn handed_off(&mut self, why: impl Into<String>) {
1140 let why = why.into();
1141 self.diagnostic = None;
1142 self.note_held();
1143 self.status = TaskStatus::Held;
1144 self.hold_source = Some(HoldSource::Machine);
1145 self.hold_reason = Some(why.clone());
1146 self.last_error = Some(why);
1147 }
1148
1149 pub fn release(&mut self) {
1153 let refused_handover = self.status == TaskStatus::Held
1154 && self.hold_source == Some(HoldSource::Machine)
1155 && self.review_branch.is_some();
1156 self.status = TaskStatus::Queued;
1157 self.held_at = None;
1158 self.attempts = 0;
1159 self.last_error = None;
1160 self.park_reason = None;
1161 self.hold_reason = None;
1164 self.hold_source = None;
1165 self.diagnostic = None;
1166 self.blocked_by.clear();
1171 self.block_reason = None;
1172 self.blocked_from = None;
1173 if !refused_handover {
1177 self.review_branch = None;
1178 }
1179 self.fresh_start = false;
1180 }
1181}
1182
1183fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
1187 use std::io::ErrorKind;
1188 let (stem, ext) = match name.rfind('.') {
1189 Some(i) if i > 0 => (&name[..i], &name[i..]),
1190 _ => (name, ""),
1191 };
1192 for n in 1u32.. {
1193 let candidate = if n == 1 {
1194 name.to_owned()
1195 } else {
1196 let suffix = format!("-{n}");
1197 let room = 64usize.saturating_sub(suffix.len() + ext.len());
1198 let stem: String = stem.chars().take(room).collect();
1199 format!("{stem}{suffix}{ext}")
1200 };
1201 if !crate::ask::valid_asset_name(&candidate) {
1202 bail!("no valid attachment name is left for `{name}`");
1203 }
1204 let path = dir.join(&candidate);
1205 match std::fs::OpenOptions::new()
1206 .write(true)
1207 .create_new(true)
1208 .open(&path)
1209 {
1210 Ok(mut out) => {
1211 let copied = std::fs::File::open(src)
1212 .and_then(|mut input| std::io::copy(&mut input, &mut out));
1213 if let Err(e) = copied {
1214 drop(out);
1215 let _ = std::fs::remove_file(&path);
1216 return Err(e).with_context(|| format!("copy {}", src.display()));
1217 }
1218 return Ok((candidate, path));
1219 }
1220 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
1221 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
1222 }
1223 }
1224 unreachable!("the counter never runs out")
1225}
1226
1227const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
1230
1231struct TaskLock {
1234 path: PathBuf,
1235 token: String,
1236}
1237
1238fn owner_token() -> String {
1241 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1242 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1243 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy() ^ n.rotate_left(32));
1244 format!("{}-{:016x}", std::process::id(), r.next_u64())
1245}
1246
1247fn break_marker(path: &Path, token: &str) -> PathBuf {
1250 let mut h: u64 = 0xcbf29ce484222325;
1251 for b in token.bytes() {
1252 h ^= u64::from(b);
1253 h = h.wrapping_mul(0x100000001b3);
1254 }
1255 let mut name = path.as_os_str().to_owned();
1256 name.push(format!(".break-{h:016x}"));
1257 PathBuf::from(name)
1258}
1259
1260fn older_than_stale(path: &Path) -> bool {
1261 std::fs::metadata(path)
1262 .and_then(|m| m.modified())
1263 .ok()
1264 .and_then(|t| t.elapsed().ok())
1265 .is_some_and(|age| age > TASK_LOCK_STALE)
1266}
1267
1268const MAX_MARKER_DEPTH: u8 = 3;
1272
1273fn take_marker(path: &Path, token: &str, depth: u8) -> Option<(PathBuf, String)> {
1279 use std::io::Write;
1280 let marker = break_marker(path, token);
1281 for _ in 0..2 {
1282 match std::fs::OpenOptions::new()
1283 .write(true)
1284 .create_new(true)
1285 .open(&marker)
1286 {
1287 Ok(mut f) => {
1288 let mine = owner_token();
1289 if f.write_all(mine.as_bytes()).is_err() {
1290 drop(f);
1291 let _ = std::fs::remove_file(&marker);
1292 return None;
1293 }
1294 return Some((marker, mine));
1295 }
1296 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1297 let judged = std::fs::read_to_string(&marker).ok();
1298 match judged.filter(|_| older_than_stale(&marker)) {
1299 Some(judged) if depth < MAX_MARKER_DEPTH => {
1300 remove_lock_if(&marker, &judged, true, depth + 1)?;
1301 }
1302 _ => return None,
1303 }
1304 }
1305 Err(_) => return None,
1306 }
1307 }
1308 None
1309}
1310
1311fn remove_lock_if(path: &Path, token: &str, require_stale: bool, depth: u8) -> Option<bool> {
1316 let (marker, mine) = take_marker(path, token, depth)?;
1317 let still = std::fs::read_to_string(path).is_ok_and(|c| c == token)
1318 && (!require_stale || older_than_stale(path));
1319 if still {
1320 let _ = std::fs::remove_file(path);
1321 }
1322 if depth >= MAX_MARKER_DEPTH {
1324 if std::fs::read_to_string(&marker).is_ok_and(|c| c == mine) {
1325 let _ = std::fs::remove_file(&marker);
1326 }
1327 } else {
1328 let _ = remove_lock_if(&marker, &mine, false, depth + 1);
1329 }
1330 Some(still)
1331}
1332
1333fn break_stale(path: &Path, judged: &str) -> bool {
1336 remove_lock_if(path, judged, true, 0).unwrap_or(false)
1337}
1338
1339impl Drop for TaskLock {
1340 fn drop(&mut self) {
1341 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
1342 loop {
1343 match remove_lock_if(&self.path, &self.token, false, 0) {
1344 Some(_) => return,
1345 None if std::time::Instant::now() > deadline => return,
1346 None => std::thread::sleep(std::time::Duration::from_millis(5)),
1347 }
1348 }
1349 }
1350}
1351
1352#[derive(Debug, Clone)]
1354pub struct Queue {
1355 root: PathBuf,
1356}
1357
1358impl Queue {
1359 pub fn open() -> Self {
1361 Self::at(crate::run::home().join("queue"))
1362 }
1363
1364 pub fn at(root: PathBuf) -> Self {
1367 Self { root }
1368 }
1369
1370 pub fn root(&self) -> &Path {
1372 &self.root
1373 }
1374
1375 pub fn path_of(&self, id: &str) -> PathBuf {
1377 self.root.join(format!("{id}.json"))
1378 }
1379
1380 pub fn attachments_dir(&self, id: &str) -> PathBuf {
1383 self.root.join(format!("{id}.attachments"))
1384 }
1385
1386 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1397 let mut wanted = Vec::new();
1398 for src in sources {
1399 let name = src
1400 .file_name()
1401 .and_then(|n| n.to_str())
1402 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1403 if !crate::ask::valid_asset_name(name) {
1404 bail!(
1405 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1406 with no `..`; rename the file and try again"
1407 );
1408 }
1409 if !src.is_file() {
1410 bail!("attachment `{}` is not a file", src.display());
1411 }
1412 wanted.push((src, name));
1413 }
1414 let dir = self.attachments_dir(&task.id);
1415 let existed = dir.is_dir();
1416 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1417 let mut created: Vec<PathBuf> = Vec::new();
1418 let mut names = Vec::new();
1419 let mut copy_all = || -> Result<()> {
1420 for (src, name) in &wanted {
1421 let (stored, path) = copy_new(&dir, src, name)?;
1422 created.push(path);
1423 names.push(stored);
1424 }
1425 Ok(())
1426 };
1427 if let Err(e) = copy_all() {
1428 for path in &created {
1429 let _ = std::fs::remove_file(path);
1430 }
1431 if !existed {
1432 let _ = std::fs::remove_dir(&dir);
1433 }
1434 return Err(e);
1435 }
1436 task.attachments.extend(names.iter().cloned());
1437 Ok(names)
1438 }
1439
1440 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1451 let dir = self.attachments_dir(&task.id);
1452 let existed = dir.is_dir();
1453 let before = task.attachments.len();
1454 let names = self.attach(task, sources)?;
1455 if let Err(e) = self.put(task) {
1456 for name in &names {
1457 let _ = std::fs::remove_file(dir.join(name));
1458 }
1459 if !existed {
1460 let _ = std::fs::remove_dir(&dir);
1461 }
1462 task.attachments.truncate(before);
1463 return Err(e);
1464 }
1465 Ok(names)
1466 }
1467
1468 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1472 let dir = self.attachments_dir(&task.id);
1473 task.attachments
1474 .iter()
1475 .map(|n| {
1476 let p = dir.join(n);
1477 std::path::absolute(&p).unwrap_or(p)
1478 })
1479 .collect()
1480 }
1481
1482 pub fn put(&self, task: &mut Task) -> Result<()> {
1491 let _lock = self.lock_task(&task.id)?;
1492 self.put_unlocked(task)
1493 }
1494
1495 pub fn create_new(&self, task: &mut Task) -> Result<bool> {
1500 let _lock = self.lock_task(&task.id)?;
1501 if self.path_of(&task.id).exists() {
1502 return Ok(false);
1503 }
1504 self.put_unlocked(task)?;
1505 Ok(true)
1506 }
1507
1508 fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1510 if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1511 for run in stored.runs {
1512 if !task.runs.contains(&run) {
1513 task.runs.push(run);
1514 }
1515 }
1516 }
1517 task.updated_at = Timestamp::now();
1518 std::fs::create_dir_all(&self.root)
1519 .with_context(|| format!("create {}", self.root.display()))?;
1520 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1521 let path = self.path_of(&task.id);
1522 let tmp = path.with_extension("json.tmp");
1523 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1524 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1525 if let (Some(notice), Some(home)) = (
1530 crate::notices::task_held(task),
1531 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1532 ) {
1533 crate::notices::raise_in(home, notice);
1534 }
1535 Ok(())
1536 }
1537
1538 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1543 let id = self.resolve_id(id)?;
1544 let _lock = self.lock_task(&id)?;
1547 let mut task = self.get(&id)?;
1548 if task.link_run(run) {
1549 self.put_unlocked(&mut task)?;
1550 }
1551 Ok(task)
1552 }
1553
1554 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1562 std::fs::create_dir_all(&self.root)
1563 .with_context(|| format!("create {}", self.root.display()))?;
1564 let path = self.root.join(format!("{id}.write-lock"));
1565 let started = std::time::Instant::now();
1566 loop {
1567 match std::fs::OpenOptions::new()
1568 .write(true)
1569 .create_new(true)
1570 .open(&path)
1571 {
1572 Ok(mut file) => {
1573 use std::io::Write;
1574 let token = owner_token();
1575 if let Err(e) = file.write_all(token.as_bytes()) {
1576 drop(file);
1577 let _ = std::fs::remove_file(&path);
1578 return Err(e).with_context(|| format!("lock {}", path.display()));
1579 }
1580 return Ok(TaskLock { path, token });
1581 }
1582 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1583 let judged = std::fs::read_to_string(&path).ok();
1586 let broken = judged
1587 .filter(|_| older_than_stale(&path))
1588 .is_some_and(|judged| break_stale(&path, &judged));
1589 if !broken {
1593 if started.elapsed() > TASK_LOCK_STALE {
1594 bail!("could not lock task {id}");
1595 }
1596 std::thread::sleep(std::time::Duration::from_millis(15));
1597 }
1598 }
1599 Err(e)
1603 if e.kind() == std::io::ErrorKind::PermissionDenied
1604 && started.elapsed() <= TASK_LOCK_STALE =>
1605 {
1606 std::thread::sleep(std::time::Duration::from_millis(15));
1607 }
1608 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1609 }
1610 }
1611 }
1612
1613 pub fn get(&self, id: &str) -> Result<Task> {
1615 let resolved = self.resolve_id(id)?;
1616 read_path(&self.path_of(&resolved))
1617 }
1618
1619 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1651 let _ = questions;
1652 let resolved = self.resolve_id(id)?;
1653 if in_flight {
1654 bail!("task {resolved} is being run by a live daemon right now");
1655 }
1656 self.write_tombstone(&resolved)?;
1657 if let Err(e) = self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))
1658 {
1659 if self.path_of(&resolved).exists() {
1663 let _ = std::fs::remove_file(self.tombstone_path(&resolved));
1664 }
1665 return Err(e);
1666 }
1667 let (released, still_blocked) = self.release_dependents_of(&resolved);
1668 Ok(Removal {
1669 id: resolved,
1670 released,
1671 still_blocked,
1672 })
1673 }
1674
1675 fn tombstone_path(&self, id: &str) -> PathBuf {
1678 self.root.join(format!("{id}.removed"))
1679 }
1680
1681 fn write_tombstone(&self, id: &str) -> Result<()> {
1682 let path = self.tombstone_path(id);
1683 let tmp = self.root.join(format!("{id}.removed.tmp"));
1684 std::fs::write(&tmp, b"")
1685 .and_then(|()| std::fs::rename(&tmp, &path))
1686 .with_context(|| format!("write {}", path.display()))
1687 }
1688
1689 fn was_deleted(&self, id: &str) -> bool {
1693 self.tombstone_path(id).is_file() && !self.path_of(id).exists()
1694 }
1695
1696 pub fn apply_deleted_blockers(&self, task: &mut Task) -> Vec<String> {
1700 let deleted = deleted_blockers(self, &task.blocked_by);
1701 deleted
1702 .into_iter()
1703 .filter(|id| task.dependency_deleted(id))
1704 .collect()
1705 }
1706
1707 pub fn note_dependency_deleted(&self, dependent: &Task, deleted: &str) {
1711 let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty()) else {
1712 return;
1713 };
1714 let ja = crate::lang::is_japanese(&crate::lang::of_repo(&dependent.repo));
1715 let waiting = dependent.status == TaskStatus::Blocked;
1716 let message = match (ja, waiting) {
1717 (true, false) => format!(
1718 "タスク {} は、待っていた {} が削除されたため待機を解除し、元の状態に戻しました",
1719 dependent.short(),
1720 short(deleted)
1721 ),
1722 (true, true) => format!(
1723 "タスク {} は、待っていた {} が削除されたため、残りの依存を待っています",
1724 dependent.short(),
1725 short(deleted)
1726 ),
1727 (false, false) => format!(
1728 "Task {} stopped waiting on {} because it was deleted, and returned to its previous state",
1729 dependent.short(),
1730 short(deleted)
1731 ),
1732 (false, true) => format!(
1733 "Task {} stopped waiting on {} because it was deleted, and is still waiting on its other dependencies",
1734 dependent.short(),
1735 short(deleted)
1736 ),
1737 };
1738 crate::notices::raise_in(
1739 home,
1740 crate::notices::Notice::info(
1741 &format!("unblocked:{}:{}", dependent.id, deleted),
1742 message,
1743 ),
1744 );
1745 }
1746
1747 fn release_dependents_of(&self, dependency: &str) -> (Vec<String>, Vec<String>) {
1752 let (mut released, mut still_blocked) = (Vec::new(), Vec::new());
1753 for listed in self.list() {
1754 if listed.status != TaskStatus::Blocked
1755 || !listed.blocked_by.iter().any(|b| b == dependency)
1756 {
1757 continue;
1758 }
1759 let Ok(_claim) = self.claim(&listed.id) else {
1760 continue;
1761 };
1762 let Ok(mut task) = self.get(&listed.id) else {
1763 continue;
1764 };
1765 if !task.dependency_deleted(dependency) {
1766 continue;
1767 }
1768 if self.put(&mut task).is_err() {
1769 continue;
1770 }
1771 self.note_dependency_deleted(&task, dependency);
1772 if task.status == TaskStatus::Blocked {
1773 still_blocked.push(task.id.clone());
1774 } else {
1775 released.push(task.id.clone());
1776 }
1777 }
1778 (released, still_blocked)
1779 }
1780
1781 fn remove_record_with_attachments(
1785 &self,
1786 resolved: &str,
1787 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1788 ) -> Result<()> {
1789 self.sweep_removed_attachments();
1796 let attachments = self.attachments_dir(resolved);
1797 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1798 let moved = match std::fs::rename(&attachments, &aside) {
1799 Ok(()) => true,
1800 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1801 Err(e) => {
1802 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1803 }
1804 };
1805 let path = self.path_of(resolved);
1806 if let Err(e) = remove_record(&path) {
1807 if moved {
1808 let _ = std::fs::rename(&aside, &attachments);
1809 }
1810 return Err(e).with_context(|| format!("remove {}", path.display()));
1811 }
1812 if moved {
1813 if let Err(e) = std::fs::remove_dir_all(&aside) {
1814 tracing::warn!("leftover attachments {}: {e}", aside.display());
1815 }
1816 }
1817 let lock = self.lock_path(resolved);
1818 if let Err(e) = std::fs::remove_file(&lock) {
1819 if e.kind() != std::io::ErrorKind::NotFound {
1820 return Err(e).with_context(|| format!("remove {}", lock.display()));
1821 }
1822 }
1823 Ok(())
1824 }
1825
1826 fn sweep_removed_attachments(&self) {
1832 let Ok(entries) = std::fs::read_dir(&self.root) else {
1833 return;
1834 };
1835 for entry in entries.flatten() {
1836 let name = entry.file_name();
1837 let name = name.to_string_lossy();
1838 let Some(id) = name.strip_suffix(".attachments.removing") else {
1839 continue;
1840 };
1841 if !self.path_of(id).exists() {
1844 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1845 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1846 }
1847 }
1848 }
1849 self.sweep_tombstones();
1850 }
1851
1852 fn sweep_tombstones(&self) {
1861 let Ok(entries) = std::fs::read_dir(&self.root) else {
1862 return;
1863 };
1864 let tasks = self.list();
1865 for entry in entries.flatten() {
1866 let name = entry.file_name();
1867 let name = name.to_string_lossy();
1868 let Some(id) = name.strip_suffix(".removed") else {
1869 continue;
1870 };
1871 if self.path_of(id).exists()
1872 || tasks.iter().any(|t| t.blocked_by.iter().any(|b| b == id))
1873 {
1874 continue;
1875 }
1876 let old = entry
1877 .metadata()
1878 .and_then(|m| m.modified())
1879 .ok()
1880 .and_then(|m| m.elapsed().ok())
1881 .is_some_and(|age| age >= TOMBSTONE_GRACE);
1882 if old {
1883 let _ = std::fs::remove_file(entry.path());
1884 }
1885 }
1886 }
1887
1888 fn lock_path(&self, id: &str) -> PathBuf {
1891 self.root.join(format!("{id}.lock"))
1892 }
1893
1894 pub fn list(&self) -> Vec<Task> {
1907 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1908 .into_iter()
1909 .flatten()
1910 .flatten()
1911 .map(|e| e.path())
1912 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1913 .filter_map(|p| read_path(&p).ok())
1914 .collect();
1915 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1916 tasks
1917 }
1918
1919 pub fn superseded(&self) -> HashMap<String, String> {
1932 let mut by = HashMap::new();
1933 for task in self.list() {
1934 for earlier in &task.runs {
1935 if let Some(later) = task.successor_of(earlier) {
1936 by.insert(earlier.clone(), later.clone());
1937 }
1938 }
1939 }
1940 by
1941 }
1942
1943 pub fn superseded_by(&self, run: &str) -> Option<String> {
1951 for task in self.list() {
1952 if task.runs.iter().any(|r| r == run) {
1953 return task.successor_of(run).cloned();
1954 }
1955 }
1956 None
1957 }
1958
1959 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1970 for task in self.list() {
1971 if task.runs.iter().any(|r| r == run) {
1972 return task.runs.last().filter(|last| **last != run).cloned();
1973 }
1974 }
1975 None
1976 }
1977
1978 pub fn next_runnable(&self) -> Option<Task> {
1983 let mut runnable: Vec<Task> = self
1984 .list()
1985 .into_iter()
1986 .filter(|t| t.status.runnable())
1987 .collect();
1988 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1989 runnable.into_iter().next()
1990 }
1991
1992 pub fn claim(&self, id: &str) -> Result<Claim> {
1999 std::fs::create_dir_all(&self.root)
2000 .with_context(|| format!("create {}", self.root.display()))?;
2001 let path = self.lock_path(id);
2002 match std::fs::OpenOptions::new()
2003 .write(true)
2004 .create_new(true)
2005 .open(&path)
2006 {
2007 Ok(mut f) => {
2008 use std::io::Write as _;
2009 let _ = writeln!(f, "{}", std::process::id());
2011 Ok(Claim { path })
2012 }
2013 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
2014 bail!("task {id} is already claimed ({} exists)", path.display())
2015 }
2016 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
2017 }
2018 }
2019
2020 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
2022 if self.path_of(prefix).is_file() {
2023 return Ok(prefix.to_owned());
2024 }
2025 let hits: Vec<String> = self
2026 .list()
2027 .into_iter()
2028 .map(|t| t.id)
2029 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
2030 .collect();
2031 match hits.len() {
2032 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
2033 0 => bail!("no task matches `{prefix}`"),
2034 _ => bail!(
2035 "`{prefix}` matches {} tasks: {}",
2036 hits.len(),
2037 hits.join(", ")
2038 ),
2039 }
2040 }
2041
2042 pub fn revision(&self) -> u64 {
2049 self.revision_excluding(&std::collections::BTreeSet::new())
2050 }
2051
2052 pub fn revision_excluding(&self, skip: &std::collections::BTreeSet<String>) -> u64 {
2055 use std::hash::{Hash as _, Hasher as _};
2056
2057 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
2058 .into_iter()
2059 .flatten()
2060 .flatten()
2061 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
2062 .filter(|e| {
2063 let path = e.path();
2064 !path
2065 .file_stem()
2066 .is_some_and(|stem| skip.contains(stem.to_string_lossy().as_ref()))
2067 })
2068 .filter_map(|e| {
2069 let name = e.file_name().to_string_lossy().into_owned();
2070 let mtime = e
2071 .metadata()
2072 .ok()?
2073 .modified()
2074 .ok()?
2075 .duration_since(std::time::UNIX_EPOCH)
2076 .ok()?
2077 .as_millis() as u64;
2078 Some((name, mtime))
2079 })
2080 .collect();
2081
2082 if entries.is_empty() {
2083 return 0;
2084 }
2085
2086 entries.sort_unstable();
2087 let mut hasher = std::hash::DefaultHasher::new();
2088 for (name, mtime) in &entries {
2089 name.hash(&mut hasher);
2090 mtime.hash(&mut hasher);
2091 }
2092 let h = hasher.finish();
2093 if h == 0 { 1 } else { h }
2094 }
2095}
2096
2097const TOMBSTONE_GRACE: std::time::Duration = std::time::Duration::from_secs(3600);
2100
2101#[derive(Debug, Clone)]
2103pub struct Removal {
2104 pub id: String,
2106 pub released: Vec<String>,
2110 pub still_blocked: Vec<String>,
2113}
2114
2115#[derive(Debug)]
2117pub struct Claim {
2118 path: PathBuf,
2119}
2120
2121impl Drop for Claim {
2122 fn drop(&mut self) {
2123 let _ = std::fs::remove_file(&self.path);
2124 }
2125}
2126
2127pub fn title_from(instruction: &str, max: usize) -> String {
2130 let Some(line) = first_line(instruction) else {
2131 return "(empty task)".to_owned();
2132 };
2133 if line.chars().count() <= max {
2134 return line.to_owned();
2135 }
2136 let head: String = line.chars().take(max.saturating_sub(1)).collect();
2137 format!("{head}…")
2138}
2139
2140pub(crate) fn first_line(instruction: &str) -> Option<&str> {
2146 let line = instruction
2147 .lines()
2148 .map(str::trim)
2149 .find(|l| !l.is_empty())?
2150 .trim_start_matches(['#', '-', '*', '>', ' '])
2151 .trim();
2152 (!line.is_empty()).then_some(line)
2153}
2154
2155fn read_path(path: &Path) -> Result<Task> {
2156 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2157 let task: Task =
2158 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
2159 if task.schema > SCHEMA {
2165 bail!(
2166 "task {} was written by a different magi (schema {}, this build \
2167 speaks {SCHEMA})",
2168 task.id,
2169 task.schema
2170 );
2171 }
2172 Ok(task)
2173}
2174
2175pub fn missing_blockers(
2189 queue: &Queue,
2190 questions: &Questions,
2191 blocked_by: &[String],
2192) -> Vec<String> {
2193 blocked_by
2194 .iter()
2195 .filter(|id| {
2196 !queue.path_of(id).is_file()
2197 && !questions.path_of(id).is_file()
2198 && !queue.was_deleted(id)
2199 })
2200 .cloned()
2201 .collect()
2202}
2203
2204pub fn deleted_blockers(queue: &Queue, blocked_by: &[String]) -> Vec<String> {
2209 blocked_by
2210 .iter()
2211 .filter(|id| queue.was_deleted(id))
2212 .cloned()
2213 .collect()
2214}
2215
2216pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
2228 missing_blocker_hold_reason_in(blocked_by, missing, "en")
2229}
2230
2231pub fn missing_blocker_hold_reason_in(
2233 blocked_by: &[String],
2234 missing: &[String],
2235 language: &str,
2236) -> String {
2237 if crate::lang::is_japanese(language) {
2238 format!(
2239 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
2240 blocked_by.join(", "),
2241 missing.join(", "),
2242 )
2243 } else {
2244 format!(
2245 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
2246 blocked_by.join(", "),
2247 missing.join(", "),
2248 )
2249 }
2250}
2251
2252pub fn short(id: &str) -> &str {
2254 id.split('-').next_back().unwrap_or(id)
2255}
2256
2257fn new_id() -> String {
2258 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2259 let seed = crate::rng::entropy();
2260 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2261}
2262
2263#[cfg(test)]
2264mod tests {
2265 #[test]
2266 fn merge_over_lets_the_later_choice_win_and_keeps_the_rest() {
2267 let mut base = RunOverrides {
2268 merge: Some("pr".to_owned()),
2269 candidates: Some(3),
2270 ..RunOverrides::default()
2271 };
2272 base.merge_over(&RunOverrides {
2273 merge: Some("none".to_owned()),
2274 seed: Some(7),
2275 ..RunOverrides::default()
2276 });
2277 assert_eq!(base.merge.as_deref(), Some("none"));
2278 assert_eq!(base.candidates, Some(3));
2279 assert_eq!(base.seed, Some(7));
2280 }
2281
2282 #[test]
2283 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
2284 let mut t = task("relanded");
2285 t.attempts = 1;
2286 t.status = TaskStatus::Running;
2287 t.already_landed("already in main as 0e368de");
2288 assert_eq!(t.status, TaskStatus::Done);
2289 assert_eq!(t.attempts, 0);
2290 assert!(t.hold_reason.is_none());
2291 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
2292 }
2293
2294 #[test]
2295 fn missing_blocker_reason_follows_the_language() {
2296 let b = vec!["a".to_owned()];
2297 let en = missing_blocker_hold_reason_in(&b, &b, "en");
2298 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
2299 assert!(en.starts_with("blocked on a"));
2300 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
2301 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
2302 }
2303
2304 use super::*;
2305
2306 #[test]
2307 fn triage_applied_survives_release_and_old_records_read_as_empty() {
2308 let mut t = Task::new(
2309 "t".to_owned(),
2310 "i".to_owned(),
2311 PathBuf::from("r"),
2312 Source::Human,
2313 );
2314 t.mark_triage_applied("q1");
2315 t.mark_triage_applied("q1");
2316 t.hold_machine(Some("x".to_owned()));
2317 t.release();
2318 assert_eq!(t.triage_applied, ["q1"]);
2319 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
2320
2321 let mut v = serde_json::to_value(&t).unwrap();
2322 v.as_object_mut().unwrap().remove("triage_applied");
2323 let old: Task = serde_json::from_value(v).unwrap();
2324 assert!(old.triage_applied.is_empty());
2325 }
2326
2327 #[test]
2328 fn task_counts_of_empty_is_all_zero() {
2329 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
2330 }
2331
2332 #[test]
2333 fn task_counts_of_tallies_every_status() {
2334 let mut queued = Task::new(
2335 "q".to_owned(),
2336 "i".to_owned(),
2337 PathBuf::from("."),
2338 Source::Human,
2339 );
2340 queued.status = TaskStatus::Queued;
2341 let mut running = queued.clone();
2342 running.status = TaskStatus::Running;
2343 let mut done = queued.clone();
2344 done.status = TaskStatus::Done;
2345 let mut failed = queued.clone();
2346 failed.status = TaskStatus::Failed;
2347 let mut held = queued.clone();
2348 held.status = TaskStatus::Held;
2349 let mut blocked = queued.clone();
2350 blocked.status = TaskStatus::Blocked;
2351
2352 let mut parked = queued.clone();
2353 parked.status = TaskStatus::Parked;
2354
2355 let counts = TaskCounts::of(&[
2356 queued,
2357 running,
2358 done.clone(),
2359 done,
2360 failed,
2361 held,
2362 blocked,
2363 parked,
2364 ]);
2365 assert_eq!(
2366 counts,
2367 TaskCounts {
2368 queued: 1,
2369 running: 1,
2370 done: 2,
2371 failed: 1,
2372 held: 1,
2373 blocked: 1,
2374 parked: 1,
2375 }
2376 );
2377 }
2378
2379 fn queue() -> (tempfile::TempDir, Queue) {
2382 let dir = tempfile::tempdir().unwrap();
2383 let q = Queue::at(dir.path().join("queue"));
2384 (dir, q)
2385 }
2386
2387 #[test]
2388 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
2389 let dir = tempfile::tempdir().unwrap();
2390 let q = Queue::at(dir.path().join("queue"));
2391 let mut t = task("held");
2392 q.put(&mut t).unwrap();
2393 assert_eq!(
2394 crate::notices::Notices::at(dir.path().join("notifications"))
2395 .list()
2396 .len(),
2397 0
2398 );
2399 t.hold_machine(Some("out of attempts".to_owned()));
2400 q.put(&mut t).unwrap();
2401 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
2402 assert_eq!(listed.len(), 1);
2403 assert!(listed[0].message.contains("out of attempts"));
2404 }
2405
2406 fn task(title: &str) -> Task {
2407 Task::new(
2408 title.to_owned(),
2409 format!("do {title}"),
2410 PathBuf::from("."),
2411 Source::Human,
2412 )
2413 }
2414
2415 #[test]
2416 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
2417 let mut t = task("retried");
2418 assert!(
2419 t.earlier_attempts().is_empty(),
2420 "a first attempt takes nothing over"
2421 );
2422 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2423 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
2424 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
2426 assert_eq!(t.successor_of("bbbb"), None);
2427 assert_eq!(t.successor_of("zzzz"), None);
2428 }
2429
2430 #[test]
2431 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
2432 let (_dir, q) = queue();
2433 let mut t = task("retried");
2434 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2435 q.put(&mut t).unwrap();
2436
2437 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
2438 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
2439 assert_eq!(
2440 q.superseded_by("cccc"),
2441 None,
2442 "the latest attempt replaces nothing"
2443 );
2444 assert_eq!(
2445 q.superseded_by("never-heard-of-it"),
2446 None,
2447 "a run belonging to no task on this queue is not superseded"
2448 );
2449
2450 let mut by = HashMap::new();
2451 by.insert("aaaa".to_owned(), "bbbb".to_owned());
2452 by.insert("bbbb".to_owned(), "cccc".to_owned());
2453 assert_eq!(
2454 q.superseded(),
2455 by,
2456 "the whole-map and single-run forms must agree"
2457 );
2458 }
2459
2460 #[test]
2461 fn superseded_attempts_is_empty_until_the_task_is_done() {
2462 let mut t = task("retried");
2463 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2464 t.status = TaskStatus::Failed;
2465 assert_eq!(
2466 t.superseded_attempts(true),
2467 &[] as &[String],
2468 "a task still retrying has no attempt yet that a later one made moot"
2469 );
2470
2471 t.status = TaskStatus::Running;
2472 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2473 }
2474
2475 #[test]
2476 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
2477 let mut t = task("retried");
2478 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2479 t.status = TaskStatus::Done;
2480 assert_eq!(
2481 t.superseded_attempts(true),
2482 &["aaaa".to_owned(), "bbbb".to_owned()],
2483 "cccc is the attempt whose success made the task done, and stays out"
2484 );
2485 }
2486
2487 #[test]
2488 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
2489 let mut t = task("first try landed");
2490 t.runs = vec!["aaaa".to_owned()];
2491 t.status = TaskStatus::Done;
2492 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2493 }
2494
2495 #[test]
2496 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
2497 let mut t = task("closed by hand after a manual merge");
2504 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2505 t.status = TaskStatus::Done;
2506 assert_eq!(
2507 t.superseded_attempts(false),
2508 &[] as &[String],
2509 "nothing here is provably why the task is done, so nothing is superseded"
2510 );
2511 }
2512
2513 #[test]
2514 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2515 let (_dir, q) = queue();
2516 let mut t = task("retried twice");
2517 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2518 q.put(&mut t).unwrap();
2519
2520 assert_eq!(
2521 q.latest_attempt("aaaa"),
2522 Some("cccc".to_owned()),
2523 "an old attempt points straight at the chain's current head, not the \
2524 next attempt in the middle of it"
2525 );
2526 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2527 assert_eq!(
2528 q.latest_attempt("cccc"),
2529 None,
2530 "the latest attempt is not superseded by anything"
2531 );
2532 assert_eq!(
2533 q.latest_attempt("never-heard-of-it"),
2534 None,
2535 "a run belonging to no task on this queue is not superseded"
2536 );
2537 }
2538
2539 #[test]
2540 fn a_markdown_heading_is_the_title_not_decoration() {
2541 assert_eq!(
2546 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2547 "Rework the config loader"
2548 );
2549 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2550 assert_eq!(title_from("> quoted task", 40), "quoted task");
2551 assert_eq!(title_from(" \n\n", 40), "(empty task)");
2553 assert_eq!(title_from("###\n", 40), "(empty task)");
2554 }
2555
2556 #[test]
2557 fn a_long_title_is_elided_by_characters_not_bytes() {
2558 let long = "課題".repeat(30);
2560 let title = title_from(&long, 10);
2561 assert_eq!(title.chars().count(), 10);
2562 assert!(title.ends_with('…'));
2563 }
2564
2565 #[test]
2566 fn priority_wins_and_ties_break_oldest_first() {
2567 let (_dir, q) = queue();
2568 let mut a = task("first");
2569 let mut b = task("second");
2570 let mut c = task("urgent");
2571 a.id = "20260101-000001-aaaa".to_owned();
2573 b.id = "20260101-000002-bbbb".to_owned();
2574 c.id = "20260101-000003-cccc".to_owned();
2575 c.priority = 5;
2576 for t in [&mut a, &mut b, &mut c] {
2577 q.put(t).unwrap();
2578 }
2579
2580 assert_eq!(q.next_runnable().unwrap().id, c.id);
2582 c.hold_machine(None);
2583 q.put(&mut c).unwrap();
2584 assert_eq!(q.next_runnable().unwrap().id, a.id);
2586 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2587 }
2588
2589 #[test]
2590 fn a_blocked_task_never_starves_another_runnable_one() {
2591 let (_dir, q) = queue();
2592 let mut blocked = task("blocked");
2593 blocked.block(vec!["something".to_owned()], None);
2594 q.put(&mut blocked).unwrap();
2595
2596 let mut runnable = task("free to go");
2597 q.put(&mut runnable).unwrap();
2598
2599 let next = q.next_runnable().expect("a runnable task is still offered");
2600 assert_eq!(next.id, runnable.id);
2601 }
2602
2603 #[test]
2604 fn a_held_task_is_never_offered_to_the_loop() {
2605 let (_dir, q) = queue();
2606 let mut t = task("held");
2607 q.put(&mut t).unwrap();
2608 assert!(q.next_runnable().is_some());
2609
2610 t.hold_machine(None);
2611 q.put(&mut t).unwrap();
2612 assert!(
2613 q.next_runnable().is_none(),
2614 "a held task must wait for a human"
2615 );
2616
2617 t.status = TaskStatus::Failed;
2619 q.put(&mut t).unwrap();
2620 assert!(q.next_runnable().is_some());
2621 }
2622
2623 #[test]
2624 fn attempts_are_capped_and_then_the_task_is_held() {
2625 let mut t = task("doomed");
2626
2627 t.start("run-1".to_owned());
2628 t.fail("gate red", 2);
2629 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2630
2631 t.start("run-2".to_owned());
2632 t.fail("gate red", 2);
2633 assert_eq!(
2634 t.status,
2635 TaskStatus::Held,
2636 "out of attempts: stop spending money on it"
2637 );
2638 assert_eq!(t.runs, ["run-1", "run-2"]);
2639 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2640 assert_eq!(
2641 t.hold_reason.as_deref(),
2642 Some("gate red"),
2643 "the hold must say why, not leave hold_reason null next to a \
2644 populated last_error"
2645 );
2646 }
2647
2648 #[test]
2649 fn handing_off_a_task_records_a_hold_reason_too() {
2650 let mut t = task("left a pull request");
2651 t.start("run-1".to_owned());
2652 t.handed_off("run ended with a pull request open [run run-1]");
2653 assert_eq!(t.status, TaskStatus::Held);
2654 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2655 assert_eq!(
2656 t.hold_reason.as_deref(),
2657 Some("run ended with a pull request open [run run-1]")
2658 );
2659 assert_eq!(t.hold_reason, t.last_error);
2660 }
2661
2662 #[test]
2663 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2664 let mut t = task("stalled by quota");
2665
2666 t.start("run-1".to_owned());
2667 assert_eq!(t.attempts, 1);
2668 t.stall("judge-1, judge-2 out of quota");
2669 assert_eq!(
2670 t.attempts, 0,
2671 "a closed quota window must not spend the task's retry budget"
2672 );
2673 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2674 assert_eq!(
2675 t.last_error.as_deref(),
2676 Some("judge-1, judge-2 out of quota")
2677 );
2678
2679 for _ in 0..20 {
2682 t.start("run-n".to_owned());
2683 t.stall("still out of quota");
2684 }
2685 t.start("run-real".to_owned());
2686 t.fail("gate red", 2);
2687 assert_eq!(
2688 t.status,
2689 TaskStatus::Failed,
2690 "the first attempt that was really judged is attempt one"
2691 );
2692 }
2693
2694 #[test]
2695 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2696 let mut t = task("retry me");
2697 t.start("run-1".to_owned());
2698 t.fail("gate red", 1);
2699 assert_eq!(t.status, TaskStatus::Held);
2700
2701 t.release();
2702 assert_eq!(t.status, TaskStatus::Queued);
2703 assert_eq!(t.attempts, 0);
2706 assert!(t.last_error.is_none());
2707 assert_eq!(
2708 t.runs.len(),
2709 1,
2710 "history is kept: attempts reset, evidence does not"
2711 );
2712 }
2713
2714 #[test]
2715 fn a_hold_reason_survives_and_a_release_clears_it() {
2716 let mut t = task("waiting on something else");
2717 t.hold_manual(Some(
2718 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2719 ));
2720 assert_eq!(t.status, TaskStatus::Held);
2721 assert_eq!(
2722 t.hold_reason.as_deref(),
2723 Some("waiting for 20260101-000000-aaaa to land first")
2724 );
2725
2726 t.hold_manual(None);
2728 assert_eq!(
2729 t.hold_reason.as_deref(),
2730 Some("waiting for 20260101-000000-aaaa to land first"),
2731 "a bare re-hold keeps whatever a human already wrote down"
2732 );
2733
2734 let mut plain = task("no reason given");
2736 plain.hold_manual(None);
2737 assert_eq!(plain.status, TaskStatus::Held);
2738 assert!(plain.hold_reason.is_none());
2739
2740 t.release();
2741 assert_eq!(t.status, TaskStatus::Queued);
2742 assert!(
2743 t.hold_reason.is_none(),
2744 "a stale reason must not greet the next person who holds this task"
2745 );
2746 }
2747
2748 #[test]
2749 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2750 let mut t = task("landed by hand while held");
2755 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2756 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2757
2758 t.succeed();
2759 assert_eq!(t.status, TaskStatus::Done);
2760 assert!(
2761 t.hold_reason.is_none(),
2762 "a done task cannot still be waiting on something"
2763 );
2764 }
2765
2766 #[test]
2767 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2768 let mut held = task("held straight out of blocked");
2775 held.block(
2776 vec!["20260101-000000-dead".to_owned()],
2777 Some("waiting on the migration script".to_owned()),
2778 );
2779 assert_eq!(held.status, TaskStatus::Blocked);
2780
2781 held.hold_manual(None);
2782 assert_eq!(held.status, TaskStatus::Held);
2783 assert!(
2784 held.blocked_by.is_empty(),
2785 "hold overrides the wait, same as release"
2786 );
2787 assert!(held.block_reason.is_none());
2788
2789 let mut done = task("closed straight out of blocked");
2790 done.block(
2791 vec!["20260101-000000-dead".to_owned()],
2792 Some("waiting on the migration script".to_owned()),
2793 );
2794 done.succeed();
2795 assert_eq!(done.status, TaskStatus::Done);
2796 assert!(
2797 done.blocked_by.is_empty(),
2798 "a done task cannot still be waiting on a dependency"
2799 );
2800 assert!(done.block_reason.is_none());
2801 }
2802
2803 #[test]
2804 fn a_blocked_task_is_never_offered_to_the_loop() {
2805 let mut t = task("blocked");
2806 assert!(t.status.runnable());
2807 t.block(
2808 vec!["dep-id".to_owned()],
2809 Some("waits on dep-id".to_owned()),
2810 );
2811 assert_eq!(t.status, TaskStatus::Blocked);
2812 assert!(!t.status.runnable());
2813 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2814 }
2815
2816 #[test]
2817 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2818 let mut t = task("blocked on two");
2819 t.block(
2820 vec!["a".to_owned(), "b".to_owned()],
2821 Some("waits on a and b".to_owned()),
2822 );
2823
2824 t.unblock("a");
2825 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2826 assert_eq!(t.blocked_by, ["b"]);
2827
2828 t.unblock("b");
2829 assert_eq!(t.status, TaskStatus::Queued);
2830 assert!(t.blocked_by.is_empty());
2831 assert!(t.block_reason.is_none());
2832 }
2833
2834 #[test]
2835 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2836 let mut t = task("never blocked");
2837 t.unblock("whatever");
2838 assert_eq!(t.status, TaskStatus::Queued);
2839 }
2840
2841 #[test]
2842 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2843 let mut t = task("held, then asked about");
2848 t.hold_machine(Some("out of attempts".to_owned()));
2849 assert_eq!(t.status, TaskStatus::Held);
2850
2851 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2852 assert_eq!(t.status, TaskStatus::Blocked);
2853
2854 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2855 t.unblock("q1");
2856 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2857 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2858 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2859 assert!(t.blocked_from.is_none(), "consumed once restored");
2860 }
2861
2862 #[test]
2863 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2864 let mut t = task("manually held, then asked about");
2865 t.hold_manual(Some("waiting on a dependency".to_owned()));
2866
2867 t.block(vec!["q1".to_owned()], None);
2868 t.unblock("q1");
2869
2870 assert_eq!(t.status, TaskStatus::Held);
2871 assert_eq!(t.hold_source, Some(HoldSource::Manual));
2872 }
2873
2874 #[test]
2875 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2876 let mut t = task("held, blocked twice");
2880 t.hold_machine(None);
2881 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2882 t.block(
2883 vec!["q1".to_owned(), "q2".to_owned()],
2884 Some("second".to_owned()),
2885 );
2886
2887 t.unblock("q1");
2888 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2889 t.unblock("q2");
2890 assert_eq!(t.status, TaskStatus::Held);
2891 }
2892
2893 #[test]
2894 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2895 let mut t = task("blocked mid-run");
2899 t.start("run-1".to_owned());
2900 assert_eq!(t.status, TaskStatus::Running);
2901
2902 t.block(vec!["q1".to_owned()], None);
2903 t.unblock("q1");
2904 assert_eq!(t.status, TaskStatus::Queued);
2905 }
2906
2907 #[test]
2908 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2909 let mut t = task("legacy record, held before it was blocked");
2915 t.hold_source = Some(HoldSource::Machine);
2916 t.hold_reason = Some("legacy hold reason".to_owned());
2917 t.status = TaskStatus::Blocked;
2918 t.blocked_by = vec!["q1".to_owned()];
2919 t.blocked_from = None;
2920
2921 t.unblock("q1");
2922 assert_eq!(t.status, TaskStatus::Held);
2923 }
2924
2925 #[test]
2926 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2927 let mut t = task("legacy record, ordinary dependency block");
2928 t.status = TaskStatus::Blocked;
2929 t.blocked_by = vec!["dep".to_owned()];
2930 t.blocked_from = None;
2931
2932 t.unblock("dep");
2933 assert_eq!(t.status, TaskStatus::Queued);
2934 }
2935
2936 #[test]
2937 fn answering_a_question_is_recorded_and_survives_a_release() {
2938 let mut t = task("asked something");
2939 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2940 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2941 t.unblock("q1");
2942 assert_eq!(t.status, TaskStatus::Queued);
2943 assert_eq!(t.answers.len(), 1);
2944 assert_eq!(t.answers[0].answer, "SQLite");
2945
2946 t.release();
2950 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2951 }
2952
2953 #[test]
2954 fn a_refused_handover_keeps_the_review_branch_across_release() {
2955 let mut t = task("refused takeover");
2956 t.start("run-1".to_owned());
2957 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
2958 assert_eq!(t.status, TaskStatus::Held);
2959 assert_eq!(t.attempts, 1);
2960 t.release();
2961 assert_eq!(t.status, TaskStatus::Queued);
2962 assert_eq!(t.attempts, 0);
2963 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2964
2965 t.hold_for_handover(None, "again".to_owned());
2967 t.requeue();
2968 assert!(t.review_branch.is_none());
2969
2970 let mut m = task("manual");
2972 m.review_branch = Some("magi/x/A".to_owned());
2973 m.hold_manual(None);
2974 m.release();
2975 assert!(m.review_branch.is_none());
2976 }
2977
2978 #[test]
2979 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2980 let mut t = task("blocked run with a surviving branch");
2981 t.start("run-1".to_owned());
2982 t.fail("blocked with major findings", 5);
2983 assert_eq!(t.status, TaskStatus::Failed);
2984
2985 t.request_review("magi/eba2/A".to_owned());
2986 assert_eq!(t.status, TaskStatus::Queued);
2987 assert_eq!(t.attempts, 0);
2988 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2989
2990 t.release();
2992 assert!(t.review_branch.is_none());
2993 }
2994
2995 #[test]
2996 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2997 let mut t = task("retry");
2998 t.start("run-1".to_owned());
2999 t.requeue();
3000 assert!(t.fresh_start);
3001
3002 t.release();
3003 assert!(!t.fresh_start);
3004 }
3005
3006 #[test]
3007 fn priority_can_be_changed_while_queued_but_not_while_running() {
3008 let mut t = task("reprioritise me");
3009 t.set_priority(5).unwrap();
3010 assert_eq!(t.priority, 5);
3011
3012 t.start("run-1".to_owned());
3013 let err = t.set_priority(9).unwrap_err().to_string();
3014 assert!(err.contains("running"), "{err}");
3015 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
3016 }
3017
3018 #[test]
3019 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
3020 let mut t = task("interrupt me");
3021 assert!(!t.interrupt, "off unless asked, same as any other task");
3022
3023 t.set_interrupt(true).unwrap();
3024 assert!(t.interrupt);
3025
3026 t.start("run-1".to_owned());
3027 assert!(
3028 !t.interrupt,
3029 "the mark is one-shot: dispatching the task fulfils it, \
3030 whatever the run that follows ends up doing"
3031 );
3032 let err = t.set_interrupt(true).unwrap_err().to_string();
3033 assert!(err.contains("running"), "{err}");
3034 t.set_interrupt(false).unwrap();
3037 assert!(!t.interrupt);
3038 }
3039
3040 #[test]
3044 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
3045 let mut t = task("interrupt me");
3046 t.set_interrupt(true).unwrap();
3047 t.start("run-1".to_owned());
3048 t.fail("mock failure", 5);
3049 assert_eq!(t.status, TaskStatus::Failed);
3050 assert!(
3051 !t.interrupt,
3052 "one attempt already spent the mark; a retry is an ordinary \
3053 requeue, not a fresh interrupt request"
3054 );
3055 }
3056
3057 #[test]
3058 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
3059 let (_dir, q) = queue();
3060 let mut a = task("first filed");
3061 let mut b = task("second filed");
3062 a.id = "20260101-000001-aaaa".to_owned();
3063 b.id = "20260101-000002-bbbb".to_owned();
3064 q.put(&mut a).unwrap();
3065 q.put(&mut b).unwrap();
3066
3067 assert_eq!(
3068 q.next_runnable().unwrap().id,
3069 a.id,
3070 "with equal priority the older task goes first, so a burst of \
3071 new work cannot starve it"
3072 );
3073 assert_eq!(
3074 q.list()[0].id,
3075 b.id,
3076 "but the list an operator reads is newest first, the same as \
3077 before priority existed - a's turn to run does not make it the \
3078 newest task"
3079 );
3080
3081 let mut a = q.get(&a.id).unwrap();
3082 a.set_priority(10).unwrap();
3083 q.put(&mut a).unwrap();
3084
3085 assert_eq!(
3086 q.next_runnable().unwrap().id,
3087 a.id,
3088 "a raised priority must be reflected the moment it is saved"
3089 );
3090 assert_eq!(
3094 q.list()[0].id,
3095 a.id,
3096 "the raised task must sort first in the list an operator reads, \
3097 not only in next_runnable's own ordering"
3098 );
3099 }
3100
3101 #[test]
3102 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
3103 let mut t = Task::new(
3104 "old title".to_owned(),
3105 "old instruction".to_owned(),
3106 PathBuf::from("/repo"),
3107 Source::Agent {
3108 run: "20260101-000000-beef".to_owned(),
3109 node: "implement".to_owned(),
3110 },
3111 );
3112 let id = t.id.clone();
3113 let created_at = t.created_at;
3114 t.runs.push("20260101-000000-beef".to_owned());
3115
3116 t.edit("new title".to_owned(), "new instruction".to_owned())
3117 .unwrap();
3118
3119 assert_eq!(t.title, "new title");
3120 assert_eq!(t.instruction, "new instruction");
3121 assert_eq!(t.id, id, "editing must not mint a new id");
3122 assert_eq!(t.created_at, created_at);
3123 assert_eq!(
3124 t.source,
3125 Source::Agent {
3126 run: "20260101-000000-beef".to_owned(),
3127 node: "implement".to_owned(),
3128 },
3129 "editing must not turn agent attribution into human"
3130 );
3131 assert_eq!(t.runs, ["20260101-000000-beef"]);
3132 }
3133
3134 #[test]
3135 fn editing_is_refused_once_a_task_is_running_or_finished() {
3136 let mut running = task("in flight");
3137 running.start("run-1".to_owned());
3138 let err = running
3139 .edit("x".to_owned(), "y".to_owned())
3140 .unwrap_err()
3141 .to_string();
3142 assert!(err.contains("running"), "{err}");
3143
3144 let mut done = task("finished");
3145 done.succeed();
3146 let err = done
3147 .edit("x".to_owned(), "y".to_owned())
3148 .unwrap_err()
3149 .to_string();
3150 assert!(err.contains("done"), "{err}");
3151
3152 let mut queued = task("waiting");
3154 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
3155 let mut held = task("parked");
3156 held.hold_machine(None);
3157 held.edit("x".to_owned(), "y".to_owned()).unwrap();
3158 }
3159
3160 #[test]
3161 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
3162 let (_dir, q) = queue();
3163 let path = q.path_of("20260101-000000-aaaa");
3164 std::fs::create_dir_all(q.root()).unwrap();
3165 std::fs::write(
3166 &path,
3167 serde_json::json!({
3168 "schema": SCHEMA,
3169 "id": "20260101-000000-aaaa",
3170 "title": "from before hold reasons existed",
3171 "instruction": "from before hold reasons existed",
3172 "repo": ".",
3173 "source": { "kind": "human" },
3174 "status": "held",
3175 "created_at": Timestamp::now().to_string(),
3176 "updated_at": Timestamp::now().to_string(),
3177 })
3178 .to_string(),
3179 )
3180 .unwrap();
3181
3182 let task = q.get("20260101-000000-aaaa").expect("must still read");
3183 assert!(task.hold_reason.is_none());
3184 assert!(task.operator_held());
3185 }
3186
3187 #[test]
3188 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
3189 let (_dir, q) = queue();
3190 let path = q.path_of("20260101-000000-bbbb");
3191 std::fs::create_dir_all(q.root()).unwrap();
3192 std::fs::write(
3193 &path,
3194 serde_json::json!({
3195 "schema": 2,
3196 "id": "20260101-000000-bbbb",
3197 "title": "old manual recovery",
3198 "instruction": "old manual recovery",
3199 "repo": ".",
3200 "source": { "kind": "human" },
3201 "status": "held",
3202 "hold_reason": "active manual recovery run20260912-224242-daf5",
3203 "created_at": Timestamp::now().to_string(),
3204 "updated_at": Timestamp::now().to_string(),
3205 })
3206 .to_string(),
3207 )
3208 .unwrap();
3209
3210 let task = q.get("20260101-000000-bbbb").expect("must still read");
3211 assert_eq!(task.hold_source, None);
3212 assert!(task.operator_held());
3213 }
3214
3215 #[test]
3216 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
3217 let (_dir, q) = queue();
3218 let path = q.path_of("20260101-000000-aaaa");
3219 std::fs::create_dir_all(q.root()).unwrap();
3220 std::fs::write(
3221 &path,
3222 serde_json::json!({
3223 "schema": SCHEMA,
3224 "id": "20260101-000000-aaaa",
3225 "title": "from before diagnostics existed",
3226 "instruction": "from before diagnostics existed",
3227 "repo": ".",
3228 "source": { "kind": "human" },
3229 "status": "held",
3230 "created_at": Timestamp::now().to_string(),
3231 "updated_at": Timestamp::now().to_string(),
3232 })
3233 .to_string(),
3234 )
3235 .unwrap();
3236
3237 let task = q.get("20260101-000000-aaaa").expect("must still read");
3238 assert!(task.diagnostic.is_none());
3239 }
3240
3241 #[test]
3242 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
3243 let (_dir, q) = queue();
3247 let path = q.path_of("20260101-000000-aaaa");
3248 std::fs::create_dir_all(q.root()).unwrap();
3249 std::fs::write(
3250 &path,
3251 serde_json::json!({
3252 "schema": 1,
3253 "id": "20260101-000000-aaaa",
3254 "title": "from before blocking existed",
3255 "instruction": "from before blocking existed",
3256 "repo": ".",
3257 "source": { "kind": "human" },
3258 "status": "queued",
3259 "created_at": Timestamp::now().to_string(),
3260 "updated_at": Timestamp::now().to_string(),
3261 })
3262 .to_string(),
3263 )
3264 .unwrap();
3265
3266 let task = q.get("20260101-000000-aaaa").expect("must still read");
3267 assert!(task.blocked_by.is_empty());
3268 assert!(task.block_reason.is_none());
3269 assert!(task.answers.is_empty());
3270 assert!(task.review_branch.is_none());
3271 }
3272
3273 #[test]
3274 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
3275 let mut held = task("diagnosed");
3280 held.start("run-1".to_owned());
3281 held.fail("gate red", 1);
3282 held.diagnostic = Some("cargo test failed: ...".to_owned());
3283 assert_eq!(held.status, TaskStatus::Held);
3284
3285 held.release();
3286 assert!(held.diagnostic.is_none());
3287
3288 held.diagnostic = Some("cargo test failed: ...".to_owned());
3289 held.succeed();
3290 assert!(held.diagnostic.is_none());
3291 }
3292
3293 #[test]
3294 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
3295 let mut t = task("retried");
3296 t.start("run-1".to_owned());
3297 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
3298 t.fail("unrelated config error", 5);
3299 assert_eq!(t.status, TaskStatus::Failed);
3300 assert!(
3301 t.diagnostic.is_none(),
3302 "fail() must not let an old diagnostic outlive the run that produced it"
3303 );
3304 }
3305
3306 #[test]
3307 fn a_claim_is_exclusive_and_releases_on_drop() {
3308 let (_dir, q) = queue();
3309 let mut t = task("contended");
3310 q.put(&mut t).unwrap();
3311
3312 let held = q.claim(&t.id).unwrap();
3313 assert!(
3314 q.claim(&t.id).is_err(),
3315 "two daemons must not drive one task into two runs"
3316 );
3317 drop(held);
3318 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
3319 }
3320
3321 #[test]
3322 fn a_round_trip_survives_disk() {
3323 let (_dir, q) = queue();
3324 let mut t = Task::new(
3325 "titled".to_owned(),
3326 "body".to_owned(),
3327 PathBuf::from("/repo"),
3328 Source::Agent {
3329 run: "20260101-000000-beef".to_owned(),
3330 node: "implement".to_owned(),
3331 },
3332 );
3333 t.priority = 3;
3334 q.put(&mut t).unwrap();
3335
3336 let back = q.get(&t.id).unwrap();
3337 assert_eq!(back.id, t.id);
3338 assert_eq!(back.priority, 3);
3339 assert_eq!(back.source.label(), "implement@beef");
3340 assert_eq!(q.get(t.short()).unwrap().id, t.id);
3342 }
3343
3344 #[test]
3345 fn an_unreadable_task_does_not_take_the_queue_down() {
3346 let (_dir, q) = queue();
3347 let mut t = task("fine");
3348 q.put(&mut t).unwrap();
3349 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
3350
3351 let listed = q.list();
3352 assert_eq!(listed.len(), 1, "the readable task still lists");
3353 assert_eq!(listed[0].id, t.id);
3354 }
3355
3356 #[test]
3357 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
3358 let (_dir, q) = queue();
3359 let path = q.path_of("20260101-000000-aaaa");
3360 std::fs::create_dir_all(q.root()).unwrap();
3361 std::fs::write(
3362 &path,
3363 serde_json::json!({
3364 "schema": SCHEMA,
3365 "id": "20260101-000000-aaaa",
3366 "title": "from before solo existed",
3367 "instruction": "from before solo existed",
3368 "repo": ".",
3369 "source": { "kind": "human" },
3370 "status": "queued",
3371 "created_at": Timestamp::now().to_string(),
3372 "updated_at": Timestamp::now().to_string(),
3373 })
3374 .to_string(),
3375 )
3376 .unwrap();
3377
3378 let task = q.get("20260101-000000-aaaa").expect("must still read");
3379 assert!(!task.solo, "a queue file with no `solo` field means false");
3380 }
3381
3382 #[test]
3383 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
3384 let (_dir, q) = queue();
3385 let path = q.path_of("20260101-000000-bbbb");
3386 std::fs::create_dir_all(q.root()).unwrap();
3387 std::fs::write(
3388 &path,
3389 serde_json::json!({
3390 "schema": SCHEMA,
3391 "id": "20260101-000000-bbbb",
3392 "title": "from before urgent existed",
3393 "instruction": "from before urgent existed",
3394 "repo": ".",
3395 "source": { "kind": "human" },
3396 "status": "queued",
3397 "created_at": Timestamp::now().to_string(),
3398 "updated_at": Timestamp::now().to_string(),
3399 })
3400 .to_string(),
3401 )
3402 .unwrap();
3403
3404 let task = q.get("20260101-000000-bbbb").expect("must still read");
3405 assert!(
3406 !task.urgent,
3407 "a queue file with no `urgent` field means false, same as `solo`"
3408 );
3409 }
3410
3411 #[test]
3412 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
3413 let (_dir, q) = queue();
3414 let mut t = task("from the future");
3415 q.put(&mut t).unwrap();
3416 let path = q.path_of(&t.id);
3417 let body = std::fs::read_to_string(&path)
3418 .unwrap()
3419 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
3420 std::fs::write(&path, body).unwrap();
3421
3422 let err = q.get(&t.id).unwrap_err().to_string();
3423 assert!(err.contains("schema 99"), "{err}");
3424 }
3425
3426 #[test]
3427 fn revision_moves_when_the_queue_changes() {
3428 let (_dir, q) = queue();
3429 assert_eq!(q.revision(), 0, "an empty queue has no revision");
3430 let mut t = task("first");
3431 q.put(&mut t).unwrap();
3432 assert!(q.revision() > 0, "a written task moves the revision");
3433 }
3434
3435 #[test]
3436 fn revision_moves_when_deleting_an_older_task() {
3437 let (dir, q) = queue();
3438 let questions = Questions::at(dir.path().join("questions"));
3439 let mut t1 = task("older");
3440 q.put(&mut t1).unwrap();
3441 std::thread::sleep(std::time::Duration::from_millis(10));
3443 let mut t2 = task("newer");
3444 q.put(&mut t2).unwrap();
3445
3446 let rev_before = q.revision();
3447 q.remove(&t1.id, false, &questions).unwrap();
3448 let rev_after = q.revision();
3449
3450 assert_ne!(
3451 rev_before, rev_after,
3452 "deleting an older task must change the revision so other clients see the deletion"
3453 );
3454 }
3455
3456 #[test]
3457 fn removing_a_task_takes_it_out_of_the_listing() {
3458 let (dir, q) = queue();
3459 let questions = Questions::at(dir.path().join("questions"));
3460 let mut t = task("delete me");
3461 q.put(&mut t).unwrap();
3462 let removed = q.remove(t.short(), false, &questions).unwrap();
3463 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
3464 assert!(removed.released.is_empty() && removed.still_blocked.is_empty());
3465 assert!(q.list().is_empty());
3466 assert!(
3467 q.remove(&t.id, false, &questions).is_err(),
3468 "removing twice is an error"
3469 );
3470 }
3471
3472 #[test]
3473 fn removing_a_task_takes_its_stale_lock_with_it() {
3474 let (dir, q) = queue();
3475 let questions = Questions::at(dir.path().join("questions"));
3476 let mut t = task("interrupted");
3477 q.put(&mut t).unwrap();
3478
3479 let claim = q.claim(&t.id).unwrap();
3482 std::mem::forget(claim);
3483 assert!(
3484 q.claim(&t.id).is_err(),
3485 "the orphaned lock is what makes the task look claimed"
3486 );
3487
3488 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
3490 assert!(err.contains("live daemon"), "{err}");
3491 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
3492
3493 q.remove(&t.id, false, &questions).unwrap();
3495 assert!(q.list().is_empty());
3496 let mut again = task("interrupted");
3497 again.id = t.id.clone();
3498 q.put(&mut again).unwrap();
3499 assert!(
3500 q.claim(&t.id).is_ok(),
3501 "a task that comes back must be claimable, which a left-behind lock would prevent"
3502 );
3503 }
3504
3505 fn notices_of(dir: &Path) -> Vec<crate::notices::Notice> {
3506 crate::notices::Notices::at(dir.join("notifications")).list()
3507 }
3508
3509 #[test]
3510 fn removing_a_sole_dependency_releases_the_dependent_without_a_hold() {
3511 let (dir, q) = queue();
3512 let questions = Questions::at(dir.path().join("questions"));
3513 let mut dep = task("dependency");
3514 q.put(&mut dep).unwrap();
3515 let mut blocked = task("waiting");
3516 blocked.block(vec![dep.id.clone()], Some("waits".to_owned()));
3517 q.put(&mut blocked).unwrap();
3518
3519 let removed = q.remove(&dep.id, false, &questions).unwrap();
3520 assert_eq!(removed.released, [blocked.id.clone()]);
3521 assert!(removed.still_blocked.is_empty());
3522
3523 let after = q.get(&blocked.id).unwrap();
3524 assert_eq!(after.status, TaskStatus::Queued);
3525 assert!(after.blocked_by.is_empty());
3526 assert!(after.block_reason.is_none());
3527 let notes = notices_of(dir.path());
3528 assert_eq!(notes.len(), 1, "{notes:?}");
3529 assert_eq!(notes[0].severity, crate::notices::Severity::Info);
3530 }
3531
3532 #[test]
3533 fn removing_one_of_two_dependencies_keeps_the_other() {
3534 let (dir, q) = queue();
3535 let questions = Questions::at(dir.path().join("questions"));
3536 let mut dep = task("dependency");
3537 q.put(&mut dep).unwrap();
3538 let mut other = task("other");
3539 q.put(&mut other).unwrap();
3540 let mut blocked = task("waiting");
3541 blocked.block(
3542 vec![dep.id.clone(), other.id.clone()],
3543 Some("waits on both".to_owned()),
3544 );
3545 q.put(&mut blocked).unwrap();
3546
3547 let removed = q.remove(&dep.id, false, &questions).unwrap();
3548 assert!(removed.released.is_empty());
3549 assert_eq!(removed.still_blocked, [blocked.id.clone()]);
3550
3551 let after = q.get(&blocked.id).unwrap();
3552 assert_eq!(after.status, TaskStatus::Blocked);
3553 assert_eq!(after.blocked_by, [other.id.clone()]);
3554 assert!(missing_blockers(&q, &questions, &after.blocked_by).is_empty());
3557 }
3558
3559 #[test]
3560 fn a_dependency_deleted_before_its_dependents_were_rewritten_is_released_later() {
3561 let (dir, q) = queue();
3562 let questions = Questions::at(dir.path().join("questions"));
3563 let mut dep = task("dependency");
3564 q.put(&mut dep).unwrap();
3565 let mut blocked = task("waiting");
3566 blocked.block(vec![dep.id.clone()], None);
3567 q.put(&mut blocked).unwrap();
3568
3569 let claim = q.claim(&blocked.id).unwrap();
3571 let removed = q.remove(&dep.id, false, &questions).unwrap();
3572 assert!(removed.released.is_empty());
3573 drop(claim);
3574
3575 let mut task = q.get(&blocked.id).unwrap();
3576 assert_eq!(task.status, TaskStatus::Blocked);
3577 assert!(missing_blockers(&q, &questions, &task.blocked_by).is_empty());
3578 assert_eq!(q.apply_deleted_blockers(&mut task), [dep.id.clone()]);
3579 assert_eq!(task.status, TaskStatus::Queued);
3580 }
3581
3582 #[test]
3583 fn a_dependency_without_a_tombstone_is_still_missing() {
3584 let (dir, q) = queue();
3585 let questions = Questions::at(dir.path().join("questions"));
3586 let mut dep = task("dependency");
3587 q.put(&mut dep).unwrap();
3588 std::fs::remove_file(q.path_of(&dep.id)).unwrap();
3589 assert_eq!(
3590 missing_blockers(&q, &questions, std::slice::from_ref(&dep.id)),
3591 [dep.id.clone()]
3592 );
3593 assert!(deleted_blockers(&q, std::slice::from_ref(&dep.id)).is_empty());
3594 }
3595
3596 #[test]
3597 fn a_held_dependent_is_left_alone_by_a_dependency_deletion() {
3598 let mut t = task("held");
3599 t.hold_machine(Some("because".to_owned()));
3600 t.blocked_by = vec!["gone".to_owned()];
3601 assert!(!t.dependency_deleted("gone"));
3602 assert_eq!(t.status, TaskStatus::Held);
3603 }
3604
3605 #[test]
3606 fn a_failed_record_removal_takes_the_tombstone_back() {
3607 let (_dir, q) = queue();
3608 let mut t = task("stays");
3609 q.put(&mut t).unwrap();
3610 q.write_tombstone(&t.id).unwrap();
3611 let err = q.remove_record_with_attachments(&t.id, |_| Err(std::io::Error::other("nope")));
3612 assert!(err.is_err());
3613 assert!(!q.was_deleted(&t.id));
3616 }
3617
3618 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3619 let p = dir.join(name);
3620 std::fs::write(&p, body).unwrap();
3621 p
3622 }
3623
3624 #[test]
3625 fn an_attachment_copy_survives_deleting_its_source() {
3626 let (dir, q) = queue();
3627 let src = source_file(dir.path(), "shot.png", "pixels");
3628 let mut t = task("with a picture");
3629 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3630 q.put(&mut t).unwrap();
3631 std::fs::remove_file(&src).unwrap();
3632 assert_eq!(names, ["shot.png"]);
3633 let loaded = q.get(&t.id).unwrap();
3634 let paths = q.attachment_paths(&loaded);
3635 assert_eq!(paths.len(), 1);
3636 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3637 }
3638
3639 #[test]
3640 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3641 let (dir, q) = queue();
3642 let mut t = task("bad names");
3643 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3644 let src = source_file(dir.path(), name, "x");
3645 assert!(
3646 q.attach(&mut t, &[src]).is_err(),
3647 "`{name}` must be refused"
3648 );
3649 }
3650 assert!(!crate::ask::valid_asset_name("C:foo.png"));
3654 #[cfg(not(windows))]
3655 {
3656 let src = source_file(dir.path(), "C:foo.png", "x");
3657 assert!(q.attach(&mut t, &[src]).is_err());
3658 }
3659 let long = format!("{}.png", "a".repeat(70));
3660 let src = source_file(dir.path(), &long, "x");
3661 assert!(q.attach(&mut t, &[src]).is_err());
3662 assert!(t.attachments.is_empty());
3663 assert!(!q.attachments_dir(&t.id).exists());
3664 }
3665
3666 #[test]
3667 fn a_taken_attachment_name_is_numbered_not_overwritten() {
3668 let (dir, q) = queue();
3669 let a = source_file(dir.path(), "shot.png", "one");
3670 let sub = dir.path().join("other");
3671 std::fs::create_dir_all(&sub).unwrap();
3672 let b = source_file(&sub, "shot.png", "two");
3673 let mut t = task("collision");
3674 q.attach(&mut t, &[a]).unwrap();
3675 q.attach(&mut t, &[b]).unwrap();
3676 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3677 let paths = q.attachment_paths(&t);
3678 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3679 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3680 }
3681
3682 #[test]
3683 fn a_renumbered_name_stays_inside_the_length_bound() {
3684 let (dir, q) = queue();
3685 let name = format!("{}.png", "a".repeat(60));
3686 assert_eq!(name.len(), 64);
3687 let a = source_file(dir.path(), &name, "one");
3688 let sub = dir.path().join("other");
3689 std::fs::create_dir_all(&sub).unwrap();
3690 let b = source_file(&sub, &name, "two");
3691 let mut t = task("long");
3692 q.attach(&mut t, &[a, b]).unwrap();
3693 assert_eq!(t.attachments.len(), 2);
3694 assert!(
3695 t.attachments
3696 .iter()
3697 .all(|n| crate::ask::valid_asset_name(n))
3698 );
3699 assert!(t.attachments[1].ends_with("-2.png"));
3700 }
3701
3702 #[test]
3703 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3704 let (dir, q) = queue();
3705 let good = source_file(dir.path(), "good.png", "ok");
3706 let mut t = task("partial");
3707 q.attach(&mut t, &[good]).unwrap();
3708 let more = source_file(dir.path(), "more.png", "ok");
3709 let missing = dir.path().join("missing.png");
3710 assert!(q.attach(&mut t, &[more, missing]).is_err());
3711 assert_eq!(t.attachments, ["good.png"]);
3712 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3713 .unwrap()
3714 .flatten()
3715 .collect();
3716 assert_eq!(on_disk.len(), 1);
3717 }
3718
3719 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3721 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3722 std::fs::create_dir_all(&tmp).unwrap();
3723 tmp
3724 }
3725
3726 #[test]
3727 fn a_failed_put_leaves_no_new_attachment_directory() {
3728 let (dir, q) = queue();
3729 let mut t = task("fresh");
3730 let tmp = block_put(&q, &t);
3731 let src = source_file(dir.path(), "shot.png", "x");
3732 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3733 assert!(t.attachments.is_empty());
3734 assert!(!q.attachments_dir(&t.id).exists());
3735 assert!(!q.path_of(&t.id).exists());
3736 std::fs::remove_dir(tmp).unwrap();
3737 }
3738
3739 #[test]
3740 fn a_failed_put_removes_only_the_copy_it_just_made() {
3741 let (dir, q) = queue();
3742 let mut t = task("edited");
3743 let first = source_file(dir.path(), "first.png", "1");
3744 q.attach_and_put(&mut t, &[first]).unwrap();
3745 block_put(&q, &t);
3746 let second = source_file(dir.path(), "second.png", "2");
3747 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3748 assert_eq!(t.attachments, ["first.png"]);
3749 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3750 .unwrap()
3751 .flatten()
3752 .map(|e| e.file_name().to_string_lossy().into_owned())
3753 .collect();
3754 assert_eq!(on_disk, ["first.png"]);
3755 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3756 }
3757
3758 #[test]
3759 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3760 let (dir, q) = queue();
3761 let questions = Questions::at(dir.path().join("questions"));
3762 let gone = task("gone");
3763 let mut other = task("other");
3764 let mut live = task("live");
3765 q.put(&mut other).unwrap();
3766 q.put(&mut live).unwrap();
3767 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3770 std::fs::create_dir_all(&orphan).unwrap();
3771 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3772 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3775 std::fs::create_dir_all(&busy).unwrap();
3776
3777 q.remove(&other.id, false, &questions).unwrap();
3778 assert!(!orphan.exists(), "an orphan is swept");
3779 assert!(busy.exists(), "a removal in progress is left alone");
3780 }
3781
3782 #[test]
3783 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3784 let (dir, q) = queue();
3785 let questions = Questions::at(dir.path().join("questions"));
3786 let mut t = task("stuck");
3787 let src = source_file(dir.path(), "shot.png", "x");
3788 q.attach_and_put(&mut t, &[src]).unwrap();
3789 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3792 std::fs::create_dir_all(&aside).unwrap();
3793 std::fs::write(aside.join("old.png"), "o").unwrap();
3794 assert!(q.remove(&t.id, false, &questions).is_err());
3795 assert!(q.path_of(&t.id).exists());
3796 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3797 }
3798
3799 #[test]
3800 fn a_failed_record_removal_puts_the_attachments_back() {
3801 let (dir, q) = queue();
3802 let mut t = task("rollback");
3803 let src = source_file(dir.path(), "shot.png", "x");
3804 q.attach_and_put(&mut t, &[src]).unwrap();
3805 let err = q
3806 .remove_record_with_attachments(&t.id, |_| {
3807 Err(std::io::Error::other("injected failure"))
3808 })
3809 .unwrap_err();
3810 assert!(format!("{err:#}").contains("injected failure"));
3811 assert!(q.path_of(&t.id).exists());
3812 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3813 assert!(
3814 !q.root
3815 .join(format!("{}.attachments.removing", t.id))
3816 .exists()
3817 );
3818 }
3819
3820 #[test]
3821 fn editing_a_task_keeps_its_attachments() {
3822 let (dir, q) = queue();
3823 let src = source_file(dir.path(), "shot.png", "x");
3824 let mut t = task("editable");
3825 q.attach(&mut t, &[src]).unwrap();
3826 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3827 q.put(&mut t).unwrap();
3828 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3829 }
3830
3831 #[test]
3832 fn removing_a_task_deletes_its_attachments() {
3833 let (dir, q) = queue();
3834 let questions = Questions::at(dir.path().join("questions"));
3835 let src = source_file(dir.path(), "shot.png", "x");
3836 let mut t = task("doomed");
3837 q.attach(&mut t, &[src]).unwrap();
3838 q.put(&mut t).unwrap();
3839 assert!(q.attachments_dir(&t.id).is_dir());
3840 q.remove(&t.id, false, &questions).unwrap();
3841 assert!(!q.attachments_dir(&t.id).exists());
3842 assert!(q.list().is_empty());
3843 }
3844
3845 #[test]
3846 fn a_task_written_before_attachments_still_reads() {
3847 let (_dir, q) = queue();
3848 let mut t = task("old");
3849 q.put(&mut t).unwrap();
3850 let path = q.path_of(&t.id);
3851 let mut v: serde_json::Value =
3852 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3853 v.as_object_mut().unwrap().remove("attachments");
3854 std::fs::write(&path, v.to_string()).unwrap();
3855 assert!(q.get(&t.id).unwrap().attachments.is_empty());
3856 }
3857
3858 #[test]
3859 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3860 let q = Queue::at(PathBuf::from("relative-queue"));
3861 let mut t = task("rel");
3862 t.attachments.push("shot.png".to_owned());
3863 let paths = q.attachment_paths(&t);
3864 assert!(paths[0].is_absolute(), "{}", paths[0].display());
3865 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3866 }
3867
3868 #[test]
3869 fn link_run_adds_a_run_once_and_touches_nothing_else() {
3870 let dir = tempfile::tempdir().unwrap();
3871 let queue = Queue::at(dir.path().join("queue"));
3872 let mut t = Task::new(
3873 "t".to_owned(),
3874 "do it".to_owned(),
3875 PathBuf::from("."),
3876 Source::Human,
3877 );
3878 queue.put(&mut t).unwrap();
3879 let before = queue.get(&t.id).unwrap();
3880
3881 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3882 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3883 assert_eq!(linked.status, before.status);
3884 assert_eq!(linked.attempts, before.attempts);
3885 assert_eq!(linked.interrupt, before.interrupt);
3886
3887 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3888 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3889 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3890 assert!(queue.link_run("no-such-task", "r").is_err());
3891 }
3892
3893 #[test]
3894 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
3895 let dir = tempfile::tempdir().unwrap();
3896 let queue = Queue::at(dir.path().join("queue"));
3897 let mut t = Task::new(
3898 "t".to_owned(),
3899 "do it".to_owned(),
3900 PathBuf::from("."),
3901 Source::Human,
3902 );
3903 queue.put(&mut t).unwrap();
3904 let mut snapshot = queue.get(&t.id).unwrap();
3906 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3907
3908 snapshot.start("20260930-000000-aaaa".to_owned());
3909 queue.put(&mut snapshot).unwrap();
3910
3911 let stored = queue.get(&t.id).unwrap();
3912 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3913 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3914 assert_eq!(stored.attempts, 1);
3915 }
3916
3917 #[test]
3918 fn concurrent_links_and_daemon_saves_lose_nothing() {
3919 let dir = tempfile::tempdir().unwrap();
3920 let queue = Queue::at(dir.path().join("queue"));
3921 let mut t = Task::new(
3922 "t".to_owned(),
3923 "do it".to_owned(),
3924 PathBuf::from("."),
3925 Source::Human,
3926 );
3927 queue.put(&mut t).unwrap();
3928 let id = t.id.clone();
3929
3930 let linkers: Vec<_> = (0..4)
3931 .map(|n| {
3932 let (queue, id) = (queue.clone(), id.clone());
3933 std::thread::spawn(move || {
3934 for k in 0..10 {
3935 queue
3936 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3937 .unwrap();
3938 }
3939 })
3940 })
3941 .collect();
3942 let mut mine = queue.get(&id).unwrap();
3945 for k in 0..10 {
3946 mine.start(format!("20260930-000009-d{k:03}"));
3947 queue.put(&mut mine).unwrap();
3948 }
3949 for l in linkers {
3950 l.join().unwrap();
3951 }
3952
3953 let stored = queue.get(&id).unwrap();
3954 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
3955 assert_eq!(
3956 stored.attempts, 10,
3957 "linking never rewinds the daemon's work"
3958 );
3959 }
3960
3961 fn age_lock(path: &Path) {
3962 let f = std::fs::OpenOptions::new().write(true).open(path).unwrap();
3963 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
3964 .unwrap();
3965 }
3966
3967 #[test]
3968 fn concurrent_stale_takeover_yields_one_holder() {
3969 use std::sync::atomic::{AtomicUsize, Ordering};
3970 use std::sync::{Arc, Barrier};
3971 for _ in 0..5 {
3972 let dir = tempfile::tempdir().unwrap();
3973 let q = Queue::at(dir.path().to_path_buf());
3974 let lock = dir.path().join("t.write-lock");
3975 std::fs::write(&lock, "dead-0000").unwrap();
3976 age_lock(&lock);
3977 let n = 6;
3978 let barrier = Arc::new(Barrier::new(n));
3979 let (now, max) = (Arc::new(AtomicUsize::new(0)), Arc::new(AtomicUsize::new(0)));
3980 let handles: Vec<_> = (0..n)
3981 .map(|_| {
3982 let (q, b, now, max) = (q.clone(), barrier.clone(), now.clone(), max.clone());
3983 std::thread::spawn(move || {
3984 b.wait();
3985 let g = q.lock_task("t").unwrap();
3986 let held = now.fetch_add(1, Ordering::SeqCst) + 1;
3987 max.fetch_max(held, Ordering::SeqCst);
3988 std::thread::sleep(std::time::Duration::from_millis(20));
3989 now.fetch_sub(1, Ordering::SeqCst);
3990 drop(g);
3991 })
3992 })
3993 .collect();
3994 for h in handles {
3995 h.join().unwrap();
3996 }
3997 assert_eq!(max.load(Ordering::SeqCst), 1);
3998 assert!(!lock.exists());
3999 }
4000 }
4001
4002 #[test]
4003 fn dropping_a_stolen_lock_leaves_the_new_holders_lock() {
4004 let dir = tempfile::tempdir().unwrap();
4005 let q = Queue::at(dir.path().to_path_buf());
4006 let lock = dir.path().join("t.write-lock");
4007 let a = q.lock_task("t").unwrap();
4008 age_lock(&lock);
4009 let b = q.lock_task("t").unwrap();
4010 assert_ne!(a.token, b.token);
4011 drop(a);
4012 assert_eq!(std::fs::read_to_string(&lock).unwrap(), b.token);
4013 drop(b);
4014 assert!(!lock.exists());
4015 }
4016
4017 #[test]
4018 fn break_stale_leaves_a_lock_that_replaced_the_one_judged() {
4019 let dir = tempfile::tempdir().unwrap();
4020 let lock = dir.path().join("t.write-lock");
4021 std::fs::write(&lock, "new-token").unwrap();
4022 assert!(!break_stale(&lock, "old-token"));
4023 assert_eq!(std::fs::read_to_string(&lock).unwrap(), "new-token");
4024 assert!(!break_stale(&lock, "new-token"));
4026 assert!(lock.exists());
4027 age_lock(&lock);
4028 assert!(break_stale(&lock, "new-token"));
4029 assert!(!lock.exists());
4030 }
4031
4032 #[test]
4033 fn a_stale_break_marker_is_recovered_and_a_fresh_one_is_respected() {
4034 let dir = tempfile::tempdir().unwrap();
4035 let lock = dir.path().join("t.write-lock");
4036 std::fs::write(&lock, "dead-1").unwrap();
4037 age_lock(&lock);
4038 let marker = break_marker(&lock, "dead-1");
4039 std::fs::write(&marker, "crashed-remover").unwrap();
4040 assert!(!break_stale(&lock, "dead-1"));
4042 assert!(lock.exists());
4043 age_lock(&marker);
4045 assert!(break_stale(&lock, "dead-1"));
4046 assert!(!lock.exists());
4047 assert!(!marker.exists());
4048 }
4049}