trusty-memory 0.23.1

MCP server (stdio + HTTP/SSE) for trusty-memory
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
//! `MemoryService` knowledge-graph / dream / activity methods.
//!
//! Why: the `MemoryService` impl exceeded the 500-SLOC production cap once the
//! former monolithic `service.rs` was split (issue #607); its KG, dream-cycle,
//! and activity-listing methods form a cohesive second half hosted here.
//! What: a continuation `impl MemoryService` block whose methods were moved
//! verbatim from the original single impl.
//! Test: covered by the corresponding `web::tests` / `service::tests`.

use crate::kg_write::{CachePolicy, KgWriteError};
use crate::{ActivityFilter, ActivitySource, DaemonEvent};
use std::sync::Arc;
use trusty_common::memory_core::dream::{DreamConfig, Dreamer, PersistedDreamStats};
use trusty_common::memory_core::palace::PalaceId;
// #4670: ExpandDirection drives the progressive `kg/graph/neighbors` route.
use trusty_common::memory_core::store::kg::{ExpandDirection, Triple};
use trusty_common::memory_core::{PalaceHandle, PalaceRegistry};

use super::core::KG_GRAPH_MAX_TRIPLES;
use super::helpers::{list_palaces_blocking, refresh_gaps_cache};
use super::types::{
    DreamStatusPayload, KgAssertBody, KgGraphPayload, KgNeighborsPayload, KgNodeView,
    KgSeedPayload, ServiceError, ServiceResult,
};
use super::MemoryService;

// ---------------------------------------------------------------------------
// KG list bounds
// ---------------------------------------------------------------------------
//
// #4776: these live on the service layer, not on either consumer, because both
// readers of the subject list must agree on them — the HTTP explorer routes in
// `web::kg_routes` (compiled only under `axum-server`) and the MCP
// `kg_list_subjects` tool in `tools::kg_ops` (always compiled). Defining them
// in `web` would put them on the far side of that feature gate from the tool.

/// Default page size for KG subject listings when the caller omits `limit`.
///
/// Why: 50 is large enough to feel responsive in the KG Explorer and to answer
/// "what is in this graph?" in one call, without dumping a full graph.
pub(crate) const DEFAULT_KG_LIST_LIMIT: usize = 50;

/// Hard ceiling on `limit` for KG subject listings.
///
/// Why: prevent a misconfigured client from asking the daemon to materialize
/// thousands of rows in one go; matches the spec's max=200.
pub(crate) const MAX_KG_LIST_LIMIT: usize = 200;

impl MemoryService {
    // -----------------------------------------------------------------
    // Knowledge graph
    // -----------------------------------------------------------------

    /// Query the KG for all active triples whose subject matches.
    pub async fn kg_query(&self, id: &str, subject: &str) -> ServiceResult<Vec<Triple>> {
        let handle = self.open_handle(id)?;
        handle
            .kg
            .query_active(subject)
            .await
            .map_err(|e| ServiceError::internal(format!("kg query: {e:#}")))
    }

    /// Assert a triple in the KG.
    ///
    /// Assert a triple through `POST /api/v1/palaces/{id}/kg`.
    ///
    /// Why: #4888 — this accepts an arbitrary predicate, so it can write a hot
    /// one and must carry the same Tier S gate as the MCP tool. Not a
    /// hypothetical path: `trusty-mpm`'s provisioner seeds its identity fact
    /// through exactly this endpoint. #5524 — it also owed the prompt-cache
    /// rebuild and never ran it, so a hot fact written here was stored and then
    /// invisible to every later turn until some unrelated write rebuilt the
    /// cache.
    /// What: delegates the whole admission → assert → refresh sequence to
    /// [`crate::kg_write::assert_triple`]. The variant split is what lets this
    /// keep answering 400 for a refused write and 500 for a failed one.
    /// Test: `http_kg_assert_endpoint_refreshes_prompt_cache`,
    /// `http_kg_assert_endpoint_rejects_over_long_tier_s_object` in
    /// `web::tests::prompt_tests`.
    pub async fn kg_assert(&self, id: &str, body: KgAssertBody) -> ServiceResult<()> {
        let handle = self.open_handle(id)?;
        let triple = Triple {
            subject: body.subject,
            predicate: body.predicate,
            object: body.object,
            valid_from: chrono::Utc::now(),
            valid_to: None,
            confidence: body.confidence.unwrap_or(1.0),
            provenance: body.provenance,
        };
        // #5524: route through the shared entry point so the prompt-cache
        // rebuild cannot be forgotten here again.
        crate::kg_write::assert_triple(&self.state, &handle, triple, CachePolicy::Inline)
            .await
            .map(|_| ())
            .map_err(|e| match e {
                KgWriteError::Admission(inner) => ServiceError::bad_request(format!("{inner:#}")),
                other => ServiceError::internal(format!("{other}")),
            })
    }

