Skip to content

Generate exact window composition and materialization candidates - #566

Draft
zzylol wants to merge 1 commit into
stack/509-20-sql-entropy-fallbackfrom
stack/509-21-window-composition
Draft

zzylol wants to merge 1 commit into
stack/509-20-sql-entropy-fallbackfrom
stack/509-21-window-composition

Conversation

@zzylol

@zzylol zzylol commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

Problem: the #509 window-composition rule and the Example 4 materialization choices are not generated automatically

#509 §Pass 2: ASAP-aware common-subexpression elimination defines the window-composition rule: computations with the same summary input data can share "one window summary feeding per-query merge (where needed) and estimation nodes". For tumbling windows it says:

The tumbling length must divide both the query window length and the evaluation interval, so that every query window starts and ends on a tumbling boundary: a 5-min window evaluated every 1 min uses 1-min tumbling windows and merges exactly 5 of them.

#509 §2 Materialization then requires stage 2 to give each sub-DAG the options "materialized at ingestion time", "materialized at query time" and "not materialized". §Example 4, Pattern B lists them for the same 5-min KLL: B1 (store 1-min panes from ingestion), B2 (rebuild all 5 panes at every query), B3 (keep panes built at query time).

Before this PR, neither step was generated by the planner. For #509 §Example 3 Pattern B:

quantile_over_time(0.99, events[5m]), repeated every 60 s

Pass 1 produces one summary over the whole lookback:

