use std::io;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, RwLock, RwLockWriteGuard};
use std::thread::JoinHandle;
use std::time::Duration;
#[cfg(feature = "persist")]
use std::time::Instant;
use crate::config::{Config, TtlReaperMode};
#[cfg(feature = "persist")]
use crate::metric::{KevyMetric, MetricSink};
use crate::store::{Inner, Shards};
#[allow(clippy::type_complexity)] pub(crate) fn spawn_reaper(
config: &Config,
shards: &Shards,
) -> io::Result<(Option<Arc<AtomicBool>>, Option<JoinHandle<()>>)> {
match config.ttl_reaper {
TtlReaperMode::Manual => Ok((None, None)),
TtlReaperMode::Background => {
let stop = Arc::new(AtomicBool::new(false));
let stop_t = stop.clone();
let shards_t = shards.clone();
let interval = config.reaper_interval;
let samples = config.reaper_samples;
let rounds = config.reaper_max_rounds;
#[cfg(feature = "persist")]
let policy = kevy_persist::RewritePolicy {
pct: config.auto_aof_rewrite_pct,
min_size: config.auto_aof_rewrite_min_size,
bytes: config.auto_aof_rewrite_bytes,
interval_secs: config.auto_aof_rewrite_interval_secs,
};
#[cfg(feature = "persist")]
let sink = config.metric_sink.clone();
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
let tier_spec = config.tier_budget;
let handle = std::thread::Builder::new()
.name(String::from("kevy-embedded-reaper"))
.spawn(move || {
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
let tier = tier_spec;
#[cfg(not(all(feature = "tier", not(target_arch = "wasm32"))))]
let tier = ();
#[cfg(feature = "persist")]
reaper_loop(shards_t, stop_t, interval, samples, rounds, tier, policy, sink);
#[cfg(not(feature = "persist"))]
reaper_loop(shards_t, stop_t, interval, samples, rounds, tier);
})?;
Ok((Some(stop), Some(handle)))
}
}
}
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
type TierSpecOpt = Option<crate::config::TierBudgetSpec>;
#[cfg(not(all(feature = "tier", not(target_arch = "wasm32"))))]
type TierSpecOpt = ();
#[allow(clippy::too_many_arguments)] fn reaper_loop(
shards: Shards,
stop: Arc<AtomicBool>,
interval: Duration,
samples: usize,
rounds: u32,
tier: TierSpecOpt,
#[cfg(feature = "persist")] policy: kevy_persist::RewritePolicy,
#[cfg(feature = "persist")] sink: Option<MetricSink>,
) {
#[cfg(not(all(feature = "tier", not(target_arch = "wasm32"))))]
let _ = tier;
while !stop.load(Ordering::Relaxed) {
std::thread::sleep(interval);
if stop.load(Ordering::Relaxed) {
break;
}
for shard in shards.iter() {
{
let mut g = lock_inner(shard);
let _ = g.store.tick_expire(samples, rounds);
let _ = g.store.tick_hash_ttl(64);
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
crate::shard::tier_tick_upkeep(&mut g, tier, shards.len());
let _ = g.store.demote_step();
let _ = g.store.tier_compact_tick();
#[cfg(feature = "persist")]
if let Some(aof) = &mut g.aof {
let _ = aof.maybe_sync();
}
}
#[cfg(feature = "persist")]
concurrent_auto_rewrite(shard, policy, sink.as_ref());
}
}
}
#[cfg(feature = "persist")]
pub(crate) fn concurrent_auto_rewrite(
inner: &Arc<RwLock<Inner>>,
policy: kevy_persist::RewritePolicy,
sink: Option<&MetricSink>,
) {
let Some((start, view, tmp, before_bytes)) = begin_rewrite(inner, policy) else {
return;
};
let keys = match kevy_persist::dump_aof(&tmp, &view) {
Ok((keys, _)) => keys,
Err(e) => {
eprintln!("kevy: embedded auto AOF rewrite (dump) failed: {e}");
let mut g = lock_inner(inner);
if let Some(aof) = &mut g.aof {
aof.abort_concurrent_rewrite();
}
let _ = std::fs::remove_file(&tmp);
return;
}
};
let mut g = lock_inner(inner);
let Some(aof) = &mut g.aof else { return };
match aof.finish_concurrent_rewrite(&tmp, keys) {
Ok(stats) => {
if let Some(sink) = sink {
sink.emit(KevyMetric::Rewrite {
keys: stats.keys,
before_bytes,
after_bytes: stats.bytes,
elapsed_ms: start.elapsed().as_millis() as u64,
});
}
}
Err(e) => {
eprintln!("kevy: embedded auto AOF rewrite (finish) failed: {e}");
aof.abort_concurrent_rewrite();
let _ = std::fs::remove_file(&tmp);
}
}
}
#[cfg(feature = "persist")]
#[allow(clippy::type_complexity)] fn begin_rewrite(
inner: &Arc<RwLock<Inner>>,
policy: kevy_persist::RewritePolicy,
) -> Option<(Instant, kevy_store::SnapshotView, std::path::PathBuf, u64)> {
let mut g = lock_inner(inner);
let ready = g.aof.as_ref().is_some_and(|a| a.rewrite_due(policy));
if !ready {
return None;
}
let start = Instant::now();
let Inner { store, aof, .. } = &mut *g;
let aof = aof.as_mut().expect("checked above");
let before = aof.size_bytes();
let view = store.collect_snapshot();
match aof.begin_view_rewrite() {
Ok(tmp) => Some((start, view, tmp, before)),
Err(e) => {
eprintln!("kevy: embedded auto AOF rewrite (begin) failed: {e}");
None
}
}
}
pub(crate) fn lock_inner(inner: &Arc<RwLock<Inner>>) -> RwLockWriteGuard<'_, Inner> {
inner.write().unwrap_or_else(std::sync::PoisonError::into_inner)
}