use std::{
ops::Range,
sync::{
Arc, Weak,
atomic::{
AtomicBool, AtomicU64,
Ordering::{Acquire, Relaxed, Release},
},
},
time::Duration,
};
use compio::{
runtime::{JoinHandle, spawn},
time::sleep,
};
use log::{info, warn};
use parking_lot::{Mutex, RwLock};
use wbase::time::now_ms;
use wcompact::CompactionType;
use wdev::Device;
use whasher::{GxBuildHasher, HashSet};
use wval::{KeyTag, NamespaceDbCodec};
use crate::{
config::GcConfig,
error::Result,
session::{SessionSlot, StoreSession},
store::WedbStore,
ttl::{TTL_VALUE_LEN, TtlProbe},
};
const MIN_SCAN_INTERVAL_MS: u64 = 10;
type ExpiredKeySet = HashSet<(u64, u64, Box<[u8]>)>;
#[derive(Default)]
struct GcStats {
expired_deleted: AtomicU64,
expired_fields_deleted: AtomicU64,
compactions: AtomicU64,
last_scan_deleted: AtomicU64,
last_scan_fields_deleted: AtomicU64,
last_scan_scanned: AtomicU64,
total_scanned: AtomicU64,
last_compact_dropped: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct GcStatsSnapshot {
pub expired_deleted: u64,
pub expired_fields_deleted: u64,
pub compactions: u64,
pub last_scan_deleted: u64,
pub last_scan_fields_deleted: u64,
pub last_scan_scanned: u64,
pub total_scanned: u64,
pub last_compact_dropped: u64,
}
struct SweepPicks {
picked: ExpiredKeySet,
fields_picked: ExpiredKeySet,
max_picks: usize,
}
impl SweepPicks {
fn len(&self) -> usize {
self.picked.len() + self.fields_picked.len()
}
}
pub struct GcManager<D: Device> {
store: Weak<WedbStore<D>>,
cfg: Arc<RwLock<GcConfig>>,
cancel: AtomicBool,
pub(crate) inflight: AtomicBool,
last_compact_ms: AtomicU64,
cold_cursor: AtomicU64,
sweep_session: SessionSlot<D>,
stats: GcStats,
}
impl<D: Device> GcManager<D> {
pub fn new(store: &Arc<WedbStore<D>>) -> Self {
Self {
store: Arc::downgrade(store),
cfg: Arc::clone(&store.gc_cfg),
cancel: AtomicBool::new(false),
inflight: AtomicBool::new(false),
last_compact_ms: AtomicU64::new(0),
cold_cursor: AtomicU64::new(0),
sweep_session: SessionSlot::new(),
stats: GcStats::default(),
}
}
fn scan_interval_ms(&self) -> u64 {
self.cfg.read().scan_interval_ms.max(MIN_SCAN_INTERVAL_MS)
}
pub fn spawn(store: Arc<WedbStore<D>>) -> GcHandle<D>
where
D: Device + 'static,
{
let mgr = Arc::new(Self::new(&store));
let weak = Arc::downgrade(&mgr);
let join = spawn(async move {
loop {
let interval_ms = match weak.upgrade() {
None => return,
Some(m) => {
if m.cancel.load(Relaxed) {
return;
}
let ms = m.scan_interval_ms();
drop(m);
ms
}
};
sleep(Duration::from_millis(interval_ms)).await;
let Some(m) = weak.upgrade() else {
return;
};
if let Err(e) = m.run_once().await {
warn!("内置 GC 轮次失败,留待下轮: err={e}");
}
}
});
GcHandle {
inner: mgr,
join: Mutex::new(Some(join)),
}
}
pub async fn drive(&self) {
loop {
sleep(Duration::from_millis(self.scan_interval_ms())).await;
if self.cancel.load(Relaxed) {
return;
}
if self.store.upgrade().is_none() {
return;
}
if let Err(e) = self.run_once().await {
warn!("内置 GC 轮次失败,留待下轮: err={e}");
}
}
}
pub fn stats(&self) -> GcStatsSnapshot {
GcStatsSnapshot {
expired_deleted: self.stats.expired_deleted.load(Relaxed),
expired_fields_deleted: self.stats.expired_fields_deleted.load(Relaxed),
compactions: self.stats.compactions.load(Relaxed),
last_scan_deleted: self.stats.last_scan_deleted.load(Relaxed),
last_scan_fields_deleted: self.stats.last_scan_fields_deleted.load(Relaxed),
last_scan_scanned: self.stats.last_scan_scanned.load(Relaxed),
total_scanned: self.stats.total_scanned.load(Relaxed),
last_compact_dropped: self.stats.last_compact_dropped.load(Relaxed),
}
}
pub async fn run_once(&self) -> Result<()> {
if self.inflight.swap(true, Acquire) {
return Ok(());
}
let _guard = RunGuard(&self.inflight);
self.tick().await
}
async fn tick(&self) -> Result<()> {
let Some(store) = self.store.upgrade() else {
return Ok(());
};
let cfg = self.cfg.read().clone();
let (deleted, fields_deleted, scanned) = self.sweep_expired(&store, &cfg).await?;
self.stats.expired_deleted.fetch_add(deleted, Relaxed);
self
.stats
.expired_fields_deleted
.fetch_add(fields_deleted, Relaxed);
self.stats.last_scan_deleted.store(deleted, Relaxed);
self
.stats
.last_scan_fields_deleted
.store(fields_deleted, Relaxed);
self.stats.last_scan_scanned.store(scanned, Relaxed);
self.stats.total_scanned.fetch_add(scanned, Relaxed);
self.try_compact(&store, &cfg).await
}
async fn sweep_expired(
&self,
store: &Arc<WedbStore<D>>,
cfg: &GcConfig,
) -> Result<(u64, u64, u64)> {
let now = now_ms();
let cap = cfg.max_batch_deletes.max(1);
let cold_cap = cfg.max_scan_records.max(1);
let session = self.sweep_session.take(store)?;
let mut picks = SweepPicks {
picked: HashSet::with_hasher(GxBuildHasher::default()),
fields_picked: HashSet::with_hasher(GxBuildHasher::default()),
max_picks: cap,
};
let mut scanned = 0u64;
let read_only = store.read_only_address();
let tail = store.tail_address();
if read_only < tail {
let (n, ..) =
Self::collect_expired(&session, store, read_only..tail, now, u64::MAX, &mut picks).await?;
scanned += n;
}
let mut cold_commit = None;
if picks.len() < cap {
let cold_from = self.cold_cursor.load(Relaxed).max(store.begin_address());
if cold_from < read_only {
let (n, next_addr, exhausted) = Self::collect_expired(
&session,
store,
cold_from..read_only,
now,
cold_cap as u64,
&mut picks,
)
.await?;
scanned += n;
cold_commit = Some(if exhausted || next_addr >= read_only {
read_only
} else {
next_addr
});
}
}
let mut deleted = 0u64;
for (ns, db, key) in &picks.picked {
session.set_context(*ns, *db);
match session.check_expired(key).await {
Ok(true) => deleted += 1,
Ok(false) => {}
Err(e) => warn!("内置 GC 过期删除失败,留待下一扫描周期重试: err={e}"),
}
}
let mut fields_deleted = 0u64;
for (ns, db, key) in &picks.fields_picked {
session.set_context(*ns, *db);
match session.collect_expired_hash_fields(key, now).await {
Ok(n) => fields_deleted += n,
Err(e) => warn!("内置 GC 字段过期收集失败,留待下一扫描周期重试: err={e}"),
}
}
if deleted + fields_deleted > 0 {
info!(
"内置 GC 过期扫描完成: 键候选={}, 物理删除键={deleted}, 字段候选={}, 物理清除字段={fields_deleted}",
picks.picked.len(),
picks.fields_picked.len()
);
}
if let Some(addr) = cold_commit {
self.cold_cursor.store(addr, Relaxed);
}
self.sweep_session.restore(session);
Ok((deleted, fields_deleted, scanned))
}
async fn collect_expired(
session: &StoreSession<D>,
store: &Arc<WedbStore<D>>,
range: Range<u64>,
now: u64,
max_records: u64,
picks: &mut SweepPicks,
) -> Result<(u64, u64, bool)> {
let mut scanned = 0u64;
let mut exhausted = false;
let mut scan = store.hlog.scan_iter(range.start, range.end);
loop {
if picks.len() >= picks.max_picks || scanned >= max_records {
break;
}
let next = scan
.next_ref(|item| {
let rec = item.rec;
if rec.is_tombstone() {
return Ok(true);
}
let Ok((ns, db, tag, user_key)) = NamespaceDbCodec::decode_tagged_key(rec.key) else {
return Ok(true);
};
match tag {
KeyTag::Ttl => {
let Ok(be) = <[u8; TTL_VALUE_LEN]>::try_from(rec.value) else {
return Ok(true);
};
if u64::from_be_bytes(be) > now {
return Ok(true);
}
session.set_context(ns, db);
if matches!(session.probe_ttl(user_key, now), TtlProbe::Pass) {
return Ok(true);
}
picks.picked.insert((ns, db, Box::from(user_key)));
}
KeyTag::Meta => {
if !StoreSession::<D>::meta_value_has_expire(rec.value) {
return Ok(true);
}
session.set_context(ns, db);
if !matches!(session.probe_meta_has_expire(user_key), Ok(true)) {
return Ok(true);
}
picks.fields_picked.insert((ns, db, Box::from(user_key)));
}
_ => {}
}
Ok(true)
})
.await?;
match next {
None => {
exhausted = true;
break;
}
Some(_) => scanned += 1,
}
}
Ok((scanned, scan.current_address(), exhausted))
}
async fn try_compact(&self, store: &Arc<WedbStore<D>>, cfg: &GcConfig) -> Result<()> {
let now = now_ms();
let last = self.last_compact_ms.load(Relaxed);
if last != 0 && now.saturating_sub(last) < cfg.compaction_interval_ms {
return Ok(());
}
self.last_compact_ms.store(now, Relaxed);
let begin = store.begin_address();
let read_only = store.read_only_address();
let seg = store
.device
.segment_size()
.unwrap_or(store.hlog.config.page_size as u64);
let max = cfg.compaction_max_segments as u64;
if seg == 0 || max == 0 || read_only.saturating_sub(begin) <= max.saturating_mul(seg) {
return Ok(());
}
let n = (cfg.compaction_num_segments.min(cfg.compaction_max_segments)) as u64;
let until = read_only
.saturating_sub(seg.saturating_mul(max - n))
.max(begin);
let outcome = store.compact(until, CompactionType::Lookup).await?;
self.stats.compactions.fetch_add(1, Relaxed);
self
.stats
.last_compact_dropped
.store(outcome.dead_dropped as u64, Relaxed);
info!(
"内置 GC 紧缩完成: until={until:#x}, 丢弃={}, 释放={}B, 新起始地址={:#x}(推进 begin 后由设备 truncate_until_address 物理回收已回收段)",
outcome.dead_dropped, outcome.bytes_freed, outcome.new_begin_address
);
Ok(())
}
}
pub struct RunGuard<'a>(&'a AtomicBool);
impl Drop for RunGuard<'_> {
fn drop(&mut self) {
self.0.store(false, Release);
}
}
pub struct GcHandle<D: Device> {
inner: Arc<GcManager<D>>,
join: Mutex<Option<JoinHandle<()>>>,
}
impl<D: Device> GcHandle<D> {
pub fn stop(&self) {
self.inner.cancel.store(true, Relaxed);
}
pub fn stats(&self) -> GcStatsSnapshot {
self.inner.stats()
}
pub fn is_finished(&self) -> bool {
self
.join
.lock()
.as_ref()
.is_none_or(JoinHandle::is_finished)
}
pub async fn run_once(&self) -> Result<()> {
self.inner.run_once().await
}
}
impl<D: Device> Drop for GcHandle<D> {
fn drop(&mut self) {
self.inner.cancel.store(true, Relaxed);
self.join.lock().take();
}
}