Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 22 additions & 12 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,27 +6,37 @@ All significant changes to this project will be documented in this file.

### Breaking changes

* `FrequentItemsSketch::is_empty` now checks whether the total stream weight is zero. A sketch whose counters were all removed by a purge remains non-empty; use `num_active_items() == 0` to check whether any items are retained.
* Remove `BloomFilter::invert`. Bit inversion has no sound set-membership interpretation; use the new `BloomFilter::difference` for the approximate set-difference (A NOT B) use case it was meant to serve.
* Move `SearchCriteria` from `req` to `common` and remove its `Default` implementation. Import `datasketches::common::SearchCriteria` and explicitly choose `Inclusive` or `Exclusive` for each query.
* `BloomFilter::invert` is removed; use `BloomFilter::difference` for approximate A-not-B. The result excludes items in the right filter, but hash collisions can also remove items unique to the left filter.
* `FrequentItemsSketch::is_empty` now returns `false` when the stream weight is nonzero, even if no items are retained. Use `num_active_items() == 0` to test for zero retained items.
* `ReqSketch` queries now use `datasketches::common::SearchCriteria` instead of `datasketches::req::SearchCriteria`. `SearchCriteria` no longer implements `Default`; explicitly choose `Inclusive` or `Exclusive`.

### New features

* Add `BloomFilter::difference` for approximate set difference: the result excludes the other filter's items exactly, while items unique to the left filter are kept unless their hash positions collide with the right filter.
* Add KLL sketches behind the `kll` feature, with rank, quantile, PMF, and CDF queries, merging, totally ordered custom item types, a `KllFloat` adapter for non-NaN floating-point values, and serialization.
* `KllSketch` is now available behind the `kll` feature, with rank, quantile, PMF, and CDF queries, merging, serialization, custom ordered item types, and a `KllFloat` adapter for non-NaN floating-point values.

### Improvements

* The crate no longer has any runtime dependencies. The `kll` and `req` features previously pulled in `rand`; compaction now draws its coin from an in-tree generator.
* Improve truncated-input diagnostics across sketch deserializers.
* Improve hash-backed sketch update performance for integer and raw-byte inputs.
* Improve Bloom filter membership-and-insert performance and simplify Theta-family hash table thresholds.
* `BloomFilter::insert` is faster for integer and raw-byte inputs. `BloomFilter::contains_and_insert` is also faster when checking already-present integer values.
* `CountMinSketch` updates are faster for integer and raw-byte inputs.
* `CpcSketch` updates are faster for integer and raw-byte inputs.
* `FrequentItemsSketch` updates are faster for integer and raw-byte keys.
* `HllSketch` updates are faster for integer and raw-byte inputs.
* `ThetaSketch` updates are faster for integer and raw-byte inputs.
* `TupleSketch` updates are faster for integer and raw-byte inputs.
* Library-wide: the crate no longer depends on `rand` and has no runtime dependencies.
* Library-wide: sketch deserializers report clearer errors for truncated input.

### Bug fixes

* `FrequentItemsSketch` updates and merges now panic before modifying the sketch if the total stream weight would overflow, including in release builds. Deserialization rejects non-empty images with zero stream weight or item weights whose sum exceeds the declared stream weight.
* Fix T-Digest `merge` so it preserves `min`/`max` from the other digest instead of re-deriving them from centroid means after compression.
* T-Digest deserialization now rejects unknown or conflicting flags, reversed extrema, out-of-range values, unsorted centroids, and non-empty images without stored values.
* `CountMinSketch` updates now panic and merges return `InvalidArgument` if the total absolute weight would exceed the counter type's maximum. Both leave the sketch unchanged, including in release builds.
* `CountMinSketch::upper_bound` now clamps to the counter type's maximum instead of overflowing.
* `CountMinSketch` deserialization now returns `InvalidData` if the total absolute weight is negative or any counter's magnitude exceeds it.
* `FrequentItemsSketch` updates and merges now panic without changing the sketch if the total stream weight would overflow, including in release builds.
* `FrequentItemsSketch` deserialization now returns `InvalidData` if a non-empty image declares zero stream weight or the item weights sum to more than the declared stream weight.
* `ReqSketch` updates now panic and merges return `InvalidArgument` if the stream weight would exceed `u64::MAX`. Both leave the sketch unchanged, including in release builds.
* `TDigestMut` updates and merges now panic without changing the digest if the total weight would exceed `u64::MAX`, including in release builds.
* `TDigestMut::merge` now preserves the true minimum and maximum from both inputs, including compressed digests.
* `TDigest` and `TDigestMut` deserialization now returns `InvalidData` for invalid flags or extrema, out-of-range values, unsorted centroids, or non-empty images with no stored values.

