1#![allow(rustdoc::bare_urls)]
2#![doc = include_str!("../README.md")]
3#![cfg_attr(not(any(feature = "std", test)), no_std)]
4
5extern crate alloc;
6
7use alloc::{boxed::Box, vec::Vec};
8use core::{fmt, iter::Peekable, ptr};
9
10#[cfg(feature = "loom")]
11use loom::sync::atomic::{AtomicPtr, AtomicU64, Ordering};
12#[cfg(not(feature = "loom"))]
13use portable_atomic::{AtomicPtr, AtomicU64, Ordering};
14
15#[cfg(all(test, not(feature = "loom")))]
16mod concurrency_tests;
17#[cfg(all(feature = "loom", test))]
18mod loom_tests;
19mod math;
20#[cfg(all(test, not(feature = "loom")))]
21mod property_tests;
22#[cfg(feature = "serde")]
23mod serde_impl;
24
25use math::CubicMapping;
26
27const BLOCK_SIZE: usize = 64;
28const DEFAULT_ERROR: f64 = 0.01;
29
30#[derive(Clone, Copy, Debug, Eq, PartialEq)]
32pub struct MergeError;
33
34impl fmt::Display for MergeError {
35 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
36 f.write_str("cannot merge sketches with different configurations")
37 }
38}
39
40#[cfg(feature = "std")]
41impl std::error::Error for MergeError {}
42
43struct Block {
45 counts: [AtomicU64; BLOCK_SIZE],
46}
47
48impl Block {
49 #[inline]
50 fn total(&self) -> u64 {
51 self.counts.iter().map(|count| count.load(Ordering::Relaxed)).sum()
52 }
53}
54
55impl Default for Block {
56 fn default() -> Self {
57 Self {
58 counts: core::array::from_fn(|_| AtomicU64::new(0)),
59 }
60 }
61}
62
63#[inline]
64fn bucket_at_rank(
65 buckets: &mut Peekable<impl Iterator<Item = ((usize, usize), u64)>>,
66 mut rank: u64,
67) -> Option<(usize, usize)> {
68 loop {
69 let (coord, remaining) = buckets.peek_mut()?;
70 if rank < *remaining {
71 *remaining -= rank;
73 return Some(*coord);
74 }
75 rank -= *remaining;
76 buckets.next();
77 }
78}
79
80#[inline]
81fn rank(q: f64, last_rank: u64) -> u64 {
82 ((q * last_rank as f64) as u64).min(last_rank)
83}
84
85fn required_bucket_count(min_index: i64, max_index: i64) -> usize {
86 max_index
87 .checked_sub(min_index)
88 .and_then(|count| count.checked_add(1))
89 .and_then(|count| usize::try_from(count).ok())
90 .expect("range requires too many buckets for this target")
91}
92
93pub struct ConcurrentDDSketch {
95 blocks: Box<[AtomicPtr<Block>]>,
96 mapping: CubicMapping,
97 min: f64,
98 max: f64,
99 min_index: i64,
100}
101
102impl ConcurrentDDSketch {
103 pub fn new() -> Self {
106 Self::with_err_and_range(DEFAULT_ERROR, 0.0, f64::MAX)
107 }
108
109 pub fn with_err(err: f64) -> Self {
117 Self::with_err_and_range(err, 0.0, f64::MAX)
118 }
119
120 pub fn with_range(min: f64, max: f64) -> Self {
127 Self::with_err_and_range(DEFAULT_ERROR, min, max)
128 }
129
130 pub fn with_err_and_range(err: f64, min: f64, max: f64) -> Self {
136 assert!(
137 min == 0.0 || min >= f64::MIN_POSITIVE,
138 "min must be zero or normal and positive"
139 );
140 assert!(min <= max, "min must not exceed max");
141 assert!(max < f64::INFINITY, "max must be finite");
142
143 let mapping = CubicMapping::new(err);
144 let min_index = if min == 0.0 {
145 mapping.index(f64::MIN_POSITIVE) - 1
146 } else {
147 mapping.index(min)
148 };
149 let max_index = if max < f64::MIN_POSITIVE {
150 min_index
151 } else {
152 mapping.index(max)
153 };
154 let bucket_count = required_bucket_count(min_index, max_index);
155 let block_count = bucket_count / BLOCK_SIZE + usize::from(bucket_count % BLOCK_SIZE != 0);
156 let blocks = (0..block_count)
157 .map(|_| AtomicPtr::new(ptr::null_mut()))
158 .collect::<Vec<_>>();
159 Self {
160 blocks: blocks.into(),
161 min_index,
162 mapping,
163 min,
164 max,
165 }
166 }
167
168 #[inline]
169 fn coord(&self, val: f64) -> (usize, usize) {
170 assert!(val >= self.min && val <= self.max, "value out of range");
171 let index = if val < f64::MIN_POSITIVE {
172 self.min_index
173 } else {
174 self.mapping.index(val)
175 };
176 let offset = (index - self.min_index) as usize;
177 (offset / BLOCK_SIZE, offset % BLOCK_SIZE)
178 }
179
180 #[inline]
181 fn allocated_blocks(&self) -> impl DoubleEndedIterator<Item = (usize, &Block)> {
182 self.blocks.iter().enumerate().filter_map(|(index, slot)| {
183 let block = slot.load(Ordering::Acquire);
184 if block.is_null() {
185 None
186 } else {
187 Some((index, unsafe { &*block }))
191 }
192 })
193 }
194
195 #[inline]
196 fn buckets(&self) -> impl DoubleEndedIterator<Item = ((usize, usize), u64)> + '_ {
197 self.allocated_blocks().flat_map(|(block_index, block)| {
198 block
199 .counts
200 .iter()
201 .enumerate()
202 .map(move |(count_index, count)| ((block_index, count_index), count.load(Ordering::Relaxed)))
203 })
204 }
205
206 #[inline]
207 fn last_rank(&self) -> Option<u64> {
208 self.allocated_blocks()
209 .map(|(_, block)| block.total())
210 .sum::<u64>()
211 .checked_sub(1)
212 }
213
214 #[inline]
215 fn value_at_coord(&self, block_index: usize, count_index: usize) -> f64 {
216 let index = self.min_index + (block_index * BLOCK_SIZE + count_index) as i64;
217 if self.min == 0.0 && index == self.min_index {
218 0.0
219 } else {
220 self.mapping.value(index).clamp(self.min, self.max)
221 }
222 }
223
224 fn write_ranked_quantiles(
225 &self,
226 ranks: impl Iterator<Item = (usize, u64)>,
227 buckets: impl Iterator<Item = ((usize, usize), u64)>,
228 output: &mut [Option<f64>],
229 ) {
230 let mut buckets = buckets.peekable();
231 let mut prev_rank = 0;
232 for (output_index, rank) in ranks {
233 let (block_index, count_index) =
234 bucket_at_rank(&mut buckets, rank - prev_rank).expect("bucket totals must match bucket counts");
235 output[output_index] = Some(self.value_at_coord(block_index, count_index));
236 prev_rank = rank;
237 }
238 }
239
240 fn write_quantiles(&self, quantiles: &[f64], last_rank: u64, reverse: bool, output: &mut [Option<f64>]) {
241 if reverse {
242 let r = quantiles.iter().map(|&q| last_rank - rank(q, last_rank)).enumerate();
243 self.write_ranked_quantiles(r.rev(), self.buckets().rev(), output);
244 } else {
245 let r = quantiles.iter().map(|&q| rank(q, last_rank)).enumerate();
246 self.write_ranked_quantiles(r, self.buckets(), output);
247 }
248 }
249
250 #[cold]
251 #[inline(never)]
252 fn init_block(slot: &AtomicPtr<Block>) -> *mut Block {
253 let new = Box::into_raw(Box::<Block>::default());
254 match slot.compare_exchange(ptr::null_mut(), new, Ordering::Release, Ordering::Acquire) {
255 Ok(_) => new,
256 Err(existing) => {
257 unsafe { drop(Box::from_raw(new)) };
261 existing
262 }
263 }
264 }
265
266 #[inline]
267 fn get_or_init_block(slot: &AtomicPtr<Block>) -> &Block {
268 let block = slot.load(Ordering::Acquire);
269 let block = if block.is_null() { Self::init_block(slot) } else { block };
270 unsafe { &*block }
274 }
275
276 #[inline]
282 pub fn insert(&self, val: f64) {
283 let (block_index, count_index) = self.coord(val);
284 Self::get_or_init_block(&self.blocks[block_index]).counts[count_index].fetch_add(1, Ordering::Relaxed);
285 }
286
287 pub fn merge(&self, other: &Self) -> Result<(), MergeError> {
293 if self.mapping.error() != other.mapping.error() || self.min != other.min || self.max != other.max {
294 return Err(MergeError);
295 }
296
297 for (block_index, source) in other.allocated_blocks() {
298 let target = Self::get_or_init_block(&self.blocks[block_index]);
299 for (target, source) in target.counts.iter().zip(&source.counts) {
300 let count = source.load(Ordering::Relaxed);
301 if count != 0 {
302 target.fetch_add(count, Ordering::Relaxed);
303 }
304 }
305 }
306 Ok(())
307 }
308
309 pub fn min(&self) -> f64 {
311 self.min
312 }
313
314 pub fn max(&self) -> f64 {
316 self.max
317 }
318
319 pub fn quantile(&self, q: f64) -> Option<f64> {
326 assert!((0.0..=1.0).contains(&q), "quantile must be between 0 and 1");
327
328 let last_rank = self.last_rank()?;
329 let rank = rank(q, last_rank);
330
331 let found = if rank <= last_rank - rank {
332 bucket_at_rank(&mut self.buckets().peekable(), rank)
333 } else {
334 bucket_at_rank(&mut self.buckets().rev().peekable(), last_rank - rank)
335 };
336
337 let (block_index, count_index) = found.expect("bucket totals must match bucket counts");
338 Some(self.value_at_coord(block_index, count_index))
339 }
340
341 pub fn quantiles(&self, quantiles: &[f64], output: &mut [Option<f64>]) {
353 assert_eq!(output.len(), quantiles.len(), "output length must match quantiles");
354 for &q in quantiles {
355 assert!((0.0..=1.0).contains(&q), "quantile must be between 0 and 1");
356 }
357 assert!(
358 quantiles.windows(2).all(|pair| pair[0] <= pair[1]),
359 "quantiles must be sorted"
360 );
361
362 output.fill(None);
363 let Some((&first, &last)) = quantiles.first().zip(quantiles.last()) else {
364 return;
365 };
366 let Some(last_rank) = self.last_rank() else {
367 return;
368 };
369
370 let reverse = rank(last, last_rank) > last_rank - rank(first, last_rank);
371 self.write_quantiles(quantiles, last_rank, reverse, output);
372 }
373
374 pub fn percentile_rank(&self, val: f64) -> Option<f64> {
385 let (target_block, target_count) = self.coord(val);
386 let mut cumulative = 0_u64;
387 let mut total = 0_u64;
388
389 for (block_index, block) in self.allocated_blocks() {
390 if block_index == target_block {
391 for (count_index, count) in block.counts.iter().enumerate() {
392 let count = count.load(Ordering::Relaxed);
393 total += count;
394 if count_index <= target_count {
395 cumulative += count;
396 }
397 }
398 continue;
399 }
400
401 let count = block.total();
402 total += count;
403 if block_index < target_block {
404 cumulative += count;
405 }
406 }
407
408 (total != 0).then_some(cumulative as f64 / total as f64)
409 }
410}
411
412impl Default for ConcurrentDDSketch {
413 fn default() -> Self {
414 Self::new()
415 }
416}
417
418impl fmt::Debug for ConcurrentDDSketch {
419 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
420 f.debug_struct("ConcurrentDDSketch")
421 .field("error", &self.mapping.error())
422 .field("min", &self.min)
423 .field("max", &self.max)
424 .finish_non_exhaustive()
425 }
426}
427
428impl Drop for ConcurrentDDSketch {
429 fn drop(&mut self) {
430 for slot in self.blocks.iter_mut() {
431 let block = slot.load(Ordering::Relaxed);
432 if !block.is_null() {
433 unsafe { drop(Box::from_raw(block)) };
437 }
438 }
439 }
440}
441
442#[cfg(test)]
443fn assert_estimate_eq(actual: Option<f64>, expected: Option<f64>) {
444 #[cfg(not(miri))]
445 assert_eq!(actual, expected);
446
447 #[cfg(miri)]
448 match (actual, expected) {
449 (Some(actual), Some(expected)) => {
450 let tolerance = 8_192.0 * f64::EPSILON * expected.abs().max(1.0);
453 assert!((actual - expected).abs() <= tolerance, "{actual} != {expected}");
454 }
455 _ => assert_eq!(actual, expected),
456 }
457}
458
459#[cfg(all(test, not(feature = "loom")))]
460mod tests {
461 use super::*;
462
463 fn assert_panics(f: impl FnOnce()) {
464 assert!(std::panic::catch_unwind(std::panic::AssertUnwindSafe(f)).is_err());
465 }
466
467 fn nonzero_bins(sketch: &ConcurrentDDSketch) -> Vec<(usize, u64)> {
468 sketch
469 .allocated_blocks()
470 .flat_map(|(block_index, block)| {
471 block.counts.iter().enumerate().filter_map(move |(count_index, count)| {
472 let count = count.load(Ordering::Relaxed);
473 (count != 0).then_some((block_index * BLOCK_SIZE + count_index, count))
474 })
475 })
476 .collect()
477 }
478
479 #[cfg(feature = "serde")]
480 #[derive(Debug, PartialEq)]
481 struct SketchState {
482 min: f64,
483 max: f64,
484 quantiles: [Option<f64>; 3],
485 }
486
487 #[cfg(feature = "serde")]
488 impl<'de> serde::Deserialize<'de> for SketchState {
489 fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
490 let sketch = ConcurrentDDSketch::deserialize(deserializer)?;
491 Ok(Self {
492 min: sketch.min(),
493 max: sketch.max(),
494 quantiles: [sketch.quantile(0.0), sketch.quantile(0.5), sketch.quantile(1.0)],
495 })
496 }
497 }
498
499 #[cfg(feature = "serde")]
500 fn serde_tokens(max_offset: u64) -> Vec<serde_test::Token> {
501 use serde_test::Token;
502
503 vec![
504 Token::Struct { name: "Repr", len: 5 },
505 Token::Str("version"),
506 Token::U8(1),
507 Token::Str("error"),
508 Token::F64(0.02),
509 Token::Str("min"),
510 Token::F64(0.0),
511 Token::Str("max"),
512 Token::F64(10.0),
513 Token::Str("bins"),
514 Token::Seq { len: Some(2) },
515 Token::Tuple { len: 2 },
516 Token::U64(0),
517 Token::U64(2),
518 Token::TupleEnd,
519 Token::Tuple { len: 2 },
520 Token::U64(max_offset),
521 Token::U64(1),
522 Token::TupleEnd,
523 Token::SeqEnd,
524 Token::StructEnd,
525 ]
526 }
527
528 #[test]
529 fn computes_number_of_blocks() {
530 for (min, max, expected) in [(1.0, 1.0, 1), (0.001, 1_000.0, 11), (1.0e-9, 1.0e9, 33)] {
531 assert_eq!(ConcurrentDDSketch::with_range(min, max).blocks.len(), expected);
532 }
533 }
534
535 #[test]
536 fn constructors_set_expected_ranges() {
537 let default = ConcurrentDDSketch::new();
538 assert_eq!((default.min, default.max), (0.0, f64::MAX));
539 assert_eq!(default.mapping.index(default.max) - default.min_index + 1, 71_609);
540 assert_eq!(default.blocks.len(), 1_119);
541 assert!(default.blocks.iter().all(|slot| slot.load(Ordering::Relaxed).is_null()));
542
543 for (sketch, error, range) in [
544 (ConcurrentDDSketch::with_err(0.02), 0.02, (0.0, f64::MAX)),
545 (ConcurrentDDSketch::with_range(1.0, 10.0), DEFAULT_ERROR, (1.0, 10.0)),
546 (
547 ConcurrentDDSketch::with_err_and_range(0.02, 1.0, 10.0),
548 0.02,
549 (1.0, 10.0),
550 ),
551 ] {
552 assert_eq!(sketch.mapping.index(10.0), CubicMapping::new(error).index(10.0));
553 assert_eq!((sketch.min, sketch.max), range);
554 }
555 }
556
557 #[test]
558 fn rejects_invalid_configurations() {
559 for (error, min, max) in [
560 (0.0, 0.0, 1.0),
561 (1.0, 0.0, 1.0),
562 (f64::NAN, 0.0, 1.0),
563 (f64::EPSILON / 4.0, 1.0, 1.0),
564 (DEFAULT_ERROR, -1.0, 1.0),
565 (DEFAULT_ERROR, f64::from_bits(1), 1.0),
566 (DEFAULT_ERROR, 2.0, 1.0),
567 (DEFAULT_ERROR, 0.0, f64::INFINITY),
568 (DEFAULT_ERROR, 0.0, f64::NAN),
569 ] {
570 assert_panics(|| drop(ConcurrentDDSketch::with_err_and_range(error, min, max)));
571 }
572 }
573
574 #[test]
575 #[should_panic(expected = "range requires too many buckets for this target")]
576 fn rejects_bucket_counts_that_cannot_be_indexed() {
577 required_bucket_count(i64::MIN, i64::MAX);
578 }
579
580 #[test]
581 fn handles_single_value_and_extreme_ranges() {
582 let zero = ConcurrentDDSketch::with_range(0.0, 0.0);
583 zero.insert(0.0);
584 assert_eq!(
585 [zero.quantile(0.0), zero.quantile(0.5), zero.quantile(1.0)],
586 [Some(0.0); 3]
587 );
588
589 let single = ConcurrentDDSketch::with_range(1.0, 1.0);
590 single.insert(1.0);
591 assert_eq!(
592 [single.quantile(0.0), single.quantile(0.5), single.quantile(1.0)],
593 [Some(1.0); 3]
594 );
595
596 let full = ConcurrentDDSketch::new();
597 for value in [0.0, f64::MIN_POSITIVE, f64::MAX] {
598 full.insert(value);
599 }
600 assert_eq!(full.quantile(0.0), Some(0.0));
601 let estimate = full.quantile(1.0).unwrap();
602 assert!((estimate - f64::MAX).abs() / f64::MAX <= DEFAULT_ERROR);
603 assert_eq!(full.percentile_rank(f64::MAX), Some(1.0));
604 }
605
606 #[test]
607 fn zero_and_underflow_have_their_own_bucket() {
608 let sketch = ConcurrentDDSketch::new();
609 sketch.insert(0.0);
610 sketch.insert(-0.0);
611 sketch.insert(f64::from_bits(1));
612 sketch.insert(f64::MIN_POSITIVE);
613 sketch.insert(1.0);
614 assert_eq!(sketch.quantile(0.0), Some(0.0));
615 assert_eq!(sketch.quantile(0.5), Some(0.0));
616 assert!(sketch.quantile(0.75).unwrap() > 0.0);
617 }
618
619 #[test]
620 fn inserts_into_the_expected_bucket() {
621 let sketch = ConcurrentDDSketch::with_err_and_range(0.01, 0.001, 1_000.0);
622 let (block_index, count_index) = sketch.coord(1.0);
623
624 sketch.insert(1.0);
625 sketch.insert(1.0);
626
627 let block = sketch.blocks[block_index].load(Ordering::Acquire);
628 let block = unsafe { &*block };
629 assert_eq!(block.counts[count_index].load(Ordering::Relaxed), 2);
630 }
631
632 #[test]
633 fn computes_percentile_ranks() {
634 let sketch = ConcurrentDDSketch::with_range(1.0, 4.0);
635 assert_eq!(sketch.percentile_rank(2.0), None);
636
637 for value in [1.0, 2.0, 2.0, 4.0] {
638 sketch.insert(value);
639 }
640
641 assert_eq!(sketch.percentile_rank(1.0), Some(0.25));
642 assert_eq!(sketch.percentile_rank(2.0), Some(0.75));
643 assert_eq!(sketch.percentile_rank(3.0), Some(0.75));
644 assert_eq!(sketch.percentile_rank(4.0), Some(1.0));
645 }
646
647 #[test]
648 fn computes_multiple_sorted_quantiles() {
649 let sketch = ConcurrentDDSketch::with_range(0.0, 4.0);
650 sketch.quantiles(&[], &mut []);
651 let mut estimates = [Some(1.0); 3];
652 sketch.quantiles(&[0.0, 0.5, 1.0], &mut estimates);
653 assert_eq!(estimates, [None; 3]);
654
655 for value in [0.0, 1.0, 2.0, 3.0, 4.0] {
656 sketch.insert(value);
657 }
658
659 let quantiles = [0.0, 0.5, 0.5, 0.75, 1.0];
660 let mut estimates = [None; 5];
661 sketch.quantiles(&quantiles, &mut estimates);
662 for (&q, estimate) in quantiles.iter().zip(estimates) {
663 assert_estimate_eq(estimate, sketch.quantile(q));
664 }
665
666 let high_quantiles = [0.75, 0.99, 0.99, 1.0];
667 let mut high_estimates = [None; 4];
668 sketch.quantiles(&high_quantiles, &mut high_estimates);
669 for (&q, estimate) in high_quantiles.iter().zip(high_estimates) {
670 assert_estimate_eq(estimate, sketch.quantile(q));
671 }
672
673 let mut ordered = [None; 4];
674 sketch.quantiles(&[0.0, 0.5, 0.99, 1.0], &mut ordered);
675 assert!(ordered.windows(2).all(|pair| pair[0] <= pair[1]));
676 }
677
678 #[test]
679 fn rejects_out_of_range_inputs() {
680 let sketch = ConcurrentDDSketch::with_range(1.0, 10.0);
681 for value in [0.0, 11.0, f64::NAN] {
682 assert_panics(|| sketch.insert(value));
683 assert_panics(|| {
684 sketch.percentile_rank(value);
685 });
686 }
687 for quantile in [-f64::EPSILON, 1.0 + f64::EPSILON, f64::NAN] {
688 assert_panics(|| {
689 sketch.quantile(quantile);
690 });
691 assert_panics(|| sketch.quantiles(&[0.5, quantile], &mut [None; 2]));
692 }
693 assert_panics(|| sketch.quantiles(&[0.9, 0.5], &mut [None; 2]));
694 assert_panics(|| sketch.quantiles(&[0.5], &mut [None; 2]));
695 assert_eq!(sketch.quantile(0.5), None);
696 }
697
698 #[test]
699 fn reports_configured_min_and_max() {
700 let default = ConcurrentDDSketch::new();
701 assert_eq!(default.min(), 0.0);
702 assert_eq!(default.max(), f64::MAX);
703
704 let bounded = ConcurrentDDSketch::with_range(1.0, 10.0);
705 assert_eq!(bounded.min(), 1.0);
706 assert_eq!(bounded.max(), 10.0);
707 }
708
709 #[test]
710 fn merges_compatible_sketches() {
711 let merged = ConcurrentDDSketch::with_range(0.0, 1_000.0);
712 let empty = ConcurrentDDSketch::with_range(0.0, 1_000.0);
713 merged.merge(&empty).unwrap();
714 assert!(merged.blocks.iter().all(|slot| slot.load(Ordering::Relaxed).is_null()));
715
716 merged.insert(1.0);
717 let other = ConcurrentDDSketch::with_range(0.0, 1_000.0);
718 for value in [0.0, f64::from_bits(1), 10.0, 1_000.0] {
719 other.insert(value);
720 }
721
722 let expected = ConcurrentDDSketch::with_range(0.0, 1_000.0);
723 for value in [0.0, f64::from_bits(1), 1.0, 10.0, 1_000.0] {
724 expected.insert(value);
725 }
726
727 merged.merge(&other).unwrap();
728 assert_eq!(nonzero_bins(&merged), nonzero_bins(&expected));
729 for q in [0.0, 0.25, 0.5, 0.75, 1.0] {
730 assert_estimate_eq(merged.quantile(q), expected.quantile(q));
731 }
732 assert_eq!(other.percentile_rank(0.0), Some(0.5));
733
734 let self_merge = ConcurrentDDSketch::with_range(1.0, 2.0);
735 self_merge.insert(1.0);
736 self_merge.insert(2.0);
737 self_merge.merge(&self_merge).unwrap();
738 assert_eq!(self_merge.percentile_rank(1.0), Some(0.5));
739 }
740
741 #[test]
742 fn rejects_incompatible_merges() {
743 let sketch = ConcurrentDDSketch::with_range(1.0, 10.0);
744 for other in [
745 ConcurrentDDSketch::with_err_and_range(0.02, 1.0, 10.0),
746 ConcurrentDDSketch::with_range(0.0, 10.0),
747 ConcurrentDDSketch::with_range(1.0, 20.0),
748 ] {
749 assert_eq!(sketch.merge(&other), Err(MergeError));
750 }
751 assert_eq!(sketch.quantile(0.5), None);
752 }
753
754 #[test]
755 fn debug_reports_configuration() {
756 let sketch = ConcurrentDDSketch::with_err_and_range(0.02, 1.0, 10.0);
757 assert_eq!(
758 format!("{sketch:?}"),
759 "ConcurrentDDSketch { error: 0.02, min: 1.0, max: 10.0, .. }"
760 );
761 }
762
763 #[cfg(feature = "serde")]
764 #[test]
765 fn serde_preserves_state() {
766 use serde_test::{assert_de_tokens, assert_ser_tokens};
767
768 let sketch = ConcurrentDDSketch::with_err_and_range(0.02, 0.0, 10.0);
769 for value in [0.0, 0.0, 10.0] {
770 sketch.insert(value);
771 }
772 let max_offset = (sketch.mapping.index(10.0) - sketch.min_index) as u64;
773 let tokens = serde_tokens(max_offset);
774 let state = SketchState {
775 min: sketch.min(),
776 max: sketch.max(),
777 quantiles: [sketch.quantile(0.0), sketch.quantile(0.5), sketch.quantile(1.0)],
778 };
779
780 assert_ser_tokens(&sketch, &tokens);
781 assert_de_tokens(&state, &tokens);
782 }
783
784 #[cfg(feature = "serde")]
785 #[test]
786 fn serde_rejects_invalid_state() {
787 use serde_test::{assert_de_tokens_error, Token};
788
789 let sketch = ConcurrentDDSketch::with_err_and_range(0.02, 0.0, 10.0);
790 let max_offset = (sketch.mapping.index(10.0) - sketch.min_index) as u64;
791 let tokens = serde_tokens(max_offset);
792 for (index, token, error) in [
793 (2, Token::U8(2), "unsupported quantile-sketch version"),
794 (
795 4,
796 Token::F64(0.0),
797 "relative error must be representable and between 0 and 1",
798 ),
799 (
800 4,
801 Token::F64(f64::EPSILON / 4.0),
802 "relative error must be representable and between 0 and 1",
803 ),
804 (6, Token::F64(f64::from_bits(1)), "invalid sketch range"),
805 (8, Token::F64(f64::INFINITY), "invalid sketch range"),
806 (13, Token::U64(0), "invalid sketch bin"),
807 (16, Token::U64(max_offset + 1), "invalid sketch bin"),
808 (16, Token::U64(0), "duplicate sketch bin"),
809 ] {
810 let mut invalid = tokens.clone();
811 invalid[index] = token;
812 assert_de_tokens_error::<ConcurrentDDSketch>(&invalid, error);
813 }
814 }
815
816 #[test]
817 fn quantiles_have_expected_relative_error() {
818 const ERROR: f64 = 0.01;
819 #[cfg(miri)]
820 const SAMPLES: usize = 8_000;
821 #[cfg(not(miri))]
822 const SAMPLES: usize = 100_000;
823 #[cfg(miri)]
824 const CHECKS: usize = 801;
825 #[cfg(not(miri))]
826 const CHECKS: usize = 1_001;
827
828 let sketch = ConcurrentDDSketch::with_err_and_range(ERROR, 1.0, SAMPLES as f64);
829 assert_eq!(sketch.quantile(0.5), None);
830 for value in 1..=SAMPLES {
831 sketch.insert(value as f64);
832 }
833
834 let mut error_sum = 0.0;
835 let mut max_error = 0.0_f64;
836 for sample in 0..CHECKS {
837 let q = sample as f64 / (CHECKS - 1) as f64;
838 let exact = 1.0 + (q * (SAMPLES - 1) as f64).floor();
839 let estimate = sketch.quantile(q).unwrap();
840 let error = (estimate - exact).abs() / exact;
841 error_sum += error;
842 max_error = max_error.max(error);
843 }
844
845 let average_error = error_sum / CHECKS as f64;
846 assert!(average_error <= ERROR * 0.55, "average error: {average_error}");
847 assert!(max_error <= ERROR * 1.001, "maximum error: {max_error}");
848 }
849}