Skip to main content

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