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