diff --git a/server/src/state.rs b/server/src/state.rs index 126b0d4..43de6fa 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -63,10 +63,10 @@ impl MemoryState { pub async fn rebuild_index(&self) { if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) { - let entities = self.graph.read().entities.values().cloned().collect(); - let tasks = self.tasks.read().clone(); - let snippets = self.snippets.read().clone(); - let adrs = self.adrs.read().clone(); + let entities = self.graph.read().entities.into_values().collect(); + let tasks = self.tasks.read(); + let snippets = self.snippets.read(); + let adrs = self.adrs.read(); let handle = new_idx.index_batch(entities, tasks, snippets, adrs); let _ = handle.await; diff --git a/server/src/store.rs b/server/src/store.rs index 700fc58..173861f 100644 --- a/server/src/store.rs +++ b/server/src/store.rs @@ -5,21 +5,32 @@ use std::sync::{Arc, RwLock}; pub const STORE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("store"); pub struct Store { - pub cache: RwLock, - tx: tokio::sync::mpsc::UnboundedSender>, + pub cache: Arc>, + tx: tokio::sync::mpsc::Sender<()>, } -impl Store { +impl Store { pub fn new(key: &str, db: Arc) -> Self { let initial_data = Self::load_from_db(key, &db); - let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::>(); + let cache = Arc::new(RwLock::new(initial_data)); + let (tx, mut rx) = tokio::sync::mpsc::channel::<()>(1); let db_clone = db.clone(); let key_clone = key.to_string(); + let cache_clone = cache.clone(); + 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 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 write_txn = db_inner.begin_write().unwrap(); { @@ -32,7 +43,7 @@ impl Store Store(&self, f: F) { - let json_data = { + { let mut lock = self.cache.write().unwrap(); f(&mut lock); - // Serialize while holding lock to avoid expensive deep clone of T - serde_json::to_vec(&*lock).unwrap() - }; - let _ = self.tx.send(json_data); + } + let _ = self.tx.try_send(()); } }