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::path::{Path, PathBuf};
34
35use anyhow::{Context, Result, bail};
36use jiff::Timestamp;
37use serde::{Deserialize, Serialize};
38
39/// On-disk format for a queued task. Bumped when a field's meaning changes.
40///
41/// 3: added [`HoldSource`] so conductor recovery cannot release a hold an
42/// operator deliberately placed. Old records default to `None` and are
43/// protected as operator-held until an explicit release; the safe direction
44/// when their author was never recorded.
45///
46/// 2: added [`TaskStatus::Blocked`], [`Task::blocked_by`] and
47/// [`Task::block_reason`] (`crate::conduct`'s decisions) and
48/// [`Task::answers`] (operator answers carried forward to the next
49/// conductor prompt and the next run's instruction). All three are
50/// `#[serde(default)]`, so [`read_path`] accepts anything up to and
51/// including this schema rather than only an exact match — a task written
52/// by a build that only knew about schema 1 has nothing to say about
53/// blocking or answers, and defaulting those fields is exactly as good a
54/// reading as a value that build never had a chance to write.
55pub const SCHEMA: u32 = 3;
56
57/// Who placed the current hold.
58#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
59#[serde(rename_all = "lowercase")]
60pub enum HoldSource {
61    /// An operator used the CLI or web UI.
62    Manual,
63    /// The daemon or conductor placed the hold as part of its own recovery.
64    Machine,
65}
66
67impl HoldSource {
68    /// Short human-facing label for reports and the CLI.
69    pub fn label(self) -> &'static str {
70        match self {
71            Self::Manual => "manual",
72            Self::Machine => "machine",
73        }
74    }
75}
76
77/// Where a task came from. Recorded because "who asked for this" is the first
78/// question about an autonomous run, and the answer is not recoverable later.
79#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
80#[serde(tag = "kind", rename_all = "lowercase")]
81pub enum Source {
82    /// A person, at a terminal or through the web UI.
83    Human,
84    /// An agent inside a run, via `magi task add`. Both ids are recorded so a
85    /// task can be traced back to the exact seat that asked for it.
86    Agent {
87        /// Run the asking agent belonged to.
88        run: String,
89        /// Node it was working in, e.g. `implement` or `review`.
90        node: String,
91    },
92    /// A GitHub issue, imported by number.
93    Issue {
94        /// Issue number.
95        number: u64,
96        /// `owner/repo`, as `gh` reports it.
97        repo: String,
98    },
99}
100
101impl Source {
102    /// Short human-facing label, for lists and the web UI.
103    pub fn label(&self) -> String {
104        match self {
105            Self::Human => "human".to_owned(),
106            Self::Agent { run, node } => format!("{node}@{}", short(run)),
107            Self::Issue { number, .. } => format!("issue #{number}"),
108        }
109    }
110}
111
112/// Where a task is in its life.
113#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
114#[serde(rename_all = "lowercase")]
115pub enum TaskStatus {
116    /// Waiting to be claimed.
117    Queued,
118    /// Claimed by a daemon; a run is in flight.
119    Running,
120    /// A run finished and its gate passed.
121    Done,
122    /// A run finished without passing, and attempts remain.
123    Failed,
124    /// Out of attempts, or held by hand. The loop will not pick it up.
125    Held,
126    /// Waiting on another task or an unanswered question. See
127    /// [`Task::blocked_by`]. Set and cleared by `crate::conduct` and
128    /// `crate::daemon`'s deterministic resolver, never by hand.
129    Blocked,
130}
131
132impl TaskStatus {
133    /// Is this task eligible for a daemon to claim?
134    pub fn runnable(self) -> bool {
135        matches!(self, Self::Queued | Self::Failed)
136    }
137
138    /// Lowercase name, as it appears on disk and in the API.
139    pub fn as_str(self) -> &'static str {
140        match self {
141            Self::Queued => "queued",
142            Self::Running => "running",
143            Self::Done => "done",
144            Self::Failed => "failed",
145            Self::Held => "held",
146            Self::Blocked => "blocked",
147        }
148    }
149}
150
151/// One unit of work.
152#[derive(Debug, Clone, Serialize, Deserialize)]
153#[serde(deny_unknown_fields)]
154pub struct Task {
155    /// On-disk format version.
156    pub schema: u32,
157    /// Task id, e.g. `20260902-140501-a1b2`.
158    pub id: String,
159    /// One line, for lists and notifications.
160    pub title: String,
161    /// The task itself, handed to the graph verbatim.
162    pub instruction: String,
163    /// Repository to work in.
164    pub repo: PathBuf,
165    /// Who asked.
166    pub source: Source,
167    /// Higher runs first; ties break oldest-first so nothing starves.
168    #[serde(default)]
169    pub priority: i32,
170    /// Run this task alone: one implementer, no panel of judges to convince.
171    ///
172    /// `#[serde(default)]` so a queue file written before this field existed
173    /// still reads, as `false` - the ordinary multi-candidate competition,
174    /// unchanged. A task set to `solo` still runs the whole graph; only the
175    /// candidate count the daemon builds it with changes, and
176    /// [`crate::graph::Runner`] already collapses a single-candidate run to
177    /// implement → review → gate → merge on its own (see
178    /// [`crate::graph::Runner::review`]'s doc), so nothing about judging,
179    /// deliberation or voting had to change to support this.
180    #[serde(default)]
181    pub solo: bool,
182    /// Current state.
183    pub status: TaskStatus,
184    /// How many times this task has been claimed.
185    #[serde(default)]
186    pub attempts: usize,
187    /// Runs this task has produced, oldest first.
188    #[serde(default)]
189    pub runs: Vec<String>,
190    /// Why the last attempt did not land.
191    #[serde(default)]
192    pub last_error: Option<String>,
193    /// What a human hold is waiting on.
194    ///
195    /// `None` covers both the ordinary cases: a hold the loop makes itself
196    /// (out of attempts, or the disk gate closed) explains itself through
197    /// [`Task::last_error`] instead, and a human hold nobody bothered to
198    /// explain is still a valid hold. The queue has no way to express a
199    /// dependency between two tasks, so on the occasions a hold really is
200    /// "wait for that other task first", this is the only place that reason
201    /// survives - see [`Task::hold_manual`] and [`Task::release`].
202    ///
203    /// `#[serde(default)]` so a queue file written before this field existed
204    /// still reads, with no reason recorded rather than a parse error.
205    #[serde(default)]
206    pub hold_reason: Option<String>,
207    /// Who placed [`Task::hold_reason`].  `None` is a compatible old record;
208    /// see [`Task::operator_held`] for its deliberately conservative meaning.
209    #[serde(default)]
210    pub hold_source: Option<HoldSource>,
211    /// Diagnostic detail excerpted from the run that led to a hold - what a
212    /// human would have found opening `artifacts/` by hand, not the one-line
213    /// reason in [`Task::last_error`]. Set only when a run's own attempts are
214    /// exhausted and the task becomes [`TaskStatus::Held`]; `daemon` computes
215    /// it from the run's own record, since this module has no notion of a
216    /// run's internals. Bounded in length by the writer - see
217    /// `daemon::diagnostic` - so a verbose run cannot make this file grow
218    /// without limit.
219    ///
220    /// `#[serde(default)]` so a queue file written before this field existed
221    /// still reads, with no diagnostic recorded rather than a parse error.
222    #[serde(default)]
223    pub diagnostic: Option<String>,
224    /// What this task is waiting on: other task ids, unanswered
225    /// `crate::ask::Question` ids, or both. Non-empty exactly when
226    /// [`TaskStatus::Blocked`]; emptying it — see [`Task::unblock`] — is what
227    /// puts the task back at [`TaskStatus::Queued`].
228    ///
229    /// Set by `crate::conduct`'s decisions and cleared deterministically by
230    /// `crate::daemon` as each dependency resolves, never by a person. Never
231    /// `#[serde(default)]` is skipped: a queue file from before this field
232    /// existed has nothing to report here, and an empty list is exactly that.
233    #[serde(default)]
234    pub blocked_by: Vec<String>,
235    /// One line explaining the current [`Task::blocked_by`], written by
236    /// `crate::conduct`. Cleared whenever `blocked_by` empties.
237    #[serde(default)]
238    pub block_reason: Option<String>,
239    /// Questions `crate::conduct` asked about this task that the operator has
240    /// since answered, oldest first — what was asked, and what they said.
241    ///
242    /// A blocking question's id leaves [`Task::blocked_by`] the moment
243    /// [`crate::ask::QuestionStatus::Answered`] is observed, but the id alone
244    /// tells nobody what was decided. This is what carries the answer's
245    /// *content* forward: into the next conductor prompt for this task, and
246    /// into the instruction handed to the next run — see `crate::daemon`'s
247    /// deterministic blocker resolution. Kept for the task's whole life, the
248    /// same as [`Task::runs`]: a release resets attempts, not evidence.
249    #[serde(default)]
250    pub answers: Vec<AnsweredQuestion>,
251    /// Set by `crate::conduct` when it chooses `Review` recovery for a task
252    /// whose branch survived a blocked run: the branch to reopen with
253    /// `crate::graph::Runner::review` instead of competing from scratch.
254    ///
255    /// Requeues the task the same way [`Task::release`] does, so it is
256    /// picked up by the ordinary loop; `crate::daemon` reads this field once,
257    /// when it actually starts the run, and clears it either way — consumed
258    /// on success, dropped if the branch no longer exists by then. Never set
259    /// from the conductor's own words: `crate::daemon` derives the branch
260    /// name itself from the task's last run, so a hallucinated branch can
261    /// never reach here.
262    #[serde(default)]
263    pub review_branch: Option<String>,
264    /// A release deliberately starts a new competition instead of resuming
265    /// the prior run. History remains as evidence in `runs`.
266    #[serde(default)]
267    pub fresh_start: bool,
268    /// When the task was filed.
269    pub created_at: Timestamp,
270    /// Last change to this file.
271    pub updated_at: Timestamp,
272}
273
274/// One question `crate::conduct` asked about a task, and what the operator
275/// said back. See [`Task::answers`].
276#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
277pub struct AnsweredQuestion {
278    /// The question as asked, e.g. [`crate::ask::Question::summary`].
279    pub question: String,
280    /// What the operator answered.
281    pub answer: String,
282}
283
284impl Task {
285    /// File a new task. Persist it with [`Queue::put`].
286    pub fn new(title: String, instruction: String, repo: PathBuf, source: Source) -> Self {
287        let now = Timestamp::now();
288        Self {
289            schema: SCHEMA,
290            id: new_id(),
291            title,
292            instruction,
293            repo,
294            source,
295            priority: 0,
296            solo: false,
297            status: TaskStatus::Queued,
298            attempts: 0,
299            runs: Vec::new(),
300            last_error: None,
301            hold_reason: None,
302            hold_source: None,
303            diagnostic: None,
304            blocked_by: Vec::new(),
305            block_reason: None,
306            answers: Vec::new(),
307            review_branch: None,
308            fresh_start: false,
309            created_at: now,
310            updated_at: now,
311        }
312    }
313
314    /// Short form used in reports, matching a run's short id.
315    pub fn short(&self) -> &str {
316        short(&self.id)
317    }
318
319    /// Record that a run has started for this task.
320    pub fn start(&mut self, run: String) {
321        self.status = TaskStatus::Running;
322        self.attempts += 1;
323        self.runs.push(run);
324        self.last_error = None;
325        self.fresh_start = false;
326    }
327
328    /// Record a successful run.
329    ///
330    /// Both `magi task done` and `POST /api/queue/{id}/done` can close a held
331    /// task directly, with no release in between, so this clears
332    /// `hold_reason` the same way [`Task::release`] does. Otherwise a task
333    /// held for "waiting on 3ed9" and then closed as done without ever being
334    /// released would still read as waiting on something in `magi task show`
335    /// and on its card, after it no longer is.
336    pub fn succeed(&mut self) {
337        self.status = TaskStatus::Done;
338        self.last_error = None;
339        self.hold_reason = None;
340        self.hold_source = None;
341        self.diagnostic = None;
342    }
343
344    /// Record a failed attempt. Out of attempts means held for a human, rather
345    /// than retried until the money runs out.
346    ///
347    /// Clears [`Task::diagnostic`] unconditionally: it belongs to whatever run
348    /// produced it, and a caller that has one for *this* attempt sets it
349    /// itself right after calling this, once it knows the task actually ended
350    /// up [`TaskStatus::Held`] - see `daemon::diagnostic`. Without the clear, a
351    /// task released after a diagnosed hold and then failed again for an
352    /// unrelated, undiagnosed reason (a config error, say) would go on
353    /// showing the previous run's diagnostic as if it explained the new one.
354    pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
355        self.last_error = Some(why.into());
356        self.diagnostic = None;
357        self.status = if self.attempts >= max_attempts {
358            self.hold_source = Some(HoldSource::Machine);
359            TaskStatus::Held
360        } else {
361            TaskStatus::Failed
362        };
363    }
364
365    /// Record an attempt that failed for a reason the task is not responsible
366    /// for - the agent CLIs ran out of quota and the judging panel collapsed.
367    ///
368    /// This refunds the attempt on purpose. A quota window closing at 4am must
369    /// not spend the backlog's retry budget: the operator would come back to a
370    /// queue of held tasks that were never actually judged, and would have to
371    /// release every one by hand to find out which had a real problem. The task
372    /// goes back to `Failed`, which the loop retries, so a reset quota picks the
373    /// work up where it stopped.
374    pub fn stall(&mut self, why: impl Into<String>) {
375        self.last_error = Some(why.into());
376        self.diagnostic = None;
377        self.attempts = self.attempts.saturating_sub(1);
378        self.status = TaskStatus::Failed;
379    }
380
381    /// Whether this held task may only be released by an operator.
382    ///
383    /// Old files did not record a source. Preserve every such hold rather
384    /// than guessing that it was automatic and risking duplicate work. New
385    /// automatic holds record [`HoldSource::Machine`] and remain recoverable.
386    pub fn operator_held(&self) -> bool {
387        self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
388    }
389
390    /// Take this task out of the loop's reach by an operator action.
391    pub fn hold_manual(&mut self, reason: Option<String>) {
392        self.status = TaskStatus::Held;
393        if reason.is_some() {
394            self.hold_reason = reason;
395        }
396        self.hold_source = Some(HoldSource::Manual);
397    }
398
399    /// Take this task out of the loop's reach during automatic recovery.
400    pub fn hold_machine(&mut self, reason: Option<String>) {
401        self.status = TaskStatus::Held;
402        if reason.is_some() {
403            self.hold_reason = reason;
404        }
405        self.hold_source = Some(HoldSource::Machine);
406    }
407
408    /// Block this task on other task ids and/or open question ids, chosen by
409    /// `crate::conduct`. Pure: the caller still owns writing it back with
410    /// [`Queue::put`].
411    pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
412        self.status = TaskStatus::Blocked;
413        self.blocked_by = blocked_by;
414        self.block_reason = reason;
415    }
416
417    /// Remove one resolved dependency (a task id that became [`TaskStatus::Done`],
418    /// or a question id that became [`crate::ask::QuestionStatus::Answered`]).
419    /// Once nothing is left in [`Task::blocked_by`], the task returns to
420    /// [`TaskStatus::Queued`] on its own - deciding *why* a task was blocked
421    /// was `crate::conduct`'s job, but noticing a dependency resolved needs no
422    /// model at all.
423    ///
424    /// A no-op, on purpose, for a task that is not [`TaskStatus::Blocked`]:
425    /// `crate::daemon`'s deterministic resolver runs over every task on every
426    /// poll, and a task that moved on for some other reason must not be
427    /// dragged back to `Queued` by a stale id it still happens to carry.
428    pub fn unblock(&mut self, resolved_id: &str) {
429        if self.status != TaskStatus::Blocked {
430            return;
431        }
432        self.blocked_by.retain(|id| id != resolved_id);
433        if self.blocked_by.is_empty() {
434            self.status = TaskStatus::Queued;
435            self.block_reason = None;
436        }
437    }
438
439    /// Record that a question `crate::conduct` asked about this task has been
440    /// answered, so the answer's content — not just the fact that the
441    /// question is gone — reaches the next conductor prompt and the next
442    /// run's instruction. See [`Task::answers`].
443    pub fn record_answer(&mut self, question: String, answer: String) {
444        self.answers.push(AnsweredQuestion { question, answer });
445    }
446
447    /// Requeue this task to reopen its last run as a review-only pass against
448    /// `branch` (`crate::graph::Runner::review`) rather than competing from
449    /// scratch. See [`Task::review_branch`].
450    pub fn request_review(&mut self, branch: String) {
451        self.release();
452        self.review_branch = Some(branch);
453    }
454
455    /// Requeue after a conductor chose a new competition. Unlike an ordinary
456    /// operator release, this deliberately does not resume the old run.
457    pub fn requeue(&mut self) {
458        self.release();
459        self.fresh_start = true;
460    }
461
462    /// Change how urgently this task should run next.
463    ///
464    /// Refused once the task is `running`: priority only feeds the sort
465    /// [`Queue::next_runnable`] does over tasks waiting to be claimed, and a
466    /// running task has already left that pool. Accepting the write anyway
467    /// would look like it worked while changing nothing until - and unless -
468    /// this attempt fails and the task becomes runnable again, which is a
469    /// surprise the phone should not hand back as a success.
470    pub fn set_priority(&mut self, priority: i32) -> Result<()> {
471        if self.status == TaskStatus::Running {
472            bail!(
473                "task {} is running; its priority cannot be changed until \
474                 this attempt finishes",
475                self.short()
476            );
477        }
478        self.priority = priority;
479        Ok(())
480    }
481
482    /// Replace this task's title and instruction wholesale.
483    ///
484    /// Restricted to `queued` and `held`. A `running` task's instruction has
485    /// already been handed to the graph, so a run in flight and the file on
486    /// disk must not be allowed to disagree about what was asked; a `done` or
487    /// `failed` task is a record of what actually happened and editing it
488    /// after the fact would falsify that record. `id`, `created_at`,
489    /// `source`, and `runs` are left untouched on purpose - an edit stands in
490    /// for "delete and refile", and keeping the id, the timestamp, the
491    /// attribution, and the run history is the entire reason it exists
492    /// instead.
493    pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
494        if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
495            bail!(
496                "task {} is {}; only a queued or held task's instruction can \
497                 be edited",
498                self.short(),
499                self.status.as_str()
500            );
501        }
502        self.title = title;
503        self.instruction = instruction;
504        Ok(())
505    }
506
507    /// Record a run that produced a pull request without merging it.
508    ///
509    /// The task is held rather than retried, and it costs no further attempt
510    /// either way. The work the task asked for exists: it is sitting on a
511    /// branch, in a pull request, waiting for CI or for a person. Retrying
512    /// would spend the whole competition budget a second time and then race a
513    /// second branch against the pull request the first one opened - which is
514    /// exactly what happened to run 01c2, whose finished and green pull request
515    /// was re-competed from scratch four seconds after it opened.
516    ///
517    /// A pull request nobody merged is a request for a person, not a failure.
518    pub fn handed_off(&mut self, why: impl Into<String>) {
519        self.last_error = Some(why.into());
520        self.diagnostic = None;
521        self.status = TaskStatus::Held;
522        self.hold_source = Some(HoldSource::Machine);
523    }
524
525    /// Put a held or finished task back in line, with its attempt count reset
526    /// so a release is a real second chance rather than an instant re-hold.
527    /// The run history is kept: attempts reset, evidence does not.
528    pub fn release(&mut self) {
529        self.status = TaskStatus::Queued;
530        self.attempts = 0;
531        self.last_error = None;
532        // Otherwise the next person who holds this task reads a reason that
533        // belonged to whatever it was waiting on last time.
534        self.hold_reason = None;
535        self.hold_source = None;
536        self.diagnostic = None;
537        // A release also un-blocks: the dependency or question `blocked_by`
538        // named may still be unresolved, but a human (or `crate::conduct`)
539        // choosing to release the task overrides that wait outright, the same
540        // as it overrides an ordinary hold.
541        self.blocked_by.clear();
542        self.block_reason = None;
543        self.review_branch = None;
544        self.fresh_start = false;
545    }
546}
547
548/// A queue on disk.
549#[derive(Debug, Clone)]
550pub struct Queue {
551    root: PathBuf,
552}
553
554impl Queue {
555    /// The operator's queue, `<home>/queue`.
556    pub fn open() -> Self {
557        Self::at(crate::run::home().join("queue"))
558    }
559
560    /// A queue at an explicit root. Tests use this; so could an operator who
561    /// wants a queue per project.
562    pub fn at(root: PathBuf) -> Self {
563        Self { root }
564    }
565
566    /// Directory holding the task files.
567    pub fn root(&self) -> &Path {
568        &self.root
569    }
570
571    /// Path for one task id.
572    pub fn path_of(&self, id: &str) -> PathBuf {
573        self.root.join(format!("{id}.json"))
574    }
575
576    /// Write a task, atomically, so a daemon killed mid-write leaves the
577    /// previous state readable rather than a truncated file.
578    pub fn put(&self, task: &mut Task) -> Result<()> {
579        task.updated_at = Timestamp::now();
580        std::fs::create_dir_all(&self.root)
581            .with_context(|| format!("create {}", self.root.display()))?;
582        let body = serde_json::to_string_pretty(task).context("serialize task")?;
583        let path = self.path_of(&task.id);
584        let tmp = path.with_extension("json.tmp");
585        std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
586        std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
587        Ok(())
588    }
589
590    /// Load a task by id or unambiguous id prefix.
591    pub fn get(&self, id: &str) -> Result<Task> {
592        let resolved = self.resolve_id(id)?;
593        read_path(&self.path_of(&resolved))
594    }
595
596    /// Remove a task, and the claim lock that belongs to it.
597    ///
598    /// `in_flight` comes from the caller — a live daemon's heartbeat naming
599    /// this task — because the task's own `running` status cannot answer the
600    /// question. A daemon killed mid-competition leaves the status at
601    /// `running` and an orphaned `.lock` behind, and a guard that trusted
602    /// either would make the task undeletable for good: the phone showed
603    /// exactly that, refusing a task whose daemon had been gone for an hour.
604    ///
605    /// So the lock is removed with the task rather than respected. Any lock
606    /// still there once no live daemon claims the task is by definition stale,
607    /// and leaving it would make a deleted task look claimed to
608    /// [`Queue::claim`] and to whoever reads the directory.
609    pub fn remove(&self, id: &str, in_flight: bool) -> Result<String> {
610        let resolved = self.resolve_id(id)?;
611        if in_flight {
612            bail!("task {resolved} is being run by a live daemon right now");
613        }
614        let path = self.path_of(&resolved);
615        std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
616        let lock = self.lock_path(&resolved);
617        if let Err(e) = std::fs::remove_file(&lock) {
618            if e.kind() != std::io::ErrorKind::NotFound {
619                return Err(e).with_context(|| format!("remove {}", lock.display()));
620            }
621        }
622        Ok(resolved)
623    }
624
625    /// Path of the claim lock for a task. One definition, so `claim` and
626    /// `remove` cannot end up naming different files.
627    fn lock_path(&self, id: &str) -> PathBuf {
628        self.root.join(format!("{id}.lock"))
629    }
630
631    /// Every task on disk, highest priority first and newest first within a
632    /// priority. This is what `magi task list` and `GET /api/queue` print, so
633    /// a raised priority has to move a task here the moment it is saved, not
634    /// only in [`Queue::next_runnable`]'s own ordering - the operator reading
635    /// the backlog and the loop about to drain it must agree on what "first"
636    /// means. Every existing task defaults to priority 0, so this is a no-op
637    /// change from the old newest-first order for a queue nobody has
638    /// reprioritised.
639    ///
640    /// Unreadable files are skipped rather than fatal: one corrupt task must
641    /// not take the queue - or the web UI, or an unattended daemon - down
642    /// with it.
643    pub fn list(&self) -> Vec<Task> {
644        let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
645            .into_iter()
646            .flatten()
647            .flatten()
648            .map(|e| e.path())
649            .filter(|p| p.extension().is_some_and(|x| x == "json"))
650            .filter_map(|p| read_path(&p).ok())
651            .collect();
652        tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
653        tasks
654    }
655
656    /// The task a daemon should run next, or `None` when the queue is idle.
657    ///
658    /// Highest priority first, oldest first within a priority, so a burst of
659    /// agent-filed work cannot starve the task a human filed this morning.
660    pub fn next_runnable(&self) -> Option<Task> {
661        let mut runnable: Vec<Task> = self
662            .list()
663            .into_iter()
664            .filter(|t| t.status.runnable())
665            .collect();
666        runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
667        runnable.into_iter().next()
668    }
669
670    /// Take exclusive ownership of a task.
671    ///
672    /// The lock is a `create_new` file next to the task, which is atomic on
673    /// every platform magi targets. It exists so two daemons - or a daemon and
674    /// a human running `magi run` - cannot drive one task into two competing
675    /// runs. The returned guard releases on drop, including on panic.
676    pub fn claim(&self, id: &str) -> Result<Claim> {
677        std::fs::create_dir_all(&self.root)
678            .with_context(|| format!("create {}", self.root.display()))?;
679        let path = self.lock_path(id);
680        match std::fs::OpenOptions::new()
681            .write(true)
682            .create_new(true)
683            .open(&path)
684        {
685            Ok(mut f) => {
686                use std::io::Write as _;
687                // Best effort: the pid is for the human looking at a stale lock.
688                let _ = writeln!(f, "{}", std::process::id());
689                Ok(Claim { path })
690            }
691            Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
692                bail!("task {id} is already claimed ({} exists)", path.display())
693            }
694            Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
695        }
696    }
697
698    /// Expand an id prefix to exactly one task id.
699    pub fn resolve_id(&self, prefix: &str) -> Result<String> {
700        if self.path_of(prefix).is_file() {
701            return Ok(prefix.to_owned());
702        }
703        let hits: Vec<String> = self
704            .list()
705            .into_iter()
706            .map(|t| t.id)
707            .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
708            .collect();
709        match hits.len() {
710            1 => Ok(hits.into_iter().next().expect("exactly one hit")),
711            0 => bail!("no task matches `{prefix}`"),
712            _ => bail!(
713                "`{prefix}` matches {} tasks: {}",
714                hits.len(),
715                hits.join(", ")
716            ),
717        }
718    }
719
720    /// Change detection token for the queue.
721    ///
722    /// Combines file names and modification times of all task files in the
723    /// queue, so adding, modifying, or deleting any task — even an older one —
724    /// moves the revision and notifies connected clients via the change stream.
725    /// Returns 0 when the queue is completely empty.
726    pub fn revision(&self) -> u64 {
727        use std::hash::{Hash as _, Hasher as _};
728
729        let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
730            .into_iter()
731            .flatten()
732            .flatten()
733            .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
734            .filter_map(|e| {
735                let name = e.file_name().to_string_lossy().into_owned();
736                let mtime = e
737                    .metadata()
738                    .ok()?
739                    .modified()
740                    .ok()?
741                    .duration_since(std::time::UNIX_EPOCH)
742                    .ok()?
743                    .as_millis() as u64;
744                Some((name, mtime))
745            })
746            .collect();
747
748        if entries.is_empty() {
749            return 0;
750        }
751
752        entries.sort_unstable();
753        let mut hasher = std::hash::DefaultHasher::new();
754        for (name, mtime) in &entries {
755            name.hash(&mut hasher);
756            mtime.hash(&mut hasher);
757        }
758        let h = hasher.finish();
759        if h == 0 { 1 } else { h }
760    }
761}
762
763/// Exclusive ownership of a task, released on drop.
764#[derive(Debug)]
765pub struct Claim {
766    path: PathBuf,
767}
768
769impl Drop for Claim {
770    fn drop(&mut self) {
771        let _ = std::fs::remove_file(&self.path);
772    }
773}
774
775/// The first line of a task, trimmed to a title. Used when the caller gives a
776/// body but no title, which is the normal case for an agent piping a file in.
777pub fn title_from(instruction: &str, max: usize) -> String {
778    // The first non-blank line, whatever it is. A markdown heading is the
779    // task's own summary - agents pipe in `# Rework the config loader` and mean
780    // exactly that - so it is preferred over the prose beneath it rather than
781    // skipped as decoration. Leading list and heading markers are stripped
782    // because they are syntax, not words.
783    let line = instruction
784        .lines()
785        .map(str::trim)
786        .find(|l| !l.is_empty())
787        .unwrap_or("(empty task)")
788        .trim_start_matches(['#', '-', '*', '>', ' '])
789        .trim();
790    if line.is_empty() {
791        return "(empty task)".to_owned();
792    }
793    if line.chars().count() <= max {
794        return line.to_owned();
795    }
796    let head: String = line.chars().take(max.saturating_sub(1)).collect();
797    format!("{head}…")
798}
799
800fn read_path(path: &Path) -> Result<Task> {
801    let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
802    let task: Task =
803        serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
804    // Greater-than, not not-equal: every field added since schema 1 carries
805    // `#[serde(default)]`, so an older task has nothing to say about it and
806    // defaulting is exactly as good a reading as a value that build never had
807    // a chance to write. Only a schema *ahead* of this build - a meaning it
808    // cannot possibly know - is refused rather than guessed at.
809    if task.schema > SCHEMA {
810        bail!(
811            "task {} was written by a different magi (schema {}, this build \
812             speaks {SCHEMA})",
813            task.id,
814            task.schema
815        );
816    }
817    Ok(task)
818}
819
820fn short(id: &str) -> &str {
821    id.split('-').next_back().unwrap_or(id)
822}
823
824fn new_id() -> String {
825    let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
826    let seed = crate::rng::entropy();
827    format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
828}
829
830#[cfg(test)]
831mod tests {
832    use super::*;
833
834    /// A queue of its own, with no process-global state - which is the point of
835    /// `Queue::at`, and why these can run in parallel.
836    fn queue() -> (tempfile::TempDir, Queue) {
837        let dir = tempfile::tempdir().unwrap();
838        let q = Queue::at(dir.path().join("queue"));
839        (dir, q)
840    }
841
842    fn task(title: &str) -> Task {
843        Task::new(
844            title.to_owned(),
845            format!("do {title}"),
846            PathBuf::from("."),
847            Source::Human,
848        )
849    }
850
851    #[test]
852    fn a_markdown_heading_is_the_title_not_decoration() {
853        // A task file's heading is the summary its author already wrote, so it
854        // beats the prose underneath. Getting this backwards was visible in the
855        // first smoke test: a task titled "# Rework the config loader" listed
856        // as "It re-reads the file on every lookup".
857        assert_eq!(
858            title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
859            "Rework the config loader"
860        );
861        assert_eq!(title_from("- fix the thing", 40), "fix the thing");
862        assert_eq!(title_from("> quoted task", 40), "quoted task");
863        // Nothing usable at all still has to produce something printable.
864        assert_eq!(title_from("   \n\n", 40), "(empty task)");
865        assert_eq!(title_from("###\n", 40), "(empty task)");
866    }
867
868    #[test]
869    fn a_long_title_is_elided_by_characters_not_bytes() {
870        // Byte truncation would split a multi-byte character and panic.
871        let long = "課題".repeat(30);
872        let title = title_from(&long, 10);
873        assert_eq!(title.chars().count(), 10);
874        assert!(title.ends_with('…'));
875    }
876
877    #[test]
878    fn priority_wins_and_ties_break_oldest_first() {
879        let (_dir, q) = queue();
880        let mut a = task("first");
881        let mut b = task("second");
882        let mut c = task("urgent");
883        // Ids carry a timestamp, so force a known order.
884        a.id = "20260101-000001-aaaa".to_owned();
885        b.id = "20260101-000002-bbbb".to_owned();
886        c.id = "20260101-000003-cccc".to_owned();
887        c.priority = 5;
888        for t in [&mut a, &mut b, &mut c] {
889            q.put(t).unwrap();
890        }
891
892        // Priority first...
893        assert_eq!(q.next_runnable().unwrap().id, c.id);
894        c.hold_machine(None);
895        q.put(&mut c).unwrap();
896        // ...then oldest, so a burst of new work cannot starve older work.
897        assert_eq!(q.next_runnable().unwrap().id, a.id);
898        assert_eq!(q.list().len(), 3, "b is still waiting its turn");
899    }
900
901    #[test]
902    fn a_blocked_task_never_starves_another_runnable_one() {
903        let (_dir, q) = queue();
904        let mut blocked = task("blocked");
905        blocked.block(vec!["something".to_owned()], None);
906        q.put(&mut blocked).unwrap();
907
908        let mut runnable = task("free to go");
909        q.put(&mut runnable).unwrap();
910
911        let next = q.next_runnable().expect("a runnable task is still offered");
912        assert_eq!(next.id, runnable.id);
913    }
914
915    #[test]
916    fn a_held_task_is_never_offered_to_the_loop() {
917        let (_dir, q) = queue();
918        let mut t = task("held");
919        q.put(&mut t).unwrap();
920        assert!(q.next_runnable().is_some());
921
922        t.hold_machine(None);
923        q.put(&mut t).unwrap();
924        assert!(
925            q.next_runnable().is_none(),
926            "a held task must wait for a human"
927        );
928
929        // A failed task, by contrast, is exactly what the loop should retry.
930        t.status = TaskStatus::Failed;
931        q.put(&mut t).unwrap();
932        assert!(q.next_runnable().is_some());
933    }
934
935    #[test]
936    fn attempts_are_capped_and_then_the_task_is_held() {
937        let mut t = task("doomed");
938
939        t.start("run-1".to_owned());
940        t.fail("gate red", 2);
941        assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
942
943        t.start("run-2".to_owned());
944        t.fail("gate red", 2);
945        assert_eq!(
946            t.status,
947            TaskStatus::Held,
948            "out of attempts: stop spending money on it"
949        );
950        assert_eq!(t.runs, ["run-1", "run-2"]);
951        assert_eq!(t.last_error.as_deref(), Some("gate red"));
952    }
953
954    #[test]
955    fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
956        let mut t = task("stalled by quota");
957
958        t.start("run-1".to_owned());
959        assert_eq!(t.attempts, 1);
960        t.stall("judge-1, judge-2 out of quota");
961        assert_eq!(
962            t.attempts, 0,
963            "a closed quota window must not spend the task's retry budget"
964        );
965        assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
966        assert_eq!(
967            t.last_error.as_deref(),
968            Some("judge-1, judge-2 out of quota")
969        );
970
971        // A task can therefore stall all night and still get its real attempts
972        // once the quota resets - which is the whole point.
973        for _ in 0..20 {
974            t.start("run-n".to_owned());
975            t.stall("still out of quota");
976        }
977        t.start("run-real".to_owned());
978        t.fail("gate red", 2);
979        assert_eq!(
980            t.status,
981            TaskStatus::Failed,
982            "the first attempt that was really judged is attempt one"
983        );
984    }
985
986    #[test]
987    fn releasing_a_held_task_gives_it_a_real_second_chance() {
988        let mut t = task("retry me");
989        t.start("run-1".to_owned());
990        t.fail("gate red", 1);
991        assert_eq!(t.status, TaskStatus::Held);
992
993        t.release();
994        assert_eq!(t.status, TaskStatus::Queued);
995        // Without resetting attempts the next failure would re-hold at once,
996        // and a release would be a no-op the operator cannot see.
997        assert_eq!(t.attempts, 0);
998        assert!(t.last_error.is_none());
999        assert_eq!(
1000            t.runs.len(),
1001            1,
1002            "history is kept: attempts reset, evidence does not"
1003        );
1004    }
1005
1006    #[test]
1007    fn a_hold_reason_survives_and_a_release_clears_it() {
1008        let mut t = task("waiting on something else");
1009        t.hold_manual(Some(
1010            "waiting for 20260101-000000-aaaa to land first".to_owned(),
1011        ));
1012        assert_eq!(t.status, TaskStatus::Held);
1013        assert_eq!(
1014            t.hold_reason.as_deref(),
1015            Some("waiting for 20260101-000000-aaaa to land first")
1016        );
1017
1018        // Holding again with no reason must not erase the one already there.
1019        t.hold_manual(None);
1020        assert_eq!(
1021            t.hold_reason.as_deref(),
1022            Some("waiting for 20260101-000000-aaaa to land first"),
1023            "a bare re-hold keeps whatever a human already wrote down"
1024        );
1025
1026        // A hold with no reason at all is still an ordinary, allowed hold.
1027        let mut plain = task("no reason given");
1028        plain.hold_manual(None);
1029        assert_eq!(plain.status, TaskStatus::Held);
1030        assert!(plain.hold_reason.is_none());
1031
1032        t.release();
1033        assert_eq!(t.status, TaskStatus::Queued);
1034        assert!(
1035            t.hold_reason.is_none(),
1036            "a stale reason must not greet the next person who holds this task"
1037        );
1038    }
1039
1040    #[test]
1041    fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
1042        // `done` can close a held task directly - neither `magi task done`
1043        // nor `POST /api/queue/{id}/done` requires a release first - so a
1044        // task held for "waiting on 3ed9" and then closed without ever being
1045        // released must not still read as waiting on it afterwards.
1046        let mut t = task("landed by hand while held");
1047        t.hold_manual(Some("waiting on 3ed9".to_owned()));
1048        assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
1049
1050        t.succeed();
1051        assert_eq!(t.status, TaskStatus::Done);
1052        assert!(
1053            t.hold_reason.is_none(),
1054            "a done task cannot still be waiting on something"
1055        );
1056    }
1057
1058    #[test]
1059    fn a_blocked_task_is_never_offered_to_the_loop() {
1060        let mut t = task("blocked");
1061        assert!(t.status.runnable());
1062        t.block(
1063            vec!["dep-id".to_owned()],
1064            Some("waits on dep-id".to_owned()),
1065        );
1066        assert_eq!(t.status, TaskStatus::Blocked);
1067        assert!(!t.status.runnable());
1068        assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
1069    }
1070
1071    #[test]
1072    fn unblocking_the_last_dependency_returns_the_task_to_queued() {
1073        let mut t = task("blocked on two");
1074        t.block(
1075            vec!["a".to_owned(), "b".to_owned()],
1076            Some("waits on a and b".to_owned()),
1077        );
1078
1079        t.unblock("a");
1080        assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
1081        assert_eq!(t.blocked_by, ["b"]);
1082
1083        t.unblock("b");
1084        assert_eq!(t.status, TaskStatus::Queued);
1085        assert!(t.blocked_by.is_empty());
1086        assert!(t.block_reason.is_none());
1087    }
1088
1089    #[test]
1090    fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
1091        let mut t = task("never blocked");
1092        t.unblock("whatever");
1093        assert_eq!(t.status, TaskStatus::Queued);
1094    }
1095
1096    #[test]
1097    fn answering_a_question_is_recorded_and_survives_a_release() {
1098        let mut t = task("asked something");
1099        t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
1100        t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
1101        t.unblock("q1");
1102        assert_eq!(t.status, TaskStatus::Queued);
1103        assert_eq!(t.answers.len(), 1);
1104        assert_eq!(t.answers[0].answer, "SQLite");
1105
1106        // A release resets attempts, not evidence - the same rule
1107        // `releasing_a_held_task_gives_it_a_real_second_chance` asserts for
1108        // `runs`.
1109        t.release();
1110        assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
1111    }
1112
1113    #[test]
1114    fn requesting_review_requeues_the_task_and_remembers_the_branch() {
1115        let mut t = task("blocked run with a surviving branch");
1116        t.start("run-1".to_owned());
1117        t.fail("blocked with major findings", 5);
1118        assert_eq!(t.status, TaskStatus::Failed);
1119
1120        t.request_review("magi/eba2/A".to_owned());
1121        assert_eq!(t.status, TaskStatus::Queued);
1122        assert_eq!(t.attempts, 0);
1123        assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
1124
1125        // An ordinary release (a human overriding the choice) drops it again.
1126        t.release();
1127        assert!(t.review_branch.is_none());
1128    }
1129
1130    #[test]
1131    fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
1132        let mut t = task("retry");
1133        t.start("run-1".to_owned());
1134        t.requeue();
1135        assert!(t.fresh_start);
1136
1137        t.release();
1138        assert!(!t.fresh_start);
1139    }
1140
1141    #[test]
1142    fn priority_can_be_changed_while_queued_but_not_while_running() {
1143        let mut t = task("reprioritise me");
1144        t.set_priority(5).unwrap();
1145        assert_eq!(t.priority, 5);
1146
1147        t.start("run-1".to_owned());
1148        let err = t.set_priority(9).unwrap_err().to_string();
1149        assert!(err.contains("running"), "{err}");
1150        assert_eq!(t.priority, 5, "the rejected write must not partially apply");
1151    }
1152
1153    #[test]
1154    fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
1155        let (_dir, q) = queue();
1156        let mut a = task("first filed");
1157        let mut b = task("second filed");
1158        a.id = "20260101-000001-aaaa".to_owned();
1159        b.id = "20260101-000002-bbbb".to_owned();
1160        q.put(&mut a).unwrap();
1161        q.put(&mut b).unwrap();
1162
1163        assert_eq!(
1164            q.next_runnable().unwrap().id,
1165            a.id,
1166            "with equal priority the older task goes first, so a burst of \
1167             new work cannot starve it"
1168        );
1169        assert_eq!(
1170            q.list()[0].id,
1171            b.id,
1172            "but the list an operator reads is newest first, the same as \
1173             before priority existed - a's turn to run does not make it the \
1174             newest task"
1175        );
1176
1177        let mut a = q.get(&a.id).unwrap();
1178        a.set_priority(10).unwrap();
1179        q.put(&mut a).unwrap();
1180
1181        assert_eq!(
1182            q.next_runnable().unwrap().id,
1183            a.id,
1184            "a raised priority must be reflected the moment it is saved"
1185        );
1186        // `magi task list` and `GET /api/queue` both print `Queue::list()`
1187        // directly, so the raised task has to lead there too - not only in
1188        // what the loop would claim next.
1189        assert_eq!(
1190            q.list()[0].id,
1191            a.id,
1192            "the raised task must sort first in the list an operator reads, \
1193             not only in next_runnable's own ordering"
1194        );
1195    }
1196
1197    #[test]
1198    fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
1199        let mut t = Task::new(
1200            "old title".to_owned(),
1201            "old instruction".to_owned(),
1202            PathBuf::from("/repo"),
1203            Source::Agent {
1204                run: "20260101-000000-beef".to_owned(),
1205                node: "implement".to_owned(),
1206            },
1207        );
1208        let id = t.id.clone();
1209        let created_at = t.created_at;
1210        t.runs.push("20260101-000000-beef".to_owned());
1211
1212        t.edit("new title".to_owned(), "new instruction".to_owned())
1213            .unwrap();
1214
1215        assert_eq!(t.title, "new title");
1216        assert_eq!(t.instruction, "new instruction");
1217        assert_eq!(t.id, id, "editing must not mint a new id");
1218        assert_eq!(t.created_at, created_at);
1219        assert_eq!(
1220            t.source,
1221            Source::Agent {
1222                run: "20260101-000000-beef".to_owned(),
1223                node: "implement".to_owned(),
1224            },
1225            "editing must not turn agent attribution into human"
1226        );
1227        assert_eq!(t.runs, ["20260101-000000-beef"]);
1228    }
1229
1230    #[test]
1231    fn editing_is_refused_once_a_task_is_running_or_finished() {
1232        let mut running = task("in flight");
1233        running.start("run-1".to_owned());
1234        let err = running
1235            .edit("x".to_owned(), "y".to_owned())
1236            .unwrap_err()
1237            .to_string();
1238        assert!(err.contains("running"), "{err}");
1239
1240        let mut done = task("finished");
1241        done.succeed();
1242        let err = done
1243            .edit("x".to_owned(), "y".to_owned())
1244            .unwrap_err()
1245            .to_string();
1246        assert!(err.contains("done"), "{err}");
1247
1248        // Both queued and held are the point of the feature and must work.
1249        let mut queued = task("waiting");
1250        queued.edit("x".to_owned(), "y".to_owned()).unwrap();
1251        let mut held = task("parked");
1252        held.hold_machine(None);
1253        held.edit("x".to_owned(), "y".to_owned()).unwrap();
1254    }
1255
1256    #[test]
1257    fn a_task_recorded_without_a_hold_reason_still_reads_as_none() {
1258        let (_dir, q) = queue();
1259        let path = q.path_of("20260101-000000-aaaa");
1260        std::fs::create_dir_all(q.root()).unwrap();
1261        std::fs::write(
1262            &path,
1263            serde_json::json!({
1264                "schema": SCHEMA,
1265                "id": "20260101-000000-aaaa",
1266                "title": "from before hold reasons existed",
1267                "instruction": "from before hold reasons existed",
1268                "repo": ".",
1269                "source": { "kind": "human" },
1270                "status": "held",
1271                "created_at": Timestamp::now().to_string(),
1272                "updated_at": Timestamp::now().to_string(),
1273            })
1274            .to_string(),
1275        )
1276        .unwrap();
1277
1278        let task = q.get("20260101-000000-aaaa").expect("must still read");
1279        assert!(task.hold_reason.is_none());
1280        assert!(task.operator_held());
1281    }
1282
1283    #[test]
1284    fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
1285        let (_dir, q) = queue();
1286        let path = q.path_of("20260101-000000-bbbb");
1287        std::fs::create_dir_all(q.root()).unwrap();
1288        std::fs::write(
1289            &path,
1290            serde_json::json!({
1291                "schema": 2,
1292                "id": "20260101-000000-bbbb",
1293                "title": "old manual recovery",
1294                "instruction": "old manual recovery",
1295                "repo": ".",
1296                "source": { "kind": "human" },
1297                "status": "held",
1298                "hold_reason": "active manual recovery run20260912-224242-daf5",
1299                "created_at": Timestamp::now().to_string(),
1300                "updated_at": Timestamp::now().to_string(),
1301            })
1302            .to_string(),
1303        )
1304        .unwrap();
1305
1306        let task = q.get("20260101-000000-bbbb").expect("must still read");
1307        assert_eq!(task.hold_source, None);
1308        assert!(task.operator_held());
1309    }
1310
1311    #[test]
1312    fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
1313        let (_dir, q) = queue();
1314        let path = q.path_of("20260101-000000-aaaa");
1315        std::fs::create_dir_all(q.root()).unwrap();
1316        std::fs::write(
1317            &path,
1318            serde_json::json!({
1319                "schema": SCHEMA,
1320                "id": "20260101-000000-aaaa",
1321                "title": "from before diagnostics existed",
1322                "instruction": "from before diagnostics existed",
1323                "repo": ".",
1324                "source": { "kind": "human" },
1325                "status": "held",
1326                "created_at": Timestamp::now().to_string(),
1327                "updated_at": Timestamp::now().to_string(),
1328            })
1329            .to_string(),
1330        )
1331        .unwrap();
1332
1333        let task = q.get("20260101-000000-aaaa").expect("must still read");
1334        assert!(task.diagnostic.is_none());
1335    }
1336
1337    #[test]
1338    fn a_schema_1_task_with_no_blocking_fields_still_reads() {
1339        // Written by a build that predates `blocked_by`, `block_reason`,
1340        // `answers` and `review_branch` entirely - literal `"schema": 1`,
1341        // not `SCHEMA`, since the whole point is a build older than this one.
1342        let (_dir, q) = queue();
1343        let path = q.path_of("20260101-000000-aaaa");
1344        std::fs::create_dir_all(q.root()).unwrap();
1345        std::fs::write(
1346            &path,
1347            serde_json::json!({
1348                "schema": 1,
1349                "id": "20260101-000000-aaaa",
1350                "title": "from before blocking existed",
1351                "instruction": "from before blocking existed",
1352                "repo": ".",
1353                "source": { "kind": "human" },
1354                "status": "queued",
1355                "created_at": Timestamp::now().to_string(),
1356                "updated_at": Timestamp::now().to_string(),
1357            })
1358            .to_string(),
1359        )
1360        .unwrap();
1361
1362        let task = q.get("20260101-000000-aaaa").expect("must still read");
1363        assert!(task.blocked_by.is_empty());
1364        assert!(task.block_reason.is_none());
1365        assert!(task.answers.is_empty());
1366        assert!(task.review_branch.is_none());
1367    }
1368
1369    #[test]
1370    fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
1371        // A diagnostic belongs to the run that produced it. Left in place
1372        // across a release, an unrelated later failure - a config error, say -
1373        // would go on showing evidence for a problem that is no longer why the
1374        // task is stuck.
1375        let mut held = task("diagnosed");
1376        held.start("run-1".to_owned());
1377        held.fail("gate red", 1);
1378        held.diagnostic = Some("cargo test failed: ...".to_owned());
1379        assert_eq!(held.status, TaskStatus::Held);
1380
1381        held.release();
1382        assert!(held.diagnostic.is_none());
1383
1384        held.diagnostic = Some("cargo test failed: ...".to_owned());
1385        held.succeed();
1386        assert!(held.diagnostic.is_none());
1387    }
1388
1389    #[test]
1390    fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
1391        let mut t = task("retried");
1392        t.start("run-1".to_owned());
1393        t.diagnostic = Some("stale evidence from a previous hold".to_owned());
1394        t.fail("unrelated config error", 5);
1395        assert_eq!(t.status, TaskStatus::Failed);
1396        assert!(
1397            t.diagnostic.is_none(),
1398            "fail() must not let an old diagnostic outlive the run that produced it"
1399        );
1400    }
1401
1402    #[test]
1403    fn a_claim_is_exclusive_and_releases_on_drop() {
1404        let (_dir, q) = queue();
1405        let mut t = task("contended");
1406        q.put(&mut t).unwrap();
1407
1408        let held = q.claim(&t.id).unwrap();
1409        assert!(
1410            q.claim(&t.id).is_err(),
1411            "two daemons must not drive one task into two runs"
1412        );
1413        drop(held);
1414        assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
1415    }
1416
1417    #[test]
1418    fn a_round_trip_survives_disk() {
1419        let (_dir, q) = queue();
1420        let mut t = Task::new(
1421            "titled".to_owned(),
1422            "body".to_owned(),
1423            PathBuf::from("/repo"),
1424            Source::Agent {
1425                run: "20260101-000000-beef".to_owned(),
1426                node: "implement".to_owned(),
1427            },
1428        );
1429        t.priority = 3;
1430        q.put(&mut t).unwrap();
1431
1432        let back = q.get(&t.id).unwrap();
1433        assert_eq!(back.id, t.id);
1434        assert_eq!(back.priority, 3);
1435        assert_eq!(back.source.label(), "implement@beef");
1436        // A prefix is enough, the way run ids work everywhere else.
1437        assert_eq!(q.get(t.short()).unwrap().id, t.id);
1438    }
1439
1440    #[test]
1441    fn an_unreadable_task_does_not_take_the_queue_down() {
1442        let (_dir, q) = queue();
1443        let mut t = task("fine");
1444        q.put(&mut t).unwrap();
1445        std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
1446
1447        let listed = q.list();
1448        assert_eq!(listed.len(), 1, "the readable task still lists");
1449        assert_eq!(listed[0].id, t.id);
1450    }
1451
1452    #[test]
1453    fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
1454        let (_dir, q) = queue();
1455        let path = q.path_of("20260101-000000-aaaa");
1456        std::fs::create_dir_all(q.root()).unwrap();
1457        std::fs::write(
1458            &path,
1459            serde_json::json!({
1460                "schema": SCHEMA,
1461                "id": "20260101-000000-aaaa",
1462                "title": "from before solo existed",
1463                "instruction": "from before solo existed",
1464                "repo": ".",
1465                "source": { "kind": "human" },
1466                "status": "queued",
1467                "created_at": Timestamp::now().to_string(),
1468                "updated_at": Timestamp::now().to_string(),
1469            })
1470            .to_string(),
1471        )
1472        .unwrap();
1473
1474        let task = q.get("20260101-000000-aaaa").expect("must still read");
1475        assert!(!task.solo, "a queue file with no `solo` field means false");
1476    }
1477
1478    #[test]
1479    fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
1480        let (_dir, q) = queue();
1481        let mut t = task("from the future");
1482        q.put(&mut t).unwrap();
1483        let path = q.path_of(&t.id);
1484        let body = std::fs::read_to_string(&path)
1485            .unwrap()
1486            .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
1487        std::fs::write(&path, body).unwrap();
1488
1489        let err = q.get(&t.id).unwrap_err().to_string();
1490        assert!(err.contains("schema 99"), "{err}");
1491    }
1492
1493    #[test]
1494    fn revision_moves_when_the_queue_changes() {
1495        let (_dir, q) = queue();
1496        assert_eq!(q.revision(), 0, "an empty queue has no revision");
1497        let mut t = task("first");
1498        q.put(&mut t).unwrap();
1499        assert!(q.revision() > 0, "a written task moves the revision");
1500    }
1501
1502    #[test]
1503    fn revision_moves_when_deleting_an_older_task() {
1504        let (_dir, q) = queue();
1505        let mut t1 = task("older");
1506        q.put(&mut t1).unwrap();
1507        // Ensure mtime ticks forward.
1508        std::thread::sleep(std::time::Duration::from_millis(10));
1509        let mut t2 = task("newer");
1510        q.put(&mut t2).unwrap();
1511
1512        let rev_before = q.revision();
1513        q.remove(&t1.id, false).unwrap();
1514        let rev_after = q.revision();
1515
1516        assert_ne!(
1517            rev_before, rev_after,
1518            "deleting an older task must change the revision so other clients see the deletion"
1519        );
1520    }
1521
1522    #[test]
1523    fn removing_a_task_takes_it_out_of_the_listing() {
1524        let (_dir, q) = queue();
1525        let mut t = task("delete me");
1526        q.put(&mut t).unwrap();
1527        let removed = q.remove(t.short(), false).unwrap();
1528        assert_eq!(removed, t.id, "a prefix resolves before deleting");
1529        assert!(q.list().is_empty());
1530        assert!(
1531            q.remove(&t.id, false).is_err(),
1532            "removing twice is an error"
1533        );
1534    }
1535
1536    #[test]
1537    fn removing_a_task_takes_its_stale_lock_with_it() {
1538        let (_dir, q) = queue();
1539        let mut t = task("interrupted");
1540        q.put(&mut t).unwrap();
1541
1542        // A daemon killed mid-run leaves this behind. Nothing holds it: the
1543        // process that would have dropped the guard is gone.
1544        let claim = q.claim(&t.id).unwrap();
1545        std::mem::forget(claim);
1546        assert!(
1547            q.claim(&t.id).is_err(),
1548            "the orphaned lock is what makes the task look claimed"
1549        );
1550
1551        // A live daemon on this task is refused, whatever the lock says.
1552        let err = q.remove(&t.id, true).unwrap_err().to_string();
1553        assert!(err.contains("live daemon"), "{err}");
1554        assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
1555
1556        // With no daemon behind it, the lock is stale and goes with the task.
1557        q.remove(&t.id, false).unwrap();
1558        assert!(q.list().is_empty());
1559        let mut again = task("interrupted");
1560        again.id = t.id.clone();
1561        q.put(&mut again).unwrap();
1562        assert!(
1563            q.claim(&t.id).is_ok(),
1564            "a task that comes back must be claimable, which a left-behind lock would prevent"
1565        );
1566    }
1567}