From f9502890a71086386a8eb403b443db0eba6d29ce Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 06:29:30 +0000 Subject: [PATCH 1/3] feat(physical): compile summaries over raw-sample precompute boundaries A precompute boundary at a raw time-series scan binds as raw sample rows ($population label map, $timestamp, value). The label map is the complete source identity, so per-series and grouped SummaryAgg lower over it with constant, column or unit-frequency (HLL) updates, and keyed heaps resolve their items from labels, the sample value or the canonical label identity. Co-Authored-By: Claude Opus 5.5 --- .../src/expressions/mod.rs | 57 ++ .../src/physical_planner/precompute.rs | 198 ++++++- .../tests/precompute_raw_samples.rs | 495 ++++++++++++++++++ 3 files changed, 733 insertions(+), 17 deletions(-) create mode 100644 crates/integration-tests/tests/precompute_raw_samples.rs diff --git a/crates/asap-physical-operators/src/expressions/mod.rs b/crates/asap-physical-operators/src/expressions/mod.rs index 2916cca7..1ff7b4db 100644 --- a/crates/asap-physical-operators/src/expressions/mod.rs +++ b/crates/asap-physical-operators/src/expressions/mod.rs @@ -23,6 +23,16 @@ pub enum Expression { labels: Vec, without: bool, }, + /// One label of a label map; an absent label reads as empty, as in PromQL. + Label { + column: usize, + name: String, + }, + /// Canonical encoding of a complete label map, identical to + /// `promql_rows::encode_series_identity`. + LabelIdentity { + column: usize, + }, Literal { value: Value, dtype: DataType, @@ -119,6 +129,21 @@ impl Expression { } Ok((expected, false)) } + Label { column, .. } | LabelIdentity { column } => { + if plain(input, *column)? + != ( + &DataType::Map { + key: Box::new(DataType::Utf8), + value: Box::new(DataType::Utf8), + value_nullable: false, + }, + false, + ) + { + return Err(invalid("label read requires a non-null Utf8 map")); + } + Ok((DataType::Utf8, false)) + } Column(i) => { let (t, n) = plain(input, *i)?; Ok((t.clone(), n)) @@ -251,6 +276,38 @@ impl Expression { } } Planner(expression) => expression.evaluate(row)?, + Label { column, name } => { + let Value::Map(entries) = &row[*column] else { + return Err(invalid("label read requires a map")); + }; + let mut found = None; + for (key, value) in entries.iter() { + let (Value::Utf8(key), Value::Utf8(value)) = (key, value) else { + return Err(invalid("label read requires Utf8 entries")); + }; + if key.as_ref() == name.as_str() && found.replace(value.clone()).is_some() { + return Err(invalid("duplicate label name")); + } + } + Value::Utf8(found.unwrap_or_else(|| "".into())) + } + LabelIdentity { column } => { + let Value::Map(entries) = &row[*column] else { + return Err(invalid("label identity requires a map")); + }; + let mut labels = std::collections::BTreeMap::new(); + for (key, value) in entries.iter() { + let (Value::Utf8(key), Value::Utf8(value)) = (key, value) else { + return Err(invalid("label identity requires Utf8 entries")); + }; + if labels.insert(key.to_string(), value.to_string()).is_some() { + return Err(invalid("duplicate label name")); + } + } + Value::Utf8( + crate::physical_planner::promql_rows::encode_series_identity(&labels)?.into(), + ) + } Column(i) => row[*i].clone(), Literal { value, .. } => value.clone(), Negate(v) => match v.evaluate(row)? { diff --git a/crates/asap-physical-operators/src/physical_planner/precompute.rs b/crates/asap-physical-operators/src/physical_planner/precompute.rs index 063bd0b8..dd66f62c 100644 --- a/crates/asap-physical-operators/src/physical_planner/precompute.rs +++ b/crates/asap-physical-operators/src/physical_planner/precompute.rs @@ -34,6 +34,83 @@ pub fn population_schema(family: SummaryFamilyType) -> Schema { }) } +/// Raw sample rows at a precompute boundary. `$population` holds the series' +/// complete label set, so it is the complete source identity of per-series +/// summaries; `$timestamp` is the sample time. Rows are what the boundary's +/// source scan selected; the deployment decides which rows and panes they are. +pub fn raw_sample_schema() -> Schema { + let mut schema = (*population_schema(SummaryFamilyType::Plain(DataType::Float64))).clone(); + schema.fields[1].name = "$timestamp".into(); + Arc::new(schema) +} + +pub fn raw_sample_row( + labels: &BTreeMap, + timestamp_ms: i64, + value: f64, +) -> Vec { + use crate::values::Value; + vec![ + Value::Map( + labels + .iter() + .map(|(k, v)| { + ( + Value::Utf8(k.as_str().into()), + Value::Utf8(v.as_str().into()), + ) + }) + .collect::>() + .into(), + ), + Value::Timestamp(timestamp_ms), + Value::Float64(value), + ] +} + +/// Input contract of a precompute boundary: raw sample rows for a raw time +/// series scan, otherwise the stored population of its summary state. +pub fn boundary_schema(node: &PostAsapDagNode) -> Result { + let Payload::Fallback { expression } = &node.payload else { + return source_schema(&node.output_schema); + }; + let scan = match expression { + planner_types::pre_asap::QueryExpr::TimeRange { child, .. } => child.as_ref(), + expression => expression, + }; + if !matches!( + scan, + planner_types::pre_asap::QueryExpr::Scan { + source: planner_types::pre_asap::Source::TimeSeries { .. }, + .. + } + ) { + return source_schema(&node.output_schema); + } + let logical = &node.output_schema; + // Labels may be absent from a series; its label map then omits them. + let valid = logical + .fields + .iter() + .enumerate() + .all(|(i, field)| match &field.dtype { + SummaryFamilyType::Plain(DataType::Timestamp) => { + Some(i) == logical.time_index && !field.nullable + } + SummaryFamilyType::Plain(DataType::Float64) => field.name == "value" && !field.nullable, + SummaryFamilyType::Plain(DataType::Utf8) => true, + _ => false, + }) + && logical.time_index.is_some() + && logical.fields.iter().filter(|f| f.name == "value").count() == 1; + if !valid { + return Err(invalid( + "raw sample boundary requires labels, a timestamp and one Float64 value", + )); + } + Ok(raw_sample_schema()) +} + /// Validate the adapter layout during installed-plan recovery without lowering operators. pub fn source_schema(logical: &SummarySchema) -> Result { let states = logical @@ -132,7 +209,7 @@ pub fn compile( for id in ordered { let node = nodes[&id]; if frontier.contains(&id) { - let schema = source_schema(&node.output_schema)?; + let schema = boundary_schema(node)?; sources.insert(id, InputContract::bounded(schema.clone())); outputs.insert(id, schema); continue; @@ -281,15 +358,39 @@ fn fragment( let [input] = schemas else { return Err(invalid("summary update requires one input")); }; - if update.item.is_some() - || !matches!(grouping, GroupingStrategy::PerSubpopulationInstance) - { + // Item identities resolve against the complete label set of raw + // samples; finalized readouts carry no such identity. + let raw = *input == raw_sample_schema(); + let unit_frequency = crate::capability::is_unit_sample_frequency(update); + let keyed = update.item.is_some() && !unit_frequency; + if (keyed && !raw) || !matches!(grouping, GroupingStrategy::PerSubpopulationInstance) { return Err(invalid( "precompute keyed/shared update needs its dedicated physical candidate", )); } crate::capability::validate_summary_kernel(family, update, grouping) .map_err(Error::Invalid)?; + if raw + && matches!( + update.weight_domain, + planner_types::post_asap::WeightDomain::NonNegative { + proof: planner_types::post_asap::NonNegativeWeightProof::ResetAwareCounterDerivative + } + ) + { + return Err(invalid( + "a counter-derivative weight cannot be read from raw cumulative samples", + )); + } + if keyed + && matches!(family, SummaryFamilyType::Sketch(kind, _) if kind.algorithm() == &planner_types::post_asap::SketchAlgorithm::CmsWithHeap) + && !matches!( + update.weight_domain, + planner_types::post_asap::WeightDomain::NonNegative { .. } + ) + { + return Err(invalid("CMS requires a nonnegative weight contract")); + } let labels = match reduction { PlannerReduction::PerEntity => Expression::Column(0), PlannerReduction::Reduce(keys) => Expression::LabelSet { @@ -302,8 +403,9 @@ fn fragment( .output_schema .fields .get(*key) + // A raw label map omits absent labels. .filter(|field| { - !field.nullable + (raw || !field.nullable) && field.dtype == SummaryFamilyType::Plain(DataType::Utf8) }) .map(|f| f.name.clone()) @@ -316,6 +418,8 @@ fn fragment( }, }; let weight = match &update.weight { + // A unit-frequency summary (HLL) observes the sample value itself. + _ if unit_frequency => Expression::Column(2), SummaryInputExpr::Constant(value) => Expression::Literal { value: crate::values::Value::Float64(*value), dtype: DataType::Float64, @@ -334,20 +438,41 @@ fn fragment( )) } }; - let project = Operator::project( - input.clone(), - vec![ - ("$population".into(), labels), - ("$window_end".into(), Expression::Column(1)), - ("value".into(), Expression::FiniteFloat64(Box::new(weight))), - ], - )? - .with_output_schema(population_schema(SummaryFamilyType::Plain( - DataType::Float64, - )))?; + let mut columns = vec![ + ("$population".into(), labels), + ("$window_end".into(), Expression::Column(1)), + ("value".into(), Expression::FiniteFloat64(Box::new(weight))), + ]; + let mut fields = population_schema(SummaryFamilyType::Plain(DataType::Float64)) + .fields + .clone(); + if keyed { + let mut items = Vec::new(); + raw_items(update.item.as_ref().expect("keyed item"), &mut items)?; + for (index, (expression, dtype)) in items.into_iter().enumerate() { + let name = format!("$item{index}"); + fields.push(planner_types::post_asap::SummaryField { + name: name.clone(), + dtype: SummaryFamilyType::Plain(dtype), + nullable: false, + }); + columns.push((name, expression)); + } + } + let item_columns = (3..fields.len()).collect::>(); + let project = Operator::project(input.clone(), columns)?.with_output_schema( + Arc::new(SummarySchema { + fields, + time_index: Some(1), + }), + )?; let projected = project.schema(); let project = add(vec![0], project)?; - let build = Operator::summary_build(projected, family.clone(), 2, Some(1), vec![0])?; + let build = if keyed { + Operator::keyed_summary_build(projected, family.clone(), 2, item_columns, vec![0])? + } else { + Operator::summary_build(projected, family.clone(), 2, Some(1), vec![0])? + }; let built = build.schema(); let build = add(vec![project], build)?; add( @@ -382,3 +507,42 @@ fn fragment( }; CompiledPhysicalDag::from_operators(sources, operators, vec![root]) } + +/// Resolve keyed item identities over raw sample rows: labels by name (absent +/// labels read as empty, as in PromQL), the sample value, or the canonical +/// encoding of the complete label set. +fn raw_items( + expr: &SummaryInputExpr, + items: &mut Vec<(Expression, DataType)>, +) -> Result<(), Error> { + match expr { + SummaryInputExpr::Column(ColumnRef::SampleValue) => { + items.push((Expression::Column(2), DataType::Float64)) + } + SummaryInputExpr::Column(ColumnRef::Named(name) | ColumnRef::Qualified { name, .. }) => { + items.push(( + Expression::Label { + column: 0, + name: name.clone(), + }, + DataType::Utf8, + )) + } + SummaryInputExpr::EntityIdentity( + planner_types::post_asap::EntityIdentity::PromqlLabelSet { excluding }, + ) if excluding.is_empty() => { + items.push((Expression::LabelIdentity { column: 0 }, DataType::Utf8)) + } + SummaryInputExpr::Tuple(parts) if !parts.is_empty() => { + for part in parts { + raw_items(part, items)?; + } + } + _ => { + return Err(invalid( + "keyed summary item does not resolve over raw samples", + )) + } + } + Ok(()) +} diff --git a/crates/integration-tests/tests/precompute_raw_samples.rs b/crates/integration-tests/tests/precompute_raw_samples.rs new file mode 100644 index 00000000..c2730be6 --- /dev/null +++ b/crates/integration-tests/tests/precompute_raw_samples.rs @@ -0,0 +1,495 @@ +//! Planner-selected summaries over raw samples compile as precompute graphs +//! and produce the same estimates as feeding their kernel sample by sample. +use std::{collections::BTreeMap, collections::BTreeSet, rc::Rc, sync::Arc}; + +use asap_aware_mapping::cost_model::DefaultCostModel; +use asap_aware_mapping::{ + search_workload, Replacement, ReplacementStrategy, ReplacementSubDAG, SketchAlgorithmStrategy, + TargetSubDAG, +}; +use asap_integration_tests::fixtures::lower_promql; +use asap_physical_operators::{ + factory::create_planner_accumulator, + operators::Operator, + physical_planner::{precompute, Source}, + runtime::{Limits, RunContext, Scope}, + summary_kernels::{exact::ExactAccumulator, weighted_frequency::WeightedFrequency}, + values::{Batch, Value}, + AggregateCore, KeyByLabelValues, Statistic, +}; +use asap_types::post_asap::{ + compile_post_asap_dag, EntityIdentity, ExactKind, PostAsapDag, PostAsapOperatorPayload, + SketchAlgorithm, SketchQuery, SummaryFamilyType, SummaryInputExpr, SummaryNode, SummaryUpdate, +}; +use asap_types::pre_asap::{expr_ir::ColumnRef, query_expr::Reduction}; +use asap_types::types::AccuracyTarget; +use futures::{executor::block_on, StreamExt}; + +type Series = BTreeMap; + +fn series(service: &str, instance: &str) -> Series { + [ + ("__name__", "m"), + ("service", service), + ("instance", instance), + ] + .into_iter() + .map(|(k, v)| (k.to_owned(), v.to_owned())) + .collect() +} + +/// Every Planner candidate for `query`: the searched selection plus each +/// summary replacement of the root. +fn candidates(query: &str, accuracy: AccuracyTarget) -> Vec> { + let root = Rc::new(lower_promql(query, accuracy).expect("lowering failed")); + let mut result = SketchAlgorithmStrategy::default_cost_model() + .replacements(&TargetSubDAG::new(&root)) + .into_iter() + .filter_map(|candidate| match candidate { + ReplacementSubDAG { + replacement: Replacement::Summary(node), + .. + } => Some(node), + _ => None, + }) + .collect::>(); + let space = search_workload(vec![("query", root)]); + if let Ok(Some(selected)) = space + .global_selection(&DefaultCostModel) + .assemble_selected_dag(&space.roots[0].1) + { + result.push(selected); + } + result +} + +/// Raw-input summary nodes: `(dag, raw source id, summary id)`. +fn raw_summaries(dag: &PostAsapDag) -> Vec<(u64, u64)> { + dag.nodes + .iter() + .filter(|node| matches!(node.payload, PostAsapOperatorPayload::SummaryAgg { .. })) + .filter_map(|node| { + let inputs = dag + .edges + .iter() + .filter(|edge| edge.consumer == node.id) + .collect::>(); + let [edge] = inputs.as_slice() else { + return None; + }; + let source = dag.nodes.iter().find(|n| n.id == edge.producer)?; + matches!(source.payload, PostAsapOperatorPayload::Fallback { .. }) + .then_some((u64::from(source.id.0), u64::from(node.id.0))) + }) + .collect() +} + +fn samples() -> Vec<(Series, i64, f64)> { + let mut rows = Vec::new(); + for (index, (service, instance)) in [("a", "1"), ("a", "2"), ("b", "1")].iter().enumerate() { + for step in 1..=5i64 { + let value = (index as f64 + 1.0) * step as f64 + (step % 2) as f64; + rows.push((series(service, instance), step * 1000, value)); + } + } + rows +} + +fn execute( + dag: &PostAsapDag, + source: u64, + root: u64, + rows: &[(Series, i64, f64)], +) -> Vec<(Series, Arc)> { + let program = precompute::compile(dag, &[source], &[root]).unwrap_or_else(|error| { + panic!( + "raw summary {root} does not compile: {error}; source {:?}", + dag.nodes + .iter() + .find(|n| u64::from(n.id.0) == source) + .map(|n| (&n.output_schema, &n.payload)) + ); + }); + let program = serde_json::from_slice::< + asap_physical_operators::physical_planner::CompiledPhysicalDag, + >(&serde_json::to_vec(&program).unwrap()) + .unwrap(); + let schema = precompute::raw_sample_schema(); + let batch = Batch::try_new( + schema.clone(), + rows.iter() + .map(|(labels, time, value)| precompute::raw_sample_row(labels, *time, *value)) + .collect(), + ) + .unwrap(); + let sources = BTreeMap::from([( + source, + Box::new(Operator::source(schema, vec![batch]).unwrap()) as Source<'_>, + )]); + let graph = program.instantiate(sources).unwrap(); + let context = RunContext::new( + Scope::Ingestion { + window_start_ms: 0, + window_end_ms: 6000, + revision: 1, + }, + Limits::default(), + ) + .unwrap(); + block_on(async { + let mut stream = graph.execute(program.roots(), context).unwrap().remove(0); + let mut result = Vec::new(); + while let Some(batch) = stream.next().await { + for row in batch.unwrap().rows() { + let [Value::Map(labels), Value::Timestamp(6000), Value::Summary { state, .. }] = + row.as_slice() + else { + panic!("unexpected population row {row:?}"); + }; + let labels = labels + .iter() + .map(|(k, v)| match (k, v) { + (Value::Utf8(k), Value::Utf8(v)) => (k.to_string(), v.to_string()), + _ => panic!("non-label population entry"), + }) + .collect(); + result.push((labels, state.clone())); + } + } + result + }) +} + +fn population(reduction: &Reduction, dag_labels: &[String], labels: &Series) -> Series { + match reduction { + Reduction::PerEntity => labels.clone(), + Reduction::Reduce(_) => labels + .iter() + .filter(|(k, v)| dag_labels.contains(k) && !v.is_empty()) + .map(|(k, v)| (k.clone(), v.clone())) + .collect(), + } +} + +/// Item identity used by a keyed summary: labels by name, the sample value, +/// or the canonical label-set identity. +fn item(expr: &SummaryInputExpr, labels: &Series, value: f64, out: &mut Vec) { + match expr { + SummaryInputExpr::Column(ColumnRef::SampleValue) => out.push(Value::Float64(value)), + SummaryInputExpr::Column(ColumnRef::Named(name)) => out.push(Value::Utf8( + labels.get(name).cloned().unwrap_or_default().into(), + )), + SummaryInputExpr::EntityIdentity(EntityIdentity::PromqlLabelSet { excluding }) + if excluding.is_empty() => + { + out.push(Value::Utf8(serde_json::to_string(labels).unwrap().into())) + } + SummaryInputExpr::Tuple(items) => items.iter().for_each(|i| item(i, labels, value, out)), + other => panic!("unsupported fixture item {other:?}"), + } +} + +fn weight(update: &SummaryUpdate, value: f64) -> f64 { + match update.weight { + SummaryInputExpr::Constant(weight) => weight, + _ => value, + } +} + +/// Estimates that identify a state's content for comparison. +fn readouts(state: &dyn AggregateCore, family: &SummaryFamilyType) -> Vec { + if let Some(exact) = state.as_any().downcast_ref::() { + let SummaryFamilyType::ExactAggregate(kind, _) = family else { + unreachable!() + }; + let statistic = match kind { + ExactKind::Sum => Statistic::Sum, + ExactKind::Count => Statistic::Count, + ExactKind::Min => Statistic::Min, + ExactKind::Max => Statistic::Max, + ExactKind::Rate => Statistic::Rate, + ExactKind::Increase => Statistic::Increase, + other => panic!("unexpected exact kind {other:?}"), + }; + return vec![exact + .readout(statistic, None, None::<&KeyByLabelValues>) + .unwrap() + .unwrap()]; + } + let SummaryFamilyType::Sketch(kind, _) = family else { + panic!("sketch state for exact family") + }; + match kind.algorithm() { + SketchAlgorithm::Kll | SketchAlgorithm::DDSketch => [0.1, 0.5, 0.9] + .into_iter() + .map(|q| state.estimate(&SketchQuery::Quantile { q }).unwrap()) + .collect(), + SketchAlgorithm::Hll => vec![state.estimate(&SketchQuery::Cardinality).unwrap()], + other => panic!("unexpected unkeyed sketch {other:?}"), + } +} + +/// Compile one raw-input summary, execute it over `rows`, and compare each +/// population with its kernel fed sample by sample. Returns the family label, +/// or the family when it has no native state. +fn check( + query: &str, + dag: &PostAsapDag, + source: u64, + root: u64, + rows: &[(Series, i64, f64)], +) -> Result { + let node = dag + .nodes + .iter() + .find(|n| u64::from(n.id.0) == root) + .unwrap(); + let PostAsapOperatorPayload::SummaryAgg { + family, + input, + reduction, + grouping, + } = &node.payload + else { + unreachable!() + }; + let source_node = dag + .nodes + .iter() + .find(|n| u64::from(n.id.0) == source) + .unwrap(); + let keys = match reduction { + Reduction::Reduce(keys) => keys + .keys() + .iter() + .map(|i| source_node.output_schema.fields[*i].name.clone()) + .collect(), + Reduction::PerEntity => vec![], + }; + if asap_physical_operators::capability::validate_native_family(family).is_err() { + // Families without a native state (e.g. plain CMS, UnivMon) + // are outside precompute execution; their compile must fail. + assert!(precompute::compile(dag, &[source], &[root]).is_err()); + return Err(format!("{family:?}")); + } + let actual = execute(dag, source, root, rows); + let label = match family { + SummaryFamilyType::ExactAggregate(kind, _) => format!("{kind:?}"), + SummaryFamilyType::Sketch(kind, _) => format!("{:?}", kind.algorithm()), + other => format!("{other:?}"), + }; + if let SummaryFamilyType::Sketch(kind, _) = family { + if let (Some(keyed), false) = (&input.item, kind.algorithm() == &SketchAlgorithm::Hll) { + // Keyed heaps: every item's estimated weight is its exact + // total at this scale (no collisions in the fixture). + let mut expected = BTreeMap::>::new(); + for (labels, _, value) in rows { + let mut items = Vec::new(); + item(keyed, labels, *value, &mut items); + *expected + .entry(population(reduction, &keys, labels)) + .or_default() + .entry(format!("{items:?}")) + .or_default() += weight(input, *value); + } + assert_eq!(actual.len(), expected.len(), "{query}"); + for (labels, state) in &actual { + let heap = state.as_any().downcast_ref::().unwrap(); + let got = heap + .rows(usize::MAX >> 1) + .into_iter() + .map(|mut row| { + let Some(Value::Float64(score)) = row.pop() else { + panic!("heap score") + }; + (format!("{row:?}"), score) + }) + .collect::>(); + assert_eq!(&got, &expected[labels], "{query}"); + } + return Ok(label); + } + } + let mut expected = + BTreeMap::>::new(); + for (labels, time, value) in rows { + let updater = expected + .entry(population(reduction, &keys, labels)) + .or_insert_with(|| create_planner_accumulator(family, input, grouping).unwrap()); + let unit = input.item.is_some(); + updater.update_single(if unit { *value } else { weight(input, *value) }, *time); + } + assert_eq!(actual.len(), expected.len(), "{query}"); + for (labels, state) in actual { + let reference = expected[&labels].snapshot_accumulator(); + assert_eq!( + readouts(state.as_ref(), family), + readouts(reference.as_ref(), family), + "{query}: {labels:?}" + ); + } + Ok(label) +} + +// Every raw-input summary selected by Planner compiles over raw sample rows, +// and each population's estimates equal feeding its kernel sample by sample. +#[test] +fn raw_sample_summaries_compile_and_match_their_kernels() { + let exact = AccuracyTarget::Exact; + let sketch = AccuracyTarget::Epsilon(0.02); + let queries = [ + ("sum_over_time(m[5m])", &exact), + ("count_over_time(m[5m])", &exact), + ("min_over_time(m[5m])", &exact), + ("max_over_time(m[5m])", &exact), + ("rate(m[5m])", &exact), + ("increase(m[5m])", &exact), + ("sum by (service) (sum_over_time(m[5m]))", &exact), + ("sum by (service) (rate(m[5m]))", &exact), + ("topk(2, sum_over_time(m[5m]))", &exact), + ("quantile_over_time(0.9, m[5m])", &sketch), + ("sum by (service) (quantile_over_time(0.9, m[5m]))", &sketch), + ("quantile by (service) (0.9, m)", &sketch), + ("distinct_over_time(m[5m])", &sketch), + ("count(m)", &sketch), + ("topk(2, m)", &sketch), + ( + "topk(2, sum by (service) (count_over_time(m[5m])))", + &sketch, + ), + ]; + let rows = samples(); + let mut families = BTreeSet::new(); + let mut unsupported = BTreeSet::new(); + for (query, accuracy) in queries { + for candidate in candidates(query, accuracy.clone()) { + let dag = compile_post_asap_dag(&candidate).unwrap(); + for (source, root) in raw_summaries(&dag) { + match check(query, &dag, source, root, &rows) { + Ok(family) => families.insert(family), + Err(family) => unsupported.insert(family), + }; + } + } + } + println!("raw summary families: {families:?}; without native state: {unsupported:?}"); + for family in [ + "Sum", "Count", "Min", "Max", "Rate", "Increase", "Kll", "DDSketch", "Hll", + ] { + assert!( + families.contains(family), + "no {family} fixture: {families:?}" + ); + } +} + +/// Replace the raw summary of `sum by (service) (sum_over_time(m[5m]))` with +/// another update, keeping its raw input and grouping. +fn grouped_raw_summary(family: SummaryFamilyType, input: SummaryUpdate) -> (PostAsapDag, u64, u64) { + let candidate = candidates( + "sum by (service) (sum_over_time(m[5m]))", + AccuracyTarget::Exact, + ) + .pop() + .unwrap(); + let mut dag = compile_post_asap_dag(&candidate).unwrap(); + let (source, root) = raw_summaries(&dag)[0]; + let node = dag + .nodes + .iter_mut() + .find(|n| u64::from(n.id.0) == root) + .unwrap(); + let PostAsapOperatorPayload::SummaryAgg { + family: old, + input: update, + .. + } = &mut node.payload + else { + unreachable!() + }; + for field in &mut node.output_schema.fields { + if field.dtype == *old { + field.dtype = family.clone(); + } + } + *old = family; + *update = input; + let schema = node.output_schema.clone(); + for edge in dag + .edges + .iter_mut() + .filter(|e| u64::from(e.producer.0) == root) + { + edge.intermediate_schema = schema.clone(); + } + (dag, source, root) +} + +// Keyed heaps over raw samples resolve items from the series label set and +// estimate each item's exact total; invalid weight contracts do not compile. +#[test] +fn raw_sample_heaps_resolve_items_from_labels() { + use asap_types::post_asap::{ + GroupingStrategy::PerSubpopulationInstance, NonNegativeWeightProof, SketchKind, + SketchParams, WeightDomain, + }; + let heap = |algorithm, params| { + SummaryFamilyType::Sketch(SketchKind::new(algorithm, params), PerSubpopulationInstance) + }; + let cms = heap( + SketchAlgorithm::CmsWithHeap, + SketchParams::CmsWithHeap { + width: 64, + depth: 3, + heap_size: 8, + }, + ); + let count_sketch = heap( + SketchAlgorithm::CountSketchWithHeap, + SketchParams::CountSketchWithHeap { + width: 64, + depth: 3, + heap_size: 8, + }, + ); + let identity = SummaryUpdate { + item: Some(SummaryInputExpr::EntityIdentity( + EntityIdentity::PromqlLabelSet { excluding: vec![] }, + )), + weight: SummaryInputExpr::Constant(1.0), + weight_domain: WeightDomain::NonNegative { + proof: NonNegativeWeightProof::UnitCount, + }, + }; + let by_instance = SummaryUpdate { + item: Some(SummaryInputExpr::Tuple(vec![SummaryInputExpr::Column( + ColumnRef::Named("instance".into()), + )])), + weight: SummaryInputExpr::Column(ColumnRef::SampleValue), + weight_domain: WeightDomain::UnknownOrSigned, + }; + let rows = samples(); + for (family, input) in [ + (cms.clone(), identity.clone()), + (count_sketch.clone(), identity), + (count_sketch, by_instance.clone()), + ] { + let (dag, source, root) = grouped_raw_summary(family, input); + assert!(check("heap", &dag, source, root, &rows).is_ok()); + } + let signed_cms = grouped_raw_summary(cms.clone(), by_instance.clone()); + assert!(precompute::compile(&signed_cms.0, &[signed_cms.1], &[signed_cms.2]).is_err()); + let derivative = grouped_raw_summary( + cms, + SummaryUpdate { + weight_domain: WeightDomain::NonNegative { + proof: NonNegativeWeightProof::ResetAwareCounterDerivative, + }, + ..by_instance + }, + ); + assert!( + precompute::compile(&derivative.0, &[derivative.1], &[derivative.2]).is_err(), + "raw samples are cumulative counters, not their derivative" + ); +} From c5bc826f2236659acce6009e0c2b89b78604b6d9 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 06:43:12 +0000 Subject: [PATCH 2/3] fix(physical): canonical raw identities and checked raw heap items Raw sample rows drop empty label values so one series has one population. Heap items resolve names against the scan (the value column, the series identity, labels; other scan columns are rejected), and identity items may exclude labels, as `topk by` emits. Unit-frequency HLL updates are raw-only and never apply to keyed families. Tests cover Planner-generated heaps, missing and empty labels, and `without` grouping. Co-Authored-By: Claude Opus 5.5 --- .../src/expressions/mod.rs | 11 +- .../src/physical_planner/precompute.rs | 102 +++++++++--- .../tests/precompute_raw_samples.rs | 151 +++++++++++++++--- 3 files changed, 217 insertions(+), 47 deletions(-) diff --git a/crates/asap-physical-operators/src/expressions/mod.rs b/crates/asap-physical-operators/src/expressions/mod.rs index 1ff7b4db..a16a3bde 100644 --- a/crates/asap-physical-operators/src/expressions/mod.rs +++ b/crates/asap-physical-operators/src/expressions/mod.rs @@ -28,10 +28,12 @@ pub enum Expression { column: usize, name: String, }, - /// Canonical encoding of a complete label map, identical to + /// Canonical encoding of a label map less `excluding`, identical to /// `promql_rows::encode_series_identity`. LabelIdentity { column: usize, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + excluding: Vec, }, Literal { value: Value, @@ -129,7 +131,7 @@ impl Expression { } Ok((expected, false)) } - Label { column, .. } | LabelIdentity { column } => { + Label { column, .. } | LabelIdentity { column, .. } => { if plain(input, *column)? != ( &DataType::Map { @@ -291,7 +293,7 @@ impl Expression { } Value::Utf8(found.unwrap_or_else(|| "".into())) } - LabelIdentity { column } => { + LabelIdentity { column, excluding } => { let Value::Map(entries) = &row[*column] else { return Err(invalid("label identity requires a map")); }; @@ -300,6 +302,9 @@ impl Expression { let (Value::Utf8(key), Value::Utf8(value)) = (key, value) else { return Err(invalid("label identity requires Utf8 entries")); }; + if excluding.iter().any(|label| label.as_str() == key.as_ref()) { + continue; + } if labels.insert(key.to_string(), value.to_string()).is_some() { return Err(invalid("duplicate label name")); } diff --git a/crates/asap-physical-operators/src/physical_planner/precompute.rs b/crates/asap-physical-operators/src/physical_planner/precompute.rs index dd66f62c..38663542 100644 --- a/crates/asap-physical-operators/src/physical_planner/precompute.rs +++ b/crates/asap-physical-operators/src/physical_planner/precompute.rs @@ -1,4 +1,5 @@ //! Compile immutable summary-input computation with explicit population and pane identity. +use super::promql_rows::SERIES_IDENTITY_COLUMN as SERIES_IDENTITY; use super::*; use planner_types::{ post_asap::{ExecutionTiming, GroupingStrategy, SummarySchema}, @@ -36,14 +37,18 @@ pub fn population_schema(family: SummaryFamilyType) -> Schema { /// Raw sample rows at a precompute boundary. `$population` holds the series' /// complete label set, so it is the complete source identity of per-series -/// summaries; `$timestamp` is the sample time. Rows are what the boundary's -/// source scan selected; the deployment decides which rows and panes they are. +/// summaries; `$timestamp` is the sample time and `value` a finite sample +/// (stale markers are not samples). Rows are what the boundary's source scan +/// selected; the deployment decides which rows and panes they are. Build rows +/// with [`raw_sample_row`], which makes the label set canonical. pub fn raw_sample_schema() -> Schema { let mut schema = (*population_schema(SummaryFamilyType::Plain(DataType::Float64))).clone(); schema.fields[1].name = "$timestamp".into(); Arc::new(schema) } +/// A raw sample row whose label set is sorted, unique and omits empty values, +/// so one series always has one population identity. pub fn raw_sample_row( labels: &BTreeMap, timestamp_ms: i64, @@ -54,6 +59,7 @@ pub fn raw_sample_row( Value::Map( labels .iter() + .filter(|(_, v)| !v.is_empty()) .map(|(k, v)| { ( Value::Utf8(k.as_str().into()), @@ -101,6 +107,10 @@ pub fn boundary_schema(node: &PostAsapDagNode) -> Result { SummaryFamilyType::Plain(DataType::Utf8) => true, _ => false, }) + && !logical + .fields + .iter() + .any(|f| f.name.starts_with('$') && f.name != SERIES_IDENTITY) && logical.time_index.is_some() && logical.fields.iter().filter(|f| f.name == "value").count() == 1; if !valid { @@ -361,7 +371,16 @@ fn fragment( // Item identities resolve against the complete label set of raw // samples; finalized readouts carry no such identity. let raw = *input == raw_sample_schema(); - let unit_frequency = crate::capability::is_unit_sample_frequency(update); + // A unit-frequency summary (HLL) observes each raw sample value. + let unit_frequency = raw + && crate::capability::is_unit_sample_frequency(update) + && !matches!(family, SummaryFamilyType::Sketch(kind, _) if matches!( + kind.algorithm(), + planner_types::post_asap::SketchAlgorithm::Cms + | planner_types::post_asap::SketchAlgorithm::CountSketch + | planner_types::post_asap::SketchAlgorithm::CmsWithHeap + | planner_types::post_asap::SketchAlgorithm::CountSketchWithHeap + )); let keyed = update.item.is_some() && !unit_frequency; if (keyed && !raw) || !matches!(grouping, GroupingStrategy::PerSubpopulationInstance) { return Err(invalid( @@ -403,9 +422,11 @@ fn fragment( .output_schema .fields .get(*key) - // A raw label map omits absent labels. + // A raw label map omits absent labels; the + // series identity is not one of its labels. .filter(|field| { (raw || !field.nullable) + && field.name != SERIES_IDENTITY && field.dtype == SummaryFamilyType::Plain(DataType::Utf8) }) .map(|f| f.name.clone()) @@ -418,7 +439,6 @@ fn fragment( }, }; let weight = match &update.weight { - // A unit-frequency summary (HLL) observes the sample value itself. _ if unit_frequency => Expression::Column(2), SummaryInputExpr::Constant(value) => Expression::Literal { value: crate::values::Value::Float64(*value), @@ -448,7 +468,11 @@ fn fragment( .clone(); if keyed { let mut items = Vec::new(); - raw_items(update.item.as_ref().expect("keyed item"), &mut items)?; + raw_items( + update.item.as_ref().expect("keyed item"), + &parents[0].output_schema, + &mut items, + )?; for (index, (expression, dtype)) in items.into_iter().enumerate() { let name = format!("$item{index}"); fields.push(planner_types::post_asap::SummaryField { @@ -508,34 +532,70 @@ fn fragment( CompiledPhysicalDag::from_operators(sources, operators, vec![root]) } -/// Resolve keyed item identities over raw sample rows: labels by name (absent -/// labels read as empty, as in PromQL), the sample value, or the canonical -/// encoding of the complete label set. +/// Resolve keyed item identities over raw sample rows: labels (absent labels +/// read as empty, as in PromQL), the sample value, or the canonical encoding +/// of the label set less excluded labels. fn raw_items( expr: &SummaryInputExpr, + scan: &SummarySchema, items: &mut Vec<(Expression, DataType)>, ) -> Result<(), Error> { + // Open PromQL scans need not list every label, so any name that is not + // another scan column (value, time, series identity) reads as a label. + let label = |column: &ColumnRef| match column { + ColumnRef::Named(name) | ColumnRef::Qualified { name, .. } + if !name.starts_with('$') + && scan.fields.iter().all(|f| { + &f.name != name || f.dtype == SummaryFamilyType::Plain(DataType::Utf8) + }) => + { + Some(name.clone()) + } + _ => None, + }; + let identity = |excluding: Vec| { + ( + Expression::LabelIdentity { + column: 0, + excluding, + }, + DataType::Utf8, + ) + }; match expr { SummaryInputExpr::Column(ColumnRef::SampleValue) => { items.push((Expression::Column(2), DataType::Float64)) } - SummaryInputExpr::Column(ColumnRef::Named(name) | ColumnRef::Qualified { name, .. }) => { - items.push(( - Expression::Label { - column: 0, - name: name.clone(), - }, - DataType::Utf8, - )) + SummaryInputExpr::Column(ColumnRef::Named(name) | ColumnRef::Qualified { name, .. }) + if name == "value" => + { + items.push((Expression::Column(2), DataType::Float64)) + } + SummaryInputExpr::Column(ColumnRef::Named(name) | ColumnRef::Qualified { name, .. }) + if name == SERIES_IDENTITY => + { + items.push(identity(vec![])) } + SummaryInputExpr::Column(column) if label(column).is_some() => items.push(( + Expression::Label { + column: 0, + name: label(column).expect("resolved label"), + }, + DataType::Utf8, + )), SummaryInputExpr::EntityIdentity( planner_types::post_asap::EntityIdentity::PromqlLabelSet { excluding }, - ) if excluding.is_empty() => { - items.push((Expression::LabelIdentity { column: 0 }, DataType::Utf8)) - } + ) => items.push(identity( + excluding + .iter() + .map(|column| { + label(column).ok_or_else(|| invalid("excluded identity label is not a label")) + }) + .collect::>()?, + )), SummaryInputExpr::Tuple(parts) if !parts.is_empty() => { for part in parts { - raw_items(part, items)?; + raw_items(part, scan, items)?; } } _ => { diff --git a/crates/integration-tests/tests/precompute_raw_samples.rs b/crates/integration-tests/tests/precompute_raw_samples.rs index c2730be6..dde29fb1 100644 --- a/crates/integration-tests/tests/precompute_raw_samples.rs +++ b/crates/integration-tests/tests/precompute_raw_samples.rs @@ -27,17 +27,27 @@ use futures::{executor::block_on, StreamExt}; type Series = BTreeMap; -fn series(service: &str, instance: &str) -> Series { +/// A series label set; `None` omits the label. An empty value is present +/// in the input but is not part of the series identity. +fn series(service: Option<&str>, instance: &str) -> Series { [ - ("__name__", "m"), + ("__name__", Some("m")), ("service", service), - ("instance", instance), + ("instance", Some(instance)), ] .into_iter() - .map(|(k, v)| (k.to_owned(), v.to_owned())) + .filter_map(|(k, v)| Some((k.to_owned(), v?.to_owned()))) .collect() } +fn canonical(labels: &Series) -> Series { + labels + .iter() + .filter(|(_, v)| !v.is_empty()) + .map(|(k, v)| (k.clone(), v.clone())) + .collect() +} + /// Every Planner candidate for `query`: the searched selection plus each /// summary replacement of the root. fn candidates(query: &str, accuracy: AccuracyTarget) -> Vec> { @@ -86,10 +96,17 @@ fn raw_summaries(dag: &PostAsapDag) -> Vec<(u64, u64)> { fn samples() -> Vec<(Series, i64, f64)> { let mut rows = Vec::new(); - for (index, (service, instance)) in [("a", "1"), ("a", "2"), ("b", "1")].iter().enumerate() { + let series_set = [ + (Some("a"), "1"), + (Some("a"), "2"), + (Some("b"), "1"), + (None, "3"), + (Some("b"), ""), + ]; + for (index, (service, instance)) in series_set.iter().enumerate() { for step in 1..=5i64 { let value = (index as f64 + 1.0) * step as f64 + (step % 2) as f64; - rows.push((series(service, instance), step * 1000, value)); + rows.push((series(*service, instance), step * 1000, value)); } } rows @@ -160,13 +177,21 @@ fn execute( }) } +/// PromQL grouping of a canonical label set: `by` keeps the named labels, +/// `without` drops them and `__name__`. fn population(reduction: &Reduction, dag_labels: &[String], labels: &Series) -> Series { + let labels = canonical(labels); match reduction { - Reduction::PerEntity => labels.clone(), - Reduction::Reduce(_) => labels - .iter() - .filter(|(k, v)| dag_labels.contains(k) && !v.is_empty()) - .map(|(k, v)| (k.clone(), v.clone())) + Reduction::PerEntity => labels, + Reduction::Reduce(keys) => labels + .into_iter() + .filter(|(k, _)| { + if keys.is_without() { + k != "__name__" && !dag_labels.contains(k) + } else { + dag_labels.contains(k) + } + }) .collect(), } } @@ -176,13 +201,27 @@ fn population(reduction: &Reduction, dag_labels: &[String], labels: &Series) -> fn item(expr: &SummaryInputExpr, labels: &Series, value: f64, out: &mut Vec) { match expr { SummaryInputExpr::Column(ColumnRef::SampleValue) => out.push(Value::Float64(value)), + SummaryInputExpr::Column(ColumnRef::Named(name)) if name == "value" => { + out.push(Value::Float64(value)) + } SummaryInputExpr::Column(ColumnRef::Named(name)) => out.push(Value::Utf8( - labels.get(name).cloned().unwrap_or_default().into(), + canonical(labels) + .get(name) + .cloned() + .unwrap_or_default() + .into(), )), - SummaryInputExpr::EntityIdentity(EntityIdentity::PromqlLabelSet { excluding }) - if excluding.is_empty() => - { - out.push(Value::Utf8(serde_json::to_string(labels).unwrap().into())) + SummaryInputExpr::EntityIdentity(EntityIdentity::PromqlLabelSet { excluding }) => { + let mut identity = canonical(labels); + for column in excluding { + let ColumnRef::Named(name) = column else { + panic!("unsupported fixture exclusion {column:?}") + }; + identity.remove(name); + } + out.push(Value::Utf8( + serde_json::to_string(&identity).unwrap().into(), + )) } SummaryInputExpr::Tuple(items) => items.iter().for_each(|i| item(i, labels, value, out)), other => panic!("unsupported fixture item {other:?}"), @@ -357,24 +396,41 @@ fn raw_sample_summaries_compile_and_match_their_kernels() { "topk(2, sum by (service) (count_over_time(m[5m])))", &sketch, ), + ("topk(2, sum_over_time(m[5m]))", &sketch), + ("topk by (service) (2, sum_over_time(m[5m]))", &sketch), ]; let rows = samples(); let mut families = BTreeSet::new(); let mut unsupported = BTreeSet::new(); + let mut checked = BTreeMap::new(); for (query, accuracy) in queries { for candidate in candidates(query, accuracy.clone()) { let dag = compile_post_asap_dag(&candidate).unwrap(); for (source, root) in raw_summaries(&dag) { match check(query, &dag, source, root, &rows) { - Ok(family) => families.insert(family), - Err(family) => unsupported.insert(family), - }; + Ok(family) => { + families.insert(family); + *checked.entry(query).or_insert(0) += 1; + } + Err(family) => { + unsupported.insert(family); + } + } } } } - println!("raw summary families: {families:?}; without native state: {unsupported:?}"); + println!("checked {checked:?}; families {families:?}; without native state {unsupported:?}"); for family in [ - "Sum", "Count", "Min", "Max", "Rate", "Increase", "Kll", "DDSketch", "Hll", + "Sum", + "Count", + "Min", + "Max", + "Rate", + "Increase", + "Kll", + "DDSketch", + "Hll", + "CountSketchWithHeap", ] { assert!( families.contains(family), @@ -384,7 +440,7 @@ fn raw_sample_summaries_compile_and_match_their_kernels() { } /// Replace the raw summary of `sum by (service) (sum_over_time(m[5m]))` with -/// another update, keeping its raw input and grouping. +/// another update, keeping its raw input and reduction. fn grouped_raw_summary(family: SummaryFamilyType, input: SummaryUpdate) -> (PostAsapDag, u64, u64) { let candidate = candidates( "sum by (service) (sum_over_time(m[5m]))", @@ -485,11 +541,60 @@ fn raw_sample_heaps_resolve_items_from_labels() { weight_domain: WeightDomain::NonNegative { proof: NonNegativeWeightProof::ResetAwareCounterDerivative, }, - ..by_instance + ..by_instance.clone() }, ); assert!( precompute::compile(&derivative.0, &[derivative.1], &[derivative.2]).is_err(), "raw samples are cumulative counters, not their derivative" ); + // A scan's time column is not a label; it cannot silently read as empty. + let time_item = grouped_raw_summary( + heap( + SketchAlgorithm::CountSketchWithHeap, + SketchParams::CountSketchWithHeap { + width: 64, + depth: 3, + heap_size: 8, + }, + ), + SummaryUpdate { + item: Some(SummaryInputExpr::Column(ColumnRef::Named("ts".into()))), + ..by_instance + }, + ); + assert!(precompute::compile(&time_item.0, &[time_item.1], &[time_item.2]).is_err()); +} + +// `without` grouping over raw samples drops the listed labels and `__name__`. +#[test] +fn raw_sample_without_grouping_drops_labels_and_name() { + use asap_types::pre_asap::query_expr::GroupKeys; + let family = + SummaryFamilyType::ExactAggregate(ExactKind::Sum, asap_types::post_asap::ExactParams::Sum); + let (mut dag, source, root) = + grouped_raw_summary(family, SummaryUpdate::column(ColumnRef::SampleValue)); + let service = dag + .nodes + .iter() + .find(|n| u64::from(n.id.0) == source) + .unwrap() + .output_schema + .fields + .iter() + .position(|f| f.name == "service") + .unwrap(); + let node = dag + .nodes + .iter_mut() + .find(|n| u64::from(n.id.0) == root) + .unwrap(); + let PostAsapOperatorPayload::SummaryAgg { reduction, .. } = &mut node.payload else { + unreachable!() + }; + *reduction = Reduction::Reduce(GroupKeys::without(vec![service])); + assert_eq!( + check("without", &dag, source, root, &samples()), + Ok("Sum".into()) + ); } From f574c00aeddcb88ddd8c1a79fa9a40252dbdf421 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 06:45:57 +0000 Subject: [PATCH 3/3] fix(physical): keep raw unit-frequency updates to unkeyed sketches State the canonical label-set obligation and assert per-query coverage. Co-Authored-By: Claude Opus 5.5 --- .../src/physical_planner/precompute.rs | 7 ++++--- .../integration-tests/tests/precompute_raw_samples.rs | 10 ++++++++++ 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/crates/asap-physical-operators/src/physical_planner/precompute.rs b/crates/asap-physical-operators/src/physical_planner/precompute.rs index 38663542..4fa40c68 100644 --- a/crates/asap-physical-operators/src/physical_planner/precompute.rs +++ b/crates/asap-physical-operators/src/physical_planner/precompute.rs @@ -39,8 +39,9 @@ pub fn population_schema(family: SummaryFamilyType) -> Schema { /// complete label set, so it is the complete source identity of per-series /// summaries; `$timestamp` is the sample time and `value` a finite sample /// (stale markers are not samples). Rows are what the boundary's source scan -/// selected; the deployment decides which rows and panes they are. Build rows -/// with [`raw_sample_row`], which makes the label set canonical. +/// selected; the deployment decides which rows and panes they are. Label sets +/// must be canonical (sorted, unique, no empty values), since they are the +/// population identity: build rows with [`raw_sample_row`]. pub fn raw_sample_schema() -> Schema { let mut schema = (*population_schema(SummaryFamilyType::Plain(DataType::Float64))).clone(); schema.fields[1].name = "$timestamp".into(); @@ -374,7 +375,7 @@ fn fragment( // A unit-frequency summary (HLL) observes each raw sample value. let unit_frequency = raw && crate::capability::is_unit_sample_frequency(update) - && !matches!(family, SummaryFamilyType::Sketch(kind, _) if matches!( + && matches!(family, SummaryFamilyType::Sketch(kind, _) if !matches!( kind.algorithm(), planner_types::post_asap::SketchAlgorithm::Cms | planner_types::post_asap::SketchAlgorithm::CountSketch diff --git a/crates/integration-tests/tests/precompute_raw_samples.rs b/crates/integration-tests/tests/precompute_raw_samples.rs index dde29fb1..ade52ec8 100644 --- a/crates/integration-tests/tests/precompute_raw_samples.rs +++ b/crates/integration-tests/tests/precompute_raw_samples.rs @@ -420,6 +420,16 @@ fn raw_sample_summaries_compile_and_match_their_kernels() { } } println!("checked {checked:?}; families {families:?}; without native state {unsupported:?}"); + for query in [ + "topk by (service) (2, sum_over_time(m[5m]))", + "quantile by (service) (0.9, m)", + "distinct_over_time(m[5m])", + ] { + assert!( + checked.contains_key(query), + "{query} has no checked raw summary" + ); + } for family in [ "Sum", "Count",