Skip to content

refactor(planner): plan one shared unified workload DAG - #542

Open
zzylol wants to merge 1 commit into
stack/528-06-nativefrom
stack/528-07-planner
Open

zzylol wants to merge 1 commit into
stack/528-06-nativefrom
stack/528-07-planner

Conversation

@zzylol

@zzylol zzylol commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor

Problem: the workload planner still builds legacy SummaryNode plans, so #511's single operator DAG never reaches selection

#511 (operator-sharing.md, §Goal and problem) asks to "use one operator model before and after ASAP optimization, so ordinary query operations and summary operations can form one visible computation DAG". Its §3 Acceptance criteria require that:

  • a projection uses the same semantics above and below summary computations;
  • a union or another ordinary operator can consume summary estimates on its inputs;
  • unifying the representation preserves existing DAG dependencies, including shared inputs;
  • scalar expressions use the same representation before and after optimization, "with no bridge nodes or hidden subplans".

The companion #511 doc (decoupling_op_and_expr.md, §2.3 and §4) adds that standalone scalar queries such as 2 or time() need no PromqlScalarBridge. #509 (§Stages and their decisions) adds that "a candidate is a DAG for the whole workload, not for one query".

The earlier layers of this stack (#537–#541) added the unified OperatorNode IR, its lowerers and the native compiler. They sit next to the production path. On the base branch (#541), workload planning still runs on the legacy types:

// base: crates/asap-aware-mapping/src/replacement.rs
pub enum Replacement {
    Summary(Rc<SummaryNode>),   // post-ASAP decision
    Rewrite(Rc<QueryExpr>),     // pre-ASAP alternative
    ExactComposition(ExactComposition),
}
pub fn keep_pre_asap(expr: &Rc<QueryExpr>) -> Result<Rc<SummaryNode>, RealizationError>;
// wraps the sub-DAG in SummaryExpr::KeepPreAsap(expr)

// base: crates/asap-aware-mapping/src/pass/mod.rs
pub struct PlanOutput { pub plans: Vec<QueryLifecyclePlan> }
pub fn dags(&self) -> Vec<Rc<SummaryNode>>;

// base: crates/types/src/parsed_workload.rs
pub fn new(workload: PlanningWorkload, exprs: Vec<Rc<QueryExpr>>) -> Result<Self, _>;

SummaryExpr repeats ordinary operators under a second name (ValueOperation::{Project, Filter, Sort, Limit}, RelationalJoin, a PromQL BinaryOp). It has no variant for set operations and the other ordinary operators. Any such operator has to be wrapped in KeepPreAsap(Rc<QueryExpr>), and a summary cannot appear below that wrapper. This is the "Today" tree from #511 §Goal and problem:

Base (#541 production path)                     After this PR
SummaryNode ValueOperation::Project             OperatorNode NonASAP Project
└─ SummaryNode SummaryEstimate                  └─ OperatorNode ASAP SummaryEstimate
   └─ SummaryNode SummaryAgg (KLL)                 └─ OperatorNode ASAP SummaryAgg (KLL)
      └─ KeepPreAsap(Rc<QueryExpr>)                   └─ OperatorNode NonASAP Project
         └─ QueryExpr Project                            └─ OperatorNode NonASAP Scan latency
            └─ QueryExpr Scan latency

Two more gaps:

  • A standalone scalar query cannot be a workload entry. ParsedWorkload holds only Rc<QueryExpr>, so the scalar is wrapped in a bridge node.
  • PlanOutput returns one SummaryNode root per query. There is no single view of the workload DAG, and no way to see that two roots share an operator.

This PR covers: the production cutover. Frontends, ParsedWorkload, Pass 1 replacement search (ASAPStrategies), the existing Pass 2 sharing (share_common_sub_dags), selection, DAG assembly, lifecycle costing and the native compiler all use OperatorNode and QueryRoot. It leaves out: new sharing rules, SummaryMerge planning rules, ASAP replacement inside scalar-root sub-DAGs, and deleting the old source files.

Proposed method

The pipeline keeps its stages. Only the type that flows between them changes.

  1. Frontend (stage 0). The canonical entry points now return the unified IR. lower_sql / lower_sql_dialect return Rc<OperatorNode>. lower_promql_query_workload returns Vec<QueryRoot>, and lower_metricsql_query returns one QueryRoot. A numeric literal query becomes QueryRoot::Scalar, with no bridge node. asap_planner::lower collects Vec<QueryRoot>.
  2. Workload boundary. ParsedWorkload::from_roots checks that there is one root per workload entry. It then splits the roots into operator roots, remembering each one's entry index in operator_indices, and (entry index, ScalarExpr) pairs. entries() yields only the operator entries.
  3. Pass 1: candidate generation. search_workload* runs asap_types::ir::cse::share_common_sub_dags over the operator roots. It asserts that no root already contains an ASAP operator. For each target it asks the strategies. SketchAlgorithmStrategy is renamed to ASAPStrategies. A candidate is now always an OperatorNode sub-DAG (Replacement::SubDAG):
    • a bound summary decision contains an ASAP operator;
    • a logical rewrite has no ASAP operator and no guarantee; is_logical_rewrite tells the two apart;
    • the "keep as is" fallback (retain_exact) returns the same sub-DAG with guarantee = exact("RetainedExact"). There is no wrapper node. A thread-local memo keyed by input pointer makes repeated calls return one Rc, so sharing is kept.
  4. Pass 2: sharing. No new rule. The existing identical-sub-DAG CSE now runs on OperatorNode before search and again after assembly (MajorPass). The second run turns identical summary producers chosen by different queries into one Rc, including producers that pre-ASAP CSE keeps apart (for example ungrouped aggregates, which have no unique key; see summary_sharing.rs::identical_ungrouped_queries_share_their_producers).
  5. Selection and assembly. GlobalSelection::assemble_selected_dag links per-target decisions into one OperatorNode DAG. A target with no chosen replacement keeps its own ordinary operator and assembles its children independently. So Project, SetOp, Join and other ordinary operators sit directly above summary nodes. assemble_selected_query adds an exact-state read boundary (finalize_query_candidate) when a query root needs one.
  6. Lifecycle and output. SummaryMaintenanceLifecyclePlan.root and every lifecycle and cost API take OperatorNode. PlanOutput keeps one QueryLifecyclePlan per operator entry. It also carries the scalar roots unchanged and exposes the whole workload DAG: roots() in entry order, and operators() with each unique node once. check_contract counts scalar roots, so entry indices stay 0..n.
  7. Physical / runtime. The native compiler rebuilds shared Rc<OperatorNode> references from the transport DAG (physical_planner/logical.rs) before lowering. "Readout" is renamed to "evaluation" in the runtime (evaluation module, ExactEvaluation, ExactAccumulator::evaluation). The exact SUM state records seen, so an empty or all-NULL SQL SUM is NULL instead of 0.

Key code interfaces

crates/types/src/parsed_workload.rs

pub struct ParsedWorkload {
    workload: PlanningWorkload,
    exprs: Vec<Rc<OperatorNode>>,
    operator_indices: Vec<usize>,
    scalars: Vec<(usize, ScalarExpr)>,
}
impl ParsedWorkload {
    pub fn new(workload: PlanningWorkload, exprs: Vec<Rc<OperatorNode>>) -> Result<Self, ParsedWorkloadError>;
    pub fn from_roots(workload: PlanningWorkload, roots: Vec<QueryRoot>) -> Result<Self, ParsedWorkloadError>;
    pub fn exprs(&self) -> &[Rc<OperatorNode>];
    pub fn entries(&self) -> impl Iterator<Item = (QueryWorkloadEntry, &Rc<OperatorNode>)> + '_;
    pub fn operator_indices(&self) -> &[usize];
    pub fn scalar_roots(&self) -> &[(usize, ScalarExpr)];
    pub fn len(&self) -> usize; // operator + scalar roots
}

crates/asap-aware-mapping/src/pass/mod.rs

pub struct PlanOutput {
    pub plans: Vec<QueryLifecyclePlan>,
    pub scalar_roots: Vec<(usize, asap_types::ir::ScalarExpr)>,
}
impl PlanOutput {
    pub fn entry_indices(&self) -> Vec<usize>;
    pub fn roots(&self) -> Vec<asap_types::ir::QueryRoot>;
    pub fn operator_roots(&self) -> Vec<Rc<OperatorNode>>;   // was dags() -> Vec<Rc<SummaryNode>>
    pub fn operators(&self) -> Vec<Rc<OperatorNode>>;
    pub fn len(&self) -> usize;
}

crates/asap-aware-mapping/src/replacement.rs

pub enum Replacement {
    SubDAG(Rc<OperatorNode>),           // replaces Summary(Rc<SummaryNode>) and Rewrite(Rc<QueryExpr>)
    ExactComposition(ExactComposition),
}
pub fn is_logical_rewrite(node: &OperatorNode) -> bool {
    node.guarantee.is_none() && !node.contains_asap()
}
pub fn retain_exact(expr: &Rc<OperatorNode>) -> Result<Rc<OperatorNode>, RealizationError>; // was keep_pre_asap

pub struct ASAPStrategies<'a> { planning_inputs: CandidatePlanningInputs<'a> } // was SketchAlgorithmStrategy

pub struct TargetSubDAG<'a> {
    pub root: &'a Rc<OperatorNode>,
    pub consumer_count: usize,
    pub strictest_sibling_accuracy: Option<&'a AccuracyTarget>,
}
pub struct CandidateLogicalASAPDAGs<Id> {
    pub roots: Vec<(Id, Rc<OperatorNode>)>,
    /* groups, order, composition_plans: private */
}
impl<'a> GlobalSelection<'a> {
    pub fn assemble_selected_dag(&self, target: &Rc<OperatorNode>)
        -> Result<Option<Rc<OperatorNode>>, RealizationError>;
    pub fn assemble_selected_query(&self, target: &Rc<OperatorNode>)
        -> Result<Option<Rc<OperatorNode>>, RealizationError>;
}
pub fn search_workload<Id>(roots: Vec<(Id, Rc<OperatorNode>)>) -> CandidateLogicalASAPDAGs<Id>;

crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs

pub struct SummaryMaintenanceLifecyclePlan {
    pub root: Rc<OperatorNode>,   // was Rc<SummaryNode>
    // other fields unchanged
}

Frontends (frontend-sql, frontend-promql, frontend-metricsql):

pub async fn lower_sql(query: &str, catalog: &SqlCatalog, accuracy: AccuracyTarget)
    -> Result<Rc<OperatorNode>, SqlError>;
pub fn lower_promql_query_workload(workload: &PlanningWorkload, now_ms: u64)
    -> Result<Vec<asap_types::ir::QueryRoot>, PromqlError>;
pub fn lower_promql_query_workload_with_histograms(workload: &PlanningWorkload, histograms: HistogramCatalog, now_ms: u64)
    -> Result<Vec<asap_types::ir::QueryRoot>, PromqlError>;
pub fn lower_metricsql_query(query: &str, accuracy: AccuracyTarget)
    -> Result<asap_types::ir::QueryRoot, MetricsqlError>;

Runtime (crates/asap-physical-operators/src/evaluation.rs, summary_kernels/exact.rs):

pub fn exact_evaluation(
    states: impl IntoIterator<Item = Arc<dyn AggregateCore>>,
    statistic: Statistic,
    range_ms: Option<(i64, i64)>,
    key: Option<&KeyByLabelValues>,
) -> Result<Option<f64>, String>;                 // was readout::exact_readout
pub fn insufficient_counter_samples(state: &dyn AggregateCore, statistic: Statistic) -> bool;

