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
58 changes: 39 additions & 19 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,18 +6,22 @@
// Writes an `asap-stage-pipeline/v1` document (tools/dag-viewer) with the
// four planner stages (#509 MVP):
// - stage0_logical: the frontends' workload DAG, one root per query;
// - stage1_logical_asap: Pass 1 workload candidates, one per choice of a
// local alternative for every target (the Cartesian product), in
// enumeration order and capped by `--max-candidates` (default 64);
// - stage1_logical_asap: Stage 1 workload candidates, one per choice of a
// local alternative for every target (Pass 1, the Cartesian product), for
// the queries as written and, when Pass 2's identical-expression rule
// merges something, again with identical sub-DAGs shared ("· shared
// input"); in enumeration order and capped by `--max-candidates`
// (default 64);
// - stage2_physical_asap: one physical candidate per logical candidate
// (operator implementation only, everything at query time), no cost;
// - stage3_selection: per-candidate costs, the selected candidate, and
// every other candidate as rejected (`valid: false`, including one that
// could not be built) or costlier.
//
// The enumeration is the library's (`plan_selection::select_exhaustive`);
// the facade's dynamic program selects the same winner when its
// assumptions hold.
// 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
// displayed candidate; the facade's dynamic program selects the same winner
// when its assumptions hold.
//
// `--promql` may repeat. `--epsilon`/`--delta` apply to every `--promql`
// query; without them the queries are exact. `--interval-ms` is the source
Expand All @@ -26,11 +30,11 @@
use std::rc::Rc;

