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