use std::sync::Arc;
use aok::{OK, Void};
use compio::runtime::Runtime;
use log::info;
use tempfile::tempdir;
use waof::{Error, RECORD_HEADER_LEN, RecordHeader, RingBuffer, WalConfig, WalLog};
use wdev::SegmentedDevice;
#[ctor::ctor(unsafe)]
fn _log_init() {
log_init::init();
}
#[test]
fn test_record_header() -> Void {
assert_eq!(RECORD_HEADER_LEN, 8);
let empty_hdr = RecordHeader::for_payload(&[]);
assert_eq!(empty_hdr.entry_len, 0);
assert_ne!(empty_hdr.crc32, 0);
assert!(!empty_hdr.is_zero());
empty_hdr.verify(&[])?;
assert!(matches!(
RecordHeader::new(0, 0).verify(&[]),
Err(Error::ChecksumMismatch { .. })
));
let payload = b"Hello, Wedb Wal Engine!";
let header = RecordHeader::for_payload(payload);
assert_eq!(header.entry_len, payload.len() as u32);
header.verify(payload)?;
let mut bytes = [0u8; 8];
header.encode(&mut bytes);
let decoded = RecordHeader::decode(&bytes)?;
assert_eq!(header, decoded);
let mut corrupted = *payload;
corrupted[0] ^= 0xFF;
assert!(matches!(
header.verify(&corrupted),
Err(Error::ChecksumMismatch { .. })
));
info!("RecordHeader 编解码与校验和测试通过");
OK
}
#[test]
fn test_record_header_corruption_and_boundary_checks() -> Void {
let payload = b"critical wal transaction payload data";
let header = RecordHeader::for_payload(payload);
assert_eq!(header.payload_len(), payload.len());
assert!(!header.is_zero());
assert!(RecordHeader::decode(&[0u8; 7]).is_err());
assert!(RecordHeader::decode(&[]).is_err());
assert!(matches!(
header.verify(&payload[..payload.len() - 1]),
Err(Error::InvalidRecordHeader)
));
let mut extended = payload.to_vec();
extended.push(0);
assert!(matches!(
header.verify(&extended),
Err(Error::InvalidRecordHeader)
));
let mut mutated = *payload;
for i in 0..payload.len() {
for bit in 0..8 {
mutated[i] ^= 1 << bit;
assert!(
matches!(header.verify(&mutated), Err(Error::ChecksumMismatch { .. })),
"第 {i} 字节第 {bit} 位翻转未被拦截"
);
mutated[i] ^= 1 << bit;
}
}
let zero_hdr = RecordHeader::new(0, 0);
assert!(zero_hdr.is_zero());
assert_eq!(zero_hdr.payload_len(), 0);
info!("RecordHeader 极端校验与位反转测试通过");
OK
}
#[test]
fn test_wal_ring_buffer_large_64bit_offset() -> Void {
let buffer_size = 64 * 1024; let align = 4096;
let ring = RingBuffer::new(buffer_size, align)?;
let huge_offset = (8u64 * 1024 * 1024 * 1024) + 123;
let test_data = b"64-bit large logical address test across 4GB boundary!";
ring.write_bytes(huge_offset, test_data);
let mut read_back = [0u8; 64];
let slice = &mut read_back[..test_data.len()];
ring.read_bytes(huge_offset, slice);
assert_eq!(slice, test_data);
let near_boundary_offset = (16u64 * 1024 * 1024 * 1024) + (buffer_size as u64 - 10);
let boundary_data = b"across_ring_boundary_large_address";
ring.write_bytes(near_boundary_offset, boundary_data);
let mut boundary_read = [0u8; 64];
let slice = &mut boundary_read[..boundary_data.len()];
ring.read_bytes(near_boundary_offset, slice);
assert_eq!(slice, boundary_data);
info!("WAL 环形写缓冲区 64 位大地址测试通过");
OK
}
#[test]
fn test_wal_end_to_end_smoke() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let db_path = dir.path().join("wal_smoke.log");
let seg_size = 16 * 1024; let config = WalConfig::new(64 * 1024);
let committed_tail;
let mut expected_payloads = Vec::with_capacity(30);
{
let device = Arc::new(SegmentedDevice::segmented(&db_path, seg_size)?);
let wal = WalLog::new(device, config)?;
for i in 0..30 {
let payload = format!("wal-smoke-record-{i:03}").into_bytes();
wal.enqueue(&payload)?;
expected_payloads.push(payload);
}
let mut mem_iter = wal.scan_all();
let mem_records = mem_iter.collect_all().await?;
assert_eq!(mem_records.len(), 30);
for (rec, expected) in mem_records.iter().zip(&expected_payloads) {
assert_eq!(&rec.payload, expected);
}
committed_tail = wal.commit().await?;
assert_eq!(committed_tail, wal.tail_address());
let mut disk_iter = wal.scan_committed();
let disk_records = disk_iter.collect_all().await?;
assert_eq!(disk_records.len(), 30);
}
{
let device = Arc::new(SegmentedDevice::segmented(&db_path, seg_size)?);
let wal = WalLog::open(device, config).await?;
assert_eq!(wal.tail_address(), committed_tail);
assert_eq!(wal.committed_until_address(), committed_tail);
let mut iter = wal.scan_committed();
let recovered_records = iter.collect_all().await?;
assert_eq!(recovered_records.len(), 30);
for (rec, expected) in recovered_records.iter().zip(&expected_payloads) {
assert_eq!(&rec.payload, expected);
rec.header.verify(&rec.payload)?;
}
let append_payload = b"wal-smoke-appended-after-recovery";
let append_addr = wal.enqueue(append_payload)?;
assert_eq!(append_addr, committed_tail);
let new_tail = wal.commit().await?;
assert!(new_tail > append_addr);
let mut post_iter = wal.scan(append_addr, new_tail);
let post_records = post_iter.collect_all().await?;
assert_eq!(post_records.len(), 1);
assert_eq!(post_records[0].payload, append_payload);
}
info!("WAL 端到端全链路冒烟测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}