Skip to main content

trusty_memory/
lib.rs

1//! MCP server (HTTP/SSE + stdio) for trusty-memory.
2//!
3//! Why: Claude Code and other MCP-aware clients integrate with trusty-memory
4//! through the standardized Model Context Protocol; we expose memory + KG
5//! tools so they can be called by name. The canonical stdio integration is
6//! `trusty-memory serve --stdio` (PR1 #919 of the #914 cutover epic), a
7//! self-contained direct MCP server that binds no HTTP port or UDS socket.
8//! The former `trusty-memory-mcp-bridge` binary and Unix-domain-socket
9//! transport were removed in PR3 (#914) once `serve --stdio` made them dead
10//! code.
11//! What: Provides `run_http` / `run_http_dynamic` / `run_http_on` (axum
12//! HTTP/SSE + REST + UI) plus an `AppState` that carries the shared
13//! `PalaceRegistry`, on-disk data root, and a lazily-initialized embedder.
14//! Test: `cargo test -p trusty-memory` validates handshake + dispatch via
15//! the in-process `handle_message` unit tests and the
16//! `tests/serve_stdio_e2e.rs` end-to-end harness.
17
18use anyhow::Result;
19use serde_json::{json, Value};
20use std::net::SocketAddr;
21use std::path::{Path, PathBuf};
22use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
23use std::sync::{Arc, OnceLock};
24use tokio::sync::{broadcast, OnceCell, RwLock};
25use trusty_common::bm25_client::Bm25Client;
26use trusty_common::mcp::initialize_response;
27use trusty_common::memory_core::embed::FastEmbedder;
28use trusty_common::memory_core::{store::ChatSessionStore, PalaceRegistry};
29use trusty_common::ChatProvider;
30
31// Why: `tracing::info` is only used by the axum HTTP-serving helpers
32//      (`run_http_on`, `spawn_uds_listener`). Pulling it in unconditionally
33//      would trigger `unused_imports` warnings when the `axum-server`
34//      feature is disabled. `SocketAddr` is still used by `bound_addr` on
35//      `AppState` so it stays unconditional.
36#[cfg(feature = "axum-server")]
37use tracing::info;
38
39/// Two-phase daemon readiness state (issues #910/#911, revised by #1970).
40///
41/// Why: The embedder cold-init (CoreML compile, 30-120 s) must never block
42/// the fast text/KG/BM25 paths that don't need it. Originally (#910/#911)
43/// this state gated a hard-error preflight that rejected every
44/// `memory_remember`/`memory_recall` call outright while `Warming` — mirrored
45/// from trusty-search's staged pipeline, #1970 replaced that with graceful
46/// degradation: writes persist immediately and defer embedding to a
47/// background task, reads return BM25 + L0/L1 results and simply omit the
48/// vector lane, all keyed off this same state.
49/// What: Two stable values stored atomically.  `Warming` (0) is the initial
50/// state; `Ready` (1) is set once the embedder has been successfully
51/// initialised by `spawn_startup_tasks`.  The transition is one-way and
52/// lock-free: a single `AtomicU8` compare-and-swap.
53/// Test: `daemon_readiness_transitions_warming_to_ready` in this module;
54///       degraded-path coverage in `tools::tests`
55///       (`remember_succeeds_and_defers_embedding_while_state_is_warming`,
56///       `recall_falls_back_to_bm25_and_l0_l1_while_warming`).
57#[derive(Debug, Clone, Copy, PartialEq, Eq)]
58pub enum DaemonReadiness {
59    /// Embedder cold-init (and/or pin scan) still in progress.
60    Warming = 0,
61    /// Embedder initialised; all handlers may proceed normally.
62    Ready = 1,
63}
64
65impl DaemonReadiness {
66    /// Decode the raw atomic value.
67    ///
68    /// Why: centralises the `0 → Warming, else Ready` mapping so every
69    /// caller loads a meaningful enum rather than comparing raw integers.
70    /// What: returns `Warming` for `0`, `Ready` for any other value (only
71    /// `1` is ever written).
72    /// Test: `daemon_readiness_from_u8` in this module.
73    pub fn from_u8(v: u8) -> Self {
74        if v == 0 {
75            Self::Warming
76        } else {
77            Self::Ready
78        }
79    }
80}
81
82pub mod activity;
83pub mod attribution;
84pub mod authz;
85pub mod bm25_supervisor;
86pub mod bootstrap;
87/// Autonomous Dreamer scheduler — spawns per-palace dream loops on daemon startup.
88///
89/// Why: issue #1529 — `Dreamer::start_with_shutdown()` was fully implemented
90/// but never called. This module wires it into the daemon so each palace gets
91/// a background dream loop that fires every 5 minutes of idle time.
92/// What: exports `spawn_dream_scheduler`, `make_shutdown_watch`, and
93/// `spawn_shutdown_bridge`. Disable with `TRUSTY_DREAM_DISABLED=1`.
94/// Test: see unit tests inside this module.
95pub mod dream_scheduler;
96/// File-descriptor usage and limit reporting for `/health`.
97///
98/// Why: expose `open_fds` / `fd_soft_limit` so operators can see the fd
99/// ceiling and current consumption without needing lsof or shell access.
100/// Test: `fd_metrics::tests::fd_metrics_returns_sane_values`.
101pub mod fd_metrics;
102/// Idle-to-disk eviction ticker: periodically drops cold palace handles to
103/// bound resident RSS. Wired in `spawn_startup_tasks` next to the dream
104/// scheduler; configured via `TRUSTY_MEMORY_IDLE_EVICT_SECS`.
105pub mod idle_evict;
106// Why (issue #226): `chat` and `web` are pure axum HTTP/SSE handler
107//      surfaces. Gating them behind the `axum-server` feature is what lets
108//      library consumers (e.g. `open-mpm` linking only `MemoryMcpService`)
109//      drop axum + tower-http entirely from their build graph.
110#[cfg(feature = "axum-server")]
111pub mod chat;
112pub mod commands;
113pub mod console_metrics;
114pub mod discovery;
115/// Supervised `serve --foreground` entry point (issue #787).
116///
117/// Why: launchd supervisors need loud failure on port collision, not silent
118/// port-walking to 7071+. Extracted to stay under the 500-line ratchet cap.
119/// What: exports `bind_foreground_port` (Fix C — abort on EADDRINUSE) and
120/// `run_http_foreground` (Fix A lock + Fix B http_addr + Fix C combined).
121/// Test: `foreground::tests::bind_foreground_port_refuses_collision`;
122/// `daemon_lock` module tests cover the lock-file logic.
123pub mod foreground;
124pub mod hook_emit;
125pub mod kg_extract;
126pub mod mcp_service;
127pub mod messaging;
128pub mod openrpc;
129/// Issue #1217: default palace-ID derivation from project identity.
130///
131/// Why: the default palace ID should reflect the project's identity
132/// (git `owner/repo`, else `parent/dir`) rather than the bare directory
133/// basename, so the same repo resolves to the same palace across checkouts.
134/// What: exports the pure `derive_palace_id` core plus
135/// `owner_repo_from_git_remote`, `parent_dir_slug`, and the
136/// `TRUSTY_MEMORY_PALACE` env-override helpers.
137/// Test: see unit tests inside this module.
138pub mod palace_id_derive;
139/// Issue #88: project-root detection and palace-slug enforcement.
140///
141/// Why: prevents unbounded palace creation by anchoring palace names to the
142/// canonical slug of the project directory that contains the CWD, or to the
143/// `personal` sentinel for non-project contexts.
144/// What: exports `find_project_root`, `project_slug_at`, `project_slug`,
145/// `validate_palace_name`, `PERSONAL_PALACE`, and `PROJECT_MARKERS`.
146/// Test: see unit tests inside this module.
147pub mod project_root;
148pub mod prompt_facts;
149pub mod prompt_log;
150pub mod service;
151pub mod startup_scan;
152pub mod tools;
153pub mod transport;
154#[cfg(feature = "axum-server")]
155pub mod web;
156
157pub use activity::{ActivityEntry, ActivityFilter, ActivityLog, ActivitySource};
158pub use attribution::{CreatorInfo, CreatorSource};
159
160/// Maximum bytes retained in the trigger-prompt excerpt embedded on a
161/// `HookFired` event.
162///
163/// Why: the full triggering prompt is sensitive and already lives in the
164/// JSONL prompt log; the activity feed only needs enough text to give an
165/// operator a glance — a single-line ~80 char preview matches the existing
166/// `drawer_content_preview` convention so dashboard rows render uniformly.
167/// What: 80 characters; longer prompts are truncated with a trailing `…`.
168/// Test: `hook_excerpt_truncates_long_prompts`.
169pub const HOOK_PROMPT_EXCERPT_CHARS: usize = 80;
170
171/// Reduce a triggering prompt to the short excerpt embedded on a
172/// `HookFired` activity event.
173///
174/// Why: see [`HOOK_PROMPT_EXCERPT_CHARS`]. Centralising the truncation rule
175/// keeps every emitter (HTTP, hook CLI handlers, future tests) producing
176/// the same preview shape so UI rendering is uniform.
177/// What: whitespace-collapses `prompt` and trims to
178/// [`HOOK_PROMPT_EXCERPT_CHARS`] chars with `…` when cut. Empty input
179/// returns an empty string.
180/// Test: `hook_excerpt_truncates_long_prompts`,
181/// `hook_excerpt_collapses_whitespace`.
182pub fn hook_prompt_excerpt(prompt: &str) -> String {
183    let normalised: String = prompt.split_whitespace().collect::<Vec<_>>().join(" ");
184    if normalised.chars().count() <= HOOK_PROMPT_EXCERPT_CHARS {
185        normalised
186    } else {
187        let kept: String = normalised
188            .chars()
189            .take(HOOK_PROMPT_EXCERPT_CHARS.saturating_sub(1))
190            .collect();
191        format!("{kept}…")
192    }
193}
194
195pub use mcp_service::MemoryMcpService;
196pub use tools::MemoryMcpServer;
197
198/// Resolve the directory that actually holds the per-palace subdirectories.
199///
200/// Why: there are two on-disk layouts in the wild. The current monorepo code
201/// treats the registry directory *itself* as the parent of per-palace dirs
202/// (`<dir>/<id>/palace.json`). The legacy standalone `trusty-memory` repo
203/// nested everything one level deeper under a `palaces/` subdirectory
204/// (`<data_dir>/palaces/<id>/palace.json`) — and that is where existing
205/// installs' data lives (e.g. 88 palaces under
206/// `~/Library/Application Support/trusty-memory/palaces/`). A daemon that uses
207/// the bare data dir as its registry root finds zero palaces because every
208/// `palace.json` sits one level below where it looked — the "palaces lost on
209/// restart" bug.
210/// What: given the standard data dir, returns `<data_dir>/palaces` when that
211/// subdirectory exists, otherwise `<data_dir>` itself. Resolving this once in
212/// `main.rs` and using the result as `AppState::data_root` keeps every call
213/// site (`status`, `palace_list`, `open_palace`, `palace_create`,
214/// `load_palaces_from_disk`) consistent without forcing a data migration.
215/// Test: `tests::resolve_palace_registry_dir_prefers_palaces_subdir` and
216/// `resolve_palace_registry_dir_falls_back_to_data_dir`.
217pub fn resolve_palace_registry_dir(data_dir: PathBuf) -> PathBuf {
218    // Issue #1939: the subdir-choice logic is hoisted into trusty-common
219    // (`palace_alias::palace_registry_dir_from`) so trusty-mpm's alias-registration
220    // path and this daemon path can never disagree on WHERE the registry (and the
221    // alias file beside it) lives. This delegates to keep a single implementation.
222    trusty_common::palace_alias::palace_registry_dir_from(data_dir)
223}
224
225/// Hook type — labels the Claude Code hook that triggered a submission.
226///
227/// Why: every hook firing produces an activity-feed entry tagged with the
228/// originating hook so operators can tell whether activity came from a user
229/// prompt (`UserPromptSubmit`), a new session (`SessionStart`), or a future
230/// hook variant. Threading this through `DaemonEvent::HookFired` lets the
231/// dashboard badge each row with the hook label.
232/// What: serde-serialised in PascalCase so the wire format matches Claude
233/// Code's own hook-name strings exactly (e.g. `"UserPromptSubmit"`).
234/// Test: `hook_type_serde_round_trips`.
235#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
236pub enum HookType {
237    /// Claude Code's `UserPromptSubmit` hook — fires on every user prompt.
238    UserPromptSubmit,
239    /// Claude Code's `SessionStart` hook — fires once at session open.
240    SessionStart,
241}
242
243impl HookType {
244    /// Stable string label used for the wire format.
245    pub fn as_str(&self) -> &'static str {
246        match self {
247            Self::UserPromptSubmit => "UserPromptSubmit",
248            Self::SessionStart => "SessionStart",
249        }
250    }
251}
252
253/// Injection kind — labels what the hook actually injected (or attempted).
254///
255/// Why: distinct from `HookType` because one hook could in principle render
256/// more than one kind of injection (e.g. SessionStart can deliver both an
257/// inbox check and bootstrap context). Tagging the rendered kind explicitly
258/// keeps the activity log searchable when that fan-out lands.
259/// What: serde-serialised as kebab-case so it matches the labels already
260/// used in the JSONL prompt log (`prompt-context-facts`,
261/// `inbox-check-messages`).
262/// Test: `injection_kind_serde_round_trips`.
263#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
264#[serde(rename_all = "kebab-case")]
265pub enum InjectionKind {
266    /// `prompt-context` hook rendered the prompt-facts block.
267    PromptContext,
268    /// `inbox-check` hook delivered unread messages.
269    InboxCheck,
270}
271
272impl InjectionKind {
273    /// Stable string label used for the wire format.
274    pub fn as_str(&self) -> &'static str {
275        match self {
276            Self::PromptContext => "prompt-context",
277            Self::InboxCheck => "inbox-check",
278        }
279    }
280}
281
282/// Live daemon events broadcast to connected SSE subscribers.
283///
284/// Why: The dashboard needs push-driven updates so palace creation, drawer
285/// add/delete, dream cycles, and aggregate status changes are visible without
286/// polling. A single broadcast channel fans out to every connected browser.
287/// What: Tagged enum serialized as `{"type": "...", ...fields}` over SSE.
288/// Test: `web::tests::sse_stream_emits_events` subscribes, triggers a
289/// mutation, and asserts the frame arrives.
290#[derive(Clone, Debug, serde::Serialize)]
291#[serde(tag = "type", rename_all = "snake_case")]
292pub enum DaemonEvent {
293    PalaceCreated {
294        id: String,
295        name: String,
296        /// Originating subsystem (HTTP, MCP, Hook). Why (issue #96): the
297        /// UI badges each row with its source so operators can tell at a
298        /// glance whether a write came from the dashboard form, an MCP
299        /// tool call, or a hook-driven path. The wire-format key is
300        /// `source` (lower-case strings via serde rename_all on
301        /// `ActivitySource`).
302        source: ActivitySource,
303    },
304    DrawerAdded {
305        palace_id: String,
306        /// Friendly palace name (Palace.name) at write time. Why: lets SSE
307        /// consumers (the dashboard activity feed) render the human-readable
308        /// label without a separate id→name lookup. Empty string if the
309        /// emitter could not resolve the name.
310        #[serde(default)]
311        palace_name: String,
312        drawer_count: usize,
313        /// Wall-clock timestamp when the drawer was added. Why: SSE
314        /// receivers want to render "just now / 2m ago" relative to the
315        /// daemon's clock, not the time the SSE frame happens to arrive.
316        timestamp: chrono::DateTime<chrono::Utc>,
317        /// Short preview of the drawer's content (whitespace-collapsed,
318        /// truncated to ~80 chars with an ellipsis when cut). Why: the TUI
319        /// activity feed and dashboard ticker want to show *what* was
320        /// stored, not just the running drawer count. Empty when the
321        /// emitter could not resolve the content (legacy clients tolerate
322        /// the missing field via `#[serde(default)]`).
323        #[serde(default)]
324        content_preview: String,
325        /// Originating subsystem (issue #96).
326        source: ActivitySource,
327    },
328    DrawerDeleted {
329        palace_id: String,
330        drawer_count: usize,
331        /// Originating subsystem (issue #96).
332        source: ActivitySource,
333    },
334    DreamCompleted {
335        palace_id: Option<String>,
336        merged: usize,
337        pruned: usize,
338        compacted: usize,
339        closets_updated: usize,
340        duration_ms: u64,
341        /// Originating subsystem (issue #96).
342        source: ActivitySource,
343    },
344    StatusChanged {
345        total_drawers: usize,
346        total_vectors: usize,
347        total_kg_triples: usize,
348    },
349    /// A Claude Code hook completed and rendered (or attempted to render) an
350    /// injection block.
351    ///
352    /// Why: pre-#XXX the activity feed only fired on drawer / palace / dream
353    /// writes, which meant a normal Claude Code session — whose only daemon
354    /// traffic is hook invocations — left the feed empty. Surfacing every
355    /// hook firing answers the user complaint "no activity in the TUI" and
356    /// gives operators a way to see how often each project palace is
357    /// actually picking up prompt-context / inbox-check work.
358    /// What: carries the resolved palace (or `None` if cwd resolution
359    /// failed), the [`HookType`] label, the [`InjectionKind`] label, the
360    /// rendered injection byte length, a short excerpt of the triggering
361    /// prompt (capped at ~80 chars; the full content stays in the JSONL
362    /// prompt log only), the timestamp, the hook's wall-clock duration,
363    /// and the [`ActivitySource`] tag (always `Hook` for this variant).
364    /// Backwards-compatible: SSE clients that do not recognise the
365    /// `hook_fired` `type` tag can safely ignore the frame.
366    HookFired {
367        /// Resolved palace id (slug) — `None` if cwd resolution failed.
368        #[serde(default)]
369        palace_id: Option<String>,
370        /// Friendly palace name at hook time — `None` if the registry
371        /// could not be consulted (HTTP path uses `palace_id` here when
372        /// no separate name is known).
373        #[serde(default)]
374        palace_name: Option<String>,
375        hook_type: HookType,
376        injection_kind: InjectionKind,
377        /// Rendered injection size in bytes (`0` when no injection was
378        /// emitted, e.g. SessionStart with an empty inbox).
379        injection_length: u64,
380        /// Short excerpt of the triggering prompt for the activity feed
381        /// display. Capped at ~80 chars with a trailing `…` when cut.
382        /// Why: the activity feed renders this directly; full prompt
383        /// content (which may be sensitive) stays in the JSONL log.
384        #[serde(default)]
385        trigger_prompt_excerpt: String,
386        timestamp: chrono::DateTime<chrono::Utc>,
387        /// Hook wall-clock duration in milliseconds.
388        duration_ms: u64,
389        /// Always `ActivitySource::Hook` for this variant; encoded explicitly
390        /// so the same dispatch path (`emit`) can persist + broadcast it.
391        source: ActivitySource,
392    },
393}
394
395/// Open the activity log under `data_root`, falling back to a per-process
396/// tempdir and finally to a no-op `Discard` variant when no writable
397/// directory is available.
398///
399/// Why (issues #96, #225): the activity log is a best-effort feature — if
400/// the data root is on a read-only mount, missing, or locked by another
401/// process, the daemon should still come up and serve every other endpoint.
402/// The first fallback is a `std::env::temp_dir()`-anchored subdirectory
403/// keyed by the daemon's process id. Issue #225: a previous version called
404/// `expect()` on the tempdir fallback, which crashed the daemon on hosts
405/// where neither `data_root` nor `std::env::temp_dir()` is writable
406/// (read-only containers, locked-down sandboxes). The contract is
407/// "best-effort", so the final fallback is now `ActivityLog::discard()` —
408/// a no-op variant that drops every append and returns empty reads. The
409/// dashboard's activity feed simply shows up empty in that degraded state.
410/// What: tries `ActivityLog::open(data_root)`; on error logs a warning and
411/// retries against `<temp>/trusty-memory-activity-<pid>/`. If both fail,
412/// emits a final warning and returns `ActivityLog::discard()`.
413/// Test: `open_activity_log_with_fallback_returns_discard_when_unwritable`
414/// covers the discard branch; existing `AppState` construction tests cover
415/// the happy and tempdir-fallback paths.
416fn open_activity_log_with_fallback(data_root: &Path) -> Arc<ActivityLog> {
417    match ActivityLog::open(data_root) {
418        Ok(log) => Arc::new(log),
419        Err(primary_err) => {
420            tracing::warn!(
421                "could not open activity log at {}: {primary_err:#}; falling back to per-process tempdir",
422                data_root.display()
423            );
424            let fallback =
425                std::env::temp_dir().join(format!("trusty-memory-activity-{}", std::process::id()));
426            match ActivityLog::open(&fallback) {
427                Ok(log) => Arc::new(log),
428                Err(fallback_err) => {
429                    tracing::warn!(
430                        "activity log tempdir fallback at {} also failed: {fallback_err:#}; \
431                         activity feed disabled for this process (no-op log)",
432                        fallback.display()
433                    );
434                    Arc::new(ActivityLog::discard())
435                }
436            }
437        }
438    }
439}
440
441impl DaemonEvent {
442    /// Short discriminant label matching the SSE `type` field.
443    ///
444    /// Why: the persisted activity log stores `event_type` as a string so
445    /// the UI can render the row without re-parsing the payload. Sharing
446    /// the same labels the SSE serializer uses keeps the wire and the
447    /// stored history consistent.
448    /// What: returns one of `palace_created`, `drawer_added`,
449    /// `drawer_deleted`, `dream_completed`, `status_changed`.
450    /// Test: `daemon_event_type_str_matches_sse_tag` in the lib tests.
451    pub fn type_str(&self) -> &'static str {
452        match self {
453            Self::PalaceCreated { .. } => "palace_created",
454            Self::DrawerAdded { .. } => "drawer_added",
455            Self::DrawerDeleted { .. } => "drawer_deleted",
456            Self::DreamCompleted { .. } => "dream_completed",
457            Self::StatusChanged { .. } => "status_changed",
458            Self::HookFired { .. } => "hook_fired",
459        }
460    }
461
462    /// `palace_id` if the event is scoped to a single palace.
463    ///
464    /// Why: the activity log indexes entries by palace id so the UI can
465    /// filter by palace; daemon-wide events (`status_changed`,
466    /// dream-across-all-palaces) return `None`.
467    /// What: returns a borrowed string when the variant carries a palace
468    /// id, otherwise `None`.
469    /// Test: `daemon_event_palace_id_extraction`.
470    pub fn palace_id(&self) -> Option<&str> {
471        match self {
472            Self::PalaceCreated { id, .. } => Some(id),
473            Self::DrawerAdded { palace_id, .. } | Self::DrawerDeleted { palace_id, .. } => {
474                Some(palace_id)
475            }
476            Self::DreamCompleted { palace_id, .. } => palace_id.as_deref(),
477            Self::HookFired { palace_id, .. } => palace_id.as_deref(),
478            Self::StatusChanged { .. } => None,
479        }
480    }
481
482    /// Originating subsystem if the event carries one.
483    ///
484    /// Why: only mutation events carry a `source`; the aggregate
485    /// `StatusChanged` is recomputed by the daemon and has no caller, so
486    /// it returns `None`.
487    /// What: returns the variant's `source` field where present.
488    /// Test: `daemon_event_source_extraction`.
489    pub fn source(&self) -> Option<ActivitySource> {
490        match self {
491            Self::PalaceCreated { source, .. }
492            | Self::DrawerAdded { source, .. }
493            | Self::DrawerDeleted { source, .. }
494            | Self::DreamCompleted { source, .. }
495            | Self::HookFired { source, .. } => Some(*source),
496            Self::StatusChanged { .. } => None,
497        }
498    }
499}
500
501/// Shared application state passed to every request handler.
502///
503/// Why: The stdio loop and HTTP server need the same handles to the registry,
504/// data root, and embedder so MCP tools can perform real reads/writes against
505/// the live trusty-memory core. The embedder is heavy (loads ONNX weights) so
506/// we hold it behind a `OnceCell` and initialize lazily on first use.
507/// What: `Clone`-able via `Arc` fields. The registry / data root are eager;
508/// `embedder` is `Arc<OnceCell<Arc<FastEmbedder>>>` so concurrent first-use
509/// races resolve to a single shared instance.
510/// Test: `app_state_default_constructs` confirms construction without panic.
511#[derive(Clone)]
512pub struct AppState {
513    pub version: String,
514    pub registry: Arc<PalaceRegistry>,
515    pub data_root: PathBuf,
516    pub embedder: Arc<OnceCell<Arc<FastEmbedder>>>,
517    /// Optional default palace applied to MCP tool calls when the caller
518    /// omits the `palace` argument. Set via `trusty-memory serve --palace`.
519    pub default_palace: Option<String>,
520    /// Active chat provider selected at startup. `None` means no upstream is
521    /// configured (no Ollama detected and no OpenRouter key) — callers must
522    /// degrade gracefully (chat endpoint returns 412).
523    pub chat_provider: Arc<OnceCell<Option<Arc<dyn ChatProvider>>>>,
524    /// Per-palace chat-session stores, opened lazily so cold-start cost is
525    /// paid only when chat-history endpoints are hit.
526    pub session_stores: Arc<dashmap::DashMap<String, Arc<ChatSessionStore>>>,
527    /// Broadcast sender for live `DaemonEvent` pushes to SSE subscribers.
528    ///
529    /// Why: Lets mutating handlers emit events that any connected dashboard
530    /// receives instantly. Cap of 128 buffers transient slow readers; if a
531    /// receiver lags it gets `RecvError::Lagged` and we emit a `lag` frame.
532    pub events: Arc<broadcast::Sender<DaemonEvent>>,
533    /// Instant the daemon started, used to compute `uptime_secs` on `/health`.
534    ///
535    /// Why (issue #35): `GET /health` reports how long the daemon has been
536    /// up. Capturing a monotonic `Instant` at `AppState` construction lets the
537    /// handler compute the elapsed seconds cheaply and without a clock-skew
538    /// hazard.
539    /// What: a wall-monotonic `Instant`; `AppState::new` stamps it at startup.
540    /// Test: `health_endpoint_includes_resource_fields`.
541    pub started_at: std::time::Instant,
542    /// In-memory ring buffer of recent tracing log lines (issue #35).
543    ///
544    /// Why: the `GET /api/v1/logs/tail` endpoint serves the last N log lines
545    /// so operators can inspect a running daemon without tailing a file. The
546    /// buffer is shared between the tracing `LogBufferLayer` (writer) and the
547    /// HTTP handler (reader).
548    /// What: a cheap `Arc`-backed clone of the buffer the subscriber writes
549    /// to. Defaults to an empty buffer for states that never install the
550    /// layer (tests, the stdio path).
551    /// Test: `logs_tail_returns_recent_lines`.
552    pub log_buffer: trusty_common::log_buffer::LogBuffer,
553    /// Bug-capture ERROR store (bug-reporting #478, Phase 1).
554    ///
555    /// Why: Phase 2 MCP / HTTP endpoints need to query captured errors; stashing
556    ///      the `ErrorStore` handle here lets any handler reach it cheaply without
557    ///      a second global or per-request construction.
558    /// What: populated by `run_serve` from the `init_tracing_with_buffer_and_capture`
559    ///      result; the layer writes to this store automatically so every
560    ///      `tracing::error!` call site contributes without any changes to call
561    ///      sites. `None` in states that do not install the layer (tests, the
562    ///      stdio path).
563    /// Test: compile-presence is verified by the `trusty-memory` build; Phase 2
564    ///      will add query tests in `web.rs`.
565    pub error_store: Option<trusty_common::error_capture::ErrorStore>,
566    /// Minimal multi-tenant authorization seam (issue #1714). `false`
567    /// (single-tenant, the default) preserves today's behaviour — every
568    /// existing caller of `palace_create force=true` keeps working. `true`
569    /// opts into `authz::authorize_force_palace_create` failing closed on
570    /// every `force=true` request until a real capability check lands. Set
571    /// via `with_multi_tenant_mode_from_env` (`TRUSTY_MEMORY_MULTI_TENANT=1`);
572    /// see the `authz` module docs for the full design rationale.
573    pub multi_tenant_mode: bool,
574    /// Most recent on-disk footprint of `data_root`, in bytes (issue #35).
575    ///
576    /// Why: `GET /health` reports `disk_bytes`. Walking the data directory on
577    /// every health request would make a frequent health poll do unbounded
578    /// I/O; a background task recomputes it every 10 s and stores it here so
579    /// the handler reads it lock-free.
580    /// What: an `AtomicU64` updated by the ticker spawned in `run_http_on`.
581    /// `0` until the first walk completes.
582    /// Test: `health_endpoint_includes_resource_fields`.
583    pub disk_bytes: Arc<std::sync::atomic::AtomicU64>,
584    /// Per-process RSS + CPU sampler, refreshed on each `/health` request
585    /// (issue #35).
586    ///
587    /// Why: CPU usage is a delta between two `sysinfo` refreshes, so the
588    /// sampler must persist between requests — hence the shared `Mutex`.
589    /// What: a `tokio::sync::Mutex<SysMetrics>` so the async health handler
590    /// can sample without blocking the runtime.
591    /// Test: `health_endpoint_includes_resource_fields`.
592    pub sys_metrics: Arc<tokio::sync::Mutex<trusty_common::sys_metrics::SysMetrics>>,
593    /// HTTP listener address the daemon bound to, once `run_http_on` is running.
594    ///
595    /// Why: clients (and `/health` responses) need to advertise the live
596    /// `host:port` even though port selection happens dynamically (7070–7079
597    /// walk + OS fallback). Stashing it on `AppState` lets request handlers
598    /// surface the discovery value without re-querying the listener.
599    /// What: a `OnceLock<SocketAddr>` so `run_http_on` writes it exactly once
600    /// at bind time and every handler reads it lock-free thereafter. Empty
601    /// (`None` from `get()`) on the stdio path where no listener exists.
602    /// Test: `health_endpoint_reports_bound_addr` (added below).
603    pub bound_addr: Arc<OnceLock<SocketAddr>>,
604    /// Cached prompt-facts surface served by the MCP `get_prompt_context`
605    /// tool (issue #42).
606    ///
607    /// Why: The original session-init `prompts/get` design loaded context
608    /// once per connection; switching to a per-message tool lets the model
609    /// pull fresh, query-filtered context on demand. The cache holds both
610    /// the raw triples (for filtered lookups) and a pre-formatted Markdown
611    /// block (for the unfiltered hot path) so neither code path re-walks
612    /// the KG. The cache is rebuilt by
613    /// `prompt_facts::rebuild_prompt_cache` after any write that touches a
614    /// hot predicate (`kg_assert`, `add_alias`, `remove_prompt_fact`).
615    /// What: An `Arc<tokio::sync::RwLock<PromptFactsCache>>` so the hot
616    /// read path takes a brief read lock and clones the cache; rebuilds
617    /// take a write lock for the assignment only. The async-aware lock
618    /// (issue #229) yields to the tokio runtime instead of blocking a
619    /// runtime thread for the rebuild duration. An empty `triples` vec ↔
620    /// "no context stored yet" (the tool handler renders a hint).
621    /// Test: `get_prompt_context_returns_cached_or_hint`,
622    /// `get_prompt_context_filters_by_query`.
623    pub prompt_context_cache: Arc<RwLock<prompt_facts::PromptFactsCache>>,
624    /// Persistent activity log (issue #96).
625    ///
626    /// Why: the dashboard activity feed used to be a pure live-stream over
627    /// `/sse` — opening the UI showed an empty feed and any mutation from
628    /// the MCP path was invisible. Holding an `ActivityLog` on `AppState`
629    /// lets `emit` record an entry on every push so the
630    /// `GET /api/v1/activity` handler can return historical rows on mount
631    /// and the live SSE stream can continue prepending events on top of
632    /// the loaded history. `None` on builds that opt out (tests that use
633    /// `AppState::new` get a real log under their tempdir so behaviour
634    /// matches production).
635    /// What: an `Arc<ActivityLog>` shared with every emitter.
636    /// Test: `web::tests::activity_endpoint_lists_recent_emits`.
637    pub activity_log: Arc<ActivityLog>,
638    /// Optional per-palace BM25 lexical search lane (issue #156).
639    ///
640    /// Why: in-process BM25 would serialise the recall hot path on disk
641    /// I/O during writes and contend with the redb/usearch locks. Delegating
642    /// to the `trusty-bm25-daemon` subprocess (one socket per palace) keeps
643    /// BM25 ingestion and search off the critical path while still feeding
644    /// hits into the recall RRF fusion.
645    /// What: `Some(client)` only when `TRUSTY_BM25_DAEMON=1` at startup —
646    /// every code path that uses this field is gated on `is_some()` and
647    /// falls back to vector-only behaviour otherwise so existing deployments
648    /// see zero behavioural change.
649    /// Test: `bm25_client_disabled_by_default`,
650    /// `bm25_client_enabled_when_env_set`.
651    pub bm25_client: Option<Arc<Bm25Client>>,
652    /// Optional per-palace BM25 daemon spawn supervisor (issue #193).
653    ///
654    /// Why: without an in-process supervisor the BM25 daemon must be
655    /// launched out-of-band (launchd, manual `trusty-bm25-daemon`), which
656    /// is the same UX trap PR #190 fixed for trusty-embedderd. Holding a
657    /// supervisor here lets us spawn the daemon on first BM25 use for a
658    /// palace, restart it if it dies, and reap it on clean shutdown.
659    /// `Some` only when `TRUSTY_BM25_DAEMON=1` at startup — the same gate
660    /// that enables `bm25_client`. When set but `TRUSTY_BM25_EXTERNAL=1`,
661    /// the supervisor's `ensure_running` becomes a no-op that just returns
662    /// the canonical socket path so operators can keep using their own
663    /// process manager.
664    /// Test: covered by `bm25_supervisor_present_when_env_set` and the
665    /// `bm25_supervisor::tests` unit tests.
666    pub bm25_supervisor: Option<Arc<bm25_supervisor::Bm25Supervisor>>,
667    /// Per-palace write serialisation locks (issue #230).
668    ///
669    /// Why: the dedup gate in `tools.rs` previously read a snapshot of
670    /// existing drawers, checked for near-duplicates via Jaro-Winkler, and
671    /// then issued the write — a classic time-of-check/time-of-use race.
672    /// Two concurrent `memory_remember` calls with the same content could
673    /// both see the pre-write snapshot, both pass the gate, and both land
674    /// duplicate drawers. Serialising the gate-then-write sequence per
675    /// palace closes the window: while one task holds the mutex, any
676    /// concurrent writer for the same palace blocks until the first write
677    /// finishes and is visible to `list_drawers`. The lock is **per
678    /// palace** (not global) so writes to different palaces continue to
679    /// run in parallel.
680    /// What: a `DashMap` keyed by palace id, where each entry is an
681    /// `Arc<tokio::sync::Mutex<()>>`. The mutex is constructed lazily by
682    /// `palace_write_lock` on first access. `Arc` lets callers hold a
683    /// clone of the lock past the lifetime of the `DashMap` entry so the
684    /// map never needs to be held across an `.await`.
685    /// Test: `tools::tests::dedup_gate_blocks_concurrent_duplicate_writes`.
686    pub palace_write_locks: Arc<dashmap::DashMap<String, Arc<tokio::sync::Mutex<()>>>>,
687    /// Counter of in-flight activity-log writes spawned by `emit`
688    /// (issue #232).
689    ///
690    /// Why: `emit` offloads the synchronous redb append to the tokio blocking
691    /// pool via `spawn_blocking` so the async runtime is never parked waiting
692    /// on fsync. The write is fire-and-forget — `emit` returns immediately
693    /// after spawning. Tests that observe the activity log right after a
694    /// burst of `emit` calls need a deterministic synchronization point;
695    /// holding an in-flight counter lets `flush_activity_writes` poll until
696    /// every spawned append has settled, which keeps the assertions
697    /// race-free without forcing every caller to `.await`.
698    /// What: an `Arc<AtomicUsize>` incremented before each `spawn_blocking`
699    /// and decremented inside the closure (after the append completes, even
700    /// if it errored). The counter is cheap (one atomic add per emit) and
701    /// stays at zero in steady-state production traffic.
702    /// Test: `web::tests::activity_endpoint_lists_recent_emits` and
703    /// `tests::emit_persists_mutations_but_skips_status_changed` call
704    /// `flush_activity_writes` to drain the counter before reading the log.
705    pub pending_activity_writes: Arc<AtomicUsize>,
706    /// In-memory cache mapping palace id → `Palace.name` (issue #228).
707    ///
708    /// Why: every `memory_remember` / `memory_note` write used to call
709    /// `PalaceRegistry::list_palaces` (a synchronous filesystem walk of the
710    /// data root) just to resolve a friendly palace name for the SSE
711    /// `DrawerAdded` event. With N palaces on disk the cost was O(N) opendirs
712    /// plus `palace.json` reads on every write, blocking the async runtime.
713    /// Caching the name in-memory turns the lookup into a `DashMap::get`.
714    /// What: `DashMap<String, String>` populated by `create_palace` and
715    /// `load_palaces_from_disk`, kept in sync by rename / delete paths.
716    /// Missing entries are treated as "name unknown" so callers fall back to
717    /// the palace id and the emit path never fails.
718    /// Test: `palace_name_cache_populated_after_hydration` and
719    /// `palace_name_cache_updates_on_create`.
720    pub palace_names: Arc<dashmap::DashMap<String, String>>,
721    /// Single-pass startup pin-file map: palace id → project root path (issue #470).
722    ///
723    /// Why: after daemon startup we have no record of which on-disk project
724    /// directories correspond to which palace ids — that information only
725    /// existed inside the pin files on disk. Eager-opening every palace on
726    /// startup is too expensive. This field captures the scan-only result of
727    /// `startup_scan::scan_pin_map` so handlers that want to locate a project
728    /// by its palace id (e.g. future cwd-inference, project-health checks)
729    /// can do a single `DashMap::get` instead of a filesystem walk.
730    /// Populated once, shortly after `load_palaces_from_disk` returns, by
731    /// `spawn_startup_tasks`. Never mutated after population — it is a
732    /// snapshot of what the filesystem looked like at startup.
733    /// What: `DashMap<String (palace_id), PathBuf (project root)>`.
734    /// The outer `Arc` lets `spawn_startup_tasks` (which holds only a clone
735    /// of `AppState`) write to the same backing map that request handlers
736    /// read. Population is asynchronous so callers must treat an absent entry
737    /// as "not yet scanned" (or "no pin found"), never as "palace unknown".
738    /// Test: `startup_scan::tests::scan_pin_map_*` validate the underlying
739    /// scanner function; the wiring in `spawn_startup_tasks` is covered by
740    /// the integration-test daemon start path.
741    pub pin_project_map: Arc<dashmap::DashMap<String, PathBuf>>,
742    /// Bounded sender for the BM25 index worker (issue #231).
743    ///
744    /// Why: the previous fire-and-forget design `tokio::spawn`ed one task per
745    /// `memory_remember` / `memory_note` call, so a write burst against a slow
746    /// or unreachable BM25 daemon grew an unbounded in-flight task queue. A
747    /// single long-lived worker draining a bounded mpsc channel caps that
748    /// back-pressure: writers `try_send` (never block), full-queue requests
749    /// are dropped with a `warn!`, and the worker exits cleanly when the last
750    /// sender is dropped on shutdown.
751    /// What: an `mpsc::Sender` cloned to every `AppState` clone (cheap). The
752    /// matching receiver is consumed by the worker spawned in
753    /// [`AppState::new`] via [`tools::spawn_bm25_index_worker`]. Capacity is
754    /// [`tools::BM25_INDEX_QUEUE_CAPACITY`] (256).
755    /// Test: `bm25_index_queue_drops_when_full` exercises the full-queue
756    /// branch via `bm25_index_enqueue`.
757    pub bm25_index_tx: tokio::sync::mpsc::Sender<tools::Bm25IndexRequest>,
758    /// Cached result of the startup update check (issue #537).
759    ///
760    /// Why: `/health` should report `update_available` without hitting crates.io
761    /// on every probe. A single background check at daemon startup stores the
762    /// result here; the health handler reads it lock-free (well, a brief mutex
763    /// lock) without a network call.
764    /// What: `None` = up-to-date or check not yet done; `Some("x.y.z")` = newer
765    /// version available. The field is populated by a `tokio::spawn` in
766    /// `spawn_startup_tasks` (main.rs) after the daemon binds.
767    /// Test: indirectly by the `/health` endpoint tests in `web.rs`.
768    pub update_available: Arc<std::sync::Mutex<Option<String>>>,
769    /// Two-phase readiness state — `Warming` until the embedder is initialised,
770    /// then `Ready` (issues #910 / #911).
771    ///
772    /// Why: `AppState::embedder()` used to call `FastEmbedder::new()` without
773    /// any timeout, so the first `memory_recall`/`memory_remember` that arrived
774    /// before CoreML finished compiling would block for 5–11 hours until the
775    /// OnceCell resolved (issue #910). Exposing this state lets the preflight
776    /// guards in `tools.rs` return an explicit fast error immediately —
777    /// `"trusty-memory is warming up, retry shortly"` — instead of queueing
778    /// behind an open-ended init.
779    /// What: An `AtomicU8` starting at `DaemonReadiness::Warming` (0) and flipped
780    /// to `DaemonReadiness::Ready` (1) by `spawn_startup_tasks` after the embedder
781    /// warm-up succeeds.  The transition is one-way and lock-free.
782    /// Test: `daemon_readiness_transitions_warming_to_ready`.
783    pub daemon_readiness: Arc<AtomicU8>,
784}
785
786impl AppState {
787    /// Construct an `AppState` rooted at the given on-disk data directory.
788    ///
789    /// Why: The CLI (`serve`) and integration tests need to point the MCP
790    /// server at different roots — production at `dirs::data_dir`, tests at a
791    /// `tempfile::tempdir()`.
792    /// What: Builds an empty `PalaceRegistry`, captures the version, and
793    /// allocates an empty `OnceCell` for the embedder. `default_palace` is
794    /// `None`; use `with_default_palace` to set it.
795    /// Test: `tools::tests::dispatch_palace_create_persists` constructs an
796    /// AppState pointed at a tempdir and round-trips a palace through it.
797    pub fn new(data_root: PathBuf) -> Self {
798        let (events_tx, _) = broadcast::channel::<DaemonEvent>(128);
799        // Issue #96: open (or create) the persistent activity log under the
800        // daemon data root. Open failure is logged but never crashes the
801        // daemon — we fall back to a per-process tempdir so emits remain
802        // best-effort and the rest of the daemon keeps working.
803        let activity_log = open_activity_log_with_fallback(&data_root);
804        // Issue #231: bounded mpsc channel + single long-lived worker
805        // replaces the per-write `tokio::spawn` fire-and-forget pattern so
806        // BM25 indexing back-pressure is capped. The worker is spawned here
807        // unconditionally so the channel always has a drain — even when
808        // `bm25_client` is `None`, the worker just consumes and discards
809        // each request so senders never block on a full queue.
810        let (bm25_index_tx, bm25_index_rx) =
811            tokio::sync::mpsc::channel::<tools::Bm25IndexRequest>(tools::BM25_INDEX_QUEUE_CAPACITY);
812        // `bm25_client` / `bm25_supervisor` start as `None`; the builder
813        // `with_bm25_client_from_env` rebuilds the worker with the real
814        // client + supervisor once env-gated opt-in is resolved.
815        tools::spawn_bm25_index_worker(bm25_index_rx, None, None);
816        Self {
817            version: env!("CARGO_PKG_VERSION").to_string(),
818            // Idle-to-disk: honour TRUSTY_MEMORY_MAX_OPEN_PALACES (default 64)
819            // so operators can bound resident-palace RAM without a rebuild.
820            registry: Arc::new(PalaceRegistry::from_env()),
821            data_root,
822            embedder: Arc::new(OnceCell::new()),
823            default_palace: None,
824            chat_provider: Arc::new(OnceCell::new()),
825            session_stores: Arc::new(dashmap::DashMap::new()),
826            events: Arc::new(events_tx),
827            started_at: std::time::Instant::now(),
828            // Default to an empty buffer — `with_log_buffer` overrides this
829            // when the daemon installs the `LogBufferLayer` (HTTP mode).
830            log_buffer: trusty_common::log_buffer::LogBuffer::new(
831                trusty_common::log_buffer::DEFAULT_LOG_CAPACITY,
832            ),
833            // Bug-reporting #478: `None` until `with_error_store` is called
834            // during daemon startup (HTTP mode). Tests keep `None` so no
835            // unexpected files are written to the OS data dir.
836            error_store: None,
837            multi_tenant_mode: false,
838            disk_bytes: Arc::new(std::sync::atomic::AtomicU64::new(0)),
839            sys_metrics: Arc::new(tokio::sync::Mutex::new(
840                trusty_common::sys_metrics::SysMetrics::new(),
841            )),
842            bound_addr: Arc::new(OnceLock::new()),
843            prompt_context_cache: Arc::new(RwLock::new(prompt_facts::PromptFactsCache::default())),
844            activity_log,
845            bm25_client: None,
846            bm25_supervisor: None,
847            palace_write_locks: Arc::new(dashmap::DashMap::new()),
848            pending_activity_writes: Arc::new(AtomicUsize::new(0)),
849            palace_names: Arc::new(dashmap::DashMap::new()),
850            pin_project_map: Arc::new(dashmap::DashMap::new()),
851            bm25_index_tx,
852            update_available: Arc::new(std::sync::Mutex::new(None)),
853            // Start in Warming state; flipped to Ready by spawn_startup_tasks
854            // once the embedder warm-up succeeds (issues #910/#911).
855            daemon_readiness: Arc::new(AtomicU8::new(DaemonReadiness::Warming as u8)),
856        }
857    }
858
859    /// Acquire (lazily, then clone) the per-palace write mutex.
860    ///
861    /// Why (issue #230): the dedup-check + `remember_with_options` write
862    /// sequence in `tools.rs` must be atomic per palace to prevent two
863    /// concurrent identical writes from both passing the dedup gate.
864    /// Callers hold the returned `Arc<Mutex<()>>`'s guard across the gate
865    /// check and the write so the second writer blocks until the first
866    /// write is visible to `list_drawers`. Returning a clone of the `Arc`
867    /// rather than a borrow into the `DashMap` lets the caller `.await`
868    /// while holding the lock without risking a deadlock against any
869    /// future map mutation (DashMap shards are sync mutexes).
870    /// What: looks up the palace id in `palace_write_locks` and returns
871    /// a clone of the existing mutex; on the first call for a palace,
872    /// inserts a freshly-constructed `tokio::sync::Mutex<()>` first. The
873    /// `DashMap::entry().or_insert_with` API guarantees the lazy
874    /// construction is racy-safe — only one mutex is ever inserted per
875    /// palace id.
876    /// Test: `tools::tests::dedup_gate_blocks_concurrent_duplicate_writes`.
877    pub fn palace_write_lock(&self, palace_id: &str) -> Arc<tokio::sync::Mutex<()>> {
878        if let Some(existing) = self.palace_write_locks.get(palace_id) {
879            return existing.clone();
880        }
881        self.palace_write_locks
882            .entry(palace_id.to_string())
883            .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
884            .clone()
885    }
886
887    /// Look up a project root path by palace id in the startup pin-scan map.
888    ///
889    /// Why: provides a stable, cheap accessor so handlers do not reach directly
890    /// into the `DashMap` field and so the accessor can be mocked in future
891    /// tests without touching `AppState` internals. The map is populated
892    /// asynchronously by `spawn_startup_tasks` — an absent entry means either
893    /// the scan has not completed yet or no pin file claimed that id.
894    /// What: returns `Some(project_path)` when the palace id was found during
895    /// startup scan; `None` otherwise.
896    /// Test: covered indirectly via the startup-scan integration path; the
897    /// underlying map data is validated by `startup_scan::tests`.
898    pub fn pinned_project_path(&self, palace_id: &str) -> Option<PathBuf> {
899        self.pin_project_map.get(palace_id).map(|e| e.clone())
900    }
901
902    /// Builder-style: opt-in to the BM25 lexical lane (issue #156).
903    ///
904    /// Why: the BM25 subprocess is gated behind `TRUSTY_BM25_DAEMON=1` so
905    /// the default `cargo install trusty-memory` / launchd plist deployment
906    /// stays vector-only and existing test fixtures keep passing without
907    /// having to provision a daemon. Reading the env var here keeps the
908    /// gating logic in one place (the helper in `main.rs` just plumbs the
909    /// result through).
910    /// What: when `TRUSTY_BM25_DAEMON=1`, constructs one `Bm25Client` per
911    /// palace by lazy-resolving the socket path the first time the palace
912    /// id is observed. Currently we install a shared `default` client up
913    /// front and re-key on the palace id at the call site — palaces with no
914    /// daemon socket simply see search/index errors which we log + ignore.
915    /// Returns `self` unchanged when the env var is unset or set to anything
916    /// other than `1`.
917    /// Test: `bm25_client_disabled_by_default`,
918    /// `bm25_client_enabled_when_env_set`.
919    #[must_use]
920    pub fn with_bm25_client_from_env(mut self) -> Self {
921        if std::env::var("TRUSTY_BM25_DAEMON").as_deref() == Ok("1") {
922            // Install the default-palace client; per-palace clients are
923            // constructed on demand via `Bm25Client::for_palace`.
924            let default_palace = self.default_palace.as_deref().unwrap_or("default");
925            self.bm25_client = Some(Arc::new(Bm25Client::for_palace(default_palace)));
926            // Issue #193: hand-in-hand with the client, attach a spawn
927            // supervisor so the BM25 daemon is auto-started on first use
928            // for any palace. Operators who want to manage daemons
929            // out-of-band (launchd, systemd, manual) set
930            // TRUSTY_BM25_EXTERNAL=1 which makes the supervisor a no-op.
931            self.bm25_supervisor = Some(Arc::new(bm25_supervisor::Bm25Supervisor::new()));
932            // Issue #231: rebuild the bounded indexer channel + worker so
933            // the worker holds the now-populated client + supervisor. The
934            // placeholder worker installed by `AppState::new` (with `None`
935            // / `None`) drained the channel into the void — replacing the
936            // sender here closes the placeholder receiver and the
937            // placeholder worker exits cleanly. The new worker takes over
938            // as the sole drain for the indexer queue.
939            let (tx, rx) = tokio::sync::mpsc::channel::<tools::Bm25IndexRequest>(
940                tools::BM25_INDEX_QUEUE_CAPACITY,
941            );
942            tools::spawn_bm25_index_worker(
943                rx,
944                self.bm25_client.clone(),
945                self.bm25_supervisor.clone(),
946            );
947            self.bm25_index_tx = tx;
948            tracing::info!(
949                palace = default_palace,
950                "BM25 daemon client + spawn supervisor enabled (TRUSTY_BM25_DAEMON=1)"
951            );
952        }
953        self
954    }
955
956    /// Scan the palace registry directory and re-register every persisted
957    /// palace into the in-memory [`PalaceRegistry`].
958    ///
959    /// Why: `AppState::new` builds an *empty* registry, so after a daemon
960    /// restart `palace_list` / the dashboard reported zero palaces even though
961    /// dozens existed on disk — palace metadata was persisted by
962    /// `palace_create` but never re-hydrated on startup. This method closes
963    /// that gap by walking the on-disk layout (each subdirectory holding a
964    /// `palace.json` is one palace) and rebuilding a live `PalaceHandle` for
965    /// each, so recall paths see the full set immediately after a restart.
966    /// What: runs the blocking filesystem walk + per-palace `PalaceHandle::open`
967    /// on a `spawn_blocking` thread (so it never stalls the async runtime),
968    /// registers each successfully opened palace via `register_arc`, logs every
969    /// load at `debug!`, and returns the count loaded. A palace that fails to
970    /// open (corrupt index, unreadable `kg.db`, etc.) is logged at `warn!` and
971    /// skipped — one bad palace must not abort startup or crash the daemon.
972    /// `data_root` is expected to already be the palace registry directory —
973    /// `main.rs` resolves it via [`resolve_palace_registry_dir`] before
974    /// constructing the `AppState`, so the flat / legacy-`palaces/` layout
975    /// difference is handled exactly once.
976    /// Test: `tests::load_palaces_from_disk_rehydrates_registry` writes two
977    /// palaces into a tempdir, constructs an `AppState`, calls this method, and
978    /// asserts the returned count and registry contents.
979    pub async fn load_palaces_from_disk(&self) -> Result<usize> {
980        let registry_dir = self.data_root.clone();
981        let registry = self.registry.clone();
982        let palace_names = self.palace_names.clone();
983        // The directory walk and each `PalaceHandle::open` perform blocking
984        // filesystem + redb/usearch I/O — run the whole hydration on the
985        // blocking pool so it never parks an async worker thread.
986        let count = tokio::task::spawn_blocking(move || -> Result<usize> {
987            let palaces = PalaceRegistry::list_palaces(&registry_dir)?;
988            let total = palaces.len();
989            let mut loaded = 0usize;
990            let mut skipped = 0usize;
991            for palace in palaces {
992                match trusty_common::memory_core::PalaceHandle::open(&palace) {
993                    Ok(handle) => {
994                        tracing::debug!(
995                            palace = %palace.id,
996                            data_dir = %palace.data_dir.display(),
997                            "loaded palace from disk"
998                        );
999                        // Issue #228: seed the in-memory name cache so write
1000                        // hot paths (memory_remember / memory_note) can resolve
1001                        // the friendly palace name without re-walking the data
1002                        // root. Insert here (during hydration) is the single
1003                        // point of truth for restart-time population.
1004                        palace_names.insert(palace.id.0.clone(), palace.name.clone());
1005                        registry.register_arc(handle);
1006                        loaded += 1;
1007                    }
1008                    Err(e) => {
1009                        // Why (issue #467): a single bad palace (corrupt kg.db,
1010                        // stale WAL, EMFILE — "Too many open files", permissions)
1011                        // must never abort startup or block the HTTP server from
1012                        // binding. Log per-palace and keep going; the summary
1013                        // below tells operators how many were skipped without
1014                        // trawling the log.
1015                        // The palace is NOT registered in the in-memory registry,
1016                        // so the next `open_palace` call for this id will attempt
1017                        // a fresh open from disk — the lazy-reopen path. If the
1018                        // root cause was EMFILE and the fd-limit fix (#462) raised
1019                        // the soft limit to 8192, that first request will succeed.
1020                        tracing::warn!(
1021                            palace = %palace.id,
1022                            data_dir = %palace.data_dir.display(),
1023                            "skipping palace during startup hydration: {e:#}; \
1024                             will retry lazily on first access"
1025                        );
1026                        skipped += 1;
1027                    }
1028                }
1029            }
1030            tracing::info!(
1031                "palace hydration summary: loaded {loaded}/{total} ({skipped} skipped due to errors)"
1032            );
1033            Ok(loaded)
1034        })
1035        .await
1036        .map_err(|e| anyhow::anyhow!("join load_palaces_from_disk: {e}"))??;
1037        Ok(count)
1038    }
1039
1040    /// Builder-style: attach the daemon's shared [`LogBuffer`] so the
1041    /// `GET /api/v1/logs/tail` endpoint serves the same lines the tracing
1042    /// subscriber captures (issue #35).
1043    ///
1044    /// Why: `main` builds the buffer (via `init_tracing_with_buffer`) before
1045    /// constructing the `AppState`, then hands a clone here so the HTTP
1046    /// handler and the tracing layer observe the same ring.
1047    /// What: replaces the empty default buffer with the supplied one.
1048    /// Test: `logs_tail_returns_recent_lines`.
1049    #[must_use]
1050    pub fn with_log_buffer(mut self, buffer: trusty_common::log_buffer::LogBuffer) -> Self {
1051        self.log_buffer = buffer;
1052        self
1053    }
1054
1055    /// Builder-style: mark this daemon as the sole palace writer so palace
1056    /// redb files open with `OpenIntent::Writer` (issue #1487).
1057    ///
1058    /// Why: The HTTP daemon owns the write lock on every palace's `kg.redb`
1059    /// and `index.usearch.redb`. Before this fix, when a *second* daemon
1060    /// instance opened the same store it silently degraded to a read-only
1061    /// snapshot and rejected every `memory_remember` for its lifetime —
1062    /// effectively silent data loss when an MCP client routed a write to the
1063    /// rogue instance. Opening as `Writer` makes the second instance fail
1064    /// loud (after a short handoff-retry window that absorbs a graceful
1065    /// launchd `bootout`→`bootstrap` overlap) instead of serving broken
1066    /// reads-only. CLI, stdio-proxy, and test code paths never call this, so
1067    /// they keep the snapshot read-fallback (issue #59).
1068    /// What: Replaces `self.registry` with a fresh `PalaceRegistry` carrying
1069    /// `OpenIntent::Writer`.
1070    ///
1071    /// Invariant: MUST be called on a fresh, unhydrated, unshared registry —
1072    /// during startup, before `spawn_startup_tasks`/`load_palaces_from_disk`
1073    /// registers any `PalaceHandle` and before the `AppState` (hence its
1074    /// `Arc<PalaceRegistry>`) is cloned to a handler. Replacing the registry
1075    /// discards the prior `Arc`; doing so after hydration would silently drop
1076    /// live handles (data loss), and doing so after the state is shared would
1077    /// leave other clones on the stale read-only registry. The guard is a
1078    /// `debug_assert!` on the strongest cheap signals the registry exposes —
1079    /// `is_empty()` (no handles hydrated) and `Arc::strong_count == 1` (not yet
1080    /// shared) — so an ordering violation fails fast as the programmer error it
1081    /// is (the call site is startup-only and fixed). Release builds elide the
1082    /// assert; the real call site (`run_serve`) always satisfies it.
1083    /// Test: `with_writer_intent_marks_registry_writer` and
1084    /// `with_writer_intent_panics_on_hydrated_registry` in `lib_tests`.
1085    #[must_use]
1086    pub fn with_writer_intent(mut self) -> Self {
1087        // Fail fast on an ordering bug: a hydrated registry (`!is_empty`) or a
1088        // shared one (`strong_count > 1`) would silently drop live handles or
1089        // strand other clones on the stale read-only registry (issue #1487).
1090        debug_assert!(self.registry.is_empty() && Arc::strong_count(&self.registry) == 1);
1091        // Idle-to-disk: preserve the configurable open-handle cap
1092        // (TRUSTY_MEMORY_MAX_OPEN_PALACES) while marking the registry a writer.
1093        self.registry = Arc::new(PalaceRegistry::from_env().with_writer_intent());
1094        self
1095    }
1096
1097    /// Builder-style: attach the bug-capture `ErrorStore` handle (bug-reporting #478).
1098    ///
1099    /// Why: Phase 2 MCP / HTTP endpoints need a handle to the in-memory error
1100    ///      ring so they can serve `recent_errors` / `errors_by_fingerprint`
1101    ///      without disk I/O on the hot path. Installing it here — rather than
1102    ///      adding it as a separate global — keeps the state graph explicit and
1103    ///      lets tests skip it by never calling this method.
1104    /// What: stores `Some(store)` in `AppState::error_store`; the `BugCaptureLayer`
1105    ///      that writes to this store is already installed in the tracing
1106    ///      subscriber by `init_tracing_with_buffer_and_capture`. The store is
1107    ///      `Clone` (cheap `Arc` clone internally) so both the layer and this
1108    ///      field share the same underlying ring.
1109    /// Test: Phase 2 will add `error_store_captures_and_queries` in `web.rs`.
1110    #[must_use]
1111    pub fn with_error_store(mut self, store: trusty_common::error_capture::ErrorStore) -> Self {
1112        self.error_store = Some(store);
1113        self
1114    }
1115
1116    /// Builder-style: opt into multi-tenant authorization mode (issue #1714).
1117    ///
1118    /// Why: mirrors `with_bm25_client_from_env`'s pattern of keeping env-var
1119    /// gating in one place. Unset (the default) preserves today's
1120    /// single-tenant behaviour with zero change for existing callers.
1121    /// What: sets `multi_tenant_mode` from `TRUSTY_MEMORY_MULTI_TENANT=1`; see
1122    /// the `authz` module for what the flag then enforces.
1123    /// Test: `authorize_force_palace_create_denies_multi_tenant_without_capability`.
1124    #[must_use]
1125    pub fn with_multi_tenant_mode_from_env(mut self) -> Self {
1126        self.multi_tenant_mode = std::env::var("TRUSTY_MEMORY_MULTI_TENANT").as_deref() == Ok("1");
1127        self
1128    }
1129
1130    /// Send a `DaemonEvent` to all connected SSE subscribers and persist
1131    /// it to the activity log when the variant carries a source.
1132    ///
1133    /// Why: Mutating handlers call this after a successful write so the
1134    /// dashboard can update without polling. The send is best-effort —
1135    /// `broadcast::Sender::send` returns `Err` only when there are no live
1136    /// receivers, which is fine (no listeners == no work to do). Issue
1137    /// #96 additionally writes the entry to the persistent activity log
1138    /// so the feed can serve historical rows on page load and so MCP /
1139    /// HTTP / Hook origins are visible to the operator. Persistence is
1140    /// also best-effort — a write failure is logged but never blocks the
1141    /// SSE broadcast.
1142    ///
1143    /// Issue #232: the activity-log append is a synchronous redb write +
1144    /// fsync. Calling it directly on the async caller's task parked a tokio
1145    /// worker thread on disk I/O for every SSE event. We now offload the
1146    /// append to the blocking thread pool via `spawn_blocking` and return
1147    /// immediately — `emit` stays synchronous so every existing caller
1148    /// (including the sync `dispatch_hook_fired` JSON-RPC handler) keeps
1149    /// compiling unchanged. The fire-and-forget pattern matches the
1150    /// pre-fix semantics (best-effort, never blocks the SSE broadcast)
1151    /// while freeing the async runtime to do real work during the write.
1152    /// What: serialises the event for the log (skipping `StatusChanged`
1153    /// which is a recomputed aggregate, not a mutation), spawns the redb
1154    /// append on `tokio::task::spawn_blocking` keyed by a clone of the
1155    /// `Arc<ActivityLog>` and the cloned event, then sends the event over
1156    /// the broadcast channel. A `pending_activity_writes` counter is bumped
1157    /// before the spawn and decremented inside the closure so
1158    /// [`Self::flush_activity_writes`] can drain in tests.
1159    /// Test: `web::tests::sse_stream_receives_palace_created` confirms a
1160    /// subscriber observes the emitted event;
1161    /// `activity_endpoint_lists_recent_emits` confirms persistence via
1162    /// `flush_activity_writes`.
1163    pub fn emit(&self, event: DaemonEvent) {
1164        if let Some(source) = event.source() {
1165            let event_type = event.type_str();
1166            let palace_id = event.palace_id().map(|s| s.to_string());
1167            let log = Arc::clone(&self.activity_log);
1168            let event_for_log = event.clone();
1169            let pending = Arc::clone(&self.pending_activity_writes);
1170            // Pre-allocate the sequence id in the emitting thread so the
1171            // persisted order matches the emission order even when blocking-pool
1172            // workers execute the writes concurrently (issue #247). Without
1173            // this, four rapid emits would assign IDs inside their respective
1174            // `spawn_blocking` closures in a non-deterministic order.
1175            let id = log.alloc_id();
1176            pending.fetch_add(1, Ordering::SeqCst);
1177            // Why: the synchronous redb append + fsync must not park an
1178            // async worker thread (issue #232). Spawn the write on the
1179            // blocking pool; the JoinHandle is intentionally dropped —
1180            // the write is best-effort and any failure is logged below.
1181            tokio::task::spawn_blocking(move || {
1182                let result = log.append_with_id(id, source, palace_id, event_type, &event_for_log);
1183                if let Err(e) = result {
1184                    tracing::warn!("activity_log.append failed for {event_type}: {e:#}");
1185                }
1186                pending.fetch_sub(1, Ordering::SeqCst);
1187            });
1188        }
1189        let _ = self.events.send(event);
1190    }
1191
1192    /// Block (asynchronously) until every in-flight activity-log write
1193    /// spawned by [`Self::emit`] has settled.
1194    ///
1195    /// Why: `emit` offloads its redb append to `tokio::task::spawn_blocking`
1196    /// and returns immediately (issue #232). Tests that observe the
1197    /// activity log right after a burst of emits would otherwise race the
1198    /// blocking-pool worker; this helper gives them a deterministic
1199    /// synchronization point. Production code never needs to call this —
1200    /// the dashboard reads through `GET /api/v1/activity`, which already
1201    /// tolerates writes settling asynchronously.
1202    /// What: spins on `pending_activity_writes` with a 1 ms yield until the
1203    /// counter is zero. Cheap: tests typically emit a handful of events
1204    /// and the loop exits within a single scheduler tick.
1205    /// Test: covered indirectly by `emit_persists_mutations_but_skips_status_changed`
1206    /// and `web::tests::activity_endpoint_lists_recent_emits`.
1207    pub async fn flush_activity_writes(&self) {
1208        while self.pending_activity_writes.load(Ordering::SeqCst) > 0 {
1209            tokio::time::sleep(std::time::Duration::from_millis(1)).await;
1210        }
1211    }
1212
1213    /// Open (or return cached) the chat-session store for a palace.
1214    ///
1215    /// Why: Chat session persistence lives in a dedicated redb file under
1216    /// the palace's data dir (`chat_sessions.redb`) so it doesn't intermingle
1217    /// with the KG's transactional load. The store is cheap to clone via
1218    /// `Arc` but the underlying connection should be reused, so cache by id.
1219    /// What: Creates the palace data dir if missing, opens (or reuses) a
1220    /// `ChatSessionStore` and stashes an `Arc` in the DashMap.
1221    /// Test: Indirectly via the session HTTP handlers in `web::tests`.
1222    pub fn session_store(&self, palace_id: &str) -> Result<Arc<ChatSessionStore>> {
1223        if let Some(entry) = self.session_stores.get(palace_id) {
1224            return Ok(entry.clone());
1225        }
1226        let dir = self.data_root.join(palace_id);
1227        std::fs::create_dir_all(&dir)
1228            .map_err(|e| anyhow::anyhow!("create palace dir {}: {e}", dir.display()))?;
1229        let store = Arc::new(ChatSessionStore::open(&dir.join("chat_sessions.db"))?);
1230        self.session_stores
1231            .insert(palace_id.to_string(), store.clone());
1232        Ok(store)
1233    }
1234
1235    /// Builder-style setter for the default palace name.
1236    ///
1237    /// Why: `serve --palace <name>` wants to bind every tool call to a
1238    /// project-scoped namespace without forcing every MCP request to repeat
1239    /// the palace argument.
1240    /// What: Returns `self` with `default_palace = Some(name)`.
1241    /// Test: `default_palace_used_when_arg_omitted` covers the resolution
1242    /// path; this setter is exercised there.
1243    pub fn with_default_palace(mut self, name: Option<String>) -> Self {
1244        self.default_palace = name;
1245        self
1246    }
1247
1248    /// Resolve (or initialize) the shared embedder.
1249    ///
1250    /// Why: FastEmbedder load is expensive — we share one instance across all
1251    /// tool calls; the `OnceCell` ensures concurrent first-use races collapse
1252    /// to a single load.
1253    /// What: Returns `Arc<FastEmbedder>` on success. Errors propagate from the
1254    /// underlying ONNX load.
1255    /// Test: Indirectly via `dispatch_remember_then_recall`.
1256    /// Resolve the active chat provider, auto-detecting on first call.
1257    ///
1258    /// Why: Provider selection depends on filesystem-loaded config plus a
1259    /// network probe (Ollama liveness), so it must be lazily initialised at
1260    /// runtime. Caching the choice in a `OnceCell` keeps it stable across
1261    /// concurrent requests without re-probing on every chat call.
1262    /// What: On first use loads `~/.trusty-memory/config.toml`, prefers an
1263    /// auto-detected Ollama instance (when `local_model.enabled`), and falls
1264    /// back to OpenRouter when an API key is set. Returns `Ok(None)` when
1265    /// neither is available so the caller can emit a 412.
1266    /// Test: `web::tests::providers_endpoint_returns_payload` covers the
1267    /// detection path indirectly through `/api/v1/chat/providers`.
1268    pub async fn chat_provider(&self) -> Option<Arc<dyn ChatProvider>> {
1269        self.chat_provider
1270            .get_or_init(|| async {
1271                // Why (issue #226): `service::load_user_config` is the
1272                //      axum-free home of the loader; the `web::load_user_config`
1273                //      re-export only exists for the HTTP handlers. Going
1274                //      direct to `service` keeps this method usable when
1275                //      the `axum-server` feature is disabled.
1276                let cfg = crate::service::load_user_config().unwrap_or_default();
1277                if cfg.local_model.enabled {
1278                    if let Some(mut p) =
1279                        trusty_common::auto_detect_local_provider(&cfg.local_model.base_url).await
1280                    {
1281                        // auto_detect returns an empty model id; callers must
1282                        // set the configured model name themselves.
1283                        p.model = cfg.local_model.model.clone();
1284                        return Some(Arc::new(p) as Arc<dyn ChatProvider>);
1285                    }
1286                }
1287                if !cfg.openrouter_api_key.is_empty() {
1288                    return Some(Arc::new(trusty_common::OpenRouterProvider::new(
1289                        cfg.openrouter_api_key,
1290                        cfg.openrouter_model,
1291                    )) as Arc<dyn ChatProvider>);
1292                }
1293                None
1294            })
1295            .await
1296            .clone()
1297    }
1298
1299    /// Spawn a fire-and-forget background task that auto-discovers project
1300    /// aliases under `project_root` and asserts new ones into `palace`.
1301    ///
1302    /// Why (issue #42): Projects carry implicit shorthand — cargo package
1303    /// names that differ from their directory, binary names that differ
1304    /// from packages, first-letter abbreviations — that should be surfaced
1305    /// without a user ever calling `add_alias`. Running discovery as a
1306    /// detached task on palace-open keeps startup latency unchanged: the
1307    /// daemon binds and starts serving immediately while the discovery scan
1308    /// completes in the background, and any newly-asserted aliases land in
1309    /// the prompt cache before the model's next `get_prompt_context` call.
1310    /// What: clones `self` (cheap; `Arc`-backed), spawns a tokio task that
1311    /// invokes the `discover_aliases` tool handler directly so the
1312    /// dedup + cache-rebuild logic runs exactly the same path as the MCP
1313    /// tool call. Errors are logged at `warn!`; one failed discovery never
1314    /// destabilises the daemon.
1315    /// Test: not unit-tested (timing-dependent fire-and-forget); the
1316    /// underlying `discover_aliases` dispatch is covered by
1317    /// `dispatch_discover_aliases_inserts_new_and_dedupes` in `tools::tests`.
1318    pub fn spawn_alias_discovery(&self, palace: String, project_root: PathBuf) {
1319        let state = self.clone();
1320        tokio::spawn(async move {
1321            let args = serde_json::json!({
1322                "palace": palace,
1323                "project_root": project_root.to_string_lossy(),
1324            });
1325            match tools::dispatch_tool(&state, "discover_aliases", args).await {
1326                Ok(result) => tracing::info!(
1327                    new = ?result.get("new"),
1328                    already_known = ?result.get("already_known"),
1329                    "alias discovery complete"
1330                ),
1331                Err(e) => tracing::warn!("alias discovery failed: {e:#}"),
1332            }
1333        });
1334    }
1335
1336    /// Return the current readiness state.
1337    ///
1338    /// Why: tool handlers and the `/health` endpoint need a cheap, lock-free
1339    /// way to check whether the embedder has been initialised yet.
1340    /// What: loads `daemon_readiness` with `Acquire` ordering so the caller
1341    /// sees all writes the startup task made before setting the state.
1342    /// Test: `daemon_readiness_transitions_warming_to_ready`.
1343    pub fn readiness(&self) -> DaemonReadiness {
1344        DaemonReadiness::from_u8(self.daemon_readiness.load(Ordering::Acquire))
1345    }
1346
1347    /// Flip the readiness state from `Warming` to `Ready`.
1348    ///
1349    /// Why: called by `spawn_startup_tasks` in `main.rs` once the embedder
1350    /// warm-up succeeds — this is the single state-transition site.
1351    /// What: `store(Ready, Release)` so subsequent `Acquire` loads in handlers
1352    /// observe a consistent state.  Idempotent: calling it multiple times is
1353    /// harmless.
1354    /// Test: `daemon_readiness_transitions_warming_to_ready`.
1355    pub fn set_ready(&self) {
1356        self.daemon_readiness
1357            .store(DaemonReadiness::Ready as u8, Ordering::Release);
1358    }
1359
1360    /// Obtain the shared `FastEmbedder` instance, initialising it on first call.
1361    ///
1362    /// Why: centralises lazy embedder access so every tool handler goes through
1363    /// one bounded init path (tracks #910 internally).
1364    /// What: wraps `OnceCell::get_or_try_init` with a timeout so a slow
1365    /// CoreML/CUDA first-compile cannot block a handler indefinitely.  On
1366    /// timeout the `OnceCell` is left unresolved and the next caller retries.
1367    ///
1368    /// **Callers on the request path SHOULD check `readiness()` before this
1369    /// method** (issue #1970) — every recall handler now checks
1370    /// `readiness() == Ready` first and only calls `embedder()` on that
1371    /// branch, falling back to a BM25/L0/L1-only path while `Warming` instead
1372    /// of paying this method's cold-init cost. Reaching this method while
1373    /// still `Warming` is not a bug (the warm-up task itself calls
1374    /// `embedder()` while in `Warming` state), just unusual on the request
1375    /// path.
1376    ///
1377    /// This timeout is a backstop against a pathological init delay (e.g. the
1378    /// warm-up task's own call, or a handler that skips the `readiness()`
1379    /// check). If this timeout fires the `OnceCell` is left in the unresolved
1380    /// state and the next call retries from scratch.
1381    pub async fn embedder(&self) -> Result<Arc<FastEmbedder>> {
1382        use trusty_common::memory_core::timeouts;
1383        let cell = self.embedder.clone();
1384        let timeout = timeouts::embedder_init_timeout();
1385        let embedder = tokio::time::timeout(
1386            timeout,
1387            cell.get_or_try_init(|| async {
1388                let e = FastEmbedder::new().await?;
1389                Ok::<Arc<FastEmbedder>, anyhow::Error>(Arc::new(e))
1390            }),
1391        )
1392        .await
1393        .map_err(|_| {
1394            anyhow::anyhow!(
1395                "AppState::embedder() timed out after {:?}; \
1396                 the CoreML/CUDA model is taking unusually long to compile — \
1397                 increase TRUSTY_EMBEDDER_INIT_TIMEOUT_SECS if needed",
1398                timeout
1399            )
1400        })??
1401        .clone();
1402        Ok(embedder)
1403    }
1404}
1405
1406impl std::fmt::Debug for AppState {
1407    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1408        f.debug_struct("AppState")
1409            .field("version", &self.version)
1410            .field("data_root", &self.data_root)
1411            .field("registry_len", &self.registry.len())
1412            .finish()
1413    }
1414}
1415
1416/// Handle a single MCP JSON-RPC message and produce its response.
1417///
1418/// Why: Pulled out of the stdio loop so unit tests can drive every method
1419/// without touching real stdin/stdout.
1420/// What: Routes `initialize`, `tools/list`, `tools/call`, `ping`, and the
1421/// `notifications/initialized` notification (which returns `Value::Null`).
1422/// Test: See unit tests below — initialize/list/call all return expected
1423/// JSON-RPC envelopes; notifications return `Null` (no response written).
1424pub async fn handle_message(state: &AppState, msg: Value) -> Value {
1425    let id = msg.get("id").cloned().unwrap_or(Value::Null);
1426    let method = msg.get("method").and_then(|m| m.as_str()).unwrap_or("");
1427
1428    match method {
1429        "initialize" => {
1430            let extra = state
1431                .default_palace
1432                .as_ref()
1433                .map(|dp| json!({ "default_palace": dp }));
1434            let result = initialize_response("trusty-memory", &state.version, extra);
1435            // Why (issue #42): prompt-facts now flow through the
1436            // per-message `get_prompt_context` tool rather than MCP
1437            // prompts, so we no longer advertise the `prompts` capability.
1438            json!({
1439                "jsonrpc": "2.0",
1440                "id": id,
1441                "result": result,
1442            })
1443        }
1444        // Notifications must NOT receive a response.
1445        "notifications/initialized" | "notifications/cancelled" => Value::Null,
1446        "tools/list" => json!({
1447            "jsonrpc": "2.0",
1448            "id": id,
1449            "result": tools::tool_definitions_with(state.default_palace.is_some())
1450        }),
1451        // OpenRPC 1.3.2 discovery — see `openrpc.rs`. Returns the full
1452        // service description so orchestrators (open-mpm, etc.) can
1453        // introspect every tool and its required `memory.read`/`memory.write`
1454        // scope without bespoke per-server adapters.
1455        "rpc.discover" => json!({
1456            "jsonrpc": "2.0",
1457            "id": id,
1458            "result": openrpc::build_discover_response(
1459                &state.version,
1460                state.default_palace.is_some(),
1461            ),
1462        }),
1463        "tools/call" => {
1464            let params = msg.get("params").cloned().unwrap_or_default();
1465            let tool_name = params
1466                .get("name")
1467                .and_then(|n| n.as_str())
1468                .unwrap_or("")
1469                .to_string();
1470            let args = params.get("arguments").cloned().unwrap_or_default();
1471            match tools::dispatch_tool(state, &tool_name, args).await {
1472                Ok(content) => {
1473                    // Why: tools that return a bare JSON string (e.g.
1474                    // `get_prompt_context` returning the formatted
1475                    // Markdown block) should surface as plain text in the
1476                    // MCP `content[0].text` field — wrapping in
1477                    // `Value::to_string()` would re-quote the payload and
1478                    // force every caller to strip outer quotes.
1479                    let text = match &content {
1480                        Value::String(s) => s.clone(),
1481                        other => other.to_string(),
1482                    };
1483                    json!({
1484                        "jsonrpc": "2.0",
1485                        "id": id,
1486                        "result": {
1487                            "content": [{"type": "text", "text": text}]
1488                        }
1489                    })
1490                }
1491                Err(e) => json!({
1492                    "jsonrpc": "2.0",
1493                    "id": id,
1494                    // Why: anyhow's `{:#}` alternate format walks the full
1495                    // `Caused by:` chain so MCP clients see actionable
1496                    // detail (e.g. "PalaceHandle::remember_with_options:
1497                    // filter rejected: too short") instead of just the
1498                    // outermost context label.
1499                    "error": {"code": -32603, "message": format!("{e:#}")}
1500                }),
1501            }
1502        }
1503        "ping" => json!({"jsonrpc": "2.0", "id": id, "result": {}}),
1504        _ => json!({
1505            "jsonrpc": "2.0",
1506            "id": id,
1507            "error": {
1508                "code": -32601,
1509                "message": format!("Method not found: {method}")
1510            }
1511        }),
1512    }
1513}
1514
1515/// Preferred starting port for the trusty-memory HTTP daemon.
1516///
1517/// Why: keeps the well-known default stable for clients that have hard-coded
1518/// `127.0.0.1:7070` in their configuration, while still allowing dynamic
1519/// walking when the port is in use (`DYNAMIC_PORT_RANGE` ports starting here).
1520/// What: `7070` — historic default, matches the launchd plist's prior value.
1521/// Test: covered indirectly by `bind_dynamic_port_returns_listener`.
1522pub const DEFAULT_HTTP_PORT: u16 = 7070;
1523
1524/// Number of consecutive ports `bind_dynamic_port` walks before falling back
1525/// to the OS-assigned port. Matches the trusty-search convention.
1526const DYNAMIC_PORT_RANGE: u16 = 10;
1527
1528/// Path to the canonical address-discovery file for the trusty-memory daemon.
1529///
1530/// Why: clients (CLI, MCP tools, dashboards) need to find the running daemon
1531/// without configuration when the port was selected dynamically. Using
1532/// `trusty_common::resolve_data_dir` aligns this path with the location
1533/// that `trusty_common::read_daemon_addr("trusty-memory")` reads from, so
1534/// `prompt-context`, `doctor`, and `start`'s probe all find the running daemon.
1535/// The old `~/.trusty-memory/http_addr` path and the new
1536/// `~/Library/Application Support/trusty-memory/http_addr` (macOS) path were
1537/// divergent — the daemon wrote one; readers expected the other.
1538/// What: returns `{resolve_data_dir("trusty-memory")}/http_addr`, or `None` if
1539/// the data dir cannot be resolved (locked-down container, no passwd entry).
1540/// Test: `http_addr_path_uses_resolve_data_dir`.
1541pub fn http_addr_path() -> Option<PathBuf> {
1542    trusty_common::resolve_data_dir("trusty-memory")
1543        .ok()
1544        .map(|d| d.join("http_addr"))
1545}
1546
1547/// Bind a `TcpListener` to `127.0.0.1`, dynamically selecting a port.
1548///
1549/// Why: the historic default `7070` is convenient for clients but a stale
1550/// process or a second daemon must not produce a noisy failure. Walking
1551/// `DEFAULT_HTTP_PORT..DEFAULT_HTTP_PORT+DYNAMIC_PORT_RANGE` first preserves
1552/// backwards compatibility for the common case; OS-assigned fallback (`:0`)
1553/// guarantees the daemon always comes up even when every preferred port is
1554/// busy.
1555/// What: returns the first successful `TcpListener` (7070..=7079, then
1556/// OS-assigned); caller inspects `local_addr()` to learn the chosen port.
1557/// Test: `bind_dynamic_port_returns_listener` confirms it always binds *some*
1558/// port even after another listener occupies the preferred one.
1559pub async fn bind_dynamic_port() -> Result<tokio::net::TcpListener> {
1560    let preferred: SocketAddr = SocketAddr::from(([127, 0, 0, 1], DEFAULT_HTTP_PORT));
1561    // First: walk the preferred range (7070..=7079).
1562    if let Ok(listener) =
1563        trusty_common::bind_with_auto_port(preferred, DYNAMIC_PORT_RANGE - 1).await
1564    {
1565        return Ok(listener);
1566    }
1567    // Last resort: ask the kernel for any free port. `bind_with_auto_port`
1568    // with `:0` resolves immediately to the OS-assigned port.
1569    tracing::warn!(
1570        "all ports {DEFAULT_HTTP_PORT}..{} in use; requesting OS-assigned port",
1571        DEFAULT_HTTP_PORT + DYNAMIC_PORT_RANGE - 1
1572    );
1573    let any: SocketAddr = SocketAddr::from(([127, 0, 0, 1], 0));
1574    trusty_common::bind_with_auto_port(any, 0).await
1575}
1576
1577/// Write the bound `host:port` to `~/.trusty-memory/http_addr` atomically.
1578///
1579/// Why: clients must read the file mid-write without observing a partial
1580/// value. Writing to a `.tmp` sibling and renaming over the target gives
1581/// POSIX atomicity, matching the trusty-search implementation.
1582/// What: creates the parent directory if missing; writes `addr` followed by a
1583/// trailing newline (avoids the "no newline at end of file" warnings from
1584/// `cat`); renames `.tmp` → `http_addr`. Best-effort: I/O errors are
1585/// returned to the caller so `run_http_on` can log without panicking.
1586/// Test: `http_addr_file_round_trip_via_helpers`.
1587#[cfg(feature = "axum-server")]
1588fn write_http_addr_file(path: &Path, addr: &SocketAddr) -> std::io::Result<()> {
1589    use std::io::Write;
1590    if let Some(parent) = path.parent() {
1591        std::fs::create_dir_all(parent)?;
1592    }
1593    let tmp = path.with_extension("addr.tmp");
1594    {
1595        let mut f = std::fs::File::create(&tmp)?;
1596        writeln!(f, "{addr}")?;
1597        f.sync_all()?;
1598    }
1599    std::fs::rename(&tmp, path)?;
1600    Ok(())
1601}
1602
1603/// Return `true` when a non-default data directory is in effect.
1604///
1605/// Why (issue #880): two startup side-effects must be suppressed when the
1606/// daemon runs with an isolated/overridden data root:
1607/// 1. The legacy `~/.trusty-memory/http_addr` dotfile write — it would
1608///    overwrite the real production daemon's discovery file with the isolated
1609///    instance's throwaway address.
1610/// 2. The startup pin-scan — it reads project pin files from the **real**
1611///    user environment (~/Projects, ~/Developer, …) and imports palaces from
1612///    the real environment into the isolated data root, defeating isolation.
1613///
1614/// A "non-default data dir" means `TRUSTY_DATA_DIR_OVERRIDE` is set to a
1615/// non-empty, non-whitespace value. Empty or whitespace-only values are
1616/// treated as unset (same rule as `resolve_data_dir`), so an accidental blank
1617/// env var does not suppress the dotfile write on real production instances.
1618/// What: reads `TRUSTY_DATA_DIR_OVERRIDE`; returns `true` when it contains a
1619/// non-empty, non-whitespace string. Returns `false` otherwise.
1620/// Test: `is_data_dir_override_active_when_set`,
1621///       `is_data_dir_override_inactive_when_unset`,
1622///       `is_data_dir_override_inactive_when_blank`.
1623#[inline]
1624pub fn is_data_dir_override_active() -> bool {
1625    matches!(
1626        std::env::var(trusty_common::DATA_DIR_OVERRIDE_ENV),
1627        Ok(v) if !v.trim().is_empty()
1628    )
1629}
1630
1631/// Resolve the dotfile discovery path `~/.trusty-memory/http_addr`.
1632///
1633/// Why (issue #498): external tooling such as claude-mpm's `migrate_trusty_autodetect`
1634/// reads `~/.trusty-memory/http_addr` to find the running daemon's port. On
1635/// macOS, `resolve_data_dir("trusty-memory")` returns
1636/// `~/Library/Application Support/trusty-memory/`, not `~/.trusty-memory/`,
1637/// so the daemon was writing to the OS-standard location while readers expected
1638/// the dotfile location. Writing to both locations keeps every reader happy
1639/// regardless of which convention they follow.
1640///
1641/// Fix #880: returns `None` when `TRUSTY_DATA_DIR_OVERRIDE` is active so an
1642/// isolated instance (test rig, CI, parallel run) never overwrites the real
1643/// production daemon's discovery dotfile.
1644///
1645/// What: returns `$HOME/.trusty-memory/http_addr` in the default (production)
1646/// case, or `None` when `dirs::home_dir()` is unavailable OR when a data-dir
1647/// override is active (see `is_data_dir_override_active`).
1648/// Test: `dotfile_http_addr_path_uses_home_dir`,
1649///       `dotfile_suppressed_when_override_active`.
1650#[cfg(feature = "axum-server")]
1651fn dotfile_http_addr_path() -> Option<PathBuf> {
1652    // Fix #880: never write to the shared dotfile when an override is active.
1653    if is_data_dir_override_active() {
1654        return None;
1655    }
1656    dirs::home_dir().map(|h| h.join(".trusty-memory").join("http_addr"))
1657}
1658
1659/// Run the optional HTTP/SSE + web admin server.
1660///
1661/// Why: A long-running daemon mode lets non-stdio clients (browsers, curl,
1662/// future remote agents) hit `/health`, the `/api/v1/*` REST surface, and the
1663/// embedded admin SPA. The Unix-domain-socket transport and the
1664/// `trusty-memory-mcp-bridge` binary were removed in PR3 of the #914
1665/// stdio-cutover epic; the canonical MCP integration is now
1666/// `trusty-memory serve --stdio` (PR1 #919).
1667/// What: axum router built from `web::router()` plus a `/sse` stub for the
1668/// existing MCP-over-SSE clients. Caller provides a pre-bound listener so
1669/// port auto-detection lives at the call site. Before accepting connections
1670/// the daemon stamps the bound `host:port` onto `AppState.bound_addr` and
1671/// writes `~/.trusty-memory/http_addr` so clients can discover the live port.
1672/// On shutdown the file is removed best-effort (a stale file with the wrong
1673/// port is worse than a missing one).
1674/// Test: `cargo test -p trusty-memory web::tests` exercises the router shape;
1675/// manual: `curl http://127.0.0.1:<port>/health` returns `ok` with `addr`.
1676#[cfg(feature = "axum-server")]
1677pub async fn run_http_on(state: AppState, listener: tokio::net::TcpListener) -> Result<()> {
1678    use axum::routing::get;
1679
1680    // Issue #35: recompute the `data_root` disk footprint every 10 s on a
1681    // background task so `GET /health` reports `disk_bytes` without doing a
1682    // recursive directory walk on the request path.
1683    spawn_disk_size_ticker(state.clone());
1684
1685    // Issue #228: emit aggregate `StatusChanged` on a fixed cadence rather
1686    // than on every drawer write. The previous design called
1687    // `aggregate_status_event` from every `memory_remember` / `memory_note`
1688    // / `memory_forget` (and the matching HTTP handlers), each of which
1689    // walked the data root + opened every palace handle. Coalescing the
1690    // emit to a 30 s ticker keeps dashboards live without dragging an
1691    // O(N palaces) recompute onto the write hot path.
1692    spawn_status_event_ticker(state.clone());
1693
1694    // Capture and advertise the bound address BEFORE serving so the first
1695    // request handler — and the http_addr discovery file — see the real port
1696    // even if `local_addr()` would otherwise be racy.
1697    let local = listener.local_addr().ok();
1698    let (written_path, written_dotfile_path) = if let Some(a) = local {
1699        // Stash on state for handlers (e.g. /health) to surface.
1700        let _ = state.bound_addr.set(a);
1701        info!("HTTP server listening on http://{a}");
1702        eprintln!("HTTP server listening on http://{a}");
1703        // Primary: write to the OS-standard data dir (`~/Library/Application
1704        // Support/trusty-memory/http_addr` on macOS, `~/.local/share/…` on
1705        // Linux). This is what `trusty_common::read_daemon_addr` reads.
1706        // Best-effort: a missing $HOME or read-only fs is non-fatal.
1707        let primary = match http_addr_path() {
1708            Some(p) => match write_http_addr_file(&p, &a) {
1709                Ok(()) => {
1710                    info!("wrote daemon address to {}", p.display());
1711                    Some(p)
1712                }
1713                Err(e) => {
1714                    tracing::warn!("could not write {}: {e}", p.display());
1715                    None
1716                }
1717            },
1718            None => {
1719                tracing::warn!("no $HOME — skipping http_addr discovery file");
1720                None
1721            }
1722        };
1723        // Issue #498: also write to `~/.trusty-memory/http_addr` so external
1724        // tools (e.g. claude-mpm's `migrate_trusty_autodetect`) that read the
1725        // dotfile path can discover the daemon's port. On macOS the OS-standard
1726        // path differs from the dotfile path; writing both ensures consumers
1727        // using either convention find the file. Best-effort: failures are
1728        // logged but do not block startup.
1729        let dotfile = match dotfile_http_addr_path() {
1730            Some(p) => match write_http_addr_file(&p, &a) {
1731                Ok(()) => {
1732                    info!("wrote daemon address to dotfile {}", p.display());
1733                    Some(p)
1734                }
1735                Err(e) => {
1736                    tracing::warn!("could not write dotfile {}: {e}", p.display());
1737                    None
1738                }
1739            },
1740            None => None,
1741        };
1742        (primary, dotfile)
1743    } else {
1744        (None, None)
1745    };
1746
1747    // Keep a handle to the BM25 supervisor (if any) so we can call
1748    // `shutdown()` on the exit path. Cloning here is cheap (`Arc`) and
1749    // detaches the lifetime of the supervisor from the `state` move into
1750    // the router below.
1751    let bm25_supervisor = state.bm25_supervisor.clone();
1752
1753    let app = web::router()
1754        .route("/sse", get(sse_handler))
1755        .with_state(state);
1756
1757    // Why (issue #534): bare axum::serve exits only on an internal error; SIGTERM
1758    // (launchctl bootout) would kill the process before the cleanup below had a
1759    // chance to run, leaving stale addr/socket files behind and dropping any
1760    // in-flight request without draining. `with_graceful_shutdown` installs a
1761    // SIGTERM + SIGINT watcher; when either fires axum stops accepting new
1762    // connections, drains active requests, then returns here so cleanup runs.
1763    let serve_result = axum::serve(listener, app)
1764        .with_graceful_shutdown(trusty_common::shutdown_signal())
1765        .await;
1766
1767    // Best-effort cleanup: remove `http_addr` files so stale clients fail fast
1768    // instead of timing out against a dead port. Remove both the OS-standard
1769    // path and the dotfile path (#498).
1770    if let Some(p) = written_path.as_ref() {
1771        let _ = std::fs::remove_file(p);
1772    }
1773    if let Some(p) = written_dotfile_path.as_ref() {
1774        let _ = std::fs::remove_file(p);
1775    }
1776
1777    // Issue #193: gracefully reap every spawned BM25 daemon before the
1778    // process exits so each one gets a chance to flush its snapshot and
1779    // unlink its socket. `kill_on_drop=true` on the children would
1780    // SIGKILL them on Drop anyway, but that skips the daemon's own
1781    // shutdown sequence and leaves stale sockets behind.
1782    if let Some(supervisor) = bm25_supervisor {
1783        supervisor.shutdown().await;
1784    }
1785
1786    serve_result?;
1787    Ok(())
1788}
1789
1790/// Convenience: bind `addr` and serve via [`run_http_on`].
1791#[cfg(feature = "axum-server")]
1792pub async fn run_http(state: AppState, addr: std::net::SocketAddr) -> Result<()> {
1793    let listener = tokio::net::TcpListener::bind(addr).await?;
1794    run_http_on(state, listener).await
1795}
1796
1797/// Convenience: bind dynamically (7070..=7079, OS fallback) and serve.
1798///
1799/// Why: `trusty-memory serve` with no `--http` flag is the canonical
1800/// launchd-managed daemon entry point. Dynamic binding lets a stale daemon
1801/// or a hand-spawned `serve --http 127.0.0.1:7070` coexist without breaking
1802/// the launchd-managed instance.
1803/// What: calls [`bind_dynamic_port`] then [`run_http_on`].
1804/// Test: integration via `trusty-memory serve` + `cat ~/.trusty-memory/http_addr`.
1805#[cfg(feature = "axum-server")]
1806pub async fn run_http_dynamic(state: AppState) -> Result<()> {
1807    let listener = bind_dynamic_port().await?;
1808    run_http_on(state, listener).await
1809}
1810
1811/// Spawn a background ticker that recomputes the `data_root` disk footprint
1812/// every 10 seconds and stores it in `state.disk_bytes` (issue #35).
1813///
1814/// Why: `GET /health` reports `disk_bytes`. Walking the data directory on
1815/// every health request would turn a frequent health poll into unbounded
1816/// recursive I/O. Computing it off the request path on a fixed cadence keeps
1817/// `/health` cheap and bounds the staleness to ~10 s — fine for an
1818/// at-a-glance footprint figure.
1819/// What: spawns a detached tokio task. `AppState` is cheap to `Clone` (all
1820/// `Arc` fields), so the task holds a full clone; the daemon process lives
1821/// for the lifetime of the server anyway, so no `Weak` downgrade is needed.
1822/// Each tick runs the blocking directory walk on `spawn_blocking` so it never
1823/// stalls the async runtime, then stores the byte total atomically.
1824/// Test: `health_endpoint_includes_resource_fields` asserts the field shape;
1825/// the ticker cadence is not unit-tested (timing-dependent).
1826#[cfg(feature = "axum-server")]
1827fn spawn_disk_size_ticker(state: AppState) {
1828    tokio::spawn(async move {
1829        let mut interval = tokio::time::interval(std::time::Duration::from_secs(10));
1830        loop {
1831            interval.tick().await;
1832            let dir = state.data_root.clone();
1833            // The directory walk is blocking filesystem I/O — run it on the
1834            // blocking pool so it never parks an async worker thread.
1835            let bytes = tokio::task::spawn_blocking(move || {
1836                trusty_common::sys_metrics::dir_size_bytes(&dir)
1837            })
1838            .await
1839            .unwrap_or(0);
1840            state
1841                .disk_bytes
1842                .store(bytes, std::sync::atomic::Ordering::Relaxed);
1843        }
1844    });
1845}
1846
1847/// Interval between aggregate-status snapshot emits on the SSE bus.
1848///
1849/// Why (issue #228): mutations used to fire `StatusChanged` synchronously on
1850/// the write path, which forced an O(N palaces) sum of drawer / vector / KG
1851/// counts on every `memory_remember`. Coalescing into a fixed-cadence ticker
1852/// lets dashboards stay current (a 30 s lag is invisible at human scale)
1853/// while keeping the write path free of aggregate work.
1854/// What: 30 seconds — short enough that the operator UI doesn't feel stale
1855/// between manual writes, long enough that the recompute cost (in-memory
1856/// registry walk plus the redb `count_active_triples` per palace) is a
1857/// rounding error on the daemon's CPU budget.
1858/// Test: covered indirectly — the math has not changed, only the cadence.
1859#[allow(dead_code)]
1860const STATUS_EVENT_TICK_SECS: u64 = 30;
1861
1862/// Spawn a background ticker that emits `DaemonEvent::StatusChanged` every
1863/// [`STATUS_EVENT_TICK_SECS`] seconds (issue #228).
1864///
1865/// Why: replaces the per-write `state.emit(self.aggregate_status_event())`
1866/// call sites that used to recompute the aggregate every time a drawer was
1867/// created or deleted. Walking N palaces on every write blocks the async
1868/// runtime; coalescing the emit onto a ticker keeps dashboards up-to-date
1869/// without that cost.
1870/// What: spawns a detached tokio task that holds a full `AppState` clone
1871/// (cheap — every field is `Arc`-backed) and ticks every
1872/// [`STATUS_EVENT_TICK_SECS`] seconds. Each tick computes
1873/// `MemoryService::aggregate_status_event` (which now iterates the
1874/// in-memory registry, not disk) and broadcasts it via `state.emit`. If
1875/// no SSE subscribers are connected the broadcast `send` is a cheap no-op,
1876/// so the ticker imposes no cost when nobody is listening.
1877/// Test: not unit-tested (timing-dependent fire-and-forget); the underlying
1878/// `aggregate_status_event` math is exercised by the existing
1879/// `status_endpoint_returns_payload` path.
1880#[allow(dead_code)]
1881fn spawn_status_event_ticker(state: AppState) {
1882    tokio::spawn(async move {
1883        let mut interval =
1884            tokio::time::interval(std::time::Duration::from_secs(STATUS_EVENT_TICK_SECS));
1885        // The first tick fires immediately, which is fine: it gives SSE
1886        // subscribers a baseline `StatusChanged` shortly after they connect.
1887        loop {
1888            interval.tick().await;
1889            let event = service::MemoryService::new(state.clone()).aggregate_status_event();
1890            state.emit(event);
1891        }
1892    });
1893}
1894
1895/// Live SSE event stream — pushes `DaemonEvent` frames to dashboard clients.
1896///
1897/// Why: The dashboard subscribes once and reacts to live pushes (palace
1898/// created, drawer added/deleted, dream completed, status changed) instead of
1899/// polling `/api/v1/*` endpoints.
1900/// What: Subscribes to `state.events`, emits an initial `connected` frame,
1901/// then forwards every `DaemonEvent` as `data: <json>\n\n`. Lagged
1902/// subscribers receive a `lag` frame indicating skipped events; channel
1903/// closure ends the stream.
1904/// Test: `web::tests::sse_stream_emits_palace_created` (covers subscribe +
1905/// emit + receive); manual: `curl -N http://.../sse`.
1906#[cfg(feature = "axum-server")]
1907pub(crate) async fn sse_handler(
1908    axum::extract::State(state): axum::extract::State<AppState>,
1909) -> impl axum::response::IntoResponse {
1910    use futures::StreamExt;
1911    use tokio_stream::wrappers::BroadcastStream;
1912
1913    let rx = state.events.subscribe();
1914    let initial = futures::stream::once(async {
1915        Ok::<axum::body::Bytes, std::io::Error>(axum::body::Bytes::from(
1916            "data: {\"type\":\"connected\"}\n\n",
1917        ))
1918    });
1919    let events = BroadcastStream::new(rx).map(|res| {
1920        let frame = match res {
1921            Ok(event) => match serde_json::to_string(&event) {
1922                Ok(json) => format!("data: {json}\n\n"),
1923                Err(e) => format!("data: {{\"type\":\"error\",\"message\":\"{e}\"}}\n\n"),
1924            },
1925            Err(tokio_stream::wrappers::errors::BroadcastStreamRecvError::Lagged(n)) => {
1926                format!("data: {{\"type\":\"lag\",\"skipped\":{n}}}\n\n")
1927            }
1928        };
1929        Ok::<axum::body::Bytes, std::io::Error>(axum::body::Bytes::from(frame))
1930    });
1931    let stream = initial.chain(events);
1932
1933    axum::response::Response::builder()
1934        .header("Content-Type", "text/event-stream")
1935        .header("Cache-Control", "no-cache")
1936        .header("X-Accel-Buffering", "no")
1937        .body(axum::body::Body::from_stream(stream))
1938        .expect("valid SSE response") // Why: invariant — SSE headers are compile-time constants; builder cannot fail
1939}
1940
1941#[cfg(test)]
1942mod lib_tests;