refactor: address 5-pass audit findings for antipatterns, bottlenecks, memory efficiency, and LLM handlers
This commit is contained in:
1 parent
626403900f
commit
924b6d09fa
30 files changed
+1120
-503
No files matched your search
+103
-49
@@ -4,10 +4,98 @@ use std::sync::{Arc, RwLock};
|
||||
|
||||
pub const STORE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("store");
|
||||
|
||||
/// Internal write request dispatched to the single database writer actor.
|
||||
struct DbWriteTask {
|
||||
key: String,
|
||||
data: Vec<u8>,
|
||||
flushed: Arc<tokio::sync::Notify>,
|
||||
}
|
||||
|
||||
/// Shared centralized write queue actor that handles all database writes serially with micro-batching.
|
||||
#[derive(Clone)]
|
||||
pub struct DbWriteQueue {
|
||||
tx: tokio::sync::mpsc::Sender<DbWriteTask>,
|
||||
}
|
||||
|
||||
static QUEUE_REGISTRY: std::sync::Mutex<Option<(Arc<Database>, DbWriteQueue)>> = std::sync::Mutex::new(None);
|
||||
|
||||
fn get_or_create_queue(db: Arc<Database>) -> DbWriteQueue {
|
||||
let mut reg = QUEUE_REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
|
||||
if let Some((ref existing_db, ref queue)) = *reg {
|
||||
if Arc::ptr_eq(existing_db, &db) && !queue.tx.is_closed() {
|
||||
return queue.clone();
|
||||
}
|
||||
}
|
||||
let new_queue = DbWriteQueue::new(db.clone());
|
||||
*reg = Some((db, new_queue.clone()));
|
||||
new_queue
|
||||
}
|
||||
|
||||
impl DbWriteQueue {
|
||||
pub fn new(db: Arc<Database>) -> Self {
|
||||
let (tx, mut rx) = tokio::sync::mpsc::channel::<DbWriteTask>(2048);
|
||||
|
||||
tokio::spawn(async move {
|
||||
while let Some(first_task) = rx.recv().await {
|
||||
let mut batch = vec![first_task];
|
||||
|
||||
// Gold Standard Micro-batching: Drain up to 100 accumulated tasks from queue without blocking
|
||||
while batch.len() < 100 {
|
||||
match rx.try_recv() {
|
||||
Ok(task) => batch.push(task),
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
|
||||
let db_inner = db.clone();
|
||||
let _ = tokio::task::spawn_blocking(move || {
|
||||
match db_inner.begin_write() {
|
||||
Ok(write_txn) => {
|
||||
if let Ok(mut table) = write_txn.open_table(STORE_TABLE) {
|
||||
for task in &batch {
|
||||
if let Err(e) = table.insert(task.key.as_str(), task.data.as_slice()) {
|
||||
tracing::error!("Failed to insert key '{}' into redb: {}", task.key, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Err(e) = write_txn.commit() {
|
||||
tracing::error!("Failed to commit batch to redb: {}", e);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!("Failed to begin write transaction on redb writer actor: {}", e);
|
||||
}
|
||||
}
|
||||
|
||||
// Event-driven notification to all waiting listeners for this micro-batch
|
||||
for task in batch {
|
||||
task.flushed.notify_waiters();
|
||||
}
|
||||
})
|
||||
.await;
|
||||
}
|
||||
});
|
||||
|
||||
Self { tx }
|
||||
}
|
||||
|
||||
pub fn push(&self, key: String, data: Vec<u8>, flushed: Arc<tokio::sync::Notify>) {
|
||||
let task = DbWriteTask { key, data, flushed };
|
||||
if let Err(e) = self.tx.try_send(task) {
|
||||
let task = e.into_inner();
|
||||
let tx = self.tx.clone();
|
||||
tokio::spawn(async move {
|
||||
let _ = tx.send(task).await;
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct Store<T> {
|
||||
pub cache: Arc<RwLock<T>>,
|
||||
pub flushed: Arc<tokio::sync::Notify>,
|
||||
tx: tokio::sync::mpsc::Sender<()>,
|
||||
key: String,
|
||||
queue: DbWriteQueue,
|
||||
}
|
||||
|
||||
impl<T: DeserializeOwned + Default + Serialize + Clone + Send + Sync + 'static> Store<T> {
|
||||
@@ -15,53 +103,14 @@ impl<T: DeserializeOwned + Default + Serialize + Clone + Send + Sync + 'static>
|
||||
let initial_data = Self::load_from_db(key, &db);
|
||||
let cache = Arc::new(RwLock::new(initial_data));
|
||||
let flushed = Arc::new(tokio::sync::Notify::new());
|
||||
let (tx, mut rx) = tokio::sync::mpsc::channel::<()>(1);
|
||||
let queue = get_or_create_queue(db);
|
||||
|
||||
let db_clone = db.clone();
|
||||
let key_clone = key.to_string();
|
||||
let cache_clone = cache.clone();
|
||||
let flushed_clone = flushed.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
while rx.recv().await.is_some() {
|
||||
// Drain any pending notifications accumulated
|
||||
while rx.try_recv().is_ok() {}
|
||||
|
||||
let db_inner = db_clone.clone();
|
||||
let key_inner = key_clone.clone();
|
||||
let flushed_inner = flushed_clone.clone();
|
||||
let json_data = {
|
||||
let lock = cache_clone.read().unwrap_or_else(|e| e.into_inner());
|
||||
serde_json::to_vec(&*lock)
|
||||
.map_err(|e| tracing::error!("Failed to serialize memory store: {}", e))
|
||||
.ok()
|
||||
};
|
||||
|
||||
if let Some(json_data) = json_data {
|
||||
let _ = tokio::task::spawn_blocking(move || {
|
||||
// Retry up to 10 times if another Store holds write transaction
|
||||
for _ in 0..10 {
|
||||
match db_inner.begin_write() {
|
||||
Ok(write_txn) => {
|
||||
if let Ok(mut table) = write_txn.open_table(STORE_TABLE) {
|
||||
let _ = table.insert(key_inner.as_str(), json_data.as_slice());
|
||||
}
|
||||
let _ = write_txn.commit();
|
||||
break;
|
||||
}
|
||||
Err(_) => {
|
||||
std::thread::yield_now();
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
.await;
|
||||
flushed_inner.notify_waiters();
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
Self { cache, flushed, tx }
|
||||
Self {
|
||||
cache,
|
||||
flushed,
|
||||
key: key.to_string(),
|
||||
queue,
|
||||
}
|
||||
}
|
||||
|
||||
fn load_from_db(key: &str, db: &Database) -> T {
|
||||
@@ -86,11 +135,16 @@ impl<T: DeserializeOwned + Default + Serialize + Clone + Send + Sync + 'static>
|
||||
}
|
||||
|
||||
pub fn modify<F: FnOnce(&mut T)>(&self, f: F) {
|
||||
{
|
||||
let cloned_data = {
|
||||
let mut lock = self.cache.write().unwrap_or_else(|e| e.into_inner());
|
||||
f(&mut lock);
|
||||
lock.clone()
|
||||
};
|
||||
|
||||
match serde_json::to_vec(&cloned_data) {
|
||||
Ok(data) => self.queue.push(self.key.clone(), data, self.flushed.clone()),
|
||||
Err(e) => tracing::error!("Failed to serialize memory store for key '{}': {}", self.key, e),
|
||||
}
|
||||
let _ = self.tx.try_send(());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user