use std::{
sync::{
Arc,
atomic::{AtomicU8, AtomicU64, AtomicUsize, Ordering},
},
thread,
};
use aok::{OK, Void};
use compio::runtime::Runtime;
use log::info;
use tempfile::tempdir;
use wdev::{BufferPool, Device, Result as DeviceResult, SegmentedDevice};
use wepoch::LightEpoch;
use whlog::{
DEFAULT_INITIAL_ADDRESS, Error, HybridLog, HybridLogConfig, PageFlushRange, PendingFlushList,
RecordOutput, SECTOR_ALIGNMENT,
};
use wram::AlignedBuf;
use wrecord::HEADER_SIZE;
#[ctor::ctor(unsafe)]
fn _log_init() {
log_init::init();
}
#[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_in_place_update_and_protection() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_test2.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 key = b"counter:1";
let val = b"value_01";
let addr = hlog.append(key, val, 0, false)?;
let new_val = b"value_99";
let updated = hlog.try_update_in_place(addr, key, new_val)?;
assert!(updated, "可变区原位更新应成功");
let out = hlog.read_record(addr).await?;
assert_eq!(out.value()?, new_val);
let bad_val = b"value_longer_than_original";
let updated_fail = hlog.try_update_in_place(addr, key, bad_val)?;
assert!(!updated_fail, "值长度不匹配时不应允许原位覆写");
let next_page_start = 64 * 1024;
hlog.shift_read_only_address(next_page_start);
assert!(hlog.is_read_only(addr));
assert!(!hlog.is_mutable(addr));
let ro_update = hlog.try_update_in_place(addr, key, b"value_02")?;
assert!(!ro_update, "只读区不可进行原位修改");
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_flush_and_cold_disk_read() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_test4.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 k0 = b"record:page0";
let v0 = b"hello_page_zero_data";
let addr0 = hlog.append(k0, v0, 0, false)?;
assert!(addr0 < page_size as u64);
let huge_val = vec![b'X'; 4000];
let addr_huge = hlog.append(b"fill", &huge_val, 0, false)?;
assert!(addr_huge >= page_size as u64, "huge 记录已换页到第 1 页");
let k1 = b"record:page1";
let v1 = b"hello_page_one_data";
let addr1 = hlog.append(k1, v1, addr0, false)?;
assert!(addr1 >= page_size as u64);
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(addr0), "第 0 页应落在磁盘区");
assert!(!hlog.is_in_memory(addr0), "第 0 页不应再被视为在内存中");
assert!(hlog.is_in_memory(addr1), "第 1 页依然驻留内存");
let cold_out = hlog.read_record(addr0).await?;
assert!(
matches!(cold_out, RecordOutput::Disk(_)),
"冷数据必须从磁盘缓冲读取"
);
assert_eq!(cold_out.key()?, k0);
assert_eq!(cold_out.value()?, v0);
assert_eq!(cold_out.prev_address()?, 0);
let hot_out = hlog.read_record(addr1).await?;
assert!(matches!(hot_out, RecordOutput::Memory(_)));
assert_eq!(hot_out.key()?, k1);
assert_eq!(hot_out.value()?, v1);
assert_eq!(hot_out.prev_address()?, addr0);
info!("换页异步落盘与冷数据磁盘读取验证通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_rcu_version_chain() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_test5.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 key = b"my_key";
let addr_v1 = hlog.append(key, b"version_1", 0, false)?;
let addr_v2 = hlog.append(key, b"version_2", addr_v1, false)?;
let addr_v3 = hlog.append(key, b"version_3", addr_v2, false)?;
let rec3 = hlog.read_record(addr_v3).await?;
assert_eq!(rec3.value()?, b"version_3");
assert_eq!(rec3.prev_address()?, addr_v2);
let rec2 = hlog.read_record(rec3.prev_address()?).await?;
assert_eq!(rec2.value()?, b"version_2");
assert_eq!(rec2.prev_address()?, addr_v1);
let rec1 = hlog.read_record(rec2.prev_address()?).await?;
assert_eq!(rec1.value()?, b"version_1");
assert_eq!(rec1.prev_address()?, 0);
info!("RCU 多版本链表反向追踪验证通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_pending_flush_list_coalesce() -> Void {
let list = PendingFlushList::new();
assert!(list.is_empty());
list.add(PageFlushRange::new(100, 200));
list.add(PageFlushRange::new(300, 400));
assert_eq!(list.len(), 2);
let merged = list.coalesce(PageFlushRange::new(200, 300));
assert_eq!(merged, PageFlushRange::new(100, 400));
assert!(list.is_empty(), "合并后队列中原本的相邻区间应已被取出");
OK
}
#[test]
fn test_flush_pages_range_coalesced_direct_io() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_flush_range.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let page_size = 4096usize;
let config = HybridLogConfig {
page_size,
num_pages: 8,
mutable_fraction: 0.5,
ro_lag_num: whlog::ro_lag_num_from_fraction(0.5),
initial_address: DEFAULT_INITIAL_ADDRESS,
};
let hlog = HybridLog::new(config, device, epoch)?;
let addr0 = hlog.append(b"k0", b"v0_page0", 0, false)?;
assert_eq!(hlog.config.page_id(addr0), 0);
let pad_val1 = vec![b'A'; 4000];
let _ = hlog.append(b"fill1", &pad_val1, 0, false)?;
let addr1 = hlog.append(b"k1", b"v1_page1", 0, false)?;
assert_eq!(hlog.config.page_id(addr1), 1);
let pad_val2 = vec![b'B'; 4000];
let _ = hlog.append(b"fill2", &pad_val2, 0, false)?;
let addr2 = hlog.append(b"k2", b"v2_page2", 0, false)?;
assert_eq!(hlog.config.page_id(addr2), 2);
hlog.flush_pages_range(0, 2).await?;
hlog.sync().await?;
let new_head = (page_size * 3) as u64;
hlog.shift_read_only_address(new_head);
hlog.shift_head_address(new_head);
assert!(hlog.is_on_disk(addr0));
assert!(hlog.is_on_disk(addr1));
assert!(hlog.is_on_disk(addr2));
let out0 = hlog.read_record(addr0).await?;
assert_eq!(out0.key()?, b"k0");
assert_eq!(out0.value()?, b"v0_page0");
let out1 = hlog.read_record(addr1).await?;
assert_eq!(out1.key()?, b"k1");
assert_eq!(out1.value()?, b"v1_page1");
let out2 = hlog.read_record(addr2).await?;
assert_eq!(out2.key()?, b"k2");
assert_eq!(out2.value()?, b"v2_page2");
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_recover_and_snapshot_invariants() -> Void {
use whlog::AddressSnapshot;
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_recover.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 valid_snapshot = AddressSnapshot::new(
64, 4096, 4096, 4096, 8192, 10000, 4096, );
assert!(valid_snapshot.validate());
let recovered_hlog = HybridLog::recover(
config.clone(),
Arc::clone(&device),
Arc::clone(&epoch),
valid_snapshot,
)
.await;
assert!(recovered_hlog.is_ok(), "合法快照恢复应成功");
let hlog = recovered_hlog.unwrap();
assert_eq!(hlog.tail_address(), 10000);
assert_eq!(hlog.head_address(), 4096);
assert!(hlog.addresses.validate_invariants());
let invalid_snapshot1 = AddressSnapshot::new(
64, 4096, 8192, 8192, 8192, 10000, 4096, );
assert!(!invalid_snapshot1.validate());
let fail1 = HybridLog::recover(
config.clone(),
Arc::clone(&device),
Arc::clone(&epoch),
invalid_snapshot1,
)
.await;
assert!(
matches!(fail1, Err(Error::InvalidState(_))),
"非法快照必须被拦截并返回 InvalidState"
);
let invalid_snapshot2 = AddressSnapshot::new(64, 8192, 4096, 4096, 8192, 10000, 4096);
assert!(!invalid_snapshot2.validate());
let fail2 = HybridLog::recover(config, device, epoch, invalid_snapshot2).await;
assert!(matches!(fail2, Err(Error::InvalidState(_))));
info!("快照恢复与状态机不变式校验测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_recovery_preload_and_resume_append() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_recover_resume.db");
let page_size = SECTOR_ALIGNMENT; let config = HybridLogConfig::new(page_size, 16, 0.5)?;
let addr0;
let addr1;
let tail_before;
{
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let hlog = HybridLog::new(config.clone(), device, epoch)?;
addr0 = hlog.append(b"k0", b"v0_before_crash", 0, false)?;
let pad_val = vec![b'P'; 4000];
let _ = hlog.append(b"pad", &pad_val, 0, false)?;
addr1 = hlog.append(b"k1", b"v1_page1", addr0, false)?;
assert_eq!(hlog.config.page_id(addr1), 1);
hlog.flush_all().await?;
hlog.sync().await?;
tail_before = hlog.tail_address();
}
{
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let snapshot = whlog::AddressSnapshot::from_bounds(
DEFAULT_INITIAL_ADDRESS,
DEFAULT_INITIAL_ADDRESS, tail_before, tail_before, tail_before, );
assert!(snapshot.validate());
let hlog = HybridLog::recover(config, device, epoch, snapshot).await?;
assert_eq!(hlog.tail_address(), tail_before);
let out0 = hlog.read_record(addr0).await?;
assert_eq!(out0.key()?, b"k0");
assert_eq!(out0.value()?, b"v0_before_crash");
let out1 = hlog.read_record(addr1).await?;
assert_eq!(out1.key()?, b"k1");
assert_eq!(out1.value()?, b"v1_page1");
let addr2 = hlog.append(b"k2", b"v2_after_recovery", addr1, false)?;
assert_eq!(addr2, tail_before);
let out2 = hlog.read_record(addr2).await?;
assert_eq!(out2.key()?, b"k2");
assert_eq!(out2.value()?, b"v2_after_recovery");
assert_eq!(out2.prev_address()?, addr1);
info!("崩溃恢复后磁盘页预热加载与无缝继续追加测试通过");
}
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_iterate_version_chain() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_version_chain.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 key = b"chain_key";
let addr1 = hlog.append(key, b"v1", 0, false)?;
let addr2 = hlog.append(key, b"v2", addr1, false)?;
let addr3 = hlog.append(key, b"v3", addr2, false)?;
let addr4 = hlog.append(key, b"v4", addr3, false)?;
let mut collected = Vec::new();
hlog
.iterate_version_chain(addr4, |_addr, rec| {
collected.push(rec.value()?.to_vec());
Ok(true)
})
.await?;
assert_eq!(
collected,
vec![
b"v4".to_vec(),
b"v3".to_vec(),
b"v2".to_vec(),
b"v1".to_vec()
]
);
let mut truncated = Vec::new();
hlog
.iterate_version_chain(addr4, |_addr, rec| {
let val = rec.value()?.to_vec();
let stop = val == b"v3";
truncated.push(val);
Ok(!stop)
})
.await?;
assert_eq!(truncated, vec![b"v4".to_vec(), b"v3".to_vec()]);
info!("历史版本反向链表回溯测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_shift_read_only_to_tail_and_flush_all() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_flush_all.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)?;
let addr = hlog.append(b"freeze_k", b"freeze_v", 0, false)?;
let tail = hlog.tail_address();
let frozen_tail = hlog.shift_read_only_to_tail();
assert_eq!(frozen_tail, tail);
assert!(hlog.is_read_only(addr));
assert!(!hlog.is_mutable(addr));
let flushed = hlog.flush_all().await?;
assert!(flushed >= tail);
assert_eq!(hlog.flushed_until_address(), flushed);
info!("shift_read_only_to_tail 与 flush_all 测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_bitcode_roundtrip() -> Void {
use whlog::AddressSnapshot;
let snap = AddressSnapshot::new(64, 4096, 4096, 8192, 8192, 16384, 8192);
let encoded = bitcode::encode(&snap);
let decoded: AddressSnapshot = bitcode::decode(&encoded)?;
assert_eq!(snap, decoded);
let range = PageFlushRange::new(4096, 8192);
let encoded_range = bitcode::encode(&range);
let decoded_range: PageFlushRange = bitcode::decode(&encoded_range)?;
assert_eq!(range, decoded_range);
info!("bitcode 序列化往返测试通过");
OK
}
#[test]
fn test_concurrent_append_stress() -> Void {
use std::sync::Arc;
use whasher::HashSet;
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_stress.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(32));
let config = HybridLogConfig::new(SECTOR_ALIGNMENT, 16, 0.5)?;
let hlog = Arc::new(HybridLog::new(config, device, epoch.clone())?);
const THREADS: usize = 8;
const PER_THREAD: usize = 50;
let mut handles = Vec::new();
for t in 0..THREADS {
let hlog = Arc::clone(&hlog);
handles.push(thread::spawn(move || {
let mut addrs = Vec::with_capacity(PER_THREAD);
for i in 0..PER_THREAD {
let key = format!("t{t}:k{i:04}");
let val = vec![(t * PER_THREAD + i) as u8; 20];
let addr = hlog.append(key.as_bytes(), &val, 0, false)?;
addrs.push((addr, key.into_bytes(), val));
}
Ok::<_, Error>(addrs)
}));
}
let mut all = Vec::new();
for h in handles {
all.extend(h.join().unwrap()?);
}
let unique: HashSet<u64> = all.iter().map(|(a, ..)| *a).collect();
assert_eq!(unique.len(), THREADS * PER_THREAD, "并发追加地址不得重叠");
assert!(
all
.iter()
.all(|&(a, ..)| a >= DEFAULT_INITIAL_ADDRESS && a < hlog.tail_address())
);
let _participant = epoch.register()?;
for (addr, key, val) in &all {
let probed = hlog.with_memory_record(*addr, |rec| {
assert_eq!(rec.key(), key.as_slice(), "回读键必须一致");
assert_eq!(rec.value(), val.as_slice(), "回读值必须与写入内容一致");
Ok(())
})?;
assert!(probed.is_some(), "记录必须仍驻留内存: addr={addr:#x}");
}
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.len(), THREADS * PER_THREAD);
let mut expected_keys: HashSet<Vec<u8>> = HashSet::default();
for (.., key, _) in &all {
expected_keys.insert(key.clone());
}
let actual_keys: HashSet<Vec<u8>> = scanned.into_iter().map(|(_, k)| k).collect();
assert_eq!(actual_keys, expected_keys, "并发追加记录内容必须完整无缺");
hlog.flush_all().await?;
assert!(hlog.addresses.validate_invariants());
info!("多线程无锁并发追加测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_ring_wraparound_eviction() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_wrap.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let config = HybridLogConfig::new(SECTOR_ALIGNMENT, 2, 0.5)?;
let hlog = HybridLog::new(config, device, epoch)?;
let big = vec![b'A'; 3900];
let addr_a = hlog.append(b"a", &big, 0, false)?;
let addr_b = hlog.append(b"b", &big, 0, false)?;
assert_eq!(hlog.config.page_id(addr_a), 0);
assert_eq!(hlog.config.page_id(addr_b), 1);
let err = hlog.append(b"c", &big, 0, false).unwrap_err();
assert!(
matches!(err, Error::PageNotReady(2)),
"回绕必须被拦截: {err:?}"
);
hlog.flush_page(0).await?;
hlog.shift_read_only_address(SECTOR_ALIGNMENT as u64);
hlog.shift_head_address(SECTOR_ALIGNMENT as u64);
while hlog.safe_head_address() < SECTOR_ALIGNMENT as u64 {
hlog.epoch.bump_epoch();
}
let addr_c = hlog.append(b"c", &big, 0, false)?;
assert_eq!(
addr_c,
2 * SECTOR_ALIGNMENT as u64,
"重试后必须落在回绕页开头"
);
assert!(hlog.addresses.validate_invariants());
let out_a = hlog.read_record(addr_a).await?;
assert!(matches!(out_a, RecordOutput::Disk(_)));
assert_eq!(out_a.key()?, b"a");
assert!(hlog.is_in_memory(addr_b));
info!("环形缓冲区回绕驱逐测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_inplace_lifecycle() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_inplace.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 key = b"life";
let addr = hlog.append(key, b"1234567890", 0, false)?;
assert_eq!(addr, DEFAULT_INITIAL_ADDRESS);
hlog.revivify_record_at(addr, 35, key, b"1234567890", 0, false)?;
let out = hlog.read_record(addr).await?;
assert_eq!(out.value()?, b"1234567890");
assert_eq!(out.header()?.filler_bytes(), 5, "富余空间必须转为松弛填充");
assert!(hlog.try_update_in_place(addr, key, b"012345678901234")?);
assert_eq!(hlog.read_record(addr).await?.value()?, b"012345678901234");
let r = hlog.try_modify_record_in_place(addr, key, |v| {
v[0] = b'X';
Some(())
})?;
assert!(r.is_some());
assert_eq!(hlog.read_record(addr).await?.value()?, b"X12345678901234");
assert!(!hlog.try_update_in_place(addr, b"wrong", b"y")?);
assert!(!hlog.try_update_in_place(addr, key, &[b'z'; 16])?);
assert!(hlog.try_mark_tombstone_in_place(addr, key)?);
assert!(hlog.read_record(addr).await?.is_tombstone()?);
assert!(!hlog.try_mark_tombstone_in_place(addr, key)?);
assert!(hlog.try_revivify_in_chain(addr, key, b"revived!")?);
let out = hlog.read_record(addr).await?;
assert!(!out.is_tombstone()?);
assert_eq!(out.value()?, b"revived!");
assert!(!hlog.try_revivify_in_chain(addr, key, b"again")?);
info!("原位更新全生命周期测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_revivify_record_at_with_pad() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_reviv.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 val21 = vec![b'x'; 21];
let addr1 = hlog.append(b"old", &val21, 0, false)?;
let addr2 = hlog.append(b"b", b"vvv", 0, false)?;
hlog.revivify_record_at(addr1, 40, b"new", b"value", 0, false)?;
let out = hlog.read_record(addr1).await?;
assert_eq!(out.key()?, b"new");
assert_eq!(out.value()?, b"value");
let addr3 = hlog.append(b"c", b"vvv", 0, false)?;
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"new".to_vec()),
(addr2, b"b".to_vec()),
(addr3, b"c".to_vec())
],
"槽内 Pad 必须按物理尺寸精确越过,不得吞并同页后续记录"
);
assert!(matches!(
hlog.revivify_record_at(addr1, 10, b"new", b"value", 0, false),
Err(Error::RecordTooLarge { .. })
));
assert!(matches!(
hlog.revivify_record_at(0, 100, b"x", b"y", 0, false),
Err(Error::AddressOutOfRange { .. })
));
assert_eq!(hlog.tail_address(), addr3 + 20);
info!("复活槽位 Pad 填充测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_shift_begin_address_and_truncate() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_shift_begin.db");
let device = Arc::new(SegmentedDevice::new(
&db_path,
Some(2 * SECTOR_ALIGNMENT as u64),
SECTOR_ALIGNMENT,
)?);
let epoch = Arc::new(LightEpoch::new(16));
let config = HybridLogConfig::new(SECTOR_ALIGNMENT, 16, 0.5)?;
let hlog = HybridLog::new(config, Arc::clone(&device), epoch)?;
let big = vec![b'B'; 3900];
let addr1 = hlog.append(b"k1", &big, 0, false)?;
let addr2 = hlog.append(b"k2", &big, 0, false)?;
let addr3 = hlog.append(b"k3", &big, 0, false)?;
assert_eq!(hlog.config.page_id(addr3), 2);
hlog.flush_all().await?;
hlog.sync().await?;
hlog.shift_read_only_to_tail();
let new_begin = 2 * SECTOR_ALIGNMENT as u64;
hlog.shift_begin_address(new_begin).await?;
assert_eq!(hlog.begin_address(), new_begin);
assert_eq!(hlog.head_address(), new_begin);
assert!(
hlog.addresses.validate_invariants(),
"推进后不变式必须保持: {:?}",
hlog.addresses.snapshot()
);
assert_eq!(device.get_file_size(0)?, 0, "段 0 必须被物理删除");
assert!(device.get_file_size(1)? > 0, "段 1 必须保留");
assert!(matches!(
hlog.read_record(addr1).await,
Err(Error::AddressOutOfRange { .. })
));
assert!(matches!(
hlog.read_record(addr2).await,
Err(Error::AddressOutOfRange { .. })
));
let out3 = hlog.read_record(addr3).await?;
assert_eq!(out3.key()?, b"k3");
info!("shift_begin_address 与设备段截断测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_recover_scrubs_non_durable_prefix() -> Void {
use whlog::AddressSnapshot;
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_scrub.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.clone(), Arc::clone(&device), Arc::clone(&epoch))?;
let big = vec![b'D'; 3900];
let addr1 = hlog.append(b"durable", &big, 0, false)?;
let addr2 = hlog.append(b"volatile", &big, 0, false)?; hlog.flush_all().await?;
hlog.sync().await?;
let tail = hlog.tail_address();
let snapshot = AddressSnapshot::from_bounds(
DEFAULT_INITIAL_ADDRESS,
DEFAULT_INITIAL_ADDRESS,
page_size as u64, tail,
tail,
);
let recovered = HybridLog::recover(config, device, epoch, snapshot).await?;
assert!(matches!(
recovered.read_record(addr2).await,
Err(Error::PadRecord(_))
));
let mut scanned = Vec::new();
recovered
.scan(0, tail, |addr, rec| {
scanned.push((addr, rec.key().to_vec()));
Ok(true)
})
.await?;
assert_eq!(scanned, vec![(addr1, b"durable".to_vec())]);
let addr3 = recovered.append(b"resumed", b"v3", addr2, false)?;
assert_eq!(addr3, tail);
let out3 = recovered.read_record(addr3).await?;
assert_eq!(out3.key()?, b"resumed");
assert!(recovered.addresses.validate_invariants());
info!("恢复非持久化前缀清洗测试通过");
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_config_validation() -> Void {
assert!(matches!(
HybridLogConfig::new(4095, 16, 0.5),
Err(Error::InvalidConfig(_))
));
assert!(matches!(
HybridLogConfig::new(4096, 0, 0.5),
Err(Error::InvalidConfig(_))
));
assert!(matches!(
HybridLogConfig::new(4096, 3, 0.5),
Err(Error::InvalidConfig(_))
));
assert!(matches!(
HybridLogConfig::new(4096, 16, 0.0),
Err(Error::InvalidConfig(_))
));
assert!(matches!(
HybridLogConfig::new(4096, 16, 1.5),
Err(Error::InvalidConfig(_))
));
assert!(
matches!(
HybridLogConfig::new(4096, 16, f64::NAN),
Err(Error::InvalidConfig(_))
),
"NaN 可变比例必须被拦截"
);
assert!(matches!(
HybridLogConfig::with_initial_address(4096, 16, 0.5, 63),
Err(Error::InvalidConfig(_))
));
assert!(HybridLogConfig::new(8192, 4, 1.0).is_ok());
info!("配置校验测试通过");
OK
}
#[test]
fn test_append_argument_guards() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_guards.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 overflow_prev = 1u64 << 48;
assert!(matches!(
hlog.append(b"k", b"v", overflow_prev, false),
Err(Error::InvalidAddress(_))
));
let huge = vec![0u8; page_size];
assert!(matches!(
hlog.append(b"k", &huge, 0, false),
Err(Error::RecordTooLarge { .. })
));
assert_eq!(hlog.tail_address(), DEFAULT_INITIAL_ADDRESS);
info!("追加参数边界校验测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_cold_read_large_record_exact_trim() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_probe.db");
let device = Arc::new(SegmentedDevice::single_file(&db_path)?);
let epoch = Arc::new(LightEpoch::new(16));
let page_size = 64 * 1024;
let config = HybridLogConfig::new(page_size, 16, 0.5)?;
let hlog = HybridLog::new(config, device, epoch)?;
let big_key = b"big:1";
let big_val = vec![b'A'; 5000];
let addr_big = hlog.append(big_key, &big_val, 0, false)?;
let small_key = b"s:1";
let small_val = b"tiny";
let addr_small = hlog.append(small_key, small_val, addr_big, false)?;
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(addr_big));
let out_big = hlog.read_record(addr_big).await?;
assert!(matches!(out_big, RecordOutput::Disk(_)));
assert_eq!(out_big.key()?, big_key);
assert_eq!(out_big.value()?, &big_val[..]);
assert_eq!(out_big.prev_address()?, 0);
assert_eq!(
out_big.as_slice().len(),
out_big.header()?.physical_size(),
"磁盘冷读缓冲必须按物理尺寸精确裁剪"
);
let out_small = hlog.read_record(addr_small).await?;
assert!(matches!(out_small, RecordOutput::Disk(_)));
assert_eq!(out_small.key()?, small_key);
assert_eq!(out_small.value()?, small_val);
assert_eq!(out_small.prev_address()?, addr_big);
info!("磁盘冷读大记录精确裁剪测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
struct CountingDevice {
inner: SegmentedDevice,
reads: AtomicUsize,
read_bytes: AtomicU64,
}
impl Device for CountingDevice {
#[inline]
fn sector_size(&self) -> usize {
self.inner.sector_size()
}
#[inline]
fn segment_size(&self) -> Option<u64> {
self.inner.segment_size()
}
#[inline]
fn direct_io(&self) -> bool {
self.inner.direct_io()
}
#[inline]
fn pool(&self) -> &Arc<BufferPool> {
self.inner.pool()
}
async fn write_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.inner.write_aligned(offset, buf).await
}
async fn read_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.reads.fetch_add(1, Ordering::Relaxed);
self
.read_bytes
.fetch_add(buf.capacity() as u64, Ordering::Relaxed);
self.inner.read_aligned(offset, buf).await
}
async fn read_raw(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.reads.fetch_add(1, Ordering::Relaxed);
self
.read_bytes
.fetch_add(buf.capacity() as u64, Ordering::Relaxed);
self.inner.read_raw(offset, buf).await
}
async fn sync(&self) -> DeviceResult<()> {
self.inner.sync().await
}
async fn truncate_until_segment(&self, segment_id: u32) -> DeviceResult<()> {
self.inner.truncate_until_segment(segment_id).await
}
}
#[test]
fn test_disk_read_page_cache() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_page_cache.db");
let device = Arc::new(CountingDevice {
inner: SegmentedDevice::single_file(&db_path)?,
reads: AtomicUsize::new(0),
read_bytes: AtomicU64::new(0),
});
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.clone(), epoch)?;
let mut recs: Vec<(u64, Vec<u8>, Vec<u8>)> = Vec::new();
for i in 0..7u8 {
let key = format!("k:{i:06}").into_bytes();
let val = vec![b'a' + i; 488];
let addr = hlog.append(&key, &val, 0, false)?;
recs.push((addr, key, val));
}
let tail_key = b"tail:000".to_vec();
let tail_val = vec![b'T'; 424]; let tail_addr = hlog.append(&tail_key, &tail_val, 0, false)?;
assert_eq!(tail_addr + 448, page_size as u64, "记录必须恰好抵达页尾");
recs.push((tail_addr, tail_key.clone(), tail_val.clone()));
let page1_end = 2 * page_size as u64;
let c_key = b"c:first1".to_vec();
let c_val = vec![b'c'; 100];
let c_addr = hlog.append(&c_key, &c_val, 0, false)?;
assert_eq!(c_addr, page_size as u64, "页尾满后新记录必须落在下一页开头");
let d_val = vec![b'd'; page_size - 148];
let d_addr = hlog.append(b"d:fill00", &d_val, 0, false)?;
assert_eq!(d_addr + 3972, page1_end, "第 1 页必须被记录精确填满");
let page2_end = 3 * page_size as u64;
let e_key = b"e:first2".to_vec();
let e_val = vec![b'e'; 100];
let e_addr = hlog.append(&e_key, &e_val, 0, false)?;
assert_eq!(e_addr, page2_end - page_size as u64);
let f_val = vec![b'f'; page_size - 148];
let f_addr = hlog.append(b"f:fill00", &f_val, 0, false)?;
assert_eq!(f_addr + 3972, page2_end);
for p in 0..3u64 {
hlog.flush_page(p).await?;
}
hlog.shift_read_only_address(page2_end);
hlog.shift_head_address(page2_end);
assert!(hlog.is_on_disk(tail_addr));
async fn assert_disk_read<D: Device>(
hlog: &HybridLog<D>,
addr: u64,
key: &[u8],
val: &[u8],
) -> Void {
let out = hlog.read_disk_record(addr).await?;
assert!(matches!(out, RecordOutput::Disk(_)));
assert_eq!(out.key()?, key);
assert_eq!(out.value()?, val);
OK
}
assert_disk_read(&hlog, e_addr, &e_key, &e_val).await?;
assert_disk_read(&hlog, f_addr, b"f:fill00", &f_val).await?;
for (addr, key, val) in &recs {
assert_disk_read(&hlog, *addr, key, val).await?;
}
assert_disk_read(&hlog, c_addr, &c_key, &c_val).await?;
assert_disk_read(&hlog, d_addr, b"d:fill00", &d_val).await?;
let reads_filled = device.reads.load(Ordering::Relaxed);
assert_eq!(
reads_filled, 6,
"三页 12 条点读应恰好产生 6 次设备 I/O(每页 1 次 probe + 1 次整页装载)"
);
assert_disk_read(&hlog, tail_addr, &tail_key, &tail_val).await?;
assert_disk_read(&hlog, c_addr, &c_key, &c_val).await?;
assert_eq!(
device.reads.load(Ordering::Relaxed),
reads_filled,
"二次回验必须全部命中缓存"
);
assert_disk_read(&hlog, e_addr, &e_key, &e_val).await?;
assert_disk_read(&hlog, f_addr, b"f:fill00", &f_val).await?;
assert_disk_read(&hlog, e_addr, &e_key, &e_val).await?;
assert_eq!(
device.reads.load(Ordering::Relaxed),
reads_filled + 2,
"页 2 重读应恰好产生 1 次 probe + 1 次整页重装载"
);
assert_disk_read(&hlog, tail_addr, &tail_key, &tail_val).await?;
assert_disk_read(&hlog, recs[0].0, &recs[0].1, &recs[0].2).await?;
assert_eq!(
device.reads.load(Ordering::Relaxed),
reads_filled + 4,
"槽位冲突驱逐后同页重读应恰好产生 1 次 probe + 1 次整页重装载"
);
info!("点读冷路径整页磁盘读缓存测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_disk_read_adaptive_install() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_adaptive_install.db");
let device = Arc::new(CountingDevice {
inner: SegmentedDevice::single_file(&db_path)?,
reads: AtomicUsize::new(0),
read_bytes: AtomicU64::new(0),
});
let epoch = Arc::new(LightEpoch::new(16));
let page_size = 64 * 1024;
const RECS_PER_PAGE: u64 = 16;
let val_len = 4092 - HEADER_SIZE - 8;
let config = HybridLogConfig::new(page_size, 16, 0.5)?;
let hlog = HybridLog::new(config, device.clone(), epoch)?;
let mut recs = Vec::new();
for page in 0..5u64 {
for i in 0..RECS_PER_PAGE {
let key = format!("k:{page:02}:{i:02}").into_bytes();
let val = vec![b'a' + (i as u8 % 26); val_len];
let addr = hlog.append(&key, &val, 0, false)?;
recs.push((addr, key, val));
}
}
let tail_key = b"tail:0000".to_vec();
let tail_addr = hlog.append(&tail_key, &vec![b'T'; val_len], 0, false)?;
assert_eq!(tail_addr / page_size as u64, 5, "收尾记录必须落在第 5 页");
for p in 0..=5u64 {
hlog.flush_page(p).await?;
}
let tail = hlog.tail_address();
hlog.shift_read_only_address(tail);
hlog.shift_head_address(tail);
let rec = |page: u64, i: u64| &recs[(page * RECS_PER_PAGE + i) as usize];
async fn assert_disk_read<D: Device>(
hlog: &HybridLog<D>,
addr: u64,
key: &[u8],
val: &[u8],
) -> Void {
let out = hlog.read_disk_record(addr).await?;
assert!(matches!(out, RecordOutput::Disk(_)));
assert_eq!(out.key()?, key);
assert_eq!(out.value()?, val);
OK
}
let reads_before = device.reads.load(Ordering::Relaxed);
let bytes_before = device.read_bytes.load(Ordering::Relaxed);
for i in 0..3u64 {
let (addr, key, val) = rec(3, i);
assert_disk_read(&hlog, *addr, key, val).await?;
}
assert_eq!(
device.reads.load(Ordering::Relaxed) - reads_before,
2,
"同页 3 条点读应恰好产生 2 次 I/O:1 次 probe + 1 次整页装载,第三条命中零 I/O"
);
let seq_bytes = device.read_bytes.load(Ordering::Relaxed) - bytes_before;
assert!(
seq_bytes >= (page_size + 4096) as u64 && seq_bytes < (page_size + 8192) as u64,
"顺序读应恰好产生 1 次整页装载({page_size} 字节)+ 1 次 probe(4KB 级),实际 {seq_bytes} 字节"
);
let (addr1, key1, val1) = rec(3, 1);
assert_disk_read(&hlog, *addr1, key1, val1).await?;
assert_eq!(
device.reads.load(Ordering::Relaxed) - reads_before,
2,
"已装载页重读必须全部命中缓存"
);
let reads_before = device.reads.load(Ordering::Relaxed);
let bytes_before = device.read_bytes.load(Ordering::Relaxed);
for page in [0u64, 2, 4, 0, 2, 4, 0] {
let (addr, key, val) = rec(page, 0);
assert_disk_read(&hlog, *addr, key, val).await?;
}
assert_eq!(
device.reads.load(Ordering::Relaxed) - reads_before,
7,
"均匀随机读每次未命中应恰好一次 probe 级设备 I/O"
);
let rand_bytes = device.read_bytes.load(Ordering::Relaxed) - bytes_before;
assert!(
rand_bytes < 7 * 8192,
"均匀随机读总 I/O 必须为 probe 级小读之和,不得出现任何整页装载(单次整页即 {page_size} 字节): {rand_bytes}"
);
assert_disk_read(&hlog, *addr1, key1, val1).await?;
assert_eq!(
device.reads.load(Ordering::Relaxed) - reads_before,
7,
"顺序装载页不得被随机 probe 驱逐"
);
let (addr2, key2, val2) = rec(2, 0);
assert_disk_read(&hlog, *addr2, key2, val2).await?;
assert_eq!(
device.reads.load(Ordering::Relaxed) - reads_before,
8,
"随机页重读仍应仅 probe,不得整页装载"
);
info!("磁盘读缓存连续性装载门槛测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
struct FaultDevice {
inner: SegmentedDevice,
mode: AtomicU8,
}
const MODE_NORMAL: u8 = 0;
const MODE_FAIL: u8 = 1;
const MODE_SHORT: u8 = 2;
impl Device for FaultDevice {
#[inline]
fn sector_size(&self) -> usize {
self.inner.sector_size()
}
#[inline]
fn segment_size(&self) -> Option<u64> {
self.inner.segment_size()
}
#[inline]
fn direct_io(&self) -> bool {
self.inner.direct_io()
}
#[inline]
fn pool(&self) -> &Arc<BufferPool> {
self.inner.pool()
}
async fn write_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
match self.mode.load(Ordering::Relaxed) {
MODE_FAIL => (
Err(wdev::Error::ReadOnly {
offset,
len: buf.len(),
}),
buf,
),
MODE_SHORT => (Ok(buf.len() / 2), buf),
_ => self.inner.write_aligned(offset, buf).await,
}
}
async fn read_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.inner.read_aligned(offset, buf).await
}
async fn read_raw(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.inner.read_raw(offset, buf).await
}
async fn sync(&self) -> DeviceResult<()> {
self.inner.sync().await
}
async fn truncate_until_segment(&self, segment_id: u32) -> DeviceResult<()> {
self.inner.truncate_until_segment(segment_id).await
}
}
#[test]
fn test_flush_stale_range_and_short_write() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_fault_flush.db");
let device = Arc::new(FaultDevice {
inner: SegmentedDevice::single_file(&db_path)?,
mode: AtomicU8::new(MODE_NORMAL),
});
let epoch = Arc::new(LightEpoch::new(16));
let config = HybridLogConfig::new(SECTOR_ALIGNMENT, 2, 0.5)?;
let hlog = HybridLog::new(config, device.clone(), epoch)?;
let big = vec![b'A'; 3900];
let a = hlog.append(b"a", &big, 0, false)?;
let b = hlog.append(b"b", &big, 0, false)?;
assert_eq!(hlog.config.page_id(a), 0);
assert_eq!(hlog.config.page_id(b), 1);
hlog.flush_page(0).await?;
hlog.shift_read_only_address(SECTOR_ALIGNMENT as u64);
hlog.shift_head_address(SECTOR_ALIGNMENT as u64);
while hlog.safe_head_address() < SECTOR_ALIGNMENT as u64 {
hlog.epoch.bump_epoch();
}
let c = hlog.append(b"c", &big, 0, false)?;
assert_eq!(c, 2 * SECTOR_ALIGNMENT as u64);
device.mode.store(MODE_FAIL, Ordering::Relaxed);
assert!(matches!(
hlog.flush_page(1).await,
Err(Error::Device(wdev::Error::ReadOnly { .. }))
));
device.mode.store(MODE_SHORT, Ordering::Relaxed);
assert!(matches!(
hlog.flush_page(1).await,
Err(Error::FlushFailed { .. })
));
assert_eq!(hlog.flushed_until_address(), SECTOR_ALIGNMENT as u64);
device.mode.store(MODE_NORMAL, Ordering::Relaxed);
hlog.flush_page(1).await?;
assert_eq!(hlog.flushed_until_address(), 2 * SECTOR_ALIGNMENT as u64);
hlog.shift_read_only_address(2 * SECTOR_ALIGNMENT as u64);
hlog.shift_head_address(2 * SECTOR_ALIGNMENT as u64);
while hlog.safe_head_address() < 2 * SECTOR_ALIGNMENT as u64 {
hlog.epoch.bump_epoch();
}
let d = hlog.append(b"d", &big, 0, false)?;
assert_eq!(d, 3 * SECTOR_ALIGNMENT as u64);
assert!(!hlog.buffer.is_page_loaded(1), "页 1 必须已驱逐出内存");
hlog.flush_page(0).await?;
hlog.flush_page(2).await?;
assert_eq!(hlog.flushed_until_address(), 3 * SECTOR_ALIGNMENT as u64);
assert!(hlog.pending_flush.is_empty(), "陈旧区间必须被钳制丢弃");
assert!(hlog.addresses.validate_invariants());
let out_a = hlog.read_record(a).await?;
assert_eq!(out_a.key()?, b"a");
let out_d = hlog.read_record(d).await?;
assert_eq!(out_d.key()?, b"d");
info!("刷盘陈旧区间钳制与短写防护测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_recover_multisegment_span_and_window_guard() -> Void {
use whlog::AddressSnapshot;
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("hlog_multiseg_recover.db");
let device = Arc::new(SegmentedDevice::new(
&db_path,
Some(2 * SECTOR_ALIGNMENT as u64),
SECTOR_ALIGNMENT,
)?);
let epoch = Arc::new(LightEpoch::new(16));
let page_size = SECTOR_ALIGNMENT;
let config = HybridLogConfig::new(page_size, 16, 0.5)?;
let mut addrs = Vec::new();
let tail_before;
{
let hlog = HybridLog::new(config.clone(), device, Arc::clone(&epoch))?;
let big = vec![b'M'; 3900];
for i in 0..5u8 {
let key = format!("m{i}").into_bytes();
let addr = hlog.append(&key, &big, 0, false)?;
assert_eq!(hlog.config.page_id(addr), i as u64, "每页恰好一条记录");
addrs.push((addr, key));
}
hlog.flush_all().await?;
hlog.sync().await?;
tail_before = hlog.tail_address();
}
{
let bad_snapshot = AddressSnapshot::from_bounds(
DEFAULT_INITIAL_ADDRESS,
DEFAULT_INITIAL_ADDRESS,
64 * page_size as u64 + DEFAULT_INITIAL_ADDRESS, 64 * page_size as u64 + DEFAULT_INITIAL_ADDRESS, 64 * page_size as u64 + DEFAULT_INITIAL_ADDRESS, );
assert!(bad_snapshot.validate(), "快照本身须合法以触达窗口守卫");
let result = HybridLog::recover(
config.clone(),
Arc::new(SegmentedDevice::segmented(
&db_path,
2 * SECTOR_ALIGNMENT as u64,
)?),
Arc::clone(&epoch),
bad_snapshot,
)
.await;
assert!(
matches!(result, Err(Error::InvalidState(msg)) if msg.contains("驻留窗口")),
"跨页窗口超限必须被环形页数守卫拦截"
);
}
let snapshot = AddressSnapshot::from_bounds(
DEFAULT_INITIAL_ADDRESS,
DEFAULT_INITIAL_ADDRESS,
tail_before,
tail_before,
tail_before,
);
let hlog = HybridLog::recover(
config,
Arc::new(SegmentedDevice::segmented(
&db_path,
2 * SECTOR_ALIGNMENT as u64,
)?),
epoch,
snapshot,
)
.await?;
for (addr, key) in &addrs {
let out = hlog.read_record(*addr).await?;
assert_eq!(out.key()?, key.as_slice());
assert_eq!(out.value()?, &[b'M'; 3900]);
}
assert_eq!(hlog.tail_address(), tail_before);
let mut scanned = Vec::new();
hlog
.scan(0, tail_before, |addr, rec| {
scanned.push((addr, rec.key().to_vec()));
Ok(true)
})
.await?;
assert_eq!(
scanned,
addrs
.iter()
.map(|(a, k)| (*a, k.clone()))
.collect::<Vec<_>>()
);
let next = hlog.append(b"m5", &[b'N'; 100], 0, false)?;
assert_eq!(next, tail_before);
info!("跨多段恢复与环形窗口守卫测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}