use std::collections::Bound;
use reifydb_cdc::storage::{CdcStorage, DropBeforeResult, memory::MemoryCdcStorage, sqlite::storage::SqliteCdcStorage};
use reifydb_codec::{key::encoded::EncodedKey, row::bytes::EncodedBytes};
use reifydb_core::{
common::CommitVersion,
interface::cdc::{Cdc, SystemChange},
};
use reifydb_sqlite::SqliteConfig;
use reifydb_value::{byte_size::ByteSize, count::Count, util::cowvec::CowVec, value::datetime::DateTime};
fn cdc_minimal(version: u64) -> Cdc {
Cdc::new(
CommitVersion(version),
DateTime::from_nanos(1_700_000_000_000_000_000),
Vec::new(),
vec![SystemChange::Insert {
key: EncodedKey::new(vec![1, 2, 3]),
post: EncodedBytes(CowVec::new(vec![10, 20, 30])),
}],
)
}
fn assert_write_read_round_trip<S: CdcStorage>(storage: S) {
let cdc = cdc_minimal(1);
storage.write(&cdc).unwrap();
let read = storage.read(CommitVersion(1)).unwrap().expect("entry should exist");
assert_eq!(read.version, CommitVersion(1));
assert_eq!(read.system_changes.len(), 1);
}
fn assert_read_nonexistent<S: CdcStorage>(storage: S) {
assert!(storage.read(CommitVersion(999)).unwrap().is_none());
}
fn assert_range_inclusive<S: CdcStorage>(storage: S) {
for v in 1..=10 {
storage.write(&cdc_minimal(v)).unwrap();
}
let batch =
storage.read_range(Bound::Included(CommitVersion(3)), Bound::Included(CommitVersion(7)), 100).unwrap();
assert_eq!(batch.items.len(), 5);
assert!(!batch.has_more);
assert_eq!(batch.items[0].version, CommitVersion(3));
assert_eq!(batch.items[4].version, CommitVersion(7));
}
fn assert_range_exclusive<S: CdcStorage>(storage: S) {
for v in 1..=5 {
storage.write(&cdc_minimal(v)).unwrap();
}
let batch =
storage.read_range(Bound::Excluded(CommitVersion(2)), Bound::Included(CommitVersion(4)), 100).unwrap();
assert_eq!(batch.items.len(), 2);
assert_eq!(batch.items[0].version, CommitVersion(3));
assert_eq!(batch.items[1].version, CommitVersion(4));
let batch =
storage.read_range(Bound::Included(CommitVersion(2)), Bound::Excluded(CommitVersion(4)), 100).unwrap();
assert_eq!(batch.items.len(), 2);
assert_eq!(batch.items[0].version, CommitVersion(2));
assert_eq!(batch.items[1].version, CommitVersion(3));
}
fn assert_range_batch_size_has_more<S: CdcStorage>(storage: S) {
for v in 1..=10 {
storage.write(&cdc_minimal(v)).unwrap();
}
let batch = storage.read_range(Bound::Unbounded, Bound::Unbounded, 3).unwrap();
assert_eq!(batch.items.len(), 3);
assert!(batch.has_more);
}
fn assert_count<S: CdcStorage>(storage: S) {
let cdc = Cdc::new(
CommitVersion(1),
DateTime::from_nanos(1),
Vec::new(),
(0..5).map(|i| SystemChange::Insert {
key: EncodedKey::new(vec![i as u8]),
post: EncodedBytes(CowVec::new(vec![])),
})
.collect(),
);
storage.write(&cdc).unwrap();
assert_eq!(storage.count(CommitVersion(1)).unwrap(), 5);
assert_eq!(storage.count(CommitVersion(2)).unwrap(), 0);
}
fn assert_min_max_version<S: CdcStorage>(storage: S) {
assert!(storage.min_version().unwrap().is_none());
assert!(storage.max_version().unwrap().is_none());
storage.write(&cdc_minimal(5)).unwrap();
storage.write(&cdc_minimal(3)).unwrap();
storage.write(&cdc_minimal(8)).unwrap();
assert_eq!(storage.min_version().unwrap(), Some(CommitVersion(3)));
assert_eq!(storage.max_version().unwrap(), Some(CommitVersion(8)));
}
fn assert_overwrite<S: CdcStorage>(storage: S) {
let cdc1 = Cdc::new(
CommitVersion(1),
DateTime::from_nanos(100),
Vec::new(),
vec![SystemChange::Insert {
key: EncodedKey::new(vec![1]),
post: EncodedBytes(CowVec::new(vec![])),
}],
);
let cdc2 = Cdc::new(
CommitVersion(1),
DateTime::from_nanos(200),
Vec::new(),
vec![
SystemChange::Insert {
key: EncodedKey::new(vec![2]),
post: EncodedBytes(CowVec::new(vec![])),
},
SystemChange::Insert {
key: EncodedKey::new(vec![3]),
post: EncodedBytes(CowVec::new(vec![])),
},
],
);
storage.write(&cdc1).unwrap();
assert_eq!(storage.count(CommitVersion(1)).unwrap(), 1);
storage.write(&cdc2).unwrap();
assert_eq!(storage.count(CommitVersion(1)).unwrap(), 2);
let read = storage.read(CommitVersion(1)).unwrap().unwrap();
assert_eq!(read.timestamp, DateTime::from_nanos(200));
}
fn assert_drop_before_empty<S: CdcStorage>(storage: S) {
let r = storage.drop_before(CommitVersion(10), usize::MAX).unwrap();
assert_eq!(r.count, Count::ZERO);
assert!(r.entries.is_empty());
}
fn assert_drop_before_some<S: CdcStorage>(storage: S) {
for v in [1u64, 3, 5, 7, 9] {
storage.write(&cdc_minimal(v)).unwrap();
}
let r = storage.drop_before(CommitVersion(5), usize::MAX).unwrap();
assert_eq!(r.count, Count::new(2));
assert_eq!(r.entries.len(), 1);
assert_eq!(r.entries[0].count, Count::new(2));
assert!(storage.read(CommitVersion(1)).unwrap().is_none());
assert!(storage.read(CommitVersion(3)).unwrap().is_none());
assert!(storage.read(CommitVersion(5)).unwrap().is_some());
assert_eq!(storage.min_version().unwrap(), Some(CommitVersion(5)));
}
fn assert_drop_before_all<S: CdcStorage>(storage: S) {
for v in 1..=3u64 {
storage.write(&cdc_minimal(v)).unwrap();
}
let r = storage.drop_before(CommitVersion(10), usize::MAX).unwrap();
assert_eq!(r.count, Count::new(3));
assert!(storage.min_version().unwrap().is_none());
}
fn assert_drop_before_none_when_too_low<S: CdcStorage>(storage: S) {
for v in 5..=7u64 {
storage.write(&cdc_minimal(v)).unwrap();
}
let r = storage.drop_before(CommitVersion(3), usize::MAX).unwrap();
assert_eq!(r.count, Count::ZERO);
assert!(r.entries.is_empty());
assert_eq!(storage.min_version().unwrap(), Some(CommitVersion(5)));
}
fn assert_drop_before_boundary<S: CdcStorage>(storage: S) {
for v in 1..=5u64 {
storage.write(&cdc_minimal(v)).unwrap();
}
let r = storage.drop_before(CommitVersion(3), usize::MAX).unwrap();
assert_eq!(r.count, Count::new(2));
assert!(storage.read(CommitVersion(3)).unwrap().is_some());
assert_eq!(storage.min_version().unwrap(), Some(CommitVersion(3)));
}
fn assert_drop_before_entry_stats<S: CdcStorage>(storage: S) {
let cdc = Cdc::new(
CommitVersion(1),
DateTime::from_nanos(12345),
Vec::new(),
vec![SystemChange::Insert {
key: EncodedKey::new(vec![1, 2, 3]),
post: EncodedBytes(CowVec::new(vec![10, 20, 30, 40, 50])),
}],
);
storage.write(&cdc).unwrap();
let r: DropBeforeResult = storage.drop_before(CommitVersion(2), usize::MAX).unwrap();
assert_eq!(r.count, Count::new(1));
assert_eq!(r.entries.len(), 1);
assert_eq!(r.entries[0].key_bytes, ByteSize::from_bytes(3));
assert_eq!(r.entries[0].value_bytes, ByteSize::from_bytes(5));
assert_eq!(r.entries[0].count, Count::new(1));
}
fn assert_drop_before_limited<S: CdcStorage>(storage: S) {
for v in 1..=12u64 {
storage.write(&cdc_minimal(v)).unwrap();
}
let cutoff = CommitVersion(11);
let first = storage.drop_before(cutoff, 4).unwrap();
assert_eq!(first.count, Count::new(4));
assert!(first.more_remaining);
assert_eq!(first.entries.len(), 1);
assert_eq!(first.entries[0].count, Count::new(4));
assert_eq!(storage.min_version().unwrap(), Some(CommitVersion(5)));
assert!(storage.read(CommitVersion(11)).unwrap().is_some());
assert!(storage.read(CommitVersion(12)).unwrap().is_some());
let mut total = first.count;
loop {
let r = storage.drop_before(cutoff, 4).unwrap();
total = total.saturating_add(r.count);
if !r.more_remaining {
break;
}
}
assert_eq!(total, Count::new(10));
assert!(storage.read(CommitVersion(10)).unwrap().is_none());
assert!(storage.read(CommitVersion(11)).unwrap().is_some());
assert_eq!(storage.min_version().unwrap(), Some(CommitVersion(11)));
}
fn assert_range_inverted_returns_empty<S: CdcStorage>(storage: S) {
for v in 1..=5 {
storage.write(&cdc_minimal(v)).unwrap();
}
let batch = storage
.read_range(Bound::Included(CommitVersion(10)), Bound::Included(CommitVersion(5)), 16)
.expect("inverted range must not error");
assert!(batch.items.is_empty(), "inverted range must return no items");
assert!(!batch.has_more, "inverted range cannot have more items");
}
fn assert_range_excluded_zero_end_returns_empty<S: CdcStorage>(storage: S) {
for v in 1..=3 {
storage.write(&cdc_minimal(v)).unwrap();
}
let batch = storage
.read_range(Bound::Unbounded, Bound::Excluded(CommitVersion(0)), 16)
.expect("Excluded(0) end must not panic");
assert!(batch.items.is_empty());
assert!(!batch.has_more);
}
fn assert_range_excluded_pair_collapsing<S: CdcStorage>(storage: S) {
for v in 1..=10 {
storage.write(&cdc_minimal(v)).unwrap();
}
let batch = storage
.read_range(Bound::Excluded(CommitVersion(5)), Bound::Excluded(CommitVersion(6)), 16)
.expect("collapsing exclusive bounds must not panic");
assert!(batch.items.is_empty());
assert!(!batch.has_more);
}
macro_rules! storage_trait_tests {
($mod_name:ident, $fresh:expr) => {
mod $mod_name {
use super::*;
#[test]
fn write_read_round_trip() {
let (storage, _guard) = $fresh();
super::assert_write_read_round_trip(storage);
}
#[test]
fn read_nonexistent() {
let (storage, _guard) = $fresh();
super::assert_read_nonexistent(storage);
}
#[test]
fn range_inclusive() {
let (storage, _guard) = $fresh();
super::assert_range_inclusive(storage);
}
#[test]
fn range_exclusive() {
let (storage, _guard) = $fresh();
super::assert_range_exclusive(storage);
}
#[test]
fn range_batch_size_has_more() {
let (storage, _guard) = $fresh();
super::assert_range_batch_size_has_more(storage);
}
#[test]
fn range_inverted_returns_empty() {
let (storage, _guard) = $fresh();
super::assert_range_inverted_returns_empty(storage);
}
#[test]
fn range_excluded_zero_end_returns_empty() {
let (storage, _guard) = $fresh();
super::assert_range_excluded_zero_end_returns_empty(storage);
}
#[test]
fn range_excluded_pair_collapsing() {
let (storage, _guard) = $fresh();
super::assert_range_excluded_pair_collapsing(storage);
}
#[test]
fn count() {
let (storage, _guard) = $fresh();
super::assert_count(storage);
}
#[test]
fn min_max_version() {
let (storage, _guard) = $fresh();
super::assert_min_max_version(storage);
}
#[test]
fn overwrite_entry() {
let (storage, _guard) = $fresh();
super::assert_overwrite(storage);
}
#[test]
fn drop_before_empty() {
let (storage, _guard) = $fresh();
super::assert_drop_before_empty(storage);
}
#[test]
fn drop_before_some() {
let (storage, _guard) = $fresh();
super::assert_drop_before_some(storage);
}
#[test]
fn drop_before_all() {
let (storage, _guard) = $fresh();
super::assert_drop_before_all(storage);
}
#[test]
fn drop_before_none_when_too_low() {
let (storage, _guard) = $fresh();
super::assert_drop_before_none_when_too_low(storage);
}
#[test]
fn drop_before_boundary() {
let (storage, _guard) = $fresh();
super::assert_drop_before_boundary(storage);
}
#[test]
fn drop_before_entry_stats() {
let (storage, _guard) = $fresh();
super::assert_drop_before_entry_stats(storage);
}
#[test]
fn drop_before_limited() {
let (storage, _guard) = $fresh();
super::assert_drop_before_limited(storage);
}
}
};
}
storage_trait_tests!(memory, || (MemoryCdcStorage::new(), ()));
storage_trait_tests!(sqlite, || {
let (config, guard) = SqliteConfig::test();
(SqliteCdcStorage::new(config), guard)
});
#[test]
fn sqlite_block_eviction_aggregates_from_stored_rollup() {
let (config, _guard) = SqliteConfig::test();
let storage = SqliteCdcStorage::new(config);
for v in 1..=20u64 {
storage.write(&cdc_minimal(v)).unwrap();
}
let blocks = storage.compact_all(1, 3, CommitVersion(1000)).unwrap();
assert_eq!(blocks.len(), 20, "each entry should compact into its own block");
assert_eq!(storage.min_version().unwrap(), Some(CommitVersion(1)), "all data now lives in blocks");
let r = storage.drop_before(CommitVersion(21), usize::MAX).unwrap();
assert_eq!(r.count, Count::new(20), "every block entry is accounted for");
assert_eq!(r.entries.len(), 1, "all 20 entries share one source, so they roll up into one entry");
let e = &r.entries[0];
assert_eq!(e.count, Count::new(20));
assert_eq!(e.key_bytes, ByteSize::from_bytes(3 * 20), "20 * 3-byte keys");
assert_eq!(e.value_bytes, ByteSize::from_bytes(3 * 20), "20 * 3-byte values");
assert_eq!(storage.min_version().unwrap(), None, "all blocks below the cutoff must be gone");
}
#[test]
fn sqlite_straddle_block_eviction_splits_rollup_by_cutoff() {
let (config, _guard) = SqliteConfig::test();
let storage = SqliteCdcStorage::new(config);
for v in 1..=6u64 {
storage.write(&cdc_minimal(v)).unwrap();
}
let blocks = storage.compact_all(6, 3, CommitVersion(1000)).unwrap();
assert_eq!(blocks.len(), 1, "all six entries compact into a single block");
let r = storage.drop_before(CommitVersion(4), usize::MAX).unwrap();
assert_eq!(r.count, Count::new(3), "only below-cutoff entries are evicted");
assert_eq!(r.entries.len(), 1);
assert_eq!(r.entries[0].count, Count::new(3));
assert_eq!(r.entries[0].value_bytes, ByteSize::from_bytes(3 * 3));
assert!(storage.read(CommitVersion(3)).unwrap().is_none(), "evicted survivor is gone");
assert!(storage.read(CommitVersion(4)).unwrap().is_some(), "survivors are kept");
assert_eq!(storage.min_version().unwrap(), Some(CommitVersion(4)));
}