Skip to main content

trusty_memory/service/
core.rs

1//! `MemoryService` — the pure business-logic facade over `AppState`.
2//!
3//! Why: lets the axum HTTP handlers stay thin one-liners and lets non-HTTP
4//! callers (chat tool dispatch, RPC bridges) reuse the same code paths without
5//! dragging axum types around (split out of the former monolithic `service.rs`,
6//! issue #607).
7//! What: the `MemoryService` struct + its full async method surface, moved
8//! verbatim. Each method returns `anyhow::Result<Value>` or a typed
9//! `ServiceResult`.
10//! Test: every method is covered by the corresponding handler test in
11//! `web::tests`.
12
13use crate::attribution::CreatorInfo;
14use crate::{ActivitySource, AppState, DaemonEvent};
15use anyhow::{anyhow, Context, Result};
16use serde_json::{json, Value};
17use std::sync::Arc;
18use trusty_common::memory_core::palace::{Palace, PalaceId, RoomType};
19use trusty_common::memory_core::retrieval::{
20    recall_across_palaces_with_default_embedder, recall_deep_with_default_embedder,
21    recall_with_default_embedder, RememberOptions,
22};
23use trusty_common::memory_core::store::PalaceStoreError;
24use trusty_common::memory_core::PalaceRegistry;
25use uuid::Uuid;
26
27use super::helpers::{
28    collect_palace_stats, drawer_content_preview, drawer_snippet, is_reserved_system_palace,
29    list_palaces_blocking, palace_info_blocking, palace_info_from, recall_entry_json,
30};
31use super::recall_stream::recall_streamed;
32use super::types::{
33    CreateDrawerBody, CreatePalaceBody, ListDrawersQuery, PalaceInfo, ServiceError, ServiceResult,
34    StatusPayload,
35};
36
37/// Hard cap on triples returned by the per-palace graph endpoint.
38pub(super) const KG_GRAPH_MAX_TRIPLES: usize = 5_000;
39
40// ---------------------------------------------------------------------------
41// MemoryService — pure business logic facade.
42// ---------------------------------------------------------------------------
43
44/// Wraps [`AppState`] and exposes one async method per logical operation.
45///
46/// Why: see module docs. Lets HTTP handlers stay thin and lets non-HTTP
47/// callers (chat tool dispatch, RPC bridges) reuse the same code paths.
48/// What: `Clone` (cheap — only the inner `AppState` is shared); construct
49/// with `MemoryService::new(state)`.
50/// Test: every method is covered by the corresponding handler test in
51/// `web::tests`.
52#[derive(Clone)]
53pub struct MemoryService {
54    pub(super) state: AppState,
55}
56
57impl MemoryService {
58    /// Construct a new service wrapper.
59    ///
60    /// Why: handlers cheaply re-wrap their `AppState` on every request; the
61    /// cost is just an `Arc` clone, so we don't bother caching the wrapper.
62    /// What: stores the `AppState` for later method calls.
63    /// Test: trivial — covered indirectly by every handler test.
64    pub fn new(state: AppState) -> Self {
65        Self { state }
66    }
67
68    /// Borrow the inner [`AppState`].
69    ///
70    /// Why: some handlers still need direct access (SSE broadcaster, session
71    /// store, etc.) while we incrementally extract code into the service.
72    /// What: returns a borrowed reference to the wrapped `AppState`.
73    /// Test: not directly tested; surface-level accessor.
74    pub fn state(&self) -> &AppState {
75        &self.state
76    }
77
78    // -----------------------------------------------------------------
79    // Status / config
80    // -----------------------------------------------------------------
81
82    /// Build the aggregate `/api/v1/status` payload.
83    ///
84    /// Why: dashboard widgets and the MCP `get_status` tool need the same
85    /// roll-up; centralising avoids drift between the two surfaces.
86    /// What: walks every persisted palace for `palace_count`, then sums
87    /// drawer/vector/triple counts across the cache-resident subset and
88    /// returns the [`StatusPayload`].
89    /// Why (issue #4637): this used to open every persisted palace to sum
90    /// those three counts. With 5,794 palaces on disk against a 64-slot LRU
91    /// that is ~5,730 cold opens of ~1s each — the endpoint measurably never
92    /// responded. `palace_count` still reflects the true on-disk total (the
93    /// directory walk is cheap and now runs on the blocking pool); the totals
94    /// cover only cache-resident palaces and say so via `cached_palace_count`.
95    /// Test: `status_endpoint_returns_payload`,
96    /// `status_does_not_open_uncached_palaces`.
97    pub async fn status(&self) -> StatusPayload {
98        // The `/status` endpoint is the one place we still want a disk view —
99        // an operator hitting this endpoint right after restart (before
100        // `load_palaces_from_disk` finishes) should still see every persisted
101        // palace counted, even if it isn't in the in-memory registry yet.
102        let palaces = list_palaces_blocking(&self.state).await.unwrap_or_default();
103        let palace_count = palaces.len();
104        // #4637: peek() not open_palace() — full-registry open is O(n) cold disk I/O
105        let stats = collect_palace_stats(&self.state, palaces.iter().map(|p| &p.id));
106        StatusPayload {
107            version: self.state.version.clone(),
108            palace_count,
109            default_palace: self.state.default_palace.clone(),
110            data_root: self.state.data_root.display().to_string(),
111            total_drawers: stats.total_drawers,
112            total_vectors: stats.total_vectors,
113            total_kg_triples: stats.total_kg_triples,
114            cached_palace_count: stats.cached_palace_count,
115        }
116    }
117
118    /// Compute the aggregate `StatusChanged` event used by SSE consumers.
119    ///
120    /// Why: mutating handlers — and the periodic status ticker — push a
121    /// refreshed status snapshot so dashboards stay in sync without an
122    /// extra `/api/v1/status` request.
123    /// Why (issue #228): this used to call `PalaceRegistry::list_palaces`
124    /// (a synchronous disk walk) + `open_palace` (more disk I/O on first
125    /// call) for every palace on every emit. Since every persisted palace
126    /// is already loaded into the in-memory registry by
127    /// `AppState::load_palaces_from_disk` at startup (and every `create_palace`
128    /// keeps it in sync), iterating the in-memory registry returns the same
129    /// counts without touching disk.
130    /// What: iterates `state.registry.list()` (a `DashMap` snapshot) and
131    /// sums the live handle stats via [`collect_palace_stats`]. Returns a
132    /// `DaemonEvent::StatusChanged`. Palaces that fail to resolve in the
133    /// registry (race during shutdown) are silently skipped — the next
134    /// emit will catch them.
135    /// Test: indirectly via SSE integration tests; the math is identical to
136    /// the disk-walk implementation and the `status_endpoint_returns_payload`
137    /// test still passes against `status()` (which keeps the disk view for
138    /// the dedicated endpoint).
139    pub fn aggregate_status_event(&self) -> DaemonEvent {
140        let ids: Vec<PalaceId> = self.state.registry.list();
141        let stats = collect_palace_stats(&self.state, ids.iter());
142        DaemonEvent::StatusChanged {
143            total_drawers: stats.total_drawers,
144            total_vectors: stats.total_vectors,
145            total_kg_triples: stats.total_kg_triples,
146        }
147    }
148
149    // -----------------------------------------------------------------
150    // Palaces
151    // -----------------------------------------------------------------
152
153    /// List every palace on disk, enriched with live handle stats.
154    ///
155    /// Why: shared between the HTTP handler and the chat tool dispatcher;
156    /// both want the same `PalaceInfo` shape. Issue #185 added the
157    /// reserved-prefix filter so internal "system" palaces (e.g. the
158    /// `__health_probe__` palace used by `/health`) never surface in the
159    /// admin UI, TUI, or any user-facing roster.
160    /// What: walks the registry, drops any palace whose id starts with the
161    /// reserved `__` prefix, and builds a `PalaceInfo` per remaining row.
162    /// Why (issue #4637): this used to call `open_palace` per row purely to
163    /// enrich it with counts. At 5,794 palaces against a 64-slot LRU that is
164    /// ~90 minutes of cold, blocking disk I/O inline on the async executor —
165    /// and it evicted the entire working set on every call. Rows now come
166    /// from `PalaceRegistry::peek` (zero I/O, no LRU promotion, mirroring the
167    /// #1924 fix in `console_metrics.rs`). Uncached rows carry `cached: false`
168    /// and zero counts; a client that needs live counts for one palace should
169    /// fetch `GET /api/v1/palaces/{id}`, which still opens it.
170    /// Test: `palace_list_includes_richer_counts`, `palace_list_includes_graph_counts`,
171    /// `health_probe_palace_is_invisible` (in `web::tests`),
172    /// `list_palaces_does_not_open_uncached_palaces`.
173    pub async fn list_palaces(&self) -> ServiceResult<Vec<PalaceInfo>> {
174        let palaces = list_palaces_blocking(&self.state)
175            .await
176            .map_err(|e| ServiceError::internal(format!("{e:#}")))?;
177        // #7106: the enrichment itself is blocking work — a resident palace's
178        // first community count after a write partitions its whole graph. One
179        // hop for the whole loop, not one per row: this list can be thousands
180        // of rows and a task each would cost more than the work.
181        let registry = Arc::clone(&self.state.registry);
182        let out = tokio::task::spawn_blocking(move || {
183            let mut out = Vec::with_capacity(palaces.len());
184            for p in palaces {
185                if is_reserved_system_palace(&p.id) {
186                    continue;
187                }
188                // #4637: peek() not open_palace() — full-registry open is O(n) cold disk I/O
189                let handle = registry.peek(&p.id);
190                out.push(palace_info_from(&p, handle.as_ref()));
191            }
192            out
193        })
194        .await
195        .map_err(|e| ServiceError::internal(format!("join list_palaces enrichment: {e}")))?;
196        Ok(out)
197    }
198
199    /// Every non-system palace with REAL counts, keeping per-palace failures.
200    ///
201    /// Why (#6286): [`Self::list_palaces`] answers placeholder zeros plus
202    /// `cached: false` for any palace not already resident, which is why the
203    /// monitor could not use it and fanned out one [`Self::get_palace`] per id
204    /// instead. That fan-out then dropped a palace whose call failed at
205    /// `debug!`, so the panel could show "12 palaces" over 9 rows and nothing
206    /// said why. This is the one call that answers what the fan-out was
207    /// assembling, and it reports a failure as a failure rather than as an
208    /// absence.
209    ///
210    /// What: one entry per non-system palace, in registry order. `Ok` carries
211    /// the same [`PalaceInfo`] `get_palace` builds — the palace is opened, so
212    /// the counts are measurements. `Err` carries the open failure's message.
213    /// A palace never silently vanishes and never becomes a row of zeros.
214    ///
215    /// **This opens every palace, and that is the point.** #4637 removed
216    /// exactly this from `list_palaces` because at 5,794 palaces a cold open per
217    /// row is ~90 minutes of blocking disk I/O. The cost is unchanged from the
218    /// N-call fan-out this replaces — the same opens, one round trip instead of
219    /// N — and after the first poll the registry is warm. A caller that wants
220    /// cheap approximate rows still has `list_palaces`.
221    ///
222    /// # Errors
223    ///
224    /// Only when the registry itself cannot be walked. A palace that will not
225    /// open is an `Err` entry, not an error for the whole call.
226    ///
227    /// **Every open runs on the blocking pool (#6836).** The opens are cold
228    /// disk I/O, so running them inline on a tokio worker parked the executor
229    /// for the whole sweep — on a many-palace install that is minutes during
230    /// which nothing else the daemon serves makes progress. One
231    /// `spawn_blocking` per palace also gives the executor a yield point
232    /// between palaces rather than one uninterruptible block.
233    ///
234    /// Test: `rpc_palaces_list_reports_counts_per_palace`,
235    /// `rpc_palaces_list_reports_an_unreadable_palace_rather_than_dropping_it`,
236    /// `list_palaces_with_counts_opens_palaces_off_the_executor`.
237    pub async fn list_palaces_with_counts(
238        &self,
239    ) -> ServiceResult<Vec<(String, Result<PalaceInfo, String>)>> {
240        let palaces = list_palaces_blocking(&self.state)
241            .await
242            .map_err(|e| ServiceError::internal(format!("{e:#}")))?;
243        let mut out = Vec::with_capacity(palaces.len());
244        for p in palaces {
245            if is_reserved_system_palace(&p.id) {
246                continue;
247            }
248            let id = p.id.0.clone();
249            // #6836: the open is cold disk I/O — hop to the blocking pool so a
250            // full-estate sweep cannot park the async executor for its duration.
251            let registry = Arc::clone(&self.state.registry);
252            let root = self.state.data_root.clone();
253            let row =
254                tokio::task::spawn_blocking(move || match registry.open_palace(&root, &p.id) {
255                    Ok(handle) => Ok(palace_info_from(&p, Some(&handle))),
256                    Err(e) => Err(format!("{e:#}")),
257                })
258                .await
259                // A join failure is still a per-palace failure: the row says why
260                // rather than vanishing, exactly as an open failure does.
261                .unwrap_or_else(|e| Err(format!("join open palace: {e}")));
262            out.push((id, row));
263        }
264        Ok(out)
265    }
266
267    /// Create a new palace and emit the corresponding activity event.
268    ///
269    /// Why: trims duplicated work between the HTTP handler and any future
270    /// non-HTTP creation flow.
271    /// What: validates the name, builds the `Palace` row, calls
272    /// `PalaceRegistry::create_palace`, and emits `PalaceCreated`. Returns
273    /// the new palace id.
274    /// Test: covered indirectly by `palace_list_includes_richer_counts` (which
275    /// posts a palace through the HTTP layer then reads it back).
276    pub async fn create_palace(
277        &self,
278        body: CreatePalaceBody,
279        source: ActivitySource,
280    ) -> ServiceResult<String> {
281        let name = body.name.trim().to_string();
282        if name.is_empty() {
283            return Err(ServiceError::bad_request("name is required"));
284        }
285        // Issue #88 / Change 2: enforce palace = project mapping for
286        // HTTP-originated palace creation. The validation cwd is, in order of
287        // preference:
288        //   a. `body.cwd` — the caller explicitly supplied their project path
289        //      (correct for any client that is not the daemon itself).
290        //   b. `std::env::current_dir()` — daemon's own cwd, the pre-Change-2
291        //      fallback (rarely meaningful when the daemon is launched from ~).
292        // This keeps older clients that omit `cwd` working without a breaking
293        // change, while letting pin-file-aware clients get accurate validation.
294        // spec-001: `force=true` lets an application bypass the project-slug
295        // gate so it can create palaces under arbitrary slugs (e.g. one per
296        // app/tenant for chat-session storage). The env-var bypass remains for
297        // test contexts; both short-circuit the same validation call.
298        //
299        // Issue #1714: `force=true` bypasses slug validation entirely, so it
300        // is gated behind the minimal authz seam in `crate::authz` before any
301        // other check runs. In the default single-tenant mode this is a
302        // no-op (unchanged behaviour); in multi-tenant mode it fails closed
303        // until a real capability check lands. See `crate::authz` module
304        // docs for the full design rationale.
305        let skip_enforcement =
306            std::env::var("TRUSTY_SKIP_PALACE_ENFORCEMENT").as_deref() == Ok("1");
307        if body.force {
308            crate::authz::authorize_force_palace_create(&self.state)
309                .map_err(|e| ServiceError::forbidden(e.to_string()))?;
310        }
311        if !skip_enforcement && !body.force {
312            let cwd = body
313                .cwd
314                .as_deref()
315                .map(std::path::Path::new)
316                .map(|p| p.to_path_buf())
317                .or_else(|| std::env::current_dir().ok())
318                .unwrap_or_else(|| self.state.data_root.clone());
319            crate::project_root::validate_palace_name(&name, &cwd)
320                .map_err(|e| ServiceError::bad_request(e.to_string()))?;
321        }
322        let id = PalaceId::new(&name);
323        let palace = Palace {
324            id: id.clone(),
325            name: name.clone(),
326            description: body.description.filter(|s| !s.is_empty()),
327            created_at: chrono::Utc::now(),
328            data_dir: self.state.data_root.join(&name),
329        };
330        self.state
331            .registry
332            .create_palace(&self.state.data_root, palace)
333            .map_err(|e| ServiceError::internal(format!("create palace: {e:#}")))?;
334        // Issue #228: keep the in-memory palace-name cache in sync so writes
335        // to this palace can resolve `Palace.name` without a disk walk.
336        self.state.palace_names.insert(name.clone(), name.clone());
337        self.state.emit(DaemonEvent::PalaceCreated {
338            id: name.clone(),
339            name: name.clone(),
340            source,
341        });
342        Ok(name)
343    }
344
345    /// Delete a palace from disk, optionally rejecting non-empty palaces.
346    ///
347    /// Why: Issue #180 — operators need a way to drop an entire palace
348    /// without going through drawer-by-drawer deletion. Defaulting to a
349    /// "must be empty" guard prevents fat-finger destruction of populated
350    /// palaces; `force=true` is the explicit opt-in to the destructive path.
351    /// What: 1) confirms the palace exists on disk (else `NotFound`),
352    /// 2) when `!force`, opens the palace and returns `Conflict` if the open
353    /// fails, if its drawer table loaded degraded, while a legacy `kg.db` holds
354    /// drawers `kg.redb` lacks or triples or a `.v2-incompatible` file remains
355    /// (#8434; checked first, and its message carries no force hint), or if it
356    /// has drawers, 3) drops the in-memory registry entry so
357    /// future opens hit the (now-missing) disk state, 4) removes
358    /// `<data_root>/<palace_id>/` recursively via `tokio::fs::remove_dir_all`,
359    /// and 5) emits an aggregate `StatusChanged` so dashboards refresh.
360    /// Test: `delete_palace_removes_dir_when_empty`,
361    /// `delete_palace_refuses_when_drawers_present`,
362    /// `delete_palace_refuses_while_legacy_kg_holds_unimported_drawers`,
363    /// `delete_palace_refuses_when_open_fails_or_load_is_degraded`,
364    /// `delete_palace_force_removes_populated_palace`,
365    /// `delete_palace_returns_not_found_for_missing_id` in `web::tests`.
366    pub async fn delete_palace(&self, palace_id: &str, force: bool) -> ServiceResult<()> {
367        let palaces = PalaceRegistry::list_palaces(&self.state.data_root)
368            .map_err(|e| ServiceError::internal(format!("list palaces: {e:#}")))?;
369        if !palaces.iter().any(|p| p.id.0 == palace_id) {
370            return Err(ServiceError::not_found(format!(
371                "palace not found: {palace_id}"
372            )));
373        }
374        if !force {
375            // Open the palace just long enough to count its drawers; we don't
376            // hold the handle past this check because the caller is about to
377            // delete the on-disk directory.
378            // #8434: fail closed — a palace that cannot be opened, or whose
379            // drawer table loaded degraded, has not been shown to be empty.
380            let handle = self
381                .state
382                .registry
383                .open_palace(&self.state.data_root, &PalaceId::new(palace_id))
384                .map_err(|e| {
385                    ServiceError::conflict(format!(
386                        "Palace could not be opened to confirm it is empty ({e:#}); refusing \
387                         to delete"
388                    ))
389                })?;
390            if handle.drawer_load_degraded {
391                return Err(ServiceError::conflict(
392                    "Palace drawer table loaded degraded, so it cannot be confirmed empty; \
393                     refusing to delete",
394                ));
395            }
396            let has_drawers = !handle.drawers.read().is_empty();
397            // #8434: dedupe legacy rows against kg.redb, the persisted store.
398            let live = handle.kg.load_drawer_ids().map_err(|e| {
399                ServiceError::conflict(format!(
400                    "Palace drawer ids could not be read ({e:#}); refusing to delete"
401                ))
402            })?;
403            drop(handle);
404            // #8434: "0 live drawers" is not "empty" while a legacy kg.db or a
405            // quarantined redb 2.x store holds data the live store never saw.
406            // Checked before the has-drawers conflict, whose force hint would
407            // steer a caller into destroying that data unimported.
408            let dir = self.state.data_root.join(palace_id);
409            let unaccounted = tokio::task::spawn_blocking(move || {
410                crate::commands::legacy_kg::unaccounted_legacy_data(&dir, &live)
411            })
412            .await
413            .map_err(|e| ServiceError::internal(format!("legacy data check: {e}")))?;
414            if let Some(reason) = unaccounted {
415                // #8434: no force hint — force destroys this data unimported.
416                return Err(ServiceError::conflict(format!(
417                    "Palace holds legacy data ({reason}); run `trusty-memory palace legacy-kg \
418                     {palace_id}` to review it"
419                )));
420            }
421            if has_drawers {
422                return Err(ServiceError::conflict(
423                    "Palace has drawers; pass force=true to delete",
424                ));
425            }
426        }
427        // Drop the cached `Arc<PalaceHandle>` and gap cache before unlinking
428        // the directory so subsequent reads can't be served from the stale
429        // in-memory state. The registry's `remove` is a no-op when the entry
430        // is absent (lazy-open palaces that no caller has touched yet).
431        self.state.registry.remove(&PalaceId::new(palace_id));
432        // #4639: drop the cached chat_sessions.redb handle too — otherwise the
433        // fd survives `remove_dir_all` and pins the deleted inode forever.
434        self.state.session_stores.remove(palace_id);
435        // Issue #228: drop the palace-name cache entry so future writes never
436        // resolve to a stale label.
437        self.state.palace_names.remove(palace_id);
438        let palace_dir = self.state.data_root.join(palace_id);
439        tokio::fs::remove_dir_all(&palace_dir).await.map_err(|e| {
440            ServiceError::internal(format!("remove palace dir {}: {e}", palace_dir.display()))
441        })?;
442        // Recompute aggregate totals so dashboards drop the deleted palace's
443        // counts. There's no dedicated `PalaceDeleted` event variant yet;
444        // `StatusChanged` is enough to keep the UI in sync.
445        self.state.emit(self.aggregate_status_event());
446        Ok(())
447    }
448
449    /// Rename a palace's display name without touching its data.
450    ///
451    /// Why: Operators need to fix typos and rebrand palaces without dropping
452    /// the underlying drawers / vectors / KG. The palace id (the directory
453    /// name on disk) is immutable — only the human-readable `name` field in
454    /// `palace.json` changes — so cached `PalaceHandle`s stay valid and no
455    /// registry invalidation is required.
456    /// What: 1) loads the palace via `PalaceStore::load_palace` (404 when the
457    /// directory or `palace.json` is genuinely missing; a probe that cannot
458    /// determine whether it is there is a 500, not a 404 — #5549), 2) trims the
459    /// new name and
460    /// returns `BadRequest` when empty, 3) mutates `palace.name` and writes
461    /// the metadata back through the atomic `PalaceStore::save_palace`
462    /// (tmp file + rename), 4) emits an aggregate `StatusChanged` so
463    /// dashboards re-render the relabelled palace, 5) returns the updated
464    /// palace as JSON (enriched with the live handle stats, so callers see
465    /// drawer/vector/KG counts in the same shape as `GET /palaces/{id}`).
466    /// Test: `update_palace_name_renames_palace`,
467    /// `update_palace_name_rejects_empty_name`,
468    /// `update_palace_name_returns_not_found_for_missing_id` in `web::tests`.
469    pub async fn update_palace_name(&self, palace_id: &str, name: &str) -> Result<Value> {
470        let trimmed = name.trim();
471        if trimmed.is_empty() {
472            return Err(anyhow!("name must be non-empty after trimming"));
473        }
474        let palace_dir = self.state.data_root.join(palace_id);
475        let mut palace = trusty_common::memory_core::store::PalaceStore::load_palace(&palace_dir)
476            .map_err(|e| {
477            // #5549: only a genuine absence may be reported as "not found".
478            if matches!(&e, PalaceStoreError::NotFound(_)) {
479                anyhow!("palace not found: {palace_id} ({e})")
480            } else {
481                anyhow!("cannot load palace {palace_id}: {e}")
482            }
483        })?;
484        palace.name = trimmed.to_string();
485        trusty_common::memory_core::store::PalaceStore::save_palace(&palace)
486            .with_context(|| format!("save palace metadata for {palace_id}"))?;
487        // Issue #228: refresh the in-memory name cache so subsequent writes
488        // surface the new label without a disk walk.
489        self.state
490            .palace_names
491            .insert(palace_id.to_string(), trimmed.to_string());
492        let handle = self
493            .state
494            .registry
495            .open_palace(&self.state.data_root, &palace.id)
496            .ok();
497        // #7106: enrichment runs on the blocking pool.
498        let info = palace_info_blocking(&palace, handle).await;
499        self.state.emit(self.aggregate_status_event());
500        serde_json::to_value(info).context("serialize palace info")
501    }
502
503    /// Typed variant of [`Self::update_palace_name`] used by the HTTP handler.
504    ///
505    /// Why: HTTP needs to distinguish 400 (empty name) from 404 (missing
506    /// palace) so the right status code is emitted; the chat / MCP tool
507    /// only cares about a `Result<Value>` because both errors are surfaced
508    /// as opaque MCP error strings. Keeping a typed variant alongside the
509    /// untyped one keeps the wire shape correct on both surfaces without
510    /// asking either caller to parse error strings.
511    /// What: same as [`Self::update_palace_name`] but returns
512    /// `ServiceError::BadRequest` for empty names and `ServiceError::NotFound`
513    /// for palace metadata that is genuinely absent. Metadata whose presence
514    /// cannot be determined — a denied or transient stat — is
515    /// `ServiceError::Internal`: a 404 would tell the client the palace does
516    /// not exist when nobody established that (#5549, ADR-0045).
517    /// Test: `update_palace_name_renames_palace`,
518    /// `update_palace_name_rejects_empty_name`,
519    /// `update_palace_name_returns_not_found_for_missing_id`,
520    /// `update_palace_name_reports_an_unstattable_palace_as_internal`.
521    pub async fn update_palace_name_typed(
522        &self,
523        palace_id: &str,
524        name: &str,
525    ) -> ServiceResult<Value> {
526        let trimmed = name.trim();
527        if trimmed.is_empty() {
528            return Err(ServiceError::bad_request(
529                "name must be non-empty after trimming",
530            ));
531        }
532        let palace_dir = self.state.data_root.join(palace_id);
533        let mut palace = trusty_common::memory_core::store::PalaceStore::load_palace(&palace_dir)
534            .map_err(|e| {
535            // #5549: `not_found` on every variant told the client the palace
536            // does not exist for a stat we were merely denied.
537            if matches!(&e, PalaceStoreError::NotFound(_)) {
538                ServiceError::not_found(format!("palace not found: {palace_id} ({e})"))
539            } else {
540                ServiceError::internal(format!("cannot load palace {palace_id}: {e}"))
541            }
542        })?;
543        palace.name = trimmed.to_string();
544        trusty_common::memory_core::store::PalaceStore::save_palace(&palace).map_err(|e| {
545            ServiceError::internal(format!("save palace metadata for {palace_id}: {e}"))
546        })?;
547        // Issue #228: refresh the in-memory name cache so subsequent writes
548        // surface the new label without a disk walk.
549        self.state
550            .palace_names
551            .insert(palace_id.to_string(), trimmed.to_string());
552        let handle = self
553            .state
554            .registry
555            .open_palace(&self.state.data_root, &palace.id)
556            .ok();
557        // #7106: enrichment runs on the blocking pool.
558        let info = palace_info_blocking(&palace, handle).await;
559        self.state.emit(self.aggregate_status_event());
560        serde_json::to_value(info)
561            .map_err(|e| ServiceError::internal(format!("serialize palace info: {e}")))
562    }
563
564    /// Look up a single palace by id and enrich with live handle stats.
565    ///
566    /// Why: distinct 404 vs. 500 path is needed by both HTTP and chat callers.
567    /// What: returns `NotFound` when the id is unknown, otherwise a fully
568    /// populated `PalaceInfo`.
569    /// Test: indirectly via `health_endpoint_round_trip_with_palace_is_ok`.
570    pub async fn get_palace(&self, id: &str) -> ServiceResult<PalaceInfo> {
571        let palaces = PalaceRegistry::list_palaces(&self.state.data_root)
572            .map_err(|e| ServiceError::internal(format!("list palaces: {e:#}")))?;
573        let palace = palaces
574            .into_iter()
575            .find(|p| p.id.0 == id)
576            .ok_or_else(|| ServiceError::not_found(format!("palace not found: {id}")))?;
577        let handle = self
578            .state
579            .registry
580            .open_palace(&self.state.data_root, &palace.id)
581            .ok();
582        // #7106: enrichment runs on the blocking pool.
583        Ok(palace_info_blocking(&palace, handle).await)
584    }
585
586    // -----------------------------------------------------------------
587    // Drawers
588    // -----------------------------------------------------------------
589
590    /// List drawers in a palace with optional room/tag filters and pagination.
591    ///
592    /// Why: deduplicates the open-handle + listing path between HTTP and chat,
593    /// and (issue #184) lets the TUI activity panel page through drawers in
594    /// creation-date order without breaking the importance-sorted default the
595    /// legacy callers rely on.
596    /// What: opens the palace handle, fetches a window of drawers, optionally
597    /// re-sorts by `created_at` descending when `sort = "created_desc"`
598    /// (leaving the importance-desc default untouched), then drops the
599    /// leading `offset` rows and keeps `limit`. For `created_desc` the
600    /// window must cover the full filtered set (otherwise the importance
601    /// pre-sort hides truly-recent low-importance drawers), so the window
602    /// is widened to a sane ceiling (`MAX_DRAWER_WINDOW`); the default
603    /// importance path keeps a tight `limit+offset` window.
604    /// Returns the serialised JSON array.
605    /// Test: `service::tests::list_drawers_creates_desc_paginates`.
606    pub async fn list_drawers(&self, id: &str, q: ListDrawersQuery) -> ServiceResult<Value> {
607        const MAX_DRAWER_WINDOW: usize = 10_000;
608        let handle = self.open_handle(id)?;
609        let room = q.room.as_deref().map(RoomType::parse);
610        let limit = q.limit.unwrap_or(50);
611        let offset = q.offset.unwrap_or(0);
612        let by_created = matches!(q.sort.as_deref(), Some("created_desc"));
613        // For created_desc the importance pre-sort would hide low-importance
614        // drawers that happen to be the most recent, so we need to fetch the
615        // full filtered set (capped at MAX_DRAWER_WINDOW). For importance
616        // ordering the legacy `limit + offset` window is sufficient.
617        let window = if by_created {
618            MAX_DRAWER_WINDOW
619        } else {
620            limit.saturating_add(offset).min(MAX_DRAWER_WINDOW)
621        };
622        let mut drawers = handle.list_drawers(room, q.tag.clone(), window);
623        if by_created {
624            drawers.sort_by_key(|d| std::cmp::Reverse(d.created_at));
625        }
626        let page: Vec<_> = drawers.into_iter().skip(offset).take(limit).collect();
627        // Issue #202: enrich every row with a short `snippet` derived from
628        // the drawer's content so the TUI activity panel can render a
629        // glanceable summary without re-parsing the full body. The
630        // snippet is whitespace-collapsed and bounded at
631        // `DRAWER_SNIPPET_MAX_CHARS` (60) — shorter than the SSE preview
632        // because the activity panel renders it on a single narrow row.
633        let payload: Vec<Value> = page
634            .into_iter()
635            .map(|drawer| {
636                let snippet = drawer_snippet(drawer.content());
637                let mut value = serde_json::to_value(&drawer).unwrap_or_else(|_| json!({}));
638                if let Value::Object(ref mut map) = value {
639                    // `null` when the drawer has no usable content so
640                    // clients can distinguish "no body" from "empty body
641                    // after whitespace collapse".
642                    let snippet_value = if snippet.is_empty() {
643                        Value::Null
644                    } else {
645                        Value::String(snippet)
646                    };
647                    map.insert("snippet".to_string(), snippet_value);
648                }
649                value
650            })
651            .collect();
652        Ok(Value::Array(payload))
653    }
654
655    /// Store a new drawer and emit the matching activity events.
656    ///
657    /// Why: HTTP and chat both need the auto-KG-extraction follow-up; this
658    /// method keeps that side-effect chain in one place.
659    /// What: opens the palace, stores the drawer via
660    /// `PalaceHandle::remember_with_options` (issue #3225: `body.force`
661    /// threads through as `RememberOptions::force`, letting a caller bypass
662    /// the QUALITY gates only — `allow_secret_like` is left at its default
663    /// `false`, so secret detection always still runs, `force` or not),
664    /// emits `DrawerAdded` + `StatusChanged`, then triggers
665    /// `tools::auto_extract_and_assert`. Returns the new drawer id.
666    /// Test: `http_create_drawer_runs_auto_kg_extraction`,
667    /// `create_drawer_rejects_json_content_without_force`,
668    /// `create_drawer_force_bypasses_quality_gate_for_json_content`.
669    pub async fn create_drawer(
670        &self,
671        id: &str,
672        body: CreateDrawerBody,
673        creator: CreatorInfo,
674        source: ActivitySource,
675    ) -> ServiceResult<Uuid> {
676        let handle = self.open_handle(id)?;
677        let room = body
678            .room
679            .as_deref()
680            .map(RoomType::parse)
681            .unwrap_or(RoomType::General);
682        let importance = body.importance.unwrap_or(0.5);
683        let force = body.force.unwrap_or(false);
684        let content_preview = drawer_content_preview(&body.content);
685        let mut tags_with_creator = body.tags;
686        // Issue #202: project a bare-UUID session tag (when the caller
687        // passed one in the request body) into the reserved
688        // `creator:session=<first-8>` slot so the activity panel can
689        // surface session attribution without bespoke parsing.
690        if let Some(session_tag) = crate::attribution::session_tag_from_tags(&tags_with_creator) {
691            tags_with_creator.push(session_tag);
692        }
693        creator.merge_into(&mut tags_with_creator);
694        let content_for_kg = body.content.clone();
695        let tags_for_kg = tags_with_creator.clone();
696        let room_label_for_kg = crate::tools::room_label(&room);
697        let drawer_id = handle
698            .remember_with_options(
699                body.content,
700                room,
701                tags_with_creator,
702                importance,
703                RememberOptions {
704                    force,
705                    ..Default::default()
706                },
707            )
708            .await
709            .map_err(|e| ServiceError::internal(format!("remember: {e:#}")))?;
710        let drawer_count = handle.drawers.read().len();
711        // Issue #228: resolve from the in-memory cache instead of re-walking
712        // the data root on every HTTP `create_drawer` call. Same cache the
713        // MCP `lookup_palace_name` helper consults.
714        let palace_name = self
715            .state
716            .palace_names
717            .get(id)
718            .map(|entry| entry.value().clone())
719            .unwrap_or_else(|| id.to_string());
720        self.state.emit(DaemonEvent::DrawerAdded {
721            palace_id: id.to_string(),
722            palace_name,
723            drawer_count,
724            timestamp: chrono::Utc::now(),
725            content_preview,
726            source,
727        });
728        // Issue #228: do NOT emit `StatusChanged` on every drawer create —
729        // the periodic ticker (`run_http_on`) refreshes aggregate totals on
730        // a fixed cadence so dashboards stay current without an O(N palaces)
731        // recompute on the write hot path.
732        crate::tools::auto_extract_and_assert(
733            &handle,
734            drawer_id,
735            &content_for_kg,
736            &tags_for_kg,
737            room_label_for_kg.as_deref(),
738        )
739        .await;
740        Ok(drawer_id)
741    }
742
743    /// Forget (delete) a drawer and emit the matching events.
744    ///
745    /// Why: same dedup story as `create_drawer`. #5231: `DELETE` on a drawer id
746    /// that was never stored used to answer `204 No Content`, the same as a
747    /// real delete — this now 404s, matching `delete_palace`.
748    /// What: parses the drawer UUID, calls `PalaceHandle::forget`, deletes the
749    /// drawer's BM25 document, maps `ForgetOutcome::NotFound` to
750    /// `ServiceError::not_found`, and emits `DrawerDeleted` only when a drawer
751    /// was actually removed. #5053: the lexical delete runs on this path for
752    /// the same reason it runs on the MCP one — `HTTP DELETE` and
753    /// `memory_forget` remove the same drawer, and the backfill indexes it
754    /// whichever way it was written, so a lexical copy left here is the same
755    /// stale document.
756    /// Test: `delete_drawer_404s_for_an_unknown_drawer_id`;
757    /// `tests/bm25_forget_delete.rs` covers the deletion contract itself.
758    pub async fn delete_drawer(
759        &self,
760        id: &str,
761        drawer_id: &str,
762        source: ActivitySource,
763    ) -> ServiceResult<()> {
764        let handle = self.open_handle(id)?;
765        let uuid = Uuid::parse_str(drawer_id)
766            .map_err(|_| ServiceError::bad_request("drawer_id must be a UUID"))?;
767        let outcome = handle
768            .forget(uuid)
769            .await
770            .map_err(|e| ServiceError::internal(format!("forget: {e:#}")))?;
771        // #5053: a drawer the user deleted must stop matching lexical queries.
772        crate::tools::bm25::bm25_delete_document(&self.state, handle.id.as_str(), uuid)
773            .await
774            .map_err(|e| ServiceError::internal(format!("{e:#}")))?;
775        if !outcome.is_deleted() {
776            return Err(ServiceError::not_found(format!(
777                "drawer '{drawer_id}' not found in palace '{id}'"
778            )));
779        }
780        let drawer_count = handle.drawers.read().len();
781        self.state.emit(DaemonEvent::DrawerDeleted {
782            palace_id: id.to_string(),
783            drawer_count,
784            source,
785        });
786        // Issue #228: skip the per-write `StatusChanged` emit — the
787        // periodic ticker handles aggregate roll-ups.
788        Ok(())
789    }
790
791    // -----------------------------------------------------------------
792    // Recall
793    // -----------------------------------------------------------------
794
795    /// Per-palace recall (semantic search), optionally with deep retrieval.
796    ///
797    /// Why: HTTP and chat tools both perform the same fan-out logic.
798    /// What: opens the palace handle and dispatches to the shallow or deep
799    /// recall helper. Returns a JSON array of flattened drawer rows (the
800    /// `recall_entry_json` shape from issue #69).
801    /// Test: `recall_entry_json_hoists_drawer_fields`.
802    pub async fn recall(
803        &self,
804        id: &str,
805        query: &str,
806        top_k: usize,
807        deep: bool,
808    ) -> ServiceResult<Value> {
809        let handle = self.open_handle(id)?;
810        let mut results = if deep {
811            recall_deep_with_default_embedder(&handle, query, top_k).await
812        } else {
813            recall_with_default_embedder(&handle, query, top_k).await
814        }
815        .map_err(|e| ServiceError::internal(format!("recall: {e:#}")))?;
816        // #5036: the lexical lane, on the path the UserPromptSubmit hook
817        // actually takes. `handle_memory_recall` has run vector and BM25 in
818        // parallel and RRF-fused them since #156; this route reached
819        // `retrieval::layers` directly and was vector-only, so a prompt with no
820        // lexical counterweight retrieved by vector centroid alone.
821        //
822        // Keyed on the RESOLVED palace id, never the caller's slug —
823        // `open_handle` follows aliases, and the corpus the backfill wrote
824        // belongs to the resolved palace.
825        //
826        // Reuses `fuse_bm25_into_recall` rather than deriving a second scorer:
827        // it only BOOSTS drawers the vector lane already returned and never
828        // promotes a BM25-only hit, so it has no scaling constant that can
829        // degenerate when the surviving set is empty — the failure that folded
830        // three earlier attempts at this wiring.
831        if let Some(hits) =
832            crate::tools::bm25::bm25_search_optional(&self.state, handle.id.as_str(), query, top_k)
833                .await
834        {
835            crate::tools::bm25::fuse_bm25_into_recall(&mut results, &hits, top_k);
836        }
837        let payload: Vec<Value> = results.into_iter().map(recall_entry_json).collect();
838        Ok(json!(payload))
839    }
840
841    /// Cross-palace recall.
842    ///
843    /// Why: shared between `/api/v1/recall` and the `memory_recall_all` chat
844    /// tool. Encapsulating the open-everything-fanout-merge dance avoids
845    /// drift.
846    /// What: lists every palace, then streams them through `recall_streamed`
847    /// in bounded batches, delegating each batch to
848    /// `recall_across_palaces_with_default_embedder`. Returns a JSON array.
849    /// Why (issue #4637): unlike `list_palaces`/`status`, this route is NOT
850    /// converted to `peek()`. A cross-palace recall that answered from
851    /// cache-resident palaces only would silently omit ~98.9% of the corpus —
852    /// a wrong answer that looks like a right one, which is strictly worse
853    /// than a slow correct one. Every palace is still opened and still
854    /// searched; what changed is when.
855    /// Why (issue #7125): opening all of them AT ONCE made peak residency and
856    /// the post-call LRU residue both scale with the palace count. The batch
857    /// walk bounds the peak at `RECALL_PALACE_BATCH` and hands back everything
858    /// the query itself brought in, so the daemon's steady state after a
859    /// recall-all matches its steady state before one.
860    /// Test: indirectly via `recall_across_palaces_merges_results` and the
861    /// MCP `memory_recall_all` integration paths;
862    /// `open_palaces_blocking_opens_every_palace` pins that uncached palaces
863    /// are still searched; `recall_all_returns_open_palaces_to_baseline` pins
864    /// the residency bound.
865    pub async fn recall_all(&self, query: &str, top_k: usize, deep: bool) -> Value {
866        let palaces = match list_palaces_blocking(&self.state).await {
867            Ok(v) => v,
868            Err(e) => return json!({ "error": format!("{e:#}") }),
869        };
870        // #7125: stream the estate in batches instead of opening all of it.
871        let streamed = recall_streamed(
872            &self.state,
873            &palaces,
874            "recall_all",
875            top_k,
876            |handles| async move {
877                recall_across_palaces_with_default_embedder(&handles, query, top_k, deep).await
878            },
879        )
880        .await;
881        match streamed {
882            Ok(results) => json!(results
883                .into_iter()
884                .map(|r| json!({
885                    "palace_id": r.palace_id,
886                    "drawer_id": r.result.drawer.id.to_string(),
887                    "content": r.result.drawer.content(),
888                    "importance": r.result.drawer.importance,
889                    "tags": r.result.drawer.tags,
890                    "score": r.result.score,
891                    "layer": r.result.layer,
892                }))
893                .collect::<Vec<_>>()),
894            Err(e) => json!({ "error": format!("recall_across_palaces: {e:#}") }),
895        }
896    }
897}