mod disk_stats;
use anyhow::Result;
use serde_json::{json, Value};
use trusty_common::console_metrics::{make_report, ServiceHealth};
use trusty_common::memory_core::store::rooms::list_room_summaries;
use trusty_common::memory_core::PalaceRegistry;
use crate::AppState;
const MAX_PALACES_IN_REPORT: usize = 20;
pub fn descriptor() -> Value {
json!({
"name": "console_metrics",
"description": "Return a ConsoleMetricsReport with palace aggregate statistics \
(palace_count, counted_palace_count, cached_palace_count, total_drawers, \
total_vectors, total_rooms, total_kg_triples) and per-palace detail \
(first 20). Counts come from the open-handle LRU cache for a resident \
palace and from a read-only pass over the palace's redb files otherwise, \
so a palace that is merely closed still reports real numbers; no palace is \
opened to be counted. Each entry carries `cached` (was it resident) and \
`stats_source` (`cache`, `disk`, or `unavailable`); an unavailable entry \
carries `stats_error` and null counts. Each entry also carries \
`last_used_unix` — when a recall, remember or note last touched that \
palace, or null when none has since the stamp shipped. Used by the \
trusty-console dashboard metrics poller.",
"inputSchema": {
"type": "object",
"properties": {},
"required": []
}
})
}
struct PalaceStats {
palace_count: usize,
counted_palace_count: usize,
cached_palace_count: usize,
total_drawers: usize,
total_vectors: usize,
total_rooms: usize,
total_kg_triples: usize,
palace_entries: Vec<Value>,
}
enum PalaceCounts {
Counted {
cached: bool,
drawers: usize,
vectors: usize,
rooms: usize,
kg_triples: usize,
},
Unavailable(String),
}
pub async fn handle_console_metrics(state: &AppState, _args: Value) -> Result<Value> {
let root = state.data_root.clone();
let palace_infos =
match tokio::task::spawn_blocking(move || PalaceRegistry::list_palaces(&root))
.await
.map_err(|e| anyhow::anyhow!("join list_palaces: {e}"))?
{
Ok(v) => v,
Err(e) => {
tracing::warn!("console_metrics: list_palaces failed: {e:#}");
Vec::new()
}
};
let registry = state.registry.clone();
let stats = tokio::task::spawn_blocking(move || collect_palace_stats(®istry, &palace_infos))
.await
.map_err(|e| anyhow::anyhow!("join collect_palace_stats: {e}"))?;
let metrics = json!({
"palace_count": stats.palace_count,
"counted_palace_count": stats.counted_palace_count,
"cached_palace_count": stats.cached_palace_count,
"total_drawers": stats.total_drawers,
"total_vectors": stats.total_vectors,
"total_rooms": stats.total_rooms,
"total_kg_triples": stats.total_kg_triples,
"palaces": stats.palace_entries,
});
let report = make_report(
"trusty-memory",
"Trusty Memory",
env!("CARGO_PKG_VERSION"),
ServiceHealth::Ok,
metrics,
4,
);
Ok(serde_json::to_value(&report)?)
}
fn count_palace(
registry: &PalaceRegistry,
info: &trusty_common::memory_core::Palace,
) -> PalaceCounts {
match registry.peek(&info.id) {
Some(handle) => cached_counts(info, &handle),
None => match disk_stats::read(&info.data_dir) {
Ok(s) => PalaceCounts::Counted {
cached: false,
drawers: s.drawer_count,
vectors: s.vector_count,
rooms: s.room_count,
kg_triples: s.kg_triple_count,
},
Err(reason) => {
tracing::debug!(palace = %info.id, "console_metrics: {reason}");
PalaceCounts::Unavailable(reason)
}
},
}
}
fn cached_counts(
info: &trusty_common::memory_core::Palace,
handle: &trusty_common::memory_core::PalaceHandle,
) -> PalaceCounts {
let drawers = handle.drawers.read().len();
let vectors = match cached_field(
&info.id,
"vector index",
handle.vector_store.try_index_size(),
) {
Ok(n) => n,
Err(reason) => return PalaceCounts::Unavailable(reason),
};
let rooms = match cached_field(
&info.id,
"room list",
list_room_summaries(&handle.kg.store()).map(|r| r.len()),
) {
Ok(n) => n,
Err(reason) => return PalaceCounts::Unavailable(reason),
};
let kg_triples = match cached_field(
&info.id,
"kg_triple count",
handle.kg.count_active_triples(),
) {
Ok(n) => n,
Err(reason) => return PalaceCounts::Unavailable(reason),
};
PalaceCounts::Counted {
cached: true,
drawers,
vectors,
rooms,
kg_triples,
}
}
fn cached_field(
palace: &trusty_common::memory_core::PalaceId,
field: &str,
read: anyhow::Result<usize>,
) -> Result<usize, String> {
read.map_err(|e| {
let reason = format!("{field} unavailable: {e:#}");
tracing::warn!(palace = %palace, "{reason}");
reason
})
}
fn collect_palace_stats(
registry: &PalaceRegistry,
palace_infos: &[trusty_common::memory_core::Palace],
) -> PalaceStats {
let palace_count = palace_infos.len();
let mut stats = PalaceStats {
palace_count,
counted_palace_count: 0,
cached_palace_count: 0,
total_drawers: 0,
total_vectors: 0,
total_rooms: 0,
total_kg_triples: 0,
palace_entries: Vec::with_capacity(palace_count.min(MAX_PALACES_IN_REPORT)),
};
for (rank, info) in palace_infos.iter().enumerate() {
let counts = count_palace(registry, info);
if let PalaceCounts::Counted {
cached,
drawers,
vectors,
rooms,
kg_triples,
} = &counts
{
stats.counted_palace_count += 1;
stats.cached_palace_count += usize::from(*cached);
stats.total_drawers += drawers;
stats.total_vectors += vectors;
stats.total_rooms += rooms;
stats.total_kg_triples += kg_triples;
}
if rank < MAX_PALACES_IN_REPORT {
stats.palace_entries.push(palace_entry(info, &counts));
}
}
stats
}
fn palace_entry(info: &trusty_common::memory_core::Palace, counts: &PalaceCounts) -> Value {
let id = info.id.as_str().to_string();
let last_used_unix = crate::palace_last_used::read(&info.data_dir);
match counts {
PalaceCounts::Counted {
cached,
drawers,
vectors,
rooms,
kg_triples,
} => json!({
"id": id,
"name": info.name,
"drawer_count": drawers,
"vector_count": vectors,
"room_count": rooms,
"kg_triple_count": kg_triples,
"cached": cached,
"stats_source": if *cached { "cache" } else { "disk" },
"last_used_unix": last_used_unix,
}),
PalaceCounts::Unavailable(reason) => json!({
"id": id,
"name": info.name,
"drawer_count": Value::Null,
"vector_count": Value::Null,
"room_count": Value::Null,
"kg_triple_count": Value::Null,
"cached": false,
"stats_source": "unavailable",
"stats_error": reason,
"last_used_unix": last_used_unix,
}),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[serial_test::serial]
#[tokio::test]
async fn handle_console_metrics_returns_valid_report() {
unsafe {
std::env::set_var("TRUSTY_SKIP_PALACE_ENFORCEMENT", "1");
}
let tmp = tempfile::tempdir().expect("tempdir");
let state = crate::AppState::new(tmp.path().to_path_buf());
let result = handle_console_metrics(&state, serde_json::json!({}))
.await
.expect("console_metrics must not return Err");
assert_eq!(result["service_id"], "trusty-memory");
assert_eq!(result["display_name"], "Trusty Memory");
assert!(result["version"].is_string());
assert!(result["status"].is_string());
assert_eq!(result["metrics_schema_version"], 4);
assert!(result["collected_at_unix"].is_number());
assert_eq!(result["metrics"]["palace_count"], 0);
assert_eq!(result["metrics"]["counted_palace_count"], 0);
assert_eq!(result["metrics"]["cached_palace_count"], 0);
assert_eq!(result["metrics"]["total_drawers"], 0);
assert_eq!(result["metrics"]["total_vectors"], 0);
assert_eq!(result["metrics"]["total_rooms"], 0);
assert_eq!(result["metrics"]["total_kg_triples"], 0);
assert!(result["metrics"]["palaces"].is_array());
assert_eq!(result["metrics"]["palaces"].as_array().unwrap().len(), 0);
}
#[tokio::test]
async fn console_metrics_uses_cache_only_and_does_not_evict() {
use trusty_common::memory_core::{Palace, PalaceId};
let tmp = tempfile::tempdir().expect("tempdir");
let data_root = tmp.path().to_path_buf();
let registry = 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: chrono::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 must hold only 2 handles after 3 creates"
);
assert!(
registry.peek(&PalaceId::new("a")).is_none(),
"'a' must already be evicted before console_metrics runs"
);
let mut state = crate::AppState::new(data_root);
state.registry = std::sync::Arc::new(registry);
let result = handle_console_metrics(&state, serde_json::json!({}))
.await
.expect("console_metrics must not return Err");
assert_eq!(
result["metrics"]["palace_count"], 3,
"palace_count reflects all 3 on-disk palaces"
);
assert_eq!(
result["metrics"]["cached_palace_count"], 2,
"cached_palace_count reflects only the 2 still-resident handles"
);
assert_eq!(
state.registry.len(),
2,
"console_metrics must not grow the LRU cache"
);
assert!(
state.registry.peek(&PalaceId::new("a")).is_none(),
"console_metrics 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());
let entries = result["metrics"]["palaces"]
.as_array()
.expect("palaces array present");
assert_eq!(entries.len(), 3);
let entry = |id: &str| {
entries
.iter()
.find(|e| e["id"] == id)
.unwrap_or_else(|| panic!("entry for '{id}' present"))
};
assert_eq!(entry("a")["cached"], false, "'a' is not cached");
assert_eq!(entry("b")["cached"], true, "'b' is cached");
assert_eq!(entry("c")["cached"], true, "'c' is cached");
assert_eq!(
entry("a")["stats_source"],
"disk",
"an evicted palace's counts come off its files: {:?}",
entry("a")
);
assert_eq!(entry("b")["stats_source"], "cache");
}
#[tokio::test]
async fn console_metrics_reports_real_counts_for_an_uncached_palace() {
use trusty_common::memory_core::{Drawer, Palace, PalaceId};
let tmp = tempfile::tempdir().expect("tempdir");
let data_root = tmp.path().to_path_buf();
{
let registry = PalaceRegistry::with_max_open(4);
let handle = registry
.create_palace(
&data_root,
Palace {
id: PalaceId::new("cold"),
name: "cold".to_string(),
description: None,
created_at: chrono::Utc::now(),
data_dir: data_root.join("cold"),
},
)
.expect("create_palace");
for i in 0..2 {
let drawer = Drawer::new(uuid::Uuid::new_v4(), format!("drawer {i}"));
handle.kg.store().upsert_drawer(&drawer).expect("upsert");
}
drop(handle);
registry.remove(&PalaceId::new("cold"));
}
let state = crate::AppState::new(data_root);
assert_eq!(state.registry.len(), 0, "nothing may be cache-resident");
let result = handle_console_metrics(&state, serde_json::json!({}))
.await
.expect("console_metrics must not return Err");
let entry = &result["metrics"]["palaces"][0];
assert_eq!(entry["id"], "cold");
assert_eq!(
entry["drawer_count"], 2,
"an uncached palace must report its real drawers, not 0: {entry:?}"
);
assert_eq!(entry["stats_source"], "disk");
assert_eq!(entry["cached"], false);
assert_eq!(
result["metrics"]["total_drawers"], 2,
"the totals must include palaces read off disk"
);
assert_eq!(result["metrics"]["counted_palace_count"], 1);
assert_eq!(
result["metrics"]["cached_palace_count"], 0,
"residency is still reported honestly"
);
}
#[tokio::test]
async fn console_metrics_reports_room_counts_for_every_palace() {
use trusty_common::memory_core::store::rooms::create_room;
use trusty_common::memory_core::{Palace, PalaceId, RoomType};
let tmp = tempfile::tempdir().expect("tempdir");
let data_root = tmp.path().to_path_buf();
let registry = PalaceRegistry::with_max_open(4);
let mut resident = None;
for name in ["hot", "cold"] {
let handle = registry
.create_palace(
&data_root,
Palace {
id: PalaceId::new(name),
name: name.to_string(),
description: None,
created_at: chrono::Utc::now(),
data_dir: data_root.join(name),
},
)
.unwrap_or_else(|e| panic!("create_palace({name}): {e:#}"));
if name == "hot" {
create_room(
&handle.kg.store(),
&RoomType::Custom("decisions".to_string()),
None,
)
.expect("create_room");
resident = Some(handle);
} else {
drop(handle);
registry.remove(&PalaceId::new(name));
}
}
let _resident = resident.expect("hot palace handle kept");
let mut state = crate::AppState::new(data_root);
state.registry = std::sync::Arc::new(registry);
let result = handle_console_metrics(&state, serde_json::json!({}))
.await
.expect("console_metrics must not return Err");
let entries = result["metrics"]["palaces"]
.as_array()
.expect("palaces array");
for e in entries {
assert!(
e["room_count"].is_number(),
"every palace row must carry a room_count: {e:?}"
);
}
let hot = entries
.iter()
.find(|e| e["id"] == "hot")
.expect("hot entry present");
assert_eq!(hot["stats_source"], "cache");
assert_eq!(
hot["room_count"], 1,
"the room that was created must be counted: {hot:?}"
);
assert_eq!(
result["metrics"]["total_rooms"], 1,
"total_rooms sums the per-palace counts"
);
}
#[tokio::test]
async fn console_metrics_marks_an_unreadable_palace_unavailable() {
use trusty_common::memory_core::{Palace, PalaceId};
let tmp = tempfile::tempdir().expect("tempdir");
let data_root = tmp.path().to_path_buf();
{
let registry = PalaceRegistry::with_max_open(4);
let handle = registry
.create_palace(
&data_root,
Palace {
id: PalaceId::new("gone"),
name: "gone".to_string(),
description: None,
created_at: chrono::Utc::now(),
data_dir: data_root.join("gone"),
},
)
.expect("create_palace");
drop(handle);
registry.remove(&PalaceId::new("gone"));
}
std::fs::remove_file(data_root.join("gone").join("kg.redb")).expect("remove kg store");
let state = crate::AppState::new(data_root);
let result = handle_console_metrics(&state, serde_json::json!({}))
.await
.expect("console_metrics must not return Err");
let entry = &result["metrics"]["palaces"][0];
assert_eq!(entry["stats_source"], "unavailable");
assert!(
entry["drawer_count"].is_null(),
"an unreadable count must be null, never 0: {entry:?}"
);
assert!(
entry["stats_error"].is_string(),
"an unavailable row must say why: {entry:?}"
);
assert_eq!(
result["metrics"]["counted_palace_count"], 0,
"a palace that could not be read did not contribute to the totals"
);
}
#[test]
fn cached_field_turns_a_read_error_into_unavailable_not_zero() {
use trusty_common::memory_core::PalaceId;
let palace = PalaceId::new("hot");
for field in ["vector index", "room list", "kg_triple count"] {
let result = cached_field(&palace, field, Err(anyhow::anyhow!("redb read failed")));
let reason = match result {
Err(reason) => reason,
Ok(n) => panic!("field {field:?} must surface its read error, not {n}"),
};
assert!(
reason.contains(field),
"the reason must name the field that failed: {reason}"
);
assert!(
reason.contains("redb read failed"),
"the reason must carry the underlying error: {reason}"
);
}
}
#[test]
fn cached_field_passes_a_successful_read_through() {
use trusty_common::memory_core::PalaceId;
let palace = PalaceId::new("hot");
assert_eq!(cached_field(&palace, "vector index", Ok(7)), Ok(7));
}
}