diff --git a/server/src/search.rs b/server/src/search.rs index be56a0c..c8cecc3 100644 --- a/server/src/search.rs +++ b/server/src/search.rs @@ -213,69 +213,52 @@ impl MemoryIndex { }) } - pub fn index_batch( - &self, - entities: &[Entity], - tasks: &[Task], - snippets: &[Snippet], - adrs: &[Adr], - ) -> tokio::task::JoinHandle> { - let writer = Arc::clone(&self.writer); - let id_field = self.id_field; - let entities = entities.to_vec(); - let tasks = tasks.to_vec(); - let snippets = snippets.to_vec(); - let title_field = self.title_field; - let adrs = adrs.to_vec(); - let body_field = self.body_field; - let type_field = self.type_field; - let namespace_field = self.namespace_field; + pub fn add_entity_sync(&self, e: &Entity) { + if let Ok(writer) = self.writer.lock() { + let _ = writer.add_document(doc!( + self.id_field => e.name.as_str(), + self.title_field => e.name.as_str(), + self.body_field => e.observations.join(" "), + self.type_field => "entity", + self.namespace_field => e.namespace.as_str() + )); + } + } - tokio::task::spawn_blocking(move || { - let writer = writer.lock().unwrap_or_else(|e| e.into_inner()); + pub fn add_task_sync(&self, t: &Task) { + if let Ok(writer) = self.writer.lock() { + let _ = writer.add_document(doc!( + self.id_field => t.id.as_str(), + self.title_field => t.title.as_str(), + self.body_field => t.description.as_str(), + self.type_field => "task", + self.namespace_field => "global" + )); + } + } - 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() - ))?; - } + pub fn add_snippet_sync(&self, s: &Snippet) { + if let Ok(writer) = self.writer.lock() { + let _ = writer.add_document(doc!( + self.id_field => s.name.as_str(), + self.title_field => s.name.as_str(), + self.body_field => format!("{} {}", s.language, s.description), + self.type_field => "snippet", + self.namespace_field => "global" + )); + } + } - 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(()) - }) + pub fn add_adr_sync(&self, a: &Adr) { + if let Ok(writer) = self.writer.lock() { + let _ = writer.add_document(doc!( + self.id_field => a.id.as_str(), + self.title_field => a.title.as_str(), + self.body_field => format!("{} {} {}", a.context, a.decision, a.consequence), + self.type_field => "adr", + self.namespace_field => "global" + )); + } } } diff --git a/server/src/state.rs b/server/src/state.rs index a0b8c6e..19624ec 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -3,7 +3,7 @@ use crate::search::MemoryIndex; use crate::store::Store; use std::collections::HashMap; use std::path::PathBuf; -use std::sync::RwLock; +use std::sync::{Arc, RwLock}; pub struct MemoryState { pub base_dir: PathBuf, @@ -60,15 +60,33 @@ impl MemoryState { self.graph.modify(update_fn); } - pub async fn rebuild_index(&self) { + pub async fn rebuild_index(self: &Arc) { if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) { - let entities: Vec = self.graph.read_with(|g| g.entities.values().cloned().collect()); - let tasks = self.tasks.read_with(|t| t.clone()); - let snippets = self.snippets.read_with(|s| s.clone()); - let adrs = self.adrs.read_with(|a| a.clone()); - - let handle = new_idx.index_batch(&entities, &tasks, &snippets, &adrs); - let _ = handle.await; + let state = Arc::clone(self); + let idx = new_idx.clone(); + + tokio::task::spawn_blocking(move || { + state.graph.read_with(|g| { + for e in g.entities.values() { + idx.add_entity_sync(e); + } + }); + state.tasks.read_with(|tasks| { + for t in tasks { + idx.add_task_sync(t); + } + }); + state.snippets.read_with(|snippets| { + for s in snippets { + idx.add_snippet_sync(s); + } + }); + state.adrs.read_with(|adrs| { + for a in adrs { + idx.add_adr_sync(a); + } + }); + }).await.unwrap(); let _ = new_idx.commit().await; if let Ok(mut w) = self.search_index.write() {