    /// Close the one active triple `(subject, predicate, object)`, leaving
    /// every sibling object at that pair active. Returns the rows closed.
    ///
    /// Why: Issue #278 — the `DELETE /kg/triples/<id>` HTTP endpoint needs a
    /// service-layer method so the HTTP handler stays a thin adapter. It took
    /// no object and called the pair-level `KnowledgeGraph::retract`, whose
    /// meaning is "close every active row at this pair", so a caller deleting
    /// one triple lost the siblings it never named. Retraction is a soft close
    /// — `close_active_row` copies the row to a `hist:` key first — so the
    /// damage was recoverable, not silent data loss.
    /// What: Opens the palace handle and calls
    /// [`trusty_common::memory_core::store::kg::KnowledgeGraph::retract_triple`],
    /// which keys on all three fields. Returns the closed count so the caller
    /// can tell a retraction (`1`) from a miss (`0`); a miss is a genuine
    /// no-op, which makes the call idempotent. Rebuilds the prompt cache when
    /// a hot-predicate row was actually closed — otherwise a retracted Tier S
    /// fact keeps being injected until the next write. This mirrors the
    /// `kg_retract_triple` MCP tool so both surfaces agree.
    /// Test: `kg_delete_triple_closes_one_object_and_keeps_siblings`,
    /// `kg_delete_triple_returns_404_for_missing`,
    /// `kg_delete_triple_rebuilds_prompt_cache_for_hot_predicate` in
    /// `web::tests`. Only the third reaches the cache rebuild — the other two
    /// retract under the predicate `is`, which is not hot.
    pub async fn kg_retract_triple(
        &self,
        id: &str,
        subject: &str,
        predicate: &str,
        object: &str,
    ) -> ServiceResult<usize> {
        let handle = self.open_handle(id)?;
        let closed = handle
            .kg
            .retract_triple(subject, predicate, object)
            .await
            .map_err(|e| ServiceError::internal(format!("kg retract_triple: {e:#}")))?;
        if closed > 0 && crate::prompt_facts::is_hot_predicate(predicate) {
            // The write landed either way and the cache is only a
            // denormalisation, so a rebuild failure is logged, not fatal.
            // No test drives this arm: `rebuild_prompt_cache` skips a palace
            // it cannot read and has no other fallible step, so it cannot
            // currently return `Err`.
            if let Err(e) = crate::prompt_facts::rebuild_prompt_cache(&self.state).await {
                tracing::warn!("rebuild_prompt_cache after kg_retract_triple failed: {e:#}");
            }
        }
        Ok(closed)
    }

    /// List distinct subjects in the KG.
    pub async fn kg_list_subjects(&self, id: &str, limit: usize) -> ServiceResult<Vec<String>> {
        let handle = self.open_handle(id)?;
        handle
            .kg
            .list_subjects(limit)
            .map_err(|e| ServiceError::internal(format!("kg list_subjects: {e:#}")))
    }

    /// List distinct subjects in the KG paired with their active-triple count.
    pub async fn kg_list_subjects_with_counts(
        &self,
        id: &str,
        limit: usize,
    ) -> ServiceResult<Vec<(String, u64)>> {
        let handle = self.open_handle(id)?;
        handle
            .kg
            .list_subjects_with_counts(limit)
            .map_err(|e| ServiceError::internal(format!("kg list_subjects_with_counts: {e:#}")))
    }

    /// Page through every active triple.
    pub async fn kg_list_all(
        &self,
        id: &str,
        limit: usize,
        offset: usize,
    ) -> ServiceResult<Vec<Triple>> {
        let handle = self.open_handle(id)?;
        handle
            .kg
            .list_active(limit, offset)
            .await
            .map_err(|e| ServiceError::internal(format!("kg list_active: {e:#}")))
    }

    /// Return the count of currently-active triples.
    ///
    /// #5384: a failed count read is a 500, not `{"active": 0}` — the badge
    /// this feeds cannot tell those apart.
    pub async fn kg_count(&self, id: &str) -> ServiceResult<usize> {
        let handle = self.open_handle(id)?;
        handle.kg.count_active_triples().map_err(|e| {
            ServiceError::internal(format!("kg count_active_triples for palace {id}: {e:#}"))
        })
    }

