use std::{
fs,
sync::{Arc, mpsc},
thread,
};
use aok::{OK, Void};
use compio::{buf::BufResult, fs::File, io::AsyncWriteAt, runtime::Runtime};
use log::info;
use tempfile::tempdir;
use wdev::{Device, SegmentedDevice};
use wram::AlignedBuf;
use crate::support::Watchdog;
const SECTOR: usize = 4096;
const SEG_SIZE: u64 = 64 * 1024;
#[test]
fn raw_cross_thread_fsync_feasibility() -> Void {
let dir = tempdir()?;
let path = dir.path().join("raw_ct_fsync.bin");
let _wd = Watchdog::start(60);
let path_a = path.clone();
let (tx, rx) = mpsc::channel::<Arc<File>>();
let ta = thread::spawn(move || -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let file = Arc::new(File::create(&path_a).await?);
let data: Vec<u8> = (0..SECTOR).map(|j| (j & 0xFF) as u8).collect();
let buf = AlignedBuf::from_slice(&data, 4096)?;
let BufResult(res, _) = (&*file).write_at(buf, 0).await;
assert_eq!(res?, SECTOR);
tx.send(file).expect("接收端存活");
aok::Result::<()>::Ok(())
})?;
OK
});
let tb = thread::spawn(move || -> Void {
let file = rx.recv().expect("发送端存活");
let rt = Runtime::new()?;
rt.block_on(async {
file.sync_all().await?;
file.sync_data().await?;
aok::Result::<()>::Ok(())
})?;
OK
});
ta.join().unwrap()?;
tb.join().unwrap()?;
let disk = fs::read(&path)?;
assert_eq!(disk.len(), SECTOR);
assert_eq!(disk[0], 0);
assert_eq!(disk[SECTOR - 1], ((SECTOR - 1) & 0xFF) as u8);
info!("跨线程 fsync 可行性实验通过:他线程 fd 的 sync_all/sync_data 全部成功");
OK
}
#[test]
fn global_sync_covers_foreign_thread_writes() -> Void {
let dir = tempdir()?;
let path = dir.path().join("global_sync.log");
let device = Arc::new(SegmentedDevice::segmented(&path, SEG_SIZE)?);
let _wd = Watchdog::start(60);
let dev_a = Arc::clone(&device);
let ta = thread::spawn(move || -> Void {
let rt = Runtime::new()?;
rt.block_on(async move {
for seg in 0u32..2 {
let data: Vec<u8> = (0..SECTOR).map(|j| (j ^ seg as usize) as u8).collect();
let wbuf = AlignedBuf::from_slice(&data, 4096)?;
let offset = u64::from(seg) * SEG_SIZE;
let (res, _) = dev_a.write_aligned(offset, wbuf).await;
assert_eq!(res?, SECTOR);
}
aok::Result::<()>::Ok(())
})?;
OK
});
ta.join().unwrap()?;
#[cfg(debug_assertions)]
assert_eq!(device.debug_dirty_segments(), vec![0, 1]);
let dev_b = Arc::clone(&device);
let tb = thread::spawn(move || -> Void {
let rt = Runtime::new()?;
rt.block_on(async move {
dev_b.sync().await?;
let fresh = SegmentedDevice::segmented(&path, SEG_SIZE)?;
for seg in 0u32..2 {
let expected: Vec<u8> = (0..SECTOR).map(|j| (j ^ seg as usize) as u8).collect();
let check = AlignedBuf::new(SECTOR, 4096)?;
let (res, check) = fresh.read_aligned(u64::from(seg) * SEG_SIZE, check).await;
assert_eq!(res?, SECTOR);
assert_eq!(check.as_slice(), &expected[..], "段 {seg} 落盘字节不匹配");
}
aok::Result::<()>::Ok(())
})?;
OK
});
tb.join().unwrap()?;
#[cfg(debug_assertions)]
assert!(
device.debug_dirty_segments().is_empty(),
"全局 sync 后仍有在册脏段"
);
info!("全局 sync 覆盖他线程写入并通过全新句柄读回复验");
OK
}
#[test]
fn global_sync_data_covers_foreign_thread_writes() -> Void {
let dir = tempdir()?;
let path = dir.path().join("global_sync_data.log");
let device = Arc::new(SegmentedDevice::segmented(&path, SEG_SIZE)?);
let _wd = Watchdog::start(60);
let dev_a = Arc::clone(&device);
let ta = thread::spawn(move || -> Void {
let rt = Runtime::new()?;
rt.block_on(async move {
let data: Vec<u8> = (0..SECTOR).map(|j| (j | 0x80) as u8).collect();
let wbuf = AlignedBuf::from_slice(&data, 4096)?;
let (res, _) = dev_a.write_aligned(0, wbuf).await;
assert_eq!(res?, SECTOR);
aok::Result::<()>::Ok(())
})?;
OK
});
ta.join().unwrap()?;
let dev_b = Arc::clone(&device);
let tb = thread::spawn(move || -> Void {
let rt = Runtime::new()?;
rt.block_on(async move { dev_b.sync_data().await })?;
OK
});
tb.join().unwrap()?;
#[cfg(debug_assertions)]
assert!(device.debug_dirty_segments().is_empty());
info!("全局 sync_data 与 sync 同覆盖口径");
OK
}
#[test]
fn sync_contract_remove_and_truncate_immunity() -> Void {
let dir = tempdir()?;
let path = dir.path().join("immunity.log");
let device = Arc::new(SegmentedDevice::segmented(&path, SEG_SIZE)?);
let _wd = Watchdog::start(60);
let rt = Runtime::new()?;
rt.block_on(async {
for seg in 0u32..2 {
let wbuf = AlignedBuf::from_slice(&vec![0xA5; SECTOR], 4096)?;
let (res, _) = device.write_aligned(u64::from(seg) * SEG_SIZE, wbuf).await;
assert_eq!(res?, SECTOR);
}
#[cfg(debug_assertions)]
assert_eq!(device.debug_dirty_segments(), vec![0, 1]);
device.remove_segment(0).await?;
#[cfg(debug_assertions)]
assert_eq!(device.debug_dirty_segments(), vec![1]);
device.truncate_until_segment(1).await?;
#[cfg(debug_assertions)]
assert_eq!(device.debug_dirty_segments(), vec![1]);
device.sync().await?;
assert_eq!(device.start_segment(), 1);
assert_eq!(device.get_file_size(0)?, 0);
aok::Result::<()>::Ok(())
})?;
info!("删段与截断免责路径通过,sync 无守护违约");
OK
}
#[test]
fn multi_segment_concurrent_sync() -> Void {
let dir = tempdir()?;
let path = dir.path().join("multi_seg_concurrent_sync.log");
let device = Arc::new(SegmentedDevice::segmented(&path, SEG_SIZE)?);
let _wd = Watchdog::start(60);
let rt = Runtime::new()?;
rt.block_on(async {
device.sync().await?;
device.sync_data().await?;
let data0 = vec![0x42u8; SECTOR];
let wbuf0 = AlignedBuf::from_slice(&data0, 4096)?;
let (res, _) = device.write_aligned(0, wbuf0).await;
assert_eq!(res?, SECTOR);
device.sync().await?;
#[cfg(debug_assertions)]
assert!(device.debug_dirty_segments().is_empty());
const SEGS: u32 = 8;
for seg in 0..SEGS {
let data = vec![(seg & 0xFF) as u8; SECTOR];
let wbuf = AlignedBuf::from_slice(&data, 4096)?;
let (res, _) = device.write_aligned(u64::from(seg) * SEG_SIZE, wbuf).await;
assert_eq!(res?, SECTOR);
}
#[cfg(debug_assertions)]
assert_eq!(device.debug_dirty_segments().len(), SEGS as usize);
device.sync_data().await?;
device.sync().await?;
#[cfg(debug_assertions)]
assert!(device.debug_dirty_segments().is_empty());
for seg in 0..SEGS {
let check = AlignedBuf::new(SECTOR, 4096)?;
let (res, check) = device.read_aligned(u64::from(seg) * SEG_SIZE, check).await;
assert_eq!(res?, SECTOR);
let expected_byte = (seg & 0xFF) as u8;
assert!(
check.as_slice().iter().all(|&b| b == expected_byte),
"段 {seg} 落盘数据不匹配"
);
}
aok::Result::<()>::Ok(())
})?;
info!("多段并发 sync / sync_data 压榨测试通过");
OK
}