From 9465a0aea49595012d61f42205f626d88fec3d11 Mon Sep 17 00:00:00 2001 From: Zach McCoy Date: Wed, 16 Sep 2026 15:39:49 -0500 Subject: [PATCH 1/3] Add files from flowstats, cut down, and add licensing attribution Signed-off-by: Zach McCoy --- flowstats/Cargo.toml | 66 ++++ flowstats/LICENSE-MIT | 21 ++ flowstats/benches/benchmarks.rs | 313 ++++++++++++++++++ flowstats/src/frequency/count_min.rs | 410 +++++++++++++++++++++++ flowstats/src/frequency/mod.rs | 37 +++ flowstats/src/frequency/space_saving.rs | 414 ++++++++++++++++++++++++ flowstats/src/lib.rs | 105 ++++++ flowstats/src/math.rs | 75 +++++ flowstats/src/mod.rs | 37 +++ flowstats/src/traits.rs | 199 ++++++++++++ 10 files changed, 1677 insertions(+) create mode 100644 flowstats/Cargo.toml create mode 100644 flowstats/LICENSE-MIT create mode 100644 flowstats/benches/benchmarks.rs create mode 100644 flowstats/src/frequency/count_min.rs create mode 100644 flowstats/src/frequency/mod.rs create mode 100644 flowstats/src/frequency/space_saving.rs create mode 100644 flowstats/src/lib.rs create mode 100644 flowstats/src/math.rs create mode 100644 flowstats/src/mod.rs create mode 100644 flowstats/src/traits.rs diff --git a/flowstats/Cargo.toml b/flowstats/Cargo.toml new file mode 100644 index 00000000..5bfce185 --- /dev/null +++ b/flowstats/Cargo.toml @@ -0,0 +1,66 @@ +[package] +name = "flowstats" +version = "0.1.2" +edition = "2021" +rust-version = "1.75.0" +authors = ["Vahid Negahdari "] +description = "Collection of stream analytics algorithms: cardinality, quantiles, frequency, sampling, and more" +license = "MIT OR Apache-2.0" +repository = "https://github.com/vnvo/flowstats" +documentation = "https://docs.rs/flowstats" +readme = "README.md" +keywords = ["streaming", "algorithms", "hyperloglog", "sketch", "probabilistic"] +categories = ["algorithms", "data-structures", "science"] + +[features] +default = [ + "std", + "cardinality", + "quantiles", + "frequency", + "membership", + "sampling", + "statistics", +] + +# All algorithms +full = [ + "cardinality", + "quantiles", + "frequency", + "membership", + "sampling", + "statistics", +] + +cardinality = [] +quantiles = [] +frequency = [] +membership = [] +sampling = [] +statistics = [] + +# Platform features +std = [] +serde = ["dep:serde"] + +[dependencies] +xxhash-rust = { version = "0.8", features = ["xxh3"] } +serde = { version = "1.0", features = ["derive"], optional = true } +libm = "0.2" + +[dev-dependencies] +criterion = "0.5" +rand = "0.9" + +[[bench]] +name = "benchmarks" +harness = false + +[profile.release] +lto = true +codegen-units = 1 + +[package.metadata.docs.rs] +all-features = true +rustdoc-args = ["--cfg", "docsrs"] diff --git a/flowstats/LICENSE-MIT b/flowstats/LICENSE-MIT new file mode 100644 index 00000000..9e3e2ee5 --- /dev/null +++ b/flowstats/LICENSE-MIT @@ -0,0 +1,21 @@ +MIT License + +Copyright (c) 2025 flowstats Contributors + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. diff --git a/flowstats/benches/benchmarks.rs b/flowstats/benches/benchmarks.rs new file mode 100644 index 00000000..047a4d89 --- /dev/null +++ b/flowstats/benches/benchmarks.rs @@ -0,0 +1,313 @@ +//! Benchmarks for flowstats algorithms +//! +//! Run with: cargo bench --features full + +// Require all features for benchmarks +#[cfg(not(all( + feature = "cardinality", + feature = "frequency", + feature = "quantiles", + feature = "membership", + feature = "sampling", + feature = "statistics" +)))] +compile_error!("Benchmarks require all features. Run: cargo bench --features full"); + +use criterion::{black_box, criterion_group, criterion_main, Criterion, Throughput}; + +use flowstats::cardinality::HyperLogLog; +use flowstats::frequency::{CountMinSketch, SpaceSaving}; +use flowstats::membership::BloomFilter; +use flowstats::quantiles::TDigest; +use flowstats::sampling::ReservoirSampler; +use flowstats::statistics::RunningStats; +use flowstats::traits::{CardinalitySketch, HeavyHitters, QuantileSketch, Sketch}; + +// ============================================================================ +// HyperLogLog Benchmarks +// ============================================================================ + +fn bench_hll(c: &mut Criterion) { + let mut group = c.benchmark_group("hyperloglog"); + group.throughput(Throughput::Elements(1)); + + for precision in [10, 12, 14, 16] { + group.bench_function(format!("insert_p{}", precision), |b| { + let mut hll = HyperLogLog::new(precision); + let mut i = 0u64; + b.iter(|| { + hll.insert(&i.to_string()); + i = i.wrapping_add(1); + }); + }); + } + + group.bench_function("estimate", |b| { + let mut hll = HyperLogLog::new(14); + for i in 0..100_000u64 { + hll.insert(&i.to_string()); + } + b.iter(|| black_box(hll.estimate())); + }); + + group.bench_function("merge", |b| { + let mut hll1 = HyperLogLog::new(14); + let mut hll2 = HyperLogLog::new(14); + for i in 0..10_000u64 { + hll1.insert(&i.to_string()); + hll2.insert(&(i + 10_000).to_string()); + } + b.iter(|| { + let mut h = hll1.clone(); + h.merge(black_box(&hll2)).unwrap(); + }); + }); + + group.finish(); +} + +// ============================================================================ +// Count-Min Sketch Benchmarks +// ============================================================================ + +fn bench_cms(c: &mut Criterion) { + let mut group = c.benchmark_group("count_min_sketch"); + group.throughput(Throughput::Elements(1)); + + group.bench_function("add", |b| { + let mut cms = CountMinSketch::new(0.001, 0.01); + let mut i = 0u64; + b.iter(|| { + cms.add(i.to_string().as_bytes(), 1); + i = i.wrapping_add(1); + }); + }); + + group.bench_function("estimate", |b| { + let mut cms = CountMinSketch::new(0.001, 0.01); + for i in 0..100_000u64 { + cms.add(i.to_string().as_bytes(), 1); + } + b.iter(|| black_box(cms.estimate(b"12345"))); + }); + + group.bench_function("merge", |b| { + let mut cms1 = CountMinSketch::new(0.001, 0.01); + let mut cms2 = CountMinSketch::new(0.001, 0.01); + for i in 0..10_000u64 { + cms1.add(i.to_string().as_bytes(), 1); + cms2.add((i + 10_000).to_string().as_bytes(), 1); + } + b.iter(|| { + let mut c = cms1.clone(); + c.merge(black_box(&cms2)).unwrap(); + }); + }); + + group.finish(); +} + +// ============================================================================ +// Space-Saving Benchmarks +// ============================================================================ + +fn bench_space_saving(c: &mut Criterion) { + let mut group = c.benchmark_group("space_saving"); + group.throughput(Throughput::Elements(1)); + + for k in [100, 1000] { + group.bench_function(format!("add_k{}", k), |b| { + let mut ss: SpaceSaving = SpaceSaving::new(k); + let mut i = 0u64; + b.iter(|| { + ss.add(i.to_string()); + i = i.wrapping_add(1); + }); + }); + } + + group.bench_function("top_k", |b| { + let mut ss: SpaceSaving = SpaceSaving::new(1000); + for i in 0..100_000u64 { + ss.add((i % 10_000).to_string()); + } + b.iter(|| black_box(ss.top_k(10))); + }); + + group.finish(); +} + +// ============================================================================ +// t-digest Benchmarks +// ============================================================================ + +fn bench_tdigest(c: &mut Criterion) { + let mut group = c.benchmark_group("tdigest"); + group.throughput(Throughput::Elements(1)); + + for compression in [50.0, 100.0, 200.0] { + group.bench_function(format!("add_c{}", compression as u32), |b| { + let mut td = TDigest::new(compression); + let mut i = 0u64; + b.iter(|| { + td.add((i as f64) * 0.001); + i = i.wrapping_add(1); + }); + }); + } + + group.bench_function("quantile", |b| { + let mut td = TDigest::new(100.0); + for i in 0..100_000u64 { + td.add(i as f64); + } + b.iter(|| black_box(td.quantile(0.99))); + }); + + group.bench_function("merge", |b| { + let mut td1 = TDigest::new(100.0); + let mut td2 = TDigest::new(100.0); + for i in 0..10_000u64 { + td1.add(i as f64); + td2.add((i + 10_000) as f64); + } + b.iter(|| { + let mut t = td1.clone(); + t.merge(black_box(&td2)).unwrap(); + }); + }); + + group.finish(); +} + +// ============================================================================ +// Bloom Filter Benchmarks +// ============================================================================ + +fn bench_bloom(c: &mut Criterion) { + let mut group = c.benchmark_group("bloom_filter"); + group.throughput(Throughput::Elements(1)); + + group.bench_function("insert", |b| { + let mut bloom = BloomFilter::new(1_000_000, 0.01); + let mut i = 0u64; + b.iter(|| { + bloom.insert(i.to_string().as_bytes()); + i = i.wrapping_add(1); + }); + }); + + group.bench_function("contains_hit", |b| { + let mut bloom = BloomFilter::new(100_000, 0.01); + for i in 0..100_000u64 { + bloom.insert(i.to_string().as_bytes()); + } + let mut i = 0u64; + b.iter(|| { + let result = bloom.contains((i % 100_000).to_string().as_bytes()); + i = i.wrapping_add(1); + black_box(result) + }); + }); + + group.bench_function("contains_miss", |b| { + let mut bloom = BloomFilter::new(100_000, 0.01); + for i in 0..100_000u64 { + bloom.insert(i.to_string().as_bytes()); + } + let mut i = 1_000_000u64; + b.iter(|| { + let result = bloom.contains(i.to_string().as_bytes()); + i = i.wrapping_add(1); + black_box(result) + }); + }); + + group.finish(); +} + +// ============================================================================ +// Reservoir Sampler Benchmarks +// ============================================================================ + +fn bench_reservoir(c: &mut Criterion) { + let mut group = c.benchmark_group("reservoir"); + group.throughput(Throughput::Elements(1)); + + for capacity in [100, 1000, 10000] { + group.bench_function(format!("add_cap{}", capacity), |b| { + let mut sampler: ReservoirSampler = ReservoirSampler::new(capacity); + let mut i = 0u64; + b.iter(|| { + sampler.add(i); + i = i.wrapping_add(1); + }); + }); + } + + group.finish(); +} + +// ============================================================================ +// Running Stats Benchmarks +// ============================================================================ + +fn bench_running_stats(c: &mut Criterion) { + let mut group = c.benchmark_group("running_stats"); + group.throughput(Throughput::Elements(1)); + + group.bench_function("add", |b| { + let mut stats = RunningStats::new(); + let mut i = 0u64; + b.iter(|| { + stats.add(i as f64); + i = i.wrapping_add(1); + }); + }); + + group.bench_function("query_all", |b| { + let mut stats = RunningStats::new(); + for i in 0..100_000u64 { + stats.add(i as f64); + } + b.iter(|| { + black_box(stats.mean()); + black_box(stats.variance()); + black_box(stats.stddev()); + black_box(stats.min()); + black_box(stats.max()); + }); + }); + + group.bench_function("merge", |b| { + let mut s1 = RunningStats::new(); + let mut s2 = RunningStats::new(); + for i in 0..10_000u64 { + s1.add(i as f64); + s2.add((i + 10_000) as f64); + } + b.iter(|| { + let mut s = s1.clone(); + s.merge(black_box(&s2)).unwrap(); + }); + }); + + group.finish(); +} + +// ============================================================================ +// Main +// ============================================================================ + +criterion_group!( + benches, + bench_hll, + bench_cms, + bench_space_saving, + bench_tdigest, + bench_bloom, + bench_reservoir, + bench_running_stats, +); + +criterion_main!(benches); diff --git a/flowstats/src/frequency/count_min.rs b/flowstats/src/frequency/count_min.rs new file mode 100644 index 00000000..0d41f10e --- /dev/null +++ b/flowstats/src/frequency/count_min.rs @@ -0,0 +1,410 @@ +// SPDX-License-Identifier: MIT +// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors +// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/frequency/count_min.rs + +//! Count-Min Sketch frequency estimator +//! +//! The Count-Min Sketch is a probabilistic data structure for estimating +//! the frequency of elements in a data stream. + +use crate::math; +use crate::traits::{FrequencySketch, MergeError, Sketch}; +use xxhash_rust::xxh3::xxh3_64_with_seed; + +#[cfg(feature = "std")] +use std::vec::Vec; + +#[cfg(not(feature = "std"))] +extern crate alloc; +#[cfg(not(feature = "std"))] +use alloc::vec::Vec; + +/// Count-Min Sketch for frequency estimation +/// +/// The Count-Min Sketch provides frequency estimates with the following guarantees: +/// - Point query: `actual_count <= estimate <= actual_count + ε * N` +/// - Where ε = e/width and N is the total count +/// - Probability of exceeding the error bound: δ = 1/2^depth +/// +/// # Example +/// +/// ``` +/// use flowstats::frequency::CountMinSketch; +/// use flowstats::traits::FrequencySketch; +/// +/// // Create with 1% error rate and 0.1% failure probability +/// let mut cms = CountMinSketch::new(0.01, 0.001); +/// +/// // Add items +/// cms.add(b"apple", 5); +/// cms.add(b"banana", 3); +/// cms.add(b"apple", 2); +/// +/// // Query frequency +/// let apple_count = cms.estimate(b"apple"); // ~7 +/// let banana_count = cms.estimate(b"banana"); // ~3 +/// ``` +#[derive(Clone, Debug)] +pub struct CountMinSketch { + /// Width of each row + width: usize, + /// Number of rows (hash functions) + depth: usize, + /// Counter table + table: Vec>, + /// Total count of all items + total_count: u64, + /// Number of updates + num_updates: u64, + /// Seeds for hash functions + seeds: Vec, +} + +impl CountMinSketch { + /// Create a new Count-Min Sketch with the given error parameters + /// + /// # Arguments + /// + /// * `epsilon` - Maximum overcount as a fraction of total (e.g., 0.01 for 1%) + /// * `delta` - Probability of exceeding the error bound (e.g., 0.001 for 0.1%) + /// + /// # Panics + /// + /// Panics if epsilon or delta are not in (0, 1) + pub fn new(epsilon: f64, delta: f64) -> Self { + assert!(epsilon > 0.0 && epsilon < 1.0, "epsilon must be in (0, 1)"); + assert!(delta > 0.0 && delta < 1.0, "delta must be in (0, 1)"); + + // width = ceil(e / epsilon) + // depth = ceil(ln(1/delta)) + let width = math::ceil(core::f64::consts::E / epsilon) as usize; + let depth = math::ceil(math::ln(1.0 / delta)) as usize; + + Self::with_dimensions(width, depth) + } + + /// Create a Count-Min Sketch with specific dimensions + /// + /// # Arguments + /// + /// * `width` - Width of each row (larger = lower error) + /// * `depth` - Number of rows (larger = lower failure probability) + pub fn with_dimensions(width: usize, depth: usize) -> Self { + assert!(width > 0, "width must be positive"); + assert!(depth > 0, "depth must be positive"); + + // Generate random seeds for hash functions + let seeds: Vec = (0..depth) + .map(|i| (i as u64).wrapping_mul(0x9e3779b97f4a7c15)) + .collect(); + + Self { + width, + depth, + table: vec![vec![0u64; width]; depth], + total_count: 0, + num_updates: 0, + seeds, + } + } + + /// Get the width of the sketch + pub fn width(&self) -> usize { + self.width + } + + /// Get the depth of the sketch + pub fn depth(&self) -> usize { + self.depth + } + + /// Get the total count of all items + pub fn total_count(&self) -> u64 { + self.total_count + } + + /// Add count to an item + pub fn add(&mut self, item: &[u8], count: u64) { + self.num_updates += 1; + self.total_count += count; + + for (row, &seed) in self.seeds.iter().enumerate() { + let hash = xxh3_64_with_seed(item, seed); + let col = (hash as usize) % self.width; + self.table[row][col] = self.table[row][col].saturating_add(count); + } + } + + /// Add count using conservative update + /// + /// Conservative update improves accuracy by only incrementing counters + /// up to the new estimated value. This reduces over-counting. + pub fn add_conservative(&mut self, item: &[u8], count: u64) { + self.num_updates += 1; + self.total_count += count; + + // First pass: find current estimate (minimum) + let min_val = self.estimate(item); + let new_val = min_val.saturating_add(count); + + // Second pass: set all counters to at least new_val + for (row, &seed) in self.seeds.iter().enumerate() { + let hash = xxh3_64_with_seed(item, seed); + let col = (hash as usize) % self.width; + + if self.table[row][col] < new_val { + self.table[row][col] = new_val; + } + } + } + + /// Estimate the frequency of an item + pub fn estimate(&self, item: &[u8]) -> u64 { + let mut min_count = u64::MAX; + + for (row, &seed) in self.seeds.iter().enumerate() { + let hash = xxh3_64_with_seed(item, seed); + let col = (hash as usize) % self.width; + min_count = min_count.min(self.table[row][col]); + } + + min_count + } + + /// Inner product of two sketches + /// + /// This can be used to estimate the dot product of two frequency distributions. + pub fn inner_product(&self, other: &Self) -> Option { + if self.width != other.width || self.depth != other.depth { + return None; + } + + let mut min_product = u64::MAX; + + for row in 0..self.depth { + let product: u64 = self.table[row] + .iter() + .zip(other.table[row].iter()) + .fold(0u64, |acc, (&a, &b)| { + acc.saturating_add(a.saturating_mul(b)) + }); + min_product = min_product.min(product); + } + + Some(min_product) + } + + /// Theoretical error bound (epsilon * total_count) + pub fn error_bound(&self) -> u64 { + let epsilon = core::f64::consts::E / self.width as f64; + (epsilon * self.total_count as f64) as u64 + } +} + +impl Sketch for CountMinSketch { + type Item = [u8]; + + fn update(&mut self, item: &[u8]) { + self.add(item, 1); + } + + fn merge(&mut self, other: &Self) -> Result<(), MergeError> { + if self.width != other.width || self.depth != other.depth { + return Err(MergeError::IncompatibleConfig { + expected: format!("{}x{}", self.width, self.depth), + found: format!("{}x{}", other.width, other.depth), + }); + } + + for row in 0..self.depth { + for col in 0..self.width { + self.table[row][col] = self.table[row][col].saturating_add(other.table[row][col]); + } + } + + self.total_count += other.total_count; + self.num_updates += other.num_updates; + + Ok(()) + } + + fn clear(&mut self) { + for row in &mut self.table { + row.fill(0); + } + self.total_count = 0; + self.num_updates = 0; + } + + fn size_bytes(&self) -> usize { + core::mem::size_of::() + + self.depth * self.width * core::mem::size_of::() + + self.seeds.len() * core::mem::size_of::() + } + + fn count(&self) -> u64 { + self.num_updates + } +} + + +impl FrequencySketch for CountMinSketch { + fn estimate_frequency(&self, item: &[u8]) -> u64 { + self.estimate(item) + } +} + +#[cfg(feature = "serde")] +impl serde::Serialize for CountMinSketch { + fn serialize(&self, serializer: S) -> Result + where + S: serde::Serializer, + { + use serde::ser::SerializeStruct; + let mut state = serializer.serialize_struct("CountMinSketch", 6)?; + state.serialize_field("width", &self.width)?; + state.serialize_field("depth", &self.depth)?; + state.serialize_field("table", &self.table)?; + state.serialize_field("total_count", &self.total_count)?; + state.serialize_field("num_updates", &self.num_updates)?; + state.serialize_field("seeds", &self.seeds)?; + state.end() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_basic() { + let mut cms = CountMinSketch::new(0.01, 0.01); + + cms.add(b"apple", 5); + cms.add(b"banana", 3); + cms.add(b"cherry", 1); + cms.add(b"apple", 2); + + // Estimates should be at least the true count + assert!(cms.estimate(b"apple") >= 7); + assert!(cms.estimate(b"banana") >= 3); + assert!(cms.estimate(b"cherry") >= 1); + } + + #[test] + fn test_empty() { + let cms = CountMinSketch::new(0.01, 0.01); + assert_eq!(cms.estimate(b"anything"), 0); + assert_eq!(cms.total_count(), 0); + } + + #[test] + fn test_conservative_update() { + let mut cms1 = CountMinSketch::new(0.001, 0.001); + let mut cms2 = CountMinSketch::new(0.001, 0.001); + + // Add many items + for i in 0..10000 { + let item = format!("item_{}", i); + cms1.add(item.as_bytes(), 1); + cms2.add_conservative(item.as_bytes(), 1); + } + + // Conservative update should generally have lower estimates for items we query + // (fewer collisions impacting the count) + let test_item = b"test_item"; + cms1.add(test_item, 100); + cms2.add_conservative(test_item, 100); + + // Both should report at least 100 + assert!(cms1.estimate(test_item) >= 100); + assert!(cms2.estimate(test_item) >= 100); + } + + #[test] + fn test_merge() { + let mut cms1 = CountMinSketch::with_dimensions(1000, 5); + let mut cms2 = CountMinSketch::with_dimensions(1000, 5); + + cms1.add(b"apple", 5); + cms2.add(b"banana", 3); + + cms1.merge(&cms2).unwrap(); + + assert!(cms1.estimate(b"apple") >= 5); + assert!(cms1.estimate(b"banana") >= 3); + assert_eq!(cms1.total_count(), 8); + } + + #[test] + fn test_merge_incompatible() { + let mut cms1 = CountMinSketch::with_dimensions(1000, 5); + let cms2 = CountMinSketch::with_dimensions(2000, 5); + + assert!(cms1.merge(&cms2).is_err()); + } + + #[test] + fn test_dimensions() { + let cms = CountMinSketch::with_dimensions(1000, 5); + assert_eq!(cms.width(), 1000); + assert_eq!(cms.depth(), 5); + } + + #[test] + fn test_clear() { + let mut cms = CountMinSketch::new(0.01, 0.01); + + cms.add(b"item", 100); + assert!(cms.estimate(b"item") >= 100); + + cms.clear(); + assert_eq!(cms.estimate(b"item"), 0); + assert_eq!(cms.total_count(), 0); + } + + #[test] + fn test_heavy_usage() { + let mut cms = CountMinSketch::new(0.01, 0.001); + + // Add 100,000 items + for i in 0..100_000 { + let item = format!("user_{}", i % 1000); // 1000 unique items + cms.add(item.as_bytes(), 1); + } + + // Each item should have been added ~100 times + // With 1% error, we allow some slack + for i in 0..10 { + let item = format!("user_{}", i); + let estimate = cms.estimate(item.as_bytes()); + // Should be at least 100, and not more than 100 + error_bound + assert!(estimate >= 100, "item {} estimate {} < 100", i, estimate); + } + } + + #[test] + fn test_error_bound() { + let cms = CountMinSketch::new(0.01, 0.001); + // Error bound formula: epsilon * total_count + // Initially 0 + assert_eq!(cms.error_bound(), 0); + } + + #[test] + fn test_inner_product() { + let mut cms1 = CountMinSketch::with_dimensions(1000, 5); + let mut cms2 = CountMinSketch::with_dimensions(1000, 5); + + cms1.add(b"a", 10); + cms1.add(b"b", 5); + + cms2.add(b"a", 3); + cms2.add(b"b", 2); + + // Inner product should give approximation of sum(f1[i] * f2[i]) + // = 10*3 + 5*2 = 40 + let ip = cms1.inner_product(&cms2).unwrap(); + assert!(ip >= 40, "inner product {} < 40", ip); + } +} diff --git a/flowstats/src/frequency/mod.rs b/flowstats/src/frequency/mod.rs new file mode 100644 index 00000000..ec1d04a8 --- /dev/null +++ b/flowstats/src/frequency/mod.rs @@ -0,0 +1,37 @@ +// SPDX-License-Identifier: MIT +// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors +// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/frequency/mod.rs + +//! Frequency estimation algorithms +//! +//! This module provides implementations of sketches for estimating item +//! frequencies in a data stream. +//! +//! # Algorithms +//! +//! - [`CountMinSketch`]: Classic count-min sketch with optional conservative update +//! - [`SpaceSaving`]: Top-K / heavy hitters tracking +//! +//! # Example +//! +//! ``` +//! use flowstats::frequency::CountMinSketch; +//! use flowstats::traits::FrequencySketch; +//! +//! let mut cms = CountMinSketch::new(0.01, 0.001); // 1% error, 0.1% probability +//! +//! cms.add(b"item1", 5); +//! cms.add(b"item2", 3); +//! +//! let count = cms.estimate(b"item1"); +//! println!("Estimated count: {}", count); +//! ``` + +mod count_min; +#[cfg(feature = "std")] +mod space_saving; + +pub use count_min::CountMinSketch; + +#[cfg(feature = "std")] +pub use space_saving::SpaceSaving; diff --git a/flowstats/src/frequency/space_saving.rs b/flowstats/src/frequency/space_saving.rs new file mode 100644 index 00000000..30fdbe45 --- /dev/null +++ b/flowstats/src/frequency/space_saving.rs @@ -0,0 +1,414 @@ +// SPDX-License-Identifier: MIT +// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors +// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/frequency/space_saving.rs + +//! Space-Saving algorithm for heavy hitters +//! +//! The Space-Saving algorithm efficiently tracks the k most frequent items +//! in a data stream using only O(k) space. +//! +//! **Note**: This module requires the `std` feature. The Space-Saving algorithm +//! does not support correct distributed merging due to its replacement strategy. +//! Attempting to merge will return an error. + +use crate::traits::{FrequencySketch, MergeError, Sketch}; +use core::hash::Hash; +use std::collections::HashMap; +use std::vec::Vec; + +/// Entry in the Space-Saving structure +#[derive(Clone, Debug)] +struct Counter { + /// The item + item: T, + /// Estimated count + count: u64, + /// Error bound (maximum overcount) + error: u64, +} + +impl Counter { + fn new(item: T, count: u64, error: u64) -> Self { + Self { item, count, error } + } +} + +/// Space-Saving algorithm for finding frequent items +/// +/// The Space-Saving algorithm maintains a summary of the k most frequent items +/// with the following guarantees: +/// +/// - Any item with true frequency > n/k is guaranteed to be in the summary +/// - The maximum overcount error for any item is at most n/k +/// +/// # Example +/// +/// ``` +/// use flowstats::frequency::SpaceSaving; +/// use flowstats::traits::HeavyHitters; +/// +/// let mut ss = SpaceSaving::new(10); // Track top 10 +/// +/// // Add some items +/// for _ in 0..100 { ss.add("apple"); } +/// for _ in 0..50 { ss.add("banana"); } +/// for _ in 0..25 { ss.add("cherry"); } +/// for _ in 0..10 { ss.add("date"); } +/// +/// // Get top 3 items +/// let top = ss.top_k(3); +/// println!("Top items: {:?}", top); +/// ``` +#[derive(Clone, Debug)] +pub struct SpaceSaving { + /// Maximum number of counters to maintain + capacity: usize, + /// Map from item to counter index + item_to_index: HashMap, + /// Counter array + counters: Vec>, + /// Total count of all items + total_count: u64, + /// Number of updates + num_updates: u64, +} + +impl SpaceSaving { + /// Create a new Space-Saving structure with the given capacity + /// + /// # Arguments + /// + /// * `capacity` - Maximum number of items to track (k) + pub fn new(capacity: usize) -> Self { + assert!(capacity > 0, "capacity must be positive"); + + Self { + capacity, + item_to_index: HashMap::with_capacity(capacity), + counters: Vec::with_capacity(capacity), + total_count: 0, + num_updates: 0, + } + } + + /// Get the capacity (k) + pub fn capacity(&self) -> usize { + self.capacity + } + + /// Get the number of distinct items currently tracked + pub fn num_tracked(&self) -> usize { + self.counters.len() + } + + /// Get the total count + pub fn total_count(&self) -> u64 { + self.total_count + } + + /// Add an item to the structure + pub fn add(&mut self, item: T) { + self.add_count(item, 1); + } + + /// Add an item with a specific count + pub fn add_count(&mut self, item: T, count: u64) { + self.num_updates += 1; + self.total_count += count; + + // Check if item is already tracked + if let Some(&idx) = self.item_to_index.get(&item) { + self.counters[idx].count += count; + return; + } + + // Not tracked - either add new or replace minimum + if self.counters.len() < self.capacity { + // Still have room, add new counter + let idx = self.counters.len(); + self.counters.push(Counter::new(item.clone(), count, 0)); + self.item_to_index.insert(item, idx); + } else { + // Find and replace minimum counter + let min_idx = self.find_min_index(); + let min_count = self.counters[min_idx].count; + + // Remove old item from map + let old_item = self.counters[min_idx].item.clone(); + self.item_to_index.remove(&old_item); + + // Replace with new item + self.counters[min_idx] = Counter::new(item.clone(), min_count + count, min_count); + self.item_to_index.insert(item, min_idx); + } + } + + /// Find the index of the counter with minimum count + fn find_min_index(&self) -> usize { + self.counters + .iter() + .enumerate() + .min_by_key(|(_, c)| c.count) + .map(|(i, _)| i) + .unwrap_or(0) + } + + /// Estimate the frequency of an item + pub fn estimate(&self, item: &T) -> u64 { + self.item_to_index + .get(item) + .map(|&idx| self.counters[idx].count) + .unwrap_or(0) + } + + /// Get the error bound for an item's estimate + pub fn error(&self, item: &T) -> u64 { + self.item_to_index + .get(item) + .map(|&idx| self.counters[idx].error) + .unwrap_or(0) + } + + /// Get guaranteed minimum count for an item + /// + /// Returns (count - error), which is guaranteed to be at most the true count. + pub fn guaranteed_count(&self, item: &T) -> u64 { + self.item_to_index + .get(item) + .map(|&idx| { + let c = &self.counters[idx]; + c.count.saturating_sub(c.error) + }) + .unwrap_or(0) + } + + /// Check if an item is currently tracked + pub fn contains(&self, item: &T) -> bool { + self.item_to_index.contains_key(item) + } +} + +impl Sketch for SpaceSaving { + type Item = T; + + fn update(&mut self, item: &T) { + self.add(item.clone()); + } + + fn merge(&mut self, _other: &Self) -> Result<(), MergeError> { + // Space-Saving does not support correct distributed merge. + // The algorithm's replacement strategy makes merging two summaries + // produce incorrect results (inflated counts, dropped heavy hitters). + // Use a single SpaceSaving instance or aggregate raw data instead. + Err(MergeError::IncompatibleConfig { + expected: "Space-Saving does not support merge".into(), + found: "merge attempted".into(), + }) + } + + fn clear(&mut self) { + self.item_to_index.clear(); + self.counters.clear(); + self.total_count = 0; + self.num_updates = 0; + } + + fn size_bytes(&self) -> usize { + core::mem::size_of::() + self.counters.capacity() * core::mem::size_of::>() + } + + fn count(&self) -> u64 { + self.num_updates + } +} + +impl FrequencySketch for SpaceSaving { + fn estimate_frequency(&self, item: &T) -> u64 { + self.estimate(item) + } +} + +#[cfg(feature = "serde")] +impl serde::Serialize + for SpaceSaving +{ + fn serialize(&self, serializer: S) -> Result + where + S: serde::Serializer, + { + use serde::ser::SerializeStruct; + let items: Vec<_> = self + .counters + .iter() + .map(|c| (&c.item, c.count, c.error)) + .collect(); + + let mut state = serializer.serialize_struct("SpaceSaving", 3)?; + state.serialize_field("capacity", &self.capacity)?; + state.serialize_field("total_count", &self.total_count)?; + state.serialize_field("items", &items)?; + state.end() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_basic() { + let mut ss = SpaceSaving::::new(10); + + ss.add("apple".to_string()); + ss.add("apple".to_string()); + ss.add("banana".to_string()); + + assert!(ss.estimate(&"apple".to_string()) >= 2); + assert!(ss.estimate(&"banana".to_string()) >= 1); + } + + #[test] + fn test_empty() { + let ss = SpaceSaving::::new(10); + assert_eq!(ss.estimate(&"anything".to_string()), 0); + assert_eq!(ss.total_count(), 0); + } + + #[test] + fn test_top_k() { + let mut ss = SpaceSaving::<&str>::new(10); + + for _ in 0..100 { + ss.add("apple"); + } + for _ in 0..50 { + ss.add("banana"); + } + for _ in 0..25 { + ss.add("cherry"); + } + + let top = ss.top_k(2); + assert_eq!(top.len(), 2); + assert_eq!(top[0].0, "apple"); + assert_eq!(top[1].0, "banana"); + } + + #[test] + fn test_heavy_hitters() { + let mut ss = SpaceSaving::<&str>::new(10); + + for _ in 0..100 { + ss.add("apple"); + } + for _ in 0..10 { + ss.add("banana"); + } + for _ in 0..1 { + ss.add("cherry"); + } + + // Items with > 5% of total + let heavy = ss.heavy_hitters(0.05); + assert!(heavy.iter().any(|(item, _)| *item == "apple")); + assert!(heavy.iter().any(|(item, _)| *item == "banana")); + // Cherry should not be included (1/111 < 5%) + } + + #[test] + fn test_replacement() { + let mut ss = SpaceSaving::::new(3); + + // Fill up capacity + ss.add(1); + ss.add(2); + ss.add(3); + + assert_eq!(ss.num_tracked(), 3); + + // Add a new item - should replace the minimum + ss.add(4); + + // Still only 3 items tracked + assert_eq!(ss.num_tracked(), 3); + } + + #[test] + fn test_contains() { + let mut ss = SpaceSaving::<&str>::new(10); + + ss.add("apple"); + + assert!(ss.contains(&"apple")); + assert!(!ss.contains(&"banana")); + } + + #[test] + fn test_merge_not_supported() { + let mut ss1 = SpaceSaving::<&str>::new(10); + let ss2 = SpaceSaving::<&str>::new(10); + + // Space-Saving does not support merge + assert!(ss1.merge(&ss2).is_err()); + } + + #[test] + fn test_guaranteed_count() { + let mut ss = SpaceSaving::<&str>::new(3); + + // Add items that will cause replacements + ss.add("a"); + ss.add("b"); + ss.add("c"); + + // Add more to force replacement + ss.add("d"); + + // Check that guaranteed count is <= estimated count + for item in ["a", "b", "c", "d"] { + let est = ss.estimate(&item); + let guar = ss.guaranteed_count(&item); + assert!( + guar <= est, + "guaranteed {} > estimate {} for {}", + guar, + est, + item + ); + } + } + + #[test] + fn test_zipf_distribution() { + // Simulate Zipf distribution (common in real data) + let mut ss = SpaceSaving::::new(10); + + // Item 1 appears 1000 times, item 2 appears 500 times, etc. + for rank in 1..=100 { + let count = 1000 / rank; + for _ in 0..count { + ss.add(rank); + } + } + + // Top items should be 1, 2, 3, ... + let top = ss.top_k(5); + assert!(!top.is_empty()); + // Item 1 should definitely be tracked + assert!(ss.contains(&1)); + } + + #[test] + fn test_clear() { + let mut ss = SpaceSaving::<&str>::new(10); + + ss.add("apple"); + ss.add("banana"); + + ss.clear(); + + assert_eq!(ss.num_tracked(), 0); + assert_eq!(ss.total_count(), 0); + assert!(!ss.contains(&"apple")); + } +} diff --git a/flowstats/src/lib.rs b/flowstats/src/lib.rs new file mode 100644 index 00000000..836296fd --- /dev/null +++ b/flowstats/src/lib.rs @@ -0,0 +1,105 @@ +// SPDX-License-Identifier: MIT +// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors +// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/lib.rs + +//! # Flowstats +//! +//! Production-grade streaming algorithms for Rust. +//! +//! Flowstats provides high-performance implementations of probabilistic data structures +//! and streaming algorithms, designed for real-time analytics and large-scale data processing. +//! +//! ## Features +//! +//! - **Cardinality Estimation**: Count distinct elements with HyperLogLog +//! - **Frequency Estimation**: Track item frequencies with Count-Min Sketch +//! - **Heavy Hitters**: Find top-K elements with Space-Saving +//! - **Quantile Estimation**: Compute percentiles with t-digest +//! - **Full Mergeability**: All sketches support distributed merge operations +//! - **Error Bounds**: Formal guarantees on approximation accuracy +//! +//! ## Quick Start +//! +//! ```rust +//! use flowstats::prelude::*; +//! +//! // Count distinct users +//! let mut hll = HyperLogLog::new(14); +//! for user_id in ["alice", "bob", "charlie", "alice"] { +//! hll.insert(user_id); +//! } +//! println!("Distinct users: ~{}", hll.estimate()); +//! +//! // Track request latencies +//! let mut digest = TDigest::new(100.0); +//! for latency in [12.5, 45.2, 23.1, 67.8, 15.3] { +//! digest.add(latency); +//! } +//! println!("p99 latency: {:?}", digest.quantile(0.99)); +//! +//! ``` +//! +//! ## Distributed Computing +//! +//! All sketches implement the [`Sketch`](traits::Sketch) trait which includes +//! a `merge` operation, allowing sketches to be combined across distributed workers: +//! +//! ```rust +//! use flowstats::cardinality::HyperLogLog; +//! use flowstats::traits::Sketch; +//! +//! let mut worker1 = HyperLogLog::new(14); +//! let mut worker2 = HyperLogLog::new(14); +//! +//! // Each worker processes its partition +//! worker1.insert("user_a"); +//! worker2.insert("user_b"); +//! +//! // Merge results +//! worker1.merge(&worker2).unwrap(); +//! ``` +//! +//! ## Feature Flags +//! +//! Algorithm families (pick what you need): +//! - `cardinality` (default): HyperLogLog for distinct counting +//! - `frequency` (default): Count-Min Sketch, Space-Saving (tbd) +//! - `quantiles` (default): t-digest for percentiles +//! - `membership`: (default) Bloom filter +//! - `sampling`: Reservoir and weighted sampling (tbd) +//! - `sets`: Theta sketch for set operations (tbd) +//! - `statistics`: Running moments, entropy (tbd) +//! - `full`: Enable all algorithm families +//! +//! Platform features: +//! - `std` (default): Standard library support +//! - `serde`: Enable serialization + +#![cfg_attr(not(feature = "std"), no_std)] +#![cfg_attr(docsrs, feature(doc_cfg))] + +#[cfg(not(feature = "std"))] +extern crate alloc; + +mod math; +pub mod traits; + +#[cfg(feature = "frequency")] +#[cfg_attr(docsrs, doc(cfg(feature = "frequency")))] +pub mod frequency; + +pub mod prelude { + pub use crate::traits::*; + + #[cfg(feature = "frequency")] + pub use crate::frequency::CountMinSketch; + + #[cfg(all(feature = "frequency", feature = "std"))] + pub use crate::frequency::SpaceSaving; +} + +#[cfg(feature = "frequency")] +pub use frequency::CountMinSketch; + +#[cfg(all(feature = "frequency", feature = "std"))] +pub use frequency::SpaceSaving; diff --git a/flowstats/src/math.rs b/flowstats/src/math.rs new file mode 100644 index 00000000..da73f726 --- /dev/null +++ b/flowstats/src/math.rs @@ -0,0 +1,75 @@ +// SPDX-License-Identifier: MIT +// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors +// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/math.rs + +//! Math function wrappers for std/no_std compatibility +//! +//! Uses standard library math when available, falls back to libm for no_std. + +#[cfg(feature = "std")] +#[inline] +pub fn ln(x: f64) -> f64 { + x.ln() +} + +#[cfg(not(feature = "std"))] +#[inline] +pub fn ln(x: f64) -> f64 { + libm::log(x) +} + + +#[cfg(not(feature = "std"))] +#[inline] +pub fn log2(x: f64) -> f64 { + libm::log2(x) +} + +#[cfg(not(feature = "std"))] +#[inline] +pub fn exp2(x: f64) -> f64 { + libm::exp2(x) +} + +#[cfg(not(feature = "std"))] +#[inline] +pub fn sqrt(x: f64) -> f64 { + libm::sqrt(x) +} + +#[cfg(not(feature = "std"))] +#[inline] +pub fn asin(x: f64) -> f64 { + libm::asin(x) +} + +#[cfg(feature = "std")] +#[inline] +pub fn ceil(x: f64) -> f64 { + x.ceil() +} + +#[cfg(not(feature = "std"))] +#[inline] +pub fn ceil(x: f64) -> f64 { + libm::ceil(x) +} + +#[cfg(feature = "std")] +#[inline] +#[allow(dead_code)] +pub fn floor(x: f64) -> f64 { + x.floor() +} + +#[cfg(not(feature = "std"))] +#[inline] +pub fn floor(x: f64) -> f64 { + libm::floor(x) +} + +#[cfg(not(feature = "std"))] +#[inline] +pub fn powi(x: f64, n: i32) -> f64 { + libm::pow(x, n as f64) +} diff --git a/flowstats/src/mod.rs b/flowstats/src/mod.rs new file mode 100644 index 00000000..fd42bbdc --- /dev/null +++ b/flowstats/src/mod.rs @@ -0,0 +1,37 @@ +// SPDX-License-Identifier: MIT +// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors +// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/mod.rs + +//! Frequency estimation algorithms +//! +//! This module provides implementations of sketches for estimating item +//! frequencies in a data stream. +//! +//! # Algorithms +//! +//! - [`CountMinSketch`]: Classic count-min sketch with optional conservative update +//! - [`SpaceSaving`]: Top-K / heavy hitters tracking +//! +//! # Example +//! +//! ``` +//! use flowstats::frequency::CountMinSketch; +//! use flowstats::traits::FrequencySketch; +//! +//! let mut cms = CountMinSketch::new(0.01, 0.001); // 1% error, 0.1% probability +//! +//! cms.add(b"item1", 5); +//! cms.add(b"item2", 3); +//! +//! let count = cms.estimate(b"item1"); +//! println!("Estimated count: {}", count); +//! ``` + +mod count_min; +#[cfg(feature = "std")] +mod space_saving; + +pub use count_min::CountMinSketch; + +#[cfg(feature = "std")] +pub use space_saving::SpaceSaving; diff --git a/flowstats/src/traits.rs b/flowstats/src/traits.rs new file mode 100644 index 00000000..43e7b35d --- /dev/null +++ b/flowstats/src/traits.rs @@ -0,0 +1,199 @@ +// SPDX-License-Identifier: MIT +// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors +// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/traits.rs + +//! Core traits for streaming algorithms +//! +//! All sketches implement the base [`Sketch`] trait, with specialized traits +//! for different algorithm families (cardinality, frequency, quantiles, etc.) + +use core::fmt::Debug; + +#[cfg(feature = "std")] +use std::{string::String}; + +#[cfg(not(feature = "std"))] +extern crate alloc; +#[cfg(not(feature = "std"))] +use alloc::{string::String, vec::Vec}; + +/// Error during sketch merge operation +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum MergeError { + /// Sketches have incompatible configurations + IncompatibleConfig { + expected: String, + found: String, + }, + /// Sketches have incompatible versions + VersionMismatch { + expected: u32, + found: u32, + }, +} + +impl core::fmt::Display for MergeError { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + MergeError::IncompatibleConfig { expected, found } => { + write!(f, "incompatible config: expected {}, found {}", expected, found) + } + MergeError::VersionMismatch { expected, found } => { + write!(f, "version mismatch: expected {}, found {}", expected, found) + } + } + } +} + +#[cfg(feature = "std")] +impl std::error::Error for MergeError {} + +/// Error during sketch decoding +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum DecodeError { + /// Input buffer too short + BufferTooShort { expected: usize, found: usize }, + /// Invalid magic number or header + InvalidHeader, + /// Unsupported version + UnsupportedVersion(u32), + /// Corrupted data + Corrupted(String), +} + +impl core::fmt::Display for DecodeError { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + DecodeError::BufferTooShort { expected, found } => { + write!(f, "buffer too short: expected {}, found {}", expected, found) + } + DecodeError::InvalidHeader => write!(f, "invalid header"), + DecodeError::UnsupportedVersion(v) => write!(f, "unsupported version: {}", v), + DecodeError::Corrupted(msg) => write!(f, "corrupted data: {}", msg), + } + } +} + +#[cfg(feature = "std")] +impl std::error::Error for DecodeError {} + +/// Error bounds for a sketch estimate +#[derive(Debug, Clone, Copy, PartialEq)] +pub struct ErrorBounds { + /// Lower bound of the estimate + pub lower: f64, + /// Point estimate + pub estimate: f64, + /// Upper bound of the estimate + pub upper: f64, + /// Confidence level (e.g., 0.95 for 95%) + pub confidence: f64, +} + +impl ErrorBounds { + /// Create new error bounds + pub fn new(lower: f64, estimate: f64, upper: f64, confidence: f64) -> Self { + Self { + lower, + estimate, + upper, + confidence, + } + } + + /// Check if a value falls within bounds + pub fn contains(&self, value: f64) -> bool { + value >= self.lower && value <= self.upper + } + + /// Width of the confidence interval + pub fn width(&self) -> f64 { + self.upper - self.lower + } + + /// Relative width (width / estimate) + pub fn relative_width(&self) -> f64 { + if self.estimate == 0.0 { + 0.0 + } else { + self.width() / self.estimate + } + } +} + +/// Core trait for all streaming sketches +pub trait Sketch: Clone + Debug { + /// The type of item this sketch processes + type Item: ?Sized; + + /// Add an item to the sketch + fn update(&mut self, item: &Self::Item); + + /// Merge another sketch into this one + /// + /// Returns an error if sketches are incompatible + fn merge(&mut self, other: &Self) -> Result<(), MergeError>; + + /// Reset sketch to empty state + fn clear(&mut self); + + /// Memory usage in bytes + fn size_bytes(&self) -> usize; + + /// Number of items processed + fn count(&self) -> u64; + + /// Check if sketch is empty + fn is_empty(&self) -> bool { + self.count() == 0 + } +} + +/// Cardinality (distinct count) estimation sketches +pub trait CardinalitySketch: Sketch { + /// Estimate number of distinct items seen + fn estimate(&self) -> f64; + + /// Get error bounds at given confidence level (0.0 to 1.0) + fn error_bounds(&self, confidence: f64) -> ErrorBounds; + + /// Relative standard error (RSE) of the estimate + /// + /// RSE = standard_error / true_value ≈ 1.04 / sqrt(m) for HLL + fn relative_error(&self) -> f64; + + /// Estimate with default 95% confidence bounds + fn estimate_with_bounds(&self) -> ErrorBounds { + self.error_bounds(0.95) + } +} + +/// Frequency estimation sketches +pub trait FrequencySketch: Sketch { + /// Estimate frequency of an item + fn estimate_frequency(&self, item: &Self::Item) -> u64; + + /// Check if frequency exceeds threshold + fn exceeds_threshold(&self, item: &Self::Item, threshold: u64) -> bool { + self.estimate_frequency(item) >= threshold + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_error_bounds() { + let bounds = ErrorBounds::new(90.0, 100.0, 110.0, 0.95); + + assert!(bounds.contains(100.0)); + assert!(bounds.contains(90.0)); + assert!(bounds.contains(110.0)); + assert!(!bounds.contains(89.0)); + assert!(!bounds.contains(111.0)); + + assert_eq!(bounds.width(), 20.0); + assert!((bounds.relative_width() - 0.2).abs() < 0.001); + } +} From 9c686de66e026f2e2faceda5ef702ce94b7adef9 Mon Sep 17 00:00:00 2001 From: Zach McCoy Date: Mon, 21 Sep 2026 15:05:11 -0500 Subject: [PATCH 2/3] Flowstats compiles internally, it is now used in the cargo.toml file Signed-off-by: Zach McCoy --- Cargo.toml | 2 +- flowstats/benches/benchmarks.rs | 256 +-------------- flowstats/src/frequency/mod.rs | 5 - flowstats/src/frequency/space_saving.rs | 414 ------------------------ flowstats/src/lib.rs | 4 - flowstats/src/mod.rs | 37 --- 6 files changed, 8 insertions(+), 710 deletions(-) delete mode 100644 flowstats/src/frequency/space_saving.rs delete mode 100644 flowstats/src/mod.rs diff --git a/Cargo.toml b/Cargo.toml index 63468b41..7e3ce4d6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -20,7 +20,7 @@ lazy_static = "1.4.0" libc = "0.2" serde = { version = "1.0", features = ["derive"] } bincode = "1.3" -flowstats = { version = "0.1.2", features = ["frequency"]} +flowstats = { path = "flowstats"} [dev-dependencies] rand = "0.8" diff --git a/flowstats/benches/benchmarks.rs b/flowstats/benches/benchmarks.rs index 047a4d89..3a4551d4 100644 --- a/flowstats/benches/benchmarks.rs +++ b/flowstats/benches/benchmarks.rs @@ -1,70 +1,22 @@ +// SPDX-License-Identifier: MIT +// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors +// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/benches/benchmarks.rs + + //! Benchmarks for flowstats algorithms //! //! Run with: cargo bench --features full // Require all features for benchmarks #[cfg(not(all( - feature = "cardinality", feature = "frequency", - feature = "quantiles", - feature = "membership", - feature = "sampling", - feature = "statistics" )))] compile_error!("Benchmarks require all features. Run: cargo bench --features full"); use criterion::{black_box, criterion_group, criterion_main, Criterion, Throughput}; -use flowstats::cardinality::HyperLogLog; -use flowstats::frequency::{CountMinSketch, SpaceSaving}; -use flowstats::membership::BloomFilter; -use flowstats::quantiles::TDigest; -use flowstats::sampling::ReservoirSampler; -use flowstats::statistics::RunningStats; -use flowstats::traits::{CardinalitySketch, HeavyHitters, QuantileSketch, Sketch}; - -// ============================================================================ -// HyperLogLog Benchmarks -// ============================================================================ - -fn bench_hll(c: &mut Criterion) { - let mut group = c.benchmark_group("hyperloglog"); - group.throughput(Throughput::Elements(1)); - - for precision in [10, 12, 14, 16] { - group.bench_function(format!("insert_p{}", precision), |b| { - let mut hll = HyperLogLog::new(precision); - let mut i = 0u64; - b.iter(|| { - hll.insert(&i.to_string()); - i = i.wrapping_add(1); - }); - }); - } - - group.bench_function("estimate", |b| { - let mut hll = HyperLogLog::new(14); - for i in 0..100_000u64 { - hll.insert(&i.to_string()); - } - b.iter(|| black_box(hll.estimate())); - }); - - group.bench_function("merge", |b| { - let mut hll1 = HyperLogLog::new(14); - let mut hll2 = HyperLogLog::new(14); - for i in 0..10_000u64 { - hll1.insert(&i.to_string()); - hll2.insert(&(i + 10_000).to_string()); - } - b.iter(|| { - let mut h = hll1.clone(); - h.merge(black_box(&hll2)).unwrap(); - }); - }); - - group.finish(); -} +use flowstats::frequency::{CountMinSketch}; +use flowstats::traits::{Sketch}; // ============================================================================ // Count-Min Sketch Benchmarks @@ -107,207 +59,13 @@ fn bench_cms(c: &mut Criterion) { group.finish(); } -// ============================================================================ -// Space-Saving Benchmarks -// ============================================================================ - -fn bench_space_saving(c: &mut Criterion) { - let mut group = c.benchmark_group("space_saving"); - group.throughput(Throughput::Elements(1)); - - for k in [100, 1000] { - group.bench_function(format!("add_k{}", k), |b| { - let mut ss: SpaceSaving = SpaceSaving::new(k); - let mut i = 0u64; - b.iter(|| { - ss.add(i.to_string()); - i = i.wrapping_add(1); - }); - }); - } - - group.bench_function("top_k", |b| { - let mut ss: SpaceSaving = SpaceSaving::new(1000); - for i in 0..100_000u64 { - ss.add((i % 10_000).to_string()); - } - b.iter(|| black_box(ss.top_k(10))); - }); - - group.finish(); -} - -// ============================================================================ -// t-digest Benchmarks -// ============================================================================ - -fn bench_tdigest(c: &mut Criterion) { - let mut group = c.benchmark_group("tdigest"); - group.throughput(Throughput::Elements(1)); - - for compression in [50.0, 100.0, 200.0] { - group.bench_function(format!("add_c{}", compression as u32), |b| { - let mut td = TDigest::new(compression); - let mut i = 0u64; - b.iter(|| { - td.add((i as f64) * 0.001); - i = i.wrapping_add(1); - }); - }); - } - - group.bench_function("quantile", |b| { - let mut td = TDigest::new(100.0); - for i in 0..100_000u64 { - td.add(i as f64); - } - b.iter(|| black_box(td.quantile(0.99))); - }); - - group.bench_function("merge", |b| { - let mut td1 = TDigest::new(100.0); - let mut td2 = TDigest::new(100.0); - for i in 0..10_000u64 { - td1.add(i as f64); - td2.add((i + 10_000) as f64); - } - b.iter(|| { - let mut t = td1.clone(); - t.merge(black_box(&td2)).unwrap(); - }); - }); - - group.finish(); -} - -// ============================================================================ -// Bloom Filter Benchmarks -// ============================================================================ - -fn bench_bloom(c: &mut Criterion) { - let mut group = c.benchmark_group("bloom_filter"); - group.throughput(Throughput::Elements(1)); - - group.bench_function("insert", |b| { - let mut bloom = BloomFilter::new(1_000_000, 0.01); - let mut i = 0u64; - b.iter(|| { - bloom.insert(i.to_string().as_bytes()); - i = i.wrapping_add(1); - }); - }); - - group.bench_function("contains_hit", |b| { - let mut bloom = BloomFilter::new(100_000, 0.01); - for i in 0..100_000u64 { - bloom.insert(i.to_string().as_bytes()); - } - let mut i = 0u64; - b.iter(|| { - let result = bloom.contains((i % 100_000).to_string().as_bytes()); - i = i.wrapping_add(1); - black_box(result) - }); - }); - - group.bench_function("contains_miss", |b| { - let mut bloom = BloomFilter::new(100_000, 0.01); - for i in 0..100_000u64 { - bloom.insert(i.to_string().as_bytes()); - } - let mut i = 1_000_000u64; - b.iter(|| { - let result = bloom.contains(i.to_string().as_bytes()); - i = i.wrapping_add(1); - black_box(result) - }); - }); - - group.finish(); -} - -// ============================================================================ -// Reservoir Sampler Benchmarks -// ============================================================================ - -fn bench_reservoir(c: &mut Criterion) { - let mut group = c.benchmark_group("reservoir"); - group.throughput(Throughput::Elements(1)); - - for capacity in [100, 1000, 10000] { - group.bench_function(format!("add_cap{}", capacity), |b| { - let mut sampler: ReservoirSampler = ReservoirSampler::new(capacity); - let mut i = 0u64; - b.iter(|| { - sampler.add(i); - i = i.wrapping_add(1); - }); - }); - } - - group.finish(); -} - -// ============================================================================ -// Running Stats Benchmarks -// ============================================================================ - -fn bench_running_stats(c: &mut Criterion) { - let mut group = c.benchmark_group("running_stats"); - group.throughput(Throughput::Elements(1)); - - group.bench_function("add", |b| { - let mut stats = RunningStats::new(); - let mut i = 0u64; - b.iter(|| { - stats.add(i as f64); - i = i.wrapping_add(1); - }); - }); - - group.bench_function("query_all", |b| { - let mut stats = RunningStats::new(); - for i in 0..100_000u64 { - stats.add(i as f64); - } - b.iter(|| { - black_box(stats.mean()); - black_box(stats.variance()); - black_box(stats.stddev()); - black_box(stats.min()); - black_box(stats.max()); - }); - }); - - group.bench_function("merge", |b| { - let mut s1 = RunningStats::new(); - let mut s2 = RunningStats::new(); - for i in 0..10_000u64 { - s1.add(i as f64); - s2.add((i + 10_000) as f64); - } - b.iter(|| { - let mut s = s1.clone(); - s.merge(black_box(&s2)).unwrap(); - }); - }); - - group.finish(); -} - // ============================================================================ // Main // ============================================================================ criterion_group!( benches, - bench_hll, bench_cms, - bench_space_saving, - bench_tdigest, - bench_bloom, - bench_reservoir, - bench_running_stats, ); criterion_main!(benches); diff --git a/flowstats/src/frequency/mod.rs b/flowstats/src/frequency/mod.rs index ec1d04a8..754c60b3 100644 --- a/flowstats/src/frequency/mod.rs +++ b/flowstats/src/frequency/mod.rs @@ -28,10 +28,5 @@ //! ``` mod count_min; -#[cfg(feature = "std")] -mod space_saving; pub use count_min::CountMinSketch; - -#[cfg(feature = "std")] -pub use space_saving::SpaceSaving; diff --git a/flowstats/src/frequency/space_saving.rs b/flowstats/src/frequency/space_saving.rs deleted file mode 100644 index 30fdbe45..00000000 --- a/flowstats/src/frequency/space_saving.rs +++ /dev/null @@ -1,414 +0,0 @@ -// SPDX-License-Identifier: MIT -// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors -// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/frequency/space_saving.rs - -//! Space-Saving algorithm for heavy hitters -//! -//! The Space-Saving algorithm efficiently tracks the k most frequent items -//! in a data stream using only O(k) space. -//! -//! **Note**: This module requires the `std` feature. The Space-Saving algorithm -//! does not support correct distributed merging due to its replacement strategy. -//! Attempting to merge will return an error. - -use crate::traits::{FrequencySketch, MergeError, Sketch}; -use core::hash::Hash; -use std::collections::HashMap; -use std::vec::Vec; - -/// Entry in the Space-Saving structure -#[derive(Clone, Debug)] -struct Counter { - /// The item - item: T, - /// Estimated count - count: u64, - /// Error bound (maximum overcount) - error: u64, -} - -impl Counter { - fn new(item: T, count: u64, error: u64) -> Self { - Self { item, count, error } - } -} - -/// Space-Saving algorithm for finding frequent items -/// -/// The Space-Saving algorithm maintains a summary of the k most frequent items -/// with the following guarantees: -/// -/// - Any item with true frequency > n/k is guaranteed to be in the summary -/// - The maximum overcount error for any item is at most n/k -/// -/// # Example -/// -/// ``` -/// use flowstats::frequency::SpaceSaving; -/// use flowstats::traits::HeavyHitters; -/// -/// let mut ss = SpaceSaving::new(10); // Track top 10 -/// -/// // Add some items -/// for _ in 0..100 { ss.add("apple"); } -/// for _ in 0..50 { ss.add("banana"); } -/// for _ in 0..25 { ss.add("cherry"); } -/// for _ in 0..10 { ss.add("date"); } -/// -/// // Get top 3 items -/// let top = ss.top_k(3); -/// println!("Top items: {:?}", top); -/// ``` -#[derive(Clone, Debug)] -pub struct SpaceSaving { - /// Maximum number of counters to maintain - capacity: usize, - /// Map from item to counter index - item_to_index: HashMap, - /// Counter array - counters: Vec>, - /// Total count of all items - total_count: u64, - /// Number of updates - num_updates: u64, -} - -impl SpaceSaving { - /// Create a new Space-Saving structure with the given capacity - /// - /// # Arguments - /// - /// * `capacity` - Maximum number of items to track (k) - pub fn new(capacity: usize) -> Self { - assert!(capacity > 0, "capacity must be positive"); - - Self { - capacity, - item_to_index: HashMap::with_capacity(capacity), - counters: Vec::with_capacity(capacity), - total_count: 0, - num_updates: 0, - } - } - - /// Get the capacity (k) - pub fn capacity(&self) -> usize { - self.capacity - } - - /// Get the number of distinct items currently tracked - pub fn num_tracked(&self) -> usize { - self.counters.len() - } - - /// Get the total count - pub fn total_count(&self) -> u64 { - self.total_count - } - - /// Add an item to the structure - pub fn add(&mut self, item: T) { - self.add_count(item, 1); - } - - /// Add an item with a specific count - pub fn add_count(&mut self, item: T, count: u64) { - self.num_updates += 1; - self.total_count += count; - - // Check if item is already tracked - if let Some(&idx) = self.item_to_index.get(&item) { - self.counters[idx].count += count; - return; - } - - // Not tracked - either add new or replace minimum - if self.counters.len() < self.capacity { - // Still have room, add new counter - let idx = self.counters.len(); - self.counters.push(Counter::new(item.clone(), count, 0)); - self.item_to_index.insert(item, idx); - } else { - // Find and replace minimum counter - let min_idx = self.find_min_index(); - let min_count = self.counters[min_idx].count; - - // Remove old item from map - let old_item = self.counters[min_idx].item.clone(); - self.item_to_index.remove(&old_item); - - // Replace with new item - self.counters[min_idx] = Counter::new(item.clone(), min_count + count, min_count); - self.item_to_index.insert(item, min_idx); - } - } - - /// Find the index of the counter with minimum count - fn find_min_index(&self) -> usize { - self.counters - .iter() - .enumerate() - .min_by_key(|(_, c)| c.count) - .map(|(i, _)| i) - .unwrap_or(0) - } - - /// Estimate the frequency of an item - pub fn estimate(&self, item: &T) -> u64 { - self.item_to_index - .get(item) - .map(|&idx| self.counters[idx].count) - .unwrap_or(0) - } - - /// Get the error bound for an item's estimate - pub fn error(&self, item: &T) -> u64 { - self.item_to_index - .get(item) - .map(|&idx| self.counters[idx].error) - .unwrap_or(0) - } - - /// Get guaranteed minimum count for an item - /// - /// Returns (count - error), which is guaranteed to be at most the true count. - pub fn guaranteed_count(&self, item: &T) -> u64 { - self.item_to_index - .get(item) - .map(|&idx| { - let c = &self.counters[idx]; - c.count.saturating_sub(c.error) - }) - .unwrap_or(0) - } - - /// Check if an item is currently tracked - pub fn contains(&self, item: &T) -> bool { - self.item_to_index.contains_key(item) - } -} - -impl Sketch for SpaceSaving { - type Item = T; - - fn update(&mut self, item: &T) { - self.add(item.clone()); - } - - fn merge(&mut self, _other: &Self) -> Result<(), MergeError> { - // Space-Saving does not support correct distributed merge. - // The algorithm's replacement strategy makes merging two summaries - // produce incorrect results (inflated counts, dropped heavy hitters). - // Use a single SpaceSaving instance or aggregate raw data instead. - Err(MergeError::IncompatibleConfig { - expected: "Space-Saving does not support merge".into(), - found: "merge attempted".into(), - }) - } - - fn clear(&mut self) { - self.item_to_index.clear(); - self.counters.clear(); - self.total_count = 0; - self.num_updates = 0; - } - - fn size_bytes(&self) -> usize { - core::mem::size_of::() + self.counters.capacity() * core::mem::size_of::>() - } - - fn count(&self) -> u64 { - self.num_updates - } -} - -impl FrequencySketch for SpaceSaving { - fn estimate_frequency(&self, item: &T) -> u64 { - self.estimate(item) - } -} - -#[cfg(feature = "serde")] -impl serde::Serialize - for SpaceSaving -{ - fn serialize(&self, serializer: S) -> Result - where - S: serde::Serializer, - { - use serde::ser::SerializeStruct; - let items: Vec<_> = self - .counters - .iter() - .map(|c| (&c.item, c.count, c.error)) - .collect(); - - let mut state = serializer.serialize_struct("SpaceSaving", 3)?; - state.serialize_field("capacity", &self.capacity)?; - state.serialize_field("total_count", &self.total_count)?; - state.serialize_field("items", &items)?; - state.end() - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_basic() { - let mut ss = SpaceSaving::::new(10); - - ss.add("apple".to_string()); - ss.add("apple".to_string()); - ss.add("banana".to_string()); - - assert!(ss.estimate(&"apple".to_string()) >= 2); - assert!(ss.estimate(&"banana".to_string()) >= 1); - } - - #[test] - fn test_empty() { - let ss = SpaceSaving::::new(10); - assert_eq!(ss.estimate(&"anything".to_string()), 0); - assert_eq!(ss.total_count(), 0); - } - - #[test] - fn test_top_k() { - let mut ss = SpaceSaving::<&str>::new(10); - - for _ in 0..100 { - ss.add("apple"); - } - for _ in 0..50 { - ss.add("banana"); - } - for _ in 0..25 { - ss.add("cherry"); - } - - let top = ss.top_k(2); - assert_eq!(top.len(), 2); - assert_eq!(top[0].0, "apple"); - assert_eq!(top[1].0, "banana"); - } - - #[test] - fn test_heavy_hitters() { - let mut ss = SpaceSaving::<&str>::new(10); - - for _ in 0..100 { - ss.add("apple"); - } - for _ in 0..10 { - ss.add("banana"); - } - for _ in 0..1 { - ss.add("cherry"); - } - - // Items with > 5% of total - let heavy = ss.heavy_hitters(0.05); - assert!(heavy.iter().any(|(item, _)| *item == "apple")); - assert!(heavy.iter().any(|(item, _)| *item == "banana")); - // Cherry should not be included (1/111 < 5%) - } - - #[test] - fn test_replacement() { - let mut ss = SpaceSaving::::new(3); - - // Fill up capacity - ss.add(1); - ss.add(2); - ss.add(3); - - assert_eq!(ss.num_tracked(), 3); - - // Add a new item - should replace the minimum - ss.add(4); - - // Still only 3 items tracked - assert_eq!(ss.num_tracked(), 3); - } - - #[test] - fn test_contains() { - let mut ss = SpaceSaving::<&str>::new(10); - - ss.add("apple"); - - assert!(ss.contains(&"apple")); - assert!(!ss.contains(&"banana")); - } - - #[test] - fn test_merge_not_supported() { - let mut ss1 = SpaceSaving::<&str>::new(10); - let ss2 = SpaceSaving::<&str>::new(10); - - // Space-Saving does not support merge - assert!(ss1.merge(&ss2).is_err()); - } - - #[test] - fn test_guaranteed_count() { - let mut ss = SpaceSaving::<&str>::new(3); - - // Add items that will cause replacements - ss.add("a"); - ss.add("b"); - ss.add("c"); - - // Add more to force replacement - ss.add("d"); - - // Check that guaranteed count is <= estimated count - for item in ["a", "b", "c", "d"] { - let est = ss.estimate(&item); - let guar = ss.guaranteed_count(&item); - assert!( - guar <= est, - "guaranteed {} > estimate {} for {}", - guar, - est, - item - ); - } - } - - #[test] - fn test_zipf_distribution() { - // Simulate Zipf distribution (common in real data) - let mut ss = SpaceSaving::::new(10); - - // Item 1 appears 1000 times, item 2 appears 500 times, etc. - for rank in 1..=100 { - let count = 1000 / rank; - for _ in 0..count { - ss.add(rank); - } - } - - // Top items should be 1, 2, 3, ... - let top = ss.top_k(5); - assert!(!top.is_empty()); - // Item 1 should definitely be tracked - assert!(ss.contains(&1)); - } - - #[test] - fn test_clear() { - let mut ss = SpaceSaving::<&str>::new(10); - - ss.add("apple"); - ss.add("banana"); - - ss.clear(); - - assert_eq!(ss.num_tracked(), 0); - assert_eq!(ss.total_count(), 0); - assert!(!ss.contains(&"apple")); - } -} diff --git a/flowstats/src/lib.rs b/flowstats/src/lib.rs index 836296fd..18d16d0c 100644 --- a/flowstats/src/lib.rs +++ b/flowstats/src/lib.rs @@ -94,12 +94,8 @@ pub mod prelude { #[cfg(feature = "frequency")] pub use crate::frequency::CountMinSketch; - #[cfg(all(feature = "frequency", feature = "std"))] - pub use crate::frequency::SpaceSaving; } #[cfg(feature = "frequency")] pub use frequency::CountMinSketch; -#[cfg(all(feature = "frequency", feature = "std"))] -pub use frequency::SpaceSaving; diff --git a/flowstats/src/mod.rs b/flowstats/src/mod.rs deleted file mode 100644 index fd42bbdc..00000000 --- a/flowstats/src/mod.rs +++ /dev/null @@ -1,37 +0,0 @@ -// SPDX-License-Identifier: MIT -// SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors -// SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/mod.rs - -//! Frequency estimation algorithms -//! -//! This module provides implementations of sketches for estimating item -//! frequencies in a data stream. -//! -//! # Algorithms -//! -//! - [`CountMinSketch`]: Classic count-min sketch with optional conservative update -//! - [`SpaceSaving`]: Top-K / heavy hitters tracking -//! -//! # Example -//! -//! ``` -//! use flowstats::frequency::CountMinSketch; -//! use flowstats::traits::FrequencySketch; -//! -//! let mut cms = CountMinSketch::new(0.01, 0.001); // 1% error, 0.1% probability -//! -//! cms.add(b"item1", 5); -//! cms.add(b"item2", 3); -//! -//! let count = cms.estimate(b"item1"); -//! println!("Estimated count: {}", count); -//! ``` - -mod count_min; -#[cfg(feature = "std")] -mod space_saving; - -pub use count_min::CountMinSketch; - -#[cfg(feature = "std")] -pub use space_saving::SpaceSaving; From 85bd0c7e347a66277bc9750ac77c283b8e88369c Mon Sep 17 00:00:00 2001 From: Zach McCoy Date: Fri, 25 Sep 2026 15:14:41 -0500 Subject: [PATCH 3/3] Fix SPDC licensing to match MIT or APACHE-2.0, add apache license, clean up comments to be around what is left of the code and nothing else Signed-off-by: Zach McCoy --- flowstats/LICENSE-APACHE | 190 +++++++++++++++++++++++++++ flowstats/src/frequency/count_min.rs | 2 +- flowstats/src/frequency/mod.rs | 3 +- flowstats/src/lib.rs | 11 +- flowstats/src/math.rs | 2 +- flowstats/src/traits.rs | 2 +- 6 files changed, 195 insertions(+), 15 deletions(-) create mode 100644 flowstats/LICENSE-APACHE diff --git a/flowstats/LICENSE-APACHE b/flowstats/LICENSE-APACHE new file mode 100644 index 00000000..4121bdfc --- /dev/null +++ b/flowstats/LICENSE-APACHE @@ -0,0 +1,190 @@ + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + +TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + +1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to the Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + +2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + +3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + +4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + +5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + +6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + +7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + +8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + +9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. + +END OF TERMS AND CONDITIONS + +Copyright 2025 flowstats Contributors + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. diff --git a/flowstats/src/frequency/count_min.rs b/flowstats/src/frequency/count_min.rs index 0d41f10e..ab57b1e5 100644 --- a/flowstats/src/frequency/count_min.rs +++ b/flowstats/src/frequency/count_min.rs @@ -1,4 +1,4 @@ -// SPDX-License-Identifier: MIT +// SPDX-License-Identifier: MIT OR Apache-2.0 // SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors // SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/frequency/count_min.rs diff --git a/flowstats/src/frequency/mod.rs b/flowstats/src/frequency/mod.rs index 754c60b3..fb58e9b1 100644 --- a/flowstats/src/frequency/mod.rs +++ b/flowstats/src/frequency/mod.rs @@ -1,4 +1,4 @@ -// SPDX-License-Identifier: MIT +// SPDX-License-Identifier: MIT OR Apache-2.0 // SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors // SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/frequency/mod.rs @@ -10,7 +10,6 @@ //! # Algorithms //! //! - [`CountMinSketch`]: Classic count-min sketch with optional conservative update -//! - [`SpaceSaving`]: Top-K / heavy hitters tracking //! //! # Example //! diff --git a/flowstats/src/lib.rs b/flowstats/src/lib.rs index 18d16d0c..a93d8337 100644 --- a/flowstats/src/lib.rs +++ b/flowstats/src/lib.rs @@ -1,4 +1,4 @@ -// SPDX-License-Identifier: MIT +// SPDX-License-Identifier: MIT OR Apache-2.0 // SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors // SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/lib.rs @@ -11,10 +11,7 @@ //! //! ## Features //! -//! - **Cardinality Estimation**: Count distinct elements with HyperLogLog //! - **Frequency Estimation**: Track item frequencies with Count-Min Sketch -//! - **Heavy Hitters**: Find top-K elements with Space-Saving -//! - **Quantile Estimation**: Compute percentiles with t-digest //! - **Full Mergeability**: All sketches support distributed merge operations //! - **Error Bounds**: Formal guarantees on approximation accuracy //! @@ -62,13 +59,7 @@ //! ## Feature Flags //! //! Algorithm families (pick what you need): -//! - `cardinality` (default): HyperLogLog for distinct counting //! - `frequency` (default): Count-Min Sketch, Space-Saving (tbd) -//! - `quantiles` (default): t-digest for percentiles -//! - `membership`: (default) Bloom filter -//! - `sampling`: Reservoir and weighted sampling (tbd) -//! - `sets`: Theta sketch for set operations (tbd) -//! - `statistics`: Running moments, entropy (tbd) //! - `full`: Enable all algorithm families //! //! Platform features: diff --git a/flowstats/src/math.rs b/flowstats/src/math.rs index da73f726..a40f4f07 100644 --- a/flowstats/src/math.rs +++ b/flowstats/src/math.rs @@ -1,4 +1,4 @@ -// SPDX-License-Identifier: MIT +// SPDX-License-Identifier: MIT OR Apache-2.0 // SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors // SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/math.rs diff --git a/flowstats/src/traits.rs b/flowstats/src/traits.rs index 43e7b35d..48e5e5ac 100644 --- a/flowstats/src/traits.rs +++ b/flowstats/src/traits.rs @@ -1,4 +1,4 @@ -// SPDX-License-Identifier: MIT +// SPDX-License-Identifier: MIT OR Apache-2.0 // SPDX-FileCopyrightText: Copyright (c) 2025 flowstats Contributors // SPDX-FileContributor: https://github.com/vnvo/flowstats/blob/v0.1.2/src/traits.rs