From 0fa8f47ae4358696a951fcdbd89b86f522081c2a Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 14:53:30 +0000 Subject: [PATCH] feat(runtime): execute exact frequency L2 and entropy --- .../src/query_physical_lowering.rs | 2 + .../src/operators/aggregate/mod.rs | 55 +++++++- .../src/physical_planner/mod.rs | 6 + .../tests/blocking_resources.rs | 43 +++++++ .../tests/physical_semantics.rs | 117 ++++++++++++++++++ docs/develop_docs/planner-layering-status.md | 47 +++++++ 6 files changed, 269 insertions(+), 1 deletion(-) create mode 100644 docs/develop_docs/planner-layering-status.md diff --git a/crates/asap-aware-mapping/src/query_physical_lowering.rs b/crates/asap-aware-mapping/src/query_physical_lowering.rs index 73c22b297..996898a94 100644 --- a/crates/asap-aware-mapping/src/query_physical_lowering.rs +++ b/crates/asap-aware-mapping/src/query_physical_lowering.rs @@ -1203,6 +1203,8 @@ fn supports_hash_aggregate( | AggIntent::PearsonCorr { .. } | AggIntent::Group | AggIntent::CountValues { .. } + | AggIntent::FrequencyL2 { .. } + | AggIntent::FrequencyEntropy { .. } ) }) } diff --git a/crates/asap-physical-operators/src/operators/aggregate/mod.rs b/crates/asap-physical-operators/src/operators/aggregate/mod.rs index 7b51b3c66..c9715a41b 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/mod.rs +++ b/crates/asap-physical-operators/src/operators/aggregate/mod.rs @@ -34,6 +34,15 @@ impl Operator { } (DataType::Float64, false) } + Reduction::FrequencyL2(i) | Reduction::FrequencyEntropy(i) => { + if !matches!( + plain(&input, *i)?.0, + DataType::Bool | DataType::Int64 | DataType::Float64 | DataType::Utf8 + ) { + return Err(invalid("frequency aggregate requires a Boolean, Int64, Float64 or Utf8 identity")); + } + (DataType::Float64, false) + } Reduction::Min(i) | Reduction::Max(i) => { let (t, nullable) = plain(&input, *i)?; if !ordered(t) { @@ -129,6 +138,10 @@ pub enum Reduction { Avg(usize), Min(usize), Max(usize), + /// L2 norm of unit-update frequencies; NULL identities are skipped. + FrequencyL2(usize), + /// Shannon entropy in bits; NULL identities are skipped. + FrequencyEntropy(usize), /// PromQL `quantile`: linear interpolation between closest ranks. Quantile { column: usize, @@ -205,7 +218,7 @@ async fn reduce( .map(|&i| rows[0][i].clone()) .collect::>(); for measure in measures { - result.push(reduce_one(&rows, measure, input, &mut work).await?); + result.push(reduce_one(&rows, measure, input, &mut work, context).await?); } workspace.grow(row_bytes(&result))?; output.push(result); @@ -243,6 +256,7 @@ async fn reduce_one( measure: &Reduction, input: &SchemaRef, work: &mut Cooperative, + context: &RunContext, ) -> Result { let column = match measure { Reduction::Count => { @@ -251,6 +265,45 @@ async fn reduce_one( )) } Reduction::Sum(i) | Reduction::Avg(i) | Reduction::Min(i) | Reduction::Max(i) => *i, + Reduction::FrequencyL2(column) | Reduction::FrequencyEntropy(column) => { + let mut workspace = Workspace::new(context)?; + let mut counts = BTreeMap::, u64>::new(); + let mut total = 0_u64; + for row in rows { + work.checkpoint().await?; + let value = &row[*column]; + if matches!(value, Value::Null) { + continue; + } + if matches!(value, Value::Float64(v) if !v.is_finite()) { + return Err(invalid( + "frequency aggregate requires a finite floating identity", + )); + } + let key = value.key()?; + if !counts.contains_key(&key) { + workspace.grow( + 64 + std::mem::size_of::>() + + key.len() + + std::mem::size_of::(), + )?; + } + *counts.entry(key).or_default() += 1; + total += 1; + } + let mut result = 0.0_f64; + for count in counts.into_values() { + work.checkpoint().await?; + let count = count as f64; + if matches!(measure, Reduction::FrequencyL2(_)) { + result = result.hypot(count); + } else { + let probability = count / total as f64; + result -= probability * probability.log2(); + } + } + return Ok(Value::Float64(result)); + } Reduction::Quantile { column, q } => { let mut values = Vec::with_capacity(rows.len()); for row in rows { diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 9e91a4bf4..59d734d8e 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -828,6 +828,12 @@ fn bind_operation(node: &PostAsapDAGNode, inputs: &[SchemaRef]) -> Result Reduction::Count, AggIntent::Sum { col } => Reduction::Sum(column(*col)?), AggIntent::Avg { col } => Reduction::Avg(column(*col)?), + AggIntent::FrequencyL2 { col, .. } => { + Reduction::FrequencyL2(column(*col)?) + } + AggIntent::FrequencyEntropy { col, .. } => { + Reduction::FrequencyEntropy(column(*col)?) + } AggIntent::Min { col } => Reduction::Min(column(*col)?), AggIntent::Max { col } => Reduction::Max(column(*col)?), _ => { diff --git a/crates/asap-physical-operators/tests/blocking_resources.rs b/crates/asap-physical-operators/tests/blocking_resources.rs index 66702afdb..579f7978d 100644 --- a/crates/asap-physical-operators/tests/blocking_resources.rs +++ b/crates/asap-physical-operators/tests/blocking_resources.rs @@ -248,3 +248,46 @@ fn weighted_summary_build_yields_within_a_batch() { drop(output); assert_eq!(run.retained_bytes(), 0); } + +// Frequency dictionaries count against the budget and release reservations on failure. +#[test] +fn frequency_dictionary_enforces_memory_budget() { + use asap_physical_operators::operators::Reduction; + let input = schema(1); + let mut sources = PhysicalDAG::default(); + sources + .add( + 0, + vec![], + Operator::source( + input.clone(), + vec![Batch::try_new( + input.clone(), + (0..64).map(|i| vec![Value::Int64(i)]).collect(), + ) + .unwrap()], + ) + .unwrap(), + ) + .unwrap(); + for (reduction, succeeds) in [ + (Reduction::Count, true), + (Reduction::FrequencyL2(0), false), + (Reduction::FrequencyEntropy(0), false), + ] { + let run = context(12_000); + let inputs = sources.execute(&[0], run.clone()).unwrap(); + let operator = + Operator::aggregate(input.clone(), vec![], vec![("result".into(), reduction)]).unwrap(); + let mut output = operator.start(inputs, run.clone()).unwrap(); + let result = block_on(output.next()).unwrap(); + if succeeds { + assert!(result.is_ok()); + } else { + assert!(matches!(result, Err(Error::MemoryLimit))); + } + drop(result); + drop(output); + assert_eq!(run.retained_bytes(), 0); + } +} diff --git a/crates/asap-physical-operators/tests/physical_semantics.rs b/crates/asap-physical-operators/tests/physical_semantics.rs index 5c58d5644..5939be253 100644 --- a/crates/asap-physical-operators/tests/physical_semantics.rs +++ b/crates/asap-physical-operators/tests/physical_semantics.rs @@ -714,3 +714,120 @@ fn empty_exact_summary_extrema_agree_with_ordinary_aggregation() { assert!(matches!(rows[0][0], Value::Null)); } } + +// Exact frequency intents bind to native reducers without a sketch or numeric key conversion. +#[test] +fn exact_frequency_intents_execute_typed_keys_and_empty_input() { + use asap_physical_operators::physical_planner::compile_node; + use planner_types::{ + post_asap::ExecutionDataState, + pre_asap::{AggIntent, GroupKeys, Reduction as PlanReduction}, + types::AccuracyTarget, + }; + for (dtype, values) in [ + ( + DataType::Utf8, + vec![Value::Utf8("a".into()), Value::Utf8("b".into())], + ), + ( + DataType::Int64, + vec![ + Value::Int64(9_007_199_254_740_992), + Value::Int64(9_007_199_254_740_993), + ], + ), + (DataType::Bool, vec![Value::Bool(false), Value::Bool(true)]), + ( + DataType::Float64, + vec![Value::Float64(-0.0), Value::Float64(1.0)], + ), + ] { + let input = schema(&[("key", dtype, true)]); + for (measure, name, expected) in [ + ( + AggIntent::FrequencyL2 { + col: Some(0), + accuracy: AccuracyTarget::Exact, + }, + "frequency_l2", + 8.0_f64.sqrt(), + ), + ( + AggIntent::FrequencyEntropy { + col: Some(0), + accuracy: AccuracyTarget::Exact, + }, + "frequency_entropy", + 1.0, + ), + ] { + let node = PostAsapDAGNode { + id: PostAsapNodeId(1), + payload: PostAsapOperatorPayload::Relational { + operator: ValueOperation::Aggregate { + reduction: PlanReduction::Reduce(GroupKeys::none()), + measures: vec![measure], + output_names: vec![name.into()], + filters: vec![], + having: None, + }, + }, + output_state: ExecutionDataState::QUERY_ROWS, + output_schema: (*schema(&[(name, DataType::Float64, false)])).clone(), + guarantee: None, + }; + let operator = compile_node(&node, std::slice::from_ref(&input)) + .expect("exact frequency intent binds"); + let rows = values + .iter() + .flat_map(|v| [vec![v.clone()], vec![v.clone()]]) + .chain([vec![Value::Null]]) + .collect(); + let result = unary(input.clone(), vec![rows], operator.clone()); + assert!(matches!(result[0][0], Value::Float64(v) if (v - expected).abs() < 1e-12)); + for batches in [vec![], vec![vec![vec![Value::Null]]]] { + let result = unary(input.clone(), batches, operator.clone()); + assert!(matches!(result[0][0], Value::Float64(0.0))); + } + } + } +} + +// Each group gets its own frequency population, including one canonical signed-zero identity. +#[test] +fn exact_frequency_grouping_and_entropy_bits() { + let input = schema(&[ + ("group", DataType::Int64, false), + ("key", DataType::Float64, true), + ]); + let operator = Operator::aggregate( + input.clone(), + vec![0], + vec![ + ("l2".into(), Reduction::FrequencyL2(1)), + ("entropy".into(), Reduction::FrequencyEntropy(1)), + ], + ) + .unwrap(); + let rows = vec![ + vec![Value::Int64(1), Value::Float64(-0.0)], + vec![Value::Int64(1), Value::Float64(0.0)], + vec![Value::Int64(1), Value::Float64(0.0)], + vec![Value::Int64(1), Value::Float64(1.0)], + vec![Value::Int64(2), Value::Float64(2.0)], + vec![Value::Int64(2), Value::Null], + vec![Value::Int64(3), Value::Null], + ]; + let result = unary(input.clone(), vec![rows], operator.clone()); + assert_eq!(result.len(), 3); + let expected_entropy = -0.75_f64 * 0.75_f64.log2() - 0.25_f64 * 0.25_f64.log2(); + for (row, l2, entropy) in [ + (&result[0], 10.0_f64.sqrt(), expected_entropy), + (&result[1], 1.0, 0.0), + (&result[2], 0.0, 0.0), + ] { + assert!(matches!(row[1], Value::Float64(v) if (v - l2).abs() < 1e-12)); + assert!(matches!(row[2], Value::Float64(v) if (v - entropy).abs() < 1e-12)); + } + assert!(unary(input, vec![], operator).is_empty()); +} diff --git a/docs/develop_docs/planner-layering-status.md b/docs/develop_docs/planner-layering-status.md new file mode 100644 index 000000000..131f2b124 --- /dev/null +++ b/docs/develop_docs/planner-layering-status.md @@ -0,0 +1,47 @@ +# Planner-layering implementation status + +Audience: planner developers. Audit baseline: PR #557 (`e0e1e2d7`), against +[the #509 proposal](../design_docs/proposals/planner-layering.md). The proposal +is a target contract, not a statement that its examples execute today. + +| Proposal contract | Evidence at #557 | Remaining scope | +| --- | --- | --- | +| Language frontends and common logical IR | SQL/PromQL/MetricsQL lower to unified operators and scalars. | Example 2 SQL frequency L2 and entropy idioms are not recognized. Preserve alias lineage, filters, NULL groups, empty inputs, count overflow and entropy units when adding recognition. | +| Local exact and summary alternatives | `replacement::summary_candidates`, realization rules and candidate inventory exist; supplied accuracy models reach Pass 1. | Specialized entropy/norm families in Example 2 are illustrative, not registered families. UnivMon certifies only unit-update total count; L2, entropy and cardinality epsilon/delta bounds need verified evidence or a deployment model. | +| Summary-capability sharing | CSE interns structurally identical producers, including states with different readers. | It does not enumerate all partial sharing partitions or resize compatible states to the strictest consumer. Example 2's 37 candidates are not an acceptance result. | +| Window composition | Mergeable state IR/native merge exists; physical pane compatibility and reuse cost helpers exist. | Automatic logical sliding/tumbling/EH alternatives over differing windows, boundary coverage and error proofs are absent. A merge kernel alone does not implement Examples 1/3. | +| Physical materialization | Ephemeral/prepared/shared/continuously maintained lifecycle alternatives, costing, capabilities and latency checks exist. | Incremental query-time pane retention, historical backfill and the complete Example 4 matrix need executable implementations and explicit state/input contracts. | +| Whole-workload selection | One unified selected DAG; shared states are interned and costed across their consumers. | `replacement.rs` documents its selection as non-exhaustive over interacting choices. The proposal's cheapest complete candidate guarantee and 54/156 inventories need a complete workload search/selection path. | +| Deployment inputs and execution | `PlanningModels` bundles cost, accuracy, evidence and capabilities; native typed UnivMon supports one build with three readouts. | At #557 native exact frequency L2/entropy fallback is absent. End-to-end SQL Example 2 is not established by the native UnivMon fixture. | +| Subtract/delete, parallelism, partitioning and resource planning | Some runtime memory/cancellation limits and maintenance capability flags exist. | These remain proposal TODOs; capability flags do not supply missing IR operators or a physical resource search. | + +## Follow-up sequence + +1. **Exact frequency execution.** Bind the existing L2/entropy intents to native + reducers with typed identities, bits for entropy, NULL skipping, zero for an + empty population, grouping, memory accounting and cooperative cancellation. + This change supplies the fallback prerequisite; it does not certify UnivMon. +2. **SQL frequency recognition.** Add narrow, proven idiom recognition while + retaining the original exact relational computation. Entropy with `LN` is in + nats; the core intent is in bits. SQL `COUNT(*) GROUP BY nullable_key` counts + a NULL group, while the frequency intents skip NULL. Do not erase these + differences or SQL's NULL result for `SUM` over no groups. +3. **Summary sharing alternatives.** Enumerate compatible consumer partitions + and size shared states against all consumers using the supplied accuracy + model. Keep independent candidates. Tests must inspect the inventory and + selected producer count, rather than manually construct a shared state. +4. **Tumbling-window composition.** Start with aligned, fixed windows and + mergeable KLL, then add query-time retention and ingestion-time lifecycle + choices. Validate offsets, boundary alignment, retention and raw fallback. +5. **Sliding/EH composition.** Separate changes for overlapping active windows + and historical bucket coverage/error. Do not reuse a merge guarantee as a + boundary-error guarantee. +6. **Complete physical candidate selection.** Explore interacting workload + choices and lifecycle assignments, charge each shared producer once and + check all consumers' accuracy/latency/capabilities. A bounded exhaustive + implementation must fail explicitly on its budget instead of truncate. +7. **Proposal TODOs.** Design subtract/delete and parallelism/partitioning/ + resource inputs before implementing their planning choices. + +The examples' numerical candidate counts depend on their stated rule sets. +Tests should establish those rule sets explicitly before asserting the counts.