diff --git a/src/trx-server/src/history_store.rs b/src/trx-server/src/history_store.rs index e5c5abdf..afb1384a 100644 --- a/src/trx-server/src/history_store.rs +++ b/src/trx-server/src/history_store.rs @@ -185,20 +185,38 @@ pub fn flush_all(db: &mut PickleDb, rig_id: &str, histories: &Arc, rig_histories: &[(String, Arc)]) { + let Ok(mut guard) = db.lock() else { + tracing::warn!("history database mutex poisoned; skipping periodic flush"); + return; + }; + for (rig_id, histories) in rig_histories { + flush_all(&mut guard, rig_id, histories); + } +} + /// Spawn a Tokio task that flushes all rigs' histories to disk every 60 seconds. +/// +/// Snapshot cloning, JSON serialization, and disk I/O run on Tokio's blocking +/// pool so a large history database cannot stall an async runtime worker. pub fn spawn_flush_task( db: Arc>, rig_histories: Vec<(String, Arc)>, ) { tokio::spawn(async move { + let rig_histories = Arc::new(rig_histories); let mut interval = tokio::time::interval(Duration::from_secs(60)); interval.tick().await; // consume the immediate first tick loop { interval.tick().await; - if let Ok(mut guard) = db.lock() { - for (rig_id, histories) in &rig_histories { - flush_all(&mut guard, rig_id, histories); - } + let db = Arc::clone(&db); + let rig_histories = Arc::clone(&rig_histories); + if let Err(err) = tokio::task::spawn_blocking(move || { + flush_all_rigs(&db, rig_histories.as_slice()); + }) + .await + { + tracing::warn!(error = %err, "history flush worker failed"); } } });