Skip to main content

trusty_memory/
lib.rs

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