perf: fix store mem leak, batch DB flushes, and remove double-clones

This commit is contained in:
Riz Ashraf committed 2026-09-21 09:32:16 +01:00
1 parent 2c5953e4f5
commit d5a05f8416
2 files changed
+24 -15

No files matched your search

+4 -4
View File
@@ -63,10 +63,10 @@ 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 entities = self.graph.read().entities.values().cloned().collect(); let entities = self.graph.read().entities.into_values().collect();
let tasks = self.tasks.read().clone(); let tasks = self.tasks.read();
let snippets = self.snippets.read().clone(); let snippets = self.snippets.read();
let adrs = self.adrs.read().clone(); let adrs = self.adrs.read();
let handle = new_idx.index_batch(entities, tasks, snippets, adrs); let handle = new_idx.index_batch(entities, tasks, snippets, adrs);
let _ = handle.await; let _ = handle.await;
+20 -11
View File
@@ -5,21 +5,32 @@ use std::sync::{Arc, RwLock};
pub const STORE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("store"); pub const STORE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("store");
pub struct Store<T> { pub struct Store<T> {
pub cache: RwLock<T>, pub cache: Arc<RwLock<T>>,
tx: tokio::sync::mpsc::UnboundedSender<Vec<u8>>, tx: tokio::sync::mpsc::Sender<()>,
} }
impl<T: DeserializeOwned + Default + Serialize + Clone + Send + 'static> Store<T> { impl<T: DeserializeOwned + Default + Serialize + Clone + Send + Sync + 'static> Store<T> {
pub fn new(key: &str, db: Arc<Database>) -> Self { pub fn new(key: &str, db: Arc<Database>) -> Self {
let initial_data = Self::load_from_db(key, &db); let initial_data = Self::load_from_db(key, &db);
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>(); let cache = Arc::new(RwLock::new(initial_data));
let (tx, mut rx) = tokio::sync::mpsc::channel::<()>(1);
let db_clone = db.clone(); let db_clone = db.clone();
let key_clone = key.to_string(); let key_clone = key.to_string();
let cache_clone = cache.clone();
tokio::spawn(async move { tokio::spawn(async move {
while let Some(json_data) = rx.recv().await { while rx.recv().await.is_some() {
// Drain any other pending notifications so we batch writes
while let Ok(_) = rx.try_recv() {}
let db_inner = db_clone.clone(); let db_inner = db_clone.clone();
let key_inner = key_clone.clone(); let key_inner = key_clone.clone();
let json_data = {
let lock = cache_clone.read().unwrap();
serde_json::to_vec(&*lock).unwrap()
};
let _ = tokio::task::spawn_blocking(move || { let _ = tokio::task::spawn_blocking(move || {
let write_txn = db_inner.begin_write().unwrap(); let write_txn = db_inner.begin_write().unwrap();
{ {
@@ -32,7 +43,7 @@ impl<T: DeserializeOwned + Default + Serialize + Clone + Send + 'static> Store<T
}); });
Self { Self {
cache: RwLock::new(initial_data), cache,
tx, tx,
} }
} }
@@ -61,13 +72,11 @@ impl<T: DeserializeOwned + Default + Serialize + Clone + Send + 'static> Store<T
} }
pub fn modify<F: FnOnce(&mut T)>(&self, f: F) { pub fn modify<F: FnOnce(&mut T)>(&self, f: F) {
let json_data = { {
let mut lock = self.cache.write().unwrap(); let mut lock = self.cache.write().unwrap();
f(&mut lock); f(&mut lock);
// Serialize while holding lock to avoid expensive deep clone of T }
serde_json::to_vec(&*lock).unwrap() let _ = self.tx.try_send(());
};
let _ = self.tx.send(json_data);
} }
} }