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