diff --git a/crates/logical-optimizer/src/accuracy/estimators/univmon.rs b/crates/logical-optimizer/src/accuracy/estimators/univmon.rs index 6a89a105e..4a0293144 100644 --- a/crates/logical-optimizer/src/accuracy/estimators/univmon.rs +++ b/crates/logical-optimizer/src/accuracy/estimators/univmon.rs @@ -1,4 +1,32 @@ -//! UnivMon currently certifies only its exact unit-update total evaluation. +//! UnivMon certifies only its unit-update total, which the kernel reads +//! exactly (`calc_l1` returns the update count). +//! +//! Distinct count, L2 norm and entropy are deliberately left uncertified, +//! so Stage 3 rejects them unless a deployment supplies accuracy evidence: +//! +//! - Liu et al., "One Sketch to Rule Them All" (SIGCOMM 2016), rely on +//! Braverman and Ostrovsky's recursive sketch. With O(log n) layers, each of +//! which returns a (g, ε)-cover of its sampled substream (every g-heavy item +//! with a (1 ± ε) frequency), the recursive G-sum is a (1 ± ε) +//! approximation with probability 1 − δ. The layer sketch's size is only +//! stated asymptotically (O(ε⁻² log(1/δ)) per CountSketch, times +//! polylogarithmic factors). No constants are given that could be inverted +//! into a concrete `(heap_size, sketch_rows, sketch_cols, layers)`. +//! - The kernel (`asap_sketchlib::UnivMon::calc_g_sum_heuristic`, behind +//! `asap-executor`'s `UnivMonAccumulator`) keeps a fixed top-`heap_size` +//! heap per layer and estimates its items with the layer's CountSketch. +//! For distinct counts it also drops items below `L2 / sqrt(heap_size)`. +//! Nothing bounds the probability that a layer's heap is such a cover, so +//! the theorem's premise does not hold for these readouts. Even when every +//! heap holds every item, a CountSketch estimate may be at most zero for a +//! present item, so the distinct count is not certified in that case +//! either. +//! - A per-layer CountSketch does give a sound F₂ (AMS) bound, but `calc_l2` +//! reads the recursive sum, not that estimate. +//! +//! A sound guarantee needs either a readout with a proven bound (for +//! example, L2 from layer 0's CountSketch) or a per-layer cover guarantee for +//! the heap. Both are kernel changes. use super::*; pub(super) fn guarantee(query: &SketchStatistic) -> Option { @@ -6,8 +34,9 @@ pub(super) fn guarantee(query: &SketchStatistic) -> Option { .then(|| ResultGuarantee::exact("univmon_unit_update_total")) } +/// One shape for every readout and requirement: no bound is inverted (see +/// the module docs), so every consumer of a key reads identical states. pub(super) fn size_params() -> SketchParams { - // Baseline dimensions are candidates, not an inverted accuracy bound. SketchParams::UnivMon { heap_size: 256, sketch_rows: 5, @@ -15,3 +44,42 @@ pub(super) fn size_params() -> SketchParams { layers: 16, } } + +#[cfg(test)] +mod tests { + use super::*; + use asap_types::ir::scalar::ColumnRef; + use asap_types::ir::schema::{GroupingStrategy, SketchKind}; + + fn family() -> FieldDataType { + FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::UnivMon, size_params()), + GroupingStrategy::default(), + ) + } + + /// The exact total is certified; distinct count, L2 and entropy have no + /// sound bound for the shipped kernel and stay uncertified. + #[test] + fn only_the_total_is_certified() { + let total = SketchStatistic::PointCount { + key: ColumnRef::SampleValue, + value: None, + }; + assert!(DefaultAccuracyModel + .local_guarantee(&family(), &total) + .is_some_and(|g| g.is_exact())); + for query in [ + SketchStatistic::Cardinality, + SketchStatistic::FrequencyL2, + SketchStatistic::FrequencyEntropy, + ] { + assert!( + DefaultAccuracyModel + .local_guarantee(&family(), &query) + .is_none(), + "{query:?}" + ); + } + } +} diff --git a/crates/logical-optimizer/src/pass2/summary_capability.rs b/crates/logical-optimizer/src/pass2/summary_capability.rs index 586f2d714..fe5f387f1 100644 --- a/crates/logical-optimizer/src/pass2/summary_capability.rs +++ b/crates/logical-optimizer/src/pass2/summary_capability.rs @@ -29,11 +29,14 @@ use crate::pass1::logical_candidates::{ }; use crate::pass1::replacement::{accuracy_budget, accuracy_target}; -/// The estimates one summary serves: any quantile of one column from one -/// quantile summary (KLL, DDSketch). +/// The estimates one summary serves over one column. #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum Capability { + /// Any quantile, from one KLL or DDSketch. Quantile, + /// Distinct count, L2 norm and entropy of the value frequencies, from one + /// UnivMon (#509 Example 2). + FrequencyMoments, } /// What a summary for a target would ingest and which estimates it would @@ -77,6 +80,13 @@ fn key(node: &OperatorNode) -> Option> { } let (capability, column) = match intent { AggIntent::Quantile { col, .. } => (Capability::Quantile, *col), + // A distinct-tuple count has no UnivMon alternative. + AggIntent::Cardinality { cols, .. } if cols.len() <= 1 => { + (Capability::FrequencyMoments, cols.first().copied()) + } + AggIntent::FrequencyL2 { col, .. } | AggIntent::FrequencyEntropy { col, .. } => { + (Capability::FrequencyMoments, *col) + } _ => return None, }; Some(Key { @@ -259,6 +269,43 @@ mod tests { assert!(!shared.resized); } + /// Distinct count, entropy and L2 over one input form one key. UnivMon's + /// shape does not depend on the requirement, so all three alternatives + /// are the same state; the distinct count's other summaries are re-sized + /// for the strictest ε. + #[test] + fn frequency_moments_share_one_univmon() { + let base = inventory(&[ + ("distinct_over_time(src[1m])", 0.02), + ("entropy_over_time(src[1m])", 0.05), + ("l2_over_time(src[1m])", 0.01), + ]); + let shared = share_summary_capability(&base).unwrap().expect("one key"); + assert!(shared.resized); + let univmon: Vec<_> = shared + .inventory + .targets + .iter() + .map(|t| { + t.alternatives + .iter() + .find(|a| matches!(a, Realization::Sketch(kind) if *kind.algorithm() == SketchAlgorithm::UnivMon)) + .cloned() + .expect("a UnivMon alternative") + }) + .collect(); + assert!(univmon.iter().all(|u| *u == univmon[0])); + let hll = |inv: &LocalLogicalCandidates| { + inv.targets[0].alternatives.iter().find_map(|a| match a { + Realization::Sketch(kind) if *kind.algorithm() == SketchAlgorithm::Hll => { + Some(kind.clone()) + } + _ => None, + }) + }; + assert_ne!(hll(&shared.inventory), hll(&base)); + } + /// A different window, selector or estimate family is a different key. #[test] fn different_summary_input_or_window_is_not_shared() { diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index 85ffdb5c4..ba5813d8a 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -503,8 +503,9 @@ impl AccuracyModel for UnivMonEvidence { /// Distinct count, entropy and L2 over one input, certified by an accuracy /// model and selected by a cost model preferring UnivMon, read one UnivMon state: #515 sharing is the summary-capability rule -/// when the states are identical. The facade's stage pipeline does not plan -/// UnivMon sharing, so this runs the legacy search with the test model. +/// when the states are identical. This runs the legacy search; the stage +/// pipeline's counterpart is +/// `frequency_moments_share_one_univmon_in_the_stage_pipeline`. #[test] fn certified_frequency_evaluations_share_one_univmon_state() { let queries = [ @@ -554,3 +555,56 @@ fn certified_frequency_evaluations_share_one_univmon_state() { assert_eq!(states.len(), 3); assert!(states.iter().all(|state| Rc::ptr_eq(state, &states[0]))); } + +/// #509 Example 2 through the stage pipeline: distinct count, entropy and L2 +/// of one input over one window, with requirements ε = 0.02, 0.05 and 0.01. +/// The summary-capability variant gives the three targets one UnivMon. +/// Certified by the synthetic model, that one state serves all three +/// queries. The built-in model has no sound bound for these readouts, so +/// Stage 3 rejects every UnivMon plan. +#[tokio::test] +async fn frequency_moments_share_one_univmon_in_the_stage_pipeline() { + let queries = [ + ("distinct_over_time(src[1m])", 0.02), + ("entropy_over_time(src[1m])", 0.05), + ("l2_over_time(src[1m])", 0.01), + ]; + let workload = promql_workload(&queries); + let plan = |models: PlanningModels<'static>| { + let workload = workload.clone(); + async move { + let input = UserInput::new( + &workload, + FrontendInput::Promql { + now_ms: NOW_MS, + histograms: None, + }, + models, + ); + e2e_plan(input).await.expect("workload plans") + } + }; + let univmon_states = |output: &PlanOutput| -> Vec> { + output + .plans + .iter() + .flat_map(plan_states) + .filter(|state| { + matches!( + &state.operator, + asap_types::ir::Operator::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) + if kind.algorithm() == &SketchAlgorithm::UnivMon + ) + }) + .collect() + }; + + let certified = plan(PlanningModels::builtin().with_accuracy(&UnivMonEvidence)).await; + let states = univmon_states(&certified); + assert_eq!(states.len(), 3, "{:?}", certified.selection); + assert!(states.iter().all(|state| Rc::ptr_eq(state, &states[0]))); + assert_eq!(unique_deployments(&certified), 1); + + let builtin = plan(PlanningModels::builtin()).await; + assert!(univmon_states(&builtin).is_empty()); +}