use crate::AppState;
use anyhow::{anyhow, Context, Result};
use serde_json::{json, Value};
use std::collections::HashSet;
use std::sync::Arc;
use trusty_common::memory_core::palace::{Palace, PalaceId};
use trusty_common::memory_core::retrieval::RecallResult;
use trusty_common::memory_core::PalaceHandle;
use uuid::Uuid;
use super::types::{PalaceInfo, ServiceResult};
#[cfg(test)]
use super::core::MemoryService;
#[cfg(test)]
use super::types::ListDrawersQuery;
pub const DRAWER_PREVIEW_MAX_CHARS: usize = 80;
pub const DRAWER_SNIPPET_MAX_CHARS: usize = 60;
pub fn drawer_content_preview(content: &str) -> String {
let normalised: String = content.split_whitespace().collect::<Vec<_>>().join(" ");
if normalised.chars().count() <= DRAWER_PREVIEW_MAX_CHARS {
normalised
} else {
let kept: String = normalised
.chars()
.take(DRAWER_PREVIEW_MAX_CHARS.saturating_sub(1))
.collect();
format!("{kept}…")
}
}
pub fn drawer_snippet(content: &str) -> String {
let normalised: String = content.split_whitespace().collect::<Vec<_>>().join(" ");
if normalised.chars().count() <= DRAWER_SNIPPET_MAX_CHARS {
normalised
} else {
let kept: String = normalised
.chars()
.take(DRAWER_SNIPPET_MAX_CHARS.saturating_sub(1))
.collect();
format!("{kept}…")
}
}
pub fn recall_entry_json(r: RecallResult) -> Value {
let mut obj = match serde_json::to_value(&r.drawer) {
Ok(Value::Object(map)) => map,
_ => serde_json::Map::new(),
};
obj.insert("score".to_string(), json!(r.score));
obj.insert("layer".to_string(), json!(r.layer));
Value::Object(obj)
}
pub(crate) fn is_reserved_system_palace(id: &PalaceId) -> bool {
id.as_str().starts_with("__")
}
pub(crate) struct PalaceStats {
pub total_drawers: usize,
pub total_vectors: usize,
pub total_kg_triples: usize,
pub cached_palace_count: usize,
}
pub(crate) fn collect_palace_stats<'a, I>(state: &AppState, ids: I) -> PalaceStats
where
I: IntoIterator<Item = &'a PalaceId>,
{
let (mut total_drawers, mut total_vectors, mut total_kg_triples): (usize, usize, usize) =
(0, 0, 0);
let mut cached_palace_count: usize = 0;
for id in ids {
if let Some(handle) = state.registry.peek(id) {
total_drawers = total_drawers.saturating_add(handle.drawers.read().len());
total_vectors = total_vectors.saturating_add(handle.vector_store.index_size());
total_kg_triples = total_kg_triples.saturating_add(kg_triple_count_or_zero(&handle));
cached_palace_count += 1;
}
}
PalaceStats {
total_drawers,
total_vectors,
total_kg_triples,
cached_palace_count,
}
}
pub(crate) async fn list_palaces_blocking(state: &AppState) -> Result<Vec<Palace>> {
let root = state.data_root.clone();
tokio::task::spawn_blocking(move || {
trusty_common::memory_core::PalaceRegistry::list_palaces(&root)
})
.await
.map_err(|e| anyhow!("join list_palaces: {e}"))?
.map_err(|e| anyhow!("list palaces: {e:#}"))
}
pub(crate) async fn open_palaces_blocking(
state: &AppState,
palaces: &[Palace],
label: &'static str,
) -> Vec<Arc<PalaceHandle>> {
let registry = Arc::clone(&state.registry);
let root = state.data_root.clone();
let ids: Vec<PalaceId> = palaces.iter().map(|p| p.id.clone()).collect();
tokio::task::spawn_blocking(move || {
let mut handles = Vec::with_capacity(ids.len());
for id in &ids {
match registry.open_palace(&root, id) {
Ok(h) => handles.push(h),
Err(e) => tracing::warn!(palace = %id, "{label}: open failed: {e:#}"),
}
}
handles
})
.await
.unwrap_or_else(|e| {
tracing::warn!("{label}: join open_palaces failed: {e}");
Vec::new()
})
}
fn room_registry_count(handle: &Arc<PalaceHandle>) -> Option<usize> {
match handle.kg.store().list_rooms() {
Ok(rooms) => Some(rooms.len()),
Err(e) => {
tracing::warn!(palace = %handle.id, "room_count unavailable: {e:#}");
None
}
}
}
fn wing_registry_count(handle: &Arc<PalaceHandle>) -> Option<usize> {
match handle.kg.store().list_wings() {
Ok(wings) => Some(wings.len()),
Err(e) => {
tracing::warn!(palace = %handle.id, "wing_count unavailable: {e:#}");
None
}
}
}
pub(crate) fn kg_triple_count_or_zero(handle: &Arc<PalaceHandle>) -> usize {
match handle.kg.count_active_triples() {
Ok(n) => n,
Err(e) => {
tracing::warn!(palace = %handle.id, "kg_triple_count unavailable: {e:#}");
0
}
}
}
pub fn palace_info_from(palace: &Palace, handle: Option<&Arc<PalaceHandle>>) -> PalaceInfo {
let (
drawer_count,
vector_count,
kg_triple_count,
room_count,
wing_count,
last_write_at,
node_count,
edge_count,
community_count,
is_compacting,
) = if let Some(h) = handle {
let drawers = h.drawers.read();
let last_write = drawers.iter().map(|d| d.created_at).max();
let rooms = room_registry_count(h).unwrap_or_else(|| {
drawers
.iter()
.map(|d| d.room_id)
.collect::<HashSet<Uuid>>()
.len()
});
let wings = wing_registry_count(h).unwrap_or(1);
(
drawers.len(),
h.vector_store.index_size(),
kg_triple_count_or_zero(h),
rooms,
wings,
last_write,
h.kg.node_count() as u64,
h.kg.edge_count() as u64,
h.kg.community_count() as u64,
h.is_compacting(),
)
} else {
(0, 0, 0, 0, 0, None, 0, 0, 0, false)
};
PalaceInfo {
id: palace.id.0.clone(),
name: palace.name.clone(),
description: palace.description.clone(),
drawer_count,
vector_count,
kg_triple_count,
room_count,
wing_count,
created_at: palace.created_at,
last_write_at,
node_count,
edge_count,
community_count,
is_compacting,
cached: handle.is_some(),
}
}
pub async fn refresh_gaps_cache(state: &AppState, handle: &Arc<PalaceHandle>) {
let mut gaps = handle.kg.knowledge_gaps();
if let Ok(api_key) = std::env::var(trusty_common::env_vars::ENV_OPENROUTER_API_KEY) {
if !api_key.is_empty() {
for gap in gaps.iter_mut() {
if let Some(enriched) = enrich_gap_exploration(&api_key, gap).await {
gap.suggested_exploration = enriched;
}
}
}
}
let gap_count = gaps.len();
state.registry.set_gaps(handle.id.clone(), gaps);
tracing::debug!(palace = %handle.id, gaps = gap_count, "community gaps updated");
}
pub async fn enrich_gap_exploration(
api_key: &str,
gap: &trusty_common::memory_core::community::KnowledgeGap,
) -> Option<String> {
let preview: Vec<&str> = gap.entities.iter().take(5).map(String::as_str).collect();
if preview.is_empty() {
return None;
}
let entities = preview.join(", ");
let user = format!(
"Given these related entities from a knowledge graph: {entities}. \
Suggest one specific research question (single sentence, under 25 words) \
that would help fill gaps in this knowledge cluster. Return only the question."
);
let messages = vec![trusty_common::ChatMessage {
role: "user".to_string(),
content: user,
tool_call_id: None,
tool_calls: None,
}];
#[allow(deprecated)]
let res = trusty_common::openrouter_chat(api_key, "openai/gpt-4o-mini", messages).await;
match res {
Ok(text) => {
let trimmed = text.trim().to_string();
if trimmed.is_empty() {
None
} else {
Some(trimmed)
}
}
Err(e) => {
tracing::debug!("openrouter gap enrichment failed (using template): {e:#}");
None
}
}
}
pub fn service_result_to_anyhow<T: serde::Serialize>(r: ServiceResult<T>) -> Result<Value> {
match r {
Ok(v) => serde_json::to_value(v).context("serialize service result"),
Err(e) => Err(anyhow!("{e}")),
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::{Duration as ChronoDuration, Utc};
use trusty_common::memory_core::palace::{Drawer, Palace};
fn test_state() -> AppState {
let tmp = tempfile::tempdir().expect("tempdir");
let root = tmp.path().to_path_buf();
std::mem::forget(tmp);
AppState::new(root)
}
#[tokio::test]
async fn list_drawers_creates_desc_paginates() {
let state = test_state();
let palace = Palace {
id: PalaceId::new("paging-test"),
name: "paging-test".to_string(),
description: None,
created_at: Utc::now(),
data_dir: state.data_root.join("paging-test"),
};
state
.registry
.create_palace(&state.data_root, palace)
.expect("create_palace");
let handle = state
.registry
.open_palace(&state.data_root, &PalaceId::new("paging-test"))
.expect("open_palace");
let room_id = Uuid::nil();
let now = Utc::now();
for (i, importance) in [0.1f32, 0.9, 0.3, 0.7, 0.5].iter().enumerate() {
let mut drawer = Drawer::new(room_id, format!("drawer-{i}"));
drawer.importance = *importance;
drawer.created_at = now - ChronoDuration::seconds(i as i64);
drawer.tags = vec![format!("idx:{i}")];
handle.add_drawer(drawer);
}
drop(handle);
let service = MemoryService::new(state.clone());
let page1 = service
.list_drawers(
"paging-test",
ListDrawersQuery {
limit: Some(2),
offset: Some(0),
sort: Some("created_desc".into()),
..Default::default()
},
)
.await
.expect("page 1");
let arr = page1.as_array().expect("array");
assert_eq!(arr.len(), 2, "page 1 must have 2 rows");
assert_eq!(arr[0]["content"].as_str(), Some("drawer-0"));
assert_eq!(arr[1]["content"].as_str(), Some("drawer-1"));
let page2 = service
.list_drawers(
"paging-test",
ListDrawersQuery {
limit: Some(2),
offset: Some(2),
sort: Some("created_desc".into()),
..Default::default()
},
)
.await
.expect("page 2");
let arr = page2.as_array().expect("array");
assert_eq!(arr.len(), 2, "page 2 must have 2 rows");
assert_eq!(arr[0]["content"].as_str(), Some("drawer-2"));
assert_eq!(arr[1]["content"].as_str(), Some("drawer-3"));
let page3 = service
.list_drawers(
"paging-test",
ListDrawersQuery {
limit: Some(2),
offset: Some(4),
sort: Some("created_desc".into()),
..Default::default()
},
)
.await
.expect("page 3");
let arr = page3.as_array().expect("array");
assert_eq!(arr.len(), 1, "page 3 (tail) must have 1 row");
assert_eq!(arr[0]["content"].as_str(), Some("drawer-4"));
let legacy = service
.list_drawers(
"paging-test",
ListDrawersQuery {
limit: Some(1),
..Default::default()
},
)
.await
.expect("legacy");
let arr = legacy.as_array().expect("array");
assert_eq!(arr.len(), 1);
assert_eq!(
arr[0]["content"].as_str(),
Some("drawer-1"),
"importance default should surface drawer with importance 0.9 first",
);
assert_eq!(
arr[0]["snippet"].as_str(),
Some("drawer-1"),
"snippet must be populated for non-empty drawer content",
);
}
#[test]
fn drawer_snippet_truncates_long_content() {
assert_eq!(drawer_snippet("hello world"), "hello world");
assert_eq!(
drawer_snippet("first line\n\nsecond\tline third"),
"first line second line third",
);
assert_eq!(drawer_snippet(" padded "), "padded");
let long = "a".repeat(200);
let snippet = drawer_snippet(&long);
assert_eq!(snippet.chars().count(), DRAWER_SNIPPET_MAX_CHARS);
assert!(
snippet.ends_with('…'),
"long body must be truncated with ellipsis",
);
let exact = "a".repeat(DRAWER_SNIPPET_MAX_CHARS);
assert_eq!(drawer_snippet(&exact), exact);
}
#[test]
fn drawer_snippet_handles_empty_content() {
assert_eq!(drawer_snippet(""), "");
assert_eq!(drawer_snippet(" \n\t "), "");
}
fn state_with_one_evicted_palace() -> AppState {
let tmp = tempfile::tempdir().expect("tempdir");
let data_root = tmp.path().to_path_buf();
std::mem::forget(tmp);
let registry = trusty_common::memory_core::PalaceRegistry::with_max_open(2);
for name in ["a", "b", "c"] {
let palace = Palace {
id: PalaceId::new(name),
name: name.to_string(),
description: None,
created_at: Utc::now(),
data_dir: data_root.join(name),
};
registry
.create_palace(&data_root, palace)
.unwrap_or_else(|e| panic!("create_palace({name}) failed: {e:#}"));
}
assert_eq!(registry.len(), 2, "capacity-2 registry holds 2 of 3");
assert!(
registry.peek(&PalaceId::new("a")).is_none(),
"'a' must be evicted before the route under test runs"
);
let mut state = AppState::new(data_root);
state.registry = Arc::new(registry);
state
}
#[tokio::test]
async fn list_palaces_does_not_open_uncached_palaces() {
let state = state_with_one_evicted_palace();
let svc = MemoryService::new(state.clone());
let rows = svc.list_palaces().await.expect("list_palaces");
assert_eq!(rows.len(), 3, "every on-disk palace is still listed");
let row = |id: &str| {
rows.iter()
.find(|r| r.id == id)
.unwrap_or_else(|| panic!("row for '{id}' present"))
};
assert!(!row("a").cached, "'a' was evicted, so it is not cached");
assert_eq!(
row("a").drawer_count,
0,
"an uncached row reports 0 (unknown), not a live count"
);
assert!(row("b").cached, "'b' is resident");
assert!(row("c").cached, "'c' is resident");
assert_eq!(
state.registry.len(),
2,
"list_palaces must not grow the LRU cache"
);
assert!(
state.registry.peek(&PalaceId::new("a")).is_none(),
"list_palaces must not reopen the evicted palace 'a'"
);
assert!(state.registry.peek(&PalaceId::new("b")).is_some());
assert!(state.registry.peek(&PalaceId::new("c")).is_some());
}
#[tokio::test]
async fn status_does_not_open_uncached_palaces() {
let state = state_with_one_evicted_palace();
let svc = MemoryService::new(state.clone());
let payload = svc.status().await;
assert_eq!(
payload.palace_count, 3,
"palace_count still reflects every palace on disk"
);
assert_eq!(
payload.cached_palace_count, 2,
"the totals cover only the 2 cache-resident palaces"
);
assert_eq!(
state.registry.len(),
2,
"status must not grow the LRU cache"
);
assert!(
state.registry.peek(&PalaceId::new("a")).is_none(),
"status must not reopen the evicted palace 'a'"
);
}
#[tokio::test]
async fn open_palaces_blocking_opens_every_palace() {
let state = state_with_one_evicted_palace();
let palaces = list_palaces_blocking(&state).await.expect("list palaces");
assert_eq!(palaces.len(), 3);
let handles = open_palaces_blocking(&state, &palaces, "test").await;
assert_eq!(
handles.len(),
3,
"recall fan-out must open every palace, including uncached ones"
);
assert!(
handles.iter().any(|h| h.id == PalaceId::new("a")),
"the evicted palace 'a' must still be opened and searched"
);
}
}