From ae5f5177b296866787290dad643e5277dc8d85a4 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Thu, 1 Oct 2026 22:33:13 +0000 Subject: [PATCH] g6 local --- .../src/query_physical_lowering.rs | 2 +- crates/asap-aware-mapping/src/replacement.rs | 34 +- .../src/expressions/arithmetic.rs | 33 +- .../src/expressions/mod.rs | 8 + .../src/operators/aggregate/mod.rs | 70 +- .../src/operators/aggregate/temporal.rs | 81 ++- .../src/operators/mod.rs | 15 +- .../src/operators/series_labels.rs | 501 ++++++++++++-- .../src/operators/unchecked.rs | 12 +- .../src/operators/vector_binary.rs | 8 +- .../src/operators/vector_window.rs | 5 +- .../src/physical_planner/mod.rs | 64 +- .../src/physical_planner/promql_fallback.rs | 152 ++--- .../src/physical_planner/row_values.rs | 162 +---- .../tests/deployment_computation.rs | 330 ++++++++- .../tests/promql_binary.rs | 52 +- .../tests/promql_fallback.rs | 639 +++++++++++++++++- crates/frontend-promql/src/promql.rs | 90 ++- .../awesome_prometheus_alerts.rs | 2 +- .../tests/observability/promql_corpus.rs | 7 +- .../tests/promql_conformance.rs | 4 +- .../frontend-promql/tests/promql_lowering.rs | 129 +++- crates/frontend-sql/src/sql/mod.rs | 6 +- crates/frontend-sql/tests/sql_lowering.rs | 3 +- .../tests/promql_to_post_asap.rs | 27 + crates/types/src/pre_asap/agg_intent.rs | 6 + crates/types/src/pre_asap/query_expr.rs | 43 +- crates/types/src/pre_asap/resolve.rs | 5 +- crates/types/src/pre_asap/schema.rs | 16 +- .../develop_docs/physical-compile-coverage.md | 131 +++- docs/develop_docs/pre-asap-ir.md | 8 +- 31 files changed, 2131 insertions(+), 514 deletions(-) diff --git a/crates/asap-aware-mapping/src/query_physical_lowering.rs b/crates/asap-aware-mapping/src/query_physical_lowering.rs index bd2aa675..1cc2b1ca 100644 --- a/crates/asap-aware-mapping/src/query_physical_lowering.rs +++ b/crates/asap-aware-mapping/src/query_physical_lowering.rs @@ -1029,7 +1029,7 @@ fn promql_binary_operation( BinaryOpKind::Set(PromQLVectorSetOpKind::And) => PromqlBinaryOperation::And, BinaryOpKind::Set(PromQLVectorSetOpKind::Or) => PromqlBinaryOperation::Or, BinaryOpKind::Set(PromQLVectorSetOpKind::Unless) => PromqlBinaryOperation::Unless, - BinaryOpKind::Arithmetic(_) | BinaryOpKind::Compare(_) => { + BinaryOpKind::Arithmetic(_) | BinaryOpKind::Compare(_) | BinaryOpKind::CompareBool(_) => { PromqlBinaryOperation::ArithmeticOrComparison } } diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index 8c826b37..3b8553e5 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -2804,13 +2804,12 @@ fn construct_summary_agg( } else { reduction.clone() }; - let per_series = matches!(reduction, Reduction::PerEntity); - let by: Vec = reduction - .group_keys() - .map(|g| g.to_vec()) - .unwrap_or_default(); let out_schema = node.output_schema()?; - let state_idx = summary_col_index(&out_schema, &by, per_series); + let measures = match node { + QueryExpr::Aggregate { measures, .. } => measures.len(), + _ => 1, + }; + let state_idx = summary_col_index(&out_schema, reduction, measures); let readout_schema = if keyed_heap && matches!(node, QueryExpr::Aggregate { child, .. } if is_snapshot_weighted_topk(intent, child)) @@ -3558,16 +3557,19 @@ fn compose_guarantee( /// cross-series output is `by ++ [agg]` (the column after the keys); /// a per-series reduction keeps every label and replaces the sample value /// (named `value` — mirror `per_series_reduction_schema`'s fallback). -/// `per_series` is the caller's already-read `Reduction` (issue #165) — -/// this never re-derives it, so it can't disagree with the caller. -fn summary_col_index(out_schema: &Schema, by: &[usize], per_series: bool) -> usize { - if per_series { - out_schema +/// `reduction` is the caller's already-read `Reduction` (issue #165). +/// `without` output is `kept labels ++ measures`, and its keys are the +/// *excluded* labels, so the state column follows the kept labels instead. +fn summary_col_index(out_schema: &Schema, reduction: &Reduction, measures: usize) -> usize { + match reduction { + Reduction::PerEntity => out_schema .column_id("value") .or_else(|| (0..out_schema.columns.len()).find(|&i| Some(i) != out_schema.time_index)) - .unwrap_or(0) - } else { - by.len() + .unwrap_or(0), + Reduction::Reduce(keys) if keys.is_without() => { + out_schema.columns.len().saturating_sub(measures) + } + Reduction::Reduce(keys) => keys.len(), } } @@ -7387,7 +7389,7 @@ mod tests { Pass, ), // classic-bucket histogram_quantile is not re-sketchable (#79) - (A::HistogramQuantile { q: 0.99 }, Pass), + (A::HistogramQuantile { q: 0.99, le: 0 }, Pass), // counter-derivative / range-vector functions (#44) (A::Changes, Pass), (A::Delta, Pass), @@ -10120,7 +10122,7 @@ mod tests { // decree. All three stay whole logical subtrees. for intent in [ AggIntent::Avg { col: None }, - AggIntent::HistogramQuantile { q: 0.99 }, + AggIntent::HistogramQuantile { q: 0.99, le: 0 }, AggIntent::Quantile { col: None, q: 0.99, diff --git a/crates/asap-physical-operators/src/expressions/arithmetic.rs b/crates/asap-physical-operators/src/expressions/arithmetic.rs index 30e277d4..e0766763 100644 --- a/crates/asap-physical-operators/src/expressions/arithmetic.rs +++ b/crates/asap-physical-operators/src/expressions/arithmetic.rs @@ -25,7 +25,7 @@ pub fn evaluate_binary( right: f64, ) -> Result { use crate::{values::Value, Error}; - use planner_types::pre_asap::{ArithmeticOpKind, BinaryOpKind, CompareOpKind}; + use planner_types::pre_asap::{ArithmeticOpKind, BinaryOpKind}; let invalid = || Error::Invalid("unsupported binary operation or invalid checked-division domain".into()); if operator.vector_match.is_some() { @@ -49,15 +49,28 @@ pub fn evaluate_binary( BinaryOpKind::Arithmetic(ref op) => { Value::Float64(evaluate_float64_arithmetic(op, left, right)) } - BinaryOpKind::Compare(ref op) => Value::Bool(match op { - CompareOpKind::Eq => left == right, - CompareOpKind::Ne => left != right, - CompareOpKind::Lt => left < right, - CompareOpKind::Le => left <= right, - CompareOpKind::Gt => left > right, - CompareOpKind::Ge => left >= right, - _ => return Err(invalid()), - }), + BinaryOpKind::Compare(ref op) => Value::Bool(compare(op, left, right).ok_or_else(invalid)?), + BinaryOpKind::CompareBool(ref op) => { + Value::Float64(if compare(op, left, right).ok_or_else(invalid)? { + 1. + } else { + 0. + }) + } _ => return Err(invalid()), }) } + +/// IEEE comparison, as Go's: NaN is unequal to everything, itself included. +fn compare(op: &planner_types::pre_asap::CompareOpKind, left: f64, right: f64) -> Option { + use planner_types::pre_asap::CompareOpKind; + Some(match op { + CompareOpKind::Eq => left == right, + CompareOpKind::Ne => left != right, + CompareOpKind::Lt => left < right, + CompareOpKind::Le => left <= right, + CompareOpKind::Gt => left > right, + CompareOpKind::Ge => left >= right, + _ => return None, + }) +} diff --git a/crates/asap-physical-operators/src/expressions/mod.rs b/crates/asap-physical-operators/src/expressions/mod.rs index a16a3bde..627fd931 100644 --- a/crates/asap-physical-operators/src/expressions/mod.rs +++ b/crates/asap-physical-operators/src/expressions/mod.rs @@ -87,6 +87,14 @@ impl Expression { | CompareOpKind::Gt | CompareOpKind::Ge, ) => DataType::Bool, + BinaryOpKind::CompareBool( + CompareOpKind::Eq + | CompareOpKind::Ne + | CompareOpKind::Lt + | CompareOpKind::Le + | CompareOpKind::Gt + | CompareOpKind::Ge, + ) => DataType::Float64, _ => return Err(invalid("unsupported binary operation")), }; Ok((dtype, n || m)) diff --git a/crates/asap-physical-operators/src/operators/aggregate/mod.rs b/crates/asap-physical-operators/src/operators/aggregate/mod.rs index d929fc07..e7e9ed22 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/mod.rs +++ b/crates/asap-physical-operators/src/operators/aggregate/mod.rs @@ -302,9 +302,9 @@ async fn reduce_one( } return Ok(best.cloned().unwrap_or(Value::Null)); } - let mut count = 0usize; let dtype = plain(input, column)?.0; if dtype == &DataType::Int64 { + let mut count = 0usize; let mut sum = 0i128; for v in values { work.checkpoint().await?; @@ -324,22 +324,80 @@ async fn reduce_one( )) }; } - let mut sum = -0.0; + let mut floats = Vec::with_capacity(rows.len()); for v in values { work.checkpoint().await?; let Value::Float64(v) = v else { return Err(invalid("floating aggregate value required")); }; - sum += v; - count += 1; + floats.push(*v); } Ok(Value::Float64(if matches!(measure, Reduction::Avg(_)) { - sum / count as f64 + promql_avg(&floats) } else { - sum + // Prometheus starts `sum` from the first value; adding it to 0 gives + // the same sum and compensation. + promql_sum(0., &floats) })) } +/// Prometheus `kahansum.Inc`: Kahan-Neumaier compensated addition, with the +/// compensation cleared once the sum is infinite. +fn kahan_inc(inc: f64, sum: f64, c: f64) -> (f64, f64) { + let t = sum + inc; + let c = if t.is_infinite() { + 0. + } else if sum.abs() >= inc.abs() { + c + ((sum - t) + inc) + } else { + c + ((inc - t) + sum) + }; + (t, c) +} + +/// Prometheus `sum` and `sum_over_time` from `start`. +pub(in crate::operators) fn promql_sum(start: f64, values: &[f64]) -> f64 { + let (sum, c) = values + .iter() + .fold((start, 0.), |(sum, c), &v| kahan_inc(v, sum, c)); + if sum.is_infinite() { + sum + } else { + sum + c + } +} + +/// Prometheus `avg` and `avg_over_time`: a compensated sum divided by the +/// count until the running sum would overflow, then an incremental mean. +/// NaN for no values. +pub(in crate::operators) fn promql_avg(values: &[f64]) -> f64 { + let Some((&first, rest)) = values.split_first() else { + return f64::NAN; + }; + let (mut sum, mut c, mut mean, mut incremental) = (first, 0., 0., false); + for (i, &v) in rest.iter().enumerate() { + let count = (i + 2) as f64; + if !incremental { + let (next, next_c) = kahan_inc(v, sum, c); + if !next.is_infinite() { + (sum, c) = (next, next_c); + continue; + } + incremental = true; + mean = sum / (count - 1.); + c /= count - 1.; + } + let q = (count - 1.) / count; + (mean, c) = kahan_inc(v / count, q * mean, q * c); + } + let count = values.len() as f64; + if incremental { + mean + c + } else { + sum / count + c / count + } +} + #[cfg(test)] mod tests { // Quantile follows Prometheus: interpolate ranks, NaN when empty, ±Inf outside [0, 1]. diff --git a/crates/asap-physical-operators/src/operators/aggregate/temporal.rs b/crates/asap-physical-operators/src/operators/aggregate/temporal.rs index 3ac0a3d1..98921852 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/temporal.rs +++ b/crates/asap-physical-operators/src/operators/aggregate/temporal.rs @@ -52,7 +52,7 @@ pub(in crate::operators) async fn reduce( let mut output = Vec::new(); for (_, (mut keys, buckets, points)) in grouped { work.checkpoint().await?; - let result = if let AggIntent::HistogramQuantile { q } = intent { + let result = if let AggIntent::HistogramQuantile { q, .. } = intent { Some(Value::Float64(bucket_quantile(*q, buckets, context).await?)) } else { let points = cooperative_sort(points, |a, b| a.0.cmp(&b.0), context).await?; @@ -92,10 +92,14 @@ pub(in crate::operators) fn window_value( AggIntent::Count { .. } => Some(Value::Int64( i64::try_from(points.len()).map_err(|_| Error::Invalid("count overflow".into()))?, )), - AggIntent::Sum { .. } => Some(Value::Float64(points.iter().map(|p| p.1).sum())), - AggIntent::Avg { .. } => Some(Value::Float64( - points.iter().map(|p| p.1).sum::() / points.len() as f64, - )), + AggIntent::Sum { .. } | AggIntent::Avg { .. } => { + let values = points.iter().map(|p| p.1).collect::>(); + Some(Value::Float64(if matches!(intent, AggIntent::Sum { .. }) { + super::promql_sum(0., &values) + } else { + super::promql_avg(&values) + })) + } AggIntent::Min { .. } => Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| { if a.is_nan() || p.1 < a { p.1 @@ -176,7 +180,7 @@ fn rate(points: &[(i64, f64)], start: i64, end: i64, counter: bool) -> Option, context: &RunContext, @@ -211,7 +215,7 @@ async fn bucket_quantile( let mut prev = buckets[0].1; for p in buckets.iter_mut().skip(1) { work.checkpoint().await?; - if p.1 < prev || (p.1 - prev).abs() <= 1e-12 * (p.1.abs() + prev.abs()) { + if p.1 < prev || almost_equal(prev, p.1) { p.1 = prev; } prev = p.1; @@ -221,7 +225,13 @@ async fn bucket_quantile( return Ok(f64::NAN); } let rank = q * count; - let idx = buckets[..buckets.len() - 1].partition_point(|p| p.1 < rank); + let idx = buckets[..buckets.len() - 1].partition_point(|p| { + // Go searches for `count >= rank`; a NaN comparison is never a match. + !matches!( + p.1.partial_cmp(&rank), + Some(std::cmp::Ordering::Greater | std::cmp::Ordering::Equal) + ) + }); if idx == buckets.len() - 1 { return Ok(buckets[idx - 1].0); } @@ -230,7 +240,21 @@ async fn bucket_quantile( } let (start, base) = if idx == 0 { (0., 0.) } else { buckets[idx - 1] }; let (end, upper) = buckets[idx]; - Ok(start + (end - start) * (rank - base) / (upper - base)) + Ok(start + (end - start) * ((rank - base) / (upper - base))) +} + +/// Prometheus `almost.Equal` with its bucket tolerance of 1e-12. +fn almost_equal(a: f64, b: f64) -> bool { + const EPSILON: f64 = 1e-12; + if a == b || (a.is_nan() && b.is_nan()) { + return true; + } + let sum = a.abs() + b.abs(); + let diff = (a - b).abs(); + if a == 0. || b == 0. || sum < f64::MIN_POSITIVE { + return diff < EPSILON * f64::MIN_POSITIVE; + } + diff / sum.min(f64::MAX) < EPSILON } #[cfg(test)] @@ -335,4 +359,43 @@ mod tests { assert!(bucket_quantile(0.5, vec![(1., 2.), (2., 4.)]).is_nan()); assert_eq!(bucket_quantile(-0.1, vec![]), f64::NEG_INFINITY); } + + // Matches Prometheus bit for bit: interpolation divides before scaling, + // and an infinite count is not "almost equal" to a finite one. + #[test] + fn histogram_matches_prometheus_arithmetic() { + let context = RunContext::new( + Scope::Query { + evaluation_time_ms: 0, + revision: 0, + }, + Limits::default(), + ) + .unwrap(); + let bucket_quantile = |q, buckets| { + futures::executor::block_on(super::bucket_quantile(q, buckets, &context)).unwrap() + }; + // 0.5 + (0.8 - 0.5) * ((19.95 - 10) / 11) in Go. + assert_eq!( + bucket_quantile(0.95, vec![(0.5, 10.), (0.8, 21.), (f64::INFINITY, 21.)]), + 0.7713636363636364 + ); + // Rank ∞ lies past every finite bucket. + assert_eq!( + bucket_quantile( + 0.5, + vec![(1., 1.), (2., 2.), (f64::INFINITY, f64::INFINITY)] + ), + 2. + ); + // A NaN rank (0 · ∞, or a NaN total) finds no bucket, as Go's sort.Search. + assert_eq!( + bucket_quantile(0., vec![(1., 1.), (2., 2.), (f64::INFINITY, f64::INFINITY)]), + 2. + ); + assert_eq!( + bucket_quantile(0.5, vec![(1., 1.), (2., 2.), (f64::INFINITY, f64::NAN)]), + 2. + ); + } } diff --git a/crates/asap-physical-operators/src/operators/mod.rs b/crates/asap-physical-operators/src/operators/mod.rs index 196e9e8a..6f44207e 100644 --- a/crates/asap-physical-operators/src/operators/mod.rs +++ b/crates/asap-physical-operators/src/operators/mod.rs @@ -81,9 +81,16 @@ enum Kind { SeriesLabels { kind: planner_types::pre_asap::VectorMatchKind, labels: Vec, + unique: bool, }, SeriesBinary { operator: planner_types::post_asap::BinaryOperator, + scalars: [bool; 2], + }, + SeriesHistogramQuantile { + /// `f64` bits: JSON cannot encode the NaN and infinite quantiles. + quantile: u64, + le: usize, }, Project(Vec), Filter(Expression), @@ -259,6 +266,7 @@ impl PhysicalOperator for Operator { | Kind::SeriesWindow { .. } | Kind::SeriesLabels { .. } | Kind::SeriesBinary { .. } + | Kind::SeriesHistogramQuantile { .. } | Kind::Aggregate { .. } | Kind::Window { .. } | Kind::Join { .. } @@ -304,6 +312,7 @@ impl PhysicalOperator for Operator { Kind::SeriesWindow { .. } => "SeriesWindow", Kind::SeriesLabels { .. } => "SeriesLabels", Kind::SeriesBinary { .. } => "SeriesBinary", + Kind::SeriesHistogramQuantile { .. } => "SeriesHistogramQuantile", Kind::Project(_) => "Project", Kind::Filter(_) => "Filter", Kind::Limit { .. } => "Limit", @@ -350,9 +359,9 @@ impl PhysicalOperator for Operator { Kind::CurrentSeries { .. } => current_series::execute(self, inputs, context), Kind::ScopeTimestamp { .. } => scope_timestamp::execute(self, inputs, context), Kind::SeriesWindow { .. } => series_window::execute(self, inputs, context), - Kind::SeriesLabels { .. } | Kind::SeriesBinary { .. } => { - series_labels::execute(self, inputs, context) - } + Kind::SeriesLabels { .. } + | Kind::SeriesBinary { .. } + | Kind::SeriesHistogramQuantile { .. } => series_labels::execute(self, inputs, context), Kind::Filter(_) => filter::execute(self, inputs, context), Kind::Limit { .. } => limit::execute(self, inputs, context), Kind::Sort { .. } => sort::execute(self, inputs, context), diff --git a/crates/asap-physical-operators/src/operators/series_labels.rs b/crates/asap-physical-operators/src/operators/series_labels.rs index be248319..3ed83491 100644 --- a/crates/asap-physical-operators/src/operators/series_labels.rs +++ b/crates/asap-physical-operators/src/operators/series_labels.rs @@ -1,5 +1,5 @@ -//! PromQL label-set rewriting and one-to-one vector matching over rows that -//! carry a series identity or plain label columns. +//! PromQL label-set rewriting and binary operators over rows that carry a +//! series identity or plain label columns. use super::*; use planner_types::{ post_asap::BinaryOperator, @@ -87,40 +87,393 @@ impl Operator { ) -> Result { layout(&input)?; Ok(Self { - kind: Kind::SeriesLabels { kind, labels }, + kind: Kind::SeriesLabels { + kind, + labels, + unique: false, + }, inputs: vec![input.clone()], output: input, }) } - /// PromQL one-to-one arithmetic between rows with equal label sets. The - /// result keeps the left row, without the metric name. + /// Drop the metric name from each row's label set, as PromQL arithmetic + /// with a literal does. Unlike matching labels, the result is itself a + /// vector, so two rows that become equal are an error, as in Prometheus. + pub fn series_without_name(input: Schema) -> Result { + layout(&input)?; + Ok(Self { + kind: Kind::SeriesLabels { + kind: VectorMatchKind::Ignoring, + labels: vec![], + unique: true, + }, + inputs: vec![input.clone()], + output: input, + }) + } + + /// PromQL `histogram_quantile` over classic buckets: one histogram per + /// label set other than the `le` column's label. The result drops `le`, + /// `__name__` and the time column; its value is the quantile. + pub fn series_histogram_quantile(input: Schema, q: f64, le: usize) -> Result { + let layout = layout(&input)?; + if !layout.labels.contains(&le) { + return Err(invalid("histogram bucket bound must be a label column")); + } + let mut fields = (0..input.fields.len()) + .filter(|&i| i != le && i != layout.value && Some(i) != input.time_index) + .map(|i| input.fields[i].clone()) + .collect::>(); + fields.push(input.fields[layout.value].clone()); + Ok(Self { + kind: Kind::SeriesHistogramQuantile { + quantile: q.to_bits(), + le, + }, + inputs: vec![input], + output: schema(fields), + }) + } + + /// A PromQL binary operator with Prometheus' matching, metric-name, and + /// duplicate rules. `scalars` marks the operands that are one-row PromQL + /// scalars, such as a literal or `scalar(x)`, rather than vectors. + /// Vectors match by `operator.vector_match`; `None` matches all labels + /// but the name. The result has the vector operand's schema, the left one + /// between vectors; label columns without a series identity must hold + /// every label the result can take from the right side. pub fn series_binary( left: Schema, right: Schema, operator: BinaryOperator, + scalars: [bool; 2], ) -> Result { - layout(&left)?; - layout(&right)?; - if !matches!(operator.kind, BinaryOpKind::Arithmetic(_)) || operator.vector_match.is_some() + use planner_types::pre_asap::{CompareOpKind::*, GroupSide, PromQLVectorSetOpKind}; + let vectors = scalars == [false, false]; + let valid = match &operator.kind { + BinaryOpKind::Arithmetic(_) => true, + BinaryOpKind::CompareBool(op) => matches!(op, Eq | Ne | Lt | Le | Gt | Ge), + // Prometheus requires `bool` between two scalars. + BinaryOpKind::Compare(op) => { + matches!(op, Eq | Ne | Lt | Le | Gt | Ge) && scalars != [true, true] + } + BinaryOpKind::Set(_) => vectors, + }; + let grouping = operator + .vector_match + .as_ref() + .and_then(|m| m.grouping.as_ref()); + // The frontend records a `bool` modifier as default matching. + let default_match = operator.vector_match.as_ref().is_none_or(|m| { + m.kind == VectorMatchKind::Ignoring && m.labels.is_empty() && m.grouping.is_none() + }); + if !valid + || (!vectors && !default_match) + || (grouping.is_some() && matches!(operator.kind, BinaryOpKind::Set(_))) { - return Err(invalid("series binary requires unmatched arithmetic")); + return Err(invalid("unsupported PromQL binary operation")); + } + for (schema, scalar) in [(&left, scalars[0]), (&right, scalars[1])] { + if scalar { + if schema.fields.len() != 1 || plain(schema, 0)? != (&DataType::Float64, false) { + return Err(invalid("PromQL scalar operand must be one Float64 value")); + } + } else { + layout(schema)?; + } } + if vectors { + let (l, r) = (layout(&left)?, layout(&right)?); + let names = |layout: &Layout, schema: &Schema| { + layout + .labels + .iter() + .map(|&i| schema.fields[i].name.clone()) + .collect::>() + }; + let right_rows = matches!(operator.kind, BinaryOpKind::Set(PromQLVectorSetOpKind::Or)) + || matches!(grouping, Some(g) if g.side == GroupSide::Right); + let fits = l.identity.is_some() + || match grouping { + _ if right_rows => { + r.identity.is_none() && names(&r, &right).is_subset(&names(&l, &left)) + } + Some(g) => g + .labels + .iter() + .all(|label| names(&l, &left).contains(label)), + None => true, + }; + let or = matches!(operator.kind, BinaryOpKind::Set(PromQLVectorSetOpKind::Or)); + // A series identity holds any label set, since `write` re-encodes it. + // A right row needs a time only if the left layout has one. + if !fits || (or && left.time_index.is_some() && right.time_index.is_none()) { + return Err(invalid( + "PromQL binary result labels do not fit the left schema", + )); + } + } + let output = if scalars == [true, false] { + right.clone() + } else { + left.clone() + }; Ok(Self { - kind: Kind::SeriesBinary { operator }, - inputs: vec![left.clone(), right], - output: left, + kind: Kind::SeriesBinary { operator, scalars }, + inputs: vec![left, right], + output, }) } } +/// PromQL's matching signature: `on` keeps only the listed labels; `ignoring` +/// drops them and the metric name. +fn signature(kind: &VectorMatchKind, names: &[String], mut labels: Labels) -> Labels { + match kind { + VectorMatchKind::On => labels.retain(|k, _| names.contains(k)), + VectorMatchKind::Ignoring => labels.retain(|k, _| k != "__name__" && !names.contains(k)), + } + labels +} + +fn label_bytes(labels: &Labels) -> usize { + labels.iter().map(|(k, v)| 64 + k.len() + v.len()).sum() +} + +fn float(layout: &Layout, row: &[Value]) -> Result { + match row[layout.value] { + Value::Float64(value) => Ok(value), + _ => Err(invalid("vector value must be Float64")), + } +} + +/// The value of a matched pair, or `None` when a comparison filters it out. +/// A filter keeps the left value. +fn apply(operator: &BinaryOperator, left: f64, right: f64) -> Result, Error> { + let unmatched = BinaryOperator { + vector_match: None, + ..operator.clone() + }; + match crate::expressions::arithmetic::evaluate_binary(&unmatched, left, right)? { + Value::Float64(value) => Ok(Some(value)), + Value::Bool(keep) => Ok(keep.then_some(left)), + _ => Err(invalid("PromQL binary result must be a number")), + } +} + +/// Evaluate `Kind::SeriesBinary` over the collected operand rows. +async fn series_binary( + operator: &Operator, + binary: &BinaryOperator, + scalars: [bool; 2], + rows: [Vec>; 2], + work: &mut Cooperative, + workspace: &mut Workspace, +) -> Result>, Error> { + use planner_types::pre_asap::{GroupSide, PromQLVectorSetOpKind}; + let drops_name = matches!( + binary.kind, + BinaryOpKind::Arithmetic(_) | BinaryOpKind::CompareBool(_) + ); + let scalar = |rows: &[Vec]| match rows { + [row] => match row[0] { + Value::Float64(value) => Ok(value), + _ => Err(invalid("PromQL scalar must be Float64")), + }, + _ => Err(invalid("PromQL scalar operand must have one row")), + }; + let [left, right] = rows; + let mut result = Vec::new(); + if scalars == [true, true] { + let value = apply(binary, scalar(&left)?, scalar(&right)?)? + .ok_or_else(|| invalid("scalar comparison requires bool"))?; + return Ok(vec![vec![Value::Float64(value)]]); + } + if scalars[0] || scalars[1] { + let (constant, vector, schema) = if scalars[0] { + (scalar(&left)?, right, &operator.inputs[1]) + } else { + (scalar(&right)?, left, &operator.inputs[0]) + }; + let layout = layout(schema)?; + let mut seen = std::collections::BTreeSet::new(); + for mut row in vector { + work.checkpoint().await?; + let value = float(&layout, &row)?; + let (l, r) = if scalars[0] { + (constant, value) + } else { + (value, constant) + }; + let Some(mut computed) = apply(binary, l, r)? else { + continue; + }; + // A filter keeps the vector's value, even on the right. + if matches!(binary.kind, BinaryOpKind::Compare(_)) { + computed = value; + } + let mut set = layout.read(schema, &row)?; + if drops_name { + set.remove("__name__"); + layout.write(schema, &mut row, &set)?; + } + row[layout.value] = Value::Float64(computed); + workspace.grow(row_bytes(&row) + label_bytes(&set))?; + if !seen.insert(set) { + return Err(invalid( + "vector cannot contain metrics with the same labelset", + )); + } + result.push(row); + } + return Ok(result); + } + let schemas = [&operator.inputs[0], &operator.inputs[1]]; + let layouts = [layout(schemas[0])?, layout(schemas[1])?]; + let (kind, names, grouping) = match &binary.vector_match { + None => (VectorMatchKind::Ignoring, &[][..], None), + Some(m) => (m.kind.clone(), m.labels.as_slice(), m.grouping.as_ref()), + }; + let read = |side: usize, row: &[Value]| layouts[side].read(schemas[side], row); + if let BinaryOpKind::Set(set) = &binary.kind { + // `and`/`unless` look up the right side; `or` adds unmatched right rows. + let lookup = if *set == PromQLVectorSetOpKind::Or { + &left + } else { + &right + }; + let side = usize::from(*set != PromQLVectorSetOpKind::Or); + let mut signatures = std::collections::BTreeSet::new(); + for row in lookup { + work.checkpoint().await?; + let labels = signature(&kind, names, read(side, row)?); + workspace.grow(label_bytes(&labels))?; + signatures.insert(labels); + } + let and = *set == PromQLVectorSetOpKind::And; + for row in &left { + work.checkpoint().await?; + let found = signatures.contains(&signature(&kind, names, read(0, row)?)); + if *set == PromQLVectorSetOpKind::Or || found == and { + workspace.grow(row_bytes(row))?; + result.push(row.clone()); + } + } + if *set == PromQLVectorSetOpKind::Or { + for row in &right { + work.checkpoint().await?; + let labels = read(1, row)?; + if signatures.contains(&signature(&kind, names, labels.clone())) { + continue; + } + let mut out = vec![Value::Null; schemas[0].fields.len()]; + if let (Some(to), Some(from)) = (schemas[0].time_index, schemas[1].time_index) { + out[to] = row[from].clone(); + } + out[layouts[0].value] = Value::Float64(float(&layouts[1], row)?); + layouts[0].write(schemas[0], &mut out, &labels)?; + workspace.grow(row_bytes(&out))?; + result.push(out); + } + } + return Ok(result); + } + // Prometheus returns before matching when either side is empty. + if left.is_empty() || right.is_empty() { + return Ok(result); + } + // `group_right` makes the left side the "one" side. + let swapped = matches!(grouping, Some(g) if g.side == GroupSide::Right); + let (one, many) = if swapped { (0, 1) } else { (1, 0) }; + let sides = [&left, &right]; + let mut ones = BTreeMap::new(); + for (index, row) in sides[one].iter().enumerate() { + work.checkpoint().await?; + let labels = read(one, row)?; + let key = signature(&kind, names, labels.clone()); + workspace.grow(2 * label_bytes(&labels))?; + if ones + .insert(key, (labels, float(&layouts[one], row)?, index)) + .is_some() + { + return Err(invalid(&format!( + "found duplicate series for the match group on the {} hand-side of the operation", + if swapped { "left" } else { "right" } + ))); + } + } + let mut matched = BTreeMap::>::new(); + for row in sides[many].iter() { + work.checkpoint().await?; + let labels = read(many, row)?; + let key = signature(&kind, names, labels.clone()); + let Some((one_labels, one_value, one_index)) = ones.get(&key) else { + continue; + }; + let value = float(&layouts[many], row)?; + let (l, r) = if swapped { + (*one_value, value) + } else { + (value, *one_value) + }; + let Some(computed) = apply(binary, l, r)? else { + continue; + }; + let mut metric = labels; + if matches!(binary.kind, BinaryOpKind::Arithmetic(_)) { + metric.remove("__name__"); + } + match grouping { + None => match kind { + VectorMatchKind::On => metric.retain(|k, _| names.contains(k)), + VectorMatchKind::Ignoring => metric.retain(|k, _| !names.contains(k)), + }, + // Included labels come from the "one" side. + Some(g) => { + for label in &g.labels { + match one_labels.get(label) { + Some(value) => metric.insert(label.clone(), value.clone()), + None => metric.remove(label), + }; + } + } + } + if matches!(binary.kind, BinaryOpKind::CompareBool(_)) { + metric.remove("__name__"); + } + workspace.grow(label_bytes(&metric))?; + let results = matched.entry(key).or_default(); + if grouping.is_none() && !results.is_empty() { + return Err(invalid( + "many-to-one matching must be explicit (group_left/group_right)", + )); + } + if !results.insert(metric.clone()) { + return Err(invalid( + "multiple matches for labels: grouping labels must ensure unique matches", + )); + } + // The output row is in the left layout. + let mut out = if swapped { + left[*one_index].clone() + } else { + row.clone() + }; + layouts[0].write(schemas[0], &mut out, &metric)?; + out[layouts[0].value] = Value::Float64(computed); + workspace.grow(row_bytes(&out))?; + result.push(out); + } + Ok(result) +} + pub(super) fn execute<'a>( operator: &'a Operator, mut inputs: Vec>, context: RunContext, ) -> Result, Error> { let output = operator.output.clone(); - let left_layout = layout(&operator.inputs[0])?; let right = match operator.kind { Kind::SeriesBinary { .. } => { Some(inputs.pop().ok_or_else(|| invalid("missing right input"))?) @@ -134,64 +487,91 @@ pub(super) fn execute<'a>( let (rows, _memory) = collect_rows(left, &context).await?; let mut result = Vec::new(); match (&operator.kind, right) { - (Kind::SeriesLabels { kind, labels }, None) => { + ( + Kind::SeriesLabels { + kind, + labels, + unique, + }, + None, + ) => { + let left_layout = layout(&operator.inputs[0])?; + let mut seen = std::collections::BTreeSet::new(); for mut row in rows { work.checkpoint().await?; - let mut set = left_layout.read(&output, &row)?; - match kind { - VectorMatchKind::On => set.retain(|k, _| labels.contains(k)), - VectorMatchKind::Ignoring => { - set.retain(|k, _| k != "__name__" && !labels.contains(k)) - } - } + let set = signature(kind, labels, left_layout.read(&output, &row)?); left_layout.write(&output, &mut row, &set)?; workspace.grow(row_bytes(&row))?; + // Matching and grouping sides may repeat a label set; + // their consumers decide whether that is an error. + if *unique { + workspace.grow(label_bytes(&set))?; + if !seen.insert(set) { + return Err(invalid( + "vector cannot contain metrics with the same labelset", + )); + } + } result.push(row); } } - (Kind::SeriesBinary { operator: binary }, Some(right)) => { - let right_schema = &operator.inputs[1]; - let right_layout = layout(right_schema)?; + ( + Kind::SeriesBinary { + operator: binary, + scalars, + }, + Some(right), + ) => { let (right, _right_memory) = collect_rows(right, &context).await?; - // Prometheus returns before matching when either side is empty. - if rows.is_empty() || right.is_empty() { - return Batch::try_new(output.clone(), vec![]); - } - let mut matches = BTreeMap::new(); - for row in &right { + result = series_binary( + operator, + binary, + *scalars, + [rows, right], + &mut work, + &mut workspace, + ) + .await?; + } + (Kind::SeriesHistogramQuantile { quantile, le }, None) => { + let input = &operator.inputs[0]; + let left_layout = layout(input)?; + let bucket = &input.fields[*le].name; + let mut histograms = BTreeMap::>::new(); + for row in rows { work.checkpoint().await?; - let set = right_layout.read(right_schema, row)?; - workspace.grow(set.iter().map(|(k, v)| 64 + k.len() + v.len()).sum())?; - let Value::Float64(value) = row[right_layout.value] else { + let mut set = left_layout.read(input, &row)?; + // Prometheus skips a series whose `le` is not a float. + let Some(bound) = set.remove(bucket).and_then(|v| parse_bound(&v)) else { + continue; + }; + let Value::Float64(count) = row[left_layout.value] else { return Err(invalid("vector value must be Float64")); }; - // Prometheus rejects a duplicate on the one side. - if matches.insert(set, (value, false)).is_some() { - return Err(invalid( - "duplicate series for a match group on the right-hand side", - )); - } + workspace.grow( + 64 + set + .iter() + .map(|(k, v)| 64 + k.len() + v.len()) + .sum::(), + )?; + // The histogram identity keeps `__name__`; only the result drops it. + histograms.entry(set).or_default().push((bound, count)); } - for mut row in rows { + let out_layout = layout(&output)?; + let q = f64::from_bits(*quantile); + let mut seen = std::collections::BTreeSet::new(); + for (mut set, buckets) in histograms { work.checkpoint().await?; - let mut set = left_layout.read(&output, &row)?; - let Some((value, matched)) = matches.get_mut(&set) else { - continue; - }; - // Only a left duplicate that finds a match is ambiguous. - if std::mem::replace(matched, true) { + set.remove("__name__"); + let value = aggregate::temporal::bucket_quantile(q, buckets, &context).await?; + let mut row = vec![Value::Null; output.fields.len()]; + out_layout.write(&output, &mut row, &set)?; + row[out_layout.value] = Value::Float64(value); + if !seen.insert(set) { return Err(invalid( - "many-to-one matching must be explicit (group_left/group_right)", + "vector cannot contain metrics with the same labelset", )); } - let Value::Float64(left_value) = row[left_layout.value] else { - return Err(invalid("vector value must be Float64")); - }; - row[left_layout.value] = crate::expressions::arithmetic::evaluate_binary( - binary, left_value, *value, - )?; - set.remove("__name__"); - left_layout.write(&output, &mut row, &set)?; workspace.grow(row_bytes(&row))?; result.push(row); } @@ -202,3 +582,14 @@ pub(super) fn execute<'a>( }) .boxed_local()) } + +/// A bucket bound as Go's `strconv.ParseFloat` reads it, except hex floats: +/// an out-of-range literal is an error, not an infinity. +fn parse_bound(text: &str) -> Option { + let bound = text.parse::().ok()?; + let unsigned = text.strip_prefix(['+', '-']).unwrap_or(text); + let infinite = ["inf", "infinity"] + .iter() + .any(|word| unsigned.eq_ignore_ascii_case(word)); + (bound.is_finite() || bound.is_nan() || infinite).then_some(bound) +} diff --git a/crates/asap-physical-operators/src/operators/unchecked.rs b/crates/asap-physical-operators/src/operators/unchecked.rs index 46bc6a0c..034eb754 100644 --- a/crates/asap-physical-operators/src/operators/unchecked.rs +++ b/crates/asap-physical-operators/src/operators/unchecked.rs @@ -65,11 +65,17 @@ impl TryFrom for Operator { at_ms, steps, )?, - Kind::SeriesLabels { kind, labels } => { + // The rebuilt kind must equal the serialized one, which rejects + // a unique rewrite with other matching labels. + Kind::SeriesLabels { unique: true, .. } => Operator::series_without_name(input(0)?)?, + Kind::SeriesLabels { kind, labels, .. } => { Operator::series_labels(input(0)?, kind, labels)? } - Kind::SeriesBinary { operator } => { - Operator::series_binary(input(0)?, input(1)?, operator)? + Kind::SeriesBinary { operator, scalars } => { + Operator::series_binary(input(0)?, input(1)?, operator, scalars)? + } + Kind::SeriesHistogramQuantile { quantile, le } => { + Operator::series_histogram_quantile(input(0)?, f64::from_bits(quantile), le)? } Kind::Project(expressions) => { if expressions.len() != output.fields.len() { diff --git a/crates/asap-physical-operators/src/operators/vector_binary.rs b/crates/asap-physical-operators/src/operators/vector_binary.rs index cc60ff4c..88f019c2 100644 --- a/crates/asap-physical-operators/src/operators/vector_binary.rs +++ b/crates/asap-physical-operators/src/operators/vector_binary.rs @@ -46,8 +46,14 @@ impl Operator { left: Schema, right: Schema, operator: BinaryOperator, - return_bool: bool, + mut return_bool: bool, ) -> Result { + let mut operator = operator; + // The IR's `bool` comparison is this operator's `return_bool` mode. + if let BinaryOpKind::CompareBool(op) = &operator.kind { + operator.kind = BinaryOpKind::Compare(op.clone()); + return_bool = true; + } let scalar = is_scalar(&left)? && is_scalar(&right)?; is_scalar(&right)?; let expression = Expression::Binary { diff --git a/crates/asap-physical-operators/src/operators/vector_window.rs b/crates/asap-physical-operators/src/operators/vector_window.rs index 407f04ea..7592d0b5 100644 --- a/crates/asap-physical-operators/src/operators/vector_window.rs +++ b/crates/asap-physical-operators/src/operators/vector_window.rs @@ -130,7 +130,10 @@ pub(super) fn execute<'a>( } let result = aggregate::temporal::reduce( rows, - &AggIntent::HistogramQuantile { q: *q }, + &AggIntent::HistogramQuantile { + q: *q, + le: ColumnRef::Named("le".into()), + }, &[0], 1, 2, diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 655c024a..b8312ef1 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -464,6 +464,25 @@ fn compile_internal( if let Payload::Binary { operator } = &node.payload { let query_time = node.output_state.timing == planner_types::post_asap::ExecutionTiming::QueryTime; + // Per-series readouts keep `__name__` even where the range + // function drops it; only name-dropping operators ignore that. + let per_series = schemas.iter().any(|schema| { + schema + .fields + .iter() + .any(|f| f.name == promql_rows::SERIES_IDENTITY_COLUMN) + }); + if per_series + && matches!( + operator.kind, + planner_types::pre_asap::BinaryOpKind::Compare(_) + | planner_types::pre_asap::BinaryOpKind::Set(_) + ) + { + return Err(invalid(format!( + "node {id}: per-series readouts do not apply the range function's __name__ rule" + ))); + } if let Some(&(value, left)) = literals.get(&id) { let [input] = schemas.as_slice() else { return Err(invalid("scalar binary requires one row input")); @@ -471,9 +490,27 @@ fn compile_internal( if !query_time { return Err(invalid("scalar literal binary must run at query time")); } - let project = row_values::scalar_binary(input, operator, value, left) + let scalar = + Operator::scalar(crate::values::Value::Float64(value), DataType::Float64)?; + let (sides, scalars, operands) = if left { + ( + [scalar.schema(), input.clone()], + [true, false], + vec![auxiliary, inputs[0]], + ) + } else { + ( + [input.clone(), scalar.schema()], + [false, true], + vec![inputs[0], auxiliary], + ) + }; + let [l, r] = sides; + let binary = Operator::series_binary(l, r, operator.clone(), scalars) .map_err(|error| invalid(format!("node {id}: {error}")))?; - graph.add(id, inputs, project.with_output_schema(output)?)?; + graph.add(auxiliary, vec![], scalar)?; + graph.add(id, operands, binary.with_output_schema(output)?)?; + auxiliary -= 1; continue; } let label_map = |schema: &Schema| { @@ -482,13 +519,26 @@ fn compile_internal( .iter() .any(|f| matches!(f.dtype, SummaryFamilyType::Plain(DataType::Map { .. }))) }; + // Grouped rows carry their labels as columns; per-series rows + // carry the series identity. if let (true, [left, right]) = (query_time, schemas.as_slice()) { if !label_map(left) && !label_map(right) { - let (join, project) = row_values::grouped_binary(left, right, operator) - .map_err(|error| invalid(format!("node {id}: {error}")))?; - graph.add(auxiliary, inputs, join)?; - graph.add(id, vec![auxiliary], project.with_output_schema(output)?)?; - auxiliary -= 1; + // A scalar-valued Fallback operand, such as `scalar(x)`, has no labels. + let scalar = |input: &NodeId| { + matches!( + nodes.get(input).map(|node| &node.payload), + Some(Payload::Fallback { expression }) + if promql_fallback::scalar(expression) + ) + }; + let binary = Operator::series_binary( + left.clone(), + right.clone(), + operator.clone(), + [scalar(&inputs[0]), scalar(&inputs[1])], + ) + .map_err(|error| invalid(format!("node {id}: {error}")))?; + graph.add(id, inputs, binary.with_output_schema(output)?)?; continue; } } diff --git a/crates/asap-physical-operators/src/physical_planner/promql_fallback.rs b/crates/asap-physical-operators/src/physical_planner/promql_fallback.rs index 12c60ad1..f9efd556 100644 --- a/crates/asap-physical-operators/src/physical_planner/promql_fallback.rs +++ b/crates/asap-physical-operators/src/physical_planner/promql_fallback.rs @@ -4,7 +4,7 @@ use super::*; use crate::operators::SubquerySteps; use planner_types::post_asap::execution_data_state::lift_plain; -use planner_types::pre_asap::{AtModifier, BinaryOpKind, VectorMatch, VectorMatchKind}; +use planner_types::pre_asap::{AtModifier, VectorMatchKind}; /// Input slot for the raw series read by the `selector`th selector (in /// [`raw_series`] order) of Fallback node `node`. The node's own ID names its @@ -84,14 +84,16 @@ fn selector(expression: &QueryExpr) -> Result<(i64, i64, Option), Error> { Ok((millis(range)?, offset, at)) } -/// PromQL scalar-valued expressions have no labels to match. -fn scalar(expression: &QueryExpr) -> bool { - matches!( - expression, +/// PromQL scalar-valued expressions have no labels to match. A binary +/// operator is scalar-valued when both operands are. +pub(super) fn scalar(expression: &QueryExpr) -> bool { + match expression { QueryExpr::PromqlScalarBridge(_) - | QueryExpr::PromqlScalarFromVector(_) - | QueryExpr::EvalTimestamp - ) + | QueryExpr::PromqlScalarFromVector(_) + | QueryExpr::EvalTimestamp => true, + QueryExpr::BinaryOp { lhs, rhs, .. } => scalar(lhs) && scalar(rhs), + _ => false, + } } impl Lowering { @@ -155,7 +157,13 @@ impl Lowering { let [function] = measures.as_slice() else { return Err(invalid("range function requires one measure")); }; - self.range_function(function, child, expression) + let step = self.range_function(function, child, expression)?; + if matches!(function, AggIntent::LastOverTime) { + return Ok(step); + } + // Other range functions drop the name; equal label sets then error. + let input = self.schema(&step); + Ok(self.add(Operator::series_without_name(input)?, vec![step])) } QueryExpr::Aggregate { reduction: planner_types::pre_asap::Reduction::Reduce(keys), @@ -206,56 +214,25 @@ impl Lowering { ) } QueryExpr::BinaryOp { - op: BinaryOpKind::Arithmetic(op), + op, lhs, rhs, vector_match, } => { - let (vector, literal, literal_left) = match ( - row_values::scalar_literal(lhs), - row_values::scalar_literal(rhs), - ) { - (None, Some(value)) => (lhs, value, false), - (Some(value), None) => (rhs, value, true), - (None, None) => return self.match_vectors(expression, vector_match), - _ => return Err(invalid("PromQL arithmetic between two literals")), - }; - let step = self.value(vector)?; - let input = self.schema(&step); - let value = named_column(&input, &ColumnRef::SampleValue)?; - let literal = Expression::Literal { - value: crate::values::Value::Float64(literal), - dtype: DataType::Float64, - }; + let sides = vec![self.value(lhs)?, self.value(rhs)?]; let operator = planner_types::post_asap::BinaryOperator { - kind: planner_types::pre_asap::BinaryOpKind::Arithmetic(op.clone()), - vector_match: None, + kind: op.clone(), + vector_match: vector_match.clone(), checked_relative_division: false, checked_finite_division: false, }; - let (left, right) = if literal_left { - (literal, Expression::Column(value)) - } else { - (Expression::Column(value), literal) - }; - let columns = input - .fields - .iter() - .enumerate() - .map(|(i, field)| { - let expression = if i == value { - Expression::Binary { - operator: operator.clone(), - left: Box::new(left.clone()), - right: Box::new(right.clone()), - } - } else { - Expression::Column(i) - }; - (field.name.clone(), expression) - }) - .collect(); - self.push(Operator::project(input, columns)?, vec![step], expression) + let binary = Operator::series_binary( + self.schema(&sides[0]), + self.schema(&sides[1]), + operator, + [scalar(lhs), scalar(rhs)], + )?; + self.push(binary, sides, expression) } QueryExpr::PromqlScalarFromVector(child) => { let step = self.value(child)?; @@ -288,54 +265,6 @@ impl Lowering { } } - /// One-to-one vector arithmetic: both sides reduce to their matching - /// labels, which are also the result's labels. - fn match_vectors( - &mut self, - logical: &QueryExpr, - vector_match: &Option, - ) -> Result { - let QueryExpr::BinaryOp { - op: kind, lhs, rhs, .. - } = logical - else { - unreachable!() - }; - if scalar(lhs) || scalar(rhs) { - return Err(invalid("PromQL arithmetic with a non-literal scalar")); - } - let (matching, labels) = match vector_match { - None => (VectorMatchKind::Ignoring, vec![]), - Some(VectorMatch { - kind, - labels, - grouping: None, - }) => (kind.clone(), labels.clone()), - Some(_) => return Err(invalid("group_left/group_right matching is unsupported")), - }; - let mut sides = Vec::new(); - for side in [lhs, rhs] { - let step = self.value(side)?; - let input = self.schema(&step); - sides.push(self.add( - Operator::series_labels(input, matching.clone(), labels.clone())?, - vec![step], - )); - } - let (left, right) = (self.schema(&sides[0]), self.schema(&sides[1])); - let operator = planner_types::post_asap::BinaryOperator { - kind: kind.clone(), - vector_match: None, - checked_relative_division: false, - checked_finite_division: false, - }; - self.push( - Operator::series_binary(left, right, operator)?, - sides, - logical, - ) - } - /// `function(matrix)`, where the matrix is a range selector or a subquery. fn range_function( &mut self, @@ -389,11 +318,25 @@ impl Lowering { let (range, inner_offset, inner_at) = selector(selected)?; let raw = self.read(selected)?; let schema = self.schema(&raw); - let step = self.push( - Operator::series_window(schema, inner, range, inner_offset, inner_at, Some(steps))?, + let mut step = self.push( + Operator::series_window( + schema, + inner.clone(), + range, + inner_offset, + inner_at, + Some(steps), + )?, vec![raw], child, )?; + // The inner function drops the name too. A series repeats across + // steps, so this rewrite does not check for equal label sets. + if inner.is_some() && !matches!(inner, Some(AggIntent::LastOverTime)) { + let input = self.schema(&step); + let relabel = Operator::series_labels(input, VectorMatchKind::Ignoring, vec![])?; + step = self.add(relabel, vec![step]); + } let input = self.schema(&step); self.push( Operator::series_window(input, Some(function), steps.range_ms, offset, at_ms, None)?, @@ -412,6 +355,13 @@ impl Lowering { logical: &QueryExpr, ) -> Result { let mut input = self.schema(&step); + if let AggIntent::HistogramQuantile { q, le } = measure { + if !keys.is_without() || keys.keys() != [*le] { + return Err(invalid("histogram_quantile must group without (le)")); + } + let operator = Operator::series_histogram_quantile(input, *q, *le)?; + return self.push(operator, vec![step], logical); + } let value = named_column(&input, &ColumnRef::SampleValue)?; let reduction = match measure { AggIntent::Sum { col: None } => Reduction::Sum(value), diff --git a/crates/asap-physical-operators/src/physical_planner/row_values.rs b/crates/asap-physical-operators/src/physical_planner/row_values.rs index 737c8065..8763437b 100644 --- a/crates/asap-physical-operators/src/physical_planner/row_values.rs +++ b/crates/asap-physical-operators/src/physical_planner/row_values.rs @@ -1,8 +1,7 @@ //! Query-time PromQL value computation over logical row schemas. use super::*; -use planner_types::post_asap::{maintained_population::PopulationReadout, BinaryOperator}; -use planner_types::pre_asap::{BinaryOpKind, DataType, Predicate, ScalarValue}; -use std::rc::Rc; +use planner_types::post_asap::maintained_population::PopulationReadout; +use planner_types::pre_asap::{DataType, ScalarValue}; /// A PromQL number literal has no row schema; its consumer folds it in. pub(super) fn scalar_literal(expression: &QueryExpr) -> Option { @@ -13,163 +12,6 @@ pub(super) fn scalar_literal(expression: &QueryExpr) -> Option { } } -/// Rows without a time column, label map, or series identity carry only -/// their group labels, so those labels are the complete PromQL identity. -fn grouped_value(input: &Schema) -> Result<(usize, Vec), Error> { - if input.time_index.is_some() - || input - .fields - .iter() - .any(|field| field.name == promql_rows::SERIES_IDENTITY_COLUMN) - { - return Err(invalid( - "row binary requires grouped rows; per-series matching needs a name-free identity", - )); - } - let mut value = None; - let mut labels = Vec::new(); - for (i, field) in input.fields.iter().enumerate() { - match &field.dtype { - SummaryFamilyType::Plain(DataType::Float64) if value.is_none() => value = Some(i), - SummaryFamilyType::Plain(DataType::Utf8) => labels.push(i), - _ => { - return Err(invalid( - "row binary requires Utf8 labels and one Float64 value", - )) - } - } - } - Ok(( - value.ok_or_else(|| invalid("row binary requires a Float64 value"))?, - labels, - )) -} - -fn arithmetic(operator: &BinaryOperator) -> Result<(), Error> { - if !matches!(operator.kind, BinaryOpKind::Arithmetic(_)) { - return Err(invalid( - "row comparison requires filter or bool semantics, which Binary does not carry", - )); - } - Ok(()) -} - -/// Apply `vector op scalar` (or `scalar op vector`) to each row's value. -pub(super) fn scalar_binary( - input: &Schema, - operator: &BinaryOperator, - literal: f64, - literal_left: bool, -) -> Result { - arithmetic(operator)?; - let (value, _) = grouped_value(input)?; - let literal = Expression::Literal { - value: crate::values::Value::Float64(literal), - dtype: DataType::Float64, - }; - let columns = input - .fields - .iter() - .enumerate() - .map(|(i, field)| { - let expression = if i != value { - Expression::Column(i) - } else if literal_left { - binary(operator, literal.clone(), Expression::Column(i)) - } else { - binary(operator, Expression::Column(i), literal.clone()) - }; - (field.name.clone(), expression) - }) - .collect(); - Operator::project(input.clone(), columns) -} - -/// One-to-one PromQL matching of grouped rows on equal label sets. Returns -/// the inner equi-join and the projection that applies the operator. -pub(super) fn grouped_binary( - left: &Schema, - right: &Schema, - operator: &BinaryOperator, -) -> Result<(Operator, Operator), Error> { - arithmetic(operator)?; - let (left_value, left_labels) = grouped_value(left)?; - let (right_value, right_labels) = grouped_value(right)?; - if left_labels.len() != right_labels.len() { - return Err(invalid("row binary inputs have different label sets")); - } - let width = left.fields.len(); - let keys = left_labels - .iter() - .map(|&l| { - let name = &left.fields[l].name; - let r = right_labels - .iter() - .copied() - .find(|&r| &right.fields[r].name == name) - .ok_or_else(|| invalid("row binary inputs have different label sets"))?; - let (a, b) = ( - Rc::new(QueryExpr::Column(l)), - Rc::new(QueryExpr::Column(width + r)), - ); - let equal = QueryExpr::Compare { - left: a.clone(), - op: CompareOpKind::Eq, - right: b.clone(), - }; - // A nullable label compares like PromQL's empty label: absent on both sides matches. - Ok(if left.fields[l].nullable || right.fields[r].nullable { - QueryExpr::BoolOr(vec![ - equal, - QueryExpr::BoolAnd(vec![QueryExpr::IsNull(a), QueryExpr::IsNull(b)]), - ]) - } else { - equal - }) - }) - .collect::, Error>>()?; - let predicate = Predicate(Rc::new(QueryExpr::BoolAnd(keys))); - let mut joined = left.fields.clone(); - joined.extend(right.fields.iter().cloned()); - let join = Operator::relational_join( - left.clone(), - right.clone(), - planner_types::pre_asap::JoinKind::Inner, - &predicate, - Arc::new(planner_types::post_asap::SummarySchema { - fields: joined, - time_index: None, - }), - )?; - let columns = left - .fields - .iter() - .enumerate() - .map(|(i, field)| { - let expression = if i == left_value { - binary( - operator, - Expression::Column(i), - Expression::Column(width + right_value), - ) - } else { - Expression::Column(i) - }; - (field.name.clone(), expression) - }) - .collect(); - let project = Operator::project(join.schema(), columns)?; - Ok((join, project)) -} - -fn binary(operator: &BinaryOperator, left: Expression, right: Expression) -> Expression { - Expression::Binary { - operator: operator.clone(), - left: Box::new(left), - right: Box::new(right), - } -} - /// Aggregate readouts of a maintained current-series population, as a chain. pub(super) fn population_aggregate( input: &Schema, diff --git a/crates/asap-physical-operators/tests/deployment_computation.rs b/crates/asap-physical-operators/tests/deployment_computation.rs index b026f47f..edef3121 100644 --- a/crates/asap-physical-operators/tests/deployment_computation.rs +++ b/crates/asap-physical-operators/tests/deployment_computation.rs @@ -97,8 +97,23 @@ fn raw_inputs(dag: &PostAsapDag) -> Vec<(u64, Arc, String)> { type Sample = (&'static str, &'static str, &'static str, i64, f64); /// Compile, round-trip, bind raw `(metric, job, instance, ts, value)` samples, -/// and return `(job, value)` rows of the root. -fn run(dag: &PostAsapDag, samples: &[Sample], end: i64) -> Result, String> { +/// and return the root's batches. +fn execute( + dag: &PostAsapDag, + samples: &[Sample], + end: i64, +) -> Result>, String> { + execute_relabeled(dag, samples, end, &BTreeMap::new()) +} + +/// [`execute`], supplying samples of each instance in `relabel` under its +/// `(__name__, instance)` instead. +fn execute_relabeled( + dag: &PostAsapDag, + samples: &[Sample], + end: i64, + relabel: &BTreeMap<&str, (&str, &str)>, +) -> Result>, String> { let inputs = raw_inputs(dag); let program = compile( dag, @@ -118,6 +133,8 @@ fn run(dag: &PostAsapDag, samples: &[Sample], end: i64) -> Result Result job.to_string(), - _ => String::new(), - }; - let value = match row.last() { - Some(Value::Float64(v)) => *v, - Some(Value::Int64(v)) => *v as f64, - other => return Err(format!("unexpected value {other:?}")), - }; - assert!(rows.insert(key, value).is_none(), "duplicate output group"); - } + batches.push(batch.map_err(|e| e.to_string())?); } - Ok(rows) + Ok(batches) }) } +/// [`execute`], returning `(job, value)` rows of the root. +fn run(dag: &PostAsapDag, samples: &[Sample], end: i64) -> Result, String> { + let mut rows = BTreeMap::new(); + for batch in execute(dag, samples, end)? { + let job = batch.schema().fields.iter().position(|f| f.name == "job"); + for row in batch.rows() { + let key = match job.map(|i| &row[i]) { + Some(Value::Utf8(job)) => job.to_string(), + _ => String::new(), + }; + let value = match row.last() { + Some(Value::Float64(v)) => *v, + Some(Value::Int64(v)) => *v as f64, + other => return Err(format!("unexpected value {other:?}")), + }; + assert!(rows.insert(key, value).is_none(), "duplicate output group"); + } + } + Ok(rows) +} + const SAMPLES: &[Sample] = &[ ("m", "api", "a", 10_000, 4.), ("m", "api", "a", 50_000, 1.), @@ -323,18 +348,269 @@ fn exact_count_finalizes_to_declared_float_value() { ); } -// Comparisons need filter/bool semantics that `Binary` does not carry, so -// they fail at compile time instead of emitting 0/1 values. +/// `dag` with its Binary operator replaced by `kind`. +fn with_kind(mut dag: PostAsapDag, kind: planner_types::pre_asap::BinaryOpKind) -> PostAsapDag { + for node in &mut dag.nodes { + if let PostAsapOperatorPayload::Binary { operator } = &mut node.payload { + operator.kind = kind.clone(); + } + } + dag +} + +// A comparison Binary over grouped values keeps the groups whose comparison +// holds, with their value, on either side of the literal; `bool` yields 1 or 0. +#[test] +fn grouped_comparisons_filter_or_return_bool() { + use planner_types::pre_asap::{BinaryOpKind::*, CompareOpKind::Gt}; + // sum_over_time over 5m per job: api = 14, db = 5. + let right = exact_dag("sum by (job) (sum_over_time(m[5m])) * 10"); + let left = exact_dag("10 - sum by (job) (sum_over_time(m[5m]))"); + for (dag, expected) in [ + ( + with_kind(right.clone(), Compare(Gt)), + reference(&[("api", 14.)]), + ), + ( + with_kind(right, CompareBool(Gt)), + reference(&[("api", 1.), ("db", 0.)]), + ), + (with_kind(left, Compare(Gt)), reference(&[("db", 5.)])), + ] { + assert_eq!(run(&dag, SAMPLES, 60_000).unwrap(), expected); + } +} + +// A `bool` comparison Binary over per-series readouts matches one-to-one and +// drops the metric name; a filter fails closed. +#[test] +fn per_series_comparisons_filter_or_return_bool() { + use planner_types::pre_asap::{BinaryOpKind::*, CompareOpKind::*}; + let samples = counter("a", "api", 10., 10.) + .chain(counter("a", "db", 10., 10.)) + .chain(counter("b", "api", 5., 5.)) + .chain(counter("b", "db", 20., 20.)) + .collect::>(); + // rate: a{api} = a{db} = 50/300, b{api} = 25/300, b{db} = 100/300. + let dag = exact_dag("rate(a[5m]) / rate(b[5m])"); + // The readout keeps `__name__`, which rate drops, so a filter would keep it. + let error = run_series(&with_kind(dag.clone(), Compare(Gt)), &samples, 300_000).unwrap_err(); + assert!(error.contains("__name__"), "{error}"); + assert_eq!( + run_series(&with_kind(dag, CompareBool(Lt)), &samples, 300_000).unwrap(), + series(&[("api", "x", 0.), ("db", "x", 1.)]) + ); +} + +/// [`execute`], returning per-series `(identity, value)` rows of the root, +/// with NaN-aware formatting for comparison. +fn run_series(dag: &PostAsapDag, samples: &[Sample], end: i64) -> Result { + let mut rows = BTreeMap::new(); + for batch in execute(dag, samples, end)? { + let schema = batch.schema(); + let column = |name: &str| schema.fields.iter().position(|f| f.name == name).unwrap(); + let (identity, value) = (column(promql_rows::SERIES_IDENTITY_COLUMN), column("value")); + for row in batch.rows() { + let (Value::Utf8(identity), Value::Float64(value)) = (&row[identity], &row[value]) + else { + return Err(format!("unexpected row {row:?}")); + }; + assert!(rows.insert(identity.to_string(), *value).is_none()); + } + } + Ok(format!("{rows:?}")) +} + +fn series(pairs: &[(&str, &str, f64)]) -> String { + let rows = pairs + .iter() + .map(|(job, instance, value)| { + ( + format!(r#"{{"instance":"{instance}","job":"{job}"}}"#), + *value, + ) + }) + .collect::>(); + format!("{rows:?}") +} + +/// Five counter samples per series, one per minute up to 300s, at `base + step * i`. +fn counter( + metric: &'static str, + job: &'static str, + base: f64, + step: f64, +) -> impl Iterator { + (1..=5).map(move |i| (metric, job, "x", i * 60_000, base + step * (i - 1) as f64)) +} + +// avg_over_time over stored per-series sum/count state divides per series and +// drops the metric name, as Prometheus does. An overflowing stored sum fails +// the checked division instead of returning +Inf. +#[test] +fn per_series_average_divides_stored_sum_by_count() { + let dag = exact_dag("avg_over_time(m[5m])"); + assert_eq!( + run_series(&dag, SAMPLES, 60_000).unwrap(), + series(&[ + ("api", "a", 2.5), + ("api", "b", 7.), + ("api", "c", 2.), + ("db", "d", 5.) + ]) + ); + let huge: &[Sample] = &[ + ("m", "api", "a", 10_000, 1.7e308), + ("m", "api", "a", 20_000, 1.7e308), + ]; + assert!(run_series(&dag, huge, 60_000).is_err()); +} + +// rate(a) / rate(b) matches series on their labels without the metric name. +// Unmatched series are dropped; x/0 is +Inf and 0/0 is NaN; an empty side +// yields an empty vector. +#[test] +fn per_series_rate_ratio_matches_prometheus() { + // Window (0, 300s]: first sample at 60s, extrapolated 60s to the start + // (durationToZero is exactly 60s too), so rate = (last - first) * 1.25 / 300. + // a{api} = 40 * 1.25 / 300, b{api} = 20 * 1.25 / 300, so the ratio is 2. + let samples = counter("a", "api", 10., 10.) + .chain(counter("a", "db", 10., 10.)) + .chain(counter("a", "web", 10., 10.)) + .chain(counter("a", "cache", 7., 0.)) + .chain(counter("b", "api", 5., 5.)) + .chain(counter("b", "web", 7., 0.)) + .chain(counter("b", "cache", 7., 0.)) + .chain(counter("b", "other", 5., 5.)) + .collect::>(); + let dag = exact_dag("rate(a[5m]) / rate(b[5m])"); + assert_eq!( + run_series(&dag, &samples, 300_000).unwrap(), + series(&[ + ("api", "x", 2.), + ("cache", "x", f64::NAN), + ("web", "x", f64::INFINITY), + ]) + ); + let only_a = samples + .iter() + .filter(|s| s.0 == "a") + .copied() + .collect::>(); + assert_eq!(run_series(&dag, &only_a, 300_000).unwrap(), series(&[])); +} + +// A literal operand applies to every stored per-series value, on either side, +// and drops the metric name. #[test] -fn row_comparison_fails_closed() { - let mut dag = exact_dag("sum by (job) (sum_over_time(m[5m])) * 2"); +fn per_series_scalar_arithmetic_applies_to_stored_readouts() { + let samples = counter("m", "api", 10., 10.).collect::>(); + // rate = 40 * 1.25 / 300 = 1/6. + for (query, expected) in [ + ("rate(m[5m]) * 2", 50. / 300. * 2.), + ("1 - rate(m[5m])", 1. - 50. / 300.), + // The stored sum readout keeps `__name__`; the arithmetic drops it. + ("sum_over_time(m[5m]) * 2", 150. * 2.), + ] { + assert_eq!( + run_series(&exact_dag(query), &samples, 300_000).unwrap(), + series(&[("api", "x", expected)]), + "{query}" + ); + } +} + +// A literal operand drops the metric name; series that then share a label set +// are an error, as in Prometheus, rather than duplicate output series. +#[test] +fn per_series_scalar_arithmetic_rejects_label_sets_equal_without_the_name() { + let dag = exact_dag("sum_over_time(m[5m]) * 2"); + let samples = counter("m", "api", 10., 10.) + .chain(counter("m", "api", 1., 1.).map(|s| (s.0, s.1, "y", s.3, s.4))) + .collect::>(); + let distinct = BTreeMap::from([("y", ("n", "y"))]); + assert!(execute_relabeled(&dag, &samples, 300_000, &distinct).is_ok()); + // Both series are then {instance="x",job="api"}, named m and n. + let equal = BTreeMap::from([("y", ("n", "x"))]); + let Err(error) = execute_relabeled(&dag, &samples, 300_000, &equal) else { + panic!("duplicate label sets were accepted"); + }; + assert!(error.contains("same labelset"), "{error}"); +} + +fn with_vector_match( + mut dag: PostAsapDag, + kind: planner_types::pre_asap::VectorMatchKind, + labels: &[&str], +) -> PostAsapDag { for node in &mut dag.nodes { if let PostAsapOperatorPayload::Binary { operator } = &mut node.payload { - operator.kind = planner_types::pre_asap::BinaryOpKind::Compare( - planner_types::pre_asap::CompareOpKind::Gt, - ); + operator.vector_match = Some(planner_types::pre_asap::VectorMatch { + kind: kind.clone(), + labels: labels.iter().map(|l| l.to_string()).collect(), + grouping: None, + }); } } - let error = run(&dag, SAMPLES, 60_000).unwrap_err(); - assert!(error.contains("comparison"), "{error}"); + dag +} + +// `on` and `ignoring` reduce each side to the matching labels, which become +// the result's labels; a duplicate match group on either side is an error. +#[test] +fn per_series_vector_matching_follows_on_and_ignoring() { + use planner_types::pre_asap::VectorMatchKind; + let samples = counter("a", "api", 10., 10.) + .chain(counter("b", "api", 5., 5.)) + .chain(counter("b", "db", 5., 5.)) + .collect::>(); + let dag = exact_dag("rate(a[5m]) / rate(b[5m])"); + for (kind, labels) in [ + (VectorMatchKind::On, ["job"]), + (VectorMatchKind::Ignoring, ["instance"]), + ] { + let dag = with_vector_match(dag.clone(), kind, &labels); + assert_eq!( + run_series(&dag, &samples, 300_000).unwrap(), + format!( + "{:?}", + BTreeMap::from([(r#"{"job":"api"}"#.to_string(), 2.)]) + ) + ); + let mut duplicate = samples.clone(); + duplicate.extend(counter("b", "api", 5., 5.).map(|s| (s.0, s.1, "y", s.3, s.4))); + let error = run_series(&dag, &duplicate, 300_000).unwrap_err(); + assert!(error.contains("duplicate series"), "{error}"); + let mut duplicate = samples.clone(); + duplicate.extend(counter("a", "api", 5., 5.).map(|s| (s.0, s.1, "y", s.3, s.4))); + let error = run_series(&dag, &duplicate, 300_000).unwrap_err(); + assert!(error.contains("many-to-one"), "{error}"); + } +} + +// Current-series sums and averages are compensated like Prometheus, and an +// overflowing running sum does not turn the average into +Inf. +#[test] +fn population_sums_and_averages_are_compensated() { + let cancel: &[Sample] = &[ + ("m", "api", "a", 50_000, 1e100), + ("m", "api", "b", 50_000, 1.), + ("m", "api", "c", 50_000, -1e100), + ]; + let huge: &[Sample] = &[ + ("m", "api", "a", 50_000, 1.7e308), + ("m", "api", "b", 50_000, 1.7e308), + ]; + for (query, samples, expected) in [ + ("sum by (job) (m)", cancel, 1.), + ("avg by (job) (m)", cancel, 1. / 3.), + ("avg by (job) (m)", huge, 1.7e308), + ] { + let dag = population_dag(query); + assert_eq!( + run(&dag, samples, 60_000).unwrap(), + reference(&[("api", expected)]), + "{query}" + ); + } } diff --git a/crates/asap-physical-operators/tests/promql_binary.rs b/crates/asap-physical-operators/tests/promql_binary.rs index 24844b5e..0a2ac35c 100644 --- a/crates/asap-physical-operators/tests/promql_binary.rs +++ b/crates/asap-physical-operators/tests/promql_binary.rs @@ -49,17 +49,18 @@ fn row(name: &str, job: &str, value: f64) -> Vec { ] } fn program() -> CompiledPhysicalDag { + program_for(BinaryOperator { + kind: BinaryOpKind::Arithmetic(ArithmeticOpKind::Div), + vector_match: None, + checked_relative_division: true, + checked_finite_division: false, + }) +} +fn program_for(operator: BinaryOperator) -> CompiledPhysicalDag { let schema = schema(); let node = PostAsapDagNode { id: PostAsapNodeId(2), - payload: PostAsapOperatorPayload::Binary { - operator: BinaryOperator { - kind: BinaryOpKind::Arithmetic(ArithmeticOpKind::Div), - vector_match: None, - checked_relative_division: true, - checked_finite_division: false, - }, - }, + payload: PostAsapOperatorPayload::Binary { operator }, output_state: ExecutionDataState::QUERY_ROWS, output_schema: (*schema).clone(), guarantee: None, @@ -80,7 +81,13 @@ fn evaluate( left: Vec>, right: Vec>, ) -> Result>, asap_physical_operators::Error> { - let graph = program(); + evaluate_with(program(), left, right) +} +fn evaluate_with( + graph: CompiledPhysicalDag, + left: Vec>, + right: Vec>, +) -> Result>, asap_physical_operators::Error> { let sources = [left, right] .into_iter() .enumerate() @@ -268,3 +275,30 @@ fn binary_obeys_memory_and_cancellation() { )); } } + +// A `bool` comparison over label-map vectors yields 1 or 0 and drops the name. +#[test] +fn label_map_bool_comparison_drops_the_name() { + let program = program_for(BinaryOperator { + kind: BinaryOpKind::CompareBool(planner_types::pre_asap::CompareOpKind::Gt), + vector_match: None, + checked_relative_division: false, + checked_finite_division: false, + }); + let rows = evaluate_with( + program, + vec![row("a", "api", 6.)], + vec![row("b", "api", 2.)], + ) + .unwrap(); + let [row] = rows.as_slice() else { + panic!("expected one row, got {}", rows.len()); + }; + let Value::Map(labels) = &row[0] else { + panic!("expected labels"); + }; + assert!(labels + .iter() + .all(|(k, _)| !matches!(k, Value::Utf8(k) if &**k == "__name__"))); + assert!(matches!(row[1], Value::Float64(v) if v == 1.)); +} diff --git a/crates/asap-physical-operators/tests/promql_fallback.rs b/crates/asap-physical-operators/tests/promql_fallback.rs index d274495c..20a2719d 100644 --- a/crates/asap-physical-operators/tests/promql_fallback.rs +++ b/crates/asap-physical-operators/tests/promql_fallback.rs @@ -18,13 +18,17 @@ use std::{collections::BTreeMap, rc::Rc}; /// Bare selectors look back one ingestion interval: 60s. fn parse(query: &str) -> QueryExpr { + parse_with(query, AccuracyTarget::Exact) +} + +fn parse_with(query: &str, accuracy: AccuracyTarget) -> QueryExpr { let workload = PlanningWorkload { query_workload: QueryWorkload { language: QueryLanguage::PromQL, query_batch: Some(vec![BatchEntry { query: Query(query.into()), requirements: QueryRequirements { - accuracy: AccuracyRequirement::Explicit(AccuracyTarget::Exact), + accuracy: AccuracyRequirement::Explicit(accuracy), ..Default::default() }, predictability: Predictability::Unknown, @@ -91,9 +95,13 @@ fn metric(selector: &QueryExpr) -> String { fn compile_query(query: &str) -> Result { let expression = lower(query); - let dag = fallback_dag(expression.clone()); + compile_dag(&expression, &fallback_dag(expression.clone())) +} + +/// Compile a DAG whose root is the Fallback computing `expression`. +fn compile_dag(expression: &QueryExpr, dag: &PostAsapDag) -> Result { let root = u64::from(dag.root.0); - let inputs = promql_fallback::raw_series(&expression) + let inputs = promql_fallback::raw_series(expression) .map_err(|e| e.to_string())? .into_iter() .enumerate() @@ -104,7 +112,7 @@ fn compile_query(query: &str) -> Result { ) }) .collect(); - let program = compile(&dag, inputs, &[root]).map_err(|e| e.to_string())?; + let program = compile(dag, inputs, &[root]).map_err(|e| e.to_string())?; Ok(serde_json::from_slice(&serde_json::to_vec(&program).unwrap()).unwrap()) } @@ -117,9 +125,19 @@ fn evaluate( at: i64, ) -> Result, i64, f64)>, String> { let expression = lower(query); - let program = compile_query(query)?; + evaluate_dag(&expression, &fallback_dag(expression.clone()), metrics, at) +} + +#[allow(clippy::type_complexity)] +fn evaluate_dag( + expression: &QueryExpr, + dag: &PostAsapDag, + metrics: &[(&str, &[Sample])], + at: i64, +) -> Result, i64, f64)>, String> { + let program = compile_dag(expression, dag)?; let mut sources = BTreeMap::new(); - let selectors = promql_fallback::raw_series(&expression).unwrap(); + let selectors = promql_fallback::raw_series(expression).unwrap(); for (i, (selector, schema)) in selectors.into_iter().enumerate() { let name = metric(&selector); let rows = metrics @@ -128,7 +146,8 @@ fn evaluate( .flat_map(|(_, samples)| samples.iter()) .map(|(spec, seconds, value)| { let mut labels = labels(spec); - labels.insert("__name__".into(), name.clone()); + // A sample may supply its own `__name__`, as a series of another metric. + labels.entry("__name__".into()).or_insert(name.clone()); promql_rows::series_row(&schema, &labels, seconds * 1000, *value).unwrap() }) .collect(); @@ -579,7 +598,6 @@ fn on_and_ignoring_select_the_matching_labels() { assert!(evaluate("a + on(job) b", &[("a", a), ("b", pair)], 60).is_err()); assert!(evaluate("a + on(job) b", &[("a", pair), ("b", b)], 60).is_err()); assert!(labeled("a + on(job) b", &[("a", pair), ("b", other)], 60).is_empty()); - assert!(promql_rows::with_series_identity(&parse("a + on(job) group_left b")).is_err()); } // without() groups by every label except the listed ones and the metric name. @@ -604,6 +622,32 @@ fn without_grouping_drops_labels_and_the_name() { vec![(String::new(), 4.)] ); assert!(labeled("sum without (inst) (a)", &[], 60).is_empty()); + // Series equal without the name share a group rather than colliding. + let named: &[Sample] = &[ + ("job=x,inst=1", 50, 1.), + ("__name__=b,job=x,inst=1", 50, 2.), + ]; + assert_eq!( + labeled("sum without (inst) (a)", &[("a", named)], 60), + vec![("job=x".into(), 3.)] + ); +} + +// Arithmetic with a literal drops the metric name; series that then share a +// label set are an error, as in Prometheus. +#[test] +fn literal_arithmetic_drops_the_name_and_rejects_equal_label_sets() { + let a: &[Sample] = &[ + ("job=x,inst=1", 50, 1.), + ("__name__=b,job=x,inst=2", 50, 2.), + ]; + assert_eq!( + labeled("a * 2", &[("a", a)], 60), + vec![("inst=1,job=x".into(), 2.), ("inst=2,job=x".into(), 4.)] + ); + let equal: &[Sample] = &[("job=x", 50, 1.), ("__name__=b,job=x", 50, 2.)]; + let error = evaluate("a * 2", &[("a", equal)], 60).unwrap_err(); + assert!(error.contains("same labelset"), "{error}"); } // An empty label value is an absent label, and an empty side yields an empty @@ -619,6 +663,581 @@ fn empty_labels_and_empty_sides_match_prometheus() { let pair: &[Sample] = &[("job=x,inst=1", 50, 1.), ("job=x,inst=2", 50, 2.)]; assert!(labeled("a + on(job) b", &[("b", pair)], 60).is_empty()); assert!(labeled("b + on(job) a", &[("b", pair)], 60).is_empty()); - // A non-literal scalar operand has no identity realization yet. - assert!(promql_rows::with_series_identity(&parse("a + scalar(b)")).is_err()); + // `time()` has no row realization yet. + assert!(promql_rows::with_series_identity(&parse("a - time()")).is_err()); +} + +// Sums and averages use Prometheus' Kahan-Neumaier compensation, and an +// average whose running sum overflows switches to an incremental mean. +#[test] +fn sums_and_averages_are_compensated_like_prometheus() { + let cancel = &[("a", 10, 1e100), ("a", 20, 1.), ("a", 30, -1e100)]; + assert_eq!(one("sum_over_time(m[1m])", cancel, 60), 1.); + assert_eq!(one("avg_over_time(m[1m])", cancel, 60), 1. / 3.); + let huge = &[("a", 10, 1.7e308), ("a", 20, 1.7e308)]; + assert_eq!(one("avg_over_time(m[1m])", huge, 60), 1.7e308); + assert_eq!(one("sum_over_time(m[1m])", huge, 60), f64::INFINITY); + let infinite = &[("a", 10, f64::INFINITY), ("a", 20, 1.)]; + assert_eq!(one("sum_over_time(m[1m])", infinite, 60), f64::INFINITY); + assert_eq!(one("avg_over_time(m[1m])", infinite, 60), f64::INFINITY); + let opposite = &[("a", 10, f64::INFINITY), ("a", 20, f64::NEG_INFINITY)]; + assert!(one("sum_over_time(m[1m])", opposite, 60).is_nan()); + assert!(one("avg_over_time(m[1m])", opposite, 60).is_nan()); + let cancel = &[("a", 50, 1e100), ("b", 50, 1.), ("c", 50, -1e100)]; + assert_eq!(one("sum(m)", cancel, 60), 1.); + assert_eq!(one("avg(m)", cancel, 60), 1. / 3.); + let huge = &[("a", 50, 1.7e308), ("b", 50, 1.7e308)]; + assert_eq!(one("avg(m)", huge, 60), 1.7e308); +} + +/// `(k=v,... sorted, value)` rows for readable expectations. +fn rows(pairs: &[(&str, f64)]) -> Vec<(String, f64)> { + let mut rows = pairs + .iter() + .map(|(spec, value)| (spec.to_string(), *value)) + .collect::>(); + rows.sort_by(|a, b| a.0.cmp(&b.0)); + rows +} + +/// `labeled`, with NaN values rendered comparable. +fn labeled_nan(query: &str, metrics: &[(&str, &[Sample])], at: i64) -> Vec<(String, String)> { + labeled(query, metrics, at) + .into_iter() + .map(|(labels, value)| (labels, format!("{value:?}"))) + .collect() +} + +const C: &[Sample] = &[ + ("job=x", 50, 10.), + ("job=y", 50, 20.), + ("job=w", 50, 0.), + ("job=n", 50, f64::NAN), +]; + +// A comparison with a scalar keeps the matching series with their value and +// metric name, whichever side the scalar is on; `bool` yields 1 or 0 for every +// series and drops the name. NaN compares unequal to everything. +#[test] +fn scalar_comparisons_filter_or_return_bool() { + let metrics = &[("a", C)]; + let kept = rows(&[("__name__=a,job=x", 10.), ("__name__=a,job=y", 20.)]); + assert_eq!(labeled("a > 5", metrics, 60), kept); + assert_eq!(labeled("5 < a", metrics, 60), kept); + assert_eq!( + labeled("a <= 10", metrics, 60), + rows(&[("__name__=a,job=w", 0.), ("__name__=a,job=x", 10.)]) + ); + assert_eq!( + labeled("a > bool 5", metrics, 60), + rows(&[("job=n", 0.), ("job=w", 0.), ("job=x", 1.), ("job=y", 1.)]) + ); + assert_eq!( + labeled("10 == bool a", metrics, 60), + rows(&[("job=n", 0.), ("job=w", 0.), ("job=x", 1.), ("job=y", 0.)]) + ); + // scalar() of no series is NaN. + assert_eq!(labeled("a != scalar(b)", metrics, 60).len(), 4); + assert!(labeled("a == scalar(b)", metrics, 60).is_empty()); + assert!(labeled("a > 5", &[], 60).is_empty()); + // Only `bool` drops the name, so only it can make label sets collide. + let equal: &[Sample] = &[("job=x", 50, 1.), ("__name__=b,job=x", 50, 2.)]; + assert_eq!(labeled("a > 0", &[("a", equal)], 60).len(), 2); + let error = evaluate("a > bool 0", &[("a", equal)], 60).unwrap_err(); + assert!(error.contains("same labelset"), "{error}"); +} + +// Vector comparisons match one-to-one like arithmetic. A filter keeps the +// left series, name included, unless `on` reduces its labels; `bool` drops the +// name. A left duplicate is an error only if more than one of it is kept. +#[test] +fn vector_comparisons_match_one_to_one() { + let metrics = &[("a", A), ("b", B)]; + assert_eq!( + labeled("a > b", metrics, 60), + rows(&[("__name__=a,job=x", 10.)]) + ); + assert_eq!( + labeled("a >= b", metrics, 60), + rows(&[("__name__=a,job=w", 0.), ("__name__=a,job=x", 10.)]) + ); + assert_eq!( + labeled("a > bool b", metrics, 60), + rows(&[("job=w", 0.), ("job=x", 1.)]) + ); + assert!(labeled("a < b", metrics, 60).is_empty()); + let a: &[Sample] = &[("job=x,inst=1", 50, 10.)]; + let b: &[Sample] = &[("job=x,inst=2", 50, 4.)]; + let metrics = &[("a", a), ("b", b)]; + assert_eq!( + labeled("a > on(job) b", metrics, 60), + rows(&[("job=x", 10.)]) + ); + assert_eq!( + labeled("a > ignoring(inst) b", metrics, 60), + rows(&[("__name__=a,job=x", 10.)]) + ); + let pair: &[Sample] = &[("job=x,inst=1", 50, 1.), ("job=x,inst=2", 50, 5.)]; + let metrics = &[("a", pair), ("b", b)]; + assert_eq!( + labeled("a > on(job) b", metrics, 60), + rows(&[("job=x", 5.)]) + ); + let error = evaluate("a > bool on(job) b", metrics, 60).unwrap_err(); + assert!(error.contains("many-to-one"), "{error}"); + let nan: &[Sample] = &[("job=x", 50, f64::NAN)]; + let metrics = &[("a", nan), ("b", nan)]; + assert_eq!(labeled("a == bool b", metrics, 60), rows(&[("job=x", 0.)])); + assert_eq!( + labeled_nan("a != b", metrics, 60), + vec![("__name__=a,job=x".into(), "NaN".into())] + ); +} + +const S: &[Sample] = &[ + ("job=x", 50, 1.), + ("job=y", 50, 2.), + ("job=z,inst=1", 50, 3.), +]; +const T: &[Sample] = &[ + ("job=x", 50, 10.), + ("job=w", 50, 20.), + ("job=z,inst=2", 50, 30.), +]; + +// Set operators match label sets many-to-many, ignoring the name by default, +// and return the original series unchanged. +#[test] +fn set_operators_match_label_sets() { + let metrics = &[("a", S), ("b", T)]; + assert_eq!( + labeled("a and b", metrics, 60), + rows(&[("__name__=a,job=x", 1.)]) + ); + assert_eq!( + labeled("a and on(job) b", metrics, 60), + rows(&[("__name__=a,job=x", 1.), ("__name__=a,inst=1,job=z", 3.)]) + ); + assert_eq!( + labeled("a and ignoring(inst) b", metrics, 60), + labeled("a and on(job) b", metrics, 60) + ); + assert_eq!( + labeled("a or b", metrics, 60), + rows(&[ + ("__name__=a,job=x", 1.), + ("__name__=a,job=y", 2.), + ("__name__=a,inst=1,job=z", 3.), + ("__name__=b,job=w", 20.), + ("__name__=b,inst=2,job=z", 30.), + ]) + ); + assert_eq!( + labeled("a or on(job) b", metrics, 60), + rows(&[ + ("__name__=a,job=x", 1.), + ("__name__=a,job=y", 2.), + ("__name__=a,inst=1,job=z", 3.), + ("__name__=b,job=w", 20.), + ]) + ); + assert_eq!( + labeled("a unless b", metrics, 60), + rows(&[("__name__=a,job=y", 2.), ("__name__=a,inst=1,job=z", 3.)]) + ); + assert_eq!( + labeled("a unless on(job) b", metrics, 60), + rows(&[("__name__=a,job=y", 2.)]) + ); + assert_eq!(labeled("a and on() b", metrics, 60).len(), 3); + // Empty sides, and duplicates on either side, which set operators allow. + let a_only = &[("a", S)]; + assert!(labeled("a and b", a_only, 60).is_empty()); + assert_eq!(labeled("a unless b", a_only, 60).len(), 3); + assert_eq!(labeled("b or a", a_only, 60).len(), 3); + let pair: &[Sample] = &[("job=x,inst=1", 50, 1.), ("job=x,inst=2", 50, f64::NAN)]; + assert_eq!( + labeled_nan("a and on(job) b", &[("a", pair), ("b", pair)], 60), + vec![ + ("__name__=a,inst=1,job=x".into(), "1.0".into()), + ("__name__=a,inst=2,job=x".into(), "NaN".into()), + ] + ); +} + +const MANY: &[Sample] = &[ + ("job=x,inst=1", 50, 2.), + ("job=x,inst=2", 50, 3.), + ("job=y,inst=1", 50, 4.), +]; +const ONE: &[Sample] = &[("job=x,team=t1", 50, 10.), ("job=y", 50, 100.)]; + +// group_left/group_right match many series to one; the result keeps the many +// side's labels plus the listed labels of the one side, which a missing label +// removes. A filter keeps the left value. +#[test] +fn group_modifiers_match_many_to_one() { + let metrics = &[("a", MANY), ("info", ONE)]; + assert_eq!( + labeled("a * on(job) group_left(team) info", metrics, 60), + rows(&[ + ("inst=1,job=x,team=t1", 20.), + ("inst=2,job=x,team=t1", 30.), + ("inst=1,job=y", 400.), + ]) + ); + assert_eq!( + labeled("info - on(job) group_right a", metrics, 60), + rows(&[ + ("inst=1,job=x", 8.), + ("inst=2,job=x", 7.), + ("inst=1,job=y", 96.) + ]) + ); + assert_eq!( + labeled("info > on(job) group_right a", metrics, 60), + rows(&[ + ("__name__=a,inst=1,job=x", 10.), + ("__name__=a,inst=2,job=x", 10.), + ("__name__=a,inst=1,job=y", 100.), + ]) + ); + assert_eq!( + labeled("a > bool ignoring(inst, team) group_left info", metrics, 60), + rows(&[ + ("inst=1,job=x", 0.), + ("inst=2,job=x", 0.), + ("inst=1,job=y", 0.) + ]) + ); + // Two "one" series for a match group, or two results with equal labels. + let two: &[Sample] = &[("job=x,team=t1", 50, 1.), ("job=x,team=t2", 50, 2.)]; + let error = evaluate( + "a * on(job) group_left info", + &[("a", MANY), ("info", two)], + 60, + ) + .unwrap_err(); + assert!(error.contains("duplicate series"), "{error}"); + let error = evaluate( + "info * on(job) group_right a", + &[("a", two), ("info", MANY)], + 60, + ) + .unwrap_err(); + assert!(error.contains("left hand-side"), "{error}"); + let named: &[Sample] = &[("job=x", 50, 1.), ("__name__=c,job=x", 50, 2.)]; + let error = evaluate( + "a * on(job) group_left info", + &[("a", named), ("info", ONE)], + 60, + ) + .unwrap_err(); + assert!(error.contains("unique matches"), "{error}"); + assert!(labeled("a * on(job) group_left info", &[("a", MANY)], 60).is_empty()); +} + +// A non-literal scalar applies like a literal; scalar-scalar arithmetic yields +// a scalar; and a literal applies to aggregated rows whose value has another name. +#[test] +fn scalar_operands_and_aggregates() { + let three: &[Sample] = &[("job=b", 50, 3.)]; + let metrics = &[("a", A), ("b", three)]; + assert_eq!( + labeled("a * scalar(b)", metrics, 60), + rows(&[("job=w", 0.), ("job=x", 30.), ("job=y", 60.)]) + ); + assert_eq!( + labeled("a > scalar(b)", metrics, 60), + rows(&[("__name__=a,job=x", 10.), ("__name__=a,job=y", 20.)]) + ); + assert_eq!(labeled("scalar(b) * 2", metrics, 60), rows(&[("", 6.)])); + assert_eq!( + labeled("scalar(b) > bool 2", metrics, 60), + rows(&[("", 1.)]) + ); + // scalar() of several series is NaN. + assert!(labeled("scalar(a) - 1", metrics, 60)[0].1.is_nan()); + assert_eq!( + labeled("sum by (job) (a) * 2", metrics, 60), + rows(&[("job=w", 0.), ("job=x", 20.), ("job=y", 40.)]) + ); + assert_eq!( + labeled("sum by (job) (a) > bool 5", metrics, 60), + rows(&[("job=w", 0.), ("job=x", 1.), ("job=y", 1.)]) + ); +} + +// Range functions other than last_over_time drop the metric name, so series +// that then share a label set are an error, as in Prometheus. +#[test] +fn range_functions_drop_the_name_and_reject_equal_label_sets() { + let equal: &[Sample] = &[ + ("job=x", 10, 1.), + ("job=x", 50, 2.), + ("__name__=b,job=x", 10, 1.), + ("__name__=b,job=x", 50, 4.), + ]; + let error = evaluate("rate(a[1m])", &[("a", equal)], 60).unwrap_err(); + assert!(error.contains("same labelset"), "{error}"); + assert_eq!( + labeled("last_over_time(a[1m])", &[("a", equal)], 60), + rows(&[("__name__=a,job=x", 2.), ("__name__=b,job=x", 4.)]) + ); + assert_eq!( + labeled("max_over_time(a[1m])", &[("a", &equal[..2])], 60), + rows(&[("job=x", 2.)]) + ); +} + +// Scalar-valued expressions are scalars too; `or vector(0)` fills an empty +// aggregate; a range function inside a subquery drops the name. +#[test] +fn scalar_expressions_or_vector_and_subquery_names() { + let three: &[Sample] = &[("job=b", 50, 3.)]; + let metrics = &[("a", A), ("b", three)]; + assert_eq!( + labeled("a + (scalar(b) * 2)", metrics, 60), + rows(&[("job=w", 6.), ("job=x", 16.), ("job=y", 26.)]) + ); + assert_eq!( + labeled("a + -scalar(b)", metrics, 60), + rows(&[("job=w", -3.), ("job=x", 7.), ("job=y", 17.)]) + ); + assert_eq!( + labeled("sum(a) or vector(0)", metrics, 60), + rows(&[("", 30.)]) + ); + assert_eq!(labeled("sum(a) or vector(0)", &[], 60), rows(&[("", 0.)])); + let counter: &[Sample] = &[("job=x", 0, 0.), ("job=x", 30, 3.), ("job=x", 60, 6.)]; + let result = labeled("last_over_time(rate(a[1m])[2m:1m])", &[("a", counter)], 60); + assert_eq!(result.len(), 1); + assert_eq!(result[0].0, "job=x"); +} + +/// Instant `x_bucket` samples at 50s: `(labels without le, [(le, count)])`. +fn buckets(series: &[(&'static str, &[(&'static str, f64)])]) -> Vec { + series + .iter() + .flat_map(|(labels, buckets)| { + buckets.iter().map(move |(le, count)| { + let spec = if labels.is_empty() { + format!("le={le}") + } else { + format!("{labels},le={le}") + }; + (&*Box::leak(spec.into_boxed_str()), 50, *count) + }) + }) + .collect() +} + +fn quantile(query: &str, samples: &[Sample]) -> Vec<(String, f64)> { + labeled(query, &[("x_bucket", samples)], 60) +} + +const HISTOGRAM: &[(&str, f64)] = &[("1", 2.), ("2", 6.), ("4", 8.), ("+Inf", 10.)]; + +// histogram_quantile interpolates linearly within the bucket holding rank q·count, +// returns the highest finite bound for the +Inf bucket, and maps q outside +// [0, 1] to ∓Inf and a NaN q to NaN. Output labels drop le and __name__. +#[test] +fn histogram_quantile_interpolates_classic_buckets() { + let samples = buckets(&[("job=a", HISTOGRAM)]); + for (q, expected) in [ + ("0", 0.), + ("0.1", 0.5), + ("0.5", 1.75), + ("0.9", 4.), + ("1", 4.), + ("-0.5", f64::NEG_INFINITY), + ("1.5", f64::INFINITY), + ] { + let query = format!("histogram_quantile({q}, x_bucket)"); + assert_eq!( + quantile(&query, &samples), + vec![("job=a".into(), expected)], + "{query}" + ); + } + let nan = quantile("histogram_quantile(NaN, x_bucket)", &samples); + assert!(matches!(nan.as_slice(), [(labels, v)] if labels == "job=a" && v.is_nan())); +} + +// Each label set other than le is its own histogram. Degenerate histograms +// yield NaN: no +Inf bucket, fewer than two buckets, or zero observations. +#[test] +fn histogram_quantile_groups_series_and_rejects_degenerate_histograms() { + let samples = buckets(&[ + ("job=a", HISTOGRAM), + ("job=b", &[("1", 1.), ("2", 2.)]), + ("job=c", &[("+Inf", 5.)]), + ("job=d", &[("1", 0.), ("+Inf", 0.)]), + ("job=e,inst=1", HISTOGRAM), + ]); + let rows = quantile("histogram_quantile(0.5, x_bucket)", &samples); + let labels: Vec<_> = rows.iter().map(|(l, _)| l.as_str()).collect(); + assert_eq!( + labels, + vec!["inst=1,job=e", "job=a", "job=b", "job=c", "job=d"] + ); + assert_eq!(rows[0].1, 1.75); + assert_eq!(rows[1].1, 1.75); + assert!(rows[2..].iter().all(|(_, v)| v.is_nan()), "{rows:?}"); +} + +// Buckets sort by bound, unparsable or missing le values are skipped, equal +// bounds merge, and decreasing cumulative counts are raised to be monotonic. +#[test] +fn histogram_quantile_normalizes_buckets_like_prometheus() { + let unordered = buckets(&[("job=a", &[("+Inf", 10.), ("4", 8.), ("1", 2.), ("2", 6.)])]); + assert_eq!( + quantile("histogram_quantile(0.5, x_bucket)", &unordered), + vec![("job=a".into(), 1.75)] + ); + let mut invalid = buckets(&[("job=a", &[("abc", 100.), ("1", 2.), ("+Inf", 4.)])]); + invalid.push(("job=a", 50, 100.)); + assert_eq!( + quantile("histogram_quantile(0.5, x_bucket)", &invalid), + vec![("job=a".into(), 1.)] + ); + // Go's ParseFloat rejects an out-of-range bound rather than rounding it to +Inf. + let overflow = buckets(&[("job=a", &[("1", 1.), ("1e400", 2.)])]); + let rows = quantile("histogram_quantile(0.5, x_bucket)", &overflow); + assert!( + matches!(rows.as_slice(), [(_, v)] if v.is_nan()), + "{rows:?}" + ); + let duplicate = buckets(&[("job=a", &[("1", 1.), ("1.0", 1.), ("+Inf", 4.)])]); + assert_eq!( + quantile("histogram_quantile(0.5, x_bucket)", &duplicate), + vec![("job=a".into(), 1.)] + ); + // Counts [6, 2→6, 8, 8]: rank 7 lies in (2, 4], 1 of its 2 observations in. + let decreasing = buckets(&[("job=a", &[("1", 6.), ("2", 2.), ("4", 8.), ("+Inf", 8.)])]); + assert_eq!( + quantile("histogram_quantile(0.875, x_bucket)", &decreasing), + vec![("job=a".into(), 3.)] + ); +} + +// A lowest bucket with a non-positive bound is returned as is, not +// interpolated from zero. +#[test] +fn histogram_quantile_non_positive_lowest_bucket() { + let samples = buckets(&[("job=a", &[("-1", 2.), ("1", 4.), ("+Inf", 4.)])]); + for (q, expected) in [("0.25", -1.), ("0.75", 0.)] { + let query = format!("histogram_quantile({q}, x_bucket)"); + assert_eq!( + quantile(&query, &samples), + vec![("job=a".into(), expected)], + "{query}" + ); + } +} + +// The common shapes: an aggregated rate keeps its by labels other than le, and +// a per-series rate keeps every label but le and __name__. +#[test] +fn histogram_quantile_over_rates_and_sums() { + // Each counter grows by c per minute, so its rate is c/60. + let counter = |labels: &'static str, le: &str, c: f64| { + let spec: &'static str = Box::leak(format!("{labels},le={le}").into_boxed_str()); + (60..=240) + .step_by(60) + .map(move |t| (spec, t as i64, c * (t / 60) as f64)) + .collect::>() + }; + let mut samples = Vec::new(); + for inst in ["job=a,inst=1", "job=a,inst=2"] { + for (le, count) in HISTOGRAM { + samples.extend(counter(inst, le, *count)); + } + } + let metrics = &[("x_bucket", samples.as_slice())]; + let close = |rows: Vec<(String, f64)>, expected: &[(&str, f64)]| { + assert_eq!(rows.len(), expected.len(), "{rows:?}"); + for ((labels, v), (want, w)) in rows.iter().zip(expected) { + assert_eq!(labels, want); + assert!((v - w).abs() < 1e-9, "{labels}: {v} vs {w}"); + } + }; + close( + labeled( + "histogram_quantile(0.5, sum by (le, job) (rate(x_bucket[5m])))", + metrics, + 300, + ), + &[("job=a", 1.75)], + ); + close( + labeled( + "histogram_quantile(0.5, sum by (le) (x_bucket))", + metrics, + 250, + ), + &[("", 1.75)], + ); + close( + labeled("histogram_quantile(0.5, rate(x_bucket[5m]))", metrics, 300), + &[("inst=1,job=a", 1.75), ("inst=2,job=a", 1.75)], + ); +} + +// Histograms that differ only in __name__ collide once it is dropped, which +// Prometheus reports as an error rather than merging them. +#[test] +fn histogram_quantile_rejects_equal_output_label_sets() { + let mut samples = buckets(&[("job=a", HISTOGRAM)]); + samples.extend(buckets(&[("__name__=y_bucket,job=a", HISTOGRAM)])); + let error = evaluate( + "histogram_quantile(0.5, x_bucket)", + &[("x_bucket", &samples)], + 60, + ) + .unwrap_err(); + assert!(error.contains("same labelset"), "{error}"); +} + +// Candidate search keeps a classic histogram_quantile whole and exact, even +// for an approximate target, and the selected DAG compiles and executes. +#[test] +fn histogram_quantile_selection_keeps_the_exact_fallback() { + use asap_aware_mapping::{ + accuracy::DefaultAccuracyModel, cost_model::DefaultCostModel, default_strategies, + search_workload_with_targets, Replacement, + }; + let samples = buckets(&[("job=a", HISTOGRAM)]); + for target in [AccuracyTarget::Exact, AccuracyTarget::Epsilon(0.01)] { + for query in [ + "histogram_quantile(0.5, x_bucket)", + "histogram_quantile(0.5, sum by (le, job) (x_bucket))", + ] { + let root = Rc::new( + promql_rows::with_series_identity(&parse_with(query, target.clone())).unwrap(), + ); + let space = search_workload_with_targets( + vec![(query, root.clone(), Some(target.clone()))], + &default_strategies(), + &DefaultAccuracyModel, + ); + let planned = &space.roots[0].1; + let candidates = &space.candidates_for_target(planned).unwrap().candidates; + assert!( + candidates.iter().all(|c| matches!(&c.replacement, + Replacement::Summary(node) if matches!(&node.expr, + SummaryExpr::KeepPreAsap(e) if **e == *root))), + "{query}: {candidates:?}" + ); + let selected = space + .global_selection(&DefaultCostModel) + .assemble_selected_dag(planned) + .unwrap() + .unwrap(); + let dag = compile_post_asap_dag(&selected).unwrap(); + let rows = evaluate_dag(&root, &dag, &[("x_bucket", &samples)], 60).unwrap(); + let values: Vec<_> = rows.iter().map(|(_, _, v)| *v).collect(); + assert_eq!(values, vec![1.75], "{query} {target:?}"); + } + } } diff --git a/crates/frontend-promql/src/promql.rs b/crates/frontend-promql/src/promql.rs index b71a8c1e..920acfa7 100644 --- a/crates/frontend-promql/src/promql.rs +++ b/crates/frontend-promql/src/promql.rs @@ -22,7 +22,7 @@ //! | PromQL | Canonical shape | //! |---|---| //! | `quantile_over_time(φ, m{f}[w])` | `Aggregate{[Quantile(φ)], TimeRange{w, Scan{predicates}}}` | -//! | `histogram_quantile(φ, )` | `Aggregate{[HistogramQuantile(φ)]}` — cumulative-bucket interpolation (classic form recognised by `by (le)` / a `_bucket` metric / an `le` matcher) | +//! | `histogram_quantile(φ, )` | `Aggregate{without(le), [HistogramQuantile(φ, le)]}` — cumulative-bucket interpolation (classic form recognised by `by (le)` / a `_bucket` metric / an `le` matcher) | //! | `histogram_quantile(φ, )` | `Aggregate{[Quantile(φ)]}` over the fully-lowered arg (generic, sketch-able with an accuracy target) | //! | `histogram_quantiles(v, "l", φ…)` | `Concat{PromqlRelabel{l=φᵢ, }…}` — one branch per φ (issue #109) | //! | `histogram_count/sum/avg/stddev/stdvar(v)`, `histogram_fraction(l,u,v)` | `Aggregate{[Histogram*]}` — per-series native-histogram accessors (issue #43) | @@ -668,18 +668,12 @@ fn build_over_subtree(outer: Outer, keys: Vec, child: Unresolved) -> }) } -/// `histogram_quantile(φ, )` lowers `` in full — preserving any -/// `sum by (le)` / `rate` structure inside it — and wraps the result in an -/// `Aggregate{[Quantile(φ)]}`. The φ-quantile reduces across the `le` buckets, -/// so the wrapper carries no grouping keys: the usage-derived schema can't -/// enumerate the non-`le` labels to group by (the same limitation that rejects -/// `without`). This handles the canonical -/// `histogram_quantile(φ, sum by (le) (rate(m_bucket[w])))` pattern, which the -/// old "extract the matrix and substitute a bare Quantile" path could not. /// The `histogram_*` function family (issues #43, histogram_quantile). /// -/// `histogram_quantile(φ, )` lowers to a `Quantile` over the fully-lowered -/// argument — it also covers the classic `le`-bucket form (`sum by (le) (…)`). +/// `histogram_quantile(φ, )` lowers `` in full — preserving any +/// `sum by (le)` / `rate` structure inside it. The classic `le`-bucket form +/// becomes [`classic_histogram_quantile`]; a native histogram or raw samples +/// become a `Quantile` over the whole argument. /// The native-histogram accessors (`histogram_count`/`sum`/`avg`/`stddev`/ /// `stdvar`/`fraction`) each extract one float per series, lowering to a /// per-series `Aggregate{[accessor]}` directly over the (instant) argument. @@ -700,14 +694,13 @@ fn walk_histogram(call: &Call) -> Result { // The true signal is the argument's sample type: a declared // `HistogramKind` (issue #79) drives the choice when available, else we // fall back to the structural `by (le)`/`_bucket` heuristic (issue #43). - let func = if histogram_arg_is_sketchable(arg_expr) { - AggIntent::Quantile { - col: None, - q: phi, - accuracy: current_accuracy(), - } - } else { - AggIntent::HistogramQuantile { q: phi } + if !histogram_arg_is_sketchable(arg_expr) { + return Ok(classic_histogram_quantile(phi, "", walk(arg_expr)?)); + } + let func = AggIntent::Quantile { + col: None, + q: phi, + accuracy: current_accuracy(), }; return Ok(outer_aggregate(vec![], func, walk(arg_expr)?)); } @@ -730,6 +723,21 @@ fn walk_histogram(call: &Call) -> Result { Ok(outer_aggregate(vec![], func, walk(arg(call, vec_idx)?)?)) } +/// Classic-bucket `histogram_quantile(φ, child)`. One histogram is the set of +/// series that differ only in `le`, so the aggregate groups `without (le)`. +/// That grouping also seeds `le` into a usage-derived source schema, even +/// when no matcher names it. An empty `output_name` keeps the intent-keyed name. +fn classic_histogram_quantile(q: f64, output_name: &str, child: Unresolved) -> Unresolved { + let le = ColumnRef::Named("le".into()); + Unresolved::Aggregate { + reduction: Reduction::Reduce(GroupKeys::without(vec![le.clone()])), + measures: vec![AggIntent::HistogramQuantile { q, le }], + output_names: vec![output_name.into()], + having: None, + child: Rc::new(child), + } +} + /// `histogram_quantiles(v, "label", φ₀, φ₁, …)` — the experimental multi-quantile /// form (issue #109). It is `histogram_quantile(φᵢ, v)` fanned out over the /// quantiles, each branch's output series tagged with `label = φᵢ`. @@ -762,26 +770,25 @@ fn walk_histogram_quantiles(call: &Call) -> Result { let branches = (2..call.args.args.len()) .map(|i| { let phi = bounded_quantile_param(num_arg(call, i)?)?; - let intent = if sketchable { - AggIntent::Quantile { + let child = walk(vec_expr)?; + // Each branch aliases its value column to "value" (not the + // intent-keyed default) so `Concat` — which derives its schema + // from the first branch — doesn't silently misdescribe the rest. + let quantile = if sketchable { + let intent = AggIntent::Quantile { col: None, q: phi, accuracy: current_accuracy(), + }; + Unresolved::Aggregate { + reduction: reduction_for(&[], intent.is_per_series()), + measures: vec![intent], + output_names: vec!["value".into()], + having: None, + child: Rc::new(child), } } else { - AggIntent::HistogramQuantile { q: phi } - }; - let child = walk(vec_expr)?; - let reduction = reduction_for(&[], intent.is_per_series()); - let quantile = Unresolved::Aggregate { - reduction, - measures: vec![intent], - // Each branch aliases its value column to "value" (not the - // intent-keyed default) so `Concat` — which derives its schema - // from the first branch — doesn't silently misdescribe the rest. - output_names: vec!["value".into()], - having: None, - child: Rc::new(child), + classic_histogram_quantile(phi, "value", child) }; Ok(Unresolved::PromqlRelabel { dst: label.clone(), @@ -1274,7 +1281,20 @@ fn selector_is_bucket(vs: &VectorSelector) -> bool { fn walk_binary(bin: &BinaryExpr) -> Result { let lhs = scalar_or_vector(&bin.lhs)?; let rhs = scalar_or_vector(&bin.rhs)?; - let op = binop(bin.op.id())?; + // `VectorMatch` has no fill field; dropping fill would change which series + // are emitted and their values, so the query must fall back to exact + // execution instead. + if let Some(m) = &bin.modifier { + if m.fill_values.lhs.is_some() || m.fill_values.rhs.is_some() { + return Err(LoweringError::UnsupportedFeature(format!( + "`fill` vector-matching modifier: `{bin}`" + ))); + } + } + let op = match (binop(bin.op.id())?, bin.return_bool()) { + (BinaryOpKind::Compare(op), true) => BinaryOpKind::CompareBool(op), + (op, _) => op, + }; let vector_match = bin.modifier.as_ref().map(|m| { let (kind, labels) = match &m.matching { Some(LabelModifier::Include(ls)) => (VectorMatchKind::On, ls.labels.clone()), diff --git a/crates/frontend-promql/tests/observability/awesome_prometheus_alerts.rs b/crates/frontend-promql/tests/observability/awesome_prometheus_alerts.rs index b2592925..12afb94d 100644 --- a/crates/frontend-promql/tests/observability/awesome_prometheus_alerts.rs +++ b/crates/frontend-promql/tests/observability/awesome_prometheus_alerts.rs @@ -213,7 +213,7 @@ fn histogram_quantile_core_lowers() { ok("histogram_quantile(0.95, sum(rate(http_request_duration_seconds_bucket[5m])) by (le))"); assert!(has( &qe, - |i| matches!(i, AggIntent::HistogramQuantile { q } if (*q - 0.95).abs() < 1e-9) + |i| matches!(i, AggIntent::HistogramQuantile { q, .. } if (*q - 0.95).abs() < 1e-9) )); assert!(has(&qe, |i| matches!(i, AggIntent::Sum { .. }))); assert!(has(&qe, |i| matches!(i, AggIntent::Rate))); diff --git a/crates/frontend-promql/tests/observability/promql_corpus.rs b/crates/frontend-promql/tests/observability/promql_corpus.rs index beb16d61..f265edb4 100644 --- a/crates/frontend-promql/tests/observability/promql_corpus.rs +++ b/crates/frontend-promql/tests/observability/promql_corpus.rs @@ -155,13 +155,16 @@ fn lowering_is_total_over_the_entire_corpus() { // rather than pin an exact count — ratchet them up as coverage lands. // // The 235 unparseable are parser-fork gaps (issue #108); the rejections are - // lowering gaps (#109). Both shrink over time, so these floors only ever rise. + // lowering gaps (#109). Both shrink over time, so these floors normally only + // rise. Exception: the testdata floor was lowered to the measured 1485 when + // the 44 `fill` vector-matching queries became rejected rather than + // silently lowered without their fill semantics. assert!( docs.lowered >= 47, "docs lowering coverage regressed: {docs:?}" ); assert!( - td.lowered >= 1495, + td.lowered >= 1485, "testdata lowering coverage regressed: {td:?}" ); } diff --git a/crates/frontend-promql/tests/promql_conformance.rs b/crates/frontend-promql/tests/promql_conformance.rs index 37712c8d..69fbb223 100644 --- a/crates/frontend-promql/tests/promql_conformance.rs +++ b/crates/frontend-promql/tests/promql_conformance.rs @@ -540,7 +540,7 @@ fn histogram_quantile_over_rate() { panic!("expected Aggregate{{HistogramQuantile}}, got {qe:?}"); }; assert!( - matches!(measures.as_slice(), [AggIntent::HistogramQuantile { q }] if (*q - 0.9).abs() < 1e-9) + matches!(measures.as_slice(), [AggIntent::HistogramQuantile { q, .. }] if (*q - 0.9).abs() < 1e-9) ); assert!(has(&qe, |i| matches!(i, AggIntent::Rate))); } @@ -1582,7 +1582,7 @@ fn histogram_quantile_classic_bucket_vs_native() { assert!( has( &qe, - |i| matches!(i, AggIntent::HistogramQuantile { q } if (*q - 0.9).abs() < 1e-9) + |i| matches!(i, AggIntent::HistogramQuantile { q, .. } if (*q - 0.9).abs() < 1e-9) ), "classic bucket form → HistogramQuantile: {classic}" ); diff --git a/crates/frontend-promql/tests/promql_lowering.rs b/crates/frontend-promql/tests/promql_lowering.rs index 3c92b7ba..9d3b80de 100644 --- a/crates/frontend-promql/tests/promql_lowering.rs +++ b/crates/frontend-promql/tests/promql_lowering.rs @@ -237,7 +237,7 @@ fn histogram_quantile_wraps_inner_in_quantile() { panic!("expected outer Aggregate{{HistogramQuantile}}, got {qe:?}"); }; assert!( - matches!(measures.as_slice(), [AggIntent::HistogramQuantile { q }] if (*q - 0.95).abs() < 1e-9) + matches!(measures.as_slice(), [AggIntent::HistogramQuantile { q, .. }] if (*q - 0.95).abs() < 1e-9) ); let QueryExpr::Aggregate { measures, child, .. @@ -274,7 +274,7 @@ fn histogram_quantile_over_sum_by_le_preserves_grouping() { }; // The `by (le)` grouping marks the classic cumulative-bucket form. assert!( - matches!(measures.as_slice(), [AggIntent::HistogramQuantile { q }] if (*q - 0.99).abs() < 1e-9) + matches!(measures.as_slice(), [AggIntent::HistogramQuantile { q, .. }] if (*q - 0.99).abs() < 1e-9) ); // `sum by (le)` survives as a positional Aggregate (by = [2], `le`) over the // inner Rate — no name-based Partition. @@ -290,6 +290,90 @@ fn histogram_quantile_over_sum_by_le_preserves_grouping() { assert!(matches!(measures.as_slice(), [AggIntent::Sum { .. }])); } +/// The classic `histogram_quantile` aggregate: its `without` keys, `le` +/// column, and output column names. +fn classic_histogram(qe: &QueryExpr) -> (Vec, usize, Vec) { + let QueryExpr::Aggregate { + reduction: Reduction::Reduce(by), + measures, + .. + } = qe + else { + panic!("expected a reducing Aggregate, got {qe:?}"); + }; + let [AggIntent::HistogramQuantile { le, .. }] = measures.as_slice() else { + panic!("expected HistogramQuantile, got {measures:?}"); + }; + assert!(by.is_without(), "histogram_quantile groups without (le)"); + let names = qe + .output_schema() + .unwrap() + .columns + .iter() + .map(|c| c.name.clone()) + .collect(); + (by.keys().to_vec(), *le, names) +} + +// A classic histogram_quantile groups `without (le)` and names the child's +// `le` column, even when no matcher or grouping mentions `le`. +#[test] +fn classic_histogram_quantile_groups_without_le() { + let qe = lower("histogram_quantile(0.9, rate(http_duration_seconds_bucket[5m]))"); + let (keys, le, names) = classic_histogram(&qe); + let QueryExpr::Aggregate { child, .. } = &qe else { + unreachable!() + }; + let child = child.output_schema().unwrap(); + assert_eq!(child.columns[le].name, "le"); + assert_eq!(keys, vec![le]); + assert_eq!(names, vec!["histogram_quantile"]); +} + +// An explicit `sum by (le, job)` argument keeps `job` and drops `le` and the +// renamed sample value from the output labels. +#[test] +fn classic_histogram_quantile_over_sum_by_keeps_other_labels() { + let qe = + lower("histogram_quantile(0.9, sum by (le, job) (rate(http_duration_seconds_bucket[5m])))"); + let (keys, le, names) = classic_histogram(&qe); + // `sum by (le, job)` outputs `[job, le, sum]`. + assert_eq!((keys, le), (vec![1], 1)); + assert_eq!(names, vec!["job", "histogram_quantile"]); +} + +// Out-of-range and NaN quantiles lower unchanged; execution returns -Inf/+Inf/NaN. +#[test] +fn classic_histogram_quantile_keeps_out_of_range_quantiles() { + for (query, expected) in [ + ("histogram_quantile(-1, x_bucket)", -1.), + ("histogram_quantile(2, x_bucket)", 2.), + ] { + let QueryExpr::Aggregate { measures, .. } = lower(query) else { + panic!("{query}"); + }; + assert!( + matches!(measures.as_slice(), [AggIntent::HistogramQuantile { q, .. }] if *q == expected) + ); + } + let QueryExpr::Aggregate { measures, .. } = lower("histogram_quantile(NaN, x_bucket)") else { + panic!("NaN"); + }; + assert!(matches!(measures.as_slice(), [AggIntent::HistogramQuantile { q, .. }] if q.is_nan())); +} + +// An argument whose closed output lacks `le` has no buckets. Prometheus +// returns an empty vector; lowering rejects it rather than guess a column. +#[test] +fn classic_histogram_quantile_rejects_an_argument_without_le() { + let error = lower_promql( + "histogram_quantile(0.9, sum by (job) (rate(x_bucket[5m])))", + AccuracyTarget::Exact, + ) + .unwrap_err(); + assert!(error.to_string().contains("le"), "{error}"); +} + // ── rate / increase carry their own window (no Window node) ───────────────────── #[test] @@ -705,6 +789,25 @@ fn binary_op_with_on_grouping() { assert_eq!(vm.labels, vec!["host".to_string()]); } +// `bool` changes a comparison from a filter to a 0/1 result, so the IR must +// carry it. +#[test] +fn bool_comparisons_are_distinct() { + let op = |q: &str| match lower(q) { + QueryExpr::BinaryOp { op, .. } => op, + other => panic!("expected BinaryOp, got {other:?}"), + }; + assert_eq!(op("a > 1"), BinaryOpKind::Compare(CompareOpKind::Gt)); + assert_eq!( + op("a > bool 1"), + BinaryOpKind::CompareBool(CompareOpKind::Gt) + ); + assert_eq!( + op("a == bool on(job) b"), + BinaryOpKind::CompareBool(CompareOpKind::Eq) + ); +} + #[test] fn binary_op_binds_each_branch_against_its_own_schema() { // Each side scans a different metric and groups by a different label. With a @@ -845,6 +948,28 @@ fn pathologically_nested_query_is_rejected_not_stack_overflow() { assert!(format!("{err}").contains("nesting"), "got {err}"); } +// Behavior: every parser-accepted `fill` modifier form is rejected with a +// fill-specific lowering error rather than silently dropped. +#[test] +fn fill_modifiers_are_rejected_not_ignored() { + for q in [ + "a + fill(0) b", + "a + fill_left(1) b", + "a + fill_right(2) b", + "a + fill_left(1) fill_right(2) b", + "a + fill_right(2) fill_left(1) b", + "a + on(job) fill(0) b", + "a * ignoring(instance) group_left(env) fill_right(0) b", + "a > bool on(job) fill(0) b", + "sum(a - on(job) group_right fill_left(0) b)", + ] { + match lower_promql(q, AccuracyTarget::Exact) { + Err(LoweringError::UnsupportedFeature(m)) if m.contains("`fill`") => {} + other => panic!("expected fill rejection for {q:?}, got {other:?}"), + } + } +} + // ── accuracy propagation ────────────────────────────────────────────────────── #[test] diff --git a/crates/frontend-sql/src/sql/mod.rs b/crates/frontend-sql/src/sql/mod.rs index 9ebf7f4a..d18e6821 100644 --- a/crates/frontend-sql/src/sql/mod.rs +++ b/crates/frontend-sql/src/sql/mod.rs @@ -702,8 +702,12 @@ impl<'a> SqlLowerer<'a> { } } PlanningBridge::HistogramQuantile { q } => Unresolved::Aggregate { + // The marker is the projection's only column: one histogram. reduction: Reduction::Reduce(GroupKeys::none()), - measures: vec![AggIntent::HistogramQuantile { q }], + measures: vec![AggIntent::HistogramQuantile { + q, + le: ColumnRef::Named("le".into()), + }], output_names: vec!["value".into()], having: None, child: Rc::new(input), diff --git a/crates/frontend-sql/tests/sql_lowering.rs b/crates/frontend-sql/tests/sql_lowering.rs index 5482e6f6..c61ddcd4 100644 --- a/crates/frontend-sql/tests/sql_lowering.rs +++ b/crates/frontend-sql/tests/sql_lowering.rs @@ -92,10 +92,11 @@ async fn planning_histogram_bridge_reuses_classic_bucket_intent() { else { panic!("expected canonical histogram aggregate"); }; + // One histogram over all rows; the bucket bound is the child's column 0. assert!(reduction.expect_reduce().keys().is_empty()); assert!(matches!( measures.as_slice(), - [AggIntent::HistogramQuantile { q }] if (*q - 0.95).abs() < 1e-12 + [AggIntent::HistogramQuantile { q, le: 0 }] if (*q - 0.95).abs() < 1e-12 )); assert!(matches!(child.as_ref(), QueryExpr::Project { .. })); } diff --git a/crates/integration-tests/tests/promql_to_post_asap.rs b/crates/integration-tests/tests/promql_to_post_asap.rs index 539db7bf..a0e71a08 100644 --- a/crates/integration-tests/tests/promql_to_post_asap.rs +++ b/crates/integration-tests/tests/promql_to_post_asap.rs @@ -1436,3 +1436,30 @@ fn ddsketch_ratio_requires_a_supported_population_size() { assert!(strategy.replacements(&TargetSubDAG::new(&pre)).is_empty()); } } + +// Every `without` aggregation candidate exports a valid DAG: its summary state +// column carries the family instead of the readout's Float64 value. +#[test] +fn without_aggregation_candidates_export_valid_dags() { + for accuracy in [ + AccuracyTarget::Exact, + AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }, + ] { + for query in ["sum without (pod) (m)", "quantile without (pod) (0.5, m)"] { + let root = Rc::new(lower_promql(query, accuracy.clone()).unwrap()); + let space = search_workload_with_targets( + vec![(0, root, Some(accuracy.clone()))], + &asap_aware_mapping::default_strategies(), + &DefaultAccuracyModel, + ); + let inventory = space.enumerate_candidate_dags_for_root(&0, 65_536).unwrap(); + assert!(!inventory.candidates.is_empty(), "{query}"); + for (_, node) in inventory.candidates.iter().flatten() { + compile_post_asap_dag(node).unwrap_or_else(|e| panic!("{query}: {e}")); + } + } + } +} diff --git a/crates/types/src/pre_asap/agg_intent.rs b/crates/types/src/pre_asap/agg_intent.rs index bead9d51..73243e6e 100644 --- a/crates/types/src/pre_asap/agg_intent.rs +++ b/crates/types/src/pre_asap/agg_intent.rs @@ -193,8 +193,14 @@ pub enum AggIntent { /// to: this is exact bucket interpolation, not a sketch-able quantile, so it /// carries no accuracy target and is a cross-series reduction over `le` /// (issue #43). + /// + /// PromQL groups the enclosing `Aggregate` `without([le])`: one histogram + /// is the set of series that differ only in `le`. The output drops `le` + /// and `__name__`. `le` names the bucket-bound column so that execution + /// need not guess it from the grouping keys. HistogramQuantile { q: f64, + le: C, }, /// A per-sample element-wise math / trig transform (issue #45) — `abs`, diff --git a/crates/types/src/pre_asap/query_expr.rs b/crates/types/src/pre_asap/query_expr.rs index 5b63e5a2..ff3e6e94 100644 --- a/crates/types/src/pre_asap/query_expr.rs +++ b/crates/types/src/pre_asap/query_expr.rs @@ -254,8 +254,13 @@ pub enum BinaryOpKind { /// Arithmetic — `Add/Sub/Mul/Div/Mod` (shared with `QueryExpr::Arithmetic`). Arithmetic(ArithmeticOpKind), /// Comparison — `Eq/Ne/Lt/Le/Gt/Ge` + `Like/ILike/Regex` family (shared - /// with `QueryExpr::Compare`). + /// with `QueryExpr::Compare`). PromQL keeps the matched series whose + /// comparison holds. Compare(CompareOpKind), + /// PromQL comparison with the `bool` modifier: every matched series + /// yields 1 or 0 and loses its metric name. A separate variant, not a + /// flag, because only comparisons take `bool`. + CompareBool(CompareOpKind), /// PromQL vector-set operation. Set(PromQLVectorSetOpKind), } @@ -265,6 +270,7 @@ impl std::fmt::Display for BinaryOpKind { match self { BinaryOpKind::Arithmetic(op) => write!(f, "{op}"), BinaryOpKind::Compare(op) => write!(f, "{op}"), + BinaryOpKind::CompareBool(op) => write!(f, "{op} bool"), BinaryOpKind::Set(PromQLVectorSetOpKind::And) => f.write_str("AND"), BinaryOpKind::Set(PromQLVectorSetOpKind::Or) => f.write_str("OR"), BinaryOpKind::Set(PromQLVectorSetOpKind::Unless) => f.write_str("unless"), @@ -1634,16 +1640,18 @@ fn without_output_schema( )); } } + // A nested aggregate renames the sample value (`sum by (le) (…)` → `sum`); + // it is still the value, not a kept label. + let value = + super::column_resolution::resolve_column_ref(&ColumnRef::SampleValue, in_schema).ok(); let mut out_cols: Vec = Vec::new(); for (i, col) in in_schema.columns.iter().enumerate() { let is_time = in_schema.time_index == Some(i); - let is_value = col.name == "value"; - if !is_time && !is_value && !excluded.contains(&i) { + if !is_time && value != Some(i) && !excluded.contains(&i) { out_cols.push(col.clone()); } } - let probe = in_schema - .column_id("value") + let probe = value .and_then(|i| in_schema.columns.get(i)) .cloned() .unwrap_or_else(|| Column::new("value", DataType::Float64, false)); @@ -2252,6 +2260,31 @@ mod tests { assert!(s.unique_keys.is_empty(), "kept set unknown → no unique key"); } + // A nested aggregate's renamed sample value is not a kept label. + #[test] + fn without_aggregate_drops_a_renamed_sample_value() { + // `sum without (inst) (sum by (inst, job) (m))` over `[inst, job, sum]`. + let inner = QueryExpr::Scan { + source: Source::TimeSeries { metric: "m".into() }, + predicates: vec![], + schema: Schema::new(vec![ + col("inst", DataType::Utf8, true), + col("job", DataType::Utf8, true), + col("sum", DataType::Float64, false), + ]), + }; + let agg = QueryExpr::Aggregate { + reduction: Reduction::Reduce(GroupKeys::without(vec![0])), + measures: vec![AggIntent::Sum { col: None }], + output_names: vec![], + having: None, + child: Rc::new(inner), + }; + let s = agg.output_schema().unwrap(); + let names: Vec<_> = s.columns.iter().map(|c| c.name.as_str()).collect(); + assert_eq!(names, vec!["job", "sum"]); + } + #[test] fn time_shift_is_schema_pass_through() { // `offset`/`@` move *when* a selector is evaluated, never its columns — diff --git a/crates/types/src/pre_asap/resolve.rs b/crates/types/src/pre_asap/resolve.rs index bec0331d..2ad60eb6 100644 --- a/crates/types/src/pre_asap/resolve.rs +++ b/crates/types/src/pre_asap/resolve.rs @@ -584,7 +584,10 @@ fn resolve_agg_intent( lower: *lower, upper: *upper, }, - AggIntent::HistogramQuantile { q } => AggIntent::HistogramQuantile { q: *q }, + AggIntent::HistogramQuantile { q, le } => AggIntent::HistogramQuantile { + q: *q, + le: resolve_column_ref(le, schema)?, + }, AggIntent::Math(f) => AggIntent::Math(f.clone()), AggIntent::Absent => AggIntent::Absent, AggIntent::AbsentOverTime => AggIntent::AbsentOverTime, diff --git a/crates/types/src/pre_asap/schema.rs b/crates/types/src/pre_asap/schema.rs index fdac93de..959e37b5 100644 --- a/crates/types/src/pre_asap/schema.rs +++ b/crates/types/src/pre_asap/schema.rs @@ -209,18 +209,10 @@ pub fn with_promql_series_identity(root: &super::QueryExpr) -> Result Ok(()), QueryExpr::PromqlVectorFromScalar(child) if scalar_literal(child).is_some() => Ok(()), - QueryExpr::BinaryOp { - op: super::BinaryOpKind::Arithmetic(_), - lhs, - rhs, - vector_match, - } if vector_match.as_ref().is_none_or(|m| m.grouping.is_none()) - && ![&*lhs, &*rhs].into_iter().any(|side| { - matches!( - side.as_ref(), - QueryExpr::PromqlScalarFromVector(_) | QueryExpr::EvalTimestamp - ) - }) => + QueryExpr::BinaryOp { lhs, rhs, .. } + if ![&*lhs, &*rhs] + .into_iter() + .any(|side| matches!(side.as_ref(), QueryExpr::EvalTimestamp)) => { visit(Rc::make_mut(lhs))?; visit(Rc::make_mut(rhs)) diff --git a/docs/develop_docs/physical-compile-coverage.md b/docs/develop_docs/physical-compile-coverage.md index 31c69782..619aaad9 100644 --- a/docs/develop_docs/physical-compile-coverage.md +++ b/docs/develop_docs/physical-compile-coverage.md @@ -116,37 +116,95 @@ Totals after this change: 19 Supported, 5 Partial, 5 Missing, 2 Backend. query's evaluation time. The raw rows must cover the window at that instant. Vector matching and `without` rewrite the series identity, so their results already lack `__name__`. Other results keep it; the query adapter still drops -it. +it. (Range functions drop it too since the comparison change below.) Totals are unchanged: 19 Supported, 5 Partial, 5 Missing, 2 Backend. +## Covered by per-series arithmetic + +| Row | Change | +|---|---| +| 5 | Query-time `Binary` over rows with a series identity, such as per-series readouts of stored state, uses the Fallback's `series_labels` and `series_binary`. Examples: `avg_over_time` as stored sum/count, and `rate(a) / rate(b)`. Matching drops `__name__` and honors `on`/`ignoring` when the payload carries them. Only one-to-one arithmetic is covered; `group_left`/`group_right` stay rejected and comparisons are row 7. Now Supported. | +| 4, 8 | A literal operand also applies to per-series rows and drops `__name__`, in the Fallback too. Series whose label sets become equal are an error, as in Prometheus. | + +Grouped `sum`/`avg`, current-series `Sum`/`Average` readouts, and +`sum_over_time`/`avg_over_time` use Prometheus' Kahan-Neumaier summation. An +average switches to an incremental mean once the running sum would overflow. +The grouped path also serves SQL `SUM`/`AVG` over Float64, which are now +compensated the same way. +Stored exact `Sum` state still sums without compensation, because a +compensation term would change the stored state layout. Its checked +`avg_over_time` division therefore fails instead of returning a finite mean. + +Totals after this change: 20 Supported, 4 Partial, 5 Missing, 2 Backend. + +## Covered by comparisons and set operators + +The IR now distinguishes `bool` comparisons: `BinaryOpKind::CompareBool(op)` +beside the filtering `BinaryOpKind::Compare(op)`. The PromQL frontend emits it +for `bool`; it previously dropped the modifier. + +`Operator::series_binary` evaluates every PromQL binary operator, in the +Fallback and in query-time `Binary` nodes, following Prometheus' +`VectorBinop`, `VectorAnd`, `VectorOr`, and `VectorUnless`: + +- Arithmetic drops `__name__`. A comparison filter keeps the matched left + series with its value and name; with a scalar on the left it keeps the + vector's value. `bool` yields 1 or 0 and drops the name. NaN compares + unequal to everything. +- Operands may be vectors, literals, `scalar()`, or scalar-valued binary + expressions such as `-scalar(x)`. Two scalars yield a scalar. The + label-map `vector_binary` also treats `CompareBool` as its `bool` mode. +- One-to-one matching and `group_left`/`group_right` with included labels. + The "one" side must not repeat a match group; many-to-one results must be + unique. A left duplicate is an error only if more than one match is kept. +- `and`, `or`, and `unless` match label sets many-to-many, with `on` or + `ignoring`, and return the original series. +- Range functions other than `last_over_time` now drop `__name__` in the + Fallback. Series whose label sets become equal are an error, as in + Prometheus. Vector-scalar results are checked the same way. Inside a + subquery the inner function also drops the name, without that check, + because each series repeats across steps. + +A result is written in the left operand's schema. Without a series identity, +its label columns must hold every label the right side can contribute +(`or`, `group_right`, and `group_left` labels); otherwise `compile` rejects it. +So `sum(a) or vector(0)` compiles, but +`sum by (job) (a) * on(job) group_left(team) info` is rejected: the +aggregate's schema has no `team` column. + +| Row | Change | +|---|---| +| 1 | Comparisons, `bool`, set operators, `group_left`/`group_right`, `scalar()` operands, and literals over aggregates whose value has another name, such as `sum by (job) (a) * 2`. Still Partial. | +| 5 | Grouped `Binary` rows use the same operator instead of a relational join. A duplicate match group is now an error instead of a cross product. | +| 7 | Fallback, grouped `Binary`, and per-series `bool` comparisons. Per-series filter comparisons and set operators on `Binary` nodes fail closed: stored readouts keep `__name__` even where the range function drops it. Now Partial. | + +Totals after this change: 20 Supported, 5 Partial, 4 Missing, 2 Backend. + +## Covered by classic histogram_quantile + +| Row | Change | +|---|---| +| 10 | A classic-bucket `histogram_quantile(q, v)` lowers to `Aggregate{Reduce(without([le])), [HistogramQuantile{q, le}]}`, where `le` is the argument's `le` column. The frontend seeds `le` into the selector's schema. The Fallback compiles it to one operator that groups rows by every label except `le` and applies Prometheus `bucketQuantile`. The output labels are the input labels without `le` and `__name__`; result label sets that become equal are an error, as in Prometheus. Covers `histogram_quantile(q, rate(x_bucket[5m]))`, `histogram_quantile(q, sum by (le, job) (…))`, and bare bucket selectors. Now Supported. | + +An argument whose output provably lacks `le`, such as +`sum by (job) (rate(x_bucket[5m]))`, is rejected at lowering. Prometheus +returns an empty vector for it. Candidate search keeps the classic form as one +exact `KeepPreAsap` subtree for every accuracy target; it has no sketch +candidate. `histogram_quantiles` lowers each branch the same way, but the +Fallback compiler does not yet accept its `Concat` of relabeled branches. +Aggregating or doing arithmetic over the result, as in +`sum(histogram_quantile(…))`, fails like any expression over a nested +aggregate's renamed value column. + +Totals after this change: 21 Supported, 5 Partial, 3 Missing, 2 Backend. + ## Remaining In order of backend usage: -1. Rows 1, 10, and 12, the remaining `Fallback` shapes: - - `histogram_quantile` (row 10). The compiler cannot recover the grouping - from today's IR, which is `Aggregate{Reduce(by ()), [HistogramQuantile{q}]}`. - It computes one quantile over all buckets and drops the output labels. The IR - must carry: - - The bucket label: the `le` column of the child's schema, named by the - intent, for example `HistogramQuantile{q, le: C}`. The frontend must - add `le` to the selector's schema even when no matcher names it. - - The grouping: `Reduce(without([le]))`, so that each histogram is - one set of series that differ only in `le`. An explicit - `sum by (x, le)` inside the argument still yields `without (le)` over - those rows. - - The output labels: every input label except `le` and `__name__`. The - `without` output schema already carries them, and the series identity - with `le` removed. `series_labels(Ignoring, [le])` computes the latter. - The operator then applies Prometheus `bucketQuantile`. It parses `le` - as a float, skips unparsable values, and requires a `+Inf` bucket, else - it returns NaN. It forces cumulative counts to be monotonic and returns - NaN for fewer than two buckets. For q < 0 it returns -Inf; for q > 1, - +Inf. - - Comparisons and set operators (`and`, `or`, `unless`); `group_left` and - `group_right`; arithmetic with a non-literal scalar, such as - `scalar(x)` or `time()`. +1. Rows 1 and 12, the remaining `Fallback` shapes: + - `time()` and other scalar functions as operands. - Subquery operands other than one per-series function; implicit subquery resolution, which is a deployment default. - `@ start()` and `@ end()`, which need the range query's bounds in the @@ -154,13 +212,22 @@ In order of backend usage: - Other functions, such as `deriv`, `predict_linear`, `stddev_over_time`, `absent`, `label_replace`, and math functions. After these shapes are covered, the backend can delete rows 28 and 30. -2. Row 7: comparison filters and `bool` comparisons. This needs `return_bool` - in the `Binary` payload. `compile` currently rejects comparisons. -3. Row 5 for per-series rows in a `Binary` payload node: its grouped-row join - still rejects `$promql_series_identity`. It could reuse the Fallback's - `series_labels` and `series_binary` operators. -4. Rows 25 and 27: constant weights and `EntityIdentity` items for precompute +2. Row 7: per-series readouts must drop `__name__` where their range + function does, and check for equal label sets. Then filter comparisons + and set operators over them can compile. +3. Rows 25 and 27: constant weights and `EntityIdentity` items for precompute `SummaryAgg`. -5. Row 16: a label-map sketch-state readout, the counterpart of +4. Row 16: a label-map sketch-state readout, the counterpart of `compile_exact_readout`, and MetricsQL `__name__` retention rules. -6. Row 20: summary join, subtract, and delete. +5. Row 20: summary join, subtract, and delete. +6. Compensated stored exact `Sum` state, a state-layout change shared with the + backend's stored-state decoding. +7. Non-finite literals (`NaN`, `Inf`) in compiled operators do not survive a + JSON round trip of the program. +8. `group_left` labels and `or`/`group_right` right-side labels onto + aggregated (label-column) rows. The logical output schema, which is the left + side's, has no column for them. +9. An equal-label-set check for inner subquery functions, per step. + +`fill`, `fill_left`, and `fill_right` matching modifiers are rejected by the +frontend (#494); they are never silently ignored. diff --git a/docs/develop_docs/pre-asap-ir.md b/docs/develop_docs/pre-asap-ir.md index 15bfc979..94677516 100644 --- a/docs/develop_docs/pre-asap-ir.md +++ b/docs/develop_docs/pre-asap-ir.md @@ -111,10 +111,16 @@ Rate, Increase // counter der Changes, Delta, IDelta, Deriv, Resets, PredictLinear(seconds), DoubleExpSmoothing(sf, tf) // range-vector functions HistogramCount, HistogramSum, HistogramAvg, HistogramStdDev, -HistogramStdVar, HistogramFraction(lo, hi), HistogramQuantile(q) // native-histogram accessors +HistogramStdVar, HistogramFraction(lo, hi) // native-histogram accessors +HistogramQuantile(q, le) // classic-bucket quantile Math(func) // element-wise transform ``` +`HistogramQuantile { q, le }` names its bucket-bound column `le`. PromQL's +`Aggregate` groups it `without([le])`, so one histogram is the set of series +that differ only in `le`. The SQL `asap_histogram_quantile` bridge reads one +histogram over all rows. + `PearsonCorr { left, right }` has two value inputs. Both references resolve to positional column IDs, and `input_cols()` exposes both dependencies. SQL lowering projects both arguments, preserving