Skip to content
Merged
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/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,8 @@ pub mod storage_io;
pub mod summary_maintenance_cost;
pub mod summary_maintenance_dag_export;
pub mod summary_maintenance_lifecycle;
#[cfg(test)]
mod test_support;
pub mod topk_reuse;

pub use accuracy::{
Expand Down
24 changes: 18 additions & 6 deletions crates/asap-aware-mapping/src/maintained_population.rs
Original file line number Diff line number Diff line change
Expand Up @@ -118,11 +118,24 @@ fn recognize(root: &QueryExpr) -> Option<(MaintainedPopulation, PopulationReadou
}
return Some((population, readout, Rc::clone(source)));
}
// A bare PromQL selector carries the declared ingestion interval as a
// temporal input scope. Membership must expire at that horizon; retain
// the wrapper as the maintained input so validation can check agreement.
let (series_source, lookback_ms) = match source.as_ref() {
QueryExpr::TimeRange { range, child } => {
let ms = u64::try_from(range.as_millis()).ok()?;
if ms == 0 || std::time::Duration::from_millis(ms) != *range {
return None;
}
(child.as_ref(), ms)
}
other => (other, 300_000),
};
let QueryExpr::Scan {
source: Source::TimeSeries { metric },
predicates,
schema,
} = source.as_ref()
} = series_source
else {
return None;
};
Expand Down Expand Up @@ -176,7 +189,7 @@ fn recognize(root: &QueryExpr) -> Option<(MaintainedPopulation, PopulationReadou
matchers,
grouping: labels,
without: grouping.is_without(),
lookback_ms: 300_000,
lookback_ms,
}),
max_k: 0,
quantiles: false,
Expand Down Expand Up @@ -285,12 +298,11 @@ impl ReplacementStrategy for MaintainedPopulationStrategy {
#[cfg(test)]
mod tests {
use super::*;
use crate::test_support::lower_promql;
use asap_types::post_asap::{compile_executable_dag, share_common_summary_subtrees};

fn lower(q: &str) -> Rc<QueryExpr> {
Rc::new(
asap_frontend_promql::lower_promql(q, asap_types::types::AccuracyTarget::Exact)
.unwrap(),
)
Rc::new(lower_promql(q, asap_types::types::AccuracyTarget::Exact))
}

// Instant scalar aggregations share the same retractable series population.
Expand Down
11 changes: 4 additions & 7 deletions crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5779,6 +5779,7 @@ mod tests {
use super::*;
use crate::accuracy::PropagationStats;
use crate::cost_model::Cost;
use crate::test_support::lower_promql;
use asap_types::pre_asap::agg_intent::{
agg_is_exact, default_cardinality, default_quantile, MathFunc, TimeFunc,
};
Expand All @@ -5798,10 +5799,7 @@ mod tests {
// Finite samples can overflow a sum although their native average is finite.
#[test]
fn temporal_average_requires_finite_division_guard() {
let root = Rc::new(
asap_frontend_promql::lower_promql("avg_over_time(a[5m])", AccuracyTarget::Exact)
.unwrap(),
);
let root = Rc::new(lower_promql("avg_over_time(a[5m])", AccuracyTarget::Exact));
let candidates =
SketchAlgorithmStrategy::default_cost_model().replacements(&TargetSubDAG::new(&root));
let operator = candidates
Expand Down Expand Up @@ -5830,8 +5828,7 @@ mod tests {
"topk(5, sum_over_time(a[5m]))",
"topk by(job)(5, count_over_time(a[5m]))",
] {
let root =
Rc::new(asap_frontend_promql::lower_promql(query, AccuracyTarget::Exact).unwrap());
let root = Rc::new(lower_promql(query, AccuracyTarget::Exact));
let models = Models::with_default_accuracy(&crate::cost_model::DefaultCostModel);
let node = exact_topk_over_temporal_values(&root, models)
.unwrap()
Expand Down Expand Up @@ -5859,7 +5856,7 @@ mod tests {
"quantile_over_time(0.5,a[5m]) / quantile_over_time(0.9,a[5m])",
"avg_over_time(a[5m]) / quantile_over_time(0.5,a[5m])",
] {
let root = Rc::new(asap_frontend_promql::lower_promql(query, target.clone()).unwrap());
let root = Rc::new(lower_promql(query, target.clone()));
let models = Models::with_default_accuracy(&crate::cost_model::DefaultCostModel);
assert!(realize_binary(&root, models, Some(&target))
.unwrap()
Expand Down
13 changes: 6 additions & 7 deletions crates/asap-aware-mapping/src/rewrite.rs
Original file line number Diff line number Diff line change
Expand Up @@ -389,8 +389,10 @@ impl ReplacementStrategy for SemanticEquivalentRewriteStrategy {
#[cfg(test)]
mod tests {
use super::*;
use crate::test_support::lower_promql;
use asap_types::pre_asap::query_expr::Source;
use asap_types::pre_asap::schema::{Column, Schema};
use asap_types::types::AccuracyTarget;
use std::time::Duration;

fn metric_scan(labels: &[&str]) -> QueryExpr {
Expand Down Expand Up @@ -419,13 +421,10 @@ mod tests {
// Temporal averages expose two single-measure children without closing labels.
#[test]
fn temporal_average_components_preserves_schema_and_exposes_sum_count() {
let root = Rc::new(
asap_frontend_promql::lower_promql(
"avg_over_time(a{job=\"api\"}[5m])",
AccuracyTarget::Exact,
)
.unwrap(),
);
let root = Rc::new(lower_promql(
"avg_over_time(a{job=\"api\"}[5m])",
AccuracyTarget::Exact,
));
assert!(SemanticEquivalentRewriteStrategy
.replacements(&TargetSubDAG::new(&root))
.is_empty());
Expand Down
10 changes: 8 additions & 2 deletions crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3083,7 +3083,10 @@ mod tests {
fn streaming_data_workload() -> DataWorkload {
DataWorkload {
arrival: DataArrival::ContinuouslyIngesting,

data_ingestion_interval: Evidence {
value: Some(asap_types::workload::DurationMs(1_000)),
..Default::default()
},
ingestion_rate: Evidence {
value: Some(Rate(2.0)),
source: EvidenceSource::Declared,
Expand Down Expand Up @@ -3491,7 +3494,10 @@ mod tests {
ComparisonScope::from_workload(
&DataWorkload {
arrival: DataArrival::ContinuouslyIngesting,

data_ingestion_interval: Evidence {
value: Some(asap_types::workload::DurationMs(1_000)),
..Default::default()
},
ingestion_rate: Evidence {
value: Some(Rate(2.0)),
source: EvidenceSource::Declared,
Expand Down
10 changes: 8 additions & 2 deletions crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1690,15 +1690,21 @@ mod tests {
fn at_rest() -> DataWorkload {
DataWorkload {
arrival: DataArrival::AtRest,

data_ingestion_interval: Evidence {
value: Some(DurationMs(1_000)),
..Default::default()
},
..Default::default()
}
}

fn continuous(observed_at_ms: u64, valid_for_ms: u64) -> DataWorkload {
DataWorkload {
arrival: DataArrival::ContinuouslyIngesting,

data_ingestion_interval: Evidence {
value: Some(DurationMs(1_000)),
..Default::default()
},
ingestion_rate: Evidence {
value: Some(Rate(1.0)),
source: EvidenceSource::Observed,
Expand Down
37 changes: 37 additions & 0 deletions crates/asap-aware-mapping/src/test_support.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
use asap_types::pre_asap::QueryExpr;
use asap_types::types::AccuracyTarget;
use asap_types::workload::{
AccuracyRequirement, BatchEntry, DataWorkload, DurationMs, Evidence, PlanningWorkload,
Predictability, Query, QueryLanguage, QueryRequirements, QueryWorkload, TimeSelection,
};

pub(crate) fn lower_promql(query: &str, accuracy: AccuracyTarget) -> QueryExpr {
let workload = PlanningWorkload {
query_workload: QueryWorkload {
language: QueryLanguage::PromQL,
query_batch: Some(vec![BatchEntry {
query: Query(query.into()),
requirements: QueryRequirements {
accuracy: AccuracyRequirement::Explicit(accuracy),
..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(1_000)),
..Default::default()
},
..Default::default()
}),
};
asap_frontend_promql::lower_promql_workload(&workload, 0)
.unwrap()
.pop()
.unwrap()
}
4 changes: 2 additions & 2 deletions crates/devtools/examples/canonical_examples.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
// One-off: pretty-print the QueryExpr for one canonical query per variant,
// plus custom Join/SetOp/Dedup/CTE probes, to eyeball the actual shape.

use asap_devtools::lower_promql;
use asap_devtools::lower_promql_with_data_ingestion_interval;
use asap_frontend_sql::{lower_sql_dialect, SqlCatalog};
use asap_types::pre_asap::schema::{Column, DataType, Schema};
use asap_types::types::AccuracyTarget;
Expand Down Expand Up @@ -70,7 +70,7 @@ async fn main() {
];
for (label, q) in promql_examples {
println!("=== {label} === promql> {q}");
match lower_promql(q, AccuracyTarget::Exact) {
match lower_promql_with_data_ingestion_interval(q, AccuracyTarget::Exact, 1_000) {
Ok(qe) => println!("{qe:#?}"),
Err(e) => println!("ERR: {e}"),
}
Expand Down
4 changes: 2 additions & 2 deletions crates/devtools/examples/topk_ir.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
// Lowers every topk-shaped query from the design discussion and prints the
// resulting pre-ASAP IR. Used for interactive exploration; not a test.

use asap_devtools::{lower_promql, lower_sql, SqlCatalog};
use asap_devtools::{lower_promql_with_data_ingestion_interval, lower_sql, SqlCatalog};
use asap_types::pre_asap::schema::{Column, DataType, Schema};
use asap_types::types::AccuracyTarget;

Expand Down Expand Up @@ -49,7 +49,7 @@ async fn show_sql(label: &str, query: &str) {
fn show_promql(label: &str, query: &str) {
println!("━━━ {label} ━━━");
println!("{query}");
match lower_promql(query, AccuracyTarget::Exact) {
match lower_promql_with_data_ingestion_interval(query, AccuracyTarget::Exact, 1_000) {
Ok(qe) => println!("{qe:#?}"),
Err(e) => println!("ERR: {e}"),
}
Expand Down
27 changes: 21 additions & 6 deletions crates/devtools/src/bin/analyze_corpora.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
// Corpus mode dumps all four PromQL corpora as JSONL and writes a heuristic
// anomaly report. The default mode remains the ad-hoc SQL/PromQL inspector.

use asap_devtools::{lower_promql, SqlCatalog};
use asap_devtools::{lower_promql_with_data_ingestion_interval, SqlCatalog};
use asap_frontend_sql::lower_sql_dialect;
use asap_types::pre_asap::schema::{Column, DataType, Schema};
use asap_types::types::AccuracyTarget;
Expand Down Expand Up @@ -201,12 +201,16 @@ fn structural_shape(expression: &str) -> String {
out
}

fn run_corpus(name: &str, source: &str) -> CorpusResult {
fn run_corpus(name: &str, source: &str, interval_ms: u64) -> CorpusResult {
let mut result = CorpusResult::default();
for (index, expression) in corpus_lines(source).into_iter().enumerate() {
let normalized_expression = normalize(expression);
let structural_shape = structural_shape(expression);
match lower_promql(expression, AccuracyTarget::Exact) {
match lower_promql_with_data_ingestion_interval(
expression,
AccuracyTarget::Exact,
interval_ms,
) {
Ok(ir) => result.lowered.push(DumpRecord {
corpus: name.to_string(),
query_number: index + 1,
Expand Down Expand Up @@ -394,7 +398,7 @@ fn anomaly_report(all: &[DumpRecord], language: &str, manual_notes: &str) -> Str
report
}

fn run_corpora(out_dir: PathBuf) {
fn run_corpora(out_dir: PathBuf, interval_ms: u64) {
std::fs::create_dir_all(&out_dir)
.unwrap_or_else(|e| panic!("failed to create {}: {e}", out_dir.display()));
let corpora = [
Expand All @@ -406,7 +410,7 @@ fn run_corpora(out_dir: PathBuf) {
let mut all = Vec::new();
let mut summary = Vec::new();
for (name, source) in corpora {
let mut result = run_corpus(name, source);
let mut result = run_corpus(name, source, interval_ms);
let total = result.lowered.len() + result.failed.len();
write_jsonl(&out_dir.join(format!("{name}.jsonl")), &result.lowered);
write_jsonl(
Expand Down Expand Up @@ -542,14 +546,25 @@ async fn main() {
let mut args = std::env::args().skip(1);
if args.next().as_deref() == Some("--corpora") {
let mut out_dir = PathBuf::from("artifacts/promql_pre_asap");
let mut interval_ms = None;
while let Some(arg) = args.next() {
if arg == "--out-dir" {
out_dir = PathBuf::from(args.next().expect("--out-dir requires a path"));
} else if arg == "--data-ingestion-interval-ms" {
interval_ms = Some(
args.next()
.expect("--data-ingestion-interval-ms requires a value")
.parse()
.expect("--data-ingestion-interval-ms must be an unsigned integer"),
);
} else {
panic!("unknown corpus-mode argument: {arg}");
}
}
run_corpora(out_dir);
run_corpora(
out_dir,
interval_ms.expect("--data-ingestion-interval-ms is required for --corpora"),
);
return;
}
if std::env::args().nth(1).as_deref() == Some("--sql-corpora") {
Expand Down
Loading
Loading