Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions crates/asap-aware-mapping/src/analytical_cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2193,6 +2193,8 @@ pub enum AnalyticalCostError {
UnsupportedCandidate,
#[error("query operator has no physical implementation in the analytical model")]
UnsupportedQueryOperator,
#[error("multi-measure per-entity aggregates have no physical implementation; lower each measure separately")]
UnsupportedMultiMeasurePerEntity,
#[error("inconsistent operator statistics: {0}")]
InconsistentOperatorStatistics(&'static str),
#[error("summary operation {0} has no lifecycle-aware cost formula")]
Expand Down
39 changes: 38 additions & 1 deletion crates/asap-aware-mapping/src/query_physical_lowering.rs
Original file line number Diff line number Diff line change
Expand Up @@ -296,7 +296,7 @@ pub fn lower_query_physical_dag(
}
if matches!(reduction, asap_types::pre_asap::Reduction::PerEntity) {
if measures.len() != 1 {
return Err(AnalyticalCostError::UnsupportedQueryOperator);
return Err(AnalyticalCostError::UnsupportedMultiMeasurePerEntity);
}
let accumulator_count = u64::try_from(measures.len())
.map_err(|_| AnalyticalCostError::Overflow)?;
Expand Down Expand Up @@ -2545,6 +2545,43 @@ mod tests {
));
}

/// Multi-measure schemas must not silently lower to a single accumulator.
#[test]
fn multi_measure_per_entity_lowering_is_explicitly_unsupported() {
use asap_types::pre_asap::{
AggIntent, Column, DataType, QueryExpr, Reduction, Schema, Source,
};
let source = Source::TimeSeries {
metric: "requests".into(),
};
let root = Rc::new(QueryExpr::Aggregate {
reduction: Reduction::PerEntity,
measures: vec![
AggIntent::Sum { col: None },
AggIntent::Count {
accuracy: asap_types::types::AccuracyTarget::Exact,
},
],
output_names: vec![],
filters: vec![],
having: None,
child: Rc::new(QueryExpr::Scan {
source: source.clone(),
predicates: vec![],
schema: Schema::new(vec![Column::new("value", DataType::Float64, false)]),
}),
});
let provided = HashMap::new();
assert!(matches!(
lower_query_physical_dag(
&root,
&scope(vec![coverage(source, vec![])]),
&scripted(&provided)
),
Err(AnalyticalCostError::UnsupportedMultiMeasurePerEntity)
));
}

#[test]
fn promql_relabel_sample_and_per_series_lower_as_a_complete_chain() {
use asap_types::pre_asap::{
Expand Down
29 changes: 27 additions & 2 deletions crates/asap-aware-mapping/src/rewrite.rs
Original file line number Diff line number Diff line change
Expand Up @@ -545,19 +545,44 @@ mod tests {
assert!(AvgToSumOverCountStrategy.replacements(&target).is_empty());
}

/// Matching schemas do not make range SUM/COUNT safe for unbounded samples.
#[test]
fn does_not_match_a_per_entity_avg_aggregate() {
fn per_entity_avg_rewrite_is_rejected_without_arithmetic_proof() {
let q = Rc::new(QueryExpr::Aggregate {
reduction: Reduction::PerEntity,
measures: vec![AggIntent::Avg { col: None }],
output_names: vec![],
filters: vec![],
having: None,
child: Rc::new(metric_scan(&[])),
child: Rc::new(QueryExpr::TimeRange {
range: Duration::from_secs(300),
child: Rc::new(metric_scan(&["job", "instance"])),
}),
});
let target = TargetSubDAG::new(&q);
assert!(!AvgToSumOverCountStrategy.matches(&target));
assert!(AvgToSumOverCountStrategy.replacements(&target).is_empty());
assert!(build_rewrite(&q).is_none());
}

/// COUNT(*) cannot replace the denominator of a nullable sample average.
#[test]
fn per_entity_nullable_or_non_sample_average_is_not_rewritten() {
for (nullable, column) in [(true, None), (false, Some(2)), (false, Some(99))] {
let mut scan = metric_scan(&["job"]);
if let QueryExpr::Scan { schema, .. } = &mut scan {
schema.columns[1].nullable = nullable;
}
let root = Rc::new(QueryExpr::Aggregate {
reduction: Reduction::PerEntity,
measures: vec![AggIntent::Avg { col: column }],
output_names: vec![],
filters: vec![],
having: None,
child: Rc::new(scan),
});
assert!(build_rewrite(&root).is_none());
}
}

#[test]
Expand Down
28 changes: 28 additions & 0 deletions crates/integration-tests/tests/avg_over_time_rewrite.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
use std::rc::Rc;

use asap_aware_mapping::replacement::{ReplacementStrategy, TargetSubDAG};
use asap_aware_mapping::rewrite::AvgToSumOverCountStrategy;
use asap_integration_tests::fixtures::lower_promql;
use asap_types::types::AccuracyTarget;

/// Unbounded Float64 ranges must retain AVG: finite samples can overflow SUM.
/// The independent Prometheus oracle is fixtures/avg_over_time_overflow.test.yml.
#[test]
fn range_average_does_not_offer_unconditional_sum_count() {
for query in [
"avg_over_time(latency[5m])",
"avg_over_time(latency{job=\"api\"}[5m])",
"avg_over_time(latency[5m:1m])",
] {
let root = Rc::new(lower_promql(query, AccuracyTarget::Exact).unwrap());
let target = TargetSubDAG::new(&root);
assert!(
!AvgToSumOverCountStrategy.matches(&target),
"unbounded average must not match: {query}"
);
assert!(
AvgToSumOverCountStrategy.replacements(&target).is_empty(),
"unbounded average must not produce a sum/count rewrite: {query}"
);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
# Official Prometheus 3.5.0 oracle: AVG remains finite when SUM overflows.
# Run: promtool test rules crates/integration-tests/tests/fixtures/avg_over_time_overflow.test.yml
evaluation_interval: 1m
tests:
- interval: 1m
input_series:
- series: 'latency{job="api"}'
values: '1e308 1e308'
promql_expr_test:
- expr: 'avg_over_time(latency[5m])'
eval_time: 1m
exp_samples:
- labels: '{job="api"}'
value: 1e308
- expr: 'sum_over_time(latency[5m]) / count_over_time(latency[5m])'
eval_time: 1m
exp_samples:
- labels: '{job="api"}'
value: .inf
- interval: 1m
input_series:
- series: 'latency{job="api"}'
values: '-1e308 -1e308'
promql_expr_test:
- expr: 'avg_over_time(latency[5m])'
eval_time: 1m
exp_samples:
- labels: '{job="api"}'
value: -1e308
- expr: 'sum_over_time(latency[5m]) / count_over_time(latency[5m])'
eval_time: 1m
exp_samples:
- labels: '{job="api"}'
value: -.inf
143 changes: 135 additions & 8 deletions crates/types/src/pre_asap/query_expr.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,8 @@ impl ColState for ColumnRef {
/// Errors from schema derivation over a canonical tree.
#[derive(Debug, Error)]
pub enum QueryExprError {
#[error("invalid per-entity aggregate: {0}")]
InvalidPerEntityAggregate(String),
#[error("invalid scalar function signature: {0}")]
InvalidScalarSignature(String),
#[error("by-column id {0} out of range (input has {1} columns)")]
Expand Down Expand Up @@ -1540,8 +1542,14 @@ fn per_series_reduction_schema(input: &Schema, agg: &AggIntent) -> Result<Schema
///
/// `Reduction::PerEntity` selects the label-preserving
/// [`per_series_reduction_schema`] (`rate`/`increase`/`*_over_time`) instead
/// of the cross-series `by ++ measures` shape. Which one applies is read directly
/// off `reduction` — decided once, at construction, by whoever built the
/// of the cross-series `by ++ measures` shape. Multiple measures replace the
/// value slot with the first measure and append the rest, preserving label and
/// timestamp positions. Names come from output overrides or the intents and
/// must be distinct from one another and retained columns; all samples are Float64.
/// Frontends resolving named multi-measure outputs should supply a scan schema:
/// usage-derived binding may otherwise infer those names as input labels, which
/// this validation rejects as collisions.
/// Which shape applies is read directly off `reduction` — decided once, at construction, by whoever built the
/// `Aggregate` node (issue #165) — not re-derived here from `by`/child shape.
pub fn aggregate_output_schema(
in_schema: &Schema,
Expand All @@ -1551,12 +1559,57 @@ pub fn aggregate_output_schema(
) -> Result<Schema, QueryExprError> {
let by = match reduction {
Reduction::PerEntity => {
debug_assert_eq!(
measures.len(),
1,
"a per-entity reduction is single-aggregate"
);
return per_series_reduction_schema(in_schema, &measures[0]);
if measures.is_empty() {
return Err(QueryExprError::InvalidPerEntityAggregate(
"at least one measure is required".into(),
));
}
if measures.len() == 1 {
return per_series_reduction_schema(in_schema, &measures[0]);
}
let invalid = |message: &str| QueryExprError::InvalidPerEntityAggregate(message.into());
let vi = in_schema
.column_id("value")
.ok_or_else(|| invalid("multiple measures require an input value column"))?;
if in_schema.time_index == Some(vi) {
return Err(invalid("the value column cannot be the timestamp"));
}
if output_names.len() > measures.len() {
return Err(invalid("more output names than measures"));
}
let mut output = in_schema.clone();
let mut names: std::collections::HashSet<String> = in_schema
.columns
.iter()
.enumerate()
.filter(|(i, _)| *i != vi)
.map(|(_, c)| c.name.clone())
.collect();
for (i, measure) in measures.iter().enumerate() {
if matches!(measure, AggIntent::CountValues { .. }) {
return Err(invalid("count_values changes series identity and cannot be combined with other measures"));
}
let input = in_schema
.columns
.get(measure.input_cols().first().copied().unwrap_or(vi))
.ok_or_else(|| invalid("measure input column is out of range"))?;
let mut column = measure.output_column(input);
if let Some(name) = output_names.get(i).filter(|n| !n.is_empty()) {
column.name = name.clone();
}
if !names.insert(column.name.clone()) {
return Err(invalid("measure names must be unique and must not collide with labels or timestamp"));
}
column.dtype = DataType::Float64;
if i == 0 {
output.columns[vi] = column;
} else {
output.columns.push(column);
}
}
// A computed sample cannot retain a uniqueness proof about raw values.
output.unique_keys.retain(|key| !key.contains(&vi));
return Ok(output);
}
Reduction::Reduce(by) => by,
};
Expand Down Expand Up @@ -2393,6 +2446,80 @@ mod tests {
}
}

/// Multiple measures retain series metadata and expose separate named float samples.
#[test]
fn per_entity_multiple_measures_preserve_schema() {
let input = Schema::with_time_index(
vec![
col("ts", DataType::Timestamp, false),
col("value", DataType::Float64, false),
col("job", DataType::Utf8, true),
],
0,
vec![vec![0, 2]],
);
let measures = vec![
AggIntent::Sum { col: None },
AggIntent::Count {
accuracy: crate::types::AccuracyTarget::Exact,
},
];
let output =
aggregate_output_schema(&input, &Reduction::PerEntity, &measures, &[]).unwrap();
assert_eq!(
output
.columns
.iter()
.map(|c| c.name.as_str())
.collect::<Vec<_>>(),
vec!["ts", "sum", "job", "count"]
);
assert_eq!(output.time_index, input.time_index);
assert_eq!(output.unique_keys, input.unique_keys);
assert_eq!(output.closed, input.closed);
assert_eq!(output.columns[1].dtype, DataType::Float64);
assert_eq!(output.columns[3].dtype, DataType::Float64);
}

/// Invalid or ambiguous multi-measure shapes fail instead of losing outputs.
#[test]
fn per_entity_measure_names_are_validated() {
let input = Schema::new(vec![
col("value", DataType::Float64, false),
col("job", DataType::Utf8, true),
]);
let measures = vec![AggIntent::Sum { col: None }, AggIntent::Sum { col: None }];
for names in [
vec![],
vec!["a".into(), "a".into()],
vec!["job".into(), "b".into()],
] {
assert!(
aggregate_output_schema(&input, &Reduction::PerEntity, &measures, &names).is_err()
);
}
let output = aggregate_output_schema(
&input,
&Reduction::PerEntity,
&measures,
&["first".into(), "second".into()],
)
.unwrap();
assert_eq!(output.column_id("first"), Some(0));
assert_eq!(output.column_id("second"), Some(2));
assert!(aggregate_output_schema(&input, &Reduction::PerEntity, &[], &[]).is_err());
assert!(aggregate_output_schema(
&input,
&Reduction::PerEntity,
&[
AggIntent::Sum { col: Some(99) },
AggIntent::Sum { col: None }
],
&[]
)
.is_err());
}

#[test]
fn per_series_rate_preserves_labels() {
// A per-series range reduction (`rate`) is label-preserving: it produces
Expand Down
Loading
Loading