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