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::{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}