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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions crates/devtools/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,9 @@ asap-frontend-promql = { path = "../frontend-promql" }
asap-frontend-sql = { path = "../frontend-sql" }
asap-plan-selection = { path = "../plan-selection" }
asap-logical-optimizer = { path = "../logical-optimizer" }
# The reference executor's capabilities, the deployment input stage_pipeline
# plans with. A devtool, not a stage crate, so the #572 guards allow it.
asap-executor = { path = "../executor" }

# Used by the show_ir / dag_export / variant_coverage bins (catalog schemas,
# async SQL path, JSON output) and by the topk_ir / canonical_examples
Expand Down
9 changes: 7 additions & 2 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,11 @@
// e.g. "· ingestion time: Kll ×5 panes"), no cost;
// - stage3_selection: per-candidate costs, the selected candidate, and
// every other candidate as rejected (`valid: false`: inaccurate, over a
// latency bound, or could not be built) or costlier.
// latency bound, needing a capability the deployment lacks, or could not
// be built) or costlier.
//
// The deployment inputs are the built-in cost and accuracy models and the
// reference executor's capabilities (`asap_executor::capabilities`).
//
// Everything is the library's `plan_selection::plan_stages`, the function the
// facade runs; this tool only serializes it. Stage 3 here is over every
Expand Down Expand Up @@ -133,11 +137,12 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result<
.map(|entry| RootDemand::from(&entry))
.collect();
let data = workload.data_workload.clone().unwrap_or_default();
let capabilities = asap_executor::capabilities();
let run = plan_stages(
roots.into_iter().enumerate().collect(),
&demand,
&data,
PlanningModels::builtin(),
PlanningModels::builtin().with_capabilities(&capabilities),
max_candidates.max(1),
)
.map_err(|e| format!("planning: {e}"))?;
Expand Down
202 changes: 202 additions & 0 deletions crates/executor/src/capability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
//!
//! Stored-state encodings belong to deployments. Full plan acceptance is
//! owned by `binding`, which also validates schemas, expressions and inputs.
//!
//! [`capabilities`] states these rules as the planner's deployment input.
use crate::Error;
use planner_types::ir::schema::{
ExactKind, ExactParams, FieldDataType as SummaryFamilyType, GroupingStrategy, SketchAlgorithm,
Expand Down Expand Up @@ -341,3 +343,203 @@ pub fn validate_exact_evaluation(
}
Ok(())
}

/// This executor's capabilities, as the planner's deployment input
/// ([`DeploymentCapabilities`]): every summary family and layout whose native
/// state this module accepts, with the readouts it evaluates. Parameters are
/// representative; sizing limits are checked when a plan is bound. The
/// executor maintains state at ingestion time, keeps no query-time result
/// across evaluations, and reads raw data from its inputs, so raw retention
/// is the deployment's and is not priced.
pub fn capabilities() -> planner_types::deployment::DeploymentCapabilities {
use planner_types::deployment::{summary_of, DeploymentCapabilities, Readout, SummarySupport};
use planner_types::ir::scalar::ColumnRef;
use planner_types::ir::schema::{default_hydra_params, HydraKind, SketchKind};
let exact = [
(ExactKind::Sum, ExactParams::Sum),
(ExactKind::Count, ExactParams::Count),
(ExactKind::Min, ExactParams::Min),
(ExactKind::Max, ExactParams::Max),
(ExactKind::Increase, ExactParams::Increase),
(ExactKind::Rate, ExactParams::Rate),
(ExactKind::IRate, ExactParams::IRate),
]
.map(|(kind, params)| SummaryFamilyType::ExactAggregate(kind, params));
let (width, depth, heap_size) = (1024, 5, 16);
let sketches = [
(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }),
(
SketchAlgorithm::DDSketch,
SketchParams::DDSketch { alpha: 0.01 },
),
(SketchAlgorithm::Hll, SketchParams::Hll { precision: 14 }),
(SketchAlgorithm::Cms, SketchParams::Cms { width, depth }),
(
SketchAlgorithm::CountSketch,
SketchParams::CountSketch { width, depth },
),
(
SketchAlgorithm::CmsWithHeap,
SketchParams::CmsWithHeap {
width,
depth,
heap_size,
},
),
(
SketchAlgorithm::CountSketchWithHeap,
SketchParams::CountSketchWithHeap {
width,
depth,
heap_size,
},
),
(
SketchAlgorithm::UnivMon,
SketchParams::UnivMon {
heap_size,
sketch_rows: depth,
sketch_cols: width,
layers: 8,
},
),
(SketchAlgorithm::Kmv, SketchParams::Kmv { k: 1024 }),
(SketchAlgorithm::Theta, SketchParams::Theta { k: 1024 }),
];
let hydra = [
HydraKind::HydraCms,
HydraKind::HydraCountSketch,
HydraKind::HydraKll,
];
let sketches = sketches.into_iter().flat_map(|(algorithm, params)| {
let layouts = std::iter::once(GroupingStrategy::PerSubpopulationInstance).chain(
hydra.iter().filter_map(|kind| {
default_hydra_params(kind.clone(), &params).map(|params| {
GroupingStrategy::SharedMultiSubpopulation {
kind: kind.clone(),
params,
}
})
}),
);
let kind = SketchKind::new(algorithm, params.clone());
layouts
.map(|layout| SummaryFamilyType::Sketch(kind.clone(), layout))
.collect::<Vec<_>>()
});
let item = |value: Option<&str>| SketchStatistic::PointCount {
key: match value {
None => ColumnRef::SampleValue,
Some(_) => ColumnRef::Named("item".into()),
},
value: value.map(Into::into),
};
let readouts = [
SketchStatistic::Quantile { q: 0.5 },
item(None),
item(Some("item")),
SketchStatistic::Cardinality,
SketchStatistic::TopK { k: 1 },
SketchStatistic::FrequencyL2,
SketchStatistic::FrequencyEntropy,
];
let summaries = exact
.into_iter()
.chain(sketches)
.filter(builds_from_rows)
.map(|family| {
let (summary, layout) = summary_of(&family).expect("a summary family");
let readouts = match &family {
SummaryFamilyType::Sketch(kind, _) => readouts
.iter()
.filter(|statistic| match statistic {
// Read by the keyed evaluation operator.
SketchStatistic::TopK { .. } => {
crate::summary_kernels::weighted_frequency::WeightedFrequency::configuration(kind)
.is_ok()
}
_ => validate_sketch_evaluation(&family, statistic).is_ok(),
})
.map(Readout::of)
.collect(),
_ => Default::default(),
};
SummarySupport {
family: summary,
layout,
readouts,
}
})
.collect();
DeploymentCapabilities {
summaries: Some(summaries),
ingestion_time: true,
query_time_retention: false,
..DeploymentCapabilities::UNRESTRICTED
}
}

