From 686fea683de65bae94894c6958ec21a591a6ffe1 Mon Sep 17 00:00:00 2001 From: Riz Ashraf Date: Sat, 12 Sep 2026 20:41:30 +0100 Subject: [PATCH] perf: completely remove disk I/O from memory access paths, optimize Tantivy indexing, and fix MemoryIndex commits --- server/src/main.rs | 75 ++++++++++++++++++++++++---------------- server/src/search.rs | 17 +++++---- server/src/state.rs | 82 +++++++++++++++++++++++++++++--------------- 3 files changed, 110 insertions(+), 64 deletions(-) diff --git a/server/src/main.rs b/server/src/main.rs index 78377d2..e34ef36 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -123,37 +123,46 @@ async fn git_sync_worker(state: Arc) { 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) { async fn reconcile_worker(state: Arc) { 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) { state_clone.rebuild_index(); }).await; } - } } @@ -754,6 +770,7 @@ fn main() -> Result<(), Box> { context_workspaces: Store::new("context_workspaces", db.clone()), }); + state.recover_wal(); state.rebuild_index(); run_server(state) diff --git a/server/src/search.rs b/server/src/search.rs index 656714a..b407efe 100644 --- a/server/src/search.rs +++ b/server/src/search.rs @@ -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(()) } } diff --git a/server/src/state.rs b/server/src/state.rs index efd1c37..b03e551 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -34,35 +34,53 @@ pub struct MemoryState { } impl MemoryState { + + pub fn unique_items(input: Vec) -> Vec { + 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(input: Vec) -> Vec { - 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::(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::(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; }