## v0.5.0 (2026-09-04)

Expand Down
45 changes: 39 additions & 6 deletions datasketches/src/countmin/sketch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,9 @@ pub struct CountMinSketch<T: CountMinValue> {
num_buckets: u32,
seed: u64,
seed_hash: u16,
// Every bucket satisfies |count| <= total_weight, so a checked total also bounds bucket
// additions during update/merge: |a + b| <= |a| + |b|. Deserialization validates this
// invariant; unsigned halving and decay preserve it.
total_weight: T,
counts: Vec<T>,
hash_seeds: Vec<u64>,
Expand Down Expand Up @@ -119,7 +122,7 @@ impl<T: CountMinValue> CountMinSketch<T> {
self.seed
}

/// Returns the total weight inserted into the sketch.
/// Returns the sum of absolute update weights, scaled by any halving or decay.
pub fn total_weight(&self) -> T {
self.total_weight
}
Expand Down Expand Up @@ -177,6 +180,10 @@ impl<T: CountMinValue> CountMinSketch<T> {

/// Updates the sketch with a single occurrence of the item.
///
/// # Panics
///
/// Panics without modifying the sketch if the total absolute weight would exceed `T::MAX`.
///
/// # Examples
///
/// ```
Expand All @@ -192,6 +199,11 @@ impl<T: CountMinValue> CountMinSketch<T> {

/// Updates the sketch with the given item and weight.
///
/// # Panics
///
/// Panics without modifying the sketch if the absolute weight or the total absolute weight
/// cannot be represented by `T`.
///
/// # Examples
///
/// ```
Expand All @@ -205,8 +217,10 @@ impl<T: CountMinValue> CountMinSketch<T> {
if weight == T::ZERO {
return;
}
let abs_weight = weight.abs();
self.total_weight = self.total_weight + abs_weight;
self.total_weight = weight
.checked_abs()
.and_then(|weight| self.total_weight.checked_add(weight))
.expect("total absolute weight overflow");
let num_buckets = self.num_buckets as usize;
for (row, seed) in self.hash_seeds.iter().enumerate() {
let bucket = self.bucket_index(&item, *seed);
Expand Down Expand Up @@ -246,17 +260,20 @@ impl<T: CountMinValue> CountMinSketch<T> {
}

/// Returns the upper bound on the true frequency of the given item.
///
/// Clamps the bound to `T::MAX` if adding the error would overflow.
pub fn upper_bound<I: Hash>(&self, item: I) -> T {
let estimate = self.estimate(item);
let error = self.total_weight.scale(self.relative_error());
estimate + error
estimate.checked_add(error).unwrap_or(T::MAX)
}

/// Merges another sketch into this one.
///
/// # Errors
///
/// Returns an error if the sketches have different numbers of hashes, bucket counts, or seeds.
/// Returns an error without modifying the sketch if the sketches have different numbers of
/// hashes, bucket counts, or seeds, or their combined total absolute weight exceeds `T::MAX`.
///
/// # Examples
///
Expand All @@ -281,10 +298,13 @@ impl<T: CountMinValue> CountMinSketch<T> {
"Count-Min sketches must have matching numbers of hashes, bucket counts, and seeds",
));
}
self.total_weight = self
.total_weight
.checked_add(other.total_weight)
.ok_or_else(|| Error::invalid_argument("total absolute weight overflow"))?;
for (count, other_count) in self.counts.iter_mut().zip(&other.counts) {
*count = *count + *other_count;
}
self.total_weight = self.total_weight + other.total_weight;
Ok(())
}

Expand Down Expand Up @@ -445,8 +465,21 @@ impl<T: CountMinValue> CountMinSketch<T> {
}

sketch.total_weight = read_value(&mut cursor, "total_weight")?;
if sketch.total_weight < T::ZERO {
return Err(Error::deserial(
"total absolute weight must be non-negative",
));
}
for count in &mut sketch.counts {
*count = read_value(&mut cursor, "counts")?;
if count
.checked_abs()
.is_none_or(|weight| weight > sketch.total_weight)
{
return Err(Error::deserial(
"counter magnitude exceeds total absolute weight",
));
}
}
Ok(sketch)
}
Expand Down
21 changes: 16 additions & 5 deletions datasketches/src/countmin/value.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,8 @@ mod private {
const ONE: Self;
const MAX: Self;

fn abs(self) -> Self;
fn checked_abs(self) -> Option<Self>;
fn checked_add(self, other: Self) -> Option<Self>;
fn scale(self, factor: f64) -> Self;
fn to_bytes(self) -> [u8; 8];
fn try_from_bytes(bytes: [u8; 8]) -> Result<Self, Error>;
Expand All @@ -56,8 +57,13 @@ macro_rules! impl_signed {
const MAX: Self = $max;

#[inline(always)]
fn abs(self) -> Self {
if self >= 0 { self } else { -self }
fn checked_abs(self) -> Option<Self> {
self.checked_abs()
}

#[inline(always)]
fn checked_add(self, other: Self) -> Option<Self> {
self.checked_add(other)
}

#[inline(always)]
Expand Down Expand Up @@ -102,8 +108,13 @@ macro_rules! impl_unsigned {
const MAX: Self = $max;

#[inline(always)]
fn abs(self) -> Self {
self
fn checked_abs(self) -> Option<Self> {
Some(self)
}

#[inline(always)]
fn checked_add(self, other: Self) -> Option<Self> {
self.checked_add(other)
}

#[inline(always)]
Expand Down
12 changes: 4 additions & 8 deletions datasketches/src/kll/sketch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -136,8 +136,9 @@ impl<T: Clone + Ord> KllSketch<T> {
///
/// # Panics
///
/// Panics if the stream weight would exceed [`u64::MAX`].
/// Panics without modifying the sketch if the stream weight would exceed [`u64::MAX`].
pub fn update(&mut self, item: T) {
assert!(self.n < u64::MAX, "total stream weight overflow");
self.update_min_max(&item);
self.internal_update(item);
}
Expand Down Expand Up @@ -708,13 +709,8 @@ impl<T: Clone + Ord> KllSketch<T> {
if self.num_retained >= self.capacity {
self.compress_while_updating();
}
self.n = self.n.checked_add(1).unwrap_or_else(|| {
panic!(
"cannot update KLL sketch: stream weight is {}, maximum is {}",
self.n,
u64::MAX
)
});
// Both update and merge check the final stream weight before modifying the sketch.
self.n += 1;
self.num_retained += 1;
self.is_level_zero_sorted = false;
self.levels[0].push(item);
Expand Down
14 changes: 11 additions & 3 deletions datasketches/src/req/sketch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,12 @@ where
}

/// Updates the sketch with a new item.
///
/// # Panics
///
/// Panics without modifying the sketch if the stream weight would exceed `u64::MAX`.
pub fn update(&mut self, item: T) {
self.n = self.n.checked_add(1).expect("total stream weight overflow");
match &mut self.min_item {
None => self.min_item = Some(item.clone()),
Some(cur) if item.cmp(cur).is_lt() => *cur = item.clone(),
Expand All @@ -146,7 +151,6 @@ where
}

self.compactors[0].append(item);
self.n += 1;
self.num_retained += 1;

if self.num_retained >= self.max_nom_size {
Expand Down Expand Up @@ -285,7 +289,8 @@ where
///
/// # Errors
///
/// Returns an error if the two sketches have different `rank_accuracy`.
/// Returns an error without modifying the sketch if the two sketches have different
/// `rank_accuracy` or their combined stream weight exceeds `u64::MAX`.
///
/// # Examples
///
Expand Down Expand Up @@ -317,7 +322,10 @@ where
return Ok(());
}

self.n += other.n;
self.n = self
.n
.checked_add(other.n)
.ok_or_else(|| Error::invalid_argument("total stream weight overflow"))?;

if let Some(m) = &other.min_item {
match &self.min_item {
Expand Down
47 changes: 24 additions & 23 deletions datasketches/src/tdigest/sketch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -264,6 +264,10 @@ impl TDigestMut {
///
/// [f64::NAN], [f64::INFINITY], and [f64::NEG_INFINITY] values are ignored.
///
/// # Panics
///
/// Panics without modifying the digest if the total weight would exceed `u64::MAX`.
///
/// # Examples
///
/// ```
Expand All @@ -277,6 +281,7 @@ impl TDigestMut {
if !value.is_finite() {
return;
}
assert!(self.total_weight() < u64::MAX, "total weight overflow");

let max_unmerged = self.max_unmerged();
if self.buffer.unmerged_len() >= max_unmerged {
Expand Down Expand Up @@ -322,6 +327,10 @@ impl TDigestMut {

/// Merges the given t-digest into this one.
///
/// # Panics
///
/// Panics without modifying the digest if the combined total weight would exceed `u64::MAX`.
///
/// # Examples
///
/// ```
Expand All @@ -338,15 +347,18 @@ impl TDigestMut {
if other.is_empty() {
return;
}
let total_weight = self
.total_weight()
.checked_add(other.total_weight())
.expect("total weight overflow");

// Preserve true extrema from `other`. Compression only sees centroid means, which can
// differ from `min`/`max` after ordinary compression or deserialization.
self.min = self.min.min(other.min);
self.max = self.max.max(other.max);

let self_unmerged_weight = self.buffer.unmerged_len() as u64;
let centroids = std::mem::take(&mut self.buffer).into_merged_centroids(&other.buffer);
self.compress_sorted_centroids(centroids, self_unmerged_weight + other.total_weight())
self.compress_sorted_centroids(centroids, total_weight);
}

/// Converts this mutable t-digest into an immutable one.
Expand Down Expand Up @@ -882,38 +894,27 @@ impl TDigestMut {

/// Processes unmerged values and merges centroids if needed.
fn compress(&mut self) {
let additional_weight = self.buffer.unmerged_len() as u64;
if additional_weight == 0 {
if self.buffer.unmerged_len() == 0 {
// Also preserves fully compressed deserialized images verbatim.
return;
}
let centroids = std::mem::take(&mut self.buffer).into_centroids_for_compression();
self.compress_centroids(centroids, additional_weight);
}

/// Compresses the given centroids into this t-digest.
///
/// # Contract
///
/// * `centroids` must contain at least one centroid.
/// * `centroids` contains every centroid to be merged, including all centroids previously
/// stored in `self`.
/// * `additional_weight` is the total weight not yet included in `self.compressed_weight`.
/// * Every centroid mean in `centroids` is finite.
/// * `self.buffer` has no unmerged values before returning.
fn compress_centroids(&mut self, mut centroids: Vec<Centroid>, additional_weight: u64) {
debug_assert!(!centroids.is_empty());
let total_weight = self.total_weight();
let mut centroids = std::mem::take(&mut self.buffer).into_centroids_for_compression();
centroids.sort_by(centroid_cmp);
self.compress_sorted_centroids(centroids, additional_weight);
self.compress_sorted_centroids(centroids, total_weight);
}

fn compress_sorted_centroids(&mut self, mut centroids: Vec<Centroid>, additional_weight: u64) {
/// Compresses nonempty, sorted centroids whose combined weight is `total_weight`.
///
/// Includes all retained and incoming values, with finite means and nonzero weights.
/// Callers ensure the total fits in `u64` before taking the buffer.
fn compress_sorted_centroids(&mut self, mut centroids: Vec<Centroid>, total_weight: u64) {
debug_assert!(!centroids.is_empty());
debug_assert!(centroids_are_sorted(&centroids));
if self.reverse_merge {
centroids.reverse();
}
self.compressed_weight += additional_weight;
self.compressed_weight = total_weight;

let mut num_centroids = 1;
let len = centroids.len();
Expand Down
Loading