/// Whether a plan can build `family` from rows: its native state is
/// accepted, except a per-group Count-Min, which is native as stored state
/// only (see [`validate_native_family`]).
fn builds_from_rows(family: &SummaryFamilyType) -> bool {
let per_group_cms = matches!(family, SummaryFamilyType::Sketch(kind, grouping)
if kind.algorithm() == &SketchAlgorithm::Cms
&& grouping == &GroupingStrategy::PerSubpopulationInstance);
!per_group_cms && validate_native_family(family).is_ok()
}

#[cfg(test)]
mod capabilities_tests {
use super::*;
use planner_types::deployment::{InstanceLayout, Readout, SummaryFamily};
use planner_types::ir::schema::HydraKind;

fn readouts(
caps: &planner_types::deployment::DeploymentCapabilities,
family: SummaryFamily,
layout: InstanceLayout,
) -> Option<Vec<Readout>> {
caps.summaries
.as_ref()
.unwrap()
.iter()
.find(|s| s.family == family && s.layout == layout)
.map(|s| s.readouts.iter().copied().collect())
}

/// The exported set follows this module's rules: e.g. KLL quantiles,
/// heap top-k, HydraCms counts, exact accumulators without IRate, and
/// no per-group Count-Min build.
#[test]
fn exported_capabilities_follow_the_validators() {
use InstanceLayout::{Hydra, PerGroup};
use SummaryFamily::{Exact, Sketch};
let caps = capabilities();
assert!(caps.ingestion_time && !caps.query_time_retention);
assert!(caps.raw_data_retained && caps.memory_budget_bytes.is_none());
let of = |family, layout| readouts(&caps, family, layout);
assert_eq!(
of(Sketch(SketchAlgorithm::Kll), PerGroup),
Some(vec![Readout::Quantile])
);
for heap in [
SketchAlgorithm::CmsWithHeap,
SketchAlgorithm::CountSketchWithHeap,
] {
assert_eq!(of(Sketch(heap), PerGroup), Some(vec![Readout::TopK]));
}
assert_eq!(
of(Sketch(SketchAlgorithm::Cms), Hydra(HydraKind::HydraCms)),
Some(vec![Readout::TotalCount, Readout::ItemCount])
);
assert_eq!(of(Sketch(SketchAlgorithm::Cms), PerGroup), None);
assert_eq!(
of(Sketch(SketchAlgorithm::Kll), Hydra(HydraKind::HydraKll)),
None
);
assert_eq!(of(Sketch(SketchAlgorithm::CountSketch), PerGroup), None);
assert_eq!(of(Exact(ExactKind::Rate), PerGroup), Some(vec![]));
assert_eq!(of(Exact(ExactKind::IRate), PerGroup), None);
}
}
1 change: 1 addition & 0 deletions crates/executor/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ pub use traits::*;

pub use expressions::arithmetic;
pub mod capability;
pub use capability::capabilities;
pub use summary_kernels::factory;

