Skip to content

feat(ir): enable compatible summary state merges - #555

Draft
zzylol wants to merge 1 commit into
stack/528-12-latencyfrom
stack/528-13-summary-merge
Draft

zzylol wants to merge 1 commit into
stack/528-12-latencyfrom
stack/528-13-summary-merge

Conversation

@zzylol

@zzylol zzylol commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

Reorganization status. The structural merge scope was extracted earlier into #560, before phase-free #537. #560 also requires and derives #567's SummaryCoverage; this PR does not. This PR stays on the preserved legacy physical stack. Its runtime / pane-materialization acceptance and physical validation will move into the later physical scopes. The new logical-foundation chain is #560 → #537 → #539 → #540 → #561.

Problem: the window-composition rule needs a merge node, but SummaryMerge is reserved

#509 §Pass 2: ASAP-aware common-subexpression elimination defines the window-composition rule:

Window-composition rule — The computations have the same summary input data, and one window summary can answer the requested windows within their accuracy requirements. — One window summary feeding per-query merge (where needed) and estimation nodes.

and the tumbling-window case: "A longer query window is answered by merging the tumbling windows it covers … a 5-min window evaluated every 1 min uses 1-min tumbling windows and merges exactly 5 of them." #509 Example 4, Pattern B then gives two physical plans for the same merge DAG: B1 stores the 1-min KLLs at ingestion time and merges the latest 5 at query time; B2 rebuilds all 5 from raw samples at query time and merges them.

#511 §1.1 Unified Operator type lists SummaryMerge { children: Vec<Rc<OperatorNode>> } under "Reserved operations; semantics and support require further design", and says "Reserved ASAP variants require further semantic and capability design before use." #511 §2.3 adds the timing constraint any such operator must keep: "ingestion-time work cannot depend on query-time results".

Before this PR, the B1/B2 DAG could not be built:

SummaryEstimate(p99)
└── SummaryMerge
    ├── SummaryAgg(KLL k=200) ← Scan pane_0
    ├── …
    └── SummaryAgg(KLL k=200) ← Scan pane_4
OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children }))
  → Err(InvalidScalarSignature("this ASAP operator is reserved: schema, accuracy, timing and export are not implemented"))
timing::validate_* on a SummaryMerge
  → Err(UnimplementedOperator { operator: "SummaryMerge" })

This was so even though a wire payload and a native merge kernel already existed.

Scope. This PR covers the merge node of #509's window-composition rule, for state inputs that already exist in the DAG: schema derivation, input checks, timing validation at both phases, and export to the native compiler. It leaves out: generating window candidates automatically in Pass 2, building Exponential Histograms, coverage (which time range / population each input holds; added by #567/#560 on the new stack), and the other reserved operators (SummarySubtract, SummaryDelete, SummaryJoin, Extension).

Proposed method

SummaryMerge stops being reserved and gets the same contracts as the other ASAP operators. All changes are in the unified IR (crates/types); no new runtime code.

  1. Input checks (ASAPOp::validate_inputs, run during node construction):
    • children must not be empty.
    • The first child's schema must have exactly one non-Plain field (the state column). Plain grouping fields may accompany it.
    • Every child must have result_kind == State.
    • Every child's schema must equal the first child's schema. This covers state family, algorithm and parameters (KLL k=200 vs k=300 fails), grouping field positions and types, names, and schema metadata.
  2. Schema (output_schema): run the input checks, then return a clone of the first child's schema. output_kind is already State for SummaryMerge.
  3. State accessor (produced_state): return the dtype of the first child's non-plain field.
  4. Not reserved (is_unimplemented): SummaryMerge is removed from the reserved list.
  5. Timing (validate_asap in ir/timing.rs): every child must have primitive == SummaryState. If the merge runs at ingestion time, every child must also run at ingestion time. A query-time merge may read ingestion-time (stored) or query-time (rebuilt) children. Violations return IllegalChildDataState { edge: "SummaryMerge.children", child }.
  6. Export: with the above, apply_lifecycle_timings + compile_post_asap_dag accept the node, and the native compiler compiles it with the existing merge kernel.

Pipeline position: the node is a Pass 2 (window-composition) building block. Its timing is assigned in physical materialization, and the timing check runs on the physical DAG.

Key code interfaces

crates/types/src/ir/asap.rs

pub enum ASAPOp {
    // … SummaryAgg, SummaryEstimate, FinalizeExactAccumulator, MaintainPopulation, EvaluatePopulation …
    /// Merge compatible partial states for the same grouping and family.
    SummaryMerge { children: Vec<Rc<OperatorNode>> },
    // ── Reserved: migrated but unimplemented ──
    SummarySubtract { left: Rc<OperatorNode>, right: Rc<OperatorNode> },
    SummaryDelete { summary_input: Rc<OperatorNode>, key: ColumnId },
    SummaryJoin { outer: Rc<OperatorNode>, inner: Rc<OperatorNode>, key: ColumnId, family: FieldDataType },
    Extension { child: Rc<OperatorNode>, name: String },
}

impl ASAPOp {
    pub fn is_unimplemented(&self) -> bool;                       // no longer true for SummaryMerge
    pub fn produced_state(&self) -> Option<&FieldDataType>;       // now Some(..) for SummaryMerge
    pub fn output_schema(&self) -> Result<Schema, SchemaDerivationError>; // now derives SummaryMerge
    pub fn validate_inputs(&self) -> Result<(), SchemaDerivationError>;   // new SummaryMerge arm
}

