diff --git a/docs/design_docs/concepts/post-asap-ir.md b/docs/design_docs/concepts/post-asap-ir.md index 6c9aa1461..08c670af4 100644 --- a/docs/design_docs/concepts/post-asap-ir.md +++ b/docs/design_docs/concepts/post-asap-ir.md @@ -49,7 +49,8 @@ summary family supports incremental maintenance. and selects the joined rows. Completeness evidence belongs to pruning, not ranking. A `SummaryNode` carries its expression, schema and optional result guarantee. -State and query values have different contracts. Exact operations over +State and query values have different contracts; see +[Schema and physical data for ASAP primitives](../proposals/asap-primitive-schema.md). Exact operations over approximate readouts still require composed accuracy guarantees. See the [accuracy implementation companion](../../develop_docs/end-to-end-accuracy-guarantees.md) and [physical-plan integration](../architecture/physical-plan-integration.md) diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index 274e4974c..00caf5d51 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -117,26 +117,14 @@ operators from logical candidates has not completed this integration. ### Input semantics and summary semantics -`source`, `filter`, `grouping` and `window` describe input-data semantics: -where records originate, which records qualify, how they are grouped and which -time interval applies. They are not a complete description of arbitrary summary -computation. In particular, the same four fields can summarize different value -expressions or produce different states. - -| Concern | Required semantic information | -| --- | --- | -| Input computation | Source identities and schemas, filters, joins/transforms and their order, or a reference to the canonical input sub-DAG | -| Values and grouping | Value expressions, item identities and weights where applicable, group keys and types, and operation-defined null/duplicate handling | -| Time | Time column and interpretation, interval bounds, evaluation alignment, and distinction between query range and maintained panes | -| Summary computation | Exact operation or sketch family, algorithm and parameters, and supported build/merge behavior | -| Output | State versus finalized value, output schema/type, and readout parameters when part of the output computation | - -For example, KLL over `latency_seconds` and KLL over `log(latency_seconds)` differ -even with identical source, filter, grouping and window. Likewise, weighted -frequency state needs both item and weight expressions. More complex inputs -must retain their computation DAG; four descriptive fields cannot replace it. - -The canonical selected computation is authoritative. These categories describe +`source`, `filter`, `grouping` and `window` describe input-data semantics but not +a complete summary computation: the same four fields can summarize different +value expressions or produce different states. The semantic information a summary +depends on, and where the IR records each part (field type, producing operator, +or coverage), is specified in +[Schema and physical data for ASAP primitives](proposals/asap-primitive-schema.md#4-proposed-node-field-design). + +The canonical selected computation is authoritative. Those categories describe what must be preserved, not a new flat IR or a second expression language. Operator-defined behavior should be referenced through its canonical contract, not independently configured in deployment metadata. Unsupported or unresolved diff --git a/docs/design_docs/proposals/README.md b/docs/design_docs/proposals/README.md index cb44fadbb..6c4ed49f8 100644 --- a/docs/design_docs/proposals/README.md +++ b/docs/design_docs/proposals/README.md @@ -11,3 +11,4 @@ extensions. A design document is not a promise of downstream runtime support. - [Operator sharing](operator-sharing.md) - [Decoupling operators from scalar expressions](decoupling_op_and_expr.md) - [ASAPPlanner layering](planner-layering.md) +- [Schema and physical data for ASAP primitives](asap-primitive-schema.md) diff --git a/docs/design_docs/proposals/asap-primitive-schema.md b/docs/design_docs/proposals/asap-primitive-schema.md new file mode 100644 index 000000000..c3610ada5 --- /dev/null +++ b/docs/design_docs/proposals/asap-primitive-schema.md @@ -0,0 +1,585 @@ +# Schema and Physical Data for ASAP Primitives + +This document is the single source of truth for the schema, and column design for ASAP Primitives. This is used in the logical stage (LogicalASAPDAG), and physical stage (PhysicalASAPDAG). + +## 1. Goal, problem, and requirements + +Unlike existing Database engines, which work on raw data or explicitly defined materialized tables with schema and column names provided by the users, ASAPPlanner is designed for querying and execution over the mix of raw data and ASAP Primitives. ASAP primitives are usually compact summaries over raw data. Therefore, it introduces new requirement when we design the schema and node definitions for LogicalASAPDAG and PhysicalASAPDAG. + +Assuming we have the Logical DAG defined for a canonicalized representation for a batch of queries (the `LogicalDAG` produced by the frontends, [planning stages §0](planner-layering.md#0-language-specific-frontends); the stages that follow are in [planner-layering.md](planner-layering.md#stages)). +The LogicalASAPDAG will share/reuse the NonASAP operator and ScalarExpr nodes in LogicalDAG ([decoupling operators from scalar expressions](decoupling_op_and_expr.md), with the unified operator type in [operator sharing §1.1](operator-sharing.md#11-unified-operator-type)), but replacing some operators in LogicalDAG with the operators operated with ASAP Primitives. The complete list is `ASAPOp` in `crates/types/src/ir/asap.rs`: + +| Operator | Input → output | Status | +|---|---|---| +| `SummaryAgg` | values → summary state (creates and updates the state; `SummaryUpdate` is its input mapping, not a separate operator) | implemented | +| `SummaryEstimate` | sketch state → value | implemented | +| `FinalizeExactAccumulator` | exact accumulator state → value | implemented | +| `MaintainPopulation` | values → maintained membership (state) | implemented | +| `EvaluatePopulation` | maintained membership → value | implemented | +| `SummaryMerge` | state × N → state | structure in #560, coverage check in #646 | +| `SummarySubtract` | state × state → state | reserved | +| `SummaryDelete` | state → state without one key | reserved | +| `SummaryJoin` | state × state → state | reserved | +| `Extension` | state → state, named by an extension | reserved | + +§5 walks through each of them. +Each of the Summary operators also require the ASAP primitive information above to inter-operate correctly, preserving semantic correctness. + +Basically, the following information should be represented to preserve the equivalent query semantics when we introduce ASAP Primitives to logical query representation, and following physical one. + +- What type of the ASAP Primitive is +- What is the ASAP Primitive parameters +- What data sources a ASAP primitive summarizes +- What query intent the summarized ASAP Primitive can support, e.g., statistical aggregation intents, time window aggregation intents + + + +And these information will be combined with relational or time series query operator information, such as group by/reduction, filtering, projection, join, time series selection, together. + +Therefore, these requirements drive the following schema and metadata, node information, and column design. + + + +## 2. Existing database terminology for schema, table, column, and physical data layout + +This section fixes the words used below. They follow relational databases and Apache Arrow / DataFusion, which ASAPPlanner's frontend already uses. + +| Term | Meaning in existing systems | In ASAPPlanner | +|---|---|---| +| **Relation / table** | A set (bag) of rows with the same columns. A base table is stored; a derived relation is the output of a query operator. | Every edge in the DAG carries a relation. A `Scan` reads a base table (SQL table or PromQL metric); every other operator outputs a derived relation. | +| **Row / tuple** | One element of a relation: one value per column. | One output row of a node. For PromQL, one sample of one series at one time. | +| **Column** | One position in every row, with a name and a type. Qualified as `table.column` when names can collide (DataFusion `Column { relation, name }`). | `ColumnId` refers to a column of the input schema; `(table, name)` identifies it across nodes (§4.3). | +| **Schema** | The ordered list of columns of a relation: name, data type, nullability (Arrow `Schema` of `Field { name, data_type, nullable }`; DataFusion `DFSchema` adds the table qualifier). The schema is *metadata*: it describes rows, it contains none. | `Schema` of `Field { name, dtype, nullable, table }` in `crates/types/src/pre_asap/schema.rs`. Unlike Arrow, `dtype` can be a summary state type (§3). | +| **Data type** | The type of a column's values (`Int64`, `Utf8`, `Timestamp`, …). | `DataType`, wrapped as `FieldDataType::Plain`. | +| **Aggregate state** | The intermediate value of an aggregate function before its final result, e.g. `(sum, count)` for `AVG` (DataFusion `Accumulator::state`, partial/final aggregation). It is never exposed as a column type to users. | Summary state *is* a column type here (`FieldDataType::Sketch`, `ExactAggregate`, …), so state can flow along edges and be merged, stored and read by later operators. | +| **View / materialized view** | A view is a named query (its *definition*). A materialized view also stores the query's result rows; a query can then be answered from it when its definition matches (view matching, §4.1). | A built summary state is a materialized aggregation view whose aggregate is a summary family. Its definition and which rows it took are its coverage (§4). | +| **Physical data layout** | How rows are stored: row-oriented or columnar (Arrow `RecordBatch`: one array per column), split into partitions (hash or range) and batches. | Decided in physical planning ([planning stages §2](planner-layering.md#2-physical-asap-aware-optimization)) and by the executing backend. The logical schema does not depend on it. | + +Two consequences for the design: + +- A schema says what *kind* of values flow along an edge, never *which* rows. Which rows a relation contains is decided by the operators below it (its definition). This is why coverage is a node property and not part of the schema (§4). +- Existing systems keep aggregate state internal to one operator. ASAPPlanner makes it a first-class column type so that one state can be shared, merged and stored across queries, which is what §3 and §4 add. + +## 3. Proposed schema design +Schema represents the **metadata** of information flow along an **edge** between two nodes in a logical or physical DAG. The schema field is associated with the node in the DAG. The consumer of the node in the DAG takes the schema from the producer node as input. + +Schema definition here is shared between LogicalDAG, LogicalASAPDAG, and PhysicalASAPDAG. The schema contain fields, and each field is mapping to a column in the physical data representation. +Based on our requirement, each field should contain the following information. +1. **What type of the ASAP Primitive is** A state column can be a raw data type (e.g., numerical number, string). It can also be a [summary type](#61-schema-and-field-types-cratestypessrcpre_asapschemars) (`FieldDataType` in `crates/types/src/pre_asap/schema.rs`), e.g., the summary family is sketch, and the sketch type is quantile KLL sketch algorithm, and KLL sketch has K as parameter as the schema. In the code these are: **family** = the `FieldDataType` variant (`ExactAggregate`, `Sketch`, `Sample`, `Wavelet`, `StatModel`; `Plain` is a raw value); for sketches, **category** = `SketchCategory` (`Quantile`, `Frequency`, `Cardinality`, `TopK`, `Universal`), **algorithm** = `SketchAlgorithm` (`Kll`, `Cms`, `Hll`, …) and **parameters** = `SketchParams` (`Kll { k }`), bundled as `SketchKind` ([§6.2](#62-state-family-parameters-cratestypessrcpost_asapsketchrs)); a sketch also carries its `GroupingStrategy`. So the example is `Sketch(SketchKind { Quantile, Kll, Kll { k: 200 } }, …)`. It has a family, an algorithm and parameters. +2. **What query intent the summarized ASAP Primitive can support, e.g., statistical aggregation intents, time window aggregation intents** This information is being mapped based on the primitive type. + + +## 4. Proposed Node field design + +A node in the physical data will represent the data or summary instance, so a node has a field for **What data sources a ASAP primitive summarizes**. + +**Why coverage is not part of the schema.** Two summary states worth merging always cover different data. `SummaryMerge` requires all inputs to have the same schema; that check is how it knows they are the same kind of state (same sketch, parameters and grouping). For example, two KLL states for "latency by job", built from minute 0–1 and minute 1–2: + +| | State A | State B | Equal? | +|---|---|---|---| +| schema | `(job: Utf8, state: KLL{k=200})` | `(job: Utf8, state: KLL{k=200})` | yes, so the merge is allowed | +| what it summarizes | time `[0,1)` | time `[1,2)` | no, which is why merging them is useful | + +If what a state summarizes were part of the schema, these two schemas would differ and the merge would be rejected; the only merge left would be a state with an exact copy of itself, which counts every observation twice. So the schema says *what kind of state* this is, and coverage says *which data it was built from*. + +A summary state summarizes the result of a whole computation, not a few columns of a raw table. A KLL over `rate(requests_total[5m])` summarizes rate outputs, and a KLL over a join summarizes join rows. So coverage describes the state by the sub-DAG below it, split into the part that says *what is computed* and the part that says *which of its rows were taken*. + +### 4.1 Design basis: view matching (Goldstein & Larson) + +The design follows the view matching algorithm of Goldstein and Larson, which decides when a query can be answered from a materialized select-project-join-group-by (SPJG) view: + +> J. Goldstein and P.-Å. Larson. *Optimizing Queries Using Materialized Views: A Practical, Scalable Solution.* SIGMOD 2001. + +The algorithm splits a view's `WHERE` into column equivalence classes, a **range** per column and **residual** predicates. A view can answer a query when the residuals match, the query's ranges lie inside the view's (§3.1.2), the columns needed by compensating predicates are in the view output (§3.3, requirement 2), and the query's `GROUP BY` is a subset of the view's, so the query's groups are further aggregations of the view's groups (§3.3, requirement 3). The SPJ part is implemented for DataFusion in [`datafusion-contrib/datafusion-materialized-views`](https://github.com/datafusion-contrib/datafusion-materialized-views), `src/rewrite/normal_form.rs` (`SpjNormalForm`, `Predicate { eq_classes, ranges_by_equivalence_class, residuals }`); it rejects `Aggregate` and `Join` input plans. + +A summary state is an aggregation view whose aggregate is a summary family. The mapping is: + +| Goldstein & Larson | Summary coverage | +|---|---| +| SPJ part: tables, joins, residual predicates | the computation `C` below the `SummaryAgg` (§4.2), part of `definition` | +| aggregate function and its argument | `family` and `input` (`SummaryUpdate`) of the `SummaryAgg`, part of `definition` | +| `GROUP BY` | the `SummaryAgg` reduction `G`, part of `definition` | +| ranges per column | `selection` (§4.3) | +| compensating predicate on view output | slice on a column of `G` only (§4.4) | +| query `GROUP BY` ⊆ view `GROUP BY` | rollup (§4.4) | + +What this design adds beyond the paper: + +- **Unions of states.** The paper considers single-view substitutes and notes that requirement 1 "is not required if substitutes containing unions of views are considered" (§3.1). `SummaryMerge` is exactly such a union, so it needs a disjointness check the paper does not have. +- **Summary families.** The paper allows `SUM` and `COUNT_BIG` only. Here each family declares how the selections of its inputs may relate (§4.4). +- **Value sets and hash partitions** next to ranges, and **evaluation-relative time** (§4.3). + +### 4.2 Coverage = definition + selection + +A state built by `SummaryAgg` means + +```text +state_g = family( input( σ( C ) ) ) for each group value g of G, restricted to G = g +``` + +- `C` is the child sub-DAG with the selection removed. Its output rows are the contributions. +- `σ` is the selection: which output rows of `C` went into the state. + +Coverage stores exactly these two things: + +- **`definition`**: the `SummaryAgg` node itself, with the selection removed from its child sub-DAG. It carries `C`, `input`, `family` and `G`. It is what the state *means*. +- **`selection`**: a union of boxes over the output columns of `C`. It is *which rows* the state took. + +If two states have the same `definition`, their contributions come from the same rows of the same computation, whatever `C` contains (join, union, `rate`, dedup). Disjoint selections then cannot share a row, so no observation is counted twice. No per-operator occurrence rule is needed. + +Examples of what ends up where: + +| Sub-DAG below `SummaryAgg` | `definition` keeps | `selection` takes | +|---|---|---| +| `Filter(region = 'us', Scan t)` | `Scan t` | `region ∈ {us}` | +| `Filter(latency < 100, Scan t)` | `Scan t` | `latency ∈ (−∞, 100)` | +| `TimeRange(1m, TimeShift(2m, Scan m))` | `Scan m` | time `(−3m, −2m]` relative to evaluation | +| `Filter(rate > 0, rate(TimeRange(5m, Scan m)))` | `rate(TimeRange(5m, Scan m))`, `rate > 0` as residual | — | +| `Filter(job = 'api', rate(TimeRange(5m, Scan m)))` | `rate(TimeRange(5m, Scan m))` | `job ∈ {api}` | + +In the last two rows the `TimeRange(5m)` stays in `definition`: it sits below `rate` and changes the rate values, so it is not a selection of output rows. The 5-minute read window is a source dependency, not coverage (§4.6). + +### 4.3 Deriving the selection + +Coverage is derived from the node, never declared. Walking down from the `SummaryAgg` (its own `filter` included), a predicate conjunct goes into `selection` when both hold: + +1. **It can be lifted to the `SummaryAgg`.** Lifting is the inverse of DataFusion's `PushDownFilter` (`datafusion-optimizer`, `push_down_filter.rs`): a predicate passes `Filter`, `TimeRange`/`TimeShift` and a direct-column `Project` (renaming the column); passes an `Aggregate` only when every column it uses is a group column; and passes a window function or a per-series temporal function such as `rate` only when every column it uses is a partition column (a series label). Anything else stops it. +2. **It is one of the box constraints.** Per column, one of: + - **value set**: `In` or `NotIn` a set of literals, from `=`, `!=`, `IN`, `NOT IN` and `OR` of equalities (as DataFusion's `LiteralGuarantee` extracts them); + - **interval**: lower and upper `std::ops::Bound` (`Included`, `Excluded` or `Unbounded`) from comparisons (as DataFusion's `Interval`); + - **hash partition**: `hash(columns) mod n = k`. + +A conjunct that fails either rule stays in `definition` as a residual, as in Goldstein & Larson. Column equalities (`a = b`) are residuals too: there are no column equivalence classes. + +Columns are identified by lineage `(table, name)`, the identity `ColumnRef::Qualified` uses, not by `Field.name`. So `shipping.region` and `billing.region` stay different columns, and a direct alias keeps the identity of the column it renames. A column whose `(table, name)` is not unique in the output (two items aliased `k`) cannot be named, so its conjuncts stay residual. Value sets compare literals by type: `1` and `1.0` are never proven different. + +**Time** is a selection like any other: + +- **Absolute** time needs nothing special: it is an interval on the timestamp column (the schema's `time_index`), for example `ts >= t0 AND ts < t1` gives `(Included(t0), Excluded(t1))` on `ts`. The IR has no timestamp literal yet, so such SQL filters stay residual until it does. +- **Relative** time has its own field, because it is not a column value: a `TimeRange(w)` over a `TimeShift(s)` on the lifted chain gives `(Excluded(−(s+w)), Included(−s))` relative to evaluation. PromQL ranges are left-open, matching the executor (`series_window.rs`). This is how Stage 2 tumbling panes are built (`window_composition.rs` in #601), so their time is derived rather than declared. Time is lifted only from a single range `TimeRange`; an instant `TimeRange` picks the latest sample per series, which is not a selection of rows, so it stays in `definition`. + +Relative time and a timestamp-column interval are different dimensions, so they are never compared: two states restricted only by different kinds of time are treated as possibly overlapping. Binding a relative pane to absolute timestamps for one evaluation (evaluation time plus the pane layout's phase) is a runtime coordinate, not coverage. + +### 4.4 Operations + +| Operation | Example | Valid when | +|---|---|---| +| merge (`SummaryMerge`, same `G`) | `[0,1m)` ⊕ `[1m,2m)`; `region='us'` ⊕ `region='eu'` | all `definition`s equal; selections related as the family requires (below) | +| rollup (`SummaryMerge` with `group_by: G'`) | `by[region, job]` → `by[job]` | `G'` ⊆ `G` and the family merges. Groups of one state are disjoint because a row has one value per group column, so no selection check is needed | +| slice | `by[region, job]` state answering `region = 'us' … by[job]` | the restricted columns are all in `G`. A sketch cannot be filtered, so a restriction on any other column is invalid | +| reuse for a query | a stored state answers a query | Goldstein & Larson containment: same `definition`, query selection inside the state's, any compensating restriction is a slice | +| subtract (`SummarySubtract`, reserved) | `[0,10) − [0,5)` | same `definition`; the right selection is contained in the left | + +One `SummaryMerge { children, group_by }` covers both merge and rollup: one child with a coarser `group_by` is a rollup, and `group_by` equal to the children's is a plain merge. Its coverage is the children's `definition` with `G'` and the union of their selections; adjacent intervals are joined, gaps stay as separate boxes. + +How selections must relate is declared by the family, next to whether it merges (`FieldDataType::family_merges`, added in #592): + +- **disjoint** for counting families (KLL, Count-Min, exact `Sum`/`Count`): an overlapping row would be counted twice; +- **overlap allowed** for idempotent families (HLL, exact `Min`/`Max`, distinct sets); +- **contained** for subtraction. + +Two `definition`s are equal when their canonical forms (`canonicalize`) are structurally equal, ignoring planning metadata: `timing`, `guarantee` and `coverage_cache`. `SummaryUpdate.weight_domain` is compared: it is derived from `C` and `input`, so it differs only if a derivation is wrong. A state built at ingestion time and one built at query time can therefore merge. + +States over different sources have different `definition`s and do not merge. To combine tables, put `UNION ALL` with a marker column below one `SummaryAgg`; the marker is then an ordinary column for `selection` or `G`. + +### 4.5 Interface + +```rust +pub struct OperatorNode { + pub operator: Operator, + pub result_kind: OperatorResultKind, + pub schema: Schema, + pub guarantee: Option, + pub timing: Option, + /// Cache for `coverage()`. Lazily filled, never serialized, ignored by + /// equality, emptied on clone. Not a source of truth: coverage is always + /// re-derivable. + coverage_cache: CoverageCache, +} + +impl OperatorNode { + /// `Some` for a `SummaryAgg` and a valid `SummaryMerge`. + pub fn coverage(&self) -> Option<&SummaryCoverage>; +} + +pub struct SummaryCoverage { + /// The `SummaryAgg` (or rolled-up equivalent) with the selection removed. + pub definition: Rc, + /// Union of boxes over the output rows of the definition's computation. + pub selection: Vec, +} + +pub struct SelectionBox { + pub columns: BTreeMap, // missing column = unrestricted + pub relative_time: Option<(Bound, Bound)>, // ms from evaluation; None = unrestricted +} + +pub struct ColumnIdentity { + pub table: Option, + pub name: String, +} + +pub enum Constraint { + In(Vec), // ScalarValue has no total order (Float64) + NotIn(Vec), + Interval { lower: Bound, upper: Bound }, + // HashPartition { columns, of, index }: added with its first producer. +} + +impl SummaryCoverage { + pub fn derive(node: &OperatorNode) -> Result; +} +``` + +`OperatorNode::new` still rejects an invalid `SummaryMerge` (different definitions, or selections the family does not allow), but it does not store the result. A `SummaryAgg` always has coverage: what cannot go into `selection` stays in `definition`. + +### 4.6 What coverage does not contain + +- **Source dependencies**: which source rows must be read to compute the contributions, such as the 5-minute window under `rate`. This is read planning and maintenance (compare `materialized/dependencies.rs` in `datafusion-materialized-views`). +- **Readiness and completeness**: whether a stored instance holds all of its rows. +- **Absolute binding** of relative time, and deployment identity. + +These belong to ASAPQuery-backend. The SDS split matches coverage: `SummaryDefinition` stores the serialized `definition` (Planner provides its serde; Backend owns the format version, definition id and hash), and a `StoredSummary`'s coordinates are the `selection` bound to one evaluation plus the group value. + +## 5. Examples on how OperatorNode, schema, and physical data information are being used with Summary operators + +Given that these information requirements are introduced by summary operators to work correctly semantically, we show the examples of how the defined OperatorNode, schema, and physical data information work with each kind of summary operators. + +Notation: an edge is written `──Kind(field Type, …)──▶`. Schemas are the ones `output_schema()` derives. Planning may rename fields through `OperatorNode::with_schema`, but types, nullability, `time_index`, `unique_keys` and `closed` must match the derivation. All examples use a table source, so values are `Relation`; with a `TimeSeries` source the value side is `InstantVector`. + +### 5.1 `SummaryAgg`: values → state + +Scenario: p99 latency by job, from KLL(k=200), over one minute of table `t`, US rows only. + +```text +Scan(t: job Utf8, region Utf8, ts Timestamp [time_index], latency Float64) + ──Relation(job Utf8, region Utf8, ts Timestamp, latency Float64)──▶ +Filter(region = 'us' AND ts >= 0 AND ts < 60_000) + ──Relation(job Utf8, region Utf8, ts Timestamp, latency Float64)──▶ +SummaryAgg(family = Sketch(KLL{k=200}, PerSubpopulationInstance), + input = SummaryUpdate::column(Named("latency")), reduction = by[job], + grouping = PerSubpopulationInstance, filter = None) + ──State(job Utf8, state Sketch(KLL{k=200}, PerSubpopulationInstance))──▶ + coverage() = { definition: this SummaryAgg over Scan(t) (the Filter removed), + selection: [{ columns: { t.region: In{'us'} }, + time: Absolute [Included(0), Excluded(60_000)) }] } +``` + +- Output schema: the `by` keys followed by one non-nullable field `state` typed `family`; `unique_keys = [[0]]`, `closed = true`, no `time_index`. With `Reduction::PerEntity` the input columns are kept and the sample-value column is replaced by `state`. +- Checks: `family` is not `Plain`; the child is not `State`; the `weight`/`item` columns resolve against the child schema; `filter`, if present, types as `Bool`. +- Coverage: **always derived**, never declared (§4.3). Both conjuncts of the `Filter` lift into `selection`, so `definition` is this node over the bare `Scan`. A KLL over `latency` for `region = 'eu'` has the same `definition` and a disjoint selection, so the two can merge. A conjunct that cannot lift (say `latency * 2 > 10`) stays in `definition` as a residual; the node still has coverage. +- Boundary: this is where values become state. The sketch family, algorithm and parameters are committed in the field type, and `guarantee` stays `None` because state is not a caller-visible value. + +### 5.2 `SummaryEstimate`: sketch state → value + +Scenario: read p99 from the state in 5.1. + +```text +──State(job Utf8, state Sketch(KLL{k=200}))──▶ +SummaryEstimate(query = SketchStatistic::Quantile { q: 0.99 }) + ──Relation(job Utf8, quantile Float64)──▶ (planner may rename to p99) +``` + +- Output schema: the input schema with the one non-plain field replaced by a non-nullable plain field. Its name and type come from the statistic: `quantile`/`frequency_l2`/`frequency_entropy` Float64, `cardinality`/`count` Int64 (Float64 if the producer is a `PerEntity` `SummaryAgg`). Keys and metadata pass through. A top-k readout is the exception: it returns the selected rows, one per ranked item, with the partition keys, the item identity columns, and a `value` Float64 score (#579). This is the same row shape as an exact Sort → Limit top-k, so the plans for one query share a root schema. +- Result kind: the value kind of the source the state was built from (`Relation` here). +- Checks: input is `State` with exactly one non-plain field, that field is `Sketch`, and its category accepts the statistic (§3). For example, `Cardinality` on KLL is rejected. +- Coverage: **none**. The output is a value; `coverage()` returns `None`. +- Boundary: state is consumed and a value is produced; `guarantee` on this node carries the readout's error bound. + +### 5.3 `FinalizeExactAccumulator`: exact state → value + +Scenario: total bytes by host with an exact Sum accumulator. + +```text +Scan(t: host Utf8, bytes Float64) + ──Relation(host Utf8, bytes Float64)──▶ +SummaryAgg(family = ExactAggregate(Sum, Sum), input = column(Named("bytes")), reduction = by[host]) + ──State(host Utf8, state ExactAggregate(Sum, Sum))──▶ coverage: derived +FinalizeExactAccumulator + ──Relation(host Utf8, state Float64)──▶ +``` + +- Output schema: each `ExactAggregate` field keeps its name (`state`) and takes the type and nullability the equivalent `NonASAPOp::Aggregate` would give: Sum/Min/Max follow the input column, Count is Int64, and Rate/IRate/Increase are Float64. If the child is not a `SummaryAgg` directly, Count falls back to Int64 and the others to Float64. `unique_keys`, `closed` and `time_index` are preserved (`schema_rebuilding.rs`). +- Checks: the input is `State` and contains an `ExactAggregate` field; a sketch is rejected (`structure_contract.rs`). +- Coverage: **none** on the output. +- Boundary: this is the explicit maintenance-to-read boundary for exact state. Exact state is never read through `SummaryEstimate`. + +### 5.4 `MaintainPopulation`: values → maintained membership (state) + +Scenario: keep the full latency population per job, so that p99 and top-10 can be evaluated later. + +```text +Scan(t: job Utf8, latency Float64) [closed schema] + ──Relation(job Utf8, latency Float64)──▶ +MaintainPopulation(population = MaintainedPopulation { + input: PopulationInput::Rows { input: , value_column: 1, grouping: by[job] }, + max_k: 10, quantiles: true }) + ──State(job Utf8, latency Float64)──▶ +``` + +- Output schema: identical to the child's, all plain. Only `result_kind = State` marks it as maintained state. +- Checks: `population.matches_node(child)`. For `Rows`, the child must be the same closed table `Scan`, the value column must be non-null Float64, and grouping must be `by` with in-range keys. For `CurrentSeries`, it must be a `TimeSeries` scan with the same metric, matchers and grouping labels, under an instant `TimeRange` of `lookback_ms` (which may be omitted only for the default 300 s lookback). +- Coverage: **none**. Maintained membership is not combined by `SummaryMerge`. If maintained populations are later materialized per pane, they derive coverage the same way as `SummaryAgg`. +- Boundary: the output is state because it must also track membership changes; downstream operators can only read it through `EvaluatePopulation`. + +### 5.5 `EvaluatePopulation`: maintained membership → value + +Scenario: p99 by job from the population in 5.4. + +```text +──State(job Utf8, latency Float64) [from MaintainPopulation]──▶ +EvaluatePopulation(evaluation = PopulationStatistic::Quantile { q: 0.99 }) + ──Relation(job Utf8, quantile_0_99 Float64)──▶ +``` + +- Output schema: the schema of `Aggregate(by grouping, measure)` over the maintained source. Quantile gives `quantile_` Float64, Sum gives `sum` (value type), Count gives `count` Int64 and Average gives `avg` Float64; `unique_keys = [[0]]`, `closed`. `TopK { k }` instead returns the source schema unchanged (the selected rows). +- Checks: the child is a `MaintainPopulation` node whose `supports(evaluation)` holds: `quantiles` must be set for `Quantile`, and `k <= max_k` for `TopK`. +- Coverage: **none**. +- Boundary: maintained membership is read as a value; the result kind is the source's (`Relation`). + +### 5.6 `SummaryMerge`: state × N → state (merge and rollup) + +On `main`, `SummaryMerge { children }` is **reserved**: `is_unimplemented()` returns true, and `output_schema()` and `validate_inputs()` return `UNIMPLEMENTED_ASAP_OP`, so `OperatorNode::new` fails. #560 enables it for children with identical schemas. This design adds `group_by`, so one operator does both merge and rollup (§4.4): + +```rust +SummaryMerge { children: Vec, group_by: Reduction } +``` + +Scenario A, time panes: two one-minute KLL panes of PromQL `quantile_over_time(0.99, m[2m])` merged into the two-minute state. Each pane reads `TimeRange(1m)` over `TimeShift(s)` over the scan, as Stage 2 builds them. + +```text +pane 0 = SummaryAgg(KLL k=200, column(SampleValue), by[]) over TimeRange(1m, TimeShift(0, Scan m)) + coverage() = { definition: SummaryAgg(...) over Scan m, selection: [{ time: Relative (−1m, 0] }] } +pane 1 = SummaryAgg(KLL k=200, column(SampleValue), by[]) over TimeRange(1m, TimeShift(1m, Scan m)) + coverage() = { definition: same, selection: [{ time: Relative (−2m, −1m] }] } +SummaryMerge(children = [pane 0, pane 1], group_by = by[]) + ──State(state Sketch(KLL{k=200}))──▶ + coverage() = { definition: same, selection: [{ time: Relative (−2m, 0] }] } +``` + +Scenario B, populations: `KLL(latency) by[job]` for `region = 'us'` and for `region = 'eu'` (as in 5.1) merge into `selection: [{ t.region: In{'us', 'eu'} }]`. + +Scenario C, rollup: one `KLL(latency) by[region, job]` state merged with `group_by = by[job]`. Every job's state is the merge of that job's per-region states. The output's `definition` is the same `SummaryAgg` with `reduction = by[job]`, and `selection` is unchanged. + +- Output schema: the children's schema with the group key fields reduced to `group_by`. +- Checks: + - at least one child, every child is `State` with exactly one state field; + - all children have equal `definition`s (§4.4), so family, parameters, `input`, `C` and grouping match. Merging k=200 with k=300, KLL over `latency` with KLL over `size`, or states over different sources fails; + - `group_by` ⊆ the children's `G`, and the family merges; + - the children's selections relate as the family requires: disjoint for KLL, so pane 0 with pane 0 is rejected; overlap is allowed for HLL. +- Coverage: **derived**: the shared `definition` with `group_by`, and the union of the children's selections. Adjacent intervals join; gaps stay as separate boxes. Nested merges work because a child merge has coverage like any other summary node. +- Boundary: state in, state out. No value is produced until a readout. + +### 5.7 Reserved operators (not implemented) + +These variants exist so that plans can name them, but `output_schema()`/`validate_inputs()` return `UNIMPLEMENTED_ASAP_OP`, so no node can be built. `output_kind()` already returns `State` for each of them. The intended edge shapes below follow from their fields; none of them is implemented. + +| Operator | Fields | Intended edge shape | +|---|---|---| +| `SummarySubtract` | `left, right` | State × State → State: remove one window's contribution, e.g. [0,10) − [0,5). Same `definition`; the right selection must lie inside the left (§4.4) | +| `SummaryDelete` | `summary_input, key: ColumnId` | State → State with the entries for `key` removed | +| `SummaryJoin` | `outer, inner, key, family` | State × State → State typed `family` (`produced_state()` returns it), e.g. join-size estimation | +| `Extension` | `child, name` | deployment-named state operator | + +### 5.8 Summary + +| Operator | Input kind | Output kind | Output carries state | Coverage on output | Status | +|---|---|---|---|---|---| +| `SummaryAgg` | value (not `State`) | `State` | yes (one `family` field) | derived: itself minus selection, plus selection | implemented | +| `SummaryEstimate` | `State` (one `Sketch` field) | source's value kind | no | none | implemented | +| `FinalizeExactAccumulator` | `State` (`ExactAggregate`) | source's value kind | no | none | implemented | +| `MaintainPopulation` | `Relation` (table) / `InstantVector` (series) | `State` | yes (by kind; fields plain) | none | implemented | +| `EvaluatePopulation` | `State` from `MaintainPopulation` | source's value kind | no | none | implemented | +| `SummaryMerge` | `State` × N | `State` | yes | derived: shared definition with `group_by`, union of selections | reserved; enabled by #560, `group_by` added by this design | +| `SummarySubtract` | `State` × 2 | `State` | yes | derived: left selection minus right (planned) | reserved | +| `SummaryDelete` | `State` | `State` | yes | — | reserved | +| `SummaryJoin` | `State` × 2 | `State` | yes | — | reserved | +| `Extension` | any | `State` | yes | — | reserved | + +## 6. Key code interfaces + +`OperatorNode`, `OperatorResultKind` and coverage are in §4. Bodies and serde/derive attributes are elided below. + +### 6.1 Schema and field types (`crates/types/src/pre_asap/schema.rs`) + +```rust +pub type ColumnId = usize; + +pub struct Schema { + pub fields: Vec, + pub time_index: Option, // must point at a plain Timestamp field + pub unique_keys: Vec>, + pub closed: bool, // true = fields enumerate every column +} +impl Schema { + pub fn new(fields: Vec) -> Self; + pub fn with_time_index(fields: Vec, time_index: ColumnId, unique_keys: Vec>) -> Self; + pub fn lifted(fields: Vec, time_index: Option) -> Self; // closed = true + pub fn is_all_plain(&self) -> bool; + pub fn column_id(&self, name: &str) -> Option; + pub fn column_id_qualified(&self, table: &str, name: &str) -> Option; +} + +pub struct Field { + pub name: String, + pub dtype: T, + pub nullable: bool, + pub table: Option, +} +impl Field { + pub fn plain(name: impl Into, dtype: DataType, nullable: bool) -> Self; + pub fn plain_dtype(&self) -> Option<&DataType>; + pub fn is_plain(&self) -> bool; +} + +/// A column's type: a plain value, or summary state of one family. +pub enum FieldDataType { + Plain(DataType), + ExactAggregate(ExactKind, ExactParams), + Sketch(SketchKind, GroupingStrategy), + Sample(SamplingKind, SamplingParams), + Wavelet(WaveletKind, WaveletParams), + StatModel(StatModelKind, StatModelParams), +} + +pub enum DataType { + Null, Int64, Float64, Utf8, Bool, Timestamp, Interval, Date, + List { element: Box> }, + Struct { fields: Vec> }, + Map { key: Box, value: Box, value_nullable: bool }, +} +``` + +### 6.2 State-family parameters (`crates/types/src/post_asap/sketch.rs`) + +```rust +pub enum ExactKind { Sum, Count, Min, Max, Increase, Rate, IRate } +pub enum ExactParams { Sum, Count, Min, Max, Increase, Rate, IRate } // no knobs; mirrors kind + +pub struct SketchKind { category: SketchCategory, algorithm: SketchAlgorithm, params: SketchParams } +impl SketchKind { + /// The only constructor; classifies the category and panics on mismatched params. + pub fn new(algorithm: SketchAlgorithm, params: SketchParams) -> Self; + pub fn category(&self) -> SketchCategory; + pub fn algorithm(&self) -> &SketchAlgorithm; + pub fn params(&self) -> &SketchParams; +} +pub enum SketchCategory { Universal, Quantile, Cardinality, Frequency, TopK } +// Universal: UnivMon | Quantile: Kll, DDSketch | Cardinality: Hll, Theta, Kmv +// Frequency: Cms, CountSketch | TopK: CmsWithHeap, CountSketchWithHeap +pub enum SketchAlgorithm { UnivMon, Kll, Cms, Hll, DDSketch, CmsWithHeap, Kmv, Theta, CountSketch, CountSketchWithHeap } +pub enum SketchParams { + UnivMon { heap_size: u32, sketch_rows: u32, sketch_cols: u32, layers: u8 }, + Kll { k: u32 }, + Cms { width: u32, depth: u32 }, + Hll { precision: u8 }, + DDSketch { alpha: f64 }, + CmsWithHeap { width: u32, depth: u32, heap_size: u32 }, + Kmv { k: u32 }, + Theta { k: u32 }, + CountSketch { width: u32, depth: u32 }, + CountSketchWithHeap { width: u32, depth: u32, heap_size: u32 }, +} + +/// How grouped state is instantiated across `by` subpopulations. Orthogonal to family. +pub enum GroupingStrategy { + PerSubpopulationInstance, // Default + SharedMultiSubpopulation { kind: HydraKind, params: HydraParams }, +} +pub enum HydraKind { HydraKll /* experimental, no error bound */, HydraCms, HydraCountSketch } +pub enum HydraParams { + HydraKll { k: u32, shared_buckets: u32 }, + HydraCms { width: u32, depth: u32, shared_rows: u32, shared_columns: u32 }, + HydraCountSketch { width: u32, depth: u32, shared_rows: u32, shared_columns: u32 }, +} +pub fn hydra_kind_for(a: &SketchAlgorithm) -> Option; // Cms, CountSketch only + +pub enum SamplingKind { Reservoir } pub enum SamplingParams { Reservoir { size: u32 } } +pub enum WaveletKind { Haar } pub enum WaveletParams { Haar { coefficients: u32 } } +pub enum StatModelKind { Parametric } pub enum StatModelParams { Parametric { family: String } } +``` + +### 6.3 Update input and readouts (`post_asap/sketch.rs`, `post_asap/maintained_population.rs`) + +```rust +/// One state update: `item` keys the update for keyed families; `weight` is applied to state. +pub struct SummaryUpdate { + pub item: Option, + pub weight: SummaryInputExpr, + pub weight_domain: WeightDomain, // serde default: UnknownOrSigned +} +impl SummaryUpdate { pub fn column(c: ColumnRef) -> Self; } // item None, UnknownOrSigned +pub enum WeightDomain { + UnknownOrSigned, // Default; never assumed non-negative + NonNegative { proof: NonNegativeWeightProof }, +} +pub enum NonNegativeWeightProof { UnitCount, ResetAwareCounterDerivative } +pub enum SummaryInputExpr { + Constant(f64), Column(ColumnRef), Tuple(Vec), EntityIdentity(EntityIdentity), +} +pub enum EntityIdentity { PromqlLabelSet { excluding: Vec } } + +/// Readout of sketch state, carried by SummaryEstimate. +pub enum SketchStatistic { + FrequencyL2, FrequencyEntropy, + Quantile { q: f64 }, + PointCount { key: ColumnRef, value: Option }, + Cardinality, + TopK { k: usize }, +} + +/// Readout of a maintained population, carried by EvaluatePopulation. +pub enum PopulationStatistic { Quantile { q: f64 }, TopK { k: usize }, Sum, Count, Average } +pub struct MaintainedPopulation { + pub input: PopulationInput, + pub max_k: usize, // largest TopK it supports + pub quantiles: bool, // whether Quantile is supported +} +pub enum PopulationInput { + CurrentSeries(CurrentSeriesInput), // metric, matchers, grouping, without, lookback_ms + Rows { input: Rc, value_column: usize, grouping: GroupKeys }, +} +``` + +A finalized value's accuracy statement is `ResultGuarantee { metric, bound, failure_probability, provenance }` (`post_asap/guarantee.rs`). It is attached to readout and finalized nodes, never to raw state. + +### 6.4 ASAP operators (`crates/types/src/ir/asap.rs`) + +```rust +pub const UNIMPLEMENTED_ASAP_OP: &str = + "this ASAP operator is reserved: schema, accuracy, timing and export are not implemented"; + +pub enum ASAPOp { + SummaryAgg { + child: Rc, + family: FieldDataType, // never Plain + input: SummaryUpdate, + reduction: Reduction, // Reduce(GroupKeys) | PerEntity + grouping: GroupingStrategy, + filter: Option, // serde default None + }, + SummaryEstimate { summary_input: Rc, query: SketchStatistic }, + FinalizeExactAccumulator { child: Rc }, + MaintainPopulation { child: Rc, population: MaintainedPopulation }, + EvaluatePopulation { child: Rc, evaluation: PopulationStatistic }, + // Reserved on this branch; #560 implements SummaryMerge. + SummaryMerge { children: Vec> }, + SummarySubtract { left: Rc, right: Rc }, + SummaryDelete { summary_input: Rc, key: ColumnId }, + SummaryJoin { outer: Rc, inner: Rc, key: ColumnId, family: FieldDataType }, + Extension { child: Rc, name: String }, +} + +impl ASAPOp { + pub fn children(&self) -> Vec<&Rc>; // SummaryAgg includes its filter's subquery nodes + pub fn map_children(&self, f: impl FnMut(&Rc) -> Rc) -> Self; + pub fn kind_name(&self) -> &'static str; + /// Merge, Subtract, Delete, Join, Extension on this branch; #560 removes Merge. + pub fn is_unimplemented(&self) -> bool; + /// SummaryAgg/SummaryJoin `family`; #560 adds SummaryMerge (its inputs' state type). + pub fn produced_state(&self) -> Option<&FieldDataType>; + pub fn output_schema(&self) -> Result; + pub fn output_kind(&self) -> OperatorResultKind; + pub fn validate_inputs(&self) -> Result<(), SchemaDerivationError>; +} +``` diff --git a/docs/design_docs/proposals/decoupling_op_and_expr.md b/docs/design_docs/proposals/decoupling_op_and_expr.md index f87456a9f..962af6c5c 100644 --- a/docs/design_docs/proposals/decoupling_op_and_expr.md +++ b/docs/design_docs/proposals/decoupling_op_and_expr.md @@ -62,7 +62,7 @@ operator inputs and scalar query-result references use `Rc`. read by expressions. Keep `ScalarExpr::Column(ColumnId)`: the ID selects a field for type checking and the corresponding input value for evaluation, independently of the executor's row/column storage layout. See the -[fields versus column references contract](operator-sharing.md#21-one-schema-model-for-values-and-state). +[fields versus column references contract](asap-primitive-schema.md#3-proposed-schema-design). Names are resolved to `ColumnId` before constructing these nodes. Parsing and unresolved `ColumnRef` handling remain frontend concerns; no alternative generic diff --git a/docs/design_docs/proposals/operator-sharing.md b/docs/design_docs/proposals/operator-sharing.md index 1851d1809..c9ab193a4 100644 --- a/docs/design_docs/proposals/operator-sharing.md +++ b/docs/design_docs/proposals/operator-sharing.md @@ -365,155 +365,13 @@ caching or mutation mechanism. ### 2.1 One schema model for values and state -Use one `Schema` for operator outputs before and after optimization. Rename today's -`SummaryFamilyType` to `FieldDataType`: it types every field, and `Plain` is not a summary -family. Rename `Column` to `Field` and `Schema.columns` to `Schema.fields`: the struct -describes a column and holds none of its data. Retain the current `Schema` metadata. -The following is the proposed resolved interface; it is not the current Rust definition. - -```rust -struct Field { - name: String, - dtype: FieldDataType, - nullable: bool, - table: Option, -} - -struct Schema { - fields: Vec, - time_index: Option, - unique_keys: Vec>, - closed: bool, -} - -// Today's `SummaryFamilyType`, renamed; variants and payloads unchanged. -enum FieldDataType { - Plain(DataType), - ExactAggregate(ExactKind, ExactParams), - Sketch(SketchKind, GroupingStrategy), - Sample(SamplingKind, SamplingParams), - Wavelet(WaveletKind, WaveletParams), - StatModel(StatModelKind, StatModelParams), -} - -// Proposed derived output classification, separate from column types. -enum OperatorResultKind { - Relation, - InstantVector, - RangeVector, - State, -} - -impl Operator { - fn output_schema(&self) -> Result; - fn output_kind(&self) -> Result; - fn validate_inputs(&self) -> Result<(), QueryExprError>; -} - -impl OperatorNode { - fn validate_structure(&self) -> Result<(), QueryExprError>; - fn validate_execution_timing(&self) -> Result<(), QueryExprError>; -} - -impl ScalarExpr { - fn scalar_type(&self, input: &Schema) -> Result<(DataType, bool), QueryExprError>; -} -``` - -**Fields versus column references.** These names describe different roles, not -competing representations of the same object: - -| Name | Role | Holds runtime values? | -|---|---|---| -| `Schema` | Ordered `Field` metadata, plus key/time/closedness information | No | -| `Field` | Name, type, nullability and optional qualifier for one output column | No | -| `ColumnRef` | Unresolved logical reference: `Named`, `Qualified`, `SampleValue`, or `Wildcard` | No | -| `ColumnId = usize` | Resolved column position in a particular input/output schema | No | -| Runtime batch | Values conforming to a schema; storage layout is executor-specific | Yes | - -Keep `ColumnRef`, `ColumnId`, and `ScalarExpr::Column(ColumnId)`. Renaming the -metadata struct `Column` to `Field` does not rename column references to field -references. The same position identifies metadata during planning and values -during execution; it is not a stable field identity across projections or joins. -Schema `unique_keys` and `time_index` also use these column positions. - -For example, resolving `t.bytes` to position `1` produces `ColumnId = 1`. -`schema.fields[1]` supplies its type and nullability; evaluating -`ScalarExpr::Column(1)` reads the corresponding value. The native executor -currently reads `row[1]` from `Batch { schema, rows: Vec> }`. A columnar -executor would select array `1` instead. No physical `Column` container is -introduced by the metadata rename, and the old metadata `Column` struct is not -retained as a second type. - -**Relationship to current types.** `Field` is today's pre-ASAP `Column` with `dtype` -widened from `DataType` to `FieldDataType`. `FieldDataType` is today's `SummaryFamilyType` -under a name that also fits its `Plain` case. The proposed common `Schema` replaces -the separate operator-edge roles of pre-ASAP `Schema` and post-ASAP `SummarySchema` / -`SummaryField`; it does not rename `DataType`. A pre-ASAP value column becomes -`Plain(dtype)`. -Frontend validation permits only ordinary value columns, preserving the current -pre-ASAP restriction even though the common schema can also express state. - -| Field | Meaning and requirement | -|---|---| -| `fields` | Ordered named fields. `Plain(DataType)` is a readable value; other variants retain the identity and parameters of summary or exact-accumulator state. | -| `Field.nullable`, `Field.table` | Preserve SQL nullability and qualified column resolution. | -| `time_index` | Identifies the time column when present; it does not by itself distinguish an instant vector from a range vector. | -| `unique_keys` | Proven column combinations identifying rows; an empty list asserts no known key. Recompute these proofs when a rewrite changes identity. | -| `closed` | Whether `fields` completely describes the output. An open PromQL schema must retain unlisted labels through the existing complete-series-identity contract. | - -`OperatorResultKind` is derived from the operation and its inputs and retained as -`OperatorNode.result_kind`. `State` describes an output carrying unfinalized state; its -schema may also contain ordinary grouping keys. `SummaryEstimate`, -`FinalizeExactAccumulator` and other readouts derive the appropriate relation or -vector kind from their operation and input context. Matching numeric columns do -not make those kinds interchangeable. - -**Interface contracts.** `Operator::output_schema` and `output_kind` derive output -metadata from the payload and validated inputs. `validate_inputs` checks local -producer/consumer compatibility, such as vector inputs for `BinaryOp` or the -required state family for a summary readout. Scalar typing checks the input-kind -contract of `PromqlScalarFromVector` and other scalar plan reads. - -| Validation entry | Scope and stage | -|---|---| -| `OperatorNode::validate_structure()` | Walks the reachable operator DAG, including scalar plan references; checks input contracts, scalar typing and agreement between retained and derived output metadata. Valid for logical and physical plans; permits `timing = None`. | -| `OperatorNode::validate_execution_timing()` | Includes structural validation, then requires assigned timing on every executable operator and checks phase dependencies. Used for executable physical candidates. | -| Existing planner assessment and selection (#509) | Establishes guarantees using the existing accuracy models and checks them against request requirements and deployment capabilities. Neither node method re-proves a guarantee or decides request feasibility. | - -The two node methods need only the DAG and its annotations. Request requirements -and deployment models remain inputs to the existing planning/selection workflow, -not implicit globals of `validate_structure`. Passing the timing check alone does -not establish that a physical candidate satisfies the query's accuracy requirement. - -`Scan.schema` declares the source columns; `Values.schema` declares the constructed -row shape. `OperatorNode.schema` is the derived output for any operation. A scan's -predicates cannot change its declared output columns; a Values row must match the -declared arity, types and nullability. These leaf outputs retain the declaration's -column layout and time/identity information, with only justified metadata changes. -The declaration and derived output therefore have distinct roles, and structural -validation rejects disagreement rather than trusting two independent schemas. - -`scalar_type` keeps the existing method name and `(DataType, nullable)` result. -Its `input` is the applicable column scope: the child schema for a projection, -both input schemas for a join predicate, or aggregate outputs for `HAVING`. -Explicit subquery/conversion expressions validate their referenced producer using -the contracts above. Numeric expressions cannot consume state columns as numbers. -A standalone scalar expression is checked with an empty column scope and needs no fabricated -relation output schema. `QueryExprError` retains the existing error-type name; -result-kind, state-family, schema and execution-phase mismatches require -corresponding validation errors. - -For example, a KLL build outputs `State` with a -`Sketch(SketchKind, GroupingStrategy)` column identifying KLL and its parameters. -Its p99 readout outputs an ordinary `Plain(Float64)` column in the appropriate -relation/vector schema. A numeric predicate can use that readout, but not the KLL -state. Exact accumulator state similarly requires `FinalizeExactAccumulator`. -An ordinary operator may pass state through only where its input/output contract -permits it. A bare-column projection can preserve the field's `FieldDataType` -directly during `output_schema` derivation; `scalar_type` applies when that column -is used as a scalar value and rejects state. Copying a state column does not turn -it into a readable scalar. +Every operator output, before and after optimization, uses one `Schema` whose +`Field`s are typed by `FieldDataType`: `Plain(DataType)` for a readable value, or +the family, algorithm and parameters of summary or exact-accumulator state. +`OperatorResultKind` marks state outputs, and state becomes a value only through +an explicit readout. The schema model, `ColumnRef` versus `ColumnId`, the +validation entry points and the readout boundary are specified in +[Schema and physical data for ASAP primitives](asap-primitive-schema.md). ### 2.2 Preserve existing accuracy semantics diff --git a/docs/design_docs/proposals/univmon-frequency-summary.md b/docs/design_docs/proposals/univmon-frequency-summary.md index 75cfd080b..a2d9f96c1 100644 --- a/docs/design_docs/proposals/univmon-frequency-summary.md +++ b/docs/design_docs/proposals/univmon-frequency-summary.md @@ -35,7 +35,8 @@ cardinality alternatives, and exact count remains the cheaper first count candidate. All four readouts have the same unit-weight update, input sub-DAG, grouping, -window, parameter identity and state schema. Existing post-ASAP structural +window, parameter identity and state schema +([ASAP primitive schema](asap-primitive-schema.md)). Existing post-ASAP structural sharing can therefore intern their state producer while preserving distinct readout nodes. Sharing is only legal within the same execution/data scope. Precompute placement, SummaryCatalog installation, retention, and runtime diff --git a/docs/develop_docs/asap-aware-mapping-contracts.md b/docs/develop_docs/asap-aware-mapping-contracts.md index 447cb3823..0e16d4116 100644 --- a/docs/develop_docs/asap-aware-mapping-contracts.md +++ b/docs/develop_docs/asap-aware-mapping-contracts.md @@ -325,24 +325,12 @@ backend inspection. Automatic selection skips those unproven ratios. Use ### Family, category, algorithm, and parameters -Sketches separate their query category from the concrete algorithm and its parameters: +A summary's identity has four levels: family (`FieldDataType` variant), sketch +category (`SketchCategory`), algorithm (`SketchAlgorithm`), and the validated +committed choice (`SketchKind`). The levels and their validation are specified in +[Schema and physical data for ASAP primitives](../design_docs/proposals/asap-primitive-schema.md#3-proposed-schema-design). -| Level | Type | Example | -| --- | --- | --- | -| **family** | `SummaryFamilyType` | `Sketch`, `Sample`, `Wavelet`, `StatModel`, `ExactAggregate` | -| **category** | `SketchCategory` | `Quantile`, `Cardinality`, `Frequency`, `TopK` | -| **algorithm** | `SketchAlgorithm` | `Kll` / `DDSketch` (both quantile); `Hll` (HyperLogLog) / `Theta` / `Kmv` (K-Minimum Values), all cardinality | -| **committed choice** | `SketchKind` | one validated category + algorithm + parameter combination | - -A `SketchKind` is a validated committed choice. Its public constructor, -`SketchKind::new(algorithm, params)`, verifies that the parameter variant belongs -to the selected algorithm and classifies the pair into its category. The public -`.category()`, `.algorithm()`, and `.params()` accessors expose the committed -values without permitting an invalid combination. - -Where this matters in practice: `CostModel::rank_candidates`, `CostModel::size_params`, and `SketchAlgorithmStrategy::replacements` operate at the **algorithm** level. `summary_candidates(intent)` returns a list of `SketchAlgorithm`s (`[Kll, DDSketch]` for a `Quantile` intent), never a bare `SketchKind` with nothing chosen underneath it. `SketchKind` appears after an algorithm has been selected and sized—on `Realization::Sketch(SketchKind)` and `SummaryFamilyType::Sketch(SketchKind)`. - -`Sample`, `Wavelet`, and `StatModel` each use a flat `(Kind, Params)` pair. `Sketch` needs the additional algorithm level because multiple algorithms can serve the same purpose—for example, KLL and DDSketch both answer quantile queries. +Where this matters in practice: `CostModel::rank_candidates`, `CostModel::size_params`, and `SketchAlgorithmStrategy::replacements` operate at the **algorithm** level. `summary_candidates(intent)` returns a list of `SketchAlgorithm`s (`[Kll, DDSketch]` for a `Quantile` intent), never a bare `SketchKind` with nothing chosen underneath it. `SketchKind` appears after an algorithm has been selected and sized—on `Realization::Sketch(SketchKind)` and `FieldDataType::Sketch(SketchKind, GroupingStrategy)`. --- @@ -378,7 +366,7 @@ The crate provides no default `Matcher` implementation because the answer depend Concretely, `explanation.rs` reports three candidate kinds from each `TargetSubDAGCandidates`: -- `ExplanationKind::SketchApproximation` — the set contains a `Replacement::Summary` that realizes `SummaryFamilyType::Sketch(..)`, not just an exact/pass-through candidate. +- `ExplanationKind::SketchApproximation` — the set contains a `Replacement::Summary` that realizes `FieldDataType::Sketch(..)`, not just an exact/pass-through candidate. - `ExplanationKind::CommonSubexpressionReuse` — `consumer_count >= 2` and the set contains `SharedSubDAGStrategy`'s "build once and share" candidate (the `Replacement::Rewrite` whose `Rc` is the set's `target`). - `ExplanationKind::ExactComposition` — the candidate set contains an exact operation diff --git a/docs/develop_docs/pre-asap-ir.md b/docs/develop_docs/pre-asap-ir.md index abf5dc50b..77d80c847 100644 --- a/docs/develop_docs/pre-asap-ir.md +++ b/docs/develop_docs/pre-asap-ir.md @@ -17,26 +17,10 @@ The pre-ASAP IR is defined using the `QueryExpr` enum. We discuss some of import ## Fields and column references -`Schema` owns `Field` metadata: name, type, nullability, and an optional table -qualifier. A `Field` contains no runtime values. The former schema `Column` -struct served this same metadata role; it was renamed to `Field`, not retained -as a second data container. - -`ColumnRef` is an unresolved logical reference (`Named`, `Qualified`, -`SampleValue`, or `Wildcard`). Resolution binds a reference to `ColumnId`, a -`usize` position within a particular schema. `QueryExpr::Column(ColumnId)` -reads that column; the same position indexes `Schema::fields` for type checking -and a runtime row for its value. Group keys, unique keys, and `time_index` also -use these column positions. They are not stable identities across projections -or joins, so the positional reference remains `ColumnId`, not `FieldId`. - -The native runtime names shared ownership `SchemaRef = Arc` and stores -`Batch { schema: SchemaRef, rows: Vec> }`. `Schema` is the same metadata -model during planning and execution; the `Ref` suffix only distinguishes ownership. -It has no physical `Column`/array container. A column reference expresses what -to read independently of whether an executor stores its data as rows or arrays. -For example, resolving `t.bytes` to `ColumnId = 1` obtains its type from -`schema.fields[1]`; native execution reads `row[1]`. +`Schema` holds `Field` metadata (name, type, nullability, qualifier) and no +values; an unresolved `ColumnRef` resolves to a positional `ColumnId` within one +schema. The design, including how the same position selects a runtime value, is +in [Schema and physical data for ASAP primitives](../design_docs/proposals/asap-primitive-schema.md#3-proposed-schema-design). ## Node index