use super::log_payload::LogPayloadStorage;
use super::traits::PayloadStorage;
use super::wal_cursor::{
watermark_allows_full_truncation, WalCursor, WalPosition, WalRecord, WalWatermarkRegistry,
};
use super::wal_cursor_reader::LogWalCursor;
use proptest::prelude::*;
use std::fs::OpenOptions;
use std::path::Path;
#[derive(Debug, Clone)]
enum WalOp {
Store { id: u64, payload_len: usize },
Delete { id: u64 },
}
fn wal_op_strategy() -> impl Strategy<Value = WalOp> {
prop_oneof![
(0u64..8, 0usize..24).prop_map(|(id, payload_len)| WalOp::Store { id, payload_len }),
(0u64..8).prop_map(|id| WalOp::Delete { id }),
]
}
fn build_wal(dir: &Path, ops: &[WalOp]) {
let mut storage = LogPayloadStorage::new(dir).expect("open log payload storage in temp dir");
for op in ops {
match op {
WalOp::Store { id, payload_len } => {
let payload = serde_json::json!({ "v": "x".repeat(*payload_len) });
storage.store(*id, &payload).expect("store payload");
}
WalOp::Delete { id } => {
storage.delete(*id).expect("delete payload");
}
}
}
storage.flush().expect("flush wal");
drop(storage);
}
fn drain_from_start(cursor: &LogWalCursor) -> Vec<WalRecord> {
let mut all = Vec::new();
let mut pos = WalPosition::START;
loop {
let batch = cursor.read_from(pos, 4).expect("read_from");
if batch.is_empty() {
break;
}
pos = batch[batch.len() - 1].next;
all.extend(batch);
}
all
}
fn reclaim_full(dir: &Path) {
let file = OpenOptions::new()
.write(true)
.open(dir.join("payloads.log"))
.expect("open wal for truncation");
file.set_len(0).expect("truncate wal to empty");
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(100))]
#[test]
fn prop_retention_never_discards_needed_position(
ops in prop::collection::vec(wal_op_strategy(), 1..30),
w_raw in 0u64..512,
) {
let dir = tempfile::tempdir().expect("tempdir");
build_wal(dir.path(), &ops);
let original = {
let cursor = LogWalCursor::new(dir.path());
drain_from_start(&cursor)
};
prop_assume!(!original.is_empty());
let tail = original.last().map_or(0, |r| r.next.offset());
let w = w_raw % (tail + 8);
let watermark = WalPosition::new(w);
let registry = WalWatermarkRegistry::new();
let consumer = registry.register();
registry.advance(consumer, watermark);
let min = registry.min_watermark();
prop_assert_eq!(min, Some(watermark));
let allows_full = watermark_allows_full_truncation(min, tail);
prop_assert_eq!(allows_full, w >= tail);
let needed: Vec<&WalRecord> = original
.iter()
.filter(|r| r.position.offset() >= w)
.collect();
if allows_full {
reclaim_full(dir.path());
}
let after = {
let cursor = LogWalCursor::new(dir.path());
drain_from_start(&cursor)
};
for record in &needed {
prop_assert!(
after.iter().any(|a| a == *record),
"record at position {} (>= W {}) was discarded",
record.position.offset(),
w
);
}
for record in &original {
let survived = after.iter().any(|a| a == record);
if !survived {
prop_assert!(
record.position.offset() < w,
"reclaimed record at position {} was not below W {}",
record.position.offset(),
w
);
}
}
}
#[test]
fn prop_no_consumer_matches_pre_cursor_baseline(
ops in prop::collection::vec(wal_op_strategy(), 1..30),
) {
let dir = tempfile::tempdir().expect("tempdir");
build_wal(dir.path(), &ops);
let original = {
let cursor = LogWalCursor::new(dir.path());
drain_from_start(&cursor)
};
prop_assume!(!original.is_empty());
let tail = original.last().map_or(0, |r| r.next.offset());
let registry = WalWatermarkRegistry::new();
prop_assert_eq!(registry.min_watermark(), None);
prop_assert!(watermark_allows_full_truncation(None, tail));
reclaim_full(dir.path());
let after = {
let cursor = LogWalCursor::new(dir.path());
drain_from_start(&cursor)
};
prop_assert!(after.is_empty());
prop_assert_eq!(LogWalCursor::new(dir.path()).tail_position(), WalPosition::START);
}
#[test]
fn prop_consumer_below_tail_holds_all(
ops in prop::collection::vec(wal_op_strategy(), 1..30),
) {
let dir = tempfile::tempdir().expect("tempdir");
build_wal(dir.path(), &ops);
let original = {
let cursor = LogWalCursor::new(dir.path());
drain_from_start(&cursor)
};
prop_assume!(!original.is_empty());
let tail = original.last().map_or(0, |r| r.next.offset());
let registry = WalWatermarkRegistry::new();
let consumer = registry.register();
let w = tail.saturating_sub(1);
registry.advance(consumer, WalPosition::new(w));
prop_assert!(!watermark_allows_full_truncation(registry.min_watermark(), tail));
let after = {
let cursor = LogWalCursor::new(dir.path());
drain_from_start(&cursor)
};
prop_assert_eq!(after.len(), original.len());
prop_assert_eq!(&after[..], &original[..]);
}
}