Skip to content

feat(ir): add canonicalization, sub-DAG sharing and flat DAGs on the unified IR - #645

Merged
zzylol merged 1 commit into
mainfrom
stack/528-03-graph
Oct 6, 2026
Merged

zzylol merged 1 commit into
mainfrom
stack/528-03-graph

Conversation

@zzylol

@zzylol zzylol commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor

Replaces #537, which GitHub closed as merged into its old stack base (stack/528-02b-merge-structure) when the stack was reordered to put this PR first; its code never reached main. Review and approval are on #537. Same code as #537 at 4ac03c34, rebased on main d4869a7, minus the summary-coverage lines that now arrive with #567.

Problem

#536 defined the unified ir::OperatorNode DAG. But canonicalization, common sub-DAG sharing (CSE) and a flat, serializable form have not been implemented on these new types. They exist only for the old split IRs:

Function on the base branch Works on
pre_asap::canonicalize::canonicalize(QueryExpr) -> QueryExpr pre-ASAP QueryExpr
pre_asap::cse::share_common_sub_dags(Vec<(Id, QueryExpr)>) pre-ASAP QueryExpr
post_asap::cse post-ASAP SummaryNode
dag_export::export / export_post_asap QueryExpr (+ post-ASAP substitutions)

So, on OperatorNode:

  1. Two equivalent spellings stay different shapes. Example: Filter { EXISTS (s) } and Join { Semi, true, left, s } mean the same thing but cannot be matched structurally.

  2. Two identical sub-DAGs stay two Rcs. Example: one query with the same grouped p50 on both sides of a comparison:

    BinaryOp(==)
    ├── Aggregate(Quantile q=0.5 col=2, by [1]) ── Scan m     (allocation A)
    └── Aggregate(Quantile q=0.5 col=2, by [1]) ── Scan m     (allocation B)
    

    Nothing collapses A and B into one node with two consumers.

  3. A shared DAG cannot be written out flat. OperatorNode serializes as a tree, so every shared sub-DAG would be repeated. The earlier version of this PR fixed that with a second copy of the operator and scalar types (wire.rs: NonASAPOpKind, WireScalarExpr, ...) whose children were ids.

Therefore, this PR adds:

  • ir::canonicalize: rewrites equivalent spellings into one shape (fixes 1).
  • ir::cse: shares structurally identical sub-DAGs across a workload batch (fixes 2). This is the identical-expression rule of docs: propose workload-wide planning, summary sharing, and materialization #509 Pass 2: "The input and computation semantics are identical" → "One common computation node serving multiple consumers."
  • Generic child references (fixes 3 without copies): Operator, NonASAPOp, ASAPOp, ScalarExpr, Predicate, ProjectItem, SortKey and QueryRoot take the child reference as a type parameter, defaulting to Rc<OperatorNode>. ir::flat::flatten writes a DAG as a list of nodes whose operators are Operator<NodeId>. No logical export document is added: the planner's stages pass Rc<OperatorNode> DAGs directly, and the physical DAG (feat(runtime): compile unified operator and scalar graphs #541) is the first consumer of the flat form.

It does not implement the summary-capability rule or the window-composition rule of #509 Pass 2, Pass 1 candidate generation, or any physical-stage decision (materialization, timing assignment, phase splitting).

Proposed method

All on bound OperatorNode DAGs in crates/types/src/ir/.

1. Canonicalize (ir::canonicalize). Bottom-up over every child (operator inputs and nodes read by scalar expressions), memoized by Rc pointer so a shared sub-DAG is rewritten once. At each node it applies local rules until none matches:

  • Heavy-hitter promotion. Limit { n: k, offset: 0, Sort { one descending column key, [passthrough Project] Aggregate(Reduce(by), [Count | Sum]) } } becomes Aggregate(Reduce(partition), [TopK { k }]) over the unchanged inner aggregate. It is skipped for ascending sorts, offsets, ranking by a group key, a Limit partition that differs from the Sort partition, and Sum over a derived child Aggregate (rate/increase).
  • Subquery lowering. In a Filter, one EXISTS (s) conjunct becomes Join { Semi, true }, NOT EXISTS (s) becomes Join { Anti, true }, and a positive x IN (s) becomes Join { Semi, x = Column(left_width) }. Remaining conjuncts stay in an outer Filter. One conjunct per round.
  • Left alone on purpose: NOT IN, an IN whose probe contains a scalar subquery, scalar subqueries (a cross join would lose their zero-row NULL and multi-row error semantics), and ROW_NUMBER() filters (removing the window column would change the visible schema).

Untouched sub-DAGs keep their pointer identity, a rewritten shared sub-DAG stays shared, and the pass is idempotent. A Concat whose first branch changes schema drops its discriminator_unique_key.

