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