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