From 2416901b8d26ec902822176fade17736ab5d0b39 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:45:46 +0000 Subject: [PATCH 1/4] fix(promql): retain subquery offset and @ as a TimeShift Co-Authored-By: Claude Opus 5.5 --- crates/frontend-promql/src/promql.rs | 22 ++++++++++++++----- .../frontend-promql/tests/promql_lowering.rs | 17 ++++++++++++++ crates/types/src/pre_asap/query_expr.rs | 3 ++- 3 files changed, 36 insertions(+), 6 deletions(-) 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, From 2cd0eecd2731b0e012165e3597e053bb7ac5a586 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:45:46 +0000 Subject: [PATCH 2/4] feat(physical): add a PromQL per-series window operator Evaluates instant selection and range functions over left-open windows at the query time or on a subquery step grid. Conflicts with earlier stack changes resolved to the integration tree: - crates/asap-physical-operators/src/operators/mod.rs: a7ff3ae Merge remote-tracking branch 'origin/feat/physical-compile-promql-fallback' into integration/planner-for-backend Co-Authored-By: Claude Opus 5.5 --- .../src/operators/aggregate/temporal.rs | 79 +++--- .../src/operators/mod.rs | 14 ++ .../src/operators/series_window.rs | 226 ++++++++++++++++++ .../src/operators/unchecked.rs | 13 + 4 files changed, 297 insertions(+), 35 deletions(-) create mode 100644 crates/asap-physical-operators/src/operators/series_window.rs diff --git a/crates/asap-physical-operators/src/operators/aggregate/temporal.rs b/crates/asap-physical-operators/src/operators/aggregate/temporal.rs index dbf66108..c14abcbb 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,47 @@ 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 + } + }))), + _ => 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 +127,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 +138,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..cb7c5433 100644 --- a/crates/asap-physical-operators/src/operators/mod.rs +++ b/crates/asap-physical-operators/src/operators/mod.rs @@ -23,6 +23,7 @@ mod joins; mod limit; mod projection; mod scope_timestamp; +mod series_window; mod sort; mod source; mod summary; @@ -30,6 +31,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 +68,14 @@ enum Kind { intent: Box>, }, HistogramQuantile, + SeriesWindow { + function: Option>>, + coordinate: usize, + value: usize, + range_ms: i64, + offset_ms: i64, + steps: Option, + }, Project(Vec), Filter(Expression), Limit { @@ -237,6 +247,7 @@ impl PhysicalOperator for Operator { | Kind::RangeWindow { .. } | Kind::HistogramQuantile | Kind::CurrentSeries { .. } + | Kind::SeriesWindow { .. } | Kind::Aggregate { .. } | Kind::Window { .. } | Kind::Join { .. } @@ -279,6 +290,7 @@ impl PhysicalOperator for Operator { Kind::AlignedBinary { .. } => "AlignedBinary", Kind::RangeWindow { .. } => "RangeWindow", Kind::HistogramQuantile => "HistogramQuantile", + Kind::SeriesWindow { .. } => "SeriesWindow", Kind::Project(_) => "Project", Kind::Filter(_) => "Filter", Kind::Limit { .. } => "Limit", @@ -295,6 +307,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 +336,7 @@ 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::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_window.rs b/crates/asap-physical-operators/src/operators/series_window.rs new file mode 100644 index 00000000..ad56db0b --- /dev/null +++ b/crates/asap-physical-operators/src/operators/series_window.rs @@ -0,0 +1,226 @@ +//! 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]`, where `T` is 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, +} + +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 + /// `(t - offset_ms - range_ms, t - offset_ms]`. `t` is the query time, or + /// each step of `steps`. `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, + 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 } + ) + ) { + return Err(invalid("unsupported PromQL range function")); + } + Ok(Self { + kind: Kind::SeriesWindow { + function: function.map(Box::new), + coordinate, + value: *value, + range_ms, + offset_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 = 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, + 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 = 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..80a4544a 100644 --- a/crates/asap-physical-operators/src/operators/unchecked.rs +++ b/crates/asap-physical-operators/src/operators/unchecked.rs @@ -50,6 +50,19 @@ 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, + steps, + .. + } => Operator::series_window( + input(0)?, + function.map(|f| *f), + range_ms, + offset_ms, + steps, + )?, Kind::Project(expressions) => { if expressions.len() != output.fields.len() { return Err(invalid("projection width mismatch")); From 64a9cb63d2d5a5863c409cdcad45fddb4c3df457 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:45:46 +0000 Subject: [PATCH 3/4] feat(physical): compile PromQL fallback subtrees from typed expressions Conflicts with earlier stack changes resolved to the integration tree: - crates/asap-physical-operators/src/physical_planner/promql_rows.rs: a7ff3ae Merge remote-tracking branch 'origin/feat/physical-compile-promql-fallback' into integration/planner-for-backend Co-Authored-By: Claude Opus 5.5 --- .../src/physical_planner/mod.rs | 55 ++- .../src/physical_planner/promql_fallback.rs | 403 ++++++++++++++++++ .../tests/physical_dag.rs | 4 +- .../tests/promql_fallback.rs | 385 +++++++++++++++++ 4 files changed, 845 insertions(+), 2 deletions(-) create mode 100644 crates/asap-physical-operators/src/physical_planner/promql_fallback.rs create mode 100644 crates/asap-physical-operators/tests/promql_fallback.rs diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index c7967288..21520150 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -32,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; @@ -173,7 +174,19 @@ fn compile_internal( .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(); @@ -228,6 +241,46 @@ 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 slot = promql_fallback::raw_series_input(id); + let (leaf, mut chain) = promql_fallback::lower(expression) + .map_err(|error| invalid(format!("node {id}: {error}")))?; + let mut inputs = match (leaf, sources.remove(&slot)) { + (Some((_, schema)), Some(contract)) if contract.schema == schema => { + graph.add_input(slot, contract)?; + vec![slot] + } + (None, None) => vec![], + (Some(_), None) => { + return Err(invalid(format!( + "node {id}: PromQL fallback requires raw series input {slot}" + ))) + } + _ => { + return Err(invalid(format!( + "node {id}: raw series input differs from the selector schema" + ))) + } + }; + let last = chain + .pop() + .ok_or_else(|| invalid("empty PromQL lowering"))?; + for operator in chain { + graph.add(auxiliary, inputs, operator)?; + inputs = vec![auxiliary]; + auxiliary -= 1; + } + graph.add(id, inputs, last.with_output_schema(output)?)?; + continue; + } if let Payload::Value { operation: ValueOperation::MaintainPopulation { population }, } = &node.payload 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..78dd1b06 --- /dev/null +++ b/crates/asap-physical-operators/src/physical_planner/promql_fallback.rs @@ -0,0 +1,403 @@ +//! Compile a retained PromQL subtree (`Fallback`) from its typed expression. +//! The deployment supplies the raw series of its one selector; the Planner +//! computes selection, range functions, subqueries and aggregation. +use super::*; +use crate::operators::SubquerySteps; +use planner_types::post_asap::execution_data_state::lift_plain; + +/// Input slot for the raw series read by Fallback node `node`'s selector. +/// The node's own ID names its computed output, so the raw rows need another. +pub fn raw_series_input(node: NodeId) -> NodeId { + node | (1 << 32) +} + +/// The Fallback node that owns a raw-series input slot. +pub(super) fn raw_series_owner(slot: NodeId) -> Option { + (slot >> 32 == 1).then_some(slot & u64::from(u32::MAX)) +} + +/// A selector expression and its raw-series row schema. +pub type Selector = (QueryExpr, Schema); + +/// The selector a Fallback expression reads, and the row schema of the raw +/// series the deployment supplies at [`raw_series_input`]. `None` means the +/// expression reads no series. The rows must cover the selector's window at +/// every evaluation instant; 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)?.0) +} + +/// Operators computing `expression`, in order, after its raw-series input. +pub(super) fn lower(expression: &QueryExpr) -> Result<(Option, Vec), Error> { + let mut chain = Chain::default(); + chain.value(expression)?; + Ok((chain.leaf, chain.operators)) +} + +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")) +} + +/// `TimeRange { range, [TimeShift { offset }], Scan }`: range and offset. +fn selector(expression: &QueryExpr) -> Result<(i64, i64), Error> { + let QueryExpr::TimeRange { range, child } = expression else { + return Err(invalid("PromQL operand must be a series selector")); + }; + let (offset, scan) = match child.as_ref() { + QueryExpr::TimeShift { shift, child } if shift.at.is_none() => { + (shift.offset_ms, child.as_ref()) + } + scan => (0, scan), + }; + if !matches!(scan, QueryExpr::Scan { .. }) { + return Err(invalid( + "PromQL selector must read one scan; @ is unsupported", + )); + } + Ok((millis(range)?, offset)) +} + +#[derive(Default)] +struct Chain { + leaf: Option, + operators: Vec, +} + +impl Chain { + fn schema(&self) -> Result { + self.operators + .last() + .map(Operator::schema) + .or_else(|| self.leaf.as_ref().map(|(_, schema)| schema.clone())) + .ok_or_else(|| invalid("PromQL operator has no input")) + } + + /// Conform `operator` to the logical schema of the expression it computes. + fn push(&mut self, operator: Operator, logical: &QueryExpr) -> Result<(), Error> { + self.operators + .push(operator.with_output_schema(declared(logical)?)?); + Ok(()) + } + + 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", + )); + } + if self + .leaf + .replace((selector.clone(), schema.clone())) + .is_some() + { + return Err(invalid("PromQL fallback reads more than one selector")); + } + Ok(schema) + } + + /// An instant vector, or a scalar for scalar-valued expressions. + fn value(&mut self, expression: &QueryExpr) -> Result<(), Error> { + match expression { + QueryExpr::TimeRange { .. } => { + let (range, offset) = selector(expression)?; + let input = self.read(expression)?; + self.push( + Operator::series_window(input, None, range, offset, None)?, + 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")); + }; + self.value(child)?; + self.aggregate(measure, keys, expression) + } + QueryExpr::Sort { + keys, + partition_by, + child, + } => { + self.value(child)?; + let input = self.schema()?; + 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)?, expression) + } + QueryExpr::Limit { n, offset, child } => { + self.value(child)?; + let input = self.schema()?; + // `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)?, + expression, + ) + } + QueryExpr::BinaryOp { + op: planner_types::pre_asap::BinaryOpKind::Arithmetic(op), + lhs, + rhs, + vector_match: None, + } => { + 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), + _ => { + return Err(invalid( + "PromQL fallback arithmetic requires one literal operand", + )) + } + }; + self.value(vector)?; + let input = self.schema()?; + 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)?, expression) + } + QueryExpr::PromqlScalarFromVector(child) => { + self.value(child)?; + let input = self.schema()?; + let value = named_column(&input, &ColumnRef::SampleValue)?; + self.push(Operator::vector_to_scalar(input, value)?, expression) + } + QueryExpr::PromqlVectorFromScalar(child) => { + self.value(child)?; + let input = self.schema()?; + self.operators + .push(Operator::scope_timestamp(input, declared(expression)?)?); + Ok(()) + } + 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)?, + expression, + ) + } + _ => Err(invalid("PromQL expression has no native fallback lowering")), + } + } + + /// `function(matrix)`, where the matrix is a range selector or a subquery. + fn range_function( + &mut self, + function: &AggIntent, + matrix: &QueryExpr, + logical: &QueryExpr, + ) -> Result<(), Error> { + let function = unbound(function)?; + let (subquery, offset) = match matrix { + QueryExpr::TimeShift { shift, child } if shift.at.is_none() => { + (child.as_ref(), shift.offset_ms) + } + other => (other, 0), + }; + let QueryExpr::PromqlSubquery { + range: outer, + resolution, + child, + } = subquery + else { + let (range, offset) = selector(matrix)?; + let input = self.read(matrix)?; + return self.push( + Operator::series_window(input, Some(function), range, offset, None)?, + 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, + }; + // 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) = selector(selected)?; + let input = self.read(selected)?; + self.push( + Operator::series_window(input, inner, range, inner_offset, Some(steps))?, + child, + )?; + let input = self.schema()?; + self.push( + Operator::series_window(input, Some(function), steps.range_ms, offset, None)?, + 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, + measure: &AggIntent, + keys: &GroupKeys, + logical: &QueryExpr, + ) -> Result<(), Error> { + let mut input = self.schema()?; + 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 = 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]; + self.operators.push(project); + } + 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(); + self.operators.push(aggregate); + // 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)?, 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 }, + _ => return Err(invalid("unsupported PromQL range function")), + }) +} 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..c31a57d4 --- /dev/null +++ b/crates/asap-physical-operators/tests/promql_fallback.rs @@ -0,0 +1,385 @@ +//! 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 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() + }), + }; + let expression = asap_frontend_promql::lower_promql_workload(&workload, 0) + .unwrap() + .remove(0); + promql_rows::with_series_identity(&expression).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() +} + +/// `(job, seconds, value)`; every sample belongs to metric `m`. +type Sample = (&'static str, i64, f64); + +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())? + .map(|(_, schema)| { + ( + promql_fallback::raw_series_input(root), + InputContract::bounded(schema), + ) + }) + .into_iter() + .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; returns `(job or "", timestamp ms, value)` rows in order. +fn run(query: &str, samples: &[Sample], at: i64) -> Result, String> { + let expression = lower(query); + let program = compile_query(query)?; + let mut sources = BTreeMap::new(); + if let Some((_, schema)) = promql_fallback::raw_series(&expression).unwrap() { + let rows = samples + .iter() + .map(|(job, seconds, value)| { + let labels = BTreeMap::from([ + ("__name__".to_string(), "m".to_string()), + ("job".into(), job.to_string()), + ]); + 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]), + 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 job = String::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)) => { + job = promql_rows::decode_series_identity(id).unwrap()["job"].clone() + } + ("job", Value::Utf8(label)) => job = label.to_string(), + (_, 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:?}")), + } + } + rows.push((job, time, value.ok_or("missing value")?)); + } + } + Ok(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. + ); + assert!(compile_query("max_over_time(m[5m:1m] @ 100)").is_err()); +} + +// 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().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), + 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().unwrap().1; + assert!(compile( + &consumed, + BTreeMap::from([( + promql_fallback::raw_series_input(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()); +} From fd0bb0a1bbb163800df6edfe949d8dfcd1fd5b93 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:45:46 +0000 Subject: [PATCH 4/4] docs: record PromQL fallback compile coverage Co-Authored-By: Claude Opus 5.5 --- .../develop_docs/physical-compile-coverage.md | 42 +++++++++++++++++-- 1 file changed, 39 insertions(+), 3 deletions(-) diff --git a/docs/develop_docs/physical-compile-coverage.md b/docs/develop_docs/physical-compile-coverage.md index 5db0edce..ce8c3817 100644 --- a/docs/develop_docs/physical-compile-coverage.md +++ b/docs/develop_docs/physical-compile-coverage.md @@ -74,13 +74,49 @@ Totals at #475: 11 Supported, 4 Partial, 14 Missing, 2 Backend. 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. +The expression must read at most one selector, realized with +`promql_rows::with_series_identity`. The deployment supplies that selector's raw +rows at `promql_fallback::raw_series_input(node)`, with the schema returned by +`promql_fallback::raw_series`. 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. + ## Remaining In order of backend usage: -1. Row 1 and rows 9–12: lower PromQL-shaped `Fallback{QueryExpr}` subtrees - (range functions over matrices, `scalar()`, `histogram_quantile`, `sort`, - subquery grids). After that, rows 28 and 30 can be deleted from the backend. +1. Rows 1, 10, and 12, the remaining `Fallback` shapes: + - Subtrees with more than one selector, such as vector-vector binaries. + - `histogram_quantile`: the frontend emits `by ()` grouping with no + output labels. The IR must group `without (le)` and keep the labels. + - Subquery operands other than one per-series function; implicit + subquery resolution, which is a deployment default. + - `@` on selectors and subqueries; `without` grouping; `irate`, + `changes`, and other range 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: matching needs a metric-name-free series