Optimize full graph rebuilds with batch indexing to prevent thread exhaustion, and fix adr/snippet search availability before restart
This commit is contained in:
1 parent
84672a00b9
commit
a9885a65d7
3 files changed
+102
-47
No files matched your search
+35
-16
@@ -978,19 +978,27 @@ impl MemoryHandler {
|
|||||||
}
|
}
|
||||||
"store_snippet" => {
|
"store_snippet" => {
|
||||||
let req = parse_tool!(args.clone(), id, StoreSnippetTool);
|
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| {
|
self.state.snippets.modify(|snippets| {
|
||||||
snippets.retain(|s| s.name != req.name);
|
snippets.retain(|s| s.name != req.name);
|
||||||
snippets.push(Snippet {
|
snippets.push(s_clone);
|
||||||
name: req.name.clone(),
|
|
||||||
language: req.language,
|
|
||||||
code: req.code,
|
|
||||||
description: req.description,
|
|
||||||
updated_at: SystemTime::now()
|
|
||||||
.duration_since(UNIX_EPOCH)
|
|
||||||
.unwrap()
|
|
||||||
.as_secs(),
|
|
||||||
});
|
|
||||||
});
|
});
|
||||||
|
|
||||||
|
if let Ok(idx) = self.state.search_index.read() {
|
||||||
|
drop(idx.index_snippet(&snippet));
|
||||||
|
}
|
||||||
|
|
||||||
Ok(format!("Snippet '{}' stored.", req.name).to_string())
|
Ok(format!("Snippet '{}' stored.", req.name).to_string())
|
||||||
}
|
}
|
||||||
"search_snippets" => {
|
"search_snippets" => {
|
||||||
@@ -1025,11 +1033,13 @@ impl MemoryHandler {
|
|||||||
}
|
}
|
||||||
"log_decision" => {
|
"log_decision" => {
|
||||||
let req = parse_tool!(args.clone(), id, LogDecisionTool);
|
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| {
|
self.state.adrs.modify(|adrs| {
|
||||||
id = format!("ADR-{:04}", adrs.len() + 1);
|
adr_id = format!("ADR-{:04}", adrs.len() + 1);
|
||||||
adrs.push(Adr {
|
let a = Adr {
|
||||||
id: id.clone(),
|
id: adr_id.clone(),
|
||||||
title: req.title,
|
title: req.title,
|
||||||
context: req.context,
|
context: req.context,
|
||||||
decision: req.decision,
|
decision: req.decision,
|
||||||
@@ -1038,9 +1048,18 @@ impl MemoryHandler {
|
|||||||
.duration_since(UNIX_EPOCH)
|
.duration_since(UNIX_EPOCH)
|
||||||
.unwrap()
|
.unwrap()
|
||||||
.as_secs(),
|
.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" => {
|
"query_decisions" => {
|
||||||
let req = parse_tool!(args.clone(), id, QueryDecisionsTool);
|
let req = parse_tool!(args.clone(), id, QueryDecisionsTool);
|
||||||
|
|||||||
@@ -180,6 +180,67 @@ impl MemoryIndex {
|
|||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn index_batch(
|
||||||
|
&self,
|
||||||
|
entities: Vec<Entity>,
|
||||||
|
tasks: Vec<Task>,
|
||||||
|
snippets: Vec<Snippet>,
|
||||||
|
adrs: Vec<Adr>,
|
||||||
|
) -> tokio::task::JoinHandle<tantivy::Result<()>> {
|
||||||
|
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)]
|
#[cfg(test)]
|
||||||
|
|||||||
+6
-31
@@ -58,39 +58,14 @@ impl MemoryState {
|
|||||||
|
|
||||||
pub async 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 mut handles = Vec::new();
|
|
||||||
|
|
||||||
{
|
let entities = self.graph.read().entities.values().cloned().collect();
|
||||||
let graph = self.graph.read();
|
let tasks = self.tasks.read().clone();
|
||||||
for e in graph.entities.values() {
|
let snippets = self.snippets.read().clone();
|
||||||
handles.push(new_idx.index_entity(e));
|
let adrs = self.adrs.read().clone();
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
{
|
let handle = new_idx.index_batch(entities, tasks, snippets, adrs);
|
||||||
let tasks = self.tasks.read();
|
let _ = handle.await;
|
||||||
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 _ = 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() {
|
||||||
|
|||||||
Reference in new issue
Block a user