pub struct ExactEvaluation { pub statistic: Statistic, pub lookback_ms: Option<i64> } // was ExactReadout
enum ScalarState { Sum { sum: f64, compensation: f64, seen: bool }, /* … */ }

Usage (from operator_design_examples.rs):

let output = e2e_plan(UserInput::new(&workload, FrontendInput::Sql { catalog: &catalog }, models, lifecycle)).await?;
assert_eq!(output.entry_indices(), [0, 1]);
let states: Vec<_> = output.operators().into_iter()
    .filter(|n| matches!(n.asap(), Some(ASAPOp::SummaryAgg { .. })))
    .collect();
assert_eq!(states.len(), 1); // one shared SUM state for both roots

Changed in the same way, by name only (the type moves from QueryExpr/SummaryNode to OperatorNode): ReplacementStrategy::propose_for_root, TargetSubDAG::{new, with_consumer_count}, bindable_intent, RecurrenceProfileMap::for_target, GlobalSelection::for_target, CostModel hooks (summary_maintenance_lifecycle_cost_inputs, raw_query_recompute_cost, summary_support_evidence, …), dag_export::{export, export_summary}, MaintainedPopulation / PopulationInput (generic parameter removed), CurrentSeriesInput::matches_node. ASAPStrategies also gains public fixed_window_rate_candidates and query_time_rate_aggregation_candidates.

