Skip to main content

magi/
talk.rs

1//! The standing conversation: a place to think out loud with an agent between
2//! tasks, reachable from a phone.
3//!
4//! This is a conversation that stays open. Ask a question, have the agent
5//! read a file or run a command to check something, talk through an idea,
6//! and when it is time to act, tell it to file the work rather than do it
7//! here. The conversation does not end; it is what the operator opens the
8//! next time something comes up.
9//!
10//! # Talking is not implementing
11//!
12//! Every turn here runs with `allow_write: false` by default, for a reason
13//! that is not security, but attribution. An agent that edits a checkout
14//! mid-conversation leaves a diff that belongs to no run and passed no
15//! review, and on a repository entered into magi's blind competition that
16//! makes every candidate's diff unjudgeable. That is why the default holds
17//! regardless of what a repository's own `magi.toml` says about anything
18//! else. When the operator wants a change made, the agent is told to run
19//! `magi task add --solo` ([`briefing`]) rather than reach for an editor: the
20//! change goes through magi's own queue, on the repository's own terms, and
21//! the operator can watch it happen instead of trusting that it did.
22//!
23//! `[talk] allow_write` ([`crate::config::Talk::allow_write`]) lets a
24//! specific repository opt out of that default - a dotfiles or personal
25//! config checkout that is never entered into a competition and never
26//! reviewed has nothing for the restriction to protect, and filing a task for
27//! a one-line edit there is pure overhead. Turning it on does not turn this
28//! conversation into an implementer: [`briefing`] still sends everything
29//! bigger than a small, operator-named edit to the queue, and still tells the
30//! agent to say what it changed.
31//!
32//! `--solo` rather than a plain `magi task add` is the point of pairing this
33//! module with [`crate::queue::Task::solo`]. A task that came out of a
34//! conversation the operator just had is a decision already made, not a
35//! design question worth three independent takes - so it runs through one
36//! implementer and straight into review, the way [`crate::graph::Runner`]
37//! already degrades a single-candidate run.
38//!
39//! # Shape
40//!
41//! The same split [`crate::queue`] uses: [`Talk`] is data plus pure helpers,
42//! [`Talks`] owns the I/O and is constructed with its root, so every test
43//! here drives a real store in a temp directory rather than the operator's
44//! own home.
45
46use std::path::{Path, PathBuf};
47use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
48use std::time::Duration;
49
50use anyhow::{Context, Result, bail};
51use jiff::Timestamp;
52use serde::{Deserialize, Serialize};
53
54use crate::agent::{self, Invocation, SeatState};
55use crate::config::{AgentSpec, Config};
56use crate::queue::{Queue, Source, Task};
57
58/// On-disk format for a conversation. Bumped when a field's meaning changes.
59pub const SCHEMA: u32 = 1;
60
61/// Wall-clock limit for one agent turn. See [`crate::config::Graph::timeout_talk`].
62///
63/// An hour by default. This turn is expected to run several shell commands
64/// and read their output before answering one - "what does this function
65/// do", "is this still true", "run the tests and tell me" - which argues for
66/// an hour rather than the five minutes a short budget once assumed, because
67/// the thing that made a short budget matter - an operator watching a
68/// spinner - is not how this conversation gets used: the operator moves on
69/// to something else while a turn runs and checks back later, so a long turn
70/// spends a held seat, not anyone's attention.
71fn turn_timeout(cfg: &Config) -> Duration {
72    Duration::from_secs(cfg.graph.timeout_talk)
73}
74
75/// Seat name for the conversation's agent, scoping its CLI-side session away
76/// from every other seat magi ever opens.
77const SEAT: &str = "talk";
78
79/// Prefix on a turn magi wrote rather than an agent.
80const MAGI_NOTE: &str = "magi: ";
81
82/// Who said something.
83#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
84#[serde(rename_all = "lowercase")]
85pub enum Who {
86    /// The operator.
87    Operator,
88    /// The conversation's agent - or magi itself, reporting that a turn
89    /// failed. See [`MAGI_NOTE`].
90    Agent,
91}
92
93/// One image the operator attached to a turn.
94///
95/// Never carries the bytes themselves: the picture lives on disk under
96/// [`Talks::attachments_dir`], named by `id` alone. `name` is the filename
97/// the operator's browser reported, kept only for display - it never
98/// contributes to a path, which is what keeps an upload from being able to
99/// traverse outside its own directory.
100#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
101#[serde(deny_unknown_fields)]
102pub struct Attachment {
103    /// Server-minted id; also the file's stem under `attachments_dir`.
104    pub id: String,
105    /// The operator's own filename, for display only.
106    pub name: String,
107    /// Validated by `web` at upload time against a closed whitelist:
108    /// `image/png`, `image/jpeg`, `image/gif`, `image/webp`.
109    pub mime: String,
110    /// Size in bytes, so the phone can show it without a second request.
111    pub bytes: u64,
112}
113
114/// One message in the conversation.
115#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
116#[serde(deny_unknown_fields)]
117pub struct Turn {
118    /// Who wrote it.
119    pub who: Who,
120    /// What they said.
121    pub body: String,
122    /// When it was said.
123    pub at: Timestamp,
124    /// Images attached to this turn. `#[serde(default)]` so a conversation
125    /// recorded before attachments existed still reads.
126    #[serde(default)]
127    pub attachments: Vec<Attachment>,
128    /// What the agent's CLI reported reading for this reply; only ever set on
129    /// an agent reply that carried usage. `#[serde(default)]` so a
130    /// conversation recorded before this field existed still reads, which is
131    /// why [`SCHEMA`] stays put: no existing field changed meaning.
132    #[serde(default, skip_serializing_if = "Option::is_none")]
133    pub usage: Option<TurnUsage>,
134}
135
136/// Raw usage of one agent reply, stored as the CLI reported it.
137///
138/// Counts only, never a percentage: the window is configuration
139/// ([`Config::context_window`]) and the conversation's model can change, so a
140/// stored percentage would go stale the moment either did. `agent` and
141/// `model` record who read that many tokens, which is how [`context_usage`]
142/// notices the figure predates a switch.
143#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
144pub struct TurnUsage {
145    /// Input-side tokens the CLI reported - see `agent::context_tokens`.
146    pub context_tokens: u64,
147    /// Roster id that answered.
148    pub agent: String,
149    /// That agent's model at the time (`None`: the CLI's own default).
150    #[serde(default)]
151    pub model: Option<String>,
152}
153
154/// How full the conversation's context window is, as the phone shows it.
155///
156/// Derived at read time and never persisted. `None` means unknown, which the
157/// UI says in words - it is never a made-up 0.
158#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
159pub struct ContextUsage {
160    /// Tokens the last agent reply's turn read, if its CLI reported any.
161    pub tokens: Option<u64>,
162    /// Window of the conversation's *current* model, if one is configured.
163    pub window: Option<u64>,
164    /// `tokens / window`, rounded down, not capped (over 100 is real news).
165    pub percent: Option<u64>,
166    /// 80% or more, judged before rounding: the conversation is getting long.
167    pub warn: bool,
168    /// The conversation's agent or model is not the one that produced
169    /// `tokens`: the figure describes the old session and the next turn
170    /// re-measures it.
171    pub since_switch: bool,
172    /// The current model, if the roster names one.
173    pub model: Option<String>,
174}
175
176/// Fraction (in percent) of the window at which [`ContextUsage::warn`] fires.
177const CONTEXT_WARN_PERCENT: u64 = 80;
178
179/// Context usage of `talk`, measured against its current model.
180///
181/// Only the latest agent reply counts (magi's own notes are skipped) and a
182/// reply without usage makes the answer unknown - older replies are never
183/// consulted, since a stale count passed off as current is worse than "unknown".
184/// The window comes from the *current* agent's model, so switching model moves
185/// the denominator at once.
186///
187/// Switching agent (or model, which is an agent change) mints a fresh CLI
188/// session, so the next turn re-sends the whole transcript and the count
189/// resets or jumps. Until that turn lands, the old figure is reported with
190/// `since_switch` set. Deterministic: same talk and config, same answer.
191pub fn context_usage(talk: &Talk, cfg: Option<&Config>) -> ContextUsage {
192    let current = cfg.and_then(|c| c.agents.iter().find(|a| a.id == talk.agent));
193    let model = current.and_then(|a| a.model.clone());
194    let window = cfg
195        .zip(model.as_deref())
196        .and_then(|(c, m)| c.context_window(m))
197        .filter(|w| *w > 0);
198    let usage = talk
199        .turns
200        .iter()
201        .rev()
202        .find(|t| t.who == Who::Agent && !t.body.starts_with(MAGI_NOTE))
203        .and_then(|t| t.usage.as_ref());
204    let tokens = usage.map(|u| u.context_tokens);
205    let since_switch =
206        usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
207    let (percent, warn) = match (tokens, window) {
208        (Some(t), Some(w)) => (
209            Some(t.saturating_mul(100) / w),
210            t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
211        ),
212        _ => (None, false),
213    };
214    ContextUsage {
215        tokens,
216        window,
217        percent,
218        warn,
219        since_switch,
220        model,
221    }
222}
223
224/// Where a conversation is in its life: this conversation can file any
225/// number of tasks without ending, so it only ever moves once, from open to
226/// closed.
227#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
228#[serde(rename_all = "lowercase")]
229pub enum TalkStatus {
230    /// Still open; the operator may say more, and may have already filed work
231    /// out of it.
232    Open,
233    /// Closed by hand. Kept on disk as a record.
234    Closed,
235}
236
237impl TalkStatus {
238    /// Is this conversation still live?
239    pub fn open(self) -> bool {
240        matches!(self, Self::Open)
241    }
242
243    /// Wire form, for the phone and for logs.
244    pub fn as_str(self) -> &'static str {
245        match self {
246            Self::Open => "open",
247            Self::Closed => "closed",
248        }
249    }
250}
251
252/// One standing conversation.
253#[derive(Debug, Clone, Serialize, Deserialize)]
254#[serde(deny_unknown_fields)]
255pub struct Talk {
256    /// On-disk format version.
257    pub schema: u32,
258    /// Conversation id, e.g. `20260904-014455-ab12`.
259    pub id: String,
260    /// Repository this conversation is about.
261    pub repo: PathBuf,
262    /// Roster agent id holding the conversation.
263    pub agent: String,
264    /// Current state.
265    pub status: TalkStatus,
266    /// Everything said, oldest first.
267    pub turns: Vec<Turn>,
268    /// Text and attachments accepted while the single CLI turn is busy.
269    /// They are durable, but become a real turn only when [`drain`] records them.
270    #[serde(default)]
271    pub pending: String,
272    /// Attachments paired with [`Self::pending`].
273    #[serde(default)]
274    pub pending_attachments: Vec<Attachment>,
275    /// May a failed turn fall back through the rest of `[roles] chatter`?
276    /// True only while the agent was chosen by that chain; an explicit
277    /// `--agent` or an operator's switch pins the conversation to its agent
278    /// even when that agent also appears in the chain.
279    #[serde(default)]
280    pub fallback: bool,
281    /// When the conversation was opened.
282    pub created_at: Timestamp,
283    /// Last change to this file.
284    pub updated_at: Timestamp,
285    /// The CLI-side conversation, so a turn after the first costs one
286    /// sentence instead of the whole transcript. Not `pub`: it is magi's
287    /// bookkeeping, and a caller that edited it would detach the record from
288    /// the conversation the model actually holds.
289    seat: SeatState,
290}
291
292impl Talk {
293    /// Short form used in lists and notifications, matching a run's short id.
294    pub fn short(&self) -> &str {
295        short(&self.id)
296    }
297}
298
299/// A conversation store on disk.
300#[derive(Debug, Clone)]
301pub struct Talks {
302    root: PathBuf,
303    /// Serializes the read-modify-write cycle that reads a talk, decides
304    /// something from its `status`, and writes the whole record back.
305    /// [`close`], [`record`] and the tail of [`turn`] all take this before
306    /// that cycle rather than after just the read: a re-read narrows the
307    /// window another writer can land in, but does not close it, since
308    /// nothing stopped that other writer's own put from landing between this
309    /// call's re-read and its own put. Shared across every clone, since every
310    /// clone is a handle onto the same files.
311    lock: Arc<Mutex<()>>,
312}
313
314impl Talks {
315    /// The operator's conversations, `<home>/talks`.
316    pub fn open() -> Self {
317        Self::at(crate::run::home().join("talks"))
318    }
319
320    /// A store at an explicit root. Tests use this, which is why none of them
321    /// need the operator's real home.
322    pub fn at(root: PathBuf) -> Self {
323        Self {
324            root,
325            lock: Arc::new(Mutex::new(())),
326        }
327    }
328
329    /// Claim the right to read-modify-write a talk's `status`. A plain
330    /// `std::sync::Mutex`, not an async one: every caller holds it across a
331    /// handful of small file operations and never across an `.await`, so
332    /// blocking the thread briefly is the right tool, not a reason to reach
333    /// for `tokio::sync::Mutex`. Poisoning recovers rather than propagates -
334    /// one panicking caller must not wedge every talk in the store the way it
335    /// would wedge the loop's own lock; see [`crate::web`]'s `lock_or_recover`,
336    /// which this mirrors.
337    fn guard(&self) -> MutexGuard<'_, ()> {
338        self.lock.lock().unwrap_or_else(PoisonError::into_inner)
339    }
340
341    /// Directory holding the conversation files.
342    pub fn root(&self) -> &Path {
343        &self.root
344    }
345
346    /// Path for one conversation id.
347    pub fn path_of(&self, id: &str) -> PathBuf {
348        self.root.join(format!("{id}.json"))
349    }
350
351    /// Where one conversation's prompts and CLI output are kept, beside the
352    /// record rather than inside it.
353    pub fn artifacts_of(&self, id: &str) -> PathBuf {
354        self.root.join(format!("{id}.artifacts"))
355    }
356
357    /// Where this conversation's attached images live: a subdirectory of
358    /// `artifacts_of`, so deleting the conversation deletes its attachments
359    /// too and nothing here needs its own cleanup path.
360    pub fn attachments_dir(&self, id: &str) -> PathBuf {
361        self.artifacts_of(id).join("attachments")
362    }
363
364    /// Persist one already-validated attachment and return its metadata.
365    ///
366    /// `web::talk_attachment_post` is the only caller: it has already
367    /// checked `mime` against the whitelist and sniffed the bytes, so an
368    /// unrecognised mime reaching here is a bug in that caller, not
369    /// something an operator did. The id is minted here and never taken
370    /// from the client; `name` is stored for display only and never used to
371    /// build a path.
372    pub fn put_attachment(
373        &self,
374        id: &str,
375        mime: &str,
376        name: &str,
377        data: &[u8],
378    ) -> Result<Attachment> {
379        let dir = self.attachments_dir(id);
380        std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
381        let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
382        let att = Attachment {
383            id: new_attachment_id(),
384            name: name.to_owned(),
385            mime: mime.to_owned(),
386            bytes: data.len() as u64,
387        };
388        std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
389            .with_context(|| format!("write attachment {}", att.id))?;
390        std::fs::write(
391            dir.join(format!("{}.json", att.id)),
392            serde_json::to_string(&att).context("serialize attachment")?,
393        )
394        .with_context(|| format!("write attachment metadata {}", att.id))?;
395        Ok(att)
396    }
397
398    /// Just the metadata, without reading the image bytes back off disk -
399    /// what `web::talk_say` uses to turn an id the operator referenced into
400    /// an [`Attachment`] before appending a [`Turn`], where the bytes
401    /// themselves are of no interest. `None` for an id this conversation
402    /// never stored - including one that merely looks plausible:
403    /// [`valid_attachment_id`] is checked here too, not only by the caller,
404    /// the same defence-in-depth `Questions::panel_asset` uses for its own
405    /// asset ids.
406    pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
407        if !valid_attachment_id(att_id) {
408            return Ok(None);
409        }
410        let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
411        if !meta_path.is_file() {
412            return Ok(None);
413        }
414        let att = serde_json::from_str(
415            &std::fs::read_to_string(&meta_path)
416                .with_context(|| format!("read {}", meta_path.display()))?,
417        )
418        .with_context(|| format!("parse {}", meta_path.display()))?;
419        Ok(Some(att))
420    }
421
422    /// A stored attachment's metadata and its bytes together, for serving it
423    /// back on `GET`. `None` under the same conditions as
424    /// [`Talks::attachment_meta`], which this is built on.
425    pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
426        let Some(att) = self.attachment_meta(id, att_id)? else {
427            return Ok(None);
428        };
429        let ext = attachment_ext(&att.mime).with_context(|| {
430            format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
431        })?;
432        let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
433        let data =
434            std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
435        Ok(Some((att, data)))
436    }
437
438    /// Absolute path of one attachment's bytes, for the prompt note [`turn`]
439    /// appends and for [`Invocation::attachments`]. `None` only for a mime
440    /// [`put_attachment`] could never have written, which means the
441    /// attachment did not come from this store.
442    ///
443    /// `self.root` (and so `attachments_dir`) is not guaranteed absolute on
444    /// its own - `run::home()` returns a bare relative `PathBuf` verbatim
445    /// when the operator sets `MAGI_HOME` to a relative path, and nothing
446    /// canonicalizes it on the way in. That is harmless for every other use
447    /// of this store, since its own I/O runs in this process against this
448    /// process's cwd - but this path is handed to a CLI invoked with `cwd:
449    /// &talk.repo`, a different directory, so a relative path here would
450    /// resolve against the wrong place once it reached the prompt.
451    /// `std::path::absolute` fixes it against *this* process's cwd before
452    /// that happens; see `disk::free_bytes_by_os` for the same function used
453    /// the same way elsewhere in this codebase.
454    fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
455        let ext = attachment_ext(&att.mime)?;
456        let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
457        std::path::absolute(&path).ok()
458    }
459
460    /// Write a conversation, atomically, so a process killed mid-write leaves
461    /// the previous state readable rather than a truncated file.
462    ///
463    /// The write-then-rename itself is retried a handful of times - see
464    /// [`write_atomic`] - because a reader with the destination file briefly
465    /// open is exactly the kind of failure that must not cost an agent's
466    /// whole reply; see `turn`'s own tail for what happens when even that is
467    /// not enough.
468    pub fn put(&self, t: &mut Talk) -> Result<()> {
469        std::fs::create_dir_all(&self.root)
470            .with_context(|| format!("create {}", self.root.display()))?;
471        t.updated_at = Timestamp::now();
472        let body = serde_json::to_string_pretty(t).context("serialize talk")?;
473        let path = self.path_of(&t.id);
474        let tmp = path.with_extension("json.tmp");
475        write_atomic(&tmp, &path, &body)
476    }
477
478    /// Load a conversation by id or unambiguous id prefix.
479    pub fn get(&self, id: &str) -> Result<Talk> {
480        let resolved = self.resolve_id(id)?;
481        read_path(&self.path_of(&resolved))
482    }
483
484    /// Every conversation on disk: open first, then newest first, so what the
485    /// operator is still using belongs above what they are done with.
486    pub fn list(&self) -> Vec<Talk> {
487        let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
488            .into_iter()
489            .flatten()
490            .flatten()
491            .map(|e| e.path())
492            .filter(|p| p.extension().is_some_and(|x| x == "json"))
493            .filter_map(|p| read_path(&p).ok())
494            .collect();
495        all.sort_unstable_by(|a, b| {
496            let rank = |t: &Talk| u8::from(!t.status.open());
497            rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
498        });
499        all
500    }
501
502    /// Expand an id prefix to exactly one conversation id.
503    pub fn resolve_id(&self, prefix: &str) -> Result<String> {
504        if self.path_of(prefix).is_file() {
505            return Ok(prefix.to_owned());
506        }
507        let hits: Vec<String> = self
508            .list()
509            .into_iter()
510            .map(|t| t.id)
511            .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
512            .collect();
513        match hits.len() {
514            1 => Ok(hits.into_iter().next().expect("exactly one hit")),
515            0 => bail!("no talk matches `{prefix}`"),
516            _ => bail!(
517                "`{prefix}` matches {} talks: {}",
518                hits.len(),
519                hits.join(", ")
520            ),
521        }
522    }
523
524    /// Change detection token: the newest modification time in the store, in
525    /// milliseconds.
526    pub fn revision(&self) -> u64 {
527        std::fs::read_dir(&self.root)
528            .into_iter()
529            .flatten()
530            .flatten()
531            .filter_map(|e| e.metadata().ok())
532            .filter_map(|m| m.modified().ok())
533            .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
534            .map(|d| d.as_millis() as u64)
535            .max()
536            .unwrap_or(0)
537    }
538
539    /// How many conversations are still open.
540    pub fn count_open(&self) -> usize {
541        self.list().iter().filter(|t| t.status.open()).count()
542    }
543
544    /// Remove a conversation from disk, record and artifacts both. The
545    /// operator's way of saying "not just done, gone" - [`close`] alone
546    /// leaves the record as history.
547    ///
548    /// Takes [`Talks::guard`] for the same reason [`close`] does: a delete
549    /// racing a [`record`] or the tail of [`turn`] must not land between
550    /// their own read and write, or the file removed here would look, to
551    /// them, like a record that simply has not been written yet. The other
552    /// half of that story is on their side - both check under this same
553    /// guard that the record they are about to write is still there, and
554    /// give up without writing if it is not, which is what stops their `put`
555    /// from resurrecting a conversation this call already removed.
556    pub fn remove(&self, id: &str) -> Result<()> {
557        let _guard = self.guard();
558        let resolved = self.resolve_id(id)?;
559        let path = self.path_of(&resolved);
560        std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
561        let artifacts = self.artifacts_of(&resolved);
562        if artifacts.is_dir() {
563            std::fs::remove_dir_all(&artifacts)
564                .with_context(|| format!("remove {}", artifacts.display()))?;
565        }
566        Ok(())
567    }
568}
569
570/// Open a conversation. Takes no agent turn: there is no idea to answer yet,
571/// and a conversation the operator has not said anything into yet is a
572/// normal, valid thing to have sitting on the phone.
573///
574/// `agent` beats `[roles] chatter`, which beats [`agent::pick`]'s own default
575/// order (a claude seat, else the first runnable agent in roster order) when
576/// nothing names a seat at all - see `[roles] chatter`'s own doc in
577/// [`crate::config`] for why a dedicated field exists rather than reusing a
578/// judge seat.
579pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
580    // Absolute: a relative path means the wrong repository once anything
581    // other than this process reads it back.
582    let repo = repo.canonicalize().unwrap_or(repo);
583    // An explicit agent is a chain of one; otherwise the first id of
584    // `[roles] chatter` that can run here (later ones are `turn`'s fallbacks).
585    let spec = match agent {
586        Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
587        None => agent::pick_chain(
588            &cfg.agents,
589            cfg.roles.chatter.as_ref(),
590            &agent::installed,
591            "chatter",
592        )?
593        .remove(0),
594    };
595
596    let now = Timestamp::now();
597    let mut talk = Talk {
598        schema: SCHEMA,
599        id: new_id(),
600        repo,
601        agent: spec.id.clone(),
602        status: TalkStatus::Open,
603        turns: Vec::new(),
604        pending: String::new(),
605        pending_attachments: Vec::new(),
606        fallback: agent.is_none(),
607        created_at: now,
608        updated_at: now,
609        seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
610    };
611    store.put(&mut talk)?;
612    Ok(talk)
613}
614
615/// Append the operator's turn and flush it, without invoking anything.
616///
617/// Split out of [`say`] so `POST /api/talks/{id}/say` can answer once the
618/// message is safely on disk, and run the agent's half in the background -
619/// holding the connection for a turn that can run fifteen minutes is the
620/// wrong shape for a phone.
621pub fn record(
622    talk: &mut Talk,
623    store: &Talks,
624    text: &str,
625    attachments: Vec<Attachment>,
626) -> Result<String> {
627    // `web::talk_say` reads the talk, then awaits config discovery before
628    // calling this - a gap a concurrent `POST /api/talks/{id}/close` can land
629    // in. The guard held for the rest of this function is what actually closes
630    // that gap: re-reading status without it only shrinks the window a
631    // concurrent `close` could land in between this call's own read and its
632    // `put`, it does not remove it. See [`Talks::guard`] and the matching
633    // guard in `turn`, which this mirrors.
634    let _guard = store.guard();
635    // A concurrent `Talks::remove` can have landed in that same gap. `put`
636    // writes unconditionally, so trusting the stale `talk` here would recreate
637    // the file a delete just removed - the record must still be there for a
638    // turn to have anywhere to append to.
639    let Ok(fresh) = store.get(&talk.id) else {
640        bail!("talk {} was deleted", talk.short());
641    };
642    talk.status = fresh.status;
643    // Do not let this older handle overwrite a draft accepted while it was
644    // waiting for configuration discovery.
645    talk.pending = fresh.pending;
646    talk.pending_attachments = fresh.pending_attachments;
647    if !talk.status.open() {
648        bail!(
649            "talk {} is {} and takes no more turns",
650            talk.short(),
651            talk.status.as_str()
652        );
653    }
654    let text = text.trim();
655    if text.is_empty() && attachments.is_empty() {
656        bail!("nothing to say");
657    }
658    talk.turns.push(Turn {
659        who: Who::Operator,
660        body: text.to_owned(),
661        at: Timestamp::now(),
662        attachments,
663        usage: None,
664    });
665    store.put(talk)?;
666    Ok(text.to_owned())
667}
668
669/// Add an unrecorded message to the durable draft while another turn runs.
670pub fn queue(
671    talk: &mut Talk,
672    store: &Talks,
673    text: &str,
674    attachments: Vec<Attachment>,
675) -> Result<()> {
676    let text = text.trim();
677    if text.is_empty() && attachments.is_empty() {
678        bail!("nothing to say");
679    }
680    let _guard = store.guard();
681    let mut fresh = store
682        .get(&talk.id)
683        .with_context(|| format!("talk {} was deleted", talk.short()))?;
684    if !fresh.status.open() {
685        bail!(
686            "talk {} is {} and takes no more turns",
687            fresh.short(),
688            fresh.status.as_str()
689        );
690    }
691    if !text.is_empty() {
692        if fresh.pending.is_empty() {
693            fresh.pending = text.to_owned();
694        } else {
695            fresh.pending.push_str("\n\n");
696            fresh.pending.push_str(text);
697        }
698    }
699    fresh.pending_attachments.extend(attachments);
700    store.put(&mut fresh)?;
701    *talk = fresh;
702    Ok(())
703}
704
705/// Promote the current durable draft to one operator turn.
706pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
707    let _guard = store.guard();
708    let mut fresh = store
709        .get(&talk.id)
710        .with_context(|| format!("talk {} was deleted", talk.short()))?;
711    if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
712        *talk = fresh;
713        return Ok(None);
714    }
715    let text = std::mem::take(&mut fresh.pending);
716    let attachments = std::mem::take(&mut fresh.pending_attachments);
717    fresh.turns.push(Turn {
718        who: Who::Operator,
719        body: text.clone(),
720        at: Timestamp::now(),
721        attachments,
722        usage: None,
723    });
724    store.put(&mut fresh)?;
725    *talk = fresh;
726    Ok(Some(text))
727}
728
729/// One operator turn and one agent turn, appended - the synchronous form, used
730/// by tests and by anything that is fine waiting out the turn itself.
731pub async fn say(
732    talk: &mut Talk,
733    store: &Talks,
734    cfg: &Config,
735    text: &str,
736    attachments: Vec<Attachment>,
737) -> Result<()> {
738    let text = record(talk, store, text, attachments)?;
739    turn(talk, store, cfg, &text).await
740}
741
742/// The agent's half of a turn: invoke, append, flush. Pairs with [`record`].
743pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
744    turn(talk, store, cfg, text).await
745}
746
747/// Close a conversation. Idempotent: closing an already-closed conversation is
748/// not an error, since the operator's intent - "I am done with this" - is
749/// already satisfied.
750///
751/// Re-reads the record under [`Talks::guard`] rather than trusting the
752/// caller's copy of `talk`, and writes that fresh copy back rather than the
753/// one passed in. `web::talk_close` loads `talk` and calls this right after
754/// with no gap of its own, but without the guard that load can still land
755/// between a `record` or `turn` elsewhere reading the file and writing it
756/// back - and a close built on the older snapshot would put it right back,
757/// silently dropping whatever turn the other call had just appended.
758///
759/// If the re-read fails, this errors rather than falling back to the
760/// caller's stale copy: `talk::begin` always `put`s the record before handing
761/// out a `Talk`, so the only way a re-read can fail is a concurrent
762/// [`Talks::remove`] having deleted it, and writing the stale copy back would
763/// resurrect exactly what that delete removed.
764pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
765    let _guard = store.guard();
766    let mut fresh = store
767        .get(&talk.id)
768        .with_context(|| format!("talk {} was deleted", talk.short()))?;
769    fresh.status = TalkStatus::Closed;
770    // A closed conversation must not replay a draft if it is reopened later.
771    fresh.pending.clear();
772    fresh.pending_attachments.clear();
773    store.put(&mut fresh)?;
774    *talk = fresh;
775    Ok(())
776}
777
778/// Reopen a closed conversation. Idempotent for the same reason [`close`] is:
779/// reopening an already-open conversation is not an error, since the
780/// operator's intent - "I want to keep talking about this" - is already
781/// satisfied.
782///
783/// Written symmetrically with [`close`]: re-reads the record under
784/// [`Talks::guard`] rather than trusting the caller's copy of `talk`, writes
785/// that fresh copy back rather than the one passed in, and errors rather than
786/// falling back to the stale copy if the re-read fails, for the same reasons
787/// `close`'s doc gives.
788pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
789    let _guard = store.guard();
790    let mut fresh = store
791        .get(&talk.id)
792        .with_context(|| format!("talk {} was deleted", talk.short()))?;
793    fresh.status = TalkStatus::Open;
794    store.put(&mut fresh)?;
795    *talk = fresh;
796    Ok(())
797}
798
799/// Hand the conversation to another roster agent.
800///
801/// A CLI session belongs to one CLI and cannot be carried to another, so the
802/// seat is minted afresh rather than edited: the next turn finds
803/// `seat.turns == 0` and re-sends the transcript, since the new agent has
804/// heard none of it. A magi-written note records the change in the
805/// transcript. Returns `false` (and writes nothing, not even a note) when the
806/// stored talk already uses `spec`.
807///
808/// Re-reads under [`Talks::guard`], like [`close`], and errors rather than
809/// resurrecting a record a concurrent delete removed. Refusing a closed talk
810/// or a turn in flight is the caller's job: only it can see the latter.
811pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
812    let _guard = store.guard();
813    let mut fresh = store
814        .get(&talk.id)
815        .with_context(|| format!("talk {} was deleted", talk.short()))?;
816    if fresh.agent == spec.id {
817        *talk = fresh;
818        return Ok(false);
819    }
820    let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
821    fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
822    // A deliberate switch pins the conversation to the agent chosen.
823    fresh.fallback = false;
824    fresh.turns.push(Turn {
825        who: Who::Agent,
826        body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
827        at: Timestamp::now(),
828        attachments: Vec::new(),
829        usage: None,
830    });
831    store.put(&mut fresh)?;
832    *talk = fresh;
833    Ok(true)
834}
835
836/// Discard the durable draft without adding a transcript turn.
837pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
838    let _guard = store.guard();
839    let mut fresh = store
840        .get(&talk.id)
841        .with_context(|| format!("talk {} was deleted", talk.short()))?;
842    fresh.pending.clear();
843    fresh.pending_attachments.clear();
844    store.put(&mut fresh)?;
845    *talk = fresh;
846    Ok(())
847}
848
849/// Clear a draft only when the caller still sees its complete snapshot.
850pub fn clear_pending_if_matches(
851    talk: &mut Talk,
852    store: &Talks,
853    expected_text: &str,
854    expected_attachments: &[String],
855) -> Result<bool> {
856    let _guard = store.guard();
857    let mut fresh = store
858        .get(&talk.id)
859        .with_context(|| format!("talk {} was deleted", talk.short()))?;
860    if !pending_matches(&fresh, expected_text, expected_attachments) {
861        *talk = fresh;
862        return Ok(false);
863    }
864    fresh.pending.clear();
865    fresh.pending_attachments.clear();
866    store.put(&mut fresh)?;
867    *talk = fresh;
868    Ok(true)
869}
870
871/// Replace just the text of the durable draft, but only if the caller's
872/// snapshot still identifies the entire draft. This refuses to overwrite a
873/// message another client queued or a draft the drain already promoted.
874pub fn edit_pending_text(
875    talk: &mut Talk,
876    store: &Talks,
877    text: &str,
878    expected_text: &str,
879    expected_attachments: &[String],
880) -> Result<bool> {
881    let _guard = store.guard();
882    let mut fresh = store
883        .get(&talk.id)
884        .with_context(|| format!("talk {} was deleted", talk.short()))?;
885    if !pending_matches(&fresh, expected_text, expected_attachments) {
886        *talk = fresh;
887        return Ok(false);
888    }
889    fresh.pending = text.trim().to_owned();
890    store.put(&mut fresh)?;
891    *talk = fresh;
892    Ok(true)
893}
894
895fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
896    talk.pending == expected_text
897        && talk
898            .pending_attachments
899            .iter()
900            .map(|attachment| &attachment.id)
901            .eq(expected_attachments.iter())
902}
903
904/// Invoke the conversation's agent once and append what it said.
905///
906/// The first turn ever taken carries the full [`briefing`], because nothing
907/// else has told the agent what this conversation is or what it may do.
908/// Every turn after that resends nothing when the CLI can resume its own
909/// session, and falls back to [`transcript`] only when it cannot.
910async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
911    let spec = cfg
912        .agents
913        .iter()
914        .find(|a| a.id == talk.agent)
915        .with_context(|| {
916            format!(
917                "talk {} was opened with agent `{}`, which is no longer in \
918                 the roster; restore it in magi.toml or start a new \
919                 conversation",
920                talk.short(),
921                talk.agent
922            )
923        })?;
924
925    // The newest turn is always the operator message this call is answering
926    // - `record` appended it before `turn` was ever called - so its own
927    // attachments are what belong at the end of *this* prompt.
928    let last_note = attachment_note(
929        store,
930        &talk.id,
931        talk.turns
932            .last()
933            .map_or(&[][..], |t| t.attachments.as_slice()),
934    );
935
936    // Every attachment this conversation has ever held, not only this
937    // turn's: a resumed session gets a fresh process every turn, so a CLI
938    // whose sandbox needs `--add-dir` (see `agent::build_command`) needs the
939    // grant again to open an image from an earlier turn, even when nothing
940    // new was attached just now.
941    let attachment_paths: Vec<PathBuf> = talk
942        .turns
943        .iter()
944        .flat_map(|t| t.attachments.iter())
945        .filter_map(|a| store.attachment_path(&talk.id, a))
946        .collect();
947
948    let artifacts = store.artifacts_of(&talk.id);
949    // From the transcript, not the seat: a switched agent's seat restarts at
950    // zero and must not overwrite an earlier turn's artifacts.
951    let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
952    let stem = format!("turn-{}", operator_turns.max(1));
953    // The chat's build cache is the same shared one the graph's seats get, so
954    // a conversation that compiles does not mint another multi-GB target dir.
955    let cache_dir = cfg.cache_dir();
956
957    // The agent holding the conversation, then - only when it came from
958    // `[roles] chatter` - the rest of that chain, each at most once.
959    let mut chain = vec![spec.clone()];
960    if let Some(choice) = cfg.roles.chatter.as_ref()
961        && talk.fallback
962    {
963        for id in choice.ids() {
964            if id == talk.agent || chain.iter().any(|s| s.id == id) {
965                continue;
966            }
967            match agent::pick(&cfg.agents, Some(id), &agent::installed) {
968                Ok(s) => chain.push(s),
969                Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
970            }
971        }
972    }
973
974    let mut outcome = None;
975    let mut fell_back_from: Option<String> = None;
976    // What the conversation looked like after the first agent's failed try,
977    // so an exhausted chain leaves exactly what a single failed seat would.
978    let mut first_try: Option<(String, SeatState)> = None;
979    for (n, spec) in chain.iter().enumerate() {
980        if n > 0 {
981            if first_try.is_none() {
982                first_try = Some((talk.agent.clone(), talk.seat.clone()));
983            }
984            tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
985            // A new CLI has none of the old one's conversation: a fresh seat
986            // puts `has_session` at false and the full transcript is re-sent.
987            fell_back_from.get_or_insert_with(|| talk.agent.clone());
988            talk.agent = spec.id.clone();
989            talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
990        }
991        let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
992        let first_ever = talk.turns.len() <= 1;
993        let body = if talk.seat.turns == 0 && first_ever {
994            format!(
995                "{}\n\n# Operator\n\n{text}{last_note}",
996                briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write)
997            )
998        } else if talk.seat.turns == 0 {
999            // A fresh seat on a conversation that already has history (the
1000            // agent was switched): the briefing, then everything said so far.
1001            format!(
1002                "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1003                briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write),
1004                transcript(talk, store)
1005            )
1006        } else if resuming {
1007            format!("{text}{last_note}")
1008        } else {
1009            format!("{}\n\n{text}{last_note}", transcript(talk, store))
1010        };
1011        let attempt_stem = if n == 0 {
1012            stem.clone()
1013        } else {
1014            format!("{stem}-{}", spec.id)
1015        };
1016        let inv = Invocation {
1017            cwd: &talk.repo,
1018            prompt: &body,
1019            timeout: turn_timeout(cfg),
1020            // Off unless this repository's own config opts in - see
1021            // `crate::config::Talk::allow_write` and this module's doc for why
1022            // the default keeps a conversational edit from landing in a checkout
1023            // no run or review can claim.
1024            allow_write: cfg.talk.allow_write,
1025            sessions: cfg.graph.sessions,
1026            artifacts: &artifacts,
1027            stem: &attempt_stem,
1028            // The conversation's own id, so `magi task add` run from inside it is
1029            // attributed to this conversation - see `Source::Agent`.
1030            run: &talk.id,
1031            node: "chat",
1032            cache_dir: cache_dir.as_deref(),
1033            attachments: &attachment_paths,
1034            writable: &[],
1035        };
1036        let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1037        let advance = agent::chain_advances(&result);
1038        if n == 0 || !advance {
1039            outcome = Some(result);
1040        } else {
1041            // A later failure is only logged; the note describes the first.
1042            tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1043        }
1044        if !advance {
1045            break;
1046        }
1047    }
1048    if outcome.as_ref().is_some_and(agent::chain_advances) {
1049        // Exhausted: back to the agent the conversation had, so the note
1050        // below names it and the next turn starts from it again.
1051        if let Some((id, seat)) = first_try {
1052            talk.agent = id;
1053            talk.seat = seat;
1054            fell_back_from = None;
1055        }
1056    }
1057    let outcome = outcome.expect("a chain holds at least one agent");
1058    let note = |why: String| Turn {
1059        who: Who::Agent,
1060        body: format!("{MAGI_NOTE}{why}"),
1061        at: Timestamp::now(),
1062        attachments: Vec::new(),
1063        usage: None,
1064    };
1065    let (reply, failure) = match outcome {
1066        Err(e) => (
1067            note(format!("could not run agent `{}`: {e}", talk.agent)),
1068            Some(format!("could not run agent `{}`: {e}", talk.agent)),
1069        ),
1070        Ok(out) if out.quota_exhausted() => {
1071            let reset = out
1072                .quota
1073                .as_ref()
1074                .and_then(|q| q.reset.clone())
1075                .map_or_else(String::new, |r| format!(" (resets {r})"));
1076            let why = format!(
1077                "agent `{}` is out of quota{reset}; your message is saved, so \
1078                 say it again when the window reopens",
1079                talk.agent
1080            );
1081            (note(why.clone()), Some(why))
1082        }
1083        Ok(out) if out.timed_out => {
1084            let why = format!(
1085                "agent `{}` did not answer within {}s; your message is saved",
1086                talk.agent,
1087                turn_timeout(cfg).as_secs()
1088            );
1089            (note(why.clone()), Some(why))
1090        }
1091        Ok(out) if !out.usable() => {
1092            let why = format!(
1093                "agent `{}` produced no answer (exit {}); your message is saved",
1094                talk.agent,
1095                out.exit_code
1096                    .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1097            );
1098            (note(why.clone()), Some(why))
1099        }
1100        Ok(out) => (
1101            Turn {
1102                who: Who::Agent,
1103                body: out.text.trim().to_owned(),
1104                at: Timestamp::now(),
1105                attachments: Vec::new(),
1106                // `talk.agent` is whoever actually answered: a fallback has
1107                // already moved it, and an exhausted chain never reaches here.
1108                usage: out.context_tokens.map(|context_tokens| TurnUsage {
1109                    context_tokens,
1110                    agent: talk.agent.clone(),
1111                    model: cfg
1112                        .agents
1113                        .iter()
1114                        .find(|a| a.id == talk.agent)
1115                        .and_then(|a| a.model.clone()),
1116                }),
1117            },
1118            None,
1119        ),
1120    };
1121
1122    // A close landed on disk while this turn was in flight is read back here
1123    // rather than trusted from the snapshot this call started with. `store`
1124    // holds nothing else this function does not itself own - the turn guard
1125    // in `web::Ui::begin_talk_turn` keeps `turns` and `seat` this call's
1126    // alone to mutate - but `status` is not behind that guard, and an
1127    // operator's close must stick: the whole point of ending a conversation
1128    // is that an agent's answer to the last message before the close cannot
1129    // silently reopen it. The guard is what makes that read-then-write
1130    // section atomic with `close`'s own - taken only for this tail and not
1131    // for the whole invocation above, so one talk's fifteen-minute turn does
1132    // not block another talk's close from proceeding.
1133    let _guard = store.guard();
1134    // A delete is the more final version of that same race: `put` writes
1135    // unconditionally, so a talk removed while this turn was in flight must
1136    // stay removed rather than being written back with this turn's reply
1137    // appended to it. The reply is simply given up on - there is no
1138    // conversation left for it to belong to.
1139    let Ok(fresh) = store.get(&talk.id) else {
1140        return Ok(());
1141    };
1142    talk.status = fresh.status;
1143    // `queue` may have accepted another operator message while the CLI was
1144    // running. This handle predates that write, so preserving only `status`
1145    // would overwrite the durable draft when the reply is appended below.
1146    talk.pending = fresh.pending;
1147    talk.pending_attachments = fresh.pending_attachments;
1148    if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1149        // The switch persists: quota coming back does not move the chat
1150        // home, an operator's switch does.
1151        talk.turns.push(note(format!(
1152            "agent changed from {from} to {} (fallback)",
1153            talk.agent
1154        )));
1155    }
1156    talk.turns.push(reply);
1157    if let Err(put_err) = store.put(talk) {
1158        // `Talks::put` already retried the write itself - reaching here
1159        // means a passing race is not what this is. An agent's answer,
1160        // possibly the result of an hour-long call, must not vanish with
1161        // nothing to show for it just because the very last step failed:
1162        // pop it back off, stash its text beside the conversation, and
1163        // replace it with a note the operator can actually see, the same
1164        // mechanism the failure branches above already use for a quota or a
1165        // timeout.
1166        let lost = talk.turns.pop().expect("just pushed above");
1167        let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1168        let why = match &stash {
1169            Ok(path) => format!(
1170                "agent `{}` answered, but the reply could not be saved to \
1171                 this conversation ({put_err:#}); the raw text was kept at \
1172                 {} - your message is saved, ask again",
1173                talk.agent,
1174                path.display()
1175            ),
1176            Err(stash_err) => format!(
1177                "agent `{}` answered, but the reply could not be saved to \
1178                 this conversation ({put_err:#}), and it could not be kept \
1179                 anywhere else either ({stash_err:#}); your message is \
1180                 saved, ask again",
1181                talk.agent
1182            ),
1183        };
1184        talk.turns.push(note(why.clone()));
1185        // Writing the note also carries the seat this call already advanced -
1186        // `agent::invoke` incremented `turns` and, for a vendor that reports
1187        // its own session id, recorded that too. That is what keeps the next
1188        // turn resuming the session the CLI is already holding instead of
1189        // re-opening it, so losing the reply costs the transcript a turn but
1190        // not the conversation.
1191        return match store.put(talk) {
1192            Ok(()) => bail!("{why}"),
1193            Err(note_err) => {
1194                // Even the short note failed to save, which means this
1195                // conversation's file cannot be written at all right now -
1196                // nothing is left for this call to retry or record. Pop the
1197                // note so `talk.turns` matches the transcript on disk, and
1198                // surface both failures for whoever reads the log.
1199                //
1200                // `talk.seat` is deliberately not wound back to match. The
1201                // CLI really did take the turn and really did consume this
1202                // seat's session id; pretending otherwise would be a second
1203                // untruth on top of the unwritable file, and the handle is
1204                // reloaded from disk by the next `drain` or `get` anyway -
1205                // see `web::drain_loop`. What the seat cannot do is reach
1206                // disk, so the record stays a turn behind the CLI until some
1207                // later write lands, and a turn taken before then re-opens a
1208                // session id the CLI already holds. That is the desync
1209                // `20260907-011805-fb57` is about, and tolerating it belongs
1210                // there rather than here: no write this branch could make
1211                // would help, since a failed write is exactly what put it in
1212                // this position twice over.
1213                talk.turns.pop();
1214                Err(note_err).context(why)
1215            }
1216        };
1217    }
1218
1219    match failure {
1220        Some(why) => bail!("{why}"),
1221        None => Ok(()),
1222    }
1223}
1224
1225/// Everything said so far, as prose, for a CLI that cannot resume its own
1226/// conversation.
1227fn transcript(talk: &Talk, store: &Talks) -> String {
1228    let mut out = String::from(
1229        "This conversation cannot resume on the CLI's side, so here is \
1230         everything said so far; answer only the last message.\n",
1231    );
1232    for t in &talk.turns {
1233        let who = match t.who {
1234            Who::Operator => "operator",
1235            Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1236            Who::Agent => "you",
1237        };
1238        out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1239        out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1240    }
1241    out
1242}
1243
1244/// The section named at the end of a turn's body, listing every attachment's
1245/// absolute path and mime so the agent knows exactly what to open. Empty
1246/// when `attachments` is, which is every turn but the rare one carrying an
1247/// image, so a turn with none changes nothing about the prompt.
1248fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1249    if attachments.is_empty() {
1250        return String::new();
1251    }
1252    let mut out = String::from(
1253        "\n\nThe operator attached the image(s) below to this message. Open \
1254         and look at each one before you answer.\n",
1255    );
1256    for att in attachments {
1257        if let Some(path) = store.attachment_path(talk_id, att) {
1258            out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1259        }
1260    }
1261    out.push('\n');
1262    out
1263}
1264
1265/// The briefing the agent opens with, sent once as part of its first turn.
1266///
1267/// Pure, so the properties that matter can be asserted without an interview:
1268/// it names `magi task add --solo` (the route this conversation always has to
1269/// changing anything) and it never tells the agent to write a task *file* of
1270/// its own - that would compete with filing through the queue.
1271/// `allow_write` only ever adds an extra permission on top of that; it never
1272/// removes the queue as an option, which is why both branches keep the same
1273/// `# When the operator wants something done` section - `write_policy` is
1274/// the only part that changes.
1275///
1276/// It also tells the agent that `--repo` is not stuck naming this
1277/// conversation's own directory: `resolve_repo` (`src/main.rs`) now accepts a
1278/// short `owner/repo` or bare `repo` name and resolves it against
1279/// `[repos] roots`, the same local checkouts `magi repos` lists. Without this
1280/// line an agent asked to change some other repository has no way to know
1281/// that option exists, and the only path it can see - asking the operator to
1282/// dictate a full path - is exactly the friction this change exists to
1283/// remove. A miss or an ambiguous name still fails the command outright, so
1284/// the instruction is to ask rather than guess when that happens - the
1285/// silent-decision line this task must not cross.
1286pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1287    let write_policy = if allow_write {
1288        "Write access is enabled for this conversation (`allow_write = \
1289         true`), so you may write files - but only a small, \
1290         already-decided edit the operator names outright in this \
1291         conversation, not an implementation. This is a permission on the \
1292         conversation as a whole, not a property of whichever repository \
1293         it happened to start in: if the operator names a different \
1294         repository for that small edit, the policy allows it there too. \
1295         Your own tool may still confine writes to the repository this \
1296         conversation started in regardless - if a write elsewhere is \
1297         refused, say so plainly rather than working around it. Once you \
1298         have made an edit, say plainly what you edited. Anything bigger, \
1299         or anything still open-ended, still goes through the queue below \
1300         rather than being done here."
1301    } else {
1302        "Do not write files. Implementing a change is not this \
1303         conversation's job; a separate, blind competition of agents does \
1304         that, and a repository this conversation has already edited would \
1305         make their diffs unjudgeable."
1306    };
1307    let mut out = format!(
1308        "You are magi's standing conversation partner for its operator, who \
1309         usually has this open on a phone. Keep replies short: no preamble, \
1310         no restating what they just said.\n\n\
1311         # Repository\n\n{repo}\n\n\
1312         You may look around: read files, run shell commands, search history, \
1313         run tests - whatever answers the question. {write_policy}\n\n\
1314         A short, command-shaped message (\"list\", \"info <id>\", \"show \
1315         3cbf\") is almost always the operator asking you to look something \
1316         up, not an instruction to file - answer it yourself with `magi \
1317         list`, `magi show <id>`, `magi task list`, or the like, the same way \
1318         you would answer any other question in this conversation.\n\n\
1319         # When the operator wants something done\n\n\
1320         Run:\n\n\
1321         magi task add --solo --repo {repo} <instruction>\n\n\
1322         and tell the operator the task id it prints, so they can follow it \
1323         from the Queue. If it refuses with a duplicate warning (the \
1324         instruction names a branch, commit or pull request that an \
1325         unfinished task, run or PR already owns), do not repeat it with \
1326         --force yourself: tell the operator what it matched and let them \
1327         decide. Write <instruction> so that an implementer who has \
1328         never seen this conversation can act on it alone - it is everything \
1329         they get. Use --solo: it runs the task through one implementer \
1330         straight into review instead of the usual multi-agent competition, \
1331         which is the right shape for a change this conversation has already \
1332         settled, rather than one still worth several independent takes.\n\n\
1333         If the operator asks for something in a different repository, \
1334         --repo does not have to be a full path: --repo owner/repo (or just \
1335         repo, when that is unambiguous) is resolved against local checkouts \
1336         the same way `magi repos` lists them. If the command fails because \
1337         nothing matches or more than one checkout shares that name, ask the \
1338         operator which repository they mean (or run `magi repos` yourself \
1339         to see the candidates) rather than guessing.\n\n\
1340         If the operator attached an image (a screenshot, say) that the task \
1341         is about, pass it with `--attach <path>`, using the absolute path \
1342         the turn's attachment note gives; repeat the flag for several. \
1343         `magi task add --solo --attach <path> <instruction>` copies the \
1344         file into the task, so the implementer receives it. Do not paste the \
1345         path into <instruction> instead: deleting this conversation deletes \
1346         its attachments, and then that path reaches no one.\n",
1347        repo = repo.display(),
1348    );
1349    out.push_str(&language_note(language));
1350    out
1351}
1352
1353/// The operator is talking, so their language matters here more than in most
1354/// prompts magi sends.
1355fn language_note(language: &str) -> String {
1356    if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1357        String::new()
1358    } else {
1359        format!("\nHold this conversation in {language}.\n")
1360    }
1361}
1362
1363/// Queue tasks this conversation has filed, oldest first.
1364///
1365/// A task is this conversation's when its [`Source::Agent`] names this
1366/// conversation's id as `run` - which is exactly what happens when
1367/// `magi task add` is run from inside a turn, because [`turn`] passes the
1368/// conversation's own id as [`Invocation::run`].
1369pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1370    let mut tasks: Vec<Task> = queue
1371        .list()
1372        .into_iter()
1373        .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1374        .collect();
1375    tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1376    tasks
1377}
1378
1379fn read_path(path: &Path) -> Result<Talk> {
1380    let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1381    serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1382}
1383
1384/// How many times [`write_atomic`] retries a failed write-then-rename before
1385/// giving up.
1386const PUT_RETRIES: u32 = 5;
1387
1388/// Write `body` to `tmp` and rename it onto `path`, retrying the whole thing
1389/// a handful of times with a short sleep in between.
1390///
1391/// The only failure this is meant to absorb is a passing one - most
1392/// concretely, a reader elsewhere in this process (or another `magi`
1393/// process) with `path` briefly open for `read_to_string` at the exact
1394/// moment this call tries to rename over it. That clears in milliseconds
1395/// once the reader lets go; a caller still failing after several short
1396/// sleeps has something more durable wrong (a full disk, a permissions
1397/// change) that a longer sleep would not fix either, and is left to report
1398/// it.
1399fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1400    let mut last_err = None;
1401    for attempt in 0..PUT_RETRIES {
1402        if attempt > 0 {
1403            std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1404        }
1405        match try_write_atomic(tmp, path, body) {
1406            Ok(()) => return Ok(()),
1407            Err(e) => last_err = Some(e),
1408        }
1409    }
1410    Err(last_err.expect("the loop above always runs at least once"))
1411}
1412
1413fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1414    #[cfg(test)]
1415    if failpoint::take_forced_put_failure() {
1416        bail!("simulated write failure (test)");
1417    }
1418    std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1419    std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1420    Ok(())
1421}
1422
1423/// Last resort when `turn`'s own `store.put` fails even after
1424/// [`write_atomic`]'s retries: keep the generated text somewhere still
1425/// findable rather than let the whole of an agent's answer disappear along
1426/// with the write that was supposed to record it.
1427fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1428    let dir = store.artifacts_of(id);
1429    std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1430    let path = dir.join(format!("{stem}-lost.txt"));
1431    std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1432    Ok(path)
1433}
1434
1435/// A test-only seam that lets [`try_write_atomic`] simulate the kind of
1436/// passing I/O race [`write_atomic`] is meant to retry through, without
1437/// depending on real OS-level file-locking behaviour, which differs across
1438/// the three platforms this crate ships on (and, on the one platform where a
1439/// reader really does block a rename, is awkward to trigger deterministically
1440/// in a unit test).
1441#[cfg(test)]
1442mod failpoint {
1443    use std::cell::Cell;
1444
1445    thread_local! {
1446        static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1447    }
1448
1449    /// Arrange for the next `count` calls into [`super::try_write_atomic`] to
1450    /// fail before touching the filesystem at all.
1451    pub(super) fn force_put_failures(count: u32) {
1452        FORCE_PUT_FAILURES.with(|c| c.set(count));
1453    }
1454
1455    /// Consumed once per attempt inside [`super::try_write_atomic`]; `true`
1456    /// means simulate this attempt failing.
1457    pub(super) fn take_forced_put_failure() -> bool {
1458        FORCE_PUT_FAILURES.with(|c| {
1459            let n = c.get();
1460            if n == 0 {
1461                false
1462            } else {
1463                c.set(n - 1);
1464                true
1465            }
1466        })
1467    }
1468}
1469
1470fn short(id: &str) -> &str {
1471    id.split('-').next_back().unwrap_or(id)
1472}
1473
1474fn new_id() -> String {
1475    let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1476    let seed = crate::rng::entropy();
1477    format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1478}
1479
1480/// Extension an attachment's bytes are stored under, from its (already
1481/// validated) mime. The one place this mapping exists on the write side;
1482/// `web`'s own whitelist is what actually decides which mimes are accepted
1483/// in the first place.
1484fn attachment_ext(mime: &str) -> Option<&'static str> {
1485    match mime {
1486        "image/png" => Some("png"),
1487        "image/jpeg" => Some("jpg"),
1488        "image/gif" => Some("gif"),
1489        "image/webp" => Some("webp"),
1490        _ => None,
1491    }
1492}
1493
1494/// Is `id` a shape [`put_attachment`](Talks::put_attachment) could have
1495/// produced? 32 lowercase hex digits and nothing else, checked before an id
1496/// that came from the client is ever allowed to build a path - so `..` and a
1497/// path separator are never even possible.
1498pub fn valid_attachment_id(id: &str) -> bool {
1499    id.len() == 32
1500        && id
1501            .bytes()
1502            .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
1503}
1504
1505/// A fresh attachment id: 128 bits of process entropy as lowercase hex - the
1506/// same "mint it, never take it from the client" rule [`new_id`] follows for
1507/// conversation ids.
1508fn new_attachment_id() -> String {
1509    let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
1510    format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
1511}
1512
1513#[cfg(test)]
1514mod tests {
1515    use std::collections::BTreeMap;
1516
1517    use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
1518    use crate::queue::{Queue, Source, Task};
1519
1520    use super::*;
1521
1522    fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
1523        AgentSpec {
1524            id: id.to_owned(),
1525            kind: AgentKind::Command,
1526            model: model.map(str::to_owned),
1527            command: Vec::new(),
1528            extra_args: Vec::new(),
1529            env: BTreeMap::new(),
1530            prompt_delivery: None,
1531        }
1532    }
1533
1534    fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
1535        Talk {
1536            schema: SCHEMA,
1537            id: "20260904-014455-ab12".to_owned(),
1538            repo: PathBuf::from("."),
1539            agent: agent.to_owned(),
1540            status: TalkStatus::Open,
1541            turns,
1542            pending: String::new(),
1543            pending_attachments: Vec::new(),
1544            fallback: false,
1545            created_at: Timestamp::now(),
1546            updated_at: Timestamp::now(),
1547            seat: SeatState::new(SEAT, agent, 1),
1548        }
1549    }
1550
1551    fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
1552        Turn {
1553            who: Who::Agent,
1554            body: body.to_owned(),
1555            at: Timestamp::now(),
1556            attachments: Vec::new(),
1557            usage: usage.map(|(t, a, m)| TurnUsage {
1558                context_tokens: t,
1559                agent: a.to_owned(),
1560                model: m.map(str::to_owned),
1561            }),
1562        }
1563    }
1564
1565    fn ctx_config(windows: &[(&str, u64)]) -> Config {
1566        Config {
1567            agents: vec![
1568                ctx_agent("small", Some("small-model")),
1569                ctx_agent("big", Some("big-model")),
1570                ctx_agent("plain", None),
1571            ],
1572            context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
1573            ..Config::default()
1574        }
1575    }
1576
1577    #[test]
1578    fn context_usage_computes_percent_and_warns_at_eighty() {
1579        let cfg = ctx_config(&[("small-model", 1000)]);
1580        let at = |tokens| {
1581            let t = ctx_talk(
1582                "small",
1583                vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
1584            );
1585            context_usage(&t, Some(&cfg))
1586        };
1587        let u = at(799);
1588        assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
1589        let u = at(800);
1590        assert_eq!((u.percent, u.warn), (Some(80), true));
1591        let u = at(1500);
1592        assert_eq!((u.percent, u.warn), (Some(150), true));
1593        assert!(!u.since_switch);
1594    }
1595
1596    #[test]
1597    fn context_usage_is_unknown_without_usage_and_never_looks_back() {
1598        let cfg = ctx_config(&[("small-model", 1000)]);
1599        let t = ctx_talk(
1600            "small",
1601            vec![
1602                reply("old", Some((900, "small", Some("small-model")))),
1603                reply("new", None),
1604            ],
1605        );
1606        let u = context_usage(&t, Some(&cfg));
1607        assert_eq!((u.tokens, u.percent, u.warn), (None, None, false));
1608        // A magi note after the reply neither hides nor replaces it.
1609        let t = ctx_talk(
1610            "small",
1611            vec![
1612                reply("old", Some((900, "small", Some("small-model")))),
1613                reply("magi: could not run agent", None),
1614            ],
1615        );
1616        assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
1617        assert_eq!(
1618            context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
1619            None
1620        );
1621    }
1622
1623    #[test]
1624    fn context_usage_without_a_window_shows_tokens_only() {
1625        let cfg = ctx_config(&[]);
1626        // No model at all, and a model nothing matches.
1627        let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
1628        let u = context_usage(&t, Some(&cfg));
1629        assert_eq!(
1630            (u.tokens, u.window, u.percent, u.warn),
1631            (Some(5000), None, None, false)
1632        );
1633        let t = ctx_talk(
1634            "small",
1635            vec![reply("hi", Some((5000, "small", Some("small-model"))))],
1636        );
1637        assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
1638        // No readable config: same, and no panic.
1639        assert_eq!(context_usage(&t, None).window, None);
1640    }
1641
1642    #[test]
1643    fn context_usage_switching_model_changes_the_denominator() {
1644        let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
1645        let used = reply("hi", Some((900, "small", Some("small-model"))));
1646        let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
1647        assert_eq!(
1648            (before.percent, before.warn, before.since_switch),
1649            (Some(90), true, false)
1650        );
1651        // Same turns, conversation now on the big model: new denominator, and
1652        // the figure is flagged as describing the previous session.
1653        let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
1654        assert_eq!(after.window, Some(10_000));
1655        assert_eq!(
1656            (after.percent, after.warn, after.since_switch),
1657            (Some(9), false, true)
1658        );
1659        assert_eq!(after.model.as_deref(), Some("big-model"));
1660    }
1661
1662    #[test]
1663    fn a_turn_recorded_before_usage_existed_still_reads() {
1664        let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
1665        let turn: Turn = serde_json::from_str(old).expect("old turn reads");
1666        assert!(turn.usage.is_none());
1667        let json = serde_json::to_string(&turn).expect("serialize");
1668        assert!(
1669            !json.contains("usage"),
1670            "absent usage is not written: {json}"
1671        );
1672    }
1673
1674    /// A store of its own, with no process-global state.
1675    fn store() -> (tempfile::TempDir, Talks) {
1676        let tmp = tempfile::tempdir().expect("tempdir");
1677        let talks = Talks::at(tmp.path().join("talks"));
1678        (tmp, talks)
1679    }
1680
1681    /// A `kind = "command"` agent whose whole behaviour is a POSIX shell
1682    /// script - see `chat`'s tests for why no test here may spawn a real
1683    /// agent CLI.
1684    fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
1685        let path = dir.join("mock-talk-agent.sh");
1686        std::fs::write(&path, script).expect("write mock");
1687        AgentSpec {
1688            id: "mock".to_owned(),
1689            kind: AgentKind::Command,
1690            model: None,
1691            command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1692            extra_args: Vec::new(),
1693            env,
1694            prompt_delivery: None,
1695        }
1696    }
1697
1698    fn config(spec: AgentSpec) -> Config {
1699        Config {
1700            agents: vec![spec],
1701            graph: Graph {
1702                language: "en".to_owned(),
1703                ..Graph::default()
1704            },
1705            ..Config::default()
1706        }
1707    }
1708
1709    /// Echo a canned reply, ignoring the prompt on stdin.
1710    const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
1711
1712    /// Say nothing and fail, the way a CLI that cannot start does.
1713    const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
1714
1715    /// Reply with the prompt it was given, so a test can inspect exactly what
1716    /// the agent received on stdin.
1717    const ECHO: &str = "#!/bin/sh\ncat\n";
1718
1719    fn env(reply: &str) -> BTreeMap<String, String> {
1720        BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
1721    }
1722
1723    #[test]
1724    fn the_frozen_json_field_names_round_trip_through_disk() {
1725        let (tmp, talks) = store();
1726        let mut talk = Talk {
1727            schema: SCHEMA,
1728            id: "20260904-014455-ab12".to_owned(),
1729            repo: tmp.path().to_owned(),
1730            agent: "sonnet".to_owned(),
1731            status: TalkStatus::Open,
1732            turns: Vec::new(),
1733            pending: String::new(),
1734            pending_attachments: Vec::new(),
1735            fallback: false,
1736            created_at: Timestamp::now(),
1737            updated_at: Timestamp::now(),
1738            seat: SeatState::new(SEAT, "sonnet", 7),
1739        };
1740        talks.put(&mut talk).expect("put");
1741
1742        let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
1743        let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
1744        for field in [
1745            "schema",
1746            "id",
1747            "repo",
1748            "agent",
1749            "status",
1750            "turns",
1751            "created_at",
1752            "updated_at",
1753        ] {
1754            assert!(v.get(field).is_some(), "missing field `{field}`");
1755        }
1756        assert_eq!(v["schema"], 1);
1757        assert_eq!(v["status"], "open");
1758
1759        let back = talks.get(&talk.id).expect("get");
1760        assert_eq!(back.id, talk.id);
1761        assert_eq!(back.status, TalkStatus::Open);
1762    }
1763
1764    #[test]
1765    fn opening_a_talk_takes_no_agent_turn() {
1766        let (tmp, talks) = store();
1767        // A script that would fail loudly if it were ever run: `begin` must
1768        // not invoke anything, since there is nothing yet for an agent to
1769        // answer.
1770        let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1771        let cfg = config(spec);
1772
1773        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1774        assert_eq!(talk.status, TalkStatus::Open);
1775        assert!(talk.turns.is_empty(), "nothing has been said yet");
1776
1777        let on_disk = talks.get(&talk.id).expect("get");
1778        assert_eq!(on_disk.turns.len(), 0);
1779    }
1780
1781    /// `[roles] chatter`, when set, decides who holds this conversation; unset,
1782    /// it falls back to [`agent::pick`]'s own default order (a claude seat,
1783    /// else the first runnable agent in roster order) rather than to any
1784    /// other role - see `[roles] chatter`'s own doc in [`crate::config`] for
1785    /// why a dedicated field exists at all: opening this against the same
1786    /// seat as a judge is what produced the `agent ... did not answer within
1787    /// 300s` timeout that led to it.
1788    #[test]
1789    fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
1790        let (tmp, talks) = store();
1791        let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1792        let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1793        chatter_spec.id = "chatter-mock".to_owned();
1794
1795        let mut cfg = Config {
1796            agents: vec![first_spec.clone(), chatter_spec.clone()],
1797            graph: Graph {
1798                language: "en".to_owned(),
1799                ..Graph::default()
1800            },
1801            ..Config::default()
1802        };
1803        cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
1804
1805        let talk =
1806            begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
1807        assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
1808
1809        cfg.roles.chatter = None;
1810        let fallback =
1811            begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
1812        assert_eq!(
1813            fallback.agent, first_spec.id,
1814            "unset chatter must fall back to agent::pick's own default order"
1815        );
1816    }
1817
1818    /// A conversation recorded before attachments existed - schema 1, no
1819    /// `attachments` key on any turn - must still read.
1820    #[test]
1821    fn a_talk_recorded_without_attachments_still_reads() {
1822        let (tmp, talks) = store();
1823        let path = talks.path_of("20260904-014455-ab12");
1824        std::fs::create_dir_all(talks.root()).expect("talks dir");
1825        std::fs::write(
1826            &path,
1827            serde_json::json!({
1828                "schema": 1,
1829                "id": "20260904-014455-ab12",
1830                "repo": tmp.path(),
1831                "agent": "sonnet",
1832                "status": "open",
1833                "turns": [
1834                    { "who": "operator", "body": "still there?",
1835                      "at": Timestamp::now().to_string() },
1836                ],
1837                "created_at": Timestamp::now().to_string(),
1838                "updated_at": Timestamp::now().to_string(),
1839                "seat": SeatState::new(SEAT, "sonnet", 7),
1840            })
1841            .to_string(),
1842        )
1843        .expect("write pre-attachments talk");
1844
1845        let talk = talks.get("20260904-014455-ab12").expect("must still read");
1846        assert!(talk.turns[0].attachments.is_empty());
1847    }
1848
1849    #[test]
1850    fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
1851        let (tmp, talks) = store();
1852        let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1853        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1854
1855        queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
1856        queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
1857        let saved = talks.get(&talk.id).expect("reload queued talk");
1858        assert_eq!(saved.pending, "first\n\nsecond");
1859        assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
1860
1861        let drained = drain(&mut talk, &talks).expect("drain");
1862        assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
1863        let saved = talks.get(&talk.id).expect("reload drained talk");
1864        assert!(saved.pending.is_empty());
1865        assert_eq!(saved.turns.len(), 1);
1866        assert_eq!(saved.turns[0].body, "first\n\nsecond");
1867    }
1868
1869    #[test]
1870    fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
1871        let (tmp, talks) = store();
1872        let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1873        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1874        let attachment = Attachment {
1875            id: "a".repeat(32),
1876            name: "shot.png".to_owned(),
1877            mime: "image/png".to_owned(),
1878            bytes: 3,
1879        };
1880
1881        queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
1882        assert!(
1883            edit_pending_text(
1884                &mut talk,
1885                &talks,
1886                "corrected",
1887                "first",
1888                std::slice::from_ref(&attachment.id),
1889            )
1890            .expect("edit")
1891        );
1892        let saved = talks.get(&talk.id).expect("reload edited draft");
1893        assert_eq!(saved.pending, "corrected");
1894        assert_eq!(saved.pending_attachments, vec![attachment]);
1895
1896        queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
1897        assert!(
1898            !edit_pending_text(
1899                &mut talk,
1900                &talks,
1901                "stale edit",
1902                "corrected",
1903                &["a".repeat(32)],
1904            )
1905            .expect("stale edit is a conflict")
1906        );
1907        assert_eq!(
1908            talks.get(&talk.id).expect("reload after conflict").pending,
1909            "corrected\n\nlater"
1910        );
1911        assert!(
1912            !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
1913                .expect("stale clear is a conflict")
1914        );
1915        assert_eq!(
1916            talks
1917                .get(&talk.id)
1918                .expect("reload after stale clear")
1919                .pending,
1920            "corrected\n\nlater"
1921        );
1922    }
1923
1924    #[tokio::test]
1925    async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
1926        let (tmp, talks) = store();
1927        let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
1928        let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
1929        let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1930        let id = running.id.clone();
1931        let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
1932
1933        let response_talks = talks.clone();
1934        let response_cfg = cfg.clone();
1935        let reply = tokio::spawn(async move {
1936            respond(&mut running, &response_talks, &response_cfg, &first).await
1937        });
1938        tokio::time::sleep(std::time::Duration::from_millis(20)).await;
1939
1940        let mut queued = talks.get(&id).expect("queued handle");
1941        queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
1942        reply.await.expect("join").expect("reply");
1943
1944        let saved = talks.get(&id).expect("reload");
1945        assert_eq!(saved.pending, "next");
1946        assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
1947    }
1948
1949    /// A mock whose script counts its own calls in `<dir>/<id>.calls` before
1950    /// running `body`.
1951    fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
1952        let calls = dir.join(format!("{id}.calls"));
1953        let script = format!(
1954            "#!/bin/sh\necho x >> '{}'\n{body}\n",
1955            calls.to_string_lossy()
1956        );
1957        let path = dir.join(format!("mock-{id}.sh"));
1958        std::fs::write(&path, script).expect("write mock");
1959        AgentSpec {
1960            id: id.to_owned(),
1961            kind: AgentKind::Command,
1962            model: None,
1963            command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1964            extra_args: Vec::new(),
1965            env: BTreeMap::new(),
1966            prompt_delivery: None,
1967        }
1968    }
1969
1970    fn calls(dir: &Path, id: &str) -> usize {
1971        std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
1972    }
1973
1974    fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
1975        let mut cfg = config(specs[0].clone());
1976        cfg.agents = specs;
1977        cfg.roles.chatter = Some(AgentChoice::Chain(
1978            ids.iter().map(|s| (*s).to_owned()).collect(),
1979        ));
1980        cfg
1981    }
1982
1983    #[tokio::test]
1984    async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
1985        let (tmp, talks) = store();
1986        let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
1987        let b = counting_agent(tmp.path(), "b", "cat");
1988        let cfg = chain_config(vec![a, b], &["a", "b"]);
1989        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1990        assert_eq!(talk.agent, "a");
1991
1992        say(&mut talk, &talks, &cfg, "hello there", Vec::new())
1993            .await
1994            .expect("turn");
1995        assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
1996        assert_eq!(calls(tmp.path(), "b"), 1);
1997        assert_eq!(talk.agent, "b", "the switch persists");
1998        assert!(talks.get(&talk.id).unwrap().agent == "b");
1999        let reply = talk.turns.last().unwrap();
2000        assert!(reply.body.contains("hello there"));
2001        assert!(
2002            reply.body.contains("magi task add --solo"),
2003            "a fresh seat gets the full briefing"
2004        );
2005        assert!(
2006            talk.turns
2007                .iter()
2008                .any(|t| t.body.contains("agent changed from a to b")),
2009            "the switch is noted"
2010        );
2011    }
2012
2013    #[tokio::test]
2014    async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2015        let (tmp, talks) = store();
2016        let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2017        let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2018        let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2019        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2020
2021        let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2022            .await
2023            .expect_err("every agent failed");
2024        assert!(err.to_string().contains("`a`"), "{err:#}");
2025        assert_eq!(calls(tmp.path(), "a"), 1);
2026        assert_eq!(calls(tmp.path(), "b"), 1);
2027        assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2028    }
2029
2030    #[test]
2031    fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2032        let (tmp, talks) = store();
2033        let b = counting_agent(tmp.path(), "b", "cat");
2034        let cfg = chain_config(vec![b], &["ghost", "b"]);
2035        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2036        assert_eq!(talk.agent, "b");
2037    }
2038
2039    #[tokio::test]
2040    async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
2041        let (tmp, talks) = store();
2042        let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2043        let b = counting_agent(tmp.path(), "b", "cat");
2044        let cfg = chain_config(vec![a, b], &["a", "b"]);
2045        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2046        say(&mut talk, &talks, &cfg, "hi", Vec::new())
2047            .await
2048            .expect_err("a alone, and it fails");
2049        assert_eq!(calls(tmp.path(), "b"), 0);
2050        assert_eq!(talk.agent, "a");
2051    }
2052
2053    #[tokio::test]
2054    async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2055        let (tmp, talks) = store();
2056        let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2057        let b = counting_agent(tmp.path(), "b", "cat");
2058        let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2059        let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2060        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2061        say(&mut talk, &talks, &cfg, "hi", Vec::new())
2062            .await
2063            .expect_err("c alone, and it fails");
2064        assert_eq!(calls(tmp.path(), "b"), 0);
2065    }
2066
2067    #[tokio::test]
2068    async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
2069        let (tmp, talks) = store();
2070        let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2071        let cfg = config(spec);
2072        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2073
2074        say(
2075            &mut talk,
2076            &talks,
2077            &cfg,
2078            "what does the queue module do?",
2079            Vec::new(),
2080        )
2081        .await
2082        .expect("first turn");
2083        let first_prompt = &talk.turns[1].body;
2084        assert!(first_prompt.contains("magi task add --solo"));
2085        assert!(first_prompt.contains("what does the queue module do?"));
2086
2087        say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
2088            .await
2089            .expect("second turn");
2090        let second_prompt = &talk.turns[3].body;
2091        assert!(
2092            !second_prompt.contains("magi task add --solo"),
2093            "the briefing is sent once, not on every turn: {second_prompt}"
2094        );
2095        assert!(second_prompt.contains("and how is it locked?"));
2096    }
2097
2098    #[tokio::test]
2099    async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
2100        let (tmp, talks) = store();
2101        let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2102        let mut b = a.clone();
2103        b.id = "other".to_owned();
2104        let mut cfg = config(a.clone());
2105        cfg.agents.push(b.clone());
2106        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
2107        say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
2108            .await
2109            .expect("first turn");
2110        let old_session = talk.seat.claude_session.clone();
2111        assert_eq!(talk.seat.turns, 1);
2112
2113        assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
2114        assert_eq!(talk.agent, "other");
2115        assert_eq!(talk.seat.turns, 0);
2116        assert_eq!(talk.seat.agent, "other");
2117        assert_ne!(talk.seat.claude_session, old_session);
2118        let note = talk.turns.last().expect("note");
2119        assert_eq!(note.who, Who::Agent);
2120        assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2121        assert!(note.body.contains("changed from"), "{}", note.body);
2122        assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
2123
2124        let before = talk.turns.len();
2125        assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
2126        assert_eq!(talk.turns.len(), before, "a no-op writes no note");
2127
2128        say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
2129            .await
2130            .expect("turn after switch");
2131        let prompt = &talk.turns.last().expect("reply").body;
2132        assert!(prompt.contains("remember the walrus"), "{prompt}");
2133        assert!(prompt.contains("## magi"), "{prompt}");
2134        assert!(prompt.contains("what did I say?"), "{prompt}");
2135    }
2136
2137    #[tokio::test]
2138    async fn say_appends_the_operator_turn_then_the_agent_turn() {
2139        let (tmp, talks) = store();
2140        let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2141        let cfg = config(spec);
2142        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2143
2144        say(
2145            &mut talk,
2146            &talks,
2147            &cfg,
2148            "can I rename this function?",
2149            Vec::new(),
2150        )
2151        .await
2152        .expect("say");
2153
2154        assert_eq!(talk.turns.len(), 2);
2155        assert_eq!(talk.turns[0].who, Who::Operator);
2156        assert_eq!(talk.turns[0].body, "can I rename this function?");
2157        assert_eq!(talk.turns[1].who, Who::Agent);
2158        assert_eq!(talk.turns[1].body, "go ahead");
2159        assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
2160    }
2161
2162    #[tokio::test]
2163    async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
2164        let (tmp, talks) = store();
2165        let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2166        let cfg = config(spec);
2167        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2168
2169        let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2170            .await
2171            .expect_err("a turn with no answer is an error");
2172        assert!(err.to_string().contains("no answer"), "{err}");
2173
2174        let on_disk = talks.get(&talk.id).expect("get");
2175        assert_eq!(on_disk.turns.len(), 2);
2176        assert_eq!(on_disk.turns[0].body, "check the tests");
2177        let note = &on_disk.turns[1];
2178        assert_eq!(note.who, Who::Agent);
2179        assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2180        assert!(note.body.contains("your message is saved"));
2181    }
2182
2183    /// The failure this stands in for: a reader elsewhere briefly has the
2184    /// talk file open right when `turn` tries to save the reply, and the
2185    /// write-then-rename fails once or twice before the reader lets go.
2186    /// `write_atomic`'s own retries must absorb that with nobody the wiser -
2187    /// no gap in the transcript, no dropped turn.
2188    #[tokio::test]
2189    async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
2190        let (tmp, talks) = store();
2191        let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2192        let cfg = config(spec);
2193        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2194
2195        let text =
2196            record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
2197        // One fewer failure than `write_atomic` will retry through, so the
2198        // very last attempt must succeed.
2199        failpoint::force_put_failures(PUT_RETRIES - 1);
2200        respond(&mut talk, &talks, &cfg, &text)
2201            .await
2202            .expect("respond must survive a write failure its own retries can outlast");
2203
2204        assert_eq!(talk.turns.len(), 2);
2205        assert_eq!(talk.turns[1].who, Who::Agent);
2206        assert_eq!(talk.turns[1].body, "go ahead");
2207        let on_disk = talks.get(&talk.id).expect("get");
2208        assert_eq!(
2209            on_disk.turns, talk.turns,
2210            "the reply must reach disk despite the early write failures"
2211        );
2212    }
2213
2214    /// When the write-then-rename never recovers - standing in for a disk
2215    /// that stays unwritable rather than a reader that eventually lets go -
2216    /// the reply must not disappear without a trace the way it did in the
2217    /// real incident this repository saw: no error on the phone, no note in
2218    /// the transcript, and the turn simply gone from `talks/<id>.json`.
2219    #[tokio::test]
2220    async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
2221        let (tmp, talks) = store();
2222        let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2223        let cfg = config(spec);
2224        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2225
2226        let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2227        // Exactly enough forced failures to exhaust the reply's own retries;
2228        // the shorter note that replaces it then saves cleanly, which is the
2229        // common case this exercises - a large write racing something,
2230        // followed by a small one that does not.
2231        failpoint::force_put_failures(PUT_RETRIES);
2232        let err = respond(&mut talk, &talks, &cfg, &text)
2233            .await
2234            .expect_err("a reply that cannot be saved must be reported, not swallowed");
2235        assert!(err.to_string().contains("could not be saved"), "{err}");
2236
2237        let on_disk = talks.get(&talk.id).expect("get");
2238        assert_eq!(
2239            on_disk.turns.len(),
2240            2,
2241            "the operator turn plus a visible note"
2242        );
2243        assert_eq!(on_disk.turns[0].body, "check the tests");
2244        let note = &on_disk.turns[1];
2245        assert_eq!(note.who, Who::Agent);
2246        assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2247        assert!(
2248            note.body.contains("could not be saved"),
2249            "the operator must be told the reply is missing, not left staring \
2250             at a gap with no explanation: {}",
2251            note.body
2252        );
2253        assert_eq!(
2254            talk.turns, on_disk.turns,
2255            "the in-memory talk must match what actually landed on disk"
2256        );
2257
2258        // The generated answer itself must still be recoverable, not merely
2259        // reported as lost.
2260        let artifacts = talks.artifacts_of(&talk.id);
2261        let stash = std::fs::read_dir(&artifacts)
2262            .expect("artifacts dir")
2263            .filter_map(|e| e.ok())
2264            .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
2265            .expect("a stash file for the lost reply");
2266        let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
2267        assert_eq!(stashed, "go ahead");
2268
2269        // Losing the reply must not also lose the seat. The CLI took a turn
2270        // and consumed this seat's session id; if the note's write left the
2271        // record claiming otherwise, the next turn would re-open a session
2272        // the CLI is already holding - the `20260907-011805-fb57` desync -
2273        // and would re-send the whole briefing besides. Both decisions read
2274        // the seat straight off disk (`agent::has_session` and `turn`'s own
2275        // `seat.turns == 0` branch), so this is the field that has to match.
2276        assert_eq!(
2277            on_disk.seat.turns, 1,
2278            "the note's write must carry the turn the CLI actually took"
2279        );
2280        assert_eq!(
2281            on_disk.seat.claude_session, talk.seat.claude_session,
2282            "the session id handed to the CLI must survive the failed reply"
2283        );
2284        assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
2285        assert!(
2286            agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
2287            "the next turn must resume, not open the same session id twice"
2288        );
2289    }
2290
2291    /// Even the note can fail to save, if the disk stays unwritable for long
2292    /// enough. `respond` must still report the failure rather than pretend
2293    /// the turn succeeded, and must not leave the in-memory `talk` claiming
2294    /// a turn that never reached disk.
2295    #[tokio::test]
2296    async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
2297        let (tmp, talks) = store();
2298        let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2299        let cfg = config(spec);
2300        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2301
2302        let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2303        // Enough forced failures to exhaust the retries for both the reply
2304        // and the note that would have replaced it.
2305        failpoint::force_put_failures(PUT_RETRIES * 2);
2306        let err = respond(&mut talk, &talks, &cfg, &text)
2307            .await
2308            .expect_err("neither the reply nor the note could be saved");
2309        assert!(err.to_string().contains("could not be saved"), "{err}");
2310
2311        assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
2312        let on_disk = talks.get(&talk.id).expect("get");
2313        assert_eq!(on_disk.turns.len(), 1);
2314
2315        // Nothing at all reached disk, so the seat could not either: the CLI
2316        // took a turn the record does not know about. That is pinned here as
2317        // the known cost of a file that cannot be written twice over, not as
2318        // something this branch could do better - the only way to record the
2319        // seat is the write that just failed. It is also the point where
2320        // this meets `20260907-011805-fb57`: a turn taken before some later
2321        // write lands would re-open a session id the CLI already holds. The
2322        // in-memory seat keeps the truth the CLI reported, which is why it is
2323        // not wound back to match.
2324        assert_eq!(
2325            on_disk.seat.turns, 0,
2326            "an unwritable file cannot record the turn the CLI took"
2327        );
2328        assert_eq!(
2329            talk.seat.turns, 1,
2330            "the in-memory seat still reports the turn the CLI actually took"
2331        );
2332        assert_eq!(
2333            on_disk.seat.claude_session, talk.seat.claude_session,
2334            "the session id was minted at `begin` and never changes here"
2335        );
2336    }
2337
2338    /// An attachment lets the operator send an otherwise-empty message, and
2339    /// its absolute path is what actually reaches the agent's prompt - here
2340    /// on the very first turn, where it has to share the briefing.
2341    #[tokio::test]
2342    async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
2343        let (tmp, talks) = store();
2344        let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2345        let cfg = config(spec);
2346        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2347
2348        let att = talks
2349            .put_attachment(
2350                &talk.id,
2351                "image/png",
2352                "screenshot.png",
2353                b"pretend-png-bytes",
2354            )
2355            .expect("put attachment");
2356
2357        say(&mut talk, &talks, &cfg, "", vec![att.clone()])
2358            .await
2359            .expect("an empty body with an attachment is still a turn");
2360
2361        let operator_turn = &talk.turns[0];
2362        assert_eq!(operator_turn.who, Who::Operator);
2363        assert_eq!(operator_turn.body, "");
2364        assert_eq!(operator_turn.attachments, vec![att.clone()]);
2365
2366        let prompt = &talk.turns[1].body;
2367        let expected_path = talks
2368            .attachments_dir(&talk.id)
2369            .join(format!("{}.png", att.id));
2370        assert!(
2371            prompt.contains(&expected_path.display().to_string()),
2372            "the agent must be told the attachment's absolute path: {prompt}"
2373        );
2374        assert!(prompt.contains("image/png"), "and its mime: {prompt}");
2375    }
2376
2377    /// See `chat`'s test of the same name: `run::home()` returns a bare
2378    /// relative `PathBuf` verbatim when `MAGI_HOME` is set to a relative
2379    /// path, so a `Talks` store built on it has a relative `root` too. That
2380    /// is fine for this store's own I/O, which runs in this process against
2381    /// this process's cwd, but `attachment_path` hands its result to a
2382    /// *different* process invoked with `cwd: &talk.repo` - an uncorrected
2383    /// relative path would resolve against the repository instead of
2384    /// wherever the attachment actually landed.
2385    #[test]
2386    fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
2387        let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
2388        let att = Attachment {
2389            id: "0".repeat(32),
2390            name: "shot.png".to_owned(),
2391            mime: "image/png".to_owned(),
2392            bytes: 3,
2393        };
2394        let path = talks
2395            .attachment_path("some-talk-id", &att)
2396            .expect("a supported mime always yields a path");
2397        assert!(
2398            path.is_absolute(),
2399            "must be absolute even off a relative store root: {}",
2400            path.display()
2401        );
2402    }
2403
2404    #[tokio::test]
2405    async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
2406        // `[graph] timeout_talk` must be the number this module actually
2407        // waits, not a leftover hardcoded fifteen minutes - so the mock
2408        // sleeps past a deliberately tiny override and the failure note is
2409        // checked against that same override, not the old default.
2410        let (tmp, talks) = store();
2411        let slow = mock_agent(
2412            tmp.path(),
2413            "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
2414            BTreeMap::new(),
2415        );
2416        let mut cfg = config(slow);
2417        cfg.graph.timeout_talk = 1;
2418        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2419
2420        let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2421            .await
2422            .expect_err("a turn that never answers is an error");
2423        assert!(
2424            err.to_string().contains("did not answer within 1s"),
2425            "{err}"
2426        );
2427
2428        let on_disk = talks.get(&talk.id).expect("get");
2429        let note = on_disk.turns.last().expect("a note turn was recorded");
2430        assert!(
2431            note.body.contains("did not answer within 1s"),
2432            "the transcript must show the configured timeout: {}",
2433            note.body
2434        );
2435    }
2436
2437    #[test]
2438    fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
2439        let (tmp, talks) = store();
2440        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2441        let cfg = config(spec);
2442        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2443
2444        close(&mut talk, &talks).expect("close");
2445        assert_eq!(talk.status, TalkStatus::Closed);
2446        close(&mut talk, &talks).expect("closing twice is not an error");
2447
2448        let err =
2449            record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
2450        assert!(err.to_string().contains("closed"));
2451        let _ = &cfg; // config kept only to build the agent above
2452    }
2453
2454    #[tokio::test]
2455    async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2456        let (tmp, talks) = store();
2457        let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2458        let cfg = config(spec);
2459        // The in-flight turn's own handle: loaded once, the way a spawned
2460        // background task in `web::talk_say` holds one for the whole turn.
2461        let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2462
2463        // The operator closes the conversation through a *different* handle
2464        // while the turn above is still running - exactly what a close typed
2465        // on the phone while an agent is mid-answer looks like.
2466        let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
2467        close(&mut closed_elsewhere, &talks).expect("close");
2468        assert_eq!(
2469            talks.get(&in_flight.id).expect("reread").status,
2470            TalkStatus::Closed,
2471            "the close landed on disk before the turn finished"
2472        );
2473
2474        // The turn's own handle still says `open` - it was loaded before the
2475        // close - and finishing it must not resurrect the conversation the
2476        // operator already ended.
2477        assert_eq!(in_flight.status, TalkStatus::Open);
2478        respond(&mut in_flight, &talks, &cfg, "one more question")
2479            .await
2480            .expect("the turn itself still completes");
2481
2482        let on_disk = talks.get(&in_flight.id).expect("reread");
2483        assert_eq!(
2484            on_disk.status,
2485            TalkStatus::Closed,
2486            "a close must stick even when a turn that started before it finishes after it"
2487        );
2488        // The reply is not lost either: a turn already in flight when the
2489        // operator closed still gets its answer recorded.
2490        assert!(
2491            on_disk.turns.iter().any(|t| t.body == "here you go"),
2492            "the in-flight turn's own reply is still recorded: {:?}",
2493            on_disk.turns
2494        );
2495    }
2496
2497    #[test]
2498    fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
2499        let (tmp, talks) = store();
2500        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2501        let cfg = config(spec);
2502        // The handle `web::talk_say` would have read before awaiting config
2503        // discovery, then carried across that await into `record`.
2504        let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2505
2506        // The operator closes the conversation through a *different* handle
2507        // in the gap between that read and the call to `record` below.
2508        let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
2509        close(&mut closed_elsewhere, &talks).expect("close");
2510        assert_eq!(
2511            talks.get(&stale.id).expect("reread").status,
2512            TalkStatus::Closed,
2513            "the close landed on disk before record was called"
2514        );
2515
2516        // The stale handle still says `open` - it was loaded before the
2517        // close - so a `record` that trusted it would append a turn and
2518        // write the conversation back open, undoing the close.
2519        assert_eq!(stale.status, TalkStatus::Open);
2520        let err = record(&mut stale, &talks, "still there?", Vec::new())
2521            .expect_err("a close that landed first must be honored, not overwritten");
2522        assert!(err.to_string().contains("closed"));
2523
2524        let on_disk = talks.get(&stale.id).expect("reread");
2525        assert_eq!(
2526            on_disk.status,
2527            TalkStatus::Closed,
2528            "record must not resurrect a conversation closed while its snapshot was stale"
2529        );
2530        assert!(
2531            on_disk.turns.is_empty(),
2532            "the rejected turn must not have been appended: {:?}",
2533            on_disk.turns
2534        );
2535        let _ = &cfg; // config kept only to build the agent above
2536    }
2537
2538    #[test]
2539    fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
2540        let (tmp, talks) = store();
2541        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2542        let cfg = config(spec);
2543        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2544
2545        // Hold the same guard `record`'s read-modify-write section holds for
2546        // the whole of its own read-then-write, standing in for `record`
2547        // being paused between its read and its `put`.
2548        let held = talks.guard();
2549
2550        let talks2 = talks.clone();
2551        let id = talk.id.clone();
2552        let closing = std::thread::spawn(move || {
2553            let mut talk = talks2.get(&id).expect("get");
2554            close(&mut talk, &talks2).expect("close");
2555        });
2556
2557        std::thread::sleep(Duration::from_millis(50));
2558        assert!(
2559            !closing.is_finished(),
2560            "close must wait for the guard, not read and write while it is held - \
2561             a re-read alone narrows this window without closing it"
2562        );
2563
2564        drop(held);
2565        closing.join().expect("close thread panicked");
2566
2567        assert_eq!(
2568            talks.get(&talk.id).expect("reread").status,
2569            TalkStatus::Closed,
2570            "once the guard is free, close still lands"
2571        );
2572        let _ = &cfg; // config kept only to build the agent above
2573    }
2574
2575    #[test]
2576    fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
2577        let (tmp, talks) = store();
2578        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2579        let cfg = config(spec);
2580        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2581
2582        close(&mut talk, &talks).expect("close");
2583        assert_eq!(talk.status, TalkStatus::Closed);
2584
2585        reopen(&mut talk, &talks).expect("reopen");
2586        assert_eq!(talk.status, TalkStatus::Open);
2587        assert_eq!(
2588            talks.get(&talk.id).expect("reread").status,
2589            TalkStatus::Open
2590        );
2591
2592        // Idempotent: reopening an already-open talk is not an error.
2593        reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
2594        assert_eq!(talk.status, TalkStatus::Open);
2595
2596        record(&mut talk, &talks, "one more thing", Vec::new())
2597            .expect("a reopened talk takes turns again");
2598        let _ = &cfg; // config kept only to build the agent above
2599    }
2600
2601    #[test]
2602    fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
2603        let (tmp, talks) = store();
2604        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2605        let cfg = config(spec);
2606        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2607
2608        let artifacts = talks.artifacts_of(&talk.id);
2609        std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
2610        std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
2611
2612        talks.remove(&talk.id).expect("remove");
2613        assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
2614        assert!(!artifacts.is_dir(), "the artifacts directory is gone");
2615        assert!(
2616            talks.get(&talk.id).is_err(),
2617            "a removed talk cannot be read back"
2618        );
2619
2620        let err = talks
2621            .remove("nonexistent-id")
2622            .expect_err("unknown id refused");
2623        assert!(err.to_string().contains("no talk matches"), "{err}");
2624        let _ = &cfg; // config kept only to build the agent above
2625    }
2626
2627    #[tokio::test]
2628    async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2629        let (tmp, talks) = store();
2630        let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2631        let cfg = config(spec);
2632        // The in-flight turn's own handle, loaded before the delete lands -
2633        // the same shape as the matching close test above.
2634        let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2635
2636        talks.remove(&in_flight.id).expect("remove");
2637        assert!(
2638            talks.get(&in_flight.id).is_err(),
2639            "the delete landed on disk before the turn finished"
2640        );
2641
2642        // The turn's own handle has no way to know the record is gone -
2643        // finishing it must not write the file back into existence.
2644        respond(&mut in_flight, &talks, &cfg, "one more question")
2645            .await
2646            .expect("the turn itself still completes rather than erroring");
2647
2648        assert!(
2649            talks.get(&in_flight.id).is_err(),
2650            "a delete must stick even when a turn that started before it finishes after it"
2651        );
2652    }
2653
2654    #[test]
2655    fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
2656        let (tmp, talks) = store();
2657        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2658        let cfg = config(spec);
2659        // The handle `web::talk_say` would have read before awaiting config
2660        // discovery, then carried across that await into `record`.
2661        let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2662
2663        talks.remove(&stale.id).expect("remove");
2664
2665        // The stale handle has no way to know the record is gone - a
2666        // `record` that trusted it would append a turn and write the
2667        // conversation back into existence.
2668        let err = record(&mut stale, &talks, "still there?", Vec::new())
2669            .expect_err("a delete that landed first must be honored, not overwritten");
2670        assert!(err.to_string().contains("deleted"), "{err}");
2671
2672        assert!(
2673            talks.get(&stale.id).is_err(),
2674            "record must not resurrect a conversation deleted while its snapshot was stale"
2675        );
2676        let _ = &cfg; // config kept only to build the agent above
2677    }
2678
2679    #[test]
2680    fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
2681        let (tmp, talks) = store();
2682        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2683        let cfg = config(spec);
2684        // `web::talk_close` loads `talk` and calls `close` right after - this
2685        // stands in for a delete landing in that gap.
2686        let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2687
2688        talks.remove(&stale.id).expect("remove");
2689
2690        // The stale handle has no way to know the record is gone - a `close`
2691        // that fell back to it would write the conversation back into
2692        // existence, closed.
2693        let err = close(&mut stale, &talks)
2694            .expect_err("a delete that landed first must be honored, not overwritten");
2695        assert!(err.to_string().contains("deleted"), "{err}");
2696
2697        assert!(
2698            talks.get(&stale.id).is_err(),
2699            "close must not resurrect a conversation deleted while its snapshot was stale"
2700        );
2701        let _ = &cfg; // config kept only to build the agent above
2702    }
2703
2704    #[test]
2705    fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
2706        let (tmp, talks) = store();
2707        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2708        let cfg = config(spec);
2709        // `web::talk_reopen` loads `talk` and calls `reopen` right after -
2710        // this stands in for a delete landing in that gap.
2711        let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2712        close(&mut stale, &talks).expect("close");
2713
2714        talks.remove(&stale.id).expect("remove");
2715
2716        // The stale handle has no way to know the record is gone - a
2717        // `reopen` that fell back to it would write the conversation back
2718        // into existence, open.
2719        let err = reopen(&mut stale, &talks)
2720            .expect_err("a delete that landed first must be honored, not overwritten");
2721        assert!(err.to_string().contains("deleted"), "{err}");
2722
2723        assert!(
2724            talks.get(&stale.id).is_err(),
2725            "reopen must not resurrect a conversation deleted while its snapshot was stale"
2726        );
2727        let _ = &cfg; // config kept only to build the agent above
2728    }
2729
2730    #[test]
2731    fn list_puts_open_talks_before_closed_ones() {
2732        let (tmp, talks) = store();
2733        let make = |id: &str, status: TalkStatus| {
2734            let mut t = Talk {
2735                schema: SCHEMA,
2736                id: id.to_owned(),
2737                repo: tmp.path().to_owned(),
2738                agent: "mock".to_owned(),
2739                status,
2740                turns: Vec::new(),
2741                pending: String::new(),
2742                pending_attachments: Vec::new(),
2743                fallback: false,
2744                created_at: Timestamp::now(),
2745                updated_at: Timestamp::now(),
2746                seat: SeatState::new(SEAT, "mock", 7),
2747            };
2748            talks.put(&mut t).expect("put");
2749        };
2750        make("20260901-000000-0001", TalkStatus::Open);
2751        make("20260902-000000-0002", TalkStatus::Open);
2752        make("20260903-000000-0003", TalkStatus::Closed);
2753
2754        let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
2755        assert_eq!(
2756            ids,
2757            [
2758                "20260902-000000-0002",
2759                "20260901-000000-0001",
2760                "20260903-000000-0003"
2761            ]
2762        );
2763        assert_eq!(talks.count_open(), 2);
2764    }
2765
2766    #[test]
2767    fn tasks_of_finds_only_this_talks_own_tasks() {
2768        let dir = tempfile::tempdir().expect("tempdir");
2769        let queue = Queue::at(dir.path().join("queue"));
2770
2771        let mut mine = Task::new(
2772            "rework the loader".to_owned(),
2773            "rework the loader".to_owned(),
2774            PathBuf::from("/repo"),
2775            Source::Agent {
2776                run: "20260904-014455-ab12".to_owned(),
2777                node: "chat".to_owned(),
2778            },
2779        );
2780        queue.put(&mut mine).expect("put mine");
2781
2782        let mut theirs = Task::new(
2783            "unrelated".to_owned(),
2784            "unrelated".to_owned(),
2785            PathBuf::from("/repo"),
2786            Source::Agent {
2787                run: "20260904-090000-zz99".to_owned(),
2788                node: "implement".to_owned(),
2789            },
2790        );
2791        queue.put(&mut theirs).expect("put theirs");
2792
2793        let mut human = Task::new(
2794            "typed by hand".to_owned(),
2795            "typed by hand".to_owned(),
2796            PathBuf::from("/repo"),
2797            Source::Human,
2798        );
2799        queue.put(&mut human).expect("put human");
2800
2801        let found = tasks_of(&queue, "20260904-014455-ab12");
2802        assert_eq!(found.len(), 1);
2803        assert_eq!(found[0].id, mine.id);
2804    }
2805
2806    #[test]
2807    fn the_briefing_names_solo_task_add() {
2808        let brief = briefing(Path::new("/repo"), "en", false);
2809        assert!(brief.contains("magi task add --solo"));
2810        assert!(brief.contains("/repo"));
2811        assert!(!brief.contains("Hold this conversation in"));
2812    }
2813
2814    /// Talk fixes `repo` at the directory the conversation was opened in, so
2815    /// an agent asked to change some other checkout has no path to it unless
2816    /// the briefing itself says `--repo` can take a short name - see
2817    /// `resolve_repo_by_name` in `src/main.rs`, which is what actually
2818    /// resolves it.
2819    #[test]
2820    fn the_briefing_explains_targeting_a_different_repository_by_name() {
2821        let brief = briefing(Path::new("/repo"), "en", false);
2822        assert!(brief.contains("--repo does not have to be a full path"));
2823        assert!(brief.contains("owner/repo"));
2824        assert!(brief.contains("magi repos"));
2825        assert!(brief.contains("ask the operator"));
2826    }
2827
2828    #[test]
2829    fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
2830        let brief = briefing(Path::new("/repo"), "en", false);
2831        assert!(brief.contains("--attach <path>"), "{brief}");
2832        assert!(brief.contains("deleting this conversation"), "{brief}");
2833    }
2834
2835    #[test]
2836    fn the_briefing_names_the_language_when_it_is_not_english() {
2837        let brief = briefing(Path::new("/repo"), "Japanese", false);
2838        assert!(brief.contains("Hold this conversation in Japanese"));
2839    }
2840
2841    #[test]
2842    fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
2843        let read_only = briefing(Path::new("/repo"), "en", false);
2844        assert!(read_only.contains("Do not write files"));
2845        assert!(!read_only.contains("allow_write"));
2846
2847        let writable = briefing(Path::new("/repo"), "en", true);
2848        assert!(!writable.contains("Do not write files"));
2849        assert!(writable.contains("allow_write = true"));
2850        // Still names the queue for anything past a small named edit, and
2851        // still tells the agent to report what it changed.
2852        assert!(writable.contains("magi task add --solo"));
2853        assert!(writable.contains("say plainly what you"));
2854    }
2855}