#![cfg(feature = "scheduler")]
use hyphae::{CellMap, MapQuery, batch, traits::InnerJoinExt};
const PAIRS: usize = 8;
const ITERATIONS: i64 = 3_000;
#[test]
fn concurrent_inner_joins_never_settle_on_a_torn_value() {
hyphae::scheduler::set_wave_threshold_for_test(4);
let mut lefts = Vec::new();
let mut rights = Vec::new();
let mut outputs = Vec::new();
for _ in 0..PAIRS {
let l = CellMap::<String, i64>::new();
let r = CellMap::<String, i64>::new();
l.insert("k".into(), 0);
r.insert("k".into(), 0);
let out = l.clone().inner_join(r.clone()).materialize();
lefts.push(l);
rights.push(r);
outputs.push(out);
}
for it in 1..=ITERATIONS {
batch(|| {
for (index, (left, right)) in lefts.iter().zip(&rights).enumerate() {
let base = it
.saturating_mul(1000)
.saturating_add(i64::try_from(index).unwrap_or(i64::MAX));
left.insert("k".into(), base);
right.insert("k".into(), base.saturating_add(500));
}
});
for (index, output) in outputs.iter().enumerate() {
let base = it
.saturating_mul(1000)
.saturating_add(i64::try_from(index).unwrap_or(i64::MAX));
assert_eq!(
output.get_value(&"k".to_string()),
Some((base, base.saturating_add(500))),
"inner_join settled on a torn/stale value at iteration {it}, pair {index}"
);
}
}
}