Skip to content

Commit c9b2aa7

Browse files
fix(promql): preserve ingestion horizons through lowering and post-ASAP IR (#417)
Co-authored-by: zz_y <zeyingz@umd.edu>
1 parent b7cfee7 commit c9b2aa7

51 files changed

Lines changed: 1102 additions & 293 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎crates/asap-aware-mapping/src/lib.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -206,6 +206,8 @@ pub mod storage_io;
206206
pub mod summary_maintenance_cost;
207207
pub mod summary_maintenance_dag_export;
208208
pub mod summary_maintenance_lifecycle;
209+
#[cfg(test)]
210+
mod test_support;
209211
pub mod topk_reuse;
210212

211213
pub use accuracy::{

‎crates/asap-aware-mapping/src/maintained_population.rs‎

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -118,11 +118,24 @@ fn recognize(root: &QueryExpr) -> Option<(MaintainedPopulation, PopulationReadou
118118
}
119119
return Some((population, readout, Rc::clone(source)));
120120
}
121+
// A bare PromQL selector carries the declared ingestion interval as a
122+
// temporal input scope. Membership must expire at that horizon; retain
123+
// the wrapper as the maintained input so validation can check agreement.
124+
let (series_source, lookback_ms) = match source.as_ref() {
125+
QueryExpr::TimeRange { range, child } => {
126+
let ms = u64::try_from(range.as_millis()).ok()?;
127+
if ms == 0 || std::time::Duration::from_millis(ms) != *range {
128+
return None;
129+
}
130+
(child.as_ref(), ms)
131+
}
132+
other => (other, 300_000),
133+
};
121134
let QueryExpr::Scan {
122135
source: Source::TimeSeries { metric },
123136
predicates,
124137
schema,
125-
} = source.as_ref()
138+
} = series_source
126139
else {
127140
return None;
128141
};
@@ -176,7 +189,7 @@ fn recognize(root: &QueryExpr) -> Option<(MaintainedPopulation, PopulationReadou
176189
matchers,
177190
grouping: labels,
178191
without: grouping.is_without(),
179-
lookback_ms: 300_000,
192+
lookback_ms,
180193
}),
181194
max_k: 0,
182195
quantiles: false,
@@ -285,12 +298,11 @@ impl ReplacementStrategy for MaintainedPopulationStrategy {
285298
#[cfg(test)]
286299
mod tests {
287300
use super::*;
301+
use crate::test_support::lower_promql;
288302
use asap_types::post_asap::{compile_executable_dag, share_common_summary_subtrees};
303+
289304
fn lower(q: &str) -> Rc<QueryExpr> {
290-
Rc::new(
291-
asap_frontend_promql::lower_promql(q, asap_types::types::AccuracyTarget::Exact)
292-
.unwrap(),
293-
)
305+
Rc::new(lower_promql(q, asap_types::types::AccuracyTarget::Exact))
294306
}
295307

296308
// Instant scalar aggregations share the same retractable series population.

‎crates/asap-aware-mapping/src/replacement.rs‎

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -5779,6 +5779,7 @@ mod tests {
57795779
use super::*;
57805780
use crate::accuracy::PropagationStats;
57815781
use crate::cost_model::Cost;
5782+
use crate::test_support::lower_promql;
57825783
use asap_types::pre_asap::agg_intent::{
57835784
agg_is_exact, default_cardinality, default_quantile, MathFunc, TimeFunc,
57845785
};
@@ -5798,10 +5799,7 @@ mod tests {
57985799
// Finite samples can overflow a sum although their native average is finite.
57995800
#[test]
58005801
fn temporal_average_requires_finite_division_guard() {
5801-
let root = Rc::new(
5802-
asap_frontend_promql::lower_promql("avg_over_time(a[5m])", AccuracyTarget::Exact)
5803-
.unwrap(),
5804-
);
5802+
let root = Rc::new(lower_promql("avg_over_time(a[5m])", AccuracyTarget::Exact));
58055803
let candidates =
58065804
SketchAlgorithmStrategy::default_cost_model().replacements(&TargetSubDAG::new(&root));
58075805
let operator = candidates
@@ -5830,8 +5828,7 @@ mod tests {
58305828
"topk(5, sum_over_time(a[5m]))",
58315829
"topk by(job)(5, count_over_time(a[5m]))",
58325830
] {
5833-
let root =
5834-
Rc::new(asap_frontend_promql::lower_promql(query, AccuracyTarget::Exact).unwrap());
5831+
let root = Rc::new(lower_promql(query, AccuracyTarget::Exact));
58355832
let models = Models::with_default_accuracy(&crate::cost_model::DefaultCostModel);
58365833
let node = exact_topk_over_temporal_values(&root, models)
58375834
.unwrap()
@@ -5859,7 +5856,7 @@ mod tests {
58595856
"quantile_over_time(0.5,a[5m]) / quantile_over_time(0.9,a[5m])",
58605857
"avg_over_time(a[5m]) / quantile_over_time(0.5,a[5m])",
58615858
] {
5862-
let root = Rc::new(asap_frontend_promql::lower_promql(query, target.clone()).unwrap());
5859+
let root = Rc::new(lower_promql(query, target.clone()));
58635860
let models = Models::with_default_accuracy(&crate::cost_model::DefaultCostModel);
58645861
assert!(realize_binary(&root, models, Some(&target))
58655862
.unwrap()

‎crates/asap-aware-mapping/src/rewrite.rs‎

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -389,8 +389,10 @@ impl ReplacementStrategy for SemanticEquivalentRewriteStrategy {
389389
#[cfg(test)]
390390
mod tests {
391391
use super::*;
392+
use crate::test_support::lower_promql;
392393
use asap_types::pre_asap::query_expr::Source;
393394
use asap_types::pre_asap::schema::{Column, Schema};
395+
use asap_types::types::AccuracyTarget;
394396
use std::time::Duration;
395397

396398
fn metric_scan(labels: &[&str]) -> QueryExpr {
@@ -419,13 +421,10 @@ mod tests {
419421
// Temporal averages expose two single-measure children without closing labels.
420422
#[test]
421423
fn temporal_average_components_preserves_schema_and_exposes_sum_count() {
422-
let root = Rc::new(
423-
asap_frontend_promql::lower_promql(
424-
"avg_over_time(a{job=\"api\"}[5m])",
425-
AccuracyTarget::Exact,
426-
)
427-
.unwrap(),
428-
);
424+
let root = Rc::new(lower_promql(
425+
"avg_over_time(a{job=\"api\"}[5m])",
426+
AccuracyTarget::Exact,
427+
));
429428
assert!(SemanticEquivalentRewriteStrategy
430429
.replacements(&TargetSubDAG::new(&root))
431430
.is_empty());

‎crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3083,7 +3083,10 @@ mod tests {
30833083
fn streaming_data_workload() -> DataWorkload {
30843084
DataWorkload {
30853085
arrival: DataArrival::ContinuouslyIngesting,
3086-
3086+
data_ingestion_interval: Evidence {
3087+
value: Some(asap_types::workload::DurationMs(1_000)),
3088+
..Default::default()
3089+
},
30873090
ingestion_rate: Evidence {
30883091
value: Some(Rate(2.0)),
30893092
source: EvidenceSource::Declared,
@@ -3491,7 +3494,10 @@ mod tests {
34913494
ComparisonScope::from_workload(
34923495
&DataWorkload {
34933496
arrival: DataArrival::ContinuouslyIngesting,
3494-
3497+
data_ingestion_interval: Evidence {
3498+
value: Some(asap_types::workload::DurationMs(1_000)),
3499+
..Default::default()
3500+
},
34953501
ingestion_rate: Evidence {
34963502
value: Some(Rate(2.0)),
34973503
source: EvidenceSource::Declared,

‎crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1690,15 +1690,21 @@ mod tests {
16901690
fn at_rest() -> DataWorkload {
16911691
DataWorkload {
16921692
arrival: DataArrival::AtRest,
1693-
1693+
data_ingestion_interval: Evidence {
1694+
value: Some(DurationMs(1_000)),
1695+
..Default::default()
1696+
},
16941697
..Default::default()
16951698
}
16961699
}
16971700

16981701
fn continuous(observed_at_ms: u64, valid_for_ms: u64) -> DataWorkload {
16991702
DataWorkload {
17001703
arrival: DataArrival::ContinuouslyIngesting,
1701-
1704+
data_ingestion_interval: Evidence {
1705+
value: Some(DurationMs(1_000)),
1706+
..Default::default()
1707+
},
17021708
ingestion_rate: Evidence {
17031709
value: Some(Rate(1.0)),
17041710
source: EvidenceSource::Observed,
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
use asap_types::pre_asap::QueryExpr;
2+
use asap_types::types::AccuracyTarget;
3+
use asap_types::workload::{
4+
AccuracyRequirement, BatchEntry, DataWorkload, DurationMs, Evidence, PlanningWorkload,
5+
Predictability, Query, QueryLanguage, QueryRequirements, QueryWorkload, TimeSelection,
6+
};
7+
8+
pub(crate) fn lower_promql(query: &str, accuracy: AccuracyTarget) -> QueryExpr {
9+
let workload = PlanningWorkload {
10+
query_workload: QueryWorkload {
11+
language: QueryLanguage::PromQL,
12+
query_batch: Some(vec![BatchEntry {
13+
query: Query(query.into()),
14+
requirements: QueryRequirements {
15+
accuracy: AccuracyRequirement::Explicit(accuracy),
16+
..Default::default()
17+
},
18+
predictability: Predictability::Unknown,
19+
invocations: 1,
20+
execute_at: None,
21+
time_selection: TimeSelection::default(),
22+
}]),
23+
repeating_queries: None,
24+
},
25+
data_workload: Some(DataWorkload {
26+
data_ingestion_interval: Evidence {
27+
value: Some(DurationMs(1_000)),
28+
..Default::default()
29+
},
30+
..Default::default()
31+
}),
32+
};
33+
asap_frontend_promql::lower_promql_workload(&workload, 0)
34+
.unwrap()
35+
.pop()
36+
.unwrap()
37+
}

‎crates/devtools/examples/canonical_examples.rs‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
// One-off: pretty-print the QueryExpr for one canonical query per variant,
44
// plus custom Join/SetOp/Dedup/CTE probes, to eyeball the actual shape.
55

6-
use asap_devtools::lower_promql;
6+
use asap_devtools::lower_promql_with_data_ingestion_interval;
77
use asap_frontend_sql::{lower_sql_dialect, SqlCatalog};
88
use asap_types::pre_asap::schema::{Column, DataType, Schema};
99
use asap_types::types::AccuracyTarget;
@@ -70,7 +70,7 @@ async fn main() {
7070
];
7171
for (label, q) in promql_examples {
7272
println!("=== {label} === promql> {q}");
73-
match lower_promql(q, AccuracyTarget::Exact) {
73+
match lower_promql_with_data_ingestion_interval(q, AccuracyTarget::Exact, 1_000) {
7474
Ok(qe) => println!("{qe:#?}"),
7575
Err(e) => println!("ERR: {e}"),
7676
}

‎crates/devtools/examples/topk_ir.rs‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
// Lowers every topk-shaped query from the design discussion and prints the
44
// resulting pre-ASAP IR. Used for interactive exploration; not a test.
55

6-
use asap_devtools::{lower_promql, lower_sql, SqlCatalog};
6+
use asap_devtools::{lower_promql_with_data_ingestion_interval, lower_sql, SqlCatalog};
77
use asap_types::pre_asap::schema::{Column, DataType, Schema};
88
use asap_types::types::AccuracyTarget;
99

@@ -49,7 +49,7 @@ async fn show_sql(label: &str, query: &str) {
4949
fn show_promql(label: &str, query: &str) {
5050
println!("━━━ {label} ━━━");
5151
println!("{query}");
52-
match lower_promql(query, AccuracyTarget::Exact) {
52+
match lower_promql_with_data_ingestion_interval(query, AccuracyTarget::Exact, 1_000) {
5353
Ok(qe) => println!("{qe:#?}"),
5454
Err(e) => println!("ERR: {e}"),
5555
}

‎crates/devtools/src/bin/analyze_corpora.rs‎

Lines changed: 21 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44
// Corpus mode dumps all four PromQL corpora as JSONL and writes a heuristic
55
// anomaly report. The default mode remains the ad-hoc SQL/PromQL inspector.
66

7-
use asap_devtools::{lower_promql, SqlCatalog};
7+
use asap_devtools::{lower_promql_with_data_ingestion_interval, SqlCatalog};
88
use asap_frontend_sql::lower_sql_dialect;
99
use asap_types::pre_asap::schema::{Column, DataType, Schema};
1010
use asap_types::types::AccuracyTarget;
@@ -201,12 +201,16 @@ fn structural_shape(expression: &str) -> String {
201201
out
202202
}
203203

204-
fn run_corpus(name: &str, source: &str) -> CorpusResult {
204+
fn run_corpus(name: &str, source: &str, interval_ms: u64) -> CorpusResult {
205205
let mut result = CorpusResult::default();
206206
for (index, expression) in corpus_lines(source).into_iter().enumerate() {
207207
let normalized_expression = normalize(expression);
208208
let structural_shape = structural_shape(expression);
209-
match lower_promql(expression, AccuracyTarget::Exact) {
209+
match lower_promql_with_data_ingestion_interval(
210+
expression,
211+
AccuracyTarget::Exact,
212+
interval_ms,
213+
) {
210214
Ok(ir) => result.lowered.push(DumpRecord {
211215
corpus: name.to_string(),
212216
query_number: index + 1,
@@ -394,7 +398,7 @@ fn anomaly_report(all: &[DumpRecord], language: &str, manual_notes: &str) -> Str
394398
report
395399
}
396400

397-
fn run_corpora(out_dir: PathBuf) {
401+
fn run_corpora(out_dir: PathBuf, interval_ms: u64) {
398402
std::fs::create_dir_all(&out_dir)
399403
.unwrap_or_else(|e| panic!("failed to create {}: {e}", out_dir.display()));
400404
let corpora = [
@@ -406,7 +410,7 @@ fn run_corpora(out_dir: PathBuf) {
406410
let mut all = Vec::new();
407411
let mut summary = Vec::new();
408412
for (name, source) in corpora {
409-
let mut result = run_corpus(name, source);
413+
let mut result = run_corpus(name, source, interval_ms);
410414
let total = result.lowered.len() + result.failed.len();
411415
write_jsonl(&out_dir.join(format!("{name}.jsonl")), &result.lowered);
412416
write_jsonl(
@@ -542,14 +546,25 @@ async fn main() {
542546
let mut args = std::env::args().skip(1);
543547
if args.next().as_deref() == Some("--corpora") {
544548
let mut out_dir = PathBuf::from("artifacts/promql_pre_asap");
549+
let mut interval_ms = None;
545550
while let Some(arg) = args.next() {
546551
if arg == "--out-dir" {
547552
out_dir = PathBuf::from(args.next().expect("--out-dir requires a path"));
553+
} else if arg == "--data-ingestion-interval-ms" {
554+
interval_ms = Some(
555+
args.next()
556+
.expect("--data-ingestion-interval-ms requires a value")
557+
.parse()
558+
.expect("--data-ingestion-interval-ms must be an unsigned integer"),
559+
);
548560
} else {
549561
panic!("unknown corpus-mode argument: {arg}");
550562
}
551563
}
552-
run_corpora(out_dir);
564+
run_corpora(
565+
out_dir,
566+
interval_ms.expect("--data-ingestion-interval-ms is required for --corpora"),
567+
);
553568
return;
554569
}
555570
if std::env::args().nth(1).as_deref() == Some("--sql-corpora") {

0 commit comments

Comments
 (0)