Files
mcp-memory/server/src/store.rs
T
Riz Ashraf 248799ca5d perf(server): offload sync disk/db IO to tokio blocking thread pool
Makes apply_sync_write, write_to_local_delta, and store modify truly async, preventing the mcp-memory-server from locking the tokio executor during disk IO
2026-09-12 08:54:12 +01:00

61 lines
1.8 KiB
Rust

use redb::{Database, ReadableDatabase, TableDefinition};
use serde::{Serialize, de::DeserializeOwned};
use std::sync::{Arc, RwLock};
pub const STORE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("store");
pub struct Store<T> {
pub key: String,
pub db: Arc<Database>,
pub cache: RwLock<T>,
}
impl<T: DeserializeOwned + Default + Serialize + Clone + Send + 'static> Store<T> {
pub fn new(key: &str, db: Arc<Database>) -> Self {
let initial_data = Self::load_from_db(key, &db);
Self {
key: key.to_string(),
db,
cache: RwLock::new(initial_data),
}
}
fn load_from_db(key: &str, db: &Database) -> T {
let read_txn = db.begin_read().unwrap();
if let Ok(table) = read_txn.open_table(STORE_TABLE) {
if let Ok(Some(value)) = table.get(key) {
if let Ok(parsed) = serde_json::from_slice::<T>(value.value()) {
return parsed;
}
}
}
T::default()
}
fn save_to_db(key: &str, db: &Database, data: &T) {
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(STORE_TABLE).unwrap();
let json_data = serde_json::to_vec(data).unwrap();
table.insert(key, json_data.as_slice()).unwrap();
}
write_txn.commit().unwrap();
}
pub fn read(&self) -> T {
let lock = self.cache.read().unwrap();
lock.clone()
}
pub fn modify<F: FnOnce(&mut T)>(&self, f: F) {
let mut lock = self.cache.write().unwrap();
f(&mut lock);
let key = self.key.clone();
let db = self.db.clone();
let data = lock.clone();
tokio::task::spawn_blocking(move || {
Self::save_to_db(&key, &db, &data);
});
}
}