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
37 changes: 28 additions & 9 deletions crates/executor/src/physical_planner/candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -364,16 +364,35 @@ mod tests {
.unwrap()
.remove(0);
let root = promql_rows::with_series_identity(&root).unwrap();
let space = asap_logical_optimizer::search_workload(vec![("q", root)]);
let selected = asap_plan_selection::candidate_selection::global_selection(
&space,
&asap_plan_selection::cost::cost_model::DefaultCostModel,
)
.assemble_selected_dag(&space.roots[0].1)
.unwrap()
.unwrap();
let demand = [RootDemand {
accuracy: Some(planner_types::types::AccuracyTarget::Exact),
recurrence: QueryRecurrence::OneTime {
invocations: 1,
execute_at: None,
},
predictability: Predictability::default(),
latency_ms: None,
}];
let data = DataWorkload {
arrival: DataArrival::ContinuouslyIngesting,
ingestion_rate: Evidence {
value: Some(Rate(1_000.0)),
source: EvidenceSource::Declared,
..Default::default()
},
..Default::default()
};
let capabilities = crate::capabilities();
let models =
asap_plan_selection::PlanningModels::builtin().with_capabilities(&capabilities);
let root = planner_types::ir::QueryRoot::Operator(root);
let run =
asap_plan_selection::plan_stages(vec![(0, root)], &demand, &data, models, 0).unwrap();
let planner_types::ir::QueryRoot::Operator(selected) = &run.plan.logical[0].1 else {
panic!("operator root")
};
let selected = planner_types::ir::apply_materialization_timings(
&selected,
selected,
&planner_types::ir::MaterializationAssignment::all_ingestion_time(),
&mut Default::default(),
)
Expand Down
46 changes: 46 additions & 0 deletions crates/executor/tests/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,3 +15,49 @@ pub fn compile_physical_asap_dag(
)?;
Ok(planner_types::ir::export::compile_physical_asap_dag(&root)?)
}

/// The logical DAG the #509 stage pipeline selects for the single query
/// `root` at `accuracy`, run once over a continuously ingested source, with
/// the built-in models and this executor's capabilities.
pub fn selected_dag(
root: Rc<OperatorNode>,
accuracy: planner_types::types::AccuracyTarget,
) -> Rc<OperatorNode> {
use planner_types::ir::QueryRoot;
use planner_types::workload::{
DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate,
RootDemand,
};
let demand = [RootDemand {
accuracy: Some(accuracy),
recurrence: QueryRecurrence::OneTime {
invocations: 1,
execute_at: None,
},
predictability: Predictability::default(),
latency_ms: None,
}];
let data = DataWorkload {
arrival: DataArrival::ContinuouslyIngesting,
ingestion_rate: Evidence {
value: Some(Rate(1_000.0)),
source: EvidenceSource::Declared,
..Default::default()
},
..Default::default()
};
let capabilities = asap_executor::capabilities();
let models = asap_plan_selection::PlanningModels::builtin().with_capabilities(&capabilities);
let run = asap_plan_selection::plan_stages(
vec![(0, QueryRoot::Operator(root))],
&demand,
&data,
models,
0,
)
.unwrap();
let QueryRoot::Operator(root) = &run.plan.logical[0].1 else {
panic!("operator root")
};
root.clone()
}
26 changes: 6 additions & 20 deletions crates/executor/tests/deployment_computation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ use asap_executor::{
runtime::{Limits, RunContext, Scope},
values::{Batch, Value},
};
use common::compile_physical_asap_dag;
use common::{compile_physical_asap_dag, selected_dag};
use futures::{executor::block_on, StreamExt};
use planner_types::ir::export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload};
use planner_types::{ir::schema::*, types::AccuracyTarget, workload::*};
Expand Down Expand Up @@ -47,19 +47,11 @@ fn lower_with(query: &str, accuracy: AccuracyTarget) -> Rc<planner_types::ir::Op
.remove(0)
}