crates/types/src/ir/timing.rs (private validate_asap, reached from validate_default and lifecycle timing):

ASAPOp::SummaryMerge { children } => {
    for child in children {
        let state = state_of(child);
        if state.primitive != DataPrimitive::SummaryState
            || (timing == ExecutionTiming::IngestionTime && state.timing != timing)
        {
            return Err(ExecutionDataStateError::IllegalChildDataState {
                edge: "SummaryMerge.children", child: state,
            });
        }
    }
    Ok(())
}

Usage (from crates/integration-tests/tests/planner_layering_merge.rs):

let merged = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children: states }))?;
let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryEstimate {
    summary_input: merged,
    query: SketchStatistic::Quantile { q: 0.99 },
}))?;
root.validate_structure()?;

Also changed: the doc comment on ExecutionDataStateError::UnimplementedOperator no longer lists SummaryMerge; a unit test in summary_maintenance_cost/model.rs that used SummaryMerge as its example of a reserved child now uses SummarySubtract.

Fields

ASAPOp::SummaryMerge

Field Type Meaning Invariants
children Vec<Rc<OperatorNode>> The partial states to merge, e.g. five 1-min KLL panes Non-empty; each child result_kind == State; all child schemas equal; the shared schema has exactly one non-Plain field

Output of a SummaryMerge node:

Property Value
schema Clone of children[0].schema
result_kind OperatorResultKind::State
produced_state() &dtype of the first child's non-plain field
Timing rule Children are SummaryState. Ingestion-time merge ⇒ children at ingestion time. Query-time merge ⇒ children at either phase.

Methods:

Method Change for SummaryMerge
is_unimplemented() Returns false. Still true for SummarySubtract, SummaryDelete, SummaryJoin, Extension.
produced_state() Metadata accessor; does not validate or merge anything
output_schema() Calls validate_inputs, then clones the first schema
validate_inputs() Errors are SchemaDerivationError::InvalidScalarSignature with: "summary merge requires at least one state input", "summary merge requires exactly one state column", "SummaryMerge requires summary state as input, got …", or "summary merge inputs must have identical state and grouping schemas"

ExecutionDataStateError::IllegalChildDataState { edge, child } (existing variant): edge = "SummaryMerge.children"; child is the offending child's ExecutionDataState (primitive and timing).

Examples

Five 1-min KLLs → one 5-min p99, rebuilt and materialized. Test five_minute_quantile_merges_five_one_minute_states in crates/integration-tests/tests/planner_layering_merge.rs (#509 Examples 3B / 4B):

  • Input: five Scans (pane_0 … pane_4), each feeding SummaryAgg with KLL k = 200 on SampleValue and no grouping; one SummaryMerge over the five; SummaryEstimate p99 on top. Each pane gets 20 rows, values pane*20 + i, so the merged input is 0 … 99.
  • What happens: validate_structure, export to wire, wire.validate(), native compile. Then two runs:
    • B2-style: run the whole plan at query time (evaluation_time_ms: 300_000).
    • B1-style: cut_candidate at the five SummaryAgg nodes; run the precompute part at ingestion time (window [0, 300_000)), then feed the stored states to the query part.
  • Output: both runs return one row, p99 = 98.0.

Contract tests. crates/types/tests/summary_merge.rs:

  • compatible_panes_merge_and_export: two KLL k=200 states merge; output schema has 1 field; timed export validates; validate_default passes at both IngestionTime and QueryTime, and planned_data_state(...).primitive == SummaryState.
  • incompatible_merge_inputs_fail: construction fails for each case below.
children Result
two KLL k=200 states accepted, schema = the KLL state schema
five KLL k=200 panes accepted (integration test)
[] rejected: at least one state input
KLL k=200 + KLL k=300 rejected: schemas differ
a raw Scan (the child of a SummaryAgg) rejected: its schema has no state column, so the "exactly one state column" check fails first
query-time child under an ingestion-time merge rejected by timing: IllegalChildDataState (rule; not in these tests)

Structural compatibility does not prove the inputs cover disjoint time ranges or populations. Two overlapping panes with equal schemas are accepted here. #567/#560 add that check on the new stack.

Out of scope

Stack and validation

Legacy physical stack: … ← #553 ← #554 ← #555 ← #556 … · Base: #554 (stack/528-12-latency) · Next: #556 · Tracker: #528. Continues the stack above #554, which is stacked on #543 through #551–#553.

Validation: type and mapping suites (745 tests/doctests passed); the new native E2E test passed; formatting; Clippy for types, mapping and integration tests with warnings denied. The regression failed before the implementation and passes afterward.

🤖 Generated with Claude Code

@zzylol

zzylol commented Oct 3, 2026

Copy link
Copy Markdown
Contributor Author

Parked as draft: PR priorities changed (see #528). Order is now (A) finish #511 operator sharing, (B) the #572 crate/module reorganization, (C) #509 end-to-end stages. This PR sits on the old stack/528-legacy-physical-base chain, and Phase B moves the files it touches. Its content will be re-scoped onto the new layout in Phase C.

🤖 Generated with Claude Code

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