feat(server): refactor handlers, router, state management, and memory tools

This commit is contained in:
Riz Ashraf committed 2026-10-02 07:27:37 +01:00
1 parent 87ddb01063
commit a083719cf1
36 files changed
+1899 -597

No files matched your search

+100 -81
View File
@@ -4,20 +4,21 @@
)]
mod api;
mod clipboard_watcher;
pub mod db;
pub mod embedding;
pub mod error;
mod handlers;
pub mod indexer;
mod mcp;
mod models;
pub mod ollama;
mod router;
mod search;
pub mod embedding;
pub mod indexer;
pub mod vector_db;
mod state;
mod store;
mod tools;
mod clipboard_watcher;
pub mod vector_db;
mod watcher;
use crate::api::rest::GateSetReq;
@@ -180,7 +181,10 @@ async fn condense_graph_worker(state: Arc<MemoryState>) {
let to_remove = snippets.len() - (threshold / 2);
let removed: Vec<_> = snippets.drain(0..to_remove).collect();
for r in removed {
condensed_snippet_content.push_str(&format!("Name: {}\nDesc: {}\nCode: {}\n", r.name, r.description, r.code));
condensed_snippet_content.push_str(&format!(
"Name: {}\nDesc: {}\nCode: {}\n",
r.name, r.description, r.code
));
}
}
});
@@ -210,7 +214,7 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
state_for_index.rebuild_index().await;
tracing::info!("Index rebuild complete.");
});
// Start the global codebase indexer
crate::indexer::start_background_indexer(Arc::clone(&state)).await;
@@ -258,31 +262,37 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
let mut buf = [0; 4096];
loop {
if let Ok((len, _addr)) = socket.recv_from(&mut buf).await
&& let Ok(payload) = serde_json::from_slice::<crate::models::TerminalHistory>(&buf[..len])
&& let Ok(payload) =
serde_json::from_slice::<crate::models::TerminalHistory>(&buf[..len])
{
udp_state.handler.state.telemetry.terminal_history.modify(|history| {
udp_state
.handler
.state
.telemetry
.terminal_history
.modify(|history| {
history.push_front(payload.clone());
if history.len() > 100 {
history.pop_back();
}
});
let ws_msg = serde_json::json!({
"type": "terminal_telemetry",
"data": payload
});
let msg_str = ws_msg.to_string();
let senders: Vec<_> = udp_state
.clients
.read()
.unwrap_or_else(|e| e.into_inner())
.values()
.cloned()
.collect();
for tx in senders {
let _ = tx.try_send(msg_str.clone());
}
let ws_msg = serde_json::json!({
"type": "terminal_telemetry",
"data": payload
});
let msg_str = ws_msg.to_string();
let senders: Vec<_> = udp_state
.clients
.read()
.unwrap_or_else(|e| e.into_inner())
.values()
.cloned()
.collect();
for tx in senders {
let _ = tx.try_send(msg_str.clone());
}
}
}
}
@@ -296,64 +306,70 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
let mut buf = [0; 4096];
loop {
if let Ok((len, _addr)) = socket.recv_from(&mut buf).await
&& let Ok(payload) = serde_json::from_slice::<crate::api::telemetry::NvimTelemetry>(&buf[..len])
&& let Ok(payload) =
serde_json::from_slice::<crate::api::telemetry::NvimTelemetry>(&buf[..len])
{
// 1. Legacy disk write for active_nvim.txt
if payload.event == "FocusGained" || payload.event == "BufEnter" || payload.event == "VimEnter" {
let session = &payload.session_id;
let is_unix_socket = session.starts_with('/') || session.starts_with('~');
if is_unix_socket {
let wsl_path = "\\\\wsl.localhost\\Ubuntu\\home\\riz\\.gemini\\active_nvim.txt";
let _ = tokio::fs::write(wsl_path, session).await;
} else {
let profile = std::env::var("USERPROFILE").unwrap_or_else(|_| "C:\\Users\\reazul.ashraf".into());
let win_path = format!("{}\\.gemini\\active_nvim.txt", profile);
let _ = tokio::fs::write(&win_path, session).await;
}
// 1. Legacy disk write for active_nvim.txt
if payload.event == "FocusGained"
|| payload.event == "BufEnter"
|| payload.event == "VimEnter"
{
let session = &payload.session_id;
let is_unix_socket = session.starts_with('/') || session.starts_with('~');
if is_unix_socket {
let wsl_path =
"\\\\wsl.localhost\\Ubuntu\\home\\riz\\.gemini\\active_nvim.txt";
let _ = tokio::fs::write(wsl_path, session).await;
} else {
let profile = std::env::var("USERPROFILE")
.unwrap_or_else(|_| "C:\\Users\\reazul.ashraf".into());
let win_path = format!("{}\\.gemini\\active_nvim.txt", profile);
let _ = tokio::fs::write(&win_path, session).await;
}
}
// 2. Broadcast to UI
let ws_msg = serde_json::json!({
"type": "nvim_telemetry",
"data": payload
});
let msg_str = ws_msg.to_string();
let senders: Vec<_> = nvim_udp_state
.clients
.read()
.unwrap_or_else(|e| e.into_inner())
.values()
.cloned()
.collect();
for tx in senders {
let _ = tx.try_send(msg_str.clone());
}
// 2. Broadcast to UI
let ws_msg = serde_json::json!({
"type": "nvim_telemetry",
"data": payload
});
let msg_str = ws_msg.to_string();
// 3. Event bus trigger for auto-save hook
if payload.event == "BufWritePost"
&& let Some(ref file_path) = payload.file
{
let normalized_file = file_path.replace("\\", "/");
let topic = format!("nvim:save:{}", normalized_file);
let event = crate::state::GenericEvent {
topic,
session_id: Some(payload.session_id.clone()),
payload: serde_json::json!(&payload),
};
let _ = nvim_udp_state.handler.state.event_bus_tx.send(event);
}
let senders: Vec<_> = nvim_udp_state
.clients
.read()
.unwrap_or_else(|e| e.into_inner())
.values()
.cloned()
.collect();
for tx in senders {
let _ = tx.try_send(msg_str.clone());
}
// 4. Interactive Agent UI Events
if payload.event.starts_with("agent_") {
let topic = format!("nvim:ui:{}", payload.event);
let event = crate::state::GenericEvent {
topic,
session_id: Some(payload.session_id.clone()),
payload: serde_json::json!(&payload),
};
let _ = nvim_udp_state.handler.state.event_bus_tx.send(event);
}
// 3. Event bus trigger for auto-save hook
if payload.event == "BufWritePost"
&& let Some(ref file_path) = payload.file
{
let normalized_file = file_path.replace("\\", "/");
let topic = format!("nvim:save:{}", normalized_file);
let event = crate::state::GenericEvent {
topic,
session_id: Some(payload.session_id.clone()),
payload: serde_json::json!(&payload),
};
let _ = nvim_udp_state.handler.state.event_bus_tx.send(event);
}
// 4. Interactive Agent UI Events
if payload.event.starts_with("agent_") {
let topic = format!("nvim:ui:{}", payload.event);
let event = crate::state::GenericEvent {
topic,
session_id: Some(payload.session_id.clone()),
payload: serde_json::json!(&payload),
};
let _ = nvim_udp_state.handler.state.event_bus_tx.send(event);
}
}
}
}
@@ -362,7 +378,10 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
let app = api::setup::create_router(app_state);
let port_str = std::env::var("MCP_PORT").unwrap_or_else(|_| "3000".to_string());
tracing::info!("MCP Memory Server running on http://127.0.0.1:{}/ws", port_str);
tracing::info!(
"MCP Memory Server running on http://127.0.0.1:{}/ws",
port_str
);
let addr: std::net::SocketAddr = format!("127.0.0.1:{}", port_str)
.parse()
.expect("Invalid bind address");
@@ -546,7 +565,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
rt.block_on(async {
let state = Arc::new(MemoryState::new(&base.to_string_lossy()));
// Initialize Qdrant VectorDB (default local URL)
match crate::vector_db::VectorDB::new("http://localhost:6334", "mcp_memory").await {
Ok(vdb) => {
@@ -557,7 +576,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing::warn!("Failed to initialize Qdrant vector database: {}. Vector search will fallback to manual embedding loop. (Is Qdrant running on localhost:6334?)", e);
}
}
if let Err(e) = run_server(state).await {
tracing::error!("Server error: {}", e);
}