perf: completely remove disk I/O from memory access paths, optimize Tantivy indexing, and fix MemoryIndex commits

This commit is contained in:
Riz Ashraf committed 2026-09-12 20:41:30 +01:00
1 parent a5011048b1
commit 686fea683d
3 files changed
+110 -64

No files matched your search

+46 -29
View File
@@ -123,37 +123,46 @@ async fn git_sync_worker(state: Arc<MemoryState>) {
loop {
tokio::time::sleep(tokio::time::Duration::from_secs(30)).await;
if let Ok(repo) = git2::Repository::discover(&repo_path) {
if let Ok(head) = repo.head() {
if let Ok(commit) = head.peel_to_commit() {
let current_id = commit.id().to_string();
if current_id != last_commit_id && !last_commit_id.is_empty() {
let repo_path_clone = repo_path.clone();
let commit_data = tokio::task::spawn_blocking(move || {
if let Ok(repo) = git2::Repository::discover(&repo_path_clone) {
if let Ok(head) = repo.head() {
if let Ok(commit) = head.peel_to_commit() {
let current_id = commit.id().to_string();
let msg = commit.message().unwrap_or("").to_string();
let branch = head.shorthand().unwrap_or("unknown").to_string();
state.ledger.modify(|changes| {
changes.push(crate::models::CodeChange {
git_commit: Some(current_id.clone()),
git_branch: Some(branch),
description: format!("Auto-synced commit: {}", msg.trim()),
timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(),
file_path: "".to_string(),
});
});
eprintln!("Git Sync: Logged new commit {}", current_id);
state.tasks.modify(|tasks| {
for task in tasks.iter_mut() {
if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) {
task.status = "completed".to_string();
eprintln!("Git Sync: Auto-completed task '{}'", task.title);
}
}
});
return Some((current_id, msg, branch));
}
last_commit_id = current_id;
}
}
None
})
.await
.unwrap_or(None);
if let Some((current_id, msg, branch)) = commit_data {
if current_id != last_commit_id && !last_commit_id.is_empty() {
state.ledger.modify(|changes| {
changes.push(crate::models::CodeChange {
git_commit: Some(current_id.clone()),
git_branch: Some(branch),
description: format!("Auto-synced commit: {}", msg.trim()),
timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(),
file_path: "".to_string(),
});
});
eprintln!("Git Sync: Logged new commit {}", current_id);
state.tasks.modify(|tasks| {
for task in tasks.iter_mut() {
if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) {
task.status = "completed".to_string();
eprintln!("Git Sync: Auto-completed task '{}'", task.title);
}
}
});
}
last_commit_id = current_id;
}
}
}
@@ -161,12 +170,20 @@ async fn git_sync_worker(state: Arc<MemoryState>) {
async fn reconcile_worker(state: Arc<MemoryState>) {
loop {
sleep(Duration::from_secs(5)).await;
let pattern = format!("{}/delta_*.json", state.base_dir.display());
let has_local = {
let session = state.session_graph.read().unwrap();
!session.entities.is_empty() || !session.relations.is_empty()
};
let has_files = glob::glob(&pattern).map(|p| p.count() > 0).unwrap_or(false);
let base_dir = state.base_dir.clone();
let has_files = tokio::task::spawn_blocking(move || {
let pattern = format!("{}/delta_*.json", base_dir.display());
glob::glob(&pattern).map(|p| p.count() > 0).unwrap_or(false)
})
.await
.unwrap_or(false);
if has_local || has_files {
state.apply_sync_write(|_master| {}).await;
let state_clone = state.clone();
@@ -174,7 +191,6 @@ async fn reconcile_worker(state: Arc<MemoryState>) {
state_clone.rebuild_index();
}).await;
}
}
}
@@ -754,6 +770,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
context_workspaces: Store::new("context_workspaces", db.clone()),
});
state.recover_wal();
state.rebuild_index();
run_server(state)
+10 -7
View File
@@ -49,7 +49,7 @@ impl MemoryIndex {
}
pub fn index_entity(&self, e: &Entity) -> tantivy::Result<()> {
let mut writer = self.writer.lock().unwrap();
let writer = self.writer.lock().unwrap();
writer.add_document(doc!(
self.id_field => e.name.clone(),
self.title_field => e.name.clone(),
@@ -57,12 +57,11 @@ impl MemoryIndex {
self.type_field => "entity",
self.namespace_field => e.namespace.clone()
))?;
writer.commit()?;
Ok(())
}
pub fn index_task(&self, t: &Task) -> tantivy::Result<()> {
let mut writer = self.writer.lock().unwrap();
let writer = self.writer.lock().unwrap();
writer.add_document(doc!(
self.id_field => t.id.clone(),
self.title_field => t.title.clone(),
@@ -70,6 +69,12 @@ impl MemoryIndex {
self.type_field => "task",
self.namespace_field => "global"
))?;
Ok(())
}
pub fn commit(&self) -> tantivy::Result<()> {
let mut writer = self.writer.lock().unwrap();
writer.commit()?;
Ok(())
}
@@ -117,7 +122,7 @@ impl MemoryIndex {
}
pub fn index_snippet(&self, s: &Snippet) -> tantivy::Result<()> {
let mut writer = self.writer.lock().unwrap();
let writer = self.writer.lock().unwrap();
writer.add_document(doc!(
self.id_field => s.name.clone(),
self.title_field => s.name.clone(),
@@ -125,12 +130,11 @@ impl MemoryIndex {
self.type_field => "snippet",
self.namespace_field => "global"
))?;
writer.commit()?;
Ok(())
}
pub fn index_adr(&self, a: &Adr) -> tantivy::Result<()> {
let mut writer = self.writer.lock().unwrap();
let writer = self.writer.lock().unwrap();
writer.add_document(doc!(
self.id_field => a.id.clone(),
self.title_field => a.title.clone(),
@@ -138,7 +142,6 @@ impl MemoryIndex {
self.type_field => "adr",
self.namespace_field => "global"
))?;
writer.commit()?;
Ok(())
}
}
+54 -28
View File
@@ -34,35 +34,53 @@ pub struct MemoryState {
}
impl MemoryState {
pub fn unique_items<T: Eq + std::hash::Hash + Clone>(input: Vec<T>) -> Vec<T> {
let mut keys = std::collections::HashSet::new();
input.into_iter().filter(|entry| keys.insert(entry.clone())).collect()
}
fn master_mtime(&self) -> SystemTime {
fs::metadata(&self.master_path)
.and_then(|m| m.modified())
.unwrap_or(SystemTime::UNIX_EPOCH)
}
pub fn unique_items<T: Eq + std::hash::Hash + Clone>(input: Vec<T>) -> Vec<T> {
let mut keys = HashSet::new();
let mut list = Vec::new();
for entry in input {
if keys.insert(entry.clone()) {
list.push(entry);
pub fn recover_wal(&self) {
let wal_path = self.base_dir.join("wal.jsonl");
if let Ok(content) = std::fs::read_to_string(&wal_path) {
let mut session = self.session_graph.write().unwrap();
for line in content.lines() {
if let Ok(d) = serde_json::from_str::<KnowledgeGraph>(line) {
Self::merge_graphs(&mut session, &d);
}
}
}
list
}
pub fn merge_graphs(dest: &mut KnowledgeGraph, src: &KnowledgeGraph) {
for (name, src_ent) in &src.entities {
let dest_ent = dest
.entities
.entry(name.clone())
.or_insert_with(|| src_ent.clone());
if dest_ent.name == src_ent.name {
dest_ent.observations.extend(src_ent.observations.clone());
dest_ent.observations = Self::unique_items(dest_ent.observations.clone());
.or_insert_with(|| crate::models::Entity {
name: src_ent.name.clone(),
entity_type: src_ent.entity_type.clone(),
observations: Vec::new(),
namespace: src_ent.namespace.clone(),
git_branch: src_ent.git_branch.clone(),
});
for obs in &src_ent.observations {
if !dest_ent.observations.contains(obs) {
dest_ent.observations.push(obs.clone());
}
}
}
for rel in &src.relations {
if !dest.relations.contains(rel) {
dest.relations.push(rel.clone());
}
}
dest.relations.extend(src.relations.clone());
dest.relations = Self::unique_items(dest.relations.clone());
}
pub fn read_master_cached(&self) -> KnowledgeGraph {
@@ -98,14 +116,6 @@ impl MemoryState {
pub fn get_full_graph(&self) -> KnowledgeGraph {
let mut master = self.read_master_cached();
let wal_path = self.base_dir.join("wal.jsonl");
if let Ok(content) = std::fs::read_to_string(&wal_path) {
for line in content.lines() {
if let Ok(d) = serde_json::from_str::<KnowledgeGraph>(line) {
Self::merge_graphs(&mut master, &d);
}
}
}
let session_graph = self.session_graph.read().unwrap();
Self::merge_graphs(&mut master, &session_graph);
master
@@ -183,14 +193,29 @@ impl MemoryState {
pub fn rebuild_index(&self) {
if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) {
let session = self.session_graph.read().unwrap();
let mut full = {
let cache = self.master_cache.read().unwrap();
cache.0.clone()
};
Self::merge_graphs(&mut full, &session);
for e in full.entities.values() {
let _ = new_idx.index_entity(e);
let cache = self.master_cache.read().unwrap();
// Index entities that are only in master, or merge if they are in both
for (name, e) in &cache.0.entities {
if let Some(session_e) = session.entities.get(name) {
let mut merged_e = e.clone();
for obs in &session_e.observations {
if !merged_e.observations.contains(obs) {
merged_e.observations.push(obs.clone());
}
}
let _ = new_idx.index_entity(&merged_e);
} else {
let _ = new_idx.index_entity(e);
}
}
// Index entities that are only in session
for (name, session_e) in &session.entities {
if !cache.0.entities.contains_key(name) {
let _ = new_idx.index_entity(session_e);
}
}
for t in self.tasks.read() {
let _ = new_idx.index_task(&t);
}
@@ -200,6 +225,7 @@ impl MemoryState {
for a in self.adrs.read() {
let _ = new_idx.index_adr(&a);
}
let _ = new_idx.commit();
if let Ok(mut w) = self.search_index.write() {
*w = new_idx;
}