2. Share common sub-DAGs (ir::cse). Hash-consing over all roots of a workload batch, bottom-up:

  • "Children" are operator inputs and operator nodes read by scalar expressions (PromqlScalarFromVector, ScalarSubquery, Exists, InSubquery). So scalar(v) in two queries can share v.
  • structural_hash only picks a bucket. The decision is typed equality (same_node): children by pointer (or memoized value comparison), then result_kind, schema, timing, guarantee, and the operator's own fields (operator.map_children(|_| ()), so children never enter the comparison twice). Guarantee and own fields must also serialize identically, so 0.0 vs -0.0 and NaN vs NaN are not shared.
  • Legality gate: a non-ASAP node is reused only if schema.has_unique_key(). ASAP nodes have no such gate.
  • Each root's value is unchanged (PartialEq-equal to its input); only its Rc structure may now alias.

The pass does not decide whether sharing is cheaper. That is selection's job (#509 stage 3).

3. Generic child references and ir::flat. Each operator and scalar type gets a parameter C for its child reference, defaulting to Rc<OperatorNode>, so existing code is unchanged. The traversal helpers (children, scalar_exprs, operator_refs, kind_name) work for any C; map_children / map_operator_refs map C to a new type D. Schema derivation and validation stay on the default Rc<OperatorNode> form.

flatten(roots) visits the DAG children-first and gives each distinct Rc (by pointer) the next index. Each node becomes a FlatNode: the same fields as OperatorNode, with operator: Operator<NodeId>. A scalar root stays a QueryRoot::Scalar(ScalarExpr<NodeId>). OperatorNode itself is not generic: a self-referential default (OperatorNode<C = Rc<OperatorNode>>) is rejected by rustc (E0391).

Key code interfaces

crates/types/src/ir/node.rs, non_asap.rs, asap.rs, scalar.rs, query.rs

pub enum Operator<C = Rc<OperatorNode>> { NonASAP(NonASAPOp<C>), ASAP(ASAPOp<C>) }
pub enum NonASAPOp<C = Rc<OperatorNode>> { /* children: C */ }
pub enum ASAPOp<C = Rc<OperatorNode>> { /* children: C */ }
pub enum ScalarExpr<C = Rc<OperatorNode>> { /* ScalarSubquery(C), Exists { subquery: C, .. }, ... */ }
pub struct Predicate<C = Rc<OperatorNode>>(pub ScalarExpr<C>);
pub struct ProjectItem<C = Rc<OperatorNode>> { pub alias: Option<String>, pub expr: ScalarExpr<C> }
pub struct SortKey<C = Rc<OperatorNode>> { pub expr: ScalarExpr<C>, pub ascending: bool, pub nulls_first: bool }
pub enum QueryRoot<C = Rc<OperatorNode>> { Operator(C), Scalar(ScalarExpr<C>) }

impl<C> Operator<C> {
    pub fn children(&self) -> Vec<&C>;   // operator inputs, then scalar-referenced nodes
    pub fn map_children<D>(&self, f: impl FnMut(&C) -> D) -> Operator<D>;
    pub fn kind_name(&self) -> &'static str;
}
impl<C> ScalarExpr<C> {
    pub fn operator_refs(&self) -> Vec<&C>;
    pub fn map_operator_refs<D>(&self, f: &mut impl FnMut(&C) -> D) -> ScalarExpr<D>;
}

crates/types/src/ir/flat.rs

pub type NodeId = usize;

pub struct FlatNode {
    pub operator: Operator<NodeId>,
    pub result_kind: OperatorResultKind,
    pub schema: Schema,
    pub guarantee: Option<ResultGuarantee>,
    pub timing: Option<ExecutionTiming>,
}

pub struct FlatDag {
    pub nodes: Vec<FlatNode>,           // node i is nodes[i]; children before parents
    pub roots: Vec<QueryRoot<NodeId>>,  // one per input root, in order
}

/// Also returns the original `Rc` for each id.
pub fn flatten(roots: &[QueryRoot]) -> (FlatDag, Vec<Rc<OperatorNode>>);

crates/types/src/ir/canonicalize.rs, cse.rs

pub fn canonicalize(root: Rc<OperatorNode>) -> Result<Rc<OperatorNode>, SchemaDerivationError>;

pub type HashCache = HashMap<*const OperatorNode, u64>;
pub fn structural_hash(node: &OperatorNode, cache: &mut HashCache) -> u64;
pub fn share_common_sub_dags<Id>(roots: Vec<(Id, Rc<OperatorNode>)>) -> Vec<(Id, Rc<OperatorNode>)>;
pub fn dag_node_count(root: &Rc<OperatorNode>) -> usize;

Usage:

let roots = share_common_sub_dags(vec![("q1", canonicalize(q1)?), ("q2", canonicalize(q2)?)]);
let (flat, originals) = flatten(&[QueryRoot::Operator(Rc::clone(&roots[0].1))]);
let json = serde_json::to_string(&flat)?;

