diff --git a/crates/asap-physical-operators/src/operators/aggregate/mod.rs b/crates/asap-physical-operators/src/operators/aggregate/mod.rs index 71e7134c..d929fc07 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/mod.rs +++ b/crates/asap-physical-operators/src/operators/aggregate/mod.rs @@ -27,6 +27,12 @@ impl Operator { false, ) } + Reduction::Quantile { column, q } => { + if plain(&input, *column)?.0 != &DataType::Float64 || q.is_nan() { + return Err(invalid("quantile requires Float64 input and a numeric q")); + } + (DataType::Float64, false) + } Reduction::Min(i) | Reduction::Max(i) => { let (t, nullable) = plain(&input, *i)?; if !ordered(t) { @@ -120,6 +126,11 @@ pub enum Reduction { Avg(usize), Min(usize), Max(usize), + /// PromQL `quantile`: linear interpolation between closest ranks. + Quantile { + column: usize, + q: f64, + }, } pub(super) fn execute<'a>( operator: &'a Operator, @@ -198,6 +209,31 @@ async fn reduce( Ok(output) } +// Matches Prometheus `quantile`: NaN for no values, ±Inf outside [0, 1], +// and NaN samples ordered first. +pub(super) fn quantile(q: f64, mut values: Vec) -> f64 { + if values.is_empty() { + return f64::NAN; + } + if q < 0. { + return f64::NEG_INFINITY; + } + if q > 1. { + return f64::INFINITY; + } + values.sort_by(|a, b| match (a.is_nan(), b.is_nan()) { + (true, true) => std::cmp::Ordering::Equal, + (true, false) => std::cmp::Ordering::Less, + (false, true) => std::cmp::Ordering::Greater, + _ => a.total_cmp(b), + }); + let rank = q * (values.len() - 1) as f64; + let low = rank.floor() as usize; + let high = (low + 1).min(values.len() - 1); + let weight = rank - low as f64; + values[low] * (1. - weight) + values[high] * weight +} + async fn reduce_one( rows: &[Vec], measure: &Reduction, @@ -211,6 +247,18 @@ async fn reduce_one( )) } Reduction::Sum(i) | Reduction::Avg(i) | Reduction::Min(i) | Reduction::Max(i) => *i, + Reduction::Quantile { column, q } => { + let mut values = Vec::with_capacity(rows.len()); + for row in rows { + work.checkpoint().await?; + match &row[*column] { + Value::Float64(value) => values.push(*value), + Value::Null => {} + _ => return Err(invalid("floating quantile value required")), + } + } + return Ok(Value::Float64(quantile(*q, values))); + } }; let values = rows .iter() @@ -291,3 +339,17 @@ async fn reduce_one( sum })) } + +#[cfg(test)] +mod tests { + // Quantile follows Prometheus: interpolate ranks, NaN when empty, ±Inf outside [0, 1]. + #[test] + fn quantile_matches_prometheus_edge_cases() { + assert!(super::quantile(0.5, vec![]).is_nan()); + assert_eq!(super::quantile(0.5, vec![3.]), 3.); + assert_eq!(super::quantile(0.75, vec![4., 1., 2., 3.]), 3.25); + assert_eq!(super::quantile(-0.1, vec![1.]), f64::NEG_INFINITY); + assert_eq!(super::quantile(1.1, vec![1.]), f64::INFINITY); + assert_eq!(super::quantile(1., vec![2., f64::NAN, 1.]), 2.); + } +} diff --git a/crates/asap-physical-operators/src/operators/aggregate/temporal.rs b/crates/asap-physical-operators/src/operators/aggregate/temporal.rs index dbf66108..3ac0a3d1 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/temporal.rs +++ b/crates/asap-physical-operators/src/operators/aggregate/temporal.rs @@ -65,38 +65,7 @@ pub(in crate::operators) async fn reduce( "duplicate or out-of-window timestamp".into(), )); } - match intent { - AggIntent::Rate => rate(&points, start, end).map(Value::Float64), - AggIntent::Increase => rate(&points, start, end) - .map(|v| Value::Float64(v * (end as f64 - start as f64) / 1000.)), - 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::Min { .. } => { - Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| { - if a.is_nan() || p.1 < a { - p.1 - } else { - a - } - }))) - } - AggIntent::Max { .. } => { - Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| { - if a.is_nan() || p.1 > a { - p.1 - } else { - a - } - }))) - } - _ => return Err(Error::Invalid("unsupported temporal intent".into())), - } + window_value(intent, &points, start, end)? }; if let Some(result) = result { keys.push(result); @@ -106,7 +75,76 @@ pub(in crate::operators) async fn reduce( Ok(output) } -fn rate(points: &[(i64, f64)], start: i64, end: i64) -> Option { +/// One series' value over its sorted samples in the window `(start, end]`. +/// `None` means PromQL emits no sample for this series. +pub(in crate::operators) fn window_value( + intent: &AggIntent, + points: &[(i64, f64)], + start: i64, + end: i64, +) -> Result, Error> { + Ok(match intent { + AggIntent::Rate => rate(points, start, end, true).map(Value::Float64), + AggIntent::Delta => rate(points, start, end, false) + .map(|v| Value::Float64(v * (end as f64 - start as f64) / 1000.)), + AggIntent::Increase => rate(points, start, end, true) + .map(|v| Value::Float64(v * (end as f64 - start as f64) / 1000.)), + 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::Min { .. } => Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| { + if a.is_nan() || p.1 < a { + p.1 + } else { + a + } + }))), + AggIntent::Max { .. } => Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| { + if a.is_nan() || p.1 > a { + p.1 + } else { + a + } + }))), + AggIntent::IRate | AggIntent::IDelta => { + let [.., (t0, v0), (t1, v1)] = points else { + return Ok(None); + }; + let rate = matches!(intent, AggIntent::IRate); + // A counter reset makes the last value the increase. + let delta = if rate && v1 < v0 { *v1 } else { v1 - v0 }; + match (rate, t1 - t0) { + (_, 0) => None, + (true, interval) => Some(Value::Float64(delta / (interval as f64 / 1000.))), + (false, _) => Some(Value::Float64(delta)), + } + } + AggIntent::Changes | AggIntent::Resets => { + let changed = |(a, b): (f64, f64)| match intent { + AggIntent::Changes => a != b && !(a.is_nan() && b.is_nan()), + _ => b < a, + }; + let count = points + .windows(2) + .filter(|pair| changed((pair[0].1, pair[1].1))) + .count(); + Some(Value::Float64(count as f64)) + } + AggIntent::LastOverTime => points.last().map(|p| Value::Float64(p.1)), + AggIntent::Quantile { col: None, q, .. } => Some(Value::Float64(super::quantile( + *q, + points.iter().map(|p| p.1).collect(), + ))), + _ => return Err(Error::Invalid("unsupported temporal intent".into())), + }) +} + +/// Prometheus `extrapolatedRate`; `counter` enables reset correction and the zero bound. +fn rate(points: &[(i64, f64)], start: i64, end: i64, counter: bool) -> Option { if points.len() < 2 { return None; } @@ -118,7 +156,7 @@ fn rate(points: &[(i64, f64)], start: i64, end: i64) -> Option { } let mut delta = last - first; for pair in points.windows(2) { - if pair[1].1 < pair[0].1 { + if counter && pair[1].1 < pair[0].1 { delta += pair[0].1; } } @@ -129,7 +167,7 @@ fn rate(points: &[(i64, f64)], start: i64, end: i64) -> Option { to_start = average / 2.; } // Apply the zero bound after the sparse-window half-interval cap. - if delta > 0. && first >= 0. { + if counter && delta > 0. && first >= 0. { to_start = to_start.min(span * first / delta); } if to_end >= average * 1.1 { diff --git a/crates/asap-physical-operators/src/operators/mod.rs b/crates/asap-physical-operators/src/operators/mod.rs index 9ae0772c..196e9e8a 100644 --- a/crates/asap-physical-operators/src/operators/mod.rs +++ b/crates/asap-physical-operators/src/operators/mod.rs @@ -23,6 +23,8 @@ mod joins; mod limit; mod projection; mod scope_timestamp; +mod series_labels; +mod series_window; mod sort; mod source; mod summary; @@ -30,6 +32,7 @@ mod unchecked; pub(crate) mod vector_binary; pub(crate) mod vector_window; pub use aggregate::Reduction; +pub use series_window::SubquerySteps; pub use sort::SortKey; pub use summary::ReadoutQuery; #[derive(Clone, serde::Serialize, serde::Deserialize)] @@ -66,6 +69,22 @@ enum Kind { intent: Box>, }, HistogramQuantile, + SeriesWindow { + function: Option>>, + coordinate: usize, + value: usize, + range_ms: i64, + offset_ms: i64, + at_ms: Option, + steps: Option, + }, + SeriesLabels { + kind: planner_types::pre_asap::VectorMatchKind, + labels: Vec, + }, + SeriesBinary { + operator: planner_types::post_asap::BinaryOperator, + }, Project(Vec), Filter(Expression), Limit { @@ -237,6 +256,9 @@ impl PhysicalOperator for Operator { | Kind::RangeWindow { .. } | Kind::HistogramQuantile | Kind::CurrentSeries { .. } + | Kind::SeriesWindow { .. } + | Kind::SeriesLabels { .. } + | Kind::SeriesBinary { .. } | Kind::Aggregate { .. } | Kind::Window { .. } | Kind::Join { .. } @@ -279,6 +301,9 @@ impl PhysicalOperator for Operator { Kind::AlignedBinary { .. } => "AlignedBinary", Kind::RangeWindow { .. } => "RangeWindow", Kind::HistogramQuantile => "HistogramQuantile", + Kind::SeriesWindow { .. } => "SeriesWindow", + Kind::SeriesLabels { .. } => "SeriesLabels", + Kind::SeriesBinary { .. } => "SeriesBinary", Kind::Project(_) => "Project", Kind::Filter(_) => "Filter", Kind::Limit { .. } => "Limit", @@ -295,6 +320,7 @@ impl PhysicalOperator for Operator { } fn validate_context(&self, context: &RunContext) -> Result<(), Error> { current_series::validate_context(self, context)?; + series_window::validate_context(self, context)?; self.readout_range(context).map(|_| ()) } fn input_schemas(&self) -> Vec { @@ -323,6 +349,10 @@ impl PhysicalOperator for Operator { Kind::Project(_) => projection::execute(self, inputs, context), 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::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 new file mode 100644 index 00000000..be248319 --- /dev/null +++ b/crates/asap-physical-operators/src/operators/series_labels.rs @@ -0,0 +1,204 @@ +//! PromQL label-set rewriting and one-to-one vector matching over rows that +//! carry a series identity or plain label columns. +use super::*; +use planner_types::{ + post_asap::BinaryOperator, + pre_asap::{schema::PROMQL_SERIES_IDENTITY, BinaryOpKind, VectorMatchKind}, +}; + +type Labels = BTreeMap; + +/// Where a row's label set lives: the encoded series identity when present, +/// otherwise the non-empty Utf8 label columns. The one Float64 column is the +/// sample value, whatever an aggregate named it. +#[derive(Clone, Debug)] +struct Layout { + identity: Option, + labels: Vec, + value: usize, +} + +fn layout(input: &Schema) -> Result { + let mut identity = None; + let mut labels = Vec::new(); + let mut value = None; + for (i, field) in input.fields.iter().enumerate() { + match plain(input, i)? { + (DataType::Utf8, false) if field.name == PROMQL_SERIES_IDENTITY => identity = Some(i), + (DataType::Utf8, _) if field.name != PROMQL_SERIES_IDENTITY => labels.push(i), + (DataType::Float64, false) if value.is_none() => value = Some(i), + (DataType::Timestamp, false) if input.time_index == Some(i) => {} + _ => return Err(invalid("PromQL vector requires labels, time and one value")), + } + } + Ok(Layout { + identity, + labels, + value: value.ok_or_else(|| invalid("PromQL vector requires one value"))?, + }) +} + +impl Layout { + fn read(&self, input: &Schema, row: &[Value]) -> Result { + if let Some(i) = self.identity { + let Value::Utf8(encoded) = &row[i] else { + return Err(invalid("series identity must be Utf8")); + }; + let mut labels: Labels = + serde_json::from_str(encoded).map_err(|e| invalid(&e.to_string()))?; + // PromQL treats an empty label value as an absent label. + labels.retain(|_, v| !v.is_empty()); + return Ok(labels); + } + let mut labels = Labels::new(); + for &i in &self.labels { + match &row[i] { + Value::Utf8(v) if !v.is_empty() => { + labels.insert(input.fields[i].name.clone(), v.to_string()); + } + Value::Utf8(_) | Value::Null => {} + _ => return Err(invalid("label must be Utf8")), + } + } + Ok(labels) + } + + /// Replace the row's labels; a label column absent from `labels` is empty. + fn write(&self, input: &Schema, row: &mut [Value], labels: &Labels) -> Result<(), Error> { + if let Some(i) = self.identity { + let encoded = serde_json::to_string(labels).map_err(|e| invalid(&e.to_string()))?; + row[i] = Value::Utf8(encoded.into()); + } + for &i in &self.labels { + let value = labels.get(&input.fields[i].name).map_or("", String::as_str); + row[i] = Value::Utf8(value.into()); + } + Ok(()) + } +} + +impl Operator { + /// Rewrite each row's label set to PromQL's matching labels: `On` keeps + /// only `labels`; `Ignoring` drops `labels` and the metric name. + pub fn series_labels( + input: Schema, + kind: VectorMatchKind, + labels: Vec, + ) -> Result { + layout(&input)?; + Ok(Self { + kind: Kind::SeriesLabels { kind, labels }, + 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. + pub fn series_binary( + left: Schema, + right: Schema, + operator: BinaryOperator, + ) -> Result { + layout(&left)?; + layout(&right)?; + if !matches!(operator.kind, BinaryOpKind::Arithmetic(_)) || operator.vector_match.is_some() + { + return Err(invalid("series binary requires unmatched arithmetic")); + } + Ok(Self { + kind: Kind::SeriesBinary { operator }, + inputs: vec![left.clone(), right], + output: left, + }) + } +} + +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"))?) + } + _ => None, + }; + let left = inputs.pop().ok_or_else(|| invalid("missing left input"))?; + Ok(futures::stream::once(async move { + let mut work = Cooperative::new(&context); + let mut workspace = Workspace::new(&context)?; + let (rows, _memory) = collect_rows(left, &context).await?; + let mut result = Vec::new(); + match (&operator.kind, right) { + (Kind::SeriesLabels { kind, labels }, None) => { + 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)) + } + } + left_layout.write(&output, &mut row, &set)?; + workspace.grow(row_bytes(&row))?; + result.push(row); + } + } + (Kind::SeriesBinary { operator: binary }, Some(right)) => { + let right_schema = &operator.inputs[1]; + let right_layout = layout(right_schema)?; + 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 { + 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 { + 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", + )); + } + } + for mut row in rows { + 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) { + return Err(invalid( + "many-to-one matching must be explicit (group_left/group_right)", + )); + } + 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); + } + } + _ => return Err(invalid("series label operator inputs mismatch")), + } + Batch::try_new(output.clone(), result) + }) + .boxed_local()) +} diff --git a/crates/asap-physical-operators/src/operators/series_window.rs b/crates/asap-physical-operators/src/operators/series_window.rs new file mode 100644 index 00000000..20ab1aed --- /dev/null +++ b/crates/asap-physical-operators/src/operators/series_window.rs @@ -0,0 +1,243 @@ +//! PromQL per-series evaluation over the samples before an evaluation instant. +use super::*; +use planner_types::pre_asap::AggIntent; + +/// A PromQL subquery grid: every multiple of `step_ms` in +/// `(T - offset_ms - range_ms, T - offset_ms]`. `T` is `at_ms` (the subquery's +/// `@`) when present, else the query time. +#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub struct SubquerySteps { + pub range_ms: i64, + pub step_ms: i64, + pub offset_ms: i64, + pub at_ms: Option, +} + +const STALE_MARKER: u64 = 0x7ff0_0000_0000_0002; +const MAX_SUBQUERY_STEPS: i64 = 100_000; + +impl Operator { + /// Evaluate each series at instant `t` over its samples in + /// `(e - offset_ms - range_ms, e - offset_ms]`. `t` is the query time, or + /// each step of `steps`; `e` is `at_ms` (the selector's `@`) when present, + /// else `t`. `function: None` is instant selection: the latest + /// sample, absent if it is a stale marker. Range functions ignore stale + /// markers. A series is every column except the time and `value` columns. + /// Output rows keep the input schema, with time `t` and the result value. + pub fn series_window( + input: Schema, + function: Option>, + range_ms: i64, + offset_ms: i64, + at_ms: Option, + steps: Option, + ) -> Result { + let coordinate = input + .time_index + .ok_or_else(|| invalid("series window requires a time column"))?; + let values = input + .fields + .iter() + .enumerate() + .filter(|(_, f)| f.name == "value") + .map(|(i, _)| i) + .collect::>(); + let [value] = values.as_slice() else { + return Err(invalid("series window requires one value column")); + }; + if plain(&input, coordinate)? != (&DataType::Timestamp, false) + || plain(&input, *value)? != (&DataType::Float64, false) + { + return Err(invalid( + "series window requires non-null time and Float64 value", + )); + } + if range_ms <= 0 || steps.is_some_and(|s| s.range_ms <= 0 || s.step_ms <= 0) { + return Err(invalid("series window ranges and steps must be positive")); + } + // Bounds per-run work independently of the data, as the backend's grid does. + if steps.is_some_and(|s| s.range_ms / s.step_ms > MAX_SUBQUERY_STEPS) { + return Err(invalid("subquery exceeds 100000 steps")); + } + if !matches!( + function, + None | Some( + AggIntent::Rate + | AggIntent::Increase + | AggIntent::Delta + | AggIntent::Count { .. } + | AggIntent::Sum { col: None } + | AggIntent::Avg { col: None } + | AggIntent::Min { col: None } + | AggIntent::Max { col: None } + | AggIntent::IRate + | AggIntent::IDelta + | AggIntent::Changes + | AggIntent::Resets + | AggIntent::LastOverTime + | AggIntent::Quantile { col: None, .. } + ) + ) { + return Err(invalid("unsupported PromQL range function")); + } + Ok(Self { + kind: Kind::SeriesWindow { + function: function.map(Box::new), + coordinate, + value: *value, + range_ms, + offset_ms, + at_ms, + steps, + }, + inputs: vec![input.clone()], + output: input, + }) + } +} + +/// The first and last evaluation instants and the step between them. The grid +/// is iterated, not allocated: its size depends only on the query. +fn evaluation_times( + context: &RunContext, + steps: Option, +) -> Result<(i64, i64, i64), Error> { + let crate::runtime::Scope::Query { + evaluation_time_ms, .. + } = context.scope + else { + return Err(invalid("series window requires a query evaluation time")); + }; + let Some(steps) = steps else { + return Ok((evaluation_time_ms, evaluation_time_ms, 1)); + }; + let overflow = || invalid("subquery grid overflows"); + let end = steps + .at_ms + .unwrap_or(evaluation_time_ms) + .checked_sub(steps.offset_ms) + .ok_or_else(overflow)?; + let start = end.checked_sub(steps.range_ms).ok_or_else(overflow)?; + let first = (start.div_euclid(steps.step_ms) + 1) + .checked_mul(steps.step_ms) + .ok_or_else(overflow)?; + Ok((first, end, steps.step_ms)) +} + +pub(super) fn execute<'a>( + operator: &'a Operator, + mut inputs: Vec>, + context: RunContext, +) -> Result, Error> { + let Kind::SeriesWindow { + function, + coordinate, + value, + range_ms, + offset_ms, + at_ms, + steps, + } = &operator.kind + else { + unreachable!() + }; + let (coordinate, value) = (*coordinate, *value); + let input = inputs + .pop() + .ok_or_else(|| invalid("series window input missing"))?; + let times = evaluation_times(&context, *steps)?; + Ok(futures::stream::once(async move { + let (rows, _memory) = collect_rows(input, &context).await?; + let mut work = Cooperative::new(&context); + let mut workspace = Workspace::new(&context)?; + let identity = (0..operator.output.fields.len()) + .filter(|&i| i != coordinate && i != value) + .collect::>(); + let mut series = BTreeMap::>, (usize, Vec<(i64, f64)>)>::new(); + for (index, row) in rows.iter().enumerate() { + work.checkpoint().await?; + let key = group_key(row, &identity)?; + let (Value::Timestamp(time), Value::Float64(sample)) = (&row[coordinate], &row[value]) + else { + return Err(invalid("series window requires time and value samples")); + }; + workspace.grow(16)?; + if !series.contains_key(&key) { + workspace.grow(key_bytes(&key) + 64)?; + } + series + .entry(key) + .or_insert_with(|| (index, Vec::new())) + .1 + .push((*time, *sample)); + } + for (_, points) in series.values_mut() { + work.checkpoint().await?; + points.sort_by_key(|p| p.0); + if points.windows(2).any(|p| p[0].0 == p[1].0) { + return Err(invalid("duplicate sample timestamp for one series")); + } + } + let mut output = Vec::new(); + let (mut time, last, step) = times; + while time <= last && !series.is_empty() { + work.checkpoint().await?; + let overflow = || invalid("series window overflows"); + let end = at_ms + .unwrap_or(time) + .checked_sub(*offset_ms) + .ok_or_else(overflow)?; + let start = end.checked_sub(*range_ms).ok_or_else(overflow)?; + for (template, points) in series.values() { + work.checkpoint().await?; + // PromQL ranges are left-open: a sample at `start` is outside. + let first = points.partition_point(|p| p.0 <= start); + let last = points.partition_point(|p| p.0 <= end); + let points = &points[first..last]; + let result = match function { + None => points + .last() + .filter(|p| p.1.to_bits() != STALE_MARKER) + .map(|p| p.1), + Some(intent) => { + let fresh = points + .iter() + .copied() + .filter(|p| p.1.to_bits() != STALE_MARKER) + .collect::>(); + if fresh.is_empty() { + None + } else { + match aggregate::temporal::window_value(intent, &fresh, start, end)? { + Some(Value::Float64(v)) => Some(v), + Some(Value::Int64(v)) => Some(v as f64), + Some(_) => return Err(invalid("invalid range function result")), + None => None, + } + } + } + }; + if let Some(result) = result { + let mut row = rows[*template].clone(); + row[coordinate] = Value::Timestamp(time); + row[value] = Value::Float64(result); + workspace.grow(row_bytes(&row))?; + output.push(row); + } + } + let Some(next) = time.checked_add(step) else { + break; + }; + time = next; + } + Batch::try_new(operator.output.clone(), output) + }) + .boxed_local()) +} + +pub(super) fn validate_context(operator: &Operator, context: &RunContext) -> Result<(), Error> { + if let Kind::SeriesWindow { steps, .. } = operator.kind { + evaluation_times(context, steps)?; + } + Ok(()) +} diff --git a/crates/asap-physical-operators/src/operators/unchecked.rs b/crates/asap-physical-operators/src/operators/unchecked.rs index f927a4e4..46bc6a0c 100644 --- a/crates/asap-physical-operators/src/operators/unchecked.rs +++ b/crates/asap-physical-operators/src/operators/unchecked.rs @@ -50,6 +50,27 @@ impl TryFrom for Operator { } => Operator::aligned_binary(input(0)?, input(1)?, keys, values, operator)?, Kind::RangeWindow { intent } => Operator::range_window(*intent)?, Kind::HistogramQuantile => Operator::histogram_quantile(), + Kind::SeriesWindow { + function, + range_ms, + offset_ms, + at_ms, + steps, + .. + } => Operator::series_window( + input(0)?, + function.map(|f| *f), + range_ms, + offset_ms, + at_ms, + steps, + )?, + Kind::SeriesLabels { kind, labels } => { + Operator::series_labels(input(0)?, kind, labels)? + } + Kind::SeriesBinary { operator } => { + Operator::series_binary(input(0)?, input(1)?, operator)? + } Kind::Project(expressions) => { if expressions.len() != output.fields.len() { return Err(invalid("projection width mismatch")); diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index fede3777..655c024a 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -14,7 +14,8 @@ use planner_types::{ SketchQuery, SummaryFamilyType, SummaryInputExpr, ValueOperation, }, pre_asap::{ - AggIntent, ColumnRef, CompareOpKind, GroupKeys, QueryExpr, Reduction as PlannerReduction, + AggIntent, ColumnRef, CompareOpKind, DataType, GroupKeys, QueryExpr, + Reduction as PlannerReduction, }, }; use std::{ @@ -31,6 +32,7 @@ fn invalid(message: impl Into) -> Error { pub type Source<'a> = Box + 'a>; pub mod precompute; +pub mod promql_fallback; pub mod promql_rows; pub mod promql_values; @@ -43,6 +45,8 @@ pub use candidates::{ mod compiled; pub use compiled::{CompiledPhysicalDag, InputContract}; +mod row_values; + /// Compile computation without opening or retaining deployment readers. /// Input contracts identify explicit boundaries selected by maintenance planning. pub fn compile( @@ -143,13 +147,46 @@ fn compile_internal( }, ) }); + // Scalar literal operands of query-time arithmetic are folded into the consumer. + let mut literals = BTreeMap::::new(); for edge in edges { + let consumer = u64::from(edge.consumer.0); + if let ( + Payload::Fallback { expression }, + Some(PostAsapDagNode { + payload: Payload::Binary { .. }, + .. + }), + ) = ( + &nodes[&u64::from(edge.producer.0)].payload, + nodes.get(&consumer), + ) { + if let Some(value) = row_values::scalar_literal(expression) { + let left = edge.role == planner_types::post_asap::EdgeRole::Left; + if literals.insert(consumer, (value, left)).is_some() { + return Err(invalid("binary with two scalar literals is not folded")); + } + continue; + } + } dependencies .entry(u64::from(edge.consumer.0)) .or_default() .push(u64::from(edge.producer.0)); } - if sources.keys().any(|id| !nodes.contains_key(id)) { + let known = |id: &NodeId| { + nodes.contains_key(id) + || promql_fallback::raw_series_owner(*id).is_some_and(|owner| { + matches!( + nodes.get(&owner), + Some(PostAsapDagNode { + payload: Payload::Fallback { .. }, + .. + }) + ) + }) + }; + if !sources.keys().all(known) { return Err(invalid("source binding names an unknown node")); } let mut ordered = Vec::new(); @@ -176,7 +213,7 @@ fn compile_internal( let mut graph = CompiledPhysicalDag::new(roots.to_vec()); for id in ordered { let node = nodes[&id]; - let auxiliary = helper_id(id, 0); + let mut auxiliary = helper_id(id, 0); let output = Arc::new(node.output_schema.clone()); crate::values::validate_schema(&output)?; if let Some(source) = sources.remove(&id) { @@ -204,6 +241,65 @@ fn compile_internal( inputs = vec![auxiliary]; schemas.truncate(1); } + // A consumed bare selector supplies raw range rows (e.g. to a + // per-entity summary), not an instant vector, so only its consumer computes. + let raw_rows = matches!( + &node.payload, + Payload::Fallback { + expression: QueryExpr::TimeRange { .. } + } + ) && dag.edges.iter().any(|e| u64::from(e.producer.0) == id); + if let (Payload::Fallback { expression }, false) = (&node.payload, raw_rows) { + let promql_fallback::Lowering { + selectors, + mut steps, + } = promql_fallback::lower(expression) + .map_err(|error| invalid(format!("node {id}: {error}")))?; + let mut slots = Vec::new(); + for (i, (_, schema)) in selectors.iter().enumerate() { + let slot = promql_fallback::raw_series_input(id, i); + match sources.remove(&slot) { + Some(contract) if &contract.schema == schema => { + graph.add_input(slot, contract)? + } + Some(_) => { + return Err(invalid(format!( + "node {id}: raw series input {slot} differs from the selector schema" + ))) + } + None => { + return Err(invalid(format!( + "node {id}: PromQL fallback requires raw series input {slot}" + ))) + } + } + slots.push(slot); + } + let (last, last_inputs) = steps + .pop() + .ok_or_else(|| invalid("empty PromQL lowering"))?; + let mut ids = Vec::new(); + let resolve = |inputs: Vec, ids: &[NodeId]| { + inputs + .into_iter() + .map(|input| match input { + promql_fallback::Input::Raw(i) => slots[i], + promql_fallback::Input::Step(i) => ids[i], + }) + .collect::>() + }; + for (operator, inputs) in steps { + graph.add(auxiliary, resolve(inputs, &ids), operator)?; + ids.push(auxiliary); + auxiliary -= 1; + } + graph.add( + id, + resolve(last_inputs, &ids), + last.with_output_schema(output)?, + )?; + continue; + } if let Payload::Value { operation: ValueOperation::MaintainPopulation { population }, } = &node.payload @@ -247,11 +343,6 @@ fn compile_internal( use planner_types::post_asap::maintained_population::{ PopulationInput, PopulationReadout, }; - let PopulationReadout::TopK { k } = readout else { - return Err(invalid( - "native population readout does not support this operation", - )); - }; let [producer] = inputs.as_slice() else { return Err(invalid("population readout requires one input")); }; @@ -272,6 +363,19 @@ fn compile_internal( )); } let input = schemas[0].clone(); + let PopulationReadout::TopK { k } = readout else { + let mut chain = + row_values::population_aggregate(&input, &spec.grouping, readout)?; + let last = chain.pop().expect("nonempty chain"); + let mut inputs = inputs; + for operator in chain { + graph.add(auxiliary, inputs, operator)?; + inputs = vec![auxiliary]; + auxiliary -= 1; + } + graph.add(id, inputs, last.with_output_schema(output)?)?; + continue; + }; let groups = spec .grouping .iter() @@ -357,6 +461,73 @@ fn compile_internal( )?; continue; } + if let Payload::Binary { operator } = &node.payload { + let query_time = node.output_state.timing + == planner_types::post_asap::ExecutionTiming::QueryTime; + if let Some(&(value, left)) = literals.get(&id) { + let [input] = schemas.as_slice() else { + return Err(invalid("scalar binary requires one row input")); + }; + if !query_time { + return Err(invalid("scalar literal binary must run at query time")); + } + let project = row_values::scalar_binary(input, operator, value, left) + .map_err(|error| invalid(format!("node {id}: {error}")))?; + graph.add(id, inputs, project.with_output_schema(output)?)?; + continue; + } + let label_map = |schema: &Schema| { + schema + .fields + .iter() + .any(|f| matches!(f.dtype, SummaryFamilyType::Plain(DataType::Map { .. }))) + }; + 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; + continue; + } + } + } + if let Payload::Value { + operation: ValueOperation::FinalizeExactAccumulator, + } = &node.payload + { + // Exact counts read out as Int64; PromQL declares a Float64 sample. + let readout = bind_operation(node, &schemas) + .map_err(|error| invalid(format!("node {id}: {error}")))?; + let actual = readout.schema(); + let converted = actual.fields.iter().zip(&output.fields).position(|(a, d)| { + a.dtype == SummaryFamilyType::Plain(DataType::Int64) + && d.dtype == SummaryFamilyType::Plain(DataType::Float64) + }); + if let Some(column) = converted { + let columns = actual + .fields + .iter() + .enumerate() + .map(|(i, field)| { + ( + field.name.clone(), + if i == column { + Expression::ExactFloat64(i) + } else { + Expression::Column(i) + }, + ) + }) + .collect(); + let project = Operator::project(actual, columns)?.with_output_schema(output)?; + graph.add(auxiliary, inputs, readout)?; + graph.add(id, vec![auxiliary], project)?; + auxiliary -= 1; + continue; + } + } let mut operator = compile_node(node, &schemas) .map_err(|error| invalid(format!("node {id}: {error}")))?; if operator.is_counter_readout() { diff --git a/crates/asap-physical-operators/src/physical_planner/promql_fallback.rs b/crates/asap-physical-operators/src/physical_planner/promql_fallback.rs new file mode 100644 index 00000000..12c60ad1 --- /dev/null +++ b/crates/asap-physical-operators/src/physical_planner/promql_fallback.rs @@ -0,0 +1,515 @@ +//! Compile a retained PromQL subtree (`Fallback`) from its typed expression. +//! The deployment supplies the raw series of each selector; the Planner +//! computes selection, range functions, subqueries, matching and aggregation. +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}; + +/// 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 +/// computed output, so the raw rows need another. +pub fn raw_series_input(node: NodeId, selector: usize) -> NodeId { + node | ((selector as u64 + 1) << 32) +} + +/// The Fallback node that owns a raw-series input slot. +pub(super) fn raw_series_owner(slot: NodeId) -> Option { + (slot >> 32 != 0).then_some(slot & u64::from(u32::MAX)) +} + +/// A selector expression and its raw-series row schema. +pub type Selector = (QueryExpr, Schema); + +/// The selectors a Fallback expression reads, left to right, and the row +/// schema of the raw series the deployment supplies for each at +/// [`raw_series_input`]. The rows must cover the selector's window at every +/// evaluation instant `T`, or at its `@` time: `(T - offset - range, T - offset]`; +/// under a subquery `[R:S] offset O` that is `(T - O - R - offset - range, T - O - offset]`. +pub fn raw_series(expression: &QueryExpr) -> Result, Error> { + Ok(lower(expression)?.selectors) +} + +/// An operator input: a selector's raw rows or an earlier step. +pub(super) enum Input { + Raw(usize), + Step(usize), +} + +/// Operators computing an expression; the last step is its result. +#[derive(Default)] +pub(super) struct Lowering { + pub selectors: Vec, + pub steps: Vec<(Operator, Vec)>, +} + +pub(super) fn lower(expression: &QueryExpr) -> Result { + let mut lowering = Lowering::default(); + lowering.value(expression)?; + Ok(lowering) +} + +fn declared(expression: &QueryExpr) -> Result { + let schema = expression + .output_schema() + .map_err(|error| invalid(error.to_string()))?; + Ok(Arc::new(lift_plain(&schema))) +} + +fn millis(duration: &std::time::Duration) -> Result { + i64::try_from(duration.as_millis()).map_err(|_| invalid("PromQL duration exceeds Int64")) +} + +/// A fixed `@` time. `start()`/`end()` depend on the deployment's range query. +fn at(shift: &planner_types::pre_asap::TimeShift) -> Result, Error> { + match shift.at { + None => Ok(None), + Some(AtModifier::Timestamp(at)) => Ok(Some(at)), + Some(_) => Err(invalid("@ start() and @ end() depend on the range query")), + } +} + +/// `TimeRange { range, [TimeShift { offset, @ }], Scan }`: range, offset, `@`. +fn selector(expression: &QueryExpr) -> Result<(i64, i64, Option), Error> { + let QueryExpr::TimeRange { range, child } = expression else { + return Err(invalid("PromQL operand must be a series selector")); + }; + let (offset, at, scan) = match child.as_ref() { + QueryExpr::TimeShift { shift, child } => (shift.offset_ms, at(shift)?, child.as_ref()), + scan => (0, None, scan), + }; + if !matches!(scan, QueryExpr::Scan { .. }) { + return Err(invalid("PromQL selector must read one scan")); + } + Ok((millis(range)?, offset, at)) +} + +/// PromQL scalar-valued expressions have no labels to match. +fn scalar(expression: &QueryExpr) -> bool { + matches!( + expression, + QueryExpr::PromqlScalarBridge(_) + | QueryExpr::PromqlScalarFromVector(_) + | QueryExpr::EvalTimestamp + ) +} + +impl Lowering { + fn schema(&self, input: &Input) -> Schema { + match input { + Input::Raw(i) => self.selectors[*i].1.clone(), + Input::Step(i) => self.steps[*i].0.schema(), + } + } + + fn add(&mut self, operator: Operator, inputs: Vec) -> Input { + self.steps.push((operator, inputs)); + Input::Step(self.steps.len() - 1) + } + + /// Conform `operator` to the logical schema of the expression it computes. + fn push( + &mut self, + operator: Operator, + inputs: Vec, + logical: &QueryExpr, + ) -> Result { + Ok(self.add(operator.with_output_schema(declared(logical)?)?, inputs)) + } + + fn read(&mut self, selector: &QueryExpr) -> Result { + let schema = declared(selector)?; + if !schema + .fields + .iter() + .any(|f| f.name == promql_rows::SERIES_IDENTITY_COLUMN) + { + return Err(invalid( + "PromQL fallback requires the complete series identity", + )); + } + self.selectors.push((selector.clone(), schema)); + Ok(Input::Raw(self.selectors.len() - 1)) + } + + /// An instant vector, or a scalar for scalar-valued expressions. + fn value(&mut self, expression: &QueryExpr) -> Result { + match expression { + QueryExpr::TimeRange { .. } => { + let (range, offset, at) = selector(expression)?; + let input = self.read(expression)?; + let schema = self.schema(&input); + self.push( + Operator::series_window(schema, None, range, offset, at, None)?, + vec![input], + expression, + ) + } + QueryExpr::Aggregate { + reduction: planner_types::pre_asap::Reduction::PerEntity, + measures, + having: None, + child, + .. + } => { + let [function] = measures.as_slice() else { + return Err(invalid("range function requires one measure")); + }; + self.range_function(function, child, expression) + } + QueryExpr::Aggregate { + reduction: planner_types::pre_asap::Reduction::Reduce(keys), + measures, + having: None, + child, + .. + } => { + let [measure] = measures.as_slice() else { + return Err(invalid("vector aggregation requires one measure")); + }; + let input = self.value(child)?; + self.aggregate(input, measure, keys, expression) + } + QueryExpr::Sort { + keys, + partition_by, + child, + } => { + let step = self.value(child)?; + let input = self.schema(&step); + let keys = keys + .iter() + .map(|key| match key.expr { + QueryExpr::Column(column) => Ok(SortKey { + column, + descending: !key.ascending, + nulls_first: key.nulls_first, + }), + _ => Err(invalid("sort key must be a column")), + }) + .collect::>()?; + let groups = groups(&input, partition_by)?; + self.push(Operator::sort(input, keys, groups)?, vec![step], expression) + } + QueryExpr::Limit { n, offset, child } => { + let step = self.value(child)?; + let input = self.schema(&step); + // `topk by (...)` partitions through the Sort it limits. + let groups = match child.as_ref() { + QueryExpr::Sort { partition_by, .. } => groups(&input, partition_by)?, + _ => vec![], + }; + self.push( + Operator::limit(input, *n as u64, *offset as u64, groups)?, + vec![step], + expression, + ) + } + QueryExpr::BinaryOp { + op: BinaryOpKind::Arithmetic(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 operator = planner_types::post_asap::BinaryOperator { + kind: planner_types::pre_asap::BinaryOpKind::Arithmetic(op.clone()), + vector_match: None, + 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) + } + QueryExpr::PromqlScalarFromVector(child) => { + let step = self.value(child)?; + let input = self.schema(&step); + let value = named_column(&input, &ColumnRef::SampleValue)?; + self.push( + Operator::vector_to_scalar(input, value)?, + vec![step], + expression, + ) + } + QueryExpr::PromqlVectorFromScalar(child) => { + let step = self.value(child)?; + let input = self.schema(&step); + Ok(self.add( + Operator::scope_timestamp(input, declared(expression)?)?, + vec![step], + )) + } + QueryExpr::PromqlScalarBridge(_) => { + let value = row_values::scalar_literal(expression) + .ok_or_else(|| invalid("PromQL scalar must be a literal"))?; + self.push( + Operator::scalar(crate::values::Value::Float64(value), DataType::Float64)?, + vec![], + expression, + ) + } + _ => Err(invalid("PromQL expression has no native fallback 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, + function: &AggIntent, + matrix: &QueryExpr, + logical: &QueryExpr, + ) -> Result { + let function = unbound(function)?; + let (subquery, offset, at_ms) = match matrix { + QueryExpr::TimeShift { shift, child } => (child.as_ref(), shift.offset_ms, at(shift)?), + other => (other, 0, None), + }; + let QueryExpr::PromqlSubquery { + range: outer, + resolution, + child, + } = subquery + else { + let (range, offset, at) = selector(matrix)?; + let input = self.read(matrix)?; + let schema = self.schema(&input); + return self.push( + Operator::series_window(schema, Some(function), range, offset, at, None)?, + vec![input], + logical, + ); + }; + let step = resolution.as_ref().ok_or_else(|| { + invalid("subquery resolution defaults to the deployment evaluation interval") + })?; + let steps = SubquerySteps { + range_ms: millis(outer)?, + step_ms: millis(step)?, + offset_ms: offset, + at_ms, + }; + // Each step evaluates a per-series selection or range function. + let (inner, selected) = match child.as_ref() { + QueryExpr::Aggregate { + reduction: planner_types::pre_asap::Reduction::PerEntity, + measures, + having: None, + child: selected, + .. + } => match measures.as_slice() { + [inner] => (Some(unbound(inner)?), selected.as_ref()), + _ => return Err(invalid("range function requires one measure")), + }, + selected => (None, selected), + }; + 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))?, + vec![raw], + child, + )?; + let input = self.schema(&step); + self.push( + Operator::series_window(input, Some(function), steps.range_ms, offset, at_ms, None)?, + vec![step], + logical, + ) + } + + /// Cross-series aggregation. A global aggregate groups by one constant so + /// that no input series yields an empty vector, not one row. + fn aggregate( + &mut self, + mut step: Input, + measure: &AggIntent, + keys: &GroupKeys, + logical: &QueryExpr, + ) -> Result { + let mut input = self.schema(&step); + let value = named_column(&input, &ColumnRef::SampleValue)?; + let reduction = match measure { + AggIntent::Sum { col: None } => Reduction::Sum(value), + AggIntent::Avg { col: None } => Reduction::Avg(value), + AggIntent::Min { col: None } => Reduction::Min(value), + AggIntent::Max { col: None } => Reduction::Max(value), + AggIntent::Count { .. } => Reduction::Count, + _ => return Err(invalid("vector aggregate has no native lowering")), + }; + let mut groups = if keys.is_without() { + // Group by every remaining label, including the rewritten identity. + let excluded = keys.keys(); + if excluded.iter().any(|&i| i >= input.fields.len()) { + return Err(invalid("grouping column out of range")); + } + let names = excluded.iter().map(|&i| input.fields[i].name.clone()); + let relabel = + Operator::series_labels(input.clone(), VectorMatchKind::Ignoring, names.collect())?; + step = self.add(relabel, vec![step]); + (0..input.fields.len()) + .filter(|&i| Some(i) != input.time_index && i != value && !excluded.contains(&i)) + .collect() + } else { + groups(&input, keys)? + }; + let global = groups.is_empty(); + if global { + let mut columns = (0..input.fields.len()) + .map(|i| (input.fields[i].name.clone(), Expression::Column(i))) + .collect::>(); + columns.push(( + "$promql_global_group".into(), + Expression::Literal { + value: crate::values::Value::Utf8("".into()), + dtype: DataType::Utf8, + }, + )); + let project = Operator::project(input, columns)?; + input = project.schema(); + groups = vec![input.fields.len() - 1]; + step = self.add(project, vec![step]); + } + let output = declared(logical)?; + let name = output + .fields + .last() + .ok_or_else(|| invalid("aggregate output lacks a value"))? + .name + .clone(); + let aggregate = Operator::aggregate(input, groups, vec![(name, reduction)])?; + let actual = aggregate.schema(); + let step = self.add(aggregate, vec![step]); + // Drop the constant group; convert counts where PromQL declares Float64. + let skip = usize::from(global); + let columns = actual.fields[skip..] + .iter() + .zip(&output.fields) + .enumerate() + .map(|(i, (field, declared))| { + let column = i + skip; + let expression = if field.dtype != declared.dtype { + Expression::ExactFloat64(column) + } else { + Expression::Column(column) + }; + (field.name.clone(), expression) + }) + .collect(); + self.push(Operator::project(actual, columns)?, vec![step], logical) + } +} + +fn unbound(intent: &AggIntent) -> Result, Error> { + Ok(match intent { + AggIntent::Rate => AggIntent::Rate, + AggIntent::Increase => AggIntent::Increase, + AggIntent::Delta => AggIntent::Delta, + AggIntent::Count { accuracy } => AggIntent::Count { + accuracy: accuracy.clone(), + }, + AggIntent::Sum { col: None } => AggIntent::Sum { col: None }, + AggIntent::Avg { col: None } => AggIntent::Avg { col: None }, + AggIntent::Min { col: None } => AggIntent::Min { col: None }, + AggIntent::Max { col: None } => AggIntent::Max { col: None }, + AggIntent::IRate => AggIntent::IRate, + AggIntent::IDelta => AggIntent::IDelta, + AggIntent::Changes => AggIntent::Changes, + AggIntent::Resets => AggIntent::Resets, + AggIntent::LastOverTime => AggIntent::LastOverTime, + AggIntent::Quantile { + col: None, + q, + accuracy, + } => AggIntent::Quantile { + col: None, + q: *q, + accuracy: accuracy.clone(), + }, + _ => return Err(invalid("unsupported PromQL range function")), + }) +} diff --git a/crates/asap-physical-operators/src/physical_planner/row_values.rs b/crates/asap-physical-operators/src/physical_planner/row_values.rs new file mode 100644 index 00000000..737c8065 --- /dev/null +++ b/crates/asap-physical-operators/src/physical_planner/row_values.rs @@ -0,0 +1,227 @@ +//! 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; + +/// A PromQL number literal has no row schema; its consumer folds it in. +pub(super) fn scalar_literal(expression: &QueryExpr) -> Option { + match expression { + QueryExpr::PromqlScalarBridge(child) => scalar_literal(child), + QueryExpr::Literal(ScalarValue::Float64(value)) => Some(*value), + _ => None, + } +} + +/// 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, + grouping: &[String], + readout: &PopulationReadout, +) -> Result, Error> { + let groups = grouping + .iter() + .map(|name| named_column(input, &ColumnRef::Named(name.clone()))) + .collect::, _>>()?; + let value = named_column(input, &ColumnRef::SampleValue)?; + let reduction = match readout { + PopulationReadout::Sum => Reduction::Sum(value), + PopulationReadout::Count => Reduction::Count, + PopulationReadout::Average => Reduction::Avg(value), + PopulationReadout::Quantile { q } => Reduction::Quantile { + column: value, + q: *q, + }, + PopulationReadout::TopK { .. } => { + return Err(invalid( + "TopK population readout ranks; it does not aggregate", + )) + } + }; + if !groups.is_empty() { + return Ok(vec![Operator::aggregate( + input.clone(), + groups, + vec![("value".into(), reduction)], + )?]); + } + // A global aggregate over no members is an empty PromQL vector, not one row. + let aggregate = Operator::aggregate( + input.clone(), + vec![], + vec![ + ("value".into(), reduction), + ("members".into(), Reduction::Count), + ], + )?; + let zero = Expression::Literal { + value: crate::values::Value::Int64(0), + dtype: DataType::Int64, + }; + let filter = Operator::filter( + aggregate.schema(), + Expression::Less(Box::new(zero), Box::new(Expression::Column(1))), + )?; + let project = Operator::project( + filter.schema(), + vec![("value".into(), Expression::Column(0))], + )?; + Ok(vec![aggregate, filter, project]) +} diff --git a/crates/asap-physical-operators/tests/deployment_computation.rs b/crates/asap-physical-operators/tests/deployment_computation.rs new file mode 100644 index 00000000..b026f47f --- /dev/null +++ b/crates/asap-physical-operators/tests/deployment_computation.rs @@ -0,0 +1,340 @@ +//! Planner-selected PromQL computation compiles from the timed DAG alone; +//! the deployment supplies only raw rows at the ingestion frontier. +use asap_physical_operators::{ + operators::Operator, + physical_planner::{compile, promql_rows, CompiledPhysicalDag, InputContract, Source}, + runtime::{Limits, RunContext, Scope}, + values::{Batch, Value}, +}; +use futures::{executor::block_on, StreamExt}; +use planner_types::{post_asap::*, pre_asap::QueryExpr, types::AccuracyTarget, workload::*}; +use std::{collections::BTreeMap, rc::Rc, sync::Arc}; + +fn lower(query: &str) -> 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), + ..Default::default() + }, + predictability: Predictability::Unknown, + invocations: 1, + execute_at: None, + time_selection: TimeSelection::default(), + }]), + repeating_queries: None, + }, + data_workload: Some(DataWorkload { + data_ingestion_interval: Evidence { + value: Some(DurationMs(60_000)), + ..Default::default() + }, + ..Default::default() + }), + }; + asap_frontend_promql::lower_promql_workload(&workload, 0) + .unwrap() + .remove(0) +} + +/// The first exact summary candidate, as Planner selection would hand it over. +fn exact_dag(query: &str) -> PostAsapDag { + use asap_aware_mapping::{Replacement, ReplacementStrategy, TargetSubDAG}; + let expression = lower(query); + let root = Rc::new(promql_rows::with_series_identity(&expression).unwrap_or(expression)); + asap_aware_mapping::SketchAlgorithmStrategy::new(&asap_aware_mapping::DefaultCostModel) + .replacements(&TargetSubDAG::new(&root)) + .into_iter() + .find_map(|candidate| match candidate.replacement { + Replacement::Summary(node) => { + let dag = compile_post_asap_dag(&node).ok()?; + dag.nodes + .iter() + .all(|n| !matches!(&n.payload, PostAsapOperatorPayload::SummaryAgg { family, .. } if !matches!(family, SummaryFamilyType::ExactAggregate(..)))) + .then_some(dag) + } + _ => None, + }) + .unwrap() +} + +fn population_dag(query: &str) -> PostAsapDag { + let root = Rc::new(promql_rows::with_series_identity(&lower(query)).unwrap()); + let selected = asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new( + std::slice::from_ref(&root), + ) + .candidate(&root) + .unwrap(); + compile_post_asap_dag(&selected).unwrap() +} + +/// Raw scan nodes are the frontier; everything above them is compiled. +fn raw_inputs(dag: &PostAsapDag) -> Vec<(u64, Arc, String)> { + dag.nodes + .iter() + .filter_map(|node| match &node.payload { + PostAsapOperatorPayload::Fallback { + expression: QueryExpr::TimeRange { child, .. }, + } => match child.as_ref() { + QueryExpr::Scan { + source: planner_types::pre_asap::Source::TimeSeries { metric }, + .. + } => Some(( + u64::from(node.id.0), + Arc::new(node.output_schema.clone()), + metric.clone(), + )), + _ => None, + }, + _ => None, + }) + .collect() +} + +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> { + let inputs = raw_inputs(dag); + let program = compile( + dag, + inputs + .iter() + .map(|(id, schema, _)| (*id, InputContract::bounded(schema.clone()))) + .collect(), + &[u64::from(dag.root.0)], + ) + .map_err(|e| e.to_string())?; + let program: CompiledPhysicalDag = + serde_json::from_slice(&serde_json::to_vec(&program).unwrap()).unwrap(); + let sources = inputs + .iter() + .map(|(id, schema, metric)| { + let rows = samples + .iter() + .filter(|sample| sample.0 == metric) + .map(|(name, job, instance, at, value)| { + let labels = BTreeMap::from([ + ("__name__".to_string(), name.to_string()), + ("job".into(), job.to_string()), + ("instance".into(), instance.to_string()), + ]); + if schema + .fields + .iter() + .any(|f| f.name == promql_rows::SERIES_IDENTITY_COLUMN) + { + promql_rows::series_row(schema, &labels, *at, *value).unwrap() + } else { + schema + .fields + .iter() + .enumerate() + .map(|(i, f)| match f.name.as_str() { + _ if Some(i) == schema.time_index => Value::Timestamp(*at), + "value" => Value::Float64(*value), + label => Value::Utf8(labels[label].clone().into()), + }) + .collect() + } + }) + .collect(); + let batch = Batch::try_new(schema.clone(), rows).unwrap(); + ( + *id, + Box::new(Operator::source(schema.clone(), vec![batch]).unwrap()) as Source<'_>, + ) + }) + .collect(); + let graph = program.instantiate(sources).map_err(|e| e.to_string())?; + let context = RunContext::new( + Scope::Query { + evaluation_time_ms: end, + revision: 0, + }, + Limits::default(), + ) + .unwrap(); + block_on(async { + let mut stream = graph + .execute(program.roots(), context) + .map_err(|e| e.to_string())? + .remove(0); + let mut rows = BTreeMap::new(); + while let Some(batch) = stream.next().await { + let batch = batch.map_err(|e| e.to_string())?; + 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.), + ("m", "api", "b", 40_000, 7.), + ("m", "api", "c", 30_000, 2.), + ("m", "db", "d", 20_000, 5.), +]; + +fn reference(pairs: &[(&str, f64)]) -> BTreeMap { + pairs.iter().map(|(k, v)| (k.to_string(), *v)).collect() +} + +// Current-series aggregates read the latest member values, matching the +// backend CurrentSeriesStore formulas (PromQL quantile interpolation). +#[test] +fn population_aggregates_match_current_series_reference() { + // Latest values: api = {a: 1, b: 7, c: 2}; db = {d: 5}. + for (query, expected) in [ + ("sum by (job) (m)", reference(&[("api", 10.), ("db", 5.)])), + ("count by (job) (m)", reference(&[("api", 3.), ("db", 1.)])), + ( + "avg by (job) (m)", + reference(&[("api", 10. / 3.), ("db", 5.)]), + ), + // Sorted api = [1, 2, 7]; rank 0.25 * 2 = 0.5 → 1.5. + ( + "quantile by (job) (0.25, m)", + reference(&[("api", 1.5), ("db", 5.)]), + ), + ] { + let dag = population_dag(query); + assert_eq!(run(&dag, SAMPLES, 60_000).unwrap(), expected, "{query}"); + } +} + +// A global readout of an empty population is an empty vector, as in PromQL. +#[test] +fn global_population_aggregate_of_no_members_is_empty() { + // Latest values are [1, 2, 5, 7] at 60s; every member has expired by 1000s. + for (query, expected) in [("sum(m)", 15.), ("count(m)", 4.), ("quantile(0.5, m)", 3.5)] { + let dag = population_dag(query); + let live = run(&dag, SAMPLES, 60_000).unwrap(); + assert_eq!(live, reference(&[("", expected)]), "{query}"); + assert!(run(&dag, SAMPLES, 1_000_000).unwrap().is_empty(), "{query}"); + } +} + +// Scalar operands on either side apply to every grouped value, including negation. +#[test] +fn scalar_literal_arithmetic_applies_to_grouped_values() { + // sum_over_time over 5m per job: api = 4 + 1 + 7 + 2 = 14, db = 5. + for (query, expected) in [ + ( + "sum by (job) (sum_over_time(m[5m])) * 2", + reference(&[("api", 28.), ("db", 10.)]), + ), + ( + "100 - sum by (job) (sum_over_time(m[5m]))", + reference(&[("api", 86.), ("db", 95.)]), + ), + ( + "-sum by (job) (sum_over_time(m[5m]))", + reference(&[("api", -14.), ("db", -5.)]), + ), + ] { + assert_eq!( + run(&exact_dag(query), SAMPLES, 60_000).unwrap(), + expected, + "{query}" + ); + } +} + +// Grouped vectors match one-to-one on labels; unmatched groups are dropped and +// unchecked division by zero yields +Inf as in PromQL. +#[test] +fn grouped_vector_arithmetic_matches_labels() { + let samples: &[Sample] = &[ + ("a", "api", "x", 10_000, 6.), + ("a", "api", "y", 20_000, 3.), + ("a", "db", "x", 10_000, 1.), + ("a", "web", "x", 10_000, 1.), + ("b", "api", "x", 10_000, 3.), + ("b", "db", "x", 10_000, 0.), + ("b", "cache", "x", 10_000, 1.), + ]; + let dag = + exact_dag("sum by (job) (sum_over_time(a[5m])) / sum by (job) (sum_over_time(b[5m]))"); + assert_eq!( + run(&dag, samples, 60_000).unwrap(), + reference(&[("api", 3.), ("db", f64::INFINITY)]) + ); +} + +// Exact observation counts finalize to the Float64 value PromQL declares, +// then roll up per job: api has 2 + 1 + 1 samples in 5m, db has 1. +#[test] +fn exact_count_finalizes_to_declared_float_value() { + let mut dag = exact_dag("sum by (job) (count_over_time(m[5m]))"); + let finalize = dag + .nodes + .iter() + .find(|node| { + matches!( + node.payload, + PostAsapOperatorPayload::Value { + operation: ValueOperation::FinalizeExactAccumulator + } + ) + }) + .unwrap() + .clone(); + let root = dag.nodes.iter().find(|n| n.id == dag.root).unwrap().clone(); + let mut edge = dag + .edges + .iter() + .find(|e| e.producer == finalize.id) + .unwrap() + .clone(); + // Read the rolled-up exact state the same way the query path does. + let mut read = finalize.clone(); + read.id = PostAsapNodeId(root.id.0 + 1); + read.output_schema = root.output_schema.clone(); + read.output_schema.fields.last_mut().unwrap().dtype = + SummaryFamilyType::Plain(planner_types::pre_asap::DataType::Float64); + edge.producer = root.id; + edge.consumer = read.id; + edge.intermediate_schema = root.output_schema.clone(); + edge.data_state = root.output_state; + dag.root = read.id; + dag.nodes.push(read); + dag.edges.push(edge); + assert_eq!( + run(&dag, SAMPLES, 60_000).unwrap(), + reference(&[("api", 4.), ("db", 1.)]) + ); +} + +// Comparisons need filter/bool semantics that `Binary` does not carry, so +// they fail at compile time instead of emitting 0/1 values. +#[test] +fn row_comparison_fails_closed() { + let mut dag = exact_dag("sum by (job) (sum_over_time(m[5m])) * 2"); + 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, + ); + } + } + let error = run(&dag, SAMPLES, 60_000).unwrap_err(); + assert!(error.contains("comparison"), "{error}"); +} diff --git a/crates/asap-physical-operators/tests/physical_dag.rs b/crates/asap-physical-operators/tests/physical_dag.rs index 96e74e54..652880e0 100644 --- a/crates/asap-physical-operators/tests/physical_dag.rs +++ b/crates/asap-physical-operators/tests/physical_dag.rs @@ -546,7 +546,9 @@ fn bind_post_asap_before_execution() { }; let native = bind(&dag, sources(), &[1]).unwrap(); assert_eq!(floats(&run(&native, 1, query()), 0), vec![3.]); - assert!(bind(&dag, BTreeMap::new(), &[1]).is_err()); + // A literal Fallback needs no deployment input. + let literal = bind(&dag, BTreeMap::new(), &[1]).unwrap(); + assert_eq!(floats(&run(&literal, 1, query()), 0), vec![3.]); dag.nodes[1].payload = PostAsapOperatorPayload::Value { operation: ValueOperation::Extension { name: "unknown".into(), diff --git a/crates/asap-physical-operators/tests/promql_fallback.rs b/crates/asap-physical-operators/tests/promql_fallback.rs new file mode 100644 index 00000000..d274495c --- /dev/null +++ b/crates/asap-physical-operators/tests/promql_fallback.rs @@ -0,0 +1,624 @@ +//! A retained PromQL subtree (`Fallback`) compiles from its typed expression. +//! The deployment supplies only its selector's raw series; expected values are +//! hand-computed with Prometheus semantics. +use asap_physical_operators::{ + operators::Operator, + physical_planner::{compile, promql_fallback, promql_rows, CompiledPhysicalDag, InputContract}, + runtime::{Limits, RunContext, Scope}, + values::{Batch, Value}, +}; +use futures::{executor::block_on, StreamExt}; +use planner_types::{ + post_asap::{execution_data_state::lift_plain, *}, + pre_asap::QueryExpr, + types::AccuracyTarget, + workload::*, +}; +use std::{collections::BTreeMap, rc::Rc}; + +/// Bare selectors look back one ingestion interval: 60s. +fn parse(query: &str) -> 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), + ..Default::default() + }, + predictability: Predictability::Unknown, + invocations: 1, + execute_at: None, + time_selection: TimeSelection::default(), + }]), + repeating_queries: None, + }, + data_workload: Some(DataWorkload { + data_ingestion_interval: Evidence { + value: Some(DurationMs(60_000)), + ..Default::default() + }, + ..Default::default() + }), + }; + asap_frontend_promql::lower_promql_workload(&workload, 0) + .unwrap() + .remove(0) +} + +fn lower(query: &str) -> QueryExpr { + promql_rows::with_series_identity(&parse(query)).unwrap() +} + +/// The whole query retained as one pre-ASAP node. +fn fallback_dag(expression: QueryExpr) -> PostAsapDag { + let schema = lift_plain(&expression.output_schema().unwrap()); + compile_post_asap_dag(&Rc::new(SummaryNode { + expr: SummaryExpr::KeepPreAsap(Rc::new(expression)), + schema, + guarantee: None, + })) + .unwrap() +} + +/// `(labels, seconds, value)`. `labels` is `k=v,...`, or a bare `job` value. +type Sample = (&'static str, i64, f64); + +fn labels(spec: &str) -> BTreeMap { + if !spec.contains('=') { + return BTreeMap::from([("job".into(), spec.into())]); + } + spec.split(',') + .map(|pair| { + let (k, v) = pair.split_once('=').unwrap(); + (k.to_string(), v.to_string()) + }) + .collect() +} + +/// The metric a selector reads. +fn metric(selector: &QueryExpr) -> String { + match selector { + QueryExpr::Scan { + source: planner_types::pre_asap::Source::TimeSeries { metric }, + .. + } => metric.clone(), + QueryExpr::TimeRange { child, .. } | QueryExpr::TimeShift { child, .. } => metric(child), + other => panic!("not a selector: {other:?}"), + } +} + +fn compile_query(query: &str) -> Result { + let expression = lower(query); + let dag = fallback_dag(expression.clone()); + let root = u64::from(dag.root.0); + let inputs = promql_fallback::raw_series(&expression) + .map_err(|e| e.to_string())? + .into_iter() + .enumerate() + .map(|(i, (_, schema))| { + ( + promql_fallback::raw_series_input(root, i), + InputContract::bounded(schema), + ) + }) + .collect(); + let program = compile(&dag, inputs, &[root]).map_err(|e| e.to_string())?; + Ok(serde_json::from_slice(&serde_json::to_vec(&program).unwrap()).unwrap()) +} + +/// Evaluate at `at` seconds over samples of each named metric; returns +/// `(output labels, timestamp ms, value)` rows in order. +#[allow(clippy::type_complexity)] +fn evaluate( + query: &str, + metrics: &[(&str, &[Sample])], + at: i64, +) -> Result, i64, f64)>, String> { + let expression = lower(query); + let program = compile_query(query)?; + let mut sources = BTreeMap::new(); + let selectors = promql_fallback::raw_series(&expression).unwrap(); + for (i, (selector, schema)) in selectors.into_iter().enumerate() { + let name = metric(&selector); + let rows = metrics + .iter() + .filter(|(m, _)| *m == name) + .flat_map(|(_, samples)| samples.iter()) + .map(|(spec, seconds, value)| { + let mut labels = labels(spec); + labels.insert("__name__".into(), name.clone()); + promql_rows::series_row(&schema, &labels, seconds * 1000, *value).unwrap() + }) + .collect(); + let batch = Batch::try_new(schema.clone(), rows).unwrap(); + sources.insert( + promql_fallback::raw_series_input(program.roots()[0], i), + Box::new(Operator::source(schema, vec![batch]).unwrap()) as _, + ); + } + let graph = program.instantiate(sources).map_err(|e| e.to_string())?; + let context = RunContext::new( + Scope::Query { + evaluation_time_ms: at * 1000, + revision: 0, + }, + Limits::default(), + ) + .unwrap(); + block_on(async { + let mut stream = graph + .execute(program.roots(), context) + .map_err(|e| e.to_string())? + .remove(0); + let mut rows = Vec::new(); + while let Some(batch) = stream.next().await { + let batch = batch.map_err(|e| e.to_string())?; + let schema = batch.schema().clone(); + for row in batch.rows() { + let mut labels = BTreeMap::new(); + let mut time = -1; + let mut value = None; + for (field, cell) in schema.fields.iter().zip(row) { + match (field.name.as_str(), cell) { + (promql_rows::SERIES_IDENTITY_COLUMN, Value::Utf8(id)) => { + labels = promql_rows::decode_series_identity(id).unwrap() + } + (_, Value::Utf8(_) | Value::Null) => {} + (_, Value::Timestamp(t)) => time = *t, + (_, Value::Float64(v)) => value = Some(*v), + (_, Value::Int64(v)) => value = Some(*v as f64), + other => return Err(format!("unexpected cell {other:?}")), + } + } + if !schema + .fields + .iter() + .any(|f| f.name == promql_rows::SERIES_IDENTITY_COLUMN) + { + for (field, cell) in schema.fields.iter().zip(row) { + if let Value::Utf8(label) = cell { + if !label.is_empty() { + labels.insert(field.name.clone(), label.to_string()); + } + } + } + } + rows.push((labels, time, value.ok_or("missing value")?)); + } + } + Ok(rows) + }) +} + +/// Evaluate at `at` seconds over metric `m`; returns `(job or "", timestamp ms, value)`. +fn run(query: &str, samples: &[Sample], at: i64) -> Result, String> { + Ok(evaluate(query, &[("m", samples)], at)? + .into_iter() + .map(|(labels, time, value)| (labels.get("job").cloned().unwrap_or_default(), time, value)) + .collect()) +} + +/// Output rows as `(k=v,... sorted, value)`, including any `__name__`. +fn labeled(query: &str, metrics: &[(&str, &[Sample])], at: i64) -> Vec<(String, f64)> { + let mut rows = evaluate(query, metrics, at) + .unwrap_or_else(|e| panic!("{query}: {e}")) + .into_iter() + .map(|(labels, _, value)| { + let spec = labels + .iter() + .map(|(k, v)| format!("{k}={v}")) + .collect::>() + .join(","); + (spec, value) + }) + .collect::>(); + rows.sort_by(|a, b| a.0.cmp(&b.0)); + rows +} + +fn values(query: &str, samples: &[Sample], at: i64) -> Vec<(String, f64)> { + run(query, samples, at) + .unwrap_or_else(|e| panic!("{query}: {e}")) + .into_iter() + .map(|(job, _, value)| (job, value)) + .collect() +} + +fn one(query: &str, samples: &[Sample], at: i64) -> f64 { + match values(query, samples, at).as_slice() { + [(_, value)] => *value, + other => panic!("{query}: expected one sample, got {other:?}"), + } +} + +const COUNTER: &[Sample] = &[ + ("a", 60, 10.), + ("a", 120, 20.), + ("a", 180, 5.), + ("a", 240, 15.), +]; + +// rate/increase correct the reset at 180s and extrapolate half an interval at +// most; delta treats the same samples as a gauge. +#[test] +fn range_functions_follow_prometheus_extrapolation_and_resets() { + // Reset-corrected increase is 25 over 180s of samples; 60s on each side extrapolates. + let increase = 25. * (180. + 60. + 60.) / 180.; + assert!((one("increase(m[5m])", COUNTER, 300) - increase).abs() < 1e-9); + assert!((one("rate(m[5m])", COUNTER, 300) - increase / 300.).abs() < 1e-12); + let delta = 5. * (180. + 60. + 60.) / 180.; + assert!((one("delta(m[5m])", COUNTER, 300) - delta).abs() < 1e-9); + // Fewer than two samples yield no rate. + assert!(values("rate(m[2m])", COUNTER, 300).is_empty()); + for (query, expected) in [ + ("sum_over_time(m[5m])", 50.), + ("avg_over_time(m[5m])", 12.5), + ("min_over_time(m[5m])", 5.), + ("max_over_time(m[5m])", 20.), + ("count_over_time(m[5m])", 4.), + ] { + assert_eq!(one(query, COUNTER, 300), expected, "{query}"); + } +} + +// Ranges are left-open: a sample at `t - range` is excluded, one at `t` is included. +#[test] +fn ranges_exclude_their_start_and_offsets_shift_them() { + let samples = &[("a", 60, 1.), ("a", 90, 1.), ("a", 120, 1.), ("a", 150, 1.)]; + assert_eq!(one("count_over_time(m[1m])", samples, 120), 2.); + // offset 1m reads (60s, 120s] at 180s; output keeps the evaluation time. + let rows = run("count_over_time(m[1m] offset 1m)", samples, 180).unwrap(); + assert_eq!(rows, vec![("a".into(), 180_000, 2.)]); +} + +// A bare selector takes the latest sample within the lookback; a stale marker +// hides the series rather than exposing an older value. +#[test] +fn instant_selection_uses_lookback_and_stale_markers() { + let stale = f64::from_bits(0x7ff0_0000_0000_0002); + let samples = &[("a", 0, 1.), ("a", 30, 2.), ("b", 30, 3.), ("b", 50, stale)]; + assert_eq!(values("m", samples, 60), vec![("a".into(), 2.)]); + // The lookback (30s, 90s] excludes the sample at 30s. + assert!(values("m", samples, 90).is_empty()); + // Range functions skip stale markers. + assert_eq!( + values("sum_over_time(m[1m])", samples, 60), + vec![("a".into(), 2.), ("b".into(), 3.)] + ); +} + +// NaN samples follow Prometheus: min/max skip them, sums propagate them. +#[test] +fn nan_samples() { + let samples = &[("a", 10, f64::NAN), ("a", 20, 3.), ("a", 30, 1.)]; + assert_eq!(one("max_over_time(m[1m])", samples, 60), 3.); + assert_eq!(one("min_over_time(m[1m])", samples, 60), 1.); + assert!(one("sum_over_time(m[1m])", samples, 60).is_nan()); +} + +// Aggregation over no series is an empty vector, not one zero or null row; +// sort_desc orders the selected series. +#[test] +fn cross_series_aggregates_and_empty_inputs() { + let samples = &[("a", 50, 1.), ("b", 40, 2.), ("b", 55, 4.)]; + assert_eq!(values("sum(m)", samples, 60), vec![(String::new(), 5.)]); + assert_eq!(values("count(m)", samples, 60), vec![(String::new(), 2.)]); + assert_eq!( + values("max by (job) (m)", samples, 60), + vec![("a".into(), 1.), ("b".into(), 4.)] + ); + assert_eq!( + values("sort_desc(m)", samples, 60), + vec![("b".into(), 4.), ("a".into(), 1.)] + ); + // topk by (job) keeps the top series of each job, not one overall. + let jobs = &[("a", 50, 1.), ("b", 50, 2.)]; + let mut top = values("topk by (job) (1, m)", jobs, 60); + top.sort_by(|x, y| x.0.cmp(&y.0)); + assert_eq!(top, vec![("a".into(), 1.), ("b".into(), 2.)]); + assert_eq!(values("topk(1, m)", jobs, 60), vec![("b".into(), 2.)]); + for query in ["sum(m)", "count(m)", "max(m)", "sum by (job) (rate(m[5m]))"] { + assert!(values(query, &[], 60).is_empty(), "{query}"); + } +} + +// scalar() is the single series' value and NaN otherwise; vector() needs no input. +#[test] +fn scalar_and_vector_bridges() { + assert_eq!(one("scalar(m)", &[("a", 50, 7.)], 60), 7.); + assert!(one("scalar(m)", &[("a", 50, 7.), ("b", 50, 8.)], 60).is_nan()); + assert!(one("scalar(m)", &[], 60).is_nan()); + assert_eq!( + run("vector(3)", &[], 60).unwrap(), + vec![(String::new(), 60_000, 3.)] + ); + assert_eq!( + values("2 - m", &[("a", 50, 7.)], 60), + vec![("a".into(), -5.)] + ); + assert_eq!( + values("m * 2", &[("a", 50, 7.)], 60), + vec![("a".into(), 14.)] + ); +} + +// Subquery steps are absolute multiples of the resolution in the left-open +// range; each step evaluates the operand, and the outer function reduces them. +#[test] +fn subqueries_evaluate_their_operand_on_the_aligned_grid() { + // Steps 60..300: selections 1, 7, 3, (none at 240s), 4. + let samples = &[ + ("a", 50, 1.), + ("a", 110, 7.), + ("a", 170, 3.), + ("a", 290, 4.), + ]; + assert_eq!(one("max_over_time(m[5m:1m])", samples, 300), 7.); + assert_eq!(one("count_over_time(m[5m:1m])", samples, 300), 4.); + // At 190s the steps are 60, 120, 180 (not 70, 130, 190): counts 1 + 2 + 2. + let samples = &[("a", 30, 1.), ("a", 90, 1.), ("a", 150, 1.), ("a", 185, 1.)]; + assert_eq!( + one("sum_over_time(count_over_time(m[2m])[3m:1m])", samples, 190), + 5. + ); + // offset 1m moves the grid to (-50s, 130s]: steps 0, 60, 120 count 0 + 1 + 2. + assert_eq!( + one( + "sum_over_time(count_over_time(m[2m])[3m:1m] offset 1m)", + samples, + 190 + ), + 3. + ); +} + +// Subquery work is bounded by the query: at most 100000 steps. +#[test] +fn dense_subquery_grids_are_rejected() { + assert!(compile_query("max_over_time(m[100s:1ms])").is_ok()); + assert!(compile_query("max_over_time(m[30d:1ms])").is_err()); +} + +// The deployment must supply the selector's raw rows under the documented slot +// with the exact selector schema; unsupported shapes stay rejected. +#[test] +fn raw_series_contract_is_explicit() { + let expression = lower("rate(m[5m])"); + let dag = fallback_dag(expression.clone()); + let root = u64::from(dag.root.0); + let [(selector, schema)] = promql_fallback::raw_series(&expression) + .unwrap() + .try_into() + .unwrap(); + assert!(matches!(selector, QueryExpr::TimeRange { .. })); + let missing = compile(&dag, BTreeMap::new(), &[root]).err().unwrap(); + assert!(missing.to_string().contains("raw series input")); + let mut wrong = (*schema).clone(); + wrong.fields.pop(); + let wrong = compile( + &dag, + BTreeMap::from([( + promql_fallback::raw_series_input(root, 0), + InputContract::bounded(std::sync::Arc::new(wrong)), + )]), + &[root], + ); + assert!(wrong.is_err()); + // A consumed bare selector is raw range rows for its consumer; it is not + // turned into instant selection. + let selector = lower("m"); + let schema = lift_plain(&selector.output_schema().unwrap()); + let node = |id, payload| PostAsapDagNode { + id: PostAsapNodeId(id), + payload, + output_state: ExecutionDataState::QUERY_ROWS, + output_schema: schema.clone(), + guarantee: None, + }; + let consumed = PostAsapDag { + nodes: vec![ + node( + 0, + PostAsapOperatorPayload::Fallback { + expression: selector.clone(), + }, + ), + node( + 1, + PostAsapOperatorPayload::Value { + operation: ValueOperation::Limit { + n: 1, + offset: 0, + partition_by: Default::default(), + }, + }, + ), + ], + edges: vec![PostAsapDagEdge { + producer: PostAsapNodeId(0), + consumer: PostAsapNodeId(1), + role: EdgeRole::Input, + intermediate_schema: schema.clone(), + data_state: ExecutionDataState::QUERY_ROWS, + grouping: GroupingEdgeCompatibility::NotApplicable, + window: WindowEdgeCompatibility::NotApplicable, + }], + root: PostAsapNodeId(1), + }; + let raw = promql_fallback::raw_series(&selector).unwrap().remove(0).1; + assert!(compile( + &consumed, + BTreeMap::from([( + promql_fallback::raw_series_input(0, 0), + InputContract::bounded(raw) + )]), + &[1], + ) + .is_err()); + // Implicit subquery resolution belongs to the deployment's evaluation interval. + assert!(compile_query("max_over_time(m[5m:])").is_err()); +} + +// irate/idelta use the last two samples (irate corrects a reset to the last +// value); changes/resets count value changes and decreases; quantile_over_time +// interpolates; all skip stale markers. +#[test] +fn instant_and_counting_range_functions() { + // COUNTER in (0s, 300s]: 10, 20, 5, 15. + assert!((one("irate(m[5m])", COUNTER, 300) - 10. / 60.).abs() < 1e-12); + assert_eq!(one("idelta(m[5m])", COUNTER, 300), 10.); + // At 200s the last pair 20 -> 5 is a reset: irate uses 5 as the increase. + assert!((one("irate(m[5m])", COUNTER, 200) - 5. / 60.).abs() < 1e-12); + assert_eq!(one("idelta(m[5m])", COUNTER, 200), -15.); + assert!(values("irate(m[1m])", COUNTER, 300).is_empty()); + assert_eq!(one("changes(m[5m])", COUNTER, 300), 3.); + assert_eq!(one("resets(m[5m])", COUNTER, 300), 1.); + assert_eq!(one("changes(m[2m])", COUNTER, 300), 0.); + // NaN to NaN is not a change; any other transition involving NaN is. + let flat = &[ + ("a", 10, 1.), + ("a", 20, 1.), + ("a", 30, 2.), + ("a", 40, f64::NAN), + ("a", 50, f64::NAN), + ("a", 55, 1.), + ]; + assert_eq!(one("changes(m[1m])", flat, 60), 3.); + let stale = f64::from_bits(0x7ff0_0000_0000_0002); + let ended = &[("a", 240, 15.), ("a", 250, stale)]; + assert_eq!(one("last_over_time(m[5m])", ended, 300), 15.); + assert!(values("m", ended, 300).is_empty()); + // Sorted 5, 10, 15, 20: rank 1.5 and 0.75; outside [0, 1] is +-Inf. + assert_eq!(one("quantile_over_time(0.5, m[5m])", COUNTER, 300), 12.5); + assert_eq!(one("quantile_over_time(0.25, m[5m])", COUNTER, 300), 8.75); + assert_eq!( + one("quantile_over_time(2, m[5m])", COUNTER, 300), + f64::INFINITY + ); + assert_eq!( + one("quantile_over_time(-1, m[5m])", COUNTER, 300), + f64::NEG_INFINITY + ); +} + +// `@ ` evaluates the selector or subquery at `t`, minus any offset, and the +// result keeps the query's evaluation time. +#[test] +fn at_modifier_fixes_the_evaluation_instant() { + let samples = &[("a", 60, 1.), ("a", 120, 2.), ("a", 180, 3.)]; + assert_eq!( + run("m @ 120", samples, 1000).unwrap(), + vec![("a".into(), 1_000_000, 2.)] + ); + assert!(values("m", samples, 1000).is_empty()); + assert_eq!(one("count_over_time(m[2m] @ 180)", samples, 1000), 2.); + assert_eq!(one("m @ 180 offset 1m", samples, 1000), 2.); + // The subquery grid is (60s, 180s]: steps 120 and 180 select 2 and 3. + assert_eq!(one("max_over_time(m[2m:1m] @ 180)", samples, 1000), 3.); + assert_eq!( + one("sum_over_time(m[2m:1m] @ 180 offset 1m)", samples, 1000), + 3. + ); + // An inner @ pins every step to the same instant. + assert_eq!(one("sum_over_time((m @ 60)[2m:1m])", samples, 180), 2.); + // start() and end() depend on the range query, which is the deployment's. + assert!(compile_query("m @ start()").is_err()); +} + +const A: &[Sample] = &[("job=x", 50, 10.), ("job=y", 50, 20.), ("job=w", 50, 0.)]; +const B: &[Sample] = &[("job=x", 50, 2.), ("job=z", 50, 5.), ("job=w", 50, 0.)]; + +// Vector-vector arithmetic matches series one-to-one on label sets without +// the metric name, and the result drops the metric name. +#[test] +fn vector_arithmetic_matches_label_sets() { + let metrics = &[("a", A), ("b", B)]; + let quotient = labeled("a / b", metrics, 60); + assert_eq!(quotient.len(), 2); + assert_eq!(quotient[0].0, "job=w"); + assert!(quotient[0].1.is_nan(), "0 / 0 is NaN"); + assert_eq!(quotient[1], ("job=x".into(), 5.)); + // Each selector reads its own raw rows, even a repeated metric. + assert_eq!( + labeled("(a - b) * a", metrics, 60), + vec![("job=w".into(), 0.), ("job=x".into(), 80.)] + ); + assert_eq!( + labeled("sum by (job) (a) - sum by (job) (b)", metrics, 60), + vec![("job=w".into(), 0.), ("job=x".into(), 8.)] + ); + // Rates of two counters over their own windows. + let up: &[Sample] = &[("job=x", 0, 0.), ("job=x", 60, 60.)]; + let down: &[Sample] = &[("job=x", 0, 0.), ("job=x", 60, 30.)]; + assert_eq!( + labeled("rate(a[2m]) / rate(b[2m])", &[("a", up), ("b", down)], 60), + vec![("job=x".into(), 2.)] + ); +} + +// on() keeps only the listed labels and ignoring() drops them; a duplicate +// match group is an error unless the left duplicates never match. +#[test] +fn on_and_ignoring_select_the_matching_labels() { + 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!(labeled("a - b", metrics, 60).is_empty()); + assert_eq!( + labeled("a - on(job) b", metrics, 60), + vec![("job=x".into(), 6.)] + ); + assert_eq!( + labeled("a - ignoring(inst) b", metrics, 60), + vec![("job=x".into(), 6.)] + ); + let pair: &[Sample] = &[("job=x,inst=1", 50, 1.), ("job=x,inst=2", 50, 2.)]; + let other: &[Sample] = &[("job=y", 50, 1.)]; + 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. +#[test] +fn without_grouping_drops_labels_and_the_name() { + let a: &[Sample] = &[ + ("job=x,inst=1", 50, 1.), + ("job=x,inst=2", 50, 2.), + ("job=y,inst=1", 50, 4.), + ]; + let metrics = &[("a", a)]; + assert_eq!( + labeled("sum without (inst) (a)", metrics, 60), + vec![("job=x".into(), 3.), ("job=y".into(), 4.)] + ); + assert_eq!( + labeled("count without (inst) (a)", metrics, 60), + vec![("job=x".into(), 2.), ("job=y".into(), 1.)] + ); + assert_eq!( + labeled("max without (job, inst) (a)", metrics, 60), + vec![(String::new(), 4.)] + ); + assert!(labeled("sum without (inst) (a)", &[], 60).is_empty()); +} + +// An empty label value is an absent label, and an empty side yields an empty +// result before any duplicate check, as in Prometheus. +#[test] +fn empty_labels_and_empty_sides_match_prometheus() { + let a: &[Sample] = &[("job=x,env=", 50, 3.)]; + let b: &[Sample] = &[("job=x", 50, 1.)]; + assert_eq!( + labeled("a + b", &[("a", a), ("b", b)], 60), + vec![("job=x".into(), 4.)] + ); + 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()); +} diff --git a/crates/frontend-promql/src/promql.rs b/crates/frontend-promql/src/promql.rs index c3df6c67..b71a8c1e 100644 --- a/crates/frontend-promql/src/promql.rs +++ b/crates/frontend-promql/src/promql.rs @@ -307,11 +307,23 @@ fn walk(expr: &Expr) -> Result { vector_match: None, }), }, - Expr::Subquery(sq) => Ok(Unresolved::PromqlSubquery { - range: sq.range, - resolution: sq.step, - child: Rc::new(walk(&sq.expr)?), - }), + Expr::Subquery(sq) => { + let subquery = Unresolved::PromqlSubquery { + range: sq.range, + resolution: sq.step, + child: Rc::new(walk(&sq.expr)?), + }; + // `offset`/`@` move the whole subquery, including its step grid. + let shift = time_shift(sq.offset.as_ref(), sq.at.as_ref())?; + Ok(if shift.is_identity() { + subquery + } else { + Unresolved::TimeShift { + shift, + child: Rc::new(subquery), + } + }) + } Expr::VectorSelector(vs) => { let (metric, matchers, shift) = vs_parts(vs)?; Ok(instant_source(metric, matchers, shift)) diff --git a/crates/frontend-promql/tests/promql_lowering.rs b/crates/frontend-promql/tests/promql_lowering.rs index 2dd56173..3c92b7ba 100644 --- a/crates/frontend-promql/tests/promql_lowering.rs +++ b/crates/frontend-promql/tests/promql_lowering.rs @@ -1202,3 +1202,20 @@ fn histogram_quantiles_rejects_an_out_of_range_quantile() { ); } } + +// A subquery's `offset`/`@` shift the whole subquery, so the tree keeps them. +#[test] +fn subquery_time_shift_is_retained() { + let QueryExpr::Aggregate { child, .. } = lower("max_over_time(m[5m:1m] offset 1m)") else { + panic!("expected a range function"); + }; + let QueryExpr::TimeShift { shift, child } = child.as_ref() else { + panic!("subquery offset was dropped: {child:?}"); + }; + assert_eq!(shift.offset_ms, 60_000); + assert!(matches!(child.as_ref(), QueryExpr::PromqlSubquery { .. })); + assert!(matches!( + lower("max_over_time(m[5m:1m] @ 100)"), + QueryExpr::Aggregate { child, .. } if matches!(child.as_ref(), QueryExpr::TimeShift { .. }) + )); +} diff --git a/crates/types/src/pre_asap/query_expr.rs b/crates/types/src/pre_asap/query_expr.rs index dec0107a..5b63e5a2 100644 --- a/crates/types/src/pre_asap/query_expr.rs +++ b/crates/types/src/pre_asap/query_expr.rs @@ -878,7 +878,8 @@ pub enum QueryExpr { /// unchanged. Wraps the shifted selector directly — `m offset 1h` → /// `TimeShift { Scan }`; a ranged selector `m[5m] offset 1h` → /// `TimeRange { 5m, TimeShift { Scan } }` (the range is taken at the shifted - /// time). Never carries the identity shift (the converter emits a bare + /// time). A shifted subquery wraps the `PromqlSubquery`, moving its step + /// grid. Never carries the identity shift (the converter emits a bare /// selector when neither modifier is present). TimeShift { shift: TimeShift, diff --git a/crates/types/src/pre_asap/schema.rs b/crates/types/src/pre_asap/schema.rs index 938b9b7f..fdac93de 100644 --- a/crates/types/src/pre_asap/schema.rs +++ b/crates/types/src/pre_asap/schema.rs @@ -168,8 +168,16 @@ pub const PROMQL_SERIES_IDENTITY: &str = "$promql_series_identity"; /// own realization; they must not accidentally treat the opaque identity as a /// user label or silently discard it. pub fn with_promql_series_identity(root: &super::QueryExpr) -> Result { - use super::{QueryExpr, Reduction, Source}; + use super::{QueryExpr, Source}; use std::rc::Rc; + fn scalar_literal(expression: &QueryExpr) -> Option { + match expression { + QueryExpr::PromqlScalarBridge(child) => scalar_literal(child), + QueryExpr::Literal(super::ScalarValue::Float64(value)) => Some(*value), + _ => None, + } + } + let mut root = root.clone(); fn visit(node: &mut QueryExpr) -> Result<(), String> { match node { QueryExpr::Scan { @@ -193,17 +201,31 @@ pub fn with_promql_series_identity(root: &super::QueryExpr) -> Result { - visit(Rc::make_mut(child)) - } - QueryExpr::Aggregate { - child, reduction, .. - } => { - if matches!(reduction, Reduction::Reduce(keys) if keys.is_without()) { - return Err("dynamic without grouping requires label-set projection".into()); - } - visit(Rc::make_mut(child)) + QueryExpr::TimeRange { child, .. } + | QueryExpr::Limit { child, .. } + | QueryExpr::TimeShift { child, .. } + | QueryExpr::PromqlSubquery { child, .. } + | QueryExpr::PromqlScalarFromVector(child) => visit(Rc::make_mut(child)), + // Constants read no series. + QueryExpr::PromqlScalarBridge(_) => 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 + ) + }) => + { + visit(Rc::make_mut(lhs))?; + visit(Rc::make_mut(rhs)) } + QueryExpr::Aggregate { child, .. } => visit(Rc::make_mut(child)), QueryExpr::Sort { child, partition_by, @@ -217,7 +239,6 @@ pub fn with_promql_series_identity(root: &super::QueryExpr) -> Result Err("operator has no dynamic series-identity realization".into()), } } - let mut root = root.clone(); visit(&mut root)?; root.output_schema().map_err(|error| error.to_string())?; Ok(root) diff --git a/docs/develop_docs/README.md b/docs/develop_docs/README.md index 2593f3f5..817722ad 100644 --- a/docs/develop_docs/README.md +++ b/docs/develop_docs/README.md @@ -15,5 +15,6 @@ formats, evidence, and verification workflows. - [Metrics-observability corpora](metrics-observability-corpora.md) - [Physical handoff cost references](physical-handoff-costs.md), [storage operations](storage-operation-costs.md) - [Replacement explanations](replacement-explanations.md) +- [Physical compile coverage for deployment computation](physical-compile-coverage.md) - [Planner vocabulary migration (#427)](planner-vocabulary-migration.md) diff --git a/docs/develop_docs/physical-compile-coverage.md b/docs/develop_docs/physical-compile-coverage.md new file mode 100644 index 00000000..31c69782 --- /dev/null +++ b/docs/develop_docs/physical-compile-coverage.md @@ -0,0 +1,166 @@ +# Physical compile coverage for deployment computation + +Audience: developers moving computation from ASAPQuery-backend into +`asap_physical_operators::physical_planner`. + +## Contract + +Logical selection decides what to compute. The maintenance lifecycle sets node +timing. `physical_planner::compile` turns a timed `PostAsapDag` into physical +operator DAGs. The backend owns ingestion, panes, storage, stored-state +readout, external exact engines, pricing/selection, and execution scheduling. + +A backend lowering is *covered* when `compile` accepts the corresponding +`PostAsapDag` node and produces operators with the same result. The backend +should then pass the timed DAG and its input contracts to `compile`. It should +not rebuild operator choices from PromQL text or construct operators itself. + +## Inventory + +Surveyed backend: `ASAPQuery-backend` branch `perf/788-startup-search`. +Planner base: `split/462-f-physical-planner` (#475). + +Status values: + +- **Supported**: `compile` or `compile_node` already covers this computation. +- **Partial**: some shapes are covered. The Notes column lists the gap. +- **Missing**: `compile` rejects this computation. +- **Backend**: not computation, or owned by the backend. + +| # | Backend site | Computation | Planner node | Status at #475 | Notes | +|---|---|---|---|---|---| +| 1 | `query_time.rs` `Lower::lower`, `compile_logical` | PromQL AST → `QueryTimeOperator` graph for a native query | `Fallback { QueryExpr }` subtrees plus value payloads | Missing | `compile` lowers `Fallback` only as a raw `Scan` source. | +| 2 | `QueryTimeOperator::Aggregate` (sum/min/max/avg/count) | Grouped value aggregation | `Value::Exact(Aggregate)`; `SummaryAgg{ExactAggregate, Reduce}` over finalized values | Supported | Also `promql_values::compile_aggregate`. | +| 3 | `QueryTimeOperator::Sort`, `Limit` (topk, sort, sort_desc) | Ordering and per-group limits | `Value::Sort`, `Value::Limit` | Supported | | +| 4 | `QueryTimeOperator::Binary`, `QueryPlanNode::Binary` (vector ⊗ scalar) | Arithmetic with a scalar operand | `Binary` whose operand is `Fallback{PromqlScalarBridge(Literal)}` | Missing | Query-time `Binary` accepts only label-map vector schemas. The literal node has no native binding. | +| 5 | `QueryTimeOperator::Binary` (vector ⊗ vector) | One-to-one label matching and arithmetic | `Binary` over grouped value rows | Missing | Only the ingestion-time `aligned_binary` and label-map `vector_binary` exist. | +| 6 | `binary_operator` CheckedDiv / FiniteDiv | Guarded division | `BinaryOperator` checked flags | Partial | Flags are evaluated, but only where rows 4/5 are covered. | +| 7 | `QueryTimeOperator::Binary` comparisons, `bool` | Filter or 0/1 comparison | `Binary{Compare}` | Missing | `Payload::Binary` does not carry `return_bool`. | +| 8 | `QueryTimeOperator::UnaryNegate` | Negation | `Binary{Mul}` by literal `-1` | Missing | The frontend emits `* -1`; same gap as row 4. | +| 9 | `QueryTimeOperator::VectorToScalar` | `scalar()` | `Fallback{PromqlScalarFromVector}` | Missing | Only `promql_values::compile_vector_to_scalar`. | +| 10 | `QueryTimeOperator::HistogramQuantile` | Bucket interpolation | `Fallback` / `AggIntent::HistogramQuantile` | Missing | Only `promql_values::compile_histogram_quantile`. | +| 11 | `QueryTimeOperator::Temporal` (rate, increase, `*_over_time`) | Per-series window functions | `SummaryAgg{PerEntity}` over `TimeRange(Scan)` | Partial | Supported with closed series identity. Not supported over `Fallback` matrices (`compile_temporal` only). | +| 12 | `logical_dag.rs` `Subquery`, `subquery_grid`, `expanded_inputs` | Re-evaluate the child on a step grid and assemble a matrix | `Fallback{PromqlSubquery}` | Missing | No Planner operator. | +| 13 | `QueryPlanNode::Scalar`, `DagCompiler::lower` scalar literal | Scalar constant | `Fallback{PromqlScalarBridge(Literal)}` | Missing | Only `promql_values::compile_scalar`. | +| 14 | `DagCompiler::lower` `ReduceSum`; `physical_values.rs` PerEntity projection | Sum over finalized values; per-entity identity | `SummaryAgg{ExactAggregate(Sum)}` | Supported | The backend builds an identity `Operator::project` itself for PerEntity. | +| 15 | `DagCompiler::lower` `ExactReadout`; `post_asap_readout.rs` ExactReadout | Finalize exact state (sum/count/min/max/rate/increase) | `Value::FinalizeExactAccumulator` | Partial | Count yields Int64 against a declared Float64 PromQL value. `compile` rejects it. | +| 16 | `post_asap_readout.rs` SummaryEstimate (`readout_bound`, `expand_item_rows`) | Sketch estimate per group; TopK item expansion | `SummaryEstimate` | Partial | The backend's label-map state layout and MetricsQL `__name__` rules have no Planner equivalent. `compile_exact_readout` has no sketch counterpart. | +| 17 | `post_asap_readout.rs` SummaryMerge (`merge_bound_states`) | Merge states by group | `SummaryMerge` | Supported | Union plus `summary_merge`. | +| 18 | `post_asap_readout.rs` counter range parameters | Counter lookback for rate/increase | `TimeRange` ancestor of finalization | Supported | Applied through `with_counter_lookback`. | +| 19 | `post_asap_readout.rs` `execute_value_fragment` | Per-timestamp binding of a value fragment | n/a | Backend | Evaluation scheduling. | +| 20 | `DagCompiler::lower` SummaryJoin / Subtract / Delete | Summary algebra | `SummaryJoin`, `SummarySubtract`, `SummaryDelete` | Missing | The backend also rejects these (`ExactFallback`). | +| 21 | `current_series.rs` Snapshot + TopK | Current-series ranking | `ReadPopulation{TopK}` | Supported | | +| 22 | `current_series.rs` Sum / Count / Average | Current-series aggregates | `ReadPopulation{Sum,Count,Average}` | Missing | `compile` accepts only TopK. | +| 23 | `current_series.rs` Quantile | Current-series quantile | `ReadPopulation{Quantile}` | Missing | No exact quantile reduction. | +| 24 | `raw_dag.rs` weight `Column` | Summary update from a sample/projected value | `SummaryAgg` | Supported | | +| 25 | `raw_dag.rs` weight `Constant` | Unit/constant-weight update | `SummaryAgg` | Missing | `compile_node` requires a column weight. | +| 26 | `raw_dag.rs` item `Column` / `Tuple` | Keyed update item | `SummaryAgg{item}` | Supported | `keyed_summary_build`. | +| 27 | `raw_dag.rs` item `EntityIdentity` | Series-identity item | `SummaryAgg{item}` | Missing | Needs the series-identity column. | +| 28 | `physical_values.rs` `compile`, `combine` | Translate `QueryTimeOperator` to `promql_values::*`; compose fragments | n/a | Supported | Exists only because of row 1. `CompiledPhysicalDag::compose` is Planner API. | +| 29 | `query_plan.rs` `compile_native_fragment` (Semi join, Exact aggregate, Sort, Limit, Filter) | Relational value ops | `RelationalJoin`, `Value::*` | Supported | Already calls `compile`. | +| 30 | `query_time.rs` `selected_query_time_nodes`, `selected_native_expression`, `selected_aggregate_operator` | Recover operator identity from original PromQL text | Payload variants (`ExactKind::Min`/`Max`, `AggIntent`) | Supported | Payloads already carry the identity. These witnesses are needed only while row 1 remains. | +| 31 | Scan, ExactSubquery, CandidateExactSubquery, CurrentSeries ingest, ReadMaterialization, ExternalExact | Storage reads and external engines | Input contracts | Backend | | + +Totals at #475: 11 Supported, 4 Partial, 14 Missing, 2 Backend. + +## Covered after this change + +| Row | Change | +|---|---| +| 4, 8, 13 | Query-time `Binary` folds a scalar-literal operand into a projection over grouped value rows. | +| 5 | Query-time `Binary` over grouped value rows performs an inner equi-join on equal label columns, then applies the operator. Per-series rows remain Partial. | +| 15 | Count finalization converts exactly to the declared Float64 value. | +| 22, 23 | `ReadPopulation` Sum/Count/Average/Quantile compile to grouped aggregation. `Reduction::Quantile` implements PromQL interpolation. | + +Totals after this change: 17 Supported, 4 Partial, 8 Missing, 2 Backend. + +## Covered by PromQL fallback compilation + +`compile` now lowers a `Fallback{QueryExpr}` node from its typed expression, +realized with `promql_rows::with_series_identity`. The deployment supplies the +raw rows of the `i`th selector returned by `promql_fallback::raw_series` at +`promql_fallback::raw_series_input(node, i)`, with that selector's schema. Each +selector has its own slot, even when two selectors read the same metric. The +node's own ID still names its output, so a deployment may instead supply the +whole result, for example from an external exact engine. A Fallback that reads +no selector, such as `vector(1)`, needs no input. + +`Operator::series_window` evaluates each series at the query time, or at each +subquery step, over the left-open window `(t - offset - range, t - offset]`. +Instant selection takes the latest sample and omits the series if that sample is +a stale marker. Range functions ignore stale markers. The raw rows must cover +every window the node evaluates; `raw_series` documents the subquery extent. +A subquery has at most 100000 steps. Output rows keep the full series identity; +the query adapter still applies PromQL's metric-name rules. A bare selector +consumed by another node, such as a per-entity summary, remains raw range rows +and is not compiled as instant selection. + +| Row | Change | +|---|---| +| 1 | Supported shapes: selectors with `offset`; `rate`, `increase`, `delta`, and `sum`/`avg`/`min`/`max`/`count_over_time`; `by` aggregation; `sort`; `topk`/`limit`; `scalar()`; `vector(literal)`; arithmetic with one literal. Now Partial. | +| 9 | `scalar()` compiles to `VectorToScalar`. | +| 11 | Range functions over the raw selector rows. | +| 12 | `f(sel[R:S])` and `f(g(sel[r])[R:S])`, with subquery `offset`, evaluate on the aligned step grid. The frontend now retains subquery `offset`/`@` as a `TimeShift`; it previously dropped them. Now Partial. | + +Totals after this change: 19 Supported, 5 Partial, 5 Missing, 2 Backend. + +## Covered by multi-selector fallback compilation + +| Row | Change | +|---|---| +| 1 | Vector-vector arithmetic with one-to-one matching, `on`, and `ignoring`. Each side is reduced to its matching labels, then matched on equal label sets. A duplicate match group is an error, as in Prometheus. The result drops `__name__`. `without` aggregation. `irate`, `idelta`, `changes`, `resets`, `last_over_time`, and exact `quantile_over_time`. `@ ` on selectors. Still Partial. | +| 11 | The range functions above, over raw selector rows. | +| 12 | `@ ` on subqueries anchors the step grid. An inner selector's `@` pins every step. Still Partial. | + +`@` fixes the instant a window ends at, before `offset`; the output keeps the +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. + +Totals are unchanged: 19 Supported, 5 Partial, 5 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()`. + - 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 + run scope. + - 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 + `SummaryAgg`. +5. 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.