use std::sync::atomic::Ordering;
use aok::{OK, Void};
use compio::runtime::Runtime;
use log::info;
use waof::{
COMMIT_FRAME_TOTAL_LEN, Error, NO_COOKIE, RECORD_HEADER_LEN, WalLog, WalScanIterator,
is_commit_frame,
};
use wbase::align::sector_bounds;
use wdev::{Device, SegmentedDevice};
use super::support::{self, WalFixture, make_pattern_payload, reopen_single_file};
const BUF: usize = 16 * 1024;
const BLOCK: usize = 2 * 1024;
type Wal = WalLog<SegmentedDevice>;
fn ring_upper(wal: &Wal) -> u64 {
let sector = wal.device().sector_size() as u64;
let flushed = wal.flushed_until_address();
sector_bounds(flushed, flushed, sector).0 + wal.config().buffer_size as u64
}
fn fill_window(wal: &Wal) -> aok::Result<Vec<usize>> {
let mut lens = Vec::new();
loop {
let room = ring_upper(wal) - wal.tail_address();
if room <= RECORD_HEADER_LEN as u64 {
break;
}
let len = (room - RECORD_HEADER_LEN as u64).min(BLOCK as u64) as usize;
wal.enqueue(&make_pattern_payload(lens.len(), len))?;
lens.push(len);
}
assert_eq!(
wal.tail_address(),
ring_upper(wal),
"前置:窗口须灌满至上界"
);
let frame_payload_len = (COMMIT_FRAME_TOTAL_LEN - RECORD_HEADER_LEN as u64) as usize;
assert!(
matches!(
wal.enqueue(&make_pattern_payload(lens.len(), frame_payload_len)),
Err(Error::BufferFull { .. })
),
"前置:帧级预留须已被拒(否则本用例并未覆盖跳帧臂)"
);
Ok(lens)
}
async fn count_frames<D: Device>(iter: WalScanIterator<D>) -> aok::Result<usize> {
let all = support::collect_iter(iter).await?;
Ok(
all
.iter()
.filter(|rec| is_commit_frame(&rec.payload))
.count(),
)
}
#[test]
fn test_commit_frame_skip_flushes_batch_and_holds_cursor() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let fixture = WalFixture::single_file("commit_frame_skip.log", BUF)?;
let wal = fixture.wal;
wal.enqueue(b"baseline-record-before-skip")?;
let base_committed = wal.commit().await?;
assert_eq!(
wal.last_commit_frame.load(Ordering::Acquire),
base_committed
);
assert_eq!(count_frames(wal.scan(0, base_committed)).await?, 1);
let lens = fill_window(&wal)?;
let batch_tail = wal.tail_address();
assert!(batch_tail > base_committed);
let committed = wal.commit().await?;
assert_eq!(committed, batch_tail);
assert_eq!(wal.committed_until_address(), batch_tail);
assert_eq!(wal.flushed_until_address(), batch_tail);
assert_eq!(
wal.last_commit_frame.load(Ordering::Acquire),
base_committed
);
assert_eq!(count_frames(wal.scan(0, batch_tail)).await?, 1);
let records = support::collect_data(wal.scan_committed()).await?;
assert_eq!(records.len(), lens.len() + 1);
assert_eq!(records[0].payload, b"baseline-record-before-skip");
for (i, rec) in records.iter().skip(1).enumerate() {
rec.header.verify(&rec.payload)?;
assert_eq!(rec.payload, make_pattern_payload(i, lens[i]));
}
info!("commit 帧跳帧降级:本批照常刷盘与帧游标不推进测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_commit_frame_skip_next_round_backfills_frame() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let fixture = WalFixture::single_file("commit_frame_backfill.log", BUF)?;
let wal = fixture.wal;
let lens = fill_window(&wal)?;
let skip_tail = wal.tail_address();
assert_eq!(wal.commit().await?, skip_tail);
assert_eq!(count_frames(wal.scan(0, skip_tail)).await?, 0);
assert_eq!(wal.last_commit_frame.load(Ordering::Acquire), 0);
wal.set_pending_cookie(11);
wal.enqueue(&make_pattern_payload(0, 100))?;
let data_end = wal.tail_address();
let committed = wal.commit().await?;
assert_eq!(committed, data_end + COMMIT_FRAME_TOTAL_LEN);
assert_eq!(wal.tail_address(), committed);
assert_eq!(wal.last_commit_frame.load(Ordering::Acquire), committed);
assert!(committed > skip_tail);
assert_eq!(count_frames(wal.scan(skip_tail, committed)).await?, 1);
let records = support::collect_data(wal.scan_committed()).await?;
assert_eq!(records.len(), lens.len() + 1);
for (i, rec) in records.iter().take(lens.len()).enumerate() {
assert_eq!(rec.payload, make_pattern_payload(i, lens[i]));
}
assert_eq!(records[lens.len()].payload, make_pattern_payload(0, 100));
info!("commit 帧跳帧降级:下一轮 commit 补写帧与游标覆盖测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_commit_frame_skip_recovery_falls_back_to_previous_frame() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let fixture = WalFixture::single_file("commit_frame_skip_recover.log", BUF)?;
let dir = fixture.dir;
let wal = fixture.wal;
let file_name = "commit_frame_skip_recover.log";
wal.enqueue(b"pre-skip-committed")?;
wal.set_pending_cookie(7);
let prev_frame_end = wal.commit().await?;
let lens = fill_window(&wal)?;
let skip_tail = wal.tail_address();
assert_eq!(wal.commit().await?, skip_tail);
drop(wal);
let wal = reopen_single_file(dir.path(), file_name, BUF).await?;
assert_eq!(wal.tail_address(), skip_tail);
assert_eq!(wal.flushed_until_address(), skip_tail);
assert_eq!(wal.committed_until_address(), prev_frame_end);
assert_eq!(
wal.last_commit_frame.load(Ordering::Acquire),
prev_frame_end
);
assert_eq!(wal.recovered_cookie(), 7);
let replayed = support::collect_data(wal.scan_committed()).await?;
assert_eq!(replayed.len(), 1);
assert_eq!(replayed[0].payload, b"pre-skip-committed");
let flushed_records = support::collect_data(wal.scan(0, skip_tail)).await?;
assert_eq!(flushed_records.len(), lens.len() + 1);
for (i, rec) in flushed_records.iter().skip(1).enumerate() {
assert_eq!(rec.payload, make_pattern_payload(i, lens[i]));
}
wal.set_pending_cookie(9);
wal.enqueue(&make_pattern_payload(0, 64))?;
let final_committed = wal.commit().await?;
assert_eq!(final_committed, wal.tail_address());
assert_eq!(
wal.last_commit_frame.load(Ordering::Acquire),
final_committed
);
assert_eq!(
count_frames(wal.scan(prev_frame_end, final_committed)).await?,
1
);
assert_eq!(
support::collect_data(wal.scan_committed()).await?.len(),
lens.len() + 2
);
drop(wal);
let wal = reopen_single_file(dir.path(), file_name, BUF).await?;
assert_eq!(wal.committed_until_address(), final_committed);
assert_eq!(wal.recovered_cookie(), 9);
let records = support::collect_data(wal.scan_committed()).await?;
assert_eq!(records.len(), lens.len() + 2);
assert_eq!(records[0].payload, b"pre-skip-committed");
for (i, rec) in records.iter().skip(1).take(lens.len()).enumerate() {
rec.header.verify(&rec.payload)?;
assert_eq!(rec.payload, make_pattern_payload(i, lens[i]));
}
assert_eq!(records[lens.len() + 1].payload, make_pattern_payload(0, 64));
info!("commit 帧跳帧降级:恢复侧回退上一有效帧与补写收敛测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}
#[test]
fn test_commit_frame_skip_first_round_without_frame_recovers_flushed_tail() -> Void {
let rt = Runtime::new()?;
rt.block_on(async {
let fixture = WalFixture::single_file("commit_frame_skip_no_frame.log", BUF)?;
let dir = fixture.dir;
let wal = fixture.wal;
let lens = fill_window(&wal)?;
let skip_tail = wal.tail_address();
assert_eq!(wal.commit().await?, skip_tail);
assert_eq!(wal.last_commit_frame.load(Ordering::Acquire), 0);
drop(wal);
let file_name = "commit_frame_skip_no_frame.log";
let wal = reopen_single_file(dir.path(), file_name, BUF).await?;
assert_eq!(wal.tail_address(), skip_tail);
assert_eq!(wal.committed_until_address(), skip_tail);
assert_eq!(wal.last_commit_frame.load(Ordering::Acquire), skip_tail);
assert_eq!(wal.recovered_cookie(), NO_COOKIE);
let records = support::collect_data(wal.scan_committed()).await?;
assert_eq!(records.len(), lens.len());
for (i, rec) in records.iter().enumerate() {
rec.header.verify(&rec.payload)?;
assert_eq!(rec.payload, make_pattern_payload(i, lens[i]));
}
info!("commit 帧跳帧降级:无帧可退时按已刷盘尾部收敛测试通过");
aok::Result::<()>::Ok(())
})?;
OK
}