use std::{sync::Arc, time::Duration};
use aok::{OK, Void};
use compio::{runtime::Runtime, time::sleep};
use tempfile::{TempDir, tempdir};
use wbase::time::now_ms;
use wdev::SegmentedDevice;
use wkv::{GcConfig, GcManager, StoreConfig, TtlOpt, WedbStore};
#[ctor::ctor(unsafe)]
fn _log_init() {
log_init::init();
}
async fn open_manual(
tag: &str,
gc: GcConfig,
) -> aok::Result<(TempDir, Arc<WedbStore<SegmentedDevice>>)> {
let dir = tempdir()?;
let device = Arc::new(SegmentedDevice::single_file(
dir.path().join(format!("gc_{tag}.db")),
)?);
let mut config = StoreConfig::new(1024, 4096, 16, 0.5)?;
config.gc = gc;
config.gc.enabled = false;
let store = Arc::new(WedbStore::open(config, device)?);
Ok((dir, store))
}
#[test]
fn test_run_once_purges_expired_keys() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let (_dir, store) = open_manual("purge", GcConfig::default()).await?;
assert!(
store.gc_handle().is_none(),
"enabled=false 时不得自动启动 GC"
);
let session = store.new_session()?;
let mgr = GcManager::new(&store);
let (dead, live, later) = (
b"gc:dead".as_slice(),
b"gc:live".as_slice(),
b"gc:later".as_slice(),
);
for k in [dead, live, later] {
session.upsert(k, b"v").await?;
}
assert_eq!(
session.expire_at(dead, now_ms() + 50, TtlOpt::NONE).await?,
1
);
assert_eq!(
session
.expire_at(later, now_ms() + 60_000, TtlOpt::NONE)
.await?,
1
);
sleep(Duration::from_millis(120)).await;
mgr.run_once().await?;
assert!(!session.contains_key(dead).await?);
assert_eq!(session.read(dead).await?, None);
assert_eq!(session.read_raw(&session.ttl_key(dead)).await?, None);
assert_eq!(session.read(live).await?, Some(b"v".to_vec()));
assert_eq!(session.read_raw(&session.ttl_key(live)).await?, None);
assert_eq!(session.read(later).await?, Some(b"v".to_vec()));
assert!(session.pttl_ms(later).await? > 0);
let st = mgr.stats();
assert_eq!(st.last_scan_deleted, 1);
assert_eq!(st.expired_deleted, 1);
assert_eq!(st.compactions, 0, "未注入紧缩端口时不得计数紧缩");
assert!(
st.last_scan_scanned >= st.last_scan_deleted,
"扫描记录数不得小于同轮物理删除数"
);
assert_eq!(st.total_scanned, st.last_scan_scanned, "首轮累计即最近一轮");
mgr.run_once().await?;
assert_eq!(mgr.stats().last_scan_deleted, 0);
assert!(
mgr.stats().total_scanned >= st.total_scanned,
"累计扫描记录数必须跨轮单调不减"
);
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_max_batch_deletes_truncation() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let gc = GcConfig {
max_batch_deletes: 2,
..GcConfig::default()
};
let (_dir, store) = open_manual("cap", gc).await?;
let session = store.new_session()?;
let mgr = GcManager::new(&store);
const N: usize = 6;
for i in 0..N {
let key = format!("gc:cap:{i}");
session.upsert(key.as_bytes(), b"v").await?;
assert_eq!(
session
.expire_at(key.as_bytes(), now_ms() + 50, TtlOpt::NONE)
.await?,
1
);
}
sleep(Duration::from_millis(120)).await;
for round in 1..=3 {
mgr.run_once().await?;
assert_eq!(
mgr.stats().last_scan_deleted,
2,
"第 {round} 轮应恰删批上限个键"
);
}
assert_eq!(mgr.stats().expired_deleted, N as u64);
for i in 0..N {
assert!(
!session
.contains_key(format!("gc:cap:{i}").as_bytes())
.await?
);
}
mgr.run_once().await?;
assert_eq!(mgr.stats().last_scan_deleted, 0);
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_compaction_threshold_trigger() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let gc = GcConfig {
compaction_interval_ms: 0,
compaction_max_segments: 1,
compaction_num_segments: 1,
..GcConfig::default()
};
let (_dir, store) = open_manual("compact", gc).await?;
let session = store.new_session()?;
let mgr = GcManager::new(&store);
let dead = b"gc:cp:dead";
session.upsert(dead, b"v").await?;
assert_eq!(
session.expire_at(dead, now_ms() + 50, TtlOpt::NONE).await?,
1
);
let gone = b"gc:cp:gone";
session.upsert(gone, b"v").await?;
assert!(session.delete(gone).await?);
session.upsert(b"gc:cp:warm", &[b'x'; 100]).await?;
mgr.run_once().await?;
assert_eq!(mgr.stats().compactions, 0, "未超阈值不得触发紧缩");
sleep(Duration::from_millis(120)).await;
for i in 0..120u32 {
let key = format!("gc:cp:bulk:{i}");
session.upsert(key.as_bytes(), &vec![b'v'; 1024]).await?;
}
let begin = store.begin_address();
let read_only = store.read_only_address();
assert!(
read_only - begin > 4096,
"前置:日志跨度应已超过阈值,实际 read_only={read_only:#x} begin={begin:#x}"
);
mgr.run_once().await?;
let st = mgr.stats();
assert_eq!(
st.expired_deleted, 1,
"只读线以下的冷区 TTL 同样必须主动回收"
);
assert_eq!(st.compactions, 1, "超阈值后应恰好触发一轮紧缩");
assert!(
st.last_compact_dropped >= 2,
"紧缩应丢弃 gone 的数据记录与墓碑"
);
assert!(
store.begin_address() > begin,
"紧缩后 begin_address 必须推进: {} -> {}",
begin,
store.begin_address()
);
assert_eq!(session.read(b"gc:cp:warm").await?, Some(vec![b'x'; 100]));
for i in 0..120u32 {
assert!(
session
.contains_key(format!("gc:cp:bulk:{i}").as_bytes())
.await?,
"紧缩后 bulk:{i} 不得丢失"
);
}
assert!(!session.contains_key(gone).await?);
mgr.run_once().await?;
assert!(!session.contains_key(dead).await?);
assert_eq!(session.read(dead).await?, None);
assert_eq!(session.read_raw(&session.ttl_key(dead)).await?, None);
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_open_shared_auto_spawn_and_stop() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let device = Arc::new(SegmentedDevice::single_file(dir.path().join("gc_auto.db"))?);
let mut config = StoreConfig::new(1024, 4096, 16, 0.5)?;
config.gc = GcConfig {
enabled: true,
scan_interval_ms: 30,
..GcConfig::default()
};
let store = WedbStore::open_shared(config, device)?;
assert!(
store.gc_handle().is_some(),
"open_shared 必须自动启动内置 GC"
);
let session = store.new_session()?;
let key = b"gc:auto:key";
session.upsert(key, b"v").await?;
assert_eq!(
session.expire_at(key, now_ms() + 40, TtlOpt::NONE).await?,
1
);
for _ in 0..150 {
if store
.gc_handle()
.is_some_and(|h| h.stats().expired_deleted >= 1)
{
break;
}
sleep(Duration::from_millis(20)).await;
}
let handle = store.gc_handle().expect("GC 句柄必须存在");
assert!(
handle.stats().expired_deleted >= 1,
"后台循环应在数个间隔内物理删除过期键"
);
assert!(!session.contains_key(key).await?);
assert_eq!(session.read_raw(&session.ttl_key(key)).await?, None);
handle.stop();
for _ in 0..200 {
if handle.is_finished() {
break;
}
sleep(Duration::from_millis(10)).await;
}
assert!(handle.is_finished(), "stop 后后台循环必须退出");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_cold_expiration_with_scan_budget() -> Void {
Runtime::new()?.block_on(async {
let gc = GcConfig {
max_scan_records: 1,
compaction_max_segments: 0,
..GcConfig::default()
};
let (_dir, store) = open_manual("cold_budget", gc).await?;
let session = store.new_session()?;
for i in 0..5 {
session
.upsert(format!("prefix:{i}").as_bytes(), b"v")
.await?;
}
let key = b"cold:expired";
session.upsert(key, b"v").await?;
session.expire_at(key, now_ms() + 50, TtlOpt::NONE).await?;
let ttl_addr_limit = store.tail_address();
for i in 0..120 {
session
.upsert(format!("bulk:{i}").as_bytes(), &[b'v'; 1024])
.await?;
}
assert!(store.read_only_address() > ttl_addr_limit);
sleep(Duration::from_millis(120)).await;
let mgr = GcManager::new(&store);
let mut scanned_sum = 0u64;
for _ in 0..5 {
mgr.run_once().await?;
assert_eq!(mgr.stats().expired_deleted, 0, "记录预算应让扫描分轮执行");
assert!(
mgr.stats().last_scan_scanned >= 1,
"冷区预算为 1 且积压未清时单轮扫描不得为空"
);
scanned_sum += mgr.stats().last_scan_scanned;
}
assert_eq!(
mgr.stats().total_scanned,
scanned_sum,
"累计扫描数必须为各轮求和"
);
for _ in 0..32 {
mgr.run_once().await?;
if mgr.stats().expired_deleted == 1 {
break;
}
}
assert_eq!(mgr.stats().expired_deleted, 1, "游标必须覆盖磁盘冷区");
assert!(
mgr.stats().total_scanned > scanned_sum,
"删除轮的扫描数必须继续计入累计值"
);
assert_eq!(session.read_raw(&session.ttl_key(key)).await?, None);
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_hot_window_priority_over_cold_backlog() -> Void {
Runtime::new()?.block_on(async {
let gc = GcConfig {
max_scan_records: 1,
..GcConfig::default()
};
let (_dir, store) = open_manual("hot_pri", gc).await?;
let session = store.new_session()?;
for i in 0..32 {
session
.upsert(format!("bulk:{i}").as_bytes(), &[b'v'; 1024])
.await?;
}
let hot = b"hot:expired";
session.upsert(hot, b"v").await?;
assert_eq!(
session.expire_at(hot, now_ms() + 50, TtlOpt::NONE).await?,
1
);
assert!(
store.read_only_address() < store.tail_address(),
"前置:热区窗口必须非空"
);
sleep(Duration::from_millis(120)).await;
let mgr = GcManager::new(&store);
mgr.run_once().await?;
assert_eq!(
mgr.stats().last_scan_deleted,
1,
"热区过期键必须一轮内清除,不受冷区积压影响"
);
assert!(!session.contains_key(hot).await?);
assert_eq!(session.read_raw(&session.ttl_key(hot)).await?, None);
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_gc_config_hot_update() -> Void {
Runtime::new()?.block_on(async {
let (_dir, store) = open_manual("hotcfg", GcConfig::default()).await?;
assert!(store.gc_handle().is_none(), "前置:GC 未启动");
assert!(!store.gc_config().enabled);
store.update_gc_config(|c| {
c.enabled = true;
c.scan_interval_ms = 30;
});
assert_eq!(store.gc_config().scan_interval_ms, 30);
assert!(store.start_gc(), "热更新 enabled 后 start_gc 必须成功");
assert!(store.gc_handle().is_some());
store.update_gc_config(|c| c.scan_interval_ms = 250);
assert_eq!(store.gc_config().scan_interval_ms, 250);
store.gc_handle().unwrap().stop();
aok::Result::<()>::Ok(())
})?;
OK
}