feat(mcp): Add 5 developer enhancements (git, logs, clipboard watcher, ast skeleton, nvim ghost text)
This commit is contained in:
1 parent
ccc268d807
commit
0e9748f797
27 files changed
+1401
-174
No files matched your search
+25
-9
@@ -15,6 +15,7 @@ pub mod embedding;
|
||||
mod state;
|
||||
mod store;
|
||||
mod tools;
|
||||
mod clipboard_watcher;
|
||||
|
||||
use crate::api::rest::GateSetReq;
|
||||
use crate::router::MemoryHandler;
|
||||
@@ -88,6 +89,7 @@ pub struct AppState {
|
||||
handler: Arc<MemoryHandler>,
|
||||
clients: RwLock<HashMap<String, mpsc::Sender<String>>>,
|
||||
next_id: AtomicUsize,
|
||||
pub shutdown_tx: std::sync::Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
|
||||
}
|
||||
|
||||
async fn ttl_sweeper_worker(state: Arc<MemoryState>) {
|
||||
@@ -124,13 +126,20 @@ async fn index_committer_worker(state: Arc<MemoryState>) {
|
||||
}
|
||||
|
||||
async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::Error>> {
|
||||
state.rebuild_index().await;
|
||||
let state_for_index = Arc::clone(&state);
|
||||
tokio::spawn(async move {
|
||||
state_for_index.rebuild_index().await;
|
||||
tracing::info!("Index rebuild complete.");
|
||||
});
|
||||
tokio::spawn(index_committer_worker(Arc::clone(&state)));
|
||||
tokio::spawn(ttl_sweeper_worker(Arc::clone(&state)));
|
||||
crate::clipboard_watcher::spawn_watcher(Arc::clone(&state));
|
||||
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
|
||||
let app_state = Arc::new(AppState {
|
||||
handler: Arc::new(MemoryHandler::new(Arc::clone(&state))),
|
||||
clients: RwLock::new(HashMap::new()),
|
||||
next_id: AtomicUsize::new(1),
|
||||
shutdown_tx: std::sync::Mutex::new(Some(shutdown_tx)),
|
||||
});
|
||||
|
||||
let app_state_clone = Arc::clone(&app_state);
|
||||
@@ -156,10 +165,11 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
|
||||
}
|
||||
});
|
||||
|
||||
// UDP Telemetry Listener on port 3001
|
||||
// UDP Telemetry Listener
|
||||
let udp_state = Arc::clone(&app_state);
|
||||
tokio::spawn(async move {
|
||||
if let Ok(socket) = tokio::net::UdpSocket::bind("127.0.0.1:3001").await {
|
||||
let port1 = std::env::var("MCP_UDP_PORT1").unwrap_or_else(|_| "3001".to_string());
|
||||
if let Ok(socket) = tokio::net::UdpSocket::bind(format!("127.0.0.1:{}", port1)).await {
|
||||
let mut buf = [0; 4096];
|
||||
loop {
|
||||
if let Ok((len, _addr)) = socket.recv_from(&mut buf).await
|
||||
@@ -193,10 +203,11 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
|
||||
}
|
||||
});
|
||||
|
||||
// UDP Neovim Telemetry Listener on port 3002
|
||||
// UDP Neovim Telemetry Listener
|
||||
let nvim_udp_state = Arc::clone(&app_state);
|
||||
tokio::spawn(async move {
|
||||
if let Ok(socket) = tokio::net::UdpSocket::bind("127.0.0.1:3002").await {
|
||||
let port2 = std::env::var("MCP_UDP_PORT2").unwrap_or_else(|_| "3002".to_string());
|
||||
if let Ok(socket) = tokio::net::UdpSocket::bind(format!("127.0.0.1:{}", port2)).await {
|
||||
let mut buf = [0; 4096];
|
||||
loop {
|
||||
if let Ok((len, _addr)) = socket.recv_from(&mut buf).await
|
||||
@@ -265,9 +276,9 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
|
||||
|
||||
let app = api::setup::create_router(app_state);
|
||||
|
||||
tracing::info!("MCP Memory Server running on http://127.0.0.1:3000/ws");
|
||||
let addr = std::env::var("MCP_PORT").unwrap_or_else(|_| "3000".to_string());
|
||||
let addr: std::net::SocketAddr = format!("127.0.0.1:{}", addr)
|
||||
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);
|
||||
let addr: std::net::SocketAddr = format!("127.0.0.1:{}", port_str)
|
||||
.parse()
|
||||
.expect("Invalid bind address");
|
||||
|
||||
@@ -282,7 +293,12 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
|
||||
return Ok(());
|
||||
}
|
||||
};
|
||||
if let Err(e) = axum::serve(listener, app.into_make_service()).await {
|
||||
if let Err(e) = axum::serve(listener, app.into_make_service())
|
||||
.with_graceful_shutdown(async move {
|
||||
let _ = shutdown_rx.await;
|
||||
})
|
||||
.await
|
||||
{
|
||||
let log_path = dirs::home_dir()
|
||||
.unwrap_or_default()
|
||||
.join(".gemini/mcp_memory/daemon_error.log");
|
||||
|
||||
Reference in new issue
Block a user