quantile-sketch 0.1.0

A fast, concurrent DDSketch for relative-error quantiles.
Documentation
use alloc::vec::Vec;

use serde::{de::Error, Deserialize, Deserializer, Serialize, Serializer};

use crate::{math::CubicMapping, ConcurrentDDSketch, Ordering, BLOCK_SIZE};

const VERSION: u8 = 1;

#[derive(Serialize, Deserialize)]
struct Repr {
    version: u8,
    error: f64,
    min: f64,
    max: f64,
    bins: Vec<(u64, u64)>,
}

impl Serialize for ConcurrentDDSketch {
    fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
        let mut bins = Vec::new();
        for (block_index, block) in self.allocated_blocks() {
            for (count_index, count) in block.counts.iter().enumerate() {
                let count = count.load(Ordering::Relaxed);
                if count != 0 {
                    bins.push(((block_index * BLOCK_SIZE + count_index) as u64, count));
                }
            }
        }
        Repr {
            version: VERSION,
            error: self.mapping.error(),
            min: self.min,
            max: self.max,
            bins,
        }
        .serialize(serializer)
    }
}

impl<'de> Deserialize<'de> for ConcurrentDDSketch {
    fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
        let repr = Repr::deserialize(deserializer)?;
        if repr.version != VERSION {
            return Err(D::Error::custom("unsupported quantile-sketch version"));
        }
        if !CubicMapping::is_valid_error(repr.error) {
            return Err(D::Error::custom(
                "relative error must be representable and between 0 and 1",
            ));
        }
        if !(repr.min == 0.0 || repr.min >= f64::MIN_POSITIVE) || repr.min > repr.max || !repr.max.is_finite() {
            return Err(D::Error::custom("invalid sketch range"));
        }

        let sketch = Self::with_err_and_range(repr.error, repr.min, repr.max);
        let max_offset = if repr.max < f64::MIN_POSITIVE {
            0
        } else {
            (sketch.mapping.index(repr.max) - sketch.min_index) as usize
        };
        for (offset, count) in repr.bins {
            let offset = usize::try_from(offset).map_err(|_| D::Error::custom("invalid sketch bin"))?;
            if count == 0 || offset > max_offset {
                return Err(D::Error::custom("invalid sketch bin"));
            }
            let block = Self::get_or_init_block(&sketch.blocks[offset / BLOCK_SIZE]);
            let count_slot = &block.counts[offset % BLOCK_SIZE];
            if count_slot.swap(count, Ordering::Relaxed) != 0 {
                return Err(D::Error::custom("duplicate sketch bin"));
            }
        }
        Ok(sketch)
    }
}