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//! With `[talk] allow_write` on, a codex turn also runs unsandboxed
33//! ([`turn_access`]), as a claude turn already does under `bypassPermissions`:
34//! otherwise `magi task add` (writes to magi's queue under the home directory)
35//! or `git fetch` is refused by the sandbox. Unattended seats never get this.
36//!
37//! `--solo` rather than a plain `magi task add` is the point of pairing this
38//! module with [`crate::queue::Task::solo`]. A task that came out of a
39//! conversation the operator just had is a decision already made, not a
40//! design question worth three independent takes - so it runs through one
41//! implementer and straight into review, the way [`crate::graph::Runner`]
42//! already degrades a single-candidate run.
43//!
44//! # Shape
45//!
46//! The same split [`crate::queue`] uses: [`Talk`] is data plus pure helpers,
47//! [`Talks`] owns the I/O and is constructed with its root, so every test
48//! here drives a real store in a temp directory rather than the operator's
49//! own home.
50
51use std::path::{Path, PathBuf};
52use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
53use std::time::Duration;
54
55use anyhow::{Context, Result, bail};
56use jiff::Timestamp;
57use serde::{Deserialize, Serialize};
58
59use crate::agent::{self, Invocation, SeatState};
60use crate::config::{AgentSpec, Config};
61use crate::queue::{Queue, Source, Task};
62
63/// On-disk format for a conversation. Bumped when a field's meaning changes.
64pub const SCHEMA: u32 = 1;
65
66/// Wall-clock limit for one agent turn. See [`crate::config::Graph::timeout_talk`].
67///
68/// An hour by default. This turn is expected to run several shell commands
69/// and read their output before answering one - "what does this function
70/// do", "is this still true", "run the tests and tell me" - which argues for
71/// an hour rather than the five minutes a short budget once assumed, because
72/// the thing that made a short budget matter - an operator watching a
73/// spinner - is not how this conversation gets used: the operator moves on
74/// to something else while a turn runs and checks back later, so a long turn
75/// spends a held seat, not anyone's attention.
76fn turn_timeout(cfg: &Config) -> Duration {
77    Duration::from_secs(cfg.graph.timeout_talk)
78}
79
80/// Seat name for the conversation's agent, scoping its CLI-side session away
81/// from every other seat magi ever opens.
82const SEAT: &str = "talk";
83
84/// Prefix on a turn magi wrote rather than an agent.
85const MAGI_NOTE: &str = "magi: ";
86
87/// Who said something.
88#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(rename_all = "lowercase")]
90pub enum Who {
91    /// The operator.
92    Operator,
93    /// The conversation's agent - or magi itself, reporting that a turn
94    /// failed. See [`MAGI_NOTE`].
95    Agent,
96}
97
98/// One image the operator attached to a turn.
99///
100/// Never carries the bytes themselves: the picture lives on disk under
101/// [`Talks::attachments_dir`], named by `id` alone. `name` is the filename
102/// the operator's browser reported, kept only for display - it never
103/// contributes to a path, which is what keeps an upload from being able to
104/// traverse outside its own directory.
105#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
106#[serde(deny_unknown_fields)]
107pub struct Attachment {
108    /// Server-minted id; also the file's stem under `attachments_dir`.
109    pub id: String,
110    /// The operator's own filename, for display only.
111    pub name: String,
112    /// Validated by `web` at upload time against a closed whitelist:
113    /// `image/png`, `image/jpeg`, `image/gif`, `image/webp`.
114    pub mime: String,
115    /// Size in bytes, so the phone can show it without a second request.
116    pub bytes: u64,
117}
118
119/// One message in the conversation.
120#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
121#[serde(deny_unknown_fields)]
122pub struct Turn {
123    /// Who wrote it.
124    pub who: Who,
125    /// What they said.
126    pub body: String,
127    /// When it was said.
128    pub at: Timestamp,
129    /// Images attached to this turn. `#[serde(default)]` so a conversation
130    /// recorded before attachments existed still reads.
131    #[serde(default)]
132    pub attachments: Vec<Attachment>,
133    /// What the agent's CLI reported reading for this reply; only ever set on
134    /// an agent reply that carried usage. `#[serde(default)]` so a
135    /// conversation recorded before this field existed still reads, which is
136    /// why [`SCHEMA`] stays put: no existing field changed meaning.
137    #[serde(default, skip_serializing_if = "Option::is_none")]
138    pub usage: Option<TurnUsage>,
139}
140
141/// Raw usage of one agent reply, stored as the CLI reported it.
142///
143/// Counts only, never a percentage: the window is configuration
144/// ([`Config::context_window`]) and the conversation's model can change, so a
145/// stored percentage would go stale the moment either did. `agent` and
146/// `model` record who read that many tokens, which is how [`context_usage`]
147/// notices the figure predates a switch.
148#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
149pub struct TurnUsage {
150    /// Input-side tokens the CLI reported - see `agent::context_tokens`.
151    pub context_tokens: u64,
152    /// Roster id that answered.
153    pub agent: String,
154    /// That agent's model at the time (`None`: the CLI's own default).
155    #[serde(default)]
156    pub model: Option<String>,
157}
158
159/// How full the conversation's context window is, as the phone shows it.
160///
161/// Derived at read time and never persisted. `None` means unknown, which the
162/// UI says in words - it is never a made-up 0.
163#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
164pub struct ContextUsage {
165    /// Tokens the last agent reply's turn read, if its CLI reported any.
166    pub tokens: Option<u64>,
167    /// Window of the conversation's *current* model, if one is configured.
168    pub window: Option<u64>,
169    /// `tokens / window`, rounded down, not capped (over 100 is real news).
170    pub percent: Option<u64>,
171    /// 80% or more, judged before rounding: the conversation is getting long.
172    pub warn: bool,
173    /// The conversation's agent or model is not the one that produced
174    /// `tokens`: the figure describes the old session and the next turn
175    /// re-measures it.
176    pub since_switch: bool,
177    /// The current model, if the roster names one.
178    pub model: Option<String>,
179    /// `tokens` is a transcript-length estimate, not a CLI measurement (see
180    /// [`estimate_context_tokens`]). A measured figure is never overridden.
181    #[serde(default)]
182    pub estimated: bool,
183}
184
185/// Fraction (in percent) of the window at which [`ContextUsage::warn`] fires.
186const CONTEXT_WARN_PERCENT: u64 = 80;
187
188/// Allowance for the standing prompt when no config is readable to render
189/// the real [`briefing`].
190const STANDING_PROMPT_FALLBACK_CHARS: u64 = 7000;
191
192/// Rough estimate of the context a conversation occupies, for CLIs whose usage
193/// is cumulative and so is never measured (claude / codex / agy).
194///
195/// Counts the characters (not bytes) of every operator and agent turn, skips
196/// magi's own notes (never sent to the agent), adds `standing_chars` for the
197/// standing prompt, and divides by ~3.5 chars per token, rounding up. `None`
198/// when no turn counts: an empty conversation stays unknown.
199///
200/// Known biases: CJK runs near one token per character, and tool output,
201/// images and the CLI's own system prompt are not counted, so this tends to
202/// under-estimate (the 80% warning comes late); a compacted session is still
203/// counted in full, which over-estimates.
204pub fn estimate_context_tokens(talk: &Talk, standing_chars: u64) -> Option<u64> {
205    let mut counted = false;
206    let mut chars = standing_chars;
207    for t in talk.turns.iter().filter(|t| !t.body.starts_with(MAGI_NOTE)) {
208        counted = true;
209        chars += t.body.chars().count() as u64;
210    }
211    // chars / 3.5, rounded up, in integers.
212    counted.then(|| (chars * 2).div_ceil(7))
213}
214
215/// Context usage of `talk`, measured against its current model.
216///
217/// Only the latest agent reply counts (magi's own notes are skipped) and a
218/// reply without usage makes the answer unknown - older replies are never
219/// consulted, since a stale count passed off as current is worse than "unknown".
220/// The window comes from the *current* agent's model, so switching model moves
221/// the denominator at once.
222///
223/// Switching agent (or model, which is an agent change) mints a fresh CLI
224/// session, so the next turn re-sends the whole transcript and the count
225/// resets or jumps. Until that turn lands, the old figure is reported with
226/// `since_switch` set. Deterministic: same talk and config, same answer.
227pub fn context_usage(talk: &Talk, cfg: Option<&Config>) -> ContextUsage {
228    let current = cfg.and_then(|c| c.agents.iter().find(|a| a.id == talk.agent));
229    let model = current.and_then(|a| a.model.clone());
230    let window = cfg
231        .zip(model.as_deref())
232        .and_then(|(c, m)| c.context_window(m))
233        .filter(|w| *w > 0);
234    let usage = talk
235        .turns
236        .iter()
237        .rev()
238        .find(|t| t.who == Who::Agent && !t.body.starts_with(MAGI_NOTE))
239        .and_then(|t| t.usage.as_ref());
240    let measured = usage.map(|u| u.context_tokens);
241    let tokens = measured.or_else(|| {
242        let standing = cfg.map_or(STANDING_PROMPT_FALLBACK_CHARS, |c| {
243            briefing_with(
244                &talk.repo,
245                &c.graph.language,
246                c.talk.allow_write,
247                crate::persona::active(&c.talk.personas, &talk.persona).as_ref(),
248                c.talk.operator_name(),
249            )
250            .chars()
251            .count() as u64
252        });
253        estimate_context_tokens(talk, standing)
254    });
255    let estimated = measured.is_none() && tokens.is_some();
256    let since_switch =
257        usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
258    let (percent, warn) = match (tokens, window) {
259        (Some(t), Some(w)) => (
260            Some(t.saturating_mul(100) / w),
261            t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
262        ),
263        _ => (None, false),
264    };
265    ContextUsage {
266        tokens,
267        window,
268        percent,
269        warn,
270        since_switch,
271        model,
272        estimated,
273    }
274}
275
276/// Where a conversation is in its life: this conversation can file any
277/// number of tasks without ending, so it only ever moves once, from open to
278/// closed.
279#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
280#[serde(rename_all = "lowercase")]
281pub enum TalkStatus {
282    /// Still open; the operator may say more, and may have already filed work
283    /// out of it.
284    Open,
285    /// Closed by hand. Kept on disk as a record.
286    Closed,
287}
288
289impl TalkStatus {
290    /// Is this conversation still live?
291    pub fn open(self) -> bool {
292        matches!(self, Self::Open)
293    }
294
295    /// Wire form, for the phone and for logs.
296    pub fn as_str(self) -> &'static str {
297        match self {
298            Self::Open => "open",
299            Self::Closed => "closed",
300        }
301    }
302}
303
304/// One standing conversation.
305#[derive(Debug, Clone, Serialize, Deserialize)]
306#[serde(deny_unknown_fields)]
307pub struct Talk {
308    /// On-disk format version.
309    pub schema: u32,
310    /// Conversation id, e.g. `20260904-014455-ab12`.
311    pub id: String,
312    /// Repository this conversation is about.
313    pub repo: PathBuf,
314    /// Roster agent id holding the conversation.
315    pub agent: String,
316    /// Current state.
317    pub status: TalkStatus,
318    /// Everything said, oldest first.
319    pub turns: Vec<Turn>,
320    /// Text and attachments accepted while the single CLI turn is busy.
321    /// They are durable, but become a real turn only when [`drain`] records them.
322    #[serde(default)]
323    pub pending: String,
324    /// Attachments paired with [`Self::pending`].
325    #[serde(default)]
326    pub pending_attachments: Vec<Attachment>,
327    /// May a failed turn fall back through the rest of `[roles] chatter`?
328    /// True only while the agent was chosen by that chain; an explicit
329    /// `--agent` or an operator's switch pins the conversation to its agent
330    /// even when that agent also appears in the chain.
331    #[serde(default)]
332    pub fallback: bool,
333    /// Persona id (see [`crate::persona`]); empty is the plain default voice.
334    /// Independent of [`Self::agent`]: it survives an agent switch.
335    #[serde(default)]
336    pub persona: String,
337    /// The persona changed after the CLI session was last told about it. The
338    /// next resumed turn carries a persona update; cleared only once a turn
339    /// has succeeded (a fresh seat's briefing already holds the new persona).
340    #[serde(default)]
341    pub persona_dirty: bool,
342    /// When the conversation was opened.
343    pub created_at: Timestamp,
344    /// Last change to this file.
345    pub updated_at: Timestamp,
346    /// The CLI-side conversation, so a turn after the first costs one
347    /// sentence instead of the whole transcript. Not `pub`: it is magi's
348    /// bookkeeping, and a caller that edited it would detach the record from
349    /// the conversation the model actually holds.
350    seat: SeatState,
351}
352
353impl Talk {
354    /// Short form used in lists and notifications, matching a run's short id.
355    pub fn short(&self) -> &str {
356        short(&self.id)
357    }
358}
359
360/// A conversation store on disk.
361#[derive(Debug, Clone)]
362pub struct Talks {
363    root: PathBuf,
364    /// Serializes the read-modify-write cycle that reads a talk, decides
365    /// something from its `status`, and writes the whole record back.
366    /// [`close`], [`record`] and the tail of [`turn`] all take this before
367    /// that cycle rather than after just the read: a re-read narrows the
368    /// window another writer can land in, but does not close it, since
369    /// nothing stopped that other writer's own put from landing between this
370    /// call's re-read and its own put. Shared across every clone, since every
371    /// clone is a handle onto the same files.
372    lock: Arc<Mutex<()>>,
373}
374
375impl Talks {
376    /// The operator's conversations, `<home>/talks`.
377    pub fn open() -> Self {
378        Self::at(crate::run::home().join("talks"))
379    }
380
381    /// A store at an explicit root. Tests use this, which is why none of them
382    /// need the operator's real home.
383    pub fn at(root: PathBuf) -> Self {
384        Self {
385            root,
386            lock: Arc::new(Mutex::new(())),
387        }
388    }
389
390    /// Claim the right to read-modify-write a talk's `status`. A plain
391    /// `std::sync::Mutex`, not an async one: every caller holds it across a
392    /// handful of small file operations and never across an `.await`, so
393    /// blocking the thread briefly is the right tool, not a reason to reach
394    /// for `tokio::sync::Mutex`. Poisoning recovers rather than propagates -
395    /// one panicking caller must not wedge every talk in the store the way it
396    /// would wedge the loop's own lock; see [`crate::web`]'s `lock_or_recover`,
397    /// which this mirrors.
398    ///
399    /// The mutex only serializes this process. The cycle is also taken under a
400    /// short file lock beside the records, so the CLI (`magi answer
401    /// --ask-chat`) and the web server cannot overwrite each other's draft.
402    /// The file lock is never skipped: a cycle that cannot get it within
403    /// longer than the lock's own expiry fails instead of writing unguarded.
404    fn guard(&self) -> Result<StoreGuard<'_>> {
405        let mutex = self.lock.lock().unwrap_or_else(PoisonError::into_inner);
406        std::fs::create_dir_all(&self.root)
407            .with_context(|| format!("create {}", self.root.display()))?;
408        let path = self.root.join(".store.turn");
409        let deadline = std::time::Instant::now() + TAKEOVER_LOCK_TTL * 2;
410        loop {
411            if let Some(file) = TurnLock::take(&path)? {
412                return Ok(StoreGuard {
413                    _file: file,
414                    _mutex: mutex,
415                });
416            }
417            if std::time::Instant::now() >= deadline {
418                bail!("the talk store is locked by another process");
419            }
420            std::thread::sleep(Duration::from_millis(10));
421        }
422    }
423
424    /// Directory holding the conversation files.
425    pub fn root(&self) -> &Path {
426        &self.root
427    }
428
429    /// Path for one conversation id.
430    pub fn path_of(&self, id: &str) -> PathBuf {
431        self.root.join(format!("{id}.json"))
432    }
433
434    /// Where one conversation's prompts and CLI output are kept, beside the
435    /// record rather than inside it.
436    pub fn artifacts_of(&self, id: &str) -> PathBuf {
437        self.root.join(format!("{id}.artifacts"))
438    }
439
440    /// Where this conversation's attached images live: a subdirectory of
441    /// `artifacts_of`, so deleting the conversation deletes its attachments
442    /// too and nothing here needs its own cleanup path.
443    pub fn attachments_dir(&self, id: &str) -> PathBuf {
444        self.artifacts_of(id).join("attachments")
445    }
446
447    /// Persist one already-validated attachment and return its metadata.
448    ///
449    /// `web::talk_attachment_post` is the only caller: it has already
450    /// checked `mime` against the whitelist and sniffed the bytes, so an
451    /// unrecognised mime reaching here is a bug in that caller, not
452    /// something an operator did. The id is minted here and never taken
453    /// from the client; `name` is stored for display only and never used to
454    /// build a path.
455    pub fn put_attachment(
456        &self,
457        id: &str,
458        mime: &str,
459        name: &str,
460        data: &[u8],
461    ) -> Result<Attachment> {
462        let dir = self.attachments_dir(id);
463        std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
464        let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
465        let att = Attachment {
466            id: new_attachment_id(),
467            name: name.to_owned(),
468            mime: mime.to_owned(),
469            bytes: data.len() as u64,
470        };
471        std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
472            .with_context(|| format!("write attachment {}", att.id))?;
473        std::fs::write(
474            dir.join(format!("{}.json", att.id)),
475            serde_json::to_string(&att).context("serialize attachment")?,
476        )
477        .with_context(|| format!("write attachment metadata {}", att.id))?;
478        Ok(att)
479    }
480
481    /// Just the metadata, without reading the image bytes back off disk -
482    /// what `web::talk_say` uses to turn an id the operator referenced into
483    /// an [`Attachment`] before appending a [`Turn`], where the bytes
484    /// themselves are of no interest. `None` for an id this conversation
485    /// never stored - including one that merely looks plausible:
486    /// [`valid_attachment_id`] is checked here too, not only by the caller,
487    /// the same defence-in-depth `Questions::panel_asset` uses for its own
488    /// asset ids.
489    pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
490        if !valid_attachment_id(att_id) {
491            return Ok(None);
492        }
493        let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
494        if !meta_path.is_file() {
495            return Ok(None);
496        }
497        let att = serde_json::from_str(
498            &std::fs::read_to_string(&meta_path)
499                .with_context(|| format!("read {}", meta_path.display()))?,
500        )
501        .with_context(|| format!("parse {}", meta_path.display()))?;
502        Ok(Some(att))
503    }
504
505    /// A stored attachment's metadata and its bytes together, for serving it
506    /// back on `GET`. `None` under the same conditions as
507    /// [`Talks::attachment_meta`], which this is built on.
508    pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
509        let Some(att) = self.attachment_meta(id, att_id)? else {
510            return Ok(None);
511        };
512        let ext = attachment_ext(&att.mime).with_context(|| {
513            format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
514        })?;
515        let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
516        let data =
517            std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
518        Ok(Some((att, data)))
519    }
520
521    /// Absolute path of one attachment's bytes, for the prompt note [`turn`]
522    /// appends and for [`Invocation::attachments`]. `None` only for a mime
523    /// [`put_attachment`] could never have written, which means the
524    /// attachment did not come from this store.
525    ///
526    /// `self.root` (and so `attachments_dir`) is not guaranteed absolute on
527    /// its own - `run::home()` returns a bare relative `PathBuf` verbatim
528    /// when the operator sets `MAGI_HOME` to a relative path, and nothing
529    /// canonicalizes it on the way in. That is harmless for every other use
530    /// of this store, since its own I/O runs in this process against this
531    /// process's cwd - but this path is handed to a CLI invoked with `cwd:
532    /// &talk.repo`, a different directory, so a relative path here would
533    /// resolve against the wrong place once it reached the prompt.
534    /// `std::path::absolute` fixes it against *this* process's cwd before
535    /// that happens; see `disk::free_bytes_by_os` for the same function used
536    /// the same way elsewhere in this codebase.
537    fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
538        let ext = attachment_ext(&att.mime)?;
539        let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
540        std::path::absolute(&path).ok()
541    }
542
543    /// Write a conversation, atomically, so a process killed mid-write leaves
544    /// the previous state readable rather than a truncated file.
545    ///
546    /// The write-then-rename itself is retried a handful of times - see
547    /// [`write_atomic`] - because a reader with the destination file briefly
548    /// open is exactly the kind of failure that must not cost an agent's
549    /// whole reply; see `turn`'s own tail for what happens when even that is
550    /// not enough.
551    pub fn put(&self, t: &mut Talk) -> Result<()> {
552        std::fs::create_dir_all(&self.root)
553            .with_context(|| format!("create {}", self.root.display()))?;
554        t.updated_at = Timestamp::now();
555        let body = serde_json::to_string_pretty(t).context("serialize talk")?;
556        let path = self.path_of(&t.id);
557        let tmp = path.with_extension("json.tmp");
558        write_atomic(&tmp, &path, &body)
559    }
560
561    /// Load a conversation by id or unambiguous id prefix.
562    pub fn get(&self, id: &str) -> Result<Talk> {
563        let resolved = self.resolve_id(id)?;
564        read_path(&self.path_of(&resolved))
565    }
566
567    /// Every conversation on disk: open first, then newest first, so what the
568    /// operator is still using belongs above what they are done with.
569    pub fn list(&self) -> Vec<Talk> {
570        self.list_counting_unreadable().0
571    }
572
573    /// [`Talks::list`], plus how many `*.json` files could not be read, so a
574    /// caller that reports to the operator can say what it skipped.
575    pub fn list_counting_unreadable(&self) -> (Vec<Talk>, usize) {
576        let mut unreadable = 0;
577        let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
578            .into_iter()
579            .flatten()
580            .flatten()
581            .map(|e| e.path())
582            .filter(|p| p.extension().is_some_and(|x| x == "json"))
583            .filter_map(|p| {
584                let talk = read_path(&p).ok();
585                if talk.is_none() {
586                    unreadable += 1;
587                }
588                talk
589            })
590            .collect();
591        all.sort_unstable_by(|a, b| {
592            let rank = |t: &Talk| u8::from(!t.status.open());
593            rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
594        });
595        (all, unreadable)
596    }
597
598    /// Expand an id prefix to exactly one conversation id.
599    pub fn resolve_id(&self, prefix: &str) -> Result<String> {
600        if self.path_of(prefix).is_file() {
601            return Ok(prefix.to_owned());
602        }
603        let hits: Vec<String> = self
604            .list()
605            .into_iter()
606            .map(|t| t.id)
607            .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
608            .collect();
609        match hits.len() {
610            1 => Ok(hits.into_iter().next().expect("exactly one hit")),
611            0 => bail!("no talk matches `{prefix}`"),
612            _ => bail!(
613                "`{prefix}` matches {} talks: {}",
614                hits.len(),
615                hits.join(", ")
616            ),
617        }
618    }
619
620    /// Change detection token: the newest modification time in the store, in
621    /// milliseconds.
622    pub fn revision(&self) -> u64 {
623        std::fs::read_dir(&self.root)
624            .into_iter()
625            .flatten()
626            .flatten()
627            .filter(|e| e.path().extension().is_none_or(|x| x != "turn"))
628            .filter(|e| !e.file_name().to_string_lossy().starts_with('.'))
629            .filter_map(|e| e.metadata().ok())
630            .filter_map(|m| m.modified().ok())
631            .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
632            .map(|d| d.as_millis() as u64)
633            .max()
634            .unwrap_or(0)
635    }
636
637    /// How many conversations are still open.
638    pub fn count_open(&self) -> usize {
639        self.list().iter().filter(|t| t.status.open()).count()
640    }
641
642    /// Where one conversation's turn lease lives. Not a `*.json`, so
643    /// [`Talks::list`] never sees it.
644    pub fn turn_path(&self, id: &str) -> PathBuf {
645        self.root.join(format!("{id}.turn"))
646    }
647
648    /// Claim the right to run one agent turn in `id`, across processes.
649    ///
650    /// `None` means somebody else holds a fresh lease - the caller reports
651    /// "a turn is already running" and must not start one. A lease is held
652    /// while its last beat is within [`crate::ask::LEASE_TTL`]; pids are never
653    /// consulted. A stale (or unreadable) lease is taken over: renamed to a
654    /// unique name first, so of several takers only the one whose rename wins
655    /// deletes it, and the claim itself is an exclusive create.
656    pub fn claim_turn(&self, id: &str) -> Result<Option<TurnLease>> {
657        self.claim_turn_at(id, Timestamp::now())
658    }
659
660    fn claim_turn_at(&self, id: &str, now: Timestamp) -> Result<Option<TurnLease>> {
661        std::fs::create_dir_all(&self.root)
662            .with_context(|| format!("create {}", self.root.display()))?;
663        let path = self.turn_path(id);
664        let token = fresh_token();
665        if create_turn(&path, &token, now)? {
666            return Ok(Some(TurnLease {
667                talk: id.to_owned(),
668                path,
669                token,
670            }));
671        }
672        if read_turn(&path).is_some_and(|r| r.fresh(now)) {
673            return Ok(None);
674        }
675        // Stale or unreadable. Every replacement happens under a short-lived
676        // exclusive lock file, so two takers cannot each delete the other's
677        // fresh lease. A taker that finds the lock held simply loses; the lock
678        // itself ages out, so a taker that died inside it cannot wedge the talk.
679        let Some(_lock) = TurnLock::take(&path)? else {
680            return Ok(None);
681        };
682        if read_turn(&path).is_some_and(|r| r.fresh(now)) {
683            return Ok(None);
684        }
685        let _ = std::fs::remove_file(&path);
686        Ok(create_turn(&path, &token, now)?.then_some(TurnLease {
687            talk: id.to_owned(),
688            path,
689            token,
690        }))
691    }
692
693    /// Is a turn running in `id` anywhere, by a fresh lease?
694    pub fn turn_held(&self, id: &str) -> bool {
695        read_turn(&self.turn_path(id)).is_some_and(|r| r.fresh(Timestamp::now()))
696    }
697
698    /// Remove a conversation from disk, record and artifacts both. The
699    /// operator's way of saying "not just done, gone" - [`close`] alone
700    /// leaves the record as history.
701    ///
702    /// Takes [`Talks::guard`] for the same reason [`close`] does: a delete
703    /// racing a [`record`] or the tail of [`turn`] must not land between
704    /// their own read and write, or the file removed here would look, to
705    /// them, like a record that simply has not been written yet. The other
706    /// half of that story is on their side - both check under this same
707    /// guard that the record they are about to write is still there, and
708    /// give up without writing if it is not, which is what stops their `put`
709    /// from resurrecting a conversation this call already removed.
710    pub fn remove(&self, id: &str) -> Result<()> {
711        let _guard = self.guard()?;
712        let resolved = self.resolve_id(id)?;
713        let path = self.path_of(&resolved);
714        std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
715        let artifacts = self.artifacts_of(&resolved);
716        if artifacts.is_dir() {
717            std::fs::remove_dir_all(&artifacts)
718                .with_context(|| format!("remove {}", artifacts.display()))?;
719        }
720        let _ = std::fs::remove_file(self.turn_path(&resolved));
721        Ok(())
722    }
723}
724
725/// What [`Talks::guard`] hands out: the in-process mutex plus the
726/// cross-process file lock. The file lock is released first.
727struct StoreGuard<'a> {
728    _file: TurnLock,
729    _mutex: MutexGuard<'a, ()>,
730}
731
732/// The body of a `<id>.turn` file. `pid` is for a human reading it; nothing
733/// decides on it.
734#[derive(Debug, Serialize, Deserialize)]
735struct TurnRecord {
736    token: String,
737    pid: u32,
738    beat_at: Timestamp,
739}
740
741impl TurnRecord {
742    fn fresh(&self, now: Timestamp) -> bool {
743        now.as_second() - self.beat_at.as_second() <= crate::ask::LEASE_TTL.as_secs() as i64
744    }
745}
746
747/// A token no other caller in this process shares: `rng::entropy` is the
748/// clock and the pid, so two threads in one clock tick would otherwise get the
749/// same token, and with it the same temp file name and the same identity.
750fn fresh_token() -> String {
751    static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
752    let n = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
753    let seed = crate::rng::entropy() ^ n.wrapping_mul(0x9E37_79B9_7F4A_7C15);
754    crate::rng::SplitMix64::new(seed).uuid_v4()
755}
756
757fn read_turn(path: &Path) -> Option<TurnRecord> {
758    serde_json::from_str(&std::fs::read_to_string(path).ok()?).ok()
759}
760
761/// A takeover lock older than this belonged to a taker that died inside it.
762const TAKEOVER_LOCK_TTL: Duration = Duration::from_secs(10);
763
764/// A break ticket older than this belonged to a taker that died holding it;
765/// the next generation's ticket may then be tried for the same dead token.
766const TICKET_TTL: Duration = Duration::from_secs(10);
767
768/// Most ticket generations tried for one dead token.
769const TICKET_GENERATIONS: u32 = 16;
770
771/// Tickets are forgotten only this long after their last write.
772const TICKET_SWEEP_AGE: Duration = Duration::from_secs(3600);
773
774/// Create `path` exclusively with `body`; `false` when it already exists.
775fn create_exclusive(path: &Path, body: &str) -> Result<bool> {
776    use std::io::Write as _;
777    match std::fs::OpenOptions::new()
778        .write(true)
779        .create_new(true)
780        .open(path)
781    {
782        Ok(mut f) => {
783            if let Err(e) = f.write_all(body.as_bytes()) {
784                drop(f);
785                let _ = std::fs::remove_file(path);
786                return Err(e).with_context(|| format!("write {}", path.display()));
787            }
788            Ok(true)
789        }
790        Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
791        Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
792    }
793}
794
795fn create_turn(path: &Path, token: &str, now: Timestamp) -> Result<bool> {
796    let record = TurnRecord {
797        token: token.to_owned(),
798        pid: std::process::id(),
799        beat_at: now,
800    };
801    let body = serde_json::to_string(&record).context("serialize turn lease")?;
802    // Written in full under a private name, then linked into place: the link
803    // fails if the lease exists, and a reader never sees a half-written one.
804    let tmp = path.with_extension(format!("turn.{token}.new"));
805    std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
806    let linked = std::fs::hard_link(&tmp, path);
807    let _ = std::fs::remove_file(&tmp);
808    match linked {
809        Ok(()) => Ok(true),
810        Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
811        Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
812    }
813}
814
815/// The short exclusive lock every change to an existing lease (takeover, beat,
816/// release) happens under, so none of them can act on a stale reading. It
817/// ages out after [`TAKEOVER_LOCK_TTL`] in case its holder died inside it.
818struct TurnLock {
819    path: PathBuf,
820    token: String,
821}
822
823impl TurnLock {
824    fn token() -> String {
825        fresh_token()
826    }
827
828    /// Publish `token` at `path` complete (written privately, then linked), so
829    /// nobody reads an empty or half-written token; `false` if `path` exists.
830    fn publish(path: &Path, token: &str) -> Result<bool> {
831        let tmp = path.with_extension(format!("lock.{token}.new"));
832        std::fs::write(&tmp, token).with_context(|| format!("write {}", tmp.display()))?;
833        let linked = std::fs::hard_link(&tmp, path);
834        let _ = std::fs::remove_file(&tmp);
835        match linked {
836            Ok(()) => Ok(true),
837            Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
838            Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
839        }
840    }
841
842    fn take(lease: &Path) -> Result<Option<Self>> {
843        let path = lease.with_extension("turn.lock");
844        let token = Self::token();
845        if Self::publish(&path, &token)? {
846            return Ok(Some(Self { path, token }));
847        }
848        let Ok(seen) = std::fs::read_to_string(&path) else {
849            return Ok(None);
850        };
851        // A lock that is not a whole token (an older build wrote it in two
852        // steps) is breakable too, under one fixed ticket name.
853        let key = if !seen.is_empty() && seen.chars().all(|c| c.is_ascii_alphanumeric() || c == '-')
854        {
855            seen.as_str()
856        } else {
857            "invalid"
858        };
859        let aged = std::fs::metadata(&path)
860            .and_then(|m| m.modified())
861            .ok()
862            .and_then(|t| t.elapsed().ok())
863            .is_some_and(|age| age > TAKEOVER_LOCK_TTL);
864        if !aged {
865            return Ok(None);
866        }
867        // Breaking a dead lock is decided by a ticket named after the token
868        // that was judged dead: exactly one taker per generation can create
869        // it. The lock itself is never moved aside, so a fresh lock made by
870        // the winner is never off the path, not even for an instant.
871        //
872        // The ticket name carries a generation, not the wall clock, so no
873        // time boundary can let a second taker in while a ticket is alive.
874        // Generation n+1 is tried only when ticket n is older than
875        // `TICKET_TTL` (its owner died). Residual risk, same kind as the lock's
876        // own TTL: an owner stalled past `TICKET_TTL` between creating its
877        // ticket and publishing can overlap with the next generation's taker.
878        let mut won = false;
879        for n in 0..TICKET_GENERATIONS {
880            let ticket = path.with_extension(format!("lock.{key}.break.{n}"));
881            if create_exclusive(&ticket, "")? {
882                won = true;
883                break;
884            }
885            let stale = std::fs::metadata(&ticket)
886                .and_then(|m| m.modified())
887                .ok()
888                .and_then(|t| t.elapsed().ok())
889                .is_some_and(|age| age > TICKET_TTL);
890            if !stale {
891                return Ok(None);
892            }
893        }
894        if !won {
895            // Every generation was abandoned. A ticket is never unlinked by
896            // name while it may still be someone's exclusion: only the sweep
897            // removes tickets, and only ones far older than any taker lives,
898            // so recovery resumes once they age out.
899            Self::sweep_tickets(&path);
900            return Ok(None);
901        }
902        Self::sweep_tickets(&path);
903        // Only the ticket's owner reaches this point for `seen`, and nobody
904        // else removes a lock that still carries it; if it changed, leave it.
905        if std::fs::read_to_string(&path).ok().as_deref() != Some(seen.as_str()) {
906            return Ok(None);
907        }
908        let _ = std::fs::remove_file(&path);
909        if Self::publish(&path, &token)? {
910            return Ok(Some(Self { path, token }));
911        }
912        Ok(None)
913    }
914
915    /// Forget tickets far older than the lock's TTL. They stay that long so a
916    /// slow taker that read the same dead token cannot break it a second time.
917    fn sweep_tickets(path: &Path) {
918        let (Some(dir), Some(name)) = (path.parent(), path.file_name().and_then(|n| n.to_str()))
919        else {
920            return;
921        };
922        let prefix = format!("{name}.");
923        let Ok(entries) = std::fs::read_dir(dir) else {
924            return;
925        };
926        for entry in entries.flatten() {
927            let file = entry.file_name();
928            let Some(file) = file.to_str() else { continue };
929            if !(file.starts_with(&prefix) && file.contains(".break.")) {
930                continue;
931            }
932            let old = entry
933                .metadata()
934                .and_then(|m| m.modified())
935                .ok()
936                .and_then(|t| t.elapsed().ok())
937                .is_some_and(|age| age > TICKET_SWEEP_AGE);
938            if old {
939                let _ = std::fs::remove_file(entry.path());
940            }
941        }
942    }
943
944    /// Wait briefly for the lock; the holders only do a read and a write.
945    fn take_patiently(lease: &Path) -> Option<Self> {
946        for _ in 0..50 {
947            match Self::take(lease) {
948                Ok(Some(lock)) => return Some(lock),
949                Ok(None) => std::thread::sleep(Duration::from_millis(10)),
950                Err(_) => return None,
951            }
952        }
953        None
954    }
955}
956
957impl Drop for TurnLock {
958    fn drop(&mut self) {
959        // Only our own lock: one that aged out and was taken over is not ours.
960        if std::fs::read_to_string(&self.path).is_ok_and(|t| t == self.token) {
961            let _ = std::fs::remove_file(&self.path);
962        }
963    }
964}
965
966/// How often a running turn renews its lease: well inside
967/// [`crate::ask::LEASE_TTL`].
968pub const TURN_BEAT: Duration = Duration::from_secs(20);
969
970/// One talk's cross-process turn slot, released on drop (success, error,
971/// panic or a dropped handler future alike). Release only removes the file
972/// while it still carries this lease's token, so a guard that outlived its
973/// own expiry cannot delete the lease of whoever took over.
974#[derive(Debug)]
975pub struct TurnLease {
976    talk: String,
977    path: PathBuf,
978    token: String,
979}
980
981impl TurnLease {
982    /// Does the lease file still carry this lease's token? A read only.
983    pub fn holds(&self) -> bool {
984        read_turn(&self.path).is_some_and(|r| r.token == self.token)
985    }
986
987    /// Renew the lease. `Ok(false)` means it was taken over or removed, so
988    /// this turn no longer owns the slot; `Err` is a transient failure (the
989    /// lock stayed busy, a write failed) and the next beat tries again.
990    pub fn beat(&self) -> Result<bool> {
991        let _lock = TurnLock::take_patiently(&self.path)
992            .with_context(|| format!("lock {} to renew it", self.path.display()))?;
993        let Some(mut record) = read_turn(&self.path).filter(|r| r.token == self.token) else {
994            return Ok(false);
995        };
996        record.beat_at = Timestamp::now();
997        let body = serde_json::to_string(&record).context("serialize turn lease")?;
998        let tmp = self.path.with_extension(format!("turn.{}.tmp", self.token));
999        write_atomic(&tmp, &self.path, &body)?;
1000        Ok(true)
1001    }
1002
1003    /// Run `fut` while renewing this lease every [`TURN_BEAT`]. A transient
1004    /// beat failure is logged and the turn goes on; a lost lease drops `fut`
1005    /// (the turn stops) and is an error, since somebody else may now be
1006    /// running the same conversation.
1007    pub async fn beating<T>(&self, fut: impl std::future::Future<Output = T>) -> Result<T> {
1008        self.beating_every(TURN_BEAT, fut).await
1009    }
1010
1011    async fn beating_every<T>(
1012        &self,
1013        period: Duration,
1014        fut: impl std::future::Future<Output = T>,
1015    ) -> Result<T> {
1016        tokio::pin!(fut);
1017        loop {
1018            match tokio::time::timeout(period, &mut fut).await {
1019                Ok(out) => return Ok(out),
1020                Err(_) => match self.beat() {
1021                    Ok(true) => {}
1022                    Ok(false) => bail!(
1023                        "the turn lease {} was taken over; this turn is stopped",
1024                        self.path.display()
1025                    ),
1026                    Err(e) => tracing::warn!("{e:#}"),
1027                },
1028            }
1029        }
1030    }
1031}
1032
1033impl Drop for TurnLease {
1034    fn drop(&mut self) {
1035        // Under the lock, so a takeover cannot slip in between the check and
1036        // the removal. If the lock stays busy the lease just ages out.
1037        if let Some(_lock) = TurnLock::take_patiently(&self.path) {
1038            if read_turn(&self.path).is_some_and(|r| r.token == self.token) {
1039                let _ = std::fs::remove_file(&self.path);
1040            }
1041        }
1042    }
1043}
1044
1045/// Open a conversation. Takes no agent turn: there is no idea to answer yet,
1046/// and a conversation the operator has not said anything into yet is a
1047/// normal, valid thing to have sitting on the phone.
1048///
1049/// `agent` beats `[roles] chatter`, which beats [`agent::pick`]'s own default
1050/// order (a claude seat, else the first runnable agent in roster order) when
1051/// nothing names a seat at all - see `[roles] chatter`'s own doc in
1052/// [`crate::config`] for why a dedicated field exists rather than reusing a
1053/// judge seat.
1054pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
1055    // Absolute: a relative path means the wrong repository once anything
1056    // other than this process reads it back.
1057    let repo = repo.canonicalize().unwrap_or(repo);
1058    // An explicit agent is a chain of one; otherwise the first id of
1059    // `[roles] chatter` that can run here (later ones are `turn`'s fallbacks).
1060    let spec = match agent {
1061        Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
1062        None => agent::pick_chain(
1063            &cfg.agents,
1064            cfg.roles.chatter.as_ref(),
1065            &agent::installed,
1066            "chatter",
1067        )?
1068        .remove(0),
1069    };
1070
1071    let now = Timestamp::now();
1072    let mut talk = Talk {
1073        schema: SCHEMA,
1074        id: new_id(),
1075        repo,
1076        agent: spec.id.clone(),
1077        status: TalkStatus::Open,
1078        turns: Vec::new(),
1079        pending: String::new(),
1080        pending_attachments: Vec::new(),
1081        fallback: agent.is_none(),
1082        persona: String::new(),
1083        persona_dirty: false,
1084        created_at: now,
1085        updated_at: now,
1086        seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
1087    };
1088    store.put(&mut talk)?;
1089    Ok(talk)
1090}
1091
1092/// Append the operator's turn and flush it, without invoking anything.
1093///
1094/// Split out of [`say`] so `POST /api/talks/{id}/say` can answer once the
1095/// message is safely on disk, and run the agent's half in the background -
1096/// holding the connection for a turn that can run fifteen minutes is the
1097/// wrong shape for a phone.
1098pub fn record(
1099    talk: &mut Talk,
1100    store: &Talks,
1101    text: &str,
1102    attachments: Vec<Attachment>,
1103) -> Result<String> {
1104    // `web::talk_say` reads the talk, then awaits config discovery before
1105    // calling this - a gap a concurrent `POST /api/talks/{id}/close` can land
1106    // in. The guard held for the rest of this function is what actually closes
1107    // that gap: re-reading status without it only shrinks the window a
1108    // concurrent `close` could land in between this call's own read and its
1109    // `put`, it does not remove it. See [`Talks::guard`] and the matching
1110    // guard in `turn`, which this mirrors.
1111    let _guard = store.guard()?;
1112    // A concurrent `Talks::remove` can have landed in that same gap. `put`
1113    // writes unconditionally, so trusting the stale `talk` here would recreate
1114    // the file a delete just removed - the record must still be there for a
1115    // turn to have anywhere to append to.
1116    let Ok(fresh) = store.get(&talk.id) else {
1117        bail!("talk {} was deleted", talk.short());
1118    };
1119    talk.status = fresh.status;
1120    // Do not let this older handle overwrite a draft accepted while it was
1121    // waiting for configuration discovery.
1122    talk.pending = fresh.pending;
1123    talk.pending_attachments = fresh.pending_attachments;
1124    if !talk.status.open() {
1125        bail!(
1126            "talk {} is {} and takes no more turns",
1127            talk.short(),
1128            talk.status.as_str()
1129        );
1130    }
1131    let text = text.trim();
1132    if text.is_empty() && attachments.is_empty() {
1133        bail!("nothing to say");
1134    }
1135    talk.turns.push(Turn {
1136        who: Who::Operator,
1137        body: text.to_owned(),
1138        at: Timestamp::now(),
1139        attachments,
1140        usage: None,
1141    });
1142    store.put(talk)?;
1143    Ok(text.to_owned())
1144}
1145
1146/// Add an unrecorded message to the durable draft while another turn runs.
1147pub fn queue(
1148    talk: &mut Talk,
1149    store: &Talks,
1150    text: &str,
1151    attachments: Vec<Attachment>,
1152) -> Result<()> {
1153    let text = text.trim();
1154    if text.is_empty() && attachments.is_empty() {
1155        bail!("nothing to say");
1156    }
1157    let _guard = store.guard()?;
1158    let mut fresh = store
1159        .get(&talk.id)
1160        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1161    if !fresh.status.open() {
1162        bail!(
1163            "talk {} is {} and takes no more turns",
1164            fresh.short(),
1165            fresh.status.as_str()
1166        );
1167    }
1168    if !text.is_empty() {
1169        if fresh.pending.is_empty() {
1170            fresh.pending = text.to_owned();
1171        } else {
1172            fresh.pending.push_str("\n\n");
1173            fresh.pending.push_str(text);
1174        }
1175    }
1176    fresh.pending_attachments.extend(attachments);
1177    store.put(&mut fresh)?;
1178    *talk = fresh;
1179    Ok(())
1180}
1181
1182/// Promote the current durable draft to one operator turn.
1183pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
1184    let _guard = store.guard()?;
1185    let mut fresh = store
1186        .get(&talk.id)
1187        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1188    if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
1189        *talk = fresh;
1190        return Ok(None);
1191    }
1192    let text = std::mem::take(&mut fresh.pending);
1193    let attachments = std::mem::take(&mut fresh.pending_attachments);
1194    fresh.turns.push(Turn {
1195        who: Who::Operator,
1196        body: text.clone(),
1197        at: Timestamp::now(),
1198        attachments,
1199        usage: None,
1200    });
1201    store.put(&mut fresh)?;
1202    *talk = fresh;
1203    Ok(Some(text))
1204}
1205
1206/// One operator turn and one agent turn, appended - the synchronous form, used
1207/// by tests and by anything that is fine waiting out the turn itself.
1208pub async fn say(
1209    lease: &TurnLease,
1210    talk: &mut Talk,
1211    store: &Talks,
1212    cfg: &Config,
1213    text: &str,
1214    attachments: Vec<Attachment>,
1215) -> Result<()> {
1216    check_lease(lease, talk)?;
1217    let text = record(talk, store, text, attachments)?;
1218    respond(lease, talk, store, cfg, &text).await
1219}
1220
1221/// The agent's half of a turn: invoke, append, flush. Pairs with [`record`].
1222///
1223/// Requires the talk's [`TurnLease`] (from [`Talks::claim_turn`]), so no
1224/// starter can run a turn without holding the cross-process slot. The lease is
1225/// renewed while the turn runs; a lease lost mid-turn stops it with an error.
1226pub async fn respond(
1227    lease: &TurnLease,
1228    talk: &mut Talk,
1229    store: &Talks,
1230    cfg: &Config,
1231    text: &str,
1232) -> Result<()> {
1233    check_lease(lease, talk)?;
1234    lease
1235        .beating(turn(talk, store, cfg, text))
1236        .await
1237        .and_then(|done| done)
1238}
1239
1240fn check_lease(lease: &TurnLease, talk: &Talk) -> Result<()> {
1241    if lease.talk != talk.id {
1242        bail!("the turn lease is for talk {}, not {}", lease.talk, talk.id);
1243    }
1244    // A lease that aged out and was taken over must not start a turn (or
1245    // write the operator's text) at all.
1246    if !lease.holds() {
1247        bail!("the turn lease for talk {} is no longer held", talk.short());
1248    }
1249    Ok(())
1250}
1251
1252/// Close a conversation. Idempotent: closing an already-closed conversation is
1253/// not an error, since the operator's intent - "I am done with this" - is
1254/// already satisfied.
1255///
1256/// Re-reads the record under [`Talks::guard`] rather than trusting the
1257/// caller's copy of `talk`, and writes that fresh copy back rather than the
1258/// one passed in. `web::talk_close` loads `talk` and calls this right after
1259/// with no gap of its own, but without the guard that load can still land
1260/// between a `record` or `turn` elsewhere reading the file and writing it
1261/// back - and a close built on the older snapshot would put it right back,
1262/// silently dropping whatever turn the other call had just appended.
1263///
1264/// If the re-read fails, this errors rather than falling back to the
1265/// caller's stale copy: `talk::begin` always `put`s the record before handing
1266/// out a `Talk`, so the only way a re-read can fail is a concurrent
1267/// [`Talks::remove`] having deleted it, and writing the stale copy back would
1268/// resurrect exactly what that delete removed.
1269pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
1270    let _guard = store.guard()?;
1271    let mut fresh = store
1272        .get(&talk.id)
1273        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1274    fresh.status = TalkStatus::Closed;
1275    // A closed conversation must not replay a draft if it is reopened later.
1276    fresh.pending.clear();
1277    fresh.pending_attachments.clear();
1278    store.put(&mut fresh)?;
1279    *talk = fresh;
1280    Ok(())
1281}
1282
1283/// Reopen a closed conversation. Idempotent for the same reason [`close`] is:
1284/// reopening an already-open conversation is not an error, since the
1285/// operator's intent - "I want to keep talking about this" - is already
1286/// satisfied.
1287///
1288/// Written symmetrically with [`close`]: re-reads the record under
1289/// [`Talks::guard`] rather than trusting the caller's copy of `talk`, writes
1290/// that fresh copy back rather than the one passed in, and errors rather than
1291/// falling back to the stale copy if the re-read fails, for the same reasons
1292/// `close`'s doc gives.
1293pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
1294    let _guard = store.guard()?;
1295    let mut fresh = store
1296        .get(&talk.id)
1297        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1298    fresh.status = TalkStatus::Open;
1299    store.put(&mut fresh)?;
1300    *talk = fresh;
1301    Ok(())
1302}
1303
1304/// Hand the conversation to another roster agent.
1305///
1306/// A CLI session belongs to one CLI and cannot be carried to another, so the
1307/// seat is minted afresh rather than edited: the next turn finds
1308/// `seat.turns == 0` and re-sends the transcript, since the new agent has
1309/// heard none of it. A magi-written note records the change in the
1310/// transcript. Returns `false` (and writes nothing, not even a note) when the
1311/// stored talk already uses `spec`.
1312///
1313/// Re-reads under [`Talks::guard`], like [`close`], and errors rather than
1314/// resurrecting a record a concurrent delete removed. Refusing a closed talk
1315/// or a turn in flight is the caller's job: only it can see the latter.
1316pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
1317    let _guard = store.guard()?;
1318    let mut fresh = store
1319        .get(&talk.id)
1320        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1321    if fresh.agent == spec.id {
1322        *talk = fresh;
1323        return Ok(false);
1324    }
1325    let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
1326    fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1327    // A deliberate switch pins the conversation to the agent chosen.
1328    fresh.fallback = false;
1329    fresh.turns.push(Turn {
1330        who: Who::Agent,
1331        body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
1332        at: Timestamp::now(),
1333        attachments: Vec::new(),
1334        usage: None,
1335    });
1336    store.put(&mut fresh)?;
1337    *talk = fresh;
1338    Ok(true)
1339}
1340
1341/// Pick the conversation's persona (`""` or `default` is the plain voice).
1342///
1343/// The CLI session keeps its seat: the change is recorded with
1344/// [`Talk::persona_dirty`] and a magi note, and the next turn tells the session
1345/// about it (or its fresh briefing already does). Returns `false` and writes
1346/// nothing when the stored persona is already `id`. Refusing a closed talk, an
1347/// unknown id or a turn in flight is the caller's job.
1348pub fn switch_persona(talk: &mut Talk, store: &Talks, id: &str) -> Result<bool> {
1349    let id = if id.trim() == crate::persona::DEFAULT_ID {
1350        ""
1351    } else {
1352        id.trim()
1353    };
1354    let _guard = store.guard()?;
1355    let mut fresh = store
1356        .get(&talk.id)
1357        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1358    if fresh.persona == id {
1359        *talk = fresh;
1360        return Ok(false);
1361    }
1362    let from = std::mem::replace(&mut fresh.persona, id.to_owned());
1363    fresh.persona_dirty = true;
1364    let label = |p: &str| {
1365        if p.is_empty() {
1366            crate::persona::DEFAULT_ID.to_owned()
1367        } else {
1368            p.to_owned()
1369        }
1370    };
1371    fresh.turns.push(Turn {
1372        who: Who::Agent,
1373        body: format!(
1374            "{MAGI_NOTE}persona changed from {} to {}",
1375            label(&from),
1376            label(id)
1377        ),
1378        at: Timestamp::now(),
1379        attachments: Vec::new(),
1380        usage: None,
1381    });
1382    store.put(&mut fresh)?;
1383    *talk = fresh;
1384    Ok(true)
1385}
1386
1387/// Discard the durable draft without adding a transcript turn.
1388pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
1389    let _guard = store.guard()?;
1390    let mut fresh = store
1391        .get(&talk.id)
1392        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1393    fresh.pending.clear();
1394    fresh.pending_attachments.clear();
1395    store.put(&mut fresh)?;
1396    *talk = fresh;
1397    Ok(())
1398}
1399
1400/// Clear a draft only when the caller still sees its complete snapshot.
1401pub fn clear_pending_if_matches(
1402    talk: &mut Talk,
1403    store: &Talks,
1404    expected_text: &str,
1405    expected_attachments: &[String],
1406) -> Result<bool> {
1407    let _guard = store.guard()?;
1408    let mut fresh = store
1409        .get(&talk.id)
1410        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1411    if !pending_matches(&fresh, expected_text, expected_attachments) {
1412        *talk = fresh;
1413        return Ok(false);
1414    }
1415    fresh.pending.clear();
1416    fresh.pending_attachments.clear();
1417    store.put(&mut fresh)?;
1418    *talk = fresh;
1419    Ok(true)
1420}
1421
1422/// Replace just the text of the durable draft, but only if the caller's
1423/// snapshot still identifies the entire draft. This refuses to overwrite a
1424/// message another client queued or a draft the drain already promoted.
1425pub fn edit_pending_text(
1426    talk: &mut Talk,
1427    store: &Talks,
1428    text: &str,
1429    expected_text: &str,
1430    expected_attachments: &[String],
1431) -> Result<bool> {
1432    let _guard = store.guard()?;
1433    let mut fresh = store
1434        .get(&talk.id)
1435        .with_context(|| format!("talk {} was deleted", talk.short()))?;
1436    if !pending_matches(&fresh, expected_text, expected_attachments) {
1437        *talk = fresh;
1438        return Ok(false);
1439    }
1440    fresh.pending = text.trim().to_owned();
1441    store.put(&mut fresh)?;
1442    *talk = fresh;
1443    Ok(true)
1444}
1445
1446fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
1447    talk.pending == expected_text
1448        && talk
1449            .pending_attachments
1450            .iter()
1451            .map(|attachment| &attachment.id)
1452            .eq(expected_attachments.iter())
1453}
1454
1455/// Invoke the conversation's agent once and append what it said.
1456///
1457/// The first turn ever taken carries the full [`briefing`], because nothing
1458/// else has told the agent what this conversation is or what it may do.
1459/// Every turn after that resends nothing when the CLI can resume its own
1460/// session, and falls back to [`transcript`] only when it cannot.
1461async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1462    let spec = cfg
1463        .agents
1464        .iter()
1465        .find(|a| a.id == talk.agent)
1466        .with_context(|| {
1467            format!(
1468                "talk {} was opened with agent `{}`, which is no longer in \
1469                 the roster; restore it in magi.toml or start a new \
1470                 conversation",
1471                talk.short(),
1472                talk.agent
1473            )
1474        })?;
1475
1476    // The newest turn is always the operator message this call is answering
1477    // - `record` appended it before `turn` was ever called - so its own
1478    // attachments are what belong at the end of *this* prompt.
1479    let last_note = attachment_note(
1480        store,
1481        &talk.id,
1482        talk.turns
1483            .last()
1484            .map_or(&[][..], |t| t.attachments.as_slice()),
1485    );
1486
1487    // Every attachment this conversation has ever held, not only this
1488    // turn's: a resumed session gets a fresh process every turn, so a CLI
1489    // whose sandbox needs `--add-dir` (see `agent::build_command`) needs the
1490    // grant again to open an image from an earlier turn, even when nothing
1491    // new was attached just now.
1492    let attachment_paths: Vec<PathBuf> = talk
1493        .turns
1494        .iter()
1495        .flat_map(|t| t.attachments.iter())
1496        .filter_map(|a| store.attachment_path(&talk.id, a))
1497        .collect();
1498
1499    // A question handed to this conversation is answered with `magi answer`,
1500    // which writes the question store; a read-only sandbox refuses that. So
1501    // while such a question is still open - the hand-over turn and the turns
1502    // in which the owner decides - the turn may write, with the question store
1503    // as a writable root. It goes back to read-only once the question is
1504    // answered or abandoned. As for the deputy, what the seat may touch beyond
1505    // that rests on the prompt, not the sandbox.
1506    let questions = crate::ask::Questions::open();
1507    let consulted = crate::consult::pending_consults(&questions, &talk.id);
1508    let consult_roots: Vec<PathBuf> = if consulted {
1509        vec![questions.root().to_path_buf()]
1510    } else {
1511        Vec::new()
1512    };
1513
1514    let (allow_write, unsandboxed) = turn_access(cfg.talk.allow_write, consulted);
1515
1516    let artifacts = store.artifacts_of(&talk.id);
1517    // From the transcript, not the seat: a switched agent's seat restarts at
1518    // zero and must not overwrite an earlier turn's artifacts.
1519    let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1520    let stem = format!("turn-{}", operator_turns.max(1));
1521    // The chat's build cache is the same shared one the graph's seats get, so
1522    // a conversation that compiles does not mint another multi-GB target dir.
1523    let cache_dir = cfg.cache_dir();
1524
1525    // The agent holding the conversation, then - only when it came from
1526    // `[roles] chatter` - the rest of that chain, each at most once.
1527    let mut chain = vec![spec.clone()];
1528    if let Some(choice) = cfg.roles.chatter.as_ref()
1529        && talk.fallback
1530    {
1531        for id in choice.ids() {
1532            if id == talk.agent || chain.iter().any(|s| s.id == id) {
1533                continue;
1534            }
1535            match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1536                Ok(s) => chain.push(s),
1537                Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1538            }
1539        }
1540    }
1541
1542    let persona = crate::persona::active(&cfg.talk.personas, &talk.persona);
1543    let operator_name = cfg.talk.operator_name();
1544    let persona_update = if talk.persona_dirty {
1545        format!(
1546            "{}\n\n",
1547            crate::persona::update_block_for(persona.as_ref(), operator_name)
1548        )
1549    } else {
1550        String::new()
1551    };
1552
1553    let mut outcome = None;
1554    let mut fell_back_from: Option<String> = None;
1555    // What the conversation looked like after the first agent's failed try,
1556    // so an exhausted chain leaves exactly what a single failed seat would.
1557    let mut first_try: Option<(String, SeatState)> = None;
1558    for (n, spec) in chain.iter().enumerate() {
1559        if n > 0 {
1560            if first_try.is_none() {
1561                first_try = Some((talk.agent.clone(), talk.seat.clone()));
1562            }
1563            tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1564            // A new CLI has none of the old one's conversation: a fresh seat
1565            // puts `has_session` at false and the full transcript is re-sent.
1566            fell_back_from.get_or_insert_with(|| talk.agent.clone());
1567            talk.agent = spec.id.clone();
1568            talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1569        }
1570        let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1571        let first_ever = talk.turns.len() <= 1;
1572        let body = if talk.seat.turns == 0 && first_ever {
1573            format!(
1574                "{}\n\n# Operator\n\n{text}{last_note}",
1575                briefing_with(
1576                    &talk.repo,
1577                    &cfg.graph.language,
1578                    cfg.talk.allow_write,
1579                    persona.as_ref(),
1580                    operator_name
1581                )
1582            )
1583        } else if talk.seat.turns == 0 {
1584            // A fresh seat on a conversation that already has history (the
1585            // agent was switched): the briefing, then everything said so far.
1586            format!(
1587                "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1588                briefing_with(
1589                    &talk.repo,
1590                    &cfg.graph.language,
1591                    cfg.talk.allow_write,
1592                    persona.as_ref(),
1593                    operator_name
1594                ),
1595                transcript(talk, store)
1596            )
1597        } else if resuming {
1598            format!("{persona_update}{text}{last_note}")
1599        } else {
1600            // No session to hold the persona: the transcript never stores it,
1601            // so a non-default persona is re-sent on every such turn (the
1602            // update block already carries it when the choice just changed).
1603            let mut standing = match (&persona, talk.persona_dirty) {
1604                (Some(p), false) => format!("{}\n", crate::persona::section_for(p, operator_name)),
1605                _ => persona_update.clone(),
1606            };
1607            // The default voice has no persona section to carry the name, and
1608            // the transcript never stores the briefing's addressing line.
1609            if let (None, Some(n)) = (&persona, operator_name) {
1610                standing.push_str(&format!(
1611                    "# Addressing the operator\n{}\n",
1612                    crate::persona::addressing(n)
1613                ));
1614            }
1615            format!("{}\n\n{standing}{text}{last_note}", transcript(talk, store))
1616        };
1617        let attempt_stem = if n == 0 {
1618            stem.clone()
1619        } else {
1620            format!("{stem}-{}", spec.id)
1621        };
1622        let inv = Invocation {
1623            cwd: &talk.repo,
1624            prompt: &body,
1625            timeout: turn_timeout(cfg),
1626            // Off unless this repository's own config opts in - see
1627            // `crate::config::Talk::allow_write` and this module's doc for why
1628            // the default keeps a conversational edit from landing in a checkout
1629            // no run or review can claim.
1630            allow_write,
1631            unsandboxed,
1632            sessions: cfg.graph.sessions,
1633            artifacts: &artifacts,
1634            stem: &attempt_stem,
1635            // The conversation's own id, so `magi task add` run from inside it is
1636            // attributed to this conversation - see `Source::Agent`.
1637            run: &talk.id,
1638            node: crate::queue::CHAT_NODE,
1639            cache_dir: cache_dir.as_deref(),
1640            attachments: &attachment_paths,
1641            writable: &consult_roots,
1642        };
1643        let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1644        let advance = agent::chain_advances(&result);
1645        if n == 0 || !advance {
1646            outcome = Some(result);
1647        } else {
1648            // A later failure is only logged; the note describes the first.
1649            tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1650        }
1651        if !advance {
1652            break;
1653        }
1654    }
1655    if outcome.as_ref().is_some_and(agent::chain_advances) {
1656        // Exhausted: back to the agent the conversation had, so the note
1657        // below names it and the next turn starts from it again.
1658        if let Some((id, seat)) = first_try {
1659            talk.agent = id;
1660            talk.seat = seat;
1661            fell_back_from = None;
1662        }
1663    }
1664    let outcome = outcome.expect("a chain holds at least one agent");
1665    let note = |why: String| Turn {
1666        who: Who::Agent,
1667        body: format!("{MAGI_NOTE}{why}"),
1668        at: Timestamp::now(),
1669        attachments: Vec::new(),
1670        usage: None,
1671    };
1672    let (reply, failure) = match outcome {
1673        Err(e) => (
1674            note(format!("could not run agent `{}`: {e}", talk.agent)),
1675            Some(format!("could not run agent `{}`: {e}", talk.agent)),
1676        ),
1677        Ok(out) if out.quota_exhausted() => {
1678            let reset = out
1679                .quota
1680                .as_ref()
1681                .and_then(|q| q.reset.clone())
1682                .map_or_else(String::new, |r| format!(" (resets {r})"));
1683            let why = format!(
1684                "agent `{}` is out of quota{reset}; your message is saved, so \
1685                 say it again when the window reopens",
1686                talk.agent
1687            );
1688            (note(why.clone()), Some(why))
1689        }
1690        Ok(out) if out.timed_out => {
1691            let why = format!(
1692                "agent `{}` did not answer within {}s; your message is saved",
1693                talk.agent,
1694                turn_timeout(cfg).as_secs()
1695            );
1696            (note(why.clone()), Some(why))
1697        }
1698        Ok(out) if !out.usable() => {
1699            let why = format!(
1700                "agent `{}` produced no answer (exit {}); your message is saved",
1701                talk.agent,
1702                out.exit_code
1703                    .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1704            );
1705            (note(why.clone()), Some(why))
1706        }
1707        Ok(out) => (
1708            Turn {
1709                who: Who::Agent,
1710                body: out.text.trim().to_owned(),
1711                at: Timestamp::now(),
1712                attachments: Vec::new(),
1713                // `talk.agent` is whoever actually answered: a fallback has
1714                // already moved it, and an exhausted chain never reaches here.
1715                usage: out.context_tokens.map(|context_tokens| TurnUsage {
1716                    context_tokens,
1717                    agent: talk.agent.clone(),
1718                    model: cfg
1719                        .agents
1720                        .iter()
1721                        .find(|a| a.id == talk.agent)
1722                        .and_then(|a| a.model.clone()),
1723                }),
1724            },
1725            None,
1726        ),
1727    };
1728
1729    // A close landed on disk while this turn was in flight is read back here
1730    // rather than trusted from the snapshot this call started with. `store`
1731    // holds nothing else this function does not itself own - the turn guard
1732    // in `web::Ui::begin_talk_turn` keeps `turns` and `seat` this call's
1733    // alone to mutate - but `status` is not behind that guard, and an
1734    // operator's close must stick: the whole point of ending a conversation
1735    // is that an agent's answer to the last message before the close cannot
1736    // silently reopen it. The guard is what makes that read-then-write
1737    // section atomic with `close`'s own - taken only for this tail and not
1738    // for the whole invocation above, so one talk's fifteen-minute turn does
1739    // not block another talk's close from proceeding.
1740    let _guard = store.guard()?;
1741    // A delete is the more final version of that same race: `put` writes
1742    // unconditionally, so a talk removed while this turn was in flight must
1743    // stay removed rather than being written back with this turn's reply
1744    // appended to it. The reply is simply given up on - there is no
1745    // conversation left for it to belong to.
1746    let Ok(fresh) = store.get(&talk.id) else {
1747        return Ok(());
1748    };
1749    talk.status = fresh.status;
1750    // `queue` may have accepted another operator message while the CLI was
1751    // running. This handle predates that write, so preserving only `status`
1752    // would overwrite the durable draft when the reply is appended below.
1753    talk.pending = fresh.pending;
1754    talk.pending_attachments = fresh.pending_attachments;
1755    if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1756        // The switch persists: quota coming back does not move the chat
1757        // home, an operator's switch does.
1758        talk.turns.push(note(format!(
1759            "agent changed from {from} to {} (fallback)",
1760            talk.agent
1761        )));
1762    }
1763    // The persona is in the session's context only once a turn has answered
1764    // (a failed one may never have reached the agent), so only then is it clean.
1765    if failure.is_none() {
1766        talk.persona_dirty = false;
1767    }
1768    talk.turns.push(reply);
1769    if let Err(put_err) = store.put(talk) {
1770        // `Talks::put` already retried the write itself - reaching here
1771        // means a passing race is not what this is. An agent's answer,
1772        // possibly the result of an hour-long call, must not vanish with
1773        // nothing to show for it just because the very last step failed:
1774        // pop it back off, stash its text beside the conversation, and
1775        // replace it with a note the operator can actually see, the same
1776        // mechanism the failure branches above already use for a quota or a
1777        // timeout.
1778        let lost = talk.turns.pop().expect("just pushed above");
1779        let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1780        let why = match &stash {
1781            Ok(path) => format!(
1782                "agent `{}` answered, but the reply could not be saved to \
1783                 this conversation ({put_err:#}); the raw text was kept at \
1784                 {} - your message is saved, ask again",
1785                talk.agent,
1786                path.display()
1787            ),
1788            Err(stash_err) => format!(
1789                "agent `{}` answered, but the reply could not be saved to \
1790                 this conversation ({put_err:#}), and it could not be kept \
1791                 anywhere else either ({stash_err:#}); your message is \
1792                 saved, ask again",
1793                talk.agent
1794            ),
1795        };
1796        talk.turns.push(note(why.clone()));
1797        // Writing the note also carries the seat this call already advanced -
1798        // `agent::invoke` incremented `turns` and, for a vendor that reports
1799        // its own session id, recorded that too. That is what keeps the next
1800        // turn resuming the session the CLI is already holding instead of
1801        // re-opening it, so losing the reply costs the transcript a turn but
1802        // not the conversation.
1803        return match store.put(talk) {
1804            Ok(()) => bail!("{why}"),
1805            Err(note_err) => {
1806                // Even the short note failed to save, which means this
1807                // conversation's file cannot be written at all right now -
1808                // nothing is left for this call to retry or record. Pop the
1809                // note so `talk.turns` matches the transcript on disk, and
1810                // surface both failures for whoever reads the log.
1811                //
1812                // `talk.seat` is deliberately not wound back to match. The
1813                // CLI really did take the turn and really did consume this
1814                // seat's session id; pretending otherwise would be a second
1815                // untruth on top of the unwritable file, and the handle is
1816                // reloaded from disk by the next `drain` or `get` anyway -
1817                // see `web::drain_loop`. What the seat cannot do is reach
1818                // disk, so the record stays a turn behind the CLI until some
1819                // later write lands, and a turn taken before then re-opens a
1820                // session id the CLI already holds. That is the desync
1821                // `20260907-011805-fb57` is about, and tolerating it belongs
1822                // there rather than here: no write this branch could make
1823                // would help, since a failed write is exactly what put it in
1824                // this position twice over.
1825                talk.turns.pop();
1826                Err(note_err).context(why)
1827            }
1828        };
1829    }
1830
1831    match failure {
1832        Some(why) => bail!("{why}"),
1833        None => Ok(()),
1834    }
1835}
1836
1837/// Everything said so far, as prose, for a CLI that cannot resume its own
1838/// conversation.
1839fn transcript(talk: &Talk, store: &Talks) -> String {
1840    let mut out = String::from(
1841        "This conversation cannot resume on the CLI's side, so here is \
1842         everything said so far; answer only the last message.\n",
1843    );
1844    for t in &talk.turns {
1845        let who = match t.who {
1846            Who::Operator => "operator",
1847            Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1848            Who::Agent => "you",
1849        };
1850        out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1851        out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1852    }
1853    out
1854}
1855
1856/// The section named at the end of a turn's body, listing every attachment's
1857/// absolute path and mime so the agent knows exactly what to open. Empty
1858/// when `attachments` is, which is every turn but the rare one carrying an
1859/// image, so a turn with none changes nothing about the prompt.
1860fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1861    if attachments.is_empty() {
1862        return String::new();
1863    }
1864    let mut out = String::from(
1865        "\n\nThe operator attached the image(s) below to this message. Open \
1866         and look at each one before you answer.\n",
1867    );
1868    for att in attachments {
1869        if let Some(path) = store.attachment_path(talk_id, att) {
1870            out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1871        }
1872    }
1873    out.push('\n');
1874    out
1875}
1876
1877/// The briefing the agent opens with, sent once as part of its first turn.
1878///
1879/// Pure, so the properties that matter can be asserted without an interview:
1880/// it names `magi task add --solo` (the route this conversation always has to
1881/// changing anything) and it never tells the agent to write a task *file* of
1882/// its own - that would compete with filing through the queue.
1883/// What a turn may do: `(allow_write, unsandboxed)`. Only the repository's own
1884/// `[talk] allow_write` lifts the CLI sandbox (codex's counterpart of claude's
1885/// `bypassPermissions`); a turn writable merely because a handed-over question
1886/// is open keeps the sandbox and gets just the question store as a writable root.
1887pub(crate) fn turn_access(talk_allow_write: bool, consulted: bool) -> (bool, bool) {
1888    (talk_allow_write || consulted, talk_allow_write)
1889}
1890
1891/// `allow_write` only ever adds an extra permission on top of that; it never
1892/// removes the queue as an option, which is why both branches keep the same
1893/// `# When the operator wants something done` section - `write_policy` is
1894/// the only part that changes.
1895///
1896/// It also tells the agent that `--repo` is not stuck naming this
1897/// conversation's own directory: `resolve_repo` (`src/main.rs`) now accepts a
1898/// short `owner/repo` or bare `repo` name and resolves it against
1899/// `[repos] roots`, the same local checkouts `magi repos` lists. Without this
1900/// line an agent asked to change some other repository has no way to know
1901/// that option exists, and the only path it can see - asking the operator to
1902/// dictate a full path - is exactly the friction this change exists to
1903/// remove. A miss or an ambiguous name still fails the command outright, so
1904/// the instruction is to ask rather than guess when that happens - the
1905/// silent-decision line this task must not cross.
1906pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1907    briefing_with(repo, language, allow_write, None, None)
1908}
1909
1910/// [`briefing`] plus the selected persona's tone-only section and, when the
1911/// operator has a configured name, how to address them. `None` for both is
1912/// byte-for-byte the plain briefing. A name changed mid-session reaches a
1913/// running CLI session only on its next fresh seat (nothing is persisted).
1914pub fn briefing_with(
1915    repo: &Path,
1916    language: &str,
1917    allow_write: bool,
1918    persona: Option<&crate::persona::Persona>,
1919    operator_name: Option<&str>,
1920) -> String {
1921    let write_policy = if allow_write {
1922        "Write access is enabled for this conversation (`allow_write = \
1923         true`), so you may write files - but only a small, \
1924         already-decided edit the operator names outright in this \
1925         conversation, not an implementation. This is a permission on the \
1926         conversation as a whole, not a property of whichever repository \
1927         it happened to start in: if the operator names a different \
1928         repository for that small edit, the policy allows it there too. \
1929         Your own tool may still confine writes to the repository this \
1930         conversation started in regardless - if a write elsewhere is \
1931         refused, say so plainly rather than working around it. Once you \
1932         have made an edit, say plainly what you edited. Anything bigger, \
1933         or anything still open-ended, still goes through the queue below \
1934         rather than being done here."
1935    } else {
1936        "Do not write files. Implementing a change is not this \
1937         conversation's job; a separate, blind competition of agents does \
1938         that, and a repository this conversation has already edited would \
1939         make their diffs unjudgeable."
1940    };
1941    let mut out = format!(
1942        "You are magi's standing conversation partner for its operator, who \
1943         usually has this open on a phone. Keep replies short: no preamble, \
1944         no restating what they just said.\n\n\
1945         # Repository\n\n{repo}\n\n\
1946         You may look around: read files, run shell commands, search history, \
1947         run tests - whatever answers the question. {write_policy}\n\n\
1948         A short, command-shaped message (\"list\", \"info <id>\", \"show \
1949         3cbf\") is almost always the operator asking you to look something \
1950         up, not an instruction to file - answer it yourself with `magi \
1951         list`, `magi show <id>`, `magi task list`, or the like, the same way \
1952         you would answer any other question in this conversation.\n\n\
1953         # When the operator wants something done\n\n\
1954         Run:\n\n\
1955         magi task add --solo --repo {repo} <instruction>\n\n\
1956         and tell the operator the task id it prints, so they can follow it \
1957         from the Queue. If it refuses with a duplicate warning (the \
1958         instruction names a branch, commit or pull request that an \
1959         unfinished task, run or PR already owns), do not repeat it with \
1960         --force yourself: tell the operator what it matched and let them \
1961         decide. Write <instruction> so that an implementer who has \
1962         never seen this conversation can act on it alone - it is everything \
1963         they get. Use --solo: it runs the task through one implementer \
1964         straight into review instead of the usual multi-agent competition, \
1965         which is the right shape for a change this conversation has already \
1966         settled, rather than one still worth several independent takes.\n\n\
1967         If the operator asks for something in a different repository, \
1968         --repo does not have to be a full path: --repo owner/repo (or just \
1969         repo, when that is unambiguous) is resolved against local checkouts \
1970         the same way `magi repos` lists them. If the command fails because \
1971         nothing matches or more than one checkout shares that name, ask the \
1972         operator which repository they mean (or run `magi repos` yourself \
1973         to see the candidates) rather than guessing.\n\n\
1974         The current state of the code is whatever origin/main holds, not \
1975         whatever a working tree shows: a primary checkout often lags \
1976         upstream, sits on a detached HEAD and carries uncommitted changes. \
1977         Before answering about code, run `git fetch origin` in that \
1978         repository if it is cheap, then read through \
1979         `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
1980         If the working tree differs, say so; if the fetch fails, say that \
1981         too, so the operator knows the answer may be stale.\n\n\
1982         If the operator attached an image (a screenshot, say) that the task \
1983         is about, pass it with `--attach <path>`, using the absolute path \
1984         the turn's attachment note gives; repeat the flag for several. \
1985         `magi task add --solo --attach <path> <instruction>` copies the \
1986         file into the task, so the implementer receives it. Do not paste the \
1987         path into <instruction> instead: deleting this conversation deletes \
1988         its attachments, and then that path reaches no one.\n",
1989        repo = repo.display(),
1990    );
1991    out.push_str(&language_note(language));
1992    match (persona, operator_name) {
1993        (Some(p), n) => out.push_str(&crate::persona::section_for(p, n)),
1994        (None, Some(n)) => {
1995            out.push_str("\n# Addressing the operator\n");
1996            out.push_str(&crate::persona::addressing(n));
1997        }
1998        (None, None) => {}
1999    }
2000    out
2001}
2002
2003/// The operator is talking, so their language matters here more than in most
2004/// prompts magi sends.
2005fn language_note(language: &str) -> String {
2006    if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
2007        String::new()
2008    } else {
2009        format!("\nHold this conversation in {language}.\n")
2010    }
2011}
2012
2013/// Queue tasks this conversation has filed, oldest first.
2014///
2015/// A task is this conversation's when its [`Source::Agent`] names this
2016/// conversation's id as `run` - which is exactly what happens when
2017/// `magi task add` is run from inside a turn, because [`turn`] passes the
2018/// conversation's own id as [`Invocation::run`].
2019pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
2020    let mut tasks: Vec<Task> = queue
2021        .list()
2022        .into_iter()
2023        .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
2024        .collect();
2025    tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
2026    tasks
2027}
2028
2029fn read_path(path: &Path) -> Result<Talk> {
2030    let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2031    serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
2032}
2033
2034/// How many times [`write_atomic`] retries a failed write-then-rename before
2035/// giving up.
2036const PUT_RETRIES: u32 = 5;
2037
2038/// Write `body` to `tmp` and rename it onto `path`, retrying the whole thing
2039/// a handful of times with a short sleep in between.
2040///
2041/// The only failure this is meant to absorb is a passing one - most
2042/// concretely, a reader elsewhere in this process (or another `magi`
2043/// process) with `path` briefly open for `read_to_string` at the exact
2044/// moment this call tries to rename over it. That clears in milliseconds
2045/// once the reader lets go; a caller still failing after several short
2046/// sleeps has something more durable wrong (a full disk, a permissions
2047/// change) that a longer sleep would not fix either, and is left to report
2048/// it.
2049fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2050    let mut last_err = None;
2051    for attempt in 0..PUT_RETRIES {
2052        if attempt > 0 {
2053            std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
2054        }
2055        match try_write_atomic(tmp, path, body) {
2056            Ok(()) => return Ok(()),
2057            Err(e) => last_err = Some(e),
2058        }
2059    }
2060    Err(last_err.expect("the loop above always runs at least once"))
2061}
2062
2063fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2064    #[cfg(test)]
2065    if failpoint::take_forced_put_failure() {
2066        bail!("simulated write failure (test)");
2067    }
2068    std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
2069    std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
2070    Ok(())
2071}
2072
2073/// Last resort when `turn`'s own `store.put` fails even after
2074/// [`write_atomic`]'s retries: keep the generated text somewhere still
2075/// findable rather than let the whole of an agent's answer disappear along
2076/// with the write that was supposed to record it.
2077fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
2078    let dir = store.artifacts_of(id);
2079    std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
2080    let path = dir.join(format!("{stem}-lost.txt"));
2081    std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
2082    Ok(path)
2083}
2084
2085/// A test-only seam that lets [`try_write_atomic`] simulate the kind of
2086/// passing I/O race [`write_atomic`] is meant to retry through, without
2087/// depending on real OS-level file-locking behaviour, which differs across
2088/// the three platforms this crate ships on (and, on the one platform where a
2089/// reader really does block a rename, is awkward to trigger deterministically
2090/// in a unit test).
2091#[cfg(test)]
2092mod failpoint {
2093    use std::cell::Cell;
2094
2095    thread_local! {
2096        static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
2097    }
2098
2099    /// Arrange for the next `count` calls into [`super::try_write_atomic`] to
2100    /// fail before touching the filesystem at all.
2101    pub(super) fn force_put_failures(count: u32) {
2102        FORCE_PUT_FAILURES.with(|c| c.set(count));
2103    }
2104
2105    /// Consumed once per attempt inside [`super::try_write_atomic`]; `true`
2106    /// means simulate this attempt failing.
2107    pub(super) fn take_forced_put_failure() -> bool {
2108        FORCE_PUT_FAILURES.with(|c| {
2109            let n = c.get();
2110            if n == 0 {
2111                false
2112            } else {
2113                c.set(n - 1);
2114                true
2115            }
2116        })
2117    }
2118}
2119
2120fn short(id: &str) -> &str {
2121    id.split('-').next_back().unwrap_or(id)
2122}
2123
2124fn new_id() -> String {
2125    let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2126    let seed = crate::rng::entropy();
2127    format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2128}
2129
2130/// Extension an attachment's bytes are stored under, from its (already
2131/// validated) mime. The one place this mapping exists on the write side;
2132/// `web`'s own whitelist is what actually decides which mimes are accepted
2133/// in the first place.
2134fn attachment_ext(mime: &str) -> Option<&'static str> {
2135    match mime {
2136        "image/png" => Some("png"),
2137        "image/jpeg" => Some("jpg"),
2138        "image/gif" => Some("gif"),
2139        "image/webp" => Some("webp"),
2140        _ => None,
2141    }
2142}
2143
2144/// Is `id` a shape [`put_attachment`](Talks::put_attachment) could have
2145/// produced? 32 lowercase hex digits and nothing else, checked before an id
2146/// that came from the client is ever allowed to build a path - so `..` and a
2147/// path separator are never even possible.
2148pub fn valid_attachment_id(id: &str) -> bool {
2149    id.len() == 32
2150        && id
2151            .bytes()
2152            .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
2153}
2154
2155/// A fresh attachment id: 128 bits of process entropy as lowercase hex - the
2156/// same "mint it, never take it from the client" rule [`new_id`] follows for
2157/// conversation ids.
2158fn new_attachment_id() -> String {
2159    let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
2160    format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
2161}
2162
2163#[cfg(test)]
2164mod tests {
2165    #[test]
2166    fn turn_access_lifts_the_sandbox_only_for_the_talk_opt_in() {
2167        use super::turn_access;
2168        assert_eq!(turn_access(false, false), (false, false));
2169        assert_eq!(turn_access(true, false), (true, true));
2170        // A handed-over question opens the question store, never the sandbox.
2171        assert_eq!(turn_access(false, true), (true, false));
2172        assert_eq!(turn_access(true, true), (true, true));
2173    }
2174
2175    #[test]
2176    fn the_briefing_points_at_origin_main_not_the_working_tree() {
2177        let b = briefing(Path::new("/r"), "en", false);
2178        assert!(b.contains("origin/main"));
2179        assert!(b.contains("git show origin/main:"));
2180    }
2181    use std::collections::BTreeMap;
2182
2183    use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
2184    use crate::queue::{Queue, Source, Task};
2185
2186    use super::*;
2187
2188    fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
2189        AgentSpec {
2190            id: id.to_owned(),
2191            kind: AgentKind::Command,
2192            model: model.map(str::to_owned),
2193            command: Vec::new(),
2194            extra_args: Vec::new(),
2195            env: BTreeMap::new(),
2196            prompt_delivery: None,
2197        }
2198    }
2199
2200    fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
2201        Talk {
2202            schema: SCHEMA,
2203            id: "20260904-014455-ab12".to_owned(),
2204            repo: PathBuf::from("."),
2205            agent: agent.to_owned(),
2206            status: TalkStatus::Open,
2207            turns,
2208            pending: String::new(),
2209            pending_attachments: Vec::new(),
2210            fallback: false,
2211            persona: String::new(),
2212            persona_dirty: false,
2213            created_at: Timestamp::now(),
2214            updated_at: Timestamp::now(),
2215            seat: SeatState::new(SEAT, agent, 1),
2216        }
2217    }
2218
2219    fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
2220        Turn {
2221            who: Who::Agent,
2222            body: body.to_owned(),
2223            at: Timestamp::now(),
2224            attachments: Vec::new(),
2225            usage: usage.map(|(t, a, m)| TurnUsage {
2226                context_tokens: t,
2227                agent: a.to_owned(),
2228                model: m.map(str::to_owned),
2229            }),
2230        }
2231    }
2232
2233    fn ctx_config(windows: &[(&str, u64)]) -> Config {
2234        Config {
2235            agents: vec![
2236                ctx_agent("small", Some("small-model")),
2237                ctx_agent("big", Some("big-model")),
2238                ctx_agent("plain", None),
2239            ],
2240            context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
2241            ..Config::default()
2242        }
2243    }
2244
2245    #[test]
2246    fn context_usage_computes_percent_and_warns_at_eighty() {
2247        let cfg = ctx_config(&[("small-model", 1000)]);
2248        let at = |tokens| {
2249            let t = ctx_talk(
2250                "small",
2251                vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
2252            );
2253            context_usage(&t, Some(&cfg))
2254        };
2255        let u = at(799);
2256        assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
2257        let u = at(800);
2258        assert_eq!((u.percent, u.warn), (Some(80), true));
2259        let u = at(1500);
2260        assert_eq!((u.percent, u.warn), (Some(150), true));
2261        assert!(!u.since_switch);
2262    }
2263
2264    #[test]
2265    fn context_usage_is_unknown_without_usage_and_never_looks_back() {
2266        let cfg = ctx_config(&[("small-model", 1000)]);
2267        let t = ctx_talk(
2268            "small",
2269            vec![
2270                reply("old", Some((900, "small", Some("small-model")))),
2271                reply("new", None),
2272            ],
2273        );
2274        let u = context_usage(&t, Some(&cfg));
2275        // No measurement in the latest reply: an estimate, never the stale 900.
2276        assert!(u.estimated);
2277        assert_ne!(u.tokens, Some(900));
2278        assert!(u.tokens.is_some());
2279        // A magi note after the reply neither hides nor replaces it.
2280        let t = ctx_talk(
2281            "small",
2282            vec![
2283                reply("old", Some((900, "small", Some("small-model")))),
2284                reply("magi: could not run agent", None),
2285            ],
2286        );
2287        assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
2288        assert_eq!(
2289            context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
2290            None
2291        );
2292    }
2293
2294    #[test]
2295    fn estimate_counts_chars_both_sides_and_standing_prompt() {
2296        let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
2297        assert_eq!(estimate_context_tokens(&t, 0), Some(2)); // 7 chars -> 2
2298        let op = Turn {
2299            who: Who::Operator,
2300            ..reply("abcdefg", None)
2301        };
2302        t.turns.push(op);
2303        assert_eq!(estimate_context_tokens(&t, 0), Some(4));
2304        assert!(
2305            estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
2306        );
2307        // Characters, not bytes: 7 kanji are 7 chars.
2308        let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
2309        assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
2310        // magi notes are not sent to the agent; with nothing else, unknown.
2311        let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
2312        assert_eq!(estimate_context_tokens(&note, 1000), None);
2313        assert_eq!(
2314            estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
2315            None
2316        );
2317    }
2318
2319    #[test]
2320    fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
2321        let cfg = ctx_config(&[("small-model", 1000)]);
2322        let t = ctx_talk(
2323            "small",
2324            vec![reply(
2325                &"x".repeat(5000),
2326                Some((10, "small", Some("small-model"))),
2327            )],
2328        );
2329        let u = context_usage(&t, Some(&cfg));
2330        assert_eq!((u.tokens, u.estimated), (Some(10), false));
2331        let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2332        let u = context_usage(&t, Some(&cfg));
2333        assert!(u.estimated && !u.since_switch);
2334        assert_eq!(u.window, Some(1000));
2335        assert!(u.warn && u.percent.unwrap() >= 80);
2336        let t = ctx_talk("small", vec![reply("hi", None)]);
2337        let u = context_usage(&t, Some(&cfg));
2338        assert!(u.estimated && u.percent.is_some());
2339    }
2340
2341    #[test]
2342    fn context_usage_without_a_window_shows_tokens_only() {
2343        let cfg = ctx_config(&[]);
2344        // No model at all, and a model nothing matches.
2345        let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2346        let u = context_usage(&t, Some(&cfg));
2347        assert_eq!(
2348            (u.tokens, u.window, u.percent, u.warn),
2349            (Some(5000), None, None, false)
2350        );
2351        let t = ctx_talk(
2352            "small",
2353            vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2354        );
2355        assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2356        // No readable config: same, and no panic.
2357        assert_eq!(context_usage(&t, None).window, None);
2358    }
2359
2360    #[test]
2361    fn context_usage_switching_model_changes_the_denominator() {
2362        let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2363        let used = reply("hi", Some((900, "small", Some("small-model"))));
2364        let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2365        assert_eq!(
2366            (before.percent, before.warn, before.since_switch),
2367            (Some(90), true, false)
2368        );
2369        // Same turns, conversation now on the big model: new denominator, and
2370        // the figure is flagged as describing the previous session.
2371        let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2372        assert_eq!(after.window, Some(10_000));
2373        assert_eq!(
2374            (after.percent, after.warn, after.since_switch),
2375            (Some(9), false, true)
2376        );
2377        assert_eq!(after.model.as_deref(), Some("big-model"));
2378    }
2379
2380    #[test]
2381    fn a_turn_recorded_before_usage_existed_still_reads() {
2382        let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2383        let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2384        assert!(turn.usage.is_none());
2385        let json = serde_json::to_string(&turn).expect("serialize");
2386        assert!(
2387            !json.contains("usage"),
2388            "absent usage is not written: {json}"
2389        );
2390    }
2391
2392    /// A store of its own, with no process-global state.
2393    fn store() -> (tempfile::TempDir, Talks) {
2394        let tmp = tempfile::tempdir().expect("tempdir");
2395        let talks = Talks::at(tmp.path().join("talks"));
2396        (tmp, talks)
2397    }
2398
2399    /// A `kind = "command"` agent whose whole behaviour is a POSIX shell
2400    /// script - see `chat`'s tests for why no test here may spawn a real
2401    /// agent CLI.
2402    fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2403        let path = dir.join("mock-talk-agent.sh");
2404        std::fs::write(&path, script).expect("write mock");
2405        AgentSpec {
2406            id: "mock".to_owned(),
2407            kind: AgentKind::Command,
2408            model: None,
2409            command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2410            extra_args: Vec::new(),
2411            env,
2412            prompt_delivery: None,
2413        }
2414    }
2415
2416    fn config(spec: AgentSpec) -> Config {
2417        Config {
2418            agents: vec![spec],
2419            graph: Graph {
2420                language: "en".to_owned(),
2421                ..Graph::default()
2422            },
2423            ..Config::default()
2424        }
2425    }
2426
2427    /// Echo a canned reply, ignoring the prompt on stdin.
2428    const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2429
2430    /// Say nothing and fail, the way a CLI that cannot start does.
2431    const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2432
2433    /// Reply with the prompt it was given, so a test can inspect exactly what
2434    /// the agent received on stdin.
2435    const ECHO: &str = "#!/bin/sh\ncat\n";
2436
2437    fn env(reply: &str) -> BTreeMap<String, String> {
2438        BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2439    }
2440
2441    #[test]
2442    fn the_frozen_json_field_names_round_trip_through_disk() {
2443        let (tmp, talks) = store();
2444        let mut talk = Talk {
2445            schema: SCHEMA,
2446            id: "20260904-014455-ab12".to_owned(),
2447            repo: tmp.path().to_owned(),
2448            agent: "sonnet".to_owned(),
2449            status: TalkStatus::Open,
2450            turns: Vec::new(),
2451            pending: String::new(),
2452            pending_attachments: Vec::new(),
2453            fallback: false,
2454            persona: String::new(),
2455            persona_dirty: false,
2456            created_at: Timestamp::now(),
2457            updated_at: Timestamp::now(),
2458            seat: SeatState::new(SEAT, "sonnet", 7),
2459        };
2460        talks.put(&mut talk).expect("put");
2461
2462        let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2463        let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2464        for field in [
2465            "schema",
2466            "id",
2467            "repo",
2468            "agent",
2469            "status",
2470            "turns",
2471            "created_at",
2472            "updated_at",
2473        ] {
2474            assert!(v.get(field).is_some(), "missing field `{field}`");
2475        }
2476        assert_eq!(v["schema"], 1);
2477        assert_eq!(v["status"], "open");
2478
2479        let back = talks.get(&talk.id).expect("get");
2480        assert_eq!(back.id, talk.id);
2481        assert_eq!(back.status, TalkStatus::Open);
2482    }
2483
2484    #[test]
2485    fn opening_a_talk_takes_no_agent_turn() {
2486        let (tmp, talks) = store();
2487        // A script that would fail loudly if it were ever run: `begin` must
2488        // not invoke anything, since there is nothing yet for an agent to
2489        // answer.
2490        let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2491        let cfg = config(spec);
2492
2493        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2494        assert_eq!(talk.status, TalkStatus::Open);
2495        assert!(talk.turns.is_empty(), "nothing has been said yet");
2496
2497        let on_disk = talks.get(&talk.id).expect("get");
2498        assert_eq!(on_disk.turns.len(), 0);
2499    }
2500
2501    /// `[roles] chatter`, when set, decides who holds this conversation; unset,
2502    /// it falls back to [`agent::pick`]'s own default order (a claude seat,
2503    /// else the first runnable agent in roster order) rather than to any
2504    /// other role - see `[roles] chatter`'s own doc in [`crate::config`] for
2505    /// why a dedicated field exists at all: opening this against the same
2506    /// seat as a judge is what produced the `agent ... did not answer within
2507    /// 300s` timeout that led to it.
2508    #[test]
2509    fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2510        let (tmp, talks) = store();
2511        let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2512        let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2513        chatter_spec.id = "chatter-mock".to_owned();
2514
2515        let mut cfg = Config {
2516            agents: vec![first_spec.clone(), chatter_spec.clone()],
2517            graph: Graph {
2518                language: "en".to_owned(),
2519                ..Graph::default()
2520            },
2521            ..Config::default()
2522        };
2523        cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2524
2525        let talk =
2526            begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2527        assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2528
2529        cfg.roles.chatter = None;
2530        let fallback =
2531            begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2532        assert_eq!(
2533            fallback.agent, first_spec.id,
2534            "unset chatter must fall back to agent::pick's own default order"
2535        );
2536    }
2537
2538    /// A conversation recorded before attachments existed - schema 1, no
2539    /// `attachments` key on any turn - must still read.
2540    #[test]
2541    fn a_talk_recorded_without_attachments_still_reads() {
2542        let (tmp, talks) = store();
2543        let path = talks.path_of("20260904-014455-ab12");
2544        std::fs::create_dir_all(talks.root()).expect("talks dir");
2545        std::fs::write(
2546            &path,
2547            serde_json::json!({
2548                "schema": 1,
2549                "id": "20260904-014455-ab12",
2550                "repo": tmp.path(),
2551                "agent": "sonnet",
2552                "status": "open",
2553                "turns": [
2554                    { "who": "operator", "body": "still there?",
2555                      "at": Timestamp::now().to_string() },
2556                ],
2557                "created_at": Timestamp::now().to_string(),
2558                "updated_at": Timestamp::now().to_string(),
2559                "seat": SeatState::new(SEAT, "sonnet", 7),
2560            })
2561            .to_string(),
2562        )
2563        .expect("write pre-attachments talk");
2564
2565        let talk = talks.get("20260904-014455-ab12").expect("must still read");
2566        assert!(talk.turns[0].attachments.is_empty());
2567    }
2568
2569    fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2570        let tmp = tempfile::TempDir::new().expect("tmp");
2571        let root = tmp.path().join("talks");
2572        (tmp, Talks::at(root.clone()), Talks::at(root))
2573    }
2574
2575    #[test]
2576    fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2577        let (_tmp, a, b) = lease_store();
2578        let won = a.claim_turn("t1").expect("claim").expect("first wins");
2579        assert!(
2580            b.claim_turn("t1").expect("claim").is_none(),
2581            "second is refused"
2582        );
2583        assert!(b.turn_held("t1"));
2584        assert!(
2585            b.claim_turn("t2").expect("claim").is_some(),
2586            "other talks are free"
2587        );
2588        drop(won);
2589    }
2590
2591    #[test]
2592    fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2593        let (_tmp, a, b) = lease_store();
2594        let old = a.claim_turn("t1").expect("claim").expect("held");
2595        let later = Timestamp::now()
2596            .checked_add(jiff::SignedDuration::from_secs(
2597                crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2598            ))
2599            .expect("later");
2600        let new = b
2601            .claim_turn_at("t1", later)
2602            .expect("claim")
2603            .expect("a stale lease is taken over");
2604        drop(old);
2605        assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2606        assert!(new.beat().expect("beat"), "the new owner still beats");
2607        drop(new);
2608        assert!(!a.turn_held("t1"));
2609    }
2610
2611    #[tokio::test]
2612    async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2613        let (_tmp, a, b) = lease_store();
2614        let old = a.claim_turn("t1").expect("claim").expect("held");
2615        let later = Timestamp::now()
2616            .checked_add(jiff::SignedDuration::from_secs(
2617                crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2618            ))
2619            .expect("later");
2620        let _new = b
2621            .claim_turn_at("t1", later)
2622            .expect("claim")
2623            .expect("taken over");
2624        let out = old
2625            .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2626            .await;
2627        assert!(out.is_err(), "the displaced turn must stop, not run on");
2628    }
2629
2630    #[tokio::test]
2631    async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2632        let (_tmp, a, _b) = lease_store();
2633        let lease = a.claim_turn("t1").expect("claim").expect("held");
2634        let out = lease
2635            .beating_every(Duration::from_millis(5), async {
2636                tokio::time::sleep(Duration::from_millis(40)).await;
2637                7
2638            })
2639            .await
2640            .expect("still ours");
2641        assert_eq!(out, 7);
2642        assert!(a.turn_held("t1"));
2643    }
2644
2645    #[test]
2646    fn an_unreadable_lease_counts_as_stale() {
2647        let (_tmp, a, b) = lease_store();
2648        std::fs::create_dir_all(a.root()).expect("dir");
2649        std::fs::write(a.turn_path("t1"), "not json").expect("write");
2650        assert!(!a.turn_held("t1"));
2651        assert!(b.claim_turn("t1").expect("claim").is_some());
2652    }
2653
2654    #[test]
2655    fn a_lease_is_released_when_the_turn_ends_or_fails() {
2656        let (_tmp, a, b) = lease_store();
2657        let lease = a.claim_turn("t1").expect("claim").expect("held");
2658        let failed: Result<()> = (|| {
2659            let _held = &lease;
2660            bail!("turn failed")
2661        })();
2662        assert!(failed.is_err());
2663        assert!(
2664            b.claim_turn("t1").expect("claim").is_none(),
2665            "held mid-turn"
2666        );
2667        drop(lease);
2668        assert!(
2669            b.claim_turn("t1").expect("claim").is_some(),
2670            "free after the turn"
2671        );
2672    }
2673
2674    #[test]
2675    fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2676        let (_tmp, a, _b) = lease_store();
2677        drop(a.claim_turn("t1").expect("claim").expect("held"));
2678        std::fs::write(
2679            a.turn_path("t1"),
2680            serde_json::to_string(&TurnRecord {
2681                token: "gone".into(),
2682                pid: 1,
2683                beat_at: Timestamp::from_second(1).expect("ts"),
2684            })
2685            .expect("json"),
2686        )
2687        .expect("write");
2688        let wins: Vec<_> = std::thread::scope(|sc| {
2689            let hs: Vec<_> = (0..8)
2690                .map(|_| {
2691                    let s = a.clone();
2692                    sc.spawn(move || s.claim_turn("t1").expect("claim"))
2693                })
2694                .collect();
2695            hs.into_iter().map(|h| h.join().expect("join")).collect()
2696        });
2697        assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
2698    }
2699
2700    fn age_file(path: &Path) {
2701        let f = std::fs::OpenOptions::new()
2702            .write(true)
2703            .open(path)
2704            .expect("open");
2705        f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
2706            .expect("age");
2707    }
2708
2709    #[test]
2710    fn a_late_taker_of_a_broken_lock_cannot_disturb_its_replacement() {
2711        let (_tmp, a, _b) = lease_store();
2712        let lease = a.turn_path("t1");
2713        let lock = lease.with_extension("turn.lock");
2714        std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2715        std::fs::write(&lock, "t1-dead").expect("dead lock");
2716        age_file(&lock);
2717        // B breaks the dead lock and holds a fresh one.
2718        let b = TurnLock::take(&lease).expect("take").expect("b wins");
2719        let fresh = std::fs::read_to_string(&lock).expect("read");
2720        assert_eq!(fresh, b.token);
2721        // C read the same dead token earlier: its ticket is already spent, and
2722        // the fresh lock is never moved, so a newcomer cannot slip in.
2723        let ticket = std::fs::read_dir(lock.parent().expect("dir"))
2724            .expect("dir")
2725            .flatten()
2726            .map(|e| e.path())
2727            .find(|p| p.to_string_lossy().ends_with(".break.0"))
2728            .expect("ticket");
2729        assert!(!create_exclusive(&ticket, "").expect("ticket"));
2730        assert!(TurnLock::take(&lease).expect("take").is_none());
2731        assert_eq!(std::fs::read_to_string(&lock).expect("read"), fresh);
2732        // A later generation is breakable once the old ticket is stale too.
2733        age_file(&lock);
2734        age_file(&ticket);
2735        std::mem::forget(b);
2736        let c = TurnLock::take(&lease)
2737            .expect("take")
2738            .expect("next generation");
2739        assert_ne!(c.token, fresh);
2740    }
2741
2742    #[test]
2743    fn a_live_ticket_blocks_and_a_stale_one_hands_over_to_the_next_generation() {
2744        let (_tmp, a, _b) = lease_store();
2745        let lease = a.turn_path("t1");
2746        let lock = lease.with_extension("turn.lock");
2747        std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2748        std::fs::write(&lock, "t1-dead").expect("dead lock");
2749        age_file(&lock);
2750        let t0 = lock.with_extension("lock.t1-dead.break.0");
2751        assert!(create_exclusive(&t0, "").expect("ticket"));
2752        // A fresh ticket 0 blocks the same dead token, whatever the clock says.
2753        assert!(TurnLock::take(&lease).expect("take").is_none());
2754        assert_eq!(std::fs::read_to_string(&lock).expect("read"), "t1-dead");
2755        // Once ticket 0 is stale, exactly one taker gets through via ticket 1.
2756        age_file(&t0);
2757        let c = TurnLock::take(&lease).expect("take").expect("generation 1");
2758        assert_eq!(std::fs::read_to_string(&lock).expect("read"), c.token);
2759        assert!(lock.with_extension("lock.t1-dead.break.1").exists());
2760        assert!(TurnLock::take(&lease).expect("take").is_none());
2761    }
2762
2763    #[test]
2764    fn exhausted_ticket_generations_recover_once_the_sweep_ages_them_out() {
2765        let (_tmp, a, _b) = lease_store();
2766        let lease = a.turn_path("t1");
2767        let lock = lease.with_extension("turn.lock");
2768        std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2769        std::fs::write(&lock, "t1-dead").expect("dead lock");
2770        age_file(&lock);
2771        let tickets: Vec<_> = (0..TICKET_GENERATIONS)
2772            .map(|n| lock.with_extension(format!("lock.t1-dead.break.{n}")))
2773            .collect();
2774        for t in &tickets {
2775            assert!(create_exclusive(t, "").expect("ticket"));
2776            age_file(t);
2777        }
2778        // Stale but not yet swept: no recovery, and no ticket is unlinked.
2779        assert!(TurnLock::take(&lease).expect("take").is_none());
2780        assert!(tickets.iter().all(|t| t.exists()));
2781        // Past the sweep age they go, and the next attempt recovers.
2782        for t in &tickets {
2783            let f = std::fs::OpenOptions::new()
2784                .write(true)
2785                .open(t)
2786                .expect("open");
2787            f.set_modified(std::time::SystemTime::now() - TICKET_SWEEP_AGE * 2)
2788                .expect("age");
2789        }
2790        assert!(TurnLock::take(&lease).expect("take").is_none());
2791        assert!(TurnLock::take(&lease).expect("take").is_some());
2792    }
2793
2794    #[test]
2795    fn an_aged_empty_lock_is_broken() {
2796        let (_tmp, a, _b) = lease_store();
2797        let lease = a.turn_path("t1");
2798        let lock = lease.with_extension("turn.lock");
2799        std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2800        std::fs::write(&lock, "").expect("empty lock");
2801        age_file(&lock);
2802        assert!(TurnLock::take(&lease).expect("take").is_some());
2803    }
2804
2805    #[test]
2806    fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
2807        let (tmp, talks) = store();
2808        let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2809        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2810
2811        queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
2812        queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
2813        let saved = talks.get(&talk.id).expect("reload queued talk");
2814        assert_eq!(saved.pending, "first\n\nsecond");
2815        assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
2816
2817        let drained = drain(&mut talk, &talks).expect("drain");
2818        assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
2819        let saved = talks.get(&talk.id).expect("reload drained talk");
2820        assert!(saved.pending.is_empty());
2821        assert_eq!(saved.turns.len(), 1);
2822        assert_eq!(saved.turns[0].body, "first\n\nsecond");
2823    }
2824
2825    #[test]
2826    fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
2827        let (tmp, talks) = store();
2828        let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2829        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2830        let attachment = Attachment {
2831            id: "a".repeat(32),
2832            name: "shot.png".to_owned(),
2833            mime: "image/png".to_owned(),
2834            bytes: 3,
2835        };
2836
2837        queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
2838        assert!(
2839            edit_pending_text(
2840                &mut talk,
2841                &talks,
2842                "corrected",
2843                "first",
2844                std::slice::from_ref(&attachment.id),
2845            )
2846            .expect("edit")
2847        );
2848        let saved = talks.get(&talk.id).expect("reload edited draft");
2849        assert_eq!(saved.pending, "corrected");
2850        assert_eq!(saved.pending_attachments, vec![attachment]);
2851
2852        queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2853        assert!(
2854            !edit_pending_text(
2855                &mut talk,
2856                &talks,
2857                "stale edit",
2858                "corrected",
2859                &["a".repeat(32)],
2860            )
2861            .expect("stale edit is a conflict")
2862        );
2863        assert_eq!(
2864            talks.get(&talk.id).expect("reload after conflict").pending,
2865            "corrected\n\nlater"
2866        );
2867        assert!(
2868            !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
2869                .expect("stale clear is a conflict")
2870        );
2871        assert_eq!(
2872            talks
2873                .get(&talk.id)
2874                .expect("reload after stale clear")
2875                .pending,
2876            "corrected\n\nlater"
2877        );
2878    }
2879
2880    #[tokio::test]
2881    async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
2882        let (tmp, talks) = store();
2883        let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
2884        let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
2885        let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2886        let id = running.id.clone();
2887        let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
2888
2889        let response_talks = talks.clone();
2890        let response_cfg = cfg.clone();
2891        let reply = tokio::spawn(async move {
2892            respond(
2893                &response_talks.claim_turn(&running.id).unwrap().unwrap(),
2894                &mut running,
2895                &response_talks,
2896                &response_cfg,
2897                &first,
2898            )
2899            .await
2900        });
2901        tokio::time::sleep(std::time::Duration::from_millis(20)).await;
2902
2903        let mut queued = talks.get(&id).expect("queued handle");
2904        queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
2905        reply.await.expect("join").expect("reply");
2906
2907        let saved = talks.get(&id).expect("reload");
2908        assert_eq!(saved.pending, "next");
2909        assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
2910    }
2911
2912    /// A mock whose script counts its own calls in `<dir>/<id>.calls` before
2913    /// running `body`.
2914    fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
2915        let calls = dir.join(format!("{id}.calls"));
2916        let script = format!(
2917            "#!/bin/sh\necho x >> '{}'\n{body}\n",
2918            calls.to_string_lossy()
2919        );
2920        let path = dir.join(format!("mock-{id}.sh"));
2921        std::fs::write(&path, script).expect("write mock");
2922        AgentSpec {
2923            id: id.to_owned(),
2924            kind: AgentKind::Command,
2925            model: None,
2926            command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2927            extra_args: Vec::new(),
2928            env: BTreeMap::new(),
2929            prompt_delivery: None,
2930        }
2931    }
2932
2933    fn calls(dir: &Path, id: &str) -> usize {
2934        std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
2935    }
2936
2937    fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
2938        let mut cfg = config(specs[0].clone());
2939        cfg.agents = specs;
2940        cfg.roles.chatter = Some(AgentChoice::Chain(
2941            ids.iter().map(|s| (*s).to_owned()).collect(),
2942        ));
2943        cfg
2944    }
2945
2946    #[tokio::test]
2947    async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
2948        let (tmp, talks) = store();
2949        let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2950        let b = counting_agent(tmp.path(), "b", "cat");
2951        let cfg = chain_config(vec![a, b], &["a", "b"]);
2952        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2953        assert_eq!(talk.agent, "a");
2954
2955        say(
2956            &talks.claim_turn(&talk.id).unwrap().unwrap(),
2957            &mut talk,
2958            &talks,
2959            &cfg,
2960            "hello there",
2961            Vec::new(),
2962        )
2963        .await
2964        .expect("turn");
2965        assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
2966        assert_eq!(calls(tmp.path(), "b"), 1);
2967        assert_eq!(talk.agent, "b", "the switch persists");
2968        assert!(talks.get(&talk.id).unwrap().agent == "b");
2969        let reply = talk.turns.last().unwrap();
2970        assert!(reply.body.contains("hello there"));
2971        assert!(
2972            reply.body.contains("magi task add --solo"),
2973            "a fresh seat gets the full briefing"
2974        );
2975        assert!(
2976            talk.turns
2977                .iter()
2978                .any(|t| t.body.contains("agent changed from a to b")),
2979            "the switch is noted"
2980        );
2981    }
2982
2983    #[tokio::test]
2984    async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2985        let (tmp, talks) = store();
2986        let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2987        let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2988        let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2989        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2990
2991        let err = say(
2992            &talks.claim_turn(&talk.id).unwrap().unwrap(),
2993            &mut talk,
2994            &talks,
2995            &cfg,
2996            "hi",
2997            Vec::new(),
2998        )
2999        .await
3000        .expect_err("every agent failed");
3001        assert!(err.to_string().contains("`a`"), "{err:#}");
3002        assert_eq!(calls(tmp.path(), "a"), 1);
3003        assert_eq!(calls(tmp.path(), "b"), 1);
3004        assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
3005    }
3006
3007    #[test]
3008    fn a_chatter_chain_skips_an_unknown_id_at_begin() {
3009        let (tmp, talks) = store();
3010        let b = counting_agent(tmp.path(), "b", "cat");
3011        let cfg = chain_config(vec![b], &["ghost", "b"]);
3012        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3013        assert_eq!(talk.agent, "b");
3014    }
3015
3016    #[tokio::test]
3017    async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
3018        let (tmp, talks) = store();
3019        let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3020        let b = counting_agent(tmp.path(), "b", "cat");
3021        let cfg = chain_config(vec![a, b], &["a", "b"]);
3022        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
3023        say(
3024            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3025            &mut talk,
3026            &talks,
3027            &cfg,
3028            "hi",
3029            Vec::new(),
3030        )
3031        .await
3032        .expect_err("a alone, and it fails");
3033        assert_eq!(calls(tmp.path(), "b"), 0);
3034        assert_eq!(talk.agent, "a");
3035    }
3036
3037    #[tokio::test]
3038    async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
3039        let (tmp, talks) = store();
3040        let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3041        let b = counting_agent(tmp.path(), "b", "cat");
3042        let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
3043        let cfg = chain_config(vec![a, b, c], &["a", "b"]);
3044        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
3045        say(
3046            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3047            &mut talk,
3048            &talks,
3049            &cfg,
3050            "hi",
3051            Vec::new(),
3052        )
3053        .await
3054        .expect_err("c alone, and it fails");
3055        assert_eq!(calls(tmp.path(), "b"), 0);
3056    }
3057
3058    #[test]
3059    fn briefing_names_the_operator_with_or_without_a_persona() {
3060        let plain = briefing(Path::new("/repo"), "en", false);
3061        let named = briefing_with(Path::new("/repo"), "en", false, None, Some("Commander"));
3062        assert!(named.starts_with(&plain));
3063        assert!(named.contains("# Addressing the operator"));
3064        assert!(named.contains("\"Commander\""));
3065        let rei = crate::persona::builtin_catalog()
3066            .into_iter()
3067            .find(|p| p.id == "rei")
3068            .expect("rei");
3069        let with = briefing_with(
3070            Path::new("/repo"),
3071            "en",
3072            false,
3073            Some(&rei),
3074            Some("Commander"),
3075        );
3076        assert!(with.contains("\"Commander\""));
3077        assert!(!with.contains("# Addressing the operator"));
3078    }
3079
3080    #[test]
3081    fn briefing_carries_a_persona_section_only_when_one_is_chosen() {
3082        let plain = briefing(Path::new("/repo"), "en", false);
3083        assert_eq!(
3084            plain,
3085            briefing_with(Path::new("/repo"), "en", false, None, None)
3086        );
3087        assert!(!plain.contains("Persona"));
3088        let rei = crate::persona::builtin_catalog()
3089            .into_iter()
3090            .find(|p| p.id == "rei")
3091            .expect("rei");
3092        let with = briefing_with(Path::new("/repo"), "en", false, Some(&rei), None);
3093        assert!(with.starts_with(&plain), "the plain briefing is untouched");
3094        assert!(with.contains("# Persona (tone only)"));
3095        assert!(with.contains("TONE ONLY"));
3096        assert!(with.contains("task ids"));
3097        assert!(with.contains("write policy"));
3098        assert!(with.contains("`magi task add`"));
3099        assert!(with.contains("Rei Ayanami"));
3100    }
3101
3102    #[test]
3103    fn a_talk_written_before_personas_still_loads() {
3104        let (tmp, talks) = store();
3105        let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3106        let cfg = config(spec);
3107        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3108        let path = talks.root.join(format!("{}.json", talk.id));
3109        let mut v: serde_json::Value =
3110            serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3111        v.as_object_mut().unwrap().remove("persona");
3112        v.as_object_mut().unwrap().remove("persona_dirty");
3113        std::fs::write(&path, v.to_string()).unwrap();
3114        let loaded = talks.get(&talk.id).expect("old record loads");
3115        assert_eq!(loaded.persona, "");
3116        assert!(!loaded.persona_dirty);
3117    }
3118
3119    #[tokio::test]
3120    async fn a_persona_switch_notes_marks_dirty_and_updates_the_next_turn_once() {
3121        let (tmp, talks) = store();
3122        let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3123        let cfg = config(spec);
3124        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3125        let lease = talks.claim_turn(&talk.id).unwrap().unwrap();
3126        say(&lease, &mut talk, &talks, &cfg, "hello", Vec::new())
3127            .await
3128            .expect("first turn");
3129        assert!(!talk.turns[1].body.contains("Persona"));
3130
3131        assert!(switch_persona(&mut talk, &talks, "misato").expect("switch"));
3132        assert!(!switch_persona(&mut talk, &talks, "misato").expect("same"));
3133        assert!(talk.persona_dirty);
3134        assert!(talk.turns.last().unwrap().body.contains("persona changed"));
3135
3136        say(&lease, &mut talk, &talks, &cfg, "next", Vec::new())
3137            .await
3138            .expect("turn");
3139        let prompt = &talk.turns.last().unwrap().body;
3140        assert!(prompt.contains("# Persona update"), "{prompt}");
3141        assert!(prompt.contains("Misato Katsuragi"));
3142        assert!(!talk.persona_dirty, "cleared after a successful turn");
3143
3144        say(&lease, &mut talk, &talks, &cfg, "again", Vec::new())
3145            .await
3146            .expect("turn");
3147        assert!(!talk.turns.last().unwrap().body.contains("# Persona update"));
3148
3149        assert!(switch_persona(&mut talk, &talks, "default").expect("back"));
3150        assert_eq!(talk.persona, "");
3151        say(&lease, &mut talk, &talks, &cfg, "plain", Vec::new())
3152            .await
3153            .expect("turn");
3154        assert!(
3155            talk.turns
3156                .last()
3157                .unwrap()
3158                .body
3159                .contains("turned the persona off")
3160        );
3161
3162        // Without a resumable session every turn re-sends the persona.
3163        assert!(switch_persona(&mut talk, &talks, "rei").expect("rei"));
3164        let mut no_sessions = cfg.clone();
3165        no_sessions.graph.sessions = false;
3166        say(&lease, &mut talk, &talks, &no_sessions, "one", Vec::new())
3167            .await
3168            .expect("turn");
3169        say(&lease, &mut talk, &talks, &no_sessions, "two", Vec::new())
3170            .await
3171            .expect("turn");
3172        let last = &talk.turns.last().unwrap().body;
3173        assert!(!talk.persona_dirty);
3174        assert!(last.contains("# Persona (tone only)"), "{last}");
3175        assert!(last.contains("Rei Ayanami"));
3176    }
3177
3178    #[tokio::test]
3179    async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
3180        let (tmp, talks) = store();
3181        let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3182        let cfg = config(spec);
3183        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3184
3185        say(
3186            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3187            &mut talk,
3188            &talks,
3189            &cfg,
3190            "what does the queue module do?",
3191            Vec::new(),
3192        )
3193        .await
3194        .expect("first turn");
3195        let first_prompt = &talk.turns[1].body;
3196        assert!(first_prompt.contains("magi task add --solo"));
3197        assert!(first_prompt.contains("what does the queue module do?"));
3198
3199        say(
3200            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3201            &mut talk,
3202            &talks,
3203            &cfg,
3204            "and how is it locked?",
3205            Vec::new(),
3206        )
3207        .await
3208        .expect("second turn");
3209        let second_prompt = &talk.turns[3].body;
3210        assert!(
3211            !second_prompt.contains("magi task add --solo"),
3212            "the briefing is sent once, not on every turn: {second_prompt}"
3213        );
3214        assert!(second_prompt.contains("and how is it locked?"));
3215    }
3216
3217    #[tokio::test]
3218    async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
3219        let (tmp, talks) = store();
3220        let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3221        let mut b = a.clone();
3222        b.id = "other".to_owned();
3223        let mut cfg = config(a.clone());
3224        cfg.agents.push(b.clone());
3225        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
3226        say(
3227            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3228            &mut talk,
3229            &talks,
3230            &cfg,
3231            "remember the walrus",
3232            Vec::new(),
3233        )
3234        .await
3235        .expect("first turn");
3236        let old_session = talk.seat.claude_session.clone();
3237        assert_eq!(talk.seat.turns, 1);
3238
3239        assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
3240        assert_eq!(talk.agent, "other");
3241        assert_eq!(talk.seat.turns, 0);
3242        assert_eq!(talk.seat.agent, "other");
3243        assert_ne!(talk.seat.claude_session, old_session);
3244        let note = talk.turns.last().expect("note");
3245        assert_eq!(note.who, Who::Agent);
3246        assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3247        assert!(note.body.contains("changed from"), "{}", note.body);
3248        assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
3249
3250        let before = talk.turns.len();
3251        assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
3252        assert_eq!(talk.turns.len(), before, "a no-op writes no note");
3253
3254        say(
3255            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3256            &mut talk,
3257            &talks,
3258            &cfg,
3259            "what did I say?",
3260            Vec::new(),
3261        )
3262        .await
3263        .expect("turn after switch");
3264        let prompt = &talk.turns.last().expect("reply").body;
3265        assert!(prompt.contains("remember the walrus"), "{prompt}");
3266        assert!(prompt.contains("## magi"), "{prompt}");
3267        assert!(prompt.contains("what did I say?"), "{prompt}");
3268    }
3269
3270    #[tokio::test]
3271    async fn say_appends_the_operator_turn_then_the_agent_turn() {
3272        let (tmp, talks) = store();
3273        let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3274        let cfg = config(spec);
3275        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3276
3277        say(
3278            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3279            &mut talk,
3280            &talks,
3281            &cfg,
3282            "can I rename this function?",
3283            Vec::new(),
3284        )
3285        .await
3286        .expect("say");
3287
3288        assert_eq!(talk.turns.len(), 2);
3289        assert_eq!(talk.turns[0].who, Who::Operator);
3290        assert_eq!(talk.turns[0].body, "can I rename this function?");
3291        assert_eq!(talk.turns[1].who, Who::Agent);
3292        assert_eq!(talk.turns[1].body, "go ahead");
3293        assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
3294    }
3295
3296    #[tokio::test]
3297    async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
3298        let (tmp, talks) = store();
3299        let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
3300        let cfg = config(spec);
3301        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3302
3303        let err = say(
3304            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3305            &mut talk,
3306            &talks,
3307            &cfg,
3308            "check the tests",
3309            Vec::new(),
3310        )
3311        .await
3312        .expect_err("a turn with no answer is an error");
3313        assert!(err.to_string().contains("no answer"), "{err}");
3314
3315        let on_disk = talks.get(&talk.id).expect("get");
3316        assert_eq!(on_disk.turns.len(), 2);
3317        assert_eq!(on_disk.turns[0].body, "check the tests");
3318        let note = &on_disk.turns[1];
3319        assert_eq!(note.who, Who::Agent);
3320        assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3321        assert!(note.body.contains("your message is saved"));
3322    }
3323
3324    /// The failure this stands in for: a reader elsewhere briefly has the
3325    /// talk file open right when `turn` tries to save the reply, and the
3326    /// write-then-rename fails once or twice before the reader lets go.
3327    /// `write_atomic`'s own retries must absorb that with nobody the wiser -
3328    /// no gap in the transcript, no dropped turn.
3329    #[tokio::test]
3330    async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
3331        let (tmp, talks) = store();
3332        let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3333        let cfg = config(spec);
3334        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3335
3336        let text =
3337            record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
3338        // One fewer failure than `write_atomic` will retry through, so the
3339        // very last attempt must succeed.
3340        failpoint::force_put_failures(PUT_RETRIES - 1);
3341        respond(
3342            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3343            &mut talk,
3344            &talks,
3345            &cfg,
3346            &text,
3347        )
3348        .await
3349        .expect("respond must survive a write failure its own retries can outlast");
3350
3351        assert_eq!(talk.turns.len(), 2);
3352        assert_eq!(talk.turns[1].who, Who::Agent);
3353        assert_eq!(talk.turns[1].body, "go ahead");
3354        let on_disk = talks.get(&talk.id).expect("get");
3355        assert_eq!(
3356            on_disk.turns, talk.turns,
3357            "the reply must reach disk despite the early write failures"
3358        );
3359    }
3360
3361    /// When the write-then-rename never recovers - standing in for a disk
3362    /// that stays unwritable rather than a reader that eventually lets go -
3363    /// the reply must not disappear without a trace the way it did in the
3364    /// real incident this repository saw: no error on the phone, no note in
3365    /// the transcript, and the turn simply gone from `talks/<id>.json`.
3366    #[tokio::test]
3367    async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
3368        let (tmp, talks) = store();
3369        let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3370        let cfg = config(spec);
3371        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3372
3373        let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3374        // Exactly enough forced failures to exhaust the reply's own retries;
3375        // the shorter note that replaces it then saves cleanly, which is the
3376        // common case this exercises - a large write racing something,
3377        // followed by a small one that does not.
3378        failpoint::force_put_failures(PUT_RETRIES);
3379        let err = respond(
3380            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3381            &mut talk,
3382            &talks,
3383            &cfg,
3384            &text,
3385        )
3386        .await
3387        .expect_err("a reply that cannot be saved must be reported, not swallowed");
3388        assert!(err.to_string().contains("could not be saved"), "{err}");
3389
3390        let on_disk = talks.get(&talk.id).expect("get");
3391        assert_eq!(
3392            on_disk.turns.len(),
3393            2,
3394            "the operator turn plus a visible note"
3395        );
3396        assert_eq!(on_disk.turns[0].body, "check the tests");
3397        let note = &on_disk.turns[1];
3398        assert_eq!(note.who, Who::Agent);
3399        assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3400        assert!(
3401            note.body.contains("could not be saved"),
3402            "the operator must be told the reply is missing, not left staring \
3403             at a gap with no explanation: {}",
3404            note.body
3405        );
3406        assert_eq!(
3407            talk.turns, on_disk.turns,
3408            "the in-memory talk must match what actually landed on disk"
3409        );
3410
3411        // The generated answer itself must still be recoverable, not merely
3412        // reported as lost.
3413        let artifacts = talks.artifacts_of(&talk.id);
3414        let stash = std::fs::read_dir(&artifacts)
3415            .expect("artifacts dir")
3416            .filter_map(|e| e.ok())
3417            .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
3418            .expect("a stash file for the lost reply");
3419        let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
3420        assert_eq!(stashed, "go ahead");
3421
3422        // Losing the reply must not also lose the seat. The CLI took a turn
3423        // and consumed this seat's session id; if the note's write left the
3424        // record claiming otherwise, the next turn would re-open a session
3425        // the CLI is already holding - the `20260907-011805-fb57` desync -
3426        // and would re-send the whole briefing besides. Both decisions read
3427        // the seat straight off disk (`agent::has_session` and `turn`'s own
3428        // `seat.turns == 0` branch), so this is the field that has to match.
3429        assert_eq!(
3430            on_disk.seat.turns, 1,
3431            "the note's write must carry the turn the CLI actually took"
3432        );
3433        assert_eq!(
3434            on_disk.seat.claude_session, talk.seat.claude_session,
3435            "the session id handed to the CLI must survive the failed reply"
3436        );
3437        assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
3438        assert!(
3439            agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
3440            "the next turn must resume, not open the same session id twice"
3441        );
3442    }
3443
3444    /// Even the note can fail to save, if the disk stays unwritable for long
3445    /// enough. `respond` must still report the failure rather than pretend
3446    /// the turn succeeded, and must not leave the in-memory `talk` claiming
3447    /// a turn that never reached disk.
3448    #[tokio::test]
3449    async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
3450        let (tmp, talks) = store();
3451        let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3452        let cfg = config(spec);
3453        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3454
3455        let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3456        // Enough forced failures to exhaust the retries for both the reply
3457        // and the note that would have replaced it.
3458        failpoint::force_put_failures(PUT_RETRIES * 2);
3459        let err = respond(
3460            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3461            &mut talk,
3462            &talks,
3463            &cfg,
3464            &text,
3465        )
3466        .await
3467        .expect_err("neither the reply nor the note could be saved");
3468        assert!(err.to_string().contains("could not be saved"), "{err}");
3469
3470        assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
3471        let on_disk = talks.get(&talk.id).expect("get");
3472        assert_eq!(on_disk.turns.len(), 1);
3473
3474        // Nothing at all reached disk, so the seat could not either: the CLI
3475        // took a turn the record does not know about. That is pinned here as
3476        // the known cost of a file that cannot be written twice over, not as
3477        // something this branch could do better - the only way to record the
3478        // seat is the write that just failed. It is also the point where
3479        // this meets `20260907-011805-fb57`: a turn taken before some later
3480        // write lands would re-open a session id the CLI already holds. The
3481        // in-memory seat keeps the truth the CLI reported, which is why it is
3482        // not wound back to match.
3483        assert_eq!(
3484            on_disk.seat.turns, 0,
3485            "an unwritable file cannot record the turn the CLI took"
3486        );
3487        assert_eq!(
3488            talk.seat.turns, 1,
3489            "the in-memory seat still reports the turn the CLI actually took"
3490        );
3491        assert_eq!(
3492            on_disk.seat.claude_session, talk.seat.claude_session,
3493            "the session id was minted at `begin` and never changes here"
3494        );
3495    }
3496
3497    /// An attachment lets the operator send an otherwise-empty message, and
3498    /// its absolute path is what actually reaches the agent's prompt - here
3499    /// on the very first turn, where it has to share the briefing.
3500    #[tokio::test]
3501    async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
3502        let (tmp, talks) = store();
3503        let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3504        let cfg = config(spec);
3505        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3506
3507        let att = talks
3508            .put_attachment(
3509                &talk.id,
3510                "image/png",
3511                "screenshot.png",
3512                b"pretend-png-bytes",
3513            )
3514            .expect("put attachment");
3515
3516        say(
3517            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3518            &mut talk,
3519            &talks,
3520            &cfg,
3521            "",
3522            vec![att.clone()],
3523        )
3524        .await
3525        .expect("an empty body with an attachment is still a turn");
3526
3527        let operator_turn = &talk.turns[0];
3528        assert_eq!(operator_turn.who, Who::Operator);
3529        assert_eq!(operator_turn.body, "");
3530        assert_eq!(operator_turn.attachments, vec![att.clone()]);
3531
3532        let prompt = &talk.turns[1].body;
3533        let expected_path = talks
3534            .attachments_dir(&talk.id)
3535            .join(format!("{}.png", att.id));
3536        assert!(
3537            prompt.contains(&expected_path.display().to_string()),
3538            "the agent must be told the attachment's absolute path: {prompt}"
3539        );
3540        assert!(prompt.contains("image/png"), "and its mime: {prompt}");
3541    }
3542
3543    /// See `chat`'s test of the same name: `run::home()` returns a bare
3544    /// relative `PathBuf` verbatim when `MAGI_HOME` is set to a relative
3545    /// path, so a `Talks` store built on it has a relative `root` too. That
3546    /// is fine for this store's own I/O, which runs in this process against
3547    /// this process's cwd, but `attachment_path` hands its result to a
3548    /// *different* process invoked with `cwd: &talk.repo` - an uncorrected
3549    /// relative path would resolve against the repository instead of
3550    /// wherever the attachment actually landed.
3551    #[test]
3552    fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
3553        let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
3554        let att = Attachment {
3555            id: "0".repeat(32),
3556            name: "shot.png".to_owned(),
3557            mime: "image/png".to_owned(),
3558            bytes: 3,
3559        };
3560        let path = talks
3561            .attachment_path("some-talk-id", &att)
3562            .expect("a supported mime always yields a path");
3563        assert!(
3564            path.is_absolute(),
3565            "must be absolute even off a relative store root: {}",
3566            path.display()
3567        );
3568    }
3569
3570    #[tokio::test]
3571    async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
3572        // `[graph] timeout_talk` must be the number this module actually
3573        // waits, not a leftover hardcoded fifteen minutes - so the mock
3574        // sleeps past a deliberately tiny override and the failure note is
3575        // checked against that same override, not the old default.
3576        let (tmp, talks) = store();
3577        let slow = mock_agent(
3578            tmp.path(),
3579            "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
3580            BTreeMap::new(),
3581        );
3582        let mut cfg = config(slow);
3583        cfg.graph.timeout_talk = 1;
3584        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3585
3586        let err = say(
3587            &talks.claim_turn(&talk.id).unwrap().unwrap(),
3588            &mut talk,
3589            &talks,
3590            &cfg,
3591            "check the tests",
3592            Vec::new(),
3593        )
3594        .await
3595        .expect_err("a turn that never answers is an error");
3596        assert!(
3597            err.to_string().contains("did not answer within 1s"),
3598            "{err}"
3599        );
3600
3601        let on_disk = talks.get(&talk.id).expect("get");
3602        let note = on_disk.turns.last().expect("a note turn was recorded");
3603        assert!(
3604            note.body.contains("did not answer within 1s"),
3605            "the transcript must show the configured timeout: {}",
3606            note.body
3607        );
3608    }
3609
3610    #[test]
3611    fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
3612        let (tmp, talks) = store();
3613        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3614        let cfg = config(spec);
3615        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3616
3617        close(&mut talk, &talks).expect("close");
3618        assert_eq!(talk.status, TalkStatus::Closed);
3619        close(&mut talk, &talks).expect("closing twice is not an error");
3620
3621        let err =
3622            record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
3623        assert!(err.to_string().contains("closed"));
3624        let _ = &cfg; // config kept only to build the agent above
3625    }
3626
3627    #[tokio::test]
3628    async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3629        let (tmp, talks) = store();
3630        let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3631        let cfg = config(spec);
3632        // The in-flight turn's own handle: loaded once, the way a spawned
3633        // background task in `web::talk_say` holds one for the whole turn.
3634        let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3635
3636        // The operator closes the conversation through a *different* handle
3637        // while the turn above is still running - exactly what a close typed
3638        // on the phone while an agent is mid-answer looks like.
3639        let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
3640        close(&mut closed_elsewhere, &talks).expect("close");
3641        assert_eq!(
3642            talks.get(&in_flight.id).expect("reread").status,
3643            TalkStatus::Closed,
3644            "the close landed on disk before the turn finished"
3645        );
3646
3647        // The turn's own handle still says `open` - it was loaded before the
3648        // close - and finishing it must not resurrect the conversation the
3649        // operator already ended.
3650        assert_eq!(in_flight.status, TalkStatus::Open);
3651        respond(
3652            &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
3653            &mut in_flight,
3654            &talks,
3655            &cfg,
3656            "one more question",
3657        )
3658        .await
3659        .expect("the turn itself still completes");
3660
3661        let on_disk = talks.get(&in_flight.id).expect("reread");
3662        assert_eq!(
3663            on_disk.status,
3664            TalkStatus::Closed,
3665            "a close must stick even when a turn that started before it finishes after it"
3666        );
3667        // The reply is not lost either: a turn already in flight when the
3668        // operator closed still gets its answer recorded.
3669        assert!(
3670            on_disk.turns.iter().any(|t| t.body == "here you go"),
3671            "the in-flight turn's own reply is still recorded: {:?}",
3672            on_disk.turns
3673        );
3674    }
3675
3676    #[test]
3677    fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
3678        let (tmp, talks) = store();
3679        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3680        let cfg = config(spec);
3681        // The handle `web::talk_say` would have read before awaiting config
3682        // discovery, then carried across that await into `record`.
3683        let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3684
3685        // The operator closes the conversation through a *different* handle
3686        // in the gap between that read and the call to `record` below.
3687        let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
3688        close(&mut closed_elsewhere, &talks).expect("close");
3689        assert_eq!(
3690            talks.get(&stale.id).expect("reread").status,
3691            TalkStatus::Closed,
3692            "the close landed on disk before record was called"
3693        );
3694
3695        // The stale handle still says `open` - it was loaded before the
3696        // close - so a `record` that trusted it would append a turn and
3697        // write the conversation back open, undoing the close.
3698        assert_eq!(stale.status, TalkStatus::Open);
3699        let err = record(&mut stale, &talks, "still there?", Vec::new())
3700            .expect_err("a close that landed first must be honored, not overwritten");
3701        assert!(err.to_string().contains("closed"));
3702
3703        let on_disk = talks.get(&stale.id).expect("reread");
3704        assert_eq!(
3705            on_disk.status,
3706            TalkStatus::Closed,
3707            "record must not resurrect a conversation closed while its snapshot was stale"
3708        );
3709        assert!(
3710            on_disk.turns.is_empty(),
3711            "the rejected turn must not have been appended: {:?}",
3712            on_disk.turns
3713        );
3714        let _ = &cfg; // config kept only to build the agent above
3715    }
3716
3717    #[test]
3718    fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
3719        let (tmp, talks) = store();
3720        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3721        let cfg = config(spec);
3722        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3723
3724        // Hold the same guard `record`'s read-modify-write section holds for
3725        // the whole of its own read-then-write, standing in for `record`
3726        // being paused between its read and its `put`.
3727        let held = talks.guard().unwrap();
3728
3729        let talks2 = talks.clone();
3730        let id = talk.id.clone();
3731        let closing = std::thread::spawn(move || {
3732            let mut talk = talks2.get(&id).expect("get");
3733            close(&mut talk, &talks2).expect("close");
3734        });
3735
3736        std::thread::sleep(Duration::from_millis(50));
3737        assert!(
3738            !closing.is_finished(),
3739            "close must wait for the guard, not read and write while it is held - \
3740             a re-read alone narrows this window without closing it"
3741        );
3742
3743        drop(held);
3744        closing.join().expect("close thread panicked");
3745
3746        assert_eq!(
3747            talks.get(&talk.id).expect("reread").status,
3748            TalkStatus::Closed,
3749            "once the guard is free, close still lands"
3750        );
3751        let _ = &cfg; // config kept only to build the agent above
3752    }
3753
3754    #[test]
3755    fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
3756        let (tmp, talks) = store();
3757        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3758        let cfg = config(spec);
3759        let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3760
3761        close(&mut talk, &talks).expect("close");
3762        assert_eq!(talk.status, TalkStatus::Closed);
3763
3764        reopen(&mut talk, &talks).expect("reopen");
3765        assert_eq!(talk.status, TalkStatus::Open);
3766        assert_eq!(
3767            talks.get(&talk.id).expect("reread").status,
3768            TalkStatus::Open
3769        );
3770
3771        // Idempotent: reopening an already-open talk is not an error.
3772        reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
3773        assert_eq!(talk.status, TalkStatus::Open);
3774
3775        record(&mut talk, &talks, "one more thing", Vec::new())
3776            .expect("a reopened talk takes turns again");
3777        let _ = &cfg; // config kept only to build the agent above
3778    }
3779
3780    #[test]
3781    fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
3782        let (tmp, talks) = store();
3783        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3784        let cfg = config(spec);
3785        let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3786
3787        let artifacts = talks.artifacts_of(&talk.id);
3788        std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
3789        std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
3790
3791        talks.remove(&talk.id).expect("remove");
3792        assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
3793        assert!(!artifacts.is_dir(), "the artifacts directory is gone");
3794        assert!(
3795            talks.get(&talk.id).is_err(),
3796            "a removed talk cannot be read back"
3797        );
3798
3799        let err = talks
3800            .remove("nonexistent-id")
3801            .expect_err("unknown id refused");
3802        assert!(err.to_string().contains("no talk matches"), "{err}");
3803        let _ = &cfg; // config kept only to build the agent above
3804    }
3805
3806    #[tokio::test]
3807    async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3808        let (tmp, talks) = store();
3809        let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3810        let cfg = config(spec);
3811        // The in-flight turn's own handle, loaded before the delete lands -
3812        // the same shape as the matching close test above.
3813        let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3814
3815        talks.remove(&in_flight.id).expect("remove");
3816        assert!(
3817            talks.get(&in_flight.id).is_err(),
3818            "the delete landed on disk before the turn finished"
3819        );
3820
3821        // The turn's own handle has no way to know the record is gone -
3822        // finishing it must not write the file back into existence.
3823        respond(
3824            &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
3825            &mut in_flight,
3826            &talks,
3827            &cfg,
3828            "one more question",
3829        )
3830        .await
3831        .expect("the turn itself still completes rather than erroring");
3832
3833        assert!(
3834            talks.get(&in_flight.id).is_err(),
3835            "a delete must stick even when a turn that started before it finishes after it"
3836        );
3837    }
3838
3839    #[test]
3840    fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
3841        let (tmp, talks) = store();
3842        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3843        let cfg = config(spec);
3844        // The handle `web::talk_say` would have read before awaiting config
3845        // discovery, then carried across that await into `record`.
3846        let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3847
3848        talks.remove(&stale.id).expect("remove");
3849
3850        // The stale handle has no way to know the record is gone - a
3851        // `record` that trusted it would append a turn and write the
3852        // conversation back into existence.
3853        let err = record(&mut stale, &talks, "still there?", Vec::new())
3854            .expect_err("a delete that landed first must be honored, not overwritten");
3855        assert!(err.to_string().contains("deleted"), "{err}");
3856
3857        assert!(
3858            talks.get(&stale.id).is_err(),
3859            "record must not resurrect a conversation deleted while its snapshot was stale"
3860        );
3861        let _ = &cfg; // config kept only to build the agent above
3862    }
3863
3864    #[test]
3865    fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
3866        let (tmp, talks) = store();
3867        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3868        let cfg = config(spec);
3869        // `web::talk_close` loads `talk` and calls `close` right after - this
3870        // stands in for a delete landing in that gap.
3871        let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3872
3873        talks.remove(&stale.id).expect("remove");
3874
3875        // The stale handle has no way to know the record is gone - a `close`
3876        // that fell back to it would write the conversation back into
3877        // existence, closed.
3878        let err = close(&mut stale, &talks)
3879            .expect_err("a delete that landed first must be honored, not overwritten");
3880        assert!(err.to_string().contains("deleted"), "{err}");
3881
3882        assert!(
3883            talks.get(&stale.id).is_err(),
3884            "close must not resurrect a conversation deleted while its snapshot was stale"
3885        );
3886        let _ = &cfg; // config kept only to build the agent above
3887    }
3888
3889    #[test]
3890    fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
3891        let (tmp, talks) = store();
3892        let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3893        let cfg = config(spec);
3894        // `web::talk_reopen` loads `talk` and calls `reopen` right after -
3895        // this stands in for a delete landing in that gap.
3896        let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3897        close(&mut stale, &talks).expect("close");
3898
3899        talks.remove(&stale.id).expect("remove");
3900
3901        // The stale handle has no way to know the record is gone - a
3902        // `reopen` that fell back to it would write the conversation back
3903        // into existence, open.
3904        let err = reopen(&mut stale, &talks)
3905            .expect_err("a delete that landed first must be honored, not overwritten");
3906        assert!(err.to_string().contains("deleted"), "{err}");
3907
3908        assert!(
3909            talks.get(&stale.id).is_err(),
3910            "reopen must not resurrect a conversation deleted while its snapshot was stale"
3911        );
3912        let _ = &cfg; // config kept only to build the agent above
3913    }
3914
3915    #[test]
3916    fn list_puts_open_talks_before_closed_ones() {
3917        let (tmp, talks) = store();
3918        let make = |id: &str, status: TalkStatus| {
3919            let mut t = Talk {
3920                schema: SCHEMA,
3921                id: id.to_owned(),
3922                repo: tmp.path().to_owned(),
3923                agent: "mock".to_owned(),
3924                status,
3925                turns: Vec::new(),
3926                pending: String::new(),
3927                pending_attachments: Vec::new(),
3928                fallback: false,
3929                persona: String::new(),
3930                persona_dirty: false,
3931                created_at: Timestamp::now(),
3932                updated_at: Timestamp::now(),
3933                seat: SeatState::new(SEAT, "mock", 7),
3934            };
3935            talks.put(&mut t).expect("put");
3936        };
3937        make("20260901-000000-0001", TalkStatus::Open);
3938        make("20260902-000000-0002", TalkStatus::Open);
3939        make("20260903-000000-0003", TalkStatus::Closed);
3940
3941        let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
3942        assert_eq!(
3943            ids,
3944            [
3945                "20260902-000000-0002",
3946                "20260901-000000-0001",
3947                "20260903-000000-0003"
3948            ]
3949        );
3950        assert_eq!(talks.count_open(), 2);
3951    }
3952
3953    #[test]
3954    fn tasks_of_finds_only_this_talks_own_tasks() {
3955        let dir = tempfile::tempdir().expect("tempdir");
3956        let queue = Queue::at(dir.path().join("queue"));
3957
3958        let mut mine = Task::new(
3959            "rework the loader".to_owned(),
3960            "rework the loader".to_owned(),
3961            PathBuf::from("/repo"),
3962            Source::Agent {
3963                run: "20260904-014455-ab12".to_owned(),
3964                node: "chat".to_owned(),
3965            },
3966        );
3967        queue.put(&mut mine).expect("put mine");
3968
3969        let mut theirs = Task::new(
3970            "unrelated".to_owned(),
3971            "unrelated".to_owned(),
3972            PathBuf::from("/repo"),
3973            Source::Agent {
3974                run: "20260904-090000-zz99".to_owned(),
3975                node: "implement".to_owned(),
3976            },
3977        );
3978        queue.put(&mut theirs).expect("put theirs");
3979
3980        let mut human = Task::new(
3981            "typed by hand".to_owned(),
3982            "typed by hand".to_owned(),
3983            PathBuf::from("/repo"),
3984            Source::Human,
3985        );
3986        queue.put(&mut human).expect("put human");
3987
3988        let found = tasks_of(&queue, "20260904-014455-ab12");
3989        assert_eq!(found.len(), 1);
3990        assert_eq!(found[0].id, mine.id);
3991    }
3992
3993    #[test]
3994    fn the_briefing_names_solo_task_add() {
3995        let brief = briefing(Path::new("/repo"), "en", false);
3996        assert!(brief.contains("magi task add --solo"));
3997        assert!(brief.contains("/repo"));
3998        assert!(!brief.contains("Hold this conversation in"));
3999    }
4000
4001    /// Talk fixes `repo` at the directory the conversation was opened in, so
4002    /// an agent asked to change some other checkout has no path to it unless
4003    /// the briefing itself says `--repo` can take a short name - see
4004    /// `resolve_repo_by_name` in `src/main.rs`, which is what actually
4005    /// resolves it.
4006    #[test]
4007    fn the_briefing_explains_targeting_a_different_repository_by_name() {
4008        let brief = briefing(Path::new("/repo"), "en", false);
4009        assert!(brief.contains("--repo does not have to be a full path"));
4010        assert!(brief.contains("owner/repo"));
4011        assert!(brief.contains("magi repos"));
4012        assert!(brief.contains("ask the operator"));
4013    }
4014
4015    #[test]
4016    fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
4017        let brief = briefing(Path::new("/repo"), "en", false);
4018        assert!(brief.contains("--attach <path>"), "{brief}");
4019        assert!(brief.contains("deleting this conversation"), "{brief}");
4020    }
4021
4022    #[test]
4023    fn the_briefing_names_the_language_when_it_is_not_english() {
4024        let brief = briefing(Path::new("/repo"), "Japanese", false);
4025        assert!(brief.contains("Hold this conversation in Japanese"));
4026    }
4027
4028    #[test]
4029    fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
4030        let read_only = briefing(Path::new("/repo"), "en", false);
4031        assert!(read_only.contains("Do not write files"));
4032        assert!(!read_only.contains("allow_write"));
4033
4034        let writable = briefing(Path::new("/repo"), "en", true);
4035        assert!(!writable.contains("Do not write files"));
4036        assert!(writable.contains("allow_write = true"));
4037        // Still names the queue for anything past a small named edit, and
4038        // still tells the agent to report what it changed.
4039        assert!(writable.contains("magi task add --solo"));
4040        assert!(writable.contains("say plainly what you"));
4041    }
4042}