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 = 13;
115
116#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
118#[serde(rename_all = "lowercase")]
119pub enum HoldSource {
120 Manual,
122 Machine,
124}
125
126impl HoldSource {
127 pub fn label(self) -> &'static str {
129 match self {
130 Self::Manual => "manual",
131 Self::Machine => "machine",
132 }
133 }
134}
135
136#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
139#[serde(tag = "kind", rename_all = "lowercase")]
140pub enum Source {
141 Human,
143 Agent {
146 run: String,
148 node: String,
150 },
151 Issue {
153 number: u64,
155 repo: String,
157 },
158}
159
160pub const CHAT_NODE: &str = "chat";
163
164impl Source {
165 pub fn label(&self) -> String {
167 match self {
168 Self::Human => "human".to_owned(),
169 Self::Agent { run, node } => format!("{node}@{}", short(run)),
170 Self::Issue { number, .. } => format!("issue #{number}"),
171 }
172 }
173}
174
175#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
177#[serde(rename_all = "lowercase")]
178pub enum TaskStatus {
179 Queued,
181 Running,
183 Done,
185 Failed,
187 Held,
189 Blocked,
193 Parked,
198}
199
200impl TaskStatus {
201 pub fn runnable(self) -> bool {
203 matches!(self, Self::Queued | Self::Failed | Self::Parked)
204 }
205
206 pub fn as_str(self) -> &'static str {
208 match self {
209 Self::Queued => "queued",
210 Self::Running => "running",
211 Self::Done => "done",
212 Self::Failed => "failed",
213 Self::Held => "held",
214 Self::Blocked => "blocked",
215 Self::Parked => "parked",
216 }
217 }
218}
219
220#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
223pub struct TaskCounts {
224 pub queued: usize,
226 pub running: usize,
228 pub done: usize,
230 pub failed: usize,
232 pub held: usize,
234 pub blocked: usize,
236 pub parked: usize,
238}
239
240impl TaskCounts {
241 pub fn of(tasks: &[Task]) -> Self {
245 let mut counts = Self::default();
246 for t in tasks {
247 match t.status {
248 TaskStatus::Queued => counts.queued += 1,
249 TaskStatus::Running => counts.running += 1,
250 TaskStatus::Done => counts.done += 1,
251 TaskStatus::Failed => counts.failed += 1,
252 TaskStatus::Held => counts.held += 1,
253 TaskStatus::Blocked => counts.blocked += 1,
254 TaskStatus::Parked => counts.parked += 1,
255 }
256 }
257 counts
258 }
259}
260
261#[derive(Debug, Clone, Serialize, Deserialize)]
263#[serde(deny_unknown_fields)]
264pub struct Task {
265 pub schema: u32,
267 pub id: String,
269 pub title: String,
271 pub instruction: String,
273 pub repo: PathBuf,
275 pub source: Source,
277 #[serde(default)]
279 pub priority: i32,
280 #[serde(default)]
291 pub solo: bool,
292 pub status: TaskStatus,
294 #[serde(default)]
296 pub attempts: usize,
297 #[serde(default)]
299 pub runs: Vec<String>,
300 #[serde(default)]
302 pub last_error: Option<String>,
303 #[serde(default)]
316 pub hold_reason: Option<String>,
317 #[serde(default)]
320 pub hold_source: Option<HoldSource>,
321 #[serde(default)]
333 pub diagnostic: Option<String>,
334 #[serde(default)]
344 pub blocked_by: Vec<String>,
345 #[serde(default)]
348 pub block_reason: Option<String>,
349 #[serde(default)]
362 pub blocked_from: Option<TaskStatus>,
363 #[serde(default)]
374 pub answers: Vec<AnsweredQuestion>,
375 #[serde(default)]
381 pub triage_applied: Vec<String>,
382 #[serde(default)]
388 pub actions_applied: Vec<String>,
389 #[serde(default)]
394 pub resume_override: Option<OperatorResume>,
395 #[serde(default)]
407 pub review_branch: Option<String>,
408 #[serde(default)]
411 pub fresh_start: bool,
412 #[serde(default)]
423 pub interrupt: bool,
424 #[serde(default)]
447 pub urgent: bool,
448 #[serde(default)]
454 pub attachments: Vec<String>,
455 #[serde(default)]
459 pub followup: Option<FollowUp>,
460 #[serde(default)]
466 pub origin_chat: Option<String>,
467 #[serde(default)]
473 pub overrides: Option<RunOverrides>,
474 #[serde(default)]
481 pub review_of: Option<String>,
482 #[serde(default)]
487 pub held_at: Option<Timestamp>,
488 #[serde(default)]
492 pub park_reason: Option<String>,
493 #[serde(default)]
496 pub chat_notice: Option<ChatNotice>,
497 pub created_at: Timestamp,
499 pub updated_at: Timestamp,
501}
502
503#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
505#[serde(rename_all = "lowercase")]
506pub enum ChatNoticeOutcome {
507 Posted,
509 Skipped,
511}
512
513#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
515pub struct ChatNotice {
516 pub key: String,
518 pub at: Timestamp,
520 pub outcome: ChatNoticeOutcome,
522}
523
524#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
526pub struct RunOverrides {
527 #[serde(default)]
529 pub merge: Option<String>,
530 #[serde(default)]
532 pub candidates: Option<usize>,
533 #[serde(default)]
535 pub judges: Option<usize>,
536 #[serde(default)]
538 pub reviewers: Option<usize>,
539 #[serde(default)]
541 pub review_rounds: Option<usize>,
542 #[serde(default)]
544 pub seed: Option<u64>,
545 #[serde(default)]
547 pub config: Option<PathBuf>,
548}
549
550impl RunOverrides {
551 pub fn merge_over(&mut self, other: &RunOverrides) {
553 macro_rules! take {
554 ($($f:ident),*) => {$(
555 if other.$f.is_some() {
556 self.$f.clone_from(&other.$f);
557 }
558 )*};
559 }
560 take!(
561 merge,
562 candidates,
563 judges,
564 reviewers,
565 review_rounds,
566 seed,
567 config
568 );
569 }
570
571 pub fn apply(&self, config: &mut crate::config::Config) {
573 if let Some(n) = self.candidates {
574 config.graph.implementers = n;
575 }
576 if let Some(n) = self.judges {
577 config.graph.judges = n;
578 }
579 if let Some(n) = self.reviewers {
580 config.graph.reviewers = n;
581 }
582 if let Some(n) = self.review_rounds {
583 config.graph.review_rounds = n;
584 }
585 if let Some(m) = &self.merge
586 && let Ok(mode) = crate::daemon::merge_mode(m)
587 {
588 config.merge.mode = mode;
589 }
590 if let Some(s) = self.seed {
591 config.blind.seed = Some(s);
592 }
593 }
594}
595
596#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
598pub struct FollowUp {
599 pub run: String,
601 #[serde(default)]
603 pub origin_task: Option<String>,
604 pub pr: String,
606 pub findings: Vec<String>,
608 pub generation: u32,
611}
612
613#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
616pub struct OperatorResume {
617 pub question_id: String,
619 pub at: Timestamp,
621 #[serde(default)]
624 pub conductor_rehold: Option<String>,
625 #[serde(default)]
628 pub forced: bool,
629 #[serde(default)]
636 pub pinned_run: Option<String>,
637}
638
639#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
642pub struct AnsweredQuestion {
643 pub question: String,
645 pub answer: String,
647}
648
649impl Task {
650 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
652 let now = Timestamp::now();
653 let origin_chat = match &source {
654 Source::Agent { run, node } if node == CHAT_NODE => Some(run.clone()),
655 _ => None,
656 };
657 Self {
658 schema: SCHEMA,
659 id: new_id(),
660 title,
661 instruction,
662 repo,
663 source,
664 priority: 0,
665 solo: false,
666 status: TaskStatus::Queued,
667 attempts: 0,
668 runs: Vec::new(),
669 last_error: None,
670 hold_reason: None,
671 hold_source: None,
672 diagnostic: None,
673 blocked_by: Vec::new(),
674 block_reason: None,
675 blocked_from: None,
676 answers: Vec::new(),
677 triage_applied: Vec::new(),
678 actions_applied: Vec::new(),
679 resume_override: None,
680 review_branch: None,
681 fresh_start: false,
682 interrupt: false,
683 urgent: false,
684 attachments: Vec::new(),
685 followup: None,
686 origin_chat,
687 overrides: None,
688 review_of: None,
689 held_at: None,
690 park_reason: None,
691 chat_notice: None,
692 created_at: now,
693 updated_at: now,
694 }
695 }
696
697 pub fn chat_talk(&self) -> Option<&str> {
700 if let Some(id) = &self.origin_chat {
701 return Some(id);
702 }
703 match &self.source {
704 Source::Agent { run, node } if node == CHAT_NODE => Some(run),
705 _ => None,
706 }
707 }
708
709 pub fn filed_by_chat(&self) -> Option<&str> {
713 match &self.source {
714 Source::Agent { run, node } if node == CHAT_NODE => Some(run),
715 _ => None,
716 }
717 }
718
719 pub fn chat_notice_key(&self) -> Option<String> {
726 match self.status {
727 TaskStatus::Done => Some(format!("done#{}", self.runs.len())),
728 TaskStatus::Held => Some(format!(
729 "held#{}",
730 self.held_at.map(|t| t.to_string()).unwrap_or_default()
731 )),
732 TaskStatus::Blocked => Some(format!(
733 "blocked#{}#{}",
734 self.blocked_by.join(","),
735 self.block_reason.as_deref().unwrap_or_default()
736 )),
737 _ => None,
738 }
739 }
740
741 pub fn chat_notice_due(&self) -> Option<String> {
745 self.filed_by_chat()?;
746 if self.schema < 13 {
747 return None;
748 }
749 let key = self.chat_notice_key()?;
750 match &self.chat_notice {
751 Some(n) if n.key == key => None,
752 _ => Some(key),
753 }
754 }
755
756 pub fn short(&self) -> &str {
758 short(&self.id)
759 }
760
761 pub fn mark_triage_applied(&mut self, question_id: &str) {
763 if !self.triage_applied(question_id) {
764 self.triage_applied.push(question_id.to_owned());
765 }
766 }
767
768 pub fn action_applied(&self, question_id: &str) -> bool {
771 self.actions_applied.iter().any(|id| id == question_id)
772 }
773
774 pub fn mark_action_applied(&mut self, question_id: &str) {
776 if !self.action_applied(question_id) {
777 self.actions_applied.push(question_id.to_owned());
778 }
779 }
780
781 pub fn triage_applied(&self, question_id: &str) -> bool {
783 self.triage_applied.iter().any(|id| id == question_id)
784 }
785
786 pub fn start(&mut self, run: String) {
797 self.stamp_schema();
798 self.status = TaskStatus::Running;
799 self.attempts += 1;
800 self.runs.push(run);
801 self.last_error = None;
802 self.park_reason = None;
803 self.fresh_start = false;
804 self.interrupt = false;
805 self.resume_override = None;
807 }
808
809 pub fn link_run(&mut self, run: &str) -> bool {
814 if self.runs.iter().any(|r| r == run) {
815 return false;
816 }
817 self.runs.push(run.to_owned());
818 true
819 }
820
821 pub fn succeed(&mut self) {
831 self.stamp_schema();
832 self.status = TaskStatus::Done;
833 self.resume_override = None;
834 self.last_error = None;
835 self.park_reason = None;
836 self.hold_reason = None;
837 self.hold_source = None;
838 self.diagnostic = None;
839 self.blocked_by.clear();
840 self.block_reason = None;
841 self.blocked_from = None;
842 }
843
844 pub fn already_landed(&mut self, note: impl Into<String>) {
852 self.succeed();
853 self.attempts = self.attempts.saturating_sub(1);
854 self.last_error = Some(note.into());
855 }
856
857 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
887 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
888 return &[];
889 }
890 &self.runs[..self.runs.len() - 1]
891 }
892
893 pub fn successor_of(&self, run: &str) -> Option<&String> {
900 let pos = self.runs.iter().position(|r| r == run)?;
901 self.runs.get(pos + 1)
902 }
903
904 pub fn earlier_attempts(&self) -> &[String] {
911 &self.runs
912 }
913
914 fn stamp_schema(&mut self) {
921 self.schema = self.schema.max(SCHEMA);
922 }
923
924 fn note_held(&mut self) {
925 self.stamp_schema();
926 if self.status != TaskStatus::Held || self.held_at.is_none() {
927 self.held_at = Some(Timestamp::now());
928 }
929 }
930
931 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
947 let why = why.into();
948 self.diagnostic = None;
949 self.park_reason = None;
950 self.status = if self.attempts >= max_attempts {
951 self.note_held();
952 self.hold_source = Some(HoldSource::Machine);
953 self.hold_reason = Some(why.clone());
954 TaskStatus::Held
955 } else {
956 TaskStatus::Failed
957 };
958 self.last_error = Some(why);
959 }
960
961 pub fn stall(&mut self, why: impl Into<String>) {
971 self.last_error = Some(why.into());
972 self.diagnostic = None;
973 self.attempts = self.attempts.saturating_sub(1);
974 self.status = TaskStatus::Failed;
975 }
976
977 pub fn park(&mut self, why: impl Into<String>) {
983 self.diagnostic = None;
984 self.attempts = self.attempts.saturating_sub(1);
985 self.status = TaskStatus::Parked;
986 self.park_reason = Some(why.into());
987 self.schema = self.schema.max(SCHEMA);
990 }
991
992 pub fn operator_held(&self) -> bool {
998 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
999 }
1000
1001 pub fn hold_manual(&mut self, reason: Option<String>) {
1011 self.park_reason = None;
1012 self.note_held();
1013 self.status = TaskStatus::Held;
1014 if reason.is_some() {
1015 self.hold_reason = reason;
1016 }
1017 self.hold_source = Some(HoldSource::Manual);
1018 self.blocked_by.clear();
1019 self.block_reason = None;
1020 self.blocked_from = None;
1021 }
1022
1023 pub fn hold_machine(&mut self, reason: Option<String>) {
1028 self.park_reason = None;
1029 self.note_held();
1030 self.status = TaskStatus::Held;
1031 if reason.is_some() {
1032 self.hold_reason = reason;
1033 }
1034 self.hold_source = Some(HoldSource::Machine);
1035 self.blocked_by.clear();
1036 self.block_reason = None;
1037 self.blocked_from = None;
1038 }
1039
1040 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
1050 self.stamp_schema();
1051 if self.status != TaskStatus::Blocked {
1052 self.blocked_from = Some(self.status);
1053 }
1054 self.status = TaskStatus::Blocked;
1055 self.blocked_by = blocked_by;
1056 self.block_reason = reason;
1057 }
1058
1059 pub fn unblock(&mut self, resolved_id: &str) {
1082 if self.status != TaskStatus::Blocked {
1083 return;
1084 }
1085 self.blocked_by.retain(|id| id != resolved_id);
1086 self.restore_if_unblocked();
1087 }
1088
1089 pub fn dependency_deleted(&mut self, deleted_id: &str) -> bool {
1100 if self.status != TaskStatus::Blocked || !self.blocked_by.iter().any(|b| b == deleted_id) {
1101 return false;
1102 }
1103 self.blocked_by.retain(|id| id != deleted_id);
1104 self.restore_if_unblocked();
1105 true
1106 }
1107
1108 fn restore_if_unblocked(&mut self) {
1112 if self.blocked_by.is_empty() {
1113 self.status = match self.blocked_from {
1114 Some(TaskStatus::Running) => TaskStatus::Queued,
1115 Some(other) => other,
1116 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
1117 TaskStatus::Held
1118 }
1119 None => TaskStatus::Queued,
1120 };
1121 self.block_reason = None;
1122 self.blocked_from = None;
1123 }
1124 }
1125
1126 pub fn record_answer(&mut self, question: String, answer: String) {
1131 self.answers.push(AnsweredQuestion { question, answer });
1132 }
1133
1134 pub fn request_review(&mut self, branch: String) {
1138 self.release();
1139 self.review_branch = Some(branch);
1140 }
1141
1142 pub fn requeue(&mut self) {
1145 self.release();
1146 self.review_branch = None;
1148 self.fresh_start = true;
1149 }
1150
1151 pub fn hold_for_handover(&mut self, branch: Option<String>, reason: String) {
1160 if branch.is_some() {
1161 self.review_branch = branch;
1162 }
1163 self.hold_machine(Some(reason));
1164 }
1165
1166 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
1175 if self.status == TaskStatus::Running {
1176 bail!(
1177 "task {} is running; its priority cannot be changed until \
1178 this attempt finishes",
1179 self.short()
1180 );
1181 }
1182 self.priority = priority;
1183 Ok(())
1184 }
1185
1186 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
1201 if interrupt && !self.status.runnable() {
1202 bail!(
1203 "task {} is {}; only a queued or failed task can be marked \
1204 to interrupt",
1205 self.short(),
1206 self.status.as_str()
1207 );
1208 }
1209 self.interrupt = interrupt;
1210 Ok(())
1211 }
1212
1213 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
1225 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
1226 bail!(
1227 "task {} is {}; only a queued or held task's instruction can \
1228 be edited",
1229 self.short(),
1230 self.status.as_str()
1231 );
1232 }
1233 self.title = title;
1234 self.instruction = instruction;
1235 Ok(())
1236 }
1237
1238 pub fn handed_off(&mut self, why: impl Into<String>) {
1258 let why = why.into();
1259 self.diagnostic = None;
1260 self.note_held();
1261 self.status = TaskStatus::Held;
1262 self.hold_source = Some(HoldSource::Machine);
1263 self.hold_reason = Some(why.clone());
1264 self.last_error = Some(why);
1265 }
1266
1267 pub fn release(&mut self) {
1271 let refused_handover = self.status == TaskStatus::Held
1272 && self.hold_source == Some(HoldSource::Machine)
1273 && self.review_branch.is_some();
1274 self.status = TaskStatus::Queued;
1275 self.held_at = None;
1276 self.attempts = 0;
1277 self.last_error = None;
1278 self.park_reason = None;
1279 self.hold_reason = None;
1282 self.hold_source = None;
1283 self.diagnostic = None;
1284 self.blocked_by.clear();
1289 self.block_reason = None;
1290 self.blocked_from = None;
1291 if !refused_handover {
1295 self.review_branch = None;
1296 }
1297 self.fresh_start = false;
1298 }
1299}
1300
1301fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
1305 use std::io::ErrorKind;
1306 let (stem, ext) = match name.rfind('.') {
1307 Some(i) if i > 0 => (&name[..i], &name[i..]),
1308 _ => (name, ""),
1309 };
1310 for n in 1u32.. {
1311 let candidate = if n == 1 {
1312 name.to_owned()
1313 } else {
1314 let suffix = format!("-{n}");
1315 let room = 64usize.saturating_sub(suffix.len() + ext.len());
1316 let stem: String = stem.chars().take(room).collect();
1317 format!("{stem}{suffix}{ext}")
1318 };
1319 if !crate::ask::valid_asset_name(&candidate) {
1320 bail!("no valid attachment name is left for `{name}`");
1321 }
1322 let path = dir.join(&candidate);
1323 match std::fs::OpenOptions::new()
1324 .write(true)
1325 .create_new(true)
1326 .open(&path)
1327 {
1328 Ok(mut out) => {
1329 let copied = std::fs::File::open(src)
1330 .and_then(|mut input| std::io::copy(&mut input, &mut out));
1331 if let Err(e) = copied {
1332 drop(out);
1333 let _ = std::fs::remove_file(&path);
1334 return Err(e).with_context(|| format!("copy {}", src.display()));
1335 }
1336 return Ok((candidate, path));
1337 }
1338 Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
1339 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
1340 }
1341 }
1342 unreachable!("the counter never runs out")
1343}
1344
1345const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
1348
1349struct TaskLock {
1352 path: PathBuf,
1353 token: String,
1354}
1355
1356fn owner_token() -> String {
1359 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1360 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1361 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy() ^ n.rotate_left(32));
1362 format!("{}-{:016x}", std::process::id(), r.next_u64())
1363}
1364
1365fn break_marker(path: &Path, token: &str) -> PathBuf {
1368 let mut h: u64 = 0xcbf29ce484222325;
1369 for b in token.bytes() {
1370 h ^= u64::from(b);
1371 h = h.wrapping_mul(0x100000001b3);
1372 }
1373 let mut name = path.as_os_str().to_owned();
1374 name.push(format!(".break-{h:016x}"));
1375 PathBuf::from(name)
1376}
1377
1378fn older_than_stale(path: &Path) -> bool {
1379 std::fs::metadata(path)
1380 .and_then(|m| m.modified())
1381 .ok()
1382 .and_then(|t| t.elapsed().ok())
1383 .is_some_and(|age| age > TASK_LOCK_STALE)
1384}
1385
1386const MAX_MARKER_DEPTH: u8 = 3;
1390
1391fn take_marker(path: &Path, token: &str, depth: u8) -> Option<(PathBuf, String)> {
1397 use std::io::Write;
1398 let marker = break_marker(path, token);
1399 for _ in 0..2 {
1400 match std::fs::OpenOptions::new()
1401 .write(true)
1402 .create_new(true)
1403 .open(&marker)
1404 {
1405 Ok(mut f) => {
1406 let mine = owner_token();
1407 if f.write_all(mine.as_bytes()).is_err() {
1408 drop(f);
1409 let _ = std::fs::remove_file(&marker);
1410 return None;
1411 }
1412 return Some((marker, mine));
1413 }
1414 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1415 let judged = std::fs::read_to_string(&marker).ok();
1416 match judged.filter(|_| older_than_stale(&marker)) {
1417 Some(judged) if depth < MAX_MARKER_DEPTH => {
1418 remove_lock_if(&marker, &judged, true, depth + 1)?;
1419 }
1420 _ => return None,
1421 }
1422 }
1423 Err(_) => return None,
1424 }
1425 }
1426 None
1427}
1428
1429fn remove_lock_if(path: &Path, token: &str, require_stale: bool, depth: u8) -> Option<bool> {
1434 let (marker, mine) = take_marker(path, token, depth)?;
1435 let still = std::fs::read_to_string(path).is_ok_and(|c| c == token)
1436 && (!require_stale || older_than_stale(path));
1437 if still {
1438 let _ = std::fs::remove_file(path);
1439 }
1440 if depth >= MAX_MARKER_DEPTH {
1442 if std::fs::read_to_string(&marker).is_ok_and(|c| c == mine) {
1443 let _ = std::fs::remove_file(&marker);
1444 }
1445 } else {
1446 let _ = remove_lock_if(&marker, &mine, false, depth + 1);
1447 }
1448 Some(still)
1449}
1450
1451fn break_stale(path: &Path, judged: &str) -> bool {
1454 remove_lock_if(path, judged, true, 0).unwrap_or(false)
1455}
1456
1457impl Drop for TaskLock {
1458 fn drop(&mut self) {
1459 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
1460 loop {
1461 match remove_lock_if(&self.path, &self.token, false, 0) {
1462 Some(_) => return,
1463 None if std::time::Instant::now() > deadline => return,
1464 None => std::thread::sleep(std::time::Duration::from_millis(5)),
1465 }
1466 }
1467 }
1468}
1469
1470#[derive(Debug, Clone)]
1472pub struct Queue {
1473 root: PathBuf,
1474}
1475
1476impl Queue {
1477 pub fn open() -> Self {
1479 Self::at(crate::run::home().join("queue"))
1480 }
1481
1482 pub fn at(root: PathBuf) -> Self {
1485 Self { root }
1486 }
1487
1488 pub fn root(&self) -> &Path {
1490 &self.root
1491 }
1492
1493 pub fn path_of(&self, id: &str) -> PathBuf {
1495 self.root.join(format!("{id}.json"))
1496 }
1497
1498 pub fn attachments_dir(&self, id: &str) -> PathBuf {
1501 self.root.join(format!("{id}.attachments"))
1502 }
1503
1504 pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1515 let mut wanted = Vec::new();
1516 for src in sources {
1517 let name = src
1518 .file_name()
1519 .and_then(|n| n.to_str())
1520 .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1521 if !crate::ask::valid_asset_name(name) {
1522 bail!(
1523 "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1524 with no `..`; rename the file and try again"
1525 );
1526 }
1527 if !src.is_file() {
1528 bail!("attachment `{}` is not a file", src.display());
1529 }
1530 wanted.push((src, name));
1531 }
1532 let dir = self.attachments_dir(&task.id);
1533 let existed = dir.is_dir();
1534 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1535 let mut created: Vec<PathBuf> = Vec::new();
1536 let mut names = Vec::new();
1537 let mut copy_all = || -> Result<()> {
1538 for (src, name) in &wanted {
1539 let (stored, path) = copy_new(&dir, src, name)?;
1540 created.push(path);
1541 names.push(stored);
1542 }
1543 Ok(())
1544 };
1545 if let Err(e) = copy_all() {
1546 for path in &created {
1547 let _ = std::fs::remove_file(path);
1548 }
1549 if !existed {
1550 let _ = std::fs::remove_dir(&dir);
1551 }
1552 return Err(e);
1553 }
1554 task.attachments.extend(names.iter().cloned());
1555 Ok(names)
1556 }
1557
1558 pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1569 let dir = self.attachments_dir(&task.id);
1570 let existed = dir.is_dir();
1571 let before = task.attachments.len();
1572 let names = self.attach(task, sources)?;
1573 if let Err(e) = self.put(task) {
1574 for name in &names {
1575 let _ = std::fs::remove_file(dir.join(name));
1576 }
1577 if !existed {
1578 let _ = std::fs::remove_dir(&dir);
1579 }
1580 task.attachments.truncate(before);
1581 return Err(e);
1582 }
1583 Ok(names)
1584 }
1585
1586 pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1590 let dir = self.attachments_dir(&task.id);
1591 task.attachments
1592 .iter()
1593 .map(|n| {
1594 let p = dir.join(n);
1595 std::path::absolute(&p).unwrap_or(p)
1596 })
1597 .collect()
1598 }
1599
1600 pub fn put(&self, task: &mut Task) -> Result<()> {
1609 let _lock = self.lock_task(&task.id)?;
1610 self.put_unlocked(task)
1611 }
1612
1613 pub fn create_new(&self, task: &mut Task) -> Result<bool> {
1618 let _lock = self.lock_task(&task.id)?;
1619 if self.path_of(&task.id).exists() {
1620 return Ok(false);
1621 }
1622 self.put_unlocked(task)?;
1623 Ok(true)
1624 }
1625
1626 fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1628 if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1629 for run in stored.runs {
1630 if !task.runs.contains(&run) {
1631 task.runs.push(run);
1632 }
1633 }
1634 if let Some(n) = stored.chat_notice
1637 && task.chat_notice.as_ref().is_none_or(|mine| mine.at < n.at)
1638 {
1639 task.chat_notice = Some(n);
1640 }
1641 }
1642 task.updated_at = Timestamp::now();
1643 std::fs::create_dir_all(&self.root)
1644 .with_context(|| format!("create {}", self.root.display()))?;
1645 let body = serde_json::to_string_pretty(task).context("serialize task")?;
1646 let path = self.path_of(&task.id);
1647 let tmp = path.with_extension("json.tmp");
1648 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1649 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1650 if let (Some(notice), Some(home)) = (
1655 crate::notices::task_held(task),
1656 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1657 ) {
1658 crate::notices::raise_in(home, notice);
1659 }
1660 Ok(())
1661 }
1662
1663 pub fn modify(&self, id: &str, f: impl FnOnce(&mut Task) -> bool) -> Result<bool> {
1668 let id = self.resolve_id(id)?;
1669 let _lock = self.lock_task(&id)?;
1670 let mut task = self.get(&id)?;
1671 if !f(&mut task) {
1672 return Ok(false);
1673 }
1674 self.put_unlocked(&mut task)?;
1675 Ok(true)
1676 }
1677
1678 pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1683 let id = self.resolve_id(id)?;
1684 let _lock = self.lock_task(&id)?;
1687 let mut task = self.get(&id)?;
1688 if task.link_run(run) {
1689 self.put_unlocked(&mut task)?;
1690 }
1691 Ok(task)
1692 }
1693
1694 fn lock_task(&self, id: &str) -> Result<TaskLock> {
1702 std::fs::create_dir_all(&self.root)
1703 .with_context(|| format!("create {}", self.root.display()))?;
1704 let path = self.root.join(format!("{id}.write-lock"));
1705 let started = std::time::Instant::now();
1706 loop {
1707 match std::fs::OpenOptions::new()
1708 .write(true)
1709 .create_new(true)
1710 .open(&path)
1711 {
1712 Ok(mut file) => {
1713 use std::io::Write;
1714 let token = owner_token();
1715 if let Err(e) = file.write_all(token.as_bytes()) {
1716 drop(file);
1717 let _ = std::fs::remove_file(&path);
1718 return Err(e).with_context(|| format!("lock {}", path.display()));
1719 }
1720 return Ok(TaskLock { path, token });
1721 }
1722 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1723 let judged = std::fs::read_to_string(&path).ok();
1726 let broken = judged
1727 .filter(|_| older_than_stale(&path))
1728 .is_some_and(|judged| break_stale(&path, &judged));
1729 if !broken {
1733 if started.elapsed() > TASK_LOCK_STALE {
1734 bail!("could not lock task {id}");
1735 }
1736 std::thread::sleep(std::time::Duration::from_millis(15));
1737 }
1738 }
1739 Err(e)
1743 if e.kind() == std::io::ErrorKind::PermissionDenied
1744 && started.elapsed() <= TASK_LOCK_STALE =>
1745 {
1746 std::thread::sleep(std::time::Duration::from_millis(15));
1747 }
1748 Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1749 }
1750 }
1751 }
1752
1753 pub fn get(&self, id: &str) -> Result<Task> {
1755 let resolved = self.resolve_id(id)?;
1756 read_path(&self.path_of(&resolved))
1757 }
1758
1759 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1791 let _ = questions;
1792 let resolved = self.resolve_id(id)?;
1793 if in_flight {
1794 bail!("task {resolved} is being run by a live daemon right now");
1795 }
1796 self.write_tombstone(&resolved)?;
1797 if let Err(e) = self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))
1798 {
1799 if self.path_of(&resolved).exists() {
1803 let _ = std::fs::remove_file(self.tombstone_path(&resolved));
1804 }
1805 return Err(e);
1806 }
1807 let (released, still_blocked) = self.release_dependents_of(&resolved);
1808 Ok(Removal {
1809 id: resolved,
1810 released,
1811 still_blocked,
1812 })
1813 }
1814
1815 fn tombstone_path(&self, id: &str) -> PathBuf {
1818 self.root.join(format!("{id}.removed"))
1819 }
1820
1821 fn write_tombstone(&self, id: &str) -> Result<()> {
1822 let path = self.tombstone_path(id);
1823 let tmp = self.root.join(format!("{id}.removed.tmp"));
1824 std::fs::write(&tmp, b"")
1825 .and_then(|()| std::fs::rename(&tmp, &path))
1826 .with_context(|| format!("write {}", path.display()))
1827 }
1828
1829 fn was_deleted(&self, id: &str) -> bool {
1833 self.tombstone_path(id).is_file() && !self.path_of(id).exists()
1834 }
1835
1836 pub fn apply_deleted_blockers(&self, task: &mut Task) -> Vec<String> {
1840 let deleted = deleted_blockers(self, &task.blocked_by);
1841 deleted
1842 .into_iter()
1843 .filter(|id| task.dependency_deleted(id))
1844 .collect()
1845 }
1846
1847 pub fn note_dependency_deleted(&self, dependent: &Task, deleted: &str) {
1851 let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty()) else {
1852 return;
1853 };
1854 let ja = crate::lang::is_japanese(&crate::lang::of_repo(&dependent.repo));
1855 let waiting = dependent.status == TaskStatus::Blocked;
1856 let message = match (ja, waiting) {
1857 (true, false) => format!(
1858 "タスク {} は、待っていた {} が削除されたため待機を解除し、元の状態に戻しました",
1859 dependent.short(),
1860 short(deleted)
1861 ),
1862 (true, true) => format!(
1863 "タスク {} は、待っていた {} が削除されたため、残りの依存を待っています",
1864 dependent.short(),
1865 short(deleted)
1866 ),
1867 (false, false) => format!(
1868 "Task {} stopped waiting on {} because it was deleted, and returned to its previous state",
1869 dependent.short(),
1870 short(deleted)
1871 ),
1872 (false, true) => format!(
1873 "Task {} stopped waiting on {} because it was deleted, and is still waiting on its other dependencies",
1874 dependent.short(),
1875 short(deleted)
1876 ),
1877 };
1878 crate::notices::raise_in(
1879 home,
1880 crate::notices::Notice::info(
1881 &format!("unblocked:{}:{}", dependent.id, deleted),
1882 message,
1883 )
1884 .link(crate::notices::Link::Task {
1885 id: dependent.id.clone(),
1886 }),
1887 );
1888 }
1889
1890 fn release_dependents_of(&self, dependency: &str) -> (Vec<String>, Vec<String>) {
1895 let (mut released, mut still_blocked) = (Vec::new(), Vec::new());
1896 for listed in self.list() {
1897 if listed.status != TaskStatus::Blocked
1898 || !listed.blocked_by.iter().any(|b| b == dependency)
1899 {
1900 continue;
1901 }
1902 let Ok(_claim) = self.claim(&listed.id) else {
1903 continue;
1904 };
1905 let Ok(mut task) = self.get(&listed.id) else {
1906 continue;
1907 };
1908 if !task.dependency_deleted(dependency) {
1909 continue;
1910 }
1911 if self.put(&mut task).is_err() {
1912 continue;
1913 }
1914 self.note_dependency_deleted(&task, dependency);
1915 if task.status == TaskStatus::Blocked {
1916 still_blocked.push(task.id.clone());
1917 } else {
1918 released.push(task.id.clone());
1919 }
1920 }
1921 (released, still_blocked)
1922 }
1923
1924 fn remove_record_with_attachments(
1928 &self,
1929 resolved: &str,
1930 remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1931 ) -> Result<()> {
1932 self.sweep_removed_attachments();
1939 let attachments = self.attachments_dir(resolved);
1940 let aside = self.root.join(format!("{resolved}.attachments.removing"));
1941 let moved = match std::fs::rename(&attachments, &aside) {
1942 Ok(()) => true,
1943 Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1944 Err(e) => {
1945 return Err(e).with_context(|| format!("remove {}", attachments.display()));
1946 }
1947 };
1948 let path = self.path_of(resolved);
1949 if let Err(e) = remove_record(&path) {
1950 if moved {
1951 let _ = std::fs::rename(&aside, &attachments);
1952 }
1953 return Err(e).with_context(|| format!("remove {}", path.display()));
1954 }
1955 if moved {
1956 if let Err(e) = std::fs::remove_dir_all(&aside) {
1957 tracing::warn!("leftover attachments {}: {e}", aside.display());
1958 }
1959 }
1960 let lock = self.lock_path(resolved);
1961 if let Err(e) = std::fs::remove_file(&lock) {
1962 if e.kind() != std::io::ErrorKind::NotFound {
1963 return Err(e).with_context(|| format!("remove {}", lock.display()));
1964 }
1965 }
1966 Ok(())
1967 }
1968
1969 fn sweep_removed_attachments(&self) {
1975 let Ok(entries) = std::fs::read_dir(&self.root) else {
1976 return;
1977 };
1978 for entry in entries.flatten() {
1979 let name = entry.file_name();
1980 let name = name.to_string_lossy();
1981 let Some(id) = name.strip_suffix(".attachments.removing") else {
1982 continue;
1983 };
1984 if !self.path_of(id).exists() {
1987 if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1988 tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1989 }
1990 }
1991 }
1992 self.sweep_tombstones();
1993 }
1994
1995 fn sweep_tombstones(&self) {
2004 let Ok(entries) = std::fs::read_dir(&self.root) else {
2005 return;
2006 };
2007 let tasks = self.list();
2008 for entry in entries.flatten() {
2009 let name = entry.file_name();
2010 let name = name.to_string_lossy();
2011 let Some(id) = name.strip_suffix(".removed") else {
2012 continue;
2013 };
2014 if self.path_of(id).exists()
2015 || tasks.iter().any(|t| t.blocked_by.iter().any(|b| b == id))
2016 {
2017 continue;
2018 }
2019 let old = entry
2020 .metadata()
2021 .and_then(|m| m.modified())
2022 .ok()
2023 .and_then(|m| m.elapsed().ok())
2024 .is_some_and(|age| age >= TOMBSTONE_GRACE);
2025 if old {
2026 let _ = std::fs::remove_file(entry.path());
2027 }
2028 }
2029 }
2030
2031 fn lock_path(&self, id: &str) -> PathBuf {
2034 self.root.join(format!("{id}.lock"))
2035 }
2036
2037 pub fn list(&self) -> Vec<Task> {
2050 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
2051 .into_iter()
2052 .flatten()
2053 .flatten()
2054 .map(|e| e.path())
2055 .filter(|p| p.extension().is_some_and(|x| x == "json"))
2056 .filter_map(|p| read_path(&p).ok())
2057 .collect();
2058 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
2059 tasks
2060 }
2061
2062 pub fn superseded(&self) -> HashMap<String, String> {
2075 let mut by = HashMap::new();
2076 for task in self.list() {
2077 for earlier in &task.runs {
2078 if let Some(later) = task.successor_of(earlier) {
2079 by.insert(earlier.clone(), later.clone());
2080 }
2081 }
2082 }
2083 by
2084 }
2085
2086 pub fn superseded_by(&self, run: &str) -> Option<String> {
2094 for task in self.list() {
2095 if task.runs.iter().any(|r| r == run) {
2096 return task.successor_of(run).cloned();
2097 }
2098 }
2099 None
2100 }
2101
2102 pub fn latest_attempt(&self, run: &str) -> Option<String> {
2113 for task in self.list() {
2114 if task.runs.iter().any(|r| r == run) {
2115 return task.runs.last().filter(|last| **last != run).cloned();
2116 }
2117 }
2118 None
2119 }
2120
2121 pub fn next_runnable(&self) -> Option<Task> {
2126 let mut runnable: Vec<Task> = self
2127 .list()
2128 .into_iter()
2129 .filter(|t| t.status.runnable())
2130 .collect();
2131 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
2132 runnable.into_iter().next()
2133 }
2134
2135 pub fn claim(&self, id: &str) -> Result<Claim> {
2142 std::fs::create_dir_all(&self.root)
2143 .with_context(|| format!("create {}", self.root.display()))?;
2144 let path = self.lock_path(id);
2145 match std::fs::OpenOptions::new()
2146 .write(true)
2147 .create_new(true)
2148 .open(&path)
2149 {
2150 Ok(mut f) => {
2151 use std::io::Write as _;
2152 let _ = writeln!(f, "{}", std::process::id());
2154 Ok(Claim { path })
2155 }
2156 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
2157 bail!("task {id} is already claimed ({} exists)", path.display())
2158 }
2159 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
2160 }
2161 }
2162
2163 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
2165 if self.path_of(prefix).is_file() {
2166 return Ok(prefix.to_owned());
2167 }
2168 let hits: Vec<String> = self
2169 .list()
2170 .into_iter()
2171 .map(|t| t.id)
2172 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
2173 .collect();
2174 match hits.len() {
2175 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
2176 0 => bail!("no task matches `{prefix}`"),
2177 _ => bail!(
2178 "`{prefix}` matches {} tasks: {}",
2179 hits.len(),
2180 hits.join(", ")
2181 ),
2182 }
2183 }
2184
2185 pub fn revision(&self) -> u64 {
2192 self.revision_excluding(&std::collections::BTreeSet::new())
2193 }
2194
2195 pub fn revision_excluding(&self, skip: &std::collections::BTreeSet<String>) -> u64 {
2198 use std::hash::{Hash as _, Hasher as _};
2199
2200 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
2201 .into_iter()
2202 .flatten()
2203 .flatten()
2204 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
2205 .filter(|e| {
2206 let path = e.path();
2207 !path
2208 .file_stem()
2209 .is_some_and(|stem| skip.contains(stem.to_string_lossy().as_ref()))
2210 })
2211 .filter_map(|e| {
2212 let name = e.file_name().to_string_lossy().into_owned();
2213 let mtime = e
2214 .metadata()
2215 .ok()?
2216 .modified()
2217 .ok()?
2218 .duration_since(std::time::UNIX_EPOCH)
2219 .ok()?
2220 .as_millis() as u64;
2221 Some((name, mtime))
2222 })
2223 .collect();
2224
2225 if entries.is_empty() {
2226 return 0;
2227 }
2228
2229 entries.sort_unstable();
2230 let mut hasher = std::hash::DefaultHasher::new();
2231 for (name, mtime) in &entries {
2232 name.hash(&mut hasher);
2233 mtime.hash(&mut hasher);
2234 }
2235 let h = hasher.finish();
2236 if h == 0 { 1 } else { h }
2237 }
2238}
2239
2240const TOMBSTONE_GRACE: std::time::Duration = std::time::Duration::from_secs(3600);
2243
2244#[derive(Debug, Clone)]
2246pub struct Removal {
2247 pub id: String,
2249 pub released: Vec<String>,
2253 pub still_blocked: Vec<String>,
2256}
2257
2258#[derive(Debug)]
2260pub struct Claim {
2261 path: PathBuf,
2262}
2263
2264impl Drop for Claim {
2265 fn drop(&mut self) {
2266 let _ = std::fs::remove_file(&self.path);
2267 }
2268}
2269
2270pub fn title_from(instruction: &str, max: usize) -> String {
2273 let Some(line) = first_line(instruction) else {
2274 return "(empty task)".to_owned();
2275 };
2276 if line.chars().count() <= max {
2277 return line.to_owned();
2278 }
2279 let head: String = line.chars().take(max.saturating_sub(1)).collect();
2280 format!("{head}…")
2281}
2282
2283pub(crate) fn first_line(instruction: &str) -> Option<&str> {
2289 let line = instruction
2290 .lines()
2291 .map(str::trim)
2292 .find(|l| !l.is_empty())?
2293 .trim_start_matches(['#', '-', '*', '>', ' '])
2294 .trim();
2295 (!line.is_empty()).then_some(line)
2296}
2297
2298fn read_path(path: &Path) -> Result<Task> {
2299 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2300 let task: Task =
2301 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
2302 if task.schema > SCHEMA {
2308 bail!(
2309 "task {} was written by a different magi (schema {}, this build \
2310 speaks {SCHEMA})",
2311 task.id,
2312 task.schema
2313 );
2314 }
2315 Ok(task)
2316}
2317
2318pub fn missing_blockers(
2332 queue: &Queue,
2333 questions: &Questions,
2334 blocked_by: &[String],
2335) -> Vec<String> {
2336 blocked_by
2337 .iter()
2338 .filter(|id| {
2339 !queue.path_of(id).is_file()
2340 && !questions.path_of(id).is_file()
2341 && !queue.was_deleted(id)
2342 })
2343 .cloned()
2344 .collect()
2345}
2346
2347pub fn deleted_blockers(queue: &Queue, blocked_by: &[String]) -> Vec<String> {
2352 blocked_by
2353 .iter()
2354 .filter(|id| queue.was_deleted(id))
2355 .cloned()
2356 .collect()
2357}
2358
2359pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
2371 missing_blocker_hold_reason_in(blocked_by, missing, "en")
2372}
2373
2374pub fn missing_blocker_hold_reason_in(
2376 blocked_by: &[String],
2377 missing: &[String],
2378 language: &str,
2379) -> String {
2380 if crate::lang::is_japanese(language) {
2381 format!(
2382 "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
2383 blocked_by.join(", "),
2384 missing.join(", "),
2385 )
2386 } else {
2387 format!(
2388 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
2389 blocked_by.join(", "),
2390 missing.join(", "),
2391 )
2392 }
2393}
2394
2395pub fn short(id: &str) -> &str {
2397 id.split('-').next_back().unwrap_or(id)
2398}
2399
2400fn new_id() -> String {
2401 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2402 let seed = crate::rng::entropy();
2403 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2404}
2405
2406#[cfg(test)]
2407mod tests {
2408 #[test]
2409 fn merge_over_lets_the_later_choice_win_and_keeps_the_rest() {
2410 let mut base = RunOverrides {
2411 merge: Some("pr".to_owned()),
2412 candidates: Some(3),
2413 ..RunOverrides::default()
2414 };
2415 base.merge_over(&RunOverrides {
2416 merge: Some("none".to_owned()),
2417 seed: Some(7),
2418 ..RunOverrides::default()
2419 });
2420 assert_eq!(base.merge.as_deref(), Some("none"));
2421 assert_eq!(base.candidates, Some(3));
2422 assert_eq!(base.seed, Some(7));
2423 }
2424
2425 #[test]
2426 fn an_already_landed_task_is_done_with_its_attempt_refunded() {
2427 let mut t = task("relanded");
2428 t.attempts = 1;
2429 t.status = TaskStatus::Running;
2430 t.already_landed("already in main as 0e368de");
2431 assert_eq!(t.status, TaskStatus::Done);
2432 assert_eq!(t.attempts, 0);
2433 assert!(t.hold_reason.is_none());
2434 assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
2435 }
2436
2437 #[test]
2438 fn missing_blocker_reason_follows_the_language() {
2439 let b = vec!["a".to_owned()];
2440 let en = missing_blocker_hold_reason_in(&b, &b, "en");
2441 assert_eq!(en, missing_blocker_hold_reason(&b, &b));
2442 assert!(en.starts_with("blocked on a"));
2443 assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
2444 assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
2445 }
2446
2447 use super::*;
2448
2449 #[test]
2450 fn triage_applied_survives_release_and_old_records_read_as_empty() {
2451 let mut t = Task::new(
2452 "t".to_owned(),
2453 "i".to_owned(),
2454 PathBuf::from("r"),
2455 Source::Human,
2456 );
2457 t.mark_triage_applied("q1");
2458 t.mark_triage_applied("q1");
2459 t.hold_machine(Some("x".to_owned()));
2460 t.release();
2461 assert_eq!(t.triage_applied, ["q1"]);
2462 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
2463
2464 let mut v = serde_json::to_value(&t).unwrap();
2465 v.as_object_mut().unwrap().remove("triage_applied");
2466 let old: Task = serde_json::from_value(v).unwrap();
2467 assert!(old.triage_applied.is_empty());
2468 }
2469
2470 #[test]
2471 fn task_counts_of_empty_is_all_zero() {
2472 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
2473 }
2474
2475 #[test]
2476 fn task_counts_of_tallies_every_status() {
2477 let mut queued = Task::new(
2478 "q".to_owned(),
2479 "i".to_owned(),
2480 PathBuf::from("."),
2481 Source::Human,
2482 );
2483 queued.status = TaskStatus::Queued;
2484 let mut running = queued.clone();
2485 running.status = TaskStatus::Running;
2486 let mut done = queued.clone();
2487 done.status = TaskStatus::Done;
2488 let mut failed = queued.clone();
2489 failed.status = TaskStatus::Failed;
2490 let mut held = queued.clone();
2491 held.status = TaskStatus::Held;
2492 let mut blocked = queued.clone();
2493 blocked.status = TaskStatus::Blocked;
2494
2495 let mut parked = queued.clone();
2496 parked.status = TaskStatus::Parked;
2497
2498 let counts = TaskCounts::of(&[
2499 queued,
2500 running,
2501 done.clone(),
2502 done,
2503 failed,
2504 held,
2505 blocked,
2506 parked,
2507 ]);
2508 assert_eq!(
2509 counts,
2510 TaskCounts {
2511 queued: 1,
2512 running: 1,
2513 done: 2,
2514 failed: 1,
2515 held: 1,
2516 blocked: 1,
2517 parked: 1,
2518 }
2519 );
2520 }
2521
2522 fn queue() -> (tempfile::TempDir, Queue) {
2525 let dir = tempfile::tempdir().unwrap();
2526 let q = Queue::at(dir.path().join("queue"));
2527 (dir, q)
2528 }
2529
2530 #[test]
2531 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
2532 let dir = tempfile::tempdir().unwrap();
2533 let q = Queue::at(dir.path().join("queue"));
2534 let mut t = task("held");
2535 q.put(&mut t).unwrap();
2536 assert_eq!(
2537 crate::notices::Notices::at(dir.path().join("notifications"))
2538 .list()
2539 .len(),
2540 0
2541 );
2542 t.hold_machine(Some("out of attempts".to_owned()));
2543 q.put(&mut t).unwrap();
2544 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
2545 assert_eq!(listed.len(), 1);
2546 assert!(listed[0].message.contains("out of attempts"));
2547 }
2548
2549 fn task(title: &str) -> Task {
2550 Task::new(
2551 title.to_owned(),
2552 format!("do {title}"),
2553 PathBuf::from("."),
2554 Source::Human,
2555 )
2556 }
2557
2558 #[test]
2559 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
2560 let mut t = task("retried");
2561 assert!(
2562 t.earlier_attempts().is_empty(),
2563 "a first attempt takes nothing over"
2564 );
2565 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2566 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
2567 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
2569 assert_eq!(t.successor_of("bbbb"), None);
2570 assert_eq!(t.successor_of("zzzz"), None);
2571 }
2572
2573 #[test]
2574 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
2575 let (_dir, q) = queue();
2576 let mut t = task("retried");
2577 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2578 q.put(&mut t).unwrap();
2579
2580 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
2581 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
2582 assert_eq!(
2583 q.superseded_by("cccc"),
2584 None,
2585 "the latest attempt replaces nothing"
2586 );
2587 assert_eq!(
2588 q.superseded_by("never-heard-of-it"),
2589 None,
2590 "a run belonging to no task on this queue is not superseded"
2591 );
2592
2593 let mut by = HashMap::new();
2594 by.insert("aaaa".to_owned(), "bbbb".to_owned());
2595 by.insert("bbbb".to_owned(), "cccc".to_owned());
2596 assert_eq!(
2597 q.superseded(),
2598 by,
2599 "the whole-map and single-run forms must agree"
2600 );
2601 }
2602
2603 #[test]
2604 fn superseded_attempts_is_empty_until_the_task_is_done() {
2605 let mut t = task("retried");
2606 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2607 t.status = TaskStatus::Failed;
2608 assert_eq!(
2609 t.superseded_attempts(true),
2610 &[] as &[String],
2611 "a task still retrying has no attempt yet that a later one made moot"
2612 );
2613
2614 t.status = TaskStatus::Running;
2615 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2616 }
2617
2618 #[test]
2619 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
2620 let mut t = task("retried");
2621 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2622 t.status = TaskStatus::Done;
2623 assert_eq!(
2624 t.superseded_attempts(true),
2625 &["aaaa".to_owned(), "bbbb".to_owned()],
2626 "cccc is the attempt whose success made the task done, and stays out"
2627 );
2628 }
2629
2630 #[test]
2631 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
2632 let mut t = task("first try landed");
2633 t.runs = vec!["aaaa".to_owned()];
2634 t.status = TaskStatus::Done;
2635 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
2636 }
2637
2638 #[test]
2639 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
2640 let mut t = task("closed by hand after a manual merge");
2647 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
2648 t.status = TaskStatus::Done;
2649 assert_eq!(
2650 t.superseded_attempts(false),
2651 &[] as &[String],
2652 "nothing here is provably why the task is done, so nothing is superseded"
2653 );
2654 }
2655
2656 #[test]
2657 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2658 let (_dir, q) = queue();
2659 let mut t = task("retried twice");
2660 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2661 q.put(&mut t).unwrap();
2662
2663 assert_eq!(
2664 q.latest_attempt("aaaa"),
2665 Some("cccc".to_owned()),
2666 "an old attempt points straight at the chain's current head, not the \
2667 next attempt in the middle of it"
2668 );
2669 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2670 assert_eq!(
2671 q.latest_attempt("cccc"),
2672 None,
2673 "the latest attempt is not superseded by anything"
2674 );
2675 assert_eq!(
2676 q.latest_attempt("never-heard-of-it"),
2677 None,
2678 "a run belonging to no task on this queue is not superseded"
2679 );
2680 }
2681
2682 #[test]
2683 fn a_markdown_heading_is_the_title_not_decoration() {
2684 assert_eq!(
2689 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2690 "Rework the config loader"
2691 );
2692 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2693 assert_eq!(title_from("> quoted task", 40), "quoted task");
2694 assert_eq!(title_from(" \n\n", 40), "(empty task)");
2696 assert_eq!(title_from("###\n", 40), "(empty task)");
2697 }
2698
2699 #[test]
2700 fn a_long_title_is_elided_by_characters_not_bytes() {
2701 let long = "課題".repeat(30);
2703 let title = title_from(&long, 10);
2704 assert_eq!(title.chars().count(), 10);
2705 assert!(title.ends_with('…'));
2706 }
2707
2708 #[test]
2709 fn priority_wins_and_ties_break_oldest_first() {
2710 let (_dir, q) = queue();
2711 let mut a = task("first");
2712 let mut b = task("second");
2713 let mut c = task("urgent");
2714 a.id = "20260101-000001-aaaa".to_owned();
2716 b.id = "20260101-000002-bbbb".to_owned();
2717 c.id = "20260101-000003-cccc".to_owned();
2718 c.priority = 5;
2719 for t in [&mut a, &mut b, &mut c] {
2720 q.put(t).unwrap();
2721 }
2722
2723 assert_eq!(q.next_runnable().unwrap().id, c.id);
2725 c.hold_machine(None);
2726 q.put(&mut c).unwrap();
2727 assert_eq!(q.next_runnable().unwrap().id, a.id);
2729 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2730 }
2731
2732 #[test]
2733 fn a_blocked_task_never_starves_another_runnable_one() {
2734 let (_dir, q) = queue();
2735 let mut blocked = task("blocked");
2736 blocked.block(vec!["something".to_owned()], None);
2737 q.put(&mut blocked).unwrap();
2738
2739 let mut runnable = task("free to go");
2740 q.put(&mut runnable).unwrap();
2741
2742 let next = q.next_runnable().expect("a runnable task is still offered");
2743 assert_eq!(next.id, runnable.id);
2744 }
2745
2746 #[test]
2747 fn a_held_task_is_never_offered_to_the_loop() {
2748 let (_dir, q) = queue();
2749 let mut t = task("held");
2750 q.put(&mut t).unwrap();
2751 assert!(q.next_runnable().is_some());
2752
2753 t.hold_machine(None);
2754 q.put(&mut t).unwrap();
2755 assert!(
2756 q.next_runnable().is_none(),
2757 "a held task must wait for a human"
2758 );
2759
2760 t.status = TaskStatus::Failed;
2762 q.put(&mut t).unwrap();
2763 assert!(q.next_runnable().is_some());
2764 }
2765
2766 #[test]
2767 fn attempts_are_capped_and_then_the_task_is_held() {
2768 let mut t = task("doomed");
2769
2770 t.start("run-1".to_owned());
2771 t.fail("gate red", 2);
2772 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2773
2774 t.start("run-2".to_owned());
2775 t.fail("gate red", 2);
2776 assert_eq!(
2777 t.status,
2778 TaskStatus::Held,
2779 "out of attempts: stop spending money on it"
2780 );
2781 assert_eq!(t.runs, ["run-1", "run-2"]);
2782 assert_eq!(t.last_error.as_deref(), Some("gate red"));
2783 assert_eq!(
2784 t.hold_reason.as_deref(),
2785 Some("gate red"),
2786 "the hold must say why, not leave hold_reason null next to a \
2787 populated last_error"
2788 );
2789 }
2790
2791 #[test]
2792 fn handing_off_a_task_records_a_hold_reason_too() {
2793 let mut t = task("left a pull request");
2794 t.start("run-1".to_owned());
2795 t.handed_off("run ended with a pull request open [run run-1]");
2796 assert_eq!(t.status, TaskStatus::Held);
2797 assert_eq!(t.hold_source, Some(HoldSource::Machine));
2798 assert_eq!(
2799 t.hold_reason.as_deref(),
2800 Some("run ended with a pull request open [run run-1]")
2801 );
2802 assert_eq!(t.hold_reason, t.last_error);
2803 }
2804
2805 #[test]
2806 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2807 let mut t = task("stalled by quota");
2808
2809 t.start("run-1".to_owned());
2810 assert_eq!(t.attempts, 1);
2811 t.stall("judge-1, judge-2 out of quota");
2812 assert_eq!(
2813 t.attempts, 0,
2814 "a closed quota window must not spend the task's retry budget"
2815 );
2816 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2817 assert_eq!(
2818 t.last_error.as_deref(),
2819 Some("judge-1, judge-2 out of quota")
2820 );
2821
2822 for _ in 0..20 {
2825 t.start("run-n".to_owned());
2826 t.stall("still out of quota");
2827 }
2828 t.start("run-real".to_owned());
2829 t.fail("gate red", 2);
2830 assert_eq!(
2831 t.status,
2832 TaskStatus::Failed,
2833 "the first attempt that was really judged is attempt one"
2834 );
2835 }
2836
2837 #[test]
2838 fn releasing_a_held_task_gives_it_a_real_second_chance() {
2839 let mut t = task("retry me");
2840 t.start("run-1".to_owned());
2841 t.fail("gate red", 1);
2842 assert_eq!(t.status, TaskStatus::Held);
2843
2844 t.release();
2845 assert_eq!(t.status, TaskStatus::Queued);
2846 assert_eq!(t.attempts, 0);
2849 assert!(t.last_error.is_none());
2850 assert_eq!(
2851 t.runs.len(),
2852 1,
2853 "history is kept: attempts reset, evidence does not"
2854 );
2855 }
2856
2857 #[test]
2858 fn a_hold_reason_survives_and_a_release_clears_it() {
2859 let mut t = task("waiting on something else");
2860 t.hold_manual(Some(
2861 "waiting for 20260101-000000-aaaa to land first".to_owned(),
2862 ));
2863 assert_eq!(t.status, TaskStatus::Held);
2864 assert_eq!(
2865 t.hold_reason.as_deref(),
2866 Some("waiting for 20260101-000000-aaaa to land first")
2867 );
2868
2869 t.hold_manual(None);
2871 assert_eq!(
2872 t.hold_reason.as_deref(),
2873 Some("waiting for 20260101-000000-aaaa to land first"),
2874 "a bare re-hold keeps whatever a human already wrote down"
2875 );
2876
2877 let mut plain = task("no reason given");
2879 plain.hold_manual(None);
2880 assert_eq!(plain.status, TaskStatus::Held);
2881 assert!(plain.hold_reason.is_none());
2882
2883 t.release();
2884 assert_eq!(t.status, TaskStatus::Queued);
2885 assert!(
2886 t.hold_reason.is_none(),
2887 "a stale reason must not greet the next person who holds this task"
2888 );
2889 }
2890
2891 #[test]
2892 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2893 let mut t = task("landed by hand while held");
2898 t.hold_manual(Some("waiting on 3ed9".to_owned()));
2899 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2900
2901 t.succeed();
2902 assert_eq!(t.status, TaskStatus::Done);
2903 assert!(
2904 t.hold_reason.is_none(),
2905 "a done task cannot still be waiting on something"
2906 );
2907 }
2908
2909 #[test]
2910 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2911 let mut held = task("held straight out of blocked");
2918 held.block(
2919 vec!["20260101-000000-dead".to_owned()],
2920 Some("waiting on the migration script".to_owned()),
2921 );
2922 assert_eq!(held.status, TaskStatus::Blocked);
2923
2924 held.hold_manual(None);
2925 assert_eq!(held.status, TaskStatus::Held);
2926 assert!(
2927 held.blocked_by.is_empty(),
2928 "hold overrides the wait, same as release"
2929 );
2930 assert!(held.block_reason.is_none());
2931
2932 let mut done = task("closed straight out of blocked");
2933 done.block(
2934 vec!["20260101-000000-dead".to_owned()],
2935 Some("waiting on the migration script".to_owned()),
2936 );
2937 done.succeed();
2938 assert_eq!(done.status, TaskStatus::Done);
2939 assert!(
2940 done.blocked_by.is_empty(),
2941 "a done task cannot still be waiting on a dependency"
2942 );
2943 assert!(done.block_reason.is_none());
2944 }
2945
2946 #[test]
2947 fn a_blocked_task_is_never_offered_to_the_loop() {
2948 let mut t = task("blocked");
2949 assert!(t.status.runnable());
2950 t.block(
2951 vec!["dep-id".to_owned()],
2952 Some("waits on dep-id".to_owned()),
2953 );
2954 assert_eq!(t.status, TaskStatus::Blocked);
2955 assert!(!t.status.runnable());
2956 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2957 }
2958
2959 #[test]
2960 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2961 let mut t = task("blocked on two");
2962 t.block(
2963 vec!["a".to_owned(), "b".to_owned()],
2964 Some("waits on a and b".to_owned()),
2965 );
2966
2967 t.unblock("a");
2968 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2969 assert_eq!(t.blocked_by, ["b"]);
2970
2971 t.unblock("b");
2972 assert_eq!(t.status, TaskStatus::Queued);
2973 assert!(t.blocked_by.is_empty());
2974 assert!(t.block_reason.is_none());
2975 }
2976
2977 #[test]
2978 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2979 let mut t = task("never blocked");
2980 t.unblock("whatever");
2981 assert_eq!(t.status, TaskStatus::Queued);
2982 }
2983
2984 #[test]
2985 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2986 let mut t = task("held, then asked about");
2991 t.hold_machine(Some("out of attempts".to_owned()));
2992 assert_eq!(t.status, TaskStatus::Held);
2993
2994 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2995 assert_eq!(t.status, TaskStatus::Blocked);
2996
2997 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2998 t.unblock("q1");
2999 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
3000 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
3001 assert_eq!(t.hold_source, Some(HoldSource::Machine));
3002 assert!(t.blocked_from.is_none(), "consumed once restored");
3003 }
3004
3005 #[test]
3006 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
3007 let mut t = task("manually held, then asked about");
3008 t.hold_manual(Some("waiting on a dependency".to_owned()));
3009
3010 t.block(vec!["q1".to_owned()], None);
3011 t.unblock("q1");
3012
3013 assert_eq!(t.status, TaskStatus::Held);
3014 assert_eq!(t.hold_source, Some(HoldSource::Manual));
3015 }
3016
3017 #[test]
3018 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
3019 let mut t = task("held, blocked twice");
3023 t.hold_machine(None);
3024 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
3025 t.block(
3026 vec!["q1".to_owned(), "q2".to_owned()],
3027 Some("second".to_owned()),
3028 );
3029
3030 t.unblock("q1");
3031 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
3032 t.unblock("q2");
3033 assert_eq!(t.status, TaskStatus::Held);
3034 }
3035
3036 #[test]
3037 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
3038 let mut t = task("blocked mid-run");
3042 t.start("run-1".to_owned());
3043 assert_eq!(t.status, TaskStatus::Running);
3044
3045 t.block(vec!["q1".to_owned()], None);
3046 t.unblock("q1");
3047 assert_eq!(t.status, TaskStatus::Queued);
3048 }
3049
3050 #[test]
3051 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
3052 let mut t = task("legacy record, held before it was blocked");
3058 t.hold_source = Some(HoldSource::Machine);
3059 t.hold_reason = Some("legacy hold reason".to_owned());
3060 t.status = TaskStatus::Blocked;
3061 t.blocked_by = vec!["q1".to_owned()];
3062 t.blocked_from = None;
3063
3064 t.unblock("q1");
3065 assert_eq!(t.status, TaskStatus::Held);
3066 }
3067
3068 #[test]
3069 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
3070 let mut t = task("legacy record, ordinary dependency block");
3071 t.status = TaskStatus::Blocked;
3072 t.blocked_by = vec!["dep".to_owned()];
3073 t.blocked_from = None;
3074
3075 t.unblock("dep");
3076 assert_eq!(t.status, TaskStatus::Queued);
3077 }
3078
3079 #[test]
3080 fn answering_a_question_is_recorded_and_survives_a_release() {
3081 let mut t = task("asked something");
3082 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
3083 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
3084 t.unblock("q1");
3085 assert_eq!(t.status, TaskStatus::Queued);
3086 assert_eq!(t.answers.len(), 1);
3087 assert_eq!(t.answers[0].answer, "SQLite");
3088
3089 t.release();
3093 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
3094 }
3095
3096 #[test]
3097 fn a_refused_handover_keeps_the_review_branch_across_release() {
3098 let mut t = task("refused takeover");
3099 t.start("run-1".to_owned());
3100 t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
3101 assert_eq!(t.status, TaskStatus::Held);
3102 assert_eq!(t.attempts, 1);
3103 t.release();
3104 assert_eq!(t.status, TaskStatus::Queued);
3105 assert_eq!(t.attempts, 0);
3106 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
3107
3108 t.hold_for_handover(None, "again".to_owned());
3110 t.requeue();
3111 assert!(t.review_branch.is_none());
3112
3113 let mut m = task("manual");
3115 m.review_branch = Some("magi/x/A".to_owned());
3116 m.hold_manual(None);
3117 m.release();
3118 assert!(m.review_branch.is_none());
3119 }
3120
3121 #[test]
3122 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
3123 let mut t = task("blocked run with a surviving branch");
3124 t.start("run-1".to_owned());
3125 t.fail("blocked with major findings", 5);
3126 assert_eq!(t.status, TaskStatus::Failed);
3127
3128 t.request_review("magi/eba2/A".to_owned());
3129 assert_eq!(t.status, TaskStatus::Queued);
3130 assert_eq!(t.attempts, 0);
3131 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
3132
3133 t.release();
3135 assert!(t.review_branch.is_none());
3136 }
3137
3138 #[test]
3139 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
3140 let mut t = task("retry");
3141 t.start("run-1".to_owned());
3142 t.requeue();
3143 assert!(t.fresh_start);
3144
3145 t.release();
3146 assert!(!t.fresh_start);
3147 }
3148
3149 #[test]
3150 fn priority_can_be_changed_while_queued_but_not_while_running() {
3151 let mut t = task("reprioritise me");
3152 t.set_priority(5).unwrap();
3153 assert_eq!(t.priority, 5);
3154
3155 t.start("run-1".to_owned());
3156 let err = t.set_priority(9).unwrap_err().to_string();
3157 assert!(err.contains("running"), "{err}");
3158 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
3159 }
3160
3161 #[test]
3162 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
3163 let mut t = task("interrupt me");
3164 assert!(!t.interrupt, "off unless asked, same as any other task");
3165
3166 t.set_interrupt(true).unwrap();
3167 assert!(t.interrupt);
3168
3169 t.start("run-1".to_owned());
3170 assert!(
3171 !t.interrupt,
3172 "the mark is one-shot: dispatching the task fulfils it, \
3173 whatever the run that follows ends up doing"
3174 );
3175 let err = t.set_interrupt(true).unwrap_err().to_string();
3176 assert!(err.contains("running"), "{err}");
3177 t.set_interrupt(false).unwrap();
3180 assert!(!t.interrupt);
3181 }
3182
3183 #[test]
3187 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
3188 let mut t = task("interrupt me");
3189 t.set_interrupt(true).unwrap();
3190 t.start("run-1".to_owned());
3191 t.fail("mock failure", 5);
3192 assert_eq!(t.status, TaskStatus::Failed);
3193 assert!(
3194 !t.interrupt,
3195 "one attempt already spent the mark; a retry is an ordinary \
3196 requeue, not a fresh interrupt request"
3197 );
3198 }
3199
3200 #[test]
3201 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
3202 let (_dir, q) = queue();
3203 let mut a = task("first filed");
3204 let mut b = task("second filed");
3205 a.id = "20260101-000001-aaaa".to_owned();
3206 b.id = "20260101-000002-bbbb".to_owned();
3207 q.put(&mut a).unwrap();
3208 q.put(&mut b).unwrap();
3209
3210 assert_eq!(
3211 q.next_runnable().unwrap().id,
3212 a.id,
3213 "with equal priority the older task goes first, so a burst of \
3214 new work cannot starve it"
3215 );
3216 assert_eq!(
3217 q.list()[0].id,
3218 b.id,
3219 "but the list an operator reads is newest first, the same as \
3220 before priority existed - a's turn to run does not make it the \
3221 newest task"
3222 );
3223
3224 let mut a = q.get(&a.id).unwrap();
3225 a.set_priority(10).unwrap();
3226 q.put(&mut a).unwrap();
3227
3228 assert_eq!(
3229 q.next_runnable().unwrap().id,
3230 a.id,
3231 "a raised priority must be reflected the moment it is saved"
3232 );
3233 assert_eq!(
3237 q.list()[0].id,
3238 a.id,
3239 "the raised task must sort first in the list an operator reads, \
3240 not only in next_runnable's own ordering"
3241 );
3242 }
3243
3244 #[test]
3245 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
3246 let mut t = Task::new(
3247 "old title".to_owned(),
3248 "old instruction".to_owned(),
3249 PathBuf::from("/repo"),
3250 Source::Agent {
3251 run: "20260101-000000-beef".to_owned(),
3252 node: "implement".to_owned(),
3253 },
3254 );
3255 let id = t.id.clone();
3256 let created_at = t.created_at;
3257 t.runs.push("20260101-000000-beef".to_owned());
3258
3259 t.edit("new title".to_owned(), "new instruction".to_owned())
3260 .unwrap();
3261
3262 assert_eq!(t.title, "new title");
3263 assert_eq!(t.instruction, "new instruction");
3264 assert_eq!(t.id, id, "editing must not mint a new id");
3265 assert_eq!(t.created_at, created_at);
3266 assert_eq!(
3267 t.source,
3268 Source::Agent {
3269 run: "20260101-000000-beef".to_owned(),
3270 node: "implement".to_owned(),
3271 },
3272 "editing must not turn agent attribution into human"
3273 );
3274 assert_eq!(t.runs, ["20260101-000000-beef"]);
3275 }
3276
3277 #[test]
3278 fn editing_is_refused_once_a_task_is_running_or_finished() {
3279 let mut running = task("in flight");
3280 running.start("run-1".to_owned());
3281 let err = running
3282 .edit("x".to_owned(), "y".to_owned())
3283 .unwrap_err()
3284 .to_string();
3285 assert!(err.contains("running"), "{err}");
3286
3287 let mut done = task("finished");
3288 done.succeed();
3289 let err = done
3290 .edit("x".to_owned(), "y".to_owned())
3291 .unwrap_err()
3292 .to_string();
3293 assert!(err.contains("done"), "{err}");
3294
3295 let mut queued = task("waiting");
3297 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
3298 let mut held = task("parked");
3299 held.hold_machine(None);
3300 held.edit("x".to_owned(), "y".to_owned()).unwrap();
3301 }
3302
3303 #[test]
3304 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
3305 let (_dir, q) = queue();
3306 let path = q.path_of("20260101-000000-aaaa");
3307 std::fs::create_dir_all(q.root()).unwrap();
3308 std::fs::write(
3309 &path,
3310 serde_json::json!({
3311 "schema": SCHEMA,
3312 "id": "20260101-000000-aaaa",
3313 "title": "from before hold reasons existed",
3314 "instruction": "from before hold reasons existed",
3315 "repo": ".",
3316 "source": { "kind": "human" },
3317 "status": "held",
3318 "created_at": Timestamp::now().to_string(),
3319 "updated_at": Timestamp::now().to_string(),
3320 })
3321 .to_string(),
3322 )
3323 .unwrap();
3324
3325 let task = q.get("20260101-000000-aaaa").expect("must still read");
3326 assert!(task.hold_reason.is_none());
3327 assert!(task.operator_held());
3328 }
3329
3330 #[test]
3331 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
3332 let (_dir, q) = queue();
3333 let path = q.path_of("20260101-000000-bbbb");
3334 std::fs::create_dir_all(q.root()).unwrap();
3335 std::fs::write(
3336 &path,
3337 serde_json::json!({
3338 "schema": 2,
3339 "id": "20260101-000000-bbbb",
3340 "title": "old manual recovery",
3341 "instruction": "old manual recovery",
3342 "repo": ".",
3343 "source": { "kind": "human" },
3344 "status": "held",
3345 "hold_reason": "active manual recovery run20260912-224242-daf5",
3346 "created_at": Timestamp::now().to_string(),
3347 "updated_at": Timestamp::now().to_string(),
3348 })
3349 .to_string(),
3350 )
3351 .unwrap();
3352
3353 let task = q.get("20260101-000000-bbbb").expect("must still read");
3354 assert_eq!(task.hold_source, None);
3355 assert!(task.operator_held());
3356 }
3357
3358 #[test]
3359 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
3360 let (_dir, q) = queue();
3361 let path = q.path_of("20260101-000000-aaaa");
3362 std::fs::create_dir_all(q.root()).unwrap();
3363 std::fs::write(
3364 &path,
3365 serde_json::json!({
3366 "schema": SCHEMA,
3367 "id": "20260101-000000-aaaa",
3368 "title": "from before diagnostics existed",
3369 "instruction": "from before diagnostics existed",
3370 "repo": ".",
3371 "source": { "kind": "human" },
3372 "status": "held",
3373 "created_at": Timestamp::now().to_string(),
3374 "updated_at": Timestamp::now().to_string(),
3375 })
3376 .to_string(),
3377 )
3378 .unwrap();
3379
3380 let task = q.get("20260101-000000-aaaa").expect("must still read");
3381 assert!(task.diagnostic.is_none());
3382 }
3383
3384 #[test]
3385 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
3386 let (_dir, q) = queue();
3390 let path = q.path_of("20260101-000000-aaaa");
3391 std::fs::create_dir_all(q.root()).unwrap();
3392 std::fs::write(
3393 &path,
3394 serde_json::json!({
3395 "schema": 1,
3396 "id": "20260101-000000-aaaa",
3397 "title": "from before blocking existed",
3398 "instruction": "from before blocking existed",
3399 "repo": ".",
3400 "source": { "kind": "human" },
3401 "status": "queued",
3402 "created_at": Timestamp::now().to_string(),
3403 "updated_at": Timestamp::now().to_string(),
3404 })
3405 .to_string(),
3406 )
3407 .unwrap();
3408
3409 let task = q.get("20260101-000000-aaaa").expect("must still read");
3410 assert!(task.blocked_by.is_empty());
3411 assert!(task.block_reason.is_none());
3412 assert!(task.answers.is_empty());
3413 assert!(task.review_branch.is_none());
3414 }
3415
3416 #[test]
3417 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
3418 let mut held = task("diagnosed");
3423 held.start("run-1".to_owned());
3424 held.fail("gate red", 1);
3425 held.diagnostic = Some("cargo test failed: ...".to_owned());
3426 assert_eq!(held.status, TaskStatus::Held);
3427
3428 held.release();
3429 assert!(held.diagnostic.is_none());
3430
3431 held.diagnostic = Some("cargo test failed: ...".to_owned());
3432 held.succeed();
3433 assert!(held.diagnostic.is_none());
3434 }
3435
3436 #[test]
3437 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
3438 let mut t = task("retried");
3439 t.start("run-1".to_owned());
3440 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
3441 t.fail("unrelated config error", 5);
3442 assert_eq!(t.status, TaskStatus::Failed);
3443 assert!(
3444 t.diagnostic.is_none(),
3445 "fail() must not let an old diagnostic outlive the run that produced it"
3446 );
3447 }
3448
3449 #[test]
3450 fn a_claim_is_exclusive_and_releases_on_drop() {
3451 let (_dir, q) = queue();
3452 let mut t = task("contended");
3453 q.put(&mut t).unwrap();
3454
3455 let held = q.claim(&t.id).unwrap();
3456 assert!(
3457 q.claim(&t.id).is_err(),
3458 "two daemons must not drive one task into two runs"
3459 );
3460 drop(held);
3461 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
3462 }
3463
3464 #[test]
3465 fn a_round_trip_survives_disk() {
3466 let (_dir, q) = queue();
3467 let mut t = Task::new(
3468 "titled".to_owned(),
3469 "body".to_owned(),
3470 PathBuf::from("/repo"),
3471 Source::Agent {
3472 run: "20260101-000000-beef".to_owned(),
3473 node: "implement".to_owned(),
3474 },
3475 );
3476 t.priority = 3;
3477 q.put(&mut t).unwrap();
3478
3479 let back = q.get(&t.id).unwrap();
3480 assert_eq!(back.id, t.id);
3481 assert_eq!(back.priority, 3);
3482 assert_eq!(back.source.label(), "implement@beef");
3483 assert_eq!(q.get(t.short()).unwrap().id, t.id);
3485 }
3486
3487 #[test]
3488 fn an_unreadable_task_does_not_take_the_queue_down() {
3489 let (_dir, q) = queue();
3490 let mut t = task("fine");
3491 q.put(&mut t).unwrap();
3492 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
3493
3494 let listed = q.list();
3495 assert_eq!(listed.len(), 1, "the readable task still lists");
3496 assert_eq!(listed[0].id, t.id);
3497 }
3498
3499 #[test]
3500 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
3501 let (_dir, q) = queue();
3502 let path = q.path_of("20260101-000000-aaaa");
3503 std::fs::create_dir_all(q.root()).unwrap();
3504 std::fs::write(
3505 &path,
3506 serde_json::json!({
3507 "schema": SCHEMA,
3508 "id": "20260101-000000-aaaa",
3509 "title": "from before solo existed",
3510 "instruction": "from before solo existed",
3511 "repo": ".",
3512 "source": { "kind": "human" },
3513 "status": "queued",
3514 "created_at": Timestamp::now().to_string(),
3515 "updated_at": Timestamp::now().to_string(),
3516 })
3517 .to_string(),
3518 )
3519 .unwrap();
3520
3521 let task = q.get("20260101-000000-aaaa").expect("must still read");
3522 assert!(!task.solo, "a queue file with no `solo` field means false");
3523 }
3524
3525 #[test]
3526 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
3527 let (_dir, q) = queue();
3528 let path = q.path_of("20260101-000000-bbbb");
3529 std::fs::create_dir_all(q.root()).unwrap();
3530 std::fs::write(
3531 &path,
3532 serde_json::json!({
3533 "schema": SCHEMA,
3534 "id": "20260101-000000-bbbb",
3535 "title": "from before urgent existed",
3536 "instruction": "from before urgent existed",
3537 "repo": ".",
3538 "source": { "kind": "human" },
3539 "status": "queued",
3540 "created_at": Timestamp::now().to_string(),
3541 "updated_at": Timestamp::now().to_string(),
3542 })
3543 .to_string(),
3544 )
3545 .unwrap();
3546
3547 let task = q.get("20260101-000000-bbbb").expect("must still read");
3548 assert!(
3549 !task.urgent,
3550 "a queue file with no `urgent` field means false, same as `solo`"
3551 );
3552 }
3553
3554 #[test]
3555 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
3556 let (_dir, q) = queue();
3557 let mut t = task("from the future");
3558 q.put(&mut t).unwrap();
3559 let path = q.path_of(&t.id);
3560 let body = std::fs::read_to_string(&path)
3561 .unwrap()
3562 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
3563 std::fs::write(&path, body).unwrap();
3564
3565 let err = q.get(&t.id).unwrap_err().to_string();
3566 assert!(err.contains("schema 99"), "{err}");
3567 }
3568
3569 #[test]
3570 fn revision_moves_when_the_queue_changes() {
3571 let (_dir, q) = queue();
3572 assert_eq!(q.revision(), 0, "an empty queue has no revision");
3573 let mut t = task("first");
3574 q.put(&mut t).unwrap();
3575 assert!(q.revision() > 0, "a written task moves the revision");
3576 }
3577
3578 #[test]
3579 fn revision_moves_when_deleting_an_older_task() {
3580 let (dir, q) = queue();
3581 let questions = Questions::at(dir.path().join("questions"));
3582 let mut t1 = task("older");
3583 q.put(&mut t1).unwrap();
3584 std::thread::sleep(std::time::Duration::from_millis(10));
3586 let mut t2 = task("newer");
3587 q.put(&mut t2).unwrap();
3588
3589 let rev_before = q.revision();
3590 q.remove(&t1.id, false, &questions).unwrap();
3591 let rev_after = q.revision();
3592
3593 assert_ne!(
3594 rev_before, rev_after,
3595 "deleting an older task must change the revision so other clients see the deletion"
3596 );
3597 }
3598
3599 #[test]
3600 fn removing_a_task_takes_it_out_of_the_listing() {
3601 let (dir, q) = queue();
3602 let questions = Questions::at(dir.path().join("questions"));
3603 let mut t = task("delete me");
3604 q.put(&mut t).unwrap();
3605 let removed = q.remove(t.short(), false, &questions).unwrap();
3606 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
3607 assert!(removed.released.is_empty() && removed.still_blocked.is_empty());
3608 assert!(q.list().is_empty());
3609 assert!(
3610 q.remove(&t.id, false, &questions).is_err(),
3611 "removing twice is an error"
3612 );
3613 }
3614
3615 #[test]
3616 fn removing_a_task_takes_its_stale_lock_with_it() {
3617 let (dir, q) = queue();
3618 let questions = Questions::at(dir.path().join("questions"));
3619 let mut t = task("interrupted");
3620 q.put(&mut t).unwrap();
3621
3622 let claim = q.claim(&t.id).unwrap();
3625 std::mem::forget(claim);
3626 assert!(
3627 q.claim(&t.id).is_err(),
3628 "the orphaned lock is what makes the task look claimed"
3629 );
3630
3631 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
3633 assert!(err.contains("live daemon"), "{err}");
3634 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
3635
3636 q.remove(&t.id, false, &questions).unwrap();
3638 assert!(q.list().is_empty());
3639 let mut again = task("interrupted");
3640 again.id = t.id.clone();
3641 q.put(&mut again).unwrap();
3642 assert!(
3643 q.claim(&t.id).is_ok(),
3644 "a task that comes back must be claimable, which a left-behind lock would prevent"
3645 );
3646 }
3647
3648 fn notices_of(dir: &Path) -> Vec<crate::notices::Notice> {
3649 crate::notices::Notices::at(dir.join("notifications")).list()
3650 }
3651
3652 #[test]
3653 fn removing_a_sole_dependency_releases_the_dependent_without_a_hold() {
3654 let (dir, q) = queue();
3655 let questions = Questions::at(dir.path().join("questions"));
3656 let mut dep = task("dependency");
3657 q.put(&mut dep).unwrap();
3658 let mut blocked = task("waiting");
3659 blocked.block(vec![dep.id.clone()], Some("waits".to_owned()));
3660 q.put(&mut blocked).unwrap();
3661
3662 let removed = q.remove(&dep.id, false, &questions).unwrap();
3663 assert_eq!(removed.released, [blocked.id.clone()]);
3664 assert!(removed.still_blocked.is_empty());
3665
3666 let after = q.get(&blocked.id).unwrap();
3667 assert_eq!(after.status, TaskStatus::Queued);
3668 assert!(after.blocked_by.is_empty());
3669 assert!(after.block_reason.is_none());
3670 let notes = notices_of(dir.path());
3671 assert_eq!(notes.len(), 1, "{notes:?}");
3672 assert_eq!(notes[0].severity, crate::notices::Severity::Info);
3673 assert_eq!(
3674 notes[0].link,
3675 Some(crate::notices::Link::Task {
3676 id: blocked.id.clone(),
3677 })
3678 );
3679 }
3680
3681 #[test]
3682 fn removing_one_of_two_dependencies_keeps_the_other() {
3683 let (dir, q) = queue();
3684 let questions = Questions::at(dir.path().join("questions"));
3685 let mut dep = task("dependency");
3686 q.put(&mut dep).unwrap();
3687 let mut other = task("other");
3688 q.put(&mut other).unwrap();
3689 let mut blocked = task("waiting");
3690 blocked.block(
3691 vec![dep.id.clone(), other.id.clone()],
3692 Some("waits on both".to_owned()),
3693 );
3694 q.put(&mut blocked).unwrap();
3695
3696 let removed = q.remove(&dep.id, false, &questions).unwrap();
3697 assert!(removed.released.is_empty());
3698 assert_eq!(removed.still_blocked, [blocked.id.clone()]);
3699
3700 let after = q.get(&blocked.id).unwrap();
3701 assert_eq!(after.status, TaskStatus::Blocked);
3702 assert_eq!(after.blocked_by, [other.id.clone()]);
3703 assert!(missing_blockers(&q, &questions, &after.blocked_by).is_empty());
3706 }
3707
3708 #[test]
3709 fn a_dependency_deleted_before_its_dependents_were_rewritten_is_released_later() {
3710 let (dir, q) = queue();
3711 let questions = Questions::at(dir.path().join("questions"));
3712 let mut dep = task("dependency");
3713 q.put(&mut dep).unwrap();
3714 let mut blocked = task("waiting");
3715 blocked.block(vec![dep.id.clone()], None);
3716 q.put(&mut blocked).unwrap();
3717
3718 let claim = q.claim(&blocked.id).unwrap();
3720 let removed = q.remove(&dep.id, false, &questions).unwrap();
3721 assert!(removed.released.is_empty());
3722 drop(claim);
3723
3724 let mut task = q.get(&blocked.id).unwrap();
3725 assert_eq!(task.status, TaskStatus::Blocked);
3726 assert!(missing_blockers(&q, &questions, &task.blocked_by).is_empty());
3727 assert_eq!(q.apply_deleted_blockers(&mut task), [dep.id.clone()]);
3728 assert_eq!(task.status, TaskStatus::Queued);
3729 }
3730
3731 #[test]
3732 fn a_dependency_without_a_tombstone_is_still_missing() {
3733 let (dir, q) = queue();
3734 let questions = Questions::at(dir.path().join("questions"));
3735 let mut dep = task("dependency");
3736 q.put(&mut dep).unwrap();
3737 std::fs::remove_file(q.path_of(&dep.id)).unwrap();
3738 assert_eq!(
3739 missing_blockers(&q, &questions, std::slice::from_ref(&dep.id)),
3740 [dep.id.clone()]
3741 );
3742 assert!(deleted_blockers(&q, std::slice::from_ref(&dep.id)).is_empty());
3743 }
3744
3745 #[test]
3746 fn a_held_dependent_is_left_alone_by_a_dependency_deletion() {
3747 let mut t = task("held");
3748 t.hold_machine(Some("because".to_owned()));
3749 t.blocked_by = vec!["gone".to_owned()];
3750 assert!(!t.dependency_deleted("gone"));
3751 assert_eq!(t.status, TaskStatus::Held);
3752 }
3753
3754 #[test]
3755 fn a_failed_record_removal_takes_the_tombstone_back() {
3756 let (_dir, q) = queue();
3757 let mut t = task("stays");
3758 q.put(&mut t).unwrap();
3759 q.write_tombstone(&t.id).unwrap();
3760 let err = q.remove_record_with_attachments(&t.id, |_| Err(std::io::Error::other("nope")));
3761 assert!(err.is_err());
3762 assert!(!q.was_deleted(&t.id));
3765 }
3766
3767 fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3768 let p = dir.join(name);
3769 std::fs::write(&p, body).unwrap();
3770 p
3771 }
3772
3773 #[test]
3774 fn an_attachment_copy_survives_deleting_its_source() {
3775 let (dir, q) = queue();
3776 let src = source_file(dir.path(), "shot.png", "pixels");
3777 let mut t = task("with a picture");
3778 let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3779 q.put(&mut t).unwrap();
3780 std::fs::remove_file(&src).unwrap();
3781 assert_eq!(names, ["shot.png"]);
3782 let loaded = q.get(&t.id).unwrap();
3783 let paths = q.attachment_paths(&loaded);
3784 assert_eq!(paths.len(), 1);
3785 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3786 }
3787
3788 #[test]
3789 fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3790 let (dir, q) = queue();
3791 let mut t = task("bad names");
3792 for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3793 let src = source_file(dir.path(), name, "x");
3794 assert!(
3795 q.attach(&mut t, &[src]).is_err(),
3796 "`{name}` must be refused"
3797 );
3798 }
3799 assert!(!crate::ask::valid_asset_name("C:foo.png"));
3803 #[cfg(not(windows))]
3804 {
3805 let src = source_file(dir.path(), "C:foo.png", "x");
3806 assert!(q.attach(&mut t, &[src]).is_err());
3807 }
3808 let long = format!("{}.png", "a".repeat(70));
3809 let src = source_file(dir.path(), &long, "x");
3810 assert!(q.attach(&mut t, &[src]).is_err());
3811 assert!(t.attachments.is_empty());
3812 assert!(!q.attachments_dir(&t.id).exists());
3813 }
3814
3815 #[test]
3816 fn a_taken_attachment_name_is_numbered_not_overwritten() {
3817 let (dir, q) = queue();
3818 let a = source_file(dir.path(), "shot.png", "one");
3819 let sub = dir.path().join("other");
3820 std::fs::create_dir_all(&sub).unwrap();
3821 let b = source_file(&sub, "shot.png", "two");
3822 let mut t = task("collision");
3823 q.attach(&mut t, &[a]).unwrap();
3824 q.attach(&mut t, &[b]).unwrap();
3825 assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3826 let paths = q.attachment_paths(&t);
3827 assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3828 assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3829 }
3830
3831 #[test]
3832 fn a_renumbered_name_stays_inside_the_length_bound() {
3833 let (dir, q) = queue();
3834 let name = format!("{}.png", "a".repeat(60));
3835 assert_eq!(name.len(), 64);
3836 let a = source_file(dir.path(), &name, "one");
3837 let sub = dir.path().join("other");
3838 std::fs::create_dir_all(&sub).unwrap();
3839 let b = source_file(&sub, &name, "two");
3840 let mut t = task("long");
3841 q.attach(&mut t, &[a, b]).unwrap();
3842 assert_eq!(t.attachments.len(), 2);
3843 assert!(
3844 t.attachments
3845 .iter()
3846 .all(|n| crate::ask::valid_asset_name(n))
3847 );
3848 assert!(t.attachments[1].ends_with("-2.png"));
3849 }
3850
3851 #[test]
3852 fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3853 let (dir, q) = queue();
3854 let good = source_file(dir.path(), "good.png", "ok");
3855 let mut t = task("partial");
3856 q.attach(&mut t, &[good]).unwrap();
3857 let more = source_file(dir.path(), "more.png", "ok");
3858 let missing = dir.path().join("missing.png");
3859 assert!(q.attach(&mut t, &[more, missing]).is_err());
3860 assert_eq!(t.attachments, ["good.png"]);
3861 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3862 .unwrap()
3863 .flatten()
3864 .collect();
3865 assert_eq!(on_disk.len(), 1);
3866 }
3867
3868 fn block_put(q: &Queue, t: &Task) -> PathBuf {
3870 let tmp = q.path_of(&t.id).with_extension("json.tmp");
3871 std::fs::create_dir_all(&tmp).unwrap();
3872 tmp
3873 }
3874
3875 #[test]
3876 fn a_failed_put_leaves_no_new_attachment_directory() {
3877 let (dir, q) = queue();
3878 let mut t = task("fresh");
3879 let tmp = block_put(&q, &t);
3880 let src = source_file(dir.path(), "shot.png", "x");
3881 assert!(q.attach_and_put(&mut t, &[src]).is_err());
3882 assert!(t.attachments.is_empty());
3883 assert!(!q.attachments_dir(&t.id).exists());
3884 assert!(!q.path_of(&t.id).exists());
3885 std::fs::remove_dir(tmp).unwrap();
3886 }
3887
3888 #[test]
3889 fn a_failed_put_removes_only_the_copy_it_just_made() {
3890 let (dir, q) = queue();
3891 let mut t = task("edited");
3892 let first = source_file(dir.path(), "first.png", "1");
3893 q.attach_and_put(&mut t, &[first]).unwrap();
3894 block_put(&q, &t);
3895 let second = source_file(dir.path(), "second.png", "2");
3896 assert!(q.attach_and_put(&mut t, &[second]).is_err());
3897 assert_eq!(t.attachments, ["first.png"]);
3898 let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3899 .unwrap()
3900 .flatten()
3901 .map(|e| e.file_name().to_string_lossy().into_owned())
3902 .collect();
3903 assert_eq!(on_disk, ["first.png"]);
3904 assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3905 }
3906
3907 #[test]
3908 fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3909 let (dir, q) = queue();
3910 let questions = Questions::at(dir.path().join("questions"));
3911 let gone = task("gone");
3912 let mut other = task("other");
3913 let mut live = task("live");
3914 q.put(&mut other).unwrap();
3915 q.put(&mut live).unwrap();
3916 let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3919 std::fs::create_dir_all(&orphan).unwrap();
3920 std::fs::write(orphan.join("shot.png"), "x").unwrap();
3921 let busy = q.root.join(format!("{}.attachments.removing", live.id));
3924 std::fs::create_dir_all(&busy).unwrap();
3925
3926 q.remove(&other.id, false, &questions).unwrap();
3927 assert!(!orphan.exists(), "an orphan is swept");
3928 assert!(busy.exists(), "a removal in progress is left alone");
3929 }
3930
3931 #[test]
3932 fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3933 let (dir, q) = queue();
3934 let questions = Questions::at(dir.path().join("questions"));
3935 let mut t = task("stuck");
3936 let src = source_file(dir.path(), "shot.png", "x");
3937 q.attach_and_put(&mut t, &[src]).unwrap();
3938 let aside = q.root.join(format!("{}.attachments.removing", t.id));
3941 std::fs::create_dir_all(&aside).unwrap();
3942 std::fs::write(aside.join("old.png"), "o").unwrap();
3943 assert!(q.remove(&t.id, false, &questions).is_err());
3944 assert!(q.path_of(&t.id).exists());
3945 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3946 }
3947
3948 #[test]
3949 fn a_failed_record_removal_puts_the_attachments_back() {
3950 let (dir, q) = queue();
3951 let mut t = task("rollback");
3952 let src = source_file(dir.path(), "shot.png", "x");
3953 q.attach_and_put(&mut t, &[src]).unwrap();
3954 let err = q
3955 .remove_record_with_attachments(&t.id, |_| {
3956 Err(std::io::Error::other("injected failure"))
3957 })
3958 .unwrap_err();
3959 assert!(format!("{err:#}").contains("injected failure"));
3960 assert!(q.path_of(&t.id).exists());
3961 assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3962 assert!(
3963 !q.root
3964 .join(format!("{}.attachments.removing", t.id))
3965 .exists()
3966 );
3967 }
3968
3969 #[test]
3970 fn editing_a_task_keeps_its_attachments() {
3971 let (dir, q) = queue();
3972 let src = source_file(dir.path(), "shot.png", "x");
3973 let mut t = task("editable");
3974 q.attach(&mut t, &[src]).unwrap();
3975 t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3976 q.put(&mut t).unwrap();
3977 assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3978 }
3979
3980 #[test]
3981 fn removing_a_task_deletes_its_attachments() {
3982 let (dir, q) = queue();
3983 let questions = Questions::at(dir.path().join("questions"));
3984 let src = source_file(dir.path(), "shot.png", "x");
3985 let mut t = task("doomed");
3986 q.attach(&mut t, &[src]).unwrap();
3987 q.put(&mut t).unwrap();
3988 assert!(q.attachments_dir(&t.id).is_dir());
3989 q.remove(&t.id, false, &questions).unwrap();
3990 assert!(!q.attachments_dir(&t.id).exists());
3991 assert!(q.list().is_empty());
3992 }
3993
3994 #[test]
3995 fn a_task_written_before_attachments_still_reads() {
3996 let (_dir, q) = queue();
3997 let mut t = task("old");
3998 q.put(&mut t).unwrap();
3999 let path = q.path_of(&t.id);
4000 let mut v: serde_json::Value =
4001 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
4002 v.as_object_mut().unwrap().remove("attachments");
4003 std::fs::write(&path, v.to_string()).unwrap();
4004 assert!(q.get(&t.id).unwrap().attachments.is_empty());
4005 }
4006
4007 #[test]
4008 fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
4009 let q = Queue::at(PathBuf::from("relative-queue"));
4010 let mut t = task("rel");
4011 t.attachments.push("shot.png".to_owned());
4012 let paths = q.attachment_paths(&t);
4013 assert!(paths[0].is_absolute(), "{}", paths[0].display());
4014 assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
4015 }
4016
4017 #[test]
4018 fn link_run_adds_a_run_once_and_touches_nothing_else() {
4019 let dir = tempfile::tempdir().unwrap();
4020 let queue = Queue::at(dir.path().join("queue"));
4021 let mut t = Task::new(
4022 "t".to_owned(),
4023 "do it".to_owned(),
4024 PathBuf::from("."),
4025 Source::Human,
4026 );
4027 queue.put(&mut t).unwrap();
4028 let before = queue.get(&t.id).unwrap();
4029
4030 let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
4031 assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
4032 assert_eq!(linked.status, before.status);
4033 assert_eq!(linked.attempts, before.attempts);
4034 assert_eq!(linked.interrupt, before.interrupt);
4035
4036 let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
4037 assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
4038 assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
4039 assert!(queue.link_run("no-such-task", "r").is_err());
4040 }
4041
4042 #[test]
4043 fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
4044 let dir = tempfile::tempdir().unwrap();
4045 let queue = Queue::at(dir.path().join("queue"));
4046 let mut t = Task::new(
4047 "t".to_owned(),
4048 "do it".to_owned(),
4049 PathBuf::from("."),
4050 Source::Human,
4051 );
4052 queue.put(&mut t).unwrap();
4053 let mut snapshot = queue.get(&t.id).unwrap();
4055 queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
4056
4057 snapshot.start("20260930-000000-aaaa".to_owned());
4058 queue.put(&mut snapshot).unwrap();
4059
4060 let stored = queue.get(&t.id).unwrap();
4061 assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
4062 assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
4063 assert_eq!(stored.attempts, 1);
4064 }
4065
4066 #[test]
4067 fn concurrent_links_and_daemon_saves_lose_nothing() {
4068 let dir = tempfile::tempdir().unwrap();
4069 let queue = Queue::at(dir.path().join("queue"));
4070 let mut t = Task::new(
4071 "t".to_owned(),
4072 "do it".to_owned(),
4073 PathBuf::from("."),
4074 Source::Human,
4075 );
4076 queue.put(&mut t).unwrap();
4077 let id = t.id.clone();
4078
4079 let linkers: Vec<_> = (0..4)
4080 .map(|n| {
4081 let (queue, id) = (queue.clone(), id.clone());
4082 std::thread::spawn(move || {
4083 for k in 0..10 {
4084 queue
4085 .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
4086 .unwrap();
4087 }
4088 })
4089 })
4090 .collect();
4091 let mut mine = queue.get(&id).unwrap();
4094 for k in 0..10 {
4095 mine.start(format!("20260930-000009-d{k:03}"));
4096 queue.put(&mut mine).unwrap();
4097 }
4098 for l in linkers {
4099 l.join().unwrap();
4100 }
4101
4102 let stored = queue.get(&id).unwrap();
4103 assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
4104 assert_eq!(
4105 stored.attempts, 10,
4106 "linking never rewinds the daemon's work"
4107 );
4108 }
4109
4110 fn age_lock(path: &Path) {
4111 let f = std::fs::OpenOptions::new().write(true).open(path).unwrap();
4112 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
4113 .unwrap();
4114 }
4115
4116 #[test]
4117 fn concurrent_stale_takeover_yields_one_holder() {
4118 use std::sync::atomic::{AtomicUsize, Ordering};
4119 use std::sync::{Arc, Barrier};
4120 for _ in 0..5 {
4121 let dir = tempfile::tempdir().unwrap();
4122 let q = Queue::at(dir.path().to_path_buf());
4123 let lock = dir.path().join("t.write-lock");
4124 std::fs::write(&lock, "dead-0000").unwrap();
4125 age_lock(&lock);
4126 let n = 6;
4127 let barrier = Arc::new(Barrier::new(n));
4128 let (now, max) = (Arc::new(AtomicUsize::new(0)), Arc::new(AtomicUsize::new(0)));
4129 let handles: Vec<_> = (0..n)
4130 .map(|_| {
4131 let (q, b, now, max) = (q.clone(), barrier.clone(), now.clone(), max.clone());
4132 std::thread::spawn(move || {
4133 b.wait();
4134 let g = q.lock_task("t").unwrap();
4135 let held = now.fetch_add(1, Ordering::SeqCst) + 1;
4136 max.fetch_max(held, Ordering::SeqCst);
4137 std::thread::sleep(std::time::Duration::from_millis(20));
4138 now.fetch_sub(1, Ordering::SeqCst);
4139 drop(g);
4140 })
4141 })
4142 .collect();
4143 for h in handles {
4144 h.join().unwrap();
4145 }
4146 assert_eq!(max.load(Ordering::SeqCst), 1);
4147 assert!(!lock.exists());
4148 }
4149 }
4150
4151 #[test]
4152 fn dropping_a_stolen_lock_leaves_the_new_holders_lock() {
4153 let dir = tempfile::tempdir().unwrap();
4154 let q = Queue::at(dir.path().to_path_buf());
4155 let lock = dir.path().join("t.write-lock");
4156 let a = q.lock_task("t").unwrap();
4157 age_lock(&lock);
4158 let b = q.lock_task("t").unwrap();
4159 assert_ne!(a.token, b.token);
4160 drop(a);
4161 assert_eq!(std::fs::read_to_string(&lock).unwrap(), b.token);
4162 drop(b);
4163 assert!(!lock.exists());
4164 }
4165
4166 #[test]
4167 fn break_stale_leaves_a_lock_that_replaced_the_one_judged() {
4168 let dir = tempfile::tempdir().unwrap();
4169 let lock = dir.path().join("t.write-lock");
4170 std::fs::write(&lock, "new-token").unwrap();
4171 assert!(!break_stale(&lock, "old-token"));
4172 assert_eq!(std::fs::read_to_string(&lock).unwrap(), "new-token");
4173 assert!(!break_stale(&lock, "new-token"));
4175 assert!(lock.exists());
4176 age_lock(&lock);
4177 assert!(break_stale(&lock, "new-token"));
4178 assert!(!lock.exists());
4179 }
4180
4181 #[test]
4182 fn a_stale_break_marker_is_recovered_and_a_fresh_one_is_respected() {
4183 let dir = tempfile::tempdir().unwrap();
4184 let lock = dir.path().join("t.write-lock");
4185 std::fs::write(&lock, "dead-1").unwrap();
4186 age_lock(&lock);
4187 let marker = break_marker(&lock, "dead-1");
4188 std::fs::write(&marker, "crashed-remover").unwrap();
4189 assert!(!break_stale(&lock, "dead-1"));
4191 assert!(lock.exists());
4192 age_lock(&marker);
4194 assert!(break_stale(&lock, "dead-1"));
4195 assert!(!lock.exists());
4196 assert!(!marker.exists());
4197 }
4198
4199 fn chat_task() -> Task {
4200 Task::new(
4201 "t".to_owned(),
4202 "do t".to_owned(),
4203 PathBuf::from("."),
4204 Source::Agent {
4205 run: "20260902-000000-beef".to_owned(),
4206 node: CHAT_NODE.to_owned(),
4207 },
4208 )
4209 }
4210
4211 #[test]
4212 fn a_chat_notice_is_due_only_for_a_final_state_and_only_once() {
4213 let mut t = chat_task();
4214 assert_eq!(t.chat_notice_due(), None, "queued");
4215 t.start("r1".to_owned());
4216 assert_eq!(t.chat_notice_due(), None, "running");
4217 t.fail("boom", 3);
4218 assert_eq!(t.status, TaskStatus::Failed);
4219 assert_eq!(t.chat_notice_due(), None, "a retried failure is not final");
4220 t.start("r2".to_owned());
4221 t.park("upgrade");
4222 assert_eq!(t.chat_notice_due(), None, "parked");
4223 t.fail("boom", 1);
4224 assert_eq!(t.status, TaskStatus::Held);
4225 let key = t
4226 .chat_notice_due()
4227 .expect("held on the last attempt is final");
4228 t.chat_notice = Some(ChatNotice {
4229 key,
4230 at: Timestamp::now(),
4231 outcome: ChatNoticeOutcome::Posted,
4232 });
4233 assert_eq!(t.chat_notice_due(), None, "same hold, no second notice");
4234 t.release();
4235 assert_eq!(t.chat_notice_due(), None, "requeued after a release");
4236 t.start("r3".to_owned());
4237 t.succeed();
4238 assert!(t.chat_notice_due().is_some(), "done is a new final state");
4239 }
4240
4241 #[test]
4242 fn a_new_hold_after_a_release_is_a_new_cause() {
4243 let mut t = chat_task();
4244 t.hold_manual(Some("wait".to_owned()));
4245 let first = t.chat_notice_due().expect("held");
4246 t.chat_notice = Some(ChatNotice {
4247 key: first.clone(),
4248 at: Timestamp::now(),
4249 outcome: ChatNoticeOutcome::Posted,
4250 });
4251 t.hold_manual(Some("still waiting".to_owned()));
4252 assert_eq!(t.chat_notice_due(), None, "a re-hold keeps its start");
4253 t.release();
4254 std::thread::sleep(std::time::Duration::from_millis(5));
4255 t.hold_manual(None);
4256 let second = t.chat_notice_due().expect("a fresh hold");
4257 assert_ne!(first, second);
4258 }
4259
4260 #[test]
4261 fn a_chat_notice_needs_a_chat_source_and_a_current_schema() {
4262 let mut human = task("h");
4263 human.succeed();
4264 assert_eq!(human.chat_notice_due(), None);
4265 let mut inherited = task("i");
4267 inherited.origin_chat = Some("20260902-000000-beef".to_owned());
4268 inherited.succeed();
4269 assert_eq!(inherited.chat_notice_due(), None);
4270 let mut old = chat_task();
4271 old.schema = 12;
4272 old.status = TaskStatus::Done;
4273 assert_eq!(old.chat_notice_due(), None, "history is not replayed");
4274 let mut live = chat_task();
4275 live.schema = 12;
4276 live.start("r1".to_owned());
4277 live.succeed();
4278 assert!(live.chat_notice_due().is_some(), "started by this build");
4279 let mut held = chat_task();
4280 held.schema = 12;
4281 held.hold_machine(Some("config".to_owned()));
4282 assert!(held.chat_notice_due().is_some(), "held before any start");
4283 }
4284
4285 #[test]
4286 fn a_stale_snapshot_does_not_rewind_a_recorded_chat_notice() {
4287 let tmp = tempfile::tempdir().expect("tempdir");
4288 let q = Queue::at(tmp.path().join("queue"));
4289 let mut t = chat_task();
4290 t.succeed();
4291 q.put(&mut t).expect("put");
4292 let mut stale = q.get(&t.id).expect("get");
4293 let key = t.chat_notice_due().expect("due");
4294 q.modify(&t.id, |x| {
4295 x.chat_notice = Some(ChatNotice {
4296 key: key.clone(),
4297 at: Timestamp::now(),
4298 outcome: ChatNoticeOutcome::Posted,
4299 });
4300 true
4301 })
4302 .expect("modify");
4303 q.put(&mut stale).expect("stale put");
4304 assert_eq!(q.get(&t.id).expect("get").chat_notice_due(), None);
4305 }
4306}