diff --git a/Cargo.lock b/Cargo.lock index 6b00e464..642b80de 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -317,6 +317,7 @@ dependencies = [ name = "asap-devtools" version = "0.1.0" dependencies = [ + "asap-executor", "asap-frontend-promql", "asap-frontend-sql", "asap-logical-optimizer", diff --git a/crates/devtools/Cargo.toml b/crates/devtools/Cargo.toml index 484ce799..edf20aef 100644 --- a/crates/devtools/Cargo.toml +++ b/crates/devtools/Cargo.toml @@ -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 diff --git a/crates/devtools/src/bin/stage_pipeline.rs b/crates/devtools/src/bin/stage_pipeline.rs index 6e305ccf..be6c47a0 100644 --- a/crates/devtools/src/bin/stage_pipeline.rs +++ b/crates/devtools/src/bin/stage_pipeline.rs @@ -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 @@ -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}"))?; diff --git a/crates/executor/src/capability.rs b/crates/executor/src/capability.rs index 89975bcc..b6e27ac3 100644 --- a/crates/executor/src/capability.rs +++ b/crates/executor/src/capability.rs @@ -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, @@ -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(), ¶ms).map(|params| { + GroupingStrategy::SharedMultiSubpopulation { + kind: kind.clone(), + params, + } + }) + }), + ); + let kind = SketchKind::new(algorithm, params.clone()); + layouts + .map(|layout| SummaryFamilyType::Sketch(kind.clone(), layout)) + .collect::>() + }); + 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> { + 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); + } +} diff --git a/crates/executor/src/lib.rs b/crates/executor/src/lib.rs index ea6b73ea..edfb1439 100644 --- a/crates/executor/src/lib.rs +++ b/crates/executor/src/lib.rs @@ -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. diff --git a/crates/integration-tests/tests/executor_models/mod.rs b/crates/integration-tests/tests/executor_models/mod.rs new file mode 100644 index 00000000..25b92e9f --- /dev/null +++ b/crates/integration-tests/tests/executor_models/mod.rs @@ -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 = LazyLock::new(asap_executor::capabilities); + +pub fn executor_models() -> PlanningModels<'static> { + PlanningModels::builtin().with_capabilities(&EXECUTOR) +} diff --git a/crates/integration-tests/tests/filtered_aggregates.rs b/crates/integration-tests/tests/filtered_aggregates.rs index 248b5bbd..e6bf8bc7 100644 --- a/crates/integration-tests/tests/filtered_aggregates.rs +++ b/crates/integration-tests/tests/filtered_aggregates.rs @@ -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; @@ -10,7 +11,7 @@ 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; @@ -18,6 +19,7 @@ use asap_types::workload::{ DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate, RootDemand, }; +use executor_models::executor_models; fn catalog() -> SqlCatalog { SqlCatalog::new().with_table( @@ -112,7 +114,7 @@ async fn selected(sql: &str) -> Rc { vec![(0, QueryRoot::Operator(root))], &demand, &data, - PlanningModels::builtin(), + executor_models(), 4096, ) .unwrap(); diff --git a/crates/integration-tests/tests/operator_design_examples.rs b/crates/integration-tests/tests/operator_design_examples.rs index 755dc61f..40001519 100644 --- a/crates/integration-tests/tests/operator_design_examples.rs +++ b/crates/integration-tests/tests/operator_design_examples.rs @@ -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 { @@ -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}; @@ -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(); diff --git a/crates/integration-tests/tests/pass1_sql_coverage.rs b/crates/integration-tests/tests/pass1_sql_coverage.rs index 2d52ef9a..0b0daf97 100644 --- a/crates/integration-tests/tests/pass1_sql_coverage.rs +++ b/crates/integration-tests/tests/pass1_sql_coverage.rs @@ -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; @@ -9,7 +10,7 @@ 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; @@ -17,6 +18,7 @@ use asap_types::workload::{ DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate, RootDemand, }; +use executor_models::executor_models; fn catalog() -> SqlCatalog { SqlCatalog::new().with_table( @@ -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 @@ -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(); diff --git a/crates/integration-tests/tests/planner_layering_common/mod.rs b/crates/integration-tests/tests/planner_layering_common/mod.rs index f51bd00e..c6de59ff 100644 --- a/crates/integration-tests/tests/planner_layering_common/mod.rs +++ b/crates/integration-tests/tests/planner_layering_common/mod.rs @@ -13,7 +13,10 @@ use std::collections::{BTreeMap, BTreeSet, HashSet}; use asap_physical_optimizer::implementation::physical_candidates::PhysicalCandidate; -use asap_plan_selection::{plan_stages, PlanningModels, Selection}; +#[path = "../executor_models/mod.rs"] +mod executor_models; + +use asap_plan_selection::{plan_stages, Selection}; use asap_types::ir::export::{ compile_logical_asap_workload, LogicalASAPDAG, LogicalASAPNodeId, LogicalASAPOperatorPayload, LogicalASAPQueryRoot, NonASAPOpKind, PhysicalASAPDAG, @@ -28,6 +31,7 @@ use asap_types::workload::{ QueryLanguage, QueryRecurrence, QueryRequirements, QueryTimeScope, QueryWorkload, Rate, RepeatedDemand, RepeatingEntry, RepetitionInterval, RootDemand, TimeSelection, TimestampMs, }; +use executor_models::executor_models; /// Every enumerated candidate is built and displayed; the largest example /// (Example 3, Pattern A) has 486 today. @@ -341,7 +345,7 @@ pub fn run_stages(workload: &PlanningWorkload, roots: Vec) -> Run { roots.into_iter().enumerate().collect(), &demand, workload.data_workload.as_ref().expect("data workload"), - PlanningModels::builtin(), + executor_models(), MAX_CANDIDATES, ) .expect("plans"); diff --git a/crates/integration-tests/tests/planner_layering_example1.rs b/crates/integration-tests/tests/planner_layering_example1.rs index e9bc9672..c1ebb211 100644 --- a/crates/integration-tests/tests/planner_layering_example1.rs +++ b/crates/integration-tests/tests/planner_layering_example1.rs @@ -19,6 +19,9 @@ use std::collections::{BTreeMap, BTreeSet, HashSet}; +mod executor_models; +use executor_models::executor_models; + use asap_plan_selection::PlanningModels; use asap_types::ir::export::{LogicalASAPDAG, LogicalASAPNodeId, LogicalASAPOperatorPayload}; use asap_types::ir::schema::state_type::{GroupingStrategy, HydraKind, SketchAlgorithm}; @@ -160,7 +163,7 @@ mod stages { lower(workload).into_iter().enumerate().collect(), &demand, workload.data_workload.as_ref().expect("data workload"), - PlanningModels::builtin(), + executor_models(), DISPLAYED, ) .expect("plans"); @@ -1014,7 +1017,7 @@ fn compile_in_runtime(p: &PhysicalCandidate) -> Result<(), String> { /// candidates invalid for their weights. fn assert_runtime_agrees_with_stage3(workload: PlanningWorkload) -> usize { let (workload, _, physical) = pipeline_for(workload); - let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selection = stage3_select(&workload, &physical, executor_models()); let invalid: BTreeMap<_, _> = selection .rejected .iter() @@ -1090,7 +1093,7 @@ fn stage2_count_sketch_heap_topk_compiles_in_the_physical_planner() { #[test] fn stage3_selected_plan_compiles_in_the_physical_planner() { let (workload, _, physical) = pipeline(); - let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selection = stage3_select(&workload, &physical, executor_models()); let selected = physical .iter() .find(|p| p.id == selection.selected) @@ -1098,13 +1101,43 @@ fn stage3_selected_plan_compiles_in_the_physical_planner() { compile_in_runtime(selected).unwrap(); } +/// Deployment inputs (C2, added by the implementer): the reference +/// executor's exported capabilities reject no candidate its physical planner +/// compiles, and planning with them selects what the unrestricted default +/// does, at the same costs (Example 1 needs nothing the executor lacks, and +/// both keep raw data, so no raw retention is priced). +#[test] +fn executor_capabilities_accept_every_compiled_candidate_and_keep_the_selection() { + let (workload, _, physical) = pipeline(); + let executor = stage3_select(&workload, &physical, executor_models()); + let default = stage3_select(&workload, &physical, PlanningModels::builtin()); + let mut compiled = 0; + for p in &physical { + if compile_in_runtime(p).is_ok() { + compiled += 1; + if let Some(r) = executor.rejected.iter().find(|r| r.id == p.id) { + assert!(!r.reason.contains("deployment"), "{}: {}", p.id, r.reason); + } + } + } + assert_eq!(compiled, physical.len()); + assert_eq!(executor.selected, default.selected); + let totals = |s: &Selection| -> BTreeMap { + s.costs + .iter() + .map(|(id, c)| (id.clone(), c.total)) + .collect() + }; + assert_eq!(totals(&executor), totals(&default)); +} + // ── Stage 3 ────────────────────────────────────────────────────────────── /// Stage 3 selects one candidate and gives every other one a reason. #[test] fn stage3_selects_one_and_explains_the_rest() { let (workload, _, physical) = pipeline(); - let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selection = stage3_select(&workload, &physical, executor_models()); let all: BTreeSet<_> = physical.iter().map(|p| p.id.clone()).collect(); let mut accounted: BTreeSet<_> = selection.rejected.iter().map(|r| r.id.clone()).collect(); assert_eq!( @@ -1124,7 +1157,7 @@ fn stage3_selects_one_and_explains_the_rest() { #[test] fn stage3_selects_cheapest_valid() { let (workload, _, physical) = pipeline(); - let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selection = stage3_select(&workload, &physical, executor_models()); let invalid: BTreeSet<_> = selection .rejected .iter() @@ -1144,7 +1177,7 @@ fn stage3_selects_cheapest_valid() { #[test] fn stage3_rejects_candidates_over_the_latency_bound() { let (workload, _, physical) = pipeline(); - let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selection = stage3_select(&workload, &physical, executor_models()); let over: BTreeSet<_> = selection .rejected .iter() @@ -1176,7 +1209,7 @@ fn stage3_rejects_candidates_over_the_latency_bound() { #[test] fn stage3_per_second_cost_keeps_the_ranking() { let (workload, _, physical) = pipeline(); - let per_second = stage3_select(&workload, &physical, PlanningModels::builtin()); + let per_second = stage3_select(&workload, &physical, executor_models()); // Evaluated once per second, a candidate's cost is its per-evaluation cost. let mut every_second = workload.clone(); for entry in every_second @@ -1187,7 +1220,7 @@ fn stage3_per_second_cost_keeps_the_ranking() { { entry.demand = RepeatedDemand::FixedInterval(RepetitionInterval(1_000)); } - let per_evaluation = stage3_select(&every_second, &physical, PlanningModels::builtin()); + let per_evaluation = stage3_select(&every_second, &physical, executor_models()); assert_eq!(per_second.selected, "P82"); assert_eq!(per_evaluation.selected, "P82"); for (id, cost) in per_second.costs.iter().filter(|(id, _)| !id.contains("-m")) { @@ -1203,7 +1236,7 @@ fn stage3_per_second_cost_keeps_the_ranking() { #[test] fn stage3_charges_each_node_once() { let (workload, _, physical) = pipeline(); - let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selection = stage3_select(&workload, &physical, executor_models()); let invalid: BTreeSet<_> = selection .rejected .iter() @@ -1238,7 +1271,7 @@ fn stage3_charges_each_node_once() { #[test] fn stage3_shared_input_is_not_costlier() { let (workload, _, physical) = pipeline(); - let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selection = stage3_select(&workload, &physical, executor_models()); // All query time: a maintained-pane candidate cannot share the scan. let by_combo: BTreeMap<_, _> = physical .iter() @@ -1278,7 +1311,7 @@ fn stage3_shared_input_is_not_costlier() { #[test] fn stage3_selects_a_shared_input_plan() { let (workload, _, physical) = pipeline(); - let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selection = stage3_select(&workload, &physical, executor_models()); let selected = physical .iter() .find(|p| p.id == selection.selected) diff --git a/crates/integration-tests/tests/planner_layering_example2.rs b/crates/integration-tests/tests/planner_layering_example2.rs index b90d69ad..12bc5389 100644 --- a/crates/integration-tests/tests/planner_layering_example2.rs +++ b/crates/integration-tests/tests/planner_layering_example2.rs @@ -1,10 +1,11 @@ //! #509 Example 2 status: one UnivMon over `src_ip` for distinct, entropy and //! L2. Records what lowers, plans through `plan_stages` and executes exactly. +mod executor_models; mod physical_common; use asap_executor::values::Value; use asap_frontend_sql::{lower_sql, SqlCatalog}; use asap_logical_optimizer::pass1::replacement::Realization; -use asap_plan_selection::{plan_stages, PlanningModels}; +use asap_plan_selection::plan_stages; use asap_types::ir::operator::AggIntent; use asap_types::ir::schema::{DataType, Field, Schema, SketchAlgorithm}; use asap_types::ir::{NonASAPOp, OperatorNode, QueryRoot, ScalarExpr}; @@ -13,6 +14,7 @@ use asap_types::workload::{ DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate, RootDemand, }; +use executor_models::executor_models; use std::rc::Rc; const WINDOW: &str = "ts >= now() - INTERVAL '1 minute'"; @@ -138,7 +140,7 @@ async fn example2_plans_through_the_stage_pipeline() { .collect(), &demand, &data, - PlanningModels::builtin(), + executor_models(), 0, ) .expect("Example 2 plans"); diff --git a/crates/integration-tests/tests/planner_layering_example4.rs b/crates/integration-tests/tests/planner_layering_example4.rs index ae1d70cb..450195e6 100644 --- a/crates/integration-tests/tests/planner_layering_example4.rs +++ b/crates/integration-tests/tests/planner_layering_example4.rs @@ -328,7 +328,7 @@ fn stage2_b_sliding_kll_has_two_materialization_options() { /// B2 rebuilds all five tumbling KLLs at every evaluation, so it costs at least B1 and B3. #[test] -#[ignore = "built-in model (#604): B2 rebuilds the panes in 310 ms, over the 200 ms latency bound, so it is not priced; and B1 retains 6 panes of 1M per-series KLLs (6.1 GB, 768 cost/s) against B2's 5.17 cost/s of rebuilds"] +#[ignore = "built-in model (#604, Q49): B2 rebuilds the panes in 310 ms, over the 200 ms latency bound, so it is not priced; B1 retains 6 panes of 1M per-series KLLs (6.1 GB) for 768.58 cost/s, and B2 costs 5.17 cost/s, or 45.17 when the deployment does not keep raw data (5 min of raw samples, 320 MB, 40.0 cost/s)"] fn stage3_b_rebuilding_every_window_costs_most() { let run = run_promql(&pattern_b()); let options = options_of(&run, tumbling(), 1); @@ -340,7 +340,7 @@ fn stage3_b_rebuilding_every_window_costs_most() { /// Repeating over arriving data, the built-in models pick B1 among the tumbling options. #[test] -#[ignore = "built-in model (#604): B2 is over the 200 ms latency bound, so it is not priced; and B1 retains 6 panes of 1M per-series KLLs (6.1 GB, 768 cost/s) against B2's 5.17 cost/s"] +#[ignore = "built-in model (#604, Q49): B2 is over the 200 ms latency bound, so it is not priced; B1 costs 768.58 cost/s against B2's 5.17, or 45.17 when the deployment does not keep raw data"] fn stage3_b_prefers_ingestion_time_tumbling_windows() { let run = run_promql(&pattern_b()); let options = options_of(&run, tumbling(), 1); diff --git a/crates/physical-optimizer/src/materialization.rs b/crates/physical-optimizer/src/materialization.rs index 4a6f2ded..32b78f6e 100644 --- a/crates/physical-optimizer/src/materialization.rs +++ b/crates/physical-optimizer/src/materialization.rs @@ -33,7 +33,7 @@ use std::collections::{BTreeSet, HashMap, HashSet}; use std::rc::Rc; use asap_types::ir::properties::ExecutionTiming; -use asap_types::ir::schema::FieldDataType; +use asap_types::ir::schema::{FieldDataType, GroupingStrategy}; use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}; use asap_types::workload::{ DataArrival, DataWorkload, Predictability, QueryRecurrence, RepeatedDemand, RootDemand, @@ -353,6 +353,10 @@ fn window_of(summary: &Rc) -> Window { fn family_name(summary: &OperatorNode) -> String { match &summary.operator { Operator::ASAP(ASAPOp::SummaryAgg { family, .. }) => match family { + // A Hydra's family is the per-group sketch it emulates. + FieldDataType::Sketch(_, GroupingStrategy::SharedMultiSubpopulation { kind, .. }) => { + format!("{kind:?}") + } FieldDataType::Sketch(kind, _) => format!("{:?}", kind.algorithm()), FieldDataType::ExactAggregate(kind, _) => format!("exact {kind:?}"), other => format!("{other:?}"), @@ -481,6 +485,47 @@ mod tests { const P99_5M: &str = "quantile_over_time(0.99, m[5m])"; + /// A Hydra summary's unit is labeled by its Hydra kind ("HydraCms"), + /// not by the per-group Count-Min it emulates. + #[test] + fn hydra_unit_is_labeled_by_its_hydra_kind() { + let demand = [every_at(300_000)]; + let approximate = AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }; + let root = crate::test_support::lower_promql("count by (job) (m)", approximate); + let root = asap_types::ir::schema_support::with_promql_series_identity(&root).unwrap(); + let variants = stage1_logical_candidates( + vec![(0, QueryRoot::Operator(root))], + &Default::default(), + &demand, + ) + .unwrap(); + let inventory = &variants[0].inventory; + let choice: Vec<_> = inventory + .targets + .iter() + .map(|t| { + t.groupings + .iter() + .position(|g| *g != GroupingStrategy::default()) + .unwrap_or(0) + }) + .collect(); + let roots: Vec<_> = compose_logical_candidate(inventory, &choice) + .unwrap() + .into_iter() + .map(|(_, root)| match root { + QueryRoot::Operator(node) => node, + QueryRoot::Scalar(_) => panic!("operator root"), + }) + .collect(); + let space = MaterializationSpace::new(&roots, &demand, &ingesting()); + let labels: Vec<_> = space.units.iter().map(|u| u.label.as_str()).collect(); + assert_eq!(labels, ["HydraCms"]); + } + /// Panes of a query repeating every minute over arriving data give a /// second candidate that builds all five panes at ingestion time and /// merges them at query time (Example 4, B1). diff --git a/crates/plan-selection/src/lib.rs b/crates/plan-selection/src/lib.rs index eda66bd8..da99b22b 100644 --- a/crates/plan-selection/src/lib.rs +++ b/crates/plan-selection/src/lib.rs @@ -25,8 +25,11 @@ //! count, weighted by [`Stage3Calibration`]; summary build and estimation are //! priced as rows × sketch depth and rows read out. These numbers are //! illustrative, not calibrated. A query's latency bound is checked against -//! the query-time work it waits for in one evaluation; deployment -//! capabilities are not checked yet. +//! the query-time work it waits for in one evaluation. A candidate needing +//! a capability the deployment lacks ([`DeploymentCapabilities`]) is +//! rejected, and so is one retaining more than the deployment's memory +//! budget. When the deployment does not keep raw data anyway, a plan reading +//! raw data at query time pays for retaining it (Q48). //! //! A logical candidate has several physical candidates, one per Stage 2 //! materialization choice; selection takes the cheapest valid one. @@ -40,6 +43,7 @@ pub mod cost; #[cfg(test)] mod test_support; +pub use asap_types::deployment::DeploymentCapabilities; pub use candidate_selection::{ CompositionDecision, CostedGlobalSelection, RankedTargetSubDAGCandidates, RecurrenceProfileMap, }; @@ -110,6 +114,7 @@ pub const MAX_ENUMERATED_CANDIDATES: usize = 64; static DEFAULT_COST_MODEL: DefaultCostModel = DefaultCostModel; static DEFAULT_ACCURACY_MODEL: DefaultAccuracyModel = DefaultAccuracyModel; static NO_ACCURACY_EVIDENCE: NoAccuracyEvidence = NoAccuracyEvidence; +static UNRESTRICTED: DeploymentCapabilities = DeploymentCapabilities::UNRESTRICTED; /// Stage 3's price coefficients and amortization horizon. Values are /// illustrative until calibrated from measurements. @@ -166,7 +171,8 @@ impl Stage3Calibration { } /// Planning logic, as opposed to the scoped facts it consumes: a model can have -/// a built-in default, evidence about a particular deployment cannot. +/// a built-in default, evidence about a particular deployment cannot. With +/// the deployment's capabilities, these are #509's deployment inputs. /// /// Stage 3 prices plans analytically, so the stage pipeline does not read /// `cost`; only the legacy replacement search does (#580). @@ -177,6 +183,9 @@ pub struct PlanningModels<'a> { pub accuracy: &'a dyn AccuracyModel, pub evidence: &'a dyn AccuracyEvidenceProvider, pub calibration: Stage3Calibration, + /// What the deployment can build, read out and keep; unrestricted by + /// default. + pub capabilities: &'a DeploymentCapabilities, } impl<'a> PlanningModels<'a> { @@ -190,6 +199,7 @@ impl<'a> PlanningModels<'a> { accuracy, evidence, calibration: Stage3Calibration::ILLUSTRATIVE, + capabilities: &UNRESTRICTED, } } @@ -202,6 +212,7 @@ impl<'a> PlanningModels<'a> { accuracy: &DEFAULT_ACCURACY_MODEL, evidence: &NO_ACCURACY_EVIDENCE, calibration: Stage3Calibration::ILLUSTRATIVE, + capabilities: &UNRESTRICTED, } } @@ -224,6 +235,11 @@ impl<'a> PlanningModels<'a> { self.calibration = calibration; self } + + pub fn with_capabilities(mut self, capabilities: &'a DeploymentCapabilities) -> Self { + self.capabilities = capabilities; + self + } } #[derive(Debug, Clone, PartialEq)] @@ -362,11 +378,36 @@ fn assess( demand.len() )); } + if let Some(reason) = capability_violation(candidate, models.capabilities) { + return Err(reason); + } if let Some(reason) = accuracy_violation(candidate, demand, models) { return Err(reason); } - let (cost, per_evaluation) = price_nodes(&candidate.dag, demand, data, &models.calibration) - .map_err(|(node, error)| format!("node {node:?}: {error}"))?; + let capabilities = models.capabilities; + let raw_bytes_per_sample = + (!capabilities.raw_data_retained).then_some(capabilities.raw_bytes_per_sample); + let priced = price_nodes( + &candidate.dag, + demand, + data, + &models.calibration, + raw_bytes_per_sample, + ) + .map_err(|(node, error)| format!("node {node:?}: {error}"))?; + let Priced { + cost, + per_evaluation, + retained_bytes, + } = priced; + if let Some(budget) = capabilities.memory_budget_bytes { + if retained_bytes > budget { + return Err(format!( + "retains {retained_bytes} bytes across evaluations, over the deployment's \ + memory budget of {budget} bytes" + )); + } + } if let Some(reason) = latency_violation(&candidate.dag, &per_evaluation, demand, &models.calibration) { @@ -1109,6 +1150,55 @@ pub fn plan_stages( }) } +/// The first capability `candidate` needs that the deployment lacks, as a +/// reason: maintaining state at ingestion time, building a summary, or +/// reading a statistic out of one. +fn capability_violation( + candidate: &PhysicalCandidate, + capabilities: &DeploymentCapabilities, +) -> Option { + if !capabilities.ingestion_time + && candidate + .dag + .nodes + .iter() + .any(|n| !n.output_state.timing.is_query_time()) + { + return Some("deployment cannot maintain state at ingestion time".into()); + } + // Summaries first: a readout of a summary that cannot be built is moot. + let reasons = |node: &Rc, readouts: bool| match &node.operator { + Operator::ASAP(ASAPOp::SummaryAgg { family, .. }) if !readouts => { + capabilities.missing_summary(family) + } + Operator::ASAP(ASAPOp::SummaryEstimate { + summary_input, + query: statistic, + }) if readouts => summary_builds(summary_input) + .into_iter() + .flatten() + .find_map(|build| match &build.operator { + Operator::ASAP(ASAPOp::SummaryAgg { family, .. }) => { + capabilities.missing_readout(family, statistic) + } + _ => None, + }), + _ => None, + }; + [false, true].into_iter().find_map(|readouts| { + candidate + .roots + .iter() + .enumerate() + .find_map(|(query, root)| { + OperatorNode::reachable(root) + .iter() + .find_map(|node| reasons(node, readouts)) + .map(|reason| format!("q{}: {reason}", query + 1)) + }) + }) +} + /// The first summary estimate that misses its query's target, as a reason. fn accuracy_violation( candidate: &PhysicalCandidate, @@ -1194,7 +1284,11 @@ fn build_violation( } fn family_name(family: &FieldDataType) -> String { + use asap_types::ir::schema::GroupingStrategy; match family { + FieldDataType::Sketch(_, GroupingStrategy::SharedMultiSubpopulation { kind, .. }) => { + format!("{kind:?}") + } FieldDataType::Sketch(kind, _) => format!("{:?}", kind.algorithm()), FieldDataType::ExactAggregate(kind, _) => format!("exact {kind:?} accumulator"), other => format!("{other:?}"), @@ -1420,19 +1514,29 @@ fn price( data: &DataWorkload, calibration: &Stage3Calibration, ) -> Result { - price_nodes(dag, demand, data, calibration).map(|(cost, _)| cost) + price_nodes(dag, demand, data, calibration, None).map(|priced| priced.cost) +} + +/// A candidate's price, and what Stage 3's checks read from pricing. +struct Priced { + cost: CandidateCost, + /// The cost of one evaluation of each query-time node. + per_evaluation: HashMap, + /// Bytes retained across evaluations: ingestion-time state that query + /// time reads, and raw data retained for query-time scans. + retained_bytes: u64, } -/// [`price`], and the cost of one evaluation of each query-time node. +/// [`price`], and what the checks read. `raw_bytes_per_sample` is `Some` +/// when the deployment does not keep raw data anyway: each query-time scan +/// then retains its raw samples, priced as memory (Q48). fn price_nodes( dag: &PhysicalASAPDAG, demand: &[RootDemand], data: &DataWorkload, calibration: &Stage3Calibration, -) -> Result< - (CandidateCost, HashMap), - (PhysicalASAPNodeId, AnalyticalCostError), -> { + raw_bytes_per_sample: Option, +) -> Result { let first = dag .roots .first() @@ -1459,6 +1563,8 @@ fn price_nodes( let reached = reaching_roots(dag); let nodes: HashMap<_, _> = dag.nodes.iter().map(|n| (n.id, n)).collect(); let roles = pane_roles(dag); + let raw_retention = raw_retention(dag, rows_per_ms, raw_bytes_per_sample); + let mut retained_bytes = 0u64; let mut output: HashMap = HashMap::new(); let mut per_node = BTreeMap::new(); let mut per_evaluation = HashMap::new(); @@ -1509,7 +1615,7 @@ fn price_nodes( true => 1_000, false => scan_extent_ms(dag, node.id, 0).unwrap_or(DEFAULT_LOOKBACK_MS), }; - let out = edge(((rows_per_ms * span_ms as f64).round() as u64).max(1)); + let out = edge(scan_rows(rows_per_ms, span_ms)); let estimate = estimate_operator( PhysicalOperator::Scan, OperatorStatistics::Scan { @@ -1709,8 +1815,9 @@ fn price_nodes( let read_at_query_time = dag.edges.iter().any(|e| { e.producer == node.id && nodes[&e.consumer].output_state.timing.is_query_time() }); - let retain = |windows: u64, cost: f64, what: &str| { + let mut retain = |windows: u64, cost: f64, what: &str| { let retained = windows * out.rows * state_bytes(node); + retained_bytes = retained_bytes.saturating_add(retained); ( cost + calibration.cost_per_retained_byte_second * retained as f64, format!("{detail}; ingestion time, {what}, retains {retained} bytes"), @@ -1743,7 +1850,17 @@ fn price_nodes( let rate = evaluation_rate(roots.filter_map(|&r| demand.get(r)), calibration.horizon_s) .map_err(|error| (node.id, error))?; per_evaluation.insert(node.id, cost); - (cost * rate, format!("{detail}; x {rate:.4} evaluations/s")) + let detail = format!("{detail}; x {rate:.4} evaluations/s"); + match raw_retention.get(&node.id) { + Some(&bytes) => { + retained_bytes = retained_bytes.saturating_add(bytes); + ( + cost * rate + calibration.cost_per_retained_byte_second * bytes as f64, + format!("{detail}; retains {bytes} bytes of raw data"), + ) + } + None => (cost * rate, detail), + } }; output.insert(node.id, out); per_node.insert( @@ -1755,8 +1872,8 @@ fn price_nodes( }, ); } - Ok(( - CandidateCost { + Ok(Priced { + cost: CandidateCost { total: per_node.values().map(|n| n.cost).sum(), unit: COST_PER_SECOND, source: format!( @@ -1766,7 +1883,46 @@ fn price_nodes( per_node, }, per_evaluation, - )) + retained_bytes, + }) +} + +/// Rows a scan over `span_ms` reads at `rows_per_ms`. +fn scan_rows(rows_per_ms: f64, span_ms: u64) -> u64 { + ((rows_per_ms * span_ms as f64).round() as u64).max(1) +} + +/// The raw bytes each query-time scan makes the deployment retain: +/// lookback × λ × bytes per sample. Empty when raw data is kept anyway +/// (`None`). A scan at ingestion time reads samples as they arrive and +/// retains none. Like its work, a scan's retention is charged once per +/// scan node, so it stays a sum over nodes; a scan shared by several +/// queries is one node. +fn raw_retention( + dag: &PhysicalASAPDAG, + rows_per_ms: f64, + raw_bytes_per_sample: Option, +) -> HashMap { + let Some(bytes_per_sample) = raw_bytes_per_sample else { + return HashMap::new(); + }; + dag.nodes + .iter() + .filter(|n| { + n.output_state.timing.is_query_time() + && matches!( + n.payload, + Payload::Relational { + operator: NonASAPOpKind::Scan { .. } + } + ) + }) + .map(|n| { + let span_ms = scan_extent_ms(dag, n.id, 0).unwrap_or(DEFAULT_LOOKBACK_MS); + let rows = scan_rows(rows_per_ms, span_ms); + (n.id, rows.saturating_mul(bytes_per_sample)) + }) + .collect() } /// Partitions a per-group ranking assumes, absent group-count evidence. @@ -2839,4 +2995,185 @@ mod tests { assert_eq!(years, [1, 2, 3, 5, 5]); assert!(shared <= separate, "shared {shared} vs separate {separate}"); } + + // ── Deployment capabilities (C3, Q48) ──────────────────────────────── + + fn reasons(selection: &Selection) -> BTreeMap<&str, &str> { + selection + .rejected + .iter() + .filter(|r| !r.valid) + .map(|r| (r.id.as_str(), r.reason.as_str())) + .collect() + } + + /// A summary the deployment cannot build, or a readout it cannot + /// compute, rejects the candidate as invalid, naming what is missing. + #[test] + fn missing_summary_or_readout_is_rejected_with_its_reason() { + use asap_types::deployment::{InstanceLayout, SummaryFamily, SummarySupport}; + let counter = [("m".to_string(), asap_types::workload::MetricType::Counter)].into(); + let target = AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.001, + }; + // Both heap sketches can be built; only Count-Min's top-k is read. + let support = |algorithm, readouts: &[_]| SummarySupport { + family: SummaryFamily::Sketch(algorithm), + layout: InstanceLayout::PerGroup, + readouts: readouts.iter().copied().collect(), + }; + use asap_types::deployment::Readout::TopK; + let caps = DeploymentCapabilities { + summaries: Some(vec![ + support(SketchAlgorithm::CmsWithHeap, &[TopK]), + support(SketchAlgorithm::CountSketchWithHeap, &[]), + SummarySupport { + family: SummaryFamily::Exact(asap_types::ir::schema::ExactKind::Sum), + ..support(SketchAlgorithm::Kll, &[]) + }, + ]), + ..DeploymentCapabilities::UNRESTRICTED + }; + let select = |caps| { + stage3_select( + &candidates_with(&counter), + &[every_10s(Some(target.clone()))], + &data(), + PlanningModels::builtin().with_capabilities(caps), + ) + .unwrap() + }; + let selection = select(&caps); + assert_eq!( + reasons(&selection), + BTreeMap::from([( + "P3", + "q1: deployment lacks a CountSketchWithHeap TopK readout" + )]) + ); + assert!(selection.costs.contains_key("P2")); + // Without the Count-Min + heap summary, P2 cannot be built either. + let without_cms = DeploymentCapabilities { + summaries: caps.summaries.clone().map(|mut s| { + s.remove(0); + s + }), + ..caps + }; + let selection = select(&without_cms); + assert_eq!( + reasons(&selection)["P2"], + "q1: deployment lacks a CmsWithHeap summary" + ); + } + + /// A deployment that cannot maintain state at ingestion time rejects + /// every candidate running a node then, and only those. + #[test] + fn ingestion_time_is_rejected_when_unsupported() { + let (run, demand, data) = pattern_b(data(), None); + let all: Vec<_> = run + .enumeration + .unwrap() + .candidates + .into_iter() + .flat_map(|c| c.physical) + .collect(); + let caps = DeploymentCapabilities { + ingestion_time: false, + ..DeploymentCapabilities::UNRESTRICTED + }; + let selection = stage3_select( + &all, + &demand, + &data, + PlanningModels::builtin().with_capabilities(&caps), + ) + .unwrap(); + let invalid = reasons(&selection); + for p in &all { + let maintained = !p.materialization.is_empty(); + assert_eq!(invalid.contains_key(p.id.as_str()), maintained, "{}", p.id); + if maintained { + assert_eq!( + invalid[p.id.as_str()], + "deployment cannot maintain state at ingestion time" + ); + } + } + } + + /// A candidate retaining more state than the deployment's memory budget + /// is rejected; at the budget it is priced. + #[test] + fn memory_budget_is_enforced_on_retained_state() { + let (run, demand, data) = pattern_b(data(), None); + let (p, _) = maintained_panes(&run); + let p = p.clone(); + // data(): 10 000 series; the newest pane keeps itself and 5 others. + let Payload::SummaryAgg { family, .. } = &node_of(&p.dag, is_build).payload else { + unreachable!() + }; + let retained = 6 * 10_000 * summary_shape(family).1; + let select = |budget| { + let caps = DeploymentCapabilities { + memory_budget_bytes: Some(budget), + ..DeploymentCapabilities::UNRESTRICTED + }; + stage3_select( + std::slice::from_ref(&p), + &demand, + &data, + PlanningModels::builtin().with_capabilities(&caps), + ) + }; + assert!(select(retained).is_ok()); + let Err(SelectionError::NoValidCandidate(rejected)) = select(retained - 1) else { + panic!("over the budget") + }; + assert_eq!( + rejected[0].reason, + format!( + "retains {retained} bytes across evaluations, over the deployment's memory \ + budget of {} bytes", + retained - 1 + ) + ); + } + + /// When the deployment does not keep raw data anyway, a query-time scan + /// pays w × lookback × λ × bytes per sample; an ingestion-time scan, and + /// any scan when raw data is kept, pays nothing for retention. + #[test] + fn raw_retention_is_charged_only_without_raw_data_and_not_for_maintained_inputs() { + let query = "quantile_over_time(0.99, m[1m])"; + let priced = |ingestion, raw| { + price_nodes( + &kll_dag(query, ingestion), + &[every_10s(None)], + &data(), + &Stage3Calibration::ILLUSTRATIVE, + raw, + ) + .unwrap() + }; + // data(): λ = 10 000 rows/s over 1 min, 16 bytes per sample. + let raw_bytes = 60 * 10_000 * 16; + let (kept, not_kept) = (priced(false, None), priced(false, Some(16))); + let scan = node_of(&kll_dag(query, false), is_scan).id; + let charge = not_kept.cost.per_node[&scan].cost - kept.cost.per_node[&scan].cost; + assert!( + (charge - 1.25e-7 * raw_bytes as f64).abs() < 1e-12, + "{charge}" + ); + assert!((not_kept.cost.total - kept.cost.total - charge).abs() < 1e-12); + assert_eq!(not_kept.retained_bytes - kept.retained_bytes, raw_bytes); + assert!(not_kept.cost.per_node[&scan] + .detail + .ends_with(&format!("retains {raw_bytes} bytes of raw data"))); + let (kept, not_kept) = (priced(true, None), priced(true, Some(16))); + assert_eq!(kept.cost, not_kept.cost); + assert_eq!(kept.retained_bytes, not_kept.retained_bytes); + } } diff --git a/crates/types/src/deployment.rs b/crates/types/src/deployment.rs new file mode 100644 index 00000000..d157e58a --- /dev/null +++ b/crates/types/src/deployment.rs @@ -0,0 +1,247 @@ +//! Deployment capabilities (#509 "Deployment inputs", #525): what the +//! deployment executing a plan can build, read out and keep. With the cost +//! and accuracy models, these are the deployment's inputs to planning; plan +//! selection rejects a candidate that needs a capability listed as missing. +//! +//! This is data only. The planner names no deployment: a deployment (for +//! example the reference executor) exports its own set, and the default is +//! unrestricted. +use std::collections::BTreeSet; + +use serde::{Deserialize, Serialize}; + +use crate::ir::schema::{ + ExactKind, FieldDataType, GroupingStrategy, HydraKind, SketchAlgorithm, SketchStatistic, +}; + +/// What a deployment can do. [`DeploymentCapabilities::UNRESTRICTED`] is the +/// default. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct DeploymentCapabilities { + /// The summaries the deployment can build, each with the readouts it can + /// compute from them; `None`: every summary and readout. + pub summaries: Option>, + /// Can maintain state at ingestion time, as data arrives. + pub ingestion_time: bool, + /// Can keep query-time results across evaluations (#509 Example 4, B3). + /// No candidate needs it until Stage 2 offers that option. + pub query_time_retention: bool, + /// Most bytes of state a plan may retain across evaluations. + pub memory_budget_bytes: Option, + /// The deployment keeps the raw data of the workload's sources anyway, + /// so a plan reading raw data at query time adds no retention (Q48). + pub raw_data_retained: bool, + /// Bytes one raw sample takes in retention, to price raw data a plan + /// makes the deployment keep when `raw_data_retained` is false. + pub raw_bytes_per_sample: u64, +} + +impl DeploymentCapabilities { + /// Every summary and readout, both timings, no memory budget, and raw + /// data kept anyway (so no raw retention is priced). One raw sample is + /// 16 bytes: an 8-byte timestamp and an 8-byte value, uncompressed. + pub const UNRESTRICTED: Self = Self { + summaries: None, + ingestion_time: true, + query_time_retention: true, + memory_budget_bytes: None, + raw_data_retained: true, + raw_bytes_per_sample: 16, + }; + + /// Why the deployment cannot build `family`, if it cannot. + pub fn missing_summary(&self, family: &FieldDataType) -> Option { + let Some(summaries) = &self.summaries else { + return None; + }; + let (summary, layout) = summary_of(family)?; + (!summaries + .iter() + .any(|s| s.family == summary && s.layout == layout)) + .then(|| format!("deployment lacks a {} summary", name(&summary, &layout))) + } + + /// Why the deployment cannot read `statistic` out of `family`, if it + /// cannot. + pub fn missing_readout( + &self, + family: &FieldDataType, + statistic: &SketchStatistic, + ) -> Option { + let Some(summaries) = &self.summaries else { + return None; + }; + let (summary, layout) = summary_of(family)?; + let readout = Readout::of(statistic); + (!summaries + .iter() + .any(|s| s.family == summary && s.layout == layout && s.readouts.contains(&readout))) + .then(|| { + format!( + "deployment lacks a {} {readout:?} readout", + name(&summary, &layout) + ) + }) + } +} + +impl Default for DeploymentCapabilities { + fn default() -> Self { + Self::UNRESTRICTED + } +} + +/// One summary the deployment can build, and what it can read out of it. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct SummarySupport { + pub family: SummaryFamily, + pub layout: InstanceLayout, + /// Sketch readouts. An exact state is finalized, not read out, so an + /// exact family lists none. + pub readouts: BTreeSet, +} + +/// A summary family, without its sizing parameters. +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +pub enum SummaryFamily { + Exact(ExactKind), + Sketch(SketchAlgorithm), +} + +/// How a summary's instances are laid out over the groups of a reduction. +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +pub enum InstanceLayout { + /// One instance per group. + PerGroup, + /// One shared Hydra structure for every group. + Hydra(HydraKind), +} + +/// A statistic read out of a sketch, without its arguments. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +pub enum Readout { + Quantile, + /// The total count of a state (a point count without an item). + TotalCount, + /// The count of one item. + ItemCount, + Cardinality, + TopK, + FrequencyL2, + FrequencyEntropy, +} + +impl Readout { + pub fn of(statistic: &SketchStatistic) -> Self { + match statistic { + SketchStatistic::Quantile { .. } => Self::Quantile, + SketchStatistic::PointCount { value: None, .. } => Self::TotalCount, + SketchStatistic::PointCount { value: Some(_), .. } => Self::ItemCount, + SketchStatistic::Cardinality => Self::Cardinality, + SketchStatistic::TopK { .. } => Self::TopK, + SketchStatistic::FrequencyL2 => Self::FrequencyL2, + SketchStatistic::FrequencyEntropy => Self::FrequencyEntropy, + } + } +} + +/// The family and layout of a summary state; `None` for other values. +pub fn summary_of(family: &FieldDataType) -> Option<(SummaryFamily, InstanceLayout)> { + match family { + FieldDataType::ExactAggregate(kind, _) => { + Some((SummaryFamily::Exact(kind.clone()), InstanceLayout::PerGroup)) + } + FieldDataType::Sketch(kind, grouping) => Some(( + SummaryFamily::Sketch(kind.algorithm().clone()), + match grouping { + GroupingStrategy::PerSubpopulationInstance => InstanceLayout::PerGroup, + GroupingStrategy::SharedMultiSubpopulation { kind, .. } => { + InstanceLayout::Hydra(kind.clone()) + } + }, + )), + _ => None, + } +} + +/// E.g. "Kll", "HydraCms" or "exact Sum". +fn name(family: &SummaryFamily, layout: &InstanceLayout) -> String { + match (family, layout) { + (_, InstanceLayout::Hydra(kind)) => format!("{kind:?}"), + (SummaryFamily::Sketch(algorithm), _) => format!("{algorithm:?}"), + (SummaryFamily::Exact(kind), _) => format!("exact {kind:?}"), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::ir::schema::{HydraParams, SketchKind, SketchParams}; + + fn hydra_cms() -> FieldDataType { + let (width, depth) = (64, 4); + FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Cms, SketchParams::Cms { width, depth }), + GroupingStrategy::SharedMultiSubpopulation { + kind: HydraKind::HydraCms, + params: HydraParams::HydraCms { + width, + depth, + shared_rows: 4, + shared_columns: 256, + }, + }, + ) + } + + /// Unrestricted capabilities accept every summary and readout. + #[test] + fn unrestricted_accepts_everything() { + let caps = DeploymentCapabilities::default(); + assert_eq!(caps, DeploymentCapabilities::UNRESTRICTED); + assert_eq!(caps.missing_summary(&hydra_cms()), None); + assert_eq!( + caps.missing_readout(&hydra_cms(), &SketchStatistic::TopK { k: 1 }), + None + ); + } + + /// A listed summary accepts only its listed readouts; an unlisted one, + /// or another layout of a listed family, is missing, named precisely. + #[test] + fn listed_summaries_and_readouts_bound_what_is_accepted() { + let caps = DeploymentCapabilities { + summaries: Some(vec![SummarySupport { + family: SummaryFamily::Sketch(SketchAlgorithm::Cms), + layout: InstanceLayout::Hydra(HydraKind::HydraCms), + readouts: BTreeSet::from([Readout::ItemCount]), + }]), + ..DeploymentCapabilities::UNRESTRICTED + }; + let item = SketchStatistic::PointCount { + key: crate::ir::scalar::ColumnRef::Named("item".into()), + value: Some("a".into()), + }; + assert_eq!(caps.missing_summary(&hydra_cms()), None); + assert_eq!(caps.missing_readout(&hydra_cms(), &item), None); + assert_eq!( + caps.missing_readout(&hydra_cms(), &SketchStatistic::TopK { k: 1 }) + .as_deref(), + Some("deployment lacks a HydraCms TopK readout") + ); + let per_group = FieldDataType::Sketch( + SketchKind::new( + SketchAlgorithm::Cms, + SketchParams::Cms { + width: 64, + depth: 4, + }, + ), + GroupingStrategy::PerSubpopulationInstance, + ); + assert_eq!( + caps.missing_summary(&per_group).as_deref(), + Some("deployment lacks a Cms summary") + ); + } +} diff --git a/crates/types/src/lib.rs b/crates/types/src/lib.rs index 3f666d8a..9ab88c3e 100644 --- a/crates/types/src/lib.rs +++ b/crates/types/src/lib.rs @@ -7,12 +7,15 @@ //! ([`ir::export`]). No execution logic lives in this crate (issue #190). //! - [`workload`] — planner inputs (#509): query and data workloads, the //! lowered [`workload::parsed_workload`], and [`workload::resources`]. +//! - [`deployment`] — the deployment's capabilities, a planner input (#509) +//! beside the cost and accuracy models. //! - [`physical`] — #509 Stage 2 decision data: exact-operator schema helpers //! and window-summary pane primitives. //! - [`types`] / [`dag_export`] / [`cost`] — accuracy targets, the generic //! DAG export, and cost annotations. pub mod cost; pub mod dag_export; +pub mod deployment; pub mod ir; pub mod physical; pub mod serde_f64; diff --git a/docs/design_docs/proposals/stage3-cost-model.md b/docs/design_docs/proposals/stage3-cost-model.md index cc37dbe0..36297c65 100644 --- a/docs/design_docs/proposals/stage3-cost-model.md +++ b/docs/design_docs/proposals/stage3-cost-model.md @@ -19,7 +19,7 @@ for the workload in steady state. | CPU work of every operator, at ingestion time or at query time | Transient query-time memory | | Bytes a scan reads | Storage tier and retention: disk versus memory, and for how long (S3, with panes) | | Memory held across evaluations by ingestion-time state that query time reads | Network, parallelism and partitioning | -| | Latency bounds (a separate check, below) and deployment capabilities | +| Raw data a query-time scan needs kept, when the deployment does not keep it anyway (Q48) | Latency bounds and deployment capabilities (separate checks, below) | | | One-time setup work such as backfill (deferred, S5) | Accuracy and latency are checks, not costs. Stage 3 rejects candidates that @@ -129,6 +129,54 @@ candidate over the bound is rejected with the query, the estimate and the bound, for example `q1: query-time work takes 310.0 ms per evaluation, over the 200 ms latency bound`. The estimate assumes one core and no queueing. +## Deployment inputs + +`PlanningModels` carries #509's three deployment inputs: the cost model +(with `Stage3Calibration`), the accuracy model, and the deployment's +capabilities, `DeploymentCapabilities` (in `asap-types`, module +`deployment`). Capabilities are data: the planner names no deployment. A +deployment exports its own set; the reference executor's is +`asap_executor::capabilities()`. Without one, the set is unrestricted. + +| Field | Unrestricted default | Meaning | +|---|---|---| +| `summaries` | `None`: every summary | Summary families (exact kind or sketch algorithm) × instance layout (per group, or a Hydra kind), each with the readouts it supports (quantile, total count, item count, cardinality, top-k, L2, entropy) | +| `ingestion_time` | `true` | Can maintain state at ingestion time | +| `query_time_retention` | `true` | Can keep query-time results across evaluations (Example 4, B3; no candidate needs it yet) | +| `memory_budget_bytes` | `None` | Most bytes a plan may retain across evaluations | +| `raw_data_retained` | `true` | The deployment keeps the sources' raw data anyway | +| `raw_bytes_per_sample` | 16 | An 8-byte timestamp and an 8-byte value, uncompressed | + +**Checks.** Stage 3 rejects a candidate that needs a capability the +deployment lacks, before checking accuracy, with the first missing one: + +* `q1: deployment lacks a CmsWithHeap summary`: a `SummaryAgg` whose + family and layout are not listed. +* `q1: deployment lacks a CountSketchWithHeap TopK readout`: a + `SummaryEstimate` whose statistic the summary it reads does not list. +* `deployment cannot maintain state at ingestion time`: any node at + ingestion time when `ingestion_time` is false. +* `retains N bytes across evaluations, over the deployment's memory budget + of M bytes`: the retained bytes below exceed `memory_budget_bytes`. + +**Raw retention (Q48).** If the deployment keeps raw data anyway, a plan +that reads raw data at query time adds nothing: that retention is sunk. +If it does not, the plan makes it keep the raw samples each query-time scan +reads, and pays for them as memory: + +```text +raw(scan) = w_mem · lookback(scan) · λ · raw_bytes_per_sample +``` + +`lookback(scan)` is the scan's extent (its longest range plus offset, as +priced above), so `lookback · λ` is the rows it reads per evaluation. The +charge is on the scan node, so cost stays a sum over nodes; a scan shared +by several queries is one node and pays once. A scan at ingestion time +reads samples as they arrive and retains none: a plan whose raw input is +maintained at ingestion time does not pay it. Retained bytes for the +memory budget are the ingestion-time state of the memory term plus, when +charged, this raw retention. + ## Calibration Every coefficient lives in `Stage3Calibration`, carried by `PlanningModels`