/// The exact Planner contract used by these kernels.
Expand Down
12 changes: 12 additions & 0 deletions crates/integration-tests/tests/executor_models/mod.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
//! The deployment input these tests plan with: the built-in models and the
//! reference executor's capabilities, as a deployment running the plans on
//! that executor would pass them (C2).
use std::sync::LazyLock;

use asap_plan_selection::{DeploymentCapabilities, PlanningModels};

static EXECUTOR: LazyLock<DeploymentCapabilities> = LazyLock::new(asap_executor::capabilities);

pub fn executor_models() -> PlanningModels<'static> {
PlanningModels::builtin().with_capabilities(&EXECUTOR)
}
6 changes: 4 additions & 2 deletions crates/integration-tests/tests/filtered_aggregates.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
//! Filtered aggregates (`FILTER (WHERE …)`, a per-measure filter, or a
//! filtered `SummaryAgg`) execute with SQL semantics: a row that fails the
//! filter contributes nothing, but its group is kept.
mod executor_models;
mod physical_common;
use std::collections::BTreeMap;
use std::rc::Rc;
Expand All @@ -10,14 +11,15 @@ use asap_frontend_sql::{lower_sql, SqlCatalog};
use asap_logical_optimizer::pass1::logical_candidates::{
compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates,
};
use asap_plan_selection::{plan_stages, PlanningModels};
use asap_plan_selection::plan_stages;
use asap_types::ir::schema::{DataType, Field, Schema};
use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, QueryRoot};
use asap_types::types::AccuracyTarget;
use asap_types::workload::{
DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate,
RootDemand,
};
use executor_models::executor_models;

fn catalog() -> SqlCatalog {
SqlCatalog::new().with_table(
Expand Down Expand Up @@ -112,7 +114,7 @@ async fn selected(sql: &str) -> Rc<OperatorNode> {
vec![(0, QueryRoot::Operator(root))],
&demand,
&data,
PlanningModels::builtin(),
executor_models(),
4096,
)
.unwrap();
Expand Down
5 changes: 3 additions & 2 deletions crates/integration-tests/tests/operator_design_examples.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ use asap_types::ir::schema::{
use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, ScalarExpr};
use asap_types::types::AccuracyTarget;
use std::rc::Rc;
mod executor_models;
mod physical_common;

fn catalog() -> SqlCatalog {
Expand Down Expand Up @@ -274,12 +275,12 @@ async fn sql_window_and_filtered_aggregate_types() {
/// constructed by the test.
#[tokio::test]
async fn batch_planning_selects_and_executes_each_plan() {
use crate::executor_models::executor_models;
use asap_executor::{
physical_planner::{compile, InputContract},
runtime::Scope,
values::{Batch, Value},
};
use asap_plan_selection::PlanningModels;
use asap_planner::{e2e_plan, FrontendInput, UserInput};
use asap_types::workload::*;
use std::{collections::BTreeMap, sync::Arc};
Expand Down Expand Up @@ -320,7 +321,7 @@ async fn batch_planning_selects_and_executes_each_plan() {
let output = e2e_plan(UserInput::new(
&workload,
FrontendInput::Sql { catalog: &catalog },
PlanningModels::builtin(),
executor_models(),
))
.await
.unwrap();
Expand Down
8 changes: 5 additions & 3 deletions crates/integration-tests/tests/pass1_sql_coverage.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
//! Pass 1 alternatives over SQL row sources compose, compile and execute.
mod executor_models;
mod physical_common;
use std::collections::BTreeMap;
use std::rc::Rc;
Expand All @@ -9,14 +10,15 @@ use asap_logical_optimizer::pass1::logical_candidates::{
compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates,
LocalLogicalCandidates,
};
use asap_plan_selection::{plan_stages, PlanningModels};
use asap_plan_selection::plan_stages;
use asap_types::ir::schema::{DataType, Field, Schema};
use asap_types::ir::{ASAPOp, Operator, OperatorNode, QueryRoot};
use asap_types::types::AccuracyTarget;
use asap_types::workload::{
DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate,
RootDemand,
};
use executor_models::executor_models;

fn catalog() -> SqlCatalog {
SqlCatalog::new().with_table(
Expand Down Expand Up @@ -134,7 +136,7 @@ async fn example2_design_candidates_all_build() {
input_cardinality: declared(10_000_000),
..Default::default()
};
let run = plan_stages(roots, &demand, &data, PlanningModels::builtin(), 4096).unwrap();
let run = plan_stages(roots, &demand, &data, executor_models(), 4096).unwrap();
let enumeration = run.enumeration.unwrap();
let unbuilt: Vec<_> = enumeration
.selection
Expand Down Expand Up @@ -203,7 +205,7 @@ async fn grouped_count_offers_a_priced_executable_hydra_plan() {
vec![(0, QueryRoot::Operator(root))],
&demand,
&data,
PlanningModels::builtin(),
executor_models(),
4096,
)
.unwrap();
Expand Down
Loading