use std::path::{Path, PathBuf};
use serde_json::{Value, json};
use crate::memory_rpc::call_memory_tool_at_with_timeout;
use crate::monitor::dashboard::{MemoryData, PalaceRow};
use super::parsers::{
parse_drawers, parse_dream_stats, parse_memory_details, parse_memory_event,
parse_palace_detail, parse_recall_hits,
};
use super::types::{DrawerInfo, DreamStats, MemoryDetail, MemoryEvent, REQUEST_TIMEOUT, RecallHit};
const EVENT_PAGE: usize = 100;
#[derive(Debug, Clone)]
pub struct MemoryClient {
pub(super) socket: PathBuf,
}
impl MemoryClient {
pub fn new(socket: impl Into<PathBuf>) -> Self {
Self {
socket: socket.into(),
}
}
pub fn socket(&self) -> &Path {
&self.socket
}
pub fn set_socket(&mut self, socket: impl Into<PathBuf>) {
self.socket = socket.into();
}
async fn call(&self, method: &str, params: serde_json::Value) -> anyhow::Result<Value> {
call_memory_tool_at_with_timeout(&self.socket, method, params, REQUEST_TIMEOUT).await
}
pub async fn fetch_all(&self) -> anyhow::Result<MemoryData> {
use super::types::StatusWire;
let raw = self.call("memory.status", json!({})).await?;
let status: StatusWire = serde_json::from_value(raw)?;
let palaces = match self.palaces().await {
Ok(rows) => rows,
Err(e) => {
tracing::warn!("palace list probe failed: {e}");
Vec::new()
}
};
Ok(MemoryData {
version: status.version,
palace_count: status.palace_count,
total_drawers: status.total_drawers,
total_vectors: status.total_vectors,
total_kg_triples: status.total_kg_triples,
palaces,
})
}
pub async fn is_healthy(&self) -> bool {
self.call("memory.health", json!({})).await.is_ok()
}
async fn palaces(&self) -> anyhow::Result<Vec<PalaceRow>> {
let listed = self.call("memory.palaces_list", json!({})).await?;
Ok(project_palaces(&listed))
}
pub async fn fetch_palace(&self, palace_id: &str) -> anyhow::Result<PalaceRow> {
let raw = self
.call("memory.palace_get", json!({ "palace_id": palace_id }))
.await?;
parse_palace_detail(&raw)
.ok_or_else(|| anyhow::anyhow!("unexpected palace payload for '{palace_id}'"))
}
pub async fn recall(&self, query: &str, top_k: usize) -> anyhow::Result<Vec<RecallHit>> {
let raw = self
.call("memory_recall_all", json!({ "q": query, "top_k": top_k }))
.await?;
Ok(parse_recall_hits(
raw.get("results").unwrap_or(&Value::Null),
))
}
pub async fn list_drawers(
&self,
palace_id: &str,
limit: usize,
offset: usize,
) -> anyhow::Result<Vec<DrawerInfo>> {
let raw = self
.call(
"memory.drawers_list",
json!({
"palace_id": palace_id,
"limit": limit,
"offset": offset,
"sort": "created_desc",
}),
)
.await?;
Ok(parse_drawers(&raw))
}
pub async fn fetch_drawer_detail(
&self,
palace_id: &str,
limit: usize,
) -> anyhow::Result<Vec<MemoryDetail>> {
let raw = self
.call(
"memory.drawers_list",
json!({
"palace_id": palace_id,
"limit": limit,
"sort": "created_desc",
}),
)
.await?;
Ok(parse_memory_details(&raw))
}
pub async fn dream_run(&self) -> anyhow::Result<DreamStats> {
let raw = self.call("memory.dream_run", json!({})).await?;
Ok(parse_dream_stats(&raw))
}
pub async fn recent_events(&self, after_id: u64) -> anyhow::Result<(u64, Vec<MemoryEvent>)> {
let raw = self
.call("memory.activity", json!({ "limit": EVENT_PAGE }))
.await?;
Ok(project_events(&raw, after_id))
}
}
pub(super) fn project_palaces(raw: &Value) -> Vec<PalaceRow> {
let Some(rows) = raw.get("palaces").and_then(|v| v.as_array()) else {
return Vec::new();
};
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let id = row
.get("id")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string();
if let Some(detail) = row.get("palace").and_then(parse_palace_detail) {
out.push(detail);
continue;
}
if let Some(error) = row.get("error").and_then(Value::as_str) {
tracing::warn!(palace = %id, "palace counts unavailable: {error}");
out.push(PalaceRow {
id: id.clone(),
name: id,
vector_count: 0,
drawer_count: 0,
last_write_at: None,
description: Some(error.to_string()),
kg_triple_count: 0,
node_count: 0,
edge_count: 0,
community_count: 0,
is_compacting: false,
counts_unknown: true,
});
continue;
}
tracing::warn!(palace = %id, "palace row carried neither counts nor a reason");
}
out
}
pub(super) fn project_events(raw: &Value, after_id: u64) -> (u64, Vec<MemoryEvent>) {
let Some(entries) = raw.get("entries").and_then(|v| v.as_array()) else {
return (after_id, Vec::new());
};
let mut highest = after_id;
let mut events = Vec::new();
for entry in entries.iter().rev() {
let id = entry
.get("id")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0);
if id <= after_id {
continue;
}
highest = highest.max(id);
if let Some(payload) = entry.get("payload")
&& let Some(event) = parse_memory_event(payload)
{
events.push(event);
}
}
(highest, events)
}