    /// Build the per-palace visual graph payload.
    ///
    /// Why (issue #4670): the `node_count` / `edge_count` / `community_count`
    /// here are computed over the FULL adjacency while `triples` is capped at
    /// [`KG_GRAPH_MAX_TRIPLES`]. That mismatch used to be invisible — the UI
    /// rendered 5,000 triples under a "9,311 nodes" badge — and because
    /// `list_active` orders by `valid_from` DESC the dropped triples were
    /// silently the oldest. The payload now reports what it actually returned
    /// alongside what exists, so truncation is machine-detectable.
    /// What: unchanged query; adds `returned_triple_count`,
    /// `active_triple_count`, and the derived `truncated` flag.
    /// Test: `kg_graph_signals_truncation`, `kg_graph_returns_active_triples`.
    pub async fn kg_graph(&self, id: &str) -> ServiceResult<KgGraphPayload> {
        self.kg_graph_with_cap(id, KG_GRAPH_MAX_TRIPLES).await
    }

    /// [`Self::kg_graph`] with an explicit triple cap.
    ///
    /// Why (issue #4670): the truncation-signalling branch is only reachable
    /// above `KG_GRAPH_MAX_TRIPLES` (5,000), and seeding 5,001 triples costs
    /// ~90 s of test time — expensive enough that the branch would in practice
    /// go untested. Taking the cap as a parameter makes it provable with five
    /// triples and a cap of three, at no cost to the production call path.
    /// What: the real implementation; `kg_graph` is a thin wrapper that passes
    /// the production constant.
    /// Test: `kg_graph_signals_truncation`.
    pub async fn kg_graph_with_cap(
        &self,
        id: &str,
        max_triples: usize,
    ) -> ServiceResult<KgGraphPayload> {
        let handle = self.open_handle(id)?;
        let triples = handle
            .kg
            .list_active(max_triples, 0)
            .await
            .map_err(|e| ServiceError::internal(format!("kg list_active: {e:#}")))?;
        // #4670: compare against the true active count, not the cap, so a
        // palace sitting exactly on the cap is not falsely flagged.
        // #5384: a failed read would come back as 0 and make `truncated` false
        // for every payload, which is the flag's exact failure mode.
        let active_triple_count = handle.kg.count_active_triples().map_err(|e| {
            ServiceError::internal(format!("kg count_active_triples for palace {id}: {e:#}"))
        })? as u64;
        let returned_triple_count = triples.len() as u64;
        Ok(KgGraphPayload {
            triples,
            node_count: handle.kg.node_count() as u64,
            edge_count: handle.kg.edge_count() as u64,
            community_count: handle.kg.community_count() as u64,
            returned_triple_count,
            active_triple_count,
            truncated: returned_triple_count < active_triple_count,
        })
    }

    /// Top-`limit` nodes by degree plus the edges among them (issue #4670).
    ///
    /// Why: first paint must show the graph's skeleton, not 9,311 nodes in an
    /// O(n²) layout. Measured on the live 8,266-triple palace, 90.2% of nodes
    /// are degree-1 leaves and only 7.2% have degree >= 5, so a
    /// top-degree slice carries essentially all of the visible structure and
    /// everything else stays one click away.
    /// What: runs `KnowledgeGraph::top_degree_subgraph` over the resident
    /// adjacency (O(V log V + E), no disk I/O) and pairs the result with the
    /// palace-wide totals the header needs to report honestly.
    /// Test: `kg_graph_seed_ranks_by_degree`, `kg_graph_seed_clamps_limit`.
    pub async fn kg_graph_seed(&self, id: &str, limit: usize) -> ServiceResult<KgSeedPayload> {
        let handle = self.open_handle(id)?;
        let (nodes, triples) = handle
            .kg
            .top_degree_subgraph(limit)
            .map_err(|e| ServiceError::internal(format!("kg top_degree_subgraph: {e:#}")))?;
        let node_count = handle.kg.node_count() as u64;
        let returned_node_count = nodes.len() as u64;
        Ok(KgSeedPayload {
            nodes: nodes.into_iter().map(KgNodeView::from).collect(),
            returned_triple_count: triples.len() as u64,
            triples,
            node_count,
            edge_count: handle.kg.edge_count() as u64,
            community_count: handle.kg.community_count() as u64,
            returned_node_count,
            limit: limit as u64,
            truncated: returned_node_count < node_count,
        })
    }

