Skip to main content

magi/
queue.rs

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