Skip to content

feat(runtime): execute UnivMon states and readouts - #552

Draft
zzylol wants to merge 2 commits into
stack/528-09-accuracyfrom
stack/528-10-univmon
Draft

zzylol wants to merge 2 commits into
stack/528-09-accuracyfrom
stack/528-10-univmon

Conversation

@zzylol

@zzylol zzylol commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

Problem: a plan with one shared UnivMon state cannot run on the native runtime

#509 Example 2 ("One summary for several computations — the summary-capability rule in Pass 2") selects one UnivMon build node feeding three estimation nodes: distinct count, entropy and L2 of src_ip. #509 §Physical operator implementation then "converts every node to physical operators", for example "a KLL node as summary build, merge and quantile estimation operators". #509 §Scenarios, "Adding a new summary family", lists what a family must provide, including "6. Its physical kernel: build, merge and estimate". #509 §4 Execution says the deployment "runs the selected plan as given".

The UnivMon kernel already existed and already answered these statistics (summary_kernels/univmon.rs):

fn estimate(&self, query: &SketchStatistic) -> Result<f64, Error> {
    Ok(match query {
        SketchStatistic::PointCount { key: ColumnRef::SampleValue, value: None } => self.inner.calc_l1(),
        SketchStatistic::Cardinality => self.inner.calc_card(),
        SketchStatistic::FrequencyL2 => self.inner.calc_l2(),
        SketchStatistic::FrequencyEntropy => self.inner.calc_entropy(),
        other => return Err(format!("UnivMon does not answer {other:?}").into()),
    })
}

But the native DAG gates did not let a UnivMon state in. Before this PR, in capability.rs:

// validate_native_family
SummaryFamilyType::Sketch(kind, _)
    if matches!(kind.algorithm(), A::Kll | A::DDSketch | A::Hll) => {}
_ => return Err(Error::Invalid("summary family has no native DAG state implementation".into())),

Every native summary constructor calls this gate (Operator::summary_build, summary_merge, evaluation, through values::validate_family), and so does Batch::try_new for summary values. So the physical DAG for Example 2 failed at its first summary node:

Source(value: Float64)
└─ SummaryBuild(UnivMon)            ← Error::Invalid("summary family has no native DAG state implementation")
   ├─ Evaluation(Cardinality)
   ├─ Evaluation(FrequencyL2)
   └─ Evaluation(FrequencyEntropy)

The precompute integration test also treated UnivMon as a family "without a native state" and expected its compile to fail.

Scope. This PR covers the physical part of #509 Example 2: native state validation and the native count, distinct, L2 and entropy evaluations for UnivMon. Planning the shared UnivMon candidate is #551. It does not add a UnivMon accuracy model or change planning.

Proposed method

All changes are in the physical runtime (asap-physical-operators). No new operators or kernels.

  1. Accept the family. validate_native_family adds A::UnivMon to the families with a native DAG state. It still runs validate_summary_kernel afterwards, as for the other families. Because the build, merge and evaluation constructors all call this gate, a UnivMon state can now be built, merged and evaluated in a PhysicalDAG.
  2. Accept four evaluations. validate_sketch_evaluation accepts, for UnivMon, the four statistics the kernel answers: a bare PointCount (value: None), Cardinality, FrequencyL2 and FrequencyEntropy. Everything else (Quantile, TopK, a PointCount with value: Some(..)) is still rejected with "evaluation is not implemented for this summary family". UnivMon results are Float64; only CMS bare counts are typed Int64.
  3. Check state shape. values::validate_state (used by Batch::try_new when a row holds a Value::Summary) now downcasts a UnivMon state to UnivMonAccumulator and requires its dimensions to equal the declared SketchParams::UnivMon. A state with a different shape is rejected, as for KLL, DDSketch, HLL and CMS.

Key code interfaces

Signatures are unchanged; the accepted sets grow. crates/asap-physical-operators/src/capability.rs:

pub fn validate_native_family(family: &SummaryFamilyType) -> Result<(), Error>;
//   now accepts: Sketch(kind, _) if kind.algorithm() is Kll | DDSketch | Hll | UnivMon

pub fn validate_sketch_evaluation(
    family: &SummaryFamilyType,
    query: &SketchStatistic,
) -> Result<(), Error>;
//   new arms:
//   (A::UnivMon, SketchStatistic::PointCount { value: None, .. })
//   | (A::UnivMon, SketchStatistic::Cardinality)
//   | (A::UnivMon, SketchStatistic::FrequencyL2)
//   | (A::UnivMon, SketchStatistic::FrequencyEntropy) => true,

crates/asap-physical-operators/src/values.rs (crate-private):

fn validate_state(family: &SummaryFamilyType, state: &dyn AggregateCore) -> Result<(), Error>;
//   new arm:
SketchParams::UnivMon { heap_size, sketch_rows, sketch_cols, layers } => state
    .as_any()
    .downcast_ref::<UnivMonAccumulator>()
    .is_some_and(|s| {
        let sketch = s.sketch();
        sketch.heap_size == *heap_size as usize
            && sketch.sketch_row == *sketch_rows as usize
            && sketch.sketch_col == *sketch_cols as usize
            && sketch.layer_size == *layers as usize
    }),

Existing types these rely on (unchanged):

// crates/types/src/post_asap/sketch.rs
SketchParams::UnivMon { heap_size: u32, sketch_rows: u32, sketch_cols: u32, layers: u8 }

pub enum SketchStatistic {
    FrequencyL2,
    FrequencyEntropy,
    Quantile { q: f64 },
    PointCount { key: ColumnRef, value: Option<String> },
    Cardinality,
    TopK { k: usize },
}

// crates/asap-physical-operators/src/summary_kernels/univmon.rs
impl UnivMonAccumulator { pub fn sketch(&self) -> &UnivMon; }

Usage, from the new test:

let state = Operator::summary_build(input, family(), 0, None, vec![])?;
dag.add(1, vec![0], state)?;
dag.add(2, vec![1], Operator::evaluation(state_schema.clone(), 0,
    SummaryEvaluation::Sketch(SketchStatistic::Cardinality))?)?;

Fields

validate_native_family:

Name Type Meaning
family &SummaryFamilyType The state type of a summary column. Ok means the native runtime has a state for it. Now also Ok for Sketch with algorithm UnivMon, provided validate_summary_kernel accepts its parameters.

validate_sketch_evaluation:

Name Type Meaning
family &SummaryFamilyType State type being read. Must first pass validate_native_family.
query &SketchStatistic The statistic to compute. For UnivMon, allowed values are listed below.

SketchStatistic variants accepted for UnivMon:

Variant Meaning Kernel call
PointCount { value: None, .. } Total number of updates (bare count). calc_l1()
Cardinality Number of distinct values. calc_card()
FrequencyL2 sqrt(Σ_v frequency(v)²) of the value frequencies. calc_l2()
FrequencyEntropy Shannon entropy of the value-frequency distribution, in bits. calc_entropy()

validate_state:

Name Type Meaning
family &SummaryFamilyType Declared state type of the column.
state &dyn AggregateCore The runtime state in the batch row.

SketchParams::UnivMon fields and the UnivMon fields they must equal:

Declared field Type Runtime field (UnivMonAccumulator::sketch()) Meaning
heap_size u32 heap_size Size of each layer's heavy-hitter heap.
sketch_rows u32 sketch_row Rows of each layer's count sketch.
sketch_cols u32 sketch_col Columns of each layer's count sketch.
layers u8 layer_size Number of sampling layers.

UnivMonAccumulator::sketch() returns the wrapped asap_sketchlib::UnivMon; it is used here only to read those dimensions.

Examples

One state, three evaluations

Test: one_univmon_state_answers_distinct_l2_and_entropy in crates/asap-physical-operators/tests/univmon_execution.rs.

Input. One batch with one non-null Float64 column value: [1.0, 1.0, 2.0, 2.0, 2.0]. Family Sketch(UnivMon, { heap_size: 64, sketch_rows: 5, sketch_cols: 128, layers: 4 }), default grouping.

DAG.

0  Source(batch)
1  SummaryBuild(UnivMon)            children [0]
2  Evaluation(Cardinality)          children [1]
3  Evaluation(FrequencyL2)          children [1]
4  Evaluation(FrequencyEntropy)     children [1]

Executed with dag.execute(&[2, 3, 4], ...) in Scope::Query { evaluation_time_ms: 0, revision: 1 }. Node 1 is built once and read by all three evaluations.

Output. Frequencies are {1.0: 2, 2.0: 3}. The test checks:

Sink Exact value Test assertion
2 Cardinality 2 `
3 FrequencyL2 sqrt(2² + 3²) = sqrt(13) `
4 FrequencyEntropy about 0.971 bits result > 0.0

Accepted vs rejected

Case Before After
summary_build / summary_merge with UnivMon family rejected: no native DAG state accepted
evaluation UnivMon + Cardinality / FrequencyL2 / FrequencyEntropy rejected (family gate) accepted, Float64
evaluation UnivMon + bare PointCount rejected (family gate) accepted, Float64
evaluation UnivMon + Quantile { q } rejected rejected: evaluation not implemented
Batch row with a UnivMon state whose dimensions differ from SketchParams::UnivMon rejected (family gate) rejected: value differs from type
Batch row with a matching UnivMon state rejected (family gate) accepted

Precompute integration test

crates/integration-tests/tests/precompute_raw_samples.rs no longer lists UnivMon as a family without native state. A UnivMon raw-sample summary now goes through the compile-and-compare path: its Cardinality, FrequencyL2 and FrequencyEntropy estimates are compared with the same kernel fed sample by sample. UnivMon is excluded from the keyed-heap per-item check, as HLL is.

Out of scope

  • An accuracy model for UnivMon readouts (fix(planner): pass deployment accuracy into Pass 1 #551 uses a synthetic test model; the built-in model still certifies only the bare count).
  • A merge test for UnivMon states. The merge constructor now accepts UnivMon through the shared gate, but this PR adds no merge test.
  • Note: the PointCount arm accepts any key with value: None. The kernel answers only key: ColumnRef::SampleValue; another key passes validation and fails at execution with "UnivMon does not answer".

Stack and validation

Stack 10 · Base: #551 · Next: #553 · Closes #524.

Validation: CARGO_TARGET_DIR=/mydata/cargo-target-412 cargo test --locked -p asap-physical-operators --test univmon_execution.

🤖 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