/// The first exact summary candidate, as Planner selection would hand it over.
/// The plan the stage pipeline selects for an exact `query`.
fn exact_dag(query: &str) -> PhysicalASAPDAG {
let expression = lower(query);
let root = promql_rows::with_series_identity(&expression).unwrap_or(expression);
let space = asap_logical_optimizer::search_workload(vec![("q", root)]);
let selected = asap_plan_selection::candidate_selection::global_selection(
&space,
&asap_plan_selection::DefaultCostModel,
)
.assemble_selected_dag(&space.roots[0].1)
.unwrap()
.unwrap();
compile_physical_asap_dag(&selected).unwrap()
compile_physical_asap_dag(&selected_dag(root, AccuracyTarget::Exact)).unwrap()
}

fn population_dag(query: &str) -> PhysicalASAPDAG {
Expand Down Expand Up @@ -446,12 +438,11 @@ fn counter(
(1..=5).map(move |i| (metric, job, "x", i * 60_000, base + step * (i - 1) as f64))
}

// avg_over_time over stored per-series sum/count state divides per series and
// drops the metric name, as Prometheus does. An overflowing stored sum fails
// the checked division instead of returning +Inf.
// Stored per-series sum and count states divide per series and drop the
// metric name, as Prometheus does.
#[test]
fn per_series_average_divides_stored_sum_by_count() {
let dag = exact_dag("avg_over_time(m[5m])");
let dag = exact_dag("sum_over_time(m[5m]) / count_over_time(m[5m])");
assert_eq!(
run_series(&dag, SAMPLES, 60_000).unwrap(),
series(&[
Expand All @@ -461,11 +452,6 @@ fn per_series_average_divides_stored_sum_by_count() {
("db", "d", 5.)
])
);
let huge: &[Sample] = &[
("m", "api", "a", 10_000, 1.7e308),
("m", "api", "a", 20_000, 1.7e308),
];
assert!(run_series(&dag, huge, 60_000).is_err());
}

// rate(a) / rate(b) matches series on their labels without the metric name.
Expand Down
19 changes: 4 additions & 15 deletions crates/executor/tests/planspace_series_identity_heap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
//! current-series TopK heaps without a caller-side series-identity pass, cost
//! ranking, or workload Cartesian expansion. Placement variants are not listed.
mod common;
use common::compile_physical_asap_dag;
use common::{compile_physical_asap_dag, selected_dag};
use planner_types::ir::OperatorNode;

use asap_executor::physical_planner::promql_rows::{
Expand All @@ -15,8 +15,6 @@ use asap_logical_optimizer::{
pass1::replacement::ReplacementProvenance, search_workload_with_targets, Proposals,
ReplacementStrategy, ReplacementSubDAG, TargetSubDAG,
};
use asap_plan_selection::candidate_selection::global_selection;
use asap_plan_selection::cost::cost_model::DefaultCostModel;
use planner_types::{
ir::properties::*,
ir::schema::*,
Expand Down Expand Up @@ -211,21 +209,12 @@ fn unrelated_queries_keep_their_inventory() {
}
}

// Default cost-based selection keeps the logical plan; deployment prices heaps.
// The stage pipeline's selection keeps the logical plan; deployment prices heaps.
#[test]
fn global_selection_never_commits_a_series_identity_heap() {
fn selection_never_commits_a_series_identity_heap() {
let accuracy = AccuracyTarget::Epsilon(0.1);
let root = lower(CURRENT_SERIES_TOPK, &accuracy);
let strategies = default_strategies_with_evidence(&Evidence);
let space = search_workload_with_targets(
vec![(0, root, Some(accuracy))],
&strategies,
&DefaultAccuracyModel,
);
let selected = global_selection(&space, &DefaultCostModel)
.assemble_selected_dag(&space.roots[0].1)
.unwrap()
.unwrap();
let selected = selected_dag(root, accuracy);
assert!(!carries_identity(&vec![(0, selected)]));
}

Expand Down
20 changes: 8 additions & 12 deletions crates/executor/tests/precompute_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,17 +12,15 @@ use asap_executor::{
values::{Batch, Value},
};
use asap_logical_optimizer::search_workload;
use asap_plan_selection::candidate_selection::global_selection;
use asap_plan_selection::cost::cost_model::DefaultCostModel;
use common::compile_physical_asap_dag;
use common::{compile_physical_asap_dag, selected_dag};
use futures::{executor::block_on, StreamExt};
use planner_types::ir::export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload};
use planner_types::ir::schema::{DataType, *};
use planner_types::types::AccuracyTarget;
use planner_types::workload::*;
use std::{collections::BTreeMap, sync::Arc};

fn grouped_rate_space() -> asap_logical_optimizer::CandidateLogicalASAPDAGs<&'static str> {
fn grouped_rate_root() -> std::rc::Rc<planner_types::ir::OperatorNode> {
let workload = PlanningWorkload {
query_workload: QueryWorkload {
language: QueryLanguage::PromQL,
Expand Down Expand Up @@ -50,17 +48,15 @@ fn grouped_rate_space() -> asap_logical_optimizer::CandidateLogicalASAPDAGs<&'st
let root = asap_frontend_promql::lower_promql_workload(&workload, 0)
.unwrap()
.remove(0);
let root = asap_executor::physical_planner::promql_rows::with_series_identity(&root).unwrap();
search_workload(vec![("grouped-rate", root)])
asap_executor::physical_planner::promql_rows::with_series_identity(&root).unwrap()
}

fn grouped_rate_space() -> asap_logical_optimizer::CandidateLogicalASAPDAGs<&'static str> {
search_workload(vec![("grouped-rate", grouped_rate_root())])
}

fn grouped_rate() -> PhysicalASAPDAG {
let space = grouped_rate_space();
let selected = global_selection(&space, &DefaultCostModel)
.assemble_selected_query(&space.roots[0].1)
.unwrap()
.unwrap();
compile_physical_asap_dag(&selected).unwrap()
compile_physical_asap_dag(&selected_dag(grouped_rate_root(), AccuracyTarget::Exact)).unwrap()
}
fn run(plan: &CompiledPhysicalDAG, inputs: BTreeMap<u64, Batch>, scope: Scope) -> Vec<Batch> {
let sources = inputs
Expand Down
56 changes: 28 additions & 28 deletions crates/executor/tests/promql_fallback.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ use asap_executor::{
runtime::{Limits, RunContext, Scope},
values::{Batch, Value},
};
use common::compile_physical_asap_dag;
use common::{compile_physical_asap_dag, selected_dag};
use futures::{executor::block_on, StreamExt};
use planner_types::ir::export::PhysicalASAPDAG;
use planner_types::physical::execution_data_state::lift_plain;
Expand Down Expand Up @@ -1225,42 +1225,42 @@ fn histogram_quantile_rejects_equal_output_label_sets() {
assert!(error.contains("same labelset"), "{error}");
}

// Candidate search keeps a classic histogram_quantile whole and exact, even
// for an approximate target, and the selected DAG compiles and executes.
// The stage pipeline keeps a classic histogram_quantile exact, with no sketch,
// even for an approximate target. Over raw buckets it keeps the whole
// expression, which compiles and executes; under a `sum by (le)` it may keep
// the sum as an exact accumulator, which this whole-expression harness does
// not bind.
#[test]
fn histogram_quantile_selection_keeps_the_exact_fallback() {
use asap_logical_optimizer::{
accuracy::DefaultAccuracyModel, default_strategies, search_workload_with_targets,
Replacement,
};
use asap_plan_selection::cost::cost_model::DefaultCostModel;
use planner_types::ir::{schema::FieldDataType, ASAPOp, OperatorNode};
let samples = buckets(&[("job=a", HISTOGRAM)]);
for target in [AccuracyTarget::Exact, AccuracyTarget::Epsilon(0.01)] {
for query in [
"histogram_quantile(0.5, x_bucket)",
"histogram_quantile(0.5, sum by (le, job) (x_bucket))",
for (query, whole) in [
("histogram_quantile(0.5, x_bucket)", true),
(
"histogram_quantile(0.5, sum by (le, job) (x_bucket))",
false,
),
] {
let root =
promql_rows::with_series_identity(&parse_with(query, target.clone())).unwrap();
let space = search_workload_with_targets(
vec![(query, root.clone(), Some(target.clone()))],
&default_strategies(),
&DefaultAccuracyModel,
);
let planned = &space.roots[0].1;
let candidates = &space.candidates_for_target(planned).unwrap().candidates;
let selected = selected_dag(root.clone(), target.clone());
assert!(
candidates.iter().all(|c| matches!(&c.replacement,
Replacement::SubDAG(node) if !node.contains_asap() && node.operator == root.operator)),
"{query}: {candidates:?}"
!OperatorNode::reachable(&selected)
.iter()
.any(|node| matches!(
&node.operator,
planner_types::ir::Operator::ASAP(ASAPOp::SummaryAgg {
family: FieldDataType::Sketch(..),
..
})
)),
"{query}: {selected:?}"
);
let selected = asap_plan_selection::candidate_selection::global_selection(
&space,
&DefaultCostModel,
)
.assemble_selected_dag(planned)
.unwrap()
.unwrap();
if !whole {
continue;
}
assert!(!selected.contains_asap(), "{query}: {selected:?}");
let dag = compile_physical_asap_dag(&selected).unwrap();
let rows = evaluate_dag(&root, &dag, &[("x_bucket", &samples)], 60).unwrap();
let values: Vec<_> = rows.iter().map(|(_, _, v)| *v).collect();
Expand Down
13 changes: 4 additions & 9 deletions crates/executor/tests/tumbling_pane_merge.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
//! Tumbling panes merged by `SummaryMerge` answer the same per-series query as
//! one build over the whole window (#509 Example 3/4, Pattern B; #580).
mod common;
use asap_executor::{
operators::Operator as PhysicalOperator,
physical_planner::{
Expand All @@ -8,9 +9,7 @@ use asap_executor::{
runtime::{Limits, RunContext, Scope},
values::{Batch, Value},
};
use asap_logical_optimizer::search_workload;
use asap_plan_selection::candidate_selection::global_selection;
use asap_plan_selection::cost::cost_model::DefaultCostModel;
use common::selected_dag;
use futures::{executor::block_on, StreamExt};
use planner_types::ir::export::PhysicalASAPDAG;
use planner_types::ir::operator::operator_properties::TimeShift;
Expand All @@ -33,7 +32,7 @@ fn single_build(query: &str, accuracy: AccuracyTarget) -> Rc<OperatorNode> {
query_batch: Some(vec![BatchEntry {
query: Query(query.into()),
requirements: QueryRequirements {
accuracy: AccuracyRequirement::Explicit(accuracy),
accuracy: AccuracyRequirement::Explicit(accuracy.clone()),
..Default::default()
},
predictability: Predictability::Unknown,
Expand All @@ -55,11 +54,7 @@ fn single_build(query: &str, accuracy: AccuracyTarget) -> Rc<OperatorNode> {
.unwrap()
.remove(0);
let root = asap_executor::physical_planner::promql_rows::with_series_identity(&root).unwrap();
let space = search_workload(vec![("q", root)]);
global_selection(&space, &DefaultCostModel)
.assemble_selected_query(&space.roots[0].1)
.unwrap()
.unwrap()
selected_dag(root, accuracy)
}

/// Rewrite the root's per-entity `SummaryAgg` over `TimeRange(5m)` into the
Expand Down
Loading