trusty_memory/service/core_kg.rs
1//! `MemoryService` knowledge-graph / dream / activity methods.
2//!
3//! Why: the `MemoryService` impl exceeded the 500-SLOC production cap once the
4//! former monolithic `service.rs` was split (issue #607); its KG, dream-cycle,
5//! and activity-listing methods form a cohesive second half hosted here.
6//! What: a continuation `impl MemoryService` block whose methods were moved
7//! verbatim from the original single impl.
8//! Test: covered by the corresponding `web::tests` / `service::tests`.
9
10use crate::kg_write::{CachePolicy, KgWriteError};
11use crate::{ActivityFilter, ActivitySource, DaemonEvent};
12use std::sync::Arc;
13use trusty_common::memory_core::dream::{DreamConfig, Dreamer, PersistedDreamStats};
14use trusty_common::memory_core::palace::PalaceId;
15// #4670: ExpandDirection drives the progressive `kg/graph/neighbors` route.
16use trusty_common::memory_core::store::kg::{ExpandDirection, Triple};
17use trusty_common::memory_core::{PalaceHandle, PalaceRegistry};
18
19use super::core::KG_GRAPH_MAX_TRIPLES;
20use super::helpers::{list_palaces_blocking, refresh_gaps_cache};
21use super::types::{
22 DreamStatusPayload, KgAssertBody, KgGraphPayload, KgNeighborsPayload, KgNodeView,
23 KgSeedPayload, ServiceError, ServiceResult,
24};
25use super::MemoryService;
26
27// ---------------------------------------------------------------------------
28// KG list bounds
29// ---------------------------------------------------------------------------
30//
31// #4776: these live on the service layer, not on either consumer, because both
32// readers of the subject list must agree on them — the HTTP explorer routes in
33// `web::kg_routes` (compiled only under `axum-server`) and the MCP
34// `kg_list_subjects` tool in `tools::kg_ops` (always compiled). Defining them
35// in `web` would put them on the far side of that feature gate from the tool.
36
37/// Default page size for KG subject listings when the caller omits `limit`.
38///
39/// Why: 50 is large enough to feel responsive in the KG Explorer and to answer
40/// "what is in this graph?" in one call, without dumping a full graph.
41pub(crate) const DEFAULT_KG_LIST_LIMIT: usize = 50;
42
43/// Hard ceiling on `limit` for KG subject listings.
44///
45/// Why: prevent a misconfigured client from asking the daemon to materialize
46/// thousands of rows in one go; matches the spec's max=200.
47pub(crate) const MAX_KG_LIST_LIMIT: usize = 200;
48
49impl MemoryService {
50 // -----------------------------------------------------------------
51 // Knowledge graph
52 // -----------------------------------------------------------------
53
54 /// Query the KG for all active triples whose subject matches.
55 pub async fn kg_query(&self, id: &str, subject: &str) -> ServiceResult<Vec<Triple>> {
56 let handle = self.open_handle(id)?;
57 handle
58 .kg
59 .query_active(subject)
60 .await
61 .map_err(|e| ServiceError::internal(format!("kg query: {e:#}")))
62 }
63
64 /// Assert a triple in the KG.
65 ///
66 /// Assert a triple through `POST /api/v1/palaces/{id}/kg`.
67 ///
68 /// Why: #4888 — this accepts an arbitrary predicate, so it can write a hot
69 /// one and must carry the same Tier S gate as the MCP tool. Not a
70 /// hypothetical path: `trusty-mpm`'s provisioner seeds its identity fact
71 /// through exactly this endpoint. #5524 — it also owed the prompt-cache
72 /// rebuild and never ran it, so a hot fact written here was stored and then
73 /// invisible to every later turn until some unrelated write rebuilt the
74 /// cache.
75 /// What: delegates the whole admission → assert → refresh sequence to
76 /// [`crate::kg_write::assert_triple`]. The variant split is what lets this
77 /// keep answering 400 for a refused write and 500 for a failed one.
78 /// Test: `http_kg_assert_endpoint_refreshes_prompt_cache`,
79 /// `http_kg_assert_endpoint_rejects_over_long_tier_s_object` in
80 /// `web::tests::prompt_tests`.
81 pub async fn kg_assert(&self, id: &str, body: KgAssertBody) -> ServiceResult<()> {
82 let handle = self.open_handle(id)?;
83 let triple = Triple {
84 subject: body.subject,
85 predicate: body.predicate,
86 object: body.object,
87 valid_from: chrono::Utc::now(),
88 valid_to: None,
89 confidence: body.confidence.unwrap_or(1.0),
90 provenance: body.provenance,
91 };
92 // #5524: route through the shared entry point so the prompt-cache
93 // rebuild cannot be forgotten here again.
94 crate::kg_write::assert_triple(&self.state, &handle, triple, CachePolicy::Inline)
95 .await
96 .map(|_| ())
97 .map_err(|e| match e {
98 KgWriteError::Admission(inner) => ServiceError::bad_request(format!("{inner:#}")),
99 other => ServiceError::internal(format!("{other}")),
100 })
101 }
102
103 /// Close the one active triple `(subject, predicate, object)`, leaving
104 /// every sibling object at that pair active. Returns the rows closed.
105 ///
106 /// Why: Issue #278 — the `DELETE /kg/triples/<id>` HTTP endpoint needs a
107 /// service-layer method so the HTTP handler stays a thin adapter. It took
108 /// no object and called the pair-level `KnowledgeGraph::retract`, whose
109 /// meaning is "close every active row at this pair", so a caller deleting
110 /// one triple lost the siblings it never named. Retraction is a soft close
111 /// — `close_active_row` copies the row to a `hist:` key first — so the
112 /// damage was recoverable, not silent data loss.
113 /// What: Opens the palace handle and calls
114 /// [`trusty_common::memory_core::store::kg::KnowledgeGraph::retract_triple`],
115 /// which keys on all three fields. Returns the closed count so the caller
116 /// can tell a retraction (`1`) from a miss (`0`); a miss is a genuine
117 /// no-op, which makes the call idempotent. Rebuilds the prompt cache when
118 /// a hot-predicate row was actually closed — otherwise a retracted Tier S
119 /// fact keeps being injected until the next write. This mirrors the
120 /// `kg_retract_triple` MCP tool so both surfaces agree.
121 /// Test: `kg_delete_triple_closes_one_object_and_keeps_siblings`,
122 /// `kg_delete_triple_returns_404_for_missing`,
123 /// `kg_delete_triple_rebuilds_prompt_cache_for_hot_predicate` in
124 /// `web::tests`. Only the third reaches the cache rebuild — the other two
125 /// retract under the predicate `is`, which is not hot.
126 pub async fn kg_retract_triple(
127 &self,
128 id: &str,
129 subject: &str,
130 predicate: &str,
131 object: &str,
132 ) -> ServiceResult<usize> {
133 let handle = self.open_handle(id)?;
134 let closed = handle
135 .kg
136 .retract_triple(subject, predicate, object)
137 .await
138 .map_err(|e| ServiceError::internal(format!("kg retract_triple: {e:#}")))?;
139 if closed > 0 && crate::prompt_facts::is_hot_predicate(predicate) {
140 // The write landed either way and the cache is only a
141 // denormalisation, so a rebuild failure is logged, not fatal.
142 // No test drives this arm: `rebuild_prompt_cache` skips a palace
143 // it cannot read and has no other fallible step, so it cannot
144 // currently return `Err`.
145 if let Err(e) = crate::prompt_facts::rebuild_prompt_cache(&self.state).await {
146 tracing::warn!("rebuild_prompt_cache after kg_retract_triple failed: {e:#}");
147 }
148 }
149 Ok(closed)
150 }
151
152 /// List distinct subjects in the KG.
153 pub async fn kg_list_subjects(&self, id: &str, limit: usize) -> ServiceResult<Vec<String>> {
154 let handle = self.open_handle(id)?;
155 handle
156 .kg
157 .list_subjects(limit)
158 .map_err(|e| ServiceError::internal(format!("kg list_subjects: {e:#}")))
159 }
160
161 /// List distinct subjects in the KG paired with their active-triple count.
162 pub async fn kg_list_subjects_with_counts(
163 &self,
164 id: &str,
165 limit: usize,
166 ) -> ServiceResult<Vec<(String, u64)>> {
167 let handle = self.open_handle(id)?;
168 handle
169 .kg
170 .list_subjects_with_counts(limit)
171 .map_err(|e| ServiceError::internal(format!("kg list_subjects_with_counts: {e:#}")))
172 }
173
174 /// Page through every active triple.
175 pub async fn kg_list_all(
176 &self,
177 id: &str,
178 limit: usize,
179 offset: usize,
180 ) -> ServiceResult<Vec<Triple>> {
181 let handle = self.open_handle(id)?;
182 handle
183 .kg
184 .list_active(limit, offset)
185 .await
186 .map_err(|e| ServiceError::internal(format!("kg list_active: {e:#}")))
187 }
188
189 /// Return the count of currently-active triples.
190 ///
191 /// #5384: a failed count read is a 500, not `{"active": 0}` — the badge
192 /// this feeds cannot tell those apart.
193 pub async fn kg_count(&self, id: &str) -> ServiceResult<usize> {
194 let handle = self.open_handle(id)?;
195 handle.kg.count_active_triples().map_err(|e| {
196 ServiceError::internal(format!("kg count_active_triples for palace {id}: {e:#}"))
197 })
198 }
199
200 /// Build the per-palace visual graph payload.
201 ///
202 /// Why (issue #4670): the `node_count` / `edge_count` / `community_count`
203 /// here are computed over the FULL adjacency while `triples` is capped at
204 /// [`KG_GRAPH_MAX_TRIPLES`]. That mismatch used to be invisible — the UI
205 /// rendered 5,000 triples under a "9,311 nodes" badge — and because
206 /// `list_active` orders by `valid_from` DESC the dropped triples were
207 /// silently the oldest. The payload now reports what it actually returned
208 /// alongside what exists, so truncation is machine-detectable.
209 /// What: unchanged query; adds `returned_triple_count`,
210 /// `active_triple_count`, and the derived `truncated` flag.
211 /// Test: `kg_graph_signals_truncation`, `kg_graph_returns_active_triples`.
212 pub async fn kg_graph(&self, id: &str) -> ServiceResult<KgGraphPayload> {
213 self.kg_graph_with_cap(id, KG_GRAPH_MAX_TRIPLES).await
214 }
215
216 /// [`Self::kg_graph`] with an explicit triple cap.
217 ///
218 /// Why (issue #4670): the truncation-signalling branch is only reachable
219 /// above `KG_GRAPH_MAX_TRIPLES` (5,000), and seeding 5,001 triples costs
220 /// ~90 s of test time — expensive enough that the branch would in practice
221 /// go untested. Taking the cap as a parameter makes it provable with five
222 /// triples and a cap of three, at no cost to the production call path.
223 /// What: the real implementation; `kg_graph` is a thin wrapper that passes
224 /// the production constant.
225 /// Test: `kg_graph_signals_truncation`.
226 pub async fn kg_graph_with_cap(
227 &self,
228 id: &str,
229 max_triples: usize,
230 ) -> ServiceResult<KgGraphPayload> {
231 let handle = self.open_handle(id)?;
232 let triples = handle
233 .kg
234 .list_active(max_triples, 0)
235 .await
236 .map_err(|e| ServiceError::internal(format!("kg list_active: {e:#}")))?;
237 // #4670: compare against the true active count, not the cap, so a
238 // palace sitting exactly on the cap is not falsely flagged.
239 // #5384: a failed read would come back as 0 and make `truncated` false
240 // for every payload, which is the flag's exact failure mode.
241 let active_triple_count = handle.kg.count_active_triples().map_err(|e| {
242 ServiceError::internal(format!("kg count_active_triples for palace {id}: {e:#}"))
243 })? as u64;
244 let returned_triple_count = triples.len() as u64;
245 Ok(KgGraphPayload {
246 triples,
247 node_count: handle.kg.node_count() as u64,
248 edge_count: handle.kg.edge_count() as u64,
249 community_count: handle.kg.community_count() as u64,
250 returned_triple_count,
251 active_triple_count,
252 truncated: returned_triple_count < active_triple_count,
253 })
254 }
255
256 /// Top-`limit` nodes by degree plus the edges among them (issue #4670).
257 ///
258 /// Why: first paint must show the graph's skeleton, not 9,311 nodes in an
259 /// O(n²) layout. Measured on the live 8,266-triple palace, 90.2% of nodes
260 /// are degree-1 leaves and only 7.2% have degree >= 5, so a
261 /// top-degree slice carries essentially all of the visible structure and
262 /// everything else stays one click away.
263 /// What: runs `KnowledgeGraph::top_degree_subgraph` over the resident
264 /// adjacency (O(V log V + E), no disk I/O) and pairs the result with the
265 /// palace-wide totals the header needs to report honestly.
266 /// Test: `kg_graph_seed_ranks_by_degree`, `kg_graph_seed_clamps_limit`.
267 pub async fn kg_graph_seed(&self, id: &str, limit: usize) -> ServiceResult<KgSeedPayload> {
268 let handle = self.open_handle(id)?;
269 let (nodes, triples) = handle
270 .kg
271 .top_degree_subgraph(limit)
272 .map_err(|e| ServiceError::internal(format!("kg top_degree_subgraph: {e:#}")))?;
273 let node_count = handle.kg.node_count() as u64;
274 let returned_node_count = nodes.len() as u64;
275 Ok(KgSeedPayload {
276 nodes: nodes.into_iter().map(KgNodeView::from).collect(),
277 returned_triple_count: triples.len() as u64,
278 triples,
279 node_count,
280 edge_count: handle.kg.edge_count() as u64,
281 community_count: handle.kg.community_count() as u64,
282 returned_node_count,
283 limit: limit as u64,
284 truncated: returned_node_count < node_count,
285 })
286 }
287
288 /// Direction-aware, hop-bounded expansion around one node (issue #4670).
289 ///
290 /// Why: click-to-expand needs "what points AT this node", which no HTTP
291 /// endpoint could answer — `kg_query` is a subject prefix scan. Bounding
292 /// the hops keeps one click on a hub from pulling the whole palace.
293 /// What: delegates to `KnowledgeGraph::expand_neighbors`. `direction`
294 /// and `max_hops` are already validated/clamped by the HTTP layer; they
295 /// are echoed back so the client can see what actually ran.
296 /// Test: `kg_neighbors_returns_incoming_edges`, `kg_neighbors_clamps_max_hops`.
297 pub async fn kg_neighbors(
298 &self,
299 id: &str,
300 node: &str,
301 direction: ExpandDirection,
302 max_hops: usize,
303 ) -> ServiceResult<KgNeighborsPayload> {
304 let handle = self.open_handle(id)?;
305 let (nodes, triples) = handle
306 .kg
307 .expand_neighbors(node, direction, max_hops)
308 .map_err(|e| ServiceError::internal(format!("kg expand_neighbors: {e:#}")))?;
309 Ok(KgNeighborsPayload {
310 origin: node.to_string(),
311 returned_node_count: nodes.len() as u64,
312 returned_triple_count: triples.len() as u64,
313 nodes: nodes.into_iter().map(KgNodeView::from).collect(),
314 triples,
315 direction: match direction {
316 ExpandDirection::In => "in",
317 ExpandDirection::Out => "out",
318 ExpandDirection::Both => "both",
319 }
320 .to_string(),
321 max_hops: max_hops as u64,
322 })
323 }
324
325 // -----------------------------------------------------------------
326 // Dream cycle
327 // -----------------------------------------------------------------
328
329 /// Aggregate dream stats across every persisted palace.
330 pub async fn dream_status_aggregate(&self) -> DreamStatusPayload {
331 let palaces = PalaceRegistry::list_palaces(&self.state.data_root).unwrap_or_default();
332 let mut out = DreamStatusPayload::default();
333 let mut latest: Option<chrono::DateTime<chrono::Utc>> = None;
334 for p in palaces {
335 let data_dir = self.state.data_root.join(p.id.as_str());
336 let snap = match PersistedDreamStats::load(&data_dir) {
337 Ok(Some(s)) => s,
338 _ => continue,
339 };
340 out.merged = out.merged.saturating_add(snap.stats.merged);
341 out.pruned = out.pruned.saturating_add(snap.stats.pruned);
342 out.compacted = out.compacted.saturating_add(snap.stats.compacted);
343 out.closets_updated = out
344 .closets_updated
345 .saturating_add(snap.stats.closets_updated);
346 out.duration_ms = out.duration_ms.saturating_add(snap.stats.duration_ms);
347 latest = match latest {
348 Some(t) if t >= snap.last_run_at => Some(t),
349 _ => Some(snap.last_run_at),
350 };
351 }
352 out.last_run_at = latest;
353 out
354 }
355
356 /// Per-palace dream stats snapshot.
357 pub async fn dream_status_for_palace(&self, id: &str) -> ServiceResult<DreamStatusPayload> {
358 let data_dir = self.state.data_root.join(id);
359 if !data_dir.exists() {
360 return Err(ServiceError::not_found(format!("palace not found: {id}")));
361 }
362 match PersistedDreamStats::load(&data_dir) {
363 Ok(Some(s)) => Ok(s.into()),
364 Ok(None) => Ok(DreamStatusPayload::default()),
365 Err(e) => Err(ServiceError::internal(format!("read dream stats: {e:#}"))),
366 }
367 }
368
369 /// Run a dream cycle across every palace.
370 ///
371 /// Why (issue #4637): like `recall_all`, this route is deliberately NOT
372 /// converted to `PalaceRegistry::peek`. Dreaming is a maintenance pass —
373 /// consolidating only the 64 palaces that happen to be cache-resident
374 /// would silently stop maintaining the other ~5,730, which is a worse
375 /// failure than a slow run because nothing reports it. Opening every
376 /// palace stays correct; what changes is that the blocking open no longer
377 /// runs inline on a tokio worker thread. A long dream run is expected —
378 /// this is an explicitly-triggered `POST`, not a page load.
379 /// What: lists palaces on the blocking pool, then per palace hops to
380 /// `spawn_blocking` for the open before awaiting the async dream cycle.
381 /// Test: `dream_run_aggregates_stats`.
382 pub async fn dream_run(&self) -> ServiceResult<DreamStatusPayload> {
383 let palaces = list_palaces_blocking(&self.state)
384 .await
385 .map_err(|e| ServiceError::internal(format!("{e:#}")))?;
386 let dreamer = Dreamer::new(DreamConfig::default());
387 let mut out = DreamStatusPayload::default();
388 for p in palaces {
389 // #4637: open_palace (not peek) is deliberate — a dream cycle must
390 // maintain every palace; the spawn_blocking hop keeps the cold open
391 // off the async executor.
392 let registry = std::sync::Arc::clone(&self.state.registry);
393 let root = self.state.data_root.clone();
394 let pid = p.id.clone();
395 let opened =
396 tokio::task::spawn_blocking(move || registry.open_palace(&root, &pid)).await;
397 let handle = match opened {
398 Ok(Ok(h)) => h,
399 Ok(Err(e)) => {
400 tracing::warn!(palace = %p.id, "dream_run: open failed: {e:#}");
401 continue;
402 }
403 Err(e) => {
404 tracing::warn!(palace = %p.id, "dream_run: join open failed: {e}");
405 continue;
406 }
407 };
408 match dreamer.dream_cycle(&handle).await {
409 Ok(stats) => {
410 out.merged = out.merged.saturating_add(stats.merged);
411 out.pruned = out.pruned.saturating_add(stats.pruned);
412 out.compacted = out.compacted.saturating_add(stats.compacted);
413 out.closets_updated = out.closets_updated.saturating_add(stats.closets_updated);
414 out.duration_ms = out.duration_ms.saturating_add(stats.duration_ms);
415 }
416 Err(e) => tracing::warn!(palace = %p.id, "dream_run: cycle failed: {e:#}"),
417 }
418 refresh_gaps_cache(&self.state, &handle).await;
419 }
420 out.last_run_at = Some(chrono::Utc::now());
421 self.state.emit(DaemonEvent::DreamCompleted {
422 palace_id: None,
423 merged: out.merged,
424 pruned: out.pruned,
425 compacted: out.compacted,
426 closets_updated: out.closets_updated,
427 duration_ms: out.duration_ms,
428 source: ActivitySource::Http,
429 });
430 self.state.emit(self.aggregate_status_event());
431 Ok(out)
432 }
433
434 // -----------------------------------------------------------------
435 // Activity log
436 // -----------------------------------------------------------------
437
438 /// Paginated activity-log read.
439 pub async fn list_activity(
440 &self,
441 filter: ActivityFilter,
442 limit: usize,
443 offset: usize,
444 ) -> ServiceResult<(Vec<crate::ActivityEntry>, u64)> {
445 let entries = self
446 .state
447 .activity_log
448 .list(&filter, limit, offset)
449 .map_err(|e| ServiceError::internal(format!("activity list: {e:#}")))?;
450 let total = self
451 .state
452 .activity_log
453 .count()
454 .map_err(|e| ServiceError::internal(format!("activity count: {e:#}")))?;
455 Ok((entries, total))
456 }
457
458 // -----------------------------------------------------------------
459 // Internal helper — open a palace handle, 404 only on a genuine absence.
460 // -----------------------------------------------------------------
461
462 /// Open the named palace.
463 ///
464 /// Why (#5549, ADR-0045): this mapped every `open_palace` failure to
465 /// `NotFound`, which the HTTP layer renders as 404. A denied or transient
466 /// read of `palace.json`, undecodable metadata, an open-queue timeout, or a
467 /// redb write-lock conflict then all reported that the palace does not
468 /// exist — erasing at the caller the distinction `load_palace` draws, and
469 /// across a much wider surface than the two rename paths: every
470 /// `/api/v1/palaces/{id}/kg*` endpoint, the drawer CRUD routes, and
471 /// per-palace recall reach this one helper.
472 /// What: returns `ServiceError::NotFound` only when
473 /// `PalaceRegistry::open_error_is_absent` confirms the palace is genuinely
474 /// not there, and `ServiceError::Internal` (500) otherwise.
475 /// Test: `unreadable_palace_is_500_not_404_at_the_service_open_handle`,
476 /// `unstattable_palace_is_500_not_404_at_the_service_open_handle`,
477 /// `absent_palace_is_still_404_at_both_open_handles`.
478 pub fn open_handle(&self, id: &str) -> ServiceResult<Arc<PalaceHandle>> {
479 self.state
480 .registry
481 .open_palace(&self.state.data_root, &PalaceId::new(id))
482 .map_err(|e| {
483 // #5549: every open failure mapped to 404, so a palace that
484 // could not be read was reported as one that is not there.
485 if PalaceRegistry::open_error_is_absent(&e) {
486 ServiceError::not_found(format!("palace not found: {id} ({e:#})"))
487 } else {
488 ServiceError::internal(format!("palace could not be loaded: {id} ({e:#})"))
489 }
490 })
491 }
492}