Skip to content

Commit 751a16d

Browse files
fix(promql): lower workloads with ingestion intervals
1 parent e96c231 commit 751a16d

4 files changed

Lines changed: 183 additions & 54 deletions

File tree

‎crates/frontend-promql/src/error.rs‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
use std::fmt;
22

33
use asap_types::pre_asap::ResolveTreeError;
4+
use asap_types::workload::WorkloadError;
45

56
/// Errors from lowering a PromQL query (parse → the canonical, unresolved
67
/// tree, built directly →
@@ -13,6 +14,8 @@ use asap_types::pre_asap::ResolveTreeError;
1314
/// shared, so neither front end pulls the other's parser.
1415
#[derive(Debug)]
1516
pub enum PromqlError {
17+
/// The workload omitted information required for plan-ready PromQL lowering.
18+
InvalidWorkload(WorkloadError),
1619
/// The `promql-parser` crate rejected the query string (parse failure).
1720
Parse(String),
1821
/// A PromQL function (`rate`, `*_over_time`, …) not supported in this version.
@@ -36,6 +39,7 @@ pub enum PromqlError {
3639
impl fmt::Display for PromqlError {
3740
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3841
match self {
42+
Self::InvalidWorkload(e) => write!(f, "invalid PromQL workload: {e}"),
3943
Self::Parse(e) => write!(f, "PromQL parse error: {e}"),
4044
Self::UnsupportedFunction(n) => write!(f, "unsupported PromQL function: {n}"),
4145
Self::UnsupportedAggregateOp(n) => write!(f, "unsupported PromQL aggregate op: {n}"),
@@ -56,6 +60,12 @@ impl From<ResolveTreeError> for PromqlError {
5660
}
5761
}
5862

63+
impl From<WorkloadError> for PromqlError {
64+
fn from(e: WorkloadError) -> Self {
65+
Self::InvalidWorkload(e)
66+
}
67+
}
68+
5969
#[cfg(test)]
6070
mod tests {
6171
use super::*;

‎crates/frontend-promql/src/lib.rs‎

Lines changed: 94 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -14,65 +14,111 @@ pub mod promql;
1414

1515
use asap_types::pre_asap::resolve_root;
1616
use asap_types::pre_asap::QueryExpr;
17-
use asap_types::types::AccuracyTarget;
18-
use asap_types::workload::{QueryLanguage, QueryWorkload};
17+
use asap_types::workload::{DurationMs, QueryLanguage, QueryWorkload};
1918

2019
pub use error::PromqlError;
2120
pub use histogram::{HistogramCatalog, HistogramKind};
2221
pub use promql::PromqlLowerer;
2322

24-
/// Lower a single PromQL query string to the canonical, resolved `QueryExpr`.
23+
/// Lower every normalized PromQL workload entry to a plan-ready `QueryExpr`.
2524
///
26-
/// `accuracy` is threaded onto every approximate intent (`Count`, `Quantile`,
27-
/// `Cardinality`, `TopK`). The returned tree carries a self-contained `Schema`
28-
/// on its `Scan`; call [`QueryExpr::output_schema`] for any node's schema.
29-
///
30-
/// `histogram_quantile` discrimination uses the structural heuristic; to drive
31-
/// it from declared sample types instead, use [`lower_promql_with_histograms`].
32-
pub fn lower_promql(query: &str, accuracy: AccuracyTarget) -> Result<QueryExpr, PromqlError> {
33-
let unresolved = PromqlLowerer::lower(query, &accuracy)?;
34-
let resolved = resolve_root(&unresolved)?;
35-
Ok(resolved)
25+
/// PromQL workloads must declare a non-zero `data_ingestion_interval`; it is
26+
/// injected around each bare instant selector. Explicit range selectors keep
27+
/// their query-specified range.
28+
pub fn lower_promql_workload(workload: &QueryWorkload) -> Result<Vec<QueryExpr>, PromqlError> {
29+
if !matches!(workload.language, QueryLanguage::PromQL) {
30+
return Err(PromqlError::WrongLanguage(format!(
31+
"{:?}",
32+
workload.language
33+
)));
34+
}
35+
workload.validate()?;
36+
let DurationMs(interval_ms) = workload
37+
.data_workload
38+
.as_ref()
39+
.expect("validated PromQL workload has data_workload")
40+
.data_ingestion_interval
41+
.value
42+
.expect("validated PromQL workload has data_ingestion_interval");
43+
workload
44+
.entries()
45+
.map(|entry| {
46+
let unresolved = PromqlLowerer::lower_with_ingestion_interval(
47+
&entry.query.0,
48+
&entry.requirements.accuracy.target(),
49+
std::time::Duration::from_millis(interval_ms),
50+
)?;
51+
Ok(resolve_root(&unresolved)?)
52+
})
53+
.collect()
3654
}
3755

38-
/// Like [`lower_promql`], but consults `histograms` to decide whether a
39-
/// `histogram_quantile` argument is sketch-able (generic `Quantile`) or a
40-
/// classic-bucket interpolation (`HistogramQuantile`) — a type-driven decision
41-
/// instead of the structural heuristic (issue #79). Metrics absent from the
42-
/// catalog still fall back to the heuristic.
43-
pub fn lower_promql_with_histograms(
44-
query: &str,
45-
accuracy: AccuracyTarget,
46-
histograms: HistogramCatalog,
47-
) -> Result<QueryExpr, PromqlError> {
48-
let _guard = histogram::CatalogGuard::install(histograms);
49-
lower_promql(query, accuracy)
50-
}
56+
#[cfg(test)]
57+
mod tests {
58+
use std::time::Duration;
5159

52-
/// Lower every PromQL batch entry in `workload` to a `QueryExpr`.
53-
///
54-
/// One `Result` per entry — errors are per-query, not fatal for the batch.
55-
/// Returns an empty `Vec` if `workload.query_batch` is absent or empty, and a
56-
/// `WrongLanguage` error for every entry if the workload language is not PromQL.
57-
pub fn lower_promql_batch(workload: &QueryWorkload) -> Vec<Result<QueryExpr, PromqlError>> {
58-
let entries = match &workload.query_batch {
59-
Some(e) if !e.is_empty() => e,
60-
_ => return vec![],
60+
use asap_types::pre_asap::QueryExpr;
61+
use asap_types::workload::{
62+
BatchEntry, DataWorkload, Evidence, Query, QueryRequirements, TimeSelection,
6163
};
6264

63-
if !matches!(workload.language, QueryLanguage::PromQL) {
64-
let lang = format!("{:?}", workload.language);
65-
return entries
66-
.iter()
67-
.map(|_| Err(PromqlError::WrongLanguage(lang.clone())))
68-
.collect();
65+
use super::*;
66+
67+
fn workload(query: &str) -> QueryWorkload {
68+
QueryWorkload {
69+
language: QueryLanguage::PromQL,
70+
query_batch: Some(vec![BatchEntry {
71+
query: Query(query.into()),
72+
requirements: QueryRequirements::default(),
73+
predictability: Default::default(),
74+
invocations: 1,
75+
execute_at: None,
76+
time_selection: TimeSelection::default(),
77+
}]),
78+
repeating_queries: None,
79+
data_workload: Some(DataWorkload {
80+
data_ingestion_interval: Evidence {
81+
value: Some(DurationMs(1_000)),
82+
..Default::default()
83+
},
84+
..Default::default()
85+
}),
86+
}
6987
}
7088

71-
entries
72-
.iter()
73-
.map(|entry| {
74-
let accuracy = entry.requirements.accuracy.target();
75-
lower_promql(&entry.query.0, accuracy)
76-
})
77-
.collect()
89+
#[test]
90+
fn instant_selector_uses_declared_ingestion_interval() {
91+
let query = lower_promql_workload(&workload("sum by (job) (data)")).unwrap();
92+
let QueryExpr::Aggregate { child, .. } = &query[0] else {
93+
panic!("expected aggregate")
94+
};
95+
assert!(
96+
matches!(child.as_ref(), QueryExpr::TimeRange { range, child }
97+
if *range == Duration::from_secs(1) && matches!(child.as_ref(), QueryExpr::Scan { .. }))
98+
);
99+
}
100+
101+
#[test]
102+
fn explicit_range_selector_keeps_its_query_range() {
103+
let query = lower_promql_workload(&workload("sum_over_time(data[5m])")).unwrap();
104+
let QueryExpr::Aggregate { child, .. } = &query[0] else {
105+
panic!("expected aggregate")
106+
};
107+
assert!(
108+
matches!(child.as_ref(), QueryExpr::TimeRange { range, child }
109+
if *range == Duration::from_secs(300) && matches!(child.as_ref(), QueryExpr::Scan { .. }))
110+
);
111+
}
112+
113+
#[test]
114+
fn workload_without_interval_fails_loudly() {
115+
let mut workload = workload("sum(data)");
116+
workload.data_workload = Some(DataWorkload::default());
117+
assert!(matches!(
118+
lower_promql_workload(&workload),
119+
Err(PromqlError::InvalidWorkload(
120+
asap_types::workload::WorkloadError::MissingDataIngestionInterval
121+
))
122+
));
123+
}
78124
}

‎crates/frontend-promql/src/promql.rs‎

Lines changed: 50 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -191,11 +191,24 @@ impl PromqlLowerer {
191191
check_depth(&ast, MAX_DEPTH)?;
192192
walk(&ast)
193193
}
194+
195+
pub(crate) fn lower_with_ingestion_interval(
196+
query: &str,
197+
accuracy: &AccuracyTarget,
198+
interval: Duration,
199+
) -> Result<Unresolved> {
200+
let _guard = AccuracyGuard::install(accuracy.clone());
201+
let _interval = IngestionIntervalGuard::install(interval);
202+
let ast = parser::parse(query).map_err(LoweringError::Parse)?;
203+
check_depth(&ast, MAX_DEPTH)?;
204+
walk(&ast)
205+
}
194206
}
195207

196208
std::thread_local! {
197209
static ACCURACY: std::cell::RefCell<AccuracyTarget> =
198210
const { std::cell::RefCell::new(AccuracyTarget::Exact) };
211+
static INGESTION_INTERVAL: std::cell::RefCell<Option<Duration>> = const { std::cell::RefCell::new(None) };
199212
}
200213

201214
/// RAII guard installing `accuracy` as the ambient accuracy target for the
@@ -221,6 +234,28 @@ fn current_accuracy() -> AccuracyTarget {
221234
ACCURACY.with(|a| a.borrow().clone())
222235
}
223236

237+
struct IngestionIntervalGuard(Option<Duration>);
238+
239+
impl IngestionIntervalGuard {
240+
fn install(interval: Duration) -> Self {
241+
Self(INGESTION_INTERVAL.with(|current| current.replace(Some(interval))))
242+
}
243+
}
244+
245+
impl Drop for IngestionIntervalGuard {
246+
fn drop(&mut self) {
247+
INGESTION_INTERVAL.with(|current| *current.borrow_mut() = self.0.take());
248+
}
249+
}
250+
251+
fn current_ingestion_interval() -> Duration {
252+
INGESTION_INTERVAL.with(|current| {
253+
current
254+
.borrow()
255+
.expect("ingestion interval is installed for workload lowering")
256+
})
257+
}
258+
224259
/// Bounded depth check over the parser AST: errors once nesting would exceed
225260
/// `budget` frames, descending into every child expression.
226261
fn check_depth(expr: &Expr, budget: usize) -> Result<()> {
@@ -299,7 +334,7 @@ fn walk(expr: &Expr) -> Result<Unresolved> {
299334
}),
300335
Expr::VectorSelector(vs) => {
301336
let (metric, matchers, shift) = vs_parts(vs)?;
302-
Ok(filtered_source(metric, matchers, shift))
337+
Ok(instant_source(metric, matchers, shift))
303338
}
304339
Expr::MatrixSelector(ms) => {
305340
let (metric, matchers, shift) = vs_parts(&ms.vs)?;
@@ -1421,7 +1456,7 @@ fn lower_inner_call(call: &Call) -> Result<Inner> {
14211456
fn build(inner: Inner, keys: Vec<ColumnRef>, outer: Outer) -> Result<Unresolved> {
14221457
match outer {
14231458
Outer::None => match &inner.func {
1424-
None => Ok(filtered_source(inner.metric, inner.matchers, inner.shift)),
1459+
None => Ok(instant_source(inner.metric, inner.matchers, inner.shift)),
14251460
Some(f) => {
14261461
let intent = inner_intent(f);
14271462
Ok(windowed_aggregate(inner, keys, intent))
@@ -1463,7 +1498,7 @@ fn build(inner: Inner, keys: Vec<ColumnRef>, outer: Outer) -> Result<Unresolved>
14631498
// wrapped in a reducing aggregate (issue #86).
14641499
let base = match inner.func.as_ref().map(inner_intent) {
14651500
Some(intent) => windowed_aggregate(inner, vec![], intent),
1466-
None => filtered_source(inner.metric, inner.matchers, inner.shift),
1501+
None => instant_source(inner.metric, inner.matchers, inner.shift),
14671502
};
14681503
Ok(Unresolved::PromqlSeriesSample {
14691504
by: keys.into(),
@@ -1526,7 +1561,7 @@ fn build(inner: Inner, keys: Vec<ColumnRef>, outer: Outer) -> Result<Unresolved>
15261561
// `Sort.partition_by` can rank within each group (issue #12).
15271562
let base = match inner.func.as_ref().map(inner_intent) {
15281563
Some(intent) => windowed_aggregate(inner, vec![], intent),
1529-
None => filtered_source(inner.metric, inner.matchers, inner.shift),
1564+
None => instant_source(inner.metric, inner.matchers, inner.shift),
15301565
};
15311566
let sorted = Unresolved::Sort {
15321567
keys: vec![SortKey {
@@ -1588,7 +1623,10 @@ fn windowed_aggregate(
15881623
range: w,
15891624
child: Rc::new(base),
15901625
},
1591-
None => base,
1626+
None => Unresolved::TimeRange {
1627+
range: current_ingestion_interval(),
1628+
child: Rc::new(base),
1629+
},
15921630
};
15931631
let reduction = reduction_for(&keys, &intent, &child);
15941632
Unresolved::Aggregate {
@@ -1641,6 +1679,13 @@ fn filtered_source(metric: String, matchers: Vec<Unresolved>, shift: TimeShift)
16411679
}
16421680
}
16431681

1682+
fn instant_source(metric: String, matchers: Vec<Unresolved>, shift: TimeShift) -> Unresolved {
1683+
Unresolved::TimeRange {
1684+
range: current_ingestion_interval(),
1685+
child: Rc::new(filtered_source(metric, matchers, shift)),
1686+
}
1687+
}
1688+
16441689
/// Count vector elements regardless of their sample values.
16451690
fn count() -> AggIntent<ColumnRef> {
16461691
AggIntent::Count {

‎crates/types/src/workload.rs‎

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -530,6 +530,9 @@ pub struct Rate(pub f64);
530530
#[serde(deny_unknown_fields)]
531531
pub struct DataWorkload {
532532
pub arrival: DataArrival,
533+
/// Cadence at which each PromQL source supplies a sample. PromQL instant
534+
/// selectors use this as their explicit selection horizon.
535+
pub data_ingestion_interval: Evidence<DurationMs>,
533536
pub ingestion_volume: Evidence<u64>,
534537
pub ingestion_rate: Evidence<Rate>,
535538
pub input_cardinality: Evidence<u64>,
@@ -586,6 +589,19 @@ impl QueryWorkload {
586589
}
587590
validate_optional_rate(data.ingestion_rate.value)?;
588591
}
592+
if matches!(self.language, QueryLanguage::PromQL) {
593+
let data = self
594+
.data_workload
595+
.as_ref()
596+
.ok_or(WorkloadError::MissingPromqlDataWorkload)?;
597+
let interval = data
598+
.data_ingestion_interval
599+
.value
600+
.ok_or(WorkloadError::MissingDataIngestionInterval)?;
601+
if interval.0 == 0 {
602+
return Err(WorkloadError::ZeroDataIngestionInterval);
603+
}
604+
}
589605
Ok(())
590606
}
591607
}
@@ -638,6 +654,12 @@ fn validate_rate(rate: Rate) -> Result<Rate, WorkloadError> {
638654

639655
#[derive(Debug, Clone, Copy, PartialEq, thiserror::Error)]
640656
pub enum WorkloadError {
657+
#[error("a PromQL workload requires data_workload")]
658+
MissingPromqlDataWorkload,
659+
#[error("a PromQL workload requires data_ingestion_interval")]
660+
MissingDataIngestionInterval,
661+
#[error("data_ingestion_interval must be greater than zero")]
662+
ZeroDataIngestionInterval,
641663
#[error("a one-time query must have at least one invocation")]
642664
ZeroInvocations,
643665
#[error("a fixed repetition interval must be greater than zero")]
@@ -665,7 +687,13 @@ mod tests {
665687
language: QueryLanguage::PromQL,
666688
query_batch: None,
667689
repeating_queries: None,
668-
data_workload: None,
690+
data_workload: Some(DataWorkload {
691+
data_ingestion_interval: Evidence {
692+
value: Some(DurationMs(1_000)),
693+
..Default::default()
694+
},
695+
..Default::default()
696+
}),
669697
}
670698
}
671699

0 commit comments

Comments
 (0)