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