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 crate::session_store_cache::SessionStoreCache;
19use anyhow::Result;
20use serde_json::{json, Value};
21use std::net::SocketAddr;
22use std::path::PathBuf;
23use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
24use std::sync::{Arc, OnceLock};
25use tokio::sync::{broadcast, OnceCell, RwLock};
26use trusty_common::bm25_client::Bm25Client;
27use trusty_common::mcp::initialize_response;
28use trusty_common::memory_core::embed::Embedder;
29use trusty_common::memory_core::{store::ChatSessionStore, PalaceRegistry};
30use trusty_common::ChatProvider;
31
32/// Two-phase daemon readiness state (issues #910/#911, revised by #1970).
33///
34/// Why: The embedder cold-init (CoreML compile, 30-120 s) must never block
35/// the fast text/KG/BM25 paths that don't need it. Originally (#910/#911)
36/// this state gated a hard-error preflight that rejected every
37/// `memory_remember`/`memory_recall` call outright while `Warming` — mirrored
38/// from trusty-search's staged pipeline, #1970 replaced that with graceful
39/// degradation: writes persist immediately and defer embedding to a
40/// background task, reads return BM25 + L0/L1 results and simply omit the
41/// vector lane, all keyed off this same state.
42/// What: Two stable values stored atomically.  `Warming` (0) is the initial
43/// state; `Ready` (1) is set once the embedder has been successfully
44/// initialised by `spawn_startup_tasks`.  The transition is one-way and
45/// lock-free: a single `AtomicU8` compare-and-swap.
46/// Test: `daemon_readiness_transitions_warming_to_ready` in this module;
47///       degraded-path coverage in `tools::tests`
48///       (`remember_succeeds_and_defers_embedding_while_state_is_warming`,
49///       `recall_falls_back_to_bm25_and_l0_l1_while_warming`).
50#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51pub enum DaemonReadiness {
52    /// Embedder cold-init (and/or pin scan) still in progress.
53    Warming = 0,
54    /// Embedder initialised; all handlers may proceed normally.
55    Ready = 1,
56}
57
58impl DaemonReadiness {
59    /// Decode the raw atomic value.
60    ///
61    /// Why: centralises the `0 → Warming, else Ready` mapping so every
62    /// caller loads a meaningful enum rather than comparing raw integers.
63    /// What: returns `Warming` for `0`, `Ready` for any other value (only
64    /// `1` is ever written).
65    /// Test: `daemon_readiness_from_u8` in this module.
66    pub fn from_u8(v: u8) -> Self {
67        if v == 0 {
68            Self::Warming
69        } else {
70            Self::Ready
71        }
72    }
73}
74
75pub mod activity;
76pub mod attribution;
77pub mod authz;
78pub mod bm25_backfill;
79pub mod bm25_repair;
80pub mod bm25_supervisor;
81pub mod bootstrap;
82/// Autonomous Dreamer scheduler — spawns per-palace dream loops on daemon startup.
83///
84/// Why: issue #1529 — `Dreamer::start_with_shutdown()` was fully implemented
85/// but never called. This module wires it into the daemon so each palace gets
86/// a background dream loop that fires every 5 minutes of idle time.
87/// What: exports `spawn_dream_scheduler`, `make_shutdown_watch`, and
88/// `spawn_shutdown_bridge`. Disable with `TRUSTY_DREAM_DISABLED=1`.
89/// Test: see unit tests inside this module.
90pub mod dream_scheduler;
91/// File-descriptor usage and limit reporting for `/health`.
92///
93/// Why: expose `open_fds` / `fd_soft_limit` so operators can see the fd
94/// ceiling and current consumption without needing lsof or shell access.
95/// Test: `fd_metrics::tests::fd_metrics_returns_sane_values`.
96pub mod fd_metrics;
97/// Idle-to-disk eviction ticker: periodically drops cold palace handles to
98/// bound resident RSS. Wired in `spawn_startup_tasks` next to the dream
99/// scheduler; configured via `TRUSTY_MEMORY_IDLE_EVICT_SECS`.
100pub mod idle_evict;
101/// Live worker-occupancy gauge for the palace open/write path.
102///
103/// Why (issue #4001): doctor could only observe *process* liveness — an HTTP
104/// listener that answers and lock files that look clean — so a daemon with
105/// every worker parked in `concurrent_open::backoff_sleep_ms` still reported
106/// HEALTHY. This tracks how long the oldest in-flight operation has been
107/// running, which is the cheapest signal that actually distinguishes "the
108/// process is up" from "work is moving".
109/// Test: `worker_liveness::tests`.
110pub mod worker_liveness;
111// Why (issue #226): `chat` and `web` are pure axum HTTP/SSE handler
112//      surfaces. Gating them behind the `axum-server` feature is what lets
113//      library consumers (e.g. `trusty-agents` linking only `MemoryMcpService`)
114//      drop axum + tower-http entirely from their build graph.
115#[cfg(feature = "axum-server")]
116pub mod chat;
117pub mod commands;
118pub mod console_metrics;
119pub mod discovery;
120/// Daemon activity events (`DaemonEvent`, `HookType`, `InjectionKind`) plus
121/// the best-effort activity-log open fallback. Split from `lib.rs` to stay
122/// under the SLOC cap (#1195); the types are re-exported at the crate root so
123/// existing `trusty_memory::DaemonEvent` paths are unchanged.
124mod events;
125/// Supervised `serve --foreground` entry point (issue #787).
126///
127/// Why: launchd supervisors need loud failure on port collision, not silent
128/// port-walking to 7071+. Extracted to stay under the 500-line ratchet cap.
129/// What: exports `bind_foreground_port` (Fix C — abort on EADDRINUSE) and
130/// `run_http_foreground` (Fix A lock + Fix B http_addr + Fix C combined).
131/// Test: `foreground::tests::bind_foreground_port_refuses_collision`;
132/// `daemon_lock` module tests cover the lock-file logic.
133pub mod foreground;
134pub mod hook_emit;
135/// HTTP/SSE serving surface (`run_http*`, `bind_dynamic_port`, address-file
136/// helpers, `sse_handler`). Split from `lib.rs` to stay under the SLOC cap
137/// (#1195); its public entry points are re-exported at the crate root so
138/// existing `trusty_memory::run_http_on` paths are unchanged.
139mod http_server;
140pub mod kg_extract;
141pub mod mcp_service;
142pub mod messaging;
143pub mod openrpc;
144/// Issue #1217: default palace-ID derivation from project identity.
145///
146/// Why: the default palace ID should reflect the project's identity
147/// (git `owner/repo`, else `parent/dir`) rather than the bare directory
148/// basename, so the same repo resolves to the same palace across checkouts.
149/// What: exports the pure `derive_palace_id` core plus
150/// `owner_repo_from_git_remote`, `parent_dir_slug`, and the
151/// `TRUSTY_MEMORY_PALACE` env-override helpers.
152/// Test: see unit tests inside this module.
153pub mod palace_id_derive;
154/// Issue #88: project-root detection and palace-slug enforcement.
155///
156/// Why: prevents unbounded palace creation by anchoring palace names to the
157/// canonical slug of the project directory that contains the CWD, or to the
158/// `personal` sentinel for non-project contexts.
159/// What: exports `find_project_root`, `project_slug_at`, `project_slug`,
160/// `validate_palace_name`, `PERSONAL_PALACE`, and `PROJECT_MARKERS`.
161/// Test: see unit tests inside this module.
162pub mod project_root;
163pub mod prompt_facts;
164pub mod prompt_log;
165pub mod service;
166/// LRU-bounded cache of per-palace `ChatSessionStore` handles (issue #4639).
167///
168/// Why: the previous unbounded `DashMap` never released a `chat_sessions.redb`
169/// file descriptor, leaking one per palace for the daemon's lifetime.
170/// What: see [`session_store_cache::SessionStoreCache`].
171/// Test: `session_store_cache::tests`.
172pub mod session_store_cache;
173pub mod startup_scan;
174pub mod tools;
175pub mod transport;
176#[cfg(feature = "axum-server")]
177pub mod web;
178/// Issue #5399: WordNet-backed part-of-speech membership for KG extraction.
179///
180/// Why: the #4678 lexical stopword filter cannot separate two ordinary content
181/// words, which is what `exhaustiveness --is-a--> hard` needs.
182/// What: exports [`wordnet_pos::WordNetPos`] and the `NOUN`/`VERB`/`ADJ`/`ADV`
183/// bitmask constants. There is no global and no load step — see the module docs.
184/// Test: see unit tests inside this module.
185pub mod wordnet_pos;
186
187pub use activity::{ActivityEntry, ActivityFilter, ActivityLog, ActivitySource};
188pub use attribution::{CreatorInfo, CreatorSource};
189
190// Re-export the event types so the crate's public API is unchanged after the
191// #1195 split (`trusty_memory::DaemonEvent`, etc.). `open_activity_log_with_fallback`
192// stays crate-internal but is re-exported so `AppState::new` and `lib_tests`
193// reach it by its bare name via `super::*`.
194pub(crate) use events::open_activity_log_with_fallback;
195// #3434: test-only seam so `lib_tests` can force the tempdir-fallback path
196// via an explicit parameter instead of mutating the process-global `TMPDIR`.
197#[cfg(test)]
198pub(crate) use events::open_activity_log_with_fallback_in;
199pub use events::{DaemonEvent, HookType, InjectionKind};
200
201// Re-export the HTTP-serving surface so the crate's public API is unchanged
202// after the #1195 split (`trusty_memory::run_http_on`, etc.). The address-file
203// helpers stay crate-internal but are re-exported for `lib_tests`' `super::*`.
204pub use http_server::{
205    bind_dynamic_port, http_addr_path, is_data_dir_override_active, DEFAULT_HTTP_PORT,
206};
207#[cfg(feature = "axum-server")]
208pub use http_server::{run_http, run_http_dynamic, run_http_on};
209// These crate-internal HTTP helpers are consumed only by test modules
210// (`lib_tests`, `web::tests`), so gate the re-export on `cfg(test)` to avoid an
211// unused-import warning in the normal (non-test) build.
212#[cfg(all(test, feature = "axum-server"))]
213pub(crate) use http_server::{dotfile_http_addr_path, write_http_addr_file};
214
215/// Maximum bytes retained in the trigger-prompt excerpt embedded on a
216/// `HookFired` event.
217///
218/// Why: the full triggering prompt is sensitive and already lives in the
219/// JSONL prompt log; the activity feed only needs enough text to give an
220/// operator a glance — a single-line ~80 char preview matches the existing
221/// `drawer_content_preview` convention so dashboard rows render uniformly.
222/// What: 80 characters; longer prompts are truncated with a trailing `…`.
223/// Test: `hook_excerpt_truncates_long_prompts`.
224pub const HOOK_PROMPT_EXCERPT_CHARS: usize = 80;
225
226/// Reduce a triggering prompt to the short excerpt embedded on a
227/// `HookFired` activity event.
228///
229/// Why: see [`HOOK_PROMPT_EXCERPT_CHARS`]. Centralising the truncation rule
230/// keeps every emitter (HTTP, hook CLI handlers, future tests) producing
231/// the same preview shape so UI rendering is uniform.
232/// What: whitespace-collapses `prompt` and trims to
233/// [`HOOK_PROMPT_EXCERPT_CHARS`] chars with `…` when cut. Empty input
234/// returns an empty string.
235/// Test: `hook_excerpt_truncates_long_prompts`,
236/// `hook_excerpt_collapses_whitespace`.
237pub fn hook_prompt_excerpt(prompt: &str) -> String {
238    let normalised: String = prompt.split_whitespace().collect::<Vec<_>>().join(" ");
239    if normalised.chars().count() <= HOOK_PROMPT_EXCERPT_CHARS {
240        normalised
241    } else {
242        let kept: String = normalised
243            .chars()
244            .take(HOOK_PROMPT_EXCERPT_CHARS.saturating_sub(1))
245            .collect();
246        format!("{kept}…")
247    }
248}
249
250pub use mcp_service::MemoryMcpService;
251pub use tools::MemoryMcpServer;
252
253/// Resolve the directory that actually holds the per-palace subdirectories.
254///
255/// Why: there are two on-disk layouts in the wild. The current monorepo code
256/// treats the registry directory *itself* as the parent of per-palace dirs
257/// (`<dir>/<id>/palace.json`). The legacy standalone `trusty-memory` repo
258/// nested everything one level deeper under a `palaces/` subdirectory
259/// (`<data_dir>/palaces/<id>/palace.json`) — and that is where existing
260/// installs' data lives (e.g. 88 palaces under
261/// `~/Library/Application Support/trusty-memory/palaces/`). A daemon that uses
262/// the bare data dir as its registry root finds zero palaces because every
263/// `palace.json` sits one level below where it looked — the "palaces lost on
264/// restart" bug.
265/// What: given the standard data dir, returns `<data_dir>/palaces` when that
266/// subdirectory exists, otherwise `<data_dir>` itself. Resolving this once in
267/// `main.rs` and using the result as `AppState::data_root` keeps every call
268/// site (`status`, `palace_list`, `open_palace`, `palace_create`,
269/// `load_palaces_from_disk`) consistent without forcing a data migration.
270/// Test: `tests::resolve_palace_registry_dir_prefers_palaces_subdir` and
271/// `resolve_palace_registry_dir_falls_back_to_data_dir`.
272pub fn resolve_palace_registry_dir(data_dir: PathBuf) -> PathBuf {
273    // Issue #1939: the subdir-choice logic is hoisted into trusty-common
274    // (`palace_alias::palace_registry_dir_from`) so trusty-mpm's alias-registration
275    // path and this daemon path can never disagree on WHERE the registry (and the
276    // alias file beside it) lives. This delegates to keep a single implementation.
277    trusty_common::palace_alias::palace_registry_dir_from(data_dir)
278}
279
280/// Shared application state passed to every request handler.
281///
282/// Why: The stdio loop and HTTP server need the same handles to the registry,
283/// data root, and embedder so MCP tools can perform real reads/writes against
284/// the live trusty-memory core. The embedder is heavy (loads ONNX weights) so
285/// it is resolved lazily through the process-wide singleton on first use.
286///
287/// #4836: this struct used to carry its own `Arc<OnceCell<Arc<FastEmbedder>>>`,
288/// a SECOND cell independent of `retrieval::shared_embedder()`. Startup warmed
289/// the shared cell and latched `daemon_readiness` off it, while every recall
290/// consumed the private cell — so the flag reported one embedder's state while
291/// the request path used another's. One cell removes the disagreement (and the
292/// duplicate ~90 MB ONNX session) by construction.
293/// What: `Clone`-able via `Arc` fields. The registry / data root are eager; the
294/// embedder is reached via [`AppState::embedder`].
295/// Test: `app_state_default_constructs` confirms construction without panic.
296#[derive(Clone)]
297pub struct AppState {
298    pub version: String,
299    pub registry: Arc<PalaceRegistry>,
300    pub data_root: PathBuf,
301    /// Optional default palace applied to MCP tool calls when the caller
302    /// omits the `palace` argument. Set via `trusty-memory serve --palace`.
303    pub default_palace: Option<String>,
304    /// Active chat provider selected at startup. `None` means no upstream is
305    /// configured (no Ollama detected and no OpenRouter key) — callers must
306    /// degrade gracefully (chat endpoint returns 412).
307    pub chat_provider: Arc<OnceCell<Option<Arc<dyn ChatProvider>>>>,
308    /// Per-palace chat-session stores, opened lazily so cold-start cost is
309    /// paid only when chat-history endpoints are hit.
310    ///
311    /// #4639: was an unbounded `DashMap` with no `remove`/TTL/cap, leaking one
312    /// `chat_sessions.redb` fd per palace for the daemon's lifetime (844
313    /// measured live, all pointing at already-unlinked files). Now an
314    /// LRU-bounded cache that evicts cold, unused stores.
315    pub session_stores: Arc<SessionStoreCache>,
316    /// Broadcast sender for live `DaemonEvent` pushes to SSE subscribers.
317    ///
318    /// Why: Lets mutating handlers emit events that any connected dashboard
319    /// receives instantly. Cap of 128 buffers transient slow readers; if a
320    /// receiver lags it gets `RecvError::Lagged` and we emit a `lag` frame.
321    pub events: Arc<broadcast::Sender<DaemonEvent>>,
322    /// Instant the daemon started, used to compute `uptime_secs` on `/health`.
323    ///
324    /// Why (issue #35): `GET /health` reports how long the daemon has been
325    /// up. Capturing a monotonic `Instant` at `AppState` construction lets the
326    /// handler compute the elapsed seconds cheaply and without a clock-skew
327    /// hazard.
328    /// What: a wall-monotonic `Instant`; `AppState::new` stamps it at startup.
329    /// Test: `health_endpoint_includes_resource_fields`.
330    pub started_at: std::time::Instant,
331    /// In-memory ring buffer of recent tracing log lines (issue #35).
332    ///
333    /// Why: the `GET /api/v1/logs/tail` endpoint serves the last N log lines
334    /// so operators can inspect a running daemon without tailing a file. The
335    /// buffer is shared between the tracing `LogBufferLayer` (writer) and the
336    /// HTTP handler (reader).
337    /// What: a cheap `Arc`-backed clone of the buffer the subscriber writes
338    /// to. Defaults to an empty buffer for states that never install the
339    /// layer (tests, the stdio path).
340    /// Test: `logs_tail_returns_recent_lines`.
341    pub log_buffer: trusty_common::log_buffer::LogBuffer,
342    /// Bug-capture ERROR store (bug-reporting #478, Phase 1).
343    ///
344    /// Why: Phase 2 MCP / HTTP endpoints need to query captured errors; stashing
345    ///      the `ErrorStore` handle here lets any handler reach it cheaply without
346    ///      a second global or per-request construction.
347    /// What: populated by `run_serve` from the `init_tracing_with_buffer_and_capture`
348    ///      result; the layer writes to this store automatically so every
349    ///      `tracing::error!` call site contributes without any changes to call
350    ///      sites. `None` in states that do not install the layer (tests, the
351    ///      stdio path).
352    /// Test: compile-presence is verified by the `trusty-memory` build; Phase 2
353    ///      will add query tests in `web.rs`.
354    pub error_store: Option<trusty_common::error_capture::ErrorStore>,
355    /// Minimal multi-tenant authorization seam (issue #1714). `false`
356    /// (single-tenant, the default) preserves today's behaviour — every
357    /// existing caller of `palace_create force=true` keeps working. `true`
358    /// opts into `authz::authorize_force_palace_create` failing closed on
359    /// every `force=true` request until a real capability check lands. Set
360    /// via `with_multi_tenant_mode_from_env` (`TRUSTY_MEMORY_MULTI_TENANT=1`);
361    /// see the `authz` module docs for the full design rationale.
362    pub multi_tenant_mode: bool,
363    /// Most recent on-disk footprint of `data_root`, in bytes (issue #35).
364    ///
365    /// Why: `GET /health` reports `disk_bytes`. Walking the data directory on
366    /// every health request would make a frequent health poll do unbounded
367    /// I/O; a background task recomputes it every 10 s and stores it here so
368    /// the handler reads it lock-free.
369    /// What: an `AtomicU64` updated by the ticker spawned in `run_http_on`.
370    /// `0` until the first walk completes.
371    /// Test: `health_endpoint_includes_resource_fields`.
372    pub disk_bytes: Arc<std::sync::atomic::AtomicU64>,
373    /// Per-process RSS + CPU sampler, refreshed on each `/health` request
374    /// (issue #35).
375    ///
376    /// Why: CPU usage is a delta between two `sysinfo` refreshes, so the
377    /// sampler must persist between requests — hence the shared `Mutex`.
378    /// What: a `tokio::sync::Mutex<SysMetrics>` so the async health handler
379    /// can sample without blocking the runtime.
380    /// Test: `health_endpoint_includes_resource_fields`.
381    pub sys_metrics: Arc<tokio::sync::Mutex<trusty_common::sys_metrics::SysMetrics>>,
382    /// HTTP listener address the daemon bound to, once `run_http_on` is running.
383    ///
384    /// Why: clients (and `/health` responses) need to advertise the live
385    /// `host:port` even though port selection happens dynamically (7070–7079
386    /// walk + OS fallback). Stashing it on `AppState` lets request handlers
387    /// surface the discovery value without re-querying the listener.
388    /// What: a `OnceLock<SocketAddr>` so `run_http_on` writes it exactly once
389    /// at bind time and every handler reads it lock-free thereafter. Empty
390    /// (`None` from `get()`) on the stdio path where no listener exists.
391    /// Test: `health_endpoint_reports_bound_addr` (added below).
392    pub bound_addr: Arc<OnceLock<SocketAddr>>,
393    /// Cached prompt-facts surface served by the MCP `get_prompt_context`
394    /// tool (issue #42).
395    ///
396    /// Why: The original session-init `prompts/get` design loaded context
397    /// once per connection; switching to a per-message tool lets the model
398    /// pull fresh, query-filtered context on demand. The cache holds both
399    /// the raw triples (for filtered lookups) and a pre-formatted Markdown
400    /// block (for the unfiltered hot path) so neither code path re-walks
401    /// the KG. The cache is rebuilt by
402    /// `prompt_facts::rebuild_prompt_cache` after any write that touches a
403    /// hot predicate (`kg_assert`, `add_alias`, `remove_prompt_fact`).
404    /// What: An `Arc<tokio::sync::RwLock<PromptFactsCache>>` so the hot
405    /// read path takes a brief read lock and clones the cache; rebuilds
406    /// take a write lock for the assignment only. The async-aware lock
407    /// (issue #229) yields to the tokio runtime instead of blocking a
408    /// runtime thread for the rebuild duration. An empty `triples` vec ↔
409    /// "no context stored yet" (the tool handler renders a hint).
410    /// Test: `get_prompt_context_returns_cached_or_hint`,
411    /// `get_prompt_context_filters_by_query`.
412    pub prompt_context_cache: Arc<RwLock<prompt_facts::PromptFactsCache>>,
413    /// Serializes the Tier S check-then-write sequence (#4888).
414    ///
415    /// Why: the 20-fact cap is only a cap if it cannot be raced past.
416    /// `check_tier_s_admission` counts active facts across every palace and
417    /// then the caller writes — two callers that both observe 19 would both
418    /// pass and the surface would land at 21. Nothing else serializes them:
419    /// the KG's single-writer actor orders writes only *within* one palace,
420    /// and the count spans all of them. "Usually 20, occasionally 21" is not
421    /// the invariant ADR-0028 D8 asks for.
422    /// What: an async mutex whose guard is acquired by
423    /// `check_tier_s_admission` and returned to the caller, which must hold it
424    /// until its `kg.assert` is enqueued. Cold predicates never acquire it, so
425    /// ordinary knowledge-graph writes stay fully concurrent; hot writes are
426    /// deliberate and rare, so serializing them costs nothing measurable.
427    /// Test: `tier_s_cap_holds_under_concurrent_writes`.
428    pub tier_s_admission_lock: Arc<tokio::sync::Mutex<()>>,
429    /// Persistent activity log (issue #96).
430    ///
431    /// Why: the dashboard activity feed used to be a pure live-stream over
432    /// `/sse` — opening the UI showed an empty feed and any mutation from
433    /// the MCP path was invisible. Holding an `ActivityLog` on `AppState`
434    /// lets `emit` record an entry on every push so the
435    /// `GET /api/v1/activity` handler can return historical rows on mount
436    /// and the live SSE stream can continue prepending events on top of
437    /// the loaded history. `None` on builds that opt out (tests that use
438    /// `AppState::new` get a real log under their tempdir so behaviour
439    /// matches production).
440    /// What: an `Arc<ActivityLog>` shared with every emitter.
441    /// Test: `web::tests::activity_endpoint_lists_recent_emits`.
442    pub activity_log: Arc<ActivityLog>,
443    /// Optional per-palace BM25 lexical search lane (issue #156).
444    ///
445    /// Why: in-process BM25 would serialise the recall hot path on disk
446    /// I/O during writes and contend with the redb/usearch locks. Delegating
447    /// to the `trusty-bm25-daemon` subprocess (one socket per palace) keeps
448    /// BM25 ingestion and search off the critical path while still feeding
449    /// hits into the recall RRF fusion.
450    /// What: `Some(client)` only when `TRUSTY_BM25_DAEMON=1` at startup —
451    /// every code path that uses this field is gated on `is_some()` and
452    /// falls back to vector-only behaviour otherwise so existing deployments
453    /// see zero behavioural change.
454    /// Test: `bm25_client_disabled_by_default`,
455    /// `bm25_client_enabled_when_env_set`.
456    pub bm25_client: Option<Arc<Bm25Client>>,
457    /// Optional per-palace BM25 daemon spawn supervisor (issue #193).
458    ///
459    /// Why: without an in-process supervisor the BM25 daemon must be
460    /// launched out-of-band (launchd, manual `trusty-bm25-daemon`), which
461    /// is the same UX trap PR #190 fixed for trusty-embedderd. Holding a
462    /// supervisor here lets us spawn the daemon on first BM25 use for a
463    /// palace, restart it if it dies, and reap it on clean shutdown.
464    /// `Some` only when `TRUSTY_BM25_DAEMON=1` at startup — the same gate
465    /// that enables `bm25_client`. When set but `TRUSTY_BM25_EXTERNAL=1`,
466    /// the supervisor's `ensure_running` becomes a no-op that just returns
467    /// the canonical socket path so operators can keep using their own
468    /// process manager.
469    /// Test: covered by `bm25_supervisor_present_when_env_set` and the
470    /// `bm25_supervisor::tests` unit tests.
471    pub bm25_supervisor: Option<Arc<bm25_supervisor::Bm25Supervisor>>,
472    /// Per-palace write serialisation locks (issue #230).
473    ///
474    /// Why: the dedup gate in `tools.rs` previously read a snapshot of
475    /// existing drawers, checked for near-duplicates via Jaro-Winkler, and
476    /// then issued the write — a classic time-of-check/time-of-use race.
477    /// Two concurrent `memory_remember` calls with the same content could
478    /// both see the pre-write snapshot, both pass the gate, and both land
479    /// duplicate drawers. Serialising the gate-then-write sequence per
480    /// palace closes the window: while one task holds the mutex, any
481    /// concurrent writer for the same palace blocks until the first write
482    /// finishes and is visible to `list_drawers`. The lock is **per
483    /// palace** (not global) so writes to different palaces continue to
484    /// run in parallel.
485    /// What: a `DashMap` keyed by palace id, where each entry is an
486    /// `Arc<tokio::sync::Mutex<()>>`. The mutex is constructed lazily by
487    /// `palace_write_lock` on first access. `Arc` lets callers hold a
488    /// clone of the lock past the lifetime of the `DashMap` entry so the
489    /// map never needs to be held across an `.await`.
490    /// Test: `tools::tests::dedup_gate_blocks_concurrent_duplicate_writes`.
491    pub palace_write_locks: Arc<dashmap::DashMap<String, Arc<tokio::sync::Mutex<()>>>>,
492    /// Counter of in-flight activity-log writes spawned by `emit`
493    /// (issue #232).
494    ///
495    /// Why: `emit` offloads the synchronous redb append to the tokio blocking
496    /// pool via `spawn_blocking` so the async runtime is never parked waiting
497    /// on fsync. The write is fire-and-forget — `emit` returns immediately
498    /// after spawning. Tests that observe the activity log right after a
499    /// burst of `emit` calls need a deterministic synchronization point;
500    /// holding an in-flight counter lets `flush_activity_writes` poll until
501    /// every spawned append has settled, which keeps the assertions
502    /// race-free without forcing every caller to `.await`.
503    /// What: an `Arc<AtomicUsize>` incremented before each `spawn_blocking`
504    /// and decremented inside the closure (after the append completes, even
505    /// if it errored). The counter is cheap (one atomic add per emit) and
506    /// stays at zero in steady-state production traffic.
507    /// Test: `web::tests::activity_endpoint_lists_recent_emits` and
508    /// `tests::emit_persists_mutations_but_skips_status_changed` call
509    /// `flush_activity_writes` to drain the counter before reading the log.
510    pub pending_activity_writes: Arc<AtomicUsize>,
511    /// Live occupancy gauge for the palace open path (issue #4001).
512    ///
513    /// Why: during the #3992 incident six daemon threads sat parked in
514    /// `concurrent_open::backoff_sleep_ms` with a `memory_remember` hung
515    /// ~1800 s, while both doctors reported HEALTHY. Every existing signal —
516    /// HTTP liveness, fastembed cache state, lock-file staleness — describes
517    /// the *process*, not the *work*. This is the one field that can answer
518    /// "is anything actually moving?", and `/health` surfaces it so an
519    /// out-of-process doctor can report what the daemon actually observed
520    /// instead of inferring health from a cheap proxy.
521    /// What: an [`worker_liveness::WorkerLiveness`] slot table; one CAS on
522    /// entry and one store on exit per tracked operation, so the gauge cannot
523    /// itself become the load problem it exists to detect.
524    /// Test: `web::tests::health_tests::health_reports_wedged_worker_pool`.
525    pub worker_liveness: Arc<worker_liveness::WorkerLiveness>,
526    /// How long an operation may run before the pool is called wedged.
527    ///
528    /// Why this is state rather than an env read on the request path: `/health`
529    /// is polled once a second, so re-reading (and re-parsing) an environment
530    /// variable per request is needless work; resolving it once at construction
531    /// also makes the value a property of the daemon instead of a process-wide
532    /// global, which is what lets tests drive the wedge condition
533    /// deterministically instead of mutating shared env state and racing each
534    /// other.
535    /// What: defaults to [`worker_liveness::wedge_threshold`].
536    /// Test: `web::tests::health_tests::health_reports_wedged_worker_pool`.
537    pub wedge_threshold: std::time::Duration,
538    /// In-memory cache mapping palace id → `Palace.name` (issue #228).
539    ///
540    /// Why: every `memory_remember` / `memory_note` write used to call
541    /// `PalaceRegistry::list_palaces` (a synchronous filesystem walk of the
542    /// data root) just to resolve a friendly palace name for the SSE
543    /// `DrawerAdded` event. With N palaces on disk the cost was O(N) opendirs
544    /// plus `palace.json` reads on every write, blocking the async runtime.
545    /// Caching the name in-memory turns the lookup into a `DashMap::get`.
546    /// What: `DashMap<String, String>` populated by `create_palace` and
547    /// `load_palaces_from_disk`, kept in sync by rename / delete paths.
548    /// Missing entries are treated as "name unknown" so callers fall back to
549    /// the palace id and the emit path never fails.
550    /// Test: `palace_name_cache_populated_after_hydration` and
551    /// `palace_name_cache_updates_on_create`.
552    pub palace_names: Arc<dashmap::DashMap<String, String>>,
553    /// Single-pass startup pin-file map: palace id → project root path (issue #470).
554    ///
555    /// Why: after daemon startup we have no record of which on-disk project
556    /// directories correspond to which palace ids — that information only
557    /// existed inside the pin files on disk. Eager-opening every palace on
558    /// startup is too expensive. This field captures the scan-only result of
559    /// `startup_scan::scan_pin_map` so handlers that want to locate a project
560    /// by its palace id (e.g. future cwd-inference, project-health checks)
561    /// can do a single `DashMap::get` instead of a filesystem walk.
562    /// Populated once, shortly after `load_palaces_from_disk` returns, by
563    /// `spawn_startup_tasks`. Never mutated after population — it is a
564    /// snapshot of what the filesystem looked like at startup.
565    /// What: `DashMap<String (palace_id), PathBuf (project root)>`.
566    /// The outer `Arc` lets `spawn_startup_tasks` (which holds only a clone
567    /// of `AppState`) write to the same backing map that request handlers
568    /// read. Population is asynchronous so callers must treat an absent entry
569    /// as "not yet scanned" (or "no pin found"), never as "palace unknown".
570    /// Test: `startup_scan::tests::scan_pin_map_*` validate the underlying
571    /// scanner function; the wiring in `spawn_startup_tasks` is covered by
572    /// the integration-test daemon start path.
573    pub pin_project_map: Arc<dashmap::DashMap<String, PathBuf>>,
574    /// Bounded sender for the BM25 index worker (issue #231).
575    ///
576    /// Why: the previous fire-and-forget design `tokio::spawn`ed one task per
577    /// `memory_remember` / `memory_note` call, so a write burst against a slow
578    /// or unreachable BM25 daemon grew an unbounded in-flight task queue. A
579    /// single long-lived worker draining a bounded mpsc channel caps that
580    /// back-pressure: writers `try_send` (never block), full-queue requests
581    /// are dropped with a `warn!`, and the worker exits cleanly when the last
582    /// sender is dropped on shutdown.
583    /// What: an `mpsc::Sender` cloned to every `AppState` clone (cheap). The
584    /// matching receiver is consumed by the worker spawned in
585    /// [`AppState::new`] via [`tools::spawn_bm25_index_worker`]. Capacity is
586    /// [`tools::BM25_INDEX_QUEUE_CAPACITY`] (256).
587    /// Test: `bm25_index_queue_drops_when_full` exercises the full-queue
588    /// branch via `bm25_index_enqueue`.
589    pub bm25_index_tx: tokio::sync::mpsc::Sender<tools::Bm25IndexRequest>,
590    /// Palaces whose BM25 coverage is known to be incomplete (#5048 review).
591    ///
592    /// Why: `bm25_index_enqueue` drops on a full queue so `memory_remember`
593    /// never waits on daemon RTT. That trade is only defensible if a drop is
594    /// actually repaired, and before this field the sole production trigger for
595    /// a backfill was daemon startup — so a drop stayed invisible until the
596    /// next restart. Every observer of lost coverage marks the palace here and
597    /// [`bm25_repair::spawn_repair_sweep`] consumes it on an interval.
598    /// What: a `DashSet` of palace ids, shared by every `AppState` clone.
599    /// Idempotent — forty drops for one palace queue one repair.
600    /// Test: `bm25_repair_tests.rs`, `bm25_index_queue_drops_when_full`.
601    pub bm25_dirty: bm25_repair::DirtyPalaces,
602    /// Cached result of the startup update check (issue #537).
603    ///
604    /// Why: `/health` should report `update_available` without hitting crates.io
605    /// on every probe. A single background check at daemon startup stores the
606    /// result here; the health handler reads it lock-free (well, a brief mutex
607    /// lock) without a network call.
608    /// What: `None` = up-to-date or check not yet done; `Some("x.y.z")` = newer
609    /// version available. The field is populated by a `tokio::spawn` in
610    /// `spawn_startup_tasks` (main.rs) after the daemon binds.
611    /// Test: indirectly by the `/health` endpoint tests in `web.rs`.
612    pub update_available: Arc<std::sync::Mutex<Option<String>>>,
613    /// Two-phase readiness state — `Warming` until the embedder is initialised,
614    /// then `Ready` (issues #910 / #911).
615    ///
616    /// Why: `AppState::embedder()` used to call `FastEmbedder::new()` without
617    /// any timeout, so the first `memory_recall`/`memory_remember` that arrived
618    /// before CoreML finished compiling would block for 5–11 hours until the
619    /// OnceCell resolved (issue #910). Exposing this state lets the preflight
620    /// guards in `tools.rs` return an explicit fast error immediately —
621    /// `"trusty-memory is warming up, retry shortly"` — instead of queueing
622    /// behind an open-ended init.
623    /// What: An `AtomicU8` starting at `DaemonReadiness::Warming` (0) and flipped
624    /// to `DaemonReadiness::Ready` (1) by `spawn_startup_tasks` after the embedder
625    /// warm-up succeeds.  The transition is one-way and lock-free.
626    /// Test: `daemon_readiness_transitions_warming_to_ready`.
627    pub daemon_readiness: Arc<AtomicU8>,
628}
629
630impl AppState {
631    /// Construct an `AppState` rooted at the given on-disk data directory.
632    ///
633    /// Why: The CLI (`serve`) and integration tests need to point the MCP
634    /// server at different roots — production at `dirs::data_dir`, tests at a
635    /// `tempfile::tempdir()`.
636    /// What: Builds an empty `PalaceRegistry`, captures the version, and
637    /// allocates an empty `OnceCell` for the embedder. `default_palace` is
638    /// `None`; use `with_default_palace` to set it.
639    /// Test: `tools::tests::dispatch_palace_create_persists` constructs an
640    /// AppState pointed at a tempdir and round-trips a palace through it.
641    pub fn new(data_root: PathBuf) -> Self {
642        let (events_tx, _) = broadcast::channel::<DaemonEvent>(128);
643        // Issue #96: open (or create) the persistent activity log under the
644        // daemon data root. Open failure is logged but never crashes the
645        // daemon — we fall back to a per-process tempdir so emits remain
646        // best-effort and the rest of the daemon keeps working.
647        let activity_log = open_activity_log_with_fallback(&data_root);
648        // Issue #231: bounded mpsc channel + single long-lived worker
649        // replaces the per-write `tokio::spawn` fire-and-forget pattern so
650        // BM25 indexing back-pressure is capped. The worker is spawned here
651        // unconditionally so the channel always has a drain — even when
652        // `bm25_client` is `None`, the worker just consumes and discards
653        // each request so senders never block on a full queue.
654        let (bm25_index_tx, bm25_index_rx) =
655            tokio::sync::mpsc::channel::<tools::Bm25IndexRequest>(tools::BM25_INDEX_QUEUE_CAPACITY);
656        // `bm25_client` / `bm25_supervisor` start as `None`; the builder
657        // `with_bm25_client_from_env` rebuilds the worker with the real
658        // client + supervisor once env-gated opt-in is resolved.
659        let bm25_dirty: bm25_repair::DirtyPalaces = Arc::new(dashmap::DashSet::new());
660        tools::spawn_bm25_index_worker(bm25_index_rx, None, None, Arc::clone(&bm25_dirty));
661        Self {
662            version: env!("CARGO_PKG_VERSION").to_string(),
663            // Idle-to-disk: honour TRUSTY_MEMORY_MAX_OPEN_PALACES (default 64)
664            // so operators can bound resident-palace RAM without a rebuild.
665            registry: Arc::new(PalaceRegistry::from_env()),
666            data_root,
667            default_palace: None,
668            chat_provider: Arc::new(OnceCell::new()),
669            // #4639: bounded LRU (TRUSTY_MEMORY_MAX_OPEN_SESSION_STORES,
670            // default 32) so chat_sessions.redb handles stop accumulating.
671            session_stores: Arc::new(SessionStoreCache::from_env()),
672            events: Arc::new(events_tx),
673            started_at: std::time::Instant::now(),
674            // Default to an empty buffer — `with_log_buffer` overrides this
675            // when the daemon installs the `LogBufferLayer` (HTTP mode).
676            log_buffer: trusty_common::log_buffer::LogBuffer::new(
677                trusty_common::log_buffer::DEFAULT_LOG_CAPACITY,
678            ),
679            // Bug-reporting #478: `None` until `with_error_store` is called
680            // during daemon startup (HTTP mode). Tests keep `None` so no
681            // unexpected files are written to the OS data dir.
682            error_store: None,
683            multi_tenant_mode: false,
684            disk_bytes: Arc::new(std::sync::atomic::AtomicU64::new(0)),
685            sys_metrics: Arc::new(tokio::sync::Mutex::new(
686                trusty_common::sys_metrics::SysMetrics::new(),
687            )),
688            bound_addr: Arc::new(OnceLock::new()),
689            prompt_context_cache: Arc::new(RwLock::new(prompt_facts::PromptFactsCache::default())),
690            tier_s_admission_lock: Arc::new(tokio::sync::Mutex::new(())),
691            activity_log,
692            bm25_client: None,
693            bm25_supervisor: None,
694            palace_write_locks: Arc::new(dashmap::DashMap::new()),
695            pending_activity_writes: Arc::new(AtomicUsize::new(0)),
696            worker_liveness: Arc::new(worker_liveness::WorkerLiveness::new()),
697            wedge_threshold: worker_liveness::wedge_threshold(),
698            palace_names: Arc::new(dashmap::DashMap::new()),
699            pin_project_map: Arc::new(dashmap::DashMap::new()),
700            bm25_index_tx,
701            bm25_dirty,
702            update_available: Arc::new(std::sync::Mutex::new(None)),
703            // Start in Warming state; flipped to Ready by spawn_startup_tasks
704            // once the embedder warm-up succeeds (issues #910/#911).
705            daemon_readiness: Arc::new(AtomicU8::new(DaemonReadiness::Warming as u8)),
706        }
707    }
708
709    /// Acquire (lazily, then clone) the per-palace write mutex.
710    ///
711    /// Why (issue #230): the dedup-check + `remember_with_options` write
712    /// sequence in `tools.rs` must be atomic per palace to prevent two
713    /// concurrent identical writes from both passing the dedup gate.
714    /// Callers hold the returned `Arc<Mutex<()>>`'s guard across the gate
715    /// check and the write so the second writer blocks until the first
716    /// write is visible to `list_drawers`. Returning a clone of the `Arc`
717    /// rather than a borrow into the `DashMap` lets the caller `.await`
718    /// while holding the lock without risking a deadlock against any
719    /// future map mutation (DashMap shards are sync mutexes).
720    /// What: looks up the palace id in `palace_write_locks` and returns
721    /// a clone of the existing mutex; on the first call for a palace,
722    /// inserts a freshly-constructed `tokio::sync::Mutex<()>` first. The
723    /// `DashMap::entry().or_insert_with` API guarantees the lazy
724    /// construction is racy-safe — only one mutex is ever inserted per
725    /// palace id.
726    /// Test: `tools::tests::dedup_gate_blocks_concurrent_duplicate_writes`.
727    pub fn palace_write_lock(&self, palace_id: &str) -> Arc<tokio::sync::Mutex<()>> {
728        if let Some(existing) = self.palace_write_locks.get(palace_id) {
729            return existing.clone();
730        }
731        self.palace_write_locks
732            .entry(palace_id.to_string())
733            .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
734            .clone()
735    }
736
737    /// Look up a project root path by palace id in the startup pin-scan map.
738    ///
739    /// Why: provides a stable, cheap accessor so handlers do not reach directly
740    /// into the `DashMap` field and so the accessor can be mocked in future
741    /// tests without touching `AppState` internals. The map is populated
742    /// asynchronously by `spawn_startup_tasks` — an absent entry means either
743    /// the scan has not completed yet or no pin file claimed that id.
744    /// What: returns `Some(project_path)` when the palace id was found during
745    /// startup scan; `None` otherwise.
746    /// Test: covered indirectly via the startup-scan integration path; the
747    /// underlying map data is validated by `startup_scan::tests`.
748    pub fn pinned_project_path(&self, palace_id: &str) -> Option<PathBuf> {
749        self.pin_project_map.get(palace_id).map(|e| e.clone())
750    }
751
752    /// Builder-style: opt-in to the BM25 lexical lane (issue #156).
753    ///
754    /// Why: the BM25 subprocess is gated behind `TRUSTY_BM25_DAEMON=1` so
755    /// the default `cargo install trusty-memory` / launchd plist deployment
756    /// stays vector-only and existing test fixtures keep passing without
757    /// having to provision a daemon. Reading the env var here keeps the
758    /// gating logic in one place (the helper in `main.rs` just plumbs the
759    /// result through).
760    /// What: when `TRUSTY_BM25_DAEMON=1`, constructs one `Bm25Client` per
761    /// palace by lazy-resolving the socket path the first time the palace
762    /// id is observed. Currently we install a shared `default` client up
763    /// front and re-key on the palace id at the call site — palaces with no
764    /// daemon socket simply see search/index errors which we log + ignore.
765    /// Returns `self` unchanged when the env var is unset or set to anything
766    /// other than `1`.
767    /// Test: `bm25_client_disabled_by_default`,
768    /// `bm25_client_enabled_when_env_set`.
769    #[must_use]
770    pub fn with_bm25_client_from_env(mut self) -> Self {
771        if std::env::var("TRUSTY_BM25_DAEMON").as_deref() == Ok("1") {
772            // Install the default-palace client; per-palace clients are
773            // constructed on demand via `Bm25Client::for_palace`.
774            let default_palace = self.default_palace.as_deref().unwrap_or("default");
775            self.bm25_client = Some(Arc::new(Bm25Client::for_palace(default_palace)));
776            // Issue #193: hand-in-hand with the client, attach a spawn
777            // supervisor so the BM25 daemon is auto-started on first use
778            // for any palace. Operators who want to manage daemons
779            // out-of-band (launchd, systemd, manual) set
780            // TRUSTY_BM25_EXTERNAL=1 which makes the supervisor a no-op.
781            self.bm25_supervisor = Some(Arc::new(bm25_supervisor::Bm25Supervisor::new()));
782            // Issue #231: rebuild the bounded indexer channel + worker so
783            // the worker holds the now-populated client + supervisor. The
784            // placeholder worker installed by `AppState::new` (with `None`
785            // / `None`) drained the channel into the void — replacing the
786            // sender here closes the placeholder receiver and the
787            // placeholder worker exits cleanly. The new worker takes over
788            // as the sole drain for the indexer queue.
789            let (tx, rx) = tokio::sync::mpsc::channel::<tools::Bm25IndexRequest>(
790                tools::BM25_INDEX_QUEUE_CAPACITY,
791            );
792            tools::spawn_bm25_index_worker(
793                rx,
794                self.bm25_client.clone(),
795                self.bm25_supervisor.clone(),
796                Arc::clone(&self.bm25_dirty),
797            );
798            self.bm25_index_tx = tx;
799            tracing::info!(
800                palace = default_palace,
801                "BM25 daemon client + spawn supervisor enabled (TRUSTY_BM25_DAEMON=1)"
802            );
803        }
804        self
805    }
806
807    /// Scan the palace registry directory and re-register every persisted
808    /// palace into the in-memory [`PalaceRegistry`].
809    ///
810    /// Why: `AppState::new` builds an *empty* registry, so after a daemon
811    /// restart `palace_list` / the dashboard reported zero palaces even though
812    /// dozens existed on disk — palace metadata was persisted by
813    /// `palace_create` but never re-hydrated on startup. This method closes
814    /// that gap by walking the on-disk layout (each subdirectory holding a
815    /// `palace.json` is one palace) and rebuilding a live `PalaceHandle` for
816    /// each, so recall paths see the full set immediately after a restart.
817    /// What: runs the blocking filesystem walk + per-palace `PalaceHandle::open`
818    /// on a `spawn_blocking` thread (so it never stalls the async runtime),
819    /// registers each successfully opened palace via `register_arc`, logs every
820    /// load at `debug!`, and returns the count loaded. A palace that fails to
821    /// open (corrupt index, unreadable `kg.db`, etc.) is logged at `warn!` and
822    /// skipped — one bad palace must not abort startup or crash the daemon.
823    /// `data_root` is expected to already be the palace registry directory —
824    /// `main.rs` resolves it via [`resolve_palace_registry_dir`] before
825    /// constructing the `AppState`, so the flat / legacy-`palaces/` layout
826    /// difference is handled exactly once.
827    /// Test: `tests::load_palaces_from_disk_rehydrates_registry` writes two
828    /// palaces into a tempdir, constructs an `AppState`, calls this method, and
829    /// asserts the returned count and registry contents.
830    pub async fn load_palaces_from_disk(&self) -> Result<usize> {
831        let registry_dir = self.data_root.clone();
832        let registry = self.registry.clone();
833        let palace_names = self.palace_names.clone();
834        // The directory walk and each `PalaceHandle::open` perform blocking
835        // filesystem + redb/usearch I/O — run the whole hydration on the
836        // blocking pool so it never parks an async worker thread.
837        let count = tokio::task::spawn_blocking(move || -> Result<usize> {
838            let palaces = PalaceRegistry::list_palaces(&registry_dir)?;
839            let total = palaces.len();
840            let mut loaded = 0usize;
841            let mut skipped = 0usize;
842            for palace in palaces {
843                match trusty_common::memory_core::PalaceHandle::open(&palace) {
844                    Ok(handle) => {
845                        tracing::debug!(
846                            palace = %palace.id,
847                            data_dir = %palace.data_dir.display(),
848                            "loaded palace from disk"
849                        );
850                        // Issue #228: seed the in-memory name cache so write
851                        // hot paths (memory_remember / memory_note) can resolve
852                        // the friendly palace name without re-walking the data
853                        // root. Insert here (during hydration) is the single
854                        // point of truth for restart-time population.
855                        palace_names.insert(palace.id.0.clone(), palace.name.clone());
856                        registry.register_arc(handle);
857                        loaded += 1;
858                    }
859                    Err(e) => {
860                        // Why (issue #467): a single bad palace (corrupt kg.db,
861                        // stale WAL, EMFILE — "Too many open files", permissions)
862                        // must never abort startup or block the HTTP server from
863                        // binding. Log per-palace and keep going; the summary
864                        // below tells operators how many were skipped without
865                        // trawling the log.
866                        // The palace is NOT registered in the in-memory registry,
867                        // so the next `open_palace` call for this id will attempt
868                        // a fresh open from disk — the lazy-reopen path. If the
869                        // root cause was EMFILE and the fd-limit fix (#462) raised
870                        // the soft limit to 8192, that first request will succeed.
871                        tracing::warn!(
872                            palace = %palace.id,
873                            data_dir = %palace.data_dir.display(),
874                            "skipping palace during startup hydration: {e:#}; \
875                             will retry lazily on first access"
876                        );
877                        skipped += 1;
878                    }
879                }
880            }
881            tracing::info!(
882                "palace hydration summary: loaded {loaded}/{total} ({skipped} skipped due to errors)"
883            );
884            Ok(loaded)
885        })
886        .await
887        .map_err(|e| anyhow::anyhow!("join load_palaces_from_disk: {e}"))??;
888        Ok(count)
889    }
890
891    /// Builder-style: attach the daemon's shared [`LogBuffer`] so the
892    /// `GET /api/v1/logs/tail` endpoint serves the same lines the tracing
893    /// subscriber captures (issue #35).
894    ///
895    /// Why: `main` builds the buffer (via `init_tracing_with_buffer`) before
896    /// constructing the `AppState`, then hands a clone here so the HTTP
897    /// handler and the tracing layer observe the same ring.
898    /// What: replaces the empty default buffer with the supplied one.
899    /// Test: `logs_tail_returns_recent_lines`.
900    #[must_use]
901    pub fn with_log_buffer(mut self, buffer: trusty_common::log_buffer::LogBuffer) -> Self {
902        self.log_buffer = buffer;
903        self
904    }
905
906    /// Builder-style: mark this daemon as the sole palace writer so palace
907    /// redb files open with `OpenIntent::Writer` (issue #1487).
908    ///
909    /// Why: The HTTP daemon owns the write lock on every palace's `kg.redb`
910    /// and `index.usearch.redb`. Before this fix, when a *second* daemon
911    /// instance opened the same store it silently degraded to a read-only
912    /// snapshot and rejected every `memory_remember` for its lifetime —
913    /// effectively silent data loss when an MCP client routed a write to the
914    /// rogue instance. Opening as `Writer` makes the second instance fail
915    /// loud (after a short handoff-retry window that absorbs a graceful
916    /// launchd `bootout`→`bootstrap` overlap) instead of serving broken
917    /// reads-only. CLI, stdio-proxy, and test code paths never call this, so
918    /// they keep the snapshot read-fallback (issue #59).
919    /// What: Replaces `self.registry` with a fresh `PalaceRegistry` carrying
920    /// `OpenIntent::Writer`.
921    ///
922    /// Invariant: MUST be called on a fresh, unhydrated, unshared registry —
923    /// during startup, before `spawn_startup_tasks`/`load_palaces_from_disk`
924    /// registers any `PalaceHandle` and before the `AppState` (hence its
925    /// `Arc<PalaceRegistry>`) is cloned to a handler. Replacing the registry
926    /// discards the prior `Arc`; doing so after hydration would silently drop
927    /// live handles (data loss), and doing so after the state is shared would
928    /// leave other clones on the stale read-only registry. The guard is a
929    /// `debug_assert!` on the strongest cheap signals the registry exposes —
930    /// `is_empty()` (no handles hydrated) and `Arc::strong_count == 1` (not yet
931    /// shared) — so an ordering violation fails fast as the programmer error it
932    /// is (the call site is startup-only and fixed). Release builds elide the
933    /// assert; the real call site (`run_serve`) always satisfies it.
934    /// Test: `with_writer_intent_marks_registry_writer` and
935    /// `with_writer_intent_panics_on_hydrated_registry` in `lib_tests`.
936    #[must_use]
937    pub fn with_writer_intent(mut self) -> Self {
938        // Fail fast on an ordering bug: a hydrated registry (`!is_empty`) or a
939        // shared one (`strong_count > 1`) would silently drop live handles or
940        // strand other clones on the stale read-only registry (issue #1487).
941        debug_assert!(self.registry.is_empty() && Arc::strong_count(&self.registry) == 1);
942        // Idle-to-disk: preserve the configurable open-handle cap
943        // (TRUSTY_MEMORY_MAX_OPEN_PALACES) while marking the registry a writer.
944        self.registry = Arc::new(PalaceRegistry::from_env().with_writer_intent());
945        self
946    }
947
948    /// Builder-style: attach the bug-capture `ErrorStore` handle (bug-reporting #478).
949    ///
950    /// Why: Phase 2 MCP / HTTP endpoints need a handle to the in-memory error
951    ///      ring so they can serve `recent_errors` / `errors_by_fingerprint`
952    ///      without disk I/O on the hot path. Installing it here — rather than
953    ///      adding it as a separate global — keeps the state graph explicit and
954    ///      lets tests skip it by never calling this method.
955    /// What: stores `Some(store)` in `AppState::error_store`; the `BugCaptureLayer`
956    ///      that writes to this store is already installed in the tracing
957    ///      subscriber by `init_tracing_with_buffer_and_capture`. The store is
958    ///      `Clone` (cheap `Arc` clone internally) so both the layer and this
959    ///      field share the same underlying ring.
960    /// Test: Phase 2 will add `error_store_captures_and_queries` in `web.rs`.
961    #[must_use]
962    pub fn with_error_store(mut self, store: trusty_common::error_capture::ErrorStore) -> Self {
963        self.error_store = Some(store);
964        self
965    }
966
967    /// Builder-style: opt into multi-tenant authorization mode (issue #1714).
968    ///
969    /// Why: mirrors `with_bm25_client_from_env`'s pattern of keeping env-var
970    /// gating in one place. Unset (the default) preserves today's
971    /// single-tenant behaviour with zero change for existing callers. Issue
972    /// #2522 review: activation is silent otherwise, which makes a
973    /// misconfigured (or unexpectedly enabled) deployment hard to diagnose
974    /// from logs alone — log once at startup when the mode flips on.
975    /// What: sets `multi_tenant_mode` from `TRUSTY_MEMORY_MULTI_TENANT=1`; see
976    /// the `authz` module for what the flag then enforces. Logs via
977    /// `tracing::info!` (stderr only) when enabled; stays silent when
978    /// disabled (the default).
979    /// Test: `authorize_force_palace_create_denies_multi_tenant_without_capability`.
980    #[must_use]
981    pub fn with_multi_tenant_mode_from_env(mut self) -> Self {
982        self.multi_tenant_mode = std::env::var("TRUSTY_MEMORY_MULTI_TENANT").as_deref() == Ok("1");
983        if self.multi_tenant_mode {
984            tracing::info!(
985                "multi-tenant mode enabled (TRUSTY_MEMORY_MULTI_TENANT=1): force=true palace_create will be refused"
986            );
987        }
988        self
989    }
990
991    /// Send a `DaemonEvent` to all connected SSE subscribers and persist
992    /// it to the activity log when the variant carries a source.
993    ///
994    /// Why: Mutating handlers call this after a successful write so the
995    /// dashboard can update without polling. The send is best-effort —
996    /// `broadcast::Sender::send` returns `Err` only when there are no live
997    /// receivers, which is fine (no listeners == no work to do). Issue
998    /// #96 additionally writes the entry to the persistent activity log
999    /// so the feed can serve historical rows on page load and so MCP /
1000    /// HTTP / Hook origins are visible to the operator. Persistence is
1001    /// also best-effort — a write failure is logged but never blocks the
1002    /// SSE broadcast.
1003    ///
1004    /// Issue #232: the activity-log append is a synchronous redb write +
1005    /// fsync. Calling it directly on the async caller's task parked a tokio
1006    /// worker thread on disk I/O for every SSE event. We now offload the
1007    /// append to the blocking thread pool via `spawn_blocking` and return
1008    /// immediately — `emit` stays synchronous so every existing caller
1009    /// (including the sync `dispatch_hook_fired` JSON-RPC handler) keeps
1010    /// compiling unchanged. The fire-and-forget pattern matches the
1011    /// pre-fix semantics (best-effort, never blocks the SSE broadcast)
1012    /// while freeing the async runtime to do real work during the write.
1013    /// What: serialises the event for the log (skipping `StatusChanged`
1014    /// which is a recomputed aggregate, not a mutation), spawns the redb
1015    /// append on `tokio::task::spawn_blocking` keyed by a clone of the
1016    /// `Arc<ActivityLog>` and the cloned event, then sends the event over
1017    /// the broadcast channel. A `pending_activity_writes` counter is bumped
1018    /// before the spawn and decremented inside the closure so
1019    /// [`Self::flush_activity_writes`] can drain in tests.
1020    /// Test: `web::tests::sse_stream_receives_palace_created` confirms a
1021    /// subscriber observes the emitted event;
1022    /// `activity_endpoint_lists_recent_emits` confirms persistence via
1023    /// `flush_activity_writes`.
1024    pub fn emit(&self, event: DaemonEvent) {
1025        if let Some(source) = event.source() {
1026            let event_type = event.type_str();
1027            let palace_id = event.palace_id().map(|s| s.to_string());
1028            let log = Arc::clone(&self.activity_log);
1029            let event_for_log = event.clone();
1030            let pending = Arc::clone(&self.pending_activity_writes);
1031            // Pre-allocate the sequence id in the emitting thread so the
1032            // persisted order matches the emission order even when blocking-pool
1033            // workers execute the writes concurrently (issue #247). Without
1034            // this, four rapid emits would assign IDs inside their respective
1035            // `spawn_blocking` closures in a non-deterministic order.
1036            let id = log.alloc_id();
1037            pending.fetch_add(1, Ordering::SeqCst);
1038            // Why: the synchronous redb append + fsync must not park an
1039            // async worker thread (issue #232). Spawn the write on the
1040            // blocking pool; the JoinHandle is intentionally dropped —
1041            // the write is best-effort and any failure is logged below.
1042            tokio::task::spawn_blocking(move || {
1043                let result = log.append_with_id(id, source, palace_id, event_type, &event_for_log);
1044                if let Err(e) = result {
1045                    tracing::warn!("activity_log.append failed for {event_type}: {e:#}");
1046                }
1047                pending.fetch_sub(1, Ordering::SeqCst);
1048            });
1049        }
1050        let _ = self.events.send(event);
1051    }
1052
1053    /// Block (asynchronously) until every in-flight activity-log write
1054    /// spawned by [`Self::emit`] has settled.
1055    ///
1056    /// Why: `emit` offloads its redb append to `tokio::task::spawn_blocking`
1057    /// and returns immediately (issue #232). Tests that observe the
1058    /// activity log right after a burst of emits would otherwise race the
1059    /// blocking-pool worker; this helper gives them a deterministic
1060    /// synchronization point. Production code never needs to call this —
1061    /// the dashboard reads through `GET /api/v1/activity`, which already
1062    /// tolerates writes settling asynchronously.
1063    /// What: spins on `pending_activity_writes` with a 1 ms yield until the
1064    /// counter is zero. Cheap: tests typically emit a handful of events
1065    /// and the loop exits within a single scheduler tick.
1066    /// Test: covered indirectly by `emit_persists_mutations_but_skips_status_changed`
1067    /// and `web::tests::activity_endpoint_lists_recent_emits`.
1068    pub async fn flush_activity_writes(&self) {
1069        while self.pending_activity_writes.load(Ordering::SeqCst) > 0 {
1070            tokio::time::sleep(std::time::Duration::from_millis(1)).await;
1071        }
1072    }
1073
1074    /// Open (or return cached) the chat-session store for a palace.
1075    ///
1076    /// Why: Chat session persistence lives in a dedicated redb file under
1077    /// the palace's data dir (`chat_sessions.redb`) so it doesn't intermingle
1078    /// with the KG's transactional load. The store is cheap to clone via
1079    /// `Arc` but the underlying connection should be reused, so cache by id.
1080    /// What: delegates to the LRU-bounded [`SessionStoreCache`], which creates
1081    /// the palace data dir if missing, opens (or reuses) a `ChatSessionStore`,
1082    /// and evicts cold, *unused* stores once more than the cap are resident.
1083    /// Callers keep the returned `Arc` for as long as they need it — eviction
1084    /// never closes a store someone still holds.
1085    /// Test: `session_store_cache::tests::open_handles_are_bounded_by_cap`,
1086    /// `session_store_cache::tests::in_use_store_is_never_evicted`; the call
1087    /// path is covered indirectly by the session HTTP handlers in `web::tests`.
1088    pub fn session_store(&self, palace_id: &str) -> Result<Arc<ChatSessionStore>> {
1089        // #4639: bounded cache replaces the unbounded, never-evicting DashMap.
1090        self.session_stores
1091            .get_or_open(palace_id, &self.data_root.join(palace_id))
1092    }
1093
1094    /// Builder-style setter for the default palace name.
1095    ///
1096    /// Why: `serve --palace <name>` wants to bind every tool call to a
1097    /// project-scoped namespace without forcing every MCP request to repeat
1098    /// the palace argument.
1099    /// What: Returns `self` with `default_palace = Some(name)`.
1100    /// Test: `default_palace_used_when_arg_omitted` covers the resolution
1101    /// path; this setter is exercised there.
1102    pub fn with_default_palace(mut self, name: Option<String>) -> Self {
1103        self.default_palace = name;
1104        self
1105    }
1106
1107    /// Resolve (or initialize) the shared embedder.
1108    ///
1109    /// Why: FastEmbedder load is expensive — we share one instance across all
1110    /// tool calls; the `OnceCell` ensures concurrent first-use races collapse
1111    /// to a single load.
1112    /// What: Returns `Arc<FastEmbedder>` on success. Errors propagate from the
1113    /// underlying ONNX load.
1114    /// Test: Indirectly via `dispatch_remember_then_recall`.
1115    /// Resolve the active chat provider, auto-detecting on first call.
1116    ///
1117    /// Why: Provider selection depends on filesystem-loaded config plus a
1118    /// network probe (Ollama liveness), so it must be lazily initialised at
1119    /// runtime. Caching the choice in a `OnceCell` keeps it stable across
1120    /// concurrent requests without re-probing on every chat call.
1121    /// What: On first use loads `~/.trusty-memory/config.toml`, prefers an
1122    /// auto-detected Ollama instance (when `local_model.enabled`), and falls
1123    /// back to OpenRouter when an API key is set. Returns `Ok(None)` when
1124    /// neither is available so the caller can emit a 412.
1125    /// Test: `web::tests::providers_endpoint_returns_payload` covers the
1126    /// detection path indirectly through `/api/v1/chat/providers`.
1127    pub async fn chat_provider(&self) -> Option<Arc<dyn ChatProvider>> {
1128        self.chat_provider
1129            .get_or_init(|| async {
1130                // Why (issue #226): `service::load_user_config` is the
1131                //      axum-free home of the loader; the `web::load_user_config`
1132                //      re-export only exists for the HTTP handlers. Going
1133                //      direct to `service` keeps this method usable when
1134                //      the `axum-server` feature is disabled.
1135                let cfg = crate::service::load_user_config().unwrap_or_default();
1136                if cfg.local_model.enabled {
1137                    if let Some(mut p) =
1138                        trusty_common::auto_detect_local_provider(&cfg.local_model.base_url).await
1139                    {
1140                        // auto_detect returns an empty model id; callers must
1141                        // set the configured model name themselves.
1142                        p.model = cfg.local_model.model.clone();
1143                        return Some(Arc::new(p) as Arc<dyn ChatProvider>);
1144                    }
1145                }
1146                if !cfg.openrouter_api_key.is_empty() {
1147                    return Some(Arc::new(trusty_common::OpenRouterProvider::new(
1148                        cfg.openrouter_api_key,
1149                        cfg.openrouter_model,
1150                    )) as Arc<dyn ChatProvider>);
1151                }
1152                None
1153            })
1154            .await
1155            .clone()
1156    }
1157
1158    /// Spawn a fire-and-forget background task that auto-discovers project
1159    /// aliases under `project_root` and asserts new ones into `palace`.
1160    ///
1161    /// Why (issue #42): Projects carry implicit shorthand — cargo package
1162    /// names that differ from their directory, binary names that differ
1163    /// from packages, first-letter abbreviations — that should be surfaced
1164    /// without a user ever calling `add_alias`. Running discovery as a
1165    /// detached task on palace-open keeps startup latency unchanged: the
1166    /// daemon binds and starts serving immediately while the discovery scan
1167    /// completes in the background, and any newly-asserted aliases land in
1168    /// the prompt cache before the model's next `get_prompt_context` call.
1169    /// What: clones `self` (cheap; `Arc`-backed), spawns a tokio task that
1170    /// invokes the `discover_aliases` tool handler directly so the
1171    /// dedup + cache-rebuild logic runs exactly the same path as the MCP
1172    /// tool call. Errors are logged at `warn!`; one failed discovery never
1173    /// destabilises the daemon.
1174    /// Test: not unit-tested (timing-dependent fire-and-forget); the
1175    /// underlying `discover_aliases` dispatch is covered by
1176    /// `dispatch_discover_aliases_inserts_new_and_dedupes` in `tools::tests`.
1177    pub fn spawn_alias_discovery(&self, palace: String, project_root: PathBuf) {
1178        let state = self.clone();
1179        tokio::spawn(async move {
1180            let args = serde_json::json!({
1181                "palace": palace,
1182                "project_root": project_root.to_string_lossy(),
1183            });
1184            match tools::dispatch_tool(&state, "discover_aliases", args).await {
1185                Ok(result) => tracing::info!(
1186                    new = ?result.get("new"),
1187                    already_known = ?result.get("already_known"),
1188                    "alias discovery complete"
1189                ),
1190                Err(e) => tracing::warn!("alias discovery failed: {e:#}"),
1191            }
1192        });
1193    }
1194
1195    /// Return the current readiness state.
1196    ///
1197    /// Why: tool handlers and the `/health` endpoint need a cheap, lock-free
1198    /// way to check whether the embedder has been initialised yet.
1199    ///
1200    /// What: loads `daemon_readiness` with `Acquire` ordering so the caller
1201    /// sees all writes the startup task made before setting the state.
1202    ///
1203    /// #4836: the value read here is only as good as the writes that reach it.
1204    /// It used to be written exactly once, by the startup warm-up task, so a
1205    /// single failed attempt pinned the daemon at `Warming` for the rest of its
1206    /// life. [`AppState::embedder`] now also flips it on success, which is the
1207    /// signal that actually proves a vector search can run.
1208    /// Test: `daemon_readiness_transitions_warming_to_ready`,
1209    /// `resolving_the_embedder_marks_a_warming_daemon_ready`.
1210    pub fn readiness(&self) -> DaemonReadiness {
1211        DaemonReadiness::from_u8(self.daemon_readiness.load(Ordering::Acquire))
1212    }
1213
1214    /// Flip the readiness state from `Warming` to `Ready`.
1215    ///
1216    /// Why: called by `spawn_startup_tasks` in `main.rs` once the embedder
1217    /// warm-up succeeds — this is the single state-transition site.
1218    /// What: `store(Ready, Release)` so subsequent `Acquire` loads in handlers
1219    /// observe a consistent state.  Idempotent: calling it multiple times is
1220    /// harmless.
1221    /// Test: `daemon_readiness_transitions_warming_to_ready`.
1222    pub fn set_ready(&self) {
1223        self.daemon_readiness
1224            .store(DaemonReadiness::Ready as u8, Ordering::Release);
1225    }
1226
1227    /// Obtain the shared `FastEmbedder` instance, initialising it on first call.
1228    ///
1229    /// Why: centralises lazy embedder access so every tool handler goes through
1230    /// one bounded init path (tracks #910 internally).
1231    /// What: wraps `OnceCell::get_or_try_init` with a timeout so a slow
1232    /// CoreML/CUDA first-compile cannot block a handler indefinitely.  On
1233    /// timeout the `OnceCell` is left unresolved and the next caller retries.
1234    ///
1235    /// **Callers on the request path SHOULD check `readiness()` before this
1236    /// method** (issue #1970) — every recall handler now checks
1237    /// `readiness() == Ready` first and only calls `embedder()` on that
1238    /// branch, falling back to a BM25/L0/L1-only path while `Warming` instead
1239    /// of paying this method's cold-init cost. Reaching this method while
1240    /// still `Warming` is not a bug (the warm-up task itself calls
1241    /// `embedder()` while in `Warming` state), just unusual on the request
1242    /// path.
1243    ///
1244    /// This timeout is a backstop against a pathological init delay (e.g. the
1245    /// warm-up task's own call, or a handler that skips the `readiness()`
1246    /// check). If this timeout fires the `OnceCell` is left in the unresolved
1247    /// state and the next call retries from scratch.
1248    pub async fn embedder(&self) -> Result<Arc<dyn Embedder + Send + Sync>> {
1249        // #4836: delegate to the ONE process-wide cell instead of a private
1250        // second one, so initialising the embedder here is the same event the
1251        // startup warm-up latches readiness off. `shared_embedder` already
1252        // applies the bounded `TRUSTY_EMBEDDER_INIT_TIMEOUT_SECS` init timeout
1253        // and the CoreML auto-fallback, so no wrapper timeout is needed here.
1254        let embedder = trusty_common::memory_core::retrieval::shared_embedder().await?;
1255        // #4836: a resolved embedder is proof a vector search can run, so it is
1256        // the authoritative readiness signal — not the startup task's one-shot
1257        // attempt. Without this the daemon stays `Warming` forever whenever that
1258        // single attempt failed, and every MCP recall serves the degraded
1259        // L0/L1 fallback that ignores the query entirely.
1260        self.set_ready();
1261        Ok(embedder)
1262    }
1263}
1264
1265impl std::fmt::Debug for AppState {
1266    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1267        f.debug_struct("AppState")
1268            .field("version", &self.version)
1269            .field("data_root", &self.data_root)
1270            .field("registry_len", &self.registry.len())
1271            .finish()
1272    }
1273}
1274
1275/// Handle a single MCP JSON-RPC message and produce its response.
1276///
1277/// Why: Pulled out of the stdio loop so unit tests can drive every method
1278/// without touching real stdin/stdout.
1279/// What: Routes `initialize`, `tools/list`, `tools/call`, `ping`, and the
1280/// `notifications/initialized` notification (which returns `Value::Null`).
1281/// Test: See unit tests below — initialize/list/call all return expected
1282/// JSON-RPC envelopes; notifications return `Null` (no response written).
1283pub async fn handle_message(state: &AppState, msg: Value) -> Value {
1284    let id = msg.get("id").cloned().unwrap_or(Value::Null);
1285    let method = msg.get("method").and_then(|m| m.as_str()).unwrap_or("");
1286
1287    match method {
1288        "initialize" => {
1289            let extra = state
1290                .default_palace
1291                .as_ref()
1292                .map(|dp| json!({ "default_palace": dp }));
1293            let result = initialize_response("trusty-memory", &state.version, extra);
1294            // Why (issue #42): prompt-facts now flow through the
1295            // per-message `get_prompt_context` tool rather than MCP
1296            // prompts, so we no longer advertise the `prompts` capability.
1297            json!({
1298                "jsonrpc": "2.0",
1299                "id": id,
1300                "result": result,
1301            })
1302        }
1303        // Notifications must NOT receive a response.
1304        "notifications/initialized" | "notifications/cancelled" => Value::Null,
1305        "tools/list" => json!({
1306            "jsonrpc": "2.0",
1307            "id": id,
1308            "result": tools::tool_definitions_with(state.default_palace.is_some())
1309        }),
1310        // OpenRPC 1.3.2 discovery — see `openrpc.rs`. Returns the full
1311        // service description so orchestrators (trusty-agents, etc.) can
1312        // introspect every tool and its required `memory.read`/`memory.write`
1313        // scope without bespoke per-server adapters.
1314        "rpc.discover" => json!({
1315            "jsonrpc": "2.0",
1316            "id": id,
1317            "result": openrpc::build_discover_response(
1318                &state.version,
1319                state.default_palace.is_some(),
1320            ),
1321        }),
1322        "tools/call" => {
1323            let params = msg.get("params").cloned().unwrap_or_default();
1324            let tool_name = params
1325                .get("name")
1326                .and_then(|n| n.as_str())
1327                .unwrap_or("")
1328                .to_string();
1329            let args = params.get("arguments").cloned().unwrap_or_default();
1330            match tools::dispatch_tool(state, &tool_name, args).await {
1331                Ok(content) => {
1332                    // Why: tools that return a bare JSON string (e.g.
1333                    // `get_prompt_context` returning the formatted
1334                    // Markdown block) should surface as plain text in the
1335                    // MCP `content[0].text` field — wrapping in
1336                    // `Value::to_string()` would re-quote the payload and
1337                    // force every caller to strip outer quotes.
1338                    let text = match &content {
1339                        Value::String(s) => s.clone(),
1340                        other => other.to_string(),
1341                    };
1342                    json!({
1343                        "jsonrpc": "2.0",
1344                        "id": id,
1345                        "result": {
1346                            "content": [{"type": "text", "text": text}]
1347                        }
1348                    })
1349                }
1350                Err(e) => json!({
1351                    "jsonrpc": "2.0",
1352                    "id": id,
1353                    // Why: anyhow's `{:#}` alternate format walks the full
1354                    // `Caused by:` chain so MCP clients see actionable
1355                    // detail (e.g. "PalaceHandle::remember_with_options:
1356                    // filter rejected: too short") instead of just the
1357                    // outermost context label.
1358                    "error": {"code": -32603, "message": format!("{e:#}")}
1359                }),
1360            }
1361        }
1362        "ping" => json!({"jsonrpc": "2.0", "id": id, "result": {}}),
1363        _ => json!({
1364            "jsonrpc": "2.0",
1365            "id": id,
1366            "error": {
1367                "code": -32601,
1368                "message": format!("Method not found: {method}")
1369            }
1370        }),
1371    }
1372}
1373
1374#[cfg(test)]
1375mod lib_tests;