use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Condvar, Mutex, Weak};
use std::time::Duration;
use super::collection_registry::dead_ranges;
use super::engine::META_CF;
use super::keyspace::SHARED_CF;
use super::RocksDb as DB;
fn delay() -> Duration {
let secs = std::env::var("SOLIDB_KS_GC_DELAY_SECS")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(10);
Duration::from_secs(secs)
}
static COMPACTIONS: AtomicU64 = AtomicU64::new(0);
pub fn compactions() -> u64 {
COMPACTIONS.load(Ordering::Relaxed)
}
struct Signal {
pending: Mutex<bool>,
cv: Condvar,
}
static SIGNAL: Signal = Signal {
pending: Mutex::new(false),
cv: Condvar::new(),
};
pub fn wake() {
if let Ok(mut p) = SIGNAL.pending.lock() {
*p = true;
SIGNAL.cv.notify_one();
}
}
pub fn run_once(db: &DB) -> usize {
let (Some(meta_cf), Some(shared)) = (db.cf_handle(META_CF), db.cf_handle(SHARED_CF)) else {
return 0;
};
let mut done = 0;
for (marker, lo, hi) in dead_ranges(db) {
db.compact_range_cf(&shared, Some(&lo), Some(&hi));
if let Err(e) = db.delete_cf(&meta_cf, &marker) {
tracing::warn!("keyspace GC: failed to clear marker: {}", e);
continue;
}
COMPACTIONS.fetch_add(1, Ordering::Relaxed);
done += 1;
}
done
}
static STARTED: std::sync::Once = std::sync::Once::new();
pub fn ensure_started(db: &Arc<DB>) {
let weak: Weak<DB> = Arc::downgrade(db);
STARTED.call_once(|| {
let _ = std::thread::Builder::new()
.name("solidb-keyspace-gc".into())
.spawn(move || loop {
{
let Ok(mut pending) = SIGNAL.pending.lock() else {
return;
};
while !*pending {
let (guard, _) = SIGNAL
.cv
.wait_timeout(pending, Duration::from_secs(60))
.unwrap_or_else(|e| e.into_inner());
pending = guard;
if weak.strong_count() == 0 {
return;
}
if !*pending {
break; }
}
*pending = false;
}
std::thread::sleep(delay());
let Some(db) = weak.upgrade() else {
return;
};
run_once(&db);
});
});
wake();
}