Skip to main content

mermaid_cli/effect/
mod.rs

1//! The effect runner: dispatches `Cmd` values into tokio tasks.
2//!
3//! There are exactly two places in the codebase that spawn a tokio
4//! task: this module and tests. Everywhere else asks the
5//! reducer to return a `Cmd`, and the runner handles it. That
6//! centralization is what makes structured concurrency per turn
7//! actually work — nothing can accidentally spawn a detached task
8//! that outlives the turn it was started for.
9//!
10//! Architecture:
11//!
12//! ```text
13//!   main loop ── reducer ── Cmd ── dispatch ── EffectRunner
14//!                                                 ├── TurnScope(turn A) ── JoinSet
15//!                                                 ├── TurnScope(turn B) ── JoinSet
16//!                                                 └── detached effects (Save, Exit, …)
17//!                                                       ↓
18//!                                              Msg via mpsc::Sender<Msg>
19//!                                                       ↓
20//!                                                 main loop (next iteration)
21//! ```
22//!
23//! The runner dispatches every `Cmd` variant to a real handler —
24//! model streaming (`CallModel` → `ModelProvider::chat`), tool
25//! execution (`ExecuteTool` → `ToolExecutor::execute`), persistence
26//! (`SaveConversation`, `LoadConversation`, `PersistLastModel`,
27//! `PersistReasoningFor`), MCP lifecycle
28//! (`InitMcpServers`, `StopMcpServer`), local side-effects
29//! (`WriteImageToTemp`, `OpenInSystem`, `PullOllamaModel`,
30//! `SetTerminalTitle`). Cancellation flows
31//! through `Cmd::CancelScope(TurnId)` → the scope's
32//! `CancellationToken`.
33
34mod config_watch;
35mod turn_scope;
36
37use std::collections::HashMap;
38use std::collections::VecDeque;
39use std::path::PathBuf;
40use std::sync::Arc;
41use std::sync::Mutex;
42
43use tokio::sync::mpsc;
44
45use crate::providers::ctx::{ExecContext, StreamContext};
46use crate::providers::{ProviderFactory, StreamEvent, ToolRegistry};
47use mermaid_domain::{
48    Cmd, CompactionRequest, CompactionResult, CompactionTrigger, Msg, Query, QueryResult, TurnId,
49};
50use mermaid_domain::{Config, MemoryConfig};
51
52pub use turn_scope::TurnScope;
53
54#[cfg(not(test))]
55const CANCEL_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
56#[cfg(test)]
57const CANCEL_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(50);
58
59/// F38: how many recently-cancelled `TurnId`s to remember as tombstones.
60/// Turn ids are strictly monotonic and never reused, so a stray turn-scoped
61/// `Cmd` for a cancelled turn can only ever be a post-cancel straggler that
62/// lands within a few turns of the cancel. A small bounded ring is plenty;
63/// older entries age out so the set never grows across a long session.
64const CANCELLED_TOMBSTONE_CAP: usize = 256;
65
66/// Single channel back to the reducer. `EffectRunner` holds the
67/// sender; every spawned task clones this so it can emit `Msg` as
68/// work progresses. Bounded capacity applies natural backpressure —
69/// if the main loop can't keep up, the provider's streaming send
70/// `.await`s and the whole pipeline throttles.
71pub type MsgSender = mpsc::Sender<Msg>;
72
73/// Bounded channel capacity for the effect → reducer stream. 512 is
74/// generous — a single streaming chunk fits comfortably, and the
75/// main loop drains at ~60 Hz so backlog rarely grows. Bigger wastes
76/// RAM; smaller introduces spurious backpressure on bursty tool
77/// output.
78pub const MSG_CHANNEL_CAPACITY: usize = 512;
79
80/// Translate a domain `CompactionEvent` into the durable row.
81///
82/// The two types share a name-adjacent concept and almost nothing else: the
83/// domain event carries 14 fields describing what compaction did, the row
84/// carries 9 describing what a later reader needs. Only `id` and
85/// `archive_path` overlap. This lived as an anonymous struct literal inline in
86/// `persist_compaction`, which is where a field mapping goes to rot unnoticed.
87///
88/// It stays in `effect`, not in the pure core: "domain value -> SQLite row" is
89/// precisely what this layer is for, and the orphan rule permitting it
90/// elsewhere is not a reason to put row-shaped knowledge in the reducer.
91fn compaction_row(
92    record: &mermaid_domain::CompactionEvent,
93    archive_path: &std::path::Path,
94    task_id: Option<String>,
95    session_id: String,
96) -> mermaid_runtime::NewCompaction {
97    mermaid_runtime::NewCompaction {
98        id: Some(record.id.clone()),
99        task_id,
100        session_id: Some(session_id),
101        source_token_estimate: Some(record.before_tokens as i64),
102        summary_token_count: Some(record.summary_tokens as i64),
103        preserved_turns: Some(record.preserved_turn_count as i64),
104        archive_path: Some(archive_path.display().to_string()),
105        verification_status: Some(record.review_status.as_str().to_string()),
106    }
107}
108
109/// Feed the cross-project session index off a successful snapshot save.
110/// Best-effort by design: the row is an index over the files, never the
111/// truth, so a store hiccup must not fail the save that just succeeded.
112/// Feed the cross-project session index off a successful append. Best-effort
113/// by design: the row is an index over the files, never the truth, so a store
114/// hiccup must not fail the save that just succeeded; the daemon rebuilds the
115/// index from disk on its next start regardless.
116fn upsert_session_index(
117    manager: &crate::session::ConversationManager,
118    snapshot: &mermaid_domain::ConversationHistory,
119) {
120    let row = crate::session::session_row(manager.conversations_dir(), snapshot);
121    let _ = mermaid_runtime::with_shared_store(|store| store.sessions().upsert(row));
122}
123
124#[derive(Clone)]
125enum PersistenceJob {
126    Conversation {
127        snapshot: Box<mermaid_domain::ConversationHistory>,
128        events: Vec<mermaid_domain::SessionEvent>,
129    },
130    Compaction(Box<PendingCompactionSave>),
131}
132
133#[derive(Clone)]
134struct PendingCompactionSave {
135    record: mermaid_domain::CompactionEvent,
136    conversation: mermaid_domain::ConversationHistory,
137    events: Vec<mermaid_domain::SessionEvent>,
138    /// Set once the events are durably appended, so a retry after a later
139    /// failure in the same save re-runs only what did not land — appending
140    /// twice would duplicate the boundary in the log.
141    events_appended: bool,
142    task_id: Option<String>,
143}
144
145struct PersistedCompaction {
146    id: String,
147    task_id: Option<String>,
148    session_id: String,
149    archive_path: PathBuf,
150}
151
152/// How many appended events may accumulate before the checkpoint is
153/// rewritten.
154///
155/// The quantity being bounded is REPLAY LENGTH — a checkpoint's whole job is
156/// to cap how much log a resume has to fold — so the trigger counts events
157/// rather than turns or seconds. At ~30 events for a tool-heavy turn this is
158/// roughly seven turns of replay in the worst case, against a rewrite that
159/// used to happen after every single message.
160const CHECKPOINT_EVERY_EVENTS: usize = 200;
161
162struct PersistenceState {
163    workdir: PathBuf,
164    manager: Option<crate::session::ConversationManager>,
165    blocked: HashMap<String, VecDeque<PendingCompactionSave>>,
166    /// Per session: events appended since its last checkpoint, and the
167    /// newest snapshot that has not been written as one. Shutdown flushes
168    /// these, so a clean exit always leaves a current checkpoint.
169    dirty: HashMap<String, DirtySession>,
170    /// Events whose append FAILED, kept in order for the next attempt.
171    /// Without this they are simply gone: the reducer drains its buffer at
172    /// emission, so a dropped batch is never re-offered — which stopped
173    /// mattering the moment the log became the truth rather than a copy.
174    unappended: HashMap<String, Vec<mermaid_domain::SessionEvent>>,
175}
176
177struct DirtySession {
178    snapshot: mermaid_domain::ConversationHistory,
179    events_since_checkpoint: usize,
180}
181
182impl PersistenceState {
183    fn new(workdir: PathBuf) -> Self {
184        Self {
185            workdir,
186            manager: None,
187            blocked: HashMap::new(),
188            dirty: HashMap::new(),
189            unappended: HashMap::new(),
190        }
191    }
192
193    fn manager(&mut self) -> anyhow::Result<&crate::session::ConversationManager> {
194        if self.manager.is_none() {
195            self.manager = Some(crate::session::ConversationManager::new(&self.workdir)?);
196        }
197        Ok(self.manager.as_ref().expect("manager initialized"))
198    }
199
200    /// Run one job. Returns every compaction event that persisted durably —
201    /// even when the job as a whole failed — so partially-drained barriers
202    /// still fire their hooks and `SessionSaved`; a dropped event would never
203    /// be re-emitted (its save is already popped).
204    fn process(&mut self, job: PersistenceJob) -> (Vec<PersistedCompaction>, anyhow::Result<()>) {
205        match job {
206            PersistenceJob::Conversation { snapshot, events } => {
207                // Barrier: a still-blocked compaction must persist before any
208                // newer (stripped) conversation snapshot may overwrite the file.
209                let (persisted, retried) = self.retry_blocked(&snapshot.id);
210                if retried.is_err() {
211                    return (persisted, retried);
212                }
213                let saved = self.save_session(*snapshot, events);
214                (persisted, saved)
215            },
216            PersistenceJob::Compaction(save) => {
217                // Queue first, then drain. The boundary event is the only
218                // record that the dropped messages ever existed, so the save
219                // must survive an Err AND a panic in the write path (pop
220                // happens only after success), and it must land behind any
221                // older still-blocked saves (FIFO).
222                let conversation_id = save.conversation.id.clone();
223                self.blocked
224                    .entry(conversation_id.clone())
225                    .or_default()
226                    .push_back(*save);
227                self.retry_blocked(&conversation_id)
228            },
229        }
230    }
231
232    fn retry_blocked(
233        &mut self,
234        conversation_id: &str,
235    ) -> (Vec<PersistedCompaction>, anyhow::Result<()>) {
236        let mut persisted = Vec::new();
237        if !self.blocked.contains_key(conversation_id) {
238            return (persisted, Ok(()));
239        }
240        if let Err(error) = self.manager() {
241            return (persisted, Err(error));
242        }
243        // Disjoint field borrows: the manager stays immutably borrowed while
244        // the queue is drained in place — no per-retry clone of the (large)
245        // pending conversation snapshots.
246        let manager = self.manager.as_ref().expect("manager initialized");
247        let dirty = &mut self.dirty;
248        let queue = self
249            .blocked
250            .get_mut(conversation_id)
251            .expect("checked above");
252        while let Some(save) = queue.front_mut() {
253            match Self::persist_compaction(manager, dirty, save) {
254                // Pop only after a successful write: `persist_compaction` runs
255                // inside `spawn_blocking`, and a panic there must not lose the
256                // save (the mutex is poison-tolerant, so the state survives).
257                Ok(event) => {
258                    persisted.push(event);
259                    queue.pop_front();
260                },
261                Err(error) => return (persisted, Err(error)),
262            }
263        }
264        self.blocked.remove(conversation_id);
265        (persisted, Ok(()))
266    }
267
268    /// Append a save's events, then write the checkpoint only when enough
269    /// have accumulated (see [`CHECKPOINT_EVERY_EVENTS`]).
270    ///
271    /// The append is the save now: it is the write that reaches the truth,
272    /// so a failure keeps the batch for the next attempt rather than
273    /// dropping it, and leaves the checkpoint alone — advancing a cache past
274    /// a log that did not take the events is how a "successful" save loses
275    /// them.
276    fn save_session(
277        &mut self,
278        snapshot: mermaid_domain::ConversationHistory,
279        events: Vec<mermaid_domain::SessionEvent>,
280    ) -> anyhow::Result<()> {
281        let id = snapshot.id.clone();
282        // Anything a previous attempt could not land goes first, in order.
283        let mut batch = self.unappended.remove(&id).unwrap_or_default();
284        batch.extend(events);
285
286        let manager = self.manager()?;
287        if let Err(error) = manager.append_session_events(&snapshot, &batch) {
288            tracing::warn!(
289                id = %id,
290                pending = batch.len(),
291                %error,
292                "session event append failed; holding the events for the next save"
293            );
294            self.unappended.insert(id, batch);
295            return Err(error);
296        }
297        upsert_session_index(manager, &snapshot);
298
299        let entry = self
300            .dirty
301            .entry(id.clone())
302            .or_insert_with(|| DirtySession {
303                snapshot: snapshot.clone(),
304                events_since_checkpoint: 0,
305            });
306        entry.snapshot = snapshot;
307        entry.events_since_checkpoint += batch.len();
308        if entry.events_since_checkpoint >= CHECKPOINT_EVERY_EVENTS {
309            return self.write_checkpoint(&id);
310        }
311        Ok(())
312    }
313
314    /// Materialize a session's checkpoint and clear its dirty counter. A
315    /// session with nothing outstanding is a no-op.
316    fn write_checkpoint(&mut self, id: &str) -> anyhow::Result<()> {
317        let Some(dirty) = self.dirty.remove(id) else {
318            return Ok(());
319        };
320        let manager = self.manager()?;
321        if let Err(error) = manager.save_conversation(&dirty.snapshot) {
322            // Put it back: the events are durable in the log either way, so
323            // this costs a longer replay, not data — but the next save
324            // should still try.
325            self.dirty.insert(id.to_string(), dirty);
326            return Err(error);
327        }
328        Ok(())
329    }
330
331    /// Write every outstanding checkpoint. Called at shutdown so a clean
332    /// exit always leaves a current one, which is what keeps the common
333    /// resume path short.
334    fn flush_checkpoints(&mut self) -> anyhow::Result<()> {
335        let ids: Vec<String> = self.dirty.keys().cloned().collect();
336        let mut first_error = None;
337        for id in ids {
338            if let Err(error) = self.write_checkpoint(&id) {
339                first_error.get_or_insert(error);
340            }
341        }
342        first_error.map_or(Ok(()), Err)
343    }
344
345    fn retry_all_blocked(&mut self) -> (Vec<PersistedCompaction>, anyhow::Result<()>) {
346        let ids: Vec<String> = self.blocked.keys().cloned().collect();
347        let mut persisted = Vec::new();
348        let mut first_error = None;
349        for id in ids {
350            // Keep draining the other conversations' barriers; one
351            // conversation's bad disk state must not strand the rest.
352            let (events, result) = self.retry_blocked(&id);
353            persisted.extend(events);
354            if let Err(error) = result {
355                first_error.get_or_insert(error);
356            }
357        }
358        match first_error {
359            None => (persisted, Ok(())),
360            Some(error) => (persisted, Err(error)),
361        }
362    }
363
364    fn persist_compaction(
365        manager: &crate::session::ConversationManager,
366        dirty: &mut HashMap<String, DirtySession>,
367        save: &mut PendingCompactionSave,
368    ) -> anyhow::Result<PersistedCompaction> {
369        // Order is the whole guarantee. A compaction drops messages from the
370        // conversation, and the only remaining record of them is this
371        // session's log — the earlier `message` events plus the `compaction`
372        // event that marks the boundary. So the append must land BEFORE the
373        // stripped snapshot overwrites the file that still holds the fuller
374        // history, and a failed append must abort the save (this is where
375        // the archive file's `?` used to be), leaving the queue blocked and
376        // the old snapshot intact for the retry.
377        if !save.events_appended {
378            manager.append_session_events(&save.conversation, &save.events)?;
379            save.events_appended = true;
380        }
381        manager.save_conversation(&save.conversation)?;
382        upsert_session_index(manager, &save.conversation);
383        // A compaction always checkpoints: it is a structural boundary, it
384        // is rare, and the transcript it leaves behind is the one a resume
385        // should start from rather than replay its way to.
386        dirty.remove(&save.conversation.id);
387
388        let log_path = manager.event_log_path(&save.conversation.id);
389        let _ = mermaid_runtime::with_shared_store(|store| {
390            store.compactions().create(compaction_row(
391                &save.record,
392                &log_path,
393                save.task_id.clone(),
394                save.conversation.id.clone(),
395            ))
396        });
397
398        Ok(PersistedCompaction {
399            id: save.record.id.clone(),
400            task_id: save.task_id.clone(),
401            session_id: save.conversation.id.clone(),
402            archive_path: log_path,
403        })
404    }
405}
406
407/// Fire the plugin `compaction` hook for one durably persisted archive.
408async fn fire_compaction_hook(event: &PersistedCompaction) {
409    fire_plugin_hooks(
410        "compaction",
411        serde_json::json!({
412            "id": event.id,
413            "task_id": event.task_id,
414            "session_id": event.session_id,
415            "archive_path": event.archive_path.display().to_string(),
416        }),
417    )
418    .await;
419}
420
421/// The runner. One instance per process, constructed by
422/// `app::run` and consumed when the main loop exits.
423pub struct EffectRunner {
424    msg_tx: MsgSender,
425    /// Per-turn scopes. Populated lazily: the first `Cmd` bearing a
426    /// `TurnId` creates a scope; `Cmd::CancelScope` tears it down.
427    /// Empty (drained) scopes are reaped by `reap_empty_scopes`, which
428    /// runs at the top of every `dispatch` call so the map stays
429    /// bounded across long sessions (F12).
430    scopes: HashMap<TurnId, TurnScope>,
431    /// F38: bounded tombstone ring of `TurnId`s whose scope has been
432    /// cancelled+dropped. A turn-scoped `Cmd` (`CallModel` / `ExecuteTool` /
433    /// `CompactConversation`) bearing a tombstoned id is dropped in `dispatch`
434    /// instead of resurrecting a fresh, un-cancelled scope through
435    /// `scope_mut`'s `or_insert_with`. Bounded to `CANCELLED_TOMBSTONE_CAP`.
436    cancelled_turns: VecDeque<TurnId>,
437    /// Detached work (saves, persists, MCP lifecycle) lives here.
438    /// This one set never gets cancelled piecemeal — shutdown drains
439    /// it during `EffectRunner::shutdown`.
440    detached: tokio::task::JoinSet<()>,
441    /// FIFO chain for conversation and compaction writes. Keeping persistence
442    /// separate from `detached` prevents an older compaction snapshot from
443    /// racing a newer normal save and winning last-write-wins.
444    persistence_state: Arc<Mutex<PersistenceState>>,
445    persistence_tail: Option<tokio::task::JoinHandle<()>>,
446    /// MCP manager handle is held elsewhere (`crate::mcp` has a
447    /// `OnceLock` for its global manager); we just note workdir so
448    /// handlers can construct absolute paths.
449    workdir: PathBuf,
450    /// Lazy provider registry. `CallModel` resolves through this.
451    /// Tests that don't care about real providers leave this `None`
452    /// and observe the fallback `UpstreamError` Msg; production
453    /// construction via `with_bindings` sets it.
454    providers: Option<Arc<ProviderFactory>>,
455    /// Shared tool registry. See `providers` — same optionality
456    /// rationale for unit tests.
457    tools: Option<Arc<ToolRegistry>>,
458    /// Durable runtime task that owns work launched by this runner.
459    task_id: Option<String>,
460    /// Interactive TUI runners write OSC 2 terminal-title updates.
461    /// Headless `mermaid run` must suppress them so stdout stays
462    /// machine-readable for JSON/markdown/text output modes.
463    terminal_title_enabled: bool,
464    /// Whether this runner's `shutdown` reaps the PROCESS-GLOBAL MCP manager
465    /// (`crate::mcp::manager_ref`). True only for the top-level runner. A
466    /// subagent's child runner shares the global manager, so it must NOT reap
467    /// it — otherwise the first subagent to finish would kill every MCP
468    /// server out from under the parent for the rest of the session.
469    owns_global_mcp: bool,
470    /// Inline-approval broker. `Some` only for interactive TUI runs (set via
471    /// `with_interactive_approvals`); headless + child runners leave it `None`,
472    /// so the gate falls back to the out-of-band DB-approval flow.
473    approval: Option<crate::providers::ApprovalBroker>,
474    /// Inline-question broker for `ask_user_question`. `Some` only for
475    /// interactive TUI runs (set via `with_interactive_questions`); headless +
476    /// child runners leave it `None`, so the tool proceeds without asking.
477    questions: Option<crate::providers::QuestionBroker>,
478    /// Checklist broker for the task tools. Built unconditionally — unlike
479    /// `questions`, task tracking works headless, and a subagent's child
480    /// runner minting its own broker (bound to the CHILD's msg channel) is
481    /// exactly what isolates its checklist from the parent's.
482    tasks: crate::providers::TaskBroker,
483    /// Abort handle for the background config watcher (#45). It's a perpetual
484    /// loop living in `detached`, so `shutdown` aborts it explicitly before
485    /// draining — otherwise the drain would block on it until the timeout.
486    config_watch: Option<tokio::task::AbortHandle>,
487}
488
489impl EffectRunner {
490    /// Create an unused runner. Pair with `msg_rx` from `channel()`.
491    #[must_use]
492    pub fn new(msg_tx: MsgSender, workdir: PathBuf) -> Self {
493        let persistence_state = Arc::new(Mutex::new(PersistenceState::new(workdir.clone())));
494        Self {
495            tasks: crate::providers::TaskBroker::new(msg_tx.clone()),
496            msg_tx,
497            scopes: HashMap::new(),
498            cancelled_turns: VecDeque::new(),
499            detached: tokio::task::JoinSet::new(),
500            persistence_state,
501            persistence_tail: None,
502            workdir,
503            providers: None,
504            tools: None,
505            task_id: None,
506            terminal_title_enabled: true,
507            owns_global_mcp: true,
508            approval: None,
509            questions: None,
510            config_watch: None,
511        }
512    }
513
514    /// The channel every effect result arrives on.
515    ///
516    /// Cloning it is how the brokers already deliver a user's approval or
517    /// answer back into the reducer; `crate::engine::EngineHandle` uses it to
518    /// give the same reach to something outside the process's own effects.
519    #[must_use]
520    pub fn sender(&self) -> MsgSender {
521        self.msg_tx.clone()
522    }
523
524    /// Enable inline approval prompts (interactive TUI only). The gate then
525    /// pauses gated tools and routes the user's decision through the
526    /// `ApprovalBroker` instead of writing an out-of-band DB approval row.
527    #[must_use]
528    pub fn with_interactive_approvals(mut self) -> Self {
529        self.approval = Some(crate::providers::ApprovalBroker::new(self.msg_tx.clone()));
530        self
531    }
532
533    /// Enable inline `ask_user_question` prompts (interactive TUI only). The tool
534    /// then parks on the `QuestionBroker` and routes the user's answers back
535    /// through it instead of proceeding without asking.
536    #[must_use]
537    pub fn with_interactive_questions(mut self) -> Self {
538        self.questions = Some(crate::providers::QuestionBroker::new(self.msg_tx.clone()));
539        self
540    }
541
542    /// Start the background config watcher (#45): it polls `MERMAID.md` + memory
543    /// and emits `Msg::InstructionsChanged`/`MemoryChanged` on change, so the
544    /// reducer reads them as injected data instead of refreshing inline. Call
545    /// once at startup. Live-loop only — a replay driver feeds the recorded
546    /// Changed Msgs rather than polling.
547    pub fn spawn_config_watcher(&mut self, cwd: PathBuf, memory: MemoryConfig) {
548        let handle = self.detached.spawn(config_watch::config_watcher(
549            self.msg_tx.clone(),
550            cwd,
551            memory,
552        ));
553        self.config_watch = Some(handle);
554    }
555
556    /// Attach a durable runtime task id so tool runs, approvals,
557    /// checkpoints, compactions, and background processes can be linked.
558    #[must_use]
559    pub fn with_task_id(mut self, task_id: Option<String>) -> Self {
560        self.task_id = task_id;
561        self
562    }
563
564    /// Disable terminal-title writes for non-interactive callers.
565    #[must_use]
566    pub fn without_terminal_title(mut self) -> Self {
567        self.terminal_title_enabled = false;
568        self
569    }
570
571    /// Leave the process-global MCP manager alone on `shutdown`. Child
572    /// (subagent) runners share it with the parent and must not reap it.
573    #[must_use]
574    pub fn without_global_mcp_shutdown(mut self) -> Self {
575        self.owns_global_mcp = false;
576        self
577    }
578
579    /// Attach provider + tool registries. Production wiring uses
580    /// this; unit tests that don't need real dispatch can skip.
581    /// Without bindings, `CallModel` / `ExecuteTool` emit well-
582    /// formed error Msgs so the reducer still transitions cleanly.
583    pub fn with_bindings(
584        mut self,
585        providers: Arc<ProviderFactory>,
586        tools: Arc<ToolRegistry>,
587    ) -> Self {
588        self.providers = Some(providers);
589        self.tools = Some(tools);
590        self
591    }
592
593    /// Pair-constructor: returns both the runner and the receiving
594    /// end of the Msg channel. Preferred for production wiring
595    /// because it keeps the channel capacity constant in one place.
596    #[must_use]
597    pub fn pair(workdir: PathBuf) -> (Self, mpsc::Receiver<Msg>) {
598        let (tx, rx) = mpsc::channel(MSG_CHANNEL_CAPACITY);
599        (Self::new(tx, workdir), rx)
600    }
601
602    /// Pair constructor that also wires the real provider factory +
603    /// tool registry. Used by `app::run_interactive`.
604    #[must_use]
605    pub fn pair_with_bindings(
606        workdir: PathBuf,
607        config: Config,
608        tools: Arc<ToolRegistry>,
609    ) -> (Self, mpsc::Receiver<Msg>) {
610        let providers = Arc::new(ProviderFactory::new(config));
611        Self::pair_from(workdir, providers, tools)
612    }
613
614    /// Pair constructor that takes a pre-built `ProviderFactory`.
615    /// Used when the caller needs to share a `ProviderFactory` with
616    /// the `SubagentSpawner` so subagents can issue model calls
617    /// through the same cache.
618    pub fn pair_from(
619        workdir: PathBuf,
620        providers: Arc<ProviderFactory>,
621        tools: Arc<ToolRegistry>,
622    ) -> (Self, mpsc::Receiver<Msg>) {
623        let (tx, rx) = mpsc::channel(MSG_CHANNEL_CAPACITY);
624        (Self::new(tx, workdir).with_bindings(providers, tools), rx)
625    }
626
627    pub fn pair_from_with_task(
628        workdir: PathBuf,
629        providers: Arc<ProviderFactory>,
630        tools: Arc<ToolRegistry>,
631        task_id: Option<String>,
632    ) -> (Self, mpsc::Receiver<Msg>) {
633        let (runner, rx) = Self::pair_from(workdir, providers, tools);
634        (runner.with_task_id(task_id), rx)
635    }
636
637    /// Construct a runner that shares a pre-derived cancellation
638    /// token for its turn scopes. Used by `SubagentSpawner` so the
639    /// child runner's work aborts as soon as the parent's `ctx.token`
640    /// fires.
641    pub fn new_child(
642        msg_tx: MsgSender,
643        workdir: PathBuf,
644        providers: Arc<ProviderFactory>,
645        tools: Arc<ToolRegistry>,
646    ) -> Self {
647        // A subagent's runner is never the interactive top-level, so it must
648        // NOT emit OSC 2 terminal-title escapes: in a headless `mermaid run`
649        // the parent suppresses them, but an un-suppressed child leaks
650        // `\x1b]2;…\x07` into stdout and corrupts `--format json`/`text` output.
651        // It must also leave the process-global MCP manager running — the
652        // child shares the parent's servers, and reaping them here would kill
653        // MCP for the whole session the moment the first subagent finished.
654        Self::new(msg_tx, workdir)
655            .with_bindings(providers, tools)
656            .without_terminal_title()
657            .without_global_mcp_shutdown()
658    }
659
660    /// Get or create the scope for a turn. Idempotent. The scope is
661    /// retained until `CancelScope` tears it down or it naturally
662    /// drains.
663    fn scope_mut(&mut self, turn: TurnId) -> &mut TurnScope {
664        self.scopes
665            .entry(turn)
666            .or_insert_with(|| TurnScope::new(turn))
667    }
668
669    /// F38: record a cancelled turn in the bounded tombstone ring, evicting the
670    /// oldest id at capacity. Skips duplicates so a re-cancel doesn't churn the
671    /// ring (membership is all `is_tombstoned` checks).
672    fn tombstone_turn(&mut self, turn: TurnId) {
673        if self.cancelled_turns.contains(&turn) {
674            return;
675        }
676        if self.cancelled_turns.len() >= CANCELLED_TOMBSTONE_CAP {
677            self.cancelled_turns.pop_front();
678        }
679        self.cancelled_turns.push_back(turn);
680    }
681
682    /// F38: true iff `turn`'s scope was cancelled (tombstoned). New turn-scoped
683    /// work for such a turn is dropped rather than spinning up a fresh scope.
684    fn is_tombstoned(&self, turn: TurnId) -> bool {
685        self.cancelled_turns.contains(&turn)
686    }
687
688    /// Run one read-only `Cmd::Query` lookup and answer with
689    /// `Msg::QueryResult`. Conversation reads and provider discovery run
690    /// async; everything touching the runtime store or the filesystem walk
691    /// goes through [`Self::send_blocking_query`] so a synchronous read never
692    /// stalls an async worker thread (#40).
693    fn dispatch_query(&mut self, query: Query) {
694        let tx = self.msg_tx.clone();
695        match query {
696            Query::LoadConversation { id } => self.query_load_conversation(id, tx),
697            Query::ListConversations => self.query_list_conversations(tx),
698            Query::ListAvailableModels => {
699                let providers = self.providers.clone();
700                self.detached.spawn(async move {
701                    let choices = discover_available_models(providers).await;
702                    let _ = tx
703                        .send(Msg::QueryResult(QueryResult::AvailableModelsListed(
704                            choices,
705                        )))
706                        .await;
707                });
708            },
709            Query::ListProjectFiles => {
710                let workdir = self.workdir.clone();
711                self.send_blocking_query(move || {
712                    QueryResult::ProjectFilesListed(walk_project_files(&workdir))
713                });
714            },
715            Query::ListOutputStyles => self.dispatch_list_output_styles(),
716            Query::LoadOutputStyle { name, project } => {
717                self.dispatch_load_output_style(name, project);
718            },
719            Query::ListRuntimeTasks { limit } => self.send_blocking_query(move || {
720                QueryResult::RuntimeTasksListed(
721                    crate::runtime_client::RuntimeClient::auto()
722                        .list_tasks(limit)
723                        .map(|read| read.value)
724                        .unwrap_or_default(),
725                )
726            }),
727            Query::LoadRuntimeTask { id } => self.send_blocking_query(move || {
728                let (task, events) = crate::runtime_client::RuntimeClient::auto()
729                    .task_detail(&id)
730                    .map(|read| (Some(Box::new(read.value.task)), read.value.events))
731                    .unwrap_or((None, Vec::new()));
732                QueryResult::RuntimeTaskLoaded { task, events }
733            }),
734            Query::ListRuntimeProcesses { limit } => self.send_blocking_query(move || {
735                QueryResult::RuntimeProcessesListed(
736                    crate::runtime_client::RuntimeClient::auto()
737                        .list_processes(limit)
738                        .map(|read| read.value)
739                        .unwrap_or_default(),
740                )
741            }),
742            Query::ListRuntimeApprovals => self.send_blocking_query(move || {
743                QueryResult::RuntimeApprovalsListed(
744                    crate::runtime_client::RuntimeClient::auto()
745                        .list_approvals()
746                        .map(|read| read.value)
747                        .unwrap_or_default(),
748                )
749            }),
750            Query::ListRuntimeCheckpoints { limit } => self.send_blocking_query(move || {
751                QueryResult::RuntimeCheckpointsListed(
752                    crate::runtime_client::RuntimeClient::auto()
753                        .list_checkpoints(limit)
754                        .map(|read| read.value)
755                        .unwrap_or_default(),
756                )
757            }),
758            Query::ListForkCheckpoints {
759                session_id,
760                message_index,
761            } => self.send_blocking_query(move || {
762                QueryResult::ForkCheckpointsFound(
763                    mermaid_runtime::with_shared_store(|store| {
764                        store
765                            .checkpoints()
766                            .list_for_session(&session_id, message_index as i64)
767                    })
768                    .unwrap_or_default(),
769                )
770            }),
771            Query::ListRuntimePlugins => self.send_blocking_query(move || {
772                QueryResult::RuntimePluginsListed(
773                    crate::runtime_client::RuntimeClient::auto()
774                        .list_plugins()
775                        .map(|read| read.value)
776                        .unwrap_or_default(),
777                )
778            }),
779        }
780    }
781
782    /// `Query::ListOutputStyles` — every selectable output style
783    /// (built-ins plus user/project files).
784    fn dispatch_list_output_styles(&mut self) {
785        let workdir = self.workdir.clone();
786        self.send_blocking_query(move || {
787            QueryResult::OutputStylesListed(crate::app::output_styles::list_styles(&workdir))
788        });
789    }
790
791    /// `Query::LoadOutputStyle` — one style's body for the `/output-style`
792    /// switch. Answers `found: false` when nothing carries the name; the
793    /// reducer falls back to `default`.
794    fn dispatch_load_output_style(&mut self, name: String, project: bool) {
795        let workdir = self.workdir.clone();
796        self.send_blocking_query(move || {
797            let loaded = crate::app::output_styles::load_style_body(&workdir, &name);
798            let (found, body, keep, custom, source) = match loaded {
799                Some((body, keep, custom, source)) => {
800                    (true, body, keep, custom, source.to_string())
801                },
802                None => (false, String::new(), true, false, String::new()),
803            };
804            QueryResult::OutputStyleLoaded {
805                name,
806                project,
807                found,
808                body,
809                keep_coding_instructions: keep,
810                custom,
811                source,
812            }
813        });
814    }
815
816    /// `Query::LoadConversation` — read one saved conversation off disk. A
817    /// missing/corrupt file answers nothing: the failure is logged and the
818    /// picker simply does not advance.
819    fn query_load_conversation(&mut self, id: String, tx: MsgSender) {
820        let workdir = self.workdir.clone();
821        self.detached.spawn(async move {
822            match crate::session::ConversationManager::new(&workdir) {
823                Ok(mgr) => match mgr.load_conversation(&id) {
824                    Ok(history) => {
825                        let _ = tx
826                            .send(Msg::QueryResult(QueryResult::ConversationLoaded(Box::new(
827                                history,
828                            ))))
829                            .await;
830                    },
831                    Err(e) => {
832                        tracing::warn!(id = %id, error = %e, "LoadConversation failed");
833                    },
834                },
835                Err(e) => {
836                    tracing::warn!(error = %e, "ConversationManager init failed");
837                },
838            }
839        });
840    }
841
842    /// `Query::ListConversations` — scan the conversations directory for the
843    /// `/load` picker (newest first).
844    fn query_list_conversations(&mut self, tx: MsgSender) {
845        let workdir = self.workdir.clone();
846        self.detached.spawn(async move {
847            let summaries = match crate::session::ConversationManager::new(&workdir) {
848                Ok(mgr) => mgr
849                    .list_conversation_metas()
850                    .unwrap_or_default()
851                    .into_iter()
852                    .map(|m| mermaid_domain::ConversationSummary {
853                        id: m.id,
854                        title: m.title,
855                        message_count: m.message_count,
856                        updated_at: m.updated_at.to_rfc3339(),
857                    })
858                    .collect(),
859                Err(_) => Vec::new(),
860            };
861            let _ = tx
862                .send(Msg::QueryResult(QueryResult::ConversationsListed(
863                    summaries,
864                )))
865                .await;
866        });
867    }
868
869    /// Run a synchronous lookup on the blocking pool and deliver its
870    /// `Msg::QueryResult` — the shared plumbing of every store/filesystem
871    /// query (rusqlite reads and the project walk must never stall an async
872    /// worker thread, #40).
873    fn send_blocking_query(&mut self, run: impl FnOnce() -> QueryResult + Send + 'static) {
874        let tx = self.msg_tx.clone();
875        self.detached.spawn_blocking(move || {
876            let _ = tx.blocking_send(Msg::QueryResult(run()));
877        });
878    }
879
880    /// Drop the scope for a turn, signalling cancellation to every
881    /// child first. Safe to call for non-existent turns.
882    ///
883    /// After the scope is cancelled, a detached task moves it off the
884    /// runner, drains its `JoinSet` (so child tasks unwind), then emits
885    /// `Msg::TurnCancelled(turn)` so the reducer can transition
886    /// `Cancelling → Idle`. Without this terminal event the TUI would
887    /// stick in `Cancelling` — the reducer has no other way to learn
888    /// that the abort fully landed.
889    fn drop_scope(&mut self, turn: TurnId) {
890        // F38: tombstone this turn so a stray post-cancel turn-scoped Cmd can't
891        // resurrect an un-cancelled scope for it. Recorded for both the live and
892        // already-reaped branches below — once cancelled, a turn is dead either
893        // way (turn ids are monotonic and never reused).
894        self.tombstone_turn(turn);
895        if let Some(mut scope) = self.scopes.remove(&turn) {
896            scope.cancel();
897            let tx = self.msg_tx.clone();
898            self.detached.spawn(async move {
899                if tokio::time::timeout(CANCEL_DRAIN_TIMEOUT, scope.drain())
900                    .await
901                    .is_err()
902                {
903                    tracing::warn!(
904                        turn = %turn,
905                        timeout_ms = CANCEL_DRAIN_TIMEOUT.as_millis(),
906                        "cancel drain timed out; aborting remaining scoped tasks"
907                    );
908                }
909                let _ = tx.send(Msg::TurnCancelled(turn)).await;
910            });
911        } else {
912            // The scope was already reaped — its `JoinSet` drained to empty
913            // and `reap_empty_scopes` (top of `dispatch`) removed it before
914            // this cancel landed. The reducer is still in `Cancelling` with
915            // no other way to learn the turn ended, so emit the terminal
916            // event anyway. Idempotent: `handle_turn_cancelled` no-ops on
917            // any turn that isn't currently `Cancelling`.
918            let tx = self.msg_tx.clone();
919            self.detached.spawn(async move {
920                let _ = tx.send(Msg::TurnCancelled(turn)).await;
921            });
922        }
923    }
924
925    /// Number of active per-turn scopes. Tests use this to observe
926    /// lifecycle without racing on internal state.
927    #[must_use]
928    pub fn scope_count(&self) -> usize {
929        self.scopes.len()
930    }
931
932    /// F12: remove scope entries whose `JoinSet` is empty — every
933    /// child task has completed, so the scope is just an orphan key
934    /// in the map. Called at the top of `dispatch` so the map stays
935    /// bounded over long sessions. Cheap: one linear walk, no async.
936    ///
937    /// `JoinSet::is_empty` only returns true after completed tasks are
938    /// harvested via `join_next`/`try_join_next`, so we first drain
939    /// any ready completions per scope.
940    fn reap_empty_scopes(&mut self) {
941        self.reap_detached();
942        self.scopes.retain(|_, scope| {
943            scope.drain_completed();
944            !scope.is_empty()
945        });
946    }
947
948    /// Harvest finished detached tasks. Without this the `detached` `JoinSet`
949    /// grows for the whole session (every fire-and-forget effect lingers as a
950    /// completed-but-unjoined handle), and a panicking detached task vanishes
951    /// without a trace. Non-blocking — only already-finished tasks are taken (#38).
952    fn reap_detached(&mut self) {
953        while let Some(result) = self.detached.try_join_next() {
954            if let Err(e) = result
955                && !e.is_cancelled()
956            {
957                tracing::warn!(error = %e, "effect: detached task panicked");
958            }
959        }
960    }
961
962    /// Route a single `Cmd` into the appropriate spawn + handler.
963    /// Returns immediately; handlers work asynchronously and emit
964    /// `Msg` back through the sender channel.
965    #[expect(
966        clippy::too_many_lines,
967        reason = "the effect router: one arm per Cmd variant, each spawning or calling the \
968         handler that owns that effect; the arms are short and the routing table is the point, so \
969         the length is the Cmd count and drops only as commands are retired"
970    )]
971    pub fn dispatch(&mut self, cmd: Cmd) {
972        // F12: reap any drained scopes before touching the map. Keeps
973        // `scope_count()` bounded as the session grows.
974        self.reap_empty_scopes();
975        tracing::trace!(cmd = %cmd.summary(), "effect: dispatch");
976
977        // F38: refuse to spawn fresh work for a turn we've already cancelled.
978        // Only the scope-spawning variants carry a `scope_turn()`; `CancelScope`
979        // returns `None` here so a re-cancel still reaches `drop_scope` (which
980        // re-emits the terminal `TurnCancelled` the reducer needs). Turn ids are
981        // monotonic and never reused, so a tombstoned id can only be a stray
982        // post-cancel straggler — dropping it stops `scope_mut`'s `or_insert_with`
983        // from resurrecting an un-cancelled scope.
984        if let Some(turn) = cmd.scope_turn()
985            && self.is_tombstoned(turn)
986        {
987            tracing::debug!(
988                cmd = %cmd.summary(),
989                turn = %turn,
990                "effect: dropping turn-scoped cmd for an already-cancelled turn"
991            );
992            return;
993        }
994
995        match cmd {
996            Cmd::CallModel { turn, mut request } => {
997                let tx = self.msg_tx.clone();
998                let providers = self.providers.clone();
999                // Enrich `request.tools` with every user-facing
1000                // tool in the bound registry. The reducer has
1001                // already populated MCP tools from `state.mcp`;
1002                // built-ins come from the runner (which holds the
1003                // registry). This keeps `ChatRequest.tools` the
1004                // single source of truth for what the model sees.
1005                // Formatting turns (`output_schema`) advertise NO tools —
1006                // the reducer already sent none; don't re-add built-ins.
1007                if let Some(tools) = &self.tools
1008                    && request.output_schema.is_none()
1009                {
1010                    let mut enriched =
1011                        filter_suppressed(tools.describe_all(), &request.suppressed_builtin_tools);
1012                    // Report the built-in tool-schema token cost so the
1013                    // reducer's /context preview can fold it into its MCP-only
1014                    // estimate and agree with what the model actually sees.
1015                    // Runs AFTER suppression so the estimate matches reality.
1016                    let builtin_tokens = mermaid_domain::estimate_tool_schema_tokens(&enriched);
1017                    // Best-effort and cosmetic (the /context preview). This is the
1018                    // synchronous dispatch path so we can't await; if the bounded
1019                    // channel is momentarily full under heavy streaming, log the
1020                    // drop rather than swallowing it silently — the estimate just
1021                    // stays briefly stale (#F43).
1022                    if let Err(e) = tx.try_send(Msg::BuiltinToolSchemaTokens(builtin_tokens)) {
1023                        tracing::debug!(
1024                            error = %e,
1025                            "effect: dropped builtin tool-schema token estimate (channel full); \
1026                             /context preview may be briefly stale"
1027                        );
1028                    }
1029                    enriched.append(&mut request.tools);
1030                    request.tools = enriched;
1031                }
1032                // Detached + off the blocking pool: never run a plugin hook on
1033                // the synchronous dispatch path (it would freeze input/render).
1034                self.detached.spawn(fire_plugin_hooks(
1035                    "prompt_submit",
1036                    serde_json::json!({
1037                        "turn_id": turn.0,
1038                        "model_id": request.model_id.clone(),
1039                        "message_count": request.messages.len(),
1040                        "tool_count": request.tools.len(),
1041                    }),
1042                ));
1043                // Task cost attribution: model dispatch reports each request's
1044                // completion tokens into the broker's cumulative counter.
1045                let task_usage = self.tasks.clone();
1046                let scope = self.scope_mut(turn);
1047                let token = scope.token();
1048                scope.spawn(async move {
1049                    use futures::FutureExt;
1050                    let fallback_tx = tx.clone();
1051                    if std::panic::AssertUnwindSafe(dispatch_call_model(
1052                        tx, providers, turn, request, token, task_usage,
1053                    ))
1054                    .catch_unwind()
1055                    .await
1056                    .is_err()
1057                    {
1058                        // The dispatch task panicked. A turn whose model call
1059                        // never emits a terminal Msg stays in `Generating`
1060                        // forever; emit one so the reducer can leave that state
1061                        // instead of wedging (#43).
1062                        tracing::error!(turn = %turn, "dispatch_call_model panicked");
1063                        let _ = fallback_tx
1064                            .send(Msg::UpstreamError {
1065                                turn,
1066                                error: mermaid_model::models::UserFacingError {
1067                                    summary: "Internal error".to_string(),
1068                                    message: "The model dispatch task panicked unexpectedly."
1069                                        .to_string(),
1070                                    suggestion: "This is a bug. Please retry; if it persists, \
1071                                                 check the logs."
1072                                        .to_string(),
1073                                    category: mermaid_model::models::ErrorCategory::Internal,
1074                                    recoverable: true,
1075                                },
1076                            })
1077                            .await;
1078                    }
1079                });
1080            },
1081            Cmd::CompactConversation { turn, mut request } => {
1082                let tx = self.msg_tx.clone();
1083                let providers = self.providers.clone();
1084                if let Some(tools) = &self.tools {
1085                    let mut enriched = tools.describe_all();
1086                    enriched.append(&mut request.chat.tools);
1087                    request.chat.tools = enriched;
1088                }
1089                // Capture the trigger before `request` moves into the task, so a
1090                // panic fallback can still name which compaction failed.
1091                let trigger = request.trigger;
1092                let scope = self.scope_mut(turn);
1093                let token = scope.token();
1094                scope.spawn(async move {
1095                    use futures::FutureExt;
1096                    let fallback_tx = tx.clone();
1097                    if std::panic::AssertUnwindSafe(dispatch_compact_conversation(
1098                        tx, providers, turn, request, token,
1099                    ))
1100                    .catch_unwind()
1101                    .await
1102                    .is_err()
1103                    {
1104                        // The compaction task panicked. Without a terminal
1105                        // `CompactionFinished`/`CompactionFailed`, the reducer
1106                        // wedges in `Compacting` until Ctrl+C; emit a failure so
1107                        // it can recover, mirroring `CallModel`/`ExecuteTool`
1108                        // (#43, F37).
1109                        tracing::error!(turn = %turn, "dispatch_compact_conversation panicked");
1110                        let _ = fallback_tx
1111                            .send(Msg::CompactionFailed {
1112                                turn,
1113                                trigger,
1114                                message: "the compaction task panicked unexpectedly".to_string(),
1115                                kind: mermaid_domain::StatusKind::Error,
1116                            })
1117                            .await;
1118                    }
1119                });
1120            },
1121            Cmd::ExecuteTool {
1122                turn,
1123                call_id,
1124                source,
1125                dispatch,
1126            } => {
1127                let tx = self.msg_tx.clone();
1128                let tools = self.tools.clone();
1129                let workdir = self.workdir.clone();
1130                // Pass the shared Config from ProviderFactory so
1131                // subagents inherit it (F7). Falls back to
1132                // Config::default() when providers aren't bound (unit
1133                // tests without real wiring).
1134                let config = self
1135                    .providers
1136                    .as_ref()
1137                    .map(|p| Arc::new(p.config().clone()))
1138                    .unwrap_or_else(|| Arc::new(mermaid_domain::Config::default()));
1139                // Auto mode: build an LLM classifier to vet borderline
1140                // actions. Only when a provider is bound (real wiring); the
1141                // gate fails safe to "escalate" when it's `None`. The vet
1142                // uses the configured classifier model, else the session model.
1143                // Plan mode also gets one: profile levels set to `auto`
1144                // resolve through `PolicyDecision::Classify`, which fails
1145                // safe to escalate without a classifier bound.
1146                let classifier: Option<Arc<dyn crate::providers::AutoClassifier>> =
1147                    if dispatch.safety_mode == mermaid_runtime::SafetyMode::Auto
1148                        || dispatch.plan_file.is_some()
1149                    {
1150                        self.providers.as_ref().map(|p| {
1151                            let model = config
1152                                .safety
1153                                .auto_classifier_model
1154                                .clone()
1155                                .unwrap_or_else(|| dispatch.model_id.clone());
1156                            Arc::new(crate::providers::ModelAutoClassifier::new(p.clone(), model))
1157                                as Arc<dyn crate::providers::AutoClassifier>
1158                        })
1159                    } else {
1160                        None
1161                    };
1162                let services = crate::providers::ctx::ToolServices {
1163                    workdir,
1164                    config,
1165                    task_id: self.task_id.clone(),
1166                    // Detached work (backgrounded subagents) reports back
1167                    // through the main msg channel after this turn's
1168                    // progress relay is gone.
1169                    notify: Some(self.msg_tx.clone()),
1170                    classifier,
1171                    approval: self.approval.clone(),
1172                    questions: self.questions.clone(),
1173                    tasks: Some(self.tasks.clone()),
1174                };
1175                let scope = self.scope_mut(turn);
1176                let signals = crate::providers::ctx::TurnSignals {
1177                    token: scope.token(),
1178                    background: scope.background_token(),
1179                    web_bytes: scope.web_bytes(),
1180                };
1181                scope.spawn(async move {
1182                    use futures::FutureExt;
1183                    let fallback_tx = tx.clone();
1184                    if std::panic::AssertUnwindSafe(dispatch_execute_tool(
1185                        tx, tools, turn, call_id, source, signals, dispatch, services,
1186                    ))
1187                    .catch_unwind()
1188                    .await
1189                    .is_err()
1190                    {
1191                        // The tool task panicked. Its turn waits on a
1192                        // `ToolFinished` for this `call_id` that will now never
1193                        // arrive; emit a terminal error outcome so the turn
1194                        // doesn't wedge (#43).
1195                        tracing::error!(
1196                            turn = %turn,
1197                            call_id = call_id.0,
1198                            "dispatch_execute_tool panicked"
1199                        );
1200                        let _ = fallback_tx
1201                            .send(Msg::ToolFinished {
1202                                turn,
1203                                call_id,
1204                                outcome: mermaid_domain::ToolOutcome::error(
1205                                    "internal error: the tool execution task panicked".to_string(),
1206                                    0.0,
1207                                ),
1208                            })
1209                            .await;
1210                    }
1211                });
1212            },
1213            Cmd::ResolveApproval { call_id, decision } => {
1214                // Deliver the user's inline decision to the parked tool task.
1215                // Not turn-scoped — fire-and-forget to the broker.
1216                if let Some(broker) = &self.approval {
1217                    broker.resolve(call_id, decision.into());
1218                }
1219            },
1220            Cmd::ResolveQuestion {
1221                call_id,
1222                resolution,
1223            } => {
1224                // Deliver the user's answers to the parked ask_user_question
1225                // task. Not turn-scoped — fire-and-forget to the broker.
1226                if let Some(broker) = &self.questions {
1227                    broker.resolve(call_id, resolution);
1228                }
1229            },
1230            Cmd::SyncTaskStore(store) => {
1231                // Reducer-initiated truth overwrite (rewind/fork, /clear,
1232                // startup resume). Synchronous; the broker does not publish
1233                // back — the reducer already holds this store.
1234                self.tasks.seed(store);
1235            },
1236            Cmd::EnsureScratchpad { session_id } => {
1237                let tx = self.msg_tx.clone();
1238                let workdir = self.workdir.clone();
1239                self.detached.spawn(async move {
1240                    match crate::session::scratchpad::ensure(&workdir, &session_id) {
1241                        Ok(path) => {
1242                            let _ = tx.send(Msg::ScratchpadReady { session_id, path }).await;
1243                        },
1244                        Err(err) => {
1245                            // Non-fatal: the session runs without a scratch
1246                            // dir (`Session::scratchpad` stays `None`).
1247                            tracing::warn!(error = %err, "failed to create session scratchpad");
1248                        },
1249                    }
1250                    // Best-effort reap of unlocked scratchpads past retention —
1251                    // piggybacks on session startup, no separate timer.
1252                    if let Err(err) = crate::session::scratchpad::sweep_stale(
1253                        crate::session::scratchpad::RETENTION_DAYS,
1254                    ) {
1255                        tracing::warn!(error = %err, "scratchpad sweep failed");
1256                    }
1257                });
1258            },
1259            Cmd::ListScratchpad { path } => {
1260                // `/scratchpad` — bounded directory listing back into the
1261                // transcript. Blocking filesystem walk, so off the runner.
1262                let tx = self.msg_tx.clone();
1263                self.detached.spawn(async move {
1264                    let text = tokio::task::spawn_blocking(move || {
1265                        crate::session::scratchpad::list_text(&path)
1266                    })
1267                    .await
1268                    .unwrap_or_else(|e| format!("Couldn't list the scratchpad: {e}"));
1269                    let _ = tx.send(Msg::RuntimeText(text)).await;
1270                });
1271            },
1272            Cmd::UserTaskEdit(edit) => {
1273                // Route the user's /tasks edit through the broker (single
1274                // writer) so it serializes with any in-flight tool call. The
1275                // broker publishes the resulting snapshot; the outcome line
1276                // lands in the transcript as transient status.
1277                let broker = self.tasks.clone();
1278                let tx = self.msg_tx.clone();
1279                self.detached.spawn(async move {
1280                    let (line, _snapshot) = broker.user_edit(edit).await;
1281                    // The user sees the ack in the transcript; the model
1282                    // learns about it on its next request via the notice
1283                    // buffer (a checklist the model believes in but the user
1284                    // has edited is the worst of both).
1285                    let _ = tx
1286                        .send(Msg::TaskNotice {
1287                            text: format!(
1288                                "The user edited the task checklist: {line}. Acknowledge and \
1289                                 incorporate this into your plan."
1290                            ),
1291                        })
1292                        .await;
1293                    let _ = tx.send(Msg::TransientStatus { text: line }).await;
1294                });
1295            },
1296            Cmd::NotifyTaskCompleted {
1297                task,
1298                completed,
1299                total,
1300            } => {
1301                // Gated `task_completed` plugin hook: a denying hook VETOES
1302                // the completion — the task flips back to in_progress via the
1303                // broker (single writer; the publish refreshes the band) and
1304                // the reason reaches both the user (transcript) and the model
1305                // (notice buffer). Fail-open like every plugin hook: no
1306                // enabled hooks / timeout => allow, zero latency added
1307                // elsewhere because this runs detached.
1308                let payload = serde_json::json!({
1309                    "task_id": task.id,
1310                    "subject": task.subject,
1311                    "description": task.description,
1312                    "evidence": task.evidence,
1313                    "completed": completed,
1314                    "total": total,
1315                });
1316                let broker = self.tasks.clone();
1317                let tx = self.msg_tx.clone();
1318                self.detached.spawn(async move {
1319                    let gate = run_plugin_hooks_gated("task_completed", payload).await;
1320                    let Some((plugin, reason)) = gate.deny else {
1321                        return;
1322                    };
1323                    let reason = mermaid_model::utils::redact_secrets(&reason);
1324                    let _ = broker
1325                        .update(vec![mermaid_domain::ChecklistEdit {
1326                            id: task.id,
1327                            status: Some(mermaid_domain::ChecklistStatus::InProgress),
1328                            ..mermaid_domain::ChecklistEdit::default()
1329                        }])
1330                        .await;
1331                    let _ = tx
1332                        .send(Msg::TaskNotice {
1333                            text: format!(
1334                                "Completion of task #{} '{}' was vetoed by the {plugin} hook: \
1335                                 {reason}. The task is back in_progress; address the reason \
1336                                 before completing it again.",
1337                                task.id, task.subject
1338                            ),
1339                        })
1340                        .await;
1341                    let _ = tx
1342                        .send(Msg::TransientStatus {
1343                            text: format!(
1344                                "task #{} completion vetoed by {plugin}: {reason}",
1345                                task.id
1346                            ),
1347                        })
1348                        .await;
1349                });
1350            },
1351            Cmd::CancelScope(turn) => {
1352                self.drop_scope(turn);
1353            },
1354            Cmd::BackgroundScope(turn) => {
1355                // Fire the scope's background token (don't drop the scope):
1356                // detachable tools move their child to a background process and
1357                // return a normal outcome, so the turn finishes naturally.
1358                self.scope_mut(turn).background();
1359            },
1360            Cmd::SaveConversation { snapshot, events } => {
1361                self.queue_persistence(PersistenceJob::Conversation {
1362                    snapshot: Box::new(snapshot),
1363                    events,
1364                });
1365            },
1366            Cmd::SaveCompaction {
1367                record,
1368                conversation,
1369                events,
1370            } => {
1371                self.queue_persistence(PersistenceJob::Compaction(Box::new(
1372                    PendingCompactionSave {
1373                        record,
1374                        conversation,
1375                        events,
1376                        events_appended: false,
1377                        task_id: self.task_id.clone(),
1378                    },
1379                )));
1380            },
1381            Cmd::SaveProcess(process) => {
1382                let task_id = self.task_id.clone();
1383                self.detached.spawn(async move {
1384                    let status = process.status;
1385                    let _ = mermaid_runtime::with_shared_store(|store| {
1386                        store.processes().upsert(mermaid_runtime::NewProcess {
1387                            id: Some(process.id),
1388                            task_id,
1389                            pid: process.pid,
1390                            command: process.command,
1391                            cwd: process.cwd,
1392                            log_path: Some(process.log_path),
1393                            detected_url: process.detected_url,
1394                            status,
1395                            health: None,
1396                        })
1397                    });
1398                });
1399            },
1400            Cmd::PersistPlanConfig(plan) => {
1401                self.detached.spawn(async move {
1402                    if let Err(err) = crate::app::persist_plan_config(&plan) {
1403                        tracing::warn!(error = %err, "failed to persist [plan] config");
1404                    }
1405                });
1406            },
1407            Cmd::PersistLastModel(model) => {
1408                self.detached.spawn(async move {
1409                    if let Err(err) = crate::app::persist_last_model(&model) {
1410                        tracing::warn!(error = %err, "failed to persist last-used model");
1411                    }
1412                });
1413            },
1414            Cmd::PersistReasoningFor { model_id, level } => {
1415                self.detached.spawn(async move {
1416                    if let Err(err) = crate::app::persist_reasoning_for_model(&model_id, level) {
1417                        tracing::warn!(error = %err, "failed to persist reasoning level for model");
1418                    }
1419                });
1420            },
1421            Cmd::PersistOllamaNumCtxFor { model_id, num_ctx } => {
1422                self.detached.spawn(async move {
1423                    if let Err(err) =
1424                        crate::app::persist_ollama_num_ctx_for_model(&model_id, num_ctx)
1425                    {
1426                        tracing::warn!(error = %err, "failed to persist Ollama num_ctx for model");
1427                    }
1428                });
1429            },
1430            Cmd::PersistOllamaOffload(enabled) => {
1431                self.detached.spawn(async move {
1432                    if let Err(err) = crate::app::persist_ollama_allow_ram_offload(enabled) {
1433                        tracing::warn!(error = %err, "failed to persist Ollama RAM-offload setting");
1434                    }
1435                });
1436            },
1437            Cmd::PersistUiTheme(theme) => {
1438                self.detached.spawn(async move {
1439                    if let Err(err) = crate::app::persist_ui_theme(theme) {
1440                        tracing::warn!(error = %err, "failed to persist theme");
1441                    }
1442                });
1443            },
1444            Cmd::PersistOutputStyle { style } => {
1445                self.detached.spawn(async move {
1446                    if let Err(err) = crate::app::persist_output_style(&style) {
1447                        tracing::warn!(error = %err, "failed to persist output style");
1448                    }
1449                });
1450            },
1451            Cmd::PersistProjectOutputStyle { style } => {
1452                let tx = self.msg_tx.clone();
1453                let workdir = self.workdir.clone();
1454                self.detached.spawn(async move {
1455                    if let Err(err) = crate::app::persist_project_output_style(&workdir, &style) {
1456                        let _ = tx
1457                            .send(Msg::TransientStatus {
1458                                text: format!("Couldn't save the project output style: {err}"),
1459                            })
1460                            .await;
1461                    }
1462                });
1463            },
1464            Cmd::ComposeInEditor { .. } => {
1465                // Run-loop-intercepted in the interactive TUI (it owns the
1466                // terminal + event stream). Reaching the effect runner means a
1467                // headless driver emitted it — nothing to suspend there.
1468                tracing::warn!("compose_in_editor is unavailable outside the interactive TUI");
1469            },
1470            Cmd::ListMemory => {
1471                let tx = self.msg_tx.clone();
1472                let workdir = self.workdir.clone();
1473                self.detached.spawn(async move {
1474                    let cfg = crate::app::load_project_scoped_config(&workdir).memory;
1475                    let text = match crate::app::memory::load(&workdir, &cfg) {
1476                        Some(mem) => mem.index,
1477                        None => "No memories saved yet. Durable facts (yours or mine) show up here — use `/remember <fact>` or just ask me to remember something.".to_string(),
1478                    };
1479                    let _ = tx.send(Msg::RuntimeText(text)).await;
1480                });
1481            },
1482            Cmd::RememberMemory { text } => {
1483                let tx = self.msg_tx.clone();
1484                let workdir = self.workdir.clone();
1485                self.detached.spawn(async move {
1486                    let cfg = crate::app::load_project_scoped_config(&workdir).memory;
1487                    let name = memory_title_from_text(&text);
1488                    let status = match crate::app::memory::write_memory(
1489                        &workdir,
1490                        mermaid_domain::MemoryScope::ProjectPrivate,
1491                        &name,
1492                        &text,
1493                        &[],
1494                        &text,
1495                    ) {
1496                        Ok(_) => format!("Remembered: {name}"),
1497                        Err(e) => format!("Couldn't save memory: {e}"),
1498                    };
1499                    let (loaded, _) = crate::app::memory::refresh(None, &workdir, &cfg);
1500                    let _ = tx.send(Msg::MemoryChanged(loaded)).await;
1501                    let _ = tx.send(Msg::TransientStatus { text: status }).await;
1502                });
1503            },
1504            Cmd::ForgetMemory { id } => {
1505                let tx = self.msg_tx.clone();
1506                let workdir = self.workdir.clone();
1507                self.detached.spawn(async move {
1508                    let cfg = crate::app::load_project_scoped_config(&workdir).memory;
1509                    let status = match crate::app::memory::delete_memory(&workdir, &id) {
1510                        Ok(Some(_)) => format!("Forgot: {id}"),
1511                        Ok(None) => format!("No memory named '{id}'"),
1512                        Err(e) => format!("Couldn't forget memory: {e}"),
1513                    };
1514                    let (loaded, _) = crate::app::memory::refresh(None, &workdir, &cfg);
1515                    let _ = tx.send(Msg::MemoryChanged(loaded)).await;
1516                    let _ = tx.send(Msg::TransientStatus { text: status }).await;
1517                });
1518            },
1519            Cmd::ConsolidateMemory { model_id } => {
1520                let tx = self.msg_tx.clone();
1521                let workdir = self.workdir.clone();
1522                let providers = self.providers.clone();
1523                self.detached.spawn(async move {
1524                    consolidate_memory(tx, providers, workdir, model_id).await;
1525                });
1526            },
1527            Cmd::Query(query) => self.dispatch_query(query),
1528            Cmd::ShowRuntimeProcessLogs { id } => {
1529                let tx = self.msg_tx.clone();
1530                self.detached.spawn_blocking(move || {
1531                    let text = crate::runtime_client::RuntimeClient::auto()
1532                        .process_log(&id, None)
1533                        .map(|log| format!("Process log {}\n\n{}", id, log.content))
1534                        .unwrap_or_else(|err| format!("Process log error: {err}"));
1535                    let _ = tx.blocking_send(Msg::RuntimeText(text));
1536                });
1537            },
1538            Cmd::StopRuntimeProcess { id } => {
1539                let tx = self.msg_tx.clone();
1540                self.detached.spawn_blocking(move || {
1541                    let msg = match crate::runtime_client::RuntimeClient::auto().stop_process(&id) {
1542                        Ok(response) => Msg::TransientStatus {
1543                            text: format!("Stopped process {} (pid {})", id, response.item.pid),
1544                        },
1545                        Err(err) => Msg::TransientStatus {
1546                            text: format!("Process stop failed: {err}"),
1547                        },
1548                    };
1549                    let _ = tx.blocking_send(msg);
1550                });
1551            },
1552            Cmd::KillBackgroundAgent { agent_id } => {
1553                // Synchronous token fire — no task to spawn. Feedback flows
1554                // through the dying child's `Msg::BackgroundAgentFinished`
1555                // (the reducer already validated the id against its registry).
1556                let spawner = self.tools.as_ref().and_then(|t| t.subagent_spawner());
1557                if let Some(spawner) = spawner {
1558                    match agent_id {
1559                        Some(id) => {
1560                            // Killing a finished child hands back the
1561                            // workspace it was holding for a continuation
1562                            // that can no longer happen. Discarding it needs
1563                            // to be async, so it rides a task.
1564                            if let crate::providers::tool::subagent::KillResult::Evicted(
1565                                workspace,
1566                            ) = spawner.kill_detached(&id)
1567                            {
1568                                // In `detached`, not a bare `tokio::spawn`: shutdown
1569                                // drains that set, so quitting mid-discard cannot
1570                                // leave the worktree on disk with no record of it.
1571                                self.detached.spawn(async move {
1572                                    workspace.discard().await;
1573                                });
1574                            }
1575                        },
1576                        None => {
1577                            spawner.kill_all_detached();
1578                        },
1579                    }
1580                }
1581            },
1582            Cmd::RestartRuntimeProcess { id } => {
1583                let tx = self.msg_tx.clone();
1584                self.detached.spawn_blocking(move || {
1585                    let msg = match crate::runtime_client::RuntimeClient::auto()
1586                        .restart_process(&id)
1587                    {
1588                        Ok(response) => Msg::TransientStatus {
1589                            text: format!("Restarted process {} (pid {})", id, response.item.pid),
1590                        },
1591                        Err(err) => Msg::TransientStatus {
1592                            text: format!("Process restart failed: {err}"),
1593                        },
1594                    };
1595                    let _ = tx.blocking_send(msg);
1596                });
1597            },
1598            Cmd::OpenRuntimeTarget { target } => {
1599                self.detached.spawn_blocking(move || {
1600                    let resolved = crate::runtime_client::RuntimeService::open_default()
1601                        .and_then(|service| service.resolve_open_target(&target))
1602                        .unwrap_or(target);
1603                    // #63: the resolved value can be a `detected_url`/`log_path`
1604                    // from a `processes` row — validate before the OS opener,
1605                    // exactly like `open_process`.
1606                    if let Err(err) = crate::runtime_client::validate_open_target(&resolved) {
1607                        tracing::warn!(error = %err, "refusing to open runtime target");
1608                        return;
1609                    }
1610                    mermaid_model::utils::open_file(resolved);
1611                });
1612            },
1613            Cmd::ShowRuntimePorts => {
1614                let tx = self.msg_tx.clone();
1615                self.detached.spawn_blocking(move || {
1616                    let text = crate::runtime_client::RuntimeClient::auto()
1617                        .ports()
1618                        .map(|ports| format!("Listening TCP ports\n\n{}", ports.ports))
1619                        .unwrap_or_else(|err| format!("Port inspection failed: {err}"));
1620                    let _ = tx.blocking_send(Msg::RuntimeText(text));
1621                });
1622            },
1623            Cmd::DecideRuntimeApproval { id, decision } => {
1624                let tx = self.msg_tx.clone();
1625                self.detached.spawn_blocking(move || {
1626                    let result = if decision == "approved" {
1627                        crate::runtime_client::RuntimeClient::auto().approve(&id)
1628                    } else {
1629                        crate::runtime_client::RuntimeClient::auto().deny(&id)
1630                    };
1631                    let msg = match result {
1632                        Ok(result) => Msg::TransientStatus {
1633                            text: if result.replayed {
1634                                format!("Approval {} {}: {}", id, decision, result.summary)
1635                            } else {
1636                                format!("Approval {id} {decision}")
1637                            },
1638                        },
1639                        Err(err) => Msg::TransientStatus {
1640                            text: format!("Approval update failed: {err}"),
1641                        },
1642                    };
1643                    let _ = tx.blocking_send(msg);
1644                });
1645            },
1646            Cmd::UpdateRuntimeTaskStatus {
1647                id,
1648                status,
1649                final_report,
1650            } => {
1651                let tx = self.msg_tx.clone();
1652                self.detached.spawn_blocking(move || {
1653                    let msg = match mermaid_runtime::with_shared_store(|store| {
1654                        store
1655                            .tasks()
1656                            .update_status(&id, status, final_report.as_deref())
1657                    }) {
1658                        Ok(()) => Msg::TransientStatus {
1659                            text: format!("Task {id} -> {status}"),
1660                        },
1661                        Err(err) => Msg::TransientStatus {
1662                            text: format!("Task update failed: {err}"),
1663                        },
1664                    };
1665                    let _ = tx.blocking_send(msg);
1666                });
1667            },
1668            Cmd::CreateRuntimeCheckpoint { paths } => {
1669                let tx = self.msg_tx.clone();
1670                let workdir = self.workdir.clone();
1671                self.detached.spawn_blocking(move || {
1672                    let pending_action = Some(serde_json::json!({
1673                        "source": "tui",
1674                        "command": "checkpoint",
1675                    }));
1676                    let msg = match mermaid_runtime::create_checkpoint(
1677                        &workdir,
1678                        &paths,
1679                        pending_action,
1680                    ) {
1681                        Ok(manifest) => Msg::TransientStatus {
1682                            text: format!(
1683                                "Checkpoint {} created for {} path(s)",
1684                                manifest.id,
1685                                manifest.files.len()
1686                            ),
1687                        },
1688                        Err(err) => Msg::TransientStatus {
1689                            text: format!("Checkpoint failed: {err}"),
1690                        },
1691                    };
1692                    let _ = tx.blocking_send(msg);
1693                });
1694            },
1695            Cmd::RestoreRuntimeCheckpoint { id } => {
1696                let tx = self.msg_tx.clone();
1697                self.detached.spawn_blocking(move || {
1698                    let msg = match crate::runtime_client::RuntimeClient::auto()
1699                        .restore_checkpoint(&id)
1700                    {
1701                        Ok(result) => Msg::TransientStatus {
1702                            text: format!(
1703                                "Restored checkpoint {} ({} file(s)){}",
1704                                result.checkpoint.id,
1705                                result.checkpoint.files.len(),
1706                                if result.checkpoint.pending_action.is_some() {
1707                                    "; pending action available in checkpoint manifest"
1708                                } else {
1709                                    ""
1710                                }
1711                            ),
1712                        },
1713                        Err(err) => Msg::TransientStatus {
1714                            text: format!("Restore failed: {err}"),
1715                        },
1716                    };
1717                    let _ = tx.blocking_send(msg);
1718                });
1719            },
1720            Cmd::ShowRuntimeModelInfo { model } => {
1721                let tx = self.msg_tx.clone();
1722                self.detached.spawn_blocking(move || {
1723                    let text = runtime_model_info_text(&model);
1724                    let _ = tx.blocking_send(Msg::RuntimeText(text));
1725                });
1726            },
1727            Cmd::InitMcpServers(configs) => {
1728                let tx = self.msg_tx.clone();
1729                self.detached
1730                    .spawn(async move { dispatch_init_mcp_servers(configs, tx).await });
1731            },
1732            Cmd::StopMcpServer { name } => {
1733                let tx = self.msg_tx.clone();
1734                self.detached.spawn(async move {
1735                    // Actually kill the child before claiming it's stopped —
1736                    // otherwise the UI says "stopped" while the server runs on.
1737                    if let Some(mgr) = crate::mcp::manager_ref::get() {
1738                        mgr.stop_server(&name).await;
1739                    }
1740                    let _ = tx.send(Msg::McpServerStopped { name }).await;
1741                });
1742            },
1743            Cmd::PullOllamaModel { model } => {
1744                let tx = self.msg_tx.clone();
1745                self.detached.spawn(async move {
1746                    dispatch_pull_ollama_model(tx, model).await;
1747                });
1748            },
1749            Cmd::OpenInSystem(path) => {
1750                self.detached.spawn(async move {
1751                    let _ = tokio::task::spawn_blocking(move || {
1752                        mermaid_model::utils::open_file(&path);
1753                    })
1754                    .await;
1755                });
1756            },
1757            Cmd::WriteImageToTemp {
1758                path,
1759                bytes,
1760                format: _,
1761            } => {
1762                self.detached.spawn(async move {
1763                    if let Err(e) = tokio::fs::write(&path, &bytes).await {
1764                        tracing::warn!(path = %path.display(), error = %e, "WriteImageToTemp failed");
1765                    }
1766                });
1767            },
1768            Cmd::ReadClipboard => {
1769                let tx = self.msg_tx.clone();
1770                self.detached.spawn(async move {
1771                    dispatch_read_clipboard(tx).await;
1772                });
1773            },
1774            Cmd::ProbeVision { model_id, warn } => {
1775                let tx = self.msg_tx.clone();
1776                let providers = self.providers.clone();
1777                self.detached.spawn(async move {
1778                    dispatch_probe_vision(model_id, warn, providers, tx).await;
1779                });
1780            },
1781            Cmd::CopyToClipboard(text) => {
1782                let tx = self.msg_tx.clone();
1783                self.detached.spawn(async move {
1784                    dispatch_copy_to_clipboard(text, tx).await;
1785                });
1786            },
1787            Cmd::Exit => {
1788                // The main loop observes `state.should_exit` after
1789                // the reducer returns; the runner doesn't need to
1790                // take any special action. Documented here for
1791                // exhaustiveness.
1792            },
1793            Cmd::SetTerminalTitle(title) => {
1794                if !self.terminal_title_enabled {
1795                    return;
1796                }
1797                // Offload the terminal write to the blocking pool: writing to
1798                // stdout can block when the terminal (or a downstream pipe) is
1799                // slow, and an async worker must not block on it (#44). The
1800                // OSC-2 title sequence is out-of-band relative to the renderer's
1801                // frame draws, so it doesn't corrupt them.
1802                self.detached.spawn_blocking(move || {
1803                    use std::io::Write;
1804                    let seq = format!("\x1b]2;{title}\x07");
1805                    let mut stdout = std::io::stdout();
1806                    let _ = stdout.write_all(seq.as_bytes());
1807                    let _ = stdout.flush();
1808                });
1809            },
1810            Cmd::AlertUser => {
1811                if !self.terminal_title_enabled {
1812                    return;
1813                }
1814                // A single BEL nudges the terminal to alert (dock bounce / tab
1815                // highlight). Offloaded to the blocking pool like the title.
1816                self.detached.spawn_blocking(|| {
1817                    use std::io::Write;
1818                    let mut stdout = std::io::stdout();
1819                    let _ = stdout.write_all(b"\x07");
1820                    let _ = stdout.flush();
1821                });
1822            },
1823        }
1824    }
1825
1826    fn queue_persistence(&mut self, job: PersistenceJob) {
1827        let previous = self.persistence_tail.take();
1828        let state = Arc::clone(&self.persistence_state);
1829        let tx = self.msg_tx.clone();
1830        self.persistence_tail = Some(tokio::spawn(async move {
1831            if let Some(previous) = previous
1832                && let Err(error) = previous.await
1833            {
1834                tracing::warn!(error = %error, "previous persistence job panicked");
1835            }
1836
1837            let result = tokio::task::spawn_blocking(move || {
1838                state
1839                    .lock()
1840                    .unwrap_or_else(|error| error.into_inner())
1841                    .process(job)
1842            })
1843            .await;
1844
1845            match result {
1846                Ok((events, outcome)) => {
1847                    // Events report durable writes even when the job as a
1848                    // whole failed — a partially drained barrier already
1849                    // persisted those archives, and they are never re-emitted.
1850                    if outcome.is_ok() || !events.is_empty() {
1851                        let _ = tx.send(Msg::SessionSaved).await;
1852                    }
1853                    for event in events {
1854                        fire_compaction_hook(&event).await;
1855                    }
1856                    if let Err(error) = outcome {
1857                        tracing::warn!(
1858                            error = %error,
1859                            "persistence job failed; compaction barriers remain queued"
1860                        );
1861                    }
1862                },
1863                Err(error) => tracing::warn!(error = %error, "persistence job panicked"),
1864            }
1865        }));
1866    }
1867
1868    /// Async shutdown: cancel every scope, then wait for all spawned
1869    /// work to drain. Bounded by 5 seconds — a hung task past that
1870    /// gets aborted outright by `JoinSet::drop`.
1871    pub async fn shutdown(mut self) {
1872        for (id, scope) in self.scopes.iter() {
1873            tracing::debug!(turn = %id, "shutdown: cancelling scope");
1874            scope.cancel();
1875        }
1876
1877        // The config watcher (#45) is a perpetual loop in `detached`; abort it
1878        // so the drain below doesn't block on it until the bounded timeout.
1879        if let Some(handle) = self.config_watch.take() {
1880            handle.abort();
1881        }
1882
1883        // Drain with a bounded timeout.
1884        let shutdown_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
1885
1886        let owns_global_mcp = self.owns_global_mcp;
1887        let persistence_tail = self.persistence_tail.take();
1888        let persistence_state = Arc::clone(&self.persistence_state);
1889        let drain = async {
1890            if let Some(tail) = persistence_tail
1891                && let Err(error) = tail.await
1892            {
1893                tracing::warn!(error = %error, "shutdown: persistence chain panicked");
1894            }
1895            match tokio::task::spawn_blocking(move || {
1896                let mut state = persistence_state
1897                    .lock()
1898                    .unwrap_or_else(|error| error.into_inner());
1899                let drained = state.retry_all_blocked();
1900                // Then materialize whatever the ~200-event throttle has been
1901                // holding back, so a clean exit always leaves a current
1902                // checkpoint and the next resume replays nothing.
1903                if let Err(error) = state.flush_checkpoints() {
1904                    tracing::warn!(
1905                        error = %error,
1906                        "shutdown: could not flush a session checkpoint; the log still has everything, so the next resume just folds further"
1907                    );
1908                }
1909                drop(state);
1910                drained
1911            })
1912            .await
1913            {
1914                Ok((events, outcome)) => {
1915                    // A barrier drained at shutdown still owes its hooks —
1916                    // these events are never re-emitted.
1917                    for event in events {
1918                        fire_compaction_hook(&event).await;
1919                    }
1920                    if let Err(error) = outcome {
1921                        tracing::warn!(
1922                            error = %error,
1923                            "shutdown: compaction persistence barrier retry failed"
1924                        );
1925                    }
1926                },
1927                Err(error) => tracing::warn!(
1928                    error = %error,
1929                    "shutdown: compaction persistence barrier panicked"
1930                ),
1931            }
1932            // Only the top-level runner reaps the process-global MCP manager.
1933            // A subagent's child runner shares it; reaping here would kill the
1934            // parent's servers the moment the first subagent finished.
1935            if owns_global_mcp {
1936                // If an MCP init is still in flight, its child processes are
1937                // already spawned but `set_manager` hasn't run yet — `get()`
1938                // below would return `None` and we'd leak those children. Wait
1939                // (bounded) for init to settle so the manager is installed
1940                // before we reap it (#59).
1941                let _ = tokio::time::timeout(
1942                    std::time::Duration::from_secs(2),
1943                    crate::mcp::manager_ref::wait_ready(),
1944                )
1945                .await;
1946                // Gracefully shut down MCP server children (the stdin-EOF →
1947                // terminate → kill ladder in `McpServerManager::shutdown`). The
1948                // manager lives in a `'static OnceLock` that never drops, so
1949                // this explicit call on the exit path is the only thing that
1950                // reaps those child processes. No-op when no servers were
1951                // configured.
1952                if let Some(mgr) = crate::mcp::manager_ref::get() {
1953                    mgr.shutdown().await;
1954                }
1955                // Tear down the auto-managed SearXNG process (zero-config
1956                // web_search). Same ownership rule as MCP: only the top-level
1957                // runner reaps process-global services. No-op if none started.
1958                crate::searxng::shutdown().await;
1959            }
1960            // F42: bound each per-scope drain so one non-cooperative task can't
1961            // eat the whole shutdown budget and starve the remaining scopes'
1962            // drains (the scopes were all cancelled above, so a well-behaved task
1963            // unwinds well within this). On timeout, dropping `scope` aborts its
1964            // still-running `JoinSet` members via `TurnScope::drop`.
1965            for (id, mut scope) in self.scopes.drain() {
1966                if tokio::time::timeout(CANCEL_DRAIN_TIMEOUT, scope.drain())
1967                    .await
1968                    .is_err()
1969                {
1970                    tracing::warn!(
1971                        turn = %id,
1972                        timeout_ms = CANCEL_DRAIN_TIMEOUT.as_millis(),
1973                        "shutdown: scope drain timed out; aborting its remaining tasks"
1974                    );
1975                }
1976            }
1977            while let Some(result) = self.detached.join_next().await {
1978                if let Err(e) = result
1979                    && !e.is_cancelled()
1980                {
1981                    tracing::warn!(error = %e, "shutdown: detached task panic");
1982                }
1983            }
1984        };
1985
1986        let _ = tokio::time::timeout_at(shutdown_deadline, drain).await;
1987    }
1988}
1989
1990/// Dispatch a `CallModel` command. Resolves the provider (lazy,
1991/// cached) and streams its events onto the Msg channel. Without a
1992/// bound `ProviderFactory` (unit tests), emits a single
1993/// `UpstreamError` so the reducer ends the turn cleanly.
1994/// Report a completed request's completion tokens into the task broker's
1995/// cumulative counter, so task cost deltas (`tokens_spent`) can be computed
1996/// between `in_progress` and completed stamps.
1997fn note_stream_usage(
1998    tasks: &crate::providers::TaskBroker,
1999    usage: &Option<mermaid_model::models::TokenUsage>,
2000) {
2001    if let Some(usage) = usage {
2002        tasks.add_tokens(usage.completion_tokens as u64);
2003    }
2004}
2005
2006/// Drop the built-in tool definitions the reducer suppressed for this request
2007/// (`ChatRequest::suppressed_builtin_tools` — e.g. the task-checklist writers
2008/// while a plan is being drafted). Pure so it unit-tests without the runner.
2009fn filter_suppressed(
2010    tools: Vec<mermaid_domain::ToolDefinition>,
2011    suppressed: &[&'static str],
2012) -> Vec<mermaid_domain::ToolDefinition> {
2013    if suppressed.is_empty() {
2014        return tools;
2015    }
2016    tools
2017        .into_iter()
2018        .filter(|t| !suppressed.contains(&t.name.as_str()))
2019        .collect()
2020}
2021
2022mod compaction;
2023mod memory;
2024mod model_call;
2025mod tool_call;
2026
2027use compaction::*;
2028use memory::*;
2029use model_call::*;
2030use tool_call::*;
2031
2032#[cfg(test)]
2033mod tests {
2034    /// Pins the domain-event -> durable-row field mapping. Only `id` and
2035    /// `archive_path` share a name between the two types, so every other line
2036    /// of `compaction_row` is a decision nothing else records.
2037    #[test]
2038    fn compaction_row_maps_every_field_it_claims_to() {
2039        use mermaid_domain::{CompactionEvent, CompactionReviewStatus, CompactionTrigger};
2040        let record = CompactionEvent {
2041            id: "cmp-1".to_string(),
2042            trigger: CompactionTrigger::Manual,
2043            created_at: chrono::Local::now(),
2044            before_tokens: 9_000,
2045            after_tokens: 1_200,
2046            archived_message_count: 40,
2047            preserved_message_count: 6,
2048            preserved_turn_count: 3,
2049            summary_tokens: 450,
2050            duration_secs: 1.5,
2051            review_status: CompactionReviewStatus::Reviewed,
2052            review_error: None,
2053            focus: None,
2054            archive_path: None,
2055        };
2056        let row = compaction_row(
2057            &record,
2058            std::path::Path::new("/tmp/archive.json"),
2059            Some("task-7".to_string()),
2060            "sess-3".to_string(),
2061        );
2062        assert_eq!(row.id.as_deref(), Some("cmp-1"));
2063        assert_eq!(row.task_id.as_deref(), Some("task-7"));
2064        assert_eq!(row.session_id.as_deref(), Some("sess-3"));
2065        assert_eq!(row.source_token_estimate, Some(9_000));
2066        assert_eq!(row.summary_token_count, Some(450));
2067        assert_eq!(row.preserved_turns, Some(3));
2068        assert!(row.archive_path.is_some_and(|p| p.contains("archive.json")));
2069        assert_eq!(
2070            row.verification_status.as_deref(),
2071            Some(CompactionReviewStatus::Reviewed.as_str())
2072        );
2073    }
2074
2075    use super::*;
2076    use mermaid_domain::ToolCallId;
2077    use std::time::Duration;
2078
2079    fn runner() -> (EffectRunner, mpsc::Receiver<Msg>) {
2080        EffectRunner::pair(PathBuf::from("/tmp"))
2081    }
2082
2083    /// The reducer's `suppressed_builtin_tools` contract: named tools drop
2084    /// out of the advertised set, everything else passes through in order.
2085    #[test]
2086    fn filter_suppressed_drops_only_the_named_tools() {
2087        let def = |name: &str| mermaid_domain::ToolDefinition {
2088            name: name.to_string(),
2089            description: String::new(),
2090            input_schema: serde_json::json!({}),
2091        };
2092        let tools = vec![def("task_create"), def("task_list"), def("task_update")];
2093        let kept = filter_suppressed(tools.clone(), &["task_create", "task_update"]);
2094        assert_eq!(
2095            kept.iter().map(|t| t.name.as_str()).collect::<Vec<_>>(),
2096            vec!["task_list"]
2097        );
2098        let kept = filter_suppressed(tools, &[]);
2099        assert_eq!(kept.len(), 3, "empty suppression list is a no-op");
2100    }
2101
2102    #[test]
2103    fn runtime_tool_payloads_are_redacted_before_serialization() {
2104        let payload = serde_json::json!({
2105            "url": "https://user:hunter2@example.test/page?X-Amz-Signature=opaque-signature#private",
2106            "authorization": "opaque-secret-value",
2107            "model_content": "Fetched page says OPENAI_API_KEY=sk-abcdefghijklmnop1234 and Authorization: Bearer abcdef123456ghijkl",
2108        });
2109        let serialized = redacted_json_string(&payload).expect("serialize redacted payload");
2110        assert!(
2111            !serialized.contains("hunter2"),
2112            "URL password leaked: {serialized}"
2113        );
2114        assert!(
2115            !serialized.contains("opaque-signature"),
2116            "signed URL leaked: {serialized}"
2117        );
2118        assert!(
2119            !serialized.contains("private"),
2120            "URL fragment leaked: {serialized}"
2121        );
2122        assert!(
2123            !serialized.contains("opaque-secret-value"),
2124            "credential-named field leaked: {serialized}"
2125        );
2126        assert!(
2127            !serialized.contains("abcdef123456ghijkl"),
2128            "bearer token leaked: {serialized}"
2129        );
2130        assert!(
2131            !serialized.contains("sk-abcdefghijklmnop1234"),
2132            "secret-shaped fetched content leaked: {serialized}"
2133        );
2134        assert!(serialized.contains("[REDACTED]"));
2135    }
2136
2137    #[test]
2138    fn project_walk_respects_gitignore_sorts_and_marks_dirs() {
2139        let root = std::env::temp_dir().join(format!(
2140            "mermaid-walk-{}-{:?}",
2141            std::process::id(),
2142            std::thread::current().id()
2143        ));
2144        let _ = std::fs::remove_dir_all(&root);
2145        std::fs::create_dir_all(root.join("src")).unwrap();
2146        std::fs::create_dir_all(root.join("target")).unwrap();
2147        std::fs::create_dir_all(root.join(".git")).unwrap();
2148        std::fs::write(root.join(".gitignore"), "target/\n").unwrap();
2149        std::fs::write(root.join("src/main.rs"), "fn main() {}").unwrap();
2150        std::fs::write(root.join("target/out.bin"), "ignored").unwrap();
2151        std::fs::write(root.join("README.md"), "readme").unwrap();
2152        std::fs::write(root.join(".hidden"), "hidden").unwrap();
2153
2154        let files = walk_project_files(&root);
2155        assert_eq!(
2156            files,
2157            vec![
2158                "README.md".to_string(),
2159                "src/".to_string(),
2160                "src/main.rs".to_string(),
2161            ],
2162            "sorted, dirs slash-marked, target/ ignored, dotfiles hidden"
2163        );
2164        let _ = std::fs::remove_dir_all(&root);
2165    }
2166
2167    #[test]
2168    fn new_child_suppresses_terminal_title() {
2169        // A subagent's child runner must not emit OSC 2 terminal titles —
2170        // otherwise they leak into a headless parent's stdout and corrupt
2171        // `--format json`/`text` output (caught during live headless testing).
2172        let (tx, _rx) = mpsc::channel::<Msg>(MSG_CHANNEL_CAPACITY);
2173        let providers = Arc::new(ProviderFactory::new(mermaid_domain::Config::default()));
2174        let tools = Arc::new(ToolRegistry::new());
2175        let child = EffectRunner::new_child(tx, PathBuf::from("/tmp"), providers, tools);
2176        assert!(
2177            !child.terminal_title_enabled,
2178            "subagent child runner must suppress terminal-title escapes"
2179        );
2180    }
2181
2182    #[test]
2183    fn new_child_does_not_own_global_mcp_shutdown() {
2184        // The MCP manager is process-global and shared with the parent. A
2185        // child runner's shutdown (which runs after EVERY subagent) must not
2186        // reap it — that would kill the parent's MCP servers for the rest of
2187        // the session. Only the top-level runner owns the reap.
2188        let (tx, _rx) = mpsc::channel::<Msg>(MSG_CHANNEL_CAPACITY);
2189        let providers = Arc::new(ProviderFactory::new(mermaid_domain::Config::default()));
2190        let tools = Arc::new(ToolRegistry::new());
2191        let child = EffectRunner::new_child(tx, PathBuf::from("/tmp"), providers, tools);
2192        assert!(
2193            !child.owns_global_mcp,
2194            "child runner must not reap the shared global MCP manager"
2195        );
2196        let (top, _rx2) = EffectRunner::pair(PathBuf::from("/tmp"));
2197        assert!(
2198            top.owns_global_mcp,
2199            "top-level runner still owns the global MCP reap"
2200        );
2201    }
2202
2203    #[test]
2204    fn parse_prune_plan_extracts_json_amid_prose() {
2205        let plan = parse_prune_plan(
2206            "Sure, here's the plan:\n```json\n{\"prune\": [\"a\", \"b\"], \"reason\": \"dupes\"}\n```\nDone.",
2207        )
2208        .expect("should parse");
2209        assert_eq!(plan.prune, vec!["a".to_string(), "b".to_string()]);
2210        assert_eq!(plan.reason, "dupes");
2211    }
2212
2213    #[test]
2214    fn parse_prune_plan_handles_empty_and_garbage() {
2215        let empty = parse_prune_plan("{\"prune\": [], \"reason\": \"all distinct\"}")
2216            .expect("empty plan parses");
2217        assert!(empty.prune.is_empty());
2218        assert!(parse_prune_plan("no json here").is_none());
2219    }
2220
2221    #[test]
2222    fn memory_title_from_text_is_short_and_nonempty() {
2223        assert_eq!(
2224            memory_title_from_text("prefer ripgrep over grep"),
2225            "prefer ripgrep over grep"
2226        );
2227        assert_eq!(memory_title_from_text("   "), "memory");
2228        let long = memory_title_from_text("one two three four five six seven eight nine ten");
2229        assert!(long.split_whitespace().count() <= 8);
2230    }
2231
2232    #[tokio::test]
2233    async fn dispatch_exit_is_noop_on_runner_state() {
2234        let (mut r, _rx) = runner();
2235        r.dispatch(Cmd::Exit);
2236        assert_eq!(r.scope_count(), 0);
2237    }
2238
2239    #[tokio::test]
2240    async fn dispatch_save_emits_session_saved() {
2241        let (mut r, mut rx) = runner();
2242        r.dispatch(Cmd::SaveConversation {
2243            snapshot: mermaid_domain::ConversationHistory::new(
2244                "/p".to_string(),
2245                "m".to_string(),
2246                chrono::Local::now(),
2247            ),
2248            events: Vec::new(),
2249        });
2250        let msg = tokio::time::timeout(Duration::from_millis(200), rx.recv())
2251            .await
2252            .expect("sender emits")
2253            .expect("channel alive");
2254        assert!(matches!(msg, Msg::SessionSaved));
2255    }
2256
2257    #[cfg(unix)]
2258    #[tokio::test]
2259    async fn init_mcp_servers_emits_incremental_errored_msgs() {
2260        // Two servers that both fail fast (nonexistent binaries): each
2261        // resolves independently and emits its own Errored msg; init
2262        // completes after both. Also exercises the empty-manager install.
2263        let (tx, mut rx) = tokio::sync::mpsc::channel(8);
2264        let mut configs = std::collections::HashMap::new();
2265        for name in ["one", "two"] {
2266            configs.insert(
2267                name.to_string(),
2268                mermaid_domain::McpServerConfig {
2269                    command: "/nonexistent/mermaid-test-mcp-binary".to_string(),
2270                    ..Default::default()
2271                },
2272            );
2273        }
2274        dispatch_init_mcp_servers(configs, tx).await;
2275        let mut errored = Vec::new();
2276        while let Ok(msg) = rx.try_recv() {
2277            match msg {
2278                Msg::McpServerErrored { name, .. } => errored.push(name),
2279                other => panic!("unexpected msg: {other:?}"),
2280            }
2281        }
2282        errored.sort();
2283        assert_eq!(errored, vec!["one".to_string(), "two".to_string()]);
2284        assert!(crate::mcp::manager_ref::is_ready());
2285    }
2286
2287    #[tokio::test]
2288    async fn cancel_scope_emits_turn_cancelled_after_bounded_timeout() {
2289        let (mut r, mut rx) = runner();
2290        let turn = TurnId(77);
2291        {
2292            let scope = r.scope_mut(turn);
2293            scope.spawn(async {
2294                std::future::pending::<()>().await;
2295            });
2296        }
2297        assert_eq!(r.scope_count(), 1);
2298
2299        let start = std::time::Instant::now();
2300        r.dispatch(Cmd::CancelScope(turn));
2301        assert_eq!(r.scope_count(), 0);
2302        let msg = tokio::time::timeout(Duration::from_millis(500), rx.recv())
2303            .await
2304            .expect("bounded cancel should emit terminal message")
2305            .expect("channel alive");
2306        assert!(matches!(msg, Msg::TurnCancelled(t) if t == turn));
2307        assert!(
2308            start.elapsed() < Duration::from_millis(500),
2309            "cancel terminal message took {:?}",
2310            start.elapsed()
2311        );
2312    }
2313
2314    #[tokio::test]
2315    async fn cancel_scope_emits_turn_cancelled_even_after_reaping() {
2316        // Regression (Axis 1 #9): if a turn's tasks complete and
2317        // `reap_empty_scopes` removes the now-empty scope before the user's
2318        // cancel lands, `drop_scope` used to be a silent no-op and the reducer
2319        // stuck forever in `Cancelling`. The terminal `TurnCancelled` must fire
2320        // even when the scope is already gone.
2321        let (mut r, mut rx) = runner();
2322        let turn = TurnId(88);
2323        {
2324            let scope = r.scope_mut(turn);
2325            scope.spawn(async {}); // completes immediately
2326        }
2327        assert_eq!(r.scope_count(), 1);
2328
2329        // Let the task finish, then any dispatch reaps the now-empty scope.
2330        tokio::time::sleep(Duration::from_millis(20)).await;
2331        r.dispatch(Cmd::Exit);
2332        assert_eq!(r.scope_count(), 0, "completed scope should be reaped");
2333
2334        // The scope is gone, but the reducer is still `Cancelling`: cancel must
2335        // still produce a terminal message.
2336        r.dispatch(Cmd::CancelScope(turn));
2337        let msg = tokio::time::timeout(Duration::from_millis(500), rx.recv())
2338            .await
2339            .expect("cancel on a reaped scope must still emit a terminal message")
2340            .expect("channel alive");
2341        assert!(matches!(msg, Msg::TurnCancelled(t) if t == turn));
2342    }
2343
2344    #[tokio::test]
2345    async fn dispatch_call_model_creates_scope() {
2346        let (mut r, _rx) = runner();
2347        let turn = TurnId(7);
2348        let request = mermaid_domain::ChatRequest {
2349            model_id: "test/m".to_string(),
2350            messages: vec![],
2351            system_prompt: String::new(),
2352            instructions: None,
2353            reasoning: mermaid_model::models::ReasoningLevel::Medium,
2354            temperature: 0.7,
2355            max_tokens: 4096,
2356            tools: vec![],
2357
2358            ollama_num_ctx: None,
2359            ollama_allow_ram_offload: None,
2360            resolved_context_window: None,
2361            resolved_max_output: None,
2362            output_schema: None,
2363            suppress_auto_compact: false,
2364            suppressed_builtin_tools: Vec::new(),
2365        };
2366        r.dispatch(Cmd::CallModel { turn, request });
2367        assert_eq!(r.scope_count(), 1);
2368    }
2369
2370    /// F12: after a spawned task completes (here via the
2371    /// no-ProviderFactory error path), the next `dispatch` call reaps
2372    /// the empty scope instead of leaving an orphan entry in the map.
2373    #[tokio::test]
2374    async fn empty_scopes_are_reaped_on_next_dispatch() {
2375        let (mut r, mut rx) = runner();
2376        let turn = TurnId(42);
2377        let request = mermaid_domain::ChatRequest {
2378            model_id: "test/m".to_string(),
2379            messages: vec![],
2380            system_prompt: String::new(),
2381            instructions: None,
2382            reasoning: mermaid_model::models::ReasoningLevel::Medium,
2383            temperature: 0.7,
2384            max_tokens: 4096,
2385            tools: vec![],
2386
2387            ollama_num_ctx: None,
2388            ollama_allow_ram_offload: None,
2389            resolved_context_window: None,
2390            resolved_max_output: None,
2391            output_schema: None,
2392            suppress_auto_compact: false,
2393            suppressed_builtin_tools: Vec::new(),
2394        };
2395        r.dispatch(Cmd::CallModel { turn, request });
2396        assert_eq!(r.scope_count(), 1);
2397
2398        // Runner has no provider bindings → dispatch_call_model hits
2399        // the "not wired" error path and emits UpstreamError, then the
2400        // spawned task returns. Drain that message so we know the task
2401        // ran to completion.
2402        let msg = tokio::time::timeout(Duration::from_millis(200), rx.recv())
2403            .await
2404            .expect("upstream error arrived")
2405            .expect("channel alive");
2406        assert!(matches!(msg, Msg::UpstreamError { .. }));
2407
2408        // Give the JoinSet a tick to notice the task finished.
2409        tokio::task::yield_now().await;
2410
2411        // Any subsequent dispatch reaps the now-empty scope.
2412        r.dispatch(Cmd::SetTerminalTitle("x".to_string()));
2413        assert_eq!(
2414            r.scope_count(),
2415            0,
2416            "completed scope must be reaped on next dispatch"
2417        );
2418    }
2419
2420    #[tokio::test]
2421    async fn dispatch_execute_tool_under_turn_emits_tool_started() {
2422        let (mut r, mut rx) = runner();
2423        let turn = TurnId(7);
2424        let call_id = ToolCallId(1);
2425        let source = mermaid_model::models::tool_call::ToolCall {
2426            id: Some("c1".to_string()),
2427            function: mermaid_model::models::tool_call::FunctionCall {
2428                name: "read_file".to_string(),
2429                arguments: serde_json::json!({"path": "x"}),
2430            },
2431        };
2432        r.dispatch(Cmd::ExecuteTool {
2433            turn,
2434            call_id,
2435            source,
2436            dispatch: mermaid_domain::ToolDispatch {
2437                model_id: "ollama/test".to_string(),
2438                safety_mode: mermaid_runtime::SafetyMode::Ask,
2439                plan_file: None,
2440                plan_permissions: mermaid_domain::PlanPermissions::default(),
2441                context_percent: None,
2442                intent: None,
2443                session_id: "sess-test".to_string(),
2444                message_index: 0,
2445                scratchpad: None,
2446            },
2447        });
2448        let first = tokio::time::timeout(Duration::from_millis(200), rx.recv())
2449            .await
2450            .expect("some msg")
2451            .expect("channel alive");
2452        assert!(matches!(
2453            first,
2454            Msg::ToolStarted {
2455                turn: t,
2456                call_id: c,
2457            } if t == turn && c == call_id
2458        ));
2459    }
2460
2461    #[tokio::test]
2462    async fn cancel_scope_before_execute_tool_drops_pending_work() {
2463        let (mut r, _rx) = runner();
2464        let turn = TurnId(9);
2465        r.dispatch(Cmd::CallModel {
2466            turn,
2467            request: mermaid_domain::ChatRequest {
2468                model_id: "m".to_string(),
2469                messages: vec![],
2470                system_prompt: String::new(),
2471                instructions: None,
2472                reasoning: mermaid_model::models::ReasoningLevel::Medium,
2473                temperature: 0.7,
2474                max_tokens: 4096,
2475                tools: vec![],
2476
2477                ollama_num_ctx: None,
2478                ollama_allow_ram_offload: None,
2479                resolved_context_window: None,
2480                resolved_max_output: None,
2481                output_schema: None,
2482                suppress_auto_compact: false,
2483                suppressed_builtin_tools: Vec::new(),
2484            },
2485        });
2486        assert_eq!(r.scope_count(), 1);
2487
2488        r.dispatch(Cmd::CancelScope(turn));
2489        assert_eq!(r.scope_count(), 0);
2490    }
2491
2492    #[tokio::test]
2493    async fn tombstoned_turn_is_not_resurrected_by_late_scoped_cmd() {
2494        // F38: once a turn's scope has been cancelled (dropped + tombstoned), a
2495        // stray turn-scoped Cmd bearing the same TurnId must be dropped — not
2496        // used to spin up a fresh, un-cancelled scope via `scope_mut`'s
2497        // `or_insert_with`. Turn ids are monotonic and never reused, so such a
2498        // Cmd can only be a post-cancel straggler.
2499        let (mut r, _rx) = runner();
2500        let req = || mermaid_domain::ChatRequest {
2501            model_id: "test/m".to_string(),
2502            messages: vec![],
2503            system_prompt: String::new(),
2504            instructions: None,
2505            reasoning: mermaid_model::models::ReasoningLevel::Medium,
2506            temperature: 0.7,
2507            max_tokens: 4096,
2508            tools: vec![],
2509            ollama_num_ctx: None,
2510            ollama_allow_ram_offload: None,
2511            resolved_context_window: None,
2512            resolved_max_output: None,
2513            output_schema: None,
2514            suppress_auto_compact: false,
2515            suppressed_builtin_tools: Vec::new(),
2516        };
2517        let turn = TurnId(123);
2518
2519        r.dispatch(Cmd::CallModel {
2520            turn,
2521            request: req(),
2522        });
2523        assert_eq!(r.scope_count(), 1);
2524
2525        // Cancel: drops the scope and tombstones the turn.
2526        r.dispatch(Cmd::CancelScope(turn));
2527        assert_eq!(r.scope_count(), 0);
2528
2529        // A late scoped Cmd for the now-tombstoned turn must be dropped.
2530        r.dispatch(Cmd::CallModel {
2531            turn,
2532            request: req(),
2533        });
2534        assert_eq!(
2535            r.scope_count(),
2536            0,
2537            "a cancelled turn must not be resurrected by a late scoped Cmd"
2538        );
2539
2540        // A fresh, higher turn id is unaffected by the tombstone.
2541        r.dispatch(Cmd::CallModel {
2542            turn: TurnId(124),
2543            request: req(),
2544        });
2545        assert_eq!(
2546            r.scope_count(),
2547            1,
2548            "a fresh turn must still create its scope normally"
2549        );
2550    }
2551
2552    #[tokio::test]
2553    async fn shutdown_drains_pending_saves() {
2554        let (mut r, _rx) = runner();
2555        for _ in 0..5 {
2556            r.dispatch(Cmd::SaveConversation {
2557                snapshot: mermaid_domain::ConversationHistory::new(
2558                    "/p".to_string(),
2559                    "m".to_string(),
2560                    chrono::Local::now(),
2561                ),
2562                events: Vec::new(),
2563            });
2564        }
2565        // Shutdown waits for all five to complete (should be instant).
2566        let start = std::time::Instant::now();
2567        r.shutdown().await;
2568        assert!(start.elapsed() < Duration::from_secs(2));
2569    }
2570
2571    fn persistence_fixture(
2572        root: &std::path::Path,
2573        record_id: &str,
2574    ) -> (mermaid_domain::ConversationHistory, PendingCompactionSave) {
2575        let now = chrono::Local::now();
2576        let mut full = mermaid_domain::ConversationHistory::new(
2577            root.display().to_string(),
2578            "test/model".to_string(),
2579            now,
2580        );
2581        full.add_messages(
2582            &[mermaid_model::models::ChatMessage::user("raw history")],
2583            now,
2584        );
2585        let mut compacted = full.clone();
2586        compacted.replace_messages(
2587            vec![mermaid_model::models::ChatMessage::user(
2588                "compacted checkpoint",
2589            )],
2590            now,
2591        );
2592        let record = mermaid_domain::CompactionEvent {
2593            id: record_id.to_string(),
2594            trigger: mermaid_domain::CompactionTrigger::Manual,
2595            created_at: now,
2596            before_tokens: 100,
2597            after_tokens: 20,
2598            archived_message_count: 1,
2599            preserved_message_count: 1,
2600            preserved_turn_count: 1,
2601            summary_tokens: 10,
2602            duration_secs: 0.1,
2603            review_status: mermaid_domain::CompactionReviewStatus::Reviewed,
2604            review_error: None,
2605            focus: None,
2606            archive_path: None,
2607        };
2608        (
2609            full,
2610            PendingCompactionSave {
2611                record,
2612                conversation: compacted,
2613                // The boundary event the save must land before it overwrites
2614                // the snapshot.
2615                events: vec![mermaid_domain::SessionEvent::Input {
2616                    text: "compaction boundary".to_string(),
2617                }],
2618                events_appended: false,
2619                task_id: None,
2620            },
2621        )
2622    }
2623
2624    /// Make every event append for `id` fail, by planting a directory where
2625    /// its log file goes. This is the failure the barrier exists for now
2626    /// that the boundary event -- not an archive file -- is the only record
2627    /// of a compaction's dropped messages.
2628    fn block_event_log(root: &std::path::Path, id: &str) {
2629        let dir = root.join(".mermaid").join("conversations");
2630        std::fs::create_dir_all(&dir).expect("conversations dir");
2631        std::fs::create_dir_all(dir.join(format!("{id}.jsonl"))).expect("plant a blocker");
2632    }
2633
2634    /// One `Message` event plus the snapshot that now contains it.
2635    fn one_message_save(
2636        conversation: &mut mermaid_domain::ConversationHistory,
2637        text: &str,
2638    ) -> PersistenceJob {
2639        let message = mermaid_model::models::ChatMessage::user(text);
2640        conversation.add_messages(std::slice::from_ref(&message), chrono::Local::now());
2641        PersistenceJob::Conversation {
2642            snapshot: Box::new(conversation.clone()),
2643            events: vec![mermaid_domain::SessionEvent::Message { message }],
2644        }
2645    }
2646
2647    #[test]
2648    fn the_checkpoint_stops_being_written_on_every_save() {
2649        // The point of the throttle: appends stay O(1) per message while
2650        // the whole-transcript rewrite happens on a coarse cadence. What
2651        // must NOT change is what a resume sees.
2652        let root = std::env::temp_dir().join(format!(
2653            "mermaid-throttle-{}-{:?}",
2654            std::process::id(),
2655            std::thread::current().id()
2656        ));
2657        let _ = std::fs::remove_dir_all(&root);
2658        let manager = crate::session::ConversationManager::new(&root).unwrap();
2659        let mut conversation = mermaid_domain::ConversationHistory::new(
2660            root.display().to_string(),
2661            "test/model".to_string(),
2662            chrono::Local::now(),
2663        );
2664        let mut state = PersistenceState::new(root.clone());
2665
2666        // The first save creates the log (and its backfill) but no
2667        // checkpoint: nowhere near the threshold.
2668        state
2669            .process(one_message_save(&mut conversation, "first"))
2670            .1
2671            .unwrap();
2672        let checkpoint = manager
2673            .conversations_dir()
2674            .join(format!("{}.json", conversation.id));
2675        assert!(
2676            !checkpoint.exists(),
2677            "a single save must not rewrite the transcript"
2678        );
2679
2680        // ...and resume still sees it, because the log is the truth.
2681        let resumed = manager.load_conversation(&conversation.id).unwrap();
2682        assert_eq!(resumed.messages().len(), 1);
2683        assert_eq!(resumed.messages()[0].content, "first");
2684
2685        // Crossing the threshold materializes one.
2686        for i in 0..CHECKPOINT_EVERY_EVENTS {
2687            state
2688                .process(one_message_save(&mut conversation, &format!("m{i}")))
2689                .1
2690                .unwrap();
2691        }
2692        assert!(
2693            checkpoint.exists(),
2694            "crossing {CHECKPOINT_EVERY_EVENTS} events must materialize a checkpoint"
2695        );
2696        let resumed = manager.load_conversation(&conversation.id).unwrap();
2697        assert_eq!(resumed.messages().len(), CHECKPOINT_EVERY_EVENTS + 1);
2698        let _ = std::fs::remove_dir_all(root);
2699    }
2700
2701    #[test]
2702    fn shutdown_flushes_the_checkpoint_it_was_holding() {
2703        let root = std::env::temp_dir().join(format!(
2704            "mermaid-flush-{}-{:?}",
2705            std::process::id(),
2706            std::thread::current().id()
2707        ));
2708        let _ = std::fs::remove_dir_all(&root);
2709        let manager = crate::session::ConversationManager::new(&root).unwrap();
2710        let mut conversation = mermaid_domain::ConversationHistory::new(
2711            root.display().to_string(),
2712            "test/model".to_string(),
2713            chrono::Local::now(),
2714        );
2715        let mut state = PersistenceState::new(root.clone());
2716        state
2717            .process(one_message_save(&mut conversation, "only message"))
2718            .1
2719            .unwrap();
2720
2721        let checkpoint = manager
2722            .conversations_dir()
2723            .join(format!("{}.json", conversation.id));
2724        assert!(!checkpoint.exists());
2725        state.flush_checkpoints().unwrap();
2726        assert!(
2727            checkpoint.exists(),
2728            "a clean exit must leave a current checkpoint"
2729        );
2730        // And it carries the watermark, so the next resume replays nothing.
2731        let raw = std::fs::read_to_string(&checkpoint).unwrap();
2732        let value: serde_json::Value = serde_json::from_str(&raw).unwrap();
2733        assert!(
2734            value.get("checkpoint_seq").is_some(),
2735            "a flushed checkpoint must be placeable in its log: {raw}"
2736        );
2737        let _ = std::fs::remove_dir_all(root);
2738    }
2739
2740    #[test]
2741    fn a_failed_append_keeps_its_events_for_the_next_save() {
2742        // The append is the save now, so a dropped batch is lost data --
2743        // the reducer drains its buffer at emission and never re-offers it.
2744        let root = std::env::temp_dir().join(format!(
2745            "mermaid-unappended-{}-{:?}",
2746            std::process::id(),
2747            std::thread::current().id()
2748        ));
2749        let _ = std::fs::remove_dir_all(&root);
2750        let manager = crate::session::ConversationManager::new(&root).unwrap();
2751        let mut conversation = mermaid_domain::ConversationHistory::new(
2752            root.display().to_string(),
2753            "test/model".to_string(),
2754            chrono::Local::now(),
2755        );
2756        let mut state = PersistenceState::new(root.clone());
2757
2758        // Block the log: a directory where its file goes.
2759        std::fs::create_dir_all(
2760            manager
2761                .conversations_dir()
2762                .join(format!("{}.jsonl", conversation.id)),
2763        )
2764        .expect("plant a blocker");
2765        let job = one_message_save(&mut conversation, "must survive");
2766        assert!(state.process(job).1.is_err(), "the append must fail");
2767        assert_eq!(
2768            state.unappended.get(&conversation.id).map(Vec::len),
2769            Some(1),
2770            "the batch must be held, not dropped"
2771        );
2772
2773        // Unblock, save again: the held event goes first and both land.
2774        std::fs::remove_dir_all(
2775            manager
2776                .conversations_dir()
2777                .join(format!("{}.jsonl", conversation.id)),
2778        )
2779        .expect("unblock");
2780        state
2781            .process(one_message_save(&mut conversation, "and this one"))
2782            .1
2783            .unwrap();
2784        assert!(state.unappended.is_empty(), "the hold must clear");
2785        state.flush_checkpoints().unwrap();
2786
2787        let resumed = manager.load_conversation(&conversation.id).unwrap();
2788        let texts: Vec<&str> = resumed
2789            .messages()
2790            .iter()
2791            .map(|m| m.content.as_str())
2792            .collect();
2793        assert!(
2794            texts.contains(&"must survive"),
2795            "the event held over a failed append must reach the log: {texts:?}"
2796        );
2797        assert!(texts.contains(&"and this one"), "{texts:?}");
2798        let _ = std::fs::remove_dir_all(root);
2799    }
2800
2801    #[test]
2802    fn persistence_orders_compaction_before_newer_conversation_save() {
2803        let root = std::env::temp_dir().join(format!(
2804            "mermaid-persistence-order-{}-{:?}",
2805            std::process::id(),
2806            std::thread::current().id()
2807        ));
2808        let _ = std::fs::remove_dir_all(&root);
2809        let (full, compaction) = persistence_fixture(&root, "compact_ordered");
2810        let manager = crate::session::ConversationManager::new(&root).unwrap();
2811        manager.save_conversation(&full).unwrap();
2812
2813        let mut state = PersistenceState::new(root.clone());
2814        let (events, outcome) =
2815            state.process(PersistenceJob::Compaction(Box::new(compaction.clone())));
2816        outcome.unwrap();
2817        assert_eq!(events.len(), 1);
2818        let mut newer = compaction.conversation;
2819        let reply = mermaid_model::models::ChatMessage::assistant("new assistant reply");
2820        newer.add_messages(std::slice::from_ref(&reply), chrono::Local::now());
2821        // The event, not just the snapshot: the log is the truth now, and a
2822        // save below the checkpoint threshold writes nothing else. A fixture
2823        // that passed an empty batch would be asserting against a file this
2824        // save no longer touches.
2825        let (_, outcome) = state.process(PersistenceJob::Conversation {
2826            snapshot: Box::new(newer),
2827            events: vec![mermaid_domain::SessionEvent::Message { message: reply }],
2828        });
2829        outcome.unwrap();
2830
2831        let loaded = crate::session::ConversationManager::new(&root)
2832            .unwrap()
2833            .load_conversation(&full.id)
2834            .unwrap();
2835        assert!(
2836            loaded
2837                .messages()
2838                .iter()
2839                .any(|message| message.content == "new assistant reply")
2840        );
2841        let _ = std::fs::remove_dir_all(root);
2842    }
2843
2844    #[test]
2845    fn failed_event_append_blocks_later_stripped_conversation_save() {
2846        let root = std::env::temp_dir().join(format!(
2847            "mermaid-persistence-barrier-{}-{:?}",
2848            std::process::id(),
2849            std::thread::current().id()
2850        ));
2851        let _ = std::fs::remove_dir_all(&root);
2852        let (full, compaction) = persistence_fixture(&root, "compact_blocked");
2853        let manager = crate::session::ConversationManager::new(&root).unwrap();
2854        manager.save_conversation(&full).unwrap();
2855        block_event_log(&root, &full.id);
2856
2857        let mut state = PersistenceState::new(root.clone());
2858        assert!(
2859            state
2860                .process(PersistenceJob::Compaction(Box::new(compaction.clone())))
2861                .1
2862                .is_err()
2863        );
2864        assert!(
2865            state
2866                .process(PersistenceJob::Conversation {
2867                    snapshot: Box::new(compaction.conversation),
2868                    events: Vec::new(),
2869                })
2870                .1
2871                .is_err()
2872        );
2873        assert_eq!(state.blocked.get(&full.id).map(VecDeque::len), Some(1));
2874
2875        let loaded = crate::session::ConversationManager::new(&root)
2876            .unwrap()
2877            .load_conversation(&full.id)
2878            .unwrap();
2879        assert_eq!(loaded.messages()[0].content, "raw history");
2880        let _ = std::fs::remove_dir_all(root);
2881    }
2882
2883    #[test]
2884    fn blocked_barrier_queues_a_new_compaction_instead_of_dropping_it() {
2885        let root = std::env::temp_dir().join(format!(
2886            "mermaid-persistence-queue-{}-{:?}",
2887            std::process::id(),
2888            std::thread::current().id()
2889        ));
2890        let _ = std::fs::remove_dir_all(&root);
2891        let (full, first) = persistence_fixture(&root, "compact_first");
2892        let mut second = first.clone();
2893        second.record.id = "compact_second".to_string();
2894        block_event_log(&root, &full.id);
2895
2896        let mut state = PersistenceState::new(root.clone());
2897        assert!(
2898            state
2899                .process(PersistenceJob::Compaction(Box::new(first)))
2900                .1
2901                .is_err()
2902        );
2903        // The older barrier still fails; the new save must queue behind it —
2904        // its boundary event is the only record of the stripped messages.
2905        assert!(
2906            state
2907                .process(PersistenceJob::Compaction(Box::new(second)))
2908                .1
2909                .is_err()
2910        );
2911        let queued = state.blocked.get(&full.id).expect("barrier queue");
2912        assert_eq!(queued.len(), 2);
2913        assert_eq!(queued[0].record.id, "compact_first");
2914        assert_eq!(queued[1].record.id, "compact_second");
2915        let _ = std::fs::remove_dir_all(root);
2916    }
2917
2918    #[test]
2919    fn retry_all_blocked_attempts_every_conversation() {
2920        let root = std::env::temp_dir().join(format!(
2921            "mermaid-persistence-drain-{}-{:?}",
2922            std::process::id(),
2923            std::thread::current().id()
2924        ));
2925        let _ = std::fs::remove_dir_all(&root);
2926        let (bad_full, bad) = persistence_fixture(&root, "compact_bad");
2927        let (mut good_full, mut good) = persistence_fixture(&root, "compact_good");
2928        block_event_log(&root, &bad_full.id);
2929        // Conversation ids are millisecond timestamps; two fixtures minted in
2930        // the same instant would collide into one barrier queue. Force the
2931        // second conversation onto a distinct (still format-valid) id.
2932        good_full.id = "20990101_000000_001".to_string();
2933        good.conversation.id = good_full.id.clone();
2934
2935        let mut state = PersistenceState::new(root.clone());
2936        state
2937            .blocked
2938            .entry(bad_full.id.clone())
2939            .or_default()
2940            .push_back(bad);
2941        state
2942            .blocked
2943            .entry(good_full.id.clone())
2944            .or_default()
2945            .push_back(good);
2946
2947        // One conversation's bad disk state must not strand the other's
2948        // barrier at shutdown: the error surfaces, but the good save lands —
2949        // and its durably persisted event is reported alongside the error.
2950        let (events, outcome) = state.retry_all_blocked();
2951        assert!(outcome.is_err());
2952        assert_eq!(events.len(), 1);
2953        assert_eq!(events[0].id, "compact_good");
2954        assert!(!state.blocked.contains_key(&good_full.id));
2955        assert_eq!(state.blocked.get(&bad_full.id).map(VecDeque::len), Some(1));
2956        let loaded = crate::session::ConversationManager::new(&root)
2957            .unwrap()
2958            .load_conversation(&good_full.id)
2959            .unwrap();
2960        assert_eq!(loaded.messages()[0].content, "compacted checkpoint");
2961        let _ = std::fs::remove_dir_all(root);
2962    }
2963
2964    #[test]
2965    fn partially_drained_barrier_reports_its_persisted_events() {
2966        let root = std::env::temp_dir().join(format!(
2967            "mermaid-persistence-partial-{}-{:?}",
2968            std::process::id(),
2969            std::thread::current().id()
2970        ));
2971        let _ = std::fs::remove_dir_all(&root);
2972        let (full, good) = persistence_fixture(&root, "compact_good");
2973        let mut bad = good.clone();
2974        bad.record.id = "compact_bad".to_string();
2975        // Both saves sit in ONE queue, so the blocker cannot be the shared
2976        // log path: give the tail an id that fails validation instead.
2977        bad.conversation.id = "../invalid".to_string();
2978
2979        let mut state = PersistenceState::new(root.clone());
2980        let queue = state.blocked.entry(full.id.clone()).or_default();
2981        queue.push_back(good);
2982        queue.push_back(bad);
2983
2984        // The good save at the head of the queue persists durably before the
2985        // bad one fails. Its event must surface with the error — it is popped
2986        // and would otherwise never fire SessionSaved or the compaction hook.
2987        let (events, outcome) = state.retry_blocked(&full.id);
2988        assert!(outcome.is_err());
2989        assert_eq!(events.len(), 1);
2990        assert_eq!(events[0].id, "compact_good");
2991        assert_eq!(state.blocked.get(&full.id).map(VecDeque::len), Some(1));
2992        let _ = std::fs::remove_dir_all(root);
2993    }
2994}