Examples

Example 1: a shared node is flattened once (crates/types/src/ir/flat.rs tests)

Input: two roots Filter(true), each over the same Values Rc.

root 0: Filter ─┐
                ├── Values (one Rc)
root 1: Filter ─┘

flatten returns 3 nodes: 0 = Values, 1 = Filter { child: 0 }, 2 = Filter { child: 0 }, and roots [Operator(1), Operator(2)] (a_shared_node_is_flattened_once). A Project whose child and ScalarSubquery are the same Values gives 2 nodes, with child: 0 and cols[0].expr == ScalarSubquery(0) (scalar_references_become_node_ids). A constant scalar root gives 0 nodes (a_constant_scalar_root_has_no_nodes), and the JSON round-trips with every timing = None (json_round_trips).

Example 2: CSE on the identical-expression rule (crates/types/src/ir/cse.rs tests)

The BinaryOp from the Problem section, after share_common_sub_dags(vec![("q", root)]): lhs and rhs are the same Rc (single_query_shares_its_own_repeated_sub_dag). The grouped aggregate has unique key by [1], so it is legal to share.

Case (test) Shared? Why
Same grouped Quantile q=0.5 col=2 in two roots (median_and_explicit_half_percentile_merge) yes equal, unique key
Same grouped quantile on col=2 vs col=3 (distinct_column_quantiles_do_not_merge) no own fields differ
Identical global (ungrouped) quantile (no_unique_keys_means_no_merge_even_when_structurally_identical) no non-ASAP, no unique key
Identical without(..) aggregate (group_keys_gate_still_prevented_when_partition_by_without_used) no no unique key
Identical Dedup { cols: [1] } (dedup_gates_sharing_the_same_as_aggregate) yes Dedup adds a unique key
vector(scalar(sum by (service)(up))) twice (scalar_referenced_vector_is_shared_across_queries) inner vector yes, bridge no vector is a scalar-referenced child with a unique key; the bridge has none
Two DDSketch SummaryAgg, alpha=0.01, same guarantee (asap_nodes_share_without_a_unique_key) yes ASAP: no unique-key gate
DDSketch alpha=0.01 vs 0.001; exact guarantee vs None (asap_nodes_with_distinct_parameters_or_guarantees_are_not_shared) no parameters / guarantee differ
SummaryEstimate p95 and p99 over equal producers (evaluations_share_their_producer_but_not_each_other) producer yes, estimates no estimates differ in q
Literal 0.0 vs -0.0, +inf vs -inf, NaN vs NaN (signed_zero_and_nonfinite_values_remain_distinct) no typed + serialized equality; +inf vs +inf is shared

Example 3: canonicalization (crates/types/src/ir/canonicalize.rs tests)

Input Output
Limit 5 { Sort desc col 1 { Aggregate(by [1], [Count]) } } Aggregate(by [], [TopK{k:5}]) { Aggregate(by [1], [Count]) } (promotes_count_ranked_limit_sort); also through a passthrough Project
Same with ascending sort, or with an offset unchanged
Filter { EXISTS (s) } Join { Semi, true, left, s } (exists_filter_becomes_semi_join)
Filter { NOT EXISTS (s) } Join { Anti, true, left, s }
Filter { service IN (s) } over a 3-column left Join { Semi, service = Column(3) }
Filter { value > 1 AND EXISTS (s) } Filter { value > 1 } { Join { Semi } }
Filter { NOT IN (s) } unchanged
Project [ScalarSubquery(s) AS a, ScalarSubquery(s) AS b] unchanged; both still point to the same s (scalar_subqueries_remain_explicit_and_shared)
Filter { rn <= 5 } { ROW_NUMBER window } unchanged

Every rewrite test also checks canonicalize(out) returns the same Rc.

Out of scope

Stack and validation

First in the stack · Next: #567 · Tracker: #528

Order: this PR → #567 → #560 → #646 → #539 → #540 → #541 → #542 → #543.

🤖 Generated with Claude Code

…unified IR

#536 defined the unified OperatorNode DAG, but canonicalization, common
sub-DAG sharing and a flat, serializable form existed only for the old
split IRs. Add them on OperatorNode:

- ir::canonicalize: heavy-hitter promotion and EXISTS / NOT EXISTS / IN
  subquery lowering, bottom-up and memoized so shared sub-DAGs stay shared.
- ir::cse: hash-consing over a workload batch (the identical-expression
  rule of #509 Pass 2), following scalar-referenced operator nodes too.
- Generic child references: Operator, NonASAPOp, ASAPOp, ScalarExpr,
  Predicate, ProjectItem, SortKey and QueryRoot take the child reference
  as a type parameter (default Rc<OperatorNode>). ir::flat::flatten writes
  a DAG as nodes whose operators are Operator<NodeId>.

Summary coverage (#567) and SummaryMerge (#560) build on this.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
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