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