use aok::{OK, Void};
use compio::runtime::Runtime;
use log::info;
use tempfile::tempdir;
use wdev::{Device, Error, SegmentedDevice};
use wram::AlignedBuf;
use crate::support::make_pattern_data;
#[test]
fn sync_persists_data_and_handles_reopen_after_reset() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let seg_size: u64 = 64 * 1024;
let device = SegmentedDevice::segmented(dir.path().join("sync.log"), seg_size)?;
let pattern0 = make_pattern_data(4096, 5, 1);
let pattern1 = make_pattern_data(4096, 7, 2);
for (seg_id, pattern) in [(0u32, &pattern0), (1, &pattern1)] {
let buf = AlignedBuf::from_slice(pattern, 4096)?;
let (res, _) = device.write_aligned((seg_id as u64) * seg_size, buf).await;
assert_eq!(res?, 4096);
}
device.sync().await?;
assert!(device.is_segment_cached(0));
assert!(device.is_segment_cached(1));
for (seg_id, pattern) in [(0u32, &pattern0), (1, &pattern1)] {
let check = AlignedBuf::new(4096, 4096)?;
let (res, check) = device.read_aligned((seg_id as u64) * seg_size, check).await;
assert_eq!(res?, 4096);
assert_eq!(check.as_slice(), &pattern[..]);
}
device.reset();
assert!(device.is_cached_empty());
let reread = AlignedBuf::new(4096, 4096)?;
let (res, reread) = device.read_aligned(0, reread).await;
assert_eq!(res?, 4096);
assert_eq!(reread.as_slice(), &pattern0[..]);
assert!(device.is_segment_cached(0), "重开后句柄应重新入缓存");
info!("sync 刷盘持久化与句柄按需重开校验通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[cfg(unix)]
#[test]
fn idevice_permission_denied_at_first_write_callback_gets_error() -> Void {
use std::{
fs::{Permissions, set_permissions},
os::unix::fs::PermissionsExt,
};
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let device = SegmentedDevice::segmented(dir.path().join("perm_denied.log"), 64 * 1024)?;
let buf = AlignedBuf::from_slice(&[0xABu8; 4096], 4096)?;
set_permissions(dir.path(), Permissions::from_mode(0o0))?;
let (res, _) = device.write_aligned(0, buf).await;
set_permissions(dir.path(), Permissions::from_mode(0o755))?;
match res {
Err(Error::Io(_)) => info!("open() 权限拒绝已正确传播为 Io 错误"),
Ok(n) => info!("特权环境 (root) 下 chmod 不生效,写入成功 {n} 字节,跳过断言"),
other => panic!("预期 Io 权限错误,实际为: {other:?}"),
}
let buf2 = AlignedBuf::from_slice(&[0xCDu8; 4096], 4096)?;
let (res, _) = device.write_aligned(0, buf2).await;
assert_eq!(res?, 4096);
let check = AlignedBuf::new(4096, 4096)?;
let (res, check) = device.read_aligned(0, check).await;
assert_eq!(res?, 4096);
assert!(check.as_slice().iter().all(|&b| b == 0xCD));
info!("权限拒绝错误传播与恢复后可用性校验通过 (PermissionDeniedAtFirstWrite)");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn sync_data_persists_data_and_handles_reopen_after_reset() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let seg_size: u64 = 64 * 1024;
let device = SegmentedDevice::segmented(dir.path().join("sync_data.log"), seg_size)?;
let pattern = make_pattern_data(4096, 13, 7);
let buf = AlignedBuf::from_slice(&pattern, 4096)?;
let (res, _) = device.write_aligned(0, buf).await;
assert_eq!(res?, 4096);
device.sync_data().await?;
device.reset();
let check = AlignedBuf::new(4096, 4096)?;
let (res, check) = device.read_aligned(0, check).await;
assert_eq!(res?, 4096);
assert_eq!(check.as_slice(), &pattern[..]);
info!("sync_data 异步数据刷盘 (fdatasync) 校验通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn read_only_device_blocks_writes_while_allowing_reads() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let log_path = dir.path().join("readonly.log");
let seg_size: u64 = 64 * 1024;
let pattern = make_pattern_data(4096, 19, 2);
{
let device = SegmentedDevice::segmented(&log_path, seg_size)?;
let buf = AlignedBuf::from_slice(&pattern, 4096)?;
let (res, _) = device.write_aligned(0, buf).await;
assert_eq!(res?, 4096);
}
let mut ro_device = SegmentedDevice::segmented(&log_path, seg_size)?;
ro_device.set_read_only(true);
assert!(ro_device.is_read_only());
let check = AlignedBuf::new(4096, 4096)?;
let (res, check) = ro_device.read_aligned(0, check).await;
assert_eq!(res?, 4096);
assert_eq!(check.as_slice(), &pattern[..]);
let write_buf = AlignedBuf::from_slice(&[0xFFu8; 4096], 4096)?;
let (res, _) = ro_device.write_aligned(0, write_buf).await;
assert!(
matches!(res, Err(Error::ReadOnly { offset, len }) if offset == 0 && len == 4096),
"只读模式下的写入必须拒绝为 Error::ReadOnly,实际为: {res:?}"
);
info!("只读设备模式 (readOnly) 保护校验通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn preallocate_sets_segment_file_size() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let seg_size: u64 = 64 * 1024;
let mut device = SegmentedDevice::segmented(dir.path().join("prealloc.log"), seg_size)?;
device.set_preallocate(true);
assert!(device.is_preallocate());
let buf = AlignedBuf::from_slice(&[0x42u8; 4096], 4096)?;
let (res, _) = device.write_aligned(0, buf).await;
assert_eq!(res?, 4096);
assert_eq!(
device.get_file_size(0)?,
seg_size,
"预分配段物理大小应达到配置段大小"
);
info!("段文件预分配 (preallocateFile) 语义校验通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn delete_on_close_cleans_up_files_on_drop() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let log_path = dir.path().join("del_close.log");
let seg_size: u64 = 64 * 1024;
let seg0_path = {
let mut device = SegmentedDevice::segmented(&log_path, seg_size)?;
device.set_delete_on_close(true);
assert!(device.is_delete_on_close());
let buf = AlignedBuf::from_slice(&[0x33u8; 4096], 4096)?;
let (res, _) = device.write_aligned(0, buf).await;
assert_eq!(res?, 4096);
let seg0_path = device.segment_path(0);
assert!(seg0_path.exists(), "析构前文件应存在");
seg0_path
};
assert!(
!seg0_path.exists(),
"delete_on_close 析构后段 0 必须已被物理清理"
);
info!("delete_on_close 自动清理段文件校验通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn storage_device_trait_dispatch() -> Void {
use wdev::StorageDevice;
async fn run_device_ops<D: StorageDevice>(device: &D, seg_size: u64) -> Result<(), wdev::Error> {
assert_eq!(device.sector_size(), 4096);
assert_eq!(device.segment_size(), Some(seg_size));
assert_eq!(device.start_segment(), 0);
let pattern = make_pattern_data(4096, 7, 1);
let buf = AlignedBuf::from_slice(&pattern, 4096)?;
let (res, _) = device.write_aligned(0, buf).await;
assert_eq!(res?, 4096);
assert!(device.get_file_size(0)? >= 4096);
device.sync_data().await?;
device.reset();
device.truncate_until_segment(1).await?;
assert_eq!(device.start_segment(), 1);
Ok(())
}
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let seg_size: u64 = 64 * 1024;
let device = SegmentedDevice::segmented(dir.path().join("dyn_dev.log"), seg_size)?;
run_device_ops(&device, seg_size).await?;
info!("StorageDevice trait 泛型静态分发校验通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn new_segment_creation_fsyncs_parent_dir() -> Void {
use std::sync::atomic::Ordering;
let rt = Runtime::new()?;
rt.block_on(async {
let dir = tempdir()?;
let seg_size: u64 = 64 * 1024;
let device = SegmentedDevice::segmented(dir.path().join("dirsync.log"), seg_size)?;
assert_eq!(device.dir_syncs.load(Ordering::Relaxed), 0);
let per_new_segment = cfg!(unix) as u64;
let mut dir_syncs_total = 0;
let pattern = make_pattern_data(4096, 3, 9);
let buf = AlignedBuf::from_slice(&pattern, 4096)?;
let (res, _) = device.write_aligned(0, buf).await;
assert_eq!(res?, 4096);
dir_syncs_total += per_new_segment;
assert_eq!(device.dir_syncs.load(Ordering::Relaxed), dir_syncs_total);
let buf = AlignedBuf::from_slice(&pattern, 4096)?;
let (res, _) = device.write_aligned(0, buf).await;
assert_eq!(res?, 4096);
assert_eq!(device.dir_syncs.load(Ordering::Relaxed), dir_syncs_total);
let buf = AlignedBuf::from_slice(&pattern, 4096)?;
let (res, _) = device.write_aligned(seg_size, buf).await;
assert_eq!(res?, 4096);
dir_syncs_total += per_new_segment;
assert_eq!(device.dir_syncs.load(Ordering::Relaxed), dir_syncs_total);
device.sync().await?;
device.reset();
let check = AlignedBuf::new(4096, 4096)?;
let (res, check) = device.read_aligned(0, check).await;
assert_eq!(res?, 4096);
assert_eq!(check.as_slice(), &pattern[..]);
assert_eq!(device.dir_syncs.load(Ordering::Relaxed), dir_syncs_total);
drop(device);
let recovered = SegmentedDevice::segmented(dir.path().join("dirsync.log"), seg_size)?;
recovered.recover()?;
assert_eq!(recovered.start_segment(), 0);
assert_eq!(recovered.end_segment(), Some(1));
info!("新建段父目录 fsync 契约验证通过");
aok::Result::<()>::Ok(())
})?;
OK
}