Skip to main content

magi/
queue.rs

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