use super::super::manifest;
use super::super::*;
use super::Fold;
use crate::test_support::{
make_batch_opk, make_batch_raw, make_schema_pk_u64_payload_string, make_schema_u64_i64, pk_payload_schema,
};
use gnitz_wire::TypeCode;
use gnitz_zset::repr::{pk_group_end, Batch, BatchBuilder};
use gnitz_zset::schema::key::probe_key;
use gnitz_zset::schema::{SchemaColumn, SchemaDescriptor};
use std::path::Path;
impl ShardIndex {
fn run(&mut self, fold: Fold) -> Result<(), StorageError> {
debug_assert!(self.running.is_none(), "one fold at a time");
self.begin(fold)?;
self.finish_fold().map(drop)
}
fn drain(&mut self) -> Result<(), StorageError> {
self.maintain(u64::MAX).map(drop)
}
fn split_overfull_guards(&mut self, level: usize) -> Result<(), StorageError> {
let keys: Vec<PkBuf> = self.levels[level].guards.iter().rev().map(|g| g.guard_key).collect();
for key in keys {
let guards = &self.levels[level].guards;
let fold = guards
.binary_search_by(|g| g.guard_key.cmp(&key))
.ok()
.and_then(|gi| self.split_of(level, &guards[gi]));
if let Some(fold) = fold {
self.run(fold)?;
}
}
Ok(())
}
fn merge_underfull_guards(&mut self, level: usize) -> Result<(), StorageError> {
while let Some(fold) = self.plan_merge(level) {
self.run(fold)?;
}
Ok(())
}
fn rebalance_guards(&mut self, level: usize) -> Result<(), StorageError> {
self.split_overfull_guards(level)?;
self.merge_underfull_guards(level)
}
fn vertical_fold(&mut self, src: usize) -> Result<(), StorageError> {
if let Some(cut) = self.plan_drain(src) {
self.run(cut)?;
}
while let Some(band) = self.bands.pop() {
if let Some(fold) = self.plan_band(band) {
self.run(fold)?;
}
}
self.rebalance_guards(TERMINAL)
}
fn run_compact(&mut self) -> Result<(), StorageError> {
let fold = self.plan_l0_fold();
self.run(fold)?;
for level in [L1, TERMINAL] {
self.rebalance_guards(level)?;
}
while self.levels[L1].bytes() > self.l1_target_bytes() {
let Some(gi) = self.cheapest_l1_guard_to_drain() else {
break;
};
self.vertical_fold(gi)?;
}
Ok(())
}
fn dehydrate_guard(&mut self, guard_idx: usize) -> Result<(), StorageError> {
let fold = self.plan_dehydration(guard_idx);
self.run(fold)
}
fn enforce_capacity(&mut self) -> Result<(), StorageError> {
self.drain()
}
}
fn open(dir: &Path, schema: SchemaDescriptor, budget: ShardBudget) -> ShardIndex {
let dir = dir.to_str().unwrap();
let shards = manifest::read(dir).unwrap().map(|m| m.shards).unwrap_or_default();
ShardIndex::open(dir, schema, budget, false, &shards).unwrap()
}
pub(super) fn fresh(dir: &Path, schema: SchemaDescriptor) -> ShardIndex {
open(dir, schema, ShardBudget::Unbounded)
}
fn dir_names(dir: &Path) -> Vec<String> {
let mut names: Vec<String> = std::fs::read_dir(dir)
.unwrap()
.map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
.collect();
names.sort();
names
}
fn shard_seqs(dir: &Path) -> Vec<u64> {
let mut seqs = manifest::shard_seqs(dir.to_str().unwrap()).unwrap();
seqs.sort_unstable();
seqs
}
fn gk(v: u64) -> PkBuf {
PkBuf::from_bytes(&v.to_be_bytes())
}
fn test_batch(pks: &[u64], values: &[i64]) -> Batch {
let rows: Vec<(u64, i64, i64)> = pks.iter().zip(values).map(|(&p, &v)| (p, 1, v)).collect();
make_batch_raw(&make_schema_u64_i64(), &rows)
}
pub(super) fn spread(i: u64) -> i64 {
i.wrapping_mul(0x9E37_79B9_7F4A_7C15) as i64
}
fn ramp(i: u64) -> i64 {
i64::MIN + ((i as i64) << 48)
}
fn dense_batch(base: u64, n: u64) -> Batch {
let pks: Vec<u64> = (base..base + n).collect();
let vals: Vec<i64> = pks.iter().map(|&p| spread(p)).collect();
test_batch(&pks, &vals)
}
fn fat_batch(base: u64, n: u64, width: usize) -> Batch {
let schema = make_schema_pk_u64_payload_string();
let mut b = BatchBuilder::new(&schema);
for pk in base..base + n {
let mut body = pk.to_string().into_bytes();
body.extend((body.len()..width).map(|i| b'a' + ((pk as usize + i) % 26) as u8));
b.begin_row(pk as u128, 1);
b.put_blob(&body);
b.end_row();
}
b.finish()
}
const STABLE_ROWS: u64 = 2800;
fn assert_stable(len: u64) {
assert!(
(MIN_GUARD_BYTES / 2..MIN_GUARD_BYTES).contains(&len),
"a stable shard must sit between the merge and split thresholds, got {len} B",
);
}
fn seed_guard(idx: &mut ShardIndex, level_idx: usize, key: PkBuf, batch: &Batch, stamp: u64) {
let entry = idx.write_shard(batch, false, Some(stamp)).unwrap();
idx.levels[level_idx].get_or_create_guard(key).entries.push(entry);
idx.mark_published();
}
fn seed_stable(idx: &mut ShardIndex, level_idx: usize, key: PkBuf, base: u64, stamp: u64) -> Vec<u64> {
seed_guard(idx, level_idx, key, &dense_batch(base, STABLE_ROWS), stamp);
assert_stable(
idx.levels[level_idx]
.get_or_create_guard(key)
.entries
.last()
.unwrap()
.shard
.file_len(),
);
(base..base + STABLE_ROWS).collect()
}
fn append_stable(idx: &mut ShardIndex, base: u64) -> u64 {
idx.append_l0_run(&dense_batch(base, STABLE_ROWS)).unwrap();
let len = idx.levels[L0].entries().last().unwrap().shard.file_len();
assert_stable(len);
len
}
fn assert_weighs(idx: &ShardIndex, weight: i64, keys: impl IntoIterator<Item = PkBuf>) {
for k in keys {
let key = k.pk_bytes();
let mut sum = 0;
idx.find_pk_bytes(key, probe_key(key), |shard, start| {
sum += (start..pk_group_end(&**shard, start))
.map(|r| shard.get_weight(r))
.sum::<i64>();
});
assert_eq!(sum, weight, "key {key:?} through the guard partition");
}
}
fn assert_all_found(idx: &ShardIndex, keys: impl IntoIterator<Item = u64>) {
assert_weighs(idx, 1, keys.into_iter().map(gk));
}
fn publish_manifest(idx: &mut ShardIndex) {
let bytes = manifest::encode(&manifest::Manifest {
checkpoint_mark: 0,
caller_record: Vec::new(),
shards: idx.shard_set(),
});
manifest::prepare(&idx.output_dir, &bytes).unwrap();
manifest::commit(&idx.output_dir).unwrap();
idx.mark_published();
idx.unlink_retired();
}
fn reopened_under(mut idx: ShardIndex, budget: ShardBudget) -> ShardIndex {
publish_manifest(&mut idx);
let (dir, schema) = (idx.output_dir.clone(), idx.schema);
drop(idx);
open(Path::new(&dir), schema, budget)
}
#[test]
fn a_fold_and_reload_keep_every_key_and_track_what_is_unsynced() {
let dir = tempfile::tempdir().unwrap();
let schema = make_schema_u64_i64();
let mut idx = fresh(dir.path(), schema);
let pks: Vec<u64> = (1..=5).map(|i| i * 10).collect();
for &pk in &pks {
idx.append_l0_run(&test_batch(&[pk], &[pk as i64])).unwrap();
}
let spills: Vec<String> = idx.unsynced_paths().collect();
assert_eq!(spills.len(), 5);
idx.run_compact().unwrap();
assert!(idx.levels[L0].guards.is_empty() && !idx.levels[L1].guards.is_empty());
assert!(
spills.iter().all(|p| !Path::new(p).exists()),
"unpublished inputs are unlinked"
);
assert!(idx.unsynced_paths().all(|p| !spills.contains(&p)));
assert!(idx.unsynced_paths().next().is_some(), "the outputs wait for a barrier");
assert_all_found(&idx, pks.iter().copied());
assert_weighs(&idx, 0, [gk(99)]);
publish_manifest(&mut idx);
let mut idx2 = fresh(dir.path(), schema);
assert!(idx2.shard_set() == idx.shard_set());
assert!(idx2.unsynced_paths().next().is_none(), "a manifest names durable files");
idx2.append_l0_run(&test_batch(&[99], &[990])).unwrap();
let mut seqs: Vec<u64> = idx2.all_entries().map(|e| e.seq).collect();
seqs.dedup();
assert_eq!(
seqs.len(),
idx2.shard_count(),
"a write after the reload takes a fresh name"
);
assert_all_found(&idx2, pks.iter().copied().chain([99]));
}
#[test]
fn guards_over_the_file_threshold_fold_in_place() {
let dir = tempfile::tempdir().unwrap();
let schema = make_schema_u64_i64();
let mut idx = fresh(dir.path(), schema);
for pk in 1..=GUARD_FILE_THRESHOLD as u64 + 1 {
seed_guard(&mut idx, L1, gk(0), &test_batch(&[pk], &[pk as i64]), 1);
}
for i in 0..GUARD_FILE_THRESHOLD as i64 + 2 {
let rows: Vec<(u64, i64, i64)> = (100..104).map(|k| (k, if i % 2 == 0 { 1 } else { -1 }, 0)).collect();
seed_guard(&mut idx, L1, gk(100), &make_batch_raw(&schema, &rows), 1);
}
idx.split_overfull_guards(L1).unwrap();
let guards: Vec<(PkBuf, usize)> = idx.levels[L1]
.guards
.iter()
.map(|g| (g.guard_key, g.entries.len()))
.collect();
assert_eq!(guards, [(gk(0), 1)]);
assert_all_found(&idx, 1..=GUARD_FILE_THRESHOLD as u64 + 1);
}
#[test]
fn a_spill_of_retractions_folds_into_the_rows_it_cancels() {
let dir = tempfile::tempdir().unwrap();
let mut idx = fresh(dir.path(), make_schema_u64_i64());
seed_guard(&mut idx, L1, gk(0), &dense_batch(0, 1000), 1);
idx.append_l0_run(&dense_batch(0, 50).negated()).unwrap();
idx.drain().unwrap();
assert_eq!(idx.level_shape(), (1, [1, 0]));
idx.append_l0_run(&dense_batch(50, 250).negated()).unwrap();
idx.drain().unwrap();
assert_eq!(idx.level_shape(), (0, [1, 0]));
let guard = &idx.levels[L1].guards[0];
assert_eq!((guard.entries.len(), guard.rows()), (1, 700));
assert_weighs(&idx, 0, (0..300).map(gk));
assert_all_found(&idx, 300..1000);
}
#[test]
fn retractions_that_cancel_nothing_stop_folding_for_them() {
let folds_of = |weight: i64| {
let dir = tempfile::tempdir().unwrap();
let mut idx = fresh(dir.path(), make_schema_u64_i64());
cstats::reset();
for spill in 0..40 {
let run = dense_batch(spill * 100, 100);
let run = if weight < 0 { run.negated() } else { run };
idx.append_l0_run(&run).unwrap();
idx.drain().unwrap();
}
assert_weighs(&idx, weight, (0..4000).map(gk));
cstats::dump().values().map(|p| p.n).sum::<usize>()
};
let (inserted, retracted) = (folds_of(1), folds_of(-1));
assert!(
retracted <= inserted + 2,
"{retracted} folds of retractions against {inserted} of insertions"
);
}
#[test]
fn a_failing_vertical_band_leaves_the_bands_before_it_folded() {
let dir = tempfile::tempdir().unwrap();
let schema = make_schema_u64_i64();
let mut idx = fresh(dir.path(), schema);
let mut dest_pks = Vec::new();
for (base, key) in [(200u64, gk(100)), (100_100, gk(100_000))] {
dest_pks.extend(seed_stable(&mut idx, TERMINAL, key, base, 80));
}
let src_pks: Vec<u64> = vec![100, 150, 100_050, 100_060];
for &pk in &src_pks {
seed_guard(&mut idx, L1, gk(100), &test_batch(&[pk], &[pk as i64]), 100);
}
let (split_outputs, first_band_fold) = (2, 1);
let second_band_fold = idx.shard_seq + split_outputs + first_band_fold + 1;
let blocker = dir.path().join(manifest::shard_name(second_band_fold));
std::fs::create_dir_all(&blocker).unwrap();
let hi_dest_seq = idx.levels[TERMINAL].guards[1].entries[0].seq;
assert!(idx.vertical_fold(0).is_err(), "the second band cannot write");
assert_eq!(idx.levels[L1].guards.len(), 1, "only the failed band is left in L1");
assert_eq!(
idx.levels[L1].guards[0].guard_key,
gk(100_000),
"the failed band keeps its own source guard",
);
assert_eq!(
idx.levels[TERMINAL].guards[1].entries[0].seq, hi_dest_seq,
"the failed band's destination guard is untouched",
);
assert_all_found(&idx, src_pks.iter().chain(&dest_pks).copied());
}
#[test]
fn an_l0_fold_routes_keys_below_the_first_l1_guard() {
let dir = tempfile::tempdir().unwrap();
let schema = make_schema_u64_i64();
let mut idx = fresh(dir.path(), schema);
seed_guard(&mut idx, L1, gk(100), &test_batch(&[100, 200], &[1000, 2000]), 1);
let low_keys = [50u64, 60, 70, 80, 90];
for &k in &low_keys {
idx.append_l0_run(&test_batch(&[k], &[k as i64 * 10])).unwrap();
}
assert!(idx.level_shape().0 > L0_COMPACT_THRESHOLD);
idx.run_compact().unwrap();
assert_all_found(&idx, low_keys.iter().copied().chain([100u64, 200]));
}
#[test]
fn find_guards_for_range_names_every_guard_the_range_meets() {
let range = |l: &FLSMLevel, lo: u64, hi: u64| l.find_guards_for_range(gk(lo).pk_bytes(), gk(hi).pk_bytes());
let mut level = FLSMLevel::default();
for k in [0u64, 100, 200, 300] {
level.guards.push(LevelGuard { guard_key: gk(k), entries: Vec::new() });
}
assert_eq!(range(&level, 10, 50), 0..1);
assert_eq!(range(&level, 100, 250), 1..3);
assert_eq!(range(&level, 0, 999), 0..4);
assert_eq!(range(&level, 200, 200), 2..3);
let mut above_zero = FLSMLevel::default();
for k in [100u64, 200] {
above_zero
.guards
.push(LevelGuard { guard_key: gk(k), entries: Vec::new() });
}
assert_eq!(range(&above_zero, 10, 50), 0..1);
let empty = FLSMLevel::default();
assert!(range(&empty, 0, 100).is_empty());
}
#[test]
fn a_schema_change_rebinds_every_shard() {
let col = |nullable| SchemaColumn::new(TypeCode::I64, nullable);
let u64_pk = SchemaColumn::new(TypeCode::U64, false);
let widened = SchemaDescriptor::new(&[u64_pk, col(false), col(true)], &[0]);
let nullable = SchemaDescriptor::new(&[u64_pk, col(true)], &[0]);
for new in [widened, nullable] {
let dir = tempfile::tempdir().unwrap();
let mut idx = fresh(dir.path(), make_schema_u64_i64());
for i in 0..3u64 {
idx.append_l0_run(&test_batch(&[i * 10 + 1], &[i as i64])).unwrap();
}
publish_manifest(&mut idx);
idx.swap_schema(new).unwrap();
for e in idx.all_entries() {
let row = e.shard.slice_to_owned_batch(0, 1);
assert!(*row.schema() == new);
let nw = row.get_null_word(0);
assert!(
(1..new.num_payload_cols()).all(|pi| gnitz_wire::null_word_get(nw, pi)),
"appended columns read NULL"
);
}
assert!(idx.unsynced_paths().next().is_none(), "a rebind moves no durability");
}
}
#[test]
fn a_corrupt_input_body_fails_the_compaction() {
let dir = tempfile::tempdir().unwrap();
let mut idx = fresh(dir.path(), make_schema_u64_i64());
for i in 0..5u64 {
idx.append_l0_run(&dense_batch(i * 100, 20)).unwrap();
}
let before = shard_seqs(dir.path());
let third = idx.unsynced_paths().nth(2).unwrap();
crate::test_support::flip_last_byte_in_place(Path::new(&third));
assert_eq!(idx.run_compact(), Err(StorageError::Corrupt("body checksum")));
assert_eq!(idx.level_shape().0, 5, "every input stays registered");
assert_eq!(
shard_seqs(dir.path()),
before,
"no input retired, no output left behind"
);
}
#[test]
fn an_ascending_run_folds_into_one_guard_and_is_not_cut_out_of_it() {
const SPILL: u64 = 14_000;
let dir = tempfile::tempdir().unwrap();
let mut idx = fresh(dir.path(), make_schema_u64_i64());
let mut guards_before = 0;
for run in 0..3u64 {
for spill in run * 5..run * 5 + 5 {
let pks: Vec<u64> = (spill * SPILL..(spill + 1) * SPILL).collect();
let vals: Vec<i64> = pks.iter().map(|&p| p as i64).collect();
idx.append_l0_run(&test_batch(&pks, &vals)).unwrap();
}
let spilled = idx.levels[L0].bytes();
idx.run_compact().unwrap();
let guards = &idx.levels[L1].guards;
if run > 0 {
assert_eq!(guards.len(), guards_before + 1, "run {run}");
let folded = guards.last().unwrap();
assert_eq!(folded.entries.len(), 1, "run {run}");
assert!(folded.bytes() > spilled, "premise: the fold's frame is the wider");
}
guards_before = guards.len();
}
assert_all_found(&idx, (0..15 * SPILL).step_by(997));
}
#[test]
fn a_failing_compaction_leaves_its_inputs_and_no_output() {
for failing in [1, 2] {
let dir = tempfile::tempdir().unwrap();
let mut idx = fresh(dir.path(), make_schema_u64_i64());
idx.append_l0_run(&test_batch(&[10, 50], &[1, 1])).unwrap();
idx.append_l0_run(&test_batch(&[150, 250], &[1, 1])).unwrap();
std::fs::create_dir_all(dir.path().join(manifest::shard_name(idx.shard_seq + failing))).unwrap();
let before = shard_seqs(dir.path());
assert!(idx.run_compact().is_err(), "output {failing} cannot be written");
assert_eq!(
shard_seqs(dir.path()),
before,
"output {failing}: no input retired, no output left behind"
);
assert_eq!(idx.level_shape().0, 2, "both inputs stay registered");
assert!(
idx.levels[L1..].iter().all(|l| l.guards.is_empty()),
"nothing was registered"
);
assert_all_found(&idx, [10, 50, 150, 250]);
}
}
#[test]
fn a_vertical_into_a_guards_lower_tail_does_not_shadow_it() {
let dir = tempfile::tempdir().unwrap();
let schema = make_schema_u64_i64();
let mut idx = fresh(dir.path(), schema);
let dest_pks = [50u64, 150, 250];
let vals: Vec<i64> = dest_pks.iter().map(|&p| p as i64).collect();
seed_guard(&mut idx, TERMINAL, gk(200), &test_batch(&dest_pks, &vals), 80);
let src_pks = [100u64, 180];
let vals: Vec<i64> = src_pks.iter().map(|&p| p as i64).collect();
seed_guard(&mut idx, L1, gk(100), &test_batch(&src_pks, &vals), 100);
idx.vertical_fold(0).unwrap();
assert_eq!(idx.levels[TERMINAL].guards.len(), 1, "no second guard was minted");
assert_all_found(&idx, dest_pks.into_iter().chain(src_pks));
}
#[test]
fn a_vertical_rewrites_exactly_the_terminal_guards_its_source_overlaps() {
let dir = tempfile::tempdir().unwrap();
let mut idx = fresh(dir.path(), make_schema_u64_i64());
let mut dest_pks = Vec::new();
for (base, key) in [(200u64, gk(100)), (100_100, gk(100_000)), (200_100, gk(200_000))] {
dest_pks.extend(seed_stable(&mut idx, TERMINAL, key, base, 80));
}
let terminal = |idx: &ShardIndex| -> Vec<(PkBuf, u64)> {
idx.levels[TERMINAL]
.guards
.iter()
.map(|g| (g.guard_key, g.entries[0].seq))
.collect()
};
let before = terminal(&idx);
let src_pks = [100u64, 150, 100_050, 100_060];
let vals: Vec<i64> = src_pks.iter().map(|&p| p as i64).collect();
seed_guard(&mut idx, L1, gk(100), &test_batch(&src_pks, &vals), 100);
idx.vertical_fold(0).unwrap();
assert!(idx.levels[L1].guards.is_empty(), "every band went down");
let after = terminal(&idx);
assert_eq!(after.len(), before.len());
for (i, (a, b)) in after.iter().zip(&before).enumerate() {
assert_eq!(a.0, b.0, "the destination partition is unchanged");
assert_eq!(a.1 != b.1, i < 2, "guard {i} is rewritten iff the source overlaps it");
}
assert_all_found(&idx, src_pks.into_iter().chain(dest_pks));
}
#[test]
fn open_removes_exactly_the_unreferenced_files() {
let dir = tempfile::tempdir().unwrap();
let stray = dir.path().join(manifest::shard_name(7));
std::fs::write(&stray, b"orphan").unwrap();
let mut idx = fresh(dir.path(), make_schema_u64_i64());
assert!(!stray.exists());
idx.append_l0_run(&test_batch(&[10], &[100])).unwrap();
publish_manifest(&mut idx);
let live = manifest::shard_name(idx.shard_seq);
for seq in [99, 7] {
std::fs::write(dir.path().join(manifest::shard_name(seq)), b"x").unwrap();
}
std::fs::write(dir.path().join("manifest.bin.tmp"), b"x").unwrap();
std::fs::write(dir.path().join("other"), b"x").unwrap();
let idx = fresh(dir.path(), make_schema_u64_i64());
assert_all_found(&idx, [10]);
let mut want = [live, "manifest.bin".to_string(), "other".to_string()];
want.sort();
assert_eq!(dir_names(dir.path()), want);
}
const OVER_TARGET_ROWS: u64 = 5000;
#[test]
fn splitting_guard_zero_mints_keys_below_its_own() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
let key = gk(OVER_TARGET_ROWS * 4 / 5);
seed_guard(&mut idx, L1, key, &dense_batch(1, OVER_TARGET_ROWS), 1);
idx.split_overfull_guards(L1).unwrap();
assert!(idx.levels[L1].guards.len() > 1, "the tail split");
assert!(
idx.levels[L1].guards[0].guard_key < key,
"the new lowest guard sits below the key it was split off",
);
assert_all_found(&idx, 1..=OVER_TARGET_ROWS);
}
#[test]
fn a_guard_of_one_distinct_key_neither_splits_nor_refolds() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
let pks = vec![7u64; 9000];
let vals: Vec<i64> = (0..9000).map(ramp).collect();
seed_guard(&mut idx, L1, gk(7), &test_batch(&pks, &vals), 1);
let target = idx.guard_target_bytes(L1);
assert!(idx.levels[L1].guards[0].bytes() > target, "premise: over target");
assert_eq!(
fold_destinations(gk(7), idx.levels[L1].guards[0].entries.iter(), target),
vec![gk(7)]
);
let before = idx.levels[L1].guards[0].entries[0].seq;
idx.split_overfull_guards(L1).unwrap();
idx.split_overfull_guards(L1).unwrap();
assert_eq!(idx.levels[L1].guards.len(), 1);
assert_eq!(
idx.levels[L1].guards[0].entries[0].seq, before,
"an uncuttable guard is not rewritten at all",
);
}
pub(super) fn stride_schema(pk_cols: usize) -> SchemaDescriptor {
pk_payload_schema(&vec![TypeCode::U64; pk_cols])
}
pub(super) fn trailing_gk(pk_cols: usize, i: u64) -> PkBuf {
let mut pk = vec![0u8; (pk_cols - 1) * 8];
pk.extend_from_slice(&i.to_be_bytes());
PkBuf::from_bytes(&pk)
}
pub(super) fn trailing_key_batch(pk_cols: usize, keys: impl IntoIterator<Item = u64>) -> Batch {
let rows: Vec<_> = keys
.into_iter()
.map(|i| (trailing_gk(pk_cols, i).pk_bytes().to_vec(), 1, spread(i)))
.collect();
make_batch_opk(&stride_schema(pk_cols), &rows)
}
#[test]
fn the_byte_target_bounds_a_guard_at_every_stride() {
for pk_cols in [1usize, 3, 4] {
let tmp = tempfile::tempdir().unwrap();
let schema = stride_schema(pk_cols);
assert_eq!(schema.pk_stride(), pk_cols * 8);
let mut idx = fresh(tmp.path(), schema);
let batch = trailing_key_batch(pk_cols, 1..=OVER_TARGET_ROWS);
seed_guard(&mut idx, L1, trailing_gk(pk_cols, 1), &batch, 1);
let target = idx.guard_target_bytes(L1);
let before = idx.levels[L1].guards[0].bytes();
assert!(before > target, "stride {}: {before} B is not over target", pk_cols * 8);
idx.split_overfull_guards(L1).unwrap();
assert!(idx.levels[L1].guards.len() > 1, "stride {}: no split", pk_cols * 8);
assert!(
idx.levels[L1].guards.iter().all(|g| g.bytes() <= target),
"stride {}: a part is still over target",
pk_cols * 8,
);
assert_weighs(
&idx,
1,
(1..=OVER_TARGET_ROWS).step_by(97).map(|i| trailing_gk(pk_cols, i)),
);
}
}
#[test]
fn underfull_guards_merge_at_every_stride() {
for pk_cols in [1usize, 3, 4] {
let tmp = tempfile::tempdir().unwrap();
let schema = stride_schema(pk_cols);
let mut idx = fresh(tmp.path(), schema);
for i in 0..4u64 {
let base = 1 + i * 1000;
let batch = trailing_key_batch(pk_cols, base..base + 100);
seed_guard(&mut idx, L1, trailing_gk(pk_cols, base), &batch, i + 1);
}
assert_eq!(idx.levels[L1].guards.len(), 4);
let bound = idx.guard_target_bytes(L1) / 2;
let total: u64 = idx.levels[L1].bytes();
assert!(
total <= bound,
"stride {}: {total} B does not fit the {bound} B run bound",
pk_cols * 8
);
idx.merge_underfull_guards(L1).unwrap();
assert_eq!(
idx.levels[L1].guards.len(),
1,
"stride {}: one run, one guard",
pk_cols * 8
);
assert_eq!(
idx.levels[L1].guards[0].guard_key,
trailing_gk(pk_cols, 1),
"stride {}: keyed by the run's lowest",
pk_cols * 8,
);
let merged = idx.levels[L1].guards[0].entries[0].seq;
idx.split_overfull_guards(L1).unwrap();
assert_eq!(
idx.levels[L1].guards[0].entries[0].seq,
merged,
"stride {}",
pk_cols * 8
);
}
}
#[test]
fn a_guard_far_over_target_splits_in_bounded_steps() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_pk_u64_payload_string());
let rows = 4_400u64;
seed_guard(&mut idx, L1, gk(1), &fat_batch(1, rows, 1000), 1);
let target = idx.guard_target_bytes(L1);
assert!(
idx.levels[L1].guards[0].bytes() > MAX_PARTS * target,
"premise: the guard must want more parts than one fold may write",
);
idx.split_overfull_guards(L1).unwrap();
assert_eq!(idx.levels[L1].guards.len(), MAX_PARTS as usize);
assert!(
idx.levels[L1].guards.iter().any(|g| g.bytes() > target),
"premise: one bounded fold cannot have finished the job, or the \
convergence below asserts nothing",
);
for _ in 0..4 {
idx.split_overfull_guards(L1).unwrap();
}
assert!(
idx.levels[L1].guards.iter().all(|g| g.bytes() <= target),
"repeated folds converge to guards at the target",
);
assert_all_found(&idx, (1..=rows).step_by(97));
}
#[test]
fn a_merge_run_does_not_cross_a_representation_boundary() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = open(tmp.path(), make_schema_u64_i64(), ShardBudget::Dehydrate(1));
for (i, base) in [1u64, 1000].into_iter().enumerate() {
seed_guard(&mut idx, TERMINAL, gk(base), &dense_batch(base, 200), i as u64 + 1);
}
idx.dehydrate_guard(0).unwrap();
let (dehy, hyd) = terminal_split(&idx);
assert_eq!((dehy.as_slice(), hyd.as_slice()), (&[0][..], &[1][..]));
idx.merge_underfull_guards(TERMINAL).unwrap();
assert_eq!(idx.levels[TERMINAL].guards.len(), 2, "the run broke at the boundary");
assert_eq!(
terminal_split(&idx),
(vec![0], vec![1]),
"neither guard changed representation"
);
}
#[test]
fn a_first_dehydration_never_splits() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = open(tmp.path(), make_schema_u64_i64(), ShardBudget::Dehydrate(1));
seed_guard(&mut idx, TERMINAL, gk(1), &dense_batch(1, OVER_TARGET_ROWS), 1);
idx.dehydrate_guard(0).unwrap();
assert_eq!(idx.levels[TERMINAL].guards.len(), 1);
assert_eq!(idx.levels[TERMINAL].guards[0].entries.len(), 1);
assert!(idx.has_skeleton_shard());
assert_all_found(&idx, 1..=OVER_TARGET_ROWS);
}
#[test]
fn a_dehydrated_guard_over_target_splits_and_stays_skeleton() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = open(tmp.path(), make_schema_u64_i64(), ShardBudget::Dehydrate(1));
let rows = 8_000u64;
seed_guard(&mut idx, TERMINAL, gk(1), &dense_batch(1, rows), 1);
idx.dehydrate_guard(0).unwrap();
let target = idx.guard_target_bytes(TERMINAL);
assert!(
idx.levels[TERMINAL].guards[0].bytes() > target,
"premise: the skeleton itself is over target",
);
idx.split_overfull_guards(TERMINAL).unwrap();
let level = &idx.levels[TERMINAL];
assert!(level.guards.len() > 1, "the skeleton split");
assert!(
level.guards.iter().all(|g| g.dehydrated()),
"a split must not re-hydrate what the sweep evicted",
);
assert!(idx.has_skeleton_shard(), "reads still route through hydration");
assert_all_found(&idx, (1..=rows).step_by(101));
let names: Vec<u64> = idx.levels[TERMINAL].guards.iter().map(|g| g.entries[0].seq).collect();
idx.split_overfull_guards(TERMINAL).unwrap();
let after: Vec<u64> = idx.levels[TERMINAL].guards.iter().map(|g| g.entries[0].seq).collect();
assert_eq!(names, after, "the trigger cleared");
}
#[test]
fn the_guard_count_comes_back_down_after_the_bytes_do() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = open(tmp.path(), make_schema_u64_i64(), ShardBudget::Dehydrate(1));
let keys = 2_500u64;
let pks: Vec<u64> = (1..=keys).flat_map(|k| std::iter::repeat_n(k, 8)).collect();
let vals: Vec<i64> = (0..pks.len() as u64).map(ramp).collect();
seed_guard(&mut idx, TERMINAL, gk(1), &test_batch(&pks, &vals), 1);
idx.split_overfull_guards(TERMINAL).unwrap();
let split_count = idx.levels[TERMINAL].guards.len();
assert!(split_count > 4, "the hydrated level really is finely partitioned");
while let Some(gi) = idx.levels[TERMINAL].guards.iter().position(|g| !g.dehydrated()) {
idx.dehydrate_guard(gi).unwrap();
}
idx.merge_underfull_guards(TERMINAL).unwrap();
let target = idx.guard_target_bytes(TERMINAL);
let count = idx.levels[TERMINAL].guards.len();
assert!(
count < split_count,
"the count followed the bytes down: {split_count} -> {count}"
);
assert!(
count <= (idx.levels[TERMINAL].bytes().div_ceil(target / 2) + 1) as usize,
"{count} guards for {} B at a {target} B target",
idx.levels[TERMINAL].bytes(),
);
assert_weighs(&idx, 8, (1..=keys).step_by(101).map(gk));
}
#[test]
fn a_range_gather_visits_only_the_guards_that_can_own_it() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
for i in 0..4u64 {
let base = 1 + i * 1000;
seed_guard(&mut idx, L1, gk(base), &dense_batch(base, 200), i + 1);
}
let count = |lo: u64, hi: Option<u64>| {
idx.shard_arcs_in_range(gk(lo), hi.map_or_else(|| PkBuf::max(8), gk), true)
.count()
};
assert_eq!(idx.all_shard_arcs_iter().count(), 4);
assert_eq!(count(1100, Some(1100)), 1, "a point read routes to one guard");
assert_eq!(count(1100, Some(2100)), 2, "a range takes the run it spans");
assert_eq!(count(0, None), 4, "an open end takes the rest of the key space");
assert_eq!(count(1500, Some(1500)), 0, "a guard's shard is rejected by its extent");
assert_eq!(count(1500, Some(2100)), 1, "an edge guard is rejected by its extent");
}
#[test]
fn the_guard_target_tracks_the_l0_folds_the_store_has_seen() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
assert_eq!(idx.l0_run_bytes, MIN_GUARD_BYTES, "before any fold");
let mut folded = 0u64;
for i in 0..5u64 {
folded += append_stable(&mut idx, 1 + i * 10_000);
}
idx.run_compact().unwrap();
assert_eq!(idx.l0_run_bytes, folded, "the fold this store actually performed");
assert_eq!(idx.guard_target_bytes(L1), folded, "every unevicted level takes R");
for i in 0..5u64 {
idx.append_l0_run(&dense_batch(500_000 + i * 100, 10)).unwrap();
}
idx.run_compact().unwrap();
assert_eq!(
idx.l0_run_bytes, folded,
"a running max never shrinks under a small fold"
);
seed_guard(
&mut idx,
TERMINAL,
gk(10_000_000),
&dense_batch(10_000_000, 50_000),
200,
);
assert!(
idx.levels.iter().flat_map(|l| &l.guards).any(|g| g.bytes() > folded),
"premise: a guard larger than R"
);
publish_manifest(&mut idx);
let reloaded = fresh(tmp.path(), make_schema_u64_i64());
assert_eq!(reloaded.l0_run_bytes, folded, "R survives a restart as it was observed");
}
#[test]
fn a_budgeted_terminal_target_is_one_sweep_step_within_the_clamp() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
for i in 0..5u64 {
append_stable(&mut idx, 1 + i * 10_000);
}
idx.run_compact().unwrap();
let r = idx.l0_run_bytes;
assert!(r > 2 * MIN_GUARD_BYTES, "premise: R leaves room inside the clamp");
let mid = (MIN_GUARD_BYTES + r) / 2;
for (cap, want) in [(mid * SWEEP_STEPS, mid), (1024, MIN_GUARD_BYTES), (u64::MAX, r)] {
idx = reopened_under(idx, ShardBudget::Dehydrate(cap));
assert_eq!(idx.guard_target_bytes(TERMINAL), want, "capacity {cap}");
assert_eq!(idx.guard_target_bytes(L1), r, "only the evicted level reads the budget");
}
}
#[test]
fn the_balanced_l1_target_computes_its_product_in_u128() {
let (l2, r) = (1u64 << 40, 1u64 << 28);
assert!(l2.checked_mul(r).is_none(), "premise: the product overflows a u64");
assert_eq!(ShardIndex::balanced_l1_target(l2, r), 1 << 35);
assert_eq!(ShardIndex::balanced_l1_target(u64::MAX, u64::MAX), u64::MAX);
}
fn index_with_l0(dir: &Path, n: u64, budget: ShardBudget) -> ShardIndex {
let mut idx = open(dir, make_schema_u64_i64(), budget);
for s in 0..n {
idx.append_l0_run(&dense_batch(s * 1000 + 1, 40)).unwrap();
}
idx
}
fn l0_keys(n: u64) -> impl Iterator<Item = u64> {
(0..n).flat_map(|s| s * 1000 + 1..s * 1000 + 41)
}
fn terminal_split(idx: &ShardIndex) -> (Vec<usize>, Vec<usize>) {
let mut dehy = Vec::new();
let mut hyd = Vec::new();
for (gi, g) in idx.levels[TERMINAL].guards.iter().enumerate() {
if g.dehydrated() {
dehy.push(gi);
} else {
hyd.push(gi);
}
}
(dehy, hyd)
}
#[test]
fn a_slack_capacity_leaves_the_store_untouched() {
let tmp = tempfile::tempdir().unwrap();
let idx = index_with_l0(tmp.path(), 3, ShardBudget::Unbounded);
let before = idx.resident_bytes();
let mut idx = reopened_under(idx, ShardBudget::Dehydrate(before * 4));
idx.enforce_capacity().unwrap();
assert_eq!(idx.resident_bytes(), before, "no compaction ran");
assert_eq!(idx.level_shape().0, 3, "L0 was not pushed down");
assert!(
idx.levels[L1..].iter().all(|l| l.guards.is_empty()),
"no guard was created"
);
assert!(idx.all_entries().all(|e| !e.shard.is_skeleton()));
}
#[test]
fn the_sweep_converges_to_the_skeleton_floor() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = index_with_l0(tmp.path(), 4, ShardBudget::Dehydrate(1));
for _ in 0..8 {
idx.enforce_capacity().unwrap();
}
let (dehy, hyd) = terminal_split(&idx);
assert!(
idx.levels[L0].guards.is_empty() && idx.levels[L1].guards.is_empty(),
"everything sank"
);
assert!(!dehy.is_empty() && hyd.is_empty(), "the floor is fully dehydrated");
assert!(idx.all_entries().all(|e| e.shard.is_skeleton()));
assert_all_found(&idx, l0_keys(4));
let floor = idx.resident_bytes();
assert!(floor > 1);
idx.enforce_capacity().unwrap();
assert_eq!(idx.resident_bytes(), floor, "the floor is the fixpoint");
}
#[test]
fn a_delta_budget_drops_its_victim_and_raises_the_floor() {
let tmp = tempfile::tempdir().unwrap();
const SHARDS: u64 = 4;
const LAST_KEY: u64 = (SHARDS - 1) * 1000 + 40;
let mut idx = index_with_l0(tmp.path(), SHARDS, ShardBudget::Drop(1));
assert_eq!(idx.dropped_max(), PkBuf::zeroed(8), "nothing dropped yet");
for _ in 0..12 {
idx.enforce_capacity().unwrap();
}
assert!(idx.levels[L0].guards.is_empty(), "L0 sank");
assert!(idx.levels[L1].guards.is_empty(), "L1 sank");
assert!(idx.levels[TERMINAL].guards.is_empty(), "a drop removes its guard");
assert_eq!(idx.resident_bytes(), 0, "a delta store has no floor to stop above");
assert_eq!(
idx.dropped_max(),
gk(LAST_KEY),
"the watermark is the HIGHEST key dropped, taken from the victim's pk_max"
);
assert!(shard_seqs(tmp.path()).is_empty(), "every dropped shard is unlinked");
}
#[test]
fn ordinary_compaction_of_a_delta_store_keeps_its_rows() {
let tmp = tempfile::tempdir().unwrap();
let idx = index_with_l0(tmp.path(), 4, ShardBudget::Unbounded);
let before = idx.resident_bytes();
let mut idx = reopened_under(idx, ShardBudget::Drop(before * 8));
idx.run_compact().unwrap(); for gi in (0..idx.levels[L1].guards.len()).rev() {
idx.vertical_fold(gi).unwrap(); }
idx.enforce_capacity().unwrap();
assert_eq!(idx.dropped_max(), PkBuf::zeroed(8), "ordinary compaction drops nothing");
assert!(idx.resident_bytes() > 0, "the rows survived");
assert!(
idx.resident_bytes() <= before,
"a fold never grows the store: {} -> {}",
before,
idx.resident_bytes()
);
assert_all_found(&idx, l0_keys(4));
}
#[test]
fn a_published_shard_a_fold_supersedes_waits_for_the_barrier() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = index_with_l0(tmp.path(), 5, ShardBudget::Unbounded);
publish_manifest(&mut idx);
let inputs: Vec<String> = idx
.all_entries()
.map(|e| manifest::shard_path(&idx.output_dir, e.seq))
.collect();
assert_eq!(inputs.len(), 5);
idx.run_compact().unwrap();
assert!(
inputs.iter().all(|p| Path::new(p).exists()),
"every published input waits for the barrier"
);
publish_manifest(&mut idx);
assert!(
!inputs.iter().any(|p| Path::new(p).exists()),
"the post-publish drain removes it"
);
assert_eq!(
shard_seqs(tmp.path()).len(),
idx.shard_count(),
"only the live shards remain"
);
assert_all_found(&fresh(tmp.path(), make_schema_u64_i64()), l0_keys(5));
}
#[test]
fn a_swept_delta_store_plateaus_under_a_steady_write_stream() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = open(tmp.path(), make_schema_u64_i64(), ShardBudget::Drop(1));
let mut early = 0u64;
for round in 0..60u64 {
idx.append_l0_run(&dense_batch(round * 1000 + 1, 40)).unwrap();
idx.drain().unwrap();
if round == 9 {
early = idx.resident_bytes();
}
}
let late = idx.resident_bytes();
assert!(
late <= early.max(1) * 2,
"footprint grew {early} -> {late} bytes over 50 further rounds of the same \
write rate — the sweep is not keeping pace",
);
assert!(idx.dropped_max() > PkBuf::zeroed(8), "the sweep dropped something");
assert_eq!(shard_seqs(tmp.path()).len(), idx.shard_count(), "no orphan shard files");
}
#[test]
fn a_drop_removes_nothing_above_the_floor_it_raises() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
let mut written: Vec<u64> = Vec::new();
let add = |idx: &mut ShardIndex, written: &mut Vec<u64>, round: u64| {
let base = round * 1000 + 1;
idx.append_l0_run(&dense_batch(base, 40)).unwrap();
written.extend(base..base + 40);
};
for round in 0..6u64 {
add(&mut idx, &mut written, round);
}
idx.run_compact().unwrap();
let budget = idx.resident_bytes();
let mut idx = reopened_under(idx, ShardBudget::Drop(budget));
for round in 6..24u64 {
add(&mut idx, &mut written, round);
idx.drain().unwrap();
}
let floor = idx.dropped_max();
let retained = idx.total_rows();
assert!(
floor > PkBuf::zeroed(8),
"the sweep dropped nothing — nothing is being tested"
);
assert!(retained > 0, "the sweep emptied the store — nothing is being tested");
let above = written.iter().filter(|&&k| gk(k) > floor).count();
assert_eq!(
retained, above,
"floor {floor:?}: every one of the {above} rows above it must survive, and \
every row at or below it must be gone — {retained} retained",
);
}
#[test]
fn dehydration_takes_the_oldest_written_terminal_guard_first() {
let tmp = tempfile::tempdir().unwrap();
let schema = make_schema_u64_i64();
let mut idx = fresh(tmp.path(), schema);
for (i, base) in [1u64, 10_000, 20_000].into_iter().enumerate() {
seed_stable(&mut idx, TERMINAL, gk(base), base, 30 - i as u64);
}
let (dehy, hyd) = terminal_split(&idx);
assert!(dehy.is_empty() && hyd.len() >= 2, "several hydrated terminal guards");
let oldest = 2;
let oldest_key = gk(20_000);
let live_before: Vec<(u128, i64)> = {
let g = &idx.levels[TERMINAL].guards[oldest];
let s = &g.entries[0].shard;
(0..s.row_count())
.map(|i| (gnitz_wire::widen_pk_be(s.get_pk_bytes(i)), s.get_weight(i)))
.collect()
};
let cap = idx.resident_bytes() - 1;
let mut idx = reopened_under(idx, ShardBudget::Dehydrate(cap));
idx.enforce_capacity().unwrap();
let (dehy, _) = terminal_split(&idx);
assert_eq!(dehy.len(), 1, "dehydration stops as soon as the cap is met");
assert_eq!(
idx.levels[TERMINAL].guards[dehy[0]].guard_key, oldest_key,
"write-recency victim",
);
let g = &idx.levels[TERMINAL].guards[dehy[0]];
let s = &g.entries[0].shard;
let live_after: Vec<(u128, i64)> = (0..s.row_count())
.map(|i| (gnitz_wire::widen_pk_be(s.get_pk_bytes(i)), s.get_weight(i)))
.collect();
assert_eq!(live_after, live_before, "keys and coarse weights survive dehydration");
}
#[test]
fn a_dehydrated_guard_stays_dehydrated_under_ordinary_compaction() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = open(tmp.path(), make_schema_u64_i64(), ShardBudget::Dehydrate(1));
for _ in 0..2 {
idx.append_l0_run(&dense_batch(1, 20)).unwrap();
}
idx.run_compact().unwrap();
idx.enforce_capacity().unwrap();
assert_eq!(terminal_split(&idx), (vec![0], vec![]), "the one guard is dehydrated");
idx.append_l0_run(&dense_batch(1, 20)).unwrap();
idx.run_compact().unwrap();
assert!(idx.levels[L1].guards.is_empty(), "the drain emptied L1");
assert_eq!(
terminal_split(&idx),
(vec![0], vec![]),
"the derived rule kept it skeleton"
);
assert_weighs(&idx, 3, (1..=20).map(gk));
}
#[test]
fn a_jump_in_r_is_merged_within_each_calls_budget() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
const GUARDS: u64 = 40;
for i in 0..GUARDS {
seed_stable(&mut idx, TERMINAL, gk(i * 10_000 + 1), i * 10_000 + 1, 1);
}
for i in 0..5u64 {
append_stable(&mut idx, 1_000_000 + i * 10_000);
}
let fold = idx.plan_l0_fold();
idx.run(fold).unwrap();
let r = idx.l0_run_bytes;
assert!(r > 2 * MIN_GUARD_BYTES, "premise: R jumped past the seeded guards");
cstats::reset();
let done = idx.maintain(r).unwrap();
assert!(idx.owed(), "premise: one budget does not finish the job");
let merged = cstats::dump()[&CompactionKind::GuardMerge].in_bytes;
assert!(
done.read < 2 * r && merged < 2 * r,
"{merged} B merged against R = {r} B"
);
let after_one = idx.levels[TERMINAL].guards.len();
assert!(after_one > GUARDS as usize / 2, "the level was not rewritten whole");
idx.drain().unwrap();
assert!(!idx.owed());
assert!(idx.levels[TERMINAL].guards.len() < after_one, "later calls converge");
assert_all_found(&idx, (0..GUARDS).map(|i| i * 10_000 + 1));
}
#[test]
fn a_call_stops_at_the_destination_past_its_budget() {
use crate::test_support::Rng;
const KEYS: u64 = 400_000;
let tmp = tempfile::tempdir().unwrap();
let schema = make_schema_u64_i64();
let mut idx = fresh(tmp.path(), schema);
let mut rng = Rng::new(11);
let (mut calls, mut part_way, mut deferred) = (0, 0, 0);
for _ in 0..120 {
let mut keys: Vec<u64> = (0..2048).map(|_| rng.gen_range(KEYS)).collect();
keys.sort_unstable();
keys.dedup();
idx.append_l0_run(&test_batch(&keys, &keys.iter().map(|&k| k as i64).collect::<Vec<_>>()))
.unwrap();
let before = idx.shard_seq;
let done = idx.maintain(1).unwrap();
assert!(done.read <= 2 * idx.l0_run_bytes + 1, "one call read {} B", done.read);
assert!(idx.shard_seq - before <= 1, "one call wrote more than one shard");
calls += 1;
part_way += usize::from(idx.running.is_some());
deferred += usize::from(idx.owed());
}
assert!(part_way > 0, "premise: no call left a fold part-way through");
assert!(deferred > calls / 2, "premise: the store rarely owed past one call");
idx.drain().unwrap();
assert!(!idx.owed() && idx.running.is_none());
assert!(idx.levels[L0].entries().count() <= L0_COMPACT_THRESHOLD);
for level in &idx.levels[L1..] {
assert!(level.guards.iter().all(|g| g.entries.len() <= GUARD_FILE_THRESHOLD));
}
}
#[test]
fn a_vertical_of_a_guard_with_a_low_tail_reads_each_terminal_guard_once() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
seed_guard(&mut idx, TERMINAL, gk(10), &test_batch(&[10, 20], &[1, 2]), 1);
seed_guard(&mut idx, TERMINAL, gk(100), &test_batch(&[100, 110], &[3, 4]), 1);
seed_guard(&mut idx, L1, gk(150), &test_batch(&[5, 105, 160], &[5, 6, 7]), 2);
cstats::reset();
idx.vertical_fold(0).unwrap();
let stats = cstats::dump();
let vertical = &stats[&CompactionKind::Vertical];
assert_eq!(
(vertical.n, vertical.in_files),
(2, 4),
"two bands, each read with the one terminal guard it meets"
);
assert_all_found(&idx, [5, 10, 20, 100, 105, 110, 160]);
}
#[test]
fn an_ascending_stream_is_written_once() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
cstats::reset();
let mut keys = Vec::new();
for i in 0..30u64 {
append_stable(&mut idx, 1 + i * 10_000);
keys.extend(1 + i * 10_000..1 + i * 10_000 + STABLE_ROWS);
idx.drain().unwrap();
}
let folded: Vec<u64> = idx.levels[L1].entries().map(|e| e.seq).collect();
assert!(folded.len() > 1, "premise: several L0 folds reached L1");
while !idx.levels[L1].guards.is_empty() {
idx.vertical_fold(0).unwrap();
}
let rewrites: Vec<CompactionKind> = cstats::dump()
.into_keys()
.filter(|k| !matches!(k, CompactionKind::L0Fold | CompactionKind::GuardMerge))
.collect();
assert_eq!(rewrites, [], "nothing but the L0 fold wrote a row");
let terminal: Vec<u64> = idx.levels[TERMINAL].entries().map(|e| e.seq).collect();
assert!(
folded.iter().any(|seq| terminal.contains(seq)),
"an L0 fold's output is the terminal shard"
);
assert_all_found(&idx, keys.into_iter().step_by(97));
}
#[test]
fn a_few_keys_above_l1_mint_no_guard() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
seed_stable(&mut idx, L1, gk(1), 1, 1);
for i in 0..5u64 {
let pks: Vec<u64> = (1..=2000).chain([1_000_000 + i]).collect();
let vals: Vec<i64> = pks.iter().map(|&p| (p + i) as i64).collect();
idx.append_l0_run(&test_batch(&pks, &vals)).unwrap();
}
idx.run_compact().unwrap();
assert!(
idx.levels[L1].guards.iter().all(|g| g.guard_key < gk(1_000_000)),
"no guard was minted for five rows"
);
assert_weighs(&idx, 1, (1_000_000..1_000_005).map(gk));
}
#[test]
fn a_fence_stays_above_a_tail_only_l1_guard() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
seed_guard(&mut idx, L1, gk(100_000), &test_batch(&[5, 6], &[5, 6]), 1);
for i in 0..5u64 {
idx.append_l0_run(&dense_batch(50 + i * 3000, 2000)).unwrap();
}
idx.run_compact().unwrap();
assert_all_found(
&idx,
[5, 6]
.into_iter()
.chain((0..5).flat_map(|i| 50 + i * 3000..2050 + i * 3000)),
);
}
#[test]
fn a_band_at_the_key_of_a_tail_only_terminal_guard_merges_with_it() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
seed_guard(&mut idx, TERMINAL, gk(100), &test_batch(&[10, 20], &[1, 2]), 1);
seed_guard(&mut idx, L1, gk(100), &test_batch(&[100, 150], &[3, 4]), 2);
idx.vertical_fold(0).unwrap();
assert_eq!(idx.levels[TERMINAL].guards.len(), 1);
assert_all_found(&idx, [10, 20, 100, 150]);
}
#[test]
fn a_vertical_over_target_writes_its_parts_directly() {
let tmp = tempfile::tempdir().unwrap();
let mut idx = fresh(tmp.path(), make_schema_u64_i64());
seed_stable(&mut idx, TERMINAL, gk(1), 1, 1);
seed_guard(&mut idx, L1, gk(2000), &dense_batch(2000, STABLE_ROWS), 2);
cstats::reset();
idx.vertical_fold(0).unwrap();
let stats = cstats::dump();
assert_eq!(stats[&CompactionKind::Vertical].n, 1);
assert!(
!stats.contains_key(&CompactionKind::GuardSplit),
"the fold cut its own output"
);
let target = idx.guard_target_bytes(TERMINAL);
let terminal = &idx.levels[TERMINAL].guards;
assert!(terminal.len() > 1 && terminal.iter().all(|g| g.bytes() <= target));
assert_weighs(&idx, 1, (1..2000).chain(2801..4800).step_by(13).map(gk));
assert_weighs(&idx, 2, (2000..=2800).step_by(13).map(gk));
}
use gnitz_zset::repr::{from_runs, Run};
use std::collections::{BTreeMap, BTreeSet};
type Elem = (Vec<u8>, i64);
fn model_key(pk_cols: usize, i: u64) -> Vec<u8> {
let mut pk = Vec::with_capacity(pk_cols * 8);
if pk_cols > 1 {
pk.extend_from_slice(&(i >> 9).to_be_bytes());
}
for _ in 2..pk_cols {
pk.extend_from_slice(&7u64.to_be_bytes());
}
pk.extend_from_slice(&i.to_be_bytes());
pk
}
#[derive(Clone, Default)]
struct Model {
rows: BTreeMap<Elem, i64>,
seen: BTreeSet<Vec<u8>>,
next_payload: i64,
next_key: u64,
}
impl Model {
fn spill(&mut self, rng: &mut crate::test_support::Rng, pk_cols: usize, monotone: bool) -> Batch {
let n = 200 + rng.gen_range(2800);
let mut delta: BTreeMap<Elem, i64> = BTreeMap::new();
for _ in 0..n {
if !monotone && !self.rows.is_empty() && rng.gen_range(4) == 0 {
let nth = rng.gen_range(self.rows.len().min(64) as u64) as usize;
let (elem, _) = self.rows.iter().nth(nth).unwrap();
let elem = elem.clone();
*delta.entry(elem.clone()).or_default() -= 1;
let w = self.rows.get_mut(&elem).unwrap();
*w -= 1;
if *w == 0 {
self.rows.remove(&elem);
}
} else {
let i = if monotone {
self.next_key += 1 + rng.gen_range(3);
self.next_key
} else {
rng.gen_range(20_000)
};
let pk = model_key(pk_cols, i);
self.next_payload += 1;
let w = 1 + rng.gen_range(2) as i64;
self.seen.insert(pk.clone());
*delta.entry((pk.clone(), self.next_payload)).or_default() += w;
*self.rows.entry((pk, self.next_payload)).or_default() += w;
}
}
let rows: Vec<(Vec<u8>, i64, i64)> = delta
.into_iter()
.filter(|&(_, w)| w != 0)
.map(|((pk, pay), w)| (pk, w, pay))
.collect();
make_batch_opk(&stride_schema(pk_cols), &rows)
}
}
fn check_model(idx: &ShardIndex, m: &Model, floor: PkBuf, what: &str) {
assert!(idx.levels[L0].guards.len() <= 1, "{what}: more than one L0 guard");
for (li, level) in idx.levels.iter().enumerate() {
assert!(
level.guards.windows(2).all(|w| w[0].guard_key < w[1].guard_key),
"{what}: L{li} guards not sorted and distinct"
);
for (gi, g) in level.guards.iter().enumerate() {
assert!(!g.entries.is_empty(), "{what}: L{li} guard {gi} is empty");
if li == TERMINAL {
assert_eq!(g.entries.len(), 1, "{what}: terminal guard {gi}");
} else {
assert!(
g.entries.iter().all(|e| !e.shard.is_skeleton()),
"{what}: skeleton above the terminal level"
);
}
if li == L0 {
continue;
}
let (lo, hi) = g.key_extent();
assert!(
gi == 0 || lo >= g.guard_key,
"{what}: L{li} guard {gi} holds a row below its key"
);
if let Some(next) = level.guards.get(gi + 1) {
assert!(hi < next.guard_key, "{what}: L{li} guard {gi} reaches into the next");
}
}
}
let skeleton = idx.all_entries().any(|e| e.shard.is_skeleton());
let dropping = matches!(idx.budget, ShardBudget::Drop(_));
if !matches!(idx.budget, ShardBudget::Dehydrate(_)) {
assert!(!skeleton, "{what}: skeleton in a store that never dehydrates");
}
assert_eq!(skeleton, idx.has_skeleton_shard(), "{what}: has_skeleton_shard");
let live = |pk: &[u8]| !dropping || PkBuf::from_bytes(pk) > floor;
let mut per_pk: BTreeMap<Vec<u8>, i64> = BTreeMap::new();
for ((pk, _), w) in &m.rows {
*per_pk.entry(pk.clone()).or_default() += w;
}
for pk in m.seen.iter().filter(|pk| live(pk)) {
let mut sum = 0;
idx.find_pk_bytes(pk, probe_key(pk), |shard, start| {
sum += (start..pk_group_end(&**shard, start))
.map(|r| shard.get_weight(r))
.sum::<i64>();
});
assert_eq!(sum, per_pk.get(pk).copied().unwrap_or(0), "{what}: key {pk:?} by probe");
}
let cursor = from_runs(idx.all_shard_arcs_iter().map(Run::Shard), idx.schema, idx.shard_count());
if skeleton {
let mut src = gnitz_zset::repr::SourceCursor::Full(Box::new(cursor));
let mut sums: BTreeMap<Vec<u8>, i64> = BTreeMap::new();
let mut skel = gnitz_zset::repr::SkeletonKeys::default();
while let Some(got) = src.drain_live_chunk(4096, &mut skel) {
for r in 0..got.len() {
*sums.entry(got.get_pk_bytes(r).to_vec()).or_default() += got.get_weight(r);
}
}
let stride = idx.schema.pk_stride();
assert_eq!(
skel.keys.len() / stride,
skel.coarse.len(),
"{what}: debug build records coarse weights"
);
for (k, w) in skel.keys.chunks_exact(stride).zip(&skel.coarse) {
*sums.entry(k.to_vec()).or_default() += w;
}
sums.retain(|_, w| *w != 0);
assert_eq!(sums, per_pk, "{what}: per-PK sums by cursor");
} else {
let got = cursor.materialize();
let mut elems: BTreeMap<Elem, i64> = BTreeMap::new();
for r in 0..got.len() {
let pk = got.get_pk_bytes(r).to_vec();
if live(&pk) {
*elems
.entry((pk, crate::test_support::payload0_i64(&*got, r)))
.or_default() += got.get_weight(r);
}
}
elems.retain(|_, w| *w != 0);
let want: BTreeMap<Elem, i64> = m
.rows
.iter()
.filter(|((pk, _), _)| live(pk))
.map(|(e, &w)| (e.clone(), w))
.collect();
assert_eq!(elems, want, "{what}: Z-set by cursor");
}
}
#[test]
fn shard_index_model() {
let configs = [
(ShardBudget::Unbounded, false),
(ShardBudget::Unbounded, true),
(ShardBudget::Dehydrate(150_000), false),
(ShardBudget::Dehydrate(1), true),
(ShardBudget::Drop(150_000), true),
];
let mut interrupted = 0;
for seed in 1..=gnitz_foundation::env::env_num("GNITZ_MODEL_SEEDS", 1u64) {
for pk_cols in [1usize, 3] {
for (ci, &(budget, monotone)) in configs.iter().enumerate() {
let tmp = tempfile::tempdir().unwrap();
let schema = stride_schema(pk_cols);
let mut rng = crate::test_support::Rng::new(seed * 1_000_003 + (pk_cols * 16 + ci) as u64);
let mut idx = open(tmp.path(), schema, budget);
let mut m = Model::default();
let mut published = m.clone();
let mut floor = PkBuf::zeroed(schema.pk_stride());
let mut published_floor = floor;
for step in 0..30 {
let blocker = (rng.gen_range(5) == 0).then(|| {
let p = tmp
.path()
.join(manifest::shard_name(idx.shard_seq + 1 + rng.gen_range(4)));
std::fs::create_dir_all(&p).unwrap();
p
});
let op = rng.gen_range(9);
let what = format!(
"seed {seed} stride {} config {ci} step {step} op {op} fail {}",
pk_cols * 8,
blocker.is_some()
);
match op {
0..=2 => {
let mut next = m.clone();
let batch = next.spill(&mut rng, pk_cols, monotone);
if !batch.is_empty() && idx.append_l0_run(&batch).is_ok() {
m = next;
let _ = match rng.gen_range(2) {
0 => idx.drain(),
_ => idx.maintain(1).map(drop),
};
interrupted += usize::from(idx.running.is_some());
}
}
3 if !idx.levels[L0].guards.is_empty() => {
let _ = idx.finish_fold();
idx.bands.clear();
let _ = idx.run_compact();
}
4 if !idx.levels[L1].guards.is_empty() => {
let _ = idx.finish_fold();
idx.bands.clear();
if !idx.levels[L1].guards.is_empty() {
let gi = rng.gen_range(idx.levels[L1].guards.len() as u64) as usize;
let _ = idx.vertical_fold(gi);
}
}
5 => {
let _ = idx.enforce_capacity();
}
6 => {
let _ = idx.finish_fold();
idx.bands.clear();
let _ = idx.rebalance_guards([L1, TERMINAL][rng.gen_range(2) as usize]);
}
7 => {
if let Some(p) = &blocker {
std::fs::remove_dir(p).unwrap();
}
let _ = idx.finish_fold();
floor = floor.max(idx.dropped_max());
idx = reopened_under(idx, budget);
(published, published_floor) = (m.clone(), floor);
assert_eq!(shard_seqs(tmp.path()).len(), idx.shard_count(), "{what}: shard files");
}
8 => {
if let Some(p) = &blocker {
std::fs::remove_dir(p).unwrap();
}
drop(idx);
idx = open(tmp.path(), schema, budget);
(m, floor) = (published.clone(), published_floor);
}
_ => {}
}
if let Some(p) = blocker.filter(|p| p.is_dir()) {
std::fs::remove_dir(&p).unwrap();
}
floor = floor.max(idx.dropped_max());
check_model(&idx, &m, floor, &what);
}
}
}
}
assert!(interrupted > 0, "premise: no fold was ever left part-way through");
}
#[test]
fn tier_folds_keep_the_zset() {
use crate::test_support::Rng;
use gnitz_zset::repr::merge_and_route;
use std::collections::BTreeMap;
const KEYS: u64 = 400_000;
for budget in [ShardBudget::Unbounded, ShardBudget::Dehydrate(4 << 20)] {
let tmp = tempfile::tempdir().unwrap();
let schema = make_schema_u64_i64();
let mut idx = open(tmp.path(), schema, budget);
let mut rng = Rng::new(7);
let mut version = vec![0u64; KEYS as usize];
let mut model: BTreeMap<(u64, i64), i64> = BTreeMap::new();
cstats::reset();
for _ in 0..200 {
let mut rows = Vec::new();
for _ in 0..2048 {
let k = if rng.gen_range(2) == 0 {
rng.gen_range(KEYS / 50)
} else {
rng.gen_range(KEYS)
};
let v = &mut version[k as usize];
if *v > 0 && rng.gen_range(4) > 0 {
rows.push((k, -1, spread(k ^ *v << 40)));
}
*v += 1;
rows.push((k, 1, spread(k ^ *v << 40)));
}
for &(k, w, val) in &rows {
*model.entry((k, val)).or_default() += w;
}
idx.append_l0_run(&make_batch_raw(&schema, &rows).into_consolidated())
.unwrap();
idx.drain().unwrap();
for level in &idx.levels[L1..] {
assert!(level.guards.iter().all(|g| g.entries.len() <= GUARD_FILE_THRESHOLD));
}
}
let stats = cstats::dump();
for kind in [
CompactionKind::TierFold,
CompactionKind::GuardSplit,
CompactionKind::Vertical,
] {
assert!(stats.contains_key(&kind), "no {kind:?} ran");
}
model.retain(|_, w| *w != 0);
let skeleton = idx.has_skeleton_shard();
assert_eq!(skeleton, matches!(budget, ShardBudget::Dehydrate(_)));
for reopen in [false, true] {
if reopen {
idx = reopened_under(idx, budget);
}
let shards: Vec<Rc<MappedShard>> = idx.all_shard_arcs_iter().collect();
let inputs: Vec<&MappedShard> = shards.iter().map(|s| &**s).collect();
let mut held: BTreeMap<(u64, i64), i64> = BTreeMap::new();
let mut per_key: BTreeMap<u64, i64> = BTreeMap::new();
merge_and_route(&inputs, &[gk(0)], false, &schema, &mut |_, _, batch| {
for row in 0..batch.len() {
let k = u64::from_be_bytes(batch.get_pk_bytes(row).try_into().unwrap());
*per_key.entry(k).or_default() += batch.get_weight(row);
if !skeleton {
let val = i64::from_le_bytes(batch.get_col_ptr(row, 0, 8).try_into().unwrap());
*held.entry((k, val)).or_default() += batch.get_weight(row);
}
}
Ok(())
})
.unwrap();
let mut expect: BTreeMap<u64, i64> = BTreeMap::new();
for (&(k, _), &w) in &model {
*expect.entry(k).or_default() += w;
}
assert_eq!(per_key, expect, "summed weight per key, reopen {reopen}");
if !skeleton {
assert_eq!(held, model, "rows, reopen {reopen}");
}
}
}
}