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