use std::path::{Path, PathBuf};
use std::time::Duration;
use anyhow::Context;
use serde::Deserialize;
use serde_json::{Value, json};
use crate::monitor::dashboard::{IndexRow, SearchData};
use crate::search_rpc::{
METHOD_HEALTH, METHOD_INDEX_REINDEX, METHOD_INDEX_STATUS, METHOD_INDEXES_LIST, call_at,
};
#[cfg(test)]
#[path = "search_client_tests.rs"]
mod tests;
const REQUEST_TIMEOUT: Duration = Duration::from_secs(3);
const STREAM_FRAME_TIMEOUT: Duration = Duration::from_secs(300);
const METHOD_GRAPH_STATS: &str = "search.graph.stats";
const METHOD_LOGS_TAIL: &str = "search.logs.tail";
const METHOD_QUERY: &str = "search.query";
const METHOD_INDEX_REINDEX_STREAM: &str = "search.index.reindex.stream";
pub fn resolve_search_socket() -> anyhow::Result<PathBuf> {
crate::search_rpc::search_socket()
}
#[derive(Debug, Deserialize)]
struct HealthWire {
version: String,
#[serde(default)]
uptime_secs: u64,
}
#[derive(Debug, Deserialize)]
struct LogsTailWire {
#[serde(default)]
lines: Vec<String>,
}
#[derive(Debug, Deserialize)]
struct IndexListWire {
#[serde(default)]
indexes: Vec<String>,
}
#[derive(Debug, Deserialize)]
struct GraphStatsWire {
#[serde(default)]
node_count: u64,
#[serde(default)]
edge_count: u64,
#[serde(default)]
edge_kinds: std::collections::HashMap<String, u64>,
}
#[derive(Debug, Deserialize)]
struct IndexStatusWire {
#[serde(default)]
root_path: String,
#[serde(default)]
chunk_count: u64,
#[serde(default)]
disk_bytes: Option<u64>,
#[serde(default)]
last_indexed: Option<chrono::DateTime<chrono::Utc>>,
}
#[derive(Debug, Clone)]
pub struct SearchClient {
socket: PathBuf,
}
impl SearchClient {
pub fn new(socket: impl Into<PathBuf>) -> Self {
Self {
socket: socket.into(),
}
}
pub fn resolve() -> anyhow::Result<Self> {
Ok(Self::new(resolve_search_socket()?))
}
pub fn socket(&self) -> &Path {
&self.socket
}
async fn call<T: serde::de::DeserializeOwned>(
&self,
method: &str,
params: Value,
) -> anyhow::Result<T> {
let raw = call_at(&self.socket, method, params, REQUEST_TIMEOUT).await?;
serde_json::from_value(raw).with_context(|| {
format!(
"decode {method} from the trusty-search daemon at {}",
self.socket.display()
)
})
}
pub async fn fetch_all(&self) -> anyhow::Result<SearchData> {
let health: HealthWire = self.call(METHOD_HEALTH, json!({})).await?;
let list: IndexListWire = self.call(METHOD_INDEXES_LIST, json!({})).await?;
let mut indexes = Vec::with_capacity(list.indexes.len());
for id in list.indexes {
let params = json!({ "index_id": id });
let status = self
.call::<IndexStatusWire>(METHOD_INDEX_STATUS, params.clone())
.await;
let mut row = match status {
Ok(status) => IndexRow {
id: id.clone(),
chunk_count: status.chunk_count,
root_path: status.root_path,
disk_bytes: status.disk_bytes,
last_indexed: status.last_indexed,
..Default::default()
},
Err(e) => {
tracing::warn!("index status probe failed for {id}: {e}");
IndexRow {
id: id.clone(),
..Default::default()
}
}
};
match self
.call::<GraphStatsWire>(METHOD_GRAPH_STATS, params)
.await
{
Ok(stats) => {
row.node_count = stats.node_count;
row.edge_count = stats.edge_count;
let mut kinds: Vec<(String, u64)> = stats.edge_kinds.into_iter().collect();
kinds.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
row.edge_kinds = kinds;
}
Err(e) => tracing::debug!("graph stats probe failed for {id}: {e}"),
}
indexes.push(row);
}
indexes.sort_by(|a, b| a.id.cmp(&b.id));
Ok(SearchData {
version: health.version,
uptime_secs: health.uptime_secs,
indexes,
})
}
pub async fn logs_tail(&self, n: usize) -> anyhow::Result<Vec<String>> {
let wire: LogsTailWire = self.call(METHOD_LOGS_TAIL, json!({ "n": n })).await?;
Ok(wire.lines)
}
pub async fn reindex(&self, id: &str) -> anyhow::Result<()> {
call_at(
&self.socket,
METHOD_INDEX_REINDEX,
json!({ "index_id": id }),
REQUEST_TIMEOUT,
)
.await?;
Ok(())
}
pub async fn search(
&self,
id: &str,
query: &str,
top_k: usize,
) -> anyhow::Result<Vec<SearchHit>> {
let params = json!({ "index_id": id, "body": { "text": query, "top_k": top_k } });
let raw = call_at(&self.socket, METHOD_QUERY, params, REQUEST_TIMEOUT).await?;
Ok(parse_search_hits(&raw))
}
pub async fn reindex_stream(&self, id: &str, tx: tokio::sync::mpsc::Sender<ReindexEvent>) {
if let Err(e) = self.reindex_stream_inner(id, &tx).await {
let _ = tx.send(ReindexEvent::Failed(format!("{e:#}"))).await;
}
}
async fn reindex_stream_inner(
&self,
id: &str,
tx: &tokio::sync::mpsc::Sender<ReindexEvent>,
) -> anyhow::Result<()> {
self.reindex(id).await?;
let request = json!({
"jsonrpc": "2.0",
"id": 1,
"method": METHOD_INDEX_REINDEX_STREAM,
"params": { "index_id": id },
"stream": true,
});
let open = crate::uds::send_framed_stream_request::<_, Value>(
&self.socket,
&request,
STREAM_FRAME_TIMEOUT,
)
.await;
let mut stream = open.with_context(|| {
format!(
"open {METHOD_INDEX_REINDEX_STREAM} on the trusty-search daemon at {}",
self.socket.display()
)
})?;
while let Some(item) = stream.next_frame().await {
let event = parse_reindex_event(&item?);
let terminal = matches!(event, ReindexEvent::Complete { .. });
if tx.send(event).await.is_err() || terminal {
return Ok(()); }
}
Ok(())
}
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct SearchHit {
pub file: String,
pub line: usize,
pub snippet: String,
}
pub fn parse_search_hits(raw: &serde_json::Value) -> Vec<SearchHit> {
let Some(results) = raw.get("results").and_then(|v| v.as_array()) else {
return Vec::new();
};
results
.iter()
.map(|item| {
let file = item
.get("file")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_string();
let line = item.get("start_line").and_then(|v| v.as_u64()).unwrap_or(0) as usize;
let snippet = item
.get("compact_snippet")
.and_then(|v| v.as_str())
.or_else(|| item.get("content").and_then(|v| v.as_str()))
.unwrap_or_default()
.lines()
.next()
.unwrap_or_default()
.trim()
.to_string();
SearchHit {
file,
line,
snippet,
}
})
.collect()
}
#[derive(Debug, Clone, PartialEq)]
pub enum ReindexEvent {
Started {
total_files: u64,
},
Progress {
indexed: u64,
total_files: u64,
},
Complete {
total_chunks: u64,
status: String,
},
Failed(String),
}
pub fn parse_reindex_event(value: &serde_json::Value) -> ReindexEvent {
let kind = value.get("event").and_then(|v| v.as_str()).unwrap_or("");
let u64_of = |key: &str| value.get(key).and_then(|v| v.as_u64()).unwrap_or(0);
match kind {
"start" => ReindexEvent::Started {
total_files: u64_of("total_files"),
},
"complete" => ReindexEvent::Complete {
total_chunks: u64_of("total_chunks"),
status: value
.get("status")
.and_then(|v| v.as_str())
.unwrap_or("complete")
.to_string(),
},
"error" => ReindexEvent::Failed(
value
.get("message")
.and_then(|v| v.as_str())
.unwrap_or("reindex error")
.to_string(),
),
_ => ReindexEvent::Progress {
indexed: u64_of("indexed"),
total_files: u64_of("total_files"),
},
}
}