use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::Duration;
use parking_lot::{Condvar, Mutex};
use surrealkv::Tree;
use super::TARGET;
use crate::kvs::err::{Error, Result};
pub(super) struct BackgroundFlusher {
notify: Arc<ShutdownSignal>,
handle: Mutex<Option<thread::JoinHandle<()>>>,
}
struct ShutdownSignal {
flag: AtomicBool,
condvar: Condvar,
mutex: Mutex<()>,
}
impl BackgroundFlusher {
pub fn new(db: Tree, interval: Duration) -> Result<Self> {
let notify = Arc::new(ShutdownSignal {
flag: AtomicBool::new(false),
condvar: Condvar::new(),
mutex: Mutex::new(()),
});
let signal = Arc::clone(¬ify);
let handle = thread::Builder::new()
.name("surrealkv-background-flusher".to_string())
.spawn(move || {
loop {
let mut guard = signal.mutex.lock();
signal.condvar.wait_for(&mut guard, interval);
drop(guard);
if signal.flag.load(Ordering::Relaxed) {
break;
}
if let Err(err) = db.flush_wal(true) {
error!(target: TARGET, "Failed to flush WAL: {err}");
}
}
})
.map_err(|_| {
Error::Datastore("failed to spawn SurrealKV background flush thread".to_string())
})?;
Ok(Self {
notify,
handle: Mutex::new(Some(handle)),
})
}
pub fn shutdown(&self) -> Result<()> {
self.notify.flag.store(true, Ordering::Relaxed);
self.notify.condvar.notify_one();
if let Some(handle) = self.handle.lock().take() {
let _ = handle.join();
}
Ok(())
}
}