    /// Direction-aware, hop-bounded expansion around one node (issue #4670).
    ///
    /// Why: click-to-expand needs "what points AT this node", which no HTTP
    /// endpoint could answer — `kg_query` is a subject prefix scan. Bounding
    /// the hops keeps one click on a hub from pulling the whole palace.
    /// What: delegates to `KnowledgeGraph::expand_neighbors`. `direction`
    /// and `max_hops` are already validated/clamped by the HTTP layer; they
    /// are echoed back so the client can see what actually ran.
    /// Test: `kg_neighbors_returns_incoming_edges`, `kg_neighbors_clamps_max_hops`.
    pub async fn kg_neighbors(
        &self,
        id: &str,
        node: &str,
        direction: ExpandDirection,
        max_hops: usize,
    ) -> ServiceResult<KgNeighborsPayload> {
        let handle = self.open_handle(id)?;
        let (nodes, triples) = handle
            .kg
            .expand_neighbors(node, direction, max_hops)
            .map_err(|e| ServiceError::internal(format!("kg expand_neighbors: {e:#}")))?;
        Ok(KgNeighborsPayload {
            origin: node.to_string(),
            returned_node_count: nodes.len() as u64,
            returned_triple_count: triples.len() as u64,
            nodes: nodes.into_iter().map(KgNodeView::from).collect(),
            triples,
            direction: match direction {
                ExpandDirection::In => "in",
                ExpandDirection::Out => "out",
                ExpandDirection::Both => "both",
            }
            .to_string(),
            max_hops: max_hops as u64,
        })
    }

    // -----------------------------------------------------------------
    // Dream cycle
    // -----------------------------------------------------------------

    /// Aggregate dream stats across every persisted palace.
    pub async fn dream_status_aggregate(&self) -> DreamStatusPayload {
        let palaces = PalaceRegistry::list_palaces(&self.state.data_root).unwrap_or_default();
        let mut out = DreamStatusPayload::default();
        let mut latest: Option<chrono::DateTime<chrono::Utc>> = None;
        for p in palaces {
            let data_dir = self.state.data_root.join(p.id.as_str());
            let snap = match PersistedDreamStats::load(&data_dir) {
                Ok(Some(s)) => s,
                _ => continue,
            };
            out.merged = out.merged.saturating_add(snap.stats.merged);
            out.pruned = out.pruned.saturating_add(snap.stats.pruned);
            out.compacted = out.compacted.saturating_add(snap.stats.compacted);
            out.closets_updated = out
                .closets_updated
                .saturating_add(snap.stats.closets_updated);
            out.duration_ms = out.duration_ms.saturating_add(snap.stats.duration_ms);
            latest = match latest {
                Some(t) if t >= snap.last_run_at => Some(t),
                _ => Some(snap.last_run_at),
            };
        }
        out.last_run_at = latest;
        out
    }

    /// Per-palace dream stats snapshot.
    pub async fn dream_status_for_palace(&self, id: &str) -> ServiceResult<DreamStatusPayload> {
        let data_dir = self.state.data_root.join(id);
        if !data_dir.exists() {
            return Err(ServiceError::not_found(format!("palace not found: {id}")));
        }
        match PersistedDreamStats::load(&data_dir) {
            Ok(Some(s)) => Ok(s.into()),
            Ok(None) => Ok(DreamStatusPayload::default()),
            Err(e) => Err(ServiceError::internal(format!("read dream stats: {e:#}"))),
        }
    }

