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 trusty_common::memory_core::palace::{Palace, PalaceId, RoomType};
18use trusty_common::memory_core::retrieval::{
19 recall_across_palaces_with_default_embedder, recall_deep_with_default_embedder,
20 recall_with_default_embedder, RememberOptions,
21};
22use trusty_common::memory_core::PalaceRegistry;
23use uuid::Uuid;
24
25use super::helpers::{
26 collect_palace_stats, drawer_content_preview, drawer_snippet, is_reserved_system_palace,
27 list_palaces_blocking, open_palaces_blocking, palace_info_from, recall_entry_json,
28};
29use super::types::{
30 CreateDrawerBody, CreatePalaceBody, ListDrawersQuery, PalaceInfo, ServiceError, ServiceResult,
31 StatusPayload,
32};
33
34/// Hard cap on triples returned by the per-palace graph endpoint.
35pub(super) const KG_GRAPH_MAX_TRIPLES: usize = 5_000;
36
37// ---------------------------------------------------------------------------
38// MemoryService — pure business logic facade.
39// ---------------------------------------------------------------------------
40
41/// Wraps [`AppState`] and exposes one async method per logical operation.
42///
43/// Why: see module docs. Lets HTTP handlers stay thin and lets non-HTTP
44/// callers (chat tool dispatch, RPC bridges) reuse the same code paths.
45/// What: `Clone` (cheap — only the inner `AppState` is shared); construct
46/// with `MemoryService::new(state)`.
47/// Test: every method is covered by the corresponding handler test in
48/// `web::tests`.
49#[derive(Clone)]
50pub struct MemoryService {
51 pub(super) state: AppState,
52}
53
54impl MemoryService {
55 /// Construct a new service wrapper.
56 ///
57 /// Why: handlers cheaply re-wrap their `AppState` on every request; the
58 /// cost is just an `Arc` clone, so we don't bother caching the wrapper.
59 /// What: stores the `AppState` for later method calls.
60 /// Test: trivial — covered indirectly by every handler test.
61 pub fn new(state: AppState) -> Self {
62 Self { state }
63 }
64
65 /// Borrow the inner [`AppState`].
66 ///
67 /// Why: some handlers still need direct access (SSE broadcaster, session
68 /// store, etc.) while we incrementally extract code into the service.
69 /// What: returns a borrowed reference to the wrapped `AppState`.
70 /// Test: not directly tested; surface-level accessor.
71 pub fn state(&self) -> &AppState {
72 &self.state
73 }
74
75 // -----------------------------------------------------------------
76 // Status / config
77 // -----------------------------------------------------------------
78
79 /// Build the aggregate `/api/v1/status` payload.
80 ///
81 /// Why: dashboard widgets and the MCP `get_status` tool need the same
82 /// roll-up; centralising avoids drift between the two surfaces.
83 /// What: walks every persisted palace for `palace_count`, then sums
84 /// drawer/vector/triple counts across the cache-resident subset and
85 /// returns the [`StatusPayload`].
86 /// Why (issue #4637): this used to open every persisted palace to sum
87 /// those three counts. With 5,794 palaces on disk against a 64-slot LRU
88 /// that is ~5,730 cold opens of ~1s each — the endpoint measurably never
89 /// responded. `palace_count` still reflects the true on-disk total (the
90 /// directory walk is cheap and now runs on the blocking pool); the totals
91 /// cover only cache-resident palaces and say so via `cached_palace_count`.
92 /// Test: `status_endpoint_returns_payload`,
93 /// `status_does_not_open_uncached_palaces`.
94 pub async fn status(&self) -> StatusPayload {
95 // The `/status` endpoint is the one place we still want a disk view —
96 // an operator hitting this endpoint right after restart (before
97 // `load_palaces_from_disk` finishes) should still see every persisted
98 // palace counted, even if it isn't in the in-memory registry yet.
99 let palaces = list_palaces_blocking(&self.state).await.unwrap_or_default();
100 let palace_count = palaces.len();
101 // #4637: peek() not open_palace() — full-registry open is O(n) cold disk I/O
102 let stats = collect_palace_stats(&self.state, palaces.iter().map(|p| &p.id));
103 StatusPayload {
104 version: self.state.version.clone(),
105 palace_count,
106 default_palace: self.state.default_palace.clone(),
107 data_root: self.state.data_root.display().to_string(),
108 total_drawers: stats.total_drawers,
109 total_vectors: stats.total_vectors,
110 total_kg_triples: stats.total_kg_triples,
111 cached_palace_count: stats.cached_palace_count,
112 }
113 }
114
115 /// Compute the aggregate `StatusChanged` event used by SSE consumers.
116 ///
117 /// Why: mutating handlers — and the periodic status ticker — push a
118 /// refreshed status snapshot so dashboards stay in sync without an
119 /// extra `/api/v1/status` request.
120 /// Why (issue #228): this used to call `PalaceRegistry::list_palaces`
121 /// (a synchronous disk walk) + `open_palace` (more disk I/O on first
122 /// call) for every palace on every emit. Since every persisted palace
123 /// is already loaded into the in-memory registry by
124 /// `AppState::load_palaces_from_disk` at startup (and every `create_palace`
125 /// keeps it in sync), iterating the in-memory registry returns the same
126 /// counts without touching disk.
127 /// What: iterates `state.registry.list()` (a `DashMap` snapshot) and
128 /// sums the live handle stats via [`collect_palace_stats`]. Returns a
129 /// `DaemonEvent::StatusChanged`. Palaces that fail to resolve in the
130 /// registry (race during shutdown) are silently skipped — the next
131 /// emit will catch them.
132 /// Test: indirectly via SSE integration tests; the math is identical to
133 /// the disk-walk implementation and the `status_endpoint_returns_payload`
134 /// test still passes against `status()` (which keeps the disk view for
135 /// the dedicated endpoint).
136 pub fn aggregate_status_event(&self) -> DaemonEvent {
137 let ids: Vec<PalaceId> = self.state.registry.list();
138 let stats = collect_palace_stats(&self.state, ids.iter());
139 DaemonEvent::StatusChanged {
140 total_drawers: stats.total_drawers,
141 total_vectors: stats.total_vectors,
142 total_kg_triples: stats.total_kg_triples,
143 }
144 }
145
146 // -----------------------------------------------------------------
147 // Palaces
148 // -----------------------------------------------------------------
149
150 /// List every palace on disk, enriched with live handle stats.
151 ///
152 /// Why: shared between the HTTP handler and the chat tool dispatcher;
153 /// both want the same `PalaceInfo` shape. Issue #185 added the
154 /// reserved-prefix filter so internal "system" palaces (e.g. the
155 /// `__health_probe__` palace used by `/health`) never surface in the
156 /// admin UI, TUI, or any user-facing roster.
157 /// What: walks the registry, drops any palace whose id starts with the
158 /// reserved `__` prefix, and builds a `PalaceInfo` per remaining row.
159 /// Why (issue #4637): this used to call `open_palace` per row purely to
160 /// enrich it with counts. At 5,794 palaces against a 64-slot LRU that is
161 /// ~90 minutes of cold, blocking disk I/O inline on the async executor —
162 /// and it evicted the entire working set on every call. Rows now come
163 /// from `PalaceRegistry::peek` (zero I/O, no LRU promotion, mirroring the
164 /// #1924 fix in `console_metrics.rs`). Uncached rows carry `cached: false`
165 /// and zero counts; a client that needs live counts for one palace should
166 /// fetch `GET /api/v1/palaces/{id}`, which still opens it.
167 /// Test: `palace_list_includes_richer_counts`, `palace_list_includes_graph_counts`,
168 /// `health_probe_palace_is_invisible` (in `web::tests`),
169 /// `list_palaces_does_not_open_uncached_palaces`.
170 pub async fn list_palaces(&self) -> ServiceResult<Vec<PalaceInfo>> {
171 let palaces = list_palaces_blocking(&self.state)
172 .await
173 .map_err(|e| ServiceError::internal(format!("{e:#}")))?;
174 let mut out = Vec::with_capacity(palaces.len());
175 for p in palaces {
176 if is_reserved_system_palace(&p.id) {
177 continue;
178 }
179 // #4637: peek() not open_palace() — full-registry open is O(n) cold disk I/O
180 let handle = self.state.registry.peek(&p.id);
181 out.push(palace_info_from(&p, handle.as_ref()));
182 }
183 Ok(out)
184 }
185
186 /// Create a new palace and emit the corresponding activity event.
187 ///
188 /// Why: trims duplicated work between the HTTP handler and any future
189 /// non-HTTP creation flow.
190 /// What: validates the name, builds the `Palace` row, calls
191 /// `PalaceRegistry::create_palace`, and emits `PalaceCreated`. Returns
192 /// the new palace id.
193 /// Test: covered indirectly by `palace_list_includes_richer_counts` (which
194 /// posts a palace through the HTTP layer then reads it back).
195 pub async fn create_palace(
196 &self,
197 body: CreatePalaceBody,
198 source: ActivitySource,
199 ) -> ServiceResult<String> {
200 let name = body.name.trim().to_string();
201 if name.is_empty() {
202 return Err(ServiceError::bad_request("name is required"));
203 }
204 // Issue #88 / Change 2: enforce palace = project mapping for
205 // HTTP-originated palace creation. The validation cwd is, in order of
206 // preference:
207 // a. `body.cwd` — the caller explicitly supplied their project path
208 // (correct for any client that is not the daemon itself).
209 // b. `std::env::current_dir()` — daemon's own cwd, the pre-Change-2
210 // fallback (rarely meaningful when the daemon is launched from ~).
211 // This keeps older clients that omit `cwd` working without a breaking
212 // change, while letting pin-file-aware clients get accurate validation.
213 // spec-001: `force=true` lets an application bypass the project-slug
214 // gate so it can create palaces under arbitrary slugs (e.g. one per
215 // app/tenant for chat-session storage). The env-var bypass remains for
216 // test contexts; both short-circuit the same validation call.
217 //
218 // Issue #1714: `force=true` bypasses slug validation entirely, so it
219 // is gated behind the minimal authz seam in `crate::authz` before any
220 // other check runs. In the default single-tenant mode this is a
221 // no-op (unchanged behaviour); in multi-tenant mode it fails closed
222 // until a real capability check lands. See `crate::authz` module
223 // docs for the full design rationale.
224 let skip_enforcement =
225 std::env::var("TRUSTY_SKIP_PALACE_ENFORCEMENT").as_deref() == Ok("1");
226 if body.force {
227 crate::authz::authorize_force_palace_create(&self.state)
228 .map_err(|e| ServiceError::forbidden(e.to_string()))?;
229 }
230 if !skip_enforcement && !body.force {
231 let cwd = body
232 .cwd
233 .as_deref()
234 .map(std::path::Path::new)
235 .map(|p| p.to_path_buf())
236 .or_else(|| std::env::current_dir().ok())
237 .unwrap_or_else(|| self.state.data_root.clone());
238 crate::project_root::validate_palace_name(&name, &cwd)
239 .map_err(|e| ServiceError::bad_request(e.to_string()))?;
240 }
241 let id = PalaceId::new(&name);
242 let palace = Palace {
243 id: id.clone(),
244 name: name.clone(),
245 description: body.description.filter(|s| !s.is_empty()),
246 created_at: chrono::Utc::now(),
247 data_dir: self.state.data_root.join(&name),
248 };
249 self.state
250 .registry
251 .create_palace(&self.state.data_root, palace)
252 .map_err(|e| ServiceError::internal(format!("create palace: {e:#}")))?;
253 // Issue #228: keep the in-memory palace-name cache in sync so writes
254 // to this palace can resolve `Palace.name` without a disk walk.
255 self.state.palace_names.insert(name.clone(), name.clone());
256 self.state.emit(DaemonEvent::PalaceCreated {
257 id: name.clone(),
258 name: name.clone(),
259 source,
260 });
261 Ok(name)
262 }
263
264 /// Delete a palace from disk, optionally rejecting non-empty palaces.
265 ///
266 /// Why: Issue #180 — operators need a way to drop an entire palace
267 /// without going through drawer-by-drawer deletion. Defaulting to a
268 /// "must be empty" guard prevents fat-finger destruction of populated
269 /// palaces; `force=true` is the explicit opt-in to the destructive path.
270 /// What: 1) confirms the palace exists on disk (else `NotFound`),
271 /// 2) when `!force`, lists drawers via the live handle and returns
272 /// `BadRequest("Palace has drawers; pass force=true to delete")` if
273 /// the palace is non-empty, 3) drops the in-memory registry entry so
274 /// future opens hit the (now-missing) disk state, 4) removes
275 /// `<data_root>/<palace_id>/` recursively via `tokio::fs::remove_dir_all`,
276 /// and 5) emits an aggregate `StatusChanged` so dashboards refresh.
277 /// Test: `delete_palace_removes_dir_when_empty`,
278 /// `delete_palace_refuses_when_drawers_present`,
279 /// `delete_palace_force_removes_populated_palace`,
280 /// `delete_palace_returns_not_found_for_missing_id` in `web::tests`.
281 pub async fn delete_palace(&self, palace_id: &str, force: bool) -> ServiceResult<()> {
282 let palaces = PalaceRegistry::list_palaces(&self.state.data_root)
283 .map_err(|e| ServiceError::internal(format!("list palaces: {e:#}")))?;
284 if !palaces.iter().any(|p| p.id.0 == palace_id) {
285 return Err(ServiceError::not_found(format!(
286 "palace not found: {palace_id}"
287 )));
288 }
289 if !force {
290 // Open the palace just long enough to count its drawers; we don't
291 // hold the handle past this check because the caller is about to
292 // delete the on-disk directory.
293 if let Ok(handle) = self
294 .state
295 .registry
296 .open_palace(&self.state.data_root, &PalaceId::new(palace_id))
297 {
298 if !handle.drawers.read().is_empty() {
299 return Err(ServiceError::conflict(
300 "Palace has drawers; pass force=true to delete",
301 ));
302 }
303 }
304 }
305 // Drop the cached `Arc<PalaceHandle>` and gap cache before unlinking
306 // the directory so subsequent reads can't be served from the stale
307 // in-memory state. The registry's `remove` is a no-op when the entry
308 // is absent (lazy-open palaces that no caller has touched yet).
309 self.state.registry.remove(&PalaceId::new(palace_id));
310 // #4639: drop the cached chat_sessions.redb handle too — otherwise the
311 // fd survives `remove_dir_all` and pins the deleted inode forever.
312 self.state.session_stores.remove(palace_id);
313 // Issue #228: drop the palace-name cache entry so future writes never
314 // resolve to a stale label.
315 self.state.palace_names.remove(palace_id);
316 let palace_dir = self.state.data_root.join(palace_id);
317 tokio::fs::remove_dir_all(&palace_dir).await.map_err(|e| {
318 ServiceError::internal(format!("remove palace dir {}: {e}", palace_dir.display()))
319 })?;
320 // Recompute aggregate totals so dashboards drop the deleted palace's
321 // counts. There's no dedicated `PalaceDeleted` event variant yet;
322 // `StatusChanged` is enough to keep the UI in sync.
323 self.state.emit(self.aggregate_status_event());
324 Ok(())
325 }
326
327 /// Rename a palace's display name without touching its data.
328 ///
329 /// Why: Operators need to fix typos and rebrand palaces without dropping
330 /// the underlying drawers / vectors / KG. The palace id (the directory
331 /// name on disk) is immutable — only the human-readable `name` field in
332 /// `palace.json` changes — so cached `PalaceHandle`s stay valid and no
333 /// registry invalidation is required.
334 /// What: 1) loads the palace via `PalaceStore::load_palace` (404 when the
335 /// directory or `palace.json` is missing), 2) trims the new name and
336 /// returns `BadRequest` when empty, 3) mutates `palace.name` and writes
337 /// the metadata back through the atomic `PalaceStore::save_palace`
338 /// (tmp file + rename), 4) emits an aggregate `StatusChanged` so
339 /// dashboards re-render the relabelled palace, 5) returns the updated
340 /// palace as JSON (enriched with the live handle stats, so callers see
341 /// drawer/vector/KG counts in the same shape as `GET /palaces/{id}`).
342 /// Test: `update_palace_name_renames_palace`,
343 /// `update_palace_name_rejects_empty_name`,
344 /// `update_palace_name_returns_not_found_for_missing_id` in `web::tests`.
345 pub async fn update_palace_name(&self, palace_id: &str, name: &str) -> Result<Value> {
346 let trimmed = name.trim();
347 if trimmed.is_empty() {
348 return Err(anyhow!("name must be non-empty after trimming"));
349 }
350 let palace_dir = self.state.data_root.join(palace_id);
351 let mut palace = trusty_common::memory_core::store::PalaceStore::load_palace(&palace_dir)
352 .map_err(|e| anyhow!("palace not found: {palace_id} ({e})"))?;
353 palace.name = trimmed.to_string();
354 trusty_common::memory_core::store::PalaceStore::save_palace(&palace)
355 .with_context(|| format!("save palace metadata for {palace_id}"))?;
356 // Issue #228: refresh the in-memory name cache so subsequent writes
357 // surface the new label without a disk walk.
358 self.state
359 .palace_names
360 .insert(palace_id.to_string(), trimmed.to_string());
361 let handle = self
362 .state
363 .registry
364 .open_palace(&self.state.data_root, &palace.id)
365 .ok();
366 let info = palace_info_from(&palace, handle.as_ref());
367 self.state.emit(self.aggregate_status_event());
368 serde_json::to_value(info).context("serialize palace info")
369 }
370
371 /// Typed variant of [`Self::update_palace_name`] used by the HTTP handler.
372 ///
373 /// Why: HTTP needs to distinguish 400 (empty name) from 404 (missing
374 /// palace) so the right status code is emitted; the chat / MCP tool
375 /// only cares about a `Result<Value>` because both errors are surfaced
376 /// as opaque MCP error strings. Keeping a typed variant alongside the
377 /// untyped one keeps the wire shape correct on both surfaces without
378 /// asking either caller to parse error strings.
379 /// What: same as [`Self::update_palace_name`] but returns
380 /// `ServiceError::BadRequest` for empty names and
381 /// `ServiceError::NotFound` for missing palace metadata.
382 /// Test: `update_palace_name_renames_palace`,
383 /// `update_palace_name_rejects_empty_name`,
384 /// `update_palace_name_returns_not_found_for_missing_id`.
385 pub async fn update_palace_name_typed(
386 &self,
387 palace_id: &str,
388 name: &str,
389 ) -> ServiceResult<Value> {
390 let trimmed = name.trim();
391 if trimmed.is_empty() {
392 return Err(ServiceError::bad_request(
393 "name must be non-empty after trimming",
394 ));
395 }
396 let palace_dir = self.state.data_root.join(palace_id);
397 let mut palace = trusty_common::memory_core::store::PalaceStore::load_palace(&palace_dir)
398 .map_err(|e| {
399 ServiceError::not_found(format!("palace not found: {palace_id} ({e})"))
400 })?;
401 palace.name = trimmed.to_string();
402 trusty_common::memory_core::store::PalaceStore::save_palace(&palace).map_err(|e| {
403 ServiceError::internal(format!("save palace metadata for {palace_id}: {e}"))
404 })?;
405 // Issue #228: refresh the in-memory name cache so subsequent writes
406 // surface the new label without a disk walk.
407 self.state
408 .palace_names
409 .insert(palace_id.to_string(), trimmed.to_string());
410 let handle = self
411 .state
412 .registry
413 .open_palace(&self.state.data_root, &palace.id)
414 .ok();
415 let info = palace_info_from(&palace, handle.as_ref());
416 self.state.emit(self.aggregate_status_event());
417 serde_json::to_value(info)
418 .map_err(|e| ServiceError::internal(format!("serialize palace info: {e}")))
419 }
420
421 /// Look up a single palace by id and enrich with live handle stats.
422 ///
423 /// Why: distinct 404 vs. 500 path is needed by both HTTP and chat callers.
424 /// What: returns `NotFound` when the id is unknown, otherwise a fully
425 /// populated `PalaceInfo`.
426 /// Test: indirectly via `health_endpoint_round_trip_with_palace_is_ok`.
427 pub async fn get_palace(&self, id: &str) -> ServiceResult<PalaceInfo> {
428 let palaces = PalaceRegistry::list_palaces(&self.state.data_root)
429 .map_err(|e| ServiceError::internal(format!("list palaces: {e:#}")))?;
430 let palace = palaces
431 .into_iter()
432 .find(|p| p.id.0 == id)
433 .ok_or_else(|| ServiceError::not_found(format!("palace not found: {id}")))?;
434 let handle = self
435 .state
436 .registry
437 .open_palace(&self.state.data_root, &palace.id)
438 .ok();
439 Ok(palace_info_from(&palace, handle.as_ref()))
440 }
441
442 // -----------------------------------------------------------------
443 // Drawers
444 // -----------------------------------------------------------------
445
446 /// List drawers in a palace with optional room/tag filters and pagination.
447 ///
448 /// Why: deduplicates the open-handle + listing path between HTTP and chat,
449 /// and (issue #184) lets the TUI activity panel page through drawers in
450 /// creation-date order without breaking the importance-sorted default the
451 /// legacy callers rely on.
452 /// What: opens the palace handle, fetches a window of drawers, optionally
453 /// re-sorts by `created_at` descending when `sort = "created_desc"`
454 /// (leaving the importance-desc default untouched), then drops the
455 /// leading `offset` rows and keeps `limit`. For `created_desc` the
456 /// window must cover the full filtered set (otherwise the importance
457 /// pre-sort hides truly-recent low-importance drawers), so the window
458 /// is widened to a sane ceiling (`MAX_DRAWER_WINDOW`); the default
459 /// importance path keeps a tight `limit+offset` window.
460 /// Returns the serialised JSON array.
461 /// Test: `service::tests::list_drawers_creates_desc_paginates`.
462 pub async fn list_drawers(&self, id: &str, q: ListDrawersQuery) -> ServiceResult<Value> {
463 const MAX_DRAWER_WINDOW: usize = 10_000;
464 let handle = self.open_handle(id)?;
465 let room = q.room.as_deref().map(RoomType::parse);
466 let limit = q.limit.unwrap_or(50);
467 let offset = q.offset.unwrap_or(0);
468 let by_created = matches!(q.sort.as_deref(), Some("created_desc"));
469 // For created_desc the importance pre-sort would hide low-importance
470 // drawers that happen to be the most recent, so we need to fetch the
471 // full filtered set (capped at MAX_DRAWER_WINDOW). For importance
472 // ordering the legacy `limit + offset` window is sufficient.
473 let window = if by_created {
474 MAX_DRAWER_WINDOW
475 } else {
476 limit.saturating_add(offset).min(MAX_DRAWER_WINDOW)
477 };
478 let mut drawers = handle.list_drawers(room, q.tag.clone(), window);
479 if by_created {
480 drawers.sort_by_key(|d| std::cmp::Reverse(d.created_at));
481 }
482 let page: Vec<_> = drawers.into_iter().skip(offset).take(limit).collect();
483 // Issue #202: enrich every row with a short `snippet` derived from
484 // the drawer's content so the TUI activity panel can render a
485 // glanceable summary without re-parsing the full body. The
486 // snippet is whitespace-collapsed and bounded at
487 // `DRAWER_SNIPPET_MAX_CHARS` (60) — shorter than the SSE preview
488 // because the activity panel renders it on a single narrow row.
489 let payload: Vec<Value> = page
490 .into_iter()
491 .map(|drawer| {
492 let snippet = drawer_snippet(&drawer.content);
493 let mut value = serde_json::to_value(&drawer).unwrap_or_else(|_| json!({}));
494 if let Value::Object(ref mut map) = value {
495 // `null` when the drawer has no usable content so
496 // clients can distinguish "no body" from "empty body
497 // after whitespace collapse".
498 let snippet_value = if snippet.is_empty() {
499 Value::Null
500 } else {
501 Value::String(snippet)
502 };
503 map.insert("snippet".to_string(), snippet_value);
504 }
505 value
506 })
507 .collect();
508 Ok(Value::Array(payload))
509 }
510
511 /// Store a new drawer and emit the matching activity events.
512 ///
513 /// Why: HTTP and chat both need the auto-KG-extraction follow-up; this
514 /// method keeps that side-effect chain in one place.
515 /// What: opens the palace, stores the drawer via
516 /// `PalaceHandle::remember_with_options` (issue #3225: `body.force`
517 /// threads through as `RememberOptions::force`, letting a caller bypass
518 /// the QUALITY gates only — `allow_secret_like` is left at its default
519 /// `false`, so secret detection always still runs, `force` or not),
520 /// emits `DrawerAdded` + `StatusChanged`, then triggers
521 /// `tools::auto_extract_and_assert`. Returns the new drawer id.
522 /// Test: `http_create_drawer_runs_auto_kg_extraction`,
523 /// `create_drawer_rejects_json_content_without_force`,
524 /// `create_drawer_force_bypasses_quality_gate_for_json_content`.
525 pub async fn create_drawer(
526 &self,
527 id: &str,
528 body: CreateDrawerBody,
529 creator: CreatorInfo,
530 source: ActivitySource,
531 ) -> ServiceResult<Uuid> {
532 let handle = self.open_handle(id)?;
533 let room = body
534 .room
535 .as_deref()
536 .map(RoomType::parse)
537 .unwrap_or(RoomType::General);
538 let importance = body.importance.unwrap_or(0.5);
539 let force = body.force.unwrap_or(false);
540 let content_preview = drawer_content_preview(&body.content);
541 let mut tags_with_creator = body.tags;
542 // Issue #202: project a bare-UUID session tag (when the caller
543 // passed one in the request body) into the reserved
544 // `creator:session=<first-8>` slot so the activity panel can
545 // surface session attribution without bespoke parsing.
546 if let Some(session_tag) = crate::attribution::session_tag_from_tags(&tags_with_creator) {
547 tags_with_creator.push(session_tag);
548 }
549 creator.merge_into(&mut tags_with_creator);
550 let content_for_kg = body.content.clone();
551 let tags_for_kg = tags_with_creator.clone();
552 let room_label_for_kg = crate::tools::room_label(&room);
553 let drawer_id = handle
554 .remember_with_options(
555 body.content,
556 room,
557 tags_with_creator,
558 importance,
559 RememberOptions {
560 force,
561 ..Default::default()
562 },
563 )
564 .await
565 .map_err(|e| ServiceError::internal(format!("remember: {e:#}")))?;
566 let drawer_count = handle.drawers.read().len();
567 // Issue #228: resolve from the in-memory cache instead of re-walking
568 // the data root on every HTTP `create_drawer` call. Same cache the
569 // MCP `lookup_palace_name` helper consults.
570 let palace_name = self
571 .state
572 .palace_names
573 .get(id)
574 .map(|entry| entry.value().clone())
575 .unwrap_or_else(|| id.to_string());
576 self.state.emit(DaemonEvent::DrawerAdded {
577 palace_id: id.to_string(),
578 palace_name,
579 drawer_count,
580 timestamp: chrono::Utc::now(),
581 content_preview,
582 source,
583 });
584 // Issue #228: do NOT emit `StatusChanged` on every drawer create —
585 // the periodic ticker (`run_http_on`) refreshes aggregate totals on
586 // a fixed cadence so dashboards stay current without an O(N palaces)
587 // recompute on the write hot path.
588 crate::tools::auto_extract_and_assert(
589 &handle,
590 drawer_id,
591 &content_for_kg,
592 &tags_for_kg,
593 room_label_for_kg.as_deref(),
594 )
595 .await;
596 Ok(drawer_id)
597 }
598
599 /// Forget (delete) a drawer and emit the matching events.
600 ///
601 /// Why: same dedup story as `create_drawer`. #5231: `DELETE` on a drawer id
602 /// that was never stored used to answer `204 No Content`, the same as a
603 /// real delete — this now 404s, matching `delete_palace`.
604 /// What: parses the drawer UUID, calls `PalaceHandle::forget`, maps
605 /// [`ForgetOutcome::NotFound`] to `ServiceError::not_found`, and emits
606 /// `DrawerDeleted` only when a drawer was actually removed.
607 /// Test: `delete_drawer_404s_for_an_unknown_drawer_id`.
608 pub async fn delete_drawer(
609 &self,
610 id: &str,
611 drawer_id: &str,
612 source: ActivitySource,
613 ) -> ServiceResult<()> {
614 let handle = self.open_handle(id)?;
615 let uuid = Uuid::parse_str(drawer_id)
616 .map_err(|_| ServiceError::bad_request("drawer_id must be a UUID"))?;
617 let outcome = handle
618 .forget(uuid)
619 .await
620 .map_err(|e| ServiceError::internal(format!("forget: {e:#}")))?;
621 if !outcome.is_deleted() {
622 return Err(ServiceError::not_found(format!(
623 "drawer '{drawer_id}' not found in palace '{id}'"
624 )));
625 }
626 let drawer_count = handle.drawers.read().len();
627 self.state.emit(DaemonEvent::DrawerDeleted {
628 palace_id: id.to_string(),
629 drawer_count,
630 source,
631 });
632 // Issue #228: skip the per-write `StatusChanged` emit — the
633 // periodic ticker handles aggregate roll-ups.
634 Ok(())
635 }
636
637 // -----------------------------------------------------------------
638 // Recall
639 // -----------------------------------------------------------------
640
641 /// Per-palace recall (semantic search), optionally with deep retrieval.
642 ///
643 /// Why: HTTP and chat tools both perform the same fan-out logic.
644 /// What: opens the palace handle and dispatches to the shallow or deep
645 /// recall helper. Returns a JSON array of flattened drawer rows (the
646 /// `recall_entry_json` shape from issue #69).
647 /// Test: `recall_entry_json_hoists_drawer_fields`.
648 pub async fn recall(
649 &self,
650 id: &str,
651 query: &str,
652 top_k: usize,
653 deep: bool,
654 ) -> ServiceResult<Value> {
655 let handle = self.open_handle(id)?;
656 let results = if deep {
657 recall_deep_with_default_embedder(&handle, query, top_k).await
658 } else {
659 recall_with_default_embedder(&handle, query, top_k).await
660 }
661 .map_err(|e| ServiceError::internal(format!("recall: {e:#}")))?;
662 let payload: Vec<Value> = results.into_iter().map(recall_entry_json).collect();
663 Ok(json!(payload))
664 }
665
666 /// Cross-palace recall.
667 ///
668 /// Why: shared between `/api/v1/recall` and the `memory_recall_all` chat
669 /// tool. Encapsulating the open-everything-fanout-merge dance avoids
670 /// drift.
671 /// What: lists every palace, opens handles (skipping failures with a
672 /// `tracing::warn!`), delegates to
673 /// `recall_across_palaces_with_default_embedder`. Returns a JSON array.
674 /// Why (issue #4637): unlike `list_palaces`/`status`, this route is NOT
675 /// converted to `peek()`. A cross-palace recall that answered from
676 /// cache-resident palaces only would silently omit ~98.9% of the corpus —
677 /// a wrong answer that looks like a right one, which is strictly worse
678 /// than a slow correct one. The open loop keeps opening every palace; it
679 /// just no longer does so inline on a tokio worker thread. Making this
680 /// route actually fast needs a different design (a shared cross-palace
681 /// index, or an explicit palace-scoped query), not a cache-only read.
682 /// Test: indirectly via `recall_across_palaces_merges_results` and the
683 /// MCP `memory_recall_all` integration paths;
684 /// `open_palaces_blocking_opens_every_palace` pins that uncached palaces
685 /// are still searched.
686 pub async fn recall_all(&self, query: &str, top_k: usize, deep: bool) -> Value {
687 let palaces = match list_palaces_blocking(&self.state).await {
688 Ok(v) => v,
689 Err(e) => return json!({ "error": format!("{e:#}") }),
690 };
691 // #4637: open_palace (not peek) is deliberate — recall must see every
692 // palace; the spawn_blocking hop keeps it off the async executor.
693 let handles = open_palaces_blocking(&self.state, &palaces, "recall_all").await;
694 if handles.is_empty() {
695 return json!([]);
696 }
697 match recall_across_palaces_with_default_embedder(&handles, query, top_k, deep).await {
698 Ok(results) => json!(results
699 .into_iter()
700 .map(|r| json!({
701 "palace_id": r.palace_id,
702 "drawer_id": r.result.drawer.id.to_string(),
703 "content": r.result.drawer.content,
704 "importance": r.result.drawer.importance,
705 "tags": r.result.drawer.tags,
706 "score": r.result.score,
707 "layer": r.result.layer,
708 }))
709 .collect::<Vec<_>>()),
710 Err(e) => json!({ "error": format!("recall_across_palaces: {e:#}") }),
711 }
712 }
713}