Skip to main content

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}