magi/queue.rs
1//! The task queue: what magi should do next, and who asked for it.
2//!
3//! The queue is what lets magi run unattended. `magi serve` takes the next
4//! task, runs the graph on it, records the outcome, and takes the next one.
5//!
6//! It is also the reason an agent can ask for work. `magi task add` is the
7//! whole interface, and it is the same command whether a human types it at a
8//! prompt, a phone posts it through the web UI, or an implementer inside a run
9//! shells out to it because it noticed something worth doing but out of scope.
10//! magi's CLI is the operating surface for both kinds of user; the queue is
11//! where their intentions meet.
12//!
13//! One task is one JSON file under [`Queue`]'s root. Files rather than a
14//! database because the operator has to be able to read, edit, and delete the
15//! backlog with the tools already on the machine, and because a crashed daemon
16//! must leave a queue the next one can pick up without recovery ceremony.
17//!
18//! # Shape
19//!
20//! [`Task`] is data plus *pure* state transitions - [`Task::fail`] decides
21//! whether an attempt was the last one, and touches no disk. [`Queue`] owns all
22//! I/O and is constructed with its root, so a test drives a real queue in a
23//! temp directory without setting a process-global home. Splitting them this
24//! way is why the retry policy below can be asserted directly.
25//!
26//! # Bounded by construction
27//!
28//! An autonomous loop that retries forever is a way to spend money on a task
29//! that cannot succeed. Every claim increments [`Task::attempts`]; a task that
30//! has burned its attempts becomes [`TaskStatus::Held`] and waits for a human
31//! rather than for another agent.
32
33use 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
42/// On-disk format for a queued task. Bumped when a field's meaning changes.
43///
44/// 6: added [`Task::resume_override`] (a field only; `#[serde(default)]`, so
45/// an older record reads as `None` and [`read_path`] still accepts it).
46///
47/// 5: added [`Task::triage_applied`], the ids of triage questions whose
48/// answer has already been applied to this task. A "resume" answer used to
49/// leave no trace ([`Task::release`] clears [`Task::hold_reason`], which was
50/// the only place the applied marker lived), so when the released task failed
51/// its attempts and went back to `held`, the next idle pass found the same
52/// answered question "not applied" and released it again with `attempts` reset
53/// to 0 - the `max_attempts` bound never held. `#[serde(default)]` so an older
54/// record reads as empty. A task already looping when this build arrives has
55/// no record, so it is released once more, recorded, and then stays held.
56///
57/// 4: added [`Task::blocked_from`], the status a task had the moment it
58/// became [`TaskStatus::Blocked`], so [`Task::unblock`] restores it instead
59/// of always landing on [`TaskStatus::Queued`]. Without it, a task a human
60/// or `crate::triage` had deliberately left [`TaskStatus::Held`] — machine
61/// or manual — would lose that the instant `crate::conduct` blocked it on a
62/// follow-up question, and come back `Queued` the moment the question was
63/// answered, regardless of what the answer said: exactly the loop where a
64/// task the operator told to stay held instead re-enters the competition
65/// queue every time someone answers a question about it. `#[serde(default)]`
66/// so an older record reads as `None`; [`Task::unblock`] then falls back to
67/// inferring `Held` from surviving hold evidence ([`Task::hold_reason`] /
68/// [`Task::hold_source`], never cleared by [`Task::block`]) rather than
69/// guessing `Queued` outright — see [`Task::unblock`]'s own doc.
70///
71/// 3: added [`HoldSource`] so conductor recovery cannot release a hold an
72/// operator deliberately placed. Old records default to `None` and are
73/// protected as operator-held until an explicit release; the safe direction
74/// when their author was never recorded.
75///
76/// 2: added [`TaskStatus::Blocked`], [`Task::blocked_by`] and
77/// [`Task::block_reason`] (`crate::conduct`'s decisions) and
78/// [`Task::answers`] (operator answers carried forward to the next
79/// conductor prompt and the next run's instruction). All three are
80/// `#[serde(default)]`, so [`read_path`] accepts anything up to and
81/// including this schema rather than only an exact match — a task written
82/// by a build that only knew about schema 1 has nothing to say about
83/// blocking or answers, and defaulting those fields is exactly as good a
84/// reading as a value that build never had a chance to write.
85pub const SCHEMA: u32 = 6;
86
87/// Who placed the current hold.
88#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(rename_all = "lowercase")]
90pub enum HoldSource {
91 /// An operator used the CLI or web UI.
92 Manual,
93 /// The daemon or conductor placed the hold as part of its own recovery.
94 Machine,
95}
96
97impl HoldSource {
98 /// Short human-facing label for reports and the CLI.
99 pub fn label(self) -> &'static str {
100 match self {
101 Self::Manual => "manual",
102 Self::Machine => "machine",
103 }
104 }
105}
106
107/// Where a task came from. Recorded because "who asked for this" is the first
108/// question about an autonomous run, and the answer is not recoverable later.
109#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
110#[serde(tag = "kind", rename_all = "lowercase")]
111pub enum Source {
112 /// A person, at a terminal or through the web UI.
113 Human,
114 /// An agent inside a run, via `magi task add`. Both ids are recorded so a
115 /// task can be traced back to the exact seat that asked for it.
116 Agent {
117 /// Run the asking agent belonged to.
118 run: String,
119 /// Node it was working in, e.g. `implement` or `review`.
120 node: String,
121 },
122 /// A GitHub issue, imported by number.
123 Issue {
124 /// Issue number.
125 number: u64,
126 /// `owner/repo`, as `gh` reports it.
127 repo: String,
128 },
129}
130
131impl Source {
132 /// Short human-facing label, for lists and the web UI.
133 pub fn label(&self) -> String {
134 match self {
135 Self::Human => "human".to_owned(),
136 Self::Agent { run, node } => format!("{node}@{}", short(run)),
137 Self::Issue { number, .. } => format!("issue #{number}"),
138 }
139 }
140}
141
142/// Where a task is in its life.
143#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
144#[serde(rename_all = "lowercase")]
145pub enum TaskStatus {
146 /// Waiting to be claimed.
147 Queued,
148 /// Claimed by a daemon; a run is in flight.
149 Running,
150 /// A run finished and its gate passed.
151 Done,
152 /// A run finished without passing, and attempts remain.
153 Failed,
154 /// Out of attempts, or held by hand. The loop will not pick it up.
155 Held,
156 /// Waiting on another task or an unanswered question. See
157 /// [`Task::blocked_by`]. Set and cleared by `crate::conduct` and
158 /// `crate::daemon`'s deterministic resolver, never by hand.
159 Blocked,
160}
161
162impl TaskStatus {
163 /// Is this task eligible for a daemon to claim?
164 pub fn runnable(self) -> bool {
165 matches!(self, Self::Queued | Self::Failed)
166 }
167
168 /// Lowercase name, as it appears on disk and in the API.
169 pub fn as_str(self) -> &'static str {
170 match self {
171 Self::Queued => "queued",
172 Self::Running => "running",
173 Self::Done => "done",
174 Self::Failed => "failed",
175 Self::Held => "held",
176 Self::Blocked => "blocked",
177 }
178 }
179}
180
181/// How many tasks sit in each [`TaskStatus`], for a dashboard tile — never a
182/// per-task view, so it carries no ids.
183#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
184pub struct TaskCounts {
185 /// Waiting to be claimed.
186 pub queued: usize,
187 /// Claimed; a run is in flight.
188 pub running: usize,
189 /// A run finished and its gate passed.
190 pub done: usize,
191 /// A run finished without passing, and attempts remain.
192 pub failed: usize,
193 /// Out of attempts, or held by hand.
194 pub held: usize,
195 /// Waiting on another task or an unanswered question.
196 pub blocked: usize,
197}
198
199impl TaskCounts {
200 /// Tally `tasks` by status. A pure function over whatever [`Queue::list`]
201 /// already read, so it needs no root of its own and stays trivially
202 /// testable against a hand-built slice.
203 pub fn of(tasks: &[Task]) -> Self {
204 let mut counts = Self::default();
205 for t in tasks {
206 match t.status {
207 TaskStatus::Queued => counts.queued += 1,
208 TaskStatus::Running => counts.running += 1,
209 TaskStatus::Done => counts.done += 1,
210 TaskStatus::Failed => counts.failed += 1,
211 TaskStatus::Held => counts.held += 1,
212 TaskStatus::Blocked => counts.blocked += 1,
213 }
214 }
215 counts
216 }
217}
218
219/// One unit of work.
220#[derive(Debug, Clone, Serialize, Deserialize)]
221#[serde(deny_unknown_fields)]
222pub struct Task {
223 /// On-disk format version.
224 pub schema: u32,
225 /// Task id, e.g. `20260902-140501-a1b2`.
226 pub id: String,
227 /// One line, for lists and notifications.
228 pub title: String,
229 /// The task itself, handed to the graph verbatim.
230 pub instruction: String,
231 /// Repository to work in.
232 pub repo: PathBuf,
233 /// Who asked.
234 pub source: Source,
235 /// Higher runs first; ties break oldest-first so nothing starves.
236 #[serde(default)]
237 pub priority: i32,
238 /// Run this task alone: one implementer, no panel of judges to convince.
239 ///
240 /// `#[serde(default)]` so a queue file written before this field existed
241 /// still reads, as `false` - the ordinary multi-candidate competition,
242 /// unchanged. A task set to `solo` still runs the whole graph; only the
243 /// candidate count the daemon builds it with changes, and
244 /// [`crate::graph::Runner`] already collapses a single-candidate run to
245 /// implement → review → gate → merge on its own (see
246 /// [`crate::graph::Runner::review`]'s doc), so nothing about judging,
247 /// deliberation or voting had to change to support this.
248 #[serde(default)]
249 pub solo: bool,
250 /// Current state.
251 pub status: TaskStatus,
252 /// How many times this task has been claimed.
253 #[serde(default)]
254 pub attempts: usize,
255 /// Runs this task has produced, oldest first.
256 #[serde(default)]
257 pub runs: Vec<String>,
258 /// Why the last attempt did not land.
259 #[serde(default)]
260 pub last_error: Option<String>,
261 /// What a human hold is waiting on.
262 ///
263 /// `None` covers both the ordinary cases: a hold the loop makes itself
264 /// (out of attempts, or the disk gate closed) explains itself through
265 /// [`Task::last_error`] instead, and a human hold nobody bothered to
266 /// explain is still a valid hold. The queue has no way to express a
267 /// dependency between two tasks, so on the occasions a hold really is
268 /// "wait for that other task first", this is the only place that reason
269 /// survives - see [`Task::hold_manual`] and [`Task::release`].
270 ///
271 /// `#[serde(default)]` so a queue file written before this field existed
272 /// still reads, with no reason recorded rather than a parse error.
273 #[serde(default)]
274 pub hold_reason: Option<String>,
275 /// Who placed [`Task::hold_reason`]. `None` is a compatible old record;
276 /// see [`Task::operator_held`] for its deliberately conservative meaning.
277 #[serde(default)]
278 pub hold_source: Option<HoldSource>,
279 /// Diagnostic detail excerpted from the run that led to a hold - what a
280 /// human would have found opening `artifacts/` by hand, not the one-line
281 /// reason in [`Task::last_error`]. Set only when a run's own attempts are
282 /// exhausted and the task becomes [`TaskStatus::Held`]; `daemon` computes
283 /// it from the run's own record, since this module has no notion of a
284 /// run's internals. Bounded in length by the writer - see
285 /// `daemon::diagnostic` - so a verbose run cannot make this file grow
286 /// without limit.
287 ///
288 /// `#[serde(default)]` so a queue file written before this field existed
289 /// still reads, with no diagnostic recorded rather than a parse error.
290 #[serde(default)]
291 pub diagnostic: Option<String>,
292 /// What this task is waiting on: other task ids, unanswered
293 /// `crate::ask::Question` ids, or both. Non-empty exactly when
294 /// [`TaskStatus::Blocked`]; emptying it — see [`Task::unblock`] — is what
295 /// puts the task back at [`TaskStatus::Queued`].
296 ///
297 /// Set by `crate::conduct`'s decisions and cleared deterministically by
298 /// `crate::daemon` as each dependency resolves, never by a person. Never
299 /// `#[serde(default)]` is skipped: a queue file from before this field
300 /// existed has nothing to report here, and an empty list is exactly that.
301 #[serde(default)]
302 pub blocked_by: Vec<String>,
303 /// One line explaining the current [`Task::blocked_by`], written by
304 /// `crate::conduct`. Cleared whenever `blocked_by` empties.
305 #[serde(default)]
306 pub block_reason: Option<String>,
307 /// The status this task had the moment [`Task::block`] most recently
308 /// moved it to [`TaskStatus::Blocked`] — what [`Task::unblock`] restores
309 /// once nothing is left in `blocked_by`, instead of always landing on
310 /// [`TaskStatus::Queued`]. See [`SCHEMA`]'s doc for schema 4 on why this
311 /// exists: an answer to a question `crate::conduct` filed about a
312 /// [`TaskStatus::Held`] task must not itself be what puts the task back
313 /// in the competition queue.
314 ///
315 /// `#[serde(default)]` so a queue file written before this field existed
316 /// reads as `None`; [`Task::unblock`] treats that the same as a task
317 /// blocked straight from `Queued`, unless surviving hold evidence says
318 /// otherwise.
319 #[serde(default)]
320 pub blocked_from: Option<TaskStatus>,
321 /// Questions `crate::conduct` asked about this task that the operator has
322 /// since answered, oldest first — what was asked, and what they said.
323 ///
324 /// A blocking question's id leaves [`Task::blocked_by`] the moment
325 /// [`crate::ask::QuestionStatus::Answered`] is observed, but the id alone
326 /// tells nobody what was decided. This is what carries the answer's
327 /// *content* forward: into the next conductor prompt for this task, and
328 /// into the instruction handed to the next run — see `crate::daemon`'s
329 /// deterministic blocker resolution. Kept for the task's whole life, the
330 /// same as [`Task::runs`]: a release resets attempts, not evidence.
331 #[serde(default)]
332 pub answers: Vec<AnsweredQuestion>,
333 /// Ids of the `crate::triage` questions whose answer has been applied to
334 /// this task. Unlike [`Task::hold_reason`], [`Task::release`] and every
335 /// hold transition leave it alone, so an answer is applied at most once
336 /// however many times the task is held again. See [`SCHEMA`]'s doc for
337 /// schema 5. `#[serde(default)]` so an older record reads as empty.
338 #[serde(default)]
339 pub triage_applied: Vec<String>,
340 /// The operator's "resume" answer to a triage question, kept until the
341 /// task actually runs (or is done) so `crate::conduct` cannot silently
342 /// undo it and `crate::triage` can tell that a hold it sees now came
343 /// *after* the answer. See [`OperatorResume`]. `#[serde(default)]`.
344 #[serde(default)]
345 pub resume_override: Option<OperatorResume>,
346 /// Set by `crate::conduct` when it chooses `Review` recovery for a task
347 /// whose branch survived a blocked run: the branch to reopen with
348 /// `crate::graph::Runner::review` instead of competing from scratch.
349 ///
350 /// Requeues the task the same way [`Task::release`] does, so it is
351 /// picked up by the ordinary loop; `crate::daemon` reads this field once,
352 /// when it actually starts the run, and clears it either way — consumed
353 /// on success, dropped if the branch no longer exists by then. Never set
354 /// from the conductor's own words: `crate::daemon` derives the branch
355 /// name itself from the task's last run, so a hallucinated branch can
356 /// never reach here.
357 #[serde(default)]
358 pub review_branch: Option<String>,
359 /// A release deliberately starts a new competition instead of resuming
360 /// the prior run. History remains as evidence in `runs`.
361 #[serde(default)]
362 pub fresh_start: bool,
363 /// Marked by an operator (`magi task interrupt`) to ask `magi serve` to
364 /// run this one ahead of whatever it already has in flight, once
365 /// `[daemon] pause_for_interrupts` is on - see
366 /// `crate::daemon::advance_interrupt`. Never set by the loop itself, and
367 /// deliberately a different operation from [`Task::set_priority`]: a
368 /// priority only reorders the queue a claim has not reached yet, while
369 /// this asks a run already in flight to park at its next safe boundary
370 /// and step aside. `#[serde(default)]` so a queue file written before
371 /// this field existed still reads, as `false` - no task interrupts
372 /// anything unless asked to, exactly as before.
373 #[serde(default)]
374 pub interrupt: bool,
375 /// Marked by `magi task add --urgent`: `crate::daemon::poll` dispatches
376 /// this task through its own one-slot `urgent_sem` the moment it is
377 /// runnable, in addition to whatever is already running under the
378 /// ordinary `[daemon] max_concurrent_runs` pool - never instead of it,
379 /// and never by pausing or otherwise touching that run. This is the
380 /// opposite direction from [`Task::interrupt`]: that one asks a run
381 /// already in flight to step aside; this one never asks anything to
382 /// step aside, it only spends one additional, temporary concurrency
383 /// slot. The two are independent and may both be set on the same task,
384 /// but this exemption stops at `[daemon] pause_for_interrupts`'s own
385 /// park/resume handoff (75dd): while an interrupt sequence is actively
386 /// parking, running, or resuming - its own, or an unrelated task's -
387 /// `crate::daemon::interrupt_gate` withholds an urgent candidate exactly
388 /// like an ordinary one, never exempted. 75dd's "at most one run, ever,
389 /// at once" guarantee takes precedence, because the alternative is a run
390 /// still genuinely in flight (only *asked* to park, not yet gone) ending
391 /// up alongside a second one this feature let through - the very thing
392 /// that guarantee exists to rule out.
393 ///
394 /// `#[serde(default)]` so a queue file written before this field existed
395 /// still reads, as `false` - no task claims the urgent slot unless asked
396 /// to, exactly as before.
397 #[serde(default)]
398 pub urgent: bool,
399 /// When the task was filed.
400 pub created_at: Timestamp,
401 /// Last change to this file.
402 pub updated_at: Timestamp,
403}
404
405/// A triage "resume" answer and what became of it. See
406/// [`Task::resume_override`].
407#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
408pub struct OperatorResume {
409 /// The triage question the operator answered.
410 pub question_id: String,
411 /// When the answer was applied.
412 pub at: Timestamp,
413 /// The reason `crate::conduct` gave for holding the task again after the
414 /// answer, if it did. The conductor may do this once.
415 #[serde(default)]
416 pub conductor_rehold: Option<String>,
417 /// The operator answered "resume" a second time, to the question about
418 /// that contradiction: the conductor may no longer hold this task.
419 #[serde(default)]
420 pub forced: bool,
421}
422
423/// One question `crate::conduct` asked about a task, and what the operator
424/// said back. See [`Task::answers`].
425#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
426pub struct AnsweredQuestion {
427 /// The question as asked, e.g. [`crate::ask::Question::summary`].
428 pub question: String,
429 /// What the operator answered.
430 pub answer: String,
431}
432
433impl Task {
434 /// File a new task. Persist it with [`Queue::put`].
435 pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
436 let now = Timestamp::now();
437 Self {
438 schema: SCHEMA,
439 id: new_id(),
440 title,
441 instruction,
442 repo,
443 source,
444 priority: 0,
445 solo: false,
446 status: TaskStatus::Queued,
447 attempts: 0,
448 runs: Vec::new(),
449 last_error: None,
450 hold_reason: None,
451 hold_source: None,
452 diagnostic: None,
453 blocked_by: Vec::new(),
454 block_reason: None,
455 blocked_from: None,
456 answers: Vec::new(),
457 triage_applied: Vec::new(),
458 resume_override: None,
459 review_branch: None,
460 fresh_start: false,
461 interrupt: false,
462 urgent: false,
463 created_at: now,
464 updated_at: now,
465 }
466 }
467
468 /// Short form used in reports, matching a run's short id.
469 pub fn short(&self) -> &str {
470 short(&self.id)
471 }
472
473 /// Record that triage question `question_id`'s answer has been applied.
474 pub fn mark_triage_applied(&mut self, question_id: &str) {
475 if !self.triage_applied(question_id) {
476 self.triage_applied.push(question_id.to_owned());
477 }
478 }
479
480 /// Has triage question `question_id`'s answer already been applied?
481 pub fn triage_applied(&self, question_id: &str) -> bool {
482 self.triage_applied.iter().any(|id| id == question_id)
483 }
484
485 /// Record that a run has started for this task.
486 ///
487 /// Clears [`Task::interrupt`]: a mark to run ahead of whatever else is
488 /// in flight is fulfilled the moment this task actually gets its turn,
489 /// dispatched same as any other. Without this, a task whose run fails
490 /// and requeues - still `runnable`, still carrying the mark from its
491 /// first attempt - would keep re-triggering `crate::daemon`'s interrupt
492 /// scheduler and re-parking whatever it interrupted on every later
493 /// boundary, for as long as its attempts hold out, instead of the
494 /// one-shot "let this go next" the mark is meant to be.
495 pub fn start(&mut self, run: String) {
496 self.status = TaskStatus::Running;
497 self.attempts += 1;
498 self.runs.push(run);
499 self.last_error = None;
500 self.fresh_start = false;
501 self.interrupt = false;
502 // The answer has been honoured: the task got its turn.
503 self.resume_override = None;
504 }
505
506 /// Record a successful run.
507 ///
508 /// Both `magi task done` and `POST /api/queue/{id}/done` can close a held
509 /// *or blocked* task directly, with no release in between, so this clears
510 /// `hold_reason` and `blocked_by`/`block_reason` the same way
511 /// [`Task::release`] does. Otherwise a task held for "waiting on 3ed9", or
512 /// blocked on a dependency that never actually finished, and then closed
513 /// as done without ever being released would still read as waiting on
514 /// something in `magi task show` and on its card, after it no longer is.
515 pub fn succeed(&mut self) {
516 self.status = TaskStatus::Done;
517 self.resume_override = None;
518 self.last_error = None;
519 self.hold_reason = None;
520 self.hold_source = None;
521 self.diagnostic = None;
522 self.blocked_by.clear();
523 self.block_reason = None;
524 self.blocked_from = None;
525 }
526
527 /// Earlier attempts at this task that a later one has since made moot —
528 /// empty unless the task is [`TaskStatus::Done`] *and* `last_run_succeeded`
529 /// says `runs.last()` is actually why.
530 ///
531 /// `runs` is oldest first, and [`Task::start`] is the only thing that
532 /// pushes to it, always right before the attempt it names either succeeds
533 /// or fails; `succeed` itself never touches `runs`. So whenever `status`
534 /// is `Done` *because* the loop itself saw that last attempt land
535 /// (`Merged`/`Ready`), everything before it in the same list is a retry
536 /// this task no longer needs. But `succeed` is also reachable directly —
537 /// `magi task done`, the web UI's equivalent, and the conductor's
538 /// `Recovery::Done` all call it on a task in *any* status, including one
539 /// whose last recorded attempt never landed at all (closed by hand after
540 /// a merge magi's own loop never saw). `runs.last()` alone cannot tell
541 /// those two cases apart — that requires the caller to have actually
542 /// looked at that run's own `status`, which this module has no way to
543 /// do — so `last_run_succeeded` is the caller's answer to exactly that
544 /// question, not something this function can derive from `Task` alone.
545 ///
546 /// Deliberately narrower than [`Queue::superseded`]'s "every earlier
547 /// attempt has a later one" walk, which fires the moment a retry starts
548 /// even though the retry itself might still fail: that reading is right
549 /// for the web UI's "a newer attempt exists, go look at that one
550 /// instead" note, but wrong for deciding a run no longer needs a human's
551 /// attention, which is only true once the task's story has actually
552 /// ended well. Two Blocked runs sitting side by side while a third
553 /// attempt is still in flight must not be touched by this — see
554 /// `daemon::supersede_prior_runs`, the caller that turns this list into
555 /// rewritten `run.json` files.
556 pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
557 if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
558 return &[];
559 }
560 &self.runs[..self.runs.len() - 1]
561 }
562
563 /// The attempt right after `run` in this task's history, if there is one.
564 ///
565 /// The one place "a later attempt replaced this run" is decided:
566 /// [`Queue::superseded`] and [`Queue::superseded_by`] (the web UI's
567 /// note) read it, and so does [`Task::earlier_attempts`], so what is shown
568 /// and what `crate::handover` acts on cannot disagree.
569 pub fn successor_of(&self, run: &str) -> Option<&String> {
570 let pos = self.runs.iter().position(|r| r == run)?;
571 self.runs.get(pos + 1)
572 }
573
574 /// Every run already recorded for this task, i.e. every run that an
575 /// attempt being minted *right now* supersedes.
576 ///
577 /// The new run is not in [`Task::runs`] yet when it decides what to take
578 /// over (`Task::start` pushes it afterwards), so "has a later attempt" is
579 /// true of everything listed here by the time the new run exists.
580 pub fn earlier_attempts(&self) -> &[String] {
581 &self.runs
582 }
583
584 /// Record a failed attempt. Out of attempts means held for a human, rather
585 /// than retried until the money runs out.
586 ///
587 /// Clears [`Task::diagnostic`] unconditionally: it belongs to whatever run
588 /// produced it, and a caller that has one for *this* attempt sets it
589 /// itself right after calling this, once it knows the task actually ended
590 /// up [`TaskStatus::Held`] - see `daemon::diagnostic`. Without the clear, a
591 /// task released after a diagnosed hold and then failed again for an
592 /// unrelated, undiagnosed reason (a config error, say) would go on
593 /// showing the previous run's diagnostic as if it explained the new one.
594 ///
595 /// Also records `why` into [`Task::hold_reason`] when this ends up
596 /// [`TaskStatus::Held`], so the notification and `magi task show` say the
597 /// same thing as [`Task::last_error`] instead of leaving the hold's own
598 /// reason blank.
599 pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
600 let why = why.into();
601 self.diagnostic = None;
602 self.status = if self.attempts >= max_attempts {
603 self.hold_source = Some(HoldSource::Machine);
604 self.hold_reason = Some(why.clone());
605 TaskStatus::Held
606 } else {
607 TaskStatus::Failed
608 };
609 self.last_error = Some(why);
610 }
611
612 /// Record an attempt that failed for a reason the task is not responsible
613 /// for - the agent CLIs ran out of quota and the judging panel collapsed.
614 ///
615 /// This refunds the attempt on purpose. A quota window closing at 4am must
616 /// not spend the backlog's retry budget: the operator would come back to a
617 /// queue of held tasks that were never actually judged, and would have to
618 /// release every one by hand to find out which had a real problem. The task
619 /// goes back to `Failed`, which the loop retries, so a reset quota picks the
620 /// work up where it stopped.
621 pub fn stall(&mut self, why: impl Into<String>) {
622 self.last_error = Some(why.into());
623 self.diagnostic = None;
624 self.attempts = self.attempts.saturating_sub(1);
625 self.status = TaskStatus::Failed;
626 }
627
628 /// Whether this held task may only be released by an operator.
629 ///
630 /// Old files did not record a source. Preserve every such hold rather
631 /// than guessing that it was automatic and risking duplicate work. New
632 /// automatic holds record [`HoldSource::Machine`] and remain recoverable.
633 pub fn operator_held(&self) -> bool {
634 self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
635 }
636
637 /// Take this task out of the loop's reach by an operator action.
638 ///
639 /// Clears `blocked_by`/`block_reason` unconditionally, the same as
640 /// [`Task::release`] and for the same reason its own comment already
641 /// gives: a human choosing to hold a *blocked* task overrides its wait
642 /// outright, the same as it overrides an ordinary hold. Without this, a
643 /// task held straight out of [`TaskStatus::Blocked`] - the web UI's "Hold"
644 /// button is reachable on a blocked task, same as "Mark done" - kept
645 /// reading as still waiting on a dependency it no longer had any claim on.
646 pub fn hold_manual(&mut self, reason: Option<String>) {
647 self.status = TaskStatus::Held;
648 if reason.is_some() {
649 self.hold_reason = reason;
650 }
651 self.hold_source = Some(HoldSource::Manual);
652 self.blocked_by.clear();
653 self.block_reason = None;
654 self.blocked_from = None;
655 }
656
657 /// Take this task out of the loop's reach during automatic recovery.
658 ///
659 /// Clears `blocked_by`/`block_reason` for the same reason
660 /// [`Task::hold_manual`] does.
661 pub fn hold_machine(&mut self, reason: Option<String>) {
662 self.status = TaskStatus::Held;
663 if reason.is_some() {
664 self.hold_reason = reason;
665 }
666 self.hold_source = Some(HoldSource::Machine);
667 self.blocked_by.clear();
668 self.block_reason = None;
669 self.blocked_from = None;
670 }
671
672 /// Block this task on other task ids and/or open question ids, chosen by
673 /// `crate::conduct`. Pure: the caller still owns writing it back with
674 /// [`Queue::put`].
675 ///
676 /// Records [`Task::blocked_from`] the first time this moves the task into
677 /// [`TaskStatus::Blocked`], and leaves it alone on a later call that adds
678 /// or replaces `blocked_by` while the task is already `Blocked` - a
679 /// second question about an already-blocked task must not overwrite the
680 /// status it should eventually return to with `Blocked` itself.
681 pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
682 if self.status != TaskStatus::Blocked {
683 self.blocked_from = Some(self.status);
684 }
685 self.status = TaskStatus::Blocked;
686 self.blocked_by = blocked_by;
687 self.block_reason = reason;
688 }
689
690 /// Remove one resolved dependency (a task id that became [`TaskStatus::Done`],
691 /// or a question id that became [`crate::ask::QuestionStatus::Answered`]).
692 /// Once nothing is left in [`Task::blocked_by`], the task returns to
693 /// whatever [`Task::blocked_from`] recorded - deciding *why* a task was
694 /// blocked was `crate::conduct`'s job, but noticing a dependency resolved
695 /// needs no model at all, and restoring the status it interrupted needs
696 /// nothing more than what `block` already wrote down.
697 ///
698 /// A task blocked while `Running` restores to [`TaskStatus::Queued`]
699 /// instead: whatever process was running it is gone by the time this
700 /// runs, so there is nothing left to resume. A task with no recorded
701 /// `blocked_from` - a pre-schema-4 record, or one blocked before this
702 /// field existed - falls back to [`TaskStatus::Held`] when it still
703 /// carries hold evidence ([`Task::hold_reason`] or [`Task::hold_source`],
704 /// neither ever cleared by `block`), and to `Queued` otherwise: the same
705 /// choice `block` itself would have recorded, reconstructed from what
706 /// survived.
707 ///
708 /// A no-op, on purpose, for a task that is not [`TaskStatus::Blocked`]:
709 /// `crate::daemon`'s deterministic resolver runs over every task on every
710 /// poll, and a task that moved on for some other reason must not be
711 /// dragged back by a stale id it still happens to carry.
712 pub fn unblock(&mut self, resolved_id: &str) {
713 if self.status != TaskStatus::Blocked {
714 return;
715 }
716 self.blocked_by.retain(|id| id != resolved_id);
717 if self.blocked_by.is_empty() {
718 self.status = match self.blocked_from {
719 Some(TaskStatus::Running) => TaskStatus::Queued,
720 Some(other) => other,
721 None if self.hold_reason.is_some() || self.hold_source.is_some() => {
722 TaskStatus::Held
723 }
724 None => TaskStatus::Queued,
725 };
726 self.block_reason = None;
727 self.blocked_from = None;
728 }
729 }
730
731 /// Record that a question `crate::conduct` asked about this task has been
732 /// answered, so the answer's content — not just the fact that the
733 /// question is gone — reaches the next conductor prompt and the next
734 /// run's instruction. See [`Task::answers`].
735 pub fn record_answer(&mut self, question: String, answer: String) {
736 self.answers.push(AnsweredQuestion { question, answer });
737 }
738
739 /// Requeue this task to reopen its last run as a review-only pass against
740 /// `branch` (`crate::graph::Runner::review`) rather than competing from
741 /// scratch. See [`Task::review_branch`].
742 pub fn request_review(&mut self, branch: String) {
743 self.release();
744 self.review_branch = Some(branch);
745 }
746
747 /// Requeue after a conductor chose a new competition. Unlike an ordinary
748 /// operator release, this deliberately does not resume the old run.
749 pub fn requeue(&mut self) {
750 self.release();
751 self.fresh_start = true;
752 }
753
754 /// Change how urgently this task should run next.
755 ///
756 /// Refused once the task is `running`: priority only feeds the sort
757 /// [`Queue::next_runnable`] does over tasks waiting to be claimed, and a
758 /// running task has already left that pool. Accepting the write anyway
759 /// would look like it worked while changing nothing until - and unless -
760 /// this attempt fails and the task becomes runnable again, which is a
761 /// surprise the phone should not hand back as a success.
762 pub fn set_priority(&mut self, priority: i32) -> Result<()> {
763 if self.status == TaskStatus::Running {
764 bail!(
765 "task {} is running; its priority cannot be changed until \
766 this attempt finishes",
767 self.short()
768 );
769 }
770 self.priority = priority;
771 Ok(())
772 }
773
774 /// Mark (or unmark) this task to interrupt whatever `magi serve` already
775 /// has in flight, once `[daemon] pause_for_interrupts` is on. See
776 /// [`Task::interrupt`].
777 ///
778 /// Setting it is restricted to a task the loop could pick up on its own
779 /// right now - [`TaskStatus::runnable`] - for the same reason as
780 /// [`Task::set_priority`]: a task already `running` has been claimed, and
781 /// a task that is `done`, `held`, or `blocked` is not going to compete
782 /// for the daemon's attention regardless of this flag. Unlike priority,
783 /// this is never silently inert while `running` - it is refused outright,
784 /// because the entire feature this flag drives (`crate::daemon`'s
785 /// interrupt scheduler) is scoped to tasks still waiting to be claimed.
786 /// Clearing it back to `false` carries no such risk and is always
787 /// allowed, including on a task that moved on since it was set.
788 pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
789 if interrupt && !self.status.runnable() {
790 bail!(
791 "task {} is {}; only a queued or failed task can be marked \
792 to interrupt",
793 self.short(),
794 self.status.as_str()
795 );
796 }
797 self.interrupt = interrupt;
798 Ok(())
799 }
800
801 /// Replace this task's title and instruction wholesale.
802 ///
803 /// Restricted to `queued` and `held`. A `running` task's instruction has
804 /// already been handed to the graph, so a run in flight and the file on
805 /// disk must not be allowed to disagree about what was asked; a `done` or
806 /// `failed` task is a record of what actually happened and editing it
807 /// after the fact would falsify that record. `id`, `created_at`,
808 /// `source`, and `runs` are left untouched on purpose - an edit stands in
809 /// for "delete and refile", and keeping the id, the timestamp, the
810 /// attribution, and the run history is the entire reason it exists
811 /// instead.
812 pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
813 if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
814 bail!(
815 "task {} is {}; only a queued or held task's instruction can \
816 be edited",
817 self.short(),
818 self.status.as_str()
819 );
820 }
821 self.title = title;
822 self.instruction = instruction;
823 Ok(())
824 }
825
826 /// Record a run that produced a pull request without merging it.
827 ///
828 /// The task is held rather than retried, and it costs no further attempt
829 /// either way. The work the task asked for exists: it is sitting on a
830 /// branch, in a pull request, waiting for CI or for a person. Retrying
831 /// would spend the whole competition budget a second time and then race a
832 /// second branch against the pull request the first one opened - which is
833 /// exactly what happened to run 01c2, whose finished and green pull request
834 /// was re-competed from scratch four seconds after it opened.
835 ///
836 /// A pull request nobody merged is a request for a person, not a failure.
837 ///
838 /// Records `why` into [`Task::hold_reason`] as well as
839 /// [`Task::last_error`], so the notification centre and `magi task show`
840 /// say why the task is held rather than "no reason recorded". This is
841 /// also the path a verified no-op with an outstanding `magi ask` question
842 /// settles through - see `daemon::settle` and `daemon::settle_and_diagnose`,
843 /// which append the question id to `hold_reason` when one is still open
844 /// for the run.
845 pub fn handed_off(&mut self, why: impl Into<String>) {
846 let why = why.into();
847 self.diagnostic = None;
848 self.status = TaskStatus::Held;
849 self.hold_source = Some(HoldSource::Machine);
850 self.hold_reason = Some(why.clone());
851 self.last_error = Some(why);
852 }
853
854 /// Put a held or finished task back in line, with its attempt count reset
855 /// so a release is a real second chance rather than an instant re-hold.
856 /// The run history is kept: attempts reset, evidence does not.
857 pub fn release(&mut self) {
858 self.status = TaskStatus::Queued;
859 self.attempts = 0;
860 self.last_error = None;
861 // Otherwise the next person who holds this task reads a reason that
862 // belonged to whatever it was waiting on last time.
863 self.hold_reason = None;
864 self.hold_source = None;
865 self.diagnostic = None;
866 // A release also un-blocks: the dependency or question `blocked_by`
867 // named may still be unresolved, but a human (or `crate::conduct`)
868 // choosing to release the task overrides that wait outright, the same
869 // as it overrides an ordinary hold.
870 self.blocked_by.clear();
871 self.block_reason = None;
872 self.blocked_from = None;
873 self.review_branch = None;
874 self.fresh_start = false;
875 }
876}
877
878/// A queue on disk.
879#[derive(Debug, Clone)]
880pub struct Queue {
881 root: PathBuf,
882}
883
884impl Queue {
885 /// The operator's queue, `<home>/queue`.
886 pub fn open() -> Self {
887 Self::at(crate::run::home().join("queue"))
888 }
889
890 /// A queue at an explicit root. Tests use this; so could an operator who
891 /// wants a queue per project.
892 pub fn at(root: PathBuf) -> Self {
893 Self { root }
894 }
895
896 /// Directory holding the task files.
897 pub fn root(&self) -> &Path {
898 &self.root
899 }
900
901 /// Path for one task id.
902 pub fn path_of(&self, id: &str) -> PathBuf {
903 self.root.join(format!("{id}.json"))
904 }
905
906 /// Write a task, atomically, so a daemon killed mid-write leaves the
907 /// previous state readable rather than a truncated file.
908 pub fn put(&self, task: &mut Task) -> Result<()> {
909 task.updated_at = Timestamp::now();
910 std::fs::create_dir_all(&self.root)
911 .with_context(|| format!("create {}", self.root.display()))?;
912 let body = serde_json::to_string_pretty(task).context("serialize task")?;
913 let path = self.path_of(&task.id);
914 let tmp = path.with_extension("json.tmp");
915 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
916 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
917 // A machine hold is news for the notification centre, filed beside
918 // this queue (`<home>/queue` -> `<home>/notifications`) rather than
919 // through a process-global, so a queue in a temp directory notifies
920 // into that directory.
921 if let (Some(notice), Some(home)) = (
922 crate::notices::task_held(task),
923 self.root.parent().filter(|p| !p.as_os_str().is_empty()),
924 ) {
925 crate::notices::raise_in(home, notice);
926 }
927 Ok(())
928 }
929
930 /// Load a task by id or unambiguous id prefix.
931 pub fn get(&self, id: &str) -> Result<Task> {
932 let resolved = self.resolve_id(id)?;
933 read_path(&self.path_of(&resolved))
934 }
935
936 /// Remove a task, and the claim lock that belongs to it.
937 ///
938 /// `in_flight` comes from the caller — a live daemon's heartbeat naming
939 /// this task — because the task's own `running` status cannot answer the
940 /// question. A daemon killed mid-competition leaves the status at
941 /// `running` and an orphaned `.lock` behind, and a guard that trusted
942 /// either would make the task undeletable for good: the phone showed
943 /// exactly that, refusing a task whose daemon had been gone for an hour.
944 ///
945 /// So the lock is removed with the task rather than respected. Any lock
946 /// still there once no live daemon claims the task is by definition stale,
947 /// and leaving it would make a deleted task look claimed to
948 /// [`Queue::claim`] and to whoever reads the directory.
949 ///
950 /// Anything still `blocked` on the id just deleted is quarantined to a
951 /// machine hold in the same call - see [`Removal::quarantined`] - rather
952 /// than left to wait on a dependency that no longer exists. Best-effort:
953 /// a dependent claimed by something else right now, or one whose write
954 /// fails, is simply left for `crate::daemon::resolve_blockers`'s own poll
955 /// (or `crate::triage::run_once`) to catch on its own next pass, and does
956 /// not fail this removal.
957 ///
958 /// `questions` is the store [`missing_blockers`] checks a `blocked_by` id
959 /// against before calling it gone - the same store the caller already
960 /// resolves `id`'s own home from, passed in rather than reopened here so
961 /// a test queue at an explicit root is never quarantined against the
962 /// operator's real questions directory.
963 pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
964 let resolved = self.resolve_id(id)?;
965 if in_flight {
966 bail!("task {resolved} is being run by a live daemon right now");
967 }
968 let path = self.path_of(&resolved);
969 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
970 let lock = self.lock_path(&resolved);
971 if let Err(e) = std::fs::remove_file(&lock) {
972 if e.kind() != std::io::ErrorKind::NotFound {
973 return Err(e).with_context(|| format!("remove {}", lock.display()));
974 }
975 }
976 let quarantined = self.quarantine_dependents_of(&resolved, questions);
977 Ok(Removal {
978 id: resolved,
979 quarantined,
980 })
981 }
982
983 /// Move every `blocked` task naming `dependency` in its own `blocked_by`
984 /// to a machine hold, now that `dependency`'s own file is gone. See
985 /// [`Queue::remove`]'s own doc for why this is best-effort.
986 fn quarantine_dependents_of(&self, dependency: &str, questions: &Questions) -> Vec<String> {
987 let mut quarantined = Vec::new();
988 for listed in self.list() {
989 if listed.status != TaskStatus::Blocked
990 || !listed.blocked_by.iter().any(|b| b == dependency)
991 {
992 continue;
993 }
994 let Ok(_claim) = self.claim(&listed.id) else {
995 continue;
996 };
997 let Ok(mut task) = self.get(&listed.id) else {
998 continue;
999 };
1000 if task.status != TaskStatus::Blocked
1001 || !task.blocked_by.iter().any(|b| b == dependency)
1002 {
1003 continue;
1004 }
1005 let missing = missing_blockers(self, questions, &task.blocked_by);
1006 task.hold_machine(Some(missing_blocker_hold_reason(
1007 &task.blocked_by,
1008 &missing,
1009 )));
1010 if self.put(&mut task).is_ok() {
1011 quarantined.push(task.id.clone());
1012 }
1013 }
1014 quarantined
1015 }
1016
1017 /// Path of the claim lock for a task. One definition, so `claim` and
1018 /// `remove` cannot end up naming different files.
1019 fn lock_path(&self, id: &str) -> PathBuf {
1020 self.root.join(format!("{id}.lock"))
1021 }
1022
1023 /// Every task on disk, highest priority first and newest first within a
1024 /// priority. This is what `magi task list` and `GET /api/queue` print, so
1025 /// a raised priority has to move a task here the moment it is saved, not
1026 /// only in [`Queue::next_runnable`]'s own ordering - the operator reading
1027 /// the backlog and the loop about to drain it must agree on what "first"
1028 /// means. Every existing task defaults to priority 0, so this is a no-op
1029 /// change from the old newest-first order for a queue nobody has
1030 /// reprioritised.
1031 ///
1032 /// Unreadable files are skipped rather than fatal: one corrupt task must
1033 /// not take the queue - or the web UI, or an unattended daemon - down
1034 /// with it.
1035 pub fn list(&self) -> Vec<Task> {
1036 let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1037 .into_iter()
1038 .flatten()
1039 .flatten()
1040 .map(|e| e.path())
1041 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1042 .filter_map(|p| read_path(&p).ok())
1043 .collect();
1044 tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1045 tasks
1046 }
1047
1048 /// Runs that a later attempt at the same task replaced, mapped to the id
1049 /// of the attempt that replaced them.
1050 ///
1051 /// A task keeps its attempts in order, and the deck showed them as two
1052 /// cards with the same title and no hint which was which: yukimemi asked
1053 /// why `stalled` and `blocked` appeared twice for one task, and the
1054 /// answer - "those are two tries, and the second one exists because of a
1055 /// bug since fixed" - was not on the screen anywhere.
1056 ///
1057 /// Read from the queue rather than stored on the run, because the
1058 /// ordering is the queue's fact: a `RunState` has no idea another attempt
1059 /// happened after it.
1060 pub fn superseded(&self) -> HashMap<String, String> {
1061 let mut by = HashMap::new();
1062 for task in self.list() {
1063 for earlier in &task.runs {
1064 if let Some(later) = task.successor_of(earlier) {
1065 by.insert(earlier.clone(), later.clone());
1066 }
1067 }
1068 }
1069 by
1070 }
1071
1072 /// Whether `run` is an earlier attempt a later one replaced, and if so
1073 /// the id of that later attempt.
1074 ///
1075 /// Same walk as [`Queue::superseded`], narrowed to one run: a run detail
1076 /// page asks about exactly one run at a time, and this keeps that call
1077 /// site from building (and discarding) the whole map's `HashMap` just to
1078 /// read one entry out of it.
1079 pub fn superseded_by(&self, run: &str) -> Option<String> {
1080 for task in self.list() {
1081 if task.runs.iter().any(|r| r == run) {
1082 return task.successor_of(run).cloned();
1083 }
1084 }
1085 None
1086 }
1087
1088 /// The task's own most recent attempt, when `run` belongs to that task
1089 /// but is not already that attempt.
1090 ///
1091 /// Distinct from [`Queue::superseded_by`], which names only the very
1092 /// next attempt: a chain of retries (A superseded by B superseded by C)
1093 /// leaves an older run pointing at an intermediate one that may itself
1094 /// be unresolved, and a run's own detail page needs to know where the
1095 /// task's story currently stands - the chain's current head, C - not an
1096 /// attempt in the middle of it that a client would otherwise have to
1097 /// walk to by hand.
1098 pub fn latest_attempt(&self, run: &str) -> Option<String> {
1099 for task in self.list() {
1100 if task.runs.iter().any(|r| r == run) {
1101 return task.runs.last().filter(|last| **last != run).cloned();
1102 }
1103 }
1104 None
1105 }
1106
1107 /// The task a daemon should run next, or `None` when the queue is idle.
1108 ///
1109 /// Highest priority first, oldest first within a priority, so a burst of
1110 /// agent-filed work cannot starve the task a human filed this morning.
1111 pub fn next_runnable(&self) -> Option<Task> {
1112 let mut runnable: Vec<Task> = self
1113 .list()
1114 .into_iter()
1115 .filter(|t| t.status.runnable())
1116 .collect();
1117 runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1118 runnable.into_iter().next()
1119 }
1120
1121 /// Take exclusive ownership of a task.
1122 ///
1123 /// The lock is a `create_new` file next to the task, which is atomic on
1124 /// every platform magi targets. It exists so two daemons - or a daemon and
1125 /// a human running `magi run` - cannot drive one task into two competing
1126 /// runs. The returned guard releases on drop, including on panic.
1127 pub fn claim(&self, id: &str) -> Result<Claim> {
1128 std::fs::create_dir_all(&self.root)
1129 .with_context(|| format!("create {}", self.root.display()))?;
1130 let path = self.lock_path(id);
1131 match std::fs::OpenOptions::new()
1132 .write(true)
1133 .create_new(true)
1134 .open(&path)
1135 {
1136 Ok(mut f) => {
1137 use std::io::Write as _;
1138 // Best effort: the pid is for the human looking at a stale lock.
1139 let _ = writeln!(f, "{}", std::process::id());
1140 Ok(Claim { path })
1141 }
1142 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1143 bail!("task {id} is already claimed ({} exists)", path.display())
1144 }
1145 Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1146 }
1147 }
1148
1149 /// Expand an id prefix to exactly one task id.
1150 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1151 if self.path_of(prefix).is_file() {
1152 return Ok(prefix.to_owned());
1153 }
1154 let hits: Vec<String> = self
1155 .list()
1156 .into_iter()
1157 .map(|t| t.id)
1158 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1159 .collect();
1160 match hits.len() {
1161 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1162 0 => bail!("no task matches `{prefix}`"),
1163 _ => bail!(
1164 "`{prefix}` matches {} tasks: {}",
1165 hits.len(),
1166 hits.join(", ")
1167 ),
1168 }
1169 }
1170
1171 /// Change detection token for the queue.
1172 ///
1173 /// Combines file names and modification times of all task files in the
1174 /// queue, so adding, modifying, or deleting any task — even an older one —
1175 /// moves the revision and notifies connected clients via the change stream.
1176 /// Returns 0 when the queue is completely empty.
1177 pub fn revision(&self) -> u64 {
1178 use std::hash::{Hash as _, Hasher as _};
1179
1180 let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
1181 .into_iter()
1182 .flatten()
1183 .flatten()
1184 .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
1185 .filter_map(|e| {
1186 let name = e.file_name().to_string_lossy().into_owned();
1187 let mtime = e
1188 .metadata()
1189 .ok()?
1190 .modified()
1191 .ok()?
1192 .duration_since(std::time::UNIX_EPOCH)
1193 .ok()?
1194 .as_millis() as u64;
1195 Some((name, mtime))
1196 })
1197 .collect();
1198
1199 if entries.is_empty() {
1200 return 0;
1201 }
1202
1203 entries.sort_unstable();
1204 let mut hasher = std::hash::DefaultHasher::new();
1205 for (name, mtime) in &entries {
1206 name.hash(&mut hasher);
1207 mtime.hash(&mut hasher);
1208 }
1209 let h = hasher.finish();
1210 if h == 0 { 1 } else { h }
1211 }
1212}
1213
1214/// What [`Queue::remove`] did, beyond deleting the named task's own file.
1215#[derive(Debug, Clone)]
1216pub struct Removal {
1217 /// The id actually removed - `id` expanded from a prefix, if it was one.
1218 pub id: String,
1219 /// Every `blocked` task that named [`Removal::id`] in its own
1220 /// `blocked_by` and was moved to a machine hold as a result, rather than
1221 /// left waiting on a dependency this call just erased.
1222 pub quarantined: Vec<String>,
1223}
1224
1225/// Exclusive ownership of a task, released on drop.
1226#[derive(Debug)]
1227pub struct Claim {
1228 path: PathBuf,
1229}
1230
1231impl Drop for Claim {
1232 fn drop(&mut self) {
1233 let _ = std::fs::remove_file(&self.path);
1234 }
1235}
1236
1237/// The first line of a task, trimmed to a title. Used when the caller gives a
1238/// body but no title, which is the normal case for an agent piping a file in.
1239pub fn title_from(instruction: &str, max: usize) -> String {
1240 // The first non-blank line, whatever it is. A markdown heading is the
1241 // task's own summary - agents pipe in `# Rework the config loader` and mean
1242 // exactly that - so it is preferred over the prose beneath it rather than
1243 // skipped as decoration. Leading list and heading markers are stripped
1244 // because they are syntax, not words.
1245 let line = instruction
1246 .lines()
1247 .map(str::trim)
1248 .find(|l| !l.is_empty())
1249 .unwrap_or("(empty task)")
1250 .trim_start_matches(['#', '-', '*', '>', ' '])
1251 .trim();
1252 if line.is_empty() {
1253 return "(empty task)".to_owned();
1254 }
1255 if line.chars().count() <= max {
1256 return line.to_owned();
1257 }
1258 let head: String = line.chars().take(max.saturating_sub(1)).collect();
1259 format!("{head}…")
1260}
1261
1262fn read_path(path: &Path) -> Result<Task> {
1263 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1264 let task: Task =
1265 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1266 // Greater-than, not not-equal: every field added since schema 1 carries
1267 // `#[serde(default)]`, so an older task has nothing to say about it and
1268 // defaulting is exactly as good a reading as a value that build never had
1269 // a chance to write. Only a schema *ahead* of this build - a meaning it
1270 // cannot possibly know - is refused rather than guessed at.
1271 if task.schema > SCHEMA {
1272 bail!(
1273 "task {} was written by a different magi (schema {}, this build \
1274 speaks {SCHEMA})",
1275 task.id,
1276 task.schema
1277 );
1278 }
1279 Ok(task)
1280}
1281
1282/// Ids inside a `blocked_by` list that name neither an existing task file nor
1283/// an existing question file - a dependency deleted (`magi task rm`, or by
1284/// hand) while something was still waiting on it.
1285///
1286/// Existence is decided by [`Queue::path_of`]/[`Questions::path_of`]
1287/// `is_file()` alone, never by [`Queue::get`]/[`Questions::get`] succeeding:
1288/// those also fail on a merely unreadable file - mid-write, corrupt, or from
1289/// a schema ahead of this build (see [`read_path`]) - and misreading "cannot
1290/// read it right now" as "it was deleted" would quarantine a task over a
1291/// transient failure. `blocked_by` always carries a full id, written by
1292/// `crate::conduct` or `crate::triage` from a real task's or question's own
1293/// `id`/`short`, never a prefix a caller typed - so the exact-path check is
1294/// complete on its own, with no [`Queue::resolve_id`] fallback needed.
1295pub fn missing_blockers(
1296 queue: &Queue,
1297 questions: &Questions,
1298 blocked_by: &[String],
1299) -> Vec<String> {
1300 blocked_by
1301 .iter()
1302 .filter(|id| !queue.path_of(id).is_file() && !questions.path_of(id).is_file())
1303 .cloned()
1304 .collect()
1305}
1306
1307/// The `hold_reason` text for a task quarantined because one or more of its
1308/// `blocked_by` ids no longer exist. Shared by `crate::daemon::resolve_blockers`,
1309/// `crate::triage::run_once`, and [`Queue::remove`]'s own dependent
1310/// quarantine, so the three call sites read as the same event to an operator
1311/// looking at `magi task show` rather than three different wordings for it.
1312///
1313/// Names the full original `blocked_by` list, not just `missing` - a task
1314/// quarantined here can also have named a dependency that was still
1315/// perfectly valid, and [`Task::hold_machine`] clears `blocked_by` on the way
1316/// in, so this text is the only place that information survives for an
1317/// operator deciding whether to release the task outright.
1318pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
1319 format!(
1320 "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
1321 blocked_by.join(", "),
1322 missing.join(", "),
1323 )
1324}
1325
1326fn short(id: &str) -> &str {
1327 id.split('-').next_back().unwrap_or(id)
1328}
1329
1330fn new_id() -> String {
1331 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1332 let seed = crate::rng::entropy();
1333 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1334}
1335
1336#[cfg(test)]
1337mod tests {
1338 use super::*;
1339
1340 #[test]
1341 fn triage_applied_survives_release_and_old_records_read_as_empty() {
1342 let mut t = Task::new(
1343 "t".to_owned(),
1344 "i".to_owned(),
1345 PathBuf::from("r"),
1346 Source::Human,
1347 );
1348 t.mark_triage_applied("q1");
1349 t.mark_triage_applied("q1");
1350 t.hold_machine(Some("x".to_owned()));
1351 t.release();
1352 assert_eq!(t.triage_applied, ["q1"]);
1353 assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
1354
1355 let mut v = serde_json::to_value(&t).unwrap();
1356 v.as_object_mut().unwrap().remove("triage_applied");
1357 let old: Task = serde_json::from_value(v).unwrap();
1358 assert!(old.triage_applied.is_empty());
1359 }
1360
1361 #[test]
1362 fn task_counts_of_empty_is_all_zero() {
1363 assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
1364 }
1365
1366 #[test]
1367 fn task_counts_of_tallies_every_status() {
1368 let mut queued = Task::new(
1369 "q".to_owned(),
1370 "i".to_owned(),
1371 PathBuf::from("."),
1372 Source::Human,
1373 );
1374 queued.status = TaskStatus::Queued;
1375 let mut running = queued.clone();
1376 running.status = TaskStatus::Running;
1377 let mut done = queued.clone();
1378 done.status = TaskStatus::Done;
1379 let mut failed = queued.clone();
1380 failed.status = TaskStatus::Failed;
1381 let mut held = queued.clone();
1382 held.status = TaskStatus::Held;
1383 let mut blocked = queued.clone();
1384 blocked.status = TaskStatus::Blocked;
1385
1386 let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
1387 assert_eq!(
1388 counts,
1389 TaskCounts {
1390 queued: 1,
1391 running: 1,
1392 done: 2,
1393 failed: 1,
1394 held: 1,
1395 blocked: 1,
1396 }
1397 );
1398 }
1399
1400 /// A queue of its own, with no process-global state - which is the point of
1401 /// `Queue::at`, and why these can run in parallel.
1402 fn queue() -> (tempfile::TempDir, Queue) {
1403 let dir = tempfile::tempdir().unwrap();
1404 let q = Queue::at(dir.path().join("queue"));
1405 (dir, q)
1406 }
1407
1408 #[test]
1409 fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
1410 let dir = tempfile::tempdir().unwrap();
1411 let q = Queue::at(dir.path().join("queue"));
1412 let mut t = task("held");
1413 q.put(&mut t).unwrap();
1414 assert_eq!(
1415 crate::notices::Notices::at(dir.path().join("notifications"))
1416 .list()
1417 .len(),
1418 0
1419 );
1420 t.hold_machine(Some("out of attempts".to_owned()));
1421 q.put(&mut t).unwrap();
1422 let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
1423 assert_eq!(listed.len(), 1);
1424 assert!(listed[0].message.contains("out of attempts"));
1425 }
1426
1427 fn task(title: &str) -> Task {
1428 Task::new(
1429 title.to_owned(),
1430 format!("do {title}"),
1431 PathBuf::from("."),
1432 Source::Human,
1433 )
1434 }
1435
1436 #[test]
1437 fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
1438 let mut t = task("retried");
1439 assert!(
1440 t.earlier_attempts().is_empty(),
1441 "a first attempt takes nothing over"
1442 );
1443 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1444 assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
1445 // What the web note shows and what a takeover acts on are one rule.
1446 assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
1447 assert_eq!(t.successor_of("bbbb"), None);
1448 assert_eq!(t.successor_of("zzzz"), None);
1449 }
1450
1451 #[test]
1452 fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
1453 let (_dir, q) = queue();
1454 let mut t = task("retried");
1455 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1456 q.put(&mut t).unwrap();
1457
1458 assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
1459 assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
1460 assert_eq!(
1461 q.superseded_by("cccc"),
1462 None,
1463 "the latest attempt replaces nothing"
1464 );
1465 assert_eq!(
1466 q.superseded_by("never-heard-of-it"),
1467 None,
1468 "a run belonging to no task on this queue is not superseded"
1469 );
1470
1471 let mut by = HashMap::new();
1472 by.insert("aaaa".to_owned(), "bbbb".to_owned());
1473 by.insert("bbbb".to_owned(), "cccc".to_owned());
1474 assert_eq!(
1475 q.superseded(),
1476 by,
1477 "the whole-map and single-run forms must agree"
1478 );
1479 }
1480
1481 #[test]
1482 fn superseded_attempts_is_empty_until_the_task_is_done() {
1483 let mut t = task("retried");
1484 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1485 t.status = TaskStatus::Failed;
1486 assert_eq!(
1487 t.superseded_attempts(true),
1488 &[] as &[String],
1489 "a task still retrying has no attempt yet that a later one made moot"
1490 );
1491
1492 t.status = TaskStatus::Running;
1493 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1494 }
1495
1496 #[test]
1497 fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
1498 let mut t = task("retried");
1499 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1500 t.status = TaskStatus::Done;
1501 assert_eq!(
1502 t.superseded_attempts(true),
1503 &["aaaa".to_owned(), "bbbb".to_owned()],
1504 "cccc is the attempt whose success made the task done, and stays out"
1505 );
1506 }
1507
1508 #[test]
1509 fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
1510 let mut t = task("first try landed");
1511 t.runs = vec!["aaaa".to_owned()];
1512 t.status = TaskStatus::Done;
1513 assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1514 }
1515
1516 #[test]
1517 fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
1518 // `succeed` is reachable directly - `magi task done`, its web
1519 // equivalent, and the conductor's `Recovery::Done` - on a task in
1520 // any status, including one whose last recorded attempt is itself
1521 // `Blocked`/`Failed`/anything but `Merged`/`Ready`. `runs.last()`
1522 // alone cannot tell that apart from the loop's own settle path, so
1523 // the caller's own read of that run's status is what decides this.
1524 let mut t = task("closed by hand after a manual merge");
1525 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1526 t.status = TaskStatus::Done;
1527 assert_eq!(
1528 t.superseded_attempts(false),
1529 &[] as &[String],
1530 "nothing here is provably why the task is done, so nothing is superseded"
1531 );
1532 }
1533
1534 #[test]
1535 fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
1536 let (_dir, q) = queue();
1537 let mut t = task("retried twice");
1538 t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1539 q.put(&mut t).unwrap();
1540
1541 assert_eq!(
1542 q.latest_attempt("aaaa"),
1543 Some("cccc".to_owned()),
1544 "an old attempt points straight at the chain's current head, not the \
1545 next attempt in the middle of it"
1546 );
1547 assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
1548 assert_eq!(
1549 q.latest_attempt("cccc"),
1550 None,
1551 "the latest attempt is not superseded by anything"
1552 );
1553 assert_eq!(
1554 q.latest_attempt("never-heard-of-it"),
1555 None,
1556 "a run belonging to no task on this queue is not superseded"
1557 );
1558 }
1559
1560 #[test]
1561 fn a_markdown_heading_is_the_title_not_decoration() {
1562 // A task file's heading is the summary its author already wrote, so it
1563 // beats the prose underneath. Getting this backwards was visible in the
1564 // first smoke test: a task titled "# Rework the config loader" listed
1565 // as "It re-reads the file on every lookup".
1566 assert_eq!(
1567 title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
1568 "Rework the config loader"
1569 );
1570 assert_eq!(title_from("- fix the thing", 40), "fix the thing");
1571 assert_eq!(title_from("> quoted task", 40), "quoted task");
1572 // Nothing usable at all still has to produce something printable.
1573 assert_eq!(title_from(" \n\n", 40), "(empty task)");
1574 assert_eq!(title_from("###\n", 40), "(empty task)");
1575 }
1576
1577 #[test]
1578 fn a_long_title_is_elided_by_characters_not_bytes() {
1579 // Byte truncation would split a multi-byte character and panic.
1580 let long = "課題".repeat(30);
1581 let title = title_from(&long, 10);
1582 assert_eq!(title.chars().count(), 10);
1583 assert!(title.ends_with('…'));
1584 }
1585
1586 #[test]
1587 fn priority_wins_and_ties_break_oldest_first() {
1588 let (_dir, q) = queue();
1589 let mut a = task("first");
1590 let mut b = task("second");
1591 let mut c = task("urgent");
1592 // Ids carry a timestamp, so force a known order.
1593 a.id = "20260101-000001-aaaa".to_owned();
1594 b.id = "20260101-000002-bbbb".to_owned();
1595 c.id = "20260101-000003-cccc".to_owned();
1596 c.priority = 5;
1597 for t in [&mut a, &mut b, &mut c] {
1598 q.put(t).unwrap();
1599 }
1600
1601 // Priority first...
1602 assert_eq!(q.next_runnable().unwrap().id, c.id);
1603 c.hold_machine(None);
1604 q.put(&mut c).unwrap();
1605 // ...then oldest, so a burst of new work cannot starve older work.
1606 assert_eq!(q.next_runnable().unwrap().id, a.id);
1607 assert_eq!(q.list().len(), 3, "b is still waiting its turn");
1608 }
1609
1610 #[test]
1611 fn a_blocked_task_never_starves_another_runnable_one() {
1612 let (_dir, q) = queue();
1613 let mut blocked = task("blocked");
1614 blocked.block(vec!["something".to_owned()], None);
1615 q.put(&mut blocked).unwrap();
1616
1617 let mut runnable = task("free to go");
1618 q.put(&mut runnable).unwrap();
1619
1620 let next = q.next_runnable().expect("a runnable task is still offered");
1621 assert_eq!(next.id, runnable.id);
1622 }
1623
1624 #[test]
1625 fn a_held_task_is_never_offered_to_the_loop() {
1626 let (_dir, q) = queue();
1627 let mut t = task("held");
1628 q.put(&mut t).unwrap();
1629 assert!(q.next_runnable().is_some());
1630
1631 t.hold_machine(None);
1632 q.put(&mut t).unwrap();
1633 assert!(
1634 q.next_runnable().is_none(),
1635 "a held task must wait for a human"
1636 );
1637
1638 // A failed task, by contrast, is exactly what the loop should retry.
1639 t.status = TaskStatus::Failed;
1640 q.put(&mut t).unwrap();
1641 assert!(q.next_runnable().is_some());
1642 }
1643
1644 #[test]
1645 fn attempts_are_capped_and_then_the_task_is_held() {
1646 let mut t = task("doomed");
1647
1648 t.start("run-1".to_owned());
1649 t.fail("gate red", 2);
1650 assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
1651
1652 t.start("run-2".to_owned());
1653 t.fail("gate red", 2);
1654 assert_eq!(
1655 t.status,
1656 TaskStatus::Held,
1657 "out of attempts: stop spending money on it"
1658 );
1659 assert_eq!(t.runs, ["run-1", "run-2"]);
1660 assert_eq!(t.last_error.as_deref(), Some("gate red"));
1661 assert_eq!(
1662 t.hold_reason.as_deref(),
1663 Some("gate red"),
1664 "the hold must say why, not leave hold_reason null next to a \
1665 populated last_error"
1666 );
1667 }
1668
1669 #[test]
1670 fn handing_off_a_task_records_a_hold_reason_too() {
1671 let mut t = task("left a pull request");
1672 t.start("run-1".to_owned());
1673 t.handed_off("run ended with a pull request open [run run-1]");
1674 assert_eq!(t.status, TaskStatus::Held);
1675 assert_eq!(t.hold_source, Some(HoldSource::Machine));
1676 assert_eq!(
1677 t.hold_reason.as_deref(),
1678 Some("run ended with a pull request open [run run-1]")
1679 );
1680 assert_eq!(t.hold_reason, t.last_error);
1681 }
1682
1683 #[test]
1684 fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
1685 let mut t = task("stalled by quota");
1686
1687 t.start("run-1".to_owned());
1688 assert_eq!(t.attempts, 1);
1689 t.stall("judge-1, judge-2 out of quota");
1690 assert_eq!(
1691 t.attempts, 0,
1692 "a closed quota window must not spend the task's retry budget"
1693 );
1694 assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
1695 assert_eq!(
1696 t.last_error.as_deref(),
1697 Some("judge-1, judge-2 out of quota")
1698 );
1699
1700 // A task can therefore stall all night and still get its real attempts
1701 // once the quota resets - which is the whole point.
1702 for _ in 0..20 {
1703 t.start("run-n".to_owned());
1704 t.stall("still out of quota");
1705 }
1706 t.start("run-real".to_owned());
1707 t.fail("gate red", 2);
1708 assert_eq!(
1709 t.status,
1710 TaskStatus::Failed,
1711 "the first attempt that was really judged is attempt one"
1712 );
1713 }
1714
1715 #[test]
1716 fn releasing_a_held_task_gives_it_a_real_second_chance() {
1717 let mut t = task("retry me");
1718 t.start("run-1".to_owned());
1719 t.fail("gate red", 1);
1720 assert_eq!(t.status, TaskStatus::Held);
1721
1722 t.release();
1723 assert_eq!(t.status, TaskStatus::Queued);
1724 // Without resetting attempts the next failure would re-hold at once,
1725 // and a release would be a no-op the operator cannot see.
1726 assert_eq!(t.attempts, 0);
1727 assert!(t.last_error.is_none());
1728 assert_eq!(
1729 t.runs.len(),
1730 1,
1731 "history is kept: attempts reset, evidence does not"
1732 );
1733 }
1734
1735 #[test]
1736 fn a_hold_reason_survives_and_a_release_clears_it() {
1737 let mut t = task("waiting on something else");
1738 t.hold_manual(Some(
1739 "waiting for 20260101-000000-aaaa to land first".to_owned(),
1740 ));
1741 assert_eq!(t.status, TaskStatus::Held);
1742 assert_eq!(
1743 t.hold_reason.as_deref(),
1744 Some("waiting for 20260101-000000-aaaa to land first")
1745 );
1746
1747 // Holding again with no reason must not erase the one already there.
1748 t.hold_manual(None);
1749 assert_eq!(
1750 t.hold_reason.as_deref(),
1751 Some("waiting for 20260101-000000-aaaa to land first"),
1752 "a bare re-hold keeps whatever a human already wrote down"
1753 );
1754
1755 // A hold with no reason at all is still an ordinary, allowed hold.
1756 let mut plain = task("no reason given");
1757 plain.hold_manual(None);
1758 assert_eq!(plain.status, TaskStatus::Held);
1759 assert!(plain.hold_reason.is_none());
1760
1761 t.release();
1762 assert_eq!(t.status, TaskStatus::Queued);
1763 assert!(
1764 t.hold_reason.is_none(),
1765 "a stale reason must not greet the next person who holds this task"
1766 );
1767 }
1768
1769 #[test]
1770 fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
1771 // `done` can close a held task directly - neither `magi task done`
1772 // nor `POST /api/queue/{id}/done` requires a release first - so a
1773 // task held for "waiting on 3ed9" and then closed without ever being
1774 // released must not still read as waiting on it afterwards.
1775 let mut t = task("landed by hand while held");
1776 t.hold_manual(Some("waiting on 3ed9".to_owned()));
1777 assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
1778
1779 t.succeed();
1780 assert_eq!(t.status, TaskStatus::Done);
1781 assert!(
1782 t.hold_reason.is_none(),
1783 "a done task cannot still be waiting on something"
1784 );
1785 }
1786
1787 #[test]
1788 fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
1789 // The web UI's "Hold" and "Mark done" buttons are both reachable on a
1790 // `blocked` task, not just on `queued`/`held` ones - neither requires
1791 // a release first. A task moved off `Blocked` that way must not still
1792 // carry the dependency it was waiting on: a dependency graph built
1793 // from `blocked_by` would otherwise keep drawing an edge for a task
1794 // that is not blocked on anything any more.
1795 let mut held = task("held straight out of blocked");
1796 held.block(
1797 vec!["20260101-000000-dead".to_owned()],
1798 Some("waiting on the migration script".to_owned()),
1799 );
1800 assert_eq!(held.status, TaskStatus::Blocked);
1801
1802 held.hold_manual(None);
1803 assert_eq!(held.status, TaskStatus::Held);
1804 assert!(
1805 held.blocked_by.is_empty(),
1806 "hold overrides the wait, same as release"
1807 );
1808 assert!(held.block_reason.is_none());
1809
1810 let mut done = task("closed straight out of blocked");
1811 done.block(
1812 vec!["20260101-000000-dead".to_owned()],
1813 Some("waiting on the migration script".to_owned()),
1814 );
1815 done.succeed();
1816 assert_eq!(done.status, TaskStatus::Done);
1817 assert!(
1818 done.blocked_by.is_empty(),
1819 "a done task cannot still be waiting on a dependency"
1820 );
1821 assert!(done.block_reason.is_none());
1822 }
1823
1824 #[test]
1825 fn a_blocked_task_is_never_offered_to_the_loop() {
1826 let mut t = task("blocked");
1827 assert!(t.status.runnable());
1828 t.block(
1829 vec!["dep-id".to_owned()],
1830 Some("waits on dep-id".to_owned()),
1831 );
1832 assert_eq!(t.status, TaskStatus::Blocked);
1833 assert!(!t.status.runnable());
1834 assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
1835 }
1836
1837 #[test]
1838 fn unblocking_the_last_dependency_returns_the_task_to_queued() {
1839 let mut t = task("blocked on two");
1840 t.block(
1841 vec!["a".to_owned(), "b".to_owned()],
1842 Some("waits on a and b".to_owned()),
1843 );
1844
1845 t.unblock("a");
1846 assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
1847 assert_eq!(t.blocked_by, ["b"]);
1848
1849 t.unblock("b");
1850 assert_eq!(t.status, TaskStatus::Queued);
1851 assert!(t.blocked_by.is_empty());
1852 assert!(t.block_reason.is_none());
1853 }
1854
1855 #[test]
1856 fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
1857 let mut t = task("never blocked");
1858 t.unblock("whatever");
1859 assert_eq!(t.status, TaskStatus::Queued);
1860 }
1861
1862 #[test]
1863 fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
1864 // The bug this guards: a task an operator (or `crate::triage`) has
1865 // deliberately held, once `crate::conduct` blocks it on a follow-up
1866 // question, must not silently re-enter the competition queue the
1867 // moment that question is answered - whatever the answer said.
1868 let mut t = task("held, then asked about");
1869 t.hold_machine(Some("out of attempts".to_owned()));
1870 assert_eq!(t.status, TaskStatus::Held);
1871
1872 t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
1873 assert_eq!(t.status, TaskStatus::Blocked);
1874
1875 t.record_answer("what now?".to_owned(), "leave it held".to_owned());
1876 t.unblock("q1");
1877 assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
1878 assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
1879 assert_eq!(t.hold_source, Some(HoldSource::Machine));
1880 assert!(t.blocked_from.is_none(), "consumed once restored");
1881 }
1882
1883 #[test]
1884 fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
1885 let mut t = task("manually held, then asked about");
1886 t.hold_manual(Some("waiting on a dependency".to_owned()));
1887
1888 t.block(vec!["q1".to_owned()], None);
1889 t.unblock("q1");
1890
1891 assert_eq!(t.status, TaskStatus::Held);
1892 assert_eq!(t.hold_source, Some(HoldSource::Manual));
1893 }
1894
1895 #[test]
1896 fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
1897 // A second `Task::block` call - `crate::conduct` adding a question on
1898 // top of an existing block - must not overwrite `blocked_from` with
1899 // `Blocked` itself, or the task would restore into itself.
1900 let mut t = task("held, blocked twice");
1901 t.hold_machine(None);
1902 t.block(vec!["q1".to_owned()], Some("first".to_owned()));
1903 t.block(
1904 vec!["q1".to_owned(), "q2".to_owned()],
1905 Some("second".to_owned()),
1906 );
1907
1908 t.unblock("q1");
1909 assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
1910 t.unblock("q2");
1911 assert_eq!(t.status, TaskStatus::Held);
1912 }
1913
1914 #[test]
1915 fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
1916 // Whatever process was driving the run is gone by the time a
1917 // conductor's question about it gets answered - there is nothing left
1918 // to resume into.
1919 let mut t = task("blocked mid-run");
1920 t.start("run-1".to_owned());
1921 assert_eq!(t.status, TaskStatus::Running);
1922
1923 t.block(vec!["q1".to_owned()], None);
1924 t.unblock("q1");
1925 assert_eq!(t.status, TaskStatus::Queued);
1926 }
1927
1928 #[test]
1929 fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
1930 // `blocked_from` is `None` for a record written before schema 4 (or,
1931 // equivalently, deserialized straight from an on-disk file that never
1932 // had the field). Held evidence surviving on the task - never cleared
1933 // by `block` - is the only way left to tell such a record apart from
1934 // one blocked straight out of `Queued`.
1935 let mut t = task("legacy record, held before it was blocked");
1936 t.hold_source = Some(HoldSource::Machine);
1937 t.hold_reason = Some("legacy hold reason".to_owned());
1938 t.status = TaskStatus::Blocked;
1939 t.blocked_by = vec!["q1".to_owned()];
1940 t.blocked_from = None;
1941
1942 t.unblock("q1");
1943 assert_eq!(t.status, TaskStatus::Held);
1944 }
1945
1946 #[test]
1947 fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
1948 let mut t = task("legacy record, ordinary dependency block");
1949 t.status = TaskStatus::Blocked;
1950 t.blocked_by = vec!["dep".to_owned()];
1951 t.blocked_from = None;
1952
1953 t.unblock("dep");
1954 assert_eq!(t.status, TaskStatus::Queued);
1955 }
1956
1957 #[test]
1958 fn answering_a_question_is_recorded_and_survives_a_release() {
1959 let mut t = task("asked something");
1960 t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
1961 t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
1962 t.unblock("q1");
1963 assert_eq!(t.status, TaskStatus::Queued);
1964 assert_eq!(t.answers.len(), 1);
1965 assert_eq!(t.answers[0].answer, "SQLite");
1966
1967 // A release resets attempts, not evidence - the same rule
1968 // `releasing_a_held_task_gives_it_a_real_second_chance` asserts for
1969 // `runs`.
1970 t.release();
1971 assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
1972 }
1973
1974 #[test]
1975 fn requesting_review_requeues_the_task_and_remembers_the_branch() {
1976 let mut t = task("blocked run with a surviving branch");
1977 t.start("run-1".to_owned());
1978 t.fail("blocked with major findings", 5);
1979 assert_eq!(t.status, TaskStatus::Failed);
1980
1981 t.request_review("magi/eba2/A".to_owned());
1982 assert_eq!(t.status, TaskStatus::Queued);
1983 assert_eq!(t.attempts, 0);
1984 assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
1985
1986 // An ordinary release (a human overriding the choice) drops it again.
1987 t.release();
1988 assert!(t.review_branch.is_none());
1989 }
1990
1991 #[test]
1992 fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
1993 let mut t = task("retry");
1994 t.start("run-1".to_owned());
1995 t.requeue();
1996 assert!(t.fresh_start);
1997
1998 t.release();
1999 assert!(!t.fresh_start);
2000 }
2001
2002 #[test]
2003 fn priority_can_be_changed_while_queued_but_not_while_running() {
2004 let mut t = task("reprioritise me");
2005 t.set_priority(5).unwrap();
2006 assert_eq!(t.priority, 5);
2007
2008 t.start("run-1".to_owned());
2009 let err = t.set_priority(9).unwrap_err().to_string();
2010 assert!(err.contains("running"), "{err}");
2011 assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2012 }
2013
2014 #[test]
2015 fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2016 let mut t = task("interrupt me");
2017 assert!(!t.interrupt, "off unless asked, same as any other task");
2018
2019 t.set_interrupt(true).unwrap();
2020 assert!(t.interrupt);
2021
2022 t.start("run-1".to_owned());
2023 assert!(
2024 !t.interrupt,
2025 "the mark is one-shot: dispatching the task fulfils it, \
2026 whatever the run that follows ends up doing"
2027 );
2028 let err = t.set_interrupt(true).unwrap_err().to_string();
2029 assert!(err.contains("running"), "{err}");
2030 // Clearing is always allowed, even on a running task - there is
2031 // nothing left for it to interrupt once it has been claimed.
2032 t.set_interrupt(false).unwrap();
2033 assert!(!t.interrupt);
2034 }
2035
2036 /// R2-1-1: a task whose run fails and requeues must not go on
2037 /// re-triggering `crate::daemon`'s interrupt scheduler on every later
2038 /// boundary, attempt after attempt, until it exhausts its budget.
2039 #[test]
2040 fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2041 let mut t = task("interrupt me");
2042 t.set_interrupt(true).unwrap();
2043 t.start("run-1".to_owned());
2044 t.fail("mock failure", 5);
2045 assert_eq!(t.status, TaskStatus::Failed);
2046 assert!(
2047 !t.interrupt,
2048 "one attempt already spent the mark; a retry is an ordinary \
2049 requeue, not a fresh interrupt request"
2050 );
2051 }
2052
2053 #[test]
2054 fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2055 let (_dir, q) = queue();
2056 let mut a = task("first filed");
2057 let mut b = task("second filed");
2058 a.id = "20260101-000001-aaaa".to_owned();
2059 b.id = "20260101-000002-bbbb".to_owned();
2060 q.put(&mut a).unwrap();
2061 q.put(&mut b).unwrap();
2062
2063 assert_eq!(
2064 q.next_runnable().unwrap().id,
2065 a.id,
2066 "with equal priority the older task goes first, so a burst of \
2067 new work cannot starve it"
2068 );
2069 assert_eq!(
2070 q.list()[0].id,
2071 b.id,
2072 "but the list an operator reads is newest first, the same as \
2073 before priority existed - a's turn to run does not make it the \
2074 newest task"
2075 );
2076
2077 let mut a = q.get(&a.id).unwrap();
2078 a.set_priority(10).unwrap();
2079 q.put(&mut a).unwrap();
2080
2081 assert_eq!(
2082 q.next_runnable().unwrap().id,
2083 a.id,
2084 "a raised priority must be reflected the moment it is saved"
2085 );
2086 // `magi task list` and `GET /api/queue` both print `Queue::list()`
2087 // directly, so the raised task has to lead there too - not only in
2088 // what the loop would claim next.
2089 assert_eq!(
2090 q.list()[0].id,
2091 a.id,
2092 "the raised task must sort first in the list an operator reads, \
2093 not only in next_runnable's own ordering"
2094 );
2095 }
2096
2097 #[test]
2098 fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2099 let mut t = Task::new(
2100 "old title".to_owned(),
2101 "old instruction".to_owned(),
2102 PathBuf::from("/repo"),
2103 Source::Agent {
2104 run: "20260101-000000-beef".to_owned(),
2105 node: "implement".to_owned(),
2106 },
2107 );
2108 let id = t.id.clone();
2109 let created_at = t.created_at;
2110 t.runs.push("20260101-000000-beef".to_owned());
2111
2112 t.edit("new title".to_owned(), "new instruction".to_owned())
2113 .unwrap();
2114
2115 assert_eq!(t.title, "new title");
2116 assert_eq!(t.instruction, "new instruction");
2117 assert_eq!(t.id, id, "editing must not mint a new id");
2118 assert_eq!(t.created_at, created_at);
2119 assert_eq!(
2120 t.source,
2121 Source::Agent {
2122 run: "20260101-000000-beef".to_owned(),
2123 node: "implement".to_owned(),
2124 },
2125 "editing must not turn agent attribution into human"
2126 );
2127 assert_eq!(t.runs, ["20260101-000000-beef"]);
2128 }
2129
2130 #[test]
2131 fn editing_is_refused_once_a_task_is_running_or_finished() {
2132 let mut running = task("in flight");
2133 running.start("run-1".to_owned());
2134 let err = running
2135 .edit("x".to_owned(), "y".to_owned())
2136 .unwrap_err()
2137 .to_string();
2138 assert!(err.contains("running"), "{err}");
2139
2140 let mut done = task("finished");
2141 done.succeed();
2142 let err = done
2143 .edit("x".to_owned(), "y".to_owned())
2144 .unwrap_err()
2145 .to_string();
2146 assert!(err.contains("done"), "{err}");
2147
2148 // Both queued and held are the point of the feature and must work.
2149 let mut queued = task("waiting");
2150 queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2151 let mut held = task("parked");
2152 held.hold_machine(None);
2153 held.edit("x".to_owned(), "y".to_owned()).unwrap();
2154 }
2155
2156 #[test]
2157 fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2158 let (_dir, q) = queue();
2159 let path = q.path_of("20260101-000000-aaaa");
2160 std::fs::create_dir_all(q.root()).unwrap();
2161 std::fs::write(
2162 &path,
2163 serde_json::json!({
2164 "schema": SCHEMA,
2165 "id": "20260101-000000-aaaa",
2166 "title": "from before hold reasons existed",
2167 "instruction": "from before hold reasons existed",
2168 "repo": ".",
2169 "source": { "kind": "human" },
2170 "status": "held",
2171 "created_at": Timestamp::now().to_string(),
2172 "updated_at": Timestamp::now().to_string(),
2173 })
2174 .to_string(),
2175 )
2176 .unwrap();
2177
2178 let task = q.get("20260101-000000-aaaa").expect("must still read");
2179 assert!(task.hold_reason.is_none());
2180 assert!(task.operator_held());
2181 }
2182
2183 #[test]
2184 fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
2185 let (_dir, q) = queue();
2186 let path = q.path_of("20260101-000000-bbbb");
2187 std::fs::create_dir_all(q.root()).unwrap();
2188 std::fs::write(
2189 &path,
2190 serde_json::json!({
2191 "schema": 2,
2192 "id": "20260101-000000-bbbb",
2193 "title": "old manual recovery",
2194 "instruction": "old manual recovery",
2195 "repo": ".",
2196 "source": { "kind": "human" },
2197 "status": "held",
2198 "hold_reason": "active manual recovery run20260912-224242-daf5",
2199 "created_at": Timestamp::now().to_string(),
2200 "updated_at": Timestamp::now().to_string(),
2201 })
2202 .to_string(),
2203 )
2204 .unwrap();
2205
2206 let task = q.get("20260101-000000-bbbb").expect("must still read");
2207 assert_eq!(task.hold_source, None);
2208 assert!(task.operator_held());
2209 }
2210
2211 #[test]
2212 fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
2213 let (_dir, q) = queue();
2214 let path = q.path_of("20260101-000000-aaaa");
2215 std::fs::create_dir_all(q.root()).unwrap();
2216 std::fs::write(
2217 &path,
2218 serde_json::json!({
2219 "schema": SCHEMA,
2220 "id": "20260101-000000-aaaa",
2221 "title": "from before diagnostics existed",
2222 "instruction": "from before diagnostics existed",
2223 "repo": ".",
2224 "source": { "kind": "human" },
2225 "status": "held",
2226 "created_at": Timestamp::now().to_string(),
2227 "updated_at": Timestamp::now().to_string(),
2228 })
2229 .to_string(),
2230 )
2231 .unwrap();
2232
2233 let task = q.get("20260101-000000-aaaa").expect("must still read");
2234 assert!(task.diagnostic.is_none());
2235 }
2236
2237 #[test]
2238 fn a_schema_1_task_with_no_blocking_fields_still_reads() {
2239 // Written by a build that predates `blocked_by`, `block_reason`,
2240 // `answers` and `review_branch` entirely - literal `"schema": 1`,
2241 // not `SCHEMA`, since the whole point is a build older than this one.
2242 let (_dir, q) = queue();
2243 let path = q.path_of("20260101-000000-aaaa");
2244 std::fs::create_dir_all(q.root()).unwrap();
2245 std::fs::write(
2246 &path,
2247 serde_json::json!({
2248 "schema": 1,
2249 "id": "20260101-000000-aaaa",
2250 "title": "from before blocking existed",
2251 "instruction": "from before blocking existed",
2252 "repo": ".",
2253 "source": { "kind": "human" },
2254 "status": "queued",
2255 "created_at": Timestamp::now().to_string(),
2256 "updated_at": Timestamp::now().to_string(),
2257 })
2258 .to_string(),
2259 )
2260 .unwrap();
2261
2262 let task = q.get("20260101-000000-aaaa").expect("must still read");
2263 assert!(task.blocked_by.is_empty());
2264 assert!(task.block_reason.is_none());
2265 assert!(task.answers.is_empty());
2266 assert!(task.review_branch.is_none());
2267 }
2268
2269 #[test]
2270 fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
2271 // A diagnostic belongs to the run that produced it. Left in place
2272 // across a release, an unrelated later failure - a config error, say -
2273 // would go on showing evidence for a problem that is no longer why the
2274 // task is stuck.
2275 let mut held = task("diagnosed");
2276 held.start("run-1".to_owned());
2277 held.fail("gate red", 1);
2278 held.diagnostic = Some("cargo test failed: ...".to_owned());
2279 assert_eq!(held.status, TaskStatus::Held);
2280
2281 held.release();
2282 assert!(held.diagnostic.is_none());
2283
2284 held.diagnostic = Some("cargo test failed: ...".to_owned());
2285 held.succeed();
2286 assert!(held.diagnostic.is_none());
2287 }
2288
2289 #[test]
2290 fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
2291 let mut t = task("retried");
2292 t.start("run-1".to_owned());
2293 t.diagnostic = Some("stale evidence from a previous hold".to_owned());
2294 t.fail("unrelated config error", 5);
2295 assert_eq!(t.status, TaskStatus::Failed);
2296 assert!(
2297 t.diagnostic.is_none(),
2298 "fail() must not let an old diagnostic outlive the run that produced it"
2299 );
2300 }
2301
2302 #[test]
2303 fn a_claim_is_exclusive_and_releases_on_drop() {
2304 let (_dir, q) = queue();
2305 let mut t = task("contended");
2306 q.put(&mut t).unwrap();
2307
2308 let held = q.claim(&t.id).unwrap();
2309 assert!(
2310 q.claim(&t.id).is_err(),
2311 "two daemons must not drive one task into two runs"
2312 );
2313 drop(held);
2314 assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
2315 }
2316
2317 #[test]
2318 fn a_round_trip_survives_disk() {
2319 let (_dir, q) = queue();
2320 let mut t = Task::new(
2321 "titled".to_owned(),
2322 "body".to_owned(),
2323 PathBuf::from("/repo"),
2324 Source::Agent {
2325 run: "20260101-000000-beef".to_owned(),
2326 node: "implement".to_owned(),
2327 },
2328 );
2329 t.priority = 3;
2330 q.put(&mut t).unwrap();
2331
2332 let back = q.get(&t.id).unwrap();
2333 assert_eq!(back.id, t.id);
2334 assert_eq!(back.priority, 3);
2335 assert_eq!(back.source.label(), "implement@beef");
2336 // A prefix is enough, the way run ids work everywhere else.
2337 assert_eq!(q.get(t.short()).unwrap().id, t.id);
2338 }
2339
2340 #[test]
2341 fn an_unreadable_task_does_not_take_the_queue_down() {
2342 let (_dir, q) = queue();
2343 let mut t = task("fine");
2344 q.put(&mut t).unwrap();
2345 std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
2346
2347 let listed = q.list();
2348 assert_eq!(listed.len(), 1, "the readable task still lists");
2349 assert_eq!(listed[0].id, t.id);
2350 }
2351
2352 #[test]
2353 fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
2354 let (_dir, q) = queue();
2355 let path = q.path_of("20260101-000000-aaaa");
2356 std::fs::create_dir_all(q.root()).unwrap();
2357 std::fs::write(
2358 &path,
2359 serde_json::json!({
2360 "schema": SCHEMA,
2361 "id": "20260101-000000-aaaa",
2362 "title": "from before solo existed",
2363 "instruction": "from before solo existed",
2364 "repo": ".",
2365 "source": { "kind": "human" },
2366 "status": "queued",
2367 "created_at": Timestamp::now().to_string(),
2368 "updated_at": Timestamp::now().to_string(),
2369 })
2370 .to_string(),
2371 )
2372 .unwrap();
2373
2374 let task = q.get("20260101-000000-aaaa").expect("must still read");
2375 assert!(!task.solo, "a queue file with no `solo` field means false");
2376 }
2377
2378 #[test]
2379 fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
2380 let (_dir, q) = queue();
2381 let path = q.path_of("20260101-000000-bbbb");
2382 std::fs::create_dir_all(q.root()).unwrap();
2383 std::fs::write(
2384 &path,
2385 serde_json::json!({
2386 "schema": SCHEMA,
2387 "id": "20260101-000000-bbbb",
2388 "title": "from before urgent existed",
2389 "instruction": "from before urgent existed",
2390 "repo": ".",
2391 "source": { "kind": "human" },
2392 "status": "queued",
2393 "created_at": Timestamp::now().to_string(),
2394 "updated_at": Timestamp::now().to_string(),
2395 })
2396 .to_string(),
2397 )
2398 .unwrap();
2399
2400 let task = q.get("20260101-000000-bbbb").expect("must still read");
2401 assert!(
2402 !task.urgent,
2403 "a queue file with no `urgent` field means false, same as `solo`"
2404 );
2405 }
2406
2407 #[test]
2408 fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
2409 let (_dir, q) = queue();
2410 let mut t = task("from the future");
2411 q.put(&mut t).unwrap();
2412 let path = q.path_of(&t.id);
2413 let body = std::fs::read_to_string(&path)
2414 .unwrap()
2415 .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
2416 std::fs::write(&path, body).unwrap();
2417
2418 let err = q.get(&t.id).unwrap_err().to_string();
2419 assert!(err.contains("schema 99"), "{err}");
2420 }
2421
2422 #[test]
2423 fn revision_moves_when_the_queue_changes() {
2424 let (_dir, q) = queue();
2425 assert_eq!(q.revision(), 0, "an empty queue has no revision");
2426 let mut t = task("first");
2427 q.put(&mut t).unwrap();
2428 assert!(q.revision() > 0, "a written task moves the revision");
2429 }
2430
2431 #[test]
2432 fn revision_moves_when_deleting_an_older_task() {
2433 let (dir, q) = queue();
2434 let questions = Questions::at(dir.path().join("questions"));
2435 let mut t1 = task("older");
2436 q.put(&mut t1).unwrap();
2437 // Ensure mtime ticks forward.
2438 std::thread::sleep(std::time::Duration::from_millis(10));
2439 let mut t2 = task("newer");
2440 q.put(&mut t2).unwrap();
2441
2442 let rev_before = q.revision();
2443 q.remove(&t1.id, false, &questions).unwrap();
2444 let rev_after = q.revision();
2445
2446 assert_ne!(
2447 rev_before, rev_after,
2448 "deleting an older task must change the revision so other clients see the deletion"
2449 );
2450 }
2451
2452 #[test]
2453 fn removing_a_task_takes_it_out_of_the_listing() {
2454 let (dir, q) = queue();
2455 let questions = Questions::at(dir.path().join("questions"));
2456 let mut t = task("delete me");
2457 q.put(&mut t).unwrap();
2458 let removed = q.remove(t.short(), false, &questions).unwrap();
2459 assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
2460 assert!(removed.quarantined.is_empty(), "nothing was blocked on it");
2461 assert!(q.list().is_empty());
2462 assert!(
2463 q.remove(&t.id, false, &questions).is_err(),
2464 "removing twice is an error"
2465 );
2466 }
2467
2468 #[test]
2469 fn removing_a_task_takes_its_stale_lock_with_it() {
2470 let (dir, q) = queue();
2471 let questions = Questions::at(dir.path().join("questions"));
2472 let mut t = task("interrupted");
2473 q.put(&mut t).unwrap();
2474
2475 // A daemon killed mid-run leaves this behind. Nothing holds it: the
2476 // process that would have dropped the guard is gone.
2477 let claim = q.claim(&t.id).unwrap();
2478 std::mem::forget(claim);
2479 assert!(
2480 q.claim(&t.id).is_err(),
2481 "the orphaned lock is what makes the task look claimed"
2482 );
2483
2484 // A live daemon on this task is refused, whatever the lock says.
2485 let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
2486 assert!(err.contains("live daemon"), "{err}");
2487 assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
2488
2489 // With no daemon behind it, the lock is stale and goes with the task.
2490 q.remove(&t.id, false, &questions).unwrap();
2491 assert!(q.list().is_empty());
2492 let mut again = task("interrupted");
2493 again.id = t.id.clone();
2494 q.put(&mut again).unwrap();
2495 assert!(
2496 q.claim(&t.id).is_ok(),
2497 "a task that comes back must be claimable, which a left-behind lock would prevent"
2498 );
2499 }
2500
2501 #[test]
2502 fn removing_a_task_quarantines_what_was_blocked_on_it() {
2503 let (dir, q) = queue();
2504 let questions = Questions::at(dir.path().join("questions"));
2505
2506 let mut dep = task("dependency");
2507 q.put(&mut dep).unwrap();
2508
2509 let mut still_valid = task("still valid");
2510 q.put(&mut still_valid).unwrap();
2511
2512 let mut blocked = task("waiting");
2513 blocked.block(
2514 vec![dep.id.clone(), still_valid.id.clone()],
2515 Some("waits on both".to_owned()),
2516 );
2517 q.put(&mut blocked).unwrap();
2518
2519 let removed = q.remove(&dep.id, false, &questions).unwrap();
2520 assert_eq!(removed.quarantined, [blocked.id.clone()]);
2521
2522 let after = q.get(&blocked.id).unwrap();
2523 assert_eq!(after.status, TaskStatus::Held);
2524 assert_eq!(after.hold_source, Some(HoldSource::Machine));
2525 assert!(after.blocked_by.is_empty());
2526 let reason = after.hold_reason.as_deref().unwrap_or_default();
2527 assert!(reason.contains(&dep.id), "{reason}");
2528 assert!(
2529 reason.contains(&still_valid.id),
2530 "the still-valid dependency must survive in the reason text: {reason}"
2531 );
2532 }
2533}