Skip to main content

quantile_sketch/
lib.rs

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/// Error returned when sketches with different configurations are merged.
31#[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
43/// Lazily allocated block
44struct 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            // Keep the selected sample as the origin for the next relative rank.
72            *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
93/// DDSketch with atomic concurrent inserts.
94pub 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    /// Creates a sketch with 1% relative error for values from zero through
104    /// `f64::MAX`.
105    pub fn new() -> Self {
106        Self::with_err_and_range(DEFAULT_ERROR, 0.0, f64::MAX)
107    }
108
109    /// Creates a sketch with the given relative error for the full positive
110    /// range.
111    ///
112    /// # Panics
113    ///
114    /// Panics if `err` is not strictly between zero and one or is too small to
115    /// produce distinct buckets.
116    pub fn with_err(err: f64) -> Self {
117        Self::with_err_and_range(err, 0.0, f64::MAX)
118    }
119
120    /// Creates a sketch with 1% relative error for the inclusive range
121    /// `min..=max`.
122    ///
123    /// # Panics
124    ///
125    /// Panics if the range is invalid or contains negative values.
126    pub fn with_range(min: f64, max: f64) -> Self {
127        Self::with_err_and_range(DEFAULT_ERROR, min, max)
128    }
129
130    /// Creates a sketch with the given relative error and inclusive range.
131    ///
132    /// # Panics
133    ///
134    /// Panics if `err` or the range is invalid.
135    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                // SAFETY: Every non-null pointer was created by
188                // `Box::into_raw` in `init_block`. Published pointers are never
189                // replaced and remain allocated until `self` is dropped.
190                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                // Another thread initialized this block first.
258                // SAFETY: `new` came from `Box::into_raw` above and was not
259                // published because the compare-exchange failed.
260                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        // SAFETY: `block` is non-null and points to a published `Block` that is
271        // not reclaimed until the sketch is dropped. The returned reference is
272        // tied to the borrow of the slot within that sketch.
273        unsafe { &*block }
274    }
275
276    /// Adds a value to the sketch.
277    ///
278    /// # Panics
279    ///
280    /// Panics if `val` is NaN or outside the configured range.
281    #[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    /// Adds the counts from a sketch with the same error and supported range.
288    ///
289    /// This is not an atomic snapshot: each source counter is read and added
290    /// once, so concurrent inserts may or may not be included and concurrent
291    /// queries may observe a partial merge.
292    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    /// Returns the configured inclusive minimum supported value.
310    pub fn min(&self) -> f64 {
311        self.min
312    }
313
314    /// Returns the configured inclusive maximum supported value.
315    pub fn max(&self) -> f64 {
316        self.max
317    }
318
319    /// Returns the lower `q`-quantile, or `None` when the sketch is empty.
320    /// Concurrent inserts may or may not be reflected in the result.
321    ///
322    /// # Panics
323    ///
324    /// Panics if `q` is NaN or outside `0.0..=1.0`.
325    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    /// Writes lower estimates for sorted `quantiles` to `output` without
342    /// allocating.
343    ///
344    /// Results are nondecreasing but are not an atomic snapshot of concurrent
345    /// inserts. Separate [`Self::quantile`] calls can return p50 > p99 if
346    /// inserts occur between them. An empty sketch fills `output` with `None`.
347    ///
348    /// # Panics
349    ///
350    /// Panics if the lengths differ or a quantile is unsorted, NaN, or outside
351    /// `0.0..=1.0`. Validation failures leave `output` unchanged.
352    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    /// Returns the estimated fraction of samples at or below `val`.
375    ///
376    /// The result is in `0.0..=1.0`, or `None` when the sketch is empty. Since
377    /// values in the same bucket are indistinguishable, this includes the
378    /// entire bucket containing `val`. Concurrent inserts may or may not be
379    /// reflected in the result.
380    ///
381    /// # Panics
382    ///
383    /// Panics if `val` is NaN or outside the configured range.
384    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                // SAFETY: Drop has exclusive access to the sketch. Each
434                // published pointer came from `Box::into_raw`, is unique to its
435                // slot, and has not previously been reclaimed.
436                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            // Miri injects permitted error into imprecise float operations, so
451            // equivalent estimates are not always bit-for-bit equal.
452            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}