From a9885a65d7eda561a5f1e2a56f883a5b34d42323 Mon Sep 17 00:00:00 2001 From: Riz Ashraf Date: Mon, 21 Sep 2026 05:56:58 +0100 Subject: [PATCH] Optimize full graph rebuilds with batch indexing to prevent thread exhaustion, and fix adr/snippet search availability before restart --- server/src/handlers.rs | 51 ++++++++++++++++++++++++----------- server/src/search.rs | 61 ++++++++++++++++++++++++++++++++++++++++++ server/src/state.rs | 37 +++++-------------------- 3 files changed, 102 insertions(+), 47 deletions(-) diff --git a/server/src/handlers.rs b/server/src/handlers.rs index dd60ac9..abf58e1 100644 --- a/server/src/handlers.rs +++ b/server/src/handlers.rs @@ -978,19 +978,27 @@ impl MemoryHandler { } "store_snippet" => { let req = parse_tool!(args.clone(), id, StoreSnippetTool); + let snippet = Snippet { + name: req.name.clone(), + language: req.language, + code: req.code, + description: req.description, + updated_at: SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_secs(), + }; + + let s_clone = snippet.clone(); self.state.snippets.modify(|snippets| { snippets.retain(|s| s.name != req.name); - snippets.push(Snippet { - name: req.name.clone(), - language: req.language, - code: req.code, - description: req.description, - updated_at: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_secs(), - }); + snippets.push(s_clone); }); + + if let Ok(idx) = self.state.search_index.read() { + drop(idx.index_snippet(&snippet)); + } + Ok(format!("Snippet '{}' stored.", req.name).to_string()) } "search_snippets" => { @@ -1025,11 +1033,13 @@ impl MemoryHandler { } "log_decision" => { let req = parse_tool!(args.clone(), id, LogDecisionTool); - let mut id = String::new(); + let mut adr_id = String::new(); + let mut new_adr = None; + self.state.adrs.modify(|adrs| { - id = format!("ADR-{:04}", adrs.len() + 1); - adrs.push(Adr { - id: id.clone(), + adr_id = format!("ADR-{:04}", adrs.len() + 1); + let a = Adr { + id: adr_id.clone(), title: req.title, context: req.context, decision: req.decision, @@ -1038,9 +1048,18 @@ impl MemoryHandler { .duration_since(UNIX_EPOCH) .unwrap() .as_secs(), - }); + }; + new_adr = Some(a.clone()); + adrs.push(a); }); - Ok(format!("Decision logged as {}", id).to_string()) + + if let Some(adr) = new_adr { + if let Ok(idx) = self.state.search_index.read() { + drop(idx.index_adr(&adr)); + } + } + + Ok(format!("Decision logged as {}", adr_id).to_string()) } "query_decisions" => { let req = parse_tool!(args.clone(), id, QueryDecisionsTool); diff --git a/server/src/search.rs b/server/src/search.rs index e89ab6f..6d5ce65 100644 --- a/server/src/search.rs +++ b/server/src/search.rs @@ -180,6 +180,67 @@ impl MemoryIndex { Ok(()) }) } + + pub fn index_batch( + &self, + entities: Vec, + tasks: Vec, + snippets: Vec, + adrs: Vec, + ) -> tokio::task::JoinHandle> { + let writer = Arc::clone(&self.writer); + let id_field = self.id_field; + let title_field = self.title_field; + let body_field = self.body_field; + let type_field = self.type_field; + let namespace_field = self.namespace_field; + + tokio::task::spawn_blocking(move || { + let writer = writer.lock().unwrap(); + + for e in entities { + writer.add_document(doc!( + id_field => e.name.clone(), + title_field => e.name.clone(), + body_field => e.observations.join(" "), + type_field => "entity", + namespace_field => e.namespace.clone() + ))?; + } + + for t in tasks { + writer.add_document(doc!( + id_field => t.id.clone(), + title_field => t.title.clone(), + body_field => t.description.clone(), + type_field => "task", + namespace_field => "global" + ))?; + } + + for s in snippets { + writer.add_document(doc!( + id_field => s.name.clone(), + title_field => s.name.clone(), + body_field => format!("{} {}", s.language, s.description), + type_field => "snippet", + namespace_field => "global" + ))?; + } + + for a in adrs { + writer.add_document(doc!( + id_field => a.id.clone(), + title_field => a.title.clone(), + body_field => format!("{} {} {}", a.context, a.decision, a.consequence), + type_field => "adr", + namespace_field => "global" + ))?; + } + + Ok(()) + }) + } } #[cfg(test)] diff --git a/server/src/state.rs b/server/src/state.rs index 6db9da3..f1c24c1 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -58,39 +58,14 @@ impl MemoryState { pub async fn rebuild_index(&self) { if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) { - let mut handles = Vec::new(); - { - let graph = self.graph.read(); - for e in graph.entities.values() { - handles.push(new_idx.index_entity(e)); - } - } + let entities = self.graph.read().entities.values().cloned().collect(); + let tasks = self.tasks.read().clone(); + let snippets = self.snippets.read().clone(); + let adrs = self.adrs.read().clone(); - { - 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 { - handles.push(new_idx.index_snippet(&s)); - } - } - - { - let adrs = self.adrs.read(); - for a in adrs { - handles.push(new_idx.index_adr(&a)); - } - } - - for handle in handles { - let _ = handle.await; - } + let handle = new_idx.index_batch(entities, tasks, snippets, adrs); + let _ = handle.await; let _ = new_idx.commit().await; if let Ok(mut w) = self.search_index.write() {