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    /// *or blocked* task directly, with no release in between, so this clears
332    /// `hold_reason` and `blocked_by`/`block_reason` the same way
333    /// [`Task::release`] does. Otherwise a task held for "waiting on 3ed9", or
334    /// blocked on a dependency that never actually finished, and then closed
335    /// as done without ever being released would still read as waiting on
336    /// something in `magi task show` and on its card, after it no longer is.
337    pub fn succeed(&mut self) {
338        self.status = TaskStatus::Done;
339        self.last_error = None;
340        self.hold_reason = None;
341        self.hold_source = None;
342        self.diagnostic = None;
343        self.blocked_by.clear();
344        self.block_reason = None;
345    }
346
347    /// Record a failed attempt. Out of attempts means held for a human, rather
348    /// than retried until the money runs out.
349    ///
350    /// Clears [`Task::diagnostic`] unconditionally: it belongs to whatever run
351    /// produced it, and a caller that has one for *this* attempt sets it
352    /// itself right after calling this, once it knows the task actually ended
353    /// up [`TaskStatus::Held`] - see `daemon::diagnostic`. Without the clear, a
354    /// task released after a diagnosed hold and then failed again for an
355    /// unrelated, undiagnosed reason (a config error, say) would go on
356    /// showing the previous run's diagnostic as if it explained the new one.
357    pub fn fail(&mut self, why: impl Into<String>, max_attempts: usize) {
358        self.last_error = Some(why.into());
359        self.diagnostic = None;
360        self.status = if self.attempts >= max_attempts {
361            self.hold_source = Some(HoldSource::Machine);
362            TaskStatus::Held
363        } else {
364            TaskStatus::Failed
365        };
366    }
367
368    /// Record an attempt that failed for a reason the task is not responsible
369    /// for - the agent CLIs ran out of quota and the judging panel collapsed.
370    ///
371    /// This refunds the attempt on purpose. A quota window closing at 4am must
372    /// not spend the backlog's retry budget: the operator would come back to a
373    /// queue of held tasks that were never actually judged, and would have to
374    /// release every one by hand to find out which had a real problem. The task
375    /// goes back to `Failed`, which the loop retries, so a reset quota picks the
376    /// work up where it stopped.
377    pub fn stall(&mut self, why: impl Into<String>) {
378        self.last_error = Some(why.into());
379        self.diagnostic = None;
380        self.attempts = self.attempts.saturating_sub(1);
381        self.status = TaskStatus::Failed;
382    }
383
384    /// Whether this held task may only be released by an operator.
385    ///
386    /// Old files did not record a source. Preserve every such hold rather
387    /// than guessing that it was automatic and risking duplicate work. New
388    /// automatic holds record [`HoldSource::Machine`] and remain recoverable.
389    pub fn operator_held(&self) -> bool {
390        self.status == TaskStatus::Held && !matches!(self.hold_source, Some(HoldSource::Machine))
391    }
392
393    /// Take this task out of the loop's reach by an operator action.
394    ///
395    /// Clears `blocked_by`/`block_reason` unconditionally, the same as
396    /// [`Task::release`] and for the same reason its own comment already
397    /// gives: a human choosing to hold a *blocked* task overrides its wait
398    /// outright, the same as it overrides an ordinary hold. Without this, a
399    /// task held straight out of [`TaskStatus::Blocked`] - the web UI's "Hold"
400    /// button is reachable on a blocked task, same as "Mark done" - kept
401    /// reading as still waiting on a dependency it no longer had any claim on.
402    pub fn hold_manual(&mut self, reason: Option<String>) {
403        self.status = TaskStatus::Held;
404        if reason.is_some() {
405            self.hold_reason = reason;
406        }
407        self.hold_source = Some(HoldSource::Manual);
408        self.blocked_by.clear();
409        self.block_reason = None;
410    }
411
412    /// Take this task out of the loop's reach during automatic recovery.
413    ///
414    /// Clears `blocked_by`/`block_reason` for the same reason
415    /// [`Task::hold_manual`] does.
416    pub fn hold_machine(&mut self, reason: Option<String>) {
417        self.status = TaskStatus::Held;
418        if reason.is_some() {
419            self.hold_reason = reason;
420        }
421        self.hold_source = Some(HoldSource::Machine);
422        self.blocked_by.clear();
423        self.block_reason = None;
424    }
425
426    /// Block this task on other task ids and/or open question ids, chosen by
427    /// `crate::conduct`. Pure: the caller still owns writing it back with
428    /// [`Queue::put`].
429    pub fn block(&mut self, blocked_by: Vec<String>, reason: Option<String>) {
430        self.status = TaskStatus::Blocked;
431        self.blocked_by = blocked_by;
432        self.block_reason = reason;
433    }
434
435    /// Remove one resolved dependency (a task id that became [`TaskStatus::Done`],
436    /// or a question id that became [`crate::ask::QuestionStatus::Answered`]).
437    /// Once nothing is left in [`Task::blocked_by`], the task returns to
438    /// [`TaskStatus::Queued`] on its own - deciding *why* a task was blocked
439    /// was `crate::conduct`'s job, but noticing a dependency resolved needs no
440    /// model at all.
441    ///
442    /// A no-op, on purpose, for a task that is not [`TaskStatus::Blocked`]:
443    /// `crate::daemon`'s deterministic resolver runs over every task on every
444    /// poll, and a task that moved on for some other reason must not be
445    /// dragged back to `Queued` by a stale id it still happens to carry.
446    pub fn unblock(&mut self, resolved_id: &str) {
447        if self.status != TaskStatus::Blocked {
448            return;
449        }
450        self.blocked_by.retain(|id| id != resolved_id);
451        if self.blocked_by.is_empty() {
452            self.status = TaskStatus::Queued;
453            self.block_reason = None;
454        }
455    }
456
457    /// Record that a question `crate::conduct` asked about this task has been
458    /// answered, so the answer's content — not just the fact that the
459    /// question is gone — reaches the next conductor prompt and the next
460    /// run's instruction. See [`Task::answers`].
461    pub fn record_answer(&mut self, question: String, answer: String) {
462        self.answers.push(AnsweredQuestion { question, answer });
463    }
464
465    /// Requeue this task to reopen its last run as a review-only pass against
466    /// `branch` (`crate::graph::Runner::review`) rather than competing from
467    /// scratch. See [`Task::review_branch`].
468    pub fn request_review(&mut self, branch: String) {
469        self.release();
470        self.review_branch = Some(branch);
471    }
472
473    /// Requeue after a conductor chose a new competition. Unlike an ordinary
474    /// operator release, this deliberately does not resume the old run.
475    pub fn requeue(&mut self) {
476        self.release();
477        self.fresh_start = true;
478    }
479
480    /// Change how urgently this task should run next.
481    ///
482    /// Refused once the task is `running`: priority only feeds the sort
483    /// [`Queue::next_runnable`] does over tasks waiting to be claimed, and a
484    /// running task has already left that pool. Accepting the write anyway
485    /// would look like it worked while changing nothing until - and unless -
486    /// this attempt fails and the task becomes runnable again, which is a
487    /// surprise the phone should not hand back as a success.
488    pub fn set_priority(&mut self, priority: i32) -> Result<()> {
489        if self.status == TaskStatus::Running {
490            bail!(
491                "task {} is running; its priority cannot be changed until \
492                 this attempt finishes",
493                self.short()
494            );
495        }
496        self.priority = priority;
497        Ok(())
498    }
499
500    /// Replace this task's title and instruction wholesale.
501    ///
502    /// Restricted to `queued` and `held`. A `running` task's instruction has
503    /// already been handed to the graph, so a run in flight and the file on
504    /// disk must not be allowed to disagree about what was asked; a `done` or
505    /// `failed` task is a record of what actually happened and editing it
506    /// after the fact would falsify that record. `id`, `created_at`,
507    /// `source`, and `runs` are left untouched on purpose - an edit stands in
508    /// for "delete and refile", and keeping the id, the timestamp, the
509    /// attribution, and the run history is the entire reason it exists
510    /// instead.
511    pub fn edit(&mut self, title: String, instruction: String) -> Result<()> {
512        if !matches!(self.status, TaskStatus::Queued | TaskStatus::Held) {
513            bail!(
514                "task {} is {}; only a queued or held task's instruction can \
515                 be edited",
516                self.short(),
517                self.status.as_str()
518            );
519        }
520        self.title = title;
521        self.instruction = instruction;
522        Ok(())
523    }
524
525    /// Record a run that produced a pull request without merging it.
526    ///
527    /// The task is held rather than retried, and it costs no further attempt
528    /// either way. The work the task asked for exists: it is sitting on a
529    /// branch, in a pull request, waiting for CI or for a person. Retrying
530    /// would spend the whole competition budget a second time and then race a
531    /// second branch against the pull request the first one opened - which is
532    /// exactly what happened to run 01c2, whose finished and green pull request
533    /// was re-competed from scratch four seconds after it opened.
534    ///
535    /// A pull request nobody merged is a request for a person, not a failure.
536    pub fn handed_off(&mut self, why: impl Into<String>) {
537        self.last_error = Some(why.into());
538        self.diagnostic = None;
539        self.status = TaskStatus::Held;
540        self.hold_source = Some(HoldSource::Machine);
541    }
542
543    /// Put a held or finished task back in line, with its attempt count reset
544    /// so a release is a real second chance rather than an instant re-hold.
545    /// The run history is kept: attempts reset, evidence does not.
546    pub fn release(&mut self) {
547        self.status = TaskStatus::Queued;
548        self.attempts = 0;
549        self.last_error = None;
550        // Otherwise the next person who holds this task reads a reason that
551        // belonged to whatever it was waiting on last time.
552        self.hold_reason = None;
553        self.hold_source = None;
554        self.diagnostic = None;
555        // A release also un-blocks: the dependency or question `blocked_by`
556        // named may still be unresolved, but a human (or `crate::conduct`)
557        // choosing to release the task overrides that wait outright, the same
558        // as it overrides an ordinary hold.
559        self.blocked_by.clear();
560        self.block_reason = None;
561        self.review_branch = None;
562        self.fresh_start = false;
563    }
564}
565
566/// A queue on disk.
567#[derive(Debug, Clone)]
568pub struct Queue {
569    root: PathBuf,
570}
571
572impl Queue {
573    /// The operator's queue, `<home>/queue`.
574    pub fn open() -> Self {
575        Self::at(crate::run::home().join("queue"))
576    }
577
578    /// A queue at an explicit root. Tests use this; so could an operator who
579    /// wants a queue per project.
580    pub fn at(root: PathBuf) -> Self {
581        Self { root }
582    }
583
584    /// Directory holding the task files.
585    pub fn root(&self) -> &Path {
586        &self.root
587    }
588
589    /// Path for one task id.
590    pub fn path_of(&self, id: &str) -> PathBuf {
591        self.root.join(format!("{id}.json"))
592    }
593
594    /// Write a task, atomically, so a daemon killed mid-write leaves the
595    /// previous state readable rather than a truncated file.
596    pub fn put(&self, task: &mut Task) -> Result<()> {
597        task.updated_at = Timestamp::now();
598        std::fs::create_dir_all(&self.root)
599            .with_context(|| format!("create {}", self.root.display()))?;
600        let body = serde_json::to_string_pretty(task).context("serialize task")?;
601        let path = self.path_of(&task.id);
602        let tmp = path.with_extension("json.tmp");
603        std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
604        std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
605        Ok(())
606    }
607
608    /// Load a task by id or unambiguous id prefix.
609    pub fn get(&self, id: &str) -> Result<Task> {
610        let resolved = self.resolve_id(id)?;
611        read_path(&self.path_of(&resolved))
612    }
613
614    /// Remove a task, and the claim lock that belongs to it.
615    ///
616    /// `in_flight` comes from the caller — a live daemon's heartbeat naming
617    /// this task — because the task's own `running` status cannot answer the
618    /// question. A daemon killed mid-competition leaves the status at
619    /// `running` and an orphaned `.lock` behind, and a guard that trusted
620    /// either would make the task undeletable for good: the phone showed
621    /// exactly that, refusing a task whose daemon had been gone for an hour.
622    ///
623    /// So the lock is removed with the task rather than respected. Any lock
624    /// still there once no live daemon claims the task is by definition stale,
625    /// and leaving it would make a deleted task look claimed to
626    /// [`Queue::claim`] and to whoever reads the directory.
627    pub fn remove(&self, id: &str, in_flight: bool) -> Result<String> {
628        let resolved = self.resolve_id(id)?;
629        if in_flight {
630            bail!("task {resolved} is being run by a live daemon right now");
631        }
632        let path = self.path_of(&resolved);
633        std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
634        let lock = self.lock_path(&resolved);
635        if let Err(e) = std::fs::remove_file(&lock) {
636            if e.kind() != std::io::ErrorKind::NotFound {
637                return Err(e).with_context(|| format!("remove {}", lock.display()));
638            }
639        }
640        Ok(resolved)
641    }
642
643    /// Path of the claim lock for a task. One definition, so `claim` and
644    /// `remove` cannot end up naming different files.
645    fn lock_path(&self, id: &str) -> PathBuf {
646        self.root.join(format!("{id}.lock"))
647    }
648
649    /// Every task on disk, highest priority first and newest first within a
650    /// priority. This is what `magi task list` and `GET /api/queue` print, so
651    /// a raised priority has to move a task here the moment it is saved, not
652    /// only in [`Queue::next_runnable`]'s own ordering - the operator reading
653    /// the backlog and the loop about to drain it must agree on what "first"
654    /// means. Every existing task defaults to priority 0, so this is a no-op
655    /// change from the old newest-first order for a queue nobody has
656    /// reprioritised.
657    ///
658    /// Unreadable files are skipped rather than fatal: one corrupt task must
659    /// not take the queue - or the web UI, or an unattended daemon - down
660    /// with it.
661    pub fn list(&self) -> Vec<Task> {
662        let mut tasks: Vec<Task> = std::fs::read_dir(&self.root)
663            .into_iter()
664            .flatten()
665            .flatten()
666            .map(|e| e.path())
667            .filter(|p| p.extension().is_some_and(|x| x == "json"))
668            .filter_map(|p| read_path(&p).ok())
669            .collect();
670        tasks.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then_with(|| b.id.cmp(&a.id)));
671        tasks
672    }
673
674    /// The task a daemon should run next, or `None` when the queue is idle.
675    ///
676    /// Highest priority first, oldest first within a priority, so a burst of
677    /// agent-filed work cannot starve the task a human filed this morning.
678    pub fn next_runnable(&self) -> Option<Task> {
679        let mut runnable: Vec<Task> = self
680            .list()
681            .into_iter()
682            .filter(|t| t.status.runnable())
683            .collect();
684        runnable.sort_unstable_by(|a, b| b.priority.cmp(&a.priority).then(a.id.cmp(&b.id)));
685        runnable.into_iter().next()
686    }
687
688    /// Take exclusive ownership of a task.
689    ///
690    /// The lock is a `create_new` file next to the task, which is atomic on
691    /// every platform magi targets. It exists so two daemons - or a daemon and
692    /// a human running `magi run` - cannot drive one task into two competing
693    /// runs. The returned guard releases on drop, including on panic.
694    pub fn claim(&self, id: &str) -> Result<Claim> {
695        std::fs::create_dir_all(&self.root)
696            .with_context(|| format!("create {}", self.root.display()))?;
697        let path = self.lock_path(id);
698        match std::fs::OpenOptions::new()
699            .write(true)
700            .create_new(true)
701            .open(&path)
702        {
703            Ok(mut f) => {
704                use std::io::Write as _;
705                // Best effort: the pid is for the human looking at a stale lock.
706                let _ = writeln!(f, "{}", std::process::id());
707                Ok(Claim { path })
708            }
709            Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
710                bail!("task {id} is already claimed ({} exists)", path.display())
711            }
712            Err(e) => Err(e).with_context(|| format!("lock {}", path.display())),
713        }
714    }
715
716    /// Expand an id prefix to exactly one task id.
717    pub fn resolve_id(&self, prefix: &str) -> Result<String> {
718        if self.path_of(prefix).is_file() {
719            return Ok(prefix.to_owned());
720        }
721        let hits: Vec<String> = self
722            .list()
723            .into_iter()
724            .map(|t| t.id)
725            .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
726            .collect();
727        match hits.len() {
728            1 => Ok(hits.into_iter().next().expect("exactly one hit")),
729            0 => bail!("no task matches `{prefix}`"),
730            _ => bail!(
731                "`{prefix}` matches {} tasks: {}",
732                hits.len(),
733                hits.join(", ")
734            ),
735        }
736    }
737
738    /// Change detection token for the queue.
739    ///
740    /// Combines file names and modification times of all task files in the
741    /// queue, so adding, modifying, or deleting any task — even an older one —
742    /// moves the revision and notifies connected clients via the change stream.
743    /// Returns 0 when the queue is completely empty.
744    pub fn revision(&self) -> u64 {
745        use std::hash::{Hash as _, Hasher as _};
746
747        let mut entries: Vec<(String, u64)> = std::fs::read_dir(&self.root)
748            .into_iter()
749            .flatten()
750            .flatten()
751            .filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
752            .filter_map(|e| {
753                let name = e.file_name().to_string_lossy().into_owned();
754                let mtime = e
755                    .metadata()
756                    .ok()?
757                    .modified()
758                    .ok()?
759                    .duration_since(std::time::UNIX_EPOCH)
760                    .ok()?
761                    .as_millis() as u64;
762                Some((name, mtime))
763            })
764            .collect();
765
766        if entries.is_empty() {
767            return 0;
768        }
769
770        entries.sort_unstable();
771        let mut hasher = std::hash::DefaultHasher::new();
772        for (name, mtime) in &entries {
773            name.hash(&mut hasher);
774            mtime.hash(&mut hasher);
775        }
776        let h = hasher.finish();
777        if h == 0 { 1 } else { h }
778    }
779}
780
781/// Exclusive ownership of a task, released on drop.
782#[derive(Debug)]
783pub struct Claim {
784    path: PathBuf,
785}
786
787impl Drop for Claim {
788    fn drop(&mut self) {
789        let _ = std::fs::remove_file(&self.path);
790    }
791}
792
793/// The first line of a task, trimmed to a title. Used when the caller gives a
794/// body but no title, which is the normal case for an agent piping a file in.
795pub fn title_from(instruction: &str, max: usize) -> String {
796    // The first non-blank line, whatever it is. A markdown heading is the
797    // task's own summary - agents pipe in `# Rework the config loader` and mean
798    // exactly that - so it is preferred over the prose beneath it rather than
799    // skipped as decoration. Leading list and heading markers are stripped
800    // because they are syntax, not words.
801    let line = instruction
802        .lines()
803        .map(str::trim)
804        .find(|l| !l.is_empty())
805        .unwrap_or("(empty task)")
806        .trim_start_matches(['#', '-', '*', '>', ' '])
807        .trim();
808    if line.is_empty() {
809        return "(empty task)".to_owned();
810    }
811    if line.chars().count() <= max {
812        return line.to_owned();
813    }
814    let head: String = line.chars().take(max.saturating_sub(1)).collect();
815    format!("{head}…")
816}
817
818fn read_path(path: &Path) -> Result<Task> {
819    let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
820    let task: Task =
821        serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
822    // Greater-than, not not-equal: every field added since schema 1 carries
823    // `#[serde(default)]`, so an older task has nothing to say about it and
824    // defaulting is exactly as good a reading as a value that build never had
825    // a chance to write. Only a schema *ahead* of this build - a meaning it
826    // cannot possibly know - is refused rather than guessed at.
827    if task.schema > SCHEMA {
828        bail!(
829            "task {} was written by a different magi (schema {}, this build \
830             speaks {SCHEMA})",
831            task.id,
832            task.schema
833        );
834    }
835    Ok(task)
836}
837
838fn short(id: &str) -> &str {
839    id.split('-').next_back().unwrap_or(id)
840}
841
842fn new_id() -> String {
843    let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
844    let seed = crate::rng::entropy();
845    format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
846}
847
848#[cfg(test)]
849mod tests {
850    use super::*;
851
852    /// A queue of its own, with no process-global state - which is the point of
853    /// `Queue::at`, and why these can run in parallel.
854    fn queue() -> (tempfile::TempDir, Queue) {
855        let dir = tempfile::tempdir().unwrap();
856        let q = Queue::at(dir.path().join("queue"));
857        (dir, q)
858    }
859
860    fn task(title: &str) -> Task {
861        Task::new(
862            title.to_owned(),
863            format!("do {title}"),
864            PathBuf::from("."),
865            Source::Human,
866        )
867    }
868
869    #[test]
870    fn a_markdown_heading_is_the_title_not_decoration() {
871        // A task file's heading is the summary its author already wrote, so it
872        // beats the prose underneath. Getting this backwards was visible in the
873        // first smoke test: a task titled "# Rework the config loader" listed
874        // as "It re-reads the file on every lookup".
875        assert_eq!(
876            title_from("# Rework the config loader\n\nIt re-reads it.\n", 40),
877            "Rework the config loader"
878        );
879        assert_eq!(title_from("- fix the thing", 40), "fix the thing");
880        assert_eq!(title_from("> quoted task", 40), "quoted task");
881        // Nothing usable at all still has to produce something printable.
882        assert_eq!(title_from("   \n\n", 40), "(empty task)");
883        assert_eq!(title_from("###\n", 40), "(empty task)");
884    }
885
886    #[test]
887    fn a_long_title_is_elided_by_characters_not_bytes() {
888        // Byte truncation would split a multi-byte character and panic.
889        let long = "課題".repeat(30);
890        let title = title_from(&long, 10);
891        assert_eq!(title.chars().count(), 10);
892        assert!(title.ends_with('…'));
893    }
894
895    #[test]
896    fn priority_wins_and_ties_break_oldest_first() {
897        let (_dir, q) = queue();
898        let mut a = task("first");
899        let mut b = task("second");
900        let mut c = task("urgent");
901        // Ids carry a timestamp, so force a known order.
902        a.id = "20260101-000001-aaaa".to_owned();
903        b.id = "20260101-000002-bbbb".to_owned();
904        c.id = "20260101-000003-cccc".to_owned();
905        c.priority = 5;
906        for t in [&mut a, &mut b, &mut c] {
907            q.put(t).unwrap();
908        }
909
910        // Priority first...
911        assert_eq!(q.next_runnable().unwrap().id, c.id);
912        c.hold_machine(None);
913        q.put(&mut c).unwrap();
914        // ...then oldest, so a burst of new work cannot starve older work.
915        assert_eq!(q.next_runnable().unwrap().id, a.id);
916        assert_eq!(q.list().len(), 3, "b is still waiting its turn");
917    }
918
919    #[test]
920    fn a_blocked_task_never_starves_another_runnable_one() {
921        let (_dir, q) = queue();
922        let mut blocked = task("blocked");
923        blocked.block(vec!["something".to_owned()], None);
924        q.put(&mut blocked).unwrap();
925
926        let mut runnable = task("free to go");
927        q.put(&mut runnable).unwrap();
928
929        let next = q.next_runnable().expect("a runnable task is still offered");
930        assert_eq!(next.id, runnable.id);
931    }
932
933    #[test]
934    fn a_held_task_is_never_offered_to_the_loop() {
935        let (_dir, q) = queue();
936        let mut t = task("held");
937        q.put(&mut t).unwrap();
938        assert!(q.next_runnable().is_some());
939
940        t.hold_machine(None);
941        q.put(&mut t).unwrap();
942        assert!(
943            q.next_runnable().is_none(),
944            "a held task must wait for a human"
945        );
946
947        // A failed task, by contrast, is exactly what the loop should retry.
948        t.status = TaskStatus::Failed;
949        q.put(&mut t).unwrap();
950        assert!(q.next_runnable().is_some());
951    }
952
953    #[test]
954    fn attempts_are_capped_and_then_the_task_is_held() {
955        let mut t = task("doomed");
956
957        t.start("run-1".to_owned());
958        t.fail("gate red", 2);
959        assert_eq!(t.status, TaskStatus::Failed, "one attempt of two: retry");
960
961        t.start("run-2".to_owned());
962        t.fail("gate red", 2);
963        assert_eq!(
964            t.status,
965            TaskStatus::Held,
966            "out of attempts: stop spending money on it"
967        );
968        assert_eq!(t.runs, ["run-1", "run-2"]);
969        assert_eq!(t.last_error.as_deref(), Some("gate red"));
970    }
971
972    #[test]
973    fn a_quota_stall_is_refunded_so_the_backlog_survives_the_night() {
974        let mut t = task("stalled by quota");
975
976        t.start("run-1".to_owned());
977        assert_eq!(t.attempts, 1);
978        t.stall("judge-1, judge-2 out of quota");
979        assert_eq!(
980            t.attempts, 0,
981            "a closed quota window must not spend the task's retry budget"
982        );
983        assert_eq!(t.status, TaskStatus::Failed, "the loop should retry it");
984        assert_eq!(
985            t.last_error.as_deref(),
986            Some("judge-1, judge-2 out of quota")
987        );
988
989        // A task can therefore stall all night and still get its real attempts
990        // once the quota resets - which is the whole point.
991        for _ in 0..20 {
992            t.start("run-n".to_owned());
993            t.stall("still out of quota");
994        }
995        t.start("run-real".to_owned());
996        t.fail("gate red", 2);
997        assert_eq!(
998            t.status,
999            TaskStatus::Failed,
1000            "the first attempt that was really judged is attempt one"
1001        );
1002    }
1003
1004    #[test]
1005    fn releasing_a_held_task_gives_it_a_real_second_chance() {
1006        let mut t = task("retry me");
1007        t.start("run-1".to_owned());
1008        t.fail("gate red", 1);
1009        assert_eq!(t.status, TaskStatus::Held);
1010
1011        t.release();
1012        assert_eq!(t.status, TaskStatus::Queued);
1013        // Without resetting attempts the next failure would re-hold at once,
1014        // and a release would be a no-op the operator cannot see.
1015        assert_eq!(t.attempts, 0);
1016        assert!(t.last_error.is_none());
1017        assert_eq!(
1018            t.runs.len(),
1019            1,
1020            "history is kept: attempts reset, evidence does not"
1021        );
1022    }
1023
1024    #[test]
1025    fn a_hold_reason_survives_and_a_release_clears_it() {
1026        let mut t = task("waiting on something else");
1027        t.hold_manual(Some(
1028            "waiting for 20260101-000000-aaaa to land first".to_owned(),
1029        ));
1030        assert_eq!(t.status, TaskStatus::Held);
1031        assert_eq!(
1032            t.hold_reason.as_deref(),
1033            Some("waiting for 20260101-000000-aaaa to land first")
1034        );
1035
1036        // Holding again with no reason must not erase the one already there.
1037        t.hold_manual(None);
1038        assert_eq!(
1039            t.hold_reason.as_deref(),
1040            Some("waiting for 20260101-000000-aaaa to land first"),
1041            "a bare re-hold keeps whatever a human already wrote down"
1042        );
1043
1044        // A hold with no reason at all is still an ordinary, allowed hold.
1045        let mut plain = task("no reason given");
1046        plain.hold_manual(None);
1047        assert_eq!(plain.status, TaskStatus::Held);
1048        assert!(plain.hold_reason.is_none());
1049
1050        t.release();
1051        assert_eq!(t.status, TaskStatus::Queued);
1052        assert!(
1053            t.hold_reason.is_none(),
1054            "a stale reason must not greet the next person who holds this task"
1055        );
1056    }
1057
1058    #[test]
1059    fn closing_a_held_task_as_done_clears_its_hold_reason_too() {
1060        // `done` can close a held task directly - neither `magi task done`
1061        // nor `POST /api/queue/{id}/done` requires a release first - so a
1062        // task held for "waiting on 3ed9" and then closed without ever being
1063        // released must not still read as waiting on it afterwards.
1064        let mut t = task("landed by hand while held");
1065        t.hold_manual(Some("waiting on 3ed9".to_owned()));
1066        assert_eq!(t.hold_reason.as_deref(), Some("waiting on 3ed9"));
1067
1068        t.succeed();
1069        assert_eq!(t.status, TaskStatus::Done);
1070        assert!(
1071            t.hold_reason.is_none(),
1072            "a done task cannot still be waiting on something"
1073        );
1074    }
1075
1076    #[test]
1077    fn holding_or_closing_a_blocked_task_clears_its_dependency_too() {
1078        // The web UI's "Hold" and "Mark done" buttons are both reachable on a
1079        // `blocked` task, not just on `queued`/`held` ones - neither requires
1080        // a release first. A task moved off `Blocked` that way must not still
1081        // carry the dependency it was waiting on: a dependency graph built
1082        // from `blocked_by` would otherwise keep drawing an edge for a task
1083        // that is not blocked on anything any more.
1084        let mut held = task("held straight out of blocked");
1085        held.block(
1086            vec!["20260101-000000-dead".to_owned()],
1087            Some("waiting on the migration script".to_owned()),
1088        );
1089        assert_eq!(held.status, TaskStatus::Blocked);
1090
1091        held.hold_manual(None);
1092        assert_eq!(held.status, TaskStatus::Held);
1093        assert!(
1094            held.blocked_by.is_empty(),
1095            "hold overrides the wait, same as release"
1096        );
1097        assert!(held.block_reason.is_none());
1098
1099        let mut done = task("closed straight out of blocked");
1100        done.block(
1101            vec!["20260101-000000-dead".to_owned()],
1102            Some("waiting on the migration script".to_owned()),
1103        );
1104        done.succeed();
1105        assert_eq!(done.status, TaskStatus::Done);
1106        assert!(
1107            done.blocked_by.is_empty(),
1108            "a done task cannot still be waiting on a dependency"
1109        );
1110        assert!(done.block_reason.is_none());
1111    }
1112
1113    #[test]
1114    fn a_blocked_task_is_never_offered_to_the_loop() {
1115        let mut t = task("blocked");
1116        assert!(t.status.runnable());
1117        t.block(
1118            vec!["dep-id".to_owned()],
1119            Some("waits on dep-id".to_owned()),
1120        );
1121        assert_eq!(t.status, TaskStatus::Blocked);
1122        assert!(!t.status.runnable());
1123        assert_eq!(TaskStatus::Blocked.as_str(), "blocked");
1124    }
1125
1126    #[test]
1127    fn unblocking_the_last_dependency_returns_the_task_to_queued() {
1128        let mut t = task("blocked on two");
1129        t.block(
1130            vec!["a".to_owned(), "b".to_owned()],
1131            Some("waits on a and b".to_owned()),
1132        );
1133
1134        t.unblock("a");
1135        assert_eq!(t.status, TaskStatus::Blocked, "b is still outstanding");
1136        assert_eq!(t.blocked_by, ["b"]);
1137
1138        t.unblock("b");
1139        assert_eq!(t.status, TaskStatus::Queued);
1140        assert!(t.blocked_by.is_empty());
1141        assert!(t.block_reason.is_none());
1142    }
1143
1144    #[test]
1145    fn unblocking_an_id_on_a_task_that_is_not_blocked_is_a_no_op() {
1146        let mut t = task("never blocked");
1147        t.unblock("whatever");
1148        assert_eq!(t.status, TaskStatus::Queued);
1149    }
1150
1151    #[test]
1152    fn answering_a_question_is_recorded_and_survives_a_release() {
1153        let mut t = task("asked something");
1154        t.block(vec!["q1".to_owned()], Some("which backend?".to_owned()));
1155        t.record_answer("Which backend?".to_owned(), "SQLite".to_owned());
1156        t.unblock("q1");
1157        assert_eq!(t.status, TaskStatus::Queued);
1158        assert_eq!(t.answers.len(), 1);
1159        assert_eq!(t.answers[0].answer, "SQLite");
1160
1161        // A release resets attempts, not evidence - the same rule
1162        // `releasing_a_held_task_gives_it_a_real_second_chance` asserts for
1163        // `runs`.
1164        t.release();
1165        assert_eq!(t.answers.len(), 1, "the answer is not lost on release");
1166    }
1167
1168    #[test]
1169    fn requesting_review_requeues_the_task_and_remembers_the_branch() {
1170        let mut t = task("blocked run with a surviving branch");
1171        t.start("run-1".to_owned());
1172        t.fail("blocked with major findings", 5);
1173        assert_eq!(t.status, TaskStatus::Failed);
1174
1175        t.request_review("magi/eba2/A".to_owned());
1176        assert_eq!(t.status, TaskStatus::Queued);
1177        assert_eq!(t.attempts, 0);
1178        assert_eq!(t.review_branch.as_deref(), Some("magi/eba2/A"));
1179
1180        // An ordinary release (a human overriding the choice) drops it again.
1181        t.release();
1182        assert!(t.review_branch.is_none());
1183    }
1184
1185    #[test]
1186    fn conductor_requeue_but_not_an_ordinary_release_forces_a_fresh_start() {
1187        let mut t = task("retry");
1188        t.start("run-1".to_owned());
1189        t.requeue();
1190        assert!(t.fresh_start);
1191
1192        t.release();
1193        assert!(!t.fresh_start);
1194    }
1195
1196    #[test]
1197    fn priority_can_be_changed_while_queued_but_not_while_running() {
1198        let mut t = task("reprioritise me");
1199        t.set_priority(5).unwrap();
1200        assert_eq!(t.priority, 5);
1201
1202        t.start("run-1".to_owned());
1203        let err = t.set_priority(9).unwrap_err().to_string();
1204        assert!(err.contains("running"), "{err}");
1205        assert_eq!(t.priority, 5, "the rejected write must not partially apply");
1206    }
1207
1208    #[test]
1209    fn changing_priority_moves_a_task_ahead_in_the_real_queue_order() {
1210        let (_dir, q) = queue();
1211        let mut a = task("first filed");
1212        let mut b = task("second filed");
1213        a.id = "20260101-000001-aaaa".to_owned();
1214        b.id = "20260101-000002-bbbb".to_owned();
1215        q.put(&mut a).unwrap();
1216        q.put(&mut b).unwrap();
1217
1218        assert_eq!(
1219            q.next_runnable().unwrap().id,
1220            a.id,
1221            "with equal priority the older task goes first, so a burst of \
1222             new work cannot starve it"
1223        );
1224        assert_eq!(
1225            q.list()[0].id,
1226            b.id,
1227            "but the list an operator reads is newest first, the same as \
1228             before priority existed - a's turn to run does not make it the \
1229             newest task"
1230        );
1231
1232        let mut a = q.get(&a.id).unwrap();
1233        a.set_priority(10).unwrap();
1234        q.put(&mut a).unwrap();
1235
1236        assert_eq!(
1237            q.next_runnable().unwrap().id,
1238            a.id,
1239            "a raised priority must be reflected the moment it is saved"
1240        );
1241        // `magi task list` and `GET /api/queue` both print `Queue::list()`
1242        // directly, so the raised task has to lead there too - not only in
1243        // what the loop would claim next.
1244        assert_eq!(
1245            q.list()[0].id,
1246            a.id,
1247            "the raised task must sort first in the list an operator reads, \
1248             not only in next_runnable's own ordering"
1249        );
1250    }
1251
1252    #[test]
1253    fn editing_replaces_title_and_instruction_but_keeps_identity_and_history() {
1254        let mut t = Task::new(
1255            "old title".to_owned(),
1256            "old instruction".to_owned(),
1257            PathBuf::from("/repo"),
1258            Source::Agent {
1259                run: "20260101-000000-beef".to_owned(),
1260                node: "implement".to_owned(),
1261            },
1262        );
1263        let id = t.id.clone();
1264        let created_at = t.created_at;
1265        t.runs.push("20260101-000000-beef".to_owned());
1266
1267        t.edit("new title".to_owned(), "new instruction".to_owned())
1268            .unwrap();
1269
1270        assert_eq!(t.title, "new title");
1271        assert_eq!(t.instruction, "new instruction");
1272        assert_eq!(t.id, id, "editing must not mint a new id");
1273        assert_eq!(t.created_at, created_at);
1274        assert_eq!(
1275            t.source,
1276            Source::Agent {
1277                run: "20260101-000000-beef".to_owned(),
1278                node: "implement".to_owned(),
1279            },
1280            "editing must not turn agent attribution into human"
1281        );
1282        assert_eq!(t.runs, ["20260101-000000-beef"]);
1283    }
1284
1285    #[test]
1286    fn editing_is_refused_once_a_task_is_running_or_finished() {
1287        let mut running = task("in flight");
1288        running.start("run-1".to_owned());
1289        let err = running
1290            .edit("x".to_owned(), "y".to_owned())
1291            .unwrap_err()
1292            .to_string();
1293        assert!(err.contains("running"), "{err}");
1294
1295        let mut done = task("finished");
1296        done.succeed();
1297        let err = done
1298            .edit("x".to_owned(), "y".to_owned())
1299            .unwrap_err()
1300            .to_string();
1301        assert!(err.contains("done"), "{err}");
1302
1303        // Both queued and held are the point of the feature and must work.
1304        let mut queued = task("waiting");
1305        queued.edit("x".to_owned(), "y".to_owned()).unwrap();
1306        let mut held = task("parked");
1307        held.hold_machine(None);
1308        held.edit("x".to_owned(), "y".to_owned()).unwrap();
1309    }
1310
1311    #[test]
1312    fn a_task_recorded_without_a_hold_reason_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 hold reasons existed",
1322                "instruction": "from before hold reasons 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.hold_reason.is_none());
1335        assert!(task.operator_held());
1336    }
1337
1338    #[test]
1339    fn a_legacy_reasoned_hold_defaults_to_operator_protection() {
1340        let (_dir, q) = queue();
1341        let path = q.path_of("20260101-000000-bbbb");
1342        std::fs::create_dir_all(q.root()).unwrap();
1343        std::fs::write(
1344            &path,
1345            serde_json::json!({
1346                "schema": 2,
1347                "id": "20260101-000000-bbbb",
1348                "title": "old manual recovery",
1349                "instruction": "old manual recovery",
1350                "repo": ".",
1351                "source": { "kind": "human" },
1352                "status": "held",
1353                "hold_reason": "active manual recovery run20260912-224242-daf5",
1354                "created_at": Timestamp::now().to_string(),
1355                "updated_at": Timestamp::now().to_string(),
1356            })
1357            .to_string(),
1358        )
1359        .unwrap();
1360
1361        let task = q.get("20260101-000000-bbbb").expect("must still read");
1362        assert_eq!(task.hold_source, None);
1363        assert!(task.operator_held());
1364    }
1365
1366    #[test]
1367    fn a_task_recorded_without_a_diagnostic_still_reads_as_none() {
1368        let (_dir, q) = queue();
1369        let path = q.path_of("20260101-000000-aaaa");
1370        std::fs::create_dir_all(q.root()).unwrap();
1371        std::fs::write(
1372            &path,
1373            serde_json::json!({
1374                "schema": SCHEMA,
1375                "id": "20260101-000000-aaaa",
1376                "title": "from before diagnostics existed",
1377                "instruction": "from before diagnostics existed",
1378                "repo": ".",
1379                "source": { "kind": "human" },
1380                "status": "held",
1381                "created_at": Timestamp::now().to_string(),
1382                "updated_at": Timestamp::now().to_string(),
1383            })
1384            .to_string(),
1385        )
1386        .unwrap();
1387
1388        let task = q.get("20260101-000000-aaaa").expect("must still read");
1389        assert!(task.diagnostic.is_none());
1390    }
1391
1392    #[test]
1393    fn a_schema_1_task_with_no_blocking_fields_still_reads() {
1394        // Written by a build that predates `blocked_by`, `block_reason`,
1395        // `answers` and `review_branch` entirely - literal `"schema": 1`,
1396        // not `SCHEMA`, since the whole point is a build older than this one.
1397        let (_dir, q) = queue();
1398        let path = q.path_of("20260101-000000-aaaa");
1399        std::fs::create_dir_all(q.root()).unwrap();
1400        std::fs::write(
1401            &path,
1402            serde_json::json!({
1403                "schema": 1,
1404                "id": "20260101-000000-aaaa",
1405                "title": "from before blocking existed",
1406                "instruction": "from before blocking existed",
1407                "repo": ".",
1408                "source": { "kind": "human" },
1409                "status": "queued",
1410                "created_at": Timestamp::now().to_string(),
1411                "updated_at": Timestamp::now().to_string(),
1412            })
1413            .to_string(),
1414        )
1415        .unwrap();
1416
1417        let task = q.get("20260101-000000-aaaa").expect("must still read");
1418        assert!(task.blocked_by.is_empty());
1419        assert!(task.block_reason.is_none());
1420        assert!(task.answers.is_empty());
1421        assert!(task.review_branch.is_none());
1422    }
1423
1424    #[test]
1425    fn releasing_or_finishing_a_task_clears_its_stale_diagnostic() {
1426        // A diagnostic belongs to the run that produced it. Left in place
1427        // across a release, an unrelated later failure - a config error, say -
1428        // would go on showing evidence for a problem that is no longer why the
1429        // task is stuck.
1430        let mut held = task("diagnosed");
1431        held.start("run-1".to_owned());
1432        held.fail("gate red", 1);
1433        held.diagnostic = Some("cargo test failed: ...".to_owned());
1434        assert_eq!(held.status, TaskStatus::Held);
1435
1436        held.release();
1437        assert!(held.diagnostic.is_none());
1438
1439        held.diagnostic = Some("cargo test failed: ...".to_owned());
1440        held.succeed();
1441        assert!(held.diagnostic.is_none());
1442    }
1443
1444    #[test]
1445    fn failing_a_task_always_clears_whatever_diagnostic_it_carried() {
1446        let mut t = task("retried");
1447        t.start("run-1".to_owned());
1448        t.diagnostic = Some("stale evidence from a previous hold".to_owned());
1449        t.fail("unrelated config error", 5);
1450        assert_eq!(t.status, TaskStatus::Failed);
1451        assert!(
1452            t.diagnostic.is_none(),
1453            "fail() must not let an old diagnostic outlive the run that produced it"
1454        );
1455    }
1456
1457    #[test]
1458    fn a_claim_is_exclusive_and_releases_on_drop() {
1459        let (_dir, q) = queue();
1460        let mut t = task("contended");
1461        q.put(&mut t).unwrap();
1462
1463        let held = q.claim(&t.id).unwrap();
1464        assert!(
1465            q.claim(&t.id).is_err(),
1466            "two daemons must not drive one task into two runs"
1467        );
1468        drop(held);
1469        assert!(q.claim(&t.id).is_ok(), "a released claim is reclaimable");
1470    }
1471
1472    #[test]
1473    fn a_round_trip_survives_disk() {
1474        let (_dir, q) = queue();
1475        let mut t = Task::new(
1476            "titled".to_owned(),
1477            "body".to_owned(),
1478            PathBuf::from("/repo"),
1479            Source::Agent {
1480                run: "20260101-000000-beef".to_owned(),
1481                node: "implement".to_owned(),
1482            },
1483        );
1484        t.priority = 3;
1485        q.put(&mut t).unwrap();
1486
1487        let back = q.get(&t.id).unwrap();
1488        assert_eq!(back.id, t.id);
1489        assert_eq!(back.priority, 3);
1490        assert_eq!(back.source.label(), "implement@beef");
1491        // A prefix is enough, the way run ids work everywhere else.
1492        assert_eq!(q.get(t.short()).unwrap().id, t.id);
1493    }
1494
1495    #[test]
1496    fn an_unreadable_task_does_not_take_the_queue_down() {
1497        let (_dir, q) = queue();
1498        let mut t = task("fine");
1499        q.put(&mut t).unwrap();
1500        std::fs::write(q.root().join("broken.json"), "{ not json").unwrap();
1501
1502        let listed = q.list();
1503        assert_eq!(listed.len(), 1, "the readable task still lists");
1504        assert_eq!(listed[0].id, t.id);
1505    }
1506
1507    #[test]
1508    fn a_task_recorded_without_a_solo_field_still_reads_as_not_solo() {
1509        let (_dir, q) = queue();
1510        let path = q.path_of("20260101-000000-aaaa");
1511        std::fs::create_dir_all(q.root()).unwrap();
1512        std::fs::write(
1513            &path,
1514            serde_json::json!({
1515                "schema": SCHEMA,
1516                "id": "20260101-000000-aaaa",
1517                "title": "from before solo existed",
1518                "instruction": "from before solo existed",
1519                "repo": ".",
1520                "source": { "kind": "human" },
1521                "status": "queued",
1522                "created_at": Timestamp::now().to_string(),
1523                "updated_at": Timestamp::now().to_string(),
1524            })
1525            .to_string(),
1526        )
1527        .unwrap();
1528
1529        let task = q.get("20260101-000000-aaaa").expect("must still read");
1530        assert!(!task.solo, "a queue file with no `solo` field means false");
1531    }
1532
1533    #[test]
1534    fn a_task_from_a_future_schema_is_refused_rather_than_guessed_at() {
1535        let (_dir, q) = queue();
1536        let mut t = task("from the future");
1537        q.put(&mut t).unwrap();
1538        let path = q.path_of(&t.id);
1539        let body = std::fs::read_to_string(&path)
1540            .unwrap()
1541            .replace(&format!("\"schema\": {SCHEMA}"), "\"schema\": 99");
1542        std::fs::write(&path, body).unwrap();
1543
1544        let err = q.get(&t.id).unwrap_err().to_string();
1545        assert!(err.contains("schema 99"), "{err}");
1546    }
1547
1548    #[test]
1549    fn revision_moves_when_the_queue_changes() {
1550        let (_dir, q) = queue();
1551        assert_eq!(q.revision(), 0, "an empty queue has no revision");
1552        let mut t = task("first");
1553        q.put(&mut t).unwrap();
1554        assert!(q.revision() > 0, "a written task moves the revision");
1555    }
1556
1557    #[test]
1558    fn revision_moves_when_deleting_an_older_task() {
1559        let (_dir, q) = queue();
1560        let mut t1 = task("older");
1561        q.put(&mut t1).unwrap();
1562        // Ensure mtime ticks forward.
1563        std::thread::sleep(std::time::Duration::from_millis(10));
1564        let mut t2 = task("newer");
1565        q.put(&mut t2).unwrap();
1566
1567        let rev_before = q.revision();
1568        q.remove(&t1.id, false).unwrap();
1569        let rev_after = q.revision();
1570
1571        assert_ne!(
1572            rev_before, rev_after,
1573            "deleting an older task must change the revision so other clients see the deletion"
1574        );
1575    }
1576
1577    #[test]
1578    fn removing_a_task_takes_it_out_of_the_listing() {
1579        let (_dir, q) = queue();
1580        let mut t = task("delete me");
1581        q.put(&mut t).unwrap();
1582        let removed = q.remove(t.short(), false).unwrap();
1583        assert_eq!(removed, t.id, "a prefix resolves before deleting");
1584        assert!(q.list().is_empty());
1585        assert!(
1586            q.remove(&t.id, false).is_err(),
1587            "removing twice is an error"
1588        );
1589    }
1590
1591    #[test]
1592    fn removing_a_task_takes_its_stale_lock_with_it() {
1593        let (_dir, q) = queue();
1594        let mut t = task("interrupted");
1595        q.put(&mut t).unwrap();
1596
1597        // A daemon killed mid-run leaves this behind. Nothing holds it: the
1598        // process that would have dropped the guard is gone.
1599        let claim = q.claim(&t.id).unwrap();
1600        std::mem::forget(claim);
1601        assert!(
1602            q.claim(&t.id).is_err(),
1603            "the orphaned lock is what makes the task look claimed"
1604        );
1605
1606        // A live daemon on this task is refused, whatever the lock says.
1607        let err = q.remove(&t.id, true).unwrap_err().to_string();
1608        assert!(err.contains("live daemon"), "{err}");
1609        assert!(q.get(&t.id).is_ok(), "a refused delete keeps the task");
1610
1611        // With no daemon behind it, the lock is stale and goes with the task.
1612        q.remove(&t.id, false).unwrap();
1613        assert!(q.list().is_empty());
1614        let mut again = task("interrupted");
1615        again.id = t.id.clone();
1616        q.put(&mut again).unwrap();
1617        assert!(
1618            q.claim(&t.id).is_ok(),
1619            "a task that comes back must be claimable, which a left-behind lock would prevent"
1620        );
1621    }
1622}