use super::*;
use crate::test_support::{arb_fold_case, fold_batch, fold_schemas, zset_of, zset_sum};
use gnitz_zset::repr::pk_group_end;
use gnitz_zset::repr::BatchBuilder;
use proptest::prelude::*;
proptest! {
#[test]
fn a_run_set_holds_the_zset_of_its_pushes(
(si, rows) in arb_fold_case(),
steps in prop::collection::vec((0usize..10, 0u8..8), 1..40),
) {
let s = fold_schemas()[si];
let mut set = RunSet::new(1 << 20);
let mut pushed: Vec<Batch> = Vec::new();
let mut rows = rows.iter().cycle();
for (n, action) in steps {
let run = fold_batch(&s, &rows.by_ref().take(n).cloned().collect::<Vec<_>>()).into_consolidated();
pushed.push(run.clone());
set.push(run, &s);
match action {
0 => set.fold(&s),
1 => {
set.clear();
pushed.clear();
}
_ => {}
}
let held: Vec<Batch> = set.runs.iter().map(|r| (**r).clone()).collect();
prop_assert_eq!(zset_sum(&held, &s), zset_sum(&pushed, &s));
prop_assert!(set.len() < FOLD_THRESHOLD);
prop_assert_eq!(set.bytes, held.iter().map(Batch::total_bytes).sum::<usize>());
let want = zset_sum(&pushed, &s);
for key in want.keys().map(|k| &k.0) {
let mut sum = 0;
set.find_pk_bytes(key, probe_key(key), |run, start| {
sum += (start..pk_group_end(&**run, start)).map(|r| run.get_weight(r)).sum::<i64>();
});
let want_sum: i64 = want.iter().filter(|(k, _)| &k.0 == key).map(|(_, w)| w).sum();
prop_assert_eq!(sum, want_sum);
}
}
let held = zset_sum(&pushed, &s);
let failed = set.spill(&s, |_| Err(()));
prop_assert_eq!(failed, if held.is_empty() { Ok(false) } else { Err(()) });
let mut written = Default::default();
let wrote = set.spill(&s, |run| {
written = zset_of(run, &s);
Ok::<(), ()>(())
});
prop_assert_eq!(wrote, Ok(!held.is_empty()));
prop_assert_eq!(written, held, "the failed write left every row held");
prop_assert_eq!((set.len(), set.bytes), (0, 0));
}
}
#[test]
fn a_fold_drops_the_filter_once_most_of_its_keys_are_gone() {
use crate::test_support::{make_batch_raw, make_schema_u64_i64, opk_pk};
let schema = make_schema_u64_i64();
let mut set = RunSet::new(0);
let push = |set: &mut RunSet, keys: std::ops::Range<u64>, w: i64| {
let rows: Vec<_> = keys.map(|k| (k, w, 7)).collect();
set.push(make_batch_raw(&schema, &rows).into_consolidated(), &schema);
};
let holds = |set: &RunSet, k: u64| {
let key = opk_pk(&schema, &[k as u128]);
let mut found = false;
set.find_pk_bytes(&key, probe_key(&key), |_, _| found = true);
found
};
push(&mut set, 0..10, 1);
assert!(holds(&set, 3), "the first probe builds the filter");
push(&mut set, 0..5, -1);
set.fold(&schema);
assert!(set.bloom.get().is_some());
assert!(holds(&set, 7) && !holds(&set, 3));
push(&mut set, 100..1100, 1);
push(&mut set, 100..1100, -1);
set.fold(&schema);
assert!(set.bloom.get().is_none());
assert!(holds(&set, 7) && !holds(&set, 3) && !holds(&set, 105));
}
#[test]
fn a_fold_merges_runs_of_different_nullability_under_its_own_schema() {
use gnitz_wire::TypeCode;
use gnitz_zset::schema::SchemaColumn;
let label = |nullable| {
SchemaDescriptor::new(
&[
SchemaColumn::new(TypeCode::U64, false),
SchemaColumn::new(TypeCode::I64, nullable),
],
&[0],
)
};
let (not_null, nullable) = (label(false), label(true));
let run = |schema: SchemaDescriptor, rows: &[(u128, Option<i64>, i64)]| {
let mut b = BatchBuilder::new(&schema);
for &(pk, v, w) in rows {
b.begin_row(pk, w);
b.put_opt_int(v.map(|v| v as u128));
b.end_row();
}
let mut b = b.finish();
b.certify_consolidated();
b
};
let mut set = RunSet::new(1 << 20);
set.push(
run(not_null, &[(1, Some(-3), 1), (2, Some(7), 1), (3, Some(8), 1)]),
¬_null,
);
set.push(run(nullable, &[(1, None, 1), (1, Some(-3), -1)]), &nullable);
set.fold(&nullable);
let folded = &set.runs[0];
let got: Vec<(u128, Option<i64>, i64)> = (0..folded.len())
.map(|i| {
let v = (folded.get_null_word(i) & 1 == 0)
.then(|| i64::from_le_bytes(folded.get_col_ptr(i, 0, 8).try_into().unwrap()));
(folded.get_pk(i), v, folded.get_weight(i))
})
.collect();
assert_eq!(got, vec![(1, None, 1), (2, Some(7), 1), (3, Some(8), 1)]);
}
fn churn_value(pk: u64, generation: u64) -> Vec<u8> {
format!("{pk:08}-{generation:08}-{}", "v".repeat(24)).into_bytes()
}
fn string_run(schema: &SchemaDescriptor, rows: &[(u64, i64, Vec<u8>)]) -> Batch {
let rows: Vec<(u64, i64, &[u8])> = rows.iter().map(|(k, w, v)| (*k, *w, &v[..])).collect();
crate::test_support::make_batch_bytes(schema, &rows)
}
#[test]
fn a_fold_under_churn_keeps_every_run_at_most_a_quarter_dead() {
let schema = crate::test_support::make_schema_pk_u64_payload_string();
const KEYS: u64 = 1000;
let mut set = RunSet::new(usize::MAX);
let mut latest = vec![0u64; KEYS as usize];
let initial: Vec<_> = (0..KEYS).map(|pk| (pk, 1, churn_value(pk, 0))).collect();
set.push(string_run(&schema, &initial), &schema);
let mut saw_dead = false;
for generation in 1..=300u64 {
let mut rows = Vec::new();
for i in 0..5 {
let pk = (i * 200 + generation * 7) % KEYS;
rows.push((pk, -1, churn_value(pk, latest[pk as usize])));
rows.push((pk, 1, churn_value(pk, generation)));
latest[pk as usize] = generation;
}
rows.sort_by(|a, b| (a.0, &a.2).cmp(&(b.0, &b.2)));
set.push(string_run(&schema, &rows), &schema);
for run in &set.runs {
saw_dead |= run.dead_heap() > 0;
assert!(
run.dead_heap() * 4 <= run.blob().len(),
"a stored run is {} of {} bytes dead",
run.dead_heap(),
run.blob().len()
);
}
}
assert!(
saw_dead,
"the churn must drive a fold that carries a heap with dead bytes"
);
set.fold(&schema);
let folded = &**set.runs.first().expect("every key survives");
let got: Vec<(u128, i64, Vec<u8>)> = (0..folded.len())
.map(|row| {
let s = gnitz_wire::payload_bytes(folded, row, 0).to_vec();
(folded.get_pk(row), folded.get_weight(row), s)
})
.collect();
let want: Vec<_> = (0..KEYS)
.map(|pk| (pk as u128, 1, churn_value(pk, latest[pk as usize])))
.collect();
assert_eq!(got, want);
}
#[test]
fn ascending_runs_fold_to_their_rows_in_order() {
use crate::test_support::{make_batch_raw, make_schema_u64_i64};
let schema = make_schema_u64_i64();
for (first_run, reach_back) in [(10u64, false), (10, true), (1000, false), (1000, true)] {
let mut set = RunSet::new(usize::MAX);
let mut pushed: Vec<Batch> = Vec::new();
let mut next = 0u64;
let mut push = |set: &mut RunSet, rows: Vec<(u64, i64, i64)>| {
let run = make_batch_raw(&schema, &rows).into_consolidated();
pushed.push(run.clone());
set.push(run, &schema);
};
for n in std::iter::once(first_run).chain(std::iter::repeat_n(10, 9)) {
push(&mut set, (next..next + n).map(|k| (k, 1, k as i64)).collect());
next += n;
}
if reach_back {
push(&mut set, vec![(3, -1, 3), (next, 1, 0)]);
}
set.fold(&schema);
let folded = &**set.runs.first().expect("rows survive");
assert!(folded.is_consolidated());
assert_eq!(zset_of(folded, &schema), zset_sum(&pushed, &schema));
assert_eq!(folded.len(), next as usize);
}
}