Fields

ParsedWorkload

Name Type Meaning
workload PlanningWorkload The input workload, unchanged.
exprs Vec<Rc<OperatorNode>> Operator roots only, in entry order.
operator_indices Vec<usize> Entry index of each element of exprs. Strictly increasing; entries() binary-searches it.
scalars Vec<(usize, ScalarExpr)> Standalone scalar roots with their entry index.
new(workload, exprs) Wraps every expr as QueryRoot::Operator and calls from_roots.
from_roots(workload, roots) LengthMismatch unless roots.len() equals the number of entries. Set by asap_planner::e2e_plan.
len() usize exprs.len() + scalars.len().

PlanOutput

Name Type Meaning
plans Vec<QueryLifecyclePlan> One per operator entry, in entry order. Each plan.root is an Rc<OperatorNode>; roots can share nodes.
scalar_roots Vec<(usize, ScalarExpr)> Copied from ParsedWorkload::scalar_roots by MajorPass. Not searched or rewritten.
entry_indices() Vec<usize> Plan and scalar indices together, sorted. check_contract requires 0..len().
roots() Vec<QueryRoot> Every query result in entry order, scalars included.
operator_roots() Vec<Rc<OperatorNode>> Only the plan roots.
operators() Vec<Rc<OperatorNode>> Every node reachable from roots(), including operators referenced by scalar expressions, each returned once (by pointer).
len() usize plans.len() + scalar_roots.len().

Replacement and helpers

Name Type Meaning
SubDAG Rc<OperatorNode> A candidate sub-DAG that replaces the target. Either a summary decision (contains an ASAP operator, or a retained exact sub-DAG with a guarantee) or a logical rewrite.
ExactComposition ExactComposition Unchanged: an exact operator over another target's own decision (#171).
is_logical_rewrite(node) bool true when node.guarantee is None and the sub-DAG has no ASAP operator. Replaces the old Summary/Rewrite split. Used for evidence and runtime-support checks.
retain_exact(expr) Result<Rc<OperatorNode>, _> Returns expr itself if it already has a guarantee or contains an ASAP operator. Otherwise returns a copy with ResultGuarantee::exact("RetainedExact"), memoized per input pointer.

ASAPStrategies<'a>

Name Type Meaning
planning_inputs CandidatePlanningInputs<'a> (private) Cost model, accuracy model and allocator. Set by default_cost_model, new, new_with_planning_inputs or new_with_planning_inputs_and_evidence. Behavior is the same as SketchAlgorithmStrategy.

TargetSubDAG<'a>

Name Type Meaning
root &Rc<OperatorNode> The pre-ASAP sub-DAG being replaced.
consumer_count usize Number of references to root in the workload. new sets 1. search_workload_with computes it.
strictest_sibling_accuracy Option<&AccuracyTarget> The strictest accuracy target among siblings that read the same summary input, if stricter than root's own. Set by search.

CandidateLogicalASAPDAGs<Id> / GlobalSelection

