use std::{
slice,
sync::{Arc, mpsc},
thread,
time::Duration,
};
use aok::{OK, Void};
use compio::runtime::Runtime;
use log::info;
use tempfile::tempdir;
use wdev::SegmentedDevice;
use wepoch::LightEpoch;
use whlog::{
DEFAULT_INITIAL_ADDRESS, Error, HybridLog, HybridLogConfig, RecordOutput, SECTOR_ALIGNMENT,
};
use wrecord::{HEADER_SIZE, encode_to_slice};
#[test]
fn test_append_and_memory_read() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_test1.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let config = HybridLogConfig::new(64 * 1024, 16, 0.5)?;
let hlog = HybridLog::new(config, device, epoch)?;
let k1 = b"user:1001";
let v1 = b"alice_data";
let addr1 = hlog.append(k1, v1, 0, false)?;
assert_eq!(addr1, DEFAULT_INITIAL_ADDRESS);
let k2 = b"user:1002";
let v2 = b"bob_payload_string";
let addr2 = hlog.append(k2, v2, addr1, false)?;
assert!(addr2 > addr1);
assert!(hlog.is_in_memory(addr1));
assert!(hlog.is_in_memory(addr2));
assert!(hlog.is_mutable(addr1));
assert!(hlog.is_mutable(addr2));
let out1 = hlog.read_record(addr1).await?;
assert!(matches!(out1, RecordOutput::Memory(_)));
assert_eq!(out1.key()?, k1);
assert_eq!(out1.value()?, v1);
assert_eq!(out1.prev_address()?, 0);
assert!(!out1.is_tombstone()?);
let out2 = hlog.read_record(addr2).await?;
assert!(matches!(out2, RecordOutput::Memory(_)));
assert_eq!(out2.key()?, k2);
assert_eq!(out2.value()?, v2);
assert_eq!(out2.prev_address()?, addr1);
info!("单页追加与内存直读测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_page_turn_and_padding() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_test3.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let page_size = SECTOR_ALIGNMENT; let config = HybridLogConfig::new(page_size, 16, 0.5)?;
let hlog = HybridLog::new(config, device, epoch)?;
let mut addrs = Vec::new();
let val_400 = vec![b'A'; 400];
for i in 0..9 {
let key = format!("k:{i:06}");
let addr = hlog.append(key.as_bytes(), &val_400, 0, false)?;
addrs.push(addr);
}
let overflow_key = b"k:overflow";
let overflow_val = vec![b'B'; 300];
let overflow_addr = hlog.append(overflow_key, &overflow_val, 0, false)?;
assert_eq!(
overflow_addr, page_size as u64,
"换页后新记录必须位于下一页开头"
);
let pad_addr = 3880u64;
let pad_res = hlog.read_record(pad_addr).await;
assert!(
matches!(pad_res, Err(Error::PadRecord(a)) if a == pad_addr),
"读取填充位置应返回 PadRecord 错误"
);
let out = hlog.read_record(overflow_addr).await?;
assert_eq!(out.key()?, overflow_key);
assert_eq!(out.value()?, &overflow_val[..]);
info!("跨页换页与 Padding 验证通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_scan_multipage_and_pull_iterator() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_scan_test.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let page_size = SECTOR_ALIGNMENT; let config = HybridLogConfig::new(page_size, 16, 0.5)?;
let hlog = HybridLog::new(config, device, epoch)?;
let mut expected_records = Vec::new();
for i in 0..30 {
let key = format!("scan_k:{i:04}");
let val = vec![(i & 0xFF) as u8; 200];
let addr = hlog.append(key.as_bytes(), &val, 0, false)?;
expected_records.push((addr, key.into_bytes(), val));
}
let mut scanned_records = Vec::new();
hlog
.scan(0, hlog.tail_address(), |addr, rec| {
scanned_records.push((addr, rec.key().to_vec(), rec.value().to_vec()));
Ok(true)
})
.await?;
assert_eq!(scanned_records.len(), expected_records.len());
for (actual, expected) in scanned_records.iter().zip(expected_records.iter()) {
assert_eq!(actual.0, expected.0, "地址不一致");
assert_eq!(actual.1, expected.1, "Key 不一致");
assert_eq!(actual.2, expected.2, "Value 不一致");
}
let mut iter = hlog.scan_iter(0, hlog.tail_address());
let mut pulled_records = Vec::new();
while let Some((addr, out)) = iter.next().await? {
pulled_records.push((addr, out.key()?.to_vec(), out.value()?.to_vec()));
}
assert_eq!(pulled_records.len(), expected_records.len());
for (actual, expected) in pulled_records.iter().zip(expected_records.iter()) {
assert_eq!(actual.0, expected.0);
assert_eq!(actual.1, expected.1);
assert_eq!(actual.2, expected.2);
}
info!("多页连续 Scan 扫描与拉模式迭代器测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_scan_hybrid_disk_and_memory() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_scan_hybrid.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let page_size = SECTOR_ALIGNMENT; let config = HybridLogConfig::new(page_size, 16, 0.5)?;
let hlog = HybridLog::new(config, device, epoch)?;
let mut all_records = Vec::new();
for i in 0..5 {
let key = format!("cold_k:{}", i);
let val = vec![0x11; 100];
let addr = hlog.append(key.as_bytes(), &val, 0, false)?;
all_records.push((addr, key.into_bytes(), val));
}
let fill_key = b"fill";
let rem = page_size - hlog.config.page_offset(hlog.tail_address());
let fill_val = vec![0x00; rem - HEADER_SIZE - fill_key.len()];
let _ = hlog.append(fill_key, &fill_val, 0, false)?;
for i in 0..5 {
let key = format!("hot_k:{}", i);
let val = vec![0x22; 100];
let addr = hlog.append(key.as_bytes(), &val, 0, false)?;
all_records.push((addr, key.into_bytes(), val));
}
hlog.flush_page(0).await?;
hlog.shift_read_only_address(page_size as u64);
hlog.shift_head_address(page_size as u64);
assert!(hlog.is_on_disk(all_records[0].0));
assert!(hlog.is_in_memory(all_records[5].0));
let mut scanned = Vec::new();
hlog
.scan(0, hlog.tail_address(), |addr, rec| {
if rec.key() != fill_key {
scanned.push((addr, rec.key().to_vec(), rec.value().to_vec()));
}
Ok(true)
})
.await?;
assert_eq!(scanned.len(), all_records.len());
for (actual, expected) in scanned.iter().zip(all_records.iter()) {
assert_eq!(actual.0, expected.0);
assert_eq!(actual.1, expected.1);
assert_eq!(actual.2, expected.2);
}
info!("混合冷热数据穿透连续 Scan 扫描测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_scan_early_termination() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_scan_early_stop.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let config = HybridLogConfig::new(SECTOR_ALIGNMENT, 16, 0.5)?;
let hlog = HybridLog::new(config, device, epoch)?;
for i in 0..20 {
let key = format!("k:{i:02}");
let val = b"data";
let _ = hlog.append(key.as_bytes(), val, 0, false)?;
}
let mut count = 0;
hlog
.scan(0, hlog.tail_address(), |_addr, _rec| {
count += 1;
if count == 5 {
Ok(false) } else {
Ok(true)
}
})
.await?;
assert_eq!(count, 5, "Scan 应在第 5 条记录处成功提前终止");
info!("Push-based Scan 提前终止测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_subheader_fragment_pad() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_fragment.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let page_size = SECTOR_ALIGNMENT;
let config = HybridLogConfig::new(page_size, 16, 0.5)?;
let hlog = HybridLog::new(config, device, epoch)?;
let v = vec![7u8; 4000];
let addr1 = hlog.append(b"0123456789", &v, 0, false)?;
assert_eq!(addr1, DEFAULT_INITIAL_ADDRESS);
let addr2 = hlog.append(b"frag", b"tail", 0, false)?;
assert_eq!(addr2, page_size as u64, "残片后新记录必须落于下一页开头");
let fragment_addr = addr1 + 4026;
assert_eq!(fragment_addr, page_size as u64 - 6);
assert!(matches!(
hlog.read_record(fragment_addr).await,
Err(Error::PadRecord(_))
));
let mut scanned = Vec::new();
hlog
.scan(0, hlog.tail_address(), |addr, rec| {
scanned.push((addr, rec.key().to_vec()));
Ok(true)
})
.await?;
assert_eq!(
scanned,
vec![(addr1, b"0123456789".to_vec()), (addr2, b"frag".to_vec())]
);
info!("页尾子头残片处理测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_scan_inflight_zero_header_respin() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_inflight.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let config = HybridLogConfig::new(SECTOR_ALIGNMENT, 16, 1.0)?;
let hlog = Arc::new(HybridLog::new(config, device, epoch)?);
let addr0 = hlog.append(b"k0", b"v0", 0, false)?;
let addr1 = hlog.append(b"k1", b"v1", 0, false)?;
let addr2 = hlog.append(b"k2", b"v2", 0, false)?;
let page1 = hlog.config.page_id(addr1);
let off1 = hlog.config.page_offset(addr1);
let rec1_len = HEADER_SIZE + 4;
{
let mut guard = hlog.buffer.write_page(page1);
guard[off1..off1 + rec1_len].fill(0);
}
let (encode_tx, encode_rx) = mpsc::channel::<()>();
let producer = {
let hlog = Arc::clone(&hlog);
thread::spawn(move || {
if encode_rx.recv().is_err() {
return;
}
thread::sleep(Duration::from_micros(50));
let slot = hlog.buffer.page_idx(page1);
let ptr = unsafe { hlog.buffer.raw_page_ptr_mut(slot) };
let dst = unsafe { slice::from_raw_parts_mut(ptr.add(off1), rec1_len) };
let _ = encode_to_slice(dst, 0, b"k1", b"v1", false);
})
};
let mut scanned = Vec::new();
hlog
.scan(0, hlog.tail_address(), |addr, rec| {
scanned.push((addr, rec.key().to_vec()));
if addr == addr0 {
let _ = encode_tx.send(());
}
Ok(true)
})
.await?;
producer.join().unwrap();
assert_eq!(
scanned,
vec![
(addr0, b"k0".to_vec()),
(addr1, b"k1".to_vec()),
(addr2, b"k2".to_vec())
],
"在途零头自旋重试后必须原址读出 k1,且同页后续记录不得漏扫"
);
info!("在途预留零头自旋重试测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}