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