Name Type Meaning
roots Vec<(Id, Rc<OperatorNode>)> Workload roots after share_common_sub_dags. None may contain an ASAP operator (asserted).
assemble_selected_dag(target) Result<Option<Rc<OperatorNode>>, _> None if target was not discovered. Memoized per target, so a shared inner summary is one Rc across roots.
assemble_selected_query(target) same assemble_selected_dag and then finalize_query_candidate, which adds a FinalizeExactAccumulator read boundary at QueryTime where exact state reaches a query result. Use this for query results.
search_workload(roots) CandidateLogicalASAPDAGs<Id> Input is now the unified operator DAG.

SummaryMaintenanceLifecyclePlan.root: Rc<OperatorNode>, the root of the assembled DAG being deployed. The other fields are unchanged.

Frontends

Name Meaning
lower_sql / lower_sql_dialect query, catalog, (dialect), accuracy are unchanged. They return a resolved Rc<OperatorNode> whose schemas were derived during binding.
lower_promql_query_workload(workload, now_ms) One QueryRoot per entry. now_ms is the evaluation time used for lowering.
…_with_histograms(workload, histograms, now_ms) Same, with a HistogramCatalog installed during lowering.
lower_metricsql_query(query, accuracy) A NumberLiteral query returns QueryRoot::Scalar. Everything else returns QueryRoot::Operator.

Runtime

Name Type Meaning
exact_evaluation(states, statistic, range_ms, key) Result<Option<f64>, String> Merges states (each must be an ExactAccumulator; empty input is an error) and reads statistic for population key (None = unkeyed). range_ms extrapolates Rate/Increase. Ok(None) means the population is absent: a counter with fewer than two samples, or an empty MIN/MAX.
insufficient_counter_samples(state, statistic) bool true only for Rate/Increase on an exact state with too few samples.
ExactEvaluation.statistic Statistic The statistic the plan reads. It must match the accumulator family.
ExactEvaluation.lookback_ms Option<i64> PromQL counter window. The evaluation range is resolved from it at run time.
ScalarState::Sum.seen bool true after any update. Merge ORs it. A nullable unkeyed SUM output with seen = false is emitted as Value::Null.

Examples

End to end: two SQL queries share one SUM state

From crates/integration-tests/tests/operator_design_examples.rs::batch_planning_replaces_and_shares_summary_operators.

Input: a batch of two exact queries over requests(bytes Float64), two invocations each, data at rest. The test cost model makes a summary build cost 1 and raw recompute cost 1000.

SELECT SUM(bytes) + 1 AS result FROM requests   -- entry 0
SELECT SUM(bytes) * 2 AS result FROM requests   -- entry 1

What happens:

  1. lower returns two QueryRoot::Operator roots. Each is a Project over SUM(bytes) over Scan requests.
  2. ASAPStrategies proposes an exact SUM SummaryAgg for the SUM target. Selection picks it over raw recompute.
  3. Assembly keeps each Project as an ordinary operator above the summary. The exact-state read boundary sits between them.
  4. MajorPass runs share_common_sub_dags over the assembled roots, so the two identical SUM producers become one Rc:
Project(+1)  root 0        Project(*2)  root 1
      │                          │
      └── (exact-state read) ────┘
                 │
      SummaryAgg ExactAggregate(Sum)   ← one Rc, one deployment
                 │
             Scan requests

Output checks: entry_indices() == [0, 1], roots().len() == 2, exactly one SummaryAgg in operators(), both roots reach it (Rc::ptr_eq), selected_raw_recompute == false. Each root compiles to the transport DAG, runs on rows [10.0, 20.0], and returns 31.0 and 60.0.

Cases covered by tests

