perf(search): eliminate full database cloning during tantivy index rebuilds by moving iteration into a synchronized blocking thread

This commit is contained in:
Riz Ashraf committed 2026-09-21 16:42:14 +01:00
1 parent 8934857635
commit 1af21b3aec
2 files changed
+69 -68

No files matched your search

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