SummaryAgg(KLL k=200, input: sample value)
└── TimeRange(range: 300 s, kind: Range)
    └── Scan(TimeSeries "events")
  • Window composition. SummaryMerge (feat(ir): define compatible logical summary merges #560) and native merge exist, but a caller had to build the five 1-min pane subtrees and the merge by hand. No rule derived the pane width from the workload's cadence.
  • Materialization. enumerate_frontiers and compile_candidates existed as two separate calls. Each lowered the DAG again, and the caller had to pair each frontier with its result.
  • Physical lowering. For a SummaryAgg with reduction: PerEntity, compilation required TimeRange directly over Scan. A pane is TimeRange over TimeShift over Scan, so a per-entity pane was rejected with per-entity summary requires a resolved source. (The KLL example above is not per-entity and did not hit this check.)

Scope covered here. The tumbling-window part of the window-composition rule, as an exact logical alternative per summary, and automatic enumeration of every legal producer/reader split of the lowered DAG (the Example 4 Pattern B options B1/B2/B3).

Left out. Sliding windows and Exponential Histograms. A rotating cache that keeps panes across evaluations. Historical EH boundary certificates. Workload-level selection over these candidates (#568). Temporal merge timestamp binding (#569). The panes do not carry #567 SummaryCoverage: this branch is on the #541 line, which predates it.

Proposed method

Two independent pieces.

1. Logical: enumerate_window_compositions (Pass 2, asap-aware-mapping).

  1. Derive the cadence from entry.recurrence:
    • Repeated(FixedInterval(i)) or Repeated(FixedIntervalAt { interval: i, .. }) → i (the phase is not used).
    • Repeated(Scheduled(times)) → GCD of the differences between consecutive times.
    • anything else (OneTime, Unknown, EstimatedRate) → 0.
  2. If limit == 0, return BudgetExceeded. If the cadence is 0, return only the original root.
  3. For each reachable node, try compose(node, cadence, limit). A node is a composition site when it is:
    • SummaryAgg whose child is TimeRange { range, kind },
    • range is a whole, nonzero number of milliseconds (lookback),
    • below the TimeRange is a Scan of Source::TimeSeries, or a TimeShift directly over such a Scan.
  4. Pane width is width = gcd(lookback, cadence), and count = lookback / width. This is the largest width that divides both the window and the evaluation interval, as docs: propose workload-wide planning, summary sharing, and materialization #509 requires. If count <= 1 there is nothing to compose. If count > limit, return BudgetExceeded.
  5. Build count panes. Pane i is a copy of the original SummaryAgg (same family, input, reduction, grouping, filter) whose input is TimeRange(width, same kind) over TimeShift { offset_ms: shift.offset_ms + i·width, at: shift.at } over the same Scan. Pane 0 is the most recent width; pane count−1 is the oldest. The panes are disjoint and together cover exactly the original [t − lookback, t) (plus any original offset/anchor).
  6. Wrap the panes in SummaryMerge. If SummaryMerge construction fails (for example, incompatible state schemas), the site is skipped. The merge node takes the original node's schema and guarantee.
  7. Return every combination: for n sites, 2^n candidates (bit k of the mask = use the composed form at site k). Mask 0 is the original root, unchanged. If 2^n > limit, return BudgetExceeded and no partial list. Ancestors of a replaced node are rebuilt with map_children and keep their guarantee.

This is candidate generation only. It does no costing and does not check the family's merge capability; physical compilation does that.

2. Physical: compile_materialization_candidates (stage 2, asap-physical-operators).

  1. Lower the PostAsapDAG once (compile).
  2. Enumerate frontiers on that compiled DAG. A frontier is a set of bounded operator nodes in which no node is an ancestor of another. The empty frontier (nothing stored) is always included. If the count would exceed max_candidates, return an error and no partial list.
  3. For each frontier, call the existing cut_candidate to split the DAG into a precompute part (produces the frontier outputs) and a query part (reads them). A failed cut is kept as Err in that candidate. It does not fail the whole call.

The physical compiler also now accepts TimeRange → TimeShift → Scan under a PerEntity summary, so generated per-entity panes lower.

Key code interfaces

crates/asap-aware-mapping/src/window_composition.rs (new public module window_composition):

#[derive(Debug, thiserror::Error)]
pub enum WindowCompositionError {
    #[error("window composition exceeds candidate budget; no partial inventory returned")]
    BudgetExceeded,
    #[error("window composition: {0}")]
    Schema(#[from] asap_types::ir::SchemaDerivationError),
}

pub fn enumerate_window_compositions(
    root: &Rc<OperatorNode>,
    entry: &QueryWorkloadEntry,
    limit: usize,
) -> Result<Vec<Rc<OperatorNode>>, WindowCompositionError>;

crates/asap-physical-operators/src/physical_planner/candidates.rs (re-exported from physical_planner):

/// One enumerated frontier with its lowering result, including any rejection.
pub struct MaterializationCandidate {
    pub frontier: Vec<NodeId>,
    pub realization: Result<PhysicalASAPDAG, Error>,
}

pub fn compile_materialization_candidates(
    dag: &PostAsapDAG,
    inputs: BTreeMap<NodeId, InputContract>,
    roots: &[NodeId],
    max_candidates: usize,
) -> Result<Vec<MaterializationCandidate>, Error>;

Existing type it returns (unchanged):

pub struct PhysicalASAPDAG {
    pub precompute: Option<CompiledPhysicalDAG>,
    pub query: CompiledPhysicalDAG,
    pub materialized_outputs: BTreeMap<NodeId, InputContract>,
}

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

let variants = enumerate_window_compositions(&state, &entry, 64)?;   // [original, 5-pane merge]
let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryEstimate {
    summary_input: variants[1].clone(),
    query: SketchStatistic::Quantile { q: 0.99 },
}))?;
let wire = compile_post_asap_dag(&root)?;
let choices = compile_materialization_candidates(&wire, inputs, &[u64::from(wire.root.0)], 1024)?;

Fields

enumerate_window_compositions

Item Type Meaning
root &Rc<OperatorNode> Logical candidate root. Every reachable SummaryAgg is a possible site.
entry &QueryWorkloadEntry Only recurrence is read, to get the cadence (see method step 1).
limit usize Budget. Must be > 0. Bounds both the panes per site and the total number of returned candidates (2^sites).
return Result<Vec<Rc<OperatorNode>>, _> All combinations. Element 0 is root itself (Rc::ptr_eq). Unknown cadence or no sites → vec![root].

WindowCompositionError

Variant When
BudgetExceeded limit == 0; a site needs more than limit panes; or 2^sites > limit. No partial list is returned.
Schema(SchemaDerivationError) Building a pane TimeShift/TimeRange node, or rebuilding a parent, failed schema derivation.

Pane construction (internal compose)

Value Source Meaning
lookback TimeRange.range in ms Original window length. Must be whole ms and nonzero.
cadence entry.recurrence Evaluation interval in ms.
width gcd(lookback, cadence) Pane length. Divides both.
count lookback / width Number of panes. Must be ≥ 2.
TimeShift.offset_ms shift.offset_ms + i·width How far back pane i ends. Positive = earlier. The original offset is added.
TimeShift.at original at @ anchor kept as is.
TimeRange.kind original kind Kept as is.

MaterializationCandidate

Field Type Meaning
frontier Vec<NodeId> (u64) Node IDs whose outputs are produced in precompute and stored. No two are ancestor/descendant. Empty = nothing stored, everything runs at query time.
realization Result<PhysicalASAPDAG, Error> The split for this frontier, or the reason it cannot be cut. Errors are kept per candidate.

compile_materialization_candidates

Item Type Meaning
dag &PostAsapDAG Wire DAG to lower.
inputs BTreeMap<NodeId, InputContract> Contract (schema, plan properties) for each external input node, e.g. each pane's raw TimeRange.
roots &[NodeId] Output nodes the query part must produce.
max_candidates usize Frontier budget. 0 or too many frontiers → Err.
return Result<Vec<MaterializationCandidate>, Error> Err only if lowering fails or the budget is exceeded.

PhysicalASAPDAG (existing)

Field Type Meaning
precompute Option<CompiledPhysicalDAG> Producer part. None for the empty frontier.
query CompiledPhysicalDAG Reader part, run per evaluation.
materialized_outputs BTreeMap<NodeId, InputContract> Contracts for stored frontier outputs that query reads.

Examples

End-to-end: #509 Example 3/4 Pattern B (crates/integration-tests/tests/automatic_window_composition.rs).

Input: SummaryAgg(KLL k=200) over TimeRange(300 s) over Scan(TimeSeries "events"); recurrence FixedInterval(60_000).

  1. enumerate_window_compositions(&state, &entry, 64) returns 2 candidates: the original, and

    SummaryMerge
    ├── SummaryAgg(KLL) ← TimeRange(60 s) ← TimeShift(offset 0)       ← Scan(events)
    ├── SummaryAgg(KLL) ← TimeRange(60 s) ← TimeShift(offset 60 000)  ← Scan(events)
    ├── …                                    offset 120 000, 180 000
    └── SummaryAgg(KLL) ← TimeRange(60 s) ← TimeShift(offset 240 000) ← Scan(events)
    
  2. A SummaryEstimate { Quantile { q: 0.99 } } is put on top and compiled to wire form. Each pane's TimeRange gets 20 rows: pane p gets values p·20 … p·20+19, so the panes together hold 0 … 99.

  3. compile_materialization_candidates(.., 1024) succeeds. With max_candidates = 1 it returns Err.

  4. Empty frontier (B2: rebuild all panes per query): query result 98.0.

  5. Frontier = the five pane SummaryAgg nodes (B1/B3: store panes): precompute runs in Scope::Ingestion { 0 .. 300_000, revision 1 }, its outputs feed query in Scope::Query { evaluation_time_ms: 300_000, revision 1 }. Result: 98.0, the same as step 4.

Pane width for other cadences (from width = gcd(lookback, cadence)):

Lookback Cadence Width Panes Result
5 min 1 min 1 min 5 merge of 5
5 min 2 min 1 min 5 merge of 5 (non-divisible, GCD)
10 min 5 min 5 min 2 merge of 2
1 min 5 min 1 min 1 not composed (count <= 1)
5 min none (OneTime) — — original only

Accepted vs rejected (unit tests in window_composition.rs):

Case Result
five_minute_window_has_original_and_five_exact_panes: 5 min / 60 s, limit = 16 2 variants; 5 panes, each 60 s, offsets 0, 60 000, …, 240 000; merge passes validate_structure
unknown_cadence_and_budget_are_explicit: same, limit = 4 BudgetExceeded (5 panes > 4)
same, recurrence OneTime { invocations: 1 }, limit = 4 1 variant (original)
SummaryAgg over a non-time-series source, or over other operators between TimeRange and Scan not a site, original kept

Out of scope

  • Sliding windows and Exponential Histograms from the window-composition rule.
  • A rotating retained-pane cache across evaluations, and historical approximate EH boundary certificates. The deployment must still give each pane's raw input the correct window and bind retained outputs to that window/revision.
  • Checking that the cadence matches a catalog's existing pane origin. FixedIntervalAt's phase is not used.
  • Costing and whole-workload selection: Select complete workload candidates with exhaustive sharing and lifecycles #568. Temporal merge timestamp binding: Preserve evaluation timestamps in native temporal pane merges #569.
  • Docs: docs/develop_docs/planner-layering-status.md updates the "Window composition" row and adds an "Automatic window/materialization acceptance" section.

Stack and validation

Stacked on #565 (stack/509-20-sql-entropy-fallback). Head: stack/509-21-window-composition. Next: #568 (complete workload selection), then #569 (temporal merge timestamp binding). Implements the exact window/materialization portion of #509.

Validation (from the current PR body): focused window tests, the native end-to-end direct/materialized quantile fixture, targeted all-feature Clippy, formatting and diff checks pass.

🤖 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

zzylol added a commit that referenced this pull request Oct 4, 2026
A tumbling pane (#580) reads TimeRange(w) over TimeShift(i*w) over Scan.
The per-entity build path required the Scan directly under the TimeRange
and rejected panes with "per-entity summary requires a resolved source".
The deployment supplies the shifted raw rows, so look through one
TimeShift to find the source schema. Ported from #566.

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