Skip to main content

magi/
queue.rs

1//! The task queue: what magi should do next, and who asked for it.
2//!
3//! The queue is what lets magi run unattended. `magi serve` takes the next
4//! task, runs the graph on it, records the outcome, and takes the next one.
5//!
6//! It is also the reason an agent can ask for work. `magi task add` is the
7//! whole interface, and it is the same command whether a human types it at a
8//! prompt, a phone posts it through the web UI, or an implementer inside a run
9//! shells out to it because it noticed something worth doing but out of scope.
10//! magi's CLI is the operating surface for both kinds of user; the queue is
11//! where their intentions meet.
12//!
13//! One task is one JSON file under [`Queue`]'s root. Files rather than a
14//! database because the operator has to be able to read, edit, and delete the
15//! backlog with the tools already on the machine, and because a crashed daemon
16//! must leave a queue the next one can pick up without recovery ceremony.
17//!
18//! # Shape
19//!
20//! [`Task`] is data plus *pure* state transitions - [`Task::fail`] decides
21//! whether an attempt was the last one, and touches no disk. [`Queue`] owns all
22//! I/O and is constructed with its root, so a test drives a real queue in a
23//! temp directory without setting a process-global home. Splitting them this
24//! way is why the retry policy below can be asserted directly.
25//!
26//! # Bounded by construction
27//!
28//! An autonomous loop that retries forever is a way to spend money on a task
29//! that cannot succeed. Every claim increments [`Task::attempts`]; a task that
30//! has burned its attempts becomes [`TaskStatus::Held`] and waits for a human
31//! rather than for another agent.
32
33use std::collections::HashMap;
34use std::path::{Path, PathBuf};
35
36use anyhow::{Context, Result, bail};
37use jiff::Timestamp;
38use serde::{Deserialize, Serialize};
39
40use crate::ask::Questions;
41
42/// On-disk format for a queued task. Bumped when a field's meaning changes.
43///
44/// 7: added [`Task::attachments`], names of files copied under
45/// `<id>.attachments/` beside the task file (a field only; `#[serde(default)]`,
46/// so an older record reads as empty and [`read_path`] still accepts it).
47///
48/// 6: added [`Task::resume_override`] (a field only; `#[serde(default)]`, so
49/// an older record reads as `None` and [`read_path`] still accepts it).
50///
51/// 5: added [`Task::triage_applied`], the ids of triage questions whose
52/// answer has already been applied to this task. A "resume" answer used to
53/// leave no trace ([`Task::release`] clears [`Task::hold_reason`], which was
54/// the only place the applied marker lived), so when the released task failed
55/// its attempts and went back to `held`, the next idle pass found the same
56/// answered question "not applied" and released it again with `attempts` reset
57/// to 0 - the `max_attempts` bound never held. `#[serde(default)]` so an older
58/// record reads as empty. A task already looping when this build arrives has
59/// no record, so it is released once more, recorded, and then stays held.
60///
61/// 8: added [`Task::followup`], the origin of a task `crate::followup` filed
62/// from a merged run's leftover findings. Field-only bump, `#[serde(default)]`.
63///
64/// 4: added [`Task::blocked_from`], the status a task had the moment it
65/// became [`TaskStatus::Blocked`], so [`Task::unblock`] restores it instead
66/// of always landing on [`TaskStatus::Queued`]. Without it, a task a human
67/// or `crate::triage` had deliberately left [`TaskStatus::Held`] — machine
68/// or manual — would lose that the instant `crate::conduct` blocked it on a
69/// follow-up question, and come back `Queued` the moment the question was
70/// answered, regardless of what the answer said: exactly the loop where a
71/// task the operator told to stay held instead re-enters the competition
72/// queue every time someone answers a question about it. `#[serde(default)]`
73/// so an older record reads as `None`; [`Task::unblock`] then falls back to
74/// inferring `Held` from surviving hold evidence ([`Task::hold_reason`] /
75/// [`Task::hold_source`], never cleared by [`Task::block`]) rather than
76/// guessing `Queued` outright — see [`Task::unblock`]'s own doc.
77///
78/// 3: added [`HoldSource`] so conductor recovery cannot release a hold an
79/// operator deliberately placed. Old records default to `None` and are
80/// protected as operator-held until an explicit release; the safe direction
81/// when their author was never recorded.
82///
83/// 2: added [`TaskStatus::Blocked`], [`Task::blocked_by`] and
84/// [`Task::block_reason`] (`crate::conduct`'s decisions) and
85/// [`Task::answers`] (operator answers carried forward to the next
86/// conductor prompt and the next run's instruction). All three are
87/// `#[serde(default)]`, so [`read_path`] accepts anything up to and
88/// including this schema rather than only an exact match — a task written
89/// by a build that only knew about schema 1 has nothing to say about
90/// blocking or answers, and defaulting those fields is exactly as good a
91/// reading as a value that build never had a chance to write.
92pub const SCHEMA: u32 = 8;
93
94/// Who placed the current hold.
95#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
96#[serde(rename_all = "lowercase")]
97pub enum HoldSource {
98    /// An operator used the CLI or web UI.
99    Manual,
100    /// The daemon or conductor placed the hold as part of its own recovery.
101    Machine,
102}
103
104impl HoldSource {
105    /// Short human-facing label for reports and the CLI.
106    pub fn label(self) -> &'static str {
107        match self {
108            Self::Manual => "manual",
109            Self::Machine => "machine",
110        }
111    }
112}
113
114/// Where a task came from. Recorded because "who asked for this" is the first
115/// question about an autonomous run, and the answer is not recoverable later.
116#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
117#[serde(tag = "kind", rename_all = "lowercase")]
118pub enum Source {
119    /// A person, at a terminal or through the web UI.
120    Human,
121    /// An agent inside a run, via `magi task add`. Both ids are recorded so a
122    /// task can be traced back to the exact seat that asked for it.
123    Agent {
124        /// Run the asking agent belonged to.
125        run: String,
126        /// Node it was working in, e.g. `implement` or `review`.
127        node: String,
128    },
129    /// A GitHub issue, imported by number.
130    Issue {
131        /// Issue number.
132        number: u64,
133        /// `owner/repo`, as `gh` reports it.
134        repo: String,
135    },
136}
137
138impl Source {
139    /// Short human-facing label, for lists and the web UI.
140    pub fn label(&self) -> String {
141        match self {
142            Self::Human => "human".to_owned(),
143            Self::Agent { run, node } => format!("{node}@{}", short(run)),
144            Self::Issue { number, .. } => format!("issue #{number}"),
145        }
146    }
147}
148
149/// Where a task is in its life.
150#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
151#[serde(rename_all = "lowercase")]
152pub enum TaskStatus {
153    /// Waiting to be claimed.
154    Queued,
155    /// Claimed by a daemon; a run is in flight.
156    Running,
157    /// A run finished and its gate passed.
158    Done,
159    /// A run finished without passing, and attempts remain.
160    Failed,
161    /// Out of attempts, or held by hand. The loop will not pick it up.
162    Held,
163    /// Waiting on another task or an unanswered question. See
164    /// [`Task::blocked_by`]. Set and cleared by `crate::conduct` and
165    /// `crate::daemon`'s deterministic resolver, never by hand.
166    Blocked,
167}
168
169impl TaskStatus {
170    /// Is this task eligible for a daemon to claim?
171    pub fn runnable(self) -> bool {
172        matches!(self, Self::Queued | Self::Failed)
173    }
174
175    /// Lowercase name, as it appears on disk and in the API.
176    pub fn as_str(self) -> &'static str {
177        match self {
178            Self::Queued => "queued",
179            Self::Running => "running",
180            Self::Done => "done",
181            Self::Failed => "failed",
182            Self::Held => "held",
183            Self::Blocked => "blocked",
184        }
185    }
186}
187
188/// How many tasks sit in each [`TaskStatus`], for a dashboard tile — never a
189/// per-task view, so it carries no ids.
190#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
191pub struct TaskCounts {
192    /// Waiting to be claimed.
193    pub queued: usize,
194    /// Claimed; a run is in flight.
195    pub running: usize,
196    /// A run finished and its gate passed.
197    pub done: usize,
198    /// A run finished without passing, and attempts remain.
199    pub failed: usize,
200    /// Out of attempts, or held by hand.
201    pub held: usize,
202    /// Waiting on another task or an unanswered question.
203    pub blocked: usize,
204}
205
206impl TaskCounts {
207    /// Tally `tasks` by status. A pure function over whatever [`Queue::list`]
208    /// already read, so it needs no root of its own and stays trivially
209    /// testable against a hand-built slice.
210    pub fn of(tasks: &[Task]) -> Self {
211        let mut counts = Self::default();
212        for t in tasks {
213            match t.status {
214                TaskStatus::Queued => counts.queued += 1,
215                TaskStatus::Running => counts.running += 1,
216                TaskStatus::Done => counts.done += 1,
217                TaskStatus::Failed => counts.failed += 1,
218                TaskStatus::Held => counts.held += 1,
219                TaskStatus::Blocked => counts.blocked += 1,
220            }
221        }
222        counts
223    }
224}
225
226/// One unit of work.
227#[derive(Debug, Clone, Serialize, Deserialize)]
228#[serde(deny_unknown_fields)]
229pub struct Task {
230    /// On-disk format version.
231    pub schema: u32,
232    /// Task id, e.g. `20260902-140501-a1b2`.
233    pub id: String,
234    /// One line, for lists and notifications.
235    pub title: String,
236    /// The task itself, handed to the graph verbatim.
237    pub instruction: String,
238    /// Repository to work in.
239    pub repo: PathBuf,
240    /// Who asked.
241    pub source: Source,
242    /// Higher runs first; ties break oldest-first so nothing starves.
243    #[serde(default)]
244    pub priority: i32,
245    /// Run this task alone: one implementer, no panel of judges to convince.
246    ///
247    /// `#[serde(default)]` so a queue file written before this field existed
248    /// still reads, as `false` - the ordinary multi-candidate competition,
249    /// unchanged. A task set to `solo` still runs the whole graph; only the
250    /// candidate count the daemon builds it with changes, and
251    /// [`crate::graph::Runner`] already collapses a single-candidate run to
252    /// implement → review → gate → merge on its own (see
253    /// [`crate::graph::Runner::review`]'s doc), so nothing about judging,
254    /// deliberation or voting had to change to support this.
255    #[serde(default)]
256    pub solo: bool,
257    /// Current state.
258    pub status: TaskStatus,
259    /// How many times this task has been claimed.
260    #[serde(default)]
261    pub attempts: usize,
262    /// Runs this task has produced, oldest first.
263    #[serde(default)]
264    pub runs: Vec<String>,
265    /// Why the last attempt did not land.
266    #[serde(default)]
267    pub last_error: Option<String>,
268    /// What a human hold is waiting on.
269    ///
270    /// `None` covers both the ordinary cases: a hold the loop makes itself
271    /// (out of attempts, or the disk gate closed) explains itself through
272    /// [`Task::last_error`] instead, and a human hold nobody bothered to
273    /// explain is still a valid hold. The queue has no way to express a
274    /// dependency between two tasks, so on the occasions a hold really is
275    /// "wait for that other task first", this is the only place that reason
276    /// survives - see [`Task::hold_manual`] and [`Task::release`].
277    ///
278    /// `#[serde(default)]` so a queue file written before this field existed
279    /// still reads, with no reason recorded rather than a parse error.
280    #[serde(default)]
281    pub hold_reason: Option<String>,
282    /// Who placed [`Task::hold_reason`].  `None` is a compatible old record;
283    /// see [`Task::operator_held`] for its deliberately conservative meaning.
284    #[serde(default)]
285    pub hold_source: Option<HoldSource>,
286    /// Diagnostic detail excerpted from the run that led to a hold - what a
287    /// human would have found opening `artifacts/` by hand, not the one-line
288    /// reason in [`Task::last_error`]. Set only when a run's own attempts are
289    /// exhausted and the task becomes [`TaskStatus::Held`]; `daemon` computes
290    /// it from the run's own record, since this module has no notion of a
291    /// run's internals. Bounded in length by the writer - see
292    /// `daemon::diagnostic` - so a verbose run cannot make this file grow
293    /// without limit.
294    ///
295    /// `#[serde(default)]` so a queue file written before this field existed
296    /// still reads, with no diagnostic recorded rather than a parse error.
297    #[serde(default)]
298    pub diagnostic: Option<String>,
299    /// What this task is waiting on: other task ids, unanswered
300    /// `crate::ask::Question` ids, or both. Non-empty exactly when
301    /// [`TaskStatus::Blocked`]; emptying it — see [`Task::unblock`] — is what
302    /// puts the task back at [`TaskStatus::Queued`].
303    ///
304    /// Set by `crate::conduct`'s decisions and cleared deterministically by
305    /// `crate::daemon` as each dependency resolves, never by a person. Never
306    /// `#[serde(default)]` is skipped: a queue file from before this field
307    /// existed has nothing to report here, and an empty list is exactly that.
308    #[serde(default)]
309    pub blocked_by: Vec<String>,
310    /// One line explaining the current [`Task::blocked_by`], written by
311    /// `crate::conduct`. Cleared whenever `blocked_by` empties.
312    #[serde(default)]
313    pub block_reason: Option<String>,
314    /// The status this task had the moment [`Task::block`] most recently
315    /// moved it to [`TaskStatus::Blocked`] — what [`Task::unblock`] restores
316    /// once nothing is left in `blocked_by`, instead of always landing on
317    /// [`TaskStatus::Queued`]. See [`SCHEMA`]'s doc for schema 4 on why this
318    /// exists: an answer to a question `crate::conduct` filed about a
319    /// [`TaskStatus::Held`] task must not itself be what puts the task back
320    /// in the competition queue.
321    ///
322    /// `#[serde(default)]` so a queue file written before this field existed
323    /// reads as `None`; [`Task::unblock`] treats that the same as a task
324    /// blocked straight from `Queued`, unless surviving hold evidence says
325    /// otherwise.
326    #[serde(default)]
327    pub blocked_from: Option<TaskStatus>,
328    /// Questions `crate::conduct` asked about this task that the operator has
329    /// since answered, oldest first — what was asked, and what they said.
330    ///
331    /// A blocking question's id leaves [`Task::blocked_by`] the moment
332    /// [`crate::ask::QuestionStatus::Answered`] is observed, but the id alone
333    /// tells nobody what was decided. This is what carries the answer's
334    /// *content* forward: into the next conductor prompt for this task, and
335    /// into the instruction handed to the next run — see `crate::daemon`'s
336    /// deterministic blocker resolution. Kept for the task's whole life, the
337    /// same as [`Task::runs`]: a release resets attempts, not evidence.
338    #[serde(default)]
339    pub answers: Vec<AnsweredQuestion>,
340    /// Ids of the `crate::triage` questions whose answer has been applied to
341    /// this task. Unlike [`Task::hold_reason`], [`Task::release`] and every
342    /// hold transition leave it alone, so an answer is applied at most once
343    /// however many times the task is held again. See [`SCHEMA`]'s doc for
344    /// schema 5. `#[serde(default)]` so an older record reads as empty.
345    #[serde(default)]
346    pub triage_applied: Vec<String>,
347    /// Ids of the `magi ask` questions whose [`crate::ask::ChoiceAction`] has
348    /// been applied to this task. Separate from [`Task::triage_applied`]
349    /// (a different producer) and, like it, untouched by [`Task::release`],
350    /// so one answer acts at most once however often the task is held again.
351    /// `#[serde(default)]` so an older record reads as empty.
352    #[serde(default)]
353    pub actions_applied: Vec<String>,
354    /// The operator's "resume" answer to a triage question, kept until the
355    /// task actually runs (or is done) so `crate::conduct` cannot silently
356    /// undo it and `crate::triage` can tell that a hold it sees now came
357    /// *after* the answer. See [`OperatorResume`]. `#[serde(default)]`.
358    #[serde(default)]
359    pub resume_override: Option<OperatorResume>,
360    /// Set by `crate::conduct` when it chooses `Review` recovery for a task
361    /// whose branch survived a blocked run: the branch to reopen with
362    /// `crate::graph::Runner::review` instead of competing from scratch.
363    ///
364    /// Requeues the task the same way [`Task::release`] does, so it is
365    /// picked up by the ordinary loop; `crate::daemon` reads this field once,
366    /// when it actually starts the run, and clears it either way — consumed
367    /// on success, dropped if the branch no longer exists by then. Never set
368    /// from the conductor's own words: `crate::daemon` derives the branch
369    /// name itself from the task's last run, so a hallucinated branch can
370    /// never reach here.
371    #[serde(default)]
372    pub review_branch: Option<String>,
373    /// A release deliberately starts a new competition instead of resuming
374    /// the prior run. History remains as evidence in `runs`.
375    #[serde(default)]
376    pub fresh_start: bool,
377    /// Marked by an operator (`magi task interrupt`) to ask `magi serve` to
378    /// run this one ahead of whatever it already has in flight, once
379    /// `[daemon] pause_for_interrupts` is on - see
380    /// `crate::daemon::advance_interrupt`. Never set by the loop itself, and
381    /// deliberately a different operation from [`Task::set_priority`]: a
382    /// priority only reorders the queue a claim has not reached yet, while
383    /// this asks a run already in flight to park at its next safe boundary
384    /// and step aside. `#[serde(default)]` so a queue file written before
385    /// this field existed still reads, as `false` - no task interrupts
386    /// anything unless asked to, exactly as before.
387    #[serde(default)]
388    pub interrupt: bool,
389    /// Marked by `magi task add --urgent`: `crate::daemon::poll` dispatches
390    /// this task through its own one-slot `urgent_sem` the moment it is
391    /// runnable, in addition to whatever is already running under the
392    /// ordinary `[daemon] max_concurrent_runs` pool - never instead of it,
393    /// and never by pausing or otherwise touching that run. This is the
394    /// opposite direction from [`Task::interrupt`]: that one asks a run
395    /// already in flight to step aside; this one never asks anything to
396    /// step aside, it only spends one additional, temporary concurrency
397    /// slot. The two are independent and may both be set on the same task,
398    /// but this exemption stops at `[daemon] pause_for_interrupts`'s own
399    /// park/resume handoff (75dd): while an interrupt sequence is actively
400    /// parking, running, or resuming - its own, or an unrelated task's -
401    /// `crate::daemon::interrupt_gate` withholds an urgent candidate exactly
402    /// like an ordinary one, never exempted. 75dd's "at most one run, ever,
403    /// at once" guarantee takes precedence, because the alternative is a run
404    /// still genuinely in flight (only *asked* to park, not yet gone) ending
405    /// up alongside a second one this feature let through - the very thing
406    /// that guarantee exists to rule out.
407    ///
408    /// `#[serde(default)]` so a queue file written before this field existed
409    /// still reads, as `false` - no task claims the urgent slot unless asked
410    /// to, exactly as before.
411    #[serde(default)]
412    pub urgent: bool,
413    /// Files (typically screenshots) copied into `<id>.attachments/` beside
414    /// this task's file by [`Queue::attach`], as bare validated names. The
415    /// copy is what makes them reach the implementer: the original may be
416    /// cleaned up long before the task runs. Untouched by [`Task::edit`].
417    /// `#[serde(default)]` so an older record reads as empty.
418    #[serde(default)]
419    pub attachments: Vec<String>,
420    /// Set on a task `crate::followup` filed from the findings a merged run
421    /// left open: which run and findings it carries, and how many follow-ups
422    /// deep it is. `None` for every ordinary task. `#[serde(default)]`.
423    #[serde(default)]
424    pub followup: Option<FollowUp>,
425    /// When the task was filed.
426    pub created_at: Timestamp,
427    /// Last change to this file.
428    pub updated_at: Timestamp,
429}
430
431/// What a follow-up task came from. See [`Task::followup`].
432#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
433pub struct FollowUp {
434    /// The merged run whose open findings this task carries.
435    pub run: String,
436    /// The queue task that run served, when it had one.
437    #[serde(default)]
438    pub origin_task: Option<String>,
439    /// The merged pull request.
440    pub pr: String,
441    /// Ids of the findings (e.g. `R3-1-1`) this task covers.
442    pub findings: Vec<String>,
443    /// Follow-up depth: 1 for a follow-up of an ordinary task's run, 2 for a
444    /// follow-up of that, and so on. An ordinary task is generation 0.
445    pub generation: u32,
446}
447
448/// A triage "resume" answer and what became of it. See
449/// [`Task::resume_override`].
450#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
451pub struct OperatorResume {
452    /// The triage question the operator answered.
453    pub question_id: String,
454    /// When the answer was applied.
455    pub at: Timestamp,
456    /// The reason `crate::conduct` gave for holding the task again after the
457    /// answer, if it did. The conductor may do this once.
458    #[serde(default)]
459    pub conductor_rehold: Option<String>,
460    /// The operator answered "resume" a second time, to the question about
461    /// that contradiction: the conductor may no longer hold this task.
462    #[serde(default)]
463    pub forced: bool,
464    /// Set when the answer was a structured `resume` action on a `magi ask`
465    /// choice ([`crate::ask::ChoiceAction::Resume`]): the run the operator
466    /// chose to continue. Its presence also stops `crate::conduct` from
467    /// requeuing the task into a fresh competition, which would discard
468    /// exactly the run that was named. Cleared with the rest of the record
469    /// once the task runs. `#[serde(default)]` so older records read as `None`.
470    #[serde(default)]
471    pub pinned_run: Option<String>,
472}
473
474/// One question `crate::conduct` asked about a task, and what the operator
475/// said back. See [`Task::answers`].
476#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
477pub struct AnsweredQuestion {
478    /// The question as asked, e.g. [`crate::ask::Question::summary`].
479    pub question: String,
480    /// What the operator answered.
481    pub answer: String,
482}
483
484impl Task {
485    /// File a new task. Persist it with [`Queue::put`].
486    pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
487        let now = Timestamp::now();
488        Self {
489            schema: SCHEMA,
490            id: new_id(),
491            title,
492            instruction,
493            repo,
494            source,
495            priority: 0,
496            solo: false,
497            status: TaskStatus::Queued,
498            attempts: 0,
499            runs: Vec::new(),
500            last_error: None,
501            hold_reason: None,
502            hold_source: None,
503            diagnostic: None,
504            blocked_by: Vec::new(),
505            block_reason: None,
506            blocked_from: None,
507            answers: Vec::new(),
508            triage_applied: Vec::new(),
509            actions_applied: Vec::new(),
510            resume_override: None,
511            review_branch: None,
512            fresh_start: false,
513            interrupt: false,
514            urgent: false,
515            attachments: Vec::new(),
516            followup: None,
517            created_at: now,
518            updated_at: now,
519        }
520    }
521
522    /// Short form used in reports, matching a run's short id.
523    pub fn short(&self) -> &str {
524        short(&self.id)
525    }
526
527    /// Record that triage question `question_id`'s answer has been applied.
528    pub fn mark_triage_applied(&mut self, question_id: &str) {
529        if !self.triage_applied(question_id) {
530            self.triage_applied.push(question_id.to_owned());
531        }
532    }
533
534    /// Has the action of `magi ask` question `question_id` already been
535    /// applied, or is it being applied now? See [`Task::actions_applied`].
536    pub fn action_applied(&self, question_id: &str) -> bool {
537        self.actions_applied.iter().any(|id| id == question_id)
538    }
539
540    /// Record that `magi ask` question `question_id`'s action was applied.
541    pub fn mark_action_applied(&mut self, question_id: &str) {
542        if !self.action_applied(question_id) {
543            self.actions_applied.push(question_id.to_owned());
544        }
545    }
546
547    /// Has triage question `question_id`'s answer already been applied?
548    pub fn triage_applied(&self, question_id: &str) -> bool {
549        self.triage_applied.iter().any(|id| id == question_id)
550    }
551
552    /// Record that a run has started for this task.
553    ///
554    /// Clears [`Task::interrupt`]: a mark to run ahead of whatever else is
555    /// in flight is fulfilled the moment this task actually gets its turn,
556    /// dispatched same as any other. Without this, a task whose run fails
557    /// and requeues - still `runnable`, still carrying the mark from its
558    /// first attempt - would keep re-triggering `crate::daemon`'s interrupt
559    /// scheduler and re-parking whatever it interrupted on every later
560    /// boundary, for as long as its attempts hold out, instead of the
561    /// one-shot "let this go next" the mark is meant to be.
562    pub fn start(&mut self, run: String) {
563        self.status = TaskStatus::Running;
564        self.attempts += 1;
565        self.runs.push(run);
566        self.last_error = None;
567        self.fresh_start = false;
568        self.interrupt = false;
569        // The answer has been honoured: the task got its turn.
570        self.resume_override = None;
571    }
572
573    /// Note that `run` serves this task without otherwise changing its life:
574    /// no status, attempt or interrupt change, unlike [`Task::start`]. For a
575    /// run somebody else started on the task's behalf (`magi run --task`);
576    /// idempotent. Returns whether anything was added.
577    pub fn link_run(&mut self, run: &str) -> bool {
578        if self.runs.iter().any(|r| r == run) {
579            return false;
580        }
581        self.runs.push(run.to_owned());
582        true
583    }
584
585    /// Record a successful run.
586    ///
587    /// Both `magi task done` and `POST /api/queue/{id}/done` can close a held
588    /// *or blocked* task directly, with no release in between, so this clears
589    /// `hold_reason` and `blocked_by`/`block_reason` the same way
590    /// [`Task::release`] does. Otherwise a task held for "waiting on 3ed9", or
591    /// blocked on a dependency that never actually finished, and then closed
592    /// as done without ever being released would still read as waiting on
593    /// something in `magi task show` and on its card, after it no longer is.
594    pub fn succeed(&mut self) {
595        self.status = TaskStatus::Done;
596        self.resume_override = None;
597        self.last_error = None;
598        self.hold_reason = None;
599        self.hold_source = None;
600        self.diagnostic = None;
601        self.blocked_by.clear();
602        self.block_reason = None;
603        self.blocked_from = None;
604    }
605
606    /// Finish the task because the run's whole change turned out to be on the
607    /// base already ([`crate::run::RunStatus::AlreadyInBase`]): `Done`, with
608    /// `note` kept in [`Task::last_error`] so `magi task show` says why, and
609    /// the attempt refunded. Nothing went wrong and nothing was spent
610    /// misbehaving, so this is neither [`Task::fail`] nor
611    /// [`Task::handed_off`]; and unlike [`Task::succeed`] it does not pretend
612    /// the run landed anything.
613    pub fn already_landed(&mut self, note: impl Into<String>) {
614        self.succeed();
615        self.attempts = self.attempts.saturating_sub(1);
616        self.last_error = Some(note.into());
617    }
618
619    /// Earlier attempts at this task that a later one has since made moot —
620    /// empty unless the task is [`TaskStatus::Done`] *and* `last_run_succeeded`
621    /// says `runs.last()` is actually why.
622    ///
623    /// `runs` is oldest first, and [`Task::start`] is the only thing that
624    /// pushes to it, always right before the attempt it names either succeeds
625    /// or fails; `succeed` itself never touches `runs`. So whenever `status`
626    /// is `Done` *because* the loop itself saw that last attempt land
627    /// (`Merged`/`Ready`), everything before it in the same list is a retry
628    /// this task no longer needs. But `succeed` is also reachable directly —
629    /// `magi task done`, the web UI's equivalent, and the conductor's
630    /// `Recovery::Done` all call it on a task in *any* status, including one
631    /// whose last recorded attempt never landed at all (closed by hand after
632    /// a merge magi's own loop never saw). `runs.last()` alone cannot tell
633    /// those two cases apart — that requires the caller to have actually
634    /// looked at that run's own `status`, which this module has no way to
635    /// do — so `last_run_succeeded` is the caller's answer to exactly that
636    /// question, not something this function can derive from `Task` alone.
637    ///
638    /// Deliberately narrower than [`Queue::superseded`]'s "every earlier
639    /// attempt has a later one" walk, which fires the moment a retry starts
640    /// even though the retry itself might still fail: that reading is right
641    /// for the web UI's "a newer attempt exists, go look at that one
642    /// instead" note, but wrong for deciding a run no longer needs a human's
643    /// attention, which is only true once the task's story has actually
644    /// ended well. Two Blocked runs sitting side by side while a third
645    /// attempt is still in flight must not be touched by this — see
646    /// `daemon::supersede_prior_runs`, the caller that turns this list into
647    /// rewritten `run.json` files.
648    pub fn superseded_attempts(&self, last_run_succeeded: bool) -> &[String] {
649        if self.status != TaskStatus::Done || !last_run_succeeded || self.runs.len() < 2 {
650            return &[];
651        }
652        &self.runs[..self.runs.len() - 1]
653    }
654
655    /// The attempt right after `run` in this task's history, if there is one.
656    ///
657    /// The one place "a later attempt replaced this run" is decided:
658    /// [`Queue::superseded`] and [`Queue::superseded_by`] (the web UI's
659    /// note) read it, and so does [`Task::earlier_attempts`], so what is shown
660    /// and what `crate::handover` acts on cannot disagree.
661    pub fn successor_of(&self, run: &str) -> Option<&String> {
662        let pos = self.runs.iter().position(|r| r == run)?;
663        self.runs.get(pos + 1)
664    }
665
666    /// Every run already recorded for this task, i.e. every run that an
667    /// attempt being minted *right now* supersedes.
668    ///
669    /// The new run is not in [`Task::runs`] yet when it decides what to take
670    /// over (`Task::start` pushes it afterwards), so "has a later attempt" is
671    /// true of everything listed here by the time the new run exists.
672    pub fn earlier_attempts(&self) -> &[String] {
673        &self.runs
674    }
675
676    /// Record a failed attempt. Out of attempts means held for a human, rather
677    /// than retried until the money runs out.
678    ///
679    /// Clears [`Task::diagnostic`] unconditionally: it belongs to whatever run
680    /// produced it, and a caller that has one for *this* attempt sets it
681    /// itself right after calling this, once it knows the task actually ended
682    /// up [`TaskStatus::Held`] - see `daemon::diagnostic`. Without the clear, a
683    /// task released after a diagnosed hold and then failed again for an
684    /// unrelated, undiagnosed reason (a config error, say) would go on
685    /// showing the previous run's diagnostic as if it explained the new one.
686    ///
687    /// Also records `why` into [`Task::hold_reason`] when this ends up
688    /// [`TaskStatus::Held`], so the notification and `magi task show` say the
689    /// same thing as [`Task::last_error`] instead of leaving the hold's own
690    /// reason blank.
691    pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
692        let why = why.into();
693        self.diagnostic = None;
694        self.status = if self.attempts >= max_attempts {
695            self.hold_source = Some(HoldSource::Machine);
696            self.hold_reason = Some(why.clone());
697            TaskStatus::Held
698        } else {
699            TaskStatus::Failed
700        };
701        self.last_error = Some(why);
702    }
703
704    /// Record an attempt that failed for a reason the task is not responsible
705    /// for - the agent CLIs ran out of quota and the judging panel collapsed.
706    ///
707    /// This refunds the attempt on purpose. A quota window closing at 4am must
708    /// not spend the backlog's retry budget: the operator would come back to a
709    /// queue of held tasks that were never actually judged, and would have to
710    /// release every one by hand to find out which had a real problem. The task
711    /// goes back to `Failed`, which the loop retries, so a reset quota picks the
712    /// work up where it stopped.
713    pub fn stall(&mut self, why: impl Into<String>) {
714        self.last_error = Some(why.into());
715        self.diagnostic = None;
716        self.attempts = self.attempts.saturating_sub(1);
717        self.status = TaskStatus::Failed;
718    }
719
720    /// Whether this held task may only be released by an operator.
721    ///
722    /// Old files did not record a source. Preserve every such hold rather
723    /// than guessing that it was automatic and risking duplicate work. New
724    /// automatic holds record [`HoldSource::Machine`] and remain recoverable.
725    pub fn operator_held(&self) -> bool {
726        self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
727    }
728
729    /// Take this task out of the loop's reach by an operator action.
730    ///
731    /// Clears `blocked_by`/`block_reason` unconditionally, the same as
732    /// [`Task::release`] and for the same reason its own comment already
733    /// gives: a human choosing to hold a *blocked* task overrides its wait
734    /// outright, the same as it overrides an ordinary hold. Without this, a
735    /// task held straight out of [`TaskStatus::Blocked`] - the web UI's "Hold"
736    /// button is reachable on a blocked task, same as "Mark done" - kept
737    /// reading as still waiting on a dependency it no longer had any claim on.
738    pub fn hold_manual(&mut self, reason: Option<String>) {
739        self.status = TaskStatus::Held;
740        if reason.is_some() {
741            self.hold_reason = reason;
742        }
743        self.hold_source = Some(HoldSource::Manual);
744        self.blocked_by.clear();
745        self.block_reason = None;
746        self.blocked_from = None;
747    }
748
749    /// Take this task out of the loop's reach during automatic recovery.
750    ///
751    /// Clears `blocked_by`/`block_reason` for the same reason
752    /// [`Task::hold_manual`] does.
753    pub fn hold_machine(&mut self, reason: Option<String>) {
754        self.status = TaskStatus::Held;
755        if reason.is_some() {
756            self.hold_reason = reason;
757        }
758        self.hold_source = Some(HoldSource::Machine);
759        self.blocked_by.clear();
760        self.block_reason = None;
761        self.blocked_from = None;
762    }
763
764    /// Block this task on other task ids and/or open question ids, chosen by
765    /// `crate::conduct`. Pure: the caller still owns writing it back with
766    /// [`Queue::put`].
767    ///
768    /// Records [`Task::blocked_from`] the first time this moves the task into
769    /// [`TaskStatus::Blocked`], and leaves it alone on a later call that adds
770    /// or replaces `blocked_by` while the task is already `Blocked` - a
771    /// second question about an already-blocked task must not overwrite the
772    /// status it should eventually return to with `Blocked` itself.
773    pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
774        if self.status != TaskStatus::Blocked {
775            self.blocked_from = Some(self.status);
776        }
777        self.status = TaskStatus::Blocked;
778        self.blocked_by = blocked_by;
779        self.block_reason = reason;
780    }
781
782    /// Remove one resolved dependency (a task id that became [`TaskStatus::Done`],
783    /// or a question id that became [`crate::ask::QuestionStatus::Answered`]).
784    /// Once nothing is left in [`Task::blocked_by`], the task returns to
785    /// whatever [`Task::blocked_from`] recorded - deciding *why* a task was
786    /// blocked was `crate::conduct`'s job, but noticing a dependency resolved
787    /// needs no model at all, and restoring the status it interrupted needs
788    /// nothing more than what `block` already wrote down.
789    ///
790    /// A task blocked while `Running` restores to [`TaskStatus::Queued`]
791    /// instead: whatever process was running it is gone by the time this
792    /// runs, so there is nothing left to resume. A task with no recorded
793    /// `blocked_from` - a pre-schema-4 record, or one blocked before this
794    /// field existed - falls back to [`TaskStatus::Held`] when it still
795    /// carries hold evidence ([`Task::hold_reason`] or [`Task::hold_source`],
796    /// neither ever cleared by `block`), and to `Queued` otherwise: the same
797    /// choice `block` itself would have recorded, reconstructed from what
798    /// survived.
799    ///
800    /// A no-op, on purpose, for a task that is not [`TaskStatus::Blocked`]:
801    /// `crate::daemon`'s deterministic resolver runs over every task on every
802    /// poll, and a task that moved on for some other reason must not be
803    /// dragged back by a stale id it still happens to carry.
804    pub fn unblock(&mut self, resolved_id: &str) {
805        if self.status != TaskStatus::Blocked {
806            return;
807        }
808        self.blocked_by.retain(|id| id != resolved_id);
809        if self.blocked_by.is_empty() {
810            self.status = match self.blocked_from {
811                Some(TaskStatus::Running) => TaskStatus::Queued,
812                Some(other) => other,
813                None if self.hold_reason.is_some() || self.hold_source.is_some() => {
814                    TaskStatus::Held
815                }
816                None => TaskStatus::Queued,
817            };
818            self.block_reason = None;
819            self.blocked_from = None;
820        }
821    }
822
823    /// Record that a question `crate::conduct` asked about this task has been
824    /// answered, so the answer's content — not just the fact that the
825    /// question is gone — reaches the next conductor prompt and the next
826    /// run's instruction. See [`Task::answers`].
827    pub fn record_answer(&mut self, question: String, answer: String) {
828        self.answers.push(AnsweredQuestion { question, answer });
829    }
830
831    /// Requeue this task to reopen its last run as a review-only pass against
832    /// `branch` (`crate::graph::Runner::review`) rather than competing from
833    /// scratch. See [`Task::review_branch`].
834    pub fn request_review(&mut self, branch: String) {
835        self.release();
836        self.review_branch = Some(branch);
837    }
838
839    /// Requeue after a conductor chose a new competition. Unlike an ordinary
840    /// operator release, this deliberately does not resume the old run.
841    pub fn requeue(&mut self) {
842        self.release();
843        // A fresh competition is the opposite of a review of the old branch.
844        self.review_branch = None;
845        self.fresh_start = true;
846    }
847
848    /// Hold this task because a takeover of `branch` was refused, keeping
849    /// [`Task::review_branch`] so that once the operator has cleaned up and
850    /// released it, the retry is a review of the same branch rather than a
851    /// resume or a fresh competition. Spends no attempt.
852    ///
853    /// This is the only way a machine hold keeps a `review_branch`: the
854    /// daemon takes it at the start of every attempt, so [`Task::release`]
855    /// relies on that to tell this hold from every other.
856    pub fn hold_for_handover(&mut self, branch: Option<String>, reason: String) {
857        if branch.is_some() {
858            self.review_branch = branch;
859        }
860        self.hold_machine(Some(reason));
861    }
862
863    /// Change how urgently this task should run next.
864    ///
865    /// Refused once the task is `running`: priority only feeds the sort
866    /// [`Queue::next_runnable`] does over tasks waiting to be claimed, and a
867    /// running task has already left that pool. Accepting the write anyway
868    /// would look like it worked while changing nothing until - and unless -
869    /// this attempt fails and the task becomes runnable again, which is a
870    /// surprise the phone should not hand back as a success.
871    pub fn set_priority(&mut self, priority: i32) -> Result<()> {
872        if self.status == TaskStatus::Running {
873            bail!(
874                "task {} is running; its priority cannot be changed until \
875                 this attempt finishes",
876                self.short()
877            );
878        }
879        self.priority = priority;
880        Ok(())
881    }
882
883    /// Mark (or unmark) this task to interrupt whatever `magi serve` already
884    /// has in flight, once `[daemon] pause_for_interrupts` is on. See
885    /// [`Task::interrupt`].
886    ///
887    /// Setting it is restricted to a task the loop could pick up on its own
888    /// right now - [`TaskStatus::runnable`] - for the same reason as
889    /// [`Task::set_priority`]: a task already `running` has been claimed, and
890    /// a task that is `done`, `held`, or `blocked` is not going to compete
891    /// for the daemon's attention regardless of this flag. Unlike priority,
892    /// this is never silently inert while `running` - it is refused outright,
893    /// because the entire feature this flag drives (`crate::daemon`'s
894    /// interrupt scheduler) is scoped to tasks still waiting to be claimed.
895    /// Clearing it back to `false` carries no such risk and is always
896    /// allowed, including on a task that moved on since it was set.
897    pub fn set_interrupt(&mut self, interrupt: bool) -> Result<()> {
898        if interrupt && !self.status.runnable() {
899            bail!(
900                "task {} is {}; only a queued or failed task can be marked \
901                 to interrupt",
902                self.short(),
903                self.status.as_str()
904            );
905        }
906        self.interrupt = interrupt;
907        Ok(())
908    }
909
910    /// Replace this task's title and instruction wholesale.
911    ///
912    /// Restricted to `queued` and `held`. A `running` task's instruction has
913    /// already been handed to the graph, so a run in flight and the file on
914    /// disk must not be allowed to disagree about what was asked; a `done` or
915    /// `failed` task is a record of what actually happened and editing it
916    /// after the fact would falsify that record. `id`, `created_at`,
917    /// `source`, and `runs` are left untouched on purpose - an edit stands in
918    /// for "delete and refile", and keeping the id, the timestamp, the
919    /// attribution, and the run history is the entire reason it exists
920    /// instead.
921    pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
922        if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
923            bail!(
924                "task {} is {}; only a queued or held task's instruction can \
925                 be edited",
926                self.short(),
927                self.status.as_str()
928            );
929        }
930        self.title = title;
931        self.instruction = instruction;
932        Ok(())
933    }
934
935    /// Record a run that produced a pull request without merging it.
936    ///
937    /// The task is held rather than retried, and it costs no further attempt
938    /// either way. The work the task asked for exists: it is sitting on a
939    /// branch, in a pull request, waiting for CI or for a person. Retrying
940    /// would spend the whole competition budget a second time and then race a
941    /// second branch against the pull request the first one opened - which is
942    /// exactly what happened to run 01c2, whose finished and green pull request
943    /// was re-competed from scratch four seconds after it opened.
944    ///
945    /// A pull request nobody merged is a request for a person, not a failure.
946    ///
947    /// Records `why` into [`Task::hold_reason`] as well as
948    /// [`Task::last_error`], so the notification centre and `magi task show`
949    /// say why the task is held rather than "no reason recorded". This is
950    /// also the path a verified no-op with an outstanding `magi ask` question
951    /// settles through - see `daemon::settle` and `daemon::settle_and_diagnose`,
952    /// which append the question id to `hold_reason` when one is still open
953    /// for the run.
954    pub fn handed_off(&mut self, why: impl Into<String>) {
955        let why = why.into();
956        self.diagnostic = None;
957        self.status = TaskStatus::Held;
958        self.hold_source = Some(HoldSource::Machine);
959        self.hold_reason = Some(why.clone());
960        self.last_error = Some(why);
961    }
962
963    /// Put a held or finished task back in line, with its attempt count reset
964    /// so a release is a real second chance rather than an instant re-hold.
965    /// The run history is kept: attempts reset, evidence does not.
966    pub fn release(&mut self) {
967        let refused_handover = self.status == TaskStatus::Held
968            && self.hold_source == Some(HoldSource::Machine)
969            && self.review_branch.is_some();
970        self.status = TaskStatus::Queued;
971        self.attempts = 0;
972        self.last_error = None;
973        // Otherwise the next person who holds this task reads a reason that
974        // belonged to whatever it was waiting on last time.
975        self.hold_reason = None;
976        self.hold_source = None;
977        self.diagnostic = None;
978        // A release also un-blocks: the dependency or question `blocked_by`
979        // named may still be unresolved, but a human (or `crate::conduct`)
980        // choosing to release the task overrides that wait outright, the same
981        // as it overrides an ordinary hold.
982        self.blocked_by.clear();
983        self.block_reason = None;
984        self.blocked_from = None;
985        // A machine hold that still names a review branch is a refused
986        // handover (see [`Task::hold_for_handover`]): the retry must review
987        // that same branch. Any other release drops it.
988        if !refused_handover {
989            self.review_branch = None;
990        }
991        self.fresh_start = false;
992    }
993}
994
995/// Copy `src` into `dir` as `name`, or `stem-2.ext`, `stem-3.ext`, ... when
996/// that is taken. `create_new` makes the collision check atomic and also
997/// catches a case-insensitive filesystem. Returns the stored name and path.
998fn copy_new(dir: &Path, src: &Path, name: &str) -> Result<(String, PathBuf)> {
999    use std::io::ErrorKind;
1000    let (stem, ext) = match name.rfind('.') {
1001        Some(i) if i > 0 => (&name[..i], &name[i..]),
1002        _ => (name, ""),
1003    };
1004    for n in 1u32.. {
1005        let candidate = if n == 1 {
1006            name.to_owned()
1007        } else {
1008            let suffix = format!("-{n}");
1009            let room = 64usize.saturating_sub(suffix.len() + ext.len());
1010            let stem: String = stem.chars().take(room).collect();
1011            format!("{stem}{suffix}{ext}")
1012        };
1013        if !crate::ask::valid_asset_name(&candidate) {
1014            bail!("no valid attachment name is left for `{name}`");
1015        }
1016        let path = dir.join(&candidate);
1017        match std::fs::OpenOptions::new()
1018            .write(true)
1019            .create_new(true)
1020            .open(&path)
1021        {
1022            Ok(mut out) => {
1023                let copied = std::fs::File::open(src)
1024                    .and_then(|mut input| std::io::copy(&mut input, &mut out));
1025                if let Err(e) = copied {
1026                    drop(out);
1027                    let _ = std::fs::remove_file(&path);
1028                    return Err(e).with_context(|| format!("copy {}", src.display()));
1029                }
1030                return Ok((candidate, path));
1031            }
1032            Err(e) if e.kind() == ErrorKind::AlreadyExists => continue,
1033            Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
1034        }
1035    }
1036    unreachable!("the counter never runs out")
1037}
1038
1039/// How long a task's write lock may stand before it is read as left behind by
1040/// a writer that died.
1041const TASK_LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
1042
1043/// Guard for [`Queue::lock_task`]; removes the lock file on drop.
1044struct TaskLock(PathBuf);
1045
1046impl Drop for TaskLock {
1047    fn drop(&mut self) {
1048        let _ = std::fs::remove_file(&self.0);
1049    }
1050}
1051
1052/// A queue on disk.
1053#[derive(Debug, Clone)]
1054pub struct Queue {
1055    root: PathBuf,
1056}
1057
1058impl Queue {
1059    /// The operator's queue, `<home>/queue`.
1060    pub fn open() -> Self {
1061        Self::at(crate::run::home().join("queue"))
1062    }
1063
1064    /// A queue at an explicit root. Tests use this; so could an operator who
1065    /// wants a queue per project.
1066    pub fn at(root: PathBuf) -> Self {
1067        Self { root }
1068    }
1069
1070    /// Directory holding the task files.
1071    pub fn root(&self) -> &Path {
1072        &self.root
1073    }
1074
1075    /// Path for one task id.
1076    pub fn path_of(&self, id: &str) -> PathBuf {
1077        self.root.join(format!("{id}.json"))
1078    }
1079
1080    /// Directory holding a task's attachments. A sibling of the task file,
1081    /// not a `*.json`, so [`Queue::list`] and [`Queue::revision`] never see it.
1082    pub fn attachments_dir(&self, id: &str) -> PathBuf {
1083        self.root.join(format!("{id}.attachments"))
1084    }
1085
1086    /// Copy each of `sources` into `task`'s attachment directory and record
1087    /// the stored names on `task` (persist with [`Queue::put`]). Returns the
1088    /// names, in order.
1089    ///
1090    /// Every source is checked (a readable file, a name passing
1091    /// [`crate::ask::valid_asset_name`]) before anything is copied, and a
1092    /// copy that fails part-way removes only what this call created. A name
1093    /// already taken is never overwritten: it becomes `stem-2.ext`,
1094    /// `stem-3.ext`, ... A name that does not validate is refused rather than
1095    /// sanitised - the operator renames the file.
1096    pub fn attach(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1097        let mut wanted = Vec::new();
1098        for src in sources {
1099            let name = src
1100                .file_name()
1101                .and_then(|n| n.to_str())
1102                .with_context(|| format!("`{}` has no usable file name", src.display()))?;
1103            if !crate::ask::valid_asset_name(name) {
1104                bail!(
1105                    "attachment name `{name}` must match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ \
1106                     with no `..`; rename the file and try again"
1107                );
1108            }
1109            if !src.is_file() {
1110                bail!("attachment `{}` is not a file", src.display());
1111            }
1112            wanted.push((src, name));
1113        }
1114        let dir = self.attachments_dir(&task.id);
1115        let existed = dir.is_dir();
1116        std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1117        let mut created: Vec<PathBuf> = Vec::new();
1118        let mut names = Vec::new();
1119        let mut copy_all = || -> Result<()> {
1120            for (src, name) in &wanted {
1121                let (stored, path) = copy_new(&dir, src, name)?;
1122                created.push(path);
1123                names.push(stored);
1124            }
1125            Ok(())
1126        };
1127        if let Err(e) = copy_all() {
1128            for path in &created {
1129                let _ = std::fs::remove_file(path);
1130            }
1131            if !existed {
1132                let _ = std::fs::remove_dir(&dir);
1133            }
1134            return Err(e);
1135        }
1136        task.attachments.extend(names.iter().cloned());
1137        Ok(names)
1138    }
1139
1140    /// [`Queue::attach`] then [`Queue::put`], undoing the copies when the
1141    /// record cannot be written. Without the undo the files sit in
1142    /// `<id>.attachments/` with no task pointing at them.
1143    ///
1144    /// Only what this call copied is removed (the names `attach` returned,
1145    /// and the directory itself only when it did not exist before and is
1146    /// empty), never an attachment an earlier save recorded. `task` is rolled
1147    /// back too, and the error the caller sees is `put`'s own. This relies on
1148    /// a failed `put` leaving the stored record untouched, which its
1149    /// write-then-rename guarantees.
1150    pub fn attach_and_put(&self, task: &mut Task, sources: &[PathBuf]) -> Result<Vec<String>> {
1151        let dir = self.attachments_dir(&task.id);
1152        let existed = dir.is_dir();
1153        let before = task.attachments.len();
1154        let names = self.attach(task, sources)?;
1155        if let Err(e) = self.put(task) {
1156            for name in &names {
1157                let _ = std::fs::remove_file(dir.join(name));
1158            }
1159            if !existed {
1160                let _ = std::fs::remove_dir(&dir);
1161            }
1162            task.attachments.truncate(before);
1163            return Err(e);
1164        }
1165        Ok(names)
1166    }
1167
1168    /// Absolute paths of `task`'s attachments, whatever shape this queue's
1169    /// root has. Absolute because the prompt hands them to an agent whose
1170    /// working directory is somewhere else entirely.
1171    pub fn attachment_paths(&self, task: &Task) -> Vec<PathBuf> {
1172        let dir = self.attachments_dir(&task.id);
1173        task.attachments
1174            .iter()
1175            .map(|n| {
1176                let p = dir.join(n);
1177                std::path::absolute(&p).unwrap_or(p)
1178            })
1179            .collect()
1180    }
1181
1182    /// Write a task, atomically, so a daemon killed mid-write leaves the
1183    /// previous state readable rather than a truncated file.
1184    ///
1185    /// Runs already recorded on disk that `task` does not carry are kept, not
1186    /// dropped: `magi run --task` links a run into a task from another
1187    /// process ([`Queue::link_run`]) while a daemon may hold a snapshot taken
1188    /// before it, and writing that snapshot back whole would erase the link.
1189    /// Nothing ever removes a run from a task, so the union loses nothing.
1190    pub fn put(&self, task: &mut Task) -> Result<()> {
1191        let _lock = self.lock_task(&task.id)?;
1192        self.put_unlocked(task)
1193    }
1194
1195    /// Write `task` only if no task with its id exists yet; `Ok(false)` when
1196    /// one does. Unlike [`Queue::put`] it never replaces a record, so a
1197    /// producer with a deterministic id cannot rewind a task that has since
1198    /// started running or finished.
1199    pub fn create_new(&self, task: &mut Task) -> Result<bool> {
1200        let _lock = self.lock_task(&task.id)?;
1201        if self.path_of(&task.id).exists() {
1202            return Ok(false);
1203        }
1204        self.put_unlocked(task)?;
1205        Ok(true)
1206    }
1207
1208    /// [`Queue::put`] for a caller already holding [`Queue::lock_task`].
1209    fn put_unlocked(&self, task: &mut Task) -> Result<()> {
1210        if let Ok(stored) = read_path(&self.path_of(&task.id)) {
1211            for run in stored.runs {
1212                if !task.runs.contains(&run) {
1213                    task.runs.push(run);
1214                }
1215            }
1216        }
1217        task.updated_at = Timestamp::now();
1218        std::fs::create_dir_all(&self.root)
1219            .with_context(|| format!("create {}", self.root.display()))?;
1220        let body = serde_json::to_string_pretty(task).context("serialize task")?;
1221        let path = self.path_of(&task.id);
1222        let tmp = path.with_extension("json.tmp");
1223        std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1224        std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1225        // A machine hold is news for the notification centre, filed beside
1226        // this queue (`<home>/queue` -> `<home>/notifications`) rather than
1227        // through a process-global, so a queue in a temp directory notifies
1228        // into that directory.
1229        if let (Some(notice), Some(home)) = (
1230            crate::notices::task_held(task),
1231            self.root.parent().filter(|p| !p.as_os_str().is_empty()),
1232        ) {
1233            crate::notices::raise_in(home, notice);
1234        }
1235        Ok(())
1236    }
1237
1238    /// Append `run` to the runs of the task `id` (or an unambiguous prefix),
1239    /// leaving everything else about it alone — see [`Task::link_run`]. Read
1240    /// and written back in one breath, because [`Queue::put`] replaces the
1241    /// whole record. Returns the task as stored.
1242    pub fn link_run(&self, id: &str, run: &str) -> Result<Task> {
1243        let id = self.resolve_id(id)?;
1244        // Read inside the lock, so the record written back is the current one
1245        // and a daemon's write cannot slip between the read and the write.
1246        let _lock = self.lock_task(&id)?;
1247        let mut task = self.get(&id)?;
1248        if task.link_run(run) {
1249            self.put_unlocked(&mut task)?;
1250        }
1251        Ok(task)
1252    }
1253
1254    /// Exclusive right to rewrite task `id`'s record, held until the guard
1255    /// drops. Every writer ([`Queue::put`], [`Queue::link_run`]) takes it, so
1256    /// a process linking a run and the daemon saving a transition are
1257    /// serialized instead of last-writer-wins. A lock older than
1258    /// [`TASK_LOCK_STALE`] belongs to a writer that died mid-update and is
1259    /// broken. Not the claim lock (`<id>.lock`), which a daemon holds for a
1260    /// whole run.
1261    fn lock_task(&self, id: &str) -> Result<TaskLock> {
1262        std::fs::create_dir_all(&self.root)
1263            .with_context(|| format!("create {}", self.root.display()))?;
1264        let path = self.root.join(format!("{id}.write-lock"));
1265        let started = std::time::Instant::now();
1266        loop {
1267            match std::fs::OpenOptions::new()
1268                .write(true)
1269                .create_new(true)
1270                .open(&path)
1271            {
1272                Ok(_) => return Ok(TaskLock(path)),
1273                Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1274                    let stale = std::fs::metadata(&path)
1275                        .and_then(|m| m.modified())
1276                        .ok()
1277                        .and_then(|t| t.elapsed().ok())
1278                        .is_some_and(|age| age > TASK_LOCK_STALE);
1279                    if stale {
1280                        let _ = std::fs::remove_file(&path);
1281                    } else if started.elapsed() > TASK_LOCK_STALE {
1282                        bail!("could not lock task {id}");
1283                    } else {
1284                        std::thread::sleep(std::time::Duration::from_millis(15));
1285                    }
1286                }
1287                Err(e) => return Err(e).with_context(|| format!("lock {}", path.display())),
1288            }
1289        }
1290    }
1291
1292    /// Load a task by id or unambiguous id prefix.
1293    pub fn get(&self, id: &str) -> Result<Task> {
1294        let resolved = self.resolve_id(id)?;
1295        read_path(&self.path_of(&resolved))
1296    }
1297
1298    /// Remove a task, and the claim lock that belongs to it.
1299    ///
1300    /// `in_flight` comes from the caller — a live daemon's heartbeat naming
1301    /// this task — because the task's own `running` status cannot answer the
1302    /// question. A daemon killed mid-competition leaves the status at
1303    /// `running` and an orphaned `.lock` behind, and a guard that trusted
1304    /// either would make the task undeletable for good: the phone showed
1305    /// exactly that, refusing a task whose daemon had been gone for an hour.
1306    ///
1307    /// So the lock is removed with the task rather than respected. Any lock
1308    /// still there once no live daemon claims the task is by definition stale,
1309    /// and leaving it would make a deleted task look claimed to
1310    /// [`Queue::claim`] and to whoever reads the directory.
1311    ///
1312    /// Anything still `blocked` on the id just deleted is quarantined to a
1313    /// machine hold in the same call - see [`Removal::quarantined`] - rather
1314    /// than left to wait on a dependency that no longer exists. Best-effort:
1315    /// a dependent claimed by something else right now, or one whose write
1316    /// fails, is simply left for `crate::daemon::resolve_blockers`'s own poll
1317    /// (or `crate::triage::run_once`) to catch on its own next pass, and does
1318    /// not fail this removal.
1319    ///
1320    /// `questions` is the store [`missing_blockers`] checks a `blocked_by` id
1321    /// against before calling it gone - the same store the caller already
1322    /// resolves `id`'s own home from, passed in rather than reopened here so
1323    /// a test queue at an explicit root is never quarantined against the
1324    /// operator's real questions directory.
1325    pub fn remove(&self, id: &str, in_flight: bool, questions: &Questions) -> Result<Removal> {
1326        let resolved = self.resolve_id(id)?;
1327        if in_flight {
1328            bail!("task {resolved} is being run by a live daemon right now");
1329        }
1330        self.remove_record_with_attachments(&resolved, |p| std::fs::remove_file(p))?;
1331        let quarantined = self.quarantine_dependents_of(&resolved, questions);
1332        Ok(Removal {
1333            id: resolved,
1334            quarantined,
1335        })
1336    }
1337
1338    /// The attachment-and-record half of [`Queue::remove`]. `remove_record` is
1339    /// how the record file is deleted; production passes `remove_file`, and a
1340    /// test passes a failing one to prove the attachments come back.
1341    fn remove_record_with_attachments(
1342        &self,
1343        resolved: &str,
1344        remove_record: impl FnOnce(&Path) -> std::io::Result<()>,
1345    ) -> Result<()> {
1346        // Neither step may strand the other. The attachments are set aside by a
1347        // rename (cheap, and undone if the record cannot be removed), then the
1348        // record goes, and only then is the set-aside directory deleted for
1349        // real. A failure before the record is gone loses nothing; a failure
1350        // after it leaves at most an orphan `.removing` directory, which the
1351        // next removal sweeps.
1352        self.sweep_removed_attachments();
1353        let attachments = self.attachments_dir(resolved);
1354        let aside = self.root.join(format!("{resolved}.attachments.removing"));
1355        let moved = match std::fs::rename(&attachments, &aside) {
1356            Ok(()) => true,
1357            Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
1358            Err(e) => {
1359                return Err(e).with_context(|| format!("remove {}", attachments.display()));
1360            }
1361        };
1362        let path = self.path_of(resolved);
1363        if let Err(e) = remove_record(&path) {
1364            if moved {
1365                let _ = std::fs::rename(&aside, &attachments);
1366            }
1367            return Err(e).with_context(|| format!("remove {}", path.display()));
1368        }
1369        if moved {
1370            if let Err(e) = std::fs::remove_dir_all(&aside) {
1371                tracing::warn!("leftover attachments {}: {e}", aside.display());
1372            }
1373        }
1374        let lock = self.lock_path(resolved);
1375        if let Err(e) = std::fs::remove_file(&lock) {
1376            if e.kind() != std::io::ErrorKind::NotFound {
1377                return Err(e).with_context(|| format!("remove {}", lock.display()));
1378            }
1379        }
1380        Ok(())
1381    }
1382
1383    /// Delete `*.attachments.removing` directories a previous [`Queue::remove`]
1384    /// could not finish deleting. The task record is already gone by then, so
1385    /// the id no longer resolves and the removal cannot be retried by name;
1386    /// every later removal sweeps them instead. A directory whose record still
1387    /// exists belongs to a removal in progress and is left alone. Best-effort.
1388    fn sweep_removed_attachments(&self) {
1389        let Ok(entries) = std::fs::read_dir(&self.root) else {
1390            return;
1391        };
1392        for entry in entries.flatten() {
1393            let name = entry.file_name();
1394            let name = name.to_string_lossy();
1395            let Some(id) = name.strip_suffix(".attachments.removing") else {
1396                continue;
1397            };
1398            // A record still on disk means a removal is mid-flight (or about
1399            // to roll back): the directory is its to delete or restore.
1400            if !self.path_of(id).exists() {
1401                if let Err(e) = std::fs::remove_dir_all(entry.path()) {
1402                    tracing::warn!("leftover attachments {}: {e}", entry.path().display());
1403                }
1404            }
1405        }
1406    }
1407
1408    /// Move every `blocked` task naming `dependency` in its own `blocked_by`
1409    /// to a machine hold, now that `dependency`'s own file is gone. See
1410    /// [`Queue::remove`]'s own doc for why this is best-effort.
1411    fn quarantine_dependents_of(&self, dependency: &str, questions: &Questions) -> Vec<String> {
1412        let mut quarantined = Vec::new();
1413        for listed in self.list() {
1414            if listed.status != TaskStatus::Blocked
1415                || !listed.blocked_by.iter().any(|b| b == dependency)
1416            {
1417                continue;
1418            }
1419            let Ok(_claim) = self.claim(&listed.id) else {
1420                continue;
1421            };
1422            let Ok(mut task) = self.get(&listed.id) else {
1423                continue;
1424            };
1425            if task.status != TaskStatus::Blocked
1426                || !task.blocked_by.iter().any(|b| b == dependency)
1427            {
1428                continue;
1429            }
1430            let missing = missing_blockers(self, questions, &task.blocked_by);
1431            let language = crate::lang::of_repo(&task.repo);
1432            task.hold_machine(Some(missing_blocker_hold_reason_in(
1433                &task.blocked_by,
1434                &missing,
1435                &language,
1436            )));
1437            if self.put(&mut task).is_ok() {
1438                quarantined.push(task.id.clone());
1439            }
1440        }
1441        quarantined
1442    }
1443
1444    /// Path of the claim lock for a task. One definition, so `claim` and
1445    /// `remove` cannot end up naming different files.
1446    fn lock_path(&self, id: &str) -> PathBuf {
1447        self.root.join(format!("{id}.lock"))
1448    }
1449
1450    /// Every task on disk, highest priority first and newest first within a
1451    /// priority. This is what `magi task list` and `GET /api/queue` print, so
1452    /// a raised priority has to move a task here the moment it is saved, not
1453    /// only in [`Queue::next_runnable`]'s own ordering - the operator reading
1454    /// the backlog and the loop about to drain it must agree on what "first"
1455    /// means. Every existing task defaults to priority 0, so this is a no-op
1456    /// change from the old newest-first order for a queue nobody has
1457    /// reprioritised.
1458    ///
1459    /// Unreadable files are skipped rather than fatal: one corrupt task must
1460    /// not take the queue - or the web UI, or an unattended daemon - down
1461    /// with it.
1462    pub fn list(&self) -> Vec<Task> {
1463        let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
1464            .into_iter()
1465            .flatten()
1466            .flatten()
1467            .map(|e| e.path())
1468            .filter(|p| p.extension().is_some_and(|x| x == "json"))
1469            .filter_map(|p| read_path(&p).ok())
1470            .collect();
1471        tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
1472        tasks
1473    }
1474
1475    /// Runs that a later attempt at the same task replaced, mapped to the id
1476    /// of the attempt that replaced them.
1477    ///
1478    /// A task keeps its attempts in order, and the deck showed them as two
1479    /// cards with the same title and no hint which was which: yukimemi asked
1480    /// why `stalled` and `blocked` appeared twice for one task, and the
1481    /// answer - "those are two tries, and the second one exists because of a
1482    /// bug since fixed" - was not on the screen anywhere.
1483    ///
1484    /// Read from the queue rather than stored on the run, because the
1485    /// ordering is the queue's fact: a `RunState` has no idea another attempt
1486    /// happened after it.
1487    pub fn superseded(&self) -> HashMap<String, String> {
1488        let mut by = HashMap::new();
1489        for task in self.list() {
1490            for earlier in &task.runs {
1491                if let Some(later) = task.successor_of(earlier) {
1492                    by.insert(earlier.clone(), later.clone());
1493                }
1494            }
1495        }
1496        by
1497    }
1498
1499    /// Whether `run` is an earlier attempt a later one replaced, and if so
1500    /// the id of that later attempt.
1501    ///
1502    /// Same walk as [`Queue::superseded`], narrowed to one run: a run detail
1503    /// page asks about exactly one run at a time, and this keeps that call
1504    /// site from building (and discarding) the whole map's `HashMap` just to
1505    /// read one entry out of it.
1506    pub fn superseded_by(&self, run: &str) -> Option<String> {
1507        for task in self.list() {
1508            if task.runs.iter().any(|r| r == run) {
1509                return task.successor_of(run).cloned();
1510            }
1511        }
1512        None
1513    }
1514
1515    /// The task's own most recent attempt, when `run` belongs to that task
1516    /// but is not already that attempt.
1517    ///
1518    /// Distinct from [`Queue::superseded_by`], which names only the very
1519    /// next attempt: a chain of retries (A superseded by B superseded by C)
1520    /// leaves an older run pointing at an intermediate one that may itself
1521    /// be unresolved, and a run's own detail page needs to know where the
1522    /// task's story currently stands - the chain's current head, C - not an
1523    /// attempt in the middle of it that a client would otherwise have to
1524    /// walk to by hand.
1525    pub fn latest_attempt(&self, run: &str) -> Option<String> {
1526        for task in self.list() {
1527            if task.runs.iter().any(|r| r == run) {
1528                return task.runs.last().filter(|last| **last != run).cloned();
1529            }
1530        }
1531        None
1532    }
1533
1534    /// The task a daemon should run next, or `None` when the queue is idle.
1535    ///
1536    /// Highest priority first, oldest first within a priority, so a burst of
1537    /// agent-filed work cannot starve the task a human filed this morning.
1538    pub fn next_runnable(&self) -> Option<Task> {
1539        let mut runnable: Vec<Task> = self
1540            .list()
1541            .into_iter()
1542            .filter(|t| t.status.runnable())
1543            .collect();
1544        runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
1545        runnable.into_iter().next()
1546    }
1547
1548    /// Take exclusive ownership of a task.
1549    ///
1550    /// The lock is a `create_new` file next to the task, which is atomic on
1551    /// every platform magi targets. It exists so two daemons - or a daemon and
1552    /// a human running `magi run` - cannot drive one task into two competing
1553    /// runs. The returned guard releases on drop, including on panic.
1554    pub fn claim(&self, id: &str) -> Result<Claim> {
1555        std::fs::create_dir_all(&self.root)
1556            .with_context(|| format!("create {}", self.root.display()))?;
1557        let path = self.lock_path(id);
1558        match std::fs::OpenOptions::new()
1559            .write(true)
1560            .create_new(true)
1561            .open(&path)
1562        {
1563            Ok(mut f) => {
1564                use std::io::Write as _;
1565                // Best effort: the pid is for the human looking at a stale lock.
1566                let _ = writeln!(f, "{}", std::process::id());
1567                Ok(Claim { path })
1568            }
1569            Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1570                bail!("task {id} is already claimed ({} exists)", path.display())
1571            }
1572            Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
1573        }
1574    }
1575
1576    /// Expand an id prefix to exactly one task id.
1577    pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1578        if self.path_of(prefix).is_file() {
1579            return Ok(prefix.to_owned());
1580        }
1581        let hits: Vec<String> = self
1582            .list()
1583            .into_iter()
1584            .map(|t| t.id)
1585            .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1586            .collect();
1587        match hits.len() {
1588            1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1589            0 => bail!("no task matches `{prefix}`"),
1590            _ => bail!(
1591                "`{prefix}` matches {} tasks: {}",
1592                hits.len(),
1593                hits.join(", ")
1594            ),
1595        }
1596    }
1597
1598    /// Change detection token for the queue.
1599    ///
1600    /// Combines file names and modification times of all task files in the
1601    /// queue, so adding, modifying, or deleting any task — even an older one —
1602    /// moves the revision and notifies connected clients via the change stream.
1603    /// Returns 0 when the queue is completely empty.
1604    pub fn revision(&self) -> u64 {
1605        use std::hash::{Hash as _, Hasher as _};
1606
1607        let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
1608            .into_iter()
1609            .flatten()
1610            .flatten()
1611            .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
1612            .filter_map(|e| {
1613                let name = e.file_name().to_string_lossy().into_owned();
1614                let mtime = e
1615                    .metadata()
1616                    .ok()?
1617                    .modified()
1618                    .ok()?
1619                    .duration_since(std::time::UNIX_EPOCH)
1620                    .ok()?
1621                    .as_millis() as u64;
1622                Some((name, mtime))
1623            })
1624            .collect();
1625
1626        if entries.is_empty() {
1627            return 0;
1628        }
1629
1630        entries.sort_unstable();
1631        let mut hasher = std::hash::DefaultHasher::new();
1632        for (name, mtime) in &entries {
1633            name.hash(&mut hasher);
1634            mtime.hash(&mut hasher);
1635        }
1636        let h = hasher.finish();
1637        if h == 0 { 1 } else { h }
1638    }
1639}
1640
1641/// What [`Queue::remove`] did, beyond deleting the named task's own file.
1642#[derive(Debug, Clone)]
1643pub struct Removal {
1644    /// The id actually removed - `id` expanded from a prefix, if it was one.
1645    pub id: String,
1646    /// Every `blocked` task that named [`Removal::id`] in its own
1647    /// `blocked_by` and was moved to a machine hold as a result, rather than
1648    /// left waiting on a dependency this call just erased.
1649    pub quarantined: Vec<String>,
1650}
1651
1652/// Exclusive ownership of a task, released on drop.
1653#[derive(Debug)]
1654pub struct Claim {
1655    path: PathBuf,
1656}
1657
1658impl Drop for Claim {
1659    fn drop(&mut self) {
1660        let _ = std::fs::remove_file(&self.path);
1661    }
1662}
1663
1664/// The first line of a task, trimmed to a title. Used when the caller gives a
1665/// body but no title, which is the normal case for an agent piping a file in.
1666pub fn title_from(instruction: &str, max: usize) -> String {
1667    // The first non-blank line, whatever it is. A markdown heading is the
1668    // task's own summary - agents pipe in `# Rework the config loader` and mean
1669    // exactly that - so it is preferred over the prose beneath it rather than
1670    // skipped as decoration. Leading list and heading markers are stripped
1671    // because they are syntax, not words.
1672    let line = instruction
1673        .lines()
1674        .map(str::trim)
1675        .find(|l| !l.is_empty())
1676        .unwrap_or("(empty task)")
1677        .trim_start_matches(['#', '-', '*', '>', ' '])
1678        .trim();
1679    if line.is_empty() {
1680        return "(empty task)".to_owned();
1681    }
1682    if line.chars().count() <= max {
1683        return line.to_owned();
1684    }
1685    let head: String = line.chars().take(max.saturating_sub(1)).collect();
1686    format!("{head}…")
1687}
1688
1689fn read_path(path: &Path) -> Result<Task> {
1690    let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1691    let task: Task =
1692        serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1693    // Greater-than, not not-equal: every field added since schema 1 carries
1694    // `#[serde(default)]`, so an older task has nothing to say about it and
1695    // defaulting is exactly as good a reading as a value that build never had
1696    // a chance to write. Only a schema *ahead* of this build - a meaning it
1697    // cannot possibly know - is refused rather than guessed at.
1698    if task.schema > SCHEMA {
1699        bail!(
1700            "task {} was written by a different magi (schema {}, this build \
1701             speaks {SCHEMA})",
1702            task.id,
1703            task.schema
1704        );
1705    }
1706    Ok(task)
1707}
1708
1709/// Ids inside a `blocked_by` list that name neither an existing task file nor
1710/// an existing question file - a dependency deleted (`magi task rm`, or by
1711/// hand) while something was still waiting on it.
1712///
1713/// Existence is decided by [`Queue::path_of`]/[`Questions::path_of`]
1714/// `is_file()` alone, never by [`Queue::get`]/[`Questions::get`] succeeding:
1715/// those also fail on a merely unreadable file - mid-write, corrupt, or from
1716/// a schema ahead of this build (see [`read_path`]) - and misreading "cannot
1717/// read it right now" as "it was deleted" would quarantine a task over a
1718/// transient failure. `blocked_by` always carries a full id, written by
1719/// `crate::conduct` or `crate::triage` from a real task's or question's own
1720/// `id`/`short`, never a prefix a caller typed - so the exact-path check is
1721/// complete on its own, with no [`Queue::resolve_id`] fallback needed.
1722pub fn missing_blockers(
1723    queue: &Queue,
1724    questions: &Questions,
1725    blocked_by: &[String],
1726) -> Vec<String> {
1727    blocked_by
1728        .iter()
1729        .filter(|id| !queue.path_of(id).is_file() && !questions.path_of(id).is_file())
1730        .cloned()
1731        .collect()
1732}
1733
1734/// The `hold_reason` text for a task quarantined because one or more of its
1735/// `blocked_by` ids no longer exist. Shared by `crate::daemon::resolve_blockers`,
1736/// `crate::triage::run_once`, and [`Queue::remove`]'s own dependent
1737/// quarantine, so the three call sites read as the same event to an operator
1738/// looking at `magi task show` rather than three different wordings for it.
1739///
1740/// Names the full original `blocked_by` list, not just `missing` - a task
1741/// quarantined here can also have named a dependency that was still
1742/// perfectly valid, and [`Task::hold_machine`] clears `blocked_by` on the way
1743/// in, so this text is the only place that information survives for an
1744/// operator deciding whether to release the task outright.
1745pub fn missing_blocker_hold_reason(blocked_by: &[String], missing: &[String]) -> String {
1746    missing_blocker_hold_reason_in(blocked_by, missing, "en")
1747}
1748
1749/// [`missing_blocker_hold_reason`] in `language` (Japanese, else English).
1750pub fn missing_blocker_hold_reason_in(
1751    blocked_by: &[String],
1752    missing: &[String],
1753    language: &str,
1754) -> String {
1755    if crate::lang::is_japanese(language) {
1756        format!(
1757            "{} を待っていましたが、{} はディスク上に存在しません - `magi task triage` を参照",
1758            blocked_by.join(", "),
1759            missing.join(", "),
1760        )
1761    } else {
1762        format!(
1763            "blocked on {} but {} no longer exist(s) on disk - see `magi task triage`",
1764            blocked_by.join(", "),
1765            missing.join(", "),
1766        )
1767    }
1768}
1769
1770/// The last dash-separated part of an id, the form people say aloud.
1771pub fn short(id: &str) -> &str {
1772    id.split('-').next_back().unwrap_or(id)
1773}
1774
1775fn new_id() -> String {
1776    let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1777    let seed = crate::rng::entropy();
1778    format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1779}
1780
1781#[cfg(test)]
1782mod tests {
1783    #[test]
1784    fn an_already_landed_task_is_done_with_its_attempt_refunded() {
1785        let mut t = task("relanded");
1786        t.attempts = 1;
1787        t.status = TaskStatus::Running;
1788        t.already_landed("already in main as 0e368de");
1789        assert_eq!(t.status, TaskStatus::Done);
1790        assert_eq!(t.attempts, 0);
1791        assert!(t.hold_reason.is_none());
1792        assert_eq!(t.last_error.as_deref(), Some("already in main as 0e368de"));
1793    }
1794
1795    #[test]
1796    fn missing_blocker_reason_follows_the_language() {
1797        let b = vec!["a".to_owned()];
1798        let en = missing_blocker_hold_reason_in(&b, &b, "en");
1799        assert_eq!(en, missing_blocker_hold_reason(&b, &b));
1800        assert!(en.starts_with("blocked on a"));
1801        assert!(missing_blocker_hold_reason_in(&b, &b, "ja").contains("存在しません"));
1802        assert_eq!(missing_blocker_hold_reason_in(&b, &b, "de"), en);
1803    }
1804
1805    use super::*;
1806
1807    #[test]
1808    fn triage_applied_survives_release_and_old_records_read_as_empty() {
1809        let mut t = Task::new(
1810            "t".to_owned(),
1811            "i".to_owned(),
1812            PathBuf::from("r"),
1813            Source::Human,
1814        );
1815        t.mark_triage_applied("q1");
1816        t.mark_triage_applied("q1");
1817        t.hold_machine(Some("x".to_owned()));
1818        t.release();
1819        assert_eq!(t.triage_applied, ["q1"]);
1820        assert!(t.triage_applied("q1") && !t.triage_applied("q2"));
1821
1822        let mut v = serde_json::to_value(&t).unwrap();
1823        v.as_object_mut().unwrap().remove("triage_applied");
1824        let old: Task = serde_json::from_value(v).unwrap();
1825        assert!(old.triage_applied.is_empty());
1826    }
1827
1828    #[test]
1829    fn task_counts_of_empty_is_all_zero() {
1830        assert_eq!(TaskCounts::of(&[]), TaskCounts::default());
1831    }
1832
1833    #[test]
1834    fn task_counts_of_tallies_every_status() {
1835        let mut queued = Task::new(
1836            "q".to_owned(),
1837            "i".to_owned(),
1838            PathBuf::from("."),
1839            Source::Human,
1840        );
1841        queued.status = TaskStatus::Queued;
1842        let mut running = queued.clone();
1843        running.status = TaskStatus::Running;
1844        let mut done = queued.clone();
1845        done.status = TaskStatus::Done;
1846        let mut failed = queued.clone();
1847        failed.status = TaskStatus::Failed;
1848        let mut held = queued.clone();
1849        held.status = TaskStatus::Held;
1850        let mut blocked = queued.clone();
1851        blocked.status = TaskStatus::Blocked;
1852
1853        let counts = TaskCounts::of(&[queued, running, done.clone(), done, failed, held, blocked]);
1854        assert_eq!(
1855            counts,
1856            TaskCounts {
1857                queued: 1,
1858                running: 1,
1859                done: 2,
1860                failed: 1,
1861                held: 1,
1862                blocked: 1,
1863            }
1864        );
1865    }
1866
1867    /// A queue of its own, with no process-global state - which is the point of
1868    /// `Queue::at`, and why these can run in parallel.
1869    fn queue() -> (tempfile::TempDir, Queue) {
1870        let dir = tempfile::tempdir().unwrap();
1871        let q = Queue::at(dir.path().join("queue"));
1872        (dir, q)
1873    }
1874
1875    #[test]
1876    fn putting_a_machine_held_task_files_a_notification_beside_the_queue() {
1877        let dir = tempfile::tempdir().unwrap();
1878        let q = Queue::at(dir.path().join("queue"));
1879        let mut t = task("held");
1880        q.put(&mut t).unwrap();
1881        assert_eq!(
1882            crate::notices::Notices::at(dir.path().join("notifications"))
1883                .list()
1884                .len(),
1885            0
1886        );
1887        t.hold_machine(Some("out of attempts".to_owned()));
1888        q.put(&mut t).unwrap();
1889        let listed = crate::notices::Notices::at(dir.path().join("notifications")).list();
1890        assert_eq!(listed.len(), 1);
1891        assert!(listed[0].message.contains("out of attempts"));
1892    }
1893
1894    fn task(title: &str) -> Task {
1895        Task::new(
1896            title.to_owned(),
1897            format!("do {title}"),
1898            PathBuf::from("."),
1899            Source::Human,
1900        )
1901    }
1902
1903    #[test]
1904    fn earlier_attempts_is_every_recorded_run_and_agrees_with_the_display() {
1905        let mut t = task("retried");
1906        assert!(
1907            t.earlier_attempts().is_empty(),
1908            "a first attempt takes nothing over"
1909        );
1910        t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1911        assert_eq!(t.earlier_attempts(), ["aaaa", "bbbb"]);
1912        // What the web note shows and what a takeover acts on are one rule.
1913        assert_eq!(t.successor_of("aaaa"), Some(&"bbbb".to_owned()));
1914        assert_eq!(t.successor_of("bbbb"), None);
1915        assert_eq!(t.successor_of("zzzz"), None);
1916    }
1917
1918    #[test]
1919    fn superseded_by_names_the_next_attempt_and_none_for_the_last() {
1920        let (_dir, q) = queue();
1921        let mut t = task("retried");
1922        t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1923        q.put(&mut t).unwrap();
1924
1925        assert_eq!(q.superseded_by("aaaa"), Some("bbbb".to_owned()));
1926        assert_eq!(q.superseded_by("bbbb"), Some("cccc".to_owned()));
1927        assert_eq!(
1928            q.superseded_by("cccc"),
1929            None,
1930            "the latest attempt replaces nothing"
1931        );
1932        assert_eq!(
1933            q.superseded_by("never-heard-of-it"),
1934            None,
1935            "a run belonging to no task on this queue is not superseded"
1936        );
1937
1938        let mut by = HashMap::new();
1939        by.insert("aaaa".to_owned(), "bbbb".to_owned());
1940        by.insert("bbbb".to_owned(), "cccc".to_owned());
1941        assert_eq!(
1942            q.superseded(),
1943            by,
1944            "the whole-map and single-run forms must agree"
1945        );
1946    }
1947
1948    #[test]
1949    fn superseded_attempts_is_empty_until_the_task_is_done() {
1950        let mut t = task("retried");
1951        t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1952        t.status = TaskStatus::Failed;
1953        assert_eq!(
1954            t.superseded_attempts(true),
1955            &[] as &[String],
1956            "a task still retrying has no attempt yet that a later one made moot"
1957        );
1958
1959        t.status = TaskStatus::Running;
1960        assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1961    }
1962
1963    #[test]
1964    fn superseded_attempts_names_every_run_before_the_one_that_succeeded() {
1965        let mut t = task("retried");
1966        t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
1967        t.status = TaskStatus::Done;
1968        assert_eq!(
1969            t.superseded_attempts(true),
1970            &["aaaa".to_owned(), "bbbb".to_owned()],
1971            "cccc is the attempt whose success made the task done, and stays out"
1972        );
1973    }
1974
1975    #[test]
1976    fn superseded_attempts_is_empty_for_a_done_task_with_only_one_attempt() {
1977        let mut t = task("first try landed");
1978        t.runs = vec!["aaaa".to_owned()];
1979        t.status = TaskStatus::Done;
1980        assert_eq!(t.superseded_attempts(true), &[] as &[String]);
1981    }
1982
1983    #[test]
1984    fn superseded_attempts_is_empty_when_the_last_run_never_actually_succeeded() {
1985        // `succeed` is reachable directly - `magi task done`, its web
1986        // equivalent, and the conductor's `Recovery::Done` - on a task in
1987        // any status, including one whose last recorded attempt is itself
1988        // `Blocked`/`Failed`/anything but `Merged`/`Ready`. `runs.last()`
1989        // alone cannot tell that apart from the loop's own settle path, so
1990        // the caller's own read of that run's status is what decides this.
1991        let mut t = task("closed by hand after a manual merge");
1992        t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned()];
1993        t.status = TaskStatus::Done;
1994        assert_eq!(
1995            t.superseded_attempts(false),
1996            &[] as &[String],
1997            "nothing here is provably why the task is done, so nothing is superseded"
1998        );
1999    }
2000
2001    #[test]
2002    fn latest_attempt_names_the_chain_s_current_head_not_just_the_next_one() {
2003        let (_dir, q) = queue();
2004        let mut t = task("retried twice");
2005        t.runs = vec!["aaaa".to_owned(), "bbbb".to_owned(), "cccc".to_owned()];
2006        q.put(&mut t).unwrap();
2007
2008        assert_eq!(
2009            q.latest_attempt("aaaa"),
2010            Some("cccc".to_owned()),
2011            "an old attempt points straight at the chain's current head, not the \
2012             next attempt in the middle of it"
2013        );
2014        assert_eq!(q.latest_attempt("bbbb"), Some("cccc".to_owned()));
2015        assert_eq!(
2016            q.latest_attempt("cccc"),
2017            None,
2018            "the latest attempt is not superseded by anything"
2019        );
2020        assert_eq!(
2021            q.latest_attempt("never-heard-of-it"),
2022            None,
2023            "a run belonging to no task on this queue is not superseded"
2024        );
2025    }
2026
2027    #[test]
2028    fn a_markdown_heading_is_the_title_not_decoration() {
2029        // A task file's heading is the summary its author already wrote, so it
2030        // beats the prose underneath. Getting this backwards was visible in the
2031        // first smoke test: a task titled "# Rework the config loader" listed
2032        // as "It re-reads the file on every lookup".
2033        assert_eq!(
2034            title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
2035            "Rework the config loader"
2036        );
2037        assert_eq!(title_from("- fix the thing", 40), "fix the thing");
2038        assert_eq!(title_from("> quoted task", 40), "quoted task");
2039        // Nothing usable at all still has to produce something printable.
2040        assert_eq!(title_from("   \n\n", 40), "(empty task)");
2041        assert_eq!(title_from("###\n", 40), "(empty task)");
2042    }
2043
2044    #[test]
2045    fn a_long_title_is_elided_by_characters_not_bytes() {
2046        // Byte truncation would split a multi-byte character and panic.
2047        let long = "課題".repeat(30);
2048        let title = title_from(&long, 10);
2049        assert_eq!(title.chars().count(), 10);
2050        assert!(title.ends_with('…'));
2051    }
2052
2053    #[test]
2054    fn priority_wins_and_ties_break_oldest_first() {
2055        let (_dir, q) = queue();
2056        let mut a = task("first");
2057        let mut b = task("second");
2058        let mut c = task("urgent");
2059        // Ids carry a timestamp, so force a known order.
2060        a.id = "20260101-000001-aaaa".to_owned();
2061        b.id = "20260101-000002-bbbb".to_owned();
2062        c.id = "20260101-000003-cccc".to_owned();
2063        c.priority = 5;
2064        for t in [&mut a, &mut b, &mut c] {
2065            q.put(t).unwrap();
2066        }
2067
2068        // Priority first...
2069        assert_eq!(q.next_runnable().unwrap().id, c.id);
2070        c.hold_machine(None);
2071        q.put(&mut c).unwrap();
2072        // ...then oldest, so a burst of new work cannot starve older work.
2073        assert_eq!(q.next_runnable().unwrap().id, a.id);
2074        assert_eq!(q.list().len(), 3, "b is still waiting its turn");
2075    }
2076
2077    #[test]
2078    fn a_blocked_task_never_starves_another_runnable_one() {
2079        let (_dir, q) = queue();
2080        let mut blocked = task("blocked");
2081        blocked.block(vec!["something".to_owned()], None);
2082        q.put(&mut blocked).unwrap();
2083
2084        let mut runnable = task("free to go");
2085        q.put(&mut runnable).unwrap();
2086
2087        let next = q.next_runnable().expect("a runnable task is still offered");
2088        assert_eq!(next.id, runnable.id);
2089    }
2090
2091    #[test]
2092    fn a_held_task_is_never_offered_to_the_loop() {
2093        let (_dir, q) = queue();
2094        let mut t = task("held");
2095        q.put(&mut t).unwrap();
2096        assert!(q.next_runnable().is_some());
2097
2098        t.hold_machine(None);
2099        q.put(&mut t).unwrap();
2100        assert!(
2101            q.next_runnable().is_none(),
2102            "a held task must wait for a human"
2103        );
2104
2105        // A failed task, by contrast, is exactly what the loop should retry.
2106        t.status = TaskStatus::Failed;
2107        q.put(&mut t).unwrap();
2108        assert!(q.next_runnable().is_some());
2109    }
2110
2111    #[test]
2112    fn attempts_are_capped_and_then_the_task_is_held() {
2113        let mut t = task("doomed");
2114
2115        t.start("run-1".to_owned());
2116        t.fail("gate red", 2);
2117        assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
2118
2119        t.start("run-2".to_owned());
2120        t.fail("gate red", 2);
2121        assert_eq!(
2122            t.status,
2123            TaskStatus::Held,
2124            "out of attempts: stop spending money on it"
2125        );
2126        assert_eq!(t.runs, ["run-1", "run-2"]);
2127        assert_eq!(t.last_error.as_deref(), Some("gate red"));
2128        assert_eq!(
2129            t.hold_reason.as_deref(),
2130            Some("gate red"),
2131            "the hold must say why, not leave hold_reason null next to a \
2132             populated last_error"
2133        );
2134    }
2135
2136    #[test]
2137    fn handing_off_a_task_records_a_hold_reason_too() {
2138        let mut t = task("left a pull request");
2139        t.start("run-1".to_owned());
2140        t.handed_off("run ended with a pull request open [run run-1]");
2141        assert_eq!(t.status, TaskStatus::Held);
2142        assert_eq!(t.hold_source, Some(HoldSource::Machine));
2143        assert_eq!(
2144            t.hold_reason.as_deref(),
2145            Some("run ended with a pull request open [run run-1]")
2146        );
2147        assert_eq!(t.hold_reason, t.last_error);
2148    }
2149
2150    #[test]
2151    fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
2152        let mut t = task("stalled by quota");
2153
2154        t.start("run-1".to_owned());
2155        assert_eq!(t.attempts, 1);
2156        t.stall("judge-1, judge-2 out of quota");
2157        assert_eq!(
2158            t.attempts, 0,
2159            "a closed quota window must not spend the task's retry budget"
2160        );
2161        assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
2162        assert_eq!(
2163            t.last_error.as_deref(),
2164            Some("judge-1, judge-2 out of quota")
2165        );
2166
2167        // A task can therefore stall all night and still get its real attempts
2168        // once the quota resets - which is the whole point.
2169        for _ in 0..20 {
2170            t.start("run-n".to_owned());
2171            t.stall("still out of quota");
2172        }
2173        t.start("run-real".to_owned());
2174        t.fail("gate red", 2);
2175        assert_eq!(
2176            t.status,
2177            TaskStatus::Failed,
2178            "the first attempt that was really judged is attempt one"
2179        );
2180    }
2181
2182    #[test]
2183    fn releasing_a_held_task_gives_it_a_real_second_chance() {
2184        let mut t = task("retry me");
2185        t.start("run-1".to_owned());
2186        t.fail("gate red", 1);
2187        assert_eq!(t.status, TaskStatus::Held);
2188
2189        t.release();
2190        assert_eq!(t.status, TaskStatus::Queued);
2191        // Without resetting attempts the next failure would re-hold at once,
2192        // and a release would be a no-op the operator cannot see.
2193        assert_eq!(t.attempts, 0);
2194        assert!(t.last_error.is_none());
2195        assert_eq!(
2196            t.runs.len(),
2197            1,
2198            "history is kept: attempts reset, evidence does not"
2199        );
2200    }
2201
2202    #[test]
2203    fn a_hold_reason_survives_and_a_release_clears_it() {
2204        let mut t = task("waiting on something else");
2205        t.hold_manual(Some(
2206            "waiting for 20260101-000000-aaaa to land first".to_owned(),
2207        ));
2208        assert_eq!(t.status, TaskStatus::Held);
2209        assert_eq!(
2210            t.hold_reason.as_deref(),
2211            Some("waiting for 20260101-000000-aaaa to land first")
2212        );
2213
2214        // Holding again with no reason must not erase the one already there.
2215        t.hold_manual(None);
2216        assert_eq!(
2217            t.hold_reason.as_deref(),
2218            Some("waiting for 20260101-000000-aaaa to land first"),
2219            "a bare re-hold keeps whatever a human already wrote down"
2220        );
2221
2222        // A hold with no reason at all is still an ordinary, allowed hold.
2223        let mut plain = task("no reason given");
2224        plain.hold_manual(None);
2225        assert_eq!(plain.status, TaskStatus::Held);
2226        assert!(plain.hold_reason.is_none());
2227
2228        t.release();
2229        assert_eq!(t.status, TaskStatus::Queued);
2230        assert!(
2231            t.hold_reason.is_none(),
2232            "a stale reason must not greet the next person who holds this task"
2233        );
2234    }
2235
2236    #[test]
2237    fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
2238        // `done` can close a held task directly - neither `magi task done`
2239        // nor `POST /api/queue/{id}/done` requires a release first - so a
2240        // task held for "waiting on 3ed9" and then closed without ever being
2241        // released must not still read as waiting on it afterwards.
2242        let mut t = task("landed by hand while held");
2243        t.hold_manual(Some("waiting on 3ed9".to_owned()));
2244        assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
2245
2246        t.succeed();
2247        assert_eq!(t.status, TaskStatus::Done);
2248        assert!(
2249            t.hold_reason.is_none(),
2250            "a done task cannot still be waiting on something"
2251        );
2252    }
2253
2254    #[test]
2255    fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
2256        // The web UI's "Hold" and "Mark done" buttons are both reachable on a
2257        // `blocked` task, not just on `queued`/`held` ones - neither requires
2258        // a release first. A task moved off `Blocked` that way must not still
2259        // carry the dependency it was waiting on: a dependency graph built
2260        // from `blocked_by` would otherwise keep drawing an edge for a task
2261        // that is not blocked on anything any more.
2262        let mut held = task("held straight out of blocked");
2263        held.block(
2264            vec!["20260101-000000-dead".to_owned()],
2265            Some("waiting on the migration script".to_owned()),
2266        );
2267        assert_eq!(held.status, TaskStatus::Blocked);
2268
2269        held.hold_manual(None);
2270        assert_eq!(held.status, TaskStatus::Held);
2271        assert!(
2272            held.blocked_by.is_empty(),
2273            "hold overrides the wait, same as release"
2274        );
2275        assert!(held.block_reason.is_none());
2276
2277        let mut done = task("closed straight out of blocked");
2278        done.block(
2279            vec!["20260101-000000-dead".to_owned()],
2280            Some("waiting on the migration script".to_owned()),
2281        );
2282        done.succeed();
2283        assert_eq!(done.status, TaskStatus::Done);
2284        assert!(
2285            done.blocked_by.is_empty(),
2286            "a done task cannot still be waiting on a dependency"
2287        );
2288        assert!(done.block_reason.is_none());
2289    }
2290
2291    #[test]
2292    fn a_blocked_task_is_never_offered_to_the_loop() {
2293        let mut t = task("blocked");
2294        assert!(t.status.runnable());
2295        t.block(
2296            vec!["dep-id".to_owned()],
2297            Some("waits on dep-id".to_owned()),
2298        );
2299        assert_eq!(t.status, TaskStatus::Blocked);
2300        assert!(!t.status.runnable());
2301        assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
2302    }
2303
2304    #[test]
2305    fn unblocking_the_last_dependency_returns_the_task_to_queued() {
2306        let mut t = task("blocked on two");
2307        t.block(
2308            vec!["a".to_owned(), "b".to_owned()],
2309            Some("waits on a and b".to_owned()),
2310        );
2311
2312        t.unblock("a");
2313        assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
2314        assert_eq!(t.blocked_by, ["b"]);
2315
2316        t.unblock("b");
2317        assert_eq!(t.status, TaskStatus::Queued);
2318        assert!(t.blocked_by.is_empty());
2319        assert!(t.block_reason.is_none());
2320    }
2321
2322    #[test]
2323    fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
2324        let mut t = task("never blocked");
2325        t.unblock("whatever");
2326        assert_eq!(t.status, TaskStatus::Queued);
2327    }
2328
2329    #[test]
2330    fn a_held_task_blocked_on_a_question_returns_to_held_not_queued() {
2331        // The bug this guards: a task an operator (or `crate::triage`) has
2332        // deliberately held, once `crate::conduct` blocks it on a follow-up
2333        // question, must not silently re-enter the competition queue the
2334        // moment that question is answered - whatever the answer said.
2335        let mut t = task("held, then asked about");
2336        t.hold_machine(Some("out of attempts".to_owned()));
2337        assert_eq!(t.status, TaskStatus::Held);
2338
2339        t.block(vec!["q1".to_owned()], Some("what now?".to_owned()));
2340        assert_eq!(t.status, TaskStatus::Blocked);
2341
2342        t.record_answer("what now?".to_owned(), "leave it held".to_owned());
2343        t.unblock("q1");
2344        assert_eq!(t.status, TaskStatus::Held, "must restore, not requeue");
2345        assert_eq!(t.hold_reason.as_deref(), Some("out of attempts"));
2346        assert_eq!(t.hold_source, Some(HoldSource::Machine));
2347        assert!(t.blocked_from.is_none(), "consumed once restored");
2348    }
2349
2350    #[test]
2351    fn a_manually_held_task_blocked_on_a_question_returns_to_held() {
2352        let mut t = task("manually held, then asked about");
2353        t.hold_manual(Some("waiting on a dependency".to_owned()));
2354
2355        t.block(vec!["q1".to_owned()], None);
2356        t.unblock("q1");
2357
2358        assert_eq!(t.status, TaskStatus::Held);
2359        assert_eq!(t.hold_source, Some(HoldSource::Manual));
2360    }
2361
2362    #[test]
2363    fn re_blocking_an_already_blocked_task_keeps_the_original_blocked_from() {
2364        // A second `Task::block` call - `crate::conduct` adding a question on
2365        // top of an existing block - must not overwrite `blocked_from` with
2366        // `Blocked` itself, or the task would restore into itself.
2367        let mut t = task("held, blocked twice");
2368        t.hold_machine(None);
2369        t.block(vec!["q1".to_owned()], Some("first".to_owned()));
2370        t.block(
2371            vec!["q1".to_owned(), "q2".to_owned()],
2372            Some("second".to_owned()),
2373        );
2374
2375        t.unblock("q1");
2376        assert_eq!(t.status, TaskStatus::Blocked, "q2 still outstanding");
2377        t.unblock("q2");
2378        assert_eq!(t.status, TaskStatus::Held);
2379    }
2380
2381    #[test]
2382    fn unblocking_a_task_blocked_while_running_lands_on_queued_not_running() {
2383        // Whatever process was driving the run is gone by the time a
2384        // conductor's question about it gets answered - there is nothing left
2385        // to resume into.
2386        let mut t = task("blocked mid-run");
2387        t.start("run-1".to_owned());
2388        assert_eq!(t.status, TaskStatus::Running);
2389
2390        t.block(vec!["q1".to_owned()], None);
2391        t.unblock("q1");
2392        assert_eq!(t.status, TaskStatus::Queued);
2393    }
2394
2395    #[test]
2396    fn a_pre_schema_4_blocked_record_with_hold_evidence_restores_to_held() {
2397        // `blocked_from` is `None` for a record written before schema 4 (or,
2398        // equivalently, deserialized straight from an on-disk file that never
2399        // had the field). Held evidence surviving on the task - never cleared
2400        // by `block` - is the only way left to tell such a record apart from
2401        // one blocked straight out of `Queued`.
2402        let mut t = task("legacy record, held before it was blocked");
2403        t.hold_source = Some(HoldSource::Machine);
2404        t.hold_reason = Some("legacy hold reason".to_owned());
2405        t.status = TaskStatus::Blocked;
2406        t.blocked_by = vec!["q1".to_owned()];
2407        t.blocked_from = None;
2408
2409        t.unblock("q1");
2410        assert_eq!(t.status, TaskStatus::Held);
2411    }
2412
2413    #[test]
2414    fn a_pre_schema_4_blocked_record_with_no_hold_evidence_restores_to_queued() {
2415        let mut t = task("legacy record, ordinary dependency block");
2416        t.status = TaskStatus::Blocked;
2417        t.blocked_by = vec!["dep".to_owned()];
2418        t.blocked_from = None;
2419
2420        t.unblock("dep");
2421        assert_eq!(t.status, TaskStatus::Queued);
2422    }
2423
2424    #[test]
2425    fn answering_a_question_is_recorded_and_survives_a_release() {
2426        let mut t = task("asked something");
2427        t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
2428        t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
2429        t.unblock("q1");
2430        assert_eq!(t.status, TaskStatus::Queued);
2431        assert_eq!(t.answers.len(), 1);
2432        assert_eq!(t.answers[0].answer, "SQLite");
2433
2434        // A release resets attempts, not evidence - the same rule
2435        // `releasing_a_held_task_gives_it_a_real_second_chance` asserts for
2436        // `runs`.
2437        t.release();
2438        assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
2439    }
2440
2441    #[test]
2442    fn a_refused_handover_keeps_the_review_branch_across_release() {
2443        let mut t = task("refused takeover");
2444        t.start("run-1".to_owned());
2445        t.hold_for_handover(Some("magi/eba2/A".to_owned()), "checked out".to_owned());
2446        assert_eq!(t.status, TaskStatus::Held);
2447        assert_eq!(t.attempts, 1);
2448        t.release();
2449        assert_eq!(t.status, TaskStatus::Queued);
2450        assert_eq!(t.attempts, 0);
2451        assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2452
2453        // A fresh competition chosen afterwards is not a review.
2454        t.hold_for_handover(None, "again".to_owned());
2455        t.requeue();
2456        assert!(t.review_branch.is_none());
2457
2458        // A manual hold never keeps one.
2459        let mut m = task("manual");
2460        m.review_branch = Some("magi/x/A".to_owned());
2461        m.hold_manual(None);
2462        m.release();
2463        assert!(m.review_branch.is_none());
2464    }
2465
2466    #[test]
2467    fn requesting_review_requeues_the_task_and_remembers_the_branch() {
2468        let mut t = task("blocked run with a surviving branch");
2469        t.start("run-1".to_owned());
2470        t.fail("blocked with major findings", 5);
2471        assert_eq!(t.status, TaskStatus::Failed);
2472
2473        t.request_review("magi/eba2/A".to_owned());
2474        assert_eq!(t.status, TaskStatus::Queued);
2475        assert_eq!(t.attempts, 0);
2476        assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
2477
2478        // An ordinary release (a human overriding the choice) drops it again.
2479        t.release();
2480        assert!(t.review_branch.is_none());
2481    }
2482
2483    #[test]
2484    fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
2485        let mut t = task("retry");
2486        t.start("run-1".to_owned());
2487        t.requeue();
2488        assert!(t.fresh_start);
2489
2490        t.release();
2491        assert!(!t.fresh_start);
2492    }
2493
2494    #[test]
2495    fn priority_can_be_changed_while_queued_but_not_while_running() {
2496        let mut t = task("reprioritise me");
2497        t.set_priority(5).unwrap();
2498        assert_eq!(t.priority, 5);
2499
2500        t.start("run-1".to_owned());
2501        let err = t.set_priority(9).unwrap_err().to_string();
2502        assert!(err.contains("running"), "{err}");
2503        assert_eq!(t.priority, 5, "the rejected write must not partially apply");
2504    }
2505
2506    #[test]
2507    fn interrupt_can_be_marked_while_queued_but_not_while_running() {
2508        let mut t = task("interrupt me");
2509        assert!(!t.interrupt, "off unless asked, same as any other task");
2510
2511        t.set_interrupt(true).unwrap();
2512        assert!(t.interrupt);
2513
2514        t.start("run-1".to_owned());
2515        assert!(
2516            !t.interrupt,
2517            "the mark is one-shot: dispatching the task fulfils it, \
2518             whatever the run that follows ends up doing"
2519        );
2520        let err = t.set_interrupt(true).unwrap_err().to_string();
2521        assert!(err.contains("running"), "{err}");
2522        // Clearing is always allowed, even on a running task - there is
2523        // nothing left for it to interrupt once it has been claimed.
2524        t.set_interrupt(false).unwrap();
2525        assert!(!t.interrupt);
2526    }
2527
2528    /// R2-1-1: a task whose run fails and requeues must not go on
2529    /// re-triggering `crate::daemon`'s interrupt scheduler on every later
2530    /// boundary, attempt after attempt, until it exhausts its budget.
2531    #[test]
2532    fn a_failed_run_does_not_leave_the_task_still_marked_to_interrupt() {
2533        let mut t = task("interrupt me");
2534        t.set_interrupt(true).unwrap();
2535        t.start("run-1".to_owned());
2536        t.fail("mock failure", 5);
2537        assert_eq!(t.status, TaskStatus::Failed);
2538        assert!(
2539            !t.interrupt,
2540            "one attempt already spent the mark; a retry is an ordinary \
2541             requeue, not a fresh interrupt request"
2542        );
2543    }
2544
2545    #[test]
2546    fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
2547        let (_dir, q) = queue();
2548        let mut a = task("first filed");
2549        let mut b = task("second filed");
2550        a.id = "20260101-000001-aaaa".to_owned();
2551        b.id = "20260101-000002-bbbb".to_owned();
2552        q.put(&mut a).unwrap();
2553        q.put(&mut b).unwrap();
2554
2555        assert_eq!(
2556            q.next_runnable().unwrap().id,
2557            a.id,
2558            "with equal priority the older task goes first, so a burst of \
2559             new work cannot starve it"
2560        );
2561        assert_eq!(
2562            q.list()[0].id,
2563            b.id,
2564            "but the list an operator reads is newest first, the same as \
2565             before priority existed - a's turn to run does not make it the \
2566             newest task"
2567        );
2568
2569        let mut a = q.get(&a.id).unwrap();
2570        a.set_priority(10).unwrap();
2571        q.put(&mut a).unwrap();
2572
2573        assert_eq!(
2574            q.next_runnable().unwrap().id,
2575            a.id,
2576            "a raised priority must be reflected the moment it is saved"
2577        );
2578        // `magi task list` and `GET /api/queue` both print `Queue::list()`
2579        // directly, so the raised task has to lead there too - not only in
2580        // what the loop would claim next.
2581        assert_eq!(
2582            q.list()[0].id,
2583            a.id,
2584            "the raised task must sort first in the list an operator reads, \
2585             not only in next_runnable's own ordering"
2586        );
2587    }
2588
2589    #[test]
2590    fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
2591        let mut t = Task::new(
2592            "old title".to_owned(),
2593            "old instruction".to_owned(),
2594            PathBuf::from("/repo"),
2595            Source::Agent {
2596                run: "20260101-000000-beef".to_owned(),
2597                node: "implement".to_owned(),
2598            },
2599        );
2600        let id = t.id.clone();
2601        let created_at = t.created_at;
2602        t.runs.push("20260101-000000-beef".to_owned());
2603
2604        t.edit("new title".to_owned(), "new instruction".to_owned())
2605            .unwrap();
2606
2607        assert_eq!(t.title, "new title");
2608        assert_eq!(t.instruction, "new instruction");
2609        assert_eq!(t.id, id, "editing must not mint a new id");
2610        assert_eq!(t.created_at, created_at);
2611        assert_eq!(
2612            t.source,
2613            Source::Agent {
2614                run: "20260101-000000-beef".to_owned(),
2615                node: "implement".to_owned(),
2616            },
2617            "editing must not turn agent attribution into human"
2618        );
2619        assert_eq!(t.runs, ["20260101-000000-beef"]);
2620    }
2621
2622    #[test]
2623    fn editing_is_refused_once_a_task_is_running_or_finished() {
2624        let mut running = task("in flight");
2625        running.start("run-1".to_owned());
2626        let err = running
2627            .edit("x".to_owned(), "y".to_owned())
2628            .unwrap_err()
2629            .to_string();
2630        assert!(err.contains("running"), "{err}");
2631
2632        let mut done = task("finished");
2633        done.succeed();
2634        let err = done
2635            .edit("x".to_owned(), "y".to_owned())
2636            .unwrap_err()
2637            .to_string();
2638        assert!(err.contains("done"), "{err}");
2639
2640        // Both queued and held are the point of the feature and must work.
2641        let mut queued = task("waiting");
2642        queued.edit("x".to_owned(), "y".to_owned()).unwrap();
2643        let mut held = task("parked");
2644        held.hold_machine(None);
2645        held.edit("x".to_owned(), "y".to_owned()).unwrap();
2646    }
2647
2648    #[test]
2649    fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
2650        let (_dir, q) = queue();
2651        let path = q.path_of("20260101-000000-aaaa");
2652        std::fs::create_dir_all(q.root()).unwrap();
2653        std::fs::write(
2654            &path,
2655            serde_json::json!({
2656                "schema": SCHEMA,
2657                "id": "20260101-000000-aaaa",
2658                "title": "from before hold reasons existed",
2659                "instruction": "from before hold reasons existed",
2660                "repo": ".",
2661                "source": { "kind": "human" },
2662                "status": "held",
2663                "created_at": Timestamp::now().to_string(),
2664                "updated_at": Timestamp::now().to_string(),
2665            })
2666            .to_string(),
2667        )
2668        .unwrap();
2669
2670        let task = q.get("20260101-000000-aaaa").expect("must still read");
2671        assert!(task.hold_reason.is_none());
2672        assert!(task.operator_held());
2673    }
2674
2675    #[test]
2676    fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
2677        let (_dir, q) = queue();
2678        let path = q.path_of("20260101-000000-bbbb");
2679        std::fs::create_dir_all(q.root()).unwrap();
2680        std::fs::write(
2681            &path,
2682            serde_json::json!({
2683                "schema": 2,
2684                "id": "20260101-000000-bbbb",
2685                "title": "old manual recovery",
2686                "instruction": "old manual recovery",
2687                "repo": ".",
2688                "source": { "kind": "human" },
2689                "status": "held",
2690                "hold_reason": "active manual recovery run20260912-224242-daf5",
2691                "created_at": Timestamp::now().to_string(),
2692                "updated_at": Timestamp::now().to_string(),
2693            })
2694            .to_string(),
2695        )
2696        .unwrap();
2697
2698        let task = q.get("20260101-000000-bbbb").expect("must still read");
2699        assert_eq!(task.hold_source, None);
2700        assert!(task.operator_held());
2701    }
2702
2703    #[test]
2704    fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
2705        let (_dir, q) = queue();
2706        let path = q.path_of("20260101-000000-aaaa");
2707        std::fs::create_dir_all(q.root()).unwrap();
2708        std::fs::write(
2709            &path,
2710            serde_json::json!({
2711                "schema": SCHEMA,
2712                "id": "20260101-000000-aaaa",
2713                "title": "from before diagnostics existed",
2714                "instruction": "from before diagnostics existed",
2715                "repo": ".",
2716                "source": { "kind": "human" },
2717                "status": "held",
2718                "created_at": Timestamp::now().to_string(),
2719                "updated_at": Timestamp::now().to_string(),
2720            })
2721            .to_string(),
2722        )
2723        .unwrap();
2724
2725        let task = q.get("20260101-000000-aaaa").expect("must still read");
2726        assert!(task.diagnostic.is_none());
2727    }
2728
2729    #[test]
2730    fn a_schema_1_task_with_no_blocking_fields_still_reads() {
2731        // Written by a build that predates `blocked_by`, `block_reason`,
2732        // `answers` and `review_branch` entirely - literal `"schema": 1`,
2733        // not `SCHEMA`, since the whole point is a build older than this one.
2734        let (_dir, q) = queue();
2735        let path = q.path_of("20260101-000000-aaaa");
2736        std::fs::create_dir_all(q.root()).unwrap();
2737        std::fs::write(
2738            &path,
2739            serde_json::json!({
2740                "schema": 1,
2741                "id": "20260101-000000-aaaa",
2742                "title": "from before blocking existed",
2743                "instruction": "from before blocking existed",
2744                "repo": ".",
2745                "source": { "kind": "human" },
2746                "status": "queued",
2747                "created_at": Timestamp::now().to_string(),
2748                "updated_at": Timestamp::now().to_string(),
2749            })
2750            .to_string(),
2751        )
2752        .unwrap();
2753
2754        let task = q.get("20260101-000000-aaaa").expect("must still read");
2755        assert!(task.blocked_by.is_empty());
2756        assert!(task.block_reason.is_none());
2757        assert!(task.answers.is_empty());
2758        assert!(task.review_branch.is_none());
2759    }
2760
2761    #[test]
2762    fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
2763        // A diagnostic belongs to the run that produced it. Left in place
2764        // across a release, an unrelated later failure - a config error, say -
2765        // would go on showing evidence for a problem that is no longer why the
2766        // task is stuck.
2767        let mut held = task("diagnosed");
2768        held.start("run-1".to_owned());
2769        held.fail("gate red", 1);
2770        held.diagnostic = Some("cargo test failed: ...".to_owned());
2771        assert_eq!(held.status, TaskStatus::Held);
2772
2773        held.release();
2774        assert!(held.diagnostic.is_none());
2775
2776        held.diagnostic = Some("cargo test failed: ...".to_owned());
2777        held.succeed();
2778        assert!(held.diagnostic.is_none());
2779    }
2780
2781    #[test]
2782    fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
2783        let mut t = task("retried");
2784        t.start("run-1".to_owned());
2785        t.diagnostic = Some("stale evidence from a previous hold".to_owned());
2786        t.fail("unrelated config error", 5);
2787        assert_eq!(t.status, TaskStatus::Failed);
2788        assert!(
2789            t.diagnostic.is_none(),
2790            "fail() must not let an old diagnostic outlive the run that produced it"
2791        );
2792    }
2793
2794    #[test]
2795    fn a_claim_is_exclusive_and_releases_on_drop() {
2796        let (_dir, q) = queue();
2797        let mut t = task("contended");
2798        q.put(&mut t).unwrap();
2799
2800        let held = q.claim(&t.id).unwrap();
2801        assert!(
2802            q.claim(&t.id).is_err(),
2803            "two daemons must not drive one task into two runs"
2804        );
2805        drop(held);
2806        assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
2807    }
2808
2809    #[test]
2810    fn a_round_trip_survives_disk() {
2811        let (_dir, q) = queue();
2812        let mut t = Task::new(
2813            "titled".to_owned(),
2814            "body".to_owned(),
2815            PathBuf::from("/repo"),
2816            Source::Agent {
2817                run: "20260101-000000-beef".to_owned(),
2818                node: "implement".to_owned(),
2819            },
2820        );
2821        t.priority = 3;
2822        q.put(&mut t).unwrap();
2823
2824        let back = q.get(&t.id).unwrap();
2825        assert_eq!(back.id, t.id);
2826        assert_eq!(back.priority, 3);
2827        assert_eq!(back.source.label(), "implement@beef");
2828        // A prefix is enough, the way run ids work everywhere else.
2829        assert_eq!(q.get(t.short()).unwrap().id, t.id);
2830    }
2831
2832    #[test]
2833    fn an_unreadable_task_does_not_take_the_queue_down() {
2834        let (_dir, q) = queue();
2835        let mut t = task("fine");
2836        q.put(&mut t).unwrap();
2837        std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
2838
2839        let listed = q.list();
2840        assert_eq!(listed.len(), 1, "the readable task still lists");
2841        assert_eq!(listed[0].id, t.id);
2842    }
2843
2844    #[test]
2845    fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
2846        let (_dir, q) = queue();
2847        let path = q.path_of("20260101-000000-aaaa");
2848        std::fs::create_dir_all(q.root()).unwrap();
2849        std::fs::write(
2850            &path,
2851            serde_json::json!({
2852                "schema": SCHEMA,
2853                "id": "20260101-000000-aaaa",
2854                "title": "from before solo existed",
2855                "instruction": "from before solo existed",
2856                "repo": ".",
2857                "source": { "kind": "human" },
2858                "status": "queued",
2859                "created_at": Timestamp::now().to_string(),
2860                "updated_at": Timestamp::now().to_string(),
2861            })
2862            .to_string(),
2863        )
2864        .unwrap();
2865
2866        let task = q.get("20260101-000000-aaaa").expect("must still read");
2867        assert!(!task.solo, "a queue file with no `solo` field means false");
2868    }
2869
2870    #[test]
2871    fn a_task_recorded_without_an_urgent_field_still_reads_as_not_urgent() {
2872        let (_dir, q) = queue();
2873        let path = q.path_of("20260101-000000-bbbb");
2874        std::fs::create_dir_all(q.root()).unwrap();
2875        std::fs::write(
2876            &path,
2877            serde_json::json!({
2878                "schema": SCHEMA,
2879                "id": "20260101-000000-bbbb",
2880                "title": "from before urgent existed",
2881                "instruction": "from before urgent existed",
2882                "repo": ".",
2883                "source": { "kind": "human" },
2884                "status": "queued",
2885                "created_at": Timestamp::now().to_string(),
2886                "updated_at": Timestamp::now().to_string(),
2887            })
2888            .to_string(),
2889        )
2890        .unwrap();
2891
2892        let task = q.get("20260101-000000-bbbb").expect("must still read");
2893        assert!(
2894            !task.urgent,
2895            "a queue file with no `urgent` field means false, same as `solo`"
2896        );
2897    }
2898
2899    #[test]
2900    fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
2901        let (_dir, q) = queue();
2902        let mut t = task("from the future");
2903        q.put(&mut t).unwrap();
2904        let path = q.path_of(&t.id);
2905        let body = std::fs::read_to_string(&path)
2906            .unwrap()
2907            .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
2908        std::fs::write(&path, body).unwrap();
2909
2910        let err = q.get(&t.id).unwrap_err().to_string();
2911        assert!(err.contains("schema 99"), "{err}");
2912    }
2913
2914    #[test]
2915    fn revision_moves_when_the_queue_changes() {
2916        let (_dir, q) = queue();
2917        assert_eq!(q.revision(), 0, "an empty queue has no revision");
2918        let mut t = task("first");
2919        q.put(&mut t).unwrap();
2920        assert!(q.revision() > 0, "a written task moves the revision");
2921    }
2922
2923    #[test]
2924    fn revision_moves_when_deleting_an_older_task() {
2925        let (dir, q) = queue();
2926        let questions = Questions::at(dir.path().join("questions"));
2927        let mut t1 = task("older");
2928        q.put(&mut t1).unwrap();
2929        // Ensure mtime ticks forward.
2930        std::thread::sleep(std::time::Duration::from_millis(10));
2931        let mut t2 = task("newer");
2932        q.put(&mut t2).unwrap();
2933
2934        let rev_before = q.revision();
2935        q.remove(&t1.id, false, &questions).unwrap();
2936        let rev_after = q.revision();
2937
2938        assert_ne!(
2939            rev_before, rev_after,
2940            "deleting an older task must change the revision so other clients see the deletion"
2941        );
2942    }
2943
2944    #[test]
2945    fn removing_a_task_takes_it_out_of_the_listing() {
2946        let (dir, q) = queue();
2947        let questions = Questions::at(dir.path().join("questions"));
2948        let mut t = task("delete me");
2949        q.put(&mut t).unwrap();
2950        let removed = q.remove(t.short(), false, &questions).unwrap();
2951        assert_eq!(removed.id, t.id, "a prefix resolves before deleting");
2952        assert!(removed.quarantined.is_empty(), "nothing was blocked on it");
2953        assert!(q.list().is_empty());
2954        assert!(
2955            q.remove(&t.id, false, &questions).is_err(),
2956            "removing twice is an error"
2957        );
2958    }
2959
2960    #[test]
2961    fn removing_a_task_takes_its_stale_lock_with_it() {
2962        let (dir, q) = queue();
2963        let questions = Questions::at(dir.path().join("questions"));
2964        let mut t = task("interrupted");
2965        q.put(&mut t).unwrap();
2966
2967        // A daemon killed mid-run leaves this behind. Nothing holds it: the
2968        // process that would have dropped the guard is gone.
2969        let claim = q.claim(&t.id).unwrap();
2970        std::mem::forget(claim);
2971        assert!(
2972            q.claim(&t.id).is_err(),
2973            "the orphaned lock is what makes the task look claimed"
2974        );
2975
2976        // A live daemon on this task is refused, whatever the lock says.
2977        let err = q.remove(&t.id, true, &questions).unwrap_err().to_string();
2978        assert!(err.contains("live daemon"), "{err}");
2979        assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
2980
2981        // With no daemon behind it, the lock is stale and goes with the task.
2982        q.remove(&t.id, false, &questions).unwrap();
2983        assert!(q.list().is_empty());
2984        let mut again = task("interrupted");
2985        again.id = t.id.clone();
2986        q.put(&mut again).unwrap();
2987        assert!(
2988            q.claim(&t.id).is_ok(),
2989            "a task that comes back must be claimable, which a left-behind lock would prevent"
2990        );
2991    }
2992
2993    #[test]
2994    fn removing_a_task_quarantines_what_was_blocked_on_it() {
2995        let (dir, q) = queue();
2996        let questions = Questions::at(dir.path().join("questions"));
2997
2998        let mut dep = task("dependency");
2999        q.put(&mut dep).unwrap();
3000
3001        let mut still_valid = task("still valid");
3002        q.put(&mut still_valid).unwrap();
3003
3004        let mut blocked = task("waiting");
3005        blocked.block(
3006            vec![dep.id.clone(), still_valid.id.clone()],
3007            Some("waits on both".to_owned()),
3008        );
3009        q.put(&mut blocked).unwrap();
3010
3011        let removed = q.remove(&dep.id, false, &questions).unwrap();
3012        assert_eq!(removed.quarantined, [blocked.id.clone()]);
3013
3014        let after = q.get(&blocked.id).unwrap();
3015        assert_eq!(after.status, TaskStatus::Held);
3016        assert_eq!(after.hold_source, Some(HoldSource::Machine));
3017        assert!(after.blocked_by.is_empty());
3018        let reason = after.hold_reason.as_deref().unwrap_or_default();
3019        assert!(reason.contains(&dep.id), "{reason}");
3020        assert!(
3021            reason.contains(&still_valid.id),
3022            "the still-valid dependency must survive in the reason text: {reason}"
3023        );
3024    }
3025
3026    fn source_file(dir: &Path, name: &str, body: &str) -> PathBuf {
3027        let p = dir.join(name);
3028        std::fs::write(&p, body).unwrap();
3029        p
3030    }
3031
3032    #[test]
3033    fn an_attachment_copy_survives_deleting_its_source() {
3034        let (dir, q) = queue();
3035        let src = source_file(dir.path(), "shot.png", "pixels");
3036        let mut t = task("with a picture");
3037        let names = q.attach(&mut t, std::slice::from_ref(&src)).unwrap();
3038        q.put(&mut t).unwrap();
3039        std::fs::remove_file(&src).unwrap();
3040        assert_eq!(names, ["shot.png"]);
3041        let loaded = q.get(&t.id).unwrap();
3042        let paths = q.attachment_paths(&loaded);
3043        assert_eq!(paths.len(), 1);
3044        assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "pixels");
3045    }
3046
3047    #[test]
3048    fn attachment_names_that_could_traverse_or_are_odd_are_refused() {
3049        let (dir, q) = queue();
3050        let mut t = task("bad names");
3051        for name in ["a..b.png", ".hidden", "with space.png", "-x.png"] {
3052            let src = source_file(dir.path(), name, "x");
3053            assert!(
3054                q.attach(&mut t, &[src]).is_err(),
3055                "`{name}` must be refused"
3056            );
3057        }
3058        // A drive-qualified name cannot exist as a file on Windows (joining it
3059        // to a directory yields a drive-relative path to `foo.png` instead),
3060        // so the rule is asserted on the name itself, everywhere.
3061        assert!(!crate::ask::valid_asset_name("C:foo.png"));
3062        #[cfg(not(windows))]
3063        {
3064            let src = source_file(dir.path(), "C:foo.png", "x");
3065            assert!(q.attach(&mut t, &[src]).is_err());
3066        }
3067        let long = format!("{}.png", "a".repeat(70));
3068        let src = source_file(dir.path(), &long, "x");
3069        assert!(q.attach(&mut t, &[src]).is_err());
3070        assert!(t.attachments.is_empty());
3071        assert!(!q.attachments_dir(&t.id).exists());
3072    }
3073
3074    #[test]
3075    fn a_taken_attachment_name_is_numbered_not_overwritten() {
3076        let (dir, q) = queue();
3077        let a = source_file(dir.path(), "shot.png", "one");
3078        let sub = dir.path().join("other");
3079        std::fs::create_dir_all(&sub).unwrap();
3080        let b = source_file(&sub, "shot.png", "two");
3081        let mut t = task("collision");
3082        q.attach(&mut t, &[a]).unwrap();
3083        q.attach(&mut t, &[b]).unwrap();
3084        assert_eq!(t.attachments, ["shot.png", "shot-2.png"]);
3085        let paths = q.attachment_paths(&t);
3086        assert_eq!(std::fs::read_to_string(&paths[0]).unwrap(), "one");
3087        assert_eq!(std::fs::read_to_string(&paths[1]).unwrap(), "two");
3088    }
3089
3090    #[test]
3091    fn a_renumbered_name_stays_inside_the_length_bound() {
3092        let (dir, q) = queue();
3093        let name = format!("{}.png", "a".repeat(60));
3094        assert_eq!(name.len(), 64);
3095        let a = source_file(dir.path(), &name, "one");
3096        let sub = dir.path().join("other");
3097        std::fs::create_dir_all(&sub).unwrap();
3098        let b = source_file(&sub, &name, "two");
3099        let mut t = task("long");
3100        q.attach(&mut t, &[a, b]).unwrap();
3101        assert_eq!(t.attachments.len(), 2);
3102        assert!(
3103            t.attachments
3104                .iter()
3105                .all(|n| crate::ask::valid_asset_name(n))
3106        );
3107        assert!(t.attachments[1].ends_with("-2.png"));
3108    }
3109
3110    #[test]
3111    fn a_failed_attach_keeps_existing_attachments_and_leaves_no_partial_copy() {
3112        let (dir, q) = queue();
3113        let good = source_file(dir.path(), "good.png", "ok");
3114        let mut t = task("partial");
3115        q.attach(&mut t, &[good]).unwrap();
3116        let more = source_file(dir.path(), "more.png", "ok");
3117        let missing = dir.path().join("missing.png");
3118        assert!(q.attach(&mut t, &[more, missing]).is_err());
3119        assert_eq!(t.attachments, ["good.png"]);
3120        let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3121            .unwrap()
3122            .flatten()
3123            .collect();
3124        assert_eq!(on_disk.len(), 1);
3125    }
3126
3127    /// Make the next `put` of `t` fail: its temp file's path is a directory.
3128    fn block_put(q: &Queue, t: &Task) -> PathBuf {
3129        let tmp = q.path_of(&t.id).with_extension("json.tmp");
3130        std::fs::create_dir_all(&tmp).unwrap();
3131        tmp
3132    }
3133
3134    #[test]
3135    fn a_failed_put_leaves_no_new_attachment_directory() {
3136        let (dir, q) = queue();
3137        let mut t = task("fresh");
3138        let tmp = block_put(&q, &t);
3139        let src = source_file(dir.path(), "shot.png", "x");
3140        assert!(q.attach_and_put(&mut t, &[src]).is_err());
3141        assert!(t.attachments.is_empty());
3142        assert!(!q.attachments_dir(&t.id).exists());
3143        assert!(!q.path_of(&t.id).exists());
3144        std::fs::remove_dir(tmp).unwrap();
3145    }
3146
3147    #[test]
3148    fn a_failed_put_removes_only_the_copy_it_just_made() {
3149        let (dir, q) = queue();
3150        let mut t = task("edited");
3151        let first = source_file(dir.path(), "first.png", "1");
3152        q.attach_and_put(&mut t, &[first]).unwrap();
3153        block_put(&q, &t);
3154        let second = source_file(dir.path(), "second.png", "2");
3155        assert!(q.attach_and_put(&mut t, &[second]).is_err());
3156        assert_eq!(t.attachments, ["first.png"]);
3157        let on_disk: Vec<_> = std::fs::read_dir(q.attachments_dir(&t.id))
3158            .unwrap()
3159            .flatten()
3160            .map(|e| e.file_name().to_string_lossy().into_owned())
3161            .collect();
3162        assert_eq!(on_disk, ["first.png"]);
3163        assert_eq!(q.get(&t.id).unwrap().attachments, ["first.png"]);
3164    }
3165
3166    #[test]
3167    fn a_leftover_removing_directory_is_swept_by_the_next_removal() {
3168        let (dir, q) = queue();
3169        let questions = Questions::at(dir.path().join("questions"));
3170        let gone = task("gone");
3171        let mut other = task("other");
3172        let mut live = task("live");
3173        q.put(&mut other).unwrap();
3174        q.put(&mut live).unwrap();
3175        // `gone` has no record: its removal finished deleting the record but
3176        // not the set-aside attachments.
3177        let orphan = q.root.join(format!("{}.attachments.removing", gone.id));
3178        std::fs::create_dir_all(&orphan).unwrap();
3179        std::fs::write(orphan.join("shot.png"), "x").unwrap();
3180        // `live` still has a record, so its `.removing` belongs to a removal
3181        // in progress and must be left alone.
3182        let busy = q.root.join(format!("{}.attachments.removing", live.id));
3183        std::fs::create_dir_all(&busy).unwrap();
3184
3185        q.remove(&other.id, false, &questions).unwrap();
3186        assert!(!orphan.exists(), "an orphan is swept");
3187        assert!(busy.exists(), "a removal in progress is left alone");
3188    }
3189
3190    #[test]
3191    fn a_blocked_aside_rename_fails_the_removal_and_loses_nothing() {
3192        let (dir, q) = queue();
3193        let questions = Questions::at(dir.path().join("questions"));
3194        let mut t = task("stuck");
3195        let src = source_file(dir.path(), "shot.png", "x");
3196        q.attach_and_put(&mut t, &[src]).unwrap();
3197        // A non-empty directory already at the set-aside name makes the rename
3198        // fail; its record is present, so the sweep leaves it in place.
3199        let aside = q.root.join(format!("{}.attachments.removing", t.id));
3200        std::fs::create_dir_all(&aside).unwrap();
3201        std::fs::write(aside.join("old.png"), "o").unwrap();
3202        assert!(q.remove(&t.id, false, &questions).is_err());
3203        assert!(q.path_of(&t.id).exists());
3204        assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3205    }
3206
3207    #[test]
3208    fn a_failed_record_removal_puts_the_attachments_back() {
3209        let (dir, q) = queue();
3210        let mut t = task("rollback");
3211        let src = source_file(dir.path(), "shot.png", "x");
3212        q.attach_and_put(&mut t, &[src]).unwrap();
3213        let err = q
3214            .remove_record_with_attachments(&t.id, |_| {
3215                Err(std::io::Error::other("injected failure"))
3216            })
3217            .unwrap_err();
3218        assert!(format!("{err:#}").contains("injected failure"));
3219        assert!(q.path_of(&t.id).exists());
3220        assert!(q.attachments_dir(&t.id).join("shot.png").is_file());
3221        assert!(
3222            !q.root
3223                .join(format!("{}.attachments.removing", t.id))
3224                .exists()
3225        );
3226    }
3227
3228    #[test]
3229    fn editing_a_task_keeps_its_attachments() {
3230        let (dir, q) = queue();
3231        let src = source_file(dir.path(), "shot.png", "x");
3232        let mut t = task("editable");
3233        q.attach(&mut t, &[src]).unwrap();
3234        t.edit("new".to_owned(), "new text".to_owned()).unwrap();
3235        q.put(&mut t).unwrap();
3236        assert_eq!(q.get(&t.id).unwrap().attachments, ["shot.png"]);
3237    }
3238
3239    #[test]
3240    fn removing_a_task_deletes_its_attachments() {
3241        let (dir, q) = queue();
3242        let questions = Questions::at(dir.path().join("questions"));
3243        let src = source_file(dir.path(), "shot.png", "x");
3244        let mut t = task("doomed");
3245        q.attach(&mut t, &[src]).unwrap();
3246        q.put(&mut t).unwrap();
3247        assert!(q.attachments_dir(&t.id).is_dir());
3248        q.remove(&t.id, false, &questions).unwrap();
3249        assert!(!q.attachments_dir(&t.id).exists());
3250        assert!(q.list().is_empty());
3251    }
3252
3253    #[test]
3254    fn a_task_written_before_attachments_still_reads() {
3255        let (_dir, q) = queue();
3256        let mut t = task("old");
3257        q.put(&mut t).unwrap();
3258        let path = q.path_of(&t.id);
3259        let mut v: serde_json::Value =
3260            serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3261        v.as_object_mut().unwrap().remove("attachments");
3262        std::fs::write(&path, v.to_string()).unwrap();
3263        assert!(q.get(&t.id).unwrap().attachments.is_empty());
3264    }
3265
3266    #[test]
3267    fn attachment_paths_are_absolute_even_when_the_root_is_relative() {
3268        let q = Queue::at(PathBuf::from("relative-queue"));
3269        let mut t = task("rel");
3270        t.attachments.push("shot.png".to_owned());
3271        let paths = q.attachment_paths(&t);
3272        assert!(paths[0].is_absolute(), "{}", paths[0].display());
3273        assert!(paths[0].ends_with(format!("{}.attachments/shot.png", t.id)));
3274    }
3275
3276    #[test]
3277    fn link_run_adds_a_run_once_and_touches_nothing_else() {
3278        let dir = tempfile::tempdir().unwrap();
3279        let queue = Queue::at(dir.path().join("queue"));
3280        let mut t = Task::new(
3281            "t".to_owned(),
3282            "do it".to_owned(),
3283            PathBuf::from("."),
3284            Source::Human,
3285        );
3286        queue.put(&mut t).unwrap();
3287        let before = queue.get(&t.id).unwrap();
3288
3289        let linked = queue.link_run(&t.id[..4], "20260930-092817-ec34").unwrap();
3290        assert_eq!(linked.runs, vec!["20260930-092817-ec34".to_owned()]);
3291        assert_eq!(linked.status, before.status);
3292        assert_eq!(linked.attempts, before.attempts);
3293        assert_eq!(linked.interrupt, before.interrupt);
3294
3295        let again = queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3296        assert_eq!(again.runs.len(), 1, "linking twice must not duplicate");
3297        assert_eq!(queue.get(&t.id).unwrap().runs.len(), 1);
3298        assert!(queue.link_run("no-such-task", "r").is_err());
3299    }
3300
3301    #[test]
3302    fn put_keeps_a_run_linked_after_the_writer_took_its_snapshot() {
3303        let dir = tempfile::tempdir().unwrap();
3304        let queue = Queue::at(dir.path().join("queue"));
3305        let mut t = Task::new(
3306            "t".to_owned(),
3307            "do it".to_owned(),
3308            PathBuf::from("."),
3309            Source::Human,
3310        );
3311        queue.put(&mut t).unwrap();
3312        // The daemon's copy, taken before the link.
3313        let mut snapshot = queue.get(&t.id).unwrap();
3314        queue.link_run(&t.id, "20260930-092817-ec34").unwrap();
3315
3316        snapshot.start("20260930-000000-aaaa".to_owned());
3317        queue.put(&mut snapshot).unwrap();
3318
3319        let stored = queue.get(&t.id).unwrap();
3320        assert!(stored.runs.contains(&"20260930-092817-ec34".to_owned()));
3321        assert!(stored.runs.contains(&"20260930-000000-aaaa".to_owned()));
3322        assert_eq!(stored.attempts, 1);
3323    }
3324
3325    #[test]
3326    fn concurrent_links_and_daemon_saves_lose_nothing() {
3327        let dir = tempfile::tempdir().unwrap();
3328        let queue = Queue::at(dir.path().join("queue"));
3329        let mut t = Task::new(
3330            "t".to_owned(),
3331            "do it".to_owned(),
3332            PathBuf::from("."),
3333            Source::Human,
3334        );
3335        queue.put(&mut t).unwrap();
3336        let id = t.id.clone();
3337
3338        let linkers: Vec<_> = (0..4)
3339            .map(|n| {
3340                let (queue, id) = (queue.clone(), id.clone());
3341                std::thread::spawn(move || {
3342                    for k in 0..10 {
3343                        queue
3344                            .link_run(&id, &format!("20260930-00000{n}-l{k:03}"))
3345                            .unwrap();
3346                    }
3347                })
3348            })
3349            .collect();
3350        // The daemon's side: transitions saved from its own, ever staler,
3351        // copy of the task.
3352        let mut mine = queue.get(&id).unwrap();
3353        for k in 0..10 {
3354            mine.start(format!("20260930-000009-d{k:03}"));
3355            queue.put(&mut mine).unwrap();
3356        }
3357        for l in linkers {
3358            l.join().unwrap();
3359        }
3360
3361        let stored = queue.get(&id).unwrap();
3362        assert_eq!(stored.runs.len(), 50, "{:?}", stored.runs);
3363        assert_eq!(
3364            stored.attempts, 10,
3365            "linking never rewinds the daemon's work"
3366        );
3367    }
3368}