use loom::{sync::Arc, thread};
use super::{assert_estimate_eq, ConcurrentDDSketch};
const ERRORS: [f64; 4] = [0.005, 0.01, 0.02, 0.05];
fn check_model(max_branches: usize, model: impl Fn() + Send + Sync + 'static) {
let mut builder = loom::model::Builder::new();
builder.max_branches = max_branches;
builder.check(model);
}
fn check_models(max_branches: usize, model: impl Fn(f64) + Copy + Send + Sync + 'static) {
for error in ERRORS {
check_model(max_branches, move || model(error));
}
}
#[test]
fn quantiles_are_independent_of_concurrent_insert_order() {
const QUANTILES: [f64; 5] = [0.0, 0.25, 0.5, 0.75, 1.0];
check_models(160_000, |error| {
let serial = ConcurrentDDSketch::with_err_and_range(error, 1.0, 4.0);
for value in [1.0, 2.0, 3.0, 2.0, 3.0, 4.0] {
serial.insert(value);
}
let expected = QUANTILES.map(|q| serial.quantile(q));
drop(serial);
let concurrent = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 4.0));
let a = Arc::clone(&concurrent);
let b = Arc::clone(&concurrent);
let first = thread::spawn(move || {
a.insert(1.0);
a.insert(2.0);
a.insert(3.0);
});
let second = thread::spawn(move || {
b.insert(2.0);
b.insert(3.0);
b.insert(4.0);
});
first.join().unwrap();
second.join().unwrap();
for (actual, expected) in QUANTILES.map(|q| concurrent.quantile(q)).into_iter().zip(expected) {
assert_estimate_eq(actual, expected);
}
});
}
#[test]
fn merges_can_initialize_the_same_block_concurrently() {
check_models(160_000, |error| {
let serial = ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0);
for value in [1.0, 1.2, 1.8, 2.0] {
serial.insert(value);
}
let expected = [serial.quantile(0.0), serial.quantile(1.0)];
drop(serial);
let merged = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0));
let low = ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0);
low.insert(1.0);
low.insert(1.2);
let high = ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0);
high.insert(1.8);
high.insert(2.0);
let a = Arc::clone(&merged);
let b = Arc::clone(&merged);
let first = thread::spawn(move || a.merge(&low).unwrap());
let second = thread::spawn(move || b.merge(&high).unwrap());
first.join().unwrap();
second.join().unwrap();
for (actual, expected) in [merged.quantile(0.0), merged.quantile(1.0)].into_iter().zip(expected) {
assert_estimate_eq(actual, expected);
}
assert_eq!(merged.percentile_rank(1.2), Some(0.5));
});
}
#[test]
fn merge_may_include_a_concurrent_source_insert() {
check_models(80_000, |error| {
let serial = ConcurrentDDSketch::with_err_and_range(error, 1.0, 1.0);
serial.insert(1.0);
let expected = serial.quantile(0.5);
drop(serial);
let source = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 1.0));
let merged = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 1.0));
let insert_source = Arc::clone(&source);
let merge_source = Arc::clone(&source);
let merge_target = Arc::clone(&merged);
let insert = thread::spawn(move || insert_source.insert(1.0));
let merge = thread::spawn(move || merge_target.merge(&merge_source).unwrap());
insert.join().unwrap();
merge.join().unwrap();
assert_estimate_eq(source.quantile(0.5), expected);
let result = merged.quantile(0.5);
if result.is_some() {
assert_estimate_eq(result, expected);
}
});
}
#[test]
fn queries_tolerate_a_concurrent_insert() {
check_models(80_000, |error| {
let sketch = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0));
sketch.insert(1.0);
let inserter = Arc::clone(&sketch);
let insert = thread::spawn(move || inserter.insert(2.0));
let quantile = sketch.quantile(1.0).unwrap();
assert!((1.0..=2.0).contains(&quantile));
let rank = sketch.percentile_rank(1.0).unwrap();
assert!(rank == 0.5 || rank == 1.0);
insert.join().unwrap();
let maximum = sketch.quantile(1.0).unwrap();
assert!((maximum - 2.0).abs() / 2.0 <= error * 1.001);
assert_eq!(sketch.percentile_rank(1.0), Some(0.5));
});
}
#[test]
fn bulk_quantiles_tolerate_a_concurrent_insert() {
const QUANTILES: [f64; 4] = [0.0, 0.5, 0.5, 1.0];
check_models(80_000, |error| {
let sketch = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0));
sketch.insert(1.0);
let inserter = Arc::clone(&sketch);
let insert = thread::spawn(move || inserter.insert(2.0));
let mut estimates = [None; 4];
sketch.quantiles(&QUANTILES, &mut estimates);
let estimates = estimates.map(Option::unwrap);
assert!(estimates.windows(2).all(|pair| pair[0] <= pair[1]));
assert!(estimates.iter().all(|value| (1.0..=2.0).contains(value)));
insert.join().unwrap();
});
}