Case Result Test
UNION ALL of two approx_distinct root is SetOp { all: true }; each side holds a SummaryEstimate operator_sharing.rs::each_side_of_union_all_holds_a_summary_estimate
SELECT SUM(bytes) + 1 … WHERE status = 200, rewritten by hand schema unchanged; DAG has SummaryAgg + FinalizeExactAccumulator; JSON has no KeepPreAsap / ScalarBridge operator_design_examples.rs::sql_sum_projection_before_and_after_summary_rewrite
Same query run natively on non-empty / empty / all-NULL input 11 / NULL / NULL operator_design_examples.rs::sql_sum_example_executes_with_sql_null_semantics
PromQL p50 and p99 over lat[5m], same ε one shared KLL Rc, one deployment planner/tests/summary_sharing.rs::quantiles_with_equal_params_share_one_producer
p50 at ε=0.01 and p99 at ε=0.001 one KLL sized for 0.001; each root meets its own target summary_sharing.rs::quantiles_share_one_producer_sized_for_the_strictest_consumer
Different window ([5m] vs [10m]) or selector (job="a" vs "b") not shared, two deployments summary_sharing.rs::different_producers_are_not_shared
SQL p50 and p99 over the same filtered column one KLL; each root keeps its own output name summary_sharing.rs::sql_p50_and_p99_share_one_producer
Same SQL percentile with a different filter or column not shared same test
Scalar subquery / NOT IN subquery producer visible through a ScalarRef edge operator_design_examples.rs::sql_scalar_subquery_retains_its_cardinality_contract
avg with a Project above it ignored: Avg has no summary realization yet operator_sharing.rs::project_above_and_scan_below_a_summary_are_both_non_asap_nodes
avg + approx_percentile_cont sharing one Scan ignored: needs the rule that splits multi-measure aggregates operator_sharing.rs::exact_aggregate_and_sketch_share_one_scan

Out of scope

  • New sharing rules (summary-capability or window-composition) and SummaryMerge planning rules from docs: propose workload-wide planning, summary sharing, and materialization #509 §Pass 2.
  • ASAP replacement inside sub-DAGs referenced by scalar roots. Scalar roots pass through exact.
  • Splitting multi-measure aggregates (see the ignored test above).
  • Deleting obsolete files. post_asap/{expr,cse,post_asap_dag}.rs, readout.rs, the frontends' unified/ modules and similar files stay on disk but are no longer exported or compiled. The cleanup PR removes them. Duplicate transitional integration tests (unified_*) are removed here as their canonical counterparts take over.
  • Viewer assets and their contract test: final layer.

Stack and validation

Stack 7/8 · Previous: #541 · Next: #543 · Reference/tracker: #528

This is the main integration review. Review the planner behavior in asap-aware-mapping and planner, and the operator_design_examples / operator_sharing integration tests. That covers accuracy, costs, maintained-population lifecycles and scalar-root handling. The frontend and compiler changes promote implementations reviewed in the earlier layers to their canonical entry points. All compiled consumers migrate together.

Validation: full workspace tests and doctests; workspace all-target check; formatting; workspace all-target/all-feature Clippy with warnings denied. Includes native execution of the document examples and real batch sharing.

🤖 Generated with Claude Code

@zzylol
zzylol force-pushed the stack/528-06-native branch from 75903b5 to b22fa53 Compare October 2, 2026 18:32
@zzylol
zzylol force-pushed the stack/528-07-planner branch 2 times, most recently from 2a3bcd0 to a03efc2 Compare October 2, 2026 19:40
@zzylol
zzylol force-pushed the stack/528-06-native branch from b22fa53 to d82bcde Compare October 2, 2026 19:40
@zzylol
zzylol force-pushed the stack/528-07-planner branch from a03efc2 to 2c708f3 Compare October 2, 2026 21:14
@zzylol
zzylol force-pushed the stack/528-06-native branch 2 times, most recently from 337a0ff to a94bd59 Compare October 2, 2026 21:22
@zzylol
zzylol force-pushed the stack/528-07-planner branch from 2c708f3 to 03166e7 Compare October 2, 2026 21:22
@zzylol
zzylol force-pushed the stack/528-06-native branch from a94bd59 to d84830e Compare October 2, 2026 21:25
@zzylol
zzylol force-pushed the stack/528-07-planner branch 2 times, most recently from 4b84314 to d9da0f9 Compare October 2, 2026 21:56
@zzylol
zzylol force-pushed the stack/528-06-native branch from d84830e to b631137 Compare October 2, 2026 21:56
@zzylol
zzylol force-pushed the stack/528-07-planner branch from d9da0f9 to ec9f8cb Compare October 3, 2026 02:31
@zzylol
zzylol force-pushed the stack/528-06-native branch from b631137 to 075d8e2 Compare October 3, 2026 02:31
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant