diff --git a/nvim-core/src/lib.rs b/nvim-core/src/lib.rs index a36e438..2d122bc 100644 --- a/nvim-core/src/lib.rs +++ b/nvim-core/src/lib.rs @@ -152,7 +152,8 @@ async fn get_nvim_connection() -> Result, String> { let (mut read_half, mut write_half) = tokio::io::split(stream); let (tx, mut rx) = mpsc::channel::(32); - let pending_requests: Arc>>>> = Arc::new(Mutex::new(HashMap::new())); + type PendingRequestsMap = Arc>>>>; + let pending_requests: PendingRequestsMap = Arc::new(Mutex::new(HashMap::new())); // Write task let pending_clone = Arc::clone(&pending_requests); diff --git a/server/src/handlers.rs b/server/src/handlers.rs index 74d5e78..f8274ef 100644 --- a/server/src/handlers.rs +++ b/server/src/handlers.rs @@ -441,7 +441,7 @@ impl MemoryHandler { for entity in req.entities { if !entity.name.is_empty() { if let Ok(idx) = self.state.search_index.read() { - let _ = idx.index_entity(&entity); + drop(idx.index_entity(&entity)); } g.entities.insert(entity.name.clone(), entity); } @@ -744,7 +744,7 @@ impl MemoryHandler { acceptance_criteria: vec![], }; if let Ok(idx) = self.state.search_index.read() { - let _ = idx.index_task(&task); + drop(idx.index_task(&task)); } self.state.tasks.modify(|tasks| { tasks.push(task); @@ -1663,7 +1663,7 @@ impl MemoryHandler { let elapsed = start_time.elapsed(); if method == "tools/call" { - let is_error = response.as_ref().map_or(false, |r| r.get("error").is_some() || r.get("result").and_then(|res| res.get("isError")).and_then(|e| e.as_bool()).unwrap_or(false)); + let is_error = response.as_ref().is_some_and(|r| r.get("error").is_some() || r.get("result").and_then(|res| res.get("isError")).and_then(|e| e.as_bool()).unwrap_or(false)); tracing::info!("<<< [Server] MCP tool call {} (id: {}) completed in {:?} [Error: {}]", tool_name, id_clone, elapsed, is_error); // Broadcast completion latency to the UI Activity Feed diff --git a/server/src/main.rs b/server/src/main.rs index d9b62a4..c5216d8 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -89,12 +89,10 @@ async fn index_committer_worker(state: Arc) { loop { tokio::time::sleep(Duration::from_secs(5)).await; // Periodically commit the search index to persist inline indexing operations - let state_clone = Arc::clone(&state); - let _ = tokio::task::spawn_blocking(move || { - if let Ok(idx) = state_clone.search_index.read() { - let _ = idx.commit(); - } - }).await; + let idx_opt = state.search_index.read().ok().map(|idx| idx.clone()); + if let Some(idx) = idx_opt { + let _ = idx.commit().await; + } } } @@ -219,6 +217,7 @@ async fn gate_set_handler( fn run_server(state: Arc) -> Result<(), Box> { let rt = tokio::runtime::Runtime::new().unwrap(); rt.block_on(async { + state.rebuild_index().await; tokio::spawn(index_committer_worker(Arc::clone(&state))); let app_state = Arc::new(AppState { handler: Arc::new(MemoryHandler { @@ -831,7 +830,5 @@ fn main() -> Result<(), Box> { activity_tx: tokio::sync::broadcast::channel(100).0, }); - state.rebuild_index(); - run_server(state) } diff --git a/server/src/search.rs b/server/src/search.rs index 8599375..e89ab6f 100644 --- a/server/src/search.rs +++ b/server/src/search.rs @@ -3,6 +3,9 @@ use std::sync::{Arc, Mutex}; use tantivy::schema::*; use tantivy::{Index, IndexReader, IndexWriter, ReloadPolicy, doc}; +pub type SearchResultTuple = (String, String, String, String, f32); + +#[derive(Clone)] pub struct MemoryIndex { index: Index, reader: IndexReader, @@ -81,17 +84,22 @@ impl MemoryIndex { }) } - pub fn commit(&self) -> tantivy::Result<()> { - let mut writer = self.writer.lock().unwrap(); - writer.commit()?; - Ok(()) + pub async fn commit(&self) -> tantivy::Result<()> { + let writer = Arc::clone(&self.writer); + tokio::task::spawn_blocking(move || { + let mut writer = writer.lock().unwrap(); + writer.commit()?; + Ok(()) + }) + .await + .unwrap_or_else(|_| Err(tantivy::TantivyError::SystemError("Commit task panicked".to_string()))) } pub fn search( &self, query: &str, namespace: Option<&str>, - ) -> tantivy::Result> { + ) -> tantivy::Result> { let searcher = self.reader.searcher(); let query_parser = tantivy::query::QueryParser::for_index( &self.index, @@ -226,7 +234,7 @@ mod tests { }; index.index_adr(&adr).await.unwrap(); - index.commit().unwrap(); + index.commit().await.unwrap(); index.reader.reload().unwrap(); // Test search diff --git a/server/src/state.rs b/server/src/state.rs index 6038495..6db9da3 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -56,26 +56,43 @@ impl MemoryState { self.graph.modify(update_fn); } - pub fn rebuild_index(&self) { + pub async fn rebuild_index(&self) { if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) { - let graph = self.graph.read(); - for e in graph.entities.values() { - let _ = new_idx.index_entity(e); + let mut handles = Vec::new(); + + { + let graph = self.graph.read(); + for e in graph.entities.values() { + handles.push(new_idx.index_entity(e)); + } } - let tasks = self.tasks.read(); - for t in tasks { - let _ = new_idx.index_task(&t); + { + let tasks = self.tasks.read(); + for t in tasks { + handles.push(new_idx.index_task(&t)); + } } - let snippets = self.snippets.read(); - for s in snippets { - let _ = new_idx.index_snippet(&s); + + { + let snippets = self.snippets.read(); + for s in snippets { + handles.push(new_idx.index_snippet(&s)); + } } - let adrs = self.adrs.read(); - for a in adrs { - let _ = new_idx.index_adr(&a); + + { + let adrs = self.adrs.read(); + for a in adrs { + handles.push(new_idx.index_adr(&a)); + } } - let _ = new_idx.commit(); + + for handle in handles { + let _ = handle.await; + } + + let _ = new_idx.commit().await; if let Ok(mut w) = self.search_index.write() { *w = new_idx; }