use asap_logical_optimizer::pass1::logical_candidates::{
choice_index, enumerate_local_logical_candidates, LocalLogicalCandidates,
choice_index, combination_count, LocalLogicalCandidates,
};
use asap_logical_optimizer::Realization;
use asap_plan_selection::PlanningModels;
use asap_plan_selection::{select_exhaustive, Selection, MAX_ENUMERATED_CANDIDATES};
use asap_plan_selection::{plan_stages, Selection, MAX_ENUMERATED_CANDIDATES};
use asap_types::ir::flat::{flatten, FlatDag};
use asap_types::ir::schema::SketchAlgorithm;
use asap_types::ir::schema_support::with_promql_series_identity;
Expand Down Expand Up @@ -109,28 +113,44 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result<
})
.collect::<Result<Vec<_>, _>>()?;
let stage0 = export(&roots)?;
let inventory = enumerate_local_logical_candidates(roots.into_iter().enumerate().collect())
.map_err(|e| format!("Pass 1: {e}"))?;
let owners = target_owners(&inventory);
let targets: Vec<_> = workload
.query_workload
.entries()
.map(|entry| Some(entry.requirements.accuracy.target()))
.collect();
let data = workload.data_workload.clone().unwrap_or_default();
let enumeration = select_exhaustive(
&inventory,
let run = plan_stages(
roots.into_iter().enumerate().collect(),
&targets,
&data,
PlanningModels::builtin(),
max_candidates,
max_candidates.max(1),
)
.map_err(|e| format!("Stage 3: {e}"))?;
.map_err(|e| format!("planning: {e}"))?;
let enumeration = run.enumeration.expect("display was requested");
let combinations = enumeration.combinations;
let mut candidates = Vec::new();
let mut stage2 = Vec::new();
for candidate in &enumeration.candidates {
let index = choice_index(&inventory, &candidate.choice) + 1;
let label = label(&inventory, &owners, &candidate.choice);
// Candidates of the shared variant are numbered after the independent ones.
let mut offset = 0;
let variant = run
.stage1
.iter()
.find(|v| {
let found = v.shared == candidate.shared;
if !found {
offset += combination_count(&v.inventory);
}
found
})
.expect("the candidate's variant");
let inventory = &variant.inventory;
let index = offset + choice_index(inventory, &candidate.choice) + 1;
let mut label = label(inventory, &target_owners(inventory), &candidate.choice);
if candidate.shared {
label += " · shared input";
}
if let Some(logical) = &candidate.logical {
let roots: Vec<_> = logical.iter().map(|(_, root)| root.clone()).collect();
candidates
Expand All @@ -147,8 +167,8 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result<
"workload": { "queries": workload_queries(workload) },
"stage0_logical": { "dag": stage0 },
"stage1_logical_asap": {
"combinations": enumeration.combinations,
"capped": enumeration.combinations > max_candidates,
"combinations": combinations,
"capped": combinations > max_candidates,
"candidates": candidates,
},
"stage2_physical_asap": { "candidates": stage2 },
Expand Down
155 changes: 94 additions & 61 deletions crates/integration-tests/tests/planner_layering_example1.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,9 @@
//! MVP scope: Stage 1 = Pass 1 + the identical-expression rule only (no
//! window-composition variants); Stage 2 = physical operator implementation
//! only (no materialization). Counts follow the planner's output (user
//! decision): 1 → 24 → 24 → 1, because Pass 1 also offers exact accumulators.
//! The doc's 1 → 54 → 156 → 1 needs window composition and materialization.
//! decision): 1 → 48 → 48 → 1, because Pass 1 also offers exact accumulators
//! (24 combinations) and Pass 2 adds a shared-input variant of each. The
//! doc's 1 → 54 → 156 → 1 needs window composition and materialization.

use std::collections::{BTreeMap, BTreeSet, HashSet};

Expand All @@ -36,9 +37,8 @@ mod stages {
use std::rc::Rc;

use super::*;
use asap_logical_optimizer::pass1::logical_candidates::{
compose_logical_candidate, enumerate_local_logical_candidates,
};

use asap_plan_selection::{plan_stages, MAX_ENUMERATED_CANDIDATES};
use asap_types::ir::{OperatorNode, QueryRoot};

/// One whole-workload candidate. `query_roots` holds one root per
Expand Down Expand Up @@ -141,39 +141,41 @@ mod stages {
candidate("S0".into(), "frontend".into(), lower(workload))
}

/// Stage 1: every combination of Pass 1 local alternatives (Pass 2 is
/// not implemented). Lowers `workload` again: Pass 1 reads the in-memory
/// Stage 1: every combination of Pass 1 local alternatives, independent
/// and with the shared input (Pass 2), as the library's stage pipeline
/// enumerates them. Lowers `workload` again: Pass 1 reads the in-memory
/// DAG, not the Stage 0 export.
pub fn stage1_logical_asap(
workload: &PlanningWorkload,
_logical: &LogicalCandidate,
) -> Vec<LogicalCandidate> {
let inventory =
enumerate_local_logical_candidates(lower(workload).into_iter().enumerate().collect())
.expect("Pass 1");
let mut choices = vec![vec![]];
for target in &inventory.targets {
choices = choices
.into_iter()
.flat_map(|prefix: Vec<usize>| {
(0..target.alternatives.len()).map(move |i| {
let mut choice = prefix.clone();
choice.push(i);
choice
})
})
.collect();
}
choices
let targets: Vec<_> = workload
.query_workload
.entries()
.map(|entry| Some(entry.requirements.accuracy.target()))
.collect();
let run = plan_stages(
lower(workload).into_iter().enumerate().collect(),
&targets,
workload.data_workload.as_ref().expect("data workload"),
PlanningModels::builtin(),
MAX_ENUMERATED_CANDIDATES,
)
.expect("plans");
let enumeration = run.enumeration.expect("enumerated");
assert_eq!(enumeration.candidates.len(), enumeration.combinations);
enumeration
.candidates
.into_iter()
.enumerate()
.map(|(index, choice)| {
let roots = compose_logical_candidate(&inventory, &choice)
.expect("composes")
.into_iter()
.map(|(_, root)| root)
.collect();
candidate(format!("L{}", index + 1), format!("{choice:?}"), roots)
.map(|c| {
let physical = c.physical.expect("every Example 1 candidate builds");
let label = format!("{:?}{}", c.choice, if c.shared { " shared" } else { "" });
let roots = c.logical.expect("composes");
candidate(
physical.from_logical,
label,
roots.into_iter().map(|(_, root)| root).collect(),
)
})
.collect()
}
Expand Down Expand Up @@ -610,23 +612,21 @@ fn stage0_queries_do_not_share_nodes() {

// ── Stage 1 ──────────────────────────────────────────────────────────────

/// Stage 1 outputs exactly the 24 Pass 1 combinations, each once, with
/// separate inputs (no Pass 2 sharing yet).
/// Stage 1 outputs the 24 Pass 1 combinations twice: L1–L24 with separate
/// inputs, then L25–L48 with the shared input (Pass 2).
#[test]
fn stage1_has_24_candidates_covering_every_combination() {
fn stage1_has_48_candidates_covering_every_combination_twice() {
let (_, logical, _) = pipeline();
assert_eq!(logical.len(), 24);
let found: BTreeSet<_> = logical
.iter()
.map(|c| choices(&c.dag, &c.query_roots))
.collect();
assert_eq!(found, expected_choices(), "each combination exactly once");
for c in &logical {
assert!(
!classify(&c.dag, &c.query_roots).1,
"{}: shared input",
c.id
);
assert_eq!(logical.len(), 48);
for (half, shared) in [(&logical[..24], false), (&logical[24..], true)] {
let found: BTreeSet<_> = half
.iter()
.map(|c| choices(&c.dag, &c.query_roots))
.collect();
assert_eq!(found, expected_choices(), "each combination exactly once");
for c in half {
assert_eq!(classify(&c.dag, &c.query_roots).1, shared, "{}", c.id);
}
}
}

Expand Down Expand Up @@ -668,9 +668,9 @@ fn stage1_q2_summary_families_are_heap_sketches_and_hydra() {
);
}

/// Sharing adds a variant and keeps the independent one, for every Q2 option.
/// Sharing adds a variant and keeps the independent one, for every Q2 option
/// Pass 1 offers (Hydra: see `stage1_q2_summary_families_are_heap_sketches_and_hydra`).
#[test]
#[ignore = "missing features: Pass 2 (identical-expression sharing) is not implemented, so there is no shared-input variant; Pass 1 has no Hydra alternative"]
fn stage1_keeps_independent_and_shared_variants() {
let (_, logical, _) = pipeline();
let found: Vec<_> = logical
Expand All @@ -681,7 +681,6 @@ fn stage1_keeps_independent_and_shared_variants() {
Q2Option::Exact,
Q2Option::CountMinHeapPerJob,
Q2Option::CountSketchHeapPerJob,
Q2Option::Hydra,
] {
assert!(
found.contains(&(option, false)),
Expand Down Expand Up @@ -725,11 +724,11 @@ fn stage1_candidates_are_valid_and_uniquely_named() {

// ── Stage 2 ──────────────────────────────────────────────────────────────

/// No candidate is discarded before Stage 3: Stage 2 maps the 24 logical candidates one-to-one.
/// No candidate is discarded before Stage 3: Stage 2 maps the 48 logical candidates one-to-one.
#[test]
fn stage2_keeps_every_logical_candidate() {
let (_, logical, physical) = pipeline();
assert_eq!(physical.len(), 24);
assert_eq!(physical.len(), 48);
let sources: BTreeSet<_> = physical.iter().map(|p| p.from_logical.as_str()).collect();
let logical_ids: BTreeSet<_> = logical.iter().map(|c| c.id.as_str()).collect();
assert_eq!(sources, logical_ids);
Expand Down Expand Up @@ -761,7 +760,7 @@ fn stage2_exact_topk_is_sort_then_limit() {
.iter()
.filter(|p| classify(&p.dag, &p.query_roots).0 == Q2Option::Exact)
.collect();
assert_eq!(exact.len(), 8);
assert_eq!(exact.len(), 16);
for p in exact {
let sort_then_limit = p.dag.edges.iter().any(|e| {
relational(payload(&p.dag, e.producer)).as_deref() == Some("sort")
Expand Down Expand Up @@ -884,7 +883,7 @@ fn stage2_runtime_compiles_exactly_the_candidates_stage3_finds_valid() {
.filter(|r| !r.valid)
.map(|r| (r.id.as_str(), r.reason.as_str()))
.collect();
assert_eq!(invalid.len(), 8, "the Count-Min + heap candidates");
assert_eq!(invalid.len(), 16, "the Count-Min + heap candidates");
for p in &physical {
let compiled = compile_in_runtime(p);
match invalid.get(p.id.as_str()) {
Expand Down Expand Up @@ -1013,9 +1012,9 @@ fn stage3_charges_each_node_once() {
}

/// Sharing the input never costs more than reading it separately, for the
/// same local choices. Count-Min + heap is invalid here and has no cost.
/// same local choices. Count-Min + heap is invalid here and has no cost
/// (Hydra: see `stage1_q2_summary_families_are_heap_sketches_and_hydra`).
#[test]
#[ignore = "missing features: Pass 2 (identical-expression sharing) is not implemented, so there are no shared-input variants to compare; Pass 1 has no Hydra alternative"]
fn stage3_shared_input_is_not_costlier() {
let (workload, _, physical) = pipeline();
let selection = stage3_select(&workload, &physical, PlanningModels::builtin());
Expand All @@ -1027,11 +1026,7 @@ fn stage3_shared_input_is_not_costlier() {
Some(((choices(&p.dag, &p.query_roots), shared), cost))
})
.collect();
for option in [
Q2Option::Exact,
Q2Option::CountSketchHeapPerJob,
Q2Option::Hydra,
] {
for option in [Q2Option::Exact, Q2Option::CountSketchHeapPerJob] {
assert!(
by_combo.keys().any(|((_, o, _), _)| *o == option),
"{option:?} has no priced candidate"
Expand All @@ -1048,3 +1043,41 @@ fn stage3_shared_input_is_not_costlier() {
"no shared variant"
);
}

/// The selected plan shares the input, the doc's "Raw with a shared input"
/// winner; its saving over the same choices read separately is exactly one
/// scan and one range node, priced once instead of twice.
#[test]
fn stage3_selects_a_shared_input_plan() {
let (workload, _, physical) = pipeline();
let selection = stage3_select(&workload, &physical, PlanningModels::builtin());
let selected = physical
.iter()
.find(|p| p.id == selection.selected)
.unwrap();
assert!(classify(&selected.dag, &selected.query_roots).1);
let separate = physical
.iter()
.find(|p| {
!classify(&p.dag, &p.query_roots).1
&& choices(&p.dag, &p.query_roots) == choices(&selected.dag, &selected.query_roots)
})
.unwrap();
let input_cost: f64 = selected
.dag
.nodes
.iter()
.filter(|n| {
matches!(
relational(&n.payload).as_deref(),
Some("scan" | "time_range")
)
})
.map(|n| selection.costs[&selected.id].per_node[&n.id])
.sum();
let saving = selection.costs[&separate.id].total - selection.costs[&selected.id].total;
assert!(
(saving - input_cost).abs() < 1e-9,
"{saving} vs {input_cost}"
);
}
8 changes: 5 additions & 3 deletions crates/logical-optimizer/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,11 @@
//! ([`accuracy::allocation`]). Candidates no analytical rule can prove
//! invalid are kept.
//!
//! **Common sub-expression elimination (CSE) of identical sub-DAGs is not this
//! crate's job.** It runs over the pre-ASAP IR itself
//! (`asap_types::ir::cse`, issues #222/#223) before search. Pass 2 recognizes
//! **Common sub-expression elimination (CSE) of identical sub-DAGs is not
//! implemented here.** It runs over the pre-ASAP IR itself
//! (`asap_types::ir::cse`, issues #222/#223). Pass 2's identical-expression
//! rule ([`pass2::identical_expressions`]) decides when to use it: the stage
//! pipeline keeps a variant with and without it. Pass 2 also recognizes
//! sharing that is invisible at that level, such as `Quantile(x, 0.99)` and
//! `Quantile(x, 0.95)` reading one built sketch.
//!
Expand Down
Loading