Resolve blocking commit deadlock and clippy warnings
This commit is contained in:
1 parent
afc7cbdad5
commit
cbd3d5801c
5 files changed
+55
-32
No files matched your search
@@ -152,7 +152,8 @@ async fn get_nvim_connection() -> Result<mpsc::Sender<NvimRequest>, String> {
|
|||||||
|
|
||||||
let (mut read_half, mut write_half) = tokio::io::split(stream);
|
let (mut read_half, mut write_half) = tokio::io::split(stream);
|
||||||
let (tx, mut rx) = mpsc::channel::<NvimRequest>(32);
|
let (tx, mut rx) = mpsc::channel::<NvimRequest>(32);
|
||||||
let pending_requests: Arc<Mutex<HashMap<String, oneshot::Sender<Result<rmpv::Value, String>>>>> = Arc::new(Mutex::new(HashMap::new()));
|
type PendingRequestsMap = Arc<Mutex<HashMap<String, oneshot::Sender<Result<rmpv::Value, String>>>>>;
|
||||||
|
let pending_requests: PendingRequestsMap = Arc::new(Mutex::new(HashMap::new()));
|
||||||
|
|
||||||
// Write task
|
// Write task
|
||||||
let pending_clone = Arc::clone(&pending_requests);
|
let pending_clone = Arc::clone(&pending_requests);
|
||||||
|
|||||||
@@ -441,7 +441,7 @@ impl MemoryHandler {
|
|||||||
for entity in req.entities {
|
for entity in req.entities {
|
||||||
if !entity.name.is_empty() {
|
if !entity.name.is_empty() {
|
||||||
if let Ok(idx) = self.state.search_index.read() {
|
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);
|
g.entities.insert(entity.name.clone(), entity);
|
||||||
}
|
}
|
||||||
@@ -744,7 +744,7 @@ impl MemoryHandler {
|
|||||||
acceptance_criteria: vec![],
|
acceptance_criteria: vec![],
|
||||||
};
|
};
|
||||||
if let Ok(idx) = self.state.search_index.read() {
|
if let Ok(idx) = self.state.search_index.read() {
|
||||||
let _ = idx.index_task(&task);
|
drop(idx.index_task(&task));
|
||||||
}
|
}
|
||||||
self.state.tasks.modify(|tasks| {
|
self.state.tasks.modify(|tasks| {
|
||||||
tasks.push(task);
|
tasks.push(task);
|
||||||
@@ -1663,7 +1663,7 @@ impl MemoryHandler {
|
|||||||
|
|
||||||
let elapsed = start_time.elapsed();
|
let elapsed = start_time.elapsed();
|
||||||
if method == "tools/call" {
|
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);
|
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
|
// Broadcast completion latency to the UI Activity Feed
|
||||||
|
|||||||
+5
-8
@@ -89,12 +89,10 @@ async fn index_committer_worker(state: Arc<MemoryState>) {
|
|||||||
loop {
|
loop {
|
||||||
tokio::time::sleep(Duration::from_secs(5)).await;
|
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||||
// Periodically commit the search index to persist inline indexing operations
|
// Periodically commit the search index to persist inline indexing operations
|
||||||
let state_clone = Arc::clone(&state);
|
let idx_opt = state.search_index.read().ok().map(|idx| idx.clone());
|
||||||
let _ = tokio::task::spawn_blocking(move || {
|
if let Some(idx) = idx_opt {
|
||||||
if let Ok(idx) = state_clone.search_index.read() {
|
let _ = idx.commit().await;
|
||||||
let _ = idx.commit();
|
}
|
||||||
}
|
|
||||||
}).await;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -219,6 +217,7 @@ async fn gate_set_handler(
|
|||||||
fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::Error>> {
|
fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::Error>> {
|
||||||
let rt = tokio::runtime::Runtime::new().unwrap();
|
let rt = tokio::runtime::Runtime::new().unwrap();
|
||||||
rt.block_on(async {
|
rt.block_on(async {
|
||||||
|
state.rebuild_index().await;
|
||||||
tokio::spawn(index_committer_worker(Arc::clone(&state)));
|
tokio::spawn(index_committer_worker(Arc::clone(&state)));
|
||||||
let app_state = Arc::new(AppState {
|
let app_state = Arc::new(AppState {
|
||||||
handler: Arc::new(MemoryHandler {
|
handler: Arc::new(MemoryHandler {
|
||||||
@@ -831,7 +830,5 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
|
|||||||
activity_tx: tokio::sync::broadcast::channel(100).0,
|
activity_tx: tokio::sync::broadcast::channel(100).0,
|
||||||
});
|
});
|
||||||
|
|
||||||
state.rebuild_index();
|
|
||||||
|
|
||||||
run_server(state)
|
run_server(state)
|
||||||
}
|
}
|
||||||
+14
-6
@@ -3,6 +3,9 @@ use std::sync::{Arc, Mutex};
|
|||||||
use tantivy::schema::*;
|
use tantivy::schema::*;
|
||||||
use tantivy::{Index, IndexReader, IndexWriter, ReloadPolicy, doc};
|
use tantivy::{Index, IndexReader, IndexWriter, ReloadPolicy, doc};
|
||||||
|
|
||||||
|
pub type SearchResultTuple = (String, String, String, String, f32);
|
||||||
|
|
||||||
|
#[derive(Clone)]
|
||||||
pub struct MemoryIndex {
|
pub struct MemoryIndex {
|
||||||
index: Index,
|
index: Index,
|
||||||
reader: IndexReader,
|
reader: IndexReader,
|
||||||
@@ -81,17 +84,22 @@ impl MemoryIndex {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn commit(&self) -> tantivy::Result<()> {
|
pub async fn commit(&self) -> tantivy::Result<()> {
|
||||||
let mut writer = self.writer.lock().unwrap();
|
let writer = Arc::clone(&self.writer);
|
||||||
writer.commit()?;
|
tokio::task::spawn_blocking(move || {
|
||||||
Ok(())
|
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(
|
pub fn search(
|
||||||
&self,
|
&self,
|
||||||
query: &str,
|
query: &str,
|
||||||
namespace: Option<&str>,
|
namespace: Option<&str>,
|
||||||
) -> tantivy::Result<Vec<(String, String, String, String, f32)>> {
|
) -> tantivy::Result<Vec<SearchResultTuple>> {
|
||||||
let searcher = self.reader.searcher();
|
let searcher = self.reader.searcher();
|
||||||
let query_parser = tantivy::query::QueryParser::for_index(
|
let query_parser = tantivy::query::QueryParser::for_index(
|
||||||
&self.index,
|
&self.index,
|
||||||
@@ -226,7 +234,7 @@ mod tests {
|
|||||||
};
|
};
|
||||||
index.index_adr(&adr).await.unwrap();
|
index.index_adr(&adr).await.unwrap();
|
||||||
|
|
||||||
index.commit().unwrap();
|
index.commit().await.unwrap();
|
||||||
index.reader.reload().unwrap();
|
index.reader.reload().unwrap();
|
||||||
|
|
||||||
// Test search
|
// Test search
|
||||||
|
|||||||
+31
-14
@@ -56,26 +56,43 @@ impl MemoryState {
|
|||||||
self.graph.modify(update_fn);
|
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) {
|
if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) {
|
||||||
let graph = self.graph.read();
|
let mut handles = Vec::new();
|
||||||
for e in graph.entities.values() {
|
|
||||||
let _ = new_idx.index_entity(e);
|
{
|
||||||
|
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 tasks = self.tasks.read();
|
||||||
let _ = new_idx.index_task(&t);
|
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() {
|
if let Ok(mut w) = self.search_index.write() {
|
||||||
*w = new_idx;
|
*w = new_idx;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in new issue
Block a user