    /// Run a dream cycle across every palace.
    ///
    /// Why (issue #4637): like `recall_all`, this route is deliberately NOT
    /// converted to `PalaceRegistry::peek`. Dreaming is a maintenance pass —
    /// consolidating only the 64 palaces that happen to be cache-resident
    /// would silently stop maintaining the other ~5,730, which is a worse
    /// failure than a slow run because nothing reports it. Opening every
    /// palace stays correct; what changes is that the blocking open no longer
    /// runs inline on a tokio worker thread. A long dream run is expected —
    /// this is an explicitly-triggered `POST`, not a page load.
    /// What: lists palaces on the blocking pool, then per palace hops to
    /// `spawn_blocking` for the open before awaiting the async dream cycle.
    /// Test: `dream_run_aggregates_stats`.
    pub async fn dream_run(&self) -> ServiceResult<DreamStatusPayload> {
        let palaces = list_palaces_blocking(&self.state)
            .await
            .map_err(|e| ServiceError::internal(format!("{e:#}")))?;
        let dreamer = Dreamer::new(DreamConfig::default());
        let mut out = DreamStatusPayload::default();
        for p in palaces {
            // #4637: open_palace (not peek) is deliberate — a dream cycle must
            // maintain every palace; the spawn_blocking hop keeps the cold open
            // off the async executor.
            let registry = std::sync::Arc::clone(&self.state.registry);
            let root = self.state.data_root.clone();
            let pid = p.id.clone();
            let opened =
                tokio::task::spawn_blocking(move || registry.open_palace(&root, &pid)).await;
            let handle = match opened {
                Ok(Ok(h)) => h,
                Ok(Err(e)) => {
                    tracing::warn!(palace = %p.id, "dream_run: open failed: {e:#}");
                    continue;
                }
                Err(e) => {
                    tracing::warn!(palace = %p.id, "dream_run: join open failed: {e}");
                    continue;
                }
            };
            match dreamer.dream_cycle(&handle).await {
                Ok(stats) => {
                    out.merged = out.merged.saturating_add(stats.merged);
                    out.pruned = out.pruned.saturating_add(stats.pruned);
                    out.compacted = out.compacted.saturating_add(stats.compacted);
                    out.closets_updated = out.closets_updated.saturating_add(stats.closets_updated);
                    out.duration_ms = out.duration_ms.saturating_add(stats.duration_ms);
                }
                Err(e) => tracing::warn!(palace = %p.id, "dream_run: cycle failed: {e:#}"),
            }
            refresh_gaps_cache(&self.state, &handle).await;
        }
        out.last_run_at = Some(chrono::Utc::now());
        self.state.emit(DaemonEvent::DreamCompleted {
            palace_id: None,
            merged: out.merged,
            pruned: out.pruned,
            compacted: out.compacted,
            closets_updated: out.closets_updated,
            duration_ms: out.duration_ms,
            source: ActivitySource::Http,
        });
        self.state.emit(self.aggregate_status_event());
        Ok(out)
    }

    // -----------------------------------------------------------------
    // Activity log
    // -----------------------------------------------------------------

    /// Paginated activity-log read.
    pub async fn list_activity(
        &self,
        filter: ActivityFilter,
        limit: usize,
        offset: usize,
    ) -> ServiceResult<(Vec<crate::ActivityEntry>, u64)> {
        let entries = self
            .state
            .activity_log
            .list(&filter, limit, offset)
            .map_err(|e| ServiceError::internal(format!("activity list: {e:#}")))?;
        let total = self
            .state
            .activity_log
            .count()
            .map_err(|e| ServiceError::internal(format!("activity count: {e:#}")))?;
        Ok((entries, total))
    }

    // -----------------------------------------------------------------
    // Internal helper — open a palace handle, 404 only on a genuine absence.
    // -----------------------------------------------------------------

    /// Open the named palace.
    ///
    /// Why (#5549, ADR-0045): this mapped every `open_palace` failure to
    /// `NotFound`, which the HTTP layer renders as 404. A denied or transient
    /// read of `palace.json`, undecodable metadata, an open-queue timeout, or a
    /// redb write-lock conflict then all reported that the palace does not
    /// exist — erasing at the caller the distinction `load_palace` draws, and
    /// across a much wider surface than the two rename paths: every
    /// `/api/v1/palaces/{id}/kg*` endpoint, the drawer CRUD routes, and
    /// per-palace recall reach this one helper.
    /// What: returns `ServiceError::NotFound` only when
    /// `PalaceRegistry::open_error_is_absent` confirms the palace is genuinely
    /// not there, and `ServiceError::Internal` (500) otherwise.
    /// Test: `unreadable_palace_is_500_not_404_at_the_service_open_handle`,
    /// `unstattable_palace_is_500_not_404_at_the_service_open_handle`,
    /// `absent_palace_is_still_404_at_both_open_handles`.
    pub fn open_handle(&self, id: &str) -> ServiceResult<Arc<PalaceHandle>> {
        self.state
            .registry
            .open_palace(&self.state.data_root, &PalaceId::new(id))
            .map_err(|e| {
                // #5549: every open failure mapped to 404, so a palace that
                // could not be read was reported as one that is not there.
                if PalaceRegistry::open_error_is_absent(&e) {
                    ServiceError::not_found(format!("palace not found: {id} ({e:#})"))
                } else {
                    ServiceError::internal(format!("palace could not be loaded: {id} ({e:#})"))
                }
            })
    }
}