From 83ccc492c60d6cd171425939d0cc66e47fbcf2c0 Mon Sep 17 00:00:00 2001 From: zzylol Date: Tue, 29 Sep 2026 19:57:13 +0000 Subject: [PATCH] feat: add shared summary kernels and typed physical values Start `asap-physical-operators` with thin summary kernels over `asap_sketchlib`, the kernel capability checks and the typed value model. The boundary is: - `asap_sketchlib` owns sketch algorithms and their state encodings. - Kernels hold one population's in-memory state. They expose `merge`, a typed sketch readout (`estimate(&SketchQuery)`) and memory accounting. Exact states answer a typed `ExactReadout`; empty MIN/MAX read as `None`. - Group-by belongs to physical operators. - Deployments own wire decoding, delta frames, edge sampling and storage statistics. So wire decoding, `SerializableToSink`, `AggregationType`, `aux_stats`, `reset_to_empty` and the keyed/sum/min/max kernels are not carried over from ASAPQuery-backend. The `asap_sketch_codec` crate is not carried over either; it moves to `asap_sketchlib`. Hydra KLL remains as the Hydra shared-grouping kernel. HLL uses sketchlib's classic estimator. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 35 +- Cargo.toml | 1 + crates/asap-physical-operators/Cargo.toml | 16 + .../asap-physical-operators/src/capability.rs | 241 ++++++ crates/asap-physical-operators/src/error.rs | 17 + .../src/key_by_label_values.rs | 126 +++ crates/asap-physical-operators/src/lib.rs | 22 + .../src/measurement.rs | 48 ++ .../asap-physical-operators/src/statistic.rs | 67 ++ .../src/summary_kernels/count_min_sketch.rs | 76 ++ .../count_min_sketch_with_heap.rs | 234 +++++ .../src/summary_kernels/count_sketch.rs | 96 +++ .../summary_kernels/count_sketch_with_heap.rs | 143 ++++ .../src/summary_kernels/datasketches_kll.rs | 106 +++ .../src/summary_kernels/dd_sketch.rs | 88 ++ .../src/summary_kernels/exact.rs | 215 +++++ .../src/summary_kernels/factory.rs | 799 ++++++++++++++++++ .../src/summary_kernels/hll_sketch.rs | 79 ++ .../src/summary_kernels/hydra_kll.rs | 55 ++ .../src/summary_kernels/increase.rs | 213 +++++ .../src/summary_kernels/mod.rs | 26 + .../src/summary_kernels/traits.rs | 35 + .../src/summary_kernels/univmon.rs | 109 +++ .../src/summary_kernels/weighted_frequency.rs | 152 ++++ crates/asap-physical-operators/src/values.rs | 381 +++++++++ .../tests/deployment.rs | 90 ++ 26 files changed, 3468 insertions(+), 2 deletions(-) create mode 100644 crates/asap-physical-operators/Cargo.toml create mode 100644 crates/asap-physical-operators/src/capability.rs create mode 100644 crates/asap-physical-operators/src/error.rs create mode 100644 crates/asap-physical-operators/src/key_by_label_values.rs create mode 100644 crates/asap-physical-operators/src/lib.rs create mode 100644 crates/asap-physical-operators/src/measurement.rs create mode 100644 crates/asap-physical-operators/src/statistic.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/count_min_sketch.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/count_min_sketch_with_heap.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/count_sketch.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/count_sketch_with_heap.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/datasketches_kll.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/dd_sketch.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/exact.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/factory.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/hll_sketch.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/hydra_kll.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/increase.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/mod.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/traits.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/univmon.rs create mode 100644 crates/asap-physical-operators/src/summary_kernels/weighted_frequency.rs create mode 100644 crates/asap-physical-operators/src/values.rs create mode 100644 crates/asap-physical-operators/tests/deployment.rs diff --git a/Cargo.lock b/Cargo.lock index d28e1bc5..b7b5e609 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -309,7 +309,7 @@ version = "0.1.0" dependencies = [ "asap-frontend-promql", "asap-types", - "asap_sketchlib", + "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib)", "serde", "serde_json", "thiserror 2.0.18", @@ -369,11 +369,24 @@ dependencies = [ "asap-frontend-promql", "asap-frontend-sql", "asap-types", - "asap_sketchlib", + "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib)", "serde_json", "tokio", ] +[[package]] +name = "asap-physical-operators" +version = "0.1.0" +dependencies = [ + "asap-aware-mapping", + "asap-frontend-promql", + "asap-types", + "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=5f03ccbd798ed5fec62bdd839bcb331123cab369)", + "serde", + "thiserror 2.0.18", + "tracing", +] + [[package]] name = "asap-sql-function-catalog" version = "0.1.0" @@ -387,6 +400,24 @@ dependencies = [ "thiserror 2.0.18", ] +[[package]] +name = "asap_sketchlib" +version = "0.3.0" +source = "git+https://github.com/ProjectASAP/asap_sketchlib?rev=5f03ccbd798ed5fec62bdd839bcb331123cab369#5f03ccbd798ed5fec62bdd839bcb331123cab369" +dependencies = [ + "bincode", + "bytes", + "prost", + "rand 0.9.5", + "rmp-serde", + "serde", + "serde-big-array", + "serde_bytes", + "smallvec", + "twox-hash 2.1.2", + "xxhash-rust", +] + [[package]] name = "asap_sketchlib" version = "0.3.0" diff --git a/Cargo.toml b/Cargo.toml index 9d44070d..e81894e3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,6 @@ [workspace] members = [ + "crates/asap-physical-operators", "crates/types", "crates/sql-function-catalog", "crates/asap-aware-mapping", diff --git a/crates/asap-physical-operators/Cargo.toml b/crates/asap-physical-operators/Cargo.toml new file mode 100644 index 00000000..51dd6a04 --- /dev/null +++ b/crates/asap-physical-operators/Cargo.toml @@ -0,0 +1,16 @@ +[package] +name = "asap-physical-operators" +version = "0.1.0" +edition = "2021" + +[dependencies] +planner-types = { package = "asap-types", path = "../types" } +asap_sketchlib = { git = "https://github.com/ProjectASAP/asap_sketchlib", rev = "5f03ccbd798ed5fec62bdd839bcb331123cab369" } +serde = { version = "1", features = ["derive", "rc"] } +tracing = "0.1" +thiserror = "2" + + +[dev-dependencies] +asap-aware-mapping = { path = "../asap-aware-mapping" } +asap-frontend-promql = { path = "../frontend-promql" } diff --git a/crates/asap-physical-operators/src/capability.rs b/crates/asap-physical-operators/src/capability.rs new file mode 100644 index 00000000..2a1cc095 --- /dev/null +++ b/crates/asap-physical-operators/src/capability.rs @@ -0,0 +1,241 @@ +//! Capability boundaries, checked without constructing accumulator state. +//! +//! `validate_summary_kernel` checks update kernels, including families without a +//! native batch representation. `validate_native_family` and +//! `validate_sketch_readout` / `validate_exact_readout` check native state and readout support. +//! Keyed weighted-frequency readouts are checked by `Operator::keyed_readout`. +//! A successful kernel check alone does not mean a physical DAG will bind. +//! +//! Stored-state encodings belong to deployments. Full plan acceptance is +//! owned by `binding`, which also validates schemas, expressions and inputs. +use crate::Error; +use planner_types::post_asap::{ + ExactKind, ExactParams, GroupingStrategy, SketchAlgorithm, SketchParams, SketchQuery, + SummaryFamilyType, SummaryUpdate, +}; + +/// Check the same contract used by `create_planner_accumulator` before a plan +/// is accepted. Execution timing is deliberately not a kernel property. +pub fn validate_summary_kernel( + family: &SummaryFamilyType, + input: &SummaryUpdate, + grouping: &GroupingStrategy, +) -> Result<(), String> { + if grouping != &GroupingStrategy::PerSubpopulationInstance { + return Err("shared summary grouping has no registered kernel".into()); + } + let keyed = match family { + SummaryFamilyType::ExactAggregate(kind, params) => { + use ExactKind as K; + use ExactParams as P; + if !matches!( + (kind, params), + (K::Sum, P::Sum) + | (K::Count, P::Count) + | (K::Min, P::Min) + | (K::Max, P::Max) + | (K::Rate, P::Rate) + | (K::Increase, P::Increase) + ) { + return Err(format!("unsupported exact kernel {family:?}")); + } + input.item.is_some() + } + SummaryFamilyType::Sketch(kind, layout) => { + if layout != grouping { + return Err("Planner family and operator grouping disagree".into()); + } + use SketchAlgorithm as A; + use SketchParams as P; + match (kind.algorithm(), kind.params()) { + (A::Kll, P::Kll { k }) if (8..=u16::MAX as u32).contains(k) => false, + (A::DDSketch, P::DDSketch { alpha }) + if alpha.is_finite() && *alpha > 0.0 && *alpha < 1.0 => + { + false + } + (A::Hll, P::Hll { precision }) if (4..=18).contains(precision) => false, + (A::Cms, P::Cms { width, depth }) + | (A::CountSketch, P::CountSketch { width, depth }) + if valid_matrix(*width, *depth) => + { + true + } + ( + A::CmsWithHeap, + P::CmsWithHeap { + width, + depth, + heap_size, + }, + ) + | ( + A::CountSketchWithHeap, + P::CountSketchWithHeap { + width, + depth, + heap_size, + }, + ) if valid_matrix(*width, *depth) && *heap_size > 0 => true, + ( + A::UnivMon, + P::UnivMon { + heap_size, + sketch_rows, + sketch_cols, + layers, + }, + ) if *heap_size > 0 + && *sketch_cols > 0 + && (1..=20).contains(sketch_rows) + && (1..=64).contains(layers) + && (*sketch_rows as usize) + .checked_mul(*sketch_cols as usize) + .and_then(|n| n.checked_mul(*layers as usize)) + .is_some() => + { + false + } + _ => { + return Err(format!( + "unsupported kernel or invalid parameters: {kind:?}" + )) + } + } + } + _ => return Err(format!("unsupported summary kernel {family:?}")), + }; + if keyed != input.item.is_some() && !is_unit_sample_frequency(input) { + return Err("Planner item expression does not match kernel layout".into()); + } + Ok(()) +} + +fn valid_matrix(width: u32, depth: u32) -> bool { + // Construction uses the kernel's native row hashing, so no encoded-size + // limit applies here. + width > 0 + && depth > 0 + && (width as usize) + .checked_mul(depth as usize) + .and_then(|n| n.checked_mul(std::mem::size_of::())) + .is_some() +} + +pub(crate) fn is_unit_sample_frequency(update: &planner_types::post_asap::SummaryUpdate) -> bool { + use planner_types::post_asap::{NonNegativeWeightProof, SummaryInputExpr, WeightDomain}; + matches!( + update.item, + Some(SummaryInputExpr::Column( + planner_types::pre_asap::ColumnRef::SampleValue + )) + ) && matches!(update.weight, SummaryInputExpr::Constant(1.0)) + && matches!( + update.weight_domain, + WeightDomain::NonNegative { + proof: NonNegativeWeightProof::UnitCount + } + ) +} + +pub fn validate_native_family(family: &SummaryFamilyType) -> Result<(), Error> { + use planner_types::post_asap::SketchAlgorithm as A; + if let SummaryFamilyType::Sketch(kind, grouping) = family { + if matches!(kind.algorithm(), A::CmsWithHeap | A::CountSketchWithHeap) { + let (_, width, depth, _) = + crate::summary_kernels::weighted_frequency::WeightedFrequency::configuration(kind)?; + return if valid_matrix(width as u32, depth as u32) && grouping == &Default::default() { + Ok(()) + } else { + Err(Error::Invalid( + "invalid weighted frequency dimensions or grouping strategy".into(), + )) + }; + } + } + match family { + SummaryFamilyType::ExactAggregate(..) => {} + SummaryFamilyType::Sketch(kind, _) + if matches!(kind.algorithm(), A::Kll | A::DDSketch | A::Hll) => {} + _ => { + return Err(Error::Invalid( + "summary family has no native DAG state implementation".into(), + )) + } + } + crate::capability::validate_summary_kernel( + family, + &planner_types::post_asap::SummaryUpdate::column( + planner_types::pre_asap::ColumnRef::SampleValue, + ), + &Default::default(), + ) + .map_err(Error::Invalid) +} + +/// A sketch readout is native only for the families Planner can read directly. +pub fn validate_sketch_readout( + family: &SummaryFamilyType, + query: &SketchQuery, +) -> Result<(), Error> { + validate_native_family(family)?; + use planner_types::post_asap::SketchAlgorithm as A; + // A point count without an item value reads the total count. + let bare_count = matches!(query, SketchQuery::PointCount { value: None, .. }); + let supported = match family { + SummaryFamilyType::Sketch(kind, _) => match (kind.algorithm(), query) { + (A::Kll, SketchQuery::Quantile { q }) | (A::DDSketch, SketchQuery::Quantile { q }) => { + if !(0.0..=1.0).contains(q) { + return Err(Error::Invalid( + "quantile readout requires quantile in [0,1]".into(), + )); + } + true + } + (A::DDSketch, _) => bare_count, + (A::Hll, SketchQuery::Cardinality) => true, + (A::Hll, _) => bare_count, + _ => false, + }, + _ => false, + }; + if !supported { + return Err(Error::Invalid( + "readout is not implemented for this summary family".into(), + )); + } + Ok(()) +} + +/// An exact readout must match the exact family it reads. +pub fn validate_exact_readout( + family: &SummaryFamilyType, + readout: &crate::summary_kernels::exact::ExactReadout, +) -> Result<(), Error> { + validate_native_family(family)?; + use crate::Statistic as S; + use planner_types::post_asap::ExactKind as E; + let supported = matches!( + (family, readout.statistic), + (SummaryFamilyType::ExactAggregate(E::Sum, _), S::Sum) + | (SummaryFamilyType::ExactAggregate(E::Count, _), S::Count) + | (SummaryFamilyType::ExactAggregate(E::Min, _), S::Min) + | (SummaryFamilyType::ExactAggregate(E::Max, _), S::Max) + | (SummaryFamilyType::ExactAggregate(E::Rate, _), S::Rate) + | ( + SummaryFamilyType::ExactAggregate(E::Increase, _), + S::Increase + ) + ); + if !supported { + return Err(Error::Invalid( + "readout is not implemented for this summary family".into(), + )); + } + if readout.lookback_ms.is_some_and(|lookback| { + lookback <= 0 || !matches!(readout.statistic, S::Rate | S::Increase) + }) { + return Err(Error::Invalid("invalid exact counter lookback".into())); + } + Ok(()) +} diff --git a/crates/asap-physical-operators/src/error.rs b/crates/asap-physical-operators/src/error.rs new file mode 100644 index 00000000..06ce16e9 --- /dev/null +++ b/crates/asap-physical-operators/src/error.rs @@ -0,0 +1,17 @@ +#[derive(Clone, Debug, PartialEq, Eq, thiserror::Error)] +pub enum Error { + #[error("invalid DAG: {0}")] + Invalid(String), + #[error("operator failed: {0}")] + Operator(String), + #[error("node {node} ({operation}) failed: {source}")] + AtNode { + node: u64, + operation: String, + source: Box, + }, + #[error("execution memory limit exceeded")] + MemoryLimit, + #[error("execution cancelled")] + Cancelled, +} diff --git a/crates/asap-physical-operators/src/key_by_label_values.rs b/crates/asap-physical-operators/src/key_by_label_values.rs new file mode 100644 index 00000000..e574da23 --- /dev/null +++ b/crates/asap-physical-operators/src/key_by_label_values.rs @@ -0,0 +1,126 @@ +use serde::{Deserialize, Serialize}; +// use std::collections::HashMap; +use std::hash::{Hash, Hasher}; + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct KeyByLabelValues { + // pub labels: HashMap, + pub labels: Vec, +} + +impl KeyByLabelValues { + pub fn new() -> Self { + Self { labels: Vec::new() } + } + + pub fn new_with_labels(labels: Vec) -> Self { + Self { labels } + } + + pub fn insert(&mut self, value: String) { + self.labels.push(value); + } + + pub fn get(&self, index: usize) -> Option<&String> { + self.labels.get(index) + } + + /// Encode labels as a semicolon-joined string — the canonical key format used + /// for sketch item hashing (CountMinSketch, CountSketch, HydraKLL). + pub fn to_semicolon_str(&self) -> String { + self.labels.join(";") + } + + #[cfg(test)] + /// Decode a semicolon-joined string back into a KeyByLabelValues. + pub fn from_semicolon_str(s: &str) -> Self { + Self { + labels: s.split(';').map(|s| s.to_string()).collect(), + } + } + + pub fn is_empty(&self) -> bool { + self.labels.is_empty() + } + + pub fn len(&self) -> usize { + self.labels.len() + } +} + +impl Hash for KeyByLabelValues { + fn hash(&self, state: &mut H) { + // Create a sorted vector of key-value pairs for consistent hashing + let mut sorted_pairs: Vec<_> = self.labels.iter().collect(); + sorted_pairs.sort(); + + for value in sorted_pairs { + value.hash(state); + } + } +} + +impl Default for KeyByLabelValues { + fn default() -> Self { + Self::new() + } +} + +impl std::fmt::Display for KeyByLabelValues { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{{")?; + let mut first = true; + for value in &self.labels { + if !first { + write!(f, ", ")?; + } + write!(f, "{value}")?; + first = false; + } + write!(f, "}}") + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_key_by_label_values() { + let mut key = KeyByLabelValues::new(); + key.insert("localhost:8080".to_string()); + key.insert("prometheus".to_string()); + + assert_eq!(key.len(), 2); + assert_eq!(key.get(0), Some(&"localhost:8080".to_string())); + assert_eq!(key.get(1), Some(&"prometheus".to_string())); + } + + #[test] + fn test_semicolon_roundtrip() { + let key = KeyByLabelValues::new_with_labels(vec!["web".to_string(), "prod".to_string()]); + assert_eq!(key.to_semicolon_str(), "web;prod"); + let roundtripped = KeyByLabelValues::from_semicolon_str("web;prod"); + assert_eq!(roundtripped, key); + } + + #[test] + fn test_hash_consistency() { + let mut key1 = KeyByLabelValues::new(); + key1.insert("a".to_string()); + key1.insert("b".to_string()); + + let mut key2 = KeyByLabelValues::new(); + key2.insert("b".to_string()); + key2.insert("a".to_string()); + + // Should hash to the same value regardless of insertion order + let mut hasher1 = std::collections::hash_map::DefaultHasher::new(); + let mut hasher2 = std::collections::hash_map::DefaultHasher::new(); + + key1.hash(&mut hasher1); + key2.hash(&mut hasher2); + + assert_eq!(hasher1.finish(), hasher2.finish()); + } +} diff --git a/crates/asap-physical-operators/src/lib.rs b/crates/asap-physical-operators/src/lib.rs new file mode 100644 index 00000000..5026cb6d --- /dev/null +++ b/crates/asap-physical-operators/src/lib.rs @@ -0,0 +1,22 @@ +//! Shared physical operators. This layer holds summary kernels and typed values. + +pub mod key_by_label_values; +pub mod measurement; +pub mod summary_kernels; +pub use summary_kernels::traits; + +mod statistic; +pub use key_by_label_values::KeyByLabelValues; +pub use measurement::Measurement; +pub use statistic::Statistic; +pub use traits::*; + +pub mod capability; +pub use summary_kernels::factory; + +/// The exact Planner contract used by these kernels. +pub use planner_types as planner; + +mod error; +pub use error::Error; +pub mod values; diff --git a/crates/asap-physical-operators/src/measurement.rs b/crates/asap-physical-operators/src/measurement.rs new file mode 100644 index 00000000..57234f01 --- /dev/null +++ b/crates/asap-physical-operators/src/measurement.rs @@ -0,0 +1,48 @@ +use serde::{Deserialize, Serialize}; +use std::ops::Add; + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct Measurement { + pub value: f64, +} + +impl Measurement { + pub fn new(value: f64) -> Self { + Self { value } + } +} + +impl Add for Measurement { + type Output = Measurement; + + fn add(self, other: Measurement) -> Measurement { + Measurement::new(self.value + other.value) + } +} + +impl Add for &Measurement { + type Output = Measurement; + + fn add(self, other: &Measurement) -> Measurement { + Measurement::new(self.value + other.value) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_measurement_creation() { + let measurement = Measurement::new(42.5); + assert_eq!(measurement.value, 42.5); + } + + #[test] + fn test_measurement_addition() { + let m1 = Measurement::new(10.0); + let m2 = Measurement::new(20.0); + let result = m1 + m2; + assert_eq!(result.value, 30.0); + } +} diff --git a/crates/asap-physical-operators/src/statistic.rs b/crates/asap-physical-operators/src/statistic.rs new file mode 100644 index 00000000..7053308d --- /dev/null +++ b/crates/asap-physical-operators/src/statistic.rs @@ -0,0 +1,67 @@ +use std::{fmt, str::FromStr}; +use tracing::debug; +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +pub enum Statistic { + Count, + Sum, + Cardinality, + FrequencyL2, + FrequencyEntropy, + Increase, + Rate, + Min, + Max, + Quantile, + Topk, +} + +impl fmt::Display for Statistic { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + debug!("Formatting Statistic: {:?}", self); + match self { + Statistic::Count => write!(f, "count"), + Statistic::Sum => write!(f, "sum"), + Statistic::Cardinality => write!(f, "cardinality"), + Statistic::FrequencyL2 => write!(f, "frequency_l2"), + Statistic::FrequencyEntropy => write!(f, "frequency_entropy"), + Statistic::Increase => write!(f, "increase"), + Statistic::Rate => write!(f, "rate"), + Statistic::Min => write!(f, "min"), + Statistic::Max => write!(f, "max"), + Statistic::Quantile => write!(f, "quantile"), + Statistic::Topk => write!(f, "topk"), + } + } +} + +#[allow(clippy::should_implement_trait)] +impl Statistic { + pub fn from_str(s: &str) -> Option { + debug!("Parsing Statistic from string: {}", s); + match s.to_lowercase().as_str() { + "count" => Some(Statistic::Count), + "sum" => Some(Statistic::Sum), + "cardinality" => Some(Statistic::Cardinality), + "frequency_l2" => Some(Statistic::FrequencyL2), + "frequency_entropy" => Some(Statistic::FrequencyEntropy), + "increase" => Some(Statistic::Increase), + "rate" => Some(Statistic::Rate), + "min" => Some(Statistic::Min), + "max" => Some(Statistic::Max), + "quantile" => Some(Statistic::Quantile), + "topk" => Some(Statistic::Topk), + _ => None, + } + } +} + +impl FromStr for Statistic { + type Err = (); + + /// Parse a statistic from a string (case-insensitive). + /// Use `s.parse::()` or `Statistic::from_str(s)`. + fn from_str(s: &str) -> Result { + debug!("FromStr trait parsing Statistic: {}", s); + Statistic::from_str(s).ok_or(()) + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/count_min_sketch.rs b/crates/asap-physical-operators/src/summary_kernels/count_min_sketch.rs new file mode 100644 index 00000000..5a78fc48 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/count_min_sketch.rs @@ -0,0 +1,76 @@ +//! Count-Min Sketch frequency summary over `asap_sketchlib::CountMinSketch`. +use crate::{AggregateCore, KernelError, KeyByLabelValues}; +use asap_sketchlib::CountMinSketch; + +#[derive(Debug, Clone)] +pub struct CountMinSketchAccumulator { + pub inner: CountMinSketch, +} + +impl CountMinSketchAccumulator { + pub fn new(row_num: usize, col_num: usize) -> Self { + Self { + inner: CountMinSketch::new(row_num, col_num), + } + } + + /// Estimated frequency of one item. + pub fn query_key(&self, key: &KeyByLabelValues) -> f64 { + self.inner.estimate(&key.to_semicolon_str()) + } +} + +impl AggregateCore for CountMinSketchAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with(&self, other: &dyn AggregateCore) -> Result, KernelError> { + let other = other + .as_any() + .downcast_ref::() + .ok_or("Count-Min Sketch merges only with Count-Min Sketch")?; + Ok(Box::new(Self { + inner: CountMinSketch::merge_refs(&[&self.inner, &other.inner])?, + })) + } + + fn approx_memory_bytes(&self) -> usize { + 16 * 1024 + } +} + +#[cfg(test)] +mod tests { + use super::*; + + // Merged point counts add item frequencies and never underestimate. + #[test] + fn merged_point_counts_add() { + let (mut a, mut b) = ( + CountMinSketchAccumulator::new(3, 128), + CountMinSketchAccumulator::new(3, 128), + ); + let key = KeyByLabelValues::new_with_labels(vec!["checkout".into()]); + a.inner.update(&key.to_semicolon_str(), 2.0); + b.inner.update(&key.to_semicolon_str(), 3.0); + let merged = a.merge_with(&b).unwrap(); + let merged = merged + .as_any() + .downcast_ref::() + .unwrap(); + assert!(merged.query_key(&key) >= 5.0); + } + + // Merge rejects a different summary family. + #[test] + fn rejects_foreign_merge() { + let cms = CountMinSketchAccumulator::new(3, 128); + let kll = crate::summary_kernels::DatasketchesKLLAccumulator::new(200); + assert!(cms.merge_with(&kll).is_err()); + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/count_min_sketch_with_heap.rs b/crates/asap-physical-operators/src/summary_kernels/count_min_sketch_with_heap.rs new file mode 100644 index 00000000..da07aeee --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/count_min_sketch_with_heap.rs @@ -0,0 +1,234 @@ +use crate::{AggregateCore, KeyByLabelValues}; +use asap_sketchlib::CountMinSketchWithHeap; + +/// Count-Min Sketch with a top-k heap over `asap_sketchlib::CountMinSketchWithHeap`. +#[derive(Debug, Clone)] +pub struct CountMinSketchWithHeapAccumulator { + pub inner: CountMinSketchWithHeap, +} + +impl CountMinSketchWithHeapAccumulator { + pub fn new(row_num: usize, col_num: usize, heap_size: usize) -> Self { + Self { + inner: CountMinSketchWithHeap::new(row_num, col_num, heap_size), + } + } + + pub fn query_key(&self, key: &KeyByLabelValues) -> f64 { + let key_string = key.labels.join(";"); + self.inner.estimate(&key_string) + } + + /// VALUE-WEIGHTED heavy-hitter update (FIX: CountSketch/CMS topk + /// recall-0). The default ingest path inserts `+1` per occurrence keyed + /// by the raw `item`, so the heap ranks groups by OCCURRENCE COUNT — the + /// wrong answer for `topk(k, sum by (label) (metric))`, which asks for + /// the top groups by SUM OF VALUE. This update adds the sample `value` + /// (not `+1`) into both the CMS matrix and the top-k heap, keyed by the + /// GROUP LABEL (e.g. the `host` / `zone` value), so the heap's ranking is + /// by summed value. Repeated calls for the same `group_label` accumulate, + /// so after folding a window the heap holds Σvalue per group. + /// + /// Delegates to the library's value-weighted `CountMinSketchWithHeap:: + /// update(key, value)` (`sketchlib_cms_heap_update` → `insert_many(key, + /// round(value))`), which is the "separate update path" the evaluation + /// plan (Fig 3c) called for. + pub fn insert_value(&mut self, group_label: &str, value: f64) { + self.inner.update(group_label, value); + } + + /// Read the top-`k` GROUPS ranked by summed VALUE (descending), keyed by + /// the group label. Pairs with [`Self::insert_value`]: the heap built by + /// value-weighted updates ranks by Σvalue, so this returns the + /// value-weighted top-k (not the occurrence-count top-k the raw `item` + /// heap would give). Sorted descending by value; ties broken by key for + /// determinism; truncated to `k`. + pub fn topk_by_value(&self, k: usize) -> Vec<(String, f64)> { + let mut items: Vec<(String, f64)> = self + .inner + .topk_heap_items() + .into_iter() + .map(|it| (it.key, it.value)) + .collect(); + items.sort_by(|a, b| { + b.1.partial_cmp(&a.1) + .unwrap_or(std::cmp::Ordering::Equal) + .then_with(|| a.0.cmp(&b.0)) + }); + items.truncate(k); + items + } + + /// Get all keys from the top-k heap. + pub fn get_topk_keys(&self) -> Vec { + self.inner + .topk_heap_items() + .iter() + .map(|item| { + let labels: Vec = item.key.split(';').map(|s| s.to_string()).collect(); + KeyByLabelValues { labels } + }) + .collect() + } +} + +impl AggregateCore for CountMinSketchWithHeapAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with( + &self, + other: &dyn AggregateCore, + ) -> Result, Box> { + let other_cms = other + .as_any() + .downcast_ref::() + .ok_or("Failed to downcast to CountMinSketchWithHeapAccumulator")?; + + let mut merged = self.clone(); + merged.inner.merge(&other_cms.inner)?; + Ok(Box::new(merged)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_count_min_sketch_with_heap_creation() { + let cms = CountMinSketchWithHeapAccumulator::new(4, 1000, 20); + assert_eq!(cms.inner.rows(), 4); + assert_eq!(cms.inner.cols(), 1000); + assert_eq!(cms.inner.heap_size, 20); + assert_eq!(cms.inner.topk_heap_items().len(), 0); + } + + #[test] + fn test_get_topk_keys() { + let mut cms = CountMinSketchWithHeapAccumulator::new(2, 3, 5); + cms.inner.update("label1;label2", 100.0); + cms.inner.update("label3;label4", 50.0); + + let keys = cms.get_topk_keys(); + assert_eq!(keys.len(), 2); + // Heap order is not part of the contract; compare as a set. + let label_sets: std::collections::HashSet<_> = + keys.iter().map(|k| k.labels.clone()).collect(); + assert!(label_sets.contains(&vec!["label1".to_string(), "label2".to_string()])); + assert!(label_sets.contains(&vec!["label3".to_string(), "label4".to_string()])); + } + + // ---------------------------------------------------------------- + // FIX 1 — VALUE-WEIGHTED top-k (recall 0 → correct). + // + // `topk(k, sum by (host) (cpu_load))` asks for the top-k hosts by + // SUM OF VALUE. The heavy-hitter heap built by the default `+1`-per- + // occurrence update ranks by COUNT keyed by `item`, so its recall + // against the value-weighted ground truth is 0 when the busiest host + // (most samples) is NOT the heaviest host (largest Σvalue). + // `insert_value(group_label, value)` adds the sample VALUE keyed by the + // GROUP LABEL, so `topk_by_value` ranks by Σvalue — correct recall. + // ---------------------------------------------------------------- + + /// Crafted adversarial dataset: the host with the MOST samples + /// (`h_chatty`, 100 tiny samples) is NOT the host with the largest + /// value-sum (`h_heavy`, a handful of huge samples). A COUNT-ranked + /// heap would surface `h_chatty`; the value-weighted top-k must surface + /// the true heavy hitters by Σvalue, giving recall 1.0 against the + /// ground-truth top-k-by-value-sum. + #[test] + fn value_weighted_topk_has_full_recall_vs_count_topk() { + // (host, per-sample value, sample count) → true Σvalue: + // h_heavy : 1000 × 3 = 3000 (few samples, huge value) + // h_mid : 200 × 5 = 1000 + // h_small : 50 × 6 = 300 + // h_chatty: 1 × 100 = 100 (MOST samples, tiny value) + let data: &[(&str, f64, usize)] = &[ + ("h_heavy", 1000.0, 3), + ("h_mid", 200.0, 5), + ("h_small", 50.0, 6), + ("h_chatty", 1.0, 100), + ]; + + // Wide CMS + heap large enough to hold every group exactly (4 groups) + // so the estimate equals the true Σvalue with no hash collisions. + let mut acc = CountMinSketchWithHeapAccumulator::new(5, 4096, 16); + let mut truth: std::collections::HashMap<&str, f64> = std::collections::HashMap::new(); + for (host, value, count) in data { + for _ in 0..*count { + acc.insert_value(host, *value); + } + *truth.entry(*host).or_insert(0.0) += value * (*count as f64); + } + + // Ground-truth top-2 by value-sum: h_heavy (3000), h_mid (1000). + let mut truth_ranked: Vec<(&str, f64)> = truth.into_iter().collect(); + truth_ranked.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap()); + let truth_top2: std::collections::HashSet<&str> = + truth_ranked.iter().take(2).map(|(k, _)| *k).collect(); + assert!( + truth_top2.contains("h_heavy") && truth_top2.contains("h_mid"), + "ground-truth top-2 by value-sum should be h_heavy + h_mid" + ); + + // Value-weighted top-2 from the heap. + let got = acc.topk_by_value(2); + assert_eq!(got.len(), 2, "k=2 → two groups: {got:?}"); + let got_keys: std::collections::HashSet<&str> = + got.iter().map(|(k, _)| k.as_str()).collect(); + + // RECALL = |got ∩ truth| / |truth| must be 1.0. + let hits = got_keys.intersection(&truth_top2).count(); + let recall = hits as f64 / truth_top2.len() as f64; + assert_eq!( + recall, 1.0, + "value-weighted top-k recall must be 1.0 (count-ranked heap would \ + surface h_chatty and miss h_heavy → recall < 1): got={got:?}" + ); + + // The busiest-by-count host (h_chatty) must NOT be in the top-2, + // proving we rank by value-sum, not occurrence count. + assert!( + !got_keys.contains("h_chatty"), + "h_chatty (most samples, smallest value-sum) must be excluded: {got:?}" + ); + + // Estimates are exact here (no collisions, heap holds all groups): + // top-1 must be h_heavy with Σvalue 3000. + assert_eq!(got[0].0, "h_heavy"); + assert!( + (got[0].1 - 3000.0).abs() < 1e-6, + "h_heavy value-sum estimate ≈ 3000, got {}", + got[0].1 + ); + assert_eq!(got[1].0, "h_mid"); + assert!( + (got[1].1 - 1000.0).abs() < 1e-6, + "h_mid value-sum estimate ≈ 1000, got {}", + got[1].1 + ); + } + + /// A single value-weighted insert must put the full value (not +1) into + /// the heap, and repeated inserts for the same group must accumulate. + #[test] + fn insert_value_accumulates_summed_value_in_heap() { + let mut acc = CountMinSketchWithHeapAccumulator::new(4, 1024, 8); + acc.insert_value("g", 10.0); + acc.insert_value("g", 25.0); + let top = acc.topk_by_value(1); + assert_eq!(top.len(), 1); + assert_eq!(top[0].0, "g"); + assert!( + (top[0].1 - 35.0).abs() < 1e-6, + "summed value should be 35 (10+25), got {}", + top[0].1 + ); + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/count_sketch.rs b/crates/asap-physical-operators/src/summary_kernels/count_sketch.rs new file mode 100644 index 00000000..2728b5f2 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/count_sketch.rs @@ -0,0 +1,96 @@ +//! CountSketch accumulator backed by `asap_sketchlib::CountSketch`. +//! +//! Per-key queries delegate to sketchlib's median-of-signed-rows estimator. +//! Top-k requires the separate heap-bearing accumulator. + +use crate::{AggregateCore, KeyByLabelValues}; +use asap_sketchlib::CountSketch; + +/// Count Sketch accumulator — inner matrix of signed counts. +#[derive(Debug, Clone)] +pub struct CountSketchAccumulator { + pub inner: CountSketch, +} + +impl CountSketchAccumulator { + pub fn new(row_num: usize, col_num: usize) -> Self { + Self { + inner: CountSketch::new(row_num, col_num), + } + } + + /// Median-of-signed-rows point estimate for `key`, via + /// `asap_sketchlib::CountSketch::estimate`. + pub fn query_key(&self, key: &KeyByLabelValues) -> f64 { + self.inner.estimate(&key.to_semicolon_str()) + } +} + +impl AggregateCore for CountSketchAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with( + &self, + other: &dyn AggregateCore, + ) -> Result, Box> { + let other_cs = other + .as_any() + .downcast_ref::() + .ok_or("Failed to downcast to CountSketchAccumulator")?; + + let merged_inner = CountSketch::merge_refs(&[&self.inner, &other_cs.inner])?; + Ok(Box::new(Self { + inner: merged_inner, + })) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_query_key_uses_real_sketchlib_estimator() { + // `query_key` must match sketchlib's estimator and hash specification. + let mut cs = CountSketchAccumulator::new(4, 1000); + let key = KeyByLabelValues::new_with_labels(vec!["web".to_string()]); + cs.inner.update(&key.to_semicolon_str(), 10.0); + assert_eq!( + cs.query_key(&key), + cs.inner.estimate(&key.to_semicolon_str()) + ); + } + + #[test] + fn test_aggregate_core_merge_matches_matrix_add() { + let a = CountSketchAccumulator { + inner: CountSketch::from_legacy_matrix(vec![vec![1.0, -2.0], vec![3.0, -4.0]], 2, 2), + }; + let b = CountSketchAccumulator { + inner: CountSketch::from_legacy_matrix(vec![vec![-1.0, 2.0], vec![-3.0, 4.0]], 2, 2), + }; + let merged_box = a.merge_with(&b).expect("merge ok"); + let merged = merged_box + .as_any() + .downcast_ref::() + .expect("downcast ok"); + let m = merged.inner.sketch(); + assert_eq!(m[0], vec![0.0, 0.0]); + assert_eq!(m[1], vec![0.0, 0.0]); + } + + #[test] + fn test_aggregate_core_merge_wrong_type_rejects() { + use crate::summary_kernels::count_min_sketch::CountMinSketchAccumulator; + let cs = CountSketchAccumulator::new(2, 3); + let cms = CountMinSketchAccumulator::new(2, 3); + let result = cs.merge_with(&cms); + assert!(result.is_err()); + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/count_sketch_with_heap.rs b/crates/asap-physical-operators/src/summary_kernels/count_sketch_with_heap.rs new file mode 100644 index 00000000..a735e4bd --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/count_sketch_with_heap.rs @@ -0,0 +1,143 @@ +//! CountSketch with a top-k heap over `asap_sketchlib::CountSketchWithHeap` +//! (median-of-signed-rows), distinct from the Count-Min heap variant. + +use crate::{AggregateCore, KeyByLabelValues}; +use asap_sketchlib::CountSketchWithHeap; + +#[derive(Debug, Clone)] +pub struct CountSketchWithHeapAccumulator { + pub inner: CountSketchWithHeap, +} + +impl CountSketchWithHeapAccumulator { + pub fn new(row_num: usize, col_num: usize, heap_size: usize) -> Self { + Self { + inner: CountSketchWithHeap::new(row_num, col_num, heap_size), + } + } + + pub fn query_key(&self, key: &KeyByLabelValues) -> f64 { + let key_string = key.labels.join(";"); + self.inner.estimate(&key_string) + } + + /// Value-weighted heavy-hitter update -- see + /// `CountMinSketchWithHeapAccumulator::insert_value`'s doc for why + /// this (not a `+1`-per-occurrence update) is the correct semantics + /// for `topk(k, sum by (label) (metric))`-shaped queries. + pub fn insert_value(&mut self, group_label: &str, value: f64) { + self.inner.update(group_label, value); + } + + /// Read the top-`k` groups ranked by summed value (descending, tie-broken + /// by key for determinism). Mirrors `CountMinSketchWithHeapAccumulator::topk_by_value`. + pub fn topk_by_value(&self, k: usize) -> Vec<(String, f64)> { + let mut items: Vec<(String, f64)> = self + .inner + .topk_heap_items() + .into_iter() + .map(|it| (it.key, it.value)) + .collect(); + items.sort_by(|a, b| { + b.1.partial_cmp(&a.1) + .unwrap_or(std::cmp::Ordering::Equal) + .then_with(|| a.0.cmp(&b.0)) + }); + items.truncate(k); + items + } + + /// Get all keys from the top-k heap. + pub fn get_topk_keys(&self) -> Vec { + self.inner + .topk_heap_items() + .iter() + .map(|item| { + let labels: Vec = item.key.split(';').map(|s| s.to_string()).collect(); + KeyByLabelValues { labels } + }) + .collect() + } +} + +impl AggregateCore for CountSketchWithHeapAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with( + &self, + other: &dyn AggregateCore, + ) -> Result, Box> { + let other_cs = other + .as_any() + .downcast_ref::() + .ok_or("Failed to downcast to CountSketchWithHeapAccumulator")?; + + let mut merged = self.clone(); + merged.inner.merge(&other_cs.inner)?; + Ok(Box::new(merged)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_count_sketch_with_heap_creation() { + let cs = CountSketchWithHeapAccumulator::new(4, 1000, 20); + assert_eq!(cs.inner.rows(), 4); + assert_eq!(cs.inner.cols(), 1000); + assert_eq!(cs.inner.heap_size, 20); + assert_eq!(cs.inner.topk_heap_items().len(), 0); + } + + #[test] + fn test_get_topk_keys() { + let mut cs = CountSketchWithHeapAccumulator::new(2, 3, 5); + cs.inner.update("label1;label2", 100.0); + cs.inner.update("label3;label4", 50.0); + + let keys = cs.get_topk_keys(); + assert_eq!(keys.len(), 2); + let label_sets: std::collections::HashSet<_> = + keys.iter().map(|k| k.labels.clone()).collect(); + assert!(label_sets.contains(&vec!["label1".to_string(), "label2".to_string()])); + assert!(label_sets.contains(&vec!["label3".to_string(), "label4".to_string()])); + } + + #[test] + fn insert_value_accumulates_summed_value_in_heap() { + let mut acc = CountSketchWithHeapAccumulator::new(4, 1024, 8); + acc.insert_value("g", 10.0); + acc.insert_value("g", 25.0); + let top = acc.topk_by_value(1); + assert_eq!(top.len(), 1); + assert_eq!(top[0].0, "g"); + assert!( + (top[0].1 - 35.0).abs() < 1e-6, + "summed value should be 35 (10+25), got {}", + top[0].1 + ); + } + + /// CountSketch and Count-Min heap states are distinct families and never merge. + #[test] + fn test_rejects_merge_with_cms_family_accumulator() { + use crate::summary_kernels::count_min_sketch_with_heap::CountMinSketchWithHeapAccumulator; + + let cs = CountSketchWithHeapAccumulator::new(4, 64, 10); + let cms = CountMinSketchWithHeapAccumulator::new(4, 64, 10); + let result = cs.merge_with(&cms); + assert!( + result.is_err(), + "CountSketchWithHeapAccumulator must not merge with CountMinSketchWithHeapAccumulator \ + -- different algorithms sharing only a storage shape" + ); + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/datasketches_kll.rs b/crates/asap-physical-operators/src/summary_kernels/datasketches_kll.rs new file mode 100644 index 00000000..0f740a9d --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/datasketches_kll.rs @@ -0,0 +1,106 @@ +//! KLL quantile summary over `asap_sketchlib::KllSketch`. +use crate::{AggregateCore, KernelError}; +use asap_sketchlib::KllSketch; +use planner_types::post_asap::SketchQuery; + +#[derive(Clone)] +pub struct DatasketchesKLLAccumulator { + pub inner: KllSketch, +} + +impl DatasketchesKLLAccumulator { + pub fn new(k: u16) -> Self { + Self { + inner: KllSketch::new(k), + } + } + + pub fn update(&mut self, value: f64) { + self.inner.update(value); + } + + pub fn get_quantile(&self, quantile: f64) -> f64 { + self.inner.quantile(quantile) + } +} + +impl std::fmt::Debug for DatasketchesKLLAccumulator { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("DatasketchesKLLAccumulator") + .field("k", &self.inner.k) + .field("sketch_n", &self.inner.count()) + .finish() + } +} + +// SAFETY: `KllSketch` owns its buffers and has no interior mutability; the +// accumulator is only mutated through `&mut self`. +unsafe impl Send for DatasketchesKLLAccumulator {} +unsafe impl Sync for DatasketchesKLLAccumulator {} + +impl AggregateCore for DatasketchesKLLAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with(&self, other: &dyn AggregateCore) -> Result, KernelError> { + let other = other + .as_any() + .downcast_ref::() + .ok_or("KLL merges only with KLL")?; + Ok(Box::new(Self { + inner: KllSketch::merge_refs(&[&self.inner, &other.inner])?, + })) + } + + fn estimate(&self, query: &SketchQuery) -> Result { + match query { + SketchQuery::Quantile { q } if (0.0..=1.0).contains(q) => Ok(self.get_quantile(*q)), + SketchQuery::Quantile { .. } => Err("quantile must be in [0, 1]".into()), + other => Err(format!("KLL does not answer {other:?}").into()), + } + } + + fn approx_memory_bytes(&self) -> usize { + // KLL with default k=200 holds ~2*k items (~3 KiB); round up for overhead. + 4 * 1024 + } +} + +#[cfg(test)] +mod tests { + use super::*; + + // Merging two KLL states reads like one state built over both inputs. + #[test] + fn merged_quantile_matches_single_build() { + let (mut a, mut b, mut all) = ( + DatasketchesKLLAccumulator::new(200), + DatasketchesKLLAccumulator::new(200), + DatasketchesKLLAccumulator::new(200), + ); + for v in 0..100 { + a.update(f64::from(v)); + all.update(f64::from(v)); + } + for v in 100..200 { + b.update(f64::from(v)); + all.update(f64::from(v)); + } + let merged = a.merge_with(&b).unwrap(); + let q = SketchQuery::Quantile { q: 0.5 }; + assert_eq!(merged.estimate(&q).unwrap(), all.estimate(&q).unwrap()); + } + + // KLL answers only quantiles in [0, 1]. + #[test] + fn rejects_unsupported_or_out_of_range_queries() { + let kll = DatasketchesKLLAccumulator::new(200); + assert!(kll.estimate(&SketchQuery::Quantile { q: 1.5 }).is_err()); + assert!(kll.estimate(&SketchQuery::Cardinality).is_err()); + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/dd_sketch.rs b/crates/asap-physical-operators/src/summary_kernels/dd_sketch.rs new file mode 100644 index 00000000..1d6be98c --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/dd_sketch.rs @@ -0,0 +1,88 @@ +//! DDSketch quantile summary over `asap_sketchlib::DdSketch`. +use crate::{AggregateCore, KernelError}; +use asap_sketchlib::DdSketch; +use planner_types::post_asap::SketchQuery; + +#[derive(Debug, Clone)] +pub struct DDSketchAccumulator { + pub inner: DdSketch, +} + +impl DDSketchAccumulator { + pub fn new(alpha: f64) -> Self { + Self { + inner: DdSketch::new(alpha), + } + } +} + +impl AggregateCore for DDSketchAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with(&self, other: &dyn AggregateCore) -> Result, KernelError> { + let other = other + .as_any() + .downcast_ref::() + .ok_or("DDSketch merges only with DDSketch")?; + Ok(Box::new(Self { + inner: DdSketch::merge_refs(&[&self.inner, &other.inner])?, + })) + } + + /// Quantiles, and the total sample count as a bare `PointCount`. + fn estimate(&self, query: &SketchQuery) -> Result { + match query { + SketchQuery::Quantile { q } if (0.0..=1.0).contains(q) => self + .inner + .quantile(*q) + .ok_or_else(|| "DDSketch quantile of an empty population".into()), + SketchQuery::Quantile { .. } => Err("quantile must be in [0, 1]".into()), + SketchQuery::PointCount { value: None, .. } => Ok(self.inner.total_count() as f64), + other => Err(format!("DDSketch does not answer {other:?}").into()), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use planner_types::pre_asap::ColumnRef; + + fn bare_count() -> SketchQuery { + SketchQuery::PointCount { + key: ColumnRef::SampleValue, + value: None, + } + } + + // A bare point count reads the total sample count, and merge adds counts. + #[test] + fn count_and_quantile_survive_merge() { + let (mut a, mut b) = ( + DDSketchAccumulator::new(0.01), + DDSketchAccumulator::new(0.01), + ); + for v in 1..=50 { + a.inner.update(f64::from(v)); + b.inner.update(f64::from(v + 50)); + } + let merged = a.merge_with(&b).unwrap(); + assert_eq!(merged.estimate(&bare_count()).unwrap(), 100.0); + let median = merged.estimate(&SketchQuery::Quantile { q: 0.5 }).unwrap(); + assert!((median - 50.0).abs() <= 1.0, "{median}"); + } + + // An empty DDSketch has no quantile, and unsupported queries are errors. + #[test] + fn empty_quantile_and_unsupported_queries_fail() { + let dd = DDSketchAccumulator::new(0.01); + assert!(dd.estimate(&SketchQuery::Quantile { q: 0.5 }).is_err()); + assert!(dd.estimate(&SketchQuery::Cardinality).is_err()); + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/exact.rs b/crates/asap-physical-operators/src/summary_kernels/exact.rs new file mode 100644 index 00000000..3343e875 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/exact.rs @@ -0,0 +1,215 @@ +//! Exact summary state identified by Planner family, independent of keyed layout. +use super::increase::IncreaseAccumulator; +use crate::Statistic; +use crate::{AggregateCore, KeyByLabelValues, Measurement}; +use planner_types::post_asap::{ExactKind, ExactParams, SummaryFamilyType}; +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; + +type Error = Box; + +#[derive(Debug, Clone, Serialize, Deserialize)] +enum ScalarState { + Sum(f64), + Count(u64), + Min(Option), + Max(Option), + Counter(Option), +} + +/// Both the family and population layout survive persistence. Sharing counter +/// arithmetic never authorizes a Rate state to answer an Increase readout. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ExactAccumulator { + family: SummaryFamilyType, + scalar: ScalarState, + keyed: Option>, +} + +/// Planned readout of an exact summary. `lookback_ms` is the logical PromQL +/// counter window; the evaluation range is resolved from it at run time. +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +pub struct ExactReadout { + pub statistic: Statistic, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub lookback_ms: Option, +} + +impl ExactAccumulator { + /// Read one population. An empty MIN/MAX population reads as `None`. + /// `range_ms` extrapolates a counter Rate/Increase to that evaluation range. + pub fn readout( + &self, + statistic: Statistic, + range_ms: Option<(i64, i64)>, + key: Option<&KeyByLabelValues>, + ) -> Result, Error> { + if statistic != self.statistic() { + return Err("readout differs from Planner exact family".into()); + } + let state = match (&self.keyed, key) { + (Some(states), Some(key)) => states.get(key).ok_or("unknown exact population")?, + (None, None) => &self.scalar, + _ => return Err("readout population differs from installed layout".into()), + }; + match state { + ScalarState::Sum(sum) => Ok(Some(*sum)), + ScalarState::Count(count) => Ok(Some(*count as f64)), + ScalarState::Min(value) | ScalarState::Max(value) => Ok(*value), + ScalarState::Counter(Some(counter)) => counter + .extrapolated_value(range_ms, statistic == Statistic::Rate) + .map(Some), + ScalarState::Counter(None) => Err("empty counter population".into()), + } + } + + /// Exact integer count of an unkeyed Count state. + pub fn count(&self) -> Option { + match (&self.keyed, &self.scalar) { + (None, ScalarState::Count(count)) => Some(*count), + _ => None, + } + } + + /// Accumulate into run-local scratch state. Persistent input states remain + /// immutable; a failed merge discards this scratch state. + pub(crate) fn merge_from(&mut self, other: &Self) -> Result<(), Error> { + if self.family != other.family || self.is_keyed() != other.is_keyed() { + return Err("cannot merge different Planner families or layouts".into()); + } + if let (Some(target), Some(source)) = (&mut self.keyed, &other.keyed) { + for (key, state) in source { + let combined = match target.get(key) { + Some(old) => merge_scalar(old, state)?, + None => state.clone(), + }; + target.insert(key.clone(), combined); + } + } else { + self.scalar = merge_scalar(&self.scalar, &other.scalar)?; + } + Ok(()) + } + + pub fn new(family: SummaryFamilyType, keyed: bool) -> Result { + use ExactKind as K; + use ExactParams as P; + let scalar = match &family { + SummaryFamilyType::ExactAggregate(K::Sum, P::Sum) => ScalarState::Sum(0.0), + SummaryFamilyType::ExactAggregate(K::Count, P::Count) => ScalarState::Count(0), + SummaryFamilyType::ExactAggregate(K::Min, P::Min) => ScalarState::Min(None), + SummaryFamilyType::ExactAggregate(K::Max, P::Max) => ScalarState::Max(None), + SummaryFamilyType::ExactAggregate(K::Rate, P::Rate) + | SummaryFamilyType::ExactAggregate(K::Increase, P::Increase) => { + ScalarState::Counter(None) + } + _ => return Err(format!("unsupported exact Planner family: {family:?}")), + }; + Ok(Self { + family, + scalar, + keyed: keyed.then(HashMap::new), + }) + } + + pub fn family(&self) -> &SummaryFamilyType { + &self.family + } + pub fn is_keyed(&self) -> bool { + self.keyed.is_some() + } + + pub fn update(&mut self, key: Option<&KeyByLabelValues>, value: f64, timestamp: i64) { + let state = match (&mut self.keyed, key) { + (Some(states), Some(key)) => states + .entry(key.clone()) + .or_insert_with(|| self.scalar.clone()), + (None, None) => &mut self.scalar, + _ => panic!("exact update population layout differs from installed DAG"), + }; + match state { + ScalarState::Sum(sum) => *sum += value, + ScalarState::Count(count) => { + *count = count.checked_add(1).expect("exact count overflow") + } + ScalarState::Min(current) => { + *current = Some(current.map_or(value, |old| old.min(value))) + } + ScalarState::Max(current) => { + *current = Some(current.map_or(value, |old| old.max(value))) + } + ScalarState::Counter(current) => match current { + Some(counter) => counter.update(Measurement::new(value), timestamp), + None => { + *current = Some(IncreaseAccumulator::new( + Measurement::new(value), + timestamp, + Measurement::new(value), + timestamp, + )) + } + }, + } + } + + fn statistic(&self) -> Statistic { + match self.family { + SummaryFamilyType::ExactAggregate(ExactKind::Sum, _) => Statistic::Sum, + SummaryFamilyType::ExactAggregate(ExactKind::Count, _) => Statistic::Count, + SummaryFamilyType::ExactAggregate(ExactKind::Min, _) => Statistic::Min, + SummaryFamilyType::ExactAggregate(ExactKind::Max, _) => Statistic::Max, + SummaryFamilyType::ExactAggregate(ExactKind::Rate, _) => Statistic::Rate, + SummaryFamilyType::ExactAggregate(ExactKind::Increase, _) => Statistic::Increase, + _ => unreachable!("validated exact family"), + } + } +} + +fn merge_scalar(left: &ScalarState, right: &ScalarState) -> Result { + Ok(match (left, right) { + (ScalarState::Sum(a), ScalarState::Sum(b)) => ScalarState::Sum(a + b), + (ScalarState::Count(a), ScalarState::Count(b)) => { + ScalarState::Count(a.checked_add(*b).ok_or("exact count overflow")?) + } + (ScalarState::Min(a), ScalarState::Min(b)) => { + ScalarState::Min(a.iter().chain(b).copied().reduce(f64::min)) + } + (ScalarState::Max(a), ScalarState::Max(b)) => { + ScalarState::Max(a.iter().chain(b).copied().reduce(f64::max)) + } + (ScalarState::Counter(a), ScalarState::Counter(b)) => ScalarState::Counter(match (a, b) { + (Some(a), Some(b)) => Some(IncreaseAccumulator::merge_pair(a, b)), + (a, b) => a.clone().or_else(|| b.clone()), + }), + _ => return Err("exact scalar state families differ".into()), + }) +} + +impl AggregateCore for ExactAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + fn as_any(&self) -> &dyn std::any::Any { + self + } + fn merge_with(&self, other: &dyn AggregateCore) -> Result, Error> { + let other = other + .as_any() + .downcast_ref::() + .ok_or("merge requires Planner exact state")?; + let mut merged = self.clone(); + merged.merge_from(other)?; + Ok(Box::new(merged)) + } + fn approx_memory_bytes(&self) -> usize { + std::mem::size_of::() + + self.keyed.as_ref().map_or(0, |m| { + m.keys() + .map(|k| { + std::mem::size_of::() + + k.labels.iter().map(String::len).sum::() + }) + .sum::() + }) + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/factory.rs b/crates/asap-physical-operators/src/summary_kernels/factory.rs new file mode 100644 index 00000000..d44e64ab --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/factory.rs @@ -0,0 +1,799 @@ +use crate::summary_kernels::hll_sketch::HllSketchAccumulator; +use crate::summary_kernels::univmon::UnivMonAccumulator; +use crate::summary_kernels::{ + CountMinSketchAccumulator, CountMinSketchWithHeapAccumulator, CountSketchAccumulator, + CountSketchWithHeapAccumulator, DDSketchAccumulator, DatasketchesKLLAccumulator, + HydraKllSketchAccumulator, +}; +use crate::{AggregateCore, KeyByLabelValues}; +use planner_types::post_asap::{SketchAlgorithm, SketchParams, SummaryFamilyType}; + +/// Generate the clone-based `AccumulatorUpdater` methods for updaters whose +/// inner `acc` field implements `Clone + AggregateCore`. +macro_rules! impl_clone_accumulator_methods { + ($acc_field:ident) => { + fn take_accumulator(&mut self) -> Box { + let result = Box::new(self.$acc_field.clone()); + self.reset(); + result + } + + fn snapshot_accumulator(&self) -> Box { + Box::new(self.$acc_field.clone()) + } + + fn into_accumulator(self: Box) -> Box { + // Consume the updater and MOVE the accumulator out — no clone. + // Avoids a clone when a pane is evicted at window close. + let this = *self; + Box::new(this.$acc_field) + } + }; +} + +/// Shared update interface for query-time and precompute-time accumulation. +/// +/// This provides a uniform interface over all accumulator types so that the +/// operators don't need to know which concrete type they're dealing with. +pub trait AccumulatorUpdater: Send { + /// Validate an immutable precompute input before an updater can silently + /// discard a value outside its representable domain. + fn validate_single_input(&self, value: f64) -> Result<(), String> { + if value.is_finite() { + Ok(()) + } else { + Err("accumulator input must be finite".into()) + } + } + + /// Feed a single (value, timestamp_ms) pair — for SingleSubpopulation types. + fn update_single(&mut self, value: f64, timestamp_ms: i64); + + /// Feed a keyed (key, value, timestamp_ms) triple, e.g. a frequency item or an exact keyed state. + fn update_keyed(&mut self, key: &KeyByLabelValues, value: f64, timestamp_ms: i64); + + /// Extract the final accumulator as a boxed `AggregateCore`. + fn take_accumulator(&mut self) -> Box; + + /// Non-destructive read of the current accumulator state (clone without reset). + /// Used by pane-based sliding windows to read shared panes. + fn snapshot_accumulator(&self) -> Box; + + /// Consume the updater and return its accumulator by move, avoiding the + /// clone that `take_accumulator`/`snapshot_accumulator` pay. Default falls + /// back to a clone for updaters that can't move their inner state out. + fn into_accumulator(self: Box) -> Box { + self.snapshot_accumulator() + } + + /// Reset internal state for reuse (avoids re-allocation). + fn reset(&mut self); + + /// Whether this updater consumes keyed updates. + fn is_keyed(&self) -> bool; + + /// Estimated memory usage in bytes. + fn memory_usage_bytes(&self) -> usize; +} + +// --------------------------------------------------------------------------- +// KllAccumulatorUpdater +// --------------------------------------------------------------------------- + +pub struct KllAccumulatorUpdater { + acc: DatasketchesKLLAccumulator, + k: u16, +} + +impl KllAccumulatorUpdater { + pub fn new(k: u16) -> Self { + Self { + acc: DatasketchesKLLAccumulator::new(k), + k, + } + } +} + +impl AccumulatorUpdater for KllAccumulatorUpdater { + fn update_single(&mut self, value: f64, _timestamp_ms: i64) { + self.acc.update(value); + } + + fn update_keyed(&mut self, _key: &KeyByLabelValues, value: f64, timestamp_ms: i64) { + self.update_single(value, timestamp_ms); + } + + impl_clone_accumulator_methods!(acc); + + fn reset(&mut self) { + self.acc = DatasketchesKLLAccumulator::new(self.k); + } + + fn is_keyed(&self) -> bool { + false + } + + fn memory_usage_bytes(&self) -> usize { + // KLL sketch size is hard to estimate precisely; use a rough estimate + std::mem::size_of::() + 4096 + } +} + +// --------------------------------------------------------------------------- +// DDSketchAccumulatorUpdater +// --------------------------------------------------------------------------- +pub struct DDSketchAccumulatorUpdater { + acc: DDSketchAccumulator, + alpha: f64, +} + +impl DDSketchAccumulatorUpdater { + pub fn new(alpha: f64) -> Self { + Self { + acc: DDSketchAccumulator::new(alpha), + alpha, + } + } +} + +impl AccumulatorUpdater for DDSketchAccumulatorUpdater { + fn validate_single_input(&self, value: f64) -> Result<(), String> { + let (minimum, maximum) = + asap_sketchlib::sketches::ddsketch::ddsketch_indexable_bounds(self.alpha); + if value.is_finite() && value > 0.0 && value >= minimum && value <= maximum { + Ok(()) + } else { + Err("DDS maintenance input is outside its positive representable domain".into()) + } + } + + fn update_single(&mut self, value: f64, _timestamp_ms: i64) { + self.acc.inner.update(value); + } + + fn update_keyed(&mut self, _key: &KeyByLabelValues, value: f64, timestamp_ms: i64) { + self.update_single(value, timestamp_ms); + } + + impl_clone_accumulator_methods!(acc); + + fn reset(&mut self) { + self.acc = DDSketchAccumulator::new(self.alpha); + } + + fn is_keyed(&self) -> bool { + false + } + + fn memory_usage_bytes(&self) -> usize { + // Bucket store is variable; rough estimate matches KLL. + std::mem::size_of::() + 4096 + } +} + +// --------------------------------------------------------------------------- +// CmsAccumulatorUpdater (CountMinSketch) +// --------------------------------------------------------------------------- + +/// Keyed weighted-frequency updater. +/// +/// A raw Prometheus sample represents the observed metric value, so a bare CMS +/// adds `value` for its key. Counting each received sample as one is a distinct +/// event-count operation and requires an explicit typed plan contract; it must +/// not be inferred from the sketch algorithm alone. +pub struct CmsAccumulatorUpdater { + acc: CountMinSketchAccumulator, + row_num: usize, + col_num: usize, +} + +impl CmsAccumulatorUpdater { + pub fn new(row_num: usize, col_num: usize) -> Self { + Self { + acc: CountMinSketchAccumulator::new(row_num, col_num), + row_num, + col_num, + } + } +} + +impl AccumulatorUpdater for CmsAccumulatorUpdater { + fn update_single(&mut self, _value: f64, _timestamp_ms: i64) { + debug_assert!( + false, + "update_single called on keyed updater; use update_keyed" + ); + } + + fn update_keyed(&mut self, key: &KeyByLabelValues, value: f64, _timestamp_ms: i64) { + self.acc.inner.update(&key.to_semicolon_str(), value); + } + + impl_clone_accumulator_methods!(acc); + + fn reset(&mut self) { + self.acc = CountMinSketchAccumulator::new(self.row_num, self.col_num); + } + + fn is_keyed(&self) -> bool { + true + } + + fn memory_usage_bytes(&self) -> usize { + std::mem::size_of::() + + self.row_num * self.col_num * std::mem::size_of::() + } +} + +// --------------------------------------------------------------------------- +// CmsHeapAccumulatorUpdater — value-weighted / count-weighted top-k +// --------------------------------------------------------------------------- + +/// What quantity the top-k heap ranks keys by. +/// +/// These are DIFFERENT query semantics and must be chosen explicitly: +/// +/// * [`TopkWeight::Value`] — accumulate **Σ of the datapoint value** per key. +/// This answers "top-k by total " (e.g. "top-k hosts by +/// total CPU"). The heap value is the summed metric value, so the read-side +/// reducer's "sort heap descending by value" yields the correct ranking. +/// +/// * [`TopkWeight::Count`] — accumulate **+1 per event** per key (occurrence +/// frequency), the textbook heavy-hitter / frequency-top-k semantics +/// ("which keys appear most often"). +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TopkWeight { + /// Σ datapoint value per key (value-weighted top-k). + Value, + /// +1 per event per key (count-weighted / frequency top-k). + Count, +} + +/// Keyed top-k updater backed by a real `CountMinSketchWithHeap` (a CMS +/// matrix PLUS a size-`heap_size` top-k heap). Unlike the heap-LESS +/// `CmsAccumulatorUpdater`, this enumerates top-k keys at read time +/// (`get_topk_keys` / `topk_heap_items`), which is what `topk(...)` queries +/// need. +/// +/// The key is the frequency item supplied by the operator (e.g. a `host` +/// value), not a group-by population. The accumulated quantity is selected +/// by [`TopkWeight`]: +/// * `Value` → `inner.update(key, value)` adds the datapoint value (Σ value). +/// * `Count` → `inner.update(key, 1.0)` adds one per event (Σ count). +/// +/// Both `CountMinSketchWithHeap` and `CountSketchWithHeap` raw-input policies +/// route here; the heap is the shared distinguishing payload. +pub struct CmsHeapAccumulatorUpdater { + acc: CountMinSketchWithHeapAccumulator, + row_num: usize, + col_num: usize, + heap_size: usize, + weight: TopkWeight, +} + +impl CmsHeapAccumulatorUpdater { + pub fn new(row_num: usize, col_num: usize, heap_size: usize, weight: TopkWeight) -> Self { + Self { + acc: CountMinSketchWithHeapAccumulator::new(row_num, col_num, heap_size), + row_num, + col_num, + heap_size, + weight, + } + } +} + +impl AccumulatorUpdater for CmsHeapAccumulatorUpdater { + fn update_single(&mut self, _value: f64, _timestamp_ms: i64) { + debug_assert!( + false, + "update_single called on keyed updater; use update_keyed" + ); + } + + fn update_keyed(&mut self, key: &KeyByLabelValues, value: f64, _timestamp_ms: i64) { + // Heap key = the group-by label-value vector (e.g. `host`), joined the + // same way the read-side `get_topk_keys` splits it back apart (`;`). + let weighted = match self.weight { + // Σ value: feed the datapoint value. sketchlib's CMS-heap + // `update(key, w)` adds `w.round()` occurrences of `key`, so the + // heap value accumulates the (rounded) summed metric value. + TopkWeight::Value => value, + // Σ count: one occurrence per event, regardless of value. + TopkWeight::Count => 1.0, + }; + self.acc.inner.update(&key.to_semicolon_str(), weighted); + } + + impl_clone_accumulator_methods!(acc); + + fn reset(&mut self) { + self.acc = + CountMinSketchWithHeapAccumulator::new(self.row_num, self.col_num, self.heap_size); + } + + fn is_keyed(&self) -> bool { + true + } + + fn memory_usage_bytes(&self) -> usize { + std::mem::size_of::() + + self.row_num * self.col_num * std::mem::size_of::() + + self.heap_size * (std::mem::size_of::() + 32) + } +} + +// --------------------------------------------------------------------------- +// CountSketchAccumulatorUpdater (real median-of-signed-rows CountSketch) +// --------------------------------------------------------------------------- + +/// Keyed point-frequency updater backed by a real `asap_sketchlib::CountSketch` +/// (signed rows, median-of-rows estimator) — distinct math from +/// `CmsAccumulatorUpdater`'s CMS (min-of-rows). +/// +/// As with bare CMS, each raw Prometheus sample contributes its `value`. +/// Unit event counting must be selected explicitly by a future typed plan +/// contract rather than being implied by `SketchAlgorithm::CountSketch`. +pub struct CountSketchAccumulatorUpdater { + acc: CountSketchAccumulator, + row_num: usize, + col_num: usize, +} + +impl CountSketchAccumulatorUpdater { + pub fn new(row_num: usize, col_num: usize) -> Self { + Self { + acc: CountSketchAccumulator::new(row_num, col_num), + row_num, + col_num, + } + } +} + +impl AccumulatorUpdater for CountSketchAccumulatorUpdater { + fn update_single(&mut self, _value: f64, _timestamp_ms: i64) { + debug_assert!( + false, + "update_single called on keyed updater; use update_keyed" + ); + } + + fn update_keyed(&mut self, key: &KeyByLabelValues, value: f64, _timestamp_ms: i64) { + self.acc.inner.update(&key.to_semicolon_str(), value); + } + + impl_clone_accumulator_methods!(acc); + + fn reset(&mut self) { + self.acc = CountSketchAccumulator::new(self.row_num, self.col_num); + } + + fn is_keyed(&self) -> bool { + true + } + + fn memory_usage_bytes(&self) -> usize { + std::mem::size_of::() + + self.row_num * self.col_num * std::mem::size_of::() + } +} + +// --------------------------------------------------------------------------- +// CountSketchWithHeapAccumulatorUpdater (real CountSketch + top-k heap) +// --------------------------------------------------------------------------- + +/// Keyed top-k updater backed by a real `CountSketchWithHeap` (signed-row +/// CountSketch matrix PLUS a size-`heap_size` top-k heap). Distinct math from +/// `CmsHeapAccumulatorUpdater`'s CMS-with-heap (min-of-rows); shares the same +/// [`TopkWeight`] semantics and heap payload shape. +pub struct CountSketchWithHeapAccumulatorUpdater { + acc: CountSketchWithHeapAccumulator, + row_num: usize, + col_num: usize, + heap_size: usize, + weight: TopkWeight, +} + +impl CountSketchWithHeapAccumulatorUpdater { + pub fn new(row_num: usize, col_num: usize, heap_size: usize, weight: TopkWeight) -> Self { + Self { + acc: CountSketchWithHeapAccumulator::new(row_num, col_num, heap_size), + row_num, + col_num, + heap_size, + weight, + } + } +} + +impl AccumulatorUpdater for CountSketchWithHeapAccumulatorUpdater { + fn update_single(&mut self, _value: f64, _timestamp_ms: i64) { + debug_assert!( + false, + "update_single called on keyed updater; use update_keyed" + ); + } + + fn update_keyed(&mut self, key: &KeyByLabelValues, value: f64, _timestamp_ms: i64) { + let weighted = match self.weight { + TopkWeight::Value => value, + TopkWeight::Count => 1.0, + }; + self.acc.inner.update(&key.to_semicolon_str(), weighted); + } + + impl_clone_accumulator_methods!(acc); + + fn reset(&mut self) { + self.acc = CountSketchWithHeapAccumulator::new(self.row_num, self.col_num, self.heap_size); + } + + fn is_keyed(&self) -> bool { + true + } + + fn memory_usage_bytes(&self) -> usize { + std::mem::size_of::() + + self.row_num * self.col_num * std::mem::size_of::() + + self.heap_size * (std::mem::size_of::() + 32) + } +} + +// --------------------------------------------------------------------------- +// HydraKllAccumulatorUpdater +// --------------------------------------------------------------------------- + +pub struct HydraKllAccumulatorUpdater { + acc: HydraKllSketchAccumulator, + row_num: usize, + col_num: usize, + k: u16, +} + +impl HydraKllAccumulatorUpdater { + pub fn new(row_num: usize, col_num: usize, k: u16) -> Self { + Self { + acc: HydraKllSketchAccumulator::new(row_num, col_num, k), + row_num, + col_num, + k, + } + } +} + +impl AccumulatorUpdater for HydraKllAccumulatorUpdater { + fn update_single(&mut self, _value: f64, _timestamp_ms: i64) { + debug_assert!( + false, + "update_single called on keyed updater; use update_keyed" + ); + } + + fn update_keyed(&mut self, key: &KeyByLabelValues, value: f64, _timestamp_ms: i64) { + self.acc.update(key, value); + } + + impl_clone_accumulator_methods!(acc); + + fn reset(&mut self) { + self.acc = HydraKllSketchAccumulator::new(self.row_num, self.col_num, self.k); + } + + fn is_keyed(&self) -> bool { + true + } + + fn memory_usage_bytes(&self) -> usize { + // Rough estimate: each cell is a KLL sketch + std::mem::size_of::() + self.row_num * self.col_num * 4096 + } +} + +// --------------------------------------------------------------------------- +// Config helpers +// --------------------------------------------------------------------------- + +fn cms_dims(params: &SketchParams) -> (usize, usize) { + match params { + SketchParams::Cms { width, depth } | SketchParams::CountSketch { width, depth } => { + (*depth as usize, *width as usize) + } + other => unreachable!( + "accumulator_spec() paired SketchAlgorithm::Cms/CountSketch with unexpected params: {other:?}" + ), + } +} + +/// Read `(rows = depth, columns = width, heap_size)` out of `SketchParams::CmsWithHeap` +/// or `::CountSketchWithHeap`. +fn cms_heap_dims(params: &SketchParams) -> (usize, usize, usize) { + match params { + SketchParams::CmsWithHeap { + width, + depth, + heap_size, + } + | SketchParams::CountSketchWithHeap { + width, + depth, + heap_size, + } => (*depth as usize, *width as usize, *heap_size as usize), + other => unreachable!( + "accumulator_spec() paired a WithHeap SketchAlgorithm with unexpected params: {other:?}" + ), + } +} + +/// Construct the kernel declared by a Planner SummaryAgg. No deployment config +/// tags participate in this dispatch and unsupported payloads are errors. +pub fn create_planner_accumulator( + family: &SummaryFamilyType, + input: &planner_types::post_asap::SummaryUpdate, + grouping: &planner_types::post_asap::GroupingStrategy, +) -> Result, String> { + if input.item.is_some() + && matches!( + input.weight_domain, + planner_types::post_asap::WeightDomain::NonNegative { + proof: + planner_types::post_asap::NonNegativeWeightProof::ResetAwareCounterDerivative + } + ) + { + return Err("window-weighted summaries require typed DAG binding; integer heap updaters cannot consume rates".into()); + } + + crate::capability::validate_summary_kernel(family, input, grouping)?; + use planner_types::post_asap::GroupingStrategy; + if grouping != &GroupingStrategy::PerSubpopulationInstance { + return Err("shared summary grouping requires a supported Planner Hydra kernel".into()); + } + if matches!(family, SummaryFamilyType::ExactAggregate(..)) { + return Ok(Box::new(PlannerExactUpdater { + acc: crate::summary_kernels::exact::ExactAccumulator::new( + family.clone(), + input.item.is_some(), + )?, + })); + } + let SummaryFamilyType::Sketch(kind, family_grouping) = family else { + return Err(format!("unsupported Planner summary family {family:?}")); + }; + if family_grouping != grouping { + return Err("Planner family and operator grouping disagree".into()); + } + let updater: Box = match (kind.algorithm(), kind.params()) { + (SketchAlgorithm::Kll, SketchParams::Kll { k }) => Box::new(KllAccumulatorUpdater::new( + u16::try_from(*k).map_err(|_| "KLL k exceeds runtime bound")?, + )), + (SketchAlgorithm::DDSketch, SketchParams::DDSketch { alpha }) => { + Box::new(DDSketchAccumulatorUpdater::new(*alpha)) + } + (SketchAlgorithm::Cms, params @ SketchParams::Cms { .. }) => { + let (r, c) = cms_dims(params); + Box::new(CmsAccumulatorUpdater::new(r, c)) + } + (SketchAlgorithm::CountSketch, params @ SketchParams::CountSketch { .. }) => { + let (r, c) = cms_dims(params); + Box::new(CountSketchAccumulatorUpdater::new(r, c)) + } + (SketchAlgorithm::CmsWithHeap, params @ SketchParams::CmsWithHeap { .. }) => { + let (r, c, h) = cms_heap_dims(params); + Box::new(CmsHeapAccumulatorUpdater::new(r, c, h, TopkWeight::Value)) + } + ( + SketchAlgorithm::CountSketchWithHeap, + params @ SketchParams::CountSketchWithHeap { .. }, + ) => { + let (r, c, h) = cms_heap_dims(params); + Box::new(CountSketchWithHeapAccumulatorUpdater::new( + r, + c, + h, + TopkWeight::Value, + )) + } + (SketchAlgorithm::Hll, SketchParams::Hll { precision }) => Box::new(HllUpdater { + acc: HllSketchAccumulator::new( + asap_sketchlib::HllVariant::Regular, + u32::from(*precision), + ), + }), + ( + SketchAlgorithm::UnivMon, + SketchParams::UnivMon { + heap_size, + sketch_rows, + sketch_cols, + layers, + }, + ) => Box::new(UnivMonUpdater { + acc: UnivMonAccumulator::new( + *heap_size as usize, + *sketch_rows as usize, + *sketch_cols as usize, + *layers as usize, + ) + .map_err(|e| e.to_string())?, + }), + _ => { + return Err(format!( + "unsupported Planner algorithm/parameters: {kind:?}" + )) + } + }; + if updater.is_keyed() != input.item.is_some() + && !crate::capability::is_unit_sample_frequency(input) + { + return Err("Planner item expression does not match the selected kernel layout".into()); + } + Ok(updater) +} + +struct PlannerExactUpdater { + acc: crate::summary_kernels::exact::ExactAccumulator, +} +impl AccumulatorUpdater for PlannerExactUpdater { + fn update_single(&mut self, value: f64, timestamp: i64) { + self.acc.update(None, value, timestamp); + } + fn update_keyed(&mut self, key: &KeyByLabelValues, value: f64, timestamp: i64) { + self.acc.update(Some(key), value, timestamp); + } + impl_clone_accumulator_methods!(acc); + fn reset(&mut self) { + self.acc = crate::summary_kernels::exact::ExactAccumulator::new( + self.acc.family().clone(), + self.acc.is_keyed(), + ) + .expect("installed exact family"); + } + fn is_keyed(&self) -> bool { + self.acc.is_keyed() + } + fn memory_usage_bytes(&self) -> usize { + self.acc.approx_memory_bytes() + } +} + +struct UnivMonUpdater { + acc: UnivMonAccumulator, +} + +struct HllUpdater { + acc: HllSketchAccumulator, +} + +impl AccumulatorUpdater for HllUpdater { + fn is_keyed(&self) -> bool { + false + } + fn memory_usage_bytes(&self) -> usize { + self.acc.approx_memory_bytes() + } + fn update_single(&mut self, value: f64, _: i64) { + if !value.is_nan() { + let bits = if value == 0.0 { 0 } else { value.to_bits() }; + self.acc.inner.update(&bits.to_le_bytes()); + } + } + fn update_keyed(&mut self, _: &KeyByLabelValues, value: f64, timestamp_ms: i64) { + self.update_single(value, timestamp_ms); + } + impl_clone_accumulator_methods!(acc); + fn reset(&mut self) { + self.acc.inner = + asap_sketchlib::HllSketch::new(self.acc.inner.variant, self.acc.inner.precision); + } +} + +impl AccumulatorUpdater for UnivMonUpdater { + fn is_keyed(&self) -> bool { + false + } + fn memory_usage_bytes(&self) -> usize { + self.acc.approx_memory_bytes() + } + fn update_single(&mut self, value: f64, _: i64) { + self.acc + .insert_sample(value) + .expect("UnivMon sample counter overflow"); + } + fn update_keyed(&mut self, _: &KeyByLabelValues, value: f64, timestamp_ms: i64) { + self.update_single(value, timestamp_ms); + } + impl_clone_accumulator_methods!(acc); + fn reset(&mut self) { + self.acc.clear(); + } +} + +#[cfg(test)] +mod planner_parameter_regression { + use super::*; + use planner_types::post_asap::{SketchKind, SummaryInputExpr, SummaryUpdate}; + + // Planner width is the bucket count; depth is the independent hash-row count. + #[test] + fn planner_sketch_dimensions_are_not_transposed() { + for (algorithm, params) in [ + ( + SketchAlgorithm::Cms, + SketchParams::Cms { + width: 128, + depth: 3, + }, + ), + ( + SketchAlgorithm::CountSketch, + SketchParams::CountSketch { + width: 128, + depth: 3, + }, + ), + ( + SketchAlgorithm::CmsWithHeap, + SketchParams::CmsWithHeap { + width: 128, + depth: 3, + heap_size: 8, + }, + ), + ( + SketchAlgorithm::CountSketchWithHeap, + SketchParams::CountSketchWithHeap { + width: 128, + depth: 3, + heap_size: 8, + }, + ), + ] { + let family = SummaryFamilyType::Sketch( + SketchKind::new(algorithm.clone(), params), + Default::default(), + ); + let update = SummaryUpdate { + item: Some(SummaryInputExpr::Column( + planner_types::pre_asap::ColumnRef::Named("host".into()), + )), + weight: SummaryInputExpr::Constant(1.0), + weight_domain: Default::default(), + }; + let state = create_planner_accumulator(&family, &update, &Default::default()) + .unwrap() + .snapshot_accumulator(); + let dims = match algorithm { + SketchAlgorithm::Cms => { + let s = state + .as_any() + .downcast_ref::() + .unwrap(); + (s.inner.rows(), s.inner.cols()) + } + SketchAlgorithm::CountSketch => { + let s = state + .as_any() + .downcast_ref::() + .unwrap(); + (s.inner.rows, s.inner.cols) + } + SketchAlgorithm::CmsWithHeap => { + let s = state + .as_any() + .downcast_ref::() + .unwrap(); + (s.inner.rows(), s.inner.cols()) + } + SketchAlgorithm::CountSketchWithHeap => { + let s = state + .as_any() + .downcast_ref::() + .unwrap(); + (s.inner.rows(), s.inner.cols()) + } + _ => unreachable!(), + }; + assert_eq!(dims, (3, 128), "{algorithm:?}"); + } + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/hll_sketch.rs b/crates/asap-physical-operators/src/summary_kernels/hll_sketch.rs new file mode 100644 index 00000000..089261c5 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/hll_sketch.rs @@ -0,0 +1,79 @@ +//! HyperLogLog distinct-count summary over `asap_sketchlib::HllSketch`. +use crate::{AggregateCore, KernelError}; +use asap_sketchlib::{HllSketch, HllVariant}; +use planner_types::post_asap::SketchQuery; + +#[derive(Debug, Clone)] +pub struct HllSketchAccumulator { + pub inner: HllSketch, +} + +impl HllSketchAccumulator { + pub fn new(variant: HllVariant, precision: u32) -> Self { + Self { + inner: HllSketch::new(variant, precision), + } + } +} + +impl AggregateCore for HllSketchAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with(&self, other: &dyn AggregateCore) -> Result, KernelError> { + let other = other + .as_any() + .downcast_ref::() + .ok_or("HLL merges only with HLL")?; + Ok(Box::new(Self { + inner: HllSketch::merge_refs(&[&self.inner, &other.inner])?, + })) + } + + /// Distinct count. A bare `PointCount` over an HLL also reads the distinct count. + fn estimate(&self, query: &SketchQuery) -> Result { + match query { + SketchQuery::Cardinality | SketchQuery::PointCount { value: None, .. } => { + Ok(self.inner.estimate()) + } + other => Err(format!("HLL does not answer {other:?}").into()), + } + } + + fn approx_memory_bytes(&self) -> usize { + std::mem::size_of::().saturating_add(self.inner.registers.capacity()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + // Distinct count after merge counts overlapping items once. + #[test] + fn merged_cardinality_deduplicates_overlap() { + let (mut a, mut b) = ( + HllSketchAccumulator::new(HllVariant::Regular, 12), + HllSketchAccumulator::new(HllVariant::Regular, 12), + ); + for v in 0..1000u32 { + a.inner.update(&v.to_le_bytes()); + b.inner.update(&(v + 500).to_le_bytes()); + } + let merged = a.merge_with(&b).unwrap(); + let estimate = merged.estimate(&SketchQuery::Cardinality).unwrap(); + assert!((estimate - 1500.0).abs() / 1500.0 < 0.05, "{estimate}"); + } + + // HLL does not answer quantiles. + #[test] + fn rejects_quantile() { + let hll = HllSketchAccumulator::new(HllVariant::Regular, 12); + assert!(hll.estimate(&SketchQuery::Quantile { q: 0.5 }).is_err()); + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/hydra_kll.rs b/crates/asap-physical-operators/src/summary_kernels/hydra_kll.rs new file mode 100644 index 00000000..c0347e69 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/hydra_kll.rs @@ -0,0 +1,55 @@ +use crate::{AggregateCore, KeyByLabelValues}; +use asap_sketchlib::HydraKllSketch; + +/// HydraKLL (shared-grouping quantiles) over `asap_sketchlib::HydraKllSketch`. +#[derive(Debug, Clone)] +pub struct HydraKllSketchAccumulator { + pub inner: HydraKllSketch, +} + +impl HydraKllSketchAccumulator { + pub fn new(row_num: usize, col_num: usize, k: u16) -> Self { + Self { + inner: HydraKllSketch::new(row_num, col_num, k), + } + } + + pub fn update(&mut self, key: &KeyByLabelValues, value: f64) { + self.inner.update(&key.to_semicolon_str(), value); + } + + pub fn query_key(&self, key: &KeyByLabelValues, quantile: f64) -> f64 { + self.inner.quantile(&key.to_semicolon_str(), quantile) + } +} + +impl AggregateCore for HydraKllSketchAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with( + &self, + other: &dyn AggregateCore, + ) -> Result, Box> { + let hk = other + .as_any() + .downcast_ref::() + .ok_or("Failed to downcast to HydraKllSketchAccumulator")?; + + let mut merged = self.clone(); + merged.inner.merge(&hk.inner)?; + Ok(Box::new(merged)) + } + + fn approx_memory_bytes(&self) -> usize { + // HydraKLL is a row*col grid of KLL sketches; typical instances + // are on the order of tens of KiB. 32 KiB is a conservative + // per-instance default. + 32 * 1024 + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/increase.rs b/crates/asap-physical-operators/src/summary_kernels/increase.rs new file mode 100644 index 00000000..f92ee582 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/increase.rs @@ -0,0 +1,213 @@ +use crate::{AggregateCore, Measurement}; +use serde::{Deserialize, Serialize}; + +/// Accumulator for tracking increases in counter metrics +/// Stores the starting and last seen measurements with timestamps +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct IncreaseAccumulator { + pub starting_measurement: Measurement, + pub starting_timestamp: i64, + pub last_seen_measurement: Measurement, + pub last_seen_timestamp: i64, + /// Sum of monotonic deltas, adding the post-reset value whenever the + /// counter decreases. This is the reset correction Prometheus applies. + #[serde(default)] + pub total_increase: f64, + #[serde(default)] + pub sample_count: u64, +} + +impl IncreaseAccumulator { + /// Merge two counter intervals without a temporary collection. Ties retain + /// the left input, matching the stable ordering of multi-pane merges. + pub(crate) fn merge_pair(left: &Self, right: &Self) -> Self { + let (first, second) = if left.starting_timestamp <= right.starting_timestamp { + (left, right) + } else { + (right, left) + }; + let mut merged = first.clone(); + if second.starting_timestamp > merged.last_seen_timestamp { + merged.total_increase += + if second.starting_measurement.value >= merged.last_seen_measurement.value { + second.starting_measurement.value - merged.last_seen_measurement.value + } else { + second.starting_measurement.value + }; + } + merged.total_increase += second.total_increase; + merged.sample_count = merged.sample_count.saturating_add(second.sample_count); + if second.last_seen_timestamp > merged.last_seen_timestamp { + merged.last_seen_measurement = second.last_seen_measurement.clone(); + merged.last_seen_timestamp = second.last_seen_timestamp; + } + + merged + } + + pub fn new( + starting_measurement: Measurement, + starting_timestamp: i64, + last_seen_measurement: Measurement, + last_seen_timestamp: i64, + ) -> Self { + let total_increase = if last_seen_timestamp <= starting_timestamp { + 0.0 + } else if last_seen_measurement.value >= starting_measurement.value { + last_seen_measurement.value - starting_measurement.value + } else { + last_seen_measurement.value + }; + let sample_count = if last_seen_timestamp > starting_timestamp { + 2 + } else { + 1 + }; + Self { + starting_measurement, + starting_timestamp, + last_seen_measurement, + last_seen_timestamp, + total_increase, + sample_count, + } + } + + pub fn update(&mut self, measurement: Measurement, timestamp: i64) { + if timestamp < self.last_seen_timestamp { + return; + } + if timestamp == self.last_seen_timestamp { + return; + } + if measurement.value >= self.last_seen_measurement.value { + self.total_increase += measurement.value - self.last_seen_measurement.value; + } else { + self.total_increase += measurement.value; + } + self.last_seen_measurement = measurement; + self.last_seen_timestamp = timestamp; + self.sample_count = self.sample_count.saturating_add(1); + } +} + +impl AggregateCore for IncreaseAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with( + &self, + other: &dyn AggregateCore, + ) -> Result, Box> { + // Downcast to IncreaseAccumulator + let other_increase = other + .as_any() + .downcast_ref::() + .ok_or("Failed to downcast to IncreaseAccumulator")?; + + let merged = Self::merge_pair(self, other_increase); + Ok(Box::new(merged)) + } + + fn approx_memory_bytes(&self) -> usize { + // Two Measurements + two i64s. Measurements are a few f64 fields. + std::mem::size_of::() + } +} + +impl IncreaseAccumulator { + /// PromQL-style increase or rate, extrapolated to `range_ms` when given. + pub(crate) fn extrapolated_value( + &self, + range_ms: Option<(i64, i64)>, + is_rate: bool, + ) -> Result> { + if self.sample_count < 2 || self.last_seen_timestamp <= self.starting_timestamp { + return Err("at least two ordered counter samples are required".into()); + } + let sampled_interval = (self.last_seen_timestamp - self.starting_timestamp) as f64 / 1000.0; + let Some((range_start, range_end)) = range_ms else { + return Ok(if is_rate { + self.total_increase / sampled_interval + } else { + self.total_increase + }); + }; + if range_end <= range_start { + return Err("invalid counter evaluation range".into()); + } + + let mut duration_to_start = + (self.starting_timestamp.saturating_sub(range_start)) as f64 / 1000.0; + let duration_to_end = (range_end.saturating_sub(self.last_seen_timestamp)) as f64 / 1000.0; + let average_sample_interval = sampled_interval / (self.sample_count - 1) as f64; + let extrapolation_threshold = average_sample_interval * 1.1; + + if self.total_increase > 0.0 && self.starting_measurement.value >= 0.0 { + let duration_to_zero = + sampled_interval * (self.starting_measurement.value / self.total_increase); + duration_to_start = duration_to_start.min(duration_to_zero); + } + let mut extrapolate_to = sampled_interval; + extrapolate_to += if duration_to_start < extrapolation_threshold { + duration_to_start.max(0.0) + } else { + average_sample_interval / 2.0 + }; + extrapolate_to += if duration_to_end < extrapolation_threshold { + duration_to_end.max(0.0) + } else { + average_sample_interval / 2.0 + }; + let mut factor = extrapolate_to / sampled_interval; + if is_rate { + factor /= (range_end - range_start) as f64 / 1000.0; + } + Ok(self.total_increase * factor) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_increase_accumulator_creation() { + let starting_measurement = Measurement::new(10.0); + let last_seen_measurement = Measurement::new(25.0); + let acc = IncreaseAccumulator::new( + starting_measurement.clone(), + 1000, + last_seen_measurement.clone(), + 2000, + ); + + assert_eq!(acc.starting_measurement.value, 10.0); + assert_eq!(acc.starting_timestamp, 1000); + assert_eq!(acc.last_seen_measurement.value, 25.0); + assert_eq!(acc.last_seen_timestamp, 2000); + } + + #[test] + fn test_increase_accumulator_update() { + let starting_measurement = Measurement::new(10.0); + let mut acc = IncreaseAccumulator::new( + starting_measurement.clone(), + 1000, + starting_measurement.clone(), + 1000, + ); + + let new_measurement = Measurement::new(25.0); + acc.update(new_measurement.clone(), 2000); + + assert_eq!(acc.last_seen_measurement.value, 25.0); + assert_eq!(acc.last_seen_timestamp, 2000); + assert_eq!(acc.starting_measurement.value, 10.0); // Should remain unchanged + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/mod.rs b/crates/asap-physical-operators/src/summary_kernels/mod.rs new file mode 100644 index 00000000..478e75d0 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/mod.rs @@ -0,0 +1,26 @@ +//! In-memory summary state: thin adapters over `asap_sketchlib` and exact Planner state. +pub mod count_min_sketch; +pub mod count_min_sketch_with_heap; +pub mod count_sketch; +pub mod count_sketch_with_heap; +pub mod datasketches_kll; +pub mod dd_sketch; +pub mod exact; +pub mod hll_sketch; +pub mod hydra_kll; +pub mod increase; +pub mod univmon; + +pub use count_min_sketch::*; +pub use count_min_sketch_with_heap::*; +pub use count_sketch::*; +pub use count_sketch_with_heap::*; +pub use datasketches_kll::*; +pub use dd_sketch::*; +pub use hll_sketch::*; +pub use hydra_kll::*; +pub use increase::*; + +pub mod factory; +pub mod traits; +pub mod weighted_frequency; diff --git a/crates/asap-physical-operators/src/summary_kernels/traits.rs b/crates/asap-physical-operators/src/summary_kernels/traits.rs new file mode 100644 index 00000000..7d866077 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/traits.rs @@ -0,0 +1,35 @@ +use planner_types::post_asap::SketchQuery; + +pub type KernelError = Box; + +/// In-memory state of one population's summary. +/// +/// Kernels adapt `asap_sketchlib` structures (or exact Planner state) to the +/// operations physical operators need: merge, typed readout and memory +/// accounting. Grouping belongs to operators; byte encodings belong to +/// `asap_sketchlib` and deployments. +pub trait AggregateCore: Send + Sync { + fn clone_boxed_core(&self) -> Box; + + fn as_any(&self) -> &dyn std::any::Any; + + /// Merge with a state of the same family and shape, leaving both inputs unchanged. + fn merge_with(&self, other: &dyn AggregateCore) -> Result, KernelError>; + + /// Answer a sketch readout. Exact states are read through + /// [`ExactAccumulator::readout`](super::exact::ExactAccumulator::readout). + fn estimate(&self, query: &SketchQuery) -> Result { + Err(format!("{query:?} is not supported by this summary").into()) + } + + /// Approximate in-memory footprint, used for execution memory reservations. + fn approx_memory_bytes(&self) -> usize { + 4096 + } +} + +impl Clone for Box { + fn clone(&self) -> Self { + self.clone_boxed_core() + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/univmon.rs b/crates/asap-physical-operators/src/summary_kernels/univmon.rs new file mode 100644 index 00000000..bffd8afa --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/univmon.rs @@ -0,0 +1,109 @@ +//! One frequency state shared by count, distinct, L2 and entropy readouts. + +use crate::AggregateCore; +use asap_sketchlib::{DataInput, UnivMon}; + +type Error = Box; + +#[derive(Debug, Clone)] +pub struct UnivMonAccumulator { + inner: UnivMon, +} + +impl UnivMonAccumulator { + /// Empty the sketch in place, keeping its shape. + pub(crate) fn clear(&mut self) { + self.inner.free(); + } + + pub fn new(heap_size: usize, rows: usize, cols: usize, layers: usize) -> Result { + if heap_size == 0 || cols == 0 || !(1..=20).contains(&rows) || !(1..=64).contains(&layers) { + return Err("invalid UnivMon dimensions".into()); + } + rows.checked_mul(cols) + .and_then(|n| n.checked_mul(layers)) + .ok_or("UnivMon dimensions overflow")?; + Ok(Self { + inner: UnivMon::init_univmon(heap_size, rows, cols, layers), + }) + } + + /// Each non-NaN sample is one occurrence. Signed zero has one identity. + pub fn insert_sample(&mut self, value: f64) -> Result<(), Error> { + if value.is_nan() { + return Ok(()); + } + self.inner + .bucket_size + .checked_add(1) + .ok_or("UnivMon count overflow")?; + let bits = if value == 0.0 { 0 } else { value.to_bits() }; + self.inner.insert(&DataInput::U64(bits), 1); + Ok(()) + } + + fn compatible(&self, other: &Self) -> bool { + ( + self.inner.heap_size, + self.inner.sketch_row, + self.inner.sketch_col, + self.inner.layer_size, + ) == ( + other.inner.heap_size, + other.inner.sketch_row, + other.inner.sketch_col, + other.inner.layer_size, + ) + } + + pub fn dimensions(&self) -> (usize, usize, usize, usize) { + ( + self.inner.heap_size, + self.inner.sketch_row, + self.inner.sketch_col, + self.inner.layer_size, + ) + } + + pub fn merge_in_place(&mut self, other: &Self) -> Result<(), Error> { + if !self.compatible(other) { + return Err("incompatible UnivMon dimensions".into()); + } + self.inner + .bucket_size + .checked_add(other.inner.bucket_size) + .ok_or("UnivMon count overflow")?; + self.inner.merge(&other.inner); + Ok(()) + } +} + +impl AggregateCore for UnivMonAccumulator { + fn approx_memory_bytes(&self) -> usize { + std::mem::size_of::().saturating_add( + self.inner.layer_size.saturating_mul( + self.inner + .sketch_row + .saturating_mul(self.inner.sketch_col) + .saturating_mul(16) + .saturating_add(self.inner.heap_size.saturating_mul(256)), + ), + ) + } + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with(&self, other: &dyn AggregateCore) -> Result, Error> { + let other = other + .as_any() + .downcast_ref::() + .ok_or("expected UnivMon state")?; + let mut merged = self.clone(); + merged.merge_in_place(other)?; + Ok(Box::new(merged)) + } +} diff --git a/crates/asap-physical-operators/src/summary_kernels/weighted_frequency.rs b/crates/asap-physical-operators/src/summary_kernels/weighted_frequency.rs new file mode 100644 index 00000000..6f008532 --- /dev/null +++ b/crates/asap-physical-operators/src/summary_kernels/weighted_frequency.rs @@ -0,0 +1,152 @@ +//! ASAP type and trait adapter for sketchlib's Float64 weighted frequency kernel. +use crate::AggregateCore; +use crate::{values::Value, Error}; +pub use asap_sketchlib::FrequencyAlgorithm; +use asap_sketchlib::{FrequencyIdentity, WeightedFrequency as Kernel, WeightedFrequencyError}; +use serde::{Deserialize, Serialize}; + +fn adapt_error(error: WeightedFrequencyError) -> Error { + match error { + WeightedFrequencyError::Invalid(message) => Error::Invalid(message), + WeightedFrequencyError::Update(message) => Error::Operator(message), + } +} +fn identity(value: &Value) -> Result { + Ok(match value { + Value::Null => FrequencyIdentity::Null, + Value::Bool(v) => FrequencyIdentity::Bool(*v), + Value::Int64(v) => FrequencyIdentity::Int64(*v), + Value::Float64(v) => FrequencyIdentity::Float64(*v), + Value::Utf8(v) => FrequencyIdentity::Utf8(v.to_string()), + _ => { + return Err(Error::Invalid( + "unsupported weighted frequency identity".into(), + )) + } + }) +} +fn value(identity: FrequencyIdentity) -> Value { + match identity { + FrequencyIdentity::Null => Value::Null, + FrequencyIdentity::Bool(v) => Value::Bool(v), + FrequencyIdentity::Int64(v) => Value::Int64(v), + FrequencyIdentity::Float64(v) => Value::Float64(v), + FrequencyIdentity::Utf8(v) => Value::Utf8(v.into()), + } +} +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(transparent)] +pub struct WeightedFrequency { + inner: Kernel, +} +impl WeightedFrequency { + pub(crate) fn configuration( + kind: &planner_types::post_asap::SketchKind, + ) -> Result<(FrequencyAlgorithm, usize, usize, usize), Error> { + use planner_types::post_asap::{SketchAlgorithm as A, SketchParams as P}; + let (algorithm, width, depth, capacity) = match (kind.algorithm(), kind.params()) { + ( + A::CmsWithHeap, + P::CmsWithHeap { + width, + depth, + heap_size, + }, + ) => (FrequencyAlgorithm::Cms, *width, *depth, *heap_size), + ( + A::CountSketchWithHeap, + P::CountSketchWithHeap { + width, + depth, + heap_size, + }, + ) if depth % 2 == 1 => (FrequencyAlgorithm::CountSketch, *width, *depth, *heap_size), + _ => { + return Err(Error::Invalid( + "unsupported weighted frequency family or depth".into(), + )) + } + }; + if width == 0 || depth == 0 || capacity == 0 { + return Err(Error::Invalid( + "invalid weighted frequency dimensions".into(), + )); + } + Ok((algorithm, width as usize, depth as usize, capacity as usize)) + } + + pub(crate) fn algorithm(&self) -> FrequencyAlgorithm { + self.inner.algorithm() + } + pub(crate) fn shape(&self) -> (usize, usize, usize) { + self.inner.shape() + } + pub fn new( + algorithm: FrequencyAlgorithm, + width: usize, + depth: usize, + capacity: usize, + ) -> Result { + Kernel::new(algorithm, width, depth, capacity) + .map(|inner| Self { inner }) + .map_err(adapt_error) + } + pub fn update(&mut self, values: &[Value], weight: f64) -> Result<(), Error> { + let values = values.iter().map(identity).collect::, _>>()?; + self.inner.update(&values, weight).map_err(adapt_error) + } + pub fn rows(&self, n: usize) -> Vec> { + self.inner + .topk(n) + .into_iter() + .map(|(items, score)| { + let mut row = items.into_iter().map(value).collect::>(); + row.push(Value::Float64(score)); + row + }) + .collect() + } +} +impl AggregateCore for WeightedFrequency { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + fn as_any(&self) -> &dyn std::any::Any { + self + } + fn merge_with( + &self, + other: &dyn AggregateCore, + ) -> Result, Box> { + let other = other + .as_any() + .downcast_ref::() + .ok_or("weighted frequency state type mismatch")?; + Ok(Box::new(Self { + inner: self.inner.merge(&other.inner)?, + })) + } + fn approx_memory_bytes(&self) -> usize { + self.inner.approx_memory_bytes() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + // Merge uses the same Float64 state representation and rejects other shapes. + #[test] + fn compatible_merge_preserves_fractional_weights() { + let mut left = WeightedFrequency::new(FrequencyAlgorithm::Cms, 4096, 5, 8).unwrap(); + let mut right = left.clone(); + left.update(&[Value::Int64(7)], 0.125).unwrap(); + right.update(&[Value::Int64(7)], 0.25).unwrap(); + let merged = left.merge_with(&right).unwrap(); + let merged = merged.as_any().downcast_ref::().unwrap(); + assert!(matches!(merged.rows(1)[0][1], Value::Float64(0.375))); + assert!(left + .merge_with(&WeightedFrequency::new(FrequencyAlgorithm::Cms, 32, 5, 8).unwrap()) + .is_err()); + } +} diff --git a/crates/asap-physical-operators/src/values.rs b/crates/asap-physical-operators/src/values.rs new file mode 100644 index 00000000..6e761c11 --- /dev/null +++ b/crates/asap-physical-operators/src/values.rs @@ -0,0 +1,381 @@ +//! Runtime values preserve Planner schemas; summary states are typed values too. +use crate::AggregateCore; +use crate::Error; +use planner_types::{ + post_asap::{SummaryFamilyType, SummarySchema}, + pre_asap::DataType, +}; +use std::{cmp::Ordering, sync::Arc}; +pub type Schema = Arc; +#[derive(Clone, serde::Serialize, serde::Deserialize)] +pub enum Value { + Null, + Bool(bool), + Int64(i64), + Float64(f64), + Utf8(Arc), + Timestamp(i64), + Date(i32), + Interval { + months: i32, + days: i32, + nanos: i64, + }, + List(Arc<[Value]>), + Struct(Arc<[Value]>), + Map(Arc<[(Value, Value)]>), + #[serde(skip)] + Summary { + family: SummaryFamilyType, + state: Arc, + }, +} +impl std::fmt::Debug for Value { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Summary { family, .. } => f.debug_tuple("Summary").field(family).finish(), + _ => write!(f, "{:?}", self.key()), + } + } +} +impl Value { + pub fn bytes(&self) -> usize { + std::mem::size_of::() + + match self { + Self::Utf8(s) => s.len(), + Self::List(v) | Self::Struct(v) => v.iter().map(Self::bytes).sum(), + Self::Map(v) => v.iter().map(|(k, v)| k.bytes() + v.bytes()).sum(), + Self::Summary { state, .. } => state.approx_memory_bytes(), + _ => 0, + } + } + pub fn matches(&self, dtype: &DataType, nullable: bool) -> bool { + if matches!(self, Self::Null) { + return nullable || matches!(dtype, DataType::Null); + } + match (self, dtype) { + (Self::Bool(_), DataType::Bool) + | (Self::Int64(_), DataType::Int64) + | (Self::Float64(_), DataType::Float64) + | (Self::Utf8(_), DataType::Utf8) + | (Self::Timestamp(_), DataType::Timestamp) + | (Self::Date(_), DataType::Date) + | (Self::Interval { .. }, DataType::Interval) => true, + (Self::List(v), DataType::List { element }) => v + .iter() + .all(|v| v.matches(&element.dtype, element.nullable)), + (Self::Struct(v), DataType::Struct { fields }) => { + v.len() == fields.len() + && v.iter() + .zip(fields) + .all(|(v, f)| v.matches(&f.dtype, f.nullable)) + } + ( + Self::Map(v), + DataType::Map { + key, + value, + value_nullable, + }, + ) => v + .iter() + .all(|(k, v)| k.matches(key, false) && v.matches(value, *value_nullable)), + _ => false, + } + } + /// Stable typed equality key. Zero signs and NaN payloads form one group. + pub fn key(&self) -> Result, Error> { + let mut out = Vec::new(); + macro_rules! number { + ($tag:expr,$v:expr) => {{ + out.push($tag); + out.extend_from_slice(&$v.to_le_bytes()); + }}; + } + match self { + Self::Null => out.push(0), + Self::Bool(v) => out.extend([1, *v as u8]), + Self::Int64(v) => number!(2, v), + Self::Float64(v) => { + let bits = if *v == 0. { + 0 + } else if v.is_nan() { + f64::NAN.to_bits() + } else { + v.to_bits() + }; + number!(3, bits); + } + Self::Utf8(v) => { + out.push(4); + out.extend(v.as_bytes()); + } + Self::Timestamp(v) => number!(5, v), + Self::Date(v) => number!(6, v), + Self::Interval { + months, + days, + nanos, + } => { + number!(7, months); + number!(8, days); + number!(9, nanos); + } + Self::List(v) | Self::Struct(v) => { + out.push(if matches!(self, Self::List(_)) { + 10 + } else { + 11 + }); + for v in v.iter() { + let key = v.key()?; + out.extend((key.len() as u64).to_le_bytes()); + out.extend(key); + } + } + Self::Map(v) => { + out.push(12); + for (k, v) in v.iter() { + for value in [k, v] { + let key = value.key()?; + out.extend((key.len() as u64).to_le_bytes()); + out.extend(key); + } + } + } + Self::Summary { .. } => { + return Err(Error::Invalid( + "summary states cannot be grouping keys".into(), + )) + } + } + Ok(out) + } + pub fn compare(&self, other: &Self) -> Result { + Ok(match (self, other) { + (Self::Null, Self::Null) => Ordering::Equal, + (Self::Int64(a), Self::Int64(b)) | (Self::Timestamp(a), Self::Timestamp(b)) => a.cmp(b), + (Self::Float64(a), Self::Float64(b)) => { + if a == b { + Ordering::Equal + } else { + a.total_cmp(b) + } + } + (Self::Utf8(a), Self::Utf8(b)) => a.cmp(b), + (Self::Bool(a), Self::Bool(b)) => a.cmp(b), + (Self::Date(a), Self::Date(b)) => a.cmp(b), + (Self::Map(left), Self::Map(right)) => { + let mut result = Ordering::Equal; + for ((lk, lv), (rk, rv)) in left.iter().zip(right.iter()) { + result = lk.compare(rk)?; + if result != Ordering::Equal { + break; + } + result = match (lv, rv) { + (Self::Null, Self::Null) => Ordering::Equal, + (Self::Null, _) => Ordering::Greater, + (_, Self::Null) => Ordering::Less, + _ => lv.compare(rv)?, + }; + if result != Ordering::Equal { + break; + } + } + if result == Ordering::Equal { + left.len().cmp(&right.len()) + } else { + result + } + } + _ => { + return Err(Error::Operator( + "values do not have a supported common ordering".into(), + )) + } + }) + } +} +#[derive(Clone, Debug)] +pub struct Batch { + schema: Schema, + rows: Vec>, +} +impl Batch { + pub fn try_new(schema: Schema, rows: Vec>) -> Result { + validate_schema(&schema)?; + for row in &rows { + if row.len() != schema.fields.len() { + return Err(Error::Invalid( + "row width differs from Planner schema".into(), + )); + } + for (value, field) in row.iter().zip(&schema.fields) { + let matches = match (&field.dtype, value) { + (SummaryFamilyType::Plain(dtype), value) => { + value.matches(dtype, field.nullable) + } + (expected, Value::Summary { family, state }) => { + expected == family && validate_state(family, state.as_ref()).is_ok() + } + _ => false, + }; + if !matches { + return Err(Error::Invalid(format!( + "value differs from type of {}", + field.name + ))); + } + } + } + Ok(Self { schema, rows }) + } + pub fn schema(&self) -> &Schema { + &self.schema + } + pub fn rows(&self) -> &[Vec] { + &self.rows + } + pub fn bytes(&self) -> usize { + std::mem::size_of::() + + self.rows.capacity() * std::mem::size_of::>() + + self + .rows + .iter() + .flat_map(|r| r.iter()) + .map(Value::bytes) + .sum::() + } +} + +pub(crate) use crate::capability::validate_native_family as validate_family; + +fn validate_state(family: &SummaryFamilyType, state: &dyn AggregateCore) -> Result<(), Error> { + use crate::summary_kernels::{ + datasketches_kll::DatasketchesKLLAccumulator, dd_sketch::DDSketchAccumulator, + exact::ExactAccumulator, hll_sketch::HllSketchAccumulator, + }; + use planner_types::post_asap::SketchParams; + validate_family(family)?; + let valid = match family { + SummaryFamilyType::Sketch(kind, _) + if matches!( + kind.params(), + SketchParams::CmsWithHeap { .. } | SketchParams::CountSketchWithHeap { .. } + ) => + { + use crate::summary_kernels::weighted_frequency::WeightedFrequency; + let (algorithm, width, depth, capacity) = WeightedFrequency::configuration(kind)?; + state + .as_any() + .downcast_ref::() + .is_some_and(|state| { + state.algorithm() == algorithm && state.shape() == (width, depth, capacity) + }) + } + + SummaryFamilyType::ExactAggregate(..) => state + .as_any() + .downcast_ref::() + .is_some_and(|s| s.family() == family && !s.is_keyed()), + SummaryFamilyType::Sketch(kind, _) => match kind.params() { + SketchParams::Kll { k } => state + .as_any() + .downcast_ref::() + .is_some_and(|s| u32::from(s.inner.k()) == *k), + SketchParams::DDSketch { alpha } => state + .as_any() + .downcast_ref::() + .is_some_and(|s| s.inner.alpha == *alpha), + SketchParams::Hll { precision } => state + .as_any() + .downcast_ref::() + .is_some_and(|s| s.inner.precision == u32::from(*precision)), + _ => false, + }, + _ => false, + }; + if valid { + Ok(()) + } else { + Err(Error::Invalid( + "state payload differs from declared family, parameters or population layout".into(), + )) + } +} + +pub(crate) fn validate_schema(schema: &Schema) -> Result<(), Error> { + if schema.time_index.is_some_and(|index| { + schema + .fields + .get(index) + .is_none_or(|field| field.dtype != SummaryFamilyType::Plain(DataType::Timestamp)) + }) { + return Err(Error::Invalid( + "time index must name a Timestamp column".into(), + )); + } + for field in &schema.fields { + if !matches!(field.dtype, SummaryFamilyType::Plain(_)) { + validate_family(&field.dtype)?; + if field.nullable { + return Err(Error::Invalid( + "nullable summary states are not supported".into(), + )); + } + } + } + Ok(()) +} + +#[cfg(test)] +mod weighted_state_tests { + use super::*; + use crate::summary_kernels::weighted_frequency::{FrequencyAlgorithm, WeightedFrequency}; + use planner_types::post_asap::{SketchAlgorithm, SketchKind, SketchParams}; + + // A state cannot acquire a different family or shape merely by relabeling its batch. + #[test] + fn weighted_state_family_and_shape_must_match() { + let cms = SummaryFamilyType::Sketch( + SketchKind::new( + SketchAlgorithm::CmsWithHeap, + SketchParams::CmsWithHeap { + width: 32, + depth: 5, + heap_size: 8, + }, + ), + Default::default(), + ); + let cs = SummaryFamilyType::Sketch( + SketchKind::new( + SketchAlgorithm::CountSketchWithHeap, + SketchParams::CountSketchWithHeap { + width: 32, + depth: 5, + heap_size: 8, + }, + ), + Default::default(), + ); + let state = WeightedFrequency::new(FrequencyAlgorithm::CountSketch, 32, 5, 8).unwrap(); + assert!(validate_state(&cs, &state).is_ok()); + assert!(validate_state(&cms, &state).is_err()); + let wrong_shape = + WeightedFrequency::new(FrequencyAlgorithm::CountSketch, 64, 5, 8).unwrap(); + assert!(validate_state(&cs, &wrong_shape).is_err()); + let even_depth = SummaryFamilyType::Sketch( + SketchKind::new( + SketchAlgorithm::CountSketchWithHeap, + SketchParams::CountSketchWithHeap { + width: 32, + depth: 4, + heap_size: 8, + }, + ), + Default::default(), + ); + assert!(validate_family(&even_depth).is_err()); + } +} diff --git a/crates/asap-physical-operators/tests/deployment.rs b/crates/asap-physical-operators/tests/deployment.rs new file mode 100644 index 00000000..0c261a03 --- /dev/null +++ b/crates/asap-physical-operators/tests/deployment.rs @@ -0,0 +1,90 @@ +//! Exercise the public library without a backend server, store, or scheduler. +use asap_physical_operators::planner::{ + post_asap::{ + GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryFamilyType, + SummaryUpdate, + }, + pre_asap::ColumnRef, +}; +use asap_physical_operators::{factory::create_planner_accumulator, AggregateCore}; + +fn family(k: u32) -> SummaryFamilyType { + SummaryFamilyType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), + GroupingStrategy::PerSubpopulationInstance, + ) +} +fn build(values: &[f64]) -> Box { + let mut operator = create_planner_accumulator( + &family(512), + &SummaryUpdate::column(ColumnRef::SampleValue), + &Default::default(), + ) + .unwrap(); + for (at, value) in values.iter().enumerate() { + operator.validate_single_input(*value).unwrap(); + operator.update_single(*value, at as i64); + } + operator.into_accumulator() +} +fn read(state: &dyn AggregateCore) -> f64 { + state + .estimate(&asap_physical_operators::planner::post_asap::SketchQuery::Quantile { q: 0.5 }) + .unwrap() +} + +// The same kernels work when every build is query-time, when only a prefix +// was precomputed, and when all state was precomputed before the readout. +#[test] +fn raw_partial_and_fully_precomputed_use_the_same_kernels() { + let raw: Vec = (0..128).map(f64::from).collect(); + let raw_only = build(&raw); + let stored_prefix = build(&raw[..64]); + let query_time_suffix = build(&raw[64..]); + let partial = stored_prefix.merge_with(&*query_time_suffix).unwrap(); + let stored_complete = build(&raw); + assert_eq!(read(&*raw_only), read(&*partial)); + assert_eq!(read(&*partial), read(&*stored_complete)); + assert!((read(&*raw_only) - 64.0).abs() <= 1.0); +} + +// A compiler must reject invalid physical parameters before starting execution. +#[test] +fn invalid_kll_parameters_are_rejected_at_binding() { + let result = create_planner_accumulator( + &family(0), + &SummaryUpdate::column(ColumnRef::SampleValue), + &Default::default(), + ); + assert!(result.is_err()); +} + +// Native CountSketch supports the confidence-sized depth used by the backend; +// a packed-wire column-bit budget must not be imposed on this constructor. +#[test] +fn native_count_sketch_dimensions_are_not_packed_wire_dimensions() { + use asap_physical_operators::planner::post_asap::SummaryInputExpr; + use asap_physical_operators::KeyByLabelValues; + let family = SummaryFamilyType::Sketch( + SketchKind::new( + SketchAlgorithm::CountSketchWithHeap, + SketchParams::CountSketchWithHeap { + width: 1200, + depth: 55, + heap_size: 3, + }, + ), + Default::default(), + ); + let mut update = SummaryUpdate::column(ColumnRef::SampleValue); + update.item = Some(SummaryInputExpr::Column(ColumnRef::Named("host".into()))); + let mut operator = create_planner_accumulator(&family, &update, &Default::default()).unwrap(); + let key = KeyByLabelValues::new_with_labels(vec!["a".into()]); + operator.update_keyed(&key, 7.0, 1000); + let state = operator.into_accumulator(); + let state = state + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(state.query_key(&key), 7.0); +}