Skip to content

feat(runtime): accept typed UnivMon frequency identities - #557

Draft
zzylol wants to merge 1 commit into
stack/528-14-raw-latencyfrom
stack/528-15-frequency-inputs
Draft

zzylol wants to merge 1 commit into
stack/528-14-raw-latencyfrom
stack/528-15-frequency-inputs

Conversation

@zzylol

@zzylol zzylol commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

Problem: the native UnivMon build cannot consume the src_ip identities that #509 Example 2 shares

#509 §Pass 2 (summary-capability rule) and §Example 2 ("One summary for several computations") describe this shared candidate:

One UnivMon over src_ip from flows in the last minute serves three queries refreshed every 10 s: COUNT(DISTINCT src_ip), the entropy of the src_ip distribution, and the L2 norm of per-src_ip counts. Each flow record updates the UnivMon once; a distinct-count, an entropy and an L2 estimation node each compute their statistic from it.

#509 §Physical operator implementation then turns that node into summary build and estimation operators. The build operator must therefore accept a src_ip value as a UnivMon update.

Before this PR, the native summary build accepts only Float64 update columns, for every family:

// crates/asap-physical-operators/src/operators/summary/mod.rs (before)
if plain(&input, value)?.0 != &DataType::Float64 {
    return Err(invalid("summary numeric update requires Float64"));
}

So the Example 2 input fails before any row is read:

flows(src_ip: Utf8)  →  Operator::summary_build(input, UnivMon{..}, value = 0, None, [])
                     →  Err("summary numeric update requires Float64")

Converting the identity to Float64 first is not a fix. Int64 identities above 2^53 collapse: 9_007_199_254_740_992 and 9_007_199_254_740_993 become the same Float64, so the distinct count, L2 and entropy all change.

Scope. This PR covers the physical build side of the Example 2 shared candidate: one native UnivMon state built from typed src_ip-like identities, read by the distinct, L2 and entropy estimators. It does not cover the frontend recognition of the Q2/Q3 SQL forms (#509 §Example 2 TODO), the Pass 2 sharing rule itself, sizing for the strictest ε, or UnivMon accuracy certification.

Proposed method

All changes are in the physical runtime (asap-physical-operators). Planning stages are unchanged.

  1. Build-time type check. Operator::summary_build still requires Float64 for every family, except UnivMon. A UnivMon build also accepts Utf8, Int64 and Bool update columns. Any other type returns "summary update type is unsupported by its family".
  2. Typed update path. The build loop no longer unwraps the row value to f64. It skips Value::Null (as before: SQL aggregates ignore NULL samples while keeping the group) and passes the typed &Value to a new trait method AccumulatorUpdater::update_value.
  3. Default keeps the numeric contract. The default update_value requires Value::Float64, runs validate_single_input, then update_single. Every non-UnivMon family therefore behaves exactly as before.
  4. UnivMon override. UnivMonUpdater::update_value calls the new UnivMonAccumulator::insert_value:
    • Null → no-op.
    • finite Float64 → the existing insert_sample path (unchanged Float64 behavior).
    • Bool / Int64 / Utf8 → encode with the existing type-tagged Value::key() bytes, check the bucket counter for overflow, and insert as DataInput::Bytes with weight 1.
    • anything else (including non-finite Float64) → error.
      Because Value::key() stores Int64 as its 8 little-endian bytes with a type tag, neighboring large integers stay distinct.
  5. Memory accounting. UnivMonAccumulator::approx_memory_bytes now adds the capacity of String/Bytes keys held in each layer's heavy-hitter heap. Long string identities count against the runtime memory reservation. After clear(), the value returns to the empty-state size.

One built state still feeds the existing SketchStatistic::Cardinality, FrequencyL2 and FrequencyEntropy evaluations.

Key code interfaces

crates/asap-physical-operators/src/summary_kernels/factory.rs

pub trait AccumulatorUpdater: Send {
    /// Update from the typed native row. Numeric kernels retain their existing
    /// Float64 contract; frequency kernels may accept nonnumeric identities.
    fn update_value(
        &mut self,
        value: &crate::values::Value,
        timestamp_ms: i64,
    ) -> Result<(), String>;      // default: Float64 only → validate_single_input + update_single
    // existing methods unchanged
}

impl AccumulatorUpdater for UnivMonUpdater {
    fn update_value(&mut self, value: &crate::values::Value, _: i64) -> Result<(), String>;
}

crates/asap-physical-operators/src/summary_kernels/univmon.rs

impl UnivMonAccumulator {
    /// SQL identities are kept in their original type: converting Int64 to
    /// Float64 would collapse neighboring keys above 2^53.
    pub fn insert_value(&mut self, value: &crate::values::Value) -> Result<(), Error>;
}

crates/asap-physical-operators/src/operators/summary/mod.rs (signature unchanged, accepted types widened)

pub fn summary_build(
    input: SchemaRef,
    family: SummaryFamilyType,
    value: usize,
    time: Option<usize>,
    groups: Vec<usize>,
) -> Result<Self, Error>;

Usage, as in tests/univmon_execution.rs:

let input = Arc::new(Schema::new(vec![Field::plain("src_ip", DataType::Utf8, true)]));
let build = Operator::summary_build(input, family(), 0, None, vec![])?; // family() = UnivMon
// node 1 = build; nodes 2,3,4 = Operator::evaluation(output, 0, SummaryEvaluation::Sketch(stat))
// for stat in [Cardinality, FrequencyL2, FrequencyEntropy]

Fields

AccumulatorUpdater::update_value

Parameter Type Meaning
self &mut self The per-group summary updater being built.
value &Value The typed update value from the input row. The build loop never passes Null; it skips NULL rows before the call.
timestamp_ms i64 Row timestamp from the time column, or 0 when time is None. Passed to update_single by the default; ignored by UnivMon.
return Result<(), String> Err for a type the family does not accept or for an input validate_single_input rejects. The build maps it to Error::Operator.

UnivMonAccumulator::insert_value

Input value Effect
Value::Null Ignored, Ok(()).
Value::Float64(x), x finite insert_sample(x), same as before this PR.
Value::Bool / Value::Int64 / Value::Utf8 Inserted once as type-tagged key bytes (Value::key()), weight 1. Errors with "UnivMon count overflow" if bucket_size + 1 overflows.
other, or non-finite Float64 Err("unsupported UnivMon identity type or nonfinite sample").

Operator::summary_build

Parameter Type Meaning / allowed values
input SchemaRef Input row schema.
family SummaryFamilyType Summary family and parameters. Validated by validate_family. UnivMon is detected with kind.algorithm() == SketchAlgorithm::UnivMon.
value usize Index of the update column. Must be plain. Float64 for all families; for UnivMon also Utf8, Int64 or Bool.
time Option<usize> Optional timestamp column; must be plain, non-null Timestamp (unchanged).
groups Vec<usize> Grouping columns; validated by validate_groups (unchanged).

Examples

End to end (tests/univmon_execution.rs::typed_frequency_keys_preserve_identity). For each key type, the input column src_ip holds two keys, each twice, and one NULL:

src_ip: k0, k0, k1, k1, NULL
PhysicalDAG: 0 source → 1 summary_build(UnivMon) → {2 Cardinality, 3 FrequencyL2, 4 FrequencyEntropy}
Key type k0, k1 Distinct L2 Entropy (bits)
Utf8 "192.0.2.1", "192.0.2.2" 2 √8 1
Int64 9_007_199_254_740_992, 9_007_199_254_740_993 2 √8 1
Bool false, true 2 √8 1

All values are checked within 0.01. The NULL row contributes nothing. The two Int64 keys are 2^53 and 2^53+1, which would be one key after a Float64 cast.

Kernel tests (src/summary_kernels/univmon.rs).

  • typed_keys_merge_and_roundtrip: two Utf8 states (192.0.2.1 ×2 and 192.0.2.2 ×2) are merged, serialized and restored. The restored state still gives distinct 2, L2 √8, entropy 1.
  • memory_accounts_for_string_identities: inserting one 4,096-byte Utf8 key raises approx_memory_bytes by at least 4,096. After clear() it equals the empty size.

Accepted vs rejected update columns

Family Column type Result
UnivMon Float64 accepted (unchanged path)
UnivMon Utf8, Int64, Bool accepted (new, typed keys)
UnivMon other types rejected at construction
UnivMon non-finite Float64 value rejected at update
any other family Float64 accepted (unchanged)
any other family non-Float64 rejected at construction, and update_value default rejects non-Float64 values

The developer note in docs/develop_docs/operator-design-acceptance.md §"Planner-layering follow-up: typed frequency inputs" records the same contract.

Out of scope

Stack and validation

Stacked on #556 · Next: #559 · Reference/tracker: #528

Validation: the typed-input regression fails before the change and passes afterward; all 225 native runtime tests/doctests passed; formatting; all-target Clippy for native runtime with warnings denied.

The full restacked tip passes all 1,608 workspace tests/doctests (two existing ignores), formatting and workspace/all-target/all-feature Clippy with warnings denied. The inherited UnivMon precompute fixture was fixed in #552, and the latency fixture lint in #554; later branches were restacked atomically with explicit leases. CI has been retriggered.

🤖 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

zzylol commented Oct 4, 2026

Copy link
Copy Markdown
Contributor Author

Ported onto the current stack in #597

zzylol added a commit that referenced this pull request Oct 4, 2026
…ntities

Port of the parked runtime PRs for #509 Example 2 onto asap-executor:

- native UnivMon build and Cardinality / FrequencyL2 / FrequencyEntropy
  readouts (the old stack's 14c8ac8, a prerequisite missing here);
- #557: a UnivMon build accepts Utf8, Int64 and Bool identities as
  type-tagged keys; heap key bytes count toward state memory;
- #559: exact FrequencyL2 / FrequencyEntropy reducers (bits, NULL skipped,
  0 for an empty population, memory-accounted, cooperative);
- #563: exact Cardinality reducer over typed tuples.

The legacy raw-cost adapter